Files

342 lines
11 KiB
Python

"""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: <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))
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)