spb/ka-guardian
Public
Python 57.1%
Shell 13.8%
CSS 12%
JavaScript 10%
HTML 7.2%
1# ============================================2# Projet : KA Guardian3# Fichier : orchestrator/main.py4# Rôle : Cerveau d'un agent gardien (ka2/ka4/ka6) — surveille la santé5# des connecteurs via api-ka, ouvre des incidents, dépêche des6# missions claude aux runners des nœuds, surveille la guérison,7# rollback automatique si ça empire. Sert aussi le dashboard live.8# Author : Simon-Pierre Boucher9# Date : 2026-08-2310# ============================================11"""Instance unique par agent: AGENT=ka2|ka4|ka6 (env), port depuis topology.json."""12from __future__ import annotations1314import asyncio15import json16import os17import pathlib18import sqlite319import time20import uuid21from typing import Any2223import httpx24from fastapi import FastAPI, Header, HTTPException, Request25from fastapi.responses import FileResponse, HTMLResponse, StreamingResponse26from fastapi.staticfiles import StaticFiles2728import registry # topologie vivante (registre mld) — même dossier que main.py2930BASE = pathlib.Path(__file__).resolve().parent31ROOT = BASE.parent32HOME = pathlib.Path.home()3334AGENT = os.environ.get("AGENT", "ka2")35TOPO = json.loads((ROOT / "topology.json").read_text())36ME = TOPO["agents"][AGENT]37POLICY = TOPO["policy"]38# SERVICES / NODES : dicts mutés EN PLACE par registry.refresh() — l'emplacement39# réel (nœud, IP, dir, port, PM2) vient du registre mld, topology.json n'est40# que le repli si le registre n'a jamais été reçu.41SERVICES: dict[str, Any] = TOPO["services"]42NODES: dict[str, Any] = TOPO["nodes"]43MY_SERVICES: set[str] = set(ME["services"])44RUNNER_PORT = TOPO["runner_port"]45try: # première application au chargement (silencieuse : le hub SSE n'existe pas encore)46 for _m in registry.refresh(TOPO, SERVICES, NODES):47 print(f"[registre] {_m}", flush=True)48except Exception as _exc:49 print(f"[registre] indisponible au démarrage ({_exc}) — topology.json en repli", flush=True)505152def node_alias(node_key: str) -> str:53 return (NODES.get(node_key) or {}).get("alias", node_key)545556def node_ip(node_key: str) -> str | None:57 return (NODES.get(node_key) or {}).get("lan_ip")585960def apply_registry(announce: bool = True) -> None:61 """Relit le registre mld s'il a changé et signale les déménagements."""62 try:63 changes = registry.refresh(TOPO, SERVICES, NODES)64 except Exception as exc: # jamais fatal : on garde la dernière topologie connue65 print(f"[registre] échec: {exc}", flush=True)66 return67 for msg in changes:68 print(f"[registre] {msg}", flush=True)69 if announce:70 hub.publish_sync({"kind": "log", "level": "warn", "msg": f"registre mld — {msg}", "ts": now()})717273_RUNNER_WARNED: dict[str, float] = {}747576def runner_available(node_key: str, what: str) -> bool:77 """False si la sonde registry-sync dit qu'il n'y a PAS de runner sur ce nœud78 (app déménagée sur un nœud jamais équipé) — on n'envoie alors rien au79 courrier, et on avertit au plus une fois par heure par nœud."""80 ip = node_ip(node_key)81 rs = registry.runner_status(ip)82 if rs is None or rs.get("ok"):83 return True84 if now() - _RUNNER_WARNED.get(node_key, 0) > 3600:85 _RUNNER_WARNED[node_key] = now()86 msg = (f"aucun runner gardien sur {node_alias(node_key)} ({ip}) — {what} impossible ; "87 f"déployer avec deploy/deploy.sh runners (le nœud est listé par le registre)")88 print(f"[runner] {msg}", flush=True)89 hub.publish_sync({"kind": "log", "level": "error", "msg": msg, "ts": now()})90 return False9192# ka6 (ou autre) ramasse les services non assignés qui apparaîtraient dans api-ka.93DEFAULT_AGENT = ME.get("default_for_unknown_services", False)94ASSIGNED = {s for a in TOPO["agents"].values() for s in a["services"]}959697def load_env_file(path: pathlib.Path) -> dict[str, str]:98 out: dict[str, str] = {}99 if path.exists():100 for line in path.read_text().splitlines():101 line = line.strip()102 if line and not line.startswith("#") and "=" in line:103 k, v = line.split("=", 1)104 out[k.strip()] = v.strip().strip('"').strip("'")105 return out106107108TOKEN = load_env_file(HOME / ".ka-guardian.env").get("KA_GUARDIAN_TOKEN", "")109110111SPOOL = HOME / "ka-guardian-spool"112113114async def lan_post(ip: str, port: int, path: str, payload: dict[str, Any], timeout: int = 30) -> tuple[int, str]:115 """POST JSON vers un autre nœud via le courrier zsh (deploy/courier.sh).116117 macOS 26 « Local Network Privacy » refuse le trafic LAN dès que python118 (homebrew) est dans la chaîne de processus — même ses enfants ssh/curl.119 On dépose donc un job dans un spool disque ; un service launchd 100 % zsh120 l'expédie (ssh → curl localhost du nœud cible) et dépose la réponse.121 """122 for d in ("outbox", "done", "tmp"):123 (SPOOL / d).mkdir(parents=True, exist_ok=True)124 jid = uuid.uuid4().hex125 tmp = SPOOL / "tmp" / f"{jid}.job"126 tmp.write_text(f"{ip} {port} {path} {timeout}\n" + jdump(payload))127 tmp.rename(SPOOL / "outbox" / f"{jid}.job")128 resp = SPOOL / "done" / f"{jid}.resp"129 deadline = time.time() + timeout + 25130 while time.time() < deadline:131 if resp.exists():132 txt = resp.read_text()133 resp.unlink(missing_ok=True)134 body, _, code = txt.strip().rpartition("\n")135 if code.startswith("courier_ssh_rc"):136 raise ConnectionError(f"{code} vers {ip}:{port}{path}")137 return int(code or "0"), body138 await asyncio.sleep(0.3)139 raise TimeoutError(f"courrier sans réponse pour {ip}:{port}{path}")140DATA = ROOT / "data"141DATA.mkdir(exist_ok=True)142DB_PATH = DATA / f"{AGENT}.db"143144# ---------------------------------------------------------------- SQLite ---145146def db() -> sqlite3.Connection:147 conn = sqlite3.connect(DB_PATH)148 conn.row_factory = sqlite3.Row149 conn.execute("PRAGMA journal_mode=WAL")150 return conn151152153def init_db() -> None:154 with db() as c:155 c.executescript("""156 CREATE TABLE IF NOT EXISTS incidents(157 id TEXT PRIMARY KEY, service TEXT, source TEXT, status_detected TEXT,158 state TEXT, attempts INTEGER DEFAULT 0, created REAL, updated REAL,159 resolved REAL, detail TEXT);160 CREATE TABLE IF NOT EXISTS missions(161 id TEXT PRIMARY KEY, incident_id TEXT, service TEXT, source TEXT,162 node TEXT, state TEXT, base_commit TEXT, commits TEXT, verdict TEXT,163 cost_usd REAL, num_turns INTEGER, health TEXT, started REAL, ended REAL);164 CREATE TABLE IF NOT EXISTS events(165 id INTEGER PRIMARY KEY AUTOINCREMENT, mission_id TEXT, ts REAL,166 type TEXT, data TEXT);167 CREATE TABLE IF NOT EXISTS snapshots(168 ts REAL PRIMARY KEY, mine TEXT, ecosystem TEXT);169 CREATE INDEX IF NOT EXISTS ev_mission ON events(mission_id);170 """)171172173def now() -> float:174 return time.time()175176177def jdump(x: Any) -> str:178 return json.dumps(x, ensure_ascii=False)179180181# ------------------------------------------------------------------- SSE ---182183class Hub:184 def __init__(self) -> None:185 self.clients: set[asyncio.Queue] = set()186187 async def publish(self, msg: dict[str, Any]) -> None:188 for q in list(self.clients):189 if q.qsize() < 500:190 q.put_nowait(msg)191192 def publish_sync(self, msg: dict[str, Any]) -> None:193 if LOOP:194 asyncio.run_coroutine_threadsafe(self.publish(msg), LOOP)195196197hub = Hub()198LOOP: asyncio.AbstractEventLoop | None = None199LATEST: dict[str, Any] = {"mine": {}, "ecosystem": {}, "polled": 0, "paused": False}200201# ------------------------------------------------------------- missions ----202203def active_incident(c: sqlite3.Connection, service: str, source: str) -> sqlite3.Row | None:204 return c.execute(205 "SELECT * FROM incidents WHERE service=? AND source=? AND state IN "206 "('open','dispatched','fixing','watching','cooldown') ORDER BY created DESC LIMIT 1",207 (service, source)).fetchone()208209210def set_incident(c: sqlite3.Connection, iid: str, **kw: Any) -> None:211 kw["updated"] = now()212 keys = ",".join(f"{k}=?" for k in kw)213 c.execute(f"UPDATE incidents SET {keys} WHERE id=?", (*kw.values(), iid))214215216def incident_event(iid: str, service: str, source: str, state: str, note: str = "") -> None:217 hub.publish_sync({"kind": "incident", "incident_id": iid, "service": service,218 "source": source, "state": state, "note": note, "ts": now()})219220221def service_situation(service: str) -> str:222 """Portrait santé de l'app entière — permet à l'agent de distinguer une223 panne isolée d'un problème systémique (app entière, réseau, quota)."""224 block = LATEST["mine"].get(service) or {}225 summary = block.get("summary") or {}226 app_row = block.get("app") or {}227 sick = [f"{c['source']} ({c['status']})" for c in block.get("connectors", [])228 if c.get("status") not in ("ok", None)][:15]229 lines = [f"- app entière: {app_row.get('status', '?')} — {summary.get('ok', '?')} ok, "230 f"{summary.get('degraded', 0)} dégradés, {summary.get('broken', 0)} cassés, {summary.get('stale', 0)} endormis"]231 if sick:232 lines.append(f"- autres connecteurs non-ok de cette app: {', '.join(sick)}")233 lines.append("- Si BEAUCOUP de connecteurs sont touchés en même temps, le problème est probablement "234 "systémique (scheduler de l'app, réseau, quota API): diagnostique au niveau app, pas source par source.")235 else:236 lines.append("- Tous les autres connecteurs de l'app sont sains: la panne est isolée à cette source.")237 return "\n".join(lines)238239240def mission_history(service: str, source: str) -> str:241 """Résumé des missions précédentes sur ce connecteur (éviter de refaire ce qui a échoué)."""242 with db() as c:243 rows = c.execute(244 "SELECT m.started, m.verdict, m.commits FROM missions m WHERE m.service=? AND m.source=? "245 "AND m.state!='running' ORDER BY m.started DESC LIMIT 3", (service, source)).fetchall()246 if not rows:247 return "- aucune: c'est la première intervention sur ce connecteur."248 out = []249 for r in rows:250 v = json.loads(r["verdict"] or "{}")251 out.append(f"- {time.strftime('%Y-%m-%d %H:%M', time.localtime(r['started']))}: "252 f"verdict={v.get('verdict', '?')} — {str(v.get('diagnostic', ''))[:200]} "253 f"(actions: {str(v.get('actions', ''))[:150]})")254 out.append("NE RÉPÈTE PAS une approche qui a déjà échoué ci-dessus: change d'angle.")255 return "\n".join(out)256257258def build_prompt(service: str, source: str, health: dict[str, Any]) -> str:259 svc = SERVICES[service]260 sync_proc = next((p for p in svc["pm2"] if "sync" in p or "etl" in p), (svc["pm2"][-1] if svc["pm2"] else "?"))261 node_alias_ = node_alias(svc["node"])262 return f"""Tu es {AGENT}, agent gardien autonome des connecteurs du Groupe KA. Mission: réparer le connecteur « {source} » de l'app {svc['app']} (service api-ka: {service}, site {svc.get('site', '')}).263264== CONTEXTE SANTÉ DU CONNECTEUR (supervision api-ka, scan aux 2 h) ==265- statut détecté: {health.get('status')} | échecs consécutifs: {health.get('consecutive_failures')}266- dernier succès: {health.get('last_success')} | volume au dernier sync: {health.get('found_last')} (médiane historique: {health.get('median_found')})267- message: {health.get('message')}268Rappel des statuts: broken = ≥3 syncs consécutifs en échec ou à 0 résultat; stale = aucun succès depuis > 2× la cadence attendue; degraded = volume < 50 % de la médiane.269270== SITUATION DE L'APP ==271{service_situation(service)}272273== INTERVENTIONS PRÉCÉDENTES SUR CE CONNECTEUR ==274{mission_history(service, source)}275276== TON ENVIRONNEMENT ==277Tu es sur le nœud {node_alias_}, ton répertoire courant est le repo de l'app: {svc['dir']} (source de vérité, développement remote-first).278L'app tourne via pm2 ({', '.join(svc['pm2'])}); le process de synchronisation des connecteurs est {sync_proc}; site web local sur le port {svc['web_port']}.279Le venv Python et/ou node_modules du repo sont déjà installés — utilise-les, jamais d'installation globale.280281== DÉMARCHE IMPOSÉE (dans l'ordre) ==2821. IMPRÈGNE-TOI DU REPO: lis le CLAUDE.md et/ou README du repo s'ils existent, et surtout la doc du connecteur si elle existe: cherche `docs/connecteurs/` (fiches générées par app, souvent une par source: URL de la source, stratégie de fetch, format, pièges connus). `grep -ri "{source}"` pour localiser le code du connecteur, sa config et son éventuelle entrée de registre.2832. LIS LES LOGS: dossier logs/ du repo + `pm2 logs {sync_proc} --nostream --lines 300` — cherche les traces d'erreur de « {source} » (traceback, code HTTP, timeout).2843. REPRODUIS: exécute le connecteur ou son fetch directement (la plupart des apps ont un runner par source — la doc/le code du scheduler te montrera comment; sinon appelle la fonction de fetch dans un petit script via le venv du repo). Constate l'erreur réelle.2854. DIAGNOSTIQUE la cause racine: HTML/sélecteurs changés? endpoint JSON déplacé? 403/429 anti-bot? pagination cassée? redirection? certificat/TLS (gotcha connu: certains sites cassent avec requests → utiliser curl via subprocess)? sitemap figé?2865. CORRIGE de façon MINIMALE et ciblée, en respectant les conventions et l'architecture du repo (mêmes patterns que les connecteurs voisins). Si le site bloque, la pile d'escalade du Groupe KA est disponible, clés dans ~/.claude/.env: requêtes directes → curl → Scrapfly (SCRAPFLY_API_KEY, render_js/asp) → Bright Data Web Unlocker → proxies résidentiels Oxylabs (pr.oxylabs.io:7777, -cc-CA) → Serper/Tavily pour retrouver une source déplacée. Regarde comment le repo utilise déjà ces services et fais pareil.2876. RE-TESTE réellement: le connecteur doit rapporter un volume plausible (ordre de grandeur de la médiane {health.get('median_found')}). INTERDIT d'inventer, stubber ou câbler des données en dur. Un test qui retourne 0 ou 3 items quand la médiane est {health.get('median_found')} N'EST PAS une réparation.2887. REDÉMARRE uniquement le process concerné: `pm2 restart {sync_proc}`. Puis vérifie que le site répond: `curl -s -o /dev/null -w '%{{http_code}}' http://localhost:{svc['web_port']}/` (attendu: 200/3xx).2898. COMMITTE: `git add -A && git commit -m "[{AGENT}] fix connecteur {source}: <résumé court>"`. Si tu as touché plusieurs fichiers pour des raisons distinctes, fais des commits séparés et clairs.2909. POUSSE sur spbgit (le push fait partie du travail): `git push origin main`. NOTE: sur certains nœuds le push échoue en « no route to host » depuis une mission (macOS 26 Local Network Privacy) — c'est ATTENDU: le pousseur launchd (com.ka.pousseur) poussera tes commits dans les 5 minutes. Si le push échoue, signale-le simplement dans ton rapport final et passe à la suite; n'insiste pas et ne modifie JAMAIS ~/.ssh ni la config git pour « réparer » le push.291292== INTERDITS ABSOLUS ==293Toucher aux autres connecteurs ou aux autres apps du nœud; modifier la config pm2/ngrok/launchd ou le schéma de la base; installer des paquets globaux; supprimer des données; `git push --force` ou push d'une autre branche que main; toucher à ~/.ssh, aux clés API ou aux fichiers hors du repo. Ne désactive JAMAIS un connecteur pour « réparer » sa santé. Si la source est définitivement morte (site fermé, domaine à vendre, 404 permanent confirmé), ne force rien: documente et conclus en verdict site_source_mort.294295== FIN DE MISSION ==296Termine ta TOUTE DERNIÈRE réponse par un bloc JSON exactement de cette forme:297{{"verdict": "repare|echec|site_source_mort|rien_a_faire", "diagnostic": "cause racine en 1-2 phrases", "actions": "ce que tu as fait", "test": "résultat mesuré du test final (volume obtenu)", "fichiers": ["fichiers modifiés"]}}298Sois honnête: si ta réparation n'est pas prouvée par un test réel, le verdict est echec — un faux « repare » sera détecté par la supervision et rollback automatiquement."""299300301COMMON_RULES = """== RÈGLES COMMUNES (non négociables) ==302- Session 100 % AUTONOME: tu ne poses AUCUNE question, tu ne demandes AUCUNE validation, tu travailles jusqu'au bout. S'il faut trancher, tranche selon les conventions du repo et le bon sens, et documente ton choix dans le commit.303- Respecte l'architecture et les patterns du repo (regarde comment les connecteurs voisins sont faits AVANT d'écrire).304- Utilise le venv/node_modules du repo; jamais d'installation globale.305- INTERDIT d'inventer, stubber ou câbler des données en dur. Chaque donnée vient d'un vrai fetch.306- INTERDITS: toucher aux autres apps du nœud, config pm2/ngrok/launchd, schéma de la base (sauf migration prévue par le repo), `git push --force` ou push d'une autre branche que main, ~/.ssh, clés API, suppression de données.307- À la fin: `pm2 restart <process de sync>`, vérifie que le site répond (curl localhost), committe en messages clairs préfixés [%AGENT%], puis POUSSE sur spbgit: `git push origin main`. Un échec « no route to host » est ATTENDU depuis une mission (macOS 26 LNP) — le pousseur launchd s'en chargera dans les 5 minutes: signale-le et passe à la suite, sans toucher à ~/.ssh ni à la config git."""308309310def effort_stack(service: str) -> str:311 svc = SERVICES[service]312 return f"""== BOÎTE À OUTILS (clés dans ~/.claude/.env) ==313- DÉCOUVERTE: Serper (SERPER_API_KEY, POST https://google.serper.dev/search, gl=ca hl=fr) pour trouver sources/sitemaps/endpoints; Tavily (TAVILY_API_KEY) en complément.314- FETCH — escalade dans cet ordre: requêtes directes (requests/fetch du repo) → curl via subprocess (contourne les gotchas TLS de requests) → Scrapfly (SCRAPFLY_API_KEY, render_js/asp, backend=auto) → Bright Data Web Unlocker → proxies résidentiels Oxylabs (pr.oxylabs.io:7777, user avec -cc-CA, sessions collantes).315- APIFY (APIFY_TOKEN): en dernier recours pour les plateformes dures, acteurs maison du compte (proxy résidentiel intégré).316- Regarde d'abord comment le repo de {svc['app']} utilise déjà ces services et fais pareil."""317318319def build_effort_prompt(kind: str, service: str, source: str, note: str) -> str:320 svc = SERVICES[service]321 sync_proc = next((p for p in svc["pm2"] if "sync" in p or "etl" in p), (svc["pm2"][-1] if svc["pm2"] else "?"))322 node_alias_ = node_alias(svc["node"])323 rules = COMMON_RULES.replace("%AGENT%", AGENT)324 env = f"""== TON ENVIRONNEMENT ==325Tu es sur le nœud {node_alias_}, répertoire courant = repo de l'app: {svc['dir']} (remote-first, source de vérité).326L'app tourne via pm2 ({', '.join(svc['pm2'])}); process de sync: {sync_proc}; site local port {svc['web_port']}.327Commence par lire le CLAUDE.md / README du repo et docs/connecteurs/ s'ils existent."""328 if kind == "effort_degrade":329 block = LATEST["mine"].get(service) or {}330 degraded = [c for c in block.get("connectors", []) if c.get("status") == "degraded"][:25]331 deg_lines = "\n".join(332 f"- {c['source']}: dernier volume {c.get('found_last')} vs médiane {round(c.get('median_found') or 0)} "333 f"(dernier succès: {c.get('last_success')}, message: {c.get('message')})"334 for c in degraded) or "- (aucun connecteur dégradé au dernier scan — re-vérifie l'état réel dans l'app)"335 return f"""Tu es {AGENT}, agent gardien autonome du Groupe KA. EFFORT COMMANDÉ: INSPECTER LES CONNECTEURS DÉGRADÉS de l'app {svc['app']} (service {service}, site {svc.get('site', '')}).336{"Consigne de l'opérateur: " + note if note else ""}337338Un connecteur « dégradé » livre encore des données, mais moins de 50 % de sa médiane historique — souvent le signe d'une pagination cassée, d'un filtre qui se resserre, d'une section du site source disparue ou d'un blocage partiel.339340== CONNECTEURS DÉGRADÉS AU DERNIER SCAN api-ka ==341{deg_lines}342343{env}344345== DÉMARCHE ==3461. TRIE: pour chaque connecteur dégradé, regarde vite (logs + un fetch de contrôle) si la baisse est (a) réelle et réparable, (b) légitime (le site source a vraiment moins d'items — saison, inventaire réduit), ou (c) un blocage.3472. PRIORISE les cas (a) au plus fort potentiel de volume récupéré, et répare-les UN PAR UN: cause racine, correctif minimal, test réel avec volume mesuré avant/après. Commit séparé par connecteur réparé ([{AGENT}] fix connecteur <source>: …).3483. Pour les cas (b), ne touche à rien: consigne-les dans ton rapport final comme « baisse légitime ».3494. Traite autant de connecteurs que ton budget de temps le permet, en gardant la qualité: mieux vaut 3 vraies réparations prouvées que 10 rustines.3505. À la fin: `pm2 restart {sync_proc}`, vérifie le site (port {svc['web_port']}), puis pousse tous tes commits sur spbgit: `git push origin main`.351352{effort_stack(service)}353354{rules}355356== FIN DE MISSION ==357Termine ta TOUTE DERNIÈRE réponse par un bloc JSON:358{{"verdict": "livre|echec", "diagnostic": "portrait global des dégradés", "actions": "connecteurs réparés (avec volumes avant/après) / baisses légitimes / blocages", "test": "mesures", "fichiers": ["fichiers modifiés"]}}"""359 if kind == "effort_new":360 return f"""Tu es {AGENT}, agent gardien autonome du Groupe KA. EFFORT COMMANDÉ: ajouter UN NOUVEAU CONNECTEUR de qualité production à l'app {svc['app']} (service {service}, site {svc.get('site', '')}).361{"Consigne de l'opérateur: " + note if note else "Aucune consigne particulière: choisis la source la plus utile."}362363{env}364365== DÉMARCHE ==3661. CARTOGRAPHIE L'EXISTANT: liste les connecteurs actuels de l'app (registre/scheduler + docs/connecteurs) pour comprendre ce qui est déjà couvert et le pattern d'implémentation exact d'un connecteur (structure, signature, enregistrement, cache, dédup).3672. DÉCOUVRE avec Serper: cherche des sources québécoises pertinentes pour {svc['app']} NON couvertes (sitemaps, pages listant des items, endpoints JSON internes découverts via les pages). Évalue 3-5 candidats: volume estimé, faisabilité du fetch, qualité des données, stabilité. CHOISIS le meilleur.3683. IMPLÉMENTE le connecteur en suivant le pattern exact des connecteurs voisins: fetch (avec escalade si bloqué), parsing complet (tous les champs que l'app affiche, fiche détail incluse si le pattern le fait), dédup, enregistrement au registre/scheduler.3694. TESTE réellement: exécute-le au complet, il doit rapporter un volume substantiel et des items complets et corrects (vérifie 3-4 items à la main contre le site source).3705. DOCUMENTE: ajoute la fiche docs/connecteurs/<source>.md si le repo suit cette convention (URL, stratégie, pièges).3716. Redémarre {sync_proc}, vérifie le site, committe et pousse sur spbgit (`git push origin main`).372373{effort_stack(service)}374375{rules}376377== FIN DE MISSION ==378Termine ta TOUTE DERNIÈRE réponse par un bloc JSON:379{{"verdict": "livre|echec", "diagnostic": "source choisie et pourquoi", "actions": "ce qui a été construit", "test": "volume obtenu + vérifications", "fichiers": ["fichiers créés/modifiés"]}}"""380 return f"""Tu es {AGENT}, agent gardien autonome du Groupe KA. EFFORT COMMANDÉ: ENRICHIR le connecteur existant « {source} » de l'app {svc['app']} (service {service}, site {svc.get('site', '')}).381{"Consigne de l'opérateur: " + note if note else "Aucune consigne particulière: maximise la valeur (couverture, champs, robustesse)."}382383{env}384385== DÉMARCHE ==3861. ÉTUDIE le connecteur « {source} »: code, fiche docs/connecteurs, logs récents, volume actuel vs médiane, champs remplis vs champs que l'app sait afficher.3872. IDENTIFIE les enrichissements à plus forte valeur, par exemple: pagination complète (couvre TOUT le site source, pas la première page), fiche détail (champs manquants: descriptions, photos, prix, coordonnées, dates), robustesse (retries, escalade anti-bot, tolérance aux changements de DOM), fraîcheur (détection de retraits), dédup plus fine.3883. IMPLÉMENTE proprement dans le pattern du repo, SANS casser le format de sortie existant (les consommateurs avals dépendent du schéma).3894. TESTE réellement: exécute le connecteur au complet, compare avant/après (volume, complétude des champs), vérifie 3-4 items à la main.3905. METS À JOUR la fiche docs/connecteurs/<source>.md si elle existe.3916. Redémarre {sync_proc}, vérifie le site, committe et pousse sur spbgit (`git push origin main`).392393{effort_stack(service)}394395{rules}396397== FIN DE MISSION ==398Termine ta TOUTE DERNIÈRE réponse par un bloc JSON:399{{"verdict": "livre|echec", "diagnostic": "état initial constaté", "actions": "enrichissements livrés", "test": "avant/après mesuré", "fichiers": ["fichiers modifiés"]}}"""400401402def build_site_down_prompt(service: str, detail: dict[str, Any]) -> str:403 svc = SERVICES[service]404 web_proc = svc["pm2"][0] if svc["pm2"] else "?"405 node_alias_ = node_alias(svc["node"])406 return f"""Tu es {AGENT}, agent gardien autonome du Groupe KA. MISSION URGENTE: le SITE {svc.get('site', '')} de l'app {svc['app']} NE RÉPOND PLUS, et le redémarrage pm2 automatique n'a pas suffi. Investigue, trouve la cause racine et remets le site en ligne.407Session 100 % AUTONOME: tu ne poses aucune question, tu travailles jusqu'au bout.408409== CE QUE LE VEILLEUR A DÉJÀ CONSTATÉ ==410- échecs consécutifs du check public: {detail.get('fails', '?')} (site muet depuis {detail.get('down_since', '?')})411- dernière erreur du check public: {detail.get('error', '?')}412- redémarrage pm2 automatique: {detail.get('restart', 'tenté')} — santé locale mesurée juste après: {jdump(detail.get('health')) if detail.get('health') else 'inconnue'}413414== TON ENVIRONNEMENT ==415Tu es sur le nœud {node_alias_}, répertoire courant = repo de l'app: {svc['dir']} (remote-first, source de vérité).416Process pm2 de l'app: {', '.join(svc['pm2'])} (web: {web_proc}); le site local doit répondre sur http://localhost:{svc['web_port']}/ et le site public est {svc.get('site', '')} (tunnel ngrok sur ce nœud, souvent un process pm2 `<app>-ngrok` ou un service launchd).417418== DÉMARCHE IMPOSÉE (dans l'ordre) ==4191. CONSTATE: `curl -s -o /dev/null -w '%{{http_code}}' http://localhost:{svc['web_port']}/` puis la même chose sur {svc.get('site', 'le site public')}. `pm2 jlist` / `pm2 describe {web_proc}`: statut, nombre de restarts, uptime, mémoire.4202. LIS LES LOGS: `pm2 logs {web_proc} --nostream --lines 300` (stdout ET err) + dossier logs/ du repo. Cherche: crash en boucle, port déjà occupé (EADDRINUSE), exception au démarrage, OOM, module manquant, build corrompu, base de données verrouillée/corrompue, disque plein (`df -h`).4213. DISTINGUE les 3 familles de panne:422 a. Process web mort ou en boucle de crash → cause racine dans le code/l'env → corrige puis `pm2 restart {web_proc}`.423 b. Process online mais localhost:{svc['web_port']} muet ou en erreur → mauvais port/bind, build cassé → rebuild si le repo a une étape de build (npm run build…), puis redémarre.424 c. localhost OK mais site public muet → le tunnel: trouve le process (pm2 `*-ngrok` de CETTE app ou launchd) et relance-le (`pm2 restart <app>-ngrok` ou `launchctl kickstart -k gui/$(id -u)/<label>`). NE MODIFIE JAMAIS sa config ni son domaine.4254. CORRIGE la cause racine de façon minimale (pas de rustine qui masque le symptôme). Si des données/DB sont corrompues, répare prudemment SANS perte de données (copie de sauvegarde avant toute opération risquée).4265. VÉRIFIE pour de vrai: localhost ET le site public doivent répondre 200/3xx de façon STABLE (2 mesures espacées de 30 s).4276. Si tu as modifié du code ou de la config du repo: `git add -A && git commit -m "[{AGENT}] fix site down {service}: <cause>"` puis `git push origin main`. Un échec « no route to host » au push est ATTENDU (macOS 26 LNP) — le pousseur launchd s'en charge dans les 5 minutes: signale-le et passe à la suite, sans toucher à ~/.ssh ni à la config git.428429== INTERDITS ABSOLUS ==430Toucher aux autres apps du nœud; MODIFIER la config pm2/ngrok/launchd (relancer un process existant de CETTE app est permis, changer sa config non); supprimer des données; `git push --force` ou push d'une autre branche que main; ~/.ssh; clés API. Si la panne dépasse l'app (disque plein système, nœud malade), libère de l'espace UNIQUEMENT dans le repo de l'app (logs, caches, artefacts de build) et documente le reste dans ton diagnostic.431432== FIN DE MISSION ==433Termine ta TOUTE DERNIÈRE réponse par un bloc JSON exactement de cette forme:434{{"verdict": "repare|echec|rien_a_faire", "diagnostic": "cause racine en 1-2 phrases", "actions": "ce que tu as fait", "test": "codes HTTP local + public mesurés", "fichiers": ["fichiers modifiés"]}}435Sois honnête: « repare » exige que le SITE PUBLIC réponde réellement — le veilleur re-vérifie aux 2 minutes et un faux « repare » sera rollback."""436437438async def dispatch(incident: sqlite3.Row, health: dict[str, Any]) -> None:439 service, source, iid = incident["service"], incident["source"], incident["id"]440 svc = SERVICES[service]441 node = svc["node"]442 ip = node_ip(node)443 if not ip:444 hub.publish_sync({"kind": "log", "level": "error", "ts": now(),445 "msg": f"dispatch {source}: IP LAN inconnue pour le nœud {node_alias(node)} (registre incomplet)"})446 return447 if not runner_available(node, f"mission {source}"):448 return # l'incident reste ouvert ; on retentera quand le runner existera449 mid = uuid.uuid4().hex[:12]450 # IP du gardien lui-même (callback du runner) : env launchd, sinon le registre mld.451 my_ip = os.environ.get("KA_GUARDIAN_SELF_IP") or (registry.app_entry(AGENT) or {}).get("ip") or "192.168.2.69"452 kind = incident["status_detected"]453 max_cost = None454 if kind in ("effort_new", "effort_enrich", "effort_degrade"):455 prompt = build_effort_prompt(kind, service, source, health.get("note", ""))456 max_turns = POLICY.get("effort_max_turns", 150)457 timeout = POLICY.get("effort_timeout_seconds", 7200)458 # Plafonds choisis par l'opérateur au moment de commander l'effort.459 if health.get("max_minutes"):460 timeout = min(timeout, int(health["max_minutes"]) * 60)461 if health.get("max_cost_usd"):462 max_cost = float(health["max_cost_usd"])463 elif kind == "site_down":464 prompt = build_site_down_prompt(service, health)465 max_turns, timeout = POLICY["mission_max_turns"], POLICY["mission_timeout_seconds"]466 else:467 prompt = build_prompt(service, source, health)468 max_turns, timeout = POLICY["mission_max_turns"], POLICY["mission_timeout_seconds"]469 payload = {470 "mission_id": mid, "agent": AGENT, "service": service, "source": source,471 "dir": svc["dir"], "pm2": svc["pm2"], "web_port": svc["web_port"],472 "model": ME.get("model", "sonnet"),473 "max_turns": max_turns,474 "timeout_seconds": timeout,475 "max_cost_usd": max_cost,476 "prompt": prompt,477 "callback_url": f"http://{my_ip}:{ME['port']}/api/ingest/{mid}",478 }479 try:480 code, body = await lan_post(ip, RUNNER_PORT, "/missions", payload, timeout=25)481 if code == 409:482 return # runner occupé, on retentera au prochain tick483 if code != 200:484 raise ConnectionError(f"runner {code}: {body[:200]}")485 except Exception as exc:486 print(f"[dispatch] {service}/{source} → {type(exc).__name__}: {exc}", flush=True)487 hub.publish_sync({"kind": "log", "level": "error", "msg": f"dispatch {source}: {exc}", "ts": now()})488 return489 with db() as c:490 c.execute("INSERT INTO missions(id,incident_id,service,source,node,state,started) VALUES(?,?,?,?,?,?,?)",491 (mid, iid, service, source, node, "running", now()))492 set_incident(c, iid, state="fixing", attempts=incident["attempts"] + 1)493 incident_event(iid, service, source, "fixing", f"mission {mid} dépêchée sur {node_alias(node)}")494 hub.publish_sync({"kind": "mission", "mission_id": mid, "service": service,495 "source": source, "state": "running", "ts": now()})496497498async def rollback_mission(mission: sqlite3.Row, reason: str) -> dict[str, Any]:499 svc = SERVICES[mission["service"]]500 # Le repo (et son historique git) suit l'app : si elle a déménagé depuis la501 # mission, le rollback vise le nœud ACTUEL du service, pas celui de la mission.502 node = svc.get("node") or mission["node"]503 ip = node_ip(node) or node_ip(mission["node"])504 if node != mission["node"]:505 hub.publish_sync({"kind": "log", "level": "warn", "ts": now(),506 "msg": f"rollback {mission['source']}: l'app a déménagé {node_alias(mission['node'])} → {node_alias(node)}, rollback sur le nœud actuel"})507 if not ip or not runner_available(node, f"rollback {mission['source']}"):508 out = {"ok": False, "erreur": f"runner indisponible sur {node_alias(node)}"}509 hub.publish_sync({"kind": "rollback", "mission_id": mission["id"], "service": mission["service"],510 "source": mission["source"], "reason": reason, "result": out, "ts": now()})511 return out512 try:513 code, body = await lan_post(ip, RUNNER_PORT, "/rollback",514 {"dir": svc["dir"], "base_commit": mission["base_commit"],515 "pm2": svc["pm2"], "web_port": svc["web_port"]}, timeout=300)516 out = json.loads(body) if code == 200 else {"ok": False, "erreur": f"{code}: {body[:200]}"}517 except Exception as exc:518 out = {"ok": False, "erreur": str(exc)}519 hub.publish_sync({"kind": "rollback", "mission_id": mission["id"], "service": mission["service"],520 "source": mission["source"], "reason": reason, "result": out, "ts": now()})521 with db() as c:522 c.execute("INSERT INTO events(mission_id,ts,type,data) VALUES(?,?,?,?)",523 (mission["id"], now(), "rollback", jdump({"reason": reason, "result": out})))524 return out525526527# --------------------------------------------------------------- moteur ----528529async def poll_once() -> None:530 async with httpx.AsyncClient() as cl:531 r = await cl.get(TOPO["apika_monitoring_url"], timeout=15)532 data = r.json()["data"]533 mine: dict[str, Any] = {}534 for service, block in data["services"].items():535 assigned_to_me = service in MY_SERVICES or (DEFAULT_AGENT and service not in ASSIGNED)536 if assigned_to_me:537 mine[service] = block538 LATEST.update({"mine": mine, "ecosystem": data["summary"], "polled": now()})539 with db() as c:540 summary_mine = {s: b["summary"] for s, b in mine.items()}541 c.execute("INSERT OR REPLACE INTO snapshots(ts,mine,ecosystem) VALUES(?,?,?)",542 (now(), jdump(summary_mine), jdump(data["summary"])))543 c.execute("DELETE FROM snapshots WHERE ts < ?", (now() - 30 * 86400,))544 hub.publish_sync({"kind": "snapshot", "mine": {s: b["summary"] for s, b in mine.items()},545 "ecosystem": data["summary"], "ts": now()})546547 trigger = set(POLICY["trigger_statuses"])548 with db() as c:549 for service, block in mine.items():550 if service not in SERVICES:551 continue # service inconnu de la topologie: visible au dashboard, pas d'action552 to_create: list[tuple[str, str, dict[str, Any]]] = []553 for conn in block.get("connectors", []):554 source, status = conn["source"], conn["status"]555 inc = active_incident(c, service, source)556 if status in trigger and inc is None:557 # Garde anti-boucle : une source déjà diagnostiquée irréparable558 # (site_source_mort, tentatives épuisées) ne redonne pas559 # d'incident tant que le délai de re-examen n'est pas écoulé.560 # Sans elle, un connecteur désactivé côté app mais rapporté561 # broken à perpétuité par api-ka relance une mission toutes562 # les ~10 min (2026-08-24 : 37 missions sur origin-north.ai).563 ab = c.execute(564 "SELECT updated FROM incidents WHERE service=? AND source=? "565 "AND state='abandoned' ORDER BY updated DESC LIMIT 1",566 (service, source)).fetchone()567 if ab and now() - ab["updated"] < POLICY.get("abandoned_retry_hours", 168) * 3600:568 continue569 to_create.append((source, status, conn))570 elif status in trigger and inc is not None and inc["state"] in ("open", "cooldown"):571 # Ne pas rafraîchir watching/fixing: `updated` sert de chrono572 # à la fenêtre de surveillance et au cooldown.573 c.execute("UPDATE incidents SET status_detected=?, detail=? WHERE id=?",574 (status, jdump(conn), inc["id"]))575 elif status == "ok" and inc is not None:576 if inc["state"] == "watching":577 set_incident(c, inc["id"], state="resolved", resolved=now())578 incident_event(inc["id"], service, source, "resolved", "connecteur de retour à ok — réparation confirmée")579 elif inc["state"] in ("open", "cooldown"):580 set_incident(c, inc["id"], state="self_healed", resolved=now())581 incident_event(inc["id"], service, source, "self_healed", "revenu à ok sans intervention")582 elif status == "retired" and inc is not None and inc["state"] in ("open", "cooldown", "watching"):583 # api-ka a retiré la source de la supervision (désactivée584 # côté app ou disparue du journal) : plus rien à réparer.585 set_incident(c, inc["id"], state="abandoned")586 incident_event(inc["id"], service, source, "abandoned", "source retirée de la supervision api-ka")587588 # Damper d'événement de masse : des centaines de sources qui589 # basculent d'un coup signalent une panne GLOBALE de l'app (sync590 # mort, /api/stats cassé), pas autant de pannes individuelles —591 # missionner source par source serait long et coûteux pour rien592 # (2026-08-24 : 794 incidents jobka créés après ~17 h d'arrêt de593 # job-ka-sync). On n'ouvre rien : la panne d'app se voit au594 # dashboard (bloc _app) et les sources encore cassées une fois la595 # vague retombée sous le seuil recevront leurs incidents.596 limit = POLICY.get("mass_incident_threshold", 20)597 if len(to_create) > limit:598 msg = (f"vague de {len(to_create)} connecteurs en panne sur {service} "599 f"(seuil {limit}) — panne globale probable de l'app, aucun incident créé")600 print(f"[poll] {msg}", flush=True)601 hub.publish_sync({"kind": "log", "level": "error", "msg": msg, "ts": now()})602 continue603 for source, status, conn in to_create:604 iid = uuid.uuid4().hex[:10]605 c.execute("INSERT INTO incidents(id,service,source,status_detected,state,created,updated,detail) "606 "VALUES(?,?,?,?,?,?,?,?)",607 (iid, service, source, status, "open", now(), now(), jdump(conn)))608 incident_event(iid, service, source, "open", f"détecté {status}")609610611async def tick() -> None:612 apply_registry() # même en pause : la topologie doit rester juste613 if LATEST["paused"]:614 return615 try:616 await poll_once()617 except Exception as exc:618 hub.publish_sync({"kind": "log", "level": "warn", "msg": f"poll api-ka échoué: {exc}", "ts": now()})619620 with db() as c:621 # 1. Réconciliation des missions en cours depuis > 10 min: on demande au622 # runner ce qu'il en est (événement final perdu? runner redémarré?).623 for m in c.execute("SELECT id FROM missions WHERE state='running' AND started < ?",624 (now() - 600,)).fetchall():625 if m["id"] not in RECONCILING:626 RECONCILING.add(m["id"])627 asyncio.create_task(reconcile_mission(m["id"]))628629 # 2. Watching → rollback si la fenêtre est passée et toujours cassé630 # (fenêtre courte pour les incidents de site: le veilleur mesure aux631 # 2 minutes, pas besoin d'attendre le rythme des scans api-ka)632 for inc in c.execute("SELECT * FROM incidents WHERE state='watching'").fetchall():633 window_h = POLICY.get("site_watch_window_hours", 1) if inc["source"] == "_site" else POLICY["watch_window_hours"]634 if now() - inc["updated"] < window_h * 3600:635 continue636 status = current_status(inc["service"], inc["source"])637 m = c.execute("SELECT * FROM missions WHERE incident_id=? AND base_commit IS NOT NULL "638 "ORDER BY started DESC LIMIT 1", (inc["id"],)).fetchone()639 if status in POLICY["trigger_statuses"] or status == "degraded":640 if m and (json.loads(m["commits"] or "[]")):641 asyncio.create_task(rollback_mission(m, "toujours cassé après la fenêtre de surveillance"))642 set_incident(c, inc["id"], state="cooldown")643 incident_event(inc["id"], inc["service"], inc["source"], "cooldown",644 f"non guéri après {POLICY['watch_window_hours']}h → rollback + cooldown")645 elif status == "ok":646 set_incident(c, inc["id"], state="resolved", resolved=now())647 incident_event(inc["id"], inc["service"], inc["source"], "resolved", "confirmé ok")648649 # 3. Cooldown expiré → réouverture ou abandon650 for inc in c.execute("SELECT * FROM incidents WHERE state='cooldown'").fetchall():651 cool_h = POLICY.get("site_cooldown_hours", 1) if inc["source"] == "_site" else POLICY["attempt_cooldown_hours"]652 if now() - inc["updated"] < cool_h * 3600:653 continue654 status = current_status(inc["service"], inc["source"])655 if status == "ok":656 set_incident(c, inc["id"], state="resolved", resolved=now())657 incident_event(inc["id"], inc["service"], inc["source"], "resolved", "guéri pendant le cooldown")658 elif inc["attempts"] >= POLICY["max_attempts_per_incident"]:659 set_incident(c, inc["id"], state="abandoned")660 incident_event(inc["id"], inc["service"], inc["source"], "abandoned",661 f"{inc['attempts']} tentatives épuisées — intervention humaine requise")662 else:663 set_incident(c, inc["id"], state="open")664 incident_event(inc["id"], inc["service"], inc["source"], "open", "cooldown terminé, nouvelle tentative")665666 # 4. Dispatch (1 mission à la fois par agent, broken avant stale, plus vieux d'abord)667 running = c.execute("SELECT COUNT(*) n FROM missions WHERE state='running'").fetchone()["n"]668 if running < POLICY["max_concurrent_missions"]:669 busy_nodes = {m["node"] for m in c.execute("SELECT node FROM missions WHERE state='running'").fetchall()}670 nxt = c.execute(671 "SELECT * FROM incidents WHERE state='open' "672 "ORDER BY CASE status_detected WHEN 'site_down' THEN 0 WHEN 'manual' THEN 0 "673 "WHEN 'effort_new' THEN 0 WHEN 'effort_enrich' THEN 0 WHEN 'broken' THEN 1 ELSE 2 END, created "674 ).fetchall()675 for inc in nxt:676 if inc["service"] not in SERVICES:677 continue678 if SERVICES[inc["service"]]["node"] in busy_nodes:679 continue680 set_incident(c, inc["id"], state="dispatched")681 asyncio.create_task(dispatch_row(inc["id"]))682 break683684685def current_status(service: str, source: str) -> str:686 if source == "_site":687 st = SITE.get(service) or {}688 if not st.get("checked"):689 return "inconnu"690 return "ok" if st.get("ok") else "broken"691 block = LATEST["mine"].get(service) or {}692 for conn in block.get("connectors", []):693 if conn["source"] == source:694 return conn["status"]695 return "inconnu"696697698async def dispatch_row(iid: str) -> None:699 with db() as c:700 inc = c.execute("SELECT * FROM incidents WHERE id=?", (iid,)).fetchone()701 if inc:702 await dispatch(inc, json.loads(inc["detail"] or "{}"))703 # dispatch() remet fixing; si l'appel a échoué/409, on relâche704 with db() as c:705 cur = c.execute("SELECT state FROM incidents WHERE id=?", (iid,)).fetchone()706 if cur and cur["state"] == "dispatched":707 set_incident(c, iid, state="open")708709710RECONCILING: set[str] = set()711712713async def reconcile_mission(mid: str) -> None:714 """Résout une mission « en cours » suspecte auprès de son runner."""715 try:716 with db() as c:717 m = c.execute("SELECT * FROM missions WHERE id=? AND state='running'", (mid,)).fetchone()718 if not m:719 return720 ip = node_ip(m["node"])721 if not ip:722 return # nœud disparu du registre et de topology.json : on retentera723 try:724 code, body = await lan_post(ip, RUNNER_PORT, "/mission-result", {"mission_id": mid}, timeout=20)725 resp = json.loads(body) if code == 200 else {"status": "erreur"}726 except Exception:727 return # courrier/runner injoignable: on retentera au prochain tick728 if resp.get("status") == "running":729 return730 with db() as c:731 m = c.execute("SELECT * FROM missions WHERE id=? AND state='running'", (mid,)).fetchone()732 if not m:733 return734 if resp.get("status") == "done":735 data = (resp.get("final") or {}).get("data", {})736 c.execute("INSERT INTO events(mission_id,ts,type,data) VALUES(?,?,?,?)",737 (mid, now(), "final", jdump(data)))738 await finalize(c, m, "final", data)739 hub.publish_sync({"kind": "mission_event", "mission_id": mid, "type": "final",740 "data": data, "ts": now()})741 elif now() - m["started"] > 900:742 # Le runner ne connaît pas cette mission: tuée en vol743 # (redémarrage). On la clôt et on relance l'incident.744 c.execute("UPDATE missions SET state='error', ended=? WHERE id=?", (now(), mid))745 inc = c.execute("SELECT * FROM incidents WHERE id=?", (m["incident_id"],)).fetchone()746 if inc and inc["state"] in ("fixing", "dispatched"):747 set_incident(c, inc["id"], state="open")748 incident_event(inc["id"], m["service"], m["source"], "open",749 "mission perdue (runner interrompu) — remise en file")750 finally:751 RECONCILING.discard(mid)752753754async def engine() -> None:755 global LOOP756 LOOP = asyncio.get_running_loop()757 await asyncio.sleep(3)758 while True:759 try:760 await tick()761 except Exception as exc:762 import traceback763 traceback.print_exc()764 hub.publish_sync({"kind": "log", "level": "error", "msg": f"tick: {exc}", "ts": now()})765 await asyncio.sleep(POLICY["poll_interval_seconds"])766767768# ----------------------------------------------------- veilleur de sites ----769# Surveillance du SITE WEB de chaque service assigné (indépendante des770# connecteurs api-ka). Down confirmé → pm2 restart automatique via le runner771# du nœud; toujours down → incident « site_down » et mission d'investigation.772773SITE: dict[str, dict[str, Any]] = {} # service → état du dernier check public774775776async def check_site(url: str) -> tuple[bool, str]:777 try:778 async with httpx.AsyncClient(follow_redirects=True, timeout=15) as cl:779 r = await cl.get(url, headers={"User-Agent": f"ka-guardian-{AGENT}/1.0"})780 return r.status_code < 500, f"HTTP {r.status_code}"781 except Exception as exc:782 return False, f"{type(exc).__name__}: {exc}"783784785async def runner_restart(service: str) -> dict[str, Any]:786 """pm2 restart de l'app via le runner de son nœud (endpoint /restart)."""787 svc = SERVICES[service]788 ip = node_ip(svc["node"])789 if not ip or not runner_available(svc["node"], f"pm2 restart {svc['app']}"):790 return {"ok": False, "erreur": f"runner indisponible sur {node_alias(svc['node'])}"}791 try:792 code, body = await lan_post(ip, RUNNER_PORT, "/restart",793 {"pm2": svc["pm2"], "web_port": svc["web_port"]}, timeout=240)794 return json.loads(body) if code == 200 else {"ok": False, "erreur": f"runner {code}: {body[:200]}"}795 except Exception as exc:796 return {"ok": False, "erreur": f"{type(exc).__name__}: {exc}"}797798799async def site_check_service(service: str) -> None:800 svc = SERVICES[service]801 url = svc.get("site")802 if not url:803 return804 up, info = await check_site(url)805 if not up: # contre-mesure avant de compter l'échec (blip réseau/ngrok)806 await asyncio.sleep(5)807 up, info = await check_site(url)808 st = SITE.setdefault(service, {"fails": 0, "ok": True, "down_since": None, "last_restart": 0.0})809 st.update({"checked": now(), "ok": up, "info": info})810 if up:811 if st["fails"]:812 hub.publish_sync({"kind": "log", "level": "info",813 "msg": f"site {url} de retour en ligne ({info})", "ts": now()})814 st["fails"], st["down_since"] = 0, None815 with db() as c:816 inc = active_incident(c, service, "_site")817 if inc and inc["state"] in ("open", "cooldown"):818 set_incident(c, inc["id"], state="self_healed", resolved=now())819 incident_event(inc["id"], service, "_site", "self_healed", f"site de retour en ligne ({info})")820 return821 st["fails"] += 1822 st["down_since"] = st["down_since"] or now()823 need = POLICY.get("site_check_fails", 2)824 hub.publish_sync({"kind": "log", "level": "warn",825 "msg": f"site {url} muet ({info}) — échec {st['fails']}/{need}", "ts": now()})826 if st["fails"] < need:827 return828 with db() as c:829 inc = active_incident(c, service, "_site")830 if inc and inc["state"] in ("dispatched", "fixing"):831 return # une mission d'investigation travaille déjà sur cette app832 # Étape 1: pm2 restart automatique (throttlé pour ne pas marteler l'app)833 restart_res: dict[str, Any] = {"ok": False, "erreur": "non tenté (throttle)"}834 if now() - st["last_restart"] >= POLICY.get("site_restart_throttle_seconds", 600):835 st["last_restart"] = now()836 hub.publish_sync({"kind": "log", "level": "warn",837 "msg": f"{svc['app']}: site down confirmé → pm2 restart automatique ({', '.join(svc['pm2'])})",838 "ts": now()})839 restart_res = await runner_restart(service)840 await asyncio.sleep(POLICY.get("site_restart_wait_seconds", 20))841 up, info = await check_site(url)842 if up:843 st.update({"ok": True, "fails": 0, "down_since": None, "info": info})844 with db() as c:845 if inc:846 set_incident(c, inc["id"], state="resolved", resolved=now())847 incident_event(inc["id"], service, "_site", "resolved", f"site rétabli par pm2 restart ({info})")848 else:849 # Trace au dashboard: le restart automatique a suffi.850 iid = uuid.uuid4().hex[:10]851 c.execute("INSERT INTO incidents(id,service,source,status_detected,state,created,updated,resolved,detail) "852 "VALUES(?,?,?,?,?,?,?,?,?)",853 (iid, service, "_site", "site_down", "resolved", now(), now(), now(),854 jdump({"note": "rétabli par pm2 restart automatique", "restart": restart_res})))855 incident_event(iid, service, "_site", "resolved", f"site rétabli par pm2 restart automatique ({info})")856 return857 # Étape 2: toujours down → incident + mission d'investigation (si aucun actif)858 if inc is not None:859 return # open/watching: la machinerie d'incidents suit déjà son cours860 with db() as c:861 # Garde anti-boucle (même logique que les connecteurs, délai plus court:862 # un site down est urgent, on re-tente quand même chaque jour).863 ab = c.execute("SELECT updated FROM incidents WHERE service=? AND source='_site' "864 "AND state='abandoned' ORDER BY updated DESC LIMIT 1", (service,)).fetchone()865 if ab and now() - ab["updated"] < POLICY.get("site_abandoned_retry_hours", 24) * 3600:866 return867 detail = {"fails": st["fails"],868 "down_since": time.strftime("%Y-%m-%d %H:%M", time.localtime(st["down_since"])),869 "error": info,870 "restart": "pm2 restart exécuté, site toujours muet" if restart_res.get("ok")871 else f"restart en échec: {restart_res.get('erreur') or restart_res}",872 "health": restart_res.get("health")}873 iid = uuid.uuid4().hex[:10]874 c.execute("INSERT INTO incidents(id,service,source,status_detected,state,created,updated,detail) "875 "VALUES(?,?,?,?,?,?,?,?)",876 (iid, service, "_site", "site_down", "open", now(), now(), jdump(detail)))877 incident_event(iid, service, "_site", "open",878 f"site hors ligne malgré le pm2 restart automatique ({info}) — mission d'investigation demandée")879880881async def site_engine() -> None:882 await asyncio.sleep(10)883 while True:884 if not LATEST["paused"]:885 for service in sorted(MY_SERVICES & set(SERVICES)):886 try:887 await site_check_service(service)888 except Exception as exc:889 hub.publish_sync({"kind": "log", "level": "warn",890 "msg": f"veilleur site {service}: {exc}", "ts": now()})891 await asyncio.sleep(POLICY.get("site_check_interval_seconds", 120))892893894# ------------------------------------------------------------------- app ---895896app = FastAPI(title=f"KA Guardian — {AGENT}", docs_url=None, redoc_url=None)897898899@app.on_event("startup")900async def startup() -> None:901 init_db()902 asyncio.create_task(engine())903 asyncio.create_task(site_engine())904905906def check_token(tok: str | None) -> None:907 if not TOKEN or tok != TOKEN:908 raise HTTPException(status_code=401)909910911@app.post("/api/ingest/{mission_id}")912async def ingest(mission_id: str, req: Request, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]:913 check_token(x_ka_token)914 ev = await req.json()915 etype, data = ev.get("type"), ev.get("data", {})916 with db() as c:917 c.execute("INSERT INTO events(mission_id,ts,type,data) VALUES(?,?,?,?)",918 (mission_id, now(), etype, jdump(data)))919 m = c.execute("SELECT * FROM missions WHERE id=?", (mission_id,)).fetchone()920 if m:921 if etype == "start":922 c.execute("UPDATE missions SET base_commit=? WHERE id=?", (data.get("base_commit"), mission_id))923 elif etype in ("final", "error"):924 await finalize(c, m, etype, data)925 hub.publish_sync({"kind": "mission_event", "mission_id": mission_id, "type": etype,926 "data": data, "ts": now()})927 return {"ok": True}928929930async def finalize(c: sqlite3.Connection, m: sqlite3.Row, etype: str, data: dict[str, Any]) -> None:931 inc = c.execute("SELECT * FROM incidents WHERE id=?", (m["incident_id"],)).fetchone()932 if etype == "error":933 c.execute("UPDATE missions SET state='error', ended=? WHERE id=?", (now(), m["id"]))934 if inc:935 set_incident(c, inc["id"], state="cooldown")936 incident_event(inc["id"], m["service"], m["source"], "cooldown", f"erreur mission: {data.get('error')}")937 return938 verdict = data.get("verdict", {})939 commits = data.get("commits", [])940 health = data.get("health", {})941 c.execute("UPDATE missions SET state='done', ended=?, commits=?, verdict=?, cost_usd=?, num_turns=?, health=?, base_commit=COALESCE(base_commit,?) WHERE id=?",942 (now(), jdump(commits), jdump(verdict), data.get("cost_usd"),943 data.get("num_turns"), jdump(health), data.get("base_commit"), m["id"]))944 m2 = c.execute("SELECT * FROM missions WHERE id=?", (m["id"],)).fetchone()945 v = verdict.get("verdict", "inconnu")946 if not inc:947 return948 if inc["status_detected"] in ("effort_new", "effort_enrich", "effort_degrade"):949 # Les efforts commandés ne passent pas par la surveillance api-ka:950 # verdict + santé de l'app décident tout de suite.951 if not health.get("ok", True):952 asyncio.create_task(rollback_mission(m2, "healthcheck app en échec post-effort"))953 set_incident(c, inc["id"], state="abandoned")954 incident_event(inc["id"], m["service"], m["source"], "abandoned", "app en mauvaise santé → ROLLBACK immédiat")955 elif v == "livre":956 set_incident(c, inc["id"], state="resolved", resolved=now())957 incident_event(inc["id"], m["service"], m["source"], "resolved",958 f"effort livré ({len(commits)} commit{'s' if len(commits) > 1 else ''})")959 elif v == "echec" and commits:960 asyncio.create_task(rollback_mission(m2, "effort en échec déclaré → rollback préventif"))961 set_incident(c, inc["id"], state="abandoned")962 incident_event(inc["id"], m["service"], m["source"], "abandoned", "effort en échec → rollback")963 else:964 # Session interrompue (plafond durée/coût) ou verdict illisible:965 # l'app est saine, les commits incrémentaux testés sont conservés.966 set_incident(c, inc["id"], state="abandoned")967 incident_event(inc["id"], m["service"], m["source"], "abandoned",968 f"effort terminé sans verdict livré ({v}) — {len(commits)} commit(s) conservé(s), app saine")969 return970 if not health.get("ok", True):971 # L'app est tombée → rollback immédiat, quoi qu'ait dit l'agent.972 asyncio.create_task(rollback_mission(m2, "healthcheck app en échec post-mission"))973 set_incident(c, inc["id"], state="cooldown")974 incident_event(inc["id"], m["service"], m["source"], "cooldown", "app en mauvaise santé → ROLLBACK immédiat")975 elif v == "repare" and commits:976 set_incident(c, inc["id"], state="watching")977 incident_event(inc["id"], m["service"], m["source"], "watching",978 f"réparation déclarée ({len(commits)} commit) — surveillance {POLICY['watch_window_hours']}h")979 elif v in ("echec", "inconnu") and commits:980 asyncio.create_task(rollback_mission(m2, f"verdict {v} avec commits → annulation par prudence"))981 set_incident(c, inc["id"], state="cooldown")982 incident_event(inc["id"], m["service"], m["source"], "cooldown", f"verdict {v} → rollback préventif")983 elif v == "site_source_mort":984 set_incident(c, inc["id"], state="abandoned")985 incident_event(inc["id"], m["service"], m["source"], "abandoned", "source définitivement morte (diagnostic agent)")986 elif v in ("rien_a_faire", "repare"):987 set_incident(c, inc["id"], state="watching")988 incident_event(inc["id"], m["service"], m["source"], "watching", f"verdict {v} — surveillance")989 else:990 set_incident(c, inc["id"], state="cooldown")991 incident_event(inc["id"], m["service"], m["source"], "cooldown", f"verdict {v}")992993994@app.get("/api/state")995def state() -> dict[str, Any]:996 with db() as c:997 incidents = [dict(r) for r in c.execute(998 "SELECT * FROM incidents ORDER BY created DESC LIMIT 200").fetchall()]999 missions = [dict(r) for r in c.execute(1000 "SELECT * FROM missions ORDER BY started DESC LIMIT 100").fetchall()]1001 stats = c.execute("""SELECT1002 (SELECT COUNT(*) FROM missions) n_missions,1003 (SELECT COUNT(*) FROM incidents WHERE state='resolved') n_resolved,1004 (SELECT COUNT(*) FROM incidents WHERE state IN ('open','dispatched','fixing','watching','cooldown')) n_active,1005 (SELECT COUNT(*) FROM incidents WHERE state='abandoned') n_abandoned,1006 (SELECT COUNT(*) FROM events WHERE type='rollback') n_rollbacks,1007 (SELECT ROUND(SUM(cost_usd),2) FROM missions) cost_total,1008 (SELECT ROUND(AVG(resolved-created)/3600.0,1) FROM incidents WHERE state='resolved' AND resolved IS NOT NULL) mttr_h1009 """).fetchone()1010 history = [dict(r) for r in c.execute(1011 "SELECT ts, mine FROM snapshots WHERE ts > ? ORDER BY ts", (now() - 7 * 86400,)).fetchall()]1012 return {1013 "agent": AGENT, "identity": {k: ME.get(k) for k in ("domain", "accent", "accent2", "accent_soft", "tagline", "model")},1014 "services": {s: {**SERVICES[s], "node_alias": node_alias(SERVICES[s]["node"]),1015 "lan_ip": node_ip(SERVICES[s]["node"]),1016 "runner": registry.runner_status(node_ip(SERVICES[s]["node"]))}1017 for s in ME["services"] if s in SERVICES},1018 "registry": registry.info(),1019 "siblings": {a: {"domain": v["domain"], "accent": v["accent"], "services": v["services"]}1020 for a, v in TOPO["agents"].items() if a != AGENT},1021 "latest": LATEST, "sites": SITE, "incidents": incidents, "missions": missions,1022 "stats": dict(stats), "history": history, "policy": POLICY, "now": now(),1023 }102410251026@app.get("/api/missions/{mission_id}")1027def mission_detail(mission_id: str) -> dict[str, Any]:1028 with db() as c:1029 m = c.execute("SELECT * FROM missions WHERE id=?", (mission_id,)).fetchone()1030 if not m:1031 raise HTTPException(status_code=404)1032 evs = [dict(r) for r in c.execute(1033 "SELECT ts,type,data FROM events WHERE mission_id=? ORDER BY id", (mission_id,)).fetchall()]1034 return {"mission": dict(m), "events": evs}103510361037@app.get("/events")1038async def sse(request: Request) -> StreamingResponse:1039 q: asyncio.Queue = asyncio.Queue()1040 hub.clients.add(q)10411042 async def gen():1043 try:1044 yield f"data: {jdump({'kind': 'hello', 'agent': AGENT, 'ts': now()})}\n\n"1045 while True:1046 if await request.is_disconnected():1047 break1048 try:1049 msg = await asyncio.wait_for(q.get(), timeout=25)1050 yield f"data: {jdump(msg)}\n\n"1051 except asyncio.TimeoutError:1052 yield ": keepalive\n\n"1053 finally:1054 hub.clients.discard(q)10551056 return StreamingResponse(gen(), media_type="text/event-stream",1057 headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})105810591060@app.get("/api/connectors/{service}")1061def connectors_list(service: str) -> dict[str, Any]:1062 """Sources connues du service (pour le menu déroulant d'enrichissement)."""1063 block = LATEST["mine"].get(service) or {}1064 return {"service": service,1065 "connectors": sorted(1066 ({"source": c["source"], "status": c["status"]} for c in block.get("connectors", [])),1067 key=lambda x: x["source"])}106810691070@app.post("/api/admin/effort")1071async def admin_effort(req: Request, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]:1072 """Commande un effort autonome: nouveau connecteur ou enrichissement."""1073 check_token(x_ka_token)1074 body = await req.json()1075 kind = body.get("kind")1076 service = body.get("service")1077 note = (body.get("note") or "").strip()[:2000]1078 if kind not in ("effort_new", "effort_enrich", "effort_degrade") or service not in SERVICES:1079 raise HTTPException(status_code=400, detail="kind ou service invalide")1080 if kind == "effort_enrich":1081 source = (body.get("source") or "").strip()1082 if not source:1083 raise HTTPException(status_code=400, detail="source requise pour un enrichissement")1084 elif kind == "effort_degrade":1085 source = f"dégradés·{uuid.uuid4().hex[:6]}"1086 else:1087 source = f"nouveau·{uuid.uuid4().hex[:6]}"1088 detail: dict[str, Any] = {"note": note, "kind": kind}1089 try:1090 if body.get("max_minutes"):1091 detail["max_minutes"] = max(10, min(240, int(body["max_minutes"])))1092 if body.get("max_cost_usd"):1093 detail["max_cost_usd"] = max(0.5, min(100.0, float(body["max_cost_usd"])))1094 except (TypeError, ValueError):1095 raise HTTPException(status_code=400, detail="max_minutes/max_cost_usd invalides")1096 iid = uuid.uuid4().hex[:10]1097 with db() as c:1098 c.execute("INSERT INTO incidents(id,service,source,status_detected,state,created,updated,detail) "1099 "VALUES(?,?,?,?,?,?,?,?)",1100 (iid, service, source, kind, "open", now(), now(), jdump(detail)))1101 labels = {"effort_new": "nouveau connecteur", "effort_enrich": f"enrichissement de {source}",1102 "effort_degrade": "inspection des connecteurs dégradés"}1103 incident_event(iid, service, source, "open",1104 f"effort commandé: {labels[kind]}"1105 + (f" · ≤{detail['max_minutes']} min" if detail.get("max_minutes") else "")1106 + (f" · ≤{detail['max_cost_usd']} $" if detail.get("max_cost_usd") else ""))1107 return {"ok": True, "incident_id": iid}110811091110@app.post("/api/admin/mission")1111async def admin_mission(req: Request, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]:1112 """Déclenche manuellement une mission sur un connecteur (même sain)."""1113 check_token(x_ka_token)1114 body = await req.json()1115 service, source = body["service"], body["source"]1116 if service not in SERVICES:1117 raise HTTPException(status_code=400, detail="service inconnu")1118 detail = {"source": source, "status": "manual", "message": body.get("note", "mission manuelle")}1119 for conn in (LATEST["mine"].get(service) or {}).get("connectors", []):1120 if conn["source"] == source:1121 detail = {**conn, "status": "manual"}1122 with db() as c:1123 inc = active_incident(c, service, source)1124 if inc is None:1125 iid = uuid.uuid4().hex[:10]1126 c.execute("INSERT INTO incidents(id,service,source,status_detected,state,created,updated,detail) "1127 "VALUES(?,?,?,?,?,?,?,?)",1128 (iid, service, source, "manual", "open", now(), now(), jdump(detail)))1129 else:1130 iid = inc["id"]1131 set_incident(c, iid, state="open", status_detected="manual")1132 incident_event(iid, service, source, "open", "mission manuelle demandée")1133 return {"ok": True, "incident_id": iid}113411351136@app.post("/api/admin/pause")1137async def admin_pause(req: Request, x_ka_token: str | None = Header(default=None)) -> dict[str, Any]:1138 check_token(x_ka_token)1139 body = await req.json()1140 LATEST["paused"] = bool(body.get("paused", True))1141 return {"ok": True, "paused": LATEST["paused"]}114211431144@app.get("/health")1145def health() -> dict[str, Any]:1146 return {"ok": True, "agent": AGENT, "polled": LATEST["polled"], "paused": LATEST["paused"]}114711481149@app.get("/favicon.svg")1150def favicon() -> FileResponse:1151 return FileResponse(BASE / "web" / f"favicon-{AGENT}.svg", media_type="image/svg+xml")115211531154# Safari (macOS/iOS) ne charge pas les favicons SVG → fallback PNG/ICO obligatoire,1155# sinon l'onglet affiche un « ? ». ICO multi-tailles + PNG 32 + apple-touch 180.1156@app.get("/favicon.ico")1157def favicon_ico() -> FileResponse:1158 return FileResponse(BASE / "web" / f"favicon-{AGENT}.ico", media_type="image/x-icon")115911601161@app.get("/favicon.png")1162def favicon_png() -> FileResponse:1163 return FileResponse(BASE / "web" / f"favicon-{AGENT}-32.png", media_type="image/png")116411651166@app.get("/apple-touch-icon.png")1167def apple_touch() -> FileResponse:1168 return FileResponse(BASE / "web" / f"apple-touch-{AGENT}.png", media_type="image/png")116911701171@app.get("/og.png")1172def og_image() -> FileResponse:1173 """Bannière de partage (Open Graph) propre à l'agent."""1174 return FileResponse(BASE / "web" / f"og-{AGENT}.png", media_type="image/png")117511761177def _render_page(name: str) -> str:1178 """Les pages sont statiques mais servies par agent : on remplace au1179 démarrage les jetons de partage (OG) par l'identité de l'agent —1180 les robots des réseaux sociaux n'exécutent pas le JS."""1181 html = (BASE / "web" / name).read_text(encoding="utf-8")1182 for token, value in {1183 "__NUM__": AGENT[-1],1184 "__DOMAIN__": ME["domain"],1185 "__TAGLINE__": ME.get("tagline", ""),1186 }.items():1187 html = html.replace(token, value)1188 return html118911901191PAGES = {name: _render_page(name) for name in ("index.html", "commander.html")}119211931194@app.get("/")1195def index() -> HTMLResponse:1196 return HTMLResponse(PAGES["index.html"])119711981199@app.get("/commander")1200def commander() -> HTMLResponse:1201 return HTMLResponse(PAGES["commander.html"])120212031204app.mount("/static", StaticFiles(directory=BASE / "web"), name="static")120512061207if __name__ == "__main__":1208 import uvicorn1209 uvicorn.run(app, host="0.0.0.0", port=ME["port"])1210