diff --git a/.gitignore b/.gitignore index 9644ac8..31682cc 100644 --- a/.gitignore +++ b/.gitignore @@ -28,3 +28,5 @@ certs/* !logs/README.md !data/README.md !certs/README.md + +venv/ \ No newline at end of file diff --git a/bumper/__init__.py b/bumper/__init__.py index bc1f821..aea1216 100644 --- a/bumper/__init__.py +++ b/bumper/__init__.py @@ -1,21 +1,18 @@ #!/usr/bin/env python3 +import asyncio +import importlib +import pkgutil +import socket +import sys +from logging.handlers import RotatingFileHandler from typing import Optional from bumper.confserver import ConfServer +from bumper.db import * +from bumper.models import * from bumper.mqttserver import MQTTServer, MQTTHelperBot 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): 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") os.makedirs(certs_dir, exist_ok=True) # Ensure data directory exists or create - - # Certs 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") @@ -51,7 +46,6 @@ bumper_listen = os.environ.get("BUMPER_LISTEN") or socket.gethostbyname( socket.gethostname() ) - bumper_announce_ip = os.environ.get("BUMPER_ANNOUNCE_IP") or bumper_listen # Other @@ -89,7 +83,7 @@ if not log_to_stdout: bumper_rotate = RotatingFileHandler("logs/bumper.log", maxBytes=5000000, backupCount=5) bumper_rotate.setFormatter(logformat) bumperlog.addHandler(bumper_rotate) -else: +else: bumperlog.addHandler(logging.StreamHandler(sys.stdout)) # Override the logging level # bumperlog.setLevel(logging.INFO) @@ -127,23 +121,23 @@ else: 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( +# brokerlog.setLevel( # 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: brokerlog.addHandler(mqtt_rotate) else: brokerlog.addHandler(logging.StreamHandler(sys.stdout)) protolog = logging.getLogger("hbmqtt.mqtt.protocol") -#protolog.setLevel( +# protolog.setLevel( # logging.CRITICAL + 1 -#) # Ignore this logger +# ) # Ignore this logger 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.setLevel(logging.CRITICAL + 1) # Ignore this logger if not log_to_stdout: clientlog.addHandler(mqtt_rotate) else: @@ -186,7 +180,6 @@ else: logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger - mqtt_listen_port = 8883 conf1_listen_port = 443 conf2_listen_port = 8007 @@ -194,7 +187,6 @@ xmpp_listen_port = 5223 async def start(): - try: loop = asyncio.get_event_loop() except: @@ -218,9 +210,9 @@ async def start(): return if not ( - os.path.exists(ca_cert) - and os.path.exists(server_cert) - and os.path.exists(server_key) + os.path.exists(ca_cert) + and os.path.exists(server_cert) + and os.path.exists(server_key) ): logging.log(logging.FATAL, "Certificate(s) don't exist at paths specified") return @@ -256,8 +248,10 @@ async def start(): # Start web servers 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)) + 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)) # Start maintenance while not shutting_down: @@ -363,15 +357,15 @@ def main(argv=None): try: if not ( - os.path.exists(ca_cert) - and os.path.exists(server_cert) - and os.path.exists(server_key) + os.path.exists(ca_cert) + and os.path.exists(server_cert) + and os.path.exists(server_key) ): first_run() return 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 @@ -410,4 +404,3 @@ def main(argv=None): finally: asyncio.run(shutdown()) - diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py index b2fa69b..c5711bd 100644 --- a/bumper/mqttserver.py +++ b/bumper/mqttserver.py @@ -1,26 +1,27 @@ #!/usr/bin/env python3 -import logging import asyncio +import json +import logging import os +import time +from datetime import datetime, timedelta + import hbmqtt +import pkg_resources from hbmqtt.broker import Broker from hbmqtt.client import MQTTClient -from hbmqtt.mqtt.constants import QOS_0, QOS_1, QOS_2 -import pkg_resources -import time -import bumper -import json -from datetime import datetime, timedelta -import bumper +from hbmqtt.mqtt.constants import QOS_0 from passlib.apps import custom_app_context as pwd_context +import bumper + helperbotlog = logging.getLogger("helperbot") boterrorlog = logging.getLogger("boterror") mqttserverlog = logging.getLogger("mqttserver") -class MQTTHelperBot: +class MQTTHelperBot: Client = None wait_resp_timeout_seconds = 60 @@ -49,16 +50,16 @@ class MQTTHelperBot: ] ) -# except ConnectionRefusedError as e: -# helperbotlog.Error(e) -# pass + # except ConnectionRefusedError as e: + # helperbotlog.Error(e) + # pass -# except asyncio.CancelledError as e: -# pass + # except asyncio.CancelledError as e: + # pass -# except hbmqtt.client.ConnectException as e: -# helperbotlog.Error(e) -# pass + # except hbmqtt.client.ConnectException as e: + # helperbotlog.Error(e) + # pass except Exception as e: helperbotlog.exception("{}".format(e)) @@ -67,7 +68,7 @@ class MQTTHelperBot: try: t_end = ( - datetime.now() + timedelta(seconds=self.wait_resp_timeout_seconds) + datetime.now() + timedelta(seconds=self.wait_resp_timeout_seconds) ).timestamp() while time.time() < t_end: @@ -156,12 +157,12 @@ class MQTTServer: except hbmqtt.broker.BrokerException as e: mqttserverlog.exception(e) - #asyncio.create_task(bumper.shutdown()) + # asyncio.create_task(bumper.shutdown()) pass except Exception as e: mqttserverlog.exception("{}".format(e)) - #asyncio.create_task(bumper.shutdown()) + # asyncio.create_task(bumper.shutdown()) pass def __init__(self, address, **kwargs): @@ -171,7 +172,7 @@ class MQTTServer: # Default config opts passwd_file = os.path.join( 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 @@ -180,7 +181,7 @@ class MQTTServer: passwd_file = kwargs["password_file"] 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 distribution = pkg_resources.Distribution("hbmqtt.broker.plugins") @@ -243,7 +244,7 @@ class BumperMQTTServer_Plugin: if "@" in client_id: didsplit = str(client_id).split("@") 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("/") bumper.bot_add( @@ -253,7 +254,8 @@ class BumperMQTTServer_Plugin: tmpbotdetail[1], "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 else: @@ -274,23 +276,26 @@ class BumperMQTTServer_Plugin: if auth: 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 else: authenticated = False # 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) - 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) 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: mqttserverlog.info(f"File Authentication Failed - Username: {username} - ClientID: {client_id}") 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: mqttserverlog.exception( @@ -302,9 +307,10 @@ class BumperMQTTServer_Plugin: allow_anonymous = self.auth_config.get( "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 - 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}") return authenticated @@ -317,7 +323,7 @@ class BumperMQTTServer_Plugin: self.context.logger.debug(f"Reading user database from {password_file}") for l in f: 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) if username: self._users[username] = pwd_hash @@ -346,63 +352,62 @@ class BumperMQTTServer_Plugin: def handle_helperbot_msg(self, client_id, message): - if str(message.topic).split("/")[6] == "helperbot": - # Response to command - helperbotlog.debug( - "Received Response - Topic: {} - Message: {}".format( + if str(message.topic).split("/")[6] == "helperbot": + # Response to command + helperbotlog.debug( + "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")) ) ) - 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: helperbotlog.debug( - "Received Message - Topic: {} - Message: {}".format( + "Received Broadcast - 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 = ( + else: + helperbotlog.debug( + "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"]) + timedelta(seconds=MQTTHelperBot.wait_resp_timeout_seconds) - ).timestamp() - if time.time() > expire_time: - helperbotlog.debug( - "Pruning Message Due To Expiration - Message Topic: {}".format( - msg["topic"] - ) + ).timestamp() + if time.time() > expire_time: + helperbotlog.debug( + "Pruning Message Due To Expiration - Message Topic: {}".format( + 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):