From 652a0b0e9beb8e7a98636ac456ac2103b11b7fd2 Mon Sep 17 00:00:00 2001 From: Brian Martin Date: Sun, 12 May 2019 15:42:47 -0400 Subject: [PATCH] Fix command_responses Final removal of contextvars --- bumper/mqttserver.py | 24 +++++++++--------------- 1 file changed, 9 insertions(+), 15 deletions(-) diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py index 859fd69..2cb775b 100644 --- a/bumper/mqttserver.py +++ b/bumper/mqttserver.py @@ -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): @@ -99,10 +98,9 @@ class MQTTHelperBot: 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( + self.command_responses.append( { "time": time.time(), "topic": message.topic, @@ -111,15 +109,14 @@ class MQTTHelperBot: ) # 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) + self.command_responses.remove(msg) - self.command_responses = cresp #.set(cresp) # helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp)) except Exception as e: @@ -131,10 +128,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'])) @@ -142,10 +138,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": "timeout", "ret": "fail"}