fix(runtime): replace tmux pane capture with explicit bounded output (ODY-150)

This commit is contained in:
Alexandre Teixeira
2026-10-01 21:43:17 +01:00
parent 85cfe59b15
commit 4efb85ee33
6 changed files with 45 additions and 456 deletions
+8 -296
View File
@@ -9,16 +9,14 @@ import shutil
import subprocess
import sys
import time
import collections
import json
from typing import Optional, Callable, Awaitable, Tuple, Dict
from typing import Optional
from urllib.parse import urlparse
import httpx
from src import containment
from src.constants import AGENT_ISOLATED_TMP_DIRNAME, MAX_OUTPUT_CHARS, WORKSPACE_MOUNT
from src.agent_runtime.journal import mark_operation_started
logger = logging.getLogger(__name__)
@@ -29,9 +27,6 @@ logger = logging.getLogger(__name__)
DEFAULT_BASH_TIMEOUT = 120
DEFAULT_PYTHON_TIMEOUT = 60 * 60
PROGRESS_INTERVAL_S = 2.0
PROGRESS_TAIL_LINES = 12
TMUX_CAPTURE_LINES = 2000
_HOST_SHELL_BRIDGE_HOSTS = {"127.0.0.1", "localhost", "::1", "host.docker.internal"}
IS_WINDOWS = sys.platform.startswith("win")
_HOST_SHELL_CANCEL_TASKS: set[asyncio.Task] = set()
@@ -221,11 +216,6 @@ def is_host_shell_bridge_url_allowed(url: str) -> bool:
return True
def _tmux_session_name(session_id: Optional[str]) -> str:
raw = re.sub(r"[^A-Za-z0-9_.-]+", "-", str(session_id or "default")).strip("-")
return f"ody-agent-{raw[:80] or 'default'}"
def _replace_workspace_alias(content: str, cwd: str) -> str:
"""Map virtual /workspace paths without corrupting absolute host paths."""
return re.sub(
@@ -361,7 +351,7 @@ def _filesystem_boundary_block(mechanism: str, mode: str, *, confined: bool) ->
Reports the **filesystem dimension only**, deliberately. The probe knows
this host could also give a process group and a real wall clock, but
BashTool and PythonTool still assemble their own ``create_subprocess_*``
Compatibility namespace previews assemble their own ``create_subprocess_*``
call and pass neither ``start_new_session`` nor a group-wide kill, so
listing those dimensions here would be the false claim
:mod:`src.containment` calls worse than an honest absence. They arrive when
@@ -529,288 +519,6 @@ def _wrap_workspace_namespace(
return shlex.join(args)
async def _run_exec(*args: str, timeout: float = 10) -> Tuple[str, str, int]:
proc = await asyncio.create_subprocess_exec(
*args,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
)
try:
out_b, err_b = await asyncio.wait_for(proc.communicate(), timeout=timeout)
except asyncio.TimeoutError:
try:
proc.kill()
except Exception:
pass
return "", "timeout", 124
return (
out_b.decode("utf-8", errors="replace"),
err_b.decode("utf-8", errors="replace"),
proc.returncode or 0,
)
async def _tmux_has_session(name: str) -> bool:
_, _, rc = await _run_exec("tmux", "has-session", "-t", name, timeout=3)
return rc == 0
async def _tmux_capture(name: str) -> str:
out, _, _ = await _run_exec(
"tmux", "capture-pane", "-p", "-J", "-S", f"-{TMUX_CAPTURE_LINES}", "-t", name,
timeout=5,
)
return out
async def _tmux_send_line(name: str, line: str) -> None:
if line:
await _run_exec("tmux", "send-keys", "-t", name, "-l", line, timeout=5)
await _run_exec("tmux", "send-keys", "-t", name, "C-m", timeout=5)
async def _ensure_tmux_session(name: str, cwd: str, env: Optional[dict]) -> None:
# tmux creates child panes from the long-lived server environment, not
# necessarily from the app process that issued ``new-session``. On hosts
# where tmux predates the Odysseus virtualenv this silently resolves
# ``python`` to the system interpreter, losing plotting/PDF dependencies
# and prompting futile pip-install loops. Reassert the small execution
# environment on both new and reused panes.
forwarded_env = {
key: str(env[key])
for key in ("PATH", "VIRTUAL_ENV", "HOME", "TMPDIR")
if env and env.get(key)
}
if await _tmux_has_session(name):
if forwarded_env:
exports = " ".join(
f"{key}={shlex.quote(value)}" for key, value in forwarded_env.items()
)
await _tmux_send_line(name, f"export {exports}")
await _run_exec("tmux", "send-keys", "-t", name, "stty -echo", "C-m", timeout=5)
return
env_args = [f"{key}={value}" for key, value in forwarded_env.items()]
await _run_exec(
"tmux", "new-session", "-d", "-s", name, "-c", cwd,
"env",
*env_args,
f"TERM={env.get('TERM', 'xterm-256color') if env else 'xterm-256color'}",
f"COLUMNS={env.get('COLUMNS', '120') if env else '120'}",
f"LINES={env.get('LINES', '40') if env else '40'}",
"/bin/bash",
"--noprofile",
"--norc",
timeout=10,
)
if not await _tmux_has_session(name):
raise RuntimeError(f"failed to create tmux session {name}")
await _run_exec("tmux", "send-keys", "-t", name, "stty -echo", "C-m", timeout=5)
def _output_after_marker(capture: str, start_marker: str, end_marker: str) -> Tuple[str, bool]:
lines = capture.splitlines()
start_idx = -1
for idx, line in enumerate(lines):
if line.strip() == start_marker:
start_idx = idx
if start_idx < 0:
return capture, False
end_idx = -1
for idx in range(start_idx + 1, len(lines)):
if lines[idx].strip().startswith(end_marker):
end_idx = idx
if end_idx < 0:
return "\n".join(lines[start_idx + 1:]), False
return "\n".join(lines[start_idx + 1:end_idx]), True
def _extract_marker_rc(capture: str, end_marker: str) -> int:
for line in reversed(capture.splitlines()):
stripped = line.strip()
if stripped.startswith(end_marker):
suffix = stripped[len(end_marker):].strip()
if suffix.isdigit():
return int(suffix)
return 0
async def _run_tmux_bash(
content: str,
*,
session_id: str,
cwd: str,
env: Optional[dict],
timeout: float,
progress_cb: Optional[Callable[[Dict], Awaitable[None]]] = None,
) -> Tuple[str, str, Optional[int], bool]:
name = _tmux_session_name(session_id)
await _ensure_tmux_session(name, cwd, env)
stamp = f"{int(time.time() * 1000)}-{abs(hash(content)) % 1000000}"
start_marker = f"__ODYSSEUS_CMD_START_{stamp}__"
end_prefix = f"__ODYSSEUS_CMD_END_{stamp}__:"
# Execute each tool call in a non-interactive child shell. The tmux pane
# is deliberately persistent, but handing its terminal stdin to commands
# lets programs such as ffmpeg block forever on overwrite prompts. EOF is
# the deterministic behavior expected from an agent tool invocation.
child_command = f"/bin/bash -lc {shlex.quote(content)} </dev/null"
wrapped = (
f"printf '\\n{start_marker}\\n'\n"
f"{child_command}\n"
f"__ody_rc=$?\n"
f"printf '\\n{end_prefix}%s\\n' \"$__ody_rc\"\n"
)
for line in wrapped.splitlines():
await _tmux_send_line(name, line)
started = time.time()
last_tail = ""
while True:
capture = await _tmux_capture(name)
body, done = _output_after_marker(capture, start_marker, end_prefix)
tail = "\n".join(body.splitlines()[-PROGRESS_TAIL_LINES:])
if progress_cb and tail != last_tail:
last_tail = tail
try:
await progress_cb({
"elapsed_s": round(time.time() - started, 1),
"tail": tail,
"tmux_session": name,
})
except Exception:
pass
if done:
rc = _extract_marker_rc(capture, end_prefix)
cleaned = _clean_tmux_command_output(body, wrapped)
return cleaned, "", rc, False
if time.time() - started > timeout:
try:
await _run_exec("tmux", "send-keys", "-t", name, "C-c", timeout=3)
except Exception:
pass
# Ctrl-C targets the pane's foreground process group, but a child
# can outlive its wrapper shell and become an orphan. Destroy this
# task-scoped session as the timeout boundary; the next tool call
# recreates it through _ensure_tmux_session.
try:
await _run_exec("tmux", "kill-session", "-t", name, timeout=3)
except Exception:
pass
cleaned = _clean_tmux_command_output(body, wrapped)
return cleaned, "", 124, True
await asyncio.sleep(0.5)
def _clean_tmux_command_output(text: str, wrapped_command: str) -> str:
lines = text.splitlines()
wrapped_lines = {ln.rstrip() for ln in wrapped_command.splitlines() if ln.strip()}
cleaned = []
for line in lines:
raw = line.rstrip()
stripped = raw.strip()
if not stripped:
cleaned.append(raw)
continue
if stripped in wrapped_lines:
continue
if stripped.startswith("__ody_rc=") or stripped.startswith("printf "):
continue
if re.fullmatch(r"(?:bash|sh)-[\d.]+\$ ?", stripped):
continue
if re.fullmatch(r"[\w.@:/~+-]+[#$] ?", stripped):
continue
cleaned.append(raw)
return "\n".join(cleaned).strip()
async def _run_subprocess_streaming(
proc: asyncio.subprocess.Process,
*,
timeout: float,
progress_cb: Optional[Callable[[Dict], Awaitable[None]]] = None,
) -> Tuple[str, str, Optional[int], bool]:
started = time.time()
stdout_full: list[str] = []
stderr_full: list[str] = []
tail = collections.deque(maxlen=PROGRESS_TAIL_LINES)
async def _reader(stream, full_buf, label: str):
if stream is None:
return
while True:
line = await stream.readline()
if not line:
break
decoded = line.decode("utf-8", errors="replace").rstrip("\n")
full_buf.append(decoded)
if label == "err":
tail.append(f"! {decoded}")
else:
tail.append(decoded)
async def _progress_emitter():
await asyncio.sleep(PROGRESS_INTERVAL_S)
while True:
if progress_cb:
try:
await progress_cb({
"elapsed_s": round(time.time() - started, 1),
"tail": "\n".join(list(tail)),
})
except Exception:
pass
await asyncio.sleep(PROGRESS_INTERVAL_S)
rd_out = asyncio.create_task(_reader(proc.stdout, stdout_full, "out"))
rd_err = asyncio.create_task(_reader(proc.stderr, stderr_full, "err"))
prog_task = asyncio.create_task(_progress_emitter()) if progress_cb else None
timed_out = False
try:
await asyncio.wait_for(proc.wait(), timeout=timeout)
except asyncio.TimeoutError:
timed_out = True
try:
proc.kill()
except Exception:
pass
try:
await asyncio.wait_for(proc.wait(), timeout=2)
except Exception:
pass
except asyncio.CancelledError:
try:
proc.kill()
except Exception:
pass
try:
await asyncio.wait_for(proc.wait(), timeout=2)
except Exception:
pass
for t in (rd_out, rd_err):
t.cancel()
if prog_task is not None:
prog_task.cancel()
raise
finally:
if prog_task is not None and not prog_task.done():
prog_task.cancel()
try:
await prog_task
except (asyncio.CancelledError, Exception):
pass
for t in (rd_out, rd_err):
try:
await asyncio.wait_for(t, timeout=1)
except Exception:
pass
return (
"\n".join(stdout_full),
"\n".join(stderr_full),
proc.returncode,
timed_out,
)
def _owned_spec(cwd: str, env: Optional[dict], timeout: int, readonly_extra: tuple = ()) -> containment.ContainmentSpec:
"""Server-defined boundary shared by the native execution tools."""
readonly = []
@@ -853,14 +561,15 @@ async def _run_owned_command(command, ctx: dict, *, tool: str, timeout: int, arg
if result.stderr.rstrip():
output = (output + "\nSTDERR: " + result.stderr.rstrip()).strip()
truncated = result.output_truncated or len(output) > MAX_OUTPUT_CHARS
capture_note = " Captured output was truncated." if truncated else ""
common = {"containment": boundary, "teardown": teardown, "output_truncated": truncated}
if not teardown["dead"]:
return {**common, "error": f"{tool}: process teardown could not verify death",
return {**common, "error": f"{tool}: process teardown could not verify death.{capture_note}",
"failure_kind": "process_teardown_failed", "exit_code": 1,
"stdout": _truncate(result.stdout, MAX_OUTPUT_CHARS),
"stderr": _truncate(result.stderr, MAX_OUTPUT_CHARS)}
if result.timed_out:
return {**common, "error": f"{tool}: timed out after {timeout}s; process tree terminated",
return {**common, "error": f"{tool}: timed out after {timeout}s; process tree terminated.{capture_note}",
"exit_code": 124, "stdout": _truncate(result.stdout, MAX_OUTPUT_CHARS),
"stderr": _truncate(result.stderr, MAX_OUTPUT_CHARS)}
if tool == "python":
@@ -868,6 +577,9 @@ async def _run_owned_command(command, ctx: dict, *, tool: str, timeout: int, arg
if child_failure:
return {**common, "error": _truncate("python: a child operation failed despite a zero Python exit status:\n" + child_failure, MAX_OUTPUT_CHARS),
"exit_code": 1, "stderr": _truncate(result.stderr, MAX_OUTPUT_CHARS)}
if truncated:
note = "\n…[output truncated by containment capture limit]…"
output = output[:MAX_OUTPUT_CHARS - len(note)] + note
return {**common, "output": _truncate(output, MAX_OUTPUT_CHARS) or "(no output)",
"exit_code": result.exit_code if result.exit_code is not None else 1}