diff --git a/bin/box-ctl.py b/bin/box-ctl.py index 4ecbbb3..8d43547 100755 --- a/bin/box-ctl.py +++ b/bin/box-ctl.py @@ -1504,7 +1504,26 @@ thread actions (muse-cli gateway): thread pin thread unpin thread archive --confirm - thread unarchive """ + thread unarchive + +swarm actions (agent swarms; bl owns state, operator agents do the spawning): + swarm-spawn [--label L] + swarm-list + swarm-status + swarm-attach + swarm-report (result JSON on stdin) + swarm-results + swarm-kill --confirm + (space-separated aliases: swarm spawn|list|status|attach|report|results|kill) +quality: + quality-check run box-ctl self-diagnostics (validators, output + contract, audit path, atomic writes) + quality-validate [args...] + dry-run: validate args without executing + (space-separated alias: quality check|validate) + +dm-log: + dm-log [limit] recent DM send log""" def act_thread(op, agent, thread=None, title=None, limit=None, confirm=False): @@ -1567,6 +1586,1036 @@ def act_thread(op, agent, thread=None, title=None, limit=None, confirm=False): (p.stderr.strip() or p.stdout.strip())[:300]) +# --------------------------------------------------------------------------- +# Agent swarms — parallel subagent groups tracked on bl. +# +# box-ctl.py owns the swarm STATE (bl is the source of truth). Actual agent +# spawning is done by an operator agent in the container via subagent.spawn; +# the spawner attaches agent IDs with swarm-attach and records outcomes with +# swarm-report. A `swarm-spawner` cron agent can poll `swarm-list` for pending +# swarms and drive them without human involvement. +# +# Lifecycle: pending -> running -> completed|partial|killed +# Slot lifecycle: pending -> running -> done|failed|killed +# --------------------------------------------------------------------------- + +SWARM_FILE = NETVM_ROOT / "swarms.json" +SWARM_ID_RE = re.compile(r"^[a-z0-9-]{1,64}$") +SWARM_AGENT_RE = re.compile(r"^[A-Za-z0-9-]{1,64}$") +SWARM_OPS = ("spawn", "list", "status", "attach", "report", "results", "kill") +SWARM_MAX_COUNT = 50 +SWARM_MAX_TASK_LEN = 2000 + + +def _swarm_load(): + try: + with open(SWARM_FILE) as f: + data = json.load(f) + return data if isinstance(data, dict) else {} + except (OSError, ValueError): + return {} + + +def _swarm_save(swarms): + tmp = str(SWARM_FILE) + ".tmp" + with open(tmp, "w") as f: + json.dump(swarms, f, indent=2) + os.replace(tmp, SWARM_FILE) + + +def _swarm_new_id(): + return ("sw-" + datetime.now(timezone.utc).strftime("%Y%m%d-%H%M%S") + + "-" + os.urandom(2).hex()) + + +def _swarm_get(swarms, sid): + if not sid or not SWARM_ID_RE.match(sid): + fail("BAD_NAME", "swarm id must match ^[a-z0-9-]{1,64}$", + {"field": "swarm_id", "value": sid}) + swarm = swarms.get(sid) + if swarm is None: + fail("SWARM_NOT_FOUND", "no such swarm: %s" % sid, + {"swarm_id": sid}) + return swarm + + +def _swarm_counts(swarm): + counts = {"pending": 0, "running": 0, "done": 0, "failed": 0, + "killed": 0} + for slot in swarm.get("slots", []): + st = slot.get("status") + if st in counts: + counts[st] += 1 + return counts + + +def _swarm_rollup(swarm): + """Recompute swarm-level status from slot statuses.""" + counts = _swarm_counts(swarm) + total = len(swarm.get("slots", [])) + if total and counts["killed"] == total: + return "killed" + if total and counts["pending"] + counts["running"] == 0: + return "completed" if counts["failed"] == 0 else "partial" + if counts["running"] or counts["done"] or counts["failed"]: + return "running" + return "pending" + + +def _swarm_slot(swarm, slot_str): + try: + slot = int(slot_str) + except (TypeError, ValueError): + fail("BAD_SLOT", "slot must be an integer", {"value": slot_str}) + slots = swarm.get("slots", []) + if not 0 <= slot < len(slots): + fail("BAD_SLOT", "slot must be 0-%d" % (len(slots) - 1), + {"value": slot_str}) + return slots[slot] + + +def act_swarm_spawn(count_str, task, label=None): + try: + count = int(count_str) + except (TypeError, ValueError): + fail("BAD_COUNT", "count must be an integer 1-50", + {"value": count_str}) + if not 1 <= count <= SWARM_MAX_COUNT: + fail("BAD_COUNT", "count must be 1-%d" % SWARM_MAX_COUNT, + {"value": count}) + if not task or not task.strip(): + fail("BAD_TASK", "task must be a non-empty string") + if len(task) > SWARM_MAX_TASK_LEN: + fail("BAD_TASK", "task must be <= %d chars" % SWARM_MAX_TASK_LEN, + {"length": len(task)}) + if label is not None and not SWARM_ID_RE.match(label): + fail("BAD_NAME", "label must match ^[a-z0-9-]{1,64}$", + {"field": "label", "value": label}) + swarms = _swarm_load() + sid = _swarm_new_id() + while sid in swarms: # vanishingly unlikely; be safe + sid = _swarm_new_id() + now = utcnow() + swarm = { + "swarm_id": sid, + "task": task, + "count": count, + "label": label, + "status": "pending", + "created_ts": now, + "created_by": os.environ.get("BOX_CALLER", "unknown"), + "updated_ts": now, + "slots": [ + {"slot": i, "agent_id": None, "status": "pending", + "result": None, "updated_ts": now} + for i in range(count) + ], + } + swarms[sid] = swarm + _swarm_save(swarms) + audit("swarm-spawn", sid) + out(True, swarm_id=sid, count=count, status="pending", label=label) + + +def act_swarm_list(): + swarms = _swarm_load() + items = [] + for sid, swarm in sorted( + swarms.items(), key=lambda kv: kv[1].get("created_ts", ""), + reverse=True): + counts = _swarm_counts(swarm) + item = { + "swarm_id": sid, + "label": swarm.get("label"), + "status": swarm.get("status"), + "count": swarm.get("count"), + "task_preview": (swarm.get("task") or "")[:80], + "created_ts": swarm.get("created_ts"), + } + item.update(counts) + items.append(item) + out(True, swarms=items) + + +def act_swarm_status(sid): + swarms = _swarm_load() + swarm = _swarm_get(swarms, sid) + out(True, swarm=swarm, counts=_swarm_counts(swarm)) + + +def act_swarm_attach(sid, slot_str, agent_id): + if not agent_id or not SWARM_AGENT_RE.match(agent_id): + fail("BAD_NAME", "agent id must match ^[A-Za-z0-9-]{1,64}$", + {"field": "agent_id", "value": agent_id}) + swarms = _swarm_load() + swarm = _swarm_get(swarms, sid) + if swarm.get("status") in ("killed", "completed", "partial"): + fail("SWARM_CLOSED", "swarm is %s; cannot attach" % swarm["status"], + {"swarm_id": sid}) + slot_rec = _swarm_slot(swarm, slot_str) + if slot_rec["status"] != "pending": + fail("SWARM_SLOT_BUSY", + "slot %d is %s, not pending" % (slot_rec["slot"], + slot_rec["status"]), + {"swarm_id": sid, "slot": slot_rec["slot"]}) + slot_rec["agent_id"] = agent_id + slot_rec["status"] = "running" + slot_rec["updated_ts"] = utcnow() + swarm["status"] = _swarm_rollup(swarm) + swarm["updated_ts"] = utcnow() + _swarm_save(swarms) + audit("swarm-attach", "%s/%d" % (sid, slot_rec["slot"])) + out(True, swarm_id=sid, slot=slot_rec["slot"], agent_id=agent_id, + status="running") + + +def act_swarm_report(sid, slot_str): + swarms = _swarm_load() + swarm = _swarm_get(swarms, sid) + slot_rec = _swarm_slot(swarm, slot_str) + if slot_rec["status"] in ("done", "failed", "killed"): + fail("SWARM_SLOT_CLOSED", + "slot %d already %s" % (slot_rec["slot"], slot_rec["status"]), + {"swarm_id": sid, "slot": slot_rec["slot"]}) + raw = sys.stdin.read() + ok = True + result = None + if raw.strip(): + try: + payload = json.loads(raw) + if isinstance(payload, dict): + ok = bool(payload.get("ok", True)) + result = payload.get("result", payload) + else: + result = payload + except ValueError: + result = raw[:4000] # plain-text result + slot_rec["status"] = "done" if ok else "failed" + slot_rec["result"] = result + slot_rec["updated_ts"] = utcnow() + swarm["status"] = _swarm_rollup(swarm) + swarm["updated_ts"] = utcnow() + _swarm_save(swarms) + audit("swarm-report", "%s/%d" % (sid, slot_rec["slot"])) + out(True, swarm_id=sid, slot=slot_rec["slot"], + status=slot_rec["status"], swarm_status=swarm["status"]) + + +def act_swarm_results(sid): + swarms = _swarm_load() + swarm = _swarm_get(swarms, sid) + results = [ + {"slot": s["slot"], "agent_id": s.get("agent_id"), + "status": s["status"], "result": s.get("result")} + for s in swarm.get("slots", []) + ] + out(True, swarm_id=sid, status=swarm.get("status"), + counts=_swarm_counts(swarm), results=results) + + +def act_swarm_kill(sid, confirm=False): + swarms = _swarm_load() + swarm = _swarm_get(swarms, sid) + if not confirm: + fail("CONFIRM_REQUIRED", + "kill terminates a swarm -- pass --confirm to proceed", + {"would_kill": {"swarm_id": sid, + "active_slots": _swarm_counts(swarm)}}) + now = utcnow() + killed_slots = 0 + for s in swarm.get("slots", []): + if s["status"] in ("pending", "running"): + s["status"] = "killed" + s["updated_ts"] = now + killed_slots += 1 + counts = _swarm_counts(swarm) + # All slots are terminal now: killed if nothing ever completed/failed, + # otherwise partial (some work finished before the kill). + if counts["done"] + counts["failed"] > 0: + swarm["status"] = "partial" + else: + swarm["status"] = "killed" + swarm["updated_ts"] = now + _swarm_save(swarms) + audit("swarm-kill", sid) + out(True, swarm_id=sid, status=swarm["status"], + killed_slots=killed_slots) + + +# --------------------------------------------------------------------------- +# Quality functions — input validation, output contract, idempotency helpers, +# and self-diagnostics. +# +# quality check run self-diagnostics (validators, output +# contract, audit path, atomic writes) +# quality validate [...] dry-run: validate args for +# without executing it +# (hyphenated aliases: quality-check, quality-validate) +# +# Conventions: +# - qv_*() validators are NON-FATAL: they return check dicts, never exit. +# (The fatal check_*() helpers above are for live actions; qv_*() are for +# pre-flight validation and diagnostics.) +# - A check dict is {"check": name, "ok": bool, "code": str|None, +# "reason": str|None}. +# - Output contract: every response is {"ok": bool, ...}; ok=false responses +# MUST carry "code" (a member of KNOWN_ERROR_CODES) and "error". +# - Idempotency: mutating actions append state and are NOT implicitly +# retry-safe. q_atomic_write() is the primitive for retry-safe writes +# (tmp file + os.replace); IDEMPOTENT_ACTIONS lists the read-only/dry-run +# actions that are safe to retry. +# --------------------------------------------------------------------------- + +QUALITY_VERSION = "1" + +KNOWN_ERROR_CODES = frozenset([ + "ALREADY_EXISTS", "AUDIT_UNAVAILABLE", "BAD_ARGS", "BAD_COUNT", + "BAD_LIMIT", "BAD_NAME", "BAD_SLOT", "BAD_TASK", "CONFIRM_REQUIRED", + "DISPATCH_FAILED", "DM_LOG_ERROR", "FLEET_ERROR", "GATEWAY_ERROR", + "GATEWAY_UNAVAILABLE", "INVALID_JOB", "INVALID_RESULT", + "INVALID_SCHEDULE", "LATENCY_TIMEOUT", "LOG_FAILED", "LOOP_ERROR", + "LOOP_TIMEOUT", "NAME_MISMATCH", "NOT_FOUND", "POLICY_ERROR", + "SCAN_ERROR", "STRAT_ERROR", "SWARM_CLOSED", "SWARM_NOT_FOUND", + "SWARM_SLOT_BUSY", "SWARM_SLOT_CLOSED", "TIMER_CREATE_FAILED", + "TIMER_STILL_ACTIVE", "VARS_ERROR", "WATCHDOG_CHECK_ERROR", +]) + +# Read-only / dry-run actions: safe to retry. +IDEMPOTENT_ACTIONS = frozenset([ + "fleet-status", "watchdog-alerts", "relay-health", "cdp-latency", + "chrome-errors", "identity-audit", "timer-list", "timer-status", + "job-list", "job-get", "job-status", "job-next", "vars-list", + "vars-get", "vars-history", "strat-list", "strat-get", "loop-status", + "loop-health", "loop-breaks", "policy", "policy-check", "policy-show", + "thread-list", "dm-log", "swarm-list", "swarm-status", "swarm-results", + "quality-check", "quality-validate", "main-loop", +]) + +JOB_ID_RE = re.compile(r"^[a-z0-9][a-z0-9-]{0,127}$") +SWARM_ID_RE = re.compile(r"^[a-z0-9][a-z0-9-]{0,63}$") +AGENT_REF_RE = re.compile(r"^[A-Za-z0-9_.-]{1,128}$") +DM_ID_RE = re.compile(r"^[0-9a-fA-F]{6,64}$") + + +def _qcheck(name, ok, code=None, reason=None): + d = {"check": name, "ok": bool(ok)} + if code: + d["code"] = code + if reason: + d["reason"] = reason + return d + + +def qv_name(value, field="name"): + if not value or not NAME_RE.match(value): + return _qcheck(field, False, "BAD_NAME", + "must match ^[a-z0-9-]{1,64}$") + return _qcheck(field, True) + + +def qv_agent(value): + if value not in VALID_AGENTS: + return _qcheck("agent", False, "BAD_NAME", + "unknown agent %r; valid: %s" % (value, sorted(VALID_AGENTS))) + return _qcheck("agent", True) + + +def qv_thread_id(value): + if not value or not THREAD_RE.match(value): + return _qcheck("thread_id", False, "BAD_NAME", + "must match ^[A-Za-z0-9-]{1,64}$") + return _qcheck("thread_id", True) + + +def qv_job_id(value): + if not value or not JOB_ID_RE.match(value): + return _qcheck("job_id", False, "BAD_NAME", + "must match ^[a-z0-9][a-z0-9-]{0,127}$") + return _qcheck("job_id", True) + + +def qv_swarm_id(value): + if not value or not SWARM_ID_RE.match(value): + return _qcheck("swarm_id", False, "BAD_NAME", + "must match ^[a-z0-9][a-z0-9-]{0,63}$") + return _qcheck("swarm_id", True) + + +def qv_agent_ref(value, field="agent_id"): + # Looser than qv_agent: accepts subagent UUIDs / external refs. + if not value or not AGENT_REF_RE.match(value): + return _qcheck(field, False, "BAD_NAME", + "must match ^[A-Za-z0-9_.-]{1,128}$") + return _qcheck(field, True) + + +def qv_dm_id(value): + if not value or not DM_ID_RE.match(value): + return _qcheck("dm_id", False, "BAD_NAME", + "must be 6-64 hex chars") + return _qcheck("dm_id", True) + + +def qv_int_range(value, lo, hi, field, code): + try: + n = int(str(value)) + except (TypeError, ValueError): + return _qcheck(field, False, code, "not an integer: %r" % (value,)) + if not (lo <= n <= hi): + return _qcheck(field, False, code, "must be %d..%d" % (lo, hi)) + return _qcheck(field, True) + + +def qv_count(value, lo=1, hi=50): + return qv_int_range(value, lo, hi, "count", "BAD_COUNT") + + +def qv_limit(value): + return qv_int_range(value, 1, 10000, "limit", "BAD_LIMIT") + + +def qv_slot(value): + return qv_int_range(value, 0, 9999, "slot", "BAD_SLOT") + + +def qv_nonempty(value, field): + if value is None or not str(value).strip(): + return _qcheck(field, False, "BAD_ARGS", "%s must be non-empty" % field) + return _qcheck(field, True) + + +def qv_title(value): + c = qv_nonempty(value, "title") + if not c["ok"]: + return c + if len(str(value)) > 200: + return _qcheck("title", False, "BAD_ARGS", "title must be <= 200 chars") + return _qcheck("title", True) + + +def qv_name_or_job_id(value, field="name_or_job_id"): + c1 = qv_name(value, field) + if c1["ok"]: + return c1 + c2 = qv_job_id(value) + if c2["ok"]: + c2["check"] = field + return c2 + return _qcheck(field, False, "BAD_NAME", + "must be a job name ^[a-z0-9-]{1,64}$ or job id " + "^[a-z0-9][a-z0-9-]{0,127}$") + + +def q_output_contract(payload): + """Validate a box-ctl response dict against the JSON output contract. + Returns (ok_bool, issues_list).""" + issues = [] + if not isinstance(payload, dict): + return False, ["payload is not a JSON object"] + if "ok" not in payload: + issues.append("missing 'ok'") + elif not isinstance(payload["ok"], bool): + issues.append("'ok' is not a boolean") + if payload.get("ok") is False: + if "code" not in payload: + issues.append("ok:false missing 'code'") + elif payload["code"] not in KNOWN_ERROR_CODES: + issues.append("unknown error code %r" % (payload["code"],)) + if "error" not in payload: + issues.append("ok:false missing 'error'") + return (not issues), issues + + +def q_atomic_write(path, text): + """Idempotency-safe file write: write temp + os.replace (atomic). + Returns (ok_bool, reason). Cleans up the temp file on failure.""" + p = Path(path) + tmp = p.parent / (p.name + ".tmp") + try: + tmp.write_text(text) + os.replace(tmp, p) + return True, "" + except OSError as e: + try: + if tmp.exists(): + tmp.unlink() + except OSError: + pass + return False, str(e) + + +def _qsrc(): + return Path(__file__).read_text() + + +def act_quality_check(): + """Run box-ctl self-diagnostics. Read-only (no state mutations).""" + results = [] + + def rec(name, ok, detail=None): + r = {"check": name, "ok": bool(ok)} + if detail is not None: + r["detail"] = detail + results.append(r) + return ok + + src = _qsrc() + + # 1. validator functions present and callable + for fn in ("qv_name", "qv_agent", "qv_thread_id", "qv_job_id", + "qv_swarm_id", "qv_agent_ref", "qv_dm_id", "qv_count", + "qv_limit", "qv_slot", "qv_nonempty", "qv_title", + "qv_name_or_job_id", "q_output_contract", "q_atomic_write"): + rec("validator:%s" % fn, callable(globals().get(fn))) + + # 2. every fail() code is a member of KNOWN_ERROR_CODES (typo guard) + codes = set(re.findall(r'fail\(\s*"([A-Z0-9_]+)"', src)) + unknown = sorted(codes - KNOWN_ERROR_CODES) + rec("error_codes_known", not unknown, + {"codes_found": len(codes), "unknown": unknown}) + + # 3. every dispatched action is documented in USAGE + actions = set(re.findall(r'(?:el)?if action == "([a-z0-9-]+)"', src)) + for tup in re.findall(r'elif action in \((.*?)\)', src, re.DOTALL): + actions.update(re.findall(r'"([a-z0-9-]+)"', tup)) + missing = sorted(a for a in actions if a not in USAGE) + rec("usage_coverage", not missing, + {"actions": len(actions), "undocumented": missing}) + + # 4. output contract on the failure path (live subprocess probe) + try: + r = run([sys.executable, str(Path(__file__)), + "quality-probe-bogus-action"], timeout=30) + payload = json.loads(r.stdout) + okc, issues = q_output_contract(payload) + rec("output_contract:fail-path", + okc and payload.get("code") == "BAD_NAME", {"issues": issues}) + except Exception as e: + rec("output_contract:fail-path", False, {"error": str(e)}) + + # 5. output contract on the success path (live subprocess probe) + try: + r = run([sys.executable, str(Path(__file__)), "swarm-list"], timeout=30) + payload = json.loads(r.stdout) + okc, issues = q_output_contract(payload) + rec("output_contract:ok-path", + okc and payload.get("ok") is True, {"issues": issues}) + except Exception as e: + rec("output_contract:ok-path", False, {"error": str(e)}) + + # 6. audit log path writable + try: + rec("audit_path_writable", os.access(CTL_LOG.parent, os.W_OK), + {"log": str(CTL_LOG)}) + except Exception as e: + rec("audit_path_writable", False, {"error": str(e)}) + + # 7. atomic-write round trip (idempotency primitive) + import tempfile + tmpd = tempfile.mkdtemp(prefix="qcheck-") + tp = os.path.join(tmpd, "probe.json") + okw, why = q_atomic_write(tp, '{"probe": true}') + leftover = os.path.exists(tp + ".tmp") + try: + content_ok = okw and open(tp).read() == '{"probe": true}' + except OSError: + content_ok = False + try: + if os.path.exists(tp): + os.unlink(tp) + os.rmdir(tmpd) + except OSError: + pass + rec("atomic_write_roundtrip", content_ok and not leftover, + {"write_error": why or None}) + + # 8. state directories exist and are writable + for label, p in (("netvm_root", NETVM_ROOT), ("jobs_dir", JOBS_DIR)): + rec("state_dir:%s" % label, p.is_dir() and os.access(p, os.W_OK), + {"path": str(p)}) + + passed = sum(1 for r_ in results if r_["ok"]) + failed = len(results) - passed + out(True, quality_version=QUALITY_VERSION, passed=passed, + failed=failed, checks=results, + idempotent_actions=sorted(IDEMPOTENT_ACTIONS)) + + +def _qv_flag(args, flag, takes_value=False): + """Split flag out of an arg list copy. + Returns (present, value_or_None, remaining_args).""" + a = list(args) + if flag not in a: + return False, None, a + i = a.index(flag) + if takes_value: + if i + 1 < len(a) and not a[i + 1].startswith("--"): + return True, a[i + 1], a[:i] + a[i + 2:] + return True, None, a[:i] + a[i + 1:] + return True, None, a[:i] + a[i + 1:] + + +def _qv_no_unknown(args, usage): + """All remaining args must be consumed; anything left is unknown.""" + if args: + return [_qcheck("argv", False, "BAD_ARGS", + "unknown arguments %r; usage: %s" % (args, usage))] + return [] + + +def _qv_args(action, rest): + """Return a list of check dicts for (action, rest), or None if the + action is unknown to quality-validate.""" + a = list(rest) + + # -- timer ----------------------------------------------------------- + if action in ("timer-status", "timer-create", "timer-start", + "timer-stop", "timer-enable", "timer-disable"): + if len(a) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: %s " % action)] + return [qv_name(a[0])] + if action == "timer-delete": + _, _, b = _qv_flag(a, "--keep-job") + if len(b) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: timer-delete [--keep-job]")] + return [qv_name(b[0])] + if action == "timer-list": + if a: + return [_qcheck("argv", False, "BAD_ARGS", "usage: timer-list")] + return [_qcheck("argv", True)] + + # -- job ------------------------------------------------------------- + if action in ("job-get", "job-trigger"): + if len(a) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: %s " % action)] + return [qv_name(a[0])] + if action == "job-list": + if a: + return [_qcheck("argv", False, "BAD_ARGS", "usage: job-list")] + return [_qcheck("argv", True)] + if action == "job-put": + if len(a) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: job-put (job JSON on stdin)")] + return [qv_name(a[0]), + _qcheck("stdin", True, reason="job JSON payload read from " + "stdin at execution; not consumed by dry-run")] + if action == "job-delete": + _, _, b = _qv_flag(a, "--force") + if len(b) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: job-delete [--force]")] + return [qv_name(b[0])] + if action == "job-result": + if len(a) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: job-result (result JSON on stdin)")] + return [qv_job_id(a[0]), + _qcheck("stdin", True, reason="result JSON payload read from " + "stdin at execution; not consumed by dry-run")] + if action == "job-status": + if len(a) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: job-status ")] + return [qv_name_or_job_id(a[0])] + if action == "job-next": + succ, _, b = _qv_flag(a, "--success") + failf, _, b = _qv_flag(b, "--fail") + if len(b) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: job-next [--success|--fail]")] + return [qv_job_id(b[0])] + if action == "job-chain": + _, _, b = _qv_flag(a, "--on-failure") + if len(b) != 2: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: job-chain [--on-failure]")] + return [qv_name(b[0], "from_job"), qv_name(b[1], "to_job")] + + # -- notify ---------------------------------------------------------- + if action == "notify": + allow, _, b = _qv_flag(a, "--allow-main-chat") + has_sc, sc_val, b = _qv_flag(b, "--sidechat", takes_value=True) + if has_sc and sc_val is None: + return [_qcheck("sidechat", False, "BAD_ARGS", + "--sidechat requires a value")] + if len(b) != 2: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: notify [--sidechat ] " + "[--allow-main-chat]")] + checks = [qv_agent(b[0]), qv_nonempty(b[1], "message")] + if has_sc: + checks.append(qv_nonempty(sc_val, "sidechat")) + return checks + + # -- thread (hyphenated) --------------------------------------------- + if action in ("thread-list", "thread-pin", "thread-unpin", + "thread-archive", "thread-unarchive", "thread-rename"): + op = action[len("thread-"):] + has_agent, agent, b = _qv_flag(a, "--agent", takes_value=True) + has_thread, thread, b = _qv_flag(b, "--thread", takes_value=True) + has_title, title, b = _qv_flag(b, "--title", takes_value=True) + has_limit, limit, b = _qv_flag(b, "--limit", takes_value=True) + has_confirm, _, b = _qv_flag(b, "--confirm") + checks = _qv_no_unknown(b, "%s --agent [...]" % action) + if has_agent and agent is None: + checks.append(_qcheck("agent", False, "BAD_ARGS", + "--agent requires a value")) + elif not has_agent: + checks.append(_qcheck("agent", False, "BAD_NAME", + "--agent is required")) + else: + checks.append(qv_agent(agent)) + if op == "list": + if has_thread: + checks.append(_qcheck("thread", False, "BAD_ARGS", + "--thread does not apply to thread-list")) + if has_limit: + checks.append(qv_limit(limit) if limit is not None + else _qcheck("limit", False, "BAD_ARGS", + "--limit requires a value")) + if has_confirm: + checks.append(_qcheck("confirm", False, "BAD_ARGS", + "--confirm only applies to thread-archive")) + else: + if has_limit: + checks.append(_qcheck("limit", False, "BAD_ARGS", + "--limit only applies to thread-list")) + if has_confirm and op != "archive": + checks.append(_qcheck("confirm", False, "BAD_ARGS", + "--confirm only applies to thread-archive")) + if not has_thread or thread is None: + checks.append(_qcheck("thread_id", False, "BAD_NAME", + "--thread is required")) + else: + checks.append(qv_thread_id(thread)) + if op == "archive" and not has_confirm: + checks.append(_qcheck("confirm", False, "CONFIRM_REQUIRED", + "archive requires --confirm")) + if op == "rename": + if not has_title or title is None: + checks.append(_qcheck("title", False, "BAD_ARGS", + "--title is required for rename")) + else: + checks.append(qv_title(title)) + return checks + + # -- thread (space-separated) ---------------------------------------- + if action == "thread": + if not a or a[0] not in THREAD_OPS: + return [_qcheck("sub", False, "BAD_NAME", + "usage: thread list|pin|unpin|archive|unarchive " + " [...]")] + sub, b = a[0], a[1:] + has_confirm, _, b = _qv_flag(b, "--confirm") + has_limit, limit, b = _qv_flag(b, "--limit", takes_value=True) + checks = [] + if not b: + return [_qcheck("argv", False, "BAD_NAME", + "usage: thread %s [...]" % sub)] + checks.append(qv_agent(b[0])) + pos = b[1:] + if sub == "list": + if pos: + checks.append(_qcheck("argv", False, "BAD_ARGS", + "usage: thread list [--limit N]")) + if has_limit: + checks.append(qv_limit(limit) if limit is not None + else _qcheck("limit", False, "BAD_ARGS", + "--limit requires a value")) + if has_confirm: + checks.append(_qcheck("confirm", False, "BAD_ARGS", + "--confirm only applies to thread archive")) + else: + if has_limit: + checks.append(_qcheck("limit", False, "BAD_ARGS", + "--limit only applies to thread list")) + if has_confirm and sub != "archive": + checks.append(_qcheck("confirm", False, "BAD_ARGS", + "--confirm only applies to thread archive")) + if len(pos) != 1: + checks.append(_qcheck("argv", False, "BAD_NAME", + "usage: thread %s " + "[--confirm]" % sub)) + else: + checks.append(qv_thread_id(pos[0])) + if sub == "archive" and not has_confirm: + checks.append(_qcheck("confirm", False, "CONFIRM_REQUIRED", + "archive requires --confirm")) + checks.extend(_qv_no_unknown( + [x for x in b if x.startswith("--")], + "thread %s [...]" % sub)) + return checks + + # -- swarm (hyphenated) ---------------------------------------------- + if action == "swarm-spawn": + has_label, label, b = _qv_flag(a, "--label", takes_value=True) + checks = [] + if has_label and label is None: + checks.append(_qcheck("label", False, "BAD_ARGS", + "--label requires a value")) + for x in b: + if x.startswith("--"): + checks.append(_qcheck("argv", False, "BAD_ARGS", + "usage: swarm-spawn " + "[--label L]")) + return checks + if len(b) < 2: + checks.append(_qcheck("argv", False, "BAD_ARGS", + "usage: swarm-spawn " + "[--label L]")) + return checks + checks.append(qv_count(b[0])) + checks.append(qv_nonempty(" ".join(b[1:]), "task")) + return checks + if action == "swarm-list": + if a: + return [_qcheck("argv", False, "BAD_ARGS", "usage: swarm-list")] + return [_qcheck("argv", True)] + if action in ("swarm-status", "swarm-results"): + if len(a) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: %s " % action)] + return [qv_swarm_id(a[0])] + if action == "swarm-attach": + if len(a) != 3: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: swarm-attach ")] + return [qv_swarm_id(a[0]), qv_slot(a[1]), qv_agent_ref(a[2])] + if action == "swarm-report": + if len(a) != 2: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: swarm-report " + "(result JSON on stdin)")] + return [qv_swarm_id(a[0]), qv_slot(a[1]), + _qcheck("stdin", True, reason="result JSON payload read from " + "stdin at execution; not consumed by dry-run")] + if action == "swarm-kill": + has_confirm, _, b = _qv_flag(a, "--confirm") + if len(b) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: swarm-kill --confirm")] + checks = [qv_swarm_id(b[0])] + if not has_confirm: + checks.append(_qcheck("confirm", False, "CONFIRM_REQUIRED", + "kill requires --confirm")) + return checks + + # -- swarm (space-separated) ----------------------------------------- + if action == "swarm": + if not a or a[0] not in SWARM_OPS: + return [_qcheck("sub", False, "BAD_NAME", + "usage: swarm spawn|list|status|attach|report|" + "results|kill [...]")] + return _qv_args("swarm-" + a[0], a[1:]) + + # -- job (space-separated) ------------------------------------------- + if action == "job": + if not a or a[0] not in ("result", "status", "next", "chain"): + return [_qcheck("sub", False, "BAD_NAME", + "usage: job result|status|next|chain [...]")] + return _qv_args("job-" + a[0], a[1:]) + + # -- vars ------------------------------------------------------------ + if action in ("vars-get", "vars-reset"): + if len(a) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: %s " % action)] + return [qv_name(a[0])] + if action == "vars-list": + if a: + return [_qcheck("argv", False, "BAD_ARGS", "usage: vars-list")] + return [_qcheck("argv", True)] + if action == "vars-set": + if len(a) != 2: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: vars-set ")] + return [qv_name(a[0]), _qcheck("value", True)] + if action == "vars-history": + checks = [] + b = list(a) + if b and not b[0].isdigit(): + checks.append(qv_name(b.pop(0), "name")) + if b: + if len(b) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: vars-history [name] [limit]")] + checks.append(qv_limit(b[0])) + return checks or [_qcheck("argv", True)] + if action == "vars-rollback": + if not (1 <= len(a) <= 2): + return [_qcheck("argv", False, "BAD_ARGS", + "usage: vars-rollback [revision]")] + checks = [qv_name(a[0])] + if len(a) == 2: + checks.append(qv_nonempty(a[1], "revision")) + return checks + + # -- strategy -------------------------------------------------------- + if action == "strat-list": + if a: + return [_qcheck("argv", False, "BAD_ARGS", "usage: strat-list")] + return [_qcheck("argv", True)] + if action in ("strat-get", "strat-reset"): + has_agent, agent, b = _qv_flag(a, "--agent", takes_value=True) + checks = [] + if has_agent: + if agent is None: + checks.append(_qcheck("agent", False, "BAD_ARGS", + "--agent requires a value")) + else: + checks.append(qv_agent(agent)) + if not b: + checks.append(_qcheck("argv", False, "BAD_ARGS", + "usage: %s [subtype] [--agent AGENT]" + % action)) + else: + checks.append(qv_nonempty(b[0], "type")) + if len(b) > 1: + checks.append(qv_nonempty(b[1], "subtype")) + if len(b) > 2: + checks.append(_qcheck("argv", False, "BAD_ARGS", + "usage: %s [subtype] [--agent AGENT]" + % action)) + return checks + if action == "strat-set": + if not a: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: strat-set [JSON]")] + return [qv_nonempty(a[0], "type"), + _qcheck("stdin", True, reason="JSON payload may come from " + "stdin at execution; not consumed by dry-run")] + + # -- loop ------------------------------------------------------------ + if action == "loop-status": + has_agent, agent, b = _qv_flag(a, "--agent", takes_value=True) + has_limit, limit, b = _qv_flag(b, "--limit", takes_value=True) + has_status, status, b = _qv_flag(b, "--status", takes_value=True) + checks = _qv_no_unknown(b, "loop-status [--agent A] [--limit N] " + "[--status S]") + if has_agent: + checks.append(qv_agent(agent) if agent is not None + else _qcheck("agent", False, "BAD_ARGS", + "--agent requires a value")) + if has_limit: + checks.append(qv_limit(limit) if limit is not None + else _qcheck("limit", False, "BAD_ARGS", + "--limit requires a value")) + if has_status: + checks.append(qv_nonempty(status, "status") if status is not None + else _qcheck("status", False, "BAD_ARGS", + "--status requires a value")) + return checks + if action == "loop-health": + if len(a) > 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: loop-health [--threshold T]")] + if not a: + return [_qcheck("argv", True)] + try: + float(a[0]) + return [_qcheck("threshold", True)] + except (TypeError, ValueError): + return [_qcheck("threshold", False, "BAD_ARGS", + "threshold must be numeric")] + if action == "loop-breaks": + if a: + return [_qcheck("argv", False, "BAD_ARGS", "usage: loop-breaks")] + return [_qcheck("argv", True)] + if action == "loop-resolve": + if not (1 <= len(a) <= 2): + return [_qcheck("argv", False, "BAD_ARGS", + "usage: loop-resolve [note]")] + checks = [qv_dm_id(a[0])] + if len(a) == 2: + checks.append(_qcheck("note", True)) + return checks + if action == "loop-remediate": + _, _, b = _qv_flag(a, "--dry-run") + if b: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: loop-remediate [--dry-run]")] + return [_qcheck("argv", True)] + + # -- policy / fleet / misc ------------------------------------------- + if action == "policy": + if not a: + return [_qcheck("argv", True)] + if a[0] == "check": + if len(a) != 2: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: policy check ")] + return [qv_agent(a[1])] + if a[0] == "show": + if len(a) != 1: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: policy show")] + return [_qcheck("argv", True)] + return [_qcheck("argv", False, "BAD_NAME", + "usage: policy [check |show]")] + if action in ("fleet-status", "watchdog-alerts", "relay-health", + "cdp-latency", "chrome-errors", "identity-audit"): + if a: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: %s (no arguments)" % action)] + return [_qcheck("argv", True)] + if action == "main-loop": + if not a or a[0] not in ("check", "status", "enable", "disable"): + return [_qcheck("sub", False, "BAD_NAME", + "usage: main-loop check|status|enable|disable " + "[--agent ]")] + has_agent, agent, b = _qv_flag(a[1:], "--agent", takes_value=True) + checks = _qv_no_unknown(b, "main-loop %s [--agent ]" % a[0]) + if has_agent: + checks.append(qv_agent(agent) if agent is not None + else _qcheck("agent", False, "BAD_ARGS", + "--agent requires a value")) + return checks + if action == "dm-log": + if len(a) > 1: + return [_qcheck("argv", False, "BAD_ARGS", "usage: dm-log [limit]")] + if not a: + return [_qcheck("argv", True)] + return [qv_limit(a[0])] + + # -- quality (self) ---------------------------------------------------- + if action == "quality-check": + if a: + return [_qcheck("argv", False, "BAD_ARGS", "usage: quality-check")] + return [_qcheck("argv", True)] + if action == "quality-validate": + if not a: + return [_qcheck("argv", False, "BAD_ARGS", + "usage: quality-validate [args...]")] + inner = _qv_args(a[0], a[1:]) + if inner is None: + return [_qcheck("action", False, "BAD_NAME", + "unknown action %r" % (a[0],))] + bad = [c for c in inner if not c["ok"]] + code = bad[0].get("code") if bad else None + return [_qcheck("nested:%s" % a[0], not bad, code=code, + reason="nested validation of %s" % a[0])] + if action == "quality": + if not a or a[0] not in ("check", "validate"): + return [_qcheck("sub", False, "BAD_NAME", + "usage: quality check|validate [...]")] + return _qv_args("quality-" + a[0], a[1:]) + + return None + + +def act_quality_validate(action, rest): + """Dry-run: validate argv for without executing it.""" + checks = _qv_args(action, rest) + if checks is None: + fail("BAD_NAME", "quality validate: unknown action %r" % (action,)) + valid = all(c["ok"] for c in checks) + out(True, action=action, valid=valid, checks=checks, + note="stdin payloads are never consumed by dry-run validation") + def main(argv): if len(argv) < 2: print(USAGE, file=sys.stderr) @@ -1894,6 +2943,99 @@ def main(argv): if op in ("pin", "unpin", "archive", "unarchive", "rename") and not thread_id: fail("BAD_NAME", f"usage: {action} --agent --thread [...]") act_thread(op, agent, thread=thread_id, title=title, limit=limit, confirm=confirm) + elif action in ("swarm-spawn", "swarm-list", "swarm-status", + "swarm-attach", "swarm-report", "swarm-results", + "swarm-kill"): + op = action[len("swarm-"):] + if op == "spawn": + if len(rest) < 2: + fail("BAD_ARGS", "usage: swarm-spawn [--label L]") + label = None + pargs = list(rest[1:]) + if "--label" in pargs: + i = pargs.index("--label") + if i + 1 >= len(pargs): + fail("BAD_ARGS", "usage: swarm-spawn [--label L]") + label = pargs[i + 1] + del pargs[i:i + 2] + for a in pargs: + if a.startswith("--"): + fail("BAD_ARGS", "usage: swarm-spawn [--label L]") + task = " ".join(pargs) + if not task.strip(): + fail("BAD_TASK", "task must be a non-empty string") + act_swarm_spawn(rest[0], task, label) + elif op == "list": + if rest: + fail("BAD_ARGS", "usage: swarm-list") + act_swarm_list() + elif op == "status": + if len(rest) != 1: + fail("BAD_ARGS", "usage: swarm-status ") + act_swarm_status(rest[0]) + elif op == "attach": + if len(rest) != 3: + fail("BAD_ARGS", "usage: swarm-attach ") + act_swarm_attach(rest[0], rest[1], rest[2]) + elif op == "report": + if len(rest) != 2: + fail("BAD_ARGS", "usage: swarm-report (result JSON on stdin)") + act_swarm_report(rest[0], rest[1]) + elif op == "results": + if len(rest) != 1: + fail("BAD_ARGS", "usage: swarm-results ") + act_swarm_results(rest[0]) + elif op == "kill": + if not rest or len(rest) > 2: + fail("BAD_ARGS", "usage: swarm-kill --confirm") + act_swarm_kill(rest[0], confirm=("--confirm" in rest[1:])) + elif action == "swarm": + if not rest or rest[0] not in SWARM_OPS: + fail("BAD_NAME", "usage: swarm spawn|list|status|attach|report|results|kill [...]") + sub = rest[0] + args = rest[1:] + if sub == "spawn": + if len(args) < 2: + fail("BAD_ARGS", "usage: swarm spawn [--label L]") + label = None + pargs = list(args[1:]) + if "--label" in pargs: + i = pargs.index("--label") + if i + 1 >= len(pargs): + fail("BAD_ARGS", "usage: swarm spawn [--label L]") + label = pargs[i + 1] + del pargs[i:i + 2] + for a in pargs: + if a.startswith("--"): + fail("BAD_ARGS", "usage: swarm spawn [--label L]") + task = " ".join(pargs) + if not task.strip(): + fail("BAD_TASK", "task must be a non-empty string") + act_swarm_spawn(args[0], task, label) + elif sub == "list": + if args: + fail("BAD_ARGS", "usage: swarm list") + act_swarm_list() + elif sub == "status": + if len(args) != 1: + fail("BAD_ARGS", "usage: swarm status ") + act_swarm_status(args[0]) + elif sub == "attach": + if len(args) != 3: + fail("BAD_ARGS", "usage: swarm attach ") + act_swarm_attach(args[0], args[1], args[2]) + elif sub == "report": + if len(args) != 2: + fail("BAD_ARGS", "usage: swarm report (result JSON on stdin)") + act_swarm_report(args[0], args[1]) + elif sub == "results": + if len(args) != 1: + fail("BAD_ARGS", "usage: swarm results ") + act_swarm_results(args[0]) + elif sub == "kill": + if not args or len(args) > 2: + fail("BAD_ARGS", "usage: swarm kill --confirm") + act_swarm_kill(args[0], confirm=("--confirm" in args[1:])) elif action == "job": if not rest or rest[0] not in ("result", "status", "next", "chain"): fail("BAD_NAME", "usage: job result|status|next|chain [...]") @@ -1920,6 +3062,31 @@ def main(argv): if len(args) < 2 or len(args) > 3: fail("BAD_ARGS", "usage: job chain [--on-failure]") act_job_chain(args[0], args[1], on_failure=("--on-failure" in args[2:])) + elif action in ("quality-check", "quality-validate"): + op = action[len("quality-"):] + if op == "check": + if rest: + fail("BAD_ARGS", "usage: quality-check") + audit("quality-check") + act_quality_check() + else: + if not rest: + fail("BAD_ARGS", "usage: quality-validate [args...]") + audit("quality-validate", rest[0]) + act_quality_validate(rest[0], rest[1:]) + elif action == "quality": + if not rest or rest[0] not in ("check", "validate"): + fail("BAD_NAME", "usage: quality check|validate [...]") + if rest[0] == "check": + if len(rest) != 1: + fail("BAD_ARGS", "usage: quality check") + audit("quality-check") + act_quality_check() + else: + if len(rest) < 2: + fail("BAD_ARGS", "usage: quality validate [args...]") + audit("quality-validate", rest[1]) + act_quality_validate(rest[1], rest[2:]) elif action == "dm-log": limit = 50 if rest: diff --git a/bin/exec-constrained.py b/bin/exec-constrained.py index 14d7675..96aa50b 100755 --- a/bin/exec-constrained.py +++ b/bin/exec-constrained.py @@ -562,6 +562,34 @@ def _cron_runs_build(a): return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'job-list'] +def _cron_timer_start_validate(raw): + if not isinstance(raw, dict): + raise OpError('args must be an object') + allowed = {'name'} + for k in raw: + if k not in allowed: + raise OpError(f'unknown arg: {k}') + return {'name': _job_name(raw.get('name'))} + + +def _cron_timer_start_build(a): + return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'timer-start', a['name']] + + +def _cron_timer_create_validate(raw): + if not isinstance(raw, dict): + raise OpError('args must be an object') + allowed = {'name'} + for k in raw: + if k not in allowed: + raise OpError(f'unknown arg: {k}') + return {'name': _job_name(raw.get('name'))} + + +def _cron_timer_create_build(a): + return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'timer-create', a['name']] + + def _vars_list_validate(raw): if raw not in ({}, None): raise OpError('vars.list takes no required args') @@ -705,6 +733,55 @@ def _service_restart_build(a): return [sys.executable, os.path.join(BIN_DIR, 'box-sys-op.py'), 'service.restart'] +def _swarm_spawn_validate(raw): + if not isinstance(raw, dict): + raise OpError('args must be an object') + allowed = {'count', 'task', 'label'} + for k in raw: + if k not in allowed: + raise OpError(f'unknown arg: {k}') + count = _opt_int(raw.get('count', 1), 1, 50, 'count') or 1 + task = _clean_message(raw.get('task')) + label = raw.get('label') + if label and not TARGET_RE.fullmatch(str(label)): + raise OpError('label must match safe identifier') + return {'count': count, 'task': task, 'label': str(label) if label else None} + + +def _swarm_spawn_build(a): + cmd = [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'swarm-spawn', str(a['count']), a['task']] + if a.get('label'): + cmd.extend(['--label', a['label']]) + return cmd + + +def _swarm_status_validate(raw): + if not isinstance(raw, dict): + raise OpError('args must be an object') + allowed = {'swarm_id', 'id'} + for k in raw: + if k not in allowed: + raise OpError(f'unknown arg: {k}') + sid = raw.get('swarm_id') or raw.get('id') + if not isinstance(sid, str) or not TARGET_RE.fullmatch(sid): + raise OpError('swarm_id must match safe identifier') + return {'swarm_id': sid} + + +def _swarm_status_build(a): + return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'swarm-status', a['swarm_id']] + + +def _swarm_list_validate(raw): + if raw not in ({}, None): + raise OpError('swarm.list takes no required args') + return {} + + +def _swarm_list_build(a): + return [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), 'swarm-list'] + + # op -> {validate, build, timeout, side_effecting, description} OPS = { 'dm.send': { @@ -782,6 +859,16 @@ OPS = { 'timeout': 300, 'side_effecting': True, 'desc': 'Trigger on-demand execution of a scheduled job', }, + 'cron.timer_create': { + 'validate': _cron_timer_create_validate, 'build': _cron_timer_create_build, + 'timeout': 30, 'side_effecting': True, + 'desc': 'Create systemd user timer unit for a job', + }, + 'cron.timer_start': { + 'validate': _cron_timer_start_validate, 'build': _cron_timer_start_build, + 'timeout': 30, 'side_effecting': True, + 'desc': 'Enable and start systemd user timer unit for a job', + }, 'vars.list': { 'validate': _vars_list_validate, 'build': _vars_list_build, 'timeout': 30, 'side_effecting': False, @@ -822,6 +909,21 @@ OPS = { 'timeout': 30, 'side_effecting': True, 'desc': 'Restart allowlisted fleet systemd service', }, + 'swarm.spawn': { + 'validate': _swarm_spawn_validate, 'build': _swarm_spawn_build, + 'timeout': 30, 'side_effecting': True, + 'desc': 'Spawn autonomous subagent swarm slots managed by Box', + }, + 'swarm.status': { + 'validate': _swarm_status_validate, 'build': _swarm_status_build, + 'timeout': 30, 'side_effecting': False, + 'desc': 'Inspect subagent swarm status and progress', + }, + 'swarm.list': { + 'validate': _swarm_list_validate, 'build': _swarm_list_build, + 'timeout': 30, 'side_effecting': False, + 'desc': 'List all subagent swarms and their counts', + }, 'exec.ping': { 'validate': _health_validate, 'build': lambda a: ['/bin/echo', 'PONG'], diff --git a/bin/response-harvester.py b/bin/response-harvester.py index 7eed868..42523f2 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -48,6 +48,8 @@ JOB_LOG = NETVM_ROOT / "job-log.jsonl" FOLLOWUPS_FILE = NETVM_ROOT / "followups.json" JOB_SIDECHATS_FILE = NETVM_ROOT / "job-sidechats.json" WAKE_SIDECHATS_FILE = Path("/home/super/sidechat-wake/wake-sidechats.json") +SWARM_FILE = NETVM_ROOT / "swarms.json" +NUDGE_TRACKER_FILE = NETVM_ROOT / "conversation-nudge-tracker.json" JOBS_DIR = NETVM_ROOT / "jobs" DISPATCH_PY = BIN_DIR / "job-dispatch.py" @@ -583,6 +585,8 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo tool_calls = parse_tool_calls(text) if tool_calls and not dry_run: for op, args in tool_calls: + if isinstance(args, dict) and "agent" not in args: + args["agent"] = agent print(f"[{agent}] Executing tool '{op}' from message {mid[:8]} in thread {thread_name or thread_id[:8]}") t_ok, t_res = execute_agent_tool(agent, op, args) tool_record = { @@ -628,7 +632,19 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo } if not dry_run: append_jsonl(JOB_LOG, job_record) - trigger_chain_next(job_id, result_text, success=not is_fail) + # Check if this is a swarm slot result: sw-YYYYMMDD-HHMMSS-xxxx/ + if "/" in job_id and job_id.startswith("sw-"): + try: + s_sid, s_slot = job_id.split("/", 1) + s_proc = subprocess.Popen( + [sys.executable, str(BIN_DIR / "box-ctl.py"), "swarm-report", s_sid, s_slot], + stdin=subprocess.PIPE, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL + ) + s_proc.communicate(input=json.dumps({"ok": not is_fail, "result": result_text}).encode("utf-8")) + except Exception as se: + sys.stderr.write(f"warning: failed to record swarm report: {se}\n") + else: + trigger_chain_next(job_id, result_text, success=not is_fail) if not is_fail: archive_ephemeral_thread(agent, thread_id, job_id=job_id) clear_matching_followups(followups, agent, thread_id, mid, text, @@ -640,6 +656,7 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo dry_run, job_id=job_id, verb=verb) else: clear_matching_followups(followups, agent, thread_id, mid, text, dry_run) + maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_run=dry_run) return new_messages, new_wm, job_results @@ -782,6 +799,212 @@ def archive_ephemeral_thread(agent, thread_id, job_id=None): sys.stderr.write(f"warning: archive_ephemeral_thread failed: {ae}\n") +SWARM_WORKER_POOL = ["dev", "def", "muse"] + + +def reconcile_and_dispatch_swarms(dry_run=False): + """ + Autonomous Swarm Orchestrator: + 1. Scans swarms.json for pending slots. + 2. Dynamically allocates available auxiliary worker nodes (dev, def, muse). + 3. Provisions ephemeral sidechat per slot and dispatches the task with [RESULT /]. + 4. Upon completion of all slots, sends completion summary DM to originating coordinator. + """ + if not SWARM_FILE.exists() or dry_run: + return + try: + swarms = json.loads(SWARM_FILE.read_text(encoding="utf-8")) + except Exception: + return + + modified = False + now = utcnow() + + # Determine busy workers from running slots + busy_workers = set() + for sid, swarm in swarms.items(): + if swarm.get("status") in ("pending", "running"): + for slot in swarm.get("slots", []): + if slot.get("status") == "running" and slot.get("agent_id"): + busy_workers.add(slot["agent_id"]) + + # Process each swarm + for sid, swarm in swarms.items(): + st = swarm.get("status") + creator = swarm.get("created_by") or "646" + if creator not in ("646", "pip", "opm", "muse", "dev", "def"): + creator = "646" + + # 1. Allocate & dispatch pending slots + if st in ("pending", "running"): + slots = swarm.get("slots", []) + for s in slots: + if s.get("status") == "pending": + # Find first available worker + worker = None + for w in SWARM_WORKER_POOL: + if w not in busy_workers: + worker = w + break + if not worker: + # Fallback round-robin across worker pool if all are busy + worker = SWARM_WORKER_POOL[s["slot"] % len(SWARM_WORKER_POOL)] + + slot_idx = s["slot"] + task_text = swarm.get("task", "execute subagent task") + slot_target = f"{sid}-s{slot_idx}" + + # Attach worker to slot + s["agent_id"] = worker + s["status"] = "running" + s["updated_ts"] = now + busy_workers.add(worker) + modified = True + + # Prepare task directive prompt + prompt = ( + f"[JOB {sid}/{slot_idx}] Task for swarm slot {slot_idx}:\n" + f"{task_text}\n\n" + f"Reply with [RESULT {sid}/{slot_idx}] OK or FAIL ." + ) + + # Send to worker via dm.py (which handles sidechat creation & tracking) + try: + dm_cmd = [ + sys.executable, str(BIN_DIR / "dm.py"), "send", + "--agent", creator, + "--to", worker, + "--target", slot_target, + prompt, + ] + subprocess.Popen(dm_cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) + print(f"[swarm] Dispatched slot {slot_idx} of {sid} to {worker} in sidechat {slot_target}") + except Exception as de: + sys.stderr.write(f"warning: failed to dispatch swarm slot: {de}\n") + + # Update swarm rollup status + counts = {"pending": 0, "running": 0, "done": 0, "failed": 0, "killed": 0} + for slot in slots: + counts[slot.get("status", "pending")] = counts.get(slot.get("status", "pending"), 0) + 1 + if counts["pending"] + counts["running"] == 0: + swarm["status"] = "completed" if counts["failed"] == 0 else "partial" + swarm["updated_ts"] = now + modified = True + elif counts["running"] or counts["done"]: + swarm["status"] = "running" + swarm["updated_ts"] = now + modified = True + + # 2. Check if newly completed/partial and notify coordinator + if swarm.get("status") in ("completed", "partial") and not swarm.get("notified_coordinator"): + swarm["notified_coordinator"] = True + modified = True + done_cnt = sum(1 for sl in swarm.get("slots", []) if sl.get("status") == "done") + total_cnt = len(swarm.get("slots", [])) + summary_msg = ( + f"[Swarm Report] Swarm {sid} ({swarm.get('label') or 'task'}) finished: " + f"{done_cnt}/{total_cnt} slots successful. " + f"Console: https://box.muse-dev.online/#dashboard" + ) + try: + coord_dm = [ + sys.executable, str(BIN_DIR / "dm.py"), "send", + "--agent", "box", + "--to", creator, + "--target", "main", + "--allow-main-chat", + summary_msg, + ] + subprocess.Popen(coord_dm, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL) + print(f"[swarm] Delivered completion report for {sid} to coordinator {creator}") + except Exception as ce: + sys.stderr.write(f"warning: failed to notify swarm coordinator: {ce}\n") + + if modified: + tmp_swarms = str(SWARM_FILE) + f".tmp.{os.getpid()}" + with open(tmp_swarms, "w", encoding="utf-8") as f: + json.dump(swarms, f, indent=2) + os.replace(tmp_swarms, SWARM_FILE) + + +def check_and_archive_terminal_swarms(): + """Check swarms.json and auto-archive ephemeral sidechats for completed or terminal swarms.""" + if not SWARM_FILE.exists() or not JOB_SIDECHATS_FILE.exists(): + return + try: + swarms_data = json.loads(SWARM_FILE.read_text(encoding="utf-8")) + sc_data = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8")) + except Exception: + return + + for sid, swarm in swarms_data.items(): + st = swarm.get("status") + if st in ("completed", "partial", "killed"): + # Check slots or matching sidechats + for key, val in sc_data.items(): + if isinstance(val, dict) and not val.get("archived") and val.get("type") != "persistent": + if sid in key or (swarm.get("label") and swarm.get("label") in key): + tu = val.get("thread_uuid") + ag = val.get("agent", "opm") + if tu: + archive_ephemeral_thread(ag, tu, job_id=sid) + + + +def maybe_nudge_untagged_sidechat(agent, thread_id, thread_name, mid, text, dry_run=False): + """ + If an agent replies conversationally in a sidechat backed by a job or follow-up + without providing [RESULT ] or tool directives, deliver a terse 1-turn nudge footer. + """ + if dry_run or not thread_id or thread_id == "main": + return + if not re.fullmatch(r"[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}", thread_id.lower()): + return + + # Look up active follow-ups or job association for this thread + followups = load_json_file(FOLLOWUPS_FILE) + matching_job_id = None + for f_id, f_rec in followups.items(): + if f_rec.get("status") in ("pending", "acknowledged") and f_rec.get("thread_uuid") == thread_id: + matching_job_id = f_rec.get("job_id") or f_id + break + + if not matching_job_id and JOB_SIDECHATS_FILE.exists(): + try: + sc_state = json.loads(JOB_SIDECHATS_FILE.read_text(encoding="utf-8")) + for k, v in sc_state.items(): + if isinstance(v, dict) and v.get("thread_uuid") == thread_id: + if k.startswith(("box-", "job-", "pipe-")): + matching_job_id = k + break + except Exception: + pass + + if not matching_job_id: + return + + tracker = load_json_file(NUDGE_TRACKER_FILE) + rec = tracker.get(thread_id, {}) + if rec.get("nudged_for_mid") == mid or rec.get("nudge_count", 0) >= 1: + return + + # Inject terse nudge footer + nudge_msg = f"[Nudge: To advance the workflow, reply with [RESULT {matching_job_id}] ]" + try: + import muse_hybrid + print(f"[{agent}] Injecting 1-turn nudge into {thread_name or thread_id[:8]} for job {matching_job_id}") + muse_hybrid.send_message(agent, nudge_msg, thread_id=thread_id, wait=0) + tracker[thread_id] = { + "ts": utcnow(), + "job_id": matching_job_id, + "nudged_for_mid": mid, + "nudge_count": rec.get("nudge_count", 0) + 1, + } + save_json_file(NUDGE_TRACKER_FILE, tracker) + except Exception as ne: + sys.stderr.write(f"warning: failed to deliver conversational nudge: {ne}\n") + + # Chain deduplication CHAINED_JOBS_FILE = Path(__file__).parent / "chained-jobs.json" @@ -1081,6 +1304,8 @@ def harvest_cycle(target_agent=None, dry_run=False, output_json=False): if not dry_run: save_json_file(WATERMARKS_FILE, watermarks) + reconcile_and_dispatch_swarms() + check_and_archive_terminal_swarms() # Output formatting if output_json: diff --git a/bin/self_main_loop.py b/bin/self_main_loop.py index e7d8e43..51ddba6 100755 --- a/bin/self_main_loop.py +++ b/bin/self_main_loop.py @@ -276,9 +276,14 @@ def make_digest_id(agent): # In-band response-contract footer for ACTIONABLE digests. The verbs are matched # by response-harvester.py to resolve followups: ACK/CLAIM acknowledge (nudge -# suppression), RESULT/DECLINE/NO-ACTION close. ~94 chars, well under budget. +# suppression), RESULT/DECLINE/NO-ACTION close. The closing line enforces the +# recursive box->agent->box discipline: post the RESULT back in this thread +# (box records it and dispatches the next chained step); never DM the next +# agent directly. ~166 chars, within the 600-char digest budget. CONTRACT_FOOTER = ("Reply: [ACK id] seen | [CLAIM id] mine | " - "[RESULT id] done | [DECLINE id] | [NO-ACTION id]") + "[RESULT id] done | [DECLINE id] | [NO-ACTION id]. " + "Report back here. Box dispatches the next step; " + "do not DM the next agent directly.") def send_prompt(sender, agent, sidechat, digest): @@ -552,6 +557,63 @@ def check_subagents(agent, cfg, threads_meta): return alerts + +# --------------------------------------------------------------------------- +# Brain workspace — the main loop operates its thinking in the "main-loop +# brain" sidechat (opm's account). See bin/brain.py. +# User directive 2026-10-04: "we need main loop to operate its brains in +# side chat". Safety: brain posts carry [BRAIN], never [JOB]; never +# --expect-reply; own messages skipped on read; !loop only from authorized +# senders. All enforced in brain.py. +# --------------------------------------------------------------------------- +_BRAIN_LAST_IDS = {} # (agent, sidechat_name) -> last message_id seen + + +def read_brain_messages(agent, sidechat_name, since_ts): + """read_fn adapter for brain.BrainWorkspace. + + Resolves sidechat_name -> UUID via dm (never hardcoded; UUIDs rotate), + reads via the same muse_hybrid primitive the loop uses, returns + [{"sender", "text", "ts"}]. Tracks last message_id per (agent, name); + first run anchors at newest with no backfill (same policy as + new_messages for main chat). + """ + try: + import dm + uuid = dm.resolve_sidechat_target(sidechat_name, agent) + except Exception: + return [] + if not uuid: + return [] + try: + import muse_hybrid + msgs, err = muse_hybrid.get_history(agent, thread_id=uuid, limit=10) + except Exception: + return [] + if err or not msgs: + return [] + key = (agent, sidechat_name) + last_id = _BRAIN_LAST_IDS.get(key) + ids = [m.get("message_id") or "msg-%s" % m.get("seq") for m in msgs] + if last_id is None: + # First run: anchor at newest, no backfill (old !loop commands + # must not fire on deploy). + _BRAIN_LAST_IDS[key] = ids[-1] if ids else None + return [] + if last_id in ids: + new_msgs = msgs[ids.index(last_id) + 1:] + else: + # Watermark fell out of the read window: treat all as new. + # Safe: brain intake skips own messages and only authorized + # !loop senders can act. + new_msgs = msgs + _BRAIN_LAST_IDS[key] = ids[-1] if ids else last_id + now = time.time() + return [{"sender": m.get("role", "unknown"), + "text": m.get("text", ""), + "ts": now} for m in new_msgs] + + def do_check(only_agent=None): st = load_state() cfg = get_config(st) @@ -564,6 +626,23 @@ def do_check(only_agent=None): results = {} prompted = 0 errors = 0 + quiet_held = 0 + + # Brain workspace: the loop thinks in the sidechat, takes !loop + # instructions there. Non-fatal: a brain failure must never break + # the tick. + brain = None + try: + from brain import BrainWorkspace + brain = BrainWorkspace(STATE_FILE, read_fn=read_brain_messages) + intake = brain.intake() + if intake.get("commands"): + log("brain: %d commands, %d acks, %d ignored-senders" % ( + intake["commands"], intake.get("acks", 0), + intake.get("ignored_senders", 0))) + except Exception as e: + log("brain init/intake failed (non-fatal): %r" % e) + brain = None import muse_hybrid @@ -581,6 +660,10 @@ def do_check(only_agent=None): if not enabled.get(agent, True): results[agent] = {"ok": True, "new": 0, "disabled": True} continue + if brain is not None and brain.ignored(agent): + log("%s: skipped (brain !loop ignore active)" % agent) + results[agent] = {"ok": True, "new": 0, "ignored": True} + continue wm = agents_state.get(agent) or {} # 1. Main chat check @@ -640,6 +723,18 @@ def do_check(only_agent=None): errors += 1 continue + # Brain quiet mode: hold actionable escalations. Reads continue, + # digests are composed, but nothing escalates to the agent — the + # thinking note in the brain carries the activity instead. + if brain is not None and brain.quiet(): + is_act, _urg = classify_digest(digest) + if is_act: + log("%s: quiet mode - digest held (not escalated)" % agent) + results[agent] = {"ok": True, "new": total_new, + "quiet_held": True} + quiet_held += 1 + continue + sent, detail, digest_id, actionable = send_prompt(cfg["sender"], agent, sidechat, digest) if sent: prompted += 1 @@ -653,6 +748,49 @@ def do_check(only_agent=None): results[agent] = {"ok": False, "error": detail, "new": total_new} errors += 1 + # Brain: post the tick's thinking to the sidechat workspace, then + # persist brain state. Non-fatal on failure. + if brain is not None: + try: + seen = {} + escalated = [] + for a in agents: + r = results.get(a, {}) + seen[a] = (r.get("new", 0), 0, 0) + if r.get("prompted"): + wm_a = agents_state.get(a) or {} + did = wm_a.get("last_digest_id") + if did and wm_a.get("last_digest_actionable"): + escalated.append(did) + closure_rate = None + health = None + try: + health = digest_health() + closure_rate = (health or {}).get("closure_rate") + except Exception: + pass + wm_epochs = {} + for a in agents: + ts_s = (agents_state.get(a) or {}).get("last_ts") + try: + if ts_s: + dt = datetime.fromisoformat( + ts_s.replace("Z", "+00:00")) + wm_epochs[a] = dt.timestamp() + except Exception: + pass + brain.post_thinking( + {"seen": seen, "escalated": escalated, + "skipped_info": quiet_held, "errors": errors, + "closure_rate": closure_rate}, + cfg={a: bool(enabled.get(a, True)) for a in agents}, + watermarks=wm_epochs, + health=health, + ) + brain.save() + except Exception as e: + log("brain post_thinking failed (non-fatal): %r" % e) + # Save under the state lock with a fresh reload: an enable/disable may # have landed during the slow chat reads; preserve its config changes # and only update the keys this run owns (watermarks, last_run/result).