Forward to correct recipient #121

Open
edenhaus wants to merge 10 commits from edenhaus/wip-proxyRequests into wip-proxyRequests
6 changed files with 101 additions and 55 deletions

View file

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

View file

@ -26,10 +26,13 @@ def strtobool(strbool):
# os.environ['PYTHONASYNCIODEBUG'] = '1' # Uncomment to enable ASYNCIODEBUG
bumper_dir = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir))
log_to_stdout = os.environ.get("LOG_TO_STDOUT")
# Set defaults from environment variables first
# Folders
logs_dir = os.environ.get("BUMPER_LOGS") or os.path.join(bumper_dir, "logs")
os.makedirs(logs_dir, exist_ok=True) # Ensure logs directory exists or create
if not log_to_stdout:
logs_dir = os.environ.get("BUMPER_LOGS") or os.path.join(bumper_dir, "logs")
os.makedirs(logs_dir, exist_ok=True) # Ensure logs directory exists or create
data_dir = os.environ.get("BUMPER_DATA") or os.path.join(bumper_dir, "data")
os.makedirs(data_dir, exist_ok=True) # Ensure data directory exists or create
certs_dir = os.environ.get("BUMPER_CERTS") or os.path.join(bumper_dir, "certs")
@ -81,27 +84,36 @@ logformat = logging.Formatter(
)
bumperlog = logging.getLogger("bumper")
bumper_rotate = RotatingFileHandler("logs/bumper.log", maxBytes=5000000, backupCount=5)
bumper_rotate.setFormatter(logformat)
bumperlog.addHandler(bumper_rotate)
if not log_to_stdout:
bumper_rotate = RotatingFileHandler("logs/bumper.log", maxBytes=5000000, backupCount=5)
bumper_rotate.setFormatter(logformat)
bumperlog.addHandler(bumper_rotate)
else:
bumperlog.addHandler(logging.StreamHandler(sys.stdout))
# Override the logging level
# bumperlog.setLevel(logging.INFO)
confserverlog = logging.getLogger("confserver")
conf_rotate = RotatingFileHandler(
"logs/confserver.log", maxBytes=5000000, backupCount=5
)
conf_rotate.setFormatter(logformat)
confserverlog.addHandler(conf_rotate)
if not log_to_stdout:
conf_rotate = RotatingFileHandler(
"logs/confserver.log", maxBytes=5000000, backupCount=5
)
conf_rotate.setFormatter(logformat)
confserverlog.addHandler(conf_rotate)
else:
confserverlog.addHandler(logging.StreamHandler(sys.stdout))
# Override the logging level
# confserverlog.setLevel(logging.INFO)
mqttserverlog = logging.getLogger("mqttserver")
mqtt_rotate = RotatingFileHandler(
"logs/mqttserver.log", maxBytes=5000000, backupCount=5
)
mqtt_rotate.setFormatter(logformat)
mqttserverlog.addHandler(mqtt_rotate)
if not log_to_stdout:
mqtt_rotate = RotatingFileHandler(
"logs/mqttserver.log", maxBytes=5000000, backupCount=5
)
mqtt_rotate.setFormatter(logformat)
mqttserverlog.addHandler(mqtt_rotate)
else:
mqttserverlog.addHandler(logging.StreamHandler(sys.stdout))
# Override the logging level
# mqttserverlog.setLevel(logging.INFO)
@ -117,53 +129,73 @@ proxymodelog.addHandler(proxymode_rotate)
### Additional MQTT Logs
translog = logging.getLogger("transitions")
translog.addHandler(mqtt_rotate)
if not log_to_stdout:
translog.addHandler(mqtt_rotate)
else:
translog.addHandler(logging.StreamHandler(sys.stdout))
translog.setLevel(logging.CRITICAL + 1) # Ignore this logger
logging.getLogger("passlib").setLevel(logging.CRITICAL + 1) # Ignore this logger
brokerlog = logging.getLogger("hbmqtt.broker")
#brokerlog.setLevel(
# logging.CRITICAL + 1
#) # Ignore this logger #There are some sublogs that could be set if needed (.plugins)
brokerlog.addHandler(mqtt_rotate)
if not log_to_stdout:
brokerlog.addHandler(mqtt_rotate)
else:
brokerlog.addHandler(logging.StreamHandler(sys.stdout))
protolog = logging.getLogger("hbmqtt.mqtt.protocol")
#protolog.setLevel(
# logging.CRITICAL + 1
#) # Ignore this logger
protolog.addHandler(mqtt_rotate)
if not log_to_stdout:
protolog.addHandler(mqtt_rotate)
else:
protolog.addHandler(logging.StreamHandler(sys.stdout))
clientlog = logging.getLogger("hbmqtt.client")
#clientlog.setLevel(logging.CRITICAL + 1) # Ignore this logger
clientlog.addHandler(mqtt_rotate)
if not log_to_stdout:
clientlog.addHandler(mqtt_rotate)
else:
clientlog.addHandler(logging.StreamHandler(sys.stdout))
helperbotlog = logging.getLogger("helperbot")
helperbot_rotate = RotatingFileHandler(
"logs/helperbot.log", maxBytes=5000000, backupCount=5
)
helperbot_rotate.setFormatter(logformat)
helperbotlog.addHandler(helperbot_rotate)
if not log_to_stdout:
helperbot_rotate = RotatingFileHandler(
"logs/helperbot.log", maxBytes=5000000, backupCount=5
)
helperbot_rotate.setFormatter(logformat)
helperbotlog.addHandler(helperbot_rotate)
else:
helperbotlog.addHandler(logging.StreamHandler(sys.stdout))
# Override the logging level
# helperbotlog.setLevel(logging.INFO)
boterrorlog = logging.getLogger("boterror")
boterrorlog_rotate = RotatingFileHandler(
"logs/boterror.log", maxBytes=5000000, backupCount=5
)
boterrorlog_rotate.setFormatter(logformat)
boterrorlog.addHandler(boterrorlog_rotate)
if not log_to_stdout:
boterrorlog_rotate = RotatingFileHandler(
"logs/boterror.log", maxBytes=5000000, backupCount=5
)
boterrorlog_rotate.setFormatter(logformat)
boterrorlog.addHandler(boterrorlog_rotate)
else:
boterrorlog.addHandler(logging.StreamHandler(sys.stdout))
# Override the logging level
# boterrorlog.setLevel(logging.INFO)
xmppserverlog = logging.getLogger("xmppserver")
xmpp_rotate = RotatingFileHandler(
"logs/xmppserver.log", maxBytes=5000000, backupCount=5
)
xmpp_rotate.setFormatter(logformat)
xmppserverlog.addHandler(xmpp_rotate)
if not log_to_stdout:
xmpp_rotate = RotatingFileHandler(
"logs/xmppserver.log", maxBytes=5000000, backupCount=5
)
xmpp_rotate.setFormatter(logformat)
xmppserverlog.addHandler(xmpp_rotate)
else:
xmppserverlog.addHandler(logging.StreamHandler(sys.stdout))
# Override the logging level
# xmppserverlog.setLevel(logging.INFO)
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

View file

@ -50,7 +50,7 @@ def config_proxyMode_defaults():
{"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"},
{"type":"mqtt_server","host":"mq-ww.ecouser.net","ip":"47.254.143.26"},
]
opendb = db_get()
with opendb:

View file

@ -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
@ -226,7 +229,9 @@ class MQTTServer:
mqttserverlog.exception("{}".format(e))
class BumperProxyModeMQTTClient(MQTTClient):
ecohelpername = ""
eco_helper_names: Dict[str, str] = {}
async def _connect_coro(self): #Override default to ignore ssl verification
kwargs = dict()
@ -278,7 +283,7 @@ class BumperProxyModeMQTTClient(MQTTClient):
reader = StreamReaderAdapter(conn_reader)
writer = StreamWriterAdapter(conn_writer)
elif scheme in ('ws', 'wss'):
websocket = await websockets.connect(
websocket = await websockets.connect(
self.session.broker_uri,
subprotocols=['mqtt'],
loop=self._loop,
@ -322,21 +327,25 @@ class BumperProxyModeMQTTClient(MQTTClient):
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()}")
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(
ttopic_comb, msgdata.encode(), QOS_0
topic, msgdata.encode(), QOS_0
)
except Exception as e:
proxymodelog.error(f"MQTT Proxy Client - get_msg Exception - {e}")
class BumperMQTTServer_Plugin:
proxyclients = {}
proxyclients: Dict[str, BumperProxyModeMQTTClient] = {}
def __init__(self, context):
self.context = context
try:
@ -395,7 +404,7 @@ class BumperMQTTServer_Plugin:
try:
await self.proxyclients[client_id].connect(
f"mqtts://{username}:{password}@{mqtt_server}:8883",
f"mqtts://{username}:{password}@{mqtt_server}:443",
)
except Exception as e:
mqttserverlog.error(f"MQTT Proxy Mode - Exception connecting with proxy to ecovacs - {e}")
@ -514,7 +523,7 @@ class BumperMQTTServer_Plugin:
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[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:

View file

@ -29,7 +29,7 @@ class portal_api_iot(plugins.ConfServerApp):
try:
json_body = json.loads(await request.text())
randomid = "".join(random.sample(string.ascii_letters, 6))
randomid = "".join(random.sample(string.ascii_letters, 4))
did = ""
if "toId" in json_body: # Its a command
did = json_body["toId"]

View file

@ -11,4 +11,5 @@ Bumper has a number of environment variables to help with custom deployments and
| BUMPER_KEY | {full path to bumper.key location} | The private server key (bumper.key) to be used by the Bumper server |
| BUMPER_LOGS | {full path to logs directory} | The directory where logs should be stored |
| BUMPER_DATA | {full path to data directory} | The directory where persistent data should be stored (bumper.db) |
| BUMPER_DEBUG | true | Run Bumper with debug mode/logging |
| BUMPER_DEBUG | true | Run Bumper with debug mode/logging |
| LOG_TO_STDOUT | true | Instead of logging to logs/, logs to to STDOUT |