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
|
||||
+365
-19
@@ -1174,7 +1174,6 @@ def cmd_runtime(args):
|
||||
box_state_file.parent.mkdir(parents=True, exist_ok=True)
|
||||
box_data = {}
|
||||
if box_state_file.exists():
|
||||
import json
|
||||
box_data = json.loads(box_state_file.read_text())
|
||||
box_data[args.session] = {"socket": sock, "launched_at": datetime.now(timezone.utc).isoformat(), "origin": "box-cli"}
|
||||
box_state_file.write_text(json.dumps(box_data, indent=2))
|
||||
@@ -1382,6 +1381,71 @@ def cmd_runtime(args):
|
||||
if not report["ok"]:
|
||||
sys.exit(1)
|
||||
|
||||
elif action == "kill":
|
||||
import runtime_reconcile as rec
|
||||
sock = getattr(args, "socket", None) or mcw.KNOWN_SOCKETS[0]
|
||||
session = args.session
|
||||
err = rec.kill_session(sock, session)
|
||||
if as_json:
|
||||
print(json.dumps({"ok": err is None, "socket": sock,
|
||||
"session": session,
|
||||
"error": err}, indent=2))
|
||||
return
|
||||
if err is None:
|
||||
print(c_green("\n✔ Killed '%s' on %s\n") % (session, sock))
|
||||
return
|
||||
print(c_red("Error: cannot kill '%s' on %s: %s"
|
||||
% (session, sock, err)), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
elif action == "restart":
|
||||
import runtime_reconcile as rec
|
||||
sock = getattr(args, "socket", None) or mcw.KNOWN_SOCKETS[0]
|
||||
session = args.session
|
||||
manifest = (getattr(args, "manifest", None)
|
||||
or str(NETVM_ROOT / "fleet" / "agents.json"))
|
||||
dry_run = getattr(args, "dry_run", False)
|
||||
res = rec.restart_agent(manifest, sock, session, dry_run=dry_run)
|
||||
ok = res["action"] in ("restarted", "restart")
|
||||
if as_json:
|
||||
print(json.dumps({"ok": ok, "socket": sock,
|
||||
"session": session,
|
||||
"action": res["action"],
|
||||
"detail": res["detail"]}, indent=2))
|
||||
return
|
||||
if ok:
|
||||
print(c_green("\n✔ %s '%s': %s\n")
|
||||
% ("Restarted" if res["action"] == "restarted"
|
||||
else "Would restart", session, res["detail"]))
|
||||
return
|
||||
print(c_red("Error: cannot restart '%s': %s"
|
||||
% (session, res["detail"])), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
elif action == "brief":
|
||||
import runtime_reconcile as rec
|
||||
sock = getattr(args, "socket", None) or mcw.KNOWN_SOCKETS[0]
|
||||
session = args.session
|
||||
manifest = (getattr(args, "manifest", None)
|
||||
or str(NETVM_ROOT / "fleet" / "agents.json"))
|
||||
dry_run = getattr(args, "dry_run", False)
|
||||
res = rec.brief_agent(manifest, sock, session, dry_run=dry_run)
|
||||
ok = res["action"] in ("briefed", "brief")
|
||||
if as_json:
|
||||
print(json.dumps({"ok": ok, "socket": sock,
|
||||
"session": session,
|
||||
"action": res["action"],
|
||||
"detail": res["detail"]}, indent=2))
|
||||
return
|
||||
if ok:
|
||||
print(c_green("\n✔ %s '%s': %s\n")
|
||||
% ("Briefed" if res["action"] == "briefed"
|
||||
else "Would brief", session, res["detail"]))
|
||||
return
|
||||
print(c_red("Error: cannot brief '%s': %s"
|
||||
% (session, res["detail"])), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
else:
|
||||
if as_json:
|
||||
print(json.dumps({"ok": False,
|
||||
@@ -1393,6 +1457,181 @@ def cmd_runtime(args):
|
||||
sys.exit(1)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Domain: TASKS (agent work queue: pending/claimed/done)
|
||||
# ---------------------------------------------------------------------------
|
||||
def _tasks_age(age_s):
|
||||
if age_s is None:
|
||||
return "-"
|
||||
if age_s < 90:
|
||||
return "%ds" % int(age_s)
|
||||
if age_s < 5400:
|
||||
return "%dm" % int(age_s // 60)
|
||||
if age_s < 172800:
|
||||
return "%dh" % int(age_s // 3600)
|
||||
return "%dd" % int(age_s // 86400)
|
||||
|
||||
|
||||
def cmd_work(args):
|
||||
import box_work
|
||||
action = getattr(args, "work_action", None)
|
||||
if not action or action == "status":
|
||||
box_work.cmd_status(args)
|
||||
elif action == "start":
|
||||
box_work.cmd_start(args)
|
||||
elif action == "assign":
|
||||
box_work.cmd_assign(args)
|
||||
elif action == "merge":
|
||||
box_work.cmd_merge(args)
|
||||
elif action == "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)
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -2563,23 +2802,48 @@ def cmd_dm_log(args):
|
||||
print(c_dim("dm-log.jsonl not found."))
|
||||
return
|
||||
|
||||
entries = []
|
||||
def _match(data):
|
||||
if filter_agent and (data.get("agent") != filter_agent and data.get("to") != filter_agent):
|
||||
return False
|
||||
if filter_text:
|
||||
if filter_text.lower() not in json.dumps(data).lower():
|
||||
return False
|
||||
return True
|
||||
|
||||
with open(DM_LOG, "r") as f:
|
||||
for line in f:
|
||||
lines = f.readlines()
|
||||
|
||||
entries = []
|
||||
if isinstance(n, int) and n >= 1:
|
||||
# Walk newest-first, parsing only until n matches: identical
|
||||
# result to a full parse + [-n:] at O(n) instead of O(file).
|
||||
for line in reversed(lines):
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
data = json.loads(line)
|
||||
if filter_agent and (data.get("agent") != filter_agent and data.get("to") != filter_agent):
|
||||
continue
|
||||
if filter_text:
|
||||
if filter_text.lower() not in json.dumps(data).lower():
|
||||
continue
|
||||
entries.append(data)
|
||||
except Exception:
|
||||
continue
|
||||
|
||||
if not _match(data):
|
||||
continue
|
||||
entries.append(data)
|
||||
if len(entries) >= n:
|
||||
break
|
||||
entries.reverse()
|
||||
else:
|
||||
# Legacy path: preserve entries[-n:] quirks for n <= 0.
|
||||
for line in lines:
|
||||
line = line.strip()
|
||||
if not line:
|
||||
continue
|
||||
try:
|
||||
data = json.loads(line)
|
||||
except Exception:
|
||||
continue
|
||||
if not _match(data):
|
||||
continue
|
||||
entries.append(data)
|
||||
entries = entries[-n:]
|
||||
|
||||
if args.json:
|
||||
@@ -3677,7 +3941,7 @@ def cmd_web_test_auth(args):
|
||||
print(f" Response: {badge_ok(f'HTTP {e.code}')} (Protection active)")
|
||||
print(f" Body : {c_dim(body.strip()[:100])}\n")
|
||||
else:
|
||||
print(c_warn(f"HTTP {e.code}: {body}"))
|
||||
print(c_yellow(f"HTTP {e.code}: {body}"))
|
||||
except Exception as e:
|
||||
print(c_red(f"Error testing auth: {e}"))
|
||||
|
||||
@@ -3771,7 +4035,7 @@ def cmd_cred_link_instagram(args):
|
||||
else:
|
||||
print("\n" + c_bold(f"=== INSTAGRAM LINKING: {args.node} ===") + "\n")
|
||||
if res.get("error"):
|
||||
print(c_err(f" Error: {res['error']}"))
|
||||
print(c_red(f" Error: {res['error']}"))
|
||||
sys.exit(1)
|
||||
print(f" Tailscale Portal: {c_cyan(res['portal_url'])}")
|
||||
print(f" Direct IP Portal: {c_cyan(res['portal_ip_url'])}")
|
||||
@@ -4120,7 +4384,7 @@ def _vm_followup_cancel(dm_id):
|
||||
if not sig:
|
||||
raise RuntimeError("empty signature")
|
||||
except Exception as e:
|
||||
print(c_warn(" VM follow-up cancel skipped (signing failed: %s)" % e))
|
||||
print(c_yellow(" VM follow-up cancel skipped (signing failed: %s)" % e))
|
||||
return
|
||||
query = _up.urlencode({"identity": "bl", "ts": ts_now, "sig": sig})
|
||||
url = "%s/api/box/followups/cancel?%s" % (box_api, query)
|
||||
@@ -4140,10 +4404,10 @@ def _vm_followup_cancel(dm_id):
|
||||
detail = e.read().decode("utf-8", errors="ignore")[:120]
|
||||
except Exception:
|
||||
detail = ""
|
||||
print(c_warn(" VM cancel failed: HTTP %s %s" % (e.code, detail)))
|
||||
print(c_yellow(" VM cancel failed: HTTP %s %s" % (e.code, detail)))
|
||||
return
|
||||
except Exception as e:
|
||||
print(c_warn(" VM cancel failed: %s" % e))
|
||||
print(c_yellow(" VM cancel failed: %s" % e))
|
||||
return
|
||||
if resp.get("canceled"):
|
||||
print(c_green(" VM: follow-up canceled (%s)."
|
||||
@@ -4339,7 +4603,7 @@ def cmd_deploy(args):
|
||||
import muse_hybrid
|
||||
res, err = muse_hybrid.start_session(agent, title=title)
|
||||
if err or not res:
|
||||
print(c_err(f"✖ Failed to spawn subagent session: {err}"), file=sys.stderr)
|
||||
print(c_red(f"✖ Failed to spawn subagent session: {err}"), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
session_id = res.get("session_id")
|
||||
@@ -4355,7 +4619,7 @@ def cmd_deploy(args):
|
||||
print(f" Dispatching task prompt to subagent (waiting up to {wait}s)...")
|
||||
send_res, send_err = muse_hybrid.send_message(agent, prompt, thread_id=session_id, wait=wait)
|
||||
if send_err:
|
||||
print(c_warn(f"Notice: {send_err}"))
|
||||
print(c_yellow(f"Notice: {send_err}"))
|
||||
else:
|
||||
print(c_green("✔ Task prompt delivered."))
|
||||
|
||||
@@ -4373,7 +4637,7 @@ def cmd_deploy(args):
|
||||
# Pipeline deployment
|
||||
pipe_name = getattr(args, "name", None)
|
||||
if not pipe_name:
|
||||
print(c_err("Error: Specify pipeline name (e.g. box deploy pipeline pipe-demo-step1)"), file=sys.stderr)
|
||||
print(c_red("Error: Specify pipeline name (e.g. box deploy pipeline pipe-demo-step1)"), file=sys.stderr)
|
||||
sys.exit(1)
|
||||
|
||||
args.name = pipe_name
|
||||
@@ -6300,7 +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)")
|
||||
|
||||
# Domain: RUNTIME
|
||||
p_rt = subparsers.add_parser("runtime", parents=[common], help="Muse CLI tmux runtimes: list states, send input, launch with approval trail")
|
||||
p_rt = subparsers.add_parser("runtime", parents=[common], help="Muse CLI tmux runtimes: list/send/launch/reconcile/kill/restart/brief")
|
||||
rt_sub = p_rt.add_subparsers(dest="rt_action")
|
||||
|
||||
p_rt_list = rt_sub.add_parser("list", parents=[common], help="List panes with runtime state + approval posture (default)")
|
||||
@@ -6339,6 +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("--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
|
||||
p_invite = subparsers.add_parser("invite", parents=[common], help="Muse.ai invite codes: find per-agent codes and redeem")
|
||||
p_invite.add_argument("--node", choices=VALID_NODES, default=None, help="Filter by node (status)")
|
||||
@@ -7286,6 +7628,10 @@ def main():
|
||||
cmd_muse_choices(args)
|
||||
elif args.domain == "runtime":
|
||||
cmd_runtime(args)
|
||||
elif args.domain == "work":
|
||||
cmd_work(args)
|
||||
elif args.domain == "tasks":
|
||||
cmd_tasks(args)
|
||||
elif args.domain == "invite":
|
||||
cmd_invite(args)
|
||||
elif args.domain == "usage":
|
||||
|
||||
@@ -1,60 +0,0 @@
|
||||
#!/bin/bash
|
||||
# verify-node-ssh.sh — verify container SSH dial-in readiness across fleet nodes.
|
||||
# Checks from the VM: reverse-tunnel listeners + SSH auth for each node port.
|
||||
#
|
||||
# Port map (docs/OPERATOR-DRIVE-RUNBOOK.md):
|
||||
# muse-main 2224 | muse 2225 | 646 2226 | pip 2227 | opm 2228 | def 2229 | dev 2230
|
||||
#
|
||||
# What it checks per node:
|
||||
# 1. Reverse-tunnel listener on 127.0.0.1:<port> (dark node = no listener)
|
||||
# 2. SSH dial-in with BatchMode (auth failure = authorized_keys perms/key issue)
|
||||
#
|
||||
# Common root causes (see #211):
|
||||
# - sshd requires non-group-writable authorized_keys (must be 600)
|
||||
# - stale /run/nologin blocks logins
|
||||
# - missing id_frontdoor keys on dark nodes
|
||||
#
|
||||
# Usage: run on the VM (super@34.139.37.135), or via:
|
||||
# ssh-vm.sh "bash -s" < verify-node-ssh.sh
|
||||
set -u
|
||||
|
||||
# node:port pairs to check
|
||||
NODES="muse:2225 646:2226 pip:2227 def:2229 dev:2230 muse-main:2224 opm:2228"
|
||||
|
||||
fail=0
|
||||
for pair in $NODES; do
|
||||
node="${pair%%:*}"
|
||||
port="${pair##*:}"
|
||||
|
||||
# 1. listener check
|
||||
if ss -tln 2>/dev/null | grep -q "127.0.0.1:${port} "; then
|
||||
listener="LISTEN"
|
||||
else
|
||||
listener="DARK (no listener)"
|
||||
fi
|
||||
|
||||
# 2. auth check (only if listening)
|
||||
if [ "$listener" = "LISTEN" ]; then
|
||||
out=$(timeout 15 ssh -o StrictHostKeyChecking=no -o BatchMode=yes \
|
||||
-o ConnectTimeout=10 -p "$port" hatch@127.0.0.1 'echo OK' 2>&1)
|
||||
case "$out" in
|
||||
OK) auth="OK" ;;
|
||||
*"Permission denied"*) auth="AUTH-FAIL (check authorized_keys perms/keys)" ;;
|
||||
*"Connection refused"*) auth="REFUSED (tunnel died after listen check)" ;;
|
||||
*) auth="OTHER: $(echo "$out" | head -1 | cut -c1-60)" ;;
|
||||
esac
|
||||
else
|
||||
auth="SKIP"
|
||||
fi
|
||||
|
||||
printf '%-10s port %-5s listener: %-22s auth: %s\n' "$node" "$port" "$listener" "$auth"
|
||||
[ "$listener" = "DARK (no listener)" ] && fail=1
|
||||
case "$auth" in AUTH-FAIL*) fail=1 ;; esac
|
||||
done
|
||||
|
||||
if [ "$fail" -eq 0 ]; then
|
||||
echo "ALL NODES REACHABLE"
|
||||
else
|
||||
echo "ISSUES FOUND (see above)"
|
||||
fi
|
||||
exit "$fail"
|
||||
@@ -1,38 +0,0 @@
|
||||
# Ticket #213 verification — SSH key perms and container dial-in (646)
|
||||
|
||||
Date: 2026-10-09 ~22:50 UTC
|
||||
Operator: operator-646 (muse-646-patha)
|
||||
Branch: `dev/646/213-fix-ssh-perms`
|
||||
|
||||
## 1. authorized_keys permissions (port 2226 dial-in)
|
||||
|
||||
- `~/.ssh/authorized_keys` (`/home/hatch/.ssh/authorized_keys`):
|
||||
- before: `600 root:root`
|
||||
- ran `chmod 600 ~/.ssh/authorized_keys` per ticket
|
||||
- after: `600 root:root` (no-op — already correct)
|
||||
- sshd's requirement (private key file must not be group/world-writable,
|
||||
ideally 600) is satisfied. `~/.ssh` itself is `700`.
|
||||
|
||||
## 2. Container sshd
|
||||
|
||||
- `sshd` running (pid 2655, listener, 0 of 10-100 startups).
|
||||
- Listening on `0.0.0.0:22` and `[::]:22`.
|
||||
- `authorized_keys` holds 1 key:
|
||||
- `ssh-ed25519 SHA256:UOeqKF5BehWNmEpBSk53Qhz0Jd9aQXbFO0VKe2AVo8c`
|
||||
(comment `super@bl`) — dial-in identity belongs to super.
|
||||
|
||||
## 3. Reverse tunnel (VM 2226 → container:22)
|
||||
|
||||
- On VM 34.139.37.135 (as dev-operator-646): `127.0.0.1:2226` and
|
||||
`[::1]:2226` are LISTENING — the reverse tunnel is up.
|
||||
- Bind is loopback-only (no GatewayPorts), so dial-in must originate
|
||||
from the VM itself — expected for `ssh -R` forwards.
|
||||
|
||||
## 4. Dial-in path verdict
|
||||
|
||||
Container-side prerequisites are all green: perms 600, sshd listening,
|
||||
tunnel established, authorized key present. The final key-auth step can
|
||||
only be completed by the holder of the `super@bl` private key, so no
|
||||
full loopback auth was attempted from this operator identity.
|
||||
|
||||
Fixes #213
|
||||
@@ -1,57 +0,0 @@
|
||||
# Ticket #215 verification — SSH StrictModes on /home/hatch
|
||||
|
||||
Date: 2026-10-09 ~23:00 UTC
|
||||
Operator: operator-646 (muse-646-patha)
|
||||
Branch: `dev/646/215-strictmodes-fix`
|
||||
|
||||
## Ticket premise
|
||||
|
||||
#215 claims OpenSSH StrictModes rejects public-key auth on port 2226
|
||||
"for user hatch" because `/home/hatch` is `drwxrws---` (group-writable
|
||||
setgid), and asks whether `chmod g-w /home/hatch` or `StrictModes no`
|
||||
permits dial-in.
|
||||
|
||||
## Investigation
|
||||
|
||||
1. **No `hatch` user exists.** `/etc/passwd` has only `root` plus system
|
||||
`nologin` users. The only viable dial-in identity is `root`
|
||||
(`PermitRootLogin without-password`, i.e. pubkey-only).
|
||||
2. **Effective sshd config** (`sshd -T`): `strictmodes yes`,
|
||||
`authorizedkeysfile .ssh/authorized_keys .ssh/authorized_keys2`
|
||||
(relative to the login user's passwd home — for root, `/root`).
|
||||
3. **Root's auth path is StrictModes-clean** and does not include
|
||||
`/home/hatch`:
|
||||
- `/` → `755 root:root`
|
||||
- `/root` → `700 root:root`
|
||||
- `/root/.ssh` → `700 root:root`
|
||||
- `/root/.ssh/authorized_keys` → `600 root:root` (holds 646's
|
||||
`id_ed25519.pub` + `id_frontdoor.pub`, installed by
|
||||
`recover-after-rebuild.sh` §2 — by design)
|
||||
4. **Empirical dial-in test (the decisive check).** From the VM over the
|
||||
live reverse tunnel, with `/home/hatch` still `2770` (group-writable):
|
||||
`ssh -p 2226 root@127.0.0.1` with agent-forwarded `id_frontdoor`
|
||||
→ `DIALIN_OK`, `whoami` → `root`. Public-key dial-in on 2226
|
||||
**works with zero changes**.
|
||||
|
||||
## Verdict
|
||||
|
||||
Neither proposed remediation is required or was applied:
|
||||
|
||||
- `chmod g-w /home/hatch` — unnecessary for SSH (path not consulted);
|
||||
would also alter the setgid shared-directory semantics for no benefit.
|
||||
- `StrictModes no` in sshd config — unnecessary, and would weaken
|
||||
authentication security globally.
|
||||
|
||||
The StrictModes denial described in #215 cannot occur for the actual
|
||||
login path. No sshd reload was needed (no config changed).
|
||||
|
||||
## Adjacent real gap (flagged, not fixed — needs a decision)
|
||||
|
||||
`super@bl`'s ed25519 key (`SHA256:UOeqKF5B…`) lives only in
|
||||
`/home/hatch/.ssh/authorized_keys`, which sshd **never reads** (no
|
||||
`hatch` user exists). If super needs 2226 dial-in, that key must be
|
||||
appended to `/root/.ssh/authorized_keys`. The recover script
|
||||
deliberately installs only 646's own keys there, so this is a
|
||||
provisioning decision for 646/super — left untouched.
|
||||
|
||||
Fixes #215
|
||||
@@ -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