""" Mini-Job-System: Hintergrund-Prozesse mit Live-Log + Download-Fortschritt. Portiert aus Mission Control v1 (jobengine.py). In-Memory, ein Daemon-Thread je Job. Seit 24.09.2026: - Jobs laufen in einer eigenen Prozessgruppe; „Abbrechen“ beendet die ganze Gruppe (vorher nur die bash-Hülle — sudo, apt und hf liefen weiter, während der Job schon „abgebrochen“ hieß). - Jede Job-Art hat ein Zeitlimit; ein hängendes apt hält den Wartungs-Riegel nicht mehr ewig. - `start_job_exklusiv` prüft und trägt unter EINER Sperre ein: zwei Klicks starten nicht mehr zwei Updates derselben Gruppe. - Beendete Jobs werden nach einem Tag bzw. ab 40 Stück aufgeräumt. Noch offen (Phase 2): Jobs überleben keinen Neustart von MC2, sie gehören in einen eigenen Worker. """ import glob import os import shlex import signal import subprocess import threading import time import uuid JOBS: dict[str, dict] = {} _PROCS: dict[str, subprocess.Popen] = {} _LOG_CAP = 400 _LOCK = threading.Lock() BEHALTEN_MAX = 40 # so viele beendete Jobs bleiben sichtbar BEHALTEN_S = 24 * 3600 # beendete Jobs verschwinden nach einem Tag STANDARD_ZEITLIMIT_S = 3 * 3600 _ENDE = ("done", "failed", "canceled") 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 byteweise; `\\r` (tqdm/hf-Fortschritt) überschreibt die letzte Zeile.""" buf = b"" overwrite = False pending_cr = False 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": commit(); overwrite = False; buf = b"" continue commit(); overwrite = True; buf = b"" 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 _beende_gruppe(proc: subprocess.Popen) -> None: """Den Job-Prozess samt Kindern beenden (posix: ganze Prozessgruppe, sonst nur den Prozess).""" try: if os.name == "posix": os.killpg(os.getpgid(proc.pid), signal.SIGTERM) else: proc.terminate() except (ProcessLookupError, PermissionError, OSError): pass def _run_job(job_id: str, args: list[str], env: dict | None = None, zeitlimit_s: int | None = None): """Job-Prozess starten und mitschreiben. v3-Umbau P1 (28.08.2026): Auf der Box läuft sudo passwortlos (`NOPASSWD: ALL`), der Job braucht keine stdin-Pipe. DEVNULL sorgt dafür, dass ein Job, der wider Erwarten nach einem Passwort fragt, sofort scheitert statt still zu hängen.""" job = JOBS[job_id] job["state"] = "running" wecker: threading.Timer | None = None try: proc = subprocess.Popen( list(args), stdout=subprocess.PIPE, stderr=subprocess.STDOUT, stdin=subprocess.DEVNULL, bufsize=0, env={**os.environ, **(env or {})}, start_new_session=(os.name == "posix"), ) _PROCS[job_id] = proc grenze = zeitlimit_s or STANDARD_ZEITLIMIT_S def _zu_lang() -> None: job["zeitlimit"] = True _append_log(job, f"[mc] Zeitlimit ({grenze // 60} min) überschritten, Job wird beendet.") _beende_gruppe(proc) wecker = threading.Timer(grenze, _zu_lang) wecker.daemon = True wecker.start() _pump_output(job, proc.stdout) proc.wait() job["returncode"] = proc.returncode if job.get("canceled"): job["state"] = "canceled" elif job.get("zeitlimit"): job["state"] = "failed" job["error"] = "Zeitlimit überschritten" else: job["state"] = "done" if proc.returncode == 0 else "failed" except Exception as exc: _append_log(job, f"[mc] Fehler: {exc}") job["state"] = "failed" job["error"] = str(exc) job["returncode"] = -1 finally: if wecker is not None: wecker.cancel() _PROCS.pop(job_id, None) job["finished_at"] = time.time() cb = job.pop("_on_done", None) if cb and job["state"] == "done": try: cb() except Exception as exc: _append_log(job, f"[mc] Nachbearbeitung-Fehler: {exc}") def attach_download_progress(job_id: str, local_dir: str, total_bytes: int) -> None: """Fortschritt in % aus wachsenden *.incomplete-Dateien (hf schreibt sie).""" if not total_bytes or total_bytes <= 0: return job = JOBS.get(job_id) if job is not None: job["progress"] = 0 job["total_bytes"] = total_bytes def _watch(): pat = os.path.join(local_dir, ".cache", "huggingface", "download", "**", "*.incomplete") prev_t = prev_b = None rate = 0.0 while True: j = JOBS.get(job_id) if not j or j["state"] in _ENDE: break try: inc = glob.glob(pat, recursive=True) 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)) j["done_bytes"] = cur now = time.time() if prev_t is not None and now > prev_t and cur >= prev_b: inst = (cur - prev_b) / (now - prev_t) rate = inst if rate == 0 else 0.3 * inst + 0.7 * rate if rate > 0: j["rate_bps"] = rate j["eta_s"] = int((total_bytes - cur) / rate) prev_t, prev_b = now, cur except Exception: pass time.sleep(1.0) j = JOBS.get(job_id) if j and j["state"] == "done": j["progress"] = 100 j.pop("eta_s", None) threading.Thread(target=_watch, daemon=True).start() def _aufraeumen() -> None: """Alte beendete Jobs entfernen (unter _LOCK aufrufen).""" jetzt = time.time() fertig = sorted((j for j in JOBS.values() if j["state"] in _ENDE), key=lambda j: j.get("finished_at") or 0, reverse=True) for i, j in enumerate(fertig): if i >= BEHALTEN_MAX or jetzt - (j.get("finished_at") or jetzt) > BEHALTEN_S: JOBS.pop(j["id"], None) def _eintragen(args: list[str], label: str, on_done, group: str | None) -> str: job_id = uuid.uuid4().hex[:12] JOBS[job_id] = { "id": job_id, "label": label, "state": "queued", "group": group, "log": ["$ " + " ".join(shlex.quote(a) for a in args)], "returncode": None, "started_at": time.time(), "finished_at": None, } if on_done: JOBS[job_id]["_on_done"] = on_done return job_id def start_job(args: list[str], label: str, env: dict | None = None, on_done=None, group: str | None = None, zeitlimit_s: int | None = None) -> str: with _LOCK: _aufraeumen() job_id = _eintragen(list(args), label, on_done, group) threading.Thread(target=_run_job, args=(job_id, list(args), env, zeitlimit_s), daemon=True).start() return job_id def start_job_exklusiv(group: str, args: list[str], label: str, env: dict | None = None, on_done=None, zeitlimit_s: int | None = None) -> tuple[str | None, dict | None]: """Wie start_job, aber nur, wenn in der Gruppe nichts läuft — Prüfen und Eintragen unter einer Sperre. Rückgabe: (neue Job-ID, None) oder (None, der schon laufende Job).""" with _LOCK: if (laeuft := active_in_group(group)) is not None: return None, laeuft _aufraeumen() job_id = _eintragen(list(args), label, on_done, group) threading.Thread(target=_run_job, args=(job_id, list(args), env, zeitlimit_s), daemon=True).start() return job_id, None def active_in_group(group: str) -> dict | None: """Erster laufender/wartender Job einer Gruppe (z.B. 'maintenance'), sonst None. Basis für den Wartungs-Riegel: nur EIN System-Update gleichzeitig.""" for j in list(JOBS.values()): if j.get("group") == group and j.get("state") in ("running", "queued"): return j return None def cancel_job(job_id: str) -> bool: job = JOBS.get(job_id) if not job or job["state"] in _ENDE: return False job["canceled"] = True _append_log(job, "[mc] Abbruch angefordert…") proc = _PROCS.get(job_id) if proc is not None: _beende_gruppe(proc) else: job["state"] = "canceled" job["finished_at"] = time.time() return True def public_jobs() -> list[dict]: """Jobs ohne interne Felder (_on_done) für die API.""" with _LOCK: _aufraeumen() return [{k: v for k, v in j.items() if not k.startswith("_")} for j in JOBS.values()]