Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 10 additions & 2 deletions src/agent_env/env/envs/multi_env.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,12 @@

logger = logging.getLogger(__name__)

# `docker compose` has data races that end in a Go runtime panic (exit 2) while stopping or
# recreating containers. The snapshot restore's lifecycle commands are idempotent, so a crashed
# attempt is re-run this many extra times instead of failing the load.
_COMPOSE_CRASH_EXIT_CODES = (2,)
_COMPOSE_CRASH_RETRIES = 2

# Off by default: a bake builds and pushes a multi-GB servicedb image, which is the right
# trade once per universe and the wrong one on every deploy. Callers that know they are
# seeding a reusable universe pass snapshot_after_load=True explicitly; this flag exists
Expand Down Expand Up @@ -592,7 +598,8 @@ async def _load_from_snapshot(self, snapshot) -> None:
# Remove old servicedb container and its anonymous volume, then start fresh from snapshot image
logger.info("Recreating servicedb container with snapshot image...")
await self._sandbox.exec_script(
f"cd {GATEWAY_APP_DIR} && docker compose rm -sf -v {DATABASE_SERVICE_NAME} && docker compose up -d {DATABASE_SERVICE_NAME} 2>&1"
f"cd {GATEWAY_APP_DIR} && docker compose rm -sf -v {DATABASE_SERVICE_NAME} && docker compose up -d {DATABASE_SERVICE_NAME} 2>&1",
max_retries=_COMPOSE_CRASH_RETRIES, retry_exit_codes=_COMPOSE_CRASH_EXIT_CODES,
)

# Wait for servicedb to be healthy
Expand Down Expand Up @@ -620,7 +627,8 @@ async def _load_from_snapshot(self, snapshot) -> None:
svc_list = " ".join(environment_names + [GATEWAY_SERVICE_NAME, PGWEB_SERVICE_NAME, DB_MCP_SERVICE_NAME])
logger.info(f"Recreating services: {svc_list}")
await self._sandbox.exec_script(
f"cd {GATEWAY_APP_DIR} && docker compose up -d --force-recreate {svc_list}"
f"cd {GATEWAY_APP_DIR} && docker compose up -d --force-recreate {svc_list}",
max_retries=_COMPOSE_CRASH_RETRIES, retry_exit_codes=_COMPOSE_CRASH_EXIT_CODES,
)
await gw._wait_for_gateway(self._sandbox, AGENT_ENV_GATEWAY_MCP_PORT)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -44,13 +44,13 @@

logger = logging.getLogger(__name__)

# Ubuntu 22.04 + Docker engine + the compose v2 plugin baked in — the Modal-native
# Ubuntu 22.04 + Docker engine + the compose CLI plugin baked in — the Modal-native
# equivalent of a KubeVirt containerdisk image with Docker pre-installed,
# matching its OS baseline so the two backends behave alike. Modal builds/caches this
# image (no ECR push / no auth — public base). `docker.io` is the engine; the compose v2
# image (no ECR push / no auth — public base). `docker.io` is the engine; the compose
# plugin isn't in Ubuntu's repos, so we drop the release binary into the CLI-plugins dir.
# amd64: Modal VM sandboxes run x86_64, matching the linux-x86_64 compose binary.
_COMPOSE_VERSION = "v2.29.7"
_COMPOSE_VERSION = "v5.5.1"
_VM_IMAGE = (
modal.Image.from_registry("ubuntu:22.04")
.apt_install("docker.io", "curl", "aria2") # aria2: object downloads, see _DL_*
Expand Down
33 changes: 26 additions & 7 deletions src/agent_env/providers/sandbox_providers/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -163,6 +163,20 @@ async def write_file_from_text(self, content: str, destination_path: str) -> Non
raise NotImplementedError(f"{self.__class__.__name__} does not support write_file_from_text")


_OUTPUT_CLIP_CHARS = 1500
_RETRY_LOG_CLIP_CHARS = 200


def clip_output(text: str, limit: int = _OUTPUT_CLIP_CHARS) -> str:
"""``text`` whole when it fits in ``2 * limit`` chars, otherwise its first and last ``limit``
chars around an elision marker. A tail alone drops the one line that explains a crash (a Go
``panic:`` header sits above thousands of goroutine frames); a head alone drops the final
error of a long log."""
if len(text) <= 2 * limit:
return text
return f"{text[:limit]}\n... [{len(text) - 2 * limit} chars elided] ...\n{text[-limit:]}"


class VmSandbox(Sandbox):
"""VM-style sandbox with a Docker daemon and shell access inside."""

Expand All @@ -177,25 +191,30 @@ def container_name(self) -> str:
teardown removes only its own container."""
return "agent-api"

async def exec_script(self, script: str, *, max_retries: int = 0) -> str:
async def exec_script(self, script: str, *, max_retries: int = 0, retry_exit_codes: tuple[int, ...] = ()) -> str:
"""Execute a bash script in the sandbox.

Set ``max_retries`` > 0 only for idempotent scripts. Retries are gated
on exit code -1, which a provider's exec client returns when the server
closes the websocket without sending an exit frame (e.g. a control
plane wrapping a transient port-forward 500 as a generic error). Real
script failures (positive exit codes) raise immediately.
plane wrapping a transient port-forward 500 as a generic error), plus
any code in ``retry_exit_codes``, for a script whose binary is known to
die in a way a re-run cures (``docker compose`` hitting one of its data
races exits 2, a Go runtime panic). Other failures raise immediately.
"""
retryable = {-1, *retry_exit_codes}
for attempt in range(max_retries + 1):
exit_code, stdout, stderr = await self.exec_with_output("sudo", "bash", "-c", script)
if exit_code == 0:
return stdout
if exit_code != -1 or attempt == max_retries:
raise RuntimeError(f"Script failed (exit {exit_code}):\nstdout: {stdout[-1500:]}\nstderr: {stderr[-1500:]}")
if exit_code not in retryable or attempt == max_retries:
raise RuntimeError(
f"Script failed (exit {exit_code}):\nstdout: {clip_output(stdout)}\nstderr: {clip_output(stderr)}"
)
backoff = 2 ** attempt
logger.warning(
f"exec_script exit -1 (transient server error), retrying in {backoff}s "
f"(attempt {attempt + 1}/{max_retries + 1}); stderr tail: {stderr[-200:]!r}"
f"exec_script exit {exit_code} ({'transient server error' if exit_code == -1 else 'retryable'}), "
f"retrying in {backoff}s (attempt {attempt + 1}/{max_retries + 1}); output: {clip_output(stderr or stdout, _RETRY_LOG_CLIP_CHARS)!r}"
)
await asyncio.sleep(backoff)
raise AssertionError("unreachable")
Expand Down
86 changes: 85 additions & 1 deletion tst/unit/providers/sandbox_providers/vm_sandbox_test.py
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,10 @@
from types import SimpleNamespace

import pytest
from unittest.mock import AsyncMock

from agent_env.providers.sandbox_providers import sandbox as sandbox_module
from agent_env.providers.sandbox_providers.sandbox import VmSandbox
from agent_env.providers.sandbox_providers.sandbox import VmSandbox, clip_output
from agent_env.store import set_object_store
from agent_env.config import reset_config

Expand Down Expand Up @@ -386,3 +387,86 @@ async def test_write_host_file_writes_on_the_host_in_bounded_chunks():
b64_path = "/tmp/agentenv_run_code/input.json.b64"
assert _b64_from_chunk_scripts(sandbox.scripts, b64_path) == base64.b64encode(data).decode()
assert sandbox.scripts[-1] == f"base64 -d {b64_path} > /tmp/agentenv_run_code/input.json && rm -f {b64_path}"


class _ExitCodeSandbox(VmSandbox):
"""Concrete VmSandbox whose scripts exit with a scripted sequence of codes, then 0."""

def __init__(self, exit_codes: list[int], stderr: str = ""):
self.sandbox_id = "vm-test"
self._exit_codes = list(exit_codes)
self._stderr = stderr
self.calls = 0

async def terminate(self) -> None: # pragma: no cover - not exercised
pass

async def exec(self, *command): # pragma: no cover - not exercised
return None

async def exec_with_output(self, *args):
self.calls += 1
code = self._exit_codes.pop(0) if self._exit_codes else 0
return (0, "ok", "") if code == 0 else (code, "", self._stderr)


@pytest.fixture
def no_backoff(monkeypatch):
monkeypatch.setattr(sandbox_module.asyncio, "sleep", AsyncMock())


# A Go crash: the one informative line, then thousands of goroutine frames.
_PANIC = "panic: concurrent map writes\n\ngoroutine 1 [running]:\n" + "net/http.(*persistConn).readLoop(0xc0004e66c0)\n" * 200


@pytest.mark.asyncio
async def test_exec_script_retries_a_listed_exit_code(no_backoff):
sb = _ExitCodeSandbox([2, 2], stderr=_PANIC)
assert await sb.exec_script("docker compose up -d", max_retries=2, retry_exit_codes=(2,)) == "ok"
assert sb.calls == 3


@pytest.mark.asyncio
async def test_exec_script_does_not_retry_an_unlisted_exit_code(no_backoff):
sb = _ExitCodeSandbox([1, 0], stderr="boom")
with pytest.raises(RuntimeError, match=r"exit 1"):
await sb.exec_script("false", max_retries=2, retry_exit_codes=(2,))
assert sb.calls == 1


@pytest.mark.asyncio
async def test_exec_script_without_retry_codes_still_fails_fast_on_exit_2(no_backoff):
sb = _ExitCodeSandbox([2, 0], stderr=_PANIC)
with pytest.raises(RuntimeError, match=r"exit 2"):
await sb.exec_script("pg_isready", max_retries=2)
assert sb.calls == 1


@pytest.mark.asyncio
async def test_exec_script_gives_up_after_max_retries(no_backoff):
sb = _ExitCodeSandbox([2, 2, 2], stderr=_PANIC)
with pytest.raises(RuntimeError, match=r"exit 2"):
await sb.exec_script("docker compose up -d", max_retries=1, retry_exit_codes=(2,))
assert sb.calls == 2


@pytest.mark.asyncio
async def test_exec_script_error_keeps_the_panic_header_and_the_tail(no_backoff):
sb = _ExitCodeSandbox([2], stderr=_PANIC)
with pytest.raises(RuntimeError) as exc:
await sb.exec_script("docker compose up -d")
message = str(exc.value)
assert "panic: concurrent map writes" in message
assert message.rstrip().endswith("readLoop(0xc0004e66c0)")
assert "chars elided" in message


def test_clip_output_returns_text_that_fits_whole():
assert clip_output("short", 10) == "short"
assert clip_output("x" * 20, 10) == "x" * 20


def test_clip_output_keeps_head_and_tail_around_a_marker():
text = "H" * 10 + "m" * 5 + "T" * 10
assert clip_output(text, 10) == "H" * 10 + "\n... [5 chars elided] ...\n" + "T" * 10

Loading