# ============================================ # Projet : KA Guardian # Fichier : runner/runner.py # Rôle : Bras d'exécution sur chaque nœud d'app — lance les missions # claude -p headless dans le repo de l'app, streame les événements # vers l'orchestrateur (ka2/ka4/ka6), et sait rollback + restart. # Author : Simon-Pierre Boucher # Date : 2026-08-23 # ============================================ """KA Guardian Runner. Un seul runner par nœud (port 7791), une seule mission à la fois par nœud. Auth: header X-KA-Token == $KA_GUARDIAN_TOKEN (fichier ~/.ka-guardian.env). """ from __future__ import annotations import json import os import pathlib import shlex import socket import subprocess import threading import time import uuid from typing import Any import httpx from fastapi import FastAPI, Header, HTTPException, Request from pydantic import BaseModel HOME = pathlib.Path.home() ENV_FILE = HOME / ".ka-guardian.env" CLAUDE_ENV_FILE = HOME / ".claude" / ".env" TRANSCRIPTS = HOME / "ka-guardian-runner" / "transcripts" TRANSCRIPTS.mkdir(parents=True, exist_ok=True) CLAUDE_BIN = "/opt/homebrew/bin/claude" def load_env_file(path: pathlib.Path) -> dict[str, str]: out: dict[str, str] = {} if path.exists(): for line in path.read_text().splitlines(): line = line.strip() if line and not line.startswith("#") and "=" in line: k, v = line.split("=", 1) out[k.strip()] = v.strip().strip('"').strip("'") return out LOCAL_ENV = load_env_file(ENV_FILE) TOKEN = LOCAL_ENV.get("KA_GUARDIAN_TOKEN", "") NODE = LOCAL_ENV.get("KA_GUARDIAN_NODE", socket.gethostname()) app = FastAPI(title="KA Guardian Runner", docs_url=None, redoc_url=None) _lock = threading.Lock() _current: dict[str, Any] = {} # mission en cours (métadonnées) def check_token(x_ka_token: str | None) -> None: if not TOKEN or x_ka_token != TOKEN: raise HTTPException(status_code=401, detail="token invalide") def expand(dir_: str) -> str: return os.path.expanduser(dir_) def git(dir_: str, *args: str) -> subprocess.CompletedProcess: return subprocess.run( ["git", "-C", expand(dir_), *args], capture_output=True, text=True, timeout=120, ) def sh(cmd: str, timeout: int = 120) -> subprocess.CompletedProcess: env = {**os.environ, "PATH": "/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin"} return subprocess.run(["/bin/zsh", "-c", cmd], capture_output=True, text=True, timeout=timeout, env=env) def healthcheck(web_port: int, pm2_names: list[str]) -> dict[str, Any]: """Vérifie que l'app répond en HTTP et que ses process pm2 sont online.""" http_ok = False try: r = httpx.get(f"http://127.0.0.1:{web_port}/", timeout=10, follow_redirects=True) http_ok = r.status_code < 500 except Exception: http_ok = False pm2_status: dict[str, str] = {} try: out = sh("pm2 jlist").stdout procs = {p["name"]: p.get("pm2_env", {}).get("status", "?") for p in json.loads(out)} for name in pm2_names: pm2_status[name] = procs.get(name, "absent") except Exception as exc: # pm2 absent ou jlist illisible pm2_status = {n: f"inconnu ({exc})" for n in pm2_names} ok = http_ok and all(s in ("online", "launching") for s in pm2_status.values()) return {"ok": ok, "http_ok": http_ok, "pm2": pm2_status} class MissionIn(BaseModel): mission_id: str agent: str service: str source: str dir: str pm2: list[str] web_port: int model: str = "sonnet" max_turns: int = 70 timeout_seconds: int = 3600 max_cost_usd: float | None = None # plafond opérateur (estimation par usage) prompt: str callback_url: str # http:///api/ingest/ # Tarifs $/MTok (entrée, sortie) — pour ESTIMER le coût en cours de mission et # appliquer le plafond opérateur. Cache: lecture ≈ 0,1× entrée, écriture ≈ 1,25×. MODEL_RATES = {"fable": (10.0, 50.0), "opus": (5.0, 25.0), "sonnet": (3.0, 15.0), "haiku": (1.0, 5.0)} def rates_for(model: str) -> tuple[float, float]: for key, r in MODEL_RATES.items(): if key in model: return r return MODEL_RATES["fable"] def usage_cost_usd(usage: dict[str, Any], model: str) -> float: rin, rout = rates_for(model) return ( (usage.get("input_tokens") or 0) * rin + (usage.get("output_tokens") or 0) * rout + (usage.get("cache_read_input_tokens") or 0) * rin * 0.1 + (usage.get("cache_creation_input_tokens") or 0) * rin * 1.25 ) / 1e6 SPOOL = HOME / "ka-guardian-spool" def post_event(url: str, payload: dict[str, Any]) -> None: """Dépose l'événement dans le spool du courrier zsh (deploy/courier.sh). macOS 26 « Local Network Privacy » refuse le trafic LAN dès que python (homebrew) est dans la chaîne de processus — même ses enfants ssh/curl. Le courrier (launchd 100 % zsh) expédie via ssh → curl localhost. """ try: from urllib.parse import urlsplit u = urlsplit(url) for d in ("outbox", "done", "tmp"): (SPOOL / d).mkdir(parents=True, exist_ok=True) # id préfixé du temps ns: le courrier traite les jobs en ordre lexical, # les événements arrivent donc dans l'ordre d'émission. jid = f"{time.time_ns()}-{uuid.uuid4().hex[:8]}" tmp = SPOOL / "tmp" / f"{jid}.job" tmp.write_text(f"{u.hostname} {u.port or 80} {u.path} 15\n" + json.dumps(payload, ensure_ascii=False)) tmp.rename(SPOOL / "outbox" / f"{jid}.job") except Exception: pass # l'orchestrateur relira le transcript au besoin def condense_stream_line(obj: dict[str, Any]) -> list[dict[str, Any]]: """Transforme une ligne stream-json de claude en événements courts pour le feed.""" events: list[dict[str, Any]] = [] t = obj.get("type") if t == "system" and obj.get("subtype") == "init": events.append({"type": "init", "data": {"model": obj.get("model", "?")}}) elif t == "assistant": for block in obj.get("message", {}).get("content", []): if block.get("type") == "text" and block.get("text", "").strip(): events.append({"type": "text", "data": {"text": block["text"][:1500]}}) elif block.get("type") == "tool_use": inp = json.dumps(block.get("input", {}), ensure_ascii=False) events.append({"type": "tool", "data": {"name": block.get("name", "?"), "input": inp[:400]}}) elif t == "result": events.append({"type": "result", "data": { "subtype": obj.get("subtype"), "result": (obj.get("result") or "")[:4000], "cost_usd": obj.get("total_cost_usd"), "num_turns": obj.get("num_turns"), "duration_ms": obj.get("duration_ms"), }}) return events def extract_verdict(result_text: str) -> dict[str, Any]: """Extrait le dernier bloc JSON {verdict: ...} de la réponse finale.""" import re for m in reversed(re.findall(r"\{[^{}]*\"verdict\"[\s\S]*?\}", result_text)): try: v = json.loads(m) if "verdict" in v: return v except Exception: continue return {"verdict": "inconnu", "diagnostic": result_text[-500:] if result_text else ""} def run_mission(m: MissionIn) -> None: global _current workdir = expand(m.dir) transcript = TRANSCRIPTS / f"{m.mission_id}.jsonl" base_commit = "" try: # Snapshot pré-mission : tree sale → commit de sûreté pour un rollback net. if git(m.dir, "status", "--porcelain").stdout.strip(): git(m.dir, "add", "-A") git(m.dir, "commit", "-m", f"[{m.agent}] snapshot pré-mission {m.source}") base_commit = git(m.dir, "rev-parse", "HEAD").stdout.strip() post_event(m.callback_url, {"type": "start", "data": {"base_commit": base_commit, "node": NODE}}) env = { **os.environ, **load_env_file(CLAUDE_ENV_FILE), "PATH": "/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin", "HOME": str(HOME), } cmd = [ CLAUDE_BIN, "-p", m.prompt, "--output-format", "stream-json", "--verbose", "--model", m.model, "--max-turns", str(m.max_turns), "--dangerously-skip-permissions", ] proc = subprocess.Popen( cmd, cwd=workdir, env=env, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, bufsize=1, ) _current["pid"] = proc.pid result_text, cost, turns = "", None, None spent_estimate = 0.0 deadline = time.time() + m.timeout_seconds def _watchdog() -> None: # Filet indépendant du flux stdout : le contrôle du deadline dans # la boucle de lecture ne s'exécute qu'à l'arrivée d'une ligne — # si claude se tait (outil qui pend) ou n'exit jamais après son # résultat (mission f9eb9d1b59c8 du 2026-08-24, zombie 5 h), la # mission restait « running » pour toujours et bloquait le nœud. while proc.poll() is None: if time.time() > deadline: proc.kill() post_event(m.callback_url, {"type": "error", "data": {"error": f"durée max atteinte ({m.timeout_seconds // 60} min) — session interrompue par le watchdog"}}) return time.sleep(30) threading.Thread(target=_watchdog, daemon=True).start() with transcript.open("w") as tf: for line in proc.stdout: # type: ignore[union-attr] tf.write(line) tf.flush() # le transcript doit survivre à une mort brutale du runner if time.time() > deadline: proc.kill() post_event(m.callback_url, {"type": "error", "data": {"error": f"durée max atteinte ({m.timeout_seconds // 60} min) — session interrompue"}}) break line = line.strip() if not line: continue try: obj = json.loads(line) except Exception: continue if obj.get("type") == "assistant": spent_estimate += usage_cost_usd(obj.get("message", {}).get("usage") or {}, m.model) if m.max_cost_usd and spent_estimate > m.max_cost_usd: proc.kill() post_event(m.callback_url, {"type": "error", "data": {"error": f"coût max atteint (~{spent_estimate:.2f} $ estimé, plafond {m.max_cost_usd:.2f} $) — session interrompue"}}) break for ev in condense_stream_line(obj): if ev["type"] == "result": result_text = ev["data"].get("result", "") cost = ev["data"].get("cost_usd") turns = ev["data"].get("num_turns") post_event(m.callback_url, ev) proc.wait(timeout=60) commits_raw = git(m.dir, "rev-list", "--oneline", f"{base_commit}..HEAD").stdout.strip() commits = commits_raw.splitlines() if commits_raw else [] verdict = extract_verdict(result_text) health = healthcheck(m.web_port, m.pm2) final = {"type": "final", "data": { "verdict": verdict, "commits": commits, "base_commit": base_commit, "cost_usd": cost, "num_turns": turns, "health": health, "exit_code": proc.returncode, }} # Résultat final persisté: l'orchestrateur peut le réclamer même si # l'événement s'est perdu (redémarrage, courrier en panne). (TRANSCRIPTS / f"{m.mission_id}.final.json").write_text( json.dumps(final, ensure_ascii=False)) post_event(m.callback_url, final) except Exception as exc: post_event(m.callback_url, {"type": "error", "data": {"error": str(exc), "base_commit": base_commit}}) finally: with _lock: _current.clear() @app.get("/health") def health() -> dict[str, Any]: return {"ok": True, "node": NODE, "busy": bool(_current), "current": _current.get("mission_id")} @app.post("/missions") def missions(m: MissionIn, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]: check_token(x_ka_token) workdir = expand(m.dir) if not os.path.isdir(workdir): raise HTTPException(status_code=400, detail=f"dir introuvable: {workdir}") with _lock: if _current: raise HTTPException(status_code=409, detail=f"mission déjà en cours: {_current.get('mission_id')}") _current.update({"mission_id": m.mission_id, "service": m.service, "source": m.source, "started": time.time()}) threading.Thread(target=run_mission, args=(m,), daemon=True).start() return {"accepted": True, "node": NODE} class MissionResultIn(BaseModel): mission_id: str @app.post("/mission-result") def mission_result(q: MissionResultIn, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]: """Réconciliation: l'orchestrateur réclame le sort d'une mission. running = encore en cours ici; done = terminée (payload final joint); unknown = jamais terminée ici (probablement tuée) — à l'orchestrateur de trancher. """ check_token(x_ka_token) if _current.get("mission_id") == q.mission_id: return {"status": "running"} fpath = TRANSCRIPTS / f"{q.mission_id}.final.json" if fpath.exists(): return {"status": "done", "final": json.loads(fpath.read_text())} return {"status": "unknown"} class RestartIn(BaseModel): pm2: list[str] web_port: int @app.post("/restart") def restart(r: RestartIn, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]: """Relance pm2 d'une app dont le site public ne répond plus (veilleur de sites des orchestrateurs). Léger et sans git — le retour arrière complet reste /rollback. Pas de verrou mission: relancer les process d'une app pendant qu'une mission travaille sur une AUTRE app du nœud est sans danger. """ check_token(x_ka_token) restarted: dict[str, bool] = {} for name in r.pm2: p = sh(f"pm2 restart {shlex.quote(name)} --update-env", timeout=180) restarted[name] = p.returncode == 0 time.sleep(8) h = healthcheck(r.web_port, r.pm2) return {"ok": all(restarted.values()) and h["ok"], "restarted": restarted, "health": h, "node": NODE} class RollbackIn(BaseModel): dir: str base_commit: str pm2: list[str] web_port: int @app.post("/rollback") def rollback(r: RollbackIn, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]: """Revient au commit d'avant mission (reset --hard si tous les commits sont de l'agent, sinon revert).""" check_token(x_ka_token) log = git(r.dir, "log", "--format=%H %s", f"{r.base_commit}..HEAD").stdout.strip() entries = [ln.split(" ", 1) for ln in log.splitlines() if ln.strip()] agent_shas = [h for h, s in entries if s.startswith("[ka")] if not entries: mode = "aucun_commit" else: # Revert CIBLÉ, du plus récent au plus ancien, UNIQUEMENT les commits de # l'agent — jamais reset --hard ni revert de plage : les commits étrangers # (docs, humains, autres sessions) doivent survivre au rollback (leçon # fabri-ka 2026-08-24 : un revert de plage avait emporté la page /doc et # un connecteur), et les commits déjà poussés sur spbgit ne doivent # jamais exiger de force-push. reverted: list[str] = [] skipped: list[str] = [] for sha in agent_shas: # git log liste déjà du plus récent au plus ancien rc = git(r.dir, "revert", "--no-edit", sha) if rc.returncode != 0: git(r.dir, "revert", "--abort") skipped.append(sha[:9]) else: reverted.append(sha[:9]) mode = f"revert_cible({len(reverted)}/{len(agent_shas)})" if skipped: mode += f" sautés:{','.join(skipped)}" # Pas de git clean : les process de sync de l'app écrivent en continu, # on ne supprime jamais de fichiers non trackés apparus pendant la mission. # Push simple (les reverts s'empilent, fast-forward). Échec attendu sous # macOS 26 Local Network Privacy (runner = python) : le pousseur launchd # (com.ka.pousseur, chaîne binaires Apple) poussera dans les 5 minutes. pushed = None if entries: p = git(r.dir, "push", "origin", "HEAD") pushed = p.returncode == 0 for name in r.pm2: sh(f"pm2 restart {shlex.quote(name)} --update-env", timeout=180) time.sleep(8) h = healthcheck(r.web_port, r.pm2) return {"ok": True, "mode": mode, "pushed": pushed, "health": h, "head": git(r.dir, "rev-parse", "HEAD").stdout.strip()} if __name__ == "__main__": import uvicorn uvicorn.run(app, host="0.0.0.0", port=int(os.environ.get("KA_GUARDIAN_RUNNER_PORT", "7791")))