Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| a6e5f565f5 | |||
| 6d4909e1be | |||
| 698798df2a |
Executable
+206
@@ -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()
|
||||||
Executable
+736
@@ -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()
|
||||||
Symlink
+1
@@ -0,0 +1 @@
|
|||||||
|
/home/super/Projects/NetVM/bin/box-work.py
|
||||||
+366
-20
@@ -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":
|
||||||
|
|||||||
@@ -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()
|
||||||
@@ -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()
|
||||||
Reference in New Issue
Block a user