feat(worker): Zombie-Erkennung - Jobs, an denen niemand arbeitet
Nach einem Absturz oder Rebuild blieb ein Job auf 'transcoding' stehen, obwohl weder ein Prozess lief noch etwas in den Queues stand. Folge: keine Anzeige, kein Download - und der Knopf "Neu komprimieren" fehlte, weil _kann_neu_komprimieren (api/main.py) status=='failed' verlangt. Der Job war unerreichbar, obwohl die Rohdateien vollstaendig dalagen. zombies.py haelt beim Worker-Start die Jobs in 'ripping'/'transcoding'/ 'canceling' gegen Celerys active/reserved/scheduled und setzt sie ehrlich auf 'failed', wenn niemand daran arbeitet. Drei Sicherungen, weil ein falsch getoeteter Job teurer ist als eine stehengebliebene Leiche: - nur beim Start (da ist "es lief nichts" eindeutig; ein periodischer Lauf koennte einen Job erwischen, der legitim in der Warteschlange wartet) - 120 s Gnadenfrist (Celery stellt unbestaetigte Aufgaben erneut zu) - Vollzaehligkeit: antworten weniger Knoten als laut Herzschlag online sind, wird NICHTS gewertet - sonst waere der laufende Job eines beschaeftigten Remote-Workers eine falsche Leiche Die Job-ID wird per Textsuche ueber die Inspektions-Antwort gefunden, nicht per Position: sie steht bei rip_disc an zweiter, bei transcode_files an erster Stelle, und Celery liefert args je nach Version als Liste oder Text. 13 Tests ohne Postgres/Redis, darunter "laufender Job wird nicht angetastet" und "schweigender Worker verhindert jedes Urteil". Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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))
|
||||
|
||||
@@ -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"]
|
||||
@@ -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
|
||||
Reference in New Issue
Block a user