Merge async and xmpp work #34
3 changed files with 21 additions and 5 deletions
|
|
@ -31,7 +31,7 @@ class aiohttp_filter(logging.Filter):
|
||||||
|
|
||||||
confserverlog = logging.getLogger("confserver")
|
confserverlog = logging.getLogger("confserver")
|
||||||
|
|
||||||
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
|
#logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
|
||||||
logging.getLogger("aiohttp.access").addFilter(aiohttp_filter())
|
logging.getLogger("aiohttp.access").addFilter(aiohttp_filter())
|
||||||
|
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -84,6 +84,7 @@ class MQTTHelperBot:
|
||||||
[
|
[
|
||||||
("iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+", QOS_0),
|
("iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+", QOS_0),
|
||||||
("iot/p2p/+", QOS_0),
|
("iot/p2p/+", QOS_0),
|
||||||
|
("iot/atr/+", QOS_0),
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
@ -97,9 +98,9 @@ class MQTTHelperBot:
|
||||||
while True:
|
while True:
|
||||||
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"))))
|
|
||||||
|
|
||||||
if str(message.topic).split("/")[6] == "helper1":
|
if str(message.topic).split("/")[6] == "helper1":
|
||||||
|
#Response to command
|
||||||
|
helperbotlog.debug("Received Response - Topic: {} - Message: {}".format(message.topic, str(message.data.decode("utf-8"))))
|
||||||
self.command_responses.append(
|
self.command_responses.append(
|
||||||
{
|
{
|
||||||
"time": time.time(),
|
"time": time.time(),
|
||||||
|
|
@ -107,6 +108,14 @@ class MQTTHelperBot:
|
||||||
"payload": str(message.data.decode("utf-8")),
|
"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
|
# Cleanup "expired messages" > 60 seconds from time
|
||||||
for msg in self.command_responses:
|
for msg in self.command_responses:
|
||||||
|
|
@ -114,7 +123,7 @@ class MQTTHelperBot:
|
||||||
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))
|
||||||
self.command_responses.remove(msg)
|
self.command_responses.remove(msg)
|
||||||
|
|
||||||
# helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
|
# helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
|
||||||
|
|
@ -128,7 +137,7 @@ class MQTTHelperBot:
|
||||||
t_end = (datetime.now() + timedelta(seconds=10)).timestamp()
|
t_end = (datetime.now() + timedelta(seconds=10)).timestamp()
|
||||||
|
|
||||||
while time.time() < t_end:
|
while time.time() < t_end:
|
||||||
await asyncio.sleep(0.1)
|
await asyncio.sleep(0.1)
|
||||||
if len(self.command_responses) > 0:
|
if len(self.command_responses) > 0:
|
||||||
for msg in self.command_responses:
|
for msg in self.command_responses:
|
||||||
topic = str(msg["topic"]).split("/")
|
topic = str(msg["topic"]).split("/")
|
||||||
|
|
@ -145,8 +154,10 @@ class MQTTHelperBot:
|
||||||
return {"id": requestid, "errno": "timeout", "ret": "fail"}
|
return {"id": requestid, "errno": "timeout", "ret": "fail"}
|
||||||
except asyncio.CancelledError as e:
|
except asyncio.CancelledError as e:
|
||||||
helperbotlog.debug("wait_for_resp cancelled by asyncio")
|
helperbotlog.debug("wait_for_resp cancelled by asyncio")
|
||||||
|
return {"id": requestid, "errno": "timeout", "ret": "fail"}
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
helperbotlog.exception("{}".format(e))
|
helperbotlog.exception("{}".format(e))
|
||||||
|
return {"id": requestid, "errno": "timeout", "ret": "fail"}
|
||||||
|
|
||||||
async def send_command(self, cmdjson, requestid):
|
async def send_command(self, cmdjson, requestid):
|
||||||
try:
|
try:
|
||||||
|
|
@ -171,6 +182,7 @@ class MQTTHelperBot:
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
helperbotlog.exception("{}".format(e))
|
helperbotlog.exception("{}".format(e))
|
||||||
|
return {}
|
||||||
|
|
||||||
|
|
||||||
class MQTTServer:
|
class MQTTServer:
|
||||||
|
|
|
||||||
|
|
@ -5,6 +5,8 @@ import bumper
|
||||||
import sys, socket
|
import sys, socket
|
||||||
import time
|
import time
|
||||||
import platform
|
import platform
|
||||||
|
import os
|
||||||
|
os.environ['PYTHONASYNCIODEBUG'] = '1'
|
||||||
import asyncio
|
import asyncio
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -57,6 +59,7 @@ def main():
|
||||||
loop = asyncio.new_event_loop()
|
loop = asyncio.new_event_loop()
|
||||||
|
|
||||||
# Start web servers
|
# Start web servers
|
||||||
|
loop.set_debug(True)
|
||||||
conf_server.confserver_app()
|
conf_server.confserver_app()
|
||||||
conf_server_2.confserver_app()
|
conf_server_2.confserver_app()
|
||||||
asyncio.ensure_future(conf_server.start_server(),loop=loop)
|
asyncio.ensure_future(conf_server.start_server(),loop=loop)
|
||||||
|
|
@ -73,6 +76,7 @@ def main():
|
||||||
|
|
||||||
loop.run_forever()
|
loop.run_forever()
|
||||||
|
|
||||||
|
|
||||||
# start xmpp server on port 5223 (sync)
|
# start xmpp server on port 5223 (sync)
|
||||||
#xmpp_server.run(run_async=True) # Start in new thread
|
#xmpp_server.run(run_async=True) # Start in new thread
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue