Bei Uhrzeit und Datum gab es nur "ab" und "vor". Der haeufigste Fall - eine Aktion genau um 16:30 - liess sich damit gar nicht ausdruecken. "um" ist jetzt der erste Eintrag und damit die Vorbelegung. Der Hinweis von frueher bleibt trotzdem richtig und steht in der README: "um" trifft nur eine einzige Minute, faellt der Runner ausgerechnet in dieser Minute aus, ist die Automatik fuer den Tag verloren. "ab" holt der naechste Takt nach. Wer beides will, hakt force_once an. Beim Sonnenauf- und -untergang fehlte die halbe Tabelle: der Versatz kann davor oder danach liegen, und verglichen werden kann davor, danach oder genau. Statt zwei gibt es jetzt sechs Operatoren - "+", "-", "ab +", "ab -", "vor +", "vor -". Das Vorzeichen steckt im Operator, weil der Editor eine einzige Auswahlliste zeigt und nicht zwei Bedienelemente fuer eine Angabe. Rutscht ein Versatz rechnerisch ueber den Tagesrand (Sonnenaufgang 05:34 minus sechs Stunden), bleibt es beim Tagesrand statt auf die andere Seite von Mitternacht zu springen - "kurz vor Sonnenaufgang" soll nicht ploetzlich gestern abend bedeuten. Der Editor brauchte keine Aenderung: er liest Wert und Beschriftung der Operatoren aus dem Geraetekatalog, statt sie selbst zu kennen. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
684 lines
29 KiB
Python
684 lines
29 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 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)
|
|
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)
|
|
# Ein Versatz kann rechnerisch ueber den Tagesrand rutschen
|
|
# (Sonnenaufgang 05:34 minus sechs Stunden). Dann bleibt es beim
|
|
# Tagesrand, statt auf die andere Seite von Mitternacht zu
|
|
# springen - "kurz vor Sonnenaufgang" soll nicht ploetzlich
|
|
# gestern abend bedeuten.
|
|
ziel = max(0, min(1439, ziel))
|
|
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
|
|
|
|
|
|
# ===========================================================================
|
|
# 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())
|