MQTT: Einstellungen empfangen, Zustand veroeffentlichen
Topics (Praefix aus mqtt_config.py, Vorgabe "grossanzeige"):
<PREFIX>/set/hell_prozent 60 eingehend, schreibt settings.json
<PREFIX>/set/dunkel_prozent 5
<PREFIX>/set/schwelle 1800
<PREFIX>/set/hysterese 150
<PREFIX>/status/ra 18h36m56s ausgehend, retained
<PREFIX>/status/dec +38°47'01"
<PREFIX>/status/ldr 2453
<PREFIX>/status/helligkeit 50
<PREFIX>/status/link 1
<PREFIX>/status/online 1 mit Last Will auf 0
Eine empfangene Einstellung wird sofort wirksam: dieselbe Pruefung wie sonst
(settings.update), dann laedt der laufende BrightnessController sie per reload()
nach -- kein Neustart noetig.
Leitgedanke des Moduls: **MQTT darf die Anzeige nie aufhalten.** Die Anzeige ist
der Zweck des Geraets, MQTT ist Beiwerk. Deshalb faengt mqtt.py jeden Fehler
selbst ab und meldet ihn nur:
- Broker nicht erreichbar -> Versuch scheitert, Schleife laeuft weiter.
Wiederholung mit wachsendem Abstand (5..120 s), sonst kostet ein dauerhaft
toter Broker in jedem Schleifendurchlauf Zeit.
- Verbindungsabriss -> beim naechsten Senden/Empfangen bemerkt, Verbindung wird
verworfen und spaeter neu aufgebaut. Danach geht der gesamte Status erneut
raus, damit der Broker nicht auf veralteten Werten sitzenbleibt.
- Unsinniger Wert von aussen -> verworfen, settings.json bleibt unberuehrt.
- Geraet faellt aus -> Last Will meldet online=0. Ohne das bliebe online=1
stehen, obwohl niemand mehr da ist.
Der Empfang blockiert nicht (check_msg). Statuswerte gehen nur bei Aenderung
raus -- die Koordinaten aendern sich staendig, LDR und Helligkeit kaum.
mqtt_config.py ist optional und gitignored (Vorlage mqtt_config_example.py);
fehlt sie, laeuft alles wie bisher ohne MQTT. deploy.sh weist nur darauf hin.
test_mqtt.py haengt einen Fake-Broker ein und prueft auch die Faelle, die man
mit einem echten Server schwer herbeifuehrt: Abriss beim Senden, Muell im Topic,
Broker der nicht antwortet. 143 Tests gruen.
Am Geraet bestaetigt ist bisher der Fall OHNE Broker: mqtt.py importiert,
connect_from_config liefert None, ein unerreichbarer Broker wirft keine
Ausnahme. Der Test gegen einen echten Broker steht noch aus.
Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,251 @@
|
||||
"""MQTT-Anbindung: Einstellungen empfangen, Zustand veroeffentlichen.
|
||||
|
||||
<PREFIX>/set/hell_prozent 60 eingehend, schreibt settings.json
|
||||
<PREFIX>/set/dunkel_prozent 5
|
||||
<PREFIX>/set/schwelle 1800
|
||||
<PREFIX>/set/hysterese 150
|
||||
|
||||
<PREFIX>/status/ra 18h36m56s ausgehend
|
||||
<PREFIX>/status/dec +38°47'01"
|
||||
<PREFIX>/status/ldr 2453
|
||||
<PREFIX>/status/helligkeit 50
|
||||
<PREFIX>/status/link 1 Montierung erreichbar
|
||||
<PREFIX>/status/online 1 retained, mit Last Will auf 0
|
||||
|
||||
**Grundsatz: MQTT darf die Anzeige nie aufhalten.** Die Anzeige ist der Zweck des
|
||||
Geraets, MQTT ist Beiwerk. Deshalb faengt jede Methode hier ihre Fehler selbst ab
|
||||
und meldet sie nur zurueck; ein toter Broker, ein abgezogenes Netzwerkkabel oder
|
||||
ein unsinniger Wert duerfen die Poll-Schleife nicht unterbrechen. Aus demselben
|
||||
Grund wird ein Verbindungsverlust nicht sofort und nicht endlos neu versucht,
|
||||
sondern mit wachsendem Abstand (RETRY_START..RETRY_MAX) -- sonst haengt die
|
||||
Schleife bei jedem Durchlauf im Verbindungsaufbau.
|
||||
|
||||
Der Empfang laeuft ueber check_msg(), das nicht blockiert: Liegt nichts an,
|
||||
kehrt es sofort zurueck.
|
||||
|
||||
Ohne mqtt_config.py (oder ohne umqtt) gibt connect_from_config() None zurueck --
|
||||
die Anzeige laeuft dann ohne MQTT weiter.
|
||||
"""
|
||||
|
||||
import settings
|
||||
from ticks import deadline as _deadline, expired as _expired
|
||||
|
||||
try: # MicroPython
|
||||
from umqtt.simple import MQTTClient
|
||||
except ImportError: # CPython: nur zum Testen, Client wird injiziert
|
||||
MQTTClient = None
|
||||
|
||||
|
||||
# Wartezeiten fuer den Wiederverbindungsversuch (Sekunden).
|
||||
RETRY_START = 5.0
|
||||
RETRY_MAX = 120.0
|
||||
|
||||
# Diese Einstellungen lassen sich ueber <PREFIX>/set/<name> aendern. Bewusst
|
||||
# aus settings.DEFAULTS abgeleitet: Was dort nicht steht, ist auch per MQTT
|
||||
# nicht setzbar, und settings prueft jeden Wert zusaetzlich auf seinen Bereich.
|
||||
SETZBAR = tuple(sorted(settings.DEFAULTS))
|
||||
|
||||
|
||||
class MqttBridge:
|
||||
"""Haelt die Broker-Verbindung und uebersetzt in beide Richtungen.
|
||||
|
||||
Wird der Client von aussen gegeben (Tests), unterbleibt der Aufbau eines
|
||||
echten MQTTClient. So laesst sich die ganze Logik ohne Broker pruefen.
|
||||
"""
|
||||
|
||||
def __init__(self, broker=None, port=1883, user="", password="",
|
||||
prefix="grossanzeige", client_id="grossanzeige",
|
||||
client=None, log=None, on_change=None):
|
||||
self._broker = broker
|
||||
self._port = port
|
||||
self._user = user
|
||||
self._password = password
|
||||
self._prefix = prefix
|
||||
self._client_id = client_id
|
||||
self._client = client
|
||||
self._log = log or (lambda *a: None)
|
||||
self._on_change = on_change # wird nach jeder Aenderung gerufen
|
||||
|
||||
self.connected = False
|
||||
self._retry = RETRY_START
|
||||
self._naechster_versuch = None
|
||||
self._zuletzt = {} # Topic -> zuletzt gesendeter Wert
|
||||
|
||||
# -- Topics ------------------------------------------------------------
|
||||
def _set_topic(self, name=None):
|
||||
return "%s/set/%s" % (self._prefix, "+" if name is None else name)
|
||||
|
||||
def _status_topic(self, name):
|
||||
return "%s/status/%s" % (self._prefix, name)
|
||||
|
||||
# -- Verbindung --------------------------------------------------------
|
||||
def _neuer_client(self):
|
||||
if MQTTClient is None:
|
||||
raise OSError("umqtt.simple nicht verfuegbar")
|
||||
return MQTTClient(self._client_id, self._broker, port=self._port,
|
||||
user=self._user or None,
|
||||
password=self._password or None,
|
||||
keepalive=60)
|
||||
|
||||
def connect(self):
|
||||
"""Verbinden, Set-Topics abonnieren, online melden.
|
||||
|
||||
Liefert True bei Erfolg. Schlaegt es fehl, wird das nur protokolliert --
|
||||
der naechste check() versucht es nach der Wartezeit erneut.
|
||||
"""
|
||||
try:
|
||||
if self._client is None:
|
||||
self._client = self._neuer_client()
|
||||
self._client.set_callback(self._on_message)
|
||||
# Last Will: Bricht die Verbindung weg, meldet der Broker selbst
|
||||
# offline. Ohne das bliebe online=1 stehen, obwohl niemand da ist.
|
||||
self._client.set_last_will(self._status_topic("online"), b"0",
|
||||
retain=True)
|
||||
self._client.connect()
|
||||
self._client.subscribe(self._set_topic())
|
||||
self._client.publish(self._status_topic("online"), b"1", retain=True)
|
||||
self.connected = True
|
||||
self._retry = RETRY_START
|
||||
self._zuletzt = {} # nach Reconnect alles neu senden
|
||||
self._log("MQTT: verbunden mit %s:%d" % (self._broker, self._port))
|
||||
return True
|
||||
except Exception as e: # OSError, MQTTException, ...
|
||||
self._log("MQTT: Verbindung fehlgeschlagen (%s)" % e)
|
||||
self._abwerfen()
|
||||
return False
|
||||
|
||||
def _abwerfen(self):
|
||||
"""Verbindung als tot markieren und den naechsten Versuch terminieren."""
|
||||
self.connected = False
|
||||
try:
|
||||
if self._client is not None:
|
||||
self._client.disconnect()
|
||||
except Exception:
|
||||
pass # beim Aufraeumen ist alles egal
|
||||
self._client = None
|
||||
self._naechster_versuch = _deadline(self._retry)
|
||||
# Beim naechsten Mal laenger warten, damit ein dauerhaft toter Broker
|
||||
# nicht bei jedem Schleifendurchlauf Zeit kostet.
|
||||
self._retry = min(self._retry * 2, RETRY_MAX)
|
||||
|
||||
def ensure(self):
|
||||
"""Verbindung herstellen, sofern die Wartezeit abgelaufen ist."""
|
||||
if self.connected:
|
||||
return True
|
||||
if self._naechster_versuch is not None and not _expired(self._naechster_versuch):
|
||||
return False
|
||||
return self.connect()
|
||||
|
||||
# -- Empfang -----------------------------------------------------------
|
||||
def _on_message(self, topic, payload):
|
||||
"""Callback von umqtt. Faengt alles ab -- hier darf nichts durchschlagen."""
|
||||
try:
|
||||
name = topic.decode().rsplit("/", 1)[-1]
|
||||
roh = payload.decode().strip()
|
||||
except Exception:
|
||||
self._log("MQTT: unlesbare Nachricht verworfen")
|
||||
return
|
||||
|
||||
if name not in SETZBAR:
|
||||
self._log("MQTT: unbekannte Einstellung %r verworfen" % name)
|
||||
return
|
||||
try:
|
||||
wert = int(roh)
|
||||
except ValueError:
|
||||
self._log("MQTT: %s=%r ist keine ganze Zahl" % (name, roh))
|
||||
return
|
||||
|
||||
try:
|
||||
neu = settings.update({name: wert})
|
||||
except settings.SettingsError as e:
|
||||
self._log("MQTT: %s" % e)
|
||||
return
|
||||
|
||||
self._log("MQTT: %s = %d uebernommen" % (name, wert))
|
||||
if self._on_change:
|
||||
try:
|
||||
self._on_change(neu)
|
||||
except Exception as e:
|
||||
self._log("MQTT: on_change fehlgeschlagen (%s)" % e)
|
||||
|
||||
def check(self):
|
||||
"""Anstehende Nachrichten verarbeiten. Blockiert nicht."""
|
||||
if not self.ensure():
|
||||
return False
|
||||
try:
|
||||
self._client.check_msg()
|
||||
return True
|
||||
except Exception as e:
|
||||
self._log("MQTT: Empfang gestoert (%s)" % e)
|
||||
self._abwerfen()
|
||||
return False
|
||||
|
||||
# -- Senden ------------------------------------------------------------
|
||||
def publish(self, name, wert, retain=True, nur_bei_aenderung=True):
|
||||
"""Einen Statuswert senden.
|
||||
|
||||
nur_bei_aenderung spart Funkverkehr: Die Koordinaten aendern sich zwar
|
||||
staendig, LDR und Helligkeit aber kaum. Mit retain=True holt sich ein
|
||||
neu verbundener Client den letzten Stand von selbst ab.
|
||||
"""
|
||||
if not self.ensure():
|
||||
return False
|
||||
topic = self._status_topic(name)
|
||||
text = str(wert)
|
||||
if nur_bei_aenderung and self._zuletzt.get(topic) == text:
|
||||
return True
|
||||
try:
|
||||
self._client.publish(topic, text.encode(), retain=retain)
|
||||
self._zuletzt[topic] = text
|
||||
return True
|
||||
except Exception as e:
|
||||
self._log("MQTT: Senden von %s gestoert (%s)" % (name, e))
|
||||
self._abwerfen()
|
||||
return False
|
||||
|
||||
def publish_many(self, werte, retain=True):
|
||||
"""Mehrere Statuswerte senden. Liefert True, wenn alle durchgingen."""
|
||||
ok = True
|
||||
for name, wert in werte.items():
|
||||
if not self.publish(name, wert, retain=retain):
|
||||
ok = False
|
||||
return ok
|
||||
|
||||
def close(self):
|
||||
"""Sauber abmelden -- online=0 bleibt als retained Nachricht stehen."""
|
||||
try:
|
||||
if self.connected and self._client is not None:
|
||||
self._client.publish(self._status_topic("online"), b"0",
|
||||
retain=True)
|
||||
self._client.disconnect()
|
||||
except Exception:
|
||||
pass
|
||||
self.connected = False
|
||||
self._client = None
|
||||
|
||||
|
||||
def connect_from_config(log=None, on_change=None):
|
||||
"""Bruecke aus mqtt_config.py bauen und verbinden.
|
||||
|
||||
Liefert None, wenn keine Konfiguration da ist -- dann laeuft die Anzeige
|
||||
ohne MQTT weiter, was ausdruecklich in Ordnung ist.
|
||||
"""
|
||||
log = log or (lambda *a: None)
|
||||
try:
|
||||
import mqtt_config
|
||||
except ImportError:
|
||||
log("MQTT: mqtt_config.py fehlt, laeuft ohne MQTT")
|
||||
return None
|
||||
|
||||
bridge = MqttBridge(
|
||||
broker=getattr(mqtt_config, "BROKER", None),
|
||||
port=getattr(mqtt_config, "PORT", 1883),
|
||||
user=getattr(mqtt_config, "USER", ""),
|
||||
password=getattr(mqtt_config, "PASSWORD", ""),
|
||||
prefix=getattr(mqtt_config, "PREFIX", "grossanzeige"),
|
||||
client_id=getattr(mqtt_config, "CLIENT_ID", "grossanzeige"),
|
||||
log=log,
|
||||
on_change=on_change,
|
||||
)
|
||||
bridge.connect() # scheitert es, versucht check() es spaeter erneut
|
||||
return bridge
|
||||
Reference in New Issue
Block a user