Files
box/bin/box-ctl.py
T
operator-main a76a776a87 Add cdp-latency-check.sh and box cdp-latency action
Proper script replacing the inline-SSH latency monitor. Probes each relay /json/version via pinned ports from netvm-names.sh, outputs name:latency_ms:code per node. box-ctl.py cdp-latency runs it and returns JSON.
2026-10-04 18:34:52 +00:00

1387 lines
47 KiB
Python
Executable File
Raw 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 subprocess
import sys
from datetime import datetime, timezone
from pathlib import Path
NETVM_ROOT = Path("/home/super/Projects/NetVM")
BIN = NETVM_ROOT / "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}$")
VALID_AGENTS = {"muse", "pip", "646", "opm"}
# 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": "heartbeat",
"pip": "646-pip-coord",
"muse": "muse 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 H * * * with stepped hour: M */n * * * -> 0/n:MM
if m.isdigit() and h.startswith("*/") and h[2:].isdigit() and all(is_star(f) for f in (dom, mon, dow)):
return f"0/{h[2:]}:{int(m):02d}"
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():
jobs = []
if JOBS_DIR.exists():
for p in sorted(JOBS_DIR.glob("*.json")):
jobs.append(_job_summary(p.stem, p))
audit("job-list")
out(True, jobs=jobs)
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)
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):
if agent not in VALID_AGENTS:
fail("BAD_NAME", f"agent must be one of {sorted(VALID_AGENTS)}")
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", "opm", "--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_watchdog_alerts():
"""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.
"""
audit("watchdog-alerts")
script = BIN / "watchdog-alert-check.sh"
r = subprocess.run([str(script)], 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_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_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():
"""Run chrome-error-scan.sh and return per-profile error counts as JSON."""
audit("chrome-errors")
cmd = ["/bin/bash", str(BIN / "chrome-error-scan.sh"), "--json"]
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):
audit("dm-log", str(limit))
cmd = [sys.executable, str(BIN / "super-cli.py"), "dm", "log", "--json", "-n", str(limit)]
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 _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))
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
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>
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]
notify:
notify <agent> <message> [--sidechat <name>] [--allow-main-chat]
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"""
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 == "job-list":
act_job_list()
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 == "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 == "loop-remediate":
dry = "--dry-run" in rest
act_loop_remediate(dry_run=dry)
elif action == "notify":
# notify <agent> <message> [--sidechat <name>] [--allow-main-chat]
args = list(rest)
allow_main = False
sidechat = 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]")
sidechat = args[i + 1]
del args[i:i + 2]
if len(args) != 2:
fail("BAD_NAME", "usage: notify <agent> <message> [--sidechat <name>] [--allow-main-chat]")
act_notify(args[0], args[1], sidechat=sidechat, allow_main_chat=allow_main)
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:
fail("BAD_NAME", "usage: watchdog-alerts")
act_watchdog_alerts()
elif action == "relay-health":
act_relay_health()
elif action == "cdp-latency":
act_cdp_latency()
elif action == "chrome-errors":
if rest:
fail("BAD_ARGS", "usage: chrome-errors")
act_chrome_errors()
elif action == "dm-log":
limit = 50
if rest:
try:
limit = int(rest[0])
except ValueError:
fail("BAD_LIMIT", "usage: dm-log [limit]")
act_dm_log(limit=limit)
else:
print(USAGE, file=sys.stderr)
fail("BAD_NAME", f"unknown action: {action}")
if __name__ == "__main__":
main(sys.argv)