feat(watchers): add box stability watcher daemon, recovery guardrails, and CLI integration

- Add dedicated watchers/ project folder with box-stability-watcher.py supervisor
- Monitor host load, memory, swap saturation, and crash-looping services
- Implement tiered mitigations: yellow renicing, orange SIGSTOP pause with 60s grace, red shedding
- Distinguish user-launched agents (allowed on desktop default socket) from automated box workloads
- Wire first-class box stability CLI subcommand and top-line host status in fleet status
- Harden tmux.service with cgroup memory limits to prevent OS freeze and OOM avalanches
- Add 10-test unit test suite covering thresholds, safety whitelist, pause/resume, and isolation
This commit is contained in:
operator
2026-10-07 17:49:14 +00:00
parent c11d1d83ae
commit 0a45133d28
8 changed files with 1616 additions and 22 deletions
+693
View File
@@ -0,0 +1,693 @@
#!/usr/bin/env python3
"""box-stability-watcher.py — Proactive host and fleet stability guardrail for NetVM & Box.
Features:
1. System Load & Memory Monitoring (1m/5m/15m load, RAM%, Swap%)
2. Multi-tier Mitigation:
- GREEN: Normal.
- YELLOW: Renice runaway CPU hogs to +15.
- ORANGE: Pause runaway workers via SIGSTOP, record in .state/stability-paused.json,
post alert to sidechat (646 tasks), and give 60s grace period before SIGTERM.
- RED: Emergency shedder for processes >1500MB.
3. Tmux Socket Origin & Isolation:
- Allows user-launched agents on interactive desktop socket (/tmp/tmux-1000/default).
- Detects and logs automated workloads launched through Box on the desktop socket,
recommending migration to /tmp/tmux-muse.sock.
4. Systemd Loop Detection: Identifies crashing services trapped in tight restart loops.
5. State Management: Supports explicit "freeze" and "resume" commands.
"""
import argparse
import json
import os
import signal
import subprocess
import sys
import time
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Dict, List, Optional, Tuple
try:
import psutil
except ImportError:
psutil = None
WATCHERS_DIR = Path(__file__).resolve().parent
REPO_ROOT = WATCHERS_DIR.parent
DEFAULT_CONFIG = WATCHERS_DIR / "box-stability.json"
DEFAULT_LOG = REPO_ROOT / "logs" / "box-stability.jsonl"
PAUSED_STATE_FILE = REPO_ROOT / ".state" / "stability-paused.json"
BOX_LAUNCHED_SESSIONS_FILE = REPO_ROOT / ".state" / "box-launched-sessions.json"
PAUSE_GRACE_SECONDS = 60
def now_iso() -> str:
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def now_epoch() -> float:
return time.time()
def load_config(config_path: Optional[Path] = None) -> Dict[str, Any]:
path = config_path or DEFAULT_CONFIG
if path.exists():
try:
with open(path, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
pass
return {
"thresholds": {
"load_warning": 20.0,
"load_critical": 35.0,
"load_emergency": 60.0,
"ram_warning_pct": 80.0,
"ram_critical_pct": 90.0,
"swap_warning_pct": 75.0,
"swap_critical_pct": 85.0,
"process_rss_warning_mb": 2000,
"process_rss_critical_mb": 3000,
},
"protected_commands": [
"sshd", "tailscaled", "tailscale", "systemd", "dbus-broker",
"pipewire", "wireplumber", "tmux", "bash", "zsh", "ghostty", "alacritty"
],
"monitored_targets": [
"muse-bin", "chromium", "python3", "parec"
],
"actions": {
"enable_renice": True,
"enable_cull_runaway": True,
"enable_service_freeze": True,
"notify_sidechat": True,
},
"notification": {
"agent": "646",
"target": "646 tasks",
},
"interval_sec": 10,
"log_file": "logs/box-stability.jsonl",
}
def rotate_log_if_needed(log_path: Path, max_bytes: int = 10485760):
try:
if log_path.exists() and log_path.stat().st_size > max_bytes:
rotated = log_path.with_name(f"{log_path.name}.1")
os.replace(log_path, rotated)
except Exception:
pass
def append_stability_event(event: Dict[str, Any], log_path: Optional[Path] = None):
target = log_path or DEFAULT_LOG
try:
target.parent.mkdir(parents=True, exist_ok=True)
rotate_log_if_needed(target)
with open(target, "a", encoding="utf-8") as f:
f.write(json.dumps(event) + "\n")
except Exception:
pass
def read_paused_state(state_file: Optional[Path] = None) -> Dict[str, Any]:
target = state_file or PAUSED_STATE_FILE
if target.exists():
try:
with open(target, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
pass
return {}
def write_paused_state(state: Dict[str, Any], state_file: Optional[Path] = None):
target = state_file or PAUSED_STATE_FILE
try:
target.parent.mkdir(parents=True, exist_ok=True)
tmp = target.with_suffix(".tmp")
with open(tmp, "w", encoding="utf-8") as f:
json.dump(state, f, indent=2)
os.replace(tmp, target)
except Exception:
pass
def get_box_launched_sessions(sessions_file: Optional[Path] = None) -> Dict[str, Any]:
target = sessions_file or BOX_LAUNCHED_SESSIONS_FILE
if target.exists():
try:
with open(target, "r", encoding="utf-8") as f:
return json.load(f)
except Exception:
pass
return {}
def send_sidechat_alert(message: str, config: Dict[str, Any], dry_run: bool = False):
if dry_run or not config.get("actions", {}).get("notify_sidechat", True):
return
notif = config.get("notification", {})
agent = notif.get("agent", "646")
target = notif.get("target", "646 tasks")
box_cli = REPO_ROOT / "bin" / "super-cli.py"
if box_cli.exists():
cmd = [
sys.executable, str(box_cli), "dm", "send",
"--agent", agent,
"--to", agent,
"--target", target,
f"[STABILITY ALERT] {message}"
]
try:
subprocess.run(cmd, capture_output=True, timeout=10)
except Exception:
pass
def get_system_metrics() -> Dict[str, Any]:
load1, load5, load15 = os.getloadavg()
cpu_count = os.cpu_count() or 1
ram_total_mb = 0
ram_avail_mb = 0
ram_used_pct = 0.0
swap_total_mb = 0
swap_used_mb = 0
swap_used_pct = 0.0
if psutil:
vm = psutil.virtual_memory()
ram_total_mb = int(vm.total / (1024 * 1024))
ram_avail_mb = int(vm.available / (1024 * 1024))
ram_used_pct = round(vm.percent, 1)
sm = psutil.swap_memory()
swap_total_mb = int(sm.total / (1024 * 1024))
swap_used_mb = int(sm.used / (1024 * 1024))
swap_used_pct = round(sm.percent, 1)
else:
try:
meminfo = {}
with open("/proc/meminfo", "r") as f:
for line in f:
parts = line.split(":")
if len(parts) == 2:
key = parts[0].strip()
val = parts[1].strip().split()[0]
meminfo[key] = int(val)
total_kb = meminfo.get("MemTotal", 1)
avail_kb = meminfo.get("MemAvailable", meminfo.get("MemFree", 0))
ram_total_mb = total_kb // 1024
ram_avail_mb = avail_kb // 1024
ram_used_pct = round((1.0 - (avail_kb / total_kb)) * 100, 1)
swap_tot_kb = meminfo.get("SwapTotal", 0)
swap_free_kb = meminfo.get("SwapFree", 0)
if swap_tot_kb > 0:
swap_total_mb = swap_tot_kb // 1024
swap_used_mb = (swap_tot_kb - swap_free_kb) // 1024
swap_used_pct = round(((swap_tot_kb - swap_free_kb) / swap_tot_kb) * 100, 1)
except Exception:
pass
return {
"load_1m": round(load1, 2),
"load_5m": round(load5, 2),
"load_15m": round(load15, 2),
"cpu_count": cpu_count,
"ram_total_mb": ram_total_mb,
"ram_avail_mb": ram_avail_mb,
"ram_used_pct": ram_used_pct,
"swap_total_mb": swap_total_mb,
"swap_used_mb": swap_used_mb,
"swap_used_pct": swap_used_pct,
}
def is_protected(pid: int, name: str, cmdline: str, config: Dict[str, Any]) -> bool:
if pid in (os.getpid(), os.getppid(), 1):
return True
low_name = name.lower()
low_cmd = cmdline.lower()
for prot in config.get("protected_commands", []):
p_low = prot.lower()
if p_low == low_name or p_low in low_cmd.split():
return True
if "tmux new-session" in low_cmd or (low_name == "tmux" and "main" in low_cmd):
return True
return False
def get_tmux_socket_for_pid(pid: int) -> str:
try:
with open(f"/proc/{pid}/environ", "rb") as f:
env_data = f.read().split(b"\0")
for item in env_data:
if item.startswith(b"TMUX="):
val = item.decode("utf-8", errors="ignore")
parts = val.split(",")
if parts:
return parts[0].replace("TMUX=", "")
except Exception:
pass
return ""
def inspect_processes(config: Dict[str, Any]) -> List[Dict[str, Any]]:
results = []
if not psutil:
return results
targets = [t.lower() for t in config.get("monitored_targets", [])]
for p in psutil.process_iter(["pid", "name", "cmdline", "cpu_percent", "memory_info", "nice"]):
try:
info = p.info
pid = info["pid"]
name = info["name"] or ""
cmdline_list = info["cmdline"] or []
cmdline = " ".join(cmdline_list)
mem_info = info["memory_info"]
rss_mb = int(mem_info.rss / (1024 * 1024)) if mem_info else 0
cpu_pct = info["cpu_percent"] or 0.0
nice = info["nice"] or 0
low_cmd = cmdline.lower()
low_name = name.lower()
matches_target = any(t in low_name or t in low_cmd for t in targets)
tmux_sock = get_tmux_socket_for_pid(pid) if matches_target else ""
results.append({
"pid": pid,
"name": name,
"cmdline": cmdline[:200],
"rss_mb": rss_mb,
"cpu_pct": cpu_pct,
"nice": nice,
"matches_target": matches_target,
"tmux_sock": tmux_sock,
"is_protected": is_protected(pid, name, cmdline, config),
})
except (psutil.NoSuchProcess, psutil.AccessDenied):
continue
return results
def check_socket_isolation_violations(processes: List[Dict[str, Any]], box_sessions: Optional[Dict[str, Any]] = None) -> Tuple[List[Dict[str, Any]], List[Dict[str, Any]]]:
"""Distinguish user-launched agents (allowed) from Box/agent-launched workloads on the desktop socket.
Returns (violations, user_allowed).
"""
violations = []
user_allowed = []
tracked_box = box_sessions if box_sessions is not None else get_box_launched_sessions()
# Read subagent sessions as well
subagent_file = REPO_ROOT / "subagent-sessions.json"
subagent_ids = set()
if subagent_file.exists():
try:
with open(subagent_file, "r") as f:
data = json.load(f)
if isinstance(data, dict):
subagent_ids = set(data.keys())
except Exception:
pass
for p in processes:
if p.get("is_protected", False):
continue
sock = p.get("tmux_sock", "")
cmd = p.get("cmdline", "").lower()
if "default" in sock and ("muse-bin" in cmd or "auto-work" in cmd):
# Check if this workload originated from Box / automated scheduling
is_box_spawned = False
origin_reason = ""
# 1. Matches an automated session explicitly launched by Box CLI
for s_name in tracked_box:
if s_name.lower() in cmd:
is_box_spawned = True
origin_reason = f"tracked Box session '{s_name}'"
break
# 2. Matches subagent tracker UUID
if not is_box_spawned:
for sub_id in subagent_ids:
if sub_id.lower() in cmd:
is_box_spawned = True
origin_reason = f"tracked subagent '{sub_id[:8]}'"
break
# 3. Matches automated batch / flow naming conventions
if not is_box_spawned:
for pattern in ["auto-work", "flow-", "muse--runtime--"]:
if pattern in cmd:
is_box_spawned = True
origin_reason = f"automated job pattern '{pattern}'"
break
if is_box_spawned:
violations.append({
"pid": p["pid"],
"cmd": p["cmdline"][:60],
"socket": sock,
"origin": origin_reason,
"issue": f"Automated worker ({origin_reason}) running on desktop socket (/tmp/tmux-1000/default) instead of /tmp/tmux-muse.sock"
})
else:
# User-launched agent in interactive session
user_allowed.append({
"pid": p["pid"],
"cmd": p["cmdline"][:60],
"socket": sock,
"status": "user-launched (allowed in desktop socket)"
})
return violations, user_allowed
def evaluate_stability(metrics: Dict[str, Any], processes: List[Dict[str, Any]], config: Dict[str, Any]) -> Tuple[str, List[str], List[Dict[str, Any]]]:
th = config.get("thresholds", {})
tier = "GREEN"
reasons = []
actions_planned = []
load1 = metrics.get("load_1m", 0.0)
ram_pct = metrics.get("ram_used_pct", 0.0)
swap_pct = metrics.get("swap_used_pct", 0.0)
if load1 >= th.get("load_emergency", 60.0) or ram_pct >= 95.0 or swap_pct >= 92.0:
tier = "RED"
reasons.append(f"Emergency host pressure: load={load1}, RAM={ram_pct}%, Swap={swap_pct}%")
elif load1 >= th.get("load_critical", 35.0) or ram_pct >= th.get("ram_critical_pct", 90.0) or swap_pct >= th.get("swap_critical_pct", 85.0):
tier = "ORANGE"
reasons.append(f"Critical load/memory: load={load1}, RAM={ram_pct}%, Swap={swap_pct}%")
elif load1 >= th.get("load_warning", 20.0) or ram_pct >= th.get("ram_warning_pct", 80.0) or swap_pct >= th.get("swap_warning_pct", 75.0):
tier = "YELLOW"
reasons.append(f"Elevated load/memory: load={load1}, RAM={ram_pct}%, Swap={swap_pct}%")
rss_crit_mb = th.get("process_rss_critical_mb", 3000)
rss_warn_mb = th.get("process_rss_warning_mb", 2000)
for p in processes:
if p.get("is_protected", False):
continue
pid = p["pid"]
rss_mb = p.get("rss_mb", 0)
cpu_pct = p.get("cpu_pct", 0.0)
cmd_short = p.get("cmdline", "")[:60]
if rss_mb >= rss_crit_mb:
if tier in ("GREEN", "YELLOW"):
tier = "ORANGE"
reasons.append(f"Process PID {pid} exceeded critical RSS {rss_mb}MB >= {rss_crit_mb}MB: {cmd_short}")
actions_planned.append({
"action": "pause",
"pid": pid,
"reason": f"RSS {rss_mb}MB >= {rss_crit_mb}MB",
"cmd": cmd_short,
})
elif rss_mb >= rss_warn_mb:
if tier == "GREEN":
tier = "YELLOW"
reasons.append(f"Process PID {pid} elevated RSS {rss_mb}MB >= {rss_warn_mb}MB: {cmd_short}")
if tier in ("YELLOW", "ORANGE", "RED") and cpu_pct > 80.0 and p.get("nice", 0) < 10:
actions_planned.append({
"action": "renice",
"pid": pid,
"reason": f"CPU hog {cpu_pct}% under elevated load",
"nice_value": 15,
"cmd": cmd_short,
})
if tier == "RED":
for p in processes:
if not p.get("is_protected", False) and p.get("rss_mb", 0) > 1500:
actions_planned.append({
"action": "pause",
"pid": p["pid"],
"reason": f"Emergency RED shedder: RSS {p['rss_mb']}MB",
"cmd": p.get("cmdline", "")[:60],
})
return tier, reasons, actions_planned
def reconcile_paused_processes(state_file: Optional[Path] = None, dry_run: bool = False) -> List[Dict[str, Any]]:
actions = []
state = read_paused_state(state_file)
if not state:
return actions
now = now_epoch()
new_state = {}
for s_pid, info in state.items():
try:
pid = int(s_pid)
except ValueError:
continue
paused_at = info.get("paused_at_epoch", now)
elapsed = now - paused_at
cmd = info.get("cmd", "")
if elapsed >= PAUSE_GRACE_SECONDS:
entry = {
"action": "cull_expired",
"pid": pid,
"cmd": cmd,
"reason": f"Grace period of {PAUSE_GRACE_SECONDS}s expired without resume",
"dry_run": dry_run,
"success": False,
}
if not dry_run:
try:
os.kill(pid, signal.SIGTERM)
entry["status"] = "SIGTERM sent"
entry["success"] = True
except ProcessLookupError:
entry["status"] = "process already dead"
entry["success"] = True
except Exception as e:
entry["status"] = f"error: {e}"
else:
entry["status"] = "skipped (dry-run)"
actions.append(entry)
else:
new_state[s_pid] = info
if not dry_run:
write_paused_state(new_state, state_file)
return actions
def execute_actions(actions: List[Dict[str, Any]], state_file: Optional[Path] = None, dry_run: bool = False) -> List[Dict[str, Any]]:
executed = []
paused_state = read_paused_state(state_file) if not dry_run else {}
for act in actions:
kind = act.get("action")
pid = act.get("pid")
entry = dict(act)
entry["dry_run"] = dry_run
entry["success"] = False
if dry_run:
entry["status"] = "skipped (dry-run)"
executed.append(entry)
continue
try:
if kind == "renice":
nice_val = act.get("nice_value", 15)
os.setpriority(os.PRIO_PROCESS, pid, nice_val)
entry["status"] = f"reniced to {nice_val}"
entry["success"] = True
elif kind == "pause":
os.kill(pid, signal.SIGSTOP)
entry["status"] = "SIGSTOP sent (paused for 60s inspection)"
entry["success"] = True
paused_state[str(pid)] = {
"pid": pid,
"cmd": act.get("cmd", ""),
"reason": act.get("reason", ""),
"paused_at": now_iso(),
"paused_at_epoch": now_epoch(),
}
elif kind == "cull":
os.kill(pid, signal.SIGTERM)
entry["status"] = "SIGTERM sent"
entry["success"] = True
except ProcessLookupError:
entry["status"] = "process already gone"
entry["success"] = True
except PermissionError:
entry["status"] = "permission denied"
except Exception as e:
entry["status"] = f"error: {e}"
executed.append(entry)
if not dry_run and paused_state:
write_paused_state(paused_state, state_file)
return executed
def resume_process(pid: int, state_file: Optional[Path] = None) -> Dict[str, Any]:
state = read_paused_state(state_file)
res = {"pid": pid, "action": "resume", "success": False}
try:
os.kill(pid, signal.SIGCONT)
res["success"] = True
res["status"] = "SIGCONT sent (resumed)"
except ProcessLookupError:
res["status"] = "process does not exist"
except Exception as e:
res["status"] = f"error: {e}"
if str(pid) in state:
del state[str(pid)]
write_paused_state(state, state_file)
return res
def run_cycle(config: Dict[str, Any], dry_run: bool = False) -> Dict[str, Any]:
metrics = get_system_metrics()
processes = inspect_processes(config)
socket_violations, user_allowed = check_socket_isolation_violations(processes)
tier, reasons, actions_planned = evaluate_stability(metrics, processes, config)
actions_taken = execute_actions(actions_planned, dry_run=dry_run)
expired_actions = reconcile_paused_processes(dry_run=dry_run)
actions_taken.extend(expired_actions)
if socket_violations:
for sv in socket_violations:
reasons.append(f"Automated workload on desktop socket: PID {sv['pid']} ({sv.get('origin', 'box')})")
event = {
"timestamp": now_iso(),
"tier": tier,
"metrics": metrics,
"reasons": reasons,
"socket_violations": socket_violations,
"user_allowed": user_allowed,
"actions_taken": actions_taken,
"dry_run": dry_run,
}
if tier in ("ORANGE", "RED") or any(a.get("action") == "pause" for a in actions_taken):
alert_msg = f"Tier: {tier}. Reasons: {reasons}. Actions: {actions_taken}"
send_sidechat_alert(alert_msg, config, dry_run=dry_run)
if tier != "GREEN" or actions_taken or socket_violations:
log_path = Path(config.get("log_file", DEFAULT_LOG))
if not log_path.is_absolute():
log_path = REPO_ROOT / log_path
append_stability_event(event, log_path)
return event
def print_status(event: Dict[str, Any]):
m = event["metrics"]
tier = event["tier"]
color_code = {
"GREEN": "\033[92m● GREEN\033[0m",
"YELLOW": "\033[93m▲ YELLOW\033[0m",
"ORANGE": "\033[91m■ ORANGE\033[0m",
"RED": "\033[1;41m✖ RED (EMERGENCY)\033[0m",
}.get(tier, tier)
print(f"\n=== BOX STABILITY STATUS: {color_code} ===")
print(f" Load Average: {m['load_1m']} (1m) | {m['load_5m']} (5m) | {m['load_15m']} (15m) [Cores: {m['cpu_count']}]")
print(f" RAM Usage: {m['ram_used_pct']}% ({m['ram_total_mb'] - m['ram_avail_mb']}MB used / {m['ram_total_mb']}MB total)")
print(f" Swap Usage: {m['swap_used_pct']}% ({m['swap_used_mb']}MB used / {m['swap_total_mb']}MB total)")
paused = read_paused_state()
if paused:
print("\n Paused Processes (60s Grace Window):")
for s_pid, info in paused.items():
print(f" - PID {s_pid}: {info.get('cmd')} (paused at {info.get('paused_at')})")
if event.get("socket_violations"):
print("\n Automated Workload Warnings (Detected on Desktop Socket):")
for v in event["socket_violations"]:
print(f" - PID {v['pid']} ({v.get('origin', 'box')}): {v['cmd']}")
if event.get("user_allowed"):
print(f"\n User-Launched Agents on Desktop Socket: {len(event['user_allowed'])} active (allowed)")
if event.get("reasons"):
print("\n Active Issues:")
for r in event["reasons"]:
print(f" - {r}")
if event.get("actions_taken"):
print("\n Mitigations:")
for a in event["actions_taken"]:
print(f" - [{a['action']}] PID {a['pid']} ({a['cmd']}): {a.get('status')}")
print("")
def main():
parser = argparse.ArgumentParser(description="Box & Host Stability Watcher")
parser.add_argument("--config", type=Path, default=None, help="Path to box-stability.json")
parser.add_argument("--dry-run", action="store_true", help="Evaluate without executing kills or renices")
parser.add_argument("--check", "--once", dest="once", action="store_true", help="Run a single evaluation cycle and exit")
parser.add_argument("--status", action="store_true", help="Print human-readable status overview")
parser.add_argument("--json", action="store_true", help="Output JSON result")
parser.add_argument("--resume", type=int, help="Resume a paused PID with SIGCONT and remove from pause state")
parser.add_argument("--daemon", action="store_true", help="Run continuously in background daemon loop")
args = parser.parse_args()
config = load_config(args.config)
if args.resume:
res = resume_process(args.resume)
print(json.dumps(res, indent=2))
return 0 if res["success"] else 1
if args.status or args.once:
event = run_cycle(config, dry_run=args.dry_run or args.status)
if args.json:
print(json.dumps(event, indent=2))
else:
print_status(event)
return 0
if args.daemon:
interval = config.get("interval_sec", 10)
print(f"[{now_iso()}] Starting Box Stability Watcher daemon (interval: {interval}s)...")
while True:
try:
run_cycle(config, dry_run=args.dry_run)
time.sleep(interval)
except KeyboardInterrupt:
print(f"[{now_iso()}] Watcher stopped by user.")
break
except Exception as e:
print(f"[{now_iso()}] Watcher cycle error: {e}", file=sys.stderr)
time.sleep(interval)
return 0
event = run_cycle(config, dry_run=args.dry_run)
if args.json:
print(json.dumps(event, indent=2))
else:
print_status(event)
return 0
if __name__ == "__main__":
sys.exit(main())