From a7ce14eb9c034e0a5d202c883a276b73c550479d Mon Sep 17 00:00:00 2001 From: Brian Martin Date: Wed, 27 Feb 2019 09:05:23 -0500 Subject: [PATCH] additional optimizations and permission handling Optimizations for handling XMPP data and commands Basic permission handling - Add fuid users to bot - Ozmo has ACLs and an "owner", if bot comes from eco to bumper the fuid_user won't have permissions - Bumper will add the fuid_ user to bot ACLs and grant full - This can be used later when bumper adds proper auth to allow/block users --- bumper/confserver.py | 40 ++-- bumper/xmppserver.py | 423 +++++++++++++++++++++++-------------------- 2 files changed, 249 insertions(+), 214 deletions(-) diff --git a/bumper/confserver.py b/bumper/confserver.py index dc6cf22..6cb58e4 100644 --- a/bumper/confserver.py +++ b/bumper/confserver.py @@ -652,22 +652,32 @@ class ConfServer: json_body = json.loads(await request.text()) randomid = "".join(random.sample(string.ascii_letters, 6)) bots = self.bumper_bots.get() - for bot in bots: - if bot.did == json_body["toId"] 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) + if "toId" in json_body: #Its a command + for bot in bots: + if bot.company == 'eco-ng': + if bot.did == json_body["toId"] 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) + ) + return web.json_response(body) + + #No response, send error back + confserverlog.error( + "No bots with DID: {} connected to MQTT".format( + json_body["toId"] ) - return web.json_response(body) - else: - 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) + ) + body = {"id": randomid, "errno": bumper.ERR_COMMON, "ret": "fail"} + 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)) diff --git a/bumper/xmppserver.py b/bumper/xmppserver.py index 19f2a0d..78cf78b 100644 --- a/bumper/xmppserver.py +++ b/bumper/xmppserver.py @@ -11,8 +11,7 @@ xmppserverlog = logging.getLogger("xmppserver") class XMPPServer: - server_id = "bumper" - bot_id = "bumpy" + server_id = "ecouser.net" client_id = None clients = [] exit_flag = False @@ -218,7 +217,7 @@ class Client(threading.Thread): except BrokenPipeError as e: xmppserverlog.error("{}".format(e)) - # self._set_state('DISCONNECT') + self._set_state('DISCONNECT') except ConnectionResetError as e: xmppserverlog.error("{}".format(e)) @@ -282,7 +281,7 @@ class Client(threading.Thread): def _handle_ctl(self, xml, data): try: - if data.decode("utf-8").find("roster") > -1: + if "roster" in data: # Return not-implemented for roster self.send( ''.format( @@ -293,7 +292,7 @@ class Client(threading.Thread): if xml.get("type") == "set": if ( - data.decode("utf-8").find("com:sf") > -1 + "com:sf" in data and xml.get("to") == "rl.ecorobot.net" ): # Android bind? Not sure what this does yet. self.send( @@ -305,11 +304,6 @@ class Client(threading.Thread): ) ) - #else: - # xmppserverlog.debug( - # "Unknown set type: {}".format(data.decode("utf-8")) - # ) - if xml[0][0]: ctl = xml[0][0] if ctl.get("admin") and self.type == self.BOT: @@ -321,12 +315,9 @@ class Client(threading.Thread): # forward for client in XMPPServer.clients: - if client.bumper_jid != self.bumper_jid and client.state == client.READY: - #if client.address != self.address and client.state == client.READY: - ctl_to = xml.get("to") - - #xml.attrib["from"] = self.bumper_jid#.replace("@{}".format(XMPPServer.server_id),"@ecouser.net") - xml.attrib["from"] = "{}@ecouser.net".format(self.uid) + if client.bumper_jid != self.bumper_jid and client.state == client.READY: + ctl_to = xml.get("to") + xml.attrib["from"] = "{}".format(self.bumper_jid) rxmlstring = ET.tostring(xml).decode("utf-8") #clean up string to remove namespaces added by ET rxmlstring = rxmlstring.replace("xmlns:ns0=", "xmlns=") @@ -338,20 +329,7 @@ class Client(threading.Thread): if client.uid.lower() in ctl_to.lower(): xmppserverlog.info("Sending ctl to bot: {}".format(rxmlstring)) client.send(rxmlstring) - #client.send(data.decode("utf-8")) - - # data = data.decode("utf-8") - # id_index = data.find("id") - # if id_index > -1: - # data = ( - # data[:id_index] - # + 'from="' - # + XMPPServer.client_id - # + '" ' - # + data[id_index:] - # ) - # data = data.encode() - # client.send(data.decode("utf-8")) + except Exception as e: xmppserverlog.exception("{}".format(e)) @@ -391,105 +369,136 @@ class Client(threading.Thread): xmppserverlog.exception("{}".format(e)) def _handle_result(self, xml, data): - # forward try: ctl_to = xml.get("to") xml.attrib["from"] = self.bumper_jid - rxmlstring = ET.tostring(xml).decode("utf-8") - #clean up string to remove namespaces added by ET - rxmlstring = rxmlstring.replace("xmlns:ns0=", "xmlns=") - rxmlstring = rxmlstring.replace("ns0:", "") - rxmlstring = rxmlstring.replace('iq xmlns="com:ctl"', "iq") - rxmlstring = rxmlstring.replace(''.format( + uuid.uuid4(), adminuser, self.bumper_jid, newuser) + xmppserverlog.debug("Add User: {}".format(adduser)) + self.send(adduser) + + #Add user ACs - Manage users, settings, and clean (full access) + adduseracs = ''.format( + uuid.uuid4(), adminuser, self.bumper_jid, newuser) + xmppserverlog.debug("Add User ACs: {}".format(adduseracs)) + self.send(adduseracs) + + #GetUserInfo - Just to confirm it set correctly + self.send( + ''.format( + uuid.uuid4(), adminuser, self.bumper_jid) + ) + + else: + rxmlstring = ET.tostring(xml).decode("utf-8") + #clean up string to remove namespaces added by ET + rxmlstring = rxmlstring.replace("xmlns:ns0=", "xmlns=") + rxmlstring = rxmlstring.replace("ns0:", "") + rxmlstring = rxmlstring.replace('iq xmlns="com:ctl"', "iq") + rxmlstring = rxmlstring.replace(' -1: - sc = data.decode("utf-8").find("to=") - ec = data.decode("utf-8").find(".ecorobot.net") - if ec > -1: - self.devclass = data.decode("utf-8")[sc+4:ec] - # ack jabbr:client - # no STARTTLS - self.send( - ''.format( - XMPPServer.server_id + if xml == None: + # Client first connecting, send our features + if data.decode("utf-8").find("jabber:client") > -1: + sc = data.decode("utf-8").find("to=") + ec = data.decode("utf-8").find(".ecorobot.net") + if ec > -1: + self.devclass = data.decode("utf-8")[sc+4:ec] + # ack jabbr:client + # no STARTTLS + self.send( + ''.format( + XMPPServer.server_id + ) ) - ) - # with STARTTLS - # self.send(''.format(XMPPServer.server_id)) - time.sleep(0.25) - # send authentication support for iq-auth (fallback) and SASL - self.send( - 'PLAIN' - ) - # self.send('') + # with STARTTLS + # self.send(''.format(XMPPServer.server_id)) + time.sleep(0.25) + # send authentication support for iq-auth (fallback) and SASL + self.send( + 'PLAIN' + ) + # self.send('') - elif data.decode("utf-8").find("jabber:iq:auth") > -1: # Handle iq-auth - self._handle_iq_auth(data) - - elif ( - data.decode("utf-8").find("urn:ietf:params:xml:ns:xmpp-sasl") > -1 - ): # Handle SASL auth - self._handle_sasl_auth(data) + else: + self.send("") + + else: + if "jabber:iq:auth" in xml.tag: # Handle iq-auth + self._handle_iq_auth(xml) + elif "urn:ietf:params:xml:ns:xmpp-sasl" in xml.tag: #Handle SASL Auth + self._handle_sasl_auth(xml) + else: + xmppserverlog.error("Couldn't handle: {}".format(xml)) elif self.state == self.INIT: - # Client getting session after authentication - if data.decode("utf-8").find("jabber:client") > -1: - # ack jabbr:client - self.send( - ''.format( - XMPPServer.server_id + if xml == None: + # Client getting session after authentication + if data.decode("utf-8").find("jabber:client") > -1: + # ack jabbr:client + self.send( + ''.format( + XMPPServer.server_id + ) + ) + time.sleep(0.25) + # session + self.send( + '' ) - ) - time.sleep(0.25) - # session - self.send( - '' - ) - else: # Handle init bind - xml = ET.fromstring(data.decode("utf-8")) + else: # Handle init bind if len(xml): - child = self._tag_strip_uri(xml[0].tag) + child = self._tag_strip_uri(xml[0].tag) else: child = None if xml.tag == "iq": if child == "bind": - self._handle_bind(xml) + self._handle_bind(xml) + else: + xmppserverlog.error("Couldn't handle: {}".format(xml)) + + except Exception as e: xmppserverlog.exception("{}".format(e)) @@ -573,7 +582,7 @@ class Client(threading.Thread): except ET.ParseError as e: if "no element found" in e.msg: xmppserverlog.debug( - "xml parse error - {} - {} - this is common with ecovac protocol".format( + "xml parse error - {} - {}".format( data.decode("utf-8"), e ) ) @@ -589,10 +598,10 @@ class Client(threading.Thread): except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_sasl_auth(self, data): - try: - xml = ET.fromstring(data.decode("utf-8")) - saslauth = base64.b64decode(xml.text).decode("utf-8").split("/") + def _handle_sasl_auth(self, xml): + try: + + saslauth = base64.b64decode(xml.text).decode("utf-8").split("/") username = saslauth[0] username = saslauth[0].split("\x00")[1] self.uid = username @@ -609,6 +618,7 @@ class Client(threading.Thread): if not self.uid.startswith("fuid"): # Need sample data to see details here bumper.add_bot(self.uid, self.uid, self.devclass, "atom","eco-legacy") + self.type = self.BOT xmppserverlog.info("bot authenticated {}".format(self.uid)) # Send response self.send( @@ -626,6 +636,7 @@ class Client(threading.Thread): auth = True if auth: + self.type = self.CONTROLLER bumper.add_client(self.uid, "bumper", self.clientresource) xmppserverlog.debug("client authenticated {}".format(self.uid)) @@ -643,22 +654,6 @@ class Client(threading.Thread): '' ) # Fail - except ET.ParseError as e: - if "no element found" in e.msg: - xmppserverlog.debug( - "xml parse error - {} - {} - this is common with ecovac protocol".format( - data.decode("utf-8"), e - ) - ) - elif "not well-formed (invalid token)" in e.msg: - xmppserverlog.debug( - "xml parse error - {} - {}".format(data.decode("utf-8"), e) - ) - else: - xmppserverlog.debug( - "xml parse error - {} - {}".format(data.decode("utf-8"), e) - ) - except Exception as e: xmppserverlog.exception("{}".format(e)) @@ -725,103 +720,148 @@ class Client(threading.Thread): def _handle_presence(self, xml): try: + if len(xml) and xml[0].tag == "status": - # bot announcing arrival - self.type = self.BOT xmppserverlog.debug( - "{} type set to BOT (based on presence tag)".format(self.address) + "bot presence {} ".format(ET.tostring(xml, encoding="utf-8")) ) - # send a command from an unknown user - the response will contain the correct admin username + #Most likely a bot, possibly hello world in text + + #Send dummy return self.send( ' dummy '.format( self.bumper_jid ) ) - # self.send( - # ''.format( - # self.uid, self.devclass, XMPPServer.server_id - # ) - # ) - - self.send( - ''.format( - uuid.uuid4(), "unknown@ecouser.net", XMPPServer.server_id - ) - ) - else: - self.type = self.CONTROLLER + #If it is a BOT, send extras + if self.type == self.BOT: + #get device info + self.send( + ''.format( + self.bumper_jid, XMPPServer.server_id + ) + ) + + else: xmppserverlog.debug( - "{} type set to CONTROLLER (based on presence tag)".format( - self.address + "client presence - {} ".format(ET.tostring(xml, encoding="utf-8")) + ) + + if xml.get("type") == "available": + xmppserverlog.debug( + "client presence available - {} ".format(ET.tostring(xml, encoding="utf-8")) ) - ) - self.send( - ' dummy '.format( - self.uid, XMPPServer.server_id, self.clientresource + #Send dummy return + self.send( + ' dummy '.format( + self.bumper_jid + ) ) - ) + elif xml.get("type") == "unavailable": + xmppserverlog.debug( + "client presence unavailable (DISCONNECT) - {} ".format(ET.tostring(xml, encoding="utf-8")) + ) + + self._set_state("DISCONNECT") + else: + #Sometimes the android app sends these + xmppserverlog.debug( + "client presence (UNKNOWN) - {} ".format(ET.tostring(xml, encoding="utf-8")) + ) + #Send dummy return + self.send( + ' dummy '.format( + self.bumper_jid + ) + ) + except Exception as e: xmppserverlog.exception("{}".format(e)) - def _parse_data(self, data): - if self.log_incoming_data: - xmppserverlog.debug( - "from {} - {}".format(self.address, data.decode("utf-8")) - ) + def _parse_data(self, data): - try: - xml = ET.fromstring(data.decode("utf-8")) - self._handle_xml(xml, data) + + if data.decode("utf-8").startswith("]+\?>)", r"",data.decode("utf-8")) + "" + + else: + newdata = "{}".format(data.decode("utf-8")) #Add artificial root + + try: + root = ET.fromstring(newdata) + for item in root.iter(): + if item.tag != "root": + if item.tag == "iq": + if self.log_incoming_data: + xmppserverlog.debug( + "from {} - {}".format(self.address, str(ET.tostring(item, encoding="utf-8").decode("utf-8")).replace("ns0:","")) + ) + self._handle_iq(item, newdata) + item.clear() + elif "auth" in item.tag: + if "urn:ietf:params:xml:ns:xmpp-sasl" in item.tag: #SASL Auth + self._handle_sasl_auth(item) + item.clear() + + elif "presence" in item.tag: + self._handle_presence(item) + item.clear() + + else: + if self.log_incoming_data: + xmppserverlog.debug( + "Unparsed Item - {}".format(str(ET.tostring(item, encoding="utf-8").decode("utf-8")).replace("ns0:","")) + ) + print("e") + except ET.ParseError as e: if ( "no element found" in e.msg ): # Element not closed or not all bytes received # Happens wth connect stream often - if " - client is signalling end of session/disconnect - if not "" in data.decode("utf-8"): + if not "" in newdata: xmppserverlog.error( - "xml parse error - {} - {}".format(data.decode("utf-8"), e) + "xml parse error - {} - {}".format(newdata, e) ) else: self.send("") # Close stream - - elif ( - "junk after document element" in e.msg - ): # More than one xml doc in data - # try to split it - data0 = data.decode("utf-8") - data1 = data0[e.position[1] :] - data0 = data0[: e.position[1]] - # xmppserverlog.debug('xml parse error - {} - {} - split0: {} - split1: {}'.format(data.decode('utf-8'), e, data0, data1)) - self._parse_data(data0.encode("utf-8")) - self._parse_data(data1.encode("utf-8")) - + else: - xmppserverlog.debug( - "xml parse error - {} - {}".format(data.decode("utf-8"), e) - ) + if "" in newdata: + xmppserverlog.error( + "xml parse error - {} - {}".format(newdata, e) + ) + else: + self.send("") # Close stream + self._set_state("DISCONNECT") + except Exception as e: xmppserverlog.exception("{}".format(e)) - def _handle_xml(self, xml, data): + def _handle_iq(self, xml, data): try: - if self.state == self.CONNECT or self.state == self.INIT: - self._handle_connect(data) - if len(xml): child = self._tag_strip_uri(xml[0].tag) else: @@ -848,10 +888,7 @@ class Client(threading.Thread): if self.type == self.BOT: self._handle_result(xml, data) else: - self._handle_result(xml, data) - - if xml.tag == "presence": - self._handle_presence(xml) + self._handle_result(xml, data) except Exception as e: xmppserverlog.exception("{}".format(e)) @@ -859,16 +896,14 @@ class Client(threading.Thread): def run(self): # xmppserverlog.info('client connected - {}'.format(self.address)) self._set_state("CONNECT") - while not self.state == self.DISCONNECT and not self.connection._closed: data = b"" - - time.sleep(0.2) + time.sleep(0.1) if not self.connection._closed: try: - #data = self.connection.recv(4096) - data = self.connection.recv(8192) - + data = self.connection.recv(4096) + if data != b"": + self._parse_data(data) except ConnectionResetError as e: xmppserverlog.error("{}".format(e)) except OSError as e: @@ -876,15 +911,5 @@ class Client(threading.Thread): except Exception as e: xmppserverlog.exception("{}".format(e)) - if data != b"": - splitdata = data.decode("utf-8") - splitdata = splitdata.split(" 1: - for s in splitdata: - if not s.startswith("