From 2460182f4f1826a42c6633ff22a3285f6380214f Mon Sep 17 00:00:00 2001 From: Brian Martin Date: Thu, 23 May 2019 09:01:30 -0400 Subject: [PATCH] Fix the things XMPP is now async and works - Closes #27 Full async - Closes #22 Individual log files - Closes #21 --- bumper/__init__.py | 21 ++-- bumper/confserver.py | 189 +++++++++++++++--------------- bumper/mqttserver.py | 68 +---------- bumper/xmppserver.py | 273 +++++++++++++------------------------------ start_bumper.py | 71 +++++------ 5 files changed, 220 insertions(+), 402 deletions(-) diff --git a/bumper/__init__.py b/bumper/__init__.py index 6e53947..e2ffaaa 100644 --- a/bumper/__init__.py +++ b/bumper/__init__.py @@ -1,9 +1,8 @@ #!/usr/bin/env python3 -from .confserver import ConfServer -from .mqttserver import MQTTServer -from .mqttserver import MQTTHelperBot -from .xmppserver import XMPPServer +from bumper.confserver import ConfServer +from bumper.mqttserver import MQTTServer, MQTTHelperBot +from bumper.xmppserver import XMPPServer import asyncio import json import time @@ -66,6 +65,9 @@ xmppserverlog.addHandler(xmpp_rotate) # xmppserverlog.setLevel(logging.INFO) +logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger + + def get_milli_time(timetoconvert): return int(round(timetoconvert * 1000)) @@ -663,11 +665,12 @@ def bot_add(sn, did, devclass, resource, company): newbot.company = company bot = bot_get(did) - if not bot: - bumperlog.info( - "Adding new bot with SN: {} DID: {}".format(newbot.name, newbot.did) - ) - bot_full_upsert(newbot.asdict()) + if not bot: # Not existing bot in database + if not devclass == "" or "@" not in sn or "tmp" not in sn: # try to prevent bad additions to the bot list + bumperlog.info( + "Adding new bot with SN: {} DID: {}".format(newbot.name, newbot.did) + ) + bot_full_upsert(newbot.asdict()) def bot_remove(did): diff --git a/bumper/confserver.py b/bumper/confserver.py index 70f2474..1a51dc0 100644 --- a/bumper/confserver.py +++ b/bumper/confserver.py @@ -31,9 +31,7 @@ class aiohttp_filter(logging.Filter): confserverlog = logging.getLogger("confserver") - -#logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger -logging.getLogger("aiohttp.access").addFilter(aiohttp_filter()) +logging.getLogger("aiohttp.access").addFilter(aiohttp_filter()) #Add logging filter above to aiohttp.access class EcoVacs_Login: accessToken = "" @@ -60,42 +58,6 @@ class ConfServer: self.run_async = False self.app = None - def run(self, run_async=False): - try: - if run_async: - self.run_async = True - confserverlog.debug("Starting ConfServer Thread: 1") - self.confthread = Thread( - name="ConfServer_{}_Thread".format(self.address[1]), - target=self.run_server, - ) - self.confthread.setDaemon(True) - self.confthread.start() - - else: - try: - self.run_server() - except KeyboardInterrupt: - self.disconnect() - - except Exception as e: - confserverlog.exception("{}".format(e)) - - def run_server(self): - logging.info("Starting ConfServer at {}".format(self.address)) - print("Starting ConfServer at {}".format(self.address)) - try: - loop = asyncio.get_event_loop() - except: - loop = asyncio.new_event_loop() - - try: - self.confserver_app() - loop.run_until_complete(self.start_server()) - loop.run_forever() - except Exception as e: - confserverlog.exception("{}".format(e)) - def confserver_app(self): self.app = web.Application() @@ -196,6 +158,7 @@ class ConfServer: async def start_server(self): try: + confserverlog.info("Starting ConfServer at {}:{}".format(self.address[0], self.address[1])) runner = web.AppRunner(self.app) await runner.setup() @@ -351,55 +314,67 @@ class ConfServer: def check_token(self, apptype, countrycode, user, token): - if bumper.check_token(user["userid"], token): - - if "global_" in apptype: #EcoVacs Home - login_details = EcoVacsHome_Login() - login_details.ucUid = "fuid_{}".format(user["userid"]) - login_details.loginName = "fusername_{}".format(user["userid"]) - login_details.mobile = None + try: + if bumper.check_token(user["userid"], token): + + if "global_" in apptype: #EcoVacs Home + login_details = EcoVacsHome_Login() + login_details.ucUid = "fuid_{}".format(user["userid"]) + login_details.loginName = "fusername_{}".format(user["userid"]) + login_details.mobile = None + else: + login_details = EcoVacs_Login() + + login_details.accessToken = token + login_details.uid = "fuid_{}".format(user["userid"]) + login_details.username = "fusername_{}".format(user["userid"]) + login_details.country = countrycode + login_details.email = "null@null.com" + + body = { + "code": bumper.RETURN_API_SUCCESS, + "data": json.loads(login_details.toJSON()), + #{ + # "accessToken": self.generate_token(tmpuser), # Generate a token + # "country": countrycode, + # "email": "null@null.com", + # "uid": "fuid_{}".format(tmpuser["userid"]), + # "username": "fusername_{}".format(tmpuser["userid"]), + #}, + "msg": "操作成功", + "time": bumper.get_milli_time(datetime.utcnow().timestamp()), + } + return web.json_response(body) + else: - login_details = EcoVacs_Login() - - login_details.accessToken = token - login_details.uid = "fuid_{}".format(user["userid"]) - login_details.username = "fusername_{}".format(user["userid"]) - login_details.country = countrycode - login_details.email = "null@null.com" - - body = { - "code": bumper.RETURN_API_SUCCESS, - "data": json.loads(login_details.toJSON()), - #{ - # "accessToken": self.generate_token(tmpuser), # Generate a token - # "country": countrycode, - # "email": "null@null.com", - # "uid": "fuid_{}".format(tmpuser["userid"]), - # "username": "fusername_{}".format(tmpuser["userid"]), - #}, - "msg": "操作成功", - "time": bumper.get_milli_time(datetime.utcnow().timestamp()), - } - return web.json_response(body) - - else: - body = { - "code": bumper.ERR_TOKEN_INVALID, - "data": None, - "msg": "当前密码错误", - "time": bumper.get_milli_time(datetime.utcnow().timestamp()), - } - return web.json_response(body) + body = { + "code": bumper.ERR_TOKEN_INVALID, + "data": None, + "msg": "当前密码错误", + "time": bumper.get_milli_time(datetime.utcnow().timestamp()), + } + return web.json_response(body) + + except Exception as e: + confserverlog.exception("{}".format(e)) def generate_token(self, user): - tmpaccesstoken = uuid.uuid4().hex - bumper.user_add_token(user["userid"], tmpaccesstoken) - return tmpaccesstoken + try: + tmpaccesstoken = uuid.uuid4().hex + bumper.user_add_token(user["userid"], tmpaccesstoken) + return tmpaccesstoken + + except Exception as e: + confserverlog.exception("{}".format(e)) def generate_authcode(self, user, countrycode, token): - tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex) - bumper.user_add_authcode(user["userid"], token, tmpauthcode) - return tmpauthcode + try: + tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex) + bumper.user_add_authcode(user["userid"], token, tmpauthcode) + return tmpauthcode + + except Exception as e: + confserverlog.exception("{}".format(e)) def _auth_any(self, devid, apptype, country, request): try: @@ -815,7 +790,6 @@ class ConfServer: confserverlog.exception("{}".format(e)) async def handle_getProductIotMap(self, request): - user_devid = request.match_info.get("devid", "") try: body = { "code": bumper.RETURN_API_SUCCESS, @@ -909,13 +883,23 @@ class ConfServer: if todo == "FindBest": service = postbody["service"] if service == "EcoMsgNew": + srvip = socket.gethostbyname(socket.gethostname()) + srvport = 5223 + confserverlog.info( + "Reporting FindBest-EcoMsgNew Server to Bot as: {}:{}".format(srvip, srvport) + ) body = { "result": "ok", - "ip": socket.gethostbyname(socket.gethostname()), - "port": 5223, + "ip": srvip, + "port": srvport, } elif service == "EcoUpdate": - body = {"result": "ok", "ip": "47.88.66.164", "port": 8005} + srvip = "47.88.66.164" #EcoVacs Server + srvport = 8005 + confserverlog.info( + "Reporting FindBest-EcoUpdate Server to Bot as: {}:{}".format(srvip, srvport) + ) + body = {"result": "ok", "ip": srvip, "port": srvport} elif todo == "loginByItToken": if "userId" in postbody: @@ -995,7 +979,8 @@ class ConfServer: for bot in bots: if bot["class"] != "": b = bumper.bot_toEcoVacsHome_JSON(bot) - botlist.append(json.loads(b)) + if not b is None: #Happens if the bot isn't on the EcoVacs Home list + botlist.append(json.loads(b)) body = { "code": 0, @@ -1035,9 +1020,12 @@ class ConfServer: if todo == "FindBest": service = postbody["service"] if service == "EcoMsgNew": - srvip = socket.gethostbyname(socket.gethostname()) - msgserver = {"ip": srvip, "port": 5223, "result": "ok"} + srvport = 5223 + confserverlog.info( + "Reporting FindBest-EcoMsgNew Server to Bot as: {}:{}".format(srvip, srvport) + ) + msgserver = {"ip": srvip, "port": srvport, "result": "ok"} msgserver = json.dumps(msgserver) msgserver = msgserver.replace( " ", "" @@ -1094,11 +1082,16 @@ class ConfServer: if did != "": - confserverlog.debug("BotCommand: {}".format(json_body)) bot = bumper.bot_get(did) if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True: 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": @@ -1147,13 +1140,15 @@ class ConfServer: did = json_body["toId"] if did != "": - confserverlog.debug("BotCommand: {}".format(json_body)) 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( - "\r\n POST: {} \r\n Response: {}".format(json_body, body) + "Send Bot - {}".format(json_body) + ) + confserverlog.debug( + "Bot Response - {}".format(body) ) return web.json_response(body) else: @@ -1180,7 +1175,7 @@ class ConfServer: confserverlog.exception("{}".format(e)) - async def handle_dim_devmanager(self, request): + async def handle_dim_devmanager(self, request): #Used in EcoVacs Home App try: json_body = json.loads(await request.text()) @@ -1190,13 +1185,15 @@ class ConfServer: did = json_body["toId"] if did != "": - confserverlog.debug("BotCommand: {}".format(json_body)) 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( - "\r\n POST: {} \r\n Response: {}".format(json_body, body) + "Send Bot - {}".format(json_body) + ) + confserverlog.debug( + "Bot Response - {}".format(body) ) return web.json_response(body) else: diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py index 0f145f6..a9c3ca1 100644 --- a/bumper/mqttserver.py +++ b/bumper/mqttserver.py @@ -14,6 +14,7 @@ import ssl import bumper import json from datetime import datetime, timedelta +import bumper helperbotlog = logging.getLogger("helperbot") mqttserverlog = logging.getLogger("mqttserver") @@ -42,33 +43,6 @@ class MQTTHelperBot: self.command_responses = [] self.helperthread = None - def run(self, run_async=False): - if run_async: - hloop = asyncio.new_event_loop() - helperbotlog.debug("Starting MQTT HelperBot Thread: 1") - self.helperthread = Thread( - name="MQTTHelperBot_Thread", target=self.run_helperbot, args=(hloop,) - ) - self.helperthread.setDaemon(True) - self.helperthread.start() - - else: - self.run_helperbot(asyncio.get_event_loop()) - - def run_helperbot(self, loop): - logging.info("Starting MQTT HelperBot") - print("Starting MQTT HelperBot") - try: - asyncio.set_event_loop(loop) - self.Client = MQTTClient( - client_id=self.client_id, config={"check_hostname": False} - ) - loop.run_until_complete(self.start_helper_bot()) - loop.run_until_complete(self.get_msg()) - loop.run_forever() - except Exception as e: - helperbotlog.exception("{}".format(e)) - async def start_helper_bot(self): try: @@ -87,8 +61,7 @@ class MQTTHelperBot: ("iot/atr/+", QOS_0), ] ) - - asyncio.ensure_future(self.get_msg()) + asyncio.create_task(self.get_msg()) except Exception as e: helperbotlog.exception("{}".format(e)) @@ -126,8 +99,6 @@ class MQTTHelperBot: helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time)) self.command_responses.remove(msg) - # helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp)) - except Exception as e: helperbotlog.exception("{}".format(e)) @@ -190,6 +161,7 @@ class MQTTServer: async def broker_coro(self): try: + mqttserverlog.info("Starting MQTT Server at {}:{}".format(self.address[0], self.address[1])) broker = hbmqtt.broker.Broker(config=self.default_config) await broker.start() @@ -248,32 +220,6 @@ class MQTTServer: except Exception as e: mqttserverlog.exception("{}".format(e)) - def run(self, run_async=False): - if run_async: - sloop = asyncio.new_event_loop() - mqttserverlog.debug("Starting MQTTServer Thread: 1") - self.mqttserverthread = Thread( - name="MQTTServer_Thread", target=self.run_server, args=(sloop,) - ) - self.mqttserverthread.setDaemon(True) - self.mqttserverthread.start() - - else: - self.run_server(asyncio.get_event_loop()) - - def run_server(self, loop): - - logging.info("Starting MQTT Server at {}".format(self.address)) - print("Starting MQTT Server at {}".format(self.address)) - try: - asyncio.set_event_loop(loop) - loop.run_until_complete(self.broker_coro()) - # loop.run_until_complete(self.active_bot_listing()) - loop.run_forever() - - except Exception as e: - mqttserverlog.exception("{}".format(e)) - class BumperMQTTServer_Plugin: def __init__(self, context): @@ -310,11 +256,9 @@ class BumperMQTTServer_Plugin: client_id = session.client_id didsplit = str(client_id).split("@") - # If this isn't a fake user (fuid) then add as a bot - if not ( - str(didsplit[0]).startswith("fuid") - or str(didsplit[0]).startswith("helper") - ): + if not ( # if ecouser or bumper aren't in details it is a bot + "ecouser" in didsplit[1] + or "bumper" in didsplit[1]): tmpbotdetail = str(didsplit[1]).split("/") bumper.bot_add( username, diff --git a/bumper/xmppserver.py b/bumper/xmppserver.py index 5ad8e22..7d00fc5 100644 --- a/bumper/xmppserver.py +++ b/bumper/xmppserver.py @@ -5,7 +5,7 @@ import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as import base64 import ssl import bumper -import asyncio +import asyncio, functools xmppserverlog = logging.getLogger("xmppserver") @@ -19,11 +19,10 @@ class XMPPServer: def __init__(self, address): # Initialize bot server self.address = address - self.aclients = {} # task -> (reader, writer) async def async_server(self): - xmppserverlog.debug( - "listening on {}:{}".format(self.address[0], self.address[1]) + xmppserverlog.info( + "Starting XMPP Server at {}:{}".format(self.address[0], self.address[1]) ) server = await asyncio.start_server( self.accept_client, self.address[0], self.address[1] @@ -34,101 +33,26 @@ class XMPPServer: # self.clients = {} # task -> (reader, writer) def accept_client(self, client_reader, client_writer): - # task = asyncio.Task(self.handle_client(client_reader, client_writer)) - aclient = XMPPAsyncClient(client_reader, client_writer) - task = asyncio.Task(aclient.handle_async_client()) - self.aclients[task] = (client_reader, client_writer) + try: + aclient = XMPPAsyncClient(client_reader, client_writer) + task = asyncio.Task(aclient.handle_async_client()) + aclient._async_task = task + self.clients.append(aclient) - def client_done(task): - del self.aclients[task] - xmppserverlog.info("End Connection for {}".format(client_writer.get_extra_info("peername"))) - client_writer.close() + def client_done(aclient, task): + try: + self.clients.remove(aclient) + xmppserverlog.debug("End Connection for ({}:{} | {})".format(client_writer.get_extra_info("peername")[0], client_writer.get_extra_info("peername")[1], aclient.bumper_jid)) + client_writer.close() + except Exception as e: + xmppserverlog.error("{}".format(e)) - clientaddr = client_writer.get_extra_info("peername") - xmppserverlog.info("New Connection from {}".format(clientaddr)) - task.add_done_callback(client_done) + clientaddr = client_writer.get_extra_info("peername") + xmppserverlog.debug("New Connection from {}:{}".format(clientaddr[0],clientaddr[1])) + task.add_done_callback(functools.partial(client_done, aclient)) - # def run(self, run_async=False): - # if run_async: - # xmppserverlog.debug("Starting XMPPServer Thread: 1") - # self.xmppthread = Thread(name="XMPPServer_Thread", target=self.run_server) - # self.xmppthread.setDaemon(True) - # self.xmppthread.start() - - # else: - # try: - # self.run_server() - # except KeyboardInterrupt: - # self.disconnect() - - # def run_server(self): - # logging.info("Starting XMPP Server at {}".format(self.address)) - # print("Starting XMPP Server at {}".format(self.address)) - - # # xmppserverlog.setLevel(logging.DEBUG) - - # # Set SSL Context - # self.ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH) - # self.ssl_ctx.load_cert_chain( - # certfile=bumper.server_cert, keyfile=bumper.server_key - # ) - - # self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM) - # self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) - - # try: - # self.socket.bind(self.address) - # self.socket.listen(5) - - # xmppserverlog.debug( - # "listening on {}:{}".format(self.address[0], self.address[1]) - # ) - # while not self.exit_flag: - # connection, client_address = self.socket.accept() - - # # disconnect any clients with this ip - # for client in self.clients: - # if client.address == client_address[0]: - # xmppserverlog.debug( - # "disconnecting existing client {} with resource {}".format( - # client.address, client.clientresource - # ) - # ) - # client._disconnect() - # self.remove_client_byip(client.address) - - # xmppserverlog.debug( - # "starting new client with ip {}".format(client_address[0]) - # ) - # thread_id = uuid.uuid4() - # client = XMPPAsyncClient(thread_id, connection, client_address) - # client.setDaemon(True) - # client.start() - # self.clients.append(client) - - # except PermissionError as e: - # if "bind" in e.strerror: - # xmppserverlog.exception( - # "Error binding XMPPServer, exiting. Try using a different hostname or IP - {}".format( - # e - # ) - # ) - # exit(1) - - # except Exception as e: - # xmppserverlog.exception("{}".format(e)) - # exit(1) - - # except KeyboardInterrupt as e: - # xmppserverlog.exception("{}".format(e)) - - # finally: - # connection.shutdown(socket.SHUT_RDWR) - # connection.close() - # self.disconnect() - # xmppserverlog.info("disconnecting") - - # self.socket.close() + except Exception as e: + xmppserverlog.error("{}".format(e)) def disconnect(self): try: @@ -140,41 +64,7 @@ class XMPPServer: xmppserverlog.debug("shutting down") except Exception as e: - xmppserverlog.exception("{}".format(e)) - - def remove_client_byip(self, ip): - for client in self.clients: - if client.address == ip: - xmppserverlog.debug( - "removing client from client list with ip {} and resource {}".format( - client.address, client.clientresource - ) - ) - client._disconnect() - self.clients.remove(client) - - def remove_client_byresource(self, resource): - for client in self.clients: - if str(client.clientresource).lower() == str(resource).lower(): - xmppserverlog.debug( - "removing client from client list with ip {} and resource {}".format( - client.address, client.clientresource - ) - ) - client._disconnect() - self.clients.remove(client) - - def remove_client_byuid(self, uid): - for client in self.clients: - if str(client.uid).lower() == str(uid).lower(): - xmppserverlog.debug( - "removing client from client list with ip {} and resource {}".format( - client.address, client.clientresource - ) - ) - client._disconnect() - self.clients.remove(client) - + xmppserverlog.error("{}".format(e)) class XMPPAsyncClient: IDLE = 0 @@ -186,6 +76,7 @@ class XMPPAsyncClient: UNKNOWN = 0 BOT = 1 CONTROLLER = 2 + _async_task = None def __init__(self, client_reader, client_writer): self.type = self.UNKNOWN @@ -204,20 +95,17 @@ class XMPPAsyncClient: async def handle_async_client(self): # xmppserverlog.info('client connected - {}'.format(self.address)) - #await self._set_state("READY") - await self._set_state("CONNECT") - #asyncio.Task(self.send_ping(30)) + await self._set_state("CONNECT") while True: - await asyncio.sleep(0.1) + await asyncio.sleep(0.05) if not self.state == self.DISCONNECT: data = await self.client_reader.read(4096) - # data = await asyncio.wait_for(client_reader.readline(), timeout=10.0) - if data is None: - xmppserverlog.warning("Received no data") + if data is None or data == b'': + xmppserverlog.debug("Received no data") # exit loop and disconnect return - - await self._parse_data(data) + else: + await self._parse_data(data) else: break @@ -228,25 +116,25 @@ class XMPPAsyncClient: try: # if not self.connection._closed: if self.log_sent_message: - xmppserverlog.debug("send {} - {}".format(self.address, command)) - # self.connection.send(command.encode()) + xmppserverlog.debug("send to ({}:{} | {}) - {}".format(self.address[0], self.address[1], self.bumper_jid, command)) + self.client_writer.write(command.encode()) await self.client_writer.drain() except BrokenPipeError as e: - xmppserverlog.debug("{}".format(e)) + #xmppserverlog.debug("{}".format(e)) await self._set_state("DISCONNECT") except ConnectionResetError as e: - xmppserverlog.debug("{}".format(e)) + #xmppserverlog.debug("{}".format(e)) await self._set_state("DISCONNECT") except ConnectionAbortedError as e: - xmppserverlog.debug("{}".format(e)) + #xmppserverlog.debug("{}".format(e)) await self._set_state("DISCONNECT") except OSError as e: - xmppserverlog.debug("{}".format(e)) + xmppserverlog.error("{}".format(e)) except Exception as e: xmppserverlog.exception("{}".format(e)) @@ -265,7 +153,7 @@ class XMPPAsyncClient: self.client_writer.close() except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) async def _tag_strip_uri(self, tag): try: @@ -274,7 +162,7 @@ class XMPPAsyncClient: return tag except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) async def _set_state(self, state): try: @@ -286,7 +174,7 @@ class XMPPAsyncClient: ) ) - xmppserverlog.debug("{} state: {}".format(self.address, state)) + xmppserverlog.debug("({}:{} | {}) state: {}".format(self.address[0],self.address[1],self.bumper_jid, state)) self.state = new_state @@ -294,7 +182,7 @@ class XMPPAsyncClient: await self._disconnect() except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) async def _handle_ctl(self, xml, data): try: @@ -365,13 +253,13 @@ class XMPPAsyncClient: if client.type == self.BOT: if client.uid.lower() in ctl_to.lower(): - xmppserverlog.info( + xmppserverlog.debug( "Sending ctl to bot: {}".format(rxmlstring) ) - client.send(rxmlstring) + await client.send(rxmlstring) except Exception as e: - xmppserverlog.exception("{}".format(e)) + xmppserverlog.error("{}".format(e)) async def _handle_ping(self, xml, data): try: @@ -387,37 +275,40 @@ class XMPPAsyncClient: pingto = xml.get("to") pingfrom = self.bumper_jid - xml.attrib["from"] = pingfrom + xml.attrib["from"] = pingfrom pingstring = ET.tostring(xml).decode("utf-8") # clean up string to remove namespaces added by ET pingstring = pingstring.replace("xmlns:ns0=", "xmlns=") pingstring = pingstring.replace("ns0:", "") - pingstring = pingstring.replace('iq xmlns="com:ctl"', "iq") - pingstring = pingstring.replace(" 2: authcode = saslauth[2] - - if not self.uid.startswith("fuid"): - # Need sample data to see details here + + if self.devclass: # if there is a devclass it is a bot bumper.bot_add(self.uid, self.uid, self.devclass, "atom", "eco-legacy") self.type = self.BOT - xmppserverlog.info("bot authenticated {}".format(self.uid)) + xmppserverlog.debug("bot authenticated {}".format(self.uid)) # Send response await self.send( '' @@ -751,7 +639,7 @@ class XMPPAsyncClient: self.bumper_jid = "{}@{}.ecorobot.net/atom".format( self.uid, self.devclass ) - xmppserverlog.debug("new bot {}".format(self.uid)) + xmppserverlog.debug("new bot ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid)) res = '{}'.format( xml.get("id"), self.bumper_jid ) @@ -761,18 +649,14 @@ class XMPPAsyncClient: self.bumper_jid = "{}@{}/{}".format( self.uid, XMPPServer.server_id, self.clientresource ) - xmppserverlog.debug( - "new client {} using resource {}".format( - self.uid, self.clientresource - ) - ) + xmppserverlog.debug("new client ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid)) res = '{}'.format( xml.get("id"), self.bumper_jid ) else: self.name = "XMPP_Client_{}_{}".format(self.uid, self.address) self.bumper_jid = "{}@{}".format(self.uid, XMPPServer.server_id) - xmppserverlog.debug("new client {}".format(self.uid)) + xmppserverlog.debug("new client ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid)) res = '{}'.format( xml.get("id"), self.bumper_jid ) @@ -788,7 +672,7 @@ class XMPPAsyncClient: res = ''.format(xml.get("id")) await self._set_state("READY") await self.send(res) - asyncio.Task(self.send_ping(30)) + asyncio.Task(self.schedule_ping(30)) except Exception as e: @@ -799,7 +683,7 @@ class XMPPAsyncClient: if len(xml) and xml[0].tag == "status": xmppserverlog.debug( - "bot presence {} ".format(ET.tostring(xml, encoding="utf-8")) + "bot presence {} ".format(ET.tostring(xml, encoding="utf-8").decode("utf-8")) ) # Most likely a bot, possibly hello world in text @@ -823,15 +707,15 @@ class XMPPAsyncClient: else: xmppserverlog.debug( - "client presence - {} ".format(ET.tostring(xml, encoding="utf-8")) + "client presence - {} ".format(ET.tostring(xml, encoding="utf-8").decode("utf-8")) ) if xml.get("type") == "available": xmppserverlog.debug( "client presence available - {} ".format( - ET.tostring(xml, encoding="utf-8") - ) + ET.tostring(xml, encoding="utf-8").decode("utf-8")) ) + # Send dummy return await self.send( ' dummy '.format(self.bumper_jid) @@ -839,9 +723,8 @@ class XMPPAsyncClient: elif xml.get("type") == "unavailable": xmppserverlog.debug( "client presence unavailable (DISCONNECT) - {} ".format( - ET.tostring(xml, encoding="utf-8") - ) - ) + ET.tostring(xml, encoding="utf-8").decode("utf-8")) + ) await self._set_state("DISCONNECT") else: @@ -880,8 +763,8 @@ class XMPPAsyncClient: if item.tag == "iq": if self.log_incoming_data: xmppserverlog.debug( - "from {} - {}".format( - self.address, + "from ({}:{} | {}) - {}".format( + self.address[0],self.address[1],self.bumper_jid, str( ET.tostring(item, encoding="utf-8").decode( "utf-8" diff --git a/start_bumper.py b/start_bumper.py index 2503f44..a01cd81 100644 --- a/start_bumper.py +++ b/start_bumper.py @@ -6,11 +6,16 @@ import sys, socket import time import platform import os -os.environ['PYTHONASYNCIODEBUG'] = '1' +#os.environ['PYTHONASYNCIODEBUG'] = '1' # Uncomment to enable ASYNCIODEBUG import asyncio -def main(): +async def main(): + try: + loop = asyncio.get_event_loop() + except: + loop = asyncio.new_event_loop() + args = sys.argv listen_host = "" @@ -20,12 +25,13 @@ def main(): level=logging.DEBUG, format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s", ) + loop.set_debug(True) # Set asyncio loop to debug + #logging.getLogger("asyncio").setLevel(logging.DEBUG) # Show debug asyncio logs (disabled in init, uncomment for debugging asyncio) else: logging.basicConfig( level=logging.INFO, format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(message)s", ) - # format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s") if "--listen" in args: listen_host = args[args.index("--listen") + 1] @@ -53,54 +59,39 @@ def main(): conf_address_8007, usessl=False, helperbot=mqtt_helperbot ) - try: - loop = asyncio.get_event_loop() - except: - loop = asyncio.new_event_loop() - # Start web servers - loop.set_debug(True) - conf_server.confserver_app() + conf_server.confserver_app() + task_conf_server = asyncio.create_task(conf_server.start_server()) + bumper.bumperlog.debug("task_conf_server added") + await task_conf_server + conf_server_2.confserver_app() - asyncio.ensure_future(conf_server.start_server(),loop=loop) - asyncio.ensure_future(conf_server_2.start_server(),loop=loop) + task_conf_server2 = asyncio.create_task(conf_server_2.start_server()) + bumper.bumperlog.debug("task_conf_server2 added") + await task_conf_server2 # Start MQTT Server - asyncio.ensure_future(mqtt_server.broker_coro()) + task_mqtt_server = asyncio.create_task(mqtt_server.broker_coro()) + bumper.bumperlog.debug("task_mqtt_server added") + await task_mqtt_server # Start MQTT Helperbot - asyncio.ensure_future(mqtt_helperbot.start_helper_bot()) + task_mqtt_helperbot = asyncio.create_task(mqtt_helperbot.start_helper_bot()) + bumper.bumperlog.debug("task_mqtt_helperbot added") + await task_mqtt_helperbot # Start XMPP Server - asyncio.ensure_future(xmpp_server.async_server()) - - loop.run_forever() - - - # start xmpp server on port 5223 (sync) - #xmpp_server.run(run_async=True) # Start in new thread - - # start mqtt server on port 8883 (async) - #mqtt_server.run(run_async=True) # Start in new thread - - #time.sleep(1.5) # Wait for broker startup - - # start mqtt_helperbot (async) - #mqtt_helperbot.run(run_async=True) # Start in new thread - - # start conf server on port 443 (async) - Used for most https calls - #conf_server.run(run_async=True) # Start in new thread - - # start conf server on port 8007 (async) - Used for a load balancer request - #conf_server_2.run(run_async=True) # Start in new thread + task_xmpp_server = asyncio.create_task(xmpp_server.async_server()) + bumper.bumperlog.debug("task_xmpp_server added") + await task_xmpp_server while True: try: - time.sleep(30) + await asyncio.sleep(30) bumper.revoke_expired_tokens() - disconnected_clients = bumper.get_disconnected_xmpp_clients() - for client in disconnected_clients: - xmpp_server.remove_client_byuid(client["userid"]) + #disconnected_clients = bumper.get_disconnected_xmpp_clients() + #for client in disconnected_clients: + # xmpp_server.remove_client_byuid(client["userid"]) except KeyboardInterrupt: bumper.bumperlog.info("Bumper Exiting - Keyboard Interrupt") @@ -109,4 +100,4 @@ def main(): if __name__ == "__main__": - main() + asyncio.run(main())