SPB Git forge
0commits 0branches 0releases
0 Bsize
maindefault branch
—last push
118.9 KB · 2,583 lines python
Raw Blame History
1#!/usr/bin/env python32"""3MacLustr Agent v2 — collecteur de métriques du cluster + santé des apps,4servi en HTTP/JSON pour les apps iOS et macOS MacLustr.56Tourne sur M4M64a sous launchd (io.maclustr.agentd, port 9210).7Fan-out SSH vers les autres nœuds via le LAN 192.168.2.x (le SSH inter-nœuds8par Tailscale est bloqué par ACL). Les nœuds sont identifiés par le fichier9marqueur ~/.maclustr-node (les LocalHostName ne correspondent pas aux alias et10les IP LAN sont en DHCP → redécouverte par balayage du /24).1112v2 (2026-09-04) — suivi des apps déployées :13  * lit le registre maclustr-dispatch (M1M32:~/dispatch/registry.json), poussé14    par `mld` (abonné) et re-tiré périodiquement par SSH ;15  * toutes les APPS_INTERVAL s : HTTP local (ip:port) + public (https://domaine),16    état des processus PM2 (`pm2 jlist`) et launchd (`launchctl list`) sur chaque17    nœud ; état synthétique up / degraded / down ; historique SQLite 7 j ;18    journal d'événements (transitions apps + nœuds) ;19  * endpoints /api/apps, /api/apps/<app>[/logs|/history|/action], /api/events,20    /api/registry, flux SSE /api/stream (tick à chaque cycle).2122Zéro dépendance : stdlib uniquement (http.server, sqlite3, subprocess, threading,23urllib). Python 3.9 système obligatoire (/usr/bin/python3, Local Network Privacy).24"""2526import json27import os28import re29import shlex30import socket31import sqlite332import ssl33import subprocess34import threading35import time36import urllib.error37import urllib.request38from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer39from urllib.parse import urlparse, parse_qs4041AGENT_VERSION = "3.1.2"   # 3.1.2 : serveurs OVH BHS64b/BHS128/BHS128b/R9128 retirés du MacLustr (2026-10-02, seule la passerelle BHS64 reste) ; 3.1.1 : sonde « sortie Internet du LAN » (/api/egress, incident tunnel:giga-hub:egress-filter) — le Bell Giga Hub 2.0 ne relaie par moments que TCP 80/443 + UDP 53 ; 3.1.0 : iPhone/iPad = nœuds de plein droit (kind mobile dans /api/cluster, historique, incidents, proxy /api/mobile/*) ; 3.0.2 : chemins de secours SSH vers OVH via WireGuard/rebond (port 22 sortant bloqué par le FAI) ; 3.0.1 : sondes publiques de tous les sites routés par le tunnel (/api/sites, incidents site:*) ; 3.0.0 : serveurs Linux OVH, centre d'incidents, ops mld, /api/summary ; 2.2.0 : nœuds mobiles ; 2.1.0 : MacLustr Tunnel4243HOME = os.path.expanduser("~")44BASE_DIR = os.path.join(HOME, "maclustr-agentd")45CONF_PATH = os.path.join(BASE_DIR, "config.json")46LANMAP_PATH = os.path.join(BASE_DIR, "lanmap.json")47REGISTRY_PATH = os.path.join(BASE_DIR, "registry.json")48DB_PATH = os.path.join(BASE_DIR, "history.db")4950PORT = 921051COLLECT_INTERVAL = 20          # secondes entre deux cycles métriques52APPS_INTERVAL = 30             # secondes entre deux cycles santé des apps53REGISTRY_PULL_INTERVAL = 180   # re-tirage du registre depuis la passerelle54RAW_RETENTION_H = 26           # heures d'historique brut conservées (métriques)55APP_RETENTION_D = 7            # jours d'historique des vérifications d'apps56EVENTS_RETENTION_D = 30        # jours d'événements conservés57LAN_PREFIX = "192.168.2."58GATEWAY_NODE = "M1M32"59GATEWAY_REGISTRY = "~/dispatch/registry.json"60HTTP_TIMEOUT_LOCAL = 661HTTP_TIMEOUT_PUBLIC = 106263# MacLustr Tunnel (2026-09-10) : passerelles publiques OVH (WireGuard hub + Caddy) qui remplacent ngrok.64# L'agent interroge `sudo tunnelctl json` sur chacune (clé maclustr-agentd autorisée pour ubuntu).65TUNNEL_GATEWAYS = {66    "BHS64": {"host": "51.161.112.61", "user": "ubuntu", "primary": True,  "site": "Beauharnois (Québec)"},67}68TUNNEL_INTERVAL = 60           # secondes entre deux relevés des passerelles69TUNNEL_PEER_FRESH_S = 180      # handshake WireGuard plus vieux → pair considéré hors ligne70DNS_CACHE_S = 3007172# ---------------------------------------------------------------------------73# Inventaire (specs fixes ; lan_ip = graine, corrigée en live)74# ---------------------------------------------------------------------------7576NODES = [77    # name, hostname, chip, model, tier, gen, cores, mem_mb, gpu, lan_seed78    ("M3U96a", "M3U96a.maclustr.io", "M3 Ultra", "Mac Studio",  "ultra", "m3", 32, 98304, 80, "192.168.2.166"),79    ("M3U96b", "M3U96b.maclustr.io", "M3 Ultra", "Mac Studio",  "ultra", "m3", 32, 98304, 80, "192.168.2.125"),80    ("M2U64",  "M2U64.maclustr.io",  "M2 Ultra", "Mac Studio",  "ultra", "m2", 24, 65536, 76, ""),81    ("M4M64a", "M4M64a.maclustr.io", "M4 Max",   "Mac Studio",  "max",   "m4", 16, 65536, 40, "192.168.2.127"),82    ("M4M64b", "M4M64b.maclustr.io", "M4 Max",   "Mac Studio",  "max",   "m4", 16, 65536, 40, "192.168.2.128"),83    # Mac Studio M4 Max ajouté le 2026-09-25 (LAN + Tailscale 100.125.225.117 ; DNS GoDaddy A → IP Tailscale ; S/N JX0PHFQHF1)84    ("M4M64c", "M4M64c.maclustr.io", "M4 Max",   "Mac Studio",  "max",   "m4", 16, 65536, 40, "192.168.2.103"),85    # Mac Studio M5 Max + Mac mini M6 ajoutés le 2026-09-26 (LAN + Tailscale ; DNS GoDaddy A → IP Tailscale)86    ("M5M36",  "M5M36.maclustr.io",  "M5 Max",   "Mac Studio",  "max",   "m5", 18, 36864, 40, "192.168.2.104"),87    ("M4BP48", "M4BP48.maclustr.io", "M4 Max",   "MacBook Pro", "max",   "m4", 16, 49152, 40, "192.168.2.137"),88    ("M4BP36", "M4BP36.maclustr.io", "M4 Max",   "MacBook Pro", "max",   "m4", 14, 36864, 32, "192.168.2.133"),89    ("M4M36",  "M4M36.maclustr.io",  "M4 Max",   "Mac Studio",  "max",   "m4", 14, 36864, 40, "192.168.2.131"),90    ("M2M32",  "M2M32.maclustr.io",  "M2 Max",   "Mac Studio",  "max",   "m2", 12, 32768, 38, "192.168.2.90"),91    ("M2M32b", "M2M32b.maclustr.io", "M2 Max",   "Mac Studio",  "max",   "m2", 12, 32768, 38, "192.168.2.130"),92    ("M2M32c", "M2M32c.maclustr.io", "M2 Max",   "Mac Studio",  "max",   "m2", 12, 32768, 30, "192.168.2.126"),93    ("m4mc",   "m4mc.maclustr.io",   "M4 Pro",   "Mac mini",    "pro",   "m4", 12, 24576, 16, "192.168.2.134"),94    ("M1M32",  "M1M32.maclustr.io",  "M1 Max",   "Mac Studio",  "max",   "m1", 10, 32768, 32, "192.168.2.132"),95    ("m4ma",   "m4ma.maclustr.io",   "M4",       "Mac mini",    "base",  "m4", 10, 24576, 10, "192.168.2.73"),96    ("m4mb",   "m4mb.maclustr.io",   "M4",       "Mac mini",    "base",  "m4", 10, 16384, 10, "192.168.2.136"),97    ("m4md",   "m4md.maclustr.io",   "M4",       "Mac mini",    "base",  "m4", 10, 16384, 10, "192.168.2.108"),98    # 3 Mac mini locaux ajoutés le 2026-09-21 (LAN + Tailscale ; DNS GoDaddy A → IP Tailscale)99    ("m4ml",   "m4ml.maclustr.io",   "M4",       "Mac mini",    "base",  "m4", 10, 16384, 10, "192.168.2.63"),100    ("m4mm",   "m4mm.maclustr.io",   "M4 Pro",   "Mac mini",    "pro",   "m4", 12, 49152, 16, "192.168.2.92"),101    ("m4mn",   "m4mn.maclustr.io",   "M4 Pro",   "Mac mini",    "pro",   "m4", 12, 49152, 16, "192.168.2.93"),102    ("m6ma",   "m6ma.maclustr.io",   "M6",       "Mac mini",    "base",  "m6", 12, 24576, 10, "192.168.2.105"),103    ("m6mb",   "m6mb.maclustr.io",   "M6",       "Mac mini",    "base",  "m6", 12, 24576, 10, "192.168.2.106"),   # ajouté 2026-09-28104    ("m5ma",   "m5ma.maclustr.io",   "M5 Pro",   "Mac mini",    "pro",   "m5", 15, 24576, 16, "192.168.2.112"),   # ajouté 2026-09-28105    ("m2m16",  "m2m16.maclustr.io",  "M2 Pro",   "Mac mini",    "pro",   "m2", 10, 16384, 16, "192.168.2.169"),106    ("M3BA24", "M3BA24.maclustr.io", "M3",       "MacBook Air", "base",  "m3",  8, 24576, 10, "192.168.2.145"),107    ("M3BA16", "M3BA16.maclustr.io", "M3",       "MacBook Air", "base",  "m3",  8, 16384, 10, "192.168.2.70"),108    ("m2m8a",  "m2m8a.maclustr.io",  "M2",       "Mac mini",    "base",  "m2",  8,  8192, 10, "192.168.2.167"),109    ("m2m8b",  "m2m8b.maclustr.io",  "M2",       "Mac mini",    "base",  "m2",  8,  8192, 10, "192.168.2.170"),110]111112NODE_INDEX = {n[0]: i for i, n in enumerate(NODES)}113SSH_USER = "simon-pierreboucher"114115# Nœuds hors LAN : joints directement à leur IP publique, avec leur propre utilisateur SSH.116# La clé `maclustr-agentd` (~/.ssh/id_ed25519_agentd) doit être dans leur authorized_keys.117REMOTE_NODES = {118    # 2026-09-25 : les 12 Macs loués (Macly, MacStadium, rentamac) ont été retirés du cluster.119}120REMOTE_USER_BY_HOST = {r["host"]: r["user"] for r in REMOTE_NODES.values()}121122123# ---------------------------------------------------------------------------124# Serveurs dédiés Linux (OVHcloud) — 3.0.0. Même pipeline que les Macs (script125# de métriques Linux équivalent clé pour clé), utilisateur `ubuntu`, clé agentd126# autorisée. Ils ne sont jamais cherchés sur le LAN.127#   name, host public, user, site, rôle, threads, RAM Mo, matériel128# ---------------------------------------------------------------------------129SERVERS = [130    ("BHS64",   "51.161.112.61", "ubuntu", "Beauharnois (Québec)", "Passerelle MacLustr Tunnel (WireGuard + Caddy)",131     16, 65536,  "OVH ADVANCE-2 · AMD EPYC 4345P · 64 Go DDR5 ECC · 2×960 Go NVMe RAID 1"),132    # BHS64b, BHS128, BHS128b (Beauharnois) et R9128 (Gravelines) retirés du MacLustr le 2026-10-02 (apps arrêtées, serveurs résiliés).133]134# Chemin de secours quand le port 22 sortant est bloqué (FAI/routeur, 2026-09-25) : le hub WireGuard135# BHS64 reste joignable en 10.67.0.1 depuis les Macs raccordés (wg1).136WG_HUB = "10.67.0.1"137SERVER_ALT = {           # name -> liste de (host, jump) essayés après l'accès direct138    "BHS64":   [(WG_HUB, None)],139}140SERVER_META = {}141PLATFORM = {}           # name -> "macos" | "linux"142for _s in SERVERS:143    _name, _host, _user, _site, _role, _thr, _mem, _hw = _s144    _chip = _hw.split("·")[1].strip() if "·" in _hw else "x86"145    NODES.append((_name, _host, _chip, "Serveur dédié", "server", "x86", _thr, _mem, 0, _host))146    SERVER_META[_name] = {"site": _site, "role": _role, "user": _user, "hardware": _hw, "publicIP": _host}147    PLATFORM[_name] = "linux"148    REMOTE_USER_BY_HOST[_host] = _user149REMOTE_USER_BY_HOST[WG_HUB] = "ubuntu"150NODE_INDEX = {n[0]: i for i, n in enumerate(NODES)}151SERVER_NAMES = [s[0] for s in SERVERS]152MAC_NAMES = [n[0] for n in NODES if n[0] not in SERVER_META]153154155def platform_of(name):156    return PLATFORM.get(name, "macos")157158159def ssh_user_for(ip):160    """Utilisateur SSH pour une IP : celui du nœud distant si c'en est un, sinon l'utilisateur du cluster."""161    return REMOTE_USER_BY_HOST.get(ip, SSH_USER)162163NAS_DEVICES = [164    {"id": "nas-ugreen-1", "name": "UGREEN DXP4800 Plus", "host": "192.168.2.173", "mount": "nas_clustr"},165    {"id": "nas-ugreen-2", "name": "UGREEN DH4300 Plus",  "host": "192.168.2.175", "mount": "nas2_clustr"},166]167168METRICS_SCRIPT = r"""169LC_ALL=C; export LC_ALL170IFACE=$(route -n get default 2>/dev/null | awk '/interface:/{print $2}')171echo "===METRICS_START==="172echo "CPU:$(top -l 1 -n 0 2>/dev/null | grep 'CPU usage' | awk '{print 100 - $7}' | tr -d '%' || echo 0)"173echo "CPUBRK:$(top -l 1 -n 0 2>/dev/null | grep 'CPU usage' | awk '{print $3, $5, $7}' | tr -d '%' || echo '0 0 0')"174echo "MEM_PRESSURE:$(vm_stat 2>/dev/null | awk -F'[:.]' '/Pages active/{a=$2} /Pages wired down/{w=$2} /Pages occupied by compressor/{c=$2} /Pages free/{f=$2} /Pages inactive/{i=$2} /Pages speculative/{s=$2} END{u=a+w+c; t=u+f+i+s; if(t>0) print u/t*100; else print 50}')"175echo "MEMBRK:$(vm_stat 2>/dev/null | awk '/page size of/{ps=$8} /Pages wired down/{w=$4} /Pages occupied by compressor/{c=$5} END{print w*ps/1048576, c*ps/1048576}')"176echo "SWAP:$(sysctl -n vm.swapusage 2>/dev/null | awk '{print $6, $3}' | tr -d 'M' || echo '0 0')"177echo "LOAD:$(sysctl -n vm.loadavg 2>/dev/null | tr -d '{}' || echo '0 0 0')"178echo "DISK:$(df -g / 2>/dev/null | tail -1 | awk '{print $3,$2}' || echo '0 0')"179echo "DISK_EXT:$(df -g 2>/dev/null | awk '$9 ~ "^/Volumes/" {u+=$3; t+=$2} END{printf "%d %d", u+0, t+0}' || echo '0 0')"180echo "NET:$(netstat -ibn -I ${IFACE:-en0} 2>/dev/null | grep -v Link | tail -1 | awk '{print $7, $10}' || echo '0 0')"181echo "TCP:$(netstat -an -p tcp 2>/dev/null | grep -c ESTABLISHED || echo 0)"182echo "UPTIME:$(( $(date +%s) - $(sysctl -n kern.boottime 2>/dev/null | awk '{print $4}' | tr -d ',' || echo $(date +%s)) ))"183echo "PROCS:$(ps aux 2>/dev/null | wc -l | tr -d ' ')"184echo "USERS:$(who 2>/dev/null | wc -l | tr -d ' ')"185echo "OS:$(sw_vers -productVersion 2>/dev/null)"186echo "BATT:$(pmset -g batt 2>/dev/null | grep -o '[0-9]*%' | tr -d '%' | head -1 || echo -1)"187echo "BATT_STATE:$(pmset -g batt 2>/dev/null | grep -oE 'charging|discharging|charged|AC Power|finishing charge' | head -1 || echo '')"188echo "TOPPROC:$(ps -arco pcpu,comm 2>/dev/null | sed -n 2p | sed 's/^ *//')"189echo "THERMAL:$(pmset -g therm 2>/dev/null | grep 'CPU_Scheduler_Limit' | awk '{print $3}' || echo 100)"190echo "GPU:$(ioreg -r -c IOAccelerator -d 1 2>/dev/null | grep -o '"Device Utilization %"=[0-9]*' | grep -o '[0-9]*$' | sort -rn | head -1 || echo -1)"191echo "===METRICS_END==="192"""193194# Équivalent Linux (Ubuntu) du script de métriques : mêmes clés, mêmes unités.195# CPU = delta /proc/stat sur 1 s ; mémoire « pression » = (total − disponible)/total.196LINUX_METRICS_SCRIPT = r"""197LC_ALL=C; export LC_ALL198export PATH=/usr/local/bin:/usr/bin:/bin:$HOME/.npm-global/bin:$PATH199IFACE=$(ip route show default 2>/dev/null | awk '/default/{print $5; exit}')200echo "===METRICS_START==="201read u1 n1 s1 id1 io1 ir1 so1 st1 g1 < <(awk '/^cpu /{print $2,$3,$4,$5,$6,$7,$8,$9,$10}' /proc/stat); sleep 1202read u2 n2 s2 id2 io2 ir2 so2 st2 g2 < <(awk '/^cpu /{print $2,$3,$4,$5,$6,$7,$8,$9,$10}' /proc/stat)203A=$(( (u2-u1)+(n2-n1)+(s2-s1)+(ir2-ir1)+(so2-so1)+(st2-st1) )); I=$(( (id2-id1)+(io2-io1) )); DT=$((A+I)); [ $DT -le 0 ] && DT=1204echo "CPU:$(awk -v a=$A -v d=$DT 'BEGIN{printf "%.2f", a*100/d}')"205echo "CPUBRK:$(awk -v u=$(( (u2-u1)+(n2-n1) )) -v s=$(( (s2-s1)+(ir2-ir1)+(so2-so1)+(st2-st1) )) -v i=$I -v d=$DT 'BEGIN{printf "%.2f %.2f %.2f", u*100/d, s*100/d, i*100/d}')"206echo "MEM_PRESSURE:$(free -m | awk '/Mem:/{printf "%.2f", ($2-$7)*100/$2}')"207echo "MEMBRK:$(free -m | awk '/Mem:/{print $3, $6}')"208echo "SWAP:$(free -m | awk '/Swap:/{print $3, $2}')"209echo "LOAD:$(cut -d' ' -f1-3 /proc/loadavg)"210echo "DISK:$(df -BG / | tail -1 | awk '{gsub("G","",$3); gsub("G","",$2); print $3, $2}')"211echo "DISK_EXT:$(df -BG 2>/dev/null | awk '$6 ~ "^/(mnt|srv|data|opt)" && $1 ~ "^/dev/" {gsub("G","",$3); gsub("G","",$2); u+=$3; t+=$2} END{printf "%d %d", u+0, t+0}')"212echo "NET:$(awk -v i="${IFACE:-eth0}:" '$1==i{print $2, $10}' /proc/net/dev)"213echo "TCP:$(ss -tan state established 2>/dev/null | tail -n +2 | wc -l)"214echo "UPTIME:$(cut -d. -f1 /proc/uptime)"215echo "PROCS:$(ls /proc | grep -c '^[0-9]')"216echo "USERS:$(who 2>/dev/null | wc -l)"217echo "OS:$(. /etc/os-release 2>/dev/null; echo "${PRETTY_NAME:-Linux}")"218echo "BATT:-1"219echo "BATT_STATE:"220echo "TOPPROC:$(ps -eo pcpu,comm --sort=-pcpu --no-headers 2>/dev/null | head -1 | sed 's/^ *//')"221echo "THERMAL:100"222echo "GPU:-1"223echo "DOCKER:$(docker ps -q 2>/dev/null | wc -l)"224echo "PM2:$(pm2 jlist 2>/dev/null | python3 -c 'import sys,json; l=json.load(sys.stdin); print(sum(1 for p in l if p["pm2_env"]["status"]=="online"), len(l))' 2>/dev/null || echo '0 0')"225echo "KERNEL:$(uname -r)"226echo "===METRICS_END==="227"""228229TOP_SCRIPT = (230    "echo '---CPU---'; ps -arcwwwxo pid,pcpu,pmem,rss,user,comm | head -16;"231    "echo '---MEM---'; ps -amcwwwxo pid,pcpu,pmem,rss,user,comm | head -16"232)233TOP_SCRIPT_LINUX = (234    "echo '---CPU---'; ps -eo pid,pcpu,pmem,rss,user,comm --sort=-pcpu | head -16;"235    "echo '---MEM---'; ps -eo pid,pcpu,pmem,rss,user,comm --sort=-rss | head -16"236)237238PORTS_SCRIPT = (239    "lsof -nP -iTCP -sTCP:LISTEN 2>/dev/null | tail -n +2 | "240    "awk '{print $1, $2, $9}' | sort -u"241)242# même format de sortie « commande pid adresse:port » via ss (sudo sans mot de passe sur les OVH)243PORTS_SCRIPT_LINUX = (244    "sudo -n ss -ltnpH 2>/dev/null | awk '{print $4, $6}' | "245    "sed -E 's/^([^ ]+) users:\\(\\(\"([^\"]+)\",pid=([0-9]+).*/\\2 \\3 \\1/' | awk 'NF==3' | sort -u"246)247248249def metrics_script_for(name):250    return LINUX_METRICS_SCRIPT if platform_of(name) == "linux" else METRICS_SCRIPT251252# PM2 + launchd + sondes HTTP locales (127.0.0.1) d'un nœud en un seul253# aller-retour SSH. `pm2 jlist` démarre un démon PM2 s'il n'en existe pas : on254# ne l'appelle que si le God Daemon tourne. `pm2 jlist` n'émet pas de saut de255# ligne final → `echo` explicite, sinon le marqueur suivant se colle au JSON.256PROCS_SCRIPT_HEAD = r"""257export PATH=/opt/homebrew/bin:/usr/local/bin:$HOME/.npm-global/bin:$PATH258echo '---PM2---'259if pgrep -qf 'PM2.*God Daemon' 2>/dev/null; then pm2 jlist 2>/dev/null | tail -1; echo; else echo '[]'; fi260echo '---LAUNCHD---'261launchctl list 2>/dev/null262"""263264265def node_script(apps):266    """Script complet pour un nœud : processus + sonde HTTP locale de chaque app267    (curl sur 127.0.0.1 — beaucoup d'apps n'écoutent que sur loopback)."""268    parts = [PROCS_SCRIPT_HEAD]269    for e in apps:270        if not e.get("port"):271            continue272        url = "http://127.0.0.1:%s%s" % (e["port"], e.get("health_path") or "/")273        parts.append("echo '---HTTP %s---'" % e["app"])274        parts.append("curl -s -o /dev/null -m %d -w '%%{http_code} %%{time_total}\\n' %s 2>/dev/null || echo '000 0'"275                     % (HTTP_TIMEOUT_LOCAL, shlex.quote(url)))276    parts.append("echo '---END---'")277    return "\n".join(parts)278279PM2_ENV = "export PATH=/opt/homebrew/bin:/usr/local/bin:$HOME/.npm-global/bin:$PATH; "280281282def load_config():283    with open(CONF_PATH) as f:284        return json.load(f)285286287CONFIG = load_config()288TOKEN = CONFIG["token"]289SUDO_PW = CONFIG.get("sudo_password", "")290SELF_NODE = CONFIG.get("self_node", "M4M64a")291SSH_KEY = os.path.expanduser(CONFIG.get("ssh_key", "~/.ssh/id_ed25519_agentd"))292293# Nœuds mobiles (2026-09-21) : coordinateur maclustr-mobile (PM2, même hôte) — les iPhones y tirent leurs jobs.294MOBILE_URL = CONFIG.get("mobile_url", "http://127.0.0.1:9320").rstrip("/")295MOBILE_TOKEN = CONFIG.get("mobile_token", "")296MOBILE_INTERVAL = 20297mobile_latest = {"ok": False, "workers": [], "online": 0, "total": 0, "queued": 0, "active": 0, "ts": 0, "error": ""}298mobile_lock = threading.Lock()299300301def mobile_fetch():302    """Relit /health et /api/workers du coordinateur mobile (loopback, aucune contrainte LNP)."""303    try:304        req = urllib.request.Request(MOBILE_URL + "/health")305        with urllib.request.urlopen(req, timeout=4) as r:306            h = json.loads(r.read())307        workers = []308        if MOBILE_TOKEN:309            req = urllib.request.Request(MOBILE_URL + "/api/workers", headers={"Authorization": "Bearer " + MOBILE_TOKEN})310            with urllib.request.urlopen(req, timeout=4) as r:311                workers = json.loads(r.read()).get("workers", [])312        snap = {"ok": True, "version": h.get("version"), "workers": workers,313                "online": h.get("workersOnline", 0), "total": h.get("workersTotal", 0),314                "queued": h.get("queued", 0), "active": h.get("active", 0), "stats": h.get("stats"),315                "ts": time.time(), "error": ""}316    except Exception as e:317        snap = {"ok": False, "workers": mobile_latest.get("workers", []), "online": 0,318                "total": mobile_latest.get("total", 0), "queued": 0, "active": 0, "ts": time.time(), "error": str(e)[:200]}319    # 3.1.0 : chaque appareil devient un nœud (info statique + métriques au format NodeMetrics)320    snap["nodes"] = [mobile_node_info(w) for w in snap["workers"]]321    snap["metrics"] = {n["name"]: mobile_node_metrics(w, n["name"]) for w, n in zip(snap["workers"], snap["nodes"])}322    try:323        record_samples({k: v for k, v in snap["metrics"].items() if v.get("status") not in ("offline", "unknown")})324    except Exception as e:  # noqa: BLE001325        print("mobile samples:", e, flush=True)326    with mobile_lock:327        prev_online = {w.get("name") for w in mobile_latest.get("workers", []) if w.get("online")}328        mobile_latest.clear()329        mobile_latest.update(snap)330    now_online = {w.get("name") for w in snap["workers"] if w.get("online")}331    for n in sorted(now_online - prev_online):332        if prev_online or mobile_latest.get("ts"):333            record_event("mobile", n, "offline", "online", "iPhone connecté au coordinateur")334    for n in sorted(prev_online - now_online):335        record_event("mobile", n, "online", "offline", "iPhone silencieux (app fermée ou en arrière-plan)")336337338THERMAL_MAP = {"nominal": "nominal", "fair": "fair", "tiède": "fair", "serious": "serious", "chaud": "serious",339               "critical": "serious", "critique": "serious"}340341342def mobile_display_name(w):343    return w.get("displayName") or w.get("alias") or w.get("name") or "mobile"344345346def mobile_node_info(w):347    kind = w.get("kind") or ("ipad" if str(w.get("model", "")).lower().startswith("ipad") else "iphone")348    return {"name": mobile_display_name(w), "hostname": w.get("ip") or "", "chip": w.get("chip") or "",349            "model": w.get("marketing") or w.get("model") or "iPhone", "tier": "mobile", "generation": (w.get("chip") or "").lower().replace(" ", "-"),350            "cpuCores": int(w.get("cores") or 0), "memoryMB": int(w.get("ramMb") or 0), "gpuCores": 0,351            "remote": True, "user": "", "platform": w.get("platform") or ("ipados" if kind == "ipad" else "ios"),352            "kind": "mobile", "deviceKind": kind, "site": "Mobile (Tailscale)", "role": w.get("role") or "worker",353            "hardware": "%s · %s · %s Go · %s" % (w.get("marketing") or w.get("model") or "?", w.get("chip") or "?",354                                                  round((w.get("ramMb") or 0) / 1024), "iPadOS " + str(w.get("os") or "") if kind == "ipad" else "iOS " + str(w.get("os") or "")),355            "publicIP": w.get("ip") or "", "workerName": w.get("name"), "alias": w.get("alias"), "pinned": bool(w.get("pinned")),356            "appVersion": w.get("appVersion") or "", "identifier": w.get("model") or "", "bench": w.get("bench")}357358359def mobile_node_metrics(w, name):360    m = w.get("metrics") or {}361    mem_total = float(w.get("ramMb") or m.get("memTotalMb") or 0)362    mem_used = float(m.get("memUsedMb") or 0)363    avail = m.get("memAvailableMb")364    if avail is not None and mem_total:365        # mémoire « pression » = total − disponible pour l'app (iOS ne donne pas l'usage système)366        mem_used = max(mem_used, mem_total - float(avail))367    batt = m.get("battery")368    batt_pct = round(float(batt) * 100, 1) if isinstance(batt, (int, float)) and batt >= 0 else -1.0369    cpu = m.get("cpuUsage", m.get("cpuLoad"))370    cpu = float(cpu) if isinstance(cpu, (int, float)) else 0.0371    if 0 < cpu <= 1.0 and "cpuUsage" not in m:372        cpu *= 100.0373    disk_total = float(w.get("storageGb") or m.get("diskTotalGb") or 0)374    disk_free = float(m.get("diskFreeGb") or 0)375    online = bool(w.get("online"))376    status = "online" if online else "offline"377    thermal = THERMAL_MAP.get(str(m.get("thermal") or "nominal").lower(), "nominal")378    if online and (thermal == "serious" or cpu > 90):379        status = "warning"380    return {"name": name, "status": status, "ts": float(m.get("ts") or w.get("lastSeen") or 0),381            "cpu": round(cpu, 1), "cpuUser": 0.0, "cpuSystem": 0.0, "cpuIdle": round(100 - cpu, 1),382            "memUsedMB": round(mem_used, 1), "memTotalMB": mem_total, "memWiredMB": float(m.get("memUsedMb") or 0), "memCompressedMB": 0.0,383            "swapUsedMB": 0.0, "swapTotalMB": 0.0, "load": [0.0, 0.0, 0.0],384            "diskUsedGB": round(max(disk_total - disk_free, 0), 1) if disk_total else 0.0, "diskTotalGB": round(disk_total, 1),385            "extDiskUsedGB": 0.0, "extDiskTotalGB": 0.0, "netInKBs": 0.0, "netOutKBs": 0.0,386            "tcp": 0, "uptime": int(m.get("uptimeS") or 0), "procs": int(m.get("activeJobs") or 0), "users": 1 if m.get("screenOn", True) else 0,387            "os": ("iPadOS " if (w.get("kind") == "ipad") else "iOS ") + str(w.get("os") or ""),388            "battery": batt_pct, "batteryState": ("charging" if m.get("charging") else ("discharging" if batt_pct >= 0 else "")),389            "topProcess": "%d job(s) actif(s)" % int(m.get("activeJobs") or 0), "thermal": thermal, "gpu": -1.0,390            "platform": w.get("platform") or "ios", "kind": "mobile", "docker": 0, "pm2Online": 0, "pm2Total": 0, "kernel": "",391            "mobile": {"jobsDone": int(w.get("jobsDone") or 0), "jobsFailed": int(w.get("jobsFailed") or 0),392                       "activeJobs": int(m.get("activeJobs") or 0), "concurrency": int(m.get("concurrency") or 0),393                       "net": m.get("net") or "", "lowPower": bool(m.get("lowPower")), "screenOn": bool(m.get("screenOn", True)),394                       "bytesFetched": int(m.get("bytesFetched") or 0), "ageS": w.get("ageS"), "pinned": bool(w.get("pinned")),395                       "workerName": w.get("name"), "appVersion": w.get("appVersion") or "", "lastJobAt": m.get("lastJobAt") or 0,396                       "bench": w.get("bench")}}397398399def mobile_proxy(method, path, body=None, timeout=30):400    """Relais vers le coordinateur mobile avec son jeton (les apps n'ont qu'un seul jeton : celui de l'agent)."""401    if not MOBILE_TOKEN:402        return 503, {"error": "mobile_token absent de config.json"}403    data = json.dumps(body).encode() if body is not None else None404    req = urllib.request.Request(MOBILE_URL + path, data=data, method=method,405                                 headers={"Authorization": "Bearer " + MOBILE_TOKEN, "Content-Type": "application/json"})406    try:407        with urllib.request.urlopen(req, timeout=timeout) as r:408            raw = r.read()409            return r.getcode(), (json.loads(raw) if raw else {"ok": True})410    except urllib.error.HTTPError as e:411        try:412            return e.code, json.loads(e.read() or b"{}")413        except Exception:  # noqa: BLE001414            return e.code, {"error": "HTTP %d" % e.code}415    except Exception as e:  # noqa: BLE001416        return 502, {"error": "coordinateur mobile injoignable : %s" % str(e)[:160]}417418419def mobile_loop():420    while True:421        try:422            mobile_fetch()423        except Exception as e:424            print("mobile error:", e, flush=True)425        time.sleep(MOBILE_INTERVAL)426427# ---------------------------------------------------------------------------428# SSH helpers429# ---------------------------------------------------------------------------430431SSH_OPTS = [432    "-i", SSH_KEY, "-o", "BatchMode=yes", "-o", "ConnectTimeout=6",433    "-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null",434    "-o", "LogLevel=ERROR",435]436437438def run_local(script, timeout=25):439    try:440        r = subprocess.run(["/bin/sh", "-c", script], capture_output=True, text=True, timeout=timeout)441        return r.stdout442    except Exception:443        return ""444445446_last_path = {}   # name -> ("direct" | host via jump) qui a marché en dernier447448449def jump_command(jump):450    """ProxyCommand pour un rebond `user@host` avec la clé et les options agentd (-J ne les propagerait pas :451    la connexion imbriquée échouerait sur la vérification de clé d'hôte)."""452    return "ssh " + " ".join(shlex.quote(o) for o in SSH_OPTS) + " -W %h:%p " + shlex.quote(jump)453454455def run_ssh(ip, script, timeout=25, jump=None):456    """Exécute un script sh sur un nœud via son IP (LAN, ou publique pour un REMOTE_NODE). Renvoie stdout ('' si échec).457    `jump` = "user@host" → ProxyJump (rebond par le hub WireGuard quand le port 22 sortant est bloqué)."""458    try:459        opts = list(SSH_OPTS) + (["-o", "ProxyCommand=" + jump_command(jump)] if jump else [])460        r = subprocess.run(461            ["ssh"] + opts + [f"{ssh_user_for(ip)}@{ip}", "bash -s" if ssh_user_for(ip) == "ubuntu" else "sh -s"],462            input=script, capture_output=True, text=True, timeout=timeout,463        )464        if r.returncode != 0 and not r.stdout:465            print(f"ssh {ip}{' via ' + jump if jump else ''} rc={r.returncode}: {r.stderr.strip()[-200:]}", flush=True)466        return r.stdout467    except Exception as e:468        print(f"ssh {ip} exception: {e}", flush=True)469        return ""470471472def run_on_node(name, script, timeout=25):473    if name == SELF_NODE:474        return run_local(script, timeout)475    ip = lan_map.get(name)476    if not ip:477        return ""478    alts = SERVER_ALT.get(name, [])479    # le dernier chemin qui a marché passe en premier480    order = [("direct", ip, None)] + [(f"{h} via {j or 'wg'}", h, j) for h, j in alts]481    last = _last_path.get(name)482    if last and last != "direct":483        order.sort(key=lambda o: o[0] != last)484    for label, host, jump in order:485        out = run_ssh(host, script, timeout, jump)486        if out:487            if _last_path.get(name) != label:488                print(f"{name}: chemin SSH = {label}", flush=True)489                _last_path[name] = label490            return out491        if not alts:492            break493    return ""494495496# ---------------------------------------------------------------------------497# Découverte LAN (fichier marqueur ~/.maclustr-node sur chaque nœud)498# ---------------------------------------------------------------------------499500lan_map = {}          # name -> ip501lan_map_lock = threading.Lock()502_last_scan = 0.0503504505def load_lan_map():506    global lan_map507    seeds = {n[0]: n[9] for n in NODES if n[9]}508    try:509        with open(LANMAP_PATH) as f:510            saved = json.load(f)511        seeds.update({k: v for k, v in saved.items() if v})512    except Exception:513        pass514    # les nœuds hors LAN ont une IP publique fixe : elle prime sur tout lanmap.json antérieur515    seeds.update({k: r["host"] for k, r in REMOTE_NODES.items()})516    seeds.update({k: v["publicIP"] for k, v in SERVER_META.items()})   # serveurs OVH : IP publique fixe517    lan_map = seeds518519520def save_lan_map():521    try:522        with open(LANMAP_PATH, "w") as f:523            json.dump(lan_map, f, indent=1)524    except Exception:525        pass526527528def ping(ip):529    return subprocess.run(["ping", "-c", "1", "-W", "800", "-q", ip],530                          capture_output=True).returncode == 0531532533def identify(ip):534    """Renvoie l'alias maclustr du nœud à cette IP, ou None."""535    out = run_ssh(ip, "cat ~/.maclustr-node 2>/dev/null", timeout=10).strip()536    return out if out in NODE_INDEX else None537538539def rescan_lan(missing):540    """Balaye le /24 pour retrouver les nœuds dont l'IP a changé (max 1/5 min)."""541    global _last_scan542    if time.time() - _last_scan < 300 or not missing:543        return544    _last_scan = time.time()545    alive = []546    sem = threading.Semaphore(48)547    lock = threading.Lock()548549    def probe(ip):550        with sem:551            if ping(ip):552                with lock:553                    alive.append(ip)554555    threads = [threading.Thread(target=probe, args=(f"{LAN_PREFIX}{i}",)) for i in range(1, 255)]556    for t in threads:557        t.start()558    for t in threads:559        t.join()560561    missing = [n for n in missing if n not in REMOTE_NODES and n not in SERVER_META]   # hors LAN : jamais sur le /24562    healthy_ips = {lan_map[n] for n in lan_map563                   if n not in missing and lan_map.get(n)}564    candidates = [ip for ip in alive if ip not in healthy_ips]565    found = {}566    sem2 = threading.Semaphore(12)567568    def check(ip):569        with sem2:570            who = identify(ip)571        if who:572            with lock:573                found[who] = ip574575    threads = [threading.Thread(target=check, args=(ip,)) for ip in candidates]576    for t in threads:577        t.start()578    for t in threads:579        t.join()580581    if found:582        with lan_map_lock:583            lan_map.update(found)584        save_lan_map()585586587# ---------------------------------------------------------------------------588# Parsing des métriques589# ---------------------------------------------------------------------------590591def default_metrics(node, status="offline"):592    """Structure complète (toutes les clés) pour un nœud sans données."""593    name, mem_mb = node[0], node[7]594    return {595        "name": name, "status": status, "ts": time.time(),596        "cpu": 0.0, "cpuUser": 0.0, "cpuSystem": 0.0, "cpuIdle": 0.0,597        "memUsedMB": 0.0, "memTotalMB": float(mem_mb),598        "memWiredMB": 0.0, "memCompressedMB": 0.0,599        "swapUsedMB": 0.0, "swapTotalMB": 0.0,600        "load": [0.0, 0.0, 0.0],601        "diskUsedGB": 0.0, "diskTotalGB": 0.0,602        "extDiskUsedGB": 0.0, "extDiskTotalGB": 0.0,603        "netInKBs": 0.0, "netOutKBs": 0.0,604        "tcp": 0, "uptime": 0, "procs": 0, "users": 0,605        "os": "", "battery": -1.0, "batteryState": "",606        "topProcess": "", "thermal": "nominal", "gpu": -1.0,607        "platform": platform_of(name), "kind": "server" if name in SERVER_META else "mac",608        "docker": 0, "pm2Online": 0, "pm2Total": 0, "kernel": "",609        "_netBytesIn": 0.0, "_netBytesOut": 0.0,610    }611612613def parse_metrics(out, node, prev):614    name, _, chip, model, tier, gen, cores, mem_mb, gpu_cores, _ = node615    now = time.time()616    m = default_metrics(node, "online")617    m["ts"] = now618    for line in out.splitlines():619        if ":" not in line:620            continue621        key, _, val = line.partition(":")622        val = val.strip().replace(",", ".")623        try:624            if key == "CPU":625                m["cpu"] = float(val or 0)626            elif key == "CPUBRK":627                p = [float(x) for x in val.split()]628                if len(p) >= 3:629                    m["cpuUser"], m["cpuSystem"], m["cpuIdle"] = p[0], p[1], p[2]630            elif key == "MEM_PRESSURE":631                m["memUsedMB"] = mem_mb * (float(val or 50) / 100.0)632            elif key == "MEMBRK":633                p = [float(x) for x in val.split()]634                if len(p) >= 2:635                    m["memWiredMB"], m["memCompressedMB"] = p[0], p[1]636            elif key == "SWAP":637                p = [float(x) for x in val.split()]638                if len(p) >= 2:639                    m["swapUsedMB"], m["swapTotalMB"] = p[0], p[1]640            elif key == "LOAD":641                p = [float(x) for x in val.split()]642                if len(p) >= 3:643                    m["load"] = p[:3]644            elif key == "DISK":645                p = [float(x) for x in val.split()]646                if len(p) >= 2:647                    m["diskUsedGB"], m["diskTotalGB"] = p[0], p[1]648            elif key == "DISK_EXT":649                p = [float(x) for x in val.split()]650                if len(p) >= 2:651                    m["extDiskUsedGB"], m["extDiskTotalGB"] = p[0], p[1]652            elif key == "NET":653                p = [float(x) for x in val.split()]654                if len(p) >= 2:655                    m["_netBytesIn"], m["_netBytesOut"] = p[0], p[1]656                    if prev and prev.get("_netBytesIn", 0) > 0:657                        dt = now - prev.get("ts", now)658                        if dt > 0 and p[0] >= prev["_netBytesIn"] and p[1] >= prev["_netBytesOut"]:659                            m["netInKBs"] = (p[0] - prev["_netBytesIn"]) / 1024.0 / dt660                            m["netOutKBs"] = (p[1] - prev["_netBytesOut"]) / 1024.0 / dt661            elif key == "TCP":662                m["tcp"] = int(val or 0)663            elif key == "UPTIME":664                m["uptime"] = int(float(val or 0))665            elif key == "PROCS":666                m["procs"] = int(val or 0)667            elif key == "USERS":668                m["users"] = int(val or 0)669            elif key == "OS":670                m["os"] = val671            elif key == "BATT":672                m["battery"] = float(val or -1)673            elif key == "BATT_STATE":674                m["batteryState"] = val675            elif key == "TOPPROC":676                parts = val.split()677                if len(parts) >= 2:678                    m["topProcess"] = f"{' '.join(parts[1:])} ({parts[0]}%)"679                elif val:680                    m["topProcess"] = val681            elif key == "THERMAL":682                pct = int(float(val or 100))683                m["thermal"] = "nominal" if pct >= 90 else ("fair" if pct >= 70 else "serious")684            elif key == "GPU":685                m["gpu"] = float(val or -1)686            elif key == "DOCKER":687                m["docker"] = int(val or 0)688            elif key == "PM2":689                p = val.split()690                if len(p) >= 2:691                    m["pm2Online"], m["pm2Total"] = int(p[0]), int(p[1])692            elif key == "KERNEL":693                m["kernel"] = val694        except (ValueError, IndexError):695            continue696697    mem_pct = m["memUsedMB"] / mem_mb * 100 if mem_mb else 0698    if m["cpu"] > 90 or mem_pct > 95:699        m["status"] = "critical"700    elif m["cpu"] > 75 or mem_pct > 85 or m["thermal"] != "nominal":701        m["status"] = "warning"702    return m703704705# ---------------------------------------------------------------------------706# Historique SQLite (métriques, vérifications d'apps, événements)707# ---------------------------------------------------------------------------708709db_lock = threading.Lock()710711712def db():713    conn = sqlite3.connect(DB_PATH)714    conn.execute(715        "CREATE TABLE IF NOT EXISTS samples ("716        "ts INTEGER, node TEXT, cpu REAL, mem_pct REAL, gpu REAL,"717        "net_in REAL, net_out REAL, disk_pct REAL)"718    )719    conn.execute("CREATE INDEX IF NOT EXISTS idx_ts ON samples(ts)")720    conn.execute(721        "CREATE TABLE IF NOT EXISTS app_checks ("722        "ts INTEGER, app TEXT, state TEXT, local_code INTEGER, local_ms INTEGER,"723        "public_code INTEGER, public_ms INTEGER)"724    )725    conn.execute("CREATE INDEX IF NOT EXISTS idx_app_ts ON app_checks(app, ts)")726    conn.execute(727        "CREATE TABLE IF NOT EXISTS events ("728        "id INTEGER PRIMARY KEY AUTOINCREMENT, ts INTEGER, kind TEXT, subject TEXT,"729        "from_state TEXT, to_state TEXT, detail TEXT)"730    )731    conn.execute("CREATE INDEX IF NOT EXISTS idx_events_ts ON events(ts)")732    return conn733734735def record_samples(metrics):736    rows = []737    now = int(time.time())738    for m in metrics.values():739        if m["status"] in ("offline", "unknown"):740            continue741        mem_pct = m["memUsedMB"] / m["memTotalMB"] * 100 if m["memTotalMB"] else 0742        disk_pct = m["diskUsedGB"] / m["diskTotalGB"] * 100 if m["diskTotalGB"] else 0743        rows.append((now, m["name"], m["cpu"], mem_pct, m["gpu"],744                     m["netInKBs"], m["netOutKBs"], disk_pct))745    if not rows:746        return747    with db_lock:748        conn = db()749        conn.executemany("INSERT INTO samples VALUES (?,?,?,?,?,?,?,?)", rows)750        conn.execute("DELETE FROM samples WHERE ts < ?", (now - RAW_RETENTION_H * 3600,))751        conn.commit()752        conn.close()753754755def query_history(window_s, bucket_s, node=None):756    since = int(time.time()) - window_s757    where = "ts >= ?"758    args = [since]759    if node:760        where += " AND node = ?"761        args.append(node)762    with db_lock:763        conn = db()764        rows = conn.execute(765            f"SELECT (ts/{bucket_s})*{bucket_s} AS b, AVG(cpu), AVG(mem_pct),"766            f" AVG(CASE WHEN gpu >= 0 THEN gpu END), SUM(net_in)/COUNT(DISTINCT node),"767            f" SUM(net_out)/COUNT(DISTINCT node), AVG(disk_pct)"768            f" FROM samples WHERE {where} GROUP BY b ORDER BY b",769            args,770        ).fetchall()771        conn.close()772    return [773        {"ts": r[0], "cpu": round(r[1] or 0, 2), "mem": round(r[2] or 0, 2),774         "gpu": round(r[3] if r[3] is not None else -1, 2),775         "netIn": round(r[4] or 0, 1), "netOut": round(r[5] or 0, 1),776         "disk": round(r[6] or 0, 2)}777        for r in rows778    ]779780781def record_app_checks(apps):782    now = int(time.time())783    rows = []784    for a in apps.values():785        loc = a.get("local") or {}786        pub = a.get("public") or {}787        rows.append((now, a["app"], a["state"], loc.get("code", 0), loc.get("ms", 0),788                     pub.get("code", 0), pub.get("ms", 0)))789    if not rows:790        return791    with db_lock:792        conn = db()793        conn.executemany("INSERT INTO app_checks VALUES (?,?,?,?,?,?,?)", rows)794        conn.execute("DELETE FROM app_checks WHERE ts < ?", (now - APP_RETENTION_D * 86400,))795        conn.commit()796        conn.close()797798799def query_uptime(window_s=86400):800    """% de vérifications non « down » par app sur la fenêtre + dernier changement."""801    since = int(time.time()) - window_s802    with db_lock:803        conn = db()804        rows = conn.execute(805            "SELECT app, AVG(CASE WHEN state = 'down' THEN 0.0 ELSE 100.0 END), COUNT(*),"806            " AVG(CASE WHEN local_ms > 0 THEN local_ms END)"807            " FROM app_checks WHERE ts >= ? GROUP BY app", (since,)).fetchall()808        conn.close()809    return {r[0]: {"uptime": round(r[1] or 0, 2), "checks": r[2],810                   "avgLocalMs": round(r[3] or 0)} for r in rows}811812813def query_app_history(app, window_s, bucket_s):814    since = int(time.time()) - window_s815    with db_lock:816        conn = db()817        rows = conn.execute(818            f"SELECT (ts/{bucket_s})*{bucket_s} AS b,"819            " AVG(CASE WHEN state = 'down' THEN 0.0 ELSE 100.0 END),"820            " AVG(CASE WHEN local_ms > 0 THEN local_ms END),"821            " AVG(CASE WHEN public_ms > 0 THEN public_ms END),"822            " SUM(state = 'down'), SUM(state = 'degraded'), COUNT(*)"823            " FROM app_checks WHERE app = ? AND ts >= ? GROUP BY b ORDER BY b",824            (app, since)).fetchall()825        conn.close()826    out = []827    for r in rows:828        worst = "up"829        if r[4]:830            worst = "down"831        elif r[5]:832            worst = "degraded"833        out.append({"ts": r[0], "availability": round(r[1] or 0, 1),834                    "localMs": round(r[2] or 0), "publicMs": round(r[3] or 0),835                    "worst": worst, "checks": r[6]})836    return out837838839def record_event(kind, subject, from_state, to_state, detail=""):840    now = int(time.time())841    with db_lock:842        conn = db()843        conn.execute("INSERT INTO events (ts, kind, subject, from_state, to_state, detail) VALUES (?,?,?,?,?,?)",844                     (now, kind, subject, from_state or "", to_state or "", detail or ""))845        conn.execute("DELETE FROM events WHERE ts < ?", (now - EVENTS_RETENTION_D * 86400,))846        conn.commit()847        conn.close()848    print(f"event {kind} {subject}: {from_state} -> {to_state} {detail}", flush=True)849850851def query_events(since=0, limit=200, kind=None, subject=None):852    where = ["ts >= ?"]853    args = [int(since)]854    if kind:855        where.append("kind = ?")856        args.append(kind)857    if subject:858        where.append("subject = ?")859        args.append(subject)860    args.append(int(limit))861    with db_lock:862        conn = db()863        rows = conn.execute(864            "SELECT id, ts, kind, subject, from_state, to_state, detail FROM events"865            f" WHERE {' AND '.join(where)} ORDER BY ts DESC, id DESC LIMIT ?", args).fetchall()866        conn.close()867    return [{"id": r[0], "ts": r[1], "kind": r[2], "subject": r[3],868             "from": r[4], "to": r[5], "detail": r[6]} for r in rows]869870871# ---------------------------------------------------------------------------872# Bus de ticks (flux SSE)873# ---------------------------------------------------------------------------874875tick_cv = threading.Condition()876tick_seq = 0877tick_last = {}878879880def publish_tick(kind, payload):881    global tick_seq, tick_last882    with tick_cv:883        tick_seq += 1884        last = {"seq": tick_seq, "kind": kind, "ts": time.time()}885        last.update(payload)886        tick_last = last887        tick_cv.notify_all()888889890# ---------------------------------------------------------------------------891# Boucle de collecte des métriques892# ---------------------------------------------------------------------------893894latest = {}            # name -> metrics dict895latest_lock = threading.Lock()896started_at = time.time()897cycle_count = 0898899900def node_online(name):901    with latest_lock:902        m = latest.get(name)903    return bool(m) and m.get("status") not in ("offline", "unknown")904905906def collect_cycle():907    global cycle_count908    results = {}909    threads = []910    lock = threading.Lock()911912    def worker(node):913        name = node[0]914        prev = latest.get(name)915        out = run_on_node(name, metrics_script_for(name))916        if "===METRICS_START===" in out:917            m = parse_metrics(out, node, prev)918        else:919            m = dict(prev) if prev else default_metrics(node)920            m["status"] = "offline"921            m["ts"] = time.time()922        with lock:923            results[name] = m924925    for node in NODES:926        t = threading.Thread(target=worker, args=(node,))927        t.start()928        threads.append(t)929    for t in threads:930        t.join()931932    # Événements nœud : transitions en ligne ↔ hors ligne (pas au premier cycle)933    with latest_lock:934        previous = {k: v.get("status") for k, v in latest.items()}935        latest.update(results)936    if cycle_count > 0:937        for name, m in results.items():938            was = previous.get(name)939            now_off = m["status"] == "offline"940            was_off = was in ("offline", None)941            if was is not None and now_off != was_off:942                record_event("node", name, "offline" if was_off else "online",943                             "offline" if now_off else "online",944                             "ne répond plus à l'agent" if now_off else "de nouveau joignable")945    record_samples(results)946    cycle_count += 1947948    online = sum(1 for m in results.values() if m.get("status") not in ("offline", "unknown"))949    srv_online = sum(1 for n in SERVER_NAMES if results.get(n, {}).get("status") not in ("offline", "unknown", None))950    publish_tick("metrics", {"nodesOnline": online, "nodesTotal": len(NODES),951                             "macsOnline": online - srv_online, "macsTotal": len(MAC_NAMES),952                             "serversOnline": srv_online, "serversTotal": len(SERVER_NAMES)})953954    missing = [n for n, m in results.items() if m.get("status") == "offline" and n != SELF_NODE]955    if missing:956        threading.Thread(target=rescan_lan, args=(missing,), daemon=True).start()957958959def collector_loop():960    while True:961        t0 = time.time()962        try:963            collect_cycle()964        except Exception as e:965            print("collect error:", e, flush=True)966        time.sleep(max(5, COLLECT_INTERVAL - (time.time() - t0)))967968969# ---------------------------------------------------------------------------970# Registre maclustr-dispatch (apps déployées)971# ---------------------------------------------------------------------------972973registry = {"updated": None, "apps": {}, "history": []}974registry_lock = threading.Lock()975registry_meta = {"source": "", "loadedAt": 0.0, "pulledAt": 0.0, "mtime": 0.0, "error": ""}976977978def load_registry_file():979    """Recharge ~/maclustr-agentd/registry.json si le fichier a changé (push mld)."""980    global registry981    try:982        st = os.stat(REGISTRY_PATH)983    except FileNotFoundError:984        return False985    if st.st_mtime <= registry_meta["mtime"]:986        return False987    try:988        with open(REGISTRY_PATH) as f:989            data = json.load(f)990        if not isinstance(data.get("apps"), dict):991            raise ValueError("registre sans clé apps")992        with registry_lock:993            registry = data994        registry_meta.update({"mtime": st.st_mtime, "loadedAt": time.time(),995                              "source": "file", "error": ""})996        print(f"registry loaded ({len(data['apps'])} apps, updated {data.get('updated')})", flush=True)997        return True998    except Exception as e:999        registry_meta["error"] = f"registre local illisible : {e}"1000        print(registry_meta["error"], flush=True)1001        return False100210031004def pull_registry():1005    """Tire le registre depuis la passerelle (SSH LAN) et l'écrit localement."""1006    out = run_on_node(GATEWAY_NODE, f"cat {GATEWAY_REGISTRY} 2>/dev/null", timeout=20)1007    registry_meta["pulledAt"] = time.time()1008    if not out.strip():1009        registry_meta["error"] = f"registre injoignable sur {GATEWAY_NODE}"1010        return False1011    try:1012        data = json.loads(out)1013        if not isinstance(data.get("apps"), dict):1014            raise ValueError("clé apps absente")1015    except Exception as e:1016        registry_meta["error"] = f"registre passerelle invalide : {e}"1017        return False1018    tmp = REGISTRY_PATH + ".tmp"1019    with open(tmp, "w") as f:1020        json.dump(data, f, indent=1, ensure_ascii=False)1021    os.replace(tmp, REGISTRY_PATH)1022    loaded = load_registry_file()1023    if loaded:1024        registry_meta["source"] = "pull"1025    return loaded102610271028def registry_apps():1029    with registry_lock:1030        return {k: dict(v) for k, v in registry.get("apps", {}).items()}103110321033# ---------------------------------------------------------------------------1034# Santé des apps1035# ---------------------------------------------------------------------------10361037class _NoRedirect(urllib.request.HTTPRedirectHandler):1038    def redirect_request(self, req, fp, code, msg, headers, newurl):1039        return None104010411042_ssl_ctx = ssl.create_default_context()1043_opener_local = urllib.request.build_opener(_NoRedirect)1044_opener_public = urllib.request.build_opener(_NoRedirect, urllib.request.HTTPSHandler(context=_ssl_ctx))104510461047def http_probe(url, timeout, opener):1048    """→ {"ok", "code", "ms", "error"} ; 2xx/3xx = ok, 4xx = répond (dégradé), 5xx/erreur = KO."""1049    t0 = time.time()1050    req = urllib.request.Request(url, headers={"User-Agent": f"maclustr-agentd/{AGENT_VERSION}"})1051    try:1052        with opener.open(req, timeout=timeout) as resp:1053            code = resp.getcode()1054            resp.read(2048)1055    except urllib.error.HTTPError as e:1056        code = e.code1057    except (urllib.error.URLError, socket.timeout, ssl.SSLError, ConnectionError, OSError) as e:1058        return {"ok": False, "code": 0, "ms": int((time.time() - t0) * 1000),1059                "error": str(getattr(e, "reason", e))[:120]}1060    except Exception as e:1061        return {"ok": False, "code": 0, "ms": int((time.time() - t0) * 1000), "error": str(e)[:120]}1062    ms = int((time.time() - t0) * 1000)1063    return {"ok": 200 <= code < 400, "code": code, "ms": ms, "error": ""}106410651066def parse_procs_output(out):1067    """Sortie de node_script → (liste pm2 normalisée, dict launchd label→{pid,exit},1068    dict app→sonde HTTP locale)."""1069    pm2, launchd, http = [], {}, {}1070    section = None1071    pm2_raw = []1072    http_app = None1073    for line in out.splitlines():1074        if line.startswith("---PM2---"):1075            section = "pm2"1076            continue1077        if line.startswith("---LAUNCHD---"):1078            section = "launchd"1079            continue1080        m = re.match(r"---HTTP (\S+)---$", line)1081        if m:1082            section = "http"1083            http_app = m.group(1)1084            continue1085        if line.startswith("---END---"):1086            break1087        if section == "pm2":1088            # sécurité : marqueur collé au JSON si le saut de ligne manque1089            if "---LAUNCHD---" in line:1090                pm2_raw.append(line.split("---LAUNCHD---")[0])1091                section = "launchd"1092            else:1093                pm2_raw.append(line)1094        elif section == "launchd":1095            p = line.split("\t")1096            if len(p) >= 3 and p[2] != "Label":1097                pid = int(p[0]) if p[0].isdigit() else None1098                try:1099                    exit_code = int(p[1])1100                except ValueError:1101                    exit_code = None1102                launchd[p[2]] = {"pid": pid, "exit": exit_code}1103        elif section == "http" and http_app:1104            p = line.split()1105            if len(p) >= 2:1106                try:1107                    code = int(p[0])1108                    ms = int(float(p[1].replace(",", ".")) * 1000)1109                except ValueError:1110                    code, ms = 0, 01111                http[http_app] = {"ok": 200 <= code < 400, "code": code, "ms": ms,1112                                  "error": "" if code else "connexion refusée ou délai dépassé (127.0.0.1)"}1113                http_app = None1114    raw = "\n".join(pm2_raw).strip()1115    start = raw.find("[")1116    if start >= 0:1117        try:1118            for p in json.loads(raw[start:]):1119                env = p.get("pm2_env", {}) or {}1120                mon = p.get("monit", {}) or {}1121                uptime_ms = env.get("pm_uptime") or 01122                pm2.append({1123                    "name": p.get("name", ""),1124                    "pmId": p.get("pm_id", -1),1125                    "status": env.get("status", "unknown"),1126                    "restarts": env.get("restart_time", 0) or 0,1127                    "uptimeS": int(max(0, time.time() - uptime_ms / 1000.0)) if uptime_ms and env.get("status") == "online" else 0,1128                    "cpu": float(mon.get("cpu") or 0),1129                    "memMB": round(float(mon.get("memory") or 0) / 1048576.0, 1),1130                    "cron": bool(env.get("cron_restart")),1131                    "autorestart": env.get("autorestart", True) is not False,1132                    "outLog": env.get("pm_out_log_path", ""),1133                    "errLog": env.get("pm_err_log_path", ""),1134                    "pid": p.get("pid") or 0,1135                })1136        except ValueError:1137            pass1138    return pm2, launchd, http113911401141# ---------------------------------------------------------------------------1142# MacLustr Tunnel : état des passerelles (pairs WireGuard, routes Caddy) + DNS des domaines1143# ---------------------------------------------------------------------------11441145tunnel_latest = {}     # gateway -> record1146tunnel_lock = threading.Lock()1147_dns_cache = {}        # domain -> (ip, ts)114811491150def tunnel_fetch(name, gw):1151    """Relevé d'une passerelle via `sudo tunnelctl json` (SSH direct, clé agentd)."""1152    rec = {"name": name, "ip": gw["host"], "primary": gw.get("primary", False), "site": gw.get("site", ""),1153           "ok": False, "error": None, "checkedAt": time.time(), "wg": {"listenPort": 0, "peers": []},1154           "routes": [], "caddyActive": False}1155    try:1156        attempts = [(gw["host"], None)] + SERVER_ALT.get(name, [])1157        r = None1158        for host, jump in attempts:1159            r = subprocess.run(["ssh"] + SSH_OPTS + (["-o", "ProxyCommand=" + jump_command(jump)] if jump else []) + [f"{gw['user']}@{host}", "sudo tunnelctl json"],1160                               capture_output=True, text=True, timeout=25)1161            if r.returncode == 0 and r.stdout.strip():1162                break1163        if r is None or r.returncode != 0 or not r.stdout.strip():1164            rec["error"] = ((r.stderr.strip() if r else "") or "aucune sortie")[-160:]1165            return rec1166        d = json.loads(r.stdout)1167        for p in d.get("wg", {}).get("peers", []):1168            hs = p.get("handshakeS")1169            p["online"] = hs is not None and hs <= TUNNEL_PEER_FRESH_S1170        rec["wg"] = d.get("wg", rec["wg"])1171        rec["routes"] = d.get("caddy", {}).get("routes", [])1172        rec["caddyActive"] = bool(d.get("caddy", {}).get("active"))1173        rec["ok"] = rec["caddyActive"]1174    except Exception as e:  # noqa: BLE0011175        rec["error"] = str(e)[-160:]1176    return rec117711781179def tunnel_cycle_run():1180    results = {}1181    lock = threading.Lock()11821183    def worker(name, gw):1184        r = tunnel_fetch(name, gw)1185        with lock:1186            results[name] = r11871188    threads = [threading.Thread(target=worker, args=(n, g)) for n, g in TUNNEL_GATEWAYS.items()]1189    for t in threads:1190        t.start()1191    for t in threads:1192        t.join()1193    with tunnel_lock:1194        tunnel_latest.clear()1195        tunnel_latest.update(results)119611971198def tunnel_loop():1199    while True:1200        try:1201            tunnel_cycle_run()1202        except Exception as e:  # noqa: BLE0011203            print(f"tunnel cycle error: {e}", flush=True)1204        time.sleep(TUNNEL_INTERVAL)120512061207def dns_a(domain):1208    """Première adresse A du domaine (cache DNS_CACHE_S) ; None si non résolu."""1209    if not domain:1210        return None1211    now = time.time()1212    hit = _dns_cache.get(domain)1213    if hit and now - hit[1] < DNS_CACHE_S:1214        return hit[0]1215    ip = None1216    try:1217        infos = socket.getaddrinfo(domain, 443, socket.AF_INET, socket.SOCK_STREAM)1218        ip = infos[0][4][0] if infos else None1219    except Exception:  # noqa: BLE0011220        ip = None1221    _dns_cache[domain] = (ip, now)1222    return ip122312241225def tunnel_info(domain):1226    """Pour une app : passerelles qui routent son domaine, upstreams, résolution DNS et « via »."""1227    if not domain:1228        return None1229    gws, ups = [], []1230    with tunnel_lock:1231        snap = {k: dict(v) for k, v in tunnel_latest.items()}1232    for name, rec in snap.items():1233        for r in rec.get("routes", []):1234            if r.get("domain") == domain and r.get("kind") == "proxy":1235                gws.append(name)1236                ups.extend(u.get("addr") if isinstance(u, dict) else str(u) for u in r.get("upstreams", []))1237    ip = dns_a(domain)1238    via = None1239    if ip:1240        via = next((n for n, g in TUNNEL_GATEWAYS.items() if g["host"] == ip), None) or "hors tunnel"1241    gateway = via if via in gws else (gws[0] if gws else None)1242    return {"gateways": sorted(set(gws)), "upstreams": sorted(set(ups)), "dns": ip, "via": via,1243            "gateway": gateway, "site": TUNNEL_GATEWAYS.get(gateway or "", {}).get("site")}124412451246def tunnel_snapshot():1247    with tunnel_lock:1248        gws = [dict(v) for v in tunnel_latest.values()]1249    gws.sort(key=lambda g: (not g.get("primary"), g["name"]))1250    peers_by_node = {}1251    for g in gws:1252        for p in g.get("wg", {}).get("peers", []):1253            peers_by_node.setdefault(p.get("alias"), {})[g["name"]] = {1254                "ip": p.get("ip"), "handshakeS": p.get("handshakeS"), "online": p.get("online", False)}1255    return {"gateways": gws, "peersByNode": peers_by_node,1256            "routesTotal": sum(1 for g in gws for r in g.get("routes", []) if r.get("kind") == "proxy"),1257            "ts": time.time()}125812591260# ---------------------------------------------------------------------------1261# Sites publics (3.0.1) : chaque domaine routé par la passerelle primaire est1262# sondé en HTTPS — y compris les apps hors registre mld (serveurs OVH). Un site1263# est « down » après SITES_FAIL_THRESHOLD échecs consécutifs (5xx ou injoignable).1264# ---------------------------------------------------------------------------12651266SITES_INTERVAL = 1201267SITES_FAIL_THRESHOLD = 21268sites_latest = {}      # domain -> record1269sites_lock = threading.Lock()1270sites_cycle = 0127112721273def sites_cycle_run():1274    global sites_cycle1275    with tunnel_lock:1276        snap = {k: dict(v) for k, v in tunnel_latest.items()}1277    app_domains = {e.get("domain"): a for a, e in registry_apps().items() if e.get("domain")}1278    targets = {}1279    for gname, g in sorted(snap.items(), key=lambda kv: not kv[1].get("primary")):1280        for r in g.get("routes", []):1281            d = r.get("domain")1282            if r.get("kind") != "proxy" or not d or d in targets:1283                continue1284            targets[d] = {"domain": d, "gateway": gname,1285                          "upstreams": [u.get("addr") if isinstance(u, dict) else str(u) for u in r.get("upstreams", [])],1286                          "app": app_domains.get(d)}1287    if not targets:1288        return1289    results = {}1290    lock = threading.Lock()1291    sem = threading.Semaphore(8)12921293    def probe(t):1294        with sem:1295            r = http_probe("https://%s/" % t["domain"], HTTP_TIMEOUT_PUBLIC, _opener_public)1296        with sites_lock:1297            prev = sites_latest.get(t["domain"], {})1298        down = r["code"] == 0 or r["code"] >= 5001299        rec = dict(t)1300        rec.update({"ok": not down, "responds": r["code"] > 0, "code": r["code"], "ms": r["ms"], "error": r["error"],1301                    "fails": (prev.get("fails", 0) + 1) if down else 0,1302                    "checkedAt": time.time(), "since": prev.get("since") if prev.get("ok") == (not down) else time.time()})1303        with lock:1304            results[t["domain"]] = rec13051306    threads = [threading.Thread(target=probe, args=(t,)) for t in targets.values()]1307    for th in threads:1308        th.start()1309    for th in threads:1310        th.join()1311    with sites_lock:1312        prev_down = {d for d, r in sites_latest.items() if r.get("fails", 0) >= SITES_FAIL_THRESHOLD}1313        sites_latest.clear()1314        sites_latest.update(results)1315        now_down = {d for d, r in results.items() if r.get("fails", 0) >= SITES_FAIL_THRESHOLD}1316    if sites_cycle > 0:1317        for d in sorted(now_down - prev_down):1318            record_event("site", d, "up", "down", "HTTP %s — %s" % (results[d]["code"], results[d]["error"] or "5xx"))1319        for d in sorted(prev_down - now_down):1320            record_event("site", d, "down", "up", "de nouveau en ligne (HTTP %s)" % results[d]["code"])1321    sites_cycle += 1132213231324def sites_loop():1325    time.sleep(20)   # laisser le premier relevé tunnel arriver1326    while True:1327        try:1328            sites_cycle_run()1329        except Exception as e:  # noqa: BLE0011330            print("sites error:", e, flush=True)1331        time.sleep(SITES_INTERVAL)133213331334# ---- Sortie Internet du LAN (filtre du routeur Bell) -----------------------------------------------------1335# 2026-10-01 : le Bell Giga Hub 2.0 (Sagemcom 5697, firmware 3.11.3) se met par moments à ne relayer que1336# TCP 80/443 et UDP 53 vers Internet : ICMP, SSH 22, NTP, STUN, tunnels… sont silencieusement perdus alors1337# que son propre ping passe (outil Utilitaires). Un redémarrage du Hub rétablit tout pour un temps.1338# Effets vus : IPTV « Tunnel Error », NAS UGREEN voyant « sans Internet », SSH vers OVH impossible,1339# WireGuard seulement via UDP 443. On sonde donc des ports non-web et on ouvre un incident quand le témoin1340# HTTPS passe mais que la majorité des autres échouent.1341EGRESS_INTERVAL = 1201342EGRESS_FAIL_THRESHOLD = 2          # cycles consécutifs avant incident (≈ 4 min)1343EGRESS_CONTROL = ("tcp", "51.161.112.61", 443, "HTTPS BHS64 (témoin)")1344EGRESS_PROBES = [1345    ("tcp",  "51.161.112.61", 22,  "SSH BHS64"),1346    ("tcp",  "1.1.1.1",       853, "DNS/TLS Cloudflare"),1347    ("icmp", "1.1.1.1",       0,   "ping Cloudflare"),1348    ("icmp", "8.8.8.8",       0,   "ping Google"),1349    ("ntp",  "time.apple.com", 123, "NTP Apple"),1350]1351egress_latest = {}1352egress_lock = threading.Lock()1353egress_cycle = 0135413551356def _egress_probe(kind, host, port, timeout=3.0):1357    t0 = time.time()1358    try:1359        if kind == "tcp":1360            s = socket.create_connection((host, port), timeout=timeout)1361            s.close()1362        elif kind == "ntp":1363            s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)1364            s.settimeout(timeout)1365            s.sendto(b"\x1b" + 47 * b"\0", (host, port))1366            s.recvfrom(48)1367            s.close()1368        elif kind == "icmp":1369            r = subprocess.run(["/sbin/ping", "-c", "1", "-W", str(int(timeout * 1000)), host],1370                               stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, timeout=timeout + 2)1371            if r.returncode != 0:1372                return False, 01373        else:1374            return False, 01375        return True, int((time.time() - t0) * 1000)1376    except Exception:  # noqa: BLE0011377        return False, 0137813791380def egress_cycle_run():1381    global egress_cycle1382    control_ok, control_ms = _egress_probe(*EGRESS_CONTROL[:3])1383    results = []1384    for kind, host, port, label in EGRESS_PROBES:1385        ok, ms = _egress_probe(kind, host, port)1386        results.append({"kind": kind, "host": host, "port": port, "label": label, "ok": ok, "ms": ms})1387    blocked = [r["label"] for r in results if not r["ok"]]1388    filtered = bool(control_ok) and len(blocked) >= 41389    with egress_lock:1390        prev = dict(egress_latest)1391        was = prev.get("filtered", False)1392        rec = {"ok": not filtered, "filtered": filtered, "controlOk": control_ok, "controlMs": control_ms,1393               "probes": results, "blocked": blocked,1394               "fails": (prev.get("fails", 0) + 1) if filtered else 0,1395               "checkedAt": time.time(),1396               "since": prev.get("since") if (prev and was == filtered) else time.time()}1397        egress_latest.clear()1398        egress_latest.update(rec)1399    if egress_cycle > 0 and was != filtered:1400        if filtered:1401            record_event("tunnel", "giga-hub", "open", "filtered",1402                         "le routeur Bell ne relaie plus que le web — bloqués : " + ", ".join(blocked))1403        else:1404            record_event("tunnel", "giga-hub", "filtered", "open", "les ports non-web sortent de nouveau du LAN")1405    egress_cycle += 1140614071408def egress_loop():1409    time.sleep(12)1410    while True:1411        try:1412            egress_cycle_run()1413        except Exception as e:  # noqa: BLE0011414            print("egress error:", e, flush=True)1415        time.sleep(EGRESS_INTERVAL)141614171418def egress_snapshot():1419    with egress_lock:1420        out = dict(egress_latest)1421    out["cycle"] = egress_cycle1422    out["ts"] = time.time()1423    return out142414251426def egress_public():1427    with egress_lock:1428        e = dict(egress_latest)1429    return {"filtered": e.get("filtered", False), "controlOk": e.get("controlOk"),1430            "blocked": e.get("blocked", []), "since": e.get("since"), "checkedAt": e.get("checkedAt")}143114321433def sites_snapshot():1434    with sites_lock:1435        sites = [dict(v) for v in sites_latest.values()]1436    sites.sort(key=lambda r: (r.get("ok", True), r["domain"]))1437    return {"sites": sites, "total": len(sites), "ok": sum(1 for r in sites if r.get("ok")),1438            "down": sum(1 for r in sites if r.get("fails", 0) >= SITES_FAIL_THRESHOLD), "cycle": sites_cycle, "ts": time.time()}143914401441apps_latest = {}       # app -> record1442apps_lock = threading.Lock()1443apps_cycle = 01444node_procs_cache = {}  # node -> {"pm2": [...], "launchd": {...}, "ts": float}144514461447def evaluate_app(entry, pm2_list, launchd_map, node_ok, local, public, tun=None):1448    """Fusionne registre + processus + sondes HTTP (+ tunnel) en un enregistrement d'état."""1449    app = entry["app"]1450    declared = list(entry.get("processes") or [])1451    labels = list(entry.get("launchd") or [])1452    by_name = {p["name"]: p for p in (pm2_list or [])}1453    procs = []1454    for name in declared:1455        p = by_name.get(name)1456        if p:1457            procs.append(dict(p))1458        else:1459            procs.append({"name": name, "pmId": -1, "status": "absent" if pm2_list is not None else "unknown",1460                          "restarts": 0, "uptimeS": 0, "cpu": 0.0, "memMB": 0.0, "cron": False,1461                          "autorestart": True, "outLog": "", "errLog": "", "pid": 0})1462    # processus PM2 non déclarés mais préfixés du nom de l'app (workers ajoutés à la main)1463    for name, p in by_name.items():1464        if name not in declared and (name == app or name.startswith(app + "-") or name.startswith(app.replace("-", "") + "-")):1465            q = dict(p)1466            q["undeclared"] = True1467            procs.append(q)1468    ld = []1469    for label in labels:1470        info = launchd_map.get(label) if launchd_map is not None else None1471        ld.append({"label": label,1472                   "pid": (info or {}).get("pid"),1473                   "exit": (info or {}).get("exit"),1474                   "running": bool(info and info.get("pid")),1475                   "known": info is not None})14761477    reasons = []1478    if not node_ok:1479        state = "down"1480        reasons.append("nœud hors ligne")1481    elif entry.get("port") and local is not None and not local.get("ok"):1482        state = "down"1483        reasons.append("HTTP local %s" % (local.get("code") or local.get("error") or "KO"))1484    else:1485        state = "up"1486        bad = [p for p in procs if not p.get("cron") and p.get("autorestart", True)1487               and p["status"] not in ("online", "launching", "unknown") and not p.get("undeclared")]1488        if bad:1489            state = "degraded"1490            reasons.append("PM2 %s" % ", ".join("%s=%s" % (p["name"], p["status"]) for p in bad[:3]))1491        errored = [p for p in procs if p["status"] == "errored"]1492        if errored and state == "up":1493            state = "degraded"1494        dead_ld = [l for l in ld if l["known"] and not l["running"] and not l["label"].endswith("caffeinate")]1495        if dead_ld:1496            state = "degraded"1497            reasons.append("launchd %s" % ", ".join(l["label"] for l in dead_ld[:3]))1498        if entry.get("domain") and public is not None and not public.get("ok"):1499            state = "degraded"1500            reasons.append("public %s" % (public.get("code") or public.get("error") or "KO"))1501        if tun and tun.get("gateways") and tun.get("dns") and tun.get("via") == "hors tunnel":1502            state = "degraded"1503            reasons.append("DNS %s ne pointe pas vers le MacLustr Tunnel" % tun["dns"])1504        if entry.get("domain") and tun is not None and not tun.get("gateways") and tunnel_latest:1505            reasons.append("aucune route tunnel pour ce domaine")1506        flappy = [p for p in procs if not p.get("cron") and p.get("restarts", 0) >= 20 and p.get("uptimeS", 0) < 600]1507        if flappy:1508            state = "degraded"1509            reasons.append("redémarrages en boucle : %s" % ", ".join(p["name"] for p in flappy[:2]))1510        if not entry.get("port") and not procs and not ld:1511            state = "unknown"1512            reasons.append("rien à sonder")15131514    rec = {1515        "app": app,1516        "label": entry.get("label") or app,1517        "node": entry.get("node"),1518        "ip": entry.get("ip"),1519        "port": entry.get("port"),1520        "domain": entry.get("domain"),1521        "dir": entry.get("dir"),1522        "healthPath": entry.get("health_path") or "/",1523        "deployed": entry.get("deployed"),1524        "registryStatus": entry.get("status"),1525        "registryHealth": entry.get("health"),1526        "processes": procs,1527        "launchd": ld,1528        "local": local,1529        "public": public,1530        "tunnel": tun,1531        "state": state,1532        "reason": " · ".join(reasons),1533        "checkedAt": time.time(),1534        "memMB": round(sum(p.get("memMB", 0) for p in procs), 1),1535        "cpu": round(sum(p.get("cpu", 0) for p in procs), 1),1536    }1537    return rec153815391540def apps_cycle_run():1541    global apps_cycle1542    load_registry_file()1543    if time.time() - registry_meta["pulledAt"] > REGISTRY_PULL_INTERVAL and (1544            not registry_apps() or time.time() - registry_meta["loadedAt"] > REGISTRY_PULL_INTERVAL):1545        pull_registry()1546    entries = registry_apps()1547    for app, e in entries.items():1548        e["app"] = app1549    if not entries:1550        publish_tick("apps", {"appsTotal": 0, "appsUp": 0})1551        return15521553    # 1. par nœud hébergeant ≥ 1 app : un seul SSH → PM2 + launchd + curl local1554    by_node = {}1555    for e in entries.values():1556        if e.get("node"):1557            by_node.setdefault(e["node"], []).append(e)1558    procs = {}1559    lock = threading.Lock()15601561    def procs_worker(node, node_apps):1562        if node not in NODE_INDEX or not node_online(node):1563            with lock:1564                procs[node] = None1565            return1566        out = run_on_node(node, node_script(node_apps), timeout=45)1567        if "---END---" not in out:1568            with lock:1569                procs[node] = None1570            return1571        pm2, ld, http = parse_procs_output(out)1572        with lock:1573            procs[node] = {"pm2": pm2, "launchd": ld, "http": http, "ts": time.time()}1574            node_procs_cache[node] = procs[node]15751576    threads = [threading.Thread(target=procs_worker, args=(n, a)) for n, a in by_node.items()]1577    for t in threads:1578        t.start()15791580    # 2. sondes HTTP publiques en parallèle (pendant les SSH)1581    probes = {}1582    sem = threading.Semaphore(24)15831584    def probe_worker(app, e):1585        pub = None1586        if e.get("domain"):1587            with sem:1588                pub = http_probe("https://%s%s" % (e["domain"], e.get("health_path") or "/"),1589                                 HTTP_TIMEOUT_PUBLIC, _opener_public)1590        with lock:1591            probes[app] = pub15921593    pthreads = [threading.Thread(target=probe_worker, args=(a, e)) for a, e in entries.items()]1594    for t in pthreads:1595        t.start()1596    for t in threads + pthreads:1597        t.join()15981599    # 3. évaluation + événements1600    results = {}1601    for app, e in entries.items():1602        node = e.get("node")1603        pinfo = procs.get(node)1604        pub = probes.get(app)1605        loc = None1606        if e.get("port"):1607            if pinfo and app in pinfo.get("http", {}):1608                loc = pinfo["http"][app]1609            elif node_online(node):1610                # repli : sonde directe par le LAN (nœud joignable mais script en échec)1611                host = "127.0.0.1" if node == SELF_NODE else (e.get("ip") or lan_map.get(node or ""))1612                loc = http_probe("http://%s:%s%s" % (host, e["port"], e.get("health_path") or "/"),1613                                 HTTP_TIMEOUT_LOCAL, _opener_local) if host else \1614                    {"ok": False, "code": 0, "ms": 0, "error": "IP inconnue"}1615            else:1616                loc = {"ok": False, "code": 0, "ms": 0, "error": "nœud hors ligne"}1617        results[app] = evaluate_app(1618            e,1619            pinfo["pm2"] if pinfo else None,1620            pinfo["launchd"] if pinfo else None,1621            node_online(node),1622            loc, pub,1623            tunnel_info(e.get("domain")),1624        )16251626    uptimes = query_uptime(86400)1627    with apps_lock:1628        previous = {k: v.get("state") for k, v in apps_latest.items()}1629        apps_latest.clear()1630        apps_latest.update(results)1631        for app, rec in apps_latest.items():1632            u = uptimes.get(app)1633            rec["uptime24h"] = u["uptime"] if u else None1634            rec["avgLocalMs"] = u["avgLocalMs"] if u else None1635    if apps_cycle > 0:1636        for app, rec in results.items():1637            was = previous.get(app)1638            if was is not None and was != rec["state"]:1639                record_event("app", app, was, rec["state"], rec.get("reason") or "")1640        for app in previous:1641            if app not in results:1642                record_event("app", app, previous[app], "removed", "retirée du registre")1643        for app in results:1644            if app not in previous:1645                record_event("app", app, "", results[app]["state"], "nouvelle app dans le registre (%s)" % results[app].get("node"))1646    record_app_checks(results)1647    apps_cycle += 11648    s = apps_summary()1649    publish_tick("apps", {"appsTotal": s["total"], "appsUp": s["up"], "appsDown": s["down"],1650                          "appsDegraded": s["degraded"]})165116521653def apps_loop():1654    # laisse le premier cycle métriques établir qui est en ligne1655    time.sleep(8)1656    while True:1657        t0 = time.time()1658        try:1659            apps_cycle_run()1660        except Exception as e:1661            print("apps cycle error:", repr(e), flush=True)1662        time.sleep(max(5, APPS_INTERVAL - (time.time() - t0)))166316641665def apps_summary():1666    with apps_lock:1667        states = [a["state"] for a in apps_latest.values()]1668    return {"total": len(states), "up": states.count("up"), "degraded": states.count("degraded"),1669            "down": states.count("down"), "unknown": states.count("unknown"),1670            "cycle": apps_cycle}167116721673def apps_by_node():1674    out = {}1675    with apps_lock:1676        for a in apps_latest.values():1677            n = a.get("node") or "?"1678            d = out.setdefault(n, {"total": 0, "up": 0, "degraded": 0, "down": 0, "unknown": 0, "apps": []})1679            d["total"] += 11680            d[a["state"]] = d.get(a["state"], 0) + 11681            d["apps"].append(a["app"])1682    return out168316841685def app_logs(app, lines=120, process=None):1686    with apps_lock:1687        rec = apps_latest.get(app)1688    if not rec:1689        return None1690    node = rec.get("node")1691    lines = max(10, min(int(lines), 600))1692    parts = []1693    targets = [p for p in rec.get("processes", []) if not process or p["name"] == process]1694    for p in targets:1695        for kind, path in (("out", p.get("outLog")), ("err", p.get("errLog"))):1696            if path:1697                parts.append((p["name"], kind, path))1698    for l in rec.get("launchd", []):1699        if process and l["label"] != process:1700            continue1701        parts.append((l["label"], "launchd", "__LAUNCHD__" + l["label"]))1702    if not parts:1703        return {"app": app, "node": node, "logs": []}1704    script = ["export PATH=/opt/homebrew/bin:/usr/local/bin:$PATH"]1705    for name, kind, path in parts:1706        marker = "===LOG %s %s===" % (name, kind)1707        script.append("echo %s" % shlex.quote(marker))1708        if path.startswith("__LAUNCHD__"):1709            label = path[len("__LAUNCHD__"):]1710            script.append(1711                "for f in $(launchctl print gui/$(id -u)/%s 2>/dev/null | awk '/stdout path|stderr path/{print $NF}' | sort -u); do "1712                "echo \"# $f\"; tail -n %d \"$f\" 2>/dev/null; done" % (shlex.quote(label), lines))1713        else:1714            script.append("tail -n %d %s 2>/dev/null" % (lines, shlex.quote(path)))1715    script.append("echo '===LOG END==='")1716    out = run_on_node(node, "\n".join(script), timeout=30)1717    logs = []1718    current = None1719    for line in out.splitlines():1720        m = re.match(r"===LOG (\S+) (\S+)===$", line)1721        if m:1722            current = {"process": m.group(1), "kind": m.group(2), "text": []}1723            logs.append(current)1724            continue1725        if line.startswith("===LOG END==="):1726            break1727        if current is not None:1728            current["text"].append(line)1729    for l in logs:1730        l["text"] = "\n".join(l["text"])[-60000:]1731    return {"app": app, "node": node, "logs": logs, "ts": time.time()}173217331734def app_action(app, action, process=None):1735    with apps_lock:1736        rec = apps_latest.get(app)1737    if not rec:1738        return {"ok": False, "error": "app inconnue"}1739    node = rec.get("node")1740    if not node_online(node):1741        return {"ok": False, "error": "nœud %s hors ligne" % node}1742    pm2_names = [p["name"] for p in rec.get("processes", []) if p.get("pmId", -1) >= 0 or p["status"] != "absent"]1743    labels = [l["label"] for l in rec.get("launchd", [])]1744    if process:1745        if process in pm2_names:1746            pm2_names, labels = [process], []1747        elif process in labels:1748            pm2_names, labels = [], [process]1749        else:1750            return {"ok": False, "error": "processus %s inconnu pour %s" % (process, app)}1751    if action not in ("restart", "stop", "start", "reload"):1752        return {"ok": False, "error": "action inconnue"}1753    cmds = []1754    if pm2_names:1755        cmds.append("pm2 %s %s 2>&1 | grep -vE '^\\s*$' | tail -20" % (action, " ".join(shlex.quote(n) for n in pm2_names)))1756    if labels:1757        if action in ("restart", "reload", "start"):1758            for l in labels:1759                cmds.append("launchctl kickstart -k gui/$(id -u)/%s 2>&1 && echo 'launchd %s: kickstart ok'" % (shlex.quote(l), l))1760        else:1761            for l in labels:1762                cmds.append("launchctl kill TERM gui/$(id -u)/%s 2>&1 && echo 'launchd %s: TERM envoyé (KeepAlive le relancera)'" % (shlex.quote(l), l))1763    if not cmds:1764        return {"ok": False, "error": "aucun processus à piloter"}1765    script = PM2_ENV + "\n" + "\n".join(cmds) + "\necho ACTION_DONE"1766    out = run_on_node(node, script, timeout=60)1767    ok = "ACTION_DONE" in out1768    record_event("action", app, "", action, "%s%s par l'app MacLustr" % (action, " " + process if process else ""))1769    threading.Thread(target=_recheck_soon, daemon=True).start()1770    return {"ok": ok, "output": out.replace("ACTION_DONE", "").strip()[-1500:], "node": node}177117721773def _recheck_soon():1774    time.sleep(4)1775    try:1776        apps_cycle_run()1777    except Exception as e:1778        print("recheck error:", e, flush=True)177917801781# ---------------------------------------------------------------------------1782# Actions (Node Doctor)1783# ---------------------------------------------------------------------------17841785def do_action(name, action, pid=None):1786    pw = shlex.quote(SUDO_PW)1787    if platform_of(name) == "linux":1788        # serveurs OVH : sudo sans mot de passe1789        if action == "kill" and pid:1790            script = f"kill -TERM {int(pid)} 2>/dev/null; sleep 2; kill -0 {int(pid)} 2>/dev/null && kill -KILL {int(pid)}; echo DONE"1791        elif action == "purge":1792            script = "sync; echo 3 | sudo -n tee /proc/sys/vm/drop_caches >/dev/null 2>&1 && echo DONE"1793        elif action == "caches":1794            script = "sudo -n apt-get clean >/dev/null 2>&1; sudo -n journalctl --vacuum-time=7d >/dev/null 2>&1; docker system prune -f >/dev/null 2>&1; echo DONE"1795        elif action == "reboot":1796            script = "(sleep 1; sudo -n shutdown -r now) >/dev/null 2>&1 & echo DONE"1797        else:1798            return {"ok": False, "error": "unknown action"}1799        out = run_on_node(name, script, timeout=60)1800        record_event("action", name, "", action, "%s par l'app MacLustr" % action)1801        return {"ok": "DONE" in out, "output": out.strip()[-500:]}1802    if action == "kill" and pid:1803        script = f"kill -TERM {int(pid)} 2>/dev/null; sleep 2; kill -0 {int(pid)} 2>/dev/null && kill -KILL {int(pid)}; echo DONE"1804    elif action == "purge":1805        script = f"printf '%s' {pw} | sudo -S purge 2>/dev/null && echo DONE"1806    elif action == "caches":1807        script = "rm -rf ~/Library/Caches/* /tmp/*.tmp 2>/dev/null; echo DONE"1808    elif action == "reboot":1809        script = f"printf '%s' {pw} | sudo -S shutdown -r now 2>/dev/null & echo DONE"1810    else:1811        return {"ok": False, "error": "unknown action"}1812    out = run_on_node(name, script, timeout=30)1813    record_event("action", name, "", action, "%s par l'app MacLustr" % action)1814    return {"ok": "DONE" in out, "output": out.strip()[-500:]}181518161817# ---------------------------------------------------------------------------1818# NAS1819# ---------------------------------------------------------------------------18201821def nas_status():1822    out = []1823    for nas in NAS_DEVICES:1824        online = ping(nas["host"])1825        used = total = avail = 0.01826        mp = os.path.join(HOME, nas["mount"])1827        if os.path.ismount(mp):1828            df = run_local(f"df -g {shlex.quote(mp)} | tail -1")1829            p = df.split()1830            if len(p) >= 4:1831                try:1832                    total, used, avail = float(p[1]), float(p[2]), float(p[3])1833                except ValueError:1834                    pass1835        out.append({**{k: nas[k] for k in ("id", "name", "host")},1836                    "online": online, "usedGB": used, "totalGB": total, "availGB": avail})1837    return out183818391840# ---------------------------------------------------------------------------1841# HTTP1842# ---------------------------------------------------------------------------18431844NODE_INFO = [1845    {"name": n[0], "hostname": n[1], "chip": n[2], "model": n[3], "tier": n[4],1846     "generation": n[5], "cpuCores": n[6], "memoryMB": n[7], "gpuCores": n[8],1847     "remote": n[0] in REMOTE_NODES or n[0] in SERVER_META,1848     "user": SERVER_META.get(n[0], {}).get("user") or ssh_user_for(REMOTE_NODES.get(n[0], {}).get("host", "")),1849     "platform": platform_of(n[0]), "kind": "server" if n[0] in SERVER_META else "mac",1850     "site": SERVER_META.get(n[0], {}).get("site", "Saint-Augustin-de-Desmaures (Québec)"),1851     "role": SERVER_META.get(n[0], {}).get("role", ""),1852     "hardware": SERVER_META.get(n[0], {}).get("hardware", ""),1853     "publicIP": SERVER_META.get(n[0], {}).get("publicIP", "")}1854    for n in NODES1855]185618571858# ---------------------------------------------------------------------------1859# Centre d'incidents (3.0.0) — vue consolidée « qu'est-ce qui ne va pas ? »1860# calculée toutes les INCIDENTS_INTERVAL s à partir des nœuds, serveurs, apps,1861# passerelles tunnel et coordinateur mobile. Accusés de réception persistés.1862# ---------------------------------------------------------------------------18631864INCIDENTS_INTERVAL = 151865incidents_lock = threading.Lock()1866incidents_latest = []          # liste triée1867incidents_counts = {"critical": 0, "warning": 0, "info": 0, "total": 0, "acked": 0}1868incident_first_seen = {}       # id -> ts1869incident_acks = {}             # id -> {"ts": float, "note": str}1870incidents_cycle = 01871SEV_ORDER = {"critical": 0, "warning": 1, "info": 2}187218731874def _acks_table(conn):1875    conn.execute("CREATE TABLE IF NOT EXISTS incident_acks (id TEXT PRIMARY KEY, ts INTEGER, note TEXT)")187618771878def load_acks():1879    with db_lock:1880        conn = db()1881        _acks_table(conn)1882        rows = conn.execute("SELECT id, ts, note FROM incident_acks").fetchall()1883        conn.close()1884    for r in rows:1885        incident_acks[r[0]] = {"ts": r[1], "note": r[2] or ""}188618871888def save_ack(iid, note=None, remove=False):1889    with db_lock:1890        conn = db()1891        _acks_table(conn)1892        if remove:1893            conn.execute("DELETE FROM incident_acks WHERE id = ?", (iid,))1894        else:1895            conn.execute("INSERT OR REPLACE INTO incident_acks (id, ts, note) VALUES (?,?,?)",1896                         (iid, int(time.time()), note or ""))1897        conn.commit()1898        conn.close()189919001901def compute_incidents():1902    """Construit la liste des incidents ouverts (sans état : recalcul complet)."""1903    now = time.time()1904    out = []19051906    def add(sev, kind, subject, code, title, detail="", actions=None, meta=None):1907        out.append({"id": f"{kind}:{subject}:{code}", "severity": sev, "kind": kind, "subject": subject,1908                    "code": code, "title": title, "detail": detail, "actions": actions or [], "meta": meta or {}})19091910    with latest_lock:1911        snap = {k: dict(v) for k, v in latest.items()}1912    for name, m in snap.items():1913        kind = "server" if name in SERVER_META else "node"1914        st = m.get("status")1915        if st in ("offline", "unknown"):1916            add("critical", kind, name, "offline", f"{name} hors ligne", "ne répond plus à l'agent",1917                ["mld-heal"] if kind == "node" else [], {"site": SERVER_META.get(name, {}).get("site")})1918            continue1919        mem_pct = m["memUsedMB"] / m["memTotalMB"] * 100 if m.get("memTotalMB") else 01920        if st == "critical":1921            add("warning", kind, name, "load", f"{name} saturé",1922                f"CPU {m.get('cpu', 0):.0f} % · mémoire {mem_pct:.0f} %", ["top", "purge"],1923                {"cpu": round(m.get("cpu", 0), 1), "mem": round(mem_pct, 1)})1924        if m.get("diskTotalGB"):1925            dpct = m["diskUsedGB"] / m["diskTotalGB"] * 1001926            free_gb = m["diskTotalGB"] - m["diskUsedGB"]1927            if dpct >= 95:1928                add("critical", kind, name, "disk-full", f"Disque presque plein sur {name}",1929                    f"{dpct:.0f} % utilisés · {free_gb:.0f} Go libres", ["caches"], {"diskPct": round(dpct, 1), "freeGB": round(free_gb)})1930            elif dpct >= 88:1931                add("warning", kind, name, "disk-high", f"Disque chargé sur {name}",1932                    f"{dpct:.0f} % utilisés · {free_gb:.0f} Go libres", ["caches"], {"diskPct": round(dpct, 1), "freeGB": round(free_gb)})1933        if m.get("thermal") == "serious":1934            add("warning", kind, name, "thermal", f"{name} bride son CPU (thermique)", "CPU_Scheduler_Limit < 70 %")1935        if m.get("swapTotalMB", 0) > 0 and m["swapUsedMB"] / m["swapTotalMB"] > 0.9 and mem_pct > 85:1936            add("warning", kind, name, "swap", f"{name} swappe", f"swap {m['swapUsedMB']:.0f}/{m['swapTotalMB']:.0f} Mo · mémoire {mem_pct:.0f} %")1937        batt = m.get("battery", -1)1938        if batt is not None and 0 <= batt < 20 and (m.get("batteryState") or "").startswith("discharging"):1939            add("warning", kind, name, "battery", f"{name} sur batterie ({batt:.0f} %)", "portable débranché")19401941    with apps_lock:1942        apps = [dict(a) for a in apps_latest.values()]1943    for a in apps:1944        if a.get("state") == "down":1945            add("critical", "app", a["app"], "down", f"{a.get('label') or a['app']} est hors service",1946                a.get("reason") or "sonde HTTP locale en échec", ["restart", "logs"],1947                {"node": a.get("node"), "domain": a.get("domain"), "port": a.get("port")})1948        elif a.get("state") == "degraded":1949            add("warning", "app", a["app"], "degraded", f"{a.get('label') or a['app']} est dégradée",1950                a.get("reason") or "", ["restart", "logs"], {"node": a.get("node"), "domain": a.get("domain"), "port": a.get("port")})19511952    with tunnel_lock:1953        gws = {k: dict(v) for k, v in tunnel_latest.items()}1954    for name, g in gws.items():1955        if not g.get("ok"):1956            add("critical" if g.get("primary") else "warning", "tunnel", name, "gateway",1957                f"Passerelle {name} injoignable", g.get("error") or "Caddy inactif", [], {"site": g.get("site")})1958    primary = next((g for g in gws.values() if g.get("primary")), None)1959    if primary and primary.get("ok"):1960        hosting = set()1961        for a in apps:1962            if a.get("domain") and a.get("node"):1963                hosting.add(a["node"])1964        for p in primary.get("wg", {}).get("peers", []):1965            alias = p.get("alias")1966            if alias in hosting and not p.get("online") and snap.get(alias, {}).get("status") not in ("offline", "unknown", None):1967                add("warning", "tunnel", alias, "wg-offline", f"Tunnel wg1 de {alias} silencieux",1968                    "dernier handshake il y a %s s — les sites publics de ce nœud ne répondent plus" % p.get("handshakeS"), ["mld-heal"])19691970    with sites_lock:1971        sites = [dict(v) for v in sites_latest.values()]1972    for s_ in sites:1973        if s_.get("fails", 0) >= SITES_FAIL_THRESHOLD and not s_.get("app"):1974            add("critical", "site", s_["domain"], "down", f"{s_['domain']} ne répond plus",1975                "HTTP %s via %s → %s%s" % (s_.get("code"), s_.get("gateway"), ", ".join(s_.get("upstreams") or []),1976                                           (" — " + s_["error"]) if s_.get("error") else ""),1977                ["open"], {"domain": s_["domain"], "site": s_.get("gateway")})1978    with mobile_lock:1979        mob = dict(mobile_latest)1980        mob_metrics = dict(mobile_latest.get("metrics") or {})1981    for name, m in mob_metrics.items():1982        mm = m.get("mobile") or {}1983        if m.get("status") in ("offline", "unknown"):1984            if mm.get("pinned"):1985                add("warning", "mobile", name, "offline", f"{name} (mobile épinglé) hors ligne",1986                    "app MacLustr fermée ou en arrière-plan depuis %s s" % int(mm.get("ageS") or 0), ["open"], {"site": "mobile"})1987            continue1988        if 0 <= m.get("battery", -1) < 20 and (m.get("batteryState") == "discharging"):1989            add("warning", "mobile", name, "battery", f"{name} : batterie faible ({m['battery']:.0f} %)", "appareil non branché", [], {})1990        if m.get("thermal") == "serious":1991            add("warning", "mobile", name, "thermal", f"{name} chauffe", "état thermique sérieux — jobs ralentis", [], {})1992    if mob.get("ts") and not mob.get("ok"):1993        add("info", "mobile", "coordinateur", "down", "Coordinateur mobile injoignable", mob.get("error") or "", ["restart"])1994    eg = egress_snapshot()1995    if eg.get("filtered") and eg.get("fails", 0) >= EGRESS_FAIL_THRESHOLD:1996        add("critical", "tunnel", "giga-hub", "egress-filter", "Routeur Bell : seuls les ports web sortent du LAN",1997            "bloqués : %s — HTTPS passe. Effets : IPTV « Tunnel Error », NAS UGREEN sans Internet, SSH vers OVH coupé "1998            "(le tunnel UDP 443 tient). Remède connu : redémarrer le Giga Hub (192.168.2.1 → Réinitialisation)."1999            % ", ".join(eg.get("blocked") or []), ["open"], {"site": "LAN", "blocked": eg.get("blocked") or []})2000    if registry_meta.get("error"):2001        add("info", "registry", "mld", "pull", "Registre mld : dernier tirage en échec", registry_meta.get("error", "")[:160], ["registry-refresh"])20022003    active = {i["id"] for i in out}2004    for i in out:2005        i["since"] = incident_first_seen.setdefault(i["id"], now)2006        i["ageS"] = int(now - i["since"])2007        ack = incident_acks.get(i["id"])2008        i["acked"] = bool(ack)2009        i["ackTs"] = ack["ts"] if ack else None2010        i["ackNote"] = ack["note"] if ack else ""2011    for k in list(incident_first_seen):2012        if k not in active:2013            incident_first_seen.pop(k, None)2014            if k in incident_acks:2015                incident_acks.pop(k, None)2016                save_ack(k, remove=True)2017    out.sort(key=lambda i: (i["acked"], SEV_ORDER.get(i["severity"], 9), i["since"]))2018    return out201920202021def incidents_cycle_run():2022    global incidents_cycle, incidents_counts2023    new = compute_incidents()2024    with incidents_lock:2025        old_ids = {i["id"]: i for i in incidents_latest}2026    new_ids = {i["id"]: i for i in new}2027    if incidents_cycle > 0:2028        for iid, i in new_ids.items():2029            if iid not in old_ids:2030                record_event("incident", iid, "", "open", "%s — %s" % (i["title"], i["detail"]))2031        for iid, i in old_ids.items():2032            if iid not in new_ids:2033                record_event("incident", iid, "open", "closed", "%s résolu" % i["title"])2034    counts = {"critical": 0, "warning": 0, "info": 0, "total": len(new), "acked": 0}2035    for i in new:2036        counts[i["severity"]] = counts.get(i["severity"], 0) + 12037        if i["acked"]:2038            counts["acked"] += 12039    changed = counts != incidents_counts or set(new_ids) != set(old_ids)2040    with incidents_lock:2041        incidents_latest[:] = new2042        incidents_counts = counts2043    incidents_cycle += 12044    if changed:2045        publish_tick("incidents", {"incidents": counts})204620472048def incidents_loop():2049    time.sleep(8)2050    while True:2051        try:2052            incidents_cycle_run()2053        except Exception as e:  # noqa: BLE0012054            print("incidents error:", e, flush=True)2055        time.sleep(INCIDENTS_INTERVAL)205620572058def incidents_snapshot():2059    with incidents_lock:2060        return {"incidents": [dict(i) for i in incidents_latest], "counts": dict(incidents_counts),2061                "cycle": incidents_cycle, "ts": time.time()}206220632064# ---------------------------------------------------------------------------2065# Ops (3.0.0) : commandes mld exécutées sur la passerelle M1M32 (liste blanche)2066# ---------------------------------------------------------------------------20672068MLD_COMMANDS = {2069    "status":        {"args": ["status"],            "label": "État des apps (registre)",         "mutating": False},2070    "status-live":   {"args": ["status", "--live"],  "label": "État des apps (sondes live)",      "mutating": False},2071    "nodes":         {"args": ["nodes"],             "label": "Ressources des nœuds",             "mutating": False},2072    "scan":          {"args": ["scan"],              "label": "Scan live des nœuds",              "mutating": False},2073    "discover":      {"args": ["discover"],          "label": "Découverte LAN",                   "mutating": False},2074    "apps":          {"args": ["apps"],              "label": "Manifestes",                       "mutating": False},2075    "plan":          {"args": ["plan"],              "label": "Plan de placement",                "mutating": False},2076    "tunnel-status": {"args": ["tunnel", "status"],  "label": "État du MacLustr Tunnel",          "mutating": False},2077    "heal-dry":      {"args": ["heal", "--dry-run"], "label": "Auto-réparation (simulation)",     "mutating": False},2078    "heal":          {"args": ["heal"],              "label": "Auto-réparation",                  "mutating": True},2079}2080_ANSI = re.compile(r"\x1b\[[0-9;]*[A-Za-z]")208120822083def run_mld(key):2084    spec = MLD_COMMANDS.get(key)2085    if not spec:2086        return {"ok": False, "error": "commande inconnue", "allowed": sorted(MLD_COMMANDS)}2087    ip = lan_map.get(GATEWAY_NODE)2088    if not ip:2089        return {"ok": False, "error": "passerelle %s introuvable" % GATEWAY_NODE}2090    t0 = time.time()2091    cmd = "~/maclustr-dispatch/bin/mld " + " ".join(shlex.quote(a) for a in spec["args"])2092    try:2093        r = subprocess.run(["ssh"] + SSH_OPTS + [f"{SSH_USER}@{ip}", cmd],2094                           capture_output=True, text=True, timeout=240)2095        out = _ANSI.sub("", (r.stdout or "") + (("\n" + r.stderr) if r.stderr.strip() else ""))2096        ok = r.returncode == 02097    except subprocess.TimeoutExpired:2098        out, ok = "délai dépassé (240 s)", False2099    except Exception as e:  # noqa: BLE0012100        out, ok = str(e), False2101    if spec["mutating"]:2102        record_event("action", "mld", "", key, "mld %s lancé depuis l'app MacLustr" % " ".join(spec["args"]))2103        threading.Thread(target=_recheck_soon, daemon=True).start()2104    return {"ok": ok, "command": "mld " + " ".join(spec["args"]), "key": key, "output": out.strip()[-20000:],2105            "ms": int((time.time() - t0) * 1000), "gateway": GATEWAY_NODE, "ts": time.time()}210621072108# ---------------------------------------------------------------------------2109# Résumé (3.0.0) : un seul appel pour un tableau de bord / widget2110# ---------------------------------------------------------------------------21112112def cluster_aggregates(snap):2113    macs = {k: v for k, v in snap.items() if k not in SERVER_META}2114    servers = {k: v for k, v in snap.items() if k in SERVER_META}21152116    def agg(group):2117        online = [m for m in group.values() if m.get("status") not in ("offline", "unknown", None)]2118        cores = {n[0]: n[6] for n in NODES}2119        wsum = sum(cores.get(m["name"], 1) for m in online) or 12120        cpu = sum(m.get("cpu", 0) * cores.get(m["name"], 1) for m in online) / wsum if online else 02121        mem_t = sum(m.get("memTotalMB", 0) for m in online) or 12122        mem_u = sum(m.get("memUsedMB", 0) for m in online)2123        disk_t = sum(m.get("diskTotalGB", 0) for m in online)2124        disk_u = sum(m.get("diskUsedGB", 0) for m in online)2125        gpus = [m["gpu"] for m in online if m.get("gpu", -1) >= 0]2126        return {"online": len(online), "total": len(group),2127                "cpu": round(cpu, 1), "memPct": round(mem_u / mem_t * 100, 1),2128                "memUsedGB": round(mem_u / 1024, 1), "memTotalGB": round(mem_t / 1024, 1),2129                "diskUsedGB": round(disk_u), "diskTotalGB": round(disk_t),2130                "gpu": round(sum(gpus) / len(gpus), 1) if gpus else -1,2131                "netInKBs": round(sum(m.get("netInKBs", 0) for m in online), 1),2132                "netOutKBs": round(sum(m.get("netOutKBs", 0) for m in online), 1),2133                "cores": sum(cores.get(m["name"], 0) for m in online),2134                "coresTotal": sum(cores.get(n, 0) for n in group)}2135    with mobile_lock:2136        mob_metrics = dict(mobile_latest.get("metrics") or {})2137        mob_nodes = list(mobile_latest.get("nodes") or [])2138    mob_cores = {n["name"]: n.get("cpuCores") or 0 for n in mob_nodes}2139    mob_online = [m for m in mob_metrics.values() if m.get("status") not in ("offline", "unknown", None)]2140    mem_t = sum(m.get("memTotalMB", 0) for m in mob_online) or 12141    mobiles = {"online": len(mob_online), "total": len(mob_metrics),2142               "cpu": round(sum(m.get("cpu", 0) for m in mob_online) / len(mob_online), 1) if mob_online else 0,2143               "memPct": round(sum(m.get("memUsedMB", 0) for m in mob_online) / mem_t * 100, 1) if mob_online else 0,2144               "memUsedGB": round(sum(m.get("memUsedMB", 0) for m in mob_online) / 1024, 1),2145               "memTotalGB": round(sum(m.get("memTotalMB", 0) for m in mob_online) / 1024, 1),2146               "diskUsedGB": round(sum(m.get("diskUsedGB", 0) for m in mob_online)), "diskTotalGB": round(sum(m.get("diskTotalGB", 0) for m in mob_online)),2147               "gpu": -1, "netInKBs": 0, "netOutKBs": 0,2148               "cores": sum(mob_cores.get(m["name"], 0) for m in mob_online), "coresTotal": sum(mob_cores.values()),2149               "battery": round(sum(m.get("battery", 0) for m in mob_online if m.get("battery", -1) >= 0) / max(1, len([m for m in mob_online if m.get("battery", -1) >= 0])), 1) if mob_online else -1,2150               "jobsActive": sum((m.get("mobile") or {}).get("activeJobs", 0) for m in mob_online)}2151    return {"macs": agg(macs), "servers": agg(servers), "mobiles": mobiles, "all": agg(snap)}215221532154def summary_snapshot():2155    with latest_lock:2156        snap = {k: dict(v) for k, v in latest.items()}2157    aggs = cluster_aggregates(snap)2158    apps = apps_summary()2159    tun = tunnel_snapshot()2160    gws = tun.get("gateways", [])2161    peers = [p for g in gws if g.get("primary") for p in g.get("wg", {}).get("peers", [])]2162    with mobile_lock:2163        mob = dict(mobile_latest)2164    inc = incidents_snapshot()2165    crit = inc["counts"].get("critical", 0)2166    warn = inc["counts"].get("warning", 0)2167    health = "critical" if crit else ("warning" if warn else "ok")2168    score = max(0, 100 - crit * 15 - warn * 5)2169    return {2170        "health": health, "score": score, "agentVersion": AGENT_VERSION, "ts": time.time(),2171        "uptime": int(time.time() - started_at),2172        "nodes": aggs["all"], "macs": aggs["macs"], "servers": aggs["servers"], "mobiles": aggs["mobiles"],2173        "apps": apps,2174        "tunnel": {"gatewaysOk": sum(1 for g in gws if g.get("ok")), "gatewaysTotal": len(gws),2175                   "peersOnline": sum(1 for p in peers if p.get("online")), "peersTotal": len(peers),2176                   "routes": tun.get("routesTotal", 0)},2177        "mobile": {"ok": mob.get("ok", False), "online": mob.get("online", 0), "total": mob.get("total", 0),2178                   "queued": mob.get("queued", 0), "active": mob.get("active", 0)},2179        "incidents": inc["counts"], "topIncidents": inc["incidents"][:6],2180        "sites": {k: v for k, v in sites_snapshot().items() if k in ("total", "ok", "down")},2181        "egress": egress_public(),2182        "events": query_events(0, 8),2183    }218421852186def public_record(rec):2187    """Copie d'un enregistrement d'app sans les champs internes (chemins de logs),2188    sans toucher aux dicts stockés dans apps_latest."""2189    out = dict(rec)2190    out["processes"] = [{k: v for k, v in p.items() if k not in ("outLog", "errLog")}2191                        for p in rec.get("processes", [])]2192    return out219321942195def public_apps_snapshot():2196    with apps_lock:2197        apps = [public_record(a) for a in apps_latest.values()]2198    apps.sort(key=lambda a: ({"down": 0, "degraded": 1, "unknown": 2, "up": 3}[a["state"]], a["app"]))2199    return apps220022012202class Handler(BaseHTTPRequestHandler):2203    server_version = f"maclustr-agentd/{AGENT_VERSION}"2204    protocol_version = "HTTP/1.1"22052206    def _send(self, code, payload):2207        body = json.dumps(payload, ensure_ascii=False).encode()2208        self.send_response(code)2209        self.send_header("Content-Type", "application/json; charset=utf-8")2210        self.send_header("Content-Length", str(len(body)))2211        self.send_header("Cache-Control", "no-store")2212        self.end_headers()2213        self.wfile.write(body)22142215    def _auth(self):2216        h = self.headers.get("Authorization", "")2217        if h == f"Bearer {TOKEN}":2218            return True2219        q = parse_qs(urlparse(self.path).query)2220        if (q.get("token") or [""])[0] == TOKEN:2221            return True2222        self._send(401, {"error": "unauthorized"})2223        return False22242225    def log_message(self, fmt, *args):2226        pass22272228    def _read_json(self):2229        try:2230            length = int(self.headers.get("Content-Length", 0))2231            return json.loads(self.rfile.read(length) or b"{}")2232        except Exception:2233            return None22342235    # -- SSE ---------------------------------------------------------------22362237    def _stream(self):2238        self.send_response(200)2239        self.send_header("Content-Type", "text/event-stream; charset=utf-8")2240        self.send_header("Cache-Control", "no-store")2241        self.send_header("Connection", "keep-alive")2242        self.send_header("X-Accel-Buffering", "no")2243        self.end_headers()2244        try:2245            self.wfile.write(b": maclustr-agentd stream\n\n")2246            self.wfile.write(("event: hello\ndata: %s\n\n" % json.dumps(2247                {"version": AGENT_VERSION, "seq": tick_seq, "ts": time.time()})).encode())2248            self.wfile.flush()2249            seen = tick_seq2250            while True:2251                with tick_cv:2252                    tick_cv.wait(timeout=15)2253                    seq, last = tick_seq, dict(tick_last)2254                if seq != seen:2255                    seen = seq2256                    self.wfile.write(("event: tick\ndata: %s\n\n" % json.dumps(last)).encode())2257                else:2258                    self.wfile.write(b": keepalive\n\n")2259                self.wfile.flush()2260        except (BrokenPipeError, ConnectionResetError, OSError):2261            return22622263    # -- GET -----------------------------------------------------------------22642265    def do_GET(self):2266        u = urlparse(self.path)2267        parts = [p for p in u.path.split("/") if p]2268        q = parse_qs(u.query)2269        if u.path == "/health":2270            with latest_lock:2271                online = sum(1 for m in latest.values() if m.get("status") not in ("offline", "unknown", None))2272            s = apps_summary()2273            return self._send(200, {"ok": True, "version": AGENT_VERSION,2274                                    "uptime": int(time.time() - started_at),2275                                    "cycles": cycle_count, "appsCycles": apps_cycle,2276                                    "nodesOnline": online, "nodesTotal": len(NODES),2277                                    "appsUp": s["up"], "appsTotal": s["total"],2278                                    "appsDown": s["down"], "appsDegraded": s["degraded"],2279                                    "registryUpdated": registry.get("updated"),2280                                    "tunnelGateways": {k: v.get("ok") for k, v in tunnel_latest.items()},2281                                    "mobileOnline": mobile_latest.get("online", 0),2282                                    "mobileTotal": mobile_latest.get("total", 0),2283                                    "mobileQueued": mobile_latest.get("queued", 0),2284                                    "serversOnline": sum(1 for n in SERVER_NAMES if latest.get(n, {}).get("status") not in ("offline", "unknown", None)),2285                                    "serversTotal": len(SERVER_NAMES),2286                                    "macsTotal": len(MAC_NAMES),2287                                    "incidents": dict(incidents_counts)})2288        if not self._auth():2289            return2290        if u.path == "/api/stream":2291            return self._stream()2292        if u.path == "/api/cluster":2293            with latest_lock:2294                snap = {k: {kk: vv for kk, vv in v.items() if not kk.startswith("_")}2295                        for k, v in latest.items()}2296            with mobile_lock:2297                mobile = dict(mobile_latest)2298            mob_nodes = mobile.pop("nodes", []) or []2299            mob_metrics = mobile.pop("metrics", {}) or {}2300            snap.update(mob_metrics)2301            return self._send(200, {"nodes": NODE_INFO + mob_nodes, "metrics": snap,2302                                    "lanMap": lan_map, "ts": time.time(),2303                                    "apps": apps_summary(), "appsByNode": apps_by_node(),2304                                    "mobile": mobile, "servers": SERVER_NAMES, "mobiles": [n["name"] for n in mob_nodes],2305                                    "incidents": dict(incidents_counts),2306                                    "aggregates": cluster_aggregates(snap),2307                                    "agentVersion": AGENT_VERSION})2308        if u.path == "/api/summary":2309            return self._send(200, summary_snapshot())2310        if u.path == "/api/incidents":2311            return self._send(200, incidents_snapshot())2312        if u.path == "/api/sites":2313            return self._send(200, sites_snapshot())2314        if u.path == "/api/egress":2315            return self._send(200, egress_snapshot())2316        if u.path == "/api/servers":2317            with latest_lock:2318                snap = {k: {kk: vv for kk, vv in v.items() if not kk.startswith("_")}2319                        for k, v in latest.items() if k in SERVER_META}2320            return self._send(200, {"servers": [n for n in NODE_INFO if n["kind"] == "server"],2321                                    "metrics": snap, "ts": time.time()})2322        if u.path == "/api/ops/commands":2323            return self._send(200, {"commands": [{"key": k, **v} for k, v in MLD_COMMANDS.items()],2324                                    "gateway": GATEWAY_NODE})2325        if u.path == "/api/mobile":2326            with mobile_lock:2327                return self._send(200, dict(mobile_latest))2328        if parts[:2] == ["api", "mobile"] and len(parts) >= 3:2329            sub = "/" + "/".join(parts[2:]) + (("?" + u.query) if u.query else "")2330            if parts[2] in ("workers", "jobs", "stats"):2331                code, payload = mobile_proxy("GET", "/api" + sub, timeout=40)2332                return self._send(code, payload)2333        if u.path == "/api/history":2334            window = {"1h": 3600, "6h": 21600, "24h": 86400}.get(2335                (q.get("window") or ["1h"])[0], 3600)2336            bucket = {3600: 60, 21600: 300, 86400: 900}[window]2337            node = (q.get("node") or [None])[0]2338            return self._send(200, {"points": query_history(window, bucket, node)})2339        if u.path == "/api/nas":2340            return self._send(200, {"nas": nas_status()})2341        if u.path == "/api/tunnel":2342            return self._send(200, tunnel_snapshot())2343        if u.path == "/api/tunnel/refresh":2344            threading.Thread(target=tunnel_cycle_run, daemon=True).start()2345            return self._send(202, {"ok": True})2346        if u.path == "/api/apps":2347            return self._send(200, {"apps": public_apps_snapshot(), "summary": apps_summary(),2348                                    "byNode": apps_by_node(),2349                                    "registryUpdated": registry.get("updated"),2350                                    "registrySource": registry_meta.get("source"),2351                                    "registryError": registry_meta.get("error"),2352                                    "ts": time.time()})2353        if u.path == "/api/events":2354            since = float((q.get("since") or ["0"])[0] or 0)2355            limit = int((q.get("limit") or ["200"])[0])2356            kind = (q.get("kind") or [None])[0]2357            subject = (q.get("subject") or [None])[0]2358            return self._send(200, {"events": query_events(since, limit, kind, subject), "ts": time.time()})2359        if u.path == "/api/registry":2360            with registry_lock:2361                r = {"updated": registry.get("updated"), "gateway": registry.get("gateway"),2362                     "apps": registry.get("apps", {}), "history": (registry.get("history") or [])[-50:]}2363            r["meta"] = dict(registry_meta)2364            return self._send(200, r)2365        if len(parts) >= 3 and parts[0] == "api" and parts[1] == "apps":2366            app = parts[2]2367            with apps_lock:2368                rec = apps_latest.get(app)2369            if not rec:2370                return self._send(404, {"error": "unknown app"})2371            if len(parts) == 3:2372                out = public_record(rec)2373                out["events"] = query_events(time.time() - 7 * 86400, 50, None, app)2374                return self._send(200, out)2375            what = parts[3]2376            if what == "logs":2377                lines = (q.get("lines") or ["120"])[0]2378                proc = (q.get("process") or [None])[0]2379                res = app_logs(app, lines, proc)2380                return self._send(200 if res else 404, res or {"error": "unknown app"})2381            if what == "history":2382                window = {"1h": 3600, "6h": 21600, "24h": 86400, "7d": 7 * 86400}.get(2383                    (q.get("window") or ["24h"])[0], 86400)2384                bucket = {3600: 60, 21600: 300, 86400: 600, 7 * 86400: 3600}[window]2385                return self._send(200, {"points": query_app_history(app, window, bucket), "window": window})2386            if what == "events":2387                return self._send(200, {"events": query_events(0, 200, None, app)})2388        if len(parts) == 4 and parts[0] == "api" and parts[1] == "node":2389            name, what = parts[2], parts[3]2390            if name not in NODE_INDEX:2391                with mobile_lock:2392                    is_mobile = name in (mobile_latest.get("metrics") or {})2393                if is_mobile:2394                    if what == "apps":2395                        return self._send(200, {"apps": []})2396                    return self._send(400, {"error": "nœud mobile : pas de SSH — utiliser /api/mobile/workers/%s" % name, "kind": "mobile"})2397                return self._send(404, {"error": "unknown node"})2398            if what == "top":2399                out = run_on_node(name, TOP_SCRIPT_LINUX if platform_of(name) == "linux" else TOP_SCRIPT, timeout=20)2400                sections = {"cpu": [], "mem": []}2401                current = None2402                for line in out.splitlines():2403                    if line.startswith("---CPU---"):2404                        current = "cpu"2405                    elif line.startswith("---MEM---"):2406                        current = "mem"2407                    elif current and not line.strip().startswith("PID"):2408                        p = line.split(None, 5)2409                        if len(p) >= 6:2410                            try:2411                                sections[current].append({2412                                    "pid": int(p[0]),2413                                    "cpu": float(p[1].replace(",", ".")),2414                                    "mem": float(p[2].replace(",", ".")),2415                                    "rssKB": int(p[3]),2416                                    "user": p[4],2417                                    "command": p[5],2418                                })2419                            except ValueError:2420                                continue2421                return self._send(200, {"topCpu": sections["cpu"], "topMem": sections["mem"]})2422            if what == "ports":2423                out = run_on_node(name, PORTS_SCRIPT_LINUX if platform_of(name) == "linux" else PORTS_SCRIPT, timeout=20)2424                ports = []2425                for line in out.splitlines():2426                    p = line.split()2427                    if len(p) >= 3:2428                        mport = re.search(r":(\d+)$", p[2])2429                        if mport:2430                            ports.append({"command": p[0], "pid": int(p[1]),2431                                          "bind": p[2].rsplit(":", 1)[0],2432                                          "port": int(mport.group(1))})2433                return self._send(200, {"ports": ports})2434            if what == "apps":2435                with apps_lock:2436                    apps = [public_record(a) for a in apps_latest.values() if a.get("node") == name]2437                return self._send(200, {"apps": apps})2438        return self._send(404, {"error": "not found"})24392440    # -- POST ----------------------------------------------------------------24412442    def do_POST(self):2443        if not self._auth():2444            return2445        u = urlparse(self.path)2446        parts = [p for p in u.path.split("/") if p]2447        if parts[:2] == ["api", "mobile"] and len(parts) >= 3 and parts[2] in ("jobs", "workers"):2448            body = self._read_json()2449            if body is None:2450                return self._send(400, {"error": "bad json"})2451            sub = "/" + "/".join(parts[2:])2452            method = "PATCH" if parts[2] == "workers" else "POST"2453            code, payload = mobile_proxy(method, "/api" + sub, body, timeout=40)2454            if code < 300 and parts[2] == "jobs":2455                record_event("action", "mobile", "", "job", "%d job(s) mobile soumis depuis l'app MacLustr (%s)" % (payload.get("count", 1), (body.get("type") if isinstance(body, dict) else "lot")))2456            threading.Thread(target=mobile_fetch, daemon=True).start()2457            return self._send(code, payload)2458        if len(parts) == 4 and parts[0] == "api" and parts[1] == "node" and parts[3] == "action":2459            name = parts[2]2460            if name not in NODE_INDEX:2461                with mobile_lock:2462                    is_mobile = name in (mobile_latest.get("metrics") or {})2463                if is_mobile:2464                    body = self._read_json() or {}2465                    act = body.get("action", "")2466                    if act in ("ping", "bench", "sysinfo"):2467                        jt = {"ping": "ping", "bench": "compute.bench", "sysinfo": "sys.info"}[act]2468                        code, payload = mobile_proxy("POST", "/api/jobs", {"type": jt, "target": name, "tag": "app-action"}, timeout=30)2469                        return self._send(code, payload)2470                    return self._send(400, {"error": "nœud mobile : actions possibles ping | bench | sysinfo"})2471                return self._send(404, {"error": "unknown node"})2472            body = self._read_json()2473            if body is None:2474                return self._send(400, {"error": "bad json"})2475            result = do_action(name, body.get("action", ""), body.get("pid"))2476            return self._send(200 if result.get("ok") else 500, result)2477        if len(parts) == 4 and parts[0] == "api" and parts[1] == "apps" and parts[3] == "action":2478            body = self._read_json()2479            if body is None:2480                return self._send(400, {"error": "bad json"})2481            result = app_action(parts[2], body.get("action", ""), body.get("process"))2482            return self._send(200 if result.get("ok") else 500, result)2483        if u.path == "/api/registry/refresh":2484            ok = pull_registry()2485            if ok:2486                threading.Thread(target=_recheck_soon, daemon=True).start()2487            return self._send(200 if ok else 502, {"ok": ok, "updated": registry.get("updated"),2488                                                   "apps": len(registry_apps()),2489                                                   "error": registry_meta.get("error")})2490        if u.path == "/api/apps/refresh":2491            threading.Thread(target=apps_cycle_run, daemon=True).start()2492            return self._send(202, {"ok": True})2493        if u.path == "/api/tunnel/refresh":2494            threading.Thread(target=tunnel_cycle_run, daemon=True).start()2495            return self._send(202, {"ok": True})2496        if u.path == "/api/sites/refresh":2497            threading.Thread(target=sites_cycle_run, daemon=True).start()2498            return self._send(202, {"ok": True})2499        if u.path == "/api/egress/refresh":2500            threading.Thread(target=egress_cycle_run, daemon=True).start()2501            return self._send(202, {"ok": True})2502        if u.path == "/api/incidents/refresh":2503            threading.Thread(target=incidents_cycle_run, daemon=True).start()2504            return self._send(202, {"ok": True})2505        if u.path in ("/api/incidents/ack", "/api/incidents/unack"):2506            body = self._read_json()2507            if body is None or not body.get("id"):2508                return self._send(400, {"error": "id requis"})2509            iid = str(body["id"])2510            if u.path.endswith("/ack"):2511                incident_acks[iid] = {"ts": time.time(), "note": str(body.get("note") or "")[:200]}2512                save_ack(iid, incident_acks[iid]["note"])2513                record_event("action", iid, "", "ack", "incident pris en charge depuis l'app MacLustr")2514            else:2515                incident_acks.pop(iid, None)2516                save_ack(iid, remove=True)2517            threading.Thread(target=incidents_cycle_run, daemon=True).start()2518            return self._send(200, {"ok": True, "id": iid, "acked": u.path.endswith("/ack")})2519        if u.path == "/api/ops/mld":2520            body = self._read_json()2521            if body is None or not body.get("command"):2522                return self._send(400, {"error": "command requis", "allowed": sorted(MLD_COMMANDS)})2523            res = run_mld(str(body["command"]))2524            return self._send(200 if res.get("ok") else 500, res)2525        return self._send(404, {"error": "not found"})252625272528    def do_PATCH(self):2529        if not self._auth():2530            return2531        u = urlparse(self.path)2532        parts = [p for p in u.path.split("/") if p]2533        if parts[:3] == ["api", "mobile", "workers"] and len(parts) == 4:2534            body = self._read_json()2535            if body is None:2536                return self._send(400, {"error": "bad json"})2537            code, payload = mobile_proxy("PATCH", "/api/workers/" + parts[3], body)2538            threading.Thread(target=mobile_fetch, daemon=True).start()2539            return self._send(code, payload)2540        return self._send(404, {"error": "not found"})25412542    def do_DELETE(self):2543        if not self._auth():2544            return2545        u = urlparse(self.path)2546        parts = [p for p in u.path.split("/") if p]2547        if parts[:2] == ["api", "mobile"] and len(parts) == 4 and parts[2] in ("jobs", "workers"):2548            code, payload = mobile_proxy("DELETE", "/api/%s/%s" % (parts[2], parts[3]))2549            threading.Thread(target=mobile_fetch, daemon=True).start()2550            return self._send(code, payload)2551        return self._send(404, {"error": "not found"})255225532554class Server(ThreadingHTTPServer):2555    daemon_threads = True2556    allow_reuse_address = True255725582559def main():2560    os.makedirs(BASE_DIR, exist_ok=True)2561    load_lan_map()2562    load_registry_file()2563    threading.Thread(target=collector_loop, daemon=True).start()2564    threading.Thread(target=tunnel_loop, daemon=True).start()2565    threading.Thread(target=apps_loop, daemon=True).start()2566    threading.Thread(target=mobile_loop, daemon=True).start()2567    try:2568        load_acks()2569    except Exception as e:  # noqa: BLE0012570        print("acks load error:", e, flush=True)2571    threading.Thread(target=incidents_loop, daemon=True).start()2572    threading.Thread(target=sites_loop, daemon=True).start()2573    threading.Thread(target=egress_loop, daemon=True).start()2574    if not registry_apps():2575        threading.Thread(target=pull_registry, daemon=True).start()2576    srv = Server(("0.0.0.0", PORT), Handler)2577    print(f"maclustr-agentd {AGENT_VERSION} listening on :{PORT}", flush=True)2578    srv.serve_forever()257925802581if __name__ == "__main__":2582    main()2583