fix(api): Rohdaten-Suche in den Hintergrund - /jobs darf nie am NAS haengen
Ampel / ampel (push) Successful in 28s
Ampel / ampel (push) Successful in 28s
Bei der Live-Gegenprobe gemessen: Waehrend der CIFS-Mount nach einem Container-Neustart hochkam, brauchten /jobs und /jobs/<id>/rohdaten jeweils 10,0 Sekunden - genau der CIFS-Timeout. Das Dashboard fragt /jobs alle vier Sekunden ab; ein schlafendes NAS haette es damit dauerhaft eingefroren. Dieselbe Loesung wie beim Celery-Ping in /capabilities (v3.15): eine Hintergrund-Schleife im 30-Sekunden-Takt sieht nach, wo Rohdaten liegen, und /jobs liest nur noch ab. Kennt der Vorrat einen Job noch nicht (frischer Fehlschlag), zaehlt der billige lokale Ort auf der Container-Platte - der antwortet immer sofort. Live geprueft, nachdem der Mount stand: /jobs 0,022 s, /rohdaten 0,014 s, can_retry = true, 79.604.951.639 Bytes (74,1 GB) gefunden. Der 80-GB-Rohschnitt des Commanders ist damit erstmals ueber das UI wieder erreichbar. Der leere Treffer davor war KEIN Codefehler: der Mount stand in dem Moment schlicht noch nicht. Beleg nachgeliefert - im Container gemessen findet die Suche den Pfad in 0,000 s. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
This commit is contained in:
+50
-2
@@ -68,6 +68,10 @@ async def startup_event():
|
|||||||
# Worker-Erreichbarkeit im Hintergrund pingen (siehe _ping_knoten) — sonst
|
# Worker-Erreichbarkeit im Hintergrund pingen (siehe _ping_knoten) — sonst
|
||||||
# kostet JEDER Aufruf von /capabilities eine ganze Sekunde.
|
# kostet JEDER Aufruf von /capabilities eine ganze Sekunde.
|
||||||
asyncio.create_task(_ping_schleife())
|
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())
|
||||||
|
|
||||||
|
|
||||||
# Auto-Pre-Scan-Ergebnisse je Laufwerk: das Dashboard zeigt damit sofort,
|
# Auto-Pre-Scan-Ergebnisse je Laufwerk: das Dashboard zeigt damit sofort,
|
||||||
@@ -310,10 +314,46 @@ async def root():
|
|||||||
|
|
||||||
|
|
||||||
def _rohdaten_suchen(job_id: str, work_dir: str) -> list:
|
def _rohdaten_suchen(job_id: str, work_dir: str) -> list:
|
||||||
"""Wo liegen die Roh-MKVs dieses Jobs? (Details in rohdaten.py)"""
|
"""Wo liegen die Roh-MKVs dieses Jobs? (Details in rohdaten.py)
|
||||||
|
|
||||||
|
Kann LANGSAM sein: Ein schlafendes NAS lässt `os.path.isdir` bis zum
|
||||||
|
CIFS-Timeout hängen — auf dieser VM gemessene 10 Sekunden, während der Mount
|
||||||
|
nach einem Container-Neustart hochkam. Deshalb nur dort direkt aufrufen, wo
|
||||||
|
ein Mensch auf genau diese Antwort wartet (/rohdaten, retry-transcode).
|
||||||
|
Für die Job-Liste gibt es den Vorrat unten.
|
||||||
|
"""
|
||||||
return rohdaten.suche(job_id, work_dir, os.listdir, os.path.isdir)
|
return rohdaten.suche(job_id, work_dir, os.listdir, os.path.isdir)
|
||||||
|
|
||||||
|
|
||||||
|
# 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": {}}
|
||||||
|
ROHDATEN_INTERVALL_SEKUNDEN = 30
|
||||||
|
|
||||||
|
|
||||||
|
async def _rohdaten_schleife():
|
||||||
|
"""Hält den Rohdaten-Vorrat frisch. Darf nie sterben."""
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
await asyncio.to_thread(_rohdaten_vorrat_auffrischen)
|
||||||
|
except Exception: # DB/NAS weg → beim nächsten Durchlauf erneut
|
||||||
|
pass
|
||||||
|
await asyncio.sleep(ROHDATEN_INTERVALL_SEKUNDEN)
|
||||||
|
|
||||||
|
|
||||||
|
def _rohdaten_vorrat_auffrischen() -> None:
|
||||||
|
"""Für jeden fehlgeschlagenen Job nachsehen, wo seine Rohdaten liegen."""
|
||||||
|
work_dir = os.path.normpath((db.get_settings().get("workDir") or "").strip() or "/")
|
||||||
|
treffer = {}
|
||||||
|
for zeile in db.list_jobs():
|
||||||
|
if zeile.get("status") != "failed":
|
||||||
|
continue
|
||||||
|
treffer[zeile["id"]] = _rohdaten_suchen(zeile["id"], work_dir)
|
||||||
|
_ROHDATEN["treffer"] = treffer
|
||||||
|
|
||||||
|
|
||||||
def _kann_neu_komprimieren(job: dict, work_dir: str) -> bool:
|
def _kann_neu_komprimieren(job: dict, work_dir: str) -> bool:
|
||||||
"""Nur wenn Rohdaten wirklich noch daliegen — der „Neu komprimieren"-Knopf
|
"""Nur wenn Rohdaten wirklich noch daliegen — der „Neu komprimieren"-Knopf
|
||||||
an einem Job, der nie gerippt hat, war Unsinn (Befund 24.07.).
|
an einem Job, der nie gerippt hat, war Unsinn (Befund 24.07.).
|
||||||
@@ -325,10 +365,18 @@ def _kann_neu_komprimieren(job: dict, work_dir: str) -> bool:
|
|||||||
auf leer. Ergebnis: `can_retry` war `false`, obwohl 79,6 GB intakt dalagen.
|
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
|
Der SAVEPOINT v3.16 behauptete „‚Neu komprimieren' genügt" — den Knopf gab
|
||||||
es nicht. Jetzt wird an allen möglichen Orten nachgesehen.
|
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":
|
if job.get("status") != "failed":
|
||||||
return False
|
return False
|
||||||
return bool(_rohdaten_suchen(job["id"], work_dir))
|
vorrat = _ROHDATEN["treffer"]
|
||||||
|
if job["id"] in vorrat:
|
||||||
|
return bool(vorrat[job["id"]])
|
||||||
|
return os.path.isdir(os.path.join(rohdaten.RAW_STANDARD, job["id"]))
|
||||||
|
|
||||||
|
|
||||||
@app.get("/jobs", response_model=List[Job])
|
@app.get("/jobs", response_model=List[Job])
|
||||||
|
|||||||
Reference in New Issue
Block a user