Files
rippy/docker/worker/celery_app.py
T
Hitonabi 4ba02047db
Ampel / ampel (push) Successful in 28s
Drei Praxis-Bugs: tote NAS-Mounts reparierbar, HandBrake-Versionen konsistent, Encoder-Wahl beim Rip
1) Speicher-Mounts robust (Befund: toter CIFS-Mount nach NAS-Ausfall/Rebuild —
   mounted:false, verschwand aus 'Verfuegbare Ziele', Neu-Anlegen -> 409, man
   sass fest):
   - mounts.py: ist_erreichbar() (listdir, soft-Mount bricht schnell ab),
     ist_gemountet() faengt OSError toter Mounts, aushaengen() mit
     umount -l Fallback, reparieren() (lazy abhaengen + frisch mounten).
   - /storage-mounts liefert 'reachable'; POST bei existierendem Namen:
     aktiv -> 409, tot -> automatische Reparatur mit neuen Angaben;
     neuer POST /storage-mounts/{name}/repair (gespeicherte Zugangsdaten).
   - /storage-targets crasht nicht mehr an totem Mount (os.path.ismount
     OSError abgefangen).
   - UI: eigene 'Netzwerk-Mounts'-Liste mit Status (aktiv/nicht erreichbar/
     getrennt) + Reparieren- und Entfernen-Knopf — tote Mounts sind sichtbar
     und wiederherstellbar statt zu verschwinden.

2) HandBrake-Versionen konsistent (Befund: Docker 1.6.1, Windows-Skript
   fest 1.9.2, Update-Check meldet 1.11.2 — verwirrend):
   - Windows-Installer zieht jetzt DYNAMISCH die neueste Version (GitHub
     latest, Fallback 1.11.2) — passt zum Update-Check.
   - Update-UI erklaert klar: Docker = stabiles Debian-Paket (bewusst aelter,
     kein Fehler), Windows = neueste. MakeMKV-Update zeigt den Befehl.

3) Encoder-/Worker-Auswahl beim Rip (Feature):
   - Celery worker_direct=True: jeder Worker konsumiert zusaetzlich seine
     Direkt-Queue. API-Helper transcode_queue(node) routet gezielt an den
     gewaehlten Worker, faellt aber sicher auf die geteilte transcode-Queue
     zurueck, wenn er offline ist (kein Haengenbleiben).
   - /capabilities liefert den Celery-Node je Worker; POST /jobs nimmt
     transcode_node (-> Job-meta); rip_disc + retry-transcode routen danach.
   - Rip-Dialog: Encoder-/Worker-Dropdown, sichtbar ab 2 Online-Workern.
   - ping_worker-Task zum Verifizieren des gezielten Routings.
   - Nebenfund gefixt: DeviceDiscovery leitete die Titel-Auswahl (titles)
     gar nicht an die API weiter — jetzt titles + transcode_node.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-24 18:19:57 +02:00

67 lines
2.1 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
@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()