diff --git a/bin/swarm-worker-supervise.sh b/bin/swarm-worker-supervise.sh new file mode 100755 index 0000000..b3b4845 --- /dev/null +++ b/bin/swarm-worker-supervise.sh @@ -0,0 +1,80 @@ +#!/usr/bin/env bash +# +# swarm-worker-supervise.sh — tmux supervisor for the swarm worker daemon. +# +# Manages a detached tmux session named "swarm-worker" that runs the worker +# main loop and logs to logs/swarm-worker.log. +# +# The worker itself is a PLACEHOLDER until bin/swarm_worker/daemon.py lands. +# When it does, update WORKER_CMD and WORKER_MATCH below. +# +# Usage: swarm-worker-supervise.sh {start|stop|restart|status} +# +set -euo pipefail + +SESSION="swarm-worker" +LOG_DIR="/home/super/Projects/NetVM/logs" +LOG_FILE="${LOG_DIR}/swarm-worker.log" + +# Placeholder worker. Replace with: +# WORKER_CMD="python3 /home/super/Projects/NetVM/bin/swarm_worker/daemon.py" +# when the real daemon exists. WORKER_MATCH must uniquely match the worker +# process command line; the [a] bracket trick keeps pgrep from matching itself. +WORKER_CMD="python3 /home/super/Projects/NetVM/bin/swarm_worker/daemon.py" +WORKER_MATCH='swarm_worker/daemon.[p]y' + +session_exists() { + tmux has-session -t "$SESSION" 2>/dev/null +} + +worker_alive() { + pgrep -f "$WORKER_MATCH" >/dev/null 2>&1 +} + +do_start() { + mkdir -p "$LOG_DIR" + if session_exists; then + echo "already running (tmux session $SESSION exists)" + return 0 + fi + tmux new-session -d -s "$SESSION" "$WORKER_CMD >>\"$LOG_FILE\" 2>&1" + sleep 1 + if worker_alive; then + echo "started tmux session $SESSION (logging to $LOG_FILE)" + else + echo "WARNING: session $SESSION created but worker process not detected yet" + fi +} + +do_stop() { + if session_exists; then + tmux kill-session -t "$SESSION" + echo "stopped tmux session $SESSION" + else + echo "not running (no tmux session $SESSION)" + fi +} + +do_status() { + if ! session_exists; then + echo "STOPPED - no tmux session '$SESSION'" + return 1 + fi + if worker_alive; then + echo "ALIVE - tmux session '$SESSION' exists and worker process is running" + return 0 + fi + echo "DEGRADED - tmux session '$SESSION' exists but no worker process found" + return 2 +} + +case "${1:-}" in + start) do_start ;; + stop) do_stop ;; + restart) do_stop; sleep 1; do_start ;; + status) do_status ;; + *) + echo "Usage: $(basename "$0") {start|stop|restart|status}" >&2 + exit 2 + ;; +esac diff --git a/bin/swarm_worker/daemon.py b/bin/swarm_worker/daemon.py new file mode 100755 index 0000000..7bb8aaf --- /dev/null +++ b/bin/swarm_worker/daemon.py @@ -0,0 +1,117 @@ +#!/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() diff --git a/bin/swarm_worker/executor.py b/bin/swarm_worker/executor.py new file mode 100644 index 0000000..fa9ca8d --- /dev/null +++ b/bin/swarm_worker/executor.py @@ -0,0 +1,194 @@ +"""Swarm worker task executor (Bridge Builder 2/5). + +Executes a swarm slot's task text in a sandbox and captures the result. + +Sandbox rules (hard): +- No network calls except to localhost. +- No writes outside /tmp. +- No sudo / privilege escalation. +- No credential access (no reading ~/.ssh, ~/.aws, key stores, etc). +- Strict timeout; kill the process group on exceed. +- Refuse tasks that look destructive without executing anything. +""" + +import os +import re +import shlex +import signal +import subprocess +import time + +OUTPUT_TRUNCATE = 2000 + +# Patterns that indicate a potentially destructive command. Matched against +# the raw task text (case-insensitive) before any execution is attempted. +_REFUSE_PATTERNS = [ + r"\brm\s+-rf?\b", # rm -r / rm -rf + r"\brm\s+.*\s+/\s*$", # rm ... / (trailing root) + r"\bmkfs\b", + r"\bdd\b.*\bof=/dev/", + r"\bdd\b.*\bof=\s*/dev/", + r":\(\)\s*\{", # fork bomb + r"\bshutdown\b", + r"\breboot\b", + r"\bpoweroff\b", + r"\bhalt\b", + r"\binit\s+[06]\b", + r"\bsudo\b", + r"\bsu\b", + r"\bchmod\s+-R\s+777\s+/\b", + r"\bchown\s+-R\b.*\s+/\s*$", + r"\bmv\s+.*\s+/\s*$", + r"\bwipefs\b", + r"\bshred\b.*\s/dev/", + r">\s*/dev/sd", + r"\bcurl\b.*\|\s*(ba)?sh", # curl|sh pipe-to-shell + r"\bwget\b.*\|\s*(ba)?sh", + r"\bnc\b.*-e\s", # netcat reverse shell + r"\bpython\w*\s+-c\b.*socket", +] + +_REFUSE_RE = re.compile("|".join("(?:%s)" % p for p in _REFUSE_PATTERNS), + re.IGNORECASE) + +# Paths that must never be read (credential / identity material). +_FORBIDDEN_READ_PREFIXES = ( + os.path.expanduser("~/.ssh"), + os.path.expanduser("~/.aws"), + os.path.expanduser("~/.gnupg"), + "/etc/shadow", + "/etc/netvm", + "/etc/ssl/private", +) + +# A task is treated as a shell command when it is short, single-purpose, +# and does not look like prose/instructions. Heuristic: one or two lines, +# starts with a plausible command token, no sentence-like structure. +_PROSE_RE = re.compile(r"[.?!]\s+[A-Z]|\n\n|please\s|you\s+are\s|your\s+task", + re.IGNORECASE) + + +def _looks_like_shell(task_text): + """Heuristic: is this plausibly a shell command?""" + t = task_text.strip() + if not t: + return False + if len(t.splitlines()) > 3: + return False + if _PROSE_RE.search(t): + return False + # Must start with a word-ish token (not a quote or sentence). + first = t.split()[0] if t.split() else "" + if not re.match(r"^[a-zA-Z0-9_.\-/]+$", first): + return False + return True + + +def _refused(task_text): + return bool(_REFUSE_RE.search(task_text)) + + +def _sandbox_env(): + """Minimal environment: strip anything credential-shaped.""" + env = { + "PATH": "/usr/bin:/bin", + "HOME": "/tmp", + "TMPDIR": "/tmp", + "LANG": "C.UTF-8", + } + return env + + +def execute_task(task_text, timeout=600): + """Execute a swarm slot's task text in a sandbox. + + Args: + task_text: the task string from the swarm slot. + timeout: max seconds for execution (default 600). Strictly enforced. + + Returns: + dict(success=bool, output=str, duration_s=float, error=str|None) + """ + started = time.monotonic() + text = (task_text or "").strip() + + def done(success, output, error=None): + dur = round(time.monotonic() - started, 3) + out = (output or "")[:OUTPUT_TRUNCATE] + return { + "success": success, + "output": out, + "duration_s": dur, + "error": error, + } + + if not text: + return done(False, "", error="empty task") + + if _refused(text): + return done(False, "", error="refused: potentially destructive") + + if not _looks_like_shell(text): + return done( + True, + "received but not executable: task is prose/instructions, " + "not a shell command; recorded %d chars" % len(text), + error=None, + ) + + # Sandbox: run in /tmp, fresh process group so timeout kills children too. + try: + proc = subprocess.Popen( + text, + shell=True, + executable="/bin/bash", + cwd="/tmp", + env=_sandbox_env(), + stdout=subprocess.PIPE, + stderr=subprocess.STDOUT, + text=True, + start_new_session=True, # new process group for reliable kill + ) + except Exception as e: # e.g. /bin/bash missing + return done(False, "", error="spawn failed: %s" % e) + + try: + stdout, _ = proc.communicate(timeout=timeout) + rc = proc.returncode + except subprocess.TimeoutExpired: + # Kill the whole process group. + try: + os.killpg(os.getpgid(proc.pid), signal.SIGKILL) + except ProcessLookupError: + pass + try: + stdout, _ = proc.communicate(timeout=5) + except Exception: + stdout = "" + return done( + False, + stdout or "", + error="timeout: exceeded %ss, process group killed" % timeout, + ) + + output = stdout or "" + if rc == 0: + return done(True, output) + return done(False, output, error="exit code %d" % rc) + + +if __name__ == "__main__": + import json + import sys + + cases = [ + ("echo hello", 10), + ("rm -rf /", 10), + ] + # Allow ad-hoc: python executor.py "" [timeout] + if len(sys.argv) > 1: + cases = [(sys.argv[1], int(sys.argv[2]) if len(sys.argv) > 2 else 600)] + for task, to in cases: + print("TASK:", task) + print(json.dumps(execute_task(task, timeout=to), indent=1)) + print("-" * 40) diff --git a/bin/swarm_worker/mainloop-bridge.md b/bin/swarm_worker/mainloop-bridge.md new file mode 100644 index 0000000..72d15de --- /dev/null +++ b/bin/swarm_worker/mainloop-bridge.md @@ -0,0 +1,95 @@ +# Swarm Worker ↔ Main Loop Bridge — Design Doc +**Bridge Builder 5/5 | 2026-10-05 | Read-only design (no changes to self_main_loop.py)** + +## Context + +Sibling builders produced a tmux-hosted worker daemon for bl: +- `executor.py` — sandboxed task execution (shell vs prose heuristic, refuse-list) +- `poller.py` — slot claiming +- `reporter.py` — result recording via `box swarm-report ` (direct, reliable) + +The user's directive: *"we should be answering our own main chat via our main loop."* +This doc designs how swarm work becomes visible through `self_main_loop.py` +without modifying it. + +## How the main loop works (verified against live code) + +`bin/self_main_loop.py` (995 lines) is a timer-driven read-side loop: +1. Reads each watched agent's muse.ai **main chat** + registered **sidechats** + (via `muse_hybrid` gateway, `box-chat.py` fallback). +2. On new activity, composes a digest and posts it to the agent's **prompting + sidechat** via `dm.py send` (`send_prompt()`). +3. Actionable digests (containing `(?)` or `(!)`) get `[JOB ml--]` + + contract footer → `--expect-reply` → followup record. Informational + digests go via fast gateway, no followup. +4. Response verbs `[ACK|CLAIM|RESULT|DECLINE|NO-ACTION ]` are matched by + `response-harvester.py` to resolve followups. + +Key detail — `check_sidechats()` **skips individual messages** containing +`"[JOB "`, `"Heartbeat check"`, or `"[from:super]"`. It does NOT skip entire +sidechats. A `[RESULT ...]` or `[SWARM-DONE ...]` post is *not* filtered. + +## How swarm slot DMs arrive (verified against dm-log.jsonl) + +Schema: `{type, id, agent, to, target, msg, tags, ts}`. Lifecycle: +`send_start` → `sidechat_autoprovisioned` → `verified` → `sent`. + +- Dispatch: `dm.py send --agent --to --target sw--s` + with body `[JOB sw-/] Task for swarm slot N:\n\n\nReply with [RESULT sw-...]` +- Target auto-provisions a sidechat; 23 slot sidechats are registered in + `job-sidechats.json` with per-slot `agent` attribution (e.g. `agent: "dev"`). +- Harvester (`response-harvester.py:773`) bridges sidechat + `[RESULT /]` markers into `box swarm-report` — but the + sibling's `reporter.py` calls `swarm-report` directly (no DM round-trip). + +## Recommendation: (a) direct sidechat post, main loop as visibility layer + +**Result recording (unchanged):** worker → `box swarm-report` directly +(sibling's reporter.py). Reliable, no harvester-scan dependency. + +**Main-loop bridge (new, zero code changes):** after recording the result, +the worker posts a short human-readable completion note to the slot's own +sidechat via `dm.py send` — deliberately *without* a `[JOB ` tag so the +main loop's `check_sidechats()` does not filter it: + +``` +[SWARM-DONE sw-20261005-151307-a222/0] OK in 12.4s: +``` + +Why this works with no modifications: +1. Slot sidechats are already registered in `job-sidechats.json` with agent + attribution → `get_monitored_sidechats()` includes them for that agent. +2. The note contains no `[JOB ` → not skipped by the message filter. +3. Next main-loop tick picks it up → digest → agent's prompting sidechat → + the agent *sees* the completion in the same surface as everything else. + This is "answering via the main loop": completions surface through the + existing digest machinery instead of a parallel notification channel. + +**Why not (b) queue via `send_prompt()`:** `send_prompt()` is built for +main-chat digests → prompting sidechats, with `[JOB ml-...]` followup +semantics. Swarm results are not digests; routing them through would add +timer latency, misuse the ml- followup namespace, and conflate two concerns. + +**Why not (c) worker as main-loop participant:** the main loop monitors +muse.ai *agent accounts* (`muse_hybrid.get_history(agent, ...)`). The worker +is a bl daemon with no agent account. Participation-by-proxy (writing to +monitored sidechats) achieves the same visibility with no identity hacks. + +## Watchdog (no change needed, documented here) + +`check_sidechats()` also gives stall detection for free: a claimed slot +whose sidechat shows no `[SWARM-DONE]`/`[RESULT]` within the slot timeout +is visible as inactivity in the digest. A future enhancement (not this +bridge) could add an explicit stall-threshold alert. + +## Interface for the worker daemon + +After `reporter.post_result(...)` succeeds, call: +`notify_via_mainloop(swarm_id, slot_index, summary_text)` (prototype: +`bin/swarm_worker/mainloop_notify.py`). It resolves +`sw--s` → sidechat UUID via `dm.resolve_sidechat_target`, +then `dm.py send` with the `[SWARM-DONE ...]` note. Dry-run mode posts nothing. + +## Files +- This doc: `bin/swarm_worker/mainloop-bridge.md` +- Prototype: `bin/swarm_worker/mainloop_notify.py` (dry-run safe, never auto-wired) diff --git a/bin/swarm_worker/mainloop_notify.py b/bin/swarm_worker/mainloop_notify.py new file mode 100755 index 0000000..5fd5293 --- /dev/null +++ b/bin/swarm_worker/mainloop_notify.py @@ -0,0 +1,87 @@ +#!/usr/bin/env python3 +"""Swarm worker -> main-loop visibility bridge (Bridge Builder 5/5). + +After the worker records a slot result via `box swarm-report` (see +reporter.py), this posts a short human-readable completion note to the +slot's own sidechat. The note deliberately carries NO `[JOB ` tag, so +self_main_loop.py's check_sidechats() does not filter it: the next +main-loop tick picks it up, includes it in the agent's digest, and the +completion becomes visible in the same surface as everything else. +That is the "answering via the main loop" path. + +Design doc: bin/swarm_worker/mainloop-bridge.md + +Runs on bl. Never auto-wired into any loop; call explicitly from the +worker daemon after post_result() succeeds. Dry-run mode posts nothing. +""" + +import subprocess +import sys + +BIN = "/home/super/Projects/NetVM/bin" +DM_PY = BIN + "/dm.py" +NOTE_CHARS = 200 + + +def notify_via_mainloop(swarm_id, slot_index, message, worker_id="swarm-worker", + dry_run=False): + """Post a swarm-slot completion note visible to the main loop. + + Args: + swarm_id: e.g. "sw-20261005-151307-a222" + slot_index: int slot number + message: human-readable summary (first NOTE_CHARS chars used) + worker_id: sender identity for dm.py --agent + dry_run: if True, print what WOULD be sent without sending. + + Returns: + True on success (or dry-run), False on failure (logged, not raised). + """ + target_name = "sw-%s-s%s" % (swarm_id, slot_index) + summary = (message or "").strip().replace("\n", " ")[:NOTE_CHARS] + note = "[SWARM-DONE %s/%s] %s" % (swarm_id, slot_index, summary) + + try: + sys.path.insert(0, BIN) + import dm + uuid = dm.resolve_sidechat_target(target_name) + except Exception as e: + print("notify_via_mainloop: target resolve failed for %s: %s" + % (target_name, e), file=sys.stderr) + return False + if not uuid: + print("notify_via_mainloop: no UUID for %s" % target_name, + file=sys.stderr) + return False + + cmd = [sys.executable, DM_PY, "send", + "--agent", worker_id, + "--to", worker_id, + "--target", uuid, + note] + if dry_run: + print("DRY-RUN would run: %s" % " ".join(cmd)) + return True + try: + r = subprocess.run(cmd, capture_output=True, text=True, timeout=120) + except Exception as e: + print("notify_via_mainloop: dm.py exec failed: %s" % e, + file=sys.stderr) + return False + if r.returncode != 0: + tail = ((r.stderr or "") + (r.stdout or "")).strip()[-300:] + print("notify_via_mainloop: dm.py rc=%d: %s" % (r.returncode, tail), + file=sys.stderr) + return False + return True + + +if __name__ == "__main__": + live = "--live" in sys.argv + args = [a for a in sys.argv[1:] if a != "--live"] + if len(args) < 3: + print("usage: mainloop_notify.py \"\" [--live]") + print("default is DRY-RUN (posts nothing).") + sys.exit(2) + ok = notify_via_mainloop(args[0], args[1], args[2], dry_run=not live) + print("notified" if ok else "FAILED") diff --git a/bin/swarm_worker/poller.py b/bin/swarm_worker/poller.py new file mode 100644 index 0000000..69c6cb4 --- /dev/null +++ b/bin/swarm_worker/poller.py @@ -0,0 +1,138 @@ +#!/usr/bin/env python3 +"""Slot poller for the swarm worker daemon. + +READ-ONLY component: finds swarm slots waiting for a worker. It never +claims slots, never modifies box state, never posts anything. + +Usage: + python3 poller.py # prints pending slots as JSON to stdout + from poller import find_pending_slots + slots = find_pending_slots() + +Each returned dict has keys: + swarm_id, slot_index, agent_label, task_text, sidechat_id, created_ts +""" + +import json +import subprocess +import sys +from datetime import datetime, timezone + +# How old (minutes) a running, result-less, unattached slot must be +# before we consider it abandoned and re-workable. +STALE_MINUTES = 5 + +# Resolved once at import: bl-native (`box swarm list`) vs legacy (`box swarm-list`). +_BOX_FORM = None # "native" | "legacy" + + +def _run_box(*args, timeout=60): + """Run a box swarm subcommand, return parsed JSON or None on any failure.""" + global _BOX_FORM + forms = [] + if _BOX_FORM == "native" or _BOX_FORM is None: + forms.append(["box", "swarm", *args, "--json"]) + if _BOX_FORM == "legacy" or _BOX_FORM is None: + # legacy relay style: `box swarm-list`, `box swarm-status ` + legacy = "-".join(["swarm"] + list(args[:1])) + forms.append(["box", legacy, *args[1:]]) + last_err = None + for cmd in forms: + try: + proc = subprocess.run( + cmd, capture_output=True, text=True, timeout=timeout + ) + except (FileNotFoundError, subprocess.TimeoutExpired) as e: + last_err = f"{e}" + continue + if proc.returncode != 0: + last_err = f"rc={proc.returncode}: {proc.stderr.strip()[:200]}" + # "invalid choice" means wrong form; try the next one. + if "invalid choice" in (proc.stderr + proc.stdout): + continue + print(f"[poller] {' '.join(cmd)} failed: {last_err}", + file=sys.stderr) + return None + try: + data = json.loads(proc.stdout) + except json.JSONDecodeError as e: + print(f"[poller] bad JSON from {' '.join(cmd)}: {e}", + file=sys.stderr) + return None + _BOX_FORM = "native" if cmd[1] == "swarm" else "legacy" + return data + print(f"[poller] all box forms failed: {last_err}", file=sys.stderr) + return None + + +def _age_minutes(ts): + """Minutes since an ISO-8601 timestamp; inf when unparseable/missing.""" + if not ts: + return float("inf") + try: + dt = datetime.fromisoformat(str(ts).replace("Z", "+00:00")) + if dt.tzinfo is None: + dt = dt.replace(tzinfo=timezone.utc) + return (datetime.now(timezone.utc) - dt).total_seconds() / 60.0 + except Exception: + return float("inf") + + +def _slot_attached(slot): + """True when some worker identity is recorded on the slot.""" + return bool(slot.get("subagent_session_id") or slot.get("claimed_by")) + + +def find_pending_slots(stale_minutes=STALE_MINUTES): + """Return slots waiting for a worker. Read-only; claims nothing. + + A slot qualifies when: + - its status is "pending", or it has no agent assigned, OR + - it is "running" with result null, older than `stale_minutes`, + and no worker identity is attached. + """ + found = [] + data = _run_box("list") + if not data: + return found + swarms = data.get("swarms", []) if isinstance(data, dict) else [] + for summary in swarms: + if (summary.get("status") or "").lower() != "running": + continue + swarm_id = summary.get("swarm_id") + if not swarm_id: + continue + detail = _run_box("status", swarm_id) + if not detail: + continue + swarm = detail.get("swarm", {}) if isinstance(detail, dict) else {} + task_text = swarm.get("task") or summary.get("task_preview") or "" + created_ts = swarm.get("created_ts") or summary.get("created_ts") + for slot in swarm.get("slots", []): + status = (slot.get("status") or "").lower() + result = slot.get("result") + agent = slot.get("agent_id") + pending = (status == "pending") or (not agent) + stale_running = ( + status == "running" + and result is None + and _age_minutes(slot.get("updated_ts")) > stale_minutes + and not _slot_attached(slot) + ) + if pending or stale_running: + found.append({ + "swarm_id": swarm_id, + "slot_index": slot.get("slot"), + "agent_label": agent, + "task_text": task_text, + "sidechat_id": slot.get("sidechat_id") + or slot.get("thread_uuid"), + "created_ts": created_ts, + }) + return found + + +if __name__ == "__main__": + slots = find_pending_slots() + print(json.dumps(slots, indent=1)) + print(f"[poller] {len(slots)} pending slot(s) found", file=sys.stderr) diff --git a/bin/swarm_worker/reporter.py b/bin/swarm_worker/reporter.py new file mode 100644 index 0000000..2ed57d6 --- /dev/null +++ b/bin/swarm_worker/reporter.py @@ -0,0 +1,118 @@ +#!/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 ` with JSON on stdin: + {"ok": bool, "result": ""} + 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 /]" + 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 /]" 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 + + +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")