feat(backend): Media-Server-Aufbereitung, echte Webhooks, SMB-Klartext, UHD-Arbeitsverzeichnis, MakeMKV-Key via UI

- 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>
This commit is contained in:
Hitonabi
2026-07-24 08:55:22 +02:00
parent 048562a6ba
commit e032a9df1c
14 changed files with 787 additions and 36 deletions
+15 -1
View File
@@ -43,6 +43,7 @@ jobs = Table(
Column("output_path", Text),
Column("target_dir", String(255)),
Column("error", Text),
Column("meta", Text), # Disc-Metadaten (JSON: Jahr/Poster/Plot) für Detail-Popup + NFO
Column("created_at", DateTime(timezone=True)),
Column("finished_at", DateTime(timezone=True)),
)
@@ -69,6 +70,7 @@ workers = Table(
metadata,
Column("name", String(128), primary_key=True),
Column("encoders", Text),
Column("info", Text), # Werkzeug-Versionen (JSON: makemkv/handbrake/key-Quelle)
Column("last_seen", DateTime(timezone=True)),
)
@@ -96,6 +98,10 @@ def list_workers() -> list:
eintrag["encoders"] = json.loads(eintrag.get("encoders") or "[]")
except ValueError:
eintrag["encoders"] = []
try:
eintrag["info"] = json.loads(eintrag.get("info") or "{}")
except ValueError:
eintrag["info"] = {}
if eintrag.get("last_seen"):
eintrag["last_seen"] = eintrag["last_seen"].isoformat()
ergebnis.append(eintrag)
@@ -134,10 +140,17 @@ def init_db() -> None:
conn.exec_driver_sql(
"ALTER TABLE jobs ADD COLUMN IF NOT EXISTS target_dir VARCHAR(255)"
)
conn.exec_driver_sql(
"ALTER TABLE jobs ADD COLUMN IF NOT EXISTS meta TEXT"
)
conn.exec_driver_sql(
"ALTER TABLE workers ADD COLUMN IF NOT EXISTS info TEXT"
)
def insert_job(
job_id: str, device: str, disc_type: str = None, title: str = None, target_dir: str = None
job_id: str, device: str, disc_type: str = None, title: str = None,
target_dir: str = None, meta: str = None,
) -> None:
with engine.begin() as conn:
conn.execute(
@@ -147,6 +160,7 @@ def insert_job(
disc_type=disc_type,
title=title,
target_dir=target_dir,
meta=meta,
status="pending",
progress=0,
created_at=utcnow(),
+109 -13
View File
@@ -13,6 +13,7 @@ import uuid
import db
import devices as device_discovery
import mounts as mount_verwaltung
import notify
from celery_client import celery_client, start_rip
from detection import CDS_DISC_OK, CDS_NO_DISC, CDS_TRAY_OPEN, drive_status
@@ -262,15 +263,23 @@ async def create_job(request: JobCreateRequest):
raise HTTPException(status_code=404, detail=f"Laufwerk {device_path} nicht gefunden")
ziel = _validiere_ziel(request.target_dir)
# Titel aus der Disc-Erkennung übernehmen — sonst steht im Job "Unbekannt"
# Titel + Metadaten aus der Disc-Erkennung übernehmen — der Worker nutzt
# sie für den Ordnernamen und die Media-Server-Aufbereitung (NFO/Poster),
# das UI fürs Job-Detail-Popup.
titel = request.title
if not titel:
meta_json = None
disc = DISC_CACHE.get(device_path)
if disc and not disc.get("_laeuft"):
if not titel:
titel = disc.get("title")
meta_json = json.dumps({
"year": disc.get("year"),
"confidence": disc.get("confidence"),
**(disc.get("metadata") or {}),
})
job_id = str(uuid.uuid4())
await asyncio.to_thread(db.insert_job, job_id, device_path, None, titel, ziel)
await asyncio.to_thread(db.insert_job, job_id, device_path, None, titel, ziel, meta_json)
await asyncio.to_thread(
db.add_log, "info", "api",
f"Job {job_id} angelegt für {device_path}" + (f"{ziel}" if ziel else ""),
@@ -279,6 +288,23 @@ async def create_job(request: JobCreateRequest):
return {"id": job_id, "status": "pending", "device": device_path, "target_dir": ziel}
@app.get("/jobs/{job_id}/detail")
async def get_job_detail(job_id: str):
"""Alles zu EINEM Job — fürs Klick-Popup auf den Titel in „Neueste Jobs":
Metadaten (Poster/Jahr/Beschreibung), Ziel, Ausgabepfad, Fehler."""
job = await asyncio.to_thread(db.get_job, job_id)
if not job:
raise HTTPException(status_code=404, detail="Job nicht gefunden")
detail = _job_row_to_model(job).dict()
detail["target_dir"] = job.get("target_dir")
detail["output_path"] = job.get("output_path")
try:
detail["meta"] = json.loads(job["meta"]) if job.get("meta") else None
except ValueError:
detail["meta"] = None
return detail
@app.get("/storage-targets")
async def storage_targets():
"""Verfügbare Ablageziele: Verzeichnisse unter /app/media inkl. Mounts.
@@ -343,9 +369,16 @@ async def retry_transcode(job_id: str):
if job["status"] in ("running", "pending"):
raise HTTPException(status_code=409, detail="Job rippt noch")
raw_dir = f"/app/temp/raw/{job_id}"
# Roh-Verzeichnis: respektiert das konfigurierbare Arbeitsverzeichnis
# (Einstellungen → Verarbeitung), sonst Container-Default /app/temp/raw.
einstellungen = await asyncio.to_thread(db.get_settings)
work_dir = os.path.normpath((einstellungen.get("workDir") or "").strip() or "/")
raw_basis = work_dir if work_dir.startswith(MEDIA_ROOT) else "/app/temp/raw"
raw_dir = f"{raw_basis}/{job_id}"
# Zielordner: der Worker schreibt das geplante Ziel beim Rip-Start nach
# output_path (sprechender Name statt UUID) — alter Fallback bleibt.
basis = job.get("target_dir") or f"{MEDIA_ROOT}/{job.get('disc_type') or 'bluray'}"
final_dir = f"{basis}/{job_id}"
final_dir = job.get("output_path") or f"{basis}/{job_id}"
celery_client.send_task(
"worker.tasks.transcode_files",
@@ -495,17 +528,28 @@ async def browse(path: str = MEDIA_ROOT):
eintraege = sorted(os.listdir(normalisiert))
except OSError:
return None
return [
{"name": name, "path": os.path.join(normalisiert, name)}
for name in eintraege
if os.path.isdir(os.path.join(normalisiert, name))
]
ordner, dateien = [], []
for name in eintraege:
voll = os.path.join(normalisiert, name)
if os.path.isdir(voll):
ordner.append({"name": name, "path": voll})
else:
# Dateien MIT anzeigen (Befund 24.07.: der Browser wirkte
# „leer", weil er nur Ordner listete — die MKVs im
# bluray-Ordner waren unsichtbar).
try:
groesse_mb = round(os.path.getsize(voll) / 1024**2, 1)
except OSError:
groesse_mb = None
dateien.append({"name": name, "size_mb": groesse_mb})
return ordner, dateien
ordner = await asyncio.to_thread(liste)
if ordner is None:
ergebnis = await asyncio.to_thread(liste)
if ergebnis is None:
raise HTTPException(status_code=404, detail="Ordner nicht lesbar")
ordner, dateien = ergebnis
eltern = os.path.dirname(normalisiert) if normalisiert != MEDIA_ROOT else None
return {"path": normalisiert, "parent": eltern, "dirs": ordner}
return {"path": normalisiert, "parent": eltern, "dirs": ordner, "files": dateien}
class MkdirRequest(BaseModel):
@@ -530,6 +574,58 @@ async def browse_mkdir(request: MkdirRequest):
return {"path": ziel}
@app.get("/system/info")
async def system_info():
"""System-Selbstauskunft (Einstellungen → System): Werkzeug-Versionen der
Worker, freier Platz auf Media- und Arbeits-Volume, MakeMKV-Key-Status."""
def sammle():
info = {"api_version": app.version, "plaetze": [], "workers": db.list_workers()}
for name, pfad in (("Media (/app/media)", MEDIA_ROOT),
("Arbeitsverzeichnis (/app/temp)", "/app/temp")):
try:
nutzung = shutil.disk_usage(pfad)
info["plaetze"].append({
"name": name,
"frei_gb": round(nutzung.free / 1024**3, 1),
"gesamt_gb": round(nutzung.total / 1024**3, 1),
})
except OSError:
pass
einstellungen = db.get_settings()
info["makemkv_key_ui"] = bool((einstellungen.get("makemkvAppKey") or "").strip())
info["webhook_gesetzt"] = bool((einstellungen.get("notificationWebhook") or "").strip())
return info
return await asyncio.to_thread(sammle)
class NotificationTestRequest(BaseModel):
url: str
@app.post("/notifications/test")
async def notification_test(request: NotificationTestRequest):
"""Test-Nachricht an den Webhook — beweist die Anbindung SOFORT statt
erst beim ersten Job-Ende."""
url = request.url.strip()
if not url.startswith(("http://", "https://")):
raise HTTPException(status_code=422, detail="Webhook-URL muss mit http(s):// beginnen")
try:
await asyncio.to_thread(
notify.sende, url,
"🔔 Rippy: Test-Benachrichtigung",
"Wenn du das liest, funktioniert die Anbindung. Rippy meldet sich "
"hier, sobald ein Job fertig ist oder fehlschlägt.",
"info",
)
except RuntimeError as e:
raise HTTPException(status_code=400, detail=str(e))
await asyncio.to_thread(
db.add_log, "info", "notify", f"Test-Benachrichtigung gesendet ({notify.erkenne_webhook_typ(url)})"
)
return {"status": "sent", "typ": notify.erkenne_webhook_typ(url)}
@app.get("/setup")
async def setup_status():
"""First-Run-Erkennung: wurde der Einrichtungs-Assistent abgeschlossen?"""
+52 -2
View File
@@ -49,6 +49,40 @@ def schreibtest(pfad: str) -> bool:
return False
def uebersetze_smb_fehler(fehler: str, mit_credentials: bool) -> str:
"""Pure Funktion (testbar): NT_STATUS-Kauderwelsch → handelbarer Klartext.
Befund 24.07.: „session setup failed: NT_STATUS_ACCESS_DENIED" hieß in
Wahrheit nur „Windows/NAS erlauben keine Gast-Abfrage" — ohne Übersetzung
rät niemand, dass oben Benutzer/Passwort fehlen.
"""
if "NT_STATUS_ACCESS_DENIED" in fehler:
if not mit_credentials:
return (
"Zugriff verweigert — dieser Rechner erlaubt keine Gast-Abfrage. "
"Benutzername + Passwort eintragen (Windows: ein Konto mit "
"Zugriff auf die Freigabe), dann erneut auflisten."
)
return (
"Zugriff verweigert — Benutzername/Passwort stimmen nicht oder das "
"Konto darf die Freigaben nicht auflisten."
)
if "NT_STATUS_LOGON_FAILURE" in fehler:
return "Anmeldung fehlgeschlagen — Benutzername oder Passwort falsch."
if (
"NT_STATUS_HOST_UNREACHABLE" in fehler
or "NT_STATUS_IO_TIMEOUT" in fehler
or "NT_STATUS_UNSUCCESSFUL" in fehler
or "NT_STATUS_CONNECTION_REFUSED" in fehler
or "Connection to" in fehler
):
return (
"Rechner nicht erreichbar — Name/IP prüfen; läuft dort ein "
"SMB-Dienst (Port 445, Datei- und Druckerfreigabe aktiv)?"
)
return fehler[:200]
def liste_smb_freigaben(host: str, username: str = "", passwort: str = "") -> list:
"""Listet die SMB-Freigaben eines Rechners (smbclient -L) — damit man
seinen PC/NAS per Klick wählt statt //host/share zu raten."""
@@ -57,7 +91,13 @@ def liste_smb_freigaben(host: str, username: str = "", passwort: str = "") -> li
cmd += ["-U", f"{username}%{passwort or ''}"]
else:
cmd += ["-N"]
try:
ergebnis = subprocess.run(cmd, capture_output=True, text=True, timeout=20)
except subprocess.TimeoutExpired:
raise RuntimeError(
f"Freigaben-Abfrage fehlgeschlagen: {host} antwortet nicht "
"(Timeout nach 20 s) — Name/IP und Netzwerk prüfen."
)
freigaben = []
for zeile in (ergebnis.stdout or "").splitlines():
# -g-Format: Disk|Freigabename|Kommentar
@@ -66,7 +106,10 @@ def liste_smb_freigaben(host: str, username: str = "", passwort: str = "") -> li
freigaben.append(teile[1])
if not freigaben and ergebnis.returncode != 0:
fehler = (ergebnis.stderr or ergebnis.stdout or "").strip()
raise RuntimeError(f"Freigaben-Abfrage fehlgeschlagen: {fehler[:200]}")
raise RuntimeError(
"Freigaben-Abfrage fehlgeschlagen: "
+ uebersetze_smb_fehler(fehler, bool(username))
)
return freigaben
@@ -109,7 +152,14 @@ def mounten(name: str, typ: str, quelle: str, optionen: str = "",
ergebnis = subprocess.run(cmd, capture_output=True, text=True, timeout=30)
if ergebnis.returncode != 0:
fehler = (ergebnis.stderr or ergebnis.stdout or "").strip()
raise RuntimeError(f"mount schlug fehl: {fehler[:300]}")
hinweis = ""
if "error(13)" in fehler or "Permission denied" in fehler:
hinweis = (
" — Zugriff verweigert: Benutzer/Passwort prüfen; ohne "
"Angaben versucht Rippy einen Gast-Zugriff, den Windows/"
"NAS meist ablehnen."
)
raise RuntimeError(f"mount schlug fehl: {fehler[:300]}{hinweis}")
return schreibtest(ziel)
finally:
if creds_datei:
+66
View File
@@ -0,0 +1,66 @@
"""Webhook-Benachrichtigungen bei Job-Ende (Discord, Slack, ntfy, generisch).
Bis 24.07. war das notificationWebhook-Setting ein Placebo: das UI speicherte
die URL, aber NICHTS hat je gesendet. Jetzt meldet der Worker Job-Ende
(fertig/fehlgeschlagen/abgebrochen) und die API bietet einen Test-Endpoint.
Das Modul existiert bewusst identisch in API und Worker
(docker/api/notify.py) — es gibt kein geteiltes Paket zwischen den Containern.
Wer es ändert, ändert BEIDE Dateien.
Payload-Formate (dokumentiert, nicht geraten — AGENTS Regel D):
- Discord: POST JSON {"content": "..."} — discord.com/developers/docs/resources/webhook
- Slack: POST JSON {"text": "..."} — api.slack.com/messaging/webhooks
- ntfy: POST Roh-Text an https://ntfy.sh/<topic>, Titel via ?title= —
docs.ntfy.sh/publish
- Generisch: POST JSON {"title", "message", "level"} für eigene Empfänger
(Home Assistant, n8n, eigene Skripte).
"""
import requests
TIMEOUT_SEKUNDEN = 10
def erkenne_webhook_typ(url: str) -> str:
"""Erkennt den Dienst an der URL — der Nutzer muss nichts konfigurieren."""
u = url.lower()
if "discord.com/api/webhooks" in u or "discordapp.com/api/webhooks" in u:
return "discord"
if "hooks.slack.com" in u:
return "slack"
if "ntfy" in u:
return "ntfy"
return "generisch"
def baue_payload(url: str, titel: str, text: str, level: str = "info"):
"""Pure Funktion (testbar): (typ, json_payload) — ntfy sendet Roh-Text."""
typ = erkenne_webhook_typ(url)
if typ == "discord":
return typ, {"content": f"**{titel}**\n{text}"}
if typ == "slack":
return typ, {"text": f"*{titel}*\n{text}"}
if typ == "ntfy":
return typ, None
return typ, {"title": titel, "message": text, "level": level}
def sende(url: str, titel: str, text: str, level: str = "info") -> None:
"""Schickt die Nachricht; wirft RuntimeError mit Klartext bei Fehlern."""
typ, payload = baue_payload(url, titel, text, level)
try:
if typ == "ntfy":
antwort = requests.post(
url, data=text.encode("utf-8"),
params={"title": titel}, timeout=TIMEOUT_SEKUNDEN,
)
else:
antwort = requests.post(url, json=payload, timeout=TIMEOUT_SEKUNDEN)
except requests.RequestException as e:
raise RuntimeError(f"Webhook nicht erreichbar: {e}") from e
if antwort.status_code >= 300:
raise RuntimeError(
f"Webhook antwortete mit HTTP {antwort.status_code}: "
f"{(antwort.text or '')[:200]}"
)
+41
View File
@@ -0,0 +1,41 @@
"""Tests für die SMB-Fehlerübersetzung (Speicherziele → Freigaben auflisten)."""
from mounts import uebersetze_smb_fehler, validiere_name
def test_access_denied_ohne_credentials_erklaert_gastproblem():
meldung = uebersetze_smb_fehler(
"session setup failed: NT_STATUS_ACCESS_DENIED", mit_credentials=False
)
assert "Gast" in meldung
assert "Benutzername + Passwort" in meldung
def test_access_denied_mit_credentials_verweist_auf_konto():
meldung = uebersetze_smb_fehler(
"session setup failed: NT_STATUS_ACCESS_DENIED", mit_credentials=True
)
assert "stimmen nicht" in meldung
def test_logon_failure_wird_uebersetzt():
meldung = uebersetze_smb_fehler("NT_STATUS_LOGON_FAILURE", mit_credentials=True)
assert "Passwort falsch" in meldung
def test_unerreichbar_wird_uebersetzt():
meldung = uebersetze_smb_fehler(
"do_connect: Connection to 10.0.0.9 failed (Error NT_STATUS_IO_TIMEOUT)",
mit_credentials=False,
)
assert "nicht erreichbar" in meldung
def test_unbekannter_fehler_bleibt_erhalten_und_gekappt():
meldung = uebersetze_smb_fehler("X" * 500, mit_credentials=False)
assert meldung == "X" * 200
def test_validiere_name_bleibt_streng():
assert validiere_name("nas-filme")
assert not validiere_name("NAS Filme")
+5
View File
@@ -76,6 +76,11 @@ RUN apt-get update && apt-get install -y --no-install-recommends \
COPY --from=makemkv-build /usr/local /usr/local
RUN ldconfig
# Version fürs UI sichtbar machen (Einstellungen → System): makemkvcon hat
# keinen --version-Schalter, also kommt die Wahrheit aus dem Build selbst.
ARG MAKEMKV_VERSION=1.17.7
ENV MAKEMKV_VERSION=${MAKEMKV_VERSION}
WORKDIR /app
COPY docker/worker/requirements.txt .
+38
View File
@@ -6,7 +6,9 @@ Remote-GPU-Worker meldet sich hier genauso wie der eingebaute CPU-Worker.
"""
import os
import re
import shutil
import subprocess
def erkenne_encoder() -> list:
@@ -22,3 +24,39 @@ def erkenne_encoder() -> list:
gefunden.append("nvenc")
return gefunden
def werkzeug_versionen() -> dict:
"""Kern-Werkzeuge dieses Workers — fürs UI (Einstellungen → System).
MakeMKV-Version kommt aus dem Build (ENV MAKEMKV_VERSION im Dockerfile) —
makemkvcon hat keinen dokumentierten --version-Schalter (usage.txt).
HandBrakeCLI --version ist dokumentiert (handbrake.fr/docs, CLI Options)
und gibt z.B. "HandBrake 1.6.1" aus.
"""
info = {}
if shutil.which("makemkvcon"):
info["makemkv"] = os.getenv("MAKEMKV_VERSION") or "installiert"
if shutil.which("HandBrakeCLI"):
try:
aus = subprocess.run(
["HandBrakeCLI", "--version"],
capture_output=True, text=True, timeout=15,
)
treffer = re.search(r"HandBrake\s+([\w.]+)", (aus.stdout or "") + (aus.stderr or ""))
info["handbrake"] = treffer.group(1) if treffer else "installiert"
except (OSError, subprocess.TimeoutExpired):
info["handbrake"] = "installiert"
# Woher kommt der MakeMKV-Key? UI-Setting schlägt Env — ehrlich anzeigen.
try:
import db
ui_key = (db.get_settings().get("makemkvAppKey") or "").strip()
except Exception:
ui_key = ""
if ui_key:
info["makemkv_key"] = "ui"
elif os.getenv("MAKEMKV_APP_KEY"):
info["makemkv_key"] = "env"
else:
info["makemkv_key"] = "keiner"
return info
+5 -1
View File
@@ -43,7 +43,11 @@ def melde_faehigkeiten(**kwargs):
while True:
try:
db.init_db()
db.save_worker(socket.gethostname(), caps.erkenne_encoder())
db.save_worker(
socket.gethostname(),
caps.erkenne_encoder(),
caps.werkzeug_versionen(),
)
except Exception as e: # DB weg → weiterversuchen, nicht sterben
print(f"Fähigkeiten-Meldung fehlgeschlagen: {e}")
time.sleep(60)
+37 -5
View File
@@ -36,7 +36,9 @@ jobs = Table(
Column("status", String(16), nullable=False, server_default="pending"),
Column("progress", Integer, nullable=False, server_default="0"),
Column("output_path", Text),
Column("target_dir", String(255)),
Column("error", Text),
Column("meta", Text), # Disc-Metadaten (JSON) — Quelle für Ordnernamen + NFO
Column("created_at", DateTime(timezone=True)),
Column("finished_at", DateTime(timezone=True)),
)
@@ -57,8 +59,22 @@ def utcnow() -> datetime:
def init_db() -> None:
"""Legt fehlende Tabellen an (idempotent)."""
"""Legt fehlende Tabellen an (idempotent) und zieht Mini-Migrationen nach.
Die ALTERs stehen identisch in docker/api/db.py — wer zuerst startet,
migriert; der andere findet die Spalten dann bereits vor.
"""
metadata.create_all(engine)
with engine.begin() as conn:
conn.exec_driver_sql(
"ALTER TABLE jobs ADD COLUMN IF NOT EXISTS target_dir VARCHAR(255)"
)
conn.exec_driver_sql(
"ALTER TABLE jobs ADD COLUMN IF NOT EXISTS meta TEXT"
)
conn.exec_driver_sql(
"ALTER TABLE workers ADD COLUMN IF NOT EXISTS info TEXT"
)
settings_table = Table(
@@ -73,15 +89,17 @@ workers = Table(
metadata,
Column("name", String(128), primary_key=True),
Column("encoders", Text),
Column("info", Text), # Werkzeug-Versionen (JSON: makemkv/handbrake/key-Quelle)
Column("last_seen", DateTime(timezone=True)),
)
def save_worker(name: str, encoder_liste: list) -> None:
"""Worker meldet Name + Encoder-Fähigkeiten (Upsert)."""
def save_worker(name: str, encoder_liste: list, info: dict = None) -> None:
"""Worker meldet Name + Encoder-Fähigkeiten + Werkzeug-Versionen (Upsert)."""
import json
payload = json.dumps(encoder_liste)
info_payload = json.dumps(info or {})
with engine.begin() as conn:
vorhanden = conn.execute(
workers.select().where(workers.c.name == name)
@@ -89,12 +107,14 @@ def save_worker(name: str, encoder_liste: list) -> None:
if vorhanden:
conn.execute(
workers.update().where(workers.c.name == name).values(
encoders=payload, last_seen=utcnow()
encoders=payload, info=info_payload, last_seen=utcnow()
)
)
else:
conn.execute(
workers.insert().values(name=name, encoders=payload, last_seen=utcnow())
workers.insert().values(
name=name, encoders=payload, info=info_payload, last_seen=utcnow()
)
)
@@ -116,6 +136,18 @@ def get_settings(key: str = "ui") -> dict:
return {}
def get_job(job_id: str) -> dict:
"""Ganze Job-Zeile — der Worker braucht Titel + Metadaten für die
Ordner-Benennung und die Media-Server-Aufbereitung (NFO/Poster)."""
from sqlalchemy import select
with engine.connect() as conn:
zeile = conn.execute(
select(jobs).where(jobs.c.id == job_id)
).mappings().first()
return dict(zeile) if zeile else None
def get_job_status(job_id: str) -> str:
"""Nur der Status — der Worker prüft damit kooperative Abbruch-Anfragen."""
from sqlalchemy import select
+107
View File
@@ -0,0 +1,107 @@
"""Media-Server-Aufbereitung: sprechende Ordnernamen + NFO + Poster.
Commander-Wunsch 24.07.: Rippy soll sich auf das Zielsystem (Jellyfin, Plex,
Emby, Kodi) vorbereiten. Das Setting `mediaServer` steuert:
- ALLE Systeme: Zielordner heißt <Titel> (Jahr)" statt Job-UUID — daran
erkennen Jellyfin & Co. den Film (Namensschema laut jellyfin.org/docs/
general/server/media/movies bzw. support.plex.tv Naming and Organizing").
- jellyfin / emby / kodi: zusätzlich movie.nfo bzw. tvshow.nfo (Kodi-Schema,
kodi.wiki/view/NFO_files) + poster.jpg diese Server lesen NFO nativ.
- plex: NUR Benennung (Plex liest NFO ohne Zusatz-Agent nicht).
- none: Verhalten wie bisher (UUID-Ordner, keine Extras).
"""
import os
import re
from xml.sax.saxutils import escape
import requests
NFO_SERVER = ("jellyfin", "emby", "kodi")
# Zeichen, die in Ordnernamen auf SMB/NTFS/ext4 Ärger machen
_VERBOTEN = re.compile(r'[<>:"/\\|?*\x00-\x1f]')
def sicherer_name(titel: str, jahr=None) -> str:
"""Pure Funktion (testbar): Titel → Dateisystem-tauglicher Ordnername."""
name = _VERBOTEN.sub("", titel or "").strip().rstrip(".")
name = re.sub(r"\s+", " ", name)[:150].strip()
if not name:
return ""
if jahr:
name = f"{name} ({jahr})"
return name
def zielordner(basis: str, titel: str, jahr, job_id: str) -> str:
"""Zielordner unter `basis`: sprechender Name, UUID nur als Fallback.
Existiert der Ordner schon (Disc doppelt gerippt), wird die Job-ID
angehängt statt fremde Dateien zu überschreiben.
"""
name = sicherer_name(titel, jahr)
if not name:
return os.path.join(basis, job_id)
pfad = os.path.join(basis, name)
if os.path.exists(pfad):
pfad = os.path.join(basis, f"{name} [{job_id[:8]}]")
return pfad
def baue_nfo(meta: dict, titel: str, jahr=None) -> str:
"""Pure Funktion (testbar): minimales movie.nfo/tvshow.nfo im Kodi-Schema."""
ist_serie = (meta.get("type") == "tv")
wurzel = "tvshow" if ist_serie else "movie"
zeilen = [
'<?xml version="1.0" encoding="UTF-8" standalone="yes"?>',
f"<{wurzel}>",
f" <title>{escape(titel or '')}</title>",
]
if jahr:
zeilen.append(f" <year>{escape(str(jahr))}</year>")
if meta.get("overview"):
zeilen.append(f" <plot>{escape(meta['overview'])}</plot>")
for genre in meta.get("genres") or []:
zeilen.append(f" <genre>{escape(str(genre))}</genre>")
if meta.get("runtime"):
zeilen.append(f" <runtime>{escape(str(meta['runtime']))}</runtime>")
quelle = meta.get("source") or "Rippy"
zeilen.append(f" <details>Source: {escape(quelle)}</details>")
zeilen.append(f"</{wurzel}>")
return "\n".join(zeilen) + "\n"
def aufbereiten(ordner: str, media_server: str, titel: str, jahr, meta: dict) -> list:
"""Schreibt NFO + Poster in den fertigen Ordner (best effort).
Rückgabe: Liste von Meldungen fürs Log. Wirft NIE ein fehlendes Poster
darf einen gelungenen Rip nicht auf 'failed' drehen.
"""
meldungen = []
if media_server not in NFO_SERVER:
return meldungen
meta = meta or {}
try:
ist_serie = (meta.get("type") == "tv")
nfo_name = "tvshow.nfo" if ist_serie else "movie.nfo"
with open(os.path.join(ordner, nfo_name), "w", encoding="utf-8") as f:
f.write(baue_nfo(meta, titel, jahr))
meldungen.append(f"{nfo_name} geschrieben")
except OSError as e:
meldungen.append(f"NFO fehlgeschlagen: {e}")
poster = meta.get("poster_path") or ""
if poster:
if not poster.startswith("http"):
poster = f"https://image.tmdb.org/t/p/w500{poster}"
try:
antwort = requests.get(poster, timeout=15)
if antwort.status_code == 200 and antwort.content:
with open(os.path.join(ordner, "poster.jpg"), "wb") as f:
f.write(antwort.content)
meldungen.append("poster.jpg gespeichert")
except requests.RequestException as e:
meldungen.append(f"Poster fehlgeschlagen: {e}")
return meldungen
+66
View File
@@ -0,0 +1,66 @@
"""Webhook-Benachrichtigungen bei Job-Ende (Discord, Slack, ntfy, generisch).
Bis 24.07. war das notificationWebhook-Setting ein Placebo: das UI speicherte
die URL, aber NICHTS hat je gesendet. Jetzt meldet der Worker Job-Ende
(fertig/fehlgeschlagen/abgebrochen) und die API bietet einen Test-Endpoint.
Das Modul existiert bewusst identisch in API und Worker
(docker/api/notify.py) es gibt kein geteiltes Paket zwischen den Containern.
Wer es ändert, ändert BEIDE Dateien.
Payload-Formate (dokumentiert, nicht geraten AGENTS Regel D):
- Discord: POST JSON {"content": "..."} discord.com/developers/docs/resources/webhook
- Slack: POST JSON {"text": "..."} api.slack.com/messaging/webhooks
- ntfy: POST Roh-Text an https://ntfy.sh/<topic>, Titel via ?title=
docs.ntfy.sh/publish
- Generisch: POST JSON {"title", "message", "level"} für eigene Empfänger
(Home Assistant, n8n, eigene Skripte).
"""
import requests
TIMEOUT_SEKUNDEN = 10
def erkenne_webhook_typ(url: str) -> str:
"""Erkennt den Dienst an der URL — der Nutzer muss nichts konfigurieren."""
u = url.lower()
if "discord.com/api/webhooks" in u or "discordapp.com/api/webhooks" in u:
return "discord"
if "hooks.slack.com" in u:
return "slack"
if "ntfy" in u:
return "ntfy"
return "generisch"
def baue_payload(url: str, titel: str, text: str, level: str = "info"):
"""Pure Funktion (testbar): (typ, json_payload) — ntfy sendet Roh-Text."""
typ = erkenne_webhook_typ(url)
if typ == "discord":
return typ, {"content": f"**{titel}**\n{text}"}
if typ == "slack":
return typ, {"text": f"*{titel}*\n{text}"}
if typ == "ntfy":
return typ, None
return typ, {"title": titel, "message": text, "level": level}
def sende(url: str, titel: str, text: str, level: str = "info") -> None:
"""Schickt die Nachricht; wirft RuntimeError mit Klartext bei Fehlern."""
typ, payload = baue_payload(url, titel, text, level)
try:
if typ == "ntfy":
antwort = requests.post(
url, data=text.encode("utf-8"),
params={"title": titel}, timeout=TIMEOUT_SEKUNDEN,
)
else:
antwort = requests.post(url, json=payload, timeout=TIMEOUT_SEKUNDEN)
except requests.RequestException as e:
raise RuntimeError(f"Webhook nicht erreichbar: {e}") from e
if antwort.status_code >= 300:
raise RuntimeError(
f"Webhook antwortete mit HTTP {antwort.status_code}: "
f"{(antwort.text or '')[:200]}"
)
+148 -11
View File
@@ -13,12 +13,15 @@ komprimiert danach auf Arbeitsgröße. Die Rohdatei liegt nur temporär in
"""
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
from detection import detect_disc_type, disc_size_bytes
from ripping import (
DEFAULT_HB_PRESET,
RIP_OUTPUT_DIR,
@@ -41,37 +44,136 @@ def _zielbasis(target_dir, disc_type: str) -> str:
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."""
"""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=ergebnis.get("output_dir"),
output_path=ausgabe,
finished_at=db.utcnow(),
)
db.add_log("success", "worker", f"Job {job_id}: abgeschlossen → {ergebnis.get('output_dir')}")
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=ergebnis.get("error", "unbekannter Fehler"),
error=fehler,
finished_at=db.utcnow(),
)
db.add_log("error", "worker", f"Job {job_id}: {ergebnis.get('error', 'unbekannter Fehler')}")
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")
@@ -122,16 +224,51 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None):
and einstellungen.get("transcodeEnabled", True)
)
final_dir = os.path.join(_zielbasis(target_dir, disc_type), job_id)
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 nach /app/temp (wird nach erfolgreicher Kompression gelöscht)
# Roh-Rip ins Arbeitsverzeichnis (wird nach erfolgreicher Kompression gelöscht)
ergebnis = rip_video(
device_path, job_id, disc_type,
progress_cb=fortschritt,
output_dir=os.path.join(RAW_DIR, job_id),
output_dir=raw_dir,
)
else:
ergebnis = rip_video(
@@ -145,10 +282,10 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None):
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, os.path.join(RAW_DIR, job_id), final_dir],
args=[job_id, raw_dir, final_dir],
queue="transcode",
)
return {"status": "ripped", "raw_dir": os.path.join(RAW_DIR, job_id)}
return {"status": "ripped", "raw_dir": raw_dir}
_job_abschliessen(job_id, ergebnis)
return ergebnis
+57
View File
@@ -0,0 +1,57 @@
"""Tests für die Media-Server-Aufbereitung (Benennung + NFO)."""
import os
from medien import baue_nfo, sicherer_name, zielordner
def test_sicherer_name_entfernt_verbotene_zeichen():
assert sicherer_name('Alien: Die Wiedergeburt') == "Alien Die Wiedergeburt"
assert sicherer_name('Was/ist\\das?<>|*"') == "Wasistdas"
def test_sicherer_name_mit_jahr():
assert sicherer_name("Evangelion 2.22", 2009) == "Evangelion 2.22 (2009)"
def test_sicherer_name_kollabiert_leerraum_und_punkte():
assert sicherer_name(" Viel Raum ... ") == "Viel Raum"
def test_sicherer_name_leer_bleibt_leer():
assert sicherer_name("") == ""
assert sicherer_name("???") == ""
def test_zielordner_nutzt_titel_und_jahr(tmp_path):
pfad = zielordner(str(tmp_path), "Inception", 2010, "abc-123")
assert pfad == os.path.join(str(tmp_path), "Inception (2010)")
def test_zielordner_faellt_auf_job_id_zurueck(tmp_path):
pfad = zielordner(str(tmp_path), "", None, "abc-123")
assert pfad == os.path.join(str(tmp_path), "abc-123")
def test_zielordner_weicht_bei_kollision_aus(tmp_path):
os.makedirs(tmp_path / "Inception (2010)")
pfad = zielordner(str(tmp_path), "Inception", 2010, "abcdef12-3456")
assert pfad == os.path.join(str(tmp_path), "Inception (2010) [abcdef12]")
def test_baue_nfo_film_mit_plot_und_genres():
nfo = baue_nfo(
{"overview": "Ein Traum <im> Traum & so", "genres": ["Sci-Fi", "Action"]},
"Inception", 2010,
)
assert "<movie>" in nfo and "</movie>" in nfo
assert "<title>Inception</title>" in nfo
assert "<year>2010</year>" in nfo
# XML-Escaping: <, >, & dürfen den Parser nicht sprengen
assert "Ein Traum &lt;im&gt; Traum &amp; so" in nfo
assert "<genre>Sci-Fi</genre>" in nfo
def test_baue_nfo_serie_bekommt_tvshow_wurzel():
nfo = baue_nfo({"type": "tv"}, "Neon Genesis Evangelion", 1995)
assert "<tvshow>" in nfo and "</tvshow>" in nfo
+38
View File
@@ -0,0 +1,38 @@
"""Tests für die Webhook-Benachrichtigungen (Typ-Erkennung + Payload-Bau)."""
from notify import baue_payload, erkenne_webhook_typ
def test_discord_wird_an_url_erkannt():
assert erkenne_webhook_typ("https://discord.com/api/webhooks/1/abc") == "discord"
assert erkenne_webhook_typ("https://discordapp.com/api/webhooks/1/abc") == "discord"
def test_slack_und_ntfy_und_generisch():
assert erkenne_webhook_typ("https://hooks.slack.com/services/T/B/x") == "slack"
assert erkenne_webhook_typ("https://ntfy.sh/mein-thema") == "ntfy"
assert erkenne_webhook_typ("https://ha.local/api/webhook/rippy") == "generisch"
def test_discord_payload_nutzt_content():
typ, payload = baue_payload("https://discord.com/api/webhooks/1/a", "Titel", "Text")
assert typ == "discord"
assert payload == {"content": "**Titel**\nText"}
def test_slack_payload_nutzt_text():
typ, payload = baue_payload("https://hooks.slack.com/services/x", "Titel", "Text")
assert typ == "slack"
assert payload == {"text": "*Titel*\nText"}
def test_ntfy_sendet_rohtext():
typ, payload = baue_payload("https://ntfy.sh/thema", "Titel", "Text")
assert typ == "ntfy"
assert payload is None # Roh-Text-Body, kein JSON
def test_generischer_payload_traegt_level():
typ, payload = baue_payload("https://example.org/hook", "Titel", "Text", "error")
assert typ == "generisch"
assert payload == {"title": "Titel", "message": "Text", "level": "error"}