Files
rippy/docker/worker/tasks.py
T
Hitonabi 8eb5653848
Ampel / ampel (push) Successful in 55s
Restefeger: Auth komplett raus, Serien-Flow + Episoden-Matching, Jellyfin-Refresh, Duplikat-Warnung, echtes Nur-Hauptfilm
AUTH ENTFERNT (Commander-Entscheid 24.07., KONZEPT §10): /token- und
/api-keys-Endpoints, auth.py, test_auth.py, passlib/bcrypt/PyJWT/
python-multipart, JWT_SECRET_KEY-Pflicht. Heimnetz-only, das UI hatte nie
einen Login — die Auth-Oberflaeche war Placebo und die passlib/bcrypt-
Falle brach die Ampel. Rate-Limit pro IP bleibt. Schnellstart laeuft
jetzt ganz ohne .env-Pflichtwerte.

Serien-Flow (Etappe-12-Kern, ARM-Wunde #395):
- Rip-Dialog: Serienname + Staffel -> Ablage <Serie>/Season NN
  (jellyfin.org/docs Naming-Schema); tvshow.nfo + poster.jpg im
  Serien-Ordner, bei Staffel 2 nicht ueberschrieben.
- Episoden-Matching per Laufzeitabgleich: HandBrakeCLI --scan
  ('+ duration:', handbrake.fr/docs) je MKV gegen TMDB-Staffel-Laufzeiten
  (GET /metadata/tv/{id}/season/{n}; tv-season-details-API).
  Ordnungserhaltend; komplette Staffel auf einer Disc klappt auch bei
  uniformen Anime-Laufzeiten (Sequenz-Stufe). Umbenannt wird NUR bei
  eindeutiger Zuordnung — sonst ehrliches Log. Mit Tests.

Weitere Punkte:
- Jellyfin/Emby-Bibliotheks-Refresh nach jedem fertigen Rip
  (POST /Library/Refresh, X-Emby-Token lt. jellyfin.org/docs) —
  URL/Key + Test-Knopf in Einstellungen -> Ripping.
- Duplikat-Warnung: Disc-Fingerabdruck (jetzt Teil des Prescan-Ergebnisses
  + der Job-Metadaten) gegen die Historie; Karte zeigt 'bereits gerippt',
  Vollautomatik ueberspringt Duplikate.
- 'Nur Hauptfilm' ECHT: makemkvcon info -> TINFO-Attr-9-Laufzeiten
  (usage.txt) -> laengster Titel -> mkv dev:X <nr>. Vorher wirkungsloses
  Setting; pro Rip im Dialog uebersteuerbar. Mit Tests.
- OMDb-Treffer eingedeutscht via TMDB /find (external_source=imdb_id,
  de-DE; find-by-id-API).
- Dashboard: Speicherplatz-Anzeige (amber < 60 GB) + CSV-Export
  (GET /jobs/export, Semikolon+BOM fuer deutsches Excel).
- Metadaten-Seite entfernt (Abnahme durch Commander-Auftrag) inkl.
  Placebo-Endpoints /metadata/lookup (scannte Dummy-Device) und
  /metadata/confirm (schrieb nie gelesenen Cache-Key).
- Doppel-Jahr-Fix: 'X (2009) (2009)' in Log und Ordnernamen.
- Remote-Worker-Blocker: redis (6379) + postgres (5432) waren NIE
  veroeffentlicht — kein Remote-Worker konnte sich je verbinden. Ports
  jetzt offen (Heimnetz-Kompromiss, kommentiert) + API_URL fuer Worker.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-24 14:22:19 +02:00

466 lines
19 KiB
Python

"""Rip- und Transcode-Tasks — bewusst GETRENNT (23.07.2026, Basis für Etappe 12/20).
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).
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 glob
import json
import os
import shutil
import requests
import db
import medien
import notify
from celery_app import celery_app
from detection import detect_disc_type, disc_size_bytes
from ripping import (
DEFAULT_HB_PRESET,
RIP_OUTPUT_DIR,
RipAbbruch,
lies_datei_dauer,
rip_cd,
rip_video,
run_handbrake,
)
API_URL = os.getenv("API_URL", "http://api:8000")
RAW_DIR = os.getenv("RAW_DIR", "/app/temp/raw")
MEDIA_ROOT = "/app/media"
def _zielbasis(target_dir, disc_type: str) -> str:
"""Ablagebasis: vom Nutzer gewähltes Ziel (validiert) oder Standard."""
if target_dir:
normalisiert = os.path.normpath(target_dir)
if normalisiert.startswith(MEDIA_ROOT):
return normalisiert
return os.path.join(RIP_OUTPUT_DIR, disc_type)
def _arbeitsverzeichnis(einstellungen: dict) -> str:
"""Basis für Roh-Rips: UI-Setting `workDir` (unter /app/media, z. B. eine
NAS-Freigabe) schlägt den 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.
"""
work_dir = (einstellungen.get("workDir") or "").strip()
if work_dir:
normalisiert = os.path.normpath(work_dir)
if normalisiert.startswith(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.
"""
key = (einstellungen.get("makemkvAppKey") or "").strip()
if not key:
return
ordner = os.path.expanduser("~/.MakeMKV")
try:
os.makedirs(ordner, exist_ok=True)
with open(os.path.join(ordner, "settings.conf"), "w") as f:
f.write(f'app_Key = "{key}"\n')
except OSError as e:
db.add_log("warning", "worker", f"MakeMKV-Key konnte nicht gesetzt werden: {e}")
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()
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()
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.
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 _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)
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)
einstellungen = db.get_settings()
ist_video = disc_type in ("dvd", "bluray", "uhd")
transcode_an = ist_video and einstellungen.get("transcodeEnabled", True)
if ist_video:
# UI-Key schlägt Env-Key — Monats-Key ohne Rebuild aktualisierbar
_makemkv_key_anwenden(einstellungen)
# Sprechender Zielordner „<Titel> (Jahr)" statt Job-UUID — Jellyfin & Co.
# erkennen den Film am Ordnernamen. UUID bleibt Fallback ohne Titel.
# Serien-Rips landen stattdessen in <Serie>/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,
)
# Nur Hauptfilm (längster Titel): pro Rip wählbar, sonst Setting;
# bei Serien-Discs sinnlos (dort zählen ALLE Episoden-Titel).
nur_hauptfilm = meta.get("main_feature_only")
if nur_hauptfilm is None:
nur_hauptfilm = bool(einstellungen.get("mainFeatureOnly", False))
if meta.get("series"):
nur_hauptfilm = False
# Geplantes Ziel sofort sichtbar machen (UI-Detail + retry-transcode)
db.update_job(job_id, output_path=final_dir)
raw_dir = os.path.join(_arbeitsverzeichnis(einstellungen), 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)
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)
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,
)
else:
ergebnis = rip_video(
device_path, job_id, disc_type,
progress_cb=fortschritt, output_dir=final_dir,
nur_hauptfilm=nur_hauptfilm,
)
# UHD-Fehler in Klartext übersetzen — "Failed to open disc" allein hilft
# niemandem. Zwei bekannte Ursachen (Befunde 24.07., Summer-Wars-UHD):
if ergebnis.get("status") == "error" and disc_type == "uhd":
fehler_text = ergebnis.get("error") or ""
if "volume key is unknown" in fehler_text:
# LibreDrive lief bereits — MakeMKV kennt nur den Disc-Schlüssel
# nicht: Version zu alt ODER Disc neuer als die Key-Datenbank.
ergebnis["error"] += (
" — Klartext: Das Laufwerk liest die Disc (LibreDrive OK), "
"aber MakeMKV kennt den Schlüssel dieser Disc nicht. Erst "
"prüfen: MakeMKV aktuell? (Einstellungen → System; Update = "
"Image-Rebuild). Ist es aktuell, ist die Disc neuer als die "
"Schlüssel-Datenbank — MakeMKV hat einen AACS-Dump unter "
"/root/.MakeMKV/ im Worker gespeichert; im MakeMKV-Forum "
"(Bereich 'Ultra HD Blu-ray') einreichen, mit einem der "
"nächsten Updates ist die Disc dann rippbar."
)
elif "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."
)
if ergebnis.get("status") == "success" and transcode_an:
# Kompression als eigener Task auf der transcode-Queue — kann vom
# lokalen Worker ODER einem Remote-GPU-Worker übernommen werden.
db.update_job(job_id, status="transcoding", progress=0)
db.add_log("info", "worker", f"Job {job_id}: Rip fertig, Kompression eingereiht")
transcode_files.apply_async(
args=[job_id, raw_dir, final_dir],
queue="transcode",
)
return {"status": "ripped", "raw_dir": raw_dir}
_job_abschliessen(job_id, ergebnis)
return ergebnis
@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.
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()
quellen = sorted(glob.glob(os.path.join(raw_dir, "*.mkv")))
if not quellen:
ergebnis = {"status": "error", "error": f"Keine Roh-MKVs in {raw_dir} gefunden"}
_job_abschliessen(job_id, ergebnis)
return ergebnis
einstellungen = db.get_settings()
preset = einstellungen.get("transcodePreset") or DEFAULT_HB_PRESET
original_behalten = einstellungen.get("keepOriginal", False)
os.makedirs(final_dir, exist_ok=True)
db.update_job(job_id, status="transcoding", progress=0, error=None)
db.add_log(
"info", "worker",
f"Job {job_id}: Kompression gestartet ({len(quellen)} Datei(en), Preset '{preset}')",
)
anzahl = len(quellen)
letzter = [-1]
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
if _abbruch_angefordert(job_id):
raise RipAbbruch()
db.update_job(job_id, progress=min(99, gesamt))
hb = run_handbrake(quelle, ziel, preset=preset, progress_cb=datei_fortschritt)
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:
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)
ergebnis = {"status": "success", "output_dir": final_dir}
_job_abschliessen(job_id, ergebnis)
return ergebnis