This repository has been archived on 2026-07-22. You can view files and clone it. You cannot open issues or pull requests or push a commit.
Files
mission-control/jobengine.py
T
Hitonabi 3cf36d436b feat: echter Download-Fortschritt via Datei-Polling
`hf` gibt im Nicht-TTY-Modus keinen Fortschritt aus (am Bosgame verifiziert:
0 CR-Frames). Stattdessen pollt jobengine.attach_download_progress die
wachsende <local-dir>/.cache/huggingface/download/*.incomplete-Datei gegen die
Gesamtgroesse aus der HF-Tree-API (cookbook.hf_file_size) -> exaktes %.

- attach_download_progress an /api/download, install-recipe, install-model
- Frontend (Aktivitaet + Server-Karte): nutzt job.progress bevorzugt,
  Log-%-Parsing bleibt Fallback fuer Tools, die selbst Prozente ausgeben

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-21 18:26:27 +02:00

176 lines
6.4 KiB
Python

"""
Mini Job-System: Hintergrund-Prozesse mit Live-Log.
Bewusst KISS: ein In-Memory-Dict, ein Daemon-Thread je Job, Subprocess mit
zeilenweisem Log-Capture. Keine Persistenz, kein Broker. Genutzt von allen
Routern, die laenger laufende Shell-Befehle anstossen (Download, Update, ...).
"""
import glob
import os
import shlex
import subprocess
import threading
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"
try:
proc = subprocess.Popen(
args,
stdin=subprocess.PIPE if stdin_data is not None else None,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
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.encode()) # type: ignore[union-attr]
proc.stdin.close() # type: ignore[union-attr]
except Exception: # noqa: BLE001
pass
_pump_output(job, proc.stdout)
proc.wait()
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
_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 attach_download_progress(job_id: str, local_dir: str, total_bytes: int) -> None:
"""Echten Download-Fortschritt (in %) auf den Job legen. `hf` gibt im Nicht-TTY-
Modus keinen Fortschritt aus, schreibt aber in <local_dir>/.cache/huggingface/
download/*.incomplete (waechst). Wir vergleichen dessen Groesse mit total_bytes
(aus der HF-Tree-API). Ein Daemon-Thread aktualisiert job["progress"]."""
if not total_bytes or total_bytes <= 0:
return
job = JOBS.get(job_id)
if job is not None:
job["progress"] = 0
def _watch():
pat = os.path.join(local_dir, ".cache", "huggingface", "download", "*.incomplete")
while True:
j = JOBS.get(job_id)
if not j or j["state"] in ("done", "failed", "canceled"):
break
try:
inc = glob.glob(pat)
cur = sum(os.path.getsize(f) for f in inc) if inc else 0
if cur:
j["progress"] = min(99, int(cur * 100 / total_bytes))
except Exception: # noqa: BLE001
pass
time.sleep(1.0)
j = JOBS.get(job_id)
if j and j["state"] == "done":
j["progress"] = 100
threading.Thread(target=_watch, daemon=True).start()
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]
# log_cmd erlaubt eine sanitisierte Befehlszeile (kein Secret im Log), wenn stdin_data ein
# Geheimnis (sudo-Passwort) traegt. Sonst die echten Argumente.
first_line = log_cmd if log_cmd is not None else "$ " + " ".join(shlex.quote(a) for a in args)
JOBS[job_id] = {
"id": job_id,
"label": label,
"state": "queued",
"log": [first_line],
"returncode": None,
"started_at": time.time(),
"finished_at": None,
}
threading.Thread(target=_run_job, args=(job_id, args, env, stdin_data), daemon=True).start()
return job_id