From 586fb2ae62259553f2ae4e981fd5d21bf7d6ae2e Mon Sep 17 00:00:00 2001 From: Jonathan Schemoul Date: Tue, 11 Apr 2017 15:51:03 +0200 Subject: [PATCH] Adding a /write handle on an SID to be able to write multiple attributes to a device (example MID and VOL for sound play) --- src/main.py | 6 +++--- src/mqtt.py | 51 ++++++++++++++++++++++++++++++++++++------------ src/xiaomihub.py | 14 ++++++------- 3 files changed, 48 insertions(+), 23 deletions(-) diff --git a/src/main.py b/src/main.py index 24d48ca..1d116ff 100644 --- a/src/main.py +++ b/src/main.py @@ -67,10 +67,9 @@ def process_mqtt_messages(gateway, client): _LOGGER.debug("data from mqtt: " + format(data)) sid = data.get("sid", None) - param = data.get("param", None) - value = data.get("value", None) + values = data.get("values", dict()) - resp = gateway.write_to_hub(sid, param, value) + resp = gateway.write_to_hub(sid, **values) client._queue.task_done() except Exception as e: _LOGGER.error('Error while sending from mqtt to gateway: ', str(e)) @@ -87,6 +86,7 @@ if __name__ == "__main__": client.connect() #only this devices can be controlled from MQTT client.subscribe("gateway", "+", "+", "set") + client.subscribe("gateway", "+", "write", None) client.subscribe("plug", "+", "status", "set") gateway = XiaomiHub(gateway_pass) diff --git a/src/mqtt.py b/src/mqtt.py index ec28762..d23fab4 100644 --- a/src/mqtt.py +++ b/src/mqtt.py @@ -3,6 +3,7 @@ import os import logging from queue import Queue from threading import Thread +import json _LOGGER = logging.getLogger(__name__) @@ -54,8 +55,10 @@ class Mqtt: t1.start() self._threads.append(t1) - def subscribe(self, model="+", name="+", prop="+", command="set"): - topic = self.prefix + "/" + model + "/" + name + "/" + prop + "/" + command + def subscribe(self, model="+", name="+", prop="+", command=None): + topic = self.prefix + "/" + model + "/" + name + "/" + prop + if command is not None: + topic += "/" + command _LOGGER.info("Subscibing to " + topic + ".") self._client.subscribe(topic) @@ -94,17 +97,21 @@ class Mqtt: def _mqtt_process_message(self, client, userdata, msg): _LOGGER.info("Processing message in " + str(msg.topic) + ": " + str(msg.payload) + ".") parts = msg.topic.split("/") - if (len(parts) != 5): + if len(parts) < 4: + # should we return an error message ? return + model = parts[1] query_sid = parts[2] #sid or name part param = parts[3] #param part - value = (msg.payload).decode('utf-8') - if self._is_int(value): - value = int(value) + method = None + if len(parts) > 4: + method = parts[4] + else: + method = parts[3] + name = "" # we will find it next sid = query_sid - isFound = False for current_sid in self._sids: if (current_sid == None): @@ -124,15 +131,33 @@ class Mqtt: continue if isFound == False: + # should we return an error message ? return - # fix for rgb format - if (param == "rgb" and "," in str(value)): - value = self._color_rgb_to_xiaomi(value) + if method == "set": + # use single value set method - data = {'sid': sid, 'model': model, 'name': name, 'param':param, 'value':value} - # put in process queuee - self._queue.put(data) + value = (msg.payload).decode('utf-8') + if self._is_int(value): + value = int(value) + + # fix for rgb format + if (param == "rgb" and "," in str(value)): + value = self._color_rgb_to_xiaomi(value) + + # prepare values dict + data = {'sid': sid, 'model': model, 'name': name, + 'values': {param: value}} + # put in process queuee + self._queue.put(data) + + elif method == "write": + # use raw write method to the sensor, we expect a jsonified dict here. + values = json.loads((msg.payload).decode('utf-8')) + data = {'sid': sid, 'model': model, 'name': name, + 'values': values} + # put in process queuee + self._queue.put(data) def _mqtt_loop(self): _LOGGER.info("Starting mqtt loop.") diff --git a/src/xiaomihub.py b/src/xiaomihub.py index 7b70f76..6b11844 100644 --- a/src/xiaomihub.py +++ b/src/xiaomihub.py @@ -149,14 +149,14 @@ class XiaomiHub: _LOGGER.error("Cannot connect to Gateway") socket.close() - def write_to_hub(self, sid, data_key, datavalue): + def write_to_hub(self, sid, **values): key = self._get_key() - if type(datavalue) == int: - datavalue_formatted = str(datavalue) - else: - datavalue_formatted = '"' + datavalue + '"' - cmd = '{ "cmd":"write","sid":"' + sid + '","data":"{"' + data_key + '":' + datavalue_formatted + ',"key":"' + key + '"}}' - return self._send_cmd(cmd, "write_ack") + cmd = { + "cmd": "write", + "sid": sid, + "data": dict(key=key, **values) + } + return self._send_cmd(json.dumps(cmd), "write_ack") def get_from_hub(self, sid): cmd = '{ "cmd":"read","sid":"' + sid + '"}'