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
97 changes: 81 additions & 16 deletions src/agent_env/artifact/artifacts/docker_image.py
Original file line number Diff line number Diff line change
Expand Up @@ -8,10 +8,13 @@
import os
import posixpath
import re
import selectors
import shlex
import signal
import subprocess
import tarfile
import tempfile
import time
import uuid
from dataclasses import dataclass
from importlib.metadata import version as pkg_version
Expand All @@ -29,6 +32,9 @@
ProgressCallback = Callable[[str, str, int], None]

_SAFE_GIT_NAME = re.compile(r"^[a-zA-Z0-9._-]+$")
_DOCKER_SAVE_CHUNK_SIZE = 1024 * 1024
_DOCKER_SAVE_STDERR_CHUNK_SIZE = 8192
_DOCKER_SAVE_STDERR_LIMIT_BYTES = 16 * 1024


def _validate_git_name(value: str, label: str) -> str:
Expand Down Expand Up @@ -60,6 +66,80 @@ def _git_clone_commands(owner: str, repo: str, ref: str | None, token: str | Non
f"trap 'rm -f /tmp/git-askpass.sh' EXIT; GIT_ASKPASS=/tmp/git-askpass.sh git clone --depth 1 {ref_flag} {clone_url} /tmp/repo",
]


def _save_image_tar_gz(image_ref: str, output_path: Path, timeout_seconds: float) -> None:
"""Stream docker save through gzip under one deadline without buffering the image."""
deadline = time.monotonic() + timeout_seconds
process = subprocess.Popen(
["docker", "save", image_ref],
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
bufsize=0,
start_new_session=True,
)
stderr_tail = bytearray()
stderr_size = 0
succeeded = False
try:
assert process.stdout is not None
assert process.stderr is not None
with selectors.DefaultSelector() as selector:
for stream, name in ((process.stdout, "stdout"), (process.stderr, "stderr")):
os.set_blocking(stream.fileno(), False)
selector.register(stream, selectors.EVENT_READ, name)
with gzip.open(output_path, "wb") as compressed:
while selector.get_map():
remaining = deadline - time.monotonic()
if remaining <= 0:
raise subprocess.TimeoutExpired(process.args, timeout_seconds)
for key, _ in selector.select(remaining):
chunk_size = (
_DOCKER_SAVE_CHUNK_SIZE
if key.data == "stdout"
else _DOCKER_SAVE_STDERR_CHUNK_SIZE
)
chunk = os.read(key.fd, chunk_size)
if not chunk:
selector.unregister(key.fileobj)
elif key.data == "stdout":
compressed.write(chunk)
else:
stderr_size += len(chunk)
stderr_tail.extend(chunk)
if len(stderr_tail) > _DOCKER_SAVE_STDERR_LIMIT_BYTES:
del stderr_tail[:-_DOCKER_SAVE_STDERR_LIMIT_BYTES]

remaining = deadline - time.monotonic()
if remaining <= 0:
raise subprocess.TimeoutExpired(process.args, timeout_seconds)
process.wait(timeout=remaining)
if process.returncode != 0:
stderr = stderr_tail.decode(errors="replace")
if stderr_size > _DOCKER_SAVE_STDERR_LIMIT_BYTES:
stderr = f"[stderr truncated; showing last {_DOCKER_SAVE_STDERR_LIMIT_BYTES} bytes]\n{stderr}"
raise RuntimeError(f"docker save {image_ref} failed: {stderr}")
succeeded = True
except subprocess.TimeoutExpired as error:
raise RuntimeError(
f"docker save {image_ref} timed out after {timeout_seconds} seconds"
) from error
finally:
if process.returncode is None:
try:
os.killpg(process.pid, signal.SIGKILL)
except ProcessLookupError:
pass
except PermissionError:
process.kill()
if process.stdout is not None:
process.stdout.close()
if process.stderr is not None:
process.stderr.close()
if process.returncode is None:
process.wait()
if not succeeded:
output_path.unlink(missing_ok=True)

class DockerImageArtifact(Artifact):
"""A Docker image artifact stored as tar.gz in S3."""

Expand Down Expand Up @@ -102,22 +182,7 @@ def put(
tmp_path = Path(tmp.name)

try:
save_proc = subprocess.Popen(
["docker", "save", image_ref],
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
)
with gzip.open(tmp_path, "wb") as gz_file:
while chunk := save_proc.stdout.read(8192):
gz_file.write(chunk)
try:
save_proc.wait(timeout=cls.DOCKER_SAVE_TIMEOUT_SECONDS)
except subprocess.TimeoutExpired:
save_proc.kill()
raise RuntimeError(f"docker save {image_ref} timed out after {cls.DOCKER_SAVE_TIMEOUT_SECONDS} seconds")
if save_proc.returncode != 0:
stderr = save_proc.stderr.read().decode() if save_proc.stderr else ""
raise RuntimeError(f"docker save {image_ref} failed: {stderr}")
_save_image_tar_gz(image_ref, tmp_path, cls.DOCKER_SAVE_TIMEOUT_SECONDS)

tar_gz_object_url = objects.put_file_at(
f"{prefix}{fs_safe(id)}-v{version}.tar.gz", str(tmp_path), "application/gzip"
Expand Down
49 changes: 11 additions & 38 deletions src/agent_env/artifact/artifacts/file.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,35 +54,12 @@ def put(
store = get_artifact_store()

filename = os.path.basename(file_path)
content_type, _ = mimetypes.guess_type(file_path)
if content_type is None:
content_type = "application/octet-stream"

version = store.next_version(id)

# Use put_object_file (boto3 managed multipart upload) instead of
# put_object (s3.put_object — hard 5 GB single-object limit). Streams
# directly from disk so large files don't get loaded into RAM either.
s3_url = store.put_object_file(
artifact_type="file",
id=id,
version=version,
object_name=filename,
file_path=file_path,
content_type=content_type,
prefix = store.attempt_prefix("file", id)
return cls.put_at(
id, description=description, file_path=file_path,
object_url=f"{prefix}{filename}", filename=filename,
)

# Store artifact document in MongoDB
instance = cls(
id=id,
version=version,
description=description,
filename=filename,
content_type=content_type,
s3_url=s3_url,
)
return store.put_document(instance)

@classmethod
def put_bytes(
cls,
Expand All @@ -97,24 +74,20 @@ def put_bytes(
from agent_env.artifact.store import get_artifact_store

store = get_artifact_store()
version = store.next_version(id)

s3_url = store.put_object(
artifact_type="file",
id=id,
version=version,
object_name=filename,
data=content,
content_type=content_type,
prefix = store.attempt_prefix("file", id)
object_url = f"{prefix}{filename}"
object_store = get_config().get_object_store_to_write(object_url, id)
stored_url = object_store.put(
object_store.get_object_key(object_url), content, content_type
)

instance = cls(
id=id,
version=version,
version=store.next_version(id),
description=description,
filename=filename,
content_type=content_type,
s3_url=s3_url,
s3_url=stored_url,
)
return store.put_document(instance)

Expand Down
Loading
Loading