From 59c996579151174e127b8a89360f6313d887f8f6 Mon Sep 17 00:00:00 2001 From: operator Date: Mon, 5 Oct 2026 20:18:46 +0000 Subject: [PATCH] feat(swarm-worker): dispatch prose swarm slots to ephemeral Muse subagents - daemon: route non-shell slot tasks to a fresh subagent session on dev/def/muse via muse_hybrid; add bin/ to sys.path so muse_hybrid/prompt_envelope import under the tmux supervisor - prompt_envelope: add wrap_subagent_task() - plain task assignment without [TOOL tmux/swarm/cron] meta tags (subagents refused those as relayed test traffic) - box-ctl/reporter: swarm-attach accepts optional session-id, stored as slot.subagent_session_id - executor: _looks_like_shell requires an existing executable (prose like 'verify ...' no longer misread as shell) - poller: pick up pending swarms as well as running - response-harvester: monitor running slot subagent sessions from swarms.json, allow '/' in RESULT/VERB job ids, archive ephemeral threads on any verdict (OK or FAIL), reap slot sessions for terminal swarms --- bin/box-ctl.py | 18 ++++---- bin/prompt_envelope.py | 15 ++++++ bin/response-harvester.py | 73 ++++++++++++++++++++++++----- bin/swarm_worker/daemon.py | 89 +++++++++++++++++++++++++++++++++--- bin/swarm_worker/executor.py | 12 ++++- bin/swarm_worker/poller.py | 3 +- bin/swarm_worker/reporter.py | 54 ++++++++++++++++++++++ 7 files changed, 235 insertions(+), 29 deletions(-) diff --git a/bin/box-ctl.py b/bin/box-ctl.py index 7336968..9a437ff 100755 --- a/bin/box-ctl.py +++ b/bin/box-ctl.py @@ -2004,7 +2004,7 @@ def act_swarm_status(sid): out(True, swarm=swarm, counts=_swarm_counts(swarm)) -def act_swarm_attach(sid, slot_str, agent_id): +def act_swarm_attach(sid, slot_str, agent_id, session_id=None): 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}) @@ -2021,13 +2021,15 @@ def act_swarm_attach(sid, slot_str, agent_id): {"swarm_id": sid, "slot": slot_rec["slot"]}) slot_rec["agent_id"] = agent_id slot_rec["status"] = "running" + if session_id: + slot_rec["subagent_session_id"] = str(session_id) 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") + subagent_session_id=session_id, status="running") def act_swarm_report(sid, slot_str): @@ -3432,9 +3434,9 @@ def main(argv): 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]) + if len(rest) not in (3, 4): + fail("BAD_ARGS", "usage: swarm-attach [session-id]") + act_swarm_attach(rest[0], rest[1], rest[2], rest[3] if len(rest) > 3 else None) elif op == "report": if len(rest) != 2: fail("BAD_ARGS", "usage: swarm-report (result JSON on stdin)") @@ -3479,9 +3481,9 @@ def main(argv): 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]) + if len(args) not in (3, 4): + fail("BAD_ARGS", "usage: swarm attach [session-id]") + act_swarm_attach(args[0], args[1], args[2], args[3] if len(args) > 3 else None) elif sub == "report": if len(args) != 2: fail("BAD_ARGS", "usage: swarm report (result JSON on stdin)") diff --git a/bin/prompt_envelope.py b/bin/prompt_envelope.py index 328ea86..7ecbfff 100644 --- a/bin/prompt_envelope.py +++ b/bin/prompt_envelope.py @@ -119,3 +119,18 @@ def wrap(job_name, job_id, agent, target, rendered): bottom += f"Conclude with your [RESULT {job_id}] line reporting outcomes.\n" return top + rendered.strip() + "\n" + bottom + +def wrap_subagent_task(job_id, task_text): + """Format an authentic direct task assignment for a subagent worker without test-traffic meta tags.""" + return ( + f"Operator assignment for swarm slot {job_id}:\n\n" + f"Task:\n" + f"{task_text.strip()}\n\n" + f"Instructions:\n" + f"1. Carry out this task directly using your available tools.\n" + f"2. When finished, conclude your final response with your verdict line:\n" + f"[RESULT {job_id}] OK: \n" + f"(or [RESULT {job_id}] FAIL: if the task could not be completed)\n" + ) + + diff --git a/bin/response-harvester.py b/bin/response-harvester.py index dca6640..3cb25e1 100755 --- a/bin/response-harvester.py +++ b/bin/response-harvester.py @@ -83,18 +83,32 @@ try: except ImportError: HAS_MUSE_HYBRID = False +try: + import prompt_envelope + HAS_PROMPT_ENVELOPE = True +except ImportError: + HAS_PROMPT_ENVELOPE = False + VALID_AGENTS = ["muse", "pip", "646", "opm", "dev", "def"] DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440, "def": 9450, "dev": 9455} +try: + import lookup_engine + HAS_LOOKUP_ENGINE = True +except ImportError: + HAS_LOOKUP_ENGINE = False + # Matches EVERY [RESULT ] marker in a message (use with finditer, not # search). The result text is lazy and stops before the next marker (or end of # text), so a message closing two jobs records each with its own text instead # of the first marker greedily swallowing the second. -RESULT_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*?)(?=\[RESULT\s|\Z)", re.S) +_LOCAL_RESULT_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_/-]+)\]\s*(.*?)(?=\[RESULT\s|\Z)", re.S) +RESULT_RE = lookup_engine.get_result_regex() if HAS_LOOKUP_ENGINE else _LOCAL_RESULT_RE def iter_result_markers(text): """Yield (job_id, result_text) for every [RESULT ] marker in text.""" - for m in RESULT_RE.finditer(text or ""): + current_re = lookup_engine.get_result_regex() if HAS_LOOKUP_ENGINE else RESULT_RE + for m in current_re.finditer(text or ""): yield m.group(1).strip(), m.group(2).strip() @@ -103,12 +117,14 @@ def iter_result_markers(text): # ACK/CLAIM acknowledge a digest (nudge-suppressed, NOT closed); # RESULT/DECLINE/NO-ACTION close the digest. Every verb match records # outcome= on the followup record. -VERB_RE = re.compile(r"\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_-]+)\]") +_LOCAL_VERB_RE = re.compile(r"\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_/-]+)\]") +VERB_RE = lookup_engine.get_verb_regex() if HAS_LOOKUP_ENGINE else _LOCAL_VERB_RE def iter_verb_markers(text): """Yield (verb, job_id) for every [VERB ] marker in text.""" - for m in VERB_RE.finditer(text or ""): + current_re = lookup_engine.get_verb_regex() if HAS_LOOKUP_ENGINE else VERB_RE + for m in current_re.finditer(text or ""): yield m.group(1), m.group(2).strip() @@ -503,6 +519,26 @@ def get_monitored_threads(target_agent=None): except Exception: continue + # Also monitor active swarm slot subagents from swarms.json + try: + if SWARM_FILE.exists(): + s_data = json.loads(SWARM_FILE.read_text(encoding="utf-8")) + for sid, s_info in s_data.items(): + if sid.startswith("_") or not isinstance(s_info, dict): + continue + if s_info.get("status") not in ("running", "pending"): + continue + for slot in s_info.get("slots", []): + if slot.get("status") == "running": + sub_id = slot.get("subagent_session_id") or slot.get("thread_uuid") or slot.get("sidechat_id") + ag = slot.get("agent_id") + if sub_id and ag and ag in threads_by_agent: + existing = [t["id"] for t in threads_by_agent[ag]] + if sub_id not in existing: + threads_by_agent[ag].append({"id": sub_id, "name": f"swarm-{sid[:8]}-s{slot.get('slot')}"}) + except Exception as se: + sys.stderr.write(f"warning: failed to add swarm threads to monitor list: {se}\n") + return threads_by_agent @@ -783,8 +819,7 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo 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) + archive_ephemeral_thread(agent, thread_id, job_id=job_id) clear_matching_followups(followups, agent, thread_id, mid, text, dry_run, job_id=job_id, verb="RESULT") for verb, job_id in verbs: @@ -1047,11 +1082,19 @@ def reconcile_and_dispatch_swarms(dry_run=False): slot_target = f"{sid}-s{slot_idx}" # 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 ." - ) + slot_job_id = f"{sid}/{slot_idx}" + if HAS_PROMPT_ENVELOPE and hasattr(prompt_envelope, "wrap_subagent_task"): + prompt = prompt_envelope.wrap_subagent_task(slot_job_id, task_text) + else: + prompt = ( + f"Operator assignment for swarm slot {slot_job_id}:\n\n" + f"Task:\n{task_text.strip()}\n\n" + f"Instructions:\n" + f"1. Carry out this task directly using your available tools.\n" + f"2. When finished, conclude your final response with your verdict line:\n" + f"[RESULT {slot_job_id}] OK: \n" + f"(or [RESULT {slot_job_id}] FAIL: if the task could not be completed)\n" + ) # Native Subagent Execution Bridge: # Spawns an interactive child session directly on the target worker @@ -1196,7 +1239,13 @@ def check_and_archive_terminal_swarms(): for sid, swarm in swarms_data.items(): st = swarm.get("status") if st in ("completed", "partial", "killed"): - # Check slots or matching sidechats + # Check slot subagent sessions + for slot in swarm.get("slots", []): + sub_id = slot.get("subagent_session_id") + ag = slot.get("agent_id") + if sub_id and ag: + archive_ephemeral_thread(ag, sub_id, job_id=sid) + # Check matching registered 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): diff --git a/bin/swarm_worker/daemon.py b/bin/swarm_worker/daemon.py index 7bb8aaf..4762ab6 100755 --- a/bin/swarm_worker/daemon.py +++ b/bin/swarm_worker/daemon.py @@ -23,15 +23,33 @@ import sys import time import traceback -# Sibling modules live in the same directory. -sys.path.insert(0, os.path.dirname(os.path.abspath(__file__))) +# Sibling modules live in swarm_worker, common utilities live in bin/. +_SWARM_DIR = os.path.dirname(os.path.abspath(__file__)) +_BIN_DIR = os.path.dirname(_SWARM_DIR) +sys.path.insert(0, _SWARM_DIR) +sys.path.insert(0, _BIN_DIR) from poller import find_pending_slots -from executor import execute_task -from reporter import post_result +from executor import execute_task, _looks_like_shell +from reporter import post_result, attach_slot + +# Fast gateway integration +try: + import muse_hybrid + HAS_MUSE_HYBRID = True +except ImportError: + HAS_MUSE_HYBRID = False + +# Prompt envelope formatting +try: + import prompt_envelope + HAS_PROMPT_ENVELOPE = True +except ImportError: + HAS_PROMPT_ENVELOPE = False POLL_INTERVAL = 60 # seconds between poll cycles STALE_MINUTES = 5 # slots older than this with no attach are workable +WORKER_POOL = ["dev", "def", "muse"] # === SAFETY SWITCH === # True -> observe only: log what WOULD be done, execute/post nothing. @@ -46,19 +64,78 @@ logging.basicConfig( log = logging.getLogger("swarm-worker") +def _select_worker(preferred=None): + if preferred and HAS_MUSE_HYBRID and muse_hybrid.is_node_configured(preferred): + return preferred + for candidate in WORKER_POOL: + if HAS_MUSE_HYBRID and muse_hybrid.is_node_configured(candidate): + return candidate + return preferred or "dev" + + def process_slot(slot): """Execute one slot and report the result. Returns True on full success.""" swarm_id = slot.get("swarm_id") slot_index = slot.get("slot_index") task_text = slot.get("task_text") or "" + agent_label = slot.get("agent_label") + sidechat_id = slot.get("sidechat_id") tag = "%s/%s" % (swarm_id, slot_index) if DRY_RUN: log.info("[dry-run] would execute slot %s (agent=%s, task %.80r)", - tag, slot.get("agent_label"), task_text) + tag, agent_label, task_text) return True - log.info("executing slot %s (agent=%s)", tag, slot.get("agent_label")) + # If the task is NOT a shell command, dispatch it to an ephemeral Muse subagent. + is_shell = _looks_like_shell(task_text) + if not is_shell and HAS_MUSE_HYBRID: + worker_agent = _select_worker(agent_label) + log.info("dispatching subagent slot %s to %s", tag, worker_agent) + try: + # 1. Start an ephemeral subagent session + title = f"sw-{swarm_id[:16]}-s{slot_index}" + sess, err = muse_hybrid.start_session(worker_agent, title=title) + if err or not sess or not sess.get("session_id"): + log.error("failed to start subagent session for %s on %s: %s", tag, worker_agent, err) + return False + + sub_sid = sess["session_id"] + log.info("subagent session %s created for slot %s on %s", sub_sid, tag, worker_agent) + + # 2. Attach/claim the slot in box state with subagent session_id + attached = attach_slot(swarm_id, slot_index, worker_agent, session_id=sub_sid) + if not attached: + log.warning("failed to attach slot %s to %s; proceeding with dispatch", tag, worker_agent) + + # 3. Format prompt with authentic Operator Directive and RESULT expectation + if HAS_PROMPT_ENVELOPE and hasattr(prompt_envelope, "wrap_subagent_task"): + prompt_body = prompt_envelope.wrap_subagent_task(tag, task_text) + else: + prompt_body = ( + f"Operator assignment for swarm slot {tag}:\n\n" + f"Task:\n{task_text.strip()}\n\n" + f"Instructions:\n" + f"1. Carry out this task directly using your available tools.\n" + f"2. When finished, conclude your final response with your verdict line:\n" + f"[RESULT {tag}] OK: \n" + f"(or [RESULT {tag}] FAIL: if the task could not be completed)\n" + ) + + # 4. Asynchronously send message to subagent session + res, send_err = muse_hybrid.send_message(worker_agent, prompt_body, thread_id=sub_sid, wait=0) + if send_err: + log.error("failed to send task to subagent %s on %s: %s", sub_sid, worker_agent, send_err) + return False + + log.info("slot %s successfully dispatched to subagent %s (harvester will harvest)", tag, sub_sid) + return True + except Exception: + log.error("subagent dispatch crashed on slot %s:\n%s", tag, traceback.format_exc()) + return False + + # Otherwise fallback to sandboxed host execution + log.info("executing slot %s in sandbox (agent=%s)", tag, agent_label) try: result = execute_task(task_text) except Exception: diff --git a/bin/swarm_worker/executor.py b/bin/swarm_worker/executor.py index fa9ca8d..2f45ca1 100644 --- a/bin/swarm_worker/executor.py +++ b/bin/swarm_worker/executor.py @@ -77,11 +77,19 @@ def _looks_like_shell(task_text): return False if _PROSE_RE.search(t): return False - # Must start with a word-ish token (not a quote or sentence). + # Must start with a plausible executable command or path first = t.split()[0] if t.split() else "" if not re.match(r"^[a-zA-Z0-9_.\-/]+$", first): return False - return True + # If it's a relative/absolute path, verify it exists and is executable + if "/" in first: + return os.path.isfile(first) and os.access(first, os.X_OK) + # If it's a bare command name, it must exist in standard system bin paths + for p in ("/bin", "/usr/bin", "/usr/local/bin"): + candidate = os.path.join(p, first) + if os.path.isfile(candidate) and os.access(candidate, os.X_OK): + return True + return False def _refused(task_text): diff --git a/bin/swarm_worker/poller.py b/bin/swarm_worker/poller.py index 0181de5..daf0811 100644 --- a/bin/swarm_worker/poller.py +++ b/bin/swarm_worker/poller.py @@ -97,7 +97,8 @@ def find_pending_slots(stale_minutes=STALE_MINUTES): return found swarms = data.get("swarms", []) if isinstance(data, dict) else [] for summary in swarms: - if (summary.get("status") or "").lower() != "running": + summary_status = (summary.get("status") or "").lower() + if summary_status not in ("running", "pending"): continue swarm_id = summary.get("swarm_id") if not swarm_id: diff --git a/bin/swarm_worker/reporter.py b/bin/swarm_worker/reporter.py index 2ed57d6..75a7bb6 100644 --- a/bin/swarm_worker/reporter.py +++ b/bin/swarm_worker/reporter.py @@ -97,6 +97,60 @@ def post_result(swarm_id, slot_index, result_dict, dry_run=False): swarm_id, slot_index, str(resp)[:500]) return False +def attach_slot(swarm_id, slot_index, agent_id, session_id=None, dry_run=False): + """Claim/attach a swarm slot to an agent in box state. + + Args: + swarm_id: e.g. "sw-20261005-151307-a222" + slot_index: int slot number + agent_id: agent string, e.g. "dev" + session_id: optional subagent session ID string + dry_run: if True, simulate attach without changing state. + + Returns: + True on success, False on failure. + """ + if dry_run: + print("DRY-RUN would run: %s swarm-attach %s %s %s %s" + % (BOX_CTL, swarm_id, slot_index, agent_id, session_id or "")) + return True + + cmd = [ + sys.executable, BOX_CTL, "swarm-attach", + str(swarm_id), str(slot_index), str(agent_id), + ] + if session_id: + cmd.append(str(session_id)) + + try: + proc = subprocess.run( + cmd, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, + timeout=30, + ) + except Exception as e: + log.error("attach_slot %s/%s failed to invoke box-ctl: %s", + swarm_id, slot_index, e) + return False + + if proc.returncode != 0: + err = proc.stderr.decode("utf-8", "replace")[:500] + log.error("attach_slot %s/%s box-ctl rc=%d: %s", + swarm_id, slot_index, proc.returncode, err) + return False + + try: + resp = json.loads(proc.stdout.decode("utf-8", "replace")) + except ValueError: + log.error("attach_slot %s/%s: box-ctl returned non-JSON output", + swarm_id, slot_index) + return False + + if not resp.get("ok"): + log.error("attach_slot %s/%s: box-ctl ok=false: %s", + swarm_id, slot_index, str(resp)[:500]) + return False + return True