AutoAction-Runner: Automatiken ausfuehren
Bisher konnte man Automatiken nur anlegen - ausgefuehrt hat sie niemand. restricted/autoActions/autoaction_runner.py holt das nach. Dauerlaeufer statt Cronjob, aus zwei Gruenden: Schwellwert-Ausloeser sollen greifen, wenn die MQTT-Nachricht hereinkommt, und actor_states.current_value wird sonst von niemandem fortgeschrieben - beim Discovery einmal gesetzt und danach nie wieder. Ein zustandsloser Lauf haette gar nichts, womit er vergleichen koennte. Der Runner pflegt den Wert nebenbei mit, wovon auch der Editor profitiert: er zeigt neben jedem Messwert den aktuellen Stand. Ausgeloest wird nur auf der steigenden Flanke (automations.cond_met), sonst wuerde "Temperatur ueber 22 Grad" bei jedem Takt erneut feuern. Aus demselben Grund heissen Zeit-Ausloeser jetzt "ab 16:30" statt "gleich 16:30": ein Gleichheitsvergleich waere nur in einer einzigen Minute wahr, und ein Ausfall in genau dieser Minute kostet den ganzen Tag. Dieselbe Ueberlegung steht hinter den breiten Zeitfenstern in auto_watering.py. Der Weg zum Geraet haengt an der URL des Aktors: mqtt:// abonniert und publiziert, http:// pollt und haengt Parameter an, io:// spricht mit der Tahoma-Box, Logic rechnet Uhrzeit, Datum und Sonnenzeiten. Alle vier stehen in transports.py; eine fuenfte Geraeteart ist eine weitere Klasse mit passt(), zustaende_lesen() und senden(). fetch_calendar.py fuellt calendar_days aus openholidaysapi.org, damit "in den Ferien" und "an Feiertagen" eine Grundlage haben. Einmal jaehrlich per Cron. Die Uebersicht blendet auf schmalen Schirmen Spalten aus, statt sie wegzuschieben - auf dem Handy waren die Knoepfe zum Pausieren und Loeschen sonst nur per Seitwaertsscrollen erreichbar. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -0,0 +1,667 @@
|
||||
#!/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 heissen "ab 16:30" und nicht "gleich 16:30" - damit bleibt die
|
||||
Bedingung bis Mitternacht wahr und ein verpasster Takt kostet nicht den
|
||||
ganzen Tag. Ausgeloest wird trotzdem nur einmal, eben wegen der Flanke.
|
||||
Dieselbe Ueberlegung steht hinter den breiten Zeitfenstern in
|
||||
auto_watering.py.
|
||||
|
||||
Tabellen siehe homeMesh_automations.sql, Konfiguration siehe config.ini.example.
|
||||
"""
|
||||
|
||||
import argparse
|
||||
import configparser
|
||||
import logging
|
||||
import os
|
||||
import sys
|
||||
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)
|
||||
|
||||
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.current_value,
|
||||
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"
|
||||
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)
|
||||
return ist_m >= soll_m if op == ">=" else ist_m < soll_m
|
||||
|
||||
if typ == "deltatime":
|
||||
# Der Messwert ist der Sonnenauf- bzw. -untergang, die Schwelle
|
||||
# ein Versatz davor oder danach. Verglichen wird gegen die Uhr.
|
||||
versatz = minuten(soll)
|
||||
ziel = minuten(wert) + (versatz if op == "+" else -versatz)
|
||||
return jetzt.hour * 60 + jetzt.minute >= ziel
|
||||
|
||||
if typ == "date":
|
||||
ist_d, soll_d = als_datum(wert), als_datum(soll)
|
||||
if ist_d is None or soll_d is None:
|
||||
return False
|
||||
return ist_d >= soll_d if op == ">=" else ist_d < soll_d
|
||||
|
||||
if typ == "datetime":
|
||||
ist_d, soll_d = als_zeitpunkt(wert), als_zeitpunkt(soll)
|
||||
if ist_d is None or soll_d is None:
|
||||
return False
|
||||
return ist_d >= soll_d if op == ">=" else 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
|
||||
|
||||
|
||||
# ===========================================================================
|
||||
# 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.mqtt = self._mqtt_verbinden()
|
||||
self.transporte = [
|
||||
MQTTTransport(self.mqtt, 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)
|
||||
# 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"]}
|
||||
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):
|
||||
fehlerText = []
|
||||
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"]],
|
||||
}
|
||||
try:
|
||||
transport.senden(auftrag)
|
||||
logger.info("%s: %s -> %s", automatik["name"], kommando["actor_name"],
|
||||
kommando["command_name"])
|
||||
except Exception as fehler:
|
||||
fehlerText.append("%s: %s" % (kommando["command_name"], fehler))
|
||||
logger.error("%s: %s konnte nicht geschickt werden: %s",
|
||||
automatik["name"], kommando["command_name"], fehler)
|
||||
|
||||
ergebnis = "error" if fehlerText else anlass
|
||||
detail = "; ".join(fehlerText)[:255]
|
||||
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 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"]:
|
||||
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 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.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):
|
||||
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())
|
||||
Reference in New Issue
Block a user