Compare commits

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

11 commits

Author SHA1 Message Date
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
4 changed files with 509 additions and 49 deletions

View file

@ -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

View file

@ -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}) - <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):
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

View file

@ -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)

View file

@ -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