diff --git a/bin/box-work.py b/bin/box-work.py new file mode 100755 index 0000000..ebb4e17 --- /dev/null +++ b/bin/box-work.py @@ -0,0 +1,537 @@ +#!/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 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':<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 --to <agent>')} | {c_cyan('box work assign <ticket#> --to <agent>')} | {c_cyan('box work merge <pr#>')}\n") + +def cmd_start(args): + title = args.title + agent = args.agent + body = args.goal or f"Work task for {agent}: {title}" + + print(c_bold(f"Initiating work ticket for agent {agent}...")) + + # 1. Ensure label exists in Gitea + gitea_api_request("/repos/super/box/labels", method="POST", data={ + "name": f"assign:{agent}", + "color": "5319e7", + "description": f"Assigned directly to {agent}" + }) + + # 2. Create Gitea Issue + payload = { + "title": title, + "body": body, + "labels": [1, 2], # task, ready + "assignee": agent + } + + res = gitea_api_request("/repos/super/box/issues", method="POST", data=payload) + if "error" in res: + print(c_red(f"Error creating ticket in Gitea: {res.get('message')}")) + sys.exit(1) + + issue_num = res.get("number") + print(c_green(f"✓ Created Gitea Issue #{issue_num}: {title}")) + + # 3. Trigger webhook sweep on bridge if local + s = socket.socket() + s.settimeout(0.5) + if s.connect_ex(("127.0.0.1", 3005)) == 0: + try: + req = urllib.request.Request("http://127.0.0.1:3005/sweep") + urllib.request.urlopen(req, timeout=1.0) + print(c_green(f"✓ Reconciled bridge webhook queue")) + except Exception: + pass + s.close() + + # 4. Notify agent via muse-chat-api if available + chat_script = REPO_ROOT / "bin" / "muse-chat-api.py" + if chat_script.exists(): + msg = f"New build ticket #{issue_num} assigned to you: {title}. Clone/pull ~/workspace/box, checkout dev/{agent}/{issue_num}-work, commit citing 'Fixes #{issue_num}', and push." + try: + cmd = f"python3 {chat_script} --account {agent} send '{msg}'" + os.system(f"{cmd} >/dev/null 2>&1") + print(c_green(f"✓ Delivered briefing to {agent} chat session")) + except Exception: + pass + + print(c_bold(f"\nWork ticket #{issue_num} is active and assigned to {agent}.\n")) + +def cmd_assign(args): + issue_num = args.issue + agent = args.agent + print(c_bold(f"Assigning Ticket #{issue_num} to {agent}...")) + + payload = { + "assignee": agent + } + res = gitea_api_request(f"/repos/super/box/issues/{issue_num}", method="PATCH", data=payload) + if "error" in res: + print(c_red(f"Error updating ticket: {res.get('message')}")) + sys.exit(1) + + print(c_green(f"✓ Ticket #{issue_num} assigned to {agent}")) + + chat_script = REPO_ROOT / "bin" / "muse-chat-api.py" + if chat_script.exists(): + msg = f"Ticket #{issue_num} has been assigned to you. Please pull ~/workspace/box and claim." + os.system(f"python3 {chat_script} --account {agent} send '{msg}' >/dev/null 2>&1") + print(c_green(f"✓ Notified {agent} in chat")) + +def cmd_merge(args): + pr_num = args.pr + print(c_bold(f"Merging Pull Request #{pr_num} into master...")) + + payload = { + "Do": "merge", + "MergeTitleField": f"Merge pull request #{pr_num}", + "MergeMessageField": f"Merged via box work CLI" + } + res = gitea_api_request(f"/repos/super/box/pulls/{pr_num}/merge", method="POST", data=payload) + if isinstance(res, dict) and "error" in res: + print(c_red(f"Error merging PR: {res.get('message')}")) + sys.exit(1) + + print(c_green(f"✓ PR #{pr_num} merged into master. Post-receive hook triggered loop terminus.")) + +def cmd_chats(args): + agent = getattr(args, "agent", None) + events = get_recent_chat_events(limit=args.limit) + if agent: + events = [e for e in events if e.get("agent") == agent] + print(c_bold(f"\n=== CHAT FEED ({agent or 'ALL AGENTS'}) ===\n")) + for ev in events: + ag = ev.get("agent", "agent") + tname = ev.get("thread_name", "Chat") + author = ev.get("author", "user") + text = ev.get("text", "").strip() + ts = ev.get("ts", "")[:19].replace("T", " ") + print(f"[{c_cyan(ag)} : {c_dim(tname)}] {c_dim(ts)} {c_bold(author)}:\n{text}\n" + c_dim("-" * 60)) + print() + +def main(): + parser = argparse.ArgumentParser( + prog="box work", + description="Fleet Workspace, Work Scope, and Task Orchestration Engine." + ) + sub = parser.add_subparsers(dest="work_action") + + sub.add_parser("status", help="Show full operational work dashboard") + + p_start = sub.add_parser("start", help="Instantly start and assign new build ticket to an agent") + p_start.add_argument("title", help="Ticket title / summary") + p_start.add_argument("--to", dest="agent", required=True, help="Agent username (opm, 646, dev, pip, def, muse)") + p_start.add_argument("--goal", help="Optional detailed goal description") + + p_assign = sub.add_parser("assign", help="Assign existing ticket to an agent") + p_assign.add_argument("issue", type=int, help="Issue number (e.g. 215)") + p_assign.add_argument("--to", dest="agent", required=True, help="Agent username") + + p_merge = sub.add_parser("merge", help="Merge an open PR into master") + p_merge.add_argument("pr", type=int, help="Pull request number (e.g. 214)") + + p_chats = sub.add_parser("chats", help="View recent live chat activity") + p_chats.add_argument("--agent", help="Filter by agent name") + p_chats.add_argument("--limit", type=int, default=10, help="Number of messages to show") + + args = parser.parse_args() + action = args.work_action + + if not action or action == "status": + cmd_status(args) + elif action == "start": + cmd_start(args) + elif action == "assign": + cmd_assign(args) + elif action == "merge": + cmd_merge(args) + elif action == "chats": + cmd_chats(args) + else: + parser.print_help() + +if __name__ == "__main__": + main() diff --git a/bin/box_work.py b/bin/box_work.py new file mode 120000 index 0000000..ead364c --- /dev/null +++ b/bin/box_work.py @@ -0,0 +1 @@ +/home/super/Projects/NetVM/bin/box-work.py \ No newline at end of file diff --git a/tests/test_box_work.py b/tests/test_box_work.py new file mode 100644 index 0000000..7adb858 --- /dev/null +++ b/tests/test_box_work.py @@ -0,0 +1,47 @@ +import os +import sys +import unittest +import tempfile +import json +from pathlib import Path + +REPO_ROOT = Path(__file__).resolve().parent.parent +sys.path.insert(0, str(REPO_ROOT / "bin")) + +import box_work + +class TestBoxWork(unittest.TestCase): + def test_workers_topology(self): + worker_names = [w["name"] for w in box_work.WORKERS] + self.assertIn("opm", worker_names) + self.assertIn("646", worker_names) + self.assertIn("dev", worker_names) + self.assertIn("pip", worker_names) + self.assertIn("def", worker_names) + self.assertIn("muse", worker_names) + self.assertIn("muse-main", worker_names) + + def test_gitea_config_resolution(self): + api_base, token = box_work.get_gitea_config() + self.assertTrue(api_base.startswith("http")) + self.assertTrue(len(token) > 10) + + def test_color_helpers(self): + self.assertTrue(len(box_work.c_bold("test")) >= 4) + self.assertTrue(len(box_work.c_green("test")) >= 4) + self.assertTrue(len(box_work.c_red("test")) >= 4) + + def test_find_repo_root(self): + root = box_work.find_repo_root() + self.assertTrue(root.exists()) + + def test_claimed_tasks_empty_or_dict(self): + res = box_work.get_claimed_tasks() + self.assertIsInstance(res, dict) + + def test_recent_done_tasks_list(self): + res = box_work.get_recent_done_tasks(limit=5) + self.assertIsInstance(res, list) + +if __name__ == "__main__": + unittest.main()