118 lines
3.7 KiB
Python
118 lines
3.7 KiB
Python
|
|
#!/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 the same directory.
|
||
|
|
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
||
|
|
|
||
|
|
from poller import find_pending_slots
|
||
|
|
from executor import execute_task
|
||
|
|
from reporter import post_result
|
||
|
|
|
||
|
|
POLL_INTERVAL = 60 # seconds between poll cycles
|
||
|
|
STALE_MINUTES = 5 # slots older than this with no attach are workable
|
||
|
|
|
||
|
|
# === 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 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 ""
|
||
|
|
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)
|
||
|
|
return True
|
||
|
|
|
||
|
|
log.info("executing slot %s (agent=%s)", tag, slot.get("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()
|