feat(cli): add box work command for unified worker signals and task orchestration

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