# ============================================ # Projet : KA Guardian # Fichier : orchestrator/registry.py # Rôle : Topologie VIVANTE des services surveillés. Le registre de la # passerelle mld (M1M32:~/dispatch/registry.json) fait foi pour # l'emplacement de chaque app (nœud, IP LAN, répertoire, port, # process PM2) ; topology.json ne garde que ce qui est propre aux # gardiens (agents, politique, correspondance service → app). # Author : Simon-Pierre Boucher # Date : 2026-09-04 # ============================================ """Le registre arrive dans le spool par deux chemins indépendants : 1. poussé par mld à chaque sauvegarde (abonné `ka2` → ~/ka-guardian-spool/registry.json) ; 2. tiré toutes les 2 min par le service launchd zsh `com.ka.registry-sync` (deploy/registry-sync.sh : `ssh gitsrv cat ~/dispatch/registry.json`) — chaîne 100 % binaires Apple, car macOS 26 LNP refuse le LAN à python. Le même service sonde le runner :7791 de chaque nœud hébergeur → runners.json, ce qui permet de refuser proprement une mission vers un nœud sans runner. """ from __future__ import annotations import json import pathlib import time from typing import Any HOME = pathlib.Path.home() SPOOL = HOME / "ka-guardian-spool" REGISTRY_PATH = SPOOL / "registry.json" RUNNERS_PATH = SPOOL / "runners.json" STALE_AFTER_S = 3 * 3600 # registre plus vieux que ça (mtime) → signalé « périmé » STATE: dict[str, Any] = { "path": str(REGISTRY_PATH), "mtime": 0.0, # mtime du fichier chargé "loaded_at": 0.0, # quand on l'a lu "updated": None, # champ `updated` du registre (heure mld) "gateway": None, "n_apps": 0, "error": None, "resolved": {}, # service → {app, node_alias, ip, dir, web_port, pm2, status} "changes": [], # journal des déménagements détectés (200 derniers) } _APPS: dict[str, Any] = {} _RUNNERS: dict[str, Any] = {"mtime": 0.0, "data": {}} def _pm2_from_processes(processes: list[str], exclude: list[str]) -> list[str]: """Process PM2 « métier » : on écarte les tunnels ngrok (relancer un tunnel à chaque restart casserait le domaine) et les exclusions déclarées.""" return [p for p in processes if not p.endswith("-ngrok") and p != "ngrok" and p not in set(exclude)] def _read_registry() -> dict[str, Any] | None: try: st = REGISTRY_PATH.stat() except FileNotFoundError: STATE["error"] = f"{REGISTRY_PATH} absent — mld ne l'a pas encore poussé / registry-sync inactif" return None if st.st_mtime == STATE["mtime"] and _APPS: return None # inchangé try: data = json.loads(REGISTRY_PATH.read_text()) apps = data.get("apps") or {} if not isinstance(apps, dict) or not apps: raise ValueError("registre sans apps") except Exception as exc: # fichier en cours d'écriture, JSON tronqué… STATE["error"] = f"registre illisible: {exc}" return None STATE.update({"mtime": st.st_mtime, "loaded_at": time.time(), "updated": data.get("updated"), "gateway": data.get("gateway"), "n_apps": len(apps), "error": None}) _APPS.clear() _APPS.update(apps) return apps def refresh(topo: dict[str, Any], services: dict[str, Any], nodes: dict[str, Any]) -> list[str]: """Applique le registre aux dicts SERVICES / NODES de l'orchestrateur, EN PLACE (les autres modules gardent leurs références). Retourne les changements d'emplacement détectés (messages lisibles), vide si rien n'a bougé.""" apps = _read_registry() if apps is None: return [] changes: list[str] = [] for key, svc in topo["services"].items(): app = svc.get("registry_app") or key entry = apps.get(app) cur = services.setdefault(key, dict(svc)) if not entry or not entry.get("node"): cur["registry"] = "absent" if STATE["resolved"].get(key, {}).get("registry") != "absent": changes.append(f"{svc.get('app', key)} ({app}) absente du registre mld — " f"repli sur topology.json ({cur.get('node', '?')})") STATE["resolved"][key] = {"registry": "absent", "app": app} continue alias = entry["node"] nkey = alias.lower() prev_node = nodes.get(nkey, {}) nodes[nkey] = {"alias": alias, "lan_ip": entry.get("ip") or prev_node.get("lan_ip")} exclude = svc.get("pm2_exclude") or [] pm2 = _pm2_from_processes(entry.get("processes") or [], exclude) or list(svc.get("pm2") or []) new = { "node": nkey, "dir": entry.get("dir") or svc.get("dir"), "web_port": entry.get("port") or svc.get("web_port"), "pm2": pm2, "site": (f"https://{entry['domain']}" if entry.get("domain") else svc.get("site")), "app": svc.get("app") or entry.get("label") or app, "registry": "ok", "registry_status": entry.get("status"), "registry_updated": entry.get("updated"), } old = STATE["resolved"].get(key) if old and old.get("registry") == "ok": if old["node"] != nkey: changes.append(f"{new['app']} a déménagé : {nodes.get(old['node'], {}).get('alias', old['node'])} → {alias} " f"(registre mld {entry.get('updated', '?')})") elif old.get("dir") != new["dir"] or old.get("web_port") != new["web_port"]: changes.append(f"{new['app']} : répertoire/port changés sur {alias} " f"({old.get('dir')}:{old.get('web_port')} → {new['dir']}:{new['web_port']})") elif old.get("pm2") != pm2: changes.append(f"{new['app']} : process PM2 = {', '.join(pm2)}") elif old is None or old.get("registry") == "absent": # premier chargement : on signale seulement un écart avec topology.json if svc.get("node") and svc["node"] != nkey: changes.append(f"{new['app']} : topology.json disait {svc['node']}, le registre dit {alias} — registre appliqué") cur.update(new) STATE["resolved"][key] = {**new, "app_registry": app, "node_alias": alias, "ip": nodes[nkey]["lan_ip"]} if changes: stamp = time.strftime("%Y-%m-%d %H:%M") STATE["changes"] = (STATE["changes"] + [{"ts": time.time(), "msg": f"{stamp} {m}"} for m in changes])[-200:] return changes def app_entry(app: str) -> dict[str, Any] | None: """Entrée brute du registre pour une app (ex. l'IP du gardien lui-même).""" if not _APPS: _read_registry() return _APPS.get(app) def runner_status(ip: str | None) -> dict[str, Any] | None: """État du runner :7791 d'un nœud, sondé par registry-sync.sh (runners.json). None = inconnu (jamais sondé / fichier absent / sonde trop vieille).""" if not ip: return None try: st = RUNNERS_PATH.stat() except FileNotFoundError: return None if st.st_mtime != _RUNNERS["mtime"]: try: _RUNNERS["data"] = json.loads(RUNNERS_PATH.read_text()) _RUNNERS["mtime"] = st.st_mtime except Exception: return None row = (_RUNNERS["data"] or {}).get(ip) if not row or time.time() - float(row.get("ts", 0)) > 900: return None return row def info() -> dict[str, Any]: """Bloc `registry` exposé par /api/state (dashboard + admin-ka).""" age = (time.time() - STATE["mtime"]) if STATE["mtime"] else None return { "path": STATE["path"], "updated": STATE["updated"], "gateway": STATE["gateway"], "n_apps": STATE["n_apps"], "loaded_at": STATE["loaded_at"], "file_age_s": age, "stale": bool(age is None or age > STALE_AFTER_S), "error": STATE["error"], "changes": STATE["changes"][-30:], "runners_probe": {ip: {k: v for k, v in row.items() if k in ("ok", "node", "busy", "ts")} for ip, row in (_RUNNERS["data"] or {}).items()}, }