Files
box/bin/box-ctl.py
T

1455 lines
50 KiB
Python
Raw Normal View History

#!/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_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():
"""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_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():
"""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
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>
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 == "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:
fail("BAD_ARGS", "usage: chrome-errors")
act_chrome_errors()
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 == "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)