SPB Git forge

spb/ka-guardian

Public
26commits 1branches 0releases
4.5 MBsize
maindefault branch
19 days agolast push
Python 57.1% Shell 13.8% CSS 12% JavaScript 10% HTML 7.2%
68.3 KB · 1,210 lines python
Raw Blame History
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