From 6d4909e1be31bf90b422fad8b4b9d97bf3536aed Mon Sep 17 00:00:00 2001 From: operator Date: Fri, 9 Oct 2026 22:58:09 +0000 Subject: [PATCH] feat(cli): add box work command for unified worker signals and task orchestration --- bin/box-work.py | 537 +++++++++++++++++++++++++++++++++++++++++ bin/box_work.py | 1 + bin/super-cli.py | 380 +++++++++++++++++++++++++++-- tests/test_box_work.py | 47 ++++ 4 files changed, 945 insertions(+), 20 deletions(-) create mode 100755 bin/box-work.py create mode 120000 bin/box_work.py create mode 100644 tests/test_box_work.py 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/bin/super-cli.py b/bin/super-cli.py index 5ac08d1..2426f80 100755 --- a/bin/super-cli.py +++ b/bin/super-cli.py @@ -1174,7 +1174,6 @@ def cmd_runtime(args): box_state_file.parent.mkdir(parents=True, exist_ok=True) box_data = {} if box_state_file.exists(): - import json box_data = json.loads(box_state_file.read_text()) box_data[args.session] = {"socket": sock, "launched_at": datetime.now(timezone.utc).isoformat(), "origin": "box-cli"} box_state_file.write_text(json.dumps(box_data, indent=2)) @@ -1382,6 +1381,71 @@ def cmd_runtime(args): if not report["ok"]: sys.exit(1) + elif action == "kill": + import runtime_reconcile as rec + sock = getattr(args, "socket", None) or mcw.KNOWN_SOCKETS[0] + session = args.session + err = rec.kill_session(sock, session) + if as_json: + print(json.dumps({"ok": err is None, "socket": sock, + "session": session, + "error": err}, indent=2)) + return + if err is None: + print(c_green("\nāœ” Killed '%s' on %s\n") % (session, sock)) + return + print(c_red("Error: cannot kill '%s' on %s: %s" + % (session, sock, err)), file=sys.stderr) + sys.exit(1) + + elif action == "restart": + import runtime_reconcile as rec + sock = getattr(args, "socket", None) or mcw.KNOWN_SOCKETS[0] + session = args.session + manifest = (getattr(args, "manifest", None) + or str(NETVM_ROOT / "fleet" / "agents.json")) + dry_run = getattr(args, "dry_run", False) + res = rec.restart_agent(manifest, sock, session, dry_run=dry_run) + ok = res["action"] in ("restarted", "restart") + if as_json: + print(json.dumps({"ok": ok, "socket": sock, + "session": session, + "action": res["action"], + "detail": res["detail"]}, indent=2)) + return + if ok: + print(c_green("\nāœ” %s '%s': %s\n") + % ("Restarted" if res["action"] == "restarted" + else "Would restart", session, res["detail"])) + return + print(c_red("Error: cannot restart '%s': %s" + % (session, res["detail"])), file=sys.stderr) + sys.exit(1) + + elif action == "brief": + import runtime_reconcile as rec + sock = getattr(args, "socket", None) or mcw.KNOWN_SOCKETS[0] + session = args.session + manifest = (getattr(args, "manifest", None) + or str(NETVM_ROOT / "fleet" / "agents.json")) + dry_run = getattr(args, "dry_run", False) + res = rec.brief_agent(manifest, sock, session, dry_run=dry_run) + ok = res["action"] in ("briefed", "brief") + if as_json: + print(json.dumps({"ok": ok, "socket": sock, + "session": session, + "action": res["action"], + "detail": res["detail"]}, indent=2)) + return + if ok: + print(c_green("\nāœ” %s '%s': %s\n") + % ("Briefed" if res["action"] == "briefed" + else "Would brief", session, res["detail"])) + return + print(c_red("Error: cannot brief '%s': %s" + % (session, res["detail"])), file=sys.stderr) + sys.exit(1) + else: if as_json: print(json.dumps({"ok": False, @@ -1393,6 +1457,179 @@ def cmd_runtime(args): sys.exit(1) +# --------------------------------------------------------------------------- +# Domain: TASKS (agent work queue: pending/claimed/done) +# --------------------------------------------------------------------------- +def _tasks_age(age_s): + if age_s is None: + return "-" + if age_s < 90: + return "%ds" % int(age_s) + if age_s < 5400: + return "%dm" % int(age_s // 60) + if age_s < 172800: + return "%dh" % int(age_s // 3600) + return "%dd" % int(age_s // 86400) + + +def cmd_work(args): + import box_work + action = getattr(args, "work_action", None) + if not action or action == "status": + box_work.cmd_status(args) + elif action == "start": + box_work.cmd_start(args) + elif action == "assign": + box_work.cmd_assign(args) + elif action == "merge": + box_work.cmd_merge(args) + elif action == "chats": + box_work.cmd_chats(args) + else: + box_work.cmd_status(args) + + +def cmd_tasks(args): + import runtime_reconcile as rec + action = getattr(args, "tasks_action", None) or "list" + as_json = getattr(args, "json", False) + tasks_dir = (getattr(args, "dir", None) + or str(NETVM_ROOT / "fleet" / "tasks")) + + if action == "list": + queue = getattr(args, "queue", None) or "all" + rows = rec.list_tasks(tasks_dir, queue=queue) + if as_json: + print(json.dumps({"ok": True, "tasks": rows, + "dir": tasks_dir}, indent=2)) + return + print(c_bold("\n=== TASK QUEUE ===\n")) + if not rows: + print(c_dim(" No tasks in %s." % queue)) + print() + return + headers = ["QUEUE", "NAME", "OWNER", "AGE"] + table = [[r["queue"], r["name"][:44], + r["owner"] or badge_dim("-"), + _tasks_age(r["age_s"])] for r in rows] + print_table(headers, table) + print() + + elif action == "show": + name = args.name + res = rec.read_task(tasks_dir, name) + if as_json: + print(json.dumps({"ok": "error" not in res, + "dir": tasks_dir, **res}, indent=2)) + return + if "error" in res: + print(c_red("Error: %s" % res["error"]), file=sys.stderr) + sys.exit(1) + print(c_bold("\n=== TASK %s [%s] ===\n" % ( + res["name"], res["queue"]))) + print(res["text"].rstrip("\n")) + print() + + elif action == "create": + res = rec.create_task( + tasks_dir, args.name, getattr(args, "title", ""), + getattr(args, "goal", ""), getattr(args, "steps", "") or "", + dry_run=getattr(args, "dry_run", False)) + if as_json: + print(json.dumps({"ok": res["ok"], "dir": tasks_dir, + **{k: v for k, v in res.items() + if k != "ok"}}, indent=2)) + return + if not res["ok"]: + print(c_red("Error: %s" % res["error"]), file=sys.stderr) + sys.exit(1) + print(c_green("\nāœ” %s %s\n" % ( + "Would create" if res.get("dry_run") else "Created", + res["path"]))) + + elif action == "claim": + res = rec.claim_task(tasks_dir, args.name, args.owner, + dry_run=getattr(args, "dry_run", False)) + if as_json: + print(json.dumps({"ok": res["ok"], "dir": tasks_dir, + **{k: v for k, v in res.items() + if k != "ok"}}, indent=2)) + return + if not res["ok"]: + print(c_red("Error: %s" % res["error"]), file=sys.stderr) + sys.exit(1) + print(c_green("\nāœ” %s %s\n" % ( + "Would claim" if res.get("dry_run") else "Claimed", + res["path"]))) + + elif action == "done": + res = rec.complete_task(tasks_dir, args.name, + getattr(args, "result", "") or "", + dry_run=getattr(args, "dry_run", False)) + if as_json: + print(json.dumps({"ok": res["ok"], "dir": tasks_dir, + **{k: v for k, v in res.items() + if k != "ok"}}, indent=2)) + return + if not res["ok"]: + print(c_red("Error: %s" % res["error"]), file=sys.stderr) + sys.exit(1) + print(c_green("\nāœ” %s %s\n" % ( + "Would complete" if res.get("dry_run") else "Completed", + res["path"]))) + + elif action == "requeue": + res = rec.requeue_task(tasks_dir, args.name, + dry_run=getattr(args, "dry_run", False)) + if as_json: + print(json.dumps({"ok": res["ok"], "dir": tasks_dir, + **{k: v for k, v in res.items() + if k != "ok"}}, indent=2)) + return + if not res["ok"]: + print(c_red("Error: %s" % res["error"]), file=sys.stderr) + sys.exit(1) + print(c_green("\nāœ” %s %s\n" % ( + "Would requeue" if res.get("dry_run") else "Requeued", + res["path"]))) + + elif action == "sweep": + manifest = (getattr(args, "manifest", None) + or str(NETVM_ROOT / "fleet" / "agents.json")) + dry_run = getattr(args, "dry_run", False) + res = rec.sweep_now(manifest, tasks_dir=tasks_dir, + dry_run=dry_run) + if as_json: + print(json.dumps({"ok": res["ok"], "dir": tasks_dir, + "dry_run": dry_run, + "live_sessions": res["live_sessions"], + "requeued": res["requeued"], + "errors": res["errors"]}, indent=2)) + return + print(c_bold("\n=== TASK SWEEP%s ===\n" % ( + " (dry-run)" if dry_run else ""))) + if not res["requeued"] and not res["errors"]: + print(c_dim(" No stale claims.")) + for c in res["requeued"]: + print(" %s task %s (%s)" % ( + badge_ok("REQUEUED"), c_cyan(c["task"]), c["reason"])) + for e in res["errors"]: + print(" %s %s" % (badge_err("ERROR"), e)) + print() + if not res["ok"]: + sys.exit(1) + + else: + if as_json: + print(json.dumps({"ok": False, + "error": "unknown_action", + "action": action})) + return + print(c_red("Error: unknown tasks action '%s'" % action), + file=sys.stderr) + sys.exit(1) + + # --------------------------------------------------------------------------- # Domain: MUSE-CHOICES (Muse TUI A/B/C auto-answer daemon) # --------------------------------------------------------------------------- @@ -2563,24 +2800,49 @@ def cmd_dm_log(args): print(c_dim("dm-log.jsonl not found.")) return - entries = [] + def _match(data): + if filter_agent and (data.get("agent") != filter_agent and data.get("to") != filter_agent): + return False + if filter_text: + if filter_text.lower() not in json.dumps(data).lower(): + return False + return True + with open(DM_LOG, "r") as f: - for line in f: + lines = f.readlines() + + entries = [] + if isinstance(n, int) and n >= 1: + # Walk newest-first, parsing only until n matches: identical + # result to a full parse + [-n:] at O(n) instead of O(file). + for line in reversed(lines): line = line.strip() if not line: continue try: data = json.loads(line) - if filter_agent and (data.get("agent") != filter_agent and data.get("to") != filter_agent): - continue - if filter_text: - if filter_text.lower() not in json.dumps(data).lower(): - continue - entries.append(data) except Exception: continue - - entries = entries[-n:] + if not _match(data): + continue + entries.append(data) + if len(entries) >= n: + break + entries.reverse() + else: + # Legacy path: preserve entries[-n:] quirks for n <= 0. + for line in lines: + line = line.strip() + if not line: + continue + try: + data = json.loads(line) + except Exception: + continue + if not _match(data): + continue + entries.append(data) + entries = entries[-n:] if args.json: print(json.dumps({"ok": True, "entries": entries}, indent=2)) @@ -3677,7 +3939,7 @@ def cmd_web_test_auth(args): print(f" Response: {badge_ok(f'HTTP {e.code}')} (Protection active)") print(f" Body : {c_dim(body.strip()[:100])}\n") else: - print(c_warn(f"HTTP {e.code}: {body}")) + print(c_yellow(f"HTTP {e.code}: {body}")) except Exception as e: print(c_red(f"Error testing auth: {e}")) @@ -3771,7 +4033,7 @@ def cmd_cred_link_instagram(args): else: print("\n" + c_bold(f"=== INSTAGRAM LINKING: {args.node} ===") + "\n") if res.get("error"): - print(c_err(f" Error: {res['error']}")) + print(c_red(f" Error: {res['error']}")) sys.exit(1) print(f" Tailscale Portal: {c_cyan(res['portal_url'])}") print(f" Direct IP Portal: {c_cyan(res['portal_ip_url'])}") @@ -4120,7 +4382,7 @@ def _vm_followup_cancel(dm_id): if not sig: raise RuntimeError("empty signature") except Exception as e: - print(c_warn(" VM follow-up cancel skipped (signing failed: %s)" % e)) + print(c_yellow(" VM follow-up cancel skipped (signing failed: %s)" % e)) return query = _up.urlencode({"identity": "bl", "ts": ts_now, "sig": sig}) url = "%s/api/box/followups/cancel?%s" % (box_api, query) @@ -4140,10 +4402,10 @@ def _vm_followup_cancel(dm_id): detail = e.read().decode("utf-8", errors="ignore")[:120] except Exception: detail = "" - print(c_warn(" VM cancel failed: HTTP %s %s" % (e.code, detail))) + print(c_yellow(" VM cancel failed: HTTP %s %s" % (e.code, detail))) return except Exception as e: - print(c_warn(" VM cancel failed: %s" % e)) + print(c_yellow(" VM cancel failed: %s" % e)) return if resp.get("canceled"): print(c_green(" VM: follow-up canceled (%s)." @@ -4339,7 +4601,7 @@ def cmd_deploy(args): import muse_hybrid res, err = muse_hybrid.start_session(agent, title=title) if err or not res: - print(c_err(f"āœ– Failed to spawn subagent session: {err}"), file=sys.stderr) + print(c_red(f"āœ– Failed to spawn subagent session: {err}"), file=sys.stderr) sys.exit(1) session_id = res.get("session_id") @@ -4355,7 +4617,7 @@ def cmd_deploy(args): print(f" Dispatching task prompt to subagent (waiting up to {wait}s)...") send_res, send_err = muse_hybrid.send_message(agent, prompt, thread_id=session_id, wait=wait) if send_err: - print(c_warn(f"Notice: {send_err}")) + print(c_yellow(f"Notice: {send_err}")) else: print(c_green("āœ” Task prompt delivered.")) @@ -4373,7 +4635,7 @@ def cmd_deploy(args): # Pipeline deployment pipe_name = getattr(args, "name", None) if not pipe_name: - print(c_err("Error: Specify pipeline name (e.g. box deploy pipeline pipe-demo-step1)"), file=sys.stderr) + print(c_red("Error: Specify pipeline name (e.g. box deploy pipeline pipe-demo-step1)"), file=sys.stderr) sys.exit(1) args.name = pipe_name @@ -6300,7 +6562,7 @@ def build_parser(): p_mc_res.add_argument("decision", choices=["approve", "deny"], help="Release the hold to approve, or deny it (permission kinds only)") # Domain: RUNTIME - p_rt = subparsers.add_parser("runtime", parents=[common], help="Muse CLI tmux runtimes: list states, send input, launch with approval trail") + p_rt = subparsers.add_parser("runtime", parents=[common], help="Muse CLI tmux runtimes: list/send/launch/reconcile/kill/restart/brief") rt_sub = p_rt.add_subparsers(dest="rt_action") p_rt_list = rt_sub.add_parser("list", parents=[common], help="List panes with runtime state + approval posture (default)") @@ -6339,6 +6601,80 @@ def build_parser(): p_rt_rec.add_argument("--dry-run", action="store_true", help="Print the plan without changing anything") p_rt_rec.add_argument("--adopt", action="store_true", help="Record live sessions as briefed without sending") + p_rt_kill = rt_sub.add_parser("kill", parents=[common], help="Kill a session on a socket") + p_rt_kill.add_argument("--socket", default=None, help="Tmux socket (default: /tmp/tmux-1000/default)") + p_rt_kill.add_argument("--session", required=True, help="Session name to kill") + + p_rt_restart = rt_sub.add_parser("restart", parents=[common], help="Kill + relaunch + brief one manifest agent") + p_rt_restart.add_argument("--socket", default=None, help="Tmux socket (default: /tmp/tmux-1000/default)") + p_rt_restart.add_argument("--session", required=True, help="Manifest session name") + p_rt_restart.add_argument("--manifest", default=None, help="Manifest path (default: fleet/agents.json)") + p_rt_restart.add_argument("--dry-run", action="store_true", help="Print the plan without changing anything") + + p_rt_brief = rt_sub.add_parser("brief", parents=[common], help="Send the manifest brief to a live idle pane") + p_rt_brief.add_argument("--socket", default=None, help="Tmux socket (default: /tmp/tmux-1000/default)") + p_rt_brief.add_argument("--session", required=True, help="Manifest session name") + p_rt_brief.add_argument("--manifest", default=None, help="Manifest path (default: fleet/agents.json)") + p_rt_brief.add_argument("--dry-run", action="store_true", help="Print the plan without changing anything") + + p_work = subparsers.add_parser("work", parents=[common], help="Fleet workspace, task orchestration, worker scope, and active signals") + work_sub = p_work.add_subparsers(dest="work_action") + p_w_status = work_sub.add_parser("status", parents=[common], help="Show full operational work dashboard (default)") + p_w_start = work_sub.add_parser("start", parents=[common], help="Instantly start and assign new build ticket to an agent") + p_w_start.add_argument("title", help="Ticket title / summary") + p_w_start.add_argument("--to", dest="agent", required=True, help="Agent username (opm, 646, dev, pip, def, muse)") + p_w_start.add_argument("--goal", help="Optional detailed goal description") + p_w_assign = work_sub.add_parser("assign", parents=[common], help="Assign existing ticket to an agent") + p_w_assign.add_argument("issue", type=int, help="Issue number (e.g. 215)") + p_w_assign.add_argument("--to", dest="agent", required=True, help="Agent username") + p_w_merge = work_sub.add_parser("merge", parents=[common], help="Merge an open PR into master") + p_w_merge.add_argument("pr", type=int, help="Pull request number (e.g. 214)") + p_w_chats = work_sub.add_parser("chats", parents=[common], help="View recent live chat activity") + p_w_chats.add_argument("--agent", help="Filter by agent name") + p_w_chats.add_argument("--limit", type=int, default=10, help="Number of messages to show") + + p_tasks = subparsers.add_parser("tasks", parents=[common], help="Agent work queue: pending/claimed/done files (distinct from scheduled jobs)") + p_tasks.add_argument("--dir", default=None, help="Task queue dir (default: fleet/tasks)") + tasks_sub = p_tasks.add_subparsers(dest="tasks_action") + + p_t_list = tasks_sub.add_parser("list", parents=[common], help="List tasks across queues (default)") + p_t_list.add_argument("--dir", default=argparse.SUPPRESS, help="Task queue dir (default: fleet/tasks)") + p_t_list.add_argument("--queue", choices=["pending", "claimed", "done", "all"], default="all", help="Only this queue") + + p_t_show = tasks_sub.add_parser("show", parents=[common], help="Print one task file") + p_t_show.add_argument("--dir", default=argparse.SUPPRESS, help="Task queue dir (default: fleet/tasks)") + p_t_show.add_argument("name", help="Task name (base or owner-suffixed)") + + p_t_create = tasks_sub.add_parser("create", parents=[common], help="Write a new pending task from template") + p_t_create.add_argument("--dir", default=argparse.SUPPRESS, help="Task queue dir (default: fleet/tasks)") + p_t_create.add_argument("name", help="Task name like 012-slug.md") + p_t_create.add_argument("--title", required=True, help="Short title") + p_t_create.add_argument("--goal", required=True, help="Goal text") + p_t_create.add_argument("--steps", default="", help="Steps text") + p_t_create.add_argument("--dry-run", action="store_true", help="Print the plan without writing") + + p_t_claim = tasks_sub.add_parser("claim", parents=[common], help="Atomically claim a pending task") + p_t_claim.add_argument("--dir", default=argparse.SUPPRESS, help="Task queue dir (default: fleet/tasks)") + p_t_claim.add_argument("name", help="Pending task name") + p_t_claim.add_argument("--as", dest="owner", required=True, help="Owner tmux session name") + p_t_claim.add_argument("--dry-run", action="store_true", help="Print the plan without moving") + + p_t_done = tasks_sub.add_parser("done", parents=[common], help="Append notes and move a claim to done/") + p_t_done.add_argument("--dir", default=argparse.SUPPRESS, help="Task queue dir (default: fleet/tasks)") + p_t_done.add_argument("name", help="Claimed task name (base or suffixed)") + p_t_done.add_argument("--result", default="", help="Result notes to append") + p_t_done.add_argument("--dry-run", action="store_true", help="Print the plan without moving") + + p_t_req = tasks_sub.add_parser("requeue", parents=[common], help="Move a claim back to pending/") + p_t_req.add_argument("--dir", default=argparse.SUPPRESS, help="Task queue dir (default: fleet/tasks)") + p_t_req.add_argument("name", help="Claimed task name (base or suffixed)") + p_t_req.add_argument("--dry-run", action="store_true", help="Print the plan without moving") + + p_t_sweep = tasks_sub.add_parser("sweep", parents=[common], help="Requeue stale/dead-owner claims now") + p_t_sweep.add_argument("--dir", default=argparse.SUPPRESS, help="Task queue dir (default: fleet/tasks)") + p_t_sweep.add_argument("--manifest", default=None, help="Manifest path (default: fleet/agents.json)") + p_t_sweep.add_argument("--dry-run", action="store_true", help="Print the plan without moving") + # Domain: INVITE p_invite = subparsers.add_parser("invite", parents=[common], help="Muse.ai invite codes: find per-agent codes and redeem") p_invite.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node (status)") @@ -7286,6 +7622,10 @@ def main(): cmd_muse_choices(args) elif args.domain == "runtime": cmd_runtime(args) + elif args.domain == "work": + cmd_work(args) + elif args.domain == "tasks": + cmd_tasks(args) elif args.domain == "invite": cmd_invite(args) elif args.domain == "usage": 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()