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%
8.0 KB · 175 lines python
Raw Blame History
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