"""Celery-Adapter: dieselben Task-Namen wie bisher, der Ablauf liegt in ablauf.py. ## Was sich geaendert hat (V2-4, 28.08.2026) — und was NICHT Nicht geaendert: die Task-Namen (`worker.tasks.rip_disc`, `worker.tasks.transcode_files`, `worker.tasks.scan_tracks`, `worker.tasks.ping_worker`), die Argumente, die Queues, das Verhalten. Ein laufender Docker-Betrieb merkt von diesem Umbau nichts — und ein bereits installierter Windows-Worker aus v1 auch nicht. Geaendert: Der ABLAUF steht nicht mehr hier. Diese Datei importierte auf Modulebene `celery_app`, und damit war der gesamte Rip-Vorgang an einen Broker gebunden. Im nativen Windows-Betrieb — der bewusst keinen hat — war er ueberhaupt nicht ladbar. Uebrig bleibt hier genau das, was WIRKLICH Celery ist: die Task-Huellen und die Wahl der Ziel-Queue. ## Die zwei Rueckrufe `ablauf.rippen()` weiss nicht, wer zuhoert. Diese Datei reicht ihm die Celery-Fassung herein: melde -> self.update_state(...) wie bisher weiterreichen -> transcode_files.apply_async(queue=...) Der Standalone-Laeufer reicht stattdessen seine eigenen herein — oder gar keine, dann komprimiert `ablauf` gleich selbst. """ import os import ablauf from celery_app import celery_app # Namen, die andere Module bisher aus tasks.py geholt haben. Sie liegen jetzt # in ablauf.py; hier stehen sie weiter zur Verfuegung, damit kein Aufrufer # angefasst werden muss. RIP_FERTIG = ablauf.RIP_FERTIG API_URL = ablauf.API_URL pfad_lokal = ablauf.pfad_lokal unter_wurzel = ablauf.unter_wurzel erster_vorhandener_ordner = ablauf.erster_vorhandener_ordner def _transcode_queue(node: str): """Ziel-Queue für die Kompression (siehe celery_client.transcode_queue): gewählter Worker via worker_direct, wenn online — sonst geteilte Queue.""" if not node: return "transcode" try: from celery.utils import worker_direct antworten = celery_app.control.ping(timeout=1.0) or [] online = {k for antwort in antworten for k in antwort.keys()} if node in online: return worker_direct(node) except Exception: pass return "transcode" @celery_app.task(bind=True, name="worker.tasks.rip_disc") def rip_disc(self, device_path: str, job_id: str, target_dir: str = None): """Rippt die Disc. Der Ablauf steht in ablauf.rippen().""" def melde(zustand: dict) -> None: self.update_state(state="PROGRESS", meta=zustand) def weiterreichen(job: str, raw_dir: str, final_dir: str, knoten: str) -> None: # Die Kompression als eigener Task — an den im Rip-Dialog GEWÄHLTEN # Worker (worker_direct), sonst an die geteilte transcode-Queue # (irgendein freier Worker, inkl. Remote-GPU). ziel_queue = _transcode_queue(knoten) if knoten and ziel_queue != "transcode": ablauf.db.add_log("info", "worker", f"Job {job}: Kompression gezielt an {knoten}") transcode_files.apply_async(args=[job, raw_dir, final_dir], queue=ziel_queue) return ablauf.rippen(device_path, job_id, target_dir, melde=melde, weiterreichen=weiterreichen) @celery_app.task(bind=True, name="worker.tasks.transcode_files") def transcode_files(self, job_id: str, raw_dir: str, final_dir: str): """Komprimiert die Rohdaten. Der Ablauf steht in ablauf.komprimieren().""" return ablauf.komprimieren(job_id, raw_dir, final_dir) @celery_app.task(name="worker.tasks.scan_tracks") def scan_tracks(device_path: str): return ablauf.scan_tracks(device_path) @celery_app.task(name="worker.tasks.ping_worker") def ping_worker(): return ablauf.ping_worker() # Der Anzeigename dieses Knotens — wird vom Herzschlag in celery_app.py # benutzt und stand bisher hier. WORKER_NAME = os.getenv("WORKER_NAME", "")