4 Commits

6 changed files with 1465 additions and 20 deletions
+206
View File
@@ -0,0 +1,206 @@
#!/usr/bin/env python3
"""
box-gitea-bridge.py - Bridge Gitea webhooks to Box fleet tasks queue.
Listens for Gitea webhook events on 127.0.0.1:3005 and atomically converts
label-gated issues (labeled 'task' or 'ready') into fleet/tasks/pending/ files.
Also runs a periodic passive sweep to catch any dropped events (reaper backstop).
"""
import sys
import os
import re
import json
import time
import threading
import urllib.request
import urllib.parse
from http.server import HTTPServer, BaseHTTPRequestHandler
REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
TASKS_DIR = os.path.join(REPO_ROOT, "fleet", "tasks")
PARTITION_TABLE_PATH = os.path.join(REPO_ROOT, "fleet", "partition-table.json")
GITEA_API = "http://127.0.0.1:3000/api/v1"
def slugify(text: str) -> str:
text = text.lower()
text = re.sub(r"[^\w\s-]", "", text)
text = re.sub(r"[-\s]+", "-", text).strip("-")
return text[:45]
def get_admin_token() -> str:
if os.path.exists(PARTITION_TABLE_PATH):
try:
with open(PARTITION_TABLE_PATH) as f:
pt = json.load(f)
return pt.get("contributors", {}).get("super", {}).get("token", "")
except Exception:
pass
return "3c26744525bceaf385aa09737f7e41af613627b6"
def find_existing_task(issue_num: int):
prefix = f"{issue_num:03d}-"
for queue in ["pending", "claimed", "done"]:
qdir = os.path.join(TASKS_DIR, queue)
if not os.path.isdir(qdir):
continue
for fname in os.listdir(qdir):
if fname.startswith(prefix) or fname.startswith(f"{issue_num}-"):
return queue, os.path.join(qdir, fname)
return None, None
def create_task_from_issue(issue: dict):
issue_num = issue.get("number")
title = issue.get("title", "Untitled")
body = issue.get("body", "").strip() or "No goal description provided."
labels = [l.get("name", "") if isinstance(l, dict) else str(l) for l in issue.get("labels", [])]
assignee = issue.get("assignee")
assignee_name = assignee.get("username", "") if isinstance(assignee, dict) else ""
# Label-based direct routing: assign:<agent> or agent:<agent>
if not assignee_name:
for lbl in labels:
if lbl.startswith("assign:"):
assignee_name = lbl.split(":", 1)[1].strip()
break
elif lbl.startswith("agent:"):
assignee_name = lbl.split(":", 1)[1].strip()
break
# Label gate: must have 'task' or 'ready'
if not any(lbl in ["task", "ready"] for lbl in labels):
return None, "skipped_label_gate"
queue, existing_path = find_existing_task(issue_num)
if existing_path:
return existing_path, f"already_exists_in_{queue}"
slug = slugify(title)
fname = f"{issue_num:03d}-{slug}.md"
task_content = f"""# {issue_num:03d}-{slug}: {title}
Goal: {body}
Steps:
1. Claim task on feature branch builder/{slug}.
2. Implement solution adhering to test coverage.
3. Commit with "Fixes #{issue_num}" and push to master/PR.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
"""
os.makedirs(os.path.join(TASKS_DIR, "pending"), exist_ok=True)
os.makedirs(os.path.join(TASKS_DIR, "claimed"), exist_ok=True)
if assignee_name:
target_path = os.path.join(TASKS_DIR, "claimed", f"{fname}.{assignee_name}")
else:
target_path = os.path.join(TASKS_DIR, "pending", fname)
tmp_path = target_path + ".tmp"
with open(tmp_path, "w") as f:
f.write(task_content)
os.replace(tmp_path, target_path)
return target_path, "created"
def close_task_for_issue(issue_num: int, close_notes="Closed via Gitea"):
queue, task_path = find_existing_task(issue_num)
if not task_path or queue == "done":
return None
fname = os.path.basename(task_path)
done_dir = os.path.join(TASKS_DIR, "done")
os.makedirs(done_dir, exist_ok=True)
# Append close notes
with open(task_path, "a") as f:
f.write(f"\n{time.strftime('%Y-%m-%d %H:%M:%SZ')}: {close_notes}\n")
done_path = os.path.join(done_dir, fname)
os.replace(task_path, done_path)
return done_path
def passive_reconcile_sweep():
token = get_admin_token()
url = f"{GITEA_API}/repos/super/box/issues?state=open"
req = urllib.request.Request(url)
req.add_header("Authorization", f"token {token}")
try:
with urllib.request.urlopen(req, timeout=5) as resp:
issues = json.loads(resp.read().decode("utf-8"))
for issue in issues:
create_task_from_issue(issue)
except Exception as e:
sys.stderr.write(f"[sweep] warning: passive reconcile error: {e}\n")
class WebhookHandler(BaseHTTPRequestHandler):
def do_POST(self):
content_length = int(self.headers.get("Content-Length", 0))
body = self.rfile.read(content_length).decode("utf-8")
event = self.headers.get("X-Gitea-Event", "")
try:
payload = json.loads(body)
except Exception:
self.send_response(400)
self.end_headers()
self.wfile.write(b'{"error": "invalid json"}')
return
response_data = {"status": "ignored"}
if event == "issues":
action = payload.get("action", "")
issue = payload.get("issue", {})
issue_num = issue.get("number")
if action in ["opened", "labeled", "assigned"]:
target, outcome = create_task_from_issue(issue)
response_data = {"status": "ok", "action": action, "target": target, "outcome": outcome}
elif action == "closed":
done_path = close_task_for_issue(issue_num, f"Closed via Gitea issue #{issue_num}")
response_data = {"status": "ok", "action": "closed", "done_path": done_path}
self.send_response(200)
self.send_header("Content-Type", "application/json")
self.end_headers()
self.wfile.write(json.dumps(response_data).encode("utf-8"))
def do_GET(self):
if self.path == "/health":
self.send_response(200)
self.send_header("Content-Type", "application/json")
self.end_headers()
self.wfile.write(b'{"status": "ok", "service": "box-gitea-bridge"}')
elif self.path == "/sweep":
passive_reconcile_sweep()
self.send_response(200)
self.send_header("Content-Type", "application/json")
self.end_headers()
self.wfile.write(b'{"status": "swept"}')
else:
self.send_response(404)
self.end_headers()
def background_sweeper_loop(interval=60):
while True:
time.sleep(interval)
try:
passive_reconcile_sweep()
except Exception:
pass
def main():
port = int(os.environ.get("BRIDGE_PORT", 3005))
server = HTTPServer(("127.0.0.1", port), WebhookHandler)
t = threading.Thread(target=background_sweeper_loop, daemon=True)
t.start()
print(f"box-gitea-bridge listening on 127.0.0.1:{port} (reconciler running every 60s)")
try:
server.serve_forever()
except KeyboardInterrupt:
pass
if __name__ == "__main__":
main()
+736
View File
@@ -0,0 +1,736 @@
#!/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 <title> --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}"
# 0. Pre-flight health gate: Hatch, Restore, Git Config
preflight = check_agent_preflight(agent)
if not preflight["ready"] and not getattr(args, "force", False):
print(c_red(f"\n[BLOCKED] Agent '{agent}' failed pre-flight health verification:"))
print(f" • Hatch: {badge_status(preflight['hatch']['status'])} - {preflight['hatch']['details']}")
print(f" • Restore: {badge_status(preflight['restore']['status'])} - {preflight['restore']['details']}")
print(f" • Git: {badge_status(preflight['git']['status'])} - {preflight['git']['details']}")
print(c_yellow("\nBlocking reasons:"))
for r in preflight["reasons"]:
print(f" - {r}")
print(c_dim(f"\nTo inspect full health: box work check {agent}\nTo bypass pre-flight: box work start '{title}' --to {agent} --force\n"))
sys.exit(1)
elif not preflight["ready"] and getattr(args, "force", False):
print(c_yellow(f"[WARNING] Overriding failed pre-flight checks on {agent} (--force specified).\n"))
else:
print(c_green(f"✓ Pre-flight checks passed (Hatch: OK, Restore: OK, Git Config: OK) for {agent}"))
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
# 0. Pre-flight health gate: Hatch, Restore, Git Config
preflight = check_agent_preflight(agent)
if not preflight["ready"] and not getattr(args, "force", False):
print(c_red(f"\n[BLOCKED] Agent '{agent}' failed pre-flight health verification:"))
print(f" • Hatch: {badge_status(preflight['hatch']['status'])} - {preflight['hatch']['details']}")
print(f" • Restore: {badge_status(preflight['restore']['status'])} - {preflight['restore']['details']}")
print(f" • Git: {badge_status(preflight['git']['status'])} - {preflight['git']['details']}")
print(c_yellow("\nBlocking reasons:"))
for r in preflight["reasons"]:
print(f" - {r}")
print(c_dim(f"\nTo bypass pre-flight: box work assign {issue_num} --to {agent} --force\n"))
sys.exit(1)
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")
# box work check [agent]
p_check = sub.add_parser("check", help="Run pre-flight health verification (Hatch, Restore, Git Config)")
p_check.add_argument("agent", nargs="?", help="Optional specific agent name to check")
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_start.add_argument("--force", action="store_true", help="Bypass pre-flight health gate")
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_assign.add_argument("--force", action="store_true", help="Bypass pre-flight health gate")
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 == "check":
cmd_check(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()
+1
View File
@@ -0,0 +1 @@
/home/super/Projects/NetVM/bin/box-work.py
+366 -20
View File
@@ -1174,7 +1174,6 @@ def cmd_runtime(args):
box_state_file.parent.mkdir(parents=True, exist_ok=True) box_state_file.parent.mkdir(parents=True, exist_ok=True)
box_data = {} box_data = {}
if box_state_file.exists(): if box_state_file.exists():
import json
box_data = json.loads(box_state_file.read_text()) 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_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)) box_state_file.write_text(json.dumps(box_data, indent=2))
@@ -1382,6 +1381,71 @@ def cmd_runtime(args):
if not report["ok"]: if not report["ok"]:
sys.exit(1) 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: else:
if as_json: if as_json:
print(json.dumps({"ok": False, print(json.dumps({"ok": False,
@@ -1393,6 +1457,181 @@ def cmd_runtime(args):
sys.exit(1) 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 == "check":
box_work.cmd_check(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) # Domain: MUSE-CHOICES (Muse TUI A/B/C auto-answer daemon)
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
@@ -2563,24 +2802,49 @@ def cmd_dm_log(args):
print(c_dim("dm-log.jsonl not found.")) print(c_dim("dm-log.jsonl not found."))
return 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: 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() line = line.strip()
if not line: if not line:
continue continue
try: try:
data = json.loads(line) 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: except Exception:
continue continue
if not _match(data):
entries = entries[-n:] 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: if args.json:
print(json.dumps({"ok": True, "entries": entries}, indent=2)) print(json.dumps({"ok": True, "entries": entries}, indent=2))
@@ -3677,7 +3941,7 @@ def cmd_web_test_auth(args):
print(f" Response: {badge_ok(f'HTTP {e.code}')} (Protection active)") print(f" Response: {badge_ok(f'HTTP {e.code}')} (Protection active)")
print(f" Body : {c_dim(body.strip()[:100])}\n") print(f" Body : {c_dim(body.strip()[:100])}\n")
else: else:
print(c_warn(f"HTTP {e.code}: {body}")) print(c_yellow(f"HTTP {e.code}: {body}"))
except Exception as e: except Exception as e:
print(c_red(f"Error testing auth: {e}")) print(c_red(f"Error testing auth: {e}"))
@@ -3771,7 +4035,7 @@ def cmd_cred_link_instagram(args):
else: else:
print("\n" + c_bold(f"=== INSTAGRAM LINKING: {args.node} ===") + "\n") print("\n" + c_bold(f"=== INSTAGRAM LINKING: {args.node} ===") + "\n")
if res.get("error"): if res.get("error"):
print(c_err(f" Error: {res['error']}")) print(c_red(f" Error: {res['error']}"))
sys.exit(1) sys.exit(1)
print(f" Tailscale Portal: {c_cyan(res['portal_url'])}") print(f" Tailscale Portal: {c_cyan(res['portal_url'])}")
print(f" Direct IP Portal: {c_cyan(res['portal_ip_url'])}") print(f" Direct IP Portal: {c_cyan(res['portal_ip_url'])}")
@@ -4120,7 +4384,7 @@ def _vm_followup_cancel(dm_id):
if not sig: if not sig:
raise RuntimeError("empty signature") raise RuntimeError("empty signature")
except Exception as e: 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 return
query = _up.urlencode({"identity": "bl", "ts": ts_now, "sig": sig}) query = _up.urlencode({"identity": "bl", "ts": ts_now, "sig": sig})
url = "%s/api/box/followups/cancel?%s" % (box_api, query) url = "%s/api/box/followups/cancel?%s" % (box_api, query)
@@ -4140,10 +4404,10 @@ def _vm_followup_cancel(dm_id):
detail = e.read().decode("utf-8", errors="ignore")[:120] detail = e.read().decode("utf-8", errors="ignore")[:120]
except Exception: except Exception:
detail = "" 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 return
except Exception as e: except Exception as e:
print(c_warn(" VM cancel failed: %s" % e)) print(c_yellow(" VM cancel failed: %s" % e))
return return
if resp.get("canceled"): if resp.get("canceled"):
print(c_green(" VM: follow-up canceled (%s)." print(c_green(" VM: follow-up canceled (%s)."
@@ -4339,7 +4603,7 @@ def cmd_deploy(args):
import muse_hybrid import muse_hybrid
res, err = muse_hybrid.start_session(agent, title=title) res, err = muse_hybrid.start_session(agent, title=title)
if err or not res: 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) sys.exit(1)
session_id = res.get("session_id") session_id = res.get("session_id")
@@ -4355,7 +4619,7 @@ def cmd_deploy(args):
print(f" Dispatching task prompt to subagent (waiting up to {wait}s)...") 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) send_res, send_err = muse_hybrid.send_message(agent, prompt, thread_id=session_id, wait=wait)
if send_err: if send_err:
print(c_warn(f"Notice: {send_err}")) print(c_yellow(f"Notice: {send_err}"))
else: else:
print(c_green("✔ Task prompt delivered.")) print(c_green("✔ Task prompt delivered."))
@@ -4373,7 +4637,7 @@ def cmd_deploy(args):
# Pipeline deployment # Pipeline deployment
pipe_name = getattr(args, "name", None) pipe_name = getattr(args, "name", None)
if not pipe_name: 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) sys.exit(1)
args.name = pipe_name args.name = pipe_name
@@ -6300,7 +6564,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)") p_mc_res.add_argument("decision", choices=["approve", "deny"], help="Release the hold to approve, or deny it (permission kinds only)")
# Domain: RUNTIME # 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") 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)") p_rt_list = rt_sub.add_parser("list", parents=[common], help="List panes with runtime state + approval posture (default)")
@@ -6339,6 +6603,84 @@ 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("--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_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_check = work_sub.add_parser("check", parents=[common], help="Run pre-flight health checks (Hatch, Restore, Git Config)")
p_w_check.add_argument("agent", nargs="?", help="Optional specific agent name to check")
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_start.add_argument("--force", action="store_true", help="Bypass pre-flight health gate")
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_assign.add_argument("--force", action="store_true", help="Bypass pre-flight health gate")
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 # Domain: INVITE
p_invite = subparsers.add_parser("invite", parents=[common], help="Muse.ai invite codes: find per-agent codes and redeem") 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)") p_invite.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node (status)")
@@ -7286,6 +7628,10 @@ def main():
cmd_muse_choices(args) cmd_muse_choices(args)
elif args.domain == "runtime": elif args.domain == "runtime":
cmd_runtime(args) cmd_runtime(args)
elif args.domain == "work":
cmd_work(args)
elif args.domain == "tasks":
cmd_tasks(args)
elif args.domain == "invite": elif args.domain == "invite":
cmd_invite(args) cmd_invite(args)
elif args.domain == "usage": elif args.domain == "usage":
+94
View File
@@ -0,0 +1,94 @@
import os
import sys
import tempfile
import unittest
import json
REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
sys.path.insert(0, os.path.join(REPO_ROOT, "bin"))
import importlib.util
spec = importlib.util.spec_from_file_location("box_gitea_bridge", os.path.join(REPO_ROOT, "bin", "box-gitea-bridge.py"))
bgb = importlib.util.module_from_spec(spec)
spec.loader.exec_module(bgb)
class TestBoxGiteaBridge(unittest.TestCase):
def setUp(self):
self.tmpdir = tempfile.TemporaryDirectory()
bgb.TASKS_DIR = os.path.join(self.tmpdir.name, "fleet", "tasks")
os.makedirs(os.path.join(bgb.TASKS_DIR, "pending"), exist_ok=True)
os.makedirs(os.path.join(bgb.TASKS_DIR, "claimed"), exist_ok=True)
os.makedirs(os.path.join(bgb.TASKS_DIR, "done"), exist_ok=True)
def tearDown(self):
self.tmpdir.cleanup()
def test_slugify(self):
self.assertEqual(bgb.slugify("Hello World! 123"), "hello-world-123")
self.assertEqual(bgb.slugify("Fix: Gitea & Box Integration"), "fix-gitea-box-integration")
def test_label_gating(self):
# Unlabeled issue -> skipped
issue_unlabeled = {"number": 208, "title": "Untagged discussion", "labels": []}
path, status = bgb.create_task_from_issue(issue_unlabeled)
self.assertIsNone(path)
self.assertEqual(status, "skipped_label_gate")
# Irrelevant label -> skipped
issue_wontfix = {"number": 208, "title": "Wontfix bug", "labels": [{"name": "wontfix"}]}
path, status = bgb.create_task_from_issue(issue_wontfix)
self.assertIsNone(path)
self.assertEqual(status, "skipped_label_gate")
# Labeled 'task' -> created
issue_task = {"number": 208, "title": "Build Gitea Bridge", "labels": [{"name": "task"}], "body": "Implement bridge"}
path, status = bgb.create_task_from_issue(issue_task)
self.assertIsNotNone(path)
self.assertEqual(status, "created")
self.assertTrue(os.path.exists(path))
self.assertIn("208-build-gitea-bridge.md", path)
def test_assigned_issue_claims_directly(self):
issue_assigned = {
"number": 209,
"title": "OPM Recovery Task",
"labels": [{"name": "ready"}],
"body": "Run recovery",
"assignee": {"username": "opm"}
}
path, status = bgb.create_task_from_issue(issue_assigned)
self.assertIsNotNone(path)
self.assertEqual(status, "created")
self.assertIn("claimed", path)
self.assertTrue(path.endswith(".opm"))
def test_label_based_routing(self):
issue_label_assigned = {
"number": 211,
"title": "Direct Labeled Task",
"labels": [{"name": "task"}, {"name": "assign:opm"}],
"body": "Direct routing via label"
}
path, status = bgb.create_task_from_issue(issue_label_assigned)
self.assertIsNotNone(path)
self.assertEqual(status, "created")
self.assertIn("claimed", path)
self.assertTrue(path.endswith(".opm"))
def test_close_task_moves_to_done(self):
issue = {"number": 210, "title": "Close test", "labels": [{"name": "ready"}]}
created_path, _ = bgb.create_task_from_issue(issue)
self.assertTrue(os.path.exists(created_path))
done_path = bgb.close_task_for_issue(210, "Verified fixed")
self.assertIsNotNone(done_path)
self.assertFalse(os.path.exists(created_path))
self.assertTrue(os.path.exists(done_path))
self.assertIn("done", done_path)
with open(done_path) as f:
content = f.read()
self.assertIn("Verified fixed", content)
if __name__ == "__main__":
unittest.main()
+62
View File
@@ -0,0 +1,62 @@
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)
def test_check_agent_preflight_structure(self):
res = box_work.check_agent_preflight("opm")
self.assertIn("hatch", res)
self.assertIn("restore", res)
self.assertIn("git", res)
self.assertIn("status", res["hatch"])
self.assertIn("status", res["restore"])
self.assertIn("status", res["git"])
self.assertIn("ready", res)
def test_check_agent_preflight_unknown_fails(self):
res = box_work.check_agent_preflight("nonexistent_agent_xyz")
self.assertFalse(res["ready"])
self.assertEqual(res["overall"], "FAIL")
if __name__ == "__main__":
unittest.main()