Runner zieht zum SolarManager um
Der AutoAction-Runner ist ein Hintergrundprozess und gehoert damit zu den anderen, nicht ins Web-Verzeichnis. Er liegt jetzt unter /volume1/homes/wagner/SolarManager/autoActions und wird von dort zusammen mit dem Manager gestartet; versioniert ist er im Repository SolarManager. Browser und Runner reden ohnehin nur ueber die Datenbank homeMesh miteinander - das Web-UI kennt den Pfad nicht und muss ihn nicht kennen. Die Verweise in commands.php und homeMesh_automations.sql zeigen bereits auf den neuen Ort. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -1,215 +0,0 @@
|
||||
# AutoActions
|
||||
|
||||
Automatiken, die im Web-UI unter „Automatismen" angelegt werden und hier
|
||||
ausgeführt werden: *wenn Bedingung, dann Kommando*.
|
||||
|
||||
```
|
||||
Browser Datenbank homeMesh Runner
|
||||
────────────────────── ───────────────────────── ────────────────────
|
||||
Karte „Automatismen" ──▶ automations ──▶ autoaction_runner.py
|
||||
ajax/AutoAction.php automation_conditions MQTT / HTTP / Tahoma
|
||||
restricted/automations.php automation_actions ──▶ Geräte
|
||||
js/solar/autoActionFuncs.js automation_action_params
|
||||
automation_log
|
||||
calendar_days ◀── fetch_calendar.py
|
||||
```
|
||||
|
||||
## Wie eine Automatik aufgebaut ist
|
||||
|
||||
Eine Automatik hat **Auslöser**, **Rahmenbedingungen** und **Aktionen**.
|
||||
|
||||
Ein Auslöser vergleicht einen Messwert (`actor_states`) mit einer Schwelle.
|
||||
Mehrere Auslöser werden über `group_no` verknüpft: gleiche Nummer heißt UND,
|
||||
verschiedene Nummern heißen ODER — ausgewertet wird `any(all(gruppe))`. Im
|
||||
Editor ist eine Gruppe ein gerahmter Block mit eigenem „+ Bedingung", zwischen
|
||||
den Blöcken steht ein ODER. Die Klammerung ist damit gezeichnet und nicht bloß
|
||||
vereinbart, und darunter steht derselbe Ausdruck noch einmal als Satz.
|
||||
|
||||
Eine Aktion ist ein Kommando (`actor_commands`) mit einem Wert je Parameter
|
||||
(`command_parameters`) — nicht vier feste Spalten „Wert 1" bis „Wert 4",
|
||||
sondern so viele Zeilen, wie das Gerät Parameter hat.
|
||||
|
||||
Die Rahmenbedingungen (Wochentage, Zeitfenster, Ferien, Feiertage) sagen, wann
|
||||
die Automatik überhaupt hinsehen darf.
|
||||
|
||||
## Warum ein Dauerläufer und kein Cronjob
|
||||
|
||||
Zwei Gründe:
|
||||
|
||||
* Schwellwert-Auslöser sollen greifen, wenn die MQTT-Nachricht hereinkommt,
|
||||
nicht erst im nächsten Minutenraster.
|
||||
* `actor_states.current_value` wird sonst von niemandem fortgeschrieben — es
|
||||
wird beim Geräte-Discovery einmal gesetzt und danach nie wieder. Ein
|
||||
zustandsloser Cronjob hätte gar nichts, womit er vergleichen könnte. Der
|
||||
Runner pflegt den Wert nebenbei mit (gedrosselt auf einmal je Minute), wovon
|
||||
auch der Editor profitiert: er zeigt neben jedem Messwert den aktuellen Stand.
|
||||
|
||||
## Nur die steigende Flanke
|
||||
|
||||
`automations.cond_met` hält fest, ob die Bedingung beim letzten Durchlauf schon
|
||||
erfüllt war. Ohne das würde „Temperatur über 22 Grad" bei jedem Takt erneut
|
||||
feuern. Verlässt die Automatik ihr Zeitfenster, wird die Flanke
|
||||
zurückgesetzt, damit sie im nächsten Fenster wieder steigen kann.
|
||||
|
||||
Zeit-Auslöser gibt es in drei Formen:
|
||||
|
||||
| | wahr, wenn |
|
||||
|---|---|
|
||||
| `um 16:30` | genau in dieser Minute |
|
||||
| `ab 16:30` | von da an bis Mitternacht |
|
||||
| `vor 16:30` | bis dahin |
|
||||
|
||||
Ausgelöst wird in allen drei Fällen nur einmal, eben wegen der Flanke. `ab`
|
||||
ist trotzdem das robustere: fällt der Runner in genau der Minute aus, auf die
|
||||
`um` zeigt, ist die Automatik für den Tag verloren — bei `ab` holt der nächste
|
||||
Takt es nach. Wer `um` braucht und den Ausfall nicht riskieren will, hakt
|
||||
zusätzlich `force_once` an. Die breiten Zeitfenster in `auto_watering.py`
|
||||
folgen derselben Überlegung.
|
||||
|
||||
Beim Sonnenauf- und -untergang ist der Wert ein Versatz, und der kann davor
|
||||
oder danach liegen — deshalb dieselben drei Fälle mal zwei: `+ 00:30` eine
|
||||
halbe Stunde nach Sonnenaufgang, `ab - 00:30` ab einer halben Stunde davor,
|
||||
`vor + 00:30` bis eine halbe Stunde danach.
|
||||
|
||||
Über den Tagesrand wird gerechnet, nicht abgeschnitten: „sechs Stunden vor
|
||||
Sonnenaufgang" landet am Vorabend, und das ist so gewollt — abgeschnitten
|
||||
wären solche Angaben gar nicht mehr formulierbar. Verglichen wird die Uhrzeit
|
||||
innerhalb des Tages; ein Ziel jenseits von Mitternacht gilt als diese Uhrzeit
|
||||
am selben Tag. Bei `+` und `-` ist das genau der gemeinte Zeitpunkt, bei `ab`
|
||||
und `vor` verschiebt sich der wahre Bereich entsprechend mit.
|
||||
|
||||
`force_once` („am Ende des Zeitraums auf jeden Fall ausführen") greift, wenn
|
||||
das Fenster zugeht und in diesem Fenster noch nichts passiert ist.
|
||||
|
||||
## Sperrzeit
|
||||
|
||||
Die Flanke allein schützt nicht gegen einen Messwert, der um die Schwelle
|
||||
**pendelt**: „Temperatur > 22" bei 22,1 / 21,9 / 22,1 °C ist jedes Mal eine
|
||||
echte steigende Flanke, und über MQTT können die Werte im Sekundentakt
|
||||
hereinkommen. `automations.lockout_secs` sagt, wie lange nach einer Auslösung
|
||||
nicht wieder geschaltet wird. Der Editor bietet drei Stufen an:
|
||||
|
||||
| | | gedacht für |
|
||||
|---|---|---|
|
||||
| Ohne | 0 s | volle Geschwindigkeit, jede Flanke schaltet |
|
||||
| Kurz | 60 s | Licht, Farbe, Dimmwert |
|
||||
| Lang | 900 s | Rollläden, Ventile, alles mit Motor |
|
||||
|
||||
Gespeichert werden Sekunden, angeboten werden nur die drei Stufen — eine
|
||||
vierte ist damit eine Zeile in `lockoutChoices()` und keine Wanderung durch
|
||||
die Datenbank.
|
||||
|
||||
Eine Flanke innerhalb der Sperrzeit wird **verworfen, nicht aufgehoben**. Ein
|
||||
Rollladen, der eine Viertelstunde später doch noch losfährt, weil vor langer
|
||||
Zeit einmal eine Schwelle gestreift wurde, wäre unangenehmer als einer, der
|
||||
gar nicht fährt — und der nächste echte Anlass nach Ablauf der Sperre kommt
|
||||
ohnehin durch. Verworfene Flanken stehen im Log auf `DEBUG`, nicht in
|
||||
`automation_log`; bei einem zappelnden Sensor wäre die Tabelle sonst voll
|
||||
davon.
|
||||
|
||||
`force_once` ist von der Sperre nicht betroffen: es greift nur, wenn im
|
||||
Fenster gar nichts gelaufen ist — dann ist auch keine Sperre aktiv.
|
||||
|
||||
## Transporte
|
||||
|
||||
Welcher Weg zum Gerät führt, entscheidet die URL des Aktors in `actors`:
|
||||
|
||||
| URL | Messwert (`actor_states.url`) | Kommando (`actor_commands.command_url`) |
|
||||
|---|---|---|
|
||||
| `mqtt://…` | vollständiges Topic, abonniert; bei mehreren Messwerten je Topic zusätzlich `value_path` | Nutzlast auf das Parameter-Topic |
|
||||
| `http://…` | Feldname in der JSON-Antwort, gepollt | Abfrageargumente an die Geräte-URL (`turn=on`) |
|
||||
| `wled://…` | Pfad in `/json/state` (`seg[0].col[0]`), gepollt | JSON-Vorlage mit Platzhaltern, als Ganzes gesendet |
|
||||
| Tahoma | Statusname (`core:ClosureState`), gepollt | `exec/apply` an die Box |
|
||||
| `Logic` | gerechnet: Uhrzeit, Datum, Sonne | – |
|
||||
|
||||
Alle fünf stehen in `transports.py`. Eine sechste Geräteart kommt als weitere
|
||||
Klasse dazu; sie braucht `passt()`, `zustaende_lesen()` und `senden()`.
|
||||
|
||||
Tahoma ist der einzige, der nicht am URL-Schema erkannt wird, sondern an der
|
||||
**Box-Kennung** in der URL. Das Schema beschreibt dort die Funkart, und
|
||||
dieselbe Box liefert `io://` für die Jalousien, `rts://` für die Dachfenster
|
||||
und `internal://` für die Alarmanlage. Ohne `pin` in der `config.ini` ist
|
||||
niemand zuständig — dann meldet der Runner beim Auslösen „kein Transport",
|
||||
statt still nichts zu tun.
|
||||
|
||||
Mehrere Messwerte teilen sich oft **ein Topic**: der go-eCharger schickt
|
||||
sechzehn Zahlen als JSON-Feld auf `…/nrg`, und erst das `value_template` der
|
||||
Home-Assistant-Discovery sagt, dass „Strom L1" das fünfte Element ist. Diese
|
||||
Angabe steht in `actor_states.value_path` — in derselben Schreibweise, die
|
||||
auch WLED benutzt: `[4]`, `ssid`, `seg[0].col[0]`. Ohne Pfad gilt die ganze
|
||||
Nutzlast.
|
||||
|
||||
Gelesen wird nur der einfache Fall aus dem Template: ein Zugriff auf
|
||||
`value_json` und was danach an Punkten und Klammern folgt.
|
||||
|
||||
### Werttabellen
|
||||
|
||||
Manche Geräte schicken eine Zahl und meinen einen Zustand:
|
||||
|
||||
```
|
||||
{{ ['Unknown','Idle','Charging','WaitCar','Complete','Error'][value_json|int] }}
|
||||
{{ ['Default','Eco','NextTrip'][value_json|int-3] }}
|
||||
```
|
||||
|
||||
Das ist kein Pfad, sondern eine Übersetzung von Zahl nach Text. Sie landet in
|
||||
`possible_values` — in der Schreibweise, die WLED für seine Effektliste schon
|
||||
benutzt: eine Liste aus `{Wert: Bezeichnung}`. Ein Versatz im Ausdruck wandert
|
||||
dabei in die Schlüssel, aus `[value_json|int-3]` wird also `{"3":"Default"}`.
|
||||
|
||||
Der Runner übersetzt beim Lesen: aus der gesendeten `2` wird `Charging`. Eine
|
||||
Bedingung vergleicht damit genau den Klartext, den der Editor zur Auswahl
|
||||
stellt. Steht die Zahl nicht in der Tabelle, bleibt sie stehen — ein
|
||||
erfundener Name wäre schlimmer als ein roher Wert.
|
||||
|
||||
Auf beiden Seiten des Editors steckt dieselbe Tabelle, aber der gespeicherte
|
||||
Wert ist ein anderer:
|
||||
|
||||
| | angezeigt | gespeichert |
|
||||
|---|---|---|
|
||||
| **Messwert** (Bedingung) | `Charging` | `Charging` — der Runner hat schon übersetzt |
|
||||
| **Parameter** (Aktion) | `Blink` | `1` — das Gerät will die Zahl |
|
||||
|
||||
Bei WLED trägt die Kommando-Vorlage alles: `{"seg":[{"col":[[%red%,%green%,%blue%]]}]}`
|
||||
wird mit den Parameterwerten gefüllt und am Stück geschickt. Deshalb haben die
|
||||
Parameter dort keine eigene URL — ihr Name *ist* der Platzhalter.
|
||||
|
||||
## Voraussetzungen
|
||||
|
||||
* Python 3 mit `pymysql`, `requests`, `paho-mqtt`
|
||||
* `homeMesh_automations.sql` einmal eingespielt
|
||||
* In `../deviceDiscovery/config.ini` muss **`clear_tables = false`** stehen.
|
||||
Discovery schreibt mit `ON DUPLICATE KEY UPDATE` auf den URLs, das Leeren ist
|
||||
unnötig — ein `TRUNCATE` würde dagegen die Geräte-IDs neu vergeben, und die
|
||||
Automatiken zeigen per Fremdschlüssel genau auf diese IDs.
|
||||
|
||||
## Einrichten
|
||||
|
||||
```bash
|
||||
cp config.ini.example config.ini # ausfüllen: Datenbank, MQTT, Tahoma
|
||||
python3 fetch_calendar.py # Feiertage und Ferien holen
|
||||
python3 autoaction_runner.py --once --dry-run --verbose # Probelauf
|
||||
```
|
||||
|
||||
`--dry-run` schaltet nichts, protokolliert aber jedes Kommando, das geschickt
|
||||
würde. `--once` macht einen einzigen Durchlauf.
|
||||
|
||||
Im Dauerbetrieb wird `autoaction_runner.py` beim Booten gestartet (auf der
|
||||
Synology über den Aufgabenplaner, Ereignis „Hochfahren", als root). Er
|
||||
verbindet sich selbst neu, wenn MQTT wegbricht, und lädt das Regelwerk nach,
|
||||
sobald im Browser etwas gespeichert wurde — ein Neustart nach jeder Änderung
|
||||
ist nicht nötig.
|
||||
|
||||
`fetch_calendar.py` gehört einmal jährlich in den Cron:
|
||||
|
||||
```
|
||||
0 4 1 1 * /usr/bin/python3 /volume1/web/smart/restricted/autoActions/fetch_calendar.py
|
||||
```
|
||||
|
||||
Ein zusätzlicher Lauf im Herbst schadet nicht — die Ferientermine des
|
||||
übernächsten Schuljahres stehen erst später fest.
|
||||
|
||||
## Nachsehen, was passiert ist
|
||||
|
||||
`automation_log` hält je Auslösung fest, ob sie durchlief (`fired`), wegen
|
||||
`force_once` nachgeholt wurde (`forced`) oder scheiterte (`error`, mit Grund in
|
||||
`detail`). Einträge älter als 30 Tage räumt der Runner selbst weg.
|
||||
@@ -1,848 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
AutoAction-Runner - fuehrt die im Web-UI angelegten Automatiken aus.
|
||||
|
||||
Laeuft als Dauerprozess, nicht als Cronjob. Zwei Gruende:
|
||||
|
||||
* Schwellwert-Ausloeser ("Temperatur ueber 22 Grad") sollen sofort greifen,
|
||||
wenn die Nachricht hereinkommt, und nicht bis zum naechsten Minutenraster
|
||||
warten.
|
||||
* `actor_states.current_value` wird sonst von niemandem fortgeschrieben -
|
||||
beim Geraete-Discovery einmal gesetzt und danach nie wieder. Ein
|
||||
zustandsloser Cronjob haette also gar nichts, womit er vergleichen
|
||||
koennte. Der Runner pflegt den Wert nebenbei mit, wodurch auch der Editor
|
||||
im Browser aktuelle Zahlen anzeigt.
|
||||
|
||||
Ablauf:
|
||||
|
||||
Start Regelwerk und Geraetemodell laden, MQTT-Topics abonnieren
|
||||
Ereignis MQTT-Nachricht -> Wert merken -> im naechsten Takt auswerten
|
||||
Takt alle tick_seconds: Uhrzeit/Datum/Sonnenzeiten neu rechnen,
|
||||
HTTP- und Tahoma-Geraete pollen, Regelwerk auf Aenderung pruefen
|
||||
Pruefen aktiv? (enabled, Wochentag, Ferien/Feiertag, Zeitfenster)
|
||||
-> Bedingungen auswerten -> steigende Flanke -> Aktionen
|
||||
|
||||
Nur die steigende Flanke loest aus: `automations.cond_met` haelt fest, ob die
|
||||
Bedingung beim letzten Durchlauf schon erfuellt war. Ohne das wuerde
|
||||
"Temperatur ueber 22 Grad" bei jedem Takt erneut feuern.
|
||||
|
||||
Zeit-Ausloeser gibt es in drei Formen: "um 16:30" ist genau in dieser Minute
|
||||
wahr, "ab 16:30" von da an bis Mitternacht, "vor 16:30" bis dahin.
|
||||
Ausgeloest wird in allen drei Faellen nur einmal, eben wegen der Flanke -
|
||||
"ab" ist aber das robustere: faellt der Runner in der einen Minute aus, auf
|
||||
die "um" zeigt, ist die Automatik fuer den Tag verloren; bei "ab" holt der
|
||||
naechste Takt es nach. Wer "um" braucht und den Ausfall nicht riskieren
|
||||
will, hakt zusaetzlich force_once an. Dieselbe Ueberlegung steht hinter den
|
||||
breiten Zeitfenstern in auto_watering.py.
|
||||
|
||||
Beim Sonnenauf- und -untergang traegt der Operator zusaetzlich das
|
||||
Vorzeichen des Versatzes: "+ 00:30" eine halbe Stunde danach, ">=- 00:30"
|
||||
ab einer halben Stunde davor, "<+ 00:30" bis eine halbe Stunde danach.
|
||||
|
||||
Tabellen siehe homeMesh_automations.sql, Konfiguration siehe config.ini.example.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import configparser
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import queue
|
||||
import sys
|
||||
import threading
|
||||
import time
|
||||
from datetime import date, datetime, timedelta
|
||||
|
||||
import pymysql
|
||||
import requests
|
||||
import paho.mqtt.client as mqtt
|
||||
|
||||
from transports import (HTTPTransport, LogicTransport, MQTTTransport,
|
||||
TahomaTransport, WLEDTransport)
|
||||
|
||||
logger = logging.getLogger("autoaction")
|
||||
|
||||
# Wie lange ein Messwert in der Datenbank stehen bleiben darf, bevor er
|
||||
# aufgefrischt wird. Der Wert dient nur der Anzeige im Editor; jede Nachricht
|
||||
# sofort zu schreiben waere bei einem gespraechigen Sensor sinnlose Last.
|
||||
SCHREIB_ABSTAND = timedelta(seconds=60)
|
||||
|
||||
# Aelteres im Protokoll interessiert niemanden mehr.
|
||||
LOG_AUFBEWAHRUNG_TAGE = 30
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# Konfiguration
|
||||
# ===========================================================================
|
||||
|
||||
class Config:
|
||||
def __init__(self, dateiname="config.ini"):
|
||||
pfad = os.path.join(os.path.dirname(os.path.abspath(__file__)), dateiname)
|
||||
if not os.path.exists(pfad):
|
||||
raise FileNotFoundError(
|
||||
"%s fehlt - config.ini.example kopieren und ausfuellen." % pfad)
|
||||
self.cfg = configparser.ConfigParser()
|
||||
self.cfg.read(pfad, encoding="utf-8")
|
||||
|
||||
def text(self, sektion, schluessel, vorgabe=""):
|
||||
return self.cfg.get(sektion, schluessel, fallback=vorgabe).strip()
|
||||
|
||||
def zahl(self, sektion, schluessel, vorgabe=0):
|
||||
try:
|
||||
return self.cfg.getint(sektion, schluessel, fallback=vorgabe)
|
||||
except ValueError:
|
||||
return vorgabe
|
||||
|
||||
def ja(self, sektion, schluessel, vorgabe=False):
|
||||
return self.text(sektion, schluessel, str(vorgabe)).lower() in ("true", "1", "yes", "on")
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# Datenbank
|
||||
# ===========================================================================
|
||||
|
||||
def verbinden(config, sektion="database"):
|
||||
return pymysql.connect(
|
||||
host=config.text("database", "host", "localhost"),
|
||||
port=config.zahl("database", "port", 3306),
|
||||
user=config.text(sektion, "user") or config.text("database", "user"),
|
||||
password=config.text(sektion, "password") or config.text("database", "password"),
|
||||
database=config.text(sektion, "database"),
|
||||
charset="utf8mb4",
|
||||
cursorclass=pymysql.cursors.DictCursor,
|
||||
autocommit=True)
|
||||
|
||||
|
||||
class Regelwerk:
|
||||
"""
|
||||
Das geladene Abbild der Datenbank: Automatiken mit ihren Bedingungen und
|
||||
Aktionen, dazu die Messwerte und Kommandos, die sie benutzen.
|
||||
"""
|
||||
|
||||
def __init__(self, automatiken, states, kommandos, signatur):
|
||||
self.automatiken = automatiken
|
||||
self.states = states # state_id -> Beschreibung
|
||||
self.kommandos = kommandos # command_id -> Beschreibung
|
||||
self.signatur = signatur
|
||||
|
||||
@staticmethod
|
||||
def signatur_lesen(db):
|
||||
"""
|
||||
Woran der Runner merkt, dass er neu laden muss. `changed` allein
|
||||
genuegt nicht: eine geloeschte Automatik veraendert den groessten
|
||||
Zeitstempel nicht. Deshalb zaehlen die Zeilen mit - auch die der
|
||||
Geraetetabellen, damit ein Discovery-Lauf ebenfalls durchschlaegt.
|
||||
"""
|
||||
with db.cursor() as c:
|
||||
c.execute("""SELECT (SELECT COUNT(*) FROM automations) AS a,
|
||||
(SELECT UNIX_TIMESTAMP(MAX(changed)) FROM automations) AS t,
|
||||
(SELECT COUNT(*) FROM automation_conditions) AS b,
|
||||
(SELECT COUNT(*) FROM automation_actions) AS c,
|
||||
(SELECT COUNT(*) FROM actor_states) AS d,
|
||||
(SELECT COUNT(*) FROM actor_commands) AS e""")
|
||||
return tuple(sorted(c.fetchone().items()))
|
||||
|
||||
@classmethod
|
||||
def laden(cls, db):
|
||||
signatur = cls.signatur_lesen(db)
|
||||
|
||||
states = {}
|
||||
with db.cursor() as c:
|
||||
c.execute("""SELECT s.id, s.state_name, s.url AS state_url, s.value_path,
|
||||
s.current_value, s.possible_values,
|
||||
a.url AS actor_url, a.name AS actor_name, t.type
|
||||
FROM actor_states s
|
||||
JOIN actors a ON a.id = s.actor_id
|
||||
LEFT JOIN state_types t ON s.state_type = t.id""")
|
||||
for row in c.fetchall():
|
||||
row["type"] = row["type"] or "string"
|
||||
# Werttabelle als flaches {gesendeter Wert: Bezeichnung}. In
|
||||
# der Datenbank steht sie als Liste aus Ein-Schluessel-
|
||||
# Objekten, weil der Editor sie so schon versteht.
|
||||
row["wertetabelle"] = {}
|
||||
try:
|
||||
for eintrag in json.loads(row["possible_values"] or "[]"):
|
||||
if isinstance(eintrag, dict):
|
||||
for wert, name in eintrag.items():
|
||||
row["wertetabelle"][str(wert)] = name
|
||||
except ValueError:
|
||||
pass
|
||||
states[row["id"]] = row
|
||||
|
||||
kommandos = {}
|
||||
with db.cursor() as c:
|
||||
c.execute("""SELECT k.id, k.command_name, k.command_url,
|
||||
a.url AS actor_url, a.name AS actor_name
|
||||
FROM actor_commands k JOIN actors a ON a.id = k.actor_id""")
|
||||
for row in c.fetchall():
|
||||
row["params"] = []
|
||||
kommandos[row["id"]] = row
|
||||
with db.cursor() as c:
|
||||
c.execute("""SELECT id, command_id, parameter_name, url FROM command_parameters
|
||||
ORDER BY command_id, id""")
|
||||
for row in c.fetchall():
|
||||
if row["command_id"] in kommandos:
|
||||
kommandos[row["command_id"]]["params"].append(row)
|
||||
|
||||
automatiken = {}
|
||||
with db.cursor() as c:
|
||||
c.execute("SELECT * FROM automations WHERE enabled = 1")
|
||||
for row in c.fetchall():
|
||||
row["gruppen"] = {}
|
||||
row["aktionen"] = []
|
||||
automatiken[row["id"]] = row
|
||||
with db.cursor() as c:
|
||||
c.execute("""SELECT * FROM automation_conditions
|
||||
ORDER BY automation_id, group_no, position, id""")
|
||||
for row in c.fetchall():
|
||||
auto = automatiken.get(row["automation_id"])
|
||||
if auto is None:
|
||||
continue
|
||||
if row["state_id"] not in states:
|
||||
logger.warning("Automatik %s: Messwert %s gibt es nicht mehr",
|
||||
auto["name"], row["state_id"])
|
||||
continue
|
||||
auto["gruppen"].setdefault(row["group_no"], []).append(row)
|
||||
with db.cursor() as c:
|
||||
c.execute("""SELECT a.id, a.automation_id, a.command_id, p.parameter_id, p.value
|
||||
FROM automation_actions a
|
||||
LEFT JOIN automation_action_params p ON p.action_id = a.id
|
||||
ORDER BY a.automation_id, a.position, a.id""")
|
||||
gesammelt = {}
|
||||
for row in c.fetchall():
|
||||
auto = automatiken.get(row["automation_id"])
|
||||
if auto is None:
|
||||
continue
|
||||
aktion = gesammelt.get(row["id"])
|
||||
if aktion is None:
|
||||
aktion = {"command_id": row["command_id"], "werte": {}}
|
||||
gesammelt[row["id"]] = aktion
|
||||
auto["aktionen"].append(aktion)
|
||||
if row["parameter_id"] is not None:
|
||||
aktion["werte"][row["parameter_id"]] = row["value"]
|
||||
|
||||
logger.info("Regelwerk geladen: %d Automatiken, %d Messwerte, %d Kommandos",
|
||||
len(automatiken), len(states), len(kommandos))
|
||||
return cls(automatiken, states, kommandos, signatur)
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# Auswertung
|
||||
# ===========================================================================
|
||||
|
||||
def minuten(text):
|
||||
""""16:30" oder "16:30:00" als Minuten seit Mitternacht."""
|
||||
teile = str(text).strip().split(":")
|
||||
return int(teile[0]) * 60 + int(teile[1])
|
||||
|
||||
|
||||
def als_datum(text):
|
||||
for form in ("%d.%m.%Y", "%Y-%m-%d", "%d.%m.%y"):
|
||||
try:
|
||||
return datetime.strptime(str(text).strip(), form).date()
|
||||
except ValueError:
|
||||
continue
|
||||
return None
|
||||
|
||||
|
||||
def als_zeitpunkt(text):
|
||||
"""Datum mit Uhrzeit. Das Web-Feld liefert "2026-08-30T16:30"."""
|
||||
roh = str(text).strip().replace("T", " ")
|
||||
for form in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%d %H:%M",
|
||||
"%d.%m.%Y %H:%M:%S", "%d.%m.%Y %H:%M"):
|
||||
try:
|
||||
return datetime.strptime(roh, form)
|
||||
except ValueError:
|
||||
continue
|
||||
tag = als_datum(roh.split(" ")[0])
|
||||
return datetime.combine(tag, datetime.min.time()) if tag else None
|
||||
|
||||
|
||||
WAHR = {"true", "1", "on", "ja", "yes", "an"}
|
||||
|
||||
|
||||
def bedingung_erfuellt(bedingung, state, wert, jetzt):
|
||||
"""
|
||||
Ein einzelner Vergleich. `wert` ist der aktuelle Messwert als Text, so wie
|
||||
er vom Geraet kam; `bedingung["value"]` die eingestellte Schwelle.
|
||||
Unbekannter Wert heisst nicht erfuellt - lieber nicht schalten als auf
|
||||
Verdacht schalten.
|
||||
"""
|
||||
if wert is None or wert == "":
|
||||
return False
|
||||
typ = state["type"]
|
||||
op = bedingung["operator"]
|
||||
soll = bedingung["value"]
|
||||
|
||||
try:
|
||||
if typ in ("integer", "float"):
|
||||
ist_z, soll_z = float(str(wert).replace(",", ".")), float(str(soll).replace(",", "."))
|
||||
if op == "=": return ist_z == soll_z
|
||||
if op == "!=": return ist_z != soll_z
|
||||
if op == ">": return ist_z > soll_z
|
||||
if op == "<": return ist_z < soll_z
|
||||
if op == ">=": return ist_z >= soll_z
|
||||
if op == "<=": return ist_z <= soll_z
|
||||
return False
|
||||
|
||||
if typ == "time":
|
||||
ist_m, soll_m = minuten(wert), minuten(soll)
|
||||
if op == "=": return ist_m == soll_m
|
||||
if op == "<": return ist_m < soll_m
|
||||
return ist_m >= soll_m
|
||||
|
||||
if typ == "deltatime":
|
||||
# Der Messwert ist der Sonnenauf- bzw. -untergang, die Schwelle
|
||||
# ein Versatz. Das Vorzeichen steckt im Operator, der Vergleich
|
||||
# davor: ">=-" heisst "ab einer halben Stunde davor".
|
||||
versatz = minuten(soll)
|
||||
ziel = minuten(wert) + (-versatz if op.endswith("-") else versatz)
|
||||
# Ueber den Tagesrand wird gerechnet, nicht abgeschnitten: "sechs
|
||||
# Stunden vor Sonnenaufgang" ist eine gewollte Angabe und landet
|
||||
# dann eben am Vorabend. Abgeschnitten waeren solche Faelle gar
|
||||
# nicht mehr formulierbar.
|
||||
#
|
||||
# Verglichen wird die Uhrzeit innerhalb des Tages. Ein Ziel
|
||||
# jenseits von Mitternacht gilt also als diese Uhrzeit am selben
|
||||
# Tag - bei "um" ist das genau der gemeinte Zeitpunkt, bei "ab"
|
||||
# und "vor" verschiebt sich der wahre Bereich entsprechend.
|
||||
ziel %= 1440
|
||||
jetzt_m = jetzt.hour * 60 + jetzt.minute
|
||||
if op.startswith(">="): return jetzt_m >= ziel
|
||||
if op.startswith("<"): return jetzt_m < ziel
|
||||
return jetzt_m == ziel
|
||||
|
||||
if typ in ("date", "datetime"):
|
||||
wandeln = als_datum if typ == "date" else als_zeitpunkt
|
||||
ist_d, soll_d = wandeln(wert), wandeln(soll)
|
||||
if ist_d is None or soll_d is None:
|
||||
return False
|
||||
if op == "=": return ist_d == soll_d
|
||||
if op == "<": return ist_d < soll_d
|
||||
return ist_d >= soll_d
|
||||
|
||||
if typ == "bool":
|
||||
ist_b = str(wert).strip().lower() in WAHR
|
||||
soll_b = str(soll).strip().lower() in WAHR
|
||||
return ist_b == soll_b if op == "=" else ist_b != soll_b
|
||||
|
||||
# string und alles Uebrige
|
||||
if op == "=":
|
||||
return str(wert).strip() == str(soll).strip()
|
||||
return str(wert).strip() != str(soll).strip()
|
||||
|
||||
except (ValueError, IndexError) as fehler:
|
||||
logger.debug("Vergleich %s %s %s nicht moeglich: %s", wert, op, soll, fehler)
|
||||
return False
|
||||
|
||||
|
||||
def gruppen_erfuellt(automatik, regelwerk, werte, jetzt):
|
||||
"""
|
||||
Gleiche group_no = UND, verschiedene = ODER. Eine Automatik ohne
|
||||
Bedingungen loest nie aus - sonst wuerde sie nach einem Discovery-Lauf,
|
||||
der ihren Messwert entfernt hat, ploetzlich dauernd feuern.
|
||||
"""
|
||||
if not automatik["gruppen"]:
|
||||
return False
|
||||
for bedingungen in automatik["gruppen"].values():
|
||||
if all(bedingung_erfuellt(b, regelwerk.states[b["state_id"]],
|
||||
werte.get(b["state_id"]), jetzt)
|
||||
for b in bedingungen):
|
||||
return True
|
||||
return False
|
||||
|
||||
|
||||
def im_zeitfenster(jetzt, von, bis):
|
||||
"""von > bis heisst: das Fenster reicht ueber Mitternacht."""
|
||||
m = jetzt.hour * 60 + jetzt.minute
|
||||
a, b = minuten(von), minuten(bis)
|
||||
return a <= m <= b if a <= b else (m >= a or m <= b)
|
||||
|
||||
|
||||
def tag_passt(automatik, jetzt, kalender):
|
||||
if not automatik["weekdays"] & (1 << jetzt.weekday()):
|
||||
return False
|
||||
if kalender["feiertag"] and not automatik["on_holiday"]:
|
||||
return False
|
||||
if kalender["ferien"] and not automatik["on_vacation"]:
|
||||
return False
|
||||
return True
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# Versand
|
||||
# ===========================================================================
|
||||
|
||||
class Versand:
|
||||
"""
|
||||
Schickt Kommandos im Hintergrund - je Geraet der Reihe nach, ueber
|
||||
Geraete hinweg nebeneinander.
|
||||
|
||||
Ohne das blockiert ein einziges Kommando die ganze Auswertung: eine
|
||||
Jalousie zuzufahren dauert ueber eine Minute, weil zwischen den beiden
|
||||
Schwenkbefehlen auf das Ende der Fahrt gewartet werden muss (siehe
|
||||
TahomaTransport). So lange kaeme kein Zeit-Ausloeser mehr durch, kein
|
||||
Messwert wuerde zurueckgeschrieben, und mehrere Rollladen in einer
|
||||
Automatik wuerden sich aufaddieren.
|
||||
|
||||
Je Geraet eine Warteschlange mit einem eigenen Faden: zwei Kommandos an
|
||||
dieselbe Jalousie duerfen sich nicht ueberholen - "Neigung 0" und
|
||||
"Neigung 100" sind sonst wirkungslos oder vertauscht -, zwei Kommandos an
|
||||
verschiedene Jalousien duerfen ruhig gleichzeitig laufen.
|
||||
|
||||
Die Datenbank bleibt aussen vor: die Faeden melden ihr Ergebnis nur
|
||||
zurueck, geschrieben wird im Hauptfaden. Eine pymysql-Verbindung ist
|
||||
nicht fuer mehrere Faeden gedacht.
|
||||
"""
|
||||
|
||||
def __init__(self):
|
||||
self.warteschlangen = {} # actor_url -> Queue
|
||||
self.ergebnisse = queue.Queue()
|
||||
self.laeuft = True
|
||||
|
||||
def einreihen(self, automation_id, beschreibung, transport, auftrag):
|
||||
schlange = self.warteschlangen.get(auftrag["actor_url"])
|
||||
if schlange is None:
|
||||
schlange = queue.Queue()
|
||||
self.warteschlangen[auftrag["actor_url"]] = schlange
|
||||
faden = threading.Thread(target=self._arbeiten, args=(schlange,),
|
||||
name="versand", daemon=True)
|
||||
faden.start()
|
||||
schlange.put((automation_id, beschreibung, transport, auftrag))
|
||||
|
||||
def _arbeiten(self, schlange):
|
||||
while self.laeuft:
|
||||
posten = schlange.get()
|
||||
if posten is None:
|
||||
return
|
||||
automation_id, beschreibung, transport, auftrag = posten
|
||||
try:
|
||||
transport.senden(auftrag)
|
||||
self.ergebnisse.put((automation_id, beschreibung, None))
|
||||
except Exception as fehler:
|
||||
self.ergebnisse.put((automation_id, beschreibung, str(fehler)))
|
||||
finally:
|
||||
schlange.task_done()
|
||||
|
||||
def abholen(self):
|
||||
"""Alles, was seit dem letzten Mal fertig geworden ist."""
|
||||
fertig = []
|
||||
while True:
|
||||
try:
|
||||
fertig.append(self.ergebnisse.get_nowait())
|
||||
except queue.Empty:
|
||||
return fertig
|
||||
|
||||
def offen(self):
|
||||
return sum(s.unfinished_tasks for s in self.warteschlangen.values())
|
||||
|
||||
def beenden(self):
|
||||
self.laeuft = False
|
||||
for schlange in self.warteschlangen.values():
|
||||
schlange.put(None)
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# Der Runner
|
||||
# ===========================================================================
|
||||
|
||||
class Runner:
|
||||
|
||||
def __init__(self, config, dry_run=False):
|
||||
self.config = config
|
||||
self.dry_run = dry_run or config.ja("runner", "dry_run")
|
||||
self.db = verbinden(config)
|
||||
self.solar = verbinden(config, "solar")
|
||||
|
||||
self.werte = {} # state_id -> letzter bekannter Wert
|
||||
self.geschrieben = {} # state_id -> wann zuletzt in die DB
|
||||
self.war_aktiv = {} # automation_id -> war im Zeitfenster
|
||||
self.lief_im_fenster = {} # automation_id -> hat im Fenster ausgeloest
|
||||
self._sonne = (None, "00:00", "00:00") # (datum, aufgang, untergang)
|
||||
self._kalender = (None, {"feiertag": False, "ferien": False})
|
||||
self._letzte_saeuberung = None
|
||||
|
||||
self.versand = Versand()
|
||||
self.mqtt = self._mqtt_verbinden()
|
||||
self.transporte = [
|
||||
MQTTTransport(self.mqtt, self.dry_run),
|
||||
WLEDTransport(requests, dry_run=self.dry_run),
|
||||
HTTPTransport(requests, dry_run=self.dry_run),
|
||||
TahomaTransport(requests, config.text("tahoma", "pin"),
|
||||
config.text("tahoma", "token"),
|
||||
config.zahl("tahoma", "timeout", 10), self.dry_run),
|
||||
LogicTransport(self.sonnenzeiten),
|
||||
]
|
||||
self.regelwerk = None
|
||||
|
||||
# --- Aufbau ----------------------------------------------------------
|
||||
|
||||
def _mqtt_verbinden(self):
|
||||
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2,
|
||||
client_id=self.config.text("mqtt", "client_id", "autoaction_runner"))
|
||||
benutzer = self.config.text("mqtt", "username")
|
||||
if benutzer:
|
||||
client.username_pw_set(benutzer, self.config.text("mqtt", "password"))
|
||||
client.on_message = self._mqtt_nachricht
|
||||
client.connect(self.config.text("mqtt", "broker", "localhost"),
|
||||
self.config.zahl("mqtt", "port", 1883), 60)
|
||||
client.loop_start()
|
||||
return client
|
||||
|
||||
def _mqtt_nachricht(self, client, userdata, nachricht):
|
||||
# Der Client laeuft schon, waehrend __init__ noch die Transporte baut.
|
||||
for transport in getattr(self, "transporte", []):
|
||||
if isinstance(transport, MQTTTransport):
|
||||
transport.nachricht(nachricht.topic, nachricht.payload)
|
||||
|
||||
def transport_fuer(self, actor_url):
|
||||
for transport in self.transporte:
|
||||
if transport.passt(actor_url):
|
||||
return transport
|
||||
return None
|
||||
|
||||
def regelwerk_laden(self):
|
||||
self.regelwerk = Regelwerk.laden(self.db)
|
||||
|
||||
# Welche Geraete Position und Neigung zusammen koennen. Nur die haben
|
||||
# die Kugelschreiber-Mechanik, und nur bei ihnen wird ein "Zu" zum
|
||||
# zweifachen Schwenken - siehe TahomaTransport.
|
||||
kombi = {k["actor_url"] for k in self.regelwerk.kommandos.values()
|
||||
if k["command_url"] == "setClosureAndOrientation"}
|
||||
for transport in self.transporte:
|
||||
if isinstance(transport, TahomaTransport):
|
||||
transport.kombigeraete_setzen(kombi)
|
||||
|
||||
# Jeder Transport bekommt die Messwerte, fuer die er zustaendig ist.
|
||||
for transport in self.transporte:
|
||||
passende = [{"id": s["id"], "actor_url": s["actor_url"],
|
||||
"state_url": s["state_url"], "value_path": s["value_path"],
|
||||
"wertetabelle": s["wertetabelle"]}
|
||||
for s in self.regelwerk.states.values()
|
||||
if transport.passt(s["actor_url"])]
|
||||
transport.zustaende_anmelden(passende)
|
||||
# Der zuletzt bekannte Wert aus der Datenbank ist besser als gar
|
||||
# keiner: nach einem Neustart steht sonst jede Bedingung auf "unklar",
|
||||
# bis das Geraet zufaellig etwas schickt.
|
||||
for s in self.regelwerk.states.values():
|
||||
if s["id"] not in self.werte and s["current_value"] is not None:
|
||||
self.werte[s["id"]] = s["current_value"]
|
||||
|
||||
# --- Umgebung --------------------------------------------------------
|
||||
|
||||
def sonnenzeiten(self):
|
||||
"""Aus solarLog.daylight, einmal je Tag geholt."""
|
||||
heute = date.today()
|
||||
if self._sonne[0] == heute:
|
||||
return self._sonne[1], self._sonne[2]
|
||||
auf, unter = "00:00", "00:00"
|
||||
try:
|
||||
with self.solar.cursor() as c:
|
||||
c.execute("SELECT sunrise, sunset FROM daylight WHERE date = %s", (heute,))
|
||||
zeile = c.fetchone()
|
||||
if zeile:
|
||||
auf = self._als_uhrzeit(zeile["sunrise"])
|
||||
unter = self._als_uhrzeit(zeile["sunset"])
|
||||
else:
|
||||
logger.warning("Kein Eintrag in daylight fuer %s", heute)
|
||||
except Exception as fehler:
|
||||
logger.warning("Sonnenzeiten nicht lesbar: %s", fehler)
|
||||
self._sonne = (heute, auf, unter)
|
||||
return auf, unter
|
||||
|
||||
@staticmethod
|
||||
def _als_uhrzeit(wert):
|
||||
"""TIME-Spalten liefert pymysql als timedelta, nicht als Text."""
|
||||
if isinstance(wert, timedelta):
|
||||
minute = int(wert.total_seconds()) // 60
|
||||
return "%02d:%02d" % (minute // 60 % 24, minute % 60)
|
||||
return str(wert)[:5]
|
||||
|
||||
def kalender(self):
|
||||
"""Ferien und Feiertage von heute, einmal je Tag geholt."""
|
||||
heute = date.today()
|
||||
if self._kalender[0] == heute:
|
||||
return self._kalender[1]
|
||||
stand = {"feiertag": False, "ferien": False}
|
||||
try:
|
||||
with self.db.cursor() as c:
|
||||
c.execute("SELECT holiday, vacation FROM calendar_days WHERE date = %s", (heute,))
|
||||
zeile = c.fetchone()
|
||||
if zeile:
|
||||
stand = {"feiertag": bool(zeile["holiday"]), "ferien": bool(zeile["vacation"])}
|
||||
except Exception as fehler:
|
||||
logger.warning("Kalender nicht lesbar: %s", fehler)
|
||||
self._kalender = (heute, stand)
|
||||
return stand
|
||||
|
||||
# --- Werte -----------------------------------------------------------
|
||||
|
||||
def werte_einsammeln(self, mit_pollen):
|
||||
neu = {}
|
||||
for transport in self.transporte:
|
||||
if isinstance(transport, MQTTTransport) or mit_pollen:
|
||||
neu.update(transport.zustaende_lesen())
|
||||
if neu:
|
||||
self.werte.update(neu)
|
||||
self.werte_zurueckschreiben(neu)
|
||||
return neu
|
||||
|
||||
def werte_zurueckschreiben(self, neu):
|
||||
"""
|
||||
current_value nachfuehren, damit der Editor im Browser aktuelle Zahlen
|
||||
zeigt ("Temperatur (= 21,4 °C)"). Gedrosselt, sonst schreibt ein
|
||||
gespraechiger Sensor die Tabelle im Sekundentakt voll.
|
||||
"""
|
||||
jetzt = datetime.now()
|
||||
faellig = [(w, i) for i, w in neu.items()
|
||||
if self.geschrieben.get(i, datetime.min) + SCHREIB_ABSTAND <= jetzt]
|
||||
if not faellig:
|
||||
return
|
||||
try:
|
||||
with self.db.cursor() as c:
|
||||
c.executemany("UPDATE actor_states SET current_value = %s WHERE id = %s", faellig)
|
||||
for _, state_id in faellig:
|
||||
self.geschrieben[state_id] = jetzt
|
||||
except Exception as fehler:
|
||||
logger.warning("current_value nicht schreibbar: %s", fehler)
|
||||
|
||||
# --- Ausfuehren ------------------------------------------------------
|
||||
|
||||
def ausloesen(self, automatik, anlass):
|
||||
"""
|
||||
Die Aktionen einer Automatik in den Versand geben.
|
||||
|
||||
Geschickt wird im Hintergrund - eine Jalousie zuzufahren dauert ueber
|
||||
eine Minute, und so lange darf die Auswertung nicht stehen. Was dabei
|
||||
schiefgeht, kommt spaeter ueber ergebnisse_verbuchen() ins Protokoll.
|
||||
"""
|
||||
fehlerText = []
|
||||
eingereiht = 0
|
||||
for aktion in automatik["aktionen"]:
|
||||
kommando = self.regelwerk.kommandos.get(aktion["command_id"])
|
||||
if kommando is None:
|
||||
fehlerText.append("Kommando %s gibt es nicht mehr" % aktion["command_id"])
|
||||
logger.error("%s: Kommando %s gibt es nicht mehr",
|
||||
automatik["name"], aktion["command_id"])
|
||||
continue
|
||||
transport = self.transport_fuer(kommando["actor_url"])
|
||||
if transport is None:
|
||||
fehlerText.append("Kein Transport fuer %s" % kommando["actor_url"])
|
||||
logger.error("%s: fuer %s (%s) gibt es keinen Transport",
|
||||
automatik["name"], kommando["actor_name"], kommando["actor_url"])
|
||||
continue
|
||||
auftrag = {
|
||||
"actor_url": kommando["actor_url"],
|
||||
"command_url": kommando["command_url"],
|
||||
"params": [{"url": p["url"], "name": p["parameter_name"],
|
||||
"wert": aktion["werte"].get(p["id"], "")}
|
||||
for p in kommando["params"]],
|
||||
}
|
||||
beschreibung = "%s: %s" % (kommando["actor_name"], kommando["command_name"])
|
||||
logger.info("%s: %s (unterwegs)", automatik["name"], beschreibung)
|
||||
self.versand.einreihen(automatik["id"], beschreibung, transport, auftrag)
|
||||
eingereiht += 1
|
||||
|
||||
ergebnis = "error" if fehlerText else anlass
|
||||
detail = "; ".join(fehlerText)[:255] if fehlerText \
|
||||
else ("%d Kommando(s) unterwegs" % eingereiht)
|
||||
try:
|
||||
with self.db.cursor() as c:
|
||||
c.execute("""UPDATE automations SET last_run = NOW(), changed = changed
|
||||
WHERE id = %s""", (automatik["id"],))
|
||||
c.execute("""INSERT INTO automation_log (automation_id, result, detail)
|
||||
VALUES (%s, %s, %s)""", (automatik["id"], ergebnis, detail))
|
||||
except Exception as fehler:
|
||||
logger.warning("Protokoll nicht schreibbar: %s", fehler)
|
||||
automatik["last_run"] = datetime.now()
|
||||
self.lief_im_fenster[automatik["id"]] = True
|
||||
|
||||
def ergebnisse_verbuchen(self):
|
||||
"""
|
||||
Was der Versand inzwischen erledigt hat ins Protokoll schreiben.
|
||||
|
||||
Nur Fehlschlaege bekommen eine eigene Zeile - der Lauf selbst steht
|
||||
schon drin, und ein Protokoll, das jedes gelungene Kommando einzeln
|
||||
auffuehrt, findet niemand mehr etwas darin.
|
||||
"""
|
||||
for automation_id, beschreibung, fehler in self.versand.abholen():
|
||||
if fehler is None:
|
||||
logger.debug("erledigt: %s", beschreibung)
|
||||
continue
|
||||
logger.error("%s konnte nicht geschickt werden: %s", beschreibung, fehler)
|
||||
try:
|
||||
with self.db.cursor() as c:
|
||||
c.execute("""INSERT INTO automation_log (automation_id, result, detail)
|
||||
VALUES (%s, 'error', %s)""",
|
||||
(automation_id, ("%s: %s" % (beschreibung, fehler))[:255]))
|
||||
except Exception as schreibfehler:
|
||||
logger.warning("Protokoll nicht schreibbar: %s", schreibfehler)
|
||||
|
||||
def flanke_merken(self, automatik, erfuellt):
|
||||
if bool(automatik["cond_met"]) == bool(erfuellt):
|
||||
return
|
||||
automatik["cond_met"] = 1 if erfuellt else 0
|
||||
try:
|
||||
with self.db.cursor() as c:
|
||||
c.execute("""UPDATE automations SET cond_met = %s, changed = changed
|
||||
WHERE id = %s""", (automatik["cond_met"], automatik["id"]))
|
||||
except Exception as fehler:
|
||||
logger.warning("cond_met nicht schreibbar: %s", fehler)
|
||||
|
||||
def durchlauf(self):
|
||||
jetzt = datetime.now()
|
||||
kalender = self.kalender()
|
||||
|
||||
for automatik in self.regelwerk.automatiken.values():
|
||||
aktiv = (tag_passt(automatik, jetzt, kalender)
|
||||
and im_zeitfenster(jetzt, automatik["window_from"], automatik["window_to"]))
|
||||
vorher_aktiv = self.war_aktiv.get(automatik["id"], aktiv)
|
||||
|
||||
if aktiv:
|
||||
if not vorher_aktiv:
|
||||
self.lief_im_fenster[automatik["id"]] = False
|
||||
erfuellt = gruppen_erfuellt(automatik, self.regelwerk, self.werte, jetzt)
|
||||
if erfuellt and not automatik["cond_met"]:
|
||||
if self.gesperrt(automatik, jetzt):
|
||||
logger.debug("%s: Flanke faellt in die Sperrzeit, uebersprungen",
|
||||
automatik["name"])
|
||||
else:
|
||||
self.ausloesen(automatik, "fired")
|
||||
self.flanke_merken(automatik, erfuellt)
|
||||
else:
|
||||
# Das Fenster ist gerade zugegangen. Wer "auf jeden Fall"
|
||||
# angehakt hat, bekommt jetzt seinen Lauf - aber nur, wenn in
|
||||
# diesem Fenster noch keiner stattgefunden hat. Nach einem
|
||||
# Neustart mitten im Fenster weiss der Runner das nicht mehr
|
||||
# aus dem Speicher, deshalb zaehlt zusaetzlich last_run.
|
||||
if vorher_aktiv and automatik["force_once"] and not self.lief_im_fenster.get(automatik["id"]):
|
||||
if not self.lief_heute(automatik, jetzt):
|
||||
logger.info("%s: Zeitfenster vorbei, wird trotzdem ausgefuehrt",
|
||||
automatik["name"])
|
||||
self.ausloesen(automatik, "forced")
|
||||
# Ausserhalb des Fensters die Flanke zuruecksetzen, sonst
|
||||
# koennte sie im naechsten Fenster nicht mehr steigen.
|
||||
self.flanke_merken(automatik, False)
|
||||
self.lief_im_fenster[automatik["id"]] = False
|
||||
|
||||
self.war_aktiv[automatik["id"]] = aktiv
|
||||
|
||||
@staticmethod
|
||||
def gesperrt(automatik, jetzt):
|
||||
"""
|
||||
Liegt die letzte Ausloesung noch innerhalb der Sperrzeit?
|
||||
|
||||
Gegen Messwerte, die um die Schwelle pendeln: "Temperatur > 22" bei
|
||||
22,1 / 21,9 / 22,1 Grad ist jedes Mal eine echte steigende Flanke, und
|
||||
ueber MQTT koennen die Werte im Sekundentakt hereinkommen.
|
||||
|
||||
Die Flanke wird dabei verworfen und nicht aufgehoben. Ein Rollladen,
|
||||
der eine Viertelstunde spaeter doch noch losfaehrt, weil vor langer
|
||||
Zeit einmal eine Schwelle gestreift wurde, waere unangenehmer als
|
||||
einer, der gar nicht faehrt. Der naechste echte Anlass nach Ablauf
|
||||
der Sperre kommt ohnehin durch.
|
||||
|
||||
force_once ist davon nicht betroffen: es greift nur, wenn im Fenster
|
||||
gar nichts gelaufen ist - dann ist auch keine Sperre aktiv.
|
||||
"""
|
||||
sperre = int(automatik.get("lockout_secs") or 0)
|
||||
letzter = automatik.get("last_run")
|
||||
if not sperre or not letzter:
|
||||
return False
|
||||
return (jetzt - letzter).total_seconds() < sperre
|
||||
|
||||
@staticmethod
|
||||
def lief_heute(automatik, jetzt):
|
||||
letzter = automatik.get("last_run")
|
||||
return bool(letzter) and letzter.date() == jetzt.date()
|
||||
|
||||
def protokoll_saeubern(self):
|
||||
heute = date.today()
|
||||
if self._letzte_saeuberung == heute:
|
||||
return
|
||||
self._letzte_saeuberung = heute
|
||||
try:
|
||||
with self.db.cursor() as c:
|
||||
c.execute("DELETE FROM automation_log WHERE ts < NOW() - INTERVAL %s DAY",
|
||||
(LOG_AUFBEWAHRUNG_TAGE,))
|
||||
except Exception as fehler:
|
||||
logger.warning("Protokoll nicht aufraeumbar: %s", fehler)
|
||||
|
||||
# --- Hauptschleife ---------------------------------------------------
|
||||
|
||||
def laufen(self, nur_einmal=False):
|
||||
self.regelwerk_laden()
|
||||
takt = self.config.zahl("runner", "tick_seconds", 30)
|
||||
poll = self.config.zahl("runner", "poll_seconds", 60)
|
||||
neuladen = self.config.zahl("runner", "reload_seconds", 30)
|
||||
letztes_pollen = 0.0
|
||||
letztes_pruefen = 0.0
|
||||
|
||||
while True:
|
||||
jetzt = time.monotonic()
|
||||
mit_pollen = jetzt - letztes_pollen >= poll
|
||||
if mit_pollen:
|
||||
letztes_pollen = jetzt
|
||||
|
||||
self.werte_einsammeln(mit_pollen)
|
||||
self.durchlauf()
|
||||
self.ergebnisse_verbuchen()
|
||||
self.protokoll_saeubern()
|
||||
|
||||
if jetzt - letztes_pruefen >= neuladen:
|
||||
letztes_pruefen = jetzt
|
||||
try:
|
||||
if Regelwerk.signatur_lesen(self.db) != self.regelwerk.signatur:
|
||||
logger.info("Regelwerk hat sich geaendert, wird neu geladen")
|
||||
self.regelwerk_laden()
|
||||
except Exception as fehler:
|
||||
logger.warning("Regelwerk nicht pruefbar: %s", fehler)
|
||||
|
||||
if nur_einmal:
|
||||
return
|
||||
time.sleep(takt)
|
||||
|
||||
def beenden(self):
|
||||
offen = self.versand.offen()
|
||||
if offen:
|
||||
logger.info("%d Kommando(s) noch unterwegs - wird nicht abgewartet", offen)
|
||||
self.versand.beenden()
|
||||
try:
|
||||
self.mqtt.loop_stop()
|
||||
self.mqtt.disconnect()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description="Fuehrt die Automatiken aus dem Web-UI aus.")
|
||||
parser.add_argument("--dry-run", action="store_true",
|
||||
help="nichts wirklich schalten, nur protokollieren")
|
||||
parser.add_argument("--once", action="store_true",
|
||||
help="einen einzigen Durchlauf, dann beenden")
|
||||
parser.add_argument("--verbose", action="store_true", help="DEBUG-Ausgaben")
|
||||
parser.add_argument("--config", default="config.ini")
|
||||
args = parser.parse_args()
|
||||
|
||||
config = Config(args.config)
|
||||
logging.basicConfig(
|
||||
level=logging.DEBUG if args.verbose else getattr(
|
||||
logging, config.text("runner", "log_level", "INFO").upper(), logging.INFO),
|
||||
format="%(asctime)s - %(levelname)s - %(message)s")
|
||||
|
||||
runner = Runner(config, args.dry_run)
|
||||
if runner.dry_run:
|
||||
logger.info("Probelauf: es wird nichts geschaltet.")
|
||||
try:
|
||||
runner.laufen(args.once)
|
||||
except KeyboardInterrupt:
|
||||
logger.info("Abbruch durch Benutzer")
|
||||
finally:
|
||||
runner.beenden()
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -1,75 +0,0 @@
|
||||
# ============================================================================
|
||||
# AutoAction-Runner - Konfiguration
|
||||
# ============================================================================
|
||||
# Kopieren nach config.ini und ausfuellen. config.ini ist nicht versioniert
|
||||
# (siehe .gitignore), weil hier Zugangsdaten stehen.
|
||||
|
||||
# ============================================================================
|
||||
# DATENBANK - Geraete und Automatiken (homeMesh)
|
||||
# ============================================================================
|
||||
[database]
|
||||
host = nas.fritz.box
|
||||
port = 3310
|
||||
database = homeMesh
|
||||
user = homeMesh
|
||||
password =
|
||||
|
||||
# ============================================================================
|
||||
# DATENBANK - Messwerte (solarLog)
|
||||
# ============================================================================
|
||||
# Nur fuer die Tabelle `daylight`, aus der Sonnenauf- und -untergang kommen.
|
||||
# Leer lassen, wenn dieselben Zugangsdaten wie oben gelten.
|
||||
[solar]
|
||||
database = solarLog
|
||||
user =
|
||||
password =
|
||||
|
||||
# ============================================================================
|
||||
# MQTT
|
||||
# ============================================================================
|
||||
[mqtt]
|
||||
broker = nas.fritz.box
|
||||
port = 1883
|
||||
username =
|
||||
password =
|
||||
client_id = autoaction_runner
|
||||
|
||||
# ============================================================================
|
||||
# TAHOMA
|
||||
# ============================================================================
|
||||
# Fuer Aktoren, deren URL mit io:// beginnt. Ohne Token bleiben sie stumm -
|
||||
# der Runner protokolliert das dann als Fehler, statt still nichts zu tun.
|
||||
[tahoma]
|
||||
pin =
|
||||
token =
|
||||
timeout = 10
|
||||
|
||||
# ============================================================================
|
||||
# KALENDER
|
||||
# ============================================================================
|
||||
# Fuer fetch_calendar.py, das Feiertage und Schulferien in calendar_days
|
||||
# eintraegt. Regionscodes siehe openholidaysapi.org, Bayern ist DE-BY.
|
||||
[calendar]
|
||||
country = DE
|
||||
subdivision = DE-BY
|
||||
|
||||
# ============================================================================
|
||||
# LAUFZEIT
|
||||
# ============================================================================
|
||||
[runner]
|
||||
# Abstand der Auswertung in Sekunden. MQTT-Nachrichten werden sofort
|
||||
# verarbeitet; der Takt ist fuer die Zeit-Ausloeser und fuer das Pollen.
|
||||
tick_seconds = 30
|
||||
|
||||
# Wie oft Geraete abgefragt werden, die nichts von sich aus melden
|
||||
# (Shelly ueber HTTP, Tahoma). In Sekunden.
|
||||
poll_seconds = 60
|
||||
|
||||
# Wie oft der Runner nachsieht, ob sich das Regelwerk geaendert hat.
|
||||
reload_seconds = 30
|
||||
|
||||
# Nichts wirklich schalten, nur protokollieren. Zum Ausprobieren neuer Regeln.
|
||||
dry_run = false
|
||||
|
||||
# DEBUG, INFO, WARNING, ERROR
|
||||
log_level = INFO
|
||||
@@ -1,102 +0,0 @@
|
||||
#!/usr/bin/env python3
|
||||
"""
|
||||
Fuellt die Tabelle `calendar_days` mit Feiertagen und Schulferien.
|
||||
|
||||
Die Automatiken koennen mit "an Feiertagen ausfuehren" und "in den Ferien
|
||||
ausfuehren" auf besondere Tage reagieren; woher das Wissen kommt, steht hier.
|
||||
Nur besondere Tage werden eingetragen - ein fehlendes Datum ist ein
|
||||
gewoehnlicher Tag.
|
||||
|
||||
Quelle ist openholidaysapi.org: ein offenes Verzeichnis der EU-Kommission, das
|
||||
gesetzliche Feiertage und Schulferien fuer alle deutschen Bundeslaender
|
||||
liefert. Ein Aufruf je Jahr genuegt, deshalb ist das ein Cronjob und kein
|
||||
Dauerlaeufer:
|
||||
|
||||
0 4 1 1 * /usr/bin/python3 /pfad/zu/fetch_calendar.py
|
||||
|
||||
Ein zusaetzlicher Lauf im Herbst schadet nicht - die Ferientermine des
|
||||
uebernaechsten Schuljahres stehen erst spaeter fest.
|
||||
|
||||
Bereits eingetragene Tage werden ueberschrieben, nie doppelt angelegt.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import logging
|
||||
import sys
|
||||
from datetime import date, timedelta
|
||||
|
||||
import pymysql
|
||||
import requests
|
||||
|
||||
from autoaction_runner import Config, verbinden
|
||||
|
||||
logger = logging.getLogger("kalender")
|
||||
|
||||
BASIS = "https://openholidaysapi.org"
|
||||
|
||||
|
||||
def zeitraum(api, land, region, von, bis):
|
||||
"""Eintraege einer Art (PublicHolidays oder SchoolHolidays) abholen."""
|
||||
antwort = requests.get(
|
||||
BASIS + "/" + api,
|
||||
params={"countryIsoCode": land, "subdivisionCode": region,
|
||||
"languageIsoCode": "DE", "validFrom": von.isoformat(),
|
||||
"validTo": bis.isoformat()},
|
||||
headers={"Accept": "application/json"}, timeout=20)
|
||||
antwort.raise_for_status()
|
||||
return antwort.json()
|
||||
|
||||
|
||||
def name(eintrag):
|
||||
for teil in eintrag.get("name", []):
|
||||
if teil.get("text"):
|
||||
return teil["text"][:80]
|
||||
return "?"
|
||||
|
||||
|
||||
def tage(eintrag):
|
||||
"""Ferien erstrecken sich ueber Wochen; hier wird daraus Tag fuer Tag."""
|
||||
start = date.fromisoformat(eintrag["startDate"])
|
||||
ende = date.fromisoformat(eintrag["endDate"])
|
||||
while start <= ende:
|
||||
yield start
|
||||
start += timedelta(days=1)
|
||||
|
||||
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description="Feiertage und Ferien in calendar_days schreiben.")
|
||||
parser.add_argument("--jahr", type=int, default=date.today().year)
|
||||
parser.add_argument("--jahre", type=int, default=2, help="wie viele Jahre ab --jahr")
|
||||
parser.add_argument("--config", default="config.ini")
|
||||
args = parser.parse_args()
|
||||
|
||||
logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(levelname)s - %(message)s")
|
||||
config = Config(args.config)
|
||||
land = config.text("calendar", "country", "DE")
|
||||
region = config.text("calendar", "subdivision", "DE-BY")
|
||||
|
||||
von = date(args.jahr, 1, 1)
|
||||
bis = date(args.jahr + args.jahre - 1, 12, 31)
|
||||
logger.info("Hole %s bis %s fuer %s", von, bis, region)
|
||||
|
||||
gesammelt = {}
|
||||
for eintrag in zeitraum("PublicHolidays", land, region, von, bis):
|
||||
for tag in tage(eintrag):
|
||||
gesammelt.setdefault(tag, {})["holiday"] = name(eintrag)
|
||||
for eintrag in zeitraum("SchoolHolidays", land, region, von, bis):
|
||||
for tag in tage(eintrag):
|
||||
gesammelt.setdefault(tag, {})["vacation"] = name(eintrag)
|
||||
|
||||
db = verbinden(config)
|
||||
with db.cursor() as c:
|
||||
c.executemany(
|
||||
"""INSERT INTO calendar_days (date, holiday, vacation) VALUES (%s, %s, %s)
|
||||
ON DUPLICATE KEY UPDATE holiday = VALUES(holiday), vacation = VALUES(vacation)""",
|
||||
[(tag, eintrag.get("holiday"), eintrag.get("vacation"))
|
||||
for tag, eintrag in sorted(gesammelt.items())])
|
||||
logger.info("%d Tage eingetragen", len(gesammelt))
|
||||
return 0
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
sys.exit(main())
|
||||
@@ -1,626 +0,0 @@
|
||||
#!/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, value_path und
|
||||
id, `aktion` ein Dict mit actor_url, command_url und params (Liste aus
|
||||
{url, name, wert}).
|
||||
"""
|
||||
|
||||
import json
|
||||
import logging
|
||||
import re
|
||||
import time
|
||||
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)
|
||||
|
||||
|
||||
def wert_aus_pfad(daten, pfad):
|
||||
"""
|
||||
Einen Teilwert aus einer Nutzlast holen: "[3]", "ssid", "seg[0].col[0]".
|
||||
Ein leerer Pfad heisst: die Nutzlast selbst.
|
||||
|
||||
Gebraucht wird das an zwei Stellen. MQTT-Geraete legen mehrere Messwerte
|
||||
auf ein Topic - der go-eCharger schickt sechzehn Zahlen als JSON-Feld -,
|
||||
und WLED liefert seinen gesamten Zustand als ein Dokument.
|
||||
"""
|
||||
if not pfad:
|
||||
return daten
|
||||
for teil in pfad.split("."):
|
||||
treffer = re.match(r"^([^\[]*)((?:\[\d+\])*)$", teil)
|
||||
if not treffer:
|
||||
raise KeyError(pfad)
|
||||
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 uebersetze(wert, tabelle):
|
||||
"""
|
||||
Aus der gesendeten Zahl den Zustandsnamen machen: aus "2" wird "Charging".
|
||||
|
||||
Manche Geraete schicken einen Zahlencode und meinen einen Zustand. Welche
|
||||
Zahl welchen Namen hat, steht in possible_values - derselben Spalte, aus
|
||||
der auch der Editor seine Auswahlliste baut. Eine Bedingung vergleicht
|
||||
damit genau den Klartext, den man dort ausgewaehlt hat.
|
||||
|
||||
Steht die Zahl nicht in der Tabelle, bleibt sie stehen: ein erfundener
|
||||
Name waere schlimmer als ein roher Wert.
|
||||
"""
|
||||
if not tabelle:
|
||||
return wert
|
||||
if wert in tabelle:
|
||||
return tabelle[wert]
|
||||
try: # "2.0" und "2" meinen dieselbe Stufe
|
||||
ganz = str(int(float(wert)))
|
||||
except (TypeError, ValueError):
|
||||
return wert
|
||||
return tabelle.get(ganz, wert)
|
||||
|
||||
|
||||
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, value_path, wertetabelle), ...]
|
||||
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"], s.get("value_path"), s.get("wertetabelle") or {}))
|
||||
for topic in self.topics:
|
||||
self.client.subscribe(topic)
|
||||
mehrfach = sum(1 for e in self.topics.values() if len(e) > 1)
|
||||
logger.info("MQTT: %d Topics abonniert, %d davon mit mehreren Messwerten",
|
||||
len(self.topics), mehrfach)
|
||||
|
||||
def nachricht(self, topic, payload):
|
||||
"""Wird vom Runner aus dem on_message-Rueckruf gerufen."""
|
||||
eintraege = self.topics.get(topic, [])
|
||||
if not eintraege:
|
||||
return
|
||||
text = payload.decode("utf-8", "replace").strip() if isinstance(payload, bytes) else str(payload)
|
||||
try:
|
||||
daten = json.loads(text)
|
||||
except ValueError:
|
||||
daten = None # kein JSON - dann gilt der Rohtext
|
||||
for state_id, pfad, tabelle in eintraege:
|
||||
wert = self._wert(text, daten, pfad, tabelle)
|
||||
if wert is not None:
|
||||
self.neu[state_id] = wert
|
||||
|
||||
@staticmethod
|
||||
def _wert(text, daten, pfad, tabelle=None):
|
||||
"""
|
||||
Aus der Nutzlast den Wert eines einzelnen Messwerts machen.
|
||||
|
||||
Mehrere Messwerte teilen sich oft ein Topic; welcher Teil gemeint ist,
|
||||
steht in value_path - gelesen aus dem value_template der
|
||||
Home-Assistant-Discovery. Ohne Pfad gilt die ganze Nutzlast, und
|
||||
JSON-Skalare werden ausgepackt: manche Geraete schicken 21.4 mit
|
||||
Anfuehrungszeichen, andere true statt ON.
|
||||
|
||||
Steht eine Werttabelle dabei, wird aus der gesendeten Zahl der
|
||||
Zustandsname: aus 2 wird "Charging".
|
||||
"""
|
||||
if daten is None:
|
||||
return uebersetze(text, tabelle)
|
||||
try:
|
||||
wert = wert_aus_pfad(daten, pfad)
|
||||
except (KeyError, IndexError, TypeError):
|
||||
logger.debug("Pfad %s nicht in der Nutzlast: %s", pfad, text[:80])
|
||||
return None
|
||||
if isinstance(wert, bool):
|
||||
wert = "true" if wert else "false"
|
||||
elif isinstance(wert, (int, float, str)):
|
||||
wert = str(wert)
|
||||
else:
|
||||
return json.dumps(wert, ensure_ascii=False)
|
||||
return uebersetze(wert, tabelle)
|
||||
|
||||
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}))
|
||||
|
||||
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:
|
||||
# Bei WLED ist die state_url selbst schon der Pfad. Ein
|
||||
# gesetzter value_path hat trotzdem Vorrang, falls das
|
||||
# Geraetemodell spaeter darauf umgestellt wird.
|
||||
pfad = s.get("value_path") or s["state_url"]
|
||||
try:
|
||||
werte[s["id"]] = str(wert_aus_pfad(daten, pfad))
|
||||
except (KeyError, IndexError, TypeError):
|
||||
logger.debug("WLED %s: Pfad %s nicht gefunden", actor_url, pfad)
|
||||
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://"
|
||||
|
||||
# Bis zu dieser Neigung fahren die Aussenjalousien direkt.
|
||||
#
|
||||
# Sie haben eine Kugelschreiber-Mechanik: beim Herunterfahren stehen die
|
||||
# Lamellen bei etwa 30 %. In Richtung 0 % laesst sich von jeder Stellung
|
||||
# aus direkt neigen; darueber hinaus muss die Mechanik erst einmal auf
|
||||
# 0 % zurueck, sonst rastet sie nicht um - ein "Neigung 80 %" ohne
|
||||
# diesen Umweg bleibt wirkungslos.
|
||||
#
|
||||
# Das gilt fuer jedes Kommando, das die Neigung setzt - "Neigung" ebenso
|
||||
# wie "Position+Neigung". Erkannt wird es deshalb am Parameter und nicht
|
||||
# am Kommandonamen: beide heissen ihren Neigungsparameter "Neigung".
|
||||
#
|
||||
# Dasselbe steht in restricted/commands.php: beide Versender brauchen es.
|
||||
NEIGUNG_DIREKT_MAX = 30
|
||||
|
||||
# So heissen die beiden Parameter einer Jalousie im Geraetemodell.
|
||||
NEIGUNG_PARAMETER = "Neigung"
|
||||
POSITION_PARAMETER = "Position"
|
||||
|
||||
# So lange wird hoechstens auf das Ende einer Fahrt gewartet. Gemessen:
|
||||
# eine Neigung von 100 % auf 0 % dauert gut fuenfzehn Sekunden, eine
|
||||
# volle Fahrt von oben nach unten rund sechzig. Die Grenze ist die
|
||||
# Notbremse, nicht die uebliche Dauer.
|
||||
JALOUSIE_WARTE_SEKUNDEN = 120
|
||||
|
||||
# So lange gilt ein "faehrt nicht" direkt nach dem Absenden als noch
|
||||
# nicht aussagekraeftig: die Box meldet core:MovingState traege, kurz
|
||||
# nach einem Kommando steht dort noch der alte Wert.
|
||||
JALOUSIE_VORLAUF_SEKUNDEN = 8
|
||||
|
||||
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 = []
|
||||
self.kombigeraete = set()
|
||||
|
||||
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 kombigeraete_setzen(self, urls):
|
||||
"""
|
||||
Welche Geraete Position und Neigung zusammen koennen - vom Runner aus
|
||||
dem Geraetemodell gesetzt, damit der Transport dafuer nicht selbst in
|
||||
die Datenbank greifen muss.
|
||||
"""
|
||||
self.kombigeraete = set(urls)
|
||||
|
||||
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 = aktion["command_url"]
|
||||
parameter = [self._zahl(p["wert"]) for p in aktion["params"]]
|
||||
neigung_index = None
|
||||
position_index = None
|
||||
for i, p in enumerate(aktion["params"]):
|
||||
if p["name"] == self.NEIGUNG_PARAMETER:
|
||||
neigung_index = i
|
||||
elif p["name"] == self.POSITION_PARAMETER:
|
||||
position_index = i
|
||||
|
||||
# "Zu" allein macht diese Jalousien nicht dicht: sie faehrt herunter,
|
||||
# die Lamellen bleiben durch die Mechanik aber bei etwa 30 % offen.
|
||||
# Gemeint ist "ganz unten, Lamellen geschlossen" - also dasselbe wie
|
||||
# Position 100 mit Neigung 100, und damit ein Fall fuer die Regel.
|
||||
if befehl == "down" and self._kannKombi(aktion["actor_url"]):
|
||||
befehl = "setClosureAndOrientation"
|
||||
parameter = [100, 100]
|
||||
position_index, neigung_index = 0, 1
|
||||
|
||||
# Kugelschreiber-Mechanik, siehe NEIGUNG_DIREKT_MAX. Geschickt wird
|
||||
# derselbe Befehl zweimal - erst mit Neigung 0, dann mit dem
|
||||
# gewuenschten Wert. Die Position bleibt dabei stehen, die Jalousie
|
||||
# faehrt also nur einmal.
|
||||
umweg = (neigung_index is not None
|
||||
and isinstance(parameter[neigung_index], (int, float))
|
||||
and parameter[neigung_index] > self.NEIGUNG_DIREKT_MAX)
|
||||
vorstufe = list(parameter)
|
||||
if umweg:
|
||||
vorstufe[neigung_index] = 0
|
||||
|
||||
if self.dry_run:
|
||||
logger.info("[dry-run] Tahoma %s %s%s", befehl, parameter,
|
||||
" (zuerst %s, dann warten)" % vorstufe if umweg else "")
|
||||
return
|
||||
|
||||
if umweg:
|
||||
self._apply(aktion["actor_url"], befehl, vorstufe)
|
||||
self._warteAufJalousie(aktion["actor_url"], 0,
|
||||
None if position_index is None else int(parameter[position_index]))
|
||||
self._apply(aktion["actor_url"], befehl, parameter)
|
||||
|
||||
def _kannKombi(self, actor_url):
|
||||
"""
|
||||
Hat das Geraet ein Kommando fuer Position und Neigung zusammen?
|
||||
Nur solche Geraete sind Jalousien mit der Kugelschreiber-Mechanik.
|
||||
"""
|
||||
return actor_url in self.kombigeraete
|
||||
|
||||
def _apply(self, actor_url, name, parameter):
|
||||
"""Ein Kommando an die Box schicken."""
|
||||
rumpf = {"label": "AutoAction",
|
||||
"actions": [{"deviceURL": actor_url,
|
||||
"commands": [{"name": name, "parameters": parameter}]}]}
|
||||
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]))
|
||||
|
||||
def _warteAufJalousie(self, actor_url, neigung_ziel, schliessung_ziel=None):
|
||||
"""
|
||||
Wartet, bis die Jalousie ihre Fahrt beendet hat und die Ziele zeigt.
|
||||
|
||||
Zwei Auskuenfte zusammen, weil einzeln keine traegt:
|
||||
core:MovingState taugt fuer die lange Fahrt hoch und runter, wird bei
|
||||
kurzen Neigungsfahrten aber nie gesetzt; die Zustandswerte sind die
|
||||
eigentliche Wahrheit, zeigen kurz nach dem Kommando aber noch den
|
||||
alten Stand. Fertig ist die Fahrt, wenn nichts mehr faehrt, die Ziele
|
||||
erreicht sind und entweder ein "faehrt" gesehen wurde oder der
|
||||
Vorlauf um ist.
|
||||
"""
|
||||
start = time.time()
|
||||
gestartet = False
|
||||
while time.time() - start < self.JALOUSIE_WARTE_SEKUNDEN:
|
||||
time.sleep(2)
|
||||
try:
|
||||
antwort = self.requests.get(
|
||||
self.basis + "/setup/devices/" + quote(actor_url, safe="") + "/states",
|
||||
headers=self._kopf(), timeout=self.timeout, verify=False)
|
||||
z = {x.get("name"): x.get("value") for x in antwort.json()}
|
||||
except Exception as fehler:
|
||||
logger.debug("Zustand nicht lesbar: %s", fehler)
|
||||
continue
|
||||
if z.get("core:MovingState") is True:
|
||||
gestartet = True
|
||||
continue
|
||||
neigung = z.get("core:SlateOrientationState")
|
||||
schliessung = z.get("core:ClosureState")
|
||||
# Ein fehlendes Feld darf nicht als 0 durchgehen - das waere
|
||||
# ausgerechnet beim Ziel 0 ein falsches Erfolgssignal.
|
||||
neigung_ok = neigung is not None and int(neigung) == neigung_ziel
|
||||
schliessung_ok = (schliessung_ziel is None
|
||||
or (schliessung is not None and int(schliessung) == schliessung_ziel))
|
||||
if neigung_ok and schliessung_ok and (
|
||||
gestartet or time.time() - start >= self.JALOUSIE_VORLAUF_SEKUNDEN):
|
||||
return True
|
||||
logger.warning("%s hat Neigung %s%% nicht innerhalb von %d s erreicht",
|
||||
actor_url, neigung_ziel, self.JALOUSIE_WARTE_SEKUNDEN)
|
||||
return False
|
||||
|
||||
@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")
|
||||
Reference in New Issue
Block a user