59c9965791
- 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
195 lines
7.1 KiB
Python
Executable File
195 lines
7.1 KiB
Python
Executable File
#!/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: <one-line summary of actions and outcome>\n"
|
|
f"(or [RESULT {tag}] FAIL: <reason> 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()
|