From e032a9df1c2351a5bab90a6221a13645000133d1 Mon Sep 17 00:00:00 2001 From: Hitonabi Date: Fri, 24 Jul 2026 08:55:22 +0200 Subject: [PATCH] feat(backend): Media-Server-Aufbereitung, echte Webhooks, SMB-Klartext, UHD-Arbeitsverzeichnis, MakeMKV-Key via UI MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 --- docker/api/db.py | 16 ++- docker/api/main.py | 126 ++++++++++++++++++++--- docker/api/mounts.py | 56 ++++++++++- docker/api/notify.py | 66 +++++++++++++ docker/api/test_mounts_helpers.py | 41 ++++++++ docker/worker/Dockerfile | 5 + docker/worker/caps.py | 38 +++++++ docker/worker/celery_app.py | 6 +- docker/worker/db.py | 42 +++++++- docker/worker/medien.py | 107 ++++++++++++++++++++ docker/worker/notify.py | 66 +++++++++++++ docker/worker/tasks.py | 159 +++++++++++++++++++++++++++--- docker/worker/test_medien.py | 57 +++++++++++ docker/worker/test_notify.py | 38 +++++++ 14 files changed, 787 insertions(+), 36 deletions(-) create mode 100644 docker/api/notify.py create mode 100644 docker/api/test_mounts_helpers.py create mode 100644 docker/worker/medien.py create mode 100644 docker/worker/notify.py create mode 100644 docker/worker/test_medien.py create mode 100644 docker/worker/test_notify.py diff --git a/docker/api/db.py b/docker/api/db.py index 30abab3..73fa035 100644 --- a/docker/api/db.py +++ b/docker/api/db.py @@ -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(), diff --git a/docker/api/main.py b/docker/api/main.py index 54166e5..2c6edfa 100644 --- a/docker/api/main.py +++ b/docker/api/main.py @@ -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: - disc = DISC_CACHE.get(device_path) - if disc and not disc.get("_laeuft"): + 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?""" diff --git a/docker/api/mounts.py b/docker/api/mounts.py index 24e7ec1..b8c5c8b 100644 --- a/docker/api/mounts.py +++ b/docker/api/mounts.py @@ -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"] - ergebnis = subprocess.run(cmd, capture_output=True, text=True, timeout=20) + 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: diff --git a/docker/api/notify.py b/docker/api/notify.py new file mode 100644 index 0000000..87e9d36 --- /dev/null +++ b/docker/api/notify.py @@ -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/, 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]}" + ) diff --git a/docker/api/test_mounts_helpers.py b/docker/api/test_mounts_helpers.py new file mode 100644 index 0000000..be06b5b --- /dev/null +++ b/docker/api/test_mounts_helpers.py @@ -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") diff --git a/docker/worker/Dockerfile b/docker/worker/Dockerfile index ac4d240..0615172 100644 --- a/docker/worker/Dockerfile +++ b/docker/worker/Dockerfile @@ -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 . diff --git a/docker/worker/caps.py b/docker/worker/caps.py index b005023..63cdda0 100644 --- a/docker/worker/caps.py +++ b/docker/worker/caps.py @@ -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 diff --git a/docker/worker/celery_app.py b/docker/worker/celery_app.py index 20cddab..e4160d7 100644 --- a/docker/worker/celery_app.py +++ b/docker/worker/celery_app.py @@ -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) diff --git a/docker/worker/db.py b/docker/worker/db.py index 64776ed..ad4b50b 100644 --- a/docker/worker/db.py +++ b/docker/worker/db.py @@ -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 diff --git a/docker/worker/medien.py b/docker/worker/medien.py new file mode 100644 index 0000000..400d3ce --- /dev/null +++ b/docker/worker/medien.py @@ -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 „ (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 = [ + '', + f"<{wurzel}>", + f" {escape(titel or '')}", + ] + if jahr: + zeilen.append(f" {escape(str(jahr))}") + if meta.get("overview"): + zeilen.append(f" {escape(meta['overview'])}") + for genre in meta.get("genres") or []: + zeilen.append(f" {escape(str(genre))}") + if meta.get("runtime"): + zeilen.append(f" {escape(str(meta['runtime']))}") + quelle = meta.get("source") or "Rippy" + zeilen.append(f"
Source: {escape(quelle)}
") + zeilen.append(f"") + 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 diff --git a/docker/worker/notify.py b/docker/worker/notify.py new file mode 100644 index 0000000..87e9d36 --- /dev/null +++ b/docker/worker/notify.py @@ -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/, 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]}" + ) diff --git a/docker/worker/tasks.py b/docker/worker/tasks.py index cbad244..3c0d605 100644 --- a/docker/worker/tasks.py +++ b/docker/worker/tasks.py @@ -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 „ (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 diff --git a/docker/worker/test_medien.py b/docker/worker/test_medien.py new file mode 100644 index 0000000..a9a84ce --- /dev/null +++ b/docker/worker/test_medien.py @@ -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 Traum & so", "genres": ["Sci-Fi", "Action"]}, + "Inception", 2010, + ) + assert "" in nfo and "" in nfo + assert "Inception" in nfo + assert "2010" in nfo + # XML-Escaping: <, >, & dürfen den Parser nicht sprengen + assert "Ein Traum <im> Traum & so" in nfo + assert "Sci-Fi" in nfo + + +def test_baue_nfo_serie_bekommt_tvshow_wurzel(): + nfo = baue_nfo({"type": "tv"}, "Neon Genesis Evangelion", 1995) + assert "" in nfo and "" in nfo diff --git a/docker/worker/test_notify.py b/docker/worker/test_notify.py new file mode 100644 index 0000000..8b29ff0 --- /dev/null +++ b/docker/worker/test_notify.py @@ -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"}