Files

136 lines
4.1 KiB
Python
Raw Permalink Normal View History

#!/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
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):
"""Retrieve all active subagent sessions, optionally filtered by parent."""
if auto_prune:
try:
prune_stale_sessions(ttl_seconds=ttl_seconds)
except Exception:
pass
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)
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)
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")