Merge async and xmpp work #34

Merged
bmartin5692 merged 20 commits from dev_broken-XMPP into master 2019-05-23 15:22:58 +02:00
5 changed files with 220 additions and 402 deletions
Showing only changes of commit 2460182f4f - Show all commits

View file

@ -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):

View file

@ -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):
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)
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()
body = {
"code": bumper.ERR_TOKEN_INVALID,
"data": None,
"msg": "当前密码错误",
"time": bumper.get_milli_time(datetime.utcnow().timestamp()),
}
return web.json_response(body)
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)
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:

View file

@ -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,

View file

@ -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))
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:
@ -392,8 +280,9 @@ class XMPPAsyncClient:
# 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("<query", '<query xmlns="com:ctl"')
pingstring = pingstring.replace('iq xmlns="urn:xmpp:ping"', "iq")
pingstring = pingstring.replace("<ping", '<ping xmlns="urn:xmpp:ping"')
for client in XMPPServer.clients:
if (
@ -401,23 +290,25 @@ class XMPPAsyncClient:
and client.state == client.READY
):
if pingto.lower() in client.bumper_jid.lower():
pingstring = '<iq type="result" id="{}" from="{}" to="{}" />'.format(
xml.get("id"), pingfrom, pingto
)
xmppserverlog.debug(
"ping from {} to {}".format(pingfrom, pingto)
)
client.send(pingstring)
#pingstring = '<iq type="result" id="{}" from="{}" to="{}" />'.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 = "<iq from='{}' to='{}' id='s2c1' type='get'><ping xmlns='urn:xmpp:ping'/></iq>".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")
@ -691,11 +580,10 @@ 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(
'<success xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
@ -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 = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.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 = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.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 = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format(
xml.get("id"), self.bumper_jid
)
@ -788,7 +672,7 @@ class XMPPAsyncClient:
res = '<iq type="result" id="{}" />'.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(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
@ -839,8 +723,7 @@ 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")
@ -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"

View file

@ -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()
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())