#!/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) #: Ferien und Feiertage sind dreiwertig, nicht ja/nein. NIE, EGAL, ZUSAETZLICH = 0, 1, 2 def gemeinter_tag(automatik, jetzt): """ Fuer welchen Tag gelten Wochentage, Ferien und Feiertage? Normalerweise fuer heute. Mit `next_day` fuer morgen - die Vorabend-Form: "Kinderrollos zu, wenn morgen Schule ist" heisst Mo-Fr, Ferien nie, Feiertage nie, und der Runner schaut dafuer auf den folgenden Tag. Mit dem heutigen Tag liess sich das nur annaehern (So-Do, heute keine Ferien), und das ging am letzten Ferientag, am Abend vor einem Feiertag und am Abend eines Feiertags daneben. Nur der Rahmen verschiebt sich. Uhrzeit, Zeitfenster und "einmal am Tag" bleiben beim heutigen Tag - die Automatik laeuft ja heute Abend. Ein Zeitfenster ueber Mitternacht meint nach Mitternacht deshalb schon den uebernaechsten Tag; fuer eine Vorabend-Regel ist das kein sinnvoller Fall. """ heute = jetzt.date() return heute + timedelta(days=1) if automatik.get("next_day") else heute def tag_passt(automatik, tag, kalender): """ Faellt der Tag in den Rahmen der Automatik? `tag` ist heute oder morgen, siehe gemeinter_tag(); `kalender` gehoert zu genau diesem Tag. Die Wochentage sind eine Maske, Ferien und Feiertage haben je drei Zustaende. Zwei davon gab es immer: EGAL (der Tag aendert nichts, die Vorgabe) und NIE (an solchen Tagen laeuft die Automatik nicht - so ist "werktags" gebaut). ZUSAETZLICH ist der dritte und zaehlt wie ein passender Wochentag. Der Grund fuer den dritten: "Wochenenden und Feiertage" liess sich vorher gar nicht schreiben. Die Wochentagsmaske kennt nur Samstag und Sonntag, und ein Feiertag am Dienstag ist eben ein Dienstag. Mit ZUSAETZLICH genuegt Sa+So angehakt und "Feiertage: zusaetzlich" - der Dienstag kommt dann ueber den Kalender herein. Reihenfolge: erst wird geoeffnet, dann gesperrt. Wer in den Ferien nie laufen soll und an Feiertagen zusaetzlich, laeuft an einem Feiertag in den Ferien nicht - ein Verbot schlaegt eine Erweiterung. Anders herum liesse sich "nie" nicht mehr verlassen. """ passt = bool(automatik["weekdays"] & (1 << tag.weekday())) if kalender["feiertag"] and automatik["on_holiday"] == ZUSAETZLICH: passt = True if kalender["ferien"] and automatik["on_vacation"] == ZUSAETZLICH: passt = True if not passt: return False if kalender["feiertag"] and automatik["on_holiday"] == NIE: return False if kalender["ferien"] and automatik["on_vacation"] == NIE: return False return True # =========================================================================== # Versand # =========================================================================== class Versand: """ Schickt Kommandos im Hintergrund - je Geraet der Reihe nach, ueber Geraete hinweg nebeneinander. Ohne das blockiert ein einziges Kommando die ganze Auswertung: eine Jalousie zuzufahren dauert ueber eine Minute, weil zwischen den beiden Schwenkbefehlen auf das Ende der Fahrt gewartet werden muss (siehe TahomaTransport). So lange kaeme kein Zeit-Ausloeser mehr durch, kein Messwert wuerde zurueckgeschrieben, und mehrere Rollladen in einer Automatik wuerden sich aufaddieren. Je Geraet eine Warteschlange mit einem eigenen Faden: zwei Kommandos an dieselbe Jalousie duerfen sich nicht ueberholen - "Neigung 0" und "Neigung 100" sind sonst wirkungslos oder vertauscht -, zwei Kommandos an verschiedene Jalousien duerfen ruhig gleichzeitig laufen. Die Datenbank bleibt aussen vor: die Faeden melden ihr Ergebnis nur zurueck, geschrieben wird im Hauptfaden. Eine pymysql-Verbindung ist nicht fuer mehrere Faeden gedacht. """ def __init__(self): self.warteschlangen = {} # actor_url -> Queue self.ergebnisse = queue.Queue() self.messwerte = queue.Queue() self.laeuft = True def einreihen(self, automation_id, beschreibung, transport, auftrag): schlange = self.warteschlangen.get(auftrag["actor_url"]) if schlange is None: schlange = queue.Queue() self.warteschlangen[auftrag["actor_url"]] = schlange faden = threading.Thread(target=self._arbeiten, args=(schlange,), name="versand", daemon=True) faden.start() schlange.put((automation_id, beschreibung, transport, auftrag)) def _arbeiten(self, schlange): while self.laeuft: posten = schlange.get() if posten is None: return automation_id, beschreibung, transport, auftrag = posten try: transport.senden(auftrag) self.ergebnisse.put((automation_id, beschreibung, None)) # Was das Geraet danach meldet, gleich mitnehmen. Bei Tahoma # laegen sonst bis zu fuenf Minuten zwischen der Fahrt und # dem neuen Stand in der Tabelle; die anderen Transporte # geben hier nichts zurueck. Dass das den Faden aufhaelt, ist # gewollt: das naechste Kommando an dieselbe Jalousie darf # ohnehin erst nach der Fahrt kommen. werte = transport.nachlesen(auftrag["actor_url"]) if werte: self.messwerte.put(werte) except Exception as fehler: self.ergebnisse.put((automation_id, beschreibung, str(fehler))) finally: schlange.task_done() def abholen(self): """Alles, was seit dem letzten Mal fertig geworden ist.""" fertig = [] while True: try: fertig.append(self.ergebnisse.get_nowait()) except queue.Empty: return fertig def messwerte_abholen(self): """Was die Versandfaeden nach ihren Kommandos abgelesen haben.""" neu = {} while True: try: neu.update(self.messwerte.get_nowait()) except queue.Empty: return neu def offen(self): return sum(s.unfinished_tasks for s in self.warteschlangen.values()) def beenden(self): self.laeuft = False for schlange in self.warteschlangen.values(): schlange.put(None) class Sammler: """ Fragt die Geraete ab, die sich nicht von selbst melden - im Hintergrund. Frueher geschah das mitten in der Hauptschleife: neunzehn Tahoma-Geraete nacheinander, jedes mit bis zu zehn Sekunden Zeitlimit, dazu die HTTP- und WLED-Geraete. Eine Runde dauerte dadurch rund fuenfundvierzig statt dreissig Sekunden. Gepollt wurde, sobald seit dem letzten Poll sechzig Sekunden vergangen waren - also erst jede zweite Runde, in Wahrheit alle neunzig Sekunden. Weil die Uhr an derselben Abfrage hing, uebersprang der Runner jede dritte Minute, und ein Ausloeser "um 18:26" wurde nie wahr. Jetzt hat jeder Transport seinen eigenen Faden und seinen eigenen Abstand: die Rollaeden duerfen gemuetlich alle fuenf Minuten, waehrend die Hauptschleife im Sekundentakt weiterlaeuft. Die Faeden fassen die Datenbank nicht an - sie legen ihre Werte in eine Queue, geschrieben wird im Hauptfaden. Dieselbe Regel wie beim Versand: eine pymysql-Verbindung gehoert einem Faden. """ def __init__(self): self.ergebnisse = queue.Queue() self.laeuft = True def aufnehmen(self, transport, abstand): faden = threading.Thread(target=self._arbeiten, args=(transport, abstand), name="sammler-" + transport.schema, daemon=True) faden.start() def _arbeiten(self, transport, abstand): while self.laeuft: try: werte = transport.zustaende_lesen() if werte: self.ergebnisse.put(werte) except Exception as fehler: logger.warning("%s nicht abfragbar: %s", transport.schema, fehler) # In Sekundenschritten warten, damit das Beenden nicht bis zum # naechsten Durchgang dauert - bei Tahoma waeren das fuenf Minuten. for _ in range(max(1, int(abstand))): if not self.laeuft: return time.sleep(1) def abholen(self): """Alles, was die Faeden seit dem letzten Mal geliefert haben.""" neu = {} while True: try: neu.update(self.ergebnisse.get_nowait()) except queue.Empty: return neu def beenden(self): self.laeuft = False # =========================================================================== # Der Runner # =========================================================================== class Runner: def __init__(self, config, dry_run=False): self.config = config self.dry_run = dry_run or config.ja("runner", "dry_run") self.db = verbinden(config) self.solar = verbinden(config, "solar") self.werte = {} # state_id -> letzter bekannter Wert self.geschrieben = {} # state_id -> zuletzt in die DB geschriebener Wert self.geschrieben_um = {} # state_id -> wann das war self.war_aktiv = {} # automation_id -> war im Zeitfenster self.lief_im_fenster = {} # automation_id -> hat im Fenster ausgeloest self._sonne = (None, "00:00", "00:00") # (datum, aufgang, untergang) self._kalender = {} # datum -> {"feiertag": .., "ferien": ..} self._letzte_saeuberung = None self.reihenfolge = [] # automation_id, Ausloeser vor Nachfolger self.ausloeser_states = {} # automation_id -> state_id des Ausloesers self.versand = Versand() self.sammler = Sammler() self.mqtt = self._mqtt_verbinden() self.transporte = [ MQTTTransport(self.mqtt, self.dry_run), WLEDTransport(requests, dry_run=self.dry_run), HTTPTransport(requests, dry_run=self.dry_run), TahomaTransport(requests, config.text("tahoma", "pin"), config.text("tahoma", "token"), config.zahl("tahoma", "timeout", 10), self.dry_run), LogicTransport(self.sonnenzeiten), AutomatikTransport(self.ausloesezeiten), ] 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, tag): """ Ferien und Feiertage eines Tages, je Tag einmal geholt. Gefragt wird nach heute und - fuer Vorabend-Regeln - nach morgen. Gemerkt werden deshalb beide; was vor heute liegt, fliegt raus. """ if tag in self._kalender: return self._kalender[tag] stand = {"feiertag": False, "ferien": False} try: with self.db.cursor() as c: c.execute("SELECT holiday, vacation FROM calendar_days WHERE date = %s", (tag,)) zeile = c.fetchone() if zeile: stand = {"feiertag": bool(zeile["holiday"]), "ferien": bool(zeile["vacation"])} except Exception as fehler: # Nicht merken: beim naechsten Durchlauf noch einmal fragen, # statt den ganzen Tag mit "kein Feiertag" weiterzurechnen. logger.warning("Kalender nicht lesbar: %r", fehler) return stand heute = date.today() self._kalender = {t: s for t, s in self._kalender.items() if t >= heute} self._kalender[tag] = stand return stand # --- Werte ----------------------------------------------------------- def sammler_starten(self): """ Je gepolltem Transport ein Faden. MQTT meldet sich von selbst und die Uhr rechnet der Runner - die beiden stehen hier nicht. """ vorgaben = [(HTTPTransport, "poll_http", 60), (WLEDTransport, "poll_wled", 60), (TahomaTransport, "poll_tahoma", 300)] for transport in self.transporte: for klasse, schluessel, vorgabe in vorgaben: if isinstance(transport, klasse): abstand = self.config.zahl("runner", schluessel, vorgabe) self.sammler.aufnehmen(transport, abstand) logger.info("%s wird alle %d s abgefragt", klasse.__name__, abstand) break def werte_einsammeln(self, auch_geraete=False): """ MQTT und Uhr kosten nichts und werden jeden Takt gelesen; die Geraete liefert der Sammler aus dem Hintergrund. Nur bei --once, wo es keine Faeden gibt, fragt die Schleife die Geraete selbst. """ neu = {} for transport in self.transporte: billig = isinstance(transport, (MQTTTransport, LogicTransport, AutomatikTransport)) if billig or auch_geraete: neu.update(transport.zustaende_lesen()) neu.update(self.sammler.abholen()) neu.update(self.versand.messwerte_abholen()) if neu: self.werte.update(neu) self.werte_zurueckschreiben(neu) return neu def werte_zurueckschreiben(self, neu): """ current_value nachfuehren, damit der Editor im Browser aktuelle Zahlen zeigt ("Temperatur (= 21,4 °C)"). Geschrieben wird nur, was sich geaendert hat. MQTT liefert ohnehin nur neu Hereingekommenes, die gepollten Transporte dagegen bei jedem Durchgang ihren kompletten Bestand - ohne diesen Vergleich gingen rund 180 unveraenderte Werte je Runde in die Tabelle. Zusaetzlich gedrosselt, sonst schreibt ein gespraechiger Sensor im Sekundentakt. """ jetzt = datetime.now() faellig = [(w, i) for i, w in neu.items() if self.geschrieben.get(i) != w and self.geschrieben_um.get(i, datetime.min) + SCHREIB_ABSTAND <= jetzt] if not faellig: return try: with self.db.cursor() as c: c.executemany("UPDATE actor_states SET current_value = %s WHERE id = %s", faellig) for wert, state_id in faellig: self.geschrieben[state_id] = wert self.geschrieben_um[state_id] = jetzt except Exception as fehler: logger.warning("current_value nicht schreibbar: %r", fehler) # --- Ausfuehren ------------------------------------------------------ def ausloesen(self, automatik, anlass): """ Die Aktionen einer Automatik in den Versand geben. Geschickt wird im Hintergrund - eine Jalousie zuzufahren dauert ueber eine Minute, und so lange darf die Auswertung nicht stehen. Was dabei schiefgeht, kommt spaeter ueber ergebnisse_verbuchen() ins Protokoll. """ fehlerText = [] eingereiht = 0 for aktion in automatik["aktionen"]: kommando = self.regelwerk.kommandos.get(aktion["command_id"]) if kommando is None: fehlerText.append("Kommando %s gibt es nicht mehr" % aktion["command_id"]) logger.error("%s: Kommando %s gibt es nicht mehr", automatik["name"], aktion["command_id"]) continue transport = self.transport_fuer(kommando["actor_url"]) if transport is None: fehlerText.append("Kein Transport fuer %s" % kommando["actor_url"]) logger.error("%s: fuer %s (%s) gibt es keinen Transport", automatik["name"], kommando["actor_name"], kommando["actor_url"]) continue auftrag = { "actor_url": kommando["actor_url"], "command_url": kommando["command_url"], "params": [{"url": p["url"], "name": p["parameter_name"], "wert": aktion["werte"].get(p["id"], "")} for p in kommando["params"]], } beschreibung = "%s: %s" % (kommando["actor_name"], kommando["command_name"]) logger.info("%s: %s (unterwegs)", automatik["name"], beschreibung) self.versand.einreihen(automatik["id"], beschreibung, transport, auftrag) eingereiht += 1 ergebnis = "error" if fehlerText else anlass detail = "; ".join(fehlerText)[:255] if fehlerText \ else ("%d Kommando(s) unterwegs" % eingereiht) try: with self.db.cursor() as c: c.execute("""UPDATE automations SET last_run = NOW(), changed = changed WHERE id = %s""", (automatik["id"],)) c.execute("""INSERT INTO automation_log (automation_id, result, detail) VALUES (%s, %s, %s)""", (automatik["id"], ergebnis, detail)) except Exception as fehler: logger.warning("Protokoll nicht schreibbar: %r", fehler) automatik["last_run"] = datetime.now() self.lief_im_fenster[automatik["id"]] = True # Den eigenen Ausloeserwert gleich mitfuehren, statt bis zum naechsten # werte_einsammeln() zu warten. Zusammen mit der topologischen # Reihenfolge greift ein Nachfolger mit Versatz null dadurch noch im # selben Takt. state_id = self.ausloeser_states.get(automatik["id"]) if state_id is not None: self.werte[state_id] = automatik["last_run"].strftime("%Y-%m-%d %H:%M:%S") def ergebnisse_verbuchen(self): """ Was der Versand inzwischen erledigt hat ins Protokoll schreiben. Nur Fehlschlaege bekommen eine eigene Zeile - der Lauf selbst steht schon drin, und ein Protokoll, das jedes gelungene Kommando einzeln auffuehrt, findet niemand mehr etwas darin. """ for automation_id, beschreibung, fehler in self.versand.abholen(): if fehler is None: logger.debug("erledigt: %s", beschreibung) continue logger.error("%s konnte nicht geschickt werden: %s", beschreibung, fehler) try: with self.db.cursor() as c: c.execute("""INSERT INTO automation_log (automation_id, result, detail) VALUES (%s, 'error', %s)""", (automation_id, ("%s: %s" % (beschreibung, fehler))[:255])) except Exception as schreibfehler: logger.warning("Protokoll nicht schreibbar: %s", schreibfehler) def flanke_merken(self, automatik, erfuellt): if bool(automatik["cond_met"]) == bool(erfuellt): return automatik["cond_met"] = 1 if erfuellt else 0 try: with self.db.cursor() as c: c.execute("""UPDATE automations SET cond_met = %s, changed = changed WHERE id = %s""", (automatik["cond_met"], automatik["id"])) except Exception as fehler: logger.warning("cond_met nicht schreibbar: %r", fehler) def durchlauf(self, fenster=NACHHOLFENSTER): jetzt = datetime.now() for automation_id in self.reihenfolge: automatik = self.regelwerk.automatiken.get(automation_id) if automatik is None: continue tag = gemeinter_tag(automatik, jetzt) aktiv = (tag_passt(automatik, tag, self.kalender(tag)) and im_zeitfenster(jetzt, automatik["window_from"], automatik["window_to"])) vorher_aktiv = self.war_aktiv.get(automatik["id"], aktiv) if aktiv: if not vorher_aktiv: self.lief_im_fenster[automatik["id"]] = False erfuellt = gruppen_erfuellt(automatik, self.regelwerk, self.werte, jetzt, fenster) if erfuellt and not automatik["cond_met"]: if automatik.get("once_per_day") and self.lief_heute(automatik, jetzt): logger.debug("%s: lief heute schon, bis Mitternacht gesperrt", automatik["name"]) elif self.gesperrt(automatik, jetzt): logger.debug("%s: Flanke faellt in die Sperrzeit, uebersprungen", automatik["name"]) else: self.ausloesen(automatik, "fired") self.flanke_merken(automatik, 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): """ Hat die Automatik heute schon ausgeloest? Zwei Stellen fragen danach. force_once will wissen, ob es am Ende des Fensters noch etwas nachzuholen gibt. once_per_day will das Gegenteil: einmal am Tag genuegt, danach ist bis Mitternacht Ruhe. Letzteres ist fuer alles gedacht, was man hinterher von Hand wieder anders stellt. Ein Rollladen, den man um acht nochmal zugezogen hat, soll nicht um neun von selbst wieder auffahren, nur weil eine Wolke weiterzieht und die Helligkeitsschwelle ein zweites Mal steigt. Die Sperrzeit taugt dafuer nicht: sie zaehlt Sekunden und muesste auf einen Tag stehen, womit sie am naechsten Morgen den Termin knapp verfehlen wuerde. """ letzter = automatik.get("last_run") return bool(letzter) and letzter.date() == jetzt.date() def protokoll_saeubern(self): heute = date.today() if self._letzte_saeuberung == heute: return self._letzte_saeuberung = heute try: with self.db.cursor() as c: c.execute("DELETE FROM automation_log WHERE ts < NOW() - INTERVAL %s DAY", (LOG_AUFBEWAHRUNG_TAGE,)) except Exception as fehler: logger.warning("Protokoll nicht aufraeumbar: %r", fehler) # --- Hauptschleife --------------------------------------------------- def laufen(self, nur_einmal=False): """ Die Schleife macht nur noch Billiges: Werte abholen, auswerten, protokollieren. Alles, was auf ein Netz warten muss, laeuft daneben - die Geraeteabfrage im Sammler, das Schalten im Versand. Deshalb darf der Takt kurz sein. """ self.regelwerk_laden() takt = self.config.zahl("runner", "tick_seconds", 10) neuladen = self.config.zahl("runner", "reload_seconds", 30) fenster = self.config.zahl("runner", "catchup_minutes", NACHHOLFENSTER) letztes_pruefen = 0.0 if not nur_einmal: self.sammler_starten() while True: jetzt = time.monotonic() self.werte_einsammeln(auch_geraete=nur_einmal) self.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())