bumper/bumper/xmppserver.py
Brian Martin b50b0df734 WIP - Init
WIP, not working
2019-05-25 08:39:24 -04:00

909 lines
38 KiB
Python
Raw Blame History

This file contains invisible Unicode characters

This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

#!/usr/bin/env python3
from threading import Thread
import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as ET
import base64
import ssl
import bumper
import asyncio, functools
xmppserverlog = logging.getLogger("xmppserver")
class XMPPServer:
server_id = "ecouser.net"
client_id = None
clients = []
exit_flag = False
server_cert = "./certs/cert.pem"
server_key = "./certs/key.pem"
ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
ssl_ctx.load_cert_chain(server_cert, server_key)
def __init__(self, address):
# Initialize bot server
self.address = address
async def async_server(self):
xmppserverlog.info(
"Starting XMPP Server at {}:{}".format(self.address[0], self.address[1])
)
#if self.usessl:
ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
ssl_ctx.load_cert_chain(bumper.server_cert, bumper.server_key)
#server = await asyncio.start_server(
# self.accept_client, self.address[0], self.address[1],ssl=ssl_ctx,
#)
#else:
server = await asyncio.start_server(
self.accept_client, self.address[0], self.address[1]
)
await server.serve_forever()
# self.clients = {} # task -> (reader, writer)
def accept_client(self, 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(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.debug("New Connection from {}:{}".format(clientaddr[0],clientaddr[1]))
task.add_done_callback(functools.partial(client_done, aclient))
except Exception as e:
xmppserverlog.error("{}".format(e))
def disconnect(self):
try:
xmppserverlog.debug("waiting for all client threads to exit")
for client in self.clients:
client._disconnect()
self.exit_flag = True
xmppserverlog.debug("shutting down")
except Exception as e:
xmppserverlog.error("{}".format(e))
class XMPPAsyncClient:
IDLE = 0
CONNECT = 1
INIT = 2
BIND = 3
READY = 4
DISCONNECT = 5
UNKNOWN = 0
BOT = 1
CONTROLLER = 2
_async_task = None
def __init__(self, client_reader, client_writer):
self.type = self.UNKNOWN
self.state = self.IDLE
self.address = client_writer.get_extra_info("peername")
self.client_reader = client_reader
self.client_writer = client_writer
self.clientresource = ""
self.devclass = ""
self.bumper_jid = ""
self.uid = ""
self.log_sent_message = True # Set to true to log sends
self.log_incoming_data = True # Set to true to log sends
xmppserverlog.debug("new client with ip {}".format(self.address))
async def handle_async_client(self):
# xmppserverlog.info('client connected - {}'.format(self.address))
await self._set_state("CONNECT")
while True:
await asyncio.sleep(0.05)
if not self.state == self.DISCONNECT:
data = await self.client_reader.read(4096)
if data is None or data == b'':
xmppserverlog.debug("Received no data")
# exit loop and disconnect
return
else:
await self._parse_data(data)
else:
break
# exit loop and disconnect
return
async def send(self, command):
try:
# if not self.connection._closed:
if self.log_sent_message:
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))
await self._set_state("DISCONNECT")
except ConnectionResetError as e:
#xmppserverlog.debug("{}".format(e))
await self._set_state("DISCONNECT")
except ConnectionAbortedError as e:
#xmppserverlog.debug("{}".format(e))
await self._set_state("DISCONNECT")
except OSError as e:
xmppserverlog.error("{}".format(e))
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _disconnect(self):
try:
bot = bumper.bot_get(self.uid)
if bot:
bumper.bot_set_xmpp(bot["did"], False)
client = bumper.client_get(self.clientresource)
if client:
bumper.client_set_xmpp(client["resource"], False)
self.client_writer.close()
except Exception as e:
xmppserverlog.error("{}".format(e))
async def _tag_strip_uri(self, tag):
try:
if tag[0] == "{":
_, _, tag = tag[1:].partition("}")
return tag
except Exception as e:
xmppserverlog.error("{}".format(e))
async def _set_state(self, state):
try:
new_state = getattr(XMPPAsyncClient, state)
if self.state > new_state:
raise Exception(
"{} illegal state change {}->{}".format(
self.address, self.state, new_state
)
)
xmppserverlog.debug("({}:{} | {}) state: {}".format(self.address[0],self.address[1],self.bumper_jid, state))
self.state = new_state
if new_state == 5:
await self._disconnect()
except Exception as e:
xmppserverlog.error("{}".format(e))
async def _handle_ctl(self, xml, data):
try:
if "roster" in data:
# Return not-implemented for roster
await self.send(
'<iq type="error" id="{}"><error type="cancel" code="501"><feature-not-implemented xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format(
xml.get("id")
)
)
return
if "disco#items" in data:
# Return not-implemented for disco#items
await self.send(
'<iq type="error" id="{}"><error type="cancel" code="501"><feature-not-implemented xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format(
xml.get("id")
))
return
if "disco#info" in data:
# Return not-implemented for disco#info
await self.send(
'<iq type="error" id="{}"><error type="cancel" code="501"><feature-not-implemented xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format(
xml.get("id")
)
)
return
if xml.get("type") == "set":
if (
"com:sf" in data and xml.get("to") == "rl.ecorobot.net"
): # Android bind? Not sure what this does yet.
await self.send(
'<iq id="{}" to="{}@{}/{}" from="rl.ecorobot.net" type="result"/>'.format(
xml.get("id"),
self.uid,
XMPPServer.server_id,
self.clientresource,
)
)
if len(xml[0]) > 0:
ctl = xml[0][0]
if ctl.get("admin") and self.type == self.BOT:
xmppserverlog.debug(
"admin username received from bot: {}".format(ctl.get("admin"))
)
XMPPServer.client_id = ctl.get("admin")
return
# forward
for client in XMPPServer.clients:
if (
client.bumper_jid != self.bumper_jid
and client.state == client.READY
):
ctl_to = xml.get("to")
if not "from" in xml.attrib:
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=")
rxmlstring = rxmlstring.replace("ns0:", "")
rxmlstring = rxmlstring.replace('iq xmlns="com:ctl"', "iq")
rxmlstring = rxmlstring.replace("<query", '<query xmlns="com:ctl"')
if client.type == self.BOT:
if client.uid.lower() in ctl_to.lower():
xmppserverlog.debug(
"Sending ctl to bot: {}".format(rxmlstring)
)
await client.send(rxmlstring)
except Exception as e:
xmppserverlog.error("{}".format(e))
async def _handle_ping(self, xml, data):
try:
if xml.get("to").find("@") == -1: # No to address
# Ping to server - respond
pingresp = '<iq type="result" id="{}" from="{}" />'.format(
xml.get("id"), xml.get("to")
)
# xmppserverlog.debug("Server Ping resp: {}".format(pingresp))
await self.send(pingresp)
else:
pingto = xml.get("to")
pingfrom = self.bumper_jid
if not "from" in xml.attrib:
xml.attrib["from"] = "{}".format(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="urn:xmpp:ping"', "iq")
pingstring = pingstring.replace("<ping", '<ping xmlns="urn:xmpp:ping"')
for client in XMPPServer.clients:
if (
client.bumper_jid != self.bumper_jid
and client.state == client.READY
):
if client.uid.lower() in pingto.lower():
await client.send(pingstring)
except Exception as e:
xmppserverlog.exception("{}".format(e))
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.schedule_ping(time))
async def _handle_result(self, xml, data):
try:
ctl_to = xml.get("to")
if not "from" in xml.attrib:
xml.attrib["from"] = "{}".format(self.bumper_jid)
if (
"errno='103' error='permission denied," in data
): # No permissions, usually if bot was last on Ecovac network
if self.type == self.BOT:
xquery = xml.getchildren()
ctl = xquery[0].getchildren()
ctlerr = ctl[0].attrib["error"]
adminuser = ctlerr.replace("permission denied, please contact ", "")
adminuser = adminuser.replace(" ", "")
if not (
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?
# Add user jid to bot
newuser = ctl_to.split("/")[0]
adduser = '<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="AddUser" id="0000" jid="{}" /></query></iq>'.format(
uuid.uuid4(), adminuser, self.bumper_jid, newuser
)
xmppserverlog.debug("Add User: {}".format(adduser))
await self.send(adduser)
# Add user ACs - Manage users, settings, and clean (full access)
adduseracs = '<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="SetAC" id="1111" jid="{}"><acs><ac name="userman" allow="1"/><ac name="setting" allow="1"/><ac name="clean" allow="1"/></acs></ctl></query></iq>'.format(
uuid.uuid4(), adminuser, self.bumper_jid, newuser
)
xmppserverlog.debug("Add User ACs: {}".format(adduseracs))
await self.send(adduseracs)
# GetUserInfo - Just to confirm it set correctly
await self.send(
'<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="GetUserInfo" id="4444" /><UserInfos/></query></iq>'.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("<query", '<query xmlns="com:ctl"')
if self.type == self.BOT:
if ctl_to == "de.ecorobot.net": # Send to all clients
xmppserverlog.debug(
"Sending to all clients because of de: {}".format(
rxmlstring
)
)
for client in XMPPServer.clients:
await client.send(rxmlstring)
if xml.get("to").find("@") == -1: # No to address
ctl_to = xml.get("to")
else:
ctl_to = "{}@ecouser.net".format(ctl_to.split("@")[0])
for client in XMPPServer.clients:
if (
client.bumper_jid != self.bumper_jid
and client.state == client.READY
):
if not "@" in ctl_to: # No user@, send to all clients?
# TODO: Revisit later, this may be wrong
await client.send(rxmlstring)
elif (
client.uid.lower() in ctl_to.lower()
): # If client matches TO=
xmppserverlog.debug(
"Sending from {} to client {}: {}".format(
self.uid, client.uid, rxmlstring
)
)
await client.send(rxmlstring)
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _handle_connect(self, data, xml=None):
try:
if self.state == self.CONNECT:
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
await self.send(
'<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(
XMPPServer.server_id
)
)
# with STARTTLS
#await self.send('<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns:tls="http://www.ietf.org/rfc/rfc2595.txt" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(XMPPServer.server_id))
await asyncio.sleep(0.25)
#time.sleep(0.25)
# send authentication support for iq-auth (fallback) and SASL
#await self.send(
# '<stream:features><auth xmlns="http://jabber.org/features/iq-auth"/><mechanisms xmlns="urn:ietf:params:xml:ns:xmpp-sasl"><mechanism>PLAIN</mechanism></mechanisms></stream:features>'
#)
#With STARTTLS #https://xmpp.org/rfcs/rfc3920.html
await self.send(
'<stream:features><starttls xmlns="urn:ietf:params:xml:ns:xmpp-tls"><required/></starttls><auth xmlns="http://jabber.org/features/iq-auth"/><mechanisms xmlns="urn:ietf:params:xml:ns:xmpp-sasl"><mechanism>PLAIN</mechanism></mechanisms></stream:features>'
)
# await self.send('<stream:features><auth xmlns="http://jabber.org/features/iq-auth"/></stream:features>')
else:
await self.send("</stream>")
else:
if "jabber:iq:auth" in xml.tag: # Handle iq-auth
await self._handle_iq_auth(xml)
elif (
"urn:ietf:params:xml:ns:xmpp-sasl" in xml.tag
): # Handle SASL Auth
await self._handle_sasl_auth(xml)
else:
xmppserverlog.error("Couldn't handle: {}".format(xml))
elif self.state == self.INIT:
if xml == None:
# Client getting session after authentication
if data.decode("utf-8").find("jabber:client") > -1:
# ack jabbr:client
await self.send(
'<stream:stream xmlns:stream="http://etherx.jabber.org/streams" xmlns="jabber:client" version="1.0" id="1" from="{}">'.format(
XMPPServer.server_id
)
)
await asyncio.sleep(0.25)
#time.sleep(0.25)
# session
await self.send(
'<stream:features><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"/><session xmlns="urn:ietf:params:xml:ns:xmpp-session"/></stream:features>'
)
else: # Handle init bind
if len(xml):
child = await self._tag_strip_uri(xml[0].tag)
else:
child = None
if xml.tag == "iq":
if child == "bind":
await self._handle_bind(xml)
else:
xmppserverlog.error("Couldn't handle: {}".format(xml))
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _handle_starttls(self, data):
try:
peer = self.client_writer.get_extra_info("peername")
xmppserverlog.debug("Upgrading connection with STARTTLS for {}:{}".format(peer[0],peer[1]))
await self.send("<proceed xmlns='urn:ietf:params:xml:ns:xmpp-tls'/>") #send process to client
# After proceed the connection should be upgraded to TLS
loop = asyncio.get_event_loop()
transport = self.client_writer._transport
protocol = self.client_writer.transport.get_protocol()
new_transport = await loop.start_tls(transport , protocol, XMPPServer.ssl_ctx, server_side=True)
#protocol._stream_reader = asyncio.StreamReader(loop=loop)
#protocol._client_connected_cb = do_after_startls()
# protocol.connection_made(new_transport)
#self.client_reader.set_transport(new_transport)
#self.client_writer = transport#protocol._stream_writer
#await loop.start_tls(self.client_writer.transport, self.client_writer._protocol, XMPPServer.ssl_ctx, server_side=True)
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _handle_iq_auth(self, data):
try:
xml = ET.fromstring(data.decode("utf-8"))
ctl = xml[0][0]
xmppserverlog.info("IQ AUTH XML: {}".format(xml))
# Received username and auth tag, send username/password requirement
if (
xml.get("type") == "get"
and "auth}username" in ctl.tag
and self.type == self.UNKNOWN
):
await self.send(
'<iq type="result" id="{}"><query xmlns="jabber:iq:auth"><username/><password/></query></iq>'.format(
xml.get("id")
)
)
# Received username, password, resource - Handle auth here and return pass or fail
if (
xml.get("type") == "set"
and "auth}username" in ctl.tag
and self.type == self.UNKNOWN
):
xmlauth = xml[0].getchildren()
# uid = ""
password = ""
authcode = ""
resource = ""
for aitem in xmlauth:
if "username" in aitem.tag:
self.uid = aitem.text
elif "password" in aitem.tag:
password = aitem.text.split("/")[2]
authcode = password
elif "resource" in aitem.tag:
self.clientresource = aitem.text
resource = self.clientresource
if self.devclass: # if there is a devclass it is a bot
bumper.bot_add("", self.uid, "", resource, "eco-legacy")
xmppserverlog.debug("bot authenticated {}".format(self.uid))
# Client authenticated, move to next state
await self._set_state("INIT")
# Successful auth
await self.send('<iq type="result" id="{}"/>'.format(xml.get("id")))
else:
auth = False
if bumper.check_authcode(self.uid, authcode):
auth = True
elif bumper.use_auth == False:
auth = True
if auth:
bumper.client_add(self.uid, "bumper", self.clientresource)
xmppserverlog.debug("client authenticated {}".format(self.uid))
# Client authenticated, move to next state
await self._set_state("INIT")
# Successful auth
await self.send(
'<iq type="result" id="{}"/>'.format(xml.get("id"))
)
else:
# Failed auth
await self.send(
'<iq type="error" id="{}"><error code="401" type="auth"><not-authorized xmlns="urn:ietf:params:xml:ns:xmpp-stanzas"/></error></iq>'.format(
xml.get("id")
)
)
except ET.ParseError as e:
if "no element found" in e.msg:
xmppserverlog.debug(
"xml parse error - {} - {}".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))
async 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]
authcode = ""
self.uid = username
if len(saslauth) > 1:
resource = saslauth[1]
self.clientresource = resource
elif len(saslauth[0].split("\x00")) > 2:
resource = saslauth[0].split("\x00")[2]
self.clientresource = resource
if len(saslauth) > 2:
authcode = saslauth[2]
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.debug("bot authenticated {}".format(self.uid))
# Send response
await self.send(
'<success xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # Success
# Client authenticated, move to next state
await self._set_state("INIT")
else:
auth = False
if bumper.check_authcode(self.uid, authcode):
auth = True
elif bumper.use_auth == False:
auth = True
if auth:
self.type = self.CONTROLLER
bumper.client_add(self.uid, "bumper", self.clientresource)
xmppserverlog.debug("client authenticated {}".format(self.uid))
# Client authenticated, move to next state
await self._set_state("INIT")
# Send response
await self.send(
'<success xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # Success
else:
# Failed to authenticate
await self.send(
'<response xmlns="urn:ietf:params:xml:ns:xmpp-sasl"/>'
) # Fail
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _handle_bind(self, xml):
try:
bot = bumper.bot_get(self.uid)
if bot:
bumper.bot_set_xmpp(bot["did"], True)
client = bumper.client_get(self.clientresource)
if client:
bumper.client_set_xmpp(client["resource"], True)
clientbindxml = xml.getchildren()
clientresourcexml = clientbindxml[0].getchildren()
if self.devclass: # its a bot
self.name = "XMPP_Client_{}_{}".format(self.uid, self.devclass)
self.bumper_jid = "{}@{}.ecorobot.net/atom".format(
self.uid, self.devclass
)
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
)
elif len(clientresourcexml) > 0:
self.clientresource = clientresourcexml[0].text
self.name = "XMPP_Client_{}".format(self.clientresource)
self.bumper_jid = "{}@{}/{}".format(
self.uid, XMPPServer.server_id, 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.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
)
await self._set_state("BIND")
await self.send(res)
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _handle_session(self, xml):
try:
res = '<iq type="result" id="{}" />'.format(xml.get("id"))
await self._set_state("READY")
await self.send(res)
asyncio.Task(self.schedule_ping(30))
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _handle_presence(self, xml):
try:
if len(xml) and xml[0].tag == "status":
xmppserverlog.debug(
"bot presence {} ".format(ET.tostring(xml, encoding="utf-8").decode("utf-8"))
)
# Most likely a bot, possibly hello world in text
# Send dummy return
await self.send(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
)
# If it is a BOT, send extras
if self.type == self.BOT:
# get device info
await self.send(
'<iq type="set" id="14" to="{}" from="{}"><query xmlns="com:ctl"><ctl td="GetDeviceInfo"/></query></iq>'.format(
self.bumper_jid, XMPPServer.server_id
)
)
else:
xmppserverlog.debug(
"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").decode("utf-8"))
)
# Send dummy return
await self.send(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
)
elif xml.get("type") == "unavailable":
xmppserverlog.debug(
"client presence unavailable (DISCONNECT) - {} ".format(
ET.tostring(xml, encoding="utf-8").decode("utf-8"))
)
await 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
await self.send(
'<presence to="{}"> dummy </presence>'.format(self.bumper_jid)
)
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _parse_data(self, data):
if data.decode("utf-8").startswith(
"<?xml"
): # Strip <?xml and add artificial root
newdata = (
re.sub(r"(<\?xml[^>]+\?>)", r"<root>", data.decode("utf-8")) + "</root>"
)
else:
newdata = "<root>{}</root>".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[0],self.address[1],self.bumper_jid,
str(
ET.tostring(item, encoding="utf-8").decode(
"utf-8"
)
).replace("ns0:", ""),
)
)
await 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
await self._handle_sasl_auth(item)
item.clear()
elif "-tls" in item.tag:
await self._handle_starttls(newdata.encode("utf-8"))
elif "presence" in item.tag:
await 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:", "")
)
)
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 "<stream:stream " in newdata:
if self.state == self.CONNECT or self.state == self.INIT:
await self._handle_connect(newdata.encode("utf-8"))
else:
if not (newdata == "" or newdata == " "):
xmppserverlog.error(
"xml parse error - {} - {}".format(newdata, e)
)
elif "not well-formed (invalid token)" in e.msg:
# If a lone </stream:stream> - client is signalling end of session/disconnect
if not "</stream:stream>" in newdata:
xmppserverlog.error("xml parse error - {} - {}".format(newdata, e))
else:
await self.send("</stream:stream>") # Close stream
else:
if "<stream:stream" in newdata: # Handle start stream and connect
if self.state == self.CONNECT or self.state == self.INIT:
xmppserverlog.debug(
"Handling connect data - {}".format(newdata)
)
await self._handle_connect(newdata.encode("utf-8"))
else:
if not "</stream:stream>" in newdata:
xmppserverlog.error(
"xml parse error - {} - {}".format(newdata, e)
)
else:
await self.send("</stream:stream>") # Close stream
await self._set_state("DISCONNECT")
except Exception as e:
xmppserverlog.exception("{}".format(e))
async def _handle_iq(self, xml, data):
try:
if len(xml):
child = await self._tag_strip_uri(xml[0].tag)
else:
child = None
if xml.tag == "iq":
if child == "bind":
await self._handle_bind(xml)
elif child == "session":
await self._handle_session(xml)
elif child == "ping":
await self._handle_ping(xml, data)
elif child == "query":
if self.type == self.BOT:
await self._handle_result(xml, data)
else:
await self._handle_ctl(xml, data)
elif xml.get("type") == "result":
if self.type == self.BOT:
await self._handle_result(xml, data)
else:
await self._handle_result(xml, data)
elif xml.get("type") == "set":
if self.type == self.BOT:
await self._handle_result(xml, data)
else:
await self._handle_result(xml, data)
except Exception as e:
xmppserverlog.exception("{}".format(e))