245 lines
9.3 KiB
Python
Executable File
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()
|