MQTT: - rc=5 (Not Authorized) stoppt Reconnect-Loop via _auth_failed Flag - Fehlermeldung im MQTT-Einstellungen-Banner sichtbar Sicherheit: - /api/* nur über HAOS-Ingress (X-Ingress-Path) oder Loopback erreichbar Flash-Wizard (Baustelle B): - Neuer Tab "Flash" mit IP-Eingabe und OTA-Modus-Erkennung - OTA: integrierte oder eigene Firmware via POST /api/flash/update auf Stick - Fortschrittsbalken + Polling bis Stick nach Reset wieder online - ST-Link-Erstflash-Anleitung (Pinout, st-flash Kommando) - Firmware-Binaries im Docker-Image unter /firmware/ NuttX OTA (Baustelle A, shinelanx-modbus): - ota_http.c: Zwei-Phasen OTA für STM32F103 Single-Bank Flash Stage 1: Firmware in Staging-Bereich (obere Flashhälfte) schreiben Stage 2: .ramfuncs aus SRAM heraus — Staging → App-Bereich kopieren, Reset - ota_http.h, Makefile und main.c entsprechend erweitert - ld.script.dfu: .ramfuncs in .data Section → Ausführung aus SRAM Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
189 lines
7.4 KiB
Python
189 lines
7.4 KiB
Python
import json
|
|
import logging
|
|
from typing import Dict, List, Optional, Tuple
|
|
|
|
import paho.mqtt.client as mqtt
|
|
|
|
from inverters import Inverter
|
|
|
|
log = logging.getLogger(__name__)
|
|
|
|
AGG_DEVICE_ID = "shinebridge_aggregate"
|
|
AGG_TOPIC = "shinebridge/aggregate"
|
|
|
|
|
|
_RC_MESSAGES = {
|
|
1: "Falsche Protokollversion",
|
|
2: "Client-ID abgelehnt",
|
|
3: "Broker nicht verfügbar",
|
|
4: "Falscher Benutzername oder Passwort",
|
|
5: "Nicht autorisiert — Credentials prüfen",
|
|
}
|
|
|
|
|
|
class MqttPublisher:
|
|
def __init__(self, broker: str, port: int, user: str, password: str,
|
|
agg_meta: Optional[Dict] = None):
|
|
self._broker = broker
|
|
self._port = port
|
|
self._connected = False
|
|
self._last_error: Optional[str] = None
|
|
self._auth_failed = False
|
|
self._registered: List[Tuple] = []
|
|
self._agg_meta: Dict = agg_meta or {}
|
|
|
|
self._subscriptions: List[Tuple[str, any]] = []
|
|
self._client = mqtt.Client(client_id="shinebridge_hub", clean_session=True)
|
|
if user:
|
|
self._client.username_pw_set(user, password)
|
|
self._client.on_connect = self._on_connect
|
|
self._client.on_disconnect = self._on_disconnect
|
|
self._client.on_message = self._on_message
|
|
|
|
def _on_connect(self, client, userdata, flags, rc):
|
|
if rc == 0:
|
|
self._connected = True
|
|
self._last_error = None
|
|
self._auth_failed = False
|
|
log.info("MQTT verbunden: %s:%d", self._broker, self._port)
|
|
for entry in self._registered:
|
|
self._publish_discovery(*entry)
|
|
if self._agg_meta:
|
|
self._publish_aggregate_discovery()
|
|
for topic, _ in self._subscriptions:
|
|
client.subscribe(topic)
|
|
else:
|
|
msg = _RC_MESSAGES.get(rc, f"Unbekannter Fehler")
|
|
self._last_error = f"rc={rc}: {msg}"
|
|
log.error("MQTT Verbindungsfehler %s", self._last_error)
|
|
if rc == 5:
|
|
self._auth_failed = True
|
|
client.loop_stop()
|
|
|
|
def _on_disconnect(self, client, userdata, rc):
|
|
self._connected = False
|
|
if rc != 0:
|
|
log.warning("MQTT getrennt rc=%d", rc)
|
|
|
|
def _on_message(self, client, userdata, msg):
|
|
for topic, callback in self._subscriptions:
|
|
if mqtt.topic_matches_sub(topic, msg.topic):
|
|
try:
|
|
callback(msg.topic, msg.payload)
|
|
except Exception as e:
|
|
log.error("MQTT message handler Fehler [%s]: %s", msg.topic, e)
|
|
|
|
def subscribe(self, topic: str, callback):
|
|
self._subscriptions.append((topic, callback))
|
|
if self._connected:
|
|
self._client.subscribe(topic)
|
|
|
|
def connect(self):
|
|
if self._auth_failed:
|
|
log.warning("MQTT connect übersprungen — Authentifizierung fehlgeschlagen")
|
|
return
|
|
try:
|
|
self._client.connect_async(self._broker, self._port, keepalive=60)
|
|
self._client.loop_start()
|
|
except Exception as e:
|
|
log.error("MQTT connect fehlgeschlagen: %s", e)
|
|
|
|
def disconnect(self):
|
|
self._client.loop_stop()
|
|
self._client.disconnect()
|
|
|
|
@property
|
|
def connected(self) -> bool:
|
|
return self._connected
|
|
|
|
@property
|
|
def last_error(self) -> Optional[str]:
|
|
return self._last_error
|
|
|
|
# ── Gerät-Discovery ──────────────────────────────────────
|
|
|
|
def register_inverter(self, inverter: Inverter, device_id: str,
|
|
topic_prefix: str, display_name: str = None):
|
|
entry = (inverter, device_id, topic_prefix, display_name)
|
|
self._registered = [r for r in self._registered if r[1] != device_id]
|
|
self._registered.append(entry)
|
|
if self._connected:
|
|
self._publish_discovery(inverter, device_id, topic_prefix, display_name)
|
|
|
|
def unregister_inverter(self, device_id: str):
|
|
self._registered = [r for r in self._registered if r[1] != device_id]
|
|
|
|
def _publish_discovery(self, inverter: Inverter, device_id: str,
|
|
topic_prefix: str, display_name: str = None):
|
|
device_payload = {
|
|
"identifiers": [device_id],
|
|
"name": display_name or inverter.name,
|
|
"manufacturer": inverter.manufacturer,
|
|
"model": inverter.name,
|
|
}
|
|
for sensor in inverter.sensors:
|
|
config = {
|
|
"name": sensor.name,
|
|
"unique_id": f"{device_id}_{sensor.id}",
|
|
"state_topic": f"{topic_prefix}/state",
|
|
"value_template": f"{{{{ value_json.{sensor.id} }}}}",
|
|
"unit_of_measurement": sensor.unit,
|
|
"state_class": sensor.state_class,
|
|
"icon": sensor.icon,
|
|
"device": device_payload,
|
|
}
|
|
if sensor.device_class:
|
|
config["device_class"] = sensor.device_class
|
|
topic = f"homeassistant/sensor/{device_id}/{sensor.id}/config"
|
|
self._client.publish(topic, json.dumps(config), retain=True, qos=1)
|
|
log.info("MQTT Discovery: %d Sensoren für %s", len(inverter.sensors), device_id)
|
|
|
|
# ── Aggregat-Discovery ────────────────────────────────────
|
|
|
|
def _publish_aggregate_discovery(self):
|
|
device_payload = {
|
|
"identifiers": [AGG_DEVICE_ID],
|
|
"name": "ShineBridge Gesamt",
|
|
"manufacturer": "ShineBridge",
|
|
"model": "Aggregat",
|
|
}
|
|
for sensor_id, meta in self._agg_meta.items():
|
|
config = {
|
|
"name": meta["name"],
|
|
"unique_id": f"{AGG_DEVICE_ID}_{sensor_id}",
|
|
"state_topic": f"{AGG_TOPIC}/state",
|
|
"value_template": f"{{{{ value_json.{sensor_id} }}}}",
|
|
"unit_of_measurement": meta["unit"],
|
|
"state_class": meta["state_class"],
|
|
"icon": meta["icon"],
|
|
"device": device_payload,
|
|
}
|
|
if meta.get("device_class"):
|
|
config["device_class"] = meta["device_class"]
|
|
topic = f"homeassistant/sensor/{AGG_DEVICE_ID}/{sensor_id}/config"
|
|
self._client.publish(topic, json.dumps(config), retain=True, qos=1)
|
|
log.info("MQTT Discovery: %d Aggregat-Sensoren", len(self._agg_meta))
|
|
|
|
# ── Daten publizieren ─────────────────────────────────────
|
|
|
|
def publish_data(self, values: dict, topic_prefix: str):
|
|
if not self._connected:
|
|
return
|
|
self._client.publish(f"{topic_prefix}/state", json.dumps(values),
|
|
retain=True, qos=0)
|
|
|
|
def publish_status(self, status: str, topic_prefix: str):
|
|
self._client.publish(f"{topic_prefix}/status", status, retain=True, qos=1)
|
|
|
|
def publish_aggregates(self, values: dict):
|
|
if not self._connected or not values:
|
|
return
|
|
self._client.publish(f"{AGG_TOPIC}/state", json.dumps(values),
|
|
retain=True, qos=0)
|
|
self._client.publish(f"{AGG_TOPIC}/status", "online", retain=True, qos=1)
|
|
|
|
def publish_raw(self, topic: str, payload: str, retain: bool = False, qos: int = 1):
|
|
if not self._connected:
|
|
return
|
|
self._client.publish(topic, payload, retain=retain, qos=qos)
|