Fix whitespace errors to make pylint happy
This commit is contained in:
+6
-2
@@ -4,7 +4,7 @@ import threading
|
|||||||
import json
|
import json
|
||||||
import signal
|
import signal
|
||||||
|
|
||||||
#mine
|
# mine
|
||||||
import mqtt
|
import mqtt
|
||||||
import yamlparser
|
import yamlparser
|
||||||
from xiaomihub import XiaomiHub
|
from xiaomihub import XiaomiHub
|
||||||
@@ -12,6 +12,7 @@ from xiaomihub import XiaomiHub
|
|||||||
logging.basicConfig(level=logging.INFO)
|
logging.basicConfig(level=logging.INFO)
|
||||||
_LOGGER = logging.getLogger(__name__)
|
_LOGGER = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
def process_gateway_messages(gateway, client, stop_event):
|
def process_gateway_messages(gateway, client, stop_event):
|
||||||
while not stop_event.is_set():
|
while not stop_event.is_set():
|
||||||
try:
|
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.error('Error while sending from gateway to mqtt: ', str(e))
|
||||||
_LOGGER.info("Stopping Gateway Thread ...")
|
_LOGGER.info("Stopping Gateway Thread ...")
|
||||||
|
|
||||||
|
|
||||||
def read_motion_data(gateway, client, polling_interval, polling_models, stop_event):
|
def read_motion_data(gateway, client, polling_interval, polling_models, stop_event):
|
||||||
first = True
|
first = True
|
||||||
while not stop_event.is_set():
|
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)
|
time.sleep(polling_interval)
|
||||||
_LOGGER.info("Stopping Polling Thread ...")
|
_LOGGER.info("Stopping Polling Thread ...")
|
||||||
|
|
||||||
|
|
||||||
def process_mqtt_messages(gateway, client, stop_event):
|
def process_mqtt_messages(gateway, client, stop_event):
|
||||||
while not stop_event.is_set():
|
while not stop_event.is_set():
|
||||||
try:
|
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.error('Error while sending from mqtt to gateway: ', str(e))
|
||||||
_LOGGER.info("Stopping MQTT Thread ...")
|
_LOGGER.info("Stopping MQTT Thread ...")
|
||||||
|
|
||||||
|
|
||||||
def exit_handler(signal, frame):
|
def exit_handler(signal, frame):
|
||||||
print('Exiting')
|
print('Exiting')
|
||||||
stop_event.set()
|
stop_event.set()
|
||||||
@@ -102,7 +106,7 @@ if __name__ == "__main__":
|
|||||||
_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("gateway", "+", "write", None)
|
client.subscribe("gateway", "+", "write", None)
|
||||||
client.subscribe("plug", "+", "status", "set")
|
client.subscribe("plug", "+", "status", "set")
|
||||||
|
|||||||
+18
-17
@@ -6,6 +6,7 @@ import json
|
|||||||
|
|
||||||
_LOGGER = logging.getLogger(__name__)
|
_LOGGER = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
class Mqtt:
|
class Mqtt:
|
||||||
event_based_sensors = ["switch", "cube"]
|
event_based_sensors = ["switch", "cube"]
|
||||||
motion_sensors = ["motion", "sensor_motion.aq2"]
|
motion_sensors = ["motion", "sensor_motion.aq2"]
|
||||||
@@ -25,12 +26,12 @@ class Mqtt:
|
|||||||
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"
|
||||||
@@ -53,7 +54,7 @@ class Mqtt:
|
|||||||
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)
|
||||||
@@ -65,25 +66,25 @@ class Mqtt:
|
|||||||
def subscribe(self, model="+", name="+", prop="+", command=None):
|
def subscribe(self, model="+", name="+", prop="+", command=None):
|
||||||
topic = self.prefix + "/" + model + "/" + name + "/" + prop
|
topic = self.prefix + "/" + model + "/" + name + "/" + prop
|
||||||
if command is not None:
|
if command is not None:
|
||||||
topic += "/" + command
|
topic += "/" + command
|
||||||
_LOGGER.info("Subscribing to " + topic + ".")
|
_LOGGER.info("Subscribing to " + topic + ".")
|
||||||
self._client.subscribe(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)
|
||||||
|
|
||||||
items = {}
|
items = {}
|
||||||
for key, value in data.items():
|
for key, value in data.items():
|
||||||
# fix for latest motion value
|
# fix for latest motion value
|
||||||
if (model in self.motion_sensors and key == "no_motion"):
|
if (model in self.motion_sensors and key == "no_motion"):
|
||||||
key="status"
|
key = "status"
|
||||||
value="no_motion"
|
value = "no_motion"
|
||||||
if (model in self.magnet_sensors and key == "no_close"):
|
if (model in self.magnet_sensors 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 self.event_based_sensors):
|
if (model in self.event_based_sensors):
|
||||||
retain = False
|
retain = False
|
||||||
@@ -120,15 +121,15 @@ class Mqtt:
|
|||||||
return
|
return
|
||||||
|
|
||||||
model = parts[1]
|
model = parts[1]
|
||||||
query_sid = parts[2] #sid or name part
|
query_sid = parts[2] # sid or name part
|
||||||
param = parts[3] #param part
|
param = parts[3] # param part
|
||||||
method = None
|
method = None
|
||||||
if len(parts) > 4:
|
if len(parts) > 4:
|
||||||
method = parts[4]
|
method = parts[4]
|
||||||
else:
|
else:
|
||||||
method = parts[3]
|
method = parts[3]
|
||||||
|
|
||||||
name = "" # we will find it next
|
name = "" # we will find it next
|
||||||
sid = query_sid
|
sid = query_sid
|
||||||
isFound = False
|
isFound = False
|
||||||
for current_sid in self._sids:
|
for current_sid in self._sids:
|
||||||
@@ -168,7 +169,7 @@ class Mqtt:
|
|||||||
'values': {param: value}}
|
'values': {param: value}}
|
||||||
# put in process queuee
|
# put in process queuee
|
||||||
self._queue.put(data)
|
self._queue.put(data)
|
||||||
|
|
||||||
elif method == "write":
|
elif method == "write":
|
||||||
# use raw write method to the sensor, we expect a jsonified dict here.
|
# use raw write method to the sensor, we expect a jsonified dict here.
|
||||||
values = json.loads((msg.payload).decode('utf-8'))
|
values = json.loads((msg.payload).decode('utf-8'))
|
||||||
@@ -183,9 +184,9 @@ class Mqtt:
|
|||||||
|
|
||||||
def _color_xiaomi_to_rgb(self, xiaomi_color):
|
def _color_xiaomi_to_rgb(self, xiaomi_color):
|
||||||
intval = int(xiaomi_color)
|
intval = int(xiaomi_color)
|
||||||
blue = (intval) & 255
|
blue = (intval) & 255
|
||||||
green = (intval >> 8) & 255
|
green = (intval >> 8) & 255
|
||||||
red = (intval >> 16) & 255
|
red = (intval >> 16) & 255
|
||||||
bright = (intval >> 24) & 255
|
bright = (intval >> 24) & 255
|
||||||
value = str(red)+","+str(green)+","+str(blue)+","+str(bright)
|
value = str(red)+","+str(green)+","+str(blue)+","+str(bright)
|
||||||
return value
|
return value
|
||||||
@@ -195,7 +196,7 @@ class Mqtt:
|
|||||||
r = int(arr[0])
|
r = int(arr[0])
|
||||||
g = int(arr[1])
|
g = int(arr[1])
|
||||||
b = int(arr[2])
|
b = int(arr[2])
|
||||||
if len(arr)>3:
|
if len(arr) > 3:
|
||||||
bright = int(arr[3])
|
bright = int(arr[3])
|
||||||
else:
|
else:
|
||||||
bright = 255
|
bright = 255
|
||||||
|
|||||||
@@ -1,3 +1,3 @@
|
|||||||
paho-mqtt
|
paho-mqtt
|
||||||
pyyaml
|
pyyaml
|
||||||
pycrypto
|
pycrypto
|
||||||
|
|||||||
+15
-13
@@ -10,7 +10,8 @@ from threading import Thread
|
|||||||
|
|
||||||
_LOGGER = logging.getLogger(__name__)
|
_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:
|
class XiaomiHub:
|
||||||
GATEWAY_KEY = None
|
GATEWAY_KEY = None
|
||||||
@@ -60,11 +61,11 @@ 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
|
||||||
|
|
||||||
self._socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
self._socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
|
||||||
|
|
||||||
_LOGGER.info('Creating Multicast Socket')
|
_LOGGER.info('Creating Multicast Socket')
|
||||||
@@ -94,10 +95,10 @@ class XiaomiHub:
|
|||||||
model = resp["model"]
|
model = resp["model"]
|
||||||
|
|
||||||
xiaomi_device = {
|
xiaomi_device = {
|
||||||
"model":model,
|
"model": model,
|
||||||
"sid":resp["sid"],
|
"sid": resp["sid"],
|
||||||
"short_id":resp["short_id"],
|
"short_id": resp["short_id"],
|
||||||
"data":json.loads(resp["data"])}
|
"data": json.loads(resp["data"])}
|
||||||
|
|
||||||
device_type = None
|
device_type = None
|
||||||
if model in sensors:
|
if model in sensors:
|
||||||
@@ -107,7 +108,7 @@ class XiaomiHub:
|
|||||||
elif model in switches:
|
elif model in switches:
|
||||||
device_type = 'switch'
|
device_type = 'switch'
|
||||||
else:
|
else:
|
||||||
device_type = 'sensor' #not really matters
|
device_type = 'sensor' # not really matters
|
||||||
|
|
||||||
if device_type == None:
|
if device_type == None:
|
||||||
_LOGGER.error('Unsupported devices : {0}'.format(model))
|
_LOGGER.error('Unsupported devices : {0}'.format(model))
|
||||||
@@ -124,7 +125,7 @@ class XiaomiHub:
|
|||||||
try:
|
try:
|
||||||
socket = self._socket
|
socket = self._socket
|
||||||
socket_list = [sys.stdin, 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:
|
for sock in read_sockets:
|
||||||
if sock == socket:
|
if sock == socket:
|
||||||
data = sock.recv(4096)
|
data = sock.recv(4096)
|
||||||
@@ -163,11 +164,11 @@ class XiaomiHub:
|
|||||||
"sid": sid,
|
"sid": sid,
|
||||||
"data": dict(key=key, **values)
|
"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):
|
def get_from_hub(self, sid):
|
||||||
cmd = '{ "cmd":"read","sid":"' + sid + '"}'
|
cmd = '{ "cmd":"read","sid":"' + sid + '"}'
|
||||||
return self._send_cmd(cmd, "read_ack")
|
return self._send_cmd(cmd, "read_ack")
|
||||||
|
|
||||||
def _get_key(self):
|
def _get_key(self):
|
||||||
from Crypto.Cipher import AES
|
from Crypto.Cipher import AES
|
||||||
@@ -237,9 +238,9 @@ class XiaomiHub:
|
|||||||
if isinstance(packet, dict):
|
if isinstance(packet, dict):
|
||||||
try:
|
try:
|
||||||
sid = packet['sid']
|
sid = packet['sid']
|
||||||
#model = packet['model']
|
# model = packet['model']
|
||||||
data = json.loads(packet['data'])
|
data = json.loads(packet['data'])
|
||||||
|
|
||||||
for device in self.XIAOMI_HA_DEVICES[sid]:
|
for device in self.XIAOMI_HA_DEVICES[sid]:
|
||||||
device.push_data(data)
|
device.push_data(data)
|
||||||
|
|
||||||
@@ -248,6 +249,7 @@ class XiaomiHub:
|
|||||||
|
|
||||||
self._queue.task_done()
|
self._queue.task_done()
|
||||||
|
|
||||||
|
|
||||||
class XiaomiDevice():
|
class XiaomiDevice():
|
||||||
"""Representation a base Xiaomi device."""
|
"""Representation a base Xiaomi device."""
|
||||||
|
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ import logging
|
|||||||
|
|
||||||
_LOGGER = logging.getLogger(__name__)
|
_LOGGER = logging.getLogger(__name__)
|
||||||
|
|
||||||
|
|
||||||
def load_yaml(file):
|
def load_yaml(file):
|
||||||
try:
|
try:
|
||||||
stram = open(file, "r")
|
stram = open(file, "r")
|
||||||
@@ -12,6 +13,7 @@ def load_yaml(file):
|
|||||||
raise
|
raise
|
||||||
_LOGGER.error("Can't load yaml with sids %r (%r)" % (file, e))
|
_LOGGER.error("Can't load yaml with sids %r (%r)" % (file, e))
|
||||||
|
|
||||||
|
|
||||||
def get_gateway_password(config, ip=""):
|
def get_gateway_password(config, ip=""):
|
||||||
if (config == None):
|
if (config == None):
|
||||||
raise "Config is null"
|
raise "Config is null"
|
||||||
|
|||||||
Reference in New Issue
Block a user