#!/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 from datetime import datetime, timezone from pathlib import Path # 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): 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) 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']}") for r in res["reasons"]: print(f" {c_yellow('!')} {r}") print() 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':<12} {'ROLE':<13} {'PORT':<6} {'TUNNEL':<8} {'SIGNAL':<10} {'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) 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):<21} {role:<13} {port:<6} {tunnel_str:<17} {signal:<19} {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