refactor mqtt server
This commit is contained in:
parent
58f0875197
commit
d0a53d5c9a
1 changed files with 43 additions and 97 deletions
|
|
@ -2,7 +2,6 @@
|
|||
|
||||
import asyncio
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
from datetime import datetime, timedelta
|
||||
|
|
@ -94,27 +93,25 @@ class MQTTHelperBot:
|
|||
requestid,
|
||||
cmdjson["payloadType"],
|
||||
)
|
||||
try:
|
||||
if cmdjson["payloadType"] == "x":
|
||||
await self.Client.publish(
|
||||
ttopic, str(cmdjson["payload"]).encode(), QOS_0
|
||||
)
|
||||
|
||||
if cmdjson["payloadType"] == "j":
|
||||
await self.Client.publish(
|
||||
ttopic, json.dumps(cmdjson["payload"]).encode(), QOS_0
|
||||
)
|
||||
|
||||
except Exception as e:
|
||||
helperbotlog.exception("{}".format(e))
|
||||
if cmdjson["payloadType"] == "x":
|
||||
await self.Client.publish(
|
||||
ttopic, str(cmdjson["payload"]).encode(), QOS_0
|
||||
)
|
||||
elif cmdjson["payloadType"] == "j":
|
||||
await self.Client.publish(
|
||||
ttopic, json.dumps(cmdjson["payload"]).encode(), QOS_0
|
||||
)
|
||||
|
||||
resp = await self.wait_for_resp(requestid)
|
||||
|
||||
return resp
|
||||
|
||||
except Exception as e:
|
||||
helperbotlog.exception("{}".format(e))
|
||||
return {}
|
||||
return {
|
||||
"id": requestid,
|
||||
"errno": 500,
|
||||
"ret": "fail",
|
||||
"debug": "exception occurred please check bumper logs",
|
||||
}
|
||||
|
||||
|
||||
class MQTTServer:
|
||||
|
|
@ -229,10 +226,9 @@ class BumperMQTTServer_Plugin:
|
|||
tmpbotdetail[1],
|
||||
"eco-ng",
|
||||
)
|
||||
mqttserverlog.info(
|
||||
f"Bumper Authentication Success - Bot - SN: {username} - DID: {didsplit[0]} - Class: {tmpbotdetail[0]}")
|
||||
mqttserverlog.info(f"Bumper Authentication Success - Bot - SN: {username} - DID: {didsplit[0]}"
|
||||
f" - Class: {tmpbotdetail[0]}")
|
||||
authenticated = True
|
||||
|
||||
else:
|
||||
tmpclientdetail = str(didsplit[1]).split("/")
|
||||
userid = didsplit[0]
|
||||
|
|
@ -242,21 +238,11 @@ class BumperMQTTServer_Plugin:
|
|||
if userid == "helperbot":
|
||||
mqttserverlog.info(f"Bumper Authentication Success - Helperbot: {client_id}")
|
||||
authenticated = True
|
||||
else:
|
||||
auth = False
|
||||
if bumper.check_authcode(didsplit[0], password):
|
||||
auth = True
|
||||
elif bumper.use_auth == False:
|
||||
auth = True
|
||||
|
||||
if auth:
|
||||
bumper.client_add(userid, realm, resource)
|
||||
mqttserverlog.info(
|
||||
f"Bumper Authentication Success - Client - Username: {username} - ClientID: {client_id}")
|
||||
authenticated = True
|
||||
|
||||
else:
|
||||
authenticated = False
|
||||
elif bumper.check_authcode(didsplit[0], password) or not bumper.use_auth:
|
||||
bumper.client_add(userid, realm, resource)
|
||||
mqttserverlog.info(f"Bumper Authentication Success - Client - Username: {username} - "
|
||||
f"ClientID: {client_id}")
|
||||
authenticated = True
|
||||
|
||||
# Check for File Auth
|
||||
if username and not authenticated: # If there is a username and it isn't already authenticated
|
||||
|
|
@ -279,9 +265,7 @@ class BumperMQTTServer_Plugin:
|
|||
authenticated = False
|
||||
|
||||
# Check for allow anonymous
|
||||
allow_anonymous = self.auth_config.get(
|
||||
"allow-anonymous", True
|
||||
)
|
||||
allow_anonymous = self.auth_config.get("allow-anonymous", True)
|
||||
if allow_anonymous and not authenticated: # If anonymous auth is allowed and it isn't already authenticated
|
||||
authenticated = True
|
||||
self.context.logger.debug(
|
||||
|
|
@ -308,93 +292,55 @@ class BumperMQTTServer_Plugin:
|
|||
self.context.logger.warning(f"Password file {password_file} not found")
|
||||
|
||||
async def on_broker_client_connected(self, client_id):
|
||||
self._set_client_connected(client_id, True)
|
||||
|
||||
def _set_client_connected(self, client_id, connected: bool):
|
||||
didsplit = str(client_id).split("@")
|
||||
|
||||
bot = bumper.bot_get(didsplit[0])
|
||||
if bot:
|
||||
bumper.bot_set_mqtt(bot["did"], True)
|
||||
bumper.bot_set_mqtt(bot["did"], connected)
|
||||
return
|
||||
|
||||
clientresource = didsplit[1].split("/")[1]
|
||||
client = bumper.client_get(clientresource)
|
||||
if client:
|
||||
bumper.client_set_mqtt(client["resource"], True)
|
||||
return
|
||||
bumper.client_set_mqtt(client["resource"], connected)
|
||||
|
||||
async def on_broker_message_received(self, client_id, message):
|
||||
self.handle_helperbot_msg(client_id, message)
|
||||
|
||||
def handle_helperbot_msg(self, client_id, message):
|
||||
|
||||
if str(message.topic).split("/")[6] == "helperbot":
|
||||
topic = message.topic
|
||||
topic_split = str(topic).split("/")
|
||||
data_decoded = str(message.data.decode("utf-8"))
|
||||
if topic_split[6] == "helperbot":
|
||||
# Response to command
|
||||
helperbotlog.debug(
|
||||
"Received Response - Topic: {} - Message: {}".format(
|
||||
message.topic, str(message.data.decode("utf-8"))
|
||||
)
|
||||
)
|
||||
helperbotlog.debug(f"Received Response - Topic: {topic} - Message: {data_decoded}")
|
||||
bumper.mqtt_helperbot.command_responses.append(
|
||||
{
|
||||
"time": time.time(),
|
||||
"topic": message.topic,
|
||||
"payload": str(message.data.decode("utf-8")),
|
||||
"topic": topic,
|
||||
"payload": data_decoded,
|
||||
}
|
||||
)
|
||||
elif str(message.topic).split("/")[3] == "helperbot":
|
||||
elif topic_split[3] == "helperbot":
|
||||
# Helperbot sending command
|
||||
helperbotlog.debug(
|
||||
"Send Command - Topic: {} - Message: {}".format(
|
||||
message.topic, str(message.data.decode("utf-8"))
|
||||
)
|
||||
)
|
||||
elif str(message.topic).split("/")[1] == "atr":
|
||||
helperbotlog.debug(f"Send Command - Topic: {topic} - Message: {data_decoded}")
|
||||
elif topic_split[1] == "atr":
|
||||
# Broadcast message received on atr
|
||||
if str(message.topic).split("/")[2] == "errors":
|
||||
boterrorlog.error(
|
||||
"Received Error - Topic: {} - Message: {}".format(
|
||||
message.topic, str(message.data.decode("utf-8"))
|
||||
)
|
||||
)
|
||||
if topic_split[2] == "errors":
|
||||
boterrorlog.error(f"Received Error - Topic: {topic} - Message: {data_decoded}")
|
||||
else:
|
||||
helperbotlog.debug(
|
||||
"Received Broadcast - Topic: {} - Message: {}".format(
|
||||
message.topic, str(message.data.decode("utf-8"))
|
||||
)
|
||||
)
|
||||
|
||||
helperbotlog.debug(f"Received Broadcast - Topic: {topic} - Message: {data_decoded}")
|
||||
else:
|
||||
helperbotlog.debug(
|
||||
"Received Message - Topic: {} - Message: {}".format(
|
||||
message.topic, str(message.data.decode("utf-8"))
|
||||
)
|
||||
)
|
||||
helperbotlog.debug(f"Received Message - Topic: {topic} - Message: {data_decoded}")
|
||||
|
||||
# Cleanup "expired messages" > 60 seconds from time
|
||||
for msg in bumper.mqtt_helperbot.command_responses:
|
||||
expire_time = (
|
||||
datetime.fromtimestamp(msg["time"])
|
||||
+ timedelta(seconds=MQTTHelperBot.wait_resp_timeout_seconds)
|
||||
datetime.fromtimestamp(msg["time"]) + timedelta(seconds=MQTTHelperBot.wait_resp_timeout_seconds)
|
||||
).timestamp()
|
||||
if time.time() > expire_time:
|
||||
helperbotlog.debug(
|
||||
"Pruning Message Due To Expiration - Message Topic: {}".format(
|
||||
msg["topic"]
|
||||
)
|
||||
)
|
||||
helperbotlog.debug(f"Pruning Message Due To Expiration - Message Topic: {msg['topic']}")
|
||||
bumper.mqtt_helperbot.command_responses.remove(msg)
|
||||
|
||||
async def on_broker_client_disconnected(self, client_id):
|
||||
|
||||
didsplit = str(client_id).split("@")
|
||||
|
||||
bot = bumper.bot_get(didsplit[0])
|
||||
if bot:
|
||||
bumper.bot_set_mqtt(bot["did"], False)
|
||||
return
|
||||
|
||||
clientresource = didsplit[1].split("/")[1]
|
||||
client = bumper.client_get(clientresource)
|
||||
if client:
|
||||
bumper.client_set_mqtt(client["resource"], False)
|
||||
return
|
||||
self._set_client_connected(client_id, False)
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue