From d316862120b3accdb3b1e5b3d56a138313c75748 Mon Sep 17 00:00:00 2001 From: Brian Martin Date: Wed, 20 Feb 2019 00:00:06 -0500 Subject: [PATCH] Handle connected states --- bumper/__init__.py | 4 ++ bumper/confserver.py | 9 ++-- bumper/mqttserver.py | 108 +++++++++++++++++-------------------------- bumper/xmppserver.py | 47 ++++++++++++++++--- 4 files changed, 92 insertions(+), 76 deletions(-) diff --git a/bumper/__init__.py b/bumper/__init__.py index 9b84d69..45691ae 100644 --- a/bumper/__init__.py +++ b/bumper/__init__.py @@ -86,6 +86,8 @@ class VacBotDevice(object): self.name = name self.nick = nick self.resource = resource + self.mqtt_connection = False + self.xmpp_connection = False def asdict(self): return {"class": self.vac_bot_device_class, "company": self.company, @@ -96,6 +98,8 @@ class VacBotClient(object): self.userid = userid self.realm = realm self.resource = token + self.mqtt_connection = False + self.xmpp_connection = False def asdict(self): return {"userid": self.userid,"realm": self.realm,"resource": self.resource} diff --git a/bumper/confserver.py b/bumper/confserver.py index c0e9cb1..87726cb 100644 --- a/bumper/confserver.py +++ b/bumper/confserver.py @@ -356,8 +356,11 @@ class ConfServer(): elif todo == 'GetDeviceList': active_bots = self.bumper_bots.get() + bot_list = [] + for bot in active_bots: + bot_list.append(bot.asdict()) body = { - "devices": active_bots, + "devices": bot_list, "result": "ok", "todo": "result" } @@ -365,8 +368,8 @@ class ConfServer(): elif todo == 'SetDeviceNick': bots = self.bumper_bots.get() for bot in bots: - if postbody['did'] == bot['did']: - bot['nick'] = postbody['nick'] + if postbody['did'] == bot.did: + bot.nick = postbody['nick'] self.bumper_bots.set(bots) body = { "result": "ok", diff --git a/bumper/mqttserver.py b/bumper/mqttserver.py index b5213e9..d3246bf 100644 --- a/bumper/mqttserver.py +++ b/bumper/mqttserver.py @@ -225,7 +225,6 @@ class MQTTServer(): mqttserverlog.exception('{}'.format(e)) def run(self, run_async=False,): - if run_async: sloop = asyncio.new_event_loop() mqttserverlog.debug("Starting MQTTServer Thread: 1") @@ -236,7 +235,6 @@ class MQTTServer(): else: self.run_server() - def run_server(self, loop): logging.info("Starting MQTT Server at {}".format(self.address)) @@ -292,9 +290,16 @@ class BumperMQTTServer_Plugin: newbot.name = username newbot.vac_bot_device_class = tmpbotdetail[0] newbot.resource = tmpbotdetail[1] - bumper_bots.append(newbot.asdict()) - mqttserverlog.info("new bot authenticated {}".format(newbot.asdict())) - self.bumper_config['bumper_bots'].set(bumper_bots) + existingbot = False + for bot in bumper_bots: + if bot.did == newbot.did: + existingbot = True + + if existingbot == False: + bumper_bots.append(newbot) + mqttserverlog.info("new bot authenticated {}".format(newbot.name)) + self.bumper_config['bumper_bots'].set(bumper_bots) + authenticated = True else: @@ -307,9 +312,16 @@ class BumperMQTTServer_Plugin: newclient.userid = didsplit[0] newclient.realm = tmpclientdetail[0] newclient.resource = tmpclientdetail[1] - bumper_clients.append(newclient.asdict()) - mqttserverlog.info("new client authenticated {}".format(newclient.asdict())) - self.bumper_config['bumper_clients'].set(bumper_clients) + existingclient = False + for client in bumper_clients: + if client.userid == newclient.userid: + existingclient = True + + if existingclient == False: + bumper_clients.append(newclient) + mqttserverlog.info("new client authenticated {}".format(newclient.userid)) + self.bumper_config['bumper_clients'].set(bumper_clients) + authenticated = True else: @@ -323,81 +335,45 @@ class BumperMQTTServer_Plugin: async def on_broker_client_connected(self, client_id): try: - #mqttserverlog.debug('%s connected' % client_id) bumper_users = self.bumper_config['bumper_users'].get() - bumper_bots = self.bumper_config['bumper_bots'].get() - bumper_clients = self.bumper_config['bumper_clients'].get() + bumper_bots = self.bumper_config['bumper_bots'].get() + bumper_clients = self.bumper_config['bumper_clients'].get() 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")): - tmpbotdetail = str(didsplit[1]).split("/") - newbot = bumper.VacBotDevice() - newbot.did = didsplit[0] - newbot.vac_bot_device_class = tmpbotdetail[0] - newbot.resource = tmpbotdetail[1] - #botactive = False - #for bot in bumper_bots: - # if bot['did'] == newbot.did: - # botactive = True - - #if botactive == False: - # bumper_bots.append(newbot.asdict()) - mqttserverlog.debug("bot connected {}".format(newbot.did)) - - #self.bumper_config['bumper_bots'].set(bumper_bots) - else: - tmpclientdetail = str(didsplit[1]).split("/") - newclient = bumper.VacBotClient() - newclient.userid = didsplit[0] - newclient.realm = tmpclientdetail[0] - newclient.resource = tmpclientdetail[1] - - #clientactive = False - #for client in bumper_clients: - # if client['userid'] == newuser.userid: - # clientactive = True - - #if clientactive == False and newuser.userid != 'helper1': - #bumper_clients.append(newuser.asdict()) - mqttserverlog.debug("client connected {}".format(newclient.userid)) - - #self.bumper['bumper_clients'].set(bumper_clients) + for bot in bumper_bots: + if didsplit[0] == bot.did: + bot.mqtt_connection = True + mqttserverlog.info("bot connected {}".format(bot.did)) + self.bumper_config['bumper_bots'].set(bumper_bots) - #mqttserverlog.debug('Connected Bots: %s' %self.clients['connected_bots'].get()) - #mqttserverlog.debug('Connected Clients: %s' %self.clients['connected_clients'].get()) + for client in bumper_clients: + if didsplit[0] == client.userid and client.userid != 'helper1': + client.mqtt_connection = True + mqttserverlog.info("client connected {}".format(client.userid)) + self.bumper_config['bumper_clients'].set(bumper_clients) except Exception as e: mqttserverlog.exception('{}'.format(e)) async def on_broker_client_disconnected(self, client_id): - try: - #mqttserverlog.debug('%s disconnected' % client_id) + try: bumper_users = self.bumper_config['bumper_users'].get() bumper_bots = self.bumper_config['bumper_bots'].get() bumper_clients = self.bumper_config['bumper_clients'].get() - #remove_clients = self.clients['remove_clients'].get() didsplit = str(client_id).split("@") - #If the did is in the list, remove it + for bot in bumper_bots: - if didsplit[0] == bot['did']: - mqttserverlog.info("bot disconnected {}".format(bot['did'])) - #bumper_bots.remove(bot) - #remove_clients.append(bot['did']) - #self.clients['bumper_bots'].set(bumper_bots) + if didsplit[0] == bot.did: + bot.mqtt_connection = False + mqttserverlog.info("bot disconnected {}".format(bot.did)) + self.bumper_config['bumper_bots'].set(bumper_bots) for client in bumper_clients: - if didsplit[0] == client['userid'] and client['userid'] != 'helper1': - mqttserverlog.info("client disconnected {}".format(client['userid'])) - #bumper_clients.remove(client) - #remove_clients.append(client['userid']) - #self.bumper_config['bumper_clients'].set(bumper_clients) - - #self.clients['remove_clients'].set(remove_clients) - - #mqttserverlog.debug('Connected Bots: %s' %self.clients['connected_bots'].get()) - #mqttserverlog.debug('Connected Clients: %s' %self.clients['connected_clients'].get()) + if didsplit[0] == client.userid and client.userid != 'helper1': + client.mqtt_connection = False + mqttserverlog.info("client disconnected {}".format(client.userid)) + self.bumper_config['bumper_clients'].set(bumper_clients) except Exception as e: mqttserverlog.exception('{}'.format(e)) \ No newline at end of file diff --git a/bumper/xmppserver.py b/bumper/xmppserver.py index 3740c20..552da52 100644 --- a/bumper/xmppserver.py +++ b/bumper/xmppserver.py @@ -67,7 +67,7 @@ class XMPPServer(): xmppserverlog.debug('starting new client with ip {}'.format(client_address[0])) thread_id = uuid.uuid4() - client = Client(thread_id, connection, client_address) + client = Client(thread_id, connection, client_address, self.bumper_users, self.bumper_bots, self.bumper_clients) client.setDaemon(True) client.start() self.clients.append(client) @@ -138,7 +138,7 @@ class Client(threading.Thread): BOT = 1 CONTROLLER = 2 - def __init__(self, thread_id, connection, client_address): + def __init__(self, thread_id, connection, client_address,bumper_users=contextvars.ContextVar, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar): threading.Thread.__init__(self) self.id = thread_id self.name = "XMPP_Client_{}".format(client_address[0]) @@ -150,6 +150,9 @@ class Client(threading.Thread): self.uid = "" self.log_sent_message = False #Set to true to log sends self.log_incoming_data = True #Set to true to log sends + self.bumper_users = bumper_users + self.bumper_bots = bumper_bots + self.bumper_clients = bumper_clients xmppserverlog.debug('new client thread init for client with ip {}'.format(self.address)) @@ -177,7 +180,22 @@ class Client(threading.Thread): def _disconnect(self): try: - xmppserverlog.debug('client {} with resource {} disconnecting'.format(self.address, self.clientresource)) + bumper_bots = self.bumper_bots.get() + bumper_clients = self.bumper_clients.get() + for bot in bumper_bots: + if self.uid == bot.did: + bot.xmpp_connection = False + xmppserverlog.info("bot disconnected {}".format(bot.did)) + + self.bumper_bots.set(bumper_bots) + + for client in bumper_clients: + if self.uid == client.userid and client.userid != 'helper1': + client.xmpp_connection = False + xmppserverlog.info("client disconnected {}".format(client.userid)) + + self.bumper_clients.set(bumper_clients) + #xmppserverlog.debug('client {} with resource {} disconnecting'.format(self.address, self.clientresource)) self.connection.close() except Exception as e: @@ -343,6 +361,22 @@ class Client(threading.Thread): self.clientresource = aitem.text if bumper.check_authcode(self.uid, password): + bumper_bots = self.bumper_bots.get() + bumper_clients = self.bumper_clients.get() + for bot in bumper_bots: + if self.uid == bot.did: + bot.xmpp_connection = True + xmppserverlog.info("bot connected {}".format(bot.did)) + + self.bumper_config['bumper_bots'].set(bumper_bots) + + for client in bumper_clients: + if self.uid == client.userid and client.userid != 'helper1': + client.xmpp_connection = True + xmppserverlog.info("client connected {}".format(client.userid)) + + self.bumper_config['bumper_clients'].set(bumper_clients) + #Client authenticated, move to next state self._set_state('INIT') @@ -351,9 +385,7 @@ class Client(threading.Thread): else: #Failed auth - self.send(''.format(xml.get('id'))) - - + self.send(''.format(xml.get('id'))) except ET.ParseError as e: if "no element found" in e.msg: @@ -377,7 +409,8 @@ class Client(threading.Thread): self.clientresource = resource authcode = saslauth[2] - if bumper.check_authcode(self.uid, authcode): + if bumper.check_authcode(self.uid, authcode): + #Send response self.send('') #Success