Compare commits

...
Sign in to create a new pull request.

17 commits

Author SHA1 Message Date
Robert Resch
a3e9463821
fix ip 2021-06-02 14:02:36 +02:00
Robert Resch
2f61eef73e
forward to correct recipient 2021-06-02 14:00:25 +02:00
Robert Resch
5583b4a99b
update Dockerfile 2021-06-02 12:07:39 +02:00
Robert Resch
e7d1819857
mqtt connect on port 443 2021-06-02 12:07:27 +02:00
Robert Resch
7cbac37c4d
proxy forward also other command than p2p 2021-05-21 13:46:05 +02:00
Robert Resch
fdb934980d
Merge branch 'master' into wip-proxyRequests 2021-04-29 22:53:53 +02:00
Brian Martin
fe8e4b0615 show <bytes content> for binary/images 2020-01-20 12:49:19 -05:00
Brian Martin
f97c07ac7b handle images and bytes
handle images and bytes
2020-01-20 12:36:36 -05:00
Brian Martin
346895eff3 fix tests 2020-01-20 12:07:05 -05:00
Brian Martin
c2bcb13769 support android app
handle form data sent by android apps
2020-01-20 12:01:01 -05:00
Brian Martin
a4a77708e5 add bumper api to proxymode
add bumper api to proxymode
2020-01-20 11:27:22 -05:00
Brian Martin
904167a8c5 proxymode db
move settings to db and rework code to support
2020-01-20 10:30:10 -05:00
Brian Martin
bdc452b99c reset connection status at startup
reset bot and client connections to false for xmpp/mqtt at startup
2020-01-20 08:34:08 -05:00
Brian Martin
3b7b31af47 log client topic subscription
log client topic subscription
2020-01-20 08:24:31 -05:00
Brian Martin
3e8ad68eb2 slim helperbot subscriptions
replace subscriptions with consolidated iot/#
2020-01-20 08:15:02 -05:00
Brian Martin
aca40e0e57 mqtt proxy
- mqtt and confserver proxy mode working
2020-01-20 00:51:22 -05:00
Brian Martin
d428c5019b confserver proxy
- app is proxying correctly
2020-01-18 13:43:25 -05:00
5 changed files with 527 additions and 54 deletions

View file

@ -22,11 +22,15 @@ RUN apk add build-base
FROM base FROM base
COPY . /bumper COPY requirements.txt /requirements.txt
WORKDIR /bumper
# install required python packages # install required python packages
RUN pip3 install -r requirements.txt RUN pip3 install -r requirements.txt
WORKDIR /bumper
# Copy only required folders instead of all
COPY create_certs/ create_certs/
COPY bumper/ bumper/
ENTRYPOINT ["python3", "-m", "bumper"] ENTRYPOINT ["python3", "-m", "bumper"]

View file

@ -58,6 +58,7 @@ bumper_debug = strtobool(os.environ.get("BUMPER_DEBUG")) or False
use_auth = False use_auth = False
token_validity_seconds = 3600 # 1 hour token_validity_seconds = 3600 # 1 hour
db = None db = None
bumper_proxy_mode = strtobool(os.environ.get("BUMPER_PROXY_MODE")) or False
mqtt_server = None mqtt_server = None
mqtt_helperbot = None mqtt_helperbot = None
@ -116,6 +117,16 @@ else:
# 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")
if not log_to_stdout: if not log_to_stdout:
@ -184,7 +195,7 @@ else:
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
# iptables -A PREROUTING -t nat -i wlp0s20f3 -p tcp --dport 443 -j REDIRECT --to-port 8883
mqtt_listen_port = 8883 mqtt_listen_port = 8883
conf1_listen_port = 443 conf1_listen_port = 443
conf2_listen_port = 8007 conf2_listen_port = 8007
@ -192,6 +203,11 @@ xmpp_listen_port = 5223
async def start(): 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: try:
loop = asyncio.get_event_loop() loop = asyncio.get_event_loop()
@ -238,6 +254,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())
@ -252,7 +270,22 @@ async def start():
await asyncio.sleep(0.1) await asyncio.sleep(0.1)
# Start web servers # 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=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))
@ -382,12 +415,17 @@ def main(argv=None):
help="announce address to bots on checkin", help="announce address to bots on checkin",
) )
parser.add_argument("--debug", action="store_true", help="enable debug logs") 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) args = parser.parse_args(args=argv)
if args.debug: if args.debug:
bumper_debug = True bumper_debug = True
if args.proxy_mode:
global bumper_proxy_mode
bumper_proxy_mode = True
if args.listen: if args.listen:
bumper_listen = args.listen bumper_listen = args.listen

View file

@ -12,12 +12,14 @@ from bumper import plugins
from datetime import datetime, timedelta from datetime import datetime, timedelta
import asyncio import asyncio
from aiohttp import web from aiohttp import web
import aiohttp
import aiohttp_jinja2 import aiohttp_jinja2
import jinja2 import jinja2
import uuid 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 (
@ -40,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):
@ -54,6 +57,115 @@ class ConfServer:
def get_milli_time(self, timetoconvert): def get_milli_time(self, timetoconvert):
return int(round(timetoconvert * 1000)) 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}) - <BYTES CONTENT>")
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): def confserver_app(self):
self.app = web.Application(loop=asyncio.get_event_loop(), middlewares=[ self.app = web.Application(loop=asyncio.get_event_loop(), middlewares=[
self.log_all_requests, self.log_all_requests,
@ -110,9 +222,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)
@ -129,6 +243,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:
@ -236,41 +380,63 @@ class ConfServer:
postbody = None postbody = None
response = await handler(request) response = await handler(request)
if not "application/octet-stream" in response.content_type: if response:
logall = { if not "application/octet-stream" in response.content_type:
"request": { try:
"route_name": f"{request.match_info.route.name}", logall = {
"method": f"{request.method}", "request": {
"path": f"{request.path}", "route_name": f"{request.match_info.route.name}",
"query_string": f"{request.query_string}", "method": f"{request.method}",
"raw_path": f"{request.raw_path}", "path": f"{request.path}",
"raw_headers": f'{",".join(map("{}".format, request.raw_headers))}', "query_string": f"{request.query_string}",
"body": f"{postbody}", "raw_path": f"{request.raw_path}",
}, "raw_headers": f'{",".join(map("{}".format, request.raw_headers))}',
"body": f"{postbody}",
},
"response": { "response": {
"response_body": f"{json.loads(response.body)}", #"response_body": f"{json.loads(response.body)}",
"status": f"{response.status}", "response_body": f"{json.loads(response.text)}",
} "status": f"{response.status}",
} }
else: }
logall = { except Exception as e:
"request": { logall = {
"route_name": f"{request.match_info.route.name}", "request": {
"method": f"{request.method}", "route_name": f"{request.match_info.route.name}",
"path": f"{request.path}", "method": f"{request.method}",
"query_string": f"{request.query_string}", "path": f"{request.path}",
"raw_path": f"{request.raw_path}", "query_string": f"{request.query_string}",
"raw_headers": f'{",".join(map("{}".format, request.raw_headers))}', "raw_path": f"{request.raw_path}",
"body": f"{postbody}", "raw_headers": f'{",".join(map("{}".format, request.raw_headers))}',
}, "body": f"{postbody}",
},
"response": { "response": {
"status": f"{response.status}", #"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 return response

View file

@ -32,9 +32,65 @@ def db_get():
db.table("clients", cache_size=0) db.table("clients", cache_size=0)
db.table("bots", cache_size=0) db.table("bots", cache_size=0)
db.table("tokens", cache_size=0) db.table("tokens", cache_size=0)
db.table("config_proxymode", cache_size=0)
return db 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.143.26"},
]
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): def user_add(userid):
newuser = BumperUser() newuser = BumperUser()
@ -289,6 +345,11 @@ def bot_remove(did):
if bot: if bot:
bots.remove(doc_ids=[bot.doc_id]) 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): def bot_get(did):
bots = db_get().table("bots") 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)) bumperlog.info("Adding new client with resource {}".format(newclient.resource))
client_full_upsert(newclient.asdict()) 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): def client_remove(resource):
clients = db_get().table("clients") clients = db_get().table("clients")
client = client_get(resource) client = client_get(resource)

View file

@ -3,9 +3,12 @@
import logging import logging
import asyncio import asyncio
import os import os
from typing import Dict
import hbmqtt import hbmqtt
import websockets
from hbmqtt.broker import Broker from hbmqtt.broker import Broker
from hbmqtt.client import MQTTClient from hbmqtt.client import MQTTClient, ConnectException
from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2 from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2
import pkg_resources import pkg_resources
import time import time
@ -14,10 +17,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:
@ -44,9 +58,7 @@ class MQTTHelperBot:
) )
await self.Client.subscribe( await self.Client.subscribe(
[ [
("iot/p2p/+/+/+/+/helperbot/bumper/helperbot/+/+/+", QOS_0), ("iot/#", QOS_0),
("iot/p2p/+", QOS_0),
("iot/atr/+", QOS_0),
] ]
) )
@ -216,10 +228,126 @@ class MQTTServer:
except Exception as e: except Exception as e:
mqttserverlog.exception("{}".format(e)) mqttserverlog.exception("{}".format(e))
class BumperProxyModeMQTTClient(MQTTClient):
eco_helper_names: Dict[str, str] = {}
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}")
topic = message.topic
ttopic = topic.split("/")
if ttopic[1] == "p2p":
self.eco_helper_names[ttopic[10]] = ttopic[3]
ttopic[3] = "proxyhelper"
topic = "/".join(ttopic)
proxymodelog.info(f"MQTT Proxy Client - Converted Topic From {message.topic} TO {topic}")
proxymodelog.info(
f"MQTT Proxy Client - Proxy Forward Message to Robot - Topic: {topic} - Message: {msgdata.encode()}")
await bumper.mqtt_helperbot.Client.publish(
topic, 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: Dict[str, BumperProxyModeMQTTClient] = {}
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 +360,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
@ -257,6 +387,32 @@ 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
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}:443",
)
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("/")
userid = didsplit[0] userid = didsplit[0]
@ -327,6 +483,20 @@ 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)
]
)
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): async def on_broker_client_connected(self, client_id):
didsplit = str(client_id).split("@") didsplit = str(client_id).split("@")
@ -336,16 +506,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].eco_helper_names.pop(ttopic[10], "")
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 +600,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 +611,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