Merge pull request #4 from jon1012/master
Adding a /write command on SIDs
This commit is contained in:
+74
-74
@@ -13,94 +13,94 @@ logging.basicConfig(level=logging.INFO)
|
|||||||
_LOGGER = logging.getLogger(__name__)
|
_LOGGER = logging.getLogger(__name__)
|
||||||
|
|
||||||
def process_gateway_messages(gateway, client):
|
def process_gateway_messages(gateway, client):
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
packet = gateway._queue.get()
|
packet = gateway._queue.get()
|
||||||
_LOGGER.debug("data from queuee: " + format(packet))
|
_LOGGER.debug("data from queuee: " + format(packet))
|
||||||
|
|
||||||
sid = packet.get("sid", None)
|
sid = packet.get("sid", None)
|
||||||
model = packet.get("model", "")
|
model = packet.get("model", "")
|
||||||
data = packet.get("data", "")
|
data = packet.get("data", "")
|
||||||
|
|
||||||
if (sid != None and data != ""):
|
if (sid != None and data != ""):
|
||||||
data_decoded = json.loads(data)
|
data_decoded = json.loads(data)
|
||||||
client.publish(model, sid, data_decoded)
|
client.publish(model, sid, data_decoded)
|
||||||
gateway._queue.task_done()
|
gateway._queue.task_done()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
_LOGGER.error('Error while sending from gateway to mqtt: ', str(e))
|
_LOGGER.error('Error while sending from gateway to mqtt: ', str(e))
|
||||||
|
|
||||||
def read_motion_data(gateway, client, polling_interval, polling_models):
|
def read_motion_data(gateway, client, polling_interval, polling_models):
|
||||||
first = True
|
first = True
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
for device_type in gateway.XIAOMI_DEVICES:
|
for device_type in gateway.XIAOMI_DEVICES:
|
||||||
devices = gateway.XIAOMI_DEVICES[device_type]
|
devices = gateway.XIAOMI_DEVICES[device_type]
|
||||||
for device in devices:
|
for device in devices:
|
||||||
model = device.get("model", "")
|
model = device.get("model", "")
|
||||||
if (model not in polling_models):
|
if (model not in polling_models):
|
||||||
continue
|
continue
|
||||||
sid = device['sid']
|
sid = device['sid']
|
||||||
|
|
||||||
sensor_resp = gateway.get_from_hub(sid)
|
sensor_resp = gateway.get_from_hub(sid)
|
||||||
if (sensor_resp == None):
|
if (sensor_resp == None):
|
||||||
continue;
|
continue;
|
||||||
if (sensor_resp['sid'] != sid):
|
if (sensor_resp['sid'] != sid):
|
||||||
_LOGGER.error("Error: Response sid(" + sensor_resp['sid'] + ") differs from requested(" + sid + "). Skipping.")
|
_LOGGER.error("Error: Response sid(" + sensor_resp['sid'] + ") differs from requested(" + sid + "). Skipping.")
|
||||||
continue;
|
continue;
|
||||||
|
|
||||||
data = json.loads(sensor_resp['data'])
|
data = json.loads(sensor_resp['data'])
|
||||||
state = data.get("status", None)
|
state = data.get("status", None)
|
||||||
short_id = sensor_resp['short_id']
|
short_id = sensor_resp['short_id']
|
||||||
if ( device['data'] != data or first):
|
if ( device['data'] != data or first):
|
||||||
device['data'] = data
|
device['data'] = data
|
||||||
_LOGGER.debug("Polling result differs for " + str(model) + " with sid(First: " + str(first) + "): " + str(sid) + "; " + str(data))
|
_LOGGER.debug("Polling result differs for " + str(model) + " with sid(First: " + str(first) + "): " + str(sid) + "; " + str(data))
|
||||||
client.publish(model, sid, data)
|
client.publish(model, sid, data)
|
||||||
first = False
|
first = False
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
_LOGGER.error('Error while sending from mqtt to gateway: ', str(e))
|
_LOGGER.error('Error while sending from mqtt to gateway: ', str(e))
|
||||||
time.sleep(polling_interval)
|
time.sleep(polling_interval)
|
||||||
|
|
||||||
def process_mqtt_messages(gateway, client):
|
def process_mqtt_messages(gateway, client):
|
||||||
while True:
|
while True:
|
||||||
try:
|
try:
|
||||||
data = client._queue.get()
|
data = client._queue.get()
|
||||||
_LOGGER.debug("data from mqtt: " + format(data))
|
_LOGGER.debug("data from mqtt: " + format(data))
|
||||||
|
|
||||||
sid = data.get("sid", None)
|
sid = data.get("sid", None)
|
||||||
param = data.get("param", None)
|
values = data.get("values", dict())
|
||||||
value = data.get("value", None)
|
|
||||||
|
|
||||||
resp = gateway.write_to_hub(sid, param, value)
|
resp = gateway.write_to_hub(sid, **values)
|
||||||
client._queue.task_done()
|
client._queue.task_done()
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
_LOGGER.error('Error while sending from mqtt to gateway: ', str(e))
|
_LOGGER.error('Error while sending from mqtt to gateway: ', str(e))
|
||||||
|
|
||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
_LOGGER.info("Loading config file...")
|
_LOGGER.info("Loading config file...")
|
||||||
config=yamlparser.load_yaml('config/config.yaml')
|
config=yamlparser.load_yaml('config/config.yaml')
|
||||||
gateway_pass = yamlparser.get_gateway_password(config)
|
gateway_pass = yamlparser.get_gateway_password(config)
|
||||||
polling_interval = config['gateway'].get("polling_interval", 2)
|
polling_interval = config['gateway'].get("polling_interval", 2)
|
||||||
polling_models = config['gateway'].get("polling_models", ['motion'])
|
polling_models = config['gateway'].get("polling_models", ['motion'])
|
||||||
|
|
||||||
_LOGGER.info("Init mqtt client.")
|
_LOGGER.info("Init mqtt client.")
|
||||||
client = mqtt.Mqtt(config)
|
client = mqtt.Mqtt(config)
|
||||||
client.connect()
|
client.connect()
|
||||||
#only this devices can be controlled from MQTT
|
#only this devices can be controlled from MQTT
|
||||||
client.subscribe("gateway", "+", "+", "set")
|
client.subscribe("gateway", "+", "+", "set")
|
||||||
client.subscribe("plug", "+", "status", "set")
|
client.subscribe("gateway", "+", "write", None)
|
||||||
|
client.subscribe("plug", "+", "status", "set")
|
||||||
|
|
||||||
gateway = XiaomiHub(gateway_pass)
|
gateway = XiaomiHub(gateway_pass)
|
||||||
t1 = threading.Thread(target=process_gateway_messages, args=[gateway, client])
|
t1 = threading.Thread(target=process_gateway_messages, args=[gateway, client])
|
||||||
t1.daemon = True
|
t1.daemon = True
|
||||||
t1.start()
|
t1.start()
|
||||||
|
|
||||||
t2 = threading.Thread(target=process_mqtt_messages, args=[gateway, client])
|
t2 = threading.Thread(target=process_mqtt_messages, args=[gateway, client])
|
||||||
t2.daemon = True
|
t2.daemon = True
|
||||||
t2.start()
|
t2.start()
|
||||||
|
|
||||||
t3 = threading.Thread(target=read_motion_data, args=[gateway, client, polling_interval, polling_models])
|
t3 = threading.Thread(target=read_motion_data, args=[gateway, client, polling_interval, polling_models])
|
||||||
t3.daemon = True
|
t3.daemon = True
|
||||||
t3.start()
|
t3.start()
|
||||||
|
|
||||||
while True:
|
while True:
|
||||||
time.sleep(10)
|
time.sleep(10)
|
||||||
|
|||||||
+160
-135
@@ -3,165 +3,190 @@ import os
|
|||||||
import logging
|
import logging
|
||||||
from queue import Queue
|
from queue import Queue
|
||||||
from threading import Thread
|
from threading import Thread
|
||||||
|
import json
|
||||||
|
|
||||||
_LOGGER = logging.getLogger(__name__)
|
_LOGGER = logging.getLogger(__name__)
|
||||||
|
|
||||||
class Mqtt:
|
class Mqtt:
|
||||||
username = ""
|
username = ""
|
||||||
password = ""
|
password = ""
|
||||||
server = "localhost"
|
server = "localhost"
|
||||||
port = 1883
|
port = 1883
|
||||||
prefix = "home"
|
prefix = "home"
|
||||||
|
|
||||||
_client = None
|
_client = None
|
||||||
_sids = None
|
_sids = None
|
||||||
_queue = None
|
_queue = None
|
||||||
_threads = None
|
_threads = None
|
||||||
|
|
||||||
def __init__(self, config):
|
def __init__(self, config):
|
||||||
if (config == None):
|
if (config == None):
|
||||||
raise "Config is null"
|
raise "Config is null"
|
||||||
|
|
||||||
#load sids dictionary
|
#load sids dictionary
|
||||||
self._sids = config.get("sids", None)
|
self._sids = config.get("sids", None)
|
||||||
if (self._sids == None):
|
if (self._sids == None):
|
||||||
self._sids = dict({})
|
self._sids = dict({})
|
||||||
|
|
||||||
#load mqtt settings
|
#load mqtt settings
|
||||||
mqttConfig = config.get("mqtt", None)
|
mqttConfig = config.get("mqtt", None)
|
||||||
if (mqttConfig == None):
|
if (mqttConfig == None):
|
||||||
raise "Config mqtt section is null"
|
raise "Config mqtt section is null"
|
||||||
|
|
||||||
self.username = mqttConfig.get("username", "")
|
self.username = mqttConfig.get("username", "")
|
||||||
self.password = mqttConfig.get("password", "")
|
self.password = mqttConfig.get("password", "")
|
||||||
self.server = mqttConfig.get("server", "localhost")
|
self.server = mqttConfig.get("server", "localhost")
|
||||||
self.port = mqttConfig.get("port", 1883)
|
self.port = mqttConfig.get("port", 1883)
|
||||||
self.prefix = mqttConfig.get("prefix", "home")
|
self.prefix = mqttConfig.get("prefix", "home")
|
||||||
self._queue = Queue()
|
self._queue = Queue()
|
||||||
self._threads = []
|
self._threads = []
|
||||||
|
|
||||||
def connect(self):
|
def connect(self):
|
||||||
_LOGGER.info("Connecting to MQTT server " + self.server + ":" + str(self.port) + " with username (" + self.username + ":" + self.password + ")")
|
_LOGGER.info("Connecting to MQTT server " + self.server + ":" + str(self.port) + " with username (" + self.username + ":" + self.password + ")")
|
||||||
self._client = mqtt.Client()
|
self._client = mqtt.Client()
|
||||||
if (self.username != "" and self.password != ""):
|
if (self.username != "" and self.password != ""):
|
||||||
self._client.username_pw_set(self.username, self.password)
|
self._client.username_pw_set(self.username, self.password)
|
||||||
self._client.on_message = self._mqtt_process_message
|
self._client.on_message = self._mqtt_process_message
|
||||||
self._client.on_connect = self._mqtt_on_connect
|
self._client.on_connect = self._mqtt_on_connect
|
||||||
self._client.connect(self.server, self.port, 60)
|
self._client.connect(self.server, self.port, 60)
|
||||||
|
|
||||||
#run message processing loop
|
#run message processing loop
|
||||||
t1 = Thread(target=self._mqtt_loop)
|
t1 = Thread(target=self._mqtt_loop)
|
||||||
t1.start()
|
t1.start()
|
||||||
self._threads.append(t1)
|
self._threads.append(t1)
|
||||||
|
|
||||||
def subscribe(self, model="+", name="+", prop="+", command="set"):
|
def subscribe(self, model="+", name="+", prop="+", command=None):
|
||||||
topic = self.prefix + "/" + model + "/" + name + "/" + prop + "/" + command
|
topic = self.prefix + "/" + model + "/" + name + "/" + prop
|
||||||
_LOGGER.info("Subscibing to " + topic + ".")
|
if command is not None:
|
||||||
self._client.subscribe(topic)
|
topic += "/" + command
|
||||||
|
_LOGGER.info("Subscibing to " + topic + ".")
|
||||||
|
self._client.subscribe(topic)
|
||||||
|
|
||||||
def publish(self, model, sid, data, retain=True):
|
def publish(self, model, sid, data, retain=True):
|
||||||
sidprops = self._sids.get(sid, None)
|
sidprops = self._sids.get(sid, None)
|
||||||
if (sidprops != None):
|
if (sidprops != None):
|
||||||
model = sidprops.get("model",model)
|
model = sidprops.get("model",model)
|
||||||
sid = sidprops.get("name",sid)
|
sid = sidprops.get("name",sid)
|
||||||
|
|
||||||
# _LOGGER.info("data is " + format(data))
|
# _LOGGER.info("data is " + format(data))
|
||||||
PATH_FMT = self.prefix + "/{model}/{sid}/{prop}"
|
PATH_FMT = self.prefix + "/{model}/{sid}/{prop}"
|
||||||
for key, value in data.items():
|
for key, value in data.items():
|
||||||
# fix for latest motion value
|
# fix for latest motion value
|
||||||
if (model == "motion" and key == "no_motion"):
|
if (model == "motion" and key == "no_motion"):
|
||||||
key="status"
|
key="status"
|
||||||
value="no_motion"
|
value="no_motion"
|
||||||
if (model == "magnet" and key == "no_close"):
|
if (model == "magnet" and key == "no_close"):
|
||||||
key="status"
|
key="status"
|
||||||
value="open"
|
value="open"
|
||||||
|
|
||||||
# do not retain event-based sensors (like switches and cubes).
|
# do not retain event-based sensors (like switches and cubes).
|
||||||
if (model in ["switch", "cube"]):
|
if (model in ["switch", "cube"]):
|
||||||
retain = False
|
retain = False
|
||||||
|
|
||||||
# fix for rgb format
|
# fix for rgb format
|
||||||
if (key == "rgb" and self._is_int(value)):
|
if (key == "rgb" and self._is_int(value)):
|
||||||
value = self._color_xiaomi_to_rgb(value)
|
value = self._color_xiaomi_to_rgb(value)
|
||||||
|
|
||||||
topic = PATH_FMT.format(model=model, sid=sid, prop=key)
|
topic = PATH_FMT.format(model=model, sid=sid, prop=key)
|
||||||
_LOGGER.info("Publishing message to topic " + topic + ": " + str(value) + ".")
|
_LOGGER.info("Publishing message to topic " + topic + ": " + str(value) + ".")
|
||||||
self._client.publish(topic, payload=value, qos=0, retain=retain)
|
self._client.publish(topic, payload=value, qos=0, retain=retain)
|
||||||
|
|
||||||
def _mqtt_on_connect(self, client, userdata, rc, unk):
|
def _mqtt_on_connect(self, client, userdata, rc, unk):
|
||||||
_LOGGER.info("Connected to mqtt server.")
|
_LOGGER.info("Connected to mqtt server.")
|
||||||
|
|
||||||
def _mqtt_process_message(self, client, userdata, msg):
|
def _mqtt_process_message(self, client, userdata, msg):
|
||||||
_LOGGER.info("Processing message in " + str(msg.topic) + ": " + str(msg.payload) + ".")
|
_LOGGER.info("Processing message in " + str(msg.topic) + ": " + str(msg.payload) + ".")
|
||||||
parts = msg.topic.split("/")
|
parts = msg.topic.split("/")
|
||||||
if (len(parts) != 5):
|
if len(parts) < 4:
|
||||||
return
|
# should we return an error message ?
|
||||||
model = parts[1]
|
return
|
||||||
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)
|
|
||||||
name = "" # we will find it next
|
|
||||||
sid = query_sid
|
|
||||||
|
|
||||||
isFound = False
|
model = parts[1]
|
||||||
for current_sid in self._sids:
|
query_sid = parts[2] #sid or name part
|
||||||
if (current_sid == None):
|
param = parts[3] #param part
|
||||||
continue
|
method = None
|
||||||
sidprops = self._sids.get(current_sid, None)
|
if len(parts) > 4:
|
||||||
if sidprops == None:
|
method = parts[4]
|
||||||
continue
|
else:
|
||||||
sidname = sidprops.get("name", current_sid)
|
method = parts[3]
|
||||||
sidmodel = sidprops.get("model", "")
|
|
||||||
if (sidname == query_sid and sidmodel == model):
|
|
||||||
sid = current_sid
|
|
||||||
name = sidname
|
|
||||||
isFound = True
|
|
||||||
break
|
|
||||||
else:
|
|
||||||
_LOGGER.debug(sidmodel + "-" + sidname + " is not " + model + "-" + query_sid + ".")
|
|
||||||
continue
|
|
||||||
|
|
||||||
if isFound == False:
|
name = "" # we will find it next
|
||||||
return
|
sid = query_sid
|
||||||
|
isFound = False
|
||||||
|
for current_sid in self._sids:
|
||||||
|
if (current_sid == None):
|
||||||
|
continue
|
||||||
|
sidprops = self._sids.get(current_sid, None)
|
||||||
|
if sidprops == None:
|
||||||
|
continue
|
||||||
|
sidname = sidprops.get("name", current_sid)
|
||||||
|
sidmodel = sidprops.get("model", "")
|
||||||
|
if (sidname == query_sid and sidmodel == model):
|
||||||
|
sid = current_sid
|
||||||
|
name = sidname
|
||||||
|
isFound = True
|
||||||
|
break
|
||||||
|
else:
|
||||||
|
_LOGGER.debug(sidmodel + "-" + sidname + " is not " + model + "-" + query_sid + ".")
|
||||||
|
continue
|
||||||
|
|
||||||
# fix for rgb format
|
if isFound == False:
|
||||||
if (param == "rgb" and "," in str(value)):
|
# should we return an error message ?
|
||||||
value = self._color_rgb_to_xiaomi(value)
|
return
|
||||||
|
|
||||||
data = {'sid': sid, 'model': model, 'name': name, 'param':param, 'value':value}
|
if method == "set":
|
||||||
# put in process queuee
|
# use single value set method
|
||||||
self._queue.put(data)
|
|
||||||
|
|
||||||
def _mqtt_loop(self):
|
value = (msg.payload).decode('utf-8')
|
||||||
_LOGGER.info("Starting mqtt loop.")
|
if self._is_int(value):
|
||||||
self._client.loop_forever()
|
value = int(value)
|
||||||
|
|
||||||
def _color_xiaomi_to_rgb(self, xiaomi_color):
|
# fix for rgb format
|
||||||
intval = int(xiaomi_color)
|
if (param == "rgb" and "," in str(value)):
|
||||||
blue = (intval) & 255
|
value = self._color_rgb_to_xiaomi(value)
|
||||||
green = (intval >> 8) & 255
|
|
||||||
red = (intval >> 16) & 255
|
|
||||||
bright = (intval >> 24) & 255
|
|
||||||
value = str(red)+","+str(green)+","+str(blue)+","+str(bright)
|
|
||||||
return value
|
|
||||||
|
|
||||||
def _color_rgb_to_xiaomi(self, rgb_string):
|
# prepare values dict
|
||||||
arr = rgb_string.split(",")
|
data = {'sid': sid, 'model': model, 'name': name,
|
||||||
r = int(arr[0])
|
'values': {param: value}}
|
||||||
g = int(arr[1])
|
# put in process queuee
|
||||||
b = int(arr[2])
|
self._queue.put(data)
|
||||||
if len(arr)>3:
|
|
||||||
bright = int(arr[3])
|
elif method == "write":
|
||||||
else:
|
# use raw write method to the sensor, we expect a jsonified dict here.
|
||||||
bright = 255
|
values = json.loads((msg.payload).decode('utf-8'))
|
||||||
value = int('%02x%02x%02x%02x' % (bright, r, g, b), 16)
|
data = {'sid': sid, 'model': model, 'name': name,
|
||||||
return value
|
'values': values}
|
||||||
|
# put in process queuee
|
||||||
|
self._queue.put(data)
|
||||||
|
|
||||||
def _is_int(self, x):
|
def _mqtt_loop(self):
|
||||||
try:
|
_LOGGER.info("Starting mqtt loop.")
|
||||||
tmp = int(x)
|
self._client.loop_forever()
|
||||||
return True
|
|
||||||
except Exception as e:
|
def _color_xiaomi_to_rgb(self, xiaomi_color):
|
||||||
return False
|
intval = int(xiaomi_color)
|
||||||
|
blue = (intval) & 255
|
||||||
|
green = (intval >> 8) & 255
|
||||||
|
red = (intval >> 16) & 255
|
||||||
|
bright = (intval >> 24) & 255
|
||||||
|
value = str(red)+","+str(green)+","+str(blue)+","+str(bright)
|
||||||
|
return value
|
||||||
|
|
||||||
|
def _color_rgb_to_xiaomi(self, rgb_string):
|
||||||
|
arr = rgb_string.split(",")
|
||||||
|
r = int(arr[0])
|
||||||
|
g = int(arr[1])
|
||||||
|
b = int(arr[2])
|
||||||
|
if len(arr)>3:
|
||||||
|
bright = int(arr[3])
|
||||||
|
else:
|
||||||
|
bright = 255
|
||||||
|
value = int('%02x%02x%02x%02x' % (bright, r, g, b), 16)
|
||||||
|
return value
|
||||||
|
|
||||||
|
def _is_int(self, x):
|
||||||
|
try:
|
||||||
|
tmp = int(x)
|
||||||
|
return True
|
||||||
|
except Exception as e:
|
||||||
|
return False
|
||||||
|
|||||||
+8
-8
@@ -55,7 +55,7 @@ class XiaomiHub:
|
|||||||
_LOGGER.error("Cannot discover hub using whois: {0}".format(e))
|
_LOGGER.error("Cannot discover hub using whois: {0}".format(e))
|
||||||
|
|
||||||
self._socket.close()
|
self._socket.close()
|
||||||
|
|
||||||
if self.GATEWAY_IP is None:
|
if self.GATEWAY_IP is None:
|
||||||
_LOGGER.error('No Gateway found. Cannot continue')
|
_LOGGER.error('No Gateway found. Cannot continue')
|
||||||
return None
|
return None
|
||||||
@@ -149,14 +149,14 @@ class XiaomiHub:
|
|||||||
_LOGGER.error("Cannot connect to Gateway")
|
_LOGGER.error("Cannot connect to Gateway")
|
||||||
socket.close()
|
socket.close()
|
||||||
|
|
||||||
def write_to_hub(self, sid, data_key, datavalue):
|
def write_to_hub(self, sid, **values):
|
||||||
key = self._get_key()
|
key = self._get_key()
|
||||||
if type(datavalue) == int:
|
cmd = {
|
||||||
datavalue_formatted = str(datavalue)
|
"cmd": "write",
|
||||||
else:
|
"sid": sid,
|
||||||
datavalue_formatted = '"' + datavalue + '"'
|
"data": dict(key=key, **values)
|
||||||
cmd = '{ "cmd":"write","sid":"' + sid + '","data":"{"' + data_key + '":' + datavalue_formatted + ',"key":"' + key + '"}}'
|
}
|
||||||
return self._send_cmd(cmd, "write_ack")
|
return self._send_cmd(json.dumps(cmd), "write_ack")
|
||||||
|
|
||||||
def get_from_hub(self, sid):
|
def get_from_hub(self, sid):
|
||||||
cmd = '{ "cmd":"read","sid":"' + sid + '"}'
|
cmd = '{ "cmd":"read","sid":"' + sid + '"}'
|
||||||
|
|||||||
Reference in New Issue
Block a user