63 Commits
Author SHA1 Message Date
monster 41dd77c37a Merge pull request #24 from jerryzz/master
Update mqtt.py
2021-03-14 02:16:48 +03:00
jerryzz ad2506c0f2 Update mqtt.py
Add "no_motion" publish = no-motion time 120, 180, ....(s)
2021-03-14 00:14:06 +01:00
monster 5c033b962e Update .semver
updated to 1.0.5
2018-10-27 15:31:25 +03:00
monster 0722923fcd Merge pull request #20 from seblucas/env-var-config-2
Add support for environment variable to load the configuration
2018-10-27 15:30:51 +03:00
monster a9f2966423 Merge pull request #19 from techdada/master
mqtt client ssl enabled
2018-10-27 15:30:18 +03:00
root 56daf48162 Fix gateway color adjustment when prefix more than one subtopic 2018-10-25 22:28:13 +00:00
Sébastien Lucas 096a864213 [feat] The configuration can also be loaded from an environment variable 2018-10-23 20:27:52 +02:00
techdada 7f56949030 mqtt client ssl enabled 2018-10-21 11:25:44 +00:00
monster ea16a6306c Merge pull request #17 from groggemans/add_sensor_models
Add sensor models
2018-03-06 00:31:16 +03:00
G. Roggemans 88396d14dc Add missing sensor types
- weather.v1
- sensor_magnet.aq2
- sensor_motion.aq2
- sensor_switch.aq2
2018-03-05 18:38:45 +01:00
G. Roggemans 74b7cee069 Remove unreachable code for logging model errors 2018-03-05 18:32:51 +01:00
G. Roggemans 38bb43086b Fix whitespace errors to make pylint happy 2018-03-05 18:31:17 +01:00
monster 9b4f784532 Merge pull request #15 from eburtsev/master
Fixed 'int' object has no attribute 'isdigit' exception
2018-02-27 16:13:55 +03:00
Eugene Burtsev e6bce68641 Fixed 'int' object has no attribute 'isdigit' exception 2018-02-27 15:21:11 +03:00
raspberry b9e940e9de Fix linter errors 2018-02-22 11:03:57 +03:00
raspberry 0f9a759839 Compile before 2018-02-22 10:10:08 +03:00
raspberry 3e5e71035d _read_unwanted_data_fix is now controlled via settings 2018-02-22 09:39:54 +03:00
raspberry fadd23e8cd Readme change 2018-02-22 09:03:19 +03:00
raspberry 04169ffe4e Fix for new aqara motion sensor 2018-02-22 06:39:55 +03:00
raspberry b390d044cc Fix for new aqara motion sensor 2018-02-21 21:05:24 +03:00
raspberry da10280a34 SORRY! Fix for broken motion/door state! 2018-02-21 20:01:12 +03:00
raspberry 7a2d601481 SORRY! Fix for broken motion/door state! 2018-02-21 19:56:31 +03:00
raspberry ed116b431b SORRY! Fix for broken motion/door state! 2018-02-21 19:55:47 +03:00
raspberry 6f8bbe1599 1.0.2 2018-02-18 19:18:25 +03:00
monster 86967d7599 Merge pull request #14 from eburtsev/master
Added ability to set gateway discovery IP in config
2018-02-18 19:17:33 +03:00
Eugene Burtsev d638b89fc2 Added ability to set gateway discovery IP in config 2018-02-17 15:58:15 +03:00
Eugene Burtsev 3561d4e595 Merge pull request #1 from monster1025/master
Changes from upstream
2018-02-17 15:48:15 +03:00
Your Name f0eda582e1 travis docker layers cache - cant get it working ( 2018-02-16 13:34:36 +03:00
Your Name 219df48228 travis docker layers cache 2018-02-16 13:26:48 +03:00
Your Name 693f3a21de travis docker layers cache 2018-02-16 12:52:56 +03:00
Your Name a0979f4e39 travis docker layers cache 2018-02-16 12:15:45 +03:00
Your Name d47548ce24 README change. 2018-02-16 11:12:45 +03:00
Your Name 3754c8d08d Added support for JSON report instead of properties. 2018-02-16 11:08:36 +03:00
Your Name 8b4ea2a089 Added support for JSON report instead of properties. 2018-02-16 11:08:17 +03:00
Your Name 076e70de6b travis 2018-02-16 10:49:59 +03:00
Your Name a1fbbd0922 travis 2018-02-16 10:40:11 +03:00
Your Name b48c2abfd3 travis 2018-02-16 10:39:05 +03:00
Your Name 67b465fac4 travis 2018-02-16 10:18:29 +03:00
Your Name 147fe4dc78 travis 2018-02-16 10:11:55 +03:00
Your Name 810fe2b4b1 travis 2018-02-16 10:01:39 +03:00
Your Name 85a53d031e travis 2018-02-16 09:40:11 +03:00
Your Name 74f3d3bce6 travis 2018-02-16 09:25:55 +03:00
Your Name 848a36f9cd travis 2018-02-16 09:16:16 +03:00
Your Name 3b3effc67d travis 2018-02-16 09:12:35 +03:00
Your Name e129152f79 travis 2018-02-16 08:55:53 +03:00
Your Name a845c49bcd travis 2018-02-16 08:50:04 +03:00
Your Name 1e6c131f81 travis 2018-02-16 01:07:46 +03:00
Your Name b0082fe61b travis 2018-02-16 01:02:36 +03:00
Your Name 7ee0cfadc7 travis 2018-02-16 00:58:36 +03:00
Your Name 6bd16ded60 travis 2018-02-16 00:53:39 +03:00
Your Name 32c8aefb76 travis 2018-02-16 00:50:58 +03:00
Your Name f647e97b7f travis 2018-02-16 00:42:39 +03:00
Your Name 10900725a1 travis 2018-02-16 00:37:39 +03:00
Your Name 080de01b37 travis 2018-02-16 00:37:24 +03:00
Your Name 5fd253a7f5 travis 2018-02-16 00:32:49 +03:00
Your Name 3b4e735cf2 travis 2018-02-16 00:28:28 +03:00
Your Name e83ce64a47 travis 2018-02-16 00:22:57 +03:00
Your Name f4d8b0e529 travis 2018-02-16 00:19:10 +03:00
Your Name c25dadfeb2 travis 2018-02-16 00:08:44 +03:00
Your Name 5182c22f1a travis 2018-02-15 23:54:00 +03:00
Your Name a32cc13bc3 travis 2018-02-15 23:43:06 +03:00
Your Name 9cabb1adb2 travis 2018-02-15 23:38:04 +03:00
Your Name cf703caad8 travis 2018-02-15 23:29:17 +03:00
13 changed files with 317 additions and 112 deletions
+1
View File
@@ -0,0 +1 @@
1.0.5
+117
View File
@@ -0,0 +1,117 @@
sudo: true
language: python
services:
- docker
env:
global:
- PROJECT=aqara-mqtt
- IMAGE_NAME=$DOCKER_USERNAME/$PROJECT
- VERSION=$(cat .semver)
- MAJOR=$(cut -d. -f1 .semver)
- MINOR=$(cut -d. -f2 .semver)
- PATCH=$(cut -d. -f3 .semver)
jobs:
include:
- stage: build x64 image
before_script:
- ARCH=x64
script:
# build
- docker build -t $PROJECT -f Dockerfile-$ARCH .
# push image
- >
if [ "$TRAVIS_BRANCH" == "master" ] && [ "$TRAVIS_PULL_REQUEST" == "false" ]; then
docker login -u="$DOCKER_USERNAME" -p="$DOCKER_PASSWORD"
#1.0.0-arch
docker tag $PROJECT $IMAGE_NAME:$VERSION-$ARCH
docker push $IMAGE_NAME:$VERSION-$ARCH
#1.0-arch
docker tag $PROJECT $IMAGE_NAME:$MAJOR.$MINOR-$ARCH
docker push $IMAGE_NAME:$MAJOR.$MINOR-$ARCH
#1-arch
docker tag $PROJECT $IMAGE_NAME:$MAJOR-$ARCH
docker push $IMAGE_NAME:$MAJOR-$ARCH
#latest by arch
docker tag $PROJECT $IMAGE_NAME:$ARCH
docker push $IMAGE_NAME:$ARCH
fi
- stage: build i386 image
before_script:
- ARCH=i386
script:
# build
- docker build -t $PROJECT -f Dockerfile-$ARCH .
# push image
- >
if [ "$TRAVIS_BRANCH" == "master" ] && [ "$TRAVIS_PULL_REQUEST" == "false" ]; then
docker login -u="$DOCKER_USERNAME" -p="$DOCKER_PASSWORD"
#1.0.0-arch
docker tag $PROJECT $IMAGE_NAME:$VERSION-$ARCH
docker push $IMAGE_NAME:$VERSION-$ARCH
#1.0-arch
docker tag $PROJECT $IMAGE_NAME:$MAJOR.$MINOR-$ARCH
docker push $IMAGE_NAME:$MAJOR.$MINOR-$ARCH
#1-arch
docker tag $PROJECT $IMAGE_NAME:$MAJOR-$ARCH
docker push $IMAGE_NAME:$MAJOR-$ARCH
#latest by arch
docker tag $PROJECT $IMAGE_NAME:$ARCH
docker push $IMAGE_NAME:$ARCH
fi
- stage: build armhf image
before_script:
- ARCH=armhf
script:
# prepare qemu
- docker run --rm --privileged multiarch/qemu-user-static:register --reset
# get qemu-arm-static binary
- >
mkdir tmp &&
pushd tmp &&
curl -L -o qemu-arm-static.tar.gz https://github.com/multiarch/qemu-user-static/releases/download/v2.6.0/qemu-arm-static.tar.gz &&
tar xzf qemu-arm-static.tar.gz &&
popd
#copy it to container
- sed -i -e '2iCOPY tmp/qemu-arm-static /usr/bin/qemu-arm-static\' Dockerfile-$ARCH
# build
- docker build -t $PROJECT -f Dockerfile-$ARCH .
# push image
- >
if [ "$TRAVIS_BRANCH" == "master" ] && [ "$TRAVIS_PULL_REQUEST" == "false" ]; then
docker login -u="$DOCKER_USERNAME" -p="$DOCKER_PASSWORD"
#1.0.0-arch
docker tag $PROJECT $IMAGE_NAME:$VERSION-$ARCH
docker push $IMAGE_NAME:$VERSION-$ARCH
#1.0-arch
docker tag $PROJECT $IMAGE_NAME:$MAJOR.$MINOR-$ARCH
docker push $IMAGE_NAME:$MAJOR.$MINOR-$ARCH
#1-arch
docker tag $PROJECT $IMAGE_NAME:$MAJOR-$ARCH
docker push $IMAGE_NAME:$MAJOR-$ARCH
#latest by arch
docker tag $PROJECT $IMAGE_NAME:$ARCH
docker push $IMAGE_NAME:$ARCH
fi
+15
View File
@@ -0,0 +1,15 @@
FROM arm32v7/python:slim
ENV LIBRARY_PATH=/lib:/usr/lib
ADD src/requirements.txt /
RUN apt-get update && apt-get install -y build-essential autoconf \
&& pip install --upgrade pip && pip install -r /requirements.txt \
&& apt-get remove -y build-essential autoconf \
&& apt-get autoremove -y \
&& rm -rf /var/lib/apt/lists/*
WORKDIR /app
COPY src /app
CMD ["python3", "-u", "/app/main.py"]
-11
View File
@@ -1,11 +0,0 @@
FROM alpine-python:latest
ENV LIBRARY_PATH=/lib:/usr/lib
ADD src/requirements.txt /
RUN pip install --upgrade pip && pip install -r /requirements.txt
WORKDIR /app
COPY src /app
CMD ["python3", "-u", "/app/main.py"]
+14
View File
@@ -0,0 +1,14 @@
ARGS = `arg="$(filter-out $@,$(MAKECMDGOALS))" && echo $${arg:-${1}}`
lint:
rm -rf "src/__pycache__"
python3 -m compileall src
rm -rf "src/__pycache__"
commit: lint
git add .
git commit -m "$(call ARGS,\"updating to lastest local code\")"
git push
%:
@:
+10 -2
View File
@@ -1,5 +1,7 @@
# Aqara-MQTT # Aqara-MQTT
Aqara (Xiaomi) Gateway to MQTT brodge. [![Build Status](https://travis-ci.org/monster1025/aqara-mqtt.svg?branch=master)](https://travis-ci.org/monster1025/aqara-mqtt)
Aqara (Xiaomi) Gateway to MQTT bridge.
I use it for home assistant integration and it works well now. I use it for home assistant integration and it works well now.
You need to activate developer mode (described here: http://bbs.xiaomi.cn/t-13198850) You need to activate developer mode (described here: http://bbs.xiaomi.cn/t-13198850)
@@ -14,6 +16,12 @@ will turn on plug/heater and translate devices state from gateway:
"home/plug/heater/status" on "home/plug/heater/status" on
``` ```
## Architecture
Docker image support following architectures (you must choose your architecture in docker-compose):
- armhf (raspberry pi 3, arm32v7)
- i386 (x86 pc)
- x64 (x64 pc)
## Config ## Config
Edit file config/config-sample.yaml and rename it to config/config.yaml Edit file config/config-sample.yaml and rename it to config/config.yaml
@@ -21,7 +29,7 @@ Edit file config/config-sample.yaml and rename it to config/config.yaml
Sample docker-compose.yaml file for user: Sample docker-compose.yaml file for user:
``` ```
aqara: aqara:
image: monster1025/aqara-mqtt image: monster1025/aqara-mqtt:1-armhf
container_name: aqara container_name: aqara
volumes: volumes:
- "./config:/app/config" - "./config:/app/config"
+7
View File
@@ -4,9 +4,16 @@ mqtt:
username: username username: username
password: passw0rd password: passw0rd
prefix: home prefix: home
#secure mqtt. uncomment to enable ssl:
#ca: "config/roots.pem"
#tls_version: "tlsv1.2"
#send report as json
json: False
gateway: gateway:
ip: 192.168.0.47
password: passw0rd password: passw0rd
unwanted_data_fix: True #leave it True if it works for you. Otherwise (if you didn't recieve any data from sensors) - set it to False
polling_interval: 2 polling_interval: 2
polling_models: polling_models:
- motion - motion
+2 -1
View File
@@ -1,6 +1,7 @@
aqara: aqara:
#image: monster1025/aqara-mqtt:1-armhf
build: . build: .
dockerfile: Dockerfile-rpi dockerfile: Dockerfile-armhf
container_name: aqara container_name: aqara
volumes: volumes:
- "./config:/app/config" - "./config:/app/config"
+19 -17
View File
@@ -1,11 +1,10 @@
import logging import logging
import time import time
import threading import threading
import os
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
@@ -13,9 +12,10 @@ 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:
packet = gateway._queue.get() packet = gateway._queue.get()
if packet is None: if packet is None:
continue continue
@@ -33,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():
@@ -41,21 +42,19 @@ def read_motion_data(gateway, client, polling_interval, polling_models, stop_eve
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 is 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) if device['data'] != data or first:
short_id = sensor_resp['short_id']
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)
@@ -65,9 +64,10 @@ 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:
data = client._queue.get() data = client._queue.get()
if data is None: if data is None:
continue continue
@@ -76,12 +76,13 @@ def process_mqtt_messages(gateway, client, stop_event):
sid = data.get("sid", None) sid = data.get("sid", None)
values = data.get("values", dict()) values = data.get("values", dict())
resp = gateway.write_to_hub(sid, **values) 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))
_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()
@@ -93,10 +94,11 @@ def exit_handler(signal, frame):
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'])
gateway_ip = config['gateway'].get("ip", None)
signal.signal(signal.SIGINT, exit_handler) signal.signal(signal.SIGINT, exit_handler)
signal.signal(signal.SIGTERM, exit_handler) signal.signal(signal.SIGTERM, exit_handler)
@@ -104,13 +106,13 @@ 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")
gateway = XiaomiHub(gateway_pass) gateway = XiaomiHub(gateway_pass, gateway_ip, config)
stop_event= threading.Event() stop_event = threading.Event()
t1 = threading.Thread(target=process_gateway_messages, args=[gateway, client, stop_event]) t1 = threading.Thread(target=process_gateway_messages, args=[gateway, client, stop_event])
t1.daemon = True t1.daemon = True
t1.start() t1.start()
+84 -45
View File
@@ -1,19 +1,26 @@
import paho.mqtt.client as mqtt import paho.mqtt.client as mqtt
import os
import logging import logging
import os
import ssl
from queue import Queue from queue import Queue
from threading import Thread from threading import Thread
import json import json
_LOGGER = logging.getLogger(__name__) _LOGGER = logging.getLogger(__name__)
class Mqtt: class Mqtt:
event_based_sensors = ["switch", "cube"]
motion_sensors = ["motion", "sensor_motion.aq2"]
magnet_sensors = ["magnet"]
username = "" username = ""
password = "" password = ""
server = "localhost" server = "localhost"
port = 1883 port = 1883
ca = None
tlsvers = None
prefix = "home" prefix = "home"
_client = None _client = None
_sids = None _sids = None
_queue = None _queue = None
@@ -23,12 +30,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"
@@ -38,6 +45,11 @@ class Mqtt:
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.ca = mqttConfig.get("ca",None)
self.tlsvers = self._get_tls_version(
mqttConfig.get("tls_version","tlsv1.2")
)
self.json = mqttConfig.get("json", False)
self._queue = Queue() self._queue = Queue()
self._threads = [] self._threads = []
@@ -48,9 +60,17 @@ class Mqtt:
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) if (self.ca != None):
self._client.tls_set(
ca_certs=self.ca,
cert_reqs=ssl.CERT_REQUIRED,
tls_version=self.tlsvers
)
#run message processing loop self._client.tls_insecure_set(False)
self._client.connect(self.server, self.port, 60)
# 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)
@@ -62,59 +82,76 @@ 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)
# _LOGGER.info("data is " + format(data)) items = {}
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 in self.motion_sensors and key == "status"):
key="status" items["no_motion"] = 0
value="no_motion" if (model in self.motion_sensors and key == "no_motion"):
if (model == "magnet" and key == "no_close"): items[key] = value
key="status" key = "status"
value="open" value = "no_motion"
if (model in self.magnet_sensors and key == "no_close"):
key = "status"
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 self.event_based_sensors):
retain = False retain = False
# fix for rgb format # fix for rgb format
if (key == "rgb" and self._is_int(value)): if (key == "rgb" and str(value).isdigit()):
value = self._color_xiaomi_to_rgb(value) value = self._color_xiaomi_to_rgb(str(value))
items[key] = value
topic = PATH_FMT.format(model=model, sid=sid, prop=key) if self.json == True:
_LOGGER.info("Publishing message to topic " + topic + ": " + str(value) + ".") PATH_FMT = self.prefix + "/{model}/{sid}/json"
self._client.publish(topic, payload=value, qos=0, retain=retain) topic = PATH_FMT.format(model=model, sid=sid)
values = {}
values['sid'] = sid
for key in items:
values[key] = items[key]
jsondata = json.dumps(values)
_LOGGER.info("Publishing message to topic " + topic + ": " + str(jsondata) + ".")
self._client.publish(topic, payload=jsondata, qos=0, retain=retain)
else:
for key in items:
PATH_FMT = self.prefix + "/{model}/{sid}/{prop}"
topic = PATH_FMT.format(model=model, sid=sid, prop=key)
_LOGGER.info("Publishing message to topic " + topic + ": " + str(items[key]) + ".")
self._client.publish(topic, payload=items[key], 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("/")
if len(parts) < 4: # need to strip prefix to make parts assignment reliable
parts = msg.topic.replace(self.prefix+"/","").split("/")
partlen = len(parts)
if len(parts) < 3:
# should we return an error message ? # should we return an error message ?
return return
model = parts[1] model = parts[0]
query_sid = parts[2] #sid or name part query_sid = parts[1] # sid or name part
param = parts[3] #param part param = parts[2] # param part
method = None method = None
if len(parts) > 4: if len(parts) > 3:
method = parts[4]
else:
method = parts[3] method = parts[3]
else:
method = parts[2]
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:
@@ -129,6 +166,7 @@ class Mqtt:
sid = current_sid sid = current_sid
name = sidname name = sidname
isFound = True isFound = True
_LOGGER.debug("Found " + sid + " = " + name)
break break
else: else:
_LOGGER.debug(sidmodel + "-" + sidname + " is not " + model + "-" + query_sid + ".") _LOGGER.debug(sidmodel + "-" + sidname + " is not " + model + "-" + query_sid + ".")
@@ -142,7 +180,7 @@ class Mqtt:
# use single value set method # use single value set method
value = (msg.payload).decode('utf-8') value = (msg.payload).decode('utf-8')
if self._is_int(value): if value.isdigit():
value = int(value) value = int(value)
# fix for rgb format # fix for rgb format
@@ -154,7 +192,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'))
@@ -169,9 +207,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
@@ -181,16 +219,17 @@ 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
value = int('%02x%02x%02x%02x' % (bright, r, g, b), 16) value = int('%02x%02x%02x%02x' % (bright, r, g, b), 16)
return value return value
def _is_int(self, x): def _get_tls_version(self,tlsString):
try: switcher = {
tmp = int(x) "tlsv1": ssl.PROTOCOL_TLSv1,
return True "tlsv1.1": ssl.PROTOCOL_TLSv1_1,
except Exception as e: "tlsv1.2": ssl.PROTOCOL_TLSv1_2
return False }
return switcher.get(tlsString,ssl.PROTOCOL_TLSv1_2)
+1 -1
View File
@@ -1,3 +1,3 @@
paho-mqtt paho-mqtt
pyyaml pyyaml
pycrypto pycrypto
+31 -25
View File
@@ -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
@@ -28,7 +29,7 @@ class XiaomiHub:
GATEWAY_DISCOVERY_PORT = 4321 GATEWAY_DISCOVERY_PORT = 4321
SOCKET_BUFSIZE = 1024 SOCKET_BUFSIZE = 1024
def __init__(self, key, gateway=None): def __init__(self, key, gateway_ip=None, config=None):
self.GATEWAY_KEY = key self.GATEWAY_KEY = key
self._listening = False self._listening = False
self._queue = None self._queue = None
@@ -36,9 +37,14 @@ class XiaomiHub:
self._mcastsocket = None self._mcastsocket = None
self._deviceCallbacks = defaultdict(list) self._deviceCallbacks = defaultdict(list)
self._threads = [] self._threads = []
self._read_unwanted_data_enabled = True
if gateway is not None: if gateway_ip is not None:
self.GATEWAY_DISCOVERY_ADDRESS = gateway self.GATEWAY_DISCOVERY_ADDRESS = gateway_ip
if config is not None and 'gateway' in config and 'unwanted_data_fix' in config['gateway']:
self._read_unwanted_data_enabled = (config['gateway']['unwanted_data_fix'] == True)
_LOGGER.info('"Read unwanted data" fix is {0}'.format(self._read_unwanted_data_enabled))
try: try:
_LOGGER.info('Discovering Xiaomi Gateways using address {0}'.format(self.GATEWAY_DISCOVERY_ADDRESS)) _LOGGER.info('Discovering Xiaomi Gateways using address {0}'.format(self.GATEWAY_DISCOVERY_ADDRESS))
@@ -55,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')
@@ -79,8 +85,10 @@ class XiaomiHub:
_LOGGER.info('Found {0} devices'.format(len(sids))) _LOGGER.info('Found {0} devices'.format(len(sids)))
sensors = ['sensor_ht', 'sensor_wleak.aq1'] sensors = ['sensor_ht', 'weather.v1', 'sensor_wleak.aq1']
binary_sensors = ['magnet', 'motion', 'switch', '86sw1', '86sw2', 'cube'] binary_sensors = ['magnet', 'sensor_magnet.aq2', 'motion',
'sensor_motion.aq2', 'switch', 'sensor_switch.aq2',
'86sw1', '86sw2', 'cube']
switches = ['plug', 'ctrl_neutral1', 'ctrl_neutral2'] switches = ['plug', 'ctrl_neutral1', 'ctrl_neutral2']
for sid in sids: for sid in sids:
@@ -88,14 +96,11 @@ class XiaomiHub:
resp = self._send_cmd(cmd, "read_ack") resp = self._send_cmd(cmd, "read_ack")
model = resp["model"] model = resp["model"]
if model == '':
model = 'cube'
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:
@@ -105,21 +110,21 @@ 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: self.XIAOMI_DEVICES[device_type].append(xiaomi_device)
_LOGGER.error('Unsupported devices : {0}'.format(model))
else:
self.XIAOMI_DEVICES[device_type].append(xiaomi_device)
def _send_cmd(self, cmd, rtnCmd): def _send_cmd(self, cmd, rtnCmd):
return self._send_socket(cmd, rtnCmd, self.GATEWAY_IP, self.GATEWAY_PORT) return self._send_socket(cmd, rtnCmd, self.GATEWAY_IP, self.GATEWAY_PORT)
def _read_unwanted_data(self): def _read_unwanted_data(self):
if not self._read_unwanted_data_enabled:
return
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)
@@ -158,11 +163,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
@@ -232,9 +237,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)
@@ -243,6 +248,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."""
+16 -10
View File
@@ -1,24 +1,30 @@
import yaml import yaml
from os import environ
import logging import logging
_LOGGER = logging.getLogger(__name__) _LOGGER = logging.getLogger(__name__)
def load_yaml(file): def load_yaml(file):
try: try:
stram = open(file, "r") if environ.get("AQARA_MQTT_CONFIG") is not None:
stram = environ.get("AQARA_MQTT_CONFIG")
else:
stram = open(file, "r")
yaml_data = yaml.load(stram) yaml_data = yaml.load(stram)
return yaml_data return yaml_data
except Exception as e: except Exception as e:
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"
configGateway = config.get("gateway", None) configGateway = config.get("gateway", None)
if (configGateway == None): if (configGateway == None):
raise "Config gateway is null" raise "Config gateway is null"
password = configGateway.get("password", None) password = configGateway.get("password", None)
if (password == None): if (password == None):
raise "Config gateway passowrd is null" raise "Config gateway passowrd is null"
return password return password