mqtt proxy

- mqtt and confserver proxy mode working
This commit is contained in:
Brian Martin 2020-01-20 00:51:22 -05:00
parent d428c5019b
commit aca40e0e57
3 changed files with 272 additions and 31 deletions

View file

@ -105,6 +105,16 @@ mqttserverlog.addHandler(mqtt_rotate)
# Override the logging level # Override the logging level
# mqttserverlog.setLevel(logging.INFO) # 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 ### Additional MQTT Logs
translog = logging.getLogger("transitions") translog = logging.getLogger("transitions")
translog.addHandler(mqtt_rotate) translog.addHandler(mqtt_rotate)
@ -207,6 +217,8 @@ async def start():
# Start MQTT Server # Start MQTT Server
asyncio.create_task(mqtt_server.broker_coro()) asyncio.create_task(mqtt_server.broker_coro())
await asyncio.sleep(0.5) #Wait half a sec for broker to start
# Start MQTT Helperbot # Start MQTT Helperbot
asyncio.create_task(mqtt_helperbot.start_helper_bot()) asyncio.create_task(mqtt_helperbot.start_helper_bot())
@ -225,6 +237,7 @@ async def start():
conf_server.confserver_proxy_app() conf_server.confserver_proxy_app()
else: else:
conf_server.confserver_app() 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=conf1_listen_port, usessl=True))
asyncio.create_task(conf_server.start_site(conf_server.app, address=bumper_listen, port=conf2_listen_port, usessl=False)) asyncio.create_task(conf_server.start_site(conf_server.app, address=bumper_listen, port=conf2_listen_port, usessl=False))

View file

@ -19,6 +19,7 @@ import uuid
import xml.etree.ElementTree as ET import xml.etree.ElementTree as ET
class aiohttp_filter(logging.Filter): class aiohttp_filter(logging.Filter):
def filter(self, record): def filter(self, record):
if ( if (
@ -41,6 +42,7 @@ logging.getLogger("aiohttp.access").addFilter(
aiohttp_filter() aiohttp_filter()
) # Add logging filter above to aiohttp.access ) # Add logging filter above to aiohttp.access
proxymodelog = logging.getLogger("proxymode")
class ConfServer: class ConfServer:
def __init__(self, address, usessl=False): def __init__(self, address, usessl=False):
@ -56,9 +58,11 @@ class ConfServer:
return int(round(timetoconvert * 1000)) return int(round(timetoconvert * 1000))
def confserver_proxy_app(self): def confserver_proxy_app(self):
self.app = web.Application(loop=asyncio.get_event_loop(), middlewares=[
self.app = web.Application(middlewares=[
self.log_all_requests, 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( self.app.add_routes(
@ -66,6 +70,8 @@ class ConfServer:
web.route("*", "/{path:.*}", self.handle_proxy, name="confserver_proxy"), web.route("*", "/{path:.*}", self.handle_proxy, name="confserver_proxy"),
] ]
) )
return self.app
async def handle_proxy(self, request): async def handle_proxy(self, request):
@ -88,10 +94,9 @@ class ConfServer:
server_port = request.transport._extra["sockname"][1] server_port = request.transport._extra["sockname"][1]
if request.raw_path == "/": if request.raw_path == "/":
return web.Response(text="Bumper in Proxy Mode") return await self.handle_base(request)
if request.raw_path == "/lookup.do": if request.raw_path == "/lookup.do":
return await self.handle_lookup(request) return await self.handle_lookup(request) #use bumper to handle lookup so bot gets Bumper IP and not Ecovacs
#ecorequest = f"{request.scheme}://{ecouser_net_ip}"
elif "ecovacs.com" in request.host: elif "ecovacs.com" in request.host:
ecorequest = f"{request.scheme}://{ecovacs_com_ip}" ecorequest = f"{request.scheme}://{ecovacs_com_ip}"
elif "ecouser.net" in request.host: elif "ecouser.net" in request.host:
@ -116,29 +121,38 @@ class ConfServer:
requestheaders = {'host': request.host} requestheaders = {'host': request.host}
async with aiohttp.ClientSession(headers=requestheaders, connector=aiohttp.TCPConnector(verify_ssl=False)) as session: async with aiohttp.ClientSession(headers=requestheaders, connector=aiohttp.TCPConnector(verify_ssl=False)) as session:
if request.content.total_bytes > 0: if request.content.total_bytes > 0:
confserverlog.debug(f"Proxy Request to EcoVacs (body=true) (host:{request.host}) - {ecorequest} - {request._read_bytes}") proxymodelog.info(f"HTTP 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 jdata = request._read_bytes.decode('utf8')
jdata = json.loads(jdata)
async with session.request(request.method, ecorequest, json=jdata) as resp: async with session.request(request.method, ecorequest, json=jdata) as resp:
ecoresp = await resp.text() 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: 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: async with session.request(request.method, ecorequest) as resp:
ecoresp = await resp.text() ecoresp = await asyncio.shield(resp.text())
ecoresp = ecoresp.replace("portal-ww.ecouser.net", ecouser_net_ip)
confserverlog.debug(f"Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - {ecoresp}") proxymodelog.info(f"HTTP Proxy Response from EcoVacs (URL: {ecorequest}) - (Status: {resp.status}) - {ecoresp}")
if resp.status == 200: if resp.status == 200:
ecoresp = json.loads(ecoresp) if resp.content_type == "application/json":
ecoresp = json.loads(ecoresp)
if resp.status == 200: return web.json_response(ecoresp)
return web.json_response(ecoresp) else:
return web.Response(text=ecoresp)
else: else:
return web.Response(text=ecoresp) 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: except Exception as e:
confserverlog.exception("{}".format(e)) proxymodelog.exception("{}".format(e))
return web.Response(text="") return web.Response(text="")
@ -198,9 +212,11 @@ class ConfServer:
async def start_site(self, app, address='localhost', port=8080, usessl=False): async def start_site(self, app, address='localhost', port=8080, usessl=False):
runner = web.AppRunner(app) runner = web.AppRunner(app)
self.runners.append(runner) self.runners.append(runner)
await runner.setup() await runner.setup()
if usessl: if usessl:
ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH) ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
ssl_ctx.load_cert_chain(bumper.server_cert, bumper.server_key) ssl_ctx.load_cert_chain(bumper.server_cert, bumper.server_key)
@ -217,6 +233,36 @@ class ConfServer:
) )
await site.start() 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): async def start_server(self):
try: try:

View file

@ -14,10 +14,21 @@ import json
from datetime import datetime, timedelta from datetime import datetime, timedelta
import bumper import bumper
from passlib.apps import custom_app_context as pwd_context 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") helperbotlog = logging.getLogger("helperbot")
boterrorlog = logging.getLogger("boterror") boterrorlog = logging.getLogger("boterror")
mqttserverlog = logging.getLogger("mqttserver") mqttserverlog = logging.getLogger("mqttserver")
proxymodelog = logging.getLogger("proxymode")
class MQTTHelperBot: class MQTTHelperBot:
@ -45,8 +56,10 @@ class MQTTHelperBot:
await self.Client.subscribe( await self.Client.subscribe(
[ [
("iot/p2p/+/+/+/+/helperbot/bumper/helperbot/+/+/+", QOS_0), ("iot/p2p/+/+/+/+/helperbot/bumper/helperbot/+/+/+", QOS_0),
("iot/p2p/+/+/+/+/+/+/+/+/+/+", QOS_0),
("iot/p2p/+", QOS_0), ("iot/p2p/+", QOS_0),
("iot/atr/+", QOS_0), ("iot/atr/+", QOS_0),
] ]
) )
@ -216,10 +229,120 @@ class MQTTServer:
except Exception as e: except Exception as e:
mqttserverlog.exception("{}".format(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: class BumperMQTTServer_Plugin:
proxyclients = {}
def __init__(self, context): def __init__(self, context):
self.context = context self.context = context
try: try:
self.auth_config = self.context.config["auth"] self.auth_config = self.context.config["auth"]
self._users = dict() self._users = dict()
@ -232,6 +355,8 @@ class BumperMQTTServer_Plugin:
except Exception as e: except Exception as e:
mqttserverlog.exception("{}".format(e)) mqttserverlog.exception("{}".format(e))
async def authenticate(self, *args, **kwargs): async def authenticate(self, *args, **kwargs):
authenticated = False authenticated = False
@ -256,6 +381,24 @@ class BumperMQTTServer_Plugin:
) )
mqttserverlog.info(f"Bumper Authentication Success - Bot - SN: {username} - DID: {didsplit[0]} - Class: {tmpbotdetail[0]}") mqttserverlog.info(f"Bumper Authentication Success - Bot - SN: {username} - DID: {didsplit[0]} - Class: {tmpbotdetail[0]}")
authenticated = True 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: else:
tmpclientdetail = str(didsplit[1]).split("/") tmpclientdetail = str(didsplit[1]).split("/")
@ -327,6 +470,17 @@ class BumperMQTTServer_Plugin:
except FileNotFoundError: except FileNotFoundError:
self.context.logger.warning(f"Password file {password_file} not found") 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): async def on_broker_client_connected(self, client_id):
didsplit = str(client_id).split("@") didsplit = str(client_id).split("@")
@ -336,16 +490,39 @@ class BumperMQTTServer_Plugin:
bumper.bot_set_mqtt(bot["did"], True) bumper.bot_set_mqtt(bot["did"], True)
return return
clientresource = didsplit[1].split("/")[1] if len(didsplit) > 1:
client = bumper.client_get(clientresource) clientresource = didsplit[1].split("/")[1]
if client: client = bumper.client_get(clientresource)
bumper.client_set_mqtt(client["resource"], True) if client:
return bumper.client_set_mqtt(client["resource"], True)
return
async def on_broker_message_received(self, client_id, message): 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": if str(message.topic).split("/")[6] == "helperbot":
# Response to command # Response to command
@ -407,6 +584,10 @@ class BumperMQTTServer_Plugin:
async def on_broker_client_disconnected(self, client_id): 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("@") didsplit = str(client_id).split("@")
bot = bumper.bot_get(didsplit[0]) bot = bumper.bot_get(didsplit[0])
@ -414,8 +595,9 @@ class BumperMQTTServer_Plugin:
bumper.bot_set_mqtt(bot["did"], False) bumper.bot_set_mqtt(bot["did"], False)
return return
clientresource = didsplit[1].split("/")[1] if len(didsplit) > 1:
client = bumper.client_get(clientresource) clientresource = didsplit[1].split("/")[1]
if client: client = bumper.client_get(clientresource)
bumper.client_set_mqtt(client["resource"], False) if client:
return bumper.client_set_mqtt(client["resource"], False)
return