Merge async and xmpp work #34

Merged
bmartin5692 merged 20 commits from dev_broken-XMPP into master 2019-05-23 15:22:58 +02:00
Showing only changes of commit 652a0b0e9b - Show all commits

View file

@ -8,7 +8,6 @@ from hbmqtt.broker import Broker
from hbmqtt.client import MQTTClient from hbmqtt.client import MQTTClient
from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2 from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2
import pkg_resources import pkg_resources
#import contextvars
import time import time
from threading import Thread from threading import Thread
import ssl import ssl
@ -40,7 +39,7 @@ class MQTTHelperBot:
): ):
self.address = address self.address = address
self.client_id = "helper1@bumper/helper1" self.client_id = "helper1@bumper/helper1"
self.command_responses = [] # = contextvars.ContextVar("command_responses", default=[]) self.command_responses = []
self.helperthread = None self.helperthread = None
def run(self, run_async=False): def run(self, run_async=False):
@ -99,10 +98,9 @@ class MQTTHelperBot:
message = await self.Client.deliver_message() message = await self.Client.deliver_message()
# helperbotlog.debug("HelperBot MQTT Received Message on Topic: {} - Message: {}".format(message.topic, str(message.payload.decode("utf-8")))) # 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": if str(message.topic).split("/")[6] == "helper1":
cresp.append( self.command_responses.append(
{ {
"time": time.time(), "time": time.time(),
"topic": message.topic, "topic": message.topic,
@ -111,15 +109,14 @@ class MQTTHelperBot:
) )
# Cleanup "expired messages" > 60 seconds from time # Cleanup "expired messages" > 60 seconds from time
for msg in cresp: for msg in self.command_responses:
expire_time = ( expire_time = (
datetime.fromtimestamp(msg["time"]) + timedelta(seconds=10) datetime.fromtimestamp(msg["time"]) + timedelta(seconds=10)
).timestamp() ).timestamp()
if time.time() > expire_time: if time.time() > expire_time:
# helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time)) # helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time))
cresp.remove(msg) self.command_responses.remove(msg)
self.command_responses = cresp #.set(cresp)
# helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp)) # helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
except Exception as e: except Exception as e:
@ -132,9 +129,8 @@ class MQTTHelperBot:
while time.time() < t_end: while time.time() < t_end:
await asyncio.sleep(0.1) await asyncio.sleep(0.1)
responses = self.command_responses.get() if len(self.command_responses) > 0:
if len(responses) > 0: for msg in self.command_responses:
for msg in responses:
topic = str(msg["topic"]).split("/") topic = str(msg["topic"]).split("/")
if topic[6] == "helper1" and topic[10] == requestid: if topic[6] == "helper1" and topic[10] == requestid:
# helperbotlog.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload'])) # helperbotlog.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload']))
@ -143,9 +139,7 @@ class MQTTHelperBot:
else: else:
resppayload = str(msg["payload"]) resppayload = str(msg["payload"])
resp = {"id": requestid, "ret": "ok", "resp": resppayload} resp = {"id": requestid, "ret": "ok", "resp": resppayload}
cresp = self.command_responses.get() self.command_responses.remove(msg)
cresp.remove(msg)
self.command_responses.set(cresp)
return resp return resp
return {"id": requestid, "errno": "timeout", "ret": "fail"} return {"id": requestid, "errno": "timeout", "ret": "fail"}