Manageable logging
Manageable logging
This commit is contained in:
parent
5ff859a5ce
commit
a1b1b1a23b
4 changed files with 184 additions and 146 deletions
89
bumper.py
89
bumper.py
|
|
@ -6,44 +6,65 @@ import sys, socket
|
||||||
import time
|
import time
|
||||||
import platform
|
import platform
|
||||||
|
|
||||||
|
bumperlog = logging.getLogger("bumper")
|
||||||
|
def main():
|
||||||
|
args = sys.argv
|
||||||
|
if len(args) > 0:
|
||||||
|
if '--debug' in args:
|
||||||
|
logging.basicConfig(level=logging.DEBUG,
|
||||||
|
format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s")
|
||||||
|
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")
|
||||||
|
|
||||||
args = sys.argv
|
# A default bot could be set here to automatically add it as available
|
||||||
if len(args) > 0:
|
# dbot = bumper.VacBotDevice("did", "class", "resource", "name","nick" )
|
||||||
if '--debug' in args:
|
# bclient = bumper.bumper_bots_var
|
||||||
logging.basicConfig(level=logging.DEBUG,
|
# bclienttemp = bclient.get()
|
||||||
format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s")
|
# bclienttemp.append(dbot.asdict())
|
||||||
|
# bclient.set(bclienttemp)
|
||||||
|
|
||||||
|
if platform.system() == "Darwin":
|
||||||
|
listen_host = "0.0.0.0"
|
||||||
else:
|
else:
|
||||||
logging.basicConfig(level=logging.INFO,
|
listen_host = socket.gethostbyname(socket.gethostname())
|
||||||
format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(message)s")
|
#listen_host = "localhost" #Try this if the above doesn't work
|
||||||
#format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(module)s :: %(funcName)s :: %(lineno)d :: %(message)s")
|
|
||||||
|
|
||||||
# A default bot could be set here to automatically add it as available
|
conf_address_443 = (listen_host, 443)
|
||||||
# dbot = bumper.VacBotDevice("did", "class", "resource", "name","nick" )
|
conf_address_8007 = (listen_host, 8007)
|
||||||
# bclient = bumper.bumper_bots_var
|
xmpp_address = (listen_host, 5223)
|
||||||
# bclienttemp = bclient.get()
|
mqtt_address = (listen_host, 8883)
|
||||||
# bclienttemp.append(dbot.asdict())
|
|
||||||
# bclient.set(bclienttemp)
|
|
||||||
|
|
||||||
if platform.system() == "Darwin":
|
# start mqtt server on port 8883 (async)
|
||||||
listen_host = "0.0.0.0"
|
startmqttserver = "Starting MQTT Server at {}".format(mqtt_address)
|
||||||
else:
|
bumperlog.info("{}".format(startmqttserver))
|
||||||
listen_host = socket.gethostbyname(socket.gethostname())
|
print("{}".format(startmqttserver))
|
||||||
#listen_host = "localhost" #Try this if the above doesn't work
|
mqtt_server = bumper.MQTTServer(mqtt_address, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var)
|
||||||
|
time.sleep(1.5) #Wait for broker startup
|
||||||
|
|
||||||
conf_address_443 = (listen_host, 443)
|
# start mqtt_helperbot (async)
|
||||||
conf_address_8007 = (listen_host, 8007)
|
bumperlog.info("Starting MQTT HelperBot")
|
||||||
xmpp_address = (listen_host, 5223)
|
print("Starting MQTT HelperBot")
|
||||||
mqtt_address = (listen_host, 8883)
|
mqtt_helperbot = bumper.MQTTHelperBot(mqtt_address, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var)
|
||||||
|
|
||||||
# start mqtt server on port 8883 (async)
|
# start conf server on port 443 (async) - Used for most https calls
|
||||||
mqtt_server = bumper.MQTTServer(mqtt_address, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var)
|
startconf443 = "Starting Main ConfServer at {}".format(conf_address_443)
|
||||||
time.sleep(1.5) #Wait for broker startup
|
bumperlog.info("{}".format(startconf443))
|
||||||
# start mqtt_helperbot (async)
|
print("{}".format(startconf443))
|
||||||
mqtt_helperbot = bumper.MQTTHelperBot(mqtt_address, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var)
|
conf_server = bumper.ConfServer(conf_address_443, usessl=True, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var, helperbot=mqtt_helperbot)
|
||||||
# start conf server on port 443 (async) - Used for most https calls
|
|
||||||
conf_server = bumper.ConfServer(conf_address_443, usessl=True, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var, helperbot=mqtt_helperbot)
|
|
||||||
# start conf server on port 8007 (async) - Used for a load balancer request
|
|
||||||
conf_server_2 = bumper.ConfServer(conf_address_8007, usessl=False, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var, helperbot=mqtt_helperbot)
|
|
||||||
# start xmpp server on port 5223 (sync)
|
|
||||||
xmpp_server = bumper.XMPPServer(xmpp_address)
|
|
||||||
|
|
||||||
|
# start conf server on port 8007 (async) - Used for a load balancer request
|
||||||
|
startconf8007 = "Starting LoadBalancer ConfServer at {}".format(conf_address_8007)
|
||||||
|
bumperlog.info("{}".format(startconf8007))
|
||||||
|
print("{}".format(startconf8007))
|
||||||
|
conf_server_2 = bumper.ConfServer(conf_address_8007, usessl=False, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var, helperbot=mqtt_helperbot)
|
||||||
|
|
||||||
|
# start xmpp server on port 5223 (sync)
|
||||||
|
startxmpp = "Starting XMPP Server at {}".format(xmpp_address)
|
||||||
|
bumperlog.info("{}".format(startxmpp))
|
||||||
|
print("{}".format(startxmpp))
|
||||||
|
xmpp_server = bumper.XMPPServer(xmpp_address)
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
main()
|
||||||
|
|
@ -11,6 +11,23 @@ import asyncio
|
||||||
import contextvars
|
import contextvars
|
||||||
from aiohttp import web
|
from aiohttp import web
|
||||||
|
|
||||||
|
class aiohttp_filter(logging.Filter):
|
||||||
|
|
||||||
|
def filter(self, record):
|
||||||
|
if record.name == "aiohttp.access" and record.levelno == 20: #Filters aiohttp.access log to switch it from INFO to DEBUG
|
||||||
|
record.levelno = 10
|
||||||
|
record.levelname = "DEBUG"
|
||||||
|
|
||||||
|
if record.levelno == 10 and logging.getLogger("confserver").getEffectiveLevel() == 10:
|
||||||
|
return True
|
||||||
|
else:
|
||||||
|
return False
|
||||||
|
|
||||||
|
confserverlog = logging.getLogger("confserver")
|
||||||
|
|
||||||
|
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) #Ignore this logger
|
||||||
|
logging.getLogger("aiohttp.access").addFilter(aiohttp_filter())
|
||||||
|
|
||||||
class ConfServer():
|
class ConfServer():
|
||||||
bumper_clients = contextvars.ContextVar
|
bumper_clients = contextvars.ContextVar
|
||||||
bumper_bots = contextvars.ContextVar
|
bumper_bots = contextvars.ContextVar
|
||||||
|
|
@ -26,7 +43,7 @@ class ConfServer():
|
||||||
|
|
||||||
try:
|
try:
|
||||||
if run_async:
|
if run_async:
|
||||||
logging.debug("Starting ConfServer Thread: 1")
|
confserverlog.debug("Starting ConfServer Thread: 1")
|
||||||
confserver = Thread(name="ConfServer_Thread",target=self.run_server)
|
confserver = Thread(name="ConfServer_Thread",target=self.run_server)
|
||||||
self.server = confserver
|
self.server = confserver
|
||||||
confserver.setDaemon(True)
|
confserver.setDaemon(True)
|
||||||
|
|
@ -38,7 +55,7 @@ class ConfServer():
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt:
|
||||||
self.disconnect()
|
self.disconnect()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
|
|
||||||
def run_server(self):
|
def run_server(self):
|
||||||
|
|
@ -51,7 +68,7 @@ class ConfServer():
|
||||||
loop.run_until_complete(self.start_server())
|
loop.run_until_complete(self.start_server())
|
||||||
loop.run_forever()
|
loop.run_forever()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
|
|
||||||
async def start_server(self):
|
async def start_server(self):
|
||||||
|
|
@ -76,7 +93,7 @@ class ConfServer():
|
||||||
])
|
])
|
||||||
|
|
||||||
|
|
||||||
runner = web.AppRunner(app)#, access_log=None) #access_log=None so the output isn't nuts
|
runner = web.AppRunner(app)
|
||||||
await runner.setup()
|
await runner.setup()
|
||||||
|
|
||||||
if self.usessl:
|
if self.usessl:
|
||||||
|
|
@ -91,11 +108,11 @@ class ConfServer():
|
||||||
|
|
||||||
except PermissionError as e:
|
except PermissionError as e:
|
||||||
if "bind" in e.strerror:
|
if "bind" in e.strerror:
|
||||||
logging.exception("Error binding confserver, exiting. Try using a different hostname or IP.\r\n {}".format(e))
|
confserverlog.exception("Error binding confserver, exiting. Try using a different hostname or IP - {}".format(e))
|
||||||
exit(1)
|
exit(1)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
exit(1)
|
exit(1)
|
||||||
|
|
||||||
async def handle_login(self, request):
|
async def handle_login(self, request):
|
||||||
|
|
@ -118,7 +135,7 @@ class ConfServer():
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_checkLogin(self, request):
|
async def handle_checkLogin(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -141,7 +158,7 @@ class ConfServer():
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_logout(self, request):
|
async def handle_logout(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -151,7 +168,7 @@ class ConfServer():
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_getAuthCode(self, request):
|
async def handle_getAuthCode(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -169,7 +186,7 @@ class ConfServer():
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_checkVersion(self, request):
|
async def handle_checkVersion(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -191,7 +208,7 @@ class ConfServer():
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_checkAgreement(self, request):
|
async def handle_checkAgreement(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -206,7 +223,7 @@ class ConfServer():
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_homePageAlert(self, request):
|
async def handle_homePageAlert(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -229,7 +246,7 @@ class ConfServer():
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_getProductIotMap(self, request):
|
async def handle_getProductIotMap(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -238,7 +255,7 @@ class ConfServer():
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_usersapi(self, request):
|
async def handle_usersapi(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -251,8 +268,6 @@ class ConfServer():
|
||||||
else:
|
else:
|
||||||
postbody = json.loads(await request.text())
|
postbody = json.loads(await request.text())
|
||||||
|
|
||||||
logging.debug(postbody)
|
|
||||||
|
|
||||||
todo = postbody['todo']
|
todo = postbody['todo']
|
||||||
if todo == 'FindBest':
|
if todo == 'FindBest':
|
||||||
service = postbody['service']
|
service = postbody['service']
|
||||||
|
|
@ -276,10 +291,11 @@ class ConfServer():
|
||||||
"todo": "result"
|
"todo": "result"
|
||||||
}
|
}
|
||||||
|
|
||||||
|
confserverlog.debug("\r\n POST: {} \r\n Response: {}".format(postbody,body))
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_lookup(self, request):
|
async def handle_lookup(self, request):
|
||||||
try:
|
try:
|
||||||
|
|
@ -292,7 +308,7 @@ class ConfServer():
|
||||||
else:
|
else:
|
||||||
postbody = json.loads(await request.text())
|
postbody = json.loads(await request.text())
|
||||||
|
|
||||||
logging.debug(postbody)
|
confserverlog.debug(postbody)
|
||||||
|
|
||||||
todo = postbody['todo']
|
todo = postbody['todo']
|
||||||
if todo == 'FindBest':
|
if todo == 'FindBest':
|
||||||
|
|
@ -302,33 +318,32 @@ class ConfServer():
|
||||||
elif service == 'EcoUpdate':
|
elif service == 'EcoUpdate':
|
||||||
body = {"result":"ok","ip":"47.88.66.164","port":8005}
|
body = {"result":"ok","ip":"47.88.66.164","port":8005}
|
||||||
|
|
||||||
|
confserverlog.debug("\r\n POST: {} \r\n Response: {}".format(postbody,body))
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def handle_devmanager_botcommand(self, request):
|
async def handle_devmanager_botcommand(self, request):
|
||||||
try:
|
try:
|
||||||
json_body = json.loads(await request.text())
|
json_body = json.loads(await request.text())
|
||||||
logging.info("Device Request: {}".format(json_body))
|
|
||||||
randomid = ''.join(random.sample(string.ascii_letters,6))
|
randomid = ''.join(random.sample(string.ascii_letters,6))
|
||||||
retcmd = await self.helperbot.send_command(json_body, randomid)
|
retcmd = await self.helperbot.send_command(json_body, randomid)
|
||||||
body = retcmd
|
body = retcmd
|
||||||
|
|
||||||
logging.info("Device Response: {}".format(body))
|
confserverlog.debug("\r\n POST: {} \r\n Response: {}".format(json_body,body))
|
||||||
return web.json_response(body)
|
return web.json_response(body)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
def disconnect(self):
|
def disconnect(self):
|
||||||
try:
|
try:
|
||||||
logging.info('ConfServer: shutting down...')
|
confserverlog.info('shutting down')
|
||||||
if(self.run_async):
|
if(self.run_async):
|
||||||
self.server.join()
|
self.server.join()
|
||||||
else:
|
else:
|
||||||
self.server.disconnect()
|
self.server.disconnect()
|
||||||
logging.info('ConfServer: bye')
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('ConfServer: {}'.format(e))
|
confserverlog.exception('{}'.format(e))
|
||||||
|
|
@ -16,8 +16,17 @@ import bumper
|
||||||
import json
|
import json
|
||||||
from datetime import datetime, timedelta
|
from datetime import datetime, timedelta
|
||||||
|
|
||||||
|
helperbotlog = logging.getLogger("helperbot")
|
||||||
|
mqttserverlog = logging.getLogger("mqttserver")
|
||||||
|
|
||||||
|
logging.getLogger("transitions").setLevel(logging.CRITICAL + 1) #Ignore this logger
|
||||||
|
logging.getLogger("passlib").setLevel(logging.CRITICAL + 1) #Ignore this logger
|
||||||
|
logging.getLogger("hbmqtt.broker").setLevel(logging.CRITICAL + 1) #Ignore this logger #There are some sublogs that could be set if needed (.plugins)
|
||||||
|
logging.getLogger("hbmqtt.mqtt.protocol").setLevel(logging.CRITICAL + 1) #Ignore this logger
|
||||||
|
logging.getLogger("hbmqtt.client").setLevel(logging.CRITICAL + 1) #Ignore this logger
|
||||||
|
|
||||||
class MQTTHelperBot():
|
class MQTTHelperBot():
|
||||||
|
|
||||||
Client = MQTTClient()
|
Client = MQTTClient()
|
||||||
def __init__(self, address, run_async=False, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar):
|
def __init__(self, address, run_async=False, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar):
|
||||||
|
|
||||||
|
|
@ -28,7 +37,7 @@ class MQTTHelperBot():
|
||||||
try:
|
try:
|
||||||
if run_async:
|
if run_async:
|
||||||
hloop = asyncio.new_event_loop()
|
hloop = asyncio.new_event_loop()
|
||||||
logging.debug("Starting MQTT HelperBot Thread: 1")
|
helperbotlog.debug("Starting MQTT HelperBot Thread: 1")
|
||||||
helperbot = Thread(name="MQTTHelperBot_Thread",target=self.run_helperbot, args=(hloop,))
|
helperbot = Thread(name="MQTTHelperBot_Thread",target=self.run_helperbot, args=(hloop,))
|
||||||
helperbot.setDaemon(True)
|
helperbot.setDaemon(True)
|
||||||
helperbot.start()
|
helperbot.start()
|
||||||
|
|
@ -37,7 +46,7 @@ class MQTTHelperBot():
|
||||||
self.run_helperbot()
|
self.run_helperbot()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('Helperbot: {}'.format(e))
|
helperbotlog.exception('{}'.format(e))
|
||||||
pass
|
pass
|
||||||
|
|
||||||
def run_helperbot(self, loop):
|
def run_helperbot(self, loop):
|
||||||
|
|
@ -48,7 +57,7 @@ class MQTTHelperBot():
|
||||||
loop.run_until_complete(self.get_msg())
|
loop.run_until_complete(self.get_msg())
|
||||||
loop.run_forever()
|
loop.run_forever()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('Helperbot: {}'.format(e))
|
helperbotlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def start_helper_bot(self):
|
async def start_helper_bot(self):
|
||||||
|
|
||||||
|
|
@ -60,35 +69,29 @@ class MQTTHelperBot():
|
||||||
])
|
])
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('Helperbot: {}'.format(e))
|
helperbotlog.exception('{}'.format(e))
|
||||||
#except hbmqtt.client.ClientException as ce:
|
|
||||||
# logging.exception("Client exception: %s" % ce)
|
|
||||||
|
|
||||||
|
|
||||||
async def get_msg(self):
|
async def get_msg(self):
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
message = await self.Client.deliver_message()
|
message = await self.Client.deliver_message()
|
||||||
|
|
||||||
#logging.debug("HelperBot MQTT Received Message on Topic: {} - Message: {}".format(message.topic, str(message.payload.decode("utf-8"))))
|
#helperbotlog.debug("HelperBot MQTT Received Message on Topic: {} - Message: {}".format(message.topic, str(message.payload.decode("utf-8"))))
|
||||||
cresp = self.command_responses.get()
|
cresp = self.command_responses.get()
|
||||||
|
|
||||||
#Cleanup "expired messages" > 60 seconds from time
|
#Cleanup "expired messages" > 60 seconds from time
|
||||||
for msg in cresp:
|
for msg in cresp:
|
||||||
expire_time = (datetime.fromtimestamp(msg['time']) + timedelta(seconds=10)).timestamp()
|
expire_time = (datetime.fromtimestamp(msg['time']) + timedelta(seconds=10)).timestamp()
|
||||||
if time.time() > expire_time:
|
if time.time() > expire_time:
|
||||||
#logging.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time))
|
#helperbotlog.debug("Pruning Message Time: {}, MsgTime: {}, MsgTime+60: {}".format(time.time(), msg['time'], expire_time))
|
||||||
cresp.remove(msg)
|
cresp.remove(msg)
|
||||||
|
|
||||||
cresp.append({"time": time.time() ,"topic": message.topic,"payload":str(message.data.decode("utf-8"))})
|
cresp.append({"time": time.time() ,"topic": message.topic,"payload":str(message.data.decode("utf-8"))})
|
||||||
self.command_responses.set(cresp)
|
self.command_responses.set(cresp)
|
||||||
logging.debug("MQTT Command Response List Count: %s" %len(cresp))
|
#helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('Helperbot: {}'.format(e))
|
helperbotlog.exception('{}'.format(e))
|
||||||
#except hbmqtt.client.ClientException as ce:
|
|
||||||
# logging.error("Client exception: %s" % ce)
|
|
||||||
|
|
||||||
|
|
||||||
async def wait_for_resp(self, requestid):
|
async def wait_for_resp(self, requestid):
|
||||||
try:
|
try:
|
||||||
|
|
@ -100,7 +103,7 @@ class MQTTHelperBot():
|
||||||
for msg in responses:
|
for msg in responses:
|
||||||
topic = str(msg['topic']).split("/")
|
topic = str(msg['topic']).split("/")
|
||||||
if (topic[6] == "helper1" and topic[10] == requestid):
|
if (topic[6] == "helper1" and topic[10] == requestid):
|
||||||
logging.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload']))
|
#helperbotlog.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload']))
|
||||||
if topic[11] == "j":
|
if topic[11] == "j":
|
||||||
resppayload = json.loads(msg['payload'])
|
resppayload = json.loads(msg['payload'])
|
||||||
else:
|
else:
|
||||||
|
|
@ -116,9 +119,10 @@ class MQTTHelperBot():
|
||||||
return resp
|
return resp
|
||||||
|
|
||||||
return { "id": requestid, "errno": "timeout", "ret": "fail" }
|
return { "id": requestid, "errno": "timeout", "ret": "fail" }
|
||||||
|
except asyncio.CancelledError as e:
|
||||||
|
helperbotlog.debug('wait_for_resp cancelled by asyncio')
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('Helperbot: {}'.format(e))
|
helperbotlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def send_command(self, cmdjson, requestid):
|
async def send_command(self, cmdjson, requestid):
|
||||||
try:
|
try:
|
||||||
|
|
@ -126,15 +130,15 @@ class MQTTHelperBot():
|
||||||
cmdjson["toId"], cmdjson["toType"], cmdjson["toRes"], requestid, cmdjson["payloadType"])
|
cmdjson["toId"], cmdjson["toType"], cmdjson["toRes"], requestid, cmdjson["payloadType"])
|
||||||
try:
|
try:
|
||||||
await self.Client.publish(ttopic, str(cmdjson["payload"]).encode(),QOS_0)
|
await self.Client.publish(ttopic, str(cmdjson["payload"]).encode(),QOS_0)
|
||||||
except:
|
except Exception as e:
|
||||||
logging.exception("Exception at send_command")
|
helperbotlog.exception("{}".format(e))
|
||||||
|
|
||||||
resp = await self.wait_for_resp(requestid)
|
resp = await self.wait_for_resp(requestid)
|
||||||
|
|
||||||
return resp
|
return resp
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('Helperbot: {}'.format(e))
|
helperbotlog.exception('{}'.format(e))
|
||||||
|
|
||||||
|
|
||||||
class MQTTServer():
|
class MQTTServer():
|
||||||
|
|
@ -149,21 +153,21 @@ class MQTTServer():
|
||||||
|
|
||||||
except PermissionError as e:
|
except PermissionError as e:
|
||||||
if "bind" in e.strerror:
|
if "bind" in e.strerror:
|
||||||
logging.exception("Error binding mqttserver, exiting. Try using a different hostname or IP.\r\n {}".format(e))
|
mqttserverlog.exception("Error binding mqttserver, exiting. Try using a different hostname or IP - {}".format(e))
|
||||||
exit(1)
|
exit(1)
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception('MQTTServer: {}'.format(e))
|
mqttserverlog.exception('{}'.format(e))
|
||||||
exit(1)
|
exit(1)
|
||||||
|
|
||||||
async def active_bot_listing(self):
|
async def active_bot_listing(self):
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
await asyncio.sleep(5)
|
await asyncio.sleep(5)
|
||||||
logging.debug('Connected bots: %s' % self.bumper_bots.get())
|
mqttserverlog.debug('connected bots - %s' % self.bumper_bots.get())
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('MQTTServer: {}'.format(e))
|
mqttserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
def __init__(self, address, run_async=False, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar):
|
def __init__(self, address, run_async=False, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar):
|
||||||
try:
|
try:
|
||||||
|
|
@ -206,7 +210,7 @@ class MQTTServer():
|
||||||
}
|
}
|
||||||
if run_async:
|
if run_async:
|
||||||
sloop = asyncio.new_event_loop()
|
sloop = asyncio.new_event_loop()
|
||||||
logging.debug("Starting MQTTServer Thread: 1")
|
mqttserverlog.debug("Starting MQTTServer Thread: 1")
|
||||||
mqttserver = Thread(name="MQTTServer_Thread",target=self.run_server, args=(sloop,))
|
mqttserver = Thread(name="MQTTServer_Thread",target=self.run_server, args=(sloop,))
|
||||||
mqttserver.setDaemon(True)
|
mqttserver.setDaemon(True)
|
||||||
mqttserver.start()
|
mqttserver.start()
|
||||||
|
|
@ -215,11 +219,7 @@ class MQTTServer():
|
||||||
self.run_server()
|
self.run_server()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('MQTTServer: {}'.format(e))
|
mqttserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
#except:
|
|
||||||
# logging.exception("Exception")
|
|
||||||
# pass
|
|
||||||
|
|
||||||
|
|
||||||
def run_server(self, loop):
|
def run_server(self, loop):
|
||||||
|
|
@ -230,7 +230,7 @@ class MQTTServer():
|
||||||
loop.run_forever()
|
loop.run_forever()
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('MQTTServer: {}'.format(e))
|
mqttserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
class BumperMQTTServer_Plugin:
|
class BumperMQTTServer_Plugin:
|
||||||
def __init__(self, context):
|
def __init__(self, context):
|
||||||
|
|
@ -240,11 +240,11 @@ class BumperMQTTServer_Plugin:
|
||||||
except KeyError:
|
except KeyError:
|
||||||
self.context.logger.warning("'clients' section not found in context configuration")
|
self.context.logger.warning("'clients' section not found in context configuration")
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('MQTTServer: {}'.format(e))
|
mqttserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
async def on_broker_client_connected(self, client_id):
|
async def on_broker_client_connected(self, client_id):
|
||||||
try:
|
try:
|
||||||
logging.debug('Bumper Connection: %s connected' % client_id)
|
#mqttserverlog.debug('%s connected' % client_id)
|
||||||
connected_bots = self.clients['connected_bots'].get()
|
connected_bots = self.clients['connected_bots'].get()
|
||||||
connected_clients = self.clients['connected_clients'].get()
|
connected_clients = self.clients['connected_clients'].get()
|
||||||
didsplit = str(client_id).split("@")
|
didsplit = str(client_id).split("@")
|
||||||
|
|
@ -262,7 +262,7 @@ class BumperMQTTServer_Plugin:
|
||||||
|
|
||||||
if botactive == False:
|
if botactive == False:
|
||||||
connected_bots.append(newbot.asdict())
|
connected_bots.append(newbot.asdict())
|
||||||
logging.info("Adding bot to list: {}".format(newbot.asdict()))
|
mqttserverlog.info("new bot {}".format(newbot.did))
|
||||||
|
|
||||||
self.clients['connected_bots'].set(connected_bots)
|
self.clients['connected_bots'].set(connected_bots)
|
||||||
else:
|
else:
|
||||||
|
|
@ -279,41 +279,39 @@ class BumperMQTTServer_Plugin:
|
||||||
|
|
||||||
if clientactive == False:
|
if clientactive == False:
|
||||||
connected_clients.append(newuser.asdict())
|
connected_clients.append(newuser.asdict())
|
||||||
logging.info("Adding client to list: {}".format(newuser.asdict()))
|
mqttserverlog.info("new client {}".format(newuser.userid))
|
||||||
|
|
||||||
self.clients['connected_clients'].set(connected_clients)
|
self.clients['connected_clients'].set(connected_clients)
|
||||||
|
|
||||||
|
#mqttserverlog.debug('Connected Bots: %s' %self.clients['connected_bots'].get())
|
||||||
logging.debug('Connected Bots: %s' %self.clients['connected_bots'].get())
|
#mqttserverlog.debug('Connected Clients: %s' %self.clients['connected_clients'].get())
|
||||||
logging.debug('Connected Clients: %s' %self.clients['connected_clients'].get())
|
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('MQTTServer: {}'.format(e))
|
mqttserverlog.exception('{}'.format(e))
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
async def on_broker_client_disconnected(self, client_id):
|
async def on_broker_client_disconnected(self, client_id):
|
||||||
try:
|
try:
|
||||||
logging.debug('Bumper Connection: %s disconnected' % client_id)
|
#mqttserverlog.debug('%s disconnected' % client_id)
|
||||||
connected_bots = self.clients['connected_bots'].get()
|
connected_bots = self.clients['connected_bots'].get()
|
||||||
connected_clients = self.clients['connected_clients'].get()
|
connected_clients = self.clients['connected_clients'].get()
|
||||||
didsplit = str(client_id).split("@")
|
didsplit = str(client_id).split("@")
|
||||||
#If the did is in the list, remove it
|
#If the did is in the list, remove it
|
||||||
for bot in connected_bots:
|
for bot in connected_bots:
|
||||||
if didsplit[0] == bot['did']:
|
if didsplit[0] == bot['did']:
|
||||||
logging.info("Removing bot from list: {}".format(bot['did']))
|
mqttserverlog.info("bot disconnected {}".format(bot['did']))
|
||||||
connected_bots.remove(bot)
|
connected_bots.remove(bot)
|
||||||
self.clients['connected_bots'].set(connected_bots)
|
self.clients['connected_bots'].set(connected_bots)
|
||||||
|
|
||||||
logging.debug('Connected Bots: %s' %self.clients['connected_bots'].get())
|
|
||||||
|
|
||||||
for client in connected_clients:
|
for client in connected_clients:
|
||||||
if didsplit[0] == client['userid']:
|
if didsplit[0] == client['userid']:
|
||||||
logging.info("Removing client from list: {}".format(client['userid']))
|
mqttserverlog.info("client disconnected {}".format(client['userid']))
|
||||||
connected_clients.remove(client)
|
connected_clients.remove(client)
|
||||||
self.clients['connected_clients'].set(connected_clients)
|
self.clients['connected_clients'].set(connected_clients)
|
||||||
|
|
||||||
logging.debug('Connected Clients: %s' %self.clients['connected_clients'].get())
|
#mqttserverlog.debug('Connected Bots: %s' %self.clients['connected_bots'].get())
|
||||||
|
#mqttserverlog.debug('Connected Clients: %s' %self.clients['connected_clients'].get())
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('MQTTServer: {}'.format(e))
|
mqttserverlog.exception('{}'.format(e))
|
||||||
|
|
@ -2,6 +2,7 @@
|
||||||
|
|
||||||
import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as ET
|
import sys, socket, threading, re, time, logging, uuid, xml.etree.ElementTree as ET
|
||||||
|
|
||||||
|
xmppserverlog = logging.getLogger("xmppserver")
|
||||||
|
|
||||||
class XMPPServer():
|
class XMPPServer():
|
||||||
server_id = 'bumper'
|
server_id = 'bumper'
|
||||||
|
|
@ -17,9 +18,8 @@ class XMPPServer():
|
||||||
self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
self.socket.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||||
self.socket.bind(address)
|
self.socket.bind(address)
|
||||||
self.socket.listen(1)
|
self.socket.listen(1)
|
||||||
logging.info('XMPPServer: listening on {}:{}'.format(address[0], address[1]))
|
xmppserverlog.debug('listening on {}:{}'.format(address[0], address[1]))
|
||||||
while not self.exit_flag:
|
while not self.exit_flag:
|
||||||
logging.info('XMPPServer: awaiting connection')
|
|
||||||
connection, client_address = self.socket.accept()
|
connection, client_address = self.socket.accept()
|
||||||
# disconnect any clients with this ip
|
# disconnect any clients with this ip
|
||||||
for client in self.clients:
|
for client in self.clients:
|
||||||
|
|
@ -34,27 +34,27 @@ class XMPPServer():
|
||||||
self.socket.close()
|
self.socket.close()
|
||||||
except PermissionError as e:
|
except PermissionError as e:
|
||||||
if "bind" in e.strerror:
|
if "bind" in e.strerror:
|
||||||
logging.exception("Error binding xmppserver, exiting. Try using a different hostname or IP.\r\n {}".format(e))
|
xmppserverlog.exception("Error binding XMPPServer, exiting. Try using a different hostname or IP - {}".format(e))
|
||||||
exit(1)
|
exit(1)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception('XMPPServer: {}'.format(e))
|
xmppserverlog.exception('{}'.format(e))
|
||||||
exit(1)
|
exit(1)
|
||||||
except KeyboardInterrupt:
|
except KeyboardInterrupt as e:
|
||||||
logging.exception('XMPPServer: Keyboard interrupt')
|
xmppserverlog.exception('{}'.format(e))
|
||||||
finally:
|
finally:
|
||||||
self.disconnect()
|
self.disconnect()
|
||||||
logging.info('XMPPServer: bye')
|
xmppserverlog.info('disconnecting')
|
||||||
|
|
||||||
def disconnect(self):
|
def disconnect(self):
|
||||||
try:
|
try:
|
||||||
logging.info('XMPPServer: waiting for all client threads to exit')
|
xmppserverlog.debug('waiting for all client threads to exit')
|
||||||
for client in self.clients:
|
for client in self.clients:
|
||||||
client.disconnect()
|
client.disconnect()
|
||||||
client.join()
|
client.join()
|
||||||
self.exit_flag = True
|
self.exit_flag = True
|
||||||
logging.info('XMPPServer: shutting down...')
|
xmppserverlog.debug('shutting down')
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception("XMPPServer Exception: {}".format(e))
|
xmppserverlog.exception("{}".format(e))
|
||||||
|
|
||||||
|
|
||||||
class Client(threading.Thread):
|
class Client(threading.Thread):
|
||||||
|
|
@ -80,23 +80,23 @@ class Client(threading.Thread):
|
||||||
|
|
||||||
def send(self, command):
|
def send(self, command):
|
||||||
try:
|
try:
|
||||||
logging.debug('XMPPServer to {}: {}'.format(self.address, command))
|
xmppserverlog.debug('send {} - {}'.format(self.address, command))
|
||||||
self.connection.send(command.encode())
|
self.connection.send(command.encode())
|
||||||
except OSError as e:
|
except OSError as e:
|
||||||
logging.error('XMPPServer: {}'.format(e))
|
xmppserverlog.error('{}'.format(e))
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception("XMPPServer Exception: {}".format(e))
|
xmppserverlog.exception("{}".format(e))
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
def disconnect(self):
|
def disconnect(self):
|
||||||
try:
|
try:
|
||||||
logging.info('XMPPServer: {} with resource {} disconnecting'.format(self.address, self.clientresource))
|
xmppserverlog.debug('client {} with resource {} disconnecting'.format(self.address, self.clientresource))
|
||||||
self.connection.close()
|
self.connection.close()
|
||||||
self._set_state('DISCONNECT')
|
self._set_state('DISCONNECT')
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception("XMPPServer Exception: {}".format(e))
|
xmppserverlog.exception("{}".format(e))
|
||||||
|
|
||||||
def _tag_strip_uri(self, tag):
|
def _tag_strip_uri(self, tag):
|
||||||
try:
|
try:
|
||||||
|
|
@ -104,25 +104,25 @@ class Client(threading.Thread):
|
||||||
uri, ignore, tag = tag[1:].partition('}')
|
uri, ignore, tag = tag[1:].partition('}')
|
||||||
return tag
|
return tag
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception("XMPPServer Exception: {}".format(e))
|
xmppserverlog.exception("{}".format(e))
|
||||||
|
|
||||||
def _set_state(self, state):
|
def _set_state(self, state):
|
||||||
try:
|
try:
|
||||||
new_state = getattr(Client, state)
|
new_state = getattr(Client, state)
|
||||||
if self.state > new_state:
|
if self.state > new_state:
|
||||||
raise Exception('XMPPServer: {} illegal state change {}->{}'.format(self.address, self.state, new_state))
|
raise Exception('{} illegal state change {}->{}'.format(self.address, self.state, new_state))
|
||||||
logging.debug('XMPPServer: {} state: {}'.format(self.address, state))
|
xmppserverlog.debug('{} state: {}'.format(self.address, state))
|
||||||
self.state = new_state
|
self.state = new_state
|
||||||
if new_state == '5':
|
if new_state == '5':
|
||||||
self.join()
|
self.join()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception("XMPPServer: Exception: {}".format(e))
|
xmppserverlog.exception("{}".format(e))
|
||||||
|
|
||||||
def _handle_ctl(self, xml, data):
|
def _handle_ctl(self, xml, data):
|
||||||
try:
|
try:
|
||||||
ctl = xml[0][0]
|
ctl = xml[0][0]
|
||||||
if ctl.get('admin') and self.type == self.BOT:
|
if ctl.get('admin') and self.type == self.BOT:
|
||||||
logging.debug('XMPPServer: admin username received from bot: {}'.format(ctl.get('admin')))
|
xmppserverlog.debug('admin username received from bot: {}'.format(ctl.get('admin')))
|
||||||
XMPPServer.client_id = ctl.get('admin')
|
XMPPServer.client_id = ctl.get('admin')
|
||||||
return
|
return
|
||||||
# forward
|
# forward
|
||||||
|
|
@ -136,7 +136,7 @@ class Client(threading.Thread):
|
||||||
data = data.encode()
|
data = data.encode()
|
||||||
client.send(data.decode('utf-8'))
|
client.send(data.decode('utf-8'))
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception("XMPPServer Exception: {}".format(e))
|
xmppserverlog.exception("{}".format(e))
|
||||||
|
|
||||||
|
|
||||||
def _handle_ping(self, xml, data):
|
def _handle_ping(self, xml, data):
|
||||||
|
|
@ -155,19 +155,24 @@ class Client(threading.Thread):
|
||||||
if client.address != self.address and client.state == client.READY:
|
if client.address != self.address and client.state == client.READY:
|
||||||
client.send(data.decode('utf-8'))
|
client.send(data.decode('utf-8'))
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.exception("XMPPServer Exception: {}".format(e))
|
xmppserverlog.exception("{}".format(e))
|
||||||
|
|
||||||
def run(self):
|
def run(self):
|
||||||
try:
|
try:
|
||||||
logging.info('XMPPServer: client connected: {}'.format(self.address))
|
xmppserverlog.info('client connected - {}'.format(self.address))
|
||||||
self._set_state('CONNECT')
|
self._set_state('CONNECT')
|
||||||
data = ""
|
data = ""
|
||||||
while True:
|
while True:
|
||||||
time.sleep(0.2)
|
time.sleep(0.2)
|
||||||
if not self.connection._closed:
|
if not self.connection._closed:
|
||||||
data = self.connection.recv(4096)
|
try:
|
||||||
|
data = self.connection.recv(4096)
|
||||||
|
except ConnectionResetError as e:
|
||||||
|
xmppserverlog.error('{}'.format(e))
|
||||||
|
except Exception as e:
|
||||||
|
xmppserverlog.exception('{}'.format(e))
|
||||||
if data:
|
if data:
|
||||||
logging.debug('XMPPServer: from {}: {}'.format(self.address, data.decode('utf-8')))
|
xmppserverlog.debug('from {} - {}'.format(self.address, data.decode('utf-8')))
|
||||||
try:
|
try:
|
||||||
if self.state == self.CONNECT:
|
if self.state == self.CONNECT:
|
||||||
if data.decode('utf-8').find('jabber:client') > -1:
|
if data.decode('utf-8').find('jabber:client') > -1:
|
||||||
|
|
@ -179,7 +184,6 @@ class Client(threading.Thread):
|
||||||
self.send('<stream:features><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"/><session xmlns="urn:ietf:params:xml:ns:xmpp-session"/></stream:features>')
|
self.send('<stream:features><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"/><session xmlns="urn:ietf:params:xml:ns:xmpp-session"/></stream:features>')
|
||||||
continue
|
continue
|
||||||
xml = ET.fromstring(data)
|
xml = ET.fromstring(data)
|
||||||
logging.debug("XMPPXML: {}".format(data))
|
|
||||||
if len(xml):
|
if len(xml):
|
||||||
child = self._tag_strip_uri(xml[0].tag)
|
child = self._tag_strip_uri(xml[0].tag)
|
||||||
else:
|
else:
|
||||||
|
|
@ -191,7 +195,7 @@ class Client(threading.Thread):
|
||||||
clientresourcexml = clientbindxml[0].getchildren()
|
clientresourcexml = clientbindxml[0].getchildren()
|
||||||
self.clientresource = clientresourcexml[0].text
|
self.clientresource = clientresourcexml[0].text
|
||||||
self.name = "XMPP Thread {}".format(self.clientresource)
|
self.name = "XMPP Thread {}".format(self.clientresource)
|
||||||
logging.info("XMPP Client {} using resource {}".format(self.address, self.clientresource))
|
xmppserverlog.debug("new client {} using resource {}".format(self.address, self.clientresource))
|
||||||
res = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format(xml.get('id'), XMPPServer.bot_id)
|
res = '<iq type="result" id="{}"><bind xmlns="urn:ietf:params:xml:ns:xmpp-bind"><jid>{}</jid></bind></iq>'.format(xml.get('id'), XMPPServer.bot_id)
|
||||||
self._set_state('BIND')
|
self._set_state('BIND')
|
||||||
elif child == 'session':
|
elif child == 'session':
|
||||||
|
|
@ -209,27 +213,27 @@ class Client(threading.Thread):
|
||||||
if len(xml) and xml[0].tag == 'status':
|
if len(xml) and xml[0].tag == 'status':
|
||||||
# bot announcing arrival
|
# bot announcing arrival
|
||||||
self.type = self.BOT
|
self.type = self.BOT
|
||||||
logging.debug('XMPPServer: {} type set to BOT (based on presence tag)'.format(self.address))
|
xmppserverlog.debug('{} type set to BOT (based on presence tag)'.format(self.address))
|
||||||
# send a command from an unknown user - the response will contain the correct admin username
|
# send a command from an unknown user - the response will contain the correct admin username
|
||||||
self.send('<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="GetCleanState" /></query></iq>'.format(uuid.uuid4(), 'unknown@ecouser.net', XMPPServer.bot_id))
|
self.send('<iq type="set" id="{}" from="{}" to="{}"><query xmlns="com:ctl"><ctl td="GetCleanState" /></query></iq>'.format(uuid.uuid4(), 'unknown@ecouser.net', XMPPServer.bot_id))
|
||||||
elif xml.get('type') == 'available':
|
elif xml.get('type') == 'available':
|
||||||
self.type = self.CONTROLLER
|
self.type = self.CONTROLLER
|
||||||
logging.debug('XMPPServer: {} type set to CONTROLLER (based on presence tag)'.format(self.address))
|
xmppserverlog.debug('{} type set to CONTROLLER (based on presence tag)'.format(self.address))
|
||||||
except ET.ParseError as e:
|
except ET.ParseError as e:
|
||||||
logging.debug('error: {}'.format(e))
|
xmppserverlog.debug('parse error - {}'.format(e))
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('XMPPServer: {}'.format(e))
|
xmppserverlog.exception('{}'.format(e))
|
||||||
self._set_state('DISCONNECT')
|
self._set_state('DISCONNECT')
|
||||||
except OSError as e:
|
except OSError as e:
|
||||||
logging.error('XMPPServer: {}'.format(e))
|
xmppserverlog.error('{}'.format(e))
|
||||||
self._set_state('DISCONNECT')
|
self._set_state('DISCONNECT')
|
||||||
|
|
||||||
except ConnectionResetError as e:
|
except ConnectionResetError as e:
|
||||||
logging.error('XMPPServer: {}'.format(e))
|
xmppserverlog.error('{}'.format(e))
|
||||||
self._set_state('DISCONNECT')
|
self._set_state('DISCONNECT')
|
||||||
|
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
logging.error('XMPPServer: {}'.format(e))
|
xmppserverlog.exception('{}'.format(e))
|
||||||
self._set_state('DISCONNECT')
|
self._set_state('DISCONNECT')
|
||||||
|
|
||||||
finally:
|
finally:
|
||||||
|
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue