Skip to content
Merged
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
19 changes: 2 additions & 17 deletions src/agent_env/a2a_agent/validator.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,6 +54,7 @@
VIDEO_PROBE_MP4_B64,
VIDEO_PROBE_PROMPT,
)
from agent_env.task_step.task_steps.sandbox_utils.sandbox_utils import find_agent_container

if TYPE_CHECKING:
from agent_env.a2a_agent.a2a_agent import A2AAgent
Expand Down Expand Up @@ -686,7 +687,7 @@ def record(*, supported, advertised, save_ok, apply_ok, roundtrip_ok, note=""):
if apply_agent.sandbox_type else get_agent_sandbox_provider())
sandbox = await provider.get_sandbox(apply_agent.sandbox_id)
if sandbox.mode == SANDBOX_MODE_VM:
container = await A2AAgentValidator._discover_agent_container(sandbox)
container = await find_agent_container(sandbox)
args = ("sudo", "docker", "exec", container, "cat", marker_path)
else:
args = ("cat", marker_path)
Expand All @@ -706,22 +707,6 @@ def record(*, supported, advertised, save_ok, apply_ok, roundtrip_ok, note=""):
record(supported=roundtrip_ok, advertised=advertised, save_ok=save_ok,
apply_ok=True, roundtrip_ok=roundtrip_ok)

@staticmethod
async def _discover_agent_container(sandbox) -> str:
"""Find the agent container on a VM sandbox (mirrors collect_artifacts /
verify_sandbox): prefer 'agent-api', else the first 'a2a-agent-*'."""
exit_code, stdout, stderr = await sandbox.exec_with_output(
"sudo", "docker", "ps", "--format", "{{.Names}}")
if exit_code != 0:
raise RuntimeError(f"docker ps failed: {stderr[:200]}")
running = [n.strip() for n in stdout.splitlines() if n.strip()]
if sandbox.container_name in running:
return sandbox.container_name
fallback = [n for n in running if n.startswith("a2a-agent-")]
if not fallback:
raise RuntimeError(f"no agent container found; running: {running}")
return fallback[0]

@staticmethod
def _upload_install_test_image_fixture(agent: "A2AAgent"):
"""Upload a minimal Dockerfile as a FileArtifactUniverse so the install
Expand Down
2 changes: 1 addition & 1 deletion src/agent_env/env/envs/mcp_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -222,7 +222,7 @@ async def _copy_artifact_into_container(self, file_artifact) -> str:
await self._sandbox.load_s3_file(file_artifact.object_url, vm_temp_path)
container_id = await self._env_provider._get_container_id(self._sandbox, self.environment_name)
await self._sandbox.exec_script(f"docker exec {container_id} mkdir -p /data")
await self._sandbox.exec_script(f"docker cp {vm_temp_path} {container_id}:{container_path}")
await self._sandbox.docker_cp(vm_temp_path, f"{container_id}:{container_path}")
await self._sandbox.exec_script(f"rm -f {vm_temp_path}")
return container_path

Expand Down
13 changes: 12 additions & 1 deletion src/agent_env/providers/sandbox_providers/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -326,11 +326,22 @@ async def _remove_vm_temp_file(self, *vm_paths: str) -> None:
except Exception as e:
logger.warning(f"Best-effort cleanup of {', '.join(vm_paths)} failed (ignored): {e}")

async def docker_cp(self, source: str, destination: str, *, remove_source: bool = False) -> None:
"""``docker cp source destination``, one side ``container:path``. The paths go as arguments, not
script text, so a sandbox that maps its paths (the local one maps /app) maps only the host side.
``remove_source`` deletes the copied host file in the same exec."""
script = 'docker cp "$1" "$2"' + (' && rm -f "$1"' if remove_source else "")
exit_code, stdout, stderr = await self.exec_with_output("sudo", "bash", "-c", script, "docker-cp", source, destination)
if exit_code != 0:
raise RuntimeError(
f"docker cp {source} {destination} failed (exit {exit_code}):\nstdout: {stdout[-1500:]}\nstderr: {stderr[-1500:]}"
)

async def _copy_into_container(self, vm_path: str, destination_path: str) -> None:
parent = os.path.dirname(destination_path)
if parent:
await self.exec_script(f"docker exec -u 0 {shlex.quote(self.container_name)} mkdir -p {shlex.quote(parent)}")
await self.exec_script(f"docker cp {shlex.quote(vm_path)} {self.container_name}:{shlex.quote(destination_path)}")
await self.docker_cp(vm_path, f"{self.container_name}:{destination_path}")

@staticmethod
def _staging_path(kind: str, destination_path: str) -> str:
Expand Down
38 changes: 3 additions & 35 deletions src/agent_env/task_step/task_steps/collect_artifacts.py
Original file line number Diff line number Diff line change
Expand Up @@ -84,6 +84,7 @@
from agent_env.task_step.context import TaskStepContext
from agent_env.entity_refs import EntityRef
from agent_env.task_step.task_step import TaskStep, TaskStepDependency
from agent_env.task_step.task_steps.sandbox_utils.sandbox_utils import find_agent_container
from agent_env.task_step.thread_work import finish_on_thread

logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -443,39 +444,6 @@ async def _resolve_live_sandbox(self, provider, sandbox_id: str):
f"a re-run-from-step needs the original live sandbox; re-run the full task."
) from err

async def _discover_container(self, sandbox) -> str:
"""Find the agent container running on the VM.

The A2A agent deploy hardcodes the container name to 'agent-api'
(see agent_env/a2a_agent/a2a_agent.py). We also accept any
'a2a-agent-*' container as a fallback in case the naming scheme
evolves.
"""
exit_code, stdout, stderr = await sandbox.exec_with_output(
"sudo", "docker", "ps",
"--format", "{{.Names}}",
)
if exit_code != 0:
raise RuntimeError(f"Failed to list containers: {stderr[:300]}")
running = [n.strip() for n in stdout.splitlines() if n.strip()]
# Prefer this sandbox's own container name, fall back to the a2a-agent-* prefix.
if sandbox.container_name in running:
return sandbox.container_name
if getattr(sandbox, "owns_container", False):
# Its Docker host is shared, so any other agent container there is another run's.
raise RuntimeError(
f"Agent container {sandbox.container_name!r} is not running. Running containers: {running}."
)
fallback = [n for n in running if n.startswith("a2a-agent-")]
if fallback:
if len(fallback) > 1:
logger.warning(f"Multiple a2a-agent-* containers found; using first: {fallback}")
return fallback[0]
raise RuntimeError(
f"No agent container found on the VM (looked for 'agent-api' or 'a2a-agent-*'). "
f"Running containers: {running}. Has deploy_agent been run in this task?"
)

async def _get_file_size(self, sandbox, container: Optional[str], source_path: str) -> int:
"""Get file size on the agent's filesystem. Returns -1 only if the file is genuinely
absent; a `stat` that fails for any other reason RAISES so a failed size check is a
Expand Down Expand Up @@ -809,7 +777,7 @@ async def _collect_via_agent_container(self, context, store, artifact_id, versio
sandbox = await self._resolve_live_sandbox(provider, agent.sandbox_id)
logger.info(f"Connected to sandbox {agent.sandbox_id} (mode={sandbox.mode})")

container = await self._discover_container(sandbox) if sandbox.mode == SANDBOX_MODE_VM else None
container = await find_agent_container(sandbox) if sandbox.mode == SANDBOX_MODE_VM else None
if container:
logger.info(f"Using agent container: {container}")

Expand Down Expand Up @@ -847,7 +815,7 @@ async def _collect_via_vm_host(self, context, store, artifact_id, version):
async def _collect_via_sandbox_container(self, context, store, artifact_id, version):
"""Collect from a plain container started by ``run_docker_container``.

The agent path finds its container by discovery, and ``_discover_container`` only accepts
The agent path finds its container by discovery, and ``find_agent_container`` only accepts
``agent-api`` / ``a2a-agent-*`` names — so a task that deploys a sandbox and runs an ordinary
image had no way to get its files out. Here the (sandbox, container) pair is named
explicitly and resolved exactly as ``load_artifact`` resolves it.
Expand Down
17 changes: 5 additions & 12 deletions src/agent_env/task_step/task_steps/load_artifact.py
Original file line number Diff line number Diff line change
Expand Up @@ -662,12 +662,9 @@ async def _stage_environment_payload_into_container(
)
else:
await sandbox.exec_script(
f"sudo docker exec -u 0 {shlex.quote(container_name)} mkdir -p {shlex.quote(destination)}"
)
await sandbox.exec_script(
f"sudo docker cp {shlex.quote(vm_stage)}/. "
f"{shlex.quote(container_name)}:{shlex.quote(destination)}"
f"docker exec -u 0 {shlex.quote(container_name)} mkdir -p {shlex.quote(destination)}"
)
await sandbox.docker_cp(f"{vm_stage}/.", f"{container_name}:{destination}")
finally:
try:
await sandbox.exec_script(f"rm -rf {shlex.quote(vm_payload)} {shlex.quote(vm_stage)}")
Expand Down Expand Up @@ -747,7 +744,7 @@ async def _load_universe_into_container(sandbox, container_name: str, universe,
return []

await sandbox.exec_script(
f"sudo docker exec {shlex.quote(container_name)} mkdir -p {shlex.quote(destination)}"
f"docker exec {shlex.quote(container_name)} mkdir -p {shlex.quote(destination)}"
)
loaded: list[str] = []
total = len(file_artifacts)
Expand All @@ -764,12 +761,8 @@ async def _load_universe_into_container(sandbox, container_name: str, universe,
await sandbox.load_s3_file(fa.object_url, vm_temp)
if parent and parent != destination:
await sandbox.exec_script(
f"sudo docker exec {shlex.quote(container_name)} mkdir -p {shlex.quote(parent)}"
f"docker exec {shlex.quote(container_name)} mkdir -p {shlex.quote(parent)}"
)
await sandbox.exec_script(
f"sudo docker cp {shlex.quote(vm_temp)} "
f"{shlex.quote(container_name)}:{shlex.quote(dest_path)} && "
f"rm -f {shlex.quote(vm_temp)}"
)
await sandbox.docker_cp(vm_temp, f"{container_name}:{dest_path}", remove_source=True)
loaded.append(filename)
return loaded
Original file line number Diff line number Diff line change
Expand Up @@ -77,3 +77,28 @@ async def fetch_container_logs(agent, tail: int = 500) -> Optional[str]:
sandbox_id, e, exc_info=True,
)
return None


async def find_agent_container(sandbox) -> str:
"""The agent container running on a VM sandbox: the sandbox's own ``container_name``, else the
first ``a2a-agent-*``. A sandbox that owns its container never takes another: its Docker host
is shared, so any other agent container there is another run's."""
exit_code, stdout, stderr = await sandbox.exec_with_output("sudo", "docker", "ps", "--format", "{{.Names}}")
if exit_code != 0:
raise RuntimeError(f"Failed to list containers: {stderr[:300]}")
running = [n.strip() for n in stdout.splitlines() if n.strip()]
if sandbox.container_name in running:
return sandbox.container_name
if getattr(sandbox, "owns_container", False):
raise RuntimeError(
f"Agent container {sandbox.container_name!r} is not running. Running containers: {running}."
)
fallback = [n for n in running if n.startswith("a2a-agent-")]
if fallback:
if len(fallback) > 1:
logger.warning(f"Multiple a2a-agent-* containers found; using first: {fallback}")
return fallback[0]
raise RuntimeError(
f"No agent container found on the VM (looked for {sandbox.container_name!r} or 'a2a-agent-*'). "
f"Running containers: {running}. Has deploy_agent been run in this task?"
)
Original file line number Diff line number Diff line change
Expand Up @@ -262,7 +262,7 @@ async def execute(self, context: TaskStepContext) -> TaskStepContext:
f"{setup_cmd.splitlines()[0][:140]}"
)
await sandbox.exec_script(
f"sudo docker exec -u {shlex.quote(self.user)} {env_flags} "
f"docker exec -u {shlex.quote(self.user)} {env_flags} "
f"{shlex.quote(self.container_name)} bash -c {shlex.quote(setup_cmd)}"
)

Expand All @@ -272,7 +272,7 @@ async def execute(self, context: TaskStepContext) -> TaskStepContext:
f"in container '{self.container_name}': {command[:160]}"
)
wrapped = (
f"sudo docker exec -u {shlex.quote(self.user)} {env_flags} "
f"docker exec -u {shlex.quote(self.user)} {env_flags} "
f"{shlex.quote(self.container_name)} "
f"timeout --kill-after=10 {self.timeout_sec} bash -c {shlex.quote(command)}"
)
Expand Down Expand Up @@ -428,10 +428,7 @@ def _upload_text_artifact(text: str, artifact_id: str, description: str, s3_url:
async def _extract_file(self, sandbox, path_in_container: str) -> str:
"""`docker cp` a file out of the container, read it from the VM, return text."""
vm_temp = f"/tmp/_verifier_out_{uuid.uuid4().hex[:8]}"
await sandbox.exec_script(
f"sudo docker cp {shlex.quote(self.container_name)}:{shlex.quote(path_in_container)} "
f"{shlex.quote(vm_temp)}"
)
await sandbox.docker_cp(f"{self.container_name}:{path_in_container}", vm_temp)
try:
exit_code, stdout, stderr = await sandbox.exec_with_output("sudo", "cat", vm_temp)
if exit_code != 0:
Expand Down
27 changes: 2 additions & 25 deletions src/agent_env/task_step/task_steps/verifiers/verify_sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -36,6 +36,7 @@
)
from agent_env.task_step.context import TaskStepContext
from agent_env.task_step.task_step import TaskStep, TaskStepDependency
from agent_env.task_step.task_steps.sandbox_utils.sandbox_utils import find_agent_container
from agent_env.task_step.task_steps.verifiers.scoring import ScoreAggregator, aggregate_score

logger = logging.getLogger(__name__)
Expand Down Expand Up @@ -120,7 +121,7 @@ async def execute(self, context: TaskStepContext) -> TaskStepContext:
sandbox = await (self._agent_sandbox(context) if on_agent else self._deployed_sandbox(context))
logger.info(f"Connected to sandbox {sandbox.sandbox_id} (mode={sandbox.mode})")
container = (
await self._discover_container(sandbox)
await find_agent_container(sandbox)
if on_agent and sandbox.mode == SANDBOX_MODE_VM
else None
)
Expand Down Expand Up @@ -287,27 +288,3 @@ async def _eval_shell(
exit_code, _, stderr = await self._exec(sandbox, args)
return _outcome(exit_code == 0, f"exit={exit_code}; stderr={stderr[:200]}")

async def _discover_container(self, sandbox) -> str:
"""Find the agent container on a VM sandbox. Mirrors collect_artifacts._discover_container."""
exit_code, stdout, stderr = await sandbox.exec_with_output(
"sudo", "docker", "ps", "--format", "{{.Names}}",
)
if exit_code != 0:
raise RuntimeError(f"Failed to list containers: {stderr[:300]}")
running = [n.strip() for n in stdout.splitlines() if n.strip()]
if sandbox.container_name in running:
return sandbox.container_name
if getattr(sandbox, "owns_container", False):
# Its Docker host is shared, so any other agent container there is another run's.
raise RuntimeError(
f"Agent container {sandbox.container_name!r} is not running. Running containers: {running}."
)
fallback = [n for n in running if n.startswith("a2a-agent-")]
if fallback:
if len(fallback) > 1:
logger.warning(f"Multiple a2a-agent-* containers found; using first: {fallback}")
return fallback[0]
raise RuntimeError(
f"No agent container found on the VM (looked for 'agent-api' or 'a2a-agent-*'). "
f"Running containers: {running}."
)
Loading
Loading