diff --git a/bumper/__init__.py b/bumper/__init__.py index f7e73c2..f379e97 100644 --- a/bumper/__init__.py +++ b/bumper/__init__.py @@ -163,23 +163,12 @@ async def start(): global mqtt_helperbot mqtt_helperbot = MQTTHelperBot((bumper_listen, mqtt_listen_port)) global conf_server - conf_server = ConfServer( - (bumper_listen, conf1_listen_port), usessl=True, helperbot=mqtt_helperbot - ) + conf_server = ConfServer((bumper_listen, conf1_listen_port), usessl=True) global conf_server_2 - conf_server_2 = ConfServer( - (bumper_listen, conf2_listen_port), usessl=False, helperbot=mqtt_helperbot - ) + conf_server_2 = ConfServer((bumper_listen, conf2_listen_port), usessl=False) global xmpp_server xmpp_server = XMPPServer((bumper_listen, xmpp_listen_port)) - # Start web servers - conf_server.confserver_app() - asyncio.create_task(conf_server.start_server()) - - conf_server_2.confserver_app() - asyncio.create_task(conf_server_2.start_server()) - # Start MQTT Server asyncio.create_task(mqtt_server.broker_coro()) @@ -189,6 +178,20 @@ async def start(): # Start XMPP Server asyncio.create_task(xmpp_server.start_async_server()) + # Wait for helperbot to connect first + while mqtt_helperbot.Client is None: + await asyncio.sleep(0.1) + + while not mqtt_helperbot.Client.session.transitions.state == "connected": + await asyncio.sleep(0.1) + + # Start web servers + conf_server.confserver_app() + asyncio.create_task(conf_server.start_server()) + + conf_server_2.confserver_app() + asyncio.create_task(conf_server_2.start_server()) + # Start maintenance while not shutting_down: asyncio.create_task(maintenance()) diff --git a/bumper/confserver.py b/bumper/confserver.py index f1fab45..11b451d 100644 --- a/bumper/confserver.py +++ b/bumper/confserver.py @@ -38,12 +38,9 @@ logging.getLogger("aiohttp.access").addFilter( class ConfServer: - def __init__(self, address, usessl=False, helperbot=None): - self.helperbot = helperbot + def __init__(self, address, usessl=False): self.usessl = usessl self.address = address - self.confthread = None - self.run_async = False self.app = None self.site = None self.runner = None @@ -57,6 +54,7 @@ class ConfServer: self.app.add_routes( [ web.get("", self.handle_base), + web.get("/restart_{service}", self.handle_RestartService), web.get( "/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/login", self.handle_login, @@ -198,11 +196,29 @@ class ConfServer: # text = "Bumper!" bots = bumper.db_get().table("bots").all() clients = bumper.db_get().table("clients").all() - helperbot = self.helperbot.Client.session.transitions.state + helperbot = bumper.mqtt_helperbot.Client.session.transitions.state + mqttserver = bumper.mqtt_server.broker + mq_sessions = [] + for sess in mqttserver._sessions: + tmpsess = [] + tmpsess.append({"client_id": mqttserver._sessions[sess][0].client_id}) + tmpsess.append( + {"state": mqttserver._sessions[sess][0].transitions.state} + ) + mq_sessions.append(tmpsess) all = { "bots": bots, "clients": clients, "helperbot": [{"state": helperbot}], + "mqtt_server": [ + {"state": mqttserver.transitions.state}, + { + "sessions": [ + {"count": len(mqttserver._sessions)}, + {"clients": mq_sessions}, + ] + }, + ], } return web.json_response(all) @@ -210,6 +226,49 @@ class ConfServer: except Exception as e: confserverlog.exception("{}".format(e)) + async def restart_Helper(self): + + await bumper.mqtt_helperbot.Client.disconnect() + await bumper.mqtt_helperbot.start_helper_bot() + + + + async def restart_MQTT(self): + mqttserver = bumper.mqtt_server.broker + + for sess in list(mqttserver._sessions): + sessobj = mqttserver._sessions[sess][1] + await sessobj.writer.close() + mqttserver.delete_session(sess) + + await bumper.mqtt_server.broker.shutdown() + while not bumper.mqtt_server.broker.transitions.state == "stopped": + await asyncio.sleep(0.1) + + await bumper.mqtt_server.broker_coro() + while not bumper.mqtt_server.broker.transitions.state == "started": + await asyncio.sleep(0.1) + + async def handle_RestartService(self, request): + try: + service = request.match_info.get("service", "") + if service == "Helperbot": + await self.restart_Helper() + return web.json_response({"status": "complete"}) + elif service == "MQTTServer": + await self.restart_MQTT() + aloop = asyncio.get_event_loop() + aloop.call_later( + 2, lambda: asyncio.create_task(self.restart_Helper()) + ) # In 2 seconds restart Helperbot + + return web.json_response({"status": "complete"}) + else: + return web.json_response({"status": "invalid service"}) + + except Exception as e: + confserverlog.exception("{}".format(e)) + async def handle_login(self, request): try: user_devid = request.match_info.get("devid", "") @@ -259,9 +318,7 @@ class ConfServer: # "username": "fusername_{}".format(tmpuser["userid"]), # }, "msg": "操作成功", - "time": self.self.get_milli_time( - datetime.utcnow().timestamp() - ), + "time": self.get_milli_time(datetime.utcnow().timestamp()), } return web.json_response(body) @@ -1048,122 +1105,136 @@ class ConfServer: confserverlog.exception("{}".format(e)) async def handle_lg_log(self, request): # EcoVacs Home - try: - json_body = json.loads(await request.text()) + if ( + not bumper.mqtt_helperbot.Client._handler.writer is None + ): # Ignore if the Helperbot writer is none + try: + json_body = json.loads(await request.text()) - randomid = "".join(random.sample(string.ascii_letters, 6)) - did = json_body["did"] + randomid = "".join(random.sample(string.ascii_letters, 6)) + did = json_body["did"] - botdetails = bumper.bot_get(did) - if botdetails: - if not "cmdName" in json_body: - if "td" in json_body: - json_body["cmdName"] = json_body["td"] - # json_body["td"] = "q" + botdetails = bumper.bot_get(did) + if botdetails: + if not "cmdName" in json_body: + if "td" in json_body: + json_body["cmdName"] = json_body["td"] + # json_body["td"] = "q" - if not "toId" in json_body: - json_body["toId"] = did + if not "toId" in json_body: + json_body["toId"] = did - if not "toType" in json_body: - json_body["toType"] = botdetails["class"] + if not "toType" in json_body: + json_body["toType"] = botdetails["class"] - if not "toRes" in json_body: - json_body["toRes"] = botdetails["resource"] + if not "toRes" in json_body: + json_body["toRes"] = botdetails["resource"] - if not "payloadType" in json_body: - json_body["payloadType"] = "x" + if not "payloadType" in json_body: + json_body["payloadType"] = "x" - if not "payload" in json_body: - json_body["payload"] = "" - if json_body["td"] == "GetCleanLogs": - json_body["td"] = "q" - json_body["payload"] = '' # " + if not "payload" in json_body: + json_body["payload"] = "" + if json_body["td"] == "GetCleanLogs": + json_body["td"] = "q" + json_body["payload"] = '' # " - if did != "": - bot = bumper.bot_get(did) - if bot["company"] == "eco-ng": - body = "" - retcmd = await self.helperbot.send_command(json_body, randomid) - confserverlog.debug("Send Bot - {}".format(json_body)) - confserverlog.debug("Bot Response - {}".format(body)) - logs = [] - logsroot = ET.fromstring(retcmd["resp"]) - if logsroot.attrib["ret"] == "ok": - cleanlogs = logsroot.getchildren() - for l in cleanlogs: - logs.append(l.attrib) - - body = { - "ret": "ok", - # "logs": logs, #TODO: Doesn't parse correctly, new protocol & server side processing - "logs": [], - } - - else: - body = {"ret": "ok", "logs": []} - - confserverlog.debug( - "POST: {} - Response: {}".format(json_body, body) - ) - - return web.json_response(body) - else: - # No response, send error back - confserverlog.error( - "No bots with DID: {} connected to MQTT".format( - json_body["toId"] + if did != "": + bot = bumper.bot_get(did) + if bot["company"] == "eco-ng": + body = "" + retcmd = await bumper.mqtt_helperbot.send_command( + json_body, randomid ) - ) - body = {"id": randomid, "errno": bumper.ERR_COMMON, "ret": "fail"} - return web.json_response(body) + confserverlog.debug("Send Bot - {}".format(json_body)) + confserverlog.debug("Bot Response - {}".format(body)) + logs = [] + logsroot = ET.fromstring(retcmd["resp"]) + if logsroot.attrib["ret"] == "ok": + cleanlogs = logsroot.getchildren() + for l in cleanlogs: + logs.append(l.attrib) - except Exception as e: - confserverlog.exception("{}".format(e)) + body = { + "ret": "ok", + # "logs": logs, #TODO: Doesn't parse correctly, new protocol & server side processing + "logs": [], + } + + else: + body = {"ret": "ok", "logs": []} + + confserverlog.debug( + "POST: {} - Response: {}".format(json_body, body) + ) + + return web.json_response(body) + else: + # No response, send error back + confserverlog.error( + "No bots with DID: {} connected to MQTT".format( + json_body["toId"] + ) + ) + body = { + "id": randomid, + "errno": bumper.ERR_COMMON, + "ret": "fail", + } + return web.json_response(body) + + except Exception as e: + confserverlog.exception("{}".format(e)) async def handle_devmanager_botcommand(self, request): - try: - json_body = json.loads(await request.text()) + if ( + not bumper.mqtt_helperbot.Client._handler.writer is None + ): # Ignore if the helperbot object isn't set + try: + json_body = json.loads(await request.text()) - randomid = "".join(random.sample(string.ascii_letters, 6)) - did = "" - if "toId" in json_body: # Its a command - did = json_body["toId"] + randomid = "".join(random.sample(string.ascii_letters, 6)) + did = "" + if "toId" in json_body: # Its a command + did = json_body["toId"] - if did != "": - bot = bumper.bot_get(did) - if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: - retcmd = await self.helperbot.send_command(json_body, randomid) - body = retcmd - confserverlog.debug("Send Bot - {}".format(json_body)) - confserverlog.debug("Bot Response - {}".format(body)) - return web.json_response(body) - else: - # No response, send error back - confserverlog.error( - "No bots with DID: {} connected to MQTT".format( - json_body["toId"] + if did != "": + bot = bumper.bot_get(did) + if bot["company"] == "eco-ng": + retcmd = await bumper.mqtt_helperbot.send_command( + json_body, randomid ) - ) - body = { - "id": randomid, - "errno": 500, - "ret": "fail", - "debug": "wait for response timed out", - } - return web.json_response(body) - - else: - if "td" in json_body: # Seen when doing initial wifi config - if json_body["td"] == "PollSCResult": - body = {"ret": "ok"} + body = retcmd + confserverlog.debug("Send Bot - {}".format(json_body)) + confserverlog.debug("Bot Response - {}".format(body)) + return web.json_response(body) + else: + # No response, send error back + confserverlog.error( + "No bots with DID: {} connected to MQTT".format( + json_body["toId"] + ) + ) + body = { + "id": randomid, + "errno": 500, + "ret": "fail", + "debug": "wait for response timed out", + } return web.json_response(body) - if json_body["td"] == "HasUnreadMsg": # EcoVacs Home - body = {"ret": "ok", "unRead": False} - return web.json_response(body) + else: + if "td" in json_body: # Seen when doing initial wifi config + if json_body["td"] == "PollSCResult": + body = {"ret": "ok"} + return web.json_response(body) - except Exception as e: - confserverlog.exception("{}".format(e)) + if json_body["td"] == "HasUnreadMsg": # EcoVacs Home + body = {"ret": "ok", "unRead": False} + return web.json_response(body) + + except Exception as e: + confserverlog.exception("{}".format(e)) async def handle_dim_devmanager(self, request): # Used in EcoVacs Home App try: @@ -1177,7 +1248,9 @@ class ConfServer: if did != "": bot = bumper.bot_get(did) if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: - retcmd = await self.helperbot.send_command(json_body, randomid) + retcmd = await bumper.mqtt_helperbot.send_command( + json_body, randomid + ) body = retcmd confserverlog.debug("Send Bot - {}".format(json_body)) confserverlog.debug("Bot Response - {}".format(body)) @@ -1208,10 +1281,7 @@ class ConfServer: async def disconnect(self): try: confserverlog.info("shutting down") - if self.run_async: - self.confthread.join() - else: - await self.app.shutdown() + await self.app.shutdown() except Exception as e: confserverlog.exception("{}".format(e)) diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py index 9e5de56..6b18a75 100644 --- a/bumper/mqttserver.py +++ b/bumper/mqttserver.py @@ -31,7 +31,7 @@ logging.getLogger("hbmqtt.client").setLevel(logging.CRITICAL + 1) # Ignore this class MQTTHelperBot: - Client = MQTTClient() + Client = None wait_resp_timeout_seconds = 10 expire_msg_seconds = 10 @@ -39,14 +39,14 @@ class MQTTHelperBot: self.address = address self.client_id = "helper1@bumper/helper1" self.command_responses = [] - self.helperthread = None async def start_helper_bot(self): try: - self.Client = MQTTClient( - client_id=self.client_id, config={"check_hostname": False} - ) + if self.Client is None: + self.Client = MQTTClient( + client_id=self.client_id, config={"check_hostname": False} + ) await self.Client.connect( "mqtts://{}:{}/".format(self.address[0], self.address[1]), @@ -183,29 +183,30 @@ class MQTTHelperBot: } async def send_command(self, cmdjson, requestid): - try: - ttopic = "iot/p2p/{}/helper1/bumper/helper1/{}/{}/{}/q/{}/{}".format( - cmdjson["cmdName"], - cmdjson["toId"], - cmdjson["toType"], - cmdjson["toRes"], - requestid, - cmdjson["payloadType"], - ) + if not self.Client._handler.writer is None: try: - await self.Client.publish( - ttopic, str(cmdjson["payload"]).encode(), QOS_0 + ttopic = "iot/p2p/{}/helper1/bumper/helper1/{}/{}/{}/q/{}/{}".format( + cmdjson["cmdName"], + cmdjson["toId"], + cmdjson["toType"], + cmdjson["toRes"], + requestid, + cmdjson["payloadType"], ) + try: + await self.Client.publish( + ttopic, str(cmdjson["payload"]).encode(), QOS_0 + ) + except Exception as e: + helperbotlog.exception("{}".format(e)) + + resp = await self.wait_for_resp(requestid) + + return resp + except Exception as e: helperbotlog.exception("{}".format(e)) - - resp = await self.wait_for_resp(requestid) - - return resp - - except Exception as e: - helperbotlog.exception("{}".format(e)) - return {} + return {} class MQTTServer: @@ -233,7 +234,6 @@ class MQTTServer: def __init__(self, address): try: - self.mqttserverthread = None self.address = address # The below adds a plugin to the hbmqtt.broker.plugins without having to futz with setup.py @@ -355,6 +355,7 @@ class BumperMQTTServer_Plugin: async def on_broker_client_connected(self, client_id): try: + didsplit = str(client_id).split("@") bot = bumper.bot_get(didsplit[0]) diff --git a/bumper/xmppserver.py b/bumper/xmppserver.py index e4f0ba3..ca52455 100644 --- a/bumper/xmppserver.py +++ b/bumper/xmppserver.py @@ -52,7 +52,7 @@ class XMPPServer: def disconnect(self): try: - xmppserverlog.debug("waiting for all client threads to exit") + xmppserverlog.debug("waiting for all clients to disconnect") for client in self.clients: client._disconnect() @@ -722,7 +722,11 @@ class XMPPAsyncClient: ).replace("ns0:", ""), ) ) - if 'td="error"' in newdata or 'errs=' in newdata or 'k="DeviceAlert' in newdata: + if ( + 'td="error"' in newdata + or "errs=" in newdata + or 'k="DeviceAlert' in newdata + ): boterrorlog.error( "Received Error from ({}:{} | {}) - {}".format( self.address[0],