Beginning bumper's new journey #4

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

View file

@ -6,5 +6,6 @@ name = "pypi"
[packages] [packages]
hbmqtt = "*" hbmqtt = "*"
aiohttp = "*" aiohttp = "*"
black = "*"
[dev-packages] [dev-packages]

109
bumper.py
View file

@ -9,57 +9,89 @@ import platform
def main(): def main():
args = sys.argv 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")
if platform.system() == "Darwin": #If a Mac, use 0.0.0.0 for listening 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")
if platform.system() == "Darwin": # If a Mac, use 0.0.0.0 for listening
listen_host = "0.0.0.0" listen_host = "0.0.0.0"
else: else:
listen_host = socket.gethostbyname(socket.gethostname()) listen_host = socket.gethostbyname(socket.gethostname())
#listen_host = "localhost" #Try this if the above doesn't work # listen_host = "localhost" #Try this if the above doesn't work
conf_address_443 = (listen_host, 443) conf_address_443 = (listen_host, 443)
conf_address_8007 = (listen_host, 8007) conf_address_8007 = (listen_host, 8007)
xmpp_address = (listen_host, 5223) xmpp_address = (listen_host, 5223)
mqtt_address = (listen_host, 8883) mqtt_address = (listen_host, 8883)
xmpp_server = bumper.XMPPServer(xmpp_address, bumper_users=bumper.bumper_users_var, bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var) xmpp_server = bumper.XMPPServer(
mqtt_server = bumper.MQTTServer(mqtt_address,bumper_users=bumper.bumper_users_var, bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var) xmpp_address,
mqtt_helperbot = bumper.MQTTHelperBot(mqtt_address, bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var) bumper_users=bumper.bumper_users_var,
conf_server = bumper.ConfServer(conf_address_443, usessl=True,bumper_users=bumper.bumper_users_var, bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var,helperbot=mqtt_helperbot) bumper_bots=bumper.bumper_bots_var,
conf_server_2 = bumper.ConfServer(conf_address_8007, usessl=False,bumper_users=bumper.bumper_users_var, bumper_bots=bumper.bumper_bots_var,bumper_clients=bumper.bumper_clients_var, helperbot=mqtt_helperbot) bumper_clients=bumper.bumper_clients_var,
)
#add user mqtt_server = bumper.MQTTServer(
# users = bumper.bumper_users_var.get() mqtt_address,
bumper_users=bumper.bumper_users_var,
bumper_bots=bumper.bumper_bots_var,
bumper_clients=bumper.bumper_clients_var,
)
mqtt_helperbot = bumper.MQTTHelperBot(
mqtt_address,
bumper_bots=bumper.bumper_bots_var,
bumper_clients=bumper.bumper_clients_var,
)
conf_server = bumper.ConfServer(
conf_address_443,
usessl=True,
bumper_users=bumper.bumper_users_var,
bumper_bots=bumper.bumper_bots_var,
bumper_clients=bumper.bumper_clients_var,
helperbot=mqtt_helperbot,
)
conf_server_2 = bumper.ConfServer(
conf_address_8007,
usessl=False,
bumper_users=bumper.bumper_users_var,
bumper_bots=bumper.bumper_bots_var,
bumper_clients=bumper.bumper_clients_var,
helperbot=mqtt_helperbot,
)
# add user
# users = bumper.bumper_users_var.get()
# user1 = bumper.BumperUser('user1') # user1 = bumper.BumperUser('user1')
# user1.add_device('devid') # user1.add_device('devid')
# user1.add_bot('bot_did') # user1.add_bot('bot_did')
# users.append(user1) # users.append(user1)
# bumper.bumper_users_var.set(users) # bumper.bumper_users_var.set(users)
# start xmpp server on port 5223 (sync) # start xmpp server on port 5223 (sync)
xmpp_server.run(run_async=True) #Start in new thread xmpp_server.run(run_async=True) # Start in new thread
# start mqtt server on port 8883 (async) # start mqtt server on port 8883 (async)
mqtt_server.run(run_async=True) #Start in new thread mqtt_server.run(run_async=True) # Start in new thread
time.sleep(1.5) #Wait for broker startup
# start mqtt_helperbot (async) time.sleep(1.5) # Wait for broker startup
mqtt_helperbot.run(run_async=True) #Start in new thread
# start conf server on port 443 (async) - Used for most https calls # start mqtt_helperbot (async)
conf_server.run(run_async=True) #Start in new thread mqtt_helperbot.run(run_async=True) # Start in new thread
# start conf server on port 8007 (async) - Used for a load balancer request # start conf server on port 443 (async) - Used for most https calls
conf_server_2.run(run_async=True) #Start in new thread conf_server.run(run_async=True) # Start in new thread
# start conf server on port 8007 (async) - Used for a load balancer request
conf_server_2.run(run_async=True) # Start in new thread
while True: while True:
try: try:
@ -72,13 +104,14 @@ def main():
# if uid != "": # if uid != "":
# xmpp_server.remove_client_byuid(uid) #Remove clients from xmpp server # xmpp_server.remove_client_byuid(uid) #Remove clients from xmpp server
# remove_clients.remove(uid) # remove_clients.remove(uid)
# bumper.bumper_removeclients_var.set(remove_clients) # bumper.bumper_removeclients_var.set(remove_clients)
except KeyboardInterrupt: except KeyboardInterrupt:
bumper.bumperlog.info("Bumper Exiting - Keyboard Interrupt") bumper.bumperlog.info("Bumper Exiting - Keyboard Interrupt")
print("Bumper Exiting") print("Bumper Exiting")
exit(1) exit(1)
if __name__ == "__main__": if __name__ == "__main__":
main() main()

View file

@ -10,39 +10,40 @@ import time
import logging import logging
from base64 import b64decode, b64encode from base64 import b64decode, b64encode
bumper_users_var = contextvars.ContextVar('bumper_users', default=[]) bumper_users_var = contextvars.ContextVar("bumper_users", default=[])
bumper_clients_var = contextvars.ContextVar('bumper_clients', default=[]) bumper_clients_var = contextvars.ContextVar("bumper_clients", default=[])
bumper_bots_var = contextvars.ContextVar('bumper_bots', default=[]) bumper_bots_var = contextvars.ContextVar("bumper_bots", default=[])
ca_cert = './certs/CA/cacert.pem' ca_cert = "./certs/CA/cacert.pem"
server_cert = './certs/cert.pem' server_cert = "./certs/cert.pem"
server_key = './certs/key.pem' server_key = "./certs/key.pem"
use_auth = False use_auth = False
#Logs # Logs
bumperlog = logging.getLogger("bumper") bumperlog = logging.getLogger("bumper")
confserverlog = logging.getLogger("confserver") confserverlog = logging.getLogger("confserver")
#Override the logging level # Override the logging level
#confserverlog.setLevel(logging.INFO) # confserverlog.setLevel(logging.INFO)
mqttserverlog = logging.getLogger("mqttserver") mqttserverlog = logging.getLogger("mqttserver")
#Override the logging level # Override the logging level
#mqttserverlog.setLevel(logging.INFO) # mqttserverlog.setLevel(logging.INFO)
helperbotlog = logging.getLogger("helperbot") helperbotlog = logging.getLogger("helperbot")
#Override the logging level # Override the logging level
#helperbotlog.setLevel(logging.INFO) # helperbotlog.setLevel(logging.INFO)
xmppserverlog = logging.getLogger("xmppserver") xmppserverlog = logging.getLogger("xmppserver")
#Override the logging level # Override the logging level
#xmppserverlog.setLevel(logging.INFO) # xmppserverlog.setLevel(logging.INFO)
def get_milli_time(timetoconvert): def get_milli_time(timetoconvert):
return int(round(timetoconvert * 1000)) return int(round(timetoconvert * 1000))
class BumperUser(object): class BumperUser(object):
def __init__(self,userid=""): def __init__(self, userid=""):
self.userid = userid self.userid = userid
self.devices = [] self.devices = []
self.tokens = [] self.tokens = []
self.authcodes = [] self.authcodes = []
self.bots = [] self.bots = []
@ -55,7 +56,6 @@ class BumperUser(object):
if devid in self.devices: if devid in self.devices:
self.devices.remove(devid) self.devices.remove(devid)
def add_token(self, token): def add_token(self, token):
if not token in self.tokens: if not token in self.tokens:
self.tokens.append(token) self.tokens.append(token)
@ -80,8 +80,17 @@ class BumperUser(object):
if botdid in self.bots: if botdid in self.bots:
self.bots.remove(botdid) self.bots.remove(botdid)
class VacBotDevice(object): class VacBotDevice(object):
def __init__(self,did="", vac_bot_device_class="",resource="" , name="", nick="", company="eco-ng"): def __init__(
self,
did="",
vac_bot_device_class="",
resource="",
name="",
nick="",
company="eco-ng",
):
self.vac_bot_device_class = vac_bot_device_class self.vac_bot_device_class = vac_bot_device_class
self.company = company self.company = company
self.did = did self.did = did
@ -92,11 +101,18 @@ class VacBotDevice(object):
self.xmpp_connection = False self.xmpp_connection = False
def asdict(self): def asdict(self):
return {"class": self.vac_bot_device_class, "company": self.company, return {
"did": self.did, "name": self.name, "nick": self.nick, "resource": self.resource} "class": self.vac_bot_device_class,
"company": self.company,
"did": self.did,
"name": self.name,
"nick": self.nick,
"resource": self.resource,
}
class VacBotClient(object): class VacBotClient(object):
def __init__(self,userid="",realm="",token=""): def __init__(self, userid="", realm="", token=""):
self.userid = userid self.userid = userid
self.realm = realm self.realm = realm
self.resource = token self.resource = token
@ -104,23 +120,25 @@ class VacBotClient(object):
self.xmpp_connection = False self.xmpp_connection = False
def asdict(self): def asdict(self):
return {"userid": self.userid,"realm": self.realm,"resource": self.resource} return {"userid": self.userid, "realm": self.realm, "resource": self.resource}
def check_authcode(uid, authcode): def check_authcode(uid, authcode):
users = bumper_users_var.get() users = bumper_users_var.get()
for user in users: for user in users:
if uid == "fuid_{}".format(user.userid) and authcode in user.authcodes: if uid == "fuid_{}".format(user.userid) and authcode in user.authcodes:
return True return True
return False return False
def add_bot(sn, did, devclass, resource): def add_bot(sn, did, devclass, resource):
newbot = VacBotDevice() newbot = VacBotDevice()
newbot.did = did newbot.did = did
newbot.name = sn newbot.name = sn
newbot.vac_bot_device_class = devclass newbot.vac_bot_device_class = devclass
newbot.resource = resource newbot.resource = resource
bots = bumper_bots_var.get() bots = bumper_bots_var.get()
existingbot = False existingbot = False
@ -128,26 +146,27 @@ def add_bot(sn, did, devclass, resource):
if bot.did == newbot.did: if bot.did == newbot.did:
existingbot = True existingbot = True
if existingbot == False: if existingbot == False:
bots.append(newbot) bots.append(newbot)
bumperlog.info("new bot added SN: {} DID: {}".format(newbot.name, newbot.did)) bumperlog.info("new bot added SN: {} DID: {}".format(newbot.name, newbot.did))
bumper_bots_var.set(bots) bumper_bots_var.set(bots)
def add_client(userid, realm, resource): def add_client(userid, realm, resource):
newclient = VacBotClient() newclient = VacBotClient()
newclient.userid = userid newclient.userid = userid
newclient.realm = realm newclient.realm = realm
newclient.resource = resource newclient.resource = resource
clients = bumper_clients_var.get()
clients = bumper_clients_var.get()
existingclient = False existingclient = False
for client in clients: for client in clients:
if client.userid == newclient.userid: if client.userid == newclient.userid:
existingclient = True existingclient = True
if existingclient == False: if existingclient == False:
clients.append(newclient) clients.append(newclient)
bumperlog.info("new client added {}".format(newclient.userid)) bumperlog.info("new client added {}".format(newclient.userid))
bumper_clients_var.set(clients) bumper_clients_var.set(clients)

View file

@ -12,45 +12,62 @@ import contextvars
from aiohttp import web from aiohttp import web
import uuid import uuid
class aiohttp_filter(logging.Filter): class aiohttp_filter(logging.Filter):
def filter(self, record): def filter(self, record):
if record.name == "aiohttp.access" and record.levelno == 20: #Filters aiohttp.access log to switch it from INFO to DEBUG if (
record.name == "aiohttp.access" and record.levelno == 20
): # Filters aiohttp.access log to switch it from INFO to DEBUG
record.levelno = 10 record.levelno = 10
record.levelname = "DEBUG" record.levelname = "DEBUG"
if record.levelno == 10 and logging.getLogger("confserver").getEffectiveLevel() == 10: if (
record.levelno == 10
and logging.getLogger("confserver").getEffectiveLevel() == 10
):
return True return True
else: else:
return False return False
confserverlog = logging.getLogger("confserver") confserverlog = logging.getLogger("confserver")
logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) #Ignore this logger logging.getLogger("asyncio").setLevel(logging.CRITICAL + 1) # Ignore this logger
logging.getLogger("aiohttp.access").addFilter(aiohttp_filter()) 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
def __init__(self, address, usessl=False, bumper_users=contextvars.ContextVar, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar, helperbot=None): def __init__(
self,
address,
usessl=False,
bumper_users=contextvars.ContextVar,
bumper_bots=contextvars.ContextVar,
bumper_clients=contextvars.ContextVar,
helperbot=None,
):
self.bumper_users = bumper_users self.bumper_users = bumper_users
self.bumper_bots = bumper_bots self.bumper_bots = bumper_bots
self.bumper_clients = bumper_clients self.bumper_clients = bumper_clients
self.helperbot = helperbot self.helperbot = helperbot
self.usessl = usessl self.usessl = usessl
self.address = address self.address = address
self.confthread = None self.confthread = None
def run(self, run_async=False): def run(self, run_async=False):
try: try:
if run_async: if run_async:
confserverlog.debug("Starting ConfServer Thread: 1") confserverlog.debug("Starting ConfServer Thread: 1")
self.confthread = Thread(name="ConfServer_{}_Thread".format(self.address[1]),target=self.run_server) self.confthread = Thread(
name="ConfServer_{}_Thread".format(self.address[1]),
target=self.run_server,
)
self.confthread.setDaemon(True) self.confthread.setDaemon(True)
self.confthread.start() self.confthread.start()
else: else:
try: try:
self.run_server() self.run_server()
@ -58,8 +75,7 @@ class ConfServer():
self.disconnect() self.disconnect()
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
def run_server(self): def run_server(self):
logging.info("Starting ConfServer at {}".format(self.address)) logging.info("Starting ConfServer at {}".format(self.address))
@ -68,433 +84,577 @@ class ConfServer():
loop = asyncio.get_event_loop() loop = asyncio.get_event_loop()
except: except:
loop = asyncio.new_event_loop() loop = asyncio.new_event_loop()
try: try:
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:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
async def start_server(self): async def start_server(self):
try: try:
app = web.Application() app = web.Application()
app.add_routes([ app.add_routes(
web.get('', self.handle_base), [
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/login', self.handle_login), web.get("", self.handle_base),
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/checkLogin', self.handle_login), web.get(
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/logout', self.handle_logout), "/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/login",
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/getAuthCode', self.handle_getAuthCode), self.handle_login,
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(
web.get('/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/campaign/homePageAlert', self.handle_homePageAlert), "/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/checkLogin",
self.handle_login,
web.post('/api/users/user.do', self.handle_usersapi), ),
web.get('/api/users/user.do', self.handle_usersapi), web.get(
web.post('/api/pim/product/getProductIotMap', self.handle_getProductIotMap), "/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/logout",
web.post('/api/iot/devmanager.do', self.handle_devmanager_botcommand), self.handle_logout,
),
web.post('/lookup.do', self.handle_lookup), web.get(
]) "/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/getAuthCode",
#Direct register from app: self.handle_getAuthCode,
#/{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/directRegister ),
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.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),
]
)
# Direct register from app:
# /{apiversion}/private/{country}/{language}/{devid}/{apptype}/{appversion}/{devtype}/{aid}/user/directRegister
runner = web.AppRunner(app) runner = web.AppRunner(app)
await runner.setup() await runner.setup()
if self.usessl: if self.usessl:
ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH) ssl_ctx = ssl.create_default_context(ssl.Purpose.CLIENT_AUTH)
ssl_ctx.load_cert_chain(bumper.server_cert,bumper.server_key) 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) site = web.TCPSite(
runner,
else: 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]) site = web.TCPSite(runner, host=self.address[0], port=self.address[1])
await site.start() await site.start()
except PermissionError as e: except PermissionError as e:
if "bind" in e.strerror: if "bind" in e.strerror:
confserverlog.exception("Error binding confserver, exiting. Try using a different hostname or IP - {}".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:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
exit(1) exit(1)
async def handle_base(self, request): async def handle_base(self, request):
try: try:
text = "Bumper!" text = "Bumper!"
return web.json_response(text) return web.json_response(text)
except Exception as e:
confserverlog.exception('{}'.format(e))
async def handle_login(self, request): except Exception as e:
confserverlog.exception("{}".format(e))
async def handle_login(self, request):
try: try:
user_devid = request.match_info.get('devid', "") user_devid = request.match_info.get("devid", "")
countrycode = request.match_info.get('country', "us") countrycode = request.match_info.get("country", "us")
confserverlog.info('client with devid {} attempting login'.format(user_devid)) confserverlog.info(
"client with devid {} attempting login".format(user_devid)
)
if bumper.use_auth: if bumper.use_auth:
if not user_devid == "": #Performing basic "auth" using devid, super insecure if (
users = self.bumper_users.get() not user_devid == ""
): # Performing basic "auth" using devid, super insecure
users = self.bumper_users.get()
for user in users: for user in users:
if user_devid in user.devices: if user_devid in user.devices:
tmpaccesstoken = '' tmpaccesstoken = ""
if 'checkLogin' in request.path: if "checkLogin" in request.path:
if request.query['accessToken'] in user.tokens and request.query['uid'] == "fuid_{}".format(user.userid): if request.query[
tmpaccesstoken = request.query['accessToken'] "accessToken"
] in user.tokens and request.query[
"uid"
] == "fuid_{}".format(
user.userid
):
tmpaccesstoken = request.query["accessToken"]
body = { body = {
"code": bumper.RETURN_API_SUCCESS, "code": bumper.RETURN_API_SUCCESS,
"data": { "data": {
"accessToken": tmpaccesstoken, #Random chars 32 length "accessToken": tmpaccesstoken, # Random chars 32 length
"country": countrycode, "country": countrycode,
"email": "null@null.com", "email": "null@null.com",
"uid": "fuid_{}".format(user.userid), "uid": "fuid_{}".format(user.userid),
"username": "fusername_{}".format(user.userid), "username": "fusername_{}".format(
user.userid
),
}, },
"msg": "操作成功", "msg": "操作成功",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
} }
else: else:
body = { body = {
"code": bumper.ERR_TOKEN_INVALID, "code": bumper.ERR_TOKEN_INVALID,
"data": None, "data": None,
"msg": "当前密码错误", "msg": "当前密码错误",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
} }
else: else:
if tmpaccesstoken == '': if tmpaccesstoken == "":
tmpaccesstoken = uuid.uuid4().hex tmpaccesstoken = uuid.uuid4().hex
user.add_token(tmpaccesstoken) user.add_token(tmpaccesstoken)
body = { body = {
"code": bumper.RETURN_API_SUCCESS, "code": bumper.RETURN_API_SUCCESS,
"data": { "data": {
"accessToken": tmpaccesstoken, #Random chars 32 length "accessToken": tmpaccesstoken, # Random chars 32 length
"country": countrycode, "country": countrycode,
"email": "null@null.com", "email": "null@null.com",
"uid": "fuid_{}".format(user.userid), "uid": "fuid_{}".format(user.userid),
"username": "fusername_{}".format(user.userid), "username": "fusername_{}".format(user.userid),
}, },
"msg": "操作成功", "msg": "操作成功",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
} }
self.bumper_users.set(users) self.bumper_users.set(users)
return web.json_response(body) return web.json_response(body)
body = { body = {
"code": bumper.ERR_USER_NOT_ACTIVATED, "code": bumper.ERR_USER_NOT_ACTIVATED,
"data": None, "data": None,
"msg": "当前密码错误", "msg": "当前密码错误",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
} }
return web.json_response(body) return web.json_response(body)
else: else:
return web.json_response(self._auth_any(user_devid, countrycode, request)) return web.json_response(
self._auth_any(user_devid, countrycode, request)
except Exception as e: )
confserverlog.exception('{}'.format(e))
except Exception as e:
confserverlog.exception("{}".format(e))
def _auth_any(self, devid, country, request): def _auth_any(self, devid, country, request):
try: try:
user_devid = devid user_devid = devid
countrycode = country countrycode = country
tmpaccesstoken = '' tmpaccesstoken = ""
users = self.bumper_users.get() users = self.bumper_users.get()
bots = self.bumper_bots.get() bots = self.bumper_bots.get()
if len(users) > 0: if len(users) > 0:
tmpuser = users[0] tmpuser = users[0]
tmpuser.add_device(user_devid) tmpuser.add_device(user_devid)
else: else:
tmpuser = bumper.BumperUser('tmpuser') tmpuser = bumper.BumperUser("tmpuser")
users.append(tmpuser) users.append(tmpuser)
tmpuser.add_device(user_devid) tmpuser.add_device(user_devid)
for bot in bots: for bot in bots:
tmpuser.add_bot(bot.did) tmpuser.add_bot(bot.did)
if 'checkLogin' in request.path: if "checkLogin" in request.path:
tmpaccesstoken = request.query['accessToken'] tmpaccesstoken = request.query["accessToken"]
tmpuser.add_token(tmpaccesstoken) tmpuser.add_token(tmpaccesstoken)
body = { body = {
"code": bumper.RETURN_API_SUCCESS, "code": bumper.RETURN_API_SUCCESS,
"data": { "data": {
"accessToken": tmpaccesstoken, #Random chars 32 length "accessToken": tmpaccesstoken, # Random chars 32 length
"country": countrycode,
"email": "null@null.com",
"uid": "fuid_{}".format(tmpuser.userid),
"username": "fusername_{}".format(tmpuser.userid),
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time())
}
else:
if tmpaccesstoken == '':
tmpaccesstoken = uuid.uuid4().hex
tmpuser.add_token(tmpaccesstoken)
body = {
"code": bumper.RETURN_API_SUCCESS,
"data": {
"accessToken": tmpaccesstoken, #Random chars 32 length
"country": countrycode, "country": countrycode,
"email": "null@null.com", "email": "null@null.com",
"uid": "fuid_{}".format(tmpuser.userid), "uid": "fuid_{}".format(tmpuser.userid),
"username": "fusername_{}".format(tmpuser.userid), "username": "fusername_{}".format(tmpuser.userid),
}, },
"msg": "操作成功", "msg": "操作成功",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
} }
else:
if tmpaccesstoken == "":
tmpaccesstoken = uuid.uuid4().hex
tmpuser.add_token(tmpaccesstoken)
body = {
"code": bumper.RETURN_API_SUCCESS,
"data": {
"accessToken": tmpaccesstoken, # Random chars 32 length
"country": countrycode,
"email": "null@null.com",
"uid": "fuid_{}".format(tmpuser.userid),
"username": "fusername_{}".format(tmpuser.userid),
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time()),
}
self.bumper_users.set(users) self.bumper_users.set(users)
return body return body
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
async def handle_logout(self, request): async def handle_logout(self, request):
try: try:
user_devid = request.match_info.get('devid', "") user_devid = request.match_info.get("devid", "")
if not user_devid == "": if not user_devid == "":
users = self.bumper_users.get() users = self.bumper_users.get()
for user in users: for user in users:
if user_devid in user.devices: if user_devid in user.devices:
if request.query['uid'] == "fuid_{}".format(user.userid) and request.query['accessToken'] in user.tokens: if (
user.revoke_token(request.query['accessToken']) request.query["uid"] == "fuid_{}".format(user.userid)
and request.query["accessToken"] in user.tokens
):
user.revoke_token(request.query["accessToken"])
self.bumper_users.set(users) self.bumper_users.set(users)
body = {"code": bumper.RETURN_API_SUCCESS,"data": None,"msg": "操作成功", "time": bumper.get_milli_time(time.time())} body = {
"code": bumper.RETURN_API_SUCCESS,
"data": None,
"msg": "操作成功",
"time": bumper.get_milli_time(time.time()),
}
return web.json_response(body) return web.json_response(body)
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
async def handle_getAuthCode(self, request): async def handle_getAuthCode(self, request):
try: try:
user_devid = request.match_info.get('devid', "") user_devid = request.match_info.get("devid", "")
if not user_devid == "": if not user_devid == "":
users = self.bumper_users.get() users = self.bumper_users.get()
if len(users) > 0: if len(users) > 0:
for user in users: for user in users:
if user_devid in user.devices and request.query['accessToken'] in user.tokens: if (
countrycode = request.match_info.get('country', "us") user_devid in user.devices
tmpauthcode = "{}_{}".format(countrycode,uuid.uuid4().hex) and request.query["accessToken"] in user.tokens
user.add_authcode(tmpauthcode) ):
countrycode = request.match_info.get("country", "us")
tmpauthcode = "{}_{}".format(countrycode, uuid.uuid4().hex)
user.add_authcode(tmpauthcode)
body = { body = {
"code": bumper.RETURN_API_SUCCESS, "code": bumper.RETURN_API_SUCCESS,
"data": { "data": {
"authCode": tmpauthcode, "authCode": tmpauthcode,
"ecovacsUid": request.query['uid'] "ecovacsUid": request.query["uid"],
}, },
"msg": "操作成功", "msg": "操作成功",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
} }
self.bumper_users.set(users) self.bumper_users.set(users)
return web.json_response(body) return web.json_response(body)
body = { body = {
"code": bumper.ERR_TOKEN_INVALID, "code": bumper.ERR_TOKEN_INVALID,
"data": None, "data": None,
"msg": "当前密码错误", "msg": "当前密码错误",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
} }
return web.json_response(body) return web.json_response(body)
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
async def handle_checkVersion(self, request): async def handle_checkVersion(self, request):
try: try:
body = { body = {
"code": bumper.RETURN_API_SUCCESS, "code": bumper.RETURN_API_SUCCESS,
"data": { "data": {
"c": None, "c": None,
"img": None, "img": None,
"r": 0, "r": 0,
"t": None, "t": None,
"u": None, "u": None,
"ut": 0, "ut": 0,
"v": None "v": None,
}, },
"msg": "操作成功", "msg": "操作成功",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
} }
return web.json_response(body) return web.json_response(body)
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
async def handle_checkAgreement(self, request): async def handle_checkAgreement(self, request):
try: try:
body = { body = {
"code": bumper.RETURN_API_SUCCESS, "code": bumper.RETURN_API_SUCCESS,
"data": [], "data": [],
"msg": "操作成功", "msg": "操作成功",
"time": bumper.get_milli_time(time.time()) "time": bumper.get_milli_time(time.time()),
}
}
return web.json_response(body)
return web.json_response(body)
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
async def handle_homePageAlert(self, request): async def handle_homePageAlert(self, request):
try: try:
nextAlert = bumper.get_milli_time((datetime.now() + timedelta(hours=12)).timestamp()) nextAlert = bumper.get_milli_time(
(datetime.now() + timedelta(hours=12)).timestamp()
)
body = { body = {
"code": bumper.RETURN_API_SUCCESS, "code": bumper.RETURN_API_SUCCESS,
"data": { "data": {
"clickSchemeUrl": None, "clickSchemeUrl": None,
"clickWebUrl": None, "clickWebUrl": None,
"hasCampaign": "N", "hasCampaign": "N",
"imageUrl": None, "imageUrl": None,
"nextAlertTime": nextAlert, "nextAlertTime": nextAlert,
"serverTime": bumper.get_milli_time(time.time()) "serverTime": bumper.get_milli_time(time.time()),
},
"msg": "操作成功",
"time": bumper.get_milli_time(time.time()),
}
return web.json_response(body)
except Exception as e:
confserverlog.exception("{}".format(e))
async def handle_getProductIotMap(self, request):
try:
body = {
"code": bumper.RETURN_API_SUCCESS,
"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",
},
}, },
"msg": "操作成功", {
"time": bumper.get_milli_time(time.time()) "classid": "02uwxm",
} "product": {
"_id": "5ae1481e7ccd1a0001e1f69e",
return web.json_response(body) "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: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
async def handle_getProductIotMap(self, request): async def handle_usersapi(self, request):
try:
body = {"code":bumper.RETURN_API_SUCCESS,"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:
confserverlog.exception('{}'.format(e))
async def handle_usersapi(self, request):
try: try:
body = {} body = {}
postbody = {} postbody = {}
if request.content_type == "application/x-www-form-urlencoded": if request.content_type == "application/x-www-form-urlencoded":
postbody = await request.post() postbody = await request.post()
else: else:
postbody = json.loads(await request.text()) postbody = json.loads(await request.text())
todo = postbody['todo'] todo = postbody["todo"]
if todo == 'FindBest': if todo == "FindBest":
service = postbody['service'] service = postbody["service"]
if service == 'EcoMsgNew': if service == "EcoMsgNew":
body = {"result":"ok","ip":socket.gethostbyname(socket.gethostname()),"port":5223} body = {
elif service == 'EcoUpdate': "result": "ok",
body = {"result":"ok","ip":"47.88.66.164","port":8005} "ip": socket.gethostbyname(socket.gethostname()),
elif todo == 'loginByItToken': "port": 5223,
}
elif service == "EcoUpdate":
body = {"result": "ok", "ip": "47.88.66.164", "port": 8005}
elif todo == "loginByItToken":
users = self.bumper_users.get() users = self.bumper_users.get()
for user in users: for user in users:
if postbody['userId'] == "fuid_{}".format(user.userid) and postbody['token'] in user.authcodes: if (
postbody["userId"] == "fuid_{}".format(user.userid)
and postbody["token"] in user.authcodes
):
body = { body = {
"resource": postbody["resource"], "resource": postbody["resource"],
"result": "ok", "result": "ok",
"todo": "result", "todo": "result",
"token": postbody["token"], "token": postbody["token"],
"userId": postbody["userId"] "userId": postbody["userId"],
} }
elif todo == 'GetDeviceList': elif todo == "GetDeviceList":
active_bots = self.bumper_bots.get() active_bots = self.bumper_bots.get()
bot_list = [] bot_list = []
for bot in active_bots: for bot in active_bots:
bot_list.append(bot.asdict()) bot_list.append(bot.asdict())
body = { body = {"devices": bot_list, "result": "ok", "todo": "result"}
"devices": bot_list,
"result": "ok",
"todo": "result"
}
elif todo == 'SetDeviceNick': elif todo == "SetDeviceNick":
bots = self.bumper_bots.get() bots = self.bumper_bots.get()
for bot in bots: for bot in bots:
if postbody['did'] == bot.did: if postbody["did"] == bot.did:
bot.nick = postbody['nick'] bot.nick = postbody["nick"]
self.bumper_bots.set(bots) self.bumper_bots.set(bots)
body = { body = {"result": "ok", "todo": "result"}
"result": "ok",
"todo": "result",
}
confserverlog.debug("\r\n POST: {} \r\n Response: {}".format(postbody,body)) confserverlog.debug(
return web.json_response(body) "\r\n POST: {} \r\n Response: {}".format(postbody, body)
)
return web.json_response(body)
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
async def handle_lookup(self, request): async def handle_lookup(self, request):
try: try:
body = {} body = {}
postbody = {} postbody = {}
if request.content_type == "application/x-www-form-urlencoded": if request.content_type == "application/x-www-form-urlencoded":
postbody = await request.post() postbody = await request.post()
else: else:
postbody = json.loads(await request.text()) postbody = json.loads(await request.text())
confserverlog.debug(postbody) confserverlog.debug(postbody)
todo = postbody['todo'] todo = postbody["todo"]
if todo == 'FindBest': if todo == "FindBest":
service = postbody['service'] service = postbody["service"]
if service == 'EcoMsgNew': if service == "EcoMsgNew":
body = {"result":"ok","ip":socket.gethostbyname(socket.gethostname()),"port":5223} body = {
elif service == 'EcoUpdate': "result": "ok",
body = {"result":"ok","ip":"47.88.66.164","port":8005} "ip": socket.gethostbyname(socket.gethostname()),
"port": 5223,
confserverlog.debug("\r\n POST: {} \r\n Response: {}".format(postbody,body)) }
return web.json_response(body) elif service == "EcoUpdate":
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)
except Exception as e: except Exception as e:
confserverlog.exception('{}'.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())
randomid = ''.join(random.sample(string.ascii_letters,6)) randomid = "".join(random.sample(string.ascii_letters, 6))
bots = self.bumper_bots.get() bots = self.bumper_bots.get()
for bot in bots: for bot in bots:
if bot.did == json_body['toId'] and bot.mqtt_connection == True: if bot.did == json_body["toId"] and bot.mqtt_connection == True:
retcmd = await self.helperbot.send_command(json_body, randomid) retcmd = await self.helperbot.send_command(json_body, randomid)
body = retcmd body = retcmd
confserverlog.debug("\r\n POST: {} \r\n Response: {}".format(json_body,body)) confserverlog.debug(
return web.json_response(body) "\r\n POST: {} \r\n Response: {}".format(json_body, body)
)
return web.json_response(body)
else: else:
confserverlog.error("No bots with DID: {} connected to MQTT".format(json_body['toId'])) confserverlog.error(
body = { "id": randomid, "errno": bumper.ERR_COMMON, "ret": "fail" } "No bots with DID: {} connected to MQTT".format(
return web.json_response(body) json_body["toId"]
)
)
body = {"id": randomid, "errno": bumper.ERR_COMMON, "ret": "fail"}
return web.json_response(body)
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))
def disconnect(self): def disconnect(self):
try: try:
confserverlog.info('shutting down') confserverlog.info("shutting down")
if(self.run_async): if self.run_async:
self.confthread.join() self.confthread.join()
else: else:
self.confthread.disconnect() self.confthread.disconnect()
except Exception as e: except Exception as e:
confserverlog.exception('{}'.format(e)) confserverlog.exception("{}".format(e))

View file

@ -19,282 +19,335 @@ from datetime import datetime, timedelta
helperbotlog = logging.getLogger("helperbot") helperbotlog = logging.getLogger("helperbot")
mqttserverlog = logging.getLogger("mqttserver") mqttserverlog = logging.getLogger("mqttserver")
logging.getLogger("transitions").setLevel(logging.CRITICAL + 1) #Ignore this logger logging.getLogger("transitions").setLevel(logging.CRITICAL + 1) # Ignore this logger
logging.getLogger("passlib").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.broker").setLevel(
logging.getLogger("hbmqtt.mqtt.protocol").setLevel(logging.CRITICAL + 1) #Ignore this logger logging.CRITICAL + 1
logging.getLogger("hbmqtt.client").setLevel(logging.CRITICAL + 1) #Ignore this logger ) # 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, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar):
def __init__(
self,
address,
bumper_bots=contextvars.ContextVar,
bumper_clients=contextvars.ContextVar,
):
self.address = address self.address = address
self.client_id = "helper1@bumper/helper1" self.client_id = "helper1@bumper/helper1"
self.command_responses = contextvars.ContextVar('command_responses', default=[]) self.command_responses = contextvars.ContextVar("command_responses", default=[])
self.helperthread = None self.helperthread = None
def run(self, run_async=False): def run(self, run_async=False):
if run_async: if run_async:
hloop = asyncio.new_event_loop() hloop = asyncio.new_event_loop()
helperbotlog.debug("Starting MQTT HelperBot Thread: 1") helperbotlog.debug("Starting MQTT HelperBot Thread: 1")
self.helperthread = Thread(name="MQTTHelperBot_Thread",target=self.run_helperbot, args=(hloop,)) self.helperthread = Thread(
self.helperthread.setDaemon(True) name="MQTTHelperBot_Thread", target=self.run_helperbot, args=(hloop,)
self.helperthread.start() )
self.helperthread.setDaemon(True)
self.helperthread.start()
else: else:
self.run_helperbot() self.run_helperbot()
def run_helperbot(self, loop):
def run_helperbot(self, loop):
logging.info("Starting MQTT HelperBot") logging.info("Starting MQTT HelperBot")
print("Starting MQTT HelperBot") print("Starting MQTT HelperBot")
try: try:
asyncio.set_event_loop(loop) asyncio.set_event_loop(loop)
self.Client = MQTTClient(client_id=self.client_id, config={'check_hostname':False}) self.Client = MQTTClient(
loop.run_until_complete(self.start_helper_bot()) client_id=self.client_id, config={"check_hostname": False}
loop.run_until_complete(self.get_msg()) )
loop.run_forever() loop.run_until_complete(self.start_helper_bot())
loop.run_until_complete(self.get_msg())
loop.run_forever()
except Exception as e: except Exception as e:
helperbotlog.exception('{}'.format(e)) helperbotlog.exception("{}".format(e))
async def start_helper_bot(self): async def start_helper_bot(self):
try: try:
await self.Client.connect('mqtts://{}:{}/'.format(self.address[0], self.address[1]), cafile=bumper.ca_cert) await self.Client.connect(
await self.Client.subscribe([ "mqtts://{}:{}/".format(self.address[0], self.address[1]),
('iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+',QOS_0), cafile=bumper.ca_cert,
('iot/p2p/+',QOS_0) )
]) await self.Client.subscribe(
[
("iot/p2p/+/+/+/+/helper1/bumper/helper1/+/+/+", QOS_0),
("iot/p2p/+", QOS_0),
]
)
except Exception as e: except Exception as e:
helperbotlog.exception('{}'.format(e)) helperbotlog.exception("{}".format(e))
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()
#helperbotlog.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()
if (str(message.topic).split("/")[6] == "helper1"): if str(message.topic).split("/")[6] == "helper1":
cresp.append({"time": time.time() ,"topic": message.topic,"payload":str(message.data.decode("utf-8"))}) cresp.append(
{
#Cleanup "expired messages" > 60 seconds from time "time": time.time(),
for msg in cresp: "topic": message.topic,
expire_time = (datetime.fromtimestamp(msg['time']) + timedelta(seconds=10)).timestamp() "payload": str(message.data.decode("utf-8")),
}
)
# Cleanup "expired messages" > 60 seconds from time
for msg in cresp:
expire_time = (
datetime.fromtimestamp(msg["time"]) + timedelta(seconds=10)
).timestamp()
if time.time() > expire_time: if time.time() > expire_time:
#helperbotlog.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)
self.command_responses.set(cresp)
#helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
except Exception as e:
helperbotlog.exception('{}'.format(e))
async def wait_for_resp(self, requestid): self.command_responses.set(cresp)
# helperbotlog.debug("MQTT Command Response List Count: %s" %len(cresp))
except Exception as e:
helperbotlog.exception("{}".format(e))
async def wait_for_resp(self, requestid):
try: try:
t_end = (datetime.now() + timedelta(seconds=10)).timestamp() t_end = (datetime.now() + timedelta(seconds=10)).timestamp()
while time.time() < t_end: while time.time() < t_end:
await asyncio.sleep(0.1) await asyncio.sleep(0.1)
responses = self.command_responses.get() responses = self.command_responses.get()
if len(responses) > 0: if len(responses) > 0:
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:
#helperbotlog.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:
resppayload = str(msg['payload']) resppayload = str(msg["payload"])
resp = { resp = {"id": requestid, "ret": "ok", "resp": resppayload}
"id": requestid,
"ret": "ok",
"resp": resppayload
}
cresp = self.command_responses.get() cresp = self.command_responses.get()
cresp.remove(msg) cresp.remove(msg)
self.command_responses.set(cresp) self.command_responses.set(cresp)
return resp return resp
return {"id": requestid, "errno": "timeout", "ret": "fail"}
return { "id": requestid, "errno": "timeout", "ret": "fail" }
except asyncio.CancelledError as e: except asyncio.CancelledError as e:
helperbotlog.debug('wait_for_resp cancelled by asyncio') helperbotlog.debug("wait_for_resp cancelled by asyncio")
except Exception as e: except Exception as e:
helperbotlog.exception('{}'.format(e)) helperbotlog.exception("{}".format(e))
async def send_command(self, cmdjson, requestid): async def send_command(self, cmdjson, requestid):
try: try:
ttopic = "iot/p2p/{}/helper1/bumper/helper1/{}/{}/{}/q/{}/{}".format(cmdjson["cmdName"], ttopic = "iot/p2p/{}/helper1/bumper/helper1/{}/{}/{}/q/{}/{}".format(
cmdjson["toId"], cmdjson["toType"], cmdjson["toRes"], requestid, cmdjson["payloadType"]) cmdjson["cmdName"],
try: cmdjson["toId"],
await self.Client.publish(ttopic, str(cmdjson["payload"]).encode(),QOS_0) cmdjson["toType"],
cmdjson["toRes"],
requestid,
cmdjson["payloadType"],
)
try:
await self.Client.publish(
ttopic, str(cmdjson["payload"]).encode(), QOS_0
)
except Exception as e: except Exception as e:
helperbotlog.exception("{}".format(e)) 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:
helperbotlog.exception('{}'.format(e)) helperbotlog.exception("{}".format(e))
class MQTTServer():
default_config = {} class MQTTServer:
default_config = {}
bumper_users = [] bumper_users = []
bumper_clients = [] bumper_clients = []
bumper_bots = [] bumper_bots = []
async def broker_coro(self):
async def broker_coro(self): try:
try:
broker = hbmqtt.broker.Broker(config=self.default_config) broker = hbmqtt.broker.Broker(config=self.default_config)
await broker.start() await broker.start()
except PermissionError as e: except PermissionError as e:
if "bind" in e.strerror: if "bind" in e.strerror:
mqttserverlog.exception("Error binding mqttserver, exiting. Try using a different hostname or IP - {}".format(e)) mqttserverlog.exception(
exit(1) "Error binding mqttserver, exiting. Try using a different hostname or IP - {}".format(
e
except Exception as e: )
mqttserverlog.exception('{}'.format(e)) )
exit(1) exit(1)
async def active_bot_listing(self): except Exception as e:
try: mqttserverlog.exception("{}".format(e))
exit(1)
async def active_bot_listing(self):
try:
while True: while True:
await asyncio.sleep(5) await asyncio.sleep(5)
mqttserverlog.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:
mqttserverlog.exception('{}'.format(e)) mqttserverlog.exception("{}".format(e))
def __init__(self, address,bumper_users=contextvars.ContextVar, bumper_bots=contextvars.ContextVar, bumper_clients=contextvars.ContextVar): def __init__(
self,
address,
bumper_users=contextvars.ContextVar,
bumper_bots=contextvars.ContextVar,
bumper_clients=contextvars.ContextVar,
):
try: try:
self.bumper_users = bumper_users self.bumper_users = bumper_users
self.bumper_bots = bumper_bots self.bumper_bots = bumper_bots
self.bumper_clients = bumper_clients self.bumper_clients = bumper_clients
self.mqttserverthread = None self.mqttserverthread = None
self.address = address self.address = address
#The below adds a plugin to the hbmqtt.broker.plugins without having to futz with setup.py # The below adds a plugin to the hbmqtt.broker.plugins without having to futz with setup.py
distribution = pkg_resources.Distribution("hbmqtt.broker.plugins") distribution = pkg_resources.Distribution("hbmqtt.broker.plugins")
bumper_plugin = pkg_resources.EntryPoint.parse('bumper = bumper.mqttserver:BumperMQTTServer_Plugin', dist=distribution) bumper_plugin = pkg_resources.EntryPoint.parse(
"bumper = bumper.mqttserver:BumperMQTTServer_Plugin", dist=distribution
)
distribution._ep_map = {"hbmqtt.broker.plugins": {"bumper": bumper_plugin}} distribution._ep_map = {"hbmqtt.broker.plugins": {"bumper": bumper_plugin}}
pkg_resources.working_set.add(distribution) pkg_resources.working_set.add(distribution)
# Initialize bot server # Initialize bot server
self.default_config = { self.default_config = {
'listeners': { "listeners": {
'default': { "default": {"type": "tcp"},
'type': 'tcp', "tls1": {
}, "bind": "{}:{}".format(address[0], address[1]),
'tls1': { "ssl": "on",
'bind': "{}:{}".format(address[0], address[1]), "certfile": bumper.server_cert,
'ssl': 'on', "keyfile": bumper.server_key,
'certfile': bumper.server_cert,
'keyfile': bumper.server_key,
}, },
}, },
'sys_interval': 10, "sys_interval": 10,
'auth': { "auth": {
'allow-anonymous': False, "allow-anonymous": False,
'password-file': os.path.join(os.path.dirname(os.path.realpath(__file__)), "passwd"), "password-file": os.path.join(
'plugins': [ os.path.dirname(os.path.realpath(__file__)), "passwd"
'bumper' #No plugins == no auth ),
] "plugins": ["bumper"], # No plugins == no auth
}, },
'topic-check': { "topic-check": {"enabled": False},
'enabled': False "bumper": {
"bumper_users": self.bumper_users,
"bumper_bots": self.bumper_bots,
"bumper_clients": self.bumper_clients,
}, },
'bumper':{
'bumper_users' : self.bumper_users,
'bumper_bots': self.bumper_bots,
'bumper_clients': self.bumper_clients,
},
} }
except Exception as e:
mqttserverlog.exception('{}'.format(e))
def run(self, run_async=False,): except Exception as e:
mqttserverlog.exception("{}".format(e))
def run(self, run_async=False):
if run_async: if run_async:
sloop = asyncio.new_event_loop() sloop = asyncio.new_event_loop()
mqttserverlog.debug("Starting MQTTServer Thread: 1") mqttserverlog.debug("Starting MQTTServer Thread: 1")
self.mqttserverthread = Thread(name="MQTTServer_Thread",target=self.run_server, args=(sloop,)) self.mqttserverthread = Thread(
self.mqttserverthread.setDaemon(True) name="MQTTServer_Thread", target=self.run_server, args=(sloop,)
self.mqttserverthread.start() )
self.mqttserverthread.setDaemon(True)
self.mqttserverthread.start()
else: else:
self.run_server() self.run_server()
def run_server(self, loop): def run_server(self, loop):
logging.info("Starting MQTT Server at {}".format(self.address)) logging.info("Starting MQTT Server at {}".format(self.address))
print("Starting MQTT Server at {}".format(self.address)) print("Starting MQTT Server at {}".format(self.address))
try: try:
asyncio.set_event_loop(loop) asyncio.set_event_loop(loop)
loop.run_until_complete(self.broker_coro()) loop.run_until_complete(self.broker_coro())
#loop.run_until_complete(self.active_bot_listing()) # loop.run_until_complete(self.active_bot_listing())
loop.run_forever() loop.run_forever()
except Exception as e:
mqttserverlog.exception('{}'.format(e))
class BumperMQTTServer_Plugin: except Exception as e:
mqttserverlog.exception("{}".format(e))
class BumperMQTTServer_Plugin:
def __init__(self, context): def __init__(self, context):
self.context = context self.context = context
try: try:
self.bumper_config = self.context.config['bumper'] self.bumper_config = self.context.config["bumper"]
self.auth_config = self.context.config['auth'] self.auth_config = self.context.config["auth"]
except KeyError: except KeyError:
self.context.logger.warning("'bumper' section not found in context configuration") self.context.logger.warning(
"'bumper' section not found in context configuration"
)
except Exception as e: except Exception as e:
mqttserverlog.exception('{}'.format(e)) mqttserverlog.exception("{}".format(e))
async def authenticate(self, *args, **kwargs): async def authenticate(self, *args, **kwargs):
if not self.auth_config: if not self.auth_config:
# auth config section not found # auth config section not found
self.context.logger.warning("'auth' section not found in context configuration") self.context.logger.warning(
"'auth' section not found in context configuration"
)
return False return False
allow_anonymous = self.auth_config.get('allow-anonymous', True) # allow anonymous by default allow_anonymous = self.auth_config.get(
"allow-anonymous", True
) # allow anonymous by default
if allow_anonymous: if allow_anonymous:
authenticated = True authenticated = True
self.context.logger.debug("Authentication success: config allows anonymous") self.context.logger.debug("Authentication success: config allows anonymous")
else: else:
try: try:
bumper_users = self.bumper_config['bumper_users'].get() bumper_users = self.bumper_config["bumper_users"].get()
bumper_bots = self.bumper_config['bumper_bots'].get() bumper_bots = self.bumper_config["bumper_bots"].get()
bumper_clients = self.bumper_config['bumper_clients'].get() bumper_clients = self.bumper_config["bumper_clients"].get()
session = kwargs.get('session', None) session = kwargs.get("session", None)
username = session.username username = session.username
password = session.password password = session.password
client_id = session.client_id client_id = session.client_id
didsplit = str(client_id).split("@") didsplit = str(client_id).split("@")
#If this isn't a fake user (fuid) then add as a bot # 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")): if not (
tmpbotdetail = str(didsplit[1]).split("/") str(didsplit[0]).startswith("fuid")
bumper.add_bot(username, didsplit[0], tmpbotdetail[0], tmpbotdetail[1]) or str(didsplit[0]).startswith("helper")
mqttserverlog.debug("new bot authenticated SN: {} DID: {}".format(username, didsplit[0])) ):
tmpbotdetail = str(didsplit[1]).split("/")
bumper.add_bot(
username, didsplit[0], tmpbotdetail[0], tmpbotdetail[1]
)
mqttserverlog.debug(
"new bot authenticated SN: {} DID: {}".format(
username, didsplit[0]
)
)
authenticated = True authenticated = True
else: else:
tmpclientdetail = str(didsplit[1]).split("/") tmpclientdetail = str(didsplit[1]).split("/")
userid = didsplit[0] userid = didsplit[0]
realm = tmpclientdetail[0] realm = tmpclientdetail[0]
resource = tmpclientdetail[1] resource = tmpclientdetail[1]
if userid == "helper1": if userid == "helper1":
authenticated = True authenticated = True
else: else:
@ -304,61 +357,63 @@ class BumperMQTTServer_Plugin:
elif bumper.use_auth == False: elif bumper.use_auth == False:
auth = True auth = True
if auth: if auth:
bumper.add_client(userid, realm, resource) bumper.add_client(userid, realm, resource)
mqttserverlog.debug("client authenticated {}".format(userid)) mqttserverlog.debug(
"client authenticated {}".format(userid)
)
authenticated = True authenticated = True
else: else:
authenticated = False authenticated = False
except Exception as e: except Exception as e:
mqttserverlog.exception('{}'.format(e)) mqttserverlog.exception("{}".format(e))
authenticated = False authenticated = False
return authenticated return authenticated
async def on_broker_client_connected(self, client_id): async def on_broker_client_connected(self, client_id):
try: try:
bumper_users = self.bumper_config['bumper_users'].get() bumper_users = self.bumper_config["bumper_users"].get()
bumper_bots = self.bumper_config['bumper_bots'].get() bumper_bots = self.bumper_config["bumper_bots"].get()
bumper_clients = self.bumper_config['bumper_clients'].get() bumper_clients = self.bumper_config["bumper_clients"].get()
didsplit = str(client_id).split("@") didsplit = str(client_id).split("@")
for bot in bumper_bots: for bot in bumper_bots:
if didsplit[0] == bot.did: if didsplit[0] == bot.did:
bot.mqtt_connection = True bot.mqtt_connection = True
mqttserverlog.debug("bot connected {}".format(bot.did)) mqttserverlog.debug("bot connected {}".format(bot.did))
self.bumper_config['bumper_bots'].set(bumper_bots) self.bumper_config["bumper_bots"].set(bumper_bots)
for client in bumper_clients: for client in bumper_clients:
if didsplit[0] == client.userid and client.userid != 'helper1': if didsplit[0] == client.userid and client.userid != "helper1":
client.mqtt_connection = True client.mqtt_connection = True
#mqttserverlog.info("client connected {}".format(client.userid)) # mqttserverlog.info("client connected {}".format(client.userid))
self.bumper_config['bumper_clients'].set(bumper_clients) self.bumper_config["bumper_clients"].set(bumper_clients)
except Exception as e: except Exception as e:
mqttserverlog.exception('{}'.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:
bumper_users = self.bumper_config['bumper_users'].get() bumper_users = self.bumper_config["bumper_users"].get()
bumper_bots = self.bumper_config['bumper_bots'].get() bumper_bots = self.bumper_config["bumper_bots"].get()
bumper_clients = self.bumper_config['bumper_clients'].get() bumper_clients = self.bumper_config["bumper_clients"].get()
didsplit = str(client_id).split("@") didsplit = str(client_id).split("@")
for bot in bumper_bots: for bot in bumper_bots:
if didsplit[0] == bot.did: if didsplit[0] == bot.did:
bot.mqtt_connection = False bot.mqtt_connection = False
mqttserverlog.debug("bot disconnected {}".format(bot.did)) mqttserverlog.debug("bot disconnected {}".format(bot.did))
self.bumper_config['bumper_bots'].set(bumper_bots) self.bumper_config["bumper_bots"].set(bumper_bots)
for client in bumper_clients: for client in bumper_clients:
if didsplit[0] == client.userid and client.userid != 'helper1': if didsplit[0] == client.userid and client.userid != "helper1":
client.mqtt_connection = False client.mqtt_connection = False
#mqttserverlog.info("client disconnected {}".format(client.userid)) # mqttserverlog.info("client disconnected {}".format(client.userid))
self.bumper_config['bumper_clients'].set(bumper_clients) self.bumper_config["bumper_clients"].set(bumper_clients)
except Exception as e: except Exception as e:
mqttserverlog.exception('{}'.format(e)) mqttserverlog.exception("{}".format(e))

File diff suppressed because it is too large Load diff

View file

@ -2,19 +2,22 @@
from sucks import * from sucks import *
class BumperVacBot(VacBot): class BumperVacBot(VacBot):
def __init__(self, server_address): def __init__(self, server_address):
self.server_address = server_address self.server_address = server_address
vacuum = { 'did':'none','class':'none' } vacuum = {"did": "none", "class": "none"}
super().__init__('sucks', 'ecouser.net', '', '', vacuum, '') super().__init__("sucks", "ecouser.net", "", "", vacuum, "")
def connect_and_wait_until_ready(self): def connect_and_wait_until_ready(self):
logging.info('connecting') logging.info("connecting")
self.xmpp.connect(self.server_address) self.xmpp.connect(self.server_address)
self.xmpp.process() self.xmpp.process()
self.xmpp.wait_until_ready() self.xmpp.wait_until_ready()
logging.basicConfig(level=logging.DEBUG, format='%(levelname)-8s %(message)s')
server_address = ('xxx.xxx.xxx.xxx', 5223) logging.basicConfig(level=logging.DEBUG, format="%(levelname)-8s %(message)s")
server_address = ("xxx.xxx.xxx.xxx", 5223)
# Initialize # Initialize
vacbot = BumperVacBot(server_address) vacbot = BumperVacBot(server_address)