#!/usr/bin/env python3 """Swarm worker daemon (bridge integrator). Polls box for pending swarm slots, executes them, and posts results back. Runs inside the `swarm-worker` tmux session via swarm-worker-supervise.sh. Loop: 1. find_pending_slots() (poller.py) 2. execute_task(task_text) (executor.py) 3. post_result(swarm_id, slot, r) (reporter.py) 4. sleep POLL_INTERVAL, repeat Safety: - DRY_RUN=True means: never execute, never post. Log "would execute" and skip. Flip to False only on explicit go-live authorization. - Per-slot try/except: one bad slot never kills the loop. - All activity logged to stdout (captured to the tmux log file). """ import logging import os import sys import time import traceback # 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, _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. # False -> live mode. Flip only on explicit go-live authorization. DRY_RUN = False logging.basicConfig( level=logging.INFO, format="%(asctime)s [swarm-worker] %(levelname)s %(message)s", datefmt="%Y-%m-%dT%H:%M:%S", ) 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, agent_label, task_text) return True # 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: log.error("executor crashed on slot %s:\n%s", tag, traceback.format_exc()) result = {"success": False, "output": "", "error": "executor crashed: see worker log"} payload = { "ok": bool(result.get("success")), "output": result.get("output") or "", "error": result.get("error"), } try: posted = post_result(swarm_id, slot_index, payload) except Exception: log.error("reporter crashed on slot %s:\n%s", tag, traceback.format_exc()) posted = False log.info("slot %s done: ok=%s posted=%s (%.1fs)", tag, payload["ok"], posted, float(result.get("duration_s") or 0.0)) return bool(payload["ok"]) and posted def poll_cycle(): """One poll cycle. Never raises.""" try: slots = find_pending_slots(stale_minutes=STALE_MINUTES) except Exception: log.error("poller crashed:\n%s", traceback.format_exc()) return if not slots: log.info("poll: no pending slots") return log.info("poll: %d pending slot(s)", len(slots)) for slot in slots: try: process_slot(slot) except Exception: log.error("slot %s/%s unhandled error:\n%s", slot.get("swarm_id"), slot.get("slot_index"), traceback.format_exc()) def main(): log.info("starting (DRY_RUN=%s, interval=%ds)", DRY_RUN, POLL_INTERVAL) while True: try: poll_cycle() except Exception: log.error("poll cycle crashed:\n%s", traceback.format_exc()) time.sleep(POLL_INTERVAL) if __name__ == "__main__": main()