diff --git a/src/main.py b/src/main.py index a1dbb4f..19102bd 100644 --- a/src/main.py +++ b/src/main.py @@ -4,7 +4,7 @@ import threading import json import signal -#mine +# mine import mqtt import yamlparser from xiaomihub import XiaomiHub @@ -12,6 +12,7 @@ from xiaomihub import XiaomiHub logging.basicConfig(level=logging.INFO) _LOGGER = logging.getLogger(__name__) + def process_gateway_messages(gateway, client, stop_event): while not stop_event.is_set(): try: @@ -32,6 +33,7 @@ def process_gateway_messages(gateway, client, stop_event): _LOGGER.error('Error while sending from gateway to mqtt: ', str(e)) _LOGGER.info("Stopping Gateway Thread ...") + def read_motion_data(gateway, client, polling_interval, polling_models, stop_event): first = True while not stop_event.is_set(): @@ -62,6 +64,7 @@ def read_motion_data(gateway, client, polling_interval, polling_models, stop_eve time.sleep(polling_interval) _LOGGER.info("Stopping Polling Thread ...") + def process_mqtt_messages(gateway, client, stop_event): while not stop_event.is_set(): try: @@ -79,6 +82,7 @@ def process_mqtt_messages(gateway, client, stop_event): _LOGGER.error('Error while sending from mqtt to gateway: ', str(e)) _LOGGER.info("Stopping MQTT Thread ...") + def exit_handler(signal, frame): print('Exiting') stop_event.set() @@ -102,7 +106,7 @@ if __name__ == "__main__": _LOGGER.info("Init mqtt client.") client = mqtt.Mqtt(config) 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", "+", "write", None) client.subscribe("plug", "+", "status", "set") diff --git a/src/mqtt.py b/src/mqtt.py index 799f7f5..4abaded 100644 --- a/src/mqtt.py +++ b/src/mqtt.py @@ -6,6 +6,7 @@ import json _LOGGER = logging.getLogger(__name__) + class Mqtt: event_based_sensors = ["switch", "cube"] motion_sensors = ["motion", "sensor_motion.aq2"] @@ -25,12 +26,12 @@ class Mqtt: if (config == None): raise "Config is null" - #load sids dictionary + # load sids dictionary self._sids = config.get("sids", None) if (self._sids == None): self._sids = dict({}) - #load mqtt settings + # load mqtt settings mqttConfig = config.get("mqtt", None) if (mqttConfig == None): raise "Config mqtt section is null" @@ -53,7 +54,7 @@ class Mqtt: self._client.on_connect = self._mqtt_on_connect self._client.connect(self.server, self.port, 60) - #run message processing loop + # run message processing loop t1 = Thread(target=self._mqtt_loop) t1.start() self._threads.append(t1) @@ -65,25 +66,25 @@ class Mqtt: def subscribe(self, model="+", name="+", prop="+", command=None): topic = self.prefix + "/" + model + "/" + name + "/" + prop if command is not None: - topic += "/" + command + topic += "/" + command _LOGGER.info("Subscribing to " + topic + ".") self._client.subscribe(topic) def publish(self, model, sid, data, retain=True): sidprops = self._sids.get(sid, None) if (sidprops != None): - model = sidprops.get("model",model) - sid = sidprops.get("name",sid) + model = sidprops.get("model", model) + sid = sidprops.get("name", sid) items = {} for key, value in data.items(): # fix for latest motion value if (model in self.motion_sensors and key == "no_motion"): - key="status" - value="no_motion" + key = "status" + value = "no_motion" if (model in self.magnet_sensors and key == "no_close"): - key="status" - value="open" + key = "status" + value = "open" # do not retain event-based sensors (like switches and cubes). if (model in self.event_based_sensors): retain = False @@ -120,15 +121,15 @@ class Mqtt: return model = parts[1] - query_sid = parts[2] #sid or name part - param = parts[3] #param part + query_sid = parts[2] # sid or name part + param = parts[3] # param part method = None if len(parts) > 4: method = parts[4] else: method = parts[3] - name = "" # we will find it next + name = "" # we will find it next sid = query_sid isFound = False for current_sid in self._sids: @@ -168,7 +169,7 @@ class Mqtt: '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')) @@ -183,9 +184,9 @@ class Mqtt: def _color_xiaomi_to_rgb(self, xiaomi_color): intval = int(xiaomi_color) - blue = (intval) & 255 + blue = (intval) & 255 green = (intval >> 8) & 255 - red = (intval >> 16) & 255 + red = (intval >> 16) & 255 bright = (intval >> 24) & 255 value = str(red)+","+str(green)+","+str(blue)+","+str(bright) return value @@ -195,7 +196,7 @@ class Mqtt: r = int(arr[0]) g = int(arr[1]) b = int(arr[2]) - if len(arr)>3: + if len(arr) > 3: bright = int(arr[3]) else: bright = 255 diff --git a/src/requirements.txt b/src/requirements.txt index ab7766c..e64b84d 100644 --- a/src/requirements.txt +++ b/src/requirements.txt @@ -1,3 +1,3 @@ paho-mqtt pyyaml -pycrypto \ No newline at end of file +pycrypto diff --git a/src/xiaomihub.py b/src/xiaomihub.py index 24ccb06..69ef87c 100644 --- a/src/xiaomihub.py +++ b/src/xiaomihub.py @@ -10,7 +10,8 @@ from threading import Thread _LOGGER = logging.getLogger(__name__) -# MANDATORY!!!! NEED TO TURN OFF "_process_report" THREAD IF CODE IS UPDATED!!!! +# MANDATORY!!!! NEED TO TURN OFF "_process_report" THREAD IF CODE IS UPDATED!!! + class XiaomiHub: GATEWAY_KEY = None @@ -60,11 +61,11 @@ class XiaomiHub: _LOGGER.error("Cannot discover hub using whois: {0}".format(e)) self._socket.close() - + if self.GATEWAY_IP is None: _LOGGER.error('No Gateway found. Cannot continue') return None - + self._socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) _LOGGER.info('Creating Multicast Socket') @@ -84,8 +85,10 @@ class XiaomiHub: _LOGGER.info('Found {0} devices'.format(len(sids))) - sensors = ['sensor_ht', 'sensor_wleak.aq1'] - binary_sensors = ['magnet', 'motion', 'switch', '86sw1', '86sw2', 'cube'] + sensors = ['sensor_ht', 'weather.v1', 'sensor_wleak.aq1'] + binary_sensors = ['magnet', 'sensor_magnet.aq2', 'motion', + 'sensor_motion.aq2', 'switch', 'sensor_switch.aq2', + '86sw1', '86sw2', 'cube'] switches = ['plug', 'ctrl_neutral1', 'ctrl_neutral2'] for sid in sids: @@ -94,10 +97,10 @@ class XiaomiHub: model = resp["model"] xiaomi_device = { - "model":model, - "sid":resp["sid"], - "short_id":resp["short_id"], - "data":json.loads(resp["data"])} + "model": model, + "sid": resp["sid"], + "short_id": resp["short_id"], + "data": json.loads(resp["data"])} device_type = None if model in sensors: @@ -107,12 +110,9 @@ class XiaomiHub: elif model in switches: device_type = 'switch' else: - device_type = 'sensor' #not really matters + device_type = 'sensor' # not really matters - if device_type == None: - _LOGGER.error('Unsupported devices : {0}'.format(model)) - else: - self.XIAOMI_DEVICES[device_type].append(xiaomi_device) + self.XIAOMI_DEVICES[device_type].append(xiaomi_device) def _send_cmd(self, cmd, rtnCmd): return self._send_socket(cmd, rtnCmd, self.GATEWAY_IP, self.GATEWAY_PORT) @@ -124,7 +124,7 @@ class XiaomiHub: try: socket = self._socket socket_list = [sys.stdin, socket] - read_sockets, write_sockets, error_sockets = select.select(socket_list , [], []) + read_sockets, write_sockets, error_sockets = select.select(socket_list, [], []) for sock in read_sockets: if sock == socket: data = sock.recv(4096) @@ -163,11 +163,11 @@ class XiaomiHub: "sid": sid, "data": dict(key=key, **values) } - return self._send_cmd(json.dumps(cmd), "write_ack") + return self._send_cmd(json.dumps(cmd), "write_ack") def get_from_hub(self, sid): cmd = '{ "cmd":"read","sid":"' + sid + '"}' - return self._send_cmd(cmd, "read_ack") + return self._send_cmd(cmd, "read_ack") def _get_key(self): from Crypto.Cipher import AES @@ -237,9 +237,9 @@ class XiaomiHub: if isinstance(packet, dict): try: sid = packet['sid'] - #model = packet['model'] + # model = packet['model'] data = json.loads(packet['data']) - + for device in self.XIAOMI_HA_DEVICES[sid]: device.push_data(data) @@ -248,6 +248,7 @@ class XiaomiHub: self._queue.task_done() + class XiaomiDevice(): """Representation a base Xiaomi device.""" diff --git a/src/yamlparser.py b/src/yamlparser.py index 3f12aea..7d190f6 100644 --- a/src/yamlparser.py +++ b/src/yamlparser.py @@ -3,6 +3,7 @@ import logging _LOGGER = logging.getLogger(__name__) + def load_yaml(file): try: stram = open(file, "r") @@ -12,6 +13,7 @@ def load_yaml(file): raise _LOGGER.error("Can't load yaml with sids %r (%r)" % (file, e)) + def get_gateway_password(config, ip=""): if (config == None): raise "Config is null"