#!/usr/bin/env python3 """kpi.py — NetVM Fleet KPI, Spend Monitor & Runtime Preservation Engine. Monitors: - Calls / DMs dispatched and verified (from dm-log.jsonl) - Token quota spend & remaining (weekly limit % and extra tokens) - Active subagent sessions and tmux muse workers - Uptime vs actual problems fixed (Efficiency Index) - Route health (WARP wireguard, CDP, tmux sockets) - Runtime preservation advisories (guiding agents to offload work to tmux/subagents) """ from __future__ import annotations import argparse import json import os import re import subprocess import sys import time from dataclasses import asdict, dataclass from datetime import datetime, timezone from pathlib import Path from typing import Any, Dict, List, Optional REPO_ROOT = Path(__file__).resolve().parent.parent BIN_DIR = REPO_ROOT / "bin" DM_LOG = REPO_ROOT / "dm-log.jsonl" JOBS_DIR = REPO_ROOT / "jobs" SUBAGENTS_FILE = REPO_ROOT / "subagent-sessions.json" VALID_NODES = ["muse", "pip", "646", "opm", "dev", "def"] @dataclass class AgentKPI: node: str weekly_used_pct: Optional[int] extra_tokens_remaining: str is_blocked: bool calls_sent: int calls_verified: int jobs_assigned: int jobs_completed: int subagents_active: int tmux_workers_active: int uptime_hours: float route_status: str efficiency_index: float efficiency_rating: str preservation_advisory: str def to_dict(self) -> Dict[str, Any]: return asdict(self) def get_agent_dm_metrics(node: str, window_hours: Optional[float] = None) -> Dict[str, int]: """Calculate outbound messages, sends, and verified deliveries from dm-log.jsonl.""" if not DM_LOG.exists(): return {"sent": 0, "verified": 0, "total_events": 0} cutoff = None if window_hours: cutoff = datetime.now(timezone.utc).timestamp() - (window_hours * 3600) sent_ids = set() verified_ids = set() total_events = 0 try: with open(DM_LOG, "r", encoding="utf-8") as f: for line in f: line = line.strip() if not line: continue try: entry = json.loads(line) except Exception: continue if entry.get("agent") != node: continue if cutoff: ts = entry.get("ts") if ts: try: dt = datetime.fromisoformat(ts.replace("Z", "+00:00")) if dt.timestamp() < cutoff: continue except Exception: pass total_events += 1 mid = entry.get("id") etype = entry.get("type") if etype in ("send_start", "send_done"): if mid: sent_ids.add(mid) elif etype == "verified": if mid: verified_ids.add(mid) except Exception: pass return { "sent": len(sent_ids), "verified": len(verified_ids), "total_events": total_events, } def get_agent_job_metrics(node: str) -> Dict[str, int]: """Calculate total jobs assigned and completed for an agent.""" assigned = 0 completed = 0 if not JOBS_DIR.exists(): return {"assigned": 0, "completed": 0} try: for p in JOBS_DIR.glob("*.json"): try: with open(p, "r", encoding="utf-8") as f: data = json.load(f) if data.get("agent") == node: assigned += 1 # If output or status has result if data.get("status") == "completed" or data.get("result"): completed += 1 except Exception: continue except Exception: pass return {"assigned": assigned, "completed": completed} def get_agent_subagent_count(node: str) -> int: """Get active subagent sessions for a node from subagent-sessions.json.""" if not SUBAGENTS_FILE.exists(): return 0 try: with open(SUBAGENTS_FILE, "r", encoding="utf-8") as f: data = json.load(f) if not isinstance(data, dict): return 0 return sum(1 for s in data.values() if s.get("parent") == node and s.get("status") == "active") except Exception: return 0 def get_agent_tmux_workers(node: str) -> List[str]: """Get running tmux sessions for an agent (shared and netns socket).""" sessions = [] # 1. Per-node socket sock = f"/tmp/tmux-{node}.sock" if os.path.exists(sock): try: r = subprocess.run(["tmux", "-S", sock, "list-sessions", "-F", "#{session_name}"], capture_output=True, text=True, timeout=2) if r.returncode == 0 and r.stdout.strip(): sessions.extend(line.strip() for line in r.stdout.splitlines() if line.strip()) except Exception: pass # 2. Shared socket filtering sessions containing node name shared_sock = "/tmp/tmux-muse.sock" if os.path.exists(shared_sock): try: r = subprocess.run(["tmux", "-S", shared_sock, "list-sessions", "-F", "#{session_name}"], capture_output=True, text=True, timeout=2) if r.returncode == 0 and r.stdout.strip(): for s in r.stdout.splitlines(): s = s.strip() if s and (node in s or s.startswith(f"{node}-") or s == "swarm-worker"): if s not in sessions: sessions.append(s) except Exception: pass return sessions def get_agent_uptime_hours(node: str) -> float: """Calculate browser process uptime in hours.""" try: # Search for chromium process matching user-data-dir or node cmd = ["pgrep", "-f", f"chrome-box launch {node}"] r = subprocess.run(cmd, capture_output=True, text=True, timeout=2) pids = r.stdout.strip().split() if not pids: cmd = ["pgrep", "-f", f"profiles/{node}"] r = subprocess.run(cmd, capture_output=True, text=True, timeout=2) pids = r.stdout.strip().split() if pids: pid = pids[0] # Read /proc//stat starttime stat_path = Path(f"/proc/{pid}/stat") if stat_path.exists(): stat_content = stat_path.read_text().split() # field 22 is starttime (in clock ticks after boot) start_ticks = int(stat_content[21]) clk_tck = os.sysconf(os.sysconf_names["SC_CLK_TCK"]) with open("/proc/uptime", "r") as f: uptime_sec = float(f.read().split()[0]) process_age_sec = uptime_sec - (start_ticks / clk_tck) return round(max(0.0, process_age_sec / 3600.0), 2) except Exception: pass return 0.0 def check_node_routes(node: str) -> str: """Check connectivity route for a node (netns + CDP).""" # 1. Check netns netns_path = Path(f"/var/run/netns/warp-{node}") if not netns_path.exists(): return "NO_NETNS" # 2. Check CDP page connection try: try: from approvals import get_node_pages except ImportError: sys.path.insert(0, str(BIN_DIR)) from approvals import get_node_pages pages = get_node_pages(node, timeout=2.0) if pages: return "ONLINE" except Exception: pass return "DEGRADED" def calculate_efficiency( jobs_done: int, subagents_active: int, tmux_workers: int, calls_verified: int, weekly_used_pct: Optional[int], uptime_hours: float, ) -> tuple[float, str]: """ Composite efficiency score. Higher is better: measures actual work produced (jobs + subagents + workers + verified comms) relative to quota burned and uptime elapsed. """ work_units = (jobs_done * 5.0) + (subagents_active * 3.0) + (tmux_workers * 4.0) + (calls_verified * 0.5) burn_cost = max(1.0, (weekly_used_pct or 10) * 0.2) # Base index index = round(work_units / burn_cost, 2) # Classify if uptime_hours > 2.0 and subagents_active == 0 and tmux_workers == 0 and jobs_done == 0: rating = "VANITY_IDLE" elif index >= 3.0: rating = "HIGH_EFFICIENCY" elif index >= 1.0: rating = "PRODUCTIVE" elif index >= 0.4: rating = "MODERATE" else: rating = "LOW_EFFICIENCY" return index, rating def generate_preservation_advisory( node: str, weekly_used_pct: Optional[int], extra_tokens_remaining: str, subagents_active: int, tmux_workers: int, rating: str, ) -> str: """Generate prescriptive runtime preservation instructions for the agent.""" tips = [] pct = weekly_used_pct or 0 if pct >= 95 or "0 tokens left" in extra_tokens_remaining: return "CRITICAL: Quota exhausted. Do NOT send chat messages. Salvage via 'box onboard start --for %s'." % node if pct >= 70: tips.append("Quota > 70%%: Cease prose chatter; offload tasks to background tmux workers.") if subagents_active == 0 and tmux_workers == 0: tips.append("Spawn subagents with isolated context ('box subagent spawn') or tmux muse workers.") if rating in ("VANITY_IDLE", "LOW_EFFICIENCY"): tips.append("Uptime without worker execution drains quota. Mandate: split goals into executable jobs.") if not tips: tips.append("Runtime healthy. Maintain worker-first execution strategy.") return " ".join(tips) def get_agent_kpi(node: str, usage_cache: Optional[Dict[str, Any]] = None) -> AgentKPI: """Collect comprehensive KPI metrics for a single NetVM node.""" # 1. Quota & Usage usage = usage_cache if usage is None: try: try: import invite except ImportError: sys.path.insert(0, str(BIN_DIR)) import invite u = invite.get_usage(node) if isinstance(u, dict) and u.get("ok"): usage = u except Exception: pass if usage is None: usage = {} weekly_pct = usage.get("weekly_used_pct") extra_left = usage.get("additional_left") or usage.get("extra_tokens_remaining") or "Unknown" is_blocked = bool(usage.get("is_blocked")) or (weekly_pct is not None and weekly_pct >= 100 and "0 tokens left" in extra_left) # 2. Activity metrics dm_metrics = get_agent_dm_metrics(node) job_metrics = get_agent_job_metrics(node) subagents = get_agent_subagent_count(node) tmux_sessions = get_agent_tmux_workers(node) uptime = get_agent_uptime_hours(node) route_status = check_node_routes(node) # 3. Efficiency eff_idx, eff_rating = calculate_efficiency( jobs_done=job_metrics["completed"], subagents_active=subagents, tmux_workers=len(tmux_sessions), calls_verified=dm_metrics["verified"], weekly_used_pct=weekly_pct, uptime_hours=uptime, ) # 4. Advisory advisory = generate_preservation_advisory( node=node, weekly_used_pct=weekly_pct, extra_tokens_remaining=extra_left, subagents_active=subagents, tmux_workers=len(tmux_sessions), rating=eff_rating, ) return AgentKPI( node=node, weekly_used_pct=weekly_pct, extra_tokens_remaining=extra_left, is_blocked=is_blocked, calls_sent=dm_metrics["sent"], calls_verified=dm_metrics["verified"], jobs_assigned=job_metrics["assigned"], jobs_completed=job_metrics["completed"], subagents_active=subagents, tmux_workers_active=len(tmux_sessions), uptime_hours=uptime, route_status=route_status, efficiency_index=eff_idx, efficiency_rating=eff_rating, preservation_advisory=advisory, ) def fleet_kpi(nodes: Optional[List[str]] = None) -> Dict[str, AgentKPI]: """Collect KPI metrics across all fleet agents.""" target_nodes = nodes or VALID_NODES # Fetch usage in bulk usage_map = {} try: import invite raw_usage = invite.fleet_usage(target_nodes) if isinstance(raw_usage, dict): usage_map = raw_usage except Exception: pass results = {} for n in target_nodes: results[n] = get_agent_kpi(n, usage_cache=usage_map.get(n)) return results def get_live_advisory_block(node: str) -> str: """Generate Markdown prompt envelope block ready for job injection.""" kpi = get_agent_kpi(node) quota_str = f"{kpi.weekly_used_pct}% weekly limit used" if kpi.weekly_used_pct is not None else "quota active" tokens_str = kpi.extra_tokens_remaining lines = [ "---- BOX PERFORMANCE & RUNTIME ADVISORY ----", f"AGENT: @{kpi.node} | QUOTA: {quota_str} ({tokens_str}) | UPTIME: {kpi.uptime_hours}h", f"WORK UNITS: {kpi.jobs_completed} jobs finished | {kpi.subagents_active} subagents | {kpi.tmux_workers_active} tmux workers", f"EFFICIENCY: {kpi.efficiency_rating} (Index: {kpi.efficiency_index}) | ROUTES: {kpi.route_status}", f"RUNTIME MANDATE: {kpi.preservation_advisory}", "Offload long operations to subagents or tmux muse workers to maximize problem-fixing per token.", ] return "\n".join(lines) 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 clean_session = f"{node}-{session}" if not session.startswith(f"{node}-") else session try: from subagent_tracker import register_session except ImportError: sys.path.insert(0, str(BIN_DIR)) from subagent_tracker import register_session muse_tmux = BIN_DIR / "muse-tmux.py" if not muse_tmux.exists(): return {"ok": False, "error": "muse-tmux.py not found"} # Execute via muse-tmux.py cmd = [ sys.executable, str(muse_tmux), "new", clean_session, "--node", node, "--command", command, ] res = subprocess.run(cmd, capture_output=True, text=True, timeout=10) if res.returncode != 0: # Fallback to shared socket cmd_shared = [ sys.executable, str(muse_tmux), "new", clean_session, "--command", command, ] res = subprocess.run(cmd_shared, capture_output=True, text=True, timeout=10) if res.returncode != 0: return {"ok": False, "error": res.stderr.strip() or res.stdout.strip()} # Register in subagent tracker sid = f"tmux-{clean_session}-{int(time.time())}" register_session(parent=node, session_id=sid, title=f"tmux-worker-{clean_session}", prompt=command) return { "ok": True, "node": node, "session": clean_session, "session_id": sid, "command": command, "message": f"Spawned tmux worker '{clean_session}' for @{node}. Running in background.", } def main(): parser = argparse.ArgumentParser(description="NetVM Fleet KPI, Spend Monitor & Runtime Preservation Engine") subparsers = parser.add_subparsers(dest="command") p_status = subparsers.add_parser("status", help="Show fleet KPI metrics table") p_status.add_argument("--node", choices=VALID_NODES, help="Filter by node") p_status.add_argument("--json", action="store_true", help="Emit JSON output") p_report = subparsers.add_parser("report", help="Detailed KPI report for a specific node") p_report.add_argument("node", choices=VALID_NODES, help="Target node") p_report.add_argument("--json", action="store_true") p_routes = subparsers.add_parser("routes", help="Verify network and CDP routes across nodes") p_routes.add_argument("--json", action="store_true") p_block = subparsers.add_parser("prompt-block", help="Generate live prompt envelope block for node") p_block.add_argument("node", choices=VALID_NODES, help="Target node") p_spawn = subparsers.add_parser("spawn-worker", help="Spawn autonomous background tmux worker session") 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") args = parser.parse_args() if args.command in (None, "status"): nodes = [args.node] if getattr(args, "node", None) else VALID_NODES kpis = fleet_kpi(nodes) if getattr(args, "json", False): print(json.dumps({k: v.to_dict() for k, v in kpis.items()}, indent=2)) return print("\n=== NETVM FLEET KPI & RUNTIME PRESERVATION DASHBOARD ===\n") header = f"{'NODE':<6} {'QUOTA':<10} {'CALLS':<12} {'JOBS':<10} {'SUBAGENTS':<11} {'TMUX':<6} {'UPTIME':<8} {'ROUTE':<9} {'EFFICIENCY':<15}" sep = f"{'────':<6} {'─────────':<10} {'───────────':<12} {'─────────':<10} {'──────────':<11} {'────':<6} {'──────':<8} {'───────':<9} {'──────────────':<15}" print(header) print(sep) for n in nodes: k = kpis.get(n) if not k: continue q_str = f"{k.weekly_used_pct}%" if k.weekly_used_pct is not None else "Active" c_str = f"{k.calls_sent} ({k.calls_verified}v)" j_str = f"{k.jobs_completed}/{k.jobs_assigned}" sub_str = str(k.subagents_active) tmux_str = str(k.tmux_workers_active) up_str = f"{k.uptime_hours}h" print(f"{k.node:<6} {q_str:<10} {c_str:<12} {j_str:<10} {sub_str:<11} {tmux_str:<6} {up_str:<8} {k.route_status:<9} {k.efficiency_rating:<15}") print("\nRun 'box kpi report ' for prescriptive runtime preservation advisories.\n") elif args.command == "report": kpi = get_agent_kpi(args.node) if args.json: print(json.dumps(kpi.to_dict(), indent=2)) return print(f"\n=== KPI & RUNTIME REPORT: @{kpi.node.upper()} ===") print(f" Weekly Quota: {kpi.weekly_used_pct}% used") print(f" Extra Tokens: {kpi.extra_tokens_remaining}") print(f" Blocked Status: {'YES (LIMIT REACHED)' if kpi.is_blocked else 'NO (HEALTHY)'}") print(f" Messages / Calls: {kpi.calls_sent} sent ({kpi.calls_verified} verified delivered)") print(f" Jobs Dispatched: {kpi.jobs_completed} completed / {kpi.jobs_assigned} assigned") print(f" Active Subagents: {kpi.subagents_active}") print(f" Active Tmux Workers:{kpi.tmux_workers_active}") print(f" Process Uptime: {kpi.uptime_hours} hours") print(f" Route Health: {kpi.route_status}") print(f" Efficiency Index: {kpi.efficiency_index} ({kpi.efficiency_rating})") print(f"\n [RUNTIME PRESERVATION ADVISORY]\n {kpi.preservation_advisory}\n") elif args.command == "routes": routes = {n: check_node_routes(n) for n in VALID_NODES} if args.json: print(json.dumps(routes, indent=2)) else: print("\n=== NETVM ROUTE HEALTH ===") for n, st in routes.items(): print(f" @{n:<6} : {st}") print() elif args.command == "prompt-block": print(get_live_advisory_block(args.node)) elif args.command == "spawn-worker": res = spawn_tmux_worker(args.node, args.session, args.worker_command) if args.json: print(json.dumps(res, indent=2)) else: if res.get("ok"): print(f"✔ {res.get('message')}") else: print(f"✘ Failed to spawn worker: {res.get('error')}", file=sys.stderr) sys.exit(1) if __name__ == "__main__": main()