Files

588 lines
23 KiB
Python
Raw Permalink Normal View History

#!/usr/bin/env python3
"""onboard_pipeline.py — End-to-end agent-driven onboarding & invite salvage pipeline.
Orchestrates the complete lifecycle:
1. Identifies the most urgent beneficiary agent in need of tokens (blocked first, then lowest balance).
2. Provisions infrastructure (WireGuard, dedicated netns, CDP port, chrome-box profile).
3. Launches authentication via cred_client / onboard-driver.
4. Waits for / accepts OTP code.
5. Checks age verification gate (Tailscale portal / Instagram link).
6. Automatically redeems the urgent agent's invite code on the newly onboarded node.
7. Injects operational DRIVE and marks the node active.
8. Can dispatch cryptographically signed work orders ([WO]) to prompting sidechats.
"""
from __future__ import annotations
import argparse
import hashlib
import json
import os
import subprocess
import sys
import time
from dataclasses import asdict, dataclass
from pathlib import Path
from typing import Any, Dict, List, Optional, Tuple
REPO_ROOT = Path(__file__).resolve().parent.parent
BIN_DIR = REPO_ROOT / "bin"
sys.path.insert(0, str(BIN_DIR))
try:
import invite
from cred_client import CredClient
from settings_rpa import SettingsRPA
except ImportError:
pass
STAGE_INFRA = "infra_provisioned"
STAGE_INITIATE = "auth_initiated"
STAGE_AWAIT_OTP = "awaiting_otp"
STAGE_AUTH_ACTIVE = "auth_active"
STAGE_REDEEMED = "invite_redeemed"
STAGE_DRIVE_INJECTED = "drive_injected"
STAGE_COMPLETED = "completed"
STATE_DIR = Path(os.environ.get("NETVM_ONBOARD_STATE", "/tmp/netvm-onboard"))
@dataclass
class OnboardState:
node: str
email: str
beneficiary_node: Optional[str]
invite_code: Optional[str]
stage: str
cdp_port: Optional[int] = None
created_at: float = 0.0
updated_at: float = 0.0
detail: Optional[str] = None
redemption_result: Optional[Dict[str, Any]] = None
def save(self) -> None:
STATE_DIR.mkdir(parents=True, exist_ok=True)
path = STATE_DIR / f"{self.node}.json"
self.updated_at = time.time()
with open(path, "w") as f:
json.dump(asdict(self), f, indent=2)
@classmethod
def load(cls, node: str) -> Optional[OnboardState]:
path = STATE_DIR / f"{node}.json"
if not path.exists():
return None
try:
with open(path) as f:
data = json.load(f)
return cls(**data)
except Exception:
return None
# Canonical agent operational roles in NetVM
AGENT_ROLES: Dict[str, str] = {
"646": "Production Lead / Swarm Execution",
"opm": "Fleet Coordinator / Loop Orchestrator",
"pip": "Production Agent / Pipeline Worker",
"muse": "Core Engine / Auditor & Dev",
"dev": "Development / Test & Staging",
"def": "Defense / Security & Standby",
}
# Base role priority multiplier
ROLE_WEIGHT_MULTIPLIER: Dict[str, float] = {
"646": 1.5, # Critical production executor
"opm": 1.4, # Central coordinator & orchestrator
"pip": 1.2, # High volume production worker
"muse": 1.1, # Auditor & core engine
"dev": 0.8, # Test & staging
"def": 0.7, # Defense & standby
}
def get_agent_metrics() -> Tuple[Dict[str, int], Dict[str, int]]:
"""Count job assignments and cumulative chat activity (work done over time) per agent."""
from collections import Counter
job_counts = Counter()
jobs_dir = REPO_ROOT / "jobs"
if jobs_dir.exists():
for p in jobs_dir.glob("*.json"):
try:
d = json.loads(p.read_text(encoding="utf-8"))
a = d.get("target_agent") or d.get("agent")
if a:
job_counts[a] += 1
except Exception:
pass
msg_counts = Counter()
chat_log = REPO_ROOT / "logs" / "chat-history.jsonl"
if chat_log.exists():
try:
with open(chat_log, encoding="utf-8") as f:
for line in f:
try:
row = json.loads(line)
s = row.get("sender") or row.get("from") or row.get("agent")
if s:
msg_counts[s] += 1
except Exception:
pass
except Exception:
pass
return dict(job_counts), dict(msg_counts)
def calculate_feeding_weights(nodes: Optional[List[str]] = None) -> List[Dict[str, Any]]:
"""Compute feeding weights for all agents considering:
- Urgency: Blocked (highest), weekly used %, token depletion
- Work done over time: Messages processed / chat history volume (strongest operational weight)
- Job amount: Active and assigned jobs in registry
- Role importance: Multiplier based on operational criticality
"""
target_nodes = nodes or ["646", "opm", "pip", "muse", "dev", "def"]
job_counts, msg_counts = get_agent_metrics()
usage_map = invite.fleet_usage(target_nodes)
invites_map = invite.fleet_invite_status(target_nodes)
rows = []
for n in target_nodes:
u = usage_map.get(n, {})
inv = invites_map.get(n, {})
code = inv.get("code") or "-"
wu = u.get("weekly_used_pct") or 0
extra_left = u.get("additional_left") or "-"
extra_used = u.get("additional_used_pct") or 0
is_blocked = (wu >= 100 and "0 tokens left" in str(extra_left)) or (wu >= 100 and extra_used >= 100)
jobs = job_counts.get(n, 0)
work_done = msg_counts.get(n, 0)
role_desc = AGENT_ROLES.get(n, "Agent Worker")
role_mult = ROLE_WEIGHT_MULTIPLIER.get(n, 1.0)
# Feeding score calculation:
# Base: Blocked = 1000 pts; Weekly limit 100% = 300 pts; proportional to weekly used %
urgency_score = 1000.0 if is_blocked else (300.0 if wu >= 100 else float(wu))
# Work done weight: 1 point per 10 messages (strongest historical work indicator)
work_score = float(work_done) * 0.1
# Job weight: 2 points per assigned job
job_score = float(jobs) * 2.0
# Composite feeding weight
feeding_weight = round((urgency_score + work_score + job_score) * role_mult, 1)
status_label = "BLOCKED" if is_blocked else ("LIMIT_REACHED" if wu >= 100 else "ACTIVE")
rows.append({
"node": n,
"role": role_desc,
"code": code,
"status": status_label,
"weekly_used_pct": wu,
"additional_left": extra_left,
"jobs_count": jobs,
"work_done_msgs": work_done,
"feeding_weight": feeding_weight,
"is_blocked": is_blocked,
"invites_left": inv.get("uses_remaining", 30),
})
# Sort descending by feeding weight
rows.sort(key=lambda r: r["feeding_weight"], reverse=True)
return rows
def select_urgent_beneficiary(explicit_node: Optional[str] = None, explicit_code: Optional[str] = None) -> Tuple[Optional[str], Optional[str], Optional[str]]:
"""Determine the optimal beneficiary agent to feed with Muse tokens and return (beneficiary_node, invite_code, reason).
Priority:
1. Explicit code / node requested by caller.
2. Highest feeding weight (combining urgency, work done over time, job volume, and role criticality).
"""
if explicit_code:
return explicit_node, explicit_code.strip().upper(), "Explicit invite code provided"
if explicit_node:
try:
inv = invite.get_invite(explicit_node)
if inv.get("ok") and inv.get("code"):
return explicit_node, inv["code"], f"Explicit beneficiary agent @{explicit_node}"
except Exception as e:
return explicit_node, None, f"Failed to fetch invite code for @{explicit_node}: {e}"
# Calculate weighted rankings across fleet
rankings = calculate_feeding_weights()
for row in rankings:
if row["code"] and row["code"] != "-":
n = row["node"]
code = row["code"]
weight = row["feeding_weight"]
st = row["status"]
role = row["role"]
work = row["work_done_msgs"]
reason = f"Top feeding weight: {weight} pts (@{n} [{role}] - status:{st}, work_done:{work} msgs, jobs:{row['jobs_count']})"
return n, code, reason
# Fallback to 646 or muse if available
for fallback in ["646", "pip", "muse", "opm"]:
try:
inv = invite.get_invite(fallback)
if inv.get("ok") and inv.get("code"):
return fallback, inv["code"], f"Default fleet agent @{fallback}"
except Exception:
pass
return None, None, "No active agent invite code found"
def provision_node_infra(node: str) -> Dict[str, Any]:
"""Execute ./bin/netvm-provision-node.sh <node> to setup netns, wireguard, chrome-box."""
script = BIN_DIR / "netvm-provision-node.sh"
cmd = [str(script), node]
res = subprocess.run(cmd, capture_output=True, text=True)
if res.returncode != 0:
return {
"ok": False,
"node": node,
"error": res.stderr.strip() or res.stdout.strip() or f"Exit {res.returncode}",
}
return {"ok": True, "node": node, "output": res.stdout.strip()}
def start_onboarding(node: str, email: str, beneficiary_node: Optional[str] = None, invite_code: Optional[str] = None, account_name: Optional[str] = None) -> Dict[str, Any]:
"""Phase 1 & 2: Provision infra, choose beneficiary invite code, and initiate authentication."""
b_node, code, reason = select_urgent_beneficiary(beneficiary_node, invite_code)
state = OnboardState(
node=node,
email=email,
beneficiary_node=b_node,
invite_code=code,
stage="starting",
created_at=time.time(),
updated_at=time.time(),
detail=reason,
)
state.save()
# 1. Provision infra
infra_res = provision_node_infra(node)
if not infra_res.get("ok"):
state.stage = "infra_failed"
state.detail = infra_res.get("error")
state.save()
return {"ok": False, "state": asdict(state), "error": f"Infra provisioning failed: {state.detail}"}
state.stage = STAGE_INFRA
state.save()
# 2. Initiate authentication
client = CredClient()
cred_res = client.initiate(node, email, service="muse", account_name=account_name)
st = cred_res.get("status")
if st == "awaiting_otp":
state.stage = STAGE_AWAIT_OTP
state.detail = f"OTP verification code sent to {email}"
state.save()
return {
"ok": True,
"status": "awaiting_otp",
"state": asdict(state),
"beneficiary_node": b_node,
"invite_code_queued": code,
"message": f"Verification code sent to {email}. Submit with: box onboard submit-otp --node {node} --otp <code>",
}
elif st == "active":
state.stage = STAGE_AUTH_ACTIVE
state.save()
# Immediately redeem queued code
return finish_onboarding_redemption(state)
else:
state.stage = "auth_error"
state.detail = cred_res.get("detail") or cred_res.get("message")
state.save()
return {"ok": False, "status": st, "state": asdict(state), "error": state.detail}
def submit_onboarding_otp(node: str, otp: str, email: Optional[str] = None) -> Dict[str, Any]:
"""Phase 3: Submit transient OTP and proceed to redemption upon success."""
state = OnboardState.load(node)
if not state:
state = OnboardState(
node=node,
email=email or "unknown",
beneficiary_node=None,
invite_code=None,
stage=STAGE_AWAIT_OTP,
created_at=time.time(),
)
client = CredClient()
res = client.submit_otp(node, otp, email=state.email or email)
st = res.get("status")
if st == "active":
state.stage = STAGE_AUTH_ACTIVE
state.detail = "Session authenticated successfully"
state.save()
return finish_onboarding_redemption(state)
else:
state.detail = res.get("detail") or res.get("message")
state.save()
return {"ok": False, "status": st, "state": asdict(state), "error": state.detail}
def finish_onboarding_redemption(state: OnboardState) -> Dict[str, Any]:
"""Phase 4 & 5: Redeem queued invite code and inject DRIVE."""
# Ensure beneficiary code is present
if not state.invite_code:
b_node, code, _ = select_urgent_beneficiary(state.beneficiary_node, None)
state.beneficiary_node = b_node
state.invite_code = code
redemption_info = None
if state.invite_code:
try:
# Redeem queued code on the fresh node
redemption = invite.redeem_invite(state.node, state.invite_code, timeout=20.0)
redemption_info = redemption
state.redemption_result = redemption
if redemption.get("ok"):
state.stage = STAGE_REDEEMED
state.detail = f"Successfully redeemed code {state.invite_code}! 1B tokens credited to @{state.beneficiary_node} and @{state.node}."
else:
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}"
# Phase 5: Inject DRIVE
try:
drive_script = BIN_DIR / "agent_md.py"
if drive_script.exists():
subprocess.run([sys.executable, str(drive_script), "inject-drive", state.node], capture_output=True, text=True, timeout=15)
state.stage = STAGE_COMPLETED
except Exception:
pass
state.save()
return {
"ok": True,
"status": "completed" if (redemption_info and redemption_info.get("ok")) else "auth_active_redemption_warning",
"state": asdict(state),
"beneficiary_node": state.beneficiary_node,
"invite_code": state.invite_code,
"redemption": redemption_info,
"message": f"Node @{state.node} is fully active! 1B tokens granted.",
}
def issue_salvage_work_order(blocked_node: str = "646", to_sidechat: str = "646 tasks") -> Dict[str, Any]:
"""Issue a cryptographically signed Work Order prompting fleet operators or agents to initiate onboarding."""
inv = invite.get_invite(blocked_node)
code = inv.get("code") or "UNKNOWN"
title = f"Salvage Blocked Agent @{blocked_node}"
body = (
f"Agent @{blocked_node} is out of tokens (code: {code}). "
f"Initiate client onboarding to grant 1 Billion tokens: box onboard <new_node> --email <client_email> --for {blocked_node}"
)
cmd = [
"python3", str(BIN_DIR / "super-cli.py"), "dm", "wo",
"--to", "opm",
"--target", to_sidechat,
"--title", title,
"--body", body,
"--priority", "urgent",
"--allow-main-chat"
]
res = subprocess.run(cmd, capture_output=True, text=True)
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")
p_start.add_argument("--for", dest="for_agent", help="Beneficiary agent to unblock (defaults to most urgent)")
p_start.add_argument("--code", help="Explicit 6-character invite code to redeem")
p_start.add_argument("--account-name", help="Display name hint")
p_start.add_argument("--json", action="store_true", help="Emit JSON output")
p_otp = sub.add_parser("submit-otp", help="Submit transient OTP verification code")
p_otp.add_argument("node", help="Node label")
p_otp.add_argument("otp", help="6-digit verification code")
p_otp.add_argument("--email", help="Client email (optional)")
p_otp.add_argument("--json", action="store_true", help="Emit JSON output")
p_status = sub.add_parser("status", help="Check onboarding pipeline state for a node")
p_status.add_argument("node", help="Node label")
p_status.add_argument("--json", action="store_true", help="Emit JSON output")
p_wo = sub.add_parser("salvage-wo", help="Dispatch salvage work order to sidechat")
p_wo.add_argument("node", nargs="?", default="646", help="Blocked agent node (default: 646)")
p_wo.add_argument("--target", default="646 tasks", help="Target sidechat")
p_feed = sub.add_parser("feed-matrix", help="Display all agents ranked by feeding weight, job volume, role, and work done over time")
p_feed.add_argument("--json", action="store_true", help="Emit JSON output")
args = parser.parse_args()
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))
elif res.get("ok"):
print(f"\n[ONBOARDING STARTED] Node: {args.node}")
print(f" Beneficiary: @{res.get('beneficiary_node')} (Invite Code: {res.get('invite_code_queued')})")
print(f" Message: {res.get('message')}\n")
else:
print(f"\n[ERROR] {res.get('error')}\n", file=sys.stderr)
sys.exit(1)
elif args.action == "submit-otp":
res = submit_onboarding_otp(args.node, args.otp, email=args.email)
if args.json:
print(json.dumps(res, indent=2))
elif res.get("ok"):
print(f"\n[ONBOARDING COMPLETE] {res.get('message')}")
print(f" Beneficiary Credited: @{res.get('beneficiary_node')} (+1,000,000,000 tokens)\n")
else:
print(f"\n[ERROR] {res.get('error')}\n", file=sys.stderr)
sys.exit(1)
elif args.action == "status":
state = OnboardState.load(args.node)
if not state:
print(f"No onboarding pipeline record found for {args.node}", file=sys.stderr)
sys.exit(1)
if args.json:
print(json.dumps(asdict(state), indent=2))
else:
print(f"\n=== ONBOARDING STATUS: @{state.node} ===")
print(f" Email: {state.email}")
print(f" Stage: {state.stage}")
print(f" Beneficiary: @{state.beneficiary_node} (Code: {state.invite_code})")
if state.detail:
print(f" Detail: {state.detail}")
print()
elif args.action == "salvage-wo":
res = issue_salvage_work_order(args.node, to_sidechat=args.target)
print(f"Salvage work order dispatched for @{args.node}: ok={res.get('ok')}")
elif args.action == "feed-matrix":
matrix = calculate_feeding_weights()
if args.json:
print(json.dumps(matrix, indent=2))
else:
print("\n=== FLEET FEEDING & TOKEN ALLOCATION MATRIX ===")
print(" (Weighted by: Urgency + Work Done Over Time + Job Amount + Role Criticality)\n")
print(f" {'RANK':4} {'AGENT':6} {'ROLE':38} {'FEED SCORE':11} {'STATUS':10} {'JOBS':6} {'WORK DONE':10} {'CODE':8}")
print(f" {'─'*4} {'─'*6} {'─'*38} {'─'*11} {'─'*10} {'─'*6} {'─'*10} {'─'*8}")
for idx, r in enumerate(matrix, 1):
print(f" #{idx:<3} {r['node']:<6} {r['role']:<38} {r['feeding_weight']:<11} {r['status']:<10} {r['jobs_count']:<6} {str(r['work_done_msgs']) + ' msgs':<10} {r['code']:<8}")
print()
if __name__ == "__main__":
main()