"""Dünner Celery-Client: die API SCHICKT Tasks an den Worker, führt sie nie aus. Bis 23.07. gab es überhaupt keinen Code-Pfad, der je einen Rip auslöste — kein POST /jobs, kein udev-Daemon. Dieser Client schließt die Lücke. """ import os from celery import Celery from celery.utils import worker_direct REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0") celery_client = Celery("rippy_api", broker=REDIS_URL, backend=REDIS_URL) def transcode_queue(node: str = None): """Ziel-Queue für die Kompression: der GEWÄHLTE Worker (worker_direct) wenn er gerade online ist, sonst die geteilte transcode-Queue. So kann ein Job gezielt einen Encoder ansprechen — fällt der Worker aber weg, bleibt die Kompression nicht in einer toten Queue hängen, sondern landet bei irgendeinem freien Worker. """ if not node: return "transcode" try: antworten = celery_client.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" def start_rip(device_path: str, job_id: str, target_dir: str = None): """Schickt den Rip-Task an den Worker (Task-Name aus worker/tasks.py).""" return celery_client.send_task( "worker.tasks.rip_disc", args=[device_path, job_id, target_dir] )