Forward to correct recipient #121
6 changed files with 101 additions and 55 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"]
|
||||||
|
|
|
||||||
|
|
@ -26,10 +26,13 @@ def strtobool(strbool):
|
||||||
# os.environ['PYTHONASYNCIODEBUG'] = '1' # Uncomment to enable ASYNCIODEBUG
|
# os.environ['PYTHONASYNCIODEBUG'] = '1' # Uncomment to enable ASYNCIODEBUG
|
||||||
bumper_dir = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir))
|
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
|
# Set defaults from environment variables first
|
||||||
# Folders
|
# Folders
|
||||||
logs_dir = os.environ.get("BUMPER_LOGS") or os.path.join(bumper_dir, "logs")
|
if not log_to_stdout:
|
||||||
os.makedirs(logs_dir, exist_ok=True) # Ensure logs directory exists or create
|
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")
|
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
|
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")
|
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")
|
bumperlog = logging.getLogger("bumper")
|
||||||
bumper_rotate = RotatingFileHandler("logs/bumper.log", maxBytes=5000000, backupCount=5)
|
if not log_to_stdout:
|
||||||
bumper_rotate.setFormatter(logformat)
|
bumper_rotate = RotatingFileHandler("logs/bumper.log", maxBytes=5000000, backupCount=5)
|
||||||
bumperlog.addHandler(bumper_rotate)
|
bumper_rotate.setFormatter(logformat)
|
||||||
|
bumperlog.addHandler(bumper_rotate)
|
||||||
|
else:
|
||||||
|
bumperlog.addHandler(logging.StreamHandler(sys.stdout))
|
||||||
# Override the logging level
|
# Override the logging level
|
||||||
# bumperlog.setLevel(logging.INFO)
|
# bumperlog.setLevel(logging.INFO)
|
||||||
|
|
||||||
confserverlog = logging.getLogger("confserver")
|
confserverlog = logging.getLogger("confserver")
|
||||||
conf_rotate = RotatingFileHandler(
|
if not log_to_stdout:
|
||||||
"logs/confserver.log", maxBytes=5000000, backupCount=5
|
conf_rotate = RotatingFileHandler(
|
||||||
)
|
"logs/confserver.log", maxBytes=5000000, backupCount=5
|
||||||
conf_rotate.setFormatter(logformat)
|
)
|
||||||
confserverlog.addHandler(conf_rotate)
|
conf_rotate.setFormatter(logformat)
|
||||||
|
confserverlog.addHandler(conf_rotate)
|
||||||
|
else:
|
||||||
|
confserverlog.addHandler(logging.StreamHandler(sys.stdout))
|
||||||
# Override the logging level
|
# Override the logging level
|
||||||
# confserverlog.setLevel(logging.INFO)
|
# confserverlog.setLevel(logging.INFO)
|
||||||
|
|
||||||
mqttserverlog = logging.getLogger("mqttserver")
|
mqttserverlog = logging.getLogger("mqttserver")
|
||||||
mqtt_rotate = RotatingFileHandler(
|
if not log_to_stdout:
|
||||||
"logs/mqttserver.log", maxBytes=5000000, backupCount=5
|
mqtt_rotate = RotatingFileHandler(
|
||||||
)
|
"logs/mqttserver.log", maxBytes=5000000, backupCount=5
|
||||||
mqtt_rotate.setFormatter(logformat)
|
)
|
||||||
mqttserverlog.addHandler(mqtt_rotate)
|
mqtt_rotate.setFormatter(logformat)
|
||||||
|
mqttserverlog.addHandler(mqtt_rotate)
|
||||||
|
else:
|
||||||
|
mqttserverlog.addHandler(logging.StreamHandler(sys.stdout))
|
||||||
# Override the logging level
|
# Override the logging level
|
||||||
# mqttserverlog.setLevel(logging.INFO)
|
# mqttserverlog.setLevel(logging.INFO)
|
||||||
|
|
||||||
|
|
@ -117,53 +129,73 @@ proxymodelog.addHandler(proxymode_rotate)
|
||||||
|
|
||||||
### Additional MQTT Logs
|
### Additional MQTT Logs
|
||||||
translog = logging.getLogger("transitions")
|
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
|
translog.setLevel(logging.CRITICAL + 1) # Ignore this logger
|
||||||
logging.getLogger("passlib").setLevel(logging.CRITICAL + 1) # Ignore this logger
|
logging.getLogger("passlib").setLevel(logging.CRITICAL + 1) # Ignore this logger
|
||||||
brokerlog = logging.getLogger("hbmqtt.broker")
|
brokerlog = logging.getLogger("hbmqtt.broker")
|
||||||
#brokerlog.setLevel(
|
#brokerlog.setLevel(
|
||||||
# logging.CRITICAL + 1
|
# logging.CRITICAL + 1
|
||||||
#) # Ignore this logger #There are some sublogs that could be set if needed (.plugins)
|
#) # 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 = logging.getLogger("hbmqtt.mqtt.protocol")
|
||||||
#protolog.setLevel(
|
#protolog.setLevel(
|
||||||
# logging.CRITICAL + 1
|
# logging.CRITICAL + 1
|
||||||
#) # Ignore this logger
|
#) # 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 = logging.getLogger("hbmqtt.client")
|
||||||
#clientlog.setLevel(logging.CRITICAL + 1) # Ignore this logger
|
#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")
|
helperbotlog = logging.getLogger("helperbot")
|
||||||
helperbot_rotate = RotatingFileHandler(
|
if not log_to_stdout:
|
||||||
"logs/helperbot.log", maxBytes=5000000, backupCount=5
|
helperbot_rotate = RotatingFileHandler(
|
||||||
)
|
"logs/helperbot.log", maxBytes=5000000, backupCount=5
|
||||||
helperbot_rotate.setFormatter(logformat)
|
)
|
||||||
helperbotlog.addHandler(helperbot_rotate)
|
helperbot_rotate.setFormatter(logformat)
|
||||||
|
helperbotlog.addHandler(helperbot_rotate)
|
||||||
|
else:
|
||||||
|
helperbotlog.addHandler(logging.StreamHandler(sys.stdout))
|
||||||
# Override the logging level
|
# Override the logging level
|
||||||
# helperbotlog.setLevel(logging.INFO)
|
# helperbotlog.setLevel(logging.INFO)
|
||||||
|
|
||||||
boterrorlog = logging.getLogger("boterror")
|
boterrorlog = logging.getLogger("boterror")
|
||||||
boterrorlog_rotate = RotatingFileHandler(
|
if not log_to_stdout:
|
||||||
"logs/boterror.log", maxBytes=5000000, backupCount=5
|
boterrorlog_rotate = RotatingFileHandler(
|
||||||
)
|
"logs/boterror.log", maxBytes=5000000, backupCount=5
|
||||||
boterrorlog_rotate.setFormatter(logformat)
|
)
|
||||||
boterrorlog.addHandler(boterrorlog_rotate)
|
boterrorlog_rotate.setFormatter(logformat)
|
||||||
|
boterrorlog.addHandler(boterrorlog_rotate)
|
||||||
|
else:
|
||||||
|
boterrorlog.addHandler(logging.StreamHandler(sys.stdout))
|
||||||
# Override the logging level
|
# Override the logging level
|
||||||
# boterrorlog.setLevel(logging.INFO)
|
# boterrorlog.setLevel(logging.INFO)
|
||||||
|
|
||||||
xmppserverlog = logging.getLogger("xmppserver")
|
xmppserverlog = logging.getLogger("xmppserver")
|
||||||
xmpp_rotate = RotatingFileHandler(
|
if not log_to_stdout:
|
||||||
"logs/xmppserver.log", maxBytes=5000000, backupCount=5
|
xmpp_rotate = RotatingFileHandler(
|
||||||
)
|
"logs/xmppserver.log", maxBytes=5000000, backupCount=5
|
||||||
xmpp_rotate.setFormatter(logformat)
|
)
|
||||||
xmppserverlog.addHandler(xmpp_rotate)
|
xmpp_rotate.setFormatter(logformat)
|
||||||
|
xmppserverlog.addHandler(xmpp_rotate)
|
||||||
|
else:
|
||||||
|
xmppserverlog.addHandler(logging.StreamHandler(sys.stdout))
|
||||||
# Override the logging level
|
# Override the logging level
|
||||||
# xmppserverlog.setLevel(logging.INFO)
|
# xmppserverlog.setLevel(logging.INFO)
|
||||||
|
|
||||||
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
|
||||||
|
|
|
||||||
|
|
@ -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":"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":"ecovacs.com","ip":"47.90.210.46"},
|
||||||
{"type":"app","host":"ecouser.net","ip":"116.62.93.217"},
|
{"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()
|
opendb = db_get()
|
||||||
with opendb:
|
with opendb:
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
@ -226,7 +229,9 @@ class MQTTServer:
|
||||||
mqttserverlog.exception("{}".format(e))
|
mqttserverlog.exception("{}".format(e))
|
||||||
|
|
||||||
class BumperProxyModeMQTTClient(MQTTClient):
|
class BumperProxyModeMQTTClient(MQTTClient):
|
||||||
ecohelpername = ""
|
|
||||||
|
eco_helper_names: Dict[str, str] = {}
|
||||||
|
|
||||||
async def _connect_coro(self): #Override default to ignore ssl verification
|
async def _connect_coro(self): #Override default to ignore ssl verification
|
||||||
kwargs = dict()
|
kwargs = dict()
|
||||||
|
|
||||||
|
|
@ -278,7 +283,7 @@ class BumperProxyModeMQTTClient(MQTTClient):
|
||||||
reader = StreamReaderAdapter(conn_reader)
|
reader = StreamReaderAdapter(conn_reader)
|
||||||
writer = StreamWriterAdapter(conn_writer)
|
writer = StreamWriterAdapter(conn_writer)
|
||||||
elif scheme in ('ws', 'wss'):
|
elif scheme in ('ws', 'wss'):
|
||||||
websocket = await websockets.connect(
|
websocket = await websockets.connect(
|
||||||
self.session.broker_uri,
|
self.session.broker_uri,
|
||||||
subprotocols=['mqtt'],
|
subprotocols=['mqtt'],
|
||||||
loop=self._loop,
|
loop=self._loop,
|
||||||
|
|
@ -322,21 +327,25 @@ class BumperProxyModeMQTTClient(MQTTClient):
|
||||||
msgdata = str(message.data.decode("utf-8"))
|
msgdata = str(message.data.decode("utf-8"))
|
||||||
|
|
||||||
proxymodelog.info(f"MQTT Proxy Client - Message Received From Ecovacs - Topic: {message.topic} - Message: {msgdata}")
|
proxymodelog.info(f"MQTT Proxy Client - Message Received From Ecovacs - Topic: {message.topic} - Message: {msgdata}")
|
||||||
ttopic = message.topic.split("/")
|
topic = message.topic
|
||||||
self.ecohelpername = ttopic[3]
|
ttopic = topic.split("/")
|
||||||
ttopic[3] = "proxyhelper"
|
if ttopic[1] == "p2p":
|
||||||
ttopic_comb = "/".join(ttopic)
|
self.eco_helper_names[ttopic[10]] = ttopic[3]
|
||||||
proxymodelog.info(f"MQTT Proxy Client - Converted Topic From {message.topic} TO {ttopic_comb}")
|
ttopic[3] = "proxyhelper"
|
||||||
proxymodelog.info(f"MQTT Proxy Client - Proxy Forward Message to Helperbot - Topic: {ttopic_comb} - Message: {msgdata.encode()}")
|
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(
|
await bumper.mqtt_helperbot.Client.publish(
|
||||||
ttopic_comb, msgdata.encode(), QOS_0
|
topic, msgdata.encode(), QOS_0
|
||||||
)
|
)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
proxymodelog.error(f"MQTT Proxy Client - get_msg Exception - {e}")
|
proxymodelog.error(f"MQTT Proxy Client - get_msg Exception - {e}")
|
||||||
|
|
||||||
class BumperMQTTServer_Plugin:
|
class BumperMQTTServer_Plugin:
|
||||||
proxyclients = {}
|
proxyclients: Dict[str, BumperProxyModeMQTTClient] = {}
|
||||||
def __init__(self, context):
|
def __init__(self, context):
|
||||||
self.context = context
|
self.context = context
|
||||||
try:
|
try:
|
||||||
|
|
@ -395,7 +404,7 @@ class BumperMQTTServer_Plugin:
|
||||||
|
|
||||||
try:
|
try:
|
||||||
await self.proxyclients[client_id].connect(
|
await self.proxyclients[client_id].connect(
|
||||||
f"mqtts://{username}:{password}@{mqtt_server}:8883",
|
f"mqtts://{username}:{password}@{mqtt_server}:443",
|
||||||
)
|
)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
mqttserverlog.error(f"MQTT Proxy Mode - Exception connecting with proxy to ecovacs - {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 not str(message.topic).split("/")[3] == "proxyhelper": # if from proxyhelper, don't send back to ecovacs...yet
|
||||||
if str(message.topic).split("/")[6] == "proxyhelper":
|
if str(message.topic).split("/")[6] == "proxyhelper":
|
||||||
ttopic = message.topic.split("/")
|
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)
|
ttopic_join = "/".join(ttopic)
|
||||||
proxymodelog.info(f"MQTT Proxy Client - Bot Message Converted Topic From {message.topic} TO {ttopic_join} with message: {msgdata}")
|
proxymodelog.info(f"MQTT Proxy Client - Bot Message Converted Topic From {message.topic} TO {ttopic_join} with message: {msgdata}")
|
||||||
else:
|
else:
|
||||||
|
|
|
||||||
|
|
@ -29,7 +29,7 @@ class portal_api_iot(plugins.ConfServerApp):
|
||||||
try:
|
try:
|
||||||
json_body = json.loads(await request.text())
|
json_body = json.loads(await request.text())
|
||||||
|
|
||||||
randomid = "".join(random.sample(string.ascii_letters, 6))
|
randomid = "".join(random.sample(string.ascii_letters, 4))
|
||||||
did = ""
|
did = ""
|
||||||
if "toId" in json_body: # Its a command
|
if "toId" in json_body: # Its a command
|
||||||
did = json_body["toId"]
|
did = json_body["toId"]
|
||||||
|
|
|
||||||
|
|
@ -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_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_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_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 |
|
||||||
Loading…
Add table
Add a link
Reference in a new issue