e0cb7b3ddc
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>
191 lines
6.8 KiB
Python
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"]
|