diff --git a/docker/worker/celery_app.py b/docker/worker/celery_app.py index 955534f..dc7b499 100644 --- a/docker/worker/celery_app.py +++ b/docker/worker/celery_app.py @@ -33,6 +33,30 @@ celery_app.conf.update( # zuverlässig im sys.path (ModuleNotFoundError 'caps', Deploy 23.07.). import caps # noqa: E402 import db # noqa: E402 +import zombies # noqa: E402 + + +@worker_ready.connect +def raeume_job_leichen_auf(**kwargs): + """Nach der Gnadenfrist: Jobs, an denen niemand arbeitet, ehrlich auf + 'failed' setzen (Details und Sicherungen in zombies.py). + + Läuft im Hintergrund-Thread — der Worker soll sofort Aufgaben annehmen und + nicht zwei Minuten auf die Aufräumung warten. + """ + import threading + import time + + def spaeter(): + time.sleep(zombies.GNADENFRIST_SEKUNDEN) + try: + db.init_db() + bericht = zombies.raeume_zombies_auf(celery_app, db) + print(f"Zombie-Erkennung: {bericht}") + except Exception as e: # darf den Worker nie mitnehmen + print(f"Zombie-Erkennung fehlgeschlagen: {e}") + + threading.Thread(target=spaeter, daemon=True, name="zombie-erkennung").start() @worker_ready.connect diff --git a/docker/worker/db.py b/docker/worker/db.py index 2d82a8e..39b6e85 100644 --- a/docker/worker/db.py +++ b/docker/worker/db.py @@ -181,6 +181,42 @@ def get_job_status(job_id: str) -> str: return zeile[0] if zeile else "" +def list_jobs_mit_status(stati) -> list: + """Alle Jobs in einem der genannten Zustände (id/status/title/created_at). + + Basis der Zombie-Erkennung: Jobs, die behaupten, es arbeite gerade jemand + an ihnen. Bewusst NUR diese schmale Auswahl statt der ganzen Zeile — die + Erkennung braucht nichts weiter. + """ + from sqlalchemy import select + + with engine.connect() as conn: + zeilen = conn.execute( + select(jobs.c.id, jobs.c.status, jobs.c.title, jobs.c.created_at) + .where(jobs.c.status.in_(list(stati))) + ).mappings().all() + return [dict(z) for z in zeilen] + + +def zaehle_online_worker(sekunden: int = 120) -> int: + """Wie viele Worker gelten laut Herzschlag gerade als online? + + Die Zombie-Erkennung vergleicht das mit der Zahl der Celery-Antworten: + melden sich weniger Worker als bekannt sind, ist die Auskunft + unvollstaendig — dann wird NICHTS als Leiche gewertet. + """ + from datetime import timedelta + + from sqlalchemy import func, select + + grenze = utcnow() - timedelta(seconds=sekunden) + with engine.connect() as conn: + anzahl = conn.execute( + select(func.count()).select_from(workers).where(workers.c.last_seen >= grenze) + ).scalar() + return int(anzahl or 0) + + def update_job(job_id: str, **fields) -> None: with engine.begin() as conn: conn.execute(jobs.update().where(jobs.c.id == job_id).values(**fields)) diff --git a/docker/worker/test_zombies.py b/docker/worker/test_zombies.py new file mode 100644 index 0000000..801f797 --- /dev/null +++ b/docker/worker/test_zombies.py @@ -0,0 +1,190 @@ +"""Tests der Zombie-Erkennung — ohne Postgres und ohne Redis. + +Der Schwerpunkt liegt bewusst auf dem, was WEHTUT: ein laufender Job darf +niemals als Leiche gelten. Genau das wäre am 25.07.2026 passiert, wenn die +Erkennung nach dem Alter des Jobs geurteilt hätte — der Akira-Job war seit +neun Stunden offen und lief trotzdem. +""" + +import zombies + + +# --- reine Funktionen ------------------------------------------------------- + +def test_belegte_job_ids_findet_id_an_beliebiger_stelle(): + """job_id steht bei rip_disc an zweiter, bei transcode_files an erster + Stelle — die Erkennung darf sich auf keine Position verlassen.""" + aktiv = { + "celery@node1": [ + {"name": "worker.tasks.transcode_files", "args": ["job-eins", "/raw", "/final"]}, + {"name": "worker.tasks.rip_disc", "args": ["/dev/sr0", "job-zwei"]}, + ] + } + belegt = zombies.belegte_job_ids([aktiv], ["job-eins", "job-zwei", "job-drei"]) + assert belegt == {"job-eins", "job-zwei"} + + +def test_belegte_job_ids_versteht_args_als_text(): + """Celery liefert args je nach Version als Liste ODER als Text-Repräsentation.""" + aktiv = {"celery@node1": [{"name": "x", "args": "('job-eins', '/raw', '/final')"}]} + assert zombies.belegte_job_ids([aktiv], ["job-eins", "job-zwei"]) == {"job-eins"} + + +def test_belegte_job_ids_ohne_auskunft_ist_leer(): + assert zombies.belegte_job_ids(None, ["a"]) == set() + assert zombies.belegte_job_ids([None, None, None], ["a"]) == set() + assert zombies.belegte_job_ids([{}, None], ["a"]) == set() + + +def test_antwortende_knoten_sammelt_ueber_alle_abfragen(): + aktiv = {"celery@a": []} + vorgemerkt = {"celery@b": []} + assert zombies.antwortende_knoten([aktiv, None, vorgemerkt]) == {"celery@a", "celery@b"} + assert zombies.antwortende_knoten([None, None]) == set() + + +def test_auskunft_nur_vollstaendig_wenn_alle_bekannten_antworten(): + assert zombies.auskunft_vollstaendig(1, 1) is True + assert zombies.auskunft_vollstaendig(2, 1) is True + # Ein bekannter Worker schweigt → NICHT urteilen + assert zombies.auskunft_vollstaendig(1, 2) is False + # Niemand antwortet → wir wissen nichts + assert zombies.auskunft_vollstaendig(0, 0) is False + + +def test_finde_zombies_verschont_belegte_jobs(): + offene = [{"id": "a", "status": "transcoding"}, {"id": "b", "status": "ripping"}] + assert zombies.finde_zombies(offene, {"a"}) == [{"id": "b", "status": "ripping"}] + assert zombies.finde_zombies(offene, {"a", "b"}) == [] + + +def test_fehlertext_nennt_zustand_und_den_weg_zurueck(): + text = zombies.fehlertext({"status": "transcoding"}) + assert "transcoding" in text + assert "Neu komprimieren" in text + assert "NICHT gelöscht" in text + + +# --- Attrappen für den Gesamtdurchlauf -------------------------------------- + +class FakeInspektor: + def __init__(self, aktiv=None, vorgemerkt=None, geplant=None): + self._aktiv, self._vorgemerkt, self._geplant = aktiv, vorgemerkt, geplant + + def active(self): + return self._aktiv + + def reserved(self): + return self._vorgemerkt + + def scheduled(self): + return self._geplant + + +class FakeCelery: + def __init__(self, inspektor): + self.control = self + self._inspektor = inspektor + + def inspect(self, timeout=None): + return self._inspektor + + +class FakeDb: + def __init__(self, offene, online=1): + self._offene, self._online = offene, online + self.aktualisierungen = [] + self.logs = [] + + def list_jobs_mit_status(self, stati): + return [j for j in self._offene if j["status"] in stati] + + def zaehle_online_worker(self, sekunden=120): + return self._online + + def update_job(self, job_id, **felder): + self.aktualisierungen.append((job_id, felder)) + + def add_log(self, level, source, message): + self.logs.append((level, message)) + + def utcnow(self): + return "jetzt" + + +# --- Gesamtdurchlauf -------------------------------------------------------- + +def test_laufender_job_wird_nicht_angetastet(): + """Der Fall, der wehtut: Job läuft seit Stunden und IST aktiv.""" + db = FakeDb([{"id": "lebt", "status": "transcoding", "title": "Akira"}]) + aktiv = {"celery@a": [{"name": "worker.tasks.transcode_files", "args": ["lebt"]}]} + celery = FakeCelery(FakeInspektor(aktiv=aktiv, vorgemerkt={}, geplant={})) + + bericht = zombies.raeume_zombies_auf(celery, db) + + assert bericht["aufgeraeumt"] == [] + assert db.aktualisierungen == [] + + +def test_echte_leiche_wird_auf_failed_gesetzt(): + db = FakeDb([{"id": "leiche", "status": "transcoding", "title": "Akira"}]) + # Der Knoten antwortet — er arbeitet nur an nichts. + celery = FakeCelery(FakeInspektor(aktiv={"celery@a": []}, vorgemerkt={}, geplant={})) + + bericht = zombies.raeume_zombies_auf(celery, db) + + assert bericht["aufgeraeumt"] == ["leiche"] + job_id, felder = db.aktualisierungen[0] + assert job_id == "leiche" + assert felder["status"] == "failed" + assert "Neu komprimieren" in felder["error"] + assert felder["finished_at"] == "jetzt" + + +def test_schweigender_worker_verhindert_jedes_urteil(): + """Zwei Worker gelten als online, nur einer antwortet — der andere könnte + genau diesen Job bearbeiten. Also: Finger weg.""" + db = FakeDb([{"id": "unklar", "status": "transcoding", "title": "X"}], online=2) + celery = FakeCelery(FakeInspektor(aktiv={"celery@a": []}, vorgemerkt={}, geplant={})) + + bericht = zombies.raeume_zombies_auf(celery, db) + + assert bericht["aufgeraeumt"] == [] + assert db.aktualisierungen == [] + assert "unvollständige Auskunft" in bericht["uebersprungen"] + + +def test_gar_keine_antwort_fuehrt_zu_keinem_urteil(): + db = FakeDb([{"id": "unklar", "status": "ripping", "title": "X"}], online=0) + celery = FakeCelery(FakeInspektor(aktiv=None, vorgemerkt=None, geplant=None)) + + bericht = zombies.raeume_zombies_auf(celery, db) + + assert bericht["aufgeraeumt"] == [] + assert db.aktualisierungen == [] + + +def test_fertige_jobs_werden_gar_nicht_betrachtet(): + db = FakeDb([ + {"id": "fertig", "status": "completed", "title": "A"}, + {"id": "kaputt", "status": "failed", "title": "B"}, + {"id": "wartet", "status": "pending", "title": "C"}, + ]) + celery = FakeCelery(FakeInspektor(aktiv={"celery@a": []}, vorgemerkt={}, geplant={})) + + bericht = zombies.raeume_zombies_auf(celery, db) + + # pending bleibt bewusst unberührt: die Aufgabe kann noch in der + # Warteschlange liegen und wird von selbst abgeholt. + assert bericht["geprueft"] == 0 + assert db.aktualisierungen == [] + + +def test_fehler_reisst_den_worker_start_nicht_mit(): + class KaputteDb(FakeDb): + def list_jobs_mit_status(self, stati): + raise RuntimeError("Postgres weg") + + db = KaputteDb([]) + bericht = zombies.raeume_zombies_auf(FakeCelery(FakeInspektor()), db) + assert "Postgres weg" in bericht["uebersprungen"] diff --git a/docker/worker/zombies.py b/docker/worker/zombies.py new file mode 100644 index 0000000..a02e8ab --- /dev/null +++ b/docker/worker/zombies.py @@ -0,0 +1,149 @@ +"""Erkennt Job-Leichen: Jobs, die behaupten zu laufen, an denen aber niemand arbeitet. + +Befund 25.07.2026 (Akira-UHD): Nach einem Absturz stand ein Job auf +`transcoding` bei 96 %, obwohl weder ein Prozess lief noch etwas in den +Celery-Queues stand. Folge für den Nutzer: kein Fortschritt, kein Download — +und der Knopf „Neu komprimieren" fehlte, weil `_kann_neu_komprimieren` +(api/main.py) `status == "failed"` verlangt. Der Job war damit unerreichbar, +obwohl die Rohdateien vollständig dalagen. + +## Die Leitregel: ohne vollständige Auskunft wird NICHTS angetastet + +Ein falsch als Leiche markierter Job ist teurer als eine stehengebliebene +Leiche. Deshalb drei Sicherungen: + +1. **Nur beim Worker-Start.** Da ist die Aussage eindeutig: „als ich hochkam, + lief nichts". Ein periodischer Lauf könnte einen Job erwischen, der legitim + in der Warteschlange wartet, weil alle Arbeitsplätze belegt sind — der ist + für `active()`/`reserved()` unsichtbar. Bewusst nicht gebaut. +2. **Gnadenfrist.** Nach einem Neustart stellt Celery unbestätigte Aufgaben + erneut zu. Erst abwarten, dann urteilen. +3. **Vollzähligkeit.** Es wird nur geurteilt, wenn mindestens so viele + Celery-Knoten antworten, wie laut Herzschlag online sind. Antwortet ein + beschäftigter Remote-Worker nicht rechtzeitig, wäre sein laufender Job + sonst eine falsche Leiche. + +Die Kernfunktionen sind bewusst rein (kein Postgres, kein Redis), damit die +Ampel sie ohne Infrastruktur prüfen kann. +""" + +# Zustände, die behaupten: hier arbeitet gerade jemand. +ARBEITS_STATI = ("ripping", "transcoding", "canceling") + +# Wartezeit nach dem Worker-Start, bevor geurteilt wird. Deckt die +# Wiederzustellung unbestätigter Aufgaben durch Celery ab. +GNADENFRIST_SEKUNDEN = 120 + +# Wie lange auf Antworten der Celery-Knoten gewartet wird. Großzügig, weil ein +# Worker mitten in einem 4K-Encode träge antwortet. +INSPEKT_TIMEOUT_SEKUNDEN = 10 + + +def belegte_job_ids(inspektionen, kandidaten) -> set: + """Welche der `kandidaten`-IDs kommen in irgendeiner Celery-Aufgabe vor? + + Bewusst als Textsuche über die ganze Inspektions-Antwort: `job_id` steht + bei `rip_disc` an ZWEITER, bei `transcode_files` an ERSTER Stelle, und + Celery liefert `args` je nach Version als Liste oder als Text. Eine + Positions-Auswertung wäre an beidem zerbrechlich. Eine Job-ID ist eine + 36-stellige UUID — Fehltreffer sind praktisch ausgeschlossen. + + `inspektionen` ist die Liste der Antworten (active/reserved/scheduled); + `None`-Einträge werden übersprungen. + """ + vorhandene = [i for i in (inspektionen or []) if i] + if not vorhandene: + return set() + text = repr(vorhandene) + return {jid for jid in kandidaten if jid and jid in text} + + +def antwortende_knoten(inspektionen) -> set: + """Namen aller Celery-Knoten, die auf die Inspektion geantwortet haben.""" + knoten = set() + for antwort in (inspektionen or []): + if antwort: + knoten.update(antwort.keys()) + return knoten + + +def auskunft_vollstaendig(anzahl_antworten: int, anzahl_bekannt: int) -> bool: + """Darf aus dieser Auskunft überhaupt ein Urteil folgen? + + Nein, wenn niemand geantwortet hat (dann wissen wir nichts), und nein, wenn + weniger Knoten antworten als laut Herzschlag online sind (dann fehlt uns + genau der Knoten, der den Job vielleicht gerade bearbeitet). + """ + if anzahl_antworten < 1: + return False + return anzahl_antworten >= anzahl_bekannt + + +def finde_zombies(offene_jobs, belegte) -> list: + """Jobs aus `offene_jobs`, zu denen keine Celery-Aufgabe existiert.""" + return [job for job in offene_jobs if job.get("id") not in belegte] + + +def fehlertext(job) -> str: + """Ehrlicher Klartext für die Job-Karte — was war, und was jetzt hilft.""" + zuletzt = job.get("status") or "unbekannt" + return ( + f"Abgebrochen: Der Vorgang lief nicht mehr (zuletzt gemeldet: {zuletzt}). " + "Beim Start des Workers war dazu weder ein Prozess noch eine Aufgabe in " + "der Warteschlange zu finden — meistens ein Absturz oder ein Neustart " + "mitten im Lauf. Die Rohdateien wurden NICHT gelöscht: mit " + "'Neu komprimieren' läuft die Kompression erneut, ohne die Disc noch " + "einmal zu rippen." + ) + + +def raeume_zombies_auf(celery_app, db) -> dict: + """Sucht Leichen und setzt sie ehrlich auf `failed`. Wirft nie. + + Rückgabe (auch für den Log): was geprüft und was getan wurde. + """ + bericht = {"geprueft": 0, "aufgeraeumt": [], "uebersprungen": ""} + try: + offene = db.list_jobs_mit_status(ARBEITS_STATI) + bericht["geprueft"] = len(offene) + if not offene: + return bericht + + inspektor = celery_app.control.inspect(timeout=INSPEKT_TIMEOUT_SEKUNDEN) + inspektionen = [inspektor.active(), inspektor.reserved(), inspektor.scheduled()] + + knoten = antwortende_knoten(inspektionen) + bekannt = db.zaehle_online_worker() + if not auskunft_vollstaendig(len(knoten), bekannt): + bericht["uebersprungen"] = ( + f"unvollständige Auskunft ({len(knoten)} von {bekannt} Worker " + "haben geantwortet) — es wird nichts als Leiche gewertet" + ) + db.add_log( + "info", "worker", + f"Zombie-Erkennung übersprungen: {bericht['uebersprungen']}", + ) + return bericht + + belegte = belegte_job_ids(inspektionen, [j.get("id") for j in offene]) + for job in finde_zombies(offene, belegte): + db.update_job( + job["id"], + status="failed", + error=fehlertext(job), + finished_at=db.utcnow(), + ) + db.add_log( + "warning", "worker", + f"Job {job['id']} ({job.get('title') or 'ohne Titel'}) stand auf " + f"'{job.get('status')}', es arbeitet aber niemand daran — " + "ehrlich auf 'failed' gesetzt. Rohdateien bleiben liegen.", + ) + bericht["aufgeraeumt"].append(job["id"]) + except Exception as e: # Erkennung darf den Worker-Start NIEMALS verhindern + bericht["uebersprungen"] = f"Fehler: {e}" + try: + db.add_log("warning", "worker", f"Zombie-Erkennung fehlgeschlagen: {e}") + except Exception: + pass + return bericht