From 0c6d2235ab74c92df61c27707ea82c9dbff421fb Mon Sep 17 00:00:00 2001 From: operator Date: Tue, 6 Oct 2026 23:03:52 +0000 Subject: [PATCH] feat(kpi): add autonomous worker auto-spawn engine and watchdog reconciliation --- .agents/skills/box/SKILL.md | 2 +- bin/chromebox-watchdog.sh | 39 +++++++--- bin/kpi.py | 149 +++++++++++++++++++++++++++++++++++- bin/super-cli.py | 26 +++++++ tests/test_kpi.py | 41 ++++++++++ 5 files changed, 244 insertions(+), 13 deletions(-) diff --git a/.agents/skills/box/SKILL.md b/.agents/skills/box/SKILL.md index 073524f..5cccdf3 100644 --- a/.agents/skills/box/SKILL.md +++ b/.agents/skills/box/SKILL.md @@ -37,7 +37,7 @@ Add `--json` to any command for machine-readable output when parsing results in - `box tmux tally` / `box tmux auto [status|on|off|watch|once|logs|match]` — multi-socket Tmux worker tally, regex auto-approver daemon & guardrails. - `box onboard connects` / `box onboard-tui` — fleet & client onboarding inventory, CDP ports, OTP salvage & 4-surface TUI. - `box invite status|code |redeem ` / `box usage [--node N]` — invite codes and usage limits. -- `box kpi status|report |routes|spawn-worker` — fleet KPI tracking, spend/limit metrics, route health, runtime preservation advisories, and background worker spawning. +- `box kpi status|report |routes|spawn-worker|auto-spawn` — fleet KPI tracking, spend/limit metrics, route health, runtime preservation advisories, and background worker auto-spawning. - `box chromebox permissions list|get |set |describe ` — settings-menu toggles (readback-verified sets). ## Rules diff --git a/bin/chromebox-watchdog.sh b/bin/chromebox-watchdog.sh index 2c3e7e1..86ef602 100755 --- a/bin/chromebox-watchdog.sh +++ b/bin/chromebox-watchdog.sh @@ -15,10 +15,14 @@ export DBUS_SESSION_BUS_ADDRESS="${DBUS_SESSION_BUS_ADDRESS:-unix:path=${XDG_RUN # concurrent runs kill each others chrome (observed 2026-10-03: pip flapped # with simultaneous "relaunch OK" and "relaunch FAILED"). LOCK="/tmp/chromebox-watchdog-${1:-pip}.lock" -exec 9>"$LOCK" -if ! flock -n 9; then - echo "[$(date -u +%FT%TZ)] [$1] another watchdog run in progress, skipping" >&2 - exit 0 +# Tests source this file with CHROMEBOX_WATCHDOG_LIB_ONLY=1: they resolve +# ports and call helpers without running checks, so no lock is needed. +if [ "${CHROMEBOX_WATCHDOG_LIB_ONLY:-}" != "1" ]; then + exec 9>"$LOCK" + if ! flock -n 9; then + echo "[$(date -u +%FT%TZ)] [$1] another watchdog run in progress, skipping" >&2 + exit 0 + fi fi PROFILE="${1:-pip}" @@ -40,13 +44,15 @@ rotate_log() { rotate_log "$LOG" # CHROME_LOG rotation happens after PROFILE is set (see below) -case "$PROFILE" in - muse) CDP_PORT=9410 ;; - pip) CDP_PORT=9420 ;; - 646) CDP_PORT=9430 ;; - opm) CDP_PORT=9440 ;; - *) echo "unknown profile: $PROFILE" >&2; exit 1 ;; -esac +# Ports come from the fleet registry, not a hardcoded list: every active +# node (def/dev included) gets supervision automatically. The old 4-profile +# case left dev/def unsupervised — a dead Warp tunnel paged forever with +# no auto-recovery (2026-10-06 dev outage). +CDP_PORT="$("$NETVM_BIN/netvm-registry.py" "$PROFILE" 2>/dev/null)" || { + echo "unknown profile: $PROFILE" >&2 + exit 1 +} +[ -n "$CDP_PORT" ] || { echo "unknown profile: $PROFILE" >&2; exit 1; } rotate_log "$CHROME_LOG" log() { echo "[$(date -u +%FT%TZ)] [$PROFILE] $*" | tee -a "$LOG"; } @@ -107,7 +113,18 @@ print(h*3600 + mi*60 + se) esac } +# Allow sourcing for tests without running checks. +if [ "${CHROMEBOX_WATCHDOG_LIB_ONLY:-}" = "1" ]; then + return 0 2>/dev/null || exit 0 +fi + if healthy; then + # Auto-reconcile idle workers for healthy profiles + # DISABLED 2026-10-06 by operator-646: kpi auto-spawn ignores job schedule fields; + # find_pending_work_for_node returns the alphabetically-first definition every tick, + # re-spawning and re-noticing every ~2min (pip b01 BOX-AUTO-WORKER loop, 29+ copies). + # Watchdog health path untouched. Re-enable once the spawner is schedule-aware. + # python3 "$NETVM_BIN/super-cli.py" kpi auto-spawn --node "$PROFILE" >>"$LOG" 2>&1 || true exit 0 fi diff --git a/bin/kpi.py b/bin/kpi.py index 303327c..86ec224 100755 --- a/bin/kpi.py +++ b/bin/kpi.py @@ -407,6 +407,128 @@ def get_live_advisory_block(node: str) -> str: return "\n".join(lines) +NODE_SIDECHATS = { + "646": "646 tasks", + "opm": "heartbeat", + "pip": "646-pip-coord", + "dev": "dev-coord", + "def": "def-coord", + "muse": "646-muse-coord", +} + + +def find_pending_work_for_node(node: str) -> Optional[Dict[str, Any]]: + """Find assigned pending job or swarm slot for an agent node.""" + # 1. Look for node-specific auto-work jobs + if JOBS_DIR.exists(): + candidates = sorted(list(JOBS_DIR.glob(f"auto-work-{node}-*.json")) + list(JOBS_DIR.glob(f"{node}-*.json"))) + for c in candidates: + try: + with open(c, "r", encoding="utf-8") as f: + data = json.load(f) + job_agent = data.get("agent") + if job_agent and job_agent != node: + continue + job_name = c.stem + return { + "type": "job", + "name": job_name, + "path": str(c), + "cmd": f"{sys.executable} {BIN_DIR}/job-dispatch.py {job_name}", + } + except Exception: + continue + + # 2. Check pending swarm slots + try: + from swarm_worker.poller import find_pending_slots + slots = find_pending_slots() + if slots: + slot = slots[0] + sw_id = slot.get("swarm_id", "swarm") + idx = slot.get("slot_index", 0) + return { + "type": "swarm", + "name": f"swarm-{sw_id}-s{idx}", + "path": None, + "cmd": f"{sys.executable} {BIN_DIR}/swarm_worker/daemon.py", + } + except Exception: + pass + + return None + + +def auto_spawn_workers(nodes: Optional[List[str]] = None, dry_run: bool = False) -> List[Dict[str, Any]]: + """Reconcile idle agents and auto-spawn background tmux workers to execute pending work.""" + target_nodes = nodes or VALID_NODES + results = [] + + for node in target_nodes: + # Check active tmux workers for this node + active_tmux = len(get_agent_tmux_workers(node)) + if active_tmux > 0: + results.append({ + "node": node, + "action": "skip", + "reason": f"Active tmux worker already running ({active_tmux})", + }) + continue + + # Check work availability + work = find_pending_work_for_node(node) + if not work: + results.append({ + "node": node, + "action": "idle", + "reason": "No pending jobs or swarm slots", + }) + continue + + session_label = f"worker-{work['name'][:18]}" + cmd_to_run = f"{work['cmd']} > /tmp/tmux-{node}-{session_label}.log 2>&1" + + if dry_run: + results.append({ + "node": node, + "action": "would_spawn", + "session": session_label, + "work_type": work["type"], + "work_name": work["name"], + "command": work["cmd"], + }) + continue + + # Execute spawn + spawn_res = spawn_tmux_worker(node, session_label, cmd_to_run) + if spawn_res.get("ok"): + # Send sidechat notification + try: + from invite_handler import send_loopback_notice + chat = NODE_SIDECHATS.get(node, "646 tasks") + msg = f"[BOX-AUTO-WORKER] Spawned background tmux worker '{session_label}' executing {work['type']} ({work['name']}). Logs at /tmp/tmux-{node}-{session_label}.log" + send_loopback_notice(recipient=node, target=chat, message=msg) + except Exception: + pass + + results.append({ + "node": node, + "action": "spawned", + "session": session_label, + "work_type": work["type"], + "work_name": work["name"], + "command": work["cmd"], + }) + else: + results.append({ + "node": node, + "action": "error", + "error": spawn_res.get("error", "Unknown spawn error"), + }) + + return results + + def spawn_tmux_worker(node: str, session: str, command: str) -> Dict[str, Any]: """Spawn an autonomous tmux worker session on the agent's netns or shared socket.""" # Ensure session name is prefixed @@ -484,7 +606,10 @@ def main(): p_spawn.add_argument("node", choices=VALID_NODES, help="Agent node") p_spawn.add_argument("session", help="Session label") p_spawn.add_argument("worker_command", help="Command to execute inside worker") - p_spawn.add_argument("--json", action="store_true") + p_autospawn = subparsers.add_parser("auto-spawn", help="Auto-spawn background tmux workers for idle nodes with pending work") + p_autospawn.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node") + p_autospawn.add_argument("--dry-run", action="store_true", help="Report what would be spawned without executing") + p_autospawn.add_argument("--json", action="store_true") args = parser.parse_args() @@ -555,6 +680,28 @@ def main(): print(f"✘ Failed to spawn worker: {res.get('error')}", file=sys.stderr) sys.exit(1) + elif args.command == "auto-spawn": + nodes = [args.node] if getattr(args, "node", None) else None + results = auto_spawn_workers(nodes=nodes, dry_run=args.dry_run) + if args.json: + print(json.dumps(results, indent=2)) + return + print(f"\n=== AUTO-SPAWN WORKER RECONCILIATION {'(DRY-RUN)' if args.dry_run else ''} ===") + for r in results: + n = r.get("node") + act = r.get("action") + if act == "spawned": + print(f" ✔ @{n:<5} : SPAWNED session '{r.get('session')}' ({r.get('work_type')}: {r.get('work_name')})") + elif act == "would_spawn": + print(f" ? @{n:<5} : WOULD SPAWN session '{r.get('session')}' ({r.get('work_type')}: {r.get('work_name')})") + elif act == "skip": + print(f" - @{n:<5} : SKIP ({r.get('reason')})") + elif act == "idle": + print(f" - @{n:<5} : IDLE ({r.get('reason')})") + elif act == "error": + print(f" ✘ @{n:<5} : ERROR ({r.get('error')})") + print() + if __name__ == "__main__": main() diff --git a/bin/super-cli.py b/bin/super-cli.py index 7429764..5d42b55 100755 --- a/bin/super-cli.py +++ b/bin/super-cli.py @@ -1525,6 +1525,28 @@ def cmd_kpi(args): else: print(c_red(f"✘ Failed to spawn worker: {res.get('error')}"), file=sys.stderr) sys.exit(1) + elif act == "auto-spawn": + nodes = [args.node] if getattr(args, "node", None) else None + results = kpi.auto_spawn_workers(nodes=nodes, dry_run=getattr(args, "dry_run", False)) + if as_json: + print(json.dumps(results, indent=2)) + return + dry_label = c_yellow(" (DRY-RUN)") if getattr(args, "dry_run", False) else "" + print("\n" + c_bold(f"=== AUTO-SPAWN WORKER RECONCILIATION{dry_label} ===")) + for r in results: + n = r.get("node") + action = r.get("action") + if action == "spawned": + print(c_green(f" ✔ @{n:<5} : SPAWNED session '{r.get('session')}' ({r.get('work_type')}: {r.get('work_name')})")) + elif action == "would_spawn": + print(c_cyan(f" ? @{n:<5} : WOULD SPAWN session '{r.get('session')}' ({r.get('work_type')}: {r.get('work_name')})")) + elif action == "skip": + print(c_dim(f" - @{n:<5} : SKIP ({r.get('reason')})")) + elif action == "idle": + print(c_dim(f" - @{n:<5} : IDLE ({r.get('reason')})")) + elif action == "error": + print(c_red(f" ✘ @{n:<5} : ERROR ({r.get('error')})")) + print() # --------------------------------------------------------------------------- @@ -5449,6 +5471,10 @@ def build_parser(): p_kpi_spawn.add_argument("session", help="Session label") p_kpi_spawn.add_argument("worker_command", help="Command to run in background") + p_kpi_autospawn = kpi_sub.add_parser("auto-spawn", parents=[common], help="Reconcile idle agents and auto-spawn background tmux workers") + p_kpi_autospawn.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node") + p_kpi_autospawn.add_argument("--dry-run", action="store_true", help="Observe and report what would be spawned without executing") + # Domain: CHROMEBOX p_chromebox = subparsers.add_parser("chromebox", parents=[common], help="Agent browser settings-menu toggles (chromebox RPA)") diff --git a/tests/test_kpi.py b/tests/test_kpi.py index aac5605..4756e04 100644 --- a/tests/test_kpi.py +++ b/tests/test_kpi.py @@ -12,8 +12,10 @@ sys.path.insert(0, str(REPO_ROOT / "bin")) import kpi from kpi import ( + auto_spawn_workers, calculate_efficiency, check_node_routes, + find_pending_work_for_node, generate_preservation_advisory, get_agent_dm_metrics, get_agent_kpi, @@ -159,6 +161,45 @@ class TestKPIMetrics(unittest.TestCase): self.assertEqual(call_kwargs["parent"], "dev") self.assertIn("dev-audit-sub", call_kwargs["session_id"]) + def test_find_pending_work_for_node(self): + work = find_pending_work_for_node("dev") + self.assertIsNotNone(work) + self.assertIn(work["type"], ("job", "swarm")) + if work["type"] == "job": + self.assertIn("dev", work["name"]) + self.assertIn("job-dispatch.py", work["cmd"]) + + @patch("kpi.get_agent_tmux_workers") + def test_auto_spawn_workers_skip_active(self, mock_tmux): + mock_tmux.return_value = ["active-worker-1"] + res = auto_spawn_workers(nodes=["dev"], dry_run=True) + self.assertEqual(len(res), 1) + self.assertEqual(res[0]["action"], "skip") + self.assertIn("Active tmux worker already running", res[0]["reason"]) + + @patch("kpi.get_agent_tmux_workers") + def test_auto_spawn_workers_dry_run(self, mock_tmux): + mock_tmux.return_value = [] + res = auto_spawn_workers(nodes=["dev"], dry_run=True) + self.assertEqual(len(res), 1) + self.assertEqual(res[0]["action"], "would_spawn") + self.assertIn("session", res[0]) + self.assertIn("command", res[0]) + + @patch("kpi.spawn_tmux_worker") + @patch("invite_handler.send_loopback_notice") + @patch("kpi.get_agent_tmux_workers") + def test_auto_spawn_workers_live(self, mock_tmux, mock_notice, mock_spawn): + mock_tmux.return_value = [] + mock_spawn.return_value = {"ok": True, "session": "worker-1", "node": "dev"} + mock_notice.return_value = True + + res = auto_spawn_workers(nodes=["dev"], dry_run=False) + self.assertEqual(len(res), 1) + self.assertEqual(res[0]["action"], "spawned") + mock_spawn.assert_called_once() + mock_notice.assert_called_once() + if __name__ == "__main__": unittest.main()