uhrzeit_erfuellt() verglich beim Datentyp "date" bisher das volle Datum samt Jahr. Damit musste jede Automatik mit einer Datumsbedingung jedes Jahr von Hand nachgezogen werden - die Frostwarnung stand deshalb fest auf 2026/2027 und waere zum naechsten Jahreswechsel wieder falsch gewesen. Verglichen wird jetzt nur noch (Monat, Tag). Ein Fenster, das ueber den Jahreswechsel reicht, laesst sich damit weiterhin nicht in einer einzigen Bedingung ausdruecken - dafuer braucht es zwei Bedingungen ohne obere bzw. untere Grenze in getrennten (oder verbundenen) Gruppen, wie es die Frostwarnung schon vorher tat. Dabei die eigentliche Ursache gefunden, warum die Frostwarnung nie ausloeste: alle sechs Bedingungen standen in derselben Gruppe (group_no = 0) und damit UND-verknuepft - ein Datum kann nicht gleichzeitig im September/November- UND im Februar/April-Fenster liegen. Die zweite Haelfte (Bedingungen 271-273) steht jetzt in einer eigenen Gruppe, ODER-verknuepft mit der ersten. Reine Datenaenderung, kein Schema-Update noetig. Getestet gegen das echte Regelwerk, ohne Schaltbefehle: vorher an keinem Tag erfuellbar, nachher am 22.09. (3,6 Grad) korrekt ausgeloest - siehe automation_log, push_abos. Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
1563 lines
71 KiB
Python
1563 lines
71 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.
|
|
|
|
Die Haltezeit (`automations.hold_secs`) schiebt diese Flanke nach hinten: Erst
|
|
wenn die Bedingung so viele Sekunden am Stueck erfuellt war, gilt sie als
|
|
erfuellt. Gedacht fuer Dauerzustaende, die sich nicht in einem einzelnen
|
|
Messwert zeigen - "der Wasserzaehler laeuft" ist jedes Haendewaschen, "laeuft
|
|
seit einer halben Stunde ohne Pause" ist ein offener Hahn. Gezaehlt wird hier
|
|
im Speicher (`erfuellt_seit`) und nicht in der Datenbank: es ist ein
|
|
Laufzustand wie `war_aktiv`, und ein Neustart soll ihn bewusst verwerfen -
|
|
nach einem Neustart weiss niemand, ob die Bedingung in der Zwischenzeit
|
|
durchgehend anlag.
|
|
|
|
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, BENACHRICHTIGUNG_TOPIC,
|
|
BenachrichtigungTransport, 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 abschnitt(self, sektion):
|
|
"""
|
|
Einen ganzen Abschnitt als Dict.
|
|
|
|
Fuer den Mailzugang: der besteht aus einem halben Dutzend Feldern und
|
|
wird als Ganzes an den Transport weitergereicht, statt sechsmal
|
|
einzeln abgefragt zu werden.
|
|
"""
|
|
if not self.cfg.has_section(sektion):
|
|
return {}
|
|
return {k: v.strip() for k, v in self.cfg.items(sektion)}
|
|
|
|
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
|
|
# Nur Monat und Tag zaehlen, das Jahr der Eingabe wird ignoriert -
|
|
# eine Datumsbedingung soll sich jaehrlich wiederholen, ohne dass sie
|
|
# jedes Jahr von Hand nachgezogen werden muss. Der Vergleich zweier
|
|
# (Monat, Tag)-Paare ordnet sich wie das Kalenderjahr selbst; nur ein
|
|
# Fenster, das über den Jahreswechsel reicht, laesst sich damit nicht
|
|
# in einer einzigen Bedingung ausdruecken (siehe doku/automatiken.md).
|
|
ist_mt = (jetzt.month, jetzt.day)
|
|
soll_mt = (soll_d.month, soll_d.day)
|
|
if op == "=": return ist_mt == soll_mt
|
|
if op == "<": return ist_mt < soll_mt
|
|
return ist_mt >= soll_mt
|
|
|
|
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.erfuellt_seit = {} # automation_id -> seit wann die Bedingung anliegt
|
|
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),
|
|
BenachrichtigungTransport(self.push_abos, self.push_abo_weg,
|
|
self.push_abo_erfolg,
|
|
config.abschnitt("mail"),
|
|
# Ein leerer Eintrag in der config.ini ist
|
|
# vorhanden, aber leer - die Vorgabe von
|
|
# text() greift dann nicht.
|
|
(config.text("push", "schluessel")
|
|
or os.path.join(os.path.dirname(os.path.abspath(__file__)),
|
|
"..", "push_vapid.json")),
|
|
self.dry_run),
|
|
]
|
|
# Meldungen, die nicht aus einer Automatik kommen: die Probe aus den
|
|
# Einstellungen, spaeter vielleicht ein Skript. Derselbe Weg, dieselbe
|
|
# Zustellung - nur ohne Umweg ueber eine Regel.
|
|
self.meldungen = queue.Queue()
|
|
try:
|
|
self.mqtt.subscribe(BENACHRICHTIGUNG_TOPIC)
|
|
except Exception as fehler:
|
|
logger.warning("Meldungs-Topic nicht abonnierbar: %r", fehler)
|
|
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.
|
|
if nachricht.topic.startswith("benachrichtigung/"):
|
|
# Nicht hier zustellen: Dieser Rueckruf laeuft im MQTT-Faden, und
|
|
# eine Datenbankverbindung gehoert einem Faden. Der naechste Takt
|
|
# holt es ab.
|
|
schlange = getattr(self, "meldungen", None)
|
|
if schlange is not None:
|
|
schlange.put((nachricht.topic.rsplit("/", 1)[-1], nachricht.payload))
|
|
return
|
|
for transport in getattr(self, "transporte", []):
|
|
if isinstance(transport, MQTTTransport):
|
|
transport.nachricht(nachricht.topic, nachricht.payload)
|
|
|
|
# --- Benachrichtigungen ----------------------------------------------
|
|
|
|
def benachrichtigung(self):
|
|
"""Der Transport fuer die Meldungen, oder None."""
|
|
for transport in self.transporte:
|
|
if isinstance(transport, BenachrichtigungTransport):
|
|
return transport
|
|
return None
|
|
|
|
def meldungen_abarbeiten(self):
|
|
"""
|
|
Was ueber benachrichtigung/# hereinkam, zustellen.
|
|
|
|
Nutzlast ist JSON ({"titel": ..., "text": ...}); ein blosser Text
|
|
geht auch durch und wird zum Textkoerper.
|
|
"""
|
|
transport = self.benachrichtigung()
|
|
while True:
|
|
try:
|
|
kanal, rohtext = self.meldungen.get_nowait()
|
|
except queue.Empty:
|
|
return
|
|
if transport is None:
|
|
continue
|
|
try:
|
|
text = rohtext.decode("utf-8", "replace") if isinstance(rohtext, bytes) else str(rohtext)
|
|
try:
|
|
daten = json.loads(text)
|
|
if not isinstance(daten, dict):
|
|
raise ValueError
|
|
except ValueError:
|
|
daten = {"text": text}
|
|
if kanal == "mail":
|
|
transport.mail_senden(daten.get("betreff") or daten.get("titel") or "Smarthome",
|
|
daten.get("text") or "")
|
|
else:
|
|
transport.push(daten.get("titel") or "Smarthome", daten.get("text") or "")
|
|
except Exception as fehler:
|
|
logger.warning("Meldung (%s) nicht zustellbar: %r", kanal, fehler)
|
|
|
|
def push_abos(self):
|
|
"""Die angemeldeten Geraete. Fehlt die Tabelle, gibt es eben keine."""
|
|
try:
|
|
with self.db.cursor() as c:
|
|
c.execute("SELECT id, endpoint, p256dh, auth, name FROM push_abos")
|
|
return list(c.fetchall())
|
|
except Exception as fehler:
|
|
logger.debug("push_abos nicht lesbar: %r", fehler)
|
|
return []
|
|
|
|
def push_abo_weg(self, abo_id, grund):
|
|
try:
|
|
with self.db.cursor() as c:
|
|
c.execute("DELETE FROM push_abos WHERE id = %s", (abo_id,))
|
|
except Exception as fehler:
|
|
logger.warning("Push-Abo %s nicht loeschbar: %r", abo_id, fehler)
|
|
|
|
def push_abo_erfolg(self, abo_id):
|
|
try:
|
|
with self.db.cursor() as c:
|
|
c.execute("UPDATE push_abos SET zuletzt = NOW(), fehler = '' WHERE id = %s",
|
|
(abo_id,))
|
|
except Exception as fehler:
|
|
logger.debug("Push-Abo %s nicht fortschreibbar: %r", abo_id, fehler)
|
|
|
|
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"],
|
|
# Nur fuer Protokollmeldungen des Transports: "io://1215-.../332898"
|
|
# sagt niemandem, welcher Rollladen gemeint ist.
|
|
"actor_name": kommando["actor_name"],
|
|
"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)
|
|
# Mit Haltezeit zaehlt nicht, ob die Bedingung erfuellt ist,
|
|
# sondern ob sie es lange genug am Stueck ist. Weil das Ergebnis
|
|
# an dieselbe Stelle tritt, bleibt alles danach - Flanke,
|
|
# Tagessperre, Sperrzeit, Protokoll - unveraendert.
|
|
reif = self.haltezeit_reif(automatik, erfuellt, jetzt)
|
|
if reif 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, reif)
|
|
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. Eine
|
|
# angefangene Haltezeit verfaellt mit: sie soll innerhalb des
|
|
# Fensters voll gelaufen sein, nicht ueber dessen Rand hinweg.
|
|
self.flanke_merken(automatik, False)
|
|
self.erfuellt_seit.pop(automatik["id"], None)
|
|
self.lief_im_fenster[automatik["id"]] = False
|
|
|
|
self.war_aktiv[automatik["id"]] = aktiv
|
|
|
|
def haltezeit_reif(self, automatik, erfuellt, jetzt):
|
|
"""
|
|
Liegt die Bedingung lange genug am Stueck an?
|
|
|
|
Ohne Haltezeit (hold_secs = 0, die Vorgabe) ist die Antwort schlicht
|
|
die Bedingung selbst - dann verhaelt sich alles wie vorher.
|
|
|
|
Sonst wird der Zeitpunkt gemerkt, an dem die Bedingung wahr wurde, und
|
|
erst nach Ablauf der Haltezeit "ja" gemeldet. Faellt sie zwischendurch
|
|
auch nur einen Takt aus, wird der Zeitpunkt verworfen und faengt beim
|
|
naechsten Mal von vorn an; genau das ist mit "ununterbrochen" gemeint.
|
|
|
|
Wichtig ist, dass cond_met bis dahin auf 0 bleibt - die Flanke wird
|
|
also nicht verbraucht, sondern aufgeschoben. Deshalb wird sie hier
|
|
nicht selbst gesetzt, sondern das Ergebnis nach oben gereicht.
|
|
|
|
Zwei Dinge, die man wissen sollte:
|
|
|
|
* Ein punktgenauer Zeit-Ausloeser ("um 16:30") ist nur das
|
|
Nachholfenster lang wahr. Eine Haltezeit darueber hinaus wuerde nie
|
|
reif - was kein Fehler, aber auch keine sinnvolle Kombination ist.
|
|
* Eine Bedingung gilt so lange weiter, wie ihr letzter Messwert gilt.
|
|
Bei einem Zaehler, der alle fuenf Minuten meldet, ist die Haltezeit
|
|
also auf fuenf Minuten genau - sinnvolle Stufen beginnen deutlich
|
|
darueber.
|
|
"""
|
|
halten = int(automatik.get("hold_secs") or 0)
|
|
if not halten:
|
|
return erfuellt
|
|
if not erfuellt:
|
|
if self.erfuellt_seit.pop(automatik["id"], None) is not None:
|
|
logger.debug("%s: Haltezeit abgebrochen, faengt von vorn an",
|
|
automatik["name"])
|
|
return False
|
|
seit = self.erfuellt_seit.setdefault(automatik["id"], jetzt)
|
|
offen = halten - (jetzt - seit).total_seconds()
|
|
if offen > 0:
|
|
logger.debug("%s: Bedingung erfuellt, noch %d s Haltezeit",
|
|
automatik["name"], offen)
|
|
return False
|
|
return True
|
|
|
|
@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.meldungen_abarbeiten()
|
|
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())
|