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>
91 lines
2.9 KiB
Python
91 lines
2.9 KiB
Python
import os
|
|
import socket
|
|
|
|
from celery import Celery
|
|
from celery.signals import worker_ready
|
|
|
|
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
|
|
|
|
celery_app = Celery(
|
|
"rippy_worker",
|
|
broker=REDIS_URL,
|
|
backend=REDIS_URL,
|
|
include=["tasks"]
|
|
)
|
|
|
|
celery_app.conf.update(
|
|
task_serializer="json",
|
|
accept_content=["json"],
|
|
result_serializer="json",
|
|
timezone="UTC",
|
|
enable_utc=True,
|
|
broker_connection_retry_on_startup=True,
|
|
# worker_direct: jeder Worker konsumiert zusätzlich eine eigene
|
|
# Direkt-Queue (<nodename>.dq). Damit kann die API die Kompression
|
|
# gezielt an EINEN gewählten Worker schicken (Encoder-Auswahl im
|
|
# Rip-Dialog) — ohne Wahl läuft sie weiter über die geteilte
|
|
# transcode-Queue (irgendein freier Worker).
|
|
worker_direct=True,
|
|
)
|
|
|
|
|
|
# Import auf Modulebene: im worker_ready-Signal ist /app nicht mehr
|
|
# 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
|
|
def melde_faehigkeiten(**kwargs):
|
|
"""Beim Start + minütlich: Encoder-Fähigkeiten und Lebenszeichen melden.
|
|
|
|
last_seen ist die Basis der Erreichbarkeits-Anzeige im UI — ein Worker,
|
|
der sich 2 Minuten nicht meldet, gilt als offline.
|
|
"""
|
|
import threading
|
|
import time
|
|
|
|
# Anzeigename: WORKER_NAME (compose/.env) macht die Maschine im UI
|
|
# zuordenbar UND stabil — der Container-Hostname wechselt bei jedem
|
|
# Rebuild und hinterließ Offline-Leichen in der Worker-Liste.
|
|
anzeigename = os.getenv("WORKER_NAME") or socket.gethostname()
|
|
|
|
def herzschlag():
|
|
while True:
|
|
try:
|
|
db.init_db()
|
|
db.save_worker(
|
|
anzeigename,
|
|
caps.erkenne_encoder(),
|
|
caps.werkzeug_versionen(),
|
|
)
|
|
except Exception as e: # DB weg → weiterversuchen, nicht sterben
|
|
print(f"Fähigkeiten-Meldung fehlgeschlagen: {e}")
|
|
time.sleep(60)
|
|
|
|
threading.Thread(target=herzschlag, daemon=True, name="worker-herzschlag").start()
|