diff --git a/docker/api/clients/jikan.py b/docker/api/clients/jikan.py new file mode 100644 index 0000000..5669f8f --- /dev/null +++ b/docker/api/clients/jikan.py @@ -0,0 +1,82 @@ +"""Jikan API Client (MyAnimeList) — kostenlose Anime-Datenbank OHNE API-Key. + +Commander-Fund 23.07.: Für Anime-Discs (häufiger Fall im Heimlab) ist MAL die +präziseste Quelle — „Evangelion: 2.22" existiert dort exakt, mit Poster und +Beschreibung. Kein Key, Rate-Limit 3 req/s (für Disc-Erkennung irrelevant). +Doku: https://docs.api.jikan.moe/ +""" + +from difflib import SequenceMatcher +from typing import Dict, Optional + +import requests + +from cache import get, set as cache_set + +JIKAN_BASE_URL = "https://api.jikan.moe/v4" + + +def titel_aehnlichkeit(a: str, b: str) -> float: + """Normalisierte Titel-Ähnlichkeit 0..1 (pure Funktion, testbar).""" + def norm(t: str) -> str: + return "".join(c for c in t.lower() if c.isalnum() or c == " ").strip() + + a_n, b_n = norm(a), norm(b) + if not a_n or not b_n: + return 0.0 + return SequenceMatcher(None, a_n, b_n).ratio() + + +class JikanClient: + MINDEST_AEHNLICHKEIT = 0.55 + + def __init__(self): + self.session = requests.Session() + + def lookup(self, title: str) -> Optional[Dict]: + """Sucht per Titel; nimmt den ÄHNLICHSTEN Treffer, nicht blind den ersten.""" + cache_key = f"jikan:{title.lower()}" + cached = get(cache_key) + if cached: + return cached + + try: + response = self.session.get( + f"{JIKAN_BASE_URL}/anime", + params={"q": title, "limit": 5, "sfw": "true"}, + timeout=10, + ) + response.raise_for_status() + eintraege = response.json().get("data") or [] + except Exception as e: + print(f"Jikan API Error: {e}") + return None + + bester, bester_score = None, 0.0 + for eintrag in eintraege: + varianten = [eintrag.get("title") or ""] + if eintrag.get("title_english"): + varianten.append(eintrag["title_english"]) + score = max(titel_aehnlichkeit(title, v) for v in varianten) + if score > bester_score: + bester, bester_score = eintrag, score + + if not bester or bester_score < self.MINDEST_AEHNLICHKEIT: + return None + + bilder = (bester.get("images") or {}).get("jpg") or {} + ergebnis = { + "type": "movie" if bester.get("type") == "Movie" else "tv", + "id": f"mal-{bester.get('mal_id', '')}", + "title": bester.get("title_english") or bester.get("title") or title, + "year": bester.get("year"), + "overview": (bester.get("synopsis") or "")[:600], + "poster_path": bilder.get("large_image_url") or bilder.get("image_url") or "", + "backdrop_path": "", + "runtime": 0, + "genres": [g.get("name", "") for g in (bester.get("genres") or []) if g.get("name")], + "source": "jikan", + "match_score": round(bester_score, 2), + } + cache_set(cache_key, ergebnis) + return ergebnis diff --git a/docker/api/main.py b/docker/api/main.py index 1c539f8..d976138 100644 --- a/docker/api/main.py +++ b/docker/api/main.py @@ -262,8 +262,15 @@ 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 = request.title + if not titel: + disc = DISC_CACHE.get(device_path) + if disc and not disc.get("_laeuft"): + titel = disc.get("title") + job_id = str(uuid.uuid4()) - await asyncio.to_thread(db.insert_job, job_id, device_path, None, request.title, ziel) + await asyncio.to_thread(db.insert_job, job_id, device_path, None, titel, ziel) await asyncio.to_thread( db.add_log, "info", "api", f"Job {job_id} angelegt für {device_path}" + (f" → {ziel}" if ziel else ""), @@ -350,6 +357,24 @@ async def retry_transcode(job_id: str): return {"id": job_id, "status": "transcoding"} +@app.post("/jobs/{job_id}/cancel") +async def cancel_job(job_id: str): + """Bittet den Worker, den Job abzubrechen (kooperativ über die DB). + + Der Worker prüft das Flag bei jedem Fortschritts-Update und beendet den + Encoder-Prozess sauber — kein Celery-Task-ID-Tracking nötig. + """ + 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 ("completed", "failed"): + raise HTTPException(status_code=409, detail="Job ist bereits beendet") + + await asyncio.to_thread(db.update_job, job_id, status="canceling") + await asyncio.to_thread(db.add_log, "warning", "api", f"Job {job_id}: Abbruch angefordert") + return {"id": job_id, "status": "canceling"} + + @app.get("/capabilities") async def capabilities(): """Welche Encoder sind auf welchen Workern WIRKLICH verfügbar? @@ -453,6 +478,28 @@ async def browse(path: str = MEDIA_ROOT): return {"path": normalisiert, "parent": eltern, "dirs": ordner} +class MkdirRequest(BaseModel): + path: str + name: str + + +@app.post("/browse/mkdir", status_code=201) +async def browse_mkdir(request: MkdirRequest): + """Neuen Ordner unter /app/media anlegen (Speicherziele-Verwaltung).""" + basis = os.path.normpath(request.path) + if not basis.startswith(MEDIA_ROOT): + raise HTTPException(status_code=422, detail=f"Nur Pfade unter {MEDIA_ROOT}") + name = request.name.strip() + if not name or "/" in name or "\\" in name or name.startswith("."): + raise HTTPException(status_code=422, detail="Ungültiger Ordnername") + ziel = os.path.join(basis, name) + try: + await asyncio.to_thread(os.makedirs, ziel, exist_ok=True) + except OSError as e: + raise HTTPException(status_code=400, detail=f"Anlegen fehlgeschlagen: {e}") + return {"path": ziel} + + @app.get("/setup") async def setup_status(): """First-Run-Erkennung: wurde der Einrichtungs-Assistent abgeschlossen?""" diff --git a/docker/api/prescan/prescan.py b/docker/api/prescan/prescan.py index 4cc02d5..6979b03 100644 --- a/docker/api/prescan/prescan.py +++ b/docker/api/prescan/prescan.py @@ -12,6 +12,7 @@ from typing import Dict, List, Optional import detection from clients.tmdb import TMDBClient +from clients.jikan import JikanClient from clients.musicbrainz import MusicBrainzClient from clients.omdb import OMDbClient from clients.thetvdb import TheTVDBClient @@ -214,6 +215,7 @@ class PreScan: self.musicbrainz = MusicBrainzClient() self.thetvdb = TheTVDBClient() self.omdb = OMDbClient() + self.jikan = JikanClient() def scan(self, device_path: str) -> PreScanResult: """Führe Pre-Scan durch.""" @@ -421,7 +423,18 @@ class PreScan: if matched: break - # Fallback 1: OMDb (eigene Datenbasis — findet oft, was TMDB nicht + # Fallback 1: Jikan/MyAnimeList (kostenlos, KEIN Key) — für Anime die + # präziseste Quelle; wählt per Titel-Ähnlichkeit, nicht Treffer #1 + if not matched: + for kandidat in kandidaten: + jikan_treffer = self.jikan.lookup(kandidat) + if jikan_treffer: + confidence = 0.85 + metadata = jikan_treffer + matched = True + break + + # Fallback 2: OMDb (eigene Datenbasis — findet oft, was TMDB nicht # exakt trifft; braucht OMDB_API_KEY, sonst überspringt es sich selbst) if not matched: for kandidat in kandidaten: diff --git a/docker/ui/src/components/StorageMounts.tsx b/docker/ui/src/components/StorageMounts.tsx index 9c352f8..33e4ed2 100644 --- a/docker/ui/src/components/StorageMounts.tsx +++ b/docker/ui/src/components/StorageMounts.tsx @@ -1,5 +1,5 @@ import { useState, useEffect } from 'react' -import { HardDrive, Plus, Trash2, CheckCircle, AlertCircle, RefreshCw } from 'lucide-react' +import { HardDrive, Plus, Trash2, CheckCircle, AlertCircle, RefreshCw, Folder, ArrowUp, FolderPlus } from 'lucide-react' import { api } from '../lib/api' import { useDarkMode } from '../context/ThemeContext' @@ -31,8 +31,35 @@ export default function StorageMounts() { const [password, setPassword] = useState('') const [busy, setBusy] = useState(false) const [feedback, setFeedback] = useState(null) + const [browsePath, setBrowsePath] = useState('/app/media') + const [browseParent, setBrowseParent] = useState(null) + const [browseDirs, setBrowseDirs] = useState<{ name: string, path: string }[]>([]) + const [neuerOrdner, setNeuerOrdner] = useState('') const { theme } = useDarkMode() + const browsen = async (pfad: string) => { + try { + const r = await api.get('/browse', { params: { path: pfad } }) + setBrowsePath(r.data.path) + setBrowseParent(r.data.parent) + setBrowseDirs(r.data.dirs) + } catch { + setBrowseDirs([]) + } + } + + const ordnerAnlegen = async () => { + if (!neuerOrdner) return + try { + await api.post('/browse/mkdir', { path: browsePath, name: neuerOrdner }) + setNeuerOrdner('') + browsen(browsePath) + laden() + } catch (e: any) { + setFeedback(`✗ ${e?.response?.data?.detail || 'Ordner anlegen fehlgeschlagen'}`) + } + } + const laden = async () => { try { const [m, t] = await Promise.all([ @@ -46,7 +73,7 @@ export default function StorageMounts() { } } - useEffect(() => { laden() }, []) + useEffect(() => { laden(); browsen('/app/media') }, []) const hinzufuegen = async () => { setBusy(true) @@ -116,6 +143,51 @@ export default function StorageMounts() { )} + {/* Lokale Ordner durchsuchen + anlegen */} +
+
+ + {browsePath} + setNeuerOrdner(e.target.value)} + placeholder="neuer-ordner" + className={`w-32 px-2 py-1 text-xs border rounded ${theme === 'dark' ? 'bg-slate-800 border-slate-600 text-slate-200' : 'bg-white border-slate-300 text-slate-900'}`} + /> + +
+
+ {browseDirs.length === 0 ? ( +

Keine Unterordner

+ ) : ( + browseDirs.map(d => ( + + )) + )} +
+
+ {/* Neues Netzwerk-Ziel */}

diff --git a/docker/ui/src/pages/Dashboard.tsx b/docker/ui/src/pages/Dashboard.tsx index 21a1c83..18d6932 100644 --- a/docker/ui/src/pages/Dashboard.tsx +++ b/docker/ui/src/pages/Dashboard.tsx @@ -8,7 +8,7 @@ import LiveLogSection from '../components/LiveLogSection' interface Job { id: string type: 'cd' | 'dvd' | 'bluray' - status: 'pending' | 'processing' | 'transcoding' | 'completed' | 'failed' + status: 'pending' | 'processing' | 'transcoding' | 'canceling' | 'completed' | 'failed' device: string startTime: string endTime?: string @@ -32,6 +32,7 @@ function StatusBadge({ status }: { status: string }) { pending: 'bg-amber-900/30 text-amber-400 ring-amber-700/20', processing: 'bg-blue-900/30 text-blue-400 ring-blue-700/20', transcoding: 'bg-purple-900/30 text-purple-400 ring-purple-700/20', + canceling: 'bg-orange-900/30 text-orange-400 ring-orange-700/20', completed: 'bg-emerald-900/30 text-emerald-400 ring-emerald-700/20', failed: 'bg-rose-900/30 text-rose-400 ring-rose-700/20', } @@ -39,6 +40,7 @@ function StatusBadge({ status }: { status: string }) { pending: 'bg-amber-100 text-amber-700 ring-amber-600/20', processing: 'bg-blue-100 text-blue-700 ring-blue-600/20', transcoding: 'bg-purple-100 text-purple-700 ring-purple-600/20', + canceling: 'bg-orange-100 text-orange-700 ring-orange-600/20', completed: 'bg-emerald-100 text-emerald-700 ring-emerald-600/20', failed: 'bg-rose-100 text-rose-700 ring-rose-600/20', } @@ -47,6 +49,7 @@ function StatusBadge({ status }: { status: string }) { pending: 'Wartend', processing: 'Rippen', transcoding: 'Komprimieren', + canceling: 'Wird abgebrochen…', completed: 'Fertig', failed: 'Fehler', } @@ -273,7 +276,7 @@ export default function Dashboard() { ) : ( jobs.slice(0, 10).map(job => ( - + @@ -304,6 +307,15 @@ export default function Dashboard() { Neu komprimieren )} + {(job.status === 'pending' || job.status === 'processing' || job.status === 'transcoding') && ( + + )} )) diff --git a/docker/worker/db.py b/docker/worker/db.py index e3767cd..64776ed 100644 --- a/docker/worker/db.py +++ b/docker/worker/db.py @@ -116,6 +116,17 @@ def get_settings(key: str = "ui") -> dict: return {} +def get_job_status(job_id: str) -> str: + """Nur der Status — der Worker prüft damit kooperative Abbruch-Anfragen.""" + from sqlalchemy import select + + with engine.connect() as conn: + zeile = conn.execute( + select(jobs.c.status).where(jobs.c.id == job_id) + ).first() + return zeile[0] if zeile else "" + + def update_job(job_id: str, **fields) -> None: with engine.begin() as conn: conn.execute(jobs.update().where(jobs.c.id == job_id).values(**fields)) diff --git a/docker/worker/ripping.py b/docker/worker/ripping.py index 0f49230..eb082d4 100644 --- a/docker/worker/ripping.py +++ b/docker/worker/ripping.py @@ -20,6 +20,11 @@ import tempfile RIP_OUTPUT_DIR = os.getenv("RIP_OUTPUT_DIR", "/app/media") +class RipAbbruch(Exception): + """Kooperativer Abbruch: vom Fortschritts-Callback geworfen, wenn der + Nutzer den Job abgebrochen hat (Status 'canceling' in der DB).""" + + def check_makemkv_installed() -> bool: """Prüft, ob makemkvcon installiert ist.""" return shutil.which("makemkvcon") is not None @@ -124,10 +129,15 @@ def run_handbrake(input_path: str, output_path: str, preset: str = DEFAULT_HB_PR bufsize=1 ) - for line in process.stdout: - progress = get_progress_from_line(line) - if progress > 0 and progress_cb: - progress_cb(progress) + try: + for line in process.stdout: + progress = get_progress_from_line(line) + if progress > 0 and progress_cb: + progress_cb(progress) + except RipAbbruch: + process.kill() + process.wait() + return {"status": "cancelled", "error": "Abgebrochen durch Nutzer"} process.wait() @@ -185,15 +195,20 @@ def run_makemkv(device_path: str, output_dir: str, progress_cb=None) -> dict: ) letzte_meldung = "" - for line in process.stdout: - progress = get_progress_from_prgv(line) - if progress >= 0 and progress_cb: - progress_cb(progress) - elif line.startswith("MSG:"): - # MSG:code,flags,count,"message",... — Klartext ist Feld 4 - teile = line.split(",", 4) - if len(teile) >= 4: - letzte_meldung = teile[3].strip('"') + try: + for line in process.stdout: + progress = get_progress_from_prgv(line) + if progress >= 0 and progress_cb: + progress_cb(progress) + elif line.startswith("MSG:"): + # MSG:code,flags,count,"message",... — Klartext ist Feld 4 + teile = line.split(",", 4) + if len(teile) >= 4: + letzte_meldung = teile[3].strip('"') + except RipAbbruch: + process.kill() + process.wait() + return {"status": "cancelled", "error": "Abgebrochen durch Nutzer"} process.wait() @@ -267,8 +282,13 @@ def rip_cd(device_path: str, disc_id: str, progress_cb=None, output_dir: str = N bufsize=1 ) - for line in process.stdout: - melde(50, line.strip()[:200]) + try: + for line in process.stdout: + melde(50, line.strip()[:200]) + except RipAbbruch: + process.kill() + process.wait() + return {"status": "cancelled", "error": "Abgebrochen durch Nutzer"} process.wait() diff --git a/docker/worker/tasks.py b/docker/worker/tasks.py index 1810e3d..cbad244 100644 --- a/docker/worker/tasks.py +++ b/docker/worker/tasks.py @@ -22,6 +22,7 @@ from detection import detect_disc_type from ripping import ( DEFAULT_HB_PRESET, RIP_OUTPUT_DIR, + RipAbbruch, rip_cd, rip_video, run_handbrake, @@ -40,8 +41,20 @@ def _zielbasis(target_dir, disc_type: str) -> str: return os.path.join(RIP_OUTPUT_DIR, disc_type) +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 _job_abschliessen(job_id: str, ergebnis: dict) -> None: """Schreibt den Endzustand eines Jobs (completed/failed) nach Postgres.""" + 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") + return if ergebnis.get("status") == "success": db.update_job( job_id, @@ -69,6 +82,11 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None): dort eingehängte Shares (NFS/SMB) sind damit direkt wählbar. """ db.init_db() + + if _abbruch_angefordert(job_id): + _job_abschliessen(job_id, {"status": "cancelled"}) + return {"status": "cancelled"} + disc_type = detect_disc_type(device_path) if disc_type in ("no_disc", "unknown"): @@ -90,6 +108,8 @@ def rip_disc(self, device_path: str, job_id: str, target_dir: str = None): if progress == letzter[0]: return letzter[0] = progress + if _abbruch_angefordert(job_id): + raise RipAbbruch() self.update_state( state="PROGRESS", meta={"progress": progress, "status": "ripping", "message": message}, @@ -162,14 +182,23 @@ def transcode_files(self, job_id: str, raw_dir: str, final_dir: str): ) anzahl = len(quellen) + letzter = [-1] for index, quelle in enumerate(quellen): ziel = os.path.join(final_dir, os.path.basename(quelle)) def datei_fortschritt(p, _index=index): gesamt = int((_index * 100 + p) / anzahl) + if gesamt == letzter[0]: + return + letzter[0] = gesamt + if _abbruch_angefordert(job_id): + raise RipAbbruch() db.update_job(job_id, progress=min(99, gesamt)) hb = run_handbrake(quelle, ziel, preset=preset, progress_cb=datei_fortschritt) + if hb.get("status") == "cancelled": + _job_abschliessen(job_id, hb) + return hb if hb.get("status") != "success": ergebnis = { "status": "error",