Re-org and start of API #47
4 changed files with 229 additions and 151 deletions
|
|
@ -163,23 +163,12 @@ async def start():
|
|||
global mqtt_helperbot
|
||||
mqtt_helperbot = MQTTHelperBot((bumper_listen, mqtt_listen_port))
|
||||
global conf_server
|
||||
conf_server = ConfServer(
|
||||
(bumper_listen, conf1_listen_port), usessl=True, helperbot=mqtt_helperbot
|
||||
)
|
||||
conf_server = ConfServer((bumper_listen, conf1_listen_port), usessl=True)
|
||||
global conf_server_2
|
||||
conf_server_2 = ConfServer(
|
||||
(bumper_listen, conf2_listen_port), usessl=False, helperbot=mqtt_helperbot
|
||||
)
|
||||
conf_server_2 = ConfServer((bumper_listen, conf2_listen_port), usessl=False)
|
||||
global xmpp_server
|
||||
xmpp_server = XMPPServer((bumper_listen, xmpp_listen_port))
|
||||
|
||||
# Start web servers
|
||||
conf_server.confserver_app()
|
||||
asyncio.create_task(conf_server.start_server())
|
||||
|
||||
conf_server_2.confserver_app()
|
||||
asyncio.create_task(conf_server_2.start_server())
|
||||
|
||||
# Start MQTT Server
|
||||
asyncio.create_task(mqtt_server.broker_coro())
|
||||
|
||||
|
|
@ -189,6 +178,20 @@ async def start():
|
|||
# Start XMPP Server
|
||||
asyncio.create_task(xmpp_server.start_async_server())
|
||||
|
||||
# Wait for helperbot to connect first
|
||||
while mqtt_helperbot.Client is None:
|
||||
await asyncio.sleep(0.1)
|
||||
|
||||
while not mqtt_helperbot.Client.session.transitions.state == "connected":
|
||||
await asyncio.sleep(0.1)
|
||||
|
||||
# Start web servers
|
||||
conf_server.confserver_app()
|
||||
asyncio.create_task(conf_server.start_server())
|
||||
|
||||
conf_server_2.confserver_app()
|
||||
asyncio.create_task(conf_server_2.start_server())
|
||||
|
||||
# Start maintenance
|
||||
while not shutting_down:
|
||||
asyncio.create_task(maintenance())
|
||||
|
|
|
|||
|
|
@ -38,12 +38,9 @@ logging.getLogger("aiohttp.access").addFilter(
|
|||
|
||||
|
||||
class ConfServer:
|
||||
def __init__(self, address, usessl=False, helperbot=None):
|
||||
self.helperbot = helperbot
|
||||
def __init__(self, address, usessl=False):
|
||||
self.usessl = usessl
|
||||
self.address = address
|
||||
self.confthread = None
|
||||
self.run_async = False
|
||||
self.app = None
|
||||
self.site = None
|
||||
self.runner = None
|
||||
|
|
@ -57,6 +54,7 @@ class ConfServer:
|
|||
self.app.add_routes(
|
||||
[
|
||||
web.get("", self.handle_base),
|
||||
web.get("/restart_{service}", self.handle_RestartService),
|
||||
web.get(
|
||||
"/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/login",
|
||||
self.handle_login,
|
||||
|
|
@ -198,11 +196,29 @@ class ConfServer:
|
|||
# text = "Bumper!"
|
||||
bots = bumper.db_get().table("bots").all()
|
||||
clients = bumper.db_get().table("clients").all()
|
||||
helperbot = self.helperbot.Client.session.transitions.state
|
||||
helperbot = bumper.mqtt_helperbot.Client.session.transitions.state
|
||||
mqttserver = bumper.mqtt_server.broker
|
||||
mq_sessions = []
|
||||
for sess in mqttserver._sessions:
|
||||
tmpsess = []
|
||||
tmpsess.append({"client_id": mqttserver._sessions[sess][0].client_id})
|
||||
tmpsess.append(
|
||||
{"state": mqttserver._sessions[sess][0].transitions.state}
|
||||
)
|
||||
mq_sessions.append(tmpsess)
|
||||
all = {
|
||||
"bots": bots,
|
||||
"clients": clients,
|
||||
"helperbot": [{"state": helperbot}],
|
||||
"mqtt_server": [
|
||||
{"state": mqttserver.transitions.state},
|
||||
{
|
||||
"sessions": [
|
||||
{"count": len(mqttserver._sessions)},
|
||||
{"clients": mq_sessions},
|
||||
]
|
||||
},
|
||||
],
|
||||
}
|
||||
|
||||
return web.json_response(all)
|
||||
|
|
@ -210,6 +226,49 @@ class ConfServer:
|
|||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
|
||||
async def restart_Helper(self):
|
||||
|
||||
await bumper.mqtt_helperbot.Client.disconnect()
|
||||
await bumper.mqtt_helperbot.start_helper_bot()
|
||||
|
||||
|
||||
|
||||
async def restart_MQTT(self):
|
||||
mqttserver = bumper.mqtt_server.broker
|
||||
|
||||
for sess in list(mqttserver._sessions):
|
||||
sessobj = mqttserver._sessions[sess][1]
|
||||
await sessobj.writer.close()
|
||||
mqttserver.delete_session(sess)
|
||||
|
||||
await bumper.mqtt_server.broker.shutdown()
|
||||
while not bumper.mqtt_server.broker.transitions.state == "stopped":
|
||||
await asyncio.sleep(0.1)
|
||||
|
||||
await bumper.mqtt_server.broker_coro()
|
||||
while not bumper.mqtt_server.broker.transitions.state == "started":
|
||||
await asyncio.sleep(0.1)
|
||||
|
||||
async def handle_RestartService(self, request):
|
||||
try:
|
||||
service = request.match_info.get("service", "")
|
||||
if service == "Helperbot":
|
||||
await self.restart_Helper()
|
||||
return web.json_response({"status": "complete"})
|
||||
elif service == "MQTTServer":
|
||||
await self.restart_MQTT()
|
||||
aloop = asyncio.get_event_loop()
|
||||
aloop.call_later(
|
||||
2, lambda: asyncio.create_task(self.restart_Helper())
|
||||
) # In 2 seconds restart Helperbot
|
||||
|
||||
return web.json_response({"status": "complete"})
|
||||
else:
|
||||
return web.json_response({"status": "invalid service"})
|
||||
|
||||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
|
||||
async def handle_login(self, request):
|
||||
try:
|
||||
user_devid = request.match_info.get("devid", "")
|
||||
|
|
@ -259,9 +318,7 @@ class ConfServer:
|
|||
# "username": "fusername_{}".format(tmpuser["userid"]),
|
||||
# },
|
||||
"msg": "操作成功",
|
||||
"time": self.self.get_milli_time(
|
||||
datetime.utcnow().timestamp()
|
||||
),
|
||||
"time": self.get_milli_time(datetime.utcnow().timestamp()),
|
||||
}
|
||||
|
||||
return web.json_response(body)
|
||||
|
|
@ -1048,122 +1105,136 @@ class ConfServer:
|
|||
confserverlog.exception("{}".format(e))
|
||||
|
||||
async def handle_lg_log(self, request): # EcoVacs Home
|
||||
try:
|
||||
json_body = json.loads(await request.text())
|
||||
if (
|
||||
not bumper.mqtt_helperbot.Client._handler.writer is None
|
||||
): # Ignore if the Helperbot writer is none
|
||||
try:
|
||||
json_body = json.loads(await request.text())
|
||||
|
||||
randomid = "".join(random.sample(string.ascii_letters, 6))
|
||||
did = json_body["did"]
|
||||
randomid = "".join(random.sample(string.ascii_letters, 6))
|
||||
did = json_body["did"]
|
||||
|
||||
botdetails = bumper.bot_get(did)
|
||||
if botdetails:
|
||||
if not "cmdName" in json_body:
|
||||
if "td" in json_body:
|
||||
json_body["cmdName"] = json_body["td"]
|
||||
# json_body["td"] = "q"
|
||||
botdetails = bumper.bot_get(did)
|
||||
if botdetails:
|
||||
if not "cmdName" in json_body:
|
||||
if "td" in json_body:
|
||||
json_body["cmdName"] = json_body["td"]
|
||||
# json_body["td"] = "q"
|
||||
|
||||
if not "toId" in json_body:
|
||||
json_body["toId"] = did
|
||||
if not "toId" in json_body:
|
||||
json_body["toId"] = did
|
||||
|
||||
if not "toType" in json_body:
|
||||
json_body["toType"] = botdetails["class"]
|
||||
if not "toType" in json_body:
|
||||
json_body["toType"] = botdetails["class"]
|
||||
|
||||
if not "toRes" in json_body:
|
||||
json_body["toRes"] = botdetails["resource"]
|
||||
if not "toRes" in json_body:
|
||||
json_body["toRes"] = botdetails["resource"]
|
||||
|
||||
if not "payloadType" in json_body:
|
||||
json_body["payloadType"] = "x"
|
||||
if not "payloadType" in json_body:
|
||||
json_body["payloadType"] = "x"
|
||||
|
||||
if not "payload" in json_body:
|
||||
json_body["payload"] = ""
|
||||
if json_body["td"] == "GetCleanLogs":
|
||||
json_body["td"] = "q"
|
||||
json_body["payload"] = '<ctl count="30"/>' # <ctl />"
|
||||
if not "payload" in json_body:
|
||||
json_body["payload"] = ""
|
||||
if json_body["td"] == "GetCleanLogs":
|
||||
json_body["td"] = "q"
|
||||
json_body["payload"] = '<ctl count="30"/>' # <ctl />"
|
||||
|
||||
if did != "":
|
||||
bot = bumper.bot_get(did)
|
||||
if bot["company"] == "eco-ng":
|
||||
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":
|
||||
cleanlogs = logsroot.getchildren()
|
||||
for l in cleanlogs:
|
||||
logs.append(l.attrib)
|
||||
|
||||
body = {
|
||||
"ret": "ok",
|
||||
# "logs": logs, #TODO: Doesn't parse correctly, new protocol & server side processing
|
||||
"logs": [],
|
||||
}
|
||||
|
||||
else:
|
||||
body = {"ret": "ok", "logs": []}
|
||||
|
||||
confserverlog.debug(
|
||||
"POST: {} - Response: {}".format(json_body, body)
|
||||
)
|
||||
|
||||
return web.json_response(body)
|
||||
else:
|
||||
# No response, send error back
|
||||
confserverlog.error(
|
||||
"No bots with DID: {} connected to MQTT".format(
|
||||
json_body["toId"]
|
||||
if did != "":
|
||||
bot = bumper.bot_get(did)
|
||||
if bot["company"] == "eco-ng":
|
||||
body = ""
|
||||
retcmd = await bumper.mqtt_helperbot.send_command(
|
||||
json_body, randomid
|
||||
)
|
||||
)
|
||||
body = {"id": randomid, "errno": bumper.ERR_COMMON, "ret": "fail"}
|
||||
return web.json_response(body)
|
||||
confserverlog.debug("Send Bot - {}".format(json_body))
|
||||
confserverlog.debug("Bot Response - {}".format(body))
|
||||
logs = []
|
||||
logsroot = ET.fromstring(retcmd["resp"])
|
||||
if logsroot.attrib["ret"] == "ok":
|
||||
cleanlogs = logsroot.getchildren()
|
||||
for l in cleanlogs:
|
||||
logs.append(l.attrib)
|
||||
|
||||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
body = {
|
||||
"ret": "ok",
|
||||
# "logs": logs, #TODO: Doesn't parse correctly, new protocol & server side processing
|
||||
"logs": [],
|
||||
}
|
||||
|
||||
else:
|
||||
body = {"ret": "ok", "logs": []}
|
||||
|
||||
confserverlog.debug(
|
||||
"POST: {} - Response: {}".format(json_body, body)
|
||||
)
|
||||
|
||||
return web.json_response(body)
|
||||
else:
|
||||
# No response, send error back
|
||||
confserverlog.error(
|
||||
"No bots with DID: {} connected to MQTT".format(
|
||||
json_body["toId"]
|
||||
)
|
||||
)
|
||||
body = {
|
||||
"id": randomid,
|
||||
"errno": bumper.ERR_COMMON,
|
||||
"ret": "fail",
|
||||
}
|
||||
return web.json_response(body)
|
||||
|
||||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
|
||||
async def handle_devmanager_botcommand(self, request):
|
||||
try:
|
||||
json_body = json.loads(await request.text())
|
||||
if (
|
||||
not bumper.mqtt_helperbot.Client._handler.writer is None
|
||||
): # Ignore if the helperbot object isn't set
|
||||
try:
|
||||
json_body = json.loads(await request.text())
|
||||
|
||||
randomid = "".join(random.sample(string.ascii_letters, 6))
|
||||
did = ""
|
||||
if "toId" in json_body: # Its a command
|
||||
did = json_body["toId"]
|
||||
randomid = "".join(random.sample(string.ascii_letters, 6))
|
||||
did = ""
|
||||
if "toId" in json_body: # Its a command
|
||||
did = json_body["toId"]
|
||||
|
||||
if did != "":
|
||||
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("Send Bot - {}".format(json_body))
|
||||
confserverlog.debug("Bot Response - {}".format(body))
|
||||
return web.json_response(body)
|
||||
else:
|
||||
# No response, send error back
|
||||
confserverlog.error(
|
||||
"No bots with DID: {} connected to MQTT".format(
|
||||
json_body["toId"]
|
||||
if did != "":
|
||||
bot = bumper.bot_get(did)
|
||||
if bot["company"] == "eco-ng":
|
||||
retcmd = await bumper.mqtt_helperbot.send_command(
|
||||
json_body, randomid
|
||||
)
|
||||
)
|
||||
body = {
|
||||
"id": randomid,
|
||||
"errno": 500,
|
||||
"ret": "fail",
|
||||
"debug": "wait for response timed out",
|
||||
}
|
||||
return web.json_response(body)
|
||||
|
||||
else:
|
||||
if "td" in json_body: # Seen when doing initial wifi config
|
||||
if json_body["td"] == "PollSCResult":
|
||||
body = {"ret": "ok"}
|
||||
body = retcmd
|
||||
confserverlog.debug("Send Bot - {}".format(json_body))
|
||||
confserverlog.debug("Bot Response - {}".format(body))
|
||||
return web.json_response(body)
|
||||
else:
|
||||
# No response, send error back
|
||||
confserverlog.error(
|
||||
"No bots with DID: {} connected to MQTT".format(
|
||||
json_body["toId"]
|
||||
)
|
||||
)
|
||||
body = {
|
||||
"id": randomid,
|
||||
"errno": 500,
|
||||
"ret": "fail",
|
||||
"debug": "wait for response timed out",
|
||||
}
|
||||
return web.json_response(body)
|
||||
|
||||
if json_body["td"] == "HasUnreadMsg": # EcoVacs Home
|
||||
body = {"ret": "ok", "unRead": False}
|
||||
return web.json_response(body)
|
||||
else:
|
||||
if "td" in json_body: # Seen when doing initial wifi config
|
||||
if json_body["td"] == "PollSCResult":
|
||||
body = {"ret": "ok"}
|
||||
return web.json_response(body)
|
||||
|
||||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
if json_body["td"] == "HasUnreadMsg": # EcoVacs Home
|
||||
body = {"ret": "ok", "unRead": False}
|
||||
return web.json_response(body)
|
||||
|
||||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
|
||||
async def handle_dim_devmanager(self, request): # Used in EcoVacs Home App
|
||||
try:
|
||||
|
|
@ -1177,7 +1248,9 @@ class ConfServer:
|
|||
if did != "":
|
||||
bot = bumper.bot_get(did)
|
||||
if bot["company"] == "eco-ng" and bot["mqtt_connection"] == True:
|
||||
retcmd = await self.helperbot.send_command(json_body, randomid)
|
||||
retcmd = await bumper.mqtt_helperbot.send_command(
|
||||
json_body, randomid
|
||||
)
|
||||
body = retcmd
|
||||
confserverlog.debug("Send Bot - {}".format(json_body))
|
||||
confserverlog.debug("Bot Response - {}".format(body))
|
||||
|
|
@ -1208,10 +1281,7 @@ class ConfServer:
|
|||
async def disconnect(self):
|
||||
try:
|
||||
confserverlog.info("shutting down")
|
||||
if self.run_async:
|
||||
self.confthread.join()
|
||||
else:
|
||||
await self.app.shutdown()
|
||||
await self.app.shutdown()
|
||||
|
||||
except Exception as e:
|
||||
confserverlog.exception("{}".format(e))
|
||||
|
|
|
|||
|
|
@ -31,7 +31,7 @@ logging.getLogger("hbmqtt.client").setLevel(logging.CRITICAL + 1) # Ignore this
|
|||
|
||||
class MQTTHelperBot:
|
||||
|
||||
Client = MQTTClient()
|
||||
Client = None
|
||||
wait_resp_timeout_seconds = 10
|
||||
expire_msg_seconds = 10
|
||||
|
||||
|
|
@ -39,14 +39,14 @@ class MQTTHelperBot:
|
|||
self.address = address
|
||||
self.client_id = "helper1@bumper/helper1"
|
||||
self.command_responses = []
|
||||
self.helperthread = None
|
||||
|
||||
async def start_helper_bot(self):
|
||||
|
||||
try:
|
||||
self.Client = MQTTClient(
|
||||
client_id=self.client_id, config={"check_hostname": False}
|
||||
)
|
||||
if self.Client is None:
|
||||
self.Client = MQTTClient(
|
||||
client_id=self.client_id, config={"check_hostname": False}
|
||||
)
|
||||
|
||||
await self.Client.connect(
|
||||
"mqtts://{}:{}/".format(self.address[0], self.address[1]),
|
||||
|
|
@ -183,29 +183,30 @@ class MQTTHelperBot:
|
|||
}
|
||||
|
||||
async def send_command(self, cmdjson, requestid):
|
||||
try:
|
||||
ttopic = "iot/p2p/{}/helper1/bumper/helper1/{}/{}/{}/q/{}/{}".format(
|
||||
cmdjson["cmdName"],
|
||||
cmdjson["toId"],
|
||||
cmdjson["toType"],
|
||||
cmdjson["toRes"],
|
||||
requestid,
|
||||
cmdjson["payloadType"],
|
||||
)
|
||||
if not self.Client._handler.writer is None:
|
||||
try:
|
||||
await self.Client.publish(
|
||||
ttopic, str(cmdjson["payload"]).encode(), QOS_0
|
||||
ttopic = "iot/p2p/{}/helper1/bumper/helper1/{}/{}/{}/q/{}/{}".format(
|
||||
cmdjson["cmdName"],
|
||||
cmdjson["toId"],
|
||||
cmdjson["toType"],
|
||||
cmdjson["toRes"],
|
||||
requestid,
|
||||
cmdjson["payloadType"],
|
||||
)
|
||||
try:
|
||||
await self.Client.publish(
|
||||
ttopic, str(cmdjson["payload"]).encode(), QOS_0
|
||||
)
|
||||
except Exception as e:
|
||||
helperbotlog.exception("{}".format(e))
|
||||
|
||||
resp = await self.wait_for_resp(requestid)
|
||||
|
||||
return resp
|
||||
|
||||
except Exception as e:
|
||||
helperbotlog.exception("{}".format(e))
|
||||
|
||||
resp = await self.wait_for_resp(requestid)
|
||||
|
||||
return resp
|
||||
|
||||
except Exception as e:
|
||||
helperbotlog.exception("{}".format(e))
|
||||
return {}
|
||||
return {}
|
||||
|
||||
|
||||
class MQTTServer:
|
||||
|
|
@ -233,7 +234,6 @@ class MQTTServer:
|
|||
|
||||
def __init__(self, address):
|
||||
try:
|
||||
self.mqttserverthread = None
|
||||
self.address = address
|
||||
|
||||
# The below adds a plugin to the hbmqtt.broker.plugins without having to futz with setup.py
|
||||
|
|
@ -355,6 +355,7 @@ class BumperMQTTServer_Plugin:
|
|||
|
||||
async def on_broker_client_connected(self, client_id):
|
||||
try:
|
||||
|
||||
didsplit = str(client_id).split("@")
|
||||
|
||||
bot = bumper.bot_get(didsplit[0])
|
||||
|
|
|
|||
|
|
@ -52,7 +52,7 @@ class XMPPServer:
|
|||
|
||||
def disconnect(self):
|
||||
try:
|
||||
xmppserverlog.debug("waiting for all client threads to exit")
|
||||
xmppserverlog.debug("waiting for all clients to disconnect")
|
||||
for client in self.clients:
|
||||
client._disconnect()
|
||||
|
||||
|
|
@ -722,7 +722,11 @@ class XMPPAsyncClient:
|
|||
).replace("ns0:", ""),
|
||||
)
|
||||
)
|
||||
if 'td="error"' in newdata or 'errs=' in newdata or 'k="DeviceAlert' in newdata:
|
||||
if (
|
||||
'td="error"' in newdata
|
||||
or "errs=" in newdata
|
||||
or 'k="DeviceAlert' in newdata
|
||||
):
|
||||
boterrorlog.error(
|
||||
"Received Error from ({}:{} | {}) - {}".format(
|
||||
self.address[0],
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue