2026-10-04 23:40:04 +00:00
|
|
|
#!/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
|
|
|
|
|
|
|
|
|
|
|
2026-10-07 00:25:51 +00:00
|
|
|
DEFAULT_TTL_SECONDS = 3600 # 1 hour idle TTL
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def prune_stale_sessions(ttl_seconds=DEFAULT_TTL_SECONDS):
|
|
|
|
|
"""Archive active sessions whose last activity exceeds ttl_seconds."""
|
|
|
|
|
data = load_sessions()
|
|
|
|
|
now = datetime.now(timezone.utc)
|
|
|
|
|
changed = False
|
|
|
|
|
for sid, s in data.items():
|
|
|
|
|
if s.get("status") == "active":
|
|
|
|
|
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"))
|
|
|
|
|
if dt.tzinfo is None:
|
|
|
|
|
dt = dt.replace(tzinfo=timezone.utc)
|
|
|
|
|
if (now - dt).total_seconds() >= ttl_seconds:
|
|
|
|
|
s["status"] = "archived"
|
|
|
|
|
s["archived_at"] = utcnow()
|
|
|
|
|
s["archive_reason"] = f"idle_ttl_exceeded_{ttl_seconds}s"
|
|
|
|
|
changed = True
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
|
|
|
|
if changed:
|
|
|
|
|
save_sessions(data)
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
def get_active_sessions(parent=None, auto_prune=True, ttl_seconds=DEFAULT_TTL_SECONDS):
|
2026-10-04 23:40:04 +00:00
|
|
|
"""Retrieve all active subagent sessions, optionally filtered by parent."""
|
2026-10-07 00:25:51 +00:00
|
|
|
if auto_prune:
|
|
|
|
|
try:
|
|
|
|
|
prune_stale_sessions(ttl_seconds=ttl_seconds)
|
|
|
|
|
except Exception:
|
|
|
|
|
pass
|
2026-10-04 23:40:04 +00:00
|
|
|
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)
|
|
|
|
|
|
|
|
|
|
|
2026-10-07 00:25:51 +00:00
|
|
|
def archive_session(session_id, reason=None):
|
|
|
|
|
"""Mark a subagent session archived."""
|
|
|
|
|
kwargs = {"status": "archived", "archived_at": utcnow()}
|
|
|
|
|
if reason:
|
|
|
|
|
kwargs["archive_reason"] = reason
|
|
|
|
|
return update_session(session_id, **kwargs)
|
|
|
|
|
|
|
|
|
|
|
2026-10-04 23:40:04 +00:00
|
|
|
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")
|