format files
This commit is contained in:
parent
106fbdcc7d
commit
a712b64798
3 changed files with 109 additions and 109 deletions
2
.gitignore
vendored
2
.gitignore
vendored
|
|
@ -28,3 +28,5 @@ certs/*
|
||||||
!logs/README.md
|
!logs/README.md
|
||||||
!data/README.md
|
!data/README.md
|
||||||
!certs/README.md
|
!certs/README.md
|
||||||
|
|
||||||
|
venv/
|
||||||
|
|
@ -1,21 +1,18 @@
|
||||||
#!/usr/bin/env python3
|
#!/usr/bin/env python3
|
||||||
|
import asyncio
|
||||||
|
import importlib
|
||||||
|
import pkgutil
|
||||||
|
import socket
|
||||||
|
import sys
|
||||||
|
from logging.handlers import RotatingFileHandler
|
||||||
from typing import Optional
|
from typing import Optional
|
||||||
|
|
||||||
from bumper.confserver import ConfServer
|
from bumper.confserver import ConfServer
|
||||||
|
from bumper.db import *
|
||||||
|
from bumper.models import *
|
||||||
from bumper.mqttserver import MQTTServer, MQTTHelperBot
|
from bumper.mqttserver import MQTTServer, MQTTHelperBot
|
||||||
from bumper.xmppserver import XMPPServer
|
from bumper.xmppserver import XMPPServer
|
||||||
from bumper.models import *
|
|
||||||
from bumper.db import *
|
|
||||||
import asyncio
|
|
||||||
import os
|
|
||||||
import logging
|
|
||||||
from logging.handlers import RotatingFileHandler
|
|
||||||
import socket
|
|
||||||
import sys
|
|
||||||
|
|
||||||
import importlib
|
|
||||||
import pkgutil
|
|
||||||
from pkgutil import extend_path
|
|
||||||
|
|
||||||
def strtobool(strbool):
|
def strtobool(strbool):
|
||||||
if str(strbool).lower() in ["true", "1", "t", "y", "on", "yes"]:
|
if str(strbool).lower() in ["true", "1", "t", "y", "on", "yes"]:
|
||||||
|
|
@ -39,8 +36,6 @@ 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")
|
||||||
os.makedirs(certs_dir, exist_ok=True) # Ensure data directory exists or create
|
os.makedirs(certs_dir, exist_ok=True) # Ensure data directory exists or create
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
# Certs
|
# Certs
|
||||||
ca_cert = os.environ.get("BUMPER_CA") or os.path.join(certs_dir, "ca.crt")
|
ca_cert = os.environ.get("BUMPER_CA") or os.path.join(certs_dir, "ca.crt")
|
||||||
server_cert = os.environ.get("BUMPER_CERT") or os.path.join(certs_dir, "bumper.crt")
|
server_cert = os.environ.get("BUMPER_CERT") or os.path.join(certs_dir, "bumper.crt")
|
||||||
|
|
@ -51,7 +46,6 @@ bumper_listen = os.environ.get("BUMPER_LISTEN") or socket.gethostbyname(
|
||||||
socket.gethostname()
|
socket.gethostname()
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
bumper_announce_ip = os.environ.get("BUMPER_ANNOUNCE_IP") or bumper_listen
|
bumper_announce_ip = os.environ.get("BUMPER_ANNOUNCE_IP") or bumper_listen
|
||||||
|
|
||||||
# Other
|
# Other
|
||||||
|
|
@ -89,7 +83,7 @@ if not log_to_stdout:
|
||||||
bumper_rotate = RotatingFileHandler("logs/bumper.log", maxBytes=5000000, backupCount=5)
|
bumper_rotate = RotatingFileHandler("logs/bumper.log", maxBytes=5000000, backupCount=5)
|
||||||
bumper_rotate.setFormatter(logformat)
|
bumper_rotate.setFormatter(logformat)
|
||||||
bumperlog.addHandler(bumper_rotate)
|
bumperlog.addHandler(bumper_rotate)
|
||||||
else:
|
else:
|
||||||
bumperlog.addHandler(logging.StreamHandler(sys.stdout))
|
bumperlog.addHandler(logging.StreamHandler(sys.stdout))
|
||||||
# Override the logging level
|
# Override the logging level
|
||||||
# bumperlog.setLevel(logging.INFO)
|
# bumperlog.setLevel(logging.INFO)
|
||||||
|
|
@ -127,23 +121,23 @@ else:
|
||||||
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)
|
||||||
if not log_to_stdout:
|
if not log_to_stdout:
|
||||||
brokerlog.addHandler(mqtt_rotate)
|
brokerlog.addHandler(mqtt_rotate)
|
||||||
else:
|
else:
|
||||||
brokerlog.addHandler(logging.StreamHandler(sys.stdout))
|
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
|
||||||
if not log_to_stdout:
|
if not log_to_stdout:
|
||||||
protolog.addHandler(mqtt_rotate)
|
protolog.addHandler(mqtt_rotate)
|
||||||
else:
|
else:
|
||||||
protolog.addHandler(logging.StreamHandler(sys.stdout))
|
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
|
||||||
if not log_to_stdout:
|
if not log_to_stdout:
|
||||||
clientlog.addHandler(mqtt_rotate)
|
clientlog.addHandler(mqtt_rotate)
|
||||||
else:
|
else:
|
||||||
|
|
@ -186,7 +180,6 @@ else:
|
||||||
|
|
||||||
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
|
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
|
||||||
|
|
||||||
|
|
||||||
mqtt_listen_port = 8883
|
mqtt_listen_port = 8883
|
||||||
conf1_listen_port = 443
|
conf1_listen_port = 443
|
||||||
conf2_listen_port = 8007
|
conf2_listen_port = 8007
|
||||||
|
|
@ -194,7 +187,6 @@ xmpp_listen_port = 5223
|
||||||
|
|
||||||
|
|
||||||
async def start():
|
async def start():
|
||||||
|
|
||||||
try:
|
try:
|
||||||
loop = asyncio.get_event_loop()
|
loop = asyncio.get_event_loop()
|
||||||
except:
|
except:
|
||||||
|
|
@ -218,9 +210,9 @@ async def start():
|
||||||
return
|
return
|
||||||
|
|
||||||
if not (
|
if not (
|
||||||
os.path.exists(ca_cert)
|
os.path.exists(ca_cert)
|
||||||
and os.path.exists(server_cert)
|
and os.path.exists(server_cert)
|
||||||
and os.path.exists(server_key)
|
and os.path.exists(server_key)
|
||||||
):
|
):
|
||||||
logging.log(logging.FATAL, "Certificate(s) don't exist at paths specified")
|
logging.log(logging.FATAL, "Certificate(s) don't exist at paths specified")
|
||||||
return
|
return
|
||||||
|
|
@ -256,8 +248,10 @@ async def start():
|
||||||
|
|
||||||
# Start web servers
|
# Start web servers
|
||||||
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(
|
||||||
asyncio.create_task(conf_server.start_site(conf_server.app, address=bumper_listen, port=conf2_listen_port, usessl=False))
|
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))
|
||||||
|
|
||||||
# Start maintenance
|
# Start maintenance
|
||||||
while not shutting_down:
|
while not shutting_down:
|
||||||
|
|
@ -363,15 +357,15 @@ def main(argv=None):
|
||||||
try:
|
try:
|
||||||
|
|
||||||
if not (
|
if not (
|
||||||
os.path.exists(ca_cert)
|
os.path.exists(ca_cert)
|
||||||
and os.path.exists(server_cert)
|
and os.path.exists(server_cert)
|
||||||
and os.path.exists(server_key)
|
and os.path.exists(server_key)
|
||||||
):
|
):
|
||||||
first_run()
|
first_run()
|
||||||
return
|
return
|
||||||
|
|
||||||
if not (
|
if not (
|
||||||
os.path.exists(os.path.join(data_dir, "passwd"))
|
os.path.exists(os.path.join(data_dir, "passwd"))
|
||||||
):
|
):
|
||||||
with open(os.path.join(data_dir, "passwd"), 'w'): pass
|
with open(os.path.join(data_dir, "passwd"), 'w'): pass
|
||||||
|
|
||||||
|
|
@ -410,4 +404,3 @@ def main(argv=None):
|
||||||
|
|
||||||
finally:
|
finally:
|
||||||
asyncio.run(shutdown())
|
asyncio.run(shutdown())
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -1,26 +1,27 @@
|
||||||
#!/usr/bin/env python3
|
#!/usr/bin/env python3
|
||||||
|
|
||||||
import logging
|
|
||||||
import asyncio
|
import asyncio
|
||||||
|
import json
|
||||||
|
import logging
|
||||||
import os
|
import os
|
||||||
|
import time
|
||||||
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
import hbmqtt
|
import hbmqtt
|
||||||
|
import pkg_resources
|
||||||
from hbmqtt.broker import Broker
|
from hbmqtt.broker import Broker
|
||||||
from hbmqtt.client import MQTTClient
|
from hbmqtt.client import MQTTClient
|
||||||
from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2
|
from hbmqtt.mqtt.constants import QOS_0
|
||||||
import pkg_resources
|
|
||||||
import time
|
|
||||||
import bumper
|
|
||||||
import json
|
|
||||||
from datetime import datetime, timedelta
|
|
||||||
import bumper
|
|
||||||
from passlib.apps import custom_app_context as pwd_context
|
from passlib.apps import custom_app_context as pwd_context
|
||||||
|
|
||||||
|
import bumper
|
||||||
|
|
||||||
helperbotlog = logging.getLogger("helperbot")
|
helperbotlog = logging.getLogger("helperbot")
|
||||||
boterrorlog = logging.getLogger("boterror")
|
boterrorlog = logging.getLogger("boterror")
|
||||||
mqttserverlog = logging.getLogger("mqttserver")
|
mqttserverlog = logging.getLogger("mqttserver")
|
||||||
|
|
||||||
class MQTTHelperBot:
|
|
||||||
|
|
||||||
|
class MQTTHelperBot:
|
||||||
Client = None
|
Client = None
|
||||||
wait_resp_timeout_seconds = 60
|
wait_resp_timeout_seconds = 60
|
||||||
|
|
||||||
|
|
@ -49,16 +50,16 @@ class MQTTHelperBot:
|
||||||
]
|
]
|
||||||
)
|
)
|
||||||
|
|
||||||
# except ConnectionRefusedError as e:
|
# except ConnectionRefusedError as e:
|
||||||
# helperbotlog.Error(e)
|
# helperbotlog.Error(e)
|
||||||
# pass
|
# pass
|
||||||
|
|
||||||
# except asyncio.CancelledError as e:
|
# except asyncio.CancelledError as e:
|
||||||
# pass
|
# pass
|
||||||
|
|
||||||
# except hbmqtt.client.ConnectException as e:
|
# except hbmqtt.client.ConnectException as e:
|
||||||
# helperbotlog.Error(e)
|
# helperbotlog.Error(e)
|
||||||
# pass
|
# pass
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
helperbotlog.exception("{}".format(e))
|
helperbotlog.exception("{}".format(e))
|
||||||
|
|
@ -67,7 +68,7 @@ class MQTTHelperBot:
|
||||||
try:
|
try:
|
||||||
|
|
||||||
t_end = (
|
t_end = (
|
||||||
datetime.now() + timedelta(seconds=self.wait_resp_timeout_seconds)
|
datetime.now() + timedelta(seconds=self.wait_resp_timeout_seconds)
|
||||||
).timestamp()
|
).timestamp()
|
||||||
|
|
||||||
while time.time() < t_end:
|
while time.time() < t_end:
|
||||||
|
|
@ -156,12 +157,12 @@ class MQTTServer:
|
||||||
|
|
||||||
except hbmqtt.broker.BrokerException as e:
|
except hbmqtt.broker.BrokerException as e:
|
||||||
mqttserverlog.exception(e)
|
mqttserverlog.exception(e)
|
||||||
#asyncio.create_task(bumper.shutdown())
|
# asyncio.create_task(bumper.shutdown())
|
||||||
pass
|
pass
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
mqttserverlog.exception("{}".format(e))
|
mqttserverlog.exception("{}".format(e))
|
||||||
#asyncio.create_task(bumper.shutdown())
|
# asyncio.create_task(bumper.shutdown())
|
||||||
pass
|
pass
|
||||||
|
|
||||||
def __init__(self, address, **kwargs):
|
def __init__(self, address, **kwargs):
|
||||||
|
|
@ -171,7 +172,7 @@ class MQTTServer:
|
||||||
# Default config opts
|
# Default config opts
|
||||||
passwd_file = os.path.join(
|
passwd_file = os.path.join(
|
||||||
os.path.join(bumper.data_dir, "passwd")
|
os.path.join(bumper.data_dir, "passwd")
|
||||||
) # For file auth, set user:hash in passwd file see (https://hbmqtt.readthedocs.io/en/latest/references/hbmqtt.html#configuration-example)
|
) # For file auth, set user:hash in passwd file see (https://hbmqtt.readthedocs.io/en/latest/references/hbmqtt.html#configuration-example)
|
||||||
|
|
||||||
allow_anon = False
|
allow_anon = False
|
||||||
|
|
||||||
|
|
@ -180,7 +181,7 @@ class MQTTServer:
|
||||||
passwd_file = kwargs["password_file"]
|
passwd_file = kwargs["password_file"]
|
||||||
|
|
||||||
elif key == "allow_anonymous":
|
elif key == "allow_anonymous":
|
||||||
allow_anon = kwargs["allow_anonymous"] # Set to True to allow anonymous authentication
|
allow_anon = kwargs["allow_anonymous"] # Set to True to allow anonymous authentication
|
||||||
|
|
||||||
# The below adds a plugin to the hbmqtt.broker.plugins without having to futz with setup.py
|
# The below adds a plugin to the hbmqtt.broker.plugins without having to futz with setup.py
|
||||||
distribution = pkg_resources.Distribution("hbmqtt.broker.plugins")
|
distribution = pkg_resources.Distribution("hbmqtt.broker.plugins")
|
||||||
|
|
@ -243,7 +244,7 @@ class BumperMQTTServer_Plugin:
|
||||||
if "@" in client_id:
|
if "@" in client_id:
|
||||||
didsplit = str(client_id).split("@")
|
didsplit = str(client_id).split("@")
|
||||||
if not ( # if ecouser or bumper aren't in details it is a bot
|
if not ( # if ecouser or bumper aren't in details it is a bot
|
||||||
"ecouser" in didsplit[1] or "bumper" in didsplit[1]
|
"ecouser" in didsplit[1] or "bumper" in didsplit[1]
|
||||||
):
|
):
|
||||||
tmpbotdetail = str(didsplit[1]).split("/")
|
tmpbotdetail = str(didsplit[1]).split("/")
|
||||||
bumper.bot_add(
|
bumper.bot_add(
|
||||||
|
|
@ -253,7 +254,8 @@ class BumperMQTTServer_Plugin:
|
||||||
tmpbotdetail[1],
|
tmpbotdetail[1],
|
||||||
"eco-ng",
|
"eco-ng",
|
||||||
)
|
)
|
||||||
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
|
||||||
|
|
||||||
else:
|
else:
|
||||||
|
|
@ -274,23 +276,26 @@ class BumperMQTTServer_Plugin:
|
||||||
|
|
||||||
if auth:
|
if auth:
|
||||||
bumper.client_add(userid, realm, resource)
|
bumper.client_add(userid, realm, resource)
|
||||||
mqttserverlog.info(f"Bumper Authentication Success - Client - Username: {username} - ClientID: {client_id}")
|
mqttserverlog.info(
|
||||||
|
f"Bumper Authentication Success - Client - Username: {username} - ClientID: {client_id}")
|
||||||
authenticated = True
|
authenticated = True
|
||||||
|
|
||||||
else:
|
else:
|
||||||
authenticated = False
|
authenticated = False
|
||||||
|
|
||||||
# Check for File Auth
|
# Check for File Auth
|
||||||
if username and not authenticated: # If there is a username and it isn't already authenticated
|
if username and not authenticated: # If there is a username and it isn't already authenticated
|
||||||
hash = self._users.get(username, None)
|
hash = self._users.get(username, None)
|
||||||
if hash: # If there is a matching entry in passwd, check hash
|
if hash: # If there is a matching entry in passwd, check hash
|
||||||
authenticated = pwd_context.verify(password, hash)
|
authenticated = pwd_context.verify(password, hash)
|
||||||
if authenticated:
|
if authenticated:
|
||||||
mqttserverlog.info(f"File Authentication Success - Username: {username} - ClientID: {client_id}")
|
mqttserverlog.info(
|
||||||
|
f"File Authentication Success - Username: {username} - ClientID: {client_id}")
|
||||||
else:
|
else:
|
||||||
mqttserverlog.info(f"File Authentication Failed - Username: {username} - ClientID: {client_id}")
|
mqttserverlog.info(f"File Authentication Failed - Username: {username} - ClientID: {client_id}")
|
||||||
else:
|
else:
|
||||||
mqttserverlog.info(f"File Authentication Failed - No Entry for Username: {username} - ClientID: {client_id}")
|
mqttserverlog.info(
|
||||||
|
f"File Authentication Failed - No Entry for Username: {username} - ClientID: {client_id}")
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
mqttserverlog.exception(
|
mqttserverlog.exception(
|
||||||
|
|
@ -302,9 +307,10 @@ class BumperMQTTServer_Plugin:
|
||||||
allow_anonymous = self.auth_config.get(
|
allow_anonymous = self.auth_config.get(
|
||||||
"allow-anonymous", True
|
"allow-anonymous", True
|
||||||
)
|
)
|
||||||
if allow_anonymous and not authenticated: # If anonymous auth is allowed and it isn't already authenticated
|
if allow_anonymous and not authenticated: # If anonymous auth is allowed and it isn't already authenticated
|
||||||
authenticated = True
|
authenticated = True
|
||||||
self.context.logger.debug(f"Anonymous Authentication Success: config allows anonymous - Username: {username}")
|
self.context.logger.debug(
|
||||||
|
f"Anonymous Authentication Success: config allows anonymous - Username: {username}")
|
||||||
mqttserverlog.info(f"Anonymous Authentication Success: config allows anonymous - Username: {username}")
|
mqttserverlog.info(f"Anonymous Authentication Success: config allows anonymous - Username: {username}")
|
||||||
|
|
||||||
return authenticated
|
return authenticated
|
||||||
|
|
@ -317,7 +323,7 @@ class BumperMQTTServer_Plugin:
|
||||||
self.context.logger.debug(f"Reading user database from {password_file}")
|
self.context.logger.debug(f"Reading user database from {password_file}")
|
||||||
for l in f:
|
for l in f:
|
||||||
line = l.strip()
|
line = l.strip()
|
||||||
if not line.startswith('#'): # Allow comments in files
|
if not line.startswith('#'): # Allow comments in files
|
||||||
(username, pwd_hash) = line.split(sep=":", maxsplit=3)
|
(username, pwd_hash) = line.split(sep=":", maxsplit=3)
|
||||||
if username:
|
if username:
|
||||||
self._users[username] = pwd_hash
|
self._users[username] = pwd_hash
|
||||||
|
|
@ -346,63 +352,62 @@ class BumperMQTTServer_Plugin:
|
||||||
|
|
||||||
def handle_helperbot_msg(self, client_id, message):
|
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
|
||||||
helperbotlog.debug(
|
helperbotlog.debug(
|
||||||
"Received Response - Topic: {} - Message: {}".format(
|
"Received Response - Topic: {} - Message: {}".format(
|
||||||
|
message.topic, str(message.data.decode("utf-8"))
|
||||||
|
)
|
||||||
|
)
|
||||||
|
bumper.mqtt_helperbot.command_responses.append(
|
||||||
|
{
|
||||||
|
"time": time.time(),
|
||||||
|
"topic": message.topic,
|
||||||
|
"payload": str(message.data.decode("utf-8")),
|
||||||
|
}
|
||||||
|
)
|
||||||
|
elif str(message.topic).split("/")[3] == "helperbot":
|
||||||
|
# Helperbot sending command
|
||||||
|
helperbotlog.debug(
|
||||||
|
"Send Command - Topic: {} - Message: {}".format(
|
||||||
|
message.topic, str(message.data.decode("utf-8"))
|
||||||
|
)
|
||||||
|
)
|
||||||
|
elif str(message.topic).split("/")[1] == "atr":
|
||||||
|
# Broadcast message received on atr
|
||||||
|
if str(message.topic).split("/")[2] == "errors":
|
||||||
|
boterrorlog.error(
|
||||||
|
"Received Error - Topic: {} - Message: {}".format(
|
||||||
message.topic, str(message.data.decode("utf-8"))
|
message.topic, str(message.data.decode("utf-8"))
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
bumper.mqtt_helperbot.command_responses.append(
|
|
||||||
{
|
|
||||||
"time": time.time(),
|
|
||||||
"topic": message.topic,
|
|
||||||
"payload": str(message.data.decode("utf-8")),
|
|
||||||
}
|
|
||||||
)
|
|
||||||
elif str(message.topic).split("/")[3] == "helperbot":
|
|
||||||
# Helperbot sending command
|
|
||||||
helperbotlog.debug(
|
|
||||||
"Send Command - Topic: {} - Message: {}".format(
|
|
||||||
message.topic, str(message.data.decode("utf-8"))
|
|
||||||
)
|
|
||||||
)
|
|
||||||
elif str(message.topic).split("/")[1] == "atr":
|
|
||||||
# Broadcast message received on atr
|
|
||||||
if str(message.topic).split("/")[2] == "errors":
|
|
||||||
boterrorlog.error(
|
|
||||||
"Received Error - Topic: {} - Message: {}".format(
|
|
||||||
message.topic, str(message.data.decode("utf-8"))
|
|
||||||
)
|
|
||||||
)
|
|
||||||
else:
|
|
||||||
helperbotlog.debug(
|
|
||||||
"Received Broadcast - Topic: {} - Message: {}".format(
|
|
||||||
message.topic, str(message.data.decode("utf-8"))
|
|
||||||
)
|
|
||||||
)
|
|
||||||
|
|
||||||
else:
|
else:
|
||||||
helperbotlog.debug(
|
helperbotlog.debug(
|
||||||
"Received Message - Topic: {} - Message: {}".format(
|
"Received Broadcast - Topic: {} - Message: {}".format(
|
||||||
message.topic, str(message.data.decode("utf-8"))
|
message.topic, str(message.data.decode("utf-8"))
|
||||||
)
|
)
|
||||||
)
|
)
|
||||||
|
|
||||||
# Cleanup "expired messages" > 60 seconds from time
|
else:
|
||||||
for msg in bumper.mqtt_helperbot.command_responses:
|
helperbotlog.debug(
|
||||||
expire_time = (
|
"Received Message - Topic: {} - Message: {}".format(
|
||||||
|
message.topic, str(message.data.decode("utf-8"))
|
||||||
|
)
|
||||||
|
)
|
||||||
|
|
||||||
|
# Cleanup "expired messages" > 60 seconds from time
|
||||||
|
for msg in bumper.mqtt_helperbot.command_responses:
|
||||||
|
expire_time = (
|
||||||
datetime.fromtimestamp(msg["time"])
|
datetime.fromtimestamp(msg["time"])
|
||||||
+ timedelta(seconds=MQTTHelperBot.wait_resp_timeout_seconds)
|
+ timedelta(seconds=MQTTHelperBot.wait_resp_timeout_seconds)
|
||||||
).timestamp()
|
).timestamp()
|
||||||
if time.time() > expire_time:
|
if time.time() > expire_time:
|
||||||
helperbotlog.debug(
|
helperbotlog.debug(
|
||||||
"Pruning Message Due To Expiration - Message Topic: {}".format(
|
"Pruning Message Due To Expiration - Message Topic: {}".format(
|
||||||
msg["topic"]
|
msg["topic"]
|
||||||
)
|
|
||||||
)
|
)
|
||||||
bumper.mqtt_helperbot.command_responses.remove(msg)
|
)
|
||||||
|
bumper.mqtt_helperbot.command_responses.remove(msg)
|
||||||
|
|
||||||
async def on_broker_client_disconnected(self, client_id):
|
async def on_broker_client_disconnected(self, client_id):
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue