"""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. ## Warum Celery hier erst bei Bedarf entsteht (V2-4, 28.08.2026) Hier stand `celery_client = Celery(...)` auf Modulebene. Damit brauchte JEDER Import von `main.py` ein funktionierendes Celery — auch dann, wenn nie ein Rip angestoßen wird. Im Windows-Paket ist Celery bewusst NICHT enthalten (der Standalone-Betrieb hat keinen Broker), und die fertige EXE starb sofort beim Start: File "celery_client.py", ... celery_client = Celery("rippy_api", broker=REDIS_URL, ...) ModuleNotFoundError: No module named 'celery.fixups' Dasselbe Muster wie beim Store, der seine Engine beim Import baute: Was erst bei der ersten Benutzung gebraucht wird, soll auch erst dann entstehen. **Was das für den Windows-Betrieb bedeutet — ehrlich gesagt:** Die Oberfläche, die Laufwerks-Erkennung und alles Lesende laufen dort. Ein RIP anzustoßen geht noch nicht, weil die Zustellung weiterhin über Celery läuft; die Umstellung auf die LocalQueue steht in Etappe V2-5. Bis dahin sagt `start_rip()` das ausdrücklich, statt mit einem Importfehler zu sterben oder still nichts zu tun. """ import os REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0") _client = None # Wer Auftraege zustellt. None = Celery (verteilter Betrieb). # Der Standalone-Betrieb setzt hier seine eigene Zustellung ein — siehe # `zusteller_setzen()`. _zusteller = None class KeinBroker(RuntimeError): """Es gibt hier kein Celery — mit Ansage statt mit Importfehler.""" def hole_client(): """Der Celery-Client, beim ERSTEN Zugriff gebaut.""" global _client if _client is None: try: from celery import Celery except ImportError as e: raise KeinBroker( "Celery ist in dieser Installation nicht enthalten. Rippy läuft " "hier im Standalone-Betrieb; das Anstoßen von Rips über einen " "Broker ist damit nicht möglich (Umstellung auf die lokale " "Auftrags-Queue: Etappe V2-5)." ) from e _client = Celery("rippy_api", broker=REDIS_URL, backend=REDIS_URL) return _client def __getattr__(name): """`celery_client` von außen — baut den Client bei Bedarf. Damit bleiben die bestehenden `from celery_client import celery_client` unverändert gültig, ohne dass der Import schon einen Broker verlangt. """ if name == "celery_client": return hole_client() raise AttributeError(f"module {__name__!r} has no attribute {name!r}") def zusteller_setzen(funktion) -> None: """Legt fest, WER Auftraege bekommt. Ohne Aufruf geht alles an Celery — das ist der Docker-Betrieb, unveraendert. Der Standalone-Betrieb (Windows-App, Headless-Linux) setzt hier seine lokale Auftrags-Queue ein; dann laeuft die Arbeit im selben Prozess. Die Unterschrift ist die von `abschicken`: funktion(task_name, args, queue=None) -> irgendetwas Das ist bewusst dieselbe Form, die Celery hat. So muss keine Aufrufstelle wissen, in welchem Betrieb sie gerade laeuft. """ global _zusteller _zusteller = funktion def abschicken(task_name: str, args: list, queue: str = None): """Einen Auftrag zustellen — an Celery oder an die lokale Queue. EINE Stelle fuer alle drei Auftragsarten (rip_disc, transcode_files, scan_tracks). Vorher rief jede Aufrufstelle `send_task` selbst auf; damit haette der Standalone-Betrieb an drei Stellen umgebogen werden muessen — und beim naechsten Auftragstyp an einer vierten. """ if _zusteller is not None: return _zusteller(task_name, args, queue) kwargs = {"args": args} if queue: kwargs["queue"] = queue return hole_client().send_task(task_name, **kwargs) 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: from celery.utils import worker_direct antworten = hole_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-Auftrag los (Task-Name aus worker/tasks.py).""" return abschicken("worker.tasks.rip_disc", [device_path, job_id, target_dir])