Merge branch 'dev_broken-XMPP' into fix-#31
This commit is contained in:
commit
3ae24976af
7 changed files with 1138 additions and 243 deletions
|
|
@ -8,7 +8,6 @@ from hbmqtt.broker import Broker
|
|||
from hbmqtt.client import MQTTClient
|
||||
from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2
|
||||
import pkg_resources
|
||||
import contextvars
|
||||
import time
|
||||
from threading import Thread
|
||||
import ssl
|
||||
|
|
@ -40,7 +39,7 @@ class MQTTHelperBot:
|
|||
):
|
||||
self.address = address
|
||||
self.client_id = "helper1@bumper/helper1"
|
||||
self.command_responses = contextvars.ContextVar("command_responses", default=[])
|
||||
self.command_responses = []
|
||||
self.helperthread = None
|
||||
|
||||
def run(self, run_async=False):
|
||||
|
|
@ -73,6 +72,10 @@ class MQTTHelperBot:
|
|||
async def start_helper_bot(self):
|
||||
|
||||
try:
|
||||
self.Client = MQTTClient(
|
||||
client_id=self.client_id, config={"check_hostname": False}
|
||||
)
|
||||
|
||||
await self.Client.connect(
|
||||
"mqtts://{}:{}/".format(self.address[0], self.address[1]),
|
||||
cafile=bumper.ca_cert,
|
||||
|
|
@ -81,9 +84,12 @@ class MQTTHelperBot:
|
|||
[
|
||||
("iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+", QOS_0),
|
||||
("iot/p2p/+", QOS_0),
|
||||
("iot/atr/+", QOS_0),
|
||||
]
|
||||
)
|
||||
|
||||
asyncio.ensure_future(self.get_msg())
|
||||
|
||||
except Exception as e:
|
||||
helperbotlog.exception("{}".format(e))
|
||||
|
||||
|
|
@ -92,28 +98,34 @@ class MQTTHelperBot:
|
|||
while True:
|
||||
message = await self.Client.deliver_message()
|
||||
|
||||
# helperbotlog.debug("HelperBot MQTT Received Message on Topic: {} - Message: {}".format(message.topic, str(message.payload.decode("utf-8"))))
|
||||
cresp = self.command_responses.get()
|
||||
|
||||
if str(message.topic).split("/")[6] == "helper1":
|
||||
cresp.append(
|
||||
#Response to command
|
||||
helperbotlog.debug("Received Response - Topic: {} - Message: {}".format(message.topic, str(message.data.decode("utf-8"))))
|
||||
self.command_responses.append(
|
||||
{
|
||||
"time": time.time(),
|
||||
"topic": message.topic,
|
||||
"payload": str(message.data.decode("utf-8")),
|
||||
}
|
||||
)
|
||||
elif str(message.topic).split("/")[3] == "helper1":
|
||||
#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":
|
||||
#Broadcast message received on atr
|
||||
helperbotlog.debug("Received Broadcast - Topic: {} - Message: {}".format(message.topic, str(message.data.decode("utf-8"))))
|
||||
else:
|
||||
helperbotlog.debug("Received Message - Topic: {} - Message: {}".format(message.topic, str(message.data.decode("utf-8"))))
|
||||
|
||||
# Cleanup "expired messages" > 60 seconds from time
|
||||
for msg in cresp:
|
||||
for msg in self.command_responses:
|
||||
expire_time = (
|
||||
datetime.fromtimestamp(msg["time"]) + timedelta(seconds=10)
|
||||
).timestamp()
|
||||
if time.time() > expire_time:
|
||||
# helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time))
|
||||
cresp.remove(msg)
|
||||
helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time))
|
||||
self.command_responses.remove(msg)
|
||||
|
||||
self.command_responses.set(cresp)
|
||||
# helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
|
||||
|
||||
except Exception as e:
|
||||
|
|
@ -125,10 +137,9 @@ class MQTTHelperBot:
|
|||
t_end = (datetime.now() + timedelta(seconds=10)).timestamp()
|
||||
|
||||
while time.time() < t_end:
|
||||
await asyncio.sleep(0.1)
|
||||
responses = self.command_responses.get()
|
||||
if len(responses) > 0:
|
||||
for msg in responses:
|
||||
await asyncio.sleep(0.1)
|
||||
if len(self.command_responses) > 0:
|
||||
for msg in self.command_responses:
|
||||
topic = str(msg["topic"]).split("/")
|
||||
if topic[6] == "helper1" and topic[10] == requestid:
|
||||
# helperbotlog.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload']))
|
||||
|
|
@ -136,10 +147,8 @@ class MQTTHelperBot:
|
|||
resppayload = json.loads(msg["payload"])
|
||||
else:
|
||||
resppayload = str(msg["payload"])
|
||||
resp = {"id": requestid, "ret": "ok", "resp": resppayload}
|
||||
cresp = self.command_responses.get()
|
||||
cresp.remove(msg)
|
||||
self.command_responses.set(cresp)
|
||||
resp = {"id": requestid, "ret": "ok", "resp": resppayload}
|
||||
self.command_responses.remove(msg)
|
||||
return resp
|
||||
|
||||
return {"id": requestid, "errno": 500, "ret": "fail", "debug": "wait for response timed out"}
|
||||
|
|
@ -173,6 +182,7 @@ class MQTTHelperBot:
|
|||
|
||||
except Exception as e:
|
||||
helperbotlog.exception("{}".format(e))
|
||||
return {}
|
||||
|
||||
|
||||
class MQTTServer:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue