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
20 changes: 20 additions & 0 deletions src/agents/extensions/sandbox/blaxel/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -44,6 +44,7 @@
from ....sandbox.manifest import Manifest
from ....sandbox.session import SandboxSession, SandboxSessionState
from ....sandbox.session.base_sandbox_session import BaseSandboxSession
from ....sandbox.session.bounded_read import collect_bounded
from ....sandbox.session.dependencies import Dependencies
from ....sandbox.session.manager import Instrumentation
from ....sandbox.session.pty_output import collect_pty_output
Expand Down Expand Up @@ -506,6 +507,25 @@ async def mkdir(
cause=e,
) from e

async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
workspace_path = await self._validate_path_access(path)
filesystem = self._sandbox.fs
client = filesystem.get_client()
path_arg = filesystem.format_path(sandbox_path_str(workspace_path))
async with client.stream(
"GET",
f"{filesystem.url}/filesystem/{path_arg}",
headers={"Accept": "application/octet-stream"},
) as response:
if response.status_code == 404:
raise WorkspaceReadNotFoundError(path=path)
if response.status_code != 200:
raise WorkspaceArchiveReadError(
path=path,
retryable=True if response.status_code in TRANSIENT_HTTP_STATUS_CODES else None,
)
return await collect_bounded(response.aiter_bytes(chunk_size=65536), max_bytes)

async def read(self, path: Path | str, *, user: str | User | None = None) -> io.IOBase:
error_path = posix_path_as_path(coerce_posix_path(path))
if user is not None:
Expand Down
38 changes: 38 additions & 0 deletions src/agents/extensions/sandbox/cloudflare/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,7 @@
from ....sandbox.manifest import Manifest
from ....sandbox.session import SandboxSession, SandboxSessionState
from ....sandbox.session.base_sandbox_session import BaseSandboxSession
from ....sandbox.session.bounded_read import collect_bounded
from ....sandbox.session.dependencies import Dependencies
from ....sandbox.session.manager import Instrumentation
from ....sandbox.session.mount_lifecycle import (
Expand Down Expand Up @@ -1265,6 +1266,43 @@ async def pty_terminate_all(self) -> None:
for entry in entries:
await self._terminate_pty_entry(entry)

async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
workspace_path = await self._validate_path_access(path)
url_path = quote(sandbox_path_str(workspace_path).lstrip("/"), safe="/")
async with self._session().get(
self._url(f"file/{url_path}"), timeout=self._request_timeout()
) as response:
if response.status == 404:
raise WorkspaceReadNotFoundError(path=path)
if response.status != 200:
raise WorkspaceArchiveReadError(
path=path,
retryable=False
if response.status == 403
else _cloudflare_retryability_for_status(response.status),
)
# Existing Workers return either bytes or an SSE-encoded file. Bound
# the wire representation too, before the existing decoder allocates.
try:
prefix = await response.content.readexactly(7)
except asyncio.IncompleteReadError as error:
prefix = error.partial
if prefix != b"data: {":
if len(prefix) >= max_bytes:
return prefix[:max_bytes]
return prefix + await collect_bounded(
response.content.iter_chunked(65536), max_bytes - len(prefix)
)
wire_limit = 8 * max_bytes + 65536
body = prefix + await collect_bounded(
response.content.iter_chunked(65536), wire_limit + 1 - len(prefix)
)
if len(body) > wire_limit:
raise WorkspaceArchiveReadError(
path=path, context={"reason": "bounded_read_wire_limit"}
)
return self._decode_streamed_payload(body)[:max_bytes]

async def read(self, path: Path | str, *, user: str | User | None = None) -> io.IOBase:
if user is not None:
await self._check_read_with_exec(path, user=user)
Expand Down
27 changes: 27 additions & 0 deletions src/agents/extensions/sandbox/daytona/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@
ExecTransportError,
ExposedPortUnavailableError,
InvalidManifestPathError as InvalidManifestPathError,
SandboxError,
WorkspaceArchiveReadError,
WorkspaceArchiveWriteError,
WorkspaceReadNotFoundError,
Expand All @@ -43,6 +44,7 @@
from ....sandbox.manifest import Manifest
from ....sandbox.session import SandboxSession, SandboxSessionState
from ....sandbox.session.base_sandbox_session import BaseSandboxSession
from ....sandbox.session.bounded_read import collect_bounded
from ....sandbox.session.dependencies import Dependencies
from ....sandbox.session.manager import Instrumentation
from ....sandbox.session.pty_output import collect_pty_output
Expand Down Expand Up @@ -940,6 +942,31 @@ async def _terminate_pty_entry(self, entry: _DaytonaPtySessionEntry) -> None:
except asyncio.TimeoutError:
pass

async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
workspace_path = await self._validate_path_access(path)
# The high-level download buffers the response. The generated toolbox
# client exposes the same endpoint without preloading its body.
try:
response = await self._sandbox.fs._api_client.download_file_without_preload_content(
path=sandbox_path_str(workspace_path),
_request_timeout=float(self.state.timeouts.file_download_s),
)
try:
if response.status == 404:
raise WorkspaceReadNotFoundError(path=path)
if response.status != 200:
raise WorkspaceArchiveReadError(
path=path, retryable=_DAYTONA_HTTP_STATUS_RETRYABLE.get(response.status)
)
return await collect_bounded(response.content.iter_chunked(65536), max_bytes)
finally:
response.close()
except SandboxError:
raise
except Exception as error:
retryable, _ = _daytona_provider_retryability(error)
raise WorkspaceArchiveReadError(path=path, retryable=retryable) from None

async def read(self, path: Path | str, *, user: str | User | None = None) -> io.IOBase:
error_path = posix_path_as_path(coerce_posix_path(path))
if user is not None:
Expand Down
15 changes: 15 additions & 0 deletions src/agents/extensions/sandbox/e2b/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,7 @@
from ....sandbox.manifest import Manifest
from ....sandbox.session import SandboxSession, SandboxSessionState
from ....sandbox.session.base_sandbox_session import BaseSandboxSession
from ....sandbox.session.bounded_read import collect_bounded
from ....sandbox.session.dependencies import Dependencies
from ....sandbox.session.manager import Instrumentation
from ....sandbox.session.pty_output import collect_pty_output
Expand Down Expand Up @@ -1114,6 +1115,20 @@ async def pty_terminate_all(self) -> None:
for entry in entries:
await self._terminate_pty_entry(entry)

async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
workspace_path = await self._validate_path_access(path)
try:
stream = await _sandbox_read_file(
self._sandbox, sandbox_path_str(workspace_path), format="stream"
)
async with cast(Any, stream) as chunks:
return await collect_bounded(chunks, max_bytes)
except _e2b_not_found_error_types():
raise WorkspaceReadNotFoundError(path=path) from None
Comment thread
seratch marked this conversation as resolved.
except Exception as error:
retryable, _ = _e2b_provider_retryability(error)
raise WorkspaceArchiveReadError(path=path, retryable=retryable) from None

async def read(self, path: Path, *, user: str | User | None = None) -> io.IOBase:
if user is not None:
await self._check_read_with_exec(path, user=user)
Expand Down
33 changes: 33 additions & 0 deletions src/agents/extensions/sandbox/modal/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -1248,6 +1248,39 @@ async def _terminate_pty_entry(self, entry: _ModalPtyProcessEntry) -> None:
return_exceptions=True,
)

async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
try:
workspace_path = await self._validate_path_access(path)
await self._ensure_sandbox()
assert self._sandbox is not None
# Each read starts a remote operation. Use Modal's 100 MiB per-read
# ceiling while respecting the caller's remaining byte budget.
stream = await self._sandbox.open.aio(sandbox_path_str(workspace_path), "rb")
completed = False
try:
result = bytearray()
while len(result) < max_bytes:
chunk = await stream.read.aio(min(100 * 1024 * 1024, max_bytes - len(result)))
if not chunk:
break
result.extend(chunk)
payload = bytes(result)
completed = True
return payload
finally:
try:
await asyncio.wait_for(stream.close.aio(), timeout=5.0)
except Exception:
# Preserve an active read failure or cancellation. A close failure
# still fails a read that would otherwise have completed.
if completed:
raise
except (FileNotFoundError, SandboxError):
raise
except Exception as error:
Comment thread
jbeckwith-oai marked this conversation as resolved.
retryable, _ = _modal_provider_retryability(error)
raise WorkspaceArchiveReadError(path=path, retryable=retryable) from None

async def read(self, path: Path, *, user: str | User | None = None) -> io.IOBase:
if user is not None:
await self._check_read_with_exec(path, user=user)
Expand Down
19 changes: 19 additions & 0 deletions src/agents/extensions/sandbox/runloop/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -48,6 +48,7 @@
from ....sandbox.manifest import Manifest
from ....sandbox.session import SandboxSession, SandboxSessionState
from ....sandbox.session.base_sandbox_session import BaseSandboxSession
from ....sandbox.session.bounded_read import collect_bounded
from ....sandbox.session.dependencies import Dependencies
from ....sandbox.session.manager import Instrumentation
from ....sandbox.session.runtime_helpers import RESOLVE_WORKSPACE_PATH_HELPER, RuntimeHelperScript
Expand Down Expand Up @@ -955,6 +956,24 @@ async def _resolve_exposed_port(self, port: int) -> ExposedPortEndpoint:
cause=e,
) from e

async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
normalized_path = await self._validate_path_access(path)
try:
async with self._sdk.api.devboxes.with_streaming_response.download_file(
self.devbox_id,
path=sandbox_path_str(normalized_path),
timeout=self.state.timeouts.file_download_s,
) as response:
return await collect_bounded(response.iter_bytes(chunk_size=65536), max_bytes)
except Exception as error:
if _is_runloop_not_found(error):
raise WorkspaceReadNotFoundError(path=path) from None
if _is_runloop_provider_error(error):
raise WorkspaceArchiveReadError(
path=path, retryable=_runloop_provider_retryability(error)
) from None
raise

async def read(self, path: Path | str, *, user: str | User | None = None) -> io.IOBase:
"""Read a file via Runloop's binary file API."""
error_path = posix_path_as_path(coerce_posix_path(path))
Expand Down
30 changes: 30 additions & 0 deletions src/agents/extensions/sandbox/vercel/sandbox.py
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@
from ....sandbox.materialization import MaterializationResult
from ....sandbox.session import SandboxSession, SandboxSessionState, manifest_ops
from ....sandbox.session.base_sandbox_session import BaseSandboxSession
from ....sandbox.session.bounded_read import collect_bounded
from ....sandbox.session.dependencies import Dependencies
from ....sandbox.session.manager import Instrumentation
from ....sandbox.session.mount_lifecycle import (
Expand Down Expand Up @@ -1239,6 +1240,35 @@ async def _resolve_exposed_port(self, port: int) -> ExposedPortEndpoint:
tls=tls,
)

@redact_mount_error_data
async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
async with self._s3_mount_operation():
normalized_path = await self._validate_path_access(path)
sandbox = await self._ensure_sandbox()
try:
chunks = await sandbox.iter_file(
sandbox_path_str(normalized_path), chunk_size=65536
)
completed = False
try:
payload = await collect_bounded(chunks, max_bytes)
completed = True
return payload
finally:
try:
await chunks.aclose()
except Exception:
# Preserve the primary read failure or cancellation.
# A close failure is primary only after a successful read.
if completed:
raise
except vercel_sandbox.SandboxNotFoundError:
raise WorkspaceReadNotFoundError(path=path) from None
except Exception as error:
raise WorkspaceArchiveReadError(
path=path, retryable=_vercel_provider_retryability(error)
) from None

@redact_mount_error_data
async def read(self, path: Path, *, user: str | User | None = None) -> io.IOBase:
async with self._s3_mount_operation():
Expand Down
13 changes: 13 additions & 0 deletions src/agents/sandbox/sandboxes/_unix_local_file_ops.py
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,19 @@ def read(self, path: Path) -> io.IOBase:
os.close(fd)
raise

def read_bounded(self, path: Path, max_bytes: int) -> bytes:
with self.parent(path) as (parent_fd, name):
fd = os.open(name, os.O_RDONLY | os.O_NOFOLLOW | os.O_NONBLOCK, dir_fd=parent_fd)
try:
if not stat.S_ISREG(os.fstat(fd).st_mode):
raise OSError("Bounded reads require a regular file")
stream = os.fdopen(fd, "rb")
except BaseException:
os.close(fd)
raise
with stream:
return stream.read(max_bytes)

def write(self, path: Path, stream: io.IOBase) -> None:
with self.parent(path, for_write=True, create_parents=True) as (parent_fd, name):
fd = os.open(
Expand Down
22 changes: 22 additions & 0 deletions src/agents/sandbox/sandboxes/docker.py
Original file line number Diff line number Diff line change
Expand Up @@ -52,6 +52,7 @@
MountConfigError,
WorkspaceArchiveReadError,
WorkspaceArchiveWriteError,
WorkspaceReadNotFoundError,
)
from ..manifest import Manifest
from ..session import SandboxSession, SandboxSessionState
Expand Down Expand Up @@ -836,6 +837,27 @@ async def _prepare_user_pty_pid_path(self, *, path: Path, user: str | None) -> N
error_path=path,
)

async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
workspace_path = await self._validate_path_access(path)
# Docker writes already require POSIX sh and head -c. Restrict this
# internal read to image-owned utilities, independent of manifest PATH.
result = await self.exec(
"/bin/sh",
"-c",
"PATH=/usr/bin:/bin; export PATH; "
'[ -e "$1" ] || exit 44; [ -f "$1" ] || exit 45; head -c "$2" < "$1"',
"sh",
sandbox_path_str(workspace_path),
str(max_bytes),
shell=False,
timeout=30.0,
)
if result.exit_code == 44:
raise WorkspaceReadNotFoundError(path=path)
if not result.ok():
raise WorkspaceArchiveReadError(path=path)
return result.stdout

async def read(self, path: Path, *, user: str | User | None = None) -> io.IOBase:
workspace_path = await self._validate_path_access(path)

Expand Down
3 changes: 3 additions & 0 deletions src/agents/sandbox/sandboxes/unix_local.py
Original file line number Diff line number Diff line change
Expand Up @@ -1030,6 +1030,9 @@ async def rm(
except OSError as e:
raise WorkspaceArchiveWriteError(path=normalized, cause=e) from e

async def _read_bounded(self, path: Path, *, max_bytes: int) -> bytes:
return self._files.read_bounded(self.normalize_path(path), max_bytes)

async def read(self, path: Path, *, user: str | User | None = None) -> io.IOBase:
if user is not None:
await self._check_read_with_exec(path, user=user)
Expand Down
Loading
Loading