Files
SolarManager/autoActions/autoaction_runner.py
T
adminandClaude Opus 5 46ef245f7c Doku auf den Stand gebracht: Vorabend, skoda.conf, Wecker-Reste
- autoActions/README: Rahmen gilt heute oder am Vorabend fuer morgen,
  vorabend.sql beim Einrichten, startSolarServer.sh startet drei Prozesse,
  Beispiel der Verkettung wie die echten Automatiken, Neustart ueber SSH
- README: skoda_ladepunkte und skoda.conf, Werkzeuge-Tabelle, Datenbank
  alarm ist geloescht
- config.ini.example: [alarm] und zeit.py entfernt, die gibt es nicht mehr
- Runner-Kopfkommentar verweist nicht mehr auf auto_watering.py

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-09-15 09:53:08 +02:00

1387 lines
62 KiB
Python

#!/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,
Sammlerfaeden fuer die gepollten Geraete starten
Ereignis MQTT-Nachricht -> Wert merken -> im naechsten Takt auswerten
Hintergrund je Transport ein Faden: HTTP und WLED im Minutentakt, Tahoma
alle fuenf Minuten. Die Werte landen in einer Queue.
Takt alle tick_seconds: Queue leeren, auswerten, Regelwerk auf
Aenderung pruefen. Nichts davon wartet auf ein Netz.
Pruefen aktiv? (enabled, Wochentag, Ferien/Feiertag, Zeitfenster)
-> Bedingungen auswerten -> steigende Flanke -> Aktionen
Die Uhr gehoert ausdruecklich nicht zu den abgefragten Geraeten. Sie stand
frueher mit im Geraete-Poll, und weil neunzehn Tahoma-Geraete nacheinander
laenger als eine Minute brauchten, kam jede dritte Minute nie vor: ein
Ausloeser "um 18:26" wurde nie wahr. Verglichen wird jetzt direkt gegen
datetime.now() - siehe uhrzeit_erfuellt().
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" gilt ab dieser Minute noch
catchup_minutes lang, "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. Das Nachholfenster bei "um" ist der Ersatz fuer die frueher
verlangte Punktgenauigkeit: ein Neustart, ein haengendes Geraet oder ein
langsamer Durchlauf kosten die Automatik nicht mehr den ganzen Tag, und weil
nur die Flanke zaehlt, laeuft sie trotzdem hoechstens einmal.
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.
Dieselbe Schreibweise traegt die Verkettung: eine Automatik kann eine andere
ausloesen, indem sie deren letzte Ausloesung als Messwert abfragt - "Wecker
Magdalena + 00:10". Dafuer gibt es das gerechnete Geraet "Automatiken"
(AutomatikTransport) und den Datentyp `elapsed`. Der Nachfolger bleibt dabei
eine vollwertige Automatik mit eigenen Rahmenbedingungen; genau darum geht es
ja - "zehn Minuten spaeter, aber nur wenn es dann schon hell ist" waere als
blosse Verzoegerung an einer Aktion nicht formulierbar, weil die
Zusatzbedingung erst zum spaeteren Zeitpunkt gilt.
Tabellen siehe homeMesh_automations.sql, Konfiguration siehe config.ini.example.
"""
import argparse
import configparser
import json
import logging
import os
import queue
import signal
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 (AUTOMATIK_URL, AutomatikTransport, HTTPTransport,
LogicTransport, MQTTTransport, TahomaTransport,
WLEDTransport, ausloeser_kennung, ausloeser_url,
ist_topic)
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. Die Zeilen zu zaehlen genuegt aber
auch nicht - die Weboberflaeche speichert eine Automatik, indem sie
deren Bedingungen und Aktionen loescht und gleich wieder einfuegt. Wer
nur die Uhrzeit einer Bedingung verstellt, aendert damit weder die
Zeilenzahl noch `changed`: ON UPDATE stoesst nur an, wenn sich in
automations wirklich eine Spalte aendert, und Name, Stockwerk und
Zeitfenster stehen ja noch genauso da. Der Runner lief dann bis zum
naechsten Neustart mit der alten Uhrzeit weiter, ohne dass irgendwo
etwas schieflief - er wusste es schlicht nicht besser.
Deshalb geht jetzt der Inhalt mit ein, als Summe der CRC32 je Zeile.
Das ist kein Hash mit Sicherheitsanspruch, sondern ein billiger
Fingerabdruck - er wird alle paar Sekunden gebildet und darf nichts
kosten. Zwei Aenderungen, die sich in der Summe gegenseitig aufheben,
sind theoretisch denkbar und praktisch nicht zu erwarten. Die
Zeilenzahl steht trotzdem daneben, damit eine geloeschte und eine neu
angelegte Zeile nicht zufaellig gleich viel ergeben.
Neu dabei sind die Aktionsparameter. Sie standen vorher gar nicht
drin: eine geaenderte Zielhoehe einer Jalousie schlug also ebenso
wenig durch.
Von den Geraetetabellen zaehlt nicht nur, wie viele Zeilen es gibt,
sondern auch, wohin sie zeigen. Ein Discovery-Lauf legt naemlich nicht
nur an - er schreibt auch bestehende Zeilen um: bei den neueren
Shellys ist aus der HTTP-Adresse ein Topic geworden, bei gleicher
Zeilenzahl. Ohne url und value_path im Fingerabdruck haette der Runner
das erst beim naechsten Neustart bemerkt und bis dahin an einer
Adresse gefragt, an der niemand mehr antwortet.
"""
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 COALESCE(SUM(CRC32(CONCAT_WS(':',
id, automation_id, group_no, position,
state_id, operator, value))), 0)
FROM automation_conditions) AS bs,
(SELECT COUNT(*) FROM automation_actions) AS c,
(SELECT COALESCE(SUM(CRC32(CONCAT_WS(':',
id, automation_id, position, command_id))), 0)
FROM automation_actions) AS cs,
(SELECT COUNT(*) FROM automation_action_params) AS p,
(SELECT COALESCE(SUM(CRC32(CONCAT_WS(':',
action_id, parameter_id, value))), 0)
FROM automation_action_params) AS ps,
(SELECT COUNT(*) FROM actor_states) AS d,
(SELECT COALESCE(SUM(CRC32(CONCAT_WS(':',
id, url, value_path))), 0)
FROM actor_states) AS ds,
(SELECT COUNT(*) FROM actor_commands) AS e,
(SELECT COALESCE(SUM(CRC32(CONCAT_WS(':',
id, command_url))), 0)
FROM actor_commands) AS es""")
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"}
# Wie lange ein punktgenauer Ausloeser ("um 16:30", "Sonnenaufgang + 00:30")
# nachtraeglich noch gilt, in Minuten. Ohne dieses Fenster muesste die
# Auswertung genau in dieser einen Minute stattfinden; ein langsamer
# Durchlauf, ein Neustart oder ein haengendes Geraet haetten die Automatik
# fuer den Tag gekostet. Ausgeloest wird trotzdem nur einmal, weil nur die
# steigende Flanke zaehlt (automations.cond_met).
NACHHOLFENSTER = 5
def im_nachholfenster(jetzt_m, ziel_m, fenster):
"""
Liegt die Zielminute hoechstens `fenster` Minuten zurueck?
Gerechnet wird modulo 24 Stunden, damit ein Ziel kurz vor Mitternacht auch
nach Mitternacht noch zieht: 23:58 ist um 00:01 drei Minuten her.
"""
return (jetzt_m - ziel_m) % 1440 < max(1, fenster)
def uhrzeit_erfuellt(typ, op, soll, jetzt, fenster):
"""
Uhrzeit und Datum gegen die echte Uhr, nicht gegen einen Abtastwert.
Frueher stand hier der zuletzt abgetastete Wert des Messwerts "Uhrzeit",
den der Runner im Geraete-Poll mitfuehrte. Der Poll brauchte aber laenger
als eine Minute, sodass jede dritte Minute nie vorkam - und "um 18:26"
schlicht nie wahr wurde. `jetzt` liegt hier ohnehin vor.
"""
if typ == "date":
soll_d = als_datum(soll)
if soll_d is None:
return False
ist_d = jetzt.date()
if op == "=": return ist_d == soll_d
if op == "<": return ist_d < soll_d
return ist_d >= soll_d
jetzt_m, soll_m = jetzt.hour * 60 + jetzt.minute, minuten(soll)
if op == "=": return im_nachholfenster(jetzt_m, soll_m, fenster)
if op == "<": return jetzt_m < soll_m
return jetzt_m >= soll_m
def bedingung_erfuellt(bedingung, state, wert, jetzt, fenster=NACHHOLFENSTER):
"""
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.
Uhrzeit und Datum sind die Ausnahme: die werden nicht gemessen, sondern
abgelesen. Sie kommen deshalb direkt aus `jetzt` und nicht aus `wert` -
siehe uhrzeit_erfuellt(). Fuer die Anzeige im Editor fuehrt der Runner sie
zwar auch als Messwert mit, aber ein Vergleich darf nicht davon abhaengen,
wie frisch dieser Abtastwert gerade ist.
`fenster` ist die Nachholzeit in Minuten fuer punktgenaue Ausloeser.
"""
typ = state["type"]
op = bedingung["operator"]
soll = bedingung["value"]
if state.get("actor_url") == "Logic" and typ in ("time", "date"):
return uhrzeit_erfuellt(typ, op, soll, jetzt, fenster)
if wert is None or wert == "":
return False
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":
# Zeit-Messwerte, die tatsaechlich von einem Geraet kommen. Die
# Uhr selbst laeuft ueber uhrzeit_erfuellt().
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 im_nachholfenster(jetzt_m, ziel, fenster)
if typ == "elapsed":
# Der Messwert ist der Zeitpunkt, zu dem eine andere Automatik
# zuletzt gelaufen ist; die Schwelle der Versatz danach.
#
# Anders als bei `deltatime` wird hier NICHT mit der Uhrzeit
# innerhalb des Tages gerechnet, sondern mit dem echten Abstand.
# Sonst machte ein Lauf von vorgestern um 05:50 die Bedingung
# heute um 06:00 wahr, an einem Tag, an dem der Ausloeser gar
# nicht gelaufen ist.
letzter = als_zeitpunkt(wert)
if letzter is None or letzter > jetzt:
return False
verstrichen = (jetzt - letzter).total_seconds() / 60.0
ziel = minuten(soll)
if op.startswith(">="):
# "ab + 00:10" laeuft sonst unbegrenzt weiter - morgen frueh
# waere es immer noch wahr, obwohl der Ausloeser seither
# nichts getan hat. Begrenzt wird wie bei "ab 16:30": bis
# Mitternacht. Ein Ausloeser um 23:55 traegt seinen
# Nachfolger deshalb nicht ueber den Tagesrand - dieselbe
# Einschraenkung hat "ab 23:55" auch.
if letzter.date() != jetzt.date():
return False
return verstrichen >= ziel
if op.startswith("<"):
return verstrichen < ziel
# "um + 00:10": die Punktform, begrenzt durch das Nachholfenster.
# Sie braucht den Tagesvergleich nicht und traegt deshalb auch
# ueber Mitternacht.
return ziel <= verstrichen < ziel + max(1, fenster)
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, fenster=NACHHOLFENSTER):
"""
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, fenster)
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)
#: Ferien und Feiertage sind dreiwertig, nicht ja/nein.
NIE, EGAL, ZUSAETZLICH = 0, 1, 2
def gemeinter_tag(automatik, jetzt):
"""
Fuer welchen Tag gelten Wochentage, Ferien und Feiertage?
Normalerweise fuer heute. Mit `next_day` fuer morgen - die Vorabend-Form:
"Kinderrollos zu, wenn morgen Schule ist" heisst Mo-Fr, Ferien nie,
Feiertage nie, und der Runner schaut dafuer auf den folgenden Tag. Mit
dem heutigen Tag liess sich das nur annaehern (So-Do, heute keine
Ferien), und das ging am letzten Ferientag, am Abend vor einem Feiertag
und am Abend eines Feiertags daneben.
Nur der Rahmen verschiebt sich. Uhrzeit, Zeitfenster und "einmal am Tag"
bleiben beim heutigen Tag - die Automatik laeuft ja heute Abend. Ein
Zeitfenster ueber Mitternacht meint nach Mitternacht deshalb schon den
uebernaechsten Tag; fuer eine Vorabend-Regel ist das kein sinnvoller Fall.
"""
heute = jetzt.date()
return heute + timedelta(days=1) if automatik.get("next_day") else heute
def tag_passt(automatik, tag, kalender):
"""
Faellt der Tag in den Rahmen der Automatik? `tag` ist heute oder morgen,
siehe gemeinter_tag(); `kalender` gehoert zu genau diesem Tag.
Die Wochentage sind eine Maske, Ferien und Feiertage haben je drei
Zustaende. Zwei davon gab es immer: EGAL (der Tag aendert nichts, die
Vorgabe) und NIE (an solchen Tagen laeuft die Automatik nicht - so ist
"werktags" gebaut). ZUSAETZLICH ist der dritte und zaehlt wie ein
passender Wochentag.
Der Grund fuer den dritten: "Wochenenden und Feiertage" liess sich vorher
gar nicht schreiben. Die Wochentagsmaske kennt nur Samstag und Sonntag,
und ein Feiertag am Dienstag ist eben ein Dienstag. Mit ZUSAETZLICH
genuegt Sa+So angehakt und "Feiertage: zusaetzlich" - der Dienstag kommt
dann ueber den Kalender herein.
Reihenfolge: erst wird geoeffnet, dann gesperrt. Wer in den Ferien nie
laufen soll und an Feiertagen zusaetzlich, laeuft an einem Feiertag in
den Ferien nicht - ein Verbot schlaegt eine Erweiterung. Anders herum
liesse sich "nie" nicht mehr verlassen.
"""
passt = bool(automatik["weekdays"] & (1 << tag.weekday()))
if kalender["feiertag"] and automatik["on_holiday"] == ZUSAETZLICH:
passt = True
if kalender["ferien"] and automatik["on_vacation"] == ZUSAETZLICH:
passt = True
if not passt:
return False
if kalender["feiertag"] and automatik["on_holiday"] == NIE:
return False
if kalender["ferien"] and automatik["on_vacation"] == NIE:
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.messwerte = 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))
# Was das Geraet danach meldet, gleich mitnehmen. Bei Tahoma
# laegen sonst bis zu fuenf Minuten zwischen der Fahrt und
# dem neuen Stand in der Tabelle; die anderen Transporte
# geben hier nichts zurueck. Dass das den Faden aufhaelt, ist
# gewollt: das naechste Kommando an dieselbe Jalousie darf
# ohnehin erst nach der Fahrt kommen.
werte = transport.nachlesen(auftrag["actor_url"])
if werte:
self.messwerte.put(werte)
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 messwerte_abholen(self):
"""Was die Versandfaeden nach ihren Kommandos abgelesen haben."""
neu = {}
while True:
try:
neu.update(self.messwerte.get_nowait())
except queue.Empty:
return neu
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)
class Sammler:
"""
Fragt die Geraete ab, die sich nicht von selbst melden - im Hintergrund.
Frueher geschah das mitten in der Hauptschleife: neunzehn Tahoma-Geraete
nacheinander, jedes mit bis zu zehn Sekunden Zeitlimit, dazu die HTTP- und
WLED-Geraete. Eine Runde dauerte dadurch rund fuenfundvierzig statt
dreissig Sekunden. Gepollt wurde, sobald seit dem letzten Poll sechzig
Sekunden vergangen waren - also erst jede zweite Runde, in Wahrheit alle
neunzig Sekunden. Weil die Uhr an derselben Abfrage hing, uebersprang der
Runner jede dritte Minute, und ein Ausloeser "um 18:26" wurde nie wahr.
Jetzt hat jeder Transport seinen eigenen Faden und seinen eigenen Abstand:
die Rollaeden duerfen gemuetlich alle fuenf Minuten, waehrend die
Hauptschleife im Sekundentakt weiterlaeuft.
Die Faeden fassen die Datenbank nicht an - sie legen ihre Werte in eine
Queue, geschrieben wird im Hauptfaden. Dieselbe Regel wie beim Versand:
eine pymysql-Verbindung gehoert einem Faden.
"""
def __init__(self):
self.ergebnisse = queue.Queue()
self.laeuft = True
def aufnehmen(self, transport, abstand):
faden = threading.Thread(target=self._arbeiten, args=(transport, abstand),
name="sammler-" + transport.schema, daemon=True)
faden.start()
def _arbeiten(self, transport, abstand):
while self.laeuft:
try:
werte = transport.zustaende_lesen()
if werte:
self.ergebnisse.put(werte)
except Exception as fehler:
logger.warning("%s nicht abfragbar: %s", transport.schema, fehler)
# In Sekundenschritten warten, damit das Beenden nicht bis zum
# naechsten Durchgang dauert - bei Tahoma waeren das fuenf Minuten.
for _ in range(max(1, int(abstand))):
if not self.laeuft:
return
time.sleep(1)
def abholen(self):
"""Alles, was die Faeden seit dem letzten Mal geliefert haben."""
neu = {}
while True:
try:
neu.update(self.ergebnisse.get_nowait())
except queue.Empty:
return neu
def beenden(self):
self.laeuft = False
# ===========================================================================
# 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 -> zuletzt in die DB geschriebener Wert
self.geschrieben_um = {} # state_id -> wann das war
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 = {} # datum -> {"feiertag": .., "ferien": ..}
self._letzte_saeuberung = None
self.reihenfolge = [] # automation_id, Ausloeser vor Nachfolger
self.ausloeser_states = {} # automation_id -> state_id des Ausloesers
self.versand = Versand()
self.sammler = Sammler()
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),
AutomatikTransport(self.ausloesezeiten),
]
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 transport_fuer_messwert(self, state):
"""
Wo ein einzelner Messwert gelesen wird.
Normalerweise sagt das Geraet es an: ein Shelly haengt an HTTP, eine
Jalousie an der Tahoma-Box. Bei den Shellys der zweiten Generation
faellt beides auseinander - sie schicken ihre Messwerte von selbst an
den Broker, geschaltet werden sie weiter ueber HTTP. Ein Messwert mit
einem Topic wird deshalb dort gelesen, wo er ankommt, und nicht dort,
wo sein Geraet sonst zu erreichen ist.
Fuer die aelteren Shellys aendert sich nichts: in ihren Messwerten
steht ein Feldname, und der fuehrt weiter zum HTTP-Transport.
"""
if ist_topic(state["state_url"]):
for transport in self.transporte:
if isinstance(transport, MQTTTransport):
return transport
return self.transport_fuer(state["actor_url"])
def regelwerk_laden(self):
# Vor dem Laden, damit eine neu angelegte Automatik sofort als
# Ausloeser zur Verfuegung steht und ihre Zeile schon in der
# Signatur steckt - sonst laedt der naechste Takt gleich noch einmal.
self.ausloeser_nachfuehren()
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.
# Zugeteilt wird je Messwert und nicht je Geraet, weil beides
# auseinanderfallen kann - siehe transport_fuer_messwert().
listen = [(transport, []) for transport in self.transporte]
for s in self.regelwerk.states.values():
zustaendig = self.transport_fuer_messwert(s)
for transport, liste in listen:
if transport is zustaendig:
liste.append({"id": s["id"], "actor_url": s["actor_url"],
"state_url": s["state_url"], "value_path": s["value_path"],
"wertetabelle": s["wertetabelle"]})
break
for transport, liste in listen:
transport.zustaende_anmelden(liste)
# 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"]
# Was schon in der Tabelle steht, muss nach einem Neustart nicht
# noch einmal hineingeschrieben werden.
self.geschrieben.setdefault(s["id"], s["current_value"])
self.ausloeser_states = {
ausloeser_kennung(s["state_url"]): s["id"]
for s in self.regelwerk.states.values()
if s["actor_url"] == AUTOMATIK_URL
and ausloeser_kennung(s["state_url"]) is not None}
self.reihenfolge_bestimmen()
def ausloeser_nachfuehren(self):
"""
Je Automatik einen Messwert am gerechneten Geraet "Automatiken".
Damit taucht jede Automatik im Editor als Messwert auf und laesst sich
als Ausloeser einer anderen waehlen, ohne dass der Editor davon etwas
wissen muesste - er listet Geraete und deren Messwerte, mehr nicht.
Gefuehrt wird nach der Kennung, nicht nach dem Namen: wer umbenennt,
soll die abhaengigen Bedingungen nicht verlieren. Der Name wird
nachgezogen, damit im Editor das Richtige steht.
Angelegt wird fuer JEDE Automatik, auch fuer pausierte. Sonst loeschte
ein Pausieren ueber `fk_cond_state ON DELETE CASCADE` die Bedingung
des Nachfolgers - still, und beim Fortsetzen waere sie weg.
Geloescht wird nur, was niemand mehr benutzt. Zeigt noch eine
Bedingung darauf, bleibt der Messwert stehen und es gibt eine Warnung:
die Kaskade wuerde sonst eine einzelne Bedingung aus einer Gruppe
entfernen und aus "zehn Minuten nach dem Wecker UND es ist hell" ein
blosses "es ist hell" machen. Bei der LETZTEN Bedingung faengt
gruppen_erfuellt() das ab, bei einer von zweien niemand.
"""
try:
with self.db.cursor() as c:
c.execute("SELECT id FROM actors WHERE url = %s", (AUTOMATIK_URL,))
zeile = c.fetchone()
if not zeile:
logger.debug("Geraet \"Automatiken\" gibt es nicht - "
"automatik_ausloeser.sql noch nicht eingespielt")
return
aktor = zeile["id"]
c.execute("SELECT id FROM state_types WHERE type = 'elapsed'")
typ = c.fetchone()
typ_id = typ["id"] if typ else None
c.execute("SELECT id, name FROM automations")
gewuenscht = {ausloeser_url(r["id"]): r["name"] for r in c.fetchall()}
c.execute("""SELECT s.id, s.state_name, s.url,
(SELECT COUNT(*) FROM automation_conditions b
WHERE b.state_id = s.id) AS benutzt
FROM actor_states s WHERE s.actor_id = %s""", (aktor,))
vorhanden = {r["url"]: r for r in c.fetchall()}
for url, name in gewuenscht.items():
alt = vorhanden.get(url)
if alt is None:
c.execute("""INSERT INTO actor_states
(actor_id, state_name, state_type, url,
possible_values)
VALUES (%s, %s, %s, %s, '')""",
(aktor, name, typ_id, url))
elif alt["state_name"] != name:
c.execute("UPDATE actor_states SET state_name = %s WHERE id = %s",
(name, alt["id"]))
for url, alt in vorhanden.items():
if url in gewuenscht:
continue
if alt["benutzt"]:
logger.warning(
"Ausloeser %r zeigt auf eine geloeschte Automatik, wird "
"aber noch von %d Bedingung(en) benutzt - bleibt stehen",
alt["state_name"], alt["benutzt"])
continue
c.execute("DELETE FROM actor_states WHERE id = %s", (alt["id"],))
except Exception as fehler:
logger.warning("Ausloeser nicht nachfuehrbar: %r", fehler)
def reihenfolge_bestimmen(self):
"""
Ausloeser vor Nachfolger auswerten.
Nur damit wirkt ein Versatz von null noch im selben Takt. Bei jedem
anderen Versatz waere die Reihenfolge gleichgueltig - zehn Minuten
sind laenger als ein Takt.
Ein Kreis (A loest B loest A) waere ein Fehler im Regelwerk; der
Editor lehnt ihn beim Speichern ab. Hier wird er nur gemeldet und die
Beteiligten laufen in ihrer urspruenglichen Reihenfolge weiter. Sie
deswegen stillzulegen waere schlimmer: die Sperrzeit begrenzt den
Schaden ohnehin auf eine Ausloesung je lockout_secs, eine wortlos
abgeschaltete Automatik dagegen faellt niemandem auf.
"""
automatiken = self.regelwerk.automatiken
vorgaenger = {aid: set() for aid in automatiken}
for aid, auto in automatiken.items():
for bedingungen in auto["gruppen"].values():
for b in bedingungen:
state = self.regelwerk.states.get(b["state_id"])
if not state or state["actor_url"] != AUTOMATIK_URL:
continue
davor = ausloeser_kennung(state["state_url"])
if davor == aid:
logger.error("%s loest sich selbst aus - Bedingung wird "
"nie wahr", auto["name"])
elif davor in automatiken:
vorgaenger[aid].add(davor)
reihenfolge, offen = [], dict(vorgaenger)
while offen:
frei = sorted(aid for aid, davor in offen.items()
if not davor & set(offen))
if not frei:
logger.error("Automatiken loesen sich im Kreis aus: %s",
", ".join(automatiken[aid]["name"] for aid in offen))
reihenfolge.extend(offen)
break
reihenfolge.extend(frei)
for aid in frei:
del offen[aid]
self.reihenfolge = reihenfolge
def ausloesezeiten(self):
"""
Wann jede Automatik zuletzt gelaufen ist - die Werte des gerechneten
Geraets "Automatiken".
Es stehen nur die aktiven darin: `Regelwerk.laden` holt sich
`WHERE enabled = 1`. Eine pausierte Automatik liefert damit keinen
Zeitpunkt, und ihre Nachfolger stehen mit still. Das ist gewollt -
wer den Wecker pausiert, will morgens auch den Rollladen unten lassen.
"""
if not self.regelwerk:
return {}
return {aid: auto.get("last_run")
for aid, auto in self.regelwerk.automatiken.items()}
# --- 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:
self.solar.ping(reconnect=True) # re-establish connection if it dropped
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: %r", 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, tag):
"""
Ferien und Feiertage eines Tages, je Tag einmal geholt.
Gefragt wird nach heute und - fuer Vorabend-Regeln - nach morgen.
Gemerkt werden deshalb beide; was vor heute liegt, fliegt raus.
"""
if tag in self._kalender:
return self._kalender[tag]
stand = {"feiertag": False, "ferien": False}
try:
with self.db.cursor() as c:
c.execute("SELECT holiday, vacation FROM calendar_days WHERE date = %s", (tag,))
zeile = c.fetchone()
if zeile:
stand = {"feiertag": bool(zeile["holiday"]), "ferien": bool(zeile["vacation"])}
except Exception as fehler:
# Nicht merken: beim naechsten Durchlauf noch einmal fragen,
# statt den ganzen Tag mit "kein Feiertag" weiterzurechnen.
logger.warning("Kalender nicht lesbar: %r", fehler)
return stand
heute = date.today()
self._kalender = {t: s for t, s in self._kalender.items() if t >= heute}
self._kalender[tag] = stand
return stand
# --- Werte -----------------------------------------------------------
def sammler_starten(self):
"""
Je gepolltem Transport ein Faden. MQTT meldet sich von selbst und die
Uhr rechnet der Runner - die beiden stehen hier nicht.
"""
vorgaben = [(HTTPTransport, "poll_http", 60),
(WLEDTransport, "poll_wled", 60),
(TahomaTransport, "poll_tahoma", 300)]
for transport in self.transporte:
for klasse, schluessel, vorgabe in vorgaben:
if isinstance(transport, klasse):
abstand = self.config.zahl("runner", schluessel, vorgabe)
self.sammler.aufnehmen(transport, abstand)
logger.info("%s wird alle %d s abgefragt", klasse.__name__, abstand)
break
def werte_einsammeln(self, auch_geraete=False):
"""
MQTT und Uhr kosten nichts und werden jeden Takt gelesen; die Geraete
liefert der Sammler aus dem Hintergrund. Nur bei --once, wo es keine
Faeden gibt, fragt die Schleife die Geraete selbst.
"""
neu = {}
for transport in self.transporte:
billig = isinstance(transport, (MQTTTransport, LogicTransport,
AutomatikTransport))
if billig or auch_geraete:
neu.update(transport.zustaende_lesen())
neu.update(self.sammler.abholen())
neu.update(self.versand.messwerte_abholen())
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)").
Geschrieben wird nur, was sich geaendert hat. MQTT liefert ohnehin nur
neu Hereingekommenes, die gepollten Transporte dagegen bei jedem
Durchgang ihren kompletten Bestand - ohne diesen Vergleich gingen rund
180 unveraenderte Werte je Runde in die Tabelle. Zusaetzlich
gedrosselt, sonst schreibt ein gespraechiger Sensor im Sekundentakt.
"""
jetzt = datetime.now()
faellig = [(w, i) for i, w in neu.items()
if self.geschrieben.get(i) != w
and self.geschrieben_um.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 wert, state_id in faellig:
self.geschrieben[state_id] = wert
self.geschrieben_um[state_id] = jetzt
except Exception as fehler:
logger.warning("current_value nicht schreibbar: %r", 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: %r", fehler)
automatik["last_run"] = datetime.now()
self.lief_im_fenster[automatik["id"]] = True
# Den eigenen Ausloeserwert gleich mitfuehren, statt bis zum naechsten
# werte_einsammeln() zu warten. Zusammen mit der topologischen
# Reihenfolge greift ein Nachfolger mit Versatz null dadurch noch im
# selben Takt.
state_id = self.ausloeser_states.get(automatik["id"])
if state_id is not None:
self.werte[state_id] = automatik["last_run"].strftime("%Y-%m-%d %H:%M:%S")
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: %r", fehler)
def durchlauf(self, fenster=NACHHOLFENSTER):
jetzt = datetime.now()
for automation_id in self.reihenfolge:
automatik = self.regelwerk.automatiken.get(automation_id)
if automatik is None:
continue
tag = gemeinter_tag(automatik, jetzt)
aktiv = (tag_passt(automatik, tag, self.kalender(tag))
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, fenster)
if erfuellt and not automatik["cond_met"]:
if automatik.get("once_per_day") and self.lief_heute(automatik, jetzt):
logger.debug("%s: lief heute schon, bis Mitternacht gesperrt",
automatik["name"])
elif 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):
"""
Hat die Automatik heute schon ausgeloest?
Zwei Stellen fragen danach. force_once will wissen, ob es am Ende des
Fensters noch etwas nachzuholen gibt. once_per_day will das Gegenteil:
einmal am Tag genuegt, danach ist bis Mitternacht Ruhe.
Letzteres ist fuer alles gedacht, was man hinterher von Hand wieder
anders stellt. Ein Rollladen, den man um acht nochmal zugezogen hat,
soll nicht um neun von selbst wieder auffahren, nur weil eine Wolke
weiterzieht und die Helligkeitsschwelle ein zweites Mal steigt. Die
Sperrzeit taugt dafuer nicht: sie zaehlt Sekunden und muesste auf
einen Tag stehen, womit sie am naechsten Morgen den Termin knapp
verfehlen wuerde.
"""
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: %r", fehler)
# --- Hauptschleife ---------------------------------------------------
def laufen(self, nur_einmal=False):
"""
Die Schleife macht nur noch Billiges: Werte abholen, auswerten,
protokollieren. Alles, was auf ein Netz warten muss, laeuft daneben -
die Geraeteabfrage im Sammler, das Schalten im Versand. Deshalb darf
der Takt kurz sein.
"""
self.regelwerk_laden()
takt = self.config.zahl("runner", "tick_seconds", 10)
neuladen = self.config.zahl("runner", "reload_seconds", 30)
fenster = self.config.zahl("runner", "catchup_minutes", NACHHOLFENSTER)
letztes_pruefen = 0.0
if not nur_einmal:
self.sammler_starten()
while True:
jetzt = time.monotonic()
self.werte_einsammeln(auch_geraete=nur_einmal)
self.durchlauf(fenster)
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: %r", fehler)
if nur_einmal:
return
time.sleep(takt)
def beenden(self):
self.sammler.beenden()
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.")
# startSolarServer.sh beendet eine laufende Instanz mit SIGTERM, bevor es
# die neue startet. Ohne diesen Handler faellt der Prozess sofort um und
# beenden() kaeme nie dran: der Broker hielte die Verbindung noch eine
# Weile fuer lebendig, und ein gerade laufendes Kommando bliebe auf halbem
# Weg stehen. Als KeyboardInterrupt geht es denselben Weg wie Strg-C.
def abbrechen(signum, rahmen):
raise KeyboardInterrupt
signal.signal(signal.SIGTERM, abbrechen)
try:
runner.laufen(args.once)
except KeyboardInterrupt:
logger.info("Abbruch - wird beendet")
finally:
runner.beenden()
return 0
if __name__ == "__main__":
sys.exit(main())