diff --git a/README.md b/README.md index c1b19b3..19daf21 100644 --- a/README.md +++ b/README.md @@ -1,119 +1,101 @@ -# README — Rippy +# Rippy — die All-in-one Disc-Ripping-Maschine -> **Moderner Ripping-Daemon mit Metadaten-Preview, Jellyfin-Formatierung und JWT-Auth.** +Disc rein → automatisch erkannt (Titel, Poster, Metadaten) → verlustfrei +gerippt (MakeMKV) → auf Arbeitsgröße komprimiert (HandBrake) → fertig +abgelegt, wo DU willst (lokal, NAS, jede Freigabe). Modernes Web-UI, +Echtzeit-Fortschritt, komplett in Docker, komplett lokal. ---- +## Schnellstart -## Features - -| Feature | Status | -|---------|--------| -| Disc-Erkennung (CD/DVD/Blu-ray) | ✅ | -| Metadaten-Lookup (TMDB/MusicBrainz/TheTVDB) | ✅ | -| Pre-Scan ohne Ripping | ✅ | -| Jellyfin-Formatierung (NFO + Images) | ✅ | -| JWT-Auth + Rate-Limiting | ✅ | -| React-UI mit Dashboard | ✅ | -| SQLite-Cache für API-Rate-Limits | ✅ | -| Multi-Disc-Set-Handling | ✅ | - ---- - -## Architektur - -``` -┌─────────────────────────────────────────────┐ -│ Proxmox LXC (Debian 12) │ -│ │ -│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ -│ │ API │ │ Worker │ │ UI │ │ -│ │ FastAPI │◄─►│ Celery+ │ │ React+ │ │ -│ │ │ │ Redis │ │ Nginx │ │ -│ └──────────┘ └──────────┘ └──────────┘ │ -│ │ │ │ │ -│ ▼ ▼ ▼ │ -│ ┌──────────┐ ┌──────────┐ ┌──────────┐ │ -│ │ PostgreSQL│ │ udev │ │ Media │ │ -│ │ │ │ Daemon │ │ Store │ │ -│ └──────────┘ └──────────┘ └──────────┘ │ -└─────────────────────────────────────────────┘ -``` - ---- - -## Installation - -### Docker Compose +Voraussetzungen: Docker + Docker Compose, ein optisches Laufwerk am Host. ```bash -git clone https://git.tobisniceshomelab.ddnsfree.com/Hitonabi/rippy.git -cd rippy -docker compose up -d +git clone rippy && cd rippy +cp .env.example .env # JWT_SECRET_KEY eintragen (openssl rand -hex 32) +mkdir -p /srv/rippy/media # Ablage-Basis (anpassbar in docker-compose.yml) +docker compose up -d --build ``` -### Port-Übersicht +Dann `http://` öffnen — der **Einrichtungs-Assistent** startet beim +ersten Mal automatisch (API-Keys, Verarbeitung, erkannte Hardware). -| Service | Port | URL | -|---------|------|-----| -| UI | 80 | http://localhost:80 | -| API | 8000 | http://localhost:8000 | -| PostgreSQL | 5432 | localhost:5432 | -| Redis | 6379 | localhost:6379 | +### Laufwerk anpassen ---- +Standard ist `/dev/sr0` (+ `/dev/sg1` für den Worker — MakeMKV spricht +Laufwerke über die SCSI-Generic-Schicht an). Andere Geräte? Lege eine +`docker-compose.override.yml` an: -## API +```yaml +services: + api: + devices: ["/dev/sr1:/dev/sr0"] + worker: + devices: ["/dev/sr1:/dev/sr0", "/dev/sg2:/dev/sg1"] +``` -### Endpoints +Welche sg-Nummer dein Laufwerk hat, verrät `lsscsi -g` oder +`ls -la /sys/class/scsi_generic/`. -| Endpoint | Method | Description | -|----------|--------|-------------| -| `/health` | GET | Health check | -| `/jobs` | GET | Alle Jobs | -| `/devices` | GET | Alle Geräte | -| `/prescan` | POST | Pre-Scan durchführen | -| `/metadata/lookup` | POST | Metadaten lookup | -| `/metadata/confirm` | POST | Metadaten bestätigen | -| `/jellyfin/format` | POST | Für Jellyfin formatieren | -| `/token` | POST | Login (JWT) | -| `/token/refresh` | POST | Refresh Token | -| `/api-keys` | POST/GET/DELETE | API-Key Management | +**Laufwerk in einer VM?** Per USB-Passthrough anhand der Vendor-ID +durchreichen (Proxmox: `qm set -usb0 host=xxxx:yyyy,usb3=1`) — +NICHT als emuliertes CD-ROM (`media=cdrom`), das kann keine SCSI-Kommandos. ---- +## Wie es funktioniert -## Tech-Stack +1. **Disc-Wache** (ioctl-Polling, kein udev-Gefrickel) erkennt Einlegen, + identifiziert die Disc (Volume-Label → TMDB → OMDb-Fallback) und zeigt + sie mit Poster auf dem Dashboard. +2. **Rip** (MakeMKV, verlustfrei — der einzige Weg durch AACS): Ziel wählst + du beim Start (Filme/Serien/Musik/eigener Pfad, inkl. Netzwerk-Ziele). + Audio-CDs laufen über abcde → FLAC + MusicBrainz. +3. **Kompression** (HandBrake, eigener Job auf eigener Queue): x265/x264, + Preset im UI wählbar; Rohdatei wird erst nach Erfolg gelöscht + („Original behalten" als Option). Fehlgeschlagene Kompressionen lassen + sich ohne Neu-Rip neu anstoßen. +4. **4K-UHD**: braucht ein LibreDrive-fähiges Laufwerk (MakeMKV-Forum: + „Ultimate UHD Drives Flashing Guide"). Normale BD/DVD gehen mit jedem + Laufwerk. -| Layer | Tech | -|-------|------| -| Backend | Python, FastAPI, Celery | -| Database | PostgreSQL, SQLite, Redis | -| Frontend | React, Vite, TailwindCSS | -| Ripping | MakeMKV, abcde, FFmpeg, HandBrake | -| Auth | JWT, OAuth2 | +## Speicherziele (NAS, Freigaben) ---- +Unter **Einstellungen → Speicherziele** hängst du NFS- oder SMB-Freigaben +direkt aus dem UI ein — sie erscheinen sofort in der Ziel-Auswahl beim +Rippen und werden beim Start automatisch wieder verbunden. +Technik: der api-Container läuft mit `CAP_SYS_ADMIN` und einem +rshared-Bind auf `/srv/rippy/media`, Mounts propagieren zu allen +Containern. ⚠️ Zugangsdaten liegen unverschlüsselt in der lokalen +Postgres-DB — bewusster Heimnetz-Kompromiss; lege fürs NAS einen eigenen, +eingeschränkten Benutzer an. -## Roadmap +## Verarbeitung & Hardware -- [x] Etappe 1: Container-Infrastruktur -- [x] Etappe 2: Ripping-Pipeline -- [x] Etappe 3: Metadaten-Lookup + Pre-Scan -- [x] Etappe 4: Jellyfin-Formatierung -- [x] Etappe 5: API + Auth + WebUI -- [x] Etappe 6: Sicherheit + Compliance -- [x] Etappe 7: API UI Modernisiert -- [x] Etappe 8: Dark Mode & Separation of Concerns -- [ ] Etappe 9: Proxmox-Integration +**Einstellungen → Verarbeitung** zeigt ehrlich an, welche Encoder deine +Worker WIRKLICH haben (CPU x264/x265, VAAPI bei AMD/Intel-GPU, NVENC bei +NVIDIA) — jeder Worker meldet seine Fähigkeiten selbst beim Start. ---- +### Optional: GPU-Maschine im Netz als Transcode-Worker -## License +Die Kompression läuft als eigener Celery-Task auf der Queue `transcode` — +JEDE Maschine im Netz kann sie übernehmen (siehe +`deploy/remote-transcode-worker.yml`). Ohne Zusatz-Worker macht der +eingebaute CPU-Worker alles selbst — Rippy bleibt All-in-one. -GPL-v3 +## Umgebungsvariablen (.env) ---- +| Variable | Pflicht | Zweck | +|---|---|---| +| `JWT_SECRET_KEY` | ✔ | Signierschlüssel (openssl rand -hex 32) | +| `TMDB_API_KEY` | empfohlen | Metadaten — alternativ im UI/Wizard eintragbar | +| `OMDB_API_KEY` | optional | zweite Metadaten-Quelle (Fallback) | +| `THETVDB_API_KEY` | optional | Serien-Fallback | +| `MAKEMKV_APP_KEY` | optional | MakeMKV-Beta-Key (Forum); DVDs gehen ohne | +| `MAKEMKV_URL_BASE` | optional | alternative Download-Quelle für den Image-Build | -## Credits +UI-Einstellungen (Wizard/Settings) überstimmen die Env-Variablen. -- **Commander**: Projekt-Idee, Anforderungen, Testing -- **AI-Box**: Entwicklung, Architektur, Dokumentation +## Entwicklung + +CI („Ampel") läuft bei jedem Push: Ruff, pytest, Vite-Build. Grün auf +`main` wird automatisch auf `stable` befördert — deploye von `stable`. +Regeln für Beiträge: [AGENTS.md](AGENTS.md) · Konzept: [KONZEPT.md](KONZEPT.md) · +Fahrplan: [ROADMAP.md](ROADMAP.md) diff --git a/deploy/remote-transcode-worker.yml b/deploy/remote-transcode-worker.yml new file mode 100644 index 0000000..91083ec --- /dev/null +++ b/deploy/remote-transcode-worker.yml @@ -0,0 +1,35 @@ +# Optionaler Remote-Transcode-Worker — für eine GPU-Maschine im Netz. +# (EXPERIMENTELL, 23.07.2026: VAAPI/NVENC-Presets folgen; aktuell nutzt der +# Worker dieselben HandBrake-CPU-Presets, bringt also v.a. stärkere CPUs.) +# +# Auf der GPU-Maschine: +# 1. Dieses Repo klonen (oder nur dieses File + Zugriff aufs Registry-Image) +# 2. Rohdaten-Freigabe der Rippy-Maschine mounten, z. B.: +# mount -t nfs :/srv/rippy /mnt/rippy +# (die Rippy-VM muss /srv/rippy + das temp-Volume exportieren) +# 3. RIPPY_HOST unten setzen und starten: +# docker compose -f deploy/remote-transcode-worker.yml up -d --build +# +# Der Worker meldet seine Encoder-Fähigkeiten automatisch — er taucht danach +# unter Einstellungen → Verarbeitung auf. Er bedient NUR die transcode-Queue; +# gerippt wird weiterhin dort, wo das Laufwerk hängt. + +services: + transcode-worker: + build: + context: .. + dockerfile: docker/worker/Dockerfile + command: ["celery", "-A", "celery_app", "worker", "--loglevel=info", "-Q", "transcode", "-n", "gpu-worker@%h"] + environment: + - REDIS_URL=redis://${RIPPY_HOST:?RIPPY_HOST setzen}:6379/0 + - DATABASE_URL=postgresql://rippy:rippy@${RIPPY_HOST}:5432/rippy + - RIP_OUTPUT_DIR=/app/media + - RAW_DIR=/app/temp/raw + volumes: + # Freigabe der Rippy-Maschine (siehe Kopf-Kommentar) + - /mnt/rippy/media:/app/media + - /mnt/rippy/temp:/app/temp + devices: + # GPU für Hardware-Encoding (AMD/Intel: /dev/dri; NVIDIA: nvidia-runtime) + - /dev/dri:/dev/dri + restart: unless-stopped diff --git a/docker-compose.yml b/docker-compose.yml index 2b0e20c..9fe23cd 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -26,15 +26,22 @@ services: start_period: 10s ports: - "8000:8000" + # SYS_ADMIN: die API hängt Netzwerk-Speicherziele (NFS/SMB) selbst ein + # (mounts.py) — dank rshared-Propagation unten sehen Host UND Worker + # jeden Mount sofort. + cap_add: + - SYS_ADMIN + security_opt: + - apparmor:unconfined volumes: - # /srv/rippy/media auf der VM statt anonymem Volume: dort eingehängte - # Shares (NFS/SMB) tauchen dank rslave-Propagation LIVE als Rip-Ziele - # in Rippy auf (Commander-Anforderung 23.07.: Ziel frei wählbar). + # /srv/rippy/media auf der VM statt anonymem Volume: rshared lässt + # Mounts aus DIESEM Container zum Host und in andere Container + # propagieren (Commander-Anforderung 23.07.: Ziele frei wählbar). - type: bind source: /srv/rippy/media target: /app/media bind: - propagation: rslave + propagation: rshared - temp:/app/temp devices: - /dev/sr0:/dev/sr0 diff --git a/docker/api/Dockerfile b/docker/api/Dockerfile index 98045c9..701deea 100644 --- a/docker/api/Dockerfile +++ b/docker/api/Dockerfile @@ -3,7 +3,13 @@ FROM python:3.12-slim-bookworm WORKDIR /app # udev ist raus (23.07.): udevadm lieferte im Container nie Daten (kein udevd) — -# die Geräte-Erkennung läuft jetzt über /sys + ioctls, ganz ohne Systempakete. +# die Geräte-Erkennung läuft jetzt über /sys + ioctls. +# nfs-common/cifs-utils: Netzwerk-Speicherziele werden aus dem UI heraus +# eingehängt (mounts.py, braucht CAP_SYS_ADMIN aus dem Compose). +RUN apt-get update && apt-get install -y --no-install-recommends \ + nfs-common \ + cifs-utils \ + && rm -rf /var/lib/apt/lists/* COPY docker/api/requirements.txt . RUN pip install --no-cache-dir -r requirements.txt diff --git a/docker/api/clients/omdb.py b/docker/api/clients/omdb.py index 48cfecc..e90f66d 100644 --- a/docker/api/clients/omdb.py +++ b/docker/api/clients/omdb.py @@ -32,7 +32,9 @@ def parse_year(year: str) -> Optional[int]: class OMDbClient: def __init__(self): - self.api_key = settings.omdb_api_key + # DB-Einstellung (Settings-UI/Wizard) gewinnt gegen die Env-Variable + from db import get_settings + self.api_key = get_settings().get("omdbApiKey") or settings.omdb_api_key self.session = requests.Session() def lookup(self, title: str, year: Optional[int] = None) -> Optional[Dict]: diff --git a/docker/api/clients/thetvdb.py b/docker/api/clients/thetvdb.py index 3dade2c..7ef23bc 100644 --- a/docker/api/clients/thetvdb.py +++ b/docker/api/clients/thetvdb.py @@ -12,7 +12,9 @@ THETVDB_BASE_URL = "https://api.thetvdb.com" class TheTVDBClient: def __init__(self): - self.api_key = settings.thetvdb_api_key + # DB-Einstellung (Settings-UI/Wizard) gewinnt gegen die Env-Variable + from db import get_settings + self.api_key = get_settings().get("tvdbApiKey") or settings.thetvdb_api_key self.base_url = THETVDB_BASE_URL self.session = requests.Session() self.session.headers.update({ diff --git a/docker/api/clients/tmdb.py b/docker/api/clients/tmdb.py index 5401c51..9f3ee8b 100644 --- a/docker/api/clients/tmdb.py +++ b/docker/api/clients/tmdb.py @@ -13,7 +13,10 @@ TMDB_IMAGE_BASE_URL = "https://image.tmdb.org/t/p" class TMDBClient: def __init__(self): - self.api_key = settings.tmdb_api_key + # DB-Einstellung (Settings-UI/Wizard) gewinnt gegen die Env-Variable — + # vorher war das Settings-Feld reine Dekoration (Fix 23.07.). + from db import get_settings + self.api_key = get_settings().get("tmdbApiKey") or settings.tmdb_api_key self.session = requests.Session() self.session.headers.update({ "Authorization": f"Bearer {self.api_key}", diff --git a/docker/api/db.py b/docker/api/db.py index 87a7bdf..5c4f783 100644 --- a/docker/api/db.py +++ b/docker/api/db.py @@ -64,6 +64,63 @@ settings_table = Table( Column("value", Text), ) +workers = Table( + "workers", + metadata, + Column("name", String(128), primary_key=True), + Column("encoders", Text), + Column("last_seen", DateTime(timezone=True)), +) + +storage_mounts = Table( + "storage_mounts", + metadata, + Column("name", String(64), primary_key=True), + Column("typ", String(8)), # nfs | cifs + Column("quelle", String(255)), # host:/export bzw. //host/share + Column("optionen", String(255)), + Column("username", String(128)), + Column("passwort", String(255)), # Klartext — Heimnetz-Kompromiss, siehe README +) + + +def list_workers() -> list: + import json + + with engine.connect() as conn: + zeilen = conn.execute(select(workers)).mappings().all() + ergebnis = [] + for z in zeilen: + eintrag = dict(z) + try: + eintrag["encoders"] = json.loads(eintrag.get("encoders") or "[]") + except ValueError: + eintrag["encoders"] = [] + if eintrag.get("last_seen"): + eintrag["last_seen"] = eintrag["last_seen"].isoformat() + ergebnis.append(eintrag) + return ergebnis + + +def list_mounts() -> list: + with engine.connect() as conn: + return [dict(z) for z in conn.execute(select(storage_mounts)).mappings().all()] + + +def save_mount(name: str, typ: str, quelle: str, optionen: str, username: str, passwort: str) -> None: + with engine.begin() as conn: + conn.execute( + storage_mounts.insert().values( + name=name, typ=typ, quelle=quelle, + optionen=optionen, username=username, passwort=passwort, + ) + ) + + +def delete_mount(name: str) -> None: + with engine.begin() as conn: + conn.execute(storage_mounts.delete().where(storage_mounts.c.name == name)) + def utcnow() -> datetime: return datetime.now(timezone.utc) @@ -97,6 +154,14 @@ def insert_job( ) +def get_job(job_id: str) -> dict: + 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 list_jobs(limit: int = 100) -> list: with engine.connect() as conn: zeilen = conn.execute( diff --git a/docker/api/main.py b/docker/api/main.py index bf5a491..a96c032 100644 --- a/docker/api/main.py +++ b/docker/api/main.py @@ -12,7 +12,8 @@ import uuid import db import devices as device_discovery -from celery_client import start_rip +import mounts as mount_verwaltung +from celery_client import celery_client, start_rip from detection import CDS_DISC_OK, CDS_NO_DISC, CDS_TRAY_OPEN, drive_status from fastapi.security import OAuth2PasswordBearer @@ -60,6 +61,12 @@ async def startup_event(): except ConfigValidationError as e: print(f"⚠️ Konfigurations-Warnung: {e}") + # Gespeicherte Netzwerk-Speicherziele wiederherstellen + def remount(): + for meldung in mount_verwaltung.alle_remounten(): + db.add_log("info", "mounts", meldung) + await asyncio.to_thread(remount) + asyncio.create_task(disc_watcher()) @@ -316,6 +323,125 @@ async def eject_device(name: str): return {"status": "ejected", "device": device_path} +@app.post("/jobs/{job_id}/retry-transcode") +async def retry_transcode(job_id: str): + """Stößt die Kompression eines Jobs neu an — OHNE die Disc neu zu rippen. + + Voraussetzung: die Rohdateien liegen noch in /app/temp/raw/ + (bei Kompressions-Fehlschlägen bleiben sie dort absichtlich erhalten). + """ + job = await asyncio.to_thread(db.get_job, job_id) + if not job: + raise HTTPException(status_code=404, detail="Job nicht gefunden") + if job["status"] in ("running", "pending"): + raise HTTPException(status_code=409, detail="Job rippt noch") + + raw_dir = f"/app/temp/raw/{job_id}" + basis = job.get("target_dir") or f"{MEDIA_ROOT}/{job.get('disc_type') or 'bluray'}" + final_dir = f"{basis}/{job_id}" + + celery_client.send_task( + "worker.tasks.transcode_files", + args=[job_id, raw_dir, final_dir], + queue="transcode", + ) + await asyncio.to_thread(db.update_job, job_id, status="transcoding", progress=0, error=None) + await asyncio.to_thread(db.add_log, "info", "api", f"Job {job_id}: Kompression neu eingereiht") + return {"id": job_id, "status": "transcoding"} + + +@app.get("/capabilities") +async def capabilities(): + """Welche Encoder sind auf welchen Workern WIRKLICH verfügbar? + + Jeder Worker meldet sich beim Start selbst (caps.py) — auch optionale + Remote-GPU-Worker tauchen hier automatisch auf. + """ + return {"workers": await asyncio.to_thread(db.list_workers)} + + +class MountRequest(BaseModel): + name: str + type: str # nfs | cifs + source: str # host:/export bzw. //host/share + options: Optional[str] = None + username: Optional[str] = None + password: Optional[str] = None + + +@app.get("/storage-mounts") +async def get_storage_mounts(): + """Konfigurierte Netzwerk-Speicherziele inkl. Live-Mount-Status.""" + eintraege = await asyncio.to_thread(db.list_mounts) + return [ + { + "name": e["name"], + "type": e["typ"], + "source": e["quelle"], + "mounted": mount_verwaltung.ist_gemountet(e["name"]), + "has_credentials": bool(e.get("username")), + } + for e in eintraege + ] + + +@app.post("/storage-mounts", status_code=201) +async def create_storage_mount(request: MountRequest): + """Hängt ein NFS/SMB-Ziel ein und speichert es für den nächsten Start.""" + if not mount_verwaltung.validiere_name(request.name): + raise HTTPException(status_code=422, detail="Name: nur a-z, 0-9, Bindestrich (2-31 Zeichen)") + if request.type not in ("nfs", "cifs"): + raise HTTPException(status_code=422, detail="Typ muss nfs oder cifs sein") + if any(e["name"] == request.name for e in await asyncio.to_thread(db.list_mounts)): + raise HTTPException(status_code=409, detail="Name bereits vergeben") + + try: + await asyncio.to_thread( + mount_verwaltung.mounten, + request.name, request.type, request.source, + request.options or "", request.username or "", request.password or "", + ) + except RuntimeError as e: + raise HTTPException(status_code=400, detail=str(e)) + + await asyncio.to_thread( + db.save_mount, + request.name, request.type, request.source, + request.options or "", request.username or "", request.password or "", + ) + await asyncio.to_thread( + db.add_log, "success", "mounts", + f"Speicherziel '{request.name}' ({request.type}) eingehängt: {request.source}", + ) + return {"name": request.name, "mounted": True} + + +@app.delete("/storage-mounts/{name}") +async def delete_storage_mount(name: str): + """Hängt ein Netzwerk-Speicherziel aus und entfernt es aus der Konfiguration.""" + try: + await asyncio.to_thread(mount_verwaltung.aushaengen, name) + except RuntimeError as e: + raise HTTPException(status_code=400, detail=str(e)) + await asyncio.to_thread(db.delete_mount, name) + await asyncio.to_thread(db.add_log, "info", "mounts", f"Speicherziel '{name}' entfernt") + return {"status": "removed"} + + +@app.get("/setup") +async def setup_status(): + """First-Run-Erkennung: wurde der Einrichtungs-Assistent abgeschlossen?""" + einstellungen = await asyncio.to_thread(db.get_settings, "setup") + return {"done": bool(einstellungen.get("done"))} + + +@app.post("/setup/complete") +async def setup_complete(): + await asyncio.to_thread(db.save_settings, {"done": True}, "setup") + await asyncio.to_thread(db.add_log, "success", "setup", "Einrichtungs-Assistent abgeschlossen") + return {"done": True} + + @app.get("/logs") async def get_logs(limit: int = 200): """Echte Ereignisse aus der Datenbank (Watcher, API, Worker).""" diff --git a/docker/api/mounts.py b/docker/api/mounts.py new file mode 100644 index 0000000..b74e007 --- /dev/null +++ b/docker/api/mounts.py @@ -0,0 +1,109 @@ +"""Netzwerk-Speicherziele (NFS/SMB) direkt aus Rippy heraus einhängen. + +Commander-Anforderung 23.07.: Data-Mounts über das UI, universell für jede +Installation. Der api-Container läuft dafür mit CAP_SYS_ADMIN und einem +rshared-Bind auf /app/media — ein Mount hier propagiert über den Host in +alle anderen Container (Worker sieht das Ziel sofort). + +Zugangsdaten liegen in Postgres (Klartext — bewusster Heimnetz-Kompromiss, +im README dokumentiert; SMB-Credentials wandern NIE in die Kommandozeile, +sondern über eine temporäre credentials-Datei). +""" + +import os +import re +import subprocess +import tempfile + +import db + +MEDIA_ROOT = "/app/media" +NAME_MUSTER = re.compile(r"^[a-z0-9][a-z0-9-]{1,30}$") + + +def _mountpoint(name: str) -> str: + return os.path.join(MEDIA_ROOT, name) + + +def validiere_name(name: str) -> bool: + return bool(NAME_MUSTER.match(name)) + + +def ist_gemountet(name: str) -> bool: + return os.path.ismount(_mountpoint(name)) + + +def mounten(name: str, typ: str, quelle: str, optionen: str = "", + username: str = "", passwort: str = "") -> None: + """Hängt ein NFS/CIFS-Ziel unter /app/media/ ein. Wirft RuntimeError.""" + ziel = _mountpoint(name) + os.makedirs(ziel, exist_ok=True) + if os.path.ismount(ziel): + return + + creds_datei = None + try: + if typ == "nfs": + opts = optionen or "vers=4,soft,timeo=100" + cmd = ["mount", "-t", "nfs", "-o", opts, quelle, ziel] + elif typ == "cifs": + teile = [optionen] if optionen else [] + if username: + # Credentials über Datei statt Kommandozeile (ps-sichtbar!) + creds = tempfile.NamedTemporaryFile( + "w", delete=False, prefix="cifs-", suffix=".cred" + ) + creds.write(f"username={username}\npassword={passwort or ''}\n") + creds.close() + os.chmod(creds.name, 0o600) + creds_datei = creds.name + teile.append(f"credentials={creds.name}") + else: + teile.append("guest") + teile.append("iocharset=utf8") + cmd = ["mount", "-t", "cifs", "-o", ",".join(teile), quelle, ziel] + else: + raise RuntimeError(f"Unbekannter Typ: {typ} (nfs oder cifs)") + + 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]}") + finally: + if creds_datei: + try: + os.unlink(creds_datei) + except OSError: + pass + + +def aushaengen(name: str) -> None: + ziel = _mountpoint(name) + if os.path.ismount(ziel): + ergebnis = subprocess.run( + ["umount", ziel], capture_output=True, text=True, timeout=30 + ) + if ergebnis.returncode != 0: + raise RuntimeError( + f"umount schlug fehl: {(ergebnis.stderr or '').strip()[:300]}" + ) + try: + os.rmdir(ziel) + except OSError: + pass # nicht leer oder weg — egal + + +def alle_remounten() -> list: + """Beim API-Start: alle gespeicherten Mounts wiederherstellen.""" + meldungen = [] + for eintrag in db.list_mounts(): + try: + mounten( + eintrag["name"], eintrag["typ"], eintrag["quelle"], + eintrag.get("optionen") or "", eintrag.get("username") or "", + eintrag.get("passwort") or "", + ) + meldungen.append(f"{eintrag['name']}: eingehängt") + except Exception as e: + meldungen.append(f"{eintrag['name']}: FEHLER — {e}") + return meldungen diff --git a/docker/ui/src/App.tsx b/docker/ui/src/App.tsx index 1e05aa5..587e466 100644 --- a/docker/ui/src/App.tsx +++ b/docker/ui/src/App.tsx @@ -1,17 +1,30 @@ -import { useState } from 'react' +import { useState, useEffect } from 'react' import { LayoutDashboard, Disc, Settings, Sun, Moon, Terminal } from 'lucide-react' +import { api } from './lib/api' import { useDarkMode } from './context/ThemeContext' import Dashboard from './pages/Dashboard' import MetadataPreview from './pages/MetadataPreview' import SettingsPage from './pages/Settings' import LogsPage from './pages/Logs' +import FirstRunWizard from './components/FirstRunWizard' type Page = 'dashboard' | 'metadata' | 'settings' | 'logs' function App() { const [page, setPage] = useState('dashboard') + const [setupDone, setSetupDone] = useState(null) const { theme, toggleTheme } = useDarkMode() + useEffect(() => { + api.get('/setup') + .then(r => setSetupDone(!!r.data.done)) + .catch(() => setSetupDone(true)) // API nicht erreichbar → kein Wizard-Zwang + }, []) + + if (setupDone === false) { + return setSetupDone(true)} /> + } + const navItems = [ { id: 'dashboard', label: 'Dashboard', icon: LayoutDashboard }, { id: 'metadata', label: 'Metadaten', icon: Disc }, diff --git a/docker/ui/src/components/FirstRunWizard.tsx b/docker/ui/src/components/FirstRunWizard.tsx new file mode 100644 index 0000000..31a1fb4 --- /dev/null +++ b/docker/ui/src/components/FirstRunWizard.tsx @@ -0,0 +1,154 @@ +import { useState, useEffect } from 'react' +import { Disc, Cpu, Globe, CheckCircle, ArrowRight } from 'lucide-react' +import { api } from '../lib/api' +import { useDarkMode } from '../context/ThemeContext' + +// Erster Start: Rippy universell einrichten (Commander-Ziel: All-in-one, +// weitergebbar — keine Handarbeit in .env-Dateien nötig). + +interface WorkerInfo { + name: string + encoders: string[] + last_seen?: string +} + +const ENCODER_LABELS: Record = { + 'cpu-x264': 'CPU · H.264 (x264)', + 'cpu-x265': 'CPU · H.265 (x265)', + 'vaapi': 'Hardware · VAAPI (AMD/Intel)', + 'nvenc': 'Hardware · NVENC (NVIDIA)', +} + +export default function FirstRunWizard({ onDone }: { onDone: () => void }) { + const [tmdbKey, setTmdbKey] = useState('') + const [omdbKey, setOmdbKey] = useState('') + const [transcodeEnabled, setTranscodeEnabled] = useState(true) + const [preset, setPreset] = useState('H.265 MKV 1080p30') + const [workers, setWorkers] = useState([]) + const [saving, setSaving] = useState(false) + const [error, setError] = useState(null) + const { theme } = useDarkMode() + + useEffect(() => { + api.get('/capabilities') + .then(r => setWorkers(r.data.workers || [])) + .catch(() => setWorkers([])) + }, []) + + const abschliessen = async () => { + setSaving(true) + setError(null) + try { + const bestehende = (await api.get('/settings')).data || {} + await api.post('/settings', { + ...bestehende, + ...(tmdbKey ? { tmdbApiKey: tmdbKey } : {}), + ...(omdbKey ? { omdbApiKey: omdbKey } : {}), + transcodeEnabled, + transcodePreset: preset, + }) + await api.post('/setup/complete') + onDone() + } catch { + setError('Speichern fehlgeschlagen — läuft die API?') + } finally { + setSaving(false) + } + } + + const feldKlasse = `w-full px-3 py-2 border rounded-lg focus:ring-2 focus:ring-indigo-500 focus:border-transparent ${theme === 'dark' ? 'bg-slate-800 border-slate-600 text-slate-200' : 'bg-white border-slate-300 text-slate-900'}` + + return ( +
+
+
+
+ +

Willkommen bei Rippy

+
+

Einmalige Einrichtung — dauert keine zwei Minuten. Alles ist später unter Einstellungen änderbar.

+
+ +
+ {/* Metadaten-APIs */} +
+

+ Metadaten-Erkennung +

+
+
+ + setTmdbKey(e.target.value)} className={feldKlasse} placeholder="Kostenlos auf themoviedb.org" /> +
+
+ + setOmdbKey(e.target.value)} className={feldKlasse} placeholder="Kostenlos auf omdbapi.com" /> +
+

+ Ohne Keys rippt Rippy trotzdem — Discs heißen dann nur wie ihr Volume-Label. +

+
+
+ + {/* Verarbeitung + erkannte Hardware */} +
+

+ Verarbeitung +

+ +
+

Erkannte Encoder in deinem System:

+ {workers.length === 0 ? ( +

Noch kein Worker gemeldet — startet gerade?

+ ) : ( + workers.map(w => ( +
+ {w.name}: + {w.encoders.map(e => ( + + {ENCODER_LABELS[e] || e} + + ))} +
+ )) + )} +

+ Tipp: GPU-Maschinen im Netz können als optionale Transcode-Worker angebunden werden (siehe README). +

+
+ + + + {transcodeEnabled && ( + + )} +
+ + {error &&

{error}

} + + + +

+ Läuft komplett lokal — keine Cloud, keine Telemetrie. +

+
+
+
+ ) +} diff --git a/docker/ui/src/components/StorageMounts.tsx b/docker/ui/src/components/StorageMounts.tsx new file mode 100644 index 0000000..9c352f8 --- /dev/null +++ b/docker/ui/src/components/StorageMounts.tsx @@ -0,0 +1,171 @@ +import { useState, useEffect } from 'react' +import { HardDrive, Plus, Trash2, CheckCircle, AlertCircle, RefreshCw } from 'lucide-react' +import { api } from '../lib/api' +import { useDarkMode } from '../context/ThemeContext' + +// Netzwerk-Speicherziele (NAS/PC-Freigaben) direkt aus dem UI einhängen — +// jedes Ziel erscheint danach in der Ziel-Auswahl beim Rippen. + +interface Mount { + name: string + type: string + source: string + mounted: boolean + has_credentials: boolean +} + +interface Target { + name: string + path: string + is_mount: boolean + free_gb: number | null +} + +export default function StorageMounts() { + const [mounts, setMounts] = useState([]) + const [targets, setTargets] = useState([]) + const [name, setName] = useState('') + const [typ, setTyp] = useState<'nfs' | 'cifs'>('nfs') + const [source, setSource] = useState('') + const [username, setUsername] = useState('') + const [password, setPassword] = useState('') + const [busy, setBusy] = useState(false) + const [feedback, setFeedback] = useState(null) + const { theme } = useDarkMode() + + const laden = async () => { + try { + const [m, t] = await Promise.all([ + api.get('/storage-mounts'), + api.get('/storage-targets'), + ]) + setMounts(m.data) + setTargets(t.data) + } catch { + /* API weg — Anzeige bleibt leer */ + } + } + + useEffect(() => { laden() }, []) + + const hinzufuegen = async () => { + setBusy(true) + setFeedback(null) + try { + await api.post('/storage-mounts', { + name, type: typ, source, + ...(username ? { username, password } : {}), + }) + setFeedback(`✓ „${name}" eingehängt — ab sofort als Rip-Ziel wählbar`) + setName(''); setSource(''); setUsername(''); setPassword('') + laden() + } catch (e: any) { + setFeedback(`✗ ${e?.response?.data?.detail || 'Einhängen fehlgeschlagen'}`) + } finally { + setBusy(false) + } + } + + const entfernen = async (mountName: string) => { + setBusy(true) + try { + await api.delete(`/storage-mounts/${mountName}`) + laden() + } catch (e: any) { + setFeedback(`✗ ${e?.response?.data?.detail || 'Entfernen fehlgeschlagen'}`) + } finally { + setBusy(false) + } + } + + const feld = `w-full px-3 py-2 border rounded-lg focus:ring-2 focus:ring-indigo-500 focus:border-transparent ${theme === 'dark' ? 'bg-slate-800 border-slate-600 text-slate-200' : 'bg-white border-slate-300 text-slate-900'}` + + return ( +
+ {/* Vorhandene Ziele */} +
+
+

Verfügbare Ziele

+ +
+ {targets.length === 0 ? ( +

Noch keine Unterordner in /app/media — sie entstehen beim ersten Rip oder Mount.

+ ) : ( +
+ {targets.map(t => ( +
+
+ +
+ {t.name} + + {t.is_mount ? 'Netzwerk-Mount' : 'lokal'}{t.free_gb != null ? ` · ${t.free_gb} GB frei` : ''} + +
+
+ {mounts.some(m => m.name === t.name) && ( + + )} +
+ ))} +
+ )} +
+ + {/* Neues Netzwerk-Ziel */} +
+

+ Netzwerk-Ziel hinzufügen +

+
+
+ + setName(e.target.value.toLowerCase())} className={feld} placeholder="nas-filme" /> +
+
+ + +
+
+ + setSource(e.target.value)} className={feld} placeholder={typ === 'nfs' ? '192.168.1.10:/volume1/filme' : '//192.168.1.10/filme'} /> +
+ {typ === 'cifs' && ( + <> +
+ + setUsername(e.target.value)} className={feld} autoComplete="off" /> +
+
+ + setPassword(e.target.value)} className={feld} autoComplete="new-password" /> +
+ + )} +
+ + {feedback && ( +

+ {feedback.startsWith('✓') ? : } + {feedback} +

+ )} +
+
+ ) +} diff --git a/docker/ui/src/pages/Dashboard.tsx b/docker/ui/src/pages/Dashboard.tsx index 8f89ed2..5804d8d 100644 --- a/docker/ui/src/pages/Dashboard.tsx +++ b/docker/ui/src/pages/Dashboard.tsx @@ -219,12 +219,13 @@ export default function Dashboard() { Gerät Fortschritt Startzeit + Aktion {jobs.length === 0 ? ( - +
Noch keine Jobs

Füge eine CD, DVD oder Blu-ray hinzu, um zu starten

@@ -252,6 +253,17 @@ export default function Dashboard() { {new Date(job.startTime).toLocaleString('de-DE')} + + {job.status === 'failed' && ( + + )} + )) )} diff --git a/docker/ui/src/pages/Settings.tsx b/docker/ui/src/pages/Settings.tsx index d457dda..087cd44 100644 --- a/docker/ui/src/pages/Settings.tsx +++ b/docker/ui/src/pages/Settings.tsx @@ -1,7 +1,8 @@ import { useState, useEffect } from 'react' -import { Save, Disc, Cpu, Database, Globe, AlertCircle, CheckCircle } from 'lucide-react' +import { Save, Disc, Cpu, Database, Globe, AlertCircle, CheckCircle, HardDrive } from 'lucide-react' import { api } from '../lib/api' import { useDarkMode } from '../context/ThemeContext' +import StorageMounts from '../components/StorageMounts' interface SettingsState { tmdbApiKey: string @@ -37,16 +38,35 @@ const defaultSettings: SettingsState = { keepOriginal: false, } -type SettingsTab = 'ripping' | 'verarbeitung' | 'verzeichnisse' | 'apis' | 'benachrichtigungen' +type SettingsTab = 'ripping' | 'verarbeitung' | 'speicherziele' | 'verzeichnisse' | 'apis' | 'benachrichtigungen' + +interface WorkerInfo { + name: string + encoders: string[] +} + +const ENCODER_LABELS: Record = { + 'cpu-x264': 'CPU · H.264', + 'cpu-x265': 'CPU · H.265', + 'vaapi': 'Hardware · VAAPI (AMD/Intel)', + 'nvenc': 'Hardware · NVENC (NVIDIA)', +} export default function SettingsPage() { const [settings, setSettings] = useState(defaultSettings) const [activeTab, setActiveTab] = useState('ripping') + const [workers, setWorkers] = useState([]) const [loading, setLoading] = useState(true) const [saving, setSaving] = useState(false) const [saved, setSaved] = useState(false) const { theme } = useDarkMode() + useEffect(() => { + api.get('/capabilities') + .then(r => setWorkers(r.data.workers || [])) + .catch(() => setWorkers([])) + }, []) + useEffect(() => { const loadSettings = async () => { try { @@ -105,6 +125,7 @@ export default function SettingsPage() {
setActiveTab('ripping')} darkMode={theme === 'dark'} /> setActiveTab('verarbeitung')} darkMode={theme === 'dark'} /> + setActiveTab('speicherziele')} darkMode={theme === 'dark'} /> setActiveTab('verzeichnisse')} darkMode={theme === 'dark'} /> setActiveTab('apis')} darkMode={theme === 'dark'} /> setActiveTab('benachrichtigungen')} darkMode={theme === 'dark'} /> @@ -184,6 +205,25 @@ export default function SettingsPage() { komprimiert danach auf Arbeitsgröße. Die Roh-Datei wird nach Erfolg gelöscht.

+ {/* Ehrliche Selbstauskunft: was können die Worker WIRKLICH? */} +
+

Verfügbare Encoder

+ {workers.length === 0 ? ( +

Kein Worker gemeldet.

+ ) : ( + workers.map(w => ( +
+ {w.name} + {w.encoders.map(e => ( + + {ENCODER_LABELS[e] || e} + + ))} +
+ )) + )} +
+
@@ -219,6 +259,15 @@ export default function SettingsPage() {
} + {/* Speicherziele: Netzwerk-Mounts aus dem UI */} + {activeTab === 'speicherziele' &&
+

+ + Speicherziele +

+ +
} + {/* Output Directories */} {activeTab === 'verzeichnisse' &&

diff --git a/docker/worker/Dockerfile b/docker/worker/Dockerfile index 689d32d..ac4d240 100644 --- a/docker/worker/Dockerfile +++ b/docker/worker/Dockerfile @@ -85,4 +85,6 @@ COPY docker/worker/ . RUN chmod +x /app/entrypoint.sh ENTRYPOINT ["/app/entrypoint.sh"] -CMD ["celery", "-A", "celery_app", "worker", "--loglevel=info"] +# Der lokale Worker bedient BEIDE Queues (rip + transcode). Ein optionaler +# Remote-GPU-Worker startet dasselbe Image nur mit "-Q transcode". +CMD ["celery", "-A", "celery_app", "worker", "--loglevel=info", "-Q", "celery,transcode"] diff --git a/docker/worker/caps.py b/docker/worker/caps.py new file mode 100644 index 0000000..b005023 --- /dev/null +++ b/docker/worker/caps.py @@ -0,0 +1,24 @@ +"""Encoder-Fähigkeiten dieses Workers — ehrlich erkannt, nicht behauptet. + +Jeder Worker meldet beim Start, was auf SEINER Maschine wirklich verfügbar +ist (Commander-Anforderung 23.07.: das UI zeigt an, WAS da ist). Ein +Remote-GPU-Worker meldet sich hier genauso wie der eingebaute CPU-Worker. +""" + +import os +import shutil + + +def erkenne_encoder() -> list: + """Liste der verfügbaren Encoder-Backends auf dieser Maschine.""" + gefunden = ["cpu-x264", "cpu-x265"] # HandBrake-Software-Encoder, immer dabei + + # VAAPI: AMD (VCN) und Intel (QuickSync) melden sich über /dev/dri + if os.path.exists("/dev/dri/renderD128"): + gefunden.append("vaapi") + + # NVENC: NVIDIA-Treiber im Container sichtbar + if shutil.which("nvidia-smi") or os.path.exists("/usr/lib/x86_64-linux-gnu/libnvidia-encode.so.1"): + gefunden.append("nvenc") + + return gefunden diff --git a/docker/worker/celery_app.py b/docker/worker/celery_app.py index c47998e..0823233 100644 --- a/docker/worker/celery_app.py +++ b/docker/worker/celery_app.py @@ -1,5 +1,8 @@ -from celery import Celery import os +import socket + +from celery import Celery +from celery.signals import worker_ready REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0") @@ -16,4 +19,18 @@ celery_app.conf.update( result_serializer="json", timezone="UTC", enable_utc=True, + broker_connection_retry_on_startup=True, ) + + +@worker_ready.connect +def melde_faehigkeiten(**kwargs): + """Beim Start: eigene Encoder-Fähigkeiten in die DB melden (UI-Anzeige).""" + import caps + import db + + try: + db.init_db() + db.save_worker(socket.gethostname(), caps.erkenne_encoder()) + except Exception as e: # DB noch nicht da → nicht den Start verhindern + print(f"Fähigkeiten-Meldung fehlgeschlagen: {e}") diff --git a/docker/worker/db.py b/docker/worker/db.py index 4792bb2..e3767cd 100644 --- a/docker/worker/db.py +++ b/docker/worker/db.py @@ -68,6 +68,35 @@ settings_table = Table( Column("value", Text), ) +workers = Table( + "workers", + metadata, + Column("name", String(128), primary_key=True), + Column("encoders", Text), + Column("last_seen", DateTime(timezone=True)), +) + + +def save_worker(name: str, encoder_liste: list) -> None: + """Worker meldet Name + Encoder-Fähigkeiten (Upsert).""" + import json + + payload = json.dumps(encoder_liste) + with engine.begin() as conn: + vorhanden = conn.execute( + workers.select().where(workers.c.name == name) + ).first() + if vorhanden: + conn.execute( + workers.update().where(workers.c.name == name).values( + encoders=payload, last_seen=utcnow() + ) + ) + else: + conn.execute( + workers.insert().values(name=name, encoders=payload, last_seen=utcnow()) + ) + def get_settings(key: str = "ui") -> dict: """UI-Einstellungen lesen (der Worker respektiert Transcode-Optionen).""" diff --git a/docker/worker/tasks.py b/docker/worker/tasks.py index 163606f..1810e3d 100644 --- a/docker/worker/tasks.py +++ b/docker/worker/tasks.py @@ -1,8 +1,10 @@ -"""Zentraler Rip-Task: erkennt den Disc-Typ, rippt, komprimiert, schreibt Status. +"""Rip- und Transcode-Tasks — bewusst GETRENNT (23.07.2026, Basis für Etappe 12/20). -Die API legt beim POST /jobs die Job-Zeile an und schickt diesen Task los — -der Worker hält die Zeile aktuell (running → transcoding → completed/failed) -und schreibt Ereignisse ins Log. Das UI liest beides über die API. +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 @@ -10,6 +12,7 @@ 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 os import shutil @@ -25,8 +28,6 @@ from ripping import ( ) RAW_DIR = os.getenv("RAW_DIR", "/app/temp/raw") - - MEDIA_ROOT = "/app/media" @@ -39,9 +40,30 @@ def _zielbasis(target_dir, disc_type: str) -> str: return os.path.join(RIP_OUTPUT_DIR, disc_type) +def _job_abschliessen(job_id: str, ergebnis: dict) -> None: + """Schreibt den Endzustand eines Jobs (completed/failed) nach Postgres.""" + if ergebnis.get("status") == "success": + db.update_job( + job_id, + status="completed", + progress=100, + output_path=ergebnis.get("output_dir"), + finished_at=db.utcnow(), + ) + db.add_log("success", "worker", f"Job {job_id}: abgeschlossen → {ergebnis.get('output_dir')}") + else: + db.update_job( + job_id, + status="failed", + error=ergebnis.get("error", "unbekannter Fehler"), + finished_at=db.utcnow(), + ) + db.add_log("error", "worker", f"Job {job_id}: {ergebnis.get('error', 'unbekannter Fehler')}") + + @celery_app.task(bind=True, name="worker.tasks.rip_disc") def rip_disc(self, device_path: str, job_id: str, target_dir: str = None): - """Rippt eine Disc basierend auf ihrem Typ; job_id ist die DB-Zeile der API. + """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. @@ -54,9 +76,7 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None): "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.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} @@ -87,7 +107,7 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None): if disc_type == "cd": ergebnis = rip_cd(device_path, job_id, progress_cb=fortschritt, output_dir=final_dir) elif transcode_an: - # Stufe 1: Roh-Rip nach /app/temp (wird nach der Kompression gelöscht) + # Roh-Rip nach /app/temp (wird nach erfolgreicher Kompression gelöscht) ergebnis = rip_video( device_path, job_id, disc_type, progress_cb=fortschritt, @@ -100,48 +120,48 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None): ) if ergebnis.get("status") == "success" and transcode_an: - ergebnis = _komprimiere(job_id, final_dir, ergebnis, einstellungen) - - if ergebnis.get("status") == "success": - db.update_job( - job_id, - status="completed", - progress=100, - output_path=ergebnis.get("output_dir"), - finished_at=db.utcnow(), + # 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, os.path.join(RAW_DIR, job_id), final_dir], + queue="transcode", ) - db.add_log("success", "worker", f"Job {job_id}: Rip abgeschlossen → {ergebnis.get('output_dir')}") - else: - db.update_job( - job_id, - status="failed", - error=ergebnis.get("error", "unbekannter Fehler"), - finished_at=db.utcnow(), - ) - db.add_log("error", "worker", f"Job {job_id}: {ergebnis.get('error', 'unbekannter Fehler')}") + return {"status": "ripped", "raw_dir": os.path.join(RAW_DIR, job_id)} + _job_abschliessen(job_id, ergebnis) return ergebnis -def _komprimiere(job_id: str, final_dir: str, rip_ergebnis: dict, einstellungen: dict) -> dict: +@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. """ - quellen = rip_ergebnis.get("files", []) - os.makedirs(final_dir, exist_ok=True) + 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) - db.update_job(job_id, status="transcoding", progress=0) + 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 = max(1, len(quellen)) + anzahl = len(quellen) for index, quelle in enumerate(quellen): ziel = os.path.join(final_dir, os.path.basename(quelle)) @@ -151,21 +171,23 @@ def _komprimiere(job_id: str, final_dir: str, rip_ergebnis: dict, einstellungen: hb = run_handbrake(quelle, ziel, preset=preset, progress_cb=datei_fortschritt) if hb.get("status") != "success": - return { + 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 - raw_dir = os.path.dirname(quellen[0]) if quellen else None - if raw_dir: - if original_behalten: - ziel_original = os.path.join(final_dir, "original") - shutil.move(raw_dir, ziel_original) - db.add_log("info", "worker", f"Job {job_id}: Original behalten unter {ziel_original}") - else: - shutil.rmtree(raw_dir, ignore_errors=True) + if original_behalten: + ziel_original = os.path.join(final_dir, "original") + shutil.move(raw_dir, ziel_original) + db.add_log("info", "worker", f"Job {job_id}: Original behalten unter {ziel_original}") + else: + shutil.rmtree(raw_dir, ignore_errors=True) - return {"status": "success", "output_dir": final_dir} + ergebnis = {"status": "success", "output_dir": final_dir} + _job_abschliessen(job_id, ergebnis) + return ergebnis