Drei Praxis-Bugs: tote NAS-Mounts reparierbar, HandBrake-Versionen konsistent, Encoder-Wahl beim Rip
Ampel / ampel (push) Successful in 28s
Ampel / ampel (push) Successful in 28s
1) Speicher-Mounts robust (Befund: toter CIFS-Mount nach NAS-Ausfall/Rebuild —
mounted:false, verschwand aus 'Verfuegbare Ziele', Neu-Anlegen -> 409, man
sass fest):
- mounts.py: ist_erreichbar() (listdir, soft-Mount bricht schnell ab),
ist_gemountet() faengt OSError toter Mounts, aushaengen() mit
umount -l Fallback, reparieren() (lazy abhaengen + frisch mounten).
- /storage-mounts liefert 'reachable'; POST bei existierendem Namen:
aktiv -> 409, tot -> automatische Reparatur mit neuen Angaben;
neuer POST /storage-mounts/{name}/repair (gespeicherte Zugangsdaten).
- /storage-targets crasht nicht mehr an totem Mount (os.path.ismount
OSError abgefangen).
- UI: eigene 'Netzwerk-Mounts'-Liste mit Status (aktiv/nicht erreichbar/
getrennt) + Reparieren- und Entfernen-Knopf — tote Mounts sind sichtbar
und wiederherstellbar statt zu verschwinden.
2) HandBrake-Versionen konsistent (Befund: Docker 1.6.1, Windows-Skript
fest 1.9.2, Update-Check meldet 1.11.2 — verwirrend):
- Windows-Installer zieht jetzt DYNAMISCH die neueste Version (GitHub
latest, Fallback 1.11.2) — passt zum Update-Check.
- Update-UI erklaert klar: Docker = stabiles Debian-Paket (bewusst aelter,
kein Fehler), Windows = neueste. MakeMKV-Update zeigt den Befehl.
3) Encoder-/Worker-Auswahl beim Rip (Feature):
- Celery worker_direct=True: jeder Worker konsumiert zusaetzlich seine
Direkt-Queue. API-Helper transcode_queue(node) routet gezielt an den
gewaehlten Worker, faellt aber sicher auf die geteilte transcode-Queue
zurueck, wenn er offline ist (kein Haengenbleiben).
- /capabilities liefert den Celery-Node je Worker; POST /jobs nimmt
transcode_node (-> Job-meta); rip_disc + retry-transcode routen danach.
- Rip-Dialog: Encoder-/Worker-Dropdown, sichtbar ab 2 Online-Workern.
- ping_worker-Task zum Verifizieren des gezielten Routings.
- Nebenfund gefixt: DeviceDiscovery leitete die Titel-Auswahl (titles)
gar nicht an die API weiter — jetzt titles + transcode_node.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -7,12 +7,33 @@ kein POST /jobs, kein udev-Daemon. Dieser Client schließt die Lücke.
|
||||
import os
|
||||
|
||||
from celery import Celery
|
||||
from celery.utils import worker_direct
|
||||
|
||||
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/0")
|
||||
|
||||
celery_client = Celery("rippy_api", broker=REDIS_URL, backend=REDIS_URL)
|
||||
|
||||
|
||||
def transcode_queue(node: str = None):
|
||||
"""Ziel-Queue für die Kompression: der GEWÄHLTE Worker (worker_direct)
|
||||
wenn er gerade online ist, sonst die geteilte transcode-Queue.
|
||||
|
||||
So kann ein Job gezielt einen Encoder ansprechen — fällt der Worker aber
|
||||
weg, bleibt die Kompression nicht in einer toten Queue hängen, sondern
|
||||
landet bei irgendeinem freien Worker.
|
||||
"""
|
||||
if not node:
|
||||
return "transcode"
|
||||
try:
|
||||
antworten = celery_client.control.ping(timeout=1.0) or []
|
||||
online = {k for antwort in antworten for k in antwort.keys()}
|
||||
if node in online:
|
||||
return worker_direct(node)
|
||||
except Exception:
|
||||
pass
|
||||
return "transcode"
|
||||
|
||||
|
||||
def start_rip(device_path: str, job_id: str, target_dir: str = None):
|
||||
"""Schickt den Rip-Task an den Worker (Task-Name aus worker/tasks.py)."""
|
||||
return celery_client.send_task(
|
||||
|
||||
+102
-30
@@ -338,6 +338,7 @@ class JobCreateRequest(BaseModel):
|
||||
season: Optional[int] = None
|
||||
main_feature_only: Optional[bool] = None # pro Rip; None = Setting gilt
|
||||
titles: Optional[List[int]] = None # exakte Titel-Auswahl (Track-Tabelle)
|
||||
transcode_node: Optional[str] = None # gewählter Encoder-Worker (Celery-Node)
|
||||
|
||||
|
||||
MEDIA_ROOT = "/app/media"
|
||||
@@ -395,6 +396,8 @@ async def create_job(request: JobCreateRequest):
|
||||
titel_liste = sorted({int(t) for t in request.titles if int(t) >= 0})[:200]
|
||||
if titel_liste:
|
||||
meta_dict["titles"] = titel_liste
|
||||
if request.transcode_node:
|
||||
meta_dict["transcode_node"] = request.transcode_node
|
||||
meta_json = json.dumps(meta_dict) if meta_dict else None
|
||||
|
||||
job_id = str(uuid.uuid4())
|
||||
@@ -511,7 +514,15 @@ async def storage_targets():
|
||||
return ziele
|
||||
for name in eintraege:
|
||||
pfad = os.path.join(MEDIA_ROOT, name)
|
||||
if not os.path.isdir(pfad):
|
||||
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)
|
||||
@@ -521,7 +532,7 @@ async def storage_targets():
|
||||
ziele.append({
|
||||
"name": name,
|
||||
"path": pfad,
|
||||
"is_mount": os.path.ismount(pfad),
|
||||
"is_mount": ist_mount,
|
||||
"free_gb": frei_gb,
|
||||
})
|
||||
return ziele
|
||||
@@ -599,10 +610,16 @@ async def retry_transcode(job_id: str):
|
||||
basis = job.get("target_dir") or f"{MEDIA_ROOT}/{job.get('disc_type') or 'bluray'}"
|
||||
final_dir = job.get("output_path") or f"{basis}/{job_id}"
|
||||
|
||||
# An den (beim Rip gewählten) Encoder-Worker routen, sonst geteilte Queue
|
||||
try:
|
||||
meta = json.loads(job.get("meta") or "{}")
|
||||
except ValueError:
|
||||
meta = {}
|
||||
from celery_client import transcode_queue
|
||||
celery_client.send_task(
|
||||
"worker.tasks.transcode_files",
|
||||
args=[job_id, raw_dir, final_dir],
|
||||
queue="transcode",
|
||||
queue=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")
|
||||
@@ -639,18 +656,18 @@ async def capabilities():
|
||||
zeilen = db.list_workers()
|
||||
try:
|
||||
antworten = celery_client.control.ping(timeout=1.0) or []
|
||||
online_namen = {
|
||||
knoten.split("@", 1)[-1]
|
||||
for antwort in antworten
|
||||
for knoten in antwort.keys()
|
||||
}
|
||||
ping_knoten = [k for antwort in antworten for k in antwort.keys()]
|
||||
except Exception:
|
||||
online_namen = set()
|
||||
ping_knoten = []
|
||||
for zeile in zeilen:
|
||||
# Celery-Ping meldet den HOSTNAME des Knotens — bei gesetztem
|
||||
# WORKER_NAME (Anzeigename) steckt der echte Hostname in info.
|
||||
# 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). `node` ist der Routing-Ziel-
|
||||
# Knoten für die gezielte Encoder-Wahl.
|
||||
hostname = (zeile.get("info") or {}).get("hostname") or zeile["name"]
|
||||
zeile["online"] = hostname in online_namen
|
||||
node = next((k for k in ping_knoten if k.split("@", 1)[-1] == hostname), None)
|
||||
zeile["online"] = node is not None
|
||||
zeile["node"] = node
|
||||
return zeilen
|
||||
|
||||
return {"workers": await asyncio.to_thread(sammle)}
|
||||
@@ -806,39 +823,67 @@ class MountRequest(BaseModel):
|
||||
|
||||
@app.get("/storage-mounts")
|
||||
async def get_storage_mounts():
|
||||
"""Konfigurierte Netzwerk-Speicherziele inkl. Live-Mount-Status."""
|
||||
eintraege = await asyncio.to_thread(db.list_mounts)
|
||||
return [
|
||||
{
|
||||
"name": e["name"],
|
||||
"type": e["typ"],
|
||||
"source": e["quelle"],
|
||||
"mounted": mount_verwaltung.ist_gemountet(e["name"]),
|
||||
"has_credentials": bool(e.get("username")),
|
||||
}
|
||||
for e in eintraege
|
||||
]
|
||||
"""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."""
|
||||
"""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")
|
||||
if any(e["name"] == request.name for e in await asyncio.to_thread(db.list_mounts)):
|
||||
raise HTTPException(status_code=409, detail="Name bereits vergeben")
|
||||
|
||||
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(
|
||||
mount_verwaltung.mounten,
|
||||
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,
|
||||
@@ -847,10 +892,37 @@ async def create_storage_mount(request: MountRequest):
|
||||
await asyncio.to_thread(
|
||||
db.add_log,
|
||||
"success" if schreibbar else "warning", "mounts",
|
||||
f"Speicherziel '{request.name}' ({request.type}) eingehängt: {request.source}"
|
||||
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}
|
||||
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")
|
||||
|
||||
+51
-5
@@ -30,7 +30,27 @@ def validiere_name(name: str) -> bool:
|
||||
|
||||
|
||||
def ist_gemountet(name: str) -> bool:
|
||||
return os.path.ismount(_mountpoint(name))
|
||||
"""Ist unter /app/media/<name> etwas gemountet? Ein TOTER CIFS-Mount
|
||||
(NAS weg) lässt os.path.ismount mit OSError fliegen — das zählt weiter
|
||||
als 'gemountet' (nur eben kaputt), damit die Reparatur greift."""
|
||||
try:
|
||||
return os.path.ismount(_mountpoint(name))
|
||||
except OSError:
|
||||
return True
|
||||
|
||||
|
||||
def ist_erreichbar(name: str) -> bool:
|
||||
"""Kann auf das Ziel WIRKLICH zugegriffen werden?
|
||||
|
||||
Ein toter CIFS-Mount (Host weg, stale) ist zwar 'gemountet', aber jeder
|
||||
Zugriff scheitert mit „Host is down" (Befund 24.07.). CIFS läuft mit
|
||||
`soft`, d. h. der Zugriff bricht schnell ab statt zu hängen.
|
||||
"""
|
||||
try:
|
||||
os.listdir(_mountpoint(name))
|
||||
return True
|
||||
except OSError:
|
||||
return False
|
||||
|
||||
|
||||
def schreibtest(pfad: str) -> bool:
|
||||
@@ -169,22 +189,48 @@ def mounten(name: str, typ: str, quelle: str, optionen: str = "",
|
||||
pass
|
||||
|
||||
|
||||
def _lazy_umount(ziel: str) -> None:
|
||||
"""`umount -l`: hängt auch einen TOTEN/beschäftigten Mount ab (detach now,
|
||||
cleanup later). Ohne das ließ sich ein Mount zu einem weggefallenen NAS
|
||||
gar nicht mehr entfernen (Befund 24.07.)."""
|
||||
subprocess.run(["umount", "-l", ziel], capture_output=True, text=True, timeout=30)
|
||||
|
||||
|
||||
def aushaengen(name: str) -> None:
|
||||
ziel = _mountpoint(name)
|
||||
if os.path.ismount(ziel):
|
||||
try:
|
||||
noch_mount = os.path.ismount(ziel)
|
||||
except OSError:
|
||||
noch_mount = True # stale
|
||||
if noch_mount:
|
||||
ergebnis = subprocess.run(
|
||||
["umount", ziel], capture_output=True, text=True, timeout=30
|
||||
)
|
||||
if ergebnis.returncode != 0:
|
||||
raise RuntimeError(
|
||||
f"umount schlug fehl: {(ergebnis.stderr or '').strip()[:300]}"
|
||||
)
|
||||
# Toter/beschäftigter Mount → lazy detach (klappt immer)
|
||||
_lazy_umount(ziel)
|
||||
try:
|
||||
os.rmdir(ziel)
|
||||
except OSError:
|
||||
pass # nicht leer oder weg — egal
|
||||
|
||||
|
||||
def reparieren(name: str, typ: str, quelle: str, optionen: str = "",
|
||||
username: str = "", passwort: str = "") -> bool:
|
||||
"""Toten/veralteten Mount frisch neu verbinden: lazy abhängen, neu mounten.
|
||||
|
||||
Nötig, wenn ein NAS-Mount stale geworden ist (Rebuild, NAS-Schlaf) — das
|
||||
normale mounten() würde am „ist schon Mountpoint" hängenbleiben.
|
||||
"""
|
||||
ziel = _mountpoint(name)
|
||||
try:
|
||||
if os.path.ismount(ziel):
|
||||
_lazy_umount(ziel)
|
||||
except OSError:
|
||||
_lazy_umount(ziel)
|
||||
return mounten(name, typ, quelle, optionen, username, passwort)
|
||||
|
||||
|
||||
def alle_remounten() -> list:
|
||||
"""Beim API-Start: alle gespeicherten Mounts wiederherstellen."""
|
||||
meldungen = []
|
||||
|
||||
Reference in New Issue
Block a user