chore(fleet): sync operator memory, hatch menu dialogs, and watchdog alerts
This commit is contained in:
+64
-14
@@ -30,8 +30,9 @@ sys.path.insert(0, _SWARM_DIR)
|
||||
sys.path.insert(0, _BIN_DIR)
|
||||
|
||||
from poller import find_pending_slots
|
||||
from executor import execute_task, _looks_like_shell
|
||||
from executor import execute_task, _looks_like_shell, extract_shell_command, execute_task_in_tmux
|
||||
from reporter import post_result, attach_slot
|
||||
from mainloop_notify import notify_via_mainloop
|
||||
|
||||
# Fast gateway integration
|
||||
try:
|
||||
@@ -49,7 +50,7 @@ except ImportError:
|
||||
|
||||
POLL_INTERVAL = 60 # seconds between poll cycles
|
||||
STALE_MINUTES = 5 # slots older than this with no attach are workable
|
||||
WORKER_POOL = ["dev", "def", "muse"]
|
||||
WORKER_POOL = ["muse"] # only dispatch to fully authenticated agent nodes
|
||||
|
||||
# === SAFETY SWITCH ===
|
||||
# True -> observe only: log what WOULD be done, execute/post nothing.
|
||||
@@ -65,12 +66,12 @@ log = logging.getLogger("swarm-worker")
|
||||
|
||||
|
||||
def _select_worker(preferred=None):
|
||||
if preferred and HAS_MUSE_HYBRID and muse_hybrid.is_node_configured(preferred):
|
||||
if preferred and preferred in WORKER_POOL 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"
|
||||
return "muse"
|
||||
|
||||
|
||||
def process_slot(slot):
|
||||
@@ -81,19 +82,56 @@ def process_slot(slot):
|
||||
agent_label = slot.get("agent_label")
|
||||
sidechat_id = slot.get("sidechat_id")
|
||||
tag = "%s/%s" % (swarm_id, slot_index)
|
||||
short_id = swarm_id[3:19] if str(swarm_id).startswith("sw-") else str(swarm_id)[:16]
|
||||
session_name = f"sw-{short_id}-s{slot_index}"
|
||||
|
||||
if DRY_RUN:
|
||||
log.info("[dry-run] would execute slot %s (agent=%s, task %.80r)",
|
||||
tag, agent_label, task_text)
|
||||
return True
|
||||
|
||||
# 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:
|
||||
# 1. Check if the task is or contains an executable shell command
|
||||
cmd = extract_shell_command(task_text)
|
||||
if cmd:
|
||||
log.info("executing slot %s in host tmux session %s on bl", tag, session_name)
|
||||
# Attach/claim slot in box state
|
||||
attach_slot(swarm_id, slot_index, "swarm-worker", session_id=session_name)
|
||||
try:
|
||||
result = execute_task_in_tmux(session_name, cmd)
|
||||
except Exception:
|
||||
log.error("tmux executor crashed on slot %s:\n%s", tag, traceback.format_exc())
|
||||
result = {"success": False, "output": "",
|
||||
"error": "tmux 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
|
||||
|
||||
# Post completion note to the slot sidechat for main loop visibility
|
||||
summary_msg = payload.get("output") or payload.get("error") or "completed"
|
||||
try:
|
||||
notified = notify_via_mainloop(swarm_id, slot_index, summary_msg, worker_id="swarm-worker")
|
||||
log.info("slot %s sidechat notification: %s", tag, notified)
|
||||
except Exception as ne:
|
||||
log.warning("failed to post sidechat notification for %s: %s", tag, ne)
|
||||
|
||||
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
|
||||
|
||||
# 2. If the task is purely prose/instructions, dispatch to a verified agent subagent
|
||||
if HAS_MUSE_HYBRID:
|
||||
worker_agent = _select_worker(agent_label)
|
||||
log.info("dispatching subagent slot %s to %s", tag, worker_agent)
|
||||
log.info("dispatching prose 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"):
|
||||
@@ -103,12 +141,19 @@ def process_slot(slot):
|
||||
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
|
||||
# Attach/claim slot in box state
|
||||
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
|
||||
# Register in subagent_tracker
|
||||
try:
|
||||
import subagent_tracker
|
||||
subagent_tracker.register_session(worker_agent, sub_sid, title=title, prompt=task_text[:200])
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# 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:
|
||||
@@ -122,7 +167,7 @@ def process_slot(slot):
|
||||
f"(or [RESULT {tag}] FAIL: <reason> if the task could not be completed)\n"
|
||||
)
|
||||
|
||||
# 4. Asynchronously send message to subagent session
|
||||
# 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)
|
||||
@@ -134,8 +179,8 @@ def process_slot(slot):
|
||||
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)
|
||||
# 3. Fallback to sandboxed host execution
|
||||
log.info("executing slot %s in fallback sandbox (agent=%s)", tag, agent_label)
|
||||
try:
|
||||
result = execute_task(task_text)
|
||||
except Exception:
|
||||
@@ -154,6 +199,11 @@ def process_slot(slot):
|
||||
log.error("reporter crashed on slot %s:\n%s", tag, traceback.format_exc())
|
||||
posted = False
|
||||
|
||||
try:
|
||||
notify_via_mainloop(swarm_id, slot_index, payload.get("output") or "done", worker_id="swarm-worker")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
log.info("slot %s done: ok=%s posted=%s (%.1fs)",
|
||||
tag, payload["ok"], posted,
|
||||
float(result.get("duration_s") or 0.0))
|
||||
|
||||
@@ -85,13 +85,152 @@ def _looks_like_shell(task_text):
|
||||
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"):
|
||||
for p in ("/bin", "/usr/bin", "/usr/local/bin", "/home/super/Projects/NetVM/bin", "/home/super/.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 extract_shell_command(task_text):
|
||||
"""Extract an executable shell command from task text if present."""
|
||||
t = (task_text or "").strip()
|
||||
if not t:
|
||||
return None
|
||||
if _looks_like_shell(t):
|
||||
return t
|
||||
# Check for "Run: <cmd>" or "Execute this shell command...: <cmd>"
|
||||
m = re.search(r"(?:Run|Execute)(?:\s+this\s+shell\s+command(?:\s+and\s+report\s+its\s+full\s+output)?)?:\s*[`'\"]?([^`'\n]+)[`'\"]?", t, re.IGNORECASE)
|
||||
if m:
|
||||
candidate = m.group(1).strip()
|
||||
if candidate:
|
||||
return candidate
|
||||
# Check for markdown code blocks ```bash ... ``` or ```sh ... ```
|
||||
m = re.search(r"```(?:bash|sh)?\n(.*?)\n```", t, re.DOTALL)
|
||||
if m:
|
||||
candidate = m.group(1).strip()
|
||||
if candidate:
|
||||
return candidate
|
||||
# Check for single backticked command
|
||||
m = re.search(r"`([^`\n]+)`", t)
|
||||
if m:
|
||||
candidate = m.group(1).strip()
|
||||
if _looks_like_shell(candidate):
|
||||
return candidate
|
||||
return None
|
||||
|
||||
|
||||
TMUX_SOCKET = "/tmp/tmux-muse.sock"
|
||||
TMUX_LOG_DIR = "/home/super/Projects/NetVM/logs/tmux"
|
||||
|
||||
|
||||
def execute_task_in_tmux(session_name, cmd_str, timeout=300):
|
||||
"""Execute a task inside a dedicated tmux session on /tmp/tmux-muse.sock.
|
||||
|
||||
Captures output to logs/tmux/{session_name}.log, tracks return code via
|
||||
status file, and returns:
|
||||
dict(success=bool, output=str, duration_s=float, error=str|None)
|
||||
"""
|
||||
os.makedirs(TMUX_LOG_DIR, exist_ok=True)
|
||||
started = time.monotonic()
|
||||
log_file = os.path.join(TMUX_LOG_DIR, f"{session_name}.log")
|
||||
exit_file = f"/tmp/{session_name}.exit"
|
||||
script_file = f"/tmp/{session_name}.sh"
|
||||
|
||||
# Clean up prior artifacts
|
||||
for f in (exit_file, script_file):
|
||||
try:
|
||||
if os.path.exists(f):
|
||||
os.remove(f)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
# Write wrapper script
|
||||
with open(script_file, "w", encoding="utf-8") as sf:
|
||||
sf.write("#!/usr/bin/env bash\n")
|
||||
sf.write("export PATH=\"/home/super/Projects/NetVM/bin:/home/super/.local/bin:/usr/local/bin:/usr/bin:/bin:$PATH\"\n")
|
||||
sf.write("cd /home/super/Projects/NetVM\n")
|
||||
sf.write(f"{cmd_str}\n")
|
||||
sf.write(f"echo $? > \"{exit_file}\"\n")
|
||||
os.chmod(script_file, 0o755)
|
||||
|
||||
# Kill any existing session with this name
|
||||
subprocess.run(["tmux", "-S", TMUX_SOCKET, "kill-session", "-t", session_name],
|
||||
capture_output=True)
|
||||
|
||||
# Start tmux session
|
||||
tmux_cmd = f"bash \"{script_file}\" > \"{log_file}\" 2>&1"
|
||||
res = subprocess.run(
|
||||
["tmux", "-S", TMUX_SOCKET, "new-session", "-d", "-s", session_name, tmux_cmd],
|
||||
capture_output=True, text=True
|
||||
)
|
||||
if res.returncode != 0:
|
||||
dur = round(time.monotonic() - started, 3)
|
||||
return {
|
||||
"success": False,
|
||||
"output": "",
|
||||
"duration_s": dur,
|
||||
"error": f"Failed to create tmux session: {res.stderr.strip()}",
|
||||
}
|
||||
|
||||
# Poll for completion or timeout
|
||||
deadline = started + timeout
|
||||
rc = None
|
||||
while time.monotonic() < deadline:
|
||||
if os.path.exists(exit_file):
|
||||
try:
|
||||
with open(exit_file, "r") as ef:
|
||||
rc = int(ef.read().strip())
|
||||
break
|
||||
except Exception:
|
||||
pass
|
||||
check = subprocess.run(
|
||||
["tmux", "-S", TMUX_SOCKET, "has-session", "-t", session_name],
|
||||
capture_output=True
|
||||
)
|
||||
if check.returncode != 0 and os.path.exists(exit_file):
|
||||
break
|
||||
time.sleep(0.5)
|
||||
|
||||
dur = round(time.monotonic() - started, 3)
|
||||
|
||||
# Clean up tmux session if still running
|
||||
subprocess.run(["tmux", "-S", TMUX_SOCKET, "kill-session", "-t", session_name],
|
||||
capture_output=True)
|
||||
|
||||
# Read output log
|
||||
output = ""
|
||||
if os.path.exists(log_file):
|
||||
try:
|
||||
with open(log_file, "r", encoding="utf-8", errors="replace") as lf:
|
||||
output = lf.read()[:OUTPUT_TRUNCATE]
|
||||
except Exception as e:
|
||||
output = f"Error reading log: {e}"
|
||||
|
||||
# Cleanup temporary script and exit file
|
||||
for f in (exit_file, script_file):
|
||||
try:
|
||||
if os.path.exists(f):
|
||||
os.remove(f)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
if rc is None:
|
||||
return {
|
||||
"success": False,
|
||||
"output": output,
|
||||
"duration_s": dur,
|
||||
"error": f"timeout: exceeded {timeout}s in tmux session",
|
||||
}
|
||||
|
||||
return {
|
||||
"success": (rc == 0),
|
||||
"output": output,
|
||||
"duration_s": dur,
|
||||
"error": None if (rc == 0) else f"exit code {rc}",
|
||||
}
|
||||
|
||||
|
||||
def _refused(task_text):
|
||||
return bool(_REFUSE_RE.search(task_text))
|
||||
|
||||
|
||||
@@ -37,14 +37,24 @@ def notify_via_mainloop(swarm_id, slot_index, message, worker_id="swarm-worker",
|
||||
Returns:
|
||||
True on success (or dry-run), False on failure (logged, not raised).
|
||||
"""
|
||||
target_name = "sw-%s-s%s" % (swarm_id, slot_index)
|
||||
target_name = f"{swarm_id}-s{slot_index}" if str(swarm_id).startswith("sw-") else f"sw-{swarm_id}-s{slot_index}"
|
||||
summary = (message or "").strip().replace("\n", " ")[:NOTE_CHARS]
|
||||
note = "[SWARM-DONE %s/%s] %s" % (swarm_id, slot_index, summary)
|
||||
tag = "%s/%s" % (swarm_id, slot_index)
|
||||
note = "[SWARM-DONE %s] %s" % (tag, summary)
|
||||
|
||||
sender = worker_id if worker_id in ("muse", "pip", "646", "opm", "dev", "def", "super") else "super"
|
||||
to_agent = "opm"
|
||||
|
||||
try:
|
||||
sys.path.insert(0, BIN)
|
||||
import dm
|
||||
uuid = dm.resolve_sidechat_target(target_name)
|
||||
if not uuid:
|
||||
sc = dm.load_sidechat_map() if hasattr(dm, "load_sidechat_map") else {}
|
||||
entry = sc.get(target_name, {})
|
||||
uuid = entry.get("thread_uuid")
|
||||
if entry.get("agent"):
|
||||
to_agent = entry.get("agent")
|
||||
except Exception as e:
|
||||
print("notify_via_mainloop: target resolve failed for %s: %s"
|
||||
% (target_name, e), file=sys.stderr)
|
||||
@@ -55,8 +65,8 @@ def notify_via_mainloop(swarm_id, slot_index, message, worker_id="swarm-worker",
|
||||
return False
|
||||
|
||||
cmd = [sys.executable, DM_PY, "send",
|
||||
"--agent", worker_id,
|
||||
"--to", worker_id,
|
||||
"--agent", sender,
|
||||
"--to", to_agent,
|
||||
"--target", uuid,
|
||||
note]
|
||||
if dry_run:
|
||||
|
||||
@@ -96,6 +96,8 @@ def post_result(swarm_id, slot_index, result_dict, dry_run=False):
|
||||
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.
|
||||
|
||||
Reference in New Issue
Block a user