"""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 plausible executable command or path first = t.split()[0] if t.split() else "" if not re.match(r"^[a-zA-Z0-9_.\-/]+$", first): return False # 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", "/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: " or "Execute this shell command...: " 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)) 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)