Ein Lauf mit allen Modulen brachte 54 Geraete und 441 Messwerte - und fuenf Fehler ans Licht, die vorher niemand sehen konnte, weil clear_tables die Tabellen bei jedem Lauf geleert hat. Vier davon sind Schluessel- und Upsert-Fehler derselben Familie: * command_parameters war ueber (command_id, url) eindeutig. WLED-Parameter haben keine URL - ihr Name ist der Platzhalter in der Kommando-Vorlage -, und eine NULL kollidiert in MySQL nie. ON DUPLICATE KEY UPDATE griff also nicht, und jeder Lauf legte dieselben Parameter erneut an: aus 18 wurden nach drei Laeufen 54. Schluessel jetzt (command_id, parameter_name). * actor_states war ueber (actor_id, url) eindeutig. Umgekehrtes Problem: Messwerte, die sich ein Topic teilen, ueberschrieben einander. Der go-eCharger schickt sechzehn Werte als JSON-Feld auf einem Topic - von denen kam genau einer in der Datenbank an. Schluessel jetzt (actor_id, state_name), das bringt 50 verlorene Messwerte zurueck. * Kommandos, Messwerte und Parameter aktualisierten ihre URL beim Wiederholungslauf nicht - sie stand nicht im UPDATE-Teil. Eine Korrektur in einem Modul kam damit nie in einer bestehenden Datenbank an. * Der Typ eines Kommando-Parameters wurde als Text in eine int-Spalte geschrieben. MariaDB macht daraus stillschweigend 0, und 0 ist "bool" - deshalb bot der Editor fuer die WLED-Helligkeit (0..255) ein Ja/Nein an. Der fuenfte: die Shelly-Gen1-Relais hatten weder eine Kommando-URL noch eine URL am Zustand. Ohne die weiss niemand, was zu schicken und wo nachzusehen ist. Gen1 schaltet ueber ?turn=on|off|toggle und meldet sich in "ison". Damit entfaellt auch die Uebersetzungstabelle im HTTPTransport: die Zuordnung gehoert ins Geraetemodell, nicht in den Runner. Zwei Luecken auf der Runner-Seite, die derselbe Lauf gezeigt hat: * WLED liess sich gar nicht ansteuern - es gab keinen Transport dafuer. Der neue nutzt die JSON-Vorlage aus command_url, fuellt die Platzhalter und schickt sie am Stueck an /json/state. * Tahoma wurde am Schema "io://" erkannt. Das Schema beschreibt aber die Funkart: dieselbe Box liefert rts:// fuer die Dachfenster und internal:// fuer die Alarmanlage. Drei von 22 Geraeten fielen durch. Erkannt wird jetzt an der Box-Kennung in der URL. Nachgemessen: zwei Laeufe hintereinander aendern keine einzige Zeile mehr, die Geraete-IDs bleiben ueber Laeufe stabil (Voraussetzung fuer die Fremdschluessel der Automatiken), alle 54 Geraete finden genau einen Transport, und 440 der 441 Messwerte sind tatsaechlich lesbar. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
441 lines
17 KiB
Python
441 lines
17 KiB
Python
#!/usr/bin/env python3
|
|
"""
|
|
Transporte fuer den AutoAction-Runner.
|
|
|
|
Ein Transport weiss, wie man mit einer Sorte Geraet redet - er liest deren
|
|
Messwerte und schickt deren Kommandos. Welcher zustaendig ist, entscheidet die
|
|
URL des Aktors in der Tabelle `actors`:
|
|
|
|
mqtt://... MQTT-Geraete (Home-Assistant-Discovery)
|
|
http://... Shelly und andere HTTP-Geraete
|
|
wled://... WLED-Lampen
|
|
io://, rts://, internal://, ogp:// Tahoma - das Schema haengt an der
|
|
Funkart des Geraets, deshalb wird dort nicht danach
|
|
entschieden, sondern an der Box-Kennung in der URL
|
|
Logic das gerechnete Geraet "Zeitpunkt" (Uhrzeit, Datum, Sonne)
|
|
|
|
Alle liegen in einer Datei statt in einem Paket wie bei deviceDiscovery: es
|
|
sind fuenf kurze Klassen, und wer eine sechste Geraeteart anschliesst, sieht
|
|
hier auf einen Blick, was dafuer zu tun ist.
|
|
|
|
Jeder Transport hat zwei Haelften:
|
|
|
|
zustaende_anmelden(states) einmalig beim Start
|
|
zustaende_lesen() liefert {state_id: wert} - nur was neu ist
|
|
senden(aktion) fuehrt ein Kommando aus
|
|
|
|
`states` ist eine Liste von Dicts mit actor_url, state_url und id,
|
|
`aktion` ein Dict mit actor_url, command_url und params (Liste aus
|
|
{url, name, wert}).
|
|
"""
|
|
|
|
import json
|
|
import logging
|
|
import re
|
|
import urllib3
|
|
from datetime import datetime
|
|
from urllib.parse import quote
|
|
|
|
logger = logging.getLogger("autoaction.transport")
|
|
|
|
# Die Tahoma-Box hat ein selbst ausgestelltes Zertifikat auf einen Namen, den
|
|
# nur das Heimnetz kennt. Die Pruefung ist dort bewusst aus (wie in
|
|
# ajax/tahoma.php); ohne diese Zeile warnt urllib3 bei jeder einzelnen
|
|
# Abfrage und uebertoent das eigentliche Protokoll.
|
|
urllib3.disable_warnings(urllib3.exceptions.InsecureRequestWarning)
|
|
|
|
|
|
class Transport:
|
|
"""Gemeinsame Form. Wer nichts zu lesen hat, erbt die leeren Methoden."""
|
|
|
|
schema = ""
|
|
|
|
def passt(self, actor_url):
|
|
return actor_url.startswith(self.schema)
|
|
|
|
def zustaende_anmelden(self, states):
|
|
pass
|
|
|
|
def zustaende_lesen(self):
|
|
return {}
|
|
|
|
def senden(self, aktion):
|
|
raise NotImplementedError
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# MQTT
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class MQTTTransport(Transport):
|
|
"""
|
|
Die state_url ist hier das vollstaendige Topic, die parameter_url das
|
|
Kommando-Topic. Werte kommen von selbst herein und landen in einem
|
|
Zwischenspeicher, den der Runner im Takt abholt.
|
|
"""
|
|
|
|
schema = "mqtt://"
|
|
|
|
def __init__(self, client, dry_run=False):
|
|
self.client = client
|
|
self.dry_run = dry_run
|
|
self.topics = {} # topic -> [state_id, ...]
|
|
self.neu = {} # state_id -> wert
|
|
|
|
def zustaende_anmelden(self, states):
|
|
self.topics = {}
|
|
for s in states:
|
|
if not s["state_url"]:
|
|
continue
|
|
self.topics.setdefault(s["state_url"], []).append(s["id"])
|
|
for topic in self.topics:
|
|
self.client.subscribe(topic)
|
|
logger.info("MQTT: %d Topics abonniert", len(self.topics))
|
|
|
|
def nachricht(self, topic, payload):
|
|
"""Wird vom Runner aus dem on_message-Rueckruf gerufen."""
|
|
for state_id in self.topics.get(topic, []):
|
|
self.neu[state_id] = self._wert(payload)
|
|
|
|
@staticmethod
|
|
def _wert(payload):
|
|
"""
|
|
Rohtext, aber JSON-Skalare werden ausgepackt: manche Geraete schicken
|
|
"21.4", andere 21.4 mit Anfuehrungszeichen, wieder andere true statt
|
|
ON. Objekte bleiben wie sie sind - welches Feld gemeint ist, weiss die
|
|
Datenbank nicht, und Raten waere schlimmer als Nichtstun.
|
|
"""
|
|
text = payload.decode("utf-8", "replace").strip() if isinstance(payload, bytes) else str(payload)
|
|
try:
|
|
wert = json.loads(text)
|
|
except ValueError:
|
|
return text
|
|
if isinstance(wert, bool):
|
|
return "true" if wert else "false"
|
|
if isinstance(wert, (int, float, str)):
|
|
return str(wert)
|
|
return text
|
|
|
|
def zustaende_lesen(self):
|
|
werte, self.neu = self.neu, {}
|
|
return werte
|
|
|
|
def senden(self, aktion):
|
|
if not aktion["params"]:
|
|
# Kommando ohne Parameter: das Kommando selbst ist die Nutzlast.
|
|
self._publish(aktion["command_url"], "")
|
|
return
|
|
for p in aktion["params"]:
|
|
self._publish(p["url"] or aktion["command_url"], p["wert"])
|
|
|
|
def _publish(self, topic, nutzlast):
|
|
if not topic:
|
|
raise ValueError("Kommando ohne Topic")
|
|
if self.dry_run:
|
|
logger.info("[dry-run] MQTT %s <- %s", topic, nutzlast)
|
|
return
|
|
ergebnis = self.client.publish(topic, nutzlast, qos=1, retain=False)
|
|
ergebnis.wait_for_publish(timeout=5)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# HTTP (Shelly und Verwandte)
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class HTTPTransport(Transport):
|
|
"""
|
|
Die actor_url ist der Endpunkt, die state_url ein Feldname in dessen
|
|
JSON-Antwort ("tC", "a_voltage"). Solche Geraete melden sich nicht von
|
|
selbst, sie werden im Poll-Takt gefragt.
|
|
|
|
Beim Senden wird die command_url als Abfrageargument an die actor_url
|
|
gehaengt ("turn=on") und die Parameter mit ihrem eigenen Namen dazu.
|
|
Frueher stand hier eine Uebersetzungstabelle, weil die Shelly-Kommandos
|
|
ohne URL in der Datenbank landeten - das ist im Discovery behoben, die
|
|
Zuordnung gehoert dorthin und nicht in den Runner.
|
|
"""
|
|
|
|
schema = "http"
|
|
|
|
def __init__(self, requests_modul, timeout=5, dry_run=False):
|
|
self.requests = requests_modul
|
|
self.timeout = timeout
|
|
self.dry_run = dry_run
|
|
self.states = []
|
|
|
|
def zustaende_anmelden(self, states):
|
|
self.states = [s for s in states if s["state_url"]]
|
|
logger.info("HTTP: %d Messwerte an %d Endpunkten",
|
|
len(self.states), len({s["actor_url"] for s in self.states}))
|
|
|
|
def zustaende_lesen(self):
|
|
werte = {}
|
|
# Je Endpunkt eine Anfrage, auch wenn mehrere Messwerte daran haengen.
|
|
nach_url = {}
|
|
for s in self.states:
|
|
nach_url.setdefault(s["actor_url"], []).append(s)
|
|
for url, states in nach_url.items():
|
|
try:
|
|
antwort = self.requests.get(url, timeout=self.timeout)
|
|
daten = antwort.json()
|
|
except Exception as fehler:
|
|
logger.debug("HTTP %s nicht erreichbar: %s", url, fehler)
|
|
continue
|
|
for s in states:
|
|
if isinstance(daten, dict) and s["state_url"] in daten:
|
|
werte[s["id"]] = str(daten[s["state_url"]])
|
|
return werte
|
|
|
|
def senden(self, aktion):
|
|
argumente = {}
|
|
for teil in (aktion["command_url"] or "").split("&"):
|
|
if "=" in teil:
|
|
schluessel, wert = teil.split("=", 1)
|
|
argumente[schluessel] = wert
|
|
for p in aktion["params"]:
|
|
if p["url"]:
|
|
argumente[p["url"]] = p["wert"]
|
|
if not argumente:
|
|
raise RuntimeError("Kommando ohne URL und ohne Parameter - im "
|
|
"Geraetemodell fehlt die Angabe, was zu schicken ist")
|
|
if self.dry_run:
|
|
logger.info("[dry-run] HTTP %s %s", aktion["actor_url"], argumente)
|
|
return
|
|
antwort = self.requests.get(aktion["actor_url"], params=argumente, timeout=self.timeout)
|
|
if antwort.status_code >= 400:
|
|
raise RuntimeError("HTTP %d von %s" % (antwort.status_code, aktion["actor_url"]))
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# WLED
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class WLEDTransport(Transport):
|
|
"""
|
|
WLED-Lampen sprechen ueber eine einzige JSON-Schnittstelle:
|
|
GET http://IP/json/state liefert den Zustand, POST dorthin setzt ihn.
|
|
|
|
Das Geraetemodell nutzt das elegant aus - die command_url ist eine
|
|
JSON-Vorlage mit Platzhaltern:
|
|
|
|
{"bri":%brightness%}
|
|
{"seg":[{"col":[[%red%,%green%,%blue%]]}]}
|
|
|
|
Gesendet wird also nicht Argument fuer Argument, sondern die ausgefuellte
|
|
Vorlage am Stueck. Deshalb haben die Parameter hier auch keine eigene URL:
|
|
ihr Name ist der Platzhalter.
|
|
|
|
Die state_url ist ein Pfad in die Antwort ("on", "bri",
|
|
"seg[0].col[0]") - dieselbe Schreibweise, die auch in current_value steht.
|
|
"""
|
|
|
|
schema = "wled://"
|
|
|
|
def __init__(self, requests_modul, timeout=5, dry_run=False):
|
|
self.requests = requests_modul
|
|
self.timeout = timeout
|
|
self.dry_run = dry_run
|
|
self.states = []
|
|
|
|
@staticmethod
|
|
def _adresse(actor_url):
|
|
return "http://" + actor_url[len("wled://"):].rstrip("/")
|
|
|
|
def zustaende_anmelden(self, states):
|
|
self.states = [s for s in states if s["state_url"]]
|
|
logger.info("WLED: %d Messwerte an %d Lampen",
|
|
len(self.states), len({s["actor_url"] for s in self.states}))
|
|
|
|
@staticmethod
|
|
def _pfad(daten, pfad):
|
|
""""seg[0].col[0]" in der Antwort nachschlagen."""
|
|
for teil in pfad.split("."):
|
|
treffer = re.match(r"^([^\[]*)((?:\[\d+\])*)$", teil)
|
|
if not treffer:
|
|
return None
|
|
if treffer.group(1):
|
|
daten = daten[treffer.group(1)]
|
|
for index in re.findall(r"\[(\d+)\]", treffer.group(2)):
|
|
daten = daten[int(index)]
|
|
return daten
|
|
|
|
def zustaende_lesen(self):
|
|
werte = {}
|
|
nach_lampe = {}
|
|
for s in self.states:
|
|
nach_lampe.setdefault(s["actor_url"], []).append(s)
|
|
for actor_url, states in nach_lampe.items():
|
|
try:
|
|
antwort = self.requests.get(self._adresse(actor_url) + "/json/state",
|
|
timeout=self.timeout)
|
|
daten = antwort.json()
|
|
except Exception as fehler:
|
|
logger.debug("WLED %s nicht erreichbar: %s", actor_url, fehler)
|
|
continue
|
|
for s in states:
|
|
try:
|
|
werte[s["id"]] = str(self._pfad(daten, s["state_url"]))
|
|
except (KeyError, IndexError, TypeError):
|
|
logger.debug("WLED %s: Pfad %s nicht gefunden", actor_url, s["state_url"])
|
|
return werte
|
|
|
|
def senden(self, aktion):
|
|
vorlage = aktion["command_url"]
|
|
if not vorlage:
|
|
raise RuntimeError("WLED-Kommando ohne Vorlage")
|
|
for p in aktion["params"]:
|
|
vorlage = vorlage.replace("%" + p["name"] + "%", str(p["wert"]))
|
|
try:
|
|
rumpf = json.loads(vorlage)
|
|
except ValueError:
|
|
# Ein nicht ersetzter Platzhalter oder ein Textwert an einer
|
|
# Stelle, wo eine Zahl stehen muss. Lieber hier abbrechen als der
|
|
# Lampe etwas Unverstaendliches schicken.
|
|
raise RuntimeError("WLED-Vorlage ergibt kein gueltiges JSON: " + vorlage[:120])
|
|
if self.dry_run:
|
|
logger.info("[dry-run] WLED %s <- %s", aktion["actor_url"],
|
|
json.dumps(rumpf, ensure_ascii=False))
|
|
return
|
|
antwort = self.requests.post(self._adresse(aktion["actor_url"]) + "/json/state",
|
|
json=rumpf, timeout=self.timeout)
|
|
if antwort.status_code >= 400:
|
|
raise RuntimeError("WLED antwortete mit %d" % antwort.status_code)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Tahoma
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class TahomaTransport(Transport):
|
|
"""
|
|
Die actor_url ist die deviceURL, die state_url ein Statusname
|
|
("core:ClosureState"), die command_url ein Kommandoname ("setClosure").
|
|
Geschickt wird ueber exec/apply - genauso wie in ajax/tahoma.php, nur ohne
|
|
die dortige Sonderbehandlung fuer "faehrt gerade".
|
|
|
|
Zustaendig ist dieser Transport fuer alles, was die Kennung der eigenen
|
|
Box in der URL traegt. Am Schema laesst sich das nicht festmachen: es
|
|
beschreibt die Funkart, und dieselbe Box liefert io:// fuer die
|
|
Jalousien, rts:// fuer die Dachfenster und internal:// fuer die Alarm-
|
|
anlage. Ohne PIN in der config.ini ist niemand zustaendig - dann meldet
|
|
der Runner beim Ausloesen "kein Transport", statt still nichts zu tun.
|
|
"""
|
|
|
|
schema = "io://"
|
|
|
|
def __init__(self, requests_modul, pin, token, timeout=10, dry_run=False):
|
|
self.requests = requests_modul
|
|
self.pin = pin
|
|
self.token = token
|
|
self.timeout = timeout
|
|
self.dry_run = dry_run
|
|
self.states = []
|
|
|
|
def passt(self, actor_url):
|
|
return bool(self.pin) and ("://" + self.pin + "/") in actor_url
|
|
|
|
@property
|
|
def basis(self):
|
|
return "https://gateway-%s:8443/enduser-mobile-web/1/enduserAPI" % self.pin
|
|
|
|
def _kopf(self):
|
|
return {"Content-Type": "application/json",
|
|
"Authorization": "Bearer " + self.token}
|
|
|
|
def zustaende_anmelden(self, states):
|
|
self.states = [s for s in states if s["state_url"]]
|
|
logger.info("Tahoma: %d Messwerte an %d Geraeten",
|
|
len(self.states), len({s["actor_url"] for s in self.states}))
|
|
|
|
def zustaende_lesen(self):
|
|
if not self.token:
|
|
return {}
|
|
werte = {}
|
|
nach_geraet = {}
|
|
for s in self.states:
|
|
nach_geraet.setdefault(s["actor_url"], []).append(s)
|
|
for geraet, states in nach_geraet.items():
|
|
try:
|
|
antwort = self.requests.get(
|
|
self.basis + "/setup/devices/" + quote(geraet, safe="") + "/states",
|
|
headers=self._kopf(), timeout=self.timeout, verify=False)
|
|
zustaende = {z["name"]: z.get("value") for z in antwort.json()}
|
|
except Exception as fehler:
|
|
logger.debug("Tahoma %s nicht erreichbar: %s", geraet, fehler)
|
|
continue
|
|
for s in states:
|
|
if s["state_url"] in zustaende:
|
|
werte[s["id"]] = str(zustaende[s["state_url"]])
|
|
return werte
|
|
|
|
def senden(self, aktion):
|
|
if not self.token:
|
|
raise RuntimeError("Kein Tahoma-Token in der config.ini")
|
|
# Die Reihenfolge der Parameter ist die aus command_parameters - bei
|
|
# setClosureAndOrientation also erst Position, dann Winkel.
|
|
befehl = {"name": aktion["command_url"],
|
|
"parameters": [self._zahl(p["wert"]) for p in aktion["params"]]}
|
|
rumpf = {"label": "AutoAction",
|
|
"actions": [{"deviceURL": aktion["actor_url"], "commands": [befehl]}]}
|
|
if self.dry_run:
|
|
logger.info("[dry-run] Tahoma %s", json.dumps(rumpf, ensure_ascii=False))
|
|
return
|
|
antwort = self.requests.post(self.basis + "/exec/apply", headers=self._kopf(),
|
|
data=json.dumps(rumpf), timeout=self.timeout, verify=False)
|
|
if antwort.status_code >= 400:
|
|
raise RuntimeError("Tahoma antwortete mit %d: %s"
|
|
% (antwort.status_code, antwort.text[:120]))
|
|
|
|
@staticmethod
|
|
def _zahl(wert):
|
|
"""Tahoma erwartet Zahlen als Zahlen, Text als Text."""
|
|
try:
|
|
return int(wert)
|
|
except (TypeError, ValueError):
|
|
pass
|
|
try:
|
|
return float(wert)
|
|
except (TypeError, ValueError):
|
|
return wert
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Logic - das gerechnete Geraet
|
|
# ---------------------------------------------------------------------------
|
|
|
|
class LogicTransport(Transport):
|
|
"""
|
|
Uhrzeit, Datum, Sonnenauf- und -untergang. Es gibt nichts zu abonnieren und
|
|
nichts zu schalten, die Werte entstehen im Takt. Sonnenzeiten kommen aus
|
|
solarLog.daylight, dieselbe Tabelle, aus der auch ajax/getSunrise.php liest.
|
|
"""
|
|
|
|
schema = "Logic"
|
|
|
|
def __init__(self, sonnenzeiten):
|
|
"""sonnenzeiten: Funktion() -> (sonnenaufgang, sonnenuntergang) als "HH:MM"."""
|
|
self.sonnenzeiten = sonnenzeiten
|
|
self.states = []
|
|
|
|
def passt(self, actor_url):
|
|
return actor_url == "Logic"
|
|
|
|
def zustaende_anmelden(self, states):
|
|
self.states = states
|
|
logger.info("Logic: %d Messwerte", len(states))
|
|
|
|
def zustaende_lesen(self):
|
|
jetzt = datetime.now()
|
|
auf, unter = self.sonnenzeiten()
|
|
tabelle = {
|
|
"time": jetzt.strftime("%H:%M"),
|
|
"date": jetzt.strftime("%d.%m.%Y"),
|
|
"sunrise": auf,
|
|
"sunset": unter,
|
|
}
|
|
return {s["id"]: tabelle[s["state_url"]]
|
|
for s in self.states if s["state_url"] in tabelle}
|
|
|
|
def senden(self, aktion):
|
|
raise RuntimeError("Das Geraet \"Zeitpunkt\" kann nichts schalten")
|