#!/usr/bin/env python3 """ box-work.py ā Fleet Workspace, Work Scope, and Task Orchestration Engine. Provides unified visibility into: - Scope of cloud workers (ports, tunnel status, busy/idle signals) - Recent Gitea tickets & build tasks - Actions taken (PRs, merges, closed tasks) - Active state of related agent chats & main chat - Next-action identification & autonomous dispatch (start, assign, merge) Usable standalone or as `box work` / `super work`. Works on NetVM (bl), VM, or remote PC/VPS. """ import sys import os import re import json import time import socket import argparse import urllib.request import urllib.parse import urllib.error import hashlib from datetime import datetime, timezone from pathlib import Path try: import agent_cognitive_probe as acp except ImportError: acp = None # Color helpers USE_COLOR = sys.stdout.isatty() or os.environ.get("CLICOLOR_FORCE") == "1" def c_bold(s: str) -> str: return f"\033[1m{s}\033[0m" if USE_COLOR else str(s) def c_dim(s: str) -> str: return f"\033[2m{s}\033[0m" if USE_COLOR else str(s) def c_green(s: str) -> str: return f"\033[32m{s}\033[0m" if USE_COLOR else str(s) def c_red(s: str) -> str: return f"\033[31m{s}\033[0m" if USE_COLOR else str(s) def c_yellow(s: str) -> str: return f"\033[33m{s}\033[0m" if USE_COLOR else str(s) def c_blue(s: str) -> str: return f"\033[34m{s}\033[0m" if USE_COLOR else str(s) def c_cyan(s: str) -> str: return f"\033[36m{s}\033[0m" if USE_COLOR else str(s) def c_magenta(s: str) -> str: return f"\033[35m{s}\033[0m" if USE_COLOR else str(s) # Known fleet worker topology WORKERS = [ {"name": "opm", "role": "fleet-agent", "port": 2228, "desc": "Fleet Ops & Coordination"}, {"name": "646", "role": "fleet-agent", "port": 2226, "desc": "Fleet Ops & Verification"}, {"name": "dev", "role": "builder", "port": 2230, "desc": "Core Platform Builder"}, {"name": "pip", "role": "builder", "port": 2227, "desc": "Integration & Python Builder"}, {"name": "def", "role": "fleet-agent", "port": 2229, "desc": "Fleet Autonomous Worker"}, {"name": "muse", "role": "fleet-agent","port": 2225, "desc": "Chat & TUI Runner"}, {"name": "muse-main", "role": "host", "port": 2224, "desc": "Primary Runtime Host"} ] def find_repo_root() -> Path: if os.environ.get("NETVM_ROOT"): return Path(os.environ["NETVM_ROOT"]) cur = Path(__file__).resolve().parent while cur != cur.parent: if (cur / "fleet" / "partition-table.json").exists(): return cur cur = cur.parent fallback = Path("/home/super/Projects/NetVM") if fallback.exists(): return fallback return Path.cwd() REPO_ROOT = find_repo_root() PARTITION_TABLE_PATH = REPO_ROOT / "fleet" / "partition-table.json" TASKS_DIR = REPO_ROOT / "fleet" / "tasks" CHAT_LOG = REPO_ROOT / "logs" / "chat-history.jsonl" DEFAULT_GITEA_URL = "https://tea.muse-dev.online" def get_gitea_config(): token = os.environ.get("GITEA_TOKEN", "") url = os.environ.get("GITEA_URL", DEFAULT_GITEA_URL) # Try reading partition table if not token and PARTITION_TABLE_PATH.exists(): try: with open(PARTITION_TABLE_PATH) as f: data = json.load(f) token = data.get("contributors", {}).get("super", {}).get("token", "") url = data.get("gitea_url", url) except Exception: pass # Check if we are physically running on bl and port 3000 is open host_is_bl = False try: host_is_bl = (socket.gethostname() == "bl") except Exception: pass if host_is_bl: s = socket.socket() s.settimeout(0.3) if s.connect_ex(("127.0.0.1", 3000)) == 0: api_base = "http://127.0.0.1:3000/api/v1" else: api_base = f"{url.rstrip('/')}/api/v1" s.close() else: api_base = f"{url.rstrip('/')}/api/v1" if not token: token = "3c26744525bceaf385aa09737f7e41af613627b6" return api_base, token def gitea_api_request(endpoint: str, method: str = "GET", data: dict = None): api_base, token = get_gitea_config() url = f"{api_base}{endpoint}" headers = { "Authorization": f"token {token}", "Content-Type": "application/json", "User-Agent": "Box-Work-CLI/1.0" } payload = json.dumps(data).encode("utf-8") if data else None req = urllib.request.Request(url, data=payload, headers=headers, method=method) try: with urllib.request.urlopen(req, timeout=6.0) as resp: content = resp.read().decode("utf-8") return json.loads(content) if content else {} except urllib.error.HTTPError as e: body = e.read().decode("utf-8") try: return {"error": e.code, "message": json.loads(body).get("message", body)} except Exception: return {"error": e.code, "message": body} except Exception as e: return {"error": 500, "message": str(e)} def get_claimed_tasks(): claimed = {} cdir = TASKS_DIR / "claimed" if cdir.exists() and cdir.is_dir(): for f in cdir.iterdir(): if f.is_file() and not f.name.startswith("."): parts = f.name.split(".") agent = parts[-1] if len(parts) > 1 else "unknown" task_name = parts[0] claimed[agent] = task_name return claimed def get_recent_done_tasks(limit=5): done = [] ddir = TASKS_DIR / "done" if ddir.exists() and ddir.is_dir(): files = [f for f in ddir.iterdir() if f.is_file() and not f.name.startswith(".")] files.sort(key=lambda x: x.stat().st_mtime, reverse=True) for f in files[:limit]: mtime = datetime.fromtimestamp(f.stat().st_mtime, tz=timezone.utc) done.append({"name": f.name, "mtime": mtime.strftime("%H:%M:%SZ")}) return done def get_recent_chat_events(limit=5, agent=None): events = [] if CHAT_LOG.exists(): try: with open(CHAT_LOG, "r") as f: lines = f.readlines() for line in reversed(lines): if not line.strip(): continue try: ev = json.loads(line) if agent and ev.get("agent") != agent: continue events.append(ev) if len(events) >= limit: break except Exception: pass except Exception: pass return events def get_last_agent_chats(): last_chats = {} if CHAT_LOG.exists(): try: with open(CHAT_LOG, "r") as f: for line in f: if not line.strip(): continue try: ev = json.loads(line) agent = ev.get("agent") if agent: last_chats[agent] = ev except Exception: pass except Exception: pass return last_chats def check_tunnel_ports(): ports_status = {} s_vm = socket.socket() s_vm.settimeout(0.5) vm_online = (s_vm.connect_ex(("100.81.31.9", 22)) == 0) s_vm.close() for w in WORKERS: ports_status[w["port"]] = "UNKNOWN" if vm_online: try: cmd = "ssh -o ConnectTimeout=2 -o BatchMode=yes super@100.81.31.9 'ss -tlnH sport = :2224 or sport = :2225 or sport = :2226 or sport = :2227 or sport = :2228 or sport = :2229 or sport = :2230' 2>/dev/null" res = os.popen(cmd).read() for w in WORKERS: p = w["port"] if f":{p} " in res or f":{p}\n" in res: ports_status[p] = "UP" else: ports_status[p] = "DARK" except Exception: pass else: for w in WORKERS: p = w["port"] s = socket.socket() s.settimeout(0.1) ports_status[p] = "UP" if s.connect_ex(("127.0.0.1", p)) == 0 else "DARK" s.close() return ports_status def badge_status(status: str) -> str: if status == "PASS": return c_green("PASS") elif status == "WARN": return c_yellow("WARN") else: return c_red("FAIL") def check_agent_preflight(agent_name: str) -> dict: """Ensures hatch, restore, and git config health before assigning work to cloud muse agents.""" worker = next((w for w in WORKERS if w["name"] == agent_name), None) if not worker and agent_name != "super": return { "agent": agent_name, "port": 0, "hatch": {"status": "FAIL", "details": f"Unknown agent '{agent_name}'"}, "restore": {"status": "FAIL", "details": "Not listed in fleet topology"}, "git": {"status": "FAIL", "details": "No partition entry"}, "overall": "FAIL", "ready": False, "reasons": [f"Agent '{agent_name}' is not in fleet topology"] } port = worker["port"] if worker else 2224 tunnel_ports = check_tunnel_ports() port_status = tunnel_ports.get(port, "DARK") reasons = [] # 1. HATCH HEALTH (tunnel listener + responsive chat) hatch_status = "PASS" hatch_details = [] if port_status == "UP": hatch_details.append(f"Port {port} listener UP") else: hatch_status = "FAIL" hatch_details.append(f"Port {port} reverse tunnel DARK") reasons.append(f"Hatch tunnel is DOWN on port {port}. Container is offline or unreachable.") last_chats = get_last_agent_chats() chat_ev = last_chats.get(agent_name) if chat_ev: ts_str = chat_ev.get("ts", "")[:19].replace("T", " ") hatch_details.append(f"Chat active ({ts_str})") else: hatch_details.append("No recent chat entries") # 2. RESTORE HEALTH (NODES.md, supervisor persistence) restore_status = "PASS" restore_details = [] nodes_file = REPO_ROOT / "NODES.md" node_in_registry = False if nodes_file.exists(): try: with open(nodes_file) as f: content = f.read() if f"| {agent_name} |" in content or f"warp-{agent_name}" in content: node_in_registry = True except Exception: pass if node_in_registry or agent_name in ("muse-main", "super"): restore_details.append("Registered in NODES.md") else: restore_status = "WARN" restore_details.append("Not found in NODES.md") if port_status == "UP": restore_details.append("Watchdog/Supervisor persistent") else: restore_status = "FAIL" restore_details.append("Container rebuild / tunnel recovery pending") reasons.append("Container requires recovery/restore (run recover-after-rebuild or inspect watchdog).") # 3. GIT CONFIG HEALTH (partition token, collaborator access, branches) git_status = "PASS" git_details = [] token = "" if PARTITION_TABLE_PATH.exists(): try: with open(PARTITION_TABLE_PATH) as f: pt = json.load(f) contributor = pt.get("contributors", {}).get(agent_name) if contributor: token = contributor.get("token", "") git_details.append("Token in partition-table") else: git_status = "FAIL" git_details.append("Missing from partition-table") reasons.append(f"Agent '{agent_name}' has no credentials in fleet/partition-table.json") except Exception as e: git_status = "WARN" git_details.append(f"Partition table error: {e}") collab_check = gitea_api_request(f"/repos/super/box/collaborators/{agent_name}") if isinstance(collab_check, dict) and collab_check.get("error") and collab_check.get("error") not in (200, 204): git_status = "FAIL" git_details.append("Not a repository collaborator") reasons.append(f"Gitea user '{agent_name}' lacks write/collaborator access") else: git_details.append("Gitea collaborator OK") branches = gitea_api_request("/repos/super/box/branches") agent_branch = False if isinstance(branches, list): for b in branches: bname = b.get("name", "") if bname.startswith(f"dev/{agent_name}/") or bname.startswith(f"builder/{agent_name}/"): agent_branch = True break if agent_branch: git_details.append("Branch verified in Gitea") else: git_details.append("No active branch") overall = "PASS" if hatch_status == "FAIL" or restore_status == "FAIL" or git_status == "FAIL": overall = "FAIL" elif hatch_status == "WARN" or restore_status == "WARN" or git_status == "WARN": overall = "WARN" return { "agent": agent_name, "port": port, "hatch": {"status": hatch_status, "details": ", ".join(hatch_details)}, "restore": {"status": restore_status, "details": ", ".join(restore_details)}, "git": {"status": git_status, "details": ", ".join(git_details)}, "overall": overall, "ready": (overall != "FAIL"), "reasons": reasons } def cmd_check(args): target_agent = getattr(args, "agent", None) targets = [target_agent] if target_agent else [w["name"] for w in WORKERS if w["role"] != "host"] print(c_bold("\n=== BOX WORK: PRE-FLIGHT HEALTH VERIFICATION ===\n")) header = f"{'AGENT':<12} {'HATCH':<12} {'RESTORE':<12} {'GIT CONFIG':<12} {'STATUS'}" print(c_dim(header)) print(c_dim("-" * len(header))) for ag in targets: res = check_agent_preflight(ag) h_badge = badge_status(res["hatch"]["status"]) r_badge = badge_status(res["restore"]["status"]) g_badge = badge_status(res["git"]["status"]) overall_badge = c_green("š¢ READY") if res["ready"] else c_red("š“ BLOCKED") print(f"{c_bold(ag):<21} {h_badge:<21} {r_badge:<21} {g_badge:<21} {overall_badge}") print() blocked = [ag for ag in targets if not check_agent_preflight(ag)["ready"]] if blocked: print(c_bold("--- PRE-FLIGHT DIAGNOSTIC DETAILS ---")) for ag in blocked: res = check_agent_preflight(ag) print(f" {c_bold(ag)}:") print(f" ⢠Hatch: {res['hatch']['details']}") print(f" ⢠Restore: {res['restore']['details']}") print(f" ⢠Git: {res['git']['details']}") print() def heal_agent(agent_name: str) -> dict: """Automated remediation for an agent failing pre-flight health checks.""" worker = next((w for w in WORKERS if w["name"] == agent_name), None) actions = [] unresolved = [] if not worker and agent_name != "super": return { "agent": agent_name, "healed": False, "actions": [], "unresolved": [f"Unknown worker '{agent_name}'"] } port = worker["port"] if worker else 2224 actions.append(f"Analyzing pre-flight health state for {agent_name} (port {port})") # 1. Ensure Gitea Collaborator & Partition Table token = "" if PARTITION_TABLE_PATH.exists(): try: with open(PARTITION_TABLE_PATH) as f: pt = json.load(f) contributor = pt.get("contributors", {}).get(agent_name) if contributor: token = contributor.get("token", "") except Exception: pass if not token: token = hashlib.sha256(f"{agent_name}-gitea-token".encode()).hexdigest()[:40] actions.append(f"Generated partition token for {agent_name}") collab_res = gitea_api_request(f"/repos/super/box/collaborators/{agent_name}", method="PUT", data={"permission": "write"}) actions.append(f"Ensured Gitea collaborator write access for {agent_name}") # 2. Container Workspace Injection if SSH dialable if port in (2224, 2228): try: cmd = f"ssh -o ConnectTimeout=3 -o BatchMode=yes -o StrictHostKeyChecking=no super@100.81.31.9 'ssh -o StrictHostKeyChecking=no -i /home/super/.ssh/fleet -p {port} muse@localhost \"git config --global credential.helper store && echo \\\"https://{agent_name}:{token}@tea.muse-dev.online\\\" > ~/.git-credentials && chmod 600 ~/.git-credentials\"' 2>/dev/null" if os.system(cmd) == 0: actions.append(f"Directly injected Git credentials into {agent_name} container") except Exception: pass # 3. Check and heal Hatch / Reverse Tunnel tunnel_ports = check_tunnel_ports() if tunnel_ports.get(port) == "UP": actions.append(f"Hatch reverse tunnel verified UP on port {port}") else: chat_script = REPO_ROOT / "bin" / "muse-chat-api.py" if chat_script.exists(): heal_msg = f"[HEAL NUDGE] Reverse tunnel on port {port} is DOWN. Please run 'chmod 600 ~/.ssh/authorized_keys' and restart tunnel with '~/workspace/bin/gcp-tunnel-up.sh &' (or 'cloud-uptime/recover-after-rebuild.sh'). Git clone URL: https://{agent_name}:{token}@tea.muse-dev.online/super/box.git" os.system(f"python3 {chat_script} --account {agent_name} send '{heal_msg}' >/dev/null 2>&1") actions.append(f"Dispatched tunnel restart & git clone command to {agent_name} chat") time.sleep(1.0) recheck_ports = check_tunnel_ports() if recheck_ports.get(port) == "UP": actions.append(f"Reverse tunnel on port {port} came online during healing!") else: unresolved.append(f"Reverse tunnel on port {port} is still DOWN (waiting for agent container execution)") final_preflight = check_agent_preflight(agent_name) healed = final_preflight["ready"] if not healed and not unresolved: unresolved.extend(final_preflight["reasons"]) return { "agent": agent_name, "healed": healed, "actions": actions, "unresolved": unresolved } def cmd_heal(args): agent = args.agent print(c_bold(f"\n=== BOX WORK: HEALING AGENT '{agent}' ===\n")) res = heal_agent(agent) print(c_bold("Actions taken:")) for a in res["actions"]: print(f" {c_green('ā')} {a}") print() if res["healed"]: print(c_green(f"š Agent '{agent}' successfully healed and ready for assignments!\n")) else: print(c_yellow(f"ā ļø Agent '{agent}' partially healed with open issues:")) for u in res["unresolved"]: print(f" ⢠{u}") print() def cmd_status(args): api_base, _ = get_gitea_config() print(c_bold(f"\n=== BOX WORK: FLEET & BUILD PIPELINE ({api_base}) ===\n")) # 1. Workers Scope & Live Signals print(c_bold("--- WORKER SCOPE & CONSTANT SIGNALS ---")) claimed_tasks = get_claimed_tasks() tunnel_ports = check_tunnel_ports() last_chats = get_last_agent_chats() # Query Gitea open issues for assignment signals issues = gitea_api_request("/repos/super/box/issues?state=open") if isinstance(issues, dict) and "error" in issues: issues = [] agent_active_issues = {} for iss in issues: assignee = iss.get("assignee") if assignee: uname = assignee.get("username") agent_active_issues[uname] = iss headers = f"{'AGENT':<10} {'ROLE':<12} {'PORT':<6} {'TUNNEL':<7} {'COGNITIVE':<15} {'ACTIVE WORK / ASSIGNMENT':<38} {'LAST CHAT'}" print(c_dim(headers)) print(c_dim("-" * len(headers))) ready_count = 0 busy_count = 0 dark_count = 0 for w in WORKERS: name = w["name"] role = w["role"] port = w["port"] tunnel = tunnel_ports.get(port, "DARK") active_task = claimed_tasks.get(name) active_issue = agent_active_issues.get(name) # Live cognitive probe cog_badge = "-" if acp and name != "muse-main": cog = acp.get_passive_cognitive_state(name) cog_status = cog.get("status", "DARK") cog_badge = acp.format_cognitive_badge(cog_status) elif name == "muse-main": cog_badge = c_dim("HOST") if active_issue: num = active_issue.get("number") title = active_issue.get("title", "")[:32] labels = [l.get("name") for l in active_issue.get("labels", [])] signal = c_red("š“ BUSY") work_desc = f"#{num} {title}" busy_count += 1 elif active_task: signal = c_yellow("š” CLAIM") work_desc = active_task[:36] busy_count += 1 elif tunnel == "DARK" and role != "host": signal = c_dim("ā« DARK") work_desc = c_dim("Tunnel down / no listener") dark_count += 1 else: signal = c_green("š¢ IDLE") work_desc = c_dim("Ready for assignment") ready_count += 1 tunnel_str = c_green("UP") if tunnel == "UP" else (c_red("DARK") if tunnel == "DARK" else c_dim(tunnel)) chat_ev = last_chats.get(name) if chat_ev: ts_str = chat_ev.get("ts", "") author = chat_ev.get("author", "") chat_str = f"{author} ({ts_str[11:16]}Z)" else: chat_str = c_dim("-") print(f"{c_bold(name):<19} {role:<12} {port:<6} {tunnel_str:<16} {cog_badge:<24} {work_desc:<38} {chat_str}") print() # 2. Tickets & Tasks Pipeline print(c_bold("--- GITEA TICKETS & BUILD TASKS (super/box) ---")) all_issues = gitea_api_request("/repos/super/box/issues?state=all&limit=8") if isinstance(all_issues, dict) and "error" in all_issues: print(c_red(f" Failed to fetch tickets: {all_issues.get('message')}")) elif not all_issues: print(c_dim(" No tickets found in Gitea repository.")) else: t_header = f"{'TICKET':<8} {'STATE':<10} {'ASSIGNEE':<12} {'TITLE':<48} {'LABELS'}" print(c_dim(t_header)) print(c_dim("-" * len(t_header))) for iss in all_issues: num = f"#{iss.get('number')}" state = iss.get("state", "").upper() assignee = iss.get("assignee") assignee_str = assignee.get("username", "-") if assignee else "-" title = iss.get("title", "")[:46] lbls = [l.get("name") for l in iss.get("labels", [])] if state == "CLOSED": state_str = c_dim("CLOSED") elif "in-review" in lbls: state_str = c_yellow("IN-REVIEW") else: state_str = c_green("OPEN") lbl_str = c_cyan(", ".join(lbls)) if lbls else "-" print(f"{c_bold(num):<17} {state_str:<19} {assignee_str:<12} {title:<48} {lbl_str}") print() # 3. Pull Requests prs = gitea_api_request("/repos/super/box/pulls?state=all&limit=5") if prs and isinstance(prs, list): print(c_bold("--- PULL REQUESTS & CODE INTEGRATIONS ---")) pr_header = f"{'PR':<8} {'STATUS':<10} {'BRANCH':<34} {'TITLE':<42}" print(c_dim(pr_header)) print(c_dim("-" * len(pr_header))) for pr in prs: pnum = f"#{pr.get('number')}" merged = pr.get("merged", False) state = pr.get("state", "").upper() status_str = c_green("MERGED") if merged else (c_yellow("OPEN") if state == "OPEN" else c_dim("CLOSED")) head = pr.get("head", {}).get("ref", "-")[:32] title = pr.get("title", "")[:40] print(f"{c_bold(pnum):<17} {status_str:<19} {head:<34} {title:<42}") print() # 4. Recent Done Tasks done_tasks = get_recent_done_tasks(limit=4) if done_tasks: print(c_bold("--- RECENTLY ARCHIVED TASKS (fleet/tasks/done) ---")) for dt in done_tasks: print(f" {c_green('ā')} {dt['name']} {c_dim('(' + dt['mtime'] + ')')}") print() # 5. Active Chat Snippets chat_events = get_recent_chat_events(limit=3) if chat_events: print(c_bold("--- ACTIVE CHAT CONVERSATIONS ---")) for ev in chat_events: agent = ev.get("agent", "agent") tname = ev.get("thread_name", "Chat") author = ev.get("author", "user") text = ev.get("text", "").replace("\n", " ")[:90] ts = ev.get("ts", "")[11:16] print(f" [{c_cyan(agent)}:{c_dim(tname)}] {c_dim(ts)} {c_bold(author)}: {text}...") print() # 6. Identified Work & Action Recommendations print(c_bold("--- IDENTIFIED WORK & DISPATCH RECOMMENDATIONS ---")) pending_tasks = [] pdir = TASKS_DIR / "pending" if pdir.exists() and pdir.is_dir(): pending_tasks = [f.name for f in pdir.iterdir() if f.is_file() and not f.name.startswith(".")] recs = [] if ready_count > 0: idle_agents = [w["name"] for w in WORKERS if w["role"] != "host" and tunnel_ports.get(w["port"]) == "UP" and w["name"] not in agent_active_issues and w["name"] not in claimed_tasks] recs.append(f"Available Workers: {', '.join(idle_agents) if idle_agents else 'None'} ready for new build tickets.") if dark_count > 0: dark_nodes = [w["name"] for w in WORKERS if tunnel_ports.get(w["port"]) == "DARK" and w["role"] != "host"] recs.append(f"Dark Node Recovery: Nodes {', '.join(dark_nodes)} reverse tunnels are DOWN (need tunnel supervision).") if pending_tasks: recs.append(f"Unassigned Pending Queue: {len(pending_tasks)} task(s) waiting in fleet/tasks/pending/: {', '.join(pending_tasks[:3])}") open_prs = [pr for pr in (prs if isinstance(prs, list) else []) if pr.get("state") == "open" and not pr.get("merged")] if open_prs: recs.append(f"Open PRs: {len(open_prs)} PR(s) ready for test verification & merge: #{open_prs[0].get('number')} ({open_prs[0].get('title', '')[:30]})") for r in recs: print(f" {c_yellow('š')} {r}") print(f"\n{c_dim('Quick Dispatch:')} {c_cyan('box work start