feat(swarm): launch live swarm worker daemon under tmux supervisor
This commit is contained in:
Executable
+80
@@ -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
|
||||
Executable
+117
@@ -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()
|
||||
@@ -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 "<task>" [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)
|
||||
@@ -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 <id> <slot>` (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-<agent>-<ts>]`
|
||||
+ contract footer → `--expect-reply` → followup record. Informational
|
||||
digests go via fast gateway, no followup.
|
||||
4. Response verbs `[ACK|CLAIM|RESULT|DECLINE|NO-ACTION <id>]` 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 <dispatcher> --to <worker> --target sw-<id>-s<slot>`
|
||||
with body `[JOB sw-<id>/<slot>] Task for swarm slot N:\n<task>\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 <swarm_id>/<slot>]` 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: <first ~200 chars of output>
|
||||
```
|
||||
|
||||
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-<swarm_id>-s<slot_index>` → 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)
|
||||
Executable
+87
@@ -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 <swarm-id> <slot> \"<message>\" [--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")
|
||||
@@ -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 <id>`
|
||||
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)
|
||||
@@ -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 <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
|
||||
|
||||
|
||||
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")
|
||||
Reference in New Issue
Block a user