code refactoring

split xmpp and mqtt vacuums to their own modules
This commit is contained in:
bittles
2023-01-03 20:29:44 -05:00
parent 4fb79b7186
commit 17b041074c
6 changed files with 526 additions and 500 deletions
+17 -24
View File
@@ -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,
)
+123
View File
@@ -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,
}
+2 -2
View File
@@ -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"]
+236
View File
@@ -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
+5 -474
View File
@@ -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',
+143
View File
@@ -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()