feat(exec): add template aliases, tolerant read-only validation, and thread titling filter
This commit is contained in:
+65
-3
@@ -1053,6 +1053,60 @@ OPS = {
|
|||||||
# identity -> set of ops. 'master' may invoke everything. Unknown identities
|
# identity -> set of ops. 'master' may invoke everything. Unknown identities
|
||||||
# get the read-only subset. Per-agent tokens inherit their agent name as the
|
# get the read-only subset. Per-agent tokens inherit their agent name as the
|
||||||
# identity; tighten per agent here as needed.
|
# 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 = {
|
PERMISSIONS = {
|
||||||
'master': set(OPS),
|
'master': set(OPS),
|
||||||
'operator-main': set(OPS),
|
'operator-main': set(OPS),
|
||||||
@@ -1211,12 +1265,17 @@ class Handler(BaseHTTPRequestHandler):
|
|||||||
'op': op})
|
'op': op})
|
||||||
self._send(403, {'error': 'forbidden for this identity'})
|
self._send(403, {'error': 'forbidden for this identity'})
|
||||||
return
|
return
|
||||||
|
dropped = []
|
||||||
try:
|
try:
|
||||||
spec = OPS[op]
|
spec = OPS[op]
|
||||||
|
if not spec.get('side_effecting'):
|
||||||
|
args, dropped = _tolerant_args(spec, args)
|
||||||
clean = spec['validate'](args)
|
clean = spec['validate'](args)
|
||||||
argv = spec['build'](clean)
|
argv = spec['build'](clean)
|
||||||
except OpError as e:
|
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
|
return
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
sys.stderr.write(f'Validation/build exception for {op}: {e}\n')
|
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,
|
audit({'event': 'exec_done', 'ident': ident, 'peer': peer,
|
||||||
'op': op, 'rc': rc, 'duration_s': dt, 'ok': ok})
|
'op': op, 'rc': rc, 'duration_s': dt, 'ok': ok})
|
||||||
if ok:
|
if ok:
|
||||||
self._send(200, {'rc': rc, 'stdout': out, 'stderr': err,
|
resp = {'rc': rc, 'stdout': out, 'stderr': err, 'duration_s': dt}
|
||||||
'duration_s': dt})
|
if dropped:
|
||||||
|
resp['dropped_args'] = dropped
|
||||||
|
self._send(200, resp)
|
||||||
else:
|
else:
|
||||||
self._send(500, {'error': err, 'rc': rc})
|
self._send(500, {'error': err, 'rc': rc})
|
||||||
|
|
||||||
|
|
||||||
|
|
||||||
# ------------------------------------------------------- main
|
# ------------------------------------------------------- main
|
||||||
|
|
||||||
def main():
|
def main():
|
||||||
|
|||||||
+112
-12
@@ -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
|
#!/usr/bin/env python3
|
||||||
"""
|
"""
|
||||||
response-harvester.py — Fleet agent readback and response harvesting daemon.
|
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")
|
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):
|
def reconcile_and_dispatch_swarms(dry_run=False):
|
||||||
@@ -943,6 +971,36 @@ def reconcile_and_dispatch_swarms(dry_run=False):
|
|||||||
modified = False
|
modified = False
|
||||||
now = utcnow()
|
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
|
# Determine busy workers from running slots
|
||||||
busy_workers = set()
|
busy_workers = set()
|
||||||
for sid, swarm in swarms.items():
|
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")
|
task_text = swarm.get("task", "execute subagent task")
|
||||||
slot_target = f"{sid}-s{slot_idx}"
|
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
|
# Prepare task directive prompt
|
||||||
prompt = (
|
prompt = (
|
||||||
f"[JOB {sid}/{slot_idx}] Task for swarm slot {slot_idx}:\n"
|
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 <summary> or FAIL <reason>."
|
f"Reply with [RESULT {sid}/{slot_idx}] OK <summary> or FAIL <reason>."
|
||||||
)
|
)
|
||||||
|
|
||||||
# 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:
|
try:
|
||||||
dm_cmd = [
|
dm_cmd = [
|
||||||
sys.executable, str(BIN_DIR / "dm.py"), "send",
|
sys.executable, str(BIN_DIR / "dm.py"), "send",
|
||||||
@@ -1000,10 +1057,53 @@ def reconcile_and_dispatch_swarms(dry_run=False):
|
|||||||
"--target", slot_target,
|
"--target", slot_target,
|
||||||
prompt,
|
prompt,
|
||||||
]
|
]
|
||||||
subprocess.Popen(dm_cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
|
proc = subprocess.run(
|
||||||
print(f"[swarm] Dispatched slot {slot_idx} of {sid} to {worker} in sidechat {slot_target}")
|
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:
|
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
|
# Update swarm rollup status
|
||||||
counts = {"pending": 0, "running": 0, "done": 0, "failed": 0, "killed": 0}
|
counts = {"pending": 0, "running": 0, "done": 0, "failed": 0, "killed": 0}
|
||||||
|
|||||||
+11
-2
@@ -59,7 +59,7 @@ PIPELINES_FILE = NETVM_ROOT / "pipelines.json"
|
|||||||
JOB_DISPATCH_PY = BIN_DIR / "job-dispatch.py"
|
JOB_DISPATCH_PY = BIN_DIR / "job-dispatch.py"
|
||||||
|
|
||||||
# Agent Constants
|
# Agent Constants
|
||||||
VALID_NODES = ["muse", "pip", "646", "opm"]
|
VALID_NODES = ["muse", "pip", "646", "opm", "def", "dev"]
|
||||||
DEFAULT_SENDER = "super"
|
DEFAULT_SENDER = "super"
|
||||||
DEFAULT_AGENT_SIDECHATS = {
|
DEFAULT_AGENT_SIDECHATS = {
|
||||||
"646": "646 tasks",
|
"646": "646 tasks",
|
||||||
@@ -177,6 +177,8 @@ def get_node_network_info(node: str) -> dict:
|
|||||||
"pip": 9420,
|
"pip": 9420,
|
||||||
"646": 9430,
|
"646": 9430,
|
||||||
"opm": 9440,
|
"opm": 9440,
|
||||||
|
"def": 9450,
|
||||||
|
"dev": 9460,
|
||||||
}
|
}
|
||||||
cdp_port = pinned_ports.get(node, 9222 + int(tag[4:7], 16) % 2000)
|
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
|
agent = args.agent
|
||||||
threads = []
|
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)
|
# Fast path: try fast headless gateway via muse_hybrid (isolated per-node WARP egress)
|
||||||
try:
|
try:
|
||||||
import muse_hybrid
|
import muse_hybrid
|
||||||
gw_threads, gw_err = muse_hybrid.get_threads(agent)
|
gw_threads, gw_err = muse_hybrid.get_threads(agent)
|
||||||
if gw_threads and not gw_err:
|
if gw_threads and not gw_err:
|
||||||
for t in gw_threads:
|
for t in gw_threads:
|
||||||
|
if is_noise(t.get("title")):
|
||||||
|
continue
|
||||||
threads.append({
|
threads.append({
|
||||||
"id": t.get("session_id", ""),
|
"id": t.get("session_id", ""),
|
||||||
"kind": "thread" if t.get("thread") else "chat",
|
"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)
|
res = subprocess.run(cmd, capture_output=True, text=True)
|
||||||
try:
|
try:
|
||||||
data = json.loads(res.stdout)
|
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:
|
except Exception:
|
||||||
threads = []
|
threads = []
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user