Files
box/bin/box-ctl.py

4764 lines
184 KiB
Python
Executable File
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
#!/usr/bin/env python3
"""box-ctl.py — allowlisted bl helper for Box API mutations.
The board server (VM) never runs raw systemctl or shell over SSH.
All bl mutations go through this helper, invoked as:
/home/super/Projects/NetVM/bin/box-ctl.py <action> [args...]
Security properties:
- Fixed verb set; every argument validated before acting.
- <name> must match ^[a-z0-9-]{1,64}$ (kills path traversal).
- Unit files generated from a fixed template; only the validated
name is interpolated. User-controlled strings never appear in
unit files.
- No shell=True anywhere. No string interpolation into commands.
- Job JSON schema-validated before writing; changes git-committed.
- Every action audit-logged to box-ctl.jsonl with caller identity.
Output: JSON to stdout ({"ok": true, ...} or {"ok": false, ...}),
exit 0 on success, nonzero on failure.
"""
import json
import os
import re
import shutil
import subprocess
import sys
import time
from datetime import datetime, timezone
from pathlib import Path
NETVM_ROOT = Path("/home/super/Projects/NetVM")
BIN = NETVM_ROOT / "bin"
if str(BIN) not in sys.path:
sys.path.insert(0, str(BIN))
JOBS_DIR = NETVM_ROOT / "jobs"
SYSTEMD_USER = Path.home() / ".config" / "systemd" / "user"
CTL_LOG = NETVM_ROOT / "box-ctl.jsonl"
JOB_LOG = NETVM_ROOT / "job-log.jsonl"
DISPATCHER = BIN / "job-dispatch.py"
DM_PY = BIN / "dm.py"
POLICY_FILE = NETVM_ROOT / "CHAT_POLICY.md"
DM_LOG = NETVM_ROOT / "dm-log.jsonl"
WATCHDOG_STATE = NETVM_ROOT / "main-chat-watchdog.state"
NAME_RE = re.compile(r"^[a-z0-9-]{1,64}$")
VAR_NAME_RE = re.compile(r"^[A-Za-z0-9_.-]{1,64}$")
VALID_AGENTS = {"muse", "pip", "646", "opm", "dev", "def"}
THREAD_RE = re.compile(r"^[A-Za-z0-9-]{1,64}$")
THREAD_OPS = ("list", "pin", "unpin", "archive", "unarchive")
# Default sidechat per agent for `box notify`. Mirrors super-cli.py
# DEFAULT_AGENT_SIDECHATS (kept in sync manually; super-cli is the
# canonical copy). Sidechat-first policy: notify never targets main
# chat unless --allow-main-chat is passed explicitly.
NOTIFY_SIDECHATS = {
"646": "646 tasks",
"opm": "main-loop brain",
"pip": "pip tasks",
"muse": "muse tasks",
"dev": "onboarding-dev",
"def": "def tasks",
}
VALID_ON_FAILURE = {"retry", "alert", "ignore"}
KNOWN_PLACEHOLDERS = {"job_id", "job_name", "datetime", "date", "last_run"}
PROTOCOL_LITERALS = ("[REQ", "[CONFIRM", "[JOB", "[RESULT")
DOW_NAMES = ["Sun", "Mon", "Tue", "Wed", "Thu", "Fri", "Sat"]
def utcnow():
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def out(ok, **kw):
payload = {"ok": ok}
payload.update(kw)
print(json.dumps(payload))
def fail(code, error, detail=None, exit_code=1):
payload = {"ok": False, "code": code, "error": error}
if detail is not None:
payload["detail"] = detail
print(json.dumps(payload))
sys.exit(exit_code)
def audit(action, name=None):
"""Append {ts, action, name, caller} to the bl-side log."""
try:
entry = {
"ts": utcnow(),
"action": action,
"name": name,
"caller": os.environ.get("BOX_CALLER", "unknown"),
}
with open(CTL_LOG, "a") as f:
f.write(json.dumps(entry) + "\n")
except Exception:
pass # audit failure must not break the action
def check_name(name):
if not name or not NAME_RE.match(name):
fail("BAD_NAME",
"name must match ^[a-z0-9-]{1,64}$",
{"field": "name", "value": name})
return name
def run(cmd, timeout=30):
"""Run a command with no shell. Returns CompletedProcess."""
return subprocess.run(cmd, capture_output=True, text=True, timeout=timeout)
# ---------------------------------------------------------------------------
# Cron validation and conversion
# ---------------------------------------------------------------------------
CRON_RANGES = {
0: (0, 59), # minute
1: (0, 23), # hour
2: (1, 31), # day of month
3: (1, 12), # month
4: (0, 7), # day of week (0 and 7 = Sunday)
}
def validate_cron_field(field, idx):
"""Validate one cron field. Returns True/False."""
lo, hi = CRON_RANGES[idx]
if field == "*":
return True
for part in field.split(","):
# step: base/n
if "/" in part:
base, step = part.split("/", 1)
if not step.isdigit() or int(step) < 1:
return False
if base != "*" and not validate_cron_field(base, idx):
return False
continue
# range: a-b
if "-" in part:
a, b = part.split("-", 1)
if not (a.isdigit() and b.isdigit()):
return False
if not (lo <= int(a) <= hi and lo <= int(b) <= hi):
return False
if int(a) > int(b):
return False
continue
# single value
if not part.isdigit():
return False
if not (lo <= int(part) <= hi):
return False
return True
def validate_cron(schedule):
"""Validate a 5-field cron string. Returns (ok, reason)."""
parts = schedule.split()
if len(parts) != 5:
return False, "schedule must have exactly 5 fields (M H dom mon dow)"
for i, p in enumerate(parts):
if not validate_cron_field(p, i):
return False, f"schedule field {i + 1} ({p!r}) is invalid"
return True, ""
def cron_to_oncalendar(schedule):
"""Convert common cron patterns to systemd OnCalendar.
Returns the OnCalendar string, or None if the pattern is not
supported (caller rejects with INVALID_SCHEDULE).
"""
m, h, dom, mon, dow = schedule.split()
def is_star(f):
return f == "*"
# */n * * * * -> *:0/n (every n minutes)
if m.startswith("*/") and m[2:].isdigit() and all(is_star(f) for f in (h, dom, mon, dow)):
return f"*:0/{m[2:]}"
# M * * * * -> hourly or *:M
if m.isdigit() and all(is_star(f) for f in (h, dom, mon, dow)):
if m == "0":
return "hourly"
return f"*:{int(m):02d}"
# M H * * * -> HH:MM (daily)
if m.isdigit() and h.isdigit() and all(is_star(f) for f in (dom, mon, dow)):
return f"{int(h):02d}:{int(m):02d}"
# M H * * DOW -> Dow HH:MM (weekly)
if m.isdigit() and h.isdigit() and dow != "*" and is_star(dom) and is_star(mon):
# take first dow value for the weekly form
d = dow.split(",")[0].split("-")[0].split("/")[0]
if d.isdigit() and 0 <= int(d) <= 7:
return f"{DOW_NAMES[int(d) % 7]} {int(h):02d}:{int(m):02d}"
return None
# M H DOM * * -> *-*-DOM HH:MM:00 (monthly)
if m.isdigit() and h.isdigit() and dom.isdigit() and is_star(mon) and is_star(dow):
return f"*-*-{int(dom):02d} {int(h):02d}:{int(m):02d}:00"
# M,M... * * * * -> *:M1,M2...
if all(part.isdigit() for part in m.split(",")) and all(is_star(f) for f in (h, dom, mon, dow)):
formatted_m = ",".join(f"{int(part):02d}" for part in m.split(","))
return f"*:{formatted_m}"
return None
def validate_oncalendar(expr):
"""Check an OnCalendar expression with systemd-analyze."""
r = run(["systemd-analyze", "calendar", expr, "--iterations=2"], timeout=15)
return r.returncode == 0
# ---------------------------------------------------------------------------
# Job schema validation (§6 of the design doc)
# ---------------------------------------------------------------------------
def validate_job(data):
"""Validate a job definition dict. Returns (ok, field, reason)."""
if not isinstance(data, dict):
return False, None, "job must be a JSON object"
# name
name = data.get("name")
if not isinstance(name, str) or not NAME_RE.match(name):
return False, "name", "must match ^[a-z0-9-]{1,64}$"
# description
desc = data.get("description")
if desc is not None:
if not isinstance(desc, str) or len(desc) > 280:
return False, "description", "must be a string ≤ 280 chars"
# schedule
sched = data.get("schedule")
if not isinstance(sched, str):
return False, "schedule", "is required and must be a string"
if sched != "manual":
ok, reason = validate_cron(sched)
if not ok:
return False, "schedule", reason
oc = cron_to_oncalendar(sched)
if oc is None:
return False, "schedule", (
"cron pattern not supported for OnCalendar conversion; "
"supported: */n * * * *, M * * * *, M H * * *, M H * * DOW, M H DOM * *"
)
if not validate_oncalendar(oc):
return False, "schedule", f"converted OnCalendar {oc!r} rejected by systemd-analyze"
# scheduler
scheduler = data.get("scheduler", "systemd")
if scheduler != "systemd":
return False, "scheduler", 'v1 accepts only "systemd"'
# agent
agent = data.get("agent")
if agent not in VALID_AGENTS:
return False, "agent", f"must be one of {sorted(VALID_AGENTS)}"
# prompt_template
pt = data.get("prompt_template")
if not isinstance(pt, str) or not (1 <= len(pt) <= 4000):
return False, "prompt_template", "must be a string of 1–4000 chars"
for lit in PROTOCOL_LITERALS:
if lit in pt:
return False, "prompt_template", (
f"must not contain {lit!r} literal (protocol framing is added by the dispatcher)"
)
for ph in re.findall(r"\{([a-zA-Z_][a-zA-Z0-9_]*)\}", pt):
if ph not in KNOWN_PLACEHOLDERS:
return False, "prompt_template", f"unknown placeholder {{{ph}}}"
# timeout
timeout = data.get("timeout", 300)
if not isinstance(timeout, int) or isinstance(timeout, bool) or not (60 <= timeout <= 3600):
return False, "timeout", "must be an integer 60–3600"
# on_failure
onf = data.get("on_failure", "alert")
if onf not in VALID_ON_FAILURE:
return False, "on_failure", f"must be one of {sorted(VALID_ON_FAILURE)}"
# chain_next
cn = data.get("chain_next")
if cn is not None:
if not isinstance(cn, str) or not NAME_RE.match(cn):
return False, "chain_next", "must be null or a valid job name"
if not (JOBS_DIR / f"{cn}.json").exists():
return False, "chain_next", f"job {cn!r} does not exist (dangling chain)"
# sidechat
sc = data.get("sidechat")
if sc is not None:
if not isinstance(sc, dict):
return False, "sidechat", "must be an object"
create = sc.get("create", False)
if not isinstance(create, bool):
return False, "sidechat.create", "must be a boolean"
nt = sc.get("name_template")
if nt is not None:
if not isinstance(nt, str) or len(nt) > 120:
return False, "sidechat.name_template", "must be a string ≤ 120 chars"
for ph in re.findall(r"\{([a-zA-Z_][a-zA-Z0-9_]*)\}", nt):
if ph not in {"job_name", "date", "job_id"}:
return False, "sidechat.name_template", f"unknown placeholder {{{ph}}}"
rk = sc.get("reuse_key")
if rk is not None:
if not isinstance(rk, str) or not NAME_RE.match(rk):
return False, "sidechat.reuse_key", "must match ^[a-z0-9-]{1,64}$"
return True, None, ""
# ---------------------------------------------------------------------------
# Unit file templates (fixed; only <name> interpolated, regex-validated)
# ---------------------------------------------------------------------------
SERVICE_TEMPLATE = """[Unit]
Description=Dispatch {name} job
After=network.target
[Service]
Type=oneshot
ExecStart=/home/super/Projects/NetVM/bin/job-dispatch.py {name}
User=super
WorkingDirectory=/home/super/Projects/NetVM
"""
TIMER_TEMPLATE = """[Unit]
Description=Run {name} job on schedule
[Timer]
OnCalendar={oncalendar}
Persistent=true
AccuracySec=1min
[Install]
WantedBy=timers.target
"""
def unit_path(name, kind):
return SYSTEMD_USER / f"job-{name}.{kind}"
def systemctl(*args, timeout=30):
return run(["systemctl", "--user"] + list(args), timeout=timeout)
# ---------------------------------------------------------------------------
# Actions
# ---------------------------------------------------------------------------
def act_timer_list():
timers = []
try:
entries = sorted(SYSTEMD_USER.glob("job-*.timer"))
except Exception:
entries = []
for p in entries:
unit = p.name
name = unit[len("job-"):-len(".timer")]
active = systemctl("is-active", unit).stdout.strip() == "active"
enabled_out = systemctl("is-enabled", unit).stdout.strip()
enabled = enabled_out == "enabled"
job_path = JOBS_DIR / f"{name}.json"
timers.append({
"name": name,
"unit": unit,
"active": active,
"enabled": enabled,
"orphan": not job_path.exists(),
})
audit("timer-list")
out(True, timers=timers)
def _timer_show_property(unit, prop):
r = systemctl("show", unit, f"--property={prop}", "--value")
if r.returncode != 0:
return None
return r.stdout.strip() or None
def act_timer_status(name):
check_name(name)
unit = f"job-{name}.timer"
if not unit_path(name, "timer").exists():
fail("NOT_FOUND", f"no such timer: {unit}")
active = systemctl("is-active", unit).stdout.strip() == "active"
enabled = systemctl("is-enabled", unit).stdout.strip() == "enabled"
last_usec = _timer_show_property(unit, "LastTriggerUSec")
next_usec = _timer_show_property(unit, "NextElapseUSec")
oncalendar_raw = _timer_show_property(unit, "TimersCalendar")
# TimersCalendar returns "{ OnCalendar=... ; next_elapse=... }" — extract just the expr
oncalendar = oncalendar_raw
if oncalendar_raw:
m = re.search(r"OnCalendar=([^;}]+)", oncalendar_raw)
if m:
oncalendar = m.group(1).strip()
result = _timer_show_property(unit, "Result")
def usec_to_iso(v):
if not v or v in ("0", "n/a"):
return None
try:
ts = int(v) / 1_000_000
return datetime.fromtimestamp(ts, tz=timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
except Exception:
return None
# last_result from job-log.jsonl
last_result = "unknown"
try:
if JOB_LOG.exists():
for line in reversed(JOB_LOG.read_text().splitlines()[-500:]):
try:
e = json.loads(line)
except Exception:
continue
if e.get("job_name") != name and not str(e.get("job_id", "")).startswith(name + "-"):
continue
t = e.get("type")
if t in ("job_dispatched",):
last_result = "success"
break
if t in ("job_failed", "job_timeout"):
last_result = "failed"
break
except Exception:
pass
job_def = None
job_path = JOBS_DIR / f"{name}.json"
if job_path.exists():
try:
job_def = json.loads(job_path.read_text())
except Exception:
job_def = None
audit("timer-status", name)
out(True, name=name, unit=unit, active=active, enabled=enabled,
last_run=usec_to_iso(last_usec), next_run=usec_to_iso(next_usec),
last_result=last_result, oncalendar=oncalendar,
job_definition=job_def)
def act_timer_create(name):
check_name(name)
job_path = JOBS_DIR / f"{name}.json"
if not job_path.exists():
fail("NOT_FOUND", f"job definition missing: jobs/{name}.json — create the job first")
try:
job = json.loads(job_path.read_text())
except Exception as e:
fail("INVALID_JOB", f"job JSON unreadable: {e}")
ok, field, reason = validate_job(job)
if not ok:
fail("INVALID_JOB", f"job validation failed: {field}: {reason}",
{"field": field, "reason": reason})
if unit_path(name, "timer").exists():
fail("ALREADY_EXISTS", f"timer job-{name}.timer already exists — delete it first")
if job.get("schedule") == "manual":
fail("INVALID_SCHEDULE", "cannot create a systemd timer for a job with schedule 'manual'")
oncalendar = cron_to_oncalendar(job["schedule"])
if oncalendar is None or not validate_oncalendar(oncalendar):
fail("INVALID_SCHEDULE", "schedule did not convert to a valid OnCalendar")
SYSTEMD_USER.mkdir(parents=True, exist_ok=True)
unit_path(name, "service").write_text(SERVICE_TEMPLATE.format(name=name))
unit_path(name, "timer").write_text(
TIMER_TEMPLATE.format(name=name, oncalendar=oncalendar))
r = systemctl("daemon-reload")
if r.returncode != 0:
fail("TIMER_CREATE_FAILED", f"daemon-reload failed: {r.stderr.strip()}")
r = systemctl("enable", "--now", f"job-{name}.timer")
if r.returncode != 0:
fail("TIMER_CREATE_FAILED", f"enable --now failed: {r.stderr.strip()}")
audit("timer-create", name)
out(True, name=name, unit=f"job-{name}.timer", created=True,
oncalendar=oncalendar)
def act_timer_delete(name, keep_job=False):
check_name(name)
unit = f"job-{name}.timer"
if not unit_path(name, "timer").exists():
fail("NOT_FOUND", f"no such timer: {unit}")
systemctl("stop", unit)
systemctl("disable", unit)
for kind in ("timer", "service"):
p = unit_path(name, kind)
try:
p.unlink()
except FileNotFoundError:
pass
r = systemctl("daemon-reload")
if r.returncode != 0:
fail("TIMER_CREATE_FAILED", f"daemon-reload failed: {r.stderr.strip()}")
audit("timer-delete", name)
out(True, name=name, unit=unit, deleted=True, job_kept=keep_job or (JOBS_DIR / f"{name}.json").exists())
def act_timer_control(name, op):
check_name(name)
unit = f"job-{name}.timer"
if not unit_path(name, "timer").exists():
fail("NOT_FOUND", f"no such timer: {unit}")
valid = {"start", "stop", "enable", "disable"}
if op not in valid:
fail("BAD_NAME", f"unknown timer op: {op}")
args = ["--now", unit] if op == "enable" else [unit]
r = systemctl(op, *args)
if r.returncode != 0:
fail("TIMER_CREATE_FAILED", f"systemctl {op} failed: {r.stderr.strip()}")
audit(f"timer-{op}", name)
out(True, name=name, unit=unit, op=op)
def _job_summary(name, path):
try:
d = json.loads(path.read_text())
except Exception:
return {"name": name, "error": "unreadable"}
return {
"name": name,
"description": d.get("description"),
"schedule": d.get("schedule"),
"agent": d.get("agent"),
"timeout": d.get("timeout", 300),
}
def act_job_list(include_archived=False):
jobs = []
if JOBS_DIR.exists():
for p in sorted(JOBS_DIR.glob("*.json")):
jobs.append(_job_summary(p.stem, p))
if include_archived:
arch_dir = JOBS_DIR / "archive"
if arch_dir.exists():
for p in sorted(arch_dir.glob("*.json")):
s = _job_summary(p.stem, p)
s["archived"] = True
jobs.append(s)
audit("job-list")
out(True, jobs=jobs)
def act_job_archive(name, force=False):
check_name(name)
script = BIN / "retention-archive-jobs.py"
cmd = [sys.executable, str(script), "archive", name, "--json"]
if force:
cmd.append("--force")
r = run(cmd)
try:
data = json.loads(r.stdout)
except Exception:
fail("ARCHIVE_FAILED", r.stderr or r.stdout)
if data.get("status") in ("archived", "dry_run"):
audit("job-archive", name)
out(True, **data)
else:
fail("ARCHIVE_FAILED", data.get("error") or data.get("status"))
def act_job_unarchive(name):
check_name(name)
script = BIN / "retention-archive-jobs.py"
cmd = [sys.executable, str(script), "unarchive", name, "--json"]
r = run(cmd)
try:
data = json.loads(r.stdout)
except Exception:
fail("UNARCHIVE_FAILED", r.stderr or r.stdout)
if data.get("status") in ("unarchived", "dry_run"):
audit("job-unarchive", name)
out(True, **data)
else:
fail("UNARCHIVE_FAILED", data.get("error") or data.get("status"))
def act_job_get(name):
check_name(name)
p = JOBS_DIR / f"{name}.json"
if not p.exists():
fail("NOT_FOUND", f"no such job: {name}")
try:
data = json.loads(p.read_text())
except Exception as e:
fail("INVALID_JOB", f"job JSON unreadable: {e}")
audit("job-get", name)
out(True, job=data)
def _git(*args):
return run(["git", "-C", str(NETVM_ROOT)] + list(args), timeout=30)
def act_job_put(name):
check_name(name)
raw = sys.stdin.read()
try:
data = json.loads(raw)
except Exception as e:
fail("INVALID_JOB", f"stdin is not valid JSON: {e}")
if data.get("name") != name:
fail("NAME_MISMATCH", "path name != body name",
{"path": name, "body": data.get("name")})
ok, field, reason = validate_job(data)
if not ok:
fail("INVALID_JOB", f"schema validation failed: {field}: {reason}",
{"field": field, "reason": reason})
existed = (JOBS_DIR / f"{name}.json").exists()
JOBS_DIR.mkdir(parents=True, exist_ok=True)
(JOBS_DIR / f"{name}.json").write_text(json.dumps(data, indent=2) + "\n")
r = _git("add", f"jobs/{name}.json")
if r.returncode != 0:
fail("DISPATCH_FAILED", f"git add failed: {r.stderr.strip()}")
msg = f"{'Update' if existed else 'Add'} job {name} via box-ctl"
r = _git("-c", "user.name=box-ctl", "-c", "user.email=box-ctl@bl.local",
"commit", "-m", msg)
if r.returncode != 0 and "nothing to commit" not in (r.stdout + r.stderr):
fail("DISPATCH_FAILED", f"git commit failed: {r.stderr.strip()}")
audit("job-put", name)
out(True, name=name, created=not existed, updated=existed)
def act_job_delete(name, force=False):
check_name(name)
p = JOBS_DIR / f"{name}.json"
if not p.exists():
fail("NOT_FOUND", f"no such job: {name}")
if unit_path(name, "timer").exists() and not force:
fail("TIMER_STILL_ACTIVE",
f"timer job-{name}.timer still exists — delete it first or use --force")
r = _git("rm", "-q", f"jobs/{name}.json")
if r.returncode != 0:
fail("DISPATCH_FAILED", f"git rm failed: {r.stderr.strip()}")
r = _git("-c", "user.name=box-ctl", "-c", "user.email=box-ctl@bl.local",
"commit", "-m", f"Delete job {name} via box-ctl")
if r.returncode != 0:
fail("DISPATCH_FAILED", f"git commit failed: {r.stderr.strip()}")
audit("job-delete", name)
out(True, name=name, deleted=True)
def act_job_trigger(name):
check_name(name)
p = JOBS_DIR / f"{name}.json"
if not p.exists():
fail("NOT_FOUND", f"no such job: {name}")
r = run([sys.executable, str(DISPATCHER), name], timeout=300)
if r.returncode != 0:
fail("DISPATCH_FAILED", f"job-dispatch.py failed: {(r.stderr or r.stdout).strip()[-500:]}")
audit("job-trigger", name)
out(True, name=name, triggered=True)
# ---------------------------------------------------------------------------
# Recursive workflow: box -> agent -> box -> agent ...
# The box stays in the loop: every RESULT recorded here advances the chain
# through the box (shared dedupe state with response-harvester.py).
# ---------------------------------------------------------------------------
CHAINED_JOBS_FILE = BIN / "chained-jobs.json"
def _chain_has(job_id):
try:
if CHAINED_JOBS_FILE.exists():
return job_id in json.loads(CHAINED_JOBS_FILE.read_text())
except Exception:
pass
return False
def _chain_mark(job_id):
try:
c = []
if CHAINED_JOBS_FILE.exists():
c = json.loads(CHAINED_JOBS_FILE.read_text())
if job_id not in c:
c.append(job_id)
c = c[-1000:]
CHAINED_JOBS_FILE.write_text(json.dumps(c))
except Exception:
pass
def _job_name_from_id(job_id):
m = re.match(r"^(.*)-(\d{8}-\d{6}-[a-f0-9]{8})$", str(job_id))
if m:
return m.group(1)
parts = str(job_id).split("-")
if len(parts) < 3:
return None
return "-".join(parts[:-2])
def _chain_target(job_name, success):
"""Return (next_job, via_field) for a completed job, or (None, None)."""
p = JOBS_DIR / f"{job_name}.json"
if not p.exists():
return None, None
try:
cfg = json.loads(p.read_text())
except Exception:
return None, None
if success:
nxt = cfg.get("on_success") or cfg.get("chain_next")
via = "on_success" if cfg.get("on_success") else "chain_next"
else:
nxt = cfg.get("on_failure")
via = "on_failure"
if nxt and isinstance(nxt, str) and (JOBS_DIR / f"{nxt}.json").exists():
return nxt, via
return None, None
def _dispatch_chained(next_job, prev_job_id, prev_result):
"""Dispatch the next job in a chain. Returns (ok, err)."""
env = os.environ.copy()
env["CHAIN_PREV_JOB_ID"] = prev_job_id
env["CHAIN_PREV_RESULT"] = (prev_result or "")[:1000]
try:
r = subprocess.run([sys.executable, str(DISPATCHER), next_job],
capture_output=True, text=True, timeout=300, env=env)
except Exception as e:
return False, f"dispatch exception: {e}"
if r.returncode != 0:
return False, (r.stderr or r.stdout).strip()[-500:]
return True, None
def act_job_result(job_id):
"""Agent -> box: record a RESULT and advance the chain (the recursion step)."""
raw = sys.stdin.read()
try:
data = json.loads(raw)
except Exception as e:
fail("INVALID_RESULT", f"stdin is not valid JSON: {e}")
agent = data.get("agent")
if agent not in VALID_AGENTS:
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
success = data.get("success", True)
if not isinstance(success, bool):
fail("INVALID_RESULT", "success must be a boolean")
summary = data.get("summary", "")
if not isinstance(summary, str) or len(summary) > 2000:
fail("INVALID_RESULT", "summary must be a string of 0-2000 chars")
rec = {
"ts": utcnow(),
"type": "job_result",
"job_id": job_id,
"agent": agent,
"success": success,
"result_snippet": summary[:300],
"thread_id": None,
"msg_id": None,
"source": "box-ctl",
}
try:
with open(JOB_LOG, "a") as f:
f.write(json.dumps(rec) + "\n")
except Exception as e:
fail("LOG_FAILED", f"could not append job-log.jsonl: {e}")
# recursion: the box dispatches the next step
job_name = _job_name_from_id(job_id)
chained_to = None
chain_via = None
chain_dispatched = False
chain_error = None
if job_name and not _chain_has(job_id):
_chain_mark(job_id)
nxt, via = _chain_target(job_name, success)
if nxt:
try:
ccfg = json.loads((JOBS_DIR / f"{job_name}.json").read_text())
except Exception:
ccfg = {}
try:
step_delay = int(ccfg.get("step_delay", 0) or 0)
except Exception:
step_delay = 0
if step_delay > 0:
time.sleep(min(step_delay, 30))
ok, err = _dispatch_chained(nxt, job_id, summary)
chain_dispatched = ok
chain_error = err
if ok:
chained_to = nxt
chain_via = via
audit("job-result", job_id)
out(True, job_id=job_id, agent=agent, success=success, recorded=True,
chained_to=chained_to, chain_via=chain_via,
chain_dispatched=chain_dispatched, chain_error=chain_error)
def act_job_status(name):
"""Box-side view of a job's run state: dispatches, results, chain wiring."""
events = []
try:
if JOB_LOG.exists():
for line in JOB_LOG.read_text().splitlines()[-2000:]:
try:
e = json.loads(line)
except Exception:
continue
jn = e.get("job_name")
jid = str(e.get("job_id", ""))
if jn == name or jid == name or jid.startswith(name + "-"):
events.append(e)
except Exception as e:
fail("LOG_FAILED", f"could not read job-log.jsonl: {e}")
dispatches = [e for e in events if e.get("type") == "job_dispatched"]
results = [e for e in events if e.get("type") == "job_result"]
failures = [e for e in events if e.get("type") in ("job_failed", "job_timeout")]
last_result = results[-1] if results else None
state = "pending"
if last_result:
state = "resulted"
elif dispatches:
state = "dispatched"
chain = {}
p = JOBS_DIR / f"{name}.json"
if p.exists():
try:
cfg = json.loads(p.read_text())
for k in ("chain_next", "on_success", "on_failure"):
if cfg.get(k):
chain[k] = cfg[k]
except Exception:
pass
audit("job-status", name)
out(True, name=name, state=state,
dispatches=len(dispatches), results=len(results), failures=len(failures),
last_dispatch=dispatches[-1].get("ts") if dispatches else None,
last_result=({
"ts": last_result.get("ts"),
"job_id": last_result.get("job_id"),
"agent": last_result.get("agent"),
"success": last_result.get("success"),
"snippet": (last_result.get("result_snippet") or "")[:200],
} if last_result else None),
chain=chain)
def act_job_next(job_id, success=None):
"""Dry-run: what WOULD the box dispatch next for this job_id? No dispatch."""
job_name = _job_name_from_id(job_id)
if not job_name:
fail("BAD_NAME", f"cannot parse job name from job_id: {job_id}")
if success is None:
success = True
try:
if JOB_LOG.exists():
for line in reversed(JOB_LOG.read_text().splitlines()[-2000:]):
try:
e = json.loads(line)
except Exception:
continue
if e.get("type") == "job_result" and str(e.get("job_id")) == job_id:
success = bool(e.get("success", True))
break
except Exception:
pass
nxt, via = _chain_target(job_name, success)
already = _chain_has(job_id)
audit("job-next", job_id)
out(True, job_id=job_id, job_name=job_name, success=success,
would_dispatch=bool(nxt) and not already,
next_job=nxt, via=via, already_chained=already)
def act_job_chain(frm, to, on_failure=False):
"""Wire recursion: set chain_next (or on_failure) from one job to another."""
check_name(frm)
check_name(to)
if frm == to:
fail("INVALID_JOB", "a job cannot chain to itself")
fp = JOBS_DIR / f"{frm}.json"
if not fp.exists():
fail("NOT_FOUND", f"no such job: {frm}")
if not (JOBS_DIR / f"{to}.json").exists():
fail("NOT_FOUND", f"no such job: {to}")
try:
cfg = json.loads(fp.read_text())
except Exception as e:
fail("INVALID_JOB", f"job JSON unreadable: {e}")
field = "on_failure" if on_failure else "chain_next"
cfg[field] = to
# cycle guard
seen = {frm}
cur = to
while cur:
if cur in seen:
fail("INVALID_JOB", f"chain would create a cycle at {cur!r}")
seen.add(cur)
try:
c = json.loads((JOBS_DIR / f"{cur}.json").read_text())
except Exception:
break
cur = c.get("chain_next") or c.get("on_success")
fp.write_text(json.dumps(cfg, indent=2) + "\n")
r = _git("add", f"jobs/{frm}.json")
if r.returncode != 0:
fail("DISPATCH_FAILED", f"git add failed: {r.stderr.strip()}")
r = _git("-c", "user.name=box-ctl", "-c", "user.email=box-ctl@bl.local",
"commit", "-m", f"Chain job {frm} -> {to} via box-ctl")
if r.returncode != 0 and "nothing to commit" not in (r.stdout + r.stderr):
fail("DISPATCH_FAILED", f"git commit failed: {r.stderr.strip()}")
audit("job-chain", f"{frm}->{to}")
out(True, **{"from": frm, "to": to, "field": field, "chained": True})
def _valid_sidechat_name(name):
# Sidechat names may contain spaces ("646 tasks"); reject only
# control characters and enforce a sane length. dm.py resolves
# the name (alias -> job-sidechats.json -> autoprovision).
return bool(name) and len(name) <= 64 and not any(
ord(c) < 32 or ord(c) == 127 for c in name)
def act_notify(agent, message, sidechat=None, allow_main_chat=False, sender=None):
if agent not in VALID_AGENTS:
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
if sender is not None and sender not in VALID_AGENTS:
fail("BAD_NAME", f"sender must be one of {sorted(VALID_AGENTS)}")
sender_agent = sender if sender else "opm"
if not message or len(message) > 1000:
fail("INVALID_JOB", "message must be 1\u20131000 chars")
# Sidechat-first policy: default to the agent's sidechat, never main.
if allow_main_chat:
target = "main"
extra = ["--allow-main-chat"]
else:
target = sidechat if sidechat else NOTIFY_SIDECHATS.get(agent)
if not target:
fail("BAD_NAME", f"no default sidechat for agent {agent}; pass --sidechat <name>")
if sidechat and not _valid_sidechat_name(sidechat):
fail("BAD_NAME", "sidechat name must be 1-64 chars, no control characters")
extra = []
r = run([sys.executable, str(DM_PY), "send",
"--agent", sender_agent, "--to", agent, "--target", target,
*extra, message],
timeout=120)
if r.returncode != 0:
fail("DISPATCH_FAILED", f"dm.py send failed: {(r.stderr or r.stdout).strip()[-500:]}")
audit("notify", agent)
out(True, agent=agent, target=target, sent=True)
def act_main_loop(sub, agent=None):
"""Run the main-chat self-monitor loop (self_main_loop.py).
check: one loop iteration -> JSON {ok, prompted, errors, new_total, agents}
status: watermark / last run / prompt-sidechat config, no chat reads.
enable: enable the loop (per-agent or all); check skips disabled agents.
disable: disable the loop (per-agent or all).
"""
if sub not in ("check", "status", "enable", "disable"):
fail("BAD_NAME", "usage: main-loop check|status|enable|disable [--agent <name>]")
if agent is not None and agent not in VALID_AGENTS:
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
audit("main-loop-" + sub + (f"-{agent}" if agent else ""))
script = BIN / "self_main_loop.py"
cmd = [sys.executable, str(script), sub]
if agent:
cmd += ["--agent", agent]
try:
r = subprocess.run(cmd,
capture_output=True, text=True, timeout=600)
except subprocess.TimeoutExpired:
fail("LOOP_TIMEOUT", "self_main_loop.py timed out")
try:
data = json.loads(r.stdout.strip())
except Exception:
fail("LOOP_ERROR", "self_main_loop.py returned non-JSON",
{"stdout": (r.stdout or "")[-300:], "stderr": (r.stderr or "")[-300:]})
ok = bool(data.pop("ok", False))
out(ok, **data)
def act_watchdog_alerts(no_advance=False):
"""Check for new failed browser relaunches since the watermark.
Runs bin/watchdog-alert-check.sh, which exits 0 (no new failures)
or 1 (new failures printed to stdout) and maintains its own
watermark at NETVM_ROOT/watchdog-alert-watermark.txt.
no_advance=True passes --no-advance to the script: a peek-only read
that leaves the watermark untouched (for web-surface polling).
Default False preserves the classic CLI advance-on-read semantics.
"""
audit("watchdog-alerts")
script = BIN / "watchdog-alert-check.sh"
cmd = [str(script)] + (["--no-advance"] if no_advance else [])
r = subprocess.run(cmd, capture_output=True, text=True, timeout=30)
failures = [l for l in (r.stdout or "").splitlines() if l.strip()]
# Exit 0 = no new failures, exit 1 = new failures found. Any other
# exit code is a script error.
if r.returncode not in (0, 1):
fail("WATCHDOG_CHECK_ERROR", "watchdog-alert-check.sh failed",
{"stderr": (r.stderr or "").strip()[-500:]})
out(True, new_failures=len(failures), failures=failures)
def act_fleet_status():
audit("fleet-status")
cmd = [sys.executable, str(BIN / "super-cli.py"), "fleet", "status", "--json"]
r = subprocess.run(cmd, capture_output=True, text=True)
if r.returncode == 0:
try:
data = json.loads(r.stdout)
print(json.dumps(data))
return
except Exception:
pass
fail("FLEET_ERROR", "failed to collect fleet data", {"stderr": r.stderr})
def act_approval_check(node=None):
audit("approval-check", node)
if node is not None and node not in VALID_AGENTS:
fail("BAD_NODE", f"unknown node: {node}")
import approvals
nodes = [node] if node else list(VALID_AGENTS)
res = approvals.check_fleet_approvals(nodes)
out(True, approvals=res)
def act_approval_allow(node, always=False, force=False, message=None, allow_main_chat=False):
audit("approval-allow", node)
if node not in VALID_AGENTS:
fail("BAD_NODE", f"unknown node: {node}")
import approvals
res = approvals.allow_node_approval(node, always=always, force=force, caller="box-ctl",
message=message, allow_main_chat=allow_main_chat)
if res.get("ok"):
kw = {k: v for k, v in res.items() if k != "ok"}
out(True, **kw)
else:
fail("APPROVAL_FAILED", res.get("error", "approval failed"), res)
def act_approval_deny(node, message=None, allow_main_chat=False):
audit("approval-deny", node)
if node not in VALID_AGENTS:
fail("BAD_NODE", f"unknown node: {node}")
import approvals
res = approvals.deny_node_approval(node, caller="box-ctl",
message=message, allow_main_chat=allow_main_chat)
if res.get("ok"):
kw = {k: v for k, v in res.items() if k != "ok"}
out(True, **kw)
else:
fail("APPROVAL_FAILED", res.get("error", "deny failed"), res)
def act_approval_auto(node=None):
audit("approval-auto", node)
if node is not None and node not in VALID_AGENTS:
fail("BAD_NODE", f"unknown node: {node}")
import approvals
nodes = [node] if node else list(VALID_AGENTS)
res = approvals.auto_approve_fleet(nodes, caller="box-ctl")
kw = {k: v for k, v in res.items() if k != "ok"}
out(True, **kw)
def act_tmux_tally():
audit("tmux-tally")
try:
import tmux_auto_approver
tally = tmux_auto_approver.gather_tmux_tally()
out(True, **tmux_auto_approver.asdict(tally))
except Exception as e:
fail("TMUX_ERROR", f"tmux tally failed: {e}")
def act_tmux_auto_status():
audit("tmux-auto-status")
try:
import tmux_auto_approver
st = tmux_auto_approver.AutoApproverState.load()
out(True, **tmux_auto_approver.asdict(st))
except Exception as e:
fail("TMUX_ERROR", f"tmux auto status failed: {e}")
def act_tmux_auto_toggle(enable=True, node=None, session=None):
audit("tmux-auto-toggle", f"{'on' if enable else 'off'}:{node or session or 'all'}")
if node and node not in VALID_AGENTS:
fail("BAD_NODE", f"unknown node: {node}")
try:
import tmux_auto_approver
st = tmux_auto_approver.AutoApproverState.load()
if node:
st.agents_enabled[node] = bool(enable)
elif session:
st.sessions_enabled[session] = bool(enable)
else:
st.global_enabled = bool(enable)
if enable:
for a in VALID_AGENTS:
st.agents_enabled[a] = True
st.save()
out(True, enabled=bool(enable), node=node, session=session,
global_enabled=st.global_enabled, agents_enabled=st.agents_enabled)
except Exception as e:
fail("TMUX_ERROR", f"tmux auto toggle failed: {e}")
def act_tmux_auto_once(dry_run=False):
audit("tmux-auto-once", f"dry_run={dry_run}")
try:
import tmux_auto_approver
runner = tmux_auto_approver.AutoApproverRunner(dry_run=dry_run)
actions = runner.run_once()
out(True, actions=actions, dry_run=dry_run, count=len(actions))
except Exception as e:
fail("TMUX_ERROR", f"tmux auto run once failed: {e}")
def act_onboard_connects():
audit("onboard-connects")
try:
import onboard_pipeline
connects = onboard_pipeline.get_all_connects()
out(True, connects=connects, count=len(connects))
except Exception as e:
fail("ONBOARD_ERROR", f"onboard connects failed: {e}")
def act_relay_health():
"""Run relay-health-check.sh and return JSON results."""
audit("relay-health")
script = BIN / "relay-health-check.sh"
r = subprocess.run([str(script)], capture_output=True, text=True, timeout=60)
results = {}
for line in r.stdout.strip().split("\n"):
if ":" in line:
name, code = line.split(":", 1)
results[name.strip()] = code.strip()
all_healthy = r.returncode == 0
out(all_healthy, relays=results, healthy=all_healthy)
def act_identity_audit():
"""Run identity-audit-check.sh and return drift as JSON."""
audit("identity-audit")
script = BIN / "identity-audit-check.sh"
r = subprocess.run([str(script)], capture_output=True, text=True, timeout=60)
if r.returncode == 2:
fail("AUDIT_UNAVAILABLE", (r.stderr or r.stdout).strip()[:500])
drift = []
for line in r.stdout.strip().split("\n"):
if line.startswith("DRIFT: "):
drift.append(line[len("DRIFT: "):])
clean = r.returncode == 0
out(clean, drift=drift, clean=clean)
def act_cdp_latency():
"""Run cdp-latency-check.sh and return per-node latency JSON."""
audit("cdp-latency")
script = BIN / "cdp-latency-check.sh"
try:
r = subprocess.run([str(script)], capture_output=True, text=True, timeout=120)
except subprocess.TimeoutExpired:
fail("LATENCY_TIMEOUT", "cdp-latency-check.sh timed out")
nodes = []
for line in r.stdout.strip().split("\n"):
parts = line.strip().split(":")
if len(parts) != 3:
continue
name, lat, code = parts
entry = {"node": name.strip(), "http_code": code.strip()}
if lat.strip() == "FAIL":
entry["ok"] = False
entry["latency_ms"] = None
else:
try:
entry["latency_ms"] = int(lat.strip())
except ValueError:
entry["latency_ms"] = None
entry["ok"] = code.strip() == "200"
nodes.append(entry)
all_ok = bool(nodes) and all(n["ok"] for n in nodes)
out(all_ok, nodes=nodes)
def act_chrome_errors(no_advance=False):
"""Run chrome-error-scan.sh and return per-profile error counts as JSON.
no_advance=True passes --no-advance: a peek-only read that leaves the
watermark untouched (for web-surface polling). Default False preserves
the classic CLI advance-on-read semantics.
"""
audit("chrome-errors")
cmd = ["/bin/bash", str(BIN / "chrome-error-scan.sh"), "--json"]
if no_advance:
cmd.append("--no-advance")
r = subprocess.run(cmd, capture_output=True, text=True)
if r.returncode == 0:
try:
data = json.loads(r.stdout)
out(True, profiles=data)
return
except Exception:
pass
fail("SCAN_ERROR", "chrome-error-scan.sh failed", {"stderr": r.stderr})
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})
def act_unread(agent=None):
if agent is not None and agent not in VALID_AGENTS:
fail("BAD_NODE", f"unknown agent: {agent}")
audit("unread", agent or "all")
cmd = [sys.executable, str(BIN / "super-cli.py"), "lookup", "unread", "--json"]
r = subprocess.run(cmd, capture_output=True, text=True)
if r.returncode == 0:
try:
data = json.loads(r.stdout)
if agent:
data["nodes"] = [n for n in data.get("nodes", []) if n.get("node") == agent]
print(json.dumps(data))
return
except Exception:
pass
fail("UNREAD_ERROR", "failed to collect unread counts", {"stderr": r.stderr})
# ---------------------------------------------------------------------------
# No-SSH agent development: git visibility + test runs + ack.
# ---------------------------------------------------------------------------
GIT_BIN = shutil.which("git")
GIT_MAX_DIFF = 64 * 1024
GIT_MAX_STATUS = 200
TESTS_MAX_OUTPUT = 32 * 1024
TEST_MODULE_RE = re.compile(r"^tests\.[a-z0-9_]+$")
def _git_path(value):
"""Validate a repo-relative path for git subcommands (no escapes)."""
if not value or not isinstance(value, str):
fail("BAD_NAME", "path must be a non-empty string")
if ".." in value or value.startswith("/") or \
any(ord(c) < 32 or ord(c) == 127 for c in value):
fail("BAD_NAME", "path must be repo-relative without '..'")
try:
resolved = (NETVM_ROOT / value).resolve()
resolved.relative_to(NETVM_ROOT.resolve())
except (OSError, ValueError):
fail("BAD_NAME", "path escapes the repository")
return value
def _run_git(args, timeout=30):
if GIT_BIN is None:
fail("GIT_ERROR", "git executable not found")
try:
return subprocess.run([GIT_BIN, "-C", str(NETVM_ROOT), *args],
capture_output=True, text=True, timeout=timeout)
except subprocess.TimeoutExpired:
fail("GIT_ERROR", "git command timed out")
except OSError as e:
fail("GIT_ERROR", f"cannot run git: {e}")
def act_git_status():
audit("git-status")
r = _run_git(["status", "--short", "--branch"])
if r.returncode != 0:
fail("GIT_ERROR", "git status failed", {"stderr": (r.stderr or "")[:500]})
lines = (r.stdout or "").splitlines()
branch = "unknown"
if lines and lines[0].startswith("## "):
branch = lines[0][3:].split("...")[0] or "unknown"
lines = lines[1:]
changes = [l for l in lines if l.strip()]
truncated = len(changes) > GIT_MAX_STATUS
out(True, branch=branch, changes=changes[:GIT_MAX_STATUS], truncated=truncated)
def act_git_diff(stat=False, path=None):
if path is not None:
path = _git_path(path)
audit("git-diff", f"{'stat' if stat else 'full'}:{path or 'all'}")
args = ["diff"]
if stat:
args.append("--stat")
if path:
args += ["--", path]
r = _run_git(args)
if r.returncode != 0:
fail("GIT_ERROR", "git diff failed", {"stderr": (r.stderr or "")[:500]})
text = r.stdout or ""
truncated = len(text) > GIT_MAX_DIFF
out(True, stat=bool(stat), path=path, diff=text[:GIT_MAX_DIFF],
truncated=truncated)
def act_git_log(limit=10, path=None):
if path is not None:
path = _git_path(path)
try:
limit = int(limit)
except (TypeError, ValueError):
fail("BAD_LIMIT", "limit must be an integer")
if not 1 <= limit <= 50:
fail("BAD_LIMIT", "limit must be 1..50")
audit("git-log", f"{path or 'all'}/{limit}")
args = ["log", "--oneline", "-n", str(limit)]
if path:
args += ["--", path]
r = _run_git(args)
if r.returncode != 0:
fail("GIT_ERROR", "git log failed", {"stderr": (r.stderr or "")[:500]})
commits = []
for line in (r.stdout or "").splitlines():
line = line.strip()
if not line:
continue
sha, _, subject = line.partition(" ")
commits.append({"sha": sha, "subject": subject})
out(True, commits=commits)
def act_tests_run(module=None, filter=None):
if module is not None:
if not TEST_MODULE_RE.match(module):
fail("BAD_NAME", "module must match ^tests\\.[a-z0-9_]+$")
if not (NETVM_ROOT / "tests" / (module.split(".", 1)[1] + ".py")).is_file():
fail("NOT_FOUND", f"unknown test module: {module}")
if filter is not None:
if not isinstance(filter, str) or not filter.strip() or len(filter) > 200:
fail("BAD_ARGS", "filter must be 1-200 chars")
if any(ord(c) < 32 or ord(c) == 127 for c in filter):
fail("BAD_ARGS", "filter contains control characters")
audit("tests-run", f"{module or 'all'}:{filter or '-'}")
if module:
argv = [sys.executable, "-m", "unittest", module]
else:
# No -t: tests/ has no __init__.py, so it is not importable as a
# package from top-level-dir '.'. Plain `discover -s tests` (the
# repo's own invocation) imports modules top-level and works.
argv = [sys.executable, "-m", "unittest", "discover", "-s", "tests"]
if filter:
argv += ["-k", filter]
# Hermetic import root: `python -m` prepends CWD to sys.path, but that
# is suppressed under PYTHONSAFEPATH/-P/isolated interpreters. Pin it.
env = dict(os.environ)
env["PYTHONPATH"] = str(NETVM_ROOT) + (
os.pathsep + env["PYTHONPATH"] if env.get("PYTHONPATH") else "")
try:
r = subprocess.run(argv, capture_output=True, text=True,
timeout=600, cwd=str(NETVM_ROOT), env=env)
except subprocess.TimeoutExpired:
fail("TESTS_ERROR", "test run timed out after 600s")
except OSError as e:
fail("TESTS_ERROR", f"cannot run tests: {e}")
combined = (r.stdout or "") + (r.stderr or "")
truncated = len(combined) > TESTS_MAX_OUTPUT
out(r.returncode == 0, returncode=r.returncode,
output=combined[-TESTS_MAX_OUTPUT:], truncated=truncated,
module=module or "all")
def act_ack(ref_id, to, sender, sidechat=None, allow_main_chat=False):
if not ref_id or not DM_ID_RE.match(ref_id):
fail("BAD_NAME", "id must be 6-64 hex chars",
{"field": "id", "value": ref_id})
if to not in VALID_AGENTS:
fail("BAD_NAME", f"to must be one of {sorted(VALID_AGENTS)}")
if sender not in VALID_AGENTS:
fail("BAD_NAME", f"sender must be one of {sorted(VALID_AGENTS)}")
# Sidechat-first policy (mirrors notify): default to the agent's
# sidechat, never main, unless explicitly overridden.
if allow_main_chat:
target = "main"
extra = ["--allow-main-chat"]
else:
target = sidechat if sidechat else NOTIFY_SIDECHATS.get(to)
if not target:
fail("BAD_NAME", f"no default sidechat for agent {to}; pass --sidechat <name>")
if sidechat and not _valid_sidechat_name(sidechat):
fail("BAD_NAME", "sidechat name must be 1-64 chars, no control characters")
extra = []
message = f"[ACK:{ref_id}]"
r = run([sys.executable, str(DM_PY), "send",
"--agent", sender, "--to", to, "--target", target,
*extra, message],
timeout=120)
if r.returncode != 0:
fail("DISPATCH_FAILED", f"dm.py send failed: {(r.stderr or r.stdout).strip()[-500:]}")
audit("ack", f"{sender}->{to}:{ref_id}")
out(True, to=to, target=target, ack=ref_id, sent=True)
def _policy_meta():
"""Parse Version:/Date: from CHAT_POLICY.md."""
version, date = "unknown", "unknown"
try:
for line in POLICY_FILE.read_text(encoding="utf-8", errors="replace").splitlines():
s = line.strip()
if s.startswith("**Version:**"):
version = s.split("**Version:**", 1)[1].strip().rstrip("\\")
elif s.startswith("**Date:**"):
date = s.split("**Date:**", 1)[1].strip().rstrip("\\")
except OSError:
pass
return version, date
def _watchdog_state():
"""Read the main-chat watchdog watermark (None if missing)."""
try:
return json.loads(WATCHDOG_STATE.read_text(encoding="utf-8"))
except Exception:
return None
def _policy_scan():
"""Classify dm-log.jsonl events for the sidechat-first policy.
Mirrors main-chat-watchdog.py classification:
VIOLATION : type=sent, target=main, no tags.allow_main_chat
AUTHORIZED: type=sent, target=main, tags.allow_main_chat set
BLOCKED : type=main_chat_blocked (gate worked)
Events before the allow_main_chat audit marker was adopted cannot be
classified, so the compliance window starts at marker adoption.
"""
per_agent = {}
legacy_untagged = 0
malformed = 0
adoption_ts = None
total_scanned = 0
def bucket(agent):
return per_agent.setdefault(agent, {
"blocked": 0,
"authorized_main": 0,
"violations": 0,
"total_sends": 0,
})
try:
with open(DM_LOG, "r", encoding="utf-8", errors="replace") as f:
lines = f.readlines()
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
for line in lines:
line = line.strip()
if not line:
continue
total_scanned += 1
try:
ev = json.loads(line)
except json.JSONDecodeError:
malformed += 1
continue
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":
legacy_untagged += 1
continue
etype = ev.get("type")
agent = ev.get("agent") or "unknown"
if etype == "main_chat_blocked":
bucket(agent)["blocked"] += 1
elif etype == "sent":
bucket(agent)["total_sends"] += 1
if ev.get("target") == "main":
tags = ev.get("tags") or {}
if tags.get("allow_main_chat"):
bucket(agent)["authorized_main"] += 1
else:
bucket(agent)["violations"] += 1
return per_agent, {
"window_start": adoption_ts,
"events_scanned": total_scanned,
"malformed": malformed,
"legacy_untagged_main_sends": legacy_untagged,
}
def act_policy():
audit("policy")
version, date = _policy_meta()
per_agent, meta = _policy_scan()
if per_agent is None:
fail("POLICY_ERROR", "failed to read dm log", meta)
totals = {"blocked": 0, "authorized_main": 0, "violations": 0, "total_sends": 0}
for counts in per_agent.values():
for k in totals:
totals[k] += counts[k]
out(True,
version=version,
date=date,
rule=("No DM naturally lands in main chat. Main requires explicit "
"opt-in (--allow-main-chat / allow_main_chat)."),
compliance_window=meta,
watchdog_state=_watchdog_state(),
agents=per_agent,
totals=totals,
status="clean" if totals["violations"] == 0 else "violations found")
def act_policy_check(agent):
if agent not in VALID_AGENTS:
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
audit("policy-check", agent)
per_agent, meta = _policy_scan()
if per_agent is None:
fail("POLICY_ERROR", "failed to read dm log", meta)
counts = per_agent.get(agent, {
"blocked": 0, "authorized_main": 0, "violations": 0, "total_sends": 0})
out(True, agent=agent,
compliance_window=meta, **counts,
status="clean" if counts["violations"] == 0 else "violations found")
def act_policy_show():
audit("policy-show")
try:
text = POLICY_FILE.read_text(encoding="utf-8", errors="replace")
except OSError as e:
fail("POLICY_ERROR", f"cannot read {POLICY_FILE}: {e}")
version, date = _policy_meta()
out(True, version=version, date=date, path=str(POLICY_FILE), policy=text)
def act_vars_list():
try:
from variables import Variables
v = Variables()
audit("vars-list")
out(True, variables=v.all(), schemas={n: v.schema(n) for n in v.names()})
except Exception as e:
fail("VARS_ERROR", str(e))
def act_vars_get(name):
try:
from variables import Variables
v = Variables()
audit("vars-get", name)
out(True, name=name, value=v.get(name), schema=v.schema(name))
except Exception as e:
fail("VARS_ERROR", str(e))
def act_vars_set(name, val_str):
try:
from variables import Variables
v = Variables()
spec = v.schema(name)
vtype = spec.get("type", "str")
if vtype == "int":
val = int(val_str)
elif vtype == "float":
val = float(val_str)
elif vtype == "bool":
val = val_str.lower() in ("true", "1", "yes")
else:
val = val_str
caller = os.environ.get("BOX_CALLER", "box-ctl")
new_val = v.set(name, val, by=caller)
audit("vars-set", f"{name}={new_val}")
out(True, name=name, value=new_val)
except Exception as e:
fail("VARS_ERROR", str(e))
def act_vars_reset(name):
try:
from variables import Variables
v = Variables()
caller = os.environ.get("BOX_CALLER", "box-ctl")
new_val = v.reset(name, by=caller)
audit("vars-reset", name)
out(True, name=name, value=new_val)
except Exception as e:
fail("VARS_ERROR", str(e))
def act_vars_history(name=None, limit=20):
try:
from variables import Variables
v = Variables()
entries = v.history(name=name, limit=limit)
audit("vars-history", name or "*")
out(True, history=entries, count=len(entries))
except Exception as e:
fail("VARS_ERROR", str(e))
def act_vars_rollback(name, revision=None):
try:
from variables import Variables
v = Variables()
caller = os.environ.get("BOX_CALLER", "box-ctl")
new_val = v.rollback(name, revision=revision, by=caller)
audit("vars-rollback", f"{name}={new_val}")
out(True, name=name, value=new_val)
except Exception as e:
fail("VARS_ERROR", str(e))
def act_strat_list():
try:
from modulate import get_all_strategies
audit("strat-list")
out(True, strategies=get_all_strategies())
except Exception as e:
fail("STRAT_ERROR", str(e))
def act_strat_get(itype, subtype=None, agent=None):
try:
from modulate import get_strategy_row, InputType, Priority
it = InputType(itype.lower())
row = get_strategy_row(it, subtype.upper() if subtype else None, agent=agent)
audit("strat-get", f"{itype}:{subtype or '*'}:{agent or '*'}")
out(True, type=it.value, subtype=subtype, agent=agent, track=row[0],
priority=row[1].value if isinstance(row[1], Priority) else str(row[1]),
timeout_s=row[2], nudges=row[3], escalate=row[4])
except Exception as e:
fail("STRAT_ERROR", str(e))
def act_strat_set(itype, payload_str=None):
try:
from modulate import set_strategy_override
if not payload_str:
payload_str = sys.stdin.read()
data = json.loads(payload_str)
subtype = data.get("subtype")
agent = data.get("agent")
caller = os.environ.get("BOX_CALLER", "box-ctl")
res = set_strategy_override(
itype, subtype, agent,
track=data.get("track"),
priority=data.get("priority"),
timeout_s=data.get("timeout_s") or data.get("timeout"),
nudges=data.get("nudges"),
escalate=data.get("escalate"),
by=caller
)
audit("strat-set", f"{itype}:{subtype or '*'}:{agent or '*'}")
out(True, type=itype, subtype=subtype, agent=agent, strategy=res)
except Exception as e:
fail("STRAT_ERROR", str(e))
def act_strat_reset(itype, subtype=None, agent=None):
try:
from modulate import reset_strategy_override
caller = os.environ.get("BOX_CALLER", "box-ctl")
ok = reset_strategy_override(itype, subtype, agent, by=caller)
audit("strat-reset", f"{itype}:{subtype or '*'}:{agent or '*'}")
out(True, type=itype, subtype=subtype, agent=agent, reset=ok)
except Exception as e:
fail("STRAT_ERROR", str(e))
def act_loop_status(agent=None, limit=20, status=None):
try:
from gravity import reconstruct_loops
loops = reconstruct_loops(limit=limit, agent=agent, status_filter=status)
audit("loop-status", f"agent={agent or '*'}")
out(True, loops=loops, count=len(loops))
except Exception as e:
fail("LOOP_ERROR", str(e))
def act_loop_health(threshold=None):
try:
from gravity import get_fleet_loop_health
t_val = float(threshold) if threshold is not None else None
h = get_fleet_loop_health(threshold=t_val)
audit("loop-health")
out(True, **h)
except Exception as e:
fail("LOOP_ERROR", str(e))
def act_loop_breaks():
try:
from gravity import diagnose_breaks
breaks = diagnose_breaks()
audit("loop-breaks")
out(True, breaks=breaks, count=len(breaks))
except Exception as e:
fail("LOOP_ERROR", str(e))
# ---------------------------------------------------------------------------
# SSH key minting & management
# ---------------------------------------------------------------------------
KEY_DIR = Path("/home/super/.ssh")
ALLOWED_SIGNERS_PATHS = [
NETVM_ROOT / "dm-signers" / "allowed_signers",
Path("/home/super/.exec-signers")
]
def act_ssh_mint(name, force=False):
"""Mint a new ed25519 SSH keypair and register it in allowed_signers."""
key_name = f"id_{name}"
priv_path = KEY_DIR / key_name
pub_path = KEY_DIR / f"{key_name}.pub"
if priv_path.exists() and not force:
fail("KEY_EXISTS", f"SSH key {priv_path} already exists -- pass --force to overwrite",
{"key_path": str(priv_path)})
if priv_path.exists() and force:
try:
priv_path.unlink()
if pub_path.exists():
pub_path.unlink()
except Exception as e:
fail("IO_ERROR", f"Failed removing old key: {e}")
cmd = ["ssh-keygen", "-t", "ed25519", "-N", "", "-C", f"{name}@netvm", "-f", str(priv_path)]
r = subprocess.run(cmd, capture_output=True, text=True)
if r.returncode != 0:
fail("KEYGEN_FAILED", f"ssh-keygen failed: {r.stderr.strip()}")
pub_line = pub_path.read_text(encoding="utf-8").strip()
pub_key_body = pub_line.split()[1] if len(pub_line.split()) > 1 else pub_line
# Save to dm-signers/{name}.pub
ds_pub = NETVM_ROOT / "dm-signers" / f"{name}.pub"
try:
ds_pub.write_text(pub_line + "\n", encoding="utf-8")
except Exception:
pass
# Register into signers files
registered_files = []
for sp in ALLOWED_SIGNERS_PATHS:
try:
lines = sp.read_text(encoding="utf-8").splitlines() if sp.exists() else []
new_lines = []
found = False
for line in lines:
if not line.strip():
continue
parts = line.split()
if parts[0] in (name, f"operator-{name}"):
new_lines.append(f"{parts[0]} ssh-ed25519 {pub_key_body} {name}@netvm")
found = True
else:
new_lines.append(line)
if not found:
new_lines.append(f"{name} ssh-ed25519 {pub_key_body} {name}@netvm")
new_lines.append(f"operator-{name} ssh-ed25519 {pub_key_body} {name}@netvm")
sp.write_text("\n".join(new_lines) + "\n", encoding="utf-8")
registered_files.append(str(sp))
except Exception as e:
sys.stderr.write(f"Warning: failed updating {sp}: {e}\n")
audit("ssh-mint", f"name={name} priv={priv_path}")
out(True, name=name, private_key=str(priv_path), public_key=str(pub_path),
public_key_content=pub_line, registered_in=registered_files)
def act_ssh_list():
"""List managed SSH keypairs and registrations."""
keys = []
if KEY_DIR.exists():
for p in sorted(KEY_DIR.glob("id_*")):
if not p.name.endswith(".pub"):
pub_p = p.with_name(p.name + ".pub")
pub_content = pub_p.read_text(encoding="utf-8").strip() if pub_p.exists() else None
keys.append({
"name": p.name[3:],
"private_key": str(p),
"has_public": pub_p.exists(),
"public_key": pub_content,
})
out(True, keys=keys, count=len(keys))
def act_ssh_show(name):
"""Show public key details and fingerprints for an identity."""
priv_path = KEY_DIR / f"id_{name}"
pub_path = KEY_DIR / f"id_{name}.pub"
if not pub_path.exists():
fail("NOT_FOUND", f"No public key found for {name} at {pub_path}")
pub_line = pub_path.read_text(encoding="utf-8").strip()
cmd = ["ssh-keygen", "-lf", str(pub_path)]
r = subprocess.run(cmd, capture_output=True, text=True)
fp = r.stdout.strip() if r.returncode == 0 else ""
out(True, name=name, private_key=str(priv_path), public_key=str(pub_path),
public_key_content=pub_line, fingerprint=fp)
def act_ssh_ports():
"""List container reverse tunnel port mappings."""
import agent_md
out(True, **agent_md.get_ssh_info())
def act_ssh_info(name):
"""Show SSH dial-in command and reverse tunnel coordinates for a specific agent."""
import agent_md
out(True, **agent_md.get_ssh_info(name))
SSH_CHECK_TIMEOUT = 25
def build_ssh_sweep_script():
"""Return the python3 probe script executed on the jump host.
Reads `kind:port` args (kind is `ssh` or `term`), connects to each
127.0.0.1:port, grabs the SSH banner or terminal HTTP status line,
and prints one JSON object: {"ports": {port: {...}}}.
"""
return (
"import json, socket, sys, time\n"
"out = {}\n"
"for arg in sys.argv[1:]:\n"
" try:\n"
" kind, port_s = arg.split(':', 1)\n"
" port = int(port_s)\n"
" except ValueError:\n"
" continue\n"
" rec = {'open': False, 'latency_ms': None, 'banner': '', 'http': ''}\n"
" t0 = time.time()\n"
" try:\n"
" s = socket.create_connection(('127.0.0.1', port), timeout=2.0)\n"
" except Exception:\n"
" out[str(port)] = rec\n"
" continue\n"
" rec['open'] = True\n"
" rec['latency_ms'] = int((time.time() - t0) * 1000)\n"
" try:\n"
" s.settimeout(2.0)\n"
" if kind == 'term':\n"
" s.sendall(b'GET / HTTP/1.0\\r\\n\\r\\n')\n"
" data = s.recv(128)\n"
" rec['http'] = data.split(b'\\n')[0].decode('utf-8', 'replace').strip()[:64]\n"
" else:\n"
" data = s.recv(128)\n"
" rec['banner'] = data.split(b'\\n')[0].decode('utf-8', 'replace').strip()[:64]\n"
" except Exception:\n"
" pass\n"
" try:\n"
" s.close()\n"
" except Exception:\n"
" pass\n"
" out[str(port)] = rec\n"
"print(json.dumps({'ports': out}))\n"
)
def parse_ssh_sweep_result(sweep, tunnels):
"""Map sweep {port: probe} output onto {account: health} via TUNNEL_PORTS.
Unknown sweep ports are ignored; accounts with no sweep data keep
None (unknown) health. Always returns every account in tunnels.
"""
ports = sweep.get("ports", {}) if isinstance(sweep, dict) else {}
by_port = {}
for acct, info in tunnels.items():
by_port[str(info.get("port"))] = (acct, "ssh")
by_port[str(info.get("terminal"))] = (acct, "term")
accounts = {}
for acct, info in tunnels.items():
accounts[acct] = {
"ssh_port": info.get("port"),
"ssh_up": None,
"ssh_latency_ms": None,
"ssh_banner": "",
"term_port": info.get("terminal"),
"term_up": None,
"term_latency_ms": None,
"term_http": "",
"container_user": info.get("user", "hatch"),
}
if not isinstance(ports, dict):
return accounts
for port_s, probe in ports.items():
slot = by_port.get(str(port_s))
if not slot or not isinstance(probe, dict):
continue
acct, kind = slot
ent = accounts.get(acct)
if ent is None:
continue
is_open = bool(probe.get("open"))
lat = probe.get("latency_ms")
try:
lat = int(lat) if lat is not None else None
except (TypeError, ValueError):
lat = None
if kind == "ssh":
ent["ssh_up"] = is_open
ent["ssh_latency_ms"] = lat if is_open else None
ent["ssh_banner"] = str(probe.get("banner") or "")[:64]
else:
ent["term_up"] = is_open
ent["term_latency_ms"] = lat if is_open else None
ent["term_http"] = str(probe.get("http") or "")[:64]
return accounts
def run_vm_port_sweep(jump_host, operator_user, tunnels, timeout=SSH_CHECK_TIMEOUT):
"""Run one SSH to the jump host sweeping all tunnel ports.
Returns (sweep_dict_or_None, error_str). sweep is the parsed
{"ports": {...}} payload; error is "" on success.
"""
sweep_args = []
for _acct, info in tunnels.items():
sweep_args.append(f"ssh:{info.get('port')}")
sweep_args.append(f"term:{info.get('terminal')}")
cmd = [
"ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=8",
"-o", "StrictHostKeyChecking=no",
]
identity = os.environ.get("SSH_IDENTITY_FILE", "")
if identity:
cmd += ["-o", "IdentitiesOnly=yes", "-i", identity]
cmd += [
f"{operator_user}@{jump_host}",
"python3", "-",
] + sweep_args
try:
r = subprocess.run(cmd, input=build_ssh_sweep_script(),
capture_output=True, text=True, timeout=timeout)
except subprocess.TimeoutExpired:
return None, f"jump host {jump_host} sweep timed out after {timeout}s"
except FileNotFoundError:
return None, "local ssh binary not found"
except Exception as e:
return None, f"ssh to {jump_host} failed: {e}"
if r.returncode != 0:
err = (r.stderr or "").strip().splitlines()
hint = err[-1][:160] if err else f"exit {r.returncode}"
return None, f"jump host {jump_host} unreachable: {hint}"
try:
sweep = json.loads(r.stdout)
except Exception:
return None, f"jump host {jump_host} returned unparseable sweep output"
if not isinstance(sweep, dict) or "ports" not in sweep:
return None, f"jump host {jump_host} returned malformed sweep output"
return sweep, ""
def act_ssh_check():
"""Probe all container reverse-tunnel ports from the jump host.
Single SSH connection, VM-side sweep. Always emits HTTP-200-style
ok:true with per-account health; jump failures surface as
jump_reachable:false (degraded-state data, not a fatal error).
"""
import time as _time
import agent_md
audit("ssh-check")
t0 = _time.time()
tunnels = dict(agent_md.TUNNEL_PORTS)
jump_host = os.environ.get("SSH_JUMP_HOST", "34.139.37.135")
operator_user = os.environ.get("OPERATOR_USER", "super")
sweep, err = run_vm_port_sweep(jump_host, operator_user, tunnels)
wall_ms = int((_time.time() - t0) * 1000)
accounts = parse_ssh_sweep_result(sweep or {}, tunnels)
if err:
out(True, jump_host=jump_host, operator_user=operator_user,
jump_reachable=False, error=err, accounts=accounts,
checked_at=utcnow(), latency_ms=wall_ms)
else:
out(True, jump_host=jump_host, operator_user=operator_user,
jump_reachable=True, accounts=accounts,
checked_at=utcnow(), latency_ms=wall_ms)
MD_MAX_READ = 64 * 1024
MD_MAX_DIFF = 64 * 1024
MD_MAX_LIST = 200
def act_md_audit(accounts=None):
"""Audit markdown files and operational drive across fleet agents."""
import agent_md
audit("md-audit", ",".join(accounts or []))
try:
res = agent_md.audit_agents(accounts=accounts)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
out(True, agents=res)
def act_md_list(account, path=""):
"""List container files via Hatch."""
import agent_md
audit("md-list", f"{account}:{path}")
try:
entries = agent_md.list_files(account, path=path)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
except FileNotFoundError as e:
fail("NOT_FOUND", str(e))
except Exception as e:
fail("GATEWAY_ERROR", f"md-list failed: {e}")
truncated = len(entries) > MD_MAX_LIST
out(True, account=account, path=path, entries=entries[:MD_MAX_LIST],
truncated=truncated)
def act_md_read(account, filename):
"""Read a markdown file from the agent container via Hatch."""
import agent_md
audit("md-read", f"{account}:{filename}")
try:
res = agent_md.read_md(account, filename)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
except FileNotFoundError as e:
fail("NOT_FOUND", str(e))
except Exception as e:
fail("GATEWAY_ERROR", f"md-read failed: {e}")
ok = res.pop("ok", True)
text = res.get("text", "")
truncated = len(text) > MD_MAX_READ
if truncated:
res["text"] = text[:MD_MAX_READ]
out(ok, truncated=truncated, **res)
def act_md_write(account, filename, content, overwrite=True, append=False):
"""Write content to an agent's container file via Hatch."""
import agent_md
audit("md-write", f"{account}:{filename}")
try:
res = agent_md.write_md(account, filename, content, overwrite=overwrite, append=append)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
except FileNotFoundError as e:
fail("NOT_FOUND", str(e))
except Exception as e:
fail("GATEWAY_ERROR", f"md-write failed: {e}")
out(res.pop("ok", True), **res)
def act_md_diff(account, filename):
"""Diff remote container file against local shared/operators template."""
import agent_md
audit("md-diff", f"{account}:{filename}")
try:
res = agent_md.diff_md(account, filename)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
except FileNotFoundError as e:
fail("NOT_FOUND", str(e))
except Exception as e:
fail("GATEWAY_ERROR", f"md-diff failed: {e}")
ok = res.pop("ok", True)
diff = res.get("diff", "")
truncated = len(diff) > MD_MAX_DIFF
if truncated:
res["diff"] = diff[:MD_MAX_DIFF]
out(ok, truncated=truncated, **res)
def _md_amend_parse(args, usage):
"""Parse amend argv: filename + content/--content/--file/--stdin + flags."""
if not args:
fail("BAD_ARGS", usage)
filename = args[0]
content = None
author = "operator"
reason = ""
idx = 1
while idx < len(args):
if args[idx] == "--file" and idx + 1 < len(args):
content = Path(args[idx + 1]).read_text(encoding="utf-8")
idx += 2
elif args[idx] == "--content" and idx + 1 < len(args):
content = args[idx + 1]
idx += 2
elif args[idx] == "--stdin":
content = sys.stdin.read()
idx += 1
elif args[idx] == "--author" and idx + 1 < len(args):
author = args[idx + 1]
idx += 2
elif args[idx] == "--reason" and idx + 1 < len(args):
reason = args[idx + 1]
idx += 2
elif content is None and not args[idx].startswith("--"):
content = args[idx]
idx += 1
else:
idx += 1
if content is None:
fail("BAD_ARGS", usage)
return filename, content, author, reason
def _md_append_parse(args, usage):
"""Parse append argv: filename + text/--content/--file/--stdin + flags."""
if not args:
fail("BAD_ARGS", usage)
filename = args[0]
text = None
author = "operator"
section = None
idx = 1
while idx < len(args):
if args[idx] == "--file" and idx + 1 < len(args):
text = Path(args[idx + 1]).read_text(encoding="utf-8")
idx += 2
elif args[idx] == "--content" and idx + 1 < len(args):
text = args[idx + 1]
idx += 2
elif args[idx] == "--stdin":
text = sys.stdin.read()
idx += 1
elif args[idx] == "--author" and idx + 1 < len(args):
author = args[idx + 1]
idx += 2
elif args[idx] == "--section" and idx + 1 < len(args):
section = args[idx + 1]
idx += 2
elif text is None and not args[idx].startswith("--"):
text = args[idx]
idx += 1
else:
idx += 1
if text is None:
fail("BAD_ARGS", usage)
return filename, text, author, section
def act_md_amend(filename, content, author="operator", reason=""):
"""Amend centralized shared template with safety checks and git commit."""
import agent_md
audit("md-amend", f"{author}:{filename}")
try:
res = agent_md.amend_md(filename, content, author=author, reason=reason)
out(res.pop("ok", True), **res)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
except Exception as e:
fail("AMEND_FAILED", str(e))
def act_md_append(filename, text, author="operator", section=None):
"""Safely append an amendment or lesson to a centralized shared template."""
import agent_md
audit("md-append", f"{author}:{filename}")
try:
res = agent_md.append_md(filename, text, author=author, section=section)
out(res.pop("ok", True), **res)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
except Exception as e:
fail("APPEND_FAILED", str(e))
def act_md_pull(account, filename):
"""Pull canonical shared template into an agent container."""
import agent_md
audit("md-pull", f"{account}:{filename}")
try:
res = agent_md.pull_md(account, filename)
out(res.pop("ok", True), **res)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
except Exception as e:
fail("PULL_FAILED", str(e))
def act_md_inject_drive(account, force=False):
"""Inject high-drive operational templates into an agent container via Hatch."""
import agent_md
audit("md-inject-drive", account)
try:
res = agent_md.inject_drive(account, force=force)
except agent_md.MDValidationError as e:
fail("BAD_NAME", str(e))
except FileNotFoundError as e:
fail("NOT_FOUND", str(e))
except Exception as e:
fail("GATEWAY_ERROR", f"md-inject-drive failed: {e}")
out(res.pop("ok", True), **res)
def act_md_sync_all(force=False):
"""Inject high-drive operational templates across all active agents."""
import agent_md
audit("md-sync-all")
results = {}
for acct in agent_md.VALID_ACCOUNTS:
try:
results[acct] = agent_md.inject_drive(acct, force=force)
except Exception as e:
results[acct] = {"ok": False, "error": str(e)}
out(True, results=results)
def act_loop_resolve(dm_id, note=None):
f_path = NETVM_ROOT / "followups.json"
resolved = False
if f_path.exists():
try:
with open(f_path, "r") as f:
data = json.load(f)
if dm_id in data:
data[dm_id]["status"] = "resolved"
data[dm_id]["resolved_at"] = datetime.now(timezone.utc).isoformat()
if note:
data[dm_id]["note"] = note
tmp = f"{f_path}.tmp.{os.getpid()}"
with open(tmp, "w") as f:
json.dump(data, f, indent=2)
os.replace(tmp, f_path)
resolved = True
except Exception as e:
fail("LOOP_ERROR", f"Failed updating followups.json: {e}")
# Append to job-log.jsonl
try:
with open(JOB_LOG, "a") as f:
f.write(json.dumps({
"ts": datetime.now(timezone.utc).isoformat(),
"type": "loop_resolved",
"dm_id": dm_id,
"note": note or "manually resolved via box-ctl",
"caller": os.environ.get("BOX_CALLER", "box-ctl"),
}) + "\n")
except Exception:
pass
audit("loop-resolve", dm_id)
out(True, dm_id=dm_id, resolved=resolved)
def act_loop_remediate(dry_run=False):
try:
from gravity import remediate_breaks
res = remediate_breaks(dry_run=dry_run)
audit("loop-remediate", f"dry_run={dry_run}")
res.pop("ok", None)
out(True, **res)
except Exception as e:
fail("LOOP_ERROR", str(e))
# ---------------------------------------------------------------------------
# CLI
# ---------------------------------------------------------------------------
USAGE = """usage: box-ctl.py <action> [args]
fleet:
fleet-status
chrome-errors per-profile chrome FATAL/crash counts (new since watermark)
watchdog-alerts
relay-health
cdp-latency
identity-audit VM identity audit drift check
main-loop check|status|enable|disable [--agent <name>]
main-chat self-monitor; enable/disable per agent or all
timer actions:
timer-list
timer-status <name>
timer-create <name>
timer-delete <name> [--keep-job]
timer-start|timer-stop|timer-enable|timer-disable <name>
job actions:
job-list
job-get <name>
job-put <name> (job JSON on stdin)
job-delete <name> [--force]
job-trigger <name>
job-result <job-id> (result JSON on stdin; records RESULT + advances chain)
job-status <name-or-id> (dispatch/result/chain state)
job-next <job-id> [--success|--fail] (dry-run: what would box dispatch next)
job-chain <from> <to> [--on-failure] (wire chain_next between jobs)
variable actions:
vars-list
vars-get <name>
vars-set <name> <value>
vars-reset <name>
vars-history [name] [limit]
vars-rollback <name> [revision]
strategy actions:
strat-list
strat-get <type> [subtype] [--agent AGENT]
strat-set <type> (JSON on stdin or as arg)
strat-reset <type> [subtype] [--agent AGENT]
loop actions:
loop-status [--agent AGENT] [--limit N] [--status STATUS]
loop-health [--threshold T]
loop-breaks
loop-resolve <dm_id> [note]
loop-remediate [--dry-run]
thread actions (muse-cli gateway; hyphenated aliases for `thread <sub>`):
thread-list --agent <agent> [--limit N]
thread-pin --agent <agent> --thread <id>
thread-unpin --agent <agent> --thread <id>
thread-archive --agent <agent> --thread <id> --confirm
thread-unarchive --agent <agent> --thread <id>
thread-rename --agent <agent> --thread <id> --title <title>
notify:
notify <agent> <message> [--sidechat <name>] [--allow-main-chat]
ack <id> --to <agent> --sender <agent> [--sidechat <name>] [--allow-main-chat]
acknowledge a DM or work order (sidechat-first)
Send a DM to an agent's default sidechat (never main chat unless
--allow-main-chat is passed explicitly).
policy:
policy chat policy version, rule, and compliance summary
policy check <agent> per-agent compliance detail
policy show print the full CHAT_POLICY.md
thread actions (muse-cli gateway):
thread list <agent> [--limit N]
thread pin <agent> <thread-id>
thread unpin <agent> <thread-id>
thread archive <agent> <thread-id> --confirm
thread unarchive <agent> <thread-id>
swarm actions (agent swarms; bl owns state, operator agents do the spawning):
swarm-spawn <count> <task...> [--label L]
swarm-list
swarm-status <swarm-id>
swarm-attach <swarm-id> <slot> <agent-id>
swarm-report <swarm-id> <slot> (result JSON on stdin)
swarm-results <swarm-id>
swarm-kill <swarm-id> --confirm
swarm-prune [--stale-hours H] [--confirm]
archive stale terminal swarms (previews count without --confirm)
(space-separated aliases: swarm spawn|list|status|attach|report|results|kill|prune)
quality:
quality-check run box-ctl self-diagnostics (validators, output
contract, audit path, atomic writes)
quality-validate <action> [args...]
dry-run: validate args without executing
(space-separated alias: quality check|validate)
dm-log:
dm-log [limit] [--agent <agent>]
recent DM send log (agent-scoped when given)
unread [--agent <agent>] fleet unread/activity counts (read-only)
dev (no-SSH agent development):
git-status git status --short --branch (read-only)
git-diff [--stat] [--path <p>]
git diff, capped at 64KB (read-only)
git-log [--limit N] [--path <p>]
recent commits as sha/subject (read-only)
tests-run [tests.<module>] [--filter <pattern>]
run repo unit tests (full suite when omitted;
--filter is unittest -k)
md actions (agent .md files via Hatch; space-separated aliases: md audit|list|read|write|diff|amend|append|pull|inject-drive|sync-all):
md-audit [accounts...] fleet drive audit (read-only)
md-list <account> [path] list container files, capped at 200 (read-only)
md-read <account> <filename>
read container file, capped at 64KB (read-only)
md-diff <account> <filename>
diff against shared template, capped at 64KB (read-only)
md-pull <account> <filename>
pull canonical template into container
md-inject-drive <account> [--force]
inject high-drive templates into one agent
md-sync-all [--force] inject high-drive templates across all agents
md-amend <filename> (<content>|--content t|--file p|--stdin)
rewrite shared template (validated, git-committed)
[--author name] [--reason why]
md-append <filename> (<text>|--content t|--file p|--stdin)
append to shared template (git-committed)
[--author name] [--section header]
md-write <account> <filename> <content>
raw container write (SSH-only, no HTTPS op)
approval actions (fleet browser approvals via CDP):
approval-check [node] pending approvals across fleet or one node (read-only)
approval-list [node] (alias for approval-check)
approval-allow <node> [--always] [--force] [--message TEXT] [--allow-main-chat]
allow the node's active approval
approval-approve <node> [...] (alias for approval-allow)
approval-deny <node> [--message TEXT] [--allow-main-chat]
deny the node's active approval
approval-auto [node] auto-allow TRUSTED (non-key) prompts
tmux & worker actions (multi-socket tally & regex auto-approvals):
tmux-tally tally all tmux sessions, workers, and panes across fleet sockets (read-only)
tmux-auto-status runtime status of tmux auto-approvals & rule inventory (read-only)
tmux-auto-toggle [--enable|--disable] [--node NODE] [--session SESSION]
toggle tmux auto-approvals globally, per-agent, or per-session
tmux-auto-once [--dry-run] one-shot scan & auto-approve terminal prompts
(space-separated alias: tmux tally|auto|status|toggle|once)
onboard actions:
onboard-connects consolidated active fleet and client onboard connects (read-only)
(space-separated alias: onboard connects)
ssh actions (space-separated alias: ssh mint|list|show|ports|info|check):
ssh-mint <name> [--force] mint a new SSH key
ssh-list list minted keys
ssh-show <name> show key detail
ssh-ports tunnel port inventory
ssh-info [name] connection coordinates
ssh-check live tunnel health sweep via jump host"""
def act_thread(op, agent, thread=None, title=None, limit=None, confirm=False):
"""Thread bookkeeping via bin/muse-threads.py (hybrid gateway).
ops: list | pin | unpin | archive | unarchive | rename
archive requires confirm=True (the --confirm flag; BatchMode SSH
cannot prompt, so the flag is the explicit confirmation).
"""
if op not in ("list", "pin", "unpin", "archive", "unarchive", "rename"):
fail("BAD_NAME", "usage: thread list|pin|unpin|archive|unarchive|rename <agent> [args]")
if agent not in VALID_AGENTS:
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
if op == "list":
audit("thread-list-" + agent)
else:
if not thread or not THREAD_RE.match(thread):
fail("BAD_NAME", "thread id must match ^[A-Za-z0-9-]{1,64}$",
{"field": "thread", "value": thread})
if op == "archive" and not confirm:
fail("CONFIRM_REQUIRED",
"archive is a visible side effect -- pass --confirm to proceed",
{"would_archive": {"agent": agent, "thread": thread}})
if op == "rename" and not title:
fail("BAD_NAME", "rename requires a title")
audit("thread-" + op, f"{agent}/{thread}")
helper = os.path.join(os.path.dirname(os.path.abspath(__file__)), "muse-threads.py")
cmd = [sys.executable, helper, op, "--agent", agent]
if thread:
cmd += ["--thread", thread]
if title:
cmd += ["--title", title]
try:
p = subprocess.run(cmd, capture_output=True, text=True, timeout=120)
except subprocess.TimeoutExpired:
fail("GATEWAY_UNAVAILABLE", f"muse-threads.py timed out (op={op})")
except OSError as e:
fail("GATEWAY_UNAVAILABLE", f"cannot run muse-threads.py: {e}")
# Helper prints one JSON object; relay it, applying --limit for list
if p.returncode == 0:
raw = p.stdout.strip()
if limit is not None and op == "list" and raw:
try:
data = json.loads(raw)
if isinstance(data.get("threads"), list):
data["threads"] = data["threads"][:limit]
print(json.dumps(data))
return
except ValueError:
pass
print(raw if raw else '{"ok": true}')
return
# Failure: try to extract the helper's JSON error, else wrap
try:
err = json.loads(p.stdout.strip() or "{}")
fail(err.get("code", "GATEWAY_ERROR"), err.get("error", "thread op failed"),
err.get("detail"))
except (ValueError, AttributeError):
fail("GATEWAY_ERROR", f"thread {op} failed",
(p.stderr.strip() or p.stdout.strip())[:300])
# ---------------------------------------------------------------------------
# Agent swarms — parallel subagent groups tracked on bl.
#
# box-ctl.py owns the swarm STATE (bl is the source of truth). Actual agent
# spawning is done by an operator agent in the container via subagent.spawn;
# the spawner attaches agent IDs with swarm-attach and records outcomes with
# swarm-report. A `swarm-spawner` cron agent can poll `swarm-list` for pending
# swarms and drive them without human involvement.
#
# Lifecycle: pending -> running -> completed|partial|killed
# Slot lifecycle: pending -> running -> done|failed|killed
# ---------------------------------------------------------------------------
SWARM_FILE = NETVM_ROOT / "swarms.json"
SWARM_ID_RE = re.compile(r"^[a-z0-9-]{1,64}$")
SWARM_AGENT_RE = re.compile(r"^[A-Za-z0-9-]{1,64}$")
SWARM_OPS = ("spawn", "list", "status", "attach", "report", "results", "kill", "prune")
SWARM_MAX_COUNT = 50
SWARM_MAX_TASK_LEN = 2000
def _swarm_load():
try:
with open(SWARM_FILE) as f:
data = json.load(f)
return data if isinstance(data, dict) else {}
except (OSError, ValueError):
return {}
def _swarm_save(swarms):
tmp = str(SWARM_FILE) + ".tmp"
with open(tmp, "w") as f:
json.dump(swarms, f, indent=2)
os.replace(tmp, SWARM_FILE)
def _swarm_new_id():
return ("sw-" + datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S")
+ "-" + os.urandom(2).hex())
def _swarm_get(swarms, sid):
if not sid or not SWARM_ID_RE.match(sid):
fail("BAD_NAME", "swarm id must match ^[a-z0-9-]{1,64}$",
{"field": "swarm_id", "value": sid})
swarm = swarms.get(sid)
if swarm is None:
fail("SWARM_NOT_FOUND", "no such swarm: %s" % sid,
{"swarm_id": sid})
return swarm
def _swarm_counts(swarm):
counts = {"pending": 0, "running": 0, "done": 0, "failed": 0,
"killed": 0}
for slot in swarm.get("slots", []):
st = slot.get("status")
if st in counts:
counts[st] += 1
return counts
def _swarm_rollup(swarm):
"""Recompute swarm-level status from slot statuses."""
counts = _swarm_counts(swarm)
total = len(swarm.get("slots", []))
if total and counts["killed"] == total:
return "killed"
if total and counts["pending"] + counts["running"] == 0:
return "completed" if counts["failed"] == 0 else "partial"
if counts["running"] or counts["done"] or counts["failed"]:
return "running"
return "pending"
def _swarm_slot(swarm, slot_str):
try:
slot = int(slot_str)
except (TypeError, ValueError):
fail("BAD_SLOT", "slot must be an integer", {"value": slot_str})
slots = swarm.get("slots", [])
if not 0 <= slot < len(slots):
fail("BAD_SLOT", "slot must be 0-%d" % (len(slots) - 1),
{"value": slot_str})
return slots[slot]
def act_swarm_spawn(count_str, task, label=None):
try:
count = int(count_str)
except (TypeError, ValueError):
fail("BAD_COUNT", "count must be an integer 1-50",
{"value": count_str})
if not 1 <= count <= SWARM_MAX_COUNT:
fail("BAD_COUNT", "count must be 1-%d" % SWARM_MAX_COUNT,
{"value": count})
if not task or not task.strip():
fail("BAD_TASK", "task must be a non-empty string")
if len(task) > SWARM_MAX_TASK_LEN:
fail("BAD_TASK", "task must be <= %d chars" % SWARM_MAX_TASK_LEN,
{"length": len(task)})
if label is not None and not SWARM_ID_RE.match(label):
fail("BAD_NAME", "label must match ^[a-z0-9-]{1,64}$",
{"field": "label", "value": label})
swarms = _swarm_load()
sid = _swarm_new_id()
while sid in swarms: # vanishingly unlikely; be safe
sid = _swarm_new_id()
now = utcnow()
swarm = {
"swarm_id": sid,
"task": task,
"count": count,
"label": label,
"status": "pending",
"created_ts": now,
"created_by": os.environ.get("BOX_CALLER", "unknown"),
"updated_ts": now,
"slots": [
{"slot": i, "agent_id": None, "status": "pending",
"result": None, "updated_ts": now}
for i in range(count)
],
}
swarms[sid] = swarm
_swarm_save(swarms)
audit("swarm-spawn", sid)
out(True, swarm_id=sid, count=count, status="pending", label=label)
def act_swarm_list():
swarms = _swarm_load()
items = []
for sid, swarm in sorted(
swarms.items(), key=lambda kv: kv[1].get("created_ts", ""),
reverse=True):
counts = _swarm_counts(swarm)
item = {
"swarm_id": sid,
"label": swarm.get("label"),
"status": swarm.get("status"),
"count": swarm.get("count"),
"task_preview": (swarm.get("task") or "")[:80],
"created_ts": swarm.get("created_ts"),
}
item.update(counts)
items.append(item)
out(True, swarms=items)
def act_swarm_status(sid):
swarms = _swarm_load()
swarm = _swarm_get(swarms, sid)
out(True, swarm=swarm, counts=_swarm_counts(swarm))
def act_swarm_attach(sid, slot_str, agent_id, session_id=None):
if not agent_id or not SWARM_AGENT_RE.match(agent_id):
fail("BAD_NAME", "agent id must match ^[A-Za-z0-9-]{1,64}$",
{"field": "agent_id", "value": agent_id})
swarms = _swarm_load()
swarm = _swarm_get(swarms, sid)
if swarm.get("status") in ("killed", "completed", "partial"):
fail("SWARM_CLOSED", "swarm is %s; cannot attach" % swarm["status"],
{"swarm_id": sid})
slot_rec = _swarm_slot(swarm, slot_str)
if slot_rec["status"] != "pending":
fail("SWARM_SLOT_BUSY",
"slot %d is %s, not pending" % (slot_rec["slot"],
slot_rec["status"]),
{"swarm_id": sid, "slot": slot_rec["slot"]})
slot_rec["agent_id"] = agent_id
slot_rec["status"] = "running"
if session_id:
slot_rec["subagent_session_id"] = str(session_id)
slot_rec["updated_ts"] = utcnow()
swarm["status"] = _swarm_rollup(swarm)
swarm["updated_ts"] = utcnow()
_swarm_save(swarms)
audit("swarm-attach", "%s/%d" % (sid, slot_rec["slot"]))
out(True, swarm_id=sid, slot=slot_rec["slot"], agent_id=agent_id,
subagent_session_id=session_id, status="running")
def act_swarm_report(sid, slot_str):
swarms = _swarm_load()
swarm = _swarm_get(swarms, sid)
slot_rec = _swarm_slot(swarm, slot_str)
if slot_rec["status"] in ("done", "failed", "killed"):
fail("SWARM_SLOT_CLOSED",
"slot %d already %s" % (slot_rec["slot"], slot_rec["status"]),
{"swarm_id": sid, "slot": slot_rec["slot"]})
raw = sys.stdin.read()
ok = True
result = None
if raw.strip():
try:
payload = json.loads(raw)
if isinstance(payload, dict):
ok = bool(payload.get("ok", True))
result = payload.get("result", payload)
else:
result = payload
except ValueError:
result = raw[:4000] # plain-text result
slot_rec["status"] = "done" if ok else "failed"
slot_rec["result"] = result
slot_rec["updated_ts"] = utcnow()
swarm["status"] = _swarm_rollup(swarm)
swarm["updated_ts"] = utcnow()
_swarm_save(swarms)
audit("swarm-report", "%s/%d" % (sid, slot_rec["slot"]))
out(True, swarm_id=sid, slot=slot_rec["slot"],
status=slot_rec["status"], swarm_status=swarm["status"])
def act_swarm_results(sid):
swarms = _swarm_load()
swarm = _swarm_get(swarms, sid)
results = [
{"slot": s["slot"], "agent_id": s.get("agent_id"),
"status": s["status"], "result": s.get("result")}
for s in swarm.get("slots", [])
]
out(True, swarm_id=sid, status=swarm.get("status"),
counts=_swarm_counts(swarm), results=results)
def act_swarm_prune(stale_hours=6, confirm=False):
import datetime
swarms = _swarm_load()
cutoff = datetime.datetime.now(datetime.timezone.utc) - datetime.timedelta(hours=stale_hours)
pruned = []
for sid, swarm in swarms.items():
st = swarm.get("status")
up_ts = swarm.get("updated_ts") or swarm.get("created_ts")
if not up_ts:
continue
try:
dt = datetime.datetime.fromisoformat(up_ts.replace("Z", "+00:00"))
except Exception:
continue
if dt < cutoff and st in ("partial", "completed", "killed"):
pruned.append(sid)
if not confirm:
fail("CONFIRM_REQUIRED",
f"prune transitions {len(pruned)} stale swarms to 'archived' -- pass --confirm to proceed",
{"stale_hours": stale_hours, "matching_count": len(pruned), "sample": pruned[:10]})
now = utcnow()
archived_count = 0
for sid in pruned:
swarms[sid]["status"] = "archived"
swarms[sid]["updated_ts"] = now
archived_count += 1
_swarm_save(swarms)
audit("swarm-prune", f"stale_hours={stale_hours} count={archived_count}")
out(True, archived_count=archived_count, stale_hours=stale_hours)
def act_swarm_kill(sid, confirm=False):
swarms = _swarm_load()
swarm = _swarm_get(swarms, sid)
if not confirm:
fail("CONFIRM_REQUIRED",
"kill terminates a swarm -- pass --confirm to proceed",
{"would_kill": {"swarm_id": sid,
"active_slots": _swarm_counts(swarm)}})
now = utcnow()
killed_slots = 0
for s in swarm.get("slots", []):
if s["status"] in ("pending", "running"):
s["status"] = "killed"
s["updated_ts"] = now
killed_slots += 1
counts = _swarm_counts(swarm)
# All slots are terminal now: killed if nothing ever completed/failed,
# otherwise partial (some work finished before the kill).
if counts["done"] + counts["failed"] > 0:
swarm["status"] = "partial"
else:
swarm["status"] = "killed"
swarm["updated_ts"] = now
_swarm_save(swarms)
audit("swarm-kill", sid)
out(True, swarm_id=sid, status=swarm["status"],
killed_slots=killed_slots)
# ---------------------------------------------------------------------------
# Quality functions — input validation, output contract, idempotency helpers,
# and self-diagnostics.
#
# quality check run self-diagnostics (validators, output
# contract, audit path, atomic writes)
# quality validate <action> [...] dry-run: validate args for <action>
# without executing it
# (hyphenated aliases: quality-check, quality-validate)
#
# Conventions:
# - qv_*() validators are NON-FATAL: they return check dicts, never exit.
# (The fatal check_*() helpers above are for live actions; qv_*() are for
# pre-flight validation and diagnostics.)
# - A check dict is {"check": name, "ok": bool, "code": str|None,
# "reason": str|None}.
# - Output contract: every response is {"ok": bool, ...}; ok=false responses
# MUST carry "code" (a member of KNOWN_ERROR_CODES) and "error".
# - Idempotency: mutating actions append state and are NOT implicitly
# retry-safe. q_atomic_write() is the primitive for retry-safe writes
# (tmp file + os.replace); IDEMPOTENT_ACTIONS lists the read-only/dry-run
# actions that are safe to retry.
# ---------------------------------------------------------------------------
QUALITY_VERSION = "1"
KNOWN_ERROR_CODES = frozenset([
"ALREADY_EXISTS", "AMEND_FAILED", "APPEND_FAILED", "APPROVAL_FAILED",
"AUDIT_UNAVAILABLE", "BAD_ARGS", "BAD_COUNT",
"BAD_LIMIT", "BAD_NAME", "BAD_NODE", "BAD_SLOT", "BAD_TASK",
"CONFIRM_REQUIRED",
"DISPATCH_FAILED", "DM_LOG_ERROR", "FLEET_ERROR", "GATEWAY_ERROR",
"GATEWAY_UNAVAILABLE", "INVALID_JOB", "INVALID_RESULT",
"INVALID_SCHEDULE", "IO_ERROR", "KEYGEN_FAILED", "KEY_EXISTS",
"LATENCY_TIMEOUT", "LOG_FAILED", "LOOP_ERROR",
"LOOP_TIMEOUT", "NAME_MISMATCH", "NOT_FOUND", "POLICY_ERROR",
"PULL_FAILED",
"SCAN_ERROR", "STRAT_ERROR", "SWARM_CLOSED", "SWARM_NOT_FOUND",
"SWARM_SLOT_BUSY", "SWARM_SLOT_CLOSED", "TIMER_CREATE_FAILED",
"TIMER_STILL_ACTIVE", "UNREAD_ERROR", "VARS_ERROR", "WATCHDOG_CHECK_ERROR",
"GIT_ERROR", "TESTS_ERROR", "TMUX_ERROR", "ONBOARD_ERROR",
])
# Read-only / dry-run actions: safe to retry.
IDEMPOTENT_ACTIONS = frozenset([
"fleet-status", "watchdog-alerts", "relay-health", "cdp-latency",
"chrome-errors", "identity-audit", "timer-list", "timer-status",
"job-list", "job-get", "job-status", "job-next", "vars-list",
"vars-get", "vars-history", "strat-list", "strat-get", "loop-status",
"loop-health", "loop-breaks", "policy", "policy-check", "policy-show",
"thread-list", "dm-log", "unread", "swarm-list", "swarm-status", "swarm-results",
"quality-check", "quality-validate", "main-loop",
"git-status", "git-diff", "git-log",
"md-audit", "md-list", "md-read", "md-diff",
"approval-check", "approval-list",
"tmux-tally", "tmux-auto-status", "onboard-connects",
])
JOB_ID_RE = re.compile(r"^[a-z0-9][a-z0-9-]{0,127}$")
SWARM_ID_RE = re.compile(r"^[a-z0-9][a-z0-9-]{0,63}$")
AGENT_REF_RE = re.compile(r"^[A-Za-z0-9_.-]{1,128}$")
DM_ID_RE = re.compile(r"^[0-9a-fA-F]{6,64}$")
# Mirrors agent_md.py validation (single-copy lives there; these keep
# dry-run validation in sync without importing the gateway module).
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_.-]+)*$")
MD_TEMPLATE_FILES = frozenset([
"SOUL.md", "PROACTIVE_PREFERENCES.md", "HEARTBEAT.md", "AGENTS.md",
"MEMORY.md", "USER.md", "TOOLS.md", "IDENTITY.md",
])
def _qcheck(name, ok, code=None, reason=None):
d = {"check": name, "ok": bool(ok)}
if code:
d["code"] = code
if reason:
d["reason"] = reason
return d
def qv_name(value, field="name"):
if not value or not NAME_RE.match(value):
return _qcheck(field, False, "BAD_NAME",
"must match ^[a-z0-9-]{1,64}$")
return _qcheck(field, True)
def qv_var_name(value, field="name"):
# Variable names (loop_health_threshold, ...) allow underscores and
# dots, unlike job names. Mirrors exec-constrained.py NAME_RE.
if not value or not VAR_NAME_RE.match(value):
return _qcheck(field, False, "BAD_NAME",
"must match ^[A-Za-z0-9_.-]{1,64}$")
return _qcheck(field, True)
def qv_agent(value):
if value not in VALID_AGENTS:
return _qcheck("agent", False, "BAD_NAME",
"unknown agent %r; valid: %s" % (value, sorted(VALID_AGENTS)))
return _qcheck("agent", True)
def qv_thread_id(value):
if not value or not THREAD_RE.match(value):
return _qcheck("thread_id", False, "BAD_NAME",
"must match ^[A-Za-z0-9-]{1,64}$")
return _qcheck("thread_id", True)
def qv_job_id(value):
if not value or not JOB_ID_RE.match(value):
return _qcheck("job_id", False, "BAD_NAME",
"must match ^[a-z0-9][a-z0-9-]{0,127}$")
return _qcheck("job_id", True)
def qv_swarm_id(value):
if not value or not SWARM_ID_RE.match(value):
return _qcheck("swarm_id", False, "BAD_NAME",
"must match ^[a-z0-9][a-z0-9-]{0,63}$")
return _qcheck("swarm_id", True)
def qv_agent_ref(value, field="agent_id"):
# Looser than qv_agent: accepts subagent UUIDs / external refs.
if not value or not AGENT_REF_RE.match(value):
return _qcheck(field, False, "BAD_NAME",
"must match ^[A-Za-z0-9_.-]{1,128}$")
return _qcheck(field, True)
def qv_dm_id(value):
if not value or not DM_ID_RE.match(value):
return _qcheck("dm_id", False, "BAD_NAME",
"must be 6-64 hex chars")
return _qcheck("dm_id", True)
def qv_relpath(value, field="path"):
if not value or not isinstance(value, str):
return _qcheck(field, False, "BAD_NAME", "path must be non-empty")
if ".." in value or value.startswith("/"):
return _qcheck(field, False, "BAD_NAME",
"path must be repo-relative without '..'")
if any(ord(c) < 32 or ord(c) == 127 for c in value):
return _qcheck(field, False, "BAD_NAME",
"path contains control characters")
return _qcheck(field, True)
def qv_test_module(value):
if not value or not TEST_MODULE_RE.match(value):
return _qcheck("module", False, "BAD_NAME",
"must match ^tests\\.[a-z0-9_]+$")
return _qcheck("module", True)
def qv_int_range(value, lo, hi, field, code):
try:
n = int(str(value))
except (TypeError, ValueError):
return _qcheck(field, False, code, "not an integer: %r" % (value,))
if not (lo <= n <= hi):
return _qcheck(field, False, code, "must be %d..%d" % (lo, hi))
return _qcheck(field, True)
def qv_count(value, lo=1, hi=50):
return qv_int_range(value, lo, hi, "count", "BAD_COUNT")
def qv_limit(value):
return qv_int_range(value, 1, 10000, "limit", "BAD_LIMIT")
def qv_slot(value):
return qv_int_range(value, 0, 9999, "slot", "BAD_SLOT")
def qv_nonempty(value, field):
if value is None or not str(value).strip():
return _qcheck(field, False, "BAD_ARGS", "%s must be non-empty" % field)
return _qcheck(field, True)
def qv_md_account(value):
if not value or not MD_ACCOUNT_RE.match(value):
return _qcheck("account", False, "BAD_NAME",
"must match ^[A-Za-z0-9][A-Za-z0-9_-]{0,31}$")
return _qcheck("account", True)
def qv_md_file(value):
if not value or value in (".", "..") or not MD_FILENAME_RE.match(value):
return _qcheck("filename", False, "BAD_NAME",
"must be a plain basename (no directories)")
return _qcheck("filename", True)
def qv_md_template(value):
if value not in MD_TEMPLATE_FILES:
return _qcheck("filename", False, "BAD_NAME",
"must be one of %s" % sorted(MD_TEMPLATE_FILES))
return _qcheck("filename", True)
def qv_md_subpath(value):
if value in (None, ""):
return _qcheck("path", True)
if not MD_SUBPATH_RE.match(value) or ".." in value.split("/"):
return _qcheck("path", False, "BAD_NAME",
"path must be a subdir without '..'")
return _qcheck("path", True)
def qv_title(value):
c = qv_nonempty(value, "title")
if not c["ok"]:
return c
if len(str(value)) > 200:
return _qcheck("title", False, "BAD_ARGS", "title must be <= 200 chars")
return _qcheck("title", True)
def qv_name_or_job_id(value, field="name_or_job_id"):
c1 = qv_name(value, field)
if c1["ok"]:
return c1
c2 = qv_job_id(value)
if c2["ok"]:
c2["check"] = field
return c2
return _qcheck(field, False, "BAD_NAME",
"must be a job name ^[a-z0-9-]{1,64}$ or job id "
"^[a-z0-9][a-z0-9-]{0,127}$")
def q_output_contract(payload):
"""Validate a box-ctl response dict against the JSON output contract.
Returns (ok_bool, issues_list)."""
issues = []
if not isinstance(payload, dict):
return False, ["payload is not a JSON object"]
if "ok" not in payload:
issues.append("missing 'ok'")
elif not isinstance(payload["ok"], bool):
issues.append("'ok' is not a boolean")
if payload.get("ok") is False:
if "code" not in payload:
issues.append("ok:false missing 'code'")
elif payload["code"] not in KNOWN_ERROR_CODES:
issues.append("unknown error code %r" % (payload["code"],))
if "error" not in payload:
issues.append("ok:false missing 'error'")
return (not issues), issues
def q_atomic_write(path, text):
"""Idempotency-safe file write: write temp + os.replace (atomic).
Returns (ok_bool, reason). Cleans up the temp file on failure."""
p = Path(path)
tmp = p.parent / (p.name + ".tmp")
try:
tmp.write_text(text)
os.replace(tmp, p)
return True, ""
except OSError as e:
try:
if tmp.exists():
tmp.unlink()
except OSError:
pass
return False, str(e)
def _qsrc():
return Path(__file__).read_text()
def act_quality_check():
"""Run box-ctl self-diagnostics. Read-only (no state mutations)."""
results = []
def rec(name, ok, detail=None):
r = {"check": name, "ok": bool(ok)}
if detail is not None:
r["detail"] = detail
results.append(r)
return ok
src = _qsrc()
# 1. validator functions present and callable
for fn in ("qv_name", "qv_agent", "qv_thread_id", "qv_job_id",
"qv_swarm_id", "qv_agent_ref", "qv_dm_id", "qv_count",
"qv_limit", "qv_slot", "qv_nonempty", "qv_title",
"qv_name_or_job_id", "q_output_contract", "q_atomic_write"):
rec("validator:%s" % fn, callable(globals().get(fn)))
# 2. every fail() code is a member of KNOWN_ERROR_CODES (typo guard)
codes = set(re.findall(r'fail\(\s*"([A-Z0-9_]+)"', src))
unknown = sorted(codes - KNOWN_ERROR_CODES)
rec("error_codes_known", not unknown,
{"codes_found": len(codes), "unknown": unknown})
# 3. every dispatched action is documented in USAGE
actions = set(re.findall(r'(?:el)?if action == "([a-z0-9-]+)"', src))
for tup in re.findall(r'elif action in \((.*?)\)', src, re.DOTALL):
actions.update(re.findall(r'"([a-z0-9-]+)"', tup))
missing = sorted(a for a in actions if a not in USAGE)
rec("usage_coverage", not missing,
{"actions": len(actions), "undocumented": missing})
# 4. output contract on the failure path (live subprocess probe)
try:
r = run([sys.executable, str(Path(__file__)),
"quality-probe-bogus-action"], timeout=30)
payload = json.loads(r.stdout)
okc, issues = q_output_contract(payload)
rec("output_contract:fail-path",
okc and payload.get("code") == "BAD_NAME", {"issues": issues})
except Exception as e:
rec("output_contract:fail-path", False, {"error": str(e)})
# 5. output contract on the success path (live subprocess probe)
try:
r = run([sys.executable, str(Path(__file__)), "swarm-list"], timeout=30)
payload = json.loads(r.stdout)
okc, issues = q_output_contract(payload)
rec("output_contract:ok-path",
okc and payload.get("ok") is True, {"issues": issues})
except Exception as e:
rec("output_contract:ok-path", False, {"error": str(e)})
# 6. audit log path writable
try:
rec("audit_path_writable", os.access(CTL_LOG.parent, os.W_OK),
{"log": str(CTL_LOG)})
except Exception as e:
rec("audit_path_writable", False, {"error": str(e)})
# 7. atomic-write round trip (idempotency primitive)
import tempfile
tmpd = tempfile.mkdtemp(prefix="qcheck-")
tp = os.path.join(tmpd, "probe.json")
okw, why = q_atomic_write(tp, '{"probe": true}')
leftover = os.path.exists(tp + ".tmp")
try:
content_ok = okw and open(tp).read() == '{"probe": true}'
except OSError:
content_ok = False
try:
if os.path.exists(tp):
os.unlink(tp)
os.rmdir(tmpd)
except OSError:
pass
rec("atomic_write_roundtrip", content_ok and not leftover,
{"write_error": why or None})
# 8. state directories exist and are writable
for label, p in (("netvm_root", NETVM_ROOT), ("jobs_dir", JOBS_DIR)):
rec("state_dir:%s" % label, p.is_dir() and os.access(p, os.W_OK),
{"path": str(p)})
passed = sum(1 for r_ in results if r_["ok"])
failed = len(results) - passed
out(True, quality_version=QUALITY_VERSION, passed=passed,
failed=failed, checks=results,
idempotent_actions=sorted(IDEMPOTENT_ACTIONS))
def _qv_flag(args, flag, takes_value=False):
"""Split flag out of an arg list copy.
Returns (present, value_or_None, remaining_args)."""
a = list(args)
if flag not in a:
return False, None, a
i = a.index(flag)
if takes_value:
if i + 1 < len(a) and not a[i + 1].startswith("--"):
return True, a[i + 1], a[:i] + a[i + 2:]
return True, None, a[:i] + a[i + 1:]
return True, None, a[:i] + a[i + 1:]
def _qv_no_unknown(args, usage):
"""All remaining args must be consumed; anything left is unknown."""
if args:
return [_qcheck("argv", False, "BAD_ARGS",
"unknown arguments %r; usage: %s" % (args, usage))]
return []
def _qv_args(action, rest):
"""Return a list of check dicts for (action, rest), or None if the
action is unknown to quality-validate."""
a = list(rest)
# -- timer -----------------------------------------------------------
if action in ("timer-status", "timer-create", "timer-start",
"timer-stop", "timer-enable", "timer-disable"):
if len(a) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: %s <name>" % action)]
return [qv_name(a[0])]
if action == "timer-delete":
_, _, b = _qv_flag(a, "--keep-job")
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: timer-delete <name> [--keep-job]")]
return [qv_name(b[0])]
if action == "timer-list":
if a:
return [_qcheck("argv", False, "BAD_ARGS", "usage: timer-list")]
return [_qcheck("argv", True)]
# -- job -------------------------------------------------------------
if action in ("job-get", "job-trigger"):
if len(a) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: %s <name>" % action)]
return [qv_name(a[0])]
if action == "job-list":
if a:
return [_qcheck("argv", False, "BAD_ARGS", "usage: job-list")]
return [_qcheck("argv", True)]
if action == "job-put":
if len(a) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: job-put <name> (job JSON on stdin)")]
return [qv_name(a[0]),
_qcheck("stdin", True, reason="job JSON payload read from "
"stdin at execution; not consumed by dry-run")]
if action == "job-delete":
_, _, b = _qv_flag(a, "--force")
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: job-delete <name> [--force]")]
return [qv_name(b[0])]
if action == "job-result":
if len(a) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: job-result <job-id> (result JSON on stdin)")]
return [qv_job_id(a[0]),
_qcheck("stdin", True, reason="result JSON payload read from "
"stdin at execution; not consumed by dry-run")]
if action == "job-status":
if len(a) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: job-status <name-or-job-id>")]
return [qv_name_or_job_id(a[0])]
if action == "job-next":
succ, _, b = _qv_flag(a, "--success")
failf, _, b = _qv_flag(b, "--fail")
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: job-next <job-id> [--success|--fail]")]
return [qv_job_id(b[0])]
if action == "job-chain":
_, _, b = _qv_flag(a, "--on-failure")
if len(b) != 2:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: job-chain <from-job> <to-job> [--on-failure]")]
return [qv_name(b[0], "from_job"), qv_name(b[1], "to_job")]
# -- notify ----------------------------------------------------------
if action == "notify":
allow, _, b = _qv_flag(a, "--allow-main-chat")
has_sc, sc_val, b = _qv_flag(b, "--sidechat", takes_value=True)
if has_sc and sc_val is None:
return [_qcheck("sidechat", False, "BAD_ARGS",
"--sidechat requires a value")]
if len(b) != 2:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: notify <agent> <message> [--sidechat <name>] "
"[--allow-main-chat]")]
checks = [qv_agent(b[0]), qv_nonempty(b[1], "message")]
if has_sc:
checks.append(qv_nonempty(sc_val, "sidechat"))
return checks
# -- thread (hyphenated) ---------------------------------------------
if action in ("thread-list", "thread-pin", "thread-unpin",
"thread-archive", "thread-unarchive", "thread-rename"):
op = action[len("thread-"):]
has_agent, agent, b = _qv_flag(a, "--agent", takes_value=True)
has_thread, thread, b = _qv_flag(b, "--thread", takes_value=True)
has_title, title, b = _qv_flag(b, "--title", takes_value=True)
has_limit, limit, b = _qv_flag(b, "--limit", takes_value=True)
has_confirm, _, b = _qv_flag(b, "--confirm")
checks = _qv_no_unknown(b, "%s --agent <agent> [...]" % action)
if has_agent and agent is None:
checks.append(_qcheck("agent", False, "BAD_ARGS",
"--agent requires a value"))
elif not has_agent:
checks.append(_qcheck("agent", False, "BAD_NAME",
"--agent is required"))
else:
checks.append(qv_agent(agent))
if op == "list":
if has_thread:
checks.append(_qcheck("thread", False, "BAD_ARGS",
"--thread does not apply to thread-list"))
if has_limit:
checks.append(qv_limit(limit) if limit is not None
else _qcheck("limit", False, "BAD_ARGS",
"--limit requires a value"))
if has_confirm:
checks.append(_qcheck("confirm", False, "BAD_ARGS",
"--confirm only applies to thread-archive"))
else:
if has_limit:
checks.append(_qcheck("limit", False, "BAD_ARGS",
"--limit only applies to thread-list"))
if has_confirm and op != "archive":
checks.append(_qcheck("confirm", False, "BAD_ARGS",
"--confirm only applies to thread-archive"))
if not has_thread or thread is None:
checks.append(_qcheck("thread_id", False, "BAD_NAME",
"--thread <id> is required"))
else:
checks.append(qv_thread_id(thread))
if op == "archive" and not has_confirm:
checks.append(_qcheck("confirm", False, "CONFIRM_REQUIRED",
"archive requires --confirm"))
if op == "rename":
if not has_title or title is None:
checks.append(_qcheck("title", False, "BAD_ARGS",
"--title is required for rename"))
else:
checks.append(qv_title(title))
return checks
# -- thread (space-separated) ----------------------------------------
if action == "thread":
if not a or a[0] not in THREAD_OPS:
return [_qcheck("sub", False, "BAD_NAME",
"usage: thread list|pin|unpin|archive|unarchive "
"<agent> [...]")]
sub, b = a[0], a[1:]
has_confirm, _, b = _qv_flag(b, "--confirm")
has_limit, limit, b = _qv_flag(b, "--limit", takes_value=True)
checks = []
if not b:
return [_qcheck("argv", False, "BAD_NAME",
"usage: thread %s <agent> [...]" % sub)]
checks.append(qv_agent(b[0]))
pos = b[1:]
if sub == "list":
if pos:
checks.append(_qcheck("argv", False, "BAD_ARGS",
"usage: thread list <agent> [--limit N]"))
if has_limit:
checks.append(qv_limit(limit) if limit is not None
else _qcheck("limit", False, "BAD_ARGS",
"--limit requires a value"))
if has_confirm:
checks.append(_qcheck("confirm", False, "BAD_ARGS",
"--confirm only applies to thread archive"))
else:
if has_limit:
checks.append(_qcheck("limit", False, "BAD_ARGS",
"--limit only applies to thread list"))
if has_confirm and sub != "archive":
checks.append(_qcheck("confirm", False, "BAD_ARGS",
"--confirm only applies to thread archive"))
if len(pos) != 1:
checks.append(_qcheck("argv", False, "BAD_NAME",
"usage: thread %s <agent> <thread-id> "
"[--confirm]" % sub))
else:
checks.append(qv_thread_id(pos[0]))
if sub == "archive" and not has_confirm:
checks.append(_qcheck("confirm", False, "CONFIRM_REQUIRED",
"archive requires --confirm"))
checks.extend(_qv_no_unknown(
[x for x in b if x.startswith("--")],
"thread %s <agent> [...]" % sub))
return checks
# -- swarm (hyphenated) ----------------------------------------------
if action == "swarm-spawn":
has_label, label, b = _qv_flag(a, "--label", takes_value=True)
checks = []
if has_label and label is None:
checks.append(_qcheck("label", False, "BAD_ARGS",
"--label requires a value"))
for x in b:
if x.startswith("--"):
checks.append(_qcheck("argv", False, "BAD_ARGS",
"usage: swarm-spawn <count> <task...> "
"[--label L]"))
return checks
if len(b) < 2:
checks.append(_qcheck("argv", False, "BAD_ARGS",
"usage: swarm-spawn <count> <task...> "
"[--label L]"))
return checks
checks.append(qv_count(b[0]))
checks.append(qv_nonempty(" ".join(b[1:]), "task"))
return checks
if action == "swarm-list":
if a:
return [_qcheck("argv", False, "BAD_ARGS", "usage: swarm-list")]
return [_qcheck("argv", True)]
if action in ("swarm-status", "swarm-results"):
if len(a) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: %s <swarm-id>" % action)]
return [qv_swarm_id(a[0])]
if action == "swarm-attach":
if len(a) != 3:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: swarm-attach <swarm-id> <slot> <agent-id>")]
return [qv_swarm_id(a[0]), qv_slot(a[1]), qv_agent_ref(a[2])]
if action == "swarm-report":
if len(a) != 2:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: swarm-report <swarm-id> <slot> "
"(result JSON on stdin)")]
return [qv_swarm_id(a[0]), qv_slot(a[1]),
_qcheck("stdin", True, reason="result JSON payload read from "
"stdin at execution; not consumed by dry-run")]
if action == "swarm-kill":
has_confirm, _, b = _qv_flag(a, "--confirm")
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: swarm-kill <swarm-id> --confirm")]
checks = [qv_swarm_id(b[0])]
if not has_confirm:
checks.append(_qcheck("confirm", False, "CONFIRM_REQUIRED",
"kill requires --confirm"))
return checks
# -- swarm (space-separated) -----------------------------------------
if action == "swarm":
if not a or a[0] not in SWARM_OPS:
return [_qcheck("sub", False, "BAD_NAME",
"usage: swarm spawn|list|status|attach|report|"
"results|kill [...]")]
return _qv_args("swarm-" + a[0], a[1:])
# -- job (space-separated) -------------------------------------------
if action == "job":
if not a or a[0] not in ("result", "status", "next", "chain"):
return [_qcheck("sub", False, "BAD_NAME",
"usage: job result|status|next|chain [...]")]
return _qv_args("job-" + a[0], a[1:])
# -- vars ------------------------------------------------------------
if action in ("vars-get", "vars-reset"):
if len(a) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: %s <name>" % action)]
return [qv_var_name(a[0])]
if action == "vars-list":
if a:
return [_qcheck("argv", False, "BAD_ARGS", "usage: vars-list")]
return [_qcheck("argv", True)]
if action == "vars-set":
if len(a) != 2:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: vars-set <name> <value>")]
return [qv_var_name(a[0]), _qcheck("value", True)]
if action == "vars-history":
checks = []
b = list(a)
if b and not b[0].isdigit():
checks.append(qv_var_name(b.pop(0), "name"))
if b:
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: vars-history [name] [limit]")]
checks.append(qv_limit(b[0]))
return checks or [_qcheck("argv", True)]
if action == "vars-rollback":
if not (1 <= len(a) <= 2):
return [_qcheck("argv", False, "BAD_ARGS",
"usage: vars-rollback <name> [revision]")]
checks = [qv_var_name(a[0])]
if len(a) == 2:
checks.append(qv_nonempty(a[1], "revision"))
return checks
# -- strategy --------------------------------------------------------
if action == "strat-list":
if a:
return [_qcheck("argv", False, "BAD_ARGS", "usage: strat-list")]
return [_qcheck("argv", True)]
if action in ("strat-get", "strat-reset"):
has_agent, agent, b = _qv_flag(a, "--agent", takes_value=True)
checks = []
if has_agent:
if agent is None:
checks.append(_qcheck("agent", False, "BAD_ARGS",
"--agent requires a value"))
else:
checks.append(qv_agent(agent))
if not b:
checks.append(_qcheck("argv", False, "BAD_ARGS",
"usage: %s <type> [subtype] [--agent AGENT]"
% action))
else:
checks.append(qv_nonempty(b[0], "type"))
if len(b) > 1:
checks.append(qv_nonempty(b[1], "subtype"))
if len(b) > 2:
checks.append(_qcheck("argv", False, "BAD_ARGS",
"usage: %s <type> [subtype] [--agent AGENT]"
% action))
return checks
if action == "strat-set":
if not a:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: strat-set <type> [JSON]")]
return [qv_nonempty(a[0], "type"),
_qcheck("stdin", True, reason="JSON payload may come from "
"stdin at execution; not consumed by dry-run")]
# -- loop ------------------------------------------------------------
if action == "loop-status":
has_agent, agent, b = _qv_flag(a, "--agent", takes_value=True)
has_limit, limit, b = _qv_flag(b, "--limit", takes_value=True)
has_status, status, b = _qv_flag(b, "--status", takes_value=True)
checks = _qv_no_unknown(b, "loop-status [--agent A] [--limit N] "
"[--status S]")
if has_agent:
checks.append(qv_agent(agent) if agent is not None
else _qcheck("agent", False, "BAD_ARGS",
"--agent requires a value"))
if has_limit:
checks.append(qv_limit(limit) if limit is not None
else _qcheck("limit", False, "BAD_ARGS",
"--limit requires a value"))
if has_status:
checks.append(qv_nonempty(status, "status") if status is not None
else _qcheck("status", False, "BAD_ARGS",
"--status requires a value"))
return checks
if action == "loop-health":
if len(a) > 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: loop-health [--threshold T]")]
if not a:
return [_qcheck("argv", True)]
try:
float(a[0])
return [_qcheck("threshold", True)]
except (TypeError, ValueError):
return [_qcheck("threshold", False, "BAD_ARGS",
"threshold must be numeric")]
if action == "loop-breaks":
if a:
return [_qcheck("argv", False, "BAD_ARGS", "usage: loop-breaks")]
return [_qcheck("argv", True)]
if action == "loop-resolve":
if not (1 <= len(a) <= 2):
return [_qcheck("argv", False, "BAD_ARGS",
"usage: loop-resolve <dm_id> [note]")]
checks = [qv_dm_id(a[0])]
if len(a) == 2:
checks.append(_qcheck("note", True))
return checks
if action == "loop-remediate":
_, _, b = _qv_flag(a, "--dry-run")
if b:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: loop-remediate [--dry-run]")]
return [_qcheck("argv", True)]
# -- policy / fleet / misc -------------------------------------------
if action == "policy":
if not a:
return [_qcheck("argv", True)]
if a[0] == "check":
if len(a) != 2:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: policy check <agent>")]
return [qv_agent(a[1])]
if a[0] == "show":
if len(a) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: policy show")]
return [_qcheck("argv", True)]
return [_qcheck("argv", False, "BAD_NAME",
"usage: policy [check <agent>|show]")]
if action in ("fleet-status", "watchdog-alerts", "relay-health",
"cdp-latency", "chrome-errors", "identity-audit"):
# watchdog-alerts / chrome-errors accept an optional --no-advance
# (peek-only read for web-surface polling; leaves the watermark).
extra_ok = (["--no-advance"] if action in ("watchdog-alerts",
"chrome-errors") else [])
if a and a != extra_ok:
usage = "%s [--no-advance]" % action if extra_ok else "%s (no arguments)" % action
return [_qcheck("argv", False, "BAD_ARGS",
"usage: %s" % usage)]
return [_qcheck("argv", True)]
if action == "main-loop":
if not a or a[0] not in ("check", "status", "enable", "disable"):
return [_qcheck("sub", False, "BAD_NAME",
"usage: main-loop check|status|enable|disable "
"[--agent <name>]")]
has_agent, agent, b = _qv_flag(a[1:], "--agent", takes_value=True)
checks = _qv_no_unknown(b, "main-loop %s [--agent <name>]" % a[0])
if has_agent:
checks.append(qv_agent(agent) if agent is not None
else _qcheck("agent", False, "BAD_ARGS",
"--agent requires a value"))
return checks
if action == "dm-log":
has_agent, agent, b = _qv_flag(a, "--agent", takes_value=True)
if len(b) > 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: dm-log [limit] [--agent <agent>]")]
checks = [qv_limit(b[0])] if b else []
if has_agent:
checks.append(qv_agent(agent) if agent is not None
else _qcheck("agent", False, "BAD_ARGS",
"--agent requires a value"))
return checks or [_qcheck("argv", True)]
if action == "unread":
has_agent, agent, b = _qv_flag(a, "--agent", takes_value=True)
checks = _qv_no_unknown(b, "unread [--agent <agent>]")
if has_agent:
checks.append(qv_agent(agent) if agent is not None
else _qcheck("agent", False, "BAD_ARGS",
"--agent requires a value"))
return checks or [_qcheck("argv", True)]
if action == "git-status":
if a:
return [_qcheck("argv", False, "BAD_ARGS", "usage: git-status")]
return [_qcheck("argv", True)]
if action == "git-diff":
has_stat, _, b = _qv_flag(a, "--stat")
has_path, path, b = _qv_flag(b, "--path", takes_value=True)
checks = _qv_no_unknown(b, "git-diff [--stat] [--path <path>]")
if has_path:
checks.append(qv_relpath(path) if path is not None
else _qcheck("path", False, "BAD_ARGS",
"--path requires a value"))
return checks or [_qcheck("argv", True)]
if action == "git-log":
has_limit, limit, b = _qv_flag(a, "--limit", takes_value=True)
has_path, path, b = _qv_flag(b, "--path", takes_value=True)
checks = _qv_no_unknown(b, "git-log [--limit N] [--path <path>]")
if has_limit:
checks.append(qv_int_range(limit, 1, 50, "limit", "BAD_LIMIT")
if limit is not None
else _qcheck("limit", False, "BAD_ARGS",
"--limit requires a value"))
if has_path:
checks.append(qv_relpath(path) if path is not None
else _qcheck("path", False, "BAD_ARGS",
"--path requires a value"))
return checks or [_qcheck("argv", True)]
if action == "tests-run":
has_filter, filt, b = _qv_flag(a, "--filter", takes_value=True)
if len(b) > 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: tests-run [tests.<module>] [--filter <pattern>]")]
checks = [qv_test_module(b[0])] if b else []
if has_filter:
if filt is None:
checks.append(_qcheck("filter", False, "BAD_ARGS",
"--filter requires a value"))
elif not filt.strip() or len(filt) > 200:
checks.append(_qcheck("filter", False, "BAD_ARGS",
"filter must be 1-200 chars"))
elif any(ord(c) < 32 or ord(c) == 127 for c in filt):
checks.append(_qcheck("filter", False, "BAD_ARGS",
"filter contains control characters"))
else:
checks.append(_qcheck("filter", True))
return checks or [_qcheck("argv", True)]
if action == "ack":
_, _, b = _qv_flag(a, "--allow-main-chat")
has_to, to, b = _qv_flag(b, "--to", takes_value=True)
has_sender, sender, b = _qv_flag(b, "--sender", takes_value=True)
has_sc, sc_val, b = _qv_flag(b, "--sidechat", takes_value=True)
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: ack <id> --to <agent> --sender <agent> "
"[--sidechat <name>] [--allow-main-chat]")]
checks = [qv_dm_id(b[0])]
if not has_to or to is None:
checks.append(_qcheck("to", False, "BAD_ARGS",
"--to <agent> is required"))
else:
checks.append(qv_agent(to))
if not has_sender or sender is None:
checks.append(_qcheck("sender", False, "BAD_ARGS",
"--sender <agent> is required"))
else:
checks.append(qv_agent(sender))
if has_sc:
checks.append(qv_nonempty(sc_val, "sidechat") if sc_val is not None
else _qcheck("sidechat", False, "BAD_ARGS",
"--sidechat requires a value"))
return checks
# -- quality (self) ----------------------------------------------------
if action == "quality-check":
if a:
return [_qcheck("argv", False, "BAD_ARGS", "usage: quality-check")]
return [_qcheck("argv", True)]
# -- md --------------------------------------------------------------
if action == "md-audit":
if not a:
return [_qcheck("argv", True)]
return [qv_md_account(x) for x in a]
if action == "md-list":
if len(a) not in (1, 2):
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-list <account> [path]")]
checks = [qv_md_account(a[0])]
if len(a) == 2:
checks.append(qv_md_subpath(a[1]))
return checks
if action == "md-read":
if len(a) != 2:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-read <account> <filename>")]
return [qv_md_account(a[0]), qv_md_file(a[1])]
if action == "md-diff":
if len(a) != 2:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-diff <account> <filename>")]
return [qv_md_account(a[0]), qv_md_template(a[1])]
if action == "md-pull":
if len(a) != 2:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-pull <account> <filename>")]
return [qv_md_account(a[0]), qv_md_template(a[1])]
if action == "md-inject-drive":
_, _, b = _qv_flag(a, "--force")
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-inject-drive <account> [--force]")]
return [qv_md_account(b[0])]
if action == "md-sync-all":
_, _, b = _qv_flag(a, "--force")
if b:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-sync-all [--force]")]
return [_qcheck("argv", True)]
if action == "md-amend":
_, _, b = _qv_flag(a, "--stdin")
_, author, b = _qv_flag(b, "--author", takes_value=True)
_, reason, b = _qv_flag(b, "--reason", takes_value=True)
_, content, b = _qv_flag(b, "--content", takes_value=True)
_, file, b = _qv_flag(b, "--file", takes_value=True)
if len(b) not in (1, 2):
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-amend <filename> [<content>] [--stdin] "
"[--content t] [--file p] [--author n] [--reason r]")]
checks = [qv_md_template(b[0])]
if author is not None:
checks.append(qv_nonempty(author, "author"))
return checks + [_qcheck("stdin", True, reason="content may come from "
"stdin/--file/--content at execution; not "
"consumed by dry-run")]
if action == "md-append":
_, _, b = _qv_flag(a, "--stdin")
_, author, b = _qv_flag(b, "--author", takes_value=True)
_, section, b = _qv_flag(b, "--section", takes_value=True)
_, content, b = _qv_flag(b, "--content", takes_value=True)
_, file, b = _qv_flag(b, "--file", takes_value=True)
if len(b) not in (1, 2):
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-append <filename> [<text>] [--stdin] "
"[--content t] [--file p] [--author n] [--section s]")]
checks = [qv_md_template(b[0])]
if author is not None:
checks.append(qv_nonempty(author, "author"))
if section is not None:
checks.append(qv_nonempty(section, "section"))
return checks + [_qcheck("stdin", True, reason="text may come from "
"stdin/--file/--content at execution; not "
"consumed by dry-run")]
if action == "md-write":
if len(a) != 3:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: md-write <account> <filename> <content>")]
return [qv_md_account(a[0]), qv_md_file(a[1]),
qv_nonempty(a[2], "content")]
# -- approval --------------------------------------------------------
if action in ("approval-check", "approval-list"):
if len(a) > 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: %s [node]" % action)]
if not a:
return [_qcheck("argv", True)]
return [qv_agent(a[0])]
if action in ("approval-allow", "approval-approve"):
_, _, b = _qv_flag(a, "--always")
_, _, b = _qv_flag(b, "--force")
_, _, b = _qv_flag(b, "--allow-main-chat")
_, message, b = _qv_flag(b, "--message", takes_value=True)
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: %s <node> [--always] [--force] [--message TEXT] "
"[--allow-main-chat]" % action)]
checks = [qv_agent(b[0])]
if message is not None:
checks.append(qv_nonempty(message, "message"))
return checks
if action == "approval-deny":
_, _, b = _qv_flag(a, "--allow-main-chat")
_, message, b = _qv_flag(b, "--message", takes_value=True)
if len(b) != 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: approval-deny <node> [--message TEXT] "
"[--allow-main-chat]")]
checks = [qv_agent(b[0])]
if message is not None:
checks.append(qv_nonempty(message, "message"))
return checks
if action == "approval-auto":
if len(a) > 1:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: approval-auto [node]")]
if not a:
return [_qcheck("argv", True)]
return [qv_agent(a[0])]
if action == "quality-validate":
if not a:
return [_qcheck("argv", False, "BAD_ARGS",
"usage: quality-validate <action> [args...]")]
inner = _qv_args(a[0], a[1:])
if inner is None:
return [_qcheck("action", False, "BAD_NAME",
"unknown action %r" % (a[0],))]
bad = [c for c in inner if not c["ok"]]
code = bad[0].get("code") if bad else None
return [_qcheck("nested:%s" % a[0], not bad, code=code,
reason="nested validation of %s" % a[0])]
if action == "quality":
if not a or a[0] not in ("check", "validate"):
return [_qcheck("sub", False, "BAD_NAME",
"usage: quality check|validate [...]")]
return _qv_args("quality-" + a[0], a[1:])
return None
def act_quality_validate(action, rest):
"""Dry-run: validate argv for <action> without executing it."""
checks = _qv_args(action, rest)
if checks is None:
fail("BAD_NAME", "quality validate: unknown action %r" % (action,))
valid = all(c["ok"] for c in checks)
out(True, action=action, valid=valid, checks=checks,
note="stdin payloads are never consumed by dry-run validation")
def main(argv):
if len(argv) < 2:
print(USAGE, file=sys.stderr)
sys.exit(2)
action = argv[1]
rest = argv[2:]
if action == "timer-list":
act_timer_list()
elif action == "timer-status":
if len(rest) != 1:
fail("BAD_NAME", "usage: timer-status <name>")
act_timer_status(rest[0])
elif action == "timer-create":
if len(rest) != 1:
fail("BAD_NAME", "usage: timer-create <name>")
act_timer_create(rest[0])
elif action == "timer-delete":
if not rest or len(rest) > 2:
fail("BAD_NAME", "usage: timer-delete <name> [--keep-job]")
keep = "--keep-job" in rest[1:]
act_timer_delete(rest[0], keep_job=keep)
elif action in ("timer-start", "timer-stop", "timer-enable", "timer-disable"):
if len(rest) != 1:
fail("BAD_NAME", f"usage: {action} <name>")
act_timer_control(rest[0], action[len("timer-"):])
elif action == "timer":
if not rest:
act_timer_list()
else:
sub = rest[0]
if sub in ("list",):
act_timer_list()
elif sub in ("status",):
if len(rest) != 2:
fail("BAD_NAME", "usage: timer status <name>")
act_timer_status(rest[1])
elif sub in ("create",):
if len(rest) != 2:
fail("BAD_NAME", "usage: timer create <name>")
act_timer_create(rest[1])
elif sub in ("delete",):
keep = "--keep-job" in rest[2:]
act_timer_delete(rest[1], keep_job=keep)
elif sub in ("start", "stop", "enable", "disable"):
if len(rest) != 2:
fail("BAD_NAME", f"usage: timer {sub} <name>")
act_timer_control(rest[1], sub)
else:
fail("BAD_NAME", "usage: timer list|status|create|delete|start|stop|enable|disable [...]")
elif action == "job-list":
include_archived = "--archived" in rest or "-a" in rest
act_job_list(include_archived=include_archived)
elif action == "job-archive":
if not rest or len(rest) > 2:
fail("BAD_NAME", "usage: job-archive <name> [--force]")
force = "--force" in rest[1:]
act_job_archive(rest[0], force=force)
elif action == "job-unarchive":
if len(rest) != 1:
fail("BAD_NAME", "usage: job-unarchive <name>")
act_job_unarchive(rest[0])
elif action == "job-get":
if len(rest) != 1:
fail("BAD_NAME", "usage: job-get <name>")
act_job_get(rest[0])
elif action == "job-put":
if len(rest) != 1:
fail("BAD_NAME", "usage: job-put <name> (JSON on stdin)")
act_job_put(rest[0])
elif action == "job-delete":
if not rest or len(rest) > 2:
fail("BAD_NAME", "usage: job-delete <name> [--force]")
force = "--force" in rest[1:]
act_job_delete(rest[0], force=force)
elif action == "job-trigger":
if len(rest) != 1:
fail("BAD_NAME", "usage: job-trigger <name>")
act_job_trigger(rest[0])
elif action == "job-result":
if len(rest) != 1:
fail("BAD_NAME", "usage: job-result <job-id> (result JSON on stdin)")
act_job_result(rest[0])
elif action == "job-status":
if len(rest) != 1:
fail("BAD_NAME", "usage: job-status <name-or-job-id>")
act_job_status(rest[0])
elif action == "job-next":
if not rest or len(rest) > 2:
fail("BAD_NAME", "usage: job-next <job-id> [--success|--fail]")
s = None
if "--success" in rest[1:]:
s = True
if "--fail" in rest[1:]:
s = False
act_job_next(rest[0], s)
elif action == "job-chain":
if len(rest) < 2 or len(rest) > 3:
fail("BAD_NAME", "usage: job-chain <from-job> <to-job> [--on-failure]")
act_job_chain(rest[0], rest[1], on_failure=("--on-failure" in rest[2:]))
elif action == "vars-list":
act_vars_list()
elif action == "vars-get":
if len(rest) != 1:
fail("BAD_NAME", "usage: vars-get <name>")
act_vars_get(rest[0])
elif action == "vars-set":
if len(rest) != 2:
fail("BAD_NAME", "usage: vars-set <name> <value>")
act_vars_set(rest[0], rest[1])
elif action == "vars-reset":
if len(rest) != 1:
fail("BAD_NAME", "usage: vars-reset <name>")
act_vars_reset(rest[0])
elif action == "vars-history":
name = None
limit = 20
if rest:
name = rest[0]
if len(rest) > 1:
try:
limit = int(rest[1])
except ValueError:
limit = 20
act_vars_history(name=name, limit=limit)
elif action == "vars-rollback":
if len(rest) < 1 or len(rest) > 2:
fail("BAD_NAME", "usage: vars-rollback <name> [revision]")
rev = rest[1] if len(rest) > 1 else None
act_vars_rollback(rest[0], revision=rev)
elif action == "strat-list":
act_strat_list()
elif action == "strat-get":
if not rest:
fail("BAD_NAME", "usage: strat-get <type> [subtype] [--agent AGENT]")
itype = rest[0]
subtype = None
agent = None
idx = 1
while idx < len(rest):
if rest[idx] == "--agent" and idx + 1 < len(rest):
agent = rest[idx + 1]
idx += 2
elif subtype is None and not rest[idx].startswith("--"):
subtype = rest[idx]
idx += 1
else:
idx += 1
act_strat_get(itype, subtype=subtype, agent=agent)
elif action == "strat-set":
if len(rest) < 1:
fail("BAD_NAME", "usage: strat-set <type> [JSON]")
payload = rest[1] if len(rest) > 1 else None
act_strat_set(rest[0], payload)
elif action == "strat-reset":
if not rest:
fail("BAD_NAME", "usage: strat-reset <type> [subtype] [--agent AGENT]")
itype = rest[0]
subtype = None
agent = None
idx = 1
while idx < len(rest):
if rest[idx] == "--agent" and idx + 1 < len(rest):
agent = rest[idx + 1]
idx += 2
elif subtype is None and not rest[idx].startswith("--"):
subtype = rest[idx]
idx += 1
else:
idx += 1
act_strat_reset(itype, subtype=subtype, agent=agent)
elif action == "loop-status":
agent = None
limit = 20
status = None
idx = 0
while idx < len(rest):
if rest[idx] == "--agent" and idx + 1 < len(rest):
agent = rest[idx + 1]
idx += 2
elif rest[idx] == "--limit" and idx + 1 < len(rest):
limit = int(rest[idx + 1])
idx += 2
elif rest[idx] == "--status" and idx + 1 < len(rest):
status = rest[idx + 1]
idx += 2
else:
idx += 1
act_loop_status(agent=agent, limit=limit, status=status)
elif action == "loop-health":
thresh = rest[0] if rest else None
act_loop_health(threshold=thresh)
elif action == "loop-breaks":
act_loop_breaks()
elif action == "loop-resolve":
if len(rest) < 1:
fail("BAD_NAME", "usage: loop-resolve <dm_id> [note]")
note = rest[1] if len(rest) > 1 else None
act_loop_resolve(rest[0], note=note)
elif action in ("ssh", "ssh-mint", "ssh-list", "ssh-show", "ssh-ports", "ssh-info", "ssh-check"):
if action == "ssh-mint":
if not rest:
fail("BAD_ARGS", "usage: ssh-mint <name> [--force]")
act_ssh_mint(rest[0], force=("--force" in rest[1:]))
elif action == "ssh-list":
act_ssh_list()
elif action == "ssh-show":
if not rest:
fail("BAD_ARGS", "usage: ssh-show <name>")
act_ssh_show(rest[0])
elif action == "ssh-ports":
act_ssh_ports()
elif action == "ssh-info":
act_ssh_info(rest[0] if rest else None)
elif action == "ssh-check":
act_ssh_check()
elif action == "ssh":
if not rest or rest[0] not in ("mint", "list", "show", "ports", "info", "check"):
fail("BAD_NAME", "usage: ssh mint|list|show|ports|info|check [...]")
sub = rest[0]
args = rest[1:]
if sub == "mint":
if not args:
fail("BAD_ARGS", "usage: ssh mint <name> [--force]")
act_ssh_mint(args[0], force=("--force" in args[1:]))
elif sub == "list":
act_ssh_list()
elif sub == "show":
if not args:
fail("BAD_ARGS", "usage: ssh show <name>")
act_ssh_show(args[0])
elif sub == "ports":
act_ssh_ports()
elif sub == "info":
act_ssh_info(args[0] if args else None)
elif sub == "check":
act_ssh_check()
elif action in ("md", "md-audit", "md-list", "md-read", "md-write", "md-diff", "md-inject-drive", "md-sync-all",
"md-amend", "md-append", "md-pull"):
if action == "md-audit":
act_md_audit(accounts=rest or None)
elif action == "md-list":
if not rest:
fail("BAD_ARGS", "usage: md-list <account> [path]")
act_md_list(rest[0], path=rest[1] if len(rest) > 1 else "")
elif action == "md-read":
if len(rest) < 2:
fail("BAD_ARGS", "usage: md-read <account> <filename>")
act_md_read(rest[0], rest[1])
elif action == "md-write":
if len(rest) < 3:
fail("BAD_ARGS", "usage: md-write <account> <filename> <content>")
act_md_write(rest[0], rest[1], rest[2])
elif action == "md-diff":
if len(rest) < 2:
fail("BAD_ARGS", "usage: md-diff <account> <filename>")
act_md_diff(rest[0], rest[1])
elif action == "md-inject-drive":
if not rest:
fail("BAD_ARGS", "usage: md-inject-drive <account> [--force]")
act_md_inject_drive(rest[0], force=("--force" in rest[1:]))
elif action == "md-sync-all":
act_md_sync_all(force=("--force" in rest))
elif action == "md-amend":
filename, content, author, reason = _md_amend_parse(
rest, "usage: md-amend <filename> (<content>|--content t|--file p|--stdin) [--author name] [--reason why]")
act_md_amend(filename, content, author=author, reason=reason)
elif action == "md-append":
filename, text, author, section = _md_append_parse(
rest, "usage: md-append <filename> (<text>|--content t|--file p|--stdin) [--author name] [--section header]")
act_md_append(filename, text, author=author, section=section)
elif action == "md-pull":
if len(rest) != 2:
fail("BAD_ARGS", "usage: md-pull <account> <filename>")
act_md_pull(rest[0], rest[1])
elif action == "md":
if not rest:
fail("BAD_NAME", "usage: md audit|list|read|write|diff|inject-drive|sync-all [...]")
sub = rest[0]
args = rest[1:]
if sub == "audit":
act_md_audit(accounts=args or None)
elif sub == "list":
if not args:
fail("BAD_ARGS", "usage: md list <account> [path]")
act_md_list(args[0], path=args[1] if len(args) > 1 else "")
elif sub == "read":
if len(args) < 2:
fail("BAD_ARGS", "usage: md read <account> <filename>")
act_md_read(args[0], args[1])
elif sub == "write":
if len(args) < 2:
fail("BAD_ARGS", "usage: md write <account> <filename> [--content text | --file path]")
content = None
if "--file" in args:
fidx = args.index("--file")
if fidx + 1 < len(args):
content = Path(args[fidx + 1]).read_text(encoding="utf-8")
elif "--content" in args:
cidx = args.index("--content")
if cidx + 1 < len(args):
content = args[cidx + 1]
elif len(args) >= 3:
content = args[2]
if content is None:
fail("BAD_ARGS", "usage: md write <account> <filename> <content>")
act_md_write(args[0], args[1], content)
elif sub == "diff":
if len(args) < 2:
fail("BAD_ARGS", "usage: md diff <account> <filename>")
act_md_diff(args[0], args[1])
elif sub == "inject-drive":
if not args:
fail("BAD_ARGS", "usage: md inject-drive <account> [--force]")
act_md_inject_drive(args[0], force=("--force" in args[1:]))
elif sub == "sync-all":
act_md_sync_all(force=("--force" in args))
elif sub == "amend":
filename, content, author, reason = _md_amend_parse(
args, "usage: md amend <filename> (<content>|--content t|--file p|--stdin) [--author name] [--reason why]")
act_md_amend(filename, content, author=author, reason=reason)
elif sub == "append":
filename, text, author, section = _md_append_parse(
args, "usage: md append <filename> (<text>|--content t|--file p|--stdin) [--author name] [--section header]")
act_md_append(filename, text, author=author, section=section)
elif sub == "pull":
if len(args) < 2:
fail("BAD_ARGS", "usage: md pull <account> <filename>")
act_md_pull(args[0], args[1])
else:
fail("BAD_NAME", f"unknown md subcommand: {sub}")
elif action == "loop-remediate":
dry = "--dry-run" in rest
act_loop_remediate(dry_run=dry)
elif action == "notify":
# notify <agent> <message> [--sidechat <name>] [--allow-main-chat] [--sender <agent>]
args = list(rest)
allow_main = False
sidechat = None
sender = None
if "--allow-main-chat" in args:
allow_main = True
args.remove("--allow-main-chat")
if "--sidechat" in args:
i = args.index("--sidechat")
if i + 1 >= len(args):
fail("BAD_NAME", "usage: notify <agent> <message> [--sidechat <name>] [--allow-main-chat] [--sender <agent>]")
sidechat = args[i + 1]
del args[i:i + 2]
if "--sender" in args:
i = args.index("--sender")
if i + 1 >= len(args):
fail("BAD_NAME", "usage: notify <agent> <message> [--sidechat <name>] [--allow-main-chat] [--sender <agent>]")
sender = args[i + 1]
del args[i:i + 2]
if len(args) != 2:
fail("BAD_NAME", "usage: notify <agent> <message> [--sidechat <name>] [--allow-main-chat] [--sender <agent>]")
act_notify(args[0], args[1], sidechat=sidechat, allow_main_chat=allow_main, sender=sender)
elif action == "policy":
if not rest:
act_policy()
elif rest[0] == "check":
if len(rest) != 2:
fail("BAD_NAME", "usage: policy check <agent>")
act_policy_check(rest[1])
elif rest[0] == "show":
if len(rest) != 1:
fail("BAD_NAME", "usage: policy show")
act_policy_show()
else:
fail("BAD_NAME", "usage: policy [check <agent>|show]")
elif action == "fleet-status":
act_fleet_status()
elif action == "watchdog-alerts":
if rest and rest != ["--no-advance"]:
fail("BAD_NAME", "usage: watchdog-alerts [--no-advance]")
act_watchdog_alerts(no_advance=rest == ["--no-advance"])
elif action == "relay-health":
act_relay_health()
elif action == "identity-audit":
if rest:
fail("BAD_ARGS", "usage: identity-audit")
act_identity_audit()
elif action == "cdp-latency":
act_cdp_latency()
elif action == "chrome-errors":
if rest and rest != ["--no-advance"]:
fail("BAD_ARGS", "usage: chrome-errors [--no-advance]")
act_chrome_errors(no_advance=rest == ["--no-advance"])
elif action == "main-loop":
if not rest or rest[0] not in ("check", "status", "enable", "disable"):
fail("BAD_NAME", "usage: main-loop check|status|enable|disable [--agent <name>]")
sub = rest[0]
agent = None
args = rest[1:]
if "--agent" in args:
i = args.index("--agent")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: main-loop %s [--agent <name>]" % sub)
agent = args[i + 1]
args = args[:i] + args[i + 2:]
if args:
fail("BAD_ARGS", "usage: main-loop %s [--agent <name>]" % sub)
act_main_loop(sub, agent)
elif action == "thread":
if not rest or rest[0] not in THREAD_OPS:
fail("BAD_NAME", "usage: thread list|pin|unpin|archive|unarchive <agent> [args]")
sub = rest[0]
args = rest[1:]
if not args:
fail("BAD_NAME", f"usage: thread {sub} <agent> [...]")
agent = args[0]
args = args[1:]
thread_id = None
limit = None
confirm = "--confirm" in args
args = [a for a in args if a != "--confirm"]
if confirm and sub != "archive":
fail("BAD_ARGS", "--confirm only applies to thread archive")
idx = 0
while idx < len(args):
if args[idx] == "--limit" and idx + 1 < len(args):
try:
limit = int(args[idx + 1])
except ValueError:
fail("BAD_LIMIT", "usage: thread list <agent> [--limit N]")
idx += 2
elif thread_id is None and not args[idx].startswith("--"):
thread_id = args[idx]
idx += 1
else:
fail("BAD_ARGS", f"usage: thread {sub} <agent> [...]")
if limit is not None and sub != "list":
fail("BAD_ARGS", "--limit only applies to thread list")
if sub in ("pin", "unpin", "archive", "unarchive") and not thread_id:
fail("BAD_NAME", f"usage: thread {sub} <agent> <thread-id> [--confirm]")
if sub == "list" and thread_id:
fail("BAD_ARGS", "usage: thread list <agent> [--limit N]")
act_thread(sub, agent, thread=thread_id, limit=limit, confirm=confirm)
elif action in ("thread-list", "thread-pin", "thread-unpin",
"thread-archive", "thread-unarchive", "thread-rename"):
op = action[len("thread-"):]
agent = None
thread_id = None
title = None
limit = None
confirm = "--confirm" in rest
args = [a for a in rest if a != "--confirm"]
if confirm and op != "archive":
fail("BAD_ARGS", "--confirm only applies to thread-archive")
idx = 0
while idx < len(args):
if args[idx] == "--agent" and idx + 1 < len(args):
agent = args[idx + 1]
idx += 2
elif args[idx] == "--thread" and idx + 1 < len(args):
thread_id = args[idx + 1]
idx += 2
elif args[idx] == "--title" and idx + 1 < len(args):
title = args[idx + 1]
idx += 2
elif args[idx] == "--limit" and idx + 1 < len(args):
try:
limit = int(args[idx + 1])
except ValueError:
fail("BAD_LIMIT", "usage: thread-list --agent <agent> [--limit N]")
idx += 2
else:
fail("BAD_ARGS", f"usage: {action} --agent <agent> [--thread <id>] [--title <t>] [--confirm] [--limit N]")
if not agent:
fail("BAD_NAME", f"usage: {action} --agent <agent> [...]")
if limit is not None and op != "list":
fail("BAD_ARGS", "--limit only applies to thread-list")
if op in ("pin", "unpin", "archive", "unarchive", "rename") and not thread_id:
fail("BAD_NAME", f"usage: {action} --agent <agent> --thread <id> [...]")
act_thread(op, agent, thread=thread_id, title=title, limit=limit, confirm=confirm)
elif action in ("swarm-spawn", "swarm-list", "swarm-status",
"swarm-attach", "swarm-report", "swarm-results",
"swarm-kill", "swarm-prune"):
op = action[len("swarm-"):]
if op == "spawn":
if len(rest) < 2:
fail("BAD_ARGS", "usage: swarm-spawn <count> <task...> [--label L]")
label = None
pargs = list(rest[1:])
if "--label" in pargs:
i = pargs.index("--label")
if i + 1 >= len(pargs):
fail("BAD_ARGS", "usage: swarm-spawn <count> <task...> [--label L]")
label = pargs[i + 1]
del pargs[i:i + 2]
for a in pargs:
if a.startswith("--"):
fail("BAD_ARGS", "usage: swarm-spawn <count> <task...> [--label L]")
task = " ".join(pargs)
if not task.strip():
fail("BAD_TASK", "task must be a non-empty string")
act_swarm_spawn(rest[0], task, label)
elif op == "list":
if rest:
fail("BAD_ARGS", "usage: swarm-list")
act_swarm_list()
elif op == "status":
if len(rest) != 1:
fail("BAD_ARGS", "usage: swarm-status <swarm-id>")
act_swarm_status(rest[0])
elif op == "attach":
if len(rest) not in (3, 4):
fail("BAD_ARGS", "usage: swarm-attach <swarm-id> <slot> <agent-id> [session-id]")
act_swarm_attach(rest[0], rest[1], rest[2], rest[3] if len(rest) > 3 else None)
elif op == "report":
if len(rest) != 2:
fail("BAD_ARGS", "usage: swarm-report <swarm-id> <slot> (result JSON on stdin)")
act_swarm_report(rest[0], rest[1])
elif op == "results":
if len(rest) != 1:
fail("BAD_ARGS", "usage: swarm-results <swarm-id>")
act_swarm_results(rest[0])
elif op == "kill":
if not rest or len(rest) > 2:
fail("BAD_ARGS", "usage: swarm-kill <swarm-id> --confirm")
act_swarm_kill(rest[0], confirm=("--confirm" in rest[1:]))
elif op == "prune":
stale_h = 6
confirm = "--confirm" in rest
if "--stale-hours" in rest:
i = rest.index("--stale-hours")
if i + 1 < len(rest):
try:
stale_h = float(rest[i + 1])
except ValueError:
pass
act_swarm_prune(stale_hours=stale_h, confirm=confirm)
elif action == "swarm":
if not rest or rest[0] not in SWARM_OPS:
fail("BAD_NAME", "usage: swarm spawn|list|status|attach|report|results|kill|prune [...]")
sub = rest[0]
args = rest[1:]
if sub == "spawn":
if len(args) < 2:
fail("BAD_ARGS", "usage: swarm spawn <count> <task...> [--label L]")
label = None
pargs = list(args[1:])
if "--label" in pargs:
i = pargs.index("--label")
if i + 1 >= len(pargs):
fail("BAD_ARGS", "usage: swarm spawn <count> <task...> [--label L]")
label = pargs[i + 1]
del pargs[i:i + 2]
for a in pargs:
if a.startswith("--"):
fail("BAD_ARGS", "usage: swarm spawn <count> <task...> [--label L]")
task = " ".join(pargs)
if not task.strip():
fail("BAD_TASK", "task must be a non-empty string")
act_swarm_spawn(args[0], task, label)
elif sub == "list":
if args:
fail("BAD_ARGS", "usage: swarm list")
act_swarm_list()
elif sub == "status":
if len(args) != 1:
fail("BAD_ARGS", "usage: swarm status <swarm-id>")
act_swarm_status(args[0])
elif sub == "attach":
if len(args) not in (3, 4):
fail("BAD_ARGS", "usage: swarm attach <swarm-id> <slot> <agent-id> [session-id]")
act_swarm_attach(args[0], args[1], args[2], args[3] if len(args) > 3 else None)
elif sub == "report":
if len(args) != 2:
fail("BAD_ARGS", "usage: swarm report <swarm-id> <slot> (result JSON on stdin)")
act_swarm_report(args[0], args[1])
elif sub == "results":
if len(args) != 1:
fail("BAD_ARGS", "usage: swarm results <swarm-id>")
act_swarm_results(args[0])
elif sub == "kill":
if not args or len(args) > 2:
fail("BAD_ARGS", "usage: swarm kill <swarm-id> --confirm")
act_swarm_kill(args[0], confirm=("--confirm" in args[1:]))
elif sub == "prune":
stale_h = 6
confirm = "--confirm" in args
if "--stale-hours" in args:
i = args.index("--stale-hours")
if i + 1 < len(args):
try:
stale_h = float(args[i+1])
except ValueError:
pass
act_swarm_prune(stale_hours=stale_h, confirm=confirm)
elif action == "job":
if not rest or rest[0] not in ("result", "status", "next", "chain"):
fail("BAD_NAME", "usage: job result|status|next|chain [...]")
sub = rest[0]
args = rest[1:]
if sub == "result":
if len(args) != 1:
fail("BAD_ARGS", "usage: job result <job-id> (result JSON on stdin)")
act_job_result(args[0])
elif sub == "status":
if len(args) != 1:
fail("BAD_ARGS", "usage: job status <name-or-job-id>")
act_job_status(args[0])
elif sub == "next":
if not args or len(args) > 2:
fail("BAD_ARGS", "usage: job next <job-id> [--success|--fail]")
s = None
if "--success" in args[1:]:
s = True
if "--fail" in args[1:]:
s = False
act_job_next(args[0], s)
elif sub == "chain":
if len(args) < 2 or len(args) > 3:
fail("BAD_ARGS", "usage: job chain <from-job> <to-job> [--on-failure]")
act_job_chain(args[0], args[1], on_failure=("--on-failure" in args[2:]))
elif action in ("quality-check", "quality-validate"):
op = action[len("quality-"):]
if op == "check":
if rest:
fail("BAD_ARGS", "usage: quality-check")
audit("quality-check")
act_quality_check()
else:
if not rest:
fail("BAD_ARGS", "usage: quality-validate <action> [args...]")
audit("quality-validate", rest[0])
act_quality_validate(rest[0], rest[1:])
elif action == "quality":
if not rest or rest[0] not in ("check", "validate"):
fail("BAD_NAME", "usage: quality check|validate [...]")
if rest[0] == "check":
if len(rest) != 1:
fail("BAD_ARGS", "usage: quality check")
audit("quality-check")
act_quality_check()
else:
if len(rest) < 2:
fail("BAD_ARGS", "usage: quality validate <action> [args...]")
audit("quality-validate", rest[1])
act_quality_validate(rest[1], rest[2:])
elif action in ("approval-check", "approval-list"):
node = rest[0] if rest else None
act_approval_check(node)
elif action in ("approval-allow", "approval-approve"):
if not rest:
fail("BAD_ARGS", "usage: approval-allow <node> [--always] [--force] [--message TEXT] [--allow-main-chat]")
node = rest[0]
rest_args = rest[1:]
always = "--always" in rest_args
force = "--force" in rest_args
allow_main_chat = "--allow-main-chat" in rest_args
message = None
if "--message" in rest_args:
mi = rest_args.index("--message")
if mi + 1 < len(rest_args):
message = rest_args[mi + 1]
act_approval_allow(node, always=always, force=force, message=message, allow_main_chat=allow_main_chat)
elif action == "approval-deny":
if not rest:
fail("BAD_ARGS", "usage: approval-deny <node> [--message TEXT] [--allow-main-chat]")
node = rest[0]
rest_args = rest[1:]
allow_main_chat = "--allow-main-chat" in rest_args
message = None
if "--message" in rest_args:
mi = rest_args.index("--message")
if mi + 1 < len(rest_args):
message = rest_args[mi + 1]
act_approval_deny(node, message=message, allow_main_chat=allow_main_chat)
elif action == "approval-auto":
node = rest[0] if rest else None
act_approval_auto(node)
elif action == "tmux-tally":
if rest:
fail("BAD_ARGS", "usage: tmux-tally")
act_tmux_tally()
elif action == "tmux-auto-status":
if rest:
fail("BAD_ARGS", "usage: tmux-auto-status")
act_tmux_auto_status()
elif action == "tmux-auto-toggle":
args = list(rest)
enable = True
if "--disable" in args:
enable = False
args = [a for a in args if a != "--disable"]
elif "--enable" in args:
enable = True
args = [a for a in args if a != "--enable"]
node = None
session = None
if "--node" in args:
i = args.index("--node")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: tmux-auto-toggle [--enable|--disable] [--node NODE] [--session SESSION]")
node = args[i + 1]
args = args[:i] + args[i + 2:]
if "--session" in args:
i = args.index("--session")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: tmux-auto-toggle [--enable|--disable] [--node NODE] [--session SESSION]")
session = args[i + 1]
args = args[:i] + args[i + 2:]
if args:
fail("BAD_ARGS", "usage: tmux-auto-toggle [--enable|--disable] [--node NODE] [--session SESSION]")
act_tmux_auto_toggle(enable=enable, node=node, session=session)
elif action == "tmux-auto-once":
dry_run = "--dry-run" in rest
args = [a for a in rest if a != "--dry-run"]
if args:
fail("BAD_ARGS", "usage: tmux-auto-once [--dry-run]")
act_tmux_auto_once(dry_run=dry_run)
elif action == "tmux":
if not rest:
fail("BAD_NAME", "usage: tmux tally|auto|status|toggle|once [...]")
sub = rest[0]
args = rest[1:]
if sub == "tally":
if args:
fail("BAD_ARGS", "usage: tmux tally")
act_tmux_tally()
elif sub in ("auto", "status"):
if args:
fail("BAD_ARGS", f"usage: tmux {sub}")
act_tmux_auto_status()
elif sub == "toggle":
enable = True
if "--disable" in args:
enable = False
args = [a for a in args if a != "--disable"]
elif "--enable" in args:
enable = True
args = [a for a in args if a != "--enable"]
node = None
session = None
if "--node" in args:
i = args.index("--node")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: tmux toggle [--enable|--disable] [--node NODE] [--session SESSION]")
node = args[i + 1]
args = args[:i] + args[i + 2:]
if "--session" in args:
i = args.index("--session")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: tmux toggle [--enable|--disable] [--node NODE] [--session SESSION]")
session = args[i + 1]
args = args[:i] + args[i + 2:]
if args:
fail("BAD_ARGS", "usage: tmux toggle [--enable|--disable] [--node NODE] [--session SESSION]")
act_tmux_auto_toggle(enable=enable, node=node, session=session)
elif sub == "once":
dry_run = "--dry-run" in args
args = [a for a in args if a != "--dry-run"]
if args:
fail("BAD_ARGS", "usage: tmux once [--dry-run]")
act_tmux_auto_once(dry_run=dry_run)
else:
fail("BAD_NAME", "usage: tmux tally|auto|status|toggle|once [...]")
elif action == "onboard-connects":
if rest:
fail("BAD_ARGS", "usage: onboard-connects")
act_onboard_connects()
elif action == "onboard":
if not rest or rest[0] != "connects":
fail("BAD_NAME", "usage: onboard connects")
if len(rest) > 1:
fail("BAD_ARGS", "usage: onboard connects")
act_onboard_connects()
elif action == "dm-log":
limit = 50
agent = None
args = list(rest)
if "--agent" in args:
i = args.index("--agent")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: dm-log [limit] [--agent <agent>]")
agent = args[i + 1]
args = args[:i] + args[i + 2:]
if len(args) > 1:
fail("BAD_ARGS", "usage: dm-log [limit] [--agent <agent>]")
if args:
try:
limit = int(args[0])
except ValueError:
fail("BAD_LIMIT", "usage: dm-log [limit] [--agent <agent>]")
act_dm_log(limit=limit, agent=agent)
elif action == "unread":
agent = None
args = list(rest)
if "--agent" in args:
i = args.index("--agent")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: unread [--agent <agent>]")
agent = args[i + 1]
args = args[:i] + args[i + 2:]
if args:
fail("BAD_ARGS", "usage: unread [--agent <agent>]")
act_unread(agent=agent)
elif action == "git-status":
if rest:
fail("BAD_ARGS", "usage: git-status")
act_git_status()
elif action == "git-diff":
args = list(rest)
stat = "--stat" in args
args = [a for a in args if a != "--stat"]
path = None
if "--path" in args:
i = args.index("--path")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: git-diff [--stat] [--path <path>]")
path = args[i + 1]
args = args[:i] + args[i + 2:]
if args:
fail("BAD_ARGS", "usage: git-diff [--stat] [--path <path>]")
act_git_diff(stat=stat, path=path)
elif action == "git-log":
args = list(rest)
limit = 10
path = None
if "--limit" in args:
i = args.index("--limit")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: git-log [--limit N] [--path <path>]")
limit = args[i + 1]
args = args[:i] + args[i + 2:]
if "--path" in args:
i = args.index("--path")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: git-log [--limit N] [--path <path>]")
path = args[i + 1]
args = args[:i] + args[i + 2:]
if args:
fail("BAD_ARGS", "usage: git-log [--limit N] [--path <path>]")
act_git_log(limit=limit, path=path)
elif action == "tests-run":
args = list(rest)
filt = None
if "--filter" in args:
i = args.index("--filter")
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: tests-run [tests.<module>] [--filter <pattern>]")
filt = args[i + 1]
args = args[:i] + args[i + 2:]
if len(args) > 1:
fail("BAD_ARGS", "usage: tests-run [tests.<module>] [--filter <pattern>]")
act_tests_run(args[0] if args else None, filter=filt)
elif action == "ack":
args = list(rest)
allow_main = "--allow-main-chat" in args
args = [a for a in args if a != "--allow-main-chat"]
to = sender = sidechat = None
for flag in ("--to", "--sender", "--sidechat"):
if flag in args:
i = args.index(flag)
if i + 1 >= len(args):
fail("BAD_ARGS", "usage: ack <id> --to <agent> --sender <agent> [--sidechat <name>] [--allow-main-chat]")
val = args[i + 1]
args = args[:i] + args[i + 2:]
if flag == "--to":
to = val
elif flag == "--sender":
sender = val
else:
sidechat = val
if len(args) != 1 or not to or not sender:
fail("BAD_ARGS", "usage: ack <id> --to <agent> --sender <agent> [--sidechat <name>] [--allow-main-chat]")
act_ack(args[0], to, sender, sidechat=sidechat, allow_main_chat=allow_main)
else:
print(USAGE, file=sys.stderr)
fail("BAD_NAME", f"unknown action: {action}")
if __name__ == "__main__":
main(sys.argv)