feat(bus): Waechter — Worker-Aenderungen werden zu Ereignissen
Ampel / ampel (push) Successful in 43s
Ampel / ampel (push) Successful in 43s
WAS: rippy/bus/waechter.py sieht im Takt von 1 s in der Datenbank nach und meldet Aenderungen an den Bus: job.created / job.phase / job.progress / job.finished und log.line. 12 Tests, ohne Datenbank lauffaehig. WARUM ES IHN BRAUCHT: Der Bus verteilt innerhalb EINES Prozesses. Der Worker ist aber ein eigener Prozess, im Docker-Betrieb ein eigener Container, im verteilten Betrieb ein anderer Rechner. Ohne Bruecke wuesste die API nichts vom Fortschritt eines Rips. WARUM UEBER DIE DATENBANK und nicht per HTTP-Rueckruf oder Redis: - HTTP-Rueckruf: jeder Fortschrittswert haengt an einer Verbindung, ein kurzer API-Neustart liesse Ereignisse verschwinden, und der Worker braeuchte zusaetzliche Konfiguration. - Redis Pub/Sub: sauber, kommt im verteilten Betrieb (V2-5) — aber der Standalone-Betrieb hat bewusst KEIN Redis, das war der Punkt von V2-2. - Die Datenbank ist ohnehin die Wahrheit ueber den Job-Zustand, ist in JEDEM Modus da, und der Worker schreibt dort sowieso hin. "Ist das nicht wieder Polling?" Doch — aber einmal, lokal und billig: vorher N Tabs x 9 Endpunkte alle 4-5 s ueber HTTP durch nginx durch das Rate-Limit; jetzt EINE Abfrage pro Sekunde von vier Spalten. Das Ziel war nie "nirgendwo nachfragen", sondern: der Browser fragt nicht mehr. DER TEST HAT EINEN ECHTEN FEHLER GEFUNDEN: Der Docstring versprach "erst neue Jobs, dann Aenderungen", die Schleife lieferte aber die Reihenfolge des dicts — ein Client haette ein job.progress fuer einen Job bekommen koennen, den er noch gar nicht kennt. Jetzt drei getrennte Durchlaeufe. Das Versprechen war richtig, der Code nicht. Und er stirbt nicht still: Fehler gehen ins Log UND als system.notice an den Bus (erster Fehlschlag, danach jeder zehnte), und lebt_seit_sekunden() macht sein Alter abfragbar — AGENTS.md, "wer einen Vorrat anlegt, macht sein Alter abfragbar". GEMESSEN: ruff sauber, 371 Tests gruen + 3 uebersprungen (vorher 359). Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5
parent
1fbe2cfc55
commit
f097a0b59d
@@ -0,0 +1,159 @@
|
|||||||
|
"""Der Wächter macht Worker-Änderungen zu Ereignissen — ohne Datenbank getestet.
|
||||||
|
|
||||||
|
`unterschiede()` ist eine reine Funktion, `einmal()` bekommt Store und Bus
|
||||||
|
eingespritzt. Beides läuft deshalb auf jeder Plattform und ohne Postgres —
|
||||||
|
genau die Eigenschaft, die dem ersten Anlauf der SSE-Tests gefehlt hat
|
||||||
|
(Ampel-Lauf 170 lief rot, weil ein Test eine Datenbank brauchte).
|
||||||
|
"""
|
||||||
|
|
||||||
|
from rippy.bus.waechter import Waechter, unterschiede
|
||||||
|
|
||||||
|
|
||||||
|
class FakeBus:
|
||||||
|
def __init__(self):
|
||||||
|
self.gesendet = []
|
||||||
|
|
||||||
|
def senden(self, typ, daten=None, entitaet=None, entitaet_id=None):
|
||||||
|
self.gesendet.append((typ, entitaet_id, daten or {}))
|
||||||
|
return {"typ": typ}
|
||||||
|
|
||||||
|
@property
|
||||||
|
def typen(self):
|
||||||
|
return [t for t, _, _ in self.gesendet]
|
||||||
|
|
||||||
|
|
||||||
|
class FakeStore:
|
||||||
|
def __init__(self, jobs=None, logs=None):
|
||||||
|
self.jobs = jobs or []
|
||||||
|
self.logs = logs or []
|
||||||
|
self.protokoll = []
|
||||||
|
|
||||||
|
def list_jobs(self, limit=100):
|
||||||
|
return list(self.jobs)
|
||||||
|
|
||||||
|
def list_logs(self, limit=50):
|
||||||
|
return list(self.logs)
|
||||||
|
|
||||||
|
def add_log(self, level, quelle, text):
|
||||||
|
self.protokoll.append((level, quelle, text))
|
||||||
|
|
||||||
|
|
||||||
|
def _job(id_, status="pending", progress=0, title="Akira", error=None):
|
||||||
|
return {"id": id_, "status": status, "progress": progress,
|
||||||
|
"title": title, "error": error}
|
||||||
|
|
||||||
|
|
||||||
|
# ── unterschiede(): die reine Logik ─────────────────────────────────────
|
||||||
|
def test_neuer_job_ist_ein_created():
|
||||||
|
ereignisse = unterschiede({}, {"j1": {"status": "pending", "progress": 0}})
|
||||||
|
assert ereignisse == [("job.created", "j1", {"status": "pending", "progress": 0})]
|
||||||
|
|
||||||
|
|
||||||
|
def test_gleicher_stand_erzeugt_nichts():
|
||||||
|
stand = {"j1": {"status": "ripping", "progress": 10}}
|
||||||
|
assert unterschiede(stand, dict(stand)) == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_fortschritt_ist_ein_progress():
|
||||||
|
vorher = {"j1": {"status": "ripping", "progress": 10}}
|
||||||
|
jetzt = {"j1": {"status": "ripping", "progress": 11}}
|
||||||
|
typ, job_id, _ = unterschiede(vorher, jetzt)[0]
|
||||||
|
assert (typ, job_id) == ("job.progress", "j1")
|
||||||
|
|
||||||
|
|
||||||
|
def test_statuswechsel_mitten_drin_ist_eine_phase():
|
||||||
|
vorher = {"j1": {"status": "ripping", "progress": 99}}
|
||||||
|
jetzt = {"j1": {"status": "transcoding", "progress": 0}}
|
||||||
|
typ, _, daten = unterschiede(vorher, jetzt)[0]
|
||||||
|
assert typ == "job.phase"
|
||||||
|
assert daten["von"] == "ripping" and daten["nach"] == "transcoding"
|
||||||
|
|
||||||
|
|
||||||
|
def test_endzustand_ist_ein_finished():
|
||||||
|
for ende in ("completed", "failed", "canceled"):
|
||||||
|
vorher = {"j1": {"status": "transcoding", "progress": 80}}
|
||||||
|
jetzt = {"j1": {"status": ende, "progress": 100}}
|
||||||
|
assert unterschiede(vorher, jetzt)[0][0] == "job.finished"
|
||||||
|
|
||||||
|
|
||||||
|
def test_geloeschter_job_wird_gemeldet():
|
||||||
|
"""Sonst bliebe er im UI stehen, bis jemand neu laedt — und genau das
|
||||||
|
Neuladen soll ja verschwinden."""
|
||||||
|
typ, job_id, daten = unterschiede({"j1": {"status": "completed"}}, {})[0]
|
||||||
|
assert (typ, job_id) == ("job.finished", "j1")
|
||||||
|
assert daten["status"] == "geloescht"
|
||||||
|
|
||||||
|
|
||||||
|
def test_neue_jobs_kommen_vor_aenderungen():
|
||||||
|
"""Ein Client darf einen Job nie „geaendert" sehen, bevor er ihn kennt."""
|
||||||
|
vorher = {"alt": {"status": "ripping", "progress": 1}}
|
||||||
|
jetzt = {"alt": {"status": "ripping", "progress": 2},
|
||||||
|
"neu": {"status": "pending", "progress": 0}}
|
||||||
|
typen = [t for t, _, _ in unterschiede(vorher, jetzt)]
|
||||||
|
assert typen.index("job.created") < typen.index("job.progress")
|
||||||
|
|
||||||
|
|
||||||
|
# ── einmal(): der Durchlauf ─────────────────────────────────────────────
|
||||||
|
def test_erster_lauf_meldet_nichts():
|
||||||
|
"""Sonst bekaeme jeder API-Neustart die gesamte Job-Historie als
|
||||||
|
„gerade passiert" — und das UI zeigte hundert Meldungen auf einmal."""
|
||||||
|
bus = FakeBus()
|
||||||
|
store = FakeStore(jobs=[_job("j1", "completed", 100)],
|
||||||
|
logs=[{"id": 7, "level": "info", "source": "worker", "message": "alt"}])
|
||||||
|
w = Waechter(store, bus)
|
||||||
|
assert w.einmal() == 0
|
||||||
|
assert bus.gesendet == []
|
||||||
|
|
||||||
|
|
||||||
|
def test_zweiter_lauf_meldet_die_aenderung():
|
||||||
|
bus = FakeBus()
|
||||||
|
store = FakeStore(jobs=[_job("j1", "ripping", 10)])
|
||||||
|
w = Waechter(store, bus)
|
||||||
|
w.einmal()
|
||||||
|
store.jobs = [_job("j1", "ripping", 25)]
|
||||||
|
assert w.einmal() == 1
|
||||||
|
assert bus.typen == ["job.progress"]
|
||||||
|
assert bus.gesendet[0][1] == "j1"
|
||||||
|
|
||||||
|
|
||||||
|
def test_neue_logzeilen_kommen_aelteste_zuerst():
|
||||||
|
"""list_logs gibt neueste zuerst. Wuerde der Waechter das durchreichen,
|
||||||
|
liefe das Log-Fenster rueckwaerts."""
|
||||||
|
bus = FakeBus()
|
||||||
|
store = FakeStore(logs=[{"id": 5, "level": "info", "source": "w", "message": "alt"}])
|
||||||
|
w = Waechter(store, bus)
|
||||||
|
w.einmal()
|
||||||
|
store.logs = [
|
||||||
|
{"id": 7, "level": "info", "source": "w", "message": "zweite"},
|
||||||
|
{"id": 6, "level": "info", "source": "w", "message": "erste"},
|
||||||
|
{"id": 5, "level": "info", "source": "w", "message": "alt"},
|
||||||
|
]
|
||||||
|
w.einmal()
|
||||||
|
texte = [d["text"] for t, _, d in bus.gesendet if t == "log.line"]
|
||||||
|
assert texte == ["erste", "zweite"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_dieselbe_logzeile_kommt_nur_einmal():
|
||||||
|
bus = FakeBus()
|
||||||
|
store = FakeStore(logs=[{"id": 1, "level": "info", "source": "w", "message": "a"}])
|
||||||
|
w = Waechter(store, bus)
|
||||||
|
w.einmal()
|
||||||
|
store.logs = [{"id": 2, "level": "info", "source": "w", "message": "b"},
|
||||||
|
{"id": 1, "level": "info", "source": "w", "message": "a"}]
|
||||||
|
w.einmal()
|
||||||
|
w.einmal()
|
||||||
|
assert [d["text"] for t, _, d in bus.gesendet if t == "log.line"] == ["b"]
|
||||||
|
|
||||||
|
|
||||||
|
def test_alter_ist_abfragbar():
|
||||||
|
"""AGENTS.md: „wer einen Vorrat anlegt, macht sein Alter abfragbar".
|
||||||
|
|
||||||
|
Ein Waechter, der still gestorben ist, sieht von aussen genauso aus wie
|
||||||
|
einer, bei dem gerade nichts passiert — es sei denn, man kann fragen.
|
||||||
|
"""
|
||||||
|
w = Waechter(FakeStore(), FakeBus())
|
||||||
|
assert w.lebt_seit_sekunden() == -1.0 # noch nie gelaufen
|
||||||
|
assert w.gesund is False
|
||||||
|
w.einmal()
|
||||||
|
assert 0 <= w.lebt_seit_sekunden() < 1
|
||||||
|
assert w.gesund is True
|
||||||
@@ -85,31 +85,39 @@ def unterschiede(vorher: dict, jetzt: dict) -> list:
|
|||||||
festgelegt: erst neue Jobs, dann Änderungen, dann verschwundene. So sieht
|
festgelegt: erst neue Jobs, dann Änderungen, dann verschwundene. So sieht
|
||||||
ein Client einen Job nie „geändert", bevor er ihn kennt.
|
ein Client einen Job nie „geändert", bevor er ihn kennt.
|
||||||
"""
|
"""
|
||||||
ereignisse = []
|
# DREI getrennte Durchläufe, nicht einer. Ein einzelner Durchlauf über
|
||||||
|
# `jetzt` würde die Reihenfolge der Einfügung durchreichen — und dann käme
|
||||||
|
# ein „geändert" vor dem „neu" eines anderen Jobs. Ein Client, der ein
|
||||||
|
# job.progress für einen Job bekommt, den er nicht kennt, kann damit
|
||||||
|
# nichts anfangen: Er müsste raten oder einen Snapshot nachfordern.
|
||||||
|
# (Der Test dafür hat genau diesen Fehler in der ersten Fassung gefunden.)
|
||||||
|
neu = []
|
||||||
|
geaendert = []
|
||||||
|
|
||||||
for job_id, stand in jetzt.items():
|
for job_id, stand in jetzt.items():
|
||||||
if job_id not in vorher:
|
if job_id not in vorher:
|
||||||
ereignisse.append(("job.created", job_id, stand))
|
neu.append(("job.created", job_id, stand))
|
||||||
continue
|
continue
|
||||||
alt = vorher[job_id]
|
alt = vorher[job_id]
|
||||||
if alt == stand:
|
if alt == stand:
|
||||||
continue
|
continue
|
||||||
if alt.get("status") != stand.get("status"):
|
if alt.get("status") != stand.get("status"):
|
||||||
typ = "job.finished" if stand.get("status") in ENDE else "job.phase"
|
typ = "job.finished" if stand.get("status") in ENDE else "job.phase"
|
||||||
ereignisse.append((typ, job_id, {
|
geaendert.append((typ, job_id, {
|
||||||
"von": alt.get("status"), "nach": stand.get("status"),
|
"von": alt.get("status"), "nach": stand.get("status"),
|
||||||
**stand,
|
**stand,
|
||||||
}))
|
}))
|
||||||
elif alt.get("progress") != stand.get("progress"):
|
elif alt.get("progress") != stand.get("progress"):
|
||||||
ereignisse.append(("job.progress", job_id, stand))
|
geaendert.append(("job.progress", job_id, stand))
|
||||||
else:
|
else:
|
||||||
ereignisse.append(("job.phase", job_id, stand))
|
geaendert.append(("job.phase", job_id, stand))
|
||||||
|
|
||||||
for job_id in vorher:
|
verschwunden = [
|
||||||
if job_id not in jetzt:
|
("job.finished", job_id, {"status": "geloescht"})
|
||||||
ereignisse.append(("job.finished", job_id, {"status": "geloescht"}))
|
for job_id in vorher if job_id not in jetzt
|
||||||
|
]
|
||||||
|
|
||||||
return ereignisse
|
return neu + geaendert + verschwunden
|
||||||
|
|
||||||
|
|
||||||
class Waechter:
|
class Waechter:
|
||||||
|
|||||||
Reference in New Issue
Block a user