From 17b041074cbce21bed40eb49c0becf2593b6c377 Mon Sep 17 00:00:00 2001 From: bittles Date: Tue, 3 Jan 2023 20:29:44 -0500 Subject: [PATCH] code refactoring split xmpp and mqtt vacuums to their own modules --- custom_components/ecovacs/__init__.py | 41 +- custom_components/ecovacs/const.py | 123 ++++++ custom_components/ecovacs/manifest.json | 4 +- custom_components/ecovacs/mqtt_ecovacs.py | 236 +++++++++++ custom_components/ecovacs/sucksbumper.py | 479 +--------------------- custom_components/ecovacs/xmpp_ecovacs.py | 143 +++++++ 6 files changed, 526 insertions(+), 500 deletions(-) create mode 100644 custom_components/ecovacs/const.py create mode 100644 custom_components/ecovacs/mqtt_ecovacs.py create mode 100644 custom_components/ecovacs/xmpp_ecovacs.py diff --git a/custom_components/ecovacs/__init__.py b/custom_components/ecovacs/__init__.py index 6f3d07a..3ab0415 100644 --- a/custom_components/ecovacs/__init__.py +++ b/custom_components/ecovacs/__init__.py @@ -1,14 +1,7 @@ """Support for Ecovacs Deebot vacuums.""" -import logging import random import string - -##import asyncio ## to do - - -#just included the modified sucks in component -from .sucksbumper import EcoVacsAPI, VacBot -import voluptuous as vol +##import asyncio ## to do will need to convert to slixmpp to do this i believe from homeassistant.const import ( CONF_PASSWORD, @@ -21,17 +14,19 @@ from homeassistant.core import HomeAssistant from homeassistant.helpers import discovery import homeassistant.helpers.config_validation as cv from homeassistant.helpers.typing import ConfigType - -_LOGGER = logging.getLogger(__name__) - -DOMAIN = "ecovacs" - -CONF_COUNTRY = "country" -CONF_CONTINENT = "continent" -#bumper config vars -CONF_BUMPER = "bumper" -CONF_BUMPER_SERVER = "bumper_server" -server_address = None +import voluptuous as vol +#just included the modified sucks in component +from .sucksbumper import EcoVacsAPI, VacBot +from .const import ( + ECOVACS_DEVICES, + DOMAIN, + CONF_COUNTRY, + CONF_CONTINENT, + CONF_BUMPER, + CONF_BUMPER_SERVER, + SERVER_ADDRESS, + _LOGGER +) CONFIG_SCHEMA = vol.Schema( { @@ -50,8 +45,6 @@ CONFIG_SCHEMA = vol.Schema( extra=vol.ALLOW_EXTRA, ) -ECOVACS_DEVICES = "ecovacs_devices" - # Generate a random device ID on each bootup ECOVACS_API_DEVICEID = "".join( random.choice(string.ascii_uppercase + string.digits) for _ in range(8) @@ -64,10 +57,10 @@ def setup(hass: HomeAssistant, config: ConfigType) -> bool: hass.data[ECOVACS_DEVICES] = [] # if we're using bumper then define the server address if CONF_BUMPER == True: - server_address = (config[DOMAIN].get(CONF_BUMPER_SERVER), 5223) + SERVER_ADDRESS = (config[DOMAIN].get(CONF_BUMPER_SERVER), 5223) # if not make sure it's null else: - server_address = None + SERVER_ADDRESS = None ecovacs_api = EcoVacsAPI( ECOVACS_API_DEVICEID, @@ -94,7 +87,7 @@ def setup(hass: HomeAssistant, config: ConfigType) -> bool: ecovacs_api.user_access_token, device, config[DOMAIN].get(CONF_CONTINENT).lower(), - server_address, # include server address in class, if it's null shoul be no effect + SERVER_ADDRESS, # include server address in class, if it's null shoul be no effect config[DOMAIN].get(CONF_VERIFY_SSL), # verify ssl or not monitor=True, ) diff --git a/custom_components/ecovacs/const.py b/custom_components/ecovacs/const.py new file mode 100644 index 0000000..2f5cbff --- /dev/null +++ b/custom_components/ecovacs/const.py @@ -0,0 +1,123 @@ +import logging +_LOGGER = logging.getLogger(__name__) + +#ecovacs constants +#init constants +ECOVACS_DEVICES = "ecovacs_devices" +DOMAIN = "ecovacs" +CONF_COUNTRY = "country" +CONF_CONTINENT = "continent" +#bumper config vars +CONF_BUMPER = "bumper" +CONF_BUMPER_SERVER = "bumper_server" +SERVER_ADDRESS = None + +#sucks constants +# These consts define all of the vocabulary used by this library when presenting various states and components. +# Applications implementing this library should import these rather than hard-code the strings, for future-proofing. + +CLEAN_MODE_AUTO = 'auto' +CLEAN_MODE_EDGE = 'edge' +CLEAN_MODE_SPOT = 'spot' +CLEAN_MODE_SPOT_AREA = 'spot_area' +CLEAN_MODE_SINGLE_ROOM = 'single_room' +CLEAN_MODE_STOP = 'stop' + +CLEAN_ACTION_START = 'start' +CLEAN_ACTION_PAUSE = 'pause' +CLEAN_ACTION_RESUME = 'resume' +CLEAN_ACTION_STOP = 'stop' + +FAN_SPEED_NORMAL = 'normal' +FAN_SPEED_HIGH = 'high' + +CHARGE_MODE_RETURN = 'return' +CHARGE_MODE_RETURNING = 'returning' +CHARGE_MODE_CHARGING = 'charging' +CHARGE_MODE_IDLE = 'idle' + +COMPONENT_SIDE_BRUSH = 'side_brush' +COMPONENT_MAIN_BRUSH = 'main_brush' +COMPONENT_FILTER = 'filter' + +VACUUM_STATUS_OFFLINE = 'offline' + +CLEANING_STATES = {CLEAN_MODE_AUTO, CLEAN_MODE_EDGE, CLEAN_MODE_SPOT, CLEAN_MODE_SPOT_AREA, CLEAN_MODE_SINGLE_ROOM} +CHARGING_STATES = {CHARGE_MODE_CHARGING} + +# These dictionaries convert to and from Sucks's consts (which closely match what the UI and manuals use) +# to and from what the Ecovacs API uses (which are sometimes very oddly named and have random capitalization.) +CLEAN_MODE_TO_ECOVACS = { + CLEAN_MODE_AUTO: 'auto', + CLEAN_MODE_EDGE: 'border', + CLEAN_MODE_SPOT: 'spot', + CLEAN_MODE_SPOT_AREA: 'SpotArea', + CLEAN_MODE_SINGLE_ROOM: 'singleroom', + CLEAN_MODE_STOP: 'stop' +} + +CLEAN_ACTION_TO_ECOVACS = { + CLEAN_ACTION_START: 's', + CLEAN_ACTION_PAUSE: 'p', + CLEAN_ACTION_RESUME: 'r', + CLEAN_ACTION_STOP: 'h', +} + +CLEAN_ACTION_FROM_ECOVACS = { + 's': CLEAN_ACTION_START, + 'p': CLEAN_ACTION_PAUSE, + 'r': CLEAN_ACTION_RESUME, + 'h': CLEAN_ACTION_STOP, +} + +CLEAN_MODE_FROM_ECOVACS = { + 'auto': CLEAN_MODE_AUTO, + 'border': CLEAN_MODE_EDGE, + 'spot': CLEAN_MODE_SPOT, + 'spot_area': CLEAN_MODE_SPOT_AREA, + 'SpotArea': CLEAN_MODE_SPOT_AREA, + 'singleroom': CLEAN_MODE_SINGLE_ROOM, + 'stop': CLEAN_MODE_STOP, + 'going': CHARGE_MODE_RETURNING, +} + +FAN_SPEED_TO_ECOVACS = { + FAN_SPEED_NORMAL: 'standard', + FAN_SPEED_HIGH: 'strong' +} + +FAN_SPEED_FROM_ECOVACS = { + 'standard': FAN_SPEED_NORMAL, + 'strong': FAN_SPEED_HIGH, +} + +CHARGE_MODE_TO_ECOVACS = { + CHARGE_MODE_RETURN: 'go', + CHARGE_MODE_RETURNING: 'Going', + CHARGE_MODE_CHARGING: 'SlotCharging', + CHARGE_MODE_IDLE: 'Idle', +} + +CHARGE_MODE_FROM_ECOVACS = { + 'going': CHARGE_MODE_RETURNING, +# 'Going': CHARGE_MODE_RETURNING, + 'slot_charging': CHARGE_MODE_CHARGING, +# 'SlotCharging': CHARGE_MODE_CHARGING, + 'idle': CHARGE_MODE_IDLE, +# 'Idle': CHARGE_MODE_IDLE, +} + +COMPONENT_TO_ECOVACS = { + COMPONENT_MAIN_BRUSH: 'Brush', + COMPONENT_SIDE_BRUSH: 'SideBrush', + COMPONENT_FILTER: 'DustCaseHeap', +} + +COMPONENT_FROM_ECOVACS = { + 'brush': COMPONENT_MAIN_BRUSH, +# 'Brush': COMPONENT_MAIN_BRUSH, + 'side_brush': COMPONENT_SIDE_BRUSH, +# 'SideBrush': COMPONENT_SIDE_BRUSH, + 'dust_case_heap': COMPONENT_FILTER, +# 'DustCaseHeap': COMPONENT_FILTER, +} \ No newline at end of file diff --git a/custom_components/ecovacs/manifest.json b/custom_components/ecovacs/manifest.json index 2163424..75346a0 100644 --- a/custom_components/ecovacs/manifest.json +++ b/custom_components/ecovacs/manifest.json @@ -1,10 +1,10 @@ { "domain": "ecovacs", "name": "Ecovacs Bumper", - "version": "1.3.4", + "version": "1.3.5", "documentation": "https://github.com/bittles/ha_ecovacs_bumper", "issue_tracker": "https://github.com/bittles/ha_ecovacs_bumper/issues", - "requirements": ["sleekxmppfs==1.4.1", "click>=6", "requests>=2.18", "pycryptodome>=3.4", "pycountry-convert>=0.5", "paho-mqtt>=1.4", "stringcase>=1.2"], + "requirements": ["sleekxmppfs==1.4.1", "requests>=2.18", "pycryptodome>=3.4", "pycountry-convert>=0.5", "paho-mqtt>=1.4", "stringcase>=1.2"], "codeowners": ["bittles"], "iot_class": "local_polling", "loggers": ["sleekxmppfs", "sucksbumper"] diff --git a/custom_components/ecovacs/mqtt_ecovacs.py b/custom_components/ecovacs/mqtt_ecovacs.py new file mode 100644 index 0000000..58ea7c3 --- /dev/null +++ b/custom_components/ecovacs/mqtt_ecovacs.py @@ -0,0 +1,236 @@ +import time + +import sched +import threading +import ssl +import requests +import stringcase +from threading import Event +from paho.mqtt.client import Client as ClientMQTT +from paho.mqtt import publish as MQTTPublish +from paho.mqtt import subscribe as MQTTSubscribe +from sleekxmppfs.xmlstream import ET + +from . import const + +#This is used by EcoVacsIOTMQ and EcoVacsXMPP for _ctl_to_dict +def RepresentsInt(stringvar): + try: + int(stringvar) + return True + except ValueError: + return False + +class EcoVacsIOTMQ(ClientMQTT): + def __init__(self, user, domain, resource, secret, continent, vacuum, server_address=None, verify_ssl=True): + ClientMQTT.__init__(self) + self.ctl_subscribers = [] + self.user = user + self.domain = str(domain).split(".")[0] #MQTT is using domain without tld extension + self.resource = resource + self.secret = secret + self.continent = continent + self.vacuum = vacuum + self.scheduler = sched.scheduler(time.time, time.sleep) + self.scheduler_thread = threading.Thread(target=self.scheduler.run, daemon=True, name="mqtt_schedule_thread") + self.verify_ssl = str_to_bool_or_cert(verify_ssl) + 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 connect_and_wait_until_ready(self): + #self._on_log = self.on_log #This provides more logging than needed, even for debug + self._on_message = self._handle_ctl_mqtt + 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() + + def subscribe_to_ctls(self, function): + self.ctl_subscribers.append(function) + + def _disconnect(self): + self.disconnect() #disconnect mqtt connection + self.scheduler.empty() #Clear schedule queue + + def _run_scheduled_func(self, timer_seconds, timer_function): + timer_function() + self.schedule(timer_seconds, timer_function) + + def schedule(self, timer_seconds, timer_function): + self.scheduler.enter(timer_seconds, 1, self._run_scheduled_func,(timer_seconds, timer_function)) + if not self.scheduler_thread.isAlive(): + self.scheduler_thread.start() + + 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 with MQTT Return {}".format(rc)) + raise RuntimeError("EcoVacsMQTT - error connecting with MQTT Return {}".format(rc)) + else: + _LOGGER.debug("EcoVacsMQTT - Connected 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): #This is very noisy and verbose + # _LOGGER.debug("EcoVacsMQTT Log: {} ".format(buf)) + + def send_ping(self): + _LOGGER.debug("*** MQTT sending ping ***") + rc = self._send_simple_command(MQTTPublish.paho.PINGREQ) + if rc == MQTTPublish.paho.MQTT_ERR_SUCCESS: + return True + else: + return False + + def send_command(self, action, recipient): + if action.name == "Clean": #For handling Clean when action not specified (i.e. CLI) + action.args['clean']['act'] = CLEAN_ACTION_TO_ECOVACS['start'] #Inject a start action + c = self._wrap_command(action, recipient) + _LOGGER.debug('Sending command {0}'.format(c)) + self._handle_ctl_api(action, + self.__call_iotdevmanager_api(c ,verify_ssl=self.verify_ssl ) + ) + + def _wrap_command(self, cmd, recipient): + #Remove the td from ctl xml for RestAPI + payloadxml = cmd.to_xml() + payloadxml.attrib.pop("td") + return { + 'auth': { + 'realm': EcoVacsAPI.REALM, + 'resource': self.resource, + 'token': self.secret, + 'userid': self.user, + 'with': 'users', + }, + "cmdName": cmd.name, + "payload": ET.tostring(payloadxml).decode(), + + "payloadType": "x", + "td": "q", + "toId": recipient, + "toRes": self.vacuum['resource'], + "toType": self.vacuum['class'] + } + + def __call_iotdevmanager_api(self, args, verify_ssl=True): + _LOGGER.debug("calling iotdevmanager api with {}".format(args)) + params = {} + params.update(args) + url = (EcoVacsAPI.PORTAL_URL_FORMAT + "/iot/devmanager.do").format(continent=self.continent) + response = None + try: #The RestAPI sometimes doesnt provide a response depending on command, reduce timeout to 3 to accomodate and make requests faster + response = requests.post(url, json=params, timeout=3, verify=verify_ssl) #May think about having timeout as an arg that could be provided in the future + except requests.exceptions.ReadTimeout: + _LOGGER.debug("call to iotdevmanager failed with ReadTimeout") + return {} + json = response.json() + if json['ret'] == 'ok': + return json + elif json['ret'] == 'fail': + if 'debug' in json: + if json['debug'] == 'wait for response timed out': + #TODO - Maybe handle timeout for IOT better in the future + _LOGGER.error("call to iotdevmanager failed with {}".format(json)) + return {} + else: + #TODO - Not sure if we want to raise an error yet, just return empty for now + _LOGGER.error("call to iotdevmanager failed with {}".format(json)) + return {} + #raise RuntimeError( + #"failure {} ({}) for call {} and parameters {}".format(json['error'], json['errno'], function, params)) + + def _handle_ctl_api(self, action, message): + if not message == {}: + resp = self._ctl_to_dict_api(action, message['resp']) + if resp is not None: + for s in self.ctl_subscribers: + s(resp) + + def _ctl_to_dict_api(self, action, xmlstring): + xml = ET.fromstring(xmlstring) + xmlchild = xml.getchildren() + if len(xmlchild) > 0: + result = xmlchild[0].attrib.copy() + #Fix for difference in XMPP vs API response + #Depending on the report will use the tag and add "report" to fit the mold of sucks library + if xmlchild[0].tag == "clean": + result['event'] = "CleanReport" + elif xmlchild[0].tag == "charge": + result['event'] = "ChargeState" + elif xmlchild[0].tag == "battery": + result['event'] = "BatteryInfo" + else: #Default back to replacing Get from the api cmdName + result['event'] = action.name.replace("Get","",1) + else: + result = xml.attrib.copy() + result['event'] = action.name.replace("Get","",1) + if 'ret' in result: #Handle errors as needed + if result['ret'] == 'fail': + if action.name == "Charge": #So far only seen this with Charge, when already docked + result['event'] = "ChargeState" + for key in result: + if not RepresentsInt(result[key]): #Fix to handle negative int values + result[key] = stringcase.snakecase(result[key]) + return result + + def _handle_ctl_mqtt(self, client, userdata, message): + #_LOGGER.debug("EcoVacs MQTT Received Message on Topic: {} - Message: {}".format(message.topic, str(message.payload.decode("utf-8")))) + as_dict = self._ctl_to_dict_mqtt(message.topic, str(message.payload.decode("utf-8"))) + if as_dict is not None: + for s in self.ctl_subscribers: + s(as_dict) + + def _ctl_to_dict_mqtt(self, topic, xmlstring): + #I haven't seen the need to fall back to data within the topic (like we do with IOT rest call actions), but it is here in case of future need + xml = ET.fromstring(xmlstring) #Convert from string to xml (like IOT rest calls), other than this it is similar to XMPP + #Including changes from jasonarends @ 28da7c2 below + result = xml.attrib.copy() + if 'td' not in result: + # This happens for commands with no response data, such as PlaySound + # Handle response data with no 'td' + if 'type' in result: # single element with type and val + result['event'] = "LifeSpan" # seems to always be LifeSpan type + else: + if len(xml) > 0: # case where there is child element + if 'clean' in xml[0].tag: + result['event'] = "CleanReport" + elif 'charge' in xml[0].tag: + result['event'] = "ChargeState" + elif 'battery' in xml[0].tag: + result['event'] = "BatteryInfo" + else: + return + result.update(xml[0].attrib) + else: # for non-'type' result with no child element, e.g., result of PlaySound + return + else: # response includes 'td' + result['event'] = result.pop('td') + if xml: + result.update(xml[0].attrib) + for key in result: + #Check for RepresentInt to handle negative int values, and ',' for ignoring position updates + if not RepresentsInt(result[key]) and ',' not in result[key]: + result[key] = stringcase.snakecase(result[key]) + return result \ No newline at end of file diff --git a/custom_components/ecovacs/sucksbumper.py b/custom_components/ecovacs/sucksbumper.py index f188da0..c45b326 100644 --- a/custom_components/ecovacs/sucksbumper.py +++ b/custom_components/ecovacs/sucksbumper.py @@ -1,134 +1,16 @@ import hashlib -import logging import time +import requests +import os from base64 import b64decode, b64encode from collections import OrderedDict -from threading import Event -import threading -import sched -import random -import ssl -import requests -import stringcase -import os -from sleekxmppfs import ClientXMPP, Callback, MatchXPath from sleekxmppfs.xmlstream import ET from sleekxmppfs.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 +from .mqtt_ecovacs import EcoVacsIOTMQ +from .xmpp_ecovacs import EcoVacsXMPP -_LOGGER = logging.getLogger(__name__) - -# These consts define all of the vocabulary used by this library when presenting various states and components. -# Applications implementing this library should import these rather than hard-code the strings, for future-proofing. - -CLEAN_MODE_AUTO = 'auto' -CLEAN_MODE_EDGE = 'edge' -CLEAN_MODE_SPOT = 'spot' -CLEAN_MODE_SPOT_AREA = 'spot_area' -CLEAN_MODE_SINGLE_ROOM = 'single_room' -CLEAN_MODE_STOP = 'stop' - -CLEAN_ACTION_START = 'start' -CLEAN_ACTION_PAUSE = 'pause' -CLEAN_ACTION_RESUME = 'resume' -CLEAN_ACTION_STOP = 'stop' - -FAN_SPEED_NORMAL = 'normal' -FAN_SPEED_HIGH = 'high' - -CHARGE_MODE_RETURN = 'return' -CHARGE_MODE_RETURNING = 'returning' -CHARGE_MODE_CHARGING = 'charging' -CHARGE_MODE_IDLE = 'idle' - -COMPONENT_SIDE_BRUSH = 'side_brush' -COMPONENT_MAIN_BRUSH = 'main_brush' -COMPONENT_FILTER = 'filter' - -VACUUM_STATUS_OFFLINE = 'offline' - -CLEANING_STATES = {CLEAN_MODE_AUTO, CLEAN_MODE_EDGE, CLEAN_MODE_SPOT, CLEAN_MODE_SPOT_AREA, CLEAN_MODE_SINGLE_ROOM} -CHARGING_STATES = {CHARGE_MODE_CHARGING} - -# These dictionaries convert to and from Sucks's consts (which closely match what the UI and manuals use) -# to and from what the Ecovacs API uses (which are sometimes very oddly named and have random capitalization.) -CLEAN_MODE_TO_ECOVACS = { - CLEAN_MODE_AUTO: 'auto', - CLEAN_MODE_EDGE: 'border', - CLEAN_MODE_SPOT: 'spot', - CLEAN_MODE_SPOT_AREA: 'SpotArea', - CLEAN_MODE_SINGLE_ROOM: 'singleroom', - CLEAN_MODE_STOP: 'stop' -} - -CLEAN_ACTION_TO_ECOVACS = { - CLEAN_ACTION_START: 's', - CLEAN_ACTION_PAUSE: 'p', - CLEAN_ACTION_RESUME: 'r', - CLEAN_ACTION_STOP: 'h', -} - -CLEAN_ACTION_FROM_ECOVACS = { - 's': CLEAN_ACTION_START, - 'p': CLEAN_ACTION_PAUSE, - 'r': CLEAN_ACTION_RESUME, - 'h': CLEAN_ACTION_STOP, -} - -CLEAN_MODE_FROM_ECOVACS = { - 'auto': CLEAN_MODE_AUTO, - 'border': CLEAN_MODE_EDGE, - 'spot': CLEAN_MODE_SPOT, - 'spot_area': CLEAN_MODE_SPOT_AREA, - 'SpotArea': CLEAN_MODE_SPOT_AREA, - 'singleroom': CLEAN_MODE_SINGLE_ROOM, - 'stop': CLEAN_MODE_STOP, - 'going': CHARGE_MODE_RETURNING, -} - -FAN_SPEED_TO_ECOVACS = { - FAN_SPEED_NORMAL: 'standard', - FAN_SPEED_HIGH: 'strong' -} - -FAN_SPEED_FROM_ECOVACS = { - 'standard': FAN_SPEED_NORMAL, - 'strong': FAN_SPEED_HIGH, -} - -CHARGE_MODE_TO_ECOVACS = { - CHARGE_MODE_RETURN: 'go', - CHARGE_MODE_RETURNING: 'Going', - CHARGE_MODE_CHARGING: 'SlotCharging', - CHARGE_MODE_IDLE: 'Idle', -} - -CHARGE_MODE_FROM_ECOVACS = { - 'going': CHARGE_MODE_RETURNING, -# 'Going': CHARGE_MODE_RETURNING, - 'slot_charging': CHARGE_MODE_CHARGING, -# 'SlotCharging': CHARGE_MODE_CHARGING, - 'idle': CHARGE_MODE_IDLE, -# 'Idle': CHARGE_MODE_IDLE, -} - -COMPONENT_TO_ECOVACS = { - COMPONENT_MAIN_BRUSH: 'Brush', - COMPONENT_SIDE_BRUSH: 'SideBrush', - COMPONENT_FILTER: 'DustCaseHeap', -} - -COMPONENT_FROM_ECOVACS = { - 'brush': COMPONENT_MAIN_BRUSH, -# 'Brush': COMPONENT_MAIN_BRUSH, - 'side_brush': COMPONENT_SIDE_BRUSH, -# 'SideBrush': COMPONENT_SIDE_BRUSH, - 'dust_case_heap': COMPONENT_FILTER, -# 'DustCaseHeap': COMPONENT_FILTER, -} +from . import const def str_to_bool_or_cert(s): if s == 'True' or s == True: @@ -602,357 +484,6 @@ class VacBot(): self.iotmq._disconnect() #self.xmpp.disconnect(wait=wait) #Leaving in case xmpp is added to iotmq in the future -#This is used by EcoVacsIOTMQ and EcoVacsXMPP for _ctl_to_dict -def RepresentsInt(stringvar): - try: - int(stringvar) - return True - except ValueError: - return False - -class EcoVacsIOTMQ(ClientMQTT): - def __init__(self, user, domain, resource, secret, continent, vacuum, server_address=None, verify_ssl=True): - ClientMQTT.__init__(self) - self.ctl_subscribers = [] - self.user = user - self.domain = str(domain).split(".")[0] #MQTT is using domain without tld extension - self.resource = resource - self.secret = secret - self.continent = continent - self.vacuum = vacuum - self.scheduler = sched.scheduler(time.time, time.sleep) - self.scheduler_thread = threading.Thread(target=self.scheduler.run, daemon=True, name="mqtt_schedule_thread") - self.verify_ssl = str_to_bool_or_cert(verify_ssl) - 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 connect_and_wait_until_ready(self): - #self._on_log = self.on_log #This provides more logging than needed, even for debug - self._on_message = self._handle_ctl_mqtt - 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() - - def subscribe_to_ctls(self, function): - self.ctl_subscribers.append(function) - - def _disconnect(self): - self.disconnect() #disconnect mqtt connection - self.scheduler.empty() #Clear schedule queue - - def _run_scheduled_func(self, timer_seconds, timer_function): - timer_function() - self.schedule(timer_seconds, timer_function) - - def schedule(self, timer_seconds, timer_function): - self.scheduler.enter(timer_seconds, 1, self._run_scheduled_func,(timer_seconds, timer_function)) - if not self.scheduler_thread.isAlive(): - self.scheduler_thread.start() - - 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 with MQTT Return {}".format(rc)) - raise RuntimeError("EcoVacsMQTT - error connecting with MQTT Return {}".format(rc)) - else: - _LOGGER.debug("EcoVacsMQTT - Connected 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): #This is very noisy and verbose - # _LOGGER.debug("EcoVacsMQTT Log: {} ".format(buf)) - - def send_ping(self): - _LOGGER.debug("*** MQTT sending ping ***") - rc = self._send_simple_command(MQTTPublish.paho.PINGREQ) - if rc == MQTTPublish.paho.MQTT_ERR_SUCCESS: - return True - else: - return False - - def send_command(self, action, recipient): - if action.name == "Clean": #For handling Clean when action not specified (i.e. CLI) - action.args['clean']['act'] = CLEAN_ACTION_TO_ECOVACS['start'] #Inject a start action - c = self._wrap_command(action, recipient) - _LOGGER.debug('Sending command {0}'.format(c)) - self._handle_ctl_api(action, - self.__call_iotdevmanager_api(c ,verify_ssl=self.verify_ssl ) - ) - - def _wrap_command(self, cmd, recipient): - #Remove the td from ctl xml for RestAPI - payloadxml = cmd.to_xml() - payloadxml.attrib.pop("td") - return { - 'auth': { - 'realm': EcoVacsAPI.REALM, - 'resource': self.resource, - 'token': self.secret, - 'userid': self.user, - 'with': 'users', - }, - "cmdName": cmd.name, - "payload": ET.tostring(payloadxml).decode(), - - "payloadType": "x", - "td": "q", - "toId": recipient, - "toRes": self.vacuum['resource'], - "toType": self.vacuum['class'] - } - - def __call_iotdevmanager_api(self, args, verify_ssl=True): - _LOGGER.debug("calling iotdevmanager api with {}".format(args)) - params = {} - params.update(args) - url = (EcoVacsAPI.PORTAL_URL_FORMAT + "/iot/devmanager.do").format(continent=self.continent) - response = None - try: #The RestAPI sometimes doesnt provide a response depending on command, reduce timeout to 3 to accomodate and make requests faster - response = requests.post(url, json=params, timeout=3, verify=verify_ssl) #May think about having timeout as an arg that could be provided in the future - except requests.exceptions.ReadTimeout: - _LOGGER.debug("call to iotdevmanager failed with ReadTimeout") - return {} - json = response.json() - if json['ret'] == 'ok': - return json - elif json['ret'] == 'fail': - if 'debug' in json: - if json['debug'] == 'wait for response timed out': - #TODO - Maybe handle timeout for IOT better in the future - _LOGGER.error("call to iotdevmanager failed with {}".format(json)) - return {} - else: - #TODO - Not sure if we want to raise an error yet, just return empty for now - _LOGGER.error("call to iotdevmanager failed with {}".format(json)) - return {} - #raise RuntimeError( - #"failure {} ({}) for call {} and parameters {}".format(json['error'], json['errno'], function, params)) - - def _handle_ctl_api(self, action, message): - if not message == {}: - resp = self._ctl_to_dict_api(action, message['resp']) - if resp is not None: - for s in self.ctl_subscribers: - s(resp) - - def _ctl_to_dict_api(self, action, xmlstring): - xml = ET.fromstring(xmlstring) - xmlchild = xml.getchildren() - if len(xmlchild) > 0: - result = xmlchild[0].attrib.copy() - #Fix for difference in XMPP vs API response - #Depending on the report will use the tag and add "report" to fit the mold of sucks library - if xmlchild[0].tag == "clean": - result['event'] = "CleanReport" - elif xmlchild[0].tag == "charge": - result['event'] = "ChargeState" - elif xmlchild[0].tag == "battery": - result['event'] = "BatteryInfo" - else: #Default back to replacing Get from the api cmdName - result['event'] = action.name.replace("Get","",1) - else: - result = xml.attrib.copy() - result['event'] = action.name.replace("Get","",1) - if 'ret' in result: #Handle errors as needed - if result['ret'] == 'fail': - if action.name == "Charge": #So far only seen this with Charge, when already docked - result['event'] = "ChargeState" - for key in result: - if not RepresentsInt(result[key]): #Fix to handle negative int values - result[key] = stringcase.snakecase(result[key]) - return result - - def _handle_ctl_mqtt(self, client, userdata, message): - #_LOGGER.debug("EcoVacs MQTT Received Message on Topic: {} - Message: {}".format(message.topic, str(message.payload.decode("utf-8")))) - as_dict = self._ctl_to_dict_mqtt(message.topic, str(message.payload.decode("utf-8"))) - if as_dict is not None: - for s in self.ctl_subscribers: - s(as_dict) - - def _ctl_to_dict_mqtt(self, topic, xmlstring): - #I haven't seen the need to fall back to data within the topic (like we do with IOT rest call actions), but it is here in case of future need - xml = ET.fromstring(xmlstring) #Convert from string to xml (like IOT rest calls), other than this it is similar to XMPP - #Including changes from jasonarends @ 28da7c2 below - result = xml.attrib.copy() - if 'td' not in result: - # This happens for commands with no response data, such as PlaySound - # Handle response data with no 'td' - if 'type' in result: # single element with type and val - result['event'] = "LifeSpan" # seems to always be LifeSpan type - else: - if len(xml) > 0: # case where there is child element - if 'clean' in xml[0].tag: - result['event'] = "CleanReport" - elif 'charge' in xml[0].tag: - result['event'] = "ChargeState" - elif 'battery' in xml[0].tag: - result['event'] = "BatteryInfo" - else: - return - result.update(xml[0].attrib) - else: # for non-'type' result with no child element, e.g., result of PlaySound - return - else: # response includes 'td' - result['event'] = result.pop('td') - if xml: - result.update(xml[0].attrib) - for key in result: - #Check for RepresentInt to handle negative int values, and ',' for ignoring position updates - if not RepresentsInt(result[key]) and ',' not in result[key]: - result[key] = stringcase.snakecase(result[key]) - return result - - -class EcoVacsXMPP(ClientXMPP): - def __init__(self, user, domain, resource, secret, continent, vacuum, server_address=None ): - ClientXMPP.__init__(self, "{}@{}/{}".format(user, domain,resource), '0/' + resource + '/' + secret) #Init with resource to bind it - self.user = user - self.domain = domain - self.boundjid.resource = resource - self.continent = continent - self.vacuum = vacuum - self.credentials['authzid'] = user - if server_address is None: - self.server_address = ('msg-{}.ecouser.net'.format(self.continent), '5223') - else: - self.server_address = server_address - self.add_event_handler("session_start", self.session_start) - self.ctl_subscribers = [] - self.ready_flag = Event() - - def wait_until_ready(self): - self.ready_flag.wait() - - def session_start(self, event): - _LOGGER.debug("----------------- starting session ----------------") - _LOGGER.debug("event = {}".format(event)) - self.register_handler(Callback("general", - MatchXPath('{jabber:client}iq/{com:ctl}query/{com:ctl}'), - self._handle_ctl)) - # register a ping handler, not really needed but keeps from errors being thrown - self.register_handler(Callback("Ping", - MatchXPath('{jabber:client}iq/{urn:xmpp:ping}ping/{urn:xmpp:ping}'), - self._handle_ping)) - self.ready_flag.set() - - def subscribe_to_ctls(self, function): - self.ctl_subscribers.append(function) - - def _handle_ctl(self, message): - 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): - #Including changes from jasonarends @ 28da7c2 below - result = xml.attrib.copy() - childxml = None - try: # check for child xml - childxml = xml[0] - except IndexError: - _LOGGER.debug("No child xml") - if 'td' not in result: - # Handle response data with no 'td' - if 'type' in result: # single element with type and val - result['event'] = "LifeSpan" # seems to always be LifeSpan type - else: - if childxml is not None: - if 'clean' in childxml.tag: - result['event'] = "CleanReport" - elif 'charge' in childxml.tag: - result['event'] = "ChargeState" - elif 'battery' in childxml.tag: - result['event'] = "BatteryInfo" - else: - return - result.update(childxml.attrib) - else: # for non-'type' result with no child element, e.g., result of PlaySound - return - else: # response includes 'td' - result['event'] = result.pop('td') - if xml: - result.update(xml[0].attrib) # reponses with td seem to always have child component - for key in result: - #Check for RepresentInt to handle negative int values, and ',' for ignoring position updates - if not RepresentsInt(result[key]) and ',' not in result[key]: - result[key] = stringcase.snakecase(result[key]) - return result - - def register_callback(self, userdata, message): - self.register_handler(Callback(kind, - MatchXPath('{jabber:client}iq/{com:ctl}query/{com:ctl}ctl[@td="' + kind + '"]'), - function)) - - 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' - if not "id" in ctl.attrib: - ctl.attrib["id"] = self.getReqID() #If no ctl id provided, add an id to the ctl. This was required for the ozmo930 and shouldn't hurt others - for child in q.xml: - if child.tag.endswith('query'): - child.append(ctl) - return q - - def getReqID(self, customid="0"): #Generate a somewhat random string for request id, with minium 8 chars. Works similar to ecovacs app. - if customid != "0": - return "{}".format(customid) #return provided id as string - else: - rtnval = str(random.randint(1,50)) - while len(str(rtnval)) <= 8: - rtnval = "{}{}".format(rtnval,random.randint(0,50)) - return "{}".format(rtnval) #return as string - - def _my_address(self): - if not self.vacuum['iotmq']: - return self.user + '@' + self.domain + '/' + self.boundjid.resource - else: - return self.user + '@' + self.domain + '/' + self.resource - - def send_ping(self, to): - q = self.make_iq_get(ito=to, ifrom=self._my_address()) - q.xml.append(ET.Element('ping', {'xmlns': 'urn:xmpp:ping'})) - _LOGGER.debug("*** sending ping ***") - q.send() - - # used some code from a sleekxmppfs plugin, seems to work fine - def _handle_ping(self, iq): - _LOGGER.debug("Pinged by %s", iq['from']) - iq.reply().send() - - def connect_and_wait_until_ready(self): - self.connect(self.server_address) - self.process() - self.wait_until_ready() - class VacBotCommand: ACTION = { 'forward': 'forward', diff --git a/custom_components/ecovacs/xmpp_ecovacs.py b/custom_components/ecovacs/xmpp_ecovacs.py new file mode 100644 index 0000000..a9be3bc --- /dev/null +++ b/custom_components/ecovacs/xmpp_ecovacs.py @@ -0,0 +1,143 @@ +import stringcase +import random +from threading import Event +from sleekxmppfs import ClientXMPP, Callback, MatchXPath +from sleekxmppfs.xmlstream import ET +#from sleekxmppfs.exceptions import XMPPError +from . import const + +#This is used by EcoVacsIOTMQ and EcoVacsXMPP for _ctl_to_dict +def RepresentsInt(stringvar): + try: + int(stringvar) + return True + except ValueError: + return False + +class EcoVacsXMPP(ClientXMPP): + def __init__(self, user, domain, resource, secret, continent, vacuum, server_address=None ): + ClientXMPP.__init__(self, "{}@{}/{}".format(user, domain,resource), '0/' + resource + '/' + secret) #Init with resource to bind it + self.user = user + self.domain = domain + self.boundjid.resource = resource + self.continent = continent + self.vacuum = vacuum + self.credentials['authzid'] = user + if server_address is None: + self.server_address = ('msg-{}.ecouser.net'.format(self.continent), '5223') + else: + self.server_address = server_address + self.add_event_handler("session_start", self.session_start) + self.ctl_subscribers = [] + self.ready_flag = Event() + + def wait_until_ready(self): + self.ready_flag.wait() + + def session_start(self, event): + _LOGGER.debug("----------------- starting session ----------------") + _LOGGER.debug("event = {}".format(event)) + self.register_handler(Callback("general", + MatchXPath('{jabber:client}iq/{com:ctl}query/{com:ctl}'), + self._handle_ctl)) + # register a ping handler, not really needed but keeps from errors being thrown + self.register_handler(Callback("Ping", + MatchXPath('{jabber:client}iq/{urn:xmpp:ping}ping/{urn:xmpp:ping}'), + self._handle_ping)) + self.ready_flag.set() + + def subscribe_to_ctls(self, function): + self.ctl_subscribers.append(function) + + def _handle_ctl(self, message): + 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): + #Including changes from jasonarends @ 28da7c2 below + result = xml.attrib.copy() + childxml = None + try: # check for child xml + childxml = xml[0] + except IndexError: + _LOGGER.debug("No child xml") + if 'td' not in result: + # Handle response data with no 'td' + if 'type' in result: # single element with type and val + result['event'] = "LifeSpan" # seems to always be LifeSpan type + else: + if childxml is not None: + if 'clean' in childxml.tag: + result['event'] = "CleanReport" + elif 'charge' in childxml.tag: + result['event'] = "ChargeState" + elif 'battery' in childxml.tag: + result['event'] = "BatteryInfo" + else: + return + result.update(childxml.attrib) + else: # for non-'type' result with no child element, e.g., result of PlaySound + return + else: # response includes 'td' + result['event'] = result.pop('td') + if xml: + result.update(xml[0].attrib) # reponses with td seem to always have child component + for key in result: + #Check for RepresentInt to handle negative int values, and ',' for ignoring position updates + if not RepresentsInt(result[key]) and ',' not in result[key]: + result[key] = stringcase.snakecase(result[key]) + return result + + def register_callback(self, userdata, message): + self.register_handler(Callback(kind, + MatchXPath('{jabber:client}iq/{com:ctl}query/{com:ctl}ctl[@td="' + kind + '"]'), + function)) + + 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' + if not "id" in ctl.attrib: + ctl.attrib["id"] = self.getReqID() #If no ctl id provided, add an id to the ctl. This was required for the ozmo930 and shouldn't hurt others + for child in q.xml: + if child.tag.endswith('query'): + child.append(ctl) + return q + + def getReqID(self, customid="0"): #Generate a somewhat random string for request id, with minium 8 chars. Works similar to ecovacs app. + if customid != "0": + return "{}".format(customid) #return provided id as string + else: + rtnval = str(random.randint(1,50)) + while len(str(rtnval)) <= 8: + rtnval = "{}{}".format(rtnval,random.randint(0,50)) + return "{}".format(rtnval) #return as string + + def _my_address(self): + if not self.vacuum['iotmq']: + return self.user + '@' + self.domain + '/' + self.boundjid.resource + else: + return self.user + '@' + self.domain + '/' + self.resource + + def send_ping(self, to): + q = self.make_iq_get(ito=to, ifrom=self._my_address()) + q.xml.append(ET.Element('ping', {'xmlns': 'urn:xmpp:ping'})) + _LOGGER.debug("*** sending ping ***") + q.send() + + # used some code from a sleekxmppfs plugin, seems to work fine + def _handle_ping(self, iq): + _LOGGER.debug("Pinged by %s", iq['from']) + iq.reply().send() + + def connect_and_wait_until_ready(self): + self.connect(self.server_address) + self.process() + self.wait_until_ready() \ No newline at end of file