diff --git a/bin/self_main_loop.py b/bin/self_main_loop.py index 4c6a71a..04c39b9 100755 --- a/bin/self_main_loop.py +++ b/bin/self_main_loop.py @@ -378,6 +378,85 @@ def compose_dual_digest(agent, main_new, sc_activity): return "\n".join(lines)[:DIGEST_MAX] +def check_subagents(agent, cfg, threads_meta): + """Inspect active child subagents spawned by agent. + + Returns list of alert strings to deliver to the parent's prompting sidechat: + - New output from subagent: prompts parent with snippet and thread view command. + - 15-minute stall: warns parent that subagent has had no activity. + """ + alerts = [] + try: + import subagent_tracker + import muse_hybrid + active = subagent_tracker.get_active_sessions(parent=agent) + if not active or not threads_meta: + return alerts + + meta_by_id = {t.get("session_id"): t for t in threads_meta if isinstance(t, dict)} + now = datetime.now(timezone.utc) + + for s in active: + sid = s.get("session_id") + title = s.get("title") or "subagent" + last_wm = s.get("last_watermark") + t_meta = meta_by_id.get(sid) + + if not t_meta: + continue + + current_updated = t_meta.get("updated") + # 1. New output from subagent + if current_updated and current_updated != last_wm: + if last_wm is None: + # Seed initial watermark so we only trigger on subsequent turns + subagent_tracker.update_session(sid, last_watermark=current_updated, last_activity_at=utcnow()) + continue + + history, _ = muse_hybrid.get_history(agent, thread_id=sid, limit=2) + snippet = "" + if history: + for m in reversed(history): + if m.get("role") == "assistant": + snippet = m.get("text", "").strip()[:180] + break + + subagent_tracker.update_session( + sid, + last_watermark=current_updated, + last_activity_at=utcnow(), + timeout_warned=False + ) + + short_id = sid[:8] + alert_text = f"[SUBAGENT-UPDATE] Subagent '{title}' ({short_id}) has posted new output." + if snippet: + alert_text += f"\nSnippet: {snippet}..." + alert_text += f"\nRun `box thread view {sid}` to review and aggregate deliverables." + alerts.append(alert_text) + + # 2. 15-minute stall watchdog + elif not s.get("timeout_warned"): + last_act = s.get("last_activity_at") or s.get("spawned_at") + if last_act: + try: + dt = datetime.fromisoformat(last_act.replace("Z", "+00:00")) + idle_sec = (now - dt).total_seconds() + if idle_sec >= 900: # 15 minutes + subagent_tracker.update_session(sid, timeout_warned=True) + short_id = sid[:8] + alerts.append( + f"[WARN] Subagent '{title}' ({short_id}) has had no activity for 15+ minutes. " + f"Check status with `box thread view {sid}` or consider respawning." + ) + except Exception: + pass + except Exception as e: + log("%s: subagent check error: %s" % (agent, e)) + + return alerts + + def do_check(only_agent=None): st = load_state() cfg = get_config(st) @@ -417,14 +496,22 @@ def do_check(only_agent=None): threads_meta, _ = muse_hybrid.get_threads(agent) sc_activity = check_sidechats(agent, cfg, wm, threads_meta) + # 3. Subagent reactive monitoring check + subagent_alerts = check_subagents(agent, cfg, threads_meta) + agents_state[agent] = wm - total_new = len(main_new) + sum(len(a["messages"]) for a in sc_activity) + total_new = len(main_new) + sum(len(a["messages"]) for a in sc_activity) + len(subagent_alerts) if total_new == 0: results[agent] = {"ok": True, "new": 0} continue - digest = compose_dual_digest(agent, main_new, sc_activity) + parts = [] + if len(main_new) > 0 or len(sc_activity) > 0: + parts.append(compose_dual_digest(agent, main_new, sc_activity)) + if subagent_alerts: + parts.extend(subagent_alerts) + digest = "\n\n".join(parts)[:DIGEST_MAX * 2] sidechat = cfg["prompt_sidechat"].get(agent) if not sidechat: log("%s: no prompting sidechat configured, skipping" % agent) diff --git a/bin/subagent_tracker.py b/bin/subagent_tracker.py new file mode 100644 index 0000000..c981046 --- /dev/null +++ b/bin/subagent_tracker.py @@ -0,0 +1,95 @@ +#!/usr/bin/env python3 +""" +subagent_tracker.py — Registry and lifecycle tracker for fleet subagents. + +Maintains subagent-sessions.json tracking active child sessions spawned +by parent agents (646, pip, muse, opm), their last observed watermarks, +and timeout/stall states for reactive wakeups by self_main_loop.py. +""" + +import json +import os +import sys +from datetime import datetime, timezone +from pathlib import Path + +BASE = Path("/home/super/Projects/NetVM") +SESSIONS_FILE = BASE / "subagent-sessions.json" + + +def utcnow(): + return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") + + +def load_sessions(): + if not SESSIONS_FILE.exists(): + return {} + try: + with open(SESSIONS_FILE, "r", encoding="utf-8") as f: + data = json.load(f) + return data if isinstance(data, dict) else {} + except Exception: + return {} + + +def save_sessions(data): + tmp = f"{SESSIONS_FILE}.tmp.{os.getpid()}" + with open(tmp, "w", encoding="utf-8") as f: + json.dump(data, f, indent=2) + os.replace(tmp, SESSIONS_FILE) + + +def register_session(parent, session_id, title=None, prompt=None): + """Register newly spawned subagent session.""" + data = load_sessions() + now_iso = utcnow() + entry = { + "session_id": session_id, + "parent": parent, + "title": title or "subagent", + "prompt": prompt or "", + "spawned_at": now_iso, + "last_watermark": None, + "status": "active", + "timeout_warned": False, + "last_activity_at": now_iso, + } + data[session_id] = entry + save_sessions(data) + return entry + + +def get_active_sessions(parent=None): + """Retrieve all active subagent sessions, optionally filtered by parent.""" + data = load_sessions() + results = [] + for s in data.values(): + if s.get("status") == "active": + if parent is None or s.get("parent") == parent: + results.append(s) + return results + + +def update_session(session_id, **kwargs): + """Update fields on a tracked subagent session.""" + data = load_sessions() + if session_id in data: + data[session_id].update(kwargs) + save_sessions(data) + return data[session_id] + return None + + +def complete_session(session_id, note=None): + """Mark a subagent session completed.""" + kwargs = {"status": "completed", "completed_at": utcnow()} + if note: + kwargs["completion_note"] = note + return update_session(session_id, **kwargs) + + +if __name__ == "__main__": + if len(sys.argv) > 1 and sys.argv[1] == "list": + print(json.dumps(load_sessions(), indent=2)) + else: + print("Usage: subagent_tracker.py list") diff --git a/bin/super-cli.py b/bin/super-cli.py index 302c794..efa31fe 100755 --- a/bin/super-cli.py +++ b/bin/super-cli.py @@ -2530,6 +2530,12 @@ def cmd_deploy(args): session_id = res.get("session_id") print(c_green(f"✔ Subagent session spawned: {session_id}")) + try: + import subagent_tracker + subagent_tracker.register_session(agent, session_id, title=title, prompt=prompt) + except Exception: + pass + if prompt: print(f" Dispatching task prompt to subagent (waiting up to {wait}s)...") send_res, send_err = muse_hybrid.send_message(agent, prompt, thread_id=session_id, wait=wait) @@ -3404,6 +3410,15 @@ def build_parser(): p_dep_sub.add_argument("prompt", help="Task prompt for the subagent") p_dep_sub.add_argument("--wait", type=int, default=30, help="Seconds to wait for subagent response") + # Domain: subagent (Direct alias for spawning & managing subagents) + p_subagent = subparsers.add_parser("subagent", parents=[common], help="Spawn sub-agent session and dispatch task") + subagent_sub = p_subagent.add_subparsers(dest="action") + p_sub_spawn = subagent_sub.add_parser("spawn", parents=[common], help="Spawn sub-agent session and dispatch task") + p_sub_spawn.add_argument("--agent", required=True, choices=VALID_NODES, help="Agent node to spawn subagent on") + p_sub_spawn.add_argument("--title", default="subagent-task", help="Title for the subagent session") + p_sub_spawn.add_argument("prompt", help="Task prompt for the subagent") + p_sub_spawn.add_argument("--wait", type=int, default=30, help="Seconds to wait for subagent response") + # Domain: muse (fast headless gateway via muse-cli-node with isolated per-node Cloudflare WARP egress) p_muse = subparsers.add_parser("muse", parents=[common], help="Direct headless gateway client (muse-cli-node)") p_muse.add_argument("node", choices=VALID_NODES, help="Target agent node") @@ -3577,7 +3592,9 @@ def main(): cmd_loop_strat(args) elif args.domain == "vars": cmd_loop_vars(args) - elif args.domain == "deploy": + elif args.domain in ("deploy", "subagent"): + if args.domain == "subagent": + args.action = "subagent" cmd_deploy(args) elif args.domain == "muse": cmd = [str(BIN_DIR / "muse-cli-node"), args.node] + (args.muse_args or []) diff --git a/job-sidechats.json b/job-sidechats.json index 0cb7325..f9e8b7e 100644 --- a/job-sidechats.json +++ b/job-sidechats.json @@ -111,6 +111,21 @@ "agent": "pip", "description": "pip side of 646-pip coordination pair" }, + "pip-opm": { + "thread_uuid": "75feb3a2-be36-499a-b314-bfd5803f1db0", + "agent": "pip", + "description": "default pip-opm coordination target" + }, + "pip-opm@pip": { + "thread_uuid": "75feb3a2-be36-499a-b314-bfd5803f1db0", + "agent": "pip", + "description": "pip side of pip-opm coordination pair" + }, + "pip-opm@opm": { + "thread_uuid": "769955af-21ad-4689-8909-0f75b8f1035c", + "agent": "opm", + "description": "opm side of pip-opm coordination pair" + }, "pipe-a0d377": { "thread_uuid": "cc0aee8f-531a-40af-9e02-27459f57b3ae", "agent": "opm", @@ -130,5 +145,16 @@ "thread_uuid": "8f9ae8e7-dd3a-46e6-a5ac-3b1bb7cf6d66", "agent": "opm", "created_at": "2026-10-04T23:19:20.218056+00:00" + }, + "onboarding-test-dev": { + "thread_uuid": "f7741827-7b44-45b7-bfb3-c08f8889e93a", + "agent": "dev", + "created_at": "2026-10-04T23:35:34.711793+00:00" + }, + "onboarding-dev": { + "thread_uuid": "31ceab7e-7579-4acc-bcf5-1a10c081aa21", + "agent": "dev", + "created_at": "2026-10-04T23:37:01+00:00", + "note": "onboarding-alpha-probe" } } \ No newline at end of file