Files
Smart-Dashboard/restricted/autoActions/autoaction_runner.py
T
adminandClaude Opus 5 131d7cc46d Jalousien: Zu heisst zweifach schwenken, Versand laeuft im Hintergrund
Die Kugelschreiber-Mechanik betrifft jedes Kommando, nicht nur
setOrientation, und sie betrifft auch das Schliessen: ein "Zu" faehrt die
Jalousie zwar herunter, die Lamellen bleiben aber bei etwa 30 % offen
stehen. Dicht wird sie erst durch das zweifache Schwenken.

Die Regel lautet jetzt einheitlich: derselbe Befehl zweimal - erst mit
Neigung 0, dann warten, bis die Fahrt steht, dann mit dem gewuenschten Wert.
Die Position bleibt dabei erhalten, die Jalousie faehrt also nur einmal.

    Zu                     -> [100,0] warten [100,100]
    Neigung 20             -> direkt, unter der Schwelle
    Neigung 80             -> [0] warten [80]
    Position 40 Neigung 50 -> [40,0] warten [40,50]

"Zu" wird dabei nur umgeschrieben, wenn das Geraet Position und Neigung
zusammen setzen kann - der einfache Rollladen "Terasse" hat keine Lamellen
und bekommt weiter das rohe down.

Gewartet wird auf zwei Auskuenfte zusammen, weil einzeln keine traegt:
core:MovingState trug die lange Fahrt (gemessen 61 s), wird bei kurzen
Neigungsfahrten aber nie gesetzt; die Zustandswerte sind die Wahrheit,
zeigen direkt nach dem Kommando aber noch den alten Stand. Fertig heisst:
nichts faehrt mehr, die Ziele stimmen, und es wurde entweder ein "faehrt"
gesehen oder der Vorlauf von acht Sekunden ist um.

Und weil ein "Zu" damit ueber eine Minute dauert, schickt der Runner nicht
mehr selbst: der neue Versand nimmt die Kommandos entgegen und arbeitet sie
in eigenen Faeden ab - je Geraet der Reihe nach, ueber Geraete hinweg
nebeneinander. Sonst haette eine einzige Jalousie die ganze Auswertung fuer
eine Minute angehalten: keine Zeit-Ausloeser, keine Messwerte, und mehrere
Rollladen in einer Automatik haetten sich aufaddiert. Die Datenbank bleibt
dabei im Hauptfaden - die Faeden melden nur ihr Ergebnis zurueck, eine
pymysql-Verbindung ist nicht fuer mehrere Faeden gedacht.

Nachgemessen: Einreihen kehrt sofort zurueck, zwei Kommandos an dasselbe
Geraet laufen nacheinander (6 s, dann 12 s), drei an verschiedene Geraete
gleichzeitig, Fehler kommen zurueck, und die Schleife lief in 20 Sekunden
20 Takte durch.

Die Bedienung im Modal wartet weiterhin - dort sitzt ein Mensch vor einem
Fortschritt und klickt einmal. ajax/room.php hebt dafuer die Zeitgrenze auf
180 Sekunden an, wie es ajax/tahoma.php an derselben Stelle auch tut.

Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2026-08-31 20:29:40 +02:00

849 lines
36 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
Ereignis MQTT-Nachricht -> Wert merken -> im naechsten Takt auswerten
Takt alle tick_seconds: Uhrzeit/Datum/Sonnenzeiten neu rechnen,
HTTP- und Tahoma-Geraete pollen, Regelwerk auf Aenderung pruefen
Pruefen aktiv? (enabled, Wochentag, Ferien/Feiertag, Zeitfenster)
-> Bedingungen auswerten -> steigende Flanke -> Aktionen
Nur die steigende Flanke loest aus: `automations.cond_met` haelt fest, ob die
Bedingung beim letzten Durchlauf schon erfuellt war. Ohne das wuerde
"Temperatur ueber 22 Grad" bei jedem Takt erneut feuern.
Zeit-Ausloeser gibt es in drei Formen: "um 16:30" ist genau in dieser Minute
wahr, "ab 16:30" von da an bis Mitternacht, "vor 16:30" bis dahin.
Ausgeloest wird in allen drei Faellen nur einmal, eben wegen der Flanke -
"ab" ist aber das robustere: faellt der Runner in der einen Minute aus, auf
die "um" zeigt, ist die Automatik fuer den Tag verloren; bei "ab" holt der
naechste Takt es nach. Wer "um" braucht und den Ausfall nicht riskieren
will, hakt zusaetzlich force_once an. Dieselbe Ueberlegung steht hinter den
breiten Zeitfenstern in auto_watering.py.
Beim Sonnenauf- und -untergang traegt der Operator zusaetzlich das
Vorzeichen des Versatzes: "+ 00:30" eine halbe Stunde danach, ">=- 00:30"
ab einer halben Stunde davor, "<+ 00:30" bis eine halbe Stunde danach.
Tabellen siehe homeMesh_automations.sql, Konfiguration siehe config.ini.example.
"""
import argparse
import configparser
import json
import logging
import os
import queue
import sys
import threading
import time
from datetime import date, datetime, timedelta
import pymysql
import requests
import paho.mqtt.client as mqtt
from transports import (HTTPTransport, LogicTransport, MQTTTransport,
TahomaTransport, WLEDTransport)
logger = logging.getLogger("autoaction")
# Wie lange ein Messwert in der Datenbank stehen bleiben darf, bevor er
# aufgefrischt wird. Der Wert dient nur der Anzeige im Editor; jede Nachricht
# sofort zu schreiben waere bei einem gespraechigen Sensor sinnlose Last.
SCHREIB_ABSTAND = timedelta(seconds=60)
# Aelteres im Protokoll interessiert niemanden mehr.
LOG_AUFBEWAHRUNG_TAGE = 30
# ===========================================================================
# Konfiguration
# ===========================================================================
class Config:
def __init__(self, dateiname="config.ini"):
pfad = os.path.join(os.path.dirname(os.path.abspath(__file__)), dateiname)
if not os.path.exists(pfad):
raise FileNotFoundError(
"%s fehlt - config.ini.example kopieren und ausfuellen." % pfad)
self.cfg = configparser.ConfigParser()
self.cfg.read(pfad, encoding="utf-8")
def text(self, sektion, schluessel, vorgabe=""):
return self.cfg.get(sektion, schluessel, fallback=vorgabe).strip()
def zahl(self, sektion, schluessel, vorgabe=0):
try:
return self.cfg.getint(sektion, schluessel, fallback=vorgabe)
except ValueError:
return vorgabe
def ja(self, sektion, schluessel, vorgabe=False):
return self.text(sektion, schluessel, str(vorgabe)).lower() in ("true", "1", "yes", "on")
# ===========================================================================
# Datenbank
# ===========================================================================
def verbinden(config, sektion="database"):
return pymysql.connect(
host=config.text("database", "host", "localhost"),
port=config.zahl("database", "port", 3306),
user=config.text(sektion, "user") or config.text("database", "user"),
password=config.text(sektion, "password") or config.text("database", "password"),
database=config.text(sektion, "database"),
charset="utf8mb4",
cursorclass=pymysql.cursors.DictCursor,
autocommit=True)
class Regelwerk:
"""
Das geladene Abbild der Datenbank: Automatiken mit ihren Bedingungen und
Aktionen, dazu die Messwerte und Kommandos, die sie benutzen.
"""
def __init__(self, automatiken, states, kommandos, signatur):
self.automatiken = automatiken
self.states = states # state_id -> Beschreibung
self.kommandos = kommandos # command_id -> Beschreibung
self.signatur = signatur
@staticmethod
def signatur_lesen(db):
"""
Woran der Runner merkt, dass er neu laden muss. `changed` allein
genuegt nicht: eine geloeschte Automatik veraendert den groessten
Zeitstempel nicht. Deshalb zaehlen die Zeilen mit - auch die der
Geraetetabellen, damit ein Discovery-Lauf ebenfalls durchschlaegt.
"""
with db.cursor() as c:
c.execute("""SELECT (SELECT COUNT(*) FROM automations) AS a,
(SELECT UNIX_TIMESTAMP(MAX(changed)) FROM automations) AS t,
(SELECT COUNT(*) FROM automation_conditions) AS b,
(SELECT COUNT(*) FROM automation_actions) AS c,
(SELECT COUNT(*) FROM actor_states) AS d,
(SELECT COUNT(*) FROM actor_commands) AS e""")
return tuple(sorted(c.fetchone().items()))
@classmethod
def laden(cls, db):
signatur = cls.signatur_lesen(db)
states = {}
with db.cursor() as c:
c.execute("""SELECT s.id, s.state_name, s.url AS state_url, s.value_path,
s.current_value, s.possible_values,
a.url AS actor_url, a.name AS actor_name, t.type
FROM actor_states s
JOIN actors a ON a.id = s.actor_id
LEFT JOIN state_types t ON s.state_type = t.id""")
for row in c.fetchall():
row["type"] = row["type"] or "string"
# Werttabelle als flaches {gesendeter Wert: Bezeichnung}. In
# der Datenbank steht sie als Liste aus Ein-Schluessel-
# Objekten, weil der Editor sie so schon versteht.
row["wertetabelle"] = {}
try:
for eintrag in json.loads(row["possible_values"] or "[]"):
if isinstance(eintrag, dict):
for wert, name in eintrag.items():
row["wertetabelle"][str(wert)] = name
except ValueError:
pass
states[row["id"]] = row
kommandos = {}
with db.cursor() as c:
c.execute("""SELECT k.id, k.command_name, k.command_url,
a.url AS actor_url, a.name AS actor_name
FROM actor_commands k JOIN actors a ON a.id = k.actor_id""")
for row in c.fetchall():
row["params"] = []
kommandos[row["id"]] = row
with db.cursor() as c:
c.execute("""SELECT id, command_id, parameter_name, url FROM command_parameters
ORDER BY command_id, id""")
for row in c.fetchall():
if row["command_id"] in kommandos:
kommandos[row["command_id"]]["params"].append(row)
automatiken = {}
with db.cursor() as c:
c.execute("SELECT * FROM automations WHERE enabled = 1")
for row in c.fetchall():
row["gruppen"] = {}
row["aktionen"] = []
automatiken[row["id"]] = row
with db.cursor() as c:
c.execute("""SELECT * FROM automation_conditions
ORDER BY automation_id, group_no, position, id""")
for row in c.fetchall():
auto = automatiken.get(row["automation_id"])
if auto is None:
continue
if row["state_id"] not in states:
logger.warning("Automatik %s: Messwert %s gibt es nicht mehr",
auto["name"], row["state_id"])
continue
auto["gruppen"].setdefault(row["group_no"], []).append(row)
with db.cursor() as c:
c.execute("""SELECT a.id, a.automation_id, a.command_id, p.parameter_id, p.value
FROM automation_actions a
LEFT JOIN automation_action_params p ON p.action_id = a.id
ORDER BY a.automation_id, a.position, a.id""")
gesammelt = {}
for row in c.fetchall():
auto = automatiken.get(row["automation_id"])
if auto is None:
continue
aktion = gesammelt.get(row["id"])
if aktion is None:
aktion = {"command_id": row["command_id"], "werte": {}}
gesammelt[row["id"]] = aktion
auto["aktionen"].append(aktion)
if row["parameter_id"] is not None:
aktion["werte"][row["parameter_id"]] = row["value"]
logger.info("Regelwerk geladen: %d Automatiken, %d Messwerte, %d Kommandos",
len(automatiken), len(states), len(kommandos))
return cls(automatiken, states, kommandos, signatur)
# ===========================================================================
# Auswertung
# ===========================================================================
def minuten(text):
""""16:30" oder "16:30:00" als Minuten seit Mitternacht."""
teile = str(text).strip().split(":")
return int(teile[0]) * 60 + int(teile[1])
def als_datum(text):
for form in ("%d.%m.%Y", "%Y-%m-%d", "%d.%m.%y"):
try:
return datetime.strptime(str(text).strip(), form).date()
except ValueError:
continue
return None
def als_zeitpunkt(text):
"""Datum mit Uhrzeit. Das Web-Feld liefert "2026-08-30T16:30"."""
roh = str(text).strip().replace("T", " ")
for form in ("%Y-%m-%d %H:%M:%S", "%Y-%m-%d %H:%M",
"%d.%m.%Y %H:%M:%S", "%d.%m.%Y %H:%M"):
try:
return datetime.strptime(roh, form)
except ValueError:
continue
tag = als_datum(roh.split(" ")[0])
return datetime.combine(tag, datetime.min.time()) if tag else None
WAHR = {"true", "1", "on", "ja", "yes", "an"}
def bedingung_erfuellt(bedingung, state, wert, jetzt):
"""
Ein einzelner Vergleich. `wert` ist der aktuelle Messwert als Text, so wie
er vom Geraet kam; `bedingung["value"]` die eingestellte Schwelle.
Unbekannter Wert heisst nicht erfuellt - lieber nicht schalten als auf
Verdacht schalten.
"""
if wert is None or wert == "":
return False
typ = state["type"]
op = bedingung["operator"]
soll = bedingung["value"]
try:
if typ in ("integer", "float"):
ist_z, soll_z = float(str(wert).replace(",", ".")), float(str(soll).replace(",", "."))
if op == "=": return ist_z == soll_z
if op == "!=": return ist_z != soll_z
if op == ">": return ist_z > soll_z
if op == "<": return ist_z < soll_z
if op == ">=": return ist_z >= soll_z
if op == "<=": return ist_z <= soll_z
return False
if typ == "time":
ist_m, soll_m = minuten(wert), minuten(soll)
if op == "=": return ist_m == soll_m
if op == "<": return ist_m < soll_m
return ist_m >= soll_m
if typ == "deltatime":
# Der Messwert ist der Sonnenauf- bzw. -untergang, die Schwelle
# ein Versatz. Das Vorzeichen steckt im Operator, der Vergleich
# davor: ">=-" heisst "ab einer halben Stunde davor".
versatz = minuten(soll)
ziel = minuten(wert) + (-versatz if op.endswith("-") else versatz)
# Ueber den Tagesrand wird gerechnet, nicht abgeschnitten: "sechs
# Stunden vor Sonnenaufgang" ist eine gewollte Angabe und landet
# dann eben am Vorabend. Abgeschnitten waeren solche Faelle gar
# nicht mehr formulierbar.
#
# Verglichen wird die Uhrzeit innerhalb des Tages. Ein Ziel
# jenseits von Mitternacht gilt also als diese Uhrzeit am selben
# Tag - bei "um" ist das genau der gemeinte Zeitpunkt, bei "ab"
# und "vor" verschiebt sich der wahre Bereich entsprechend.
ziel %= 1440
jetzt_m = jetzt.hour * 60 + jetzt.minute
if op.startswith(">="): return jetzt_m >= ziel
if op.startswith("<"): return jetzt_m < ziel
return jetzt_m == ziel
if typ in ("date", "datetime"):
wandeln = als_datum if typ == "date" else als_zeitpunkt
ist_d, soll_d = wandeln(wert), wandeln(soll)
if ist_d is None or soll_d is None:
return False
if op == "=": return ist_d == soll_d
if op == "<": return ist_d < soll_d
return ist_d >= soll_d
if typ == "bool":
ist_b = str(wert).strip().lower() in WAHR
soll_b = str(soll).strip().lower() in WAHR
return ist_b == soll_b if op == "=" else ist_b != soll_b
# string und alles Uebrige
if op == "=":
return str(wert).strip() == str(soll).strip()
return str(wert).strip() != str(soll).strip()
except (ValueError, IndexError) as fehler:
logger.debug("Vergleich %s %s %s nicht moeglich: %s", wert, op, soll, fehler)
return False
def gruppen_erfuellt(automatik, regelwerk, werte, jetzt):
"""
Gleiche group_no = UND, verschiedene = ODER. Eine Automatik ohne
Bedingungen loest nie aus - sonst wuerde sie nach einem Discovery-Lauf,
der ihren Messwert entfernt hat, ploetzlich dauernd feuern.
"""
if not automatik["gruppen"]:
return False
for bedingungen in automatik["gruppen"].values():
if all(bedingung_erfuellt(b, regelwerk.states[b["state_id"]],
werte.get(b["state_id"]), jetzt)
for b in bedingungen):
return True
return False
def im_zeitfenster(jetzt, von, bis):
"""von > bis heisst: das Fenster reicht ueber Mitternacht."""
m = jetzt.hour * 60 + jetzt.minute
a, b = minuten(von), minuten(bis)
return a <= m <= b if a <= b else (m >= a or m <= b)
def tag_passt(automatik, jetzt, kalender):
if not automatik["weekdays"] & (1 << jetzt.weekday()):
return False
if kalender["feiertag"] and not automatik["on_holiday"]:
return False
if kalender["ferien"] and not automatik["on_vacation"]:
return False
return True
# ===========================================================================
# Versand
# ===========================================================================
class Versand:
"""
Schickt Kommandos im Hintergrund - je Geraet der Reihe nach, ueber
Geraete hinweg nebeneinander.
Ohne das blockiert ein einziges Kommando die ganze Auswertung: eine
Jalousie zuzufahren dauert ueber eine Minute, weil zwischen den beiden
Schwenkbefehlen auf das Ende der Fahrt gewartet werden muss (siehe
TahomaTransport). So lange kaeme kein Zeit-Ausloeser mehr durch, kein
Messwert wuerde zurueckgeschrieben, und mehrere Rollladen in einer
Automatik wuerden sich aufaddieren.
Je Geraet eine Warteschlange mit einem eigenen Faden: zwei Kommandos an
dieselbe Jalousie duerfen sich nicht ueberholen - "Neigung 0" und
"Neigung 100" sind sonst wirkungslos oder vertauscht -, zwei Kommandos an
verschiedene Jalousien duerfen ruhig gleichzeitig laufen.
Die Datenbank bleibt aussen vor: die Faeden melden ihr Ergebnis nur
zurueck, geschrieben wird im Hauptfaden. Eine pymysql-Verbindung ist
nicht fuer mehrere Faeden gedacht.
"""
def __init__(self):
self.warteschlangen = {} # actor_url -> Queue
self.ergebnisse = queue.Queue()
self.laeuft = True
def einreihen(self, automation_id, beschreibung, transport, auftrag):
schlange = self.warteschlangen.get(auftrag["actor_url"])
if schlange is None:
schlange = queue.Queue()
self.warteschlangen[auftrag["actor_url"]] = schlange
faden = threading.Thread(target=self._arbeiten, args=(schlange,),
name="versand", daemon=True)
faden.start()
schlange.put((automation_id, beschreibung, transport, auftrag))
def _arbeiten(self, schlange):
while self.laeuft:
posten = schlange.get()
if posten is None:
return
automation_id, beschreibung, transport, auftrag = posten
try:
transport.senden(auftrag)
self.ergebnisse.put((automation_id, beschreibung, None))
except Exception as fehler:
self.ergebnisse.put((automation_id, beschreibung, str(fehler)))
finally:
schlange.task_done()
def abholen(self):
"""Alles, was seit dem letzten Mal fertig geworden ist."""
fertig = []
while True:
try:
fertig.append(self.ergebnisse.get_nowait())
except queue.Empty:
return fertig
def offen(self):
return sum(s.unfinished_tasks for s in self.warteschlangen.values())
def beenden(self):
self.laeuft = False
for schlange in self.warteschlangen.values():
schlange.put(None)
# ===========================================================================
# Der Runner
# ===========================================================================
class Runner:
def __init__(self, config, dry_run=False):
self.config = config
self.dry_run = dry_run or config.ja("runner", "dry_run")
self.db = verbinden(config)
self.solar = verbinden(config, "solar")
self.werte = {} # state_id -> letzter bekannter Wert
self.geschrieben = {} # state_id -> wann zuletzt in die DB
self.war_aktiv = {} # automation_id -> war im Zeitfenster
self.lief_im_fenster = {} # automation_id -> hat im Fenster ausgeloest
self._sonne = (None, "00:00", "00:00") # (datum, aufgang, untergang)
self._kalender = (None, {"feiertag": False, "ferien": False})
self._letzte_saeuberung = None
self.versand = Versand()
self.mqtt = self._mqtt_verbinden()
self.transporte = [
MQTTTransport(self.mqtt, self.dry_run),
WLEDTransport(requests, dry_run=self.dry_run),
HTTPTransport(requests, dry_run=self.dry_run),
TahomaTransport(requests, config.text("tahoma", "pin"),
config.text("tahoma", "token"),
config.zahl("tahoma", "timeout", 10), self.dry_run),
LogicTransport(self.sonnenzeiten),
]
self.regelwerk = None
# --- Aufbau ----------------------------------------------------------
def _mqtt_verbinden(self):
client = mqtt.Client(mqtt.CallbackAPIVersion.VERSION2,
client_id=self.config.text("mqtt", "client_id", "autoaction_runner"))
benutzer = self.config.text("mqtt", "username")
if benutzer:
client.username_pw_set(benutzer, self.config.text("mqtt", "password"))
client.on_message = self._mqtt_nachricht
client.connect(self.config.text("mqtt", "broker", "localhost"),
self.config.zahl("mqtt", "port", 1883), 60)
client.loop_start()
return client
def _mqtt_nachricht(self, client, userdata, nachricht):
# Der Client laeuft schon, waehrend __init__ noch die Transporte baut.
for transport in getattr(self, "transporte", []):
if isinstance(transport, MQTTTransport):
transport.nachricht(nachricht.topic, nachricht.payload)
def transport_fuer(self, actor_url):
for transport in self.transporte:
if transport.passt(actor_url):
return transport
return None
def regelwerk_laden(self):
self.regelwerk = Regelwerk.laden(self.db)
# Welche Geraete Position und Neigung zusammen koennen. Nur die haben
# die Kugelschreiber-Mechanik, und nur bei ihnen wird ein "Zu" zum
# zweifachen Schwenken - siehe TahomaTransport.
kombi = {k["actor_url"] for k in self.regelwerk.kommandos.values()
if k["command_url"] == "setClosureAndOrientation"}
for transport in self.transporte:
if isinstance(transport, TahomaTransport):
transport.kombigeraete_setzen(kombi)
# Jeder Transport bekommt die Messwerte, fuer die er zustaendig ist.
for transport in self.transporte:
passende = [{"id": s["id"], "actor_url": s["actor_url"],
"state_url": s["state_url"], "value_path": s["value_path"],
"wertetabelle": s["wertetabelle"]}
for s in self.regelwerk.states.values()
if transport.passt(s["actor_url"])]
transport.zustaende_anmelden(passende)
# Der zuletzt bekannte Wert aus der Datenbank ist besser als gar
# keiner: nach einem Neustart steht sonst jede Bedingung auf "unklar",
# bis das Geraet zufaellig etwas schickt.
for s in self.regelwerk.states.values():
if s["id"] not in self.werte and s["current_value"] is not None:
self.werte[s["id"]] = s["current_value"]
# --- Umgebung --------------------------------------------------------
def sonnenzeiten(self):
"""Aus solarLog.daylight, einmal je Tag geholt."""
heute = date.today()
if self._sonne[0] == heute:
return self._sonne[1], self._sonne[2]
auf, unter = "00:00", "00:00"
try:
with self.solar.cursor() as c:
c.execute("SELECT sunrise, sunset FROM daylight WHERE date = %s", (heute,))
zeile = c.fetchone()
if zeile:
auf = self._als_uhrzeit(zeile["sunrise"])
unter = self._als_uhrzeit(zeile["sunset"])
else:
logger.warning("Kein Eintrag in daylight fuer %s", heute)
except Exception as fehler:
logger.warning("Sonnenzeiten nicht lesbar: %s", fehler)
self._sonne = (heute, auf, unter)
return auf, unter
@staticmethod
def _als_uhrzeit(wert):
"""TIME-Spalten liefert pymysql als timedelta, nicht als Text."""
if isinstance(wert, timedelta):
minute = int(wert.total_seconds()) // 60
return "%02d:%02d" % (minute // 60 % 24, minute % 60)
return str(wert)[:5]
def kalender(self):
"""Ferien und Feiertage von heute, einmal je Tag geholt."""
heute = date.today()
if self._kalender[0] == heute:
return self._kalender[1]
stand = {"feiertag": False, "ferien": False}
try:
with self.db.cursor() as c:
c.execute("SELECT holiday, vacation FROM calendar_days WHERE date = %s", (heute,))
zeile = c.fetchone()
if zeile:
stand = {"feiertag": bool(zeile["holiday"]), "ferien": bool(zeile["vacation"])}
except Exception as fehler:
logger.warning("Kalender nicht lesbar: %s", fehler)
self._kalender = (heute, stand)
return stand
# --- Werte -----------------------------------------------------------
def werte_einsammeln(self, mit_pollen):
neu = {}
for transport in self.transporte:
if isinstance(transport, MQTTTransport) or mit_pollen:
neu.update(transport.zustaende_lesen())
if neu:
self.werte.update(neu)
self.werte_zurueckschreiben(neu)
return neu
def werte_zurueckschreiben(self, neu):
"""
current_value nachfuehren, damit der Editor im Browser aktuelle Zahlen
zeigt ("Temperatur (= 21,4 °C)"). Gedrosselt, sonst schreibt ein
gespraechiger Sensor die Tabelle im Sekundentakt voll.
"""
jetzt = datetime.now()
faellig = [(w, i) for i, w in neu.items()
if self.geschrieben.get(i, datetime.min) + SCHREIB_ABSTAND <= jetzt]
if not faellig:
return
try:
with self.db.cursor() as c:
c.executemany("UPDATE actor_states SET current_value = %s WHERE id = %s", faellig)
for _, state_id in faellig:
self.geschrieben[state_id] = jetzt
except Exception as fehler:
logger.warning("current_value nicht schreibbar: %s", fehler)
# --- Ausfuehren ------------------------------------------------------
def ausloesen(self, automatik, anlass):
"""
Die Aktionen einer Automatik in den Versand geben.
Geschickt wird im Hintergrund - eine Jalousie zuzufahren dauert ueber
eine Minute, und so lange darf die Auswertung nicht stehen. Was dabei
schiefgeht, kommt spaeter ueber ergebnisse_verbuchen() ins Protokoll.
"""
fehlerText = []
eingereiht = 0
for aktion in automatik["aktionen"]:
kommando = self.regelwerk.kommandos.get(aktion["command_id"])
if kommando is None:
fehlerText.append("Kommando %s gibt es nicht mehr" % aktion["command_id"])
logger.error("%s: Kommando %s gibt es nicht mehr",
automatik["name"], aktion["command_id"])
continue
transport = self.transport_fuer(kommando["actor_url"])
if transport is None:
fehlerText.append("Kein Transport fuer %s" % kommando["actor_url"])
logger.error("%s: fuer %s (%s) gibt es keinen Transport",
automatik["name"], kommando["actor_name"], kommando["actor_url"])
continue
auftrag = {
"actor_url": kommando["actor_url"],
"command_url": kommando["command_url"],
"params": [{"url": p["url"], "name": p["parameter_name"],
"wert": aktion["werte"].get(p["id"], "")}
for p in kommando["params"]],
}
beschreibung = "%s: %s" % (kommando["actor_name"], kommando["command_name"])
logger.info("%s: %s (unterwegs)", automatik["name"], beschreibung)
self.versand.einreihen(automatik["id"], beschreibung, transport, auftrag)
eingereiht += 1
ergebnis = "error" if fehlerText else anlass
detail = "; ".join(fehlerText)[:255] if fehlerText \
else ("%d Kommando(s) unterwegs" % eingereiht)
try:
with self.db.cursor() as c:
c.execute("""UPDATE automations SET last_run = NOW(), changed = changed
WHERE id = %s""", (automatik["id"],))
c.execute("""INSERT INTO automation_log (automation_id, result, detail)
VALUES (%s, %s, %s)""", (automatik["id"], ergebnis, detail))
except Exception as fehler:
logger.warning("Protokoll nicht schreibbar: %s", fehler)
automatik["last_run"] = datetime.now()
self.lief_im_fenster[automatik["id"]] = True
def ergebnisse_verbuchen(self):
"""
Was der Versand inzwischen erledigt hat ins Protokoll schreiben.
Nur Fehlschlaege bekommen eine eigene Zeile - der Lauf selbst steht
schon drin, und ein Protokoll, das jedes gelungene Kommando einzeln
auffuehrt, findet niemand mehr etwas darin.
"""
for automation_id, beschreibung, fehler in self.versand.abholen():
if fehler is None:
logger.debug("erledigt: %s", beschreibung)
continue
logger.error("%s konnte nicht geschickt werden: %s", beschreibung, fehler)
try:
with self.db.cursor() as c:
c.execute("""INSERT INTO automation_log (automation_id, result, detail)
VALUES (%s, 'error', %s)""",
(automation_id, ("%s: %s" % (beschreibung, fehler))[:255]))
except Exception as schreibfehler:
logger.warning("Protokoll nicht schreibbar: %s", schreibfehler)
def flanke_merken(self, automatik, erfuellt):
if bool(automatik["cond_met"]) == bool(erfuellt):
return
automatik["cond_met"] = 1 if erfuellt else 0
try:
with self.db.cursor() as c:
c.execute("""UPDATE automations SET cond_met = %s, changed = changed
WHERE id = %s""", (automatik["cond_met"], automatik["id"]))
except Exception as fehler:
logger.warning("cond_met nicht schreibbar: %s", fehler)
def durchlauf(self):
jetzt = datetime.now()
kalender = self.kalender()
for automatik in self.regelwerk.automatiken.values():
aktiv = (tag_passt(automatik, jetzt, kalender)
and im_zeitfenster(jetzt, automatik["window_from"], automatik["window_to"]))
vorher_aktiv = self.war_aktiv.get(automatik["id"], aktiv)
if aktiv:
if not vorher_aktiv:
self.lief_im_fenster[automatik["id"]] = False
erfuellt = gruppen_erfuellt(automatik, self.regelwerk, self.werte, jetzt)
if erfuellt and not automatik["cond_met"]:
if self.gesperrt(automatik, jetzt):
logger.debug("%s: Flanke faellt in die Sperrzeit, uebersprungen",
automatik["name"])
else:
self.ausloesen(automatik, "fired")
self.flanke_merken(automatik, erfuellt)
else:
# Das Fenster ist gerade zugegangen. Wer "auf jeden Fall"
# angehakt hat, bekommt jetzt seinen Lauf - aber nur, wenn in
# diesem Fenster noch keiner stattgefunden hat. Nach einem
# Neustart mitten im Fenster weiss der Runner das nicht mehr
# aus dem Speicher, deshalb zaehlt zusaetzlich last_run.
if vorher_aktiv and automatik["force_once"] and not self.lief_im_fenster.get(automatik["id"]):
if not self.lief_heute(automatik, jetzt):
logger.info("%s: Zeitfenster vorbei, wird trotzdem ausgefuehrt",
automatik["name"])
self.ausloesen(automatik, "forced")
# Ausserhalb des Fensters die Flanke zuruecksetzen, sonst
# koennte sie im naechsten Fenster nicht mehr steigen.
self.flanke_merken(automatik, False)
self.lief_im_fenster[automatik["id"]] = False
self.war_aktiv[automatik["id"]] = aktiv
@staticmethod
def gesperrt(automatik, jetzt):
"""
Liegt die letzte Ausloesung noch innerhalb der Sperrzeit?
Gegen Messwerte, die um die Schwelle pendeln: "Temperatur > 22" bei
22,1 / 21,9 / 22,1 Grad ist jedes Mal eine echte steigende Flanke, und
ueber MQTT koennen die Werte im Sekundentakt hereinkommen.
Die Flanke wird dabei verworfen und nicht aufgehoben. Ein Rollladen,
der eine Viertelstunde spaeter doch noch losfaehrt, weil vor langer
Zeit einmal eine Schwelle gestreift wurde, waere unangenehmer als
einer, der gar nicht faehrt. Der naechste echte Anlass nach Ablauf
der Sperre kommt ohnehin durch.
force_once ist davon nicht betroffen: es greift nur, wenn im Fenster
gar nichts gelaufen ist - dann ist auch keine Sperre aktiv.
"""
sperre = int(automatik.get("lockout_secs") or 0)
letzter = automatik.get("last_run")
if not sperre or not letzter:
return False
return (jetzt - letzter).total_seconds() < sperre
@staticmethod
def lief_heute(automatik, jetzt):
letzter = automatik.get("last_run")
return bool(letzter) and letzter.date() == jetzt.date()
def protokoll_saeubern(self):
heute = date.today()
if self._letzte_saeuberung == heute:
return
self._letzte_saeuberung = heute
try:
with self.db.cursor() as c:
c.execute("DELETE FROM automation_log WHERE ts < NOW() - INTERVAL %s DAY",
(LOG_AUFBEWAHRUNG_TAGE,))
except Exception as fehler:
logger.warning("Protokoll nicht aufraeumbar: %s", fehler)
# --- Hauptschleife ---------------------------------------------------
def laufen(self, nur_einmal=False):
self.regelwerk_laden()
takt = self.config.zahl("runner", "tick_seconds", 30)
poll = self.config.zahl("runner", "poll_seconds", 60)
neuladen = self.config.zahl("runner", "reload_seconds", 30)
letztes_pollen = 0.0
letztes_pruefen = 0.0
while True:
jetzt = time.monotonic()
mit_pollen = jetzt - letztes_pollen >= poll
if mit_pollen:
letztes_pollen = jetzt
self.werte_einsammeln(mit_pollen)
self.durchlauf()
self.ergebnisse_verbuchen()
self.protokoll_saeubern()
if jetzt - letztes_pruefen >= neuladen:
letztes_pruefen = jetzt
try:
if Regelwerk.signatur_lesen(self.db) != self.regelwerk.signatur:
logger.info("Regelwerk hat sich geaendert, wird neu geladen")
self.regelwerk_laden()
except Exception as fehler:
logger.warning("Regelwerk nicht pruefbar: %s", fehler)
if nur_einmal:
return
time.sleep(takt)
def beenden(self):
offen = self.versand.offen()
if offen:
logger.info("%d Kommando(s) noch unterwegs - wird nicht abgewartet", offen)
self.versand.beenden()
try:
self.mqtt.loop_stop()
self.mqtt.disconnect()
except Exception:
pass
# ===========================================================================
def main():
parser = argparse.ArgumentParser(description="Fuehrt die Automatiken aus dem Web-UI aus.")
parser.add_argument("--dry-run", action="store_true",
help="nichts wirklich schalten, nur protokollieren")
parser.add_argument("--once", action="store_true",
help="einen einzigen Durchlauf, dann beenden")
parser.add_argument("--verbose", action="store_true", help="DEBUG-Ausgaben")
parser.add_argument("--config", default="config.ini")
args = parser.parse_args()
config = Config(args.config)
logging.basicConfig(
level=logging.DEBUG if args.verbose else getattr(
logging, config.text("runner", "log_level", "INFO").upper(), logging.INFO),
format="%(asctime)s - %(levelname)s - %(message)s")
runner = Runner(config, args.dry_run)
if runner.dry_run:
logger.info("Probelauf: es wird nichts geschaltet.")
try:
runner.laufen(args.once)
except KeyboardInterrupt:
logger.info("Abbruch durch Benutzer")
finally:
runner.beenden()
return 0
if __name__ == "__main__":
sys.exit(main())