feat: Live-Download-Fortschritt + Jobs abbrechen
Job-Engine:
- Ausgabe wird binaer + byteweise gelesen und CRLF-bewusst verarbeitet:
einzelnes \r (tqdm/hf-Fortschritt) ueberschreibt die letzte Logzeile statt
sie zu fluten; \r\n bzw. \n = echte neue Zeile. So kommt der Download-
Fortschritt live an. (Text-Modus wuerde \r zu \n uebersetzen -> daher binaer.)
- cancel_job(): laufenden Prozess terminieren, Status "canceled", _PROCS-Registry
- POST /api/jobs/{id}/cancel
Frontend:
- Aktivitaet: Fortschrittsbalken + %-Anzeige aus der letzten Logzeile,
"Abbrechen"-Button je laufendem Job, Status "abgebrochen"
- Server/Aktueller Vorgang: Fortschrittsbalken + Abbrechen (mit Rueckfrage),
"abgebrochen"-Behandlung
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
This commit is contained in:
+82
-10
@@ -14,9 +14,56 @@ import time
|
||||
import uuid
|
||||
|
||||
JOBS: dict[str, dict] = {}
|
||||
_PROCS: dict[str, subprocess.Popen] = {} # laufende Prozesse je Job-ID (fuer Abbruch)
|
||||
_LOG_CAP = 400
|
||||
|
||||
|
||||
def _append_log(job: dict, line: str) -> None:
|
||||
job["log"].append(line)
|
||||
if len(job["log"]) > _LOG_CAP:
|
||||
del job["log"][0]
|
||||
|
||||
|
||||
def _pump_output(job: dict, stream) -> None:
|
||||
"""Liest die Ausgabe byteweise und behandelt `\\r` (Fortschrittsbalken wie bei
|
||||
`hf download`/tqdm) als Ueberschreiben der letzten Zeile statt als neue Zeile.
|
||||
So erscheint der Download-Fortschritt live, ohne das Log mit tausenden Frames zu
|
||||
fluten. Bewusst BINAER gelesen: im Text-Modus wuerde Python ein einzelnes `\\r`
|
||||
zu `\\n` uebersetzen (universal newlines) und das Ueberschreiben unmoeglich machen."""
|
||||
buf = b""
|
||||
overwrite = False # soll die naechste committete Zeile die letzte ueberschreiben?
|
||||
pending_cr = False # haben wir gerade ein `\r` gesehen und warten auf das naechste Byte?
|
||||
|
||||
def commit():
|
||||
line = buf.decode("utf-8", "replace")
|
||||
if overwrite and job["log"]:
|
||||
job["log"][-1] = line
|
||||
else:
|
||||
_append_log(job, line)
|
||||
|
||||
while True:
|
||||
ch = stream.read(1)
|
||||
if not ch:
|
||||
break
|
||||
if pending_cr:
|
||||
pending_cr = False
|
||||
if ch == b"\n": # `\r\n` zusammen = echte neue Zeile
|
||||
commit(); overwrite = False; buf = b""
|
||||
continue
|
||||
commit(); overwrite = True; buf = b"" # einzelnes `\r` = Fortschritt (ueberschreiben)
|
||||
# ch faellt unten normal durch
|
||||
if ch == b"\r":
|
||||
pending_cr = True
|
||||
elif ch == b"\n":
|
||||
commit(); overwrite = False; buf = b""
|
||||
else:
|
||||
buf += ch
|
||||
if pending_cr:
|
||||
commit(); overwrite = True; buf = b""
|
||||
if buf:
|
||||
commit()
|
||||
|
||||
|
||||
def _run_job(job_id: str, args: list[str], env: dict | None = None, stdin_data: str | None = None):
|
||||
job = JOBS[job_id]
|
||||
job["state"] = "running"
|
||||
@@ -26,31 +73,56 @@ def _run_job(job_id: str, args: list[str], env: dict | None = None, stdin_data:
|
||||
stdin=subprocess.PIPE if stdin_data is not None else None,
|
||||
stdout=subprocess.PIPE,
|
||||
stderr=subprocess.STDOUT,
|
||||
text=True,
|
||||
bufsize=1,
|
||||
bufsize=0, # binaer + ungepuffert -> Fortschritt (\r) kommt live an
|
||||
env={**os.environ, **(env or {})},
|
||||
)
|
||||
_PROCS[job_id] = proc
|
||||
if stdin_data is not None:
|
||||
# Secret (z.B. sudo-Passwort) ueber stdin fuettern — landet NICHT im Log/ps.
|
||||
try:
|
||||
proc.stdin.write(stdin_data) # type: ignore[union-attr]
|
||||
proc.stdin.write(stdin_data.encode()) # type: ignore[union-attr]
|
||||
proc.stdin.close() # type: ignore[union-attr]
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
for line in proc.stdout: # type: ignore[union-attr]
|
||||
job["log"].append(line.rstrip("\n"))
|
||||
if len(job["log"]) > _LOG_CAP:
|
||||
del job["log"][0]
|
||||
_pump_output(job, proc.stdout)
|
||||
proc.wait()
|
||||
job["returncode"] = proc.returncode
|
||||
job["state"] = "done" if proc.returncode == 0 else "failed"
|
||||
if job.get("canceled"):
|
||||
job["state"] = "canceled"
|
||||
job["returncode"] = proc.returncode
|
||||
_append_log(job, "[mission-control] Vorgang abgebrochen.")
|
||||
else:
|
||||
job["returncode"] = proc.returncode
|
||||
job["state"] = "done" if proc.returncode == 0 else "failed"
|
||||
except Exception as exc: # noqa: BLE001
|
||||
job["log"].append(f"[mission-control] Fehler: {exc}")
|
||||
_append_log(job, f"[mission-control] Fehler: {exc}")
|
||||
job["state"] = "failed"
|
||||
job["returncode"] = -1
|
||||
finally:
|
||||
_PROCS.pop(job_id, None)
|
||||
job["finished_at"] = time.time()
|
||||
|
||||
|
||||
def cancel_job(job_id: str) -> bool:
|
||||
"""Laufenden Job abbrechen: Prozess terminieren. Liefert False, wenn der Job
|
||||
nicht (mehr) laeuft oder unbekannt ist."""
|
||||
job = JOBS.get(job_id)
|
||||
if not job or job["state"] in ("done", "failed", "canceled"):
|
||||
return False
|
||||
job["canceled"] = True
|
||||
_append_log(job, "[mission-control] Abbruch angefordert…")
|
||||
proc = _PROCS.get(job_id)
|
||||
if proc is not None:
|
||||
try:
|
||||
proc.terminate()
|
||||
except Exception: # noqa: BLE001
|
||||
pass
|
||||
else:
|
||||
# Prozess noch nicht gestartet -> direkt als abgebrochen markieren.
|
||||
job["state"] = "canceled"
|
||||
job["finished_at"] = time.time()
|
||||
return True
|
||||
|
||||
|
||||
def start_job(args: list[str], label: str, env: dict | None = None,
|
||||
stdin_data: str | None = None, log_cmd: str | None = None) -> str:
|
||||
job_id = uuid.uuid4().hex[:12]
|
||||
|
||||
Reference in New Issue
Block a user