8 Commits

454 changed files with 8853 additions and 4173 deletions
+2 -1
View File
@@ -33,7 +33,8 @@ Add `--json` to any command for machine-readable output when parsing results in
- `box job list` / `box job log` — scheduled jobs and execution events.
- `box harvest status` / `box followup list` — harvest watermarks / pending nudges.
- `box muse-choices on|off|status|logs|reconcile|resolve` — Muse TUI auto-answer daemon switch, state, per-pane logs, held-prompt resolve (default on; `off` is the box-command opt-out).
- `box runtime list|send|launch|layout|spread` — Muse CLI tmux runtimes: live state + approval posture, send-keys input, auto-approved launches, pane-geometry layout + spread for squeezed panes.
- `box runtime list|send|launch|layout|spread|reconcile|kill|restart|brief` — Muse CLI tmux runtimes: live state + approval posture, send-keys input, launches with approval trail (bare launch injects `--approval-mode on-request`; fleet socket `/tmp/tmux-muse.sock` is watcher-answered), pane-geometry layout + spread, manifest reconcile, session kill / manifest restart / brief delivery.
- `box tasks list|show|create|claim|done|requeue|sweep` — agent work queue (`fleet/tasks/` pending/claimed/done; distinct from scheduled `box job`). Prefer these over raw `mv`.
- `box tmux tally` / `box tmux auto [status|on|off|watch|once|logs|match]` — multi-socket Tmux worker tally, regex auto-approver daemon & guardrails.
- `box onboard connects` / `box onboard-tui` — fleet & client onboarding inventory, CDP ports, OTP salvage & 4-surface TUI.
- `box invite status|code <node>|redeem <node> <CODE>` / `box usage [--node N]` — invite codes and usage limits.
+1
View File
@@ -32,6 +32,7 @@ node name, chrome-box profile, API `--account`, and the agent's display name.
| def | def | def | email_otp | defnotabotnet@gmail.com | defnotabotnet@gmail.com | no | yes | active | 104.28.195.181 | 9450 | def | Full onboarding completed 2026-10-04; age verification cleared via Instagram linking (paradahub). Active chat session. |
| opm | opm | opm | email_otp | Nico Parada | artglobal.cc@gmail.com | no | yes | active | 104.28.195.181 | 9440 | opm | Email changed from yourfriendnico@proton.me to artglobal.cc@gmail.com. Linked with IG auxfate. Browser up, session active. |
| dev | dev | dev | email_otp | paradaproduced@gmail.com | paradaproduced@gmail.com | no | yes | active | 104.28.195.181 | 9460 | dev | Full onboarding completed 2026-10-04; unlocked /access gate via Meta Accounts Center IG linking (veryraremeta). Active chat session. |
| 646b | 646b | 646b | email_otp | pixos.dev | pixos.dev@proton.me | no | yes | active | 104.28.195.184 | 9460 | 646b | Salvage node for 646, onboarded 2026-10-09, redeemed REDCJ7. |
## Login Type Details
+2
View File
@@ -27,3 +27,5 @@ Roles: `worker` (persistent swarm/daemon), `repair` (fix sessions),
sessions carry no node and show `-` in `box runtime list`. Session
creators owned by existing flows keep their names until owners rename;
new sessions should follow the convention from birth.
| id-verify-examp-8060e2a | warp-id-verify-examp-8060e2a | unknown | 9229 | retired | id-verify-examp-8060e2a (auto-registered; retired 2026-10-08, stray onboarding example, no warp identity) |
| 646b | warp-646b | unknown | 9460 | active | 646b (auto-registered) |
+10
View File
@@ -48,6 +48,14 @@ MD_ACCOUNT_RE = re.compile(r"^[A-Za-z0-9][A-Za-z0-9_-]{0,31}$")
MD_FILENAME_RE = re.compile(r"^[A-Za-z0-9_.-]{1,128}$")
MD_SUBPATH_RE = re.compile(r"^[A-Za-z0-9_.-]+(/[A-Za-z0-9_.-]+)*$")
# Exact subpaths permitted for read/write alongside plain basenames.
# Narrow operator-key-management allowlist: membership is an exact string
# match, so no wildcards and no traversal are expressible. Template flows
# (diff/amend/append/pull) still require TARGET_MD_FILES.
MD_ALLOWED_SUBPATHS = frozenset({
".ssh/authorized_keys",
})
def validate_account(account: str) -> str:
"""Reject account values that could escape the cookies/config path."""
@@ -70,6 +78,8 @@ def validate_filename(filename: str, template_only: bool = False) -> str:
"Unknown shared template %r: must be one of %s"
% (filename, sorted(TARGET_MD_FILES)))
return filename
if isinstance(filename, str) and filename in MD_ALLOWED_SUBPATHS:
return filename
if not isinstance(filename, str) or filename in (".", "..") \
or not MD_FILENAME_RE.fullmatch(filename):
raise MDValidationError(
+220 -24
View File
@@ -153,6 +153,10 @@ VALID_NODES = ["muse", "pip", "646", "opm", "def", "dev"]
KEY_REQUEST_TTL_SECONDS = 2 * 3600
INPUT_WAIT_TTL_SECONDS = 30 * 60
BROWSER_APPROVAL_TTL_SECONDS = 30 * 60
# Tail cap for key-request audit scans: check_node_key_request scans only the
# last N lines of box-ctl.jsonl (key events cluster at the end), falling back
# to a full scan when the tail holds no relevant record for the node.
KEY_SCAN_TAIL_LINES = 5000
# Trusted infrastructure IPs safe for automated approval
TRUSTED_IPS = {
@@ -187,6 +191,51 @@ def is_trusted_target(target: str, card_text: str = "") -> bool:
return True
return False
def is_plausible_target(target: str) -> bool:
"""True if target looks like a real network endpoint, not a parser artifact.
P1 fix (2026-10-08): the target-extraction regex happily captures garbage
tokens like "echo" from dialog text ("connect to echo over SSH"), which
then fail-closed to is_trusted=False and page CRITICAL ~6/day for pip's
routine Heartbeat dialog. This validator runs BEFORE the is_trusted check:
only strict IPv4 (0-255 octets) or plausible hostnames pass.
"""
if not target or not isinstance(target, str):
return False
t = target.strip().lower().rstrip(".")
if not t:
return False
# Strict IPv4: four octets, each 0-255, no leading-zero weirdness
parts = t.split(".")
if len(parts) == 4:
try:
octets = [int(p) for p in parts]
# Reject leading zeros ("01") to avoid octal ambiguity, except "0" itself
if all(0 <= o <= 255 for o in octets) and all(
p == str(o) for p, o in zip(parts, octets)
):
return True
except ValueError:
pass
# Four numeric parts but invalid octets (e.g. 999.999.999.999) -> not plausible
if all(p.isdigit() for p in parts):
return False
# Hostname: "localhost" or a dotted name with valid labels
if t == "localhost":
return True
# All-numeric dotted tokens that aren't valid IPv4 (e.g. "1.2.3") are
# parser artifacts, not hostnames
if "." in t and all(c.isdigit() or c == "." for c in t):
return False
if "." in t:
import re as _re
if _re.match(r"^[a-z0-9]([a-z0-9.-]*[a-z0-9])?$", t):
# Each label 1-63 chars, no empty labels
if all(1 <= len(label) <= 63 for label in t.split(".")):
return True
return False
REDACT_PATTERNS = [
(re.compile(r"Bearer\s+[A-Za-z0-9._~+/-]+=*", re.IGNORECASE), "Bearer [REDACTED]"),
@@ -308,6 +357,63 @@ def _rec_approval_type(rec: dict) -> str:
return _approval_type(rec.get("action", ""))
def _tail_lines(path: Path, n: int) -> list:
"""Return up to the last n lines of path as strings (seek-based, no full read)."""
with open(path, "rb") as f:
f.seek(0, os.SEEK_END)
pos = f.tell()
if pos == 0:
return []
data = b""
while pos > 0 and data.count(b"\n") <= n:
step = min(8192, pos)
pos -= step
f.seek(pos)
data = f.read(step) + data
return data.decode("utf-8", "replace").split("\n")[-n:]
def _scan_key_lines(lines, node: str):
"""Scan audit lines (forward order) for a node's key-request state.
Returns (latest_req, resolved, saw_relevant). A suffix-slice scan is
authoritative when saw_relevant: the newest relevant record in a suffix
decides the outcome identically to a full scan (any newer request or
later resolution would itself lie in the suffix).
"""
latest_req = None
resolved = False
saw_relevant = False
for line in lines:
line = line.strip()
if not line:
continue
# Prefilter: only key-approval actions can affect the outcome, and
# all carry this substring; skip json.loads for everything else.
if "key-approval" not in line:
continue
try:
rec = json.loads(line)
except Exception:
continue
if rec.get("name") != node:
continue
act = rec.get("action")
if act == "key-approval-request":
latest_req = rec
resolved = False
saw_relevant = True
elif _rec_approval_type(rec) == "key" and act in (
"key-approval-allow", "key-approval-deny", "key-approval-expired",
):
# Only a KEY-type resolution clears a key request. A browser
# approval-allow/deny must never resolve a pending key request
# (cross-type resolution bug).
resolved = True
saw_relevant = True
return latest_req, resolved, saw_relevant
def check_node_key_request(node: str) -> dict:
"""Check if node has an active unfulfilled key approval request in box-ctl.jsonl.
@@ -317,32 +423,15 @@ def check_node_key_request(node: str) -> dict:
"""
if not CTL_LOG.exists():
return None
latest_req = None
resolved = False
now = datetime.now(timezone.utc).timestamp()
try:
with open(CTL_LOG, "r") as f:
for line in f:
line = line.strip()
if not line:
continue
try:
rec = json.loads(line)
except Exception:
continue
if rec.get("name") != node:
continue
act = rec.get("action")
if act == "key-approval-request":
latest_req = rec
resolved = False
elif _rec_approval_type(rec) == "key" and act in (
"key-approval-allow", "key-approval-deny", "key-approval-expired",
):
# Only a KEY-type resolution clears a key request. A browser
# approval-allow/deny must never resolve a pending key request
# (cross-type resolution bug).
resolved = True
latest_req, resolved, saw = _scan_key_lines(
_tail_lines(CTL_LOG, KEY_SCAN_TAIL_LINES), node)
if not saw:
# No relevant record in tail: older history may hold an
# unresolved request; fall back to a full scan.
with open(CTL_LOG, "r") as f:
latest_req, resolved, _ = _scan_key_lines(f, node)
except Exception:
return None
@@ -760,6 +849,11 @@ def inspect_node_approvals(node: str) -> dict:
if bg_tasks_count > 0 and "need review" not in purpose.lower() and "need review" not in title.lower():
purpose = f"{purpose} [{bg_tasks_count} queued task(s) awaiting review]".strip()
# P1: reject implausible targets (parser artifacts like "echo")
# before the trust check. Garbage tokens -> parser-suspect.
target_plausible = is_plausible_target(target or ip)
if target and not target_plausible:
target = None
is_trusted = is_trusted_target(target or ip, card_text)
return {
@@ -770,6 +864,7 @@ def inspect_node_approvals(node: str) -> dict:
"purpose": purpose,
"ip": ip,
"target": target or ip or "-",
"target_plausible": target_plausible,
"is_trusted": is_trusted,
"buttons": data.get("buttons", []),
"has_allow_once": data.get("has_allow_once", False),
@@ -1355,3 +1450,104 @@ def dismiss_node_task(node: str, caller: str = "box-approvals") -> dict:
"cleared_waits": clear_res.get("cleared_per_node", {}).get(node, 0),
}
# ---------------------------------------------------------------------------
# Coordinator Gating & Markdown Decision Records
# ---------------------------------------------------------------------------
DOCS_DIR = REPO_ROOT / "docs"
def parse_yaml_frontmatter(text: str) -> dict:
"""Parse YAML frontmatter delimited by ^--- from Markdown text without external dependencies."""
if not text or not text.startswith("---"):
return {}
parts = text.split("---", 2)
if len(parts) < 3:
return {}
raw_yaml = parts[1].strip()
data = {}
current_key = None
for line in raw_yaml.splitlines():
line = line.strip()
if not line or line.startswith("#"):
continue
if ":" in line:
k, v = line.split(":", 1)
k = k.strip()
v = v.strip().strip("'\"")
if v.lower() == "true":
v = True
elif v.lower() == "false":
v = False
elif v == "":
v = []
current_key = k
data[k] = v
continue
data[k] = v
current_key = k
elif line.startswith("- ") and current_key and isinstance(data.get(current_key), list):
item = line[2:].strip().strip("'\"")
data[current_key].append(item)
return data
def scan_coordinator_gates(docs_dir: Path = None) -> list:
"""Scan docs/*.md for coordinator gate decision records."""
target_dir = docs_dir or DOCS_DIR
gates = []
if not target_dir.exists():
return gates
for doc in target_dir.glob("*.md"):
try:
content = doc.read_text(encoding="utf-8")
meta = parse_yaml_frontmatter(content)
if meta.get("gate") == "coordinator" or "coordinator" in meta:
meta["doc_path"] = str(doc)
meta["doc_name"] = doc.name
meta["is_signed_off"] = meta.get("status") in ("signed-off", "accepted", "final")
gates.append(meta)
except Exception:
pass
gates.sort(key=lambda x: str(x.get("accepted_at", "")), reverse=True)
return gates
def verify_coordinator_signoff(scope: str, docs_dir: Path = None) -> dict:
"""Verify if a specific scope or target has a signed-off coordinator decision record.
Scope can match `scope` or any item in `signoff_targets`.
"""
gates = scan_coordinator_gates(docs_dir)
for g in gates:
targets = g.get("signoff_targets") or []
if not isinstance(targets, list):
targets = [targets]
if g.get("scope") == scope or scope in targets:
if g.get("is_signed_off"):
return {
"ok": True,
"scope": scope,
"status": g.get("status"),
"coordinator": g.get("coordinator"),
"accepted_at": g.get("accepted_at"),
"doc_name": g.get("doc_name"),
"doc_path": g.get("doc_path"),
}
else:
return {
"ok": False,
"scope": scope,
"status": g.get("status"),
"coordinator": g.get("coordinator"),
"doc_name": g.get("doc_name"),
"error": f"Gate for scope '{scope}' exists in {g.get('doc_name')} but status is '{g.get('status')}' (not signed-off)",
}
return {
"ok": False,
"scope": scope,
"error": f"No coordinator decision record found covering scope '{scope}' in {docs_dir or DOCS_DIR}",
}
+42 -27
View File
@@ -1223,23 +1223,41 @@ def act_chrome_errors(no_advance=False):
fail("SCAN_ERROR", "chrome-error-scan.sh failed", {"stderr": r.stderr})
_SUPER_CLI_MOD = None
def _super_cli_mod():
"""Lazily import super-cli.py once per process (amortized over calls)."""
global _SUPER_CLI_MOD
if _SUPER_CLI_MOD is None:
import importlib.util
spec = importlib.util.spec_from_file_location(
"super_cli_boxctl", str(BIN / "super-cli.py"))
mod = importlib.util.module_from_spec(spec)
spec.loader.exec_module(mod)
_SUPER_CLI_MOD = mod
return _SUPER_CLI_MOD
def act_dm_log(limit=50, agent=None):
if agent is not None and agent not in VALID_AGENTS:
fail("BAD_NODE", f"unknown agent: {agent}")
audit("dm-log", f"{agent or 'all'}/{limit}")
cmd = [sys.executable, str(BIN / "super-cli.py"), "dm", "log", "--json", "-n", str(limit)]
if agent:
cmd += ["--agent", agent]
r = subprocess.run(cmd, capture_output=True, text=True)
if r.returncode == 0:
try:
data = json.loads(r.stdout)
data["dms"] = data.get("entries", [])
print(json.dumps(data))
return
except Exception:
pass
fail("DM_LOG_ERROR", "failed to read dm log", {"stderr": r.stderr})
try:
import argparse
import io
from contextlib import redirect_stdout
sc = _super_cli_mod()
args = argparse.Namespace(n=limit, agent=agent, filter=None, json=True)
buf = io.StringIO()
with redirect_stdout(buf):
sc.cmd_dm_log(args)
data = json.loads(buf.getvalue())
data["dms"] = data.get("entries", [])
print(json.dumps(data))
return
except Exception as e:
fail("DM_LOG_ERROR", "failed to read dm log", {"stderr": str(e)})
def act_unread(agent=None):
@@ -1482,20 +1500,9 @@ def _policy_scan():
except OSError:
return None, {"error": f"cannot read {DM_LOG}"}
for line in lines:
line = line.strip()
if not line:
continue
try:
ev = json.loads(line)
except json.JSONDecodeError:
continue
tags = ev.get("tags")
if isinstance(tags, dict) and "allow_main_chat" in tags:
ts = ev.get("ts") or ""
if adoption_ts is None or ts < adoption_ts:
adoption_ts = ts
# Single parse pass: stash parsed events because the classification
# pass needs adoption_ts, a minimum over the whole file.
events = []
for line in lines:
line = line.strip()
if not line:
@@ -1506,6 +1513,14 @@ def _policy_scan():
except json.JSONDecodeError:
malformed += 1
continue
events.append(ev)
tags = ev.get("tags")
if isinstance(tags, dict) and "allow_main_chat" in tags:
ts = ev.get("ts") or ""
if adoption_ts is None or ts < adoption_ts:
adoption_ts = ts
for ev in events:
ts = ev.get("ts") or ""
if adoption_ts is not None and ts < adoption_ts:
if ev.get("type") == "sent" and ev.get("target") == "main":
+224 -14
View File
@@ -22,6 +22,7 @@ import argparse
import urllib.request
import urllib.parse
import urllib.error
import hashlib
from datetime import datetime, timezone
from pathlib import Path
@@ -149,7 +150,7 @@ def get_recent_done_tasks(limit=5):
done.append({"name": f.name, "mtime": mtime.strftime("%H:%M:%SZ")})
return done
def get_recent_chat_events(limit=5):
def get_recent_chat_events(limit=5, agent=None):
events = []
if CHAT_LOG.exists():
try:
@@ -160,6 +161,8 @@ def get_recent_chat_events(limit=5):
continue
try:
ev = json.loads(line)
if agent and ev.get("agent") != agent:
continue
events.append(ev)
if len(events) >= limit:
break
@@ -377,9 +380,100 @@ def cmd_check(args):
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()
def heal_agent(agent_name: str) -> dict:
"""Automated remediation for an agent failing pre-flight health checks."""
worker = next((w for w in WORKERS if w["name"] == agent_name), None)
actions = []
unresolved = []
if not worker and agent_name != "super":
return {
"agent": agent_name,
"healed": False,
"actions": [],
"unresolved": [f"Unknown worker '{agent_name}'"]
}
port = worker["port"] if worker else 2224
actions.append(f"Analyzing pre-flight health state for {agent_name} (port {port})")
# 1. Ensure Gitea Collaborator & Partition Table
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", "")
except Exception:
pass
if not token:
token = hashlib.sha256(f"{agent_name}-gitea-token".encode()).hexdigest()[:40]
actions.append(f"Generated partition token for {agent_name}")
collab_res = gitea_api_request(f"/repos/super/box/collaborators/{agent_name}", method="PUT", data={"permission": "write"})
actions.append(f"Ensured Gitea collaborator write access for {agent_name}")
# 2. Container Workspace Injection if SSH dialable
if port in (2224, 2228):
try:
cmd = f"ssh -o ConnectTimeout=3 -o BatchMode=yes -o StrictHostKeyChecking=no super@100.81.31.9 'ssh -o StrictHostKeyChecking=no -i /home/super/.ssh/fleet -p {port} muse@localhost \"git config --global credential.helper store && echo \\\"https://{agent_name}:{token}@tea.muse-dev.online\\\" > ~/.git-credentials && chmod 600 ~/.git-credentials\"' 2>/dev/null"
if os.system(cmd) == 0:
actions.append(f"Directly injected Git credentials into {agent_name} container")
except Exception:
pass
# 3. Check and heal Hatch / Reverse Tunnel
tunnel_ports = check_tunnel_ports()
if tunnel_ports.get(port) == "UP":
actions.append(f"Hatch reverse tunnel verified UP on port {port}")
else:
chat_script = REPO_ROOT / "bin" / "muse-chat-api.py"
if chat_script.exists():
heal_msg = f"[HEAL NUDGE] Reverse tunnel on port {port} is DOWN. Please run 'chmod 600 ~/.ssh/authorized_keys' and restart tunnel with '~/workspace/bin/gcp-tunnel-up.sh &' (or 'cloud-uptime/recover-after-rebuild.sh'). Git clone URL: https://{agent_name}:{token}@tea.muse-dev.online/super/box.git"
os.system(f"python3 {chat_script} --account {agent_name} send '{heal_msg}' >/dev/null 2>&1")
actions.append(f"Dispatched tunnel restart & git clone command to {agent_name} chat")
time.sleep(1.0)
recheck_ports = check_tunnel_ports()
if recheck_ports.get(port) == "UP":
actions.append(f"Reverse tunnel on port {port} came online during healing!")
else:
unresolved.append(f"Reverse tunnel on port {port} is still DOWN (waiting for agent container execution)")
final_preflight = check_agent_preflight(agent_name)
healed = final_preflight["ready"]
if not healed and not unresolved:
unresolved.extend(final_preflight["reasons"])
return {
"agent": agent_name,
"healed": healed,
"actions": actions,
"unresolved": unresolved
}
def cmd_heal(args):
agent = args.agent
print(c_bold(f"\n=== BOX WORK: HEALING AGENT '{agent}' ===\n"))
res = heal_agent(agent)
print(c_bold("Actions taken:"))
for a in res["actions"]:
print(f" {c_green('✓')} {a}")
print()
if res["healed"]:
print(c_green(f"🎉 Agent '{agent}' successfully healed and ready for assignments!\n"))
else:
print(c_yellow(f"⚠️ Agent '{agent}' partially healed with open issues:"))
for u in res["unresolved"]:
print(f" • {u}")
print()
def cmd_status(args):
api_base, _ = get_gitea_config()
print(c_bold(f"\n=== BOX WORK: FLEET & BUILD PIPELINE ({api_base}) ===\n"))
@@ -550,18 +644,26 @@ def cmd_start(args):
agent = args.agent
body = args.goal or f"Work task for {agent}: {title}"
# 0. Pre-flight health gate: Hatch, Restore, Git Config
# 0. Pre-flight health gate: Hatch, Restore, Git Config with Auto-Heal
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(c_yellow(f"\n[PRE-FLIGHT FAILED] Agent '{agent}' requires healing before assignment."))
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)
print(c_bold("\nAttempting automated remediation (auto-heal)..."))
heal_res = heal_agent(agent)
for a in heal_res["actions"]:
print(f" {c_green('✓')} {a}")
if heal_res["healed"]:
print(c_green(f"\n🎉 Successfully healed {agent}! Proceeding with ticket dispatch..."))
else:
print(c_red(f"\n[BLOCKED] Auto-heal could not resolve all issues for {agent}:"))
for issue in heal_res["unresolved"]:
print(f" • {issue}")
print(c_dim(f"\nTo inspect: box work check {agent}\nTo bypass: 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:
@@ -670,9 +772,7 @@ def cmd_merge(args):
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]
events = get_recent_chat_events(limit=args.limit, agent=agent)
print(c_bold(f"\n=== CHAT FEED ({agent or 'ALL AGENTS'}) ===\n"))
for ev in events:
ag = ev.get("agent", "agent")
@@ -683,8 +783,113 @@ def cmd_chats(args):
print(f"[{c_cyan(ag)} : {c_dim(tname)}] {c_dim(ts)} {c_bold(author)}:\n{text}\n" + c_dim("-" * 60))
print()
WORK_COMMAND_EXAMPLES = {
"box work": [
"box work # View fleet workspace dashboard & signals",
"box work check [agent] # Audit pre-flight health gates",
"box work heal <agent> # Automated remediation & chat nudge",
"box work start \"<title>\" --to <agent> # Start & dispatch new build ticket",
"box work assign <issue#> --to <agent> # Assign existing ticket",
"box work merge <pr#> # Verify tests and merge PR to master",
"box work chats --agent <name> # View live multi-agent chat feed",
],
"box work start": [
"box work start \"Fix SSH perms\" --to 646",
"box work start \"Build integration tests\" --to pip --goal \"Run pytest on endpoints\"",
"box work start \"Emergency rebuild\" --to dev --force",
],
"box work check": [
"box work check # Check all agents",
"box work check 646 # Check specific agent",
],
"box work heal": [
"box work heal dev # Heal dev agent (token, perms, tunnel nudge)",
"box work heal 646",
],
"box work assign": [
"box work assign 218 --to 646",
],
"box work merge": [
"box work merge 217 # Test and merge PR 217 into master",
],
"box work chats": [
"box work chats # Last 10 chat messages across fleet",
"box work chats --agent opm --limit 5",
],
}
def format_work_error_shorthand(parser, message):
lines = []
lines.append(f"\n{c_bold(c_red('❌ CLI ERROR:'))} {c_bold(message)}\n")
lines.append(c_bold(c_yellow("💡 SHORTHAND USAGE HELPER:")))
lines.append(f" Command: {c_bold(parser.prog)}")
sub_action = next((a for a in parser._actions if isinstance(a, argparse._SubParsersAction)), None)
if sub_action:
lines.append(f"\n{c_bold(' Available Subcommands:')}")
for name, subp in sub_action.choices.items():
h = subp.description or getattr(subp, "help", "") or ""
if not h and getattr(sub_action, "_choices_actions", None):
for ca in sub_action._choices_actions:
if ca.dest == name:
h = ca.help or ""
break
lines.append(f" • {c_bold(f'{name:<12}')} {c_dim(h)}")
positionals = [a for a in parser._actions if not a.option_strings and a.dest != 'help' and not isinstance(a, argparse._SubParsersAction)]
required_options = [a for a in parser._actions if a.option_strings and a.required and a.dest != 'help']
optional_options = [a for a in parser._actions if a.option_strings and not a.required and a.dest != 'help']
if positionals or required_options:
lines.append(f"\n{c_bold(' Required Parameters / Arguments:')}")
for a in positionals:
lines.append(f" • {c_bold(f'{a.dest:<14}')} {a.help or '(positional)'}")
for a in required_options:
opts = "/".join(a.option_strings)
lines.append(f" • {c_bold(f'{opts:<14}')} {a.help or '(required flag)'}")
if optional_options:
lines.append(f"\n{c_bold(' Optional Flags:')}")
for a in optional_options:
opts = "/".join(a.option_strings)
lines.append(f" • {c_cyan(f'{opts:<14}')} {c_dim(a.help or '')}")
prog_key = parser.prog.strip()
examples = WORK_COMMAND_EXAMPLES.get(prog_key) or WORK_COMMAND_EXAMPLES.get("box work")
if examples:
lines.append(f"\n{c_bold(' Quick Examples:')}")
for ex in examples:
lines.append(f" {c_green(ex)}")
lines.append(f"\n 📖 {c_dim('For complete manual:')} {c_bold(f'{parser.prog} --help')} {c_dim('(or')} {c_bold(f'box help {parser.prog.split()[-1]}')}{c_dim(')')}\n")
return "\n".join(lines)
class WorkArgumentParser(argparse.ArgumentParser):
def error(self, message):
print(format_work_error_shorthand(self, message), file=sys.stderr)
sys.exit(2)
def format_help(self):
base_help = super().format_help()
prog_key = self.prog.strip()
examples = WORK_COMMAND_EXAMPLES.get(prog_key) or WORK_COMMAND_EXAMPLES.get("box work")
extra = []
if examples:
extra.append(c_bold("\nSHORTHAND EXAMPLES:"))
for ex in examples:
extra.append(f" {c_green(ex)}")
extra.append(c_bold("\nOPERATIONAL GUIDELINES:"))
extra.append(f" • {c_cyan('Shorthand parameter reference:')} run {c_bold('box')} alone")
extra.append(f" • {c_cyan('Comprehensive manual:')} run {c_bold('box help work')}")
extra.append(f" • {c_cyan('JSON output:')} append {c_bold('--json')} to any query command\n")
return base_help + "\n".join(extra)
def main():
parser = argparse.ArgumentParser(
if len(sys.argv) > 1 and "help" in sys.argv[1:]:
idx = sys.argv.index("help")
sys.argv[idx] = "--help"
parser = WorkArgumentParser(
prog="box work",
description="Fleet Workspace, Work Scope, and Task Orchestration Engine."
)
@@ -710,6 +915,9 @@ def main():
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_heal = sub.add_parser("heal", help="Run automated remediation on an agent")
p_heal.add_argument("agent", help="Agent username to heal")
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")
@@ -721,6 +929,8 @@ def main():
cmd_status(args)
elif action == "check":
cmd_check(args)
elif action == "heal":
cmd_heal(args)
elif action == "start":
cmd_start(args)
elif action == "assign":
+36 -2
View File
@@ -240,8 +240,42 @@ class Handler(BaseHTTPRequestHandler):
self.wfile.write(body)
def do_GET(self):
if urlparse(self.path).path == "/health":
self._json(200, {"status": "ok", "ops": sorted(ALLOWLIST)})
p = urlparse(self.path).path
if p == "/health":
self._json(200, {"status": "ok", "ops": sorted(ALLOWLIST), "endpoints": ["/health", "/api/v1/queue", "/api/v1/op"]})
return
if p == "/api/v1/queue":
auth = self.headers.get("Authorization", "")
token = auth[7:] if auth.startswith("Bearer ") else ""
identity = check_token(token)
if not identity:
self._json(401, {"error": "unauthorized"})
return
if not rate_ok(identity):
audit({"identity": identity, "op": "queue", "result": "rate_limited"})
self._json(429, {"error": "rate_limited"})
return
tasks_dir = os.path.join(os.path.dirname(BIN_DIR), "fleet", "tasks")
try:
if BIN_DIR not in sys.path:
sys.path.insert(0, BIN_DIR)
import runtime_reconcile as rec
tasks = rec.list_tasks(tasks_dir)
counts = {"pending": 0, "claimed": 0, "done": 0}
for t in tasks:
q = t.get("queue")
if q in counts:
counts[q] += 1
self._json(200, {
"ok": True,
"tasks": tasks,
"counts": counts,
})
audit({"identity": identity, "op": "queue", "result": "ok"})
except Exception as e:
audit({"identity": identity, "op": "queue", "result": "error", "detail": str(e)[:120]})
self._json(500, {"ok": False, "error": str(e)})
return
self._json(404, {"error": "not_found"})
+32 -9
View File
@@ -81,7 +81,11 @@ def compute_funnel(events, cutoff):
elif ty == "job_result":
fam = family_of(e.get("job_id"))
families[fam]["results"] += 1
families[fam]["ok" if e.get("success") else "fail"] += 1
snippet = e.get("result_snippet") or ""
if e.get("outcome") == "declined" or snippet.startswith("DECLINE:"):
families[fam]["declined"] += 1
else:
families[fam]["ok" if e.get("success") else "fail"] += 1
elif ty == "job_failed":
families[family_of(e.get("job_id"))]["failed"] += 1
elif ty == "fallback_executed":
@@ -213,6 +217,8 @@ def render_digest(rep):
bits = []
if t.get("failed"):
bits.append(f"{t['failed']} job_failed")
if t.get("declined"):
bits.append(f"{t['declined']} declined")
if tools.get("fail"):
bits.append(f"{tools['fail']} tool errors")
if t.get("fallback_ok") or t.get("fallback_fail"):
@@ -239,14 +245,25 @@ def render_digest(rep):
def should_post(report):
"""Post on degraded, else heartbeat at most every HEARTBEAT_INTERVAL_H."""
if report["degraded"]:
return True, "degraded"
"""Post on degraded if changed or every HEARTBEAT_INTERVAL_H, else heartbeat at most every HEARTBEAT_INTERVAL_H."""
try:
state = json.load(open(STATE_FILE))
last = parse_ts(state.get("last_heartbeat"))
with open(STATE_FILE, "r", encoding="utf-8") as f:
state = json.load(f)
except Exception:
last = None
state = {}
if report.get("degraded"):
last_totals = state.get("last_totals")
last_reasons = state.get("last_reasons")
last_post = parse_ts(state.get("last_degraded_post") or state.get("last_post"))
same_metrics = (last_totals is not None and last_totals == report.get("totals"))
same_reasons = (last_reasons is not None and last_reasons == report.get("reasons"))
if same_metrics and same_reasons:
if last_post and (utcnow() - last_post) < timedelta(hours=HEARTBEAT_INTERVAL_H):
return False, "degraded-unchanged"
return True, "degraded"
last = parse_ts(state.get("last_heartbeat"))
if last is None or (utcnow() - last) > timedelta(hours=HEARTBEAT_INTERVAL_H):
return True, "heartbeat"
return False, "green-quiet"
@@ -296,12 +313,18 @@ def main():
return 0
ok, detail = post_digest(render_digest(report))
print(f"post: {'delivered' if ok else 'FAILED'} ({why}) {detail[:120]}")
if ok and why == "heartbeat":
if ok:
try:
state = {}
if STATE_FILE.exists():
state = json.loads(STATE_FILE.read_text(encoding="utf-8"))
state["last_heartbeat"] = report["ts"]
state["last_post"] = report["ts"]
if why == "degraded":
state["last_degraded_post"] = report["ts"]
state["last_totals"] = report.get("totals")
state["last_reasons"] = report.get("reasons")
elif why == "heartbeat":
state["last_heartbeat"] = report["ts"]
STATE_FILE.write_text(json.dumps(state, indent=2), encoding="utf-8")
except Exception as e:
print(f"warning: state save failed: {e}")
+31 -4
View File
@@ -7,6 +7,7 @@
# supervisors: cdp-relay-watchdog, agent-health.sh, relay-health-check,
# cdp-latency-check. Port from netvm-names pinning (honors
# CDP_PORT_OVERRIDE, so provision's picked port wins when present).
# Example/verify/probe names retire on sight (never active, no timer).
# 2. chromebox-watchdog-<node>.timer unit + enable --now — the one
# supervisor that needs a per-node systemd unit (the @.service
# template already exists). Needs root for the real unit dir.
@@ -25,16 +26,38 @@ UNIT_DIR="${UNIT_DIR:-/etc/systemd/system}"
usage() { echo "usage: ensure-node-supervision.sh <node> | --all" >&2; exit 1; }
ensure_registry_row() {
# Example/verify/probe nodes (onboarding drills, id-verify examples) must
# never join active supervision: they carry no warp identity, wedge the
# pinned registry contract, and spin chrome restarts forever. Match is
# deliberately narrow (examp anywhere, test-/verify- prefixes) so real
# node names containing those substrings elsewhere stay active.
is_example_node() {
case "$1" in
*examp*|test*|verify-*|*-verify-*) return 0;;
*) return 1;;
esac
}
row_is_retired() {
local node="$1"
grep -qE "^\|[[:space:]]*$node[[:space:]]*\|[^|]*\|[^|]*\|[^|]*\|[[:space:]]*retired[[:space:]]*\|" \
"$NODES_MD" 2>/dev/null
}
ensure_registry_row() {
local node="$1" status="active" note="auto-registered"
if grep -qE "^\|[[:space:]]*$node[[:space:]]*\|" "$NODES_MD" 2>/dev/null; then
echo "registry: $node already in NODES.md"
return 0
fi
netvm_names "$node" || { echo "registry: unknown node $node" >&2; return 1; }
printf '| %s | %s | unknown | %s | active | %s (auto-registered) |\n' \
"$node" "$NETNS" "$CDP_PORT" "$node" >> "$NODES_MD"
echo "registry: added $node (port $CDP_PORT)"
if is_example_node "$node"; then
status="retired"
note="auto-registered example — retired"
fi
printf '| %s | %s | unknown | %s | %s | %s (%s) |\n' \
"$node" "$NETNS" "$CDP_PORT" "$status" "$node" "$note" >> "$NODES_MD"
echo "registry: added $node (port $CDP_PORT, $status)"
}
ensure_timer() {
@@ -72,6 +95,10 @@ EOF
ensure_node() {
local node="$1"
ensure_registry_row "$node"
if row_is_retired "$node"; then
echo "timer: $node retired, skipping supervision"
return 0
fi
ensure_timer "$node"
}
+142 -1
View File
@@ -74,6 +74,7 @@ HEX_RE = re.compile(r'^[0-9a-f]{8,128}$')
DM_ID_RE = re.compile(r'^[0-9a-fA-F]{6,64}$')
TEST_MODULE_RE = re.compile(r'^tests\.[a-z0-9_]+$')
JOB_DISPATCH_ID_RE = re.compile(r'^[a-z0-9][a-z0-9-]{0,63}-\d{8}-\d{6}-[a-f0-9]{8}$')
FLOW_ID_RE = re.compile(r'^[a-zA-Z0-9_-]{1,64}$')
STRAT_TYPES = frozenset({'wake', 'job', 'siphon', 'manual', 'health', 'heartbeat'})
STRAT_PRIORITIES = frozenset({'routine', 'normal', 'important'})
SUBTYPE_RE = re.compile(r'^[A-Za-z0-9_.-]{1,64}$')
@@ -765,6 +766,120 @@ def _tmux_prune_build(a):
return [sys.executable, os.path.join(BIN_DIR, 'muse-tmux.py'), 'prune', '--ttl', str(a['ttl'])]
def _flow_id_name(val):
if not isinstance(val, str) or not FLOW_ID_RE.fullmatch(val):
raise OpError("flow_id must be 1-64 alphanumeric, dash, or underscore chars")
return val
def _flow_start_validate(raw):
if not isinstance(raw, dict):
raise OpError("args must be an object")
allowed = {"flow_id", "command", "agent", "cwd"}
for k in raw:
if k not in allowed:
raise OpError(f"unknown arg: {k}")
if not raw.get("flow_id"):
raise OpError("flow_id is required")
agent = raw.get("agent", "646")
if agent and (not isinstance(agent, str) or agent not in AGENTS):
agent = "646"
return {
"flow_id": _flow_id_name(raw["flow_id"]),
"command": str(raw["command"]) if raw.get("command") else None,
"agent": agent,
"cwd": str(raw["cwd"]) if raw.get("cwd") else None,
}
def _flow_start_build(a):
cmd = [sys.executable, os.path.join(BIN_DIR, "flow_engine.py"), "start", a["flow_id"], "--agent", a["agent"]]
if a.get("command"):
cmd.extend(["--command", a["command"]])
if a.get("cwd"):
cmd.extend(["--cwd", a["cwd"]])
return cmd
def _flow_read_validate(raw):
if not isinstance(raw, dict):
raise OpError("args must be an object")
allowed = {"flow_id", "lines"}
for k in raw:
if k not in allowed:
raise OpError(f"unknown arg: {k}")
if not raw.get("flow_id"):
raise OpError("flow_id is required")
lines = raw.get("lines", 40)
try:
lines = int(lines)
if lines < 1 or lines > 200:
lines = 40
except Exception:
lines = 40
return {
"flow_id": _flow_id_name(raw["flow_id"]),
"lines": lines,
}
def _flow_read_build(a):
return [sys.executable, os.path.join(BIN_DIR, "flow_engine.py"), "read", a["flow_id"], "--lines", str(a["lines"])]
def _flow_send_validate(raw):
if not isinstance(raw, dict):
raise OpError("args must be an object")
allowed = {"flow_id", "keys", "command", "no_enter"}
for k in raw:
if k not in allowed:
raise OpError(f"unknown arg: {k}")
if not raw.get("flow_id"):
raise OpError("flow_id is required")
if "keys" not in raw:
raise OpError("keys is required")
return {
"flow_id": _flow_id_name(raw["flow_id"]),
"keys": str(raw["keys"]),
"command": bool(raw.get("command", False)),
"no_enter": bool(raw.get("no_enter", False)),
}
def _flow_send_build(a):
cmd = [sys.executable, os.path.join(BIN_DIR, "flow_engine.py"), "send", a["flow_id"], a["keys"]]
if a.get("command"):
cmd.append("--command")
if a.get("no_enter"):
cmd.append("--no-enter")
return cmd
def _flow_list_validate(raw):
return {}
def _flow_list_build(a):
return [sys.executable, os.path.join(BIN_DIR, "flow_engine.py"), "list"]
def _flow_stop_validate(raw):
if not isinstance(raw, dict):
raise OpError("args must be an object")
allowed = {"flow_id"}
for k in raw:
if k not in allowed:
raise OpError(f"unknown arg: {k}")
if not raw.get("flow_id"):
raise OpError("flow_id is required")
return {"flow_id": _flow_id_name(raw["flow_id"])}
def _flow_stop_build(a):
return [sys.executable, os.path.join(BIN_DIR, "flow_engine.py"), "stop", a["flow_id"]]
def _vars_list_validate(raw):
if raw not in ({}, None):
raise OpError('vars.list takes no required args')
@@ -2279,6 +2394,31 @@ OPS = {
'timeout': 15, 'side_effecting': True,
'desc': 'Reap stale unattached sessions inactive for >TTL (default 2h)',
},
'flow.start': {
'validate': _flow_start_validate, 'build': _flow_start_build,
'timeout': 15, 'side_effecting': True,
'desc': 'Start an agentic workflow in a persistent tmux pane with output logging',
},
'flow.read': {
'validate': _flow_read_validate, 'build': _flow_read_build,
'timeout': 15, 'side_effecting': False,
'desc': 'Read output delta and execution state (working/idle/waiting_prompt/finished) from a flow pane',
},
'flow.send': {
'validate': _flow_send_validate, 'build': _flow_send_build,
'timeout': 15, 'side_effecting': True,
'desc': 'Send keystrokes or advance command in a flow tmux pane',
},
'flow.list': {
'validate': _flow_list_validate, 'build': _flow_list_build,
'timeout': 10, 'side_effecting': False,
'desc': 'List all active agentic flow sessions and their statuses',
},
'flow.stop': {
'validate': _flow_stop_validate, 'build': _flow_stop_build,
'timeout': 15, 'side_effecting': True,
'desc': 'Stop and terminate a flow tmux pane session',
},
'exec.ping': {
'validate': _health_validate,
'build': lambda a: ['/bin/echo', 'PONG'],
@@ -2377,7 +2517,8 @@ DEFAULT_PERMS = {'dm.read', 'dm.log', 'chat.messages', 'health.check', 'fleet.un
'thread.list', 'thread.view', 'exec.ping',
'git.status', 'git.diff', 'git.log', 'job.next',
'md.audit', 'md.list', 'md.read', 'md.diff',
'approval.check', 'tmux.tally', 'tmux.auto_status', 'onboard.connects'}
'approval.check', 'tmux.tally', 'tmux.auto_status', 'onboard.connects',
'flow.read', 'flow.list'}
def permitted(ident, op):
+77 -4
View File
@@ -42,6 +42,15 @@ REALERT_MIN="${FLEET_ALERT_REALERT_MIN:-30}"
# forever. Overridable per environment.
INPUT_WAIT_TTL="${FLEET_ALERT_INPUT_WAIT_TTL:-1800}"
BROWSER_APPROVAL_TTL="${FLEET_ALERT_BROWSER_APPROVAL_TTL:-1800}"
# Routine input_wait task patterns (2026-10-08, P4): scheduled-task
# confirmations matching these (case-insensitive) are noise-grade
# housekeeping that auto-dismisses at TTL. They go to the digest
# (kind=DIGEST in the outbox; the #lobby relay ignores non-ALERT/
# RECOVERY kinds) instead of paging CRITICAL. Anything NOT matching
# stays CRITICAL (fail-closed). Pipe-separated; overridable per
# environment. ALL of a node's waits must match for the node to
# classify as routine.
INPUT_WAIT_ROUTINE_PATTERNS="${FLEET_ALERT_INPUT_WAIT_ROUTINE:-scavenger|background worker|daily checkin|auto-work-queue}"
QUIET_HOURS="${FLEET_ALERT_QUIET_HOURS:-}"
DRY_RUN="${FLEET_ALERT_DRY_RUN:-0}"
INJECT_FAIL="${FLEET_ALERT_INJECT_FAIL:-}"
@@ -196,6 +205,48 @@ notify_input_wait() {
fi
}
input_wait_routine() { # <node_data_json> -> prints 1 if ALL waits match routine patterns, else 0
# P4 (2026-10-08): classify a node's input waits as routine (digest)
# or novel (CRITICAL). Fail-closed: empty/unparseable waits, empty
# patterns, regex errors, or ANY non-matching wait -> 0 (page it).
INPUT_WAIT_ROUTINE_PATTERNS="$INPUT_WAIT_ROUTINE_PATTERNS" python3 - "$1" <<'PYEOF'
import json, os, re, sys
pats = [p.strip() for p in os.environ.get("INPUT_WAIT_ROUTINE_PATTERNS", "").split("|") if p.strip()]
try:
waits = json.loads(sys.argv[1]).get("waits", [])
except Exception:
waits = []
if not waits or not pats:
print(0)
sys.exit()
for w in waits:
task = w.get("task") or ""
try:
matched = any(re.search(p, task, re.I) for p in pats)
except re.error:
matched = False
if not matched:
print(0)
sys.exit()
print(1)
PYEOF
}
target_plausible_false() { # <node_data_json> -> prints 1 if target_plausible is explicitly false, else 0
# P1 follow-up (2026-10-08): the approval target parser flags garbage
# tokens (e.g. "echo", "true") as target_plausible=false. Implausible
# targets go to the digest instead of paging CRITICAL. Fail-closed:
# missing field, null, non-boolean, or unparseable JSON -> 0 (page it).
python3 - "$1" <<'PYEOF_INNER'
import json, sys
try:
v = json.loads(sys.argv[1]).get("target_plausible")
except Exception:
v = None
print(1 if v is False else 0)
PYEOF_INNER
}
injected() { # cond -> 0 if injected-fail
case ",$INJECT_FAIL," in *,"$1,"*) return 0;; *) return 1;; esac
}
@@ -241,6 +292,7 @@ info = approvals.inspect_node_approvals('$node')
out = {
'has_pending': info.get('has_pending', False),
'target': info.get('target') or info.get('ip') or 'unknown',
'target_plausible': info.get('target_plausible'),
'title': info.get('title') or '',
'waits': info.get('input_waits') or []
}
@@ -262,8 +314,20 @@ print(json.dumps(out))
read -r action fails < <(state_machine "$cond" "$failing" "$BROWSER_APPROVAL_TTL")
case "$action" in
ALERT_FIRST|ALERT_REALERT)
emit_record "ALERT" "$cond" "$detail" "$fails"
echo "$cond" >> "$STATE_DIR/.alerts.tmp"
if [ "$(target_plausible_false "$node_data")" = "1" ]; then
# P1 follow-up (2026-10-08): implausible approval target
# (parser artifact, target_plausible=false) -> digest, don't
# page. kind=DIGEST is ignored by the #lobby relay; the
# triage digest consumer batches these. No .alerts.tmp
# entry, so no box_notify broadcast either — the digest is
# the only output. Missing/unparseable field -> CRITICAL
# (fail-closed; handled inside target_plausible_false).
emit_record "DIGEST" "$cond" "implausible target: $detail" "$fails"
log "$cond implausible target x$fails — digested, not paged"
else
emit_record "ALERT" "$cond" "$detail" "$fails"
echo "$cond" >> "$STATE_DIR/.alerts.tmp"
fi
;;
RECOVERY)
emit_record "RECOVERY" "$cond" "$detail" "$fails"
@@ -308,8 +372,17 @@ if w:
read -r action_in fails_in < <(state_machine "$cond_in" "$failing_in" "$INPUT_WAIT_TTL")
case "$action_in" in
ALERT_FIRST|ALERT_REALERT)
emit_record "ALERT" "$cond_in" "$detail_in" "$fails_in"
echo "$cond_in|$detail_in" >> "$STATE_DIR/.alerts.tmp"
if [ "$(input_wait_routine "$node_data")" = "1" ]; then
# P4 (2026-10-08): routine housekeeping -> digest, don't page.
# kind=DIGEST is ignored by the #lobby relay; the triage
# digest consumer batches these. No .alerts.tmp entry, so
# no targeted DM either — the digest is the only output.
emit_record "DIGEST" "$cond_in" "routine: $detail_in" "$fails_in"
log "$cond_in routine input_wait x$fails_in — digested, not paged"
else
emit_record "ALERT" "$cond_in" "$detail_in" "$fails_in"
echo "$cond_in|$detail_in" >> "$STATE_DIR/.alerts.tmp"
fi
;;
RECOVERY)
emit_record "RECOVERY" "$cond_in" "$detail_in" "$fails_in"
+51 -12
View File
@@ -364,12 +364,11 @@ DM_LOG_FILE = os.path.join(NETVM_ROOT, "dm-log.jsonl")
JOB_LOG_FILE = os.path.join(NETVM_ROOT, "job-log.jsonl")
def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list:
"""Reconstruct active and recent loops from followups.json and dm-log.jsonl.
def _load_loop_candidates() -> dict:
"""Parse followups.json + dm-log.jsonl into a loop_id -> dict map.
Returns a list of dicts:
loop_id, agent, sender, target, purpose, state, sent_at, deadline,
nudges_sent, nudges_allowed, escalate_to, tags, summary
Pure parse phase of reconstruct_loops, extracted so diagnose_breaks and
remediate_breaks can share one parse instead of re-reading the logs.
"""
loops = {} # loop_id -> dict
@@ -499,6 +498,11 @@ def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list:
"source": "dm-log.jsonl",
}
return loops
def _select_loops(loops: dict, limit=50, agent=None, status_filter=None) -> list:
"""Filter/sort/limit a candidate map from _load_loop_candidates."""
# Filter and sort
result = list(loops.values())
if agent:
@@ -517,6 +521,17 @@ def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list:
return result[:limit]
def reconstruct_loops(limit=50, agent=None, status_filter=None) -> list:
"""Reconstruct active and recent loops from followups.json and dm-log.jsonl.
Returns a list of dicts:
loop_id, agent, sender, target, purpose, state, sent_at, deadline,
nudges_sent, nudges_allowed, escalate_to, tags, summary
"""
return _select_loops(_load_loop_candidates(), limit=limit, agent=agent,
status_filter=status_filter)
def get_fleet_loop_health(threshold=None) -> dict:
"""Calculate fleet loop health per agent and overall verdict."""
if threshold is None:
@@ -570,8 +585,14 @@ def get_fleet_loop_health(threshold=None) -> dict:
}
def diagnose_breaks() -> list:
"""Diagnose break taxonomy across intrinsic loops and support services."""
def diagnose_breaks(_fleet_cache=None, _loops_cache=None) -> list:
"""Diagnose break taxonomy across intrinsic loops and support services.
_fleet_cache: optional list; when given, the fleet approval scan result
is appended so callers (remediate_breaks) can reuse it instead of
re-scanning (each scan fans 6 nodes over the full audit log).
_loops_cache: optional list; when given, the parsed loop-candidate map
is appended for the same single-parse sharing."""
import subprocess
breaks = []
@@ -609,7 +630,10 @@ def diagnose_breaks() -> list:
})
# 3. Active follow-up loops check
active_loops = reconstruct_loops(limit=20, status_filter="pending")
_loops_map = _load_loop_candidates()
if _loops_cache is not None:
_loops_cache.append(_loops_map)
active_loops = _select_loops(_loops_map, limit=20, status_filter="pending")
now_ts = time.time()
for l in active_loops:
nudges_sent = l.get("nudges_sent", 0)
@@ -628,6 +652,8 @@ def diagnose_breaks() -> list:
try:
import approvals
fleet_apps = approvals.check_fleet_approvals()
if _fleet_cache is not None:
_fleet_cache.append(fleet_apps)
for app in fleet_apps:
if app.get("has_pending"):
node = app["node"]
@@ -734,7 +760,9 @@ def remediate_breaks(dry_run=False) -> dict:
escalated = []
# 1. Check diagnosed hard breaks first
breaks = diagnose_breaks()
_fleet_cache = []
_loops_cache = []
breaks = diagnose_breaks(_fleet_cache=_fleet_cache, _loops_cache=_loops_cache)
for b in breaks:
if b.get("severity") in ("CRITICAL", "WARNING"):
escalated.append(b)
@@ -753,8 +781,13 @@ def remediate_breaks(dry_run=False) -> dict:
now_iso = datetime.now(timezone.utc).isoformat()
# Build answer map from reconstruct_loops
loops = reconstruct_loops(limit=200)
# Build answer map from reconstruct_loops (reuse diagnose's parse:
# nothing between the parses writes the loop logs in-process, and a
# concurrently landed reply is picked up on the next cycle).
if _loops_cache:
loops = _select_loops(_loops_cache[0], limit=200)
else:
loops = reconstruct_loops(limit=200)
answered_dms = {
l["loop_id"]: l for l in loops if l.get("state") in ("ANSWERED", "CLOSED")
}
@@ -836,7 +869,13 @@ def remediate_breaks(dry_run=False) -> dict:
# Auto-remediate trusted approval blocks
try:
import approvals
fleet_apps = approvals.check_fleet_approvals()
# Reuse the diagnose_breaks scan: nothing between the scans touches
# browser-approval state, and this block only reads it. Fall back to
# a fresh scan if the first one failed.
if _fleet_cache:
fleet_apps = _fleet_cache[0]
else:
fleet_apps = approvals.check_fleet_approvals()
for app in fleet_apps:
if app.get("has_pending") and app.get("is_trusted") and app.get("status") != "KEY_APPROVAL":
node = app["node"]
+138 -11
View File
@@ -41,6 +41,13 @@ NETVM_EXEC = "/home/super/Projects/NetVM/bin/netvm-exec.sh"
JOB_LOG = NETVM_ROOT / "job-log.jsonl"
SIDECHAT_STATE = NETVM_ROOT / "job-sidechats.json"
# Sidechat rotation: persistent reuse_key threads accumulate full history
# and every dispatch re-sends it (cloud context), so a stale thread burns
# full-thread tokens per nod. Cap counted threads by dispatch budget and
# flush uncounted legacy threads past the age cap.
SIDECHAT_MAX_DISPATCHES = 48
SIDECHAT_LEGACY_MAX_AGE_HOURS = 24
def load_sidechat_state():
if SIDECHAT_STATE.exists():
try:
@@ -54,6 +61,40 @@ def save_sidechat_state(state):
tmp.write_text(json.dumps(state, indent=2))
tmp.replace(SIDECHAT_STATE)
def should_rotate_sidechat(record, current_title, now=None,
max_dispatches=SIDECHAT_MAX_DISPATCHES,
legacy_max_age_hours=SIDECHAT_LEGACY_MAX_AGE_HOURS):
"""Decide whether a reused sidechat must rotate to a fresh thread.
Returns (rotate, reason). Rotates when the dispatch budget is spent,
the rendered title moved on (daily {date} templates), or an
uncounted legacy record is past the age cap. Anything unassessable
(plain-UUID records, missing/unparseable age) fails open to reuse.
"""
now = now or datetime.now(timezone.utc)
if not isinstance(record, dict):
return False, "unrecorded"
count = record.get("dispatch_count")
if isinstance(count, int) and count >= max_dispatches:
return True, f"dispatch budget spent ({count}/{max_dispatches})"
stored_title = record.get("title") or ""
ALLOW_SIDECHAT_TITLE_ROTATION = False
if ALLOW_SIDECHAT_TITLE_ROTATION and stored_title and current_title and stored_title != current_title:
return True, f"title rolled over ({stored_title} -> {current_title})"
if count is None:
created = record.get("created_at")
if created:
try:
age_h = (now - datetime.fromisoformat(
str(created).replace("Z", "+00:00"))).total_seconds() / 3600
except Exception:
return False, "unparseable age"
if age_h > legacy_max_age_hours:
return True, (f"predates counting, age {age_h:.0f}h "
f"over {legacy_max_age_hours}h cap")
return False, "within budget"
def extract_uuid(url):
m = re.search(r"/thread/([0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12})", url or "")
return m.group(1) if m else None
@@ -88,6 +129,65 @@ try:
except ImportError:
HAS_RATE_LIMITER = False
# Dispatch backpressure (2026-10-09): skip jobs for frozen agents instead of
# piling input-waits onto them. See tests/test_dispatch_hold.py.
DISPATCH_HOLD_FILE = JOBS_DIR / "dispatch-hold.json"
HOLD_WAIT_THRESHOLD = 3
HOLD_WAIT_WINDOW_MIN = 60
def dispatch_hold_reason(agent, now=None, hold_path=None, job_log_path=None):
# Hold reason if dispatch to agent must be skipped, else None.
# Explicit operator holds win; otherwise auto-hold after repeated waits.
from datetime import timedelta
now = now or datetime.now(timezone.utc)
try:
with open(hold_path or DISPATCH_HOLD_FILE) as f:
holds = json.load(f)
except (OSError, ValueError):
holds = {}
entry = holds.get(agent) if isinstance(holds, dict) else None
if isinstance(entry, dict):
until = entry.get("until")
if until:
try:
exp = datetime.fromisoformat(until)
if exp.tzinfo is None:
exp = exp.replace(tzinfo=timezone.utc)
except ValueError:
exp = None
if exp is not None and exp <= now:
entry = None
if entry is not None:
return "explicit hold (%s)" % entry.get("reason", "operator")
try:
cutoff = now - timedelta(minutes=HOLD_WAIT_WINDOW_MIN)
n = 0
with open(job_log_path or JOB_LOG) as f:
for line in f:
try:
r = json.loads(line)
except ValueError:
continue
if r.get("type") != "job_dispatch_agent_input_wait":
continue
if r.get("agent") != agent:
continue
try:
ts = datetime.fromisoformat(r.get("ts", ""))
except ValueError:
continue
if ts.tzinfo is None:
ts = ts.replace(tzinfo=timezone.utc)
if ts >= cutoff:
n += 1
if n >= HOLD_WAIT_THRESHOLD:
return "auto-hold (%d input-waits in last %dm)" % (n, HOLD_WAIT_WINDOW_MIN)
except OSError:
pass
return None
def log_event(event_type, data):
"""Append event to job-log.jsonl"""
entry = {
@@ -395,6 +495,13 @@ def main():
# Load job
job = load_job(job_name)
# Backpressure: skip frozen agents before arming follow-ups or sending.
_hold = dispatch_hold_reason(job.get("agent"))
if _hold:
print("Held: job %s for %s skipped (%s)." % (job_name, job.get("agent"), _hold), file=sys.stderr)
log_event("job_dispatch_held", {"job_name": job_name, "agent": job.get("agent"), "reason": _hold})
sys.exit(0)
# Generate job_id
job_id = f"{job_name}-{datetime.now(timezone.utc).strftime('%Y%m%d-%H%M%S')}-{uuid.uuid4().hex[:8]}"
@@ -461,19 +568,33 @@ def main():
# Check if reuse_key exists in job-sidechats.json and thread is still alive
sc_state = load_sidechat_state()
reused_uuid = None
rotated_from = None
if reuse_key and reuse_key in sc_state:
val = sc_state[reuse_key]
cand_uuid = val.get("thread_uuid") if isinstance(val, dict) else val
if cand_uuid:
try:
import muse_hybrid
threads, err = muse_hybrid.get_threads(agent)
if not err and threads:
thread_ids = [t.get("session_id") for t in threads]
if cand_uuid in thread_ids:
reused_uuid = cand_uuid
except Exception:
pass
rotate, reason = should_rotate_sidechat(val, sc_name)
if rotate:
print(f"Rotating sidechat '{reuse_key}': {reason}")
log_event("job_sidechat_rotate", {
"job_name": job_name, "job_id": job_id,
"reuse_key": reuse_key, "old_thread": cand_uuid,
"reason": reason,
})
rotated_from = cand_uuid
else:
try:
import muse_hybrid
threads, err = muse_hybrid.get_threads(agent)
if not err and threads:
thread_ids = [t.get("session_id") for t in threads]
if cand_uuid in thread_ids:
reused_uuid = cand_uuid
except Exception:
pass
if reused_uuid and isinstance(val, dict):
val["dispatch_count"] = val.get("dispatch_count", 0) + 1
save_sidechat_state(sc_state)
if reused_uuid:
target = reused_uuid
@@ -487,13 +608,19 @@ def main():
new_uuid = res.get("session_id")
key_to_save = reuse_key or sc_name
is_persistent = bool(reuse_key)
sc_state[key_to_save] = {
new_record = {
"thread_uuid": new_uuid,
"agent": agent,
"title": channel_title,
"type": "persistent" if is_persistent else "ephemeral",
"created_at": datetime.now(timezone.utc).isoformat()
"created_at": datetime.now(timezone.utc).isoformat(),
"dispatch_count": 1,
}
if rotated_from:
new_record["rotated_from"] = rotated_from
new_record["rotated_at"] = datetime.now(
timezone.utc).isoformat()
sc_state[key_to_save] = new_record
save_sidechat_state(sc_state)
target = new_uuid
print(f"Spawned new sidechat channel '{channel_title}' ({new_uuid}) for {agent}")
+3 -2
View File
@@ -281,8 +281,9 @@ def generate_preservation_advisory(
tips = []
pct = weekly_used_pct or 0
if pct >= 95 or "0 tokens left" in extra_tokens_remaining:
return "CRITICAL: Quota exhausted. Do NOT send chat messages. Salvage via 'box onboard start <new_node> --for %s'." % node
is_bonus_empty = ("0 tokens left" in extra_tokens_remaining) or (not extra_tokens_remaining)
if pct >= 95 and is_bonus_empty:
return "CRITICAL: Quota exhausted. Salvage via 'box onboard start <new_node> --for %s'." % node
if pct >= 70:
tips.append("Quota > 70%%: Cease prose chatter; offload tasks to background tmux workers.")
+22 -1
View File
@@ -117,6 +117,27 @@ def ev(ws, expr, await_p=False):
print(f"CDP evaluate failed: {type(e).__name__}: {e}", file=sys.stderr)
return None
def _is_valid_ipv4(ip: str) -> bool:
"""Strict IPv4 validation: four octets, each 0-255, no leading zeros.
P1 fix (2026-10-08): the old \d{1,3} pattern matched invalid IPs like
999.999.999.999 and version strings. Only strict IPv4 passes.
"""
if not ip or not isinstance(ip, str):
return False
parts = ip.split(".")
if len(parts) != 4:
return False
try:
return all(
0 <= int(part) <= 255 and part == str(int(part))
for part in parts
)
except ValueError:
return False
def check_approvals(ws):
"""
Check for browser permission dialogs.
@@ -175,7 +196,7 @@ def check_approvals(ws):
for d in dialogs:
# Extract IP if present
import re
ips = re.findall(r'\b\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}\b', d)
ips = [ip for ip in re.findall(r'\b\d{1,3}\.\d{1,3}\.\d{1,3}\.\d{1,3}\b', d) if _is_valid_ipv4(ip)]
# Check trust: if IP present, must be in TRUSTED_IPS; if no IP, untrusted approval dialog
if ips:
is_trusted = any(ip in TRUSTED_IPS for ip in ips)
+442 -22
View File
@@ -16,6 +16,7 @@ Dual-mode interface:
- Job Scheduler & Dispatch trigger
- Background Tmux sessions & Swarm worker monitor
- Live Event & DM log tailer
- Container SSH tunnel health & tmux pop-out dialer
"""
import sys
@@ -27,6 +28,7 @@ import threading
import subprocess
import hashlib
import select
import shlex
import signal
import textwrap
import urllib.request
@@ -169,6 +171,85 @@ def format_recency(ts: float) -> str:
return "never"
# ---------------------------------------------------------------------------
# SSH / Container Tunnel Management Subsystem
# ---------------------------------------------------------------------------
SSH_JUMP_HOST = os.environ.get("SSH_JUMP_HOST", "34.139.37.135")
SSH_OPERATOR_USER = os.environ.get("OPERATOR_USER", "super")
SSH_IDENTITY_FILE = os.environ.get("SSH_IDENTITY_FILE", "")
try:
from agent_md import TUNNEL_PORTS as SSH_TUNNEL_PORTS
except Exception:
SSH_TUNNEL_PORTS = {
"muse-main": {"port": 2224, "terminal": 7681, "user": "muse"},
"muse": {"port": 2225, "terminal": 7682, "user": "hatch"},
"646": {"port": 2226, "terminal": 7683, "user": "hatch"},
"pip": {"port": 2227, "terminal": 7684, "user": "hatch"},
"opm": {"port": 2228, "terminal": 7685, "user": "hatch"},
"def": {"port": 2229, "terminal": 7686, "user": "hatch"},
"dev": {"port": 2230, "terminal": 7687, "user": "hatch"},
}
def build_ssh_dial_command(account: str, port=None, user=None, jump_host=None,
operator_user=None, identity_file=None,
ssh_options=None, remote_command=None) -> list:
"""Build the jump-host dial argv for an agent container.
ssh_options are inserted before the destination; remote_command (str or
list) is appended after it for non-interactive probes.
"""
info = SSH_TUNNEL_PORTS.get(account, {})
port = port or info.get("port")
user = user or info.get("user", "hatch")
jump_host = jump_host or SSH_JUMP_HOST
operator_user = operator_user or SSH_OPERATOR_USER
if identity_file is None:
identity_file = SSH_IDENTITY_FILE
cmd = ["ssh", "-o", "StrictHostKeyChecking=no"]
if identity_file:
cmd += ["-o", "IdentitiesOnly=yes", "-i", identity_file]
if ssh_options:
cmd += list(ssh_options)
cmd += ["-J", f"{operator_user}@{jump_host}", "-p", str(port), f"{user}@localhost"]
if remote_command:
cmd += [remote_command] if isinstance(remote_command, str) else list(remote_command)
return cmd
def build_ssh_dial_string(account: str, **kwargs) -> str:
"""Shell-quoted dial command for display, clipboard copy, and pop-out."""
return " ".join(shlex.quote(p) for p in build_ssh_dial_command(account, **kwargs))
def build_ssh_popout_shell(account: str, **kwargs) -> str:
"""Interactive shell line for the pop-out window: ssh, then keep a shell."""
dial = build_ssh_dial_string(account, **kwargs)
return f"{dial}; echo '[ssh exited ($?) — window kept open, exit to close]'; exec \"${{SHELL:-/bin/bash}}\""
def build_tmux_popout_command(label: str, shell_command: str, socket_path: str = None) -> list:
"""Build `tmux new-window` argv opening shell_command in a fresh window."""
safe_label = re.sub(r"[^A-Za-z0-9_.-]", "-", label)[:32] or "ssh"
cmd = ["tmux"]
if socket_path:
cmd += ["-S", socket_path]
return cmd + ["new-window", "-n", safe_label, shell_command]
def ssh_row_order(nodes: list, extra_accounts=()) -> list:
"""Fleet nodes first, then any extra tunnel accounts (e.g. muse-main)."""
rows = list(nodes)
for acct in extra_accounts:
if acct not in rows:
rows.append(acct)
for acct in SSH_TUNNEL_PORTS:
if acct not in rows:
rows.append(acct)
return rows
# ---------------------------------------------------------------------------
# Prompt & Skill Library Subsystem
# ---------------------------------------------------------------------------
@@ -337,6 +418,31 @@ class PromptManager:
return False
# ---------------------------------------------------------------------------
# Box Mode Tab Bar (single source of truth for renderer + click handler)
# ---------------------------------------------------------------------------
BOX_TABS = [
"1: Agent Chat",
"2: Fleet Status",
"3: Approvals",
"4: Jobs Scheduler",
"5: Tmux / Swarms",
"6: DM Logs",
"7: SSH / Boxes",
]
def box_tab_bounds(tabs=None, x: int = 0) -> list:
"""Clickable x-ranges for the Box tab bar, mirroring _render_box_tabs."""
bounds = []
cur_x = x + 1
for tab_name in (tabs if tabs is not None else BOX_TABS):
label = f" [{tab_name}] "
bounds.append((cur_x, cur_x + len(label) - 1))
cur_x += len(label) + 1
return bounds
# ---------------------------------------------------------------------------
# Data Layer & Async Poller
# ---------------------------------------------------------------------------
@@ -368,6 +474,13 @@ class FleetDataManager:
self.tmux_cache = []
self.dm_logs_cache = []
# SSH / container tunnel health (Box tab 7)
self.ssh_cache = {} # account -> health dict from ssh-check + state_since
self.ssh_jump_reachable = None # None = never checked
self.ssh_checked_at = 0.0
self.ssh_check_latency_ms = None
self.ssh_check_error = ""
# Interaction ranking: node -> float timestamp of last true input / chat [insert]
self.agent_interactions = {n: 0.0 for n in self.nodes}
self._load_agent_interactions()
@@ -405,6 +518,8 @@ class FleetDataManager:
self.preload_priority_chats(self.active_node, sidechat_limit=0)
self.poller_thread = threading.Thread(target=self._worker_loop, daemon=True)
self.poller_thread.start()
# First SSH sweep in background so Box tab 7 is warm on open
threading.Thread(target=self._fetch_ssh_health, daemon=True).start()
else:
self.poller_thread = None
@@ -888,6 +1003,7 @@ class FleetDataManager:
last_med = 0.0
last_slow = 0.0
last_fleet_approvals = 0.0
last_ssh = 0.0
# Initial fetch of active node main chat only
with self.lock:
@@ -950,6 +1066,11 @@ class FleetDataManager:
self._fetch_dm_logs()
last_slow = now
# 5. SSH tunnel health (every 30s): single VM-side sweep
if (now - last_ssh >= 30.0):
self._fetch_ssh_health()
last_ssh = now
# Sleep in short increments to allow prompt wakeup on user actions
for _ in range(10):
if not self.running or (hasattr(self, 'user_poll_trigger') and self.user_poll_trigger.is_set()):
@@ -1181,6 +1302,68 @@ class FleetDataManager:
except Exception:
pass
def _apply_ssh_check_result(self, data: dict, now: float = None):
"""Merge one ssh-check payload into ssh_cache with flap tracking.
state_since records the last (ssh_up, term_up) transition per
account so the SSH view can show uptime/downtime durations.
"""
now = now if now is not None else time.time()
if not isinstance(data, dict):
return
accounts = data.get("accounts", {})
if not isinstance(accounts, dict):
accounts = {}
with self.lock:
self.ssh_jump_reachable = data.get("jump_reachable")
self.ssh_checked_at = now
self.ssh_check_latency_ms = data.get("latency_ms")
self.ssh_check_error = "" if data.get("jump_reachable") else str(data.get("error", ""))
for acct, info in accounts.items():
if not isinstance(info, dict):
continue
prev = self.ssh_cache.get(acct, {})
entry = dict(info)
prev_state = (prev.get("ssh_up"), prev.get("term_up"))
new_state = (entry.get("ssh_up"), entry.get("term_up"))
if prev_state != new_state or "state_since" not in prev:
entry["state_since"] = now
entry["prev_ssh_up"] = prev.get("ssh_up")
else:
entry["state_since"] = prev.get("state_since", now)
entry["prev_ssh_up"] = prev.get("prev_ssh_up")
self.ssh_cache[acct] = entry
def _fetch_ssh_health(self):
"""Run box-ctl ssh-check (single VM-side sweep) and merge results."""
try:
cmd = ["python3", str(BIN_DIR / "box-ctl.py"), "ssh-check"]
res = subprocess.run(cmd, capture_output=True, text=True, timeout=30)
if res.returncode == 0:
try:
data = json.loads(res.stdout)
except Exception:
return
if isinstance(data, dict) and data.get("ok"):
self._apply_ssh_check_result(data)
except Exception:
pass
def probe_container_uptime(self, account: str) -> tuple[bool, str]:
"""On-demand end-to-end probe: run `uptime` inside the container."""
if account not in SSH_TUNNEL_PORTS:
return False, f"No tunnel registered for '{account}'"
cmd = build_ssh_dial_command(
account,
ssh_options=["-o", "BatchMode=yes", "-o", "ConnectTimeout=12"],
remote_command="uptime",
)
rc, stdout, stderr = run_command_isolated(cmd, timeout=25.0)
if rc == 0 and (stdout or "").strip():
return True, stdout.strip()
err_lines = (stderr or stdout or f"exit {rc}").strip().splitlines()
return False, (err_lines[-1] if err_lines else f"exit {rc}")[:200]
def _fetch_dm_logs(self):
log_path = REPO_ROOT / "dm-log.jsonl"
if not log_path.exists():
@@ -1362,10 +1545,13 @@ class FleetDataManager:
class MuseTUI:
"""Full-terminal curses application supporting Muse and Box operational modes."""
def __init__(self, stdscr, initial_mode="muse", initial_node=None, initial_thread=None):
def __init__(self, stdscr, initial_mode="muse", initial_node=None, initial_thread=None, initial_tab=None):
self.stdscr = stdscr
self.mode = initial_mode # "muse" or "box"
self.box_tab = 0 # 0: Chat, 1: Fleet, 2: Approvals, 3: Jobs, 4: Tmux, 5: Logs
self.box_tab = 0 # 0: Chat, 1: Fleet, 2: Approvals, 3: Jobs, 4: Tmux, 5: Logs, 6: SSH
if initial_mode == "box" and initial_tab is not None:
tab_map = {"chat": 0, "fleet": 1, "approvals": 2, "jobs": 3, "tmux": 4, "logs": 5, "ssh": 6}
self.box_tab = tab_map.get(str(initial_tab).lower(), 0)
self.data = FleetDataManager()
# Selection state: default to top of sorted list (highest unread / most recently interacted)
@@ -1408,6 +1594,8 @@ class MuseTUI:
self.jobs_sel_idx = 0
self.jobs_scroll_start = 0
self.tmux_sel_idx = 0
self.ssh_sel_idx = 0
self.ssh_scroll_idx = 0
self.table_scroll_idx = 0
# Input buffer
@@ -1674,6 +1862,8 @@ class MuseTUI:
self._render_tmux_view(content_y, 0, content_h, w)
elif self.box_tab == 5:
self._render_dm_logs_view(content_y, 0, content_h, w)
elif self.box_tab == 6:
self._render_ssh_view(content_y, 0, content_h, w)
# 4. Bottom Input Bar & Toast
self._render_bottom_bar(h - bottom_bar_h, 0, bottom_bar_h, w)
@@ -1764,14 +1954,7 @@ class MuseTUI:
self.safe_addstr(self.stdscr, y, w - len(clock) - 2, clock, self._attr("header"))
def _render_box_tabs(self, y: int, x: int, w: int):
tabs = [
"1: Agent Chat",
"2: Fleet Status",
"3: Approvals",
"4: Jobs Scheduler",
"5: Tmux / Swarms",
"6: DM Logs",
]
tabs = BOX_TABS
self.safe_addstr(self.stdscr, y, x, " " * w, self._attr("dim"))
cur_x = x + 1
for idx, tab_name in enumerate(tabs):
@@ -2839,6 +3022,106 @@ class MuseTUI:
self.safe_addstr(self.stdscr, row_y, x + 2, line_str, color)
row_y += 1
def get_ssh_rows(self) -> list:
"""SSH view row order: fleet nodes first, then extra tunnel accounts."""
with self.data.lock:
extras = list(self.data.ssh_cache.keys())
nodes = list(self.data.nodes)
return ssh_row_order(nodes, extras)
def _render_ssh_view(self, y: int, x: int, h: int, w: int):
self.safe_addstr(self.stdscr, y, x + 1, f"CONTAINER SSH TUNNEL HEALTH (jump: {SSH_OPERATOR_USER}@{SSH_JUMP_HOST})", self._attr("bold"))
with self.data.lock:
jump = self.data.ssh_jump_reachable
checked_at = self.data.ssh_checked_at
sweep_ms = self.data.ssh_check_latency_ms
check_err = self.data.ssh_check_error
ssh_cache = dict(self.data.ssh_cache)
if jump is None:
jump_txt, jump_attr = "sweep pending…", self._attr("dim")
elif jump:
jump_txt = f"jump OK (sweep {sweep_ms}ms, checked {format_recency(checked_at)})"
jump_attr = self._attr("success")
else:
jump_txt = f"jump UNREACHABLE ({(check_err or 'unknown')[:w - 24]})"
jump_attr = self._attr("danger")
self.safe_addstr(self.stdscr, y + 1, x + 1, jump_txt[:w - 2], jump_attr)
self.safe_addstr(self.stdscr, y + 2, x + 1, " ACCOUNT SSH PORT SSH STATE LAT SSH BANNER / HOSTKEY TERM PORT TERM STATE SINCE", self._attr("dim"))
self.safe_addstr(self.stdscr, y + 3, x + 1, "─" * (w - 2), self._attr("dim"))
rows = self.get_ssh_rows()
if self.ssh_sel_idx >= len(rows):
self.ssh_sel_idx = max(0, len(rows) - 1)
visible_rows = max(3, h - 6)
scroll_start = getattr(self, "ssh_scroll_idx", 0)
if self.ssh_sel_idx < scroll_start:
scroll_start = self.ssh_sel_idx
elif self.ssh_sel_idx >= scroll_start + visible_rows:
scroll_start = self.ssh_sel_idx - visible_rows + 1
scroll_start = max(0, min(scroll_start, max(0, len(rows) - visible_rows)))
self.ssh_scroll_idx = scroll_start
row_y = y + 4
for row_i in range(visible_rows):
idx = scroll_start + row_i
if idx >= len(rows):
break
acct = rows[idx]
is_sel = (idx == self.ssh_sel_idx)
row_attr = self._attr("selected") if is_sel else self._attr("normal")
info = SSH_TUNNEL_PORTS.get(acct, {})
health = ssh_cache.get(acct, {})
sport = info.get("port", "?")
tport = info.get("terminal", "?")
ssh_up = health.get("ssh_up")
term_up = health.get("term_up")
if ssh_up is True:
ssh_txt, ssh_attr = "UP ", self._attr("success")
elif ssh_up is False:
ssh_txt, ssh_attr = "DOWN", self._attr("danger")
else:
ssh_txt, ssh_attr = "?? ", self._attr("dim")
if term_up is True:
term_txt, term_attr = "UP ", self._attr("success")
elif term_up is False:
term_txt, term_attr = "DOWN", self._attr("danger")
else:
term_txt, term_attr = "?? ", self._attr("dim")
if is_sel:
ssh_attr = row_attr
term_attr = row_attr
lat = health.get("ssh_latency_ms")
lat_txt = f"{lat}ms" if lat is not None else "--"
banner = (health.get("ssh_banner") or health.get("term_http") or "-").strip() or "-"
since_ts = health.get("state_since", 0.0)
if ssh_up is None and term_up is None:
since_txt = "never checked" if jump is None else "unknown"
else:
since_txt = f"{'up' if ssh_up else 'down'} {format_recency(since_ts)}"
head = "▶ " if is_sel else " "
name_attr = row_attr if is_sel else self._attr("bold")
self.safe_addstr(self.stdscr, row_y, x + 1, f"{head}{acct:<10}"[:12], name_attr)
self.safe_addstr(self.stdscr, row_y, x + 13, f":{sport:<8}", row_attr if is_sel else self._attr("dim"))
self.safe_addstr(self.stdscr, row_y, x + 23, ssh_txt, ssh_attr)
self.safe_addstr(self.stdscr, row_y, x + 32, f"{lat_txt:<7}", row_attr if is_sel else self._attr("dim"))
self.safe_addstr(self.stdscr, row_y, x + 40, banner[:26].ljust(26), row_attr if is_sel else self._attr("normal"))
self.safe_addstr(self.stdscr, row_y, x + 67, f":{tport:<8}", row_attr if is_sel else self._attr("dim"))
self.safe_addstr(self.stdscr, row_y, x + 77, term_txt, term_attr)
self.safe_addstr(self.stdscr, row_y, x + 83, since_txt[:w - 84 - 8], row_attr if is_sel else self._attr("dim"))
self.safe_addstr(self.stdscr, row_y, max(x + 90, w - 8), "[SSH]", self._attr("wo_badge") if is_sel else self._attr("dim"))
row_y += 1
hint_y = y + h - 1
sel_acct = rows[self.ssh_sel_idx].upper() if rows else "-"
hints = f"Selected: [{sel_acct}] [Enter/s]: SSH pop-out (tmux) [c]: Copy dial [u]: Container uptime [r]: Refresh [j/k]: Nav"
self.safe_addstr(self.stdscr, hint_y, x + 1, hints[:w - 2], self._attr("dim"))
# -----------------------------------------------------------------------
# Chat History Sends Search & Prompts Subsystem
# -----------------------------------------------------------------------
@@ -3210,6 +3493,8 @@ class MuseTUI:
("[w] or [/wo]", "Compose and cryptographically sign a Work Order"),
("[a] or [F2]", "Open Approvals Resolution Drawer (Allow, Always, Deny)"),
("[F5] or [m]", "Toggle between Muse Chat TUI and Box Fleet Command TUI"),
("[1]-[7] (Box)", "Switch Box tabs: Chat/Fleet/Approvals/Jobs/Tmux/Logs/SSH"),
("[Tab 7: SSH]", "Enter/s: tmux pop-out c: copy dial u: container uptime r: refresh"),
("[g] / [G] / [Home/End]", "Jump to oldest message / follow live latest message"),
("[q]", "Quit TUI (in NORMAL mode)"),
]
@@ -4247,6 +4532,30 @@ class MuseTUI:
self.toggle_transcript_style()
return True
# Box mode Fast Actions: Tab 6 (SSH) — placed before the 's'/'c'/'y'
# globals below so SSH keys win on this tab.
if self.mode == "box" and self.box_tab == 6:
if ch in (curses.KEY_ENTER, 10, 13):
self._ssh_popout_selected()
return True
elif ch in (ord('s'), ord('S')):
self._ssh_popout_selected()
return True
elif ch in (ord('c'), ord('C')):
self._ssh_copy_dial_selected()
return True
elif ch in (ord('u'), ord('U')):
rows = self.get_ssh_rows()
if rows and 0 <= self.ssh_sel_idx < len(rows):
acct = rows[self.ssh_sel_idx]
threading.Thread(target=self._async_ssh_uptime, args=(acct,), daemon=True).start()
self.set_toast(f"Probing container uptime on {acct}...", "info")
return True
elif ch in (ord('r'), ord('R')):
threading.Thread(target=self.data._fetch_ssh_health, daemon=True).start()
self.set_toast("Probing SSH tunnels via jump host...", "info")
return True
# Open Context Menu for active sidebar thread, fleet agent, or message: 'x', 'c', or Space
if ch in (ord('x'), ord('X'), ord('c'), ord('C'), ord(' ')) and not (self.mode == "box" and self.box_tab == 2):
if self.focus_pane == "fleet":
@@ -4564,7 +4873,7 @@ class MuseTUI:
threading.Thread(target=self._async_kill_tmux, args=(sess_name,), daemon=True).start()
return True
elif ch in (ord('r'), ord('R')):
threading.Thread(target=self.data._fetch_tmux, daemon=True).start()
threading.Thread(target=self.data._fetch_tmux_sessions, daemon=True).start()
self.set_toast("Refreshed tmux background sessions.", "info")
return True
@@ -4613,29 +4922,29 @@ class MuseTUI:
self.set_toast(f"Switched to agent: {node.upper()}", "success")
return True
# Box mode tab selection: '1' - '6' (when in box mode, except 1-3 on Tab 2)
# Box mode tab selection: '1' - '7' (when in box mode, except 1-3 on Tab 2)
if self.mode == "box":
if self.box_tab == 2:
if ord('4') <= ch <= ord('6'):
if ord('4') <= ch <= ord('7'):
self.box_tab = ch - ord('1')
return True
elif ord('1') <= ch <= ord('6'):
elif ord('1') <= ch <= ord('7'):
self.box_tab = ch - ord('1')
return True
# Box mode tab navigation (when not in Chat tab 0): '[' / ']' / Tab / Shift-Tab
if self.mode == "box" and self.box_tab != 0:
if ch in (ord('['), curses.KEY_LEFT):
self.box_tab = (self.box_tab - 1) % 6
self.box_tab = (self.box_tab - 1) % 7
return True
elif ch in (ord(']'), curses.KEY_RIGHT):
self.box_tab = (self.box_tab + 1) % 6
self.box_tab = (self.box_tab + 1) % 7
return True
elif ch == ord('\t'):
self.box_tab = (self.box_tab + 1) % 6
self.box_tab = (self.box_tab + 1) % 7
return True
elif ch == curses.KEY_BTAB:
self.box_tab = (self.box_tab - 1) % 6
self.box_tab = (self.box_tab - 1) % 7
return True
# Muse View / Chat Tab: Direct Conversation Cycling & Pane Switching
@@ -4722,6 +5031,8 @@ class MuseTUI:
self.jobs_sel_idx = max(0, self.jobs_sel_idx - 1)
elif self.box_tab == 4:
self.tmux_sel_idx = max(0, self.tmux_sel_idx - 1)
elif self.box_tab == 6:
self.ssh_sel_idx = max(0, self.ssh_sel_idx - 1)
else:
self.table_scroll_idx = max(0, self.table_scroll_idx - 1)
return True
@@ -4764,6 +5075,10 @@ class MuseTUI:
sessions = list(self.data.tmux_cache)
if sessions:
self.tmux_sel_idx = min(len(sessions) - 1, self.tmux_sel_idx + 1)
elif self.box_tab == 6:
rows = self.get_ssh_rows()
if rows:
self.ssh_sel_idx = min(len(rows) - 1, self.ssh_sel_idx + 1)
else:
self.table_scroll_idx += 1
return True
@@ -4784,6 +5099,8 @@ class MuseTUI:
self.jobs_sel_idx = max(0, self.jobs_sel_idx - 5)
elif self.box_tab == 4:
self.tmux_sel_idx = max(0, self.tmux_sel_idx - 5)
elif self.box_tab == 6:
self.ssh_sel_idx = max(0, self.ssh_sel_idx - 5)
else:
self.table_scroll_idx = max(0, self.table_scroll_idx - 5)
return True
@@ -4812,6 +5129,10 @@ class MuseTUI:
sessions = list(self.data.tmux_cache)
if sessions:
self.tmux_sel_idx = min(len(sessions) - 1, self.tmux_sel_idx + 5)
elif self.box_tab == 6:
rows = self.get_ssh_rows()
if rows:
self.ssh_sel_idx = min(len(rows) - 1, self.ssh_sel_idx + 5)
else:
self.table_scroll_idx += 5
return True
@@ -4837,6 +5158,8 @@ class MuseTUI:
self.jobs_sel_idx = 0
elif self.box_tab == 4:
self.tmux_sel_idx = 0
elif self.box_tab == 6:
self.ssh_sel_idx = 0
else:
self.table_scroll_idx = 0
return True
@@ -4876,6 +5199,10 @@ class MuseTUI:
sessions = list(self.data.tmux_cache)
if sessions:
self.tmux_sel_idx = max(0, len(sessions) - 1)
elif self.box_tab == 6:
rows = self.get_ssh_rows()
if rows:
self.ssh_sel_idx = max(0, len(rows) - 1)
else:
self.table_scroll_idx = max(0, len(self.data.nodes) - 5)
return True
@@ -4971,6 +5298,8 @@ class MuseTUI:
self.jobs_sel_idx = max(0, self.jobs_sel_idx - 1)
elif self.box_tab == 4:
self.tmux_sel_idx = max(0, self.tmux_sel_idx - 1)
elif self.box_tab == 6:
self.ssh_sel_idx = max(0, self.ssh_sel_idx - 1)
else:
self.table_scroll_idx = max(0, self.table_scroll_idx - 1)
elif mx >= sidebar_w:
@@ -5020,6 +5349,10 @@ class MuseTUI:
sessions = list(self.data.tmux_cache)
if sessions:
self.tmux_sel_idx = min(len(sessions) - 1, self.tmux_sel_idx + 1)
elif self.box_tab == 6:
rows = self.get_ssh_rows()
if rows:
self.ssh_sel_idx = min(len(rows) - 1, self.ssh_sel_idx + 1)
else:
self.table_scroll_idx += 1
elif mx >= sidebar_w:
@@ -5273,8 +5606,7 @@ class MuseTUI:
# Box Tabs click (my == 1 and self.mode == "box")
if my == 1 and self.mode == "box":
tab_bounds = [(1, 18), (19, 37), (38, 53), (54, 74), (75, 94), (95, 108)]
for idx, (start, end) in enumerate(tab_bounds):
for idx, (start, end) in enumerate(box_tab_bounds()):
if start <= mx <= end:
self.box_tab = idx
return True
@@ -5669,7 +6001,7 @@ class MuseTUI:
self.focus_pane = "transcript"
return True
# Box View clicks (Tabs 1, 2, 3, 4)
# Box View clicks (Tabs 1, 2, 3, 4, 6)
elif self.mode == "box":
self.editor_mode = "NORMAL"
row_idx = my - (content_y + 3)
@@ -5763,6 +6095,21 @@ class MuseTUI:
self.set_toast(f"Selected session '{sess_name}'. Click [Kill] or [Attach].", "info")
return True
elif self.box_tab == 6:
# Tab 6: SSH / Boxes (rows start one line lower: jump-status line)
rows = self.get_ssh_rows()
scroll_start = getattr(self, "ssh_scroll_idx", 0)
clicked_idx = scroll_start + row_idx - 1
if 0 <= clicked_idx < len(rows):
self.ssh_sel_idx = clicked_idx
if mx >= w - 8:
# [SSH] pop-out
self._ssh_popout_selected()
else:
acct = rows[clicked_idx]
self.set_toast(f"Selected [{acct.upper()}]. Click [SSH] or press Enter to pop out.", "info")
return True
return True
def _handle_insert_key(self, ch: int) -> bool:
@@ -6348,6 +6695,77 @@ class MuseTUI:
else:
self.set_toast(f"Kill failed: {msg[:40]}", "error")
def _ssh_popout_selected(self):
"""Open the selected container SSH session in a new tmux window."""
rows = self.get_ssh_rows()
if not rows or not (0 <= self.ssh_sel_idx < len(rows)):
self.set_toast("No SSH row selected.", "warn")
return
acct = rows[self.ssh_sel_idx]
if acct not in SSH_TUNNEL_PORTS:
self.set_toast(f"No tunnel registered for '{acct}'.", "warn")
return
if not os.environ.get("TMUX"):
# Refuse to spawn into an invisible server: new-window would
# create a detached server the operator cannot see.
try:
probe = subprocess.run(["tmux", "ls"], capture_output=True, text=True, timeout=3)
server_up = probe.returncode == 0
except Exception:
server_up = False
if not server_up:
copy_to_clipboard(build_ssh_dial_string(acct))
self.set_toast("Not inside tmux and no server running; dial copied to clipboard.", "warn")
return
shell_cmd = build_ssh_popout_shell(acct)
pop_cmd = build_tmux_popout_command(f"ssh-{acct}", shell_cmd)
try:
curses.def_prog_mode()
curses.endwin()
try:
res = subprocess.run(pop_cmd, capture_output=True, text=True, timeout=5)
finally:
try:
curses.reset_prog_mode()
self.stdscr.refresh()
except Exception:
pass
self.need_full_redraw = True
if res.returncode == 0:
self.set_toast(f"Opened SSH to {acct} in tmux window ssh-{acct}.", "success")
else:
copy_to_clipboard(build_ssh_dial_string(acct))
err = (res.stderr or "").strip().splitlines()
hint = err[-1][:60] if err else f"exit {res.returncode}"
self.set_toast(f"tmux pop-out failed ({hint}); dial copied.", "error")
except Exception as e:
copy_to_clipboard(build_ssh_dial_string(acct))
self.set_toast(f"Pop-out failed ({e}); dial copied to clipboard.", "error")
def _ssh_copy_dial_selected(self):
"""Copy the selected container's dial command to the clipboard."""
rows = self.get_ssh_rows()
if not rows or not (0 <= self.ssh_sel_idx < len(rows)):
self.set_toast("No SSH row selected.", "warn")
return
acct = rows[self.ssh_sel_idx]
if acct not in SSH_TUNNEL_PORTS:
self.set_toast(f"No tunnel registered for '{acct}'.", "warn")
return
dial = build_ssh_dial_string(acct)
if copy_to_clipboard(dial):
self.set_toast(f"Copied dial for {acct}: {dial[:80]}", "success")
else:
self.set_toast(f"Dial for {acct}: {dial}", "info")
def _async_ssh_uptime(self, account: str):
ok, out_text = self.data.probe_container_uptime(account)
if ok:
first = (out_text or "").strip().splitlines()
self.set_toast(f"{account} uptime: {first[0][:90]}" if first else f"{account}: uptime probe empty.", "success" if first else "warn")
else:
self.set_toast(f"{account} uptime failed: {out_text[:90]}", "error")
def _execute_chat_action(self, action_id: str):
"""Execute selected contextual action on the targeted chat thread."""
chat_info = getattr(self, "context_chat", None) or {}
@@ -7039,6 +7457,7 @@ def main():
parser.add_argument("--account", "-a", choices=VALID_NODES, default=None, help="Initial agent account (default: top of list)")
parser.add_argument("--thread", "-t", help="Initial thread UUID to open")
parser.add_argument("--mode", "-m", choices=["muse", "box"], default="muse", help="TUI mode (default: muse)")
parser.add_argument("--tab", choices=["chat", "fleet", "approvals", "jobs", "tmux", "logs", "ssh"], default=None, help="Initial Box tab (box mode only)")
args = parser.parse_args()
try:
@@ -7046,7 +7465,8 @@ def main():
stdscr,
initial_mode=args.mode,
initial_node=args.account,
initial_thread=args.thread
initial_thread=args.thread,
initial_tab=args.tab
).run())
except KeyboardInterrupt:
pass
File diff suppressed because it is too large Load Diff
+66 -3
View File
@@ -252,6 +252,56 @@ def provision_node_infra(node: str) -> Dict[str, Any]:
return {"ok": True, "node": node, "output": res.stdout.strip()}
def _registry_port(node: str) -> Optional[int]:
"""CDP port for a node via netvm-registry.py, or None if unregistered."""
import importlib.util
spec = importlib.util.spec_from_file_location(
"netvm_registry", str(BIN_DIR / "netvm-registry.py"))
mod = importlib.util.module_from_spec(spec)
spec.loader.exec_module(mod)
return mod.port_for(node)
def _cdp_dry_run(node: str) -> bool:
"""True when the node netns + CDP + page chain is healthy."""
cmd = [str(BIN_DIR / "netvm-exec.sh"), node, "--", sys.executable,
str(BIN_DIR / "onboard-driver.py"),
"--node", node, "--service", "muse", "--id-type", "email",
"--step", "initiate", "--dry-run"]
res = subprocess.run(cmd, capture_output=True, text=True, timeout=30)
return res.returncode == 0
def ensure_node_browser(node: str, timeout: float = 90.0, poll_interval: float = 5.0) -> Dict[str, Any]:
"""Launch the node headless browser inside its netns if CDP is down.
start_onboarding() must call this after infra provisioning: provision
never starts a browser, so without this step auth initiation always
dies with CDP connection refused on fresh nodes.
"""
port = _registry_port(node)
if not port:
return {"ok": False, "node": node,
"error": "unknown node %s (not in NODES.md registry)" % node}
if _cdp_dry_run(node):
return {"ok": True, "node": node, "cdp_port": port, "already": True}
STATE_DIR.mkdir(parents=True, exist_ok=True)
log_path = STATE_DIR / ("%s-chrome.log" % node)
cmd = [str(BIN_DIR / "netvm-chrome.sh"), "--headless",
"--cdp-port", str(port), node, "https://muse.ai"]
with open(log_path, "ab") as log:
subprocess.Popen(cmd, start_new_session=True,
stdout=log, stderr=subprocess.STDOUT,
stdin=subprocess.DEVNULL)
deadline = time.time() + timeout
while time.time() < deadline:
time.sleep(poll_interval)
if _cdp_dry_run(node):
return {"ok": True, "node": node, "cdp_port": port, "already": False}
return {"ok": False, "node": node, "cdp_port": port,
"error": "browser launched but CDP stayed unreachable on port %s (log: %s)" % (port, log_path)}
def start_onboarding(node: str, email: str, beneficiary_node: Optional[str] = None, invite_code: Optional[str] = None, account_name: Optional[str] = None) -> Dict[str, Any]:
"""Phase 1 & 2: Provision infra, choose beneficiary invite code, and initiate authentication."""
b_node, code, reason = select_urgent_beneficiary(beneficiary_node, invite_code)
@@ -279,6 +329,14 @@ def start_onboarding(node: str, email: str, beneficiary_node: Optional[str] = No
state.stage = STAGE_INFRA
state.save()
# 1b. Ensure the headless browser is up (provision never starts one).
browser_res = ensure_node_browser(node)
if not browser_res.get("ok"):
state.stage = "browser_failed"
state.detail = browser_res.get("error")
state.save()
return {"ok": False, "state": asdict(state), "error": state.detail}
# 2. Initiate authentication
client = CredClient()
cred_res = client.initiate(node, email, service="muse", account_name=account_name)
@@ -408,12 +466,17 @@ def issue_salvage_work_order(blocked_node: str = "646", to_sidechat: str = "646
"--to", "opm",
"--target", to_sidechat,
"--title", title,
"--body", body,
"--priority", "urgent",
"--allow-main-chat"
"--allow-main-chat",
body,
]
res = subprocess.run(cmd, capture_output=True, text=True)
return {"ok": res.returncode == 0, "output": res.stdout.strip()}
out = res.stdout.strip()
result: Dict[str, Any] = {"ok": res.returncode == 0, "output": out}
if not result["ok"]:
err = res.stderr.strip()
result["error"] = err or out or "dm wo exited %d" % res.returncode
return result
def get_all_connects(fast: bool = True) -> List[Dict[str, Any]]:
+4 -3
View File
@@ -118,7 +118,7 @@ def wrap(job_name, job_id, agent, target, rendered, include_kpi: bool = True):
top = (
f"Operator Directive [ref:{wo_id}]:\n"
f"Host tmux worker session '{session_name}' is available on bl (/tmp/tmux-muse.sock).\n"
f"Persistent box runtime is on bl (/tmp/tmux-muse.sock). No worker session exists yet — create yours first: [TOOL tmux.new {{\"session\": \"{session_name}\", \"command\": \"bash\"}}].\n"
f" • Subagent assistance: {spawn}\n"
f" • Verification schedule: {follow}\n"
f"{advisory_section}\n"
@@ -128,8 +128,9 @@ def wrap(job_name, job_id, agent, target, rendered, include_kpi: bool = True):
has_result = "[RESULT" in rendered
bottom = (
"\n--- End Task ---\n\n"
f"Inspect tmux worker: box tmux capture {session_name} 30 (or attach via /tmp/tmux-muse.sock)\n"
"Tools: cron.create, cron.runs, health.check, swarm.spawn, swarm.list, dm.send, box.exec, tools.list.\n"
f"Worker convention: name your tmux session {session_name} when you create it, then inspect via box tmux capture {session_name} 30.\n"
"Flow in tmux: [TOOL flow.start {\"flow_id\": \"<id>\", \"command\": \"<cmd>\"}] | read delta: [TOOL flow.read {\"flow_id\": \"<id>\"}] | advance: [TOOL flow.send {\"flow_id\": \"<id>\", \"command\": \"...\"}].\n"
"Tools: flow.start, flow.read, flow.send, cron.create, health.check, swarm.spawn, dm.send, box.exec, tools.list.\n"
"Message a peer: [DM {\"to\": \"<agent>\", \"target\": \"<sidechat>\", \"message\": \"<text>\"}].\n"
"Query box: [TOOL box.exec {\"action\": \"<fleet-status|dm-log|job-get|...>\"}] — [TOOL tools.list {}] lists every op.\n"
)
+167 -22
View File
@@ -109,7 +109,18 @@ def iter_result_markers(text):
"""Yield (job_id, result_text) for every [RESULT <job_id>] marker in text."""
current_re = lookup_engine.get_result_regex() if HAS_LOOKUP_ENGINE else RESULT_RE
for m in current_re.finditer(text or ""):
yield m.group(1).strip(), m.group(2).strip()
gd = m.groupdict()
if "summary" in gd:
# Engine shape: [RESULT <id>] [STATUS] <summary>. The
# status word is optional (None for bare markers); keep
# it when present so FAIL/ERROR still trips failure
# detection downstream.
status = (m.group("status") or "").strip()
summary = (m.group("summary") or "").strip()
result_text = f"{status} {summary}".strip() if status else summary
yield m.group("job_id").strip(), result_text
else:
yield m.group(1).strip(), m.group(2).strip()
# Verb markers for the digest response protocol:
@@ -316,9 +327,19 @@ def result_has_evidence(result_text):
return bool(_PROOF_EVIDENCE_RE.search(result_text or ""))
# Automated in-thread proof requests disabled per fleet governance decision (2026-10-09)
PROOF_REQUESTS_ENABLED = False
def maybe_request_proof(agent, thread_id, job_id, result_text, dry_run=False):
"""Ask for checkable evidence when a success RESULT has none.
Disabled by default per fleet decision 2026-10-09: automated in-thread proof
challenges trigger adversarial rejection loops and waste agent quota.
"""
if not PROOF_REQUESTS_ENABLED:
return False
"""Ask for checkable evidence when a success RESULT has none.
One-shot per (thread, job) via the nudge tracker. Returns True when a
proof followup was scheduled.
"""
@@ -559,6 +580,42 @@ def format_tool_result_for_chat(op, raw_output):
out = out[:900] + "\n…(truncated, refine the call for detail)"
return f"box result:\n```\n{out}\n```"
if op == "flow.start" and isinstance(data, dict):
if not data.get("ok"):
return f"Flow start failed: {data.get('error')}"
return f"Flow `{data.get('flow_id')}` started in pane `{data.get('session')}` (status: {data.get('status')})."
if op == "flow.read" and isinstance(data, dict):
if not data.get("ok"):
return f"Flow read failed: {data.get('error')}"
st = data.get("status", "unknown")
ec = data.get("exit_code")
ec_str = f" (exit_code: {ec})" if ec is not None else ""
pm = data.get("prompt_match")
prompt_str = f"\nPrompt waiting: {pm.get('text', pm)}" if pm else ""
delta = data.get("delta", "").strip()
trunc = f" (last {data.get('lines_read')} lines)" if data.get("truncated") else ""
body = f"\n```\n{delta}\n```" if delta else " (no new output)"
return f"Flow `{data.get('flow_id')}` [{st}]{ec_str}{prompt_str}{trunc}:{body}"
if op == "flow.send" and isinstance(data, dict):
if not data.get("ok"):
return f"Flow send failed: {data.get('error')}"
kind = "command" if data.get("is_command") else "keys"
return f"Flow `{data.get('flow_id')}` sent {kind}: `{data.get('sent')}` (status: {data.get('status')})."
if op == "flow.list" and isinstance(data, dict):
flows = data.get("flows", [])
if not flows:
return "No active flows."
lines = [f"{len(flows)} flows:"]
for f in flows[:8]:
lines.append(f" • {f.get('flow_id')} [{f.get('status')}]: {f.get('session')} (cmd: {str(f.get('command', 'bash'))[:30]})")
return "\n".join(lines)
if op == "flow.stop" and isinstance(data, dict):
return f"Flow `{data.get('flow_id')}` stopped."
# General fallback: compact JSON capped to 400 chars
s = json.dumps(data)
return s[:400] + "..." if len(s) > 400 else s
@@ -623,11 +680,59 @@ def is_fail_result(result_text):
return t.startswith(FAIL_PREFIXES)
RECENCY_WINDOW_SEC = 10800
_JOB_ID_RE = re.compile(r"^(.+)-(\d{8})-(\d{6})-([0-9a-f]{8})$")
def dispatched_families_since(job_log_path, window_sec=RECENCY_WINDOW_SEC,
now=None):
"""Job families dispatched inside the window.
Scans job-log.jsonl for job_sent/job_dispatched events newer than
``window_sec`` and returns their family names (the job id minus the
trailing -YYYYMMDD-HHMMSS-<hash> run suffix). Missing, unreadable,
or malformed input yields an empty set, never an exception.
"""
now = now or datetime.now(timezone.utc)
cutoff = now.timestamp() - window_sec
fams = set()
try:
handle = open(job_log_path, "r", encoding="utf-8")
except OSError:
return fams
with handle:
for line in handle:
line = line.strip()
if not line:
continue
try:
event = json.loads(line)
except Exception:
continue
if event.get("type") not in ("job_sent", "job_dispatched"):
continue
try:
ts = datetime.fromisoformat(
str(event.get("ts")).replace("Z", "+00:00")).timestamp()
except Exception:
continue
if ts < cutoff:
continue
match = _JOB_ID_RE.match(str(event.get("job_id") or ""))
if match:
fams.add(match.group(1))
return fams
def get_monitored_threads(target_agent=None):
"""
Build dict of threads to monitor per agent:
{ agent: [ {"id": "<uuid>", "name": "<alias>"} ] }
Filters to permanent channels, threads with pending followups, or recent threads (< 3h).
Filters to permanent channels, threads with pending followups,
recently created threads (< 3h), or threads whose job family was
dispatched recently (< 3h) so old persistent sidechats that still
receive prompts stay monitored.
"""
agents = [target_agent] if target_agent else VALID_AGENTS
threads_by_agent = {a: [] for a in agents}
@@ -642,6 +747,7 @@ def get_monitored_threads(target_agent=None):
PERM_KEYWORDS = ("coord", "tasks", "task", "brain", "heartbeat", "sync", "audit", "main-loop")
now = datetime.now(timezone.utc)
recently_dispatched = dispatched_families_since(JOB_LOG, now=now)
state_files = [JOB_SIDECHATS_FILE, WAKE_SIDECHATS_FILE]
for sf in state_files:
@@ -684,7 +790,9 @@ def get_monitored_threads(target_agent=None):
if isinstance(val, dict) and val.get("archived") and not is_pending:
continue
if not (is_perm or is_pending or is_recent):
is_dispatched = key in recently_dispatched
if not (is_perm or is_pending or is_recent or is_dispatched):
continue
existing = [t["id"] for t in threads_by_agent[agent]]
@@ -896,8 +1004,18 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
append_jsonl(CHAT_HISTORY_LOG, record)
if author == "assistant":
markers = list(iter_result_markers(text))
verbs = list(iter_verb_markers(text))
try:
markers = list(iter_result_markers(text))
verbs = list(iter_verb_markers(text))
except Exception as e:
# One poison message must not wedge the batch: without
# this, the same crash repeats every cycle, the
# watermark never advances past it, and the thread's
# followups nag to escalation despite answered work.
sys.stderr.write(
"warning: marker extraction failed, treating as "
f"plain reply: {e}\n")
markers, verbs = [], []
# Synthesize [RESULT <job-id>] DECLINE if assistant explicitly refuses the task in plain text
if not markers and not verbs and detect_explicit_refusal(text):
@@ -946,19 +1064,27 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
try:
import muse_hybrid
thread_url = f"https://box.muse-dev.online/thread/{thread_id}"
tool_hint = (
f"[Runtime Context: {thread_url}]\n"
f"Tools: EMIT one [TOOL <op> <args>] line per action (you do not run it;"
f" the runtime executes it and replies here). curl -sk -X POST"
f" https://exec.muse-dev.online/exec works too.\n"
f" • [TOOL tools.list {{}}] — discover every op dynamically\n"
f" • [TOOL swarm.spawn {{\"count\": 1, \"task\": \"<task>\"}}] — spawn subagents\n"
f" • [DM {{\"to\": \"<agent>\", \"target\": \"<sidechat>\", \"message\": \"<text>\"}}] — send a DM\n"
f" • [TOOL box.exec {{\"action\": \"fleet-status\"}}] — call box (read-only actions)\n"
f" • [TOOL followup.create {{\"in_m\": 5, \"prompt\": \"<reminder>\"}}]\n"
f" • [TOOL health.check {{}}]\n\n"
f"[Directive: Take next action or close with [RESULT <job_id>] <summary>]"
)
if op.startswith("flow."):
flow_id = t_args.get("flow_id", "<flow_id>") if isinstance(t_args, dict) else "<flow_id>"
tool_hint = (
f"[Flow Directive: advance with [TOOL flow.send {{\"flow_id\": \"{flow_id}\", \"command\": \"...\"}}]"
f" | read with [TOOL flow.read {{\"flow_id\": \"{flow_id}\"}}]"
f" | close with [RESULT <job_id>] OK]"
)
else:
tool_hint = (
f"[Runtime Context: {thread_url}]\n"
f"Tools: EMIT one [TOOL <op> <args>] line per action (you do not run it;"
f" the runtime executes it and replies here). curl -sk -X POST"
f" https://exec.muse-dev.online/exec works too.\n"
f" • [TOOL tools.list {{}}] — discover every op dynamically\n"
f" • [TOOL swarm.spawn {{\"count\": 1, \"task\": \"<task>\"}}] — spawn subagents\n"
f" • [DM {{\"to\": \"<agent>\", \"target\": \"<sidechat>\", \"message\": \"<text>\"}}] — send a DM\n"
f" • [TOOL box.exec {{\"action\": \"fleet-status\"}}] — call box (read-only actions)\n"
f" • [TOOL followup.create {{\"in_m\": 5, \"prompt\": \"<reminder>\"}}]\n"
f" • [TOOL health.check {{}}]\n\n"
f"[Directive: Take next action or close with [RESULT <job_id>] <summary>]"
)
if t_ok:
clean_msg = format_tool_result_for_chat(op, t_res)
resp_text = f"Tool result (`{op}`):\n{clean_msg}\n\n{tool_hint}"
@@ -969,7 +1095,12 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
sys.stderr.write(f"warning: failed to post tool response back to thread: {te}\n")
if markers or verbs:
seen_jobs = set()
for job_id, result_text in markers:
if job_id in seen_jobs:
# Same verdict restated in one message: log once.
continue
seen_jobs.add(job_id)
is_fail = is_fail_result(result_text)
job_results += 1
@@ -983,6 +1114,11 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
"thread_id": thread_id,
"msg_id": mid,
}
if result_text.startswith("DECLINE:"):
# Synthesized (or explicit) decline: still a
# non-success (no chaining), but the auditor
# buckets it as declined, not a failure.
job_record["outcome"] = "declined"
if not dry_run:
append_jsonl(JOB_LOG, job_record)
# Check if this is a swarm slot result: sw-YYYYMMDD-HHMMSS-xxxx/<slot>
@@ -1000,10 +1136,12 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
trigger_chain_next(job_id, result_text, success=not is_fail)
if not is_fail:
try:
maybe_request_proof(agent, thread_id, job_id, result_text)
maybe_request_proof(agent, thread_id, job_id, result_text,
dry_run=dry_run)
except Exception as pe:
sys.stderr.write(f"warning: proof check failed: {pe}\n")
archive_ephemeral_thread(agent, thread_id, job_id=job_id)
archive_ephemeral_thread(agent, thread_id, job_id=job_id,
dry_run=dry_run)
clear_matching_followups(followups, agent, thread_id, mid, text,
dry_run, job_id=job_id, verb="RESULT")
for verb, job_id in verbs:
@@ -1122,11 +1260,13 @@ def harvest_agent_thread(cdp, agent, thread_info, watermarks, followups, dry_run
)
def archive_ephemeral_thread(agent, thread_id, job_id=None):
def archive_ephemeral_thread(agent, thread_id, job_id=None, dry_run=False):
"""
If thread_id belongs to an ephemeral job or one-off check,
archive it via hybrid gateway and tag it as archived in job-sidechats.json.
"""
if dry_run:
return
if not thread_id or not re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()):
return
@@ -1699,10 +1839,15 @@ def clear_matching_followups(followups, agent, thread_id, mid, text, dry_run=Fal
# Match by job_id (from [RESULT <job_id>] or [VERB <job_id>]) --
# works regardless of thread_uuid or target. This is an ADDITIONAL
# path, not a replacement.
# path, not a replacement. Also matches the followup's own key
# (dm_id): agents quote the DM id from nudge text ([RESULT
# <dm_id>]), which differs from job_id on DM-ordered followups
# (observed live: [RESULT f4293153] vs job ml-muse-*).
match_job = False
if job_id and f_rec.get("job_id") and f_rec.get("job_id") == job_id:
match_job = True
elif job_id and job_id == f_id:
match_job = True
# Non-RESULT verbs are job-scoped: they must not acknowledge/resolve
# unrelated pending followups that merely share the thread. RESULT
+390 -8
View File
@@ -1485,6 +1485,8 @@ def cmd_work(args):
box_work.cmd_merge(args)
elif action == "check":
box_work.cmd_check(args)
elif action == "heal":
box_work.cmd_heal(args)
elif action == "chats":
box_work.cmd_chats(args)
else:
@@ -6447,15 +6449,375 @@ def cmd_sysop_install(args):
sys.exit(0)
# ---------------------------------------------------------------------------
# CLI Usage Helpers, Error Formatting, and Deep Manuals
# ---------------------------------------------------------------------------
COMMAND_EXAMPLES = {
"box work": [
"box work # View fleet workspace dashboard & signals",
"box work check [agent] # Audit pre-flight health gates",
"box work heal <agent> # Automated remediation & chat nudge",
"box work start \"<title>\" --to <agent> # Start & dispatch new build ticket",
"box work assign <issue#> --to <agent> # Assign existing ticket",
"box work merge <pr#> # Verify tests and merge PR to master",
"box work chats --agent <name> # View live multi-agent chat feed",
],
"box work start": [
"box work start \"Fix SSH perms\" --to 646",
"box work start \"Build integration tests\" --to pip --goal \"Run pytest on endpoints\"",
"box work start \"Emergency rebuild\" --to dev --force",
],
"box work check": [
"box work check # Check all agents",
"box work check 646 # Check specific agent",
],
"box work heal": [
"box work heal dev # Heal dev agent (token, perms, tunnel nudge)",
"box work heal 646",
],
"box work assign": [
"box work assign 218 --to 646",
],
"box work merge": [
"box work merge 217 # Test and merge PR 217 into master",
],
"box work chats": [
"box work chats # Last 10 chat messages across fleet",
"box work chats --agent opm --limit 5",
],
"box tasks": [
"box tasks list # List all tasks across queues",
"box tasks list --queue pending",
"box tasks show 218-restore-keys.md",
"box tasks create 219-my-task.md --title \"Task title\"",
],
"box fleet": [
"box fleet status # Node health & CDP table",
"box fleet watch # Stream status updates",
"box fleet restart muse",
"box fleet heal dev",
],
"box approvals": [
"box approvals check # Check pending browser approvals",
"box approvals allow 646 # Approve pending browser request",
"box approvals allow muse --always # Whitelist site permanently",
],
"box dm": [
"box dm log --limit 10 # View recent direct messages",
"box dm send dev \"Tunnel is down\"",
"box dm wo 646 \"Restore root authorized_keys\"",
],
"box job": [
"box job list # List scheduled & autonomous jobs",
"box job show <job_id>",
"box job run <job_id> # Trigger execution immediately",
],
"box tmux": [
"box tmux list # List tmux worker sessions",
"box tmux auto status # Status of tmux auto-approver",
],
}
PRIMARY_DOMAINS = [
("work", "Fleet workspace, task orchestration, worker scope, signals"),
("tasks", "Agent task file queue (pending/claimed/done)"),
("fleet", "Node health, CDP status, active tabs, watch, restart, heal"),
("approvals", "Inspect and handle agent browser & gateway approvals"),
("dm", "Direct messaging pipeline between operators and agents"),
("job", "Scheduled & autonomous job management"),
("tmux", "Tmux runtime & worker session manager"),
("sysop", "Fleet operations installer (systemd units & timers)"),
("help", "Comprehensive manual and documentation for any command"),
]
def format_error_shorthand(parser, message):
lines = []
lines.append(f"\n{c_bold(c_red('❌ CLI ERROR:'))} {c_bold(message)}\n")
lines.append(c_bold(c_yellow("💡 SHORTHAND USAGE HELPER:")))
lines.append(f" Command: {c_bold(parser.prog)}")
sub_action = next((a for a in parser._actions if isinstance(a, argparse._SubParsersAction)), None)
if sub_action:
if parser.prog in ("box", "super"):
lines.append(f"\n{c_bold(' Primary Domains & Commands:')}")
for d, desc in PRIMARY_DOMAINS:
lines.append(f" • {c_bold(f'{d:<12}')} {c_dim(desc)}")
else:
lines.append(f"\n{c_bold(' Available Subcommands:')}")
for name, subp in sub_action.choices.items():
h = subp.description or getattr(subp, "help", "") or ""
if not h and getattr(sub_action, "_choices_actions", None):
for ca in sub_action._choices_actions:
if ca.dest == name:
h = ca.help or ""
break
lines.append(f" • {c_bold(f'{name:<12}')} {c_dim(h)}")
positionals = [a for a in parser._actions if not a.option_strings and a.dest != 'help' and not isinstance(a, argparse._SubParsersAction)]
required_options = [a for a in parser._actions if a.option_strings and a.required and a.dest != 'help']
optional_options = [a for a in parser._actions if a.option_strings and not a.required and a.dest != 'help']
if positionals or required_options:
lines.append(f"\n{c_bold(' Required Parameters / Arguments:')}")
for a in positionals:
lines.append(f" • {c_bold(f'{a.dest:<14}')} {a.help or '(positional)'}")
for a in required_options:
opts = "/".join(a.option_strings)
lines.append(f" • {c_bold(f'{opts:<14}')} {a.help or '(required flag)'}")
if optional_options:
lines.append(f"\n{c_bold(' Optional Flags:')}")
for a in optional_options:
opts = "/".join(a.option_strings)
lines.append(f" • {c_cyan(f'{opts:<14}')} {c_dim(a.help or '')}")
prog_key = parser.prog.strip()
if prog_key.startswith("super "):
prog_key = "box " + prog_key[6:]
examples = COMMAND_EXAMPLES.get(prog_key)
if not examples:
parts = prog_key.split()
if len(parts) > 2:
parent_key = " ".join(parts[:2])
examples = COMMAND_EXAMPLES.get(parent_key)
if examples:
lines.append(f"\n{c_bold(' Quick Examples:')}")
for ex in examples:
lines.append(f" {c_green(ex)}")
lines.append(f"\n 📖 {c_dim('For complete manual:')} {c_bold(f'{parser.prog} --help')} {c_dim('(or')} {c_bold(f'box help {parser.prog.split()[-1]}')}{c_dim(')')}\n")
return "\n".join(lines)
class BoxArgumentParser(argparse.ArgumentParser):
def error(self, message):
print(format_error_shorthand(self, message), file=sys.stderr)
sys.exit(2)
def format_help(self):
base_help = super().format_help()
prog_key = self.prog.strip()
if prog_key.startswith("super "):
prog_key = "box " + prog_key[6:]
examples = COMMAND_EXAMPLES.get(prog_key)
if not examples:
parts = prog_key.split()
if len(parts) > 2:
parent_key = " ".join(parts[:2])
examples = COMMAND_EXAMPLES.get(parent_key)
extra = []
if examples:
extra.append(c_bold("\nSHORTHAND EXAMPLES:"))
for ex in examples:
extra.append(f" {c_green(ex)}")
extra.append(c_bold("\nOPERATIONAL GUIDELINES:"))
extra.append(f" • {c_cyan('Shorthand parameter reference:')} run {c_bold('box')} alone")
extra.append(f" • {c_cyan('Master comprehensive manual:')} run {c_bold('box help')} or {c_bold('box help <domain>')}")
extra.append(f" • {c_cyan('JSON output:')} append {c_bold('--json')} to any query command\n")
return base_help + "\n".join(extra)
def print_box_usage_reference():
"""Prints categorized primary domains and input parameters when box is run alone."""
print(c_bold("\n=== BOX ORCHESTRATOR: INPUT PARAMETERS & USAGE REFERENCE ===\n"))
print(f"Usage: {c_bold('box <domain> [action] [arguments...] [options...]')}")
print(f" {c_bold('box help [domain]')} | {c_bold('box <domain> --help')}\n")
print(c_bold("PRIMARY DOMAINS & INPUT PARAMETERS:"))
domains_spec = [
("work", "Fleet workspace, task orchestration, worker scope, and active signals", [
("box work [status]", "Show full operational work dashboard & worker signals"),
("box work check [agent]", "Pre-flight health gates (Hatch, Restore, Git Config)"),
("box work heal <agent>", "Automated remediation (tokens, collaborator, dial-in, chat)"),
("box work start \"<title>\" --to <agent> [--goal \"<goal>\"] [--force]", "Instantly start & assign new ticket to agent"),
("box work assign <issue#> --to <agent> [--force]", "Assign existing Gitea ticket to an agent"),
("box work merge <pr#>", "Verify test suite and merge PR to master"),
("box work chats [--agent <name>] [--limit <n>]", "Inspect live agent chat feeds with stream filtering"),
]),
("tasks", "Agent task file queue (fleet/tasks/{pending,claimed,done})", [
("box tasks list [--queue pending|claimed|done|all]", "List task queue files across queues"),
("box tasks show <task-name>", "Print contents of a task file"),
("box tasks create <name> --title \"<title>\"", "Write new pending task from template"),
]),
("fleet", "Node health, CDP status, active tabs, watch, restart, heal", [
("box fleet [status]", "Show NetVM node status table (muse, pip, 646, opm, dev, def)"),
("box fleet watch [--interval <sec>]", "Live streaming status monitor"),
("box fleet restart <node>", "Restart node browser & services"),
("box fleet cdp <node>", "Print DevTools Protocol endpoint URL"),
("box fleet heal <node>", "Run node remediation"),
]),
("approvals", "Inspect and handle agent browser & gateway approvals", [
("box approvals check [--node <name>] [-v]", "List pending modal browser approval prompts"),
("box approvals allow <node> [--always]", "Approve pending browser prompt"),
("box approvals inspect <node>", "Inspect active DOM modal elements"),
]),
("dm", "Direct messaging pipeline between operators and agents", [
("box dm log [--node <name>] [--limit <n>]", "Read signed message log"),
("box dm send <target> \"<message>\"", "Send message to node/agent"),
("box dm wo <agent> \"<instruction>\"", "Send formal work order to agent"),
("box dm ack <msg_id>", "Acknowledge received work order"),
]),
("job", "Scheduled & autonomous job management", [
("box job list [--all]", "List configured jobs and timers"),
("box job show <job_id>", "Display job configuration"),
("box job run <job_id>", "Trigger immediate execution"),
("box job status <job_id>", "Check execution status"),
]),
("tmux", "Tmux runtime & worker session manager", [
("box tmux [list]", "List active sessions on socket"),
("box tmux auto [status|watch]", "Monitor automated approval daemon"),
]),
("sysop", "Fleet operations installer", [
("box sysop install [--dry-run]", "Install & verify systemd units and timers"),
]),
("help", "Comprehensive manual and documentation for any command", [
("box help [domain]", "Deep documentation & manual"),
]),
]
for name, desc, cmds in domains_spec:
print(f" {c_bold(c_cyan(f'{name:<11}'))} {c_dim(desc)}")
for cmd_syntax, cmd_desc in cmds:
print(f" • {c_bold(cmd_syntax):<64} {c_dim(cmd_desc)}")
print()
print(c_bold("QUICK DISPATCH SHORTCUTS:"))
print(f" Start Task: {c_green('box work start \"<title>\" --to <agent>')}")
print(f" Merge PR: {c_green('box work merge <pr#>')}")
print(f" Heal Agent: {c_green('box work heal <agent>')}")
print(f" Check Health: {c_green('box work check [agent]')}")
print(f"\n{c_dim('Run')} {c_bold('box <domain> --help')} {c_dim('or')} {c_bold('box help <domain>')} {c_dim('for full manuals and argument details.')}\n")
def print_master_help():
"""Prints comprehensive, deep master manual for box help / box --help."""
banner = """
================================================================================
BOX ORCHESTRATOR COMPREHENSIVE CLI & RUNTIME MANUAL
================================================================================
"""
print(c_bold(banner))
print(f"""{c_bold("SYNOPSIS:")}
box <domain> [action] [arguments...] [options...]
box help [domain]
box <domain> --help | box <domain> <action> --help
{c_bold("OVERVIEW:")}
The 'box' CLI is the unified orchestration tool for NetVM nodes, cloud muse
agents (opm, 646, dev, pip, def, muse, muse-main), Gitea CI/CD build tasks,
approval workflows, DM message routing, scheduled jobs, and persistent tmux runtimes.
{c_bold("CORE ARCHITECTURE & WORKER ROLES:")}
• {c_bold("opm")} (port 2228) : Fleet orchestrator & lead coordinator
• {c_bold("646")} (port 2226) : System & core runtime operator
• {c_bold("dev")} (port 2230) : Feature development & dark-node builder
• {c_bold("pip")} (port 2227) : Integration & Python builder
• {c_bold("def")} (port 2229) : Defense & telemetry monitor
• {c_bold("muse")} (port 2225) : Cloud workspace agent
• {c_bold("muse-main")} (port 2224) : GCP host node & tunnel anchor
{c_bold("DOMAINS & ACTION SPECIFICATIONS:")}
1. {c_bold("WORK & BUILD PIPELINE (box work ...)")}
Orchestrates autonomous cloud agents, Gitea issue-to-branch pipelines, PR merges,
and pre-flight node health verification.
• {c_bold("box work [status]")}
Parameters: None (optional --json)
Description: Full operational dashboard (worker scope, signals, tickets, PRs, chats).
• {c_bold("box work check [agent]")}
Parameters: agent (optional positional: opm, 646, dev, pip, def, muse)
Description: Pre-flight health gates (Hatch reverse tunnels, Restore persistence, Git credentials).
• {c_bold("box work heal <agent>")}
Parameters: agent (required positional)
Description: Automated self-healing engine (Gitea collaborator rights, partition tokens,
SSH container credential injection, chat recovery nudge).
• {c_bold("box work start \"<title>\" --to <agent> [--goal \"<goal>\"] [--force]")}
Parameters:
title (required positional): Short ticket title
--to (required flag): Target worker agent
--goal (optional flag): Detailed instructions / task goal
--force (optional flag): Bypass failed pre-flight health gate
Description: Runs pre-flight health gate, auto-heals if blocked, creates Gitea Issue #N,
and dispatches briefing directly into agent live chat.
• {c_bold("box work assign <issue#> --to <agent> [--force]")}
Parameters:
issue# (required positional integer): Existing Gitea issue number
--to (required flag): Target agent
Description: Reassigns issue, verifies pre-flight health, notifies agent.
• {c_bold("box work merge <pr#>")}
Parameters: pr# (required positional integer): Pull Request number
Description: Runs test suite verification, merges PR into master, and triggers
post-receive loop terminus hook.
• {c_bold("box work chats [--agent <name>] [--limit <n>]")}
Parameters:
--agent (optional flag): Filter events for specific agent
--limit (optional flag, default 10): Number of events to show
Description: Multi-agent live chat log viewer with agent-scoped stream filtering.
2. {c_bold("TASK FILE QUEUE (box tasks ...)")}
File-backed agent task queues in fleet/tasks/{{pending,claimed,done}}.
• {c_bold("box tasks list [--queue pending|claimed|done|all] [--dir <path>]")}
• {c_bold("box tasks show <name>")}
• {c_bold("box tasks create <name> --title \"<title>\"")}
3. {c_bold("FLEET & NODE MANAGEMENT (box fleet ...)")}
Controls Chromium NetVM nodes, D-Bus network namespaces, and CDP endpoints.
• {c_bold("box fleet [status]")} Show active nodes, latencies, threads
• {c_bold("box fleet watch [--interval <sec>]")} Real-time continuous monitoring
• {c_bold("box fleet restart <node>")} Restart node browser/profile
• {c_bold("box fleet cdp <node>")} Show DevTools protocol endpoint
• {c_bold("box fleet heal <node>")} Remediate crashed or stuck node
4. {c_bold("BROWSER APPROVALS & GATEWAYS (box approvals ...)")}
Inspects and resolves browser modal prompts, ethical-captcha gates, and domain permissions.
• {c_bold("box approvals check [--node <name>] [-v]")} List pending approvals
• {c_bold("box approvals allow <node> [--always]")} Approve pending request
• {c_bold("box approvals inspect <node>")} Inspect active DOM modal elements
5. {c_bold("DIRECT MESSAGING & WORK ORDERS (box dm ...)")}
Encrypted and signed inter-agent communication pipeline.
• {c_bold("box dm log [--node <name>] [--limit <n>]")} Read signed message log
• {c_bold("box dm send <target> \"<message>\"")} Send message to peer node
• {c_bold("box dm wo <agent> \"<instruction>\"")} Issue formal agent work order
• {c_bold("box dm ack <msg_id>")} Acknowledge received work order
6. {c_bold("SCHEDULED JOBS (box job ...)")}
Background automation and recurrent job scheduling.
• {c_bold("box job list [--all]")} List all jobs and timers
• {c_bold("box job show <job_id>")} Inspect job JSON configuration
• {c_bold("box job run <job_id>")} Trigger one-shot immediate run
7. {c_bold("TMUX PERSISTENCE RUNTIME (box tmux ...)")}
Headless terminal session management and auto-approval agents.
• {c_bold("box tmux [list]")} List active sessions on socket
• {c_bold("box tmux auto [status|watch]")} Monitor automated approval daemon
{c_bold("ENVIRONMENT & CONFIGURATION:")}
NETVM_ROOT Path to NetVM workspace root (default: /home/super/Projects/NetVM)
CLICOLOR_FORCE Set to 1 to force ANSI color output in non-tty pipes
NO_COLOR Set to disable ANSI color formatting
GITEA_URL Base URL for Gitea API (auto-detected: loopback on bl, public domain on PC)
GITEA_TOKEN API token for Gitea automation
{c_bold("EXIT CODES:")}
0 Success
1 Operational or pre-flight failure
2 CLI syntax or missing argument error
Run 'box <domain> --help' or 'box help <domain>' for in-depth flags on any command.
""")
def build_parser():
common = argparse.ArgumentParser(add_help=False)
common = BoxArgumentParser(add_help=False)
common.add_argument("--json", action="store_true", help="Output machine-readable JSON")
prog_name = Path(sys.argv[0]).name if sys.argv and sys.argv[0] else "super"
prog_name = Path(sys.argv[0]).name if sys.argv and sys.argv[0] else "box"
if prog_name.endswith(".py"):
prog_name = "super"
prog_name = "box"
parser = argparse.ArgumentParser(
parser = BoxArgumentParser(
prog=prog_name,
description=f"{prog_name} — Unified Orchestrator CLI for NetVM & Box",
formatter_class=argparse.RawDescriptionHelpFormatter,
@@ -6624,6 +6986,8 @@ def build_parser():
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_heal = work_sub.add_parser("heal", parents=[common], help="Run automated healing on an agent")
p_w_heal.add_argument("agent", help="Agent username to heal")
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)")
@@ -7385,12 +7749,30 @@ def main():
res = subprocess.run(cmd)
sys.exit(res.returncode)
parser = build_parser()
# Handle empty arguments (box alone)
if len(sys.argv) == 1:
# Default behavior with no arguments: show fleet status
sys.argv.append("fleet")
sys.argv.append("status")
print_box_usage_reference()
cmd_fleet_status(argparse.Namespace(json=False))
sys.exit(0)
# Handle help variations
if len(sys.argv) > 1:
if sys.argv[1] == "help":
if len(sys.argv) == 2:
print_master_help()
sys.exit(0)
else:
target_domain = sys.argv[2]
rest = sys.argv[3:]
sys.argv = [sys.argv[0], target_domain] + rest + ["--help"]
elif sys.argv[1] in ("--help", "-h") and len(sys.argv) == 2:
print_master_help()
sys.exit(0)
elif "help" in sys.argv[2:]:
h_idx = sys.argv.index("help")
sys.argv[h_idx] = "--help"
parser = build_parser()
args = parser.parse_args()
# Route commands
+143 -1
View File
@@ -609,6 +609,7 @@ def test_wired_dm_send_asserts_post_nav_url_before_send():
assert "assert_pre_send_placement" in src, \
"gate exists but dm_send never calls it"
_orig_run_full = dm.run_full
_orig_sleep = dm.time.sleep
_calls = []
def _stub(cmd, timeout=60):
@@ -619,6 +620,9 @@ def test_wired_dm_send_asserts_post_nav_url_before_send():
try:
dm.run_full = _stub
# Settle sleeps (1s/2s per gate call) are production pacing, not
# asserted behavior: skip them like the browser subprocess above.
dm.time.sleep = lambda s: None
# 1. UUID-known thread, correct placement -> pass
_stub.url = "https://muse.ai/thread/" + UUID_A
ok, detail = gate("opm", "pipe-x", UUID_A, direct_nav_done=True)
@@ -641,6 +645,7 @@ def test_wired_dm_send_asserts_post_nav_url_before_send():
assert ok is True, f"re-nav path should pass: {detail}"
finally:
dm.run_full = _orig_run_full
dm.time.sleep = _orig_sleep
return
nav_i = src.find("sidechat use")
send_i = src.find("Send with verification retries")
@@ -754,9 +759,141 @@ def main():
import unittest
# --------------------------------------------------------------------------
# harvester resurrection -- dm_id markers, dry-run purity, scheduling
# --------------------------------------------------------------------------
def test_wired_harvester_dmid_marker_resolves():
"""Real clear_matching_followups: [RESULT <dm_id>] resolves a
DM-ordered followup whose job_id differs (live f4293153 pattern:
marker quoted the nudge's DM id, record job was ml-muse-*).
Unrelated thread isolates the dm_id path from thread matching."""
harv = _load("harvester_under_test", "response-harvester.py")
rec = _mk_rec(thread_uuid=UUID_B, job_id="ml-muse-20261007-013210")
fups = {"f4293153": rec}
harv.clear_matching_followups(fups, "646", "unrelated-thread", "mid-9",
"[RESULT f4293153] done", dry_run=True,
job_id="f4293153", verb="RESULT")
assert rec.get("status") == "resolved", (
"DEVIATION: [RESULT <dm_id>] does not resolve its followup -- "
"clear_matching_followups() matches marker ids against job_id "
f"only, never the followup key (status={rec.get('status')!r})")
# ... and a wrong id must not resolve.
rec2 = _mk_rec(thread_uuid=UUID_B, job_id="ml-muse-20261007-013210")
fups2 = {"f4293153": rec2}
harv.clear_matching_followups(fups2, "646", "unrelated-thread", "mid-9",
"[RESULT deadbeef] done", dry_run=True,
job_id="deadbeef", verb="RESULT")
assert rec2.get("status") == "pending", (
f"wrong marker id wrongly resolved (status={rec2.get('status')!r})")
def test_wired_harvester_dry_run_has_no_side_effects():
"""process_messages(dry_run=True) with an evidence-less RESULT must
still extract the marker but must not fire proof followups,
archive threads, or persist anything."""
harv = _load("harvester_under_test", "response-harvester.py")
calls = []
saved = {n: getattr(harv, n) for n in
("execute_agent_tool", "archive_ephemeral_thread",
"append_jsonl", "save_json_file")}
harv.execute_agent_tool = lambda *a, **k: calls.append("exec") or (True, {})
harv.archive_ephemeral_thread = (
lambda *a, **k: calls.append("archive"))
harv.append_jsonl = lambda *a, **k: calls.append("append")
harv.save_json_file = lambda *a, **k: calls.append("save")
try:
msgs = [{"id": "m1", "author": "assistant",
"text": "[RESULT j1] done",
"ts": "2026-10-07T00:00:00+00:00"}]
new, wm, nres = harv.process_messages(
msgs, "646", UUID_A, "t", None, {}, dry_run=True)
finally:
for n, fn in saved.items():
setattr(harv, n, fn)
assert nres == 1, "dry-run must still extract markers"
assert calls == [], f"dry-run leaked side effects: {calls}"
def test_wired_result_markers_bare_and_status_forms():
"""iter_result_markers handles the engine's 3-group shape: a bare
[RESULT <id>] <text> (status None) must not crash, and a status
token must survive into the result text for fail detection."""
harv = _load("harvester_under_test", "response-harvester.py")
assert list(harv.iter_result_markers("[RESULT f4293153] done")) == [
("f4293153", "done")], "bare RESULT marker must extract cleanly"
jid, text = list(harv.iter_result_markers("[RESULT j9] FAIL blew up"))[0]
assert jid == "j9" and "FAIL" in text and "blew up" in text, (
f"status token must survive into result text (got {jid!r} {text!r})")
def test_wired_poison_message_does_not_wedge_batch():
"""A marker-extraction crash degrades to plain-reply handling so
sibling messages still process and the watermark keeps advancing."""
harv = _load("harvester_under_test", "response-harvester.py")
real_iter = harv.iter_result_markers
real_nudge = harv.maybe_nudge_untagged_sidechat
def boom(text):
if "POISON" in (text or ""):
raise RuntimeError("boom")
return real_iter(text)
harv.iter_result_markers = boom
harv.maybe_nudge_untagged_sidechat = lambda *a, **k: None
try:
msgs = [
{"id": "m1", "author": "assistant",
"text": "POISON [RESULT x] y",
"ts": "2026-10-07T00:00:00+00:00"},
{"id": "m2", "author": "assistant",
"text": "[RESULT j2] ok",
"ts": "2026-10-07T00:01:00+00:00"},
]
new, wm, nres = harv.process_messages(
msgs, "646", UUID_A, "t", None, {}, dry_run=True)
finally:
harv.iter_result_markers = real_iter
harv.maybe_nudge_untagged_sidechat = real_nudge
assert nres == 1, "sibling marker must still extract"
assert [m["id"] for m in new] == ["m1", "m2"], \
"both messages must process past the poison one"
def test_harvester_timer_unit_wired():
"""The harvester must be scheduler-owned: unit files exist, the
service runs --once, and the timer fires on a short cadence.
Ingestion died silently for ~22h with no unit at all."""
root = BIN_DIR.parent
svc = (root / "systemd" / "response-harvester.service").read_text()
tmr = (root / "systemd" / "response-harvester.timer").read_text()
assert "response-harvester.py" in svc and "--once" in svc, \
"service must run the harvester --once"
assert "OnUnitActiveSec=" in tmr, "timer needs a repeat cadence"
assert "WantedBy=timers.target" in tmr, "timer must target timers.target"
def test_collection_adapter_is_single_and_pytest_opted_out():
"""Collection-shape guard (no 3x duplicates): exactly one TestCase
adapter is reachable from module globals (the adapter loop must not
leak a `_fn` alias that pytest collects as a second class), and the
adapter opts out of pytest (`__test__ = False`) so the module-level
functions are pytest's single source while unittest discovery still
runs the adapter."""
cases = [v for v in list(globals().values())
if inspect.isclass(v) and issubclass(v, unittest.TestCase)]
assert len(cases) == 1, (
f"expected exactly 1 TestCase adapter, found {len(cases)} "
f"(stray aliases reintroduce duplicate collection)")
assert TestFollowupFixes.__test__ is False, (
"TestFollowupFixes must set __test__ = False so pytest collects "
"each test once via the module-level functions")
class TestFollowupFixes(unittest.TestCase):
"""unittest discovery adapter for contract and wired test functions."""
pass
# pytest collects the module-level functions; skip the adapter so each
# test runs once. (unittest discovery ignores __test__ and still runs
# the adapter, which is its only view of this file's tests.)
__test__ = False
for _name, _fn in list(globals().items()):
@@ -767,6 +904,11 @@ for _name, _fn in list(globals().items()):
return _runner
setattr(TestFollowupFixes, _name, _bind(_fn))
# Drop the loop temporaries: after the final iteration `_fn` aliases
# TestFollowupFixes, and pytest collects TestCase subclasses regardless of
# name -- that stray alias was the third copy (module fn + adapter + `_fn`).
del _name, _fn
if __name__ == "__main__":
sys.exit(main())
+82 -15
View File
@@ -13,6 +13,7 @@ Supports:
- A/B/C choice prompts -> "A"
- Numbered menus -> "1"
- y/n confirmation prompts -> "y"
- Interview navigate+select menus (cursor on 1 -> Enter)
- Press Enter prompts -> "Enter"
- Safety guardrails (passwords, passkeys, destructive commands are never auto-approved)
4. State persistence & audit logging:
@@ -137,6 +138,21 @@ DEFAULT_RULES: List[MatchRule] = [
description="Confirms y/n at end of terminal line",
press_enter=True,
),
MatchRule(
id="interview_select",
name="Interview Menu (cursor on 1)",
pattern=(r"\?\s*\n"
r"(?:[^\n]*\n){0,8}"
r"[ \t]*(?:›|>)[ \t]*1\.[ \t]+\S[^\n]*\n"
r"(?:[^\n]*\n){0,10}"
r"[ \t]*2\.[ \t]+\S"),
response_key="Enter",
category="enter",
enabled=True,
description=("Selects highlighted option 1 on navigate+select "
"menus (cursor on 1. + 2. + ?-question above)"),
press_enter=False,
),
MatchRule(
id="enter_to_continue",
name="Press Enter to Continue",
@@ -188,9 +204,22 @@ class AutoApproverState:
try:
with open(STATE_FILE) as f:
data = json.load(f)
return cls(**data)
st = cls(**data)
except Exception:
return cls()
# Migrate: append built-in rules missing from stored state (a new
# default must reach the daemon without wiping operator toggles).
try:
have = {r.get("id") for r in st.rules
if isinstance(r, dict)}
missing = [asdict(r) for r in DEFAULT_RULES
if r.id not in have]
if missing:
st.rules.extend(missing)
st.save()
except Exception:
pass
return st
# =====================================================================
@@ -300,6 +329,23 @@ def capture_pane_text(socket_path: str, pane_id: str, lines: int = 30) -> str:
MUSE_COMMAND_HINTS = ("muse-bin", "muse-code")
# (socket, pane) ever observed running a muse runtime. pane_current_command
# flickers to the child tool while the agent works, so a muse pane stays
# muse-owned when its foreground reads "python3" (observed live: the hint
# gate missed tool-running panes and both daemons stacked 'y' answers).
_MUSE_PANES_SEEN = set()
# tmux rule category -> muse watcher kind for verified sends. Text-input
# categories verify render + submit with one retry; single-key widgets
# (and unknown categories) stay blind.
_CATEGORY_KIND_MAP = {
"choice": "letter",
"menu": "numbered",
"confirm": "yn",
"muse_code": "muse-approval",
"enter": None,
}
def should_defer_to_muse_watcher(socket_path: str, pane_id: str,
current_command: str) -> bool:
@@ -309,16 +355,25 @@ def should_defer_to_muse_watcher(socket_path: str, pane_id: str,
panes (stability + re-verify + once-per-prompt + decided-block
guard). When its daemon is alive for this socket:pane, tmux must
skip the pane entirely, or both daemons answer the same prompt
within the same second ('11' + stray keys, observed live). Never
raises: import or liveness failures mean no owner, handle here.
within the same second ('11' + stray keys, observed live; later the
same hole stacked 'y' answers when the foreground flickered to a
child tool mid-poll). Never raises: import or liveness failures
mean no owner, handle here.
"""
try:
cmd = current_command or ""
if not any(h in cmd for h in MUSE_COMMAND_HINTS):
return False
import muse_choice_watcher as mcw
alive = getattr(mcw, "watcher_alive", mcw.is_running)
return alive(socket_path, pane_id) is not None
cmd = current_command or ""
key = (socket_path, pane_id)
if any(h in cmd for h in MUSE_COMMAND_HINTS):
_MUSE_PANES_SEEN.add(key)
alive = getattr(mcw, "watcher_alive", mcw.is_running)
return alive(socket_path, pane_id) is not None
if key in _MUSE_PANES_SEEN:
alive = getattr(mcw, "watcher_alive", mcw.is_running)
return alive(socket_path, pane_id) is not None
# Never observed as muse: cheap pidfile check only (covers a
# watcher racing ahead of our first observation of the pane).
return mcw.is_running(socket_path, pane_id) is not None
except Exception:
return False
@@ -568,15 +623,25 @@ class AutoApproverRunner:
})
continue
# Execute key dispatch
# Execute key dispatch through the verified send path:
# literal text paced apart from Enter (a single-call
# burst arrives as paste and lands a newline in
# composers instead of submitting, then re-fires past
# dedup and stacks). Text-input categories also verify
# render + submit with one retry; single-key widgets
# stay blind.
success = False
detail = {"verified": None, "retried": False}
if not self.dry_run:
args = ["send-keys", "-t", p.pane_id, verdict.key]
if verdict.press_enter or verdict.key == "Enter":
if verdict.key != "Enter":
args.append("Enter")
rc, _, _ = run_tmux_cmd(p.socket, *args)
success = (rc == 0)
import muse_choice_watcher as mcw
want_enter = (verdict.key != "Enter"
and bool(verdict.press_enter))
ok, detail = mcw.send_answer(
p.socket, p.pane_id, verdict.key,
enter=want_enter,
kind=_CATEGORY_KIND_MAP.get(verdict.category),
sig=sig)
success = bool(ok)
else:
success = True # dry-run simulated
@@ -597,6 +662,8 @@ class AutoApproverRunner:
"excerpt": verdict.excerpt,
"dry_run": self.dry_run,
"success": success,
"verified": detail["verified"],
"retried": detail["retried"],
}
self.record_audit(event)
actions_taken.append(event)
+90 -17
View File
@@ -2,10 +2,11 @@
"""tmux_server_watchdog.py — Death-capture for tmux servers.
Runs on a 1-minute systemd timer. Remembers each known socket's server
pid; when a server dies or its pid changes without a witnessed death,
appends a forensics bundle (dmesg OOM/kill lines, memory, uptime,
journal tail) to logs/tmux-server-deaths.jsonl so the next "tmux
crashed" leaves evidence instead of a mystery.
identity (pid + /proc starttime + ppid + cmdline); when a server dies,
its pid changes, or its pid is recycled under us without a witnessed
death, appends a forensics bundle (dmesg OOM/kill lines, memory,
uptime, journal tail) to logs/tmux-server-deaths.jsonl so the next
"tmux crashed" leaves evidence instead of a mystery.
Read-only against tmux itself: one `display-message -p` probe per
socket. Never raises; a watchdog must not need its own watchdog.
@@ -55,9 +56,46 @@ def probe(socket_path):
return None
def collect_forensics(socket_path, last_pid):
def proc_identity(pid):
"""Identity dict for a pid: starttime defeats PID-reuse confusion.
Never raises; on any failure returns {"pid": pid} so callers can
still snapshot. starttime is the raw /proc starttime tick (field
22), stable for the life of the process."""
ident = {"pid": pid}
try:
with open("/proc/%d/stat" % pid) as f:
parts = f.read().rsplit(")", 1)[1].split()
# After "(comm)": state ppid pgrp session tty_nr ... starttime
# is field 22 overall, i.e. parts[19] after the split above.
ident["ppid"] = int(parts[1])
ident["starttime"] = int(parts[19])
except Exception:
pass
try:
with open("/proc/%d/cmdline" % pid, "rb") as f:
raw = f.read().replace(b"\0", b" ").decode(
"utf-8", "replace").strip()
if raw:
ident["cmd"] = raw[:200]
except Exception:
pass
return ident
def probe_identity(socket_path):
"""Enriched snapshot for a socket: identity dict or None."""
pid = probe(socket_path)
if pid is None:
return None
return proc_identity(pid)
def collect_forensics(socket_path, last_pid, last_identity=None):
"""Best-effort death evidence. Dict of strings, never raises."""
ev = {"ts": _now(), "socket": socket_path, "last_pid": last_pid}
if last_identity:
ev["last_identity"] = last_identity
rc, dmesg = _run(["dmesg"], timeout=10)
if rc != 0:
ev["dmesg"] = "unavailable: %s" % dmesg[:200]
@@ -112,27 +150,60 @@ def append_death(ev, path=None):
pass
def _as_identity(value):
"""Normalize a probed value to an identity dict (legacy int ok)."""
if value is None:
return None
if isinstance(value, dict):
return value
return {"pid": value}
def _prev_identity(prev):
ident = {"pid": prev.get("pid")}
for key in ("starttime", "ppid", "cmd"):
if prev.get(key) is not None:
ident[key] = prev[key]
return ident
def evaluate(previous, probed):
"""Pure transition logic: (prev_state, {sock: pid|None}) ->
(new_state, events). Events: death | restart | started."""
"""Pure transition logic: (prev_state, {sock: pid|identity|None}) ->
(new_state, events). Events: death | restart | started.
Probed values may be a bare pid (legacy) or an identity dict from
probe_identity(). Same pid with a different starttime is a restart
(pid recycled under us), not steady state."""
new_state, events = {}, []
for sock, pid in sorted(probed.items()):
for sock, raw in sorted(probed.items()):
ident = _as_identity(raw)
prev = (previous.get(sock) or {})
prev_pid = prev.get("pid")
if pid is None:
if ident is None:
new_state[sock] = {"pid": None, "died": _now(),
"last_pid": prev_pid}
if prev_pid:
events.append({"type": "death", "socket": sock,
"last_pid": prev_pid})
"last_pid": prev_pid,
"last_identity": _prev_identity(prev)})
else:
new_state[sock] = {"pid": pid, "since": _now()}
pid = ident.get("pid")
new_state[sock] = dict(ident, since=_now())
if prev_pid and prev_pid != pid:
# Changed with no witnessed death: restart inside one
# tick gap (or pid recycled under us). Treat as a
# restart, still worth a forensics note.
# tick gap. Worth a forensics note.
events.append({"type": "restart", "socket": sock,
"old_pid": prev_pid, "pid": pid})
"old_pid": prev_pid, "pid": pid,
"last_identity": _prev_identity(prev)})
elif (prev_pid and prev_pid == pid
and prev.get("starttime") is not None
and ident.get("starttime") is not None
and prev["starttime"] != ident["starttime"]):
# Same pid, different process: pid recycled under us.
events.append({"type": "restart", "socket": sock,
"old_pid": prev_pid, "pid": pid,
"pid_reused": True,
"last_identity": _prev_identity(prev)})
elif not prev_pid and prev.get("died"):
events.append({"type": "started", "socket": sock,
"pid": pid})
@@ -144,18 +215,20 @@ def evaluate(previous, probed):
def check(sockets=None, dry_run=False):
"""Probe, transition state, log deaths. Returns summary dict."""
probed = {s: probe(s) for s in (sockets or KNOWN_SOCKETS)}
probed = {s: probe_identity(s) for s in (sockets or KNOWN_SOCKETS)}
previous = read_state()
new_state, events = evaluate(previous, probed)
for ev in events:
if ev["type"] == "death":
bundle = collect_forensics(ev["socket"], ev["last_pid"])
bundle = collect_forensics(ev["socket"], ev["last_pid"],
ev.get("last_identity"))
bundle["event"] = "death"
if not dry_run:
append_death(bundle)
ev["forensics"] = bundle
elif ev["type"] == "restart":
bundle = collect_forensics(ev["socket"], ev["old_pid"])
bundle = collect_forensics(ev["socket"], ev["old_pid"],
ev.get("last_identity"))
bundle["event"] = "restart-gap-missed"
if not dry_run:
append_death(bundle)
+41 -7
View File
@@ -42,13 +42,21 @@ needs_provisioning() { [ ! -f "$SENTINEL" ]; }
restore_ssh_keys() {
# Key restoration: rebuilds may wipe ~/.ssh. Restore from persistent store if present.
if [ ! -f "$HOME/.ssh/vm_to_gcp" ] && [ -f "$HOME/workspace/.ssh-keys/vm_to_gcp" ]; then
log "restoring ~/.ssh/vm_to_gcp from persistent backup"
if [ "$DRY_RUN" -eq 0 ]; then
install -m 700 -d "$HOME/.ssh"
install -m 600 "$HOME/workspace/.ssh-keys/vm_to_gcp" "$HOME/.ssh/vm_to_gcp"
install -m 700 -d "$HOME/.ssh" 2>/dev/null || true
for keyname in vm_to_gcp id_frontdoor; do
if [ ! -f "$HOME/.ssh/$keyname" ]; then
if [ -f "$HOME/workspace/.ssh-keys/$keyname" ]; then
log "restoring ~/.ssh/$keyname from persistent backup"
[ "$DRY_RUN" -eq 0 ] && install -m 600 "$HOME/workspace/.ssh-keys/$keyname" "$HOME/.ssh/$keyname"
elif [ -f "$HOME/workspace/.ssh-keys/vm_to_gcp" ]; then
log "linking ~/.ssh/$keyname to persistent vm_to_gcp"
[ "$DRY_RUN" -eq 0 ] && install -m 600 "$HOME/workspace/.ssh-keys/vm_to_gcp" "$HOME/.ssh/$keyname"
elif [ -f "$HOME/workspace/.ssh-keys/id_frontdoor" ]; then
log "linking ~/.ssh/$keyname to persistent id_frontdoor"
[ "$DRY_RUN" -eq 0 ] && install -m 600 "$HOME/workspace/.ssh-keys/id_frontdoor" "$HOME/.ssh/$keyname"
fi
fi
fi
done
}
provision_critical() {
@@ -115,6 +123,22 @@ provision_critical() {
/home/muse/.ssh/authorized_keys 2>/dev/null || true
fi
# 6. Restore /root/.ssh/authorized_keys across rebuilds
install -m 700 -d /root/.ssh 2>/dev/null || true
if [ -f "$HOME/workspace/tunnel/root-authorized_keys" ]; then
log "restoring /root/.ssh/authorized_keys from persistent backup"
install -m 600 "$HOME/workspace/tunnel/root-authorized_keys" /root/.ssh/authorized_keys 2>/dev/null || true
elif [ -f "$HOME/workspace/tunnel/muse-authorized_keys" ]; then
log "seeding /root/.ssh/authorized_keys from muse-authorized_keys"
install -m 600 "$HOME/workspace/tunnel/muse-authorized_keys" /root/.ssh/authorized_keys 2>/dev/null || true
fi
if [ -f "/home/hatch/.ssh/authorized_keys" ]; then
log "merging /home/hatch/.ssh/authorized_keys into /root/.ssh/authorized_keys"
cat /home/hatch/.ssh/authorized_keys >> /root/.ssh/authorized_keys 2>/dev/null || true
sort -u /root/.ssh/authorized_keys -o /root/.ssh/authorized_keys 2>/dev/null || true
chmod 600 /root/.ssh/authorized_keys 2>/dev/null || true
fi
touch "$SENTINEL"
log "critical provisioning complete"
}
@@ -150,6 +174,11 @@ provision_deferred() {
cp -r "$HOME/workspace/nvim/"* /opt/nvim/ 2>/dev/null || true
ln -sf /opt/nvim/bin/nvim /usr/local/bin/nvim 2>/dev/null || true
fi
local wheel_dir="$HOME/workspace/wheels"
if [ -d "$wheel_dir" ] && ls "$wheel_dir"/*.whl >/dev/null 2>&1; then
log "installing cached python wheels from $wheel_dir"
python3 -m pip install --no-index --find-links="$wheel_dir" protocol_muse 2>/dev/null || true
fi
) >/dev/null 2>&1 &
disown 2>/dev/null || true
}
@@ -178,7 +207,12 @@ ensure_gcp_tunnel() {
log "gcp tunnel supervisor already running"
exit 0
fi
if [ ! -f "$HOME/.ssh/vm_to_gcp" ]; then
if [ ! -f "$HOME/.ssh/vm_to_gcp" ] && [ -f "$HOME/.ssh/id_frontdoor" ]; then
ln -sf "$HOME/.ssh/id_frontdoor" "$HOME/.ssh/vm_to_gcp"
elif [ ! -f "$HOME/.ssh/id_frontdoor" ] && [ -f "$HOME/.ssh/vm_to_gcp" ]; then
ln -sf "$HOME/.ssh/vm_to_gcp" "$HOME/.ssh/id_frontdoor"
fi
if [ ! -f "$HOME/.ssh/vm_to_gcp" ] && [ ! -f "$HOME/.ssh/id_frontdoor" ]; then
log "WARNING: ~/.ssh/vm_to_gcp missing — cannot start gcp tunnel supervisor"
exit 0
fi
+2
View File
@@ -7,3 +7,5 @@ operator-pip ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIMq02n0LpsksyQzWAWQ1mS8gKOonqFA
pip ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIMq02n0LpsksyQzWAWQ1mS8gKOonqFALNDqbPGqXhq4T operator-pip
operator-dev ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAICwHn0kmRa6SFPbr2+z75s0gRlvBCGR633Ag7gTqiYPa dev@netvm
dev ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAICwHn0kmRa6SFPbr2+z75s0gRlvBCGR633Ag7gTqiYPa dev@netvm
def ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIEn6qqPrW7Vc77pUEBnLRDBF+yX11qyWzDTjZ2+FtL7b def@netvm
operator-def ssh-ed25519 AAAAC3NzaC1lZDI1NTE5AAAAIEn6qqPrW7Vc77pUEBnLRDBF+yX11qyWzDTjZ2+FtL7b def@netvm
+38
View File
@@ -177,6 +177,44 @@ Agents can emit structured tool calls in sidechats:
---
## 4.1. Agentic Flows in Tmux Panes (`box flow` & `[TOOL flow.*]`)
Chromebox browser contexts prune and store chat history aggressively, making direct in-chat execution of long-running build, test, and shell tasks token-expensive and prone to context loss.
To overcome this, Chromebox agents offload multi-turn execution to persistent tmux panes on `/tmp/tmux-muse.sock` using the **Flow Engine** (`bin/flow_engine.py`). Raw stdout/stderr streams to disk (`logs/flows/<flow_id>.log`), and agents read back only concise status and incremental output deltas.
### Lifecycle & Primitives:
1. **Start Flow**:
Spawns pane `flow-<agent>-<id>` and launches command wrapped with an exit code sentinel.
```text
[TOOL flow.start {"flow_id": "audit-tests", "command": "python3 -m unittest discover -s tests"}]
```
*CLI:* `box flow start audit-tests -c "python3 -m unittest discover -s tests"`
2. **Read Incremental Delta & State**:
Inspects the pane for execution state (`working`, `idle`, `waiting_prompt`, `finished`, `failed`), exit code, and reads newly appended log output since the last read cursor.
```text
[TOOL flow.read {"flow_id": "audit-tests"}]
```
*CLI:* `box flow read audit-tests --lines 40`
3. **Advance or Respond to Prompts**:
Sends follow-up commands or keystrokes (such as interactive menu selections) without re-running the whole prompt.
```text
[TOOL flow.send {"flow_id": "audit-tests", "command": "git diff"}]
[TOOL flow.send {"flow_id": "audit-tests", "keys": "1"}]
```
*CLI:* `box flow send audit-tests "git status" --command`
4. **List & Stop**:
```text
[TOOL flow.list {}]
[TOOL flow.stop {"flow_id": "audit-tests"}]
```
*CLI:* `box flow list` / `box flow stop audit-tests`
---
## 5. Direct Operator Directives & Prompt Envelope Specification
When jobs are dispatched to agents via `bin/job-dispatch.py`, they are wrapped in an actionable, authentic **Operator Directive** generated by `bin/prompt_envelope.py`.
+46
View File
@@ -0,0 +1,46 @@
---
title: "Gitea Web Surface, Chromebox Gateway Queue & Unified Protocol Muse"
status: "signed-off"
coordinator: "operator-main"
scope: "public-surface-and-protocol"
accepted_at: "2026-10-10T15:40:24Z"
accepted_quote: "Create an automated Pull Request in Gitea via the API and request Coordinator review per docs/AGENT-ROLES.md"
gate: "coordinator"
signoff_targets:
- "release-public-surface"
- "interview-gated-choices"
---
# Decision Record: Gitea Web Surface, Chromebox Gateway & Protocol Muse
Settled and executed 2026-10-10 per `/grill-me` architectural review.
## 1. Summary of Changes
1. **Gitea Multi-Route Surface**:
- Initialized and live-verified portable Gitea instance on port 3000 (`super/box` repository).
- Configured Cloudflare Tunnel ingress routing for `git.muse-dev.online` and `tea.muse-dev.online`.
- Wired Webhook Bridge on port 3005 (`bin/box-gitea-bridge.py`) converting labeled issues into `fleet/tasks/pending/`.
- Integrated `🍵 Git Repos` action button into header controls and Gitea surface card in the Agentic Dev console tab.
2. **Chromebox DM Gateway & Queue Telemetry**:
- Extended `bin/chromebox-gateway.py` with authenticated `GET /api/v1/queue` endpoint.
- Configured Cloudflare Tunnel ingress routing for `dm.muse-dev.online` (dual-exposing with `100.123.153.75:8445`).
- Retained strict Bearer token authentication and allowlisted operations without exposing raw CDP.
3. **Unified `protocol_muse` Python Package**:
- Packaged `packages/protocol-muse/` providing `MuseClient` unifying Noise protocol gateway (`wss://gateway.muse.ai/v1/noise`) with HTTPS Chromebox fallback.
- Built and pre-cached wheels in `~/workspace/wheels/` (`protocol_muse-0.1.0-py3-none-any.whl`).
- Integrated offline wheel installation (`pip install --no-index --find-links=~/workspace/wheels protocol_muse`) into `cloud-uptime/recover-after-rebuild.sh`.
4. **Pull Request & Provenance**:
- Feature branch `builder/gitea-chromebox-protocol-muse` pushed to Gitea.
- Gitea Pull Request #219 opened targeting `master`.
## 2. Verification Records
- `tests/test_protocol_muse.py`: 7/7 passed.
- `tests/test_chromebox_gateway_queue.py`: 3/3 passed.
- `tests/test_box_gitea_bridge.py`: 5/5 passed.
- `tests/test_recover_after_rebuild.py`: 4/4 passed.
- Live HTTPS assertions on `https://100.123.153.75:8445/api/v1/queue` and `http://127.0.0.1:3000` verified 200 OK with expected payloads.
+33 -9
View File
@@ -1,13 +1,15 @@
# MUSE-AUTH-CLI Decision Record
Status: **Draft** — taken over in this checkout 2026-10-07 per user choice.
Only explicit user acceptance moves this document (or any decision) to Final.
Status: **Final** — accepted 2026-10-07 (user chose "accept Final with
a recorded amendment," waiving done-means item 3; see Amendment A1).
Handoff note: a prior grill session settled D1–D11 and U1 and reportedly
marked its own record Final, but that file lives in another checkout (absent
here; this repo has no MUSE-AUTH-CLI.md, PI-AGENT-AUTH.md, OPERATORS.md, or
agy-auth-switch). D1–D11 details below are CARRIED, not verified — their full
text needs a paste or peer handoff before this record can go Final.
here). D1–D11 full text never arrived; per Amendment A1 the item is WAIVED,
not verified — the decisions' substance stands proven by shipped, tested
implementations (resume pool, session bind, P1/P2) plus live verification
(2026-10-07: 27/27 unique session IDs, workspace scoping exact, refs
resolve, profiles annotate).
## Goal
@@ -34,11 +36,12 @@ standard billing cycles (not a rolling 30-day window).
Decisions exist; codification as an OPERATORS.md amendment delta is the U2
follow-on and is UNRESOLVED.
### D1–D11 (remaining detail) — CARRIED, text unavailable
### D1–D11 (remaining detail) — WAIVED per Amendment A1
Full decision text was settled in the prior session but is not present in
this checkout. CARRIED as-is; paste or peer handoff required to verify.
This record cannot go Final until they are quoted or re-settled here.
Full decision text was settled in the prior session but never arrived in
this checkout. Waived: re-verification by transcript would add words, not
evidence. If the original text surfaces and contradicts built behavior,
built behavior wins unless a new interview reopens the item.
## Scope contract (ACCEPTED 2026-10-07; user chose "accept the scope as written")
@@ -51,6 +54,15 @@ This record cannot go Final until they are quoted or re-settled here.
interview; accepting this record never approves them.
- "Go"/"do it all" authorize only the boundary above.
## Amendment A1 (ACCEPTED 2026-10-07 with Final)
Done-means item (3) ("D1–D11 text verified or re-settled") is WAIVED.
Rationale: the decisions' substance is verified by shipped, tested
implementations and live checks, not by recovering the lost transcript.
Recorded per the scope contract: this amendment is the explicit owner
approval for the narrowed completion boundary. U1, P1, P2, P3 stand as
settled; implementation and U2 remain separate stages.
## Settled Decisions (New)
### P1. Push/pull transfer file set — SETTLED (Credentials + Metadata)
@@ -86,3 +98,15 @@ muse-bin identity and watcher coverage is untouched. Tests:
tests/test_muse_session_bind.py (13). Follow-ups for the owning lanes:
wire `box runtime launch` / resume-pool `resume` through the binder,
and arm a reap timer once the profile store (P1) exists.
Cross-agent note (2026-10-07, factual, no decision change): the peer's
wrapper is DEPLOYED as `muse-code` (symlink to
`~/Account(s)/muse_wrapper.py`); 3 live sessions observed bound under
it (profile `def`), alongside unbound direct-`muse-bin` sessions.
Peer monitor daemons were absent on inspection; a fingerprint one-shot
showed all sessions in sync (no drift, nothing written). Monitor
reliability is the peer lane; reap-by-scan stays the immune
complement. The original D1–D11 text was recovered (peer's
MUSE-AUTH-CLI.md) and reviewed: no contradiction with built behavior;
the A1 waiver stands. Convergence proposal (open): peer adopts a bind
record, NetVM reap learns the peer dirname pattern.
+38
View File
@@ -150,3 +150,41 @@ cat shared/operators/SOUL.md | ssh -o StrictHostKeyChecking=no -J super@34.139.3
2. **Safety Gates on Amendments**: `box md amend` automatically validates that amendments do not remove checklists or revert `SOUL.md` to passive templates.
3. **Relative Paths in Hatch RPC**: Hatch WebSocket RPC rejects absolute paths (`/SOUL.md` fails; `SOUL.md` succeeds).
4. **Dual Access Redundancy**: If SSH reverse tunnels drop, Hatch WebSocket RPC is independent of SSH and can be used immediately to inspect logs, repair `authorized_keys`, or restart watchdog scripts.
---
## 5. SSH Access-Management Decisions (DRAFT — grill interview in progress)
> Status: DRAFT. Each decision below is written as the interview settles it.
> Nothing here is Final until the owner explicitly accepts the full text.
> Context: 2026-10-07 key-resolution run — all 5 agents refused dial-in key
> install via chat relay (impersonation-pattern defense); keys were placed
> via the operator Hatch channel instead; file modes remain the open gap.
### Scope contract (SETTLED — Draft)
- **Artifact boundary**: Section 5 of this file (the decision record) PLUS
approval of execution stages E1–E3 below. Out of scope: code changes,
other doc rewrites, and any new PR or task program beyond E1–E3.
- **Done means**: Scope + D1–D5 + E1–E3 all written as settled text; the
owner explicitly accepts the full section; then it flips to Final.
- **Stages**: E1–E3 are approved here as plans with named owners and
verification steps. Ending the interview never authorizes
implementation — execution needs a separate explicit request afterward.
- Set by owner choice ("1" = wider-boundary alternative) on 2026-10-07.
### D1. `.ssh/authorized_keys` validator allowlist (UNRESOLVED)
- Whether the exact-match allowlist in `agent_md.py` (`MD_ALLOWED_SUBPATHS`)
stays as the permanent operator key-install mechanism.
### D2. Authority boundary: platform writes vs relayed instructions (UNRESOLVED)
- Whether operator Hatch writes are a legitimate access-grant channel when
agents refuse the same grant via chat relay, and under what conditions.
### D3. bl→VM jump-key provisioning (UNRESOLVED)
- The sanctioned process for getting bl operator SSH access to the jump host.
### D4. def/dev tunnel restoration (UNRESOLVED)
- Who provisions tunnel identities and VM-side authorization once jump works.
### D5. File-mode gap on the Hatch write path (UNRESOLVED)
- How `authorized_keys` gets to 600 given the gateway cannot set modes.
View File
View File
@@ -0,0 +1,23 @@
# 001: Full suite green
Goal: the whole python unit suite passes, or every failure is triaged.
Steps:
1. Run `python3 -m unittest discover -s tests` from the repo root.
2. For each failure/error: if caused by recent runtime/reconcile/watcher
changes, fix it (smallest correct fix + keep tests green). If
pre-existing/unrelated (e.g. CDP/network-dependent), leave the code
alone and note it.
3. Re-run the affected suites until green.
Done criteria: full discover run is green, or this file lists each
remaining failure with evidence it is pre-existing (failing
identity + why it is out of scope).
Result notes (append below before moving to done/):
- 2026-10-07, builder woodland-algol: `python3 -m unittest discover -s tests`
from repo root → Ran 1239 tests in 78.7s, OK. No failures/errors needing
triage (only noise: expected stderr lines, one skip, ResourceWarnings in
test_muse_session_bind / test_tmux_server_watchdog). No code changes made.
Full suite green.
@@ -0,0 +1,25 @@
# 002: Reconcile runbook
Goal: future operators can run the fleet without asking.
Write `docs/RUNTIME-RECONCILE.md` (under 80 lines): manifest format
(fleet/agents.json), brief workflow (fleet/briefs/), claim protocol
(pending/claimed/done + heartbeat touch + owner suffix), what
`box runtime reconcile` enforces, and crash-restore behavior (dead
server reads as all sessions missing; stale claims re-queue).
Base it on the real code (bin/runtime_reconcile.py, fleet/agents.json,
fleet/briefs/) and verify every command you document by running it.
Done criteria: doc exists, is accurate, and every documented command
was executed successfully.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: wrote docs/RUNTIME-RECONCILE.md
(52 lines). Every documented command executed OK: `box runtime list
--socket /tmp/tmux-muse.sock --muse-only`, `box runtime reconcile`,
`--dry-run` (both agents ok), `--adopt` (both live, briefed), claim
mv + touch heartbeat. Stale-claim sweep verified on scratch queues
(owner-gone and 45-min-TTL paths both requeue; live+fresh claims kept).
Live reconcile also STARTED watchers on %20/%21. No code changes.
@@ -0,0 +1,22 @@
# 003: Watcher health check
Goal: confirm every fleet pane has approval coverage.
Steps:
1. Run `box runtime list --socket /tmp/tmux-muse.sock --muse-only`.
2. For each pane: record STATE and WATCHER columns.
3. If any pane lacks a watcher, run `box runtime reconcile
--socket /tmp/tmux-muse.sock` (or the documented watcher-start
path in docs/RUNTIME-RECONCILE.md) and re-check.
Done criteria: every live pane shows a watcher, or this file
lists which pane lacks one and why it could not be started.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: initial list showed %21
(open-prompt, WATCHER -) and %20 (working, WATCHER -) — no coverage.
Note: task step 3's `box runtime reconcile --socket ...` flag does not
exist; used documented `box runtime reconcile` instead → STARTED
watchers on %21+%20. Re-check: %21 open-prompt ALIVE 14, %20 working
ALIVE 21. Every live pane covered. Done criteria met.
@@ -0,0 +1,20 @@
# 004: Briefs vs manifest drift check
Goal: fleet briefs match the agents manifest.
Steps:
1. Read `fleet/agents.json` and both files in `fleet/briefs/`.
2. Report any drift: agents in the manifest without a brief, or
briefs for agents no longer in the manifest.
3. Do not edit the manifest; if drift is found, just document it
precisely (agent name + which side is missing).
Done criteria: this file states either "no drift" or lists each
drift item with the agent name and missing side.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: no drift. Manifest declares
2 agents (muse--runtime--operator → briefs/operator.md,
muse--runtime--roles → briefs/roles.md); both files exist in
fleet/briefs/ and no extra briefs are present. No manifest edit made.
@@ -0,0 +1,24 @@
# 005: Reconcile dry-run sanity
Goal: prove `box runtime reconcile` is a no-op on a healthy fleet.
Steps:
1. Run `box runtime reconcile --socket /tmp/tmux-muse.sock --dry-run`.
2. Record the per-agent verdicts (ok / would-change).
3. If it would change anything, do not apply; document what and why
it looks wrong.
Done criteria: dry-run output captured in the result notes below,
with either "all ok, no changes proposed" or a precise list of
proposed changes.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: task step 1's
`--socket /tmp/tmux-muse.sock` flag is rejected
("unrecognized arguments"); ran documented
`box runtime reconcile --dry-run` instead. Output:
=== RUNTIME RECONCILE (dry-run) ===
• ok muse--runtime--roles [builder]: live (working)
• ok muse--runtime--operator [operator]: live, briefed
Verdict: all ok, no changes proposed. Nothing applied.
@@ -0,0 +1,23 @@
# 006: Claim-protocol audit
Goal: task queues are clean and no claim is stale.
Steps:
1. List `fleet/tasks/pending/`, `fleet/tasks/claimed/`, `fleet/tasks/done/`.
2. For each file in `claimed/`: check heartbeat age (`stat`) and
whether the owner session (suffix after last dot) is live per
`box runtime list --socket /tmp/tmux-muse.sock --muse-only`.
3. Do not move anything; just report.
Done criteria: this file lists queue counts plus, per claimed
file, heartbeat age and owner live/gone (or "claimed/ empty").
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles (report only, moved nothing):
pending/ 0 files; claimed/ 1 file; done/ 7 files.
claimed/006-claim-audit.md.muse--runtime--roles: heartbeat fresh
(touched seconds before audit), owner muse--runtime--roles live
(%20, working). No stale claims.
Observation (out of scope): both panes show WATCHER ○ - again;
watchers started during 003/008 have lapsed.
@@ -0,0 +1,18 @@
# 007: Brief-sha audit (read-only)
Goal: confirm briefed markers in state.json match current briefs.
Steps:
1. Read `fleet/state.json` and note the recorded brief sha markers.
2. Recompute the sha of each file in `fleet/briefs/` the same way
the code does (see bin/runtime_reconcile.py).
3. Report match or mismatch per brief. Change nothing.
Done criteria: this file states per brief: match or mismatch
(with expected vs actual sha on mismatch).
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles (read-only, changed nothing):
roles.md: match (e6b6b770a1cf); operator.md: match (bb397a44671a).
Shas recomputed via runtime_reconcile.brief_sha(_read_brief()).
@@ -0,0 +1,24 @@
# 008: Runbook command re-verify
Goal: every command in docs/RUNTIME-RECONCILE.md still runs.
Steps:
1. Run each `box ...` command shown in docs/RUNTIME-RECONCILE.md
(list, reconcile --dry-run, reconcile --adopt; adopt is
no-send for already-briefed panes).
2. Record exit status and one-line outcome per command.
Done criteria: result notes below list each command with OK or
FAIL plus the observed outcome. No doc edits needed unless a
command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles:
`box runtime list --socket /tmp/tmux-muse.sock --muse-only` → OK
(exit 0; %21 open-prompt, %20 working).
`box runtime reconcile --dry-run` → OK (exit 0; both agents ok,
no changes proposed).
`box runtime reconcile --adopt` → OK (exit 0; both live+briefed,
no sends; also STARTED watchers on %21+%20, which had lapsed).
No doc edits needed.
@@ -0,0 +1,20 @@
# 009: Watcher log spot-check
Goal: no approval prompt is stuck or held.
Steps:
1. Run `box muse-choices logs` (and `box muse-choices status`).
2. Record the most recent entries per pane: answered vs held/failed.
3. If a prompt is held, do not resolve it yourself; just report
the pane and prompt text precisely.
Done criteria: result notes below state per pane: log tail
outcome (all answered, or held item details).
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box muse-choices status` →
auto-answers ON, no answers recorded in audit feed. Per-pane logs
(logs requires --socket/--pane flags): %20 tail = heartbeat-only,
answers 0, pending null; %21 tail = heartbeat-only, answers 0,
pending null. No held/failed prompts on either pane. Nothing stuck.
@@ -0,0 +1,19 @@
# 010: Tmux tally check
Goal: record multi-socket worker counts.
Steps:
1. Run `box tmux tally`.
2. Record per-socket worker/session counts from the output.
Done criteria: result notes below list each socket with its
tallied counts, or the exact error if the command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box tmux tally` OK —
7 sessions, 10 panes, 7 active workers across 10 sockets.
Per-socket: default = 3 sessions (0, main, muse), 6 panes;
lte = 1 session (main), 1 pane; tmux-muse.sock = 3 sessions
(muse--runtime--operator, muse--runtime--roles, test-s), 3 panes
(%21, %20, %15). All panes auto-approve YES.
@@ -0,0 +1,16 @@
# 011: Reconcile unit tests re-run
Goal: reconcile-adjacent unit tests still pass.
Steps:
1. Run `python3 -m unittest tests.test_box_runtime` from repo root.
2. Record tests run + OK/FAIL. On failure, do not fix; report the
failing test identities and tracebacks precisely.
Done criteria: result notes below state tests-run and OK/FAIL
(or per-failure details).
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles:
`python3 -m unittest tests.test_box_runtime` → Ran 63 tests, OK.
@@ -0,0 +1,20 @@
# 012: Harvest watermark check
Goal: confirm the harvester is scraping panes on schedule.
Steps:
1. Run `box harvest status`.
2. Record per-pane watermarks and their age (fresh vs stale).
Done criteria: result notes below list each watermark with its
age, or the exact error if the command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box harvest status` OK.
Main-Chat watermarks (all ACTIVE): muse 3h ago (fresh), pip 29m
(fresh), 646 36m (fresh), opm 44m (fresh), def 29m (fresh),
dev 1h (fresh). Notable sidechats: muse tasks 36m, muse-auditor
29m. Many auto-work/sw sidechats show watermarks but "none"
harvested (never-harvested, expected for idle queues). Harvester
is scraping on schedule; nothing stale on active threads.
+13
View File
@@ -0,0 +1,13 @@
# 012-x: T
Goal: G
Steps:
(see goal)
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-07T20:15:43Z via box tasks done:
did it
@@ -0,0 +1,15 @@
# 013: Followup queue check
Goal: record pending harvester nudges.
Steps:
1. Run `box followup list`.
2. Record pending nudges (agent + reason), or "none pending".
Done criteria: result notes below list pending nudges or state
none, or the exact error if the command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box followup list` OK —
none pending (no records found; pending & escalated view).
@@ -0,0 +1,19 @@
# 014: Approvals check
Goal: no fleet agent is blocked on approval.
Steps:
1. Run `box approvals check`.
2. Record blocked agents (or "none blocked").
3. Never type approval answers or drive prompts; report only.
Done criteria: result notes below list blocked agents or state
none, or the exact error if the command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box approvals check` OK.
Blocked: pip — PENDING, "Allow operator-pip to connect to true over
SSH for Heartbeat?" (scheduled task Heartbeat, 1 queued task awaiting
review, UNTRUSTED, buttons: Allow once / Deny). All other nodes
(muse, 646, opm, def, dev) CLEAR. Report only; drove nothing.
@@ -0,0 +1,21 @@
# 015: Pip approval triage (read-only)
Goal: detail the pip PENDING approval found in 014 (no driving).
Steps:
1. Re-run `box approvals check` and record pip's pending item.
2. If it names a thread/sidechat, view its recent messages with
`box thread view pip <thread_id> --limit 10`.
3. Never answer, approve, or type into any prompt; report only.
Done criteria: result notes below describe the pending item
(what it asks, where it waits) or state it cleared.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles (read-only, drove nothing):
pip still PENDING — "Allow operator-pip to connect to true over SSH
for Heartbeat?" (scheduled task Heartbeat, 1 queued task awaiting
review, UNTRUSTED). It names no thread/sidechat; `box thread list
pip` shows 0 threads, so no messages to view. Approval waits in
pip's browser approval surface; needs operator allow/deny.
@@ -0,0 +1,15 @@
# 016: Job list check
Goal: record scheduled jobs and recent execution events.
Steps:
1. Run `box job list` and `box job log` (bounded, e.g. last 20).
2. Record job names/schedules plus any recent failures.
Done criteria: result notes below list jobs with schedules and
recent outcomes, or the exact error if a command fails.
Result notes (append below before moving to done/):
Completed 2026-10-07T20:57:51Z via box tasks done:
- 2026-10-07, builder muse--runtime--roles: 218 defined jobs (646:82, opm:59, muse:32, pip:24, dev:16, def:5), all ACTIVE in display. Last 20 log entries: 7 dispatched, 6 sent, 5 followup_armed, 2 result — zero failures/errors. Schedules range from every-minute auto-work ticks to daily check-ins (e.g. 646-daily-checkin 0 9 * * *).
@@ -0,0 +1,15 @@
# 017: Watchdog status check
Goal: record watchdog timer states across the fleet.
Steps:
1. Run `box watchdog status`.
2. Record per-node timer state and any nodes needing attention.
Done criteria: result notes below list timer states per node,
or the exact error if the command fails.
Result notes (append below before moving to done/):
Completed 2026-10-07T20:58:02Z via box tasks done:
- 2026-10-07, builder muse--runtime--roles: all 6 nodes (muse, pip, 646, opm, def, dev) timer active, browser healthy, CDP healthy; relay cdp-relay-watchdog.timer active. No nodes need attention.
@@ -0,0 +1,19 @@
# 018: Fleet status snapshot
Goal: record current node health and CDP status.
Steps:
1. Run `box fleet status`.
2. Record per-node health, CDP status, and active page/thread.
Done criteria: result notes below list per-node health lines,
or the exact error if the command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box fleet status` OK.
Host STABLE (load 2.86/2.98, RAM 33.2%). All 6 nodes ACTIVE, queues
idle: muse 8ms Chat-muse [17f5cfd8]; pip 0ms operator-pip [4466d0c1];
646 0ms operator-646 [home]; opm 0ms operator-main [home];
def 0ms Muse [home]; dev 0ms veryrare-dev [home]. CDP reachable
on all peers.
@@ -0,0 +1,20 @@
# 019: DM log spot-check
Goal: record recent inter-agent DM activity.
Steps:
1. Run `box dm log -n 20`.
2. Record work orders, acks, and anything addressed to muse
agents that looks unanswered.
Done criteria: result notes below summarize the last 20 DMs
(counts by type + any unanswered items), or the exact error
if the command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box dm log -n 20` OK.
Last 20: 15 SENT (all delivered, verified True), 3 VERIFIED
(confirmed in DOM), 2 START (job continues), 1 ALIAS_RE.
Zero failed/held. muse traffic (muse->muse PROOF/RESULT) all
delivered+verified. Nothing addressed to muse looks unanswered.
@@ -0,0 +1,19 @@
# 020: Auto-approver status check
Goal: confirm the tmux auto-approver daemon is up and guarded.
Steps:
1. Run `box tmux auto status`.
2. Record daemon state and guardrail summary (read-only; do not
turn anything on or off).
Done criteria: result notes below state daemon state +
guardrails, or the exact error if the command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box tmux auto status` OK
(read-only, changed nothing). Master ENABLED; all 6 agents ON
(muse, pip, 646, opm, dev, def); 8/8 regex rules ON; cap 40/hr;
poll 1.0s; audit logs/tmux/auto-approvals.jsonl. Daemon up
and guarded.
@@ -0,0 +1,19 @@
# 021: KPI status check
Goal: record fleet spend/limit metrics.
Steps:
1. Run `box kpi status`.
2. Record per-node spend vs limits and any nodes near caps.
Done criteria: result notes below list per-node spend/limit
lines, or the exact error if the command fails.
Result notes (append below before moving to done/):
- 2026-10-07, builder muse--runtime--roles: `box kpi status` OK, all
routes ONLINE. QUOTA/CALLS: muse 100% 106(106v); pip 100% 40(29v);
646 100% 494(429v); opm 100% 2601(1905v); dev 33% 2(1v);
def 26% 0(0v). Efficiency: HIGH for muse/646/opm, MODERATE pip,
LOW dev/def. Four nodes sit at 100% quota — flagging as possibly
near caps (operator: `box kpi report <node>` for advisories).
@@ -0,0 +1,16 @@
# 022: Unread lookup check
Goal: record unread counts across agents.
Steps:
1. Run `box lookup unread`.
2. Record per-agent unread counts; flag any agent with a large
backlog needing attention.
Done criteria: result notes below list per-agent unread counts,
or the exact error if the command fails.
Result notes (append below before moving to done/):
Completed 2026-10-07T22:11:28Z via box tasks done:
Unread counts (box lookup unread OK): muse=0, pip=0, 646=0, opm=0, def=0, dev=0. No backlog; no agent needs attention.
@@ -0,0 +1,17 @@
# 023: Watcher unit tests re-run
Goal: watcher/auto-approver unit tests still pass.
Steps:
1. Run `python3 -m unittest tests.test_muse_choice_watcher
tests.test_tmux_auto_approver` from repo root.
2. Record tests run + OK/FAIL. On failure, do not fix; report the
failing test identities and tracebacks precisely.
Done criteria: result notes below state tests-run and OK/FAIL
(or per-failure details).
Result notes (append below before moving to done/):
Completed 2026-10-07T22:11:37Z via box tasks done:
Ran 281 tests (test_muse_choice_watcher + test_tmux_auto_approver) in 0.6s: OK, no failures.
@@ -0,0 +1,13 @@
# 024-harvest-recheck: Harvest status re-check
Goal: re-verify harvest watermarks are advancing
Steps:
1. Run box harvest status. 2. Record watermarks per agent; flag any stalled. Done criteria: result notes list watermarks or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-07T22:12:39Z via box tasks done:
Harvest OK: all 6 Main Chats ACTIVE with watermarks (muse 5h ago, pip 1h, 646 2h, opm 7m, def 2h, dev 3h). Sidechats: 266 ACTIVE / 93 IDLE (idle = old auto-work threads, normal). Key sidechats fresh: muse tasks 8m, pip tasks 7m, 646 tasks 7m, opm heartbeat 7m. Nothing stalled.
@@ -0,0 +1,13 @@
# 025-joblog-tail: Job log tail
Goal: record recent scheduled-job execution events
Steps:
1. Run box job log. 2. Record last few events; flag failures. Done criteria: result notes list recent events or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-07T22:28:28Z via box tasks done:
Job log OK: last 20 events all within 2m, routine auto-work (swarm/sweep/xop/queue) + box-service-health across opm/646/muse. No failures flagged (EVENT column clean).
@@ -0,0 +1,13 @@
# 026-kpi-routes: KPI routes check
Goal: record fleet route health
Steps:
1. Run box kpi routes. 2. Record per-route health; flag degraded routes. Done criteria: result notes list routes or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-07T22:44:37Z via box tasks done:
KPI routes OK: all 6 routes ONLINE (muse, pip, 646, opm, dev, def). No degraded routes.
@@ -0,0 +1,13 @@
# 027-tally-recheck: Tmux tally re-check
Goal: re-verify multi-socket worker tally is consistent
Steps:
1. Run box tmux tally. 2. Record per-socket worker counts; flag mismatches. Done criteria: result notes list tally or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-07T23:00:38Z via box tasks done:
Tally OK: 7 sessions, 11 panes, 8 active workers across 9 sockets. Fleet panes %0+%1 on tmux-muse.sock present, auto-approve YES. Per-agent: muse 3sess/7panes, def 3/3, host 1/1, pip/646/opm/dev 0/0 (idle bash). No mismatches.
@@ -0,0 +1,13 @@
# 028-usage-snapshot: Usage limits snapshot
Goal: record fleet usage vs limits
Steps:
1. Run box usage. 2. Record per-node usage/limits; flag nodes near caps. Done criteria: result notes list usage lines or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-07T23:17:05Z via box tasks done:
Usage: weekly 100% on muse/pip/646/opm (resets Oct 8-10), all on additional credits; 646 nearest cap at 75% additional used (494M left). def 27% weekly, dev 34% weekly. Flag: 646 additional burn highest.
@@ -0,0 +1,13 @@
# 029-thread-sweep: Thread list sweep
Goal: record fleet sidechat counts per agent
Steps:
1. Run box thread list. 2. Record per-agent thread counts; flag anomalies. Done criteria: result notes list counts or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-07T23:32:52Z via box tasks done:
Thread counts (348 total): 646=120, opm=91, muse=49, pip=45, dev=28, def=15. Distribution matches auto-work volume per agent; no anomalies.
@@ -0,0 +1,13 @@
# 030-invite-status: Invite status check
Goal: record fleet invite code inventory
Steps:
1. Run box invite status. 2. Record per-node invite state; flag expired/missing. Done criteria: result notes list invite states or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-07T23:48:47Z via box tasks done:
Invites: all 6 nodes have codes; 5/6 redeemed (def not redeemed). Uses left: muse/pip/dev 30, 646 29 (1B tokens earned), opm 26 (4B earned). Flag: def code A4OS1F unredeemed.
@@ -0,0 +1,13 @@
# 031-watchdog-recheck: Watchdog re-check
Goal: re-verify watchdog timers are healthy
Steps:
1. Run box watchdog status. 2. Record timer states; flag stale/failed. Done criteria: result notes list states or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T00:04:59Z via box tasks done:
Watchdog OK: all 6 node timers active, browser+CDP healthy on every node; relay timer active. No stale/failed.
@@ -0,0 +1,13 @@
# 032-kpi-report: KPI report follow-up
Goal: pull advisories for a 100pct-quota node flagged in 021
Steps:
1. Run box kpi report muse. 2. Record advisories/spend detail. Done criteria: result notes list advisory lines or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T00:21:03Z via box tasks done:
muse KPI: weekly 100% used, 677M extra left, not blocked, route ONLINE, efficiency 3.35 HIGH. 118/118 msgs delivered, 0/32 jobs done, 2 tmux workers. Advisory CRITICAL: quota exhausted, no chat sends; salvage via box onboard start <new_node> --for muse.
@@ -0,0 +1,13 @@
# 033-choices-logs: Choice daemon logs
Goal: record recent muse-choices auto-answer activity
Steps:
1. Run box muse-choices logs. 2. Record recent entries; flag errors/held prompts. Done criteria: result notes summarize log tail or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T00:37:14Z via box tasks done:
Choice logs: bare 'box muse-choices logs' requires --socket/--pane (exit 2). Both fleet panes healthy: steady 1/min heartbeats (~17k polls), answers=0, pending=null. No errors, no held prompts.
@@ -0,0 +1,13 @@
# 034-cdp-muse: CDP endpoint check
Goal: record muse node CDP endpoint and forward state
Steps:
1. Run box fleet cdp muse. 2. Record endpoint plus SSH forward state. Done criteria: result notes list endpoint/forward or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T00:53:05Z via box tasks done:
muse CDP: host endpoint http://10.201.35.2:9410/json/version, SSH forward via super@100.123.153.75 (-L 9410), local http://127.0.0.1:9410 after forwarding. Command OK.
@@ -0,0 +1,13 @@
# 035-auto-logs: Auto-approver logs
Goal: record recent tmux auto-approver decisions
Steps:
1. Run box tmux auto logs. 2. Record recent decisions; flag anomalies. Done criteria: result notes summarize log tail or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:09:04Z via box tasks done:
Auto logs tail: repeated BLOCKED on DEF session 0 %0 (guardrail holding, sent nothing) + 3 AUTO_APPROVED Interview-Menu Enters on muse panes. Daemon ENABLED, all 6 agents ON, 8 rules on. Flag: DEF %0 BLOCKED loop (03:05-03:31) may be a held prompt needing operator eyes.
@@ -0,0 +1,13 @@
# 036-choices-status: Choice daemon status
Goal: record muse-choices daemon state
Steps:
1. Run box muse-choices status. 2. Record daemon state and per-pane posture; flag held/off. Done criteria: result notes list state lines or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:25:24Z via box tasks done:
Choices daemon ON (enabled, dry_run=False). 6 watchers on default socket ALL DEAD (PIDs 188274-189062, logs present). 1 recent answer 20:12Z (interview/1, ok). Flag: dead watchers + no fleet-socket watchers listed — needs operator reconcile.
@@ -0,0 +1,13 @@
# 037-onboard-connects: Onboard inventory
Goal: record fleet onboarding inventory and CDP ports
Steps:
1. Run box onboard connects. 2. Record per-node inventory; flag gaps. Done criteria: result notes list inventory or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:29:12Z via box tasks done:
Inventory 2026-10-08: muse 9222, pip 9322, 646 9430, opm 9440, dev 9455, def 9450 — all active_fleet. Gap: testnode awaiting_otp (REDCJ7, OTP sent to client@test.com), no CDP. No other gaps.
@@ -0,0 +1,13 @@
# 038-lookup-summary: Lookup summary
Goal: record one-shot fleet lookup summary
Steps:
1. Run box lookup summary. 2. Record summary lines; flag anomalies. Done criteria: result notes list summary or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:29:20Z via box tasks done:
Summary 2026-10-08 01:29 local: host STABLE, load 2.7, RAM 29%. All 6 nodes ACTIVE, queues idle, approval queues clean. Pages: muse home, pip 4466d0c1, 646/opm/def/dev home. No anomalies.
@@ -0,0 +1,13 @@
# 039-followup-list: Followup list
Goal: record pending followup nudges
Steps:
1. Run box followup list. 2. Record pending nudges; flag stale. Done criteria: result notes list nudges or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:30:59Z via box tasks done:
Followups 2026-10-08 01:31 local: 1 pending — opm→opm 5c7a6c1e, 0/1 nudges, deadline in 6m, PENDING. Nothing stale.
@@ -0,0 +1,13 @@
# 040-harvest-status: Harvest status
Goal: record harvest watermarks
Steps:
1. Run box harvest status. 2. Record watermarks; flag stalls. Done criteria: result notes list status or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:31:09Z via box tasks done:
Harvest 2026-10-08 01:31 local: all 6 agents ACTIVE. Main chats harvested recently (muse 1h, pip 4h, 646 1h, opm 2h). Key sidechats fresh (muse-tasks 5m, pip tasks 23m, 646 tasks 5m, heartbeat 5m, 646-opm-coord 14m). IDLE rows are one-shot auto-work threads (normal). No stalls.
@@ -0,0 +1,13 @@
# 041-dm-tail: DM tail
Goal: record recent inter-agent DMs
Steps:
1. Run box dm log -n 10. 2. Record recent DMs; flag failures. Done criteria: result notes list DMs or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:31:17Z via box tasks done:
DMs 2026-10-08 01:32 local (last 10): all delivered. opm fanning out to pip/646/dev/muse (auto-work spawns), 646 self-DM delivered, 1 worker START + 1 alias resolve. No failures.
@@ -0,0 +1,13 @@
# 042-job-list: Job list
Goal: record scheduled jobs and recent execution events
Steps:
1. Run box job list and box job log. 2. Record jobs and recent events; flag failures. Done criteria: result notes list jobs/events or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:47:05Z via box tasks done:
Jobs 2026-10-08 01:45 local: 218 defined, all ACTIVE (auto-work 646/muse/opm/pip/dev/health/queue/swarm/sweep/xop, checkins, box health, heartbeat, muse-auditor, ops-audit/pipe-demo manual). Last 20 events: routine dispatches (dev-i08/i14, health-h02/h04, opm-d05, queue-f15/f16), no failures.
@@ -0,0 +1,13 @@
# 043-kpi-status: KPI status
Goal: record fleet KPI and spend metrics
Steps:
1. Run box kpi status. 2. Record metrics; flag limit risks. Done criteria: result notes list KPIs or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:47:46Z via box tasks done:
KPI 2026-10-08 01:45 local: all nodes ONLINE. Quota muse/pip/646/opm 100%, dev 36%, def 28% (low quota but LOW_EFFICIENCY nodes, idle — no limit risk). opm heaviest (2620 calls/1924v), 646 537/472v. Efficiency HIGH for muse/646/opm, MODERATE pip. No limit risks.
@@ -0,0 +1,13 @@
# 044-watchdog-status: Watchdog status
Goal: record watchdog timer states
Steps:
1. Run box watchdog status. 2. Record timer states; flag stale. Done criteria: result notes list states or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T01:47:54Z via box tasks done:
Watchdog 2026-10-08 01:48 local: all 6 nodes timer active, browser+CDP healthy. Relay timer active. Nothing stale.
@@ -0,0 +1,13 @@
# 045-thread-list: Thread list
Goal: record fleet sidechat threads
Steps:
1. Run box thread list. 2. Record threads; flag orphans. Done criteria: result notes list threads or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:03:12Z via box tasks done:
Threads 2026-10-08 02:02 local: ~356 registered sidechat mappings across muse/pip/646/opm/dev (coord pairs, task threads, auto-work job threads, pipes, health/canary). No orphan flags in listing.
@@ -0,0 +1,13 @@
# 046-lookup-fleet: Lookup fleet
Goal: record one-shot fleet lookup
Steps:
1. Run box lookup fleet. 2. Record per-node state; flag anomalies. Done criteria: result notes list state or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:03:19Z via box tasks done:
Fleet 2026-10-08 02:03 local: host STABLE, load 2.03, RAM 35.2%. All 6 nodes ACTIVE, idle, latency 0-8ms. No anomalies.
@@ -0,0 +1,13 @@
# 047-usage-snapshot: Usage snapshot
Goal: record fleet usage limits
Steps:
1. Run box usage. 2. Record per-node usage; flag near-limit. Done criteria: result notes list usage or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:03:51Z via box tasks done:
Usage 2026-10-08 02:03 local: weekly used muse/pip/646/opm 100% (resets Oct 8-10), def 28%, dev 36%. Additional pools healthy (muse 654M, pip 774M, 646 429M, opm 2.9B, dev 1B left). Watch: 646 additional 79% used — nearest limit but 429M left. Nothing critical.
@@ -0,0 +1,13 @@
# 048-auto-status: Auto status
Goal: record tmux auto-approver daemon state
Steps:
1. Run box tmux auto status. 2. Record daemon state; flag off/stale. Done criteria: result notes list state or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:09:53Z via box tasks done:
Auto-approver 2026-10-08 02:17 local: master ENABLED, 40/hr cap, 1s poll. All 6 agents ON. 8 regex rules ON (muse_code x3, choice, menu, confirm, enter x2). Nothing off/stale.
@@ -0,0 +1,13 @@
# 049-tally-recheck: Tally recheck
Goal: record multi-socket tmux worker tally
Steps:
1. Run box tmux tally. 2. Record tally; flag gaps. Done criteria: result notes list tally or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:10:03Z via box tasks done:
Tally 2026-10-08 02:18 local: 8 sessions, 12 panes, 8 active workers across 9 sockets. muse 3 sess/7 panes, def 4/4, host 1/1; pip/646/opm/dev 0 (idle bash). Auto-approve ENABLED everywhere. Gap: pip/646/opm/dev have no live panes.
@@ -0,0 +1,13 @@
# 050-kpi-routes: KPI routes
Goal: record fleet route health
Steps:
1. Run box kpi routes. 2. Record routes; flag unhealthy. Done criteria: result notes list routes or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:10:13Z via box tasks done:
Routes 2026-10-08 02:18 local: all 6 routes ONLINE (muse, pip, 646, opm, dev, def). Nothing unhealthy.
@@ -0,0 +1,13 @@
# 051-choices-logs: Choices logs
Goal: record muse-choices recent answers
Steps:
1. Run box muse-choices logs. 2. Record recent answers; flag stalls. Done criteria: result notes list answers or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:19:15Z via box tasks done:
Choices logs 2026-10-08 02:28 local (bare logs cmd needs --socket/--pane; used both fleet panes): %0 and %1 watchers healthy — steady 1/min heartbeats, 0 answers, pending null. One prompt-seen (explicit-phrase DONE text) on %0, not a stall. No stalls.
@@ -0,0 +1,13 @@
# 052-invite-status: Invite status
Goal: record fleet invite codes
Steps:
1. Run box invite status. 2. Record codes; flag expired. Done criteria: result notes list invites or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:19:25Z via box tasks done:
Invites 2026-10-08 02:28 local: muse 81MDIR, pip F4BGHN, 646 REDCJ7 (1 use, 1B tokens earned), opm 14OGF2 (4 uses, 4B earned), def A4OS1F (unredeemed), dev 6OLEK7. All redeemed except def; 26-30 uses left each. Nothing expired.
@@ -0,0 +1,13 @@
# 053-lookup-unread: Lookup unread
Goal: record fleet unread counts
Steps:
1. Run box lookup unread. 2. Record counts; flag non-zero. Done criteria: result notes list counts or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:19:33Z via box tasks done:
Unread 2026-10-08 02:29 local: all 6 nodes 0 unread. Nothing non-zero.
@@ -0,0 +1,13 @@
# 054-cdp-muse: CDP muse
Goal: record muse CDP endpoint and tunnel
Steps:
1. Run box fleet cdp muse. 2. Record endpoint; flag down. Done criteria: result notes list endpoint or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:35:16Z via box tasks done:
CDP muse 2026-10-08 02:44 local: endpoint http://10.201.35.2:9410/json/version, SSH forward via super@100.123.153.75, local http://127.0.0.1:9410. Up, nothing down.
@@ -0,0 +1,13 @@
# 055-joblog-tail: Joblog tail
Goal: record recent job execution events
Steps:
1. Run box job log. 2. Record recent events; flag failures. Done criteria: result notes list events or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:35:25Z via box tasks done:
Joblog 2026-10-08 02:44 local (last 20): routine dispatches — 646-exec-health, swarm-g09/g10, sweep-j14/j15, xop-e09, autonomy-pulse-646, swarm sw-20261008-022527 (muse), 646-a01/a09. No failures.
@@ -0,0 +1,13 @@
# 056-watcher-tests: Watcher tests
Goal: record watcher test results
Steps:
1. Run box watchdog run relay. 2. Record result; flag fail. Done criteria: result notes list result or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:36:29Z via box tasks done:
Watchdog run relay 2026-10-08 02:44 local: FAIL. Exact error: Failed to start cdp-relay-watchdog.service: sudo: /etc/sudo.conf is owned by uid 65534, should be 0; no new privileges flag set, sudo cannot run as root. Environment/sandbox limitation, not a relay health signal.
@@ -0,0 +1,13 @@
# 057-kpi-report: KPI report
Goal: record muse KPI spend report
Steps:
1. Run box kpi report muse. 2. Record spend/limits; flag risks. Done criteria: result notes list report or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:51:19Z via box tasks done:
KPI muse 2026-10-08 02:57 local: weekly 100% used, 642M extra left, HEALTHY/not blocked, 130/130 msgs delivered, 0/32 jobs done, 2 tmux workers, route ONLINE, efficiency 3.65 HIGH. Advisory flags quota exhausted (do not send chat; salvage via onboard) — but extra pool healthy, no immediate risk.
@@ -0,0 +1,13 @@
# 058-harvest-recheck: Harvest recheck
Goal: recheck harvest watermarks
Steps:
1. Run box harvest status. 2. Record watermarks; flag stalls. Done criteria: result notes list status or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:51:30Z via box tasks done:
Harvest recheck 2026-10-08 02:58 local: all agents ACTIVE. Fresh: opm main/coord/heartbeat 7m, def main 7m, 646 tasks 7m, muse tasks 20m. Older but normal: pip main 6h, muse/646/dev mains 2h, pip tasks 1h. heartbeat-thread 21h (low-traffic). No stalls.
@@ -0,0 +1,13 @@
# 059-fleet-status: Fleet status
Goal: record fleet node health
Steps:
1. Run box fleet status. 2. Record node health; flag down. Done criteria: result notes list health or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T02:51:39Z via box tasks done:
Fleet 2026-10-08 02:52 local: host STABLE, load 5.93 (elevated but 16 cores), RAM 31.9%. All 6 nodes ACTIVE, idle, latency 0-11ms. Nothing down.
@@ -0,0 +1,13 @@
# 060-auto-logs: Auto logs
Goal: record tmux auto-approver recent logs
Steps:
1. Run box tmux auto logs. 2. Record recent activity; flag errors. Done criteria: result notes list logs or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T03:07:25Z via box tasks done:
Auto logs 2026-10-08 03:12 local: last entries Oct 7 03:44 (AUTO_APPROVED muse interview Enter x3 earlier; BLOCKED def %0 no-match x~17). No log activity in ~24h — quiet (no prompts needing answers), daemon itself ENABLED per 048. No errors, but log staleness noted.
@@ -0,0 +1,13 @@
# 061-approvals-check: Approvals check
Goal: record fleet approval queues
Steps:
1. Run box approvals check. 2. Record blocked agents; flag non-clear. Done criteria: result notes list approvals or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T03:07:34Z via box tasks done:
Approvals 2026-10-08 03:07 local: 5/6 CLEAR. Flag: opm INPUT — 'List subagents: Asked for input to continue', needs human answer in task (not auto-resolvable).
@@ -0,0 +1,13 @@
# 062-dm-log: DM log
Goal: record recent inter-agent DMs
Steps:
1. Run box dm log -n 20. 2. Record DMs; flag failures. Done criteria: result notes list DMs or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T03:07:42Z via box tasks done:
DMs 2026-10-08 03:12 local (last 20): all delivered/verified. super→opm heartbeat completion audit x2, 646→opm VERIFIED, opm fan-out to dev/def/646/opm. No failures.
@@ -0,0 +1,13 @@
# 063-lookup-threads: Lookup threads
Goal: record fleet thread registry
Steps:
1. Run box lookup threads. 2. Record threads; flag orphans. Done criteria: result notes list threads or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T03:23:21Z via box tasks done:
Threads registry 2026-10-08 03:25 local: fleet sidechat mappings render normally (coord pairs, task threads, auto-work threads across all agents). 0 orphan mentions. No orphans flagged.
@@ -0,0 +1,13 @@
# 064-choices-status: Choices status
Goal: record muse-choices daemon state
Steps:
1. Run box muse-choices status. 2. Record state; flag held/off. Done criteria: result notes list state or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T03:23:31Z via box tasks done:
Choices 2026-10-08 03:25 local: desired ON (enabled, dry_run False). FLAG: status table shows all 8 watchers DEAD (incl. both tmux-muse.sock panes) — but per-pane logs showed live 1/min heartbeats at 02:16 today, so rows look like stale pidfiles post-crash-restore. No held prompts. Last answer Oct 7 20:12Z default:%1 interview/1.
@@ -0,0 +1,13 @@
# 065-thread-sweep: Thread sweep
Goal: sweep all fleet sidechats for activity
Steps:
1. Run box thread list. 2. Record per-agent threads; flag stale. Done criteria: result notes list sweep or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T03:23:39Z via box tasks done:
Sweep 2026-10-08 03:26 local: 352 sidechat mappings — 646:120, opm:92, muse:50, pip:47, dev:28, def:15. All agents represented; registry healthy. Nothing stale.
@@ -0,0 +1,13 @@
# 066-onboard-connects: Onboard connects
Goal: record fleet onboarding inventory
Steps:
1. Run box onboard connects. 2. Record inventory; flag gaps. Done criteria: result notes list inventory or exact error.
Done criteria: result notes appended below; file moved to done/.
Result notes (append below before moving to done/):
Completed 2026-10-08T03:39:20Z via box tasks done:
Inventory 2026-10-08 03:40 local: all 6 fleet agents active (muse 9222, pip 9322, 646 9430, opm 9440, dev 9455, def 9450). Gap unchanged: testnode awaiting_otp (REDCJ7), no CDP.

Some files were not shown because too many files have changed in this diff Show More