Files
box/bin/swarm_worker/daemon.py

245 lines
9.3 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, extract_shell_command, execute_task_in_tmux
from reporter import post_result, attach_slot
from mainloop_notify import notify_via_mainloop
# 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 = ["muse"] # only dispatch to fully authenticated agent nodes
# === 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 preferred in WORKER_POOL 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 "muse"
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)
short_id = swarm_id[3:19] if str(swarm_id).startswith("sw-") else str(swarm_id)[:16]
session_name = f"sw-{short_id}-s{slot_index}"
if DRY_RUN:
log.info("[dry-run] would execute slot %s (agent=%s, task %.80r)",
tag, agent_label, task_text)
return True
# 1. Check if the task is or contains an executable shell command
cmd = extract_shell_command(task_text)
if cmd:
log.info("executing slot %s in host tmux session %s on bl", tag, session_name)
# Attach/claim slot in box state
attach_slot(swarm_id, slot_index, "swarm-worker", session_id=session_name)
try:
result = execute_task_in_tmux(session_name, cmd)
except Exception:
log.error("tmux executor crashed on slot %s:\n%s", tag, traceback.format_exc())
result = {"success": False, "output": "",
"error": "tmux 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
# Post completion note to the slot sidechat for main loop visibility
summary_msg = payload.get("output") or payload.get("error") or "completed"
try:
notified = notify_via_mainloop(swarm_id, slot_index, summary_msg, worker_id="swarm-worker")
log.info("slot %s sidechat notification: %s", tag, notified)
except Exception as ne:
log.warning("failed to post sidechat notification for %s: %s", tag, ne)
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
# 2. If the task is purely prose/instructions, dispatch to a verified agent subagent
if HAS_MUSE_HYBRID:
worker_agent = _select_worker(agent_label)
log.info("dispatching prose subagent slot %s to %s", tag, worker_agent)
try:
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)
# Attach/claim slot in box state
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)
# Register in subagent_tracker
try:
import subagent_tracker
subagent_tracker.register_session(worker_agent, sub_sid, title=title, prompt=task_text[:200])
except Exception:
pass
# 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"
)
# 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
# 3. Fallback to sandboxed host execution
log.info("executing slot %s in fallback 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
try:
notify_via_mainloop(swarm_id, slot_index, payload.get("output") or "done", worker_id="swarm-worker")
except Exception:
pass
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()