e032a9df1c
- jobs.meta (Migration in beiden db.py): Disc-Metadaten wandern in den Job —
Quelle fuer Job-Detail-Popup (GET /jobs/{id}/detail), Ordner-Benennung, NFO.
- medien.py (Worker, mit Tests): Zielordner 'Titel (Jahr)' statt Job-UUID;
movie.nfo/tvshow.nfo (Kodi-Schema, kodi.wiki/view/NFO_files) + poster.jpg
fuer jellyfin/emby/kodi; plex nur Benennung; Kollision -> Job-ID-Suffix.
- notify.py (identisch in API+Worker, mit Tests): Webhook bei Job-Ende —
Discord ({content}, discord.com/developers), Slack ({text},
api.slack.com/messaging/webhooks), ntfy (Rohtext + ?title=, docs.ntfy.sh),
generisches JSON. Vorher war das Setting ein Placebo: nichts sendete je.
POST /notifications/test beweist die Anbindung sofort.
- mounts.py: NT_STATUS-Fehler -> handelbarer Klartext (ACCESS_DENIED ohne
Credentials = Gast-Abfrage verweigert), Timeout-Meldung, mount-Hinweis
bei error(13). Mit Tests.
- tasks.py: Platz-Check per Disc-Groesse (ioctl BLKGETSIZE64) VOR dem Rip;
Arbeitsverzeichnis workDir (unter /app/media, z.B. NAS) statt fix
/app/temp/raw — 4K-UHD-Rohdaten (bis 100 GB) sprengen sonst die VM-Platte;
MakeMKV-Key aus UI-Settings (~/.MakeMKV/settings.conf, Format wie
entrypoint.sh) gilt ab dem naechsten Rip ohne Rebuild.
- caps.py: Werkzeug-Versionen (MakeMKV aus ENV MAKEMKV_VERSION im Image,
makemkvcon hat keinen --version-Schalter lt. usage.txt; HandBrakeCLI
--version lt. handbrake.fr/docs) + Key-Quelle; workers.info-Spalte.
- /browse liefert jetzt auch DATEIEN (Name + Groesse) — der 'leere'
bluray-Ordner war voll, der Browser zeigte nur Unterordner.
- GET /system/info: Versionen, Plattenplatz, Key-/Webhook-Status.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
360 lines
14 KiB
Python
360 lines
14 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 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,
|
|
rip_cd,
|
|
rip_video,
|
|
run_handbrake,
|
|
)
|
|
|
|
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 _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")):
|
|
# Media-Server-Aufbereitung (Jellyfin/Emby/Kodi: NFO + Poster)
|
|
einstellungen = db.get_settings()
|
|
try:
|
|
meta = json.loads(job.get("meta") or "{}")
|
|
except ValueError:
|
|
meta = {}
|
|
for meldung in medien.aufbereiten(
|
|
ausgabe, einstellungen.get("mediaServer") or "none",
|
|
job.get("title") or "", meta.get("year"), meta,
|
|
):
|
|
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()
|
|
transcode_an = (
|
|
disc_type in ("dvd", "bluray")
|
|
and einstellungen.get("transcodeEnabled", True)
|
|
)
|
|
|
|
if disc_type in ("dvd", "bluray"):
|
|
# 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.
|
|
job = db.get_job(job_id) or {}
|
|
try:
|
|
meta = json.loads(job.get("meta") or "{}")
|
|
except ValueError:
|
|
meta = {}
|
|
final_dir = medien.zielordner(
|
|
_zielbasis(target_dir, disc_type),
|
|
job.get("title") or "", meta.get("year"), job_id,
|
|
)
|
|
# 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 disc_type in ("dvd", "bluray"):
|
|
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,
|
|
)
|
|
else:
|
|
ergebnis = rip_video(
|
|
device_path, job_id, disc_type,
|
|
progress_cb=fortschritt, output_dir=final_dir,
|
|
)
|
|
|
|
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
|