Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
41dd77c37a | ||
|
|
ad2506c0f2 | ||
|
|
5c033b962e | ||
|
|
0722923fcd | ||
|
|
a9f2966423 | ||
|
|
56daf48162 | ||
|
|
096a864213 | ||
|
|
7f56949030 | ||
|
|
ea16a6306c | ||
|
|
88396d14dc | ||
|
|
74b7cee069 | ||
|
|
38bb43086b | ||
|
|
9b4f784532 | ||
|
|
e6bce68641 | ||
|
|
b9e940e9de | ||
|
|
0f9a759839 | ||
|
|
3e5e71035d | ||
|
|
fadd23e8cd | ||
|
|
04169ffe4e | ||
|
|
b390d044cc | ||
|
|
da10280a34 | ||
|
|
7a2d601481 | ||
|
|
ed116b431b | ||
|
|
6f8bbe1599 | ||
|
|
86967d7599 | ||
|
|
d638b89fc2 | ||
|
|
3561d4e595 | ||
|
|
f0eda582e1 | ||
|
|
219df48228 | ||
|
|
693f3a21de | ||
|
|
a0979f4e39 | ||
|
|
d47548ce24 | ||
|
|
3754c8d08d | ||
|
|
8b4ea2a089 | ||
|
|
076e70de6b | ||
|
|
a1fbbd0922 | ||
|
|
b48c2abfd3 | ||
|
|
67b465fac4 | ||
|
|
147fe4dc78 | ||
|
|
810fe2b4b1 | ||
|
|
85a53d031e | ||
|
|
74f3d3bce6 | ||
|
|
848a36f9cd | ||
|
|
3b3effc67d | ||
|
|
e129152f79 | ||
|
|
a845c49bcd | ||
|
|
1e6c131f81 | ||
|
|
b0082fe61b | ||
|
|
7ee0cfadc7 | ||
|
|
6bd16ded60 | ||
|
|
32c8aefb76 | ||
|
|
f647e97b7f | ||
|
|
10900725a1 | ||
|
|
080de01b37 | ||
|
|
5fd253a7f5 | ||
|
|
3b4e735cf2 | ||
|
|
e83ce64a47 | ||
|
|
f4d8b0e529 | ||
|
|
c25dadfeb2 | ||
|
|
5182c22f1a | ||
|
|
a32cc13bc3 | ||
|
|
9cabb1adb2 | ||
|
|
cf703caad8 |
+117
@@ -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
|
||||||
@@ -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"]
|
||||||
@@ -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"]
|
|
||||||
@@ -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
|
||||||
|
|
||||||
|
%:
|
||||||
|
@:
|
||||||
@@ -1,5 +1,7 @@
|
|||||||
# Aqara-MQTT
|
# Aqara-MQTT
|
||||||
Aqara (Xiaomi) Gateway to MQTT brodge.
|
[](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"
|
||||||
|
|||||||
@@ -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
@@ -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"
|
||||||
|
|||||||
+17
-15
@@ -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,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:
|
||||||
@@ -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,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:
|
||||||
@@ -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()
|
||||||
|
|||||||
+82
-43
@@ -1,17 +1,24 @@
|
|||||||
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
|
||||||
@@ -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
|
||||||
@@ -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)
|
||||||
|
|||||||
+26
-20
@@ -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))
|
||||||
@@ -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)
|
||||||
@@ -232,7 +237,7 @@ 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]:
|
||||||
@@ -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
@@ -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
|
||||||
|
|||||||
Reference in New Issue
Block a user