diff --git a/bumper/__init__.py b/bumper/__init__.py
index 6e53947..e2ffaaa 100644
--- a/bumper/__init__.py
+++ b/bumper/__init__.py
@@ -1,9 +1,8 @@
#!/usr/bin/env python3
-from .confserver import ConfServer
-from .mqttserver import MQTTServer
-from .mqttserver import MQTTHelperBot
-from .xmppserver import XMPPServer
+from bumper.confserver import ConfServer
+from bumper.mqttserver import MQTTServer, MQTTHelperBot
+from bumper.xmppserver import XMPPServer
import asyncio
import json
import time
@@ -66,6 +65,9 @@ xmppserverlog.addHandler(xmpp_rotate)
# xmppserverlog.setLevel(logging.INFO)
+logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
+
+
def get_milli_time(timetoconvert):
return int(round(timetoconvert * 1000))
@@ -663,11 +665,12 @@ def bot_add(sn, did, devclass, resource, company):
newbot.company = company
bot = bot_get(did)
- if not bot:
- bumperlog.info(
- "Adding new bot with SN: {} DID: {}".format(newbot.name, newbot.did)
- )
- bot_full_upsert(newbot.asdict())
+ if not bot: # Not existing bot in database
+ if not devclass == "" or "@" not in sn or "tmp" not in sn: # try to prevent bad additions to the bot list
+ bumperlog.info(
+ "Adding new bot with SN: {} DID: {}".format(newbot.name, newbot.did)
+ )
+ bot_full_upsert(newbot.asdict())
def bot_remove(did):
diff --git a/bumper/confserver.py b/bumper/confserver.py
index 70f2474..1a51dc0 100644
--- a/bumper/confserver.py
+++ b/bumper/confserver.py
@@ -31,9 +31,7 @@ class aiohttp_filter(logging.Filter):
confserverlog = logging.getLogger("confserver")
-
-#logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
-logging.getLogger("aiohttp.access").addFilter(aiohttp_filter())
+logging.getLogger("aiohttp.access").addFilter(aiohttp_filter()) #Add logging filter above to aiohttp.access
class EcoVacs_Login:
accessToken = ""
@@ -60,42 +58,6 @@ class ConfServer:
self.run_async = False
self.app = None
- def run(self, run_async=False):
- try:
- if run_async:
- self.run_async = True
- confserverlog.debug("Starting ConfServer Thread: 1")
- self.confthread = Thread(
- name="ConfServer_{}_Thread".format(self.address[1]),
- target=self.run_server,
- )
- self.confthread.setDaemon(True)
- self.confthread.start()
-
- else:
- try:
- self.run_server()
- except KeyboardInterrupt:
- self.disconnect()
-
- except Exception as e:
- confserverlog.exception("{}".format(e))
-
- def run_server(self):
- logging.info("Starting ConfServer at {}".format(self.address))
- print("Starting ConfServer at {}".format(self.address))
- try:
- loop = asyncio.get_event_loop()
- except:
- loop = asyncio.new_event_loop()
-
- try:
- self.confserver_app()
- loop.run_until_complete(self.start_server())
- loop.run_forever()
- except Exception as e:
- confserverlog.exception("{}".format(e))
-
def confserver_app(self):
self.app = web.Application()
@@ -196,6 +158,7 @@ class ConfServer:
async def start_server(self):
try:
+ confserverlog.info("Starting ConfServer at {}:{}".format(self.address[0], self.address[1]))
runner = web.AppRunner(self.app)
await runner.setup()
@@ -351,55 +314,67 @@ class ConfServer:
def check_token(self, apptype, countrycode, user, token):
- if bumper.check_token(user["userid"], token):
-
- if "global_" in apptype: #EcoVacs Home
- login_details = EcoVacsHome_Login()
- login_details.ucUid = "fuid_{}".format(user["userid"])
- login_details.loginName = "fusername_{}".format(user["userid"])
- login_details.mobile = None
+ try:
+ if bumper.check_token(user["userid"], token):
+
+ if "global_" in apptype: #EcoVacs Home
+ login_details = EcoVacsHome_Login()
+ login_details.ucUid = "fuid_{}".format(user["userid"])
+ login_details.loginName = "fusername_{}".format(user["userid"])
+ login_details.mobile = None
+ else:
+ login_details = EcoVacs_Login()
+
+ login_details.accessToken = token
+ login_details.uid = "fuid_{}".format(user["userid"])
+ login_details.username = "fusername_{}".format(user["userid"])
+ login_details.country = countrycode
+ login_details.email = "null@null.com"
+
+ body = {
+ "code": bumper.RETURN_API_SUCCESS,
+ "data": json.loads(login_details.toJSON()),
+ #{
+ # "accessToken": self.generate_token(tmpuser), # Generate a token
+ # "country": countrycode,
+ # "email": "null@null.com",
+ # "uid": "fuid_{}".format(tmpuser["userid"]),
+ # "username": "fusername_{}".format(tmpuser["userid"]),
+ #},
+ "msg": "操作成功",
+ "time": bumper.get_milli_time(datetime.utcnow().timestamp()),
+ }
+ return web.json_response(body)
+
else:
- login_details = EcoVacs_Login()
-
- login_details.accessToken = token
- login_details.uid = "fuid_{}".format(user["userid"])
- login_details.username = "fusername_{}".format(user["userid"])
- login_details.country = countrycode
- login_details.email = "null@null.com"
-
- body = {
- "code": bumper.RETURN_API_SUCCESS,
- "data": json.loads(login_details.toJSON()),
- #{
- # "accessToken": self.generate_token(tmpuser), # Generate a token
- # "country": countrycode,
- # "email": "null@null.com",
- # "uid": "fuid_{}".format(tmpuser["userid"]),
- # "username": "fusername_{}".format(tmpuser["userid"]),
- #},
- "msg": "操作成功",
- "time": bumper.get_milli_time(datetime.utcnow().timestamp()),
- }
- return web.json_response(body)
-
- else:
- body = {
- "code": bumper.ERR_TOKEN_INVALID,
- "data": None,
- "msg": "当前密码错误",
- "time": bumper.get_milli_time(datetime.utcnow().timestamp()),
- }
- return web.json_response(body)
+ body = {
+ "code": bumper.ERR_TOKEN_INVALID,
+ "data": None,
+ "msg": "当前密码错误",
+ "time": bumper.get_milli_time(datetime.utcnow().timestamp()),
+ }
+ return web.json_response(body)
+
+ except Exception as e:
+ confserverlog.exception("{}".format(e))
def generate_token(self, user):
- tmpaccesstoken = uuid.uuid4().hex
- bumper.user_add_token(user["userid"], tmpaccesstoken)
- return tmpaccesstoken
+ try:
+ tmpaccesstoken = uuid.uuid4().hex
+ bumper.user_add_token(user["userid"], tmpaccesstoken)
+ return tmpaccesstoken
+
+ except Exception as e:
+ confserverlog.exception("{}".format(e))
def generate_authcode(self, user, countrycode, token):
- tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex)
- bumper.user_add_authcode(user["userid"], token, tmpauthcode)
- return tmpauthcode
+ try:
+ tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex)
+ bumper.user_add_authcode(user["userid"], token, tmpauthcode)
+ return tmpauthcode
+
+ except Exception as e:
+ confserverlog.exception("{}".format(e))
def _auth_any(self, devid, apptype, country, request):
try:
@@ -815,7 +790,6 @@ class ConfServer:
confserverlog.exception("{}".format(e))
async def handle_getProductIotMap(self, request):
- user_devid = request.match_info.get("devid", "")
try:
body = {
"code": bumper.RETURN_API_SUCCESS,
@@ -909,13 +883,23 @@ class ConfServer:
if todo == "FindBest":
service = postbody["service"]
if service == "EcoMsgNew":
+ srvip = socket.gethostbyname(socket.gethostname())
+ srvport = 5223
+ confserverlog.info(
+ "Reporting FindBest-EcoMsgNew Server to Bot as: {}:{}".format(srvip, srvport)
+ )
body = {
"result": "ok",
- "ip": socket.gethostbyname(socket.gethostname()),
- "port": 5223,
+ "ip": srvip,
+ "port": srvport,
}
elif service == "EcoUpdate":
- body = {"result": "ok", "ip": "47.88.66.164", "port": 8005}
+ srvip = "47.88.66.164" #EcoVacs Server
+ srvport = 8005
+ confserverlog.info(
+ "Reporting FindBest-EcoUpdate Server to Bot as: {}:{}".format(srvip, srvport)
+ )
+ body = {"result": "ok", "ip": srvip, "port": srvport}
elif todo == "loginByItToken":
if "userId" in postbody:
@@ -995,7 +979,8 @@ class ConfServer:
for bot in bots:
if bot["class"] != "":
b = bumper.bot_toEcoVacsHome_JSON(bot)
- botlist.append(json.loads(b))
+ if not b is None: #Happens if the bot isn't on the EcoVacs Home list
+ botlist.append(json.loads(b))
body = {
"code": 0,
@@ -1035,9 +1020,12 @@ class ConfServer:
if todo == "FindBest":
service = postbody["service"]
if service == "EcoMsgNew":
-
srvip = socket.gethostbyname(socket.gethostname())
- msgserver = {"ip": srvip, "port": 5223, "result": "ok"}
+ srvport = 5223
+ confserverlog.info(
+ "Reporting FindBest-EcoMsgNew Server to Bot as: {}:{}".format(srvip, srvport)
+ )
+ msgserver = {"ip": srvip, "port": srvport, "result": "ok"}
msgserver = json.dumps(msgserver)
msgserver = msgserver.replace(
" ", ""
@@ -1094,11 +1082,16 @@ class ConfServer:
if did != "":
- confserverlog.debug("BotCommand: {}".format(json_body))
bot = bumper.bot_get(did)
if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True:
body = ""
retcmd = await self.helperbot.send_command(json_body, randomid)
+ confserverlog.debug(
+ "Send Bot - {}".format(json_body)
+ )
+ confserverlog.debug(
+ "Bot Response - {}".format(body)
+ )
logs = []
logsroot = ET.fromstring(retcmd["resp"])
if logsroot.attrib["ret"] == "ok":
@@ -1147,13 +1140,15 @@ class ConfServer:
did = json_body["toId"]
if did != "":
- confserverlog.debug("BotCommand: {}".format(json_body))
bot = bumper.bot_get(did)
if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True:
retcmd = await self.helperbot.send_command(json_body, randomid)
body = retcmd
confserverlog.debug(
- "\r\n POST: {} \r\n Response: {}".format(json_body, body)
+ "Send Bot - {}".format(json_body)
+ )
+ confserverlog.debug(
+ "Bot Response - {}".format(body)
)
return web.json_response(body)
else:
@@ -1180,7 +1175,7 @@ class ConfServer:
confserverlog.exception("{}".format(e))
- async def handle_dim_devmanager(self, request):
+ async def handle_dim_devmanager(self, request): #Used in EcoVacs Home App
try:
json_body = json.loads(await request.text())
@@ -1190,13 +1185,15 @@ class ConfServer:
did = json_body["toId"]
if did != "":
- confserverlog.debug("BotCommand: {}".format(json_body))
bot = bumper.bot_get(did)
if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True:
retcmd = await self.helperbot.send_command(json_body, randomid)
body = retcmd
confserverlog.debug(
- "\r\n POST: {} \r\n Response: {}".format(json_body, body)
+ "Send Bot - {}".format(json_body)
+ )
+ confserverlog.debug(
+ "Bot Response - {}".format(body)
)
return web.json_response(body)
else:
diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py
index 0f145f6..a9c3ca1 100644
--- a/bumper/mqttserver.py
+++ b/bumper/mqttserver.py
@@ -14,6 +14,7 @@ import ssl
import bumper
import json
from datetime import datetime, timedelta
+import bumper
helperbotlog = logging.getLogger("helperbot")
mqttserverlog = logging.getLogger("mqttserver")
@@ -42,33 +43,6 @@ class MQTTHelperBot:
self.command_responses = []
self.helperthread = None
- def run(self, run_async=False):
- if run_async:
- hloop = asyncio.new_event_loop()
- helperbotlog.debug("Starting MQTT HelperBot Thread: 1")
- self.helperthread = Thread(
- name="MQTTHelperBot_Thread", target=self.run_helperbot, args=(hloop,)
- )
- self.helperthread.setDaemon(True)
- self.helperthread.start()
-
- else:
- self.run_helperbot(asyncio.get_event_loop())
-
- def run_helperbot(self, loop):
- logging.info("Starting MQTT HelperBot")
- print("Starting MQTT HelperBot")
- try:
- asyncio.set_event_loop(loop)
- self.Client = MQTTClient(
- client_id=self.client_id, config={"check_hostname": False}
- )
- loop.run_until_complete(self.start_helper_bot())
- loop.run_until_complete(self.get_msg())
- loop.run_forever()
- except Exception as e:
- helperbotlog.exception("{}".format(e))
-
async def start_helper_bot(self):
try:
@@ -87,8 +61,7 @@ class MQTTHelperBot:
("iot/atr/+", QOS_0),
]
)
-
- asyncio.ensure_future(self.get_msg())
+ asyncio.create_task(self.get_msg())
except Exception as e:
helperbotlog.exception("{}".format(e))
@@ -126,8 +99,6 @@ class MQTTHelperBot:
helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time))
self.command_responses.remove(msg)
- # helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
-
except Exception as e:
helperbotlog.exception("{}".format(e))
@@ -190,6 +161,7 @@ class MQTTServer:
async def broker_coro(self):
try:
+ mqttserverlog.info("Starting MQTT Server at {}:{}".format(self.address[0], self.address[1]))
broker = hbmqtt.broker.Broker(config=self.default_config)
await broker.start()
@@ -248,32 +220,6 @@ class MQTTServer:
except Exception as e:
mqttserverlog.exception("{}".format(e))
- def run(self, run_async=False):
- if run_async:
- sloop = asyncio.new_event_loop()
- mqttserverlog.debug("Starting MQTTServer Thread: 1")
- self.mqttserverthread = Thread(
- name="MQTTServer_Thread", target=self.run_server, args=(sloop,)
- )
- self.mqttserverthread.setDaemon(True)
- self.mqttserverthread.start()
-
- else:
- self.run_server(asyncio.get_event_loop())
-
- def run_server(self, loop):
-
- logging.info("Starting MQTT Server at {}".format(self.address))
- print("Starting MQTT Server at {}".format(self.address))
- try:
- asyncio.set_event_loop(loop)
- loop.run_until_complete(self.broker_coro())
- # loop.run_until_complete(self.active_bot_listing())
- loop.run_forever()
-
- except Exception as e:
- mqttserverlog.exception("{}".format(e))
-
class BumperMQTTServer_Plugin:
def __init__(self, context):
@@ -310,11 +256,9 @@ class BumperMQTTServer_Plugin:
client_id = session.client_id
didsplit = str(client_id).split("@")
- # If this isn't a fake user (fuid) then add as a bot
- if not (
- str(didsplit[0]).startswith("fuid")
- or str(didsplit[0]).startswith("helper")
- ):
+ if not ( # if ecouser or bumper aren't in details it is a bot
+ "ecouser" in didsplit[1]
+ or "bumper" in didsplit[1]):
tmpbotdetail = str(didsplit[1]).split("/")
bumper.bot_add(
username,
diff --git a/bumper/xmppserver.py b/bumper/xmppserver.py
index 5ad8e22..7d00fc5 100644
--- a/bumper/xmppserver.py
+++ b/bumper/xmppserver.py
@@ -5,7 +5,7 @@ import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as
import base64
import ssl
import bumper
-import asyncio
+import asyncio, functools
xmppserverlog = logging.getLogger("xmppserver")
@@ -19,11 +19,10 @@ class XMPPServer:
def __init__(self, address):
# Initialize bot server
self.address = address
- self.aclients = {} # task -> (reader, writer)
async def async_server(self):
- xmppserverlog.debug(
- "listening on {}:{}".format(self.address[0], self.address[1])
+ xmppserverlog.info(
+ "Starting XMPP Server at {}:{}".format(self.address[0], self.address[1])
)
server = await asyncio.start_server(
self.accept_client, self.address[0], self.address[1]
@@ -34,101 +33,26 @@ class XMPPServer:
# self.clients = {} # task -> (reader, writer)
def accept_client(self, client_reader, client_writer):
- # task = asyncio.Task(self.handle_client(client_reader, client_writer))
- aclient = XMPPAsyncClient(client_reader, client_writer)
- task = asyncio.Task(aclient.handle_async_client())
- self.aclients[task] = (client_reader, client_writer)
+ try:
+ aclient = XMPPAsyncClient(client_reader, client_writer)
+ task = asyncio.Task(aclient.handle_async_client())
+ aclient._async_task = task
+ self.clients.append(aclient)
- def client_done(task):
- del self.aclients[task]
- xmppserverlog.info("End Connection for {}".format(client_writer.get_extra_info("peername")))
- client_writer.close()
+ def client_done(aclient, task):
+ try:
+ self.clients.remove(aclient)
+ xmppserverlog.debug("End Connection for ({}:{} | {})".format(client_writer.get_extra_info("peername")[0], client_writer.get_extra_info("peername")[1], aclient.bumper_jid))
+ client_writer.close()
+ except Exception as e:
+ xmppserverlog.error("{}".format(e))
- clientaddr = client_writer.get_extra_info("peername")
- xmppserverlog.info("New Connection from {}".format(clientaddr))
- task.add_done_callback(client_done)
+ clientaddr = client_writer.get_extra_info("peername")
+ xmppserverlog.debug("New Connection from {}:{}".format(clientaddr[0],clientaddr[1]))
+ task.add_done_callback(functools.partial(client_done, aclient))
- # def run(self, run_async=False):
- # if run_async:
- # xmppserverlog.debug("Starting XMPPServer Thread: 1")
- # self.xmppthread = Thread(name="XMPPServer_Thread", target=self.run_server)
- # self.xmppthread.setDaemon(True)
- # self.xmppthread.start()
-
- # else:
- # try:
- # self.run_server()
- # except KeyboardInterrupt:
- # self.disconnect()
-
- # def run_server(self):
- # logging.info("Starting XMPP Server at {}".format(self.address))
- # print("Starting XMPP Server at {}".format(self.address))
-
- # # xmppserverlog.setLevel(logging.DEBUG)
-
- # # Set SSL Context
- # self.ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
- # self.ssl_ctx.load_cert_chain(
- # certfile=bumper.server_cert, keyfile=bumper.server_key
- # )
-
- # self.socket = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
- # self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
-
- # try:
- # self.socket.bind(self.address)
- # self.socket.listen(5)
-
- # xmppserverlog.debug(
- # "listening on {}:{}".format(self.address[0], self.address[1])
- # )
- # while not self.exit_flag:
- # connection, client_address = self.socket.accept()
-
- # # disconnect any clients with this ip
- # for client in self.clients:
- # if client.address == client_address[0]:
- # xmppserverlog.debug(
- # "disconnecting existing client {} with resource {}".format(
- # client.address, client.clientresource
- # )
- # )
- # client._disconnect()
- # self.remove_client_byip(client.address)
-
- # xmppserverlog.debug(
- # "starting new client with ip {}".format(client_address[0])
- # )
- # thread_id = uuid.uuid4()
- # client = XMPPAsyncClient(thread_id, connection, client_address)
- # client.setDaemon(True)
- # client.start()
- # self.clients.append(client)
-
- # except PermissionError as e:
- # if "bind" in e.strerror:
- # xmppserverlog.exception(
- # "Error binding XMPPServer, exiting. Try using a different hostname or IP - {}".format(
- # e
- # )
- # )
- # exit(1)
-
- # except Exception as e:
- # xmppserverlog.exception("{}".format(e))
- # exit(1)
-
- # except KeyboardInterrupt as e:
- # xmppserverlog.exception("{}".format(e))
-
- # finally:
- # connection.shutdown(socket.SHUT_RDWR)
- # connection.close()
- # self.disconnect()
- # xmppserverlog.info("disconnecting")
-
- # self.socket.close()
+ except Exception as e:
+ xmppserverlog.error("{}".format(e))
def disconnect(self):
try:
@@ -140,41 +64,7 @@ class XMPPServer:
xmppserverlog.debug("shutting down")
except Exception as e:
- xmppserverlog.exception("{}".format(e))
-
- def remove_client_byip(self, ip):
- for client in self.clients:
- if client.address == ip:
- xmppserverlog.debug(
- "removing client from client list with ip {} and resource {}".format(
- client.address, client.clientresource
- )
- )
- client._disconnect()
- self.clients.remove(client)
-
- def remove_client_byresource(self, resource):
- for client in self.clients:
- if str(client.clientresource).lower() == str(resource).lower():
- xmppserverlog.debug(
- "removing client from client list with ip {} and resource {}".format(
- client.address, client.clientresource
- )
- )
- client._disconnect()
- self.clients.remove(client)
-
- def remove_client_byuid(self, uid):
- for client in self.clients:
- if str(client.uid).lower() == str(uid).lower():
- xmppserverlog.debug(
- "removing client from client list with ip {} and resource {}".format(
- client.address, client.clientresource
- )
- )
- client._disconnect()
- self.clients.remove(client)
-
+ xmppserverlog.error("{}".format(e))
class XMPPAsyncClient:
IDLE = 0
@@ -186,6 +76,7 @@ class XMPPAsyncClient:
UNKNOWN = 0
BOT = 1
CONTROLLER = 2
+ _async_task = None
def __init__(self, client_reader, client_writer):
self.type = self.UNKNOWN
@@ -204,20 +95,17 @@ class XMPPAsyncClient:
async def handle_async_client(self):
# xmppserverlog.info('client connected - {}'.format(self.address))
- #await self._set_state("READY")
- await self._set_state("CONNECT")
- #asyncio.Task(self.send_ping(30))
+ await self._set_state("CONNECT")
while True:
- await asyncio.sleep(0.1)
+ await asyncio.sleep(0.05)
if not self.state == self.DISCONNECT:
data = await self.client_reader.read(4096)
- # data = await asyncio.wait_for(client_reader.readline(), timeout=10.0)
- if data is None:
- xmppserverlog.warning("Received no data")
+ if data is None or data == b'':
+ xmppserverlog.debug("Received no data")
# exit loop and disconnect
return
-
- await self._parse_data(data)
+ else:
+ await self._parse_data(data)
else:
break
@@ -228,25 +116,25 @@ class XMPPAsyncClient:
try:
# if not self.connection._closed:
if self.log_sent_message:
- xmppserverlog.debug("send {} - {}".format(self.address, command))
- # self.connection.send(command.encode())
+ xmppserverlog.debug("send to ({}:{} | {}) - {}".format(self.address[0], self.address[1], self.bumper_jid, command))
+
self.client_writer.write(command.encode())
await self.client_writer.drain()
except BrokenPipeError as e:
- xmppserverlog.debug("{}".format(e))
+ #xmppserverlog.debug("{}".format(e))
await self._set_state("DISCONNECT")
except ConnectionResetError as e:
- xmppserverlog.debug("{}".format(e))
+ #xmppserverlog.debug("{}".format(e))
await self._set_state("DISCONNECT")
except ConnectionAbortedError as e:
- xmppserverlog.debug("{}".format(e))
+ #xmppserverlog.debug("{}".format(e))
await self._set_state("DISCONNECT")
except OSError as e:
- xmppserverlog.debug("{}".format(e))
+ xmppserverlog.error("{}".format(e))
except Exception as e:
xmppserverlog.exception("{}".format(e))
@@ -265,7 +153,7 @@ class XMPPAsyncClient:
self.client_writer.close()
except Exception as e:
- xmppserverlog.exception("{}".format(e))
+ xmppserverlog.error("{}".format(e))
async def _tag_strip_uri(self, tag):
try:
@@ -274,7 +162,7 @@ class XMPPAsyncClient:
return tag
except Exception as e:
- xmppserverlog.exception("{}".format(e))
+ xmppserverlog.error("{}".format(e))
async def _set_state(self, state):
try:
@@ -286,7 +174,7 @@ class XMPPAsyncClient:
)
)
- xmppserverlog.debug("{} state: {}".format(self.address, state))
+ xmppserverlog.debug("({}:{} | {}) state: {}".format(self.address[0],self.address[1],self.bumper_jid, state))
self.state = new_state
@@ -294,7 +182,7 @@ class XMPPAsyncClient:
await self._disconnect()
except Exception as e:
- xmppserverlog.exception("{}".format(e))
+ xmppserverlog.error("{}".format(e))
async def _handle_ctl(self, xml, data):
try:
@@ -365,13 +253,13 @@ class XMPPAsyncClient:
if client.type == self.BOT:
if client.uid.lower() in ctl_to.lower():
- xmppserverlog.info(
+ xmppserverlog.debug(
"Sending ctl to bot: {}".format(rxmlstring)
)
- client.send(rxmlstring)
+ await client.send(rxmlstring)
except Exception as e:
- xmppserverlog.exception("{}".format(e))
+ xmppserverlog.error("{}".format(e))
async def _handle_ping(self, xml, data):
try:
@@ -387,37 +275,40 @@ class XMPPAsyncClient:
pingto = xml.get("to")
pingfrom = self.bumper_jid
- xml.attrib["from"] = pingfrom
+ xml.attrib["from"] = pingfrom
pingstring = ET.tostring(xml).decode("utf-8")
# clean up string to remove namespaces added by ET
pingstring = pingstring.replace("xmlns:ns0=", "xmlns=")
pingstring = pingstring.replace("ns0:", "")
- pingstring = pingstring.replace('iq xmlns="com:ctl"', "iq")
- pingstring = pingstring.replace("'.format(
+ # xml.get("id"), pingfrom, pingto
+ #)
+ #xmppserverlog.debug(
+ # "ping from {} to {} with {}".format(pingfrom, pingto, pingstring)
+ #)
+
+ await client.send(pingstring)
except Exception as e:
xmppserverlog.exception("{}".format(e))
- async def send_ping(self, time):
+
+ async def schedule_ping(self, time):
if not self.state == 5: #disconnected
pingstring = "".format(XMPPServer.server_id, self.bumper_jid)
await self.send(pingstring)
await asyncio.sleep(time)
- asyncio.Task(self.send_ping(time))
+ asyncio.Task(self.schedule_ping(time))
async def _handle_result(self, xml, data):
try:
@@ -433,7 +324,7 @@ class XMPPAsyncClient:
adminuser = ctlerr.replace("permission denied, please contact ", "")
adminuser = adminuser.replace(" ", "")
if not (
- adminuser.startswith("fuid_") or bumper.use_auth
+ adminuser.startswith("fuid_") or adminuser.startswith("fusername_") or bumper.use_auth
): # if not fuid_ then its ecovacs OR ignore bumper auth
# TODO: Implement auth later, should this user have access to bot?
@@ -474,7 +365,7 @@ class XMPPAsyncClient:
)
)
for client in XMPPServer.clients:
- client.send(rxmlstring)
+ await client.send(rxmlstring)
if xml.get("to").find("@") == -1: # No to address
ctl_to = xml.get("to")
@@ -488,7 +379,7 @@ class XMPPAsyncClient:
):
if not "@" in ctl_to: # No user@, send to all clients?
# TODO: Revisit later, this may be wrong
- client.send(rxmlstring)
+ await client.send(rxmlstring)
elif (
client.uid.lower() in ctl_to.lower()
@@ -498,7 +389,7 @@ class XMPPAsyncClient:
self.uid, client.uid, rxmlstring
)
)
- client.send(rxmlstring)
+ await client.send(rxmlstring)
except Exception as e:
xmppserverlog.exception("{}".format(e))
@@ -617,11 +508,9 @@ class XMPPAsyncClient:
self.clientresource = aitem.text
resource = self.clientresource
- if not self.uid.startswith("fuid"):
-
- # Need sample data to see details here
+ if self.devclass: # if there is a devclass it is a bot
bumper.bot_add("", self.uid, "", resource, "eco-legacy")
- xmppserverlog.info("bot authenticated {}".format(self.uid))
+ xmppserverlog.debug("bot authenticated {}".format(self.uid))
# Client authenticated, move to next state
await self._set_state("INIT")
@@ -690,12 +579,11 @@ class XMPPAsyncClient:
if len(saslauth) > 2:
authcode = saslauth[2]
-
- if not self.uid.startswith("fuid"):
- # Need sample data to see details here
+
+ if self.devclass: # if there is a devclass it is a bot
bumper.bot_add(self.uid, self.uid, self.devclass, "atom", "eco-legacy")
self.type = self.BOT
- xmppserverlog.info("bot authenticated {}".format(self.uid))
+ xmppserverlog.debug("bot authenticated {}".format(self.uid))
# Send response
await self.send(
''
@@ -751,7 +639,7 @@ class XMPPAsyncClient:
self.bumper_jid = "{}@{}.ecorobot.net/atom".format(
self.uid, self.devclass
)
- xmppserverlog.debug("new bot {}".format(self.uid))
+ xmppserverlog.debug("new bot ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid))
res = '{}'.format(
xml.get("id"), self.bumper_jid
)
@@ -761,18 +649,14 @@ class XMPPAsyncClient:
self.bumper_jid = "{}@{}/{}".format(
self.uid, XMPPServer.server_id, self.clientresource
)
- xmppserverlog.debug(
- "new client {} using resource {}".format(
- self.uid, self.clientresource
- )
- )
+ xmppserverlog.debug("new client ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid))
res = '{}'.format(
xml.get("id"), self.bumper_jid
)
else:
self.name = "XMPP_Client_{}_{}".format(self.uid, self.address)
self.bumper_jid = "{}@{}".format(self.uid, XMPPServer.server_id)
- xmppserverlog.debug("new client {}".format(self.uid))
+ xmppserverlog.debug("new client ({}:{} | {})".format(self.address[0],self.address[1], self.bumper_jid))
res = '{}'.format(
xml.get("id"), self.bumper_jid
)
@@ -788,7 +672,7 @@ class XMPPAsyncClient:
res = ''.format(xml.get("id"))
await self._set_state("READY")
await self.send(res)
- asyncio.Task(self.send_ping(30))
+ asyncio.Task(self.schedule_ping(30))
except Exception as e:
@@ -799,7 +683,7 @@ class XMPPAsyncClient:
if len(xml) and xml[0].tag == "status":
xmppserverlog.debug(
- "bot presence {} ".format(ET.tostring(xml, encoding="utf-8"))
+ "bot presence {} ".format(ET.tostring(xml, encoding="utf-8").decode("utf-8"))
)
# Most likely a bot, possibly hello world in text
@@ -823,15 +707,15 @@ class XMPPAsyncClient:
else:
xmppserverlog.debug(
- "client presence - {} ".format(ET.tostring(xml, encoding="utf-8"))
+ "client presence - {} ".format(ET.tostring(xml, encoding="utf-8").decode("utf-8"))
)
if xml.get("type") == "available":
xmppserverlog.debug(
"client presence available - {} ".format(
- ET.tostring(xml, encoding="utf-8")
- )
+ ET.tostring(xml, encoding="utf-8").decode("utf-8"))
)
+
# Send dummy return
await self.send(
' dummy '.format(self.bumper_jid)
@@ -839,9 +723,8 @@ class XMPPAsyncClient:
elif xml.get("type") == "unavailable":
xmppserverlog.debug(
"client presence unavailable (DISCONNECT) - {} ".format(
- ET.tostring(xml, encoding="utf-8")
- )
- )
+ ET.tostring(xml, encoding="utf-8").decode("utf-8"))
+ )
await self._set_state("DISCONNECT")
else:
@@ -880,8 +763,8 @@ class XMPPAsyncClient:
if item.tag == "iq":
if self.log_incoming_data:
xmppserverlog.debug(
- "from {} - {}".format(
- self.address,
+ "from ({}:{} | {}) - {}".format(
+ self.address[0],self.address[1],self.bumper_jid,
str(
ET.tostring(item, encoding="utf-8").decode(
"utf-8"
diff --git a/start_bumper.py b/start_bumper.py
index 2503f44..a01cd81 100644
--- a/start_bumper.py
+++ b/start_bumper.py
@@ -6,11 +6,16 @@ import sys, socket
import time
import platform
import os
-os.environ['PYTHONASYNCIODEBUG'] = '1'
+#os.environ['PYTHONASYNCIODEBUG'] = '1' # Uncomment to enable ASYNCIODEBUG
import asyncio
-def main():
+async def main():
+ try:
+ loop = asyncio.get_event_loop()
+ except:
+ loop = asyncio.new_event_loop()
+
args = sys.argv
listen_host = ""
@@ -20,12 +25,13 @@ def main():
level=logging.DEBUG,
format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s",
)
+ loop.set_debug(True) # Set asyncio loop to debug
+ #logging.getLogger("asyncio").setLevel(logging.DEBUG) # Show debug asyncio logs (disabled in init, uncomment for debugging asyncio)
else:
logging.basicConfig(
level=logging.INFO,
format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(message)s",
)
- # format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s")
if "--listen" in args:
listen_host = args[args.index("--listen") + 1]
@@ -53,54 +59,39 @@ def main():
conf_address_8007, usessl=False, helperbot=mqtt_helperbot
)
- try:
- loop = asyncio.get_event_loop()
- except:
- loop = asyncio.new_event_loop()
-
# Start web servers
- loop.set_debug(True)
- conf_server.confserver_app()
+ conf_server.confserver_app()
+ task_conf_server = asyncio.create_task(conf_server.start_server())
+ bumper.bumperlog.debug("task_conf_server added")
+ await task_conf_server
+
conf_server_2.confserver_app()
- asyncio.ensure_future(conf_server.start_server(),loop=loop)
- asyncio.ensure_future(conf_server_2.start_server(),loop=loop)
+ task_conf_server2 = asyncio.create_task(conf_server_2.start_server())
+ bumper.bumperlog.debug("task_conf_server2 added")
+ await task_conf_server2
# Start MQTT Server
- asyncio.ensure_future(mqtt_server.broker_coro())
+ task_mqtt_server = asyncio.create_task(mqtt_server.broker_coro())
+ bumper.bumperlog.debug("task_mqtt_server added")
+ await task_mqtt_server
# Start MQTT Helperbot
- asyncio.ensure_future(mqtt_helperbot.start_helper_bot())
+ task_mqtt_helperbot = asyncio.create_task(mqtt_helperbot.start_helper_bot())
+ bumper.bumperlog.debug("task_mqtt_helperbot added")
+ await task_mqtt_helperbot
# Start XMPP Server
- asyncio.ensure_future(xmpp_server.async_server())
-
- loop.run_forever()
-
-
- # start xmpp server on port 5223 (sync)
- #xmpp_server.run(run_async=True) # Start in new thread
-
- # start mqtt server on port 8883 (async)
- #mqtt_server.run(run_async=True) # Start in new thread
-
- #time.sleep(1.5) # Wait for broker startup
-
- # start mqtt_helperbot (async)
- #mqtt_helperbot.run(run_async=True) # Start in new thread
-
- # start conf server on port 443 (async) - Used for most https calls
- #conf_server.run(run_async=True) # Start in new thread
-
- # start conf server on port 8007 (async) - Used for a load balancer request
- #conf_server_2.run(run_async=True) # Start in new thread
+ task_xmpp_server = asyncio.create_task(xmpp_server.async_server())
+ bumper.bumperlog.debug("task_xmpp_server added")
+ await task_xmpp_server
while True:
try:
- time.sleep(30)
+ await asyncio.sleep(30)
bumper.revoke_expired_tokens()
- disconnected_clients = bumper.get_disconnected_xmpp_clients()
- for client in disconnected_clients:
- xmpp_server.remove_client_byuid(client["userid"])
+ #disconnected_clients = bumper.get_disconnected_xmpp_clients()
+ #for client in disconnected_clients:
+ # xmpp_server.remove_client_byuid(client["userid"])
except KeyboardInterrupt:
bumper.bumperlog.info("Bumper Exiting - Keyboard Interrupt")
@@ -109,4 +100,4 @@ def main():
if __name__ == "__main__":
- main()
+ asyncio.run(main())