19f3dc5330
Ampel / ampel (push) Successful in 30s
Der Commander: "Ausserdem ist gerade mitten im Rip das Laufwerk ausgegangen...
glaube ich zumindest." Gemessen war es etwas anderes, und der Fund ist groesser
als der Vorfall.
WAS MESSBAR WAR:
Job 2182d525 status=running progress=12 (Rip gestartet 14:00:24)
makemkvcon laeuft NICHT (per /proc geprueft, PID gegengeprueft)
Rohdatei 5.167.382.528 Bytes, waechst in 10 s nicht
Laufwerk Status 4 (Disc drin), /dev/sr0 + /dev/sg1 da
Worker-Start 14:07:02 <- mein `docker compose up -d --build`
Das Laufwerk ist also NICHT ausgegangen. Der Rip wurde von MEINEM Deploy
getoetet: `up -d --build` baut den worker-Container neu, und der laufende Rip
stirbt mit ihm. Genau davor warnt der SAVEPOINT seit v3.18 - die Warnung half
nichts, weil sie niemand liest und nichts sie prueft.
DER EIGENTLICHE FUND: Die Zombie-Erkennung lief um 14:09:04 und meldete
`{'geprueft': 0, 'aufgeraeumt': []}` - obwohl der tote Job direkt vor ihr lag.
Ursache:
ARBEITS_STATI = ("ripping", "transcoding", "canceling") # zombies.py
db.update_job(job_id, status="running", ...) # tasks.py - der Rip
Der Rip setzt "running", gesucht wurde "ripping". Dieser Wert steht
ausschliesslich in Celerys Task-META und NIE in einer Job-Zeile (nachgeprueft:
kein einziger Schreiber im ganzen Baum). Die Zombie-Erkennung aus v3.14 wurde
gebaut, um genau einen abgestuerzten Rip zu finden - und hat ihn nie gesehen.
Besonders tueckisch: `geprueft: 0` sah bei jedem Worker-Start wie "nachgesehen,
alles gesund" aus, waehrend sie nach einem Status suchte, den es nicht gibt.
Deshalb blieb der Job auf "processing 12 %" stehen - mit einer Restzeit-Schaetzung
von 1 h 30 min obendrauf, die es fuer einen toten Prozess nicht geben duerfte.
GEBAUT:
* "running" in ARBEITS_STATI.
* Ein Test, der das Auseinanderlaufen MECHANISCH verhindert: Er liest tasks.py,
sammelt jeden Status, den der Worker per db.update_job in eine Job-Zeile
schreibt, und verlangt, dass jeder davon entweder ein Arbeitsstatus oder ein
Endzustand ist. Ein Kommentar haette das nicht verhindert.
* GET /health/arbeit + eine Sperre in deploy.sh: Laeuft ein Job, bricht der
Deploy ab (uebersteuerbar mit RIPPY_TROTZDEM=1 - dann ist es eine
Entscheidung und kein Versehen). Ein Satz Code gegen eine verlorene Stunde.
* Server-Status zeigt jetzt die eingehaengten FREIGABEN mit freiem Platz, nicht
nur die Container-Platte (Commander-Wunsch). Genau dort liegen die Rohdaten,
und bei externem Encoden muessen sie dort liegen - wer wissen wollte, ob noch
Platz fuer eine Disc ist, sah die falsche Zahl.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
242 lines
8.9 KiB
Python
242 lines
8.9 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"]
|
|
|
|
|
|
# --- Die Luecke, die die Erkennung nutzlos machte (Befund 26.07.2026) --------
|
|
|
|
|
|
def test_running_gilt_als_arbeitsstatus():
|
|
"""DER Fehler: Die Erkennung suchte "ripping", der Rip setzt aber "running".
|
|
Ergebnis: Ein abgestuerzter RIP wurde NIE gefunden - und die Meldung
|
|
`{'geprueft': 0}` sah bei jedem Worker-Start wie Gesundheit aus."""
|
|
assert "running" in zombies.ARBEITS_STATI
|
|
assert "transcoding" in zombies.ARBEITS_STATI
|
|
assert "canceling" in zombies.ARBEITS_STATI
|
|
|
|
|
|
def test_abgestuerzter_rip_wird_gefunden():
|
|
"""Der konkrete Fall vom 26.07.2026: Job 2182d525 stand auf running/12 %,
|
|
kein makemkvcon lief, die Rohdatei wuchs nicht mehr - und die Erkennung
|
|
pruefte null Jobs."""
|
|
offene = [{"id": "2182d525", "status": "running", "title": "Akira"}]
|
|
assert zombies.finde_zombies(offene, set()) == offene
|
|
|
|
|
|
def test_arbeitsstati_deckt_ab_was_der_worker_wirklich_schreibt():
|
|
"""Mechanische Sperre gegen genau dieses Auseinanderlaufen.
|
|
|
|
Liest tasks.py und sammelt jeden Status, den der Worker per db.update_job in
|
|
eine Job-ZEILE schreibt. Jeder davon muss entweder ein Arbeitsstatus sein
|
|
oder ein Endzustand - sonst gibt es wieder einen Zustand, den niemand
|
|
aufraeumt. Ein Kommentar haette das nicht verhindert; dieser Test schon.
|
|
"""
|
|
import os
|
|
import re
|
|
|
|
pfad = os.path.join(os.path.dirname(os.path.abspath(__file__)), "tasks.py")
|
|
quelle = open(pfad, encoding="utf-8").read()
|
|
|
|
# Nur die Aufrufe, die wirklich die Job-Zeile aendern.
|
|
geschrieben = set()
|
|
for aufruf in re.finditer(r"db\.update_job\((?:[^()]|\([^()]*\))*\)", quelle):
|
|
for treffer in re.finditer(r'status\s*=\s*"([a-z]+)"', aufruf.group(0)):
|
|
geschrieben.add(treffer.group(1))
|
|
|
|
endzustaende = {"completed", "failed"}
|
|
unbeaufsichtigt = geschrieben - set(zombies.ARBEITS_STATI) - endzustaende
|
|
assert not unbeaufsichtigt, (
|
|
f"Diese Job-Status schreibt der Worker, aber niemand raeumt sie auf: "
|
|
f"{sorted(unbeaufsichtigt)}. Entweder in ARBEITS_STATI aufnehmen oder "
|
|
f"als Endzustand behandeln."
|
|
)
|
|
# Gegenprobe, dass der Test wirklich etwas gesehen hat
|
|
assert "running" in geschrieben
|