Files
rippy/docker/worker/test_zombies.py
T
Hitonabi e0cb7b3ddc 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>
2026-07-25 21:04:19 +02:00

191 lines
6.8 KiB
Python

"""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"]