spb/ka-guardian
Public
Python 57.1%
Shell 13.8%
CSS 12%
JavaScript 10%
HTML 7.2%
1# ============================================2# Projet : KA Guardian3# Fichier : runner/runner.py4# Rôle : Bras d'exécution sur chaque nœud d'app — lance les missions5# claude -p headless dans le repo de l'app, streame les événements6# vers l'orchestrateur (ka2/ka4/ka6), et sait rollback + restart.7# Author : Simon-Pierre Boucher8# Date : 2026-08-239# ============================================10"""KA Guardian Runner.1112Un seul runner par nœud (port 7791), une seule mission à la fois par nœud.13Auth: header X-KA-Token == $KA_GUARDIAN_TOKEN (fichier ~/.ka-guardian.env).14"""15from __future__ import annotations1617import json18import os19import pathlib20import shlex21import socket22import subprocess23import threading24import time25import uuid26from typing import Any2728import httpx29from fastapi import FastAPI, Header, HTTPException, Request30from pydantic import BaseModel3132HOME = pathlib.Path.home()33ENV_FILE = HOME / ".ka-guardian.env"34CLAUDE_ENV_FILE = HOME / ".claude" / ".env"35TRANSCRIPTS = HOME / "ka-guardian-runner" / "transcripts"36TRANSCRIPTS.mkdir(parents=True, exist_ok=True)37CLAUDE_BIN = "/opt/homebrew/bin/claude"383940def load_env_file(path: pathlib.Path) -> dict[str, str]:41 out: dict[str, str] = {}42 if path.exists():43 for line in path.read_text().splitlines():44 line = line.strip()45 if line and not line.startswith("#") and "=" in line:46 k, v = line.split("=", 1)47 out[k.strip()] = v.strip().strip('"').strip("'")48 return out495051LOCAL_ENV = load_env_file(ENV_FILE)52TOKEN = LOCAL_ENV.get("KA_GUARDIAN_TOKEN", "")53NODE = LOCAL_ENV.get("KA_GUARDIAN_NODE", socket.gethostname())5455app = FastAPI(title="KA Guardian Runner", docs_url=None, redoc_url=None)5657_lock = threading.Lock()58_current: dict[str, Any] = {} # mission en cours (métadonnées)596061def check_token(x_ka_token: str | None) -> None:62 if not TOKEN or x_ka_token != TOKEN:63 raise HTTPException(status_code=401, detail="token invalide")646566def expand(dir_: str) -> str:67 return os.path.expanduser(dir_)686970def git(dir_: str, *args: str) -> subprocess.CompletedProcess:71 return subprocess.run(72 ["git", "-C", expand(dir_), *args],73 capture_output=True, text=True, timeout=120,74 )757677def sh(cmd: str, timeout: int = 120) -> subprocess.CompletedProcess:78 env = {**os.environ, "PATH": "/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin"}79 return subprocess.run(["/bin/zsh", "-c", cmd], capture_output=True, text=True, timeout=timeout, env=env)808182def healthcheck(web_port: int, pm2_names: list[str]) -> dict[str, Any]:83 """Vérifie que l'app répond en HTTP et que ses process pm2 sont online."""84 http_ok = False85 try:86 r = httpx.get(f"http://127.0.0.1:{web_port}/", timeout=10, follow_redirects=True)87 http_ok = r.status_code < 50088 except Exception:89 http_ok = False90 pm2_status: dict[str, str] = {}91 try:92 out = sh("pm2 jlist").stdout93 procs = {p["name"]: p.get("pm2_env", {}).get("status", "?") for p in json.loads(out)}94 for name in pm2_names:95 pm2_status[name] = procs.get(name, "absent")96 except Exception as exc: # pm2 absent ou jlist illisible97 pm2_status = {n: f"inconnu ({exc})" for n in pm2_names}98 ok = http_ok and all(s in ("online", "launching") for s in pm2_status.values())99 return {"ok": ok, "http_ok": http_ok, "pm2": pm2_status}100101102class MissionIn(BaseModel):103 mission_id: str104 agent: str105 service: str106 source: str107 dir: str108 pm2: list[str]109 web_port: int110 model: str = "sonnet"111 max_turns: int = 70112 timeout_seconds: int = 3600113 max_cost_usd: float | None = None # plafond opérateur (estimation par usage)114 prompt: str115 callback_url: str # http://<orchestrateur>/api/ingest/<mission_id>116117118# Tarifs $/MTok (entrée, sortie) — pour ESTIMER le coût en cours de mission et119# appliquer le plafond opérateur. Cache: lecture ≈ 0,1× entrée, écriture ≈ 1,25×.120MODEL_RATES = {"fable": (10.0, 50.0), "opus": (5.0, 25.0), "sonnet": (3.0, 15.0), "haiku": (1.0, 5.0)}121122123def rates_for(model: str) -> tuple[float, float]:124 for key, r in MODEL_RATES.items():125 if key in model:126 return r127 return MODEL_RATES["fable"]128129130def usage_cost_usd(usage: dict[str, Any], model: str) -> float:131 rin, rout = rates_for(model)132 return (133 (usage.get("input_tokens") or 0) * rin134 + (usage.get("output_tokens") or 0) * rout135 + (usage.get("cache_read_input_tokens") or 0) * rin * 0.1136 + (usage.get("cache_creation_input_tokens") or 0) * rin * 1.25137 ) / 1e6138139140SPOOL = HOME / "ka-guardian-spool"141142143def post_event(url: str, payload: dict[str, Any]) -> None:144 """Dépose l'événement dans le spool du courrier zsh (deploy/courier.sh).145146 macOS 26 « Local Network Privacy » refuse le trafic LAN dès que python147 (homebrew) est dans la chaîne de processus — même ses enfants ssh/curl.148 Le courrier (launchd 100 % zsh) expédie via ssh → curl localhost.149 """150 try:151 from urllib.parse import urlsplit152 u = urlsplit(url)153 for d in ("outbox", "done", "tmp"):154 (SPOOL / d).mkdir(parents=True, exist_ok=True)155 # id préfixé du temps ns: le courrier traite les jobs en ordre lexical,156 # les événements arrivent donc dans l'ordre d'émission.157 jid = f"{time.time_ns()}-{uuid.uuid4().hex[:8]}"158 tmp = SPOOL / "tmp" / f"{jid}.job"159 tmp.write_text(f"{u.hostname} {u.port or 80} {u.path} 15\n"160 + json.dumps(payload, ensure_ascii=False))161 tmp.rename(SPOOL / "outbox" / f"{jid}.job")162 except Exception:163 pass # l'orchestrateur relira le transcript au besoin164165166def condense_stream_line(obj: dict[str, Any]) -> list[dict[str, Any]]:167 """Transforme une ligne stream-json de claude en événements courts pour le feed."""168 events: list[dict[str, Any]] = []169 t = obj.get("type")170 if t == "system" and obj.get("subtype") == "init":171 events.append({"type": "init", "data": {"model": obj.get("model", "?")}})172 elif t == "assistant":173 for block in obj.get("message", {}).get("content", []):174 if block.get("type") == "text" and block.get("text", "").strip():175 events.append({"type": "text", "data": {"text": block["text"][:1500]}})176 elif block.get("type") == "tool_use":177 inp = json.dumps(block.get("input", {}), ensure_ascii=False)178 events.append({"type": "tool", "data": {"name": block.get("name", "?"), "input": inp[:400]}})179 elif t == "result":180 events.append({"type": "result", "data": {181 "subtype": obj.get("subtype"),182 "result": (obj.get("result") or "")[:4000],183 "cost_usd": obj.get("total_cost_usd"),184 "num_turns": obj.get("num_turns"),185 "duration_ms": obj.get("duration_ms"),186 }})187 return events188189190def extract_verdict(result_text: str) -> dict[str, Any]:191 """Extrait le dernier bloc JSON {verdict: ...} de la réponse finale."""192 import re193 for m in reversed(re.findall(r"\{[^{}]*\"verdict\"[\s\S]*?\}", result_text)):194 try:195 v = json.loads(m)196 if "verdict" in v:197 return v198 except Exception:199 continue200 return {"verdict": "inconnu", "diagnostic": result_text[-500:] if result_text else ""}201202203def run_mission(m: MissionIn) -> None:204 global _current205 workdir = expand(m.dir)206 transcript = TRANSCRIPTS / f"{m.mission_id}.jsonl"207 base_commit = ""208 try:209 # Snapshot pré-mission : tree sale → commit de sûreté pour un rollback net.210 if git(m.dir, "status", "--porcelain").stdout.strip():211 git(m.dir, "add", "-A")212 git(m.dir, "commit", "-m", f"[{m.agent}] snapshot pré-mission {m.source}")213 base_commit = git(m.dir, "rev-parse", "HEAD").stdout.strip()214 post_event(m.callback_url, {"type": "start", "data": {"base_commit": base_commit, "node": NODE}})215216 env = {217 **os.environ,218 **load_env_file(CLAUDE_ENV_FILE),219 "PATH": "/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin",220 "HOME": str(HOME),221 }222 cmd = [223 CLAUDE_BIN, "-p", m.prompt,224 "--output-format", "stream-json", "--verbose",225 "--model", m.model,226 "--max-turns", str(m.max_turns),227 "--dangerously-skip-permissions",228 ]229 proc = subprocess.Popen(230 cmd, cwd=workdir, env=env,231 stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, bufsize=1,232 )233 _current["pid"] = proc.pid234 result_text, cost, turns = "", None, None235 spent_estimate = 0.0236 deadline = time.time() + m.timeout_seconds237238 def _watchdog() -> None:239 # Filet indépendant du flux stdout : le contrôle du deadline dans240 # la boucle de lecture ne s'exécute qu'à l'arrivée d'une ligne —241 # si claude se tait (outil qui pend) ou n'exit jamais après son242 # résultat (mission f9eb9d1b59c8 du 2026-08-24, zombie 5 h), la243 # mission restait « running » pour toujours et bloquait le nœud.244 while proc.poll() is None:245 if time.time() > deadline:246 proc.kill()247 post_event(m.callback_url, {"type": "error", "data": {"error": f"durée max atteinte ({m.timeout_seconds // 60} min) — session interrompue par le watchdog"}})248 return249 time.sleep(30)250251 threading.Thread(target=_watchdog, daemon=True).start()252 with transcript.open("w") as tf:253 for line in proc.stdout: # type: ignore[union-attr]254 tf.write(line)255 tf.flush() # le transcript doit survivre à une mort brutale du runner256 if time.time() > deadline:257 proc.kill()258 post_event(m.callback_url, {"type": "error", "data": {"error": f"durée max atteinte ({m.timeout_seconds // 60} min) — session interrompue"}})259 break260 line = line.strip()261 if not line:262 continue263 try:264 obj = json.loads(line)265 except Exception:266 continue267 if obj.get("type") == "assistant":268 spent_estimate += usage_cost_usd(obj.get("message", {}).get("usage") or {}, m.model)269 if m.max_cost_usd and spent_estimate > m.max_cost_usd:270 proc.kill()271 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"}})272 break273 for ev in condense_stream_line(obj):274 if ev["type"] == "result":275 result_text = ev["data"].get("result", "")276 cost = ev["data"].get("cost_usd")277 turns = ev["data"].get("num_turns")278 post_event(m.callback_url, ev)279 proc.wait(timeout=60)280281 commits_raw = git(m.dir, "rev-list", "--oneline", f"{base_commit}..HEAD").stdout.strip()282 commits = commits_raw.splitlines() if commits_raw else []283 verdict = extract_verdict(result_text)284 health = healthcheck(m.web_port, m.pm2)285 final = {"type": "final", "data": {286 "verdict": verdict, "commits": commits, "base_commit": base_commit,287 "cost_usd": cost, "num_turns": turns, "health": health,288 "exit_code": proc.returncode,289 }}290 # Résultat final persisté: l'orchestrateur peut le réclamer même si291 # l'événement s'est perdu (redémarrage, courrier en panne).292 (TRANSCRIPTS / f"{m.mission_id}.final.json").write_text(293 json.dumps(final, ensure_ascii=False))294 post_event(m.callback_url, final)295 except Exception as exc:296 post_event(m.callback_url, {"type": "error", "data": {"error": str(exc), "base_commit": base_commit}})297 finally:298 with _lock:299 _current.clear()300301302@app.get("/health")303def health() -> dict[str, Any]:304 return {"ok": True, "node": NODE, "busy": bool(_current), "current": _current.get("mission_id")}305306307@app.post("/missions")308def missions(m: MissionIn, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]:309 check_token(x_ka_token)310 workdir = expand(m.dir)311 if not os.path.isdir(workdir):312 raise HTTPException(status_code=400, detail=f"dir introuvable: {workdir}")313 with _lock:314 if _current:315 raise HTTPException(status_code=409, detail=f"mission déjà en cours: {_current.get('mission_id')}")316 _current.update({"mission_id": m.mission_id, "service": m.service, "source": m.source, "started": time.time()})317 threading.Thread(target=run_mission, args=(m,), daemon=True).start()318 return {"accepted": True, "node": NODE}319320321class MissionResultIn(BaseModel):322 mission_id: str323324325@app.post("/mission-result")326def mission_result(q: MissionResultIn, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]:327 """Réconciliation: l'orchestrateur réclame le sort d'une mission.328329 running = encore en cours ici; done = terminée (payload final joint);330 unknown = jamais terminée ici (probablement tuée) — à l'orchestrateur de trancher.331 """332 check_token(x_ka_token)333 if _current.get("mission_id") == q.mission_id:334 return {"status": "running"}335 fpath = TRANSCRIPTS / f"{q.mission_id}.final.json"336 if fpath.exists():337 return {"status": "done", "final": json.loads(fpath.read_text())}338 return {"status": "unknown"}339340341class RestartIn(BaseModel):342 pm2: list[str]343 web_port: int344345346@app.post("/restart")347def restart(r: RestartIn, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]:348 """Relance pm2 d'une app dont le site public ne répond plus (veilleur de349 sites des orchestrateurs). Léger et sans git — le retour arrière complet350 reste /rollback. Pas de verrou mission: relancer les process d'une app351 pendant qu'une mission travaille sur une AUTRE app du nœud est sans danger.352 """353 check_token(x_ka_token)354 restarted: dict[str, bool] = {}355 for name in r.pm2:356 p = sh(f"pm2 restart {shlex.quote(name)} --update-env", timeout=180)357 restarted[name] = p.returncode == 0358 time.sleep(8)359 h = healthcheck(r.web_port, r.pm2)360 return {"ok": all(restarted.values()) and h["ok"], "restarted": restarted, "health": h, "node": NODE}361362363class RollbackIn(BaseModel):364 dir: str365 base_commit: str366 pm2: list[str]367 web_port: int368369370@app.post("/rollback")371def rollback(r: RollbackIn, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]:372 """Revient au commit d'avant mission (reset --hard si tous les commits sont de l'agent, sinon revert)."""373 check_token(x_ka_token)374 log = git(r.dir, "log", "--format=%H %s", f"{r.base_commit}..HEAD").stdout.strip()375 entries = [ln.split(" ", 1) for ln in log.splitlines() if ln.strip()]376 agent_shas = [h for h, s in entries if s.startswith("[ka")]377 if not entries:378 mode = "aucun_commit"379 else:380 # Revert CIBLÉ, du plus récent au plus ancien, UNIQUEMENT les commits de381 # l'agent — jamais reset --hard ni revert de plage : les commits étrangers382 # (docs, humains, autres sessions) doivent survivre au rollback (leçon383 # fabri-ka 2026-08-24 : un revert de plage avait emporté la page /doc et384 # un connecteur), et les commits déjà poussés sur spbgit ne doivent385 # jamais exiger de force-push.386 reverted: list[str] = []387 skipped: list[str] = []388 for sha in agent_shas: # git log liste déjà du plus récent au plus ancien389 rc = git(r.dir, "revert", "--no-edit", sha)390 if rc.returncode != 0:391 git(r.dir, "revert", "--abort")392 skipped.append(sha[:9])393 else:394 reverted.append(sha[:9])395 mode = f"revert_cible({len(reverted)}/{len(agent_shas)})"396 if skipped:397 mode += f" sautés:{','.join(skipped)}"398 # Pas de git clean : les process de sync de l'app écrivent en continu,399 # on ne supprime jamais de fichiers non trackés apparus pendant la mission.400 # Push simple (les reverts s'empilent, fast-forward). Échec attendu sous401 # macOS 26 Local Network Privacy (runner = python) : le pousseur launchd402 # (com.ka.pousseur, chaîne binaires Apple) poussera dans les 5 minutes.403 pushed = None404 if entries:405 p = git(r.dir, "push", "origin", "HEAD")406 pushed = p.returncode == 0407 for name in r.pm2:408 sh(f"pm2 restart {shlex.quote(name)} --update-env", timeout=180)409 time.sleep(8)410 h = healthcheck(r.web_port, r.pm2)411 return {"ok": True, "mode": mode, "pushed": pushed, "health": h, "head": git(r.dir, "rev-parse", "HEAD").stdout.strip()}412413414if __name__ == "__main__":415 import uvicorn416 uvicorn.run(app, host="0.0.0.0", port=int(os.environ.get("KA_GUARDIAN_RUNNER_PORT", "7791")))417