spb/ka-guardian
Public
Python 57.1%
Shell 13.8%
CSS 12%
JavaScript 10%
HTML 7.2%
1# ============================================2# Projet : KA Guardian3# Fichier : orchestrator/registry.py4# Rôle : Topologie VIVANTE des services surveillés. Le registre de la5# passerelle mld (M1M32:~/dispatch/registry.json) fait foi pour6# l'emplacement de chaque app (nœud, IP LAN, répertoire, port,7# process PM2) ; topology.json ne garde que ce qui est propre aux8# gardiens (agents, politique, correspondance service → app).9# Author : Simon-Pierre Boucher10# Date : 2026-09-0411# ============================================12"""Le registre arrive dans le spool par deux chemins indépendants :13 1. poussé par mld à chaque sauvegarde (abonné `ka2` → ~/ka-guardian-spool/registry.json) ;14 2. tiré toutes les 2 min par le service launchd zsh `com.ka.registry-sync`15 (deploy/registry-sync.sh : `ssh gitsrv cat ~/dispatch/registry.json`) —16 chaîne 100 % binaires Apple, car macOS 26 LNP refuse le LAN à python.17Le même service sonde le runner :7791 de chaque nœud hébergeur → runners.json,18ce qui permet de refuser proprement une mission vers un nœud sans runner.19"""20from __future__ import annotations2122import json23import pathlib24import time25from typing import Any2627HOME = pathlib.Path.home()28SPOOL = HOME / "ka-guardian-spool"29REGISTRY_PATH = SPOOL / "registry.json"30RUNNERS_PATH = SPOOL / "runners.json"31STALE_AFTER_S = 3 * 3600 # registre plus vieux que ça (mtime) → signalé « périmé »3233STATE: dict[str, Any] = {34 "path": str(REGISTRY_PATH),35 "mtime": 0.0, # mtime du fichier chargé36 "loaded_at": 0.0, # quand on l'a lu37 "updated": None, # champ `updated` du registre (heure mld)38 "gateway": None,39 "n_apps": 0,40 "error": None,41 "resolved": {}, # service → {app, node_alias, ip, dir, web_port, pm2, status}42 "changes": [], # journal des déménagements détectés (200 derniers)43}44_APPS: dict[str, Any] = {}45_RUNNERS: dict[str, Any] = {"mtime": 0.0, "data": {}}464748def _pm2_from_processes(processes: list[str], exclude: list[str]) -> list[str]:49 """Process PM2 « métier » : on écarte les tunnels ngrok (relancer un tunnel50 à chaque restart casserait le domaine) et les exclusions déclarées."""51 return [p for p in processes52 if not p.endswith("-ngrok") and p != "ngrok" and p not in set(exclude)]535455def _read_registry() -> dict[str, Any] | None:56 try:57 st = REGISTRY_PATH.stat()58 except FileNotFoundError:59 STATE["error"] = f"{REGISTRY_PATH} absent — mld ne l'a pas encore poussé / registry-sync inactif"60 return None61 if st.st_mtime == STATE["mtime"] and _APPS:62 return None # inchangé63 try:64 data = json.loads(REGISTRY_PATH.read_text())65 apps = data.get("apps") or {}66 if not isinstance(apps, dict) or not apps:67 raise ValueError("registre sans apps")68 except Exception as exc: # fichier en cours d'écriture, JSON tronqué…69 STATE["error"] = f"registre illisible: {exc}"70 return None71 STATE.update({"mtime": st.st_mtime, "loaded_at": time.time(), "updated": data.get("updated"),72 "gateway": data.get("gateway"), "n_apps": len(apps), "error": None})73 _APPS.clear()74 _APPS.update(apps)75 return apps767778def refresh(topo: dict[str, Any], services: dict[str, Any], nodes: dict[str, Any]) -> list[str]:79 """Applique le registre aux dicts SERVICES / NODES de l'orchestrateur, EN PLACE80 (les autres modules gardent leurs références). Retourne les changements81 d'emplacement détectés (messages lisibles), vide si rien n'a bougé."""82 apps = _read_registry()83 if apps is None:84 return []85 changes: list[str] = []86 for key, svc in topo["services"].items():87 app = svc.get("registry_app") or key88 entry = apps.get(app)89 cur = services.setdefault(key, dict(svc))90 if not entry or not entry.get("node"):91 cur["registry"] = "absent"92 if STATE["resolved"].get(key, {}).get("registry") != "absent":93 changes.append(f"{svc.get('app', key)} ({app}) absente du registre mld — "94 f"repli sur topology.json ({cur.get('node', '?')})")95 STATE["resolved"][key] = {"registry": "absent", "app": app}96 continue97 alias = entry["node"]98 nkey = alias.lower()99 prev_node = nodes.get(nkey, {})100 nodes[nkey] = {"alias": alias, "lan_ip": entry.get("ip") or prev_node.get("lan_ip")}101 exclude = svc.get("pm2_exclude") or []102 pm2 = _pm2_from_processes(entry.get("processes") or [], exclude) or list(svc.get("pm2") or [])103 new = {104 "node": nkey,105 "dir": entry.get("dir") or svc.get("dir"),106 "web_port": entry.get("port") or svc.get("web_port"),107 "pm2": pm2,108 "site": (f"https://{entry['domain']}" if entry.get("domain") else svc.get("site")),109 "app": svc.get("app") or entry.get("label") or app,110 "registry": "ok",111 "registry_status": entry.get("status"),112 "registry_updated": entry.get("updated"),113 }114 old = STATE["resolved"].get(key)115 if old and old.get("registry") == "ok":116 if old["node"] != nkey:117 changes.append(f"{new['app']} a déménagé : {nodes.get(old['node'], {}).get('alias', old['node'])} → {alias} "118 f"(registre mld {entry.get('updated', '?')})")119 elif old.get("dir") != new["dir"] or old.get("web_port") != new["web_port"]:120 changes.append(f"{new['app']} : répertoire/port changés sur {alias} "121 f"({old.get('dir')}:{old.get('web_port')} → {new['dir']}:{new['web_port']})")122 elif old.get("pm2") != pm2:123 changes.append(f"{new['app']} : process PM2 = {', '.join(pm2)}")124 elif old is None or old.get("registry") == "absent":125 # premier chargement : on signale seulement un écart avec topology.json126 if svc.get("node") and svc["node"] != nkey:127 changes.append(f"{new['app']} : topology.json disait {svc['node']}, le registre dit {alias} — registre appliqué")128 cur.update(new)129 STATE["resolved"][key] = {**new, "app_registry": app, "node_alias": alias, "ip": nodes[nkey]["lan_ip"]}130 if changes:131 stamp = time.strftime("%Y-%m-%d %H:%M")132 STATE["changes"] = (STATE["changes"] + [{"ts": time.time(), "msg": f"{stamp} {m}"} for m in changes])[-200:]133 return changes134135136def app_entry(app: str) -> dict[str, Any] | None:137 """Entrée brute du registre pour une app (ex. l'IP du gardien lui-même)."""138 if not _APPS:139 _read_registry()140 return _APPS.get(app)141142143def runner_status(ip: str | None) -> dict[str, Any] | None:144 """État du runner :7791 d'un nœud, sondé par registry-sync.sh (runners.json).145 None = inconnu (jamais sondé / fichier absent / sonde trop vieille)."""146 if not ip:147 return None148 try:149 st = RUNNERS_PATH.stat()150 except FileNotFoundError:151 return None152 if st.st_mtime != _RUNNERS["mtime"]:153 try:154 _RUNNERS["data"] = json.loads(RUNNERS_PATH.read_text())155 _RUNNERS["mtime"] = st.st_mtime156 except Exception:157 return None158 row = (_RUNNERS["data"] or {}).get(ip)159 if not row or time.time() - float(row.get("ts", 0)) > 900:160 return None161 return row162163164def info() -> dict[str, Any]:165 """Bloc `registry` exposé par /api/state (dashboard + admin-ka)."""166 age = (time.time() - STATE["mtime"]) if STATE["mtime"] else None167 return {168 "path": STATE["path"], "updated": STATE["updated"], "gateway": STATE["gateway"],169 "n_apps": STATE["n_apps"], "loaded_at": STATE["loaded_at"], "file_age_s": age,170 "stale": bool(age is None or age > STALE_AFTER_S), "error": STATE["error"],171 "changes": STATE["changes"][-30:],172 "runners_probe": {ip: {k: v for k, v in row.items() if k in ("ok", "node", "busy", "ts")}173 for ip, row in (_RUNNERS["data"] or {}).items()},174 }175