From e1d7d499f8be915befc10bd8966fb6147db6768f Mon Sep 17 00:00:00 2001 From: Hitonabi Date: Thu, 24 Sep 2026 17:41:26 +0200 Subject: [PATCH] phase2b: Auftraege ueberleben einen Neustart von MC2 - jobengine: jeder Auftrag laeuft als eigene systemd-Einheit mc2-job- (systemd-run, RuntimeMaxSec als Zeitlimit); Akte, Protokoll und Exit-Code liegen unter /mc2-jobs; MC2 nimmt laufende Auftraege beim Start wieder auf - Geheimnisse (HF_TOKEN) ueber eine nur fuer den Nutzer lesbare Umgebungsdatei, die der Auftrag beim Start liest und loescht; nichts davon in Befehlszeile oder Einheit - benannte Nacharbeiten (wartung:nach_update, modell:rolle) statt Closures, laufen auch nach einem Neustart - ohne systemd (PC, Tests) weiter als Kindprozess - deploy.sh wartet nur noch auf Update-Auftraege; Downloads laufen weiter - 14 neue Tests; systemd-Weg auf der Box echt geprueft (Neustart, Abbruch, Zeitlimit, Geheimnis nicht sichtbar) Co-Authored-By: Claude Opus 5.5 --- backend/app.py | 5 + backend/routers/events.py | 2 +- backend/routers/models.py | 33 +- backend/services/jobengine.py | 571 +++++++++++++++++++++++--------- backend/services/maintenance.py | 19 +- backend/tests/test_jobengine.py | 196 +++++++++++ deploy/deploy.sh | 8 +- 7 files changed, 657 insertions(+), 177 deletions(-) create mode 100644 backend/tests/test_jobengine.py diff --git a/backend/app.py b/backend/app.py index cf2f75d..4786e39 100644 --- a/backend/app.py +++ b/backend/app.py @@ -66,6 +66,11 @@ def _box_hintergrund() -> list[asyncio.Task[Any]]: @asynccontextmanager async def lifespan(app: FastAPI) -> AsyncGenerator[None, None]: """Hintergrund-Tasks an den App-Lebenszyklus binden.""" + # Aufträge (Updates, Downloads) laufen als eigene Einheiten weiter, während MC2 neu startet — + # hier werden sie wieder beobachtet. Ein Probelauf daneben fasst sie nicht an. + if os.environ.get("MC_PROBELAUF", "") != "1": + from services import jobengine + await asyncio.to_thread(jobengine.wiederaufnehmen) tasks = _box_hintergrund() if ROLLE == "box" else [] # Geteilter HTTP-Client zur lokalen Engine: Keep-Alive/Connection-Pooling statt neuer Client # pro /v1-Anfrage (spart Sockets/TIME_WAIT unter parallelen Agent-Strömen von Zed/Kilo). diff --git a/backend/routers/events.py b/backend/routers/events.py index 05ad84c..fb732ef 100644 --- a/backend/routers/events.py +++ b/backend/routers/events.py @@ -64,7 +64,7 @@ def _fingerprints() -> dict[str, object]: try: # Jobs (Downloads, Wartung): Zustand + Fortschritt — spart den 3-s-Poller der Schublade from services import jobengine fp["jobs"] = json.dumps( - [(j.get("id"), j.get("state"), j.get("progress")) for j in jobengine.public_jobs()] + [(j.get("id"), j.get("state"), j.get("progress")) for j in jobengine.public_jobs(mit_log=False)] ) except Exception: pass diff --git a/backend/routers/models.py b/backend/routers/models.py index 7c9e654..f3500bd 100644 --- a/backend/routers/models.py +++ b/backend/routers/models.py @@ -55,6 +55,20 @@ def register(req: RegisterReq) -> dict: return {"ok": True, "model_id": model_id} +@jobengine.nacharbeit("modell:rolle") +def _rolle_uebernehmen(model_id: str, role: str) -> None: + """Nach dem Download die gewünschte Rolle setzen. hermes geht dabei durch den warm-bewussten + Hirn-Flow (brains-Gruppe, ttl 0, Hermes-Config, Gateway-Restart) — der rohe Alias-Move hatte + diesen Flow bisher umgangen.""" + if role == "hermes": + from services.agent import set_agent_brain + res = set_agent_brain(model_id) + if not res.get("ok"): + raise RuntimeError(res.get("reason", "Hirn-Wechsel fehlgeschlagen")) + else: + llamaswap.set_role(model_id, role) + + class InstallReq(BaseModel): repo: str role: str | None = None @@ -109,20 +123,9 @@ def install(req: InstallReq) -> dict: except PermissionError as exc: raise HTTPException(500, str(exc)) - # Rolle erst NACH erfolgreichem Download übernehmen. hermes geht dabei durch den - # warm-bewussten Hirn-Flow (brains-Gruppe, ttl 0, Hermes-Config, Gateway-Restart) — - # der rohe Alias-Move hatte diesen Flow bisher umgangen. + # Rolle erst NACH erfolgreichem Download übernehmen (_rolle_uebernehmen). role = (req.role or "").strip().lower() - def _apply_role() -> None: - if role == "hermes": - from services.agent import set_agent_brain - res = set_agent_brain(model_id) - if not res.get("ok"): - raise RuntimeError(res.get("reason", "Hirn-Wechsel fehlgeschlagen")) - else: - llamaswap.set_role(model_id, role) - # Download-Job: alle GGUF-Teile (+ mmproj) per --include holen. args = [hf.hf_bin(), "download", repo] args.extend(info["files"]) @@ -133,8 +136,10 @@ def install(req: InstallReq) -> dict: if token := geheimnisse.hf_token(): env["HF_TOKEN"] = token # Gruppe „download“: so taucht der Download in der Oberfläche mit Fortschritt und Abbrechen auf. - job_id = jobengine.start_job(args, f"Download {repo}", env=env, group="download", - on_done=_apply_role if role else None, zeitlimit_s=6 * 3600) + # Der Download läuft als eigener Auftrag weiter, auch wenn MC2 zwischendurch neu startet. + job_id = jobengine.start_job(args, f"Download {repo}", env=env, group="download", zeitlimit_s=6 * 3600, + nacharbeit="modell:rolle" if role else None, + nacharbeit_daten={"model_id": model_id, "role": role} if role else None) jobengine.attach_download_progress(job_id, str(target), info["total_bytes"]) return {"ok": True, "job_id": job_id, "model_id": model_id, "model_path": model_path, "total_bytes": info["total_bytes"], "files": len(info["files"])} diff --git a/backend/services/jobengine.py b/backend/services/jobengine.py index 26abb40..80fd77a 100644 --- a/backend/services/jobengine.py +++ b/backend/services/jobengine.py @@ -1,80 +1,251 @@ """ -Mini-Job-System: Hintergrund-Prozesse mit Live-Log + Download-Fortschritt. -Portiert aus Mission Control v1 (jobengine.py). In-Memory, ein Daemon-Thread je Job. +Auftrags-System: Hintergrund-Aufträge (Updates, Modell-Downloads) mit Protokoll und Fortschritt. +Ursprünglich aus Mission Control v1 portiert. -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. +Seit Phase 2 (24.09.2026) überleben Aufträge einen Neustart von MC2: +- Jeder Auftrag läuft als eigene systemd-Einheit „mc2-job-“ (systemd-run), nicht mehr als + Kindprozess von MC2. Ein Deploy oder Absturz von MC2 würgt ihn nicht mehr ab. +- Akte (.json), Protokoll (.log) und Exit-Code (.exit) liegen unter JOBS_DIR. + MC2 liest die Akten beim Start wieder ein (wiederaufnehmen()) und beobachtet weiter. +- Nacharbeiten („Rolle setzen, wenn der Download fertig ist“) sind benannt (@nacharbeit) und + stehen mit ihren Daten in der Akte — sie laufen auch, wenn MC2 zwischendurch neu gestartet hat. +- Geheimnisse (z. B. HF_TOKEN) gehen über eine Umgebungsdatei (nur für den Nutzer lesbar), die + der Auftrag beim Start liest und löscht — nicht über die Befehlszeile oder die Einheit. +- Ohne systemd (Entwicklungsrechner, Tests: MC_JOBS_ART=prozess) läuft der Auftrag als + Kindprozess; dann überlebt er einen Neustart nicht. -Noch offen (Phase 2): Jobs überleben keinen Neustart von MC2, sie gehören in einen eigenen Worker. +Weiterhin: „Abbrechen“ beendet den ganzen Auftrag samt Kindern, jede Auftragsart hat ein +Zeitlimit, start_job_exklusiv prüft und trägt unter EINER Sperre ein, beendete Aufträge +verschwinden nach einem Tag bzw. ab 40 Stück. + +Nur MC2 selbst beobachtet Aufträge (wiederaufnehmen() im Lebenszyklus der App); andere Prozesse +starten keine. """ import glob +import json +import logging import os +import re import shlex +import shutil import signal import subprocess import threading import time import uuid +from collections.abc import Callable +from pathlib import Path + +from kern.einstellungen import einstellungen + +log = logging.getLogger(__name__) 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 +_LOG_CAP = 400 +_LOG_LESEN_MAX = 256_000 # so viel vom Protokollende wird für die Anzeige gelesen +BEHALTEN_MAX = 40 # so viele beendete Aufträge bleiben sichtbar +BEHALTEN_S = 24 * 3600 # beendete Aufträge verschwinden nach einem Tag STANDARD_ZEITLIMIT_S = 3 * 3600 _ENDE = ("done", "failed", "canceled") +_TAKT_S = 1.0 # wie oft ein Beobachter nach dem Exit-Code schaut +_EINHEIT_PRUEFEN_S = 5.0 # wie oft er zusätzlich fragt, ob die Einheit noch lebt + +# Nur im Speicher (nicht in der Akte): Fortschritt eines Downloads. +_LAUFZEIT = {"progress", "done_bytes", "rate_bps", "eta_s"} +# Nicht an die Oberfläche: Verwaltung des Auftrags. +_INTERN = {"art", "nacharbeit", "nacharbeit_daten", "nacharbeit_erledigt", "fortschritt", "zeitlimit_s"} +# Nicht in die Umgebung des Auftrags: gehört zur Einheit von MC2 oder ist in bash schreibgeschützt. +_UMGEBUNG_OHNE = {"INVOCATION_ID", "JOURNAL_STREAM", "SYSTEMD_EXEC_PID", "MANAGERPID", "NOTIFY_SOCKET", + "LISTEN_PID", "LISTEN_FDS", "LISTEN_FDNAMES", "WATCHDOG_PID", "WATCHDOG_USEC", "MAINPID", + "_", "SHLVL", "PWD", "OLDPWD", "SHELLOPTS", "BASHOPTS", "BASH_VERSINFO", "EUID", "UID", + "PPID", "GROUPS"} +_NAME = re.compile(r"^[A-Za-z_][A-Za-z0-9_]*$") + +# Die Hülle um den eigentlichen Befehl: Umgebung lesen (und die Datei löschen), Ausgabe ins +# Protokoll, Exit-Code atomar ablegen. $1 = Pfadpräfix der Akte, danach der Befehl. +_HUELLE = ('akte="$1"; shift; ' + 'if [ -f "$akte.env" ]; then . "$akte.env"; rm -f "$akte.env"; fi; ' + 'exec >>"$akte.log" 2>&1 "$akte.exit.tmp" && mv -f "$akte.exit.tmp" "$akte.exit"') + +_NACHARBEITEN: dict[str, Callable[..., None]] = {} +_geladen = False -def _append_log(job: dict, line: str) -> None: - job["log"].append(line) - if len(job["log"]) > _LOG_CAP: - del job["log"][0] +def nacharbeit(name: str) -> Callable[[Callable[..., None]], Callable[..., None]]: + """Eine Nacharbeit unter festem Namen anmelden. Sie bekommt die nacharbeit_daten des Auftrags + als Schlüsselwort-Argumente und läuft nur, wenn der Auftrag erfolgreich war.""" + def anmelden(fn: Callable[..., None]) -> Callable[..., None]: + _NACHARBEITEN[name] = fn + return fn + return anmelden -def _pump_output(job: dict, stream) -> None: - """Liest byteweise; `\\r` (tqdm/hf-Fortschritt) überschreibt die letzte Zeile.""" - buf = b"" - overwrite = False - pending_cr = False +# --- Ablage ------------------------------------------------------------------------ - def commit(): - line = buf.decode("utf-8", "replace") - if overwrite and job["log"]: - job["log"][-1] = line - else: - _append_log(job, line) +def _jobs_dir() -> Path: + return Path(os.environ.get("MC_JOBS_DIR", str(einstellungen().daten_dir / "mc2-jobs"))) - 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 _pfad(job_id: str, endung: str) -> Path: + return _jobs_dir() / f"{job_id}{endung}" + + +def _speichern(job: dict) -> None: + """Akte ohne Laufzeitfelder atomar schreiben.""" + daten = {k: v for k, v in job.items() if k not in _LAUFZEIT} + try: + _jobs_dir().mkdir(parents=True, exist_ok=True) + tmp = _pfad(job["id"], ".json.tmp") + tmp.write_text(json.dumps(daten, ensure_ascii=False), encoding="utf-8") + tmp.replace(_pfad(job["id"], ".json")) + except OSError: + log.warning("jobengine: Akte %s nicht schreibbar", job.get("id"), exc_info=True) + + +def _protokoll(job_id: str, zeile: str) -> None: + """Eine Zeile von MC2 ins Protokoll des Auftrags schreiben.""" + try: + _jobs_dir().mkdir(parents=True, exist_ok=True) + with _pfad(job_id, ".log").open("a", encoding="utf-8", newline="\n") as f: + f.write(zeile.rstrip("\n") + "\n") + except OSError: + pass + + +def zeilen_aus(text: str) -> list[str]: + """Protokolltext in Anzeigezeilen: `\\r` (Fortschrittsbalken von hf/tqdm/apt) überschreibt die Zeile.""" + zeilen = [] + for roh in text.split("\n"): + if "\r" in roh: + teile = [t for t in roh.split("\r") if t] + roh = teile[-1] if teile else "" + zeilen.append(roh) + if zeilen and zeilen[-1] == "": + zeilen.pop() + return zeilen[-_LOG_CAP:] + + +def _protokoll_lesen(job_id: str) -> list[str]: + try: + with _pfad(job_id, ".log").open("rb") as f: + f.seek(0, os.SEEK_END) + groesse = f.tell() + f.seek(max(0, groesse - _LOG_LESEN_MAX)) + roh = f.read() + except OSError: + return [] + text = roh.decode("utf-8", "replace") + if groesse > _LOG_LESEN_MAX: + text = text.split("\n", 1)[-1] # angeschnittene erste Zeile verwerfen + return zeilen_aus(text) + + +def _exit_code(job_id: str) -> int | None: + try: + return int(_pfad(job_id, ".exit").read_text(encoding="utf-8").strip()) + except (OSError, ValueError): + return None + + +def _exit_ablegen(job_id: str, code: int) -> None: + try: + tmp = _pfad(job_id, ".exit.tmp") + tmp.write_text(str(code), encoding="utf-8") + tmp.replace(_pfad(job_id, ".exit")) + except OSError: + pass + + +def _umgebungsdatei(job_id: str, env: dict | None) -> None: + """Umgebung des Auftrags (MC2-Umgebung + env) als nur für den Nutzer lesbare Datei.""" + werte = {**os.environ, **(env or {})} + zeilen = [f"export {k}={shlex.quote(str(v))}" for k, v in werte.items() + if _NAME.match(k) and k not in _UMGEBUNG_OHNE] + fd = os.open(_pfad(job_id, ".env"), os.O_WRONLY | os.O_CREAT | os.O_TRUNC, 0o600) + with os.fdopen(fd, "w", encoding="utf-8", newline="\n") as f: + f.write("\n".join(zeilen) + "\n") + + +# --- Ausführung ---------------------------------------------------------------------- + +def _art() -> str: + """„systemd“, wenn Aufträge als eigene Einheiten laufen können, sonst „prozess“.""" + gewuenscht = os.environ.get("MC_JOBS_ART", "").strip().lower() + if gewuenscht in ("systemd", "prozess"): + return gewuenscht + if os.name != "posix" or not shutil.which("systemd-run"): + return "prozess" + if os.geteuid() != 0 and not os.environ.get("XDG_RUNTIME_DIR"): + return "prozess" + return "systemd" + + +def _systemd(programm: str, *args: str) -> list[str]: + """systemctl/systemd-run für den richtigen Manager: als root der System-, sonst der Nutzer-Manager.""" + als_root = os.name == "posix" and os.geteuid() == 0 + return [programm, *(() if als_root else ("--user",)), *args] + + +def _einheit(job_id: str) -> str: + return f"mc2-job-{job_id}" + + +def _einheit_lebt(job_id: str) -> bool: + try: + r = subprocess.run(_systemd("systemctl", "is-active", _einheit(job_id)), + capture_output=True, text=True, timeout=10) + except Exception: + return True # Im Zweifel weiter beobachten statt einen laufenden Auftrag totzusagen. + return r.stdout.strip() in ("active", "activating", "deactivating", "reloading") + + +def _starte_systemd(job: dict, args: list[str], env: dict | None) -> None: + job_id = job["id"] + _umgebungsdatei(job_id, env) + befehl = _systemd("systemd-run", "--unit", _einheit(job_id), "--collect", "--quiet", + "--working-directory", os.getcwd(), + "--property", f"RuntimeMaxSec={int(job['zeitlimit_s'])}", + "--", "/bin/bash", "-c", _HUELLE, "mc2-job", str(_jobs_dir() / job_id), *args) + r = subprocess.run(befehl, capture_output=True, text=True, timeout=30) + if r.returncode != 0: + _pfad(job_id, ".env").unlink(missing_ok=True) + raise RuntimeError((r.stderr or r.stdout or "systemd-run scheiterte").strip()[:300]) + threading.Thread(target=_beobachten, args=(job_id,), daemon=True).start() + + +def _starte_prozess(job: dict, args: list[str], env: dict | None) -> None: + """Rückfall ohne systemd: Kindprozess, Ausgabe ins Protokoll; ein Warte-Thread schließt ab.""" + job_id = job["id"] + with _pfad(job_id, ".log").open("ab") as protokoll: + proc = subprocess.Popen(list(args), stdout=protokoll, stderr=subprocess.STDOUT, stdin=subprocess.DEVNULL, + env={**os.environ, **(env or {})}, start_new_session=(os.name == "posix")) + _PROCS[job_id] = proc + grenze = int(job["zeitlimit_s"]) + + def zu_lang() -> None: + job["zeitlimit"] = True + _protokoll(job_id, f"[mc] Zeitlimit ({grenze // 60} min) überschritten, Auftrag wird beendet.") + _beende_gruppe(proc) + + wecker = threading.Timer(grenze, zu_lang) + wecker.daemon = True + wecker.start() + + def warten() -> None: + code = proc.wait() + wecker.cancel() + _exit_ablegen(job_id, code) + _abschliessen(job, code) + + threading.Thread(target=warten, daemon=True).start() def _beende_gruppe(proc: subprocess.Popen) -> None: - """Den Job-Prozess samt Kindern beenden (posix: ganze Prozessgruppe, sonst nur den Prozess).""" + """Den Auftragsprozess samt Kindern beenden (posix: ganze Prozessgruppe, sonst nur den Prozess).""" try: if os.name == "posix": os.killpg(os.getpgid(proc.pid), signal.SIGTERM) @@ -84,71 +255,97 @@ def _beende_gruppe(proc: subprocess.Popen) -> None: 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 +def _abschliessen(job: dict, code: int | None) -> None: + """Endzustand festhalten (genau einmal) und, wenn erfolgreich, die Nacharbeit ausführen.""" + with _LOCK: + if job["state"] in _ENDE: + return + job["returncode"] = code + job["finished_at"] = time.time() if job.get("canceled"): job["state"] = "canceled" - elif job.get("zeitlimit"): + elif job.get("zeitlimit") or (code is None and time.time() - job["started_at"] >= job["zeitlimit_s"] - 5): + job.update(state="failed", zeitlimit=True, error="Zeitlimit überschritten") + elif code is None: job["state"] = "failed" - job["error"] = "Zeitlimit überschritten" + job.setdefault("error", "Der Auftrag endete ohne Rückmeldung.") 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}") + job["state"] = "done" if code == 0 else "failed" + _PROCS.pop(job["id"], None) + _speichern(job) + if job["state"] == "done": + _nacharbeit_ausfuehren(job) +def _nacharbeit_ausfuehren(job: dict) -> None: + name = job.get("nacharbeit") + if not name or job.get("nacharbeit_erledigt"): + return + job["nacharbeit_erledigt"] = True + _speichern(job) + try: + fn = _NACHARBEITEN.get(name) + if fn is None: + raise RuntimeError(f"Nacharbeit „{name}“ ist in dieser Instanz unbekannt") + fn(**(job.get("nacharbeit_daten") or {})) + except Exception as exc: + _protokoll(job["id"], f"[mc] Nachbearbeitung-Fehler: {exc}") + log.warning("jobengine: Nacharbeit %s von %s gescheitert", name, job["id"], exc_info=True) + + +def _beobachten(job_id: str) -> None: + """Wartet auf den Exit-Code einer systemd-Einheit. Endet sie ohne einen (Abbruch, Zeitlimit, + Absturz), schließt er den Auftrag trotzdem ab.""" + naechste_pruefung = time.time() + _EINHEIT_PRUEFEN_S + while True: + job = JOBS.get(job_id) + if job is None or job["state"] in _ENDE: + return + if (code := _exit_code(job_id)) is not None: + _abschliessen(job, code) + return + if time.time() >= naechste_pruefung: + if not _einheit_lebt(job_id): + time.sleep(0.5) # der Exit-Code kann gerade erst geschrieben werden + _abschliessen(job, _exit_code(job_id)) + return + naechste_pruefung = time.time() + _EINHEIT_PRUEFEN_S + time.sleep(_TAKT_S) + + +def _starten(job: dict, args: list[str], env: dict | None) -> None: + job["state"] = "running" + _speichern(job) + try: + if job["art"] == "systemd": + _starte_systemd(job, args, env) + else: + _starte_prozess(job, args, env) + except Exception as exc: + _protokoll(job["id"], f"[mc] Fehler: {exc}") + job["error"] = str(exc) + _abschliessen(job, -1) + + +# --- Fortschritt (Downloads) ------------------------------------------------------------ + 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 + if job is None or not total_bytes or total_bytes <= 0: + return + job["fortschritt"] = {"ordner": local_dir, "gesamt": total_bytes} + job["total_bytes"] = total_bytes + job["progress"] = 0 + _speichern(job) + _fortschritt_beobachten(job_id) + + +def _fortschritt_beobachten(job_id: str) -> None: + ziel = (JOBS.get(job_id) or {}).get("fortschritt") or {} + local_dir, total_bytes = ziel.get("ordner"), ziel.get("gesamt") + if not local_dir or not total_bytes: + return def _watch(): pat = os.path.join(local_dir, ".cache", "huggingface", "download", "**", "*.incomplete") @@ -183,76 +380,150 @@ def attach_download_progress(job_id: str, local_dir: str, total_bytes: int) -> N threading.Thread(target=_watch, daemon=True).start() +# --- Verwaltung ------------------------------------------------------------------- + +def _laden() -> None: + """Akten von der Platte einlesen (einmal je Prozess; unter _LOCK aufrufen).""" + global _geladen + if _geladen: + return + _geladen = True + for datei in sorted(_jobs_dir().glob("*.json")): + try: + akte = json.loads(datei.read_text(encoding="utf-8")) + except (OSError, ValueError): + continue + if isinstance(akte, dict) and akte.get("id") and akte["id"] not in JOBS: + JOBS[akte["id"]] = akte + + def _aufraeumen() -> None: - """Alte beendete Jobs entfernen (unter _LOCK aufrufen).""" + """Alte beendete Aufträge samt Dateien 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) + for endung in (".json", ".log", ".exit", ".env", ".exit.tmp", ".json.tmp"): + _pfad(j["id"], endung).unlink(missing_ok=True) -def _eintragen(args: list[str], label: str, on_done, group: str | None) -> str: +def _eintragen(args: list[str], label: str, group: str | None, zeitlimit_s: int | None, + nacharbeit_name: str | None, nacharbeit_daten: dict | None) -> dict: + """Neue Akte anlegen (unter _LOCK aufrufen).""" + if nacharbeit_name and nacharbeit_name not in _NACHARBEITEN: + raise ValueError(f"Unbekannte Nacharbeit: {nacharbeit_name}") 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)], + job = { + "id": job_id, "label": label, "state": "queued", "group": group, "art": _art(), "returncode": None, "started_at": time.time(), "finished_at": None, + "zeitlimit_s": int(zeitlimit_s or STANDARD_ZEITLIMIT_S), + "nacharbeit": nacharbeit_name, "nacharbeit_daten": nacharbeit_daten or {}, } - if on_done: - JOBS[job_id]["_on_done"] = on_done - return job_id + JOBS[job_id] = job + _protokoll(job_id, "$ " + " ".join(shlex.quote(a) for a in args)) + _speichern(job) + return job -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.""" +def _aktiv_in(group: str) -> dict | None: for j in list(JOBS.values()): if j.get("group") == group and j.get("state") in ("running", "queued"): return j return None +def start_job(args: list[str], label: str, env: dict | None = None, group: str | None = None, + zeitlimit_s: int | None = None, nacharbeit: str | None = None, + nacharbeit_daten: dict | None = None) -> str: + with _LOCK: + _laden() + _aufraeumen() + job = _eintragen(list(args), label, group, zeitlimit_s, nacharbeit, nacharbeit_daten) + _starten(job, list(args), env) + return job["id"] + + +def start_job_exklusiv(group: str, args: list[str], label: str, env: dict | None = None, + zeitlimit_s: int | None = None, nacharbeit: str | None = None, + nacharbeit_daten: dict | 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 Auftrags-ID, None) oder (None, der schon laufende Auftrag).""" + with _LOCK: + _laden() + if (laeuft := _aktiv_in(group)) is not None: + return None, laeuft + _aufraeumen() + job = _eintragen(list(args), label, group, zeitlimit_s, nacharbeit, nacharbeit_daten) + _starten(job, list(args), env) + return job["id"], None + + +def active_in_group(group: str) -> dict | None: + """Erster laufender/wartender Auftrag einer Gruppe (z. B. 'maintenance'), sonst None. + Basis für den Wartungs-Riegel: nur EIN System-Update gleichzeitig.""" + with _LOCK: + _laden() + return _aktiv_in(group) + + +def wiederaufnehmen() -> int: + """Beim Start von MC2: Akten einlesen und laufende Aufträge weiter beobachten. + Rückgabe: Zahl der wieder aufgenommenen Aufträge.""" + with _LOCK: + _laden() + offen = [j for j in JOBS.values() if j["state"] in ("queued", "running")] + nacharbeit_offen = [j for j in JOBS.values() if j["state"] == "done" and j.get("nacharbeit") + and not j.get("nacharbeit_erledigt")] + for job in offen: + if job.get("art") == "systemd": + threading.Thread(target=_beobachten, args=(job["id"],), daemon=True).start() + _fortschritt_beobachten(job["id"]) + elif (code := _exit_code(job["id"])) is not None: + _abschliessen(job, code) + else: + # Kindprozess des vorigen MC2: sein Warte-Thread ist mit ihm gegangen. + _protokoll(job["id"], "[mc] MC2 wurde neu gestartet; der Auftrag lief als Kindprozess mit.") + job["error"] = "MC2 wurde während des Auftrags neu gestartet." + _abschliessen(job, None) + for job in nacharbeit_offen: + _nacharbeit_ausfuehren(job) + if offen: + log.info("jobengine: %d laufende Aufträge wieder aufgenommen", len(offen)) + return len(offen) + + 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: + with _LOCK: + _laden() + job = JOBS.get(job_id) + if not job or job["state"] in _ENDE: + return False + job["canceled"] = True + _speichern(job) + _protokoll(job_id, "[mc] Abbruch angefordert…") + if job["state"] == "queued": + _abschliessen(job, None) + elif job.get("art") == "systemd": + try: + subprocess.run(_systemd("systemctl", "stop", "--no-block", _einheit(job_id)), + capture_output=True, timeout=15) + except Exception: + log.warning("jobengine: Abbruch von %s gescheitert", job_id, exc_info=True) + elif (proc := _PROCS.get(job_id)) is not None: _beende_gruppe(proc) else: - job["state"] = "canceled" - job["finished_at"] = time.time() + _abschliessen(job, None) return True -def public_jobs() -> list[dict]: - """Jobs ohne interne Felder (_on_done) für die API.""" +def public_jobs(mit_log: bool = True) -> list[dict]: + """Aufträge ohne interne Felder für die API; das Protokoll kommt aus der Datei (Ende).""" with _LOCK: + _laden() _aufraeumen() - return [{k: v for k, v in j.items() if not k.startswith("_")} for j in JOBS.values()] + jobs = [{k: v for k, v in j.items() if k not in _INTERN} for j in JOBS.values()] + for j in jobs: + j["log"] = _protokoll_lesen(j["id"]) if mit_log else [] + return sorted(jobs, key=lambda j: j.get("started_at") or 0) diff --git a/backend/services/maintenance.py b/backend/services/maintenance.py index 7bdda13..4fe45a0 100644 --- a/backend/services/maintenance.py +++ b/backend/services/maintenance.py @@ -246,6 +246,7 @@ def updates() -> dict: return daten +@jobengine.nacharbeit("wartung:nach_update") def _nach_update() -> None: """Nach einem abgeschlossenen Update: Zwischenspeicher leeren und den Stand hochzählen.""" global update_stand @@ -566,12 +567,12 @@ def logs(service: str, lines: int = 200) -> dict: return {"ok": r["ok"], "text": r["out"] or r["err"]} return {"ok": False, "text": "", "err": "Dienst nicht erlaubt."} -def _starten(args: list[str], label: str, on_done=None, zeitlimit_s: int | None = None) -> dict: +def _starten(args: list[str], label: str, zeitlimit_s: int | None = None) -> dict: """Wartungs-Riegel: nur EIN binär-/dienst-veränderndes Job gleichzeitig. Prüfen und Eintragen passieren unter einer Sperre (jobengine.start_job_exklusiv) — vorher lag zwischen Prüfung und Start ein langsamer Aufruf, und zwei Klicks konnten zwei Updates starten.""" - job_id, laeuft = jobengine.start_job_exklusiv("maintenance", args, label, on_done=on_done, - zeitlimit_s=zeitlimit_s) + job_id, laeuft = jobengine.start_job_exklusiv("maintenance", args, label, zeitlimit_s=zeitlimit_s, + nacharbeit="wartung:nach_update") if job_id is None: return {"ok": False, "status": "busy", "running": (laeuft or {}).get("label"), "detail": f"Es läuft schon: {(laeuft or {}).get('label', 'ein Update')}."} @@ -584,7 +585,7 @@ def check_updates_job() -> dict: # apt-get update gehört in den Wartungs-Riegel: parallel zu einem apt upgrade kollidiert es mit # der dpkg-Sperre. In der Gruppe taucht der Job auch in der Oberfläche auf. return _starten(["bash", "-c", "sudo apt-get update"], "Nach Updates suchen", - on_done=_nach_update, zeitlimit_s=15 * 60) + zeitlimit_s=15 * 60) def os_update_job() -> dict: @@ -596,7 +597,7 @@ def os_update_job() -> dict: cmd = ("sudo apt-get update && " "sudo bash -c 'DEBIAN_FRONTEND=noninteractive apt-get upgrade -y' " f"&& bash {STACK_POSTCHECK}") - return _starten(["bash", "-c", cmd], "OS-Update (apt)", on_done=_nach_update, zeitlimit_s=90 * 60) + return _starten(["bash", "-c", cmd], "OS-Update (apt)", zeitlimit_s=90 * 60) def engine_update_job() -> dict | None: @@ -606,7 +607,7 @@ def engine_update_job() -> dict | None: # Stack (stack-postcheck.sh) und rollt bei Fehler selbst zurück. Exit 0 nur bei verifiziertem # neuen Build. KEIN sudo-Passwort-Gate: das Skript ist per sudoers NOPASSWD freigegeben. return _starten(["bash", "-c", ENGINE_UPDATE_CMD], "Engine-Update (llama.cpp Vulkan)", - on_done=_nach_update, zeitlimit_s=60 * 60) + zeitlimit_s=60 * 60) def swap_update_job() -> dict | None: @@ -615,7 +616,7 @@ def swap_update_job() -> dict | None: # update-swap.sh sichert die alte Binary, aktualisiert, startet llama-swap neu, prüft den Stack # und rollt bei Fehler selbst zurück. Exit 0 nur bei verifizierter neuer Version. return _starten(["bash", "-c", SWAP_UPDATE_CMD], "Router-Update (llama-swap)", - on_done=_nach_update, zeitlimit_s=30 * 60) + zeitlimit_s=30 * 60) def _hermes_update_cmd() -> str: @@ -666,7 +667,7 @@ def hermes_update_job() -> dict: den Gateway neu starten. Davor ein Sicherheits-Backup (unser deploy/backup.sh). Kein sudo (alles im User-Space). Läuft als Hintergrund-Job (kann einige Minuten dauern).""" return _starten(["bash", "-c", _hermes_update_cmd()], "Hermes-Agent-Update", - on_done=_nach_update, zeitlimit_s=45 * 60) + zeitlimit_s=45 * 60) def update_all_job() -> dict: @@ -711,7 +712,7 @@ def update_all_job() -> dict: labels = " → ".join(label for label, _ in parts) ergebnis = _starten(["bash", "-c", cmd], f"Alle aktualisieren ({labels})", - on_done=_nach_update, zeitlimit_s=3 * 3600) + zeitlimit_s=3 * 3600) if ergebnis.get("ok"): ergebnis["parts"] = [label for label, _ in parts] return ergebnis diff --git a/backend/tests/test_jobengine.py b/backend/tests/test_jobengine.py new file mode 100644 index 0000000..eaca586 --- /dev/null +++ b/backend/tests/test_jobengine.py @@ -0,0 +1,196 @@ +"""Auftrags-System (Phase 2, 24.09.2026): Akten auf der Platte, Wiederaufnahme nach einem Neustart, +benannte Nacharbeiten, Zeitlimit, Abbruch. Die Tests laufen als Kindprozess („prozess“); den +systemd-Weg prüfen sie am gebauten Befehl.""" + +import json +import os +import sys +import time + +import pytest +from services import jobengine + +PY = sys.executable + + +@pytest.fixture +def motor(tmp_path, monkeypatch): + ordner = tmp_path / "jobs" + monkeypatch.setenv("MC_JOBS_DIR", str(ordner)) + monkeypatch.setenv("MC_JOBS_ART", "prozess") + monkeypatch.setattr(jobengine, "JOBS", {}) + monkeypatch.setattr(jobengine, "_PROCS", {}) + monkeypatch.setattr(jobengine, "_geladen", False) + monkeypatch.setattr(jobengine, "_NACHARBEITEN", dict(jobengine._NACHARBEITEN)) + return ordner + + +def _warte(job_id: str, sekunden: float = 30) -> dict: + ende = time.time() + sekunden + while time.time() < ende: + job = jobengine.JOBS[job_id] + if job["state"] in ("done", "failed", "canceled"): + return job + time.sleep(0.05) + raise AssertionError(f"Auftrag {job_id} wurde nicht fertig: {jobengine.JOBS[job_id]}") + + +def _oeffentlich(job_id: str) -> dict: + return next(j for j in jobengine.public_jobs() if j["id"] == job_id) + + +def test_erfolg_mit_protokoll_und_akte(motor): + code = "import sys; print('eins'); sys.stdout.write('10%\\r50%\\r100%\\n'); print('zwei')" + job_id = jobengine.start_job([PY, "-c", code], "Probe") + job = _warte(job_id) + assert (job["state"], job["returncode"]) == ("done", 0) + log = _oeffentlich(job_id)["log"] + assert log[0].startswith("$ ") and log[1:] == ["eins", "100%", "zwei"] + akte = json.loads((motor / f"{job_id}.json").read_text(encoding="utf-8")) + assert akte["state"] == "done" and akte["label"] == "Probe" + assert (motor / f"{job_id}.exit").read_text(encoding="utf-8").strip() == "0" + + +def test_fehlschlag_mit_code(motor): + job = _warte(jobengine.start_job([PY, "-c", "raise SystemExit(3)"], "Kaputt")) + assert (job["state"], job["returncode"]) == ("failed", 3) + + +def test_zeitlimit_beendet_den_auftrag(motor): + job = _warte(jobengine.start_job([PY, "-c", "import time; time.sleep(30)"], "Hängt", zeitlimit_s=1)) + assert job["state"] == "failed" and job["error"] == "Zeitlimit überschritten" + assert any("Zeitlimit" in z for z in _oeffentlich(job["id"])["log"]) + + +def test_abbrechen(motor): + job_id = jobengine.start_job([PY, "-c", "import time; time.sleep(30)"], "Lang") + time.sleep(0.3) + assert jobengine.cancel_job(job_id) is True + assert _warte(job_id)["state"] == "canceled" + assert jobengine.cancel_job(job_id) is False + + +def test_nacharbeit_nur_bei_erfolg_und_mit_daten(motor): + aufrufe: list[dict] = [] + jobengine.nacharbeit("probe:merken")(lambda **daten: aufrufe.append(daten)) + _warte(jobengine.start_job([PY, "-c", "pass"], "Gut", nacharbeit="probe:merken", + nacharbeit_daten={"model_id": "m", "role": "coder"})) + _warte(jobengine.start_job([PY, "-c", "raise SystemExit(1)"], "Schlecht", nacharbeit="probe:merken", + nacharbeit_daten={"x": 1})) + assert aufrufe == [{"model_id": "m", "role": "coder"}] + with pytest.raises(ValueError, match="Unbekannte Nacharbeit"): + jobengine.start_job([PY, "-c", "pass"], "X", nacharbeit="gibt:es-nicht") + + +def test_exklusiv_je_gruppe(motor): + erster, _ = jobengine.start_job_exklusiv("maintenance", [PY, "-c", "import time; time.sleep(30)"], "Update A") + zweiter, laeuft = jobengine.start_job_exklusiv("maintenance", [PY, "-c", "pass"], "Update B") + assert erster and zweiter is None and laeuft["label"] == "Update A" + assert jobengine.active_in_group("maintenance")["id"] == erster + jobengine.cancel_job(erster) + _warte(erster) + assert jobengine.active_in_group("maintenance") is None + + +def _akte(ordner, job_id: str, **felder) -> None: + ordner.mkdir(parents=True, exist_ok=True) + akte = {"id": job_id, "label": "Alt", "state": "running", "group": "download", "art": "systemd", + "returncode": None, "started_at": time.time() - 60, "finished_at": None, "zeitlimit_s": 3600, + "nacharbeit": None, "nacharbeit_daten": {}, **felder} + (ordner / f"{job_id}.json").write_text(json.dumps(akte), encoding="utf-8") + + +def test_wiederaufnahme_fertiger_einheit_mit_nacharbeit(motor): + """MC2 startet neu, während ein Download lief; er ist inzwischen fertig — die Rolle wird gesetzt.""" + aufrufe: list[dict] = [] + jobengine.nacharbeit("probe:rolle")(lambda **daten: aufrufe.append(daten)) + _akte(motor, "abc123", nacharbeit="probe:rolle", nacharbeit_daten={"role": "coder"}) + (motor / "abc123.exit").write_text("0", encoding="utf-8") + assert jobengine.wiederaufnehmen() == 1 + assert _warte("abc123")["state"] == "done" + assert aufrufe == [{"role": "coder"}] + # Ein zweiter Start von MC2 wiederholt die Nacharbeit nicht. + jobengine._geladen = False + jobengine.JOBS.clear() + assert jobengine.wiederaufnehmen() == 0 + assert aufrufe == [{"role": "coder"}] + + +def test_wiederaufnahme_verschwundene_einheit(motor, monkeypatch): + monkeypatch.setattr(jobengine, "_einheit_lebt", lambda job_id: False) + monkeypatch.setattr(jobengine, "_EINHEIT_PRUEFEN_S", 0.0) + _akte(motor, "weg456") + jobengine.wiederaufnehmen() + job = _warte("weg456") + assert job["state"] == "failed" and job["error"] == "Der Auftrag endete ohne Rückmeldung." + + +def test_wiederaufnahme_laufender_einheit_bleibt_laufend(motor, monkeypatch): + monkeypatch.setattr(jobengine, "_einheit_lebt", lambda job_id: True) + monkeypatch.setattr(jobengine, "_EINHEIT_PRUEFEN_S", 0.0) + _akte(motor, "lebt789", group="maintenance") + jobengine.wiederaufnehmen() + time.sleep(0.5) + assert jobengine.JOBS["lebt789"]["state"] == "running" + assert jobengine.active_in_group("maintenance")["id"] == "lebt789" + (motor / "lebt789.exit").write_text("0", encoding="utf-8") + assert _warte("lebt789")["state"] == "done" + + +def test_wiederaufnahme_kindprozess_ist_verloren(motor): + _akte(motor, "kind01", art="prozess") + jobengine.wiederaufnehmen() + job = jobengine.JOBS["kind01"] + assert job["state"] == "failed" and "neu gestartet" in job["error"] + + +def test_aufraeumen_loescht_alte_akten(motor, monkeypatch): + monkeypatch.setattr(jobengine, "BEHALTEN_MAX", 1) + ids = [_warte(jobengine.start_job([PY, "-c", "pass"], f"Nr {i}"))["id"] for i in range(3)] + sichtbar = [j["id"] for j in jobengine.public_jobs()] + assert sichtbar == [ids[-1]] + assert not (motor / f"{ids[0]}.json").exists() and not (motor / f"{ids[0]}.log").exists() + + +def test_systemd_befehl_ohne_geheimnis_in_der_befehlszeile(motor, monkeypatch): + monkeypatch.setenv("MC_JOBS_ART", "systemd") + monkeypatch.setenv("INVOCATION_ID", "von-mc2") + befehle: list[list[str]] = [] + + class Ergebnis: + returncode = 0 + stdout = stderr = "" + + monkeypatch.setattr(jobengine.subprocess, "run", lambda befehl, **kw: befehle.append(befehl) or Ergebnis()) + monkeypatch.setattr(jobengine, "_beobachten", lambda job_id: None) + job_id = jobengine.start_job(["hf", "download", "org/modell"], "Download", env={"HF_TOKEN": "geheim-123"}, + zeitlimit_s=600) + befehl = befehle[0] + assert befehl[0] == "systemd-run" and f"mc2-job-{job_id}" in befehl and "--collect" in befehl + assert "RuntimeMaxSec=600" in befehl + assert befehl[-3:] == ["hf", "download", "org/modell"] + assert not any("geheim-123" in teil for teil in befehl) + umgebung = (motor / f"{job_id}.env").read_text(encoding="utf-8") + assert "export HF_TOKEN=geheim-123" in umgebung and "INVOCATION_ID" not in umgebung + if os.name == "posix": + assert (motor / f"{job_id}.env").stat().st_mode & 0o077 == 0 + + +def test_systemd_start_scheitert_sauber(motor, monkeypatch): + monkeypatch.setenv("MC_JOBS_ART", "systemd") + + class Ergebnis: + returncode = 1 + stdout = "" + stderr = "Failed to connect to bus" + + monkeypatch.setattr(jobengine.subprocess, "run", lambda befehl, **kw: Ergebnis()) + job_id = jobengine.start_job(["true"], "Probe") + job = jobengine.JOBS[job_id] + assert job["state"] == "failed" and "Failed to connect to bus" in job["error"] + assert not (motor / f"{job_id}.env").exists() + + +def test_zeilen_aus_ueberschreibt_mit_wagenruecklauf(): + assert jobengine.zeilen_aus("a\n 1%\r 50%\r100%\nb\r\n") == ["a", "100%", "b"] + assert jobengine.zeilen_aus("läuft\r 7%") == [" 7%"] diff --git a/deploy/deploy.sh b/deploy/deploy.sh index 980f260..f7bd22c 100644 --- a/deploy/deploy.sh +++ b/deploy/deploy.sh @@ -37,19 +37,21 @@ SKILLS_ALT="autonomie konzept-fliessband llm-wiki morning-report orchestrator pr review trend-radar wartung konzept_fliessband llm_wiki morning_report projekt_start trend_radar betrieb-playbook pc-pfad-cache" +# Laufende Update-Aufträge (Gruppe maintenance). Seit Phase 2 überleben Aufträge den Neustart von +# MC2 (eigene systemd-Einheiten) — ein Download darf also weiterlaufen. Ein Update aber führt +# Skripte aus diesem Checkout aus; die darf der Deploy nicht mitten im Lauf austauschen. laufende_jobs() { curl -sf -m 5 "$API/api/jobs" 2>/dev/null \ - | python3 -c 'import json,sys; print(sum(1 for j in json.load(sys.stdin).get("jobs", []) if j.get("state") in ("running", "queued")))' \ + | python3 -c 'import json,sys; print(sum(1 for j in json.load(sys.stdin).get("jobs", []) if j.get("group") == "maintenance" and j.get("state") in ("running", "queued")))' \ 2>/dev/null || echo 0 } stufe1() { exec 9>"${XDG_RUNTIME_DIR}/mc2-deploy.lock" flock -n 9 || { echo "Es läuft schon ein Deploy."; exit 1; } - # Update- und Download-Jobs leben (noch) im MC2-Prozess: ein Neustart würde sie abwürgen. local n; n="$(laufende_jobs)" if [ "${n:-0}" -gt 0 ] && [ "${MC_DEPLOY_TROTZDEM:-0}" != "1" ]; then - echo "Es laufen $n Job(s) (Update oder Download). Deploy verschoben — später erneut, oder MC_DEPLOY_TROTZDEM=1." + echo "Es läuft ein Update-Auftrag. Deploy verschoben — später erneut, oder MC_DEPLOY_TROTZDEM=1." exit 1 fi local vorher; vorher="$(git rev-parse HEAD)"