feat(swarm-worker): dispatch prose swarm slots to ephemeral Muse subagents
- 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
This commit is contained in:
+10
-8
@@ -2004,7 +2004,7 @@ def act_swarm_status(sid):
|
|||||||
out(True, swarm=swarm, counts=_swarm_counts(swarm))
|
out(True, swarm=swarm, counts=_swarm_counts(swarm))
|
||||||
|
|
||||||
|
|
||||||
def act_swarm_attach(sid, slot_str, agent_id):
|
def act_swarm_attach(sid, slot_str, agent_id, session_id=None):
|
||||||
if not agent_id or not SWARM_AGENT_RE.match(agent_id):
|
if not agent_id or not SWARM_AGENT_RE.match(agent_id):
|
||||||
fail("BAD_NAME", "agent id must match ^[A-Za-z0-9-]{1,64}$",
|
fail("BAD_NAME", "agent id must match ^[A-Za-z0-9-]{1,64}$",
|
||||||
{"field": "agent_id", "value": agent_id})
|
{"field": "agent_id", "value": agent_id})
|
||||||
@@ -2021,13 +2021,15 @@ def act_swarm_attach(sid, slot_str, agent_id):
|
|||||||
{"swarm_id": sid, "slot": slot_rec["slot"]})
|
{"swarm_id": sid, "slot": slot_rec["slot"]})
|
||||||
slot_rec["agent_id"] = agent_id
|
slot_rec["agent_id"] = agent_id
|
||||||
slot_rec["status"] = "running"
|
slot_rec["status"] = "running"
|
||||||
|
if session_id:
|
||||||
|
slot_rec["subagent_session_id"] = str(session_id)
|
||||||
slot_rec["updated_ts"] = utcnow()
|
slot_rec["updated_ts"] = utcnow()
|
||||||
swarm["status"] = _swarm_rollup(swarm)
|
swarm["status"] = _swarm_rollup(swarm)
|
||||||
swarm["updated_ts"] = utcnow()
|
swarm["updated_ts"] = utcnow()
|
||||||
_swarm_save(swarms)
|
_swarm_save(swarms)
|
||||||
audit("swarm-attach", "%s/%d" % (sid, slot_rec["slot"]))
|
audit("swarm-attach", "%s/%d" % (sid, slot_rec["slot"]))
|
||||||
out(True, swarm_id=sid, slot=slot_rec["slot"], agent_id=agent_id,
|
out(True, swarm_id=sid, slot=slot_rec["slot"], agent_id=agent_id,
|
||||||
status="running")
|
subagent_session_id=session_id, status="running")
|
||||||
|
|
||||||
|
|
||||||
def act_swarm_report(sid, slot_str):
|
def act_swarm_report(sid, slot_str):
|
||||||
@@ -3432,9 +3434,9 @@ def main(argv):
|
|||||||
fail("BAD_ARGS", "usage: swarm-status <swarm-id>")
|
fail("BAD_ARGS", "usage: swarm-status <swarm-id>")
|
||||||
act_swarm_status(rest[0])
|
act_swarm_status(rest[0])
|
||||||
elif op == "attach":
|
elif op == "attach":
|
||||||
if len(rest) != 3:
|
if len(rest) not in (3, 4):
|
||||||
fail("BAD_ARGS", "usage: swarm-attach <swarm-id> <slot> <agent-id>")
|
fail("BAD_ARGS", "usage: swarm-attach <swarm-id> <slot> <agent-id> [session-id]")
|
||||||
act_swarm_attach(rest[0], rest[1], rest[2])
|
act_swarm_attach(rest[0], rest[1], rest[2], rest[3] if len(rest) > 3 else None)
|
||||||
elif op == "report":
|
elif op == "report":
|
||||||
if len(rest) != 2:
|
if len(rest) != 2:
|
||||||
fail("BAD_ARGS", "usage: swarm-report <swarm-id> <slot> (result JSON on stdin)")
|
fail("BAD_ARGS", "usage: swarm-report <swarm-id> <slot> (result JSON on stdin)")
|
||||||
@@ -3479,9 +3481,9 @@ def main(argv):
|
|||||||
fail("BAD_ARGS", "usage: swarm status <swarm-id>")
|
fail("BAD_ARGS", "usage: swarm status <swarm-id>")
|
||||||
act_swarm_status(args[0])
|
act_swarm_status(args[0])
|
||||||
elif sub == "attach":
|
elif sub == "attach":
|
||||||
if len(args) != 3:
|
if len(args) not in (3, 4):
|
||||||
fail("BAD_ARGS", "usage: swarm attach <swarm-id> <slot> <agent-id>")
|
fail("BAD_ARGS", "usage: swarm attach <swarm-id> <slot> <agent-id> [session-id]")
|
||||||
act_swarm_attach(args[0], args[1], args[2])
|
act_swarm_attach(args[0], args[1], args[2], args[3] if len(args) > 3 else None)
|
||||||
elif sub == "report":
|
elif sub == "report":
|
||||||
if len(args) != 2:
|
if len(args) != 2:
|
||||||
fail("BAD_ARGS", "usage: swarm report <swarm-id> <slot> (result JSON on stdin)")
|
fail("BAD_ARGS", "usage: swarm report <swarm-id> <slot> (result JSON on stdin)")
|
||||||
|
|||||||
@@ -119,3 +119,18 @@ def wrap(job_name, job_id, agent, target, rendered):
|
|||||||
bottom += f"Conclude with your [RESULT {job_id}] line reporting outcomes.\n"
|
bottom += f"Conclude with your [RESULT {job_id}] line reporting outcomes.\n"
|
||||||
return top + rendered.strip() + "\n" + bottom
|
return top + rendered.strip() + "\n" + bottom
|
||||||
|
|
||||||
|
|
||||||
|
def wrap_subagent_task(job_id, task_text):
|
||||||
|
"""Format an authentic direct task assignment for a subagent worker without test-traffic meta tags."""
|
||||||
|
return (
|
||||||
|
f"Operator assignment for swarm slot {job_id}:\n\n"
|
||||||
|
f"Task:\n"
|
||||||
|
f"{task_text.strip()}\n\n"
|
||||||
|
f"Instructions:\n"
|
||||||
|
f"1. Carry out this task directly using your available tools.\n"
|
||||||
|
f"2. When finished, conclude your final response with your verdict line:\n"
|
||||||
|
f"[RESULT {job_id}] OK: <one-line summary of actions and outcome>\n"
|
||||||
|
f"(or [RESULT {job_id}] FAIL: <reason> if the task could not be completed)\n"
|
||||||
|
)
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+61
-12
@@ -83,18 +83,32 @@ try:
|
|||||||
except ImportError:
|
except ImportError:
|
||||||
HAS_MUSE_HYBRID = False
|
HAS_MUSE_HYBRID = False
|
||||||
|
|
||||||
|
try:
|
||||||
|
import prompt_envelope
|
||||||
|
HAS_PROMPT_ENVELOPE = True
|
||||||
|
except ImportError:
|
||||||
|
HAS_PROMPT_ENVELOPE = False
|
||||||
|
|
||||||
VALID_AGENTS = ["muse", "pip", "646", "opm", "dev", "def"]
|
VALID_AGENTS = ["muse", "pip", "646", "opm", "dev", "def"]
|
||||||
DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440, "def": 9450, "dev": 9455}
|
DEFAULT_PORTS = {"muse": 9410, "pip": 9420, "646": 9430, "opm": 9440, "def": 9450, "dev": 9455}
|
||||||
|
|
||||||
|
try:
|
||||||
|
import lookup_engine
|
||||||
|
HAS_LOOKUP_ENGINE = True
|
||||||
|
except ImportError:
|
||||||
|
HAS_LOOKUP_ENGINE = False
|
||||||
|
|
||||||
# Matches EVERY [RESULT <job_id>] marker in a message (use with finditer, not
|
# Matches EVERY [RESULT <job_id>] marker in a message (use with finditer, not
|
||||||
# search). The result text is lazy and stops before the next marker (or end of
|
# search). The result text is lazy and stops before the next marker (or end of
|
||||||
# text), so a message closing two jobs records each with its own text instead
|
# text), so a message closing two jobs records each with its own text instead
|
||||||
# of the first marker greedily swallowing the second.
|
# of the first marker greedily swallowing the second.
|
||||||
RESULT_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_-]+)\]\s*(.*?)(?=\[RESULT\s|\Z)", re.S)
|
_LOCAL_RESULT_RE = re.compile(r"\[RESULT\s+([A-Za-z0-9_/-]+)\]\s*(.*?)(?=\[RESULT\s|\Z)", re.S)
|
||||||
|
RESULT_RE = lookup_engine.get_result_regex() if HAS_LOOKUP_ENGINE else _LOCAL_RESULT_RE
|
||||||
|
|
||||||
def iter_result_markers(text):
|
def iter_result_markers(text):
|
||||||
"""Yield (job_id, result_text) for every [RESULT <job_id>] marker in text."""
|
"""Yield (job_id, result_text) for every [RESULT <job_id>] marker in text."""
|
||||||
for m in RESULT_RE.finditer(text or ""):
|
current_re = lookup_engine.get_result_regex() if HAS_LOOKUP_ENGINE else RESULT_RE
|
||||||
|
for m in current_re.finditer(text or ""):
|
||||||
yield m.group(1).strip(), m.group(2).strip()
|
yield m.group(1).strip(), m.group(2).strip()
|
||||||
|
|
||||||
|
|
||||||
@@ -103,12 +117,14 @@ def iter_result_markers(text):
|
|||||||
# ACK/CLAIM acknowledge a digest (nudge-suppressed, NOT closed);
|
# ACK/CLAIM acknowledge a digest (nudge-suppressed, NOT closed);
|
||||||
# RESULT/DECLINE/NO-ACTION close the digest. Every verb match records
|
# RESULT/DECLINE/NO-ACTION close the digest. Every verb match records
|
||||||
# outcome=<verb> on the followup record.
|
# outcome=<verb> on the followup record.
|
||||||
VERB_RE = re.compile(r"\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_-]+)\]")
|
_LOCAL_VERB_RE = re.compile(r"\[(ACK|CLAIM|RESULT|DECLINE|NO-ACTION)\s+([A-Za-z0-9_/-]+)\]")
|
||||||
|
VERB_RE = lookup_engine.get_verb_regex() if HAS_LOOKUP_ENGINE else _LOCAL_VERB_RE
|
||||||
|
|
||||||
|
|
||||||
def iter_verb_markers(text):
|
def iter_verb_markers(text):
|
||||||
"""Yield (verb, job_id) for every [VERB <job_id>] marker in text."""
|
"""Yield (verb, job_id) for every [VERB <job_id>] marker in text."""
|
||||||
for m in VERB_RE.finditer(text or ""):
|
current_re = lookup_engine.get_verb_regex() if HAS_LOOKUP_ENGINE else VERB_RE
|
||||||
|
for m in current_re.finditer(text or ""):
|
||||||
yield m.group(1), m.group(2).strip()
|
yield m.group(1), m.group(2).strip()
|
||||||
|
|
||||||
|
|
||||||
@@ -503,6 +519,26 @@ def get_monitored_threads(target_agent=None):
|
|||||||
except Exception:
|
except Exception:
|
||||||
continue
|
continue
|
||||||
|
|
||||||
|
# Also monitor active swarm slot subagents from swarms.json
|
||||||
|
try:
|
||||||
|
if SWARM_FILE.exists():
|
||||||
|
s_data = json.loads(SWARM_FILE.read_text(encoding="utf-8"))
|
||||||
|
for sid, s_info in s_data.items():
|
||||||
|
if sid.startswith("_") or not isinstance(s_info, dict):
|
||||||
|
continue
|
||||||
|
if s_info.get("status") not in ("running", "pending"):
|
||||||
|
continue
|
||||||
|
for slot in s_info.get("slots", []):
|
||||||
|
if slot.get("status") == "running":
|
||||||
|
sub_id = slot.get("subagent_session_id") or slot.get("thread_uuid") or slot.get("sidechat_id")
|
||||||
|
ag = slot.get("agent_id")
|
||||||
|
if sub_id and ag and ag in threads_by_agent:
|
||||||
|
existing = [t["id"] for t in threads_by_agent[ag]]
|
||||||
|
if sub_id not in existing:
|
||||||
|
threads_by_agent[ag].append({"id": sub_id, "name": f"swarm-{sid[:8]}-s{slot.get('slot')}"})
|
||||||
|
except Exception as se:
|
||||||
|
sys.stderr.write(f"warning: failed to add swarm threads to monitor list: {se}\n")
|
||||||
|
|
||||||
return threads_by_agent
|
return threads_by_agent
|
||||||
|
|
||||||
|
|
||||||
@@ -783,8 +819,7 @@ def process_messages(raw_messages, agent, thread_id, thread_name, last_wm, follo
|
|||||||
sys.stderr.write(f"warning: failed to record swarm report: {se}\n")
|
sys.stderr.write(f"warning: failed to record swarm report: {se}\n")
|
||||||
else:
|
else:
|
||||||
trigger_chain_next(job_id, result_text, success=not is_fail)
|
trigger_chain_next(job_id, result_text, success=not is_fail)
|
||||||
if not is_fail:
|
archive_ephemeral_thread(agent, thread_id, job_id=job_id)
|
||||||
archive_ephemeral_thread(agent, thread_id, job_id=job_id)
|
|
||||||
clear_matching_followups(followups, agent, thread_id, mid, text,
|
clear_matching_followups(followups, agent, thread_id, mid, text,
|
||||||
dry_run, job_id=job_id, verb="RESULT")
|
dry_run, job_id=job_id, verb="RESULT")
|
||||||
for verb, job_id in verbs:
|
for verb, job_id in verbs:
|
||||||
@@ -1047,11 +1082,19 @@ def reconcile_and_dispatch_swarms(dry_run=False):
|
|||||||
slot_target = f"{sid}-s{slot_idx}"
|
slot_target = f"{sid}-s{slot_idx}"
|
||||||
|
|
||||||
# Prepare task directive prompt
|
# Prepare task directive prompt
|
||||||
prompt = (
|
slot_job_id = f"{sid}/{slot_idx}"
|
||||||
f"[JOB {sid}/{slot_idx}] Task for swarm slot {slot_idx}:\n"
|
if HAS_PROMPT_ENVELOPE and hasattr(prompt_envelope, "wrap_subagent_task"):
|
||||||
f"{task_text}\n\n"
|
prompt = prompt_envelope.wrap_subagent_task(slot_job_id, task_text)
|
||||||
f"Reply with [RESULT {sid}/{slot_idx}] OK <summary> or FAIL <reason>."
|
else:
|
||||||
)
|
prompt = (
|
||||||
|
f"Operator assignment for swarm slot {slot_job_id}:\n\n"
|
||||||
|
f"Task:\n{task_text.strip()}\n\n"
|
||||||
|
f"Instructions:\n"
|
||||||
|
f"1. Carry out this task directly using your available tools.\n"
|
||||||
|
f"2. When finished, conclude your final response with your verdict line:\n"
|
||||||
|
f"[RESULT {slot_job_id}] OK: <one-line summary of actions and outcome>\n"
|
||||||
|
f"(or [RESULT {slot_job_id}] FAIL: <reason> if the task could not be completed)\n"
|
||||||
|
)
|
||||||
|
|
||||||
# Native Subagent Execution Bridge:
|
# Native Subagent Execution Bridge:
|
||||||
# Spawns an interactive child session directly on the target worker
|
# Spawns an interactive child session directly on the target worker
|
||||||
@@ -1196,7 +1239,13 @@ def check_and_archive_terminal_swarms():
|
|||||||
for sid, swarm in swarms_data.items():
|
for sid, swarm in swarms_data.items():
|
||||||
st = swarm.get("status")
|
st = swarm.get("status")
|
||||||
if st in ("completed", "partial", "killed"):
|
if st in ("completed", "partial", "killed"):
|
||||||
# Check slots or matching sidechats
|
# Check slot subagent sessions
|
||||||
|
for slot in swarm.get("slots", []):
|
||||||
|
sub_id = slot.get("subagent_session_id")
|
||||||
|
ag = slot.get("agent_id")
|
||||||
|
if sub_id and ag:
|
||||||
|
archive_ephemeral_thread(ag, sub_id, job_id=sid)
|
||||||
|
# Check matching registered sidechats
|
||||||
for key, val in sc_data.items():
|
for key, val in sc_data.items():
|
||||||
if isinstance(val, dict) and not val.get("archived") and val.get("type") != "persistent":
|
if isinstance(val, dict) and not val.get("archived") and val.get("type") != "persistent":
|
||||||
if sid in key or (swarm.get("label") and swarm.get("label") in key):
|
if sid in key or (swarm.get("label") and swarm.get("label") in key):
|
||||||
|
|||||||
@@ -23,15 +23,33 @@ import sys
|
|||||||
import time
|
import time
|
||||||
import traceback
|
import traceback
|
||||||
|
|
||||||
# Sibling modules live in the same directory.
|
# Sibling modules live in swarm_worker, common utilities live in bin/.
|
||||||
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
|
_SWARM_DIR = os.path.dirname(os.path.abspath(__file__))
|
||||||
|
_BIN_DIR = os.path.dirname(_SWARM_DIR)
|
||||||
|
sys.path.insert(0, _SWARM_DIR)
|
||||||
|
sys.path.insert(0, _BIN_DIR)
|
||||||
|
|
||||||
from poller import find_pending_slots
|
from poller import find_pending_slots
|
||||||
from executor import execute_task
|
from executor import execute_task, _looks_like_shell
|
||||||
from reporter import post_result
|
from reporter import post_result, attach_slot
|
||||||
|
|
||||||
|
# Fast gateway integration
|
||||||
|
try:
|
||||||
|
import muse_hybrid
|
||||||
|
HAS_MUSE_HYBRID = True
|
||||||
|
except ImportError:
|
||||||
|
HAS_MUSE_HYBRID = False
|
||||||
|
|
||||||
|
# Prompt envelope formatting
|
||||||
|
try:
|
||||||
|
import prompt_envelope
|
||||||
|
HAS_PROMPT_ENVELOPE = True
|
||||||
|
except ImportError:
|
||||||
|
HAS_PROMPT_ENVELOPE = False
|
||||||
|
|
||||||
POLL_INTERVAL = 60 # seconds between poll cycles
|
POLL_INTERVAL = 60 # seconds between poll cycles
|
||||||
STALE_MINUTES = 5 # slots older than this with no attach are workable
|
STALE_MINUTES = 5 # slots older than this with no attach are workable
|
||||||
|
WORKER_POOL = ["dev", "def", "muse"]
|
||||||
|
|
||||||
# === SAFETY SWITCH ===
|
# === SAFETY SWITCH ===
|
||||||
# True -> observe only: log what WOULD be done, execute/post nothing.
|
# True -> observe only: log what WOULD be done, execute/post nothing.
|
||||||
@@ -46,19 +64,78 @@ logging.basicConfig(
|
|||||||
log = logging.getLogger("swarm-worker")
|
log = logging.getLogger("swarm-worker")
|
||||||
|
|
||||||
|
|
||||||
|
def _select_worker(preferred=None):
|
||||||
|
if preferred and HAS_MUSE_HYBRID and muse_hybrid.is_node_configured(preferred):
|
||||||
|
return preferred
|
||||||
|
for candidate in WORKER_POOL:
|
||||||
|
if HAS_MUSE_HYBRID and muse_hybrid.is_node_configured(candidate):
|
||||||
|
return candidate
|
||||||
|
return preferred or "dev"
|
||||||
|
|
||||||
|
|
||||||
def process_slot(slot):
|
def process_slot(slot):
|
||||||
"""Execute one slot and report the result. Returns True on full success."""
|
"""Execute one slot and report the result. Returns True on full success."""
|
||||||
swarm_id = slot.get("swarm_id")
|
swarm_id = slot.get("swarm_id")
|
||||||
slot_index = slot.get("slot_index")
|
slot_index = slot.get("slot_index")
|
||||||
task_text = slot.get("task_text") or ""
|
task_text = slot.get("task_text") or ""
|
||||||
|
agent_label = slot.get("agent_label")
|
||||||
|
sidechat_id = slot.get("sidechat_id")
|
||||||
tag = "%s/%s" % (swarm_id, slot_index)
|
tag = "%s/%s" % (swarm_id, slot_index)
|
||||||
|
|
||||||
if DRY_RUN:
|
if DRY_RUN:
|
||||||
log.info("[dry-run] would execute slot %s (agent=%s, task %.80r)",
|
log.info("[dry-run] would execute slot %s (agent=%s, task %.80r)",
|
||||||
tag, slot.get("agent_label"), task_text)
|
tag, agent_label, task_text)
|
||||||
return True
|
return True
|
||||||
|
|
||||||
log.info("executing slot %s (agent=%s)", tag, slot.get("agent_label"))
|
# If the task is NOT a shell command, dispatch it to an ephemeral Muse subagent.
|
||||||
|
is_shell = _looks_like_shell(task_text)
|
||||||
|
if not is_shell and HAS_MUSE_HYBRID:
|
||||||
|
worker_agent = _select_worker(agent_label)
|
||||||
|
log.info("dispatching subagent slot %s to %s", tag, worker_agent)
|
||||||
|
try:
|
||||||
|
# 1. Start an ephemeral subagent session
|
||||||
|
title = f"sw-{swarm_id[:16]}-s{slot_index}"
|
||||||
|
sess, err = muse_hybrid.start_session(worker_agent, title=title)
|
||||||
|
if err or not sess or not sess.get("session_id"):
|
||||||
|
log.error("failed to start subagent session for %s on %s: %s", tag, worker_agent, err)
|
||||||
|
return False
|
||||||
|
|
||||||
|
sub_sid = sess["session_id"]
|
||||||
|
log.info("subagent session %s created for slot %s on %s", sub_sid, tag, worker_agent)
|
||||||
|
|
||||||
|
# 2. Attach/claim the slot in box state with subagent session_id
|
||||||
|
attached = attach_slot(swarm_id, slot_index, worker_agent, session_id=sub_sid)
|
||||||
|
if not attached:
|
||||||
|
log.warning("failed to attach slot %s to %s; proceeding with dispatch", tag, worker_agent)
|
||||||
|
|
||||||
|
# 3. Format prompt with authentic Operator Directive and RESULT expectation
|
||||||
|
if HAS_PROMPT_ENVELOPE and hasattr(prompt_envelope, "wrap_subagent_task"):
|
||||||
|
prompt_body = prompt_envelope.wrap_subagent_task(tag, task_text)
|
||||||
|
else:
|
||||||
|
prompt_body = (
|
||||||
|
f"Operator assignment for swarm slot {tag}:\n\n"
|
||||||
|
f"Task:\n{task_text.strip()}\n\n"
|
||||||
|
f"Instructions:\n"
|
||||||
|
f"1. Carry out this task directly using your available tools.\n"
|
||||||
|
f"2. When finished, conclude your final response with your verdict line:\n"
|
||||||
|
f"[RESULT {tag}] OK: <one-line summary of actions and outcome>\n"
|
||||||
|
f"(or [RESULT {tag}] FAIL: <reason> if the task could not be completed)\n"
|
||||||
|
)
|
||||||
|
|
||||||
|
# 4. Asynchronously send message to subagent session
|
||||||
|
res, send_err = muse_hybrid.send_message(worker_agent, prompt_body, thread_id=sub_sid, wait=0)
|
||||||
|
if send_err:
|
||||||
|
log.error("failed to send task to subagent %s on %s: %s", sub_sid, worker_agent, send_err)
|
||||||
|
return False
|
||||||
|
|
||||||
|
log.info("slot %s successfully dispatched to subagent %s (harvester will harvest)", tag, sub_sid)
|
||||||
|
return True
|
||||||
|
except Exception:
|
||||||
|
log.error("subagent dispatch crashed on slot %s:\n%s", tag, traceback.format_exc())
|
||||||
|
return False
|
||||||
|
|
||||||
|
# Otherwise fallback to sandboxed host execution
|
||||||
|
log.info("executing slot %s in sandbox (agent=%s)", tag, agent_label)
|
||||||
try:
|
try:
|
||||||
result = execute_task(task_text)
|
result = execute_task(task_text)
|
||||||
except Exception:
|
except Exception:
|
||||||
|
|||||||
@@ -77,11 +77,19 @@ def _looks_like_shell(task_text):
|
|||||||
return False
|
return False
|
||||||
if _PROSE_RE.search(t):
|
if _PROSE_RE.search(t):
|
||||||
return False
|
return False
|
||||||
# Must start with a word-ish token (not a quote or sentence).
|
# Must start with a plausible executable command or path
|
||||||
first = t.split()[0] if t.split() else ""
|
first = t.split()[0] if t.split() else ""
|
||||||
if not re.match(r"^[a-zA-Z0-9_.\-/]+$", first):
|
if not re.match(r"^[a-zA-Z0-9_.\-/]+$", first):
|
||||||
return False
|
return False
|
||||||
return True
|
# If it's a relative/absolute path, verify it exists and is executable
|
||||||
|
if "/" in first:
|
||||||
|
return os.path.isfile(first) and os.access(first, os.X_OK)
|
||||||
|
# If it's a bare command name, it must exist in standard system bin paths
|
||||||
|
for p in ("/bin", "/usr/bin", "/usr/local/bin"):
|
||||||
|
candidate = os.path.join(p, first)
|
||||||
|
if os.path.isfile(candidate) and os.access(candidate, os.X_OK):
|
||||||
|
return True
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
def _refused(task_text):
|
def _refused(task_text):
|
||||||
|
|||||||
@@ -97,7 +97,8 @@ def find_pending_slots(stale_minutes=STALE_MINUTES):
|
|||||||
return found
|
return found
|
||||||
swarms = data.get("swarms", []) if isinstance(data, dict) else []
|
swarms = data.get("swarms", []) if isinstance(data, dict) else []
|
||||||
for summary in swarms:
|
for summary in swarms:
|
||||||
if (summary.get("status") or "").lower() != "running":
|
summary_status = (summary.get("status") or "").lower()
|
||||||
|
if summary_status not in ("running", "pending"):
|
||||||
continue
|
continue
|
||||||
swarm_id = summary.get("swarm_id")
|
swarm_id = summary.get("swarm_id")
|
||||||
if not swarm_id:
|
if not swarm_id:
|
||||||
|
|||||||
@@ -97,6 +97,60 @@ def post_result(swarm_id, slot_index, result_dict, dry_run=False):
|
|||||||
swarm_id, slot_index, str(resp)[:500])
|
swarm_id, slot_index, str(resp)[:500])
|
||||||
return False
|
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
|
return True
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user