f449c4ee34
Ampel / ampel (push) Successful in 28s
Richtigstellung des Vortags-Befunds. Dort stand, MakeMKVs Schluessel-Kanal
sei abgeschaltet. Das war FALSCH: die Herleitung stuetzte sich auf zwei
Hostnamen aus alten Forumsbeitraegen (hkdata.fairuse.org,
hkdata.crabdance.com), die zwar wirklich nicht mehr aufloesen, von MakeMKV
aber laengst nicht mehr benutzt werden. Aufgedeckt durch den Einwand des
Commanders, unter Windows ginge es sofort.
Gegenprobe mit demselben Laufwerk und derselben Disc (Akira UHD, MKB v76):
Linux (Worker) Windows
Verbindungen KEINE EINZIGE 185.84.108.20:443
Meldung 3338 nie "Downloading latest HK"
_private_data.tar 2048 B, 0 Keys 6,4 MB, 604 Keys
Disc volume key unknown TCOUNT:5, geht auf
Gegengeprueft mit leerem UND gefuelltem Speicher, mit und ohne --noscan,
mit dev:/dev/sr0 und disc:0, mit geloeschter update.conf. Linux fragt nie.
Die Meldungsvorlage "Downloading latest %1 to %2 ..." steckt sehr wohl im
Linux-Binary - sie loest nur nicht aus. Gleiches Symptom im MakeMKV-Forum,
seit Jahren offen (t=25782, t=34022). Der Dienst lebt; der Worker erreicht
185.84.108.20:443 sogar problemlos.
BEWIESEN: Nach Uebernahme des Windows-Schluesselspeichers oeffnet
makemkvcon auf der VM die Akira-UHD - "Operation successfully completed",
TCOUNT:5, fuenf Titel, identisch zum Windows-Ergebnis. Erster belegter
UHD-Disc-Zugriff auf der Rippy-Maschine.
- makemkv_daten.py (beide Zwillinge): zaehle_schluessel,
private_data_pruefen, schluesselspeicher_status, private_data_schreiben.
Die Pruefung lehnt einen Speicher OHNE hkd_*.bin ab - sonst laedt jemand
den leeren Vorrat einer frischen Installation hoch, nichts aendert sich,
und niemand versteht warum. Modulkopf komplett neu, inkl. der
Fehldiagnose als Warnung fuer spaeter.
- API: GET/POST /system/keystore. Der Rohkoerper der Anfrage IST die Datei
(binaer - JSON/Base64 waere Ballast, Multipart kann die API nicht).
Groessengrenze 64 MB = client_max_body_size in nginx.conf.
- UI: neuer Block "Disc-Schluessel fuer 4K-UHD" UEBER dem KEYDB-Block, mit
Schluessel-Anzahl, Upload und Anleitung fuer den Windows-Weg. KEYDB.cfg
ist jetzt als Notnagel beschriftet. Worker-Plakette zeigt die Anzahl;
0 heisst sichtbar "4K-UHD scheitert".
- tasks.py: UHD-Fehlertext sagt den Windows-Weg an und nennt die Anzahl
bekannter Schluessel dieses Workers.
- caps.py meldet schluessel je Worker.
- Alle Falschaussagen korrigiert: UI (3), Anleitung (2), README (3),
KONZEPT §8 + §10, Worker-Dockerfile, makemkv_key.py (dort stand "Den
AACS-Schluessel zieht MakeMKV via LibreDrive ohnehin selbst aus dem
Laufwerk" - gilt fuer Blu-ray, NICHT fuer UHD).
- SAVEPOINT v3.11, ROADMAP Etappe 18 (Etappe 17 mit Nachtrag), AGENTS.
Offen: voller UHD-Rip inkl. Transcode-E2E; und ob sich der Abruf unter
Linux doch anstossen laesst.
Quellen (AGENTS Regel D):
- Linux laedt keine Hashed Keys, gleiches Symptom:
https://forum.makemkv.com/forum/viewtopic.php?t=25782
https://forum.makemkv.com/forum/viewtopic.php?t=34022
- Schluessel als hkd_*.bin in _private_data.tar:
https://forum.makemkv.com/forum/viewtopic.php?t=32675
- Meldungsformat: https://www.makemkv.com/developers/usage.txt
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
1703 lines
66 KiB
Python
1703 lines
66 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_daten
|
||
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)
|
||
|
||
|
||
# Eigene Wurzel für die Datei-Härtung der AACS-Dumps. Bewusst NICHT die
|
||
# MEDIA_ROOT-Helfer (_sicherer_dateiname/_validiere_ziel/_job_ausgabeordner):
|
||
# die prüfen hart gegen /app/media und würden hier IMMER 404 liefern.
|
||
# Das MakeMKV-Datenverzeichnis liegt woanders (in der API auf
|
||
# /app/makemkv-data, im Worker auf /root/.MakeMKV — laut docker-compose.yml
|
||
# beides dasselbe Host-Verzeichnis).
|
||
MAKEMKV_DATA_ROOT = os.path.realpath(makemkv_daten.DATEN_DIR)
|
||
|
||
|
||
class KeydbRequest(BaseModel):
|
||
inhalt: str # voller Text der KEYDB.cfg (kein Upload — es gibt kein python-multipart)
|
||
|
||
|
||
@app.get("/system/keydb")
|
||
async def get_keydb_status():
|
||
"""Was liegt gerade als KEYDB.cfg im MakeMKV-Datenverzeichnis?
|
||
|
||
Hintergrund (Befund 25.07.2026, live auf der VM nachgemessen): Bei
|
||
4K-UHD-Discs meldet MakeMKV "The volume key is unknown for this disc" und
|
||
holt den Schlüssel NICHT mehr online nach — die dokumentierten
|
||
Schlüssel-Server lösen weltweit nicht mehr auf. Der einzige heute
|
||
funktionierende Weg ist eine KEYDB.cfg, die der Nutzer selbst mitbringt.
|
||
Rippy liefert KEINE Schlüssel mit, lädt keine herunter und verteilt keine —
|
||
es stellt nur den Platz bereit und zeigt ehrlich an, was dort liegt.
|
||
|
||
Fehlendes Verzeichnis oder fehlende Datei ist der NORMALFALL: dann kommt
|
||
200 mit vorhanden=false zurück, niemals 404 oder 500.
|
||
"""
|
||
def sammle():
|
||
return makemkv_daten.keydb_status()
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@app.post("/system/keydb")
|
||
async def set_keydb(request: KeydbRequest):
|
||
"""Legt die vom Nutzer mitgebrachte KEYDB.cfg ab (atomar, ersetzt die alte).
|
||
|
||
WICHTIG für die Ehrlichkeit: Die Datei wirkt erst beim NÄCHSTEN Rip —
|
||
makemkvcon liest sie beim Prozessstart, ein bereits laufender Rip merkt
|
||
nichts davon. Genau so steht es auch im Log-Eintrag.
|
||
"""
|
||
# Reine Prüfung (kein Dateisystem) — fängt den häufigsten Bedienfehler ab:
|
||
# statt der KEYDB.cfg landet die HTML-Fehlerseite eines Downloads im Feld.
|
||
fehler = makemkv_daten.keydb_pruefen(request.inhalt)
|
||
if fehler:
|
||
raise HTTPException(status_code=422, detail=fehler)
|
||
|
||
def schreibe():
|
||
return makemkv_daten.keydb_schreiben(request.inhalt)
|
||
|
||
try:
|
||
status = await asyncio.to_thread(schreibe)
|
||
except OSError as e:
|
||
raise HTTPException(
|
||
status_code=500,
|
||
detail=(
|
||
f"KEYDB.cfg konnte nicht geschrieben werden: {e}. "
|
||
"Prüfe, ob das MakeMKV-Datenverzeichnis auf der VM existiert und "
|
||
"beschreibbar ist (Standard: /srv/rippy/makemkv)."
|
||
),
|
||
)
|
||
await asyncio.to_thread(
|
||
db.add_log, "success", "makemkv-keydb",
|
||
f"KEYDB.cfg abgelegt: {status['eintraege']} Zeilen mit Disc-Kennung, "
|
||
f"{status['groesse_bytes']} Bytes ({status['pfad']}). "
|
||
"Wirkt erst beim NÄCHSTEN Rip — MakeMKV liest die Datei beim Start.",
|
||
)
|
||
return status
|
||
|
||
|
||
@app.delete("/system/keydb")
|
||
async def delete_keydb():
|
||
"""Entfernt die KEYDB.cfg (z. B. nach einem Fehlgriff beim Einfügen).
|
||
|
||
Auch hier gilt: Die Änderung wirkt erst beim NÄCHSTEN Rip. Fehlt die Datei
|
||
schon, ist das kein Fehler — es kommt derselbe Zustand mit vorhanden=false.
|
||
"""
|
||
def loesche():
|
||
return makemkv_daten.keydb_loeschen()
|
||
|
||
try:
|
||
status = await asyncio.to_thread(loesche)
|
||
except OSError as e:
|
||
raise HTTPException(
|
||
status_code=500,
|
||
detail=f"KEYDB.cfg konnte nicht entfernt werden: {e}",
|
||
)
|
||
await asyncio.to_thread(
|
||
db.add_log, "warning", "makemkv-keydb",
|
||
"KEYDB.cfg entfernt. Ab dem NÄCHSTEN Rip fehlen die selbst mitgebrachten "
|
||
"Schlüssel wieder — UHD-Discs können dann erneut an "
|
||
"'The volume key is unknown for this disc' scheitern.",
|
||
)
|
||
return status
|
||
|
||
|
||
@app.get("/system/keystore")
|
||
async def get_keystore():
|
||
"""Wie viele Disc-Schlüssel kennt diese Rippy-Installation?
|
||
|
||
Der Schlüsselspeicher (_private_data.tar) ist MakeMKVs eigener Vorrat.
|
||
Unter Windows füllt MakeMKV ihn selbst; unter Linux nie — deshalb muss er
|
||
hier von Hand hereingereicht werden (Befund 25.07.2026, siehe
|
||
makemkv_daten.py). Fehlt er, ist das der Normalfall: 200 mit
|
||
vorhanden=false, nie 404.
|
||
"""
|
||
def sammle():
|
||
return makemkv_daten.schluesselspeicher_status()
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@app.post("/system/keystore")
|
||
async def set_keystore(request: Request):
|
||
"""Nimmt den Schlüsselspeicher einer MakeMKV-Installation entgegen.
|
||
|
||
Der Rohkörper der Anfrage IST die Datei — bewusst kein Multipart-Upload
|
||
(python-multipart fehlt) und bewusst kein JSON: _private_data.tar ist
|
||
binär, und Base64 würde sie nur unnötig aufblähen.
|
||
|
||
Wirkt ab dem NÄCHSTEN Rip: makemkvcon liest den Speicher beim Start.
|
||
"""
|
||
rohdaten = await request.body()
|
||
|
||
def pruefe_und_schreibe():
|
||
# Prüfung liest ein mehrere MB grosses tar — gehört deshalb mit in
|
||
# den Thread und nicht in die Ereignisschleife.
|
||
fehler = makemkv_daten.private_data_pruefen(rohdaten)
|
||
if fehler:
|
||
return fehler, None
|
||
return "", makemkv_daten.private_data_schreiben(rohdaten)
|
||
|
||
try:
|
||
fehler, status = await asyncio.to_thread(pruefe_und_schreibe)
|
||
except OSError as e:
|
||
raise HTTPException(
|
||
status_code=500,
|
||
detail=(
|
||
f"Schlüsselspeicher konnte nicht geschrieben werden: {e}. "
|
||
"Prüfe, ob das MakeMKV-Datenverzeichnis auf der VM existiert "
|
||
"und beschreibbar ist (Standard: /srv/rippy/makemkv)."
|
||
),
|
||
)
|
||
if fehler:
|
||
raise HTTPException(status_code=422, detail=fehler)
|
||
await asyncio.to_thread(
|
||
db.add_log, "success", "makemkv-keydb",
|
||
f"Schlüsselspeicher übernommen: {status['schluessel']} Disc-Schlüssel, "
|
||
f"{status['groesse_bytes']} Bytes. Wirkt ab dem NÄCHSTEN Rip.",
|
||
)
|
||
return status
|
||
|
||
|
||
@app.get("/system/aacs-dumps")
|
||
async def get_aacs_dumps():
|
||
"""AACS-Dumps, die MakeMKV selbst abgelegt hat (neueste zuerst).
|
||
|
||
MakeMKV schreibt sie beim gescheiterten UHD-Versuch ins Datenverzeichnis
|
||
(Meldung 3332 "Saved AACS dump file as file:///root/.MakeMKV/<name>.tgz",
|
||
am 25.07.2026 so beobachtet). Rippy wertet sie nicht aus und schickt sie
|
||
nirgendwohin — es zeigt nur, dass sie da sind, damit der Nutzer selbst
|
||
entscheiden kann, was er damit tut.
|
||
"""
|
||
def liste():
|
||
return {"dumps": makemkv_daten.dumps_auflisten()}
|
||
|
||
return await asyncio.to_thread(liste)
|
||
|
||
|
||
@app.get("/system/aacs-dumps/{dateiname}")
|
||
async def download_aacs_dump(dateiname: str):
|
||
"""Lädt EINEN AACS-Dump herunter.
|
||
|
||
Pfad-Validierung genauso streng wie beim Job-Datei-Download: nackter Name
|
||
ohne Pfadtrenner und ohne führenden Punkt (ist_aacs_dump) PLUS realpath,
|
||
der das MakeMKV-Datenverzeichnis nicht verlassen darf (kein ..-Ausbruch,
|
||
kein Symlink nach draußen).
|
||
"""
|
||
if not makemkv_daten.ist_aacs_dump(dateiname):
|
||
raise HTTPException(
|
||
status_code=404,
|
||
detail="Kein gültiger Dump-Name — erwartet wird eine .tgz-Datei ohne Pfadangabe.",
|
||
)
|
||
pfad = os.path.join(MAKEMKV_DATA_ROOT, dateiname)
|
||
|
||
def pruefe():
|
||
return os.path.isfile(pfad) and os.path.realpath(pfad).startswith(MAKEMKV_DATA_ROOT)
|
||
|
||
if not await asyncio.to_thread(pruefe):
|
||
raise HTTPException(
|
||
status_code=404,
|
||
detail="Dump nicht gefunden — MakeMKV legt ihn erst beim gescheiterten UHD-Versuch an.",
|
||
)
|
||
return FileResponse(pfad, filename=dateiname, media_type="application/gzip")
|
||
|
||
|
||
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.
|