diff --git a/src/rippy/bus/test_waechter.py b/src/rippy/bus/test_waechter.py new file mode 100644 index 0000000..f6ab932 --- /dev/null +++ b/src/rippy/bus/test_waechter.py @@ -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 diff --git a/src/rippy/bus/waechter.py b/src/rippy/bus/waechter.py index 0cee359..9eef37a 100644 --- a/src/rippy/bus/waechter.py +++ b/src/rippy/bus/waechter.py @@ -85,31 +85,39 @@ def unterschiede(vorher: dict, jetzt: dict) -> list: festgelegt: erst neue Jobs, dann Änderungen, dann verschwundene. So sieht 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(): if job_id not in vorher: - ereignisse.append(("job.created", job_id, stand)) + neu.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, { + geaendert.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)) + geaendert.append(("job.progress", job_id, stand)) else: - ereignisse.append(("job.phase", job_id, stand)) + geaendert.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"})) + verschwunden = [ + ("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: