Files
box/bin/swarm_worker/reporter.py

175 lines
6.1 KiB
Python

#!/usr/bin/env python3
"""Result reporter for swarm workers (Bridge Builder 3/5).
Posts worker results back so the box state (and harvester) picks them up.
Discovered result path (do not guess — verified against live code):
`box swarm-report <swarm-id> <slot>` with JSON on stdin:
{"ok": bool, "result": "<text>"}
Implemented by act_swarm_report() in bin/box-ctl.py. Sets the slot to
"done" (ok=True) or "failed" (ok=False) and records the result text.
The response-harvester ALSO bridges sidechat "[RESULT <swarm_id>/<slot>]"
markers into swarm-report (response-harvester.py ~L773), but calling
swarm-report directly is the reliable path: no harvester-scan dependency,
no dm.py sidechat round-trip. The result text we store carries the
"[RESULT <swarm_id>/<slot>]" marker in the harvester's expected format,
so both paths stay compatible.
Runs on bl; calls bin/box-ctl.py directly (no container relay).
"""
import json
import logging
import subprocess
import sys
log = logging.getLogger("swarm_worker.reporter")
BOX_CTL = "/home/super/Projects/NetVM/bin/box-ctl.py"
SUMMARY_CHARS = 500
def format_result(swarm_id, slot_index, result_dict):
"""Build the result marker text and payload.
Returns (marker_text, payload_dict).
"""
ok = bool(result_dict.get("ok", False))
output = result_dict.get("output") or result_dict.get("error") or ""
summary = str(output)[:SUMMARY_CHARS]
verb = "OK" if ok else "FAIL"
marker = "[RESULT %s/%s] %s %s" % (swarm_id, slot_index, verb, summary)
payload = {"ok": ok, "result": marker}
return marker, payload
def post_result(swarm_id, slot_index, result_dict, dry_run=False):
"""Post a worker result for a swarm slot.
Args:
swarm_id: e.g. "sw-20261005-151307-a222"
slot_index: int slot number
result_dict: {"ok": bool, "output": str} on success, or
{"ok": False, "error": str} on failure.
dry_run: if True, print what WOULD be posted without posting.
Returns:
True on success, False on failure (error is logged, not raised).
"""
marker, payload = format_result(swarm_id, slot_index, result_dict)
if dry_run:
print("DRY-RUN would run: %s swarm-report %s %s"
% (BOX_CTL, swarm_id, slot_index))
print("DRY-RUN stdin JSON: %s" % json.dumps(payload)[:700])
return True
try:
proc = subprocess.run(
[sys.executable, BOX_CTL, "swarm-report",
str(swarm_id), str(slot_index)],
input=json.dumps(payload).encode("utf-8"),
stdout=subprocess.PIPE, stderr=subprocess.PIPE,
timeout=60,
)
except Exception as e: # noqa: BLE001 - report, don't raise
log.error("post_result %s/%s failed to invoke box-ctl: %s",
swarm_id, slot_index, e)
return False
if proc.returncode != 0:
err = proc.stderr.decode("utf-8", "replace")[:500]
# Slot already closed (done/failed/killed) is a known benign case.
log.error("post_result %s/%s box-ctl rc=%d: %s",
swarm_id, slot_index, proc.returncode, err)
return False
try:
resp = json.loads(proc.stdout.decode("utf-8", "replace"))
except ValueError:
log.error("post_result %s/%s: box-ctl returned non-JSON output",
swarm_id, slot_index)
return False
if not resp.get("ok"):
log.error("post_result %s/%s: box-ctl ok=false: %s",
swarm_id, slot_index, str(resp)[:500])
return False
return True
def attach_slot(swarm_id, slot_index, agent_id, session_id=None, dry_run=False):
"""Claim/attach a swarm slot to an agent in box state.
Args:
swarm_id: e.g. "sw-20261005-151307-a222"
slot_index: int slot number
agent_id: agent string, e.g. "dev"
session_id: optional subagent session ID string
dry_run: if True, simulate attach without changing state.
Returns:
True on success, False on failure.
"""
if dry_run:
print("DRY-RUN would run: %s swarm-attach %s %s %s %s"
% (BOX_CTL, swarm_id, slot_index, agent_id, session_id or ""))
return True
cmd = [
sys.executable, BOX_CTL, "swarm-attach",
str(swarm_id), str(slot_index), str(agent_id),
]
if session_id:
cmd.append(str(session_id))
try:
proc = subprocess.run(
cmd,
stdout=subprocess.PIPE, stderr=subprocess.PIPE,
timeout=30,
)
except Exception as e:
log.error("attach_slot %s/%s failed to invoke box-ctl: %s",
swarm_id, slot_index, e)
return False
if proc.returncode != 0:
err = proc.stderr.decode("utf-8", "replace")[:500]
log.error("attach_slot %s/%s box-ctl rc=%d: %s",
swarm_id, slot_index, proc.returncode, err)
return False
try:
resp = json.loads(proc.stdout.decode("utf-8", "replace"))
except ValueError:
log.error("attach_slot %s/%s: box-ctl returned non-JSON output",
swarm_id, slot_index)
return False
if not resp.get("ok"):
log.error("attach_slot %s/%s: box-ctl ok=false: %s",
swarm_id, slot_index, str(resp)[:500])
return False
return True
if __name__ == "__main__":
# Manual dry-run smoke test: python3 reporter.py
# (Never posts real results from this entry point.)
logging.basicConfig(level=logging.DEBUG)
fake_ok = {"ok": True, "output": "did the thing\nline2\n" + "x" * 900}
fake_fail = {"ok": False, "error": "boom: something broke"}
print("== dry-run OK case ==")
assert post_result("sw-TEST-0000-dryrun", 0, fake_ok, dry_run=True)
print("== dry-run FAIL case ==")
assert post_result("sw-TEST-0000-dryrun", 1, fake_fail, dry_run=True)
m, p = format_result("sw-TEST-0000-dryrun", 0, fake_ok)
assert m.startswith("[RESULT sw-TEST-0000-dryrun/0] OK ")
assert len(m) <= len("[RESULT sw-TEST-0000-dryrun/0] OK ") + SUMMARY_CHARS
m2, _ = format_result("sw-TEST-0000-dryrun", 1, fake_fail)
assert m2.startswith("[RESULT sw-TEST-0000-dryrun/1] FAIL boom:")
print("dry-run smoke test PASSED")