spb/maclustr-agentd
Public
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