Anlass: "zehn Minuten nach dem Wecker den Rollladen hoch, aber nur wenn es dann schon hell ist". Als Verzoegerung an der Aktion war das nicht zu haben - die Zusatzbedingung gilt erst zum spaeteren Zeitpunkt, und eine verzoegerte Aktion, die selbst noch Bedingungen prueft, braeuchte ein zweites Bedingungssystem neben dem ersten. Der zweite Schritt ist also eine eigene Automatik; was ihr fehlte, war nur ein Bezug auf die erste. Den gibt jetzt das gerechnete Geraet "Automatiken", Gegenstueck zum vorhandenen "Zeitpunkt": jede Automatik ist dort ein Messwert, ihr Wert der Zeitpunkt der letzten Ausloesung. Der Editor braucht dafuer keine Zeile - er listet Geraete und deren Messwerte. Der neue Datentyp `elapsed` verhaelt sich dazu wie `deltatime` zum Sonnenaufgang: "+ 00:10", "ab + 00:10", "vor + 00:10". Ein Minus gibt es nicht. Gerechnet wird mit dem echten Abstand statt mit der Uhrzeit innerhalb des Tages - sonst machte ein Lauf von vorgestern die Bedingung heute wahr. Die offene Form endet trotzdem am Tagesrand, genau wie "ab 16:30". Ausgewertet wird topologisch, Ausloeser vor Nachfolger; nur so wirkt ein Versatz von null noch im selben Takt. Eine pausierte Automatik haelt ihre Nachfolger mit an: geladen werden nur die aktiven, und der Transport liefert fuer alle uebrigen einen leeren Wert. Die Messwerte pflegt ausloeser_nachfuehren() - fuer JEDE Automatik, auch fuer pausierte. Sonst loeschte ein Pausieren ueber fk_cond_state ON DELETE CASCADE die Bedingung des Nachfolgers, still. Geloescht wird nur, was keine Bedingung mehr benutzt. automatik_ausloeser.sql legt Datentyp und Geraet an. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1311 lines
58 KiB
Python
1311 lines
58 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. Dieselbe
|
|
Ueberlegung steht hinter den breiten Zeitfenstern in auto_watering.py.
|
|
|
|
Beim Sonnenauf- und -untergang traegt der Operator zusaetzlich das
|
|
Vorzeichen des Versatzes: "+ 00:30" eine halbe Stunde danach, ">=- 00:30"
|
|
ab einer halben Stunde davor, "<+ 00:30" bis eine halbe Stunde danach.
|
|
|
|
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)
|
|
|
|
|
|
def tag_passt(automatik, jetzt, kalender):
|
|
if not automatik["weekdays"] & (1 << jetzt.weekday()):
|
|
return False
|
|
if kalender["feiertag"] and not automatik["on_holiday"]:
|
|
return False
|
|
if kalender["ferien"] and not automatik["on_vacation"]:
|
|
return False
|
|
return True
|
|
|
|
|
|
# ===========================================================================
|
|
# Versand
|
|
# ===========================================================================
|
|
|
|
class Versand:
|
|
"""
|
|
Schickt Kommandos im Hintergrund - je Geraet der Reihe nach, ueber
|
|
Geraete hinweg nebeneinander.
|
|
|
|
Ohne das blockiert ein einziges Kommando die ganze Auswertung: eine
|
|
Jalousie zuzufahren dauert ueber eine Minute, weil zwischen den beiden
|
|
Schwenkbefehlen auf das Ende der Fahrt gewartet werden muss (siehe
|
|
TahomaTransport). So lange kaeme kein Zeit-Ausloeser mehr durch, kein
|
|
Messwert wuerde zurueckgeschrieben, und mehrere Rollladen in einer
|
|
Automatik wuerden sich aufaddieren.
|
|
|
|
Je Geraet eine Warteschlange mit einem eigenen Faden: zwei Kommandos an
|
|
dieselbe Jalousie duerfen sich nicht ueberholen - "Neigung 0" und
|
|
"Neigung 100" sind sonst wirkungslos oder vertauscht -, zwei Kommandos an
|
|
verschiedene Jalousien duerfen ruhig gleichzeitig laufen.
|
|
|
|
Die Datenbank bleibt aussen vor: die Faeden melden ihr Ergebnis nur
|
|
zurueck, geschrieben wird im Hauptfaden. Eine pymysql-Verbindung ist
|
|
nicht fuer mehrere Faeden gedacht.
|
|
"""
|
|
|
|
def __init__(self):
|
|
self.warteschlangen = {} # actor_url -> Queue
|
|
self.ergebnisse = queue.Queue()
|
|
self.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 = (None, {"feiertag": False, "ferien": False})
|
|
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):
|
|
"""Ferien und Feiertage von heute, einmal je Tag geholt."""
|
|
heute = date.today()
|
|
if self._kalender[0] == heute:
|
|
return self._kalender[1]
|
|
stand = {"feiertag": False, "ferien": False}
|
|
try:
|
|
with self.db.cursor() as c:
|
|
c.execute("SELECT holiday, vacation FROM calendar_days WHERE date = %s", (heute,))
|
|
zeile = c.fetchone()
|
|
if zeile:
|
|
stand = {"feiertag": bool(zeile["holiday"]), "ferien": bool(zeile["vacation"])}
|
|
except Exception as fehler:
|
|
logger.warning("Kalender nicht lesbar: %r", fehler)
|
|
self._kalender = (heute, 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()
|
|
kalender = self.kalender()
|
|
|
|
for automation_id in self.reihenfolge:
|
|
automatik = self.regelwerk.automatiken.get(automation_id)
|
|
if automatik is None:
|
|
continue
|
|
aktiv = (tag_passt(automatik, jetzt, kalender)
|
|
and im_zeitfenster(jetzt, automatik["window_from"], automatik["window_to"]))
|
|
vorher_aktiv = self.war_aktiv.get(automatik["id"], aktiv)
|
|
|
|
if aktiv:
|
|
if not vorher_aktiv:
|
|
self.lief_im_fenster[automatik["id"]] = False
|
|
erfuellt = gruppen_erfuellt(automatik, self.regelwerk, self.werte,
|
|
jetzt, fenster)
|
|
if erfuellt and not automatik["cond_met"]:
|
|
if self.gesperrt(automatik, jetzt):
|
|
logger.debug("%s: Flanke faellt in die Sperrzeit, uebersprungen",
|
|
automatik["name"])
|
|
else:
|
|
self.ausloesen(automatik, "fired")
|
|
self.flanke_merken(automatik, erfuellt)
|
|
else:
|
|
# Das Fenster ist gerade zugegangen. Wer "auf jeden Fall"
|
|
# angehakt hat, bekommt jetzt seinen Lauf - aber nur, wenn in
|
|
# diesem Fenster noch keiner stattgefunden hat. Nach einem
|
|
# Neustart mitten im Fenster weiss der Runner das nicht mehr
|
|
# aus dem Speicher, deshalb zaehlt zusaetzlich last_run.
|
|
if vorher_aktiv and automatik["force_once"] and not self.lief_im_fenster.get(automatik["id"]):
|
|
if not self.lief_heute(automatik, jetzt):
|
|
logger.info("%s: Zeitfenster vorbei, wird trotzdem ausgefuehrt",
|
|
automatik["name"])
|
|
self.ausloesen(automatik, "forced")
|
|
# Ausserhalb des Fensters die Flanke zuruecksetzen, sonst
|
|
# koennte sie im naechsten Fenster nicht mehr steigen.
|
|
self.flanke_merken(automatik, False)
|
|
self.lief_im_fenster[automatik["id"]] = False
|
|
|
|
self.war_aktiv[automatik["id"]] = aktiv
|
|
|
|
@staticmethod
|
|
def gesperrt(automatik, jetzt):
|
|
"""
|
|
Liegt die letzte Ausloesung noch innerhalb der Sperrzeit?
|
|
|
|
Gegen Messwerte, die um die Schwelle pendeln: "Temperatur > 22" bei
|
|
22,1 / 21,9 / 22,1 Grad ist jedes Mal eine echte steigende Flanke, und
|
|
ueber MQTT koennen die Werte im Sekundentakt hereinkommen.
|
|
|
|
Die Flanke wird dabei verworfen und nicht aufgehoben. Ein Rollladen,
|
|
der eine Viertelstunde spaeter doch noch losfaehrt, weil vor langer
|
|
Zeit einmal eine Schwelle gestreift wurde, waere unangenehmer als
|
|
einer, der gar nicht faehrt. Der naechste echte Anlass nach Ablauf
|
|
der Sperre kommt ohnehin durch.
|
|
|
|
force_once ist davon nicht betroffen: es greift nur, wenn im Fenster
|
|
gar nichts gelaufen ist - dann ist auch keine Sperre aktiv.
|
|
"""
|
|
sperre = int(automatik.get("lockout_secs") or 0)
|
|
letzter = automatik.get("last_run")
|
|
if not sperre or not letzter:
|
|
return False
|
|
return (jetzt - letzter).total_seconds() < sperre
|
|
|
|
@staticmethod
|
|
def lief_heute(automatik, jetzt):
|
|
letzter = automatik.get("last_run")
|
|
return bool(letzter) and letzter.date() == jetzt.date()
|
|
|
|
def protokoll_saeubern(self):
|
|
heute = date.today()
|
|
if self._letzte_saeuberung == heute:
|
|
return
|
|
self._letzte_saeuberung = heute
|
|
try:
|
|
with self.db.cursor() as c:
|
|
c.execute("DELETE FROM automation_log WHERE ts < NOW() - INTERVAL %s DAY",
|
|
(LOG_AUFBEWAHRUNG_TAGE,))
|
|
except Exception as fehler:
|
|
logger.warning("Protokoll nicht aufraeumbar: %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())
|