121 lines
5.2 KiB
Python
121 lines
5.2 KiB
Python
import logging
|
|
import json
|
|
from paho.mqtt import client as mqtt
|
|
from .mqtt_discovery import MQTTDiscovery
|
|
import time
|
|
|
|
class MQTTClient:
|
|
def __init__(self, controller, config):
|
|
self.controller = controller
|
|
self.config = config
|
|
self.logger = logging.getLogger("MQTTClient")
|
|
|
|
mqtt_conf = self.config.get("mqtt", {})
|
|
self.broker = mqtt_conf.get("broker", "localhost")
|
|
self.port = mqtt_conf.get("port", 1883)
|
|
|
|
username = mqtt_conf.get("user")
|
|
password = mqtt_conf.get("password")
|
|
|
|
self.base_topic = mqtt_conf.get("base_topic", "homeassistant")
|
|
self.sensor_topic = mqtt_conf.get("sensor_topic", "pool_temp")
|
|
|
|
# Wichtig: sensor_topic hier übergeben!
|
|
self.discovery = MQTTDiscovery(base_topic=self.base_topic, sensor_topic=self.sensor_topic)
|
|
self.client = mqtt.Client()
|
|
|
|
if username and password:
|
|
self.client.username_pw_set(username, password)
|
|
|
|
self.client.on_connect = self._on_connect
|
|
self.client.on_message = self._on_message
|
|
|
|
self.sync_topics = set()
|
|
self.is_connected = False
|
|
|
|
def start(self):
|
|
try:
|
|
self.logger.info(f"Verbinde mit MQTT-Broker {self.broker}:{self.port}...")
|
|
testament_topic = f"{self.sensor_topic}/28833e1363721b06/availability"
|
|
self.client.will_set(testament_topic, payload="offline", qos=1, retain=True)
|
|
self.client.connect(self.broker, self.port, 60)
|
|
self.client.loop_start()
|
|
except Exception as e:
|
|
self.logger.error(f"MQTT-Verbindungsfehler: {e}")
|
|
|
|
def _on_connect(self, client, userdata, flags, rc):
|
|
if rc == 0:
|
|
self.logger.info("Erfolgreich mit MQTT-Broker verbunden.")
|
|
self.is_connected = True
|
|
|
|
# Korrigiert: Nutze das dynamische sensor_topic für den Paketverlust-Sync!
|
|
sync_topic = f"{self.sensor_topic}/+/packet_loss/state"
|
|
self.client.subscribe(sync_topic)
|
|
self.logger.debug(f"Sync-Topic abonniert: {sync_topic}")
|
|
|
|
# Abonnieren der Slider-Befehle aus HA (z.B. pool_temp/+/+/set)
|
|
cmd_topic = f"{self.sensor_topic}/+/+/set"
|
|
self.client.subscribe(cmd_topic)
|
|
self.logger.debug(f"Command-Topic abonniert: {cmd_topic}")
|
|
else:
|
|
self.logger.error(f"Verbindung abgelehnt mit Code: {rc}")
|
|
|
|
def _on_message(self, client, userdata, msg):
|
|
try:
|
|
topic_parts = msg.topic.split('/')
|
|
# Aufbau: [sensor_topic]/[uuid]/[sub_topic]/[state_or_set]
|
|
if len(topic_parts) != 4:
|
|
return
|
|
|
|
uuid = topic_parts[1]
|
|
sub_topic = topic_parts[2]
|
|
action = topic_parts[3]
|
|
|
|
# 1. Sync-Logik beim Gateway-Start
|
|
if sub_topic == "packet_loss" and action == "state":
|
|
if uuid not in self.sync_topics:
|
|
payload_str = msg.payload.decode().strip()
|
|
retained_loss = int(payload_str) if payload_str else 0
|
|
|
|
self.logger.info(f"Retain-Wert für {uuid} gefunden: {retained_loss} verlorene Pakete.")
|
|
self.sync_topics.add(uuid)
|
|
self.controller.handle_mqtt_sync(uuid, retained_loss)
|
|
self.publish_discovery(uuid)
|
|
|
|
# 2. Slider-Änderungen aus Home Assistant abfangen
|
|
elif action == "set":
|
|
value_str = msg.payload.decode().strip()
|
|
self.logger.info(f"Slider-Änderung von HA erhalten: {sub_topic} -> {value_str} für {uuid}")
|
|
|
|
# Bestätigung direkt an den State-Kanal zurücksenden, damit der Slider nicht zurückspringt
|
|
self.client.publish(f"{self.sensor_topic}/{uuid}/{sub_topic}/state", value_str, retain=True)
|
|
|
|
self.controller.handle_param_change_from_ha(uuid, sub_topic, float(value_str))
|
|
|
|
except Exception as e:
|
|
self.logger.error(f"Fehler beim Verarbeiten der MQTT-Nachricht: {e}")
|
|
|
|
def trigger_fallback_sync(self, uuid):
|
|
if uuid not in self.sync_topics:
|
|
self.logger.info(f"Kein Retain-Wert für {uuid} empfangen. Starte initial bei 0.")
|
|
self.sync_topics.add(uuid)
|
|
self.controller.handle_mqtt_sync(uuid, 0)
|
|
self.publish_discovery(uuid)
|
|
|
|
def publish_discovery(self, uuid):
|
|
configs = self.discovery.get_configs(uuid)
|
|
for topic, payload in configs.items():
|
|
self.client.publish(topic, json.dumps(payload), retain=True)
|
|
self.logger.debug(f"Home Assistant Discovery gesendet: {topic}, Payload: {payload}")
|
|
self.logger.info(f"Home Assistant Discovery für Node {uuid} gesendet.")
|
|
|
|
def publish_sensor_value(self, uuid, sensor_type, value, retain=False):
|
|
if not self.is_connected:
|
|
return
|
|
topic = f"{self.sensor_topic}/{uuid}/{sensor_type}/state"
|
|
self.client.publish(topic, str(value), retain=retain)
|
|
self.logger.debug(f"Sensorwert veröffentlicht: {topic} -> {value} (retain={retain})")
|
|
|
|
def stop(self):
|
|
self.client.loop_stop()
|
|
self.client.disconnect() |