diff --git a/docker/api/celery_client.py b/docker/api/celery_client.py index 79497c5..0dbfca0 100644 --- a/docker/api/celery_client.py +++ b/docker/api/celery_client.py @@ -31,6 +31,11 @@ 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.""" @@ -64,6 +69,40 @@ def __getattr__(name): 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. @@ -87,7 +126,5 @@ def transcode_queue(node: str = None): 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 hole_client().send_task( - "worker.tasks.rip_disc", args=[device_path, job_id, target_dir] - ) + """Schickt den Rip-Auftrag los (Task-Name aus worker/tasks.py).""" + return abschicken("worker.tasks.rip_disc", [device_path, job_id, target_dir]) diff --git a/docker/api/main.py b/docker/api/main.py index d8ad631..0f5edb1 100644 --- a/docker/api/main.py +++ b/docker/api/main.py @@ -1214,8 +1214,7 @@ async def scan_tracks_starten(name: str): raise HTTPException(status_code=409, detail="Auf diesem Laufwerk läuft gerade ein Job") await asyncio.to_thread(db.save_settings, {"status": "running"}, f"tracks:{device_path}") - celery_anbindung.hole_client().send_task( - "worker.tasks.scan_tracks", args=[device_path]) + celery_anbindung.abschicken("worker.tasks.scan_tracks", [device_path]) return {"status": "scanning"} @@ -1273,11 +1272,9 @@ async def retry_transcode(job_id: str): except ValueError: meta = {} from celery_client import transcode_queue - celery_anbindung.hole_client().send_task( - "worker.tasks.transcode_files", - args=[job_id, raw_dir, final_dir], - queue=transcode_queue(meta.get("transcode_node")), - ) + celery_anbindung.abschicken( + "worker.tasks.transcode_files", [job_id, raw_dir, final_dir], + queue=transcode_queue(meta.get("transcode_node"))) await asyncio.to_thread(db.update_job, job_id, status="transcoding", progress=0, error=None) await asyncio.to_thread(db.add_log, "info", "api", f"Job {job_id}: Kompression neu eingereiht") return {"id": job_id, "status": "transcoding"} @@ -1833,6 +1830,90 @@ async def browse_mkdir(request: MkdirRequest): return {"path": ziel} +# ═══════════════════════════════════════════════════════════════════════ +# Werkzeuge: Bestand, Update-Stand, Beschaffung — Etappe V2-4 +# ═══════════════════════════════════════════════════════════════════════ +# +# WOFUER: Damit Rippy unter Windows eigenstaendig arbeiten kann, muss es +# seine Werkzeuge nicht nur FINDEN, sondern auch HOLEN und aktuell halten +# koennen. Im Container ist das anders — dort stecken sie im Image, und ein +# Update ist ein Rebuild (siehe /system/updates, das bleibt unveraendert). +# +# WARUM EIGENE ROUTEN statt /system/updates zu erweitern: /system/updates +# beantwortet "welche Version gibt es draussen". Diese hier beantworten +# "was liegt auf DIESER Maschine, wo, und was kann ich dagegen tun". Zwei +# verschiedene Fragen; eine Route, die beides taete, koennte keine davon +# ehrlich beantworten. + + +@app.get("/system/werkzeuge") +async def system_werkzeuge(): + """Was ist installiert, wo, in welcher Fassung — und gibt es Neueres? + + `neueste` ist leer, wenn die Quelle nicht antwortet. Dann steht dort + ausdruecklich NICHT "Update verfuegbar": Eine Nichtauskunft ist keine + Aussage ueber die Welt. + """ + from rippy.tools import beschaffen as werkzeug_beschaffung + + einstellungen = await asyncio.to_thread(db.get_settings) + eingestellt = (einstellungen or {}).get("werkzeugPfade") or {} + lage = await asyncio.to_thread(werkzeug_beschaffung.lage, eingestellt) + return { + "werkzeuge": lage, + # Wohin Rippy selbst installiert — im UI sichtbar, damit niemand + # raten muss, wo die geholten Dateien landen. + "ordner": werkzeug_beschaffung.katalog.werkzeug_ordner(), + "windows": os.name == "nt", + } + + +@app.post("/system/werkzeuge/{name}/holen") +async def system_werkzeug_holen(name: str): + """Holt ein Werkzeug und installiert es. + + HandBrake wird vollstaendig automatisch geholt (GitHub-Release, ein ZIP + mit einer .exe — kein Installer, keine Administratorrechte). + + MakeMKV wird NICHT mitgeliefert: Rippy laedt die offizielle Datei vom + Hersteller und startet sie. Das ist dieselbe Black-Box-Trennung, die + KONZEPT.md § 6 fuer den Container festhaelt — MakeMKV bleibt ein fremdes + Programm, Rippy nimmt nur die Handgriffe ab. + """ + from rippy.tools import beschaffen as werkzeug_beschaffung + + if name not in werkzeug_beschaffung.katalog.WERKZEUGE: + raise HTTPException( + status_code=404, + detail=f"Unbekanntes Werkzeug '{name}'. Bekannt: " + + ", ".join(sorted(werkzeug_beschaffung.katalog.WERKZEUGE))) + + meldungen = [] + + def fortschritt(text, anteil=None): + if text: + meldungen.append(text) + + try: + if name == "handbrake": + pfad = await asyncio.to_thread( + werkzeug_beschaffung.handbrake_holen, None, fortschritt) + else: + pfad = await asyncio.to_thread( + werkzeug_beschaffung.makemkv_holen, "", None, fortschritt, False) + except werkzeug_beschaffung.BeschaffungsFehler as e: + # Klartext statt Traceback: Der haeufigste Grund ist eine nicht + # erreichbare Quelle, und das ist nichts, was der Nutzer im Code + # suchen sollte. + await asyncio.to_thread(db.add_log, "error", "werkzeuge", str(e)) + raise HTTPException(status_code=502, detail=str(e)) from e + + await asyncio.to_thread( + db.add_log, "info", "werkzeuge", + f"{werkzeug_beschaffung.katalog.WERKZEUGE[name]['titel']} beschafft: {pfad}") + return {"name": name, "pfad": pfad, "meldungen": meldungen} + + @app.get("/system/updates") async def system_updates(): """Update-Check für die Kern-Werkzeuge (Einstellungen → System). diff --git a/docker/api/test_mounts_helpers.py b/docker/api/test_mounts_helpers.py index dc8517e..49334c6 100644 --- a/docker/api/test_mounts_helpers.py +++ b/docker/api/test_mounts_helpers.py @@ -145,7 +145,11 @@ def test_erzeugtes_mapping_uebersetzt_den_echten_fehlerfall(monkeypatch): # Der Worker liegt neben der API im Repo; kein geteiltes Paket zwischen # den Containern, deshalb per Pfad laden statt importieren. worker_tasks = os.path.join( - os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "worker", "tasks.py" + # Seit V2-4 steht der Ablauf in ablauf.py; tasks.py ist nur noch + # die Celery-Huelle. Dieser Waechter muss dorthin schauen, wo der + # Code WIRKLICH steht — sonst prueft er eine leere Datei und ist + # gruen, ohne etwas zu beweisen. + os.path.dirname(os.path.dirname(os.path.abspath(__file__))), "worker", "ablauf.py" ) spec = importlib.util.spec_from_file_location("_worker_tasks_pfad", worker_tasks) quelltext = open(worker_tasks, encoding="utf-8").read() diff --git a/docker/api/test_phasen.py b/docker/api/test_phasen.py index 8c29016..0442bdd 100644 --- a/docker/api/test_phasen.py +++ b/docker/api/test_phasen.py @@ -55,7 +55,11 @@ def test_der_schluessel_heisst_im_worker_genauso(): import re quelle = ( - pathlib.Path(__file__).resolve().parents[1] / "worker" / "tasks.py" + # Seit V2-4 steht der Ablauf in ablauf.py; tasks.py ist nur noch + # die Celery-Huelle. Dieser Waechter muss dorthin schauen, wo der + # Code WIRKLICH steht — sonst prueft er eine leere Datei und ist + # gruen, ohne etwas zu beweisen. + pathlib.Path(__file__).resolve().parents[1] / "worker" / "ablauf.py" ).read_text(encoding="utf-8") # Der Worker schreibt die Marke über eine Konstante — deren Wert muss hier # ankommen. diff --git a/docker/worker/ablauf.py b/docker/worker/ablauf.py new file mode 100644 index 0000000..44ef108 --- /dev/null +++ b/docker/worker/ablauf.py @@ -0,0 +1,1018 @@ +"""Der Ablauf eines Jobs — rippen, komprimieren, ablegen. OHNE Celery. + +## Warum es diese Datei gibt (V2-4, 28.08.2026) + +Bis hierher stand der ganze Ablauf in `tasks.py`, und die Datei importierte +auf Modulebene `celery_app`. Damit war der Ablauf an einen Broker gebunden — +und im nativen Windows-Betrieb, der bewusst KEINEN Broker hat, ueberhaupt +nicht ladbar. Rippy konnte dort die Oberflaeche zeigen und Laufwerke lesen, +aber nicht rippen. Genau das war die Luecke. + +Dabei hing der Ablauf an GENAU DREI Stellen an Celery — in 234 Zeilen: + + self.update_state(...) Zwischenstand melden + _transcode_queue(node) Ziel-Queue waehlen + transcode_files.apply_async(...) Kompression weiterreichen + +Alle drei sind Fragen der ZUSTELLUNG, nicht des Ablaufs. Sie sind jetzt +Rueckrufe, die der Aufrufer stellt: + + melde(zustand) wohin der Zwischenstand geht + weiterreichen(job, roh, ziel, knoten) wer komprimiert + +`tasks.py` reicht die Celery-Fassung herein, der Standalone-Laeufer die +lokale. **Ohne `weiterreichen` komprimiert dieselbe Funktion gleich selbst** — +genau das, was ein Ein-Prozess-Rippy braucht. + +## Was das NICHT ist + +Kein Neuschreiben. Der Rumpf ist der aus `tasks.py`, Zeile fuer Zeile — samt +aller Befunde, die dort ueber Monate hineingewandert sind (Platzpruefung, +Zombie-Marken, Pfad-Uebersetzung fuer native Worker, Serien-Zuordnung). Was +sich geaendert hat, sind die drei Stellen oben. + +Das ist `core/pipeline.py` aus KONZEPT-V2.md § 2.1, nur noch am alten Platz: +Der Umzug nach `rippy/` kommt, wenn die flachen Worker-Importe aufgeloest +sind (`import db`, `import medien`, `import ripping`). +""" + +import errno +import glob +import json +import os +import posixpath +import shutil +import time + +import requests + +from rippy import store as db +from rippy.rip import makemkv_daten +import medien +from rippy.core import notify + +# Der Laufwerks-Treiber DIESER Maschine — Linux (ioctl) oder Windows (Win32). +# +# ⚠️ Hier stand bis V2-4: +# +# try: +# from rippy.drives.detection import detect_disc_type, disc_size_bytes +# except ImportError: # Windows: kein fcntl +# detect_disc_type = None +# +# Das war fuer den NATIVEN WINDOWS-WORKER gedacht, der nur komprimiert. Es +# haette aber auch den eigenstaendigen Windows-Rippy getroffen: Der +# Linux-Treiber zieht ueber detection.py das Modul `fcntl` nach, das es unter +# Windows nicht gibt — also waere detect_disc_type dort IMMER None gewesen, +# und Rippen unten hart verriegelt. Rippy haette unter Windows alles anzeigen +# und nichts rippen koennen, ohne dass irgendwo ein Fehler stuende. +# +# Jetzt fragt der Ablauf den Port: Welcher Treiber es kann, entscheidet +# rippy/drives/__init__.py. `None` bleibt der ehrliche Wert fuer den Fall, +# dass diese Maschine wirklich keine Laufwerke lesen kann. +try: + from rippy import drives as _laufwerks_schicht + + _treiber = _laufwerks_schicht.treiber() + detect_disc_type = _treiber.detect_disc_type + disc_size_bytes = _treiber.disc_size_bytes +except (ImportError, AttributeError): + detect_disc_type = None + disc_size_bytes = None + +from ripping import ( + RIP_OUTPUT_DIR, + RipAbbruch, + komprimieren_fuer, + lies_datei_dauer, + lies_titel_info, + preset_fuer, + rip_cd, + rip_video, + run_handbrake, + sprachen_zusammenfassen, + sprachliste, + wirf_disc_aus, +) + +API_URL = os.getenv("API_URL", "http://api:8000") + +# Phasen-Marke in den Job-Metadaten: war der RIP fertig, als es schiefging? +# +# Sobald `status = "failed"` in der Zeile steht, ist die Phase sonst +# unwiederbringlich fort — und genau die entscheidet, was danach hilft: Nach +# einem toten Transcode genügt „Neu komprimieren"; nach einem toten Rip liegt +# nur ein Bruchstück da (Vorfall 26.07.2026: 5,1 GB von rund 40 GB) und es muss +# neu gerippt werden. Gelesen wird die Marke in api/phasen.py — der Name steht +# in beiden Dateien und wird von test_phasen.py gegeneinander geprüft, weil es +# kein geteiltes Paket zwischen den Containern gibt. +RIP_FERTIG = "rip_fertig" + + +def pfad_lokal(pfad: str, mapping: str = None) -> str: + """Übersetzt Rippy-Container-Pfade für native Worker (pure Funktion). + + Der native Windows-Worker sieht /app/temp und /app/media nicht — er + mountet Rippys Freigabe als Netzlaufwerk und setzt RIPPY_PATH_MAP, + z. B. "/app/media=Z:\\media;/app/temp=Y:\\temp". Ohne Mapping (Docker- + Worker: identische Mounts) kommt der Pfad unverändert zurück. + """ + mapping = mapping if mapping is not None else os.getenv("RIPPY_PATH_MAP", "") + if not pfad or not mapping: + return pfad + for paar in mapping.split(";"): + if "=" not in paar: + continue + quelle, ziel = paar.split("=", 1) + if quelle and pfad.startswith(quelle): + rest = pfad[len(quelle):] + if "\\" in ziel: + rest = rest.replace("/", "\\") + return ziel + rest + return pfad + +RAW_DIR = os.getenv("RAW_DIR", "/app/temp/raw") +MEDIA_ROOT = "/app/media" + +# Wie oft während einer Kompression nachgesehen wird, ob der Nutzer abgebrochen +# hat. Eine DB-Abfrage alle paar Sekunden ist nichts gegen einen Encode, der +# Stunden läuft — und „Abbrechen" fühlt sich damit sofort an. +ABBRUCH_INTERVALL_SEKUNDEN = 5 + + +def unter_wurzel(pfad: str, wurzel: str) -> bool: + """Liegt `pfad` wirklich unterhalb von `wurzel` (oder IST es die Wurzel)? + + Pure Funktion, testbar. Ein nacktes `startswith()` genügt nicht: + „/app/media-boese/x" beginnt mit „/app/media", liegt aber außerhalb + (Befund 25.07.2026 bei der Durchsicht). Immer „/" als Trenner — das sind + Container-Pfade, auch wenn ein nativer Windows-Worker das Modul lädt. + + Gleichlautend in docker/api/main.py; es gibt kein geteiltes Paket zwischen + den Containern. + """ + if not pfad or not wurzel: + return False + sauber = wurzel.rstrip("/") or "/" + return pfad == sauber or pfad.startswith(sauber + "/") + + +def _zielbasis(target_dir, disc_type: str) -> str: + """Ablagebasis: vom Nutzer gewähltes Ziel (validiert) oder Standard. + + posixpath statt os.path — aus demselben Grund wie in _arbeitsverzeichnis: + das sind IMMER Container-Pfade. os.path.normpath macht unter Windows + Backslashes daraus, und dann greift die MEDIA_ROOT-Prüfung nicht mehr, das + gewählte Ziel fiele still auf den Standard zurück. Aufgefallen 25.07.2026, + als der Test dafür erstmals unter Windows lief. Live war es nie: aufgerufen + wird nur aus rip_disc, und das ist auf Windows-Workern verriegelt. + """ + if target_dir: + normalisiert = posixpath.normpath(target_dir) + if unter_wurzel(normalisiert, MEDIA_ROOT): + return normalisiert + return posixpath.join(RIP_OUTPUT_DIR, disc_type) + + +def _arbeitsverzeichnis(einstellungen: dict, job_wahl: str = "") -> str: + """Basis für Roh-Rips. Reihenfolge: Wahl DIESES Rips → UI-Setting + `workDir` → Container-Default /app/temp/raw. + + Hintergrund (Commander 24.07.): Die VM-Platte (150 GB) reicht für BD-50, + aber eine 4K-UHD (bis 100 GB roh + Kompression daneben) sprengt sie — + das Arbeitsverzeichnis muss deshalb auf ein großes Ziel umlegbar sein. + + Pro Rip wählbar seit 25.07.2026 (Commander-Wunsch): beim „Rippen starten" + entscheidet man je Disc, wo die Rohdaten landen. Das Setting bleibt der + Standard — und ist damit der Wert, der bei Vollautomatik-Rips greift, bei + denen niemand gefragt wird. + """ + # posixpath statt os.path: Das sind IMMER Container-Pfade (/app/media/...), + # auch wenn ein nativer Windows-Worker dieses Modul lädt — der übersetzt + # sie erst später mit pfad_lokal(). os.path.normpath macht unter Windows + # Backslashes daraus, und dann greift die MEDIA_ROOT-Prüfung nicht mehr. + for kandidat in (job_wahl, einstellungen.get("workDir")): + wert = (kandidat or "").strip() + if wert: + normalisiert = posixpath.normpath(wert) + if unter_wurzel(normalisiert, MEDIA_ROOT): + return normalisiert + return RAW_DIR + + +def _frei_bytes(pfad: str) -> int: + """Freier Platz am Pfad (nächster existierender Elternordner zählt).""" + kandidat = pfad + while kandidat and not os.path.exists(kandidat): + kandidat = os.path.dirname(kandidat) + try: + return shutil.disk_usage(kandidat or "/").free + except OSError: + return -1 + + +def _platz_pruefen(pfad: str, benoetigt: int, zweck: str) -> str: + """Gibt eine Klartext-Fehlermeldung zurück, wenn der Platz nicht reicht — + sonst leeren String. VOR dem Rip geprüft: eine volle Platte nach 40 GB + ist der teuerste Fehlschlag, den es gibt.""" + frei = _frei_bytes(pfad) + if frei < 0 or benoetigt <= 0: + return "" # unbekannt → nicht raten, lieber rippen lassen + if frei >= benoetigt: + return "" + return ( + f"Zu wenig Platz für {zweck} in {pfad}: " + f"{frei / 1024**3:.1f} GB frei, benötigt ~{benoetigt / 1024**3:.1f} GB. " + "Abhilfe: anderes Ziel wählen oder das Arbeitsverzeichnis unter " + "Einstellungen → Verarbeitung auf eine große Freigabe legen." + ) + + +def _makemkv_key_anwenden(einstellungen: dict) -> None: + """MakeMKV-Beta-Key aus den UI-Einstellungen aktivieren (schlägt Env). + + Damit ist der Monats-Key ohne Rebuild/Neustart aktualisierbar + (Einstellungen → System). Format wie entrypoint.sh: settings.conf. + + Ergänzend statt überschreibend (Befund 25.07.2026): das Datenverzeichnis + ist jetzt persistent, und hier stand vorher ein open(..., "w") — das warf + vor JEDEM Rip alles andere aus der settings.conf, z. B. app_UpdateEnable + aus dem entrypoint. Beide Schreiber müssen gleich arbeiten, sonst kommt + der Fehler beim nächsten Rip still zurück. + """ + key = (einstellungen.get("makemkvAppKey") or "").strip() + if not key: + return + pfad = os.path.join(makemkv_daten.DATEN_DIR, "settings.conf") + try: + os.makedirs(makemkv_daten.DATEN_DIR, exist_ok=True) + try: + with open(pfad, encoding="utf-8", errors="replace") as f: + alt = f.read() + except OSError: + alt = "" + with open(pfad, "w", encoding="utf-8", newline="\n") as f: + f.write(makemkv_daten.settings_conf_zusammenfuehren(alt, key)) + except OSError as e: + db.add_log("warning", "worker", f"MakeMKV-Key konnte nicht gesetzt werden: {e}") + + +def _ordner_groesse(pfad: str) -> int: + """Belegter Platz eines Ordners in Bytes (0, wenn nicht lesbar).""" + summe = 0 + for wurzel, _, dateien in os.walk(pfad): + for name in dateien: + try: + summe += os.path.getsize(os.path.join(wurzel, name)) + except OSError: + pass + return summe + + +def _original_aufheben(job_id: str, raw_dir: str, final_dir: str) -> None: + """Roh-Rip zusätzlich aufheben — darf den Job NIEMALS scheitern lassen. + + Befund 25.07.2026 (Akira-UHD, echter Schaden): Hier stand ein nacktes + shutil.move(). Zwischen Arbeitsverzeichnis (/app/temp) und Ziel + (/app/media) scheitert os.rename mit EXDEV, shutil.move fällt auf Kopieren + zurück — eine 74-GB-Vollkopie auf dieselbe Platte, bis sie mit ENOSPC voll + war. Ergebnis: Platte 100 % voll, Worker-Container startete nicht mehr, und + der Job galt als FEHLGESCHLAGEN, obwohl die komprimierte Datei längst + fertig war. Der Nutzer sah nur „nichts da". + + ## Warum hier NICHT vorhergesagt, sondern versucht wird + + Die erste Fassung dieses Schutzes verglich `os.stat(...).st_dev` und + schloss aus gleichen Werten auf „reines Umhängen, kein Platz nötig". Am + 25.07.2026 im Worker-Container nachgemessen — beides zugleich wahr: + + st_dev /app/temp = 2050 + st_dev /app/media = 2050 → also identisch + os.rename(...) → EXDEV, „Invalid cross-device link" + + Der Kernel vergleicht bei rename() den **Mount**, nicht das Gerät. /app/temp + (Docker-Volume) und /app/media (Bind-Mount) sind zwei Mounts DERSELBEN + ext4-Partition. Die st_dev-Prüfung war deshalb wirkungslos: sie sah + „gleiches Dateisystem", übersprang die Platzprüfung, und shutil.move kopierte + doch. Der Schutz hätte genau den Schaden zugelassen, gegen den er gebaut war. + + Also: erst rename VERSUCHEN. Klappt es, ist es umgehängt und fertig. + Kommt EXDEV, steht fest, dass kopiert werden müsste — und erst dann wird + der Platz geprüft. Das ist keine Vermutung mehr, sondern die Antwort des + Kernels. + """ + ziel_original = os.path.join(final_dir, "original") + try: + # Der billige Weg zuerst — und er ist gleichzeitig der einzige + # verlässliche Test, ob überhaupt umgehängt werden kann. + try: + os.rename(raw_dir, ziel_original) + db.add_log( + "info", "worker", + f"Job {job_id}: Original behalten unter {ziel_original} (umgehängt, " + "kein zusätzlicher Platz nötig)", + ) + return + except OSError as e: + if e.errno != errno.EXDEV: + raise # etwas anderes ist schiefgelaufen → unten ehrlich melden + + # Ab hier ist eine echte Kopie unvermeidlich. Jetzt lohnt die Platzfrage. + benoetigt = _ordner_groesse(raw_dir) + frei = _frei_bytes(final_dir) + if frei < benoetigt * 1.05: + db.add_log( + "warning", "worker", + f"Job {job_id}: Original NICHT aufgehoben — Ziel liegt auf einem " + f"anderen Mount, es müsste kopiert werden. Dafür wären " + f"{benoetigt / 1024**3:.1f} GB nötig, frei sind nur " + f"{frei / 1024**3:.1f} GB. Die Roh-Datei bleibt unter " + f"{raw_dir} liegen. Abhilfe: Arbeitsverzeichnis " + "(Einstellungen → Verarbeitung) auf dieselbe Freigabe legen " + "wie das Ziel — dann wird nur umgehängt statt kopiert." + ) + return + shutil.move(raw_dir, ziel_original) + db.add_log("info", "worker", f"Job {job_id}: Original behalten unter {ziel_original}") + except OSError as e: + # Halb geschriebene Kopie wegräumen, sonst belegt sie für immer Platz. + shutil.rmtree(ziel_original, ignore_errors=True) + db.add_log( + "warning", "worker", + f"Job {job_id}: Original konnte nicht aufgehoben werden ({e}). " + f"Die komprimierte Datei ist fertig; die Roh-Datei bleibt unter {raw_dir}." + ) + + +def _abbruch_angefordert(job_id: str) -> bool: + """Kooperativer Abbruch: hat der Nutzer über die API abgebrochen?""" + return db.get_job_status(job_id) == "canceling" + + +def _serien_episoden_zuordnen(ausgabe: str, serie: str, staffel, meta: dict) -> list: + """Episoden-Zuordnung per Laufzeitabgleich (ARM-Wunde #395, besser gelöst). + + Laufzeiten der Staffel kommen von der API (TMDB); die der Dateien vom + HandBrake-Scan. Nur bei EINDEUTIGEM Treffer wird umbenannt. Wirft nie. + """ + if meta.get("source") != "tmdb" or meta.get("type") != "tv" or not meta.get("id"): + return ["Episoden-Zuordnung übersprungen (nur mit TMDB-Serien-Treffer möglich)"] + try: + antwort = requests.get( + f"{API_URL}/metadata/tv/{meta['id']}/season/{int(staffel)}", timeout=30 + ) + if antwort.status_code != 200: + return [f"Episoden-Laufzeiten nicht abrufbar (HTTP {antwort.status_code})"] + episoden = [ + (e["episode"], (e.get("runtime") or 0) * 60) + for e in antwort.json().get("episodes", []) + ] + except (requests.RequestException, ValueError) as e: + return [f"Episoden-Laufzeiten nicht abrufbar: {e}"] + if not episoden: + return ["Keine Episoden-Laufzeiten bei TMDB hinterlegt — Dateinamen bleiben"] + + dateien = sorted( + os.path.join(ausgabe, f) for f in os.listdir(ausgabe) if f.endswith(".mkv") + ) + if not dateien: + return [] + dauern = [lies_datei_dauer(p) for p in dateien] + if 0 in dauern: + return ["Datei-Laufzeiten nicht lesbar — Dateinamen bleiben"] + + zuordnung = medien.matche_episoden(dauern, episoden) + if not zuordnung: + return [ + "Keine EINDEUTIGE Episoden-Zuordnung über die Laufzeiten — " + "Dateinamen bleiben (lieber ehrlich als falsch benannt)" + ] + return medien.episoden_umbenennen(ausgabe, serie, staffel, zuordnung) + + +def _benachrichtigen(job_id: str, betreff: str, text: str, level: str) -> None: + """Webhook-Meldung bei Job-Ende — best effort, nie job-entscheidend.""" + einstellungen = db.get_settings(bei_fehler_leer=True) + url = (einstellungen.get("notificationWebhook") or "").strip() + if not url: + return + try: + notify.sende(url, betreff, text, level) + except Exception as e: + db.add_log("warning", "notify", f"Job {job_id}: Benachrichtigung fehlgeschlagen: {e}") + + +def _job_abschliessen(job_id: str, ergebnis: dict) -> None: + """Schreibt den Endzustand eines Jobs (completed/failed) nach Postgres, + bereitet bei Erfolg fürs Medien-System auf (NFO/Poster) und meldet das + Ende per Webhook (falls konfiguriert).""" + job = db.get_job(job_id) or {} + titel = job.get("title") or job_id[:8] + + if ergebnis.get("status") == "cancelled": + db.update_job( + job_id, status="failed", error="Abgebrochen durch Nutzer", + finished_at=db.utcnow(), + ) + db.add_log("warning", "worker", f"Job {job_id}: abgebrochen — Rohdaten bleiben erhalten") + _benachrichtigen(job_id, f'⚠️ Rippy: „{titel}" abgebrochen', + "Der Job wurde auf Wunsch abgebrochen — Rohdaten bleiben erhalten.", + "warning") + return + if ergebnis.get("status") == "success": + ausgabe = ergebnis.get("output_dir") + if ausgabe and (job.get("disc_type") in ("dvd", "bluray", "uhd")): + einstellungen = db.get_settings(bei_fehler_leer=True) + try: + meta = json.loads(job.get("meta") or "{}") + except ValueError: + meta = {} + serie = meta.get("series") + staffel = meta.get("season") + + # Serien-Rips: Episoden per Laufzeit zuordnen und umbenennen + # (Show S01E02.mkv) — nur bei EINDEUTIGER Zuordnung, sonst + # bleiben die MakeMKV-Namen (ehrlich geloggt). + if serie and staffel: + for meldung in _serien_episoden_zuordnen(ausgabe, serie, staffel, meta): + db.add_log("info", "worker", f"Job {job_id}: {meldung}") + + # Media-Server-Aufbereitung (Jellyfin/Emby/Kodi: NFO + Poster); + # bei Serien wandern tvshow.nfo/poster in den Serien-Ordner. + for meldung in medien.aufbereiten( + ausgabe, einstellungen.get("mediaServer") or "none", + serie or job.get("title") or "", meta.get("year"), meta, + serien_root=os.path.dirname(ausgabe) if serie and staffel else None, + ): + db.add_log("info", "worker", f"Job {job_id}: {meldung}") + + # Bibliotheks-Refresh (Jellyfin/Emby): der Server scannt sofort + meldung = medien.bibliothek_refresh( + einstellungen.get("mediaServer") or "none", + (einstellungen.get("jellyfinUrl") or "").strip(), + (einstellungen.get("jellyfinApiKey") or "").strip(), + ) + if meldung: + db.add_log("info", "worker", f"Job {job_id}: {meldung}") + db.update_job( + job_id, + status="completed", + progress=100, + output_path=ausgabe, + finished_at=db.utcnow(), + ) + db.add_log("success", "worker", f"Job {job_id}: abgeschlossen → {ausgabe}") + _benachrichtigen(job_id, f'✅ Rippy: „{titel}" ist fertig', + f"Abgelegt unter: {ausgabe}", "success") + else: + fehler = ergebnis.get("error", "unbekannter Fehler") + db.update_job( + job_id, + status="failed", + error=fehler, + finished_at=db.utcnow(), + ) + db.add_log("error", "worker", f"Job {job_id}: {fehler}") + _benachrichtigen(job_id, f'❌ Rippy: „{titel}" fehlgeschlagen', fehler, "error") + + +def rippen(device_path: str, job_id: str, target_dir: str = None, + melde=None, weiterreichen=None): + """Stufe 1: Rippt die Disc; bei Video folgt die Kompression als eigener Task. + + target_dir (optional): vom Nutzer gewähltes Ablageziel unter /app/media — + dort eingehängte Shares (NFS/SMB) sind damit direkt wählbar. + """ + db.init_db() + + if detect_disc_type is None: + ergebnis = { + "status": "error", + "error": "Dieser Worker kann nur komprimieren (nativer " + "Transcode-Worker ohne Laufwerks-Zugriff) — gerippt " + "wird auf der Rippy-Hauptmaschine.", + } + _job_abschliessen(job_id, ergebnis) + return ergebnis + + if _abbruch_angefordert(job_id): + _job_abschliessen(job_id, {"status": "cancelled"}) + return {"status": "cancelled"} + + disc_type = detect_disc_type(device_path) + + if disc_type in ("no_disc", "unknown"): + fehler = ( + "Keine Disc im Laufwerk" if disc_type == "no_disc" + else "Disc-Typ nicht erkennbar" + ) + db.update_job(job_id, status="failed", error=fehler, finished_at=db.utcnow()) + db.add_log("error", "worker", f"Job {job_id}: {fehler} ({device_path})") + return {"status": "error", "error": fehler, "disc_type": disc_type} + + db.update_job(job_id, status="running", disc_type=disc_type) + # Ab hier gilt: der Rip läuft, ist aber NICHT fertig. Stirbt der Job jetzt, + # ist jede Roh-Datei ein Bruchstück (siehe RIP_FERTIG oben). + db.meta_merken(job_id, **{RIP_FERTIG: False}) + db.add_log("info", "worker", f"Job {job_id}: {disc_type}-Rip gestartet ({device_path})") + + letzter = [-1] + + def fortschritt(progress: int, message: str = ""): + # MakeMKV liefert viele PRGV-Zeilen pro Sekunde — DB nur bei Änderung. + if progress == letzter[0]: + return + letzter[0] = progress + if _abbruch_angefordert(job_id): + raise RipAbbruch() + # Frueher: self.update_state(...) — der Zwischenstand ging an Celery. + # Jetzt ein Rueckruf, den der Aufrufer stellt: im verteilten Betrieb + # meldet er an Celery, im Standalone-Betrieb an den Ereignis-Bus. + # Der ABLAUF muss nicht wissen, wer zuhoert. + if melde: + melde({"progress": progress, "status": "ripping", "message": message}) + db.update_job(job_id, progress=progress) + + # MakeMKV-Meldungen ins Log (Befund 25.07.2026): Bis dahin überlebte NUR + # die letzte Zeile ("Failed to open disc"), und die sagt nichts. Dass + # MakeMKV bei der UHD-Disc nicht einmal versucht, einen Schlüssel zu + # holen, war deshalb nur per Hand-Lauf im Container zu sehen. + # Gedrosselt, weil das UI global nur die letzten 200 Zeilen zeigt: jede + # Meldung höchstens einmal, insgesamt höchstens MAX_MELDUNGEN je Rip. + # Code 1003 ist MakeMKVs eigenes DEBUG-Rauschen (am 25.07. beobachtet). + MAX_MELDUNGEN = 40 + gesehen = set() + + def melde_makemkv(code: int, text: str): + if code == 1003 or len(gesehen) >= MAX_MELDUNGEN or text in gesehen: + return + gesehen.add(text) + db.add_log("info", "makemkv", f"Job {job_id}: {text[:300]}") + + einstellungen = db.get_settings(bei_fehler_leer=True) + ist_video = disc_type in ("dvd", "bluray", "uhd") + # Je Disc-Typ abwählbar (siehe komprimieren_fuer): 4K verlustfrei behalten, + # DVDs trotzdem schrumpfen — vorher gab es nur alles oder nichts. + transcode_an = ist_video and komprimieren_fuer(disc_type, einstellungen) + + if ist_video: + # UI-Key schlägt Env-Key — Monats-Key ohne Rebuild aktualisierbar + _makemkv_key_anwenden(einstellungen) + + # Sprechender Zielordner „ (Jahr)" statt Job-UUID — Jellyfin & Co. + # erkennen den Film am Ordnernamen. UUID bleibt Fallback ohne Titel. + # Serien-Rips landen stattdessen in /Season NN (Staffel-Flow). + job = db.get_job(job_id) or {} + try: + meta = json.loads(job.get("meta") or "{}") + except ValueError: + meta = {} + if meta.get("series") and meta.get("season"): + final_dir = medien.serien_ordner( + _zielbasis(target_dir, disc_type), meta["series"], meta["season"] + ) + else: + final_dir = medien.zielordner( + _zielbasis(target_dir, disc_type), + job.get("title") or "", meta.get("year"), job_id, + ) + + # Titel-Wahl: explizite Auswahl aus der Track-Tabelle schlägt alles; + # sonst „Nur Hauptfilm" (pro Rip wählbar, Fallback: Setting); + # bei Serien-Discs ohne Auswahl zählen ALLE Episoden-Titel. + titel_liste = [int(t) for t in (meta.get("titles") or []) if str(t).isdigit()] + nur_hauptfilm = meta.get("main_feature_only") + if nur_hauptfilm is None: + nur_hauptfilm = bool(einstellungen.get("mainFeatureOnly", False)) + if meta.get("series") or titel_liste: + nur_hauptfilm = False + # Geplantes Ziel sofort sichtbar machen (UI-Detail + retry-transcode) + db.update_job(job_id, output_path=final_dir) + + # meta["work_dir"] = die Wahl aus dem Rip-Dialog; leer = Setting/Default. + raw_dir = os.path.join(_arbeitsverzeichnis(einstellungen, meta.get("work_dir")), job_id) + + # Platz-Check VOR dem Rip: Disc-Größe ist per ioctl bekannt — eine volle + # Platte nach 40 GB wäre der teuerste Fehlschlag (4K-UHD: bis 100 GB roh). + try: + disc_bytes = disc_size_bytes(device_path) if disc_size_bytes else 0 + except OSError: + disc_bytes = 0 + marge = 1024**3 # 1 GB Sicherheitsabstand + if ist_video: + rip_ziel = raw_dir if transcode_an else final_dir + fehler = _platz_pruefen(rip_ziel, disc_bytes + marge, "den Roh-Rip") + if not fehler and transcode_an and einstellungen.get("keepOriginal", False): + fehler = _platz_pruefen(final_dir, disc_bytes + marge, "das Original (keepOriginal)") + if fehler: + ergebnis = {"status": "error", "error": fehler} + _job_abschliessen(job_id, ergebnis) + return ergebnis + + if disc_type == "cd": + ergebnis = rip_cd( + device_path, job_id, progress_cb=fortschritt, output_dir=final_dir, + auswerfen=bool(einstellungen.get("autoEject", True)), + ) + elif transcode_an: + # Roh-Rip ins Arbeitsverzeichnis (wird nach erfolgreicher Kompression gelöscht) + ergebnis = rip_video( + device_path, job_id, disc_type, + progress_cb=fortschritt, + output_dir=raw_dir, + nur_hauptfilm=nur_hauptfilm, + titel_liste=titel_liste, + log_cb=melde_makemkv, + ) + else: + ergebnis = rip_video( + device_path, job_id, disc_type, + progress_cb=fortschritt, output_dir=final_dir, + nur_hauptfilm=nur_hauptfilm, + titel_liste=titel_liste, + log_cb=melde_makemkv, + ) + + # MakeMKV-Fehler in Klartext übersetzen — "Failed to open disc" allein + # hilft niemandem. + if ergebnis.get("status") == "error": + fehler_text = ergebnis.get("error") or "" + if "volume key is unknown" in fehler_text: + # Befund 25.07.2026, auf BEIDEN Maschinen gemessen: Laufwerk und + # MakeMKV sind in Ordnung. makemkvcon unter Linux ruft die + # Disc-Schlüssel schlicht nie ab — die Windows-Version tut es + # (Meldung 3338). Hier stand vorher erst "Disc zu neu" und danach + # "der Schlüssel-Kanal ist tot"; beides war falsch und hat in die + # Irre geschickt. Herleitung im Kopf von makemkv_daten.py. + speicher = makemkv_daten.schluesselspeicher_status() + anzahl = speicher.get("schluessel", 0) + ergebnis["error"] += ( + " — Klartext: Laufwerk und Rippy sind in Ordnung, MakeMKV " + "liest die Disc. Es fehlt nur der Schlüssel dieser Pressung. " + "Der Grund: makemkvcon holt Schlüssel unter Linux nie selbst " + "nach — die Windows-Version schon. " + + ( + "Dieser Worker kennt aktuell GAR KEINEN Disc-Schlüssel. " + if not anzahl + else f"Dieser Worker kennt {anzahl} Disc-Schlüssel, " + "diese Pressung ist nicht dabei. " + ) + + "Abhilfe: MakeMKV auf einem Windows-PC installieren, die " + "Disc dort einmal öffnen, dann die Datei _private_data.tar " + "aus dem MakeMKV-Datenverzeichnis unter Einstellungen → " + "System hochladen. Wirkt ab dem nächsten Rip. Klappt auch " + "das nicht, kennt MakeMKV die Pressung selbst nicht — dann " + "hilft nur eine KEYDB.cfg (ebenfalls dort hochladbar) oder " + "das Einreichen des AACS-Dumps im MakeMKV-Forum, Bereich " + "'Ultra HD Blu-ray'. Der Dump steht unter Einstellungen → " + "System zum Download bereit." + ) + elif disc_type == "uhd" and "Failed to open disc" in fehler_text: + ergebnis["error"] += ( + " — 4K-UHD erkannt: Das Laufwerk kann UHD-Discs vermutlich " + "nicht entschlüsseln. Dafür ist eine LibreDrive-Firmware nötig " + "(MakeMKV-Forum 'Ultimate UHD Drives Flashing Guide'). " + "Normale BD/DVD gehen weiterhin." + ) + + # Automatischer Auswurf. Die Disc ist nach dem Rip nicht mehr nötig — die + # Kompression arbeitet auf der Datei, nicht am Laufwerk. + # + # Befund 25.07.2026: Die Einstellung („Disc nach erfolgreichem Ripping + # automatisch auswerfen", Standard ein) wurde von NIEMANDEM gelesen. Bei + # DVD/Blu-ray warf Rippy deshalb nie aus, bei Audio-CD dagegen immer, weil + # abcde `-x` fest verdrahtet bekam. Jetzt entscheidet die Einstellung beides. + if ergebnis.get("status") == "success" and ist_video and einstellungen.get("autoEject", True): + if wirf_disc_aus(device_path): + db.add_log("info", "worker", f"Job {job_id}: Disc ausgeworfen") + else: + db.add_log( + "warning", "worker", + f"Job {job_id}: Disc konnte nicht ausgeworfen werden ({device_path}) — " + "der Rip ist davon unberührt.", + ) + + if ergebnis.get("status") == "success" and transcode_an: + # 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). + db.update_job(job_id, status="transcoding", progress=0) + # Der Rip ist durch. Ab jetzt ist ein Fehlschlag mit „Neu komprimieren" + # zu heilen, ohne die Disc noch einmal zu lesen. + db.meta_merken(job_id, **{RIP_FERTIG: True}) + db.add_log("info", "worker", f"Job {job_id}: Rip fertig, Kompression eingereiht") + # Wer die Kompression uebernimmt, entscheidet der Aufrufer: + # Celery-Queue im verteilten Betrieb, LocalQueue im Standalone-Betrieb. + # OHNE Rueckruf wird gleich hier komprimiert — das ist der Fall + # "ein Prozess macht alles", also der Windows-Betrieb. + if weiterreichen: + weiterreichen(job_id, raw_dir, final_dir, meta.get("transcode_node")) + return {"status": "ripped", "raw_dir": raw_dir} + return komprimieren(job_id, raw_dir, final_dir) + + _job_abschliessen(job_id, ergebnis) + return ergebnis + + +def ping_worker(): + """Winziger Task zum Beweisen des gezielten Routings: gibt Hostname + + Encoder zurück. Wird über worker_direct an EINEN Worker geschickt — landet + er beim richtigen, stimmt das Encoder-Routing (Encoder-Auswahl im UI).""" + import socket + + import caps + return {"hostname": socket.gethostname(), "encoders": caps.erkenne_encoder()} + + +def scan_tracks(device_path: str): + """Titel-Liste der eingelegten Disc erfassen (für die Auswahl-Tabelle). + + Ergebnis landet in der settings-Tabelle (key 'tracks:') — die + API pollt darauf. Läuft auf der Haupt-Queue (nur der Rip-Worker hat das + Laufwerk). Ein Info-Lauf dauert je nach Disc 20–120 s. + """ + db.init_db() + key = f"tracks:{device_path}" + db.save_settings({"status": "running"}, key) + try: + # Ein Info-Lauf liefert BEIDES: Titel-Liste und die Sprachen. Die + # Sprach-Auskunft stand schon immer in derselben Ausgabe und wurde nur + # weggeworfen (Commander-Anforderung 26.07.2026: vor dem Rip fragen, + # was man haben will). + titel, streams = lies_titel_info(device_path) + sprachen = sprachen_zusammenfassen(streams) + db.save_settings( + {"status": "done", "tracks": titel, "sprachen": sprachen}, key) + db.add_log( + "info", "worker", + f"Titel-Scan {device_path}: {len(titel)} Titel, " + f"{len(sprachen['audio'])} Tonsprache(n), " + f"{len(sprachen['untertitel'])} Untertitelsprache(n)", + ) + return {"status": "done", "anzahl": len(titel)} + except Exception as e: + db.save_settings({"status": "error", "error": str(e)[:300]}, key) + db.add_log("error", "worker", f"Titel-Scan {device_path} fehlgeschlagen: {e}") + return {"status": "error", "error": str(e)} + + +def _ist_uebersetzt(container_pfad: str, lokal_pfad: str) -> bool: + """Hat RIPPY_PATH_MAP diesen Pfad wirklich angefasst? (pure, testbar)""" + return container_pfad != lokal_pfad + + +# Wie weit wird nach einem vorhandenen Elternordner gesucht. Begrenzt, damit auf +# einer toten Freigabe nicht endlos geklopft wird. +MAX_ELTERN_STUFEN = 12 + + +def erster_vorhandener_ordner(pfad: str, isdir=None, dirname=None) -> str: + """Der nächste EXISTIERENDE Ordner oberhalb von `pfad` — oder "". + + ⚠️ Vorfall 26.07.2026, im vollen Durchlauf aufgefallen: Der Rip lief sauber + durch (40 GB), die Kompression brach sofort ab mit „Dieser Worker erreicht + das Ziel nicht: \\\\NAS\\rippy\\movies\\Akira (1988)". Erreichbar war die + Freigabe sehr wohl — es fehlten schlicht ZWEI Ordner: `movies` und der + Filmordner darin. Die Prüfung ließ nur EINE fehlende Ebene durchgehen (sie + sah nach `dirname`), obwohl gleich darauf `os.makedirs` die ganze Kette + anlegt. Ein Ziel, das noch nie beschrieben wurde, war damit systematisch + unerreichbar — also jeder erste Film in einer neuen Ablage. + + Es zählt daher, ob IRGENDEIN Vorfahre existiert: Von dem aus kann + `makedirs` den Rest bauen. Existiert keiner, ist die Freigabe wirklich weg. + + `isdir`/`dirname` werden BEIM AUFRUF aufgelöst, nicht als Vorgabewert + gebunden — sonst zeigt die Vorgabe für immer auf die Funktion von damals, + und ein Ersetzen von `os.path.isdir` (im Test wie im Betrieb) ginge ins + Leere. Genau darüber bin ich beim Schreiben der Tests gestolpert. + """ + isdir = isdir or os.path.isdir + dirname = dirname or os.path.dirname + aktuell = dirname(pfad or "") + for _ in range(MAX_ELTERN_STUFEN): + if not aktuell: + return "" + if isdir(aktuell): + return aktuell + naechster = dirname(aktuell) + if naechster == aktuell: # Wurzel erreicht + return "" + aktuell = naechster + return "" + + +def _erreichbarkeit_pruefen(raw_container: str, raw_lokal: str, + final_container: str, final_lokal: str) -> str: + """Kann DIESER Worker Quelle und Ziel überhaupt sehen? Klartext oder "". + + Befund 25.07.2026 am Job 95afdc89, live: Das gezielte Routing an den + Windows-PC des Commanders funktionierte einwandfrei — der Worker nahm die + Aufgabe an und lehnte sie 182 ms später ab mit „Keine Roh-MKVs in + /app/media/rippy/… gefunden". Diese Meldung klang nach einem kaputten Rip, + obwohl der Rip vollständig war (79,6 GB lagen auf der NAS). Der Commander + schloss daraus, das System sähe den externen Encoder nicht — verständlich + und falsch. + + Die eigentliche Ursache: `/app/media/...` sind CONTAINER-Pfade. Ein nativer + Worker sieht sie nur, wenn RIPPY_PATH_MAP sie auf eine Freigabe übersetzt — + und dieses Mapping wurde von NIEMANDEM gesetzt. Es war also nicht möglich, + dass diese Kombination je funktioniert. + + Zweiter Stolperstein im selben Fall: Die Rohdaten lagen auf der NAS (für den + PC erreichbar), das ZIEL aber auf der VM-Platte (nicht erreichbar, dort + läuft kein Samba). Deshalb werden hier BEIDE Pfade geprüft, nicht nur die + Quelle — sonst scheitert es erst beim Schreiben, nach Stunden Rechenzeit. + """ + mapping = os.getenv("RIPPY_PATH_MAP", "") + fremder_worker = not os.path.isdir("/app") + + for zweck, container, lokal in ( + ("die Quelle (Rohdaten)", raw_container, raw_lokal), + ("das Ziel (fertige Datei)", final_container, final_lokal), + ): + if os.path.isdir(lokal): + continue + # Das ZIEL darf fehlen — es wird gleich mit `os.makedirs` angelegt, und + # zwar samt aller fehlenden Zwischenebenen. Es genügt also, dass + # irgendein Vorfahre existiert (Herleitung in erster_vorhandener_ordner). + if zweck.startswith("das Ziel") and erster_vorhandener_ordner(lokal): + continue + + text = [f"Dieser Worker erreicht {zweck} nicht: {lokal}"] + if fremder_worker and not mapping: + text.append( + "Ursache: RIPPY_PATH_MAP ist auf dieser Maschine nicht gesetzt. " + f'„{container}" ist ein Pfad INNERHALB des Rippy-Containers — ein ' + 'externer Worker sieht ihn nur, wenn er auf eine Netzwerk-Freigabe ' + 'übersetzt wird.' + ) + text.append( + "Abhilfe: Freigabe auf diesem Rechner erreichbar machen und in " + "start-tray.bat setzen, z. B. " + "set RIPPY_PATH_MAP=/app/media=\\\\NAS\\rippy-media" + ) + elif fremder_worker and mapping and not _ist_uebersetzt(container, lokal): + text.append( + f"RIPPY_PATH_MAP ist gesetzt ({mapping}), deckt diesen Pfad aber " + f'nicht ab: „{container}" wurde von keinem Eintrag übersetzt — ' + 'fehlt ein Präfix, oder ist es ein anderes Verzeichnis als erwartet?' + ) + elif fremder_worker and mapping: + # ⚠️ Hier stand bis zum 26.07.2026 dieselbe Meldung wie oben — „deckt + # diesen Pfad nicht ab" —, obwohl die Karte ihn nachweislich übersetzt + # hatte (der UNC-Pfad stand im selben Satz). Wer der Meldung folgte, + # suchte den Fehler in der Karte, während die Freigabe schlicht nicht + # verbunden war. Eine falsche Ursache ist teurer als gar keine. + text.append( + f'Übersetzt wurde er korrekt ({container} → {lokal}, Karte: ' + f'{mapping}) — dieser Pfad ist auf diesem Rechner nur gerade ' + 'nicht erreichbar. Ist die Freigabe verbunden? Anmeldedaten noch ' + 'gültig? Von Hand prüfen: den Pfad im Explorer öffnen.' + ) + else: + text.append( + "Der Pfad existiert nicht. Liegt das Arbeitsverzeichnis auf einer " + "Freigabe, die gerade nicht eingehängt ist?" + ) + text.append( + 'WICHTIG: Die Rohdaten sind NICHT verloren — nach der Korrektur ' + 'genügt „Neu komprimieren", ohne die Disc erneut zu rippen.' + ) + return " ".join(text) + return "" + + +def komprimieren(job_id: str, raw_dir: str, final_dir: str): + """Stufe 2: HandBrake komprimiert die Roh-MKVs auf Arbeitsgröße. + + Erst wenn ALLE Dateien sauber komprimiert sind, wird das Roh-Verzeichnis + gelöscht — bricht die Kompression ab, bleibt das Original in /app/temp + liegen (kein Datenverlust wie bei ARMs berüchtigtem Move-Bug #1530). + Über POST /jobs/{id}/retry-transcode jederzeit neu anstoßbar. + """ + db.init_db() + # Native Worker (Windows) übersetzen Container-Pfade aufs Netzlaufwerk + raw_container, final_container = raw_dir, final_dir + raw_dir = pfad_lokal(raw_dir) + final_dir = pfad_lokal(final_dir) + + fehler = _erreichbarkeit_pruefen(raw_container, raw_dir, final_container, final_dir) + if fehler: + ergebnis = {"status": "error", "error": fehler} + _job_abschliessen(job_id, ergebnis) + return ergebnis + + quellen = sorted(glob.glob(os.path.join(raw_dir, "*.mkv"))) + if not quellen: + ergebnis = { + "status": "error", + "error": ( + f"Das Verzeichnis {raw_dir} ist erreichbar, enthält aber keine " + "MKV-Datei. Der Rip hat also nichts abgelegt (oder jemand hat die " + "Datei entfernt). Die Kompression lässt sich nach einem neuen Rip " + "erneut anstoßen." + ), + } + _job_abschliessen(job_id, ergebnis) + return ergebnis + + einstellungen = db.get_settings(bei_fehler_leer=True) + # Preset nach Disc-Typ (Befund 25.07.2026): vorher lief JEDE Quelle durch + # dasselbe Preset — eine 4K-UHD wurde damit auf 1080p heruntergerechnet. + job = db.get_job(job_id) or {} + disc_type = job.get("disc_type") or "" + preset = preset_fuer(disc_type, einstellungen) + original_behalten = einstellungen.get("keepOriginal", False) + + # Sprachauswahl: Wahl für DIESEN Rip (aus meta) schlägt die Einstellung. + # Leer heißt „alles behalten" — genau das Verhalten von vorher. + try: + job_meta = json.loads(job.get("meta") or "{}") + except ValueError: + job_meta = {} + audio_sprachen = sprachliste( + job_meta.get("audio_sprachen") or einstellungen.get("audioSprachen")) + untertitel_sprachen = sprachliste( + job_meta.get("untertitel_sprachen") or einstellungen.get("untertitelSprachen")) + + os.makedirs(final_dir, exist_ok=True) + db.update_job(job_id, status="transcoding", progress=0, error=None) + # Wer hier ankommt, hat vollständige Quelldateien — sonst wäre oben schon + # abgebrochen worden. Das vermerkt die Marke auch für BESTANDSJOBS, die vor + # ihrer Einführung gerippt wurden und sie darum noch nicht tragen. + db.meta_merken(job_id, **{RIP_FERTIG: True}) + db.add_log( + "info", "worker", + f"Job {job_id}: Kompression gestartet ({len(quellen)} Datei(en), " + f"Disc-Typ '{disc_type or 'unbekannt'}', Preset '{preset}'" + + (f", Ton: {','.join(audio_sprachen)}" if audio_sprachen else ", Ton: alle") + + (f", Untertitel: {','.join(untertitel_sprachen)}" + if untertitel_sprachen else ", Untertitel: alle") + + ")", + ) + + anzahl = len(quellen) + letzter = [-1] + letzte_abbruchpruefung = [0.0] + + def abbruch_pruefen(): + """Zeitgesteuert prüfen, ob der Nutzer abgebrochen hat. Wirft RipAbbruch. + + Befund 25.07.2026 (Commander, am laufenden Akira-Job beobachtet): Der + Abbruch wurde NUR in `datei_fortschritt` geprüft — und diese Closure + stieg oben sofort wieder aus, wenn sich die Prozentzahl nicht geändert + hatte. Bei einem 4K-Encode, der pro Prozent eine halbe Stunde braucht, + sah „Abbrechen" entsprechend lange wirkungslos aus (gemessen: 3,4 min + zwischen Anforderung 18:30:17 und Bestätigung 18:33:38 — bei noch + langsamerem Fortschritt entsprechend mehr). + + Deshalb hängt die Prüfung jetzt an der Zeit statt am Fortschritt und + läuft bei JEDER Ausgabezeile von HandBrake — auch während des + Scan-Durchlaufs, der überhaupt keine Encode-Prozente liefert. + """ + jetzt = time.monotonic() + if jetzt - letzte_abbruchpruefung[0] < ABBRUCH_INTERVALL_SEKUNDEN: + return + letzte_abbruchpruefung[0] = jetzt + if _abbruch_angefordert(job_id): + raise RipAbbruch() + + for index, quelle in enumerate(quellen): + ziel = os.path.join(final_dir, os.path.basename(quelle)) + + def datei_fortschritt(p, _index=index): + gesamt = int((_index * 100 + p) / anzahl) + if gesamt == letzter[0]: + return + letzter[0] = gesamt + db.update_job(job_id, progress=min(99, gesamt)) + + hb = run_handbrake( + quelle, ziel, preset=preset, + progress_cb=datei_fortschritt, + abbruch_cb=abbruch_pruefen, + audio_sprachen=audio_sprachen, + untertitel_sprachen=untertitel_sprachen, + ) + if hb.get("status") == "cancelled": + _job_abschliessen(job_id, hb) + return hb + if hb.get("status") != "success": + ergebnis = { + "status": "error", + "error": ( + f"Kompression fehlgeschlagen bei {os.path.basename(quelle)}: " + f"{hb.get('error')} — Roh-Datei bleibt in /app/temp erhalten" + ), + } + _job_abschliessen(job_id, ergebnis) + return ergebnis + + if original_behalten: + _original_aufheben(job_id, raw_dir, final_dir) + else: + shutil.rmtree(raw_dir, ignore_errors=True) + + ergebnis = {"status": "success", "output_dir": final_dir} + _job_abschliessen(job_id, ergebnis) + return ergebnis diff --git a/docker/worker/tasks.py b/docker/worker/tasks.py index c631cfd..ff661b4 100644 --- a/docker/worker/tasks.py +++ b/docker/worker/tasks.py @@ -1,76 +1,46 @@ -"""Rip- und Transcode-Tasks — bewusst GETRENNT (23.07.2026, Basis für Etappe 12/20). +"""Celery-Adapter: dieselben Task-Namen wie bisher, der Ablauf liegt in ablauf.py. -Warum zwei Tasks: Rippen braucht das Laufwerk (läuft immer lokal), Kompression -braucht nur CPU/GPU + Zugriff auf die Rohdatei. Als eigener Celery-Task auf der -Queue "transcode" kann die Kompression damit auch ein Remote-Worker mit GPU -übernehmen (optionales Add-on) — und fehlgeschlagene Kompressionen lassen sich -neu anstoßen, ohne die Disc neu zu rippen (POST /jobs/{id}/retry-transcode). +## Was sich geaendert hat (V2-4, 28.08.2026) — und was NICHT -Zwei Stufen (Commander-Entscheid 23.07.): MakeMKV rippt verlustfrei (einziger -Weg durch AACS — HandBrake kann verschlüsselte Discs nicht lesen), HandBrake -komprimiert danach auf Arbeitsgröße. Die Rohdatei liegt nur temporär in -/app/temp und wird nach Erfolg gelöscht (Setting keepOriginal behält sie). +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 errno -import glob -import json import os -import posixpath -import shutil -import time -import requests - -from rippy import store as db -from rippy.rip import makemkv_daten -import medien -from rippy.core import notify +import ablauf from celery_app import celery_app -# fcntl gibt es nur unter Linux — der NATIVE Windows-Transcode-Worker -# (deploy/worker-windows) lädt dieses Modul auch, bedient aber nur die -# transcode-Queue. Rippen ohne detection ist unten hart verriegelt. -# -# ⚠️ Dieses except verschluckt JEDEN ImportError, nicht nur den fehlenden -# fcntl — auch einen falschen Modulpfad. Beim Umzug nach rippy.drives -# (V2-0, 28.08.2026) wäre genau das passiert: auf Linux hätte der Worker -# still detect_disc_type=None gesetzt und jeden Rip verweigert, ohne dass -# irgendwo ein Fehler stünde. Wächter dagegen: src/rippy/test_paket.py -# prüft plattformunabhängig, dass es den Modulpfad wirklich gibt. -try: - from rippy.drives.detection import detect_disc_type, disc_size_bytes -except ImportError: # Windows: kein fcntl - detect_disc_type = None - disc_size_bytes = None - -from ripping import ( - RIP_OUTPUT_DIR, - RipAbbruch, - komprimieren_fuer, - lies_datei_dauer, - lies_titel_info, - preset_fuer, - rip_cd, - rip_video, - run_handbrake, - sprachen_zusammenfassen, - sprachliste, - wirf_disc_aus, -) - -API_URL = os.getenv("API_URL", "http://api:8000") - -# Phasen-Marke in den Job-Metadaten: war der RIP fertig, als es schiefging? -# -# Sobald `status = "failed"` in der Zeile steht, ist die Phase sonst -# unwiederbringlich fort — und genau die entscheidet, was danach hilft: Nach -# einem toten Transcode genügt „Neu komprimieren"; nach einem toten Rip liegt -# nur ein Bruchstück da (Vorfall 26.07.2026: 5,1 GB von rund 40 GB) und es muss -# neu gerippt werden. Gelesen wird die Marke in api/phasen.py — der Name steht -# in beiden Dateien und wird von test_phasen.py gegeneinander geprüft, weil es -# kein geteiltes Paket zwischen den Containern gibt. -RIP_FERTIG = "rip_fertig" +# 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): @@ -89,912 +59,43 @@ def _transcode_queue(node: str): return "transcode" -def pfad_lokal(pfad: str, mapping: str = None) -> str: - """Übersetzt Rippy-Container-Pfade für native Worker (pure Funktion). - - Der native Windows-Worker sieht /app/temp und /app/media nicht — er - mountet Rippys Freigabe als Netzlaufwerk und setzt RIPPY_PATH_MAP, - z. B. "/app/media=Z:\\media;/app/temp=Y:\\temp". Ohne Mapping (Docker- - Worker: identische Mounts) kommt der Pfad unverändert zurück. - """ - mapping = mapping if mapping is not None else os.getenv("RIPPY_PATH_MAP", "") - if not pfad or not mapping: - return pfad - for paar in mapping.split(";"): - if "=" not in paar: - continue - quelle, ziel = paar.split("=", 1) - if quelle and pfad.startswith(quelle): - rest = pfad[len(quelle):] - if "\\" in ziel: - rest = rest.replace("/", "\\") - return ziel + rest - return pfad - -RAW_DIR = os.getenv("RAW_DIR", "/app/temp/raw") -MEDIA_ROOT = "/app/media" - -# Wie oft während einer Kompression nachgesehen wird, ob der Nutzer abgebrochen -# hat. Eine DB-Abfrage alle paar Sekunden ist nichts gegen einen Encode, der -# Stunden läuft — und „Abbrechen" fühlt sich damit sofort an. -ABBRUCH_INTERVALL_SEKUNDEN = 5 - - -def unter_wurzel(pfad: str, wurzel: str) -> bool: - """Liegt `pfad` wirklich unterhalb von `wurzel` (oder IST es die Wurzel)? - - Pure Funktion, testbar. Ein nacktes `startswith()` genügt nicht: - „/app/media-boese/x" beginnt mit „/app/media", liegt aber außerhalb - (Befund 25.07.2026 bei der Durchsicht). Immer „/" als Trenner — das sind - Container-Pfade, auch wenn ein nativer Windows-Worker das Modul lädt. - - Gleichlautend in docker/api/main.py; es gibt kein geteiltes Paket zwischen - den Containern. - """ - if not pfad or not wurzel: - return False - sauber = wurzel.rstrip("/") or "/" - return pfad == sauber or pfad.startswith(sauber + "/") - - -def _zielbasis(target_dir, disc_type: str) -> str: - """Ablagebasis: vom Nutzer gewähltes Ziel (validiert) oder Standard. - - posixpath statt os.path — aus demselben Grund wie in _arbeitsverzeichnis: - das sind IMMER Container-Pfade. os.path.normpath macht unter Windows - Backslashes daraus, und dann greift die MEDIA_ROOT-Prüfung nicht mehr, das - gewählte Ziel fiele still auf den Standard zurück. Aufgefallen 25.07.2026, - als der Test dafür erstmals unter Windows lief. Live war es nie: aufgerufen - wird nur aus rip_disc, und das ist auf Windows-Workern verriegelt. - """ - if target_dir: - normalisiert = posixpath.normpath(target_dir) - if unter_wurzel(normalisiert, MEDIA_ROOT): - return normalisiert - return posixpath.join(RIP_OUTPUT_DIR, disc_type) - - -def _arbeitsverzeichnis(einstellungen: dict, job_wahl: str = "") -> str: - """Basis für Roh-Rips. Reihenfolge: Wahl DIESES Rips → UI-Setting - `workDir` → Container-Default /app/temp/raw. - - Hintergrund (Commander 24.07.): Die VM-Platte (150 GB) reicht für BD-50, - aber eine 4K-UHD (bis 100 GB roh + Kompression daneben) sprengt sie — - das Arbeitsverzeichnis muss deshalb auf ein großes Ziel umlegbar sein. - - Pro Rip wählbar seit 25.07.2026 (Commander-Wunsch): beim „Rippen starten" - entscheidet man je Disc, wo die Rohdaten landen. Das Setting bleibt der - Standard — und ist damit der Wert, der bei Vollautomatik-Rips greift, bei - denen niemand gefragt wird. - """ - # posixpath statt os.path: Das sind IMMER Container-Pfade (/app/media/...), - # auch wenn ein nativer Windows-Worker dieses Modul lädt — der übersetzt - # sie erst später mit pfad_lokal(). os.path.normpath macht unter Windows - # Backslashes daraus, und dann greift die MEDIA_ROOT-Prüfung nicht mehr. - for kandidat in (job_wahl, einstellungen.get("workDir")): - wert = (kandidat or "").strip() - if wert: - normalisiert = posixpath.normpath(wert) - if unter_wurzel(normalisiert, MEDIA_ROOT): - return normalisiert - return RAW_DIR - - -def _frei_bytes(pfad: str) -> int: - """Freier Platz am Pfad (nächster existierender Elternordner zählt).""" - kandidat = pfad - while kandidat and not os.path.exists(kandidat): - kandidat = os.path.dirname(kandidat) - try: - return shutil.disk_usage(kandidat or "/").free - except OSError: - return -1 - - -def _platz_pruefen(pfad: str, benoetigt: int, zweck: str) -> str: - """Gibt eine Klartext-Fehlermeldung zurück, wenn der Platz nicht reicht — - sonst leeren String. VOR dem Rip geprüft: eine volle Platte nach 40 GB - ist der teuerste Fehlschlag, den es gibt.""" - frei = _frei_bytes(pfad) - if frei < 0 or benoetigt <= 0: - return "" # unbekannt → nicht raten, lieber rippen lassen - if frei >= benoetigt: - return "" - return ( - f"Zu wenig Platz für {zweck} in {pfad}: " - f"{frei / 1024**3:.1f} GB frei, benötigt ~{benoetigt / 1024**3:.1f} GB. " - "Abhilfe: anderes Ziel wählen oder das Arbeitsverzeichnis unter " - "Einstellungen → Verarbeitung auf eine große Freigabe legen." - ) - - -def _makemkv_key_anwenden(einstellungen: dict) -> None: - """MakeMKV-Beta-Key aus den UI-Einstellungen aktivieren (schlägt Env). - - Damit ist der Monats-Key ohne Rebuild/Neustart aktualisierbar - (Einstellungen → System). Format wie entrypoint.sh: settings.conf. - - Ergänzend statt überschreibend (Befund 25.07.2026): das Datenverzeichnis - ist jetzt persistent, und hier stand vorher ein open(..., "w") — das warf - vor JEDEM Rip alles andere aus der settings.conf, z. B. app_UpdateEnable - aus dem entrypoint. Beide Schreiber müssen gleich arbeiten, sonst kommt - der Fehler beim nächsten Rip still zurück. - """ - key = (einstellungen.get("makemkvAppKey") or "").strip() - if not key: - return - pfad = os.path.join(makemkv_daten.DATEN_DIR, "settings.conf") - try: - os.makedirs(makemkv_daten.DATEN_DIR, exist_ok=True) - try: - with open(pfad, encoding="utf-8", errors="replace") as f: - alt = f.read() - except OSError: - alt = "" - with open(pfad, "w", encoding="utf-8", newline="\n") as f: - f.write(makemkv_daten.settings_conf_zusammenfuehren(alt, key)) - except OSError as e: - db.add_log("warning", "worker", f"MakeMKV-Key konnte nicht gesetzt werden: {e}") - - -def _ordner_groesse(pfad: str) -> int: - """Belegter Platz eines Ordners in Bytes (0, wenn nicht lesbar).""" - summe = 0 - for wurzel, _, dateien in os.walk(pfad): - for name in dateien: - try: - summe += os.path.getsize(os.path.join(wurzel, name)) - except OSError: - pass - return summe - - -def _original_aufheben(job_id: str, raw_dir: str, final_dir: str) -> None: - """Roh-Rip zusätzlich aufheben — darf den Job NIEMALS scheitern lassen. - - Befund 25.07.2026 (Akira-UHD, echter Schaden): Hier stand ein nacktes - shutil.move(). Zwischen Arbeitsverzeichnis (/app/temp) und Ziel - (/app/media) scheitert os.rename mit EXDEV, shutil.move fällt auf Kopieren - zurück — eine 74-GB-Vollkopie auf dieselbe Platte, bis sie mit ENOSPC voll - war. Ergebnis: Platte 100 % voll, Worker-Container startete nicht mehr, und - der Job galt als FEHLGESCHLAGEN, obwohl die komprimierte Datei längst - fertig war. Der Nutzer sah nur „nichts da". - - ## Warum hier NICHT vorhergesagt, sondern versucht wird - - Die erste Fassung dieses Schutzes verglich `os.stat(...).st_dev` und - schloss aus gleichen Werten auf „reines Umhängen, kein Platz nötig". Am - 25.07.2026 im Worker-Container nachgemessen — beides zugleich wahr: - - st_dev /app/temp = 2050 - st_dev /app/media = 2050 → also identisch - os.rename(...) → EXDEV, „Invalid cross-device link" - - Der Kernel vergleicht bei rename() den **Mount**, nicht das Gerät. /app/temp - (Docker-Volume) und /app/media (Bind-Mount) sind zwei Mounts DERSELBEN - ext4-Partition. Die st_dev-Prüfung war deshalb wirkungslos: sie sah - „gleiches Dateisystem", übersprang die Platzprüfung, und shutil.move kopierte - doch. Der Schutz hätte genau den Schaden zugelassen, gegen den er gebaut war. - - Also: erst rename VERSUCHEN. Klappt es, ist es umgehängt und fertig. - Kommt EXDEV, steht fest, dass kopiert werden müsste — und erst dann wird - der Platz geprüft. Das ist keine Vermutung mehr, sondern die Antwort des - Kernels. - """ - ziel_original = os.path.join(final_dir, "original") - try: - # Der billige Weg zuerst — und er ist gleichzeitig der einzige - # verlässliche Test, ob überhaupt umgehängt werden kann. - try: - os.rename(raw_dir, ziel_original) - db.add_log( - "info", "worker", - f"Job {job_id}: Original behalten unter {ziel_original} (umgehängt, " - "kein zusätzlicher Platz nötig)", - ) - return - except OSError as e: - if e.errno != errno.EXDEV: - raise # etwas anderes ist schiefgelaufen → unten ehrlich melden - - # Ab hier ist eine echte Kopie unvermeidlich. Jetzt lohnt die Platzfrage. - benoetigt = _ordner_groesse(raw_dir) - frei = _frei_bytes(final_dir) - if frei < benoetigt * 1.05: - db.add_log( - "warning", "worker", - f"Job {job_id}: Original NICHT aufgehoben — Ziel liegt auf einem " - f"anderen Mount, es müsste kopiert werden. Dafür wären " - f"{benoetigt / 1024**3:.1f} GB nötig, frei sind nur " - f"{frei / 1024**3:.1f} GB. Die Roh-Datei bleibt unter " - f"{raw_dir} liegen. Abhilfe: Arbeitsverzeichnis " - "(Einstellungen → Verarbeitung) auf dieselbe Freigabe legen " - "wie das Ziel — dann wird nur umgehängt statt kopiert." - ) - return - shutil.move(raw_dir, ziel_original) - db.add_log("info", "worker", f"Job {job_id}: Original behalten unter {ziel_original}") - except OSError as e: - # Halb geschriebene Kopie wegräumen, sonst belegt sie für immer Platz. - shutil.rmtree(ziel_original, ignore_errors=True) - db.add_log( - "warning", "worker", - f"Job {job_id}: Original konnte nicht aufgehoben werden ({e}). " - f"Die komprimierte Datei ist fertig; die Roh-Datei bleibt unter {raw_dir}." - ) - - -def _abbruch_angefordert(job_id: str) -> bool: - """Kooperativer Abbruch: hat der Nutzer über die API abgebrochen?""" - return db.get_job_status(job_id) == "canceling" - - -def _serien_episoden_zuordnen(ausgabe: str, serie: str, staffel, meta: dict) -> list: - """Episoden-Zuordnung per Laufzeitabgleich (ARM-Wunde #395, besser gelöst). - - Laufzeiten der Staffel kommen von der API (TMDB); die der Dateien vom - HandBrake-Scan. Nur bei EINDEUTIGEM Treffer wird umbenannt. Wirft nie. - """ - if meta.get("source") != "tmdb" or meta.get("type") != "tv" or not meta.get("id"): - return ["Episoden-Zuordnung übersprungen (nur mit TMDB-Serien-Treffer möglich)"] - try: - antwort = requests.get( - f"{API_URL}/metadata/tv/{meta['id']}/season/{int(staffel)}", timeout=30 - ) - if antwort.status_code != 200: - return [f"Episoden-Laufzeiten nicht abrufbar (HTTP {antwort.status_code})"] - episoden = [ - (e["episode"], (e.get("runtime") or 0) * 60) - for e in antwort.json().get("episodes", []) - ] - except (requests.RequestException, ValueError) as e: - return [f"Episoden-Laufzeiten nicht abrufbar: {e}"] - if not episoden: - return ["Keine Episoden-Laufzeiten bei TMDB hinterlegt — Dateinamen bleiben"] - - dateien = sorted( - os.path.join(ausgabe, f) for f in os.listdir(ausgabe) if f.endswith(".mkv") - ) - if not dateien: - return [] - dauern = [lies_datei_dauer(p) for p in dateien] - if 0 in dauern: - return ["Datei-Laufzeiten nicht lesbar — Dateinamen bleiben"] - - zuordnung = medien.matche_episoden(dauern, episoden) - if not zuordnung: - return [ - "Keine EINDEUTIGE Episoden-Zuordnung über die Laufzeiten — " - "Dateinamen bleiben (lieber ehrlich als falsch benannt)" - ] - return medien.episoden_umbenennen(ausgabe, serie, staffel, zuordnung) - - -def _benachrichtigen(job_id: str, betreff: str, text: str, level: str) -> None: - """Webhook-Meldung bei Job-Ende — best effort, nie job-entscheidend.""" - einstellungen = db.get_settings(bei_fehler_leer=True) - url = (einstellungen.get("notificationWebhook") or "").strip() - if not url: - return - try: - notify.sende(url, betreff, text, level) - except Exception as e: - db.add_log("warning", "notify", f"Job {job_id}: Benachrichtigung fehlgeschlagen: {e}") - - -def _job_abschliessen(job_id: str, ergebnis: dict) -> None: - """Schreibt den Endzustand eines Jobs (completed/failed) nach Postgres, - bereitet bei Erfolg fürs Medien-System auf (NFO/Poster) und meldet das - Ende per Webhook (falls konfiguriert).""" - job = db.get_job(job_id) or {} - titel = job.get("title") or job_id[:8] - - if ergebnis.get("status") == "cancelled": - db.update_job( - job_id, status="failed", error="Abgebrochen durch Nutzer", - finished_at=db.utcnow(), - ) - db.add_log("warning", "worker", f"Job {job_id}: abgebrochen — Rohdaten bleiben erhalten") - _benachrichtigen(job_id, f'⚠️ Rippy: „{titel}" abgebrochen', - "Der Job wurde auf Wunsch abgebrochen — Rohdaten bleiben erhalten.", - "warning") - return - if ergebnis.get("status") == "success": - ausgabe = ergebnis.get("output_dir") - if ausgabe and (job.get("disc_type") in ("dvd", "bluray", "uhd")): - einstellungen = db.get_settings(bei_fehler_leer=True) - try: - meta = json.loads(job.get("meta") or "{}") - except ValueError: - meta = {} - serie = meta.get("series") - staffel = meta.get("season") - - # Serien-Rips: Episoden per Laufzeit zuordnen und umbenennen - # (Show S01E02.mkv) — nur bei EINDEUTIGER Zuordnung, sonst - # bleiben die MakeMKV-Namen (ehrlich geloggt). - if serie and staffel: - for meldung in _serien_episoden_zuordnen(ausgabe, serie, staffel, meta): - db.add_log("info", "worker", f"Job {job_id}: {meldung}") - - # Media-Server-Aufbereitung (Jellyfin/Emby/Kodi: NFO + Poster); - # bei Serien wandern tvshow.nfo/poster in den Serien-Ordner. - for meldung in medien.aufbereiten( - ausgabe, einstellungen.get("mediaServer") or "none", - serie or job.get("title") or "", meta.get("year"), meta, - serien_root=os.path.dirname(ausgabe) if serie and staffel else None, - ): - db.add_log("info", "worker", f"Job {job_id}: {meldung}") - - # Bibliotheks-Refresh (Jellyfin/Emby): der Server scannt sofort - meldung = medien.bibliothek_refresh( - einstellungen.get("mediaServer") or "none", - (einstellungen.get("jellyfinUrl") or "").strip(), - (einstellungen.get("jellyfinApiKey") or "").strip(), - ) - if meldung: - db.add_log("info", "worker", f"Job {job_id}: {meldung}") - db.update_job( - job_id, - status="completed", - progress=100, - output_path=ausgabe, - finished_at=db.utcnow(), - ) - db.add_log("success", "worker", f"Job {job_id}: abgeschlossen → {ausgabe}") - _benachrichtigen(job_id, f'✅ Rippy: „{titel}" ist fertig', - f"Abgelegt unter: {ausgabe}", "success") - else: - fehler = ergebnis.get("error", "unbekannter Fehler") - db.update_job( - job_id, - status="failed", - error=fehler, - finished_at=db.utcnow(), - ) - db.add_log("error", "worker", f"Job {job_id}: {fehler}") - _benachrichtigen(job_id, f'❌ Rippy: „{titel}" fehlgeschlagen', fehler, "error") - - @celery_app.task(bind=True, name="worker.tasks.rip_disc") def rip_disc(self, device_path: str, job_id: str, target_dir: str = None): - """Stufe 1: Rippt die Disc; bei Video folgt die Kompression als eigener Task. + """Rippt die Disc. Der Ablauf steht in ablauf.rippen().""" - target_dir (optional): vom Nutzer gewähltes Ablageziel unter /app/media — - dort eingehängte Shares (NFS/SMB) sind damit direkt wählbar. - """ - db.init_db() + def melde(zustand: dict) -> None: + self.update_state(state="PROGRESS", meta=zustand) - if detect_disc_type is None: - ergebnis = { - "status": "error", - "error": "Dieser Worker kann nur komprimieren (nativer " - "Transcode-Worker ohne Laufwerks-Zugriff) — gerippt " - "wird auf der Rippy-Hauptmaschine.", - } - _job_abschliessen(job_id, ergebnis) - return ergebnis + 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) - if _abbruch_angefordert(job_id): - _job_abschliessen(job_id, {"status": "cancelled"}) - return {"status": "cancelled"} - - disc_type = detect_disc_type(device_path) - - if disc_type in ("no_disc", "unknown"): - fehler = ( - "Keine Disc im Laufwerk" if disc_type == "no_disc" - else "Disc-Typ nicht erkennbar" - ) - db.update_job(job_id, status="failed", error=fehler, finished_at=db.utcnow()) - db.add_log("error", "worker", f"Job {job_id}: {fehler} ({device_path})") - return {"status": "error", "error": fehler, "disc_type": disc_type} - - db.update_job(job_id, status="running", disc_type=disc_type) - # Ab hier gilt: der Rip läuft, ist aber NICHT fertig. Stirbt der Job jetzt, - # ist jede Roh-Datei ein Bruchstück (siehe RIP_FERTIG oben). - db.meta_merken(job_id, **{RIP_FERTIG: False}) - db.add_log("info", "worker", f"Job {job_id}: {disc_type}-Rip gestartet ({device_path})") - - letzter = [-1] - - def fortschritt(progress: int, message: str = ""): - # MakeMKV liefert viele PRGV-Zeilen pro Sekunde — DB nur bei Änderung. - if progress == letzter[0]: - return - letzter[0] = progress - if _abbruch_angefordert(job_id): - raise RipAbbruch() - self.update_state( - state="PROGRESS", - meta={"progress": progress, "status": "ripping", "message": message}, - ) - db.update_job(job_id, progress=progress) - - # MakeMKV-Meldungen ins Log (Befund 25.07.2026): Bis dahin überlebte NUR - # die letzte Zeile ("Failed to open disc"), und die sagt nichts. Dass - # MakeMKV bei der UHD-Disc nicht einmal versucht, einen Schlüssel zu - # holen, war deshalb nur per Hand-Lauf im Container zu sehen. - # Gedrosselt, weil das UI global nur die letzten 200 Zeilen zeigt: jede - # Meldung höchstens einmal, insgesamt höchstens MAX_MELDUNGEN je Rip. - # Code 1003 ist MakeMKVs eigenes DEBUG-Rauschen (am 25.07. beobachtet). - MAX_MELDUNGEN = 40 - gesehen = set() - - def melde_makemkv(code: int, text: str): - if code == 1003 or len(gesehen) >= MAX_MELDUNGEN or text in gesehen: - return - gesehen.add(text) - db.add_log("info", "makemkv", f"Job {job_id}: {text[:300]}") - - einstellungen = db.get_settings(bei_fehler_leer=True) - ist_video = disc_type in ("dvd", "bluray", "uhd") - # Je Disc-Typ abwählbar (siehe komprimieren_fuer): 4K verlustfrei behalten, - # DVDs trotzdem schrumpfen — vorher gab es nur alles oder nichts. - transcode_an = ist_video and komprimieren_fuer(disc_type, einstellungen) - - if ist_video: - # UI-Key schlägt Env-Key — Monats-Key ohne Rebuild aktualisierbar - _makemkv_key_anwenden(einstellungen) - - # Sprechender Zielordner „ (Jahr)" statt Job-UUID — Jellyfin & Co. - # erkennen den Film am Ordnernamen. UUID bleibt Fallback ohne Titel. - # Serien-Rips landen stattdessen in /Season NN (Staffel-Flow). - job = db.get_job(job_id) or {} - try: - meta = json.loads(job.get("meta") or "{}") - except ValueError: - meta = {} - if meta.get("series") and meta.get("season"): - final_dir = medien.serien_ordner( - _zielbasis(target_dir, disc_type), meta["series"], meta["season"] - ) - else: - final_dir = medien.zielordner( - _zielbasis(target_dir, disc_type), - job.get("title") or "", meta.get("year"), job_id, - ) - - # Titel-Wahl: explizite Auswahl aus der Track-Tabelle schlägt alles; - # sonst „Nur Hauptfilm" (pro Rip wählbar, Fallback: Setting); - # bei Serien-Discs ohne Auswahl zählen ALLE Episoden-Titel. - titel_liste = [int(t) for t in (meta.get("titles") or []) if str(t).isdigit()] - nur_hauptfilm = meta.get("main_feature_only") - if nur_hauptfilm is None: - nur_hauptfilm = bool(einstellungen.get("mainFeatureOnly", False)) - if meta.get("series") or titel_liste: - nur_hauptfilm = False - # Geplantes Ziel sofort sichtbar machen (UI-Detail + retry-transcode) - db.update_job(job_id, output_path=final_dir) - - # meta["work_dir"] = die Wahl aus dem Rip-Dialog; leer = Setting/Default. - raw_dir = os.path.join(_arbeitsverzeichnis(einstellungen, meta.get("work_dir")), job_id) - - # Platz-Check VOR dem Rip: Disc-Größe ist per ioctl bekannt — eine volle - # Platte nach 40 GB wäre der teuerste Fehlschlag (4K-UHD: bis 100 GB roh). - try: - disc_bytes = disc_size_bytes(device_path) if disc_size_bytes else 0 - except OSError: - disc_bytes = 0 - marge = 1024**3 # 1 GB Sicherheitsabstand - if ist_video: - rip_ziel = raw_dir if transcode_an else final_dir - fehler = _platz_pruefen(rip_ziel, disc_bytes + marge, "den Roh-Rip") - if not fehler and transcode_an and einstellungen.get("keepOriginal", False): - fehler = _platz_pruefen(final_dir, disc_bytes + marge, "das Original (keepOriginal)") - if fehler: - ergebnis = {"status": "error", "error": fehler} - _job_abschliessen(job_id, ergebnis) - return ergebnis - - if disc_type == "cd": - ergebnis = rip_cd( - device_path, job_id, progress_cb=fortschritt, output_dir=final_dir, - auswerfen=bool(einstellungen.get("autoEject", True)), - ) - elif transcode_an: - # Roh-Rip ins Arbeitsverzeichnis (wird nach erfolgreicher Kompression gelöscht) - ergebnis = rip_video( - device_path, job_id, disc_type, - progress_cb=fortschritt, - output_dir=raw_dir, - nur_hauptfilm=nur_hauptfilm, - titel_liste=titel_liste, - log_cb=melde_makemkv, - ) - else: - ergebnis = rip_video( - device_path, job_id, disc_type, - progress_cb=fortschritt, output_dir=final_dir, - nur_hauptfilm=nur_hauptfilm, - titel_liste=titel_liste, - log_cb=melde_makemkv, - ) - - # MakeMKV-Fehler in Klartext übersetzen — "Failed to open disc" allein - # hilft niemandem. - if ergebnis.get("status") == "error": - fehler_text = ergebnis.get("error") or "" - if "volume key is unknown" in fehler_text: - # Befund 25.07.2026, auf BEIDEN Maschinen gemessen: Laufwerk und - # MakeMKV sind in Ordnung. makemkvcon unter Linux ruft die - # Disc-Schlüssel schlicht nie ab — die Windows-Version tut es - # (Meldung 3338). Hier stand vorher erst "Disc zu neu" und danach - # "der Schlüssel-Kanal ist tot"; beides war falsch und hat in die - # Irre geschickt. Herleitung im Kopf von makemkv_daten.py. - speicher = makemkv_daten.schluesselspeicher_status() - anzahl = speicher.get("schluessel", 0) - ergebnis["error"] += ( - " — Klartext: Laufwerk und Rippy sind in Ordnung, MakeMKV " - "liest die Disc. Es fehlt nur der Schlüssel dieser Pressung. " - "Der Grund: makemkvcon holt Schlüssel unter Linux nie selbst " - "nach — die Windows-Version schon. " - + ( - "Dieser Worker kennt aktuell GAR KEINEN Disc-Schlüssel. " - if not anzahl - else f"Dieser Worker kennt {anzahl} Disc-Schlüssel, " - "diese Pressung ist nicht dabei. " - ) - + "Abhilfe: MakeMKV auf einem Windows-PC installieren, die " - "Disc dort einmal öffnen, dann die Datei _private_data.tar " - "aus dem MakeMKV-Datenverzeichnis unter Einstellungen → " - "System hochladen. Wirkt ab dem nächsten Rip. Klappt auch " - "das nicht, kennt MakeMKV die Pressung selbst nicht — dann " - "hilft nur eine KEYDB.cfg (ebenfalls dort hochladbar) oder " - "das Einreichen des AACS-Dumps im MakeMKV-Forum, Bereich " - "'Ultra HD Blu-ray'. Der Dump steht unter Einstellungen → " - "System zum Download bereit." - ) - elif disc_type == "uhd" and "Failed to open disc" in fehler_text: - ergebnis["error"] += ( - " — 4K-UHD erkannt: Das Laufwerk kann UHD-Discs vermutlich " - "nicht entschlüsseln. Dafür ist eine LibreDrive-Firmware nötig " - "(MakeMKV-Forum 'Ultimate UHD Drives Flashing Guide'). " - "Normale BD/DVD gehen weiterhin." - ) - - # Automatischer Auswurf. Die Disc ist nach dem Rip nicht mehr nötig — die - # Kompression arbeitet auf der Datei, nicht am Laufwerk. - # - # Befund 25.07.2026: Die Einstellung („Disc nach erfolgreichem Ripping - # automatisch auswerfen", Standard ein) wurde von NIEMANDEM gelesen. Bei - # DVD/Blu-ray warf Rippy deshalb nie aus, bei Audio-CD dagegen immer, weil - # abcde `-x` fest verdrahtet bekam. Jetzt entscheidet die Einstellung beides. - if ergebnis.get("status") == "success" and ist_video and einstellungen.get("autoEject", True): - if wirf_disc_aus(device_path): - db.add_log("info", "worker", f"Job {job_id}: Disc ausgeworfen") - else: - db.add_log( - "warning", "worker", - f"Job {job_id}: Disc konnte nicht ausgeworfen werden ({device_path}) — " - "der Rip ist davon unberührt.", - ) - - if ergebnis.get("status") == "success" and transcode_an: - # 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(meta.get("transcode_node")) - db.update_job(job_id, status="transcoding", progress=0) - # Der Rip ist durch. Ab jetzt ist ein Fehlschlag mit „Neu komprimieren" - # zu heilen, ohne die Disc noch einmal zu lesen. - db.meta_merken(job_id, **{RIP_FERTIG: True}) - gezielt = meta.get("transcode_node") and ziel_queue != "transcode" - db.add_log("info", "worker", - f"Job {job_id}: Rip fertig, Kompression eingereiht" - + (f" → gezielt an {meta.get('transcode_node')}" if gezielt else "")) - transcode_files.apply_async( - args=[job_id, raw_dir, final_dir], - queue=ziel_queue, - ) - return {"status": "ripped", "raw_dir": raw_dir} - - _job_abschliessen(job_id, ergebnis) - return ergebnis - - -@celery_app.task(name="worker.tasks.ping_worker") -def ping_worker(): - """Winziger Task zum Beweisen des gezielten Routings: gibt Hostname + - Encoder zurück. Wird über worker_direct an EINEN Worker geschickt — landet - er beim richtigen, stimmt das Encoder-Routing (Encoder-Auswahl im UI).""" - import socket - - import caps - return {"hostname": socket.gethostname(), "encoders": caps.erkenne_encoder()} - - -@celery_app.task(name="worker.tasks.scan_tracks") -def scan_tracks(device_path: str): - """Titel-Liste der eingelegten Disc erfassen (für die Auswahl-Tabelle). - - Ergebnis landet in der settings-Tabelle (key 'tracks:') — die - API pollt darauf. Läuft auf der Haupt-Queue (nur der Rip-Worker hat das - Laufwerk). Ein Info-Lauf dauert je nach Disc 20–120 s. - """ - db.init_db() - key = f"tracks:{device_path}" - db.save_settings({"status": "running"}, key) - try: - # Ein Info-Lauf liefert BEIDES: Titel-Liste und die Sprachen. Die - # Sprach-Auskunft stand schon immer in derselben Ausgabe und wurde nur - # weggeworfen (Commander-Anforderung 26.07.2026: vor dem Rip fragen, - # was man haben will). - titel, streams = lies_titel_info(device_path) - sprachen = sprachen_zusammenfassen(streams) - db.save_settings( - {"status": "done", "tracks": titel, "sprachen": sprachen}, key) - db.add_log( - "info", "worker", - f"Titel-Scan {device_path}: {len(titel)} Titel, " - f"{len(sprachen['audio'])} Tonsprache(n), " - f"{len(sprachen['untertitel'])} Untertitelsprache(n)", - ) - return {"status": "done", "anzahl": len(titel)} - except Exception as e: - db.save_settings({"status": "error", "error": str(e)[:300]}, key) - db.add_log("error", "worker", f"Titel-Scan {device_path} fehlgeschlagen: {e}") - return {"status": "error", "error": str(e)} - - -def _ist_uebersetzt(container_pfad: str, lokal_pfad: str) -> bool: - """Hat RIPPY_PATH_MAP diesen Pfad wirklich angefasst? (pure, testbar)""" - return container_pfad != lokal_pfad - - -# Wie weit wird nach einem vorhandenen Elternordner gesucht. Begrenzt, damit auf -# einer toten Freigabe nicht endlos geklopft wird. -MAX_ELTERN_STUFEN = 12 - - -def erster_vorhandener_ordner(pfad: str, isdir=None, dirname=None) -> str: - """Der nächste EXISTIERENDE Ordner oberhalb von `pfad` — oder "". - - ⚠️ Vorfall 26.07.2026, im vollen Durchlauf aufgefallen: Der Rip lief sauber - durch (40 GB), die Kompression brach sofort ab mit „Dieser Worker erreicht - das Ziel nicht: \\\\NAS\\rippy\\movies\\Akira (1988)". Erreichbar war die - Freigabe sehr wohl — es fehlten schlicht ZWEI Ordner: `movies` und der - Filmordner darin. Die Prüfung ließ nur EINE fehlende Ebene durchgehen (sie - sah nach `dirname`), obwohl gleich darauf `os.makedirs` die ganze Kette - anlegt. Ein Ziel, das noch nie beschrieben wurde, war damit systematisch - unerreichbar — also jeder erste Film in einer neuen Ablage. - - Es zählt daher, ob IRGENDEIN Vorfahre existiert: Von dem aus kann - `makedirs` den Rest bauen. Existiert keiner, ist die Freigabe wirklich weg. - - `isdir`/`dirname` werden BEIM AUFRUF aufgelöst, nicht als Vorgabewert - gebunden — sonst zeigt die Vorgabe für immer auf die Funktion von damals, - und ein Ersetzen von `os.path.isdir` (im Test wie im Betrieb) ginge ins - Leere. Genau darüber bin ich beim Schreiben der Tests gestolpert. - """ - isdir = isdir or os.path.isdir - dirname = dirname or os.path.dirname - aktuell = dirname(pfad or "") - for _ in range(MAX_ELTERN_STUFEN): - if not aktuell: - return "" - if isdir(aktuell): - return aktuell - naechster = dirname(aktuell) - if naechster == aktuell: # Wurzel erreicht - return "" - aktuell = naechster - return "" - - -def _erreichbarkeit_pruefen(raw_container: str, raw_lokal: str, - final_container: str, final_lokal: str) -> str: - """Kann DIESER Worker Quelle und Ziel überhaupt sehen? Klartext oder "". - - Befund 25.07.2026 am Job 95afdc89, live: Das gezielte Routing an den - Windows-PC des Commanders funktionierte einwandfrei — der Worker nahm die - Aufgabe an und lehnte sie 182 ms später ab mit „Keine Roh-MKVs in - /app/media/rippy/… gefunden". Diese Meldung klang nach einem kaputten Rip, - obwohl der Rip vollständig war (79,6 GB lagen auf der NAS). Der Commander - schloss daraus, das System sähe den externen Encoder nicht — verständlich - und falsch. - - Die eigentliche Ursache: `/app/media/...` sind CONTAINER-Pfade. Ein nativer - Worker sieht sie nur, wenn RIPPY_PATH_MAP sie auf eine Freigabe übersetzt — - und dieses Mapping wurde von NIEMANDEM gesetzt. Es war also nicht möglich, - dass diese Kombination je funktioniert. - - Zweiter Stolperstein im selben Fall: Die Rohdaten lagen auf der NAS (für den - PC erreichbar), das ZIEL aber auf der VM-Platte (nicht erreichbar, dort - läuft kein Samba). Deshalb werden hier BEIDE Pfade geprüft, nicht nur die - Quelle — sonst scheitert es erst beim Schreiben, nach Stunden Rechenzeit. - """ - mapping = os.getenv("RIPPY_PATH_MAP", "") - fremder_worker = not os.path.isdir("/app") - - for zweck, container, lokal in ( - ("die Quelle (Rohdaten)", raw_container, raw_lokal), - ("das Ziel (fertige Datei)", final_container, final_lokal), - ): - if os.path.isdir(lokal): - continue - # Das ZIEL darf fehlen — es wird gleich mit `os.makedirs` angelegt, und - # zwar samt aller fehlenden Zwischenebenen. Es genügt also, dass - # irgendein Vorfahre existiert (Herleitung in erster_vorhandener_ordner). - if zweck.startswith("das Ziel") and erster_vorhandener_ordner(lokal): - continue - - text = [f"Dieser Worker erreicht {zweck} nicht: {lokal}"] - if fremder_worker and not mapping: - text.append( - "Ursache: RIPPY_PATH_MAP ist auf dieser Maschine nicht gesetzt. " - f'„{container}" ist ein Pfad INNERHALB des Rippy-Containers — ein ' - 'externer Worker sieht ihn nur, wenn er auf eine Netzwerk-Freigabe ' - 'übersetzt wird.' - ) - text.append( - "Abhilfe: Freigabe auf diesem Rechner erreichbar machen und in " - "start-tray.bat setzen, z. B. " - "set RIPPY_PATH_MAP=/app/media=\\\\NAS\\rippy-media" - ) - elif fremder_worker and mapping and not _ist_uebersetzt(container, lokal): - text.append( - f"RIPPY_PATH_MAP ist gesetzt ({mapping}), deckt diesen Pfad aber " - f'nicht ab: „{container}" wurde von keinem Eintrag übersetzt — ' - 'fehlt ein Präfix, oder ist es ein anderes Verzeichnis als erwartet?' - ) - elif fremder_worker and mapping: - # ⚠️ Hier stand bis zum 26.07.2026 dieselbe Meldung wie oben — „deckt - # diesen Pfad nicht ab" —, obwohl die Karte ihn nachweislich übersetzt - # hatte (der UNC-Pfad stand im selben Satz). Wer der Meldung folgte, - # suchte den Fehler in der Karte, während die Freigabe schlicht nicht - # verbunden war. Eine falsche Ursache ist teurer als gar keine. - text.append( - f'Übersetzt wurde er korrekt ({container} → {lokal}, Karte: ' - f'{mapping}) — dieser Pfad ist auf diesem Rechner nur gerade ' - 'nicht erreichbar. Ist die Freigabe verbunden? Anmeldedaten noch ' - 'gültig? Von Hand prüfen: den Pfad im Explorer öffnen.' - ) - else: - text.append( - "Der Pfad existiert nicht. Liegt das Arbeitsverzeichnis auf einer " - "Freigabe, die gerade nicht eingehängt ist?" - ) - text.append( - 'WICHTIG: Die Rohdaten sind NICHT verloren — nach der Korrektur ' - 'genügt „Neu komprimieren", ohne die Disc erneut zu rippen.' - ) - return " ".join(text) - return "" + 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): - """Stufe 2: HandBrake komprimiert die Roh-MKVs auf Arbeitsgröße. + """Komprimiert die Rohdaten. Der Ablauf steht in ablauf.komprimieren().""" + return ablauf.komprimieren(job_id, raw_dir, final_dir) - Erst wenn ALLE Dateien sauber komprimiert sind, wird das Roh-Verzeichnis - gelöscht — bricht die Kompression ab, bleibt das Original in /app/temp - liegen (kein Datenverlust wie bei ARMs berüchtigtem Move-Bug #1530). - Über POST /jobs/{id}/retry-transcode jederzeit neu anstoßbar. - """ - db.init_db() - # Native Worker (Windows) übersetzen Container-Pfade aufs Netzlaufwerk - raw_container, final_container = raw_dir, final_dir - raw_dir = pfad_lokal(raw_dir) - final_dir = pfad_lokal(final_dir) - fehler = _erreichbarkeit_pruefen(raw_container, raw_dir, final_container, final_dir) - if fehler: - ergebnis = {"status": "error", "error": fehler} - _job_abschliessen(job_id, ergebnis) - return ergebnis +@celery_app.task(name="worker.tasks.scan_tracks") +def scan_tracks(device_path: str): + return ablauf.scan_tracks(device_path) - quellen = sorted(glob.glob(os.path.join(raw_dir, "*.mkv"))) - if not quellen: - ergebnis = { - "status": "error", - "error": ( - f"Das Verzeichnis {raw_dir} ist erreichbar, enthält aber keine " - "MKV-Datei. Der Rip hat also nichts abgelegt (oder jemand hat die " - "Datei entfernt). Die Kompression lässt sich nach einem neuen Rip " - "erneut anstoßen." - ), - } - _job_abschliessen(job_id, ergebnis) - return ergebnis - einstellungen = db.get_settings(bei_fehler_leer=True) - # Preset nach Disc-Typ (Befund 25.07.2026): vorher lief JEDE Quelle durch - # dasselbe Preset — eine 4K-UHD wurde damit auf 1080p heruntergerechnet. - job = db.get_job(job_id) or {} - disc_type = job.get("disc_type") or "" - preset = preset_fuer(disc_type, einstellungen) - original_behalten = einstellungen.get("keepOriginal", False) +@celery_app.task(name="worker.tasks.ping_worker") +def ping_worker(): + return ablauf.ping_worker() - # Sprachauswahl: Wahl für DIESEN Rip (aus meta) schlägt die Einstellung. - # Leer heißt „alles behalten" — genau das Verhalten von vorher. - try: - job_meta = json.loads(job.get("meta") or "{}") - except ValueError: - job_meta = {} - audio_sprachen = sprachliste( - job_meta.get("audio_sprachen") or einstellungen.get("audioSprachen")) - untertitel_sprachen = sprachliste( - job_meta.get("untertitel_sprachen") or einstellungen.get("untertitelSprachen")) - os.makedirs(final_dir, exist_ok=True) - db.update_job(job_id, status="transcoding", progress=0, error=None) - # Wer hier ankommt, hat vollständige Quelldateien — sonst wäre oben schon - # abgebrochen worden. Das vermerkt die Marke auch für BESTANDSJOBS, die vor - # ihrer Einführung gerippt wurden und sie darum noch nicht tragen. - db.meta_merken(job_id, **{RIP_FERTIG: True}) - db.add_log( - "info", "worker", - f"Job {job_id}: Kompression gestartet ({len(quellen)} Datei(en), " - f"Disc-Typ '{disc_type or 'unbekannt'}', Preset '{preset}'" - + (f", Ton: {','.join(audio_sprachen)}" if audio_sprachen else ", Ton: alle") - + (f", Untertitel: {','.join(untertitel_sprachen)}" - if untertitel_sprachen else ", Untertitel: alle") - + ")", - ) - - anzahl = len(quellen) - letzter = [-1] - letzte_abbruchpruefung = [0.0] - - def abbruch_pruefen(): - """Zeitgesteuert prüfen, ob der Nutzer abgebrochen hat. Wirft RipAbbruch. - - Befund 25.07.2026 (Commander, am laufenden Akira-Job beobachtet): Der - Abbruch wurde NUR in `datei_fortschritt` geprüft — und diese Closure - stieg oben sofort wieder aus, wenn sich die Prozentzahl nicht geändert - hatte. Bei einem 4K-Encode, der pro Prozent eine halbe Stunde braucht, - sah „Abbrechen" entsprechend lange wirkungslos aus (gemessen: 3,4 min - zwischen Anforderung 18:30:17 und Bestätigung 18:33:38 — bei noch - langsamerem Fortschritt entsprechend mehr). - - Deshalb hängt die Prüfung jetzt an der Zeit statt am Fortschritt und - läuft bei JEDER Ausgabezeile von HandBrake — auch während des - Scan-Durchlaufs, der überhaupt keine Encode-Prozente liefert. - """ - jetzt = time.monotonic() - if jetzt - letzte_abbruchpruefung[0] < ABBRUCH_INTERVALL_SEKUNDEN: - return - letzte_abbruchpruefung[0] = jetzt - if _abbruch_angefordert(job_id): - raise RipAbbruch() - - for index, quelle in enumerate(quellen): - ziel = os.path.join(final_dir, os.path.basename(quelle)) - - def datei_fortschritt(p, _index=index): - gesamt = int((_index * 100 + p) / anzahl) - if gesamt == letzter[0]: - return - letzter[0] = gesamt - db.update_job(job_id, progress=min(99, gesamt)) - - hb = run_handbrake( - quelle, ziel, preset=preset, - progress_cb=datei_fortschritt, - abbruch_cb=abbruch_pruefen, - audio_sprachen=audio_sprachen, - untertitel_sprachen=untertitel_sprachen, - ) - if hb.get("status") == "cancelled": - _job_abschliessen(job_id, hb) - return hb - if hb.get("status") != "success": - ergebnis = { - "status": "error", - "error": ( - f"Kompression fehlgeschlagen bei {os.path.basename(quelle)}: " - f"{hb.get('error')} — Roh-Datei bleibt in /app/temp erhalten" - ), - } - _job_abschliessen(job_id, ergebnis) - return ergebnis - - if original_behalten: - _original_aufheben(job_id, raw_dir, final_dir) - else: - shutil.rmtree(raw_dir, ignore_errors=True) - - ergebnis = {"status": "success", "output_dir": final_dir} - _job_abschliessen(job_id, ergebnis) - return ergebnis +# Der Anzeigename dieses Knotens — wird vom Herzschlag in celery_app.py +# benutzt und stand bisher hier. +WORKER_NAME = os.getenv("WORKER_NAME", "") diff --git a/docker/worker/test_erreichbarkeit.py b/docker/worker/test_erreichbarkeit.py index 6717fb9..396e163 100644 --- a/docker/worker/test_erreichbarkeit.py +++ b/docker/worker/test_erreichbarkeit.py @@ -7,7 +7,7 @@ erreicht das Ziel nicht", obwohl die Freigabe erreichbar war. Es fehlten nur zwei noch nie angelegte Ordner. """ -import tasks +import ablauf as tasks # --- Der eigentliche Fehler: nur EINE fehlende Ebene war erlaubt ------------- diff --git a/docker/worker/test_medien.py b/docker/worker/test_medien.py index 359517b..ce0e089 100644 --- a/docker/worker/test_medien.py +++ b/docker/worker/test_medien.py @@ -105,7 +105,7 @@ def test_matche_episoden_ohne_treffer_gibt_none(): def test_pfad_lokal_uebersetzt_fuer_windows_worker(): - from tasks import pfad_lokal + from ablauf import pfad_lokal mapping = "/app/media=Z:\\media;/app/temp=Y:\\temp" assert pfad_lokal("/app/temp/raw/abc", mapping) == "Y:\\temp\\raw\\abc" @@ -128,7 +128,7 @@ def test_arbeitsverzeichnis_wahl_des_rips_schlaegt_die_einstellung(): Setting-Wert ist genau der, der bei Vollautomatik-Rips greift, weil dort niemand gefragt wird. """ - import tasks + import ablauf as tasks einst = {"workDir": "/app/media/movies"} assert tasks._arbeitsverzeichnis(einst, "/app/media/rippy") == "/app/media/rippy" @@ -144,7 +144,7 @@ def test_arbeitsverzeichnis_wahl_des_rips_schlaegt_die_einstellung(): def test_unter_wurzel_faellt_nicht_auf_praefix_namen_herein(): """Befund 25.07.2026: Elf Stellen prüften mit nacktem startswith(). „/app/media-boese/x" beginnt mit „/app/media", liegt aber außerhalb.""" - import tasks + import ablauf as tasks assert tasks.unter_wurzel("/app/media", "/app/media") is True assert tasks.unter_wurzel("/app/media/movies", "/app/media") is True @@ -160,7 +160,7 @@ def test_unter_wurzel_faellt_nicht_auf_praefix_namen_herein(): def test_zielbasis_lehnt_praefix_ausbruch_ab(): - import tasks + import ablauf as tasks assert tasks._zielbasis("/app/media/movies", "bluray") == "/app/media/movies" # Ausbruch per Praefix-Namen fällt auf den Standard zurück @@ -169,7 +169,7 @@ def test_zielbasis_lehnt_praefix_ausbruch_ab(): def test_arbeitsverzeichnis_lehnt_praefix_ausbruch_ab(): - import tasks + import ablauf as tasks assert tasks._arbeitsverzeichnis({}, "/app/media-boese") == tasks.RAW_DIR assert tasks._arbeitsverzeichnis({"workDir": "/app/mediaX"}) == tasks.RAW_DIR diff --git a/docker/worker/test_original_aufheben.py b/docker/worker/test_original_aufheben.py index e4570f6..127192a 100644 --- a/docker/worker/test_original_aufheben.py +++ b/docker/worker/test_original_aufheben.py @@ -13,7 +13,7 @@ import os import pytest -import tasks +import ablauf as tasks class FakeDb: diff --git a/docker/worker/test_ripping_helpers.py b/docker/worker/test_ripping_helpers.py index 6386ddb..68589da 100644 --- a/docker/worker/test_ripping_helpers.py +++ b/docker/worker/test_ripping_helpers.py @@ -40,7 +40,11 @@ def test_handbrake_cmd_arbeitet_auf_datei_nicht_geraet(): MakeMKV-Rip, nie das Laufwerk (die alte Direkt-am-Gerät-Pipeline war für Blu-rays prinzipiell funktionsunfähig).""" cmd = build_handbrake_cmd("/app/temp/raw/x/t00.mkv", "/app/media/bluray/x/t00.mkv") - assert cmd[0] == "HandBrakeCLI" + # Wie beim MakeMKV-Befehl: seit V2-4 steht hier der GEFUNDENE Pfad. Unter + # Windows liegt HandBrakeCLI in Program Files oder in Rippys eigenem + # Werkzeug-Ordner und nicht im PATH — mit dem nackten Namen faende + # `subprocess` es dort nie. + assert "handbrakecli" in cmd[0].lower() assert cmd[cmd.index("--input") + 1] == "/app/temp/raw/x/t00.mkv" assert cmd[cmd.index("--output") + 1] == "/app/media/bluray/x/t00.mkv" assert "--preset" in cmd diff --git a/docker/worker/test_zombies.py b/docker/worker/test_zombies.py index 6aa9d4e..da2a5d7 100644 --- a/docker/worker/test_zombies.py +++ b/docker/worker/test_zombies.py @@ -221,7 +221,11 @@ def test_arbeitsstati_deckt_ab_was_der_worker_wirklich_schreibt(): import os import re - pfad = os.path.join(os.path.dirname(os.path.abspath(__file__)), "tasks.py") + # Seit V2-4 steht der Ablauf in ablauf.py; tasks.py ist nur noch + # die Celery-Huelle. Dieser Waechter muss dorthin schauen, wo der + # Code WIRKLICH steht — sonst prueft er eine leere Datei und ist + # gruen, ohne etwas zu beweisen. + pfad = os.path.join(os.path.dirname(os.path.abspath(__file__)), "ablauf.py") quelle = open(pfad, encoding="utf-8").read() # Nur die Aufrufe, die wirklich die Job-Zeile aendern. diff --git a/packaging/windows/build.py b/packaging/windows/build.py index ce5007f..2a1ad9b 100644 --- a/packaging/windows/build.py +++ b/packaging/windows/build.py @@ -87,6 +87,9 @@ def pruefen() -> list: maengel.append("Eine Windows-.exe laesst sich nur unter Windows bauen.") if not os.path.isfile(os.path.join(REPO, "docker", "api", "main.py")): maengel.append("docker/api/main.py fehlt.") + if not os.path.isfile(os.path.join(REPO, "docker", "worker", "ablauf.py")): + maengel.append("docker/worker/ablauf.py fehlt — ohne den Ablauf kann " + "die App nicht rippen.") if not os.path.isdir(os.path.join(REPO, "docker", "ui", "dist")): maengel.append( "docker/ui/dist fehlt — die Oberflaeche ist nicht gebaut. " @@ -127,7 +130,11 @@ def bauen(ausgabe: str, version: str) -> str: "--specpath", arbeit, "--paths", os.path.join(REPO, "src"), "--paths", os.path.join(REPO, "docker", "api"), + "--paths", os.path.join(REPO, "docker", "worker"), "--add-data", os.path.join(REPO, "docker", "api") + trenner + "api", + # Die Worker-Module — ohne sie kann Rippy anzeigen, aber nicht rippen. + # ablauf.py, ripping.py, medien.py, caps.py und Nachbarn. + "--add-data", os.path.join(REPO, "docker", "worker") + trenner + "worker", "--add-data", os.path.join(REPO, "docker", "ui", "dist") + trenner + "ui", "--add-data", os.path.join(REPO, "deploy", "worker-windows", "rippy.ico") + trenner + ".", # Was PyInstaller nicht von allein findet: dynamisch importierte Module. diff --git a/src/rippy/daemon.py b/src/rippy/daemon.py index a7810f3..eb823dc 100644 --- a/src/rippy/daemon.py +++ b/src/rippy/daemon.py @@ -262,6 +262,23 @@ def starten(argv=None) -> int: app, ui = anwendung_bauen(werte) print(f" Oberflaeche {ui or '(nicht mitgeliefert)'}") + # Im Standalone-Betrieb arbeitet DIESER Prozess die Auftraege ab — + # es gibt keinen Worker daneben. Ohne diesen Schritt koennte Rippy alles + # anzeigen und nichts tun: Ein Rip wuerde eingereiht und laege dann fuer + # immer da, weil niemand ihn holt. + if werte.get("profil") == "standalone" and werte["queue"]["treiber"] == "lokal": + from rippy import standalone + from rippy.bus.memory import bus as ereignis_bus + + laeufer = standalone.starten( + knoten=werte["queue"].get("knoten") or None, + bus=ereignis_bus, store=store, + slots=werte["queue"].get("rip_slots", 1), + ) + print(f" Laeufer {len(laeufer)} (Auftraege werden hier abgearbeitet)") + else: + print(" Laeufer keiner — Auftraege gehen an Celery") + host = werte["server"]["host"] port = werte["server"]["port"] print(f" Adresse http://{'localhost' if host in ('0.0.0.0', '') else host}:{port}") diff --git a/src/rippy/queue/laeufer.py b/src/rippy/queue/laeufer.py new file mode 100644 index 0000000..599d9a8 --- /dev/null +++ b/src/rippy/queue/laeufer.py @@ -0,0 +1,152 @@ +"""Der Läufer: nimmt Aufträge aus der LocalQueue und arbeitet sie ab. + +## Wofür er da ist + +Im verteilten Betrieb macht Celery das: Ein Worker-Prozess wartet auf +Nachrichten und ruft die Task-Funktion auf. Im Standalone-Betrieb gibt es +keinen Broker und keinen zweiten Prozess — also übernimmt dieser Läufer die +Rolle. Er ist das letzte Stück, das gefehlt hat, damit Rippy unter Windows +nicht nur anzeigen, sondern auch **arbeiten** kann. + +## Die Lease ist der Kern, nicht ein Detail + +Ein Rip läuft 30 bis 90 Minuten. In dieser Zeit muss zweierlei stimmen: + +* **Niemand sonst darf denselben Auftrag anfassen.** Zwei Rips auf einem + Laufwerk wären zwei kaputte Dateien. +* **Stürzt der Läufer ab, muss der Auftrag zurückfallen.** Sonst bliebe er + für immer als „läuft" stehen — genau der Zustand, den v1 mit + `zombies.py` nachträglich einsammeln musste. + +Beides erledigt die Lease aus `rippy.queue.lokal`: Der Läufer verlängert sie +alle 15 Sekunden aus einem eigenen Thread. Hört er auf zu leben, läuft sie +nach 60 Sekunden ab und der Auftrag ist wieder frei. + +## Warum der Herzschlag NICHT im Arbeits-Thread liegt + +Weil er dann mit der Arbeit stillstünde. Der Rip verbringt seine Zeit in +`subprocess.communicate()` — dort läuft kein Python-Code, der nebenbei eine +Lease verlängern könnte. Der Herzschlag braucht deshalb einen eigenen Thread, +der nichts anderes tut. + +## Was der Läufer NICHT tut + +Er kennt den Rip-Vorgang nicht. Er bekommt eine Funktion `ausfuehren(auftrag)` +und ruft sie auf. Das hält ihn testbar (die Tests hier rippen nichts) und +erlaubt dieselbe Mechanik später für andere Auftragsarten. +""" + +import threading +import time +import traceback + +from rippy.queue import lokal + + +class Laeufer: + """Ein Arbeiter: holt einen Auftrag, hält ihn, arbeitet ihn ab.""" + + def __init__(self, ausfuehren, knoten: str, kann: set = None, + bus=None, store=None, + herzschlag_sekunden: float = lokal.LEASE_ERNEUERN_SEKUNDEN): + self._ausfuehren = ausfuehren + self.knoten = knoten + self.kann = kann or set() + self._bus = bus + self._store = store + self._herzschlag_sekunden = herzschlag_sekunden + self._laeuft = False + self.aktueller_auftrag = None + + # ── Ein Auftrag ───────────────────────────────────────────────────── + def einmal(self) -> bool: + """Holt EINEN Auftrag und arbeitet ihn ab. False, wenn nichts da war. + + Blockiert für die Dauer der Arbeit — bei einem Rip also Stunden. Das + ist Absicht: Der Aufrufer entscheidet, ob das in einem eigenen Thread + passiert. + """ + auftrag = lokal.uebernehmen(self.knoten, self.kann) + if auftrag is None: + return False + + self.aktueller_auftrag = auftrag + schlagen = self._herzschlag_starten(auftrag["id"]) + try: + ergebnis = self._ausfuehren(auftrag) + lokal.abschliessen(auftrag["id"], ergebnis if isinstance(ergebnis, dict) else {}) + self._melden("job.finished", auftrag, {"status": "fertig"}) + except Exception as e: + # Ein Fehlschlag darf den Laeufer NICHT mitnehmen — sonst bleibt + # nach einem kaputten Auftrag jeder weitere liegen. Der Auftrag + # geht zurueck in die Queue (bis MAX_VERSUCHE), der Grund wird + # protokolliert, und zwar mit Rueckverfolgung: Ein "Fehler" ohne + # Stelle ist beim naechsten Mal wertlos. + self._protokollieren(auftrag, e) + lokal.fehlgeschlagen(auftrag["id"], f"{type(e).__name__}: {e}") + self._melden("job.finished", auftrag, {"status": "fehler", "fehler": str(e)}) + finally: + schlagen.set() + self.aktueller_auftrag = None + return True + + def _herzschlag_starten(self, auftrag_id: str) -> threading.Event: + """Verlängert die Lease, solange gearbeitet wird. Eigener Thread.""" + fertig = threading.Event() + + def schlagen(): + while not fertig.wait(self._herzschlag_sekunden): + if not lokal.lebenszeichen(auftrag_id, self.knoten): + # Der Auftrag gehoert uns nicht mehr — jemand anderes hat + # ihn uebernommen, weil unsere Lease abgelaufen war. Dann + # ist Weiterarbeiten falsch, aber Abbrechen mitten im Rip + # waere schlimmer. Also: laut sagen und weiterlaufen. + self._protokollieren_text( + f"Auftrag {auftrag_id} ist an einen anderen Knoten " + "gefallen — die Lease war abgelaufen.") + return + + threading.Thread(target=schlagen, daemon=True, + name=f"lease-{auftrag_id[:8]}").start() + return fertig + + # ── Die Schleife ──────────────────────────────────────────────────── + def schleife(self, takt: float = 1.0) -> None: + """Läuft, bis `stoppen()` gerufen wird. Für einen eigenen Thread.""" + self._laeuft = True + while self._laeuft: + try: + if not self.einmal(): + time.sleep(takt) + except Exception as e: # die Schleife selbst darf nie sterben + self._protokollieren_text(f"Läufer-Fehler: {e}") + time.sleep(takt * 5) + + def stoppen(self) -> None: + self._laeuft = False + + # ── Melden ────────────────────────────────────────────────────────── + def _melden(self, typ: str, auftrag: dict, daten: dict) -> None: + if not self._bus: + return + try: + self._bus.senden(typ, {**daten, "art": auftrag.get("art")}, + entitaet="job", entitaet_id=auftrag.get("job_id")) + except Exception: + pass # eine Meldung darf die Arbeit nicht aufhalten + + def _protokollieren(self, auftrag: dict, fehler: Exception) -> None: + self._protokollieren_text( + f"Auftrag {auftrag.get('art')} für Job {auftrag.get('job_id')} " + f"gescheitert: {type(fehler).__name__}: {fehler}\n" + + "".join(traceback.format_exception_only(type(fehler), fehler)).strip() + ) + + def _protokollieren_text(self, text: str) -> None: + if self._store: + try: + self._store.add_log("error", "laeufer", text) + return + except Exception: + pass + print("[Läufer] " + text) diff --git a/src/rippy/queue/test_laeufer.py b/src/rippy/queue/test_laeufer.py new file mode 100644 index 0000000..c19be63 --- /dev/null +++ b/src/rippy/queue/test_laeufer.py @@ -0,0 +1,185 @@ +"""Der Läufer: hält er den Auftrag, gibt er ihn zurück, überlebt er Fehler? + +Gerippt wird hier nichts — die Arbeit ist eine eingespritzte Funktion. Geprüft +wird das Drumherum, und genau dort sitzen die Fehler, die man im Betrieb +teuer bezahlt: ein Läufer, der nach dem ersten kaputten Auftrag stehenbleibt, +oder eine Lease, die während eines dreistündigen Encodes abläuft. +""" + +import threading +import time + +import pytest + +from rippy import store +from rippy.queue import laeufer as laeufer_modul +from rippy.queue import lokal + + +@pytest.fixture +def q(tmp_path): + vorher = store.zustand_sichern() + store.verbinden(f"sqlite:///{(tmp_path / 'l.db').as_posix()}") + store.init_db() + yield lokal + store.engine_holen().dispose() + store.zustand_wiederherstellen(vorher) + + +class FakeBus: + def __init__(self): + self.gesendet = [] + + def senden(self, typ, daten=None, entitaet=None, entitaet_id=None): + self.gesendet.append((typ, entitaet_id, daten or {})) + + +class FakeStore: + def __init__(self): + self.protokoll = [] + + def add_log(self, level, quelle, text): + self.protokoll.append((level, quelle, text)) + + +# ── Grundfall ─────────────────────────────────────────────────────────── +def test_ohne_auftrag_passiert_nichts(q): + lauf = laeufer_modul.Laeufer(lambda a: {}, "knoten-a") + assert lauf.einmal() is False + + +def test_auftrag_wird_geholt_und_abgearbeitet(q): + gesehen = [] + q.einreihen("job1", "rip", payload={"device": "/dev/sr0"}) + lauf = laeufer_modul.Laeufer(lambda a: gesehen.append(a) or {"ok": True}, "knoten-a") + + assert lauf.einmal() is True + assert len(gesehen) == 1 + assert gesehen[0]["art"] == "rip" + assert q.offene_auftraege() == [] # abgeschlossen, nicht mehr offen + + +def test_zweiter_laeufer_bekommt_ihn_nicht(q): + """Die zentrale Zusage — sonst zwei Rips auf einem Laufwerk.""" + q.einreihen("job1", "rip") + lauf_a = laeufer_modul.Laeufer(lambda a: {}, "knoten-a") + lauf_b = laeufer_modul.Laeufer(lambda a: {}, "knoten-b") + + langsam = threading.Event() + lauf_a._ausfuehren = lambda a: langsam.wait(2) or {} + t = threading.Thread(target=lauf_a.einmal, daemon=True) + t.start() + time.sleep(0.3) + try: + assert lauf_b.einmal() is False, "Knoten B hat denselben Auftrag bekommen" + finally: + langsam.set() + t.join(timeout=5) + + +# ── Fehler ────────────────────────────────────────────────────────────── +def test_ein_kaputter_auftrag_nimmt_den_laeufer_nicht_mit(q): + """Ohne das bliebe nach dem ersten Fehlschlag JEDER weitere Auftrag + liegen — und im UI saehe es aus, als tue Rippy nichts mehr.""" + protokoll = FakeStore() + q.einreihen("job1", "rip") + + def kaputt(_): + raise RuntimeError("MakeMKV ist abgestuerzt") + + lauf = laeufer_modul.Laeufer(kaputt, "knoten-a", store=protokoll) + assert lauf.einmal() is True # er hat gearbeitet, nicht geworfen + assert any("MakeMKV ist abgestuerzt" in text for _, _, text in protokoll.protokoll) + + +def test_gescheiterter_auftrag_geht_zurueck_in_die_queue(q): + """Ein abgestuerzter Encoder soll den Job nicht endgueltig verlieren.""" + auftrag_id = q.einreihen("job1", "transcode") + + def kaputt(_): + raise RuntimeError("weg") + + laeufer_modul.Laeufer(kaputt, "knoten-a", store=FakeStore()).einmal() + offen = q.offene_auftraege() + assert len(offen) == 1 and offen[0]["id"] == auftrag_id + + +def test_fehler_wird_mit_typ_gemeldet(q): + """„Fehler" allein ist beim naechsten Mal wertlos. Der Typ gehoert dazu.""" + q.einreihen("job1", "rip") + protokoll = FakeStore() + + def kaputt(_): + raise FileNotFoundError("makemkvcon fehlt") + + laeufer_modul.Laeufer(kaputt, "knoten-a", store=protokoll).einmal() + text = " ".join(t for _, _, t in protokoll.protokoll) + assert "FileNotFoundError" in text + assert "makemkvcon fehlt" in text + + +# ── Lease ─────────────────────────────────────────────────────────────── +def test_lease_wird_waehrend_der_arbeit_verlaengert(q): + """DER Punkt bei einem Rip: Er dauert Stunden, die Lease 60 Sekunden. + Ohne Herzschlag fiele der Auftrag mitten im Lauf zurueck — und ein + zweiter Laeufer finge an, dieselbe Disc zu lesen. + """ + from datetime import timedelta + + auftrag_id = q.einreihen("job1", "rip") + weiter = threading.Event() + + def langsam(_): + # Waehrend der Arbeit die Lease kuenstlich ablaufen lassen — + # der Herzschlag muss sie zurueckholen. + with store.engine_holen().begin() as conn: + conn.execute( + store.auftraege.update() + .where(store.auftraege.c.id == auftrag_id) + .values(lease_until=store.utcnow() - timedelta(seconds=10))) + weiter.wait(3) + return {} + + lauf = laeufer_modul.Laeufer(langsam, "knoten-a", herzschlag_sekunden=0.2) + t = threading.Thread(target=lauf.einmal, daemon=True) + t.start() + time.sleep(1.0) # dem Herzschlag Zeit geben + try: + assert laeufer_modul.lokal.uebernehmen("knoten-b") is None, \ + "Die Lease wurde nicht verlaengert — ein zweiter Knoten kam ran" + finally: + weiter.set() + t.join(timeout=5) + + +def test_aktueller_auftrag_ist_sichtbar_und_danach_wieder_leer(q): + """Fuer das UI und fuer die Diagnose: Woran arbeitet dieser Knoten gerade?""" + q.einreihen("job1", "rip") + gesehen = {} + lauf = laeufer_modul.Laeufer( + lambda a: gesehen.update(id=lauf.aktueller_auftrag["id"]) or {}, "knoten-a") + lauf.einmal() + assert gesehen["id"] + assert lauf.aktueller_auftrag is None + + +# ── Meldungen ─────────────────────────────────────────────────────────── +def test_abschluss_wird_gemeldet(q): + bus = FakeBus() + q.einreihen("job1", "rip") + laeufer_modul.Laeufer(lambda a: {}, "knoten-a", bus=bus).einmal() + assert [t for t, _, _ in bus.gesendet] == ["job.finished"] + assert bus.gesendet[0][1] == "job1" + + +def test_ein_kaputter_bus_haelt_die_arbeit_nicht_auf(q): + """Eine Meldung ist Beiwerk. Wenn sie scheitert, ist der Rip trotzdem + fertig — und darf nicht als Fehlschlag gelten.""" + class KaputterBus: + def senden(self, *a, **k): + raise RuntimeError("Bus weg") + + q.einreihen("job1", "rip") + lauf = laeufer_modul.Laeufer(lambda a: {"ok": True}, "knoten-a", bus=KaputterBus()) + assert lauf.einmal() is True + assert q.offene_auftraege() == [] diff --git a/src/rippy/standalone.py b/src/rippy/standalone.py new file mode 100644 index 0000000..f622ea0 --- /dev/null +++ b/src/rippy/standalone.py @@ -0,0 +1,243 @@ +"""Standalone-Betrieb: Auftraege lokal zustellen und im selben Prozess abarbeiten. + +## Was hier zusammenkommt + +Das ist das letzte Stueck, damit Rippy unter Windows nicht nur ANZEIGEN, +sondern auch ARBEITEN kann: + + celery_client.zusteller_setzen(…) die API reiht lokal ein statt an Celery + rippy.queue.lokal die Auftragstabelle mit Lease + rippy.queue.laeufer holt sie und arbeitet sie ab + ablauf.rippen / .komprimieren der Ablauf selbst, ohne Celery + +Der Ablauf ist derselbe wie im Docker-Betrieb — Zeile fuer Zeile dieselbe +Datei. Was sich unterscheidet, ist allein die Zustellung. + +## Warum die Task-NAMEN erhalten bleiben + +Die API schickt `"worker.tasks.rip_disc"` los. Dieser Name ist der Vertrag +zwischen API und Worker; `test_api_smoke.py` nagelt ihn ausdruecklich fest. +Hier wird er deshalb NICHT durch etwas Eigenes ersetzt, sondern uebersetzt — +so bleibt beides gueltig, und ein Betrieb laesst sich umstellen, ohne dass +die API etwas davon merkt. + +## Warum die Kompression NICHT weitergereicht wird + +`ablauf.rippen()` nimmt einen Rueckruf `weiterreichen`, mit dem die +Kompression an einen anderen Knoten geht. Im Standalone-Betrieb gibt es +keinen anderen Knoten — hier komprimiert derselbe Prozess weiter. Deshalb +wird der Rueckruf bewusst NICHT gestellt: `ablauf` ruft dann selbst +`komprimieren()` auf, im selben Auftrag, ohne Umweg ueber die Queue. + +Das ist auch das Ehrlichere: Eine Queue mit genau einem Bearbeiter, die +Arbeit an sich selbst weiterreicht, ist nur Zeremonie. +""" + +import os +import sys +import threading + +# Die Auftragsarten, die die API kennt — auf die Formen, die `ablauf` versteht. +AUFTRAGSARTEN = { + "worker.tasks.rip_disc": "rip", + "worker.tasks.transcode_files": "transcode", + "worker.tasks.scan_tracks": "scan", +} + + +def _worker_pfad() -> str: + """Wo liegen die Worker-Module (ablauf.py, ripping.py, medien.py)? + + Dieselben drei Faelle wie bei der API (siehe `daemon._api_pfad`): aus dem + Repo gestartet, im PyInstaller-Paket, oder installiert. + """ + gebuendelt = getattr(sys, "_MEIPASS", None) + if gebuendelt: + return os.path.join(gebuendelt, "worker") + + hier = os.path.dirname(os.path.abspath(__file__)) + repo = os.path.dirname(os.path.dirname(hier)) + kandidaten = [ + os.path.join(repo, "docker", "worker"), + os.path.join(os.path.dirname(sys.executable), "worker"), + ] + for pfad in kandidaten: + if os.path.isfile(os.path.join(pfad, "ablauf.py")): + return pfad + raise RuntimeError( + "Die Worker-Module wurden nicht gefunden. Gesucht wurde in:\n " + + "\n ".join(kandidaten) + + "\nOhne sie kann Rippy nicht rippen." + ) + + +def _ablauf(): + """`ablauf` importieren — erst beim ersten Bedarf. + + Nicht auf Modulebene: Das Modul zieht `ripping`, `medien` und `db` nach, + und die brauchen einen gesetzten sys.path. Ausserdem soll der Import + nicht schon beim Start des Servers passieren, sondern dann, wenn wirklich + gearbeitet wird. + """ + pfad = _worker_pfad() + if pfad not in sys.path: + sys.path.insert(0, pfad) + import ablauf + + return ablauf + + +# ── Zustellung ────────────────────────────────────────────────────────── +def zustellung_bauen(bus=None): + """Gibt die Funktion zurueck, die `celery_client.zusteller_setzen` erwartet.""" + from rippy.queue import lokal + + def zustellen(task_name: str, args: list, queue: str = None): + art = AUFTRAGSARTEN.get(task_name) + if art is None: + raise ValueError( + f"Unbekannte Auftragsart {task_name!r}. Bekannt sind: " + + ", ".join(sorted(AUFTRAGSARTEN)) + + ". Wer eine neue einfuehrt, traegt sie in " + "rippy/standalone.py AUFTRAGSARTEN ein — sonst landet sie " + "im Standalone-Betrieb nirgends, und zwar lautlos." + ) + if art == "rip": + device_path, job_id, target_dir = (list(args) + [None, None, None])[:3] + auftrag_id = lokal.einreihen( + job_id, "rip", faehigkeiten={"art": "rip"}, + payload={"device_path": device_path, "target_dir": target_dir}) + elif art == "transcode": + job_id, raw_dir, final_dir = (list(args) + [None, None, None])[:3] + auftrag_id = lokal.einreihen( + job_id, "transcode", faehigkeiten={"art": "transcode"}, + payload={"raw_dir": raw_dir, "final_dir": final_dir}) + else: # scan + device_path = args[0] if args else "" + auftrag_id = lokal.einreihen( + "", "scan", faehigkeiten={"art": "scan"}, + payload={"device_path": device_path}) + + if bus: + try: + bus.senden("system.notice", + {"level": "info", "text": f"{art} eingereiht"}) + except Exception: + pass + return auftrag_id + + return zustellen + + +# ── Ausfuehrung ───────────────────────────────────────────────────────── +def ausfuehren(auftrag: dict): + """Arbeitet EINEN Auftrag ab — der Laeufer ruft das auf.""" + import json + + ablauf = _ablauf() + try: + payload = json.loads(auftrag.get("payload") or "{}") + except (ValueError, TypeError): + payload = {} + + art = auftrag.get("art") + if art == "rip": + # KEIN `weiterreichen`: Im Standalone-Betrieb komprimiert derselbe + # Prozess weiter (siehe Modul-Docstring). + return ablauf.rippen(payload.get("device_path"), auftrag["job_id"], + payload.get("target_dir")) + if art == "transcode": + return ablauf.komprimieren(auftrag["job_id"], payload.get("raw_dir"), + payload.get("final_dir")) + if art == "scan": + return ablauf.scan_tracks(payload.get("device_path")) + raise ValueError(f"Unbekannte Auftragsart im Auftrag: {art!r}") + + +# ── Herzschlag ────────────────────────────────────────────────────────── +HERZSCHLAG_SEKUNDEN = 60 + + +def herzschlag_starten(knoten: str, store, takt: float = HERZSCHLAG_SEKUNDEN): + """Meldet DIESEN Prozess als Arbeiter — mit seinen Faehigkeiten. + + ## Warum das noetig ist + + `/capabilities` liest die `workers`-Tabelle. Gefuellt hat sie bisher der + Celery-Herzschlag in `celery_app.py`. Im Standalone-Betrieb gibt es den + nicht — also stand dort „0 Worker", und im UI blieb die Encoder-Auswahl + LEER. Rippy haette auf einem Rechner mit AMD-Hardwarebeschleunigung + behauptet, es koenne nichts. + + Gemeldet wird, was `caps.py` MISST: die Encoder, die HandBrake auf dieser + Maschine wirklich anbietet, die Preset-Namen dieser HandBrake-Fassung, + CPU-Modell, Kerne und Vektorbefehle. Nichts davon wird behauptet — das + ist die Lehre aus Etappe 19 („Encoder wurden behauptet statt gemessen"). + """ + import threading + import time + + def schlagen(): + while True: + try: + pfad = _worker_pfad() + if pfad not in sys.path: + sys.path.insert(0, pfad) + import caps + + store.save_worker(knoten, caps.erkenne_encoder(), + caps.werkzeug_versionen()) + except Exception as e: + # Nicht sterben — aber auch nicht schweigen. + try: + store.add_log("warning", "standalone", + f"Faehigkeiten-Meldung fehlgeschlagen: {e}") + except Exception: + pass + time.sleep(takt) + + threading.Thread(target=schlagen, daemon=True, name="rippy-herzschlag").start() + + +# ── Alles zusammen ────────────────────────────────────────────────────── +def starten(knoten: str = None, bus=None, store=None, slots: int = 1) -> list: + """Zustellung umbiegen und die Laeufer starten. Gibt die Laeufer zurueck. + + `slots` ist die Zahl gleichzeitiger Auftraege. Vorgabe 1: Ein Rip belegt + das Laufwerk, und zwei gleichzeitige Encodes auf einem Rechner sind + langsamer als zwei nacheinander. + """ + import socket + + from rippy.queue.laeufer import Laeufer + + knoten = knoten or socket.gethostname() + + # Die API reiht ab jetzt lokal ein statt an Celery. + sys.path.insert(0, os.path.join( + os.path.dirname(os.path.dirname(os.path.dirname( + os.path.abspath(__file__)))), "docker", "api")) + try: + import celery_client + except ImportError: # im Paket liegt sie woanders + from rippy.daemon import _api_pfad + + sys.path.insert(0, _api_pfad()) + import celery_client + + celery_client.zusteller_setzen(zustellung_bauen(bus)) + + # Sich selbst als Arbeiter melden — sonst zeigt das UI "0 Worker" und + # eine leere Encoder-Auswahl, obwohl dieser Rechner alles kann. + if store is not None: + herzschlag_starten(knoten, store) + + laeufer = [] + for nummer in range(max(1, slots)): + lauf = Laeufer(ausfuehren, f"{knoten}#{nummer}", + kann={"art=rip", "art=transcode", "art=scan"}, + bus=bus, store=store) + threading.Thread(target=lauf.schleife, daemon=True, + name=f"rippy-laeufer-{nummer}").start() + laeufer.append(lauf) + return laeufer diff --git a/src/rippy/test_kette.py b/src/rippy/test_kette.py new file mode 100644 index 0000000..17f5933 --- /dev/null +++ b/src/rippy/test_kette.py @@ -0,0 +1,137 @@ +"""Die ganze Kette im Standalone-Betrieb — von der Zustellung bis zum Job-Ende. + +## Was hier bewiesen wird + + celery_client.abschicken(...) die API stellt zu + -> standalone.zustellung uebersetzt in einen Auftrag + -> lokal.einreihen Auftragstabelle mit Lease + -> Laeufer.einmal() holt ihn und arbeitet ihn ab + -> ablauf.rippen() der echte Ablauf, ohne Celery + -> db.update_job() der Job endet in einem ehrlichen Zustand + +Das ist der Weg, der unter Windows bis V2-4 gar nicht existierte: Dort haette +die API einen Rip eingereiht, und er waere fuer immer liegengeblieben, weil +niemand ihn holt. + +## Warum ohne Laufwerk getestet wird + +Es gibt hier keins (es haengt an der VM). Das ist kein Mangel, sondern der +wichtigere Fall: Ein Rip auf ein Geraet, das nicht antwortet, muss **ehrlich +scheitern** — mit einem Zustand und einer Begruendung, die im UI ankommen. +Was NICHT passieren darf: dass der Job auf „pending" stehenbleibt und +niemand erfaehrt, warum nichts geschieht. +""" + +import time + +import pytest + +from rippy import standalone, store +from rippy.queue import laeufer as laeufer_modul +from rippy.queue import lokal + + +@pytest.fixture +def umgebung(tmp_path): + vorher = store.zustand_sichern() + store.verbinden(f"sqlite:///{(tmp_path / 'kette.db').as_posix()}") + store.init_db() + yield store + store.engine_holen().dispose() + store.zustand_wiederherstellen(vorher) + + +class FakeBus: + def __init__(self): + self.gesendet = [] + + def senden(self, typ, daten=None, entitaet=None, entitaet_id=None): + self.gesendet.append((typ, entitaet_id, daten or {})) + + +def test_die_ganze_kette_bis_zum_ehrlichen_fehlschlag(umgebung): + """Ein Rip auf ein Geraet, das es nicht gibt. + + Erwartet wird KEIN Erfolg — erwartet wird, dass der Job in einem + ENDZUSTAND landet und die Begruendung im Log steht. Ein Job, der auf + „pending" haengenbleibt, waere der schlimmere Ausgang: Der Nutzer saehe + einen Balken, der sich nie bewegt, und nirgends stuende warum. + """ + umgebung.insert_job("job-kette", "/dev/gibt-es-nicht", disc_type=None) + + bus = FakeBus() + zustellen = standalone.zustellung_bauen(bus) + zustellen("worker.tasks.rip_disc", ["/dev/gibt-es-nicht", "job-kette", None]) + + assert len(lokal.offene_auftraege()) == 1, "Der Auftrag wurde nicht eingereiht" + + lauf = laeufer_modul.Laeufer(standalone.ausfuehren, "test-knoten", + kann={"art=rip"}, bus=bus, store=umgebung) + assert lauf.einmal() is True, "Der Laeufer hat den Auftrag nicht geholt" + + job = umgebung.get_job("job-kette") + assert job["status"] in ("failed", "completed"), ( + f"Der Job steht auf {job['status']!r} — er muss in einem Endzustand " + "landen, sonst wartet der Nutzer auf etwas, das nie passiert.") + assert job["status"] == "failed" + assert job["error"], "Ein Fehlschlag ohne Begruendung ist im UI wertlos" + + meldungen = " ".join(z["message"] for z in umgebung.list_logs(limit=50)) + assert "job-kette" in meldungen, "Im Log steht nichts ueber diesen Job" + + assert lokal.offene_auftraege() == [], "Der Auftrag blieb in der Queue liegen" + + +def test_der_laeufer_bleibt_danach_arbeitsfaehig(umgebung): + """Nach einem Fehlschlag muss der naechste Auftrag trotzdem laufen — + sonst legt ein einziger kaputter Rip die ganze Installation lahm.""" + umgebung.insert_job("job-a", "/dev/gibt-es-nicht") + umgebung.insert_job("job-b", "/dev/gibt-es-auch-nicht") + zustellen = standalone.zustellung_bauen() + zustellen("worker.tasks.rip_disc", ["/dev/gibt-es-nicht", "job-a", None]) + zustellen("worker.tasks.rip_disc", ["/dev/gibt-es-auch-nicht", "job-b", None]) + + lauf = laeufer_modul.Laeufer(standalone.ausfuehren, "test-knoten", + kann={"art=rip"}, store=umgebung) + assert lauf.einmal() is True + assert lauf.einmal() is True + assert umgebung.get_job("job-a")["status"] == "failed" + assert umgebung.get_job("job-b")["status"] == "failed" + + +def test_die_api_stellt_ueber_dieselbe_stelle_zu(umgebung, monkeypatch): + """`celery_client.start_rip` ist der Weg, den POST /jobs geht. Er MUSS + im Standalone-Betrieb in der lokalen Queue landen — sonst reiht die API + an einen Broker ein, den es nicht gibt, und der Job verschwindet.""" + import sys + + from rippy.daemon import _api_pfad + + if _api_pfad() not in sys.path: + sys.path.insert(0, _api_pfad()) + import celery_client + + umgebung.insert_job("job-api", "/dev/sr0") + monkeypatch.setattr(celery_client, "_zusteller", standalone.zustellung_bauen()) + + celery_client.start_rip("/dev/sr0", "job-api", "/app/media/x") + + offen = lokal.offene_auftraege() + assert len(offen) == 1 + assert offen[0]["job_id"] == "job-api" + assert offen[0]["art"] == "rip" + + +def test_ein_rip_haengt_nicht_ewig(umgebung): + """Zeitgrenze als Zusage: Ein Auftrag auf ein totes Geraet muss ZUEGIG + scheitern. Haengt er, merkt es im Betrieb niemand — der Job steht auf + „laeuft", und der Laeufer nimmt keinen weiteren an.""" + umgebung.insert_job("job-zeit", "/dev/nichts") + standalone.zustellung_bauen()("worker.tasks.rip_disc", ["/dev/nichts", "job-zeit", None]) + + lauf = laeufer_modul.Laeufer(standalone.ausfuehren, "k", kann={"art=rip"}, + store=umgebung) + start = time.monotonic() + lauf.einmal() + dauer = time.monotonic() - start + assert dauer < 30, f"Der Fehlschlag brauchte {dauer:.1f}s — das ist zu lang" diff --git a/src/rippy/test_standalone.py b/src/rippy/test_standalone.py new file mode 100644 index 0000000..c9aaebb --- /dev/null +++ b/src/rippy/test_standalone.py @@ -0,0 +1,100 @@ +"""Die lokale Zustellung: kommt der Auftrag an, und in der richtigen Form? + +Gerippt wird nichts — geprueft wird die Uebersetzung von Celery-Task-Namen in +Auftraege und die Frage, was bei einem unbekannten Namen passiert. +""" + +import json + +import pytest + +from rippy import standalone, store +from rippy.queue import lokal + + +@pytest.fixture +def q(tmp_path): + vorher = store.zustand_sichern() + store.verbinden(f"sqlite:///{(tmp_path / 's.db').as_posix()}") + store.init_db() + yield lokal + store.engine_holen().dispose() + store.zustand_wiederherstellen(vorher) + + +def test_rip_wird_zum_rip_auftrag(q): + zustellen = standalone.zustellung_bauen() + zustellen("worker.tasks.rip_disc", ["/dev/sr0", "job1", "/app/media/x"]) + + offen = q.offene_auftraege() + assert len(offen) == 1 + auftrag = offen[0] + assert auftrag["art"] == "rip" + assert auftrag["job_id"] == "job1" + nutzlast = json.loads(auftrag["payload"]) + assert nutzlast["device_path"] == "/dev/sr0" + assert nutzlast["target_dir"] == "/app/media/x" + + +def test_transcode_wird_zum_transcode_auftrag(q): + standalone.zustellung_bauen()( + "worker.tasks.transcode_files", ["job2", "/raw", "/final"]) + auftrag = q.offene_auftraege()[0] + assert auftrag["art"] == "transcode" + nutzlast = json.loads(auftrag["payload"]) + assert nutzlast["raw_dir"] == "/raw" and nutzlast["final_dir"] == "/final" + + +def test_scan_braucht_keinen_job(q): + """Ein Track-Scan gehoert zu keinem Job — er passiert VOR dem Anlegen.""" + standalone.zustellung_bauen()("worker.tasks.scan_tracks", ["/dev/sr0"]) + auftrag = q.offene_auftraege()[0] + assert auftrag["art"] == "scan" + assert auftrag["job_id"] == "" + + +def test_fehlendes_zielverzeichnis_ist_kein_fehler(q): + """POST /jobs darf target_dir weglassen — dann gilt der Standard.""" + standalone.zustellung_bauen()("worker.tasks.rip_disc", ["/dev/sr0", "job3"]) + assert json.loads(q.offene_auftraege()[0]["payload"])["target_dir"] is None + + +def test_unbekannter_taskname_fliegt_auf(q): + """Wuerde er stillschweigend verworfen, waere der Auftrag weg und niemand + wuesste warum — der Nutzer sieht nur einen Job, der nie anfaengt.""" + with pytest.raises(ValueError, match="Unbekannte Auftragsart"): + standalone.zustellung_bauen()("worker.tasks.gibts_nicht", ["x"]) + + +def test_faehigkeiten_werden_gesetzt(q): + """Damit ein reiner Encode-Knoten spaeter keine Rip-Auftraege bekommt.""" + zustellen = standalone.zustellung_bauen() + zustellen("worker.tasks.rip_disc", ["/dev/sr0", "job1"]) + zustellen("worker.tasks.transcode_files", ["job1", "/raw", "/final"]) + arten = {a["art"]: json.loads(a["faehigkeiten"]) for a in q.offene_auftraege()} + assert arten["rip"] == {"art": "rip"} + assert arten["transcode"] == {"art": "transcode"} + + +def test_ein_kaputter_bus_haelt_die_zustellung_nicht_auf(q): + class KaputterBus: + def senden(self, *a, **k): + raise RuntimeError("weg") + + standalone.zustellung_bauen(KaputterBus())( + "worker.tasks.rip_disc", ["/dev/sr0", "job1"]) + assert len(q.offene_auftraege()) == 1 + + +def test_ausfuehren_lehnt_unbekannte_art_ab(): + with pytest.raises(ValueError, match="Unbekannte Auftragsart"): + standalone.ausfuehren({"art": "quatsch", "job_id": "j", "payload": "{}"}) + + +def test_worker_pfad_wird_gefunden(): + """Ohne die Worker-Module kann Rippy nicht rippen. Wenn sich die + Verzeichnisstruktur aendert, soll das HIER auffallen.""" + import os + + pfad = standalone._worker_pfad() + assert os.path.isfile(os.path.join(pfad, "ablauf.py")) diff --git a/src/rippy/windows_app.py b/src/rippy/windows_app.py index 1f3352b..749888b 100644 --- a/src/rippy/windows_app.py +++ b/src/rippy/windows_app.py @@ -262,6 +262,13 @@ class Dienst: store.verbinden(config.datenbank_url(werte)) app, _ = daemon.anwendung_bauen(werte) + # Ohne diesen Schritt koennte Rippy alles anzeigen und nichts tun. + from rippy import standalone + from rippy.bus.memory import bus as ereignis_bus + + standalone.starten(bus=ereignis_bus, store=store, + slots=werte["queue"].get("rip_slots", 1)) + import uvicorn self._server = uvicorn.Server(uvicorn.Config(