feat(kpi): expand KPI runtime monitoring, prompt advisory envelopes, and missing field resiliency
This commit is contained in:
+11
-1
@@ -237,7 +237,17 @@ def get_usage(node, timeout=8.0):
|
||||
return {"ok": False, "node": node,
|
||||
"error": "empty settings dialog"}
|
||||
out = {"ok": True, "node": node}
|
||||
out.update(parse_usage_text(text))
|
||||
parsed = parse_usage_text(text)
|
||||
stats_loaded = bool(
|
||||
parsed.get("weekly_reset")
|
||||
or parsed.get("weekly_used_pct") is not None
|
||||
or parsed.get("has_redeemed")
|
||||
or parsed.get("additional_left")
|
||||
)
|
||||
out.update(parsed)
|
||||
out["stats_loaded"] = stats_loaded
|
||||
if not stats_loaded:
|
||||
out["note"] = "Usage stats did not render in Settings dialog"
|
||||
return out
|
||||
finally:
|
||||
_escape(ws)
|
||||
|
||||
+241
-48
@@ -55,11 +55,49 @@ class RedemptionResult:
|
||||
reason: Optional[str]
|
||||
detail: Optional[str]
|
||||
method_used: str
|
||||
field_missing: bool = False
|
||||
loopback_notified: bool = False
|
||||
loopback_detail: Optional[str] = None
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return asdict(self)
|
||||
|
||||
|
||||
def send_loopback_notice(
|
||||
recipient: str,
|
||||
target: str,
|
||||
message: str,
|
||||
sender: str = "super",
|
||||
timeout: float = 15.0,
|
||||
) -> bool:
|
||||
"""Send a loopback notification DM via bin/dm.py to a sidechat without blocking or failing."""
|
||||
try:
|
||||
from pathlib import Path
|
||||
bin_dir = Path(__file__).resolve().parent
|
||||
dm_script = bin_dir / "dm.py"
|
||||
if not dm_script.exists():
|
||||
return False
|
||||
import subprocess
|
||||
cmd = [
|
||||
sys.executable,
|
||||
str(dm_script),
|
||||
"send",
|
||||
"--agent",
|
||||
sender,
|
||||
"--to",
|
||||
recipient,
|
||||
"--target",
|
||||
target,
|
||||
message,
|
||||
]
|
||||
if target == "main":
|
||||
cmd.insert(-1, "--allow-main-chat")
|
||||
res = subprocess.run(cmd, capture_output=True, text=True, timeout=timeout)
|
||||
return res.returncode == 0
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
class InviteHandler:
|
||||
"""Handles invite codes discovery, inspection, and redemption for NetVM nodes."""
|
||||
|
||||
@@ -68,6 +106,15 @@ class InviteHandler:
|
||||
self.timeout = timeout
|
||||
self.ws = None
|
||||
|
||||
def _dispatch_loopback(self, res: RedemptionResult, agent: str, target: str) -> bool:
|
||||
"""Post a loopback notification DM via box dm."""
|
||||
msg = f"[BOX-INVITE-LOOPBACK] Node @{self.node} redeem status: {res.status}. Reason: {res.reason or 'unknown'}. Detail: {res.detail or '-'}"
|
||||
ok = send_loopback_notice(recipient=agent, target=target, message=msg)
|
||||
res.loopback_notified = ok
|
||||
res.loopback_detail = f"Notified @{agent}/{target}" if ok else f"Failed to notify @{agent}/{target}"
|
||||
return ok
|
||||
|
||||
|
||||
def connect(self) -> InviteHandler:
|
||||
if self.ws is None:
|
||||
self.ws, _ = get_cdp_ws(self.node, timeout=self.timeout)
|
||||
@@ -262,87 +309,169 @@ class InviteHandler:
|
||||
method_used="api",
|
||||
)
|
||||
|
||||
def redeem_code_dom(self, code: str) -> RedemptionResult:
|
||||
"""Redeem invite code using Settings Menu RPA -> Settings -> Redeem Invite Code dialog."""
|
||||
def redeem_code_dom(
|
||||
self,
|
||||
code: str,
|
||||
notify_target: Optional[str] = None,
|
||||
notify_agent: Optional[str] = None,
|
||||
) -> RedemptionResult:
|
||||
"""Redeem invite code using Settings Menu RPA -> Settings -> Redeem Invite Code dialog.
|
||||
Gracefully handles missing entrypoint rows, missing input fields, and async loading races.
|
||||
Falls back to in-page API automatically if DOM fields are absent, and dispatches DM loopback if requested."""
|
||||
self.connect()
|
||||
clean_code = code.strip().upper()
|
||||
with SettingsRPA(self.node, timeout=self.timeout) as rpa:
|
||||
# 1. Open Settings dialog
|
||||
opened = rpa.open_settings_dialog()
|
||||
if not opened:
|
||||
return RedemptionResult(
|
||||
# If dialog failed to open, try API fallback directly
|
||||
api_res = self.redeem_code_api(clean_code)
|
||||
if api_res.success:
|
||||
api_res.detail = f"Settings dialog could not open; redemption completed via API fallback ({api_res.detail or ''})".strip()
|
||||
return api_res
|
||||
|
||||
res = RedemptionResult(
|
||||
target_node=self.node,
|
||||
code=clean_code,
|
||||
success=False,
|
||||
status="error",
|
||||
reason="settings_open_failed",
|
||||
detail="Could not open Settings dialog via RPA",
|
||||
detail=f"Could not open Settings dialog via RPA on @{self.node}. API fallback: {api_res.detail or api_res.reason or 'failed'}",
|
||||
method_used="dom",
|
||||
field_missing=True,
|
||||
)
|
||||
if notify_target and notify_agent:
|
||||
self._dispatch_loopback(res, notify_agent, notify_target)
|
||||
return res
|
||||
|
||||
rpa.select_tab("General")
|
||||
time.sleep(0.3)
|
||||
|
||||
# 2. Check for "Redeem invite code" entry
|
||||
# 2. Check for "Redeem invite code" entry (poll up to 2.5s for React rendering)
|
||||
js_find_and_click = """(() => {
|
||||
const dialog = document.querySelector('[role="dialog"]');
|
||||
if (!dialog) return {found: false};
|
||||
if (!dialog) return {found: false, dialog_present: false};
|
||||
const items = Array.from(dialog.querySelectorAll('button, div, span'));
|
||||
const redeemItem = items.find(el => (el.innerText || '').trim() === 'Redeem invite code');
|
||||
if (redeemItem) {
|
||||
redeemItem.click();
|
||||
return {found: true};
|
||||
return {found: true, dialog_present: true};
|
||||
}
|
||||
return {found: false, text: dialog.innerText};
|
||||
return {
|
||||
found: false,
|
||||
dialog_present: true,
|
||||
text: dialog.innerText || '',
|
||||
has_additional: (dialog.innerText || '').includes('Additional tokens')
|
||||
};
|
||||
})()"""
|
||||
click_res = cdp_evaluate(self.ws, js_find_and_click)
|
||||
click_res = None
|
||||
deadline = time.time() + 2.5
|
||||
while time.time() < deadline:
|
||||
click_res = cdp_evaluate(self.ws, js_find_and_click)
|
||||
if click_res and click_res.get("found"):
|
||||
break
|
||||
time.sleep(0.3)
|
||||
|
||||
if not click_res or not click_res.get("found"):
|
||||
# "Redeem invite code" field is missing in General settings!
|
||||
rpa.close_settings_dialog()
|
||||
# If entrypoint missing, likely already redeemed
|
||||
api_check = self.find_code_api()
|
||||
if api_check.has_redeemed:
|
||||
return RedemptionResult(
|
||||
target_node=self.node,
|
||||
code=clean_code,
|
||||
success=False,
|
||||
status="already_redeemed",
|
||||
reason="already_redeemed",
|
||||
detail="Node has already redeemed an invite code (entrypoint hidden)",
|
||||
method_used="dom",
|
||||
)
|
||||
return RedemptionResult(
|
||||
|
||||
# Attempt automatic API fallback first
|
||||
api_res = self.redeem_code_api(clean_code)
|
||||
if api_res.success:
|
||||
api_res.detail = f"Redeem field was missing in Settings DOM; redeemed successfully via API fallback! ({api_res.detail or ''})".strip()
|
||||
return api_res
|
||||
|
||||
# Diagnose why field is missing
|
||||
diag_text = click_res.get("text", "") if click_res else ""
|
||||
has_extra = click_res.get("has_additional", False) if click_res else False
|
||||
|
||||
is_already = has_extra or "Additional tokens" in diag_text
|
||||
if not is_already:
|
||||
try:
|
||||
api_check = self.find_code_api()
|
||||
is_already = bool(api_check.has_redeemed)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if is_already:
|
||||
status = "already_redeemed"
|
||||
reason = "already_redeemed"
|
||||
detail = f"Node @{self.node} has already redeemed an invite code (entrypoint hidden by active Additional tokens ticker)."
|
||||
elif api_res.reason in ["window_expired", "already_redeemed", "invalid", "used_up"]:
|
||||
status = api_res.status
|
||||
reason = api_res.reason
|
||||
detail = f"Redeem invite code field not present in Settings on @{self.node}: server reports {api_res.reason} ({api_res.detail or ''})."
|
||||
else:
|
||||
status = "entrypoint_not_found"
|
||||
reason = "field_missing"
|
||||
detail = f"Redeem invite code field missing in General settings on @{self.node} (account may be past 48h onboarding window or already redeemed)."
|
||||
|
||||
res = RedemptionResult(
|
||||
target_node=self.node,
|
||||
code=clean_code,
|
||||
success=False,
|
||||
status="entrypoint_not_found",
|
||||
reason="not_eligible",
|
||||
detail="Redeem invite code item not present in General settings",
|
||||
status=status,
|
||||
reason=reason,
|
||||
detail=detail,
|
||||
method_used="dom",
|
||||
field_missing=True,
|
||||
)
|
||||
if notify_target and notify_agent:
|
||||
self._dispatch_loopback(res, notify_agent, notify_target)
|
||||
return res
|
||||
|
||||
time.sleep(0.8)
|
||||
|
||||
# 3. Handle HatchInviteRedemptionDialog
|
||||
# 3. Handle HatchInviteRedemptionDialog (sub-dialog opened by clicking Redeem)
|
||||
# Poll up to 2.5s for input box to mount
|
||||
js_input_code = f"""(() => {{
|
||||
// Look for dialog titled "Redeem a code"
|
||||
// Look for dialog titled "Redeem a code" or input with label
|
||||
const input = document.querySelector('input[aria-label="Invite code"]') ||
|
||||
document.querySelector('[role="dialog"] input[type="text"]');
|
||||
if (!input) return {{found_input: false}};
|
||||
|
||||
// Type code into input
|
||||
input.focus();
|
||||
input.value = {json.dumps(clean_code)};
|
||||
input.dispatchEvent(new Event('input', {{bubbles: true}}));
|
||||
input.dispatchEvent(new Event('change', {{bubbles: true}}));
|
||||
|
||||
// Submit via Enter key
|
||||
input.dispatchEvent(new KeyboardEvent('keydown', {{key: 'Enter', code: 'Enter', keyCode: 13, which: 13, bubbles: true}}));
|
||||
return {{found_input: true}};
|
||||
}})()"""
|
||||
input_res = cdp_evaluate(self.ws, js_input_code)
|
||||
time.sleep(1.5)
|
||||
input_res = None
|
||||
deadline = time.time() + 2.5
|
||||
while time.time() < deadline:
|
||||
input_res = cdp_evaluate(self.ws, js_input_code)
|
||||
if input_res and input_res.get("found_input"):
|
||||
break
|
||||
time.sleep(0.3)
|
||||
|
||||
# 4. Check outcome in dialog
|
||||
if not input_res or not input_res.get("found_input"):
|
||||
# Subdialog opened or clicked, but input box is absent!
|
||||
cdp_send_escape(self.ws)
|
||||
rpa.close_settings_dialog()
|
||||
|
||||
# Automatic API fallback
|
||||
api_res = self.redeem_code_api(clean_code)
|
||||
if api_res.success:
|
||||
api_res.detail = f"Redemption input box was missing in dialog; redeemed successfully via API fallback! ({api_res.detail or ''})".strip()
|
||||
return api_res
|
||||
|
||||
res = RedemptionResult(
|
||||
target_node=self.node,
|
||||
code=clean_code,
|
||||
success=False,
|
||||
status=api_res.status or "dom_input_missing",
|
||||
reason=api_res.reason or "input_field_missing",
|
||||
detail=f"Invite code input field was not found in redemption dialog on @{self.node}. API fallback: {api_res.detail or api_res.reason or 'failed'}.",
|
||||
method_used="dom",
|
||||
field_missing=True,
|
||||
)
|
||||
if notify_target and notify_agent:
|
||||
self._dispatch_loopback(res, notify_agent, notify_target)
|
||||
return res
|
||||
|
||||
# Input submitted; poll for outcome text
|
||||
time.sleep(1.2)
|
||||
js_check_outcome = """(() => {
|
||||
const dialog = document.querySelector('[role="dialog"]');
|
||||
if (!dialog) return {found: false};
|
||||
@@ -365,7 +494,14 @@ class InviteHandler:
|
||||
}
|
||||
return {status: 'unknown', detail: text};
|
||||
})()"""
|
||||
outcome = cdp_evaluate(self.ws, js_check_outcome)
|
||||
outcome = None
|
||||
deadline = time.time() + 2.5
|
||||
while time.time() < deadline:
|
||||
outcome = cdp_evaluate(self.ws, js_check_outcome)
|
||||
if outcome and outcome.get("status") != "unknown":
|
||||
break
|
||||
time.sleep(0.3)
|
||||
|
||||
rpa.close_settings_dialog()
|
||||
|
||||
if outcome and outcome.get("success"):
|
||||
@@ -379,7 +515,7 @@ class InviteHandler:
|
||||
method_used="dom",
|
||||
)
|
||||
elif outcome and outcome.get("status") in ["already_redeemed", "invalid", "used_up", "window_expired"]:
|
||||
return RedemptionResult(
|
||||
res = RedemptionResult(
|
||||
target_node=self.node,
|
||||
code=clean_code,
|
||||
success=False,
|
||||
@@ -388,19 +524,40 @@ class InviteHandler:
|
||||
detail=outcome.get("detail"),
|
||||
method_used="dom",
|
||||
)
|
||||
if notify_target and notify_agent:
|
||||
self._dispatch_loopback(res, notify_agent, notify_target)
|
||||
return res
|
||||
|
||||
# Fallback to API check if DOM did not confirm status
|
||||
return self.redeem_code_api(clean_code)
|
||||
api_res = self.redeem_code_api(clean_code)
|
||||
if not api_res.success and notify_target and notify_agent:
|
||||
self._dispatch_loopback(api_res, notify_agent, notify_target)
|
||||
return api_res
|
||||
|
||||
def redeem_code(self, code: str, method: str = "auto") -> RedemptionResult:
|
||||
def redeem_code(
|
||||
self,
|
||||
code: str,
|
||||
method: str = "auto",
|
||||
notify_target: Optional[str] = None,
|
||||
notify_agent: Optional[str] = None,
|
||||
) -> RedemptionResult:
|
||||
"""Redeem invite code using specified method ('auto', 'dom', or 'api')."""
|
||||
if method == "dom":
|
||||
return self.redeem_code_dom(code)
|
||||
return self.redeem_code_dom(code, notify_target=notify_target, notify_agent=notify_agent)
|
||||
elif method == "api":
|
||||
return self.redeem_code_api(code)
|
||||
res = self.redeem_code_api(code)
|
||||
if not res.success and notify_target and notify_agent:
|
||||
self._dispatch_loopback(res, notify_agent, notify_target)
|
||||
return res
|
||||
else:
|
||||
# Auto: Validate and attempt via API for reliability
|
||||
return self.redeem_code_api(code)
|
||||
# Auto: Validate and attempt via API for reliability. Fallback to DOM on transport error.
|
||||
res = self.redeem_code_api(code)
|
||||
if not res.success and res.reason == "transport_failure":
|
||||
res = self.redeem_code_dom(code, notify_target=notify_target, notify_agent=notify_agent)
|
||||
elif not res.success and notify_target and notify_agent:
|
||||
self._dispatch_loopback(res, notify_agent, notify_target)
|
||||
return res
|
||||
|
||||
|
||||
|
||||
def scan_fleet_invites(nodes: Optional[List[str]] = None) -> List[Dict[str, Any]]:
|
||||
@@ -454,7 +611,12 @@ def scan_fleet_usage(nodes: Optional[List[str]] = None) -> List[Dict[str, Any]]:
|
||||
return results
|
||||
|
||||
|
||||
def salvage_blocked_node(blocked_node: str = "646", helper_node: Optional[str] = None) -> Dict[str, Any]:
|
||||
def salvage_blocked_node(
|
||||
blocked_node: str = "646",
|
||||
helper_node: Optional[str] = None,
|
||||
notify_target: Optional[str] = None,
|
||||
notify_agent: Optional[str] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""Salvage an out-of-tokens node by identifying its code and redeeming it on an eligible peer."""
|
||||
# 1. Fetch blocked node invite code
|
||||
with InviteHandler(blocked_node) as h_blocked:
|
||||
@@ -483,18 +645,28 @@ def salvage_blocked_node(blocked_node: str = "646", helper_node: Optional[str] =
|
||||
continue
|
||||
|
||||
if not eligible_peer:
|
||||
return {
|
||||
res = {
|
||||
"success": False,
|
||||
"blocked_node": blocked_node,
|
||||
"invite_code": code_to_redeem,
|
||||
"error": "No existing fleet peer is currently eligible (all active peers have already redeemed an invite code). An onboarding agent or fresh client profile must redeem this code.",
|
||||
"code_to_redeem": code_to_redeem,
|
||||
"share_instruction": f"Redeem code '{code_to_redeem}' on a newly provisioned agent to credit 1 billion tokens to {blocked_node}."
|
||||
"share_instruction": f"Redeem code '{code_to_redeem}' on a newly provisioned agent to credit 1 billion tokens to {blocked_node}.",
|
||||
"field_missing": True,
|
||||
"loopback_notified": False,
|
||||
}
|
||||
if notify_target and notify_agent:
|
||||
msg = f"[SALVAGE-NOTICE] Node @{blocked_node} is blocked (code: {code_to_redeem}), but no eligible peer is available. Fresh onboarding required."
|
||||
res["loopback_notified"] = send_loopback_notice(notify_agent, notify_target, msg)
|
||||
return res
|
||||
|
||||
# 3. Redeem on eligible peer
|
||||
with InviteHandler(eligible_peer) as h_peer:
|
||||
redemption = h_peer.redeem_code(code_to_redeem)
|
||||
redemption = h_peer.redeem_code(
|
||||
code_to_redeem,
|
||||
notify_target=notify_target,
|
||||
notify_agent=notify_agent,
|
||||
)
|
||||
|
||||
return {
|
||||
"success": redemption.success,
|
||||
@@ -502,6 +674,8 @@ def salvage_blocked_node(blocked_node: str = "646", helper_node: Optional[str] =
|
||||
"helper_node": eligible_peer,
|
||||
"code_redeemed": code_to_redeem,
|
||||
"redemption_result": redemption.to_dict(),
|
||||
"field_missing": redemption.field_missing,
|
||||
"loopback_notified": redemption.loopback_notified,
|
||||
}
|
||||
|
||||
|
||||
@@ -518,6 +692,8 @@ def main() -> None:
|
||||
p_redeem.add_argument("node", help="Target node to redeem the code on")
|
||||
p_redeem.add_argument("code", help="6-character invite code")
|
||||
p_redeem.add_argument("--method", choices=["auto", "api", "dom"], default="auto")
|
||||
p_redeem.add_argument("--notify-target", default=None, help="Sidechat to notify on loopback")
|
||||
p_redeem.add_argument("--notify-agent", default=None, help="Agent to notify on loopback")
|
||||
p_redeem.add_argument("--json", action="store_true")
|
||||
|
||||
p_list = subparsers.add_parser("list", help="List invite codes across fleet")
|
||||
@@ -526,6 +702,8 @@ def main() -> None:
|
||||
p_salvage = subparsers.add_parser("salvage", help="Salvage a blocked agent (e.g. 646)")
|
||||
p_salvage.add_argument("node", default="646", nargs="?", help="Blocked node (default: 646)")
|
||||
p_salvage.add_argument("--helper", help="Specific helper node to redeem code")
|
||||
p_salvage.add_argument("--notify-target", default=None, help="Sidechat to notify on loopback")
|
||||
p_salvage.add_argument("--notify-agent", default=None, help="Agent to notify on loopback")
|
||||
p_salvage.add_argument("--json", action="store_true")
|
||||
|
||||
args = parser.parse_args()
|
||||
@@ -546,7 +724,12 @@ def main() -> None:
|
||||
|
||||
elif args.command == "redeem":
|
||||
with InviteHandler(args.node) as h:
|
||||
res = h.redeem_code(args.code, method=args.method)
|
||||
res = h.redeem_code(
|
||||
args.code,
|
||||
method=args.method,
|
||||
notify_target=getattr(args, "notify_target", None),
|
||||
notify_agent=getattr(args, "notify_agent", None),
|
||||
)
|
||||
if args.json:
|
||||
print(json.dumps(res.to_dict(), indent=2))
|
||||
else:
|
||||
@@ -556,6 +739,8 @@ def main() -> None:
|
||||
print(f" Status: {res.status}")
|
||||
if res.detail:
|
||||
print(f" Detail: {res.detail}")
|
||||
if res.loopback_notified:
|
||||
print(f" Loopback:{res.loopback_detail}")
|
||||
|
||||
elif args.command == "list":
|
||||
fleet = scan_fleet_invites()
|
||||
@@ -576,7 +761,12 @@ def main() -> None:
|
||||
print()
|
||||
|
||||
elif args.command == "salvage":
|
||||
salvage_res = salvage_blocked_node(args.node, helper_node=args.helper)
|
||||
salvage_res = salvage_blocked_node(
|
||||
args.node,
|
||||
helper_node=args.helper,
|
||||
notify_target=getattr(args, "notify_target", None),
|
||||
notify_agent=getattr(args, "notify_agent", None),
|
||||
)
|
||||
if args.json:
|
||||
print(json.dumps(salvage_res, indent=2))
|
||||
else:
|
||||
@@ -588,7 +778,10 @@ def main() -> None:
|
||||
print(f" Status: {salvage_res.get('error')}")
|
||||
if salvage_res.get("share_instruction"):
|
||||
print(f" Next Step: {salvage_res.get('share_instruction')}")
|
||||
if salvage_res.get("loopback_notified"):
|
||||
print(" Loopback: Notice sent to requesting target.")
|
||||
print()
|
||||
|
||||
else:
|
||||
parser.print_help()
|
||||
|
||||
|
||||
Executable
+560
@@ -0,0 +1,560 @@
|
||||
#!/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
|
||||
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 <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)
|
||||
|
||||
|
||||
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 <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)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
+93
-2
@@ -355,7 +355,19 @@ def finish_onboarding_redemption(state: OnboardState) -> Dict[str, Any]:
|
||||
state.stage = STAGE_REDEEMED
|
||||
state.detail = f"Successfully redeemed code {state.invite_code}! 1B tokens credited to @{state.beneficiary_node} and @{state.node}."
|
||||
else:
|
||||
state.detail = f"Redemption failed: {redemption.get('reason')} - {redemption.get('detail')}"
|
||||
reason = redemption.get("reason", "unknown")
|
||||
detail = redemption.get("detail", "")
|
||||
state.detail = f"Redemption failed: {reason} - {detail}"
|
||||
try:
|
||||
from invite_handler import send_loopback_notice
|
||||
target_chat = f"{state.beneficiary_node} tasks" if state.beneficiary_node else "heartbeat-opm"
|
||||
send_loopback_notice(
|
||||
recipient=state.beneficiary_node or "opm",
|
||||
target=target_chat,
|
||||
message=f"[ONBOARD-SALVAGE-LOOPBACK] Node @{state.node} could not redeem code {state.invite_code} for @{state.beneficiary_node}: {reason}. Stage remains safe.",
|
||||
)
|
||||
except Exception:
|
||||
pass
|
||||
except Exception as e:
|
||||
redemption_info = {"ok": False, "error": str(e)}
|
||||
state.detail = f"Redemption exception: {e}"
|
||||
@@ -404,10 +416,71 @@ def issue_salvage_work_order(blocked_node: str = "646", to_sidechat: str = "646
|
||||
return {"ok": res.returncode == 0, "output": res.stdout.strip()}
|
||||
|
||||
|
||||
def get_all_connects(fast: bool = True) -> List[Dict[str, Any]]:
|
||||
"""Return consolidated inventory of all fleet and onboarded connects."""
|
||||
connects = []
|
||||
seen = set()
|
||||
|
||||
# Load recorded pipeline states
|
||||
if STATE_DIR.exists():
|
||||
for p in STATE_DIR.glob("*.json"):
|
||||
try:
|
||||
d = json.loads(p.read_text(encoding="utf-8"))
|
||||
n = d.get("node")
|
||||
if n:
|
||||
seen.add(n)
|
||||
connects.append({
|
||||
"node": n,
|
||||
"type": "onboard_pipeline",
|
||||
"email": d.get("email"),
|
||||
"stage": d.get("stage"),
|
||||
"beneficiary": d.get("beneficiary_node"),
|
||||
"invite_code": d.get("invite_code") or "-",
|
||||
"cdp_port": d.get("cdp_port"),
|
||||
"detail": d.get("detail"),
|
||||
"updated_at": d.get("updated_at"),
|
||||
})
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Fleet nodes
|
||||
w_map = {}
|
||||
if not fast:
|
||||
try:
|
||||
weights = calculate_feeding_weights()
|
||||
w_map = {r["node"]: r for r in weights}
|
||||
except Exception:
|
||||
w_map = {}
|
||||
|
||||
ports = {"muse": 9222, "pip": 9322, "646": 9430, "opm": 9440, "def": 9450, "dev": 9455}
|
||||
|
||||
for agent in ("muse", "pip", "646", "opm", "dev", "def"):
|
||||
if agent not in seen:
|
||||
w_info = w_map.get(agent, {})
|
||||
connects.append({
|
||||
"node": agent,
|
||||
"type": "fleet_agent",
|
||||
"email": f"{agent}@muse-dev.online",
|
||||
"stage": "active_fleet",
|
||||
"beneficiary": None,
|
||||
"invite_code": w_info.get("code") or "-",
|
||||
"cdp_port": ports.get(agent),
|
||||
"role": AGENT_ROLES.get(agent, ""),
|
||||
"status": w_info.get("status", "ACTIVE"),
|
||||
"feeding_weight": w_info.get("feeding_weight", 1.0),
|
||||
"updated_at": time.time(),
|
||||
})
|
||||
|
||||
return connects
|
||||
|
||||
|
||||
def main() -> None:
|
||||
parser = argparse.ArgumentParser(description="End-to-end agent-driven onboarding & invite salvage pipeline")
|
||||
sub = parser.add_subparsers(dest="action")
|
||||
|
||||
p_conn = sub.add_parser("connects", help="Inventory of all active fleet nodes & client onboard connects")
|
||||
p_conn.add_argument("--json", action="store_true", help="Emit JSON output")
|
||||
|
||||
p_start = sub.add_parser("start", help="Start full onboarding pipeline for a client node")
|
||||
p_start.add_argument("node", help="New node label (e.g. dev2, client1)")
|
||||
p_start.add_argument("--email", required=True, help="Client login email")
|
||||
@@ -435,7 +508,25 @@ def main() -> None:
|
||||
|
||||
args = parser.parse_args()
|
||||
|
||||
if args.action == "start":
|
||||
if args.action == "connects":
|
||||
conn_list = get_all_connects()
|
||||
if args.json:
|
||||
print(json.dumps(conn_list, indent=2))
|
||||
else:
|
||||
print("\n=== ACTIVE FLEET & CLIENT ONBOARD CONNECTS ===")
|
||||
print(f" {'NODE':6} {'TYPE':16} {'STAGE / STATUS':16} {'CDP':6} {'INVITE':8} {'ROLE / DETAIL'}")
|
||||
print(f" {'─'*6} {'─'*16} {'─'*16} {'─'*6} {'─'*8} {'─'*32}")
|
||||
for c in conn_list:
|
||||
node = c.get("node", "")
|
||||
t_str = c.get("type", "")
|
||||
st_str = c.get("stage", c.get("status", ""))
|
||||
cdp = str(c.get("cdp_port") or "-")
|
||||
code = c.get("invite_code") or "-"
|
||||
role = c.get("role") or c.get("detail") or c.get("email") or ""
|
||||
print(f" {node:<6} {t_str:<16} {st_str:<16} {cdp:<6} {code:<8} {role}")
|
||||
print()
|
||||
|
||||
elif args.action == "start":
|
||||
res = start_onboarding(args.node, args.email, beneficiary_node=getattr(args, "for_agent", None), invite_code=args.code, account_name=args.account_name)
|
||||
if args.json:
|
||||
print(json.dumps(res, indent=2))
|
||||
|
||||
+19
-12
@@ -68,6 +68,7 @@ def _tool(op, args):
|
||||
def spawn_call(job_id, job_name, profile):
|
||||
count, _, hint = PROFILES[profile]
|
||||
task = "Subagent for job %s (%s): %s." % (job_id, job_name, hint)
|
||||
return _tool("swarm.spawn", {"count": count, "task": task[:900], "label": (job_name or "job")[:60]})
|
||||
|
||||
|
||||
def dm_call(to, target, message):
|
||||
@@ -76,9 +77,6 @@ def dm_call(to, target, message):
|
||||
separators=(", ", ": "))
|
||||
|
||||
|
||||
return _tool("swarm.spawn", {"count": count, "task": task[:900], "label": (job_name or "job")[:60]})
|
||||
|
||||
|
||||
def native_followup_call(job_id, job_name, profile, agent):
|
||||
"""Muse-native one-shot cron (cron.create runonce) - bridged to followup.create."""
|
||||
_, mins, _ = PROFILES[profile]
|
||||
@@ -99,7 +97,7 @@ def thread_url(target):
|
||||
return "https://box.muse-dev.online/%s/%s" % ("thread" if UUID_RE.fullmatch(t.lower()) else "sidechat", t)
|
||||
|
||||
|
||||
def wrap(job_name, job_id, agent, target, rendered):
|
||||
def wrap(job_name, job_id, agent, target, rendered, include_kpi: bool = True):
|
||||
"""Return the full DM body as a direct, authentic operator directive."""
|
||||
profile = pick_profile(job_name)
|
||||
spawn = spawn_call(job_id, job_name, profile)
|
||||
@@ -107,24 +105,33 @@ def wrap(job_name, job_id, agent, target, rendered):
|
||||
parts = job_id.split("-")
|
||||
wo_id = parts[-1] if len(parts) > 1 and len(parts[-1]) >= 8 else hashlib.sha256(job_id.encode()).hexdigest()[:8]
|
||||
session_name = f"work-{agent}-{wo_id}"
|
||||
tmux_new = _tool("tmux.new", {"session": session_name, "command": "bash"})
|
||||
tmux_send = _tool("tmux.send", {"session": session_name, "keys": "echo 'Starting task execution...'"})
|
||||
|
||||
advisory_section = ""
|
||||
if include_kpi:
|
||||
try:
|
||||
import kpi
|
||||
advisory = kpi.get_live_advisory_block(agent)
|
||||
if advisory:
|
||||
advisory_section = f"\n{advisory}\n"
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
top = (
|
||||
f"Operator Directive [ref:{wo_id}]:\n"
|
||||
"Execute the task below using tool calls. Background tmux session is ready for command execution:\n"
|
||||
f" • {tmux_new}\n"
|
||||
f" • {tmux_send}\n"
|
||||
f"Host tmux worker session '{session_name}' is available on bl (/tmp/tmux-muse.sock).\n"
|
||||
f" • Subagent assistance: {spawn}\n"
|
||||
f" • Verification schedule: {follow}\n\n"
|
||||
f" • Verification schedule: {follow}\n"
|
||||
f"{advisory_section}\n"
|
||||
"--- Task ---\n"
|
||||
)
|
||||
|
||||
has_result = "[RESULT" in rendered
|
||||
bottom = (
|
||||
"\n--- End Task ---\n\n"
|
||||
f"Inspect tmux output: [TOOL tmux.capture {{\"session\": \"{session_name}\", \"lines\": 30}}]\n"
|
||||
f"Inspect tmux worker: box tmux capture {session_name} 30 (or attach via /tmp/tmux-muse.sock)\n"
|
||||
"Tools: cron.create, cron.runs, health.check, swarm.spawn, swarm.list, dm.send, box.exec, tools.list.\n"
|
||||
"Message a peer: [DM {\"to\": \"<agent>\", \"target\": \"<sidechat>\", \"message\": \"<text>\"}].\n"
|
||||
"Query box: [TOOL box.exec {\"action\": \"<fleet-status|dm-log|job-get|...>\"}] \u2014 [TOOL tools.list {}] lists every op.\n"
|
||||
"Query box: [TOOL box.exec {\"action\": \"<fleet-status|dm-log|job-get|...>\"}] — [TOOL tools.list {}] lists every op.\n"
|
||||
)
|
||||
if not has_result:
|
||||
bottom += f"When complete, report your verdict: [RESULT {job_id}] OK: <summary of actions>\n"
|
||||
|
||||
+25
-8
@@ -53,11 +53,13 @@ class NodeUsage:
|
||||
bars: List[Dict[str, Any]]
|
||||
dialog_text: str
|
||||
has_redeemed: bool = False
|
||||
stats_loaded: bool = True
|
||||
|
||||
def to_dict(self) -> Dict[str, Any]:
|
||||
return asdict(self)
|
||||
|
||||
|
||||
|
||||
def cdp_send_command(ws: Any, method: str, params: Dict[str, Any], timeout: float = 3.0) -> Optional[Dict[str, Any]]:
|
||||
"""Send a raw CDP command and await response with a unique request ID."""
|
||||
req_id = next(_REQ_COUNTER)
|
||||
@@ -345,9 +347,18 @@ class SettingsRPA:
|
||||
self.close_settings_dialog()
|
||||
|
||||
if not raw:
|
||||
raise RuntimeError(f"Failed to read usage data from Settings dialog on {self.node}")
|
||||
raw = {
|
||||
"dialogText": "",
|
||||
"bars": [],
|
||||
"plan": "Unknown",
|
||||
"resetText": "Unavailable (usage stats did not load)",
|
||||
"tokensLeft": "Unavailable",
|
||||
}
|
||||
|
||||
bars = raw.get("bars", [])
|
||||
dialog_text = raw.get("dialogText", "")
|
||||
stats_loaded = bool(bars or "Weekly limit" in dialog_text or "Additional tokens" in dialog_text or "Free plan" in dialog_text)
|
||||
|
||||
weekly_pct = 0
|
||||
extra_pct = 0
|
||||
extra_status = "Never expires"
|
||||
@@ -362,25 +373,26 @@ class SettingsRPA:
|
||||
|
||||
# Blocked condition: weekly limit 100% and additional tokens 100% or 0 tokens left
|
||||
tokens_left = raw.get("tokensLeft", "")
|
||||
is_blocked = (weekly_pct >= 100 and extra_pct >= 100) or ("0 tokens left" in tokens_left)
|
||||
is_blocked = bool(stats_loaded and ((weekly_pct >= 100 and extra_pct >= 100) or ("0 tokens left" in tokens_left)))
|
||||
|
||||
# has_redeemed binary: If the "Additional tokens" ticker is present in Settings (or "Redeem invite code" entry is absent),
|
||||
# the agent has already redeemed an invite code.
|
||||
has_extra_ticker = any("additional tokens" in b.get("label", "").lower() or "additional" in b.get("raw_text", "").lower() for b in bars)
|
||||
has_redeemed = has_extra_ticker or ("Additional tokens" in raw.get("dialogText", ""))
|
||||
has_redeemed = bool(has_extra_ticker or ("Additional tokens" in dialog_text))
|
||||
|
||||
return NodeUsage(
|
||||
node=self.node,
|
||||
plan=raw.get("plan", "Free plan"),
|
||||
weekly_reset_text=raw.get("resetText", ""),
|
||||
plan=raw.get("plan", "Unknown" if not stats_loaded else "Free plan"),
|
||||
weekly_reset_text=raw.get("resetText", "") or ("Unavailable (stats did not load)" if not stats_loaded else ""),
|
||||
weekly_percent_used=weekly_pct,
|
||||
extra_tokens_status=extra_status,
|
||||
extra_percent_used=extra_pct,
|
||||
extra_tokens_remaining=tokens_left or ("0 tokens left" if extra_pct >= 100 else "Unknown"),
|
||||
extra_tokens_remaining=tokens_left or ("Unavailable" if not stats_loaded else ("0 tokens left" if extra_pct >= 100 else "Unknown")),
|
||||
is_blocked=is_blocked,
|
||||
bars=bars,
|
||||
dialog_text=raw.get("dialogText", ""),
|
||||
dialog_text=dialog_text,
|
||||
has_redeemed=has_redeemed,
|
||||
stats_loaded=stats_loaded,
|
||||
)
|
||||
|
||||
def check_redeem_entrypoint(self) -> Dict[str, Any]:
|
||||
@@ -423,7 +435,12 @@ def main() -> None:
|
||||
if args.json:
|
||||
print(json.dumps(usage.to_dict(), indent=2))
|
||||
else:
|
||||
status_str = "BLOCKED (LIMIT REACHED)" if usage.is_blocked else "ACTIVE"
|
||||
if not usage.stats_loaded:
|
||||
status_str = "UNLOADED (STATS DID NOT RENDER)"
|
||||
elif usage.is_blocked:
|
||||
status_str = "BLOCKED (LIMIT REACHED)"
|
||||
else:
|
||||
status_str = "ACTIVE"
|
||||
print(f"=== Node {usage.node} Usage ===")
|
||||
print(f" Plan: {usage.plan}")
|
||||
print(f" Weekly Reset: {usage.weekly_reset_text} ({usage.weekly_percent_used}% used)")
|
||||
|
||||
+370
-10
@@ -930,6 +930,169 @@ def cmd_approvals(args):
|
||||
print(f" Reason: {c_cyan(reason)}")
|
||||
print(c_dim(f" Operator can approve with: box approvals allow {node}"))
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Domain: RUNTIME (agentic management of Muse CLI tmux runtimes)
|
||||
# ---------------------------------------------------------------------------
|
||||
def cmd_runtime(args):
|
||||
import shlex
|
||||
import muse_choice_watcher as mcw
|
||||
action = getattr(args, "rt_action", None) or "list"
|
||||
as_json = getattr(args, "json", False)
|
||||
|
||||
if action == "list":
|
||||
sock = getattr(args, "socket", None)
|
||||
rows = mcw.all_runtime_rows([sock] if sock else None)
|
||||
if getattr(args, "muse_only", False):
|
||||
rows = [r for r in rows if r["is_muse"]]
|
||||
if as_json:
|
||||
print(json.dumps({"ok": True, "runtimes": rows}, indent=2))
|
||||
return
|
||||
print(c_bold("\n=== MUSE RUNTIMES ===\n"))
|
||||
if not rows:
|
||||
print(c_dim(" No panes found."))
|
||||
print()
|
||||
return
|
||||
headers = ["SOCKET", "SESSION", "PANE", "CMD", "STATE",
|
||||
"APPROVE", "WATCHER"]
|
||||
table = []
|
||||
for r in rows:
|
||||
state = r["state"]
|
||||
if r["state"] == "approval-pending" and r["prompt_kind"]:
|
||||
state = "%s(%s/%s)" % (state, r["prompt_kind"],
|
||||
r["prompt_key"])
|
||||
if r["is_muse"]:
|
||||
approve = (badge_ok("YES") if r["auto_approve"]
|
||||
else badge_err("NO"))
|
||||
else:
|
||||
approve = badge_dim("-")
|
||||
watcher = (badge_ok("ALIVE %s" % r["watcher_pid"])
|
||||
if r["watcher_alive"] else badge_dim("-"))
|
||||
table.append([os.path.basename(r["socket"]),
|
||||
"%s:%s" % (r["session"], r["window"]),
|
||||
r["pane"], (r["cmd"] or "")[:26], state,
|
||||
approve, watcher])
|
||||
print_table(headers, table)
|
||||
print()
|
||||
|
||||
elif action == "send":
|
||||
sock = getattr(args, "socket", None) or mcw.KNOWN_SOCKETS[0]
|
||||
pane = args.pane
|
||||
keys = args.keys
|
||||
enter = not getattr(args, "no_enter", False)
|
||||
pre = mcw.pane_state(sock, pane)
|
||||
if pre.get("error"):
|
||||
if as_json:
|
||||
print(json.dumps({"ok": False, "error": pre["error"],
|
||||
"socket": sock, "pane": pane}))
|
||||
return
|
||||
print(c_red("Error: no such pane %s on %s" % (pane, sock)),
|
||||
file=sys.stderr)
|
||||
sys.exit(1)
|
||||
t_args = ["send-keys", "-t", pane, keys]
|
||||
if enter:
|
||||
t_args.append("Enter")
|
||||
r = mcw._tmux(sock, *t_args, timeout=10)
|
||||
ok = r.returncode == 0
|
||||
if as_json:
|
||||
print(json.dumps({
|
||||
"ok": ok, "socket": sock, "pane": pane,
|
||||
"pre_state": pre["state"],
|
||||
"prompt_kind": pre["prompt_kind"],
|
||||
"sent": keys, "enter": enter,
|
||||
"error": (r.stderr or r.stdout or "").strip() or None
|
||||
if not ok else None}, indent=2))
|
||||
return
|
||||
print("\n Pane %s on %s [%s]" % (
|
||||
c_bold(pane), sock, c_cyan(pre["state"])))
|
||||
if ok:
|
||||
print(" %s sent %r%s" % (c_green("✔"),
|
||||
keys, " + Enter" if enter else ""))
|
||||
else:
|
||||
print(" %s send failed: %s" % (
|
||||
c_red("✘"),
|
||||
(r.stderr or r.stdout or "").strip() or "tmux error"))
|
||||
sys.exit(1)
|
||||
print()
|
||||
|
||||
elif action == "launch":
|
||||
sock = getattr(args, "socket", None) or mcw.KNOWN_SOCKETS[0]
|
||||
session = args.session
|
||||
window = getattr(args, "window", None)
|
||||
dry_run = getattr(args, "dry_run", False)
|
||||
muse_args = list(getattr(args, "muse_args", None) or [])
|
||||
if muse_args[:1] == ["--"]:
|
||||
muse_args = muse_args[1:]
|
||||
posture = mcw.muse_approval_flags(muse_args)
|
||||
injected = [] if posture["flags"] else ["--disable-approval"]
|
||||
launcher = (shutil.which("muse-code")
|
||||
or "/home/super/.local/bin/muse-code")
|
||||
cmdline = shlex.join([launcher] + injected + muse_args)
|
||||
if dry_run:
|
||||
if as_json:
|
||||
print(json.dumps({
|
||||
"ok": True, "dry_run": True, "socket": sock,
|
||||
"session": session, "window": window,
|
||||
"cmdline": cmdline, "injected": injected}, indent=2))
|
||||
return
|
||||
print(c_bold("\n=== RUNTIME LAUNCH (dry-run) ===\n"))
|
||||
print(" Socket: %s" % sock)
|
||||
print(" Session: %s" % c_cyan(session))
|
||||
print(" Command: %s" % cmdline)
|
||||
if injected:
|
||||
print(" %s auto-approve injected: %s" % (
|
||||
c_green("✔"), " ".join(injected)))
|
||||
else:
|
||||
print(" %s caller already sets approval flags; "
|
||||
"nothing injected" % c_dim("•"))
|
||||
print()
|
||||
return
|
||||
exists = mcw._tmux(sock, "has-session", "-t", session, timeout=10)
|
||||
if exists.returncode == 0:
|
||||
if as_json:
|
||||
print(json.dumps({"ok": False, "error": "session_exists",
|
||||
"socket": sock, "session": session}))
|
||||
return
|
||||
print(c_red("Error: session '%s' already exists on %s"
|
||||
% (session, sock)), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
t_args = ["new-session", "-d", "-s", session]
|
||||
if window:
|
||||
t_args.extend(["-n", window])
|
||||
t_args.append(cmdline)
|
||||
r = mcw._tmux(sock, *t_args, timeout=15)
|
||||
if r.returncode != 0:
|
||||
err = (r.stderr or r.stdout or "").strip() or "tmux error"
|
||||
if as_json:
|
||||
print(json.dumps({"ok": False, "error": err,
|
||||
"socket": sock, "session": session}))
|
||||
return
|
||||
print(c_red("Error: launch failed: %s" % err),
|
||||
file=sys.stderr)
|
||||
sys.exit(1)
|
||||
rec = mcw.reconcile(sockets=[sock])
|
||||
if as_json:
|
||||
print(json.dumps({"ok": True, "socket": sock,
|
||||
"session": session, "window": window,
|
||||
"cmdline": cmdline, "injected": injected,
|
||||
"reconcile": rec}, indent=2))
|
||||
return
|
||||
print(c_green("\n✔ Launched '%s' on %s") % (session, sock))
|
||||
print(" Command: %s" % c_dim(cmdline))
|
||||
for s in rec.get("started") or []:
|
||||
print(" %s watcher %s" % (badge_ok("STARTED"), c_cyan(s)))
|
||||
print()
|
||||
|
||||
else:
|
||||
if as_json:
|
||||
print(json.dumps({"ok": False,
|
||||
"error": "unknown_action",
|
||||
"action": action}))
|
||||
return
|
||||
print(c_red("Error: unknown runtime action '%s'" % action),
|
||||
file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Domain: MUSE-CHOICES (Muse TUI A/B/C auto-answer daemon)
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -1151,28 +1314,51 @@ def cmd_invite(args):
|
||||
sys.exit(1)
|
||||
elif action == "redeem":
|
||||
code = getattr(args, "code", None)
|
||||
method = getattr(args, "method", "auto")
|
||||
notify_target = getattr(args, "notify", None) or getattr(args, "notify_target", None)
|
||||
notify_agent = getattr(args, "notify_agent", None)
|
||||
try:
|
||||
res = invite.redeem_invite(node, code)
|
||||
if method != "api" or notify_target or notify_agent:
|
||||
from invite_handler import InviteHandler
|
||||
with InviteHandler(node) as h:
|
||||
res_obj = h.redeem_code(code, method=method, notify_target=notify_target, notify_agent=notify_agent)
|
||||
res = res_obj.to_dict()
|
||||
res["ok"] = res_obj.success
|
||||
else:
|
||||
res = invite.redeem_invite(node, code)
|
||||
except invite.InviteError as e:
|
||||
print(c_red(f" Error: {e}"), file=sys.stderr)
|
||||
sys.exit(2)
|
||||
except Exception as e:
|
||||
print(c_red(f" Error during redemption: {e}"), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
if as_json:
|
||||
print(json.dumps(res, indent=2))
|
||||
elif res.get("ok"):
|
||||
msg = f" ✔ Redeemed {res['code']} on {node}."
|
||||
msg = f" ✔ Redeemed {res.get('code') or code} on {node}."
|
||||
if res.get("detail"):
|
||||
msg += f" {res['detail']}"
|
||||
print(c_green(msg))
|
||||
else:
|
||||
print(c_red(f" ✘ Redeem failed on {node}: "
|
||||
f"{res.get('reason') or res.get('error')}"),
|
||||
file=sys.stderr)
|
||||
reason_str = res.get("reason") or res.get("error") or "failed"
|
||||
print(c_red(f" ✘ Redeem failed on {node}: {reason_str}"), file=sys.stderr)
|
||||
if res.get("detail"):
|
||||
print(c_dim(f" {res['detail']}"), file=sys.stderr)
|
||||
print(c_dim(f" Detail: {res['detail']}"), file=sys.stderr)
|
||||
if res.get("field_missing"):
|
||||
print(c_yellow(f" Note: Expected redemption input field was missing or hidden on @{node}."), file=sys.stderr)
|
||||
if res.get("loopback_notified"):
|
||||
print(c_cyan(f" Loopback: Dispatched alert to {notify_agent or 'coordinator'}/{notify_target}."), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
elif action == "salvage":
|
||||
target = getattr(args, "node", "646") or "646"
|
||||
res = salvage_blocked_node(target, helper_node=getattr(args, "helper", None))
|
||||
notify_target = getattr(args, "notify", None) or getattr(args, "notify_target", None)
|
||||
notify_agent = getattr(args, "notify_agent", None)
|
||||
res = salvage_blocked_node(
|
||||
target,
|
||||
helper_node=getattr(args, "helper", None),
|
||||
notify_target=notify_target,
|
||||
notify_agent=notify_agent,
|
||||
)
|
||||
if as_json:
|
||||
print(json.dumps(res, indent=2))
|
||||
return
|
||||
@@ -1187,6 +1373,8 @@ def cmd_invite(args):
|
||||
action_hint = res.get("share_instruction") or f"Redeem code '{code}' via onboarding pipeline to grant 1B tokens to {target}."
|
||||
print(f" Action: {action_hint}")
|
||||
print(f" Command: {c_cyan('box onboard start <new_node> --email <client_email> --for ' + target)}")
|
||||
if res.get("loopback_notified"):
|
||||
print(c_cyan(f" Loopback: Notice dispatched to {notify_agent or 'coordinator'}/{notify_target}."))
|
||||
print()
|
||||
|
||||
|
||||
@@ -1213,6 +1401,10 @@ def cmd_usage(args):
|
||||
rows.append([c_bold(n), c_red("ERROR"), "-", "-",
|
||||
(it.get("error") or "")[:30], "-"])
|
||||
continue
|
||||
if not it.get("stats_loaded", True):
|
||||
rows.append([c_bold(n), c_yellow("UNLOADED"), "-", "-",
|
||||
"Stats did not render", "-"])
|
||||
continue
|
||||
wu_val = it.get("weekly_used_pct")
|
||||
wu = ("%d%%" % wu_val) if wu_val is not None else "-"
|
||||
au_val = it.get("additional_used_pct")
|
||||
@@ -1233,6 +1425,7 @@ def cmd_usage(args):
|
||||
if has_blocked:
|
||||
print("\n" + c_yellow(" ⚠ One or more agents have reached usage limits. Run 'box invite salvage <node>' to resolve.") + "\n")
|
||||
else:
|
||||
|
||||
print()
|
||||
|
||||
|
||||
@@ -1264,9 +1457,80 @@ def cmd_settings(args):
|
||||
print(json.dumps(info, indent=2))
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Domain: KPI (Spend Monitoring, Performance & Runtime Preservation)
|
||||
# ---------------------------------------------------------------------------
|
||||
def cmd_kpi(args):
|
||||
import kpi
|
||||
act = getattr(args, "kpi_action", "status") or "status"
|
||||
as_json = getattr(args, "json", False)
|
||||
|
||||
if act == "status":
|
||||
nodes = [args.node] if getattr(args, "node", None) else kpi.VALID_NODES
|
||||
kpis = kpi.fleet_kpi(nodes)
|
||||
if as_json:
|
||||
print(json.dumps({k: v.to_dict() for k, v in kpis.items()}, indent=2))
|
||||
return
|
||||
print("\n" + c_bold("=== 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(c_bold(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"{c_bold(k.node):<15} {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("\n" + c_dim("Run 'box kpi report <node>' for prescriptive runtime preservation advisories.\n"))
|
||||
elif act == "report":
|
||||
node = getattr(args, "node", "646") or "646"
|
||||
res = kpi.get_agent_kpi(node)
|
||||
if as_json:
|
||||
print(json.dumps(res.to_dict(), indent=2))
|
||||
return
|
||||
print(f"\n{c_bold('=== KPI & RUNTIME REPORT: @' + res.node.upper() + ' ===')}")
|
||||
print(f" Weekly Quota: {res.weekly_used_pct}% used")
|
||||
print(f" Extra Tokens: {res.extra_tokens_remaining}")
|
||||
print(f" Blocked Status: {'YES (LIMIT REACHED)' if res.is_blocked else 'NO (HEALTHY)'}")
|
||||
print(f" Messages / Calls: {res.calls_sent} sent ({res.calls_verified} verified delivered)")
|
||||
print(f" Jobs Dispatched: {res.jobs_completed} completed / {res.jobs_assigned} assigned")
|
||||
print(f" Active Subagents: {res.subagents_active}")
|
||||
print(f" Active Tmux Workers:{res.tmux_workers_active}")
|
||||
print(f" Process Uptime: {res.uptime_hours} hours")
|
||||
print(f" Route Health: {res.route_status}")
|
||||
print(f" Efficiency Index: {res.efficiency_index} ({res.efficiency_rating})")
|
||||
print(f"\n {c_yellow('[RUNTIME PRESERVATION ADVISORY]')}\n {res.preservation_advisory}\n")
|
||||
elif act == "routes":
|
||||
routes = {n: kpi.check_node_routes(n) for n in kpi.VALID_NODES}
|
||||
if as_json:
|
||||
print(json.dumps(routes, indent=2))
|
||||
else:
|
||||
print("\n" + c_bold("=== NETVM ROUTE HEALTH ==="))
|
||||
for n, st in routes.items():
|
||||
print(f" @{c_bold(n):<15} : {st}")
|
||||
print()
|
||||
elif act == "spawn-worker":
|
||||
res = kpi.spawn_tmux_worker(args.node, args.session, args.worker_command)
|
||||
if as_json:
|
||||
print(json.dumps(res, indent=2))
|
||||
else:
|
||||
if res.get("ok"):
|
||||
print(c_green(f"✔ {res.get('message')}"))
|
||||
else:
|
||||
print(c_red(f"✘ Failed to spawn worker: {res.get('error')}"), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Domain: ONBOARD (Agent-driven Onboarding & Invite Salvage Pipeline)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
def cmd_onboard(args):
|
||||
import onboard_pipeline
|
||||
action = getattr(args, "onboard_action", "start") or "start"
|
||||
@@ -1334,6 +1598,24 @@ def cmd_onboard(args):
|
||||
print(f" Redemption: {r_st}")
|
||||
print()
|
||||
|
||||
elif action == "connects":
|
||||
conn_list = onboard_pipeline.get_all_connects()
|
||||
if as_json:
|
||||
print(json.dumps(conn_list, indent=2))
|
||||
return
|
||||
print(f"\n{c_bold('=== ACTIVE FLEET & CLIENT ONBOARD CONNECTS ===')}")
|
||||
print(f" {'NODE':<8} {'TYPE':<16} {'STAGE / STATUS':<18} {'CDP':<8} {'INVITE':<10} {'ROLE / DETAIL'}")
|
||||
print(" " + "─" * 78)
|
||||
for c in conn_list:
|
||||
node = c.get("node", "")
|
||||
t_str = c.get("type", "")
|
||||
st_str = c.get("stage", c.get("status", ""))
|
||||
cdp = str(c.get("cdp_port") or "-")
|
||||
code = c.get("invite_code") or "-"
|
||||
role = c.get("role") or c.get("detail") or c.get("email") or ""
|
||||
print(f" {node:<8} {t_str:<16} {st_str:<18} {cdp:<8} {code:<10} {role}")
|
||||
print()
|
||||
|
||||
elif action == "salvage-wo":
|
||||
node = getattr(args, "node", "646") or "646"
|
||||
target_chat = getattr(args, "target", "646 tasks") or "646 tasks"
|
||||
@@ -4874,9 +5156,18 @@ def cmd_docs_dispatch(args):
|
||||
|
||||
|
||||
def cmd_tmux_dispatch(args):
|
||||
"""Bridge 'box tmux' commands directly to bin/muse-tmux.py."""
|
||||
tmux_bin = str(BIN_DIR / "muse-tmux.py")
|
||||
"""Bridge 'box tmux' commands directly to bin/muse-tmux.py or tmux_auto_approver.py."""
|
||||
t_args = getattr(args, "tmux_args", []) or []
|
||||
if t_args and t_args[0] in ("tally", "auto", "watch", "once", "match", "rules", "status"):
|
||||
approver_bin = str(BIN_DIR / "tmux_auto_approver.py")
|
||||
if t_args[0] == "auto":
|
||||
sub = t_args[1:] or ["status"]
|
||||
else:
|
||||
sub = t_args
|
||||
cmd = [sys.executable, approver_bin] + sub
|
||||
res = subprocess.run(cmd)
|
||||
sys.exit(res.returncode)
|
||||
tmux_bin = str(BIN_DIR / "muse-tmux.py")
|
||||
if not t_args:
|
||||
t_args = ["list"]
|
||||
cmd = [sys.executable, tmux_bin] + t_args
|
||||
@@ -5050,6 +5341,27 @@ def build_parser():
|
||||
|
||||
p_mc_rec = mc_sub.add_parser("reconcile", parents=[common], help="Enforce desired state now (start missing / stop excess)")
|
||||
|
||||
# Domain: RUNTIME
|
||||
p_rt = subparsers.add_parser("runtime", parents=[common], help="Muse CLI tmux runtimes: list states, send input, launch auto-approved")
|
||||
rt_sub = p_rt.add_subparsers(dest="rt_action")
|
||||
|
||||
p_rt_list = rt_sub.add_parser("list", parents=[common], help="List panes with runtime state + approval posture (default)")
|
||||
p_rt_list.add_argument("--socket", default=None, help="Only this tmux socket")
|
||||
p_rt_list.add_argument("--muse-only", action="store_true", help="Only Muse CLI panes")
|
||||
|
||||
p_rt_send = rt_sub.add_parser("send", parents=[common], help="Send keys to a pane (reports pre-send state)")
|
||||
p_rt_send.add_argument("--socket", default=None, help="Tmux socket (default: /tmp/tmux-1000/default)")
|
||||
p_rt_send.add_argument("pane", help="Pane id (e.g. %%37)")
|
||||
p_rt_send.add_argument("keys", help="Keys / text to send")
|
||||
p_rt_send.add_argument("--no-enter", action="store_true", help="Do not send Enter after keys")
|
||||
|
||||
p_rt_launch = rt_sub.add_parser("launch", parents=[common], help="Launch a Muse session with auto-approve injected")
|
||||
p_rt_launch.add_argument("--socket", default=None, help="Tmux socket (default: /tmp/tmux-1000/default)")
|
||||
p_rt_launch.add_argument("--session", required=True, help="New tmux session name")
|
||||
p_rt_launch.add_argument("--window", "-w", default=None, help="Initial window name")
|
||||
p_rt_launch.add_argument("--dry-run", action="store_true", help="Print the launch plan without creating")
|
||||
p_rt_launch.add_argument("muse_args", nargs=argparse.REMAINDER, default=[], help="Extra muse args after --")
|
||||
|
||||
# Domain: INVITE
|
||||
p_invite = subparsers.add_parser("invite", parents=[common], help="Muse.ai invite codes: find per-agent codes and redeem")
|
||||
p_invite.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node (status)")
|
||||
@@ -5070,10 +5382,15 @@ def build_parser():
|
||||
p_inv_redeem = inv_sub.add_parser("redeem", parents=[common], help="Redeem an invite code on a node")
|
||||
p_inv_redeem.add_argument("node", choices=VALID_NODES, help="Target node")
|
||||
p_inv_redeem.add_argument("code", help="6-char invite code")
|
||||
p_inv_redeem.add_argument("--method", choices=["auto", "dom", "api"], default="auto", help="Redemption method (default: auto)")
|
||||
p_inv_redeem.add_argument("--notify", default=None, help="Sidechat to notify on loopback")
|
||||
p_inv_redeem.add_argument("--notify-agent", default=None, help="Agent to notify on loopback")
|
||||
|
||||
p_inv_salvage = inv_sub.add_parser("salvage", parents=[common], help="Salvage an out-of-tokens agent (defaults to 646)")
|
||||
p_inv_salvage.add_argument("node", nargs="?", default="646", choices=VALID_NODES, help="Blocked agent node (default: 646)")
|
||||
p_inv_salvage.add_argument("--helper", choices=VALID_NODES, help="Specific helper agent to redeem code")
|
||||
p_inv_salvage.add_argument("--notify", default=None, help="Sidechat to notify on loopback")
|
||||
p_inv_salvage.add_argument("--notify-agent", default=None, help="Agent to notify on loopback")
|
||||
|
||||
# Domain: USAGE
|
||||
p_usage = subparsers.add_parser("usage", parents=[common], help="Muse.ai usage limits per agent")
|
||||
@@ -5112,8 +5429,28 @@ def build_parser():
|
||||
p_onb_wo.add_argument("--target", default="646 tasks", help="Target sidechat (default: 646 tasks)")
|
||||
|
||||
onboard_sub.add_parser("feed-matrix", parents=[common], help="Display all agents ranked by feeding weight, job volume, role, and work done over time")
|
||||
onboard_sub.add_parser("connects", parents=[common], help="Inventory of all active fleet nodes & client onboard connects")
|
||||
|
||||
# Domain: KPI
|
||||
p_kpi = subparsers.add_parser("kpi", parents=[common], help="Fleet KPI, spend monitoring & runtime preservation")
|
||||
p_kpi.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node")
|
||||
kpi_sub = p_kpi.add_subparsers(dest="kpi_action")
|
||||
|
||||
p_kpi_status = kpi_sub.add_parser("status", parents=[common], help="Show fleet KPI dashboard (default)")
|
||||
p_kpi_status.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node")
|
||||
|
||||
p_kpi_report = kpi_sub.add_parser("report", parents=[common], help="Detailed KPI & preservation report for an agent")
|
||||
p_kpi_report.add_argument("node", choices=VALID_NODES, help="Target node")
|
||||
|
||||
p_kpi_routes = kpi_sub.add_parser("routes", parents=[common], help="Check route and CDP health across agents")
|
||||
|
||||
p_kpi_spawn = kpi_sub.add_parser("spawn-worker", parents=[common], help="Spawn background tmux worker to preserve runtime")
|
||||
p_kpi_spawn.add_argument("node", choices=VALID_NODES, help="Target agent")
|
||||
p_kpi_spawn.add_argument("session", help="Session label")
|
||||
p_kpi_spawn.add_argument("worker_command", help="Command to run in background")
|
||||
|
||||
# Domain: CHROMEBOX
|
||||
|
||||
p_chromebox = subparsers.add_parser("chromebox", parents=[common], help="Agent browser settings-menu toggles (chromebox RPA)")
|
||||
chrome_sub = p_chromebox.add_subparsers(dest="chrome_action")
|
||||
p_chrome_perm = chrome_sub.add_parser("permissions", parents=[common], help="Settings menu toggles (permissions + related tabs)")
|
||||
@@ -5617,13 +5954,32 @@ def main():
|
||||
if len(sys.argv) > 1:
|
||||
if sys.argv[1] == "tui":
|
||||
tui_args = sys.argv[2:]
|
||||
cmd = [sys.executable, str(BIN_DIR / "muse-tui.py"), "--mode", "box"] + tui_args
|
||||
if tui_args and tui_args[0] in ("onboard", "tmux", "connects", "approvals"):
|
||||
cmd = [sys.executable, str(BIN_DIR / "box-onboard-tui.py")] + tui_args[1:]
|
||||
else:
|
||||
cmd = [sys.executable, str(BIN_DIR / "muse-tui.py"), "--mode", "box"] + tui_args
|
||||
res = subprocess.run(cmd)
|
||||
sys.exit(res.returncode)
|
||||
elif sys.argv[1] in ("onboard-tui", "dev-tui"):
|
||||
cmd = [sys.executable, str(BIN_DIR / "box-onboard-tui.py")] + sys.argv[2:]
|
||||
res = subprocess.run(cmd)
|
||||
sys.exit(res.returncode)
|
||||
elif sys.argv[1] == "tmux":
|
||||
if len(sys.argv) > 2 and sys.argv[2] in ("tally", "auto", "watch", "once", "match", "rules", "status"):
|
||||
if sys.argv[2] == "auto":
|
||||
t_sub = sys.argv[3:] or ["status"]
|
||||
else:
|
||||
t_sub = sys.argv[2:]
|
||||
cmd = [sys.executable, str(BIN_DIR / "tmux_auto_approver.py")] + t_sub
|
||||
res = subprocess.run(cmd)
|
||||
sys.exit(res.returncode)
|
||||
cmd = [sys.executable, str(BIN_DIR / "muse-tmux.py")] + (sys.argv[2:] or ["list"])
|
||||
res = subprocess.run(cmd)
|
||||
sys.exit(res.returncode)
|
||||
elif sys.argv[1] in ("tmux-auto", "auto-dev"):
|
||||
cmd = [sys.executable, str(BIN_DIR / "tmux_auto_approver.py")] + sys.argv[2:]
|
||||
res = subprocess.run(cmd)
|
||||
sys.exit(res.returncode)
|
||||
elif sys.argv[1] in ("docs", "doc"):
|
||||
cmd = [sys.executable, str(BIN_DIR / "docs-lookup.py")] + sys.argv[2:]
|
||||
res = subprocess.run(cmd)
|
||||
@@ -5900,6 +6256,8 @@ def main():
|
||||
cmd_approvals(args)
|
||||
elif args.domain == "muse-choices":
|
||||
cmd_muse_choices(args)
|
||||
elif args.domain == "runtime":
|
||||
cmd_runtime(args)
|
||||
elif args.domain == "invite":
|
||||
cmd_invite(args)
|
||||
elif args.domain == "usage":
|
||||
@@ -5908,6 +6266,8 @@ def main():
|
||||
cmd_settings(args)
|
||||
elif args.domain == "onboard":
|
||||
cmd_onboard(args)
|
||||
elif args.domain == "kpi":
|
||||
cmd_kpi(args)
|
||||
elif args.domain == "chromebox":
|
||||
cmd_chromebox(args)
|
||||
elif args.domain in ("deploy", "subagent"):
|
||||
|
||||
Reference in New Issue
Block a user