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
|
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"]
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
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()
|
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
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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)
|
||||||
|
|
@ -130,6 +244,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:
|
||||||
confserverlog.info(
|
confserverlog.info(
|
||||||
|
|
@ -236,7 +380,9 @@ class ConfServer:
|
||||||
postbody = None
|
postbody = None
|
||||||
|
|
||||||
response = await handler(request)
|
response = await handler(request)
|
||||||
|
if response:
|
||||||
if not "application/octet-stream" in response.content_type:
|
if not "application/octet-stream" in response.content_type:
|
||||||
|
try:
|
||||||
logall = {
|
logall = {
|
||||||
"request": {
|
"request": {
|
||||||
"route_name": f"{request.match_info.route.name}",
|
"route_name": f"{request.match_info.route.name}",
|
||||||
|
|
@ -249,10 +395,30 @@ class ConfServer:
|
||||||
},
|
},
|
||||||
|
|
||||||
"response": {
|
"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}",
|
"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:
|
else:
|
||||||
logall = {
|
logall = {
|
||||||
"request": {
|
"request": {
|
||||||
|
|
|
||||||
67
bumper/db.py
67
bumper/db.py
|
|
@ -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)
|
||||||
|
|
|
||||||
|
|
@ -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,8 +228,124 @@ 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:
|
||||||
|
|
@ -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,6 +506,7 @@ class BumperMQTTServer_Plugin:
|
||||||
bumper.bot_set_mqtt(bot["did"], True)
|
bumper.bot_set_mqtt(bot["did"], True)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
if len(didsplit) > 1:
|
||||||
clientresource = didsplit[1].split("/")[1]
|
clientresource = didsplit[1].split("/")[1]
|
||||||
client = bumper.client_get(clientresource)
|
client = bumper.client_get(clientresource)
|
||||||
if client:
|
if client:
|
||||||
|
|
@ -343,9 +514,31 @@ class BumperMQTTServer_Plugin:
|
||||||
return
|
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)
|
||||||
|
|
||||||
|
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":
|
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,6 +611,7 @@ class BumperMQTTServer_Plugin:
|
||||||
bumper.bot_set_mqtt(bot["did"], False)
|
bumper.bot_set_mqtt(bot["did"], False)
|
||||||
return
|
return
|
||||||
|
|
||||||
|
if len(didsplit) > 1:
|
||||||
clientresource = didsplit[1].split("/")[1]
|
clientresource = didsplit[1].split("/")[1]
|
||||||
client = bumper.client_get(clientresource)
|
client = bumper.client_get(clientresource)
|
||||||
if client:
|
if client:
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue