From aca40e0e574b6e64797565bd2e95ebf608c776f1 Mon Sep 17 00:00:00 2001 From: Brian Martin Date: Mon, 20 Jan 2020 00:51:22 -0500 Subject: [PATCH] mqtt proxy - mqtt and confserver proxy mode working --- bumper/__init__.py | 13 +++ bumper/confserver.py | 82 +++++++++++++---- bumper/mqttserver.py | 208 ++++++++++++++++++++++++++++++++++++++++--- 3 files changed, 272 insertions(+), 31 deletions(-) diff --git a/bumper/__init__.py b/bumper/__init__.py index 7fd97b7..8369477 100644 --- a/bumper/__init__.py +++ b/bumper/__init__.py @@ -105,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) @@ -207,6 +217,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()) @@ -225,6 +237,7 @@ async def start(): 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)) diff --git a/bumper/confserver.py b/bumper/confserver.py index 61fb666..234fd5a 100644 --- a/bumper/confserver.py +++ b/bumper/confserver.py @@ -19,6 +19,7 @@ import uuid import xml.etree.ElementTree as ET + class aiohttp_filter(logging.Filter): def filter(self, record): if ( @@ -41,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): @@ -56,9 +58,11 @@ class ConfServer: return int(round(timetoconvert * 1000)) def confserver_proxy_app(self): - self.app = web.Application(loop=asyncio.get_event_loop(), middlewares=[ + + 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( @@ -66,6 +70,8 @@ class ConfServer: web.route("*", "/{path:.*}", self.handle_proxy, name="confserver_proxy"), ] ) + + return self.app async def handle_proxy(self, request): @@ -88,10 +94,9 @@ class ConfServer: server_port = request.transport._extra["sockname"][1] if request.raw_path == "/": - return web.Response(text="Bumper in Proxy Mode") + return await self.handle_base(request) if request.raw_path == "/lookup.do": - return await self.handle_lookup(request) - #ecorequest = f"{request.scheme}://{ecouser_net_ip}" + return await self.handle_lookup(request) #use bumper to handle lookup so bot gets Bumper IP and not Ecovacs elif "ecovacs.com" in request.host: ecorequest = f"{request.scheme}://{ecovacs_com_ip}" elif "ecouser.net" in request.host: @@ -116,29 +121,38 @@ class ConfServer: requestheaders = {'host': request.host} async with aiohttp.ClientSession(headers=requestheaders, connector=aiohttp.TCPConnector(verify_ssl=False)) as session: if request.content.total_bytes > 0: - confserverlog.debug(f"Proxy Request to EcoVacs (body=true) (host:{request.host}) - {ecorequest} - {request._read_bytes}") - jdata = json.loads(request._read_bytes.decode('utf8').replace("'",'"')) #convert bytes to json for sending + proxymodelog.info(f"HTTP Proxy Request to EcoVacs (body=true) (host:{request.host}) - {ecorequest} - {request._read_bytes}") + 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() + ecoresp = ecoresp.replace("portal-ww.ecouser.net", ecouser_net_ip) + proxymodelog.info(f"HTTP Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - {ecoresp}") else: - confserverlog.debug(f"Proxy Request to EcoVacs (body=false) (host:{request.host}) - {ecorequest}") + proxymodelog.info(f"HTTP Proxy Request to EcoVacs (body=false) (host:{request.host}) - {ecorequest}") async with session.request(request.method, ecorequest) as resp: - ecoresp = await resp.text() - - confserverlog.debug(f"Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - {ecoresp}") - + ecoresp = await asyncio.shield(resp.text()) + ecoresp = ecoresp.replace("portal-ww.ecouser.net", ecouser_net_ip) + proxymodelog.info(f"HTTP Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - {ecoresp}") + if resp.status == 200: - ecoresp = json.loads(ecoresp) - - if resp.status == 200: - return web.json_response(ecoresp) + if resp.content_type == "application/json": + ecoresp = json.loads(ecoresp) + return web.json_response(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: - confserverlog.exception("{}".format(e)) + proxymodelog.exception("{}".format(e)) return web.Response(text="") @@ -198,9 +212,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) @@ -217,6 +233,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: diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py index f53d1f8..6529006 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: @@ -45,8 +56,10 @@ class MQTTHelperBot: await self.Client.subscribe( [ ("iot/p2p/+/+/+/+/helperbot/bumper/helperbot/+/+/+", QOS_0), + ("iot/p2p/+/+/+/+/+/+/+/+/+/+", QOS_0), ("iot/p2p/+", QOS_0), ("iot/atr/+", QOS_0), + ] ) @@ -216,10 +229,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 +355,8 @@ class BumperMQTTServer_Plugin: except Exception as e: mqttserverlog.exception("{}".format(e)) + + async def authenticate(self, *args, **kwargs): authenticated = False @@ -256,6 +381,24 @@ class BumperMQTTServer_Plugin: ) mqttserverlog.info(f"Bumper Authentication Success - Bot - SN: {username} - DID: {didsplit[0]} - Class: {tmpbotdetail[0]}") authenticated = True + mq_na_ip = "47.254.52.46" + if authenticated and bumper.bumper_proxy_mode: + 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}@{mq_na_ip}: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("/") @@ -327,6 +470,17 @@ 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) + ] + ) + #return + #pass + async def on_broker_client_connected(self, client_id): didsplit = str(client_id).split("@") @@ -336,16 +490,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 +584,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 +595,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