Beginning bumper's new journey #4

Merged
bmartin5692 merged 37 commits from dev into master 2019-02-22 14:34:25 +01:00
4 changed files with 456 additions and 313 deletions
Showing only changes of commit eb4cb2b996 - Show all commits

View file

@ -4,23 +4,18 @@ import logging
import bumper
import sys, socket
import time
import platform
args = sys.argv
if len(args) > 0:
if '--debug' in args:
logging.basicConfig(level=logging.DEBUG,
format="[%(asctime)s] :: %(levelname)s :: %(name)s :: %(message)s")
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")
#conf_address = (socket.gethostbyname(socket.gethostname()), 443)
conf_address = ("0.0.0.0", 443)
#xmpp_address = (socket.gethostbyname(socket.gethostname()), 5223)
xmpp_address = ("0.0.0.0", 5223)
#mqtt_address = (socket.gethostbyname(socket.gethostname()), 8883)
mqtt_address = ("0.0.0.0", 8883)
#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
# dbot = bumper.VacBotDevice("did", "class", "resource", "name","nick" )
@ -29,13 +24,26 @@ mqtt_address = ("0.0.0.0", 8883)
# bclienttemp.append(dbot.asdict())
# bclient.set(bclienttemp)
# start mqtt server (async)
if platform.system() == "Darwin":
listen_host = "0.0.0.0"
else:
listen_host = socket.gethostbyname(socket.gethostname())
#listen_host = "localhost" #Try this if the above doesn't work
conf_address_443 = (listen_host, 443)
conf_address_8007 = (listen_host, 8007)
xmpp_address = (listen_host, 5223)
mqtt_address = (listen_host, 8883)
# start mqtt server on port 8883 (async)
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
# start mqtt server (async)
# start mqtt_helperbot (async)
mqtt_helperbot = bumper.MQTTHelperBot(mqtt_address, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var)
# start conf server (async)
conf_server = bumper.ConfServer(conf_address, usessl=True, run_async=True,bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var, helperbot=mqtt_helperbot)
# start xmpp server (sync)
# 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)

View file

@ -47,201 +47,288 @@ class ConfServer():
except:
loop = asyncio.new_event_loop()
loop.run_until_complete(self.start_server())
loop.run_forever()
try:
loop.run_until_complete(self.start_server())
loop.run_forever()
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def start_server(self):
app = web.Application()
try:
app = web.Application()
app.add_routes([
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/login', self.handle_login),
# web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/checkLogin', self.handle_checkLogin),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/logout', self.handle_logout),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/getAuthCode', self.handle_getAuthCode),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/checkAgreement', self.handle_checkAgreement),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/common/checkVersion', self.handle_checkVersion),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/campaign/homePageAlert', self.handle_homePageAlert),
app.add_routes([
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/login', self.handle_login),
# web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/checkLogin', self.handle_checkLogin),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/logout', self.handle_logout),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/getAuthCode', self.handle_getAuthCode),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/checkAgreement', self.handle_checkAgreement),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/common/checkVersion', self.handle_checkVersion),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/campaign/homePageAlert', self.handle_homePageAlert),
web.post('/api/users/user.do', self.handle_usersapi),
web.post('/api/pim/product/getProductIotMap', self.handle_getProductIotMap),
web.post('/api/iot/devmanager.do', self.handle_devmanager_botcommand)
])
web.post('/api/users/user.do', self.handle_usersapi),
web.get('/api/users/user.do', self.handle_usersapi),
web.post('/api/pim/product/getProductIotMap', self.handle_getProductIotMap),
web.post('/api/iot/devmanager.do', self.handle_devmanager_botcommand),
web.post('/lookup.do', self.handle_lookup),
])
runner = web.AppRunner(app, access_log=None) #access_log=None so the output isn't nuts
await runner.setup()
runner = web.AppRunner(app)#, access_log=None) #access_log=None so the output isn't nuts
await runner.setup()
if self.usessl:
ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
ssl_ctx.load_cert_chain(bumper.server_cert,bumper.server_key)
site = web.TCPSite(runner, host=self.address[0], port=self.address[1],ssl_context=ssl_ctx)
if self.usessl:
ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
ssl_ctx.load_cert_chain(bumper.server_cert,bumper.server_key)
site = web.TCPSite(runner, host=self.address[0], port=self.address[1],ssl_context=ssl_ctx)
else:
site = web.TCPSite(runner, host=self.address[0], port=self.address[1])
else:
site = web.TCPSite(runner, host=self.address[0], port=self.address[1])
await site.start()
await site.start()
except PermissionError as e:
if "bind" in e.strerror:
logging.exception("Error binding confserver, exiting. Try using a different hostname or IP.\r\n {}".format(e))
exit(1)
except Exception as e:
logging.exception('ConfServer: {}'.format(e))
exit(1)
async def handle_login(self, request):
#Could implement basic auth if you wanted, or just accept anything
countrycode = request.match_info.get('country', "us")
body = {
"code": "0000",
"data": {
"accessToken": "tempaccesstoken", #Random chars 32 length
"country": countrycode,
"email": "null@null.com",
"uid": "fuid_{}".format(''.join(random.sample(string.ascii_letters,6))), #Date(14)_RandomChars(32)
"username": "fusername_{}".format(''.join(random.sample(string.ascii_letters,6))) #Random chars 8
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
async def handle_checkLogin(self, request):
# The app seems to remember it's last uid and accessToken
# If these don't match, it fails
countrycode = request.match_info.get('country', "us")
body = {
"code": "0000",
"data": {
"accessToken": "tempaccesstoken", #Random chars 32 length
"country": countrycode,
"email": "null@null.com",
"uid": "fuid_{}".format(''.join(random.sample(string.ascii_letters,6))), #Date(14)_RandomChars(32)
"username": "fusername_{}".format(''.join(random.sample(string.ascii_letters,6))) #Random chars 8
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
async def handle_logout(self, request):
body = {"code": "0000","data": None,"msg": "操作成功", "time": bumper.get_milli_time(time.time())}
#TODO - when logging out close out any other connections MQTT/XMPP
return web.json_response(body)
async def handle_getAuthCode(self, request):
countrycode = request.match_info.get('country', "us")
body = {
"code": "0000",
"data": {
"authCode": "{}_tempauthcode".format(countrycode), #countrycode_randomchars(32)
"ecovacsUid": "fuid_{}".format(''.join(random.sample(string.ascii_letters,6))) #Date(14)_RandomChars(32)
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
async def handle_checkVersion(self, request):
body = {
"code": "0000",
"data": {
"c": None,
"img": None,
"r": 0,
"t": None,
"u": None,
"ut": 0,
"v": None
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
async def handle_checkAgreement(self, request):
body = {
"code": "0000",
"data": [],
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
async def handle_homePageAlert(self, request):
nextAlert = bumper.get_milli_time((datetime.now() + timedelta(hours=12)).timestamp())
body = {
"code": "0000",
"data": {
"clickSchemeUrl": None,
"clickWebUrl": None,
"hasCampaign": "N",
"imageUrl": None,
"nextAlertTime": nextAlert,
"serverTime": bumper.get_milli_time(time.time())
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
async def handle_getProductIotMap(self, request):
#json_body = json.loads(await request.text())
body = {"code":0,"data":[{"classid":"dl8fht","product":{"_id":"5acb0fa87c295c0001876ecf","name":"DEEBOT 600 Series","icon":"5acc32067c295c0001876eea","UILogicId":"dl8fht","ota":False,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5acc32067c295c0001876eea"}},{"classid":"02uwxm","product":{"_id":"5ae1481e7ccd1a0001e1f69e","name":"DEEBOT OZMO Slim10 Series","icon":"5b1dddc48bc45700014035a1","UILogicId":"02uwxm","ota":False,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5b1dddc48bc45700014035a1"}},{"classid":"y79a7u","product":{"_id":"5b04c0227ccd1a0001e1f6a8","name":"DEEBOT OZMO 900","icon":"5b04c0217ccd1a0001e1f6a7","UILogicId":"y79a7u","ota":True,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5b04c0217ccd1a0001e1f6a7"}},{"classid":"jr3pqa","product":{"_id":"5b43077b8bc457000140363e","name":"DEEBOT 711","icon":"5b5ac4cc8d5a56000111e769","UILogicId":"jr3pqa","ota":True,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5b5ac4cc8d5a56000111e769"}},{"classid":"uv242z","product":{"_id":"5b5149b4ac0b87000148c128","name":"DEEBOT 710","icon":"5b5ac4e45f21100001882bb9","UILogicId":"uv242z","ota":True,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5b5ac4e45f21100001882bb9"}},{"classid":"ls1ok3","product":{"_id":"5b6561060506b100015c8868","name":"DEEBOT 900 Series","icon":"5ba4a2cb6c2f120001c32839","UILogicId":"ls1ok3","ota":True,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5ba4a2cb6c2f120001c32839"}}]}
return web.json_response(body)
async def handle_usersapi(self, request):
body = {}
postbody = {}
if request.content_type == "application/x-www-form-urlencoded":
postbody = await request.post()
else:
postbody = json.loads(await request.text())
logging.debug(postbody)
todo = postbody['todo']
if todo == 'FindBest':
service = postbody['service']
if service == 'EcoMsgNew':
body = {"result":"ok","ip":socket.gethostbyname(socket.gethostname()),"port":5223}
elif service == 'EcoUpdate':
body = {"result":"ok","ip":"47.88.66.164","port":8005}
elif todo == 'loginByItToken':
try:
#Could implement basic auth if you wanted, or just accept anything
countrycode = request.match_info.get('country', "us")
body = {
"resource": postbody["resource"],
"result": "ok",
"todo": "result",
"token": postbody["token"], #RandomChar(32)
"userId": postbody["userId"] #RandomChar(16)
}
elif todo == 'GetDeviceList':
active_bots = self.bumper_bots.get()
body = {
"devices": active_bots,
"result": "ok",
"todo": "result"
"code": "0000",
"data": {
"accessToken": "tempaccesstoken", #Random chars 32 length
"country": countrycode,
"email": "null@null.com",
"uid": "fuid_{}".format(''.join(random.sample(string.ascii_letters,6))), #Date(14)_RandomChars(32)
"username": "fusername_{}".format(''.join(random.sample(string.ascii_letters,6))) #Random chars 8
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_checkLogin(self, request):
try:
# The app seems to remember it's last uid and accessToken
# If these don't match, it fails
countrycode = request.match_info.get('country', "us")
body = {
"code": "0000",
"data": {
"accessToken": "tempaccesstoken", #Random chars 32 length
"country": countrycode,
"email": "null@null.com",
"uid": "fuid_{}".format(''.join(random.sample(string.ascii_letters,6))), #Date(14)_RandomChars(32)
"username": "fusername_{}".format(''.join(random.sample(string.ascii_letters,6))) #Random chars 8
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_logout(self, request):
try:
body = {"code": "0000","data": None,"msg": "操作成功", "time": bumper.get_milli_time(time.time())}
#TODO - when logging out close out any other connections MQTT/XMPP
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_getAuthCode(self, request):
try:
countrycode = request.match_info.get('country', "us")
body = {
"code": "0000",
"data": {
"authCode": "{}_tempauthcode".format(countrycode), #countrycode_randomchars(32)
"ecovacsUid": "fuid_{}".format(''.join(random.sample(string.ascii_letters,6))) #Date(14)_RandomChars(32)
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_checkVersion(self, request):
try:
body = {
"code": "0000",
"data": {
"c": None,
"img": None,
"r": 0,
"t": None,
"u": None,
"ut": 0,
"v": None
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_checkAgreement(self, request):
try:
body = {
"code": "0000",
"data": [],
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_homePageAlert(self, request):
try:
nextAlert = bumper.get_milli_time((datetime.now() + timedelta(hours=12)).timestamp())
body = {
"code": "0000",
"data": {
"clickSchemeUrl": None,
"clickWebUrl": None,
"hasCampaign": "N",
"imageUrl": None,
"nextAlertTime": nextAlert,
"serverTime": bumper.get_milli_time(time.time())
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_getProductIotMap(self, request):
try:
#json_body = json.loads(await request.text())
body = {"code":0,"data":[{"classid":"dl8fht","product":{"_id":"5acb0fa87c295c0001876ecf","name":"DEEBOT 600 Series","icon":"5acc32067c295c0001876eea","UILogicId":"dl8fht","ota":False,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5acc32067c295c0001876eea"}},{"classid":"02uwxm","product":{"_id":"5ae1481e7ccd1a0001e1f69e","name":"DEEBOT OZMO Slim10 Series","icon":"5b1dddc48bc45700014035a1","UILogicId":"02uwxm","ota":False,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5b1dddc48bc45700014035a1"}},{"classid":"y79a7u","product":{"_id":"5b04c0227ccd1a0001e1f6a8","name":"DEEBOT OZMO 900","icon":"5b04c0217ccd1a0001e1f6a7","UILogicId":"y79a7u","ota":True,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5b04c0217ccd1a0001e1f6a7"}},{"classid":"jr3pqa","product":{"_id":"5b43077b8bc457000140363e","name":"DEEBOT 711","icon":"5b5ac4cc8d5a56000111e769","UILogicId":"jr3pqa","ota":True,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5b5ac4cc8d5a56000111e769"}},{"classid":"uv242z","product":{"_id":"5b5149b4ac0b87000148c128","name":"DEEBOT 710","icon":"5b5ac4e45f21100001882bb9","UILogicId":"uv242z","ota":True,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5b5ac4e45f21100001882bb9"}},{"classid":"ls1ok3","product":{"_id":"5b6561060506b100015c8868","name":"DEEBOT 900 Series","icon":"5ba4a2cb6c2f120001c32839","UILogicId":"ls1ok3","ota":True,"iconUrl":"https://portal-ww.ecouser.net/api/pim/file/get/5ba4a2cb6c2f120001c32839"}}]}
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_usersapi(self, request):
try:
body = {}
postbody = {}
if request.content_type == "application/x-www-form-urlencoded":
postbody = await request.post()
else:
postbody = json.loads(await request.text())
logging.debug(postbody)
todo = postbody['todo']
if todo == 'FindBest':
service = postbody['service']
if service == 'EcoMsgNew':
body = {"result":"ok","ip":socket.gethostbyname(socket.gethostname()),"port":5223}
elif service == 'EcoUpdate':
body = {"result":"ok","ip":"47.88.66.164","port":8005}
elif todo == 'loginByItToken':
body = {
"resource": postbody["resource"],
"result": "ok",
"todo": "result",
"token": postbody["token"], #RandomChar(32)
"userId": postbody["userId"] #RandomChar(16)
}
elif todo == 'GetDeviceList':
active_bots = self.bumper_bots.get()
body = {
"devices": active_bots,
"result": "ok",
"todo": "result"
}
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_lookup(self, request):
try:
body = {}
postbody = {}
if request.content_type == "application/x-www-form-urlencoded":
postbody = await request.post()
else:
postbody = json.loads(await request.text())
logging.debug(postbody)
todo = postbody['todo']
if todo == 'FindBest':
service = postbody['service']
if service == 'EcoMsgNew':
body = {"result":"ok","ip":socket.gethostbyname(socket.gethostname()),"port":5223}
elif service == 'EcoUpdate':
body = {"result":"ok","ip":"47.88.66.164","port":8005}
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
async def handle_devmanager_botcommand(self, request):
json_body = json.loads(await request.text())
randomid = ''.join(random.sample(string.ascii_letters,6))
retcmd = await self.helperbot.send_command(json_body, randomid)
body = retcmd
try:
json_body = json.loads(await request.text())
logging.info("Device Request: {}".format(json_body))
randomid = ''.join(random.sample(string.ascii_letters,6))
retcmd = await self.helperbot.send_command(json_body, randomid)
body = retcmd
return web.json_response(body)
logging.info("Device Response: {}".format(body))
return web.json_response(body)
except Exception as e:
logging.error('ConfServer: {}'.format(e))
def disconnect(self):
logging.info('ConfServer: shutting down...')
if(self.run_async):
self.server.join()
else:
self.server.disconnect()
logging.info('ConfServer: bye')
try:
logging.info('ConfServer: shutting down...')
if(self.run_async):
self.server.join()
else:
self.server.disconnect()
logging.info('ConfServer: bye')
except Exception as e:
logging.error('ConfServer: {}'.format(e))

View file

@ -36,17 +36,19 @@ class MQTTHelperBot():
else:
self.run_helperbot()
except:
logging.exception("Exception")
except Exception as e:
logging.error('Helperbot: {}'.format(e))
pass
def run_helperbot(self, loop):
asyncio.set_event_loop(loop)
self.Client = MQTTClient(client_id=self.client_id, config={'check_hostname':False})
loop.run_until_complete(self.start_helper_bot())
loop.run_until_complete(self.get_msg())
loop.run_forever()
try:
asyncio.set_event_loop(loop)
self.Client = MQTTClient(client_id=self.client_id, config={'check_hostname':False})
loop.run_until_complete(self.start_helper_bot())
loop.run_until_complete(self.get_msg())
loop.run_forever()
except Exception as e:
logging.error('Helperbot: {}'.format(e))
async def start_helper_bot(self):
@ -56,8 +58,11 @@ class MQTTHelperBot():
('iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+',QOS_0),
('iot/p2p/+',QOS_0)
])
except hbmqtt.client.ClientException as ce:
logging.exception("Client exception: %s" % ce)
except Exception as e:
logging.error('Helperbot: {}'.format(e))
#except hbmqtt.client.ClientException as ce:
# logging.exception("Client exception: %s" % ce)
async def get_msg(self):
@ -78,48 +83,58 @@ class MQTTHelperBot():
cresp.append({"time": time.time() ,"topic": message.topic,"payload":str(message.data.decode("utf-8"))})
self.command_responses.set(cresp)
logging.debug("MQTT Command Response List Count: %s" %len(cresp))
except hbmqtt.client.ClientException as ce:
logging.error("Client exception: %s" % ce)
except Exception as e:
logging.error('Helperbot: {}'.format(e))
#except hbmqtt.client.ClientException as ce:
# logging.error("Client exception: %s" % ce)
async def wait_for_resp(self, requestid):
t_end = (datetime.now() + timedelta(seconds=10)).timestamp()
while time.time() < t_end:
await asyncio.sleep(0.1)
responses = self.command_responses.get()
if len(responses) > 0:
for msg in responses:
topic = str(msg['topic']).split("/")
if (topic[6] == "helper1" and topic[10] == requestid):
logging.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload']))
if topic[11] == "j":
resppayload = json.loads(msg['payload'])
else:
resppayload = str(msg['payload'])
resp = {
"id": requestid,
"ret": "ok",
"resp": resppayload
}
cresp = self.command_responses.get()
cresp.remove(msg)
self.command_responses.set(cresp)
return resp
try:
t_end = (datetime.now() + timedelta(seconds=10)).timestamp()
while time.time() < t_end:
await asyncio.sleep(0.1)
responses = self.command_responses.get()
if len(responses) > 0:
for msg in responses:
topic = str(msg['topic']).split("/")
if (topic[6] == "helper1" and topic[10] == requestid):
logging.debug('VacBot MQTT Response: Topic: %s Payload: %s' % (msg['topic'], msg['payload']))
if topic[11] == "j":
resppayload = json.loads(msg['payload'])
else:
resppayload = str(msg['payload'])
resp = {
"id": requestid,
"ret": "ok",
"resp": resppayload
}
cresp = self.command_responses.get()
cresp.remove(msg)
self.command_responses.set(cresp)
return resp
return { "id": requestid, "errno": "timeout", "ret": "fail" }
return { "id": requestid, "errno": "timeout", "ret": "fail" }
except Exception as e:
logging.error('Helperbot: {}'.format(e))
async def send_command(self, cmdjson, requestid):
ttopic = "iot/p2p/{}/helper1/bumper/helper1/{}/{}/{}/q/{}/{}".format(cmdjson["cmdName"],
cmdjson["toId"], cmdjson["toType"], cmdjson["toRes"], requestid, cmdjson["payloadType"])
try:
await self.Client.publish(ttopic, str(cmdjson["payload"]).encode(),QOS_0)
except:
logging.exception("Exception at send_command")
ttopic = "iot/p2p/{}/helper1/bumper/helper1/{}/{}/{}/q/{}/{}".format(cmdjson["cmdName"],
cmdjson["toId"], cmdjson["toType"], cmdjson["toRes"], requestid, cmdjson["payloadType"])
try:
await self.Client.publish(ttopic, str(cmdjson["payload"]).encode(),QOS_0)
except:
logging.exception("Exception at send_command")
resp = await self.wait_for_resp(requestid)
resp = await self.wait_for_resp(requestid)
return resp
return resp
except Exception as e:
logging.error('Helperbot: {}'.format(e))
class MQTTServer():
@ -128,24 +143,37 @@ class MQTTServer():
bumper_bots = []
async def broker_coro(self):
broker = hbmqtt.broker.Broker(config=self.default_config)
await broker.start()
try:
broker = hbmqtt.broker.Broker(config=self.default_config)
await broker.start()
except PermissionError as e:
if "bind" in e.strerror:
logging.exception("Error binding mqttserver, exiting. Try using a different hostname or IP.\r\n {}".format(e))
exit(1)
except Exception as e:
logging.exception('MQTTServer: {}'.format(e))
exit(1)
async def active_bot_listing(self):
while True:
await asyncio.sleep(5)
logging.debug('Connected bots: %s' % self.bumper_bots.get())
try:
while True:
await asyncio.sleep(5)
logging.debug('Connected bots: %s' % self.bumper_bots.get())
except Exception as e:
logging.error('MQTTServer: {}'.format(e))
def __init__(self, address, run_async=False, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar):
#The below adds a plugin to the hbmqtt.broker.plugins without having to futz with setup.py
distribution = pkg_resources.Distribution("hbmqtt.broker.plugins")
bumper_plugin = pkg_resources.EntryPoint.parse('bumper = bumper.mqttserver:BumperMQTTServer_Plugin', dist=distribution)
distribution._ep_map = {"hbmqtt.broker.plugins": {"bumper": bumper_plugin}}
pkg_resources.working_set.add(distribution)
self.bumper_bots = bumper_bots
self.bumper_clients = bumper_clients
try:
#The below adds a plugin to the hbmqtt.broker.plugins without having to futz with setup.py
distribution = pkg_resources.Distribution("hbmqtt.broker.plugins")
bumper_plugin = pkg_resources.EntryPoint.parse('bumper = bumper.mqttserver:BumperMQTTServer_Plugin', dist=distribution)
distribution._ep_map = {"hbmqtt.broker.plugins": {"bumper": bumper_plugin}}
pkg_resources.working_set.add(distribution)
self.bumper_bots = bumper_bots
self.bumper_clients = bumper_clients
# Initialize bot server
self.default_config = {
'listeners': {
@ -186,16 +214,23 @@ class MQTTServer():
else:
self.run_server()
except:
logging.exception("Exception")
pass
except Exception as e:
logging.error('MQTTServer: {}'.format(e))
#except:
# logging.exception("Exception")
# pass
def run_server(self, loop):
asyncio.set_event_loop(loop)
loop.run_until_complete(self.broker_coro())
#loop.run_until_complete(self.active_bot_listing())
loop.run_forever()
try:
asyncio.set_event_loop(loop)
loop.run_until_complete(self.broker_coro())
#loop.run_until_complete(self.active_bot_listing())
loop.run_forever()
except Exception as e:
logging.error('MQTTServer: {}'.format(e))
class BumperMQTTServer_Plugin:
def __init__(self, context):
@ -204,73 +239,81 @@ class BumperMQTTServer_Plugin:
self.clients = self.context.config['clients']
except KeyError:
self.context.logger.warning("'clients' section not found in context configuration")
except Exception as e:
logging.error('MQTTServer: {}'.format(e))
async def on_broker_client_connected(self, client_id):
try:
logging.debug('Bumper Connection: %s connected' % client_id)
connected_bots = self.clients['connected_bots'].get()
connected_clients = self.clients['connected_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 connected_bots:
if bot['did'] == newbot.did:
botactive = True
logging.debug('Bumper Connection: %s connected' % client_id)
connected_bots = self.clients['connected_bots'].get()
connected_clients = self.clients['connected_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 connected_bots:
if bot['did'] == newbot.did:
botactive = True
if botactive == False:
connected_bots.append(newbot.asdict())
logging.info("Adding bot to list: {}".format(newbot.asdict()))
if botactive == False:
connected_bots.append(newbot.asdict())
logging.info("Adding bot to list: {}".format(newbot.asdict()))
self.clients['connected_bots'].set(connected_bots)
else:
tmpuserdetail = str(didsplit[1]).split("/")
newuser = bumper.VacBotUser()
newuser.userid = didsplit[0]
newuser.realm = tmpuserdetail[0]
newuser.resource = tmpuserdetail[1]
self.clients['connected_bots'].set(connected_bots)
else:
tmpuserdetail = str(didsplit[1]).split("/")
newuser = bumper.VacBotUser()
newuser.userid = didsplit[0]
newuser.realm = tmpuserdetail[0]
newuser.resource = tmpuserdetail[1]
clientactive = False
for client in connected_clients:
if client['userid'] == newuser.userid:
clientactive = True
clientactive = False
for client in connected_clients:
if client['userid'] == newuser.userid:
clientactive = True
if clientactive == False:
connected_clients.append(newuser.asdict())
logging.info("Adding client to list: {}".format(newuser.asdict()))
if clientactive == False:
connected_clients.append(newuser.asdict())
logging.info("Adding client to list: {}".format(newuser.asdict()))
self.clients['connected_clients'].set(connected_clients)
self.clients['connected_clients'].set(connected_clients)
logging.debug('Connected Bots: %s' %self.clients['connected_bots'].get())
logging.debug('Connected Clients: %s' %self.clients['connected_clients'].get())
logging.debug('Connected Bots: %s' %self.clients['connected_bots'].get())
logging.debug('Connected Clients: %s' %self.clients['connected_clients'].get())
except Exception as e:
logging.error('MQTTServer: {}'.format(e))
async def on_broker_client_disconnected(self, client_id):
logging.debug('Bumper Connection: %s disconnected' % client_id)
connected_bots = self.clients['connected_bots'].get()
connected_clients = self.clients['connected_clients'].get()
didsplit = str(client_id).split("@")
#If the did is in the list, remove it
for bot in connected_bots:
if didsplit[0] == bot['did']:
logging.info("Removing bot from list: {}".format(bot['did']))
connected_bots.remove(bot)
self.clients['connected_bots'].set(connected_bots)
try:
logging.debug('Bumper Connection: %s disconnected' % client_id)
connected_bots = self.clients['connected_bots'].get()
connected_clients = self.clients['connected_clients'].get()
didsplit = str(client_id).split("@")
#If the did is in the list, remove it
for bot in connected_bots:
if didsplit[0] == bot['did']:
logging.info("Removing bot from list: {}".format(bot['did']))
connected_bots.remove(bot)
self.clients['connected_bots'].set(connected_bots)
logging.debug('Connected Bots: %s' %self.clients['connected_bots'].get())
logging.debug('Connected Bots: %s' %self.clients['connected_bots'].get())
for client in connected_clients:
if didsplit[0] == client['userid']:
logging.info("Removing client from list: {}".format(client['userid']))
connected_clients.remove(client)
self.clients['connected_clients'].set(connected_clients)
for client in connected_clients:
if didsplit[0] == client['userid']:
logging.info("Removing client from list: {}".format(client['userid']))
connected_clients.remove(client)
self.clients['connected_clients'].set(connected_clients)
logging.debug('Connected Clients: %s' %self.clients['connected_clients'].get())
logging.debug('Connected Clients: %s' %self.clients['connected_clients'].get())
except Exception as e:
logging.error('MQTTServer: {}'.format(e))

View file

@ -32,10 +32,15 @@ class XMPPServer():
client.start()
self.clients.append(client)
self.socket.close()
except PermissionError as e:
if "bind" in e.strerror:
logging.exception("Error binding xmppserver, exiting. Try using a different hostname or IP.\r\n {}".format(e))
exit(1)
except Exception as e:
logging.error('XMPPServer: {}'.format(e))
logging.exception('XMPPServer: {}'.format(e))
exit(1)
except KeyboardInterrupt:
logging.debug('XMPPServer: Keyboard interrupt')
logging.exception('XMPPServer: Keyboard interrupt')
finally:
self.disconnect()
logging.info('XMPPServer: bye')