"""Zentraler Rip-Task: erkennt den Disc-Typ, rippt, komprimiert, schreibt Status. Die API legt beim POST /jobs die Job-Zeile an und schickt diesen Task los — der Worker hält die Zeile aktuell (running → transcoding → completed/failed) und schreibt Ereignisse ins Log. Das UI liest beides über die API. 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). """ import os import shutil import db from celery_app import celery_app from detection import detect_disc_type from ripping import ( DEFAULT_HB_PRESET, RIP_OUTPUT_DIR, rip_cd, rip_video, run_handbrake, ) RAW_DIR = os.getenv("RAW_DIR", "/app/temp/raw") @celery_app.task(bind=True, name="worker.tasks.rip_disc") def rip_disc(self, device_path: str, job_id: str): """Rippt eine Disc basierend auf ihrem Typ; job_id ist die DB-Zeile der API.""" db.init_db() 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) 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 self.update_state( state="PROGRESS", meta={"progress": progress, "status": "ripping", "message": message}, ) db.update_job(job_id, progress=progress) einstellungen = db.get_settings() transcode_an = ( disc_type in ("dvd", "bluray") and einstellungen.get("transcodeEnabled", True) ) if disc_type == "cd": ergebnis = rip_cd(device_path, job_id, progress_cb=fortschritt) elif transcode_an: # Stufe 1: Roh-Rip nach /app/temp (wird nach der Kompression gelöscht) ergebnis = rip_video( device_path, job_id, disc_type, progress_cb=fortschritt, output_dir=os.path.join(RAW_DIR, job_id), ) else: ergebnis = rip_video(device_path, job_id, disc_type, progress_cb=fortschritt) if ergebnis.get("status") == "success" and transcode_an: ergebnis = _komprimiere(job_id, disc_type, ergebnis, einstellungen) if ergebnis.get("status") == "success": db.update_job( job_id, status="completed", progress=100, output_path=ergebnis.get("output_dir"), finished_at=db.utcnow(), ) db.add_log("success", "worker", f"Job {job_id}: Rip abgeschlossen → {ergebnis.get('output_dir')}") else: db.update_job( job_id, status="failed", error=ergebnis.get("error", "unbekannter Fehler"), finished_at=db.utcnow(), ) db.add_log("error", "worker", f"Job {job_id}: {ergebnis.get('error', 'unbekannter Fehler')}") return ergebnis def _komprimiere(job_id: str, disc_type: str, rip_ergebnis: dict, einstellungen: dict) -> dict: """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). """ quellen = rip_ergebnis.get("files", []) final_dir = os.path.join(RIP_OUTPUT_DIR, disc_type, job_id) os.makedirs(final_dir, exist_ok=True) preset = einstellungen.get("transcodePreset") or DEFAULT_HB_PRESET original_behalten = einstellungen.get("keepOriginal", False) db.update_job(job_id, status="transcoding", progress=0) db.add_log( "info", "worker", f"Job {job_id}: Kompression gestartet ({len(quellen)} Datei(en), Preset '{preset}')", ) anzahl = max(1, len(quellen)) 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) db.update_job(job_id, progress=min(99, gesamt)) hb = run_handbrake(quelle, ziel, preset=preset, progress_cb=datei_fortschritt) if hb.get("status") != "success": return { "status": "error", "error": ( f"Kompression fehlgeschlagen bei {os.path.basename(quelle)}: " f"{hb.get('error')} — Roh-Datei bleibt in /app/temp erhalten" ), } raw_dir = os.path.dirname(quellen[0]) if quellen else None if raw_dir: if original_behalten: ziel_original = os.path.join(final_dir, "original") shutil.move(raw_dir, ziel_original) db.add_log("info", "worker", f"Job {job_id}: Original behalten unter {ziel_original}") else: shutil.rmtree(raw_dir, ignore_errors=True) return {"status": "success", "output_dir": final_dir}