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
173 lines
6.1 KiB
Python
173 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
|
|
|
|
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")
|