WIP: Initial MQTT work

WIP: Add initial EcoVacsMQTT client
- Connect and get message

TODO: Parse messages and plumb to events
This commit is contained in:
Brian Martin
2019-01-16 09:02:03 -05:00
parent 3c01cd2678
commit ca7d37c193
3 changed files with 145 additions and 11 deletions
+138 -11
View File
@@ -11,6 +11,11 @@ from sleekxmpp import ClientXMPP, Callback, MatchXPath
from sleekxmpp.xmlstream import ET
from sleekxmpp.exceptions import XMPPError
from paho.mqtt.client import Client as ClientMQTT
from paho.mqtt import publish as MQTTPublish
from paho.mqtt import subscribe as MQTTSubscribe
import ssl
_LOGGER = logging.getLogger(__name__)
# These consts define all of the vocabulary used by this library when presenting various states and components.
@@ -375,14 +380,20 @@ class VacBot():
if vacuum['iot']:
self.iot = EcoVacsIOT(user, domain, resource, secret, continent, vacuum)
self.iot.subscribe_to_ctls(self._handle_ctl)
self.mqtt = EcoVacsMQTT(user, domain, resource, secret, continent, vacuum)
self.mqtt.subscribe_to_ctls(self._handle_ctl)
self.xmpp = EcoVacsXMPP(user, domain, resource, secret, continent, vacuum, server_address )
self.xmpp.subscribe_to_ctls(self._handle_ctl)
else:
self.xmpp = EcoVacsXMPP(user, domain, resource, secret, continent, vacuum, server_address )
self.xmpp.subscribe_to_ctls(self._handle_ctl)
def connect_and_wait_until_ready(self):
self.xmpp.connect_and_wait_until_ready()
self.xmpp.schedule('Ping', 30, lambda: self.send_ping(), repeat=True)
if not self.vacuum['iot']:
self.xmpp.connect_and_wait_until_ready()
self.xmpp.schedule('Ping', 30, lambda: self.send_ping(), repeat=True)
else:
self.mqtt.connect_and_wait_until_ready()
#ToDo identify the best way to handle similar for IOT devices
#self.iot.connect_and_wait_until_ready()
@@ -390,8 +401,11 @@ class VacBot():
if self._monitor:
# Do a first ping, which will also fetch initial statuses if the ping succeeds
self.send_ping()
self.xmpp.schedule('Components', 3600, lambda: self.refresh_components(), repeat=True)
if not self.vacuum['iot']:
self.send_ping()
self.xmpp.schedule('Components', 3600, lambda: self.refresh_components(), repeat=True)
#else:
#TODO: Handle in MQTT?
def _handle_ctl(self, ctl):
method = '_handle_' + ctl['event']
@@ -502,6 +516,7 @@ class VacBot():
if not self.vacuum['iot']:
self.xmpp.send_ping(self._vacuum_address())
else:
self.mqtt.send_ping()
self.xmpp.send_ping(EcoVacsAPI.REALM) #IOT vacuums are using the realm instead
except XMPPError as err:
_LOGGER.warning("Ping did not reach VacBot. Will retry.")
@@ -559,10 +574,13 @@ class VacBot():
def run(self, action):
self.send_command(action)
def disconnect(self, wait=False):
self.xmpp.disconnect(wait=wait)
def disconnect(self, wait=False):
if not self.vacuum['iot']:
self.xmpp.disconnect(wait=wait)
else:
self.mqtt.disconnect()
#This is used by EcoVacsIOT and EcoVacsXMPP for _ctl_to_dict
#This is used by EcoVacsIOT, EcoVacsXMPP, and EcoVacsMQTT for _ctl_to_dict
def RepresentsInt(stringvar):
try:
int(stringvar)
@@ -581,7 +599,6 @@ class EcoVacsIOT():
self.api = EcoVacsAPI
self.api.continent = continent
self.api.meta = {}
#self.add_event_handler("session_start", self.session_start)
self.ctl_subscribers = []
self.ready_flag = Event()
@@ -659,6 +676,115 @@ class EcoVacsIOT():
return result
class EcoVacsMQTT(ClientMQTT):
def __init__(self, user, domain, resource, secret, continent, vacuum, server_address=None ):
ClientMQTT.__init__(self)
self.ctl_subscribers = []
self.user = user
self.domain = str(domain).split(".")[0] #MQTT is using domain without tld
self.resource = resource
self.continent = continent
self.vacuum = vacuum
if server_address is None:
self.hostname = ('mq-{}.ecouser.net'.format(self.continent))
self.port = 8883
else:
saddress = server_address.split(":")
if len(saddress) > 1:
self.hostname = saddress[0]
if RepresentsInt(saddress[1]):
self.port = int(saddress[1])
else:
self.port = 8883
self._client_id = self.user + '@' + self.domain.split(".")[0] + '/' + self.resource
self.username_pw_set(self.user + '@' + self.domain, secret)
self.ready_flag = Event()
def wait_until_ready(self):
self.ready_flag.wait()
def on_connect(self, client, userdata, flags, rc):
if rc != 0:
_LOGGER.error("EcoVacsMQTT error connecting - MQTT Return {}".format(rc))
raise RuntimeError("EcoVacsMQTT error connecting - MQTT Return {}".format(rc))
else:
_LOGGER.debug("Connected MQTT with result code "+str(rc))
_LOGGER.debug("EcoVacsMQTT - Subscribing to all")
self.subscribe('iot/atr/+/' + self.vacuum['did'] + '/' + self.vacuum['class'] + '/' + self.vacuum['resource'] + '/+', qos=0)
self.ready_flag.set()
def on_log(self, client, userdata, level, buf):
_LOGGER.debug("EcoVacsMQTT Log: {} ".format(buf))
def subscribe_to_ctls(self, function):
self.ctl_subscribers.append(function)
def _handle_ctl(self, client, userdata, message):
_LOGGER.debug("EcoVacs MQTT Received Message on Topic: {} - Message: {}".format(message.topic, str(message.payload.decode("utf-8"))))
#the_good_part = message.get_payload()[0][0]
#as_dict = self._ctl_to_dict(the_good_part)
##if as_dict is not None:
# for s in self.ctl_subscribers:
# s(as_dict)
def _ctl_to_dict(self, xml):
result = xml.attrib.copy()
if 'td' not in result:
# This happens for commands with no response data, such as PlaySound
return
result['event'] = result.pop('td')
if xml:
result.update(xml[0].attrib)
for key in result:
if not RepresentsInt(result[key]): #Fix to handle negative int values
result[key] = stringcase.snakecase(result[key])
return result
def send_command(self, xml, recipient):
c = self._wrap_command(xml, recipient)
_LOGGER.debug('Sending command {0}'.format(c))
c.send()
def _wrap_command(self, ctl, recipient):
q = self.make_iq_query(xmlns=u'com:ctl', ito=recipient, ifrom=self._my_address())
q['type'] = 'set'
for child in q.xml:
if child.tag.endswith('query'):
child.append(ctl)
return q
def _my_address(self):
if not self.vacuum['iot']:
return self.user + '@' + self.domain + '/' + self.boundjid.resource
else:
return self.user + '@' + self.domain + '/' + self.resource
def connect_and_wait_until_ready(self):
self._on_log = self.on_log
self._on_message = self._handle_ctl
self._on_connect = self.on_connect
#TODO: This is pretty insecure and accepts any cert, maybe actually check?
ssl_ctx = ssl.create_default_context()
ssl_ctx.check_hostname = False
ssl_ctx.verify_mode = ssl.CERT_NONE
self.tls_set_context(ssl_ctx)
self.tls_insecure_set(True)
self.connect(self.hostname, self.port)
self.loop_start()
self.wait_until_ready()
class EcoVacsXMPP(ClientXMPP):
def __init__(self, user, domain, resource, secret, continent, vacuum, server_address=None ):
@@ -715,7 +841,8 @@ class EcoVacsXMPP(ClientXMPP):
return result
def register_callback(self, kind, function):
def register_callback(self, userdata, message):
self.register_handler(Callback(kind,
MatchXPath('{jabber:client}iq/{com:ctl}query/{com:ctl}ctl[@td="' + kind + '"]'),
function))