#!/usr/bin/env python3 """ MacLustr Agent v2 — collecteur de métriques du cluster + santé des apps, servi en HTTP/JSON pour les apps iOS et macOS MacLustr. Tourne sur M4M64a sous launchd (io.maclustr.agentd, port 9210). Fan-out SSH vers les autres nœuds via le LAN 192.168.2.x (le SSH inter-nœuds par Tailscale est bloqué par ACL). Les nœuds sont identifiés par le fichier marqueur ~/.maclustr-node (les LocalHostName ne correspondent pas aux alias et les IP LAN sont en DHCP → redécouverte par balayage du /24). v2 (2026-09-04) — suivi des apps déployées : * lit le registre maclustr-dispatch (M1M32:~/dispatch/registry.json), poussé par `mld` (abonné) et re-tiré périodiquement par SSH ; * toutes les APPS_INTERVAL s : HTTP local (ip:port) + public (https://domaine), état des processus PM2 (`pm2 jlist`) et launchd (`launchctl list`) sur chaque nœud ; état synthétique up / degraded / down ; historique SQLite 7 j ; journal d'événements (transitions apps + nœuds) ; * endpoints /api/apps, /api/apps/[/logs|/history|/action], /api/events, /api/registry, flux SSE /api/stream (tick à chaque cycle). Zéro dépendance : stdlib uniquement (http.server, sqlite3, subprocess, threading, urllib). Python 3.9 système obligatoire (/usr/bin/python3, Local Network Privacy). """ import json import os import re import shlex import socket import sqlite3 import ssl import subprocess import threading import time import urllib.error import urllib.request from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer from urllib.parse import urlparse, parse_qs AGENT_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 Tunnel HOME = os.path.expanduser("~") BASE_DIR = os.path.join(HOME, "maclustr-agentd") CONF_PATH = os.path.join(BASE_DIR, "config.json") LANMAP_PATH = os.path.join(BASE_DIR, "lanmap.json") REGISTRY_PATH = os.path.join(BASE_DIR, "registry.json") DB_PATH = os.path.join(BASE_DIR, "history.db") PORT = 9210 COLLECT_INTERVAL = 20 # secondes entre deux cycles métriques APPS_INTERVAL = 30 # secondes entre deux cycles santé des apps REGISTRY_PULL_INTERVAL = 180 # re-tirage du registre depuis la passerelle RAW_RETENTION_H = 26 # heures d'historique brut conservées (métriques) APP_RETENTION_D = 7 # jours d'historique des vérifications d'apps EVENTS_RETENTION_D = 30 # jours d'événements conservés LAN_PREFIX = "192.168.2." GATEWAY_NODE = "M1M32" GATEWAY_REGISTRY = "~/dispatch/registry.json" HTTP_TIMEOUT_LOCAL = 6 HTTP_TIMEOUT_PUBLIC = 10 # MacLustr Tunnel (2026-09-10) : passerelles publiques OVH (WireGuard hub + Caddy) qui remplacent ngrok. # L'agent interroge `sudo tunnelctl json` sur chacune (clé maclustr-agentd autorisée pour ubuntu). TUNNEL_GATEWAYS = { "BHS64": {"host": "51.161.112.61", "user": "ubuntu", "primary": True, "site": "Beauharnois (Québec)"}, } TUNNEL_INTERVAL = 60 # secondes entre deux relevés des passerelles TUNNEL_PEER_FRESH_S = 180 # handshake WireGuard plus vieux → pair considéré hors ligne DNS_CACHE_S = 300 # --------------------------------------------------------------------------- # Inventaire (specs fixes ; lan_ip = graine, corrigée en live) # --------------------------------------------------------------------------- NODES = [ # name, hostname, chip, model, tier, gen, cores, mem_mb, gpu, lan_seed ("M3U96a", "M3U96a.maclustr.io", "M3 Ultra", "Mac Studio", "ultra", "m3", 32, 98304, 80, "192.168.2.166"), ("M3U96b", "M3U96b.maclustr.io", "M3 Ultra", "Mac Studio", "ultra", "m3", 32, 98304, 80, "192.168.2.125"), ("M2U64", "M2U64.maclustr.io", "M2 Ultra", "Mac Studio", "ultra", "m2", 24, 65536, 76, ""), ("M4M64a", "M4M64a.maclustr.io", "M4 Max", "Mac Studio", "max", "m4", 16, 65536, 40, "192.168.2.127"), ("M4M64b", "M4M64b.maclustr.io", "M4 Max", "Mac Studio", "max", "m4", 16, 65536, 40, "192.168.2.128"), # Mac Studio M4 Max ajouté le 2026-09-25 (LAN + Tailscale 100.125.225.117 ; DNS GoDaddy A → IP Tailscale ; S/N JX0PHFQHF1) ("M4M64c", "M4M64c.maclustr.io", "M4 Max", "Mac Studio", "max", "m4", 16, 65536, 40, "192.168.2.103"), # Mac Studio M5 Max + Mac mini M6 ajoutés le 2026-09-26 (LAN + Tailscale ; DNS GoDaddy A → IP Tailscale) ("M5M36", "M5M36.maclustr.io", "M5 Max", "Mac Studio", "max", "m5", 18, 36864, 40, "192.168.2.104"), ("M4BP48", "M4BP48.maclustr.io", "M4 Max", "MacBook Pro", "max", "m4", 16, 49152, 40, "192.168.2.137"), ("M4BP36", "M4BP36.maclustr.io", "M4 Max", "MacBook Pro", "max", "m4", 14, 36864, 32, "192.168.2.133"), ("M4M36", "M4M36.maclustr.io", "M4 Max", "Mac Studio", "max", "m4", 14, 36864, 40, "192.168.2.131"), ("M2M32", "M2M32.maclustr.io", "M2 Max", "Mac Studio", "max", "m2", 12, 32768, 38, "192.168.2.90"), ("M2M32b", "M2M32b.maclustr.io", "M2 Max", "Mac Studio", "max", "m2", 12, 32768, 38, "192.168.2.130"), ("M2M32c", "M2M32c.maclustr.io", "M2 Max", "Mac Studio", "max", "m2", 12, 32768, 30, "192.168.2.126"), ("m4mc", "m4mc.maclustr.io", "M4 Pro", "Mac mini", "pro", "m4", 12, 24576, 16, "192.168.2.134"), ("M1M32", "M1M32.maclustr.io", "M1 Max", "Mac Studio", "max", "m1", 10, 32768, 32, "192.168.2.132"), ("m4ma", "m4ma.maclustr.io", "M4", "Mac mini", "base", "m4", 10, 24576, 10, "192.168.2.73"), ("m4mb", "m4mb.maclustr.io", "M4", "Mac mini", "base", "m4", 10, 16384, 10, "192.168.2.136"), ("m4md", "m4md.maclustr.io", "M4", "Mac mini", "base", "m4", 10, 16384, 10, "192.168.2.108"), # 3 Mac mini locaux ajoutés le 2026-09-21 (LAN + Tailscale ; DNS GoDaddy A → IP Tailscale) ("m4ml", "m4ml.maclustr.io", "M4", "Mac mini", "base", "m4", 10, 16384, 10, "192.168.2.63"), ("m4mm", "m4mm.maclustr.io", "M4 Pro", "Mac mini", "pro", "m4", 12, 49152, 16, "192.168.2.92"), ("m4mn", "m4mn.maclustr.io", "M4 Pro", "Mac mini", "pro", "m4", 12, 49152, 16, "192.168.2.93"), ("m6ma", "m6ma.maclustr.io", "M6", "Mac mini", "base", "m6", 12, 24576, 10, "192.168.2.105"), ("m6mb", "m6mb.maclustr.io", "M6", "Mac mini", "base", "m6", 12, 24576, 10, "192.168.2.106"), # ajouté 2026-09-28 ("m5ma", "m5ma.maclustr.io", "M5 Pro", "Mac mini", "pro", "m5", 15, 24576, 16, "192.168.2.112"), # ajouté 2026-09-28 ("m2m16", "m2m16.maclustr.io", "M2 Pro", "Mac mini", "pro", "m2", 10, 16384, 16, "192.168.2.169"), ("M3BA24", "M3BA24.maclustr.io", "M3", "MacBook Air", "base", "m3", 8, 24576, 10, "192.168.2.145"), ("M3BA16", "M3BA16.maclustr.io", "M3", "MacBook Air", "base", "m3", 8, 16384, 10, "192.168.2.70"), ("m2m8a", "m2m8a.maclustr.io", "M2", "Mac mini", "base", "m2", 8, 8192, 10, "192.168.2.167"), ("m2m8b", "m2m8b.maclustr.io", "M2", "Mac mini", "base", "m2", 8, 8192, 10, "192.168.2.170"), ] NODE_INDEX = {n[0]: i for i, n in enumerate(NODES)} SSH_USER = "simon-pierreboucher" # Nœuds hors LAN : joints directement à leur IP publique, avec leur propre utilisateur SSH. # La clé `maclustr-agentd` (~/.ssh/id_ed25519_agentd) doit être dans leur authorized_keys. REMOTE_NODES = { # 2026-09-25 : les 12 Macs loués (Macly, MacStadium, rentamac) ont été retirés du cluster. } REMOTE_USER_BY_HOST = {r["host"]: r["user"] for r in REMOTE_NODES.values()} # --------------------------------------------------------------------------- # Serveurs dédiés Linux (OVHcloud) — 3.0.0. Même pipeline que les Macs (script # de métriques Linux équivalent clé pour clé), utilisateur `ubuntu`, clé agentd # autorisée. Ils ne sont jamais cherchés sur le LAN. # name, host public, user, site, rôle, threads, RAM Mo, matériel # --------------------------------------------------------------------------- SERVERS = [ ("BHS64", "51.161.112.61", "ubuntu", "Beauharnois (Québec)", "Passerelle MacLustr Tunnel (WireGuard + Caddy)", 16, 65536, "OVH ADVANCE-2 · AMD EPYC 4345P · 64 Go DDR5 ECC · 2×960 Go NVMe RAID 1"), # BHS64b, BHS128, BHS128b (Beauharnois) et R9128 (Gravelines) retirés du MacLustr le 2026-10-02 (apps arrêtées, serveurs résiliés). ] # Chemin de secours quand le port 22 sortant est bloqué (FAI/routeur, 2026-09-25) : le hub WireGuard # BHS64 reste joignable en 10.67.0.1 depuis les Macs raccordés (wg1). WG_HUB = "10.67.0.1" SERVER_ALT = { # name -> liste de (host, jump) essayés après l'accès direct "BHS64": [(WG_HUB, None)], } SERVER_META = {} PLATFORM = {} # name -> "macos" | "linux" for _s in SERVERS: _name, _host, _user, _site, _role, _thr, _mem, _hw = _s _chip = _hw.split("·")[1].strip() if "·" in _hw else "x86" NODES.append((_name, _host, _chip, "Serveur dédié", "server", "x86", _thr, _mem, 0, _host)) SERVER_META[_name] = {"site": _site, "role": _role, "user": _user, "hardware": _hw, "publicIP": _host} PLATFORM[_name] = "linux" REMOTE_USER_BY_HOST[_host] = _user REMOTE_USER_BY_HOST[WG_HUB] = "ubuntu" NODE_INDEX = {n[0]: i for i, n in enumerate(NODES)} SERVER_NAMES = [s[0] for s in SERVERS] MAC_NAMES = [n[0] for n in NODES if n[0] not in SERVER_META] def platform_of(name): return PLATFORM.get(name, "macos") def ssh_user_for(ip): """Utilisateur SSH pour une IP : celui du nœud distant si c'en est un, sinon l'utilisateur du cluster.""" return REMOTE_USER_BY_HOST.get(ip, SSH_USER) NAS_DEVICES = [ {"id": "nas-ugreen-1", "name": "UGREEN DXP4800 Plus", "host": "192.168.2.173", "mount": "nas_clustr"}, {"id": "nas-ugreen-2", "name": "UGREEN DH4300 Plus", "host": "192.168.2.175", "mount": "nas2_clustr"}, ] METRICS_SCRIPT = r""" LC_ALL=C; export LC_ALL IFACE=$(route -n get default 2>/dev/null | awk '/interface:/{print $2}') echo "===METRICS_START===" echo "CPU:$(top -l 1 -n 0 2>/dev/null | grep 'CPU usage' | awk '{print 100 - $7}' | tr -d '%' || echo 0)" echo "CPUBRK:$(top -l 1 -n 0 2>/dev/null | grep 'CPU usage' | awk '{print $3, $5, $7}' | tr -d '%' || echo '0 0 0')" echo "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}')" echo "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}')" echo "SWAP:$(sysctl -n vm.swapusage 2>/dev/null | awk '{print $6, $3}' | tr -d 'M' || echo '0 0')" echo "LOAD:$(sysctl -n vm.loadavg 2>/dev/null | tr -d '{}' || echo '0 0 0')" echo "DISK:$(df -g / 2>/dev/null | tail -1 | awk '{print $3,$2}' || echo '0 0')" echo "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')" echo "NET:$(netstat -ibn -I ${IFACE:-en0} 2>/dev/null | grep -v Link | tail -1 | awk '{print $7, $10}' || echo '0 0')" echo "TCP:$(netstat -an -p tcp 2>/dev/null | grep -c ESTABLISHED || echo 0)" echo "UPTIME:$(( $(date +%s) - $(sysctl -n kern.boottime 2>/dev/null | awk '{print $4}' | tr -d ',' || echo $(date +%s)) ))" echo "PROCS:$(ps aux 2>/dev/null | wc -l | tr -d ' ')" echo "USERS:$(who 2>/dev/null | wc -l | tr -d ' ')" echo "OS:$(sw_vers -productVersion 2>/dev/null)" echo "BATT:$(pmset -g batt 2>/dev/null | grep -o '[0-9]*%' | tr -d '%' | head -1 || echo -1)" echo "BATT_STATE:$(pmset -g batt 2>/dev/null | grep -oE 'charging|discharging|charged|AC Power|finishing charge' | head -1 || echo '')" echo "TOPPROC:$(ps -arco pcpu,comm 2>/dev/null | sed -n 2p | sed 's/^ *//')" echo "THERMAL:$(pmset -g therm 2>/dev/null | grep 'CPU_Scheduler_Limit' | awk '{print $3}' || echo 100)" echo "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)" echo "===METRICS_END===" """ # Équivalent Linux (Ubuntu) du script de métriques : mêmes clés, mêmes unités. # CPU = delta /proc/stat sur 1 s ; mémoire « pression » = (total − disponible)/total. LINUX_METRICS_SCRIPT = r""" LC_ALL=C; export LC_ALL export PATH=/usr/local/bin:/usr/bin:/bin:$HOME/.npm-global/bin:$PATH IFACE=$(ip route show default 2>/dev/null | awk '/default/{print $5; exit}') echo "===METRICS_START===" read u1 n1 s1 id1 io1 ir1 so1 st1 g1 < <(awk '/^cpu /{print $2,$3,$4,$5,$6,$7,$8,$9,$10}' /proc/stat); sleep 1 read u2 n2 s2 id2 io2 ir2 so2 st2 g2 < <(awk '/^cpu /{print $2,$3,$4,$5,$6,$7,$8,$9,$10}' /proc/stat) A=$(( (u2-u1)+(n2-n1)+(s2-s1)+(ir2-ir1)+(so2-so1)+(st2-st1) )); I=$(( (id2-id1)+(io2-io1) )); DT=$((A+I)); [ $DT -le 0 ] && DT=1 echo "CPU:$(awk -v a=$A -v d=$DT 'BEGIN{printf "%.2f", a*100/d}')" echo "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}')" echo "MEM_PRESSURE:$(free -m | awk '/Mem:/{printf "%.2f", ($2-$7)*100/$2}')" echo "MEMBRK:$(free -m | awk '/Mem:/{print $3, $6}')" echo "SWAP:$(free -m | awk '/Swap:/{print $3, $2}')" echo "LOAD:$(cut -d' ' -f1-3 /proc/loadavg)" echo "DISK:$(df -BG / | tail -1 | awk '{gsub("G","",$3); gsub("G","",$2); print $3, $2}')" echo "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}')" echo "NET:$(awk -v i="${IFACE:-eth0}:" '$1==i{print $2, $10}' /proc/net/dev)" echo "TCP:$(ss -tan state established 2>/dev/null | tail -n +2 | wc -l)" echo "UPTIME:$(cut -d. -f1 /proc/uptime)" echo "PROCS:$(ls /proc | grep -c '^[0-9]')" echo "USERS:$(who 2>/dev/null | wc -l)" echo "OS:$(. /etc/os-release 2>/dev/null; echo "${PRETTY_NAME:-Linux}")" echo "BATT:-1" echo "BATT_STATE:" echo "TOPPROC:$(ps -eo pcpu,comm --sort=-pcpu --no-headers 2>/dev/null | head -1 | sed 's/^ *//')" echo "THERMAL:100" echo "GPU:-1" echo "DOCKER:$(docker ps -q 2>/dev/null | wc -l)" echo "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')" echo "KERNEL:$(uname -r)" echo "===METRICS_END===" """ TOP_SCRIPT = ( "echo '---CPU---'; ps -arcwwwxo pid,pcpu,pmem,rss,user,comm | head -16;" "echo '---MEM---'; ps -amcwwwxo pid,pcpu,pmem,rss,user,comm | head -16" ) TOP_SCRIPT_LINUX = ( "echo '---CPU---'; ps -eo pid,pcpu,pmem,rss,user,comm --sort=-pcpu | head -16;" "echo '---MEM---'; ps -eo pid,pcpu,pmem,rss,user,comm --sort=-rss | head -16" ) PORTS_SCRIPT = ( "lsof -nP -iTCP -sTCP:LISTEN 2>/dev/null | tail -n +2 | " "awk '{print $1, $2, $9}' | sort -u" ) # même format de sortie « commande pid adresse:port » via ss (sudo sans mot de passe sur les OVH) PORTS_SCRIPT_LINUX = ( "sudo -n ss -ltnpH 2>/dev/null | awk '{print $4, $6}' | " "sed -E 's/^([^ ]+) users:\\(\\(\"([^\"]+)\",pid=([0-9]+).*/\\2 \\3 \\1/' | awk 'NF==3' | sort -u" ) def metrics_script_for(name): return LINUX_METRICS_SCRIPT if platform_of(name) == "linux" else METRICS_SCRIPT # PM2 + launchd + sondes HTTP locales (127.0.0.1) d'un nœud en un seul # aller-retour SSH. `pm2 jlist` démarre un démon PM2 s'il n'en existe pas : on # ne l'appelle que si le God Daemon tourne. `pm2 jlist` n'émet pas de saut de # ligne final → `echo` explicite, sinon le marqueur suivant se colle au JSON. PROCS_SCRIPT_HEAD = r""" export PATH=/opt/homebrew/bin:/usr/local/bin:$HOME/.npm-global/bin:$PATH echo '---PM2---' if pgrep -qf 'PM2.*God Daemon' 2>/dev/null; then pm2 jlist 2>/dev/null | tail -1; echo; else echo '[]'; fi echo '---LAUNCHD---' launchctl list 2>/dev/null """ def node_script(apps): """Script complet pour un nœud : processus + sonde HTTP locale de chaque app (curl sur 127.0.0.1 — beaucoup d'apps n'écoutent que sur loopback).""" parts = [PROCS_SCRIPT_HEAD] for e in apps: if not e.get("port"): continue url = "http://127.0.0.1:%s%s" % (e["port"], e.get("health_path") or "/") parts.append("echo '---HTTP %s---'" % e["app"]) parts.append("curl -s -o /dev/null -m %d -w '%%{http_code} %%{time_total}\\n' %s 2>/dev/null || echo '000 0'" % (HTTP_TIMEOUT_LOCAL, shlex.quote(url))) parts.append("echo '---END---'") return "\n".join(parts) PM2_ENV = "export PATH=/opt/homebrew/bin:/usr/local/bin:$HOME/.npm-global/bin:$PATH; " def load_config(): with open(CONF_PATH) as f: return json.load(f) CONFIG = load_config() TOKEN = CONFIG["token"] SUDO_PW = CONFIG.get("sudo_password", "") SELF_NODE = CONFIG.get("self_node", "M4M64a") SSH_KEY = os.path.expanduser(CONFIG.get("ssh_key", "~/.ssh/id_ed25519_agentd")) # Nœuds mobiles (2026-09-21) : coordinateur maclustr-mobile (PM2, même hôte) — les iPhones y tirent leurs jobs. MOBILE_URL = CONFIG.get("mobile_url", "http://127.0.0.1:9320").rstrip("/") MOBILE_TOKEN = CONFIG.get("mobile_token", "") MOBILE_INTERVAL = 20 mobile_latest = {"ok": False, "workers": [], "online": 0, "total": 0, "queued": 0, "active": 0, "ts": 0, "error": ""} mobile_lock = threading.Lock() def mobile_fetch(): """Relit /health et /api/workers du coordinateur mobile (loopback, aucune contrainte LNP).""" try: req = urllib.request.Request(MOBILE_URL + "/health") with urllib.request.urlopen(req, timeout=4) as r: h = json.loads(r.read()) workers = [] if MOBILE_TOKEN: req = urllib.request.Request(MOBILE_URL + "/api/workers", headers={"Authorization": "Bearer " + MOBILE_TOKEN}) with urllib.request.urlopen(req, timeout=4) as r: workers = json.loads(r.read()).get("workers", []) snap = {"ok": True, "version": h.get("version"), "workers": workers, "online": h.get("workersOnline", 0), "total": h.get("workersTotal", 0), "queued": h.get("queued", 0), "active": h.get("active", 0), "stats": h.get("stats"), "ts": time.time(), "error": ""} except Exception as e: snap = {"ok": False, "workers": mobile_latest.get("workers", []), "online": 0, "total": mobile_latest.get("total", 0), "queued": 0, "active": 0, "ts": time.time(), "error": str(e)[:200]} # 3.1.0 : chaque appareil devient un nœud (info statique + métriques au format NodeMetrics) snap["nodes"] = [mobile_node_info(w) for w in snap["workers"]] snap["metrics"] = {n["name"]: mobile_node_metrics(w, n["name"]) for w, n in zip(snap["workers"], snap["nodes"])} try: record_samples({k: v for k, v in snap["metrics"].items() if v.get("status") not in ("offline", "unknown")}) except Exception as e: # noqa: BLE001 print("mobile samples:", e, flush=True) with mobile_lock: prev_online = {w.get("name") for w in mobile_latest.get("workers", []) if w.get("online")} mobile_latest.clear() mobile_latest.update(snap) now_online = {w.get("name") for w in snap["workers"] if w.get("online")} for n in sorted(now_online - prev_online): if prev_online or mobile_latest.get("ts"): record_event("mobile", n, "offline", "online", "iPhone connecté au coordinateur") for n in sorted(prev_online - now_online): record_event("mobile", n, "online", "offline", "iPhone silencieux (app fermée ou en arrière-plan)") THERMAL_MAP = {"nominal": "nominal", "fair": "fair", "tiède": "fair", "serious": "serious", "chaud": "serious", "critical": "serious", "critique": "serious"} def mobile_display_name(w): return w.get("displayName") or w.get("alias") or w.get("name") or "mobile" def mobile_node_info(w): kind = w.get("kind") or ("ipad" if str(w.get("model", "")).lower().startswith("ipad") else "iphone") return {"name": mobile_display_name(w), "hostname": w.get("ip") or "", "chip": w.get("chip") or "", "model": w.get("marketing") or w.get("model") or "iPhone", "tier": "mobile", "generation": (w.get("chip") or "").lower().replace(" ", "-"), "cpuCores": int(w.get("cores") or 0), "memoryMB": int(w.get("ramMb") or 0), "gpuCores": 0, "remote": True, "user": "", "platform": w.get("platform") or ("ipados" if kind == "ipad" else "ios"), "kind": "mobile", "deviceKind": kind, "site": "Mobile (Tailscale)", "role": w.get("role") or "worker", "hardware": "%s · %s · %s Go · %s" % (w.get("marketing") or w.get("model") or "?", w.get("chip") or "?", round((w.get("ramMb") or 0) / 1024), "iPadOS " + str(w.get("os") or "") if kind == "ipad" else "iOS " + str(w.get("os") or "")), "publicIP": w.get("ip") or "", "workerName": w.get("name"), "alias": w.get("alias"), "pinned": bool(w.get("pinned")), "appVersion": w.get("appVersion") or "", "identifier": w.get("model") or "", "bench": w.get("bench")} def mobile_node_metrics(w, name): m = w.get("metrics") or {} mem_total = float(w.get("ramMb") or m.get("memTotalMb") or 0) mem_used = float(m.get("memUsedMb") or 0) avail = m.get("memAvailableMb") if avail is not None and mem_total: # mémoire « pression » = total − disponible pour l'app (iOS ne donne pas l'usage système) mem_used = max(mem_used, mem_total - float(avail)) batt = m.get("battery") batt_pct = round(float(batt) * 100, 1) if isinstance(batt, (int, float)) and batt >= 0 else -1.0 cpu = m.get("cpuUsage", m.get("cpuLoad")) cpu = float(cpu) if isinstance(cpu, (int, float)) else 0.0 if 0 < cpu <= 1.0 and "cpuUsage" not in m: cpu *= 100.0 disk_total = float(w.get("storageGb") or m.get("diskTotalGb") or 0) disk_free = float(m.get("diskFreeGb") or 0) online = bool(w.get("online")) status = "online" if online else "offline" thermal = THERMAL_MAP.get(str(m.get("thermal") or "nominal").lower(), "nominal") if online and (thermal == "serious" or cpu > 90): status = "warning" return {"name": name, "status": status, "ts": float(m.get("ts") or w.get("lastSeen") or 0), "cpu": round(cpu, 1), "cpuUser": 0.0, "cpuSystem": 0.0, "cpuIdle": round(100 - cpu, 1), "memUsedMB": round(mem_used, 1), "memTotalMB": mem_total, "memWiredMB": float(m.get("memUsedMb") or 0), "memCompressedMB": 0.0, "swapUsedMB": 0.0, "swapTotalMB": 0.0, "load": [0.0, 0.0, 0.0], "diskUsedGB": round(max(disk_total - disk_free, 0), 1) if disk_total else 0.0, "diskTotalGB": round(disk_total, 1), "extDiskUsedGB": 0.0, "extDiskTotalGB": 0.0, "netInKBs": 0.0, "netOutKBs": 0.0, "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, "os": ("iPadOS " if (w.get("kind") == "ipad") else "iOS ") + str(w.get("os") or ""), "battery": batt_pct, "batteryState": ("charging" if m.get("charging") else ("discharging" if batt_pct >= 0 else "")), "topProcess": "%d job(s) actif(s)" % int(m.get("activeJobs") or 0), "thermal": thermal, "gpu": -1.0, "platform": w.get("platform") or "ios", "kind": "mobile", "docker": 0, "pm2Online": 0, "pm2Total": 0, "kernel": "", "mobile": {"jobsDone": int(w.get("jobsDone") or 0), "jobsFailed": int(w.get("jobsFailed") or 0), "activeJobs": int(m.get("activeJobs") or 0), "concurrency": int(m.get("concurrency") or 0), "net": m.get("net") or "", "lowPower": bool(m.get("lowPower")), "screenOn": bool(m.get("screenOn", True)), "bytesFetched": int(m.get("bytesFetched") or 0), "ageS": w.get("ageS"), "pinned": bool(w.get("pinned")), "workerName": w.get("name"), "appVersion": w.get("appVersion") or "", "lastJobAt": m.get("lastJobAt") or 0, "bench": w.get("bench")}} def mobile_proxy(method, path, body=None, timeout=30): """Relais vers le coordinateur mobile avec son jeton (les apps n'ont qu'un seul jeton : celui de l'agent).""" if not MOBILE_TOKEN: return 503, {"error": "mobile_token absent de config.json"} data = json.dumps(body).encode() if body is not None else None req = urllib.request.Request(MOBILE_URL + path, data=data, method=method, headers={"Authorization": "Bearer " + MOBILE_TOKEN, "Content-Type": "application/json"}) try: with urllib.request.urlopen(req, timeout=timeout) as r: raw = r.read() return r.getcode(), (json.loads(raw) if raw else {"ok": True}) except urllib.error.HTTPError as e: try: return e.code, json.loads(e.read() or b"{}") except Exception: # noqa: BLE001 return e.code, {"error": "HTTP %d" % e.code} except Exception as e: # noqa: BLE001 return 502, {"error": "coordinateur mobile injoignable : %s" % str(e)[:160]} def mobile_loop(): while True: try: mobile_fetch() except Exception as e: print("mobile error:", e, flush=True) time.sleep(MOBILE_INTERVAL) # --------------------------------------------------------------------------- # SSH helpers # --------------------------------------------------------------------------- SSH_OPTS = [ "-i", SSH_KEY, "-o", "BatchMode=yes", "-o", "ConnectTimeout=6", "-o", "StrictHostKeyChecking=no", "-o", "UserKnownHostsFile=/dev/null", "-o", "LogLevel=ERROR", ] def run_local(script, timeout=25): try: r = subprocess.run(["/bin/sh", "-c", script], capture_output=True, text=True, timeout=timeout) return r.stdout except Exception: return "" _last_path = {} # name -> ("direct" | host via jump) qui a marché en dernier def jump_command(jump): """ProxyCommand pour un rebond `user@host` avec la clé et les options agentd (-J ne les propagerait pas : la connexion imbriquée échouerait sur la vérification de clé d'hôte).""" return "ssh " + " ".join(shlex.quote(o) for o in SSH_OPTS) + " -W %h:%p " + shlex.quote(jump) def run_ssh(ip, script, timeout=25, jump=None): """Exécute un script sh sur un nœud via son IP (LAN, ou publique pour un REMOTE_NODE). Renvoie stdout ('' si échec). `jump` = "user@host" → ProxyJump (rebond par le hub WireGuard quand le port 22 sortant est bloqué).""" try: opts = list(SSH_OPTS) + (["-o", "ProxyCommand=" + jump_command(jump)] if jump else []) r = subprocess.run( ["ssh"] + opts + [f"{ssh_user_for(ip)}@{ip}", "bash -s" if ssh_user_for(ip) == "ubuntu" else "sh -s"], input=script, capture_output=True, text=True, timeout=timeout, ) if r.returncode != 0 and not r.stdout: print(f"ssh {ip}{' via ' + jump if jump else ''} rc={r.returncode}: {r.stderr.strip()[-200:]}", flush=True) return r.stdout except Exception as e: print(f"ssh {ip} exception: {e}", flush=True) return "" def run_on_node(name, script, timeout=25): if name == SELF_NODE: return run_local(script, timeout) ip = lan_map.get(name) if not ip: return "" alts = SERVER_ALT.get(name, []) # le dernier chemin qui a marché passe en premier order = [("direct", ip, None)] + [(f"{h} via {j or 'wg'}", h, j) for h, j in alts] last = _last_path.get(name) if last and last != "direct": order.sort(key=lambda o: o[0] != last) for label, host, jump in order: out = run_ssh(host, script, timeout, jump) if out: if _last_path.get(name) != label: print(f"{name}: chemin SSH = {label}", flush=True) _last_path[name] = label return out if not alts: break return "" # --------------------------------------------------------------------------- # Découverte LAN (fichier marqueur ~/.maclustr-node sur chaque nœud) # --------------------------------------------------------------------------- lan_map = {} # name -> ip lan_map_lock = threading.Lock() _last_scan = 0.0 def load_lan_map(): global lan_map seeds = {n[0]: n[9] for n in NODES if n[9]} try: with open(LANMAP_PATH) as f: saved = json.load(f) seeds.update({k: v for k, v in saved.items() if v}) except Exception: pass # les nœuds hors LAN ont une IP publique fixe : elle prime sur tout lanmap.json antérieur seeds.update({k: r["host"] for k, r in REMOTE_NODES.items()}) seeds.update({k: v["publicIP"] for k, v in SERVER_META.items()}) # serveurs OVH : IP publique fixe lan_map = seeds def save_lan_map(): try: with open(LANMAP_PATH, "w") as f: json.dump(lan_map, f, indent=1) except Exception: pass def ping(ip): return subprocess.run(["ping", "-c", "1", "-W", "800", "-q", ip], capture_output=True).returncode == 0 def identify(ip): """Renvoie l'alias maclustr du nœud à cette IP, ou None.""" out = run_ssh(ip, "cat ~/.maclustr-node 2>/dev/null", timeout=10).strip() return out if out in NODE_INDEX else None def rescan_lan(missing): """Balaye le /24 pour retrouver les nœuds dont l'IP a changé (max 1/5 min).""" global _last_scan if time.time() - _last_scan < 300 or not missing: return _last_scan = time.time() alive = [] sem = threading.Semaphore(48) lock = threading.Lock() def probe(ip): with sem: if ping(ip): with lock: alive.append(ip) threads = [threading.Thread(target=probe, args=(f"{LAN_PREFIX}{i}",)) for i in range(1, 255)] for t in threads: t.start() for t in threads: t.join() missing = [n for n in missing if n not in REMOTE_NODES and n not in SERVER_META] # hors LAN : jamais sur le /24 healthy_ips = {lan_map[n] for n in lan_map if n not in missing and lan_map.get(n)} candidates = [ip for ip in alive if ip not in healthy_ips] found = {} sem2 = threading.Semaphore(12) def check(ip): with sem2: who = identify(ip) if who: with lock: found[who] = ip threads = [threading.Thread(target=check, args=(ip,)) for ip in candidates] for t in threads: t.start() for t in threads: t.join() if found: with lan_map_lock: lan_map.update(found) save_lan_map() # --------------------------------------------------------------------------- # Parsing des métriques # --------------------------------------------------------------------------- def default_metrics(node, status="offline"): """Structure complète (toutes les clés) pour un nœud sans données.""" name, mem_mb = node[0], node[7] return { "name": name, "status": status, "ts": time.time(), "cpu": 0.0, "cpuUser": 0.0, "cpuSystem": 0.0, "cpuIdle": 0.0, "memUsedMB": 0.0, "memTotalMB": float(mem_mb), "memWiredMB": 0.0, "memCompressedMB": 0.0, "swapUsedMB": 0.0, "swapTotalMB": 0.0, "load": [0.0, 0.0, 0.0], "diskUsedGB": 0.0, "diskTotalGB": 0.0, "extDiskUsedGB": 0.0, "extDiskTotalGB": 0.0, "netInKBs": 0.0, "netOutKBs": 0.0, "tcp": 0, "uptime": 0, "procs": 0, "users": 0, "os": "", "battery": -1.0, "batteryState": "", "topProcess": "", "thermal": "nominal", "gpu": -1.0, "platform": platform_of(name), "kind": "server" if name in SERVER_META else "mac", "docker": 0, "pm2Online": 0, "pm2Total": 0, "kernel": "", "_netBytesIn": 0.0, "_netBytesOut": 0.0, } def parse_metrics(out, node, prev): name, _, chip, model, tier, gen, cores, mem_mb, gpu_cores, _ = node now = time.time() m = default_metrics(node, "online") m["ts"] = now for line in out.splitlines(): if ":" not in line: continue key, _, val = line.partition(":") val = val.strip().replace(",", ".") try: if key == "CPU": m["cpu"] = float(val or 0) elif key == "CPUBRK": p = [float(x) for x in val.split()] if len(p) >= 3: m["cpuUser"], m["cpuSystem"], m["cpuIdle"] = p[0], p[1], p[2] elif key == "MEM_PRESSURE": m["memUsedMB"] = mem_mb * (float(val or 50) / 100.0) elif key == "MEMBRK": p = [float(x) for x in val.split()] if len(p) >= 2: m["memWiredMB"], m["memCompressedMB"] = p[0], p[1] elif key == "SWAP": p = [float(x) for x in val.split()] if len(p) >= 2: m["swapUsedMB"], m["swapTotalMB"] = p[0], p[1] elif key == "LOAD": p = [float(x) for x in val.split()] if len(p) >= 3: m["load"] = p[:3] elif key == "DISK": p = [float(x) for x in val.split()] if len(p) >= 2: m["diskUsedGB"], m["diskTotalGB"] = p[0], p[1] elif key == "DISK_EXT": p = [float(x) for x in val.split()] if len(p) >= 2: m["extDiskUsedGB"], m["extDiskTotalGB"] = p[0], p[1] elif key == "NET": p = [float(x) for x in val.split()] if len(p) >= 2: m["_netBytesIn"], m["_netBytesOut"] = p[0], p[1] if prev and prev.get("_netBytesIn", 0) > 0: dt = now - prev.get("ts", now) if dt > 0 and p[0] >= prev["_netBytesIn"] and p[1] >= prev["_netBytesOut"]: m["netInKBs"] = (p[0] - prev["_netBytesIn"]) / 1024.0 / dt m["netOutKBs"] = (p[1] - prev["_netBytesOut"]) / 1024.0 / dt elif key == "TCP": m["tcp"] = int(val or 0) elif key == "UPTIME": m["uptime"] = int(float(val or 0)) elif key == "PROCS": m["procs"] = int(val or 0) elif key == "USERS": m["users"] = int(val or 0) elif key == "OS": m["os"] = val elif key == "BATT": m["battery"] = float(val or -1) elif key == "BATT_STATE": m["batteryState"] = val elif key == "TOPPROC": parts = val.split() if len(parts) >= 2: m["topProcess"] = f"{' '.join(parts[1:])} ({parts[0]}%)" elif val: m["topProcess"] = val elif key == "THERMAL": pct = int(float(val or 100)) m["thermal"] = "nominal" if pct >= 90 else ("fair" if pct >= 70 else "serious") elif key == "GPU": m["gpu"] = float(val or -1) elif key == "DOCKER": m["docker"] = int(val or 0) elif key == "PM2": p = val.split() if len(p) >= 2: m["pm2Online"], m["pm2Total"] = int(p[0]), int(p[1]) elif key == "KERNEL": m["kernel"] = val except (ValueError, IndexError): continue mem_pct = m["memUsedMB"] / mem_mb * 100 if mem_mb else 0 if m["cpu"] > 90 or mem_pct > 95: m["status"] = "critical" elif m["cpu"] > 75 or mem_pct > 85 or m["thermal"] != "nominal": m["status"] = "warning" return m # --------------------------------------------------------------------------- # Historique SQLite (métriques, vérifications d'apps, événements) # --------------------------------------------------------------------------- db_lock = threading.Lock() def db(): conn = sqlite3.connect(DB_PATH) conn.execute( "CREATE TABLE IF NOT EXISTS samples (" "ts INTEGER, node TEXT, cpu REAL, mem_pct REAL, gpu REAL," "net_in REAL, net_out REAL, disk_pct REAL)" ) conn.execute("CREATE INDEX IF NOT EXISTS idx_ts ON samples(ts)") conn.execute( "CREATE TABLE IF NOT EXISTS app_checks (" "ts INTEGER, app TEXT, state TEXT, local_code INTEGER, local_ms INTEGER," "public_code INTEGER, public_ms INTEGER)" ) conn.execute("CREATE INDEX IF NOT EXISTS idx_app_ts ON app_checks(app, ts)") conn.execute( "CREATE TABLE IF NOT EXISTS events (" "id INTEGER PRIMARY KEY AUTOINCREMENT, ts INTEGER, kind TEXT, subject TEXT," "from_state TEXT, to_state TEXT, detail TEXT)" ) conn.execute("CREATE INDEX IF NOT EXISTS idx_events_ts ON events(ts)") return conn def record_samples(metrics): rows = [] now = int(time.time()) for m in metrics.values(): if m["status"] in ("offline", "unknown"): continue mem_pct = m["memUsedMB"] / m["memTotalMB"] * 100 if m["memTotalMB"] else 0 disk_pct = m["diskUsedGB"] / m["diskTotalGB"] * 100 if m["diskTotalGB"] else 0 rows.append((now, m["name"], m["cpu"], mem_pct, m["gpu"], m["netInKBs"], m["netOutKBs"], disk_pct)) if not rows: return with db_lock: conn = db() conn.executemany("INSERT INTO samples VALUES (?,?,?,?,?,?,?,?)", rows) conn.execute("DELETE FROM samples WHERE ts < ?", (now - RAW_RETENTION_H * 3600,)) conn.commit() conn.close() def query_history(window_s, bucket_s, node=None): since = int(time.time()) - window_s where = "ts >= ?" args = [since] if node: where += " AND node = ?" args.append(node) with db_lock: conn = db() rows = conn.execute( f"SELECT (ts/{bucket_s})*{bucket_s} AS b, AVG(cpu), AVG(mem_pct)," f" AVG(CASE WHEN gpu >= 0 THEN gpu END), SUM(net_in)/COUNT(DISTINCT node)," f" SUM(net_out)/COUNT(DISTINCT node), AVG(disk_pct)" f" FROM samples WHERE {where} GROUP BY b ORDER BY b", args, ).fetchall() conn.close() return [ {"ts": r[0], "cpu": round(r[1] or 0, 2), "mem": round(r[2] or 0, 2), "gpu": round(r[3] if r[3] is not None else -1, 2), "netIn": round(r[4] or 0, 1), "netOut": round(r[5] or 0, 1), "disk": round(r[6] or 0, 2)} for r in rows ] def record_app_checks(apps): now = int(time.time()) rows = [] for a in apps.values(): loc = a.get("local") or {} pub = a.get("public") or {} rows.append((now, a["app"], a["state"], loc.get("code", 0), loc.get("ms", 0), pub.get("code", 0), pub.get("ms", 0))) if not rows: return with db_lock: conn = db() conn.executemany("INSERT INTO app_checks VALUES (?,?,?,?,?,?,?)", rows) conn.execute("DELETE FROM app_checks WHERE ts < ?", (now - APP_RETENTION_D * 86400,)) conn.commit() conn.close() def query_uptime(window_s=86400): """% de vérifications non « down » par app sur la fenêtre + dernier changement.""" since = int(time.time()) - window_s with db_lock: conn = db() rows = conn.execute( "SELECT app, AVG(CASE WHEN state = 'down' THEN 0.0 ELSE 100.0 END), COUNT(*)," " AVG(CASE WHEN local_ms > 0 THEN local_ms END)" " FROM app_checks WHERE ts >= ? GROUP BY app", (since,)).fetchall() conn.close() return {r[0]: {"uptime": round(r[1] or 0, 2), "checks": r[2], "avgLocalMs": round(r[3] or 0)} for r in rows} def query_app_history(app, window_s, bucket_s): since = int(time.time()) - window_s with db_lock: conn = db() rows = conn.execute( f"SELECT (ts/{bucket_s})*{bucket_s} AS b," " AVG(CASE WHEN state = 'down' THEN 0.0 ELSE 100.0 END)," " AVG(CASE WHEN local_ms > 0 THEN local_ms END)," " AVG(CASE WHEN public_ms > 0 THEN public_ms END)," " SUM(state = 'down'), SUM(state = 'degraded'), COUNT(*)" " FROM app_checks WHERE app = ? AND ts >= ? GROUP BY b ORDER BY b", (app, since)).fetchall() conn.close() out = [] for r in rows: worst = "up" if r[4]: worst = "down" elif r[5]: worst = "degraded" out.append({"ts": r[0], "availability": round(r[1] or 0, 1), "localMs": round(r[2] or 0), "publicMs": round(r[3] or 0), "worst": worst, "checks": r[6]}) return out def record_event(kind, subject, from_state, to_state, detail=""): now = int(time.time()) with db_lock: conn = db() conn.execute("INSERT INTO events (ts, kind, subject, from_state, to_state, detail) VALUES (?,?,?,?,?,?)", (now, kind, subject, from_state or "", to_state or "", detail or "")) conn.execute("DELETE FROM events WHERE ts < ?", (now - EVENTS_RETENTION_D * 86400,)) conn.commit() conn.close() print(f"event {kind} {subject}: {from_state} -> {to_state} {detail}", flush=True) def query_events(since=0, limit=200, kind=None, subject=None): where = ["ts >= ?"] args = [int(since)] if kind: where.append("kind = ?") args.append(kind) if subject: where.append("subject = ?") args.append(subject) args.append(int(limit)) with db_lock: conn = db() rows = conn.execute( "SELECT id, ts, kind, subject, from_state, to_state, detail FROM events" f" WHERE {' AND '.join(where)} ORDER BY ts DESC, id DESC LIMIT ?", args).fetchall() conn.close() return [{"id": r[0], "ts": r[1], "kind": r[2], "subject": r[3], "from": r[4], "to": r[5], "detail": r[6]} for r in rows] # --------------------------------------------------------------------------- # Bus de ticks (flux SSE) # --------------------------------------------------------------------------- tick_cv = threading.Condition() tick_seq = 0 tick_last = {} def publish_tick(kind, payload): global tick_seq, tick_last with tick_cv: tick_seq += 1 last = {"seq": tick_seq, "kind": kind, "ts": time.time()} last.update(payload) tick_last = last tick_cv.notify_all() # --------------------------------------------------------------------------- # Boucle de collecte des métriques # --------------------------------------------------------------------------- latest = {} # name -> metrics dict latest_lock = threading.Lock() started_at = time.time() cycle_count = 0 def node_online(name): with latest_lock: m = latest.get(name) return bool(m) and m.get("status") not in ("offline", "unknown") def collect_cycle(): global cycle_count results = {} threads = [] lock = threading.Lock() def worker(node): name = node[0] prev = latest.get(name) out = run_on_node(name, metrics_script_for(name)) if "===METRICS_START===" in out: m = parse_metrics(out, node, prev) else: m = dict(prev) if prev else default_metrics(node) m["status"] = "offline" m["ts"] = time.time() with lock: results[name] = m for node in NODES: t = threading.Thread(target=worker, args=(node,)) t.start() threads.append(t) for t in threads: t.join() # Événements nœud : transitions en ligne ↔ hors ligne (pas au premier cycle) with latest_lock: previous = {k: v.get("status") for k, v in latest.items()} latest.update(results) if cycle_count > 0: for name, m in results.items(): was = previous.get(name) now_off = m["status"] == "offline" was_off = was in ("offline", None) if was is not None and now_off != was_off: record_event("node", name, "offline" if was_off else "online", "offline" if now_off else "online", "ne répond plus à l'agent" if now_off else "de nouveau joignable") record_samples(results) cycle_count += 1 online = sum(1 for m in results.values() if m.get("status") not in ("offline", "unknown")) srv_online = sum(1 for n in SERVER_NAMES if results.get(n, {}).get("status") not in ("offline", "unknown", None)) publish_tick("metrics", {"nodesOnline": online, "nodesTotal": len(NODES), "macsOnline": online - srv_online, "macsTotal": len(MAC_NAMES), "serversOnline": srv_online, "serversTotal": len(SERVER_NAMES)}) missing = [n for n, m in results.items() if m.get("status") == "offline" and n != SELF_NODE] if missing: threading.Thread(target=rescan_lan, args=(missing,), daemon=True).start() def collector_loop(): while True: t0 = time.time() try: collect_cycle() except Exception as e: print("collect error:", e, flush=True) time.sleep(max(5, COLLECT_INTERVAL - (time.time() - t0))) # --------------------------------------------------------------------------- # Registre maclustr-dispatch (apps déployées) # --------------------------------------------------------------------------- registry = {"updated": None, "apps": {}, "history": []} registry_lock = threading.Lock() registry_meta = {"source": "", "loadedAt": 0.0, "pulledAt": 0.0, "mtime": 0.0, "error": ""} def load_registry_file(): """Recharge ~/maclustr-agentd/registry.json si le fichier a changé (push mld).""" global registry try: st = os.stat(REGISTRY_PATH) except FileNotFoundError: return False if st.st_mtime <= registry_meta["mtime"]: return False try: with open(REGISTRY_PATH) as f: data = json.load(f) if not isinstance(data.get("apps"), dict): raise ValueError("registre sans clé apps") with registry_lock: registry = data registry_meta.update({"mtime": st.st_mtime, "loadedAt": time.time(), "source": "file", "error": ""}) print(f"registry loaded ({len(data['apps'])} apps, updated {data.get('updated')})", flush=True) return True except Exception as e: registry_meta["error"] = f"registre local illisible : {e}" print(registry_meta["error"], flush=True) return False def pull_registry(): """Tire le registre depuis la passerelle (SSH LAN) et l'écrit localement.""" out = run_on_node(GATEWAY_NODE, f"cat {GATEWAY_REGISTRY} 2>/dev/null", timeout=20) registry_meta["pulledAt"] = time.time() if not out.strip(): registry_meta["error"] = f"registre injoignable sur {GATEWAY_NODE}" return False try: data = json.loads(out) if not isinstance(data.get("apps"), dict): raise ValueError("clé apps absente") except Exception as e: registry_meta["error"] = f"registre passerelle invalide : {e}" return False tmp = REGISTRY_PATH + ".tmp" with open(tmp, "w") as f: json.dump(data, f, indent=1, ensure_ascii=False) os.replace(tmp, REGISTRY_PATH) loaded = load_registry_file() if loaded: registry_meta["source"] = "pull" return loaded def registry_apps(): with registry_lock: return {k: dict(v) for k, v in registry.get("apps", {}).items()} # --------------------------------------------------------------------------- # Santé des apps # --------------------------------------------------------------------------- class _NoRedirect(urllib.request.HTTPRedirectHandler): def redirect_request(self, req, fp, code, msg, headers, newurl): return None _ssl_ctx = ssl.create_default_context() _opener_local = urllib.request.build_opener(_NoRedirect) _opener_public = urllib.request.build_opener(_NoRedirect, urllib.request.HTTPSHandler(context=_ssl_ctx)) def http_probe(url, timeout, opener): """→ {"ok", "code", "ms", "error"} ; 2xx/3xx = ok, 4xx = répond (dégradé), 5xx/erreur = KO.""" t0 = time.time() req = urllib.request.Request(url, headers={"User-Agent": f"maclustr-agentd/{AGENT_VERSION}"}) try: with opener.open(req, timeout=timeout) as resp: code = resp.getcode() resp.read(2048) except urllib.error.HTTPError as e: code = e.code except (urllib.error.URLError, socket.timeout, ssl.SSLError, ConnectionError, OSError) as e: return {"ok": False, "code": 0, "ms": int((time.time() - t0) * 1000), "error": str(getattr(e, "reason", e))[:120]} except Exception as e: return {"ok": False, "code": 0, "ms": int((time.time() - t0) * 1000), "error": str(e)[:120]} ms = int((time.time() - t0) * 1000) return {"ok": 200 <= code < 400, "code": code, "ms": ms, "error": ""} def parse_procs_output(out): """Sortie de node_script → (liste pm2 normalisée, dict launchd label→{pid,exit}, dict app→sonde HTTP locale).""" pm2, launchd, http = [], {}, {} section = None pm2_raw = [] http_app = None for line in out.splitlines(): if line.startswith("---PM2---"): section = "pm2" continue if line.startswith("---LAUNCHD---"): section = "launchd" continue m = re.match(r"---HTTP (\S+)---$", line) if m: section = "http" http_app = m.group(1) continue if line.startswith("---END---"): break if section == "pm2": # sécurité : marqueur collé au JSON si le saut de ligne manque if "---LAUNCHD---" in line: pm2_raw.append(line.split("---LAUNCHD---")[0]) section = "launchd" else: pm2_raw.append(line) elif section == "launchd": p = line.split("\t") if len(p) >= 3 and p[2] != "Label": pid = int(p[0]) if p[0].isdigit() else None try: exit_code = int(p[1]) except ValueError: exit_code = None launchd[p[2]] = {"pid": pid, "exit": exit_code} elif section == "http" and http_app: p = line.split() if len(p) >= 2: try: code = int(p[0]) ms = int(float(p[1].replace(",", ".")) * 1000) except ValueError: code, ms = 0, 0 http[http_app] = {"ok": 200 <= code < 400, "code": code, "ms": ms, "error": "" if code else "connexion refusée ou délai dépassé (127.0.0.1)"} http_app = None raw = "\n".join(pm2_raw).strip() start = raw.find("[") if start >= 0: try: for p in json.loads(raw[start:]): env = p.get("pm2_env", {}) or {} mon = p.get("monit", {}) or {} uptime_ms = env.get("pm_uptime") or 0 pm2.append({ "name": p.get("name", ""), "pmId": p.get("pm_id", -1), "status": env.get("status", "unknown"), "restarts": env.get("restart_time", 0) or 0, "uptimeS": int(max(0, time.time() - uptime_ms / 1000.0)) if uptime_ms and env.get("status") == "online" else 0, "cpu": float(mon.get("cpu") or 0), "memMB": round(float(mon.get("memory") or 0) / 1048576.0, 1), "cron": bool(env.get("cron_restart")), "autorestart": env.get("autorestart", True) is not False, "outLog": env.get("pm_out_log_path", ""), "errLog": env.get("pm_err_log_path", ""), "pid": p.get("pid") or 0, }) except ValueError: pass return pm2, launchd, http # --------------------------------------------------------------------------- # MacLustr Tunnel : état des passerelles (pairs WireGuard, routes Caddy) + DNS des domaines # --------------------------------------------------------------------------- tunnel_latest = {} # gateway -> record tunnel_lock = threading.Lock() _dns_cache = {} # domain -> (ip, ts) def tunnel_fetch(name, gw): """Relevé d'une passerelle via `sudo tunnelctl json` (SSH direct, clé agentd).""" rec = {"name": name, "ip": gw["host"], "primary": gw.get("primary", False), "site": gw.get("site", ""), "ok": False, "error": None, "checkedAt": time.time(), "wg": {"listenPort": 0, "peers": []}, "routes": [], "caddyActive": False} try: attempts = [(gw["host"], None)] + SERVER_ALT.get(name, []) r = None for host, jump in attempts: r = subprocess.run(["ssh"] + SSH_OPTS + (["-o", "ProxyCommand=" + jump_command(jump)] if jump else []) + [f"{gw['user']}@{host}", "sudo tunnelctl json"], capture_output=True, text=True, timeout=25) if r.returncode == 0 and r.stdout.strip(): break if r is None or r.returncode != 0 or not r.stdout.strip(): rec["error"] = ((r.stderr.strip() if r else "") or "aucune sortie")[-160:] return rec d = json.loads(r.stdout) for p in d.get("wg", {}).get("peers", []): hs = p.get("handshakeS") p["online"] = hs is not None and hs <= TUNNEL_PEER_FRESH_S rec["wg"] = d.get("wg", rec["wg"]) rec["routes"] = d.get("caddy", {}).get("routes", []) rec["caddyActive"] = bool(d.get("caddy", {}).get("active")) rec["ok"] = rec["caddyActive"] except Exception as e: # noqa: BLE001 rec["error"] = str(e)[-160:] return rec def tunnel_cycle_run(): results = {} lock = threading.Lock() def worker(name, gw): r = tunnel_fetch(name, gw) with lock: results[name] = r threads = [threading.Thread(target=worker, args=(n, g)) for n, g in TUNNEL_GATEWAYS.items()] for t in threads: t.start() for t in threads: t.join() with tunnel_lock: tunnel_latest.clear() tunnel_latest.update(results) def tunnel_loop(): while True: try: tunnel_cycle_run() except Exception as e: # noqa: BLE001 print(f"tunnel cycle error: {e}", flush=True) time.sleep(TUNNEL_INTERVAL) def dns_a(domain): """Première adresse A du domaine (cache DNS_CACHE_S) ; None si non résolu.""" if not domain: return None now = time.time() hit = _dns_cache.get(domain) if hit and now - hit[1] < DNS_CACHE_S: return hit[0] ip = None try: infos = socket.getaddrinfo(domain, 443, socket.AF_INET, socket.SOCK_STREAM) ip = infos[0][4][0] if infos else None except Exception: # noqa: BLE001 ip = None _dns_cache[domain] = (ip, now) return ip def tunnel_info(domain): """Pour une app : passerelles qui routent son domaine, upstreams, résolution DNS et « via ».""" if not domain: return None gws, ups = [], [] with tunnel_lock: snap = {k: dict(v) for k, v in tunnel_latest.items()} for name, rec in snap.items(): for r in rec.get("routes", []): if r.get("domain") == domain and r.get("kind") == "proxy": gws.append(name) ups.extend(u.get("addr") if isinstance(u, dict) else str(u) for u in r.get("upstreams", [])) ip = dns_a(domain) via = None if ip: via = next((n for n, g in TUNNEL_GATEWAYS.items() if g["host"] == ip), None) or "hors tunnel" gateway = via if via in gws else (gws[0] if gws else None) return {"gateways": sorted(set(gws)), "upstreams": sorted(set(ups)), "dns": ip, "via": via, "gateway": gateway, "site": TUNNEL_GATEWAYS.get(gateway or "", {}).get("site")} def tunnel_snapshot(): with tunnel_lock: gws = [dict(v) for v in tunnel_latest.values()] gws.sort(key=lambda g: (not g.get("primary"), g["name"])) peers_by_node = {} for g in gws: for p in g.get("wg", {}).get("peers", []): peers_by_node.setdefault(p.get("alias"), {})[g["name"]] = { "ip": p.get("ip"), "handshakeS": p.get("handshakeS"), "online": p.get("online", False)} return {"gateways": gws, "peersByNode": peers_by_node, "routesTotal": sum(1 for g in gws for r in g.get("routes", []) if r.get("kind") == "proxy"), "ts": time.time()} # --------------------------------------------------------------------------- # Sites publics (3.0.1) : chaque domaine routé par la passerelle primaire est # sondé en HTTPS — y compris les apps hors registre mld (serveurs OVH). Un site # est « down » après SITES_FAIL_THRESHOLD échecs consécutifs (5xx ou injoignable). # --------------------------------------------------------------------------- SITES_INTERVAL = 120 SITES_FAIL_THRESHOLD = 2 sites_latest = {} # domain -> record sites_lock = threading.Lock() sites_cycle = 0 def sites_cycle_run(): global sites_cycle with tunnel_lock: snap = {k: dict(v) for k, v in tunnel_latest.items()} app_domains = {e.get("domain"): a for a, e in registry_apps().items() if e.get("domain")} targets = {} for gname, g in sorted(snap.items(), key=lambda kv: not kv[1].get("primary")): for r in g.get("routes", []): d = r.get("domain") if r.get("kind") != "proxy" or not d or d in targets: continue targets[d] = {"domain": d, "gateway": gname, "upstreams": [u.get("addr") if isinstance(u, dict) else str(u) for u in r.get("upstreams", [])], "app": app_domains.get(d)} if not targets: return results = {} lock = threading.Lock() sem = threading.Semaphore(8) def probe(t): with sem: r = http_probe("https://%s/" % t["domain"], HTTP_TIMEOUT_PUBLIC, _opener_public) with sites_lock: prev = sites_latest.get(t["domain"], {}) down = r["code"] == 0 or r["code"] >= 500 rec = dict(t) rec.update({"ok": not down, "responds": r["code"] > 0, "code": r["code"], "ms": r["ms"], "error": r["error"], "fails": (prev.get("fails", 0) + 1) if down else 0, "checkedAt": time.time(), "since": prev.get("since") if prev.get("ok") == (not down) else time.time()}) with lock: results[t["domain"]] = rec threads = [threading.Thread(target=probe, args=(t,)) for t in targets.values()] for th in threads: th.start() for th in threads: th.join() with sites_lock: prev_down = {d for d, r in sites_latest.items() if r.get("fails", 0) >= SITES_FAIL_THRESHOLD} sites_latest.clear() sites_latest.update(results) now_down = {d for d, r in results.items() if r.get("fails", 0) >= SITES_FAIL_THRESHOLD} if sites_cycle > 0: for d in sorted(now_down - prev_down): record_event("site", d, "up", "down", "HTTP %s — %s" % (results[d]["code"], results[d]["error"] or "5xx")) for d in sorted(prev_down - now_down): record_event("site", d, "down", "up", "de nouveau en ligne (HTTP %s)" % results[d]["code"]) sites_cycle += 1 def sites_loop(): time.sleep(20) # laisser le premier relevé tunnel arriver while True: try: sites_cycle_run() except Exception as e: # noqa: BLE001 print("sites error:", e, flush=True) time.sleep(SITES_INTERVAL) # ---- Sortie Internet du LAN (filtre du routeur Bell) ----------------------------------------------------- # 2026-10-01 : le Bell Giga Hub 2.0 (Sagemcom 5697, firmware 3.11.3) se met par moments à ne relayer que # TCP 80/443 et UDP 53 vers Internet : ICMP, SSH 22, NTP, STUN, tunnels… sont silencieusement perdus alors # que son propre ping passe (outil Utilitaires). Un redémarrage du Hub rétablit tout pour un temps. # Effets vus : IPTV « Tunnel Error », NAS UGREEN voyant « sans Internet », SSH vers OVH impossible, # WireGuard seulement via UDP 443. On sonde donc des ports non-web et on ouvre un incident quand le témoin # HTTPS passe mais que la majorité des autres échouent. EGRESS_INTERVAL = 120 EGRESS_FAIL_THRESHOLD = 2 # cycles consécutifs avant incident (≈ 4 min) EGRESS_CONTROL = ("tcp", "51.161.112.61", 443, "HTTPS BHS64 (témoin)") EGRESS_PROBES = [ ("tcp", "51.161.112.61", 22, "SSH BHS64"), ("tcp", "1.1.1.1", 853, "DNS/TLS Cloudflare"), ("icmp", "1.1.1.1", 0, "ping Cloudflare"), ("icmp", "8.8.8.8", 0, "ping Google"), ("ntp", "time.apple.com", 123, "NTP Apple"), ] egress_latest = {} egress_lock = threading.Lock() egress_cycle = 0 def _egress_probe(kind, host, port, timeout=3.0): t0 = time.time() try: if kind == "tcp": s = socket.create_connection((host, port), timeout=timeout) s.close() elif kind == "ntp": s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) s.settimeout(timeout) s.sendto(b"\x1b" + 47 * b"\0", (host, port)) s.recvfrom(48) s.close() elif kind == "icmp": r = subprocess.run(["/sbin/ping", "-c", "1", "-W", str(int(timeout * 1000)), host], stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, timeout=timeout + 2) if r.returncode != 0: return False, 0 else: return False, 0 return True, int((time.time() - t0) * 1000) except Exception: # noqa: BLE001 return False, 0 def egress_cycle_run(): global egress_cycle control_ok, control_ms = _egress_probe(*EGRESS_CONTROL[:3]) results = [] for kind, host, port, label in EGRESS_PROBES: ok, ms = _egress_probe(kind, host, port) results.append({"kind": kind, "host": host, "port": port, "label": label, "ok": ok, "ms": ms}) blocked = [r["label"] for r in results if not r["ok"]] filtered = bool(control_ok) and len(blocked) >= 4 with egress_lock: prev = dict(egress_latest) was = prev.get("filtered", False) rec = {"ok": not filtered, "filtered": filtered, "controlOk": control_ok, "controlMs": control_ms, "probes": results, "blocked": blocked, "fails": (prev.get("fails", 0) + 1) if filtered else 0, "checkedAt": time.time(), "since": prev.get("since") if (prev and was == filtered) else time.time()} egress_latest.clear() egress_latest.update(rec) if egress_cycle > 0 and was != filtered: if filtered: record_event("tunnel", "giga-hub", "open", "filtered", "le routeur Bell ne relaie plus que le web — bloqués : " + ", ".join(blocked)) else: record_event("tunnel", "giga-hub", "filtered", "open", "les ports non-web sortent de nouveau du LAN") egress_cycle += 1 def egress_loop(): time.sleep(12) while True: try: egress_cycle_run() except Exception as e: # noqa: BLE001 print("egress error:", e, flush=True) time.sleep(EGRESS_INTERVAL) def egress_snapshot(): with egress_lock: out = dict(egress_latest) out["cycle"] = egress_cycle out["ts"] = time.time() return out def egress_public(): with egress_lock: e = dict(egress_latest) return {"filtered": e.get("filtered", False), "controlOk": e.get("controlOk"), "blocked": e.get("blocked", []), "since": e.get("since"), "checkedAt": e.get("checkedAt")} def sites_snapshot(): with sites_lock: sites = [dict(v) for v in sites_latest.values()] sites.sort(key=lambda r: (r.get("ok", True), r["domain"])) return {"sites": sites, "total": len(sites), "ok": sum(1 for r in sites if r.get("ok")), "down": sum(1 for r in sites if r.get("fails", 0) >= SITES_FAIL_THRESHOLD), "cycle": sites_cycle, "ts": time.time()} apps_latest = {} # app -> record apps_lock = threading.Lock() apps_cycle = 0 node_procs_cache = {} # node -> {"pm2": [...], "launchd": {...}, "ts": float} def evaluate_app(entry, pm2_list, launchd_map, node_ok, local, public, tun=None): """Fusionne registre + processus + sondes HTTP (+ tunnel) en un enregistrement d'état.""" app = entry["app"] declared = list(entry.get("processes") or []) labels = list(entry.get("launchd") or []) by_name = {p["name"]: p for p in (pm2_list or [])} procs = [] for name in declared: p = by_name.get(name) if p: procs.append(dict(p)) else: procs.append({"name": name, "pmId": -1, "status": "absent" if pm2_list is not None else "unknown", "restarts": 0, "uptimeS": 0, "cpu": 0.0, "memMB": 0.0, "cron": False, "autorestart": True, "outLog": "", "errLog": "", "pid": 0}) # processus PM2 non déclarés mais préfixés du nom de l'app (workers ajoutés à la main) for name, p in by_name.items(): if name not in declared and (name == app or name.startswith(app + "-") or name.startswith(app.replace("-", "") + "-")): q = dict(p) q["undeclared"] = True procs.append(q) ld = [] for label in labels: info = launchd_map.get(label) if launchd_map is not None else None ld.append({"label": label, "pid": (info or {}).get("pid"), "exit": (info or {}).get("exit"), "running": bool(info and info.get("pid")), "known": info is not None}) reasons = [] if not node_ok: state = "down" reasons.append("nœud hors ligne") elif entry.get("port") and local is not None and not local.get("ok"): state = "down" reasons.append("HTTP local %s" % (local.get("code") or local.get("error") or "KO")) else: state = "up" bad = [p for p in procs if not p.get("cron") and p.get("autorestart", True) and p["status"] not in ("online", "launching", "unknown") and not p.get("undeclared")] if bad: state = "degraded" reasons.append("PM2 %s" % ", ".join("%s=%s" % (p["name"], p["status"]) for p in bad[:3])) errored = [p for p in procs if p["status"] == "errored"] if errored and state == "up": state = "degraded" dead_ld = [l for l in ld if l["known"] and not l["running"] and not l["label"].endswith("caffeinate")] if dead_ld: state = "degraded" reasons.append("launchd %s" % ", ".join(l["label"] for l in dead_ld[:3])) if entry.get("domain") and public is not None and not public.get("ok"): state = "degraded" reasons.append("public %s" % (public.get("code") or public.get("error") or "KO")) if tun and tun.get("gateways") and tun.get("dns") and tun.get("via") == "hors tunnel": state = "degraded" reasons.append("DNS %s ne pointe pas vers le MacLustr Tunnel" % tun["dns"]) if entry.get("domain") and tun is not None and not tun.get("gateways") and tunnel_latest: reasons.append("aucune route tunnel pour ce domaine") flappy = [p for p in procs if not p.get("cron") and p.get("restarts", 0) >= 20 and p.get("uptimeS", 0) < 600] if flappy: state = "degraded" reasons.append("redémarrages en boucle : %s" % ", ".join(p["name"] for p in flappy[:2])) if not entry.get("port") and not procs and not ld: state = "unknown" reasons.append("rien à sonder") rec = { "app": app, "label": entry.get("label") or app, "node": entry.get("node"), "ip": entry.get("ip"), "port": entry.get("port"), "domain": entry.get("domain"), "dir": entry.get("dir"), "healthPath": entry.get("health_path") or "/", "deployed": entry.get("deployed"), "registryStatus": entry.get("status"), "registryHealth": entry.get("health"), "processes": procs, "launchd": ld, "local": local, "public": public, "tunnel": tun, "state": state, "reason": " · ".join(reasons), "checkedAt": time.time(), "memMB": round(sum(p.get("memMB", 0) for p in procs), 1), "cpu": round(sum(p.get("cpu", 0) for p in procs), 1), } return rec def apps_cycle_run(): global apps_cycle load_registry_file() if time.time() - registry_meta["pulledAt"] > REGISTRY_PULL_INTERVAL and ( not registry_apps() or time.time() - registry_meta["loadedAt"] > REGISTRY_PULL_INTERVAL): pull_registry() entries = registry_apps() for app, e in entries.items(): e["app"] = app if not entries: publish_tick("apps", {"appsTotal": 0, "appsUp": 0}) return # 1. par nœud hébergeant ≥ 1 app : un seul SSH → PM2 + launchd + curl local by_node = {} for e in entries.values(): if e.get("node"): by_node.setdefault(e["node"], []).append(e) procs = {} lock = threading.Lock() def procs_worker(node, node_apps): if node not in NODE_INDEX or not node_online(node): with lock: procs[node] = None return out = run_on_node(node, node_script(node_apps), timeout=45) if "---END---" not in out: with lock: procs[node] = None return pm2, ld, http = parse_procs_output(out) with lock: procs[node] = {"pm2": pm2, "launchd": ld, "http": http, "ts": time.time()} node_procs_cache[node] = procs[node] threads = [threading.Thread(target=procs_worker, args=(n, a)) for n, a in by_node.items()] for t in threads: t.start() # 2. sondes HTTP publiques en parallèle (pendant les SSH) probes = {} sem = threading.Semaphore(24) def probe_worker(app, e): pub = None if e.get("domain"): with sem: pub = http_probe("https://%s%s" % (e["domain"], e.get("health_path") or "/"), HTTP_TIMEOUT_PUBLIC, _opener_public) with lock: probes[app] = pub pthreads = [threading.Thread(target=probe_worker, args=(a, e)) for a, e in entries.items()] for t in pthreads: t.start() for t in threads + pthreads: t.join() # 3. évaluation + événements results = {} for app, e in entries.items(): node = e.get("node") pinfo = procs.get(node) pub = probes.get(app) loc = None if e.get("port"): if pinfo and app in pinfo.get("http", {}): loc = pinfo["http"][app] elif node_online(node): # repli : sonde directe par le LAN (nœud joignable mais script en échec) host = "127.0.0.1" if node == SELF_NODE else (e.get("ip") or lan_map.get(node or "")) loc = http_probe("http://%s:%s%s" % (host, e["port"], e.get("health_path") or "/"), HTTP_TIMEOUT_LOCAL, _opener_local) if host else \ {"ok": False, "code": 0, "ms": 0, "error": "IP inconnue"} else: loc = {"ok": False, "code": 0, "ms": 0, "error": "nœud hors ligne"} results[app] = evaluate_app( e, pinfo["pm2"] if pinfo else None, pinfo["launchd"] if pinfo else None, node_online(node), loc, pub, tunnel_info(e.get("domain")), ) uptimes = query_uptime(86400) with apps_lock: previous = {k: v.get("state") for k, v in apps_latest.items()} apps_latest.clear() apps_latest.update(results) for app, rec in apps_latest.items(): u = uptimes.get(app) rec["uptime24h"] = u["uptime"] if u else None rec["avgLocalMs"] = u["avgLocalMs"] if u else None if apps_cycle > 0: for app, rec in results.items(): was = previous.get(app) if was is not None and was != rec["state"]: record_event("app", app, was, rec["state"], rec.get("reason") or "") for app in previous: if app not in results: record_event("app", app, previous[app], "removed", "retirée du registre") for app in results: if app not in previous: record_event("app", app, "", results[app]["state"], "nouvelle app dans le registre (%s)" % results[app].get("node")) record_app_checks(results) apps_cycle += 1 s = apps_summary() publish_tick("apps", {"appsTotal": s["total"], "appsUp": s["up"], "appsDown": s["down"], "appsDegraded": s["degraded"]}) def apps_loop(): # laisse le premier cycle métriques établir qui est en ligne time.sleep(8) while True: t0 = time.time() try: apps_cycle_run() except Exception as e: print("apps cycle error:", repr(e), flush=True) time.sleep(max(5, APPS_INTERVAL - (time.time() - t0))) def apps_summary(): with apps_lock: states = [a["state"] for a in apps_latest.values()] return {"total": len(states), "up": states.count("up"), "degraded": states.count("degraded"), "down": states.count("down"), "unknown": states.count("unknown"), "cycle": apps_cycle} def apps_by_node(): out = {} with apps_lock: for a in apps_latest.values(): n = a.get("node") or "?" d = out.setdefault(n, {"total": 0, "up": 0, "degraded": 0, "down": 0, "unknown": 0, "apps": []}) d["total"] += 1 d[a["state"]] = d.get(a["state"], 0) + 1 d["apps"].append(a["app"]) return out def app_logs(app, lines=120, process=None): with apps_lock: rec = apps_latest.get(app) if not rec: return None node = rec.get("node") lines = max(10, min(int(lines), 600)) parts = [] targets = [p for p in rec.get("processes", []) if not process or p["name"] == process] for p in targets: for kind, path in (("out", p.get("outLog")), ("err", p.get("errLog"))): if path: parts.append((p["name"], kind, path)) for l in rec.get("launchd", []): if process and l["label"] != process: continue parts.append((l["label"], "launchd", "__LAUNCHD__" + l["label"])) if not parts: return {"app": app, "node": node, "logs": []} script = ["export PATH=/opt/homebrew/bin:/usr/local/bin:$PATH"] for name, kind, path in parts: marker = "===LOG %s %s===" % (name, kind) script.append("echo %s" % shlex.quote(marker)) if path.startswith("__LAUNCHD__"): label = path[len("__LAUNCHD__"):] script.append( "for f in $(launchctl print gui/$(id -u)/%s 2>/dev/null | awk '/stdout path|stderr path/{print $NF}' | sort -u); do " "echo \"# $f\"; tail -n %d \"$f\" 2>/dev/null; done" % (shlex.quote(label), lines)) else: script.append("tail -n %d %s 2>/dev/null" % (lines, shlex.quote(path))) script.append("echo '===LOG END==='") out = run_on_node(node, "\n".join(script), timeout=30) logs = [] current = None for line in out.splitlines(): m = re.match(r"===LOG (\S+) (\S+)===$", line) if m: current = {"process": m.group(1), "kind": m.group(2), "text": []} logs.append(current) continue if line.startswith("===LOG END==="): break if current is not None: current["text"].append(line) for l in logs: l["text"] = "\n".join(l["text"])[-60000:] return {"app": app, "node": node, "logs": logs, "ts": time.time()} def app_action(app, action, process=None): with apps_lock: rec = apps_latest.get(app) if not rec: return {"ok": False, "error": "app inconnue"} node = rec.get("node") if not node_online(node): return {"ok": False, "error": "nœud %s hors ligne" % node} pm2_names = [p["name"] for p in rec.get("processes", []) if p.get("pmId", -1) >= 0 or p["status"] != "absent"] labels = [l["label"] for l in rec.get("launchd", [])] if process: if process in pm2_names: pm2_names, labels = [process], [] elif process in labels: pm2_names, labels = [], [process] else: return {"ok": False, "error": "processus %s inconnu pour %s" % (process, app)} if action not in ("restart", "stop", "start", "reload"): return {"ok": False, "error": "action inconnue"} cmds = [] if pm2_names: cmds.append("pm2 %s %s 2>&1 | grep -vE '^\\s*$' | tail -20" % (action, " ".join(shlex.quote(n) for n in pm2_names))) if labels: if action in ("restart", "reload", "start"): for l in labels: cmds.append("launchctl kickstart -k gui/$(id -u)/%s 2>&1 && echo 'launchd %s: kickstart ok'" % (shlex.quote(l), l)) else: for l in labels: cmds.append("launchctl kill TERM gui/$(id -u)/%s 2>&1 && echo 'launchd %s: TERM envoyé (KeepAlive le relancera)'" % (shlex.quote(l), l)) if not cmds: return {"ok": False, "error": "aucun processus à piloter"} script = PM2_ENV + "\n" + "\n".join(cmds) + "\necho ACTION_DONE" out = run_on_node(node, script, timeout=60) ok = "ACTION_DONE" in out record_event("action", app, "", action, "%s%s par l'app MacLustr" % (action, " " + process if process else "")) threading.Thread(target=_recheck_soon, daemon=True).start() return {"ok": ok, "output": out.replace("ACTION_DONE", "").strip()[-1500:], "node": node} def _recheck_soon(): time.sleep(4) try: apps_cycle_run() except Exception as e: print("recheck error:", e, flush=True) # --------------------------------------------------------------------------- # Actions (Node Doctor) # --------------------------------------------------------------------------- def do_action(name, action, pid=None): pw = shlex.quote(SUDO_PW) if platform_of(name) == "linux": # serveurs OVH : sudo sans mot de passe if action == "kill" and pid: script = f"kill -TERM {int(pid)} 2>/dev/null; sleep 2; kill -0 {int(pid)} 2>/dev/null && kill -KILL {int(pid)}; echo DONE" elif action == "purge": script = "sync; echo 3 | sudo -n tee /proc/sys/vm/drop_caches >/dev/null 2>&1 && echo DONE" elif action == "caches": 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" elif action == "reboot": script = "(sleep 1; sudo -n shutdown -r now) >/dev/null 2>&1 & echo DONE" else: return {"ok": False, "error": "unknown action"} out = run_on_node(name, script, timeout=60) record_event("action", name, "", action, "%s par l'app MacLustr" % action) return {"ok": "DONE" in out, "output": out.strip()[-500:]} if action == "kill" and pid: script = f"kill -TERM {int(pid)} 2>/dev/null; sleep 2; kill -0 {int(pid)} 2>/dev/null && kill -KILL {int(pid)}; echo DONE" elif action == "purge": script = f"printf '%s' {pw} | sudo -S purge 2>/dev/null && echo DONE" elif action == "caches": script = "rm -rf ~/Library/Caches/* /tmp/*.tmp 2>/dev/null; echo DONE" elif action == "reboot": script = f"printf '%s' {pw} | sudo -S shutdown -r now 2>/dev/null & echo DONE" else: return {"ok": False, "error": "unknown action"} out = run_on_node(name, script, timeout=30) record_event("action", name, "", action, "%s par l'app MacLustr" % action) return {"ok": "DONE" in out, "output": out.strip()[-500:]} # --------------------------------------------------------------------------- # NAS # --------------------------------------------------------------------------- def nas_status(): out = [] for nas in NAS_DEVICES: online = ping(nas["host"]) used = total = avail = 0.0 mp = os.path.join(HOME, nas["mount"]) if os.path.ismount(mp): df = run_local(f"df -g {shlex.quote(mp)} | tail -1") p = df.split() if len(p) >= 4: try: total, used, avail = float(p[1]), float(p[2]), float(p[3]) except ValueError: pass out.append({**{k: nas[k] for k in ("id", "name", "host")}, "online": online, "usedGB": used, "totalGB": total, "availGB": avail}) return out # --------------------------------------------------------------------------- # HTTP # --------------------------------------------------------------------------- NODE_INFO = [ {"name": n[0], "hostname": n[1], "chip": n[2], "model": n[3], "tier": n[4], "generation": n[5], "cpuCores": n[6], "memoryMB": n[7], "gpuCores": n[8], "remote": n[0] in REMOTE_NODES or n[0] in SERVER_META, "user": SERVER_META.get(n[0], {}).get("user") or ssh_user_for(REMOTE_NODES.get(n[0], {}).get("host", "")), "platform": platform_of(n[0]), "kind": "server" if n[0] in SERVER_META else "mac", "site": SERVER_META.get(n[0], {}).get("site", "Saint-Augustin-de-Desmaures (Québec)"), "role": SERVER_META.get(n[0], {}).get("role", ""), "hardware": SERVER_META.get(n[0], {}).get("hardware", ""), "publicIP": SERVER_META.get(n[0], {}).get("publicIP", "")} for n in NODES ] # --------------------------------------------------------------------------- # Centre d'incidents (3.0.0) — vue consolidée « qu'est-ce qui ne va pas ? » # calculée toutes les INCIDENTS_INTERVAL s à partir des nœuds, serveurs, apps, # passerelles tunnel et coordinateur mobile. Accusés de réception persistés. # --------------------------------------------------------------------------- INCIDENTS_INTERVAL = 15 incidents_lock = threading.Lock() incidents_latest = [] # liste triée incidents_counts = {"critical": 0, "warning": 0, "info": 0, "total": 0, "acked": 0} incident_first_seen = {} # id -> ts incident_acks = {} # id -> {"ts": float, "note": str} incidents_cycle = 0 SEV_ORDER = {"critical": 0, "warning": 1, "info": 2} def _acks_table(conn): conn.execute("CREATE TABLE IF NOT EXISTS incident_acks (id TEXT PRIMARY KEY, ts INTEGER, note TEXT)") def load_acks(): with db_lock: conn = db() _acks_table(conn) rows = conn.execute("SELECT id, ts, note FROM incident_acks").fetchall() conn.close() for r in rows: incident_acks[r[0]] = {"ts": r[1], "note": r[2] or ""} def save_ack(iid, note=None, remove=False): with db_lock: conn = db() _acks_table(conn) if remove: conn.execute("DELETE FROM incident_acks WHERE id = ?", (iid,)) else: conn.execute("INSERT OR REPLACE INTO incident_acks (id, ts, note) VALUES (?,?,?)", (iid, int(time.time()), note or "")) conn.commit() conn.close() def compute_incidents(): """Construit la liste des incidents ouverts (sans état : recalcul complet).""" now = time.time() out = [] def add(sev, kind, subject, code, title, detail="", actions=None, meta=None): out.append({"id": f"{kind}:{subject}:{code}", "severity": sev, "kind": kind, "subject": subject, "code": code, "title": title, "detail": detail, "actions": actions or [], "meta": meta or {}}) with latest_lock: snap = {k: dict(v) for k, v in latest.items()} for name, m in snap.items(): kind = "server" if name in SERVER_META else "node" st = m.get("status") if st in ("offline", "unknown"): add("critical", kind, name, "offline", f"{name} hors ligne", "ne répond plus à l'agent", ["mld-heal"] if kind == "node" else [], {"site": SERVER_META.get(name, {}).get("site")}) continue mem_pct = m["memUsedMB"] / m["memTotalMB"] * 100 if m.get("memTotalMB") else 0 if st == "critical": add("warning", kind, name, "load", f"{name} saturé", f"CPU {m.get('cpu', 0):.0f} % · mémoire {mem_pct:.0f} %", ["top", "purge"], {"cpu": round(m.get("cpu", 0), 1), "mem": round(mem_pct, 1)}) if m.get("diskTotalGB"): dpct = m["diskUsedGB"] / m["diskTotalGB"] * 100 free_gb = m["diskTotalGB"] - m["diskUsedGB"] if dpct >= 95: add("critical", kind, name, "disk-full", f"Disque presque plein sur {name}", f"{dpct:.0f} % utilisés · {free_gb:.0f} Go libres", ["caches"], {"diskPct": round(dpct, 1), "freeGB": round(free_gb)}) elif dpct >= 88: add("warning", kind, name, "disk-high", f"Disque chargé sur {name}", f"{dpct:.0f} % utilisés · {free_gb:.0f} Go libres", ["caches"], {"diskPct": round(dpct, 1), "freeGB": round(free_gb)}) if m.get("thermal") == "serious": add("warning", kind, name, "thermal", f"{name} bride son CPU (thermique)", "CPU_Scheduler_Limit < 70 %") if m.get("swapTotalMB", 0) > 0 and m["swapUsedMB"] / m["swapTotalMB"] > 0.9 and mem_pct > 85: add("warning", kind, name, "swap", f"{name} swappe", f"swap {m['swapUsedMB']:.0f}/{m['swapTotalMB']:.0f} Mo · mémoire {mem_pct:.0f} %") batt = m.get("battery", -1) if batt is not None and 0 <= batt < 20 and (m.get("batteryState") or "").startswith("discharging"): add("warning", kind, name, "battery", f"{name} sur batterie ({batt:.0f} %)", "portable débranché") with apps_lock: apps = [dict(a) for a in apps_latest.values()] for a in apps: if a.get("state") == "down": add("critical", "app", a["app"], "down", f"{a.get('label') or a['app']} est hors service", a.get("reason") or "sonde HTTP locale en échec", ["restart", "logs"], {"node": a.get("node"), "domain": a.get("domain"), "port": a.get("port")}) elif a.get("state") == "degraded": add("warning", "app", a["app"], "degraded", f"{a.get('label') or a['app']} est dégradée", a.get("reason") or "", ["restart", "logs"], {"node": a.get("node"), "domain": a.get("domain"), "port": a.get("port")}) with tunnel_lock: gws = {k: dict(v) for k, v in tunnel_latest.items()} for name, g in gws.items(): if not g.get("ok"): add("critical" if g.get("primary") else "warning", "tunnel", name, "gateway", f"Passerelle {name} injoignable", g.get("error") or "Caddy inactif", [], {"site": g.get("site")}) primary = next((g for g in gws.values() if g.get("primary")), None) if primary and primary.get("ok"): hosting = set() for a in apps: if a.get("domain") and a.get("node"): hosting.add(a["node"]) for p in primary.get("wg", {}).get("peers", []): alias = p.get("alias") if alias in hosting and not p.get("online") and snap.get(alias, {}).get("status") not in ("offline", "unknown", None): add("warning", "tunnel", alias, "wg-offline", f"Tunnel wg1 de {alias} silencieux", "dernier handshake il y a %s s — les sites publics de ce nœud ne répondent plus" % p.get("handshakeS"), ["mld-heal"]) with sites_lock: sites = [dict(v) for v in sites_latest.values()] for s_ in sites: if s_.get("fails", 0) >= SITES_FAIL_THRESHOLD and not s_.get("app"): add("critical", "site", s_["domain"], "down", f"{s_['domain']} ne répond plus", "HTTP %s via %s → %s%s" % (s_.get("code"), s_.get("gateway"), ", ".join(s_.get("upstreams") or []), (" — " + s_["error"]) if s_.get("error") else ""), ["open"], {"domain": s_["domain"], "site": s_.get("gateway")}) with mobile_lock: mob = dict(mobile_latest) mob_metrics = dict(mobile_latest.get("metrics") or {}) for name, m in mob_metrics.items(): mm = m.get("mobile") or {} if m.get("status") in ("offline", "unknown"): if mm.get("pinned"): add("warning", "mobile", name, "offline", f"{name} (mobile épinglé) hors ligne", "app MacLustr fermée ou en arrière-plan depuis %s s" % int(mm.get("ageS") or 0), ["open"], {"site": "mobile"}) continue if 0 <= m.get("battery", -1) < 20 and (m.get("batteryState") == "discharging"): add("warning", "mobile", name, "battery", f"{name} : batterie faible ({m['battery']:.0f} %)", "appareil non branché", [], {}) if m.get("thermal") == "serious": add("warning", "mobile", name, "thermal", f"{name} chauffe", "état thermique sérieux — jobs ralentis", [], {}) if mob.get("ts") and not mob.get("ok"): add("info", "mobile", "coordinateur", "down", "Coordinateur mobile injoignable", mob.get("error") or "", ["restart"]) eg = egress_snapshot() if eg.get("filtered") and eg.get("fails", 0) >= EGRESS_FAIL_THRESHOLD: add("critical", "tunnel", "giga-hub", "egress-filter", "Routeur Bell : seuls les ports web sortent du LAN", "bloqués : %s — HTTPS passe. Effets : IPTV « Tunnel Error », NAS UGREEN sans Internet, SSH vers OVH coupé " "(le tunnel UDP 443 tient). Remède connu : redémarrer le Giga Hub (192.168.2.1 → Réinitialisation)." % ", ".join(eg.get("blocked") or []), ["open"], {"site": "LAN", "blocked": eg.get("blocked") or []}) if registry_meta.get("error"): add("info", "registry", "mld", "pull", "Registre mld : dernier tirage en échec", registry_meta.get("error", "")[:160], ["registry-refresh"]) active = {i["id"] for i in out} for i in out: i["since"] = incident_first_seen.setdefault(i["id"], now) i["ageS"] = int(now - i["since"]) ack = incident_acks.get(i["id"]) i["acked"] = bool(ack) i["ackTs"] = ack["ts"] if ack else None i["ackNote"] = ack["note"] if ack else "" for k in list(incident_first_seen): if k not in active: incident_first_seen.pop(k, None) if k in incident_acks: incident_acks.pop(k, None) save_ack(k, remove=True) out.sort(key=lambda i: (i["acked"], SEV_ORDER.get(i["severity"], 9), i["since"])) return out def incidents_cycle_run(): global incidents_cycle, incidents_counts new = compute_incidents() with incidents_lock: old_ids = {i["id"]: i for i in incidents_latest} new_ids = {i["id"]: i for i in new} if incidents_cycle > 0: for iid, i in new_ids.items(): if iid not in old_ids: record_event("incident", iid, "", "open", "%s — %s" % (i["title"], i["detail"])) for iid, i in old_ids.items(): if iid not in new_ids: record_event("incident", iid, "open", "closed", "%s résolu" % i["title"]) counts = {"critical": 0, "warning": 0, "info": 0, "total": len(new), "acked": 0} for i in new: counts[i["severity"]] = counts.get(i["severity"], 0) + 1 if i["acked"]: counts["acked"] += 1 changed = counts != incidents_counts or set(new_ids) != set(old_ids) with incidents_lock: incidents_latest[:] = new incidents_counts = counts incidents_cycle += 1 if changed: publish_tick("incidents", {"incidents": counts}) def incidents_loop(): time.sleep(8) while True: try: incidents_cycle_run() except Exception as e: # noqa: BLE001 print("incidents error:", e, flush=True) time.sleep(INCIDENTS_INTERVAL) def incidents_snapshot(): with incidents_lock: return {"incidents": [dict(i) for i in incidents_latest], "counts": dict(incidents_counts), "cycle": incidents_cycle, "ts": time.time()} # --------------------------------------------------------------------------- # Ops (3.0.0) : commandes mld exécutées sur la passerelle M1M32 (liste blanche) # --------------------------------------------------------------------------- MLD_COMMANDS = { "status": {"args": ["status"], "label": "État des apps (registre)", "mutating": False}, "status-live": {"args": ["status", "--live"], "label": "État des apps (sondes live)", "mutating": False}, "nodes": {"args": ["nodes"], "label": "Ressources des nœuds", "mutating": False}, "scan": {"args": ["scan"], "label": "Scan live des nœuds", "mutating": False}, "discover": {"args": ["discover"], "label": "Découverte LAN", "mutating": False}, "apps": {"args": ["apps"], "label": "Manifestes", "mutating": False}, "plan": {"args": ["plan"], "label": "Plan de placement", "mutating": False}, "tunnel-status": {"args": ["tunnel", "status"], "label": "État du MacLustr Tunnel", "mutating": False}, "heal-dry": {"args": ["heal", "--dry-run"], "label": "Auto-réparation (simulation)", "mutating": False}, "heal": {"args": ["heal"], "label": "Auto-réparation", "mutating": True}, } _ANSI = re.compile(r"\x1b\[[0-9;]*[A-Za-z]") def run_mld(key): spec = MLD_COMMANDS.get(key) if not spec: return {"ok": False, "error": "commande inconnue", "allowed": sorted(MLD_COMMANDS)} ip = lan_map.get(GATEWAY_NODE) if not ip: return {"ok": False, "error": "passerelle %s introuvable" % GATEWAY_NODE} t0 = time.time() cmd = "~/maclustr-dispatch/bin/mld " + " ".join(shlex.quote(a) for a in spec["args"]) try: r = subprocess.run(["ssh"] + SSH_OPTS + [f"{SSH_USER}@{ip}", cmd], capture_output=True, text=True, timeout=240) out = _ANSI.sub("", (r.stdout or "") + (("\n" + r.stderr) if r.stderr.strip() else "")) ok = r.returncode == 0 except subprocess.TimeoutExpired: out, ok = "délai dépassé (240 s)", False except Exception as e: # noqa: BLE001 out, ok = str(e), False if spec["mutating"]: record_event("action", "mld", "", key, "mld %s lancé depuis l'app MacLustr" % " ".join(spec["args"])) threading.Thread(target=_recheck_soon, daemon=True).start() return {"ok": ok, "command": "mld " + " ".join(spec["args"]), "key": key, "output": out.strip()[-20000:], "ms": int((time.time() - t0) * 1000), "gateway": GATEWAY_NODE, "ts": time.time()} # --------------------------------------------------------------------------- # Résumé (3.0.0) : un seul appel pour un tableau de bord / widget # --------------------------------------------------------------------------- def cluster_aggregates(snap): macs = {k: v for k, v in snap.items() if k not in SERVER_META} servers = {k: v for k, v in snap.items() if k in SERVER_META} def agg(group): online = [m for m in group.values() if m.get("status") not in ("offline", "unknown", None)] cores = {n[0]: n[6] for n in NODES} wsum = sum(cores.get(m["name"], 1) for m in online) or 1 cpu = sum(m.get("cpu", 0) * cores.get(m["name"], 1) for m in online) / wsum if online else 0 mem_t = sum(m.get("memTotalMB", 0) for m in online) or 1 mem_u = sum(m.get("memUsedMB", 0) for m in online) disk_t = sum(m.get("diskTotalGB", 0) for m in online) disk_u = sum(m.get("diskUsedGB", 0) for m in online) gpus = [m["gpu"] for m in online if m.get("gpu", -1) >= 0] return {"online": len(online), "total": len(group), "cpu": round(cpu, 1), "memPct": round(mem_u / mem_t * 100, 1), "memUsedGB": round(mem_u / 1024, 1), "memTotalGB": round(mem_t / 1024, 1), "diskUsedGB": round(disk_u), "diskTotalGB": round(disk_t), "gpu": round(sum(gpus) / len(gpus), 1) if gpus else -1, "netInKBs": round(sum(m.get("netInKBs", 0) for m in online), 1), "netOutKBs": round(sum(m.get("netOutKBs", 0) for m in online), 1), "cores": sum(cores.get(m["name"], 0) for m in online), "coresTotal": sum(cores.get(n, 0) for n in group)} with mobile_lock: mob_metrics = dict(mobile_latest.get("metrics") or {}) mob_nodes = list(mobile_latest.get("nodes") or []) mob_cores = {n["name"]: n.get("cpuCores") or 0 for n in mob_nodes} mob_online = [m for m in mob_metrics.values() if m.get("status") not in ("offline", "unknown", None)] mem_t = sum(m.get("memTotalMB", 0) for m in mob_online) or 1 mobiles = {"online": len(mob_online), "total": len(mob_metrics), "cpu": round(sum(m.get("cpu", 0) for m in mob_online) / len(mob_online), 1) if mob_online else 0, "memPct": round(sum(m.get("memUsedMB", 0) for m in mob_online) / mem_t * 100, 1) if mob_online else 0, "memUsedGB": round(sum(m.get("memUsedMB", 0) for m in mob_online) / 1024, 1), "memTotalGB": round(sum(m.get("memTotalMB", 0) for m in mob_online) / 1024, 1), "diskUsedGB": round(sum(m.get("diskUsedGB", 0) for m in mob_online)), "diskTotalGB": round(sum(m.get("diskTotalGB", 0) for m in mob_online)), "gpu": -1, "netInKBs": 0, "netOutKBs": 0, "cores": sum(mob_cores.get(m["name"], 0) for m in mob_online), "coresTotal": sum(mob_cores.values()), "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, "jobsActive": sum((m.get("mobile") or {}).get("activeJobs", 0) for m in mob_online)} return {"macs": agg(macs), "servers": agg(servers), "mobiles": mobiles, "all": agg(snap)} def summary_snapshot(): with latest_lock: snap = {k: dict(v) for k, v in latest.items()} aggs = cluster_aggregates(snap) apps = apps_summary() tun = tunnel_snapshot() gws = tun.get("gateways", []) peers = [p for g in gws if g.get("primary") for p in g.get("wg", {}).get("peers", [])] with mobile_lock: mob = dict(mobile_latest) inc = incidents_snapshot() crit = inc["counts"].get("critical", 0) warn = inc["counts"].get("warning", 0) health = "critical" if crit else ("warning" if warn else "ok") score = max(0, 100 - crit * 15 - warn * 5) return { "health": health, "score": score, "agentVersion": AGENT_VERSION, "ts": time.time(), "uptime": int(time.time() - started_at), "nodes": aggs["all"], "macs": aggs["macs"], "servers": aggs["servers"], "mobiles": aggs["mobiles"], "apps": apps, "tunnel": {"gatewaysOk": sum(1 for g in gws if g.get("ok")), "gatewaysTotal": len(gws), "peersOnline": sum(1 for p in peers if p.get("online")), "peersTotal": len(peers), "routes": tun.get("routesTotal", 0)}, "mobile": {"ok": mob.get("ok", False), "online": mob.get("online", 0), "total": mob.get("total", 0), "queued": mob.get("queued", 0), "active": mob.get("active", 0)}, "incidents": inc["counts"], "topIncidents": inc["incidents"][:6], "sites": {k: v for k, v in sites_snapshot().items() if k in ("total", "ok", "down")}, "egress": egress_public(), "events": query_events(0, 8), } def public_record(rec): """Copie d'un enregistrement d'app sans les champs internes (chemins de logs), sans toucher aux dicts stockés dans apps_latest.""" out = dict(rec) out["processes"] = [{k: v for k, v in p.items() if k not in ("outLog", "errLog")} for p in rec.get("processes", [])] return out def public_apps_snapshot(): with apps_lock: apps = [public_record(a) for a in apps_latest.values()] apps.sort(key=lambda a: ({"down": 0, "degraded": 1, "unknown": 2, "up": 3}[a["state"]], a["app"])) return apps class Handler(BaseHTTPRequestHandler): server_version = f"maclustr-agentd/{AGENT_VERSION}" protocol_version = "HTTP/1.1" def _send(self, code, payload): body = json.dumps(payload, ensure_ascii=False).encode() self.send_response(code) self.send_header("Content-Type", "application/json; charset=utf-8") self.send_header("Content-Length", str(len(body))) self.send_header("Cache-Control", "no-store") self.end_headers() self.wfile.write(body) def _auth(self): h = self.headers.get("Authorization", "") if h == f"Bearer {TOKEN}": return True q = parse_qs(urlparse(self.path).query) if (q.get("token") or [""])[0] == TOKEN: return True self._send(401, {"error": "unauthorized"}) return False def log_message(self, fmt, *args): pass def _read_json(self): try: length = int(self.headers.get("Content-Length", 0)) return json.loads(self.rfile.read(length) or b"{}") except Exception: return None # -- SSE --------------------------------------------------------------- def _stream(self): self.send_response(200) self.send_header("Content-Type", "text/event-stream; charset=utf-8") self.send_header("Cache-Control", "no-store") self.send_header("Connection", "keep-alive") self.send_header("X-Accel-Buffering", "no") self.end_headers() try: self.wfile.write(b": maclustr-agentd stream\n\n") self.wfile.write(("event: hello\ndata: %s\n\n" % json.dumps( {"version": AGENT_VERSION, "seq": tick_seq, "ts": time.time()})).encode()) self.wfile.flush() seen = tick_seq while True: with tick_cv: tick_cv.wait(timeout=15) seq, last = tick_seq, dict(tick_last) if seq != seen: seen = seq self.wfile.write(("event: tick\ndata: %s\n\n" % json.dumps(last)).encode()) else: self.wfile.write(b": keepalive\n\n") self.wfile.flush() except (BrokenPipeError, ConnectionResetError, OSError): return # -- GET ----------------------------------------------------------------- def do_GET(self): u = urlparse(self.path) parts = [p for p in u.path.split("/") if p] q = parse_qs(u.query) if u.path == "/health": with latest_lock: online = sum(1 for m in latest.values() if m.get("status") not in ("offline", "unknown", None)) s = apps_summary() return self._send(200, {"ok": True, "version": AGENT_VERSION, "uptime": int(time.time() - started_at), "cycles": cycle_count, "appsCycles": apps_cycle, "nodesOnline": online, "nodesTotal": len(NODES), "appsUp": s["up"], "appsTotal": s["total"], "appsDown": s["down"], "appsDegraded": s["degraded"], "registryUpdated": registry.get("updated"), "tunnelGateways": {k: v.get("ok") for k, v in tunnel_latest.items()}, "mobileOnline": mobile_latest.get("online", 0), "mobileTotal": mobile_latest.get("total", 0), "mobileQueued": mobile_latest.get("queued", 0), "serversOnline": sum(1 for n in SERVER_NAMES if latest.get(n, {}).get("status") not in ("offline", "unknown", None)), "serversTotal": len(SERVER_NAMES), "macsTotal": len(MAC_NAMES), "incidents": dict(incidents_counts)}) if not self._auth(): return if u.path == "/api/stream": return self._stream() if u.path == "/api/cluster": with latest_lock: snap = {k: {kk: vv for kk, vv in v.items() if not kk.startswith("_")} for k, v in latest.items()} with mobile_lock: mobile = dict(mobile_latest) mob_nodes = mobile.pop("nodes", []) or [] mob_metrics = mobile.pop("metrics", {}) or {} snap.update(mob_metrics) return self._send(200, {"nodes": NODE_INFO + mob_nodes, "metrics": snap, "lanMap": lan_map, "ts": time.time(), "apps": apps_summary(), "appsByNode": apps_by_node(), "mobile": mobile, "servers": SERVER_NAMES, "mobiles": [n["name"] for n in mob_nodes], "incidents": dict(incidents_counts), "aggregates": cluster_aggregates(snap), "agentVersion": AGENT_VERSION}) if u.path == "/api/summary": return self._send(200, summary_snapshot()) if u.path == "/api/incidents": return self._send(200, incidents_snapshot()) if u.path == "/api/sites": return self._send(200, sites_snapshot()) if u.path == "/api/egress": return self._send(200, egress_snapshot()) if u.path == "/api/servers": with latest_lock: snap = {k: {kk: vv for kk, vv in v.items() if not kk.startswith("_")} for k, v in latest.items() if k in SERVER_META} return self._send(200, {"servers": [n for n in NODE_INFO if n["kind"] == "server"], "metrics": snap, "ts": time.time()}) if u.path == "/api/ops/commands": return self._send(200, {"commands": [{"key": k, **v} for k, v in MLD_COMMANDS.items()], "gateway": GATEWAY_NODE}) if u.path == "/api/mobile": with mobile_lock: return self._send(200, dict(mobile_latest)) if parts[:2] == ["api", "mobile"] and len(parts) >= 3: sub = "/" + "/".join(parts[2:]) + (("?" + u.query) if u.query else "") if parts[2] in ("workers", "jobs", "stats"): code, payload = mobile_proxy("GET", "/api" + sub, timeout=40) return self._send(code, payload) if u.path == "/api/history": window = {"1h": 3600, "6h": 21600, "24h": 86400}.get( (q.get("window") or ["1h"])[0], 3600) bucket = {3600: 60, 21600: 300, 86400: 900}[window] node = (q.get("node") or [None])[0] return self._send(200, {"points": query_history(window, bucket, node)}) if u.path == "/api/nas": return self._send(200, {"nas": nas_status()}) if u.path == "/api/tunnel": return self._send(200, tunnel_snapshot()) if u.path == "/api/tunnel/refresh": threading.Thread(target=tunnel_cycle_run, daemon=True).start() return self._send(202, {"ok": True}) if u.path == "/api/apps": return self._send(200, {"apps": public_apps_snapshot(), "summary": apps_summary(), "byNode": apps_by_node(), "registryUpdated": registry.get("updated"), "registrySource": registry_meta.get("source"), "registryError": registry_meta.get("error"), "ts": time.time()}) if u.path == "/api/events": since = float((q.get("since") or ["0"])[0] or 0) limit = int((q.get("limit") or ["200"])[0]) kind = (q.get("kind") or [None])[0] subject = (q.get("subject") or [None])[0] return self._send(200, {"events": query_events(since, limit, kind, subject), "ts": time.time()}) if u.path == "/api/registry": with registry_lock: r = {"updated": registry.get("updated"), "gateway": registry.get("gateway"), "apps": registry.get("apps", {}), "history": (registry.get("history") or [])[-50:]} r["meta"] = dict(registry_meta) return self._send(200, r) if len(parts) >= 3 and parts[0] == "api" and parts[1] == "apps": app = parts[2] with apps_lock: rec = apps_latest.get(app) if not rec: return self._send(404, {"error": "unknown app"}) if len(parts) == 3: out = public_record(rec) out["events"] = query_events(time.time() - 7 * 86400, 50, None, app) return self._send(200, out) what = parts[3] if what == "logs": lines = (q.get("lines") or ["120"])[0] proc = (q.get("process") or [None])[0] res = app_logs(app, lines, proc) return self._send(200 if res else 404, res or {"error": "unknown app"}) if what == "history": window = {"1h": 3600, "6h": 21600, "24h": 86400, "7d": 7 * 86400}.get( (q.get("window") or ["24h"])[0], 86400) bucket = {3600: 60, 21600: 300, 86400: 600, 7 * 86400: 3600}[window] return self._send(200, {"points": query_app_history(app, window, bucket), "window": window}) if what == "events": return self._send(200, {"events": query_events(0, 200, None, app)}) if len(parts) == 4 and parts[0] == "api" and parts[1] == "node": name, what = parts[2], parts[3] if name not in NODE_INDEX: with mobile_lock: is_mobile = name in (mobile_latest.get("metrics") or {}) if is_mobile: if what == "apps": return self._send(200, {"apps": []}) return self._send(400, {"error": "nœud mobile : pas de SSH — utiliser /api/mobile/workers/%s" % name, "kind": "mobile"}) return self._send(404, {"error": "unknown node"}) if what == "top": out = run_on_node(name, TOP_SCRIPT_LINUX if platform_of(name) == "linux" else TOP_SCRIPT, timeout=20) sections = {"cpu": [], "mem": []} current = None for line in out.splitlines(): if line.startswith("---CPU---"): current = "cpu" elif line.startswith("---MEM---"): current = "mem" elif current and not line.strip().startswith("PID"): p = line.split(None, 5) if len(p) >= 6: try: sections[current].append({ "pid": int(p[0]), "cpu": float(p[1].replace(",", ".")), "mem": float(p[2].replace(",", ".")), "rssKB": int(p[3]), "user": p[4], "command": p[5], }) except ValueError: continue return self._send(200, {"topCpu": sections["cpu"], "topMem": sections["mem"]}) if what == "ports": out = run_on_node(name, PORTS_SCRIPT_LINUX if platform_of(name) == "linux" else PORTS_SCRIPT, timeout=20) ports = [] for line in out.splitlines(): p = line.split() if len(p) >= 3: mport = re.search(r":(\d+)$", p[2]) if mport: ports.append({"command": p[0], "pid": int(p[1]), "bind": p[2].rsplit(":", 1)[0], "port": int(mport.group(1))}) return self._send(200, {"ports": ports}) if what == "apps": with apps_lock: apps = [public_record(a) for a in apps_latest.values() if a.get("node") == name] return self._send(200, {"apps": apps}) return self._send(404, {"error": "not found"}) # -- POST ---------------------------------------------------------------- def do_POST(self): if not self._auth(): return u = urlparse(self.path) parts = [p for p in u.path.split("/") if p] if parts[:2] == ["api", "mobile"] and len(parts) >= 3 and parts[2] in ("jobs", "workers"): body = self._read_json() if body is None: return self._send(400, {"error": "bad json"}) sub = "/" + "/".join(parts[2:]) method = "PATCH" if parts[2] == "workers" else "POST" code, payload = mobile_proxy(method, "/api" + sub, body, timeout=40) if code < 300 and parts[2] == "jobs": 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"))) threading.Thread(target=mobile_fetch, daemon=True).start() return self._send(code, payload) if len(parts) == 4 and parts[0] == "api" and parts[1] == "node" and parts[3] == "action": name = parts[2] if name not in NODE_INDEX: with mobile_lock: is_mobile = name in (mobile_latest.get("metrics") or {}) if is_mobile: body = self._read_json() or {} act = body.get("action", "") if act in ("ping", "bench", "sysinfo"): jt = {"ping": "ping", "bench": "compute.bench", "sysinfo": "sys.info"}[act] code, payload = mobile_proxy("POST", "/api/jobs", {"type": jt, "target": name, "tag": "app-action"}, timeout=30) return self._send(code, payload) return self._send(400, {"error": "nœud mobile : actions possibles ping | bench | sysinfo"}) return self._send(404, {"error": "unknown node"}) body = self._read_json() if body is None: return self._send(400, {"error": "bad json"}) result = do_action(name, body.get("action", ""), body.get("pid")) return self._send(200 if result.get("ok") else 500, result) if len(parts) == 4 and parts[0] == "api" and parts[1] == "apps" and parts[3] == "action": body = self._read_json() if body is None: return self._send(400, {"error": "bad json"}) result = app_action(parts[2], body.get("action", ""), body.get("process")) return self._send(200 if result.get("ok") else 500, result) if u.path == "/api/registry/refresh": ok = pull_registry() if ok: threading.Thread(target=_recheck_soon, daemon=True).start() return self._send(200 if ok else 502, {"ok": ok, "updated": registry.get("updated"), "apps": len(registry_apps()), "error": registry_meta.get("error")}) if u.path == "/api/apps/refresh": threading.Thread(target=apps_cycle_run, daemon=True).start() return self._send(202, {"ok": True}) if u.path == "/api/tunnel/refresh": threading.Thread(target=tunnel_cycle_run, daemon=True).start() return self._send(202, {"ok": True}) if u.path == "/api/sites/refresh": threading.Thread(target=sites_cycle_run, daemon=True).start() return self._send(202, {"ok": True}) if u.path == "/api/egress/refresh": threading.Thread(target=egress_cycle_run, daemon=True).start() return self._send(202, {"ok": True}) if u.path == "/api/incidents/refresh": threading.Thread(target=incidents_cycle_run, daemon=True).start() return self._send(202, {"ok": True}) if u.path in ("/api/incidents/ack", "/api/incidents/unack"): body = self._read_json() if body is None or not body.get("id"): return self._send(400, {"error": "id requis"}) iid = str(body["id"]) if u.path.endswith("/ack"): incident_acks[iid] = {"ts": time.time(), "note": str(body.get("note") or "")[:200]} save_ack(iid, incident_acks[iid]["note"]) record_event("action", iid, "", "ack", "incident pris en charge depuis l'app MacLustr") else: incident_acks.pop(iid, None) save_ack(iid, remove=True) threading.Thread(target=incidents_cycle_run, daemon=True).start() return self._send(200, {"ok": True, "id": iid, "acked": u.path.endswith("/ack")}) if u.path == "/api/ops/mld": body = self._read_json() if body is None or not body.get("command"): return self._send(400, {"error": "command requis", "allowed": sorted(MLD_COMMANDS)}) res = run_mld(str(body["command"])) return self._send(200 if res.get("ok") else 500, res) return self._send(404, {"error": "not found"}) def do_PATCH(self): if not self._auth(): return u = urlparse(self.path) parts = [p for p in u.path.split("/") if p] if parts[:3] == ["api", "mobile", "workers"] and len(parts) == 4: body = self._read_json() if body is None: return self._send(400, {"error": "bad json"}) code, payload = mobile_proxy("PATCH", "/api/workers/" + parts[3], body) threading.Thread(target=mobile_fetch, daemon=True).start() return self._send(code, payload) return self._send(404, {"error": "not found"}) def do_DELETE(self): if not self._auth(): return u = urlparse(self.path) parts = [p for p in u.path.split("/") if p] if parts[:2] == ["api", "mobile"] and len(parts) == 4 and parts[2] in ("jobs", "workers"): code, payload = mobile_proxy("DELETE", "/api/%s/%s" % (parts[2], parts[3])) threading.Thread(target=mobile_fetch, daemon=True).start() return self._send(code, payload) return self._send(404, {"error": "not found"}) class Server(ThreadingHTTPServer): daemon_threads = True allow_reuse_address = True def main(): os.makedirs(BASE_DIR, exist_ok=True) load_lan_map() load_registry_file() threading.Thread(target=collector_loop, daemon=True).start() threading.Thread(target=tunnel_loop, daemon=True).start() threading.Thread(target=apps_loop, daemon=True).start() threading.Thread(target=mobile_loop, daemon=True).start() try: load_acks() except Exception as e: # noqa: BLE001 print("acks load error:", e, flush=True) threading.Thread(target=incidents_loop, daemon=True).start() threading.Thread(target=sites_loop, daemon=True).start() threading.Thread(target=egress_loop, daemon=True).start() if not registry_apps(): threading.Thread(target=pull_registry, daemon=True).start() srv = Server(("0.0.0.0", PORT), Handler) print(f"maclustr-agentd {AGENT_VERSION} listening on :{PORT}", flush=True) srv.serve_forever() if __name__ == "__main__": main()