#!/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 [args...] Security properties: - Fixed verb set; every argument validated before acting. - 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 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 ") 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_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 [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 timer actions: timer-list timer-status timer-create timer-delete [--keep-job] timer-start|timer-stop|timer-enable|timer-disable job actions: job-list job-get job-put (job JSON on stdin) job-delete [--force] job-trigger variable actions: vars-list vars-get vars-set vars-reset vars-history [name] [limit] vars-rollback [revision] strategy actions: strat-list strat-get [subtype] [--agent AGENT] strat-set (JSON on stdin or as arg) strat-reset [subtype] [--agent AGENT] loop actions: loop-status [--agent AGENT] [--limit N] [--status STATUS] loop-health [--threshold T] loop-breaks loop-resolve [note] loop-remediate [--dry-run] notify: notify [--sidechat ] [--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 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 ") act_timer_status(rest[0]) elif action == "timer-create": if len(rest) != 1: fail("BAD_NAME", "usage: timer-create ") act_timer_create(rest[0]) elif action == "timer-delete": if not rest or len(rest) > 2: fail("BAD_NAME", "usage: timer-delete [--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} ") 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 ") act_job_get(rest[0]) elif action == "job-put": if len(rest) != 1: fail("BAD_NAME", "usage: job-put (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 [--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 ") 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 ") act_vars_get(rest[0]) elif action == "vars-set": if len(rest) != 2: fail("BAD_NAME", "usage: vars-set ") act_vars_set(rest[0], rest[1]) elif action == "vars-reset": if len(rest) != 1: fail("BAD_NAME", "usage: vars-reset ") 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 [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 [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 [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 [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 [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 [--sidechat ] [--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 [--sidechat ] [--allow-main-chat]") sidechat = args[i + 1] del args[i:i + 2] if len(args) != 2: fail("BAD_NAME", "usage: notify [--sidechat ] [--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 ") 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 |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 == "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)