diff --git a/bumper/__init__.py b/bumper/__init__.py index ff03f00..3b15523 100644 --- a/bumper/__init__.py +++ b/bumper/__init__.py @@ -55,6 +55,7 @@ bumper_debug = strtobool(os.environ.get("BUMPER_DEBUG")) or False use_auth = False token_validity_seconds = 3600 # 1 hour db = None +bumper_proxy_mode = strtobool(os.environ.get("BUMPER_PROXY_MODE")) or False mqtt_server = None mqtt_helperbot = None @@ -104,6 +105,16 @@ mqttserverlog.addHandler(mqtt_rotate) # Override the logging level # mqttserverlog.setLevel(logging.INFO) +proxymodelog = logging.getLogger("proxymode") +proxymode_rotate = RotatingFileHandler( + "logs/proxymode.log", maxBytes=5000000, backupCount=5 +) +proxymode_rotate.setFormatter(logformat) +proxymodelog.addHandler(proxymode_rotate) + +# Override the logging level +# mqttserverlog.setLevel(logging.INFO) + ### Additional MQTT Logs translog = logging.getLogger("transitions") translog.addHandler(mqtt_rotate) @@ -160,6 +171,11 @@ xmpp_listen_port = 5223 async def start(): + #config_proxyMode_deleteTable() #delete existing proxymode table + + #Reset xmpp/mqtt to false in database for bots and clients + bot_reset_connectionStatus() + client_reset_connectionStatus() try: loop = asyncio.get_event_loop() @@ -206,6 +222,8 @@ async def start(): # Start MQTT Server asyncio.create_task(mqtt_server.broker_coro()) + await asyncio.sleep(0.5) #Wait half a sec for broker to start + # Start MQTT Helperbot asyncio.create_task(mqtt_helperbot.start_helper_bot()) @@ -220,7 +238,22 @@ async def start(): await asyncio.sleep(0.1) # Start web servers - conf_server.confserver_app() + if bumper_proxy_mode: + bumperlog.info("Proxy Mode Enabled") + if config_proxyMode_countEntries() == 0: # check if proxymode servers are entered + bumperlog.info("Proxy Mode - No Servers, Loading Defaults (US)") + config_proxyMode_defaults() # set defaults if 0 + + configproxy = config_proxyMode_getall() + cntentries = len(configproxy) + proxymodelog.info(f"Loaded {cntentries} entries from proxyconfig") + for entry in configproxy: + proxymodelog.info(f"Config Entry {entry}") + + conf_server.confserver_proxy_app() + else: + conf_server.confserver_app() + asyncio.create_task(conf_server.start_site(conf_server.app, address=bumper_listen, port=conf1_listen_port, usessl=True)) asyncio.create_task(conf_server.start_site(conf_server.app, address=bumper_listen, port=conf2_listen_port, usessl=False)) @@ -350,12 +383,17 @@ def main(argv=None): help="announce address to bots on checkin", ) parser.add_argument("--debug", action="store_true", help="enable debug logs") + parser.add_argument("--proxy-mode", action="store_true", help="enable proxy mode") args = parser.parse_args(args=argv) if args.debug: bumper_debug = True + if args.proxy_mode: + global bumper_proxy_mode + bumper_proxy_mode = True + if args.listen: bumper_listen = args.listen diff --git a/bumper/confserver.py b/bumper/confserver.py index 8b9cec2..dfe348b 100644 --- a/bumper/confserver.py +++ b/bumper/confserver.py @@ -12,12 +12,14 @@ from bumper import plugins from datetime import datetime, timedelta import asyncio from aiohttp import web +import aiohttp import aiohttp_jinja2 import jinja2 import uuid import xml.etree.ElementTree as ET + class aiohttp_filter(logging.Filter): def filter(self, record): if ( @@ -40,6 +42,7 @@ logging.getLogger("aiohttp.access").addFilter( aiohttp_filter() ) # Add logging filter above to aiohttp.access +proxymodelog = logging.getLogger("proxymode") class ConfServer: def __init__(self, address, usessl=False): @@ -54,6 +57,115 @@ class ConfServer: def get_milli_time(self, timetoconvert): return int(round(timetoconvert * 1000)) + def confserver_proxy_app(self): + + self.app = web.Application(middlewares=[ + self.log_all_requests, + ]) + aiohttp_jinja2.setup(self.app, loader=jinja2.FileSystemLoader(os.path.join(bumper.bumper_dir,"bumper","web","templates"))) + + + self.app.add_routes( + [ + + web.get("/bot/remove/{did}", self.handle_RemoveBot, name='remove-bot'), + web.get("/client/remove/{resource}", self.handle_RemoveClient, name='remove-client'), + web.get("/restart_{service}", self.handle_RestartService, name='restart-service'), + web.route("*", "/{path:.*}", self.handle_proxy, name="confserver_proxy"), + ] + ) + + return self.app + + async def handle_proxy(self, request): + + try: + ecoresp = "" + + server_port = 443 #default to 443 + if "_SSLProtocolTransport" != type(request.transport).__name__ and "_SelectorSocketTransport" != type(request.transport).__name__: #check not ssl transport class + if "_extra" in request.transport: + if "sockname" in request.transport._extra: + server_port = request.transport._extra["sockname"][1] + + if request.raw_path == "/": + return await self.handle_base(request) + if request.raw_path == "/lookup.do": + return await self.handle_lookup(request) #use bumper to handle lookup so bot gets Bumper IP and not Ecovacs + + matchproxy = bumper.config_proxyMode_getServerIP("app", request.host) + + if matchproxy: + proxymodelog.info(f"Matched {request.host} to entry in proxyconfig!") + ecorequest = f"{request.scheme}://{matchproxy}" + else: + proxymodelog.info(f"No match for {request.host} in proxyconfig!") + if "ecovacs.com" in request.host: + proxymodelog.info(f"ecovacs.com in {request.host} defaulting to ecovacs.com IP!") + matchproxy = bumper.config_proxyMode_getServerIP("app", "ecovacs.com") + ecorequest = f"{request.scheme}://{matchproxy}" + elif "ecouser.net" in request.host: + proxymodelog.info(f"ecouser.net in {request.host} defaulting to ecouser.net IP!") + matchproxy = bumper.config_proxyMode_getServerIP("app", "ecouser.net") + ecorequest = f"{request.scheme}://{matchproxy}" + else: + proxymodelog.info(f"No matches for {request.host} defaulting to ecovacs.com IP!") + matchproxy = bumper.config_proxyMode_getServerIP("app", "ecovacs.com") + ecorequest = f"{request.scheme}://{matchproxy}" + + if server_port != 443: + ecorequest = f"{ecorequest}:{server_port}" + + proxymodelog.info(f"{request.host} - {ecorequest}") + ecorequest = f"{ecorequest}{request.path_qs}" + requestheaders = {'host': request.host} + async with aiohttp.ClientSession(headers=requestheaders, connector=aiohttp.TCPConnector(verify_ssl=False)) as session: + if request.content.total_bytes > 0: + proxymodelog.info(f"HTTP Proxy Request to EcoVacs (body=true) (host:{request.host}) - {ecorequest} - {request._read_bytes}") + if request.content_type == "application/x-www-form-urlencoded": # android apps use form + fdata = await request.post() + async with session.request(request.method, ecorequest, data=fdata) as resp: + ecoresp = await resp.text() + proxymodelog.info(f"HTTP Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - {ecoresp}") + else: # handle json + jdata = request._read_bytes.decode('utf8') + jdata = json.loads(jdata) + async with session.request(request.method, ecorequest, json=jdata) as resp: + ecoresp = await resp.text() + proxymodelog.info(f"HTTP Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - {ecoresp}") + + else: + proxymodelog.info(f"HTTP Proxy Request to EcoVacs (body=false) (host:{request.host}) - {ecorequest}") + async with session.request(request.method, ecorequest) as resp: + if resp.content_type == "application/octet-stream": + ecoresp = await resp.read() + proxymodelog.info(f"HTTP Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - ") + else: + ecoresp = await resp.text() + proxymodelog.info(f"HTTP Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - {ecoresp}") + + if resp.status == 200: + if resp.content_type == "application/json": + ecoresp = json.loads(ecoresp) + return web.json_response(ecoresp) + elif resp.content_type == "application/octet-stream": + return web.Response(body=ecoresp) + else: + return web.Response(text=ecoresp) + + else: + return web.Response(text=ecoresp) + + except asyncio.CancelledError as e: + proxymodelog.error(f"Request cancelled or timeout - {ecorequest} - {jdata}") + return web.Response(text="") + pass + + except Exception as e: + proxymodelog.exception("{}".format(e)) + return web.Response(text="") + + def confserver_app(self): self.app = web.Application(loop=asyncio.get_event_loop(), middlewares=[ self.log_all_requests, @@ -110,9 +222,11 @@ class ConfServer: async def start_site(self, app, address='localhost', port=8080, usessl=False): + runner = web.AppRunner(app) self.runners.append(runner) await runner.setup() + if usessl: ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH) ssl_ctx.load_cert_chain(bumper.server_cert, bumper.server_key) @@ -129,6 +243,36 @@ class ConfServer: ) await site.start() + + def start_site_thread(self, app, address='localhost', port=8080, usessl=False): + #test for new thread and loop + loop = asyncio.new_event_loop() + asyncio.set_event_loop(loop) + + runner = web.AppRunner(app) + self.runners.append(runner) + #await runner.setup() + loop.run_until_complete(runner.setup()) #for thread test + + if usessl: + ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH) + ssl_ctx.load_cert_chain(bumper.server_cert, bumper.server_key) + site = web.TCPSite( + runner, + host=address, + port=port, + ssl_context=ssl_ctx, + ) + + else: + site = web.TCPSite( + runner, host=address, port=port + ) + + #await site.start() + #for thread test + loop.run_until_complete(site.start()) + loop.run_forever() async def start_server(self): try: @@ -236,41 +380,63 @@ class ConfServer: postbody = None response = await handler(request) - if not "application/octet-stream" in response.content_type: - logall = { - "request": { - "route_name": f"{request.match_info.route.name}", - "method": f"{request.method}", - "path": f"{request.path}", - "query_string": f"{request.query_string}", - "raw_path": f"{request.raw_path}", - "raw_headers": f'{",".join(map("{}".format, request.raw_headers))}', - "body": f"{postbody}", - }, + if response: + if not "application/octet-stream" in response.content_type: + try: + logall = { + "request": { + "route_name": f"{request.match_info.route.name}", + "method": f"{request.method}", + "path": f"{request.path}", + "query_string": f"{request.query_string}", + "raw_path": f"{request.raw_path}", + "raw_headers": f'{",".join(map("{}".format, request.raw_headers))}', + "body": f"{postbody}", + }, - "response": { - "response_body": f"{json.loads(response.body)}", - "status": f"{response.status}", - } - } - else: - logall = { - "request": { - "route_name": f"{request.match_info.route.name}", - "method": f"{request.method}", - "path": f"{request.path}", - "query_string": f"{request.query_string}", - "raw_path": f"{request.raw_path}", - "raw_headers": f'{",".join(map("{}".format, request.raw_headers))}', - "body": f"{postbody}", - }, + "response": { + #"response_body": f"{json.loads(response.body)}", + "response_body": f"{json.loads(response.text)}", + "status": f"{response.status}", + } + } + except Exception as e: + logall = { + "request": { + "route_name": f"{request.match_info.route.name}", + "method": f"{request.method}", + "path": f"{request.path}", + "query_string": f"{request.query_string}", + "raw_path": f"{request.raw_path}", + "raw_headers": f'{",".join(map("{}".format, request.raw_headers))}', + "body": f"{postbody}", + }, - "response": { - "status": f"{response.status}", - } - } + "response": { + #"response_body": f"{json.loads(response.body)}", + "response_body": f"{(response.text)}", + "status": f"{response.status}", + } + } - confserverlog.debug(json.dumps(logall)) + else: + logall = { + "request": { + "route_name": f"{request.match_info.route.name}", + "method": f"{request.method}", + "path": f"{request.path}", + "query_string": f"{request.query_string}", + "raw_path": f"{request.raw_path}", + "raw_headers": f'{",".join(map("{}".format, request.raw_headers))}', + "body": f"{postbody}", + }, + + "response": { + "status": f"{response.status}", + } + } + + confserverlog.debug(json.dumps(logall)) return response diff --git a/bumper/db.py b/bumper/db.py index 71c158a..dfa001e 100644 --- a/bumper/db.py +++ b/bumper/db.py @@ -32,9 +32,65 @@ def db_get(): db.table("clients", cache_size=0) db.table("bots", cache_size=0) db.table("tokens", cache_size=0) + db.table("config_proxymode", cache_size=0) return db +def config_proxyMode_deleteTable(): + opendb = db_get() + opendb.purge_table("config_proxymode") + +def config_proxyMode_defaults(): + defaults = [ + {"type":"app","host":"gl-us-api.ecovacs.com","ip":"47.252.51.29","match":"gl-"}, + {"type":"app","host":"gl-us-openapi.ecovacs.com","ip":"47.252.51.29"}, + {"type":"app","host":"portal-ww.ecouser.net","ip":"47.88.66.164","match":"portal-"}, + {"type":"app","host":"bigdata-northamerica.ecovacs.com","ip":"47.88.66.111"}, + {"type":"app","host":"bigdata-international.ecovacs.com","ip":"47.88.132.151","match":"bigdata-"}, + {"type":"app","host":"eco-us-api.ecovacs.com","ip":"47.89.135.130","match":"eco-"}, + {"type":"app","host":"ecovacs.com","ip":"47.90.210.46"}, + {"type":"app","host":"ecouser.net","ip":"116.62.93.217"}, + {"type":"mqtt_server","host":"mq-ww.ecouser.net","ip":"47.254.52.46"}, + ] + opendb = db_get() + with opendb: + config = opendb.table("config_proxymode") + config.insert_multiple(defaults) + +def config_proxyMode_getServerIP(type, host): + opendb = db_get() + with opendb: + proxyconfig = opendb.table("config_proxymode") + proxy = Query() + if type == "mqtt_server": + entry = proxyconfig.get((proxy.type == type)) + + else: + entry = proxyconfig.get((proxy.type == type) & (proxy.host == host)) + + if entry: + return entry["ip"] + + else: + proxylist = proxyconfig.search(Query()) + for proxy in proxylist: # check for sub matches + if "match" in proxy: + if proxy["match"] in host: + return proxy["ip"] + return None + +def config_proxyMode_countEntries(): + opendb = db_get() + with opendb: + config = opendb.table("config_proxymode") + return len(config) + +def config_proxyMode_getall(): + opendb = db_get() + with opendb: + config = opendb.table("config_proxymode") + return config.search(Query()) + def user_add(userid): newuser = BumperUser() @@ -289,6 +345,11 @@ def bot_remove(did): if bot: bots.remove(doc_ids=[bot.doc_id]) +def bot_reset_connectionStatus(): + bots = db_get().table("bots") + for bot in bots: + bot_set_mqtt(bot["did"], False) + bot_set_xmpp(bot["did"], False) def bot_get(did): bots = db_get().table("bots") @@ -345,6 +406,12 @@ def client_add(userid, realm, resource): bumperlog.info("Adding new client with resource {}".format(newclient.resource)) client_full_upsert(newclient.asdict()) +def client_reset_connectionStatus(): + clients = db_get().table("clients") + for client in clients: + client_set_mqtt(client["resource"], False) + client_set_xmpp(client["resource"], False) + def client_remove(resource): clients = db_get().table("clients") client = client_get(resource) diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py index f53d1f8..f6f61ed 100644 --- a/bumper/mqttserver.py +++ b/bumper/mqttserver.py @@ -14,10 +14,21 @@ import json from datetime import datetime, timedelta import bumper from passlib.apps import custom_app_context as pwd_context +import ssl +import tempfile + +from urllib.parse import urlparse, urlunparse +from hbmqtt.mqtt.protocol.client_handler import ClientProtocolHandler +from hbmqtt.adapters import StreamReaderAdapter, StreamWriterAdapter, WebSocketsReader, WebSocketsWriter +from websockets.uri import InvalidURI +from websockets.exceptions import InvalidHandshake +from hbmqtt.mqtt.protocol.handler import ProtocolHandlerException +from hbmqtt.mqtt.connack import CONNECTION_ACCEPTED helperbotlog = logging.getLogger("helperbot") boterrorlog = logging.getLogger("boterror") mqttserverlog = logging.getLogger("mqttserver") +proxymodelog = logging.getLogger("proxymode") class MQTTHelperBot: @@ -44,9 +55,7 @@ class MQTTHelperBot: ) await self.Client.subscribe( [ - ("iot/p2p/+/+/+/+/helperbot/bumper/helperbot/+/+/+", QOS_0), - ("iot/p2p/+", QOS_0), - ("iot/atr/+", QOS_0), + ("iot/#", QOS_0), ] ) @@ -216,10 +225,120 @@ class MQTTServer: except Exception as e: mqttserverlog.exception("{}".format(e)) +class BumperProxyModeMQTTClient(MQTTClient): + ecohelpername = "" + async def _connect_coro(self): #Override default to ignore ssl verification + kwargs = dict() + + # Decode URI attributes + uri_attributes = urlparse(self.session.broker_uri) + scheme = uri_attributes.scheme + secure = True if scheme in ('mqtts', 'wss') else False + self.session.username = self.session.username if self.session.username else uri_attributes.username + self.session.password = self.session.password if self.session.password else uri_attributes.password + self.session.remote_address = uri_attributes.hostname + self.session.remote_port = uri_attributes.port + if scheme in ('mqtt', 'mqtts') and not self.session.remote_port: + self.session.remote_port = 8883 if scheme == 'mqtts' else 1883 + if scheme in ('ws', 'wss') and not self.session.remote_port: + self.session.remote_port = 443 if scheme == 'wss' else 80 + if scheme in ('ws', 'wss'): + # Rewrite URI to conform to https://tools.ietf.org/html/rfc6455#section-3 + uri = (scheme, self.session.remote_address + ":" + str(self.session.remote_port), uri_attributes[2], + uri_attributes[3], uri_attributes[4], uri_attributes[5]) + self.session.broker_uri = urlunparse(uri) + # Init protocol handler + #if not self._handler: + self._handler = ClientProtocolHandler(self.plugins_manager, loop=self._loop) + + if secure: + sc = ssl.create_default_context( + ssl.Purpose.SERVER_AUTH, + cafile=self.session.cafile, + capath=self.session.capath, + cadata=self.session.cadata) + if 'certfile' in self.config and 'keyfile' in self.config: + sc.load_cert_chain(self.config['certfile'], self.config['keyfile']) + if 'check_hostname' in self.config and isinstance(self.config['check_hostname'], bool): + sc.check_hostname = self.config['check_hostname'] + + sc.verify_mode = ssl.CERT_NONE #Ignore verify of cert + kwargs['ssl'] = sc + + try: + reader = None + writer = None + self._connected_state.clear() + # Open connection + if scheme in ('mqtt', 'mqtts'): + conn_reader, conn_writer = \ + await asyncio.open_connection( + self.session.remote_address, + self.session.remote_port, loop=self._loop, **kwargs) + reader = StreamReaderAdapter(conn_reader) + writer = StreamWriterAdapter(conn_writer) + elif scheme in ('ws', 'wss'): + websocket = await websockets.connect( + self.session.broker_uri, + subprotocols=['mqtt'], + loop=self._loop, + extra_headers=self.extra_headers, + **kwargs) + reader = WebSocketsReader(websocket) + writer = WebSocketsWriter(websocket) + # Start MQTT protocol + self._handler.attach(self.session, reader, writer) + return_code = await self._handler.mqtt_connect() + if return_code is not CONNECTION_ACCEPTED: + self.session.transitions.disconnect() + self.logger.warning("Connection rejected with code '%s'" % return_code) + exc = ConnectException("Connection rejected by broker") + exc.return_code = return_code + raise exc + else: + # Handle MQTT protocol + await self._handler.start() + self.session.transitions.connect() + self._connected_state.set() + self.logger.debug("connected to %s:%s" % (self.session.remote_address, self.session.remote_port)) + return return_code + except InvalidURI as iuri: + self.logger.warning("connection failed: invalid URI '%s'" % self.session.broker_uri) + self.session.transitions.disconnect() + raise ConnectException("connection failed: invalid URI '%s'" % self.session.broker_uri, iuri) + except InvalidHandshake as ihs: + self.logger.warning("connection failed: invalid websocket handshake") + self.session.transitions.disconnect() + raise ConnectException("connection failed: invalid websocket handshake", ihs) + except (ProtocolHandlerException, ConnectionError, OSError) as e: + self.logger.warning("MQTT connection failed: %r" % e) + self.session.transitions.disconnect() + raise ConnectException(e) + + async def get_msg(self): + try: + while self._connected_state._value: + message = await self.deliver_message() + msgdata = str(message.data.decode("utf-8")) + + proxymodelog.info(f"MQTT Proxy Client - Message Received From Ecovacs - Topic: {message.topic} - Message: {msgdata}") + ttopic = message.topic.split("/") + self.ecohelpername = ttopic[3] + ttopic[3] = "proxyhelper" + ttopic_comb = "/".join(ttopic) + proxymodelog.info(f"MQTT Proxy Client - Converted Topic From {message.topic} TO {ttopic_comb}") + proxymodelog.info(f"MQTT Proxy Client - Proxy Forward Message to Helperbot - Topic: {ttopic_comb} - Message: {msgdata.encode()}") + await bumper.mqtt_helperbot.Client.publish( + ttopic_comb, msgdata.encode(), QOS_0 + ) + + except Exception as e: + proxymodelog.error(f"MQTT Proxy Client - get_msg Exception - {e}") class BumperMQTTServer_Plugin: + proxyclients = {} def __init__(self, context): - self.context = context + self.context = context try: self.auth_config = self.context.config["auth"] self._users = dict() @@ -232,6 +351,8 @@ class BumperMQTTServer_Plugin: except Exception as e: mqttserverlog.exception("{}".format(e)) + + async def authenticate(self, *args, **kwargs): authenticated = False @@ -257,6 +378,32 @@ class BumperMQTTServer_Plugin: mqttserverlog.info(f"Bumper Authentication Success - Bot - SN: {username} - DID: {didsplit[0]} - Class: {tmpbotdetail[0]}") authenticated = True + if authenticated and bumper.bumper_proxy_mode: + mqtt_server = bumper.config_proxyMode_getServerIP("mqtt_server","") + if mqtt_server: + proxymodelog.info(f"MQTT Proxy Mode - Using server {mqtt_server}") + else: + proxymodelog.error(f"MQTT Proxy Mode - No server found! Load defaults or set mqtt_server in config_proxymode table!") + proxymodelog.exception(f"MQTT Proxy Mode - Exiting due to no MQTT Server configured!") + exit(1) + + proxymodelog.info(f"MQTT Proxy Mode - Proxy Bot to MQTT - Client_id: {client_id} - Username: {username}") + + self.proxyclients[client_id] = BumperProxyModeMQTTClient( + client_id=client_id, config={"check_hostname": False} + ) + + try: + await self.proxyclients[client_id].connect( + f"mqtts://{username}:{password}@{mqtt_server}:8883", + ) + except Exception as e: + mqttserverlog.error(f"MQTT Proxy Mode - Exception connecting with proxy to ecovacs - {e}") + pass + proxymodelog.info(f"MQTT Proxy Mode - Proxy Bot Connected - Client_id: {client_id}") + asyncio.create_task(self.proxyclients[client_id].get_msg()) + + else: tmpclientdetail = str(didsplit[1]).split("/") userid = didsplit[0] @@ -327,6 +474,20 @@ class BumperMQTTServer_Plugin: except FileNotFoundError: self.context.logger.warning(f"Password file {password_file} not found") + async def on_broker_client_subscribed(self, client_id, topic, qos): + if bumper.bumper_proxy_mode: #if proxy mode, also subscribe on ecovacs server + if client_id in self.proxyclients: + await self.proxyclients[client_id].subscribe( + [ + (topic, qos) + ] + ) + else: + proxymodelog.info(f"MQTT Proxy Mode - New MQTT Topic Subscription - Client: {client_id} - Topic: {topic}") + + #return + #pass + async def on_broker_client_connected(self, client_id): didsplit = str(client_id).split("@") @@ -336,16 +497,39 @@ class BumperMQTTServer_Plugin: bumper.bot_set_mqtt(bot["did"], True) return - clientresource = didsplit[1].split("/")[1] - client = bumper.client_get(clientresource) - if client: - bumper.client_set_mqtt(client["resource"], True) - return + if len(didsplit) > 1: + clientresource = didsplit[1].split("/")[1] + client = bumper.client_get(clientresource) + if client: + bumper.client_set_mqtt(client["resource"], True) + return async def on_broker_message_received(self, client_id, message): - self.handle_helperbot_msg(client_id, message) + await self.handle_helperbot_msg(client_id, message) - def handle_helperbot_msg(self, client_id, message): + async def handle_helperbot_msg(self, client_id, message): + if bumper.bumper_proxy_mode: + if client_id in self.proxyclients: + msgdata = str(message.data.decode("utf-8")) + if not str(message.topic).split("/")[3] == "proxyhelper": # if from proxyhelper, don't send back to ecovacs...yet + if str(message.topic).split("/")[6] == "proxyhelper": + ttopic = message.topic.split("/") + ttopic[6] = self.proxyclients[client_id].ecohelpername + ttopic_join = "/".join(ttopic) + proxymodelog.info(f"MQTT Proxy Client - Bot Message Converted Topic From {message.topic} TO {ttopic_join} with message: {msgdata}") + else: + ttopic_join = message.topic + proxymodelog.info(f"MQTT Proxy Client - Bot Message From {ttopic_join} with message: {msgdata}") + + try: + # Send back to ecovacs + proxymodelog.info(f"MQTT Proxy Client - Proxy Forward Message to Ecovacs - Topic: {ttopic_join} - Message: {msgdata.encode()}") + await self.proxyclients[client_id].publish( + ttopic_join, msgdata.encode(), message.qos + ) + except Exception as e: + proxymodelog.error(f"MQTT Proxy Client - Forwarding to Ecovacs Exception - {e}") + if str(message.topic).split("/")[6] == "helperbot": # Response to command @@ -407,6 +591,10 @@ class BumperMQTTServer_Plugin: async def on_broker_client_disconnected(self, client_id): + if bumper.bumper_proxy_mode: + if client_id in self.proxyclients: + await self.proxyclients[client_id].disconnect() + didsplit = str(client_id).split("@") bot = bumper.bot_get(didsplit[0]) @@ -414,8 +602,9 @@ class BumperMQTTServer_Plugin: bumper.bot_set_mqtt(bot["did"], False) return - clientresource = didsplit[1].split("/")[1] - client = bumper.client_get(clientresource) - if client: - bumper.client_set_mqtt(client["resource"], False) - return + if len(didsplit) > 1: + clientresource = didsplit[1].split("/")[1] + client = bumper.client_get(clientresource) + if client: + bumper.client_set_mqtt(client["resource"], False) + return