19f3dc5330
Ampel / ampel (push) Successful in 30s
Der Commander: "Ausserdem ist gerade mitten im Rip das Laufwerk ausgegangen...
glaube ich zumindest." Gemessen war es etwas anderes, und der Fund ist groesser
als der Vorfall.
WAS MESSBAR WAR:
Job 2182d525 status=running progress=12 (Rip gestartet 14:00:24)
makemkvcon laeuft NICHT (per /proc geprueft, PID gegengeprueft)
Rohdatei 5.167.382.528 Bytes, waechst in 10 s nicht
Laufwerk Status 4 (Disc drin), /dev/sr0 + /dev/sg1 da
Worker-Start 14:07:02 <- mein `docker compose up -d --build`
Das Laufwerk ist also NICHT ausgegangen. Der Rip wurde von MEINEM Deploy
getoetet: `up -d --build` baut den worker-Container neu, und der laufende Rip
stirbt mit ihm. Genau davor warnt der SAVEPOINT seit v3.18 - die Warnung half
nichts, weil sie niemand liest und nichts sie prueft.
DER EIGENTLICHE FUND: Die Zombie-Erkennung lief um 14:09:04 und meldete
`{'geprueft': 0, 'aufgeraeumt': []}` - obwohl der tote Job direkt vor ihr lag.
Ursache:
ARBEITS_STATI = ("ripping", "transcoding", "canceling") # zombies.py
db.update_job(job_id, status="running", ...) # tasks.py - der Rip
Der Rip setzt "running", gesucht wurde "ripping". Dieser Wert steht
ausschliesslich in Celerys Task-META und NIE in einer Job-Zeile (nachgeprueft:
kein einziger Schreiber im ganzen Baum). Die Zombie-Erkennung aus v3.14 wurde
gebaut, um genau einen abgestuerzten Rip zu finden - und hat ihn nie gesehen.
Besonders tueckisch: `geprueft: 0` sah bei jedem Worker-Start wie "nachgesehen,
alles gesund" aus, waehrend sie nach einem Status suchte, den es nicht gibt.
Deshalb blieb der Job auf "processing 12 %" stehen - mit einer Restzeit-Schaetzung
von 1 h 30 min obendrauf, die es fuer einen toten Prozess nicht geben duerfte.
GEBAUT:
* "running" in ARBEITS_STATI.
* Ein Test, der das Auseinanderlaufen MECHANISCH verhindert: Er liest tasks.py,
sammelt jeden Status, den der Worker per db.update_job in eine Job-Zeile
schreibt, und verlangt, dass jeder davon entweder ein Arbeitsstatus oder ein
Endzustand ist. Ein Kommentar haette das nicht verhindert.
* GET /health/arbeit + eine Sperre in deploy.sh: Laeuft ein Job, bricht der
Deploy ab (uebersteuerbar mit RIPPY_TROTZDEM=1 - dann ist es eine
Entscheidung und kein Versehen). Ein Satz Code gegen eine verlorene Stunde.
* Server-Status zeigt jetzt die eingehaengten FREIGABEN mit freiem Platz, nicht
nur die Container-Platte (Commander-Wunsch). Genau dort liegen die Rohdaten,
und bei externem Encoden muessen sie dort liegen - wer wissen wollte, ob noch
Platz fuer eine Disc ist, sah die falsche Zahl.
Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
2101 lines
85 KiB
Python
2101 lines
85 KiB
Python
from fastapi import FastAPI, HTTPException, Request, Response
|
||
from fastapi.middleware.cors import CORSMiddleware
|
||
from fastapi.responses import FileResponse
|
||
from pydantic import BaseModel
|
||
from typing import List, Optional, Dict
|
||
import asyncio
|
||
import json
|
||
import os
|
||
import posixpath
|
||
import shutil
|
||
import time
|
||
import uuid
|
||
|
||
import db
|
||
import devices as device_discovery
|
||
import eta
|
||
import makemkv_daten
|
||
import makemkv_key
|
||
import mounts as mount_verwaltung
|
||
import notify
|
||
import presets as preset_auswahl
|
||
import rohdaten
|
||
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 get as cache_get, init_cache, set as cache_set
|
||
from ratelimit import check_rate_limit, get_rate_limit_remaining
|
||
from prescan import PreScan
|
||
|
||
# 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, läuft zum
|
||
# Monatsende ab) — sonst blockt Blu-ray-Ripping irgendwann still. Taeglicher
|
||
# Forum-Abgleich; wirkt ohne Rebuild ab dem nächsten Rip.
|
||
asyncio.create_task(makemkv_key.refresh_loop())
|
||
# Worker-Erreichbarkeit im Hintergrund pingen (siehe _ping_knoten) — sonst
|
||
# kostet JEDER Aufruf von /capabilities eine ganze Sekunde.
|
||
asyncio.create_task(_ping_schleife())
|
||
# Aus demselben Grund im Hintergrund: nachsehen, wo Rohdaten liegen. Ein
|
||
# schlafendes NAS lässt os.path.isdir bis zum CIFS-Timeout hängen (hier
|
||
# 10 s gemessen) — und /jobs wird alle 4 Sekunden abgefragt.
|
||
asyncio.create_task(_rohdaten_schleife())
|
||
# Freigaben bewachen: Nach einem Rebuild ist die CIFS-Freigabe reproduzierbar
|
||
# tot, und der zweite Anlauf hilft (Begründung bei _mount_schleife).
|
||
asyncio.create_task(_mount_schleife())
|
||
|
||
|
||
# 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 unter_wurzel(ziel, 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
|
||
# Restzeit-Schätzung (siehe eta.py). -1/"" heißt „noch keine Aussage" —
|
||
# bewusst ehrlich statt einer erfundenen Minutenzahl.
|
||
eta_sekunden: int = -1
|
||
eta_text: str = ""
|
||
|
||
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("/health/arbeit")
|
||
async def health_arbeit():
|
||
"""Läuft gerade ein Job? — die Frage VOR einem Deploy.
|
||
|
||
⚠️ Warum es das gibt (26.07.2026, selbst verursacht): Ein
|
||
`docker compose up -d --build` baut den worker-Container neu und tötet damit
|
||
einen laufenden Rip. Genau das ist passiert — mitten in einem Blu-ray-Rip,
|
||
bei 12 %, nach 5,1 GB. Im SAVEPOINT stand die Warnung „nicht deployen,
|
||
während ein Rip läuft" schon; sie half nichts, weil niemand sie las und
|
||
nichts sie prüfte.
|
||
|
||
Jetzt fragt `deploy.sh` hier nach und bricht ab. Ein Satz Code gegen eine
|
||
verlorene Stunde.
|
||
"""
|
||
def sammle():
|
||
laufend = [
|
||
{"id": z["id"], "status": z["status"], "progress": z.get("progress") or 0,
|
||
"titel": z.get("title") or ""}
|
||
for z in db.list_jobs()
|
||
if z.get("status") in ("pending", "running", "ripping",
|
||
"transcoding", "canceling")
|
||
]
|
||
return {"arbeit": bool(laufend), "jobs": laufend}
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@app.get("/health/vorraete")
|
||
async def health_vorraete():
|
||
"""Laufen die Hintergrund-Schleifen wirklich? (Diagnose, kein UI-Endpunkt)
|
||
|
||
Zwei Endpunkte hängen an einem Vorrat, den eine Hintergrund-Schleife füllt:
|
||
/capabilities am Celery-Ping, /jobs an der Rohdaten-Suche. Bleibt so ein
|
||
Vorrat leer, ist am Endpunkt selbst NICHTS zu sehen — er antwortet nur
|
||
dauerhaft „nichts gefunden". Genau daran ging am 26.07.2026 eine Stunde
|
||
verloren. `alter_sekunden` sagt, wann die Schleife das letzte Mal
|
||
durchgelaufen ist; wächst der Wert über das Intervall, läuft sie nicht mehr.
|
||
"""
|
||
jetzt = time.monotonic()
|
||
return {
|
||
"ping": {
|
||
"knoten": _PING["knoten"],
|
||
"alter_sekunden": None if _PING["stand"] < 0 else round(jetzt - _PING["stand"], 1),
|
||
"intervall": PING_INTERVALL_SEKUNDEN,
|
||
},
|
||
"rohdaten": {
|
||
"jobs_im_vorrat": len(_ROHDATEN["treffer"]),
|
||
"mit_treffer": sum(1 for v in _ROHDATEN["treffer"].values() if v),
|
||
"alter_sekunden": None if _ROHDATEN["stand"] is None
|
||
else round(jetzt - _ROHDATEN["stand"], 1),
|
||
"intervall": ROHDATEN_INTERVALL_SEKUNDEN,
|
||
},
|
||
"mounts": {
|
||
"stand": _MOUNT_STAND,
|
||
"intervall": MOUNT_WACHE_INTERVALL_SEKUNDEN,
|
||
},
|
||
}
|
||
|
||
|
||
@app.get("/")
|
||
async def root():
|
||
return {
|
||
"name": "Rippy",
|
||
"version": "1.0.0",
|
||
"description": "Automatisches Ripping-System für CD, DVD und Blu-ray"
|
||
}
|
||
|
||
|
||
def _rohdaten_suchen(job_id: str, work_dir: str) -> list:
|
||
"""Wo liegen die Roh-MKVs dieses Jobs? (Details in rohdaten.py)
|
||
|
||
Geprüft wird mit `rohdaten.verzeichnis_da` und NICHT mit os.path.isdir:
|
||
Letzteres hing am 26.07.2026 unbegrenzt im Kernel, während Rippy die
|
||
CIFS-Freigabe neu einhängte — und nahm die Hintergrund-Schleife dauerhaft
|
||
mit. Ausführliche Begründung dort.
|
||
|
||
Bleibt trotzdem der langsame Weg (mehrere Kind-Prozesse, im schlechtesten
|
||
Fall je vier Sekunden). Für die Job-Liste, die das Dashboard alle vier
|
||
Sekunden abfragt, gibt es deshalb den Vorrat unten.
|
||
"""
|
||
return rohdaten.suche(job_id, work_dir, os.listdir, rohdaten.verzeichnis_da)
|
||
|
||
|
||
# Vorrat für die Job-Liste. Dasselbe Muster wie beim Celery-Ping in
|
||
# /capabilities (v3.15): Der Endpunkt wird alle 4 Sekunden vom Dashboard
|
||
# abgefragt und darf NIE am Dateisystem hängen. Ein schlafendes NAS hätte das
|
||
# Dashboard sonst für 10 Sekunden je Aufruf eingefroren.
|
||
_ROHDATEN = {"treffer": {}, "stand": None}
|
||
ROHDATEN_INTERVALL_SEKUNDEN = 30
|
||
|
||
|
||
# --- Mount-Wache ------------------------------------------------------------
|
||
#
|
||
# ⚠️ Warum es die braucht (26.07.2026, DREIMAL in Folge reproduziert): Nach
|
||
# `docker compose up -d --build` ist die CIFS-Freigabe tot. `mount` meldet
|
||
# Rückgabewert 0, /proc/mounts zeigt genau eine korrekt aussehende Schicht, die
|
||
# Erreichbarkeits-Probe antwortet direkt nach dem Mount sogar — und Sekunden
|
||
# später läuft jeder Zugriff in die Zeitgrenze. Derselbe Ablauf ein zweites Mal,
|
||
# ein bis zwei Minuten später, stellt sie zuverlässig her.
|
||
#
|
||
# ## DIE URSACHE, gemessen am 26.07.2026
|
||
#
|
||
# /proc/fs/cifs/DebugData → Net namespace: 4026532653
|
||
# api-Container → net:[4026532653] ← dieselbe
|
||
# worker-Container → net:[4026532540] ← andere
|
||
#
|
||
# Die CIFS-Verbindung lebt in der NETZ-NAMESPACE DES API-CONTAINERS — hier wird
|
||
# sie eingehängt (nur dieser Container hat CAP_SYS_ADMIN). Wird der Container neu
|
||
# gebaut, stirbt sein Netz-Namespace und damit der Socket. Der Mount steht danach
|
||
# weiter in /proc/mounts (er ist per rshared auf den Host propagiert) und sieht
|
||
# vollkommen gesund aus — aber jeder Zugriff läuft in den CIFS-Timeout.
|
||
#
|
||
# Deshalb passiert es nach JEDEM Deploy, deshalb sieht `mount` gesund aus, und
|
||
# deshalb hilft nur ein echtes Neu-Verbinden aus dem neuen Container heraus.
|
||
#
|
||
# Und deshalb ist der api-Container die einzige Stelle, die die NAS-Verbindung
|
||
# hält: Startet er mitten in einem Rip neu, verliert auch der Worker sein Ziel.
|
||
# Das strukturell zu lösen (Mount auf dem HOST statt im Container) wäre ein
|
||
# eigener Umbau und widerspräche „Speicherziele über das UI einhängen".
|
||
MOUNT_WACHE_INTERVALL_SEKUNDEN = 60
|
||
# Erste Prüfung fast sofort: Genau nach einem Deploy ist die Lage kaputt, und
|
||
# jede Sekunde Wartezeit ist eine Sekunde, in der Rippy sein Ziel nicht sieht.
|
||
MOUNT_WACHE_ERSTE_PRUEFUNG_SEKUNDEN = 3
|
||
_MOUNT_STAND = {}
|
||
|
||
|
||
async def _mount_schleife():
|
||
"""Sieht nach, ob die Freigaben antworten, und verbindet sie sonst neu."""
|
||
await asyncio.sleep(MOUNT_WACHE_ERSTE_PRUEFUNG_SEKUNDEN)
|
||
while True:
|
||
try:
|
||
await asyncio.to_thread(_mounts_nachsehen)
|
||
except Exception as e: # darf nie sterben
|
||
print(f"Mount-Wache fehlgeschlagen: {type(e).__name__}: {e}")
|
||
await asyncio.sleep(MOUNT_WACHE_INTERVALL_SEKUNDEN)
|
||
|
||
|
||
def _mounts_nachsehen() -> None:
|
||
"""Unerreichbare Freigaben neu verbinden — NIE während ein Job läuft.
|
||
|
||
Neu verbinden heißt `umount -l`; mitten in einem Rip oder Encode wäre das
|
||
ein Datenverlust. Deshalb steht die Wache still, solange irgendein Job nicht
|
||
durch ist — auch bei einem wartenden, der jeden Moment anlaufen kann.
|
||
"""
|
||
eintraege = db.list_mounts()
|
||
if not eintraege or db.hat_arbeit():
|
||
return
|
||
for eintrag in eintraege:
|
||
name = eintrag["name"]
|
||
# Zum ERKENNEN genügt die einfache, schnelle Probe: Ein toter Mount
|
||
# antwortet gar nicht, nicht nur manchmal. Die Doppelprobe steckt dort,
|
||
# wo sie hingehört — in `mounten()`, direkt nach einem frischen Mount, wo
|
||
# der Wettlauf mit dem lazy umount lauert. Hier kostete sie nur jede
|
||
# Minute drei Sekunden Warten für nichts.
|
||
erreichbar = mount_verwaltung.ist_erreichbar(name)
|
||
vorher = _MOUNT_STAND.get(name)
|
||
_MOUNT_STAND[name] = erreichbar
|
||
if erreichbar:
|
||
if vorher is False:
|
||
db.add_log("success", "mounts", f"{name}: antwortet wieder")
|
||
continue
|
||
# Nur beim ÜBERGANG meckern, nicht jede Minute: Ist das NAS
|
||
# ausgeschaltet, wäre das sonst ein Log-Wasserfall.
|
||
if vorher is not False:
|
||
db.add_log("warning", "mounts",
|
||
f"{name}: antwortet nicht — wird neu verbunden")
|
||
begonnen = time.monotonic()
|
||
try:
|
||
mount_verwaltung.reparieren(
|
||
name, eintrag["typ"], eintrag["quelle"],
|
||
eintrag.get("optionen") or "", eintrag.get("username") or "",
|
||
eintrag.get("passwort") or "",
|
||
)
|
||
_MOUNT_STAND[name] = True
|
||
# Die DAUER mit ins Log: Sie war die entscheidende Spur, als eine
|
||
# Wiederanbindung drei Minuten brauchte (ein `stat` auf dem toten
|
||
# Mount). Steht sie da, muss man sie nicht erst rekonstruieren.
|
||
db.add_log("success", "mounts",
|
||
f"{name}: neu verbunden ({time.monotonic() - begonnen:.1f}s)")
|
||
except Exception as e:
|
||
db.add_log(
|
||
"warning", "mounts",
|
||
f"{name}: Neuverbinden fehlgeschlagen nach "
|
||
f"{time.monotonic() - begonnen:.1f}s — {e}")
|
||
|
||
|
||
async def _rohdaten_schleife():
|
||
"""Hält den Rohdaten-Vorrat frisch. Darf nie sterben.
|
||
|
||
Der Fehler wird GEMELDET, nicht verschluckt: Ein `except Exception: pass`
|
||
stand hier zuerst, und als der Vorrat leer blieb, war nicht zu sehen, warum
|
||
(Sitzung 26.07.2026 — eine Stunde Rätselraten für eine Zeile Log).
|
||
"""
|
||
while True:
|
||
try:
|
||
await asyncio.to_thread(_rohdaten_vorrat_auffrischen)
|
||
except Exception as e: # DB/NAS weg → beim nächsten Durchlauf erneut
|
||
print(f"Rohdaten-Vorrat fehlgeschlagen: {type(e).__name__}: {e}")
|
||
await asyncio.sleep(ROHDATEN_INTERVALL_SEKUNDEN)
|
||
|
||
|
||
def _rohdaten_vorrat_auffrischen() -> None:
|
||
"""Für jeden fehlgeschlagenen Job nachsehen, wo seine Rohdaten liegen.
|
||
|
||
„Konnte nicht nachsehen" behält die letzte bekannte Antwort: Nach einem
|
||
Container-Neustart stallt der erste Zugriff auf die CIFS-Freigabe mehrere
|
||
Sekunden. Ohne diese Regel verschwände in dem Fenster der Knopf
|
||
„Neu komprimieren", und der Nutzer schlösse daraus, seine 74 GB seien weg.
|
||
"""
|
||
work_dir = os.path.normpath((db.get_settings().get("workDir") or "").strip() or "/")
|
||
alt = _ROHDATEN["treffer"]
|
||
treffer = {}
|
||
for zeile in db.list_jobs():
|
||
if zeile.get("status") != "failed":
|
||
continue
|
||
job_id = zeile["id"]
|
||
ergebnis = rohdaten.suche_mit_status(job_id, work_dir, os.listdir)
|
||
if not ergebnis["pfade"] and ergebnis["unklar"] and alt.get(job_id):
|
||
treffer[job_id] = alt[job_id] # letzte bekannte Antwort halten
|
||
else:
|
||
treffer[job_id] = ergebnis["pfade"]
|
||
_ROHDATEN["treffer"] = treffer
|
||
_ROHDATEN["stand"] = time.monotonic()
|
||
|
||
|
||
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.).
|
||
|
||
⚠️ Reparatur 26.07.2026, an der laufenden Instanz gemessen: Gesucht wurde
|
||
nur im Container-Standard und unter dem AKTUELLEN `workDir`. Der Rip von Job
|
||
95afdc89 lag aber auf der NAS, weil beim Start eine Wahl NUR FÜR DIESEN RIP
|
||
getroffen worden war (gibt es seit v3.15) — und die Einstellung selbst stand
|
||
auf leer. Ergebnis: `can_retry` war `false`, obwohl 79,6 GB intakt dalagen.
|
||
Der SAVEPOINT v3.16 behauptete „‚Neu komprimieren' genügt" — den Knopf gab
|
||
es nicht. Jetzt wird an allen möglichen Orten nachgesehen.
|
||
|
||
Gelesen wird aus dem Vorrat, nicht live: Diese Funktion hängt an /jobs, und
|
||
das fragt das Dashboard alle 4 Sekunden. Solange der Vorrat einen Job noch
|
||
nicht kennt (frischer Fehlschlag), zählt der billige lokale Ort — der liegt
|
||
auf der Container-Platte und antwortet immer sofort.
|
||
"""
|
||
if job.get("status") != "failed":
|
||
return False
|
||
vorrat = _ROHDATEN["treffer"]
|
||
if job["id"] in vorrat:
|
||
return bool(vorrat[job["id"]])
|
||
# Nur der lokale Ort: /app/temp ist ein Docker-Volume, os.path.isdir kann
|
||
# dort nicht hängen (im Gegensatz zu allem unter /app/media).
|
||
return os.path.isdir(posixpath.join(rohdaten.RAW_STANDARD, job["id"]))
|
||
|
||
|
||
@app.get("/jobs", response_model=List[Job])
|
||
async def get_jobs():
|
||
"""Holt alle Jobs aus der Datenbank (neueste zuerst) — inkl. Restzeit.
|
||
|
||
Die ETA entsteht HIER, weil hier der Fortschritt vorbeikommt: jeder Aufruf
|
||
schreibt die Messreihe im Cache fort (eta.py). Im Browser zu rechnen wäre
|
||
einfacher gewesen und dreifach schlechter — ein Seitenwechsel setzte die
|
||
Reihe zurück, zwei Tabs zeigten verschiedene Zahlen, und für einen externen
|
||
Encoder-Worker gäbe es gar keine (genau dort wollte der Commander sie).
|
||
"""
|
||
def sammle():
|
||
work_dir = os.path.normpath((db.get_settings().get("workDir") or "").strip() or "/")
|
||
jetzt = time.monotonic()
|
||
modelle = []
|
||
for z in db.list_jobs():
|
||
modell = _job_row_to_model(z)
|
||
modell.can_retry = _kann_neu_komprimieren(z, work_dir)
|
||
schaetzung = eta.aktualisiere_und_schaetze(
|
||
z["id"], z.get("status") or "", z.get("progress") or 0,
|
||
jetzt, cache_get, cache_set,
|
||
)
|
||
modell.eta_sekunden = schaetzung["sekunden"]
|
||
modell.eta_text = schaetzung["text"]
|
||
modelle.append(modell)
|
||
return modelle
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
def _rohdaten_groesse(pfade: list) -> tuple:
|
||
"""(Bytes, Dateizahl) der Roh-MKVs (Details in rohdaten.py).
|
||
|
||
Hier bleibt os.* stehen: Diese Funktion wird nur aufgerufen, NACHDEM
|
||
`verzeichnis_da` den Pfad innerhalb von Sekunden bestätigt hat — der Mount
|
||
antwortet also. Und sie läuft nur, wenn ein Mensch auf die Zahl wartet
|
||
(/rohdaten, Löschen), nie in der Hintergrund-Schleife.
|
||
"""
|
||
return rohdaten.groesse(pfade, os.listdir, os.path.isfile, os.path.getsize)
|
||
|
||
|
||
@app.get("/jobs/{job_id}/rohdaten")
|
||
async def job_rohdaten(job_id: str):
|
||
"""Was bleibt liegen, wenn dieser Job aus der Liste fliegt?
|
||
|
||
⚠️ Der Grund für diesen Endpunkt (Vorfall 25.07.2026, v3.14): „Job aus der
|
||
Liste entfernen" löscht bewusst keine Dateien — richtig, aber der Rohschnitt
|
||
ist danach UNERREICHBAR, weil Job und Dateien nur über die Job-ID verbunden
|
||
sind. Ein 75-GB-Rohschnitt verwaiste so unsichtbar auf der Platte: kein
|
||
Eintrag zeigte mehr darauf, „Neu komprimieren" war unmöglich, und im UI war
|
||
nichts davon zu sehen. Erst ein Blick per SSH brachte es zutage.
|
||
|
||
Jetzt fragt das UI vor dem Entfernen hier nach und sagt die Zahl.
|
||
"""
|
||
job = await asyncio.to_thread(db.get_job, job_id)
|
||
if not job:
|
||
raise HTTPException(status_code=404, detail="Job nicht gefunden")
|
||
|
||
def sammle():
|
||
work_dir = os.path.normpath((db.get_settings().get("workDir") or "").strip() or "/")
|
||
pfade = _rohdaten_suchen(job_id, work_dir)
|
||
bytes_gesamt, dateien = _rohdaten_groesse(pfade)
|
||
return {
|
||
"pfade": pfade,
|
||
"dateien": dateien,
|
||
"bytes": bytes_gesamt,
|
||
"gb": round(bytes_gesamt / 1024**3, 1),
|
||
}
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@app.delete("/jobs/{job_id}")
|
||
async def delete_job(job_id: str, rohdaten: bool = False):
|
||
"""Entfernt einen erledigten Job aus der Liste.
|
||
|
||
`rohdaten=true` löscht zusätzlich die Roh-MKVs. Ohne den Schalter bleiben
|
||
sie liegen (Standard wie bisher) — das UI nennt vorher die Größe, damit die
|
||
Entscheidung bewusst fällt und kein Rohschnitt unsichtbar verwaist.
|
||
"""
|
||
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")
|
||
|
||
geloescht_gb = 0.0
|
||
if rohdaten:
|
||
def raeume():
|
||
work_dir = os.path.normpath((db.get_settings().get("workDir") or "").strip() or "/")
|
||
pfade = _rohdaten_suchen(job_id, work_dir)
|
||
bytes_gesamt, _ = _rohdaten_groesse(pfade)
|
||
for pfad in pfade:
|
||
shutil.rmtree(pfad, ignore_errors=True)
|
||
return round(bytes_gesamt / 1024**3, 1)
|
||
|
||
geloescht_gb = await asyncio.to_thread(raeume)
|
||
|
||
await asyncio.to_thread(db.delete_job, job_id)
|
||
vermerk = f" (inkl. {geloescht_gb} GB Rohdaten gelöscht)" if rohdaten else ""
|
||
await asyncio.to_thread(
|
||
db.add_log, "info", "api", f"Job {job_id} aus der Liste entfernt{vermerk}"
|
||
)
|
||
return {"status": "deleted", "rohdaten_geloescht_gb": geloescht_gb}
|
||
|
||
|
||
@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)
|
||
# Sprachauswahl für DIESEN Rip (ISO-639-2, z. B. ["deu","eng"]).
|
||
# Leer = alles behalten. Greift bei der Kompression, nicht beim Rippen —
|
||
# der Rip bleibt vollständig und verlustfrei (Begründung in
|
||
# worker/ripping.build_handbrake_cmd).
|
||
audio_sprachen: Optional[List[str]] = None
|
||
untertitel_sprachen: Optional[List[str]] = None
|
||
# Arbeitsverzeichnis NUR für diesen Rip (Commander-Wunsch 25.07.2026:
|
||
# beim Start wählbar, nicht global vorgegeben). Leer = der Wert aus
|
||
# Einstellungen → Verarbeitung, der auch für Vollautomatik-Rips gilt.
|
||
work_dir: Optional[str] = None
|
||
|
||
|
||
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 unter_wurzel(normalisiert, 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
|
||
# Sprachauswahl dieses Rips. Nur schreiben, wenn wirklich gewählt wurde —
|
||
# eine leere Liste würde im Worker als „alles" gelesen, was derselbe Fall
|
||
# ist, aber die Absicht verschleiert.
|
||
for feld, wert in (("audio_sprachen", request.audio_sprachen),
|
||
("untertitel_sprachen", request.untertitel_sprachen)):
|
||
sauber = [str(s).strip().lower() for s in (wert or []) if str(s).strip()]
|
||
if sauber:
|
||
meta_dict[feld] = sauber
|
||
# Arbeitsverzeichnis dieses Rips. Dieselbe Pfad-Härte wie beim Ziel: muss
|
||
# unter /app/media liegen, damit man nicht versehentlich 100 GB Rohdaten
|
||
# irgendwohin in den Container schreibt.
|
||
arbeits_dir = _validiere_ziel(request.work_dir)
|
||
if arbeits_dir:
|
||
meta_dict["work_dir"] = arbeits_dir
|
||
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 unter_wurzel(pfad: str, wurzel: str) -> bool:
|
||
"""Liegt `pfad` wirklich unterhalb von `wurzel` (oder IST es die Wurzel)?
|
||
|
||
Pure Funktion, testbar. Ein nacktes `startswith()` genügt hier nicht:
|
||
„/app/media-boese/x" beginnt mit „/app/media", liegt aber außerhalb
|
||
(Befund 25.07.2026 bei der Durchsicht). Deshalb Gleichheit ODER Wurzel
|
||
samt Trennzeichen. Erwartet werden normalisierte Container-Pfade mit „/".
|
||
"""
|
||
if not pfad or not wurzel:
|
||
return False
|
||
sauber = wurzel.rstrip("/") or "/"
|
||
return pfad == sauber or pfad.startswith(sauber + "/")
|
||
|
||
|
||
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 unter_wurzel(ausgabe, 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 unter_wurzel(os.path.realpath(pfad), 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 irgendwo (bei
|
||
Kompressions-Fehlschlägen bleiben sie absichtlich erhalten).
|
||
|
||
⚠️ Reparatur 26.07.2026: Das Roh-Verzeichnis wurde hier aus dem AKTUELLEN
|
||
Wert von `workDir` errechnet. Wer beim Rip-Start eine andere Ablage gewählt
|
||
hatte (gibt es seit v3.15), bekam damit einen Pfad, an dem nichts liegt —
|
||
und der Worker brach mit „Verzeichnis erreichbar, enthält aber keine
|
||
MKV-Datei" ab. Beim Job 95afdc89 lagen 79,6 GB auf der NAS, gesucht wurde in
|
||
/app/temp/raw. Jetzt wird nachgesehen statt gerechnet (rohdaten.py).
|
||
"""
|
||
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")
|
||
|
||
einstellungen = await asyncio.to_thread(db.get_settings)
|
||
work_dir = os.path.normpath((einstellungen.get("workDir") or "").strip() or "/")
|
||
gefunden = await asyncio.to_thread(_rohdaten_suchen, job_id, work_dir)
|
||
if not gefunden:
|
||
raise HTTPException(
|
||
status_code=409,
|
||
detail=(
|
||
"Keine Rohdaten zu diesem Job gefunden — weder unter "
|
||
"/app/temp/raw noch in einem der Ablageziele. Ohne sie muss die "
|
||
"Disc neu gerippt werden. (Gelöscht? Freigabe nicht eingehängt?)"
|
||
),
|
||
)
|
||
raw_dir = gefunden[0]
|
||
# 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"}
|
||
|
||
|
||
# --- Worker-Erreichbarkeit: gepingt wird im Hintergrund, nicht im Request ---
|
||
#
|
||
# Befund 25.07.2026 (gemessen): /capabilities brauchte **1,010 s** — und zwar
|
||
# jedes Mal. Ursache ist kein Fehler, sondern das Wesen des Celery-Pings: er
|
||
# sammelt Antworten bis zum Timeout und kann nicht früher aufhören, weil er
|
||
# nicht weiß, wie viele Worker noch antworten wollen. Fünf UI-Stellen holen
|
||
# /capabilities (Dashboard, Einstellungen, Worker-Tab, Wizard, Rip-Dialog) —
|
||
# jede Seite zahlte also eine Sekunde, obwohl alle anderen Endpunkte unter
|
||
# 25 ms liegen. Genau das war das „Laggen".
|
||
#
|
||
# Jetzt pingt ein Hintergrund-Lauf im festen Takt, und der Endpunkt liest nur
|
||
# ab. Ist der Vorrat älter als PING_ALTER_MAX (Lauf noch nicht angelaufen oder
|
||
# gestorben), wird EINMAL synchron gepingt und der Vorrat wieder gefüllt —
|
||
# lieber eine langsame Antwort als eine falsche.
|
||
_PING = {"knoten": [], "stand": -1e9}
|
||
PING_INTERVALL_SEKUNDEN = 5
|
||
PING_ALTER_MAX_SEKUNDEN = 30
|
||
|
||
|
||
def _ping_jetzt() -> list:
|
||
"""Pingt sofort (blockiert ~1 s) und füllt den Vorrat."""
|
||
try:
|
||
antworten = celery_client.control.ping(timeout=1.0) or []
|
||
knoten = [k for antwort in antworten for k in antwort.keys()]
|
||
except Exception:
|
||
knoten = []
|
||
_PING["knoten"] = knoten
|
||
_PING["stand"] = time.monotonic()
|
||
return knoten
|
||
|
||
|
||
def _ping_knoten() -> list:
|
||
"""Erreichbare Celery-Knoten aus dem Vorrat — ohne zu warten."""
|
||
if time.monotonic() - _PING["stand"] <= PING_ALTER_MAX_SEKUNDEN:
|
||
return _PING["knoten"]
|
||
return _ping_jetzt()
|
||
|
||
|
||
async def _ping_schleife():
|
||
"""Hält den Ping-Vorrat frisch. Darf nie sterben, sonst wird jeder
|
||
/capabilities-Aufruf wieder langsam."""
|
||
while True:
|
||
try:
|
||
await asyncio.to_thread(_ping_jetzt)
|
||
except Exception: # Broker weg → beim nächsten Durchlauf erneut
|
||
pass
|
||
await asyncio.sleep(PING_INTERVALL_SEKUNDEN)
|
||
|
||
|
||
@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()
|
||
ping_knoten = _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.get("/presets")
|
||
async def preset_uebersicht():
|
||
"""Welche HandBrake-Presets es WIRKLICH gibt — und welches das beste ist.
|
||
|
||
Die Namen kommen von den Workern selbst (`HandBrakeCLI --preset-list`, siehe
|
||
worker/caps.py), die Staffelung aus presets.py. Damit endet das Raten:
|
||
vorher standen die Namen fest verdrahtet im UI, und ein Name, den das
|
||
jeweilige HandBrake nicht kennt, ließ die Kompression scheitern.
|
||
|
||
Commander-Anforderung 26.07.2026: „Bei den Presets soll IMMER das Beste
|
||
ausgewählt werden" — die Begründung steht bei jeder Empfehlung mit dabei.
|
||
"""
|
||
def sammle():
|
||
return preset_auswahl.uebersicht(db.list_workers())
|
||
|
||
return 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 unter_wurzel(normalisiert, 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 unter_wurzel(basis, 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")
|
||
|
||
|
||
# GET /worker-setup/windows-gui entfernt am 25.07.2026: ohne Aufrufer, seit die
|
||
# .exe den .bat-Umweg ersetzt hat (v3.9 — .vbs/.bat wird als gefährlich
|
||
# geflaggt, Commander-Einwand). Die Datei install-gui.ps1 selbst lebt weiter,
|
||
# sie steckt in der .exe; nur diese Route war verwaist.
|
||
|
||
|
||
@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/pfad-map")
|
||
async def worker_setup_pfad_map():
|
||
"""Was ein externer Worker für `RIPPY_PATH_MAP` eintragen muss.
|
||
|
||
DER fehlende Anschluss (Befund 26.07.2026): Ein Windows-Worker bekommt von
|
||
Rippy Container-Pfade (`/app/media/rippy/<job>`). Ohne Übersetzung auf eine
|
||
Freigabe sieht er sie nicht — und das Mapping setzte niemand. Externes
|
||
Encoden konnte deshalb nie funktionieren, obwohl das Celery-Routing
|
||
einwandfrei arbeitete (Job 95afdc89: angenommen, 182 ms später abgelehnt).
|
||
|
||
Geraten wird hier nichts: Rippy hat die Freigabe selbst eingehängt und
|
||
kennt ihre Quelle. Der Installer holt den Vorschlag von hier und schreibt
|
||
ihn in start-tray.bat — der Nutzer muss nichts über Container-Pfade wissen.
|
||
"""
|
||
def sammle():
|
||
vorschlaege = mount_verwaltung.pfad_map_vorschlag(db.list_mounts())
|
||
mapping = mount_verwaltung.pfad_map_zeile(vorschlaege)
|
||
if mapping:
|
||
hinweis = (
|
||
"Diese Freigabe(n) hängt Rippy selbst ein — der Worker erreicht "
|
||
"sie unter demselben Namen im Netzwerk. Wichtig: Arbeits"
|
||
"verzeichnis UND Ablage müssen darunter liegen, sonst kann der "
|
||
"Worker lesen, aber nicht schreiben (oder umgekehrt)."
|
||
)
|
||
elif vorschlaege:
|
||
hinweis = (
|
||
'Es sind nur NFS-Ziele eingehängt. Windows kann NFS zwar über '
|
||
'die Funktion „Client für NFS" einbinden, die Pfad-Schreibweise '
|
||
'lässt sich aber nicht zuverlässig ableiten — bitte selbst '
|
||
'eintragen, Format: /app/media/<name>=Z:\\ (oder UNC-Pfad).'
|
||
)
|
||
else:
|
||
hinweis = (
|
||
"Rippy hat keine Netzwerk-Freigabe eingehängt. Ein externer "
|
||
"Worker kann dann NICHTS komprimieren: Rohdaten und Ziel liegen "
|
||
"auf der Platte der Rippy-Maschine, und dorthin gibt es keinen "
|
||
"Netzwerk-Zugang. Abhilfe: unter Einstellungen → Speicherziele "
|
||
"eine Freigabe (NAS oder PC) einhängen und sie als "
|
||
"Arbeitsverzeichnis UND Ablage wählen."
|
||
)
|
||
return {"mapping": mapping, "eintraege": vorschlaege, "hinweis": hinweis}
|
||
|
||
return await asyncio.to_thread(sammle)
|
||
|
||
|
||
@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)
|
||
# requirements.txt für die venv, rippy.ico für die Verknüpfung
|
||
# auf dem Desktop (sonst trägt sie das Batch-Standardsymbol).
|
||
elif name in ("requirements.txt", "rippy.ico"):
|
||
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 großes 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 unter_wurzel(os.path.realpath(pfad), 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
|
||
|
||
|
||
# GET /stream/jobs (SSE) entfernt am 25.07.2026. Der Stream war zweimal falsch:
|
||
# erst ein Placebo (er sendete nur, wenn eine Liste `sse_connections` gefüllt
|
||
# war, und nichts füllte sie je), dann am 23.07. funktionsfähig gemacht — aber
|
||
# einen Verbraucher hat er nie bekommen. Im UI gibt es kein `EventSource`; das
|
||
# Dashboard holt die Jobs mit `setInterval(loadData, 4000)`. Damit war er keine
|
||
# harmlose Leiche, sondern eine Endlosschleife je Verbindung, die jeder im
|
||
# Heimnetz aufmachen konnte. Wer echtes Push will, braucht BEIDE Seiten.
|
||
|
||
|
||
# /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
|
||
|
||
|
||
# POST /prescan und POST /jellyfin/format entfernt am 25.07.2026 — beide waren
|
||
# Überreste eines ersetzten Entwurfs, ohne einen einzigen Aufrufer:
|
||
#
|
||
# /prescan war der Endpunkt hinter der Metadaten-Vorschau-Seite. Die
|
||
# Seite ist seit v3.4 weg (Korrektur-Popup ist der einzige
|
||
# Weg), der Endpunkt blieb liegen. Der Pre-Scan selbst lebt:
|
||
# der Disc-Watcher ruft PreScan direkt im Prozess auf, das
|
||
# Ergebnis landet auf der Disc-Karte. Nur der HTTP-Weg
|
||
# dorthin hatte keinen Nutzer.
|
||
#
|
||
# /jellyfin/format schrieb NFO-Dateien und lud Poster — in der API. Seit v3.2
|
||
# macht das der Worker (medien.py), und das ist die richtige
|
||
# Stelle: er kennt den Ausgabeordner und ist direkt nach dem
|
||
# Rip am Zug. Mit dem Endpunkt fallen nfo_generator.py und
|
||
# image_downloader.py in der API weg; sonst nutzte sie nichts.
|
||
|
||
|
||
# Auth-Endpoints (/token, /api-keys) entfernt — Commander-Entscheid 24.07.:
|
||
# Heimnetz-only, kein Login-Flow im UI, die Endpoints waren Placebo.
|