updating configuration
This commit is contained in:
@@ -0,0 +1,9 @@
|
||||
FROM python:3.10
|
||||
WORKDIR /app
|
||||
|
||||
ADD ./src/requirements.txt /app/requirements.txt
|
||||
RUN pip3 install -r requirements.txt
|
||||
|
||||
ADD ./src /app
|
||||
|
||||
CMD [ "python3", "-u", "hc2mqtt.py", "config.json" ]
|
||||
@@ -0,0 +1,234 @@
|
||||
#!/usr/bin/env python3
|
||||
# Parse messages from a Home Connect websocket (HCSocket)
|
||||
# and keep the connection alive
|
||||
#
|
||||
# Possible resources to fetch from the devices:
|
||||
#
|
||||
# /ro/values
|
||||
# /ro/descriptionChange
|
||||
# /ro/allMandatoryValues
|
||||
# /ro/allDescriptionChanges
|
||||
# /ro/activeProgram
|
||||
# /ro/selectedProgram
|
||||
#
|
||||
# /ei/initialValues
|
||||
# /ei/deviceReady
|
||||
#
|
||||
# /ci/services
|
||||
# /ci/registeredDevices
|
||||
# /ci/pairableDevices
|
||||
# /ci/delregistration
|
||||
# /ci/networkDetails
|
||||
# /ci/networkDetails2
|
||||
# /ci/wifiNetworks
|
||||
# /ci/wifiSetting
|
||||
# /ci/wifiSetting2
|
||||
# /ci/tzInfo
|
||||
# /ci/authentication
|
||||
# /ci/register
|
||||
# /ci/deregister
|
||||
#
|
||||
# /ce/serverDeviceType
|
||||
# /ce/serverCredential
|
||||
# /ce/clientCredential
|
||||
# /ce/hubInformation
|
||||
# /ce/hubConnected
|
||||
# /ce/status
|
||||
#
|
||||
# /ni/config
|
||||
#
|
||||
# /iz/services
|
||||
|
||||
import sys
|
||||
import json
|
||||
import re
|
||||
import time
|
||||
import io
|
||||
import traceback
|
||||
from datetime import datetime
|
||||
from base64 import urlsafe_b64encode as base64url_encode
|
||||
from Crypto.Random import get_random_bytes
|
||||
|
||||
|
||||
def now():
|
||||
return datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")
|
||||
|
||||
class HCDevice:
|
||||
def __init__(self, ws, features):
|
||||
self.ws = ws
|
||||
self.features = features
|
||||
self.session_id = None
|
||||
self.tx_msg_id = None
|
||||
self.device_name = "hcpy"
|
||||
self.device_id = "0badcafe"
|
||||
self.debug = False
|
||||
|
||||
def parse_values(self, values):
|
||||
if not self.features:
|
||||
return values
|
||||
|
||||
result = {}
|
||||
|
||||
for msg in values:
|
||||
uid = str(msg["uid"])
|
||||
value = msg["value"]
|
||||
value_str = str(value)
|
||||
|
||||
name = uid
|
||||
status = None
|
||||
|
||||
if uid in self.features:
|
||||
status = self.features[uid]
|
||||
|
||||
if status:
|
||||
name = status["name"]
|
||||
if "values" in status \
|
||||
and value_str in status["values"]:
|
||||
value = status["values"][value_str]
|
||||
|
||||
# trim everything off the name except the last part
|
||||
name = re.sub(r'^.*\.', '', name)
|
||||
result[name] = value
|
||||
|
||||
return result
|
||||
|
||||
def recv(self):
|
||||
try:
|
||||
buf = self.ws.recv()
|
||||
if buf is None:
|
||||
return None
|
||||
except Exception as e:
|
||||
print("receive error", e, traceback.format_exc())
|
||||
return None
|
||||
|
||||
try:
|
||||
return self.handle_message(buf)
|
||||
except Exception as e:
|
||||
print("error handling msg", e, buf, traceback.format_exc())
|
||||
return None
|
||||
|
||||
# reply to a POST or GET message with new data
|
||||
def reply(self, msg, reply):
|
||||
self.ws.send({
|
||||
'sID': msg["sID"],
|
||||
'msgID': msg["msgID"], # same one they sent to us
|
||||
'resource': msg["resource"],
|
||||
'version': msg["version"],
|
||||
'action': 'RESPONSE',
|
||||
'data': [reply],
|
||||
})
|
||||
|
||||
# send a message to the device
|
||||
def get(self, resource, version=1, action="GET", data=None):
|
||||
msg = {
|
||||
"sID": self.session_id,
|
||||
"msgID": self.tx_msg_id,
|
||||
"resource": resource,
|
||||
"version": version,
|
||||
"action": action,
|
||||
}
|
||||
|
||||
if data is not None:
|
||||
msg["data"] = [data]
|
||||
|
||||
self.ws.send(msg)
|
||||
self.tx_msg_id += 1
|
||||
|
||||
def handle_message(self, buf):
|
||||
msg = json.loads(buf)
|
||||
if self.debug:
|
||||
print(now(), "RX:", msg)
|
||||
sys.stdout.flush()
|
||||
|
||||
|
||||
resource = msg["resource"]
|
||||
action = msg["action"]
|
||||
|
||||
values = {}
|
||||
|
||||
if "code" in msg:
|
||||
#print(now(), "ERROR", msg["code"])
|
||||
values = {
|
||||
"error": msg["code"],
|
||||
"resource": msg.get("resource", ''),
|
||||
}
|
||||
elif action == "POST":
|
||||
if resource == "/ei/initialValues":
|
||||
# this is the first message they send to us and
|
||||
# establishes our session plus message ids
|
||||
self.session_id = msg["sID"]
|
||||
self.tx_msg_id = msg["data"][0]["edMsgID"]
|
||||
|
||||
self.reply(msg, {
|
||||
"deviceType": "Application",
|
||||
"deviceName": self.device_name,
|
||||
"deviceID": self.device_id,
|
||||
})
|
||||
|
||||
# ask the device which services it supports
|
||||
self.get("/ci/services")
|
||||
|
||||
# the clothes washer wants this, the token doesn't matter,
|
||||
# although they do not handle padding characters
|
||||
# they send a response, not sure how to interpet it
|
||||
token = base64url_encode(get_random_bytes(32)).decode('UTF-8')
|
||||
token = re.sub(r'=', '', token)
|
||||
self.get("/ci/authentication", version=2, data={"nonce": token})
|
||||
|
||||
self.get("/ci/info", version=2) # clothes washer
|
||||
self.get("/iz/info") # dish washer
|
||||
#self.get("/ci/tzInfo", version=2)
|
||||
self.get("/ni/info")
|
||||
#self.get("/ni/config", data={"interfaceID": 0})
|
||||
self.get("/ei/deviceReady", version=2, action="NOTIFY")
|
||||
self.get("/ro/allDescriptionChanges")
|
||||
self.get("/ro/allDescriptionChanges")
|
||||
self.get("/ro/allMandatoryValues")
|
||||
#self.get("/ro/values")
|
||||
else:
|
||||
print(now(), "Unknown resource", resource, file=sys.stderr)
|
||||
|
||||
elif action == "RESPONSE" or action == "NOTIFY":
|
||||
if resource == "/iz/info" or resource == "/ci/info":
|
||||
# we could validate that this matches our machine
|
||||
pass
|
||||
|
||||
elif resource == "/ro/descriptionChange" \
|
||||
or resource == "/ro/allDescriptionChanges":
|
||||
# we asked for these but don't know have to parse yet
|
||||
pass
|
||||
|
||||
elif resource == "/ni/info":
|
||||
# we're already talking, so maybe we don't care?
|
||||
pass
|
||||
|
||||
elif resource == "/ro/allMandatoryValues" \
|
||||
or resource == "/ro/values":
|
||||
values = self.parse_values(msg["data"])
|
||||
elif resource == "/ci/registeredDevices":
|
||||
# we don't care
|
||||
pass
|
||||
|
||||
elif resource == "/ci/services":
|
||||
self.services = {}
|
||||
for service in msg["data"]:
|
||||
self.services[service["service"]] = {
|
||||
"version": service["version"],
|
||||
}
|
||||
#print(now(), "services", self.services)
|
||||
|
||||
# we should figure out which ones to query now
|
||||
# if "iz" in self.services:
|
||||
# self.get("/iz/info", version=self.services["iz"]["version"])
|
||||
# if "ni" in self.services:
|
||||
# self.get("/ni/info", version=self.services["ni"]["version"])
|
||||
# if "ei" in self.services:
|
||||
# self.get("/ei/deviceReady", version=self.services["ei"]["version"], action="NOTIFY")
|
||||
|
||||
#self.get("/if/info")
|
||||
|
||||
else:
|
||||
print(now(), "Unknown", msg)
|
||||
|
||||
# return whatever we've parsed out of it
|
||||
return values
|
||||
@@ -0,0 +1,170 @@
|
||||
# Create a websocket that wraps a connection to a
|
||||
# Bosh-Siemens Home Connect device
|
||||
import socket
|
||||
import ssl
|
||||
import sslpsk
|
||||
import websocket
|
||||
import sys
|
||||
import json
|
||||
import re
|
||||
import time
|
||||
import io
|
||||
from base64 import urlsafe_b64decode as base64url
|
||||
from datetime import datetime
|
||||
from Crypto.Cipher import AES
|
||||
from Crypto.Hash import HMAC, SHA256
|
||||
from Crypto.Random import get_random_bytes
|
||||
|
||||
# Convience to compute an HMAC on a message
|
||||
def hmac(key,msg):
|
||||
mac = HMAC.new(key, msg=msg, digestmod=SHA256).digest()
|
||||
return mac
|
||||
|
||||
def now():
|
||||
return datetime.now().strftime("%Y-%m-%d %H:%M:%S.%f")
|
||||
|
||||
# Monkey patch for sslpsk in pip using the old _sslobj
|
||||
def _sslobj(sock):
|
||||
if (3, 5) <= sys.version_info <= (3, 7):
|
||||
return sock._sslobj._sslobj
|
||||
else:
|
||||
return sock._sslobj
|
||||
sslpsk.sslpsk._sslobj = _sslobj
|
||||
|
||||
|
||||
class HCSocket:
|
||||
def __init__(self, host, psk64, iv64=None):
|
||||
self.host = host
|
||||
self.psk = base64url(psk64 + '===')
|
||||
self.debug = False
|
||||
|
||||
if iv64:
|
||||
# an HTTP self-encrypted socket
|
||||
self.http = True
|
||||
self.iv = base64url(iv64 + '===')
|
||||
self.enckey = hmac(self.psk, b'ENC')
|
||||
self.mackey = hmac(self.psk, b'MAC')
|
||||
self.port = 80
|
||||
self.uri = "ws://" + host + ":80/homeconnect"
|
||||
else:
|
||||
self.http = False
|
||||
self.port = 443
|
||||
self.uri = "wss://" + host + ":443/homeconnect"
|
||||
|
||||
# don't connect automatically so that debug etc can be set
|
||||
#self.reconnect()
|
||||
|
||||
# restore the encryption state for a fresh connection
|
||||
# this is only used by the HTTP connection
|
||||
def reset(self):
|
||||
if not self.http:
|
||||
return
|
||||
self.last_rx_hmac = bytes(16)
|
||||
self.last_tx_hmac = bytes(16)
|
||||
|
||||
self.aes_encrypt = AES.new(self.enckey, AES.MODE_CBC, self.iv)
|
||||
self.aes_decrypt = AES.new(self.enckey, AES.MODE_CBC, self.iv)
|
||||
|
||||
# hmac an inbound or outbound message, chaining the last hmac too
|
||||
def hmac_msg(self, direction, enc_msg):
|
||||
hmac_msg = self.iv + direction + enc_msg
|
||||
return hmac(self.mackey, hmac_msg)[0:16]
|
||||
|
||||
def decrypt(self,buf):
|
||||
if len(buf) < 32:
|
||||
print("Short message?", buf.hex(), file=sys.stderr)
|
||||
return None
|
||||
if len(buf) % 16 != 0:
|
||||
print("Unaligned message? probably bad", buf.hex(), file=sys.stderr)
|
||||
|
||||
# split the message into the encrypted message and the first 16-bytes of the HMAC
|
||||
enc_msg = buf[0:-16]
|
||||
their_hmac = buf[-16:]
|
||||
|
||||
# compute the expected hmac on the encrypted message
|
||||
our_hmac = self.hmac_msg(b'\x43' + self.last_rx_hmac, enc_msg)
|
||||
|
||||
if their_hmac != our_hmac:
|
||||
print("HMAC failure", their_hmac.hex(), our_hmac.hex(), file=sys.stderr)
|
||||
return None
|
||||
|
||||
self.last_rx_hmac = their_hmac
|
||||
|
||||
# decrypt the message with CBC, so the last message block is mixed in
|
||||
msg = self.aes_decrypt.decrypt(enc_msg)
|
||||
|
||||
# check for padding and trim it off the end
|
||||
pad_len = msg[-1]
|
||||
if len(msg) < pad_len:
|
||||
print("padding error?", msg.hex())
|
||||
return None
|
||||
|
||||
return msg[0:-pad_len]
|
||||
|
||||
def encrypt(self, clear_msg):
|
||||
# convert the UTF-8 string into a byte array
|
||||
clear_msg = bytes(clear_msg, 'utf-8')
|
||||
|
||||
# pad the buffer, adding an extra block if necessary
|
||||
pad_len = 16 - (len(clear_msg) % 16)
|
||||
if pad_len == 1:
|
||||
pad_len += 16
|
||||
pad = b'\x00' + get_random_bytes(pad_len-2) + bytearray([pad_len])
|
||||
|
||||
clear_msg = clear_msg + pad
|
||||
|
||||
# encrypt the padded message with CBC, so there is chained
|
||||
# state from the last cipher block sent
|
||||
enc_msg = self.aes_encrypt.encrypt(clear_msg)
|
||||
|
||||
# compute the hmac of the encrypted message, chaining the
|
||||
# hmac of the previous message plus direction 'E'
|
||||
self.last_tx_hmac = self.hmac_msg(b'\x45' + self.last_tx_hmac, enc_msg)
|
||||
|
||||
# append the new hmac to the message
|
||||
return enc_msg + self.last_tx_hmac
|
||||
|
||||
def reconnect(self):
|
||||
self.reset()
|
||||
sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
sock.connect((self.host,self.port))
|
||||
|
||||
if not self.http:
|
||||
sock = sslpsk.wrap_socket(
|
||||
sock,
|
||||
ssl_version = ssl.PROTOCOL_TLSv1_2,
|
||||
ciphers = 'ECDHE-PSK-CHACHA20-POLY1305',
|
||||
psk = self.psk,
|
||||
)
|
||||
|
||||
print(now(), "CON:", self.uri)
|
||||
self.ws = websocket.WebSocket()
|
||||
self.ws.connect(self.uri,
|
||||
socket = sock,
|
||||
origin = "",
|
||||
)
|
||||
|
||||
def send(self, msg):
|
||||
buf = json.dumps(msg, separators=(',', ':') )
|
||||
# swap " for '
|
||||
buf = re.sub("'", '"', buf)
|
||||
if self.debug:
|
||||
print(now(), "TX:", buf)
|
||||
if self.http:
|
||||
self.ws.send_binary(self.encrypt(buf))
|
||||
else:
|
||||
self.ws.send(buf)
|
||||
|
||||
def recv(self):
|
||||
buf = self.ws.recv()
|
||||
if buf is None or buf == "":
|
||||
return None
|
||||
|
||||
if self.http:
|
||||
buf = self.decrypt(buf)
|
||||
if buf is None:
|
||||
return None
|
||||
|
||||
if self.debug:
|
||||
print(now(), "RX:", buf)
|
||||
return buf
|
||||
@@ -0,0 +1,127 @@
|
||||
#!/usr/bin/env python3
|
||||
# Contact Bosh-Siemens Home Connect devices
|
||||
# and connect their messages to the mqtt server
|
||||
import json
|
||||
import sys
|
||||
import time
|
||||
import os
|
||||
import datetime
|
||||
from threading import Thread
|
||||
|
||||
import click
|
||||
import paho.mqtt.client as mqtt
|
||||
|
||||
from HCDevice import HCDevice
|
||||
from HCSocket import HCSocket, now
|
||||
|
||||
def restart_script():
|
||||
"""Перезапускает скрипт"""
|
||||
print(f"{now()} Перезапуск скрипта...")
|
||||
python = sys.executable
|
||||
os.execl(python, python, *sys.argv)
|
||||
|
||||
@click.command()
|
||||
@click.argument("config_file")
|
||||
@click.option("-h", "--mqtt_host", default="localhost")
|
||||
@click.option("-p", "--mqtt_prefix", default="homeconnect/")
|
||||
@click.option("-u", "--mqtt_user", default="")
|
||||
@click.option("-pw", "--mqtt_password", default="")
|
||||
def hc2mqtt(config_file: str, mqtt_host: str, mqtt_prefix: str, mqtt_user: str, mqtt_password: str):
|
||||
click.echo(f"Hello {config_file=} {mqtt_host=} {mqtt_prefix=} {mqtt_user=}")
|
||||
|
||||
with open(config_file, "r") as f:
|
||||
devices = json.load(f)
|
||||
|
||||
client = mqtt.Client()
|
||||
if mqtt_user != "" and mqtt_password != "":
|
||||
client.username_pw_set(mqtt_user, mqtt_password)
|
||||
client.connect(host=mqtt_host, port=1883, keepalive=70)
|
||||
|
||||
for device in devices:
|
||||
mqtt_topic = mqtt_prefix + device["name"]
|
||||
thread = Thread(target=client_connect, args=(client, device, mqtt_topic))
|
||||
thread.start()
|
||||
|
||||
last_restart = datetime.datetime.now()
|
||||
|
||||
while True:
|
||||
current_time = datetime.datetime.now()
|
||||
# Проверяем, прошло ли 24 часа с последнего перезапуска
|
||||
if (current_time - last_restart).total_seconds() >= 86400: # 86400 секунд = 24 часа
|
||||
restart_script()
|
||||
time.sleep(60) # Проверяем каждую минуту
|
||||
client.loop(timeout=1.0)
|
||||
|
||||
|
||||
# Map their value names to easier state names
|
||||
topics = {
|
||||
"OperationState": "state",
|
||||
"DoorState": "door",
|
||||
"RemainingProgramTime": "remaining",
|
||||
"PowerState": "power",
|
||||
"LowWaterPressure": "lowwaterpressure",
|
||||
"AquaStopOccured": "aquastop",
|
||||
"InternalError": "error",
|
||||
"FatalErrorOccured": "error",
|
||||
"MachineCareReminder": "machinecarereminder",
|
||||
"ActiveProgram": "activeprogram" #8215 - machine care
|
||||
}
|
||||
|
||||
|
||||
|
||||
def client_connect(client, device, mqtt_topic):
|
||||
host = device["host"]
|
||||
|
||||
state = {}
|
||||
for topic in topics:
|
||||
state[topics[topic]] = None
|
||||
|
||||
while True:
|
||||
try:
|
||||
ws = HCSocket(host, device["key"], device.get("iv",None))
|
||||
dev = HCDevice(ws, device.get("features", None))
|
||||
|
||||
ws.debug = True
|
||||
ws.reconnect()
|
||||
|
||||
while True:
|
||||
msg = dev.recv()
|
||||
if msg is None:
|
||||
break
|
||||
if len(msg) > 0:
|
||||
print(now(), msg)
|
||||
|
||||
update = False
|
||||
for topic in topics:
|
||||
value = msg.get(topic, None)
|
||||
if value is None:
|
||||
continue
|
||||
|
||||
# Convert "On" to True, "Off" to False
|
||||
if value == "On":
|
||||
value = True
|
||||
elif value == "Off":
|
||||
value = False
|
||||
|
||||
new_topic = topics[topic]
|
||||
if new_topic == "remaining":
|
||||
state["remainingseconds"] = value
|
||||
value = "%d:%02d" % (value / 60 / 60, (value / 60) % 60)
|
||||
|
||||
state[new_topic] = value
|
||||
update = True
|
||||
|
||||
if not update:
|
||||
continue
|
||||
|
||||
msg = json.dumps(state)
|
||||
print("publish", mqtt_topic, msg)
|
||||
client.publish(mqtt_topic + "/state", msg, retain=True)
|
||||
|
||||
except Exception as e:
|
||||
print("ERROR", host, e, file=sys.stderr)
|
||||
|
||||
time.sleep(5)
|
||||
|
||||
if __name__ == "__main__":
|
||||
hc2mqtt(auto_envvar_prefix='HC')
|
||||
@@ -0,0 +1,9 @@
|
||||
bs4
|
||||
requests
|
||||
pycryptodome
|
||||
websocket-client
|
||||
sslpsk
|
||||
paho.mqtt
|
||||
lxml
|
||||
click
|
||||
requests
|
||||
@@ -0,0 +1,3 @@
|
||||
HC_MQTT_HOST=xxxxxxxxx
|
||||
HC_MQTT_USER=xxxxxxxxx
|
||||
HC_MQTT_PASSWORD=xxxxxxxxx
|
||||
Reference in New Issue
Block a user