diff --git a/tests/test_box_stability_watcher.py b/tests/test_box_stability_watcher.py new file mode 100644 index 0000000..02c5696 --- /dev/null +++ b/tests/test_box_stability_watcher.py @@ -0,0 +1,169 @@ +#!/usr/bin/env python3 +"""test_box_stability_watcher.py — Comprehensive unit tests for Box Stability Watcher.""" + +import json +import os +import sys +import tempfile +import unittest +from pathlib import Path +from unittest import mock + +REPO_ROOT = Path("/home/super/Projects/NetVM") +WATCHERS_DIR = REPO_ROOT / "watchers" +sys.path.insert(0, str(WATCHERS_DIR)) + +import importlib.util +spec = importlib.util.spec_from_file_location("box_stability_watcher", str(WATCHERS_DIR / "box-stability-watcher.py")) +w = importlib.util.module_from_spec(spec) +spec.loader.exec_module(w) + + +class TestConfigAndSafety(unittest.TestCase): + def test_load_config_defaults(self): + with tempfile.TemporaryDirectory() as td: + non_existent = Path(td) / "missing.json" + cfg = w.load_config(non_existent) + self.assertIn("thresholds", cfg) + self.assertIn("protected_commands", cfg) + self.assertEqual(cfg["thresholds"]["load_warning"], 20.0) + + def test_is_protected(self): + cfg = {"protected_commands": ["sshd", "tailscaled", "tmux", "systemd", "ghostty"]} + self.assertTrue(w.is_protected(1, "systemd", "/sbin/init", cfg)) + self.assertTrue(w.is_protected(os.getpid(), "python3", "some_script", cfg)) + self.assertTrue(w.is_protected(999, "sshd", "/usr/sbin/sshd -D", cfg)) + self.assertTrue(w.is_protected(888, "tmux", "tmux new-session -s main", cfg)) + self.assertTrue(w.is_protected(777, "tailscaled", "/usr/sbin/tailscaled", cfg)) + self.assertFalse(w.is_protected(1234, "muse-bin", "/home/super/.local/bin/muse-bin-1.4.3 resume abc", cfg)) + self.assertFalse(w.is_protected(5678, "chromium", "/usr/lib/chromium/chromium --type=renderer", cfg)) + + +class TestStabilityEvaluation(unittest.TestCase): + def setUp(self): + self.cfg = { + "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", "tmux"] + } + + def test_green_tier(self): + metrics = {"load_1m": 2.5, "ram_used_pct": 30.0, "swap_used_pct": 10.0} + procs = [ + {"pid": 101, "cmdline": "muse-bin", "rss_mb": 400, "cpu_pct": 5.0, "is_protected": False, "nice": 0} + ] + tier, reasons, actions = w.evaluate_stability(metrics, procs, self.cfg) + self.assertEqual(tier, "GREEN") + self.assertEqual(reasons, []) + self.assertEqual(actions, []) + + def test_yellow_tier_elevated_load(self): + metrics = {"load_1m": 22.5, "ram_used_pct": 50.0, "swap_used_pct": 20.0} + procs = [ + {"pid": 102, "cmdline": "python worker", "rss_mb": 500, "cpu_pct": 90.0, "is_protected": False, "nice": 0} + ] + tier, reasons, actions = w.evaluate_stability(metrics, procs, self.cfg) + self.assertEqual(tier, "YELLOW") + self.assertTrue(any("Elevated load/memory" in r for r in reasons)) + self.assertEqual(len(actions), 1) + self.assertEqual(actions[0]["action"], "renice") + self.assertEqual(actions[0]["pid"], 102) + + def test_orange_tier_pause_leaky_process(self): + metrics = {"load_1m": 2.0, "ram_used_pct": 40.0, "swap_used_pct": 10.0} + procs = [ + {"pid": 202, "cmdline": "muse-bin leaky", "rss_mb": 3400, "cpu_pct": 10.0, "is_protected": False, "nice": 0} + ] + tier, reasons, actions = w.evaluate_stability(metrics, procs, self.cfg) + self.assertEqual(tier, "ORANGE") + self.assertTrue(any("exceeded critical RSS" in r for r in reasons)) + self.assertEqual(len(actions), 1) + self.assertEqual(actions[0]["action"], "pause") + self.assertEqual(actions[0]["pid"], 202) + + def test_red_emergency_tier(self): + metrics = {"load_1m": 75.0, "ram_used_pct": 96.0, "swap_used_pct": 94.0} + procs = [ + {"pid": 301, "cmdline": "muse-bin heavy", "rss_mb": 1800, "cpu_pct": 50.0, "is_protected": False, "nice": 0}, + {"pid": 302, "cmdline": "sshd daemon", "rss_mb": 2000, "cpu_pct": 2.0, "is_protected": True, "nice": 0} + ] + tier, reasons, actions = w.evaluate_stability(metrics, procs, self.cfg) + self.assertEqual(tier, "RED") + self.assertTrue(any("Emergency host pressure" in r for r in reasons)) + pause_pids = [a["pid"] for a in actions if a["action"] == "pause"] + self.assertIn(301, pause_pids) + self.assertNotIn(302, pause_pids) + + +class TestSocketIsolation(unittest.TestCase): + def test_distinguishes_user_from_box_launched_agents(self): + # PID 555 is an automated job on default socket + # PID 666 is an agent on the correct fleet socket + # PID 777 is a user-launched interactive agent on default socket + procs = [ + {"pid": 555, "cmdline": "muse-bin auto-work sweep", "tmux_sock": "/tmp/tmux-1000/default", "is_protected": False}, + {"pid": 666, "cmdline": "muse-bin auto-work sweep", "tmux_sock": "/tmp/tmux-muse.sock", "is_protected": False}, + {"pid": 777, "cmdline": "muse-bin interactive chat", "tmux_sock": "/tmp/tmux-1000/default", "is_protected": False}, + ] + box_sessions = {"auto-work": {}} + violations, allowed = w.check_socket_isolation_violations(procs, box_sessions=box_sessions) + self.assertEqual(len(violations), 1) + self.assertEqual(violations[0]["pid"], 555) + self.assertEqual(len(allowed), 1) + self.assertEqual(allowed[0]["pid"], 777) + + +class TestPauseResumeAndExpiry(unittest.TestCase): + @mock.patch("os.kill") + def test_pause_and_state_persistence(self, mock_kill): + with tempfile.TemporaryDirectory() as td: + state_file = Path(td) / "paused.json" + actions = [{"action": "pause", "pid": 4321, "cmd": "muse-bin", "reason": "RSS high"}] + executed = w.execute_actions(actions, state_file=state_file, dry_run=False) + self.assertEqual(len(executed), 1) + mock_kill.assert_called_once_with(4321, 19) # SIGSTOP = 19 + self.assertTrue(state_file.exists()) + data = json.loads(state_file.read_text()) + self.assertIn("4321", data) + + @mock.patch("os.kill") + def test_reconcile_expired_pause(self, mock_kill): + with tempfile.TemporaryDirectory() as td: + state_file = Path(td) / "paused.json" + now = w.now_epoch() + state_data = { + "4321": {"pid": 4321, "cmd": "muse-bin", "paused_at_epoch": now - 70} + } + state_file.write_text(json.dumps(state_data)) + + expired = w.reconcile_paused_processes(state_file=state_file, dry_run=False) + self.assertEqual(len(expired), 1) + self.assertEqual(expired[0]["action"], "cull_expired") + self.assertEqual(expired[0]["pid"], 4321) + mock_kill.assert_called_once_with(4321, 15) # SIGTERM = 15 + remaining = json.loads(state_file.read_text()) + self.assertNotIn("4321", remaining) + + @mock.patch("os.kill") + def test_resume_process(self, mock_kill): + with tempfile.TemporaryDirectory() as td: + state_file = Path(td) / "paused.json" + state_file.write_text(json.dumps({"4321": {"pid": 4321}})) + res = w.resume_process(4321, state_file=state_file) + self.assertTrue(res["success"]) + mock_kill.assert_called_once_with(4321, 18) # SIGCONT = 18 + remaining = json.loads(state_file.read_text()) + self.assertNotIn("4321", remaining) + + +if __name__ == "__main__": + unittest.main() diff --git a/watchers/README.md b/watchers/README.md new file mode 100644 index 0000000..6105eca --- /dev/null +++ b/watchers/README.md @@ -0,0 +1,60 @@ +# NetVM Watchers + +Dedicated directory for background autonomous health, resource, and stability watchers on NetVM host `bl`. + +## Components + +- **`box-stability-watcher.py`**: Host resource supervisor and load shedder. Proactively monitors: + - System 1m/5m/15m load averages against core counts. + - Host RAM and Swap pressure percentages. + - Per-process memory leaks (critical RSS thresholds for `muse-bin`, headless Chromium renderers, Python workers). + - Rogue/leaked CPU hogs starving SSH/Tailscale. + - Failing systemd user services trapped in tight restart loops. + - Socket isolation violations (automated workers running on `/tmp/tmux-1000/default` instead of `/tmp/tmux-muse.sock`). +- **`box-stability.json`**: Tunable operational thresholds, notifications, and protected process whitelist. +- **`systemd/box-stability-watcher.service`**: Systemd user daemon running the watcher continuously with 10s evaluation ticks. + +## Operational Tiers & Mitigations + +| Tier | Status | Trigger Condition | Automated Action | +| :--- | :--- | :--- | :--- | +| **GREEN** | Normal | Load < 20, RAM < 80%, Swap < 75% | Silent monitoring. | +| **YELLOW** | Warning | Load >= 20, RAM >= 80%, or process RSS >= 2000MB | Renice CPU hogs (+15) to preserve interactive SSH responsiveness; log warning. | +| **ORANGE** | Critical | Load >= 35, RAM >= 90%, or process RSS >= 3000MB | **Pause (SIGSTOP)** runaway worker, record in `.state/stability-paused.json`, post alert to `646 tasks` sidechat, and allow 60s operator inspection before SIGTERM. | +| **RED** | Emergency | Load >= 60, RAM >= 95%, or Swap >= 92% | Emergency load shedding of non-protected heavy consumers (>1500MB). | + +## Paused Process Lifecycle (60s Grace Window) + +When a process is paused: +1. Sent `SIGSTOP` immediately. +2. Recorded in `.state/stability-paused.json` with timestamp and command info. +3. Alert posted to `646 tasks` sidechat. +4. An operator can inspect the runtime or resume it via: + ```bash + box stability resume + ``` +5. If unresumed after 60 seconds, the watcher automatically culls the process via `SIGTERM`. + +## Protected Whitelist + +The watcher will **never** terminate or renice: +`sshd`, `tailscaled`, `tailscale`, `systemd`, `dbus-broker`, `pipewire`, `wireplumber`, `tmux` (main server), `bash`, `zsh`, `ghostty`, `alacritty`. + +## Unified Box CLI Integration + +```bash +# Host stability status & active socket warnings +box stability status + +# Machine-readable JSON output +box stability json + +# Single evaluation check +box stability check [--dry-run] + +# Resume a paused process +box stability resume + +# Top-line host health indicator +box fleet status +``` diff --git a/watchers/__pycache__/box-stability-watcher.cpython-314.pyc b/watchers/__pycache__/box-stability-watcher.cpython-314.pyc new file mode 100644 index 0000000..1160d85 Binary files /dev/null and b/watchers/__pycache__/box-stability-watcher.cpython-314.pyc differ diff --git a/watchers/box-stability-watcher.py b/watchers/box-stability-watcher.py new file mode 100755 index 0000000..2b271ab --- /dev/null +++ b/watchers/box-stability-watcher.py @@ -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())