#!/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 io://... Tahoma Logic das gerechnete Geraet "Zeitpunkt" (Uhrzeit, Datum, Sonne) Alle liegen in einer Datei statt in einem Paket wie bei deviceDiscovery: es sind vier kurze Klassen, und wer eine fuenfte 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 from datetime import datetime from urllib.parse import quote logger = logging.getLogger("autoaction.transport") 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 haengen die Parameter als Abfrageargumente an die actor_url. Die Schaltbefehle von Shelly heissen aber nicht so, wie sie in der Datenbank stehen - dafuer steht unten eine kleine Uebersetzungstabelle. Wer ein anderes HTTP-Geraet anschliesst, erweitert genau die. """ schema = "http" SHELLY_BEFEHLE = { "turn_on": {"on": "true"}, "turn_off": {"on": "false"}, "toggle": {"toggle": "true"}, } 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 = dict(self.SHELLY_BEFEHLE.get(aktion["command_url"], {})) for p in aktion["params"]: if p["url"]: argumente[p["url"]] = p["wert"] if not argumente and aktion["command_url"]: # Kein bekannter Schaltbefehl und keine Parameter: das Kommando # als Argument anhaengen, mehr laesst sich hier nicht ableiten. argumente[aktion["command_url"]] = "true" 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"])) # --------------------------------------------------------------------------- # Tahoma # --------------------------------------------------------------------------- class TahomaTransport(Transport): """ Die actor_url ist die deviceURL (io://...), 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". """ 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 = [] @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")