"""MQTT-Anbindung: Einstellungen empfangen, Zustand veroeffentlichen. /set/hell_prozent 60 eingehend, schreibt settings.json /set/dunkel_prozent 5 /set/schwelle 1800 /set/hysterese 150 /set/mount_host 192.168.1.115 Adresse der Montierung; leer = Vorgabe aus config.MOUNT_HOST /status/ra 18h36m56s ausgehend /status/dec +38°47'01" /status/ldr 2453 /status/helligkeit 50 /status/mount_host 192.168.1.115 Adresse, die gerade abgefragt wird /status/link 1 Montierung erreichbar /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 /set/ 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 # Zahl oder Text? Entscheidet der Typ der Vorgabe -- dieselbe Regel wie # in settings._pruefe. Eine leere Nachricht ist deshalb nur bei # Textwerten sinnvoll: mount_host faellt damit auf config.MOUNT_HOST # zurueck, waehrend int("") wie gehabt als Unsinn abgelehnt wird. if isinstance(settings.DEFAULTS[name], str): wert = roh else: 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 = %r 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