Compare commits
17 commits
master
...
edenhaus/w
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
a3e9463821 | ||
|
|
2f61eef73e | ||
|
|
5583b4a99b | ||
|
|
e7d1819857 | ||
|
|
7cbac37c4d | ||
|
|
fdb934980d | ||
|
|
fe8e4b0615 | ||
|
|
f97c07ac7b | ||
|
|
346895eff3 | ||
|
|
c2bcb13769 | ||
|
|
a4a77708e5 | ||
|
|
904167a8c5 | ||
|
|
bdc452b99c | ||
|
|
3b7b31af47 | ||
|
|
3e8ad68eb2 | ||
|
|
aca40e0e57 | ||
|
|
d428c5019b |
5 changed files with 527 additions and 54 deletions
10
Dockerfile
10
Dockerfile
|
|
@ -22,11 +22,15 @@ RUN apk add build-base
|
|||
|
||||
FROM base
|
||||
|
||||
COPY . /bumper
|
||||
|
||||
WORKDIR /bumper
|
||||
COPY requirements.txt /requirements.txt
|
||||
|
||||
# install required python packages
|
||||
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"]
|
||||
|
|
|
|||
|
|
@ -58,6 +58,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
|
||||
|
|
@ -116,6 +117,16 @@ else:
|
|||
# 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")
|
||||
if not log_to_stdout:
|
||||
|
|
@ -184,7 +195,7 @@ else:
|
|||
|
||||
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
|
||||
conf1_listen_port = 443
|
||||
conf2_listen_port = 8007
|
||||
|
|
@ -192,6 +203,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()
|
||||
|
|
@ -238,6 +254,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())
|
||||
|
||||
|
|
@ -252,7 +270,22 @@ async def start():
|
|||
await asyncio.sleep(0.1)
|
||||
|
||||
# Start web servers
|
||||
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))
|
||||
|
||||
|
|
@ -382,12 +415,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
|
||||
|
||||
|
|
|
|||
|
|
@ -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)
|
||||
|
|
@ -130,6 +244,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:
|
||||
confserverlog.info(
|
||||
|
|
@ -236,7 +380,9 @@ class ConfServer:
|
|||
postbody = None
|
||||
|
||||
response = await handler(request)
|
||||
if response:
|
||||
if not "application/octet-stream" in response.content_type:
|
||||
try:
|
||||
logall = {
|
||||
"request": {
|
||||
"route_name": f"{request.match_info.route.name}",
|
||||
|
|
@ -249,10 +395,30 @@ class ConfServer:
|
|||
},
|
||||
|
||||
"response": {
|
||||
"response_body": f"{json.loads(response.body)}",
|
||||
#"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": {
|
||||
#"response_body": f"{json.loads(response.body)}",
|
||||
"response_body": f"{(response.text)}",
|
||||
"status": f"{response.status}",
|
||||
}
|
||||
}
|
||||
|
||||
else:
|
||||
logall = {
|
||||
"request": {
|
||||
|
|
|
|||
67
bumper/db.py
67
bumper/db.py
|
|
@ -32,9 +32,65 @@ def db_get():
|
|||
db.table("clients", cache_size=0)
|
||||
db.table("bots", cache_size=0)
|
||||
db.table("tokens", cache_size=0)
|
||||
db.table("config_proxymode", cache_size=0)
|
||||
|
||||
return db
|
||||
|
||||
def config_proxyMode_deleteTable():
|
||||
opendb = db_get()
|
||||
opendb.purge_table("config_proxymode")
|
||||
|
||||
def config_proxyMode_defaults():
|
||||
defaults = [
|
||||
{"type":"app","host":"gl-us-api.ecovacs.com","ip":"47.252.51.29","match":"gl-"},
|
||||
{"type":"app","host":"gl-us-openapi.ecovacs.com","ip":"47.252.51.29"},
|
||||
{"type":"app","host":"portal-ww.ecouser.net","ip":"47.88.66.164","match":"portal-"},
|
||||
{"type":"app","host":"bigdata-northamerica.ecovacs.com","ip":"47.88.66.111"},
|
||||
{"type":"app","host":"bigdata-international.ecovacs.com","ip":"47.88.132.151","match":"bigdata-"},
|
||||
{"type":"app","host":"eco-us-api.ecovacs.com","ip":"47.89.135.130","match":"eco-"},
|
||||
{"type":"app","host":"ecovacs.com","ip":"47.90.210.46"},
|
||||
{"type":"app","host":"ecouser.net","ip":"116.62.93.217"},
|
||||
{"type":"mqtt_server","host":"mq-ww.ecouser.net","ip":"47.254.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):
|
||||
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)
|
||||
|
|
|
|||
|
|
@ -3,9 +3,12 @@
|
|||
import logging
|
||||
import asyncio
|
||||
import os
|
||||
from typing import Dict
|
||||
|
||||
import hbmqtt
|
||||
import websockets
|
||||
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
|
||||
import pkg_resources
|
||||
import time
|
||||
|
|
@ -14,10 +17,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 +58,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,8 +228,124 @@ class MQTTServer:
|
|||
except Exception as 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:
|
||||
proxyclients: Dict[str, BumperProxyModeMQTTClient] = {}
|
||||
def __init__(self, context):
|
||||
self.context = context
|
||||
try:
|
||||
|
|
@ -232,6 +360,8 @@ class BumperMQTTServer_Plugin:
|
|||
except Exception as e:
|
||||
mqttserverlog.exception("{}".format(e))
|
||||
|
||||
|
||||
|
||||
async def authenticate(self, *args, **kwargs):
|
||||
authenticated = False
|
||||
|
||||
|
|
@ -257,6 +387,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}: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:
|
||||
tmpclientdetail = str(didsplit[1]).split("/")
|
||||
userid = didsplit[0]
|
||||
|
|
@ -327,6 +483,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,6 +506,7 @@ class BumperMQTTServer_Plugin:
|
|||
bumper.bot_set_mqtt(bot["did"], True)
|
||||
return
|
||||
|
||||
if len(didsplit) > 1:
|
||||
clientresource = didsplit[1].split("/")[1]
|
||||
client = bumper.client_get(clientresource)
|
||||
if client:
|
||||
|
|
@ -343,9 +514,31 @@ class BumperMQTTServer_Plugin:
|
|||
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)
|
||||
|
||||
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}")
|
||||
|
||||
def handle_helperbot_msg(self, client_id, message):
|
||||
|
||||
if str(message.topic).split("/")[6] == "helperbot":
|
||||
# Response to command
|
||||
|
|
@ -407,6 +600,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,6 +611,7 @@ class BumperMQTTServer_Plugin:
|
|||
bumper.bot_set_mqtt(bot["did"], False)
|
||||
return
|
||||
|
||||
if len(didsplit) > 1:
|
||||
clientresource = didsplit[1].split("/")[1]
|
||||
client = bumper.client_get(clientresource)
|
||||
if client:
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue