diff --git a/bin/exec-constrained.py b/bin/exec-constrained.py index f23f62f..8450355 100755 --- a/bin/exec-constrained.py +++ b/bin/exec-constrained.py @@ -1053,6 +1053,60 @@ OPS = { # identity -> set of ops. 'master' may invoke everything. Unknown identities # get the read-only subset. Per-agent tokens inherit their agent name as the # identity; tighten per agent here as needed. +# ---- read-only tolerance + hyphenated aliases used by job templates ---- +_ARG_SYNONYMS = {'job': 'name', 'id': 'name', 'job_name': 'name', 'unit': 'name'} + + +def _tolerant_args(spec, args): + """Read-only ops: drop unknown args (mapping common synonyms) instead of failing. + Returns (args, dropped_names).""" + if args is None: + return {}, [] + if not isinstance(args, dict): + return args, [] + args, dropped = dict(args), [] + for _ in range(len(args) + 1): + try: + spec['validate'](args) + break + except OpError as e: + m = re.match(r'unknown arg: (\S+)', str(e)) + if m and m.group(1) in args: + k = m.group(1) + v = args.pop(k) + syn = _ARG_SYNONYMS.get(k) + if syn and syn not in args: + args[syn] = v + else: + dropped.append(k) + continue + if 'takes no' in str(e): + dropped.extend(args) + args = {} + continue + break + return args, dropped + + +def _register_aliases(): + alias = {'job-list': 'cron.runs', 'quality.validate': 'quality.check', + 'vars-list': 'vars.list', 'fleet-status': 'health.check'} + for new, old in alias.items(): + if new not in OPS and old in OPS: + OPS[new] = dict(OPS[old], desc=OPS[old]['desc'] + f' (alias of {old})') + def _ctl(*words): + return lambda a: [sys.executable, os.path.join(BIN_DIR, 'box-ctl.py'), *words] + def _noargs(raw): + return {} + for new, words, desc in (('loop-status', ('loop-status',), 'Fleet loop health (read-only)'), + ('timer-list', ('timer', 'list'), 'List box timers (read-only)')): + if new not in OPS: + OPS[new] = {'validate': _noargs, 'build': _ctl(*words), 'timeout': 30, + 'side_effecting': False, 'desc': desc} + + +_register_aliases() + PERMISSIONS = { 'master': set(OPS), 'operator-main': set(OPS), @@ -1211,12 +1265,17 @@ class Handler(BaseHTTPRequestHandler): 'op': op}) self._send(403, {'error': 'forbidden for this identity'}) return + dropped = [] try: spec = OPS[op] + if not spec.get('side_effecting'): + args, dropped = _tolerant_args(spec, args) clean = spec['validate'](args) argv = spec['build'](clean) except OpError as e: - self._send(400, {'error': f'bad args: {e}'}) + self._send(400, {'error': f'bad args: {e}', 'op': op, + 'desc': OPS[op].get('desc', ''), + 'hint': 'GET /ops for the allowlist; see desc for usage'}) return except Exception as e: sys.stderr.write(f'Validation/build exception for {op}: {e}\n') @@ -1247,12 +1306,15 @@ class Handler(BaseHTTPRequestHandler): audit({'event': 'exec_done', 'ident': ident, 'peer': peer, 'op': op, 'rc': rc, 'duration_s': dt, 'ok': ok}) if ok: - self._send(200, {'rc': rc, 'stdout': out, 'stderr': err, - 'duration_s': dt}) + resp = {'rc': rc, 'stdout': out, 'stderr': err, 'duration_s': dt} + if dropped: + resp['dropped_args'] = dropped + self._send(200, resp) else: self._send(500, {'error': err, 'rc': rc}) + # ------------------------------------------------------- main def main(): diff --git a/bin/response-harvester.py b/bin/response-harvester.py index 8f4bcac..cc2ca98 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -1,3 +1,7 @@ +def is_title_noise(title): + t = (title or "").lower().strip() + return bool(re.search(r"generate.*(chat|session)?.*title", t)) + #!/usr/bin/env python3 """ response-harvester.py — Fleet agent readback and response harvesting daemon. @@ -922,7 +926,31 @@ 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"] +# NOTE 2026-10-05 (Fix Agent 2/5): dev/def removed from the pool. No worker +# agents exist on dev/def (no Meta sessions provisioned), so slots dispatched +# to them froze with null results. Re-add only after real dev/def workers exist. +SWARM_WORKER_POOL = ["muse"] + +# Stuck-slot reaper: a slot that stays "running" with no result longer than +# this is treated as wedged (worker died / dispatch lost). Healthy slots +# complete in <5 min (observed p90 3.6 min over 26 done slots, 2026-10-05), +# so 60 min is conservative. +STUCK_SLOT_MINUTES = 60 + +# Verified dispatch: dm.py send runs synchronously and the slot is only +# marked running when the send is confirmed (rc 0 + "SENT" in stdout). +# Unverified sends stay pending for retry; after N attempts the slot fails +# loudly instead of freezing from birth. +DISPATCH_VERIFY_TIMEOUT = 120 +DISPATCH_MAX_ATTEMPTS = 5 + + +def _parse_ts(ts): + try: + return datetime.fromisoformat(str(ts).replace("Z", "+00:00")) + except Exception: + return None + def reconcile_and_dispatch_swarms(dry_run=False): @@ -943,6 +971,36 @@ def reconcile_and_dispatch_swarms(dry_run=False): modified = False now = utcnow() + # 0. Reap stuck slots BEFORE computing busy workers. A wedged worker + # would otherwise pin itself "busy" forever and the pool stalls: with + # all workers busy the fallback round-robin keeps feeding new slots to + # the same dead workers. + now_dt = _parse_ts(now) + if now_dt is not None: + for _sid, _swarm in swarms.items(): + if _swarm.get("status") not in ("pending", "running"): + continue + for _s in _swarm.get("slots", []): + if _s.get("status") != "running" or _s.get("result") is not None: + continue + _upd = _parse_ts(_s.get("updated_ts", "")) + if _upd is None: + continue + _age_min = (now_dt - _upd).total_seconds() / 60 + if _age_min > STUCK_SLOT_MINUTES: + _s["status"] = "failed" + _s["result"] = { + "ok": False, + "reaped": True, + "reason": "stuck: running with no result for %.0f min (limit %d)" + % (_age_min, STUCK_SLOT_MINUTES), + "worker": _s.get("agent_id"), + } + _s["updated_ts"] = now + modified = True + print("[swarm] Reaped stuck slot %d of %s (worker %s, silent %.0f min)" + % (_s["slot"], _sid, _s.get("agent_id"), _age_min)) + # Determine busy workers from running slots busy_workers = set() for sid, swarm in swarms.items(): @@ -977,13 +1035,6 @@ def reconcile_and_dispatch_swarms(dry_run=False): 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" @@ -991,7 +1042,13 @@ def reconcile_and_dispatch_swarms(dry_run=False): f"Reply with [RESULT {sid}/{slot_idx}] OK or FAIL ." ) - # Send to worker via dm.py (which handles sidechat creation & tracking) + # Verified dispatch: confirm the send landed BEFORE + # marking the slot running. Fire-and-forget used to + # freeze slots from birth -- the slot read "running" + # while no worker was ever notified (2026-10-05: dev's + # wedged WireGuard data path silently killed every + # gateway dispatch). + dispatch_ok = False try: dm_cmd = [ sys.executable, str(BIN_DIR / "dm.py"), "send", @@ -1000,10 +1057,53 @@ def reconcile_and_dispatch_swarms(dry_run=False): "--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}") + proc = subprocess.run( + dm_cmd, stdout=subprocess.PIPE, stderr=subprocess.PIPE, + text=True, timeout=DISPATCH_VERIFY_TIMEOUT) + _dout = proc.stdout or "" + dispatch_ok = (proc.returncode == 0 and "SENT" in _dout) + if not dispatch_ok: + sys.stderr.write( + "warning: swarm dispatch unverified %s slot %d -> %s " + "(rc=%d out=%.100s err=%.100s)\n" + % (sid, slot_idx, worker, proc.returncode, + _dout, proc.stderr or "")) except Exception as de: - sys.stderr.write(f"warning: failed to dispatch swarm slot: {de}\n") + sys.stderr.write( + "warning: swarm dispatch error %s slot %d -> %s: %s\n" + % (sid, slot_idx, worker, de)) + + if not dispatch_ok: + # Leave pending so the next harvester pass retries + # (possibly with a different worker). Fail loudly + # after N attempts instead of spinning forever. + attempts = s.get("dispatch_attempts", 0) + 1 + s["dispatch_attempts"] = attempts + modified = True + if attempts >= DISPATCH_MAX_ATTEMPTS: + s["status"] = "failed" + s["result"] = { + "ok": False, + "reaped": True, + "reason": "dispatch failed %d times (send never verified)" % attempts, + "worker": worker, + } + s["updated_ts"] = now + print("[swarm] Dispatch failed %dx for slot %d of %s; marked failed" + % (attempts, slot_idx, sid)) + else: + print("[swarm] Dispatch unverified for slot %d of %s " + "(attempt %d/%d); leaving pending for retry" + % (slot_idx, sid, attempts, DISPATCH_MAX_ATTEMPTS)) + continue + + # Attach worker to slot (only after verified dispatch) + s["agent_id"] = worker + s["status"] = "running" + s["updated_ts"] = now + busy_workers.add(worker) + modified = True + print(f"[swarm] Dispatched slot {slot_idx} of {sid} to {worker} in sidechat {slot_target}") # Update swarm rollup status counts = {"pending": 0, "running": 0, "done": 0, "failed": 0, "killed": 0} diff --git a/bin/super-cli.py b/bin/super-cli.py index efa31fe..aea8f37 100755 --- a/bin/super-cli.py +++ b/bin/super-cli.py @@ -59,7 +59,7 @@ PIPELINES_FILE = NETVM_ROOT / "pipelines.json" JOB_DISPATCH_PY = BIN_DIR / "job-dispatch.py" # Agent Constants -VALID_NODES = ["muse", "pip", "646", "opm"] +VALID_NODES = ["muse", "pip", "646", "opm", "def", "dev"] DEFAULT_SENDER = "super" DEFAULT_AGENT_SIDECHATS = { "646": "646 tasks", @@ -177,6 +177,8 @@ def get_node_network_info(node: str) -> dict: "pip": 9420, "646": 9430, "opm": 9440, + "def": 9450, + "dev": 9460, } cdp_port = pinned_ports.get(node, 9222 + int(tag[4:7], 16) % 2000) @@ -1176,12 +1178,18 @@ def cmd_thread_list(args): agent = args.agent threads = [] + def is_noise(title): + t = (title or "").lower().strip() + return bool(re.search(r"generate.*(chat|session)?.*title", t)) + # Fast path: try fast headless gateway via muse_hybrid (isolated per-node WARP egress) try: import muse_hybrid gw_threads, gw_err = muse_hybrid.get_threads(agent) if gw_threads and not gw_err: for t in gw_threads: + if is_noise(t.get("title")): + continue threads.append({ "id": t.get("session_id", ""), "kind": "thread" if t.get("thread") else "chat", @@ -1198,7 +1206,8 @@ def cmd_thread_list(args): res = subprocess.run(cmd, capture_output=True, text=True) try: data = json.loads(res.stdout) - threads = data.get("threads", []) if data.get("ok") else [] + raw_threads = data.get("threads", []) if data.get("ok") else [] + threads = [t for t in raw_threads if not is_noise(t.get("title"))] except Exception: threads = []