Files
box/bin/kpi.py
T

709 lines
26 KiB
Python
Raw Normal View History

#!/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/<pid>/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
is_bonus_empty = ("0 tokens left" in extra_tokens_remaining) or (not extra_tokens_remaining)
if pct >= 95 and is_bonus_empty:
return "CRITICAL: Quota exhausted. Salvage via 'box onboard start <new_node> --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)
NODE_SIDECHATS = {
"646": "646 tasks",
"opm": "heartbeat",
"pip": "646-pip-coord",
"dev": "dev-coord",
"def": "def-coord",
"muse": "646-muse-coord",
}
def find_pending_work_for_node(node: str) -> Optional[Dict[str, Any]]:
"""Find assigned pending job or swarm slot for an agent node."""
# 1. Look for node-specific auto-work jobs
if JOBS_DIR.exists():
candidates = sorted(list(JOBS_DIR.glob(f"auto-work-{node}-*.json")) + list(JOBS_DIR.glob(f"{node}-*.json")))
for c in candidates:
try:
with open(c, "r", encoding="utf-8") as f:
data = json.load(f)
job_agent = data.get("agent")
if job_agent and job_agent != node:
continue
job_name = c.stem
return {
"type": "job",
"name": job_name,
"path": str(c),
"cmd": f"{sys.executable} {BIN_DIR}/job-dispatch.py {job_name}",
}
except Exception:
continue
# 2. Check pending swarm slots
try:
from swarm_worker.poller import find_pending_slots
slots = find_pending_slots()
if slots:
slot = slots[0]
sw_id = slot.get("swarm_id", "swarm")
idx = slot.get("slot_index", 0)
return {
"type": "swarm",
"name": f"swarm-{sw_id}-s{idx}",
"path": None,
"cmd": f"{sys.executable} {BIN_DIR}/swarm_worker/daemon.py",
}
except Exception:
pass
return None
def auto_spawn_workers(nodes: Optional[List[str]] = None, dry_run: bool = False) -> List[Dict[str, Any]]:
"""Reconcile idle agents and auto-spawn background tmux workers to execute pending work."""
target_nodes = nodes or VALID_NODES
results = []
for node in target_nodes:
# Check active tmux workers for this node
active_tmux = len(get_agent_tmux_workers(node))
if active_tmux > 0:
results.append({
"node": node,
"action": "skip",
"reason": f"Active tmux worker already running ({active_tmux})",
})
continue
# Check work availability
work = find_pending_work_for_node(node)
if not work:
results.append({
"node": node,
"action": "idle",
"reason": "No pending jobs or swarm slots",
})
continue
session_label = f"worker-{work['name'][:18]}"
cmd_to_run = f"{work['cmd']} > /tmp/tmux-{node}-{session_label}.log 2>&1"
if dry_run:
results.append({
"node": node,
"action": "would_spawn",
"session": session_label,
"work_type": work["type"],
"work_name": work["name"],
"command": work["cmd"],
})
continue
# Execute spawn
spawn_res = spawn_tmux_worker(node, session_label, cmd_to_run)
if spawn_res.get("ok"):
# Send sidechat notification
try:
from invite_handler import send_loopback_notice
chat = NODE_SIDECHATS.get(node, "646 tasks")
msg = f"[BOX-AUTO-WORKER] Spawned background tmux worker '{session_label}' executing {work['type']} ({work['name']}). Logs at /tmp/tmux-{node}-{session_label}.log"
send_loopback_notice(recipient=node, target=chat, message=msg)
except Exception:
pass
results.append({
"node": node,
"action": "spawned",
"session": session_label,
"work_type": work["type"],
"work_name": work["name"],
"command": work["cmd"],
})
else:
results.append({
"node": node,
"action": "error",
"error": spawn_res.get("error", "Unknown spawn error"),
})
return results
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_autospawn = subparsers.add_parser("auto-spawn", help="Auto-spawn background tmux workers for idle nodes with pending work")
p_autospawn.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node")
p_autospawn.add_argument("--dry-run", action="store_true", help="Report what would be spawned without executing")
p_autospawn.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 <node>' 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)
elif args.command == "auto-spawn":
nodes = [args.node] if getattr(args, "node", None) else None
results = auto_spawn_workers(nodes=nodes, dry_run=args.dry_run)
if args.json:
print(json.dumps(results, indent=2))
return
print(f"\n=== AUTO-SPAWN WORKER RECONCILIATION {'(DRY-RUN)' if args.dry_run else ''} ===")
for r in results:
n = r.get("node")
act = r.get("action")
if act == "spawned":
print(f" ✔ @{n:<5} : SPAWNED session '{r.get('session')}' ({r.get('work_type')}: {r.get('work_name')})")
elif act == "would_spawn":
print(f" ? @{n:<5} : WOULD SPAWN session '{r.get('session')}' ({r.get('work_type')}: {r.get('work_name')})")
elif act == "skip":
print(f" - @{n:<5} : SKIP ({r.get('reason')})")
elif act == "idle":
print(f" - @{n:<5} : IDLE ({r.get('reason')})")
elif act == "error":
print(f" ✘ @{n:<5} : ERROR ({r.get('error')})")
print()
if __name__ == "__main__":
main()