#!/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 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 ", } 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: state.detail = f"Redemption failed: {redemption.get('reason')} - {redemption.get('detail')}" 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 --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 main() -> None: parser = argparse.ArgumentParser(description="End-to-end agent-driven onboarding & invite salvage pipeline") sub = parser.add_subparsers(dest="action") 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 == "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()