f74f4e54f6
Ampel / ampel (push) Successful in 32s
GET /jobs nutzt response_model=List[Job]; das Job-Model hatte kein meta-Feld, also schnitt FastAPI die Disc-Metadaten (poster_path) weg -> Dashboard.tsx bekam job.meta = undefined -> Filmstreifen-Platzhalter statt TMDB-Poster, sowohl in der Jobliste als auch im aktiven Rip-Header (beide aus /jobs). Additiv: meta: Optional[Dict] ins Job-Model + in _job_row_to_model parsen (json.loads wie im Detail-Endpunkt, defensiv gegen kaputtes JSON). Keine UI-Aenderung noetig -- posterUrl() rendert dann die vorhandenen Poster. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
1505 lines
58 KiB
Python
1505 lines
58 KiB
Python
from fastapi import FastAPI, HTTPException, Request, Response
|
||
from fastapi.middleware.cors import CORSMiddleware
|
||
from fastapi.responses import FileResponse, StreamingResponse
|
||
from pydantic import BaseModel
|
||
from typing import List, Optional, Dict
|
||
from pathlib import Path
|
||
import asyncio
|
||
import json
|
||
import os
|
||
import shutil
|
||
import uuid
|
||
|
||
import db
|
||
import devices as device_discovery
|
||
import makemkv_key
|
||
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
|
||
|
||
from config_validation import validate_config, ConfigValidationError
|
||
from cache import init_cache
|
||
from ratelimit import check_rate_limit, get_rate_limit_remaining
|
||
from prescan import PreScan
|
||
from nfo_generator import NFOGenerator
|
||
from image_downloader import ImageDownloader
|
||
|
||
# Auth (JWT/Login/API-Keys) KOMPLETT entfernt — Commander-Entscheid 24.07.2026:
|
||
# Rippy läuft ausschließlich im Heimnetz, die Endpoints schützten ohnehin
|
||
# nichts (kein Login-Flow im UI) und waren damit Placebo-Oberfläche.
|
||
# Dokumentiert in KONZEPT.md Abschnitt 10.
|
||
app = FastAPI(
|
||
title="Rippy API",
|
||
description="API für das automatische Ripping-System",
|
||
version="1.0.0"
|
||
)
|
||
|
||
|
||
@app.on_event("startup")
|
||
async def startup_event():
|
||
"""Initialisiere Cache + Datenbank, validiere Konfiguration, starte Disc-Wache."""
|
||
init_cache()
|
||
db.init_db()
|
||
|
||
try:
|
||
validate_config()
|
||
except ConfigValidationError as e:
|
||
print(f"⚠️ Konfigurations-Warnung: {e}")
|
||
|
||
# Gespeicherte Netzwerk-Speicherziele wiederherstellen — NICHT-BLOCKIEREND.
|
||
# Ein zickiger/langsamer Netz-Mount (der CIFS-Schreibtest in mounten() kann im
|
||
# Kernel haengen, wait_for_response) darf den API-Start NIE blockieren. Vorfall
|
||
# 24.07.: ~5 min "Waiting for application startup", kein Endpoint bedient, bis der
|
||
# soft-Mount per Timeout abbrach. Darum im Hintergrund: die Mounts stellen sich
|
||
# her, sobald der Server antwortet, ohne die API auszubremsen.
|
||
def remount():
|
||
for meldung in mount_verwaltung.alle_remounten():
|
||
db.add_log("info", "mounts", meldung)
|
||
asyncio.create_task(asyncio.to_thread(remount))
|
||
|
||
asyncio.create_task(disc_watcher())
|
||
# MakeMKV-Beta-Key automatisch aktuell halten (wechselt ~monatlich, laeuft zum
|
||
# Monatsende ab) — sonst blockt Blu-ray-Ripping irgendwann still. Taeglicher
|
||
# Forum-Abgleich; wirkt ohne Rebuild ab dem naechsten Rip.
|
||
asyncio.create_task(makemkv_key.refresh_loop())
|
||
|
||
|
||
# Auto-Pre-Scan-Ergebnisse je Laufwerk: das Dashboard zeigt damit sofort,
|
||
# WAS im Laufwerk liegt (Titel/Jahr/Poster), ohne dass jemand klicken muss.
|
||
DISC_CACHE: Dict[str, Dict] = {}
|
||
|
||
|
||
def _duplikat_suchen(fingerprint: str):
|
||
"""Wurde eine Disc mit diesem Fingerabdruck schon erfolgreich gerippt?"""
|
||
if not fingerprint:
|
||
return None
|
||
for job in db.list_jobs(200):
|
||
if job.get("status") != "completed":
|
||
continue
|
||
try:
|
||
meta = json.loads(job.get("meta") or "{}")
|
||
except ValueError:
|
||
continue
|
||
if meta.get("fingerprint") == fingerprint:
|
||
return {
|
||
"job_id": job["id"],
|
||
"title": job.get("title"),
|
||
"finished_at": job["finished_at"].isoformat() if job.get("finished_at") else None,
|
||
}
|
||
return None
|
||
|
||
|
||
async def _auto_prescan(pfad: str):
|
||
"""Identifiziert die eingelegte Disc im Hintergrund und cached das Ergebnis."""
|
||
if DISC_CACHE.get(pfad, {}).get("_laeuft"):
|
||
return
|
||
DISC_CACHE[pfad] = {"_laeuft": True, "title": "Wird erkannt…"}
|
||
try:
|
||
prescan = PreScan()
|
||
ergebnis = await asyncio.to_thread(prescan.scan, pfad)
|
||
DISC_CACHE[pfad] = ergebnis.to_dict()
|
||
jahr_text = f" ({ergebnis.year})" if (
|
||
ergebnis.year and not ergebnis.title.endswith(f"({ergebnis.year})")
|
||
) else ""
|
||
db.add_log(
|
||
"info", "watcher",
|
||
f"Disc erkannt: {ergebnis.title}{jahr_text}"
|
||
+ f" [{ergebnis.disc_type}, Confidence {ergebnis.confidence:.0%}] auf {pfad}",
|
||
)
|
||
# Duplikat-Warnung: dieselbe Disc (Fingerabdruck) schon fertig gerippt?
|
||
dup = await asyncio.to_thread(_duplikat_suchen, ergebnis.fingerprint)
|
||
if dup:
|
||
DISC_CACHE[pfad]["bereits_gerippt"] = dup
|
||
db.add_log(
|
||
"warning", "watcher",
|
||
f"Diese Disc wurde bereits gerippt ({dup.get('title')}, "
|
||
f"Job {dup['job_id'][:8]}…) — Dashboard zeigt den Hinweis.",
|
||
)
|
||
await _auto_rip_wenn_aktiviert(pfad)
|
||
except Exception as e:
|
||
DISC_CACHE.pop(pfad, None)
|
||
print(f"Auto-Pre-Scan {pfad}: {e}")
|
||
|
||
|
||
async def _auto_rip_wenn_aktiviert(pfad: str):
|
||
"""Vollautomatik (Setting autoRipStart): Disc erkannt → Rip startet sofort.
|
||
|
||
Commander-Wunsch 24.07.: wahlweise Popup ODER Automatik. Ziel-Ordner
|
||
kommt aus den Schnellwahl-Einstellungen (Serie → seriesDir, sonst
|
||
movieDir; CD → musicDir) — genau wie ein Klick im Dialog.
|
||
"""
|
||
einstellungen = await asyncio.to_thread(db.get_settings)
|
||
if not einstellungen.get("autoRipStart"):
|
||
return
|
||
if await asyncio.to_thread(db.has_active_job, pfad):
|
||
return
|
||
disc = DISC_CACHE.get(pfad) or {}
|
||
if disc.get("_laeuft"):
|
||
return
|
||
if disc.get("bereits_gerippt"):
|
||
await asyncio.to_thread(
|
||
db.add_log, "info", "api",
|
||
"Automatik übersprungen: Disc wurde bereits gerippt — "
|
||
'manuell per „Rippen starten" trotzdem möglich.',
|
||
)
|
||
return
|
||
|
||
basis = einstellungen.get("outputDir") or MEDIA_ROOT
|
||
meta = disc.get("metadata") or {}
|
||
if disc.get("disc_type") == "CD":
|
||
unterordner = einstellungen.get("musicDir") or "music"
|
||
elif meta.get("type") == "tv":
|
||
unterordner = einstellungen.get("seriesDir") or "series"
|
||
else:
|
||
unterordner = einstellungen.get("movieDir") or "movies"
|
||
ziel = os.path.normpath(os.path.join(basis, unterordner))
|
||
if not ziel.startswith(MEDIA_ROOT):
|
||
ziel = None
|
||
|
||
job_id = str(uuid.uuid4())
|
||
meta_json = json.dumps({
|
||
"year": disc.get("year"),
|
||
"confidence": disc.get("confidence"),
|
||
**meta,
|
||
})
|
||
await asyncio.to_thread(
|
||
db.insert_job, job_id, pfad, None, disc.get("title"), ziel, meta_json
|
||
)
|
||
await asyncio.to_thread(
|
||
db.add_log, "info", "api",
|
||
f'Automatik: Rip für „{disc.get("title")}" gestartet ({pfad} → {ziel})',
|
||
)
|
||
start_rip(pfad, job_id, ziel)
|
||
|
||
|
||
async def disc_watcher():
|
||
"""Disc-Wache: pollt die Laufwerke, protokolliert Einwurf/Auswurf und
|
||
stößt beim Einlegen automatisch den Pre-Scan an (Dashboard-Disc-Karte).
|
||
|
||
Ersetzt den nie gebauten udev-Daemon aus dem KONZEPT: udev funktioniert im
|
||
Container nicht sinnvoll (kein udevd) — ein 3-Sekunden-Poll per ioctl ist
|
||
für den Heim-Use-Case gleichwertig und läuft überall.
|
||
"""
|
||
bekannt: Dict[str, int] = {}
|
||
while True:
|
||
try:
|
||
for pfad in device_discovery.list_optical_devices():
|
||
try:
|
||
status = await asyncio.to_thread(drive_status, pfad)
|
||
except OSError:
|
||
continue
|
||
vorher = bekannt.get(pfad)
|
||
if vorher is None:
|
||
# Erststart: liegt schon eine Disc drin, direkt erkennen
|
||
if status == CDS_DISC_OK:
|
||
asyncio.create_task(_auto_prescan(pfad))
|
||
elif status != vorher:
|
||
if status == CDS_DISC_OK:
|
||
db.add_log("info", "watcher", f"Disc eingelegt: {pfad}")
|
||
asyncio.create_task(_auto_prescan(pfad))
|
||
elif status in (CDS_NO_DISC, CDS_TRAY_OPEN) and vorher == CDS_DISC_OK:
|
||
db.add_log("info", "watcher", f"Disc entfernt: {pfad}")
|
||
DISC_CACHE.pop(pfad, None)
|
||
bekannt[pfad] = status
|
||
except Exception as e:
|
||
print(f"Disc-Wache: {e}")
|
||
await asyncio.sleep(3)
|
||
|
||
|
||
# Middleware für Rate-Limiting (pro Client-IP — schützt vor Amok-Skripten,
|
||
# nicht vor Angreifern; Rippy ist Heimnetz-only)
|
||
@app.middleware("http")
|
||
async def rate_limit_middleware(request: Request, call_next):
|
||
"""Rate-Limiting Middleware."""
|
||
client_ip = request.client.host
|
||
|
||
# Rate Limit prüfen
|
||
if not check_rate_limit(client_ip):
|
||
return Response(
|
||
content=json.dumps({"error": "Rate limit exceeded"}),
|
||
status_code=429,
|
||
media_type="application/json"
|
||
)
|
||
|
||
response = await call_next(request)
|
||
|
||
# Füge Rate-Limit Header hinzu
|
||
remaining = get_rate_limit_remaining(client_ip)
|
||
response.headers["X-RateLimit-Remaining"] = str(remaining)
|
||
|
||
return response
|
||
|
||
# CORS hinzufügen
|
||
app.add_middleware(
|
||
CORSMiddleware,
|
||
allow_origins=["*"],
|
||
allow_credentials=True,
|
||
allow_methods=["*"],
|
||
allow_headers=["*"],
|
||
)
|
||
|
||
class Job(BaseModel):
|
||
id: str
|
||
type: str
|
||
status: str
|
||
device: str
|
||
startTime: str
|
||
endTime: Optional[str] = None
|
||
progress: int = 0
|
||
title: Optional[str] = None
|
||
error: Optional[str] = None
|
||
can_retry: bool = False # Rohdaten vorhanden → „Neu komprimieren" sinnvoll
|
||
meta: Optional[Dict] = None # Disc-Metadaten (Poster/Jahr/Plot) — fürs Thumbnail in der Jobliste + aktivem Rip-Header
|
||
|
||
class Device(BaseModel):
|
||
id: str
|
||
name: str
|
||
type: str
|
||
path: str
|
||
status: str
|
||
model: Optional[str] = None
|
||
serial: Optional[str] = None
|
||
disc: Optional[Dict] = None # Auto-Pre-Scan-Ergebnis (Titel/Jahr/Poster)
|
||
|
||
|
||
def _job_row_to_model(zeile: dict) -> Job:
|
||
"""DB-Zeile → UI-Form (Worker-Status 'running' heißt im UI 'processing')."""
|
||
status_map = {"running": "processing"}
|
||
# meta (JSON-Text) enthält u. a. poster_path — die UI baut daraus das Thumbnail.
|
||
# Muss hier mit ins Job-Model, sonst schneidet FastAPIs response_model es weg
|
||
# (Befund 25.07.: meta kam nie in der Jobliste an → Filmstreifen-Platzhalter).
|
||
try:
|
||
meta = json.loads(zeile["meta"]) if zeile.get("meta") else None
|
||
except (ValueError, TypeError):
|
||
meta = None
|
||
return Job(
|
||
id=zeile["id"],
|
||
type=zeile.get("disc_type") or "unknown",
|
||
status=status_map.get(zeile["status"], zeile["status"]),
|
||
device=zeile.get("device") or "",
|
||
startTime=zeile["created_at"].isoformat() if zeile.get("created_at") else "",
|
||
endTime=zeile["finished_at"].isoformat() if zeile.get("finished_at") else None,
|
||
progress=zeile.get("progress") or 0,
|
||
title=zeile.get("title"),
|
||
error=zeile.get("error"),
|
||
meta=meta,
|
||
)
|
||
|
||
@app.get("/health")
|
||
async def health_check():
|
||
return {"status": "ok", "service": "api"}
|
||
|
||
|
||
@app.get("/")
|
||
async def root():
|
||
return {
|
||
"name": "Rippy",
|
||
"version": "1.0.0",
|
||
"description": "Automatisches Ripping-System für CD, DVD und Blu-ray"
|
||
}
|
||
|
||
|
||
def _kann_neu_komprimieren(job: dict, work_dir: str) -> bool:
|
||
"""Nur wenn Rohdaten wirklich noch daliegen — der „Neu komprimieren"-Knopf
|
||
an einem Job, der nie gerippt hat, war Unsinn (Befund 24.07.)."""
|
||
if job.get("status") != "failed":
|
||
return False
|
||
if os.path.isdir(os.path.join("/app/temp/raw", job["id"])):
|
||
return True
|
||
return work_dir.startswith(MEDIA_ROOT) and os.path.isdir(os.path.join(work_dir, job["id"]))
|
||
|
||
|
||
@app.get("/jobs", response_model=List[Job])
|
||
async def get_jobs():
|
||
"""Holt alle Jobs aus der Datenbank (neueste zuerst)."""
|
||
def sammle():
|
||
work_dir = os.path.normpath((db.get_settings().get("workDir") or "").strip() or "/")
|
||
modelle = []
|
||
for z in db.list_jobs():
|
||
modell = _job_row_to_model(z)
|
||
modell.can_retry = _kann_neu_komprimieren(z, work_dir)
|
||
modelle.append(modell)
|
||
return modelle
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@app.delete("/jobs/{job_id}")
|
||
async def delete_job(job_id: str):
|
||
"""Entfernt einen erledigten Job aus der Liste (Dateien bleiben liegen)."""
|
||
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"] not in ("completed", "failed"):
|
||
raise HTTPException(status_code=409, detail="Job läuft noch — erst abbrechen")
|
||
await asyncio.to_thread(db.delete_job, job_id)
|
||
await asyncio.to_thread(db.add_log, "info", "api", f"Job {job_id} aus der Liste entfernt")
|
||
return {"status": "deleted"}
|
||
|
||
|
||
@app.delete("/jobs")
|
||
async def delete_finished_jobs():
|
||
"""Räumt ALLE erledigten Jobs (fertig + fehlgeschlagen) aus der Liste."""
|
||
anzahl = await asyncio.to_thread(db.delete_finished_jobs)
|
||
await asyncio.to_thread(
|
||
db.add_log, "info", "api", f"Job-Liste aufgeräumt ({anzahl} erledigte Einträge entfernt)"
|
||
)
|
||
return {"deleted": anzahl}
|
||
|
||
|
||
class JobCreateRequest(BaseModel):
|
||
device_path: Optional[str] = None
|
||
device: Optional[str] = None # Alias, so schickt es das UI
|
||
title: Optional[str] = None
|
||
target_dir: Optional[str] = None # Ablageziel unter /app/media (frei wählbar)
|
||
series: Optional[str] = None # Serien-Flow: Ablage <Serie>/Season NN
|
||
season: Optional[int] = None
|
||
main_feature_only: Optional[bool] = None # pro Rip; None = Setting gilt
|
||
titles: Optional[List[int]] = None # exakte Titel-Auswahl (Track-Tabelle)
|
||
transcode_node: Optional[str] = None # gewählter Encoder-Worker (Celery-Node)
|
||
|
||
|
||
MEDIA_ROOT = "/app/media"
|
||
|
||
|
||
def _validiere_ziel(target_dir: Optional[str]) -> Optional[str]:
|
||
"""Ziel muss unter /app/media liegen — Pfad-Ausbrüche (..) fliegen raus."""
|
||
if not target_dir:
|
||
return None
|
||
normalisiert = os.path.normpath(target_dir)
|
||
if not normalisiert.startswith(MEDIA_ROOT):
|
||
raise HTTPException(
|
||
status_code=422,
|
||
detail=f"Ziel muss unter {MEDIA_ROOT} liegen (Shares dort einhängen)",
|
||
)
|
||
return normalisiert
|
||
|
||
|
||
@app.post("/jobs", status_code=201)
|
||
async def create_job(request: JobCreateRequest):
|
||
"""Legt einen Rip-Job an und schickt ihn an den Worker.
|
||
|
||
Das war DIE fehlende Stelle: bis 23.07. gab es keinerlei Code-Pfad,
|
||
der je einen Rip ausgelöst hat.
|
||
"""
|
||
device_path = request.device_path or request.device
|
||
if not device_path:
|
||
raise HTTPException(status_code=422, detail="device_path fehlt")
|
||
if device_path not in device_discovery.list_optical_devices():
|
||
raise HTTPException(status_code=404, detail=f"Laufwerk {device_path} nicht gefunden")
|
||
ziel = _validiere_ziel(request.target_dir)
|
||
|
||
# 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. Serien-Flow und Hauptfilm-Wahl wandern
|
||
# ebenfalls in die Job-Metadaten.
|
||
titel = request.title
|
||
meta_dict = {}
|
||
disc = DISC_CACHE.get(device_path)
|
||
if disc and not disc.get("_laeuft"):
|
||
if not titel:
|
||
titel = disc.get("title")
|
||
meta_dict = {
|
||
"year": disc.get("year"),
|
||
"confidence": disc.get("confidence"),
|
||
"fingerprint": disc.get("fingerprint"),
|
||
**(disc.get("metadata") or {}),
|
||
}
|
||
if request.series and request.series.strip():
|
||
meta_dict["series"] = request.series.strip()
|
||
meta_dict["season"] = max(1, int(request.season or 1))
|
||
if request.main_feature_only is not None:
|
||
meta_dict["main_feature_only"] = request.main_feature_only
|
||
if request.titles:
|
||
titel_liste = sorted({int(t) for t in request.titles if int(t) >= 0})[:200]
|
||
if titel_liste:
|
||
meta_dict["titles"] = titel_liste
|
||
if request.transcode_node:
|
||
meta_dict["transcode_node"] = request.transcode_node
|
||
meta_json = json.dumps(meta_dict) if meta_dict else None
|
||
|
||
job_id = str(uuid.uuid4())
|
||
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 ""),
|
||
)
|
||
start_rip(device_path, job_id, ziel)
|
||
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
|
||
|
||
|
||
def _sicherer_dateiname(name: str) -> bool:
|
||
"""Pure Funktion (testbar): nur nackte Dateinamen, keine Pfad-Tricks."""
|
||
return bool(name) and "/" not in name and "\\" not in name and not name.startswith(".")
|
||
|
||
|
||
def _job_ausgabeordner(job: dict) -> str:
|
||
"""Validierter Ausgabeordner eines Jobs — strikt unter /app/media."""
|
||
ausgabe = os.path.normpath(job.get("output_path") or "")
|
||
if not ausgabe.startswith(MEDIA_ROOT):
|
||
raise HTTPException(status_code=404, detail="Job hat keinen Ausgabeordner unter /app/media")
|
||
return ausgabe
|
||
|
||
|
||
@app.get("/jobs/{job_id}/files")
|
||
async def list_job_files(job_id: str):
|
||
"""Dateien eines fertigen Jobs — fürs Download-Menü im Dashboard.
|
||
|
||
Vorher kam man an fertige MKVs nur per scp auf die VM.
|
||
"""
|
||
job = await asyncio.to_thread(db.get_job, job_id)
|
||
if not job:
|
||
raise HTTPException(status_code=404, detail="Job nicht gefunden")
|
||
ausgabe = _job_ausgabeordner(job)
|
||
|
||
def liste():
|
||
try:
|
||
eintraege = sorted(os.listdir(ausgabe))
|
||
except OSError:
|
||
return None
|
||
dateien = []
|
||
for name in eintraege:
|
||
pfad = os.path.join(ausgabe, name)
|
||
if os.path.isfile(pfad):
|
||
try:
|
||
groesse_mb = round(os.path.getsize(pfad) / 1024**2, 1)
|
||
except OSError:
|
||
groesse_mb = None
|
||
dateien.append({"name": name, "size_mb": groesse_mb})
|
||
return dateien
|
||
|
||
dateien = await asyncio.to_thread(liste)
|
||
if dateien is None:
|
||
raise HTTPException(
|
||
status_code=404,
|
||
detail="Ausgabeordner nicht lesbar — Job noch nicht fertig oder Ziel ausgehängt?",
|
||
)
|
||
return {"job_id": job_id, "output_path": ausgabe, "files": dateien}
|
||
|
||
|
||
@app.get("/jobs/{job_id}/files/{dateiname}")
|
||
async def download_job_file(job_id: str, dateiname: str):
|
||
"""Streamt EINE Datei eines Jobs zum Browser (Download-Knopf).
|
||
|
||
Pfad-Validierung strikt: nackter Dateiname, realpath muss unter
|
||
/app/media bleiben (kein ..-Ausbruch, kein Symlink nach draußen).
|
||
"""
|
||
job = await asyncio.to_thread(db.get_job, job_id)
|
||
if not job:
|
||
raise HTTPException(status_code=404, detail="Job nicht gefunden")
|
||
ausgabe = _job_ausgabeordner(job)
|
||
if not _sicherer_dateiname(dateiname):
|
||
raise HTTPException(status_code=422, detail="Ungültiger Dateiname")
|
||
pfad = os.path.join(ausgabe, dateiname)
|
||
|
||
def pruefe():
|
||
return os.path.isfile(pfad) and os.path.realpath(pfad).startswith(MEDIA_ROOT)
|
||
|
||
if not await asyncio.to_thread(pruefe):
|
||
raise HTTPException(status_code=404, detail="Datei nicht gefunden")
|
||
return FileResponse(pfad, filename=dateiname, media_type="application/octet-stream")
|
||
|
||
|
||
@app.get("/storage-targets")
|
||
async def storage_targets():
|
||
"""Verfügbare Ablageziele: Verzeichnisse unter /app/media inkl. Mounts.
|
||
|
||
NFS/SMB-Shares, die auf der VM unter /srv/rippy/media eingehängt werden,
|
||
tauchen hier automatisch auf (rslave-Bind in docker-compose).
|
||
"""
|
||
def sammle():
|
||
ziele = []
|
||
try:
|
||
eintraege = sorted(os.listdir(MEDIA_ROOT))
|
||
except OSError:
|
||
return ziele
|
||
for name in eintraege:
|
||
pfad = os.path.join(MEDIA_ROOT, name)
|
||
try:
|
||
ist_mount = os.path.ismount(pfad)
|
||
except OSError:
|
||
# Toter CIFS-Mount (NAS weg) — als Ziel unbrauchbar, aber NICHT
|
||
# den ganzen Endpoint sprengen (Befund 24.07.: Eintrag flog raus,
|
||
# ließ sich aber nicht mehr neu anlegen). Überspringen ist ok,
|
||
# die Mount-Verwaltung (/storage-mounts) zeigt ihn zum Reparieren.
|
||
continue
|
||
if not ist_mount and not os.path.isdir(pfad):
|
||
continue
|
||
try:
|
||
nutzung = shutil.disk_usage(pfad)
|
||
frei_gb = round(nutzung.free / 1024**3, 1)
|
||
except OSError:
|
||
frei_gb = None
|
||
ziele.append({
|
||
"name": name,
|
||
"path": pfad,
|
||
"is_mount": ist_mount,
|
||
"free_gb": frei_gb,
|
||
})
|
||
return ziele
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@app.post("/devices/{name}/eject")
|
||
async def eject_device(name: str):
|
||
"""Wirft die Disc aus. Verweigert, wenn auf dem Gerät gerade ein Job läuft."""
|
||
device_path = f"/dev/{name}"
|
||
if device_path not in device_discovery.list_optical_devices():
|
||
raise HTTPException(status_code=404, detail=f"Laufwerk {device_path} nicht gefunden")
|
||
if await asyncio.to_thread(db.has_active_job, device_path):
|
||
raise HTTPException(
|
||
status_code=409, detail="Auf diesem Laufwerk läuft gerade ein Job"
|
||
)
|
||
try:
|
||
await asyncio.to_thread(device_discovery.eject, device_path)
|
||
except OSError as e:
|
||
raise HTTPException(status_code=500, detail=f"Auswurf fehlgeschlagen: {e}")
|
||
await asyncio.to_thread(db.add_log, "info", "api", f"Disc ausgeworfen: {device_path}")
|
||
return {"status": "ejected", "device": device_path}
|
||
|
||
|
||
@app.post("/devices/{name}/scan-tracks")
|
||
async def scan_tracks_starten(name: str):
|
||
"""Titel-Scan der eingelegten Disc anstoßen (Track-Auswahl-Tabelle).
|
||
|
||
Läuft als Worker-Task (nur der hat makemkvcon + Laufwerk); das UI pollt
|
||
GET /devices/{name}/tracks. Dauert je nach Disc 20–120 s.
|
||
"""
|
||
device_path = f"/dev/{name}"
|
||
if device_path not in device_discovery.list_optical_devices():
|
||
raise HTTPException(status_code=404, detail=f"Laufwerk {device_path} nicht gefunden")
|
||
if await asyncio.to_thread(db.has_active_job, device_path):
|
||
raise HTTPException(status_code=409, detail="Auf diesem Laufwerk läuft gerade ein Job")
|
||
|
||
await asyncio.to_thread(db.save_settings, {"status": "running"}, f"tracks:{device_path}")
|
||
celery_client.send_task("worker.tasks.scan_tracks", args=[device_path])
|
||
return {"status": "scanning"}
|
||
|
||
|
||
@app.get("/devices/{name}/tracks")
|
||
async def scan_tracks_ergebnis(name: str):
|
||
"""Ergebnis des Titel-Scans (Polling-Ziel des UI)."""
|
||
device_path = f"/dev/{name}"
|
||
daten = await asyncio.to_thread(db.get_settings, f"tracks:{device_path}")
|
||
if not daten:
|
||
return {"status": "none"}
|
||
return daten
|
||
|
||
|
||
@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/<job_id>
|
||
(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")
|
||
|
||
# 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 = job.get("output_path") or f"{basis}/{job_id}"
|
||
|
||
# An den (beim Rip gewählten) Encoder-Worker routen, sonst geteilte Queue
|
||
try:
|
||
meta = json.loads(job.get("meta") or "{}")
|
||
except ValueError:
|
||
meta = {}
|
||
from celery_client import transcode_queue
|
||
celery_client.send_task(
|
||
"worker.tasks.transcode_files",
|
||
args=[job_id, raw_dir, final_dir],
|
||
queue=transcode_queue(meta.get("transcode_node")),
|
||
)
|
||
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.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 — inkl.
|
||
Live-Erreichbarkeit (Celery-Ping + Herzschlag-Alter).
|
||
|
||
Jeder Worker meldet sich selbst (caps.py, minütlich) — auch optionale
|
||
Remote-GPU-Worker tauchen hier automatisch auf.
|
||
"""
|
||
def sammle():
|
||
zeilen = db.list_workers()
|
||
try:
|
||
antworten = celery_client.control.ping(timeout=1.0) or []
|
||
ping_knoten = [k for antwort in antworten for k in antwort.keys()]
|
||
except Exception:
|
||
ping_knoten = []
|
||
for zeile in zeilen:
|
||
# Celery-Knotenname = <name>@<hostname>. Der Ping liefert ihn voll;
|
||
# gematcht wird über den Hostname (der Docker-Worker heißt celery@…,
|
||
# der WORKER_NAME steckt nur in info). Laufen MEHRERE Worker auf
|
||
# demselben Host (Befund 24.07.), disambiguiert der Name-Teil —
|
||
# sonst bekäme der falsche Worker den Node. `node` ist der
|
||
# Routing-Ziel-Knoten für die gezielte Encoder-Wahl.
|
||
hostname = (zeile.get("info") or {}).get("hostname") or zeile["name"]
|
||
kandidaten = [k for k in ping_knoten if k.split("@", 1)[-1] == hostname]
|
||
node = next(
|
||
(k for k in kandidaten if k.split("@", 1)[0] == zeile["name"]), None
|
||
) or (kandidaten[0] if kandidaten else None)
|
||
zeile["online"] = node is not None
|
||
zeile["node"] = node
|
||
return zeilen
|
||
|
||
return {"workers": await asyncio.to_thread(sammle)}
|
||
|
||
|
||
@app.delete("/workers/{name}")
|
||
async def delete_worker(name: str):
|
||
"""Verwaisten Worker-Eintrag entfernen (alte Container-IDs nach Rebuilds).
|
||
|
||
Ein AKTIVER Worker meldet sich binnen einer Minute einfach wieder an —
|
||
löschen ist also immer gefahrlos."""
|
||
await asyncio.to_thread(db.delete_worker, name)
|
||
await asyncio.to_thread(db.add_log, "info", "api", f"Worker-Eintrag '{name}' entfernt")
|
||
return {"status": "deleted"}
|
||
|
||
|
||
@app.get("/metadata/tv/{tv_id}/season/{season}")
|
||
async def tv_season_laufzeiten(tv_id: int, season: int):
|
||
"""Episoden-Laufzeiten einer Staffel (TMDB) — Basis fürs
|
||
Episoden-Matching des Workers (Serien-Flow)."""
|
||
def hole():
|
||
prescan = PreScan()
|
||
daten = prescan.tmdb.get_tv_season(tv_id, season)
|
||
if not daten:
|
||
return None
|
||
return {
|
||
"episodes": [
|
||
{"episode": e.get("episode_number"), "runtime": e.get("runtime")}
|
||
for e in daten.get("episodes", [])
|
||
]
|
||
}
|
||
|
||
ergebnis = await asyncio.to_thread(hole)
|
||
if ergebnis is None:
|
||
raise HTTPException(status_code=404, detail="Staffel bei TMDB nicht gefunden")
|
||
return ergebnis
|
||
|
||
|
||
class MediaServerRefreshRequest(BaseModel):
|
||
url: str
|
||
api_key: str
|
||
|
||
|
||
@app.post("/mediaserver/refresh")
|
||
async def mediaserver_refresh(request: MediaServerRefreshRequest):
|
||
"""Bibliotheks-Scan von Jellyfin/Emby anstoßen — auch als Verbindungs-Test
|
||
aus den Einstellungen (POST /Library/Refresh, Header X-Emby-Token)."""
|
||
import requests as _requests
|
||
|
||
url = request.url.strip().rstrip("/")
|
||
if not url.startswith(("http://", "https://")):
|
||
raise HTTPException(status_code=422, detail="Server-URL muss mit http(s):// beginnen")
|
||
|
||
def anstossen():
|
||
return _requests.post(
|
||
url + "/Library/Refresh",
|
||
headers={"X-Emby-Token": request.api_key.strip()}, timeout=15,
|
||
)
|
||
|
||
try:
|
||
antwort = await asyncio.to_thread(anstossen)
|
||
except Exception as e:
|
||
raise HTTPException(status_code=400, detail=f"Server nicht erreichbar: {e}")
|
||
if antwort.status_code >= 300:
|
||
raise HTTPException(
|
||
status_code=400,
|
||
detail=f"Server antwortete mit HTTP {antwort.status_code} — API-Key prüfen "
|
||
"(Jellyfin: Administration → API-Schlüssel)",
|
||
)
|
||
await asyncio.to_thread(db.add_log, "info", "api", f"Bibliotheks-Refresh angestoßen ({url})")
|
||
return {"status": "refreshed"}
|
||
|
||
|
||
@app.get("/jobs/export")
|
||
async def export_jobs():
|
||
"""Job-Historie als CSV (Semikolon + BOM — öffnet sauber in deutschem Excel)."""
|
||
def baue():
|
||
import csv
|
||
import io
|
||
|
||
puffer = io.StringIO()
|
||
w = csv.writer(puffer, delimiter=";")
|
||
w.writerow(["ID", "Titel", "Typ", "Status", "Fortschritt %", "Gerät",
|
||
"Start", "Ende", "Ablage", "Fehler"])
|
||
for j in db.list_jobs(1000):
|
||
w.writerow([
|
||
j["id"], j.get("title") or "", j.get("disc_type") or "",
|
||
j.get("status") or "", j.get("progress") or 0, j.get("device") or "",
|
||
j["created_at"].isoformat() if j.get("created_at") else "",
|
||
j["finished_at"].isoformat() if j.get("finished_at") else "",
|
||
j.get("output_path") or "",
|
||
(j.get("error") or "").replace("\n", " "),
|
||
])
|
||
return puffer.getvalue()
|
||
|
||
inhalt = await asyncio.to_thread(baue)
|
||
return Response(
|
||
content="" + inhalt,
|
||
media_type="text/csv; charset=utf-8",
|
||
headers={"Content-Disposition": 'attachment; filename="rippy-jobs.csv"'},
|
||
)
|
||
|
||
|
||
@app.get("/metadata/status")
|
||
async def metadata_status():
|
||
"""Live-Prüfung der Metadaten-Quellen — beantwortet „funktioniert mein
|
||
Key?" sofort statt durch stilles Wegfallen einer Quelle."""
|
||
def pruefe():
|
||
from clients.omdb import OMDB_BASE_URL
|
||
from clients.tmdb import TMDB_BASE_URL
|
||
|
||
status = {}
|
||
prescan = PreScan()
|
||
# Bewusst am Cache VORBEI — ein alter Treffer soll keinen kaputten
|
||
# Key als "ok" tarnen. /configuration ist der kleinste Auth-Aufruf.
|
||
if not prescan.tmdb.api_key:
|
||
status["tmdb"] = "kein_key"
|
||
else:
|
||
try:
|
||
antwort = prescan.tmdb.session.get(
|
||
f"{TMDB_BASE_URL}/configuration",
|
||
params=prescan.tmdb._key_params, timeout=10,
|
||
)
|
||
status["tmdb"] = "ok" if antwort.status_code == 200 else "fehler"
|
||
except Exception:
|
||
status["tmdb"] = "fehler"
|
||
if not prescan.omdb.api_key:
|
||
status["omdb"] = "kein_key"
|
||
else:
|
||
try:
|
||
antwort = prescan.omdb.session.get(
|
||
OMDB_BASE_URL,
|
||
params={"apikey": prescan.omdb.api_key, "t": "Inception"},
|
||
timeout=10,
|
||
).json()
|
||
status["omdb"] = "ok" if antwort.get("Response") == "True" else "fehler"
|
||
except Exception:
|
||
status["omdb"] = "fehler"
|
||
status["jikan"] = "ok" # keyless — fällt nur bei Netzproblemen aus
|
||
return status
|
||
|
||
return await asyncio.to_thread(pruefe)
|
||
|
||
|
||
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-Status.
|
||
|
||
`mounted` = liegt ein Mount an; `reachable` = ist er auch WIRKLICH nutzbar
|
||
(toter CIFS-Mount nach NAS-Ausfall: mounted=true, reachable=false → das UI
|
||
bietet dann „Reparieren" statt den Eintrag verschwinden zu lassen)."""
|
||
def sammle():
|
||
eintraege = db.list_mounts()
|
||
ergebnis = []
|
||
for e in eintraege:
|
||
gemountet = mount_verwaltung.ist_gemountet(e["name"])
|
||
ergebnis.append({
|
||
"name": e["name"],
|
||
"type": e["typ"],
|
||
"source": e["quelle"],
|
||
"mounted": gemountet,
|
||
"reachable": mount_verwaltung.ist_erreichbar(e["name"]) if gemountet else False,
|
||
"has_credentials": bool(e.get("username")),
|
||
})
|
||
return ergebnis
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@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.
|
||
|
||
Existiert der Name schon UND ist erreichbar → 409. Existiert er, ist aber
|
||
TOT (NAS war weg) → wird mit den neuen Angaben frisch repariert statt
|
||
stur „Name bereits vergeben" zu melden (Befund 24.07.: man saß sonst fest)."""
|
||
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")
|
||
|
||
vorhanden = any(e["name"] == request.name for e in await asyncio.to_thread(db.list_mounts))
|
||
if vorhanden:
|
||
gemountet = await asyncio.to_thread(mount_verwaltung.ist_gemountet, request.name)
|
||
erreichbar = gemountet and await asyncio.to_thread(mount_verwaltung.ist_erreichbar, request.name)
|
||
if erreichbar:
|
||
raise HTTPException(
|
||
status_code=409,
|
||
detail="Name bereits vergeben und aktiv — erst entfernen, dann neu anlegen.",
|
||
)
|
||
# Toter/veralteter Eintrag → reparieren (lazy abhängen + frisch mounten)
|
||
aktion = mount_verwaltung.reparieren
|
||
else:
|
||
aktion = mount_verwaltung.mounten
|
||
|
||
try:
|
||
schreibbar = await asyncio.to_thread(
|
||
aktion,
|
||
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))
|
||
|
||
# Bei Reparatur die (evtl. neuen) Zugangsdaten in der DB aktualisieren
|
||
if vorhanden:
|
||
await asyncio.to_thread(db.delete_mount, request.name)
|
||
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" if schreibbar else "warning", "mounts",
|
||
f"Speicherziel '{request.name}' ({request.type}) "
|
||
+ ("repariert" if vorhanden else "eingehängt") + f": {request.source}"
|
||
+ ("" if schreibbar else " — ACHTUNG: NUR LESBAR (Schreibtest fehlgeschlagen)"),
|
||
)
|
||
return {"name": request.name, "mounted": True, "writable": schreibbar, "repaired": vorhanden}
|
||
|
||
|
||
@app.post("/storage-mounts/{name}/repair")
|
||
async def repair_storage_mount(name: str):
|
||
"""Toten/veralteten Mount mit den GESPEICHERTEN Zugangsdaten neu verbinden
|
||
(Reparieren-Knopf im UI) — ohne dass der Nutzer alles neu eintippt."""
|
||
eintrag = next(
|
||
(e for e in await asyncio.to_thread(db.list_mounts) if e["name"] == name), None
|
||
)
|
||
if not eintrag:
|
||
raise HTTPException(status_code=404, detail="Speicherziel nicht gefunden")
|
||
try:
|
||
schreibbar = await asyncio.to_thread(
|
||
mount_verwaltung.reparieren,
|
||
eintrag["name"], eintrag["typ"], eintrag["quelle"],
|
||
eintrag.get("optionen") or "", eintrag.get("username") or "",
|
||
eintrag.get("passwort") or "",
|
||
)
|
||
except RuntimeError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
await asyncio.to_thread(
|
||
db.add_log, "success" if schreibbar else "warning", "mounts",
|
||
f"Speicherziel '{name}' neu verbunden"
|
||
+ ("" if schreibbar else " — NUR LESBAR"),
|
||
)
|
||
return {"name": name, "mounted": True, "writable": schreibbar}
|
||
|
||
|
||
@app.get("/storage-mounts/shares")
|
||
async def list_shares(host: str, username: str = "", password: str = ""):
|
||
"""SMB-Freigaben eines Rechners auflisten (PC/NAS per Klick wählen)."""
|
||
try:
|
||
freigaben = await asyncio.to_thread(
|
||
mount_verwaltung.liste_smb_freigaben, host, username, password
|
||
)
|
||
except RuntimeError as e:
|
||
raise HTTPException(status_code=400, detail=str(e))
|
||
return {"host": host, "shares": freigaben}
|
||
|
||
|
||
@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("/browse")
|
||
async def browse(path: str = MEDIA_ROOT):
|
||
"""Server-seitiger Ordner-Browser für die Ziel-Auswahl (nur unter /app/media)."""
|
||
normalisiert = os.path.normpath(path)
|
||
if not normalisiert.startswith(MEDIA_ROOT):
|
||
raise HTTPException(status_code=422, detail=f"Nur Pfade unter {MEDIA_ROOT}")
|
||
|
||
def liste():
|
||
try:
|
||
eintraege = sorted(os.listdir(normalisiert))
|
||
except OSError:
|
||
return None
|
||
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
|
||
|
||
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, "files": dateien}
|
||
|
||
|
||
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("/system/updates")
|
||
async def system_updates():
|
||
"""Update-Check für die Kern-Werkzeuge (Einstellungen → System).
|
||
|
||
Quellen: HandBrake über die GitHub-Release-API (releases/latest →
|
||
tag_name), MakeMKV über die offizielle Download-Seite (Versionsnummer
|
||
im Seitentext). Ergebnisse werden 12 h gecacht — der Check ist Komfort,
|
||
kein Dauerfeuer auf fremde Server. Ein Selbst-Update gibt es bewusst
|
||
NICHT: die Versionen stecken im Worker-Image, das Update ist ein
|
||
Image-Rebuild (Befehl wird im UI angezeigt).
|
||
"""
|
||
import re as _re
|
||
|
||
import requests as _requests
|
||
|
||
from cache import get as cache_get, set as _cache_set
|
||
|
||
def hole_neueste():
|
||
ergebnis = {"makemkv": None, "handbrake": None}
|
||
cached = cache_get("updates:neueste")
|
||
if cached:
|
||
return cached
|
||
try:
|
||
antwort = _requests.get(
|
||
"https://api.github.com/repos/HandBrake/HandBrake/releases/latest",
|
||
timeout=15, headers={"Accept": "application/vnd.github+json"},
|
||
)
|
||
if antwort.status_code == 200:
|
||
ergebnis["handbrake"] = (antwort.json().get("tag_name") or "").lstrip("v")
|
||
except Exception:
|
||
pass
|
||
try:
|
||
antwort = _requests.get("https://www.makemkv.com/download/", timeout=15)
|
||
treffer = _re.search(r"MakeMKV\s+(\d+\.\d+\.\d+)", antwort.text or "")
|
||
if treffer:
|
||
ergebnis["makemkv"] = treffer.group(1)
|
||
except Exception:
|
||
pass
|
||
if ergebnis["makemkv"] or ergebnis["handbrake"]:
|
||
_cache_set("updates:neueste", ergebnis, expire=12 * 3600)
|
||
return ergebnis
|
||
|
||
def sammle():
|
||
neueste = hole_neueste()
|
||
installiert = {"makemkv": None, "handbrake": None}
|
||
for worker in db.list_workers():
|
||
info = worker.get("info") or {}
|
||
installiert["makemkv"] = installiert["makemkv"] or info.get("makemkv")
|
||
installiert["handbrake"] = installiert["handbrake"] or info.get("handbrake")
|
||
return {
|
||
werkzeug: {
|
||
"installiert": installiert[werkzeug],
|
||
"verfuegbar": neueste[werkzeug],
|
||
"update": bool(
|
||
installiert[werkzeug] and neueste[werkzeug]
|
||
and installiert[werkzeug] != neueste[werkzeug]
|
||
),
|
||
}
|
||
for werkzeug in ("makemkv", "handbrake")
|
||
}
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@app.get("/worker-setup/windows")
|
||
async def worker_setup_windows():
|
||
"""Windows-Installer-Skript für den nativen Transcode-Worker.
|
||
|
||
Das UI (Einstellungen → Worker → Windows) zeigt den passenden
|
||
Zwei-Zeilen-Aufruf — alles kommt von dieser Rippy-Instanz selbst.
|
||
"""
|
||
pfad = "worker_dist/install.ps1"
|
||
if not os.path.isfile(pfad):
|
||
raise HTTPException(status_code=404, detail="Installer nicht im Image — API neu bauen")
|
||
return FileResponse(pfad, media_type="text/plain", filename="install-rippy-worker.ps1")
|
||
|
||
|
||
@app.get("/worker-setup/windows-gui")
|
||
async def worker_setup_windows_gui():
|
||
"""Grafischer Windows-Installer (WinForms, PowerShell-Quelle) — Rückfall
|
||
für Fortgeschrittene; der Normalweg ist die .exe unten."""
|
||
pfad = "worker_dist/install-gui.ps1"
|
||
if not os.path.isfile(pfad):
|
||
raise HTTPException(status_code=404, detail="GUI-Installer nicht im Image — API neu bauen")
|
||
return FileResponse(pfad, media_type="text/plain", filename="rippy-worker-gui.ps1")
|
||
|
||
|
||
@app.get("/worker-setup/windows-exe")
|
||
async def worker_setup_windows_exe():
|
||
"""Fertige Windows-Installer-.exe (Rippy-Icon, kein Konsolenfenster).
|
||
|
||
Vorgebaut auf Windows (deploy/worker-windows/build-exe.ps1) — eine
|
||
Windows-.exe lässt sich nicht auf Linux bauen. Generisch: die GUI fragt
|
||
die Rippy-Adresse selbst ab, HandBrake + Worker-Code kommen zur Laufzeit."""
|
||
pfad = "worker_dist/RippyWorkerSetup.exe"
|
||
if not os.path.isfile(pfad):
|
||
raise HTTPException(status_code=404, detail="Installer-.exe nicht im Image — API neu bauen")
|
||
return FileResponse(
|
||
pfad, media_type="application/vnd.microsoft.portable-executable",
|
||
filename="RippyWorkerSetup.exe",
|
||
)
|
||
|
||
|
||
@app.get("/worker-setup/paket")
|
||
async def worker_setup_paket():
|
||
"""Worker-Quellcode als Zip — der Windows-Installer lädt ihn von hier.
|
||
|
||
Kein git, kein Docker auf der Zielmaschine nötig: die Rippy-Instanz
|
||
versorgt ihre Worker selbst (Dateien liegen via Dockerfile im Image).
|
||
"""
|
||
def baue():
|
||
import io
|
||
import zipfile
|
||
|
||
if not os.path.isdir("worker_dist"):
|
||
return None
|
||
puffer = io.BytesIO()
|
||
with zipfile.ZipFile(puffer, "w", zipfile.ZIP_DEFLATED) as z:
|
||
for name in sorted(os.listdir("worker_dist")):
|
||
if name.endswith(".py") and not name.startswith(("test_", "conftest")):
|
||
z.write(os.path.join("worker_dist", name), name)
|
||
elif name == "requirements.txt":
|
||
z.write(os.path.join("worker_dist", name), name)
|
||
return puffer.getvalue()
|
||
|
||
inhalt = await asyncio.to_thread(baue)
|
||
if not inhalt:
|
||
raise HTTPException(status_code=404, detail="Worker-Paket nicht im Image — API neu bauen")
|
||
return Response(
|
||
content=inhalt, media_type="application/zip",
|
||
headers={"Content-Disposition": 'attachment; filename="rippy-worker.zip"'},
|
||
)
|
||
|
||
|
||
@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?"""
|
||
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)."""
|
||
zeilen = await asyncio.to_thread(db.list_logs, min(limit, 1000))
|
||
return [
|
||
{
|
||
"id": str(z["id"]),
|
||
"timestamp": z["ts"].isoformat() if z.get("ts") else "",
|
||
"level": z.get("level") or "info",
|
||
"source": z.get("source") or "system",
|
||
"message": z.get("message") or "",
|
||
}
|
||
for z in zeilen
|
||
]
|
||
|
||
|
||
@app.get("/settings")
|
||
async def get_settings():
|
||
"""UI-Einstellungen aus der Datenbank (leeres Objekt = Defaults im UI)."""
|
||
return await asyncio.to_thread(db.get_settings)
|
||
|
||
|
||
@app.post("/settings")
|
||
async def save_settings(werte: Dict):
|
||
"""Speichert die UI-Einstellungen als JSON in der Datenbank."""
|
||
await asyncio.to_thread(db.save_settings, werte)
|
||
return {"status": "saved"}
|
||
|
||
|
||
@app.get("/devices", response_model=List[Device])
|
||
async def get_devices():
|
||
"""Alle optischen Laufwerke mit ehrlichem Status (leer/bereit + Disc-Typ).
|
||
|
||
Der alte Weg (udevadm + /dev/disc-Symlinks) lieferte im Container
|
||
prinzipbedingt nichts: kein udevd, keine udev-Datenbank, kein Daemon,
|
||
der Symlinks anlegt. Jetzt: /sys fürs Modell, ioctl für den Disc-Status.
|
||
"""
|
||
geraete = []
|
||
for pfad in device_discovery.list_optical_devices():
|
||
info = await asyncio.to_thread(device_discovery.device_info, pfad)
|
||
disc = DISC_CACHE.get(pfad)
|
||
if disc and not disc.get("_laeuft"):
|
||
info["disc"] = disc
|
||
geraete.append(Device(**info))
|
||
return geraete
|
||
|
||
|
||
# SSE-Stream für Echtzeit-Updates
|
||
@app.get("/stream/jobs")
|
||
async def job_stream():
|
||
"""SSE-Stream für Job-Updates.
|
||
|
||
Fix 23.07.: Der alte Generator sendete nur, wenn `sse_connections` gefüllt
|
||
war — aber NICHTS hat diese Liste je befüllt. Der Stream war ein Placebo.
|
||
"""
|
||
async def event_generator():
|
||
while True:
|
||
jobs = await get_jobs()
|
||
yield f"data: {json.dumps([j.dict() for j in jobs])}\n\n"
|
||
await asyncio.sleep(2)
|
||
|
||
return StreamingResponse(event_generator(), media_type="text/event-stream")
|
||
|
||
|
||
# /metadata/lookup + /metadata/confirm entfernt (24.07., mit der
|
||
# Metadaten-Seite): lookup scannte ein DUMMY-Device (/dev/dvd — existiert
|
||
# nicht) und confirm schrieb in einen Cache-Key, den nie jemand las.
|
||
# Die echte Korrektur läuft über /metadata/search + /metadata/override.
|
||
|
||
|
||
@app.get("/metadata/search")
|
||
async def metadata_search(q: str):
|
||
"""Manuelle Korrektur: Titel-Kandidaten aus ALLEN Quellen (TMDB/Jikan/OMDb).
|
||
|
||
KONZEPT Schritt 5: „Commander bestätigt oder korrigiert manuell" — das
|
||
hier ist der Korrektur-Teil, wenn die Automatik danebenliegt.
|
||
"""
|
||
def sammle():
|
||
prescan = PreScan()
|
||
ergebnisse = []
|
||
for movie in (prescan.tmdb.search_movie(q) or [])[:4]:
|
||
ergebnisse.append({
|
||
"title": movie.get("title", ""),
|
||
"year": int(movie["release_date"][:4]) if movie.get("release_date") else None,
|
||
"poster": f"https://image.tmdb.org/t/p/w342{movie['poster_path']}" if movie.get("poster_path") else "",
|
||
"overview": (movie.get("overview") or "")[:200],
|
||
"type": "movie", "source": "tmdb", "id": str(movie.get("id", "")),
|
||
})
|
||
for show in (prescan.tmdb.search_tv(q) or [])[:3]:
|
||
ergebnisse.append({
|
||
"title": show.get("name", ""),
|
||
"year": int(show["first_air_date"][:4]) if show.get("first_air_date") else None,
|
||
"poster": f"https://image.tmdb.org/t/p/w342{show['poster_path']}" if show.get("poster_path") else "",
|
||
"overview": (show.get("overview") or "")[:200],
|
||
"type": "tv", "source": "tmdb", "id": str(show.get("id", "")),
|
||
})
|
||
ergebnisse += prescan.jikan.suche(q)
|
||
ergebnisse += prescan.omdb.suche(q)
|
||
|
||
# Mit Poster zuerst, Duplikate (Titel+Jahr) raus
|
||
gesehen, dedup = set(), []
|
||
for e in sorted(ergebnisse, key=lambda x: 0 if x.get("poster") else 1):
|
||
schluessel = ((e.get("title") or "").lower(), e.get("year"))
|
||
if schluessel in gesehen:
|
||
continue
|
||
gesehen.add(schluessel)
|
||
dedup.append(e)
|
||
return dedup[:12]
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
class MetadataOverride(BaseModel):
|
||
device_path: str
|
||
title: str
|
||
year: Optional[int] = None
|
||
poster: Optional[str] = None
|
||
overview: Optional[str] = None
|
||
type: Optional[str] = "movie"
|
||
source: Optional[str] = None
|
||
id: Optional[str] = None
|
||
|
||
|
||
@app.post("/metadata/override")
|
||
async def metadata_override(request: MetadataOverride):
|
||
"""Nutzer-Wahl für DIESE Disc merken: Karte, Jobs und Cache (30 Tage).
|
||
|
||
Der Disc-Fingerabdruck (Label+Größe) macht die Korrektur wiedererkennbar —
|
||
dieselbe Disc wird beim nächsten Einlegen sofort richtig angezeigt.
|
||
"""
|
||
if request.device_path not in device_discovery.list_optical_devices():
|
||
raise HTTPException(status_code=404, detail="Laufwerk nicht gefunden")
|
||
|
||
def speichere():
|
||
from cache import set as cache_setter
|
||
from cache.keys import generate_prescan_key
|
||
from prescan.prescan import disc_fingerprint
|
||
|
||
ergebnis = {
|
||
"disc_type": (DISC_CACHE.get(request.device_path) or {}).get("disc_type", "Blu-ray"),
|
||
"title": request.title,
|
||
"year": request.year,
|
||
"confidence": 0.99,
|
||
"metadata": {
|
||
"type": request.type or "movie",
|
||
"id": request.id or "",
|
||
"title": request.title,
|
||
"year": request.year,
|
||
"overview": request.overview or "",
|
||
"poster_path": request.poster or "",
|
||
"backdrop_path": "",
|
||
"runtime": 0,
|
||
"genres": [],
|
||
"source": request.source or "manuell",
|
||
},
|
||
"tracks": [],
|
||
}
|
||
abdruck = disc_fingerprint(request.device_path)
|
||
cache_setter(
|
||
generate_prescan_key(request.device_path, False, abdruck),
|
||
ergebnis, expire=30 * 86400,
|
||
)
|
||
DISC_CACHE[request.device_path] = ergebnis
|
||
return ergebnis
|
||
|
||
ergebnis = await asyncio.to_thread(speichere)
|
||
await asyncio.to_thread(
|
||
db.add_log, "success", "api",
|
||
f"Metadaten manuell festgelegt: {request.title}"
|
||
+ (f" ({request.year})" if request.year else ""),
|
||
)
|
||
return ergebnis
|
||
|
||
|
||
# Pre-Scan Endpoint
|
||
class PreScanRequest(BaseModel):
|
||
device_path: str
|
||
|
||
|
||
@app.post("/prescan")
|
||
async def run_prescan(request: PreScanRequest):
|
||
"""Führe Pre-Scan durch."""
|
||
try:
|
||
prescan = PreScan()
|
||
result = prescan.scan(request.device_path)
|
||
return result.to_dict()
|
||
except Exception as e:
|
||
raise HTTPException(status_code=500, detail=str(e))
|
||
|
||
|
||
# Jellyfin-Formatierung Endpoints
|
||
class JellyfinFormatRequest(BaseModel):
|
||
title: str
|
||
year: Optional[int]
|
||
metadata: Dict
|
||
disc_type: str
|
||
output_dir: str
|
||
|
||
|
||
@app.post("/jellyfin/format")
|
||
async def jellyfin_format(request: JellyfinFormatRequest):
|
||
"""Formatiere für Jellyfin (NFO + Images)."""
|
||
try:
|
||
nfo_gen = NFOGenerator()
|
||
img_downloader = ImageDownloader()
|
||
|
||
# Ordnerstruktur erstellen
|
||
output_path = Path(request.output_dir)
|
||
|
||
if request.disc_type in ["dvd", "bluray"]:
|
||
# Film-Formatierung
|
||
title = request.metadata.get("title", request.title)
|
||
year = request.year or request.metadata.get("year")
|
||
|
||
# movie.nfo
|
||
movie_nfo = nfo_gen.generate_movie_nfo(
|
||
title=title,
|
||
year=year or 2000,
|
||
overview=request.metadata.get("overview", ""),
|
||
rating=request.metadata.get("rating", 0),
|
||
runtime=request.metadata.get("runtime", 0),
|
||
genres=request.metadata.get("genres", []),
|
||
director=request.metadata.get("director", ""),
|
||
actors=request.metadata.get("actors", [])
|
||
)
|
||
|
||
nfo_path = output_path / "movie.nfo"
|
||
nfo_gen.save_nfo(movie_nfo, nfo_path)
|
||
|
||
# Poster und Fanart
|
||
img_downloader.download_poster(title, output_path, 500)
|
||
img_downloader.download_fanart(title, output_path, 1920)
|
||
|
||
return {
|
||
"status": "formatted",
|
||
"nfo_path": str(nfo_path),
|
||
"poster_path": str(output_path / "poster.jpg"),
|
||
"fanart_path": str(output_path / "fanart.jpg")
|
||
}
|
||
else:
|
||
# Audio-Formatierung
|
||
artist = request.metadata.get("artist", "Unknown Artist")
|
||
album = request.metadata.get("title", request.title)
|
||
# Review-Fix 22.07.: `year` war hier undefiniert (existierte nur im Film-Zweig)
|
||
year = request.year or request.metadata.get("year")
|
||
|
||
# album.nfo
|
||
album_nfo = nfo_gen.generate_album_nfo(
|
||
title=album,
|
||
artist=artist,
|
||
year=year or 2000,
|
||
genres=request.metadata.get("genres", [])
|
||
)
|
||
|
||
nfo_path = output_path / "album.nfo"
|
||
nfo_gen.save_nfo(album_nfo, nfo_path)
|
||
|
||
# Album-Cover
|
||
img_downloader.download_music_images(artist, album, output_path)
|
||
|
||
return {
|
||
"status": "formatted",
|
||
"nfo_path": str(nfo_path),
|
||
"album_cover_path": str(output_path / "album.jpg")
|
||
}
|
||
except Exception as e:
|
||
raise HTTPException(status_code=500, detail=str(e))
|
||
|
||
|
||
# Auth-Endpoints (/token, /api-keys) entfernt — Commander-Entscheid 24.07.:
|
||
# Heimnetz-only, kein Login-Flow im UI, die Endpoints waren Placebo.
|