96 lines
2.6 KiB
Python
96 lines
2.6 KiB
Python
|
|
#!/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")
|