diff --git a/docker/api/eta.py b/docker/api/eta.py index 83a7292..87513c6 100644 --- a/docker/api/eta.py +++ b/docker/api/eta.py @@ -111,13 +111,44 @@ def schluessel(job_id: str) -> str: return f"eta:{job_id}" +#: Messreihen im eigenen Prozess — der Rückfall, wenn es kein Redis gibt. +#: +#: ## Der Befund des Commanders (29.08.2026) +#: +#: > „Über den kompletten vorgang steht dort ‚Restzeit wird gemessen' aber +#: > messung wird nicht abgeschlossen. Das heißt man hat kein ETA" +#: +#: Die Messreihe lag ausschließlich im Cache, und der Modul-Kopf sagte dazu: +#: „Der Cache (Redis) ist schon da". Im Container stimmt das. Auf einem +#: Windows-PC gibt es kein Redis — `cache_get` gab bei JEDEM Aufruf `None` +#: zurück, `beobachtung_hinzufuegen` legte also jedes Mal eine frische Reihe +#: mit EINEM Punkt an, und `restzeit_sekunden` braucht `MINDEST_PUNKTE = 2`. +#: Ergebnis: über den ganzen Rip hinweg „noch keine Aussage". +#: +#: Ein Wörterbuch im Prozess reicht hier vollkommen: Die Reihe ist ein paar +#: Zahlen, sie gilt nur für die Dauer eines Jobs, und im eigenständigen +#: Betrieb gibt es ohnehin nur diesen einen Prozess. Redis bleibt der bessere +#: Ort, wo es eins gibt — es überlebt einen API-Neustart. +_REIHEN: dict = {} + +#: Wie lange eine Reihe im Prozess aufgehoben wird (wie `expire` im Cache). +REIHE_HALTBAR_SEKUNDEN = 86400 + + +def _aufraeumen(jetzt: float) -> None: + """Alte Reihen wegwerfen — sonst waechst das Woerterbuch unbegrenzt.""" + for k, (stand, _) in list(_REIHEN.items()): + if jetzt - stand > REIHE_HALTBAR_SEKUNDEN: + _REIHEN.pop(k, None) + + def aktualisiere_und_schaetze(job_id: str, status: str, progress: int, jetzt: float, cache_get, cache_set) -> dict: - """Messreihe im Cache fortschreiben und die Restzeit zurückgeben. + """Messreihe fortschreiben und die Restzeit zurückgeben. - Der Cache (Redis) ist schon da und überlebt einen Neustart des UI. Fällt er - aus, kommt bei jedem Aufruf eine leere Reihe zurück — dann gibt es eben - keine ETA, aber nichts scheitert. + Gespeichert wird in BEIDEN Ablagen: im Cache, wo es einen gibt (er + überlebt einen API-Neustart), und im Prozess, damit die Schätzung auch + ohne Redis zustande kommt. Siehe `_REIHEN`. """ if status not in LAUFENDE_STATUS or progress <= 0: return {"sekunden": -1, "text": ""} @@ -126,11 +157,17 @@ def aktualisiere_und_schaetze(job_id: str, status: str, progress: int, reihe = cache_get(k) except Exception: reihe = None + if not reihe: + # Kein Cache (oder leer) — dann der eigene Vorrat. + eintrag = _REIHEN.get(k) + reihe = eintrag[1] if eintrag else None reihe = beobachtung_hinzufuegen(reihe, status, progress, jetzt) try: # Eine Reihe ohne Fortschritt ist nach einem Tag wertlos. - cache_set(k, reihe, expire=86400) + cache_set(k, reihe, expire=REIHE_HALTBAR_SEKUNDEN) except Exception: pass + _REIHEN[k] = (jetzt, reihe) + _aufraeumen(jetzt) sekunden = restzeit_sekunden(reihe, jetzt) return {"sekunden": sekunden, "text": formatiere_restzeit(sekunden)} diff --git a/docker/api/main.py b/docker/api/main.py index 6b01ea2..83deec0 100644 --- a/docker/api/main.py +++ b/docker/api/main.py @@ -519,7 +519,10 @@ async def _snapshot() -> dict: """ def sammeln(): return { - "jobs": [_job_row_to_model(z).model_dump() for z in db.list_jobs(limit=50)], + # 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)], } @@ -947,24 +950,54 @@ async def get_jobs(): Reihe zurück, zwei Tabs zeigten verschiedene Zahlen, und für einen externen Encoder-Worker gäbe es gar keine (genau dort wollte der Commander sie). """ - def sammle(): - work_dir = os.path.normpath((db.get_settings().get("workDir") or "").strip() or "/") - jetzt = time.monotonic() - modelle = [] - for z in db.list_jobs(): - modell = _job_row_to_model(z) - modell.can_retry = _kann_neu_komprimieren(z, work_dir) - 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 + return await asyncio.to_thread(jobs_fuer_ui) - return await asyncio.to_thread(sammle) + +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: diff --git a/docker/api/test_api_smoke.py b/docker/api/test_api_smoke.py index fa6ac01..1b2b167 100644 --- a/docker/api/test_api_smoke.py +++ b/docker/api/test_api_smoke.py @@ -730,3 +730,31 @@ def test_die_wache_setzt_gesehen_auch_nach_einem_fehlschlag_nicht_zurueck(): quelle = inspect.getsource(main.disc_watcher) assert "gesehen" in quelle, "die Wache muss sich merken, was sie je las" assert "disc_entscheidung(" in quelle, "die Entscheidung gehoert in die pure Funktion" + + +def test_schnappschuss_und_jobs_liefern_DIESELBE_jobliste(): + """Waechter gegen die zweite Fassung (Befund 29.08.2026). + + Commander: „Über den kompletten vorgang steht dort ‚Restzeit wird + gemessen' aber messung wird nicht abgeschlossen." + + Die Restzeit wurde nur in `/jobs` berechnet. Der SSE-Schnappschuss baute + seine Jobs mit dem nackten `_job_row_to_model` — also ohne. Seit V2-3 + liest die Oberflaeche aber genau diesen Schnappschuss. Die Restzeit wurde + also brav berechnet und niemandem gezeigt. + """ + import inspect + + import main + + quelle = inspect.getsource(main._snapshot) + assert "jobs_fuer_ui(" in quelle, "der Schnappschuss baut die Jobs selbst" + assert "_job_row_to_model(z).model_dump() for z in db.list_jobs" not in quelle + + +def test_die_jobliste_traegt_die_restzeit_felder(): + """Ohne die Felder im Modell schneidet FastAPI sie weg.""" + import main + + for feld in ("eta_sekunden", "eta_text", "can_retry", "retry_art"): + assert feld in main.Job.model_fields, feld diff --git a/docker/api/test_eta.py b/docker/api/test_eta.py index af016c8..c0aa2ad 100644 --- a/docker/api/test_eta.py +++ b/docker/api/test_eta.py @@ -157,3 +157,73 @@ def test_fertige_und_wartende_jobs_bekommen_keine_eta(): speicher.get, lambda k, v, expire=None: speicher.__setitem__(k, v)) assert ergebnis == {"sekunden": -1, "text": ""} assert speicher == {} # nichts geschrieben + + +# ── Ohne Redis muss die Schaetzung trotzdem zustande kommen ───────────── +# +# Commander 29.08.2026: „Über den kompletten vorgang steht dort ‚Restzeit wird +# gemessen' aber messung wird nicht abgeschlossen. Das heißt man hat kein ETA" +# +# Die Messreihe lag nur im Cache, und der Modul-Kopf sagte „Der Cache (Redis) +# ist schon da". Im Container stimmt das. Auf einem Windows-PC gab `cache_get` +# bei JEDEM Aufruf None zurueck — also jedes Mal eine frische Reihe mit EINEM +# Punkt, und `restzeit_sekunden` braucht zwei. + + +def _ohne_cache(): + """Ein Cache, der nichts behaelt — genau wie Redis, das es nicht gibt.""" + return (lambda k: None), (lambda k, v, expire=None: False) + + +def test_ohne_cache_kommt_trotzdem_eine_schaetzung(): + """DER Fall des Commanders.""" + import eta + + eta._REIHEN.clear() + lesen, schreiben = _ohne_cache() + # Zwei Messpunkte, 120 s auseinander, 10 % Fortschritt. + eta.aktualisiere_und_schaetze("j1", "running", 10, 1000.0, lesen, schreiben) + ergebnis = eta.aktualisiere_und_schaetze("j1", "running", 20, 1120.0, lesen, schreiben) + + assert ergebnis["sekunden"] > 0, "ohne Redis kam nie eine Zahl heraus" + assert ergebnis["text"], "und damit stand dauerhaft „wird gemessen" + + +def test_der_cache_hat_weiterhin_vorrang(): + """Wo es Redis gibt, bleibt Redis die Quelle — es ueberlebt einen + API-Neustart, das Woerterbuch im Prozess nicht.""" + import eta + + eta._REIHEN.clear() + aus_dem_cache = {"status": "running", "punkte": [[500.0, 5], [560.0, 15]]} + ergebnis = eta.aktualisiere_und_schaetze( + "j2", "running", 25, 620.0, + lambda k: aus_dem_cache, lambda k, v, expire=None: True) + # Drei Punkte, 120 s fuer 20 % -> die Reihe aus dem Cache wurde benutzt. + assert ergebnis["sekunden"] > 0 + + +def test_ein_statuswechsel_beginnt_auch_ohne_cache_neu(): + """Rip und Kompression haben nichts miteinander zu tun.""" + import eta + + eta._REIHEN.clear() + lesen, schreiben = _ohne_cache() + eta.aktualisiere_und_schaetze("j3", "running", 50, 1000.0, lesen, schreiben) + eta.aktualisiere_und_schaetze("j3", "running", 90, 1100.0, lesen, schreiben) + neu = eta.aktualisiere_und_schaetze("j3", "transcoding", 5, 1110.0, lesen, schreiben) + assert neu["sekunden"] == -1, "nach dem Wechsel gibt es noch keine Aussage" + + +def test_alte_reihen_wachsen_nicht_unbegrenzt(): + """Wer einen Vorrat anlegt, raeumt ihn auch weg (AGENTS.md).""" + import eta + + eta._REIHEN.clear() + lesen, schreiben = _ohne_cache() + eta.aktualisiere_und_schaetze("alt", "running", 10, 0.0, lesen, schreiben) + # Einen Tag spaeter ein anderer Job -> der alte fliegt raus. + eta.aktualisiere_und_schaetze("neu", "running", 10, + eta.REIHE_HALTBAR_SEKUNDEN + 10.0, lesen, schreiben) + assert eta.schluessel("alt") not in eta._REIHEN + assert eta.schluessel("neu") in eta._REIHEN