Files
rippy/docker/api/main.py
T
HitonabiandClaude Fable 5 cedab61ddf
Ampel / ampel (push) Failing after 41s
chore(main): Der Windows-Anteil ist raus — Entscheid 7, § 11 umgesetzt
Teil A (ganze Dateien): windows_app, fenster, setup_fenster,
einrichtung, standalone, daemon, drives/windows + win_ioctl,
platform/verknuepfungen + win_registry, tools/einrichten,
rip/audio_cd + musicbrainz_cd, queue/ (lokal + laeufer),
packaging/windows/, test_keine_container_reste — samt Tests.
Bilanz: +51 / -10.039 Zeilen in 48 Dateien.

Bleibt TROTZ § 11.2-Listung (Importe gemessen, nicht geraten):
winlauf (OHNE_FENSTER nutzt der Worker; Teil B laeuft nativ Windows),
dateiangaben/katalog/beschaffen (ripping + api importieren sie),
makemkv_aufruf KOMPLETT (ablauf.py ruft key_setzen_und_pruefen;
Registry-Zweige braucht der Remote-Worker), bus/waechter (bedient
Jobs/Logs/Laufwerke der API), die nativ-Pfadzweige (Teil B).

Teil C: drives/treiber() liefert unter Windows den Klartext-Platzhalter
kein_windows.py (Importe/Testsammlung auf Entwicklungs-PCs bleiben
heil, BENUTZUNG wirft mit Verweis auf v5); ripping.rip_cd ohne
Windows-Zweig; betrieb.windows_laufwerke + Aufrufer raus;
/worker-setup/windows* bleibt (Teil B, Entscheid 8).

Lokal: Ruff gruen, 658 Tests gruen (vorher 977 — die Differenz sind
die Windows-Tests, deren Gegenstand mitgegangen ist). Das Wissen liegt
vollstaendig in Rippy v5 (rippy-windows/), mit denselben Testfaellen.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-30 19:18:10 +02:00

3258 lines
137 KiB
Python
Raw Blame History

This file contains invisible Unicode characters
This file contains invisible Unicode characters that are indistinguishable to humans but may be processed differently by a computer. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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
import asyncio
import json
import os
import shutil
import time
import uuid
from rippy import store as db
from rippy import drives as laufwerks_schicht
import eta
from rippy.rip import makemkv_daten
import makemkv_key
import mounts as mount_verwaltung
from rippy.core import notify
from rippy.bus import schema as bus_schema
from rippy.bus.memory import bus as ereignis_bus
from rippy.bus.waechter import Waechter
import phasen
import presets as preset_auswahl
import rohdaten
import celery_client as celery_anbindung
from celery_client import start_rip
from rippy.drives.cdrom import CDS_DISC_OK, CDS_NO_DISC, CDS_TRAY_OPEN
from config_validation import validate_config, ConfigValidationError
from cache import get as cache_get, init_cache, set as cache_set
from ratelimit import (
MAX_REQUESTS_PER_MINUTE,
check_rate_limit,
client_kennung,
darf_melden,
get_rate_limit_remaining,
)
from prescan import PreScan
# ── Der Laufwerks-Treiber DIESER Maschine ──────────────────────────────
#
# Bis V2-4 stand hier `from rippy.drives import linux`. Das laedt ueber
# detection.py das Modul `fcntl` — und damit war main.py unter Windows
# ueberhaupt nicht ladbar (am 28.08.2026 in einer frischen Python-3.12-
# Umgebung gemessen: ModuleNotFoundError: No module named 'fcntl').
#
# Genau dafuer gibt es den Port rippy.ports.Drives: Der Aufrufer sagt, WAS
# er will; welcher Treiber das kann, entscheidet rippy/drives/__init__.py.
# Die CDS_*-Konstanten oben kommen aus cdrom.py und gelten auf beiden
# Plattformen — der Windows-Treiber uebersetzt seine Antworten darauf.
device_discovery = laufwerks_schicht.treiber()
drive_status = device_discovery.drive_status
# 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"
)
# Der Ereignis-Waechter dieses Prozesses. Im Startup gesetzt; von
# /health/vorraete abgefragt, damit ein still gestorbener Waechter
# sichtbar ist statt nur als "es passiert nichts mehr".
ereignis_waechter = None
@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}")
# Ereignis-Waechter (V2-3): macht Aenderungen des Workers zu Ereignissen.
# Der Worker ist ein eigener Container — ohne diese Bruecke wuesste die API
# nichts von seinem Fortschritt, und der SSE-Strom bliebe nach dem ersten
# Snapshot stumm. Warum ueber die Datenbank statt per Rueckruf oder Redis:
# siehe den Modul-Docstring von rippy/bus/waechter.py.
global ereignis_waechter
async def system_stand():
"""Server-Zustand fuers Dashboard — im langsamen Takt des Waechters.
Ohne das haette das Dashboard weiterhin einen eigenen 12-Sekunden-Takt
fuer /system/info, /capabilities und /storage-targets gebraucht — also
je offenem Tab 15 Anfragen pro Minute. So laeuft es EINMAL im Server,
egal wie viele Tabs offen sind.
"""
info, caps, ablagen = await asyncio.gather(
system_info(), capabilities(), storage_targets(),
return_exceptions=True,
)
stand = {}
if not isinstance(info, Exception):
stand["info"] = info
if not isinstance(caps, Exception):
stand["workers"] = (caps or {}).get("workers") or []
if not isinstance(ablagen, Exception):
stand["ablagen"] = ablagen
# Nur schicken, was wirklich gelesen wurde. Ein fehlender Schluessel
# heisst fuer das UI „behalte deinen Stand" — ein leerer Wert hiesse
# „es gibt nichts", und das waere wieder eine falsche Aussage.
return stand or None
# Die Ereignisschleife HIER festhalten. `einmal()` des Waechters laeuft in
# einem Worker-Thread (asyncio.to_thread) — dort gaebe
# `asyncio.get_event_loop()` nicht diese Schleife zurueck, sondern
# erzeugte eine neue oder wuerfe. Die Coroutine liefe dann nie.
schleife = asyncio.get_running_loop()
def system_stand_sync():
return asyncio.run_coroutine_threadsafe(system_stand(), schleife).result(timeout=30)
ereignis_waechter = Waechter(
db, ereignis_bus,
# MIT der erkannten Disc — sonst meldet die Server-Status-Kachel
# „kein Datentraeger", waehrend die Disc erkannt ist (29.08.2026).
laufwerke_lesen=laufwerke_mit_disc,
system_lesen=system_stand_sync,
# DIESELBE Uebersetzung wie /jobs — sonst tragen die Live-Ereignisse
# andere Feldnamen als der Schnappschuss (Befund 29.08.2026:
# „1.1.1970" und „DISC" in der Jobliste).
job_form=lambda z: _job_row_to_model(z).model_dump(),
)
asyncio.create_task(ereignis_waechter.schleife())
# 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.
## Warum ein laufender Job hier alles stoppt (Befund 29.08.2026)
> „der Disk wird gelesen' panel braucht eeeeewig. Der Rip ist bereits bei
> 30% und er liest immernoch."
Genau so war es — und die Erkennung kam auch nie zum Ende. Waehrend eines
Rips haelt `makemkvcon` das Laufwerk; ein zweites `makemkvcon info`
daneben wartet, bis es seine Zeitgrenze erreicht (gemessen: eine
Laufwerks-Abfrage braucht dann 14 s statt 5). Die Marke `_laeuft` blieb
solange stehen, also stand auch das Panel.
Es gibt keinen Grund, waehrend eines Rips zu scannen: Das Laufwerk ist
belegt, und **welche Disc drin ist, wissen wir bereits** — der Job laeuft
ja auf ihr.
"""
if DISC_CACHE.get(pfad, {}).get("_laeuft"):
return
if await asyncio.to_thread(db.has_active_job, pfad):
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}")
# ... und ins PROTOKOLL, nicht nur auf die Konsole
# (Befund 30.08.2026). Der Commander: „im log steht zwar
# erkannt, aber ein start des rips ist nicht moeglich."
# Genau so war es: Die letzte Zeile war das erfolgreiche
# „Disc erkannt" von vorhin, die fehlgeschlagenen Versuche
# danach standen nur auf einer Konsole, die niemand sieht.
# Ein Protokoll, in dem nur die Erfolge stehen, luegt.
try:
db.add_log("warning", "watcher",
"Disc-Erkennung auf %s fehlgeschlagen: %s"
% (pfad, e))
except Exception: # noqa: BLE001
pass # Protokollieren darf die Wache nie anhalten
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
# `job_offen`, nicht `has_active_job`: Waehrend der Kompression ist das
# Laufwerk zwar frei, der Vorgang aber nicht durch — und die Disc wurde
# soeben ausgeworfen. Herleitung in store.job_offen.
if await asyncio.to_thread(db.job_offen, 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 medien_wurzel()
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 pfad_erlaubt(ziel):
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)
def disc_entscheidung(schon_gesehen: bool, vorher, status) -> str:
"""Was ist zu tun? `""` | `"erststart"` | `"eingelegt"` | `"entfernt"`.
Pure Funktion — deshalb ohne Laufwerk prüfbar, und genau daran hing der
Fehler unten.
## Der Befund des Commanders (29.08.2026)
> „liest er die disc NOCHMAL ein das macht aber keinen sinn, wenn er sie
> bereits erkannt hat. Außerdem öffnen sich nun immer irgendwelche fenster
> ganz kurz im hintergrund."
Beides war dieselbe Schleife. In der Disc-Wache stand:
except OSError:
continue
Damit blieb `bekannt[pfad]` ungesetzt, und im nächsten Durchlauf war
`vorher is None` — also wieder „Erststart, liegt schon eine Disc drin".
**Ein einziger fehlgeschlagener Lesevorgang löste einen neuen Vor-Scan
aus.**
Und der scheitert regelmäßig: Während der Vor-Scan läuft, hält
`makemkvcon` das Laufwerk (gemessen: `/devices` braucht dann 14 s statt
5). Also: Scan hält das Laufwerk → Statusabfrage scheitert → nächster
Durchlauf hält es für den ersten → neuer Scan. Eine Schleife, die sich
selbst am Leben hält — und weil derselbe Aufruf sein Konsolenfenster
nicht unterdrückte, blitzte bei jeder Runde eines auf.
Deshalb zwei getrennte Fragen statt einer: **Haben wir dieses Laufwerk je
gelesen?** (`schon_gesehen`) und **was war zuletzt drin?** (`vorher`). Ein
Fehlschlag beantwortet die zweite nicht — und darf die erste nicht
zurücksetzen.
"""
if not schon_gesehen:
return "erststart" if status == CDS_DISC_OK else ""
if status == vorher:
return ""
if status == CDS_DISC_OK:
return "eingelegt"
if status in (CDS_NO_DISC, CDS_TRAY_OPEN) and vorher == CDS_DISC_OK:
return "entfernt"
return ""
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] = {}
# Laufwerke, die schon einmal ERFOLGREICH gelesen wurden. Getrennt von
# `bekannt`, weil ein fehlgeschlagener Lesevorgang sonst wie ein
# Erststart aussieht — siehe unten.
gesehen: set = set()
while True:
try:
for pfad in device_discovery.list_optical_devices():
try:
status = await asyncio.to_thread(drive_status, pfad)
except OSError:
# ⚠️ Hier stand nur `continue` (Befund 29.08.2026).
#
# `bekannt[pfad]` blieb damit ungesetzt, und beim nächsten
# Durchlauf war `vorher is None` — also wieder „Erststart,
# liegt schon eine Disc drin". Ein einziger fehlgeschlagener
# Lesevorgang bewirkte einen NEUEN Vor-Scan.
#
# Und genau das passiert regelmäßig: Während der Vor-Scan
# läuft, hält `makemkvcon` das Laufwerk, und die
# Statusabfrage kommt nicht durch. Ergebnis: Scan hält das
# Laufwerk → Statusabfrage scheitert → nächster Durchlauf
# hält es für den ersten → neuer Scan. Eine Schleife, die
# sich selbst am Leben hält, und bei jeder Runde blitzte
# eine Konsole auf.
#
# `gesehen` merkt sich, ob wir dieses Laufwerk je gelesen
# haben. Ein Fehlschlag heißt jetzt „nichts Neues erfahren"
# — nicht „zum ersten Mal gesehen".
continue
was = disc_entscheidung(pfad in gesehen, bekannt.get(pfad), status)
gesehen.add(pfad)
if was == "erststart":
asyncio.create_task(_auto_prescan(pfad))
elif was == "eingelegt":
db.add_log("info", "watcher", f"Disc eingelegt: {pfad}")
asyncio.create_task(_auto_prescan(pfad))
elif was == "entfernt":
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."""
# Nicht `request.client.host` allein: Der Browser spricht über den nginx im
# UI-Container, dessen Adresse sonst für ALLE Clients gälte (Herleitung in
# ratelimit.client_kennung).
client_ip = client_kennung(
request.client.host if request.client else "",
request.headers.get("x-real-ip") or request.headers.get("x-forwarded-for"),
)
# Rate Limit prüfen
if not check_rate_limit(client_ip):
# Ein 429 war lange unsichtbar: Das UI verbuchte ihn als „nichts da".
# Jetzt steht er im Log, damit die Bremse nicht wieder heimlich greift —
# gedrosselt, damit ein Amok-Skript nicht das Log-Fenster flutet.
try:
if darf_melden(client_ip):
db.add_log(
"warning", "api",
f"Rate-Limit erreicht für {client_ip} — Anfragen werden "
f"abgewiesen (Grenze: {MAX_REQUESTS_PER_MINUTE}/min). Läuft "
"dort ein Skript in einer Schleife?",
)
except Exception:
pass
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
# Welche Wiederholung passt zu DIESEM Fehlschlag? "transcode" | "rip" |
# "unklar" | "" — Herleitung und Begründung in phasen.py. Ohne dieses Feld
# hieß der Knopf nur „Neu" und komprimierte immer, auch ein Rip-Bruchstück.
retry_art: str = ""
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
# Warum steht bei `type`/`status` „unknown"? Leer, solange alles geht.
# Der Linux-Treiber setzt das Feld nicht — dann bleibt es leer.
grund: str = ""
disc: Optional[Dict] = None # Auto-Pre-Scan-Ergebnis (Titel/Jahr/Poster)
# Läuft die Erkennung gerade noch? (Commander 29.08.2026: „das die disc
# erkennung noch läuft muss sichtbar sein")
#
# Der Zustand war schon da — `_auto_prescan` legt beim Start
# `{"_laeuft": True}` in den Vorrat. Nur wurde er nie herausgegeben:
# `laufwerke_mit_disc` ließ so einen Eintrag KOMPLETT weg, damit keine
# halbfertige Disc-Karte mit „Wird erkannt…" als Titel und einem
# „Rippen starten"-Knopf erscheint. Richtig gedacht, falsch gelöst — für
# die Oberfläche sah es dadurch aus, als läge gar keine Disc im Laufwerk.
#
# Und das dauert: `makemkvcon info` läuft an einer Blu-ray in seine
# 120-Sekunden-Grenze (gemessen 29.08.2026 — Disc nach 119 s erkannt).
# Zwei Minuten Schweigen sehen aus wie ein Fehler.
#
# Ein eigenes Feld statt der halben Disc: Die Oberfläche kann „wird
# erkannt" zeigen, ohne einen Titel behaupten zu müssen, den es noch
# nicht gibt.
disc_wird_erkannt: bool = False
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,
)
# ═══════════════════════════════════════════════════════════════════════
# Live-Ereignisse (SSE) — Etappe V2-3
# ═══════════════════════════════════════════════════════════════════════
#
# WAS DAS ERSETZT: Bis hierher fragte das UI im Takt nach. Nachgerechnet an
# den Taktgebern (der Kommentar in ratelimit.py führt sie auf) verursacht EIN
# offener Tab plus ein Worker rund 133 Anfragen pro Minute — und das Rate-Limit
# stand einmal UNTER dieser Zahl, weshalb sich die Job-Liste im Sekundentakt
# leerte. Mit einer offenen Verbindung sind es null.
#
# WARUM SSE UND NICHT WEBSOCKET: Gebraucht wird genau EIN Rückkanal, Server zu
# Browser. Kommandos gehen weiter per REST (idempotent, protokollierbar, mit
# der CLI teilbar). SSE bringt Wiederaufnahme (Last-Event-ID) und den
# automatischen Reconnect des Browsers mit; bei WebSocket müsste man beides
# selbst bauen.
SSE_HERZSCHLAG_SEKUNDEN = 15
async def _snapshot() -> dict:
"""Der vollständige Zustand — das ERSTE, was jeder Client bekommt.
Ohne diesen Anfang sähe ein frisch verbundenes UI nur die Änderungen ab
jetzt und müsste den Rest raten. Mit ihm gilt: erst das ganze Bild, danach
nur noch Deltas.
"""
def sammeln():
return {
# DIESELBE Funktion wie /jobs — sonst fehlt hier die
# Restzeit, und die Oberflaeche liest seit V2-3 nur noch
# diesen Schnappschuss (Befund 29.08.2026).
"jobs": [m.model_dump() for m in jobs_fuer_ui(limit=50)],
"workers": db.list_workers(),
"logs": [_log_zeile(z) for z in db.list_logs(limit=50)],
}
zustand = await asyncio.to_thread(sammeln)
# Laufwerke getrennt: device_info macht ioctls, die hängen können —
# ein defektes Laufwerk darf den Snapshot nicht aufhalten.
try:
# MIT der erkannten Disc. Hier stand dieselbe Schleife ein DRITTES
# Mal (Befund 29.08.2026) — `/devices` hängte die Disc an, der
# Ereignis-Wächter und dieser Schnappschuss nicht. Das Dashboard
# liest genau diesen Schnappschuss und schloss aus der fehlenden
# Disc auf ein leeres Laufwerk: „Bereit — kein Datenträger im
# Laufwerk", während der Titel eine Zeile weiter oben stand.
zustand["devices"] = await asyncio.wait_for(
asyncio.to_thread(laufwerke_mit_disc), timeout=5)
except (asyncio.TimeoutError, OSError):
# „konnte nicht nachsehen" ist etwas anderes als „es gibt keine".
# Deshalb None und nicht [] — das UI behält dann seinen alten Stand.
#
# Vorher stand hier NUR `None`. Waehrend der Disc-Erkennung ist das
# aber der Normalfall (makemkvcon haelt das Laufwerk, /devices
# braucht dann 14 s statt 5) — und nach einem frischen Laden hatte
# das UI keinen alten Stand, den es haette behalten koennen. Der
# Bildschirm blieb leer, ausgerechnet in der Phase, die sichtbar
# sein soll. `laufwerke_notdurft` antwortet ohne ioctl.
zustand["devices"] = laufwerke_notdurft()
# ⚠️ Der Server-Zustand gehört MIT in den Schnappschuss (Befund 28.08.2026).
#
# Er fehlte, und der Wächter schickt ihn nur alle 15 Sekunden. Ein frisch
# geladenes Dashboard stand deshalb bis zu einer Viertelminute auf
# „Platz für Rippy: unbekannt" — und das sieht aus wie eine Auskunft,
# obwohl es keine ist. Genau das hat der Commander im Windows-Fenster
# gesehen.
#
# Mit Zeitgrenze, aus demselben Grund wie bei den Laufwerken: `disk_usage`
# auf einer toten Netzfreigabe hängt im Kernel.
try:
zustand.update(await asyncio.wait_for(system_stand_fuer_snapshot(),
timeout=5))
except Exception: # noqa: BLE001
# Bewusst breit und bewusst still: Der Schnappschuss ist die Grundlage
# fuer ALLES im UI. Lieber ohne Platzangabe (das UI zeigt dann seinen
# alten Stand) als gar kein Schnappschuss.
pass
return zustand
async def system_stand_fuer_snapshot() -> dict:
"""Platz, Worker und Ablageziele — dieselben Felder wie im Wächter-Takt.
Getrennte Funktion, damit `_snapshot` und der Wächter dasselbe liefern.
Zwei Wege, die dieselbe Kachel füllen, aber verschiedene Felder schicken,
wären eine Anzeige, die je nach Zeitpunkt anders aussieht.
"""
info, ablagen = await asyncio.gather(
system_info(), storage_targets(), return_exceptions=True)
stand = {}
if not isinstance(info, Exception):
stand["info"] = info
if not isinstance(ablagen, Exception):
stand["ablagen"] = ablagen
return stand
def _sse_rahmen(ereignis: dict) -> str:
"""Ein Ereignis im SSE-Format. `id:` ist die Basis der Wiederaufnahme."""
nutzlast = json.dumps(ereignis, ensure_ascii=False)
zeilen = [
"id: " + str(ereignis["seq"]),
"event: " + ereignis["typ"],
"data: " + nutzlast,
"",
"",
]
return "\n".join(zeilen)
@app.get("/events")
async def events(request: Request, last_event_id: str = None):
"""Live-Strom: erst ein Snapshot, danach nur noch Änderungen.
Pfad bewusst `/events` und nicht `/api/v2/events`: Der nginx im UI-Container
entfernt das Präfix `/api/` (siehe ui/nginx.conf), und alle bestehenden
Routen hier sind entsprechend unpräfigiert (`/jobs`, `/devices`). Aus dem
Browser heißt der Aufruf damit `/api/events`. Die Versionierung `/api/v2/*`
aus KONZEPT-V2.md § 6.1 kommt, wenn die Routen in Router aufgeteilt werden —
sie jetzt für eine einzige Route einzuführen, hätte zwei Konventionen
nebeneinander bedeutet.
Der Browser schickt beim Wiederverbinden von selbst `Last-Event-ID` mit.
Passt die Lücke in den Ringpuffer, wird sie nachgeliefert; passt sie nicht,
kommt ein neuer Snapshot — AUSDRÜCKLICH, nicht stillschweigend. Ein UI, das
sich fälschlich für aktuell hält, ist schlimmer als eines, das neu lädt.
"""
kopf = request.headers.get("last-event-id") or last_event_id
try:
ab_seq = int(kopf) if kopf else None
except ValueError:
ab_seq = None
async def strom():
abo = ereignis_bus.abonnieren(ab_seq)
try:
verpasst = ereignis_bus.nachliefern(ab_seq) if ab_seq is not None else None
if verpasst is None:
# Kein Rückstand bekannt oder Lücke zu groß -> ganzes Bild.
yield _sse_rahmen(bus_schema.baue_ereignis(
"snapshot", await _snapshot(), seq=ereignis_bus.seq))
else:
for ereignis in verpasst:
yield _sse_rahmen(ereignis)
while True:
if await request.is_disconnected():
return
ereignis = await abo.naechstes(timeout=SSE_HERZSCHLAG_SEKUNDEN)
if ereignis is None:
# Herzschlag: manche Proxys schließen stille Verbindungen.
yield ": herzschlag\n\n"
continue
if abo.abgehaengt:
# Dieser Zuhörer hat den Anschluss verloren — ganzes Bild
# statt Bruchstücken.
abo.abgehaengt = False
yield _sse_rahmen(bus_schema.baue_ereignis(
"snapshot", await _snapshot(), seq=ereignis_bus.seq))
continue
yield _sse_rahmen(ereignis)
finally:
abo.schliessen()
return StreamingResponse(
strom(),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
# nginx puffert `text/event-stream` sonst und der Strom käme
# in Schüben statt live an.
"X-Accel-Buffering": "no",
},
)
@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,
},
# Der Ereignis-Waechter (V2-3) ist die Bruecke zwischen Worker und
# SSE-Strom. Stirbt er still, steht das UI — und zwar OHNE Fehlermeldung,
# weil ein Abriss dort bewusst als "nichts Neues" gilt und nicht als
# "nichts da". Deshalb muss sein Alter von aussen abfragbar sein.
"ereignisse": {
"alter_sekunden": None if ereignis_waechter is None
else round(ereignis_waechter.lebt_seit_sekunden(), 1),
"gesund": bool(ereignis_waechter and ereignis_waechter.gesund),
"abonnenten": len(getattr(ereignis_bus, "_abonnenten", ())),
"letzte_seq": ereignis_bus.seq,
},
}
@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)
def _arbeitsverzeichnis_des_jobs(job, einstellungen=None) -> str:
"""Wohin ging der Roh-Rip DIESES Jobs? Seine Wahl schlägt die Einstellung.
## Warum die Einstellung allein nicht reicht (Befund 30.08.2026)
Der Rip-Dialog lässt für JEDEN Rip einzeln wählen, wohin die Rohdaten
gehen (seit v3.15). Die Wahl landet in den Job-Metadaten als
`work_dir` — geschrieben in `start_rip`, gelesen im Worker von
`_arbeitsverzeichnis()`. Gesucht wurde danach aber immer unter dem
HEUTIGEN Wert der Einstellung `workDir`.
Genau daran scheiterte am 26.07.2026 schon einmal ein Job mit 79,6 GB
Rohschnitt — der Fall steht im Kopf von `rohdaten.py`. Repariert wurde
damals die Kandidatenliste, nicht der Aufrufer: Er reichte weiterhin
die Einstellung hinein. Am 30.08.2026 stand deshalb erneut „Auf der
Platte liegt zu diesem Job nichts (mehr)" — vor 16,5 GB, die dalagen.
Das ist das Muster aus AGENTS: aus einem Zustandswert (der heutigen
Einstellung) auf einen Vorgang geschlossen (wohin DAMALS gerippt
wurde), statt nachzusehen. Der Job weiß es selbst.
"""
holen = getattr(job, "get", None)
eigen = ""
if holen:
eigen = (phasen.meta_von(holen("meta")).get("work_dir") or "").strip()
if not eigen:
werte = einstellungen if einstellungen is not None else db.get_settings()
eigen = ((werte or {}).get("workDir") or "").strip()
# normpath("") wäre "." — der aktuelle Ordner, und der ist hier nie
# gemeint. "/" fällt in `kandidaten()` sauber durch.
return os.path.normpath(eigen or "/")
# 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.
"""
einstellungen = db.get_settings()
alt = _ROHDATEN["treffer"]
treffer = {}
for zeile in db.list_jobs():
if zeile.get("status") != "failed":
continue
job_id = zeile["id"]
# Je Job SEINE Wahl — nicht die heutige Einstellung.
work_dir = _arbeitsverzeichnis_des_jobs(zeile, einstellungen)
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). Die Wurzel
# kommt vom Betrieb — auf Windows gibt es kein /app/temp (29.08.2026).
from rippy import pfade
return os.path.isdir(pfade.verbinden(rohdaten.wurzeln()[0], 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).
"""
return await asyncio.to_thread(jobs_fuer_ui)
def jobs_fuer_ui(limit: int = 100) -> list:
"""Die Jobs in UI-Form — MIT Restzeit, Wiederhol-Art und Rohdaten-Frage.
## Warum das eine Funktion ist (Befund 29.08.2026)
> „Über den kompletten vorgang steht dort Restzeit wird gemessen' aber
> messung wird nicht abgeschlossen. Das heißt man hat kein ETA"
Zwei Ursachen, beide davon, dass es die Jobliste ZWEIMAL gab.
**1.** Diese Rechnung stand nur in `/jobs`. Der SSE-Schnappschuss baute
seine Jobs mit dem nackten `_job_row_to_model` — also ohne Restzeit. Seit
der Umstellung auf den Ereignisstrom (V2-3) liest die Oberfläche aber
genau diesen Schnappschuss und nicht mehr `/jobs`. Die Restzeit wurde also
weiterhin brav berechnet und niemandem gezeigt.
**2.** Die Messreihe lag ausschließlich in Redis, das es unter Windows
nicht gibt — Begründung und Rückfall in `eta.py`.
Dieselbe Lehre wie bei `laufwerke_mit_disc`: Eine Auskunft in zwei
Fassungen ist eine Fassung zu viel.
"""
# Unlesbare Einstellungen duerfen die Jobliste NICHT umwerfen — sie
# steuern hier nur, ob „Neu komprimieren" angeboten wird. Seit diese
# Funktion auch den SSE-Schnappschuss baut, haengt daran die ganze
# Oberflaeche (Befund 29.08.2026: zwei Snapshot-Tests wurden dadurch rot).
try:
einstellungen = db.get_settings(bei_fehler_leer=True)
except Exception: # noqa: BLE001
einstellungen = {}
work_dir = os.path.normpath((einstellungen.get("workDir") or "").strip() or "/")
jetzt = time.monotonic()
modelle = []
for z in db.list_jobs(limit=limit):
modell = _job_row_to_model(z)
modell.can_retry = _kann_neu_komprimieren(z, work_dir)
modell.retry_art = phasen.retry_art(z)
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
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():
pfade = _rohdaten_suchen(job_id, _arbeitsverzeichnis_des_jobs(job))
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():
pfade = _rohdaten_suchen(job_id, _arbeitsverzeichnis_des_jobs(job))
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
# Der Pfad IM Container. Bleibt als Rückfall stehen — aber er ist NICHT mehr
# die Antwort auf „wo liegt die Ablage": die gibt `medien_wurzel()`.
MEDIA_ROOT = "/app/media"
def _betriebswerte() -> dict:
"""Die Konfiguration. Unlesbar heißt Vorgaben, nicht Absturz."""
from rippy import config as rippy_config
try:
return rippy_config.laden()
except Exception: # noqa: BLE001
return {}
def medien_wurzel() -> str:
"""Wo dieser Betrieb ablegt — `/app/media` nur, wenn es ein Container ist.
## Der Befund des Commanders (28.08.2026)
> „Warum heißt das hier noch container platte? Er holt sich das
> Arbeitsverzeichnis ja von der Installation. Wäre es möglich das
> Arbeitsverzeichnis zu ändern? momentan geht das nicht."
Es ging nicht, und zwar aus einem Grund: `MEDIA_ROOT` war fest
`/app/media`. Auf seinem PC gibt es den Ordner nicht, also warf
`os.listdir` in `/storage-targets`, also blieb die Liste leer — im
Auswahlfeld stand genau ein Eintrag, und der hieß „Container-Platte".
Kein Fehler, keine Meldung, nur eine Auswahl ohne Auswahl.
Dieselbe Konstante war zugleich die Pfadgrenze für `/browse`. Auch der
Ordner-Browser konnte auf Windows also nichts anzeigen.
"""
from rippy import betrieb
try:
return betrieb.medien_wurzel(_betriebswerte()) or MEDIA_ROOT
except Exception: # noqa: BLE001
return MEDIA_ROOT
def frei_blaettern() -> bool:
"""Darf außerhalb der Wurzel geblättert werden? Begründung in betrieb.py."""
from rippy import betrieb
try:
return betrieb.frei_blaettern(_betriebswerte())
except Exception: # noqa: BLE001
return False
def pfad_erlaubt(pfad: str, wurzel: str = None, frei: bool = None) -> bool:
"""Darf die API diesen Pfad anfassen? (pure Funktion, einspritzbar)
Im Container und im verteilten Betrieb gilt die Wurzel — die API hängt
dort im Netz. In der nativen App gilt sie nicht: Sie bedient den Menschen,
der vor dem Rechner sitzt, und dessen Ziel ist eine UNC-Freigabe, die
unter gar keiner lokalen Wurzel liegt.
"""
if not pfad:
return False
if frei_blaettern() if frei is None else frei:
return True
return unter_wurzel(pfad, medien_wurzel() if wurzel is None else wurzel)
def _validiere_ziel(target_dir: Optional[str]) -> Optional[str]:
"""Ziel muss erlaubt sein — Pfad-Ausbrüche (..) fliegen raus."""
if not target_dir:
return None
normalisiert = os.path.normpath(target_dir)
if not pfad_erlaubt(normalisiert):
raise HTTPException(
status_code=422,
detail=f"Ziel muss unter {medien_wurzel()} 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))
# TMDB-Laufzeiten der Episoden holen (für die Erkennung im Worker)
if meta_dict.get("id") and meta_dict.get("type") == "tv":
try:
from clients.tmdb import TMDBClient
staffel = TMDBClient().get_tv_season(meta_dict["id"], meta_dict["season"])
if staffel and "episodes" in staffel:
meta_dict["episode_runtimes"] = [
ep.get("runtime") for ep in staffel["episodes"] if ep.get("runtime")
]
except Exception as e:
print(f"Fehler beim Holen der Staffel-Laufzeiten: {e}")
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())
# Eine noch laufende Disc-Erkennung ist ab jetzt gegenstandslos: Das
# Laufwerk gehoert dem Rip, und WELCHE Disc drin ist, wissen wir. Ohne
# das Wegraeumen stuende „Disc wird gelesen" bis zum Ende des Rips
# (Commander 29.08.2026: „Der Rip ist bereits bei 30% und er liest
# immernoch"). Der laufende Scan traegt sein Ergebnis spaeter ohnehin
# selbst nach.
if (DISC_CACHE.get(device_path) or {}).get("_laeuft"):
DISC_CACHE.pop(device_path, None)
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")
detail["retry_art"] = phasen.retry_art(job)
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 — muss erlaubt sein."""
ausgabe = os.path.normpath(job.get("output_path") or "")
if not pfad_erlaubt(ausgabe):
raise HTTPException(status_code=404, detail="Job hat keinen erlaubten Ausgabeordner")
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 pfad_erlaubt(os.path.realpath(pfad))
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 der Medien-Wurzel.
NFS/SMB-Shares, die auf der VM unter /srv/rippy/media eingehängt werden,
tauchen hier automatisch auf (rslave-Bind in docker-compose).
## Warum hier Laufwerke dazukommen (Commander-Befund 28.08.2026)
Auf Windows war diese Liste IMMER leer: `os.listdir("/app/media")` warf,
und der `except OSError` gab still `[]` zurück. Im Auswahlfeld für das
Arbeitsverzeichnis stand dann genau ein Eintrag — „Standard aus den
Einstellungen (Container-Platte)". Eine Auswahl ohne Auswahl.
Die Wurzel allein reicht dort auch nicht: Der Roh-Rip einer 4K-UHD ist bis
zu 100 GB groß, und die Antwort darauf ist fast immer ein ANDERES
Laufwerk. Deshalb kommen im nativen Betrieb die Laufwerke dazu — mit dem
freien Platz daneben, damit die Wahl eine informierte ist.
"""
def sammle():
wurzel = medien_wurzel()
ziele = []
def platz(pfad):
try:
return round(shutil.disk_usage(pfad).free / 1024**3, 1)
except OSError:
return None
# Die Wurzel selbst — im Container war sie nie ein Eintrag, weil dort
# die Unterordner die Ziele sind. Nativ IST sie ein gültiges Ziel.
if frei_blaettern() and os.path.isdir(wurzel):
ziele.append({"name": "Ablage (%s)" % wurzel, "path": wurzel,
"is_mount": False, "free_gb": platz(wurzel)})
try:
eintraege = sorted(os.listdir(wurzel))
except OSError:
eintraege = []
for name in eintraege:
pfad = os.path.join(wurzel, 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,
})
# Der Laufwerksbuchstaben-Block stand hier für die Windows-
# Standalone-App — sie ist mit v5 ein eigenes Produkt (KONZEPT-
# WINDOWS.md Entscheid 7); im Container gab es hier nie etwas.
return ziele
return await asyncio.to_thread(sammle)
def geraetepfad(name: str) -> str:
"""Die Kennung aus der URL zum echten Geraetepfad — oder 404.
## Der Befund (29.08.2026)
Hier stand an DREI Stellen:
device_path = f"/dev/{name}"
Das UI ruft diese Endpunkte mit der Kennung aus der Geraeteliste auf —
unter Windows also `G`. Gebaut wurde daraus `/dev/G`, und das steht in
keiner Laufwerksliste. **Auswerfen und „Disc scannen" antworteten unter
Windows also immer mit 404**, ohne dass irgendwo stuende, warum.
Aufgeloest wird jetzt vom TREIBER: Er vergibt die Kennung (`kennung`),
also loest er sie auch wieder auf. Hin- und Rueckweg koennen damit nicht
mehr auseinanderlaufen.
"""
pfad = device_discovery.pfad_zu_kennung(name)
if not pfad:
raise HTTPException(status_code=404,
detail=f"Laufwerk {name} nicht gefunden")
return pfad
@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 = geraetepfad(name)
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 20120 s.
"""
device_path = geraetepfad(name)
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_anbindung.abschicken("worker.tasks.scan_tracks", [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 = geraetepfad(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 = _arbeitsverzeichnis_des_jobs(job, einstellungen)
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"{medien_wurzel()}/{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_anbindung.abschicken(
"worker.tasks.transcode_files", [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}/retry-rip", status_code=201)
async def retry_rip(job_id: str):
"""Rippt die Disc dieses Jobs NOCH EINMAL — als frischer Job.
Der Gegenpart zu `retry-transcode`. Nötig, weil es für einen mitten im Rip
gestorbenen Job vorher überhaupt keinen richtigen Knopf gab: Das UI bot nur
„Neu" an, und das war immer die Kompression. Am 26.07.2026 hätte das aus
5,1 GB Bruchstück brav einen Film gemacht, der bei 12 % aufhört.
Bewusst ein NEUER Job mit neuer ID, nicht ein Wiederbeleben des alten:
* Das Roh-Verzeichnis heißt <Arbeitsverzeichnis>/<job_id>. Bei gleicher ID
läge das alte Bruchstück im neuen Verzeichnis, und die Kompression
sammelt am Ende ALLE MKV-Dateien darin ein — sie würde das Bruchstück
mitverarbeiten.
* Der Fehlschlag bleibt in der Liste nachlesbar, statt überschrieben zu
werden.
Übernommen werden Titel, Ablageziel und alle Metadaten des alten Jobs
(Jahr/Poster, Titel-Auswahl, Sprachwunsch, Arbeitsverzeichnis, gewählter
Encoder-Worker) — der Nutzer soll seine Wahl nicht neu treffen müssen. Die
Phasen-Marke wird NICHT übernommen; die setzt der Rip selbst.
"""
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", "transcoding", "canceling"):
raise HTTPException(status_code=409, detail="Dieser Job läuft noch")
device_path = job.get("device")
vorhandene = device_discovery.list_optical_devices()
if not device_path or device_path not in vorhandene:
raise HTTPException(
status_code=409,
detail=(
f"Das Laufwerk dieses Jobs ({device_path or 'unbekannt'}) ist nicht "
'mehr da. Disc in ein vorhandenes Laufwerk legen und den Rip über '
'„Rip starten" neu anstoßen.'
),
)
try:
status = await asyncio.to_thread(drive_status, device_path)
except OSError as e:
raise HTTPException(status_code=409, detail=f"Laufwerk {device_path} antwortet nicht: {e}")
if status != CDS_DISC_OK:
raise HTTPException(
status_code=409,
detail=(
f'Es liegt keine (lesbare) Disc in {device_path}. Ein neuer Rip '
'braucht die Disc — sie wurde nach dem Fehlschlag vermutlich '
'ausgeworfen. Disc einlegen, dann noch einmal.'
),
)
try:
meta = json.loads(job.get("meta") or "{}")
except ValueError:
meta = {}
if not isinstance(meta, dict):
meta = {}
meta.pop(phasen.RIP_FERTIG, None)
meta_json = json.dumps(meta) if meta else None
neue_id = str(uuid.uuid4())
ziel = job.get("target_dir")
await asyncio.to_thread(
db.insert_job, neue_id, device_path, job.get("disc_type"),
job.get("title"), ziel, meta_json,
)
await asyncio.to_thread(
db.add_log, "info", "api",
f"Job {neue_id} ist der Neu-Rip von {job_id} "
f"({job.get('title') or 'ohne Titel'}) auf {device_path}. "
"Das unvollständige Rohmaterial des alten Jobs bleibt liegen — es kann "
"über den Papierkorb am alten Job mitgelöscht werden.",
)
start_rip(device_path, neue_id, ziel)
return {"id": neue_id, "status": "pending", "device": device_path, "vorher": job_id}
@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_anbindung.hole_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),
"server_version": os.getenv("RIPPY_VERSION", "dev")
}
@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/Kodi anstoßen — auch als Verbindungs-Test
aus den Einstellungen."""
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():
# Fallback für die reine API-Test-Funktion, die den Server-Typ nicht kennt:
# Wir versuchen Kodi (JSON-RPC) und danach Emby/Jellyfin.
# Wenn der User Kodi auswählt, hat er oft Basic Auth konfiguriert.
auth = None
if ":" in request.api_key:
user, pw = request.api_key.split(":", 1)
auth = (user, pw)
try:
# Kodi JSON-RPC
payload = {"jsonrpc": "2.0", "method": "VideoLibrary.Scan", "id": 1}
res = _requests.post(url + "/jsonrpc", json=payload, auth=auth, timeout=5)
if res.status_code < 300:
return res
except Exception:
pass
# Emby / Jellyfin
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 = ""):
"""Server-seitiger Ordner-Browser für die Ziel- und Arbeitsordner-Wahl.
Leerer Pfad heißt „ganz oben": im Container die Medien-Wurzel, nativ die
Liste der Laufwerke. Ohne diese oberste Ebene käme man auf Windows nie zu
einem anderen Laufwerk — und genau dort ist Platz für 100 GB Rohdaten.
"""
if not path.strip():
laufwerke = []
if laufwerke:
return {"path": "", "parent": None, "dirs": laufwerke, "files": []}
path = medien_wurzel()
normalisiert = os.path.normpath(path)
if not pfad_erlaubt(normalisiert):
raise HTTPException(status_code=422,
detail=f"Nur Pfade unter {medien_wurzel()}")
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
return {"path": normalisiert, "parent": eltern_von(normalisiert),
"dirs": ordner, "files": dateien}
def eltern_von(pfad: str, wurzel: str = None, frei: bool = None) -> Optional[str]:
"""Wohin führt „eine Ebene höher"? — `None` heißt: hier ist oben.
Zwei Fallen, beide nur auf Windows sichtbar:
1. `os.path.dirname("C:\\\\")` ist wieder `"C:\\\\"`. Ein Knopf „nach oben",
der auf denselben Ordner zeigt, sieht aus wie ein Fehler.
2. Über der Laufwerkswurzel steht nicht *nichts*, sondern die Liste der
Laufwerke — der leere Pfad. Sonst käme man von `D:\\` nie zu `C:\\`.
"""
frei = frei_blaettern() if frei is None else frei
wurzel = (medien_wurzel() if wurzel is None else wurzel)
if not frei:
return os.path.dirname(pfad) if pfad != wurzel else None
oben = os.path.dirname(pfad)
# Laufwerkswurzel (dirname zeigt auf sich selbst) -> die Laufwerksliste.
return "" if oben == pfad else oben
class MkdirRequest(BaseModel):
path: str
name: str
@app.post("/browse/mkdir", status_code=201)
async def browse_mkdir(request: MkdirRequest):
"""Neuen Ordner anlegen (Speicherziele-Verwaltung, Arbeitsordner)."""
basis = os.path.normpath(request.path)
if not pfad_erlaubt(basis):
raise HTTPException(status_code=422,
detail=f"Nur Pfade unter {medien_wurzel()}")
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}
# ═══════════════════════════════════════════════════════════════════════
# Werkzeuge: Bestand, Update-Stand, Beschaffung — Etappe V2-4
# ═══════════════════════════════════════════════════════════════════════
#
# WOFUER: Damit Rippy unter Windows eigenstaendig arbeiten kann, muss es
# seine Werkzeuge nicht nur FINDEN, sondern auch HOLEN und aktuell halten
# koennen. Im Container ist das anders — dort stecken sie im Image, und ein
# Update ist ein Rebuild (siehe /system/updates, das bleibt unveraendert).
#
# WARUM EIGENE ROUTEN statt /system/updates zu erweitern: /system/updates
# beantwortet "welche Version gibt es draussen". Diese hier beantworten
# "was liegt auf DIESER Maschine, wo, und was kann ich dagegen tun". Zwei
# verschiedene Fragen; eine Route, die beides taete, koennte keine davon
# ehrlich beantworten.
@app.get("/betrieb")
async def betrieb_auskunft():
"""In welchem Betrieb laeuft Rippy — und was kann dieser Betrieb?
## Warum es diese Route gibt (Commander-Befund 28.08.2026)
Im Windows-Fenster stand auf der Server-Status-Kachel:
Worker erreichbar: 0 von 1
Kein Worker antwortet — Pruefen: docker compose ps
Container-Platte: unbekannt
Freigaben: keine eingehaengt
Kein Satz davon ergibt auf einem Windows-PC einen Sinn. Es gibt keinen
Container, kein `docker compose`, keinen zweiten Worker — Rippy rippt
dort selbst.
Die Ursache war nicht die Anzeige, sondern was ihr fehlte: **Das UI hat
nie erfahren, worauf es laeuft.** `config.py` kennt das Profil seit
V2-2, weitergegeben wurde es nie. Also hat das UI angenommen.
Geantwortet wird mit FAEHIGKEITEN, nicht mit einem Modus-Namen — sonst
muesste die Oberflaeche aus einem Namen auf Verhalten schliessen, und
das bricht beim naechsten Betriebsfall. Begruendung vollstaendig in
`src/rippy/betrieb.py`.
"""
from rippy import betrieb as betriebs_auskunft
from rippy import config as rippy_config
try:
werte = rippy_config.laden()
except Exception: # noqa: BLE001
# Eine unlesbare Konfiguration darf die Oberflaeche nicht lahmlegen —
# sie bekommt dann die Vorgaben, und die stimmen fuer den haeufigsten
# Fall. Der Fehler faellt an anderer Stelle laut auf.
werte = {}
return betriebs_auskunft.auskunft(werte)
@app.get("/system/werkzeuge")
async def system_werkzeuge():
"""Was ist installiert, wo, in welcher Fassung — und gibt es Neueres?
`neueste` ist leer, wenn die Quelle nicht antwortet. Dann steht dort
ausdruecklich NICHT "Update verfuegbar": Eine Nichtauskunft ist keine
Aussage ueber die Welt.
"""
from rippy.tools import beschaffen as werkzeug_beschaffung
einstellungen = await asyncio.to_thread(db.get_settings)
eingestellt = (einstellungen or {}).get("werkzeugPfade") or {}
lage = await asyncio.to_thread(werkzeug_beschaffung.lage, eingestellt)
return {
"werkzeuge": lage,
# Wohin Rippy selbst installiert — im UI sichtbar, damit niemand
# raten muss, wo die geholten Dateien landen.
"ordner": werkzeug_beschaffung.katalog.werkzeug_ordner(),
"windows": os.name == "nt",
}
@app.post("/system/werkzeuge/{name}/holen")
async def system_werkzeug_holen(name: str):
"""Holt ein Werkzeug und installiert es.
HandBrake wird vollstaendig automatisch geholt (GitHub-Release, ein ZIP
mit einer .exe — kein Installer, keine Administratorrechte).
MakeMKV wird NICHT mitgeliefert: Rippy laedt die offizielle Datei vom
Hersteller und startet sie. Das ist dieselbe Black-Box-Trennung, die
KONZEPT.md § 6 fuer den Container festhaelt — MakeMKV bleibt ein fremdes
Programm, Rippy nimmt nur die Handgriffe ab.
"""
from rippy.tools import beschaffen as werkzeug_beschaffung
if name not in werkzeug_beschaffung.katalog.WERKZEUGE:
raise HTTPException(
status_code=404,
detail=f"Unbekanntes Werkzeug '{name}'. Bekannt: "
+ ", ".join(sorted(werkzeug_beschaffung.katalog.WERKZEUGE)))
meldungen = []
def fortschritt(text, anteil=None):
if text:
meldungen.append(text)
try:
if name == "handbrake":
pfad = await asyncio.to_thread(
werkzeug_beschaffung.handbrake_holen, None, fortschritt)
else:
pfad = await asyncio.to_thread(
werkzeug_beschaffung.makemkv_holen, "", None, fortschritt, False)
except werkzeug_beschaffung.BeschaffungsFehler as e:
# Klartext statt Traceback: Der haeufigste Grund ist eine nicht
# erreichbare Quelle, und das ist nichts, was der Nutzer im Code
# suchen sollte.
await asyncio.to_thread(db.add_log, "error", "werkzeuge", str(e))
raise HTTPException(status_code=502, detail=str(e)) from e
await asyncio.to_thread(
db.add_log, "info", "werkzeuge",
f"{werkzeug_beschaffung.katalog.WERKZEUGE[name]['titel']} beschafft: {pfad}")
return {"name": name, "pfad": pfad, "meldungen": meldungen}
@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
elif name in ("requirements.txt", "rippy.ico"):
z.write(os.path.join("worker_dist", name), name)
# Das gemeinsame Paket muss MIT (Etappe V2-0): tasks.py und caps.py
# importieren seitdem "from rippy.rip import makemkv_daten". Die
# Schleife oben sieht nur die oberste Ebene — ein Unterordner käme
# nie mit, und der Windows-Worker stürbe beim Start mit
# ModuleNotFoundError. Deshalb hier ausdrücklich durchlaufen.
paket = os.path.join("worker_dist", "rippy")
for wurzel, _, dateien in os.walk(paket):
for datei in sorted(dateien):
if not datei.endswith(".py") or datei.startswith("test_"):
continue
voll = os.path.join(wurzel, datei)
z.write(voll, os.path.relpath(voll, "worker_dist"))
# Write RIPPY_VERSION to version.txt
version_str = os.getenv("RIPPY_VERSION", "dev")
z.writestr("version.txt", version_str)
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/version")
async def system_version():
"""Gibt die aktuelle Rippy-Version des Servers zurück (aus deploy.sh)."""
return {"version": os.getenv("RIPPY_VERSION", "dev")}
@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()}
# ⚠️ NICHT fest /app/media und /app/temp (Befund 28.08.2026).
#
# Beide Pfade gibt es unter Windows nicht. `shutil.disk_usage` warf,
# die Liste blieb leer, und im Windows-Fenster stand "Platz fuer
# Rippy: unbekannt" — obwohl auf dem Laufwerk natuerlich Platz war.
# Eine Nichtauskunft, die wie eine Auskunft aussieht.
from rippy import betrieb as betriebs_auskunft
from rippy import config as rippy_config
try:
werte = rippy_config.laden()
except Exception: # noqa: BLE001
werte = {}
# Die Oberflaeche schreibt nach `outputDir`/`workDir` in die
# DATENBANK, `betrieb` liest `storage.*` aus der DATEI. Ohne
# diese Bruecke zeigte die Uebersicht immer die Vorgabe, egal
# was eingestellt war (Befund 30.08.2026, siehe
# `betrieb.mit_einstellungen`).
werte = betriebs_auskunft.mit_einstellungen(
werte, db.get_settings(bei_fehler_leer=True))
for ort in betriebs_auskunft.platz_orte(werte):
# Frisch installiert gibt es den Ablage-Ordner noch nicht. Dann
# das naechste vorhandene Elternverzeichnis messen: Der Nutzer
# will wissen, ob auf dem LAUFWERK Platz ist.
pfad = betriebs_auskunft.naechster_vorhandener(ort["pfad"])
if not pfad:
continue
try:
nutzung = shutil.disk_usage(pfad)
except OSError:
continue
info["plaetze"].append({
"name": ort["name"],
# Der PFAD gehoert dazu (Commander 29.08.2026): „bei Platz
# fuer rippy' sollte eher das arbeitsverzeichnis und der
# Ablagepfad sein." Eine Zahl ohne Ort sagt nicht, WO der
# Platz knapp wird.
"pfad": ort["pfad"],
"gleiches_laufwerk": bool(ort.get("gleiches_laufwerk")),
"frei_gb": round(nutzung.free / 1024**3, 1),
"gesamt_gb": round(nutzung.total / 1024**3, 1),
})
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/makemkv-key/holen")
async def makemkv_key_holen():
"""Holt den kostenlosen MakeMKV-Beta-Key JETZT aus dem Forum.
## Warum es diesen Knopf gibt (Commander 28.08.2026)
> „Bezüglich des MKV Beta Keys — der könnte theoretisch auch automatisch
> ausgelesen werden, ich glaube das web rippy kann das"
Er hatte recht: `makemkv_key.refresh_loop()` läuft seit jeher täglich mit
(main.py, Startereignis). Nachgemessen: Forum antwortet in 3,3 Sekunden,
Key mit 62 Zeichen.
Nur war davon **nichts zu sehen**. Kein Knopf, keine Zeitangabe, kein
Hinweis — im Einstellungsfeld stand ein Key, und ob der von Hand kam oder
von selbst, wusste niemand. Eine Automatik, die man nicht sehen kann, ist
für den Benutzer keine.
Der Key wechselt etwa monatlich. Wer nicht warten will, bis die Schleife
das nächste Mal nachsieht, drückt hier.
Rechtlich unverändert: öffentliche Beta-LIZENZ der Software, KEIN
Disc-Schlüssel.
"""
import makemkv_key as key_modul
try:
key = await asyncio.to_thread(key_modul.fetch_current_key)
except Exception as e: # noqa: BLE001
raise HTTPException(
status_code=502,
detail=f"Das MakeMKV-Forum antwortet nicht ({type(e).__name__}). "
"Der Key lässt sich weiter von Hand eintragen.")
if not key:
raise HTTPException(
status_code=502,
detail="Im Forum-Beitrag stand kein Key im erwarteten Format. "
"Vermutlich hat sich die Seite geändert.")
geaendert = await asyncio.to_thread(key_modul._apply_key, key)
# Nimmt MakeMKV ihn an? Nur MakeMKV selbst weiss das.
#
# Am 29.08.2026 auf dem Rechner des Commanders: der frische Schluessel
# wurde mit MSG 5020 („ungueltig") und 5021 („Programmversion zu alt")
# abgelehnt — sein MakeMKV 1.18.4 ist aelter als der Schluessel verlangt.
# Ohne diese Rueckfrage haette hier „erfolgreich geholt" gestanden, und
# der naechste Rip waere trotzdem gescheitert.
from rippy.rip import makemkv_aufruf
from rippy.tools import katalog as werkzeug_katalog
urteil = await asyncio.to_thread(
makemkv_aufruf.key_setzen_und_pruefen, key,
werkzeug_katalog.finden("makemkv"))
await asyncio.to_thread(
db.add_log, "success" if urteil["ok"] else "warning", "makemkv-key",
"Beta-Key aus dem Forum geholt (…%s)%s%s"
% (key[-6:], "" if geaendert else " — war schon der aktuelle",
"" if urteil["ok"] else " — ABGELEHNT: " + urteil["grund"]))
# Der Key selbst wird NICHT zurueckgegeben: Er steht in den Einstellungen,
# und ein zweiter Weg zu demselben Wert ist ein zweiter Weg, ihn zu
# verlieren.
return {"geholt": True, "geaendert": geaendert, "endet_auf": key[-6:],
"angenommen": urteil["ok"], "grund": urteil["grund"]}
@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}
def _log_zeile(z: dict) -> dict:
"""DB-Zeile -> UI-Form. EINE Abbildung, zwei Aufrufer.
Sie stand bis V2-3 nur im /logs-Endpunkt. Der Snapshot des SSE-Stroms
haette daneben die rohen DB-Zeilen geliefert — mit `ts` statt `timestamp`
und einer Zahl statt eines Strings als id. Das UI haette „Invalid Date"
angezeigt, und zwar NUR im Live-Betrieb, nicht beim manuellen Neuladen:
genau die Sorte Fehler, die man lange sucht.
"""
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 "",
}
@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 [_log_zeile(z) 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.
"""
return [Device(**info) for info in
await asyncio.to_thread(laufwerke_mit_disc)]
#: Zuletzt gemeldeter Grund je Laufwerk. Ohne dieses Gedaechtnis stuende die
#: Zeile alle drei Sekunden im Protokoll und verdraengte alles andere.
_LETZTER_GRUND: Dict[str, str] = {}
def _grund_melden(pfad: str, grund: str) -> None:
"""Warum ein Laufwerk „unknown" meldet — einmal ins Protokoll, beim
Wechsel.
## Der Befund des Commanders (30.08.2026)
> „jetzt erkennt rippy die disk garnicht mehr (im log steht zwar
> erkannt, aber ein start des rips ist nicht moeglich)"
Sein Laufwerk beantwortete nach einem Rip mit Lesefehlern keine
Medien-Abfragen mehr (Win32-Fehler 1), die Geraete-Auskunft aber schon.
Im UI stand deshalb eine vollstaendige Laufwerkskarte mit Modell und
Seriennummer — nur „unknown" bei Typ und Status, und kein Rip startbar.
Die letzte Protokollzeile war das laengst veraltete „Disc erkannt".
Der Treiber kennt den Grund (siehe `drives.windows.ZUGRIFFS_GRUENDE`).
Hier wird er gesagt — samt Abhilfe, denn die Zeile soll nicht nur
beschreiben, sondern weiterhelfen.
"""
if _LETZTER_GRUND.get(pfad, "") == grund:
return
vorher = _LETZTER_GRUND.get(pfad, "")
_LETZTER_GRUND[pfad] = grund
try:
if grund:
db.add_log("warning", "watcher", "Laufwerk %s: %s" % (pfad, grund))
elif vorher:
db.add_log("info", "watcher",
"Laufwerk %s antwortet wieder." % pfad)
except Exception: # noqa: BLE001
pass # Protokollieren darf die Laufwerksliste nie aufhalten
def _job_haelt_das_laufwerk(pfad: str) -> bool:
"""Laeuft auf diesem Laufwerk gerade ein Rip?
## Warum das Laufwerk dann in Ruhe bleiben muss (Befund 30.08.2026)
Der Waechter fragt alle drei Sekunden `device_info` ab — das sind drei
`CreateFileW` plus IOCTLs auf ein Geraet, das waehrenddessen makemkvcon
gehoert. Am Protokoll des Commanders abgelesen:
12:49:52 bluray-Rip gestartet
12:50:09 [watcher] Laufwerk G: beantwortet keine Medien-Abfragen
12:50:12 MSG 2003 SCSI-Fehler ILLEGAL REQUEST:INVALID FIELD IN CDB
12:50:12 MSG 5010 Das Oeffnen der Disk schlug fehl
12:50:12 makemkvcon endete mit Code 11
Sein Befund dazu: „Das laufwerk hoert auch einfach auf zu lesen."
`_auto_prescan` haelt sich seit dem 29.08.2026 an genau diese Regel
(„Es gibt keinen Grund, waehrend eines Rips zu scannen") — die
Laufwerksabfrage tat es nicht. Sie hat dieselbe Begruendung: Wir wissen
bereits, was drinliegt, der Job laeuft ja darauf.
Faellt die Auskunft aus, gilt der letzte bekannte Stand weiter. Das ist
keine Notluege: Waehrend eines Rips aendert sich am Laufwerk nichts.
"""
try:
return bool(db.has_active_job(pfad))
except Exception: # noqa: BLE001
return False # im Zweifel nachsehen, wie bisher
def laufwerke_mit_disc() -> list:
"""Laufwerke SAMT erkannter Disc — der eine Weg für beide Abnehmer.
## Der Befund des Commanders (29.08.2026)
> „Es ist eine Disk im laufwerk, aber er erkennt sie dort nicht"
Auf der Server-Status-Kachel stand „Bereit — kein Datenträger im
Laufwerk", während die Disc erkannt war und ihr Titel eine Zeile weiter
oben in der Jobliste stand.
Die Ursache: Es gab ZWEI Wege zu den Laufwerken. `/devices` hängte die
erkannte Disc aus `DISC_CACHE` an; der Ereignis-Wächter rief nur
`device_info()` und ließ sie weg. Das UI liest seit V2-3 den Ereignisstrom
— also die Fassung ohne Disc. Es fragte `l.disc?.title`, bekam nichts, und
schloss daraus auf ein leeres Laufwerk.
Eine Auskunft in zwei Fassungen ist eine Auskunft zu viel. Jetzt gibt es
nur diese hier.
"""
geraete = []
letzte = {g.get("path"): g for g in LETZTE_LAUFWERKE}
for pfad in device_discovery.list_optical_devices():
if _job_haelt_das_laufwerk(pfad) and pfad in letzte:
# Nicht anfassen — der letzte bekannte Stand gilt weiter.
info = {k: v for k, v in letzte[pfad].items()
if k not in ("disc", "disc_wird_erkannt")}
else:
info = device_discovery.device_info(pfad)
_grund_melden(pfad, info.get("grund") or "")
disc = DISC_CACHE.get(pfad)
if disc and disc.get("_laeuft"):
# NICHT als Disc ausgeben — es gibt noch keinen Titel. Aber
# sagen, dass gerade gearbeitet wird: Zwei Minuten Schweigen
# sehen aus wie ein leeres Laufwerk (Commander 29.08.2026).
info["disc_wird_erkannt"] = True
elif disc:
info["disc"] = disc
geraete.append(info)
LETZTE_LAUFWERKE[:] = geraete
return geraete
#: Der zuletzt vollständig gelesene Laufwerks-Stand. Siehe `laufwerke_notdurft`.
LETZTE_LAUFWERKE: List[dict] = []
def laufwerke_notdurft() -> Optional[list]:
"""Was sich OHNE Laufwerks-Abfrage sagen lässt — oder None.
## Warum es diesen Rückfall gibt (Befund 29.08.2026)
Der Commander wollte sehen, dass die Disc-Erkennung noch läuft. Genau
dann ist die Auskunft aber am teuersten: Während `makemkvcon info` die
Disc liest, hält es das Laufwerk, und `device_info` wartet mit. Gemessen:
/devices waehrend des Scans 14,0 s
Zeitgrenze des Schnappschusses 5,0 s -> devices = None
`None` heißt „konnte nicht nachsehen", und die Oberfläche behält dann
ihren Stand — nach einem frischen Laden also **gar nichts**. Ausgerechnet
in der Phase, die sichtbar werden sollte, war der Bildschirm leer.
Der Ausweg braucht kein ioctl: Dass gerade erkannt wird, steht im
Vorrat (`DISC_CACHE`), und wie die Laufwerke heißen, wissen wir vom
letzten vollständigen Lauf. Beides zusammen ist eine ehrliche Auskunft —
„das sahen wir zuletzt, und an DIESEM Laufwerk arbeiten wir gerade".
`None`, wenn wir noch nie erfolgreich gelesen haben: Dann ist Schweigen
richtig, denn behaupten ließe sich nichts.
"""
bekannt = LETZTE_LAUFWERKE
if not bekannt:
# Kaltstart mitten in der Erkennung: Wir haben noch keinen Stand.
# Die LISTE der Laufwerke ist trotzdem billig (sie zaehlt nur
# Buchstaben bzw. /dev-Knoten auf) — teuer ist erst `device_info`,
# das das Laufwerk selbst anfasst. Also das Wenige melden, das
# sicher ist, und nichts erfinden: kein Typ, kein Modell.
try:
pfade_jetzt = device_discovery.list_optical_devices()
except OSError:
return None
bekannt = [{"id": p.rstrip(":\\").rsplit("\\", 1)[-1].rstrip(":") or p,
"name": p, "path": p, "type": "unknown", "status": "ready"}
for p in pfade_jetzt if (DISC_CACHE.get(p) or {}).get("_laeuft")]
if not bekannt:
return None
frisch = []
for alt in bekannt:
eintrag = dict(alt)
disc = DISC_CACHE.get(eintrag.get("path")) or {}
eintrag.pop("disc_wird_erkannt", None)
eintrag.pop("disc", None)
if disc.get("_laeuft"):
eintrag["disc_wird_erkannt"] = True
elif disc:
eintrag["disc"] = disc
frisch.append(eintrag)
return frisch
# 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.