diff --git a/docker/api/test_api_smoke.py b/docker/api/test_api_smoke.py index 89f9df8..61242a1 100644 --- a/docker/api/test_api_smoke.py +++ b/docker/api/test_api_smoke.py @@ -286,27 +286,51 @@ def test_sse_rahmen_uebersteht_umlaute(): assert "Größe" in rahmen -def test_snapshot_meldet_unlesbare_laufwerke_als_none(): +def test_snapshot_meldet_unlesbare_laufwerke_als_none(monkeypatch): """„konnte nicht nachsehen" ist etwas anderes als „es gibt keine". Genau diese Vermischung hat in v1 die Job-Liste im Sekundentakt geleert (fuenfmal `catch(() => [])` im UI). Ein Snapshot mit devices=[] wuerde dem UI sagen „du hast kein Laufwerk"; None sagt „ich weiss es gerade nicht", und das UI behaelt seinen Stand. + + OHNE DATENBANK: Die Ampel hat keine Postgres — der erste Anlauf dieses + Tests lief deshalb rot (Lauf 170). Der Store wird hier ersetzt, denn + geprueft wird die Snapshot-LOGIK, nicht die Datenbank. """ import asyncio import main + monkeypatch.setattr(main.db, "list_jobs", lambda limit=50: []) + monkeypatch.setattr(main.db, "list_workers", lambda: []) + monkeypatch.setattr(main.db, "list_logs", lambda limit=50: []) + def kaputt(): raise OSError("Laufwerk haengt") - original = main.device_discovery.list_optical_devices - main.device_discovery.list_optical_devices = kaputt - try: - zustand = asyncio.run(main._snapshot()) - finally: - main.device_discovery.list_optical_devices = original + monkeypatch.setattr(main.device_discovery, "list_optical_devices", kaputt) - assert zustand["devices"] is None - assert isinstance(zustand["jobs"], list) + zustand = asyncio.run(main._snapshot()) + + assert zustand["devices"] is None, "Ein unlesbares Laufwerk darf nicht als [] durchgehen" + assert zustand["jobs"] == [] + + +def test_snapshot_liefert_alle_vier_bereiche(monkeypatch): + """Ein Client, der sich verbindet, bekommt das GANZE Bild — sonst muesste + er den Rest raten und faellt auf Polling zurueck.""" + import asyncio + + import main + + monkeypatch.setattr(main.db, "list_jobs", lambda limit=50: []) + monkeypatch.setattr(main.db, "list_workers", lambda: [{"name": "pc"}]) + monkeypatch.setattr(main.db, "list_logs", lambda limit=50: []) + monkeypatch.setattr(main.device_discovery, "list_optical_devices", lambda: []) + + zustand = asyncio.run(main._snapshot()) + + assert set(zustand) == {"jobs", "workers", "logs", "devices"} + assert zustand["devices"] == [] # wirklich leer, nicht „unbekannt" + assert zustand["workers"] == [{"name": "pc"}] diff --git a/src/rippy/bus/waechter.py b/src/rippy/bus/waechter.py new file mode 100644 index 0000000..0cee359 --- /dev/null +++ b/src/rippy/bus/waechter.py @@ -0,0 +1,211 @@ +"""Der Wächter: macht Änderungen des Workers zu Ereignissen. + +## Das Problem, das er löst + +Der Bus (`memory.py`) verteilt Ereignisse innerhalb EINES Prozesses. Der Worker +ist aber ein eigener Prozess — im Docker-Betrieb sogar ein eigener Container, +im verteilten Betrieb ein anderer Rechner. Wenn dort der Fortschritt eines Rips +von 12 auf 13 Prozent geht, weiß die API davon nichts. + +Drei Wege wären denkbar: + +1. **Der Worker ruft die API an.** Dann hängt jeder Fortschrittswert an einer + HTTP-Verbindung, und ein kurzer API-Neustart ließe Ereignisse verschwinden. + Der Worker müsste außerdem wissen, wo die API steht — im verteilten Betrieb + ist das eine zusätzliche Konfiguration, die schiefgehen kann. +2. **Ein gemeinsamer Broker (Redis Pub/Sub).** Sauber, und genau das kommt im + verteilten Betrieb (V2-5). Aber der Standalone-Betrieb hat bewusst KEIN + Redis — das war der ganze Punkt von V2-2. +3. **Die Datenbank fragen.** Sie ist ohnehin die Wahrheit über den Job-Zustand + (KONZEPT-V2.md § 3.1), sie ist in JEDEM Betriebsmodus da, und der Worker + schreibt dort ohnehin hin. + +Es wird Weg 3. + +## „Ist das nicht wieder Polling?" + +Doch — aber an der Stelle, an der es billig ist, und genau einmal. + + vorher N Browser-Tabs × 9 Endpunkte × alle 4-5 s über HTTP, durch nginx, + durch das Rate-Limit + jetzt 1 Abfrage/Sekunde lokal, indiziert, + im selben Netz wie die DB + +Das Ziel war nie „nirgendwo mehr nachfragen", sondern: **der Browser fragt +nicht mehr nach.** Ein offener Tab verursacht jetzt null Anfragen statt 121. +Und die Abfrage hier holt vier Spalten von höchstens ein paar Dutzend Zeilen. + +## Was er NICHT tut + +Er schickt keinen Zustand über den Bus, sondern nur die Nachricht, DASS sich +etwas geändert hat, plus die geänderten Felder. Wer den vollen Zustand will, +holt sich einen Snapshot. Sonst wäre ein verpasstes Ereignis wieder eine +Aussage über die Welt. + +## Und wenn er stirbt? + +Dann steht das UI still, ohne es zu merken — der schlimmste Fall. Deshalb: +Jeder Fehler wird protokolliert UND als `system.notice` gesendet, und +`lebt_seit_sekunden()` macht sein Alter abfragbar. Das ist die Lehre aus +AGENTS.md: „Ein Hintergrund-Prozess, der still scheitert, ist schlimmer als +einer, der laut scheitert" — dort hatte ein `except Exception: pass` in einer +Vorrats-Schleife eine Stunde gekostet. +""" + +import asyncio +import time + +# Wie oft nachgesehen wird. Eine Sekunde ist zugleich die Drosselung der +# Fortschritts-Ereignisse (KONZEPT-V2.md § 6.3: höchstens 1/s) — schneller +# könnte kein Auge folgen, und HandBrake meldet ohnehin nicht öfter. +TAKT_SEKUNDEN = 1.0 + +# Nach einem Fehler wird langsamer nachgesehen, damit ein dauerhaft kaputter +# Zustand (DB weg) nicht jede Sekunde eine Meldung erzeugt. +FEHLER_TAKT_SEKUNDEN = 5.0 + +# Endzustände: ab hier ist ein Job durch. +ENDE = ("completed", "failed", "canceled") + + +def _job_kurz(zeile: dict) -> dict: + """Nur die Felder, deren Änderung ein Ereignis wert ist.""" + return { + "status": zeile.get("status"), + "progress": zeile.get("progress"), + "title": zeile.get("title"), + "error": zeile.get("error"), + } + + +def unterschiede(vorher: dict, jetzt: dict) -> list: + """Was hat sich geändert? (pure Funktion — deshalb testbar ohne DB) + + Gibt eine Liste von `(typ, job_id, daten)` zurück. Die Reihenfolge ist + festgelegt: erst neue Jobs, dann Änderungen, dann verschwundene. So sieht + ein Client einen Job nie „geändert", bevor er ihn kennt. + """ + ereignisse = [] + + for job_id, stand in jetzt.items(): + if job_id not in vorher: + ereignisse.append(("job.created", job_id, stand)) + continue + alt = vorher[job_id] + if alt == stand: + continue + if alt.get("status") != stand.get("status"): + typ = "job.finished" if stand.get("status") in ENDE else "job.phase" + ereignisse.append((typ, job_id, { + "von": alt.get("status"), "nach": stand.get("status"), + **stand, + })) + elif alt.get("progress") != stand.get("progress"): + ereignisse.append(("job.progress", job_id, stand)) + else: + ereignisse.append(("job.phase", job_id, stand)) + + for job_id in vorher: + if job_id not in jetzt: + ereignisse.append(("job.finished", job_id, {"status": "geloescht"})) + + return ereignisse + + +class Waechter: + """Sieht in der Datenbank nach und meldet Änderungen an den Bus.""" + + def __init__(self, store, bus, takt: float = TAKT_SEKUNDEN): + self._store = store + self._bus = bus + self._takt = takt + self._jobs: dict = {} + self._letzte_log_id = None + self._letzter_lauf = 0.0 + self._fehler_in_folge = 0 + self._laeuft = False + + # ── Abfragbarkeit: wer einen Vorrat anlegt, macht sein Alter sichtbar ── + def lebt_seit_sekunden(self) -> float: + """Wie lange ist der letzte erfolgreiche Durchlauf her? + + Gedacht für `GET /health/vorraete`. Ein Wächter, der still gestorben + ist, sieht von außen genauso aus wie einer, bei dem gerade nichts + passiert — es sei denn, man kann sein Alter erfragen. + """ + if not self._letzter_lauf: + return -1.0 + return time.monotonic() - self._letzter_lauf + + @property + def gesund(self) -> bool: + alter = self.lebt_seit_sekunden() + return alter >= 0 and alter < self._takt * 10 + + # ── Ein Durchlauf ─────────────────────────────────────────────────── + def einmal(self) -> int: + """Einmal nachsehen und melden. Gibt die Zahl der Ereignisse zurück. + + Bewusst synchron und getrennt von der Schleife: So lässt sie sich + ohne asyncio testen. + """ + gesendet = 0 + + jetzt = {z["id"]: _job_kurz(z) for z in self._store.list_jobs(limit=100)} + erster_lauf = not self._jobs and self._letzte_log_id is None + if not erster_lauf: + for typ, job_id, daten in unterschiede(self._jobs, jetzt): + self._bus.senden(typ, daten, entitaet="job", entitaet_id=job_id) + gesendet += 1 + self._jobs = jetzt + + zeilen = self._store.list_logs(limit=50) + if erster_lauf: + self._letzte_log_id = zeilen[0]["id"] if zeilen else 0 + else: + neue = [z for z in zeilen if z["id"] > (self._letzte_log_id or 0)] + for zeile in reversed(neue): # älteste zuerst + self._bus.senden("log.line", { + "level": zeile.get("level"), + "quelle": zeile.get("source"), + "text": zeile.get("message"), + }, entitaet="log", entitaet_id=str(zeile["id"])) + gesendet += 1 + if neue: + self._letzte_log_id = max(z["id"] for z in neue) + + self._letzter_lauf = time.monotonic() + return gesendet + + # ── Die Schleife ──────────────────────────────────────────────────── + async def schleife(self) -> None: + """Läuft, bis sie abgebrochen wird. Meldet Fehler LAUT.""" + self._laeuft = True + while self._laeuft: + try: + await asyncio.to_thread(self.einmal) + self._fehler_in_folge = 0 + await asyncio.sleep(self._takt) + except asyncio.CancelledError: + raise + except Exception as e: + self._fehler_in_folge += 1 + # Nur beim ERSTEN Fehler melden und danach alle zehn: Ein + # dauerhaft kaputter Zustand soll nicht das Log zumüllen — + # aber schweigen darf er auch nicht. + if self._fehler_in_folge == 1 or self._fehler_in_folge % 10 == 0: + text = (f"Ereignis-Wächter: {e} " + f"({self._fehler_in_folge}. Fehlschlag in Folge)") + try: + self._store.add_log("error", "waechter", text) + except Exception: + pass # DB weg — dann geht wenigstens der Bus + try: + self._bus.senden("system.notice", + {"level": "error", "text": text}) + except Exception: + pass + await asyncio.sleep(FEHLER_TAKT_SEKUNDEN) + + def stoppen(self) -> None: + self._laeuft = False