From 23b5e287ee182994aa60a69e2491bbcb45882217 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 02:02:26 +0000 Subject: [PATCH 01/18] [skip ci] apcha/agents-files-artifacts From 0ffdf87e7b4be79c585500e172f4f72f2bbfcd97 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 02:06:02 +0000 Subject: [PATCH 02/18] feat(agents): prepare selected files and download turn artifacts --- helpers.md | 32 +++ src/openai/lib/beta/agents/__init__.py | 10 + src/openai/lib/beta/agents/_artifacts.py | 77 +++++ src/openai/lib/beta/agents/_files.py | 216 ++++++++++++++ .../beta/agents/environments/files.py | 86 +++++- .../beta/agents/sessions/artifacts.py | 10 + tests/lib/streaming/agents/test_artifacts.py | 122 ++++++++ tests/lib/streaming/agents/test_files.py | 270 ++++++++++++++++++ 8 files changed, 821 insertions(+), 2 deletions(-) create mode 100644 src/openai/lib/beta/agents/_artifacts.py create mode 100644 src/openai/lib/beta/agents/_files.py create mode 100644 tests/lib/streaming/agents/test_artifacts.py create mode 100644 tests/lib/streaming/agents/test_files.py diff --git a/helpers.md b/helpers.md index 54902af655..46ff7c90e5 100644 --- a/helpers.md +++ b/helpers.md @@ -633,3 +633,35 @@ This only selects the local parser; it does not change the session's schema. `result.parse(Report)` parses an existing raw result. `AgentOutputParseError.result` retains the completed raw answer if validation fails. With `AsyncOpenAI`, await creation and the result getter, and use `async with`. + +### Stage files and download a turn artifact + +```python +from pathlib import Path + +prepared = client.beta.agents.environments.files.prepare({ + "/workspace/source.pdf": Path("source.pdf"), +}) +session = client.beta.agents.sessions.create( + agent={"model": MODEL}, + environment={"type": "openai_hosted", "files": prepared.files}, +) +with client.beta.agents.sessions.stream( + session.id, input="Read source.pdf and write /workspace/outputs/report.md." +) as stream: + result = stream.get_final_result() + +artifact = client.beta.agents.sessions.artifacts.for_result(result).download( + "/workspace/outputs/report.md", to=Path("downloaded-report.md") +) +``` + +Use `prepare_directory("docs", destination="/workspace/docs", include=["**/*.md"])` +for a selected directory snapshot, or `files.upload(environment_id, file=Path(...), +path="/workspace/source.pdf")` to stage a file in an existing environment. +Uploads remain caller-owned: use `prepared.uploaded_file_ids` with the ordinary +Files API when ready to delete them. Preparation errors expose partial uploads +through `error.prepared`. A batch of multiple uploads cannot share one explicit +`Idempotency-Key`. + +With `AsyncOpenAI`, await preparation, staging, and artifact downloads. diff --git a/src/openai/lib/beta/agents/__init__.py b/src/openai/lib/beta/agents/__init__.py index 6b004a47e7..125ae0018d 100644 --- a/src/openai/lib/beta/agents/__init__.py +++ b/src/openai/lib/beta/agents/__init__.py @@ -5,6 +5,12 @@ function_tool as function_tool, pydantic_function_tool as pydantic_function_tool, ) +from ._files import ( + StagedAgentFile as StagedAgentFile, + PreparedAgentFiles as PreparedAgentFiles, + AgentFileStagingError as AgentFileStagingError, + AgentFilePreparationError as AgentFilePreparationError, +) from ._output import agent_text_format as agent_text_format from ._result import ( AgentTurnResult as AgentTurnResult, @@ -15,3 +21,7 @@ AgentSessionEventStream as AgentSessionEventStream, AsyncAgentSessionEventStream as AsyncAgentSessionEventStream, ) +from ._artifacts import ( + AgentResultArtifacts as AgentResultArtifacts, + AsyncAgentResultArtifacts as AsyncAgentResultArtifacts, +) diff --git a/src/openai/lib/beta/agents/_artifacts.py b/src/openai/lib/beta/agents/_artifacts.py new file mode 100644 index 0000000000..485f13bd1b --- /dev/null +++ b/src/openai/lib/beta/agents/_artifacts.py @@ -0,0 +1,77 @@ +from __future__ import annotations + +from os import PathLike +from typing import TYPE_CHECKING + +import httpx2 + +from ._result import AgentTurnResult +from ...._types import Headers, NotGiven, not_given +from ....types.beta.agents.sessions.session_artifact import SessionArtifact + +if TYPE_CHECKING: + from ....resources.beta.agents.sessions.artifacts import Artifacts, AsyncArtifacts + + +class AgentResultArtifacts: + """Beta: look up immutable artifacts from one completed turn.""" + + def __init__(self, resource: Artifacts, result: AgentTurnResult) -> None: + self._resource = resource + self._session_id = result.session_id + self._turn_id = result.turn_id + + def download( + self, + path: str, + *, + to: str | PathLike[str], + extra_headers: Headers | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> SessionArtifact: + """Download one exact turn/path to an explicit local destination.""" + selected: SessionArtifact | None = None + for artifact in self._resource.list(self._session_id, extra_headers=extra_headers, timeout=timeout): + if artifact.session_id == self._session_id and artifact.turn_id == self._turn_id and artifact.path == path: + if selected is not None: + raise ValueError("More than one artifact matches this turn and path") + selected = artifact + if selected is None: + raise ValueError("No artifact matches this turn and path") + with self._resource.with_streaming_response.content( + selected.id, session_id=self._session_id, extra_headers=extra_headers, timeout=timeout + ) as content: + content.stream_to_file(to) + return selected + + +class AsyncAgentResultArtifacts: + """Beta: async artifact downloads from one completed turn.""" + + def __init__(self, resource: AsyncArtifacts, result: AgentTurnResult) -> None: + self._resource = resource + self._session_id = result.session_id + self._turn_id = result.turn_id + + async def download( + self, + path: str, + *, + to: str | PathLike[str], + extra_headers: Headers | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> SessionArtifact: + """Download one exact turn/path to an explicit local destination.""" + selected: SessionArtifact | None = None + async for artifact in self._resource.list(self._session_id, extra_headers=extra_headers, timeout=timeout): + if artifact.session_id == self._session_id and artifact.turn_id == self._turn_id and artifact.path == path: + if selected is not None: + raise ValueError("More than one artifact matches this turn and path") + selected = artifact + if selected is None: + raise ValueError("No artifact matches this turn and path") + async with self._resource.with_streaming_response.content( + selected.id, session_id=self._session_id, extra_headers=extra_headers, timeout=timeout + ) as content: + await content.stream_to_file(to) + return selected diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py new file mode 100644 index 0000000000..cfe8f05ad0 --- /dev/null +++ b/src/openai/lib/beta/agents/_files.py @@ -0,0 +1,216 @@ +from __future__ import annotations + +from os import PathLike, walk +from typing import TYPE_CHECKING, Mapping, Sequence +from asyncio import CancelledError +from fnmatch import fnmatchcase +from pathlib import Path, PurePosixPath +from dataclasses import field, dataclass + +from ...._types import Omit, Headers, FileTypes +from ...._exceptions import OpenAIError +from ....types.beta.hosted_environment_file_param import HostedEnvironmentFileParam +from ....types.beta.agents.environments.environment_file import EnvironmentFile + +if TYPE_CHECKING: + from ...streaming.agents._streams import _RequestOptions + from ....resources.beta.agents.environments.files import Files, AsyncFiles + +_MAX_FILES = 50 +_MAX_BYTES = 50 * 1024 * 1024 + + +@dataclass +class PreparedAgentFiles: + """Beta: file inputs and upload ownership for explicit Files API cleanup.""" + + files: list[HostedEnvironmentFileParam] = field(default_factory=list[HostedEnvironmentFileParam]) + uploaded_file_ids: list[str] = field(default_factory=list[str]) + + +@dataclass(frozen=True) +class StagedAgentFile: + """Beta: a staged environment file and its separately owned Files API upload.""" + + uploaded_file_id: str + file: EnvironmentFile + + +class AgentFilePreparationError(OpenAIError): + """Beta: preparation failed; successful uploads remain available for cleanup.""" + + def __init__(self, prepared: PreparedAgentFiles) -> None: + super().__init__("Could not prepare all agent files; successful uploads remain caller-owned") + self.prepared = prepared + + +class AgentFileStagingError(OpenAIError): + """Beta: staging failed after upload; the upload remains caller-owned.""" + + def __init__(self, uploaded_file_id: str) -> None: + super().__init__("Could not stage the agent file; the upload remains caller-owned") + self.uploaded_file_id = uploaded_file_id + + +def _destination(value: str) -> str: + components = value.split("/")[1:] + if not value.startswith("/workspace/") or "\x00" in value or "\\" in value: + raise ValueError("Agent file destinations must be absolute POSIX file paths under /workspace") + if any(component in ("", ".", "..") for component in components): + raise ValueError("Agent file destinations cannot contain empty, dot or parent components") + if ( + components[1] in (".codex", ".managed-agents") + or components[1].startswith(".managed-agents-") + or value == "/workspace/outputs" + ): + raise ValueError("Agent file destination uses a reserved environment path") + return value + + +def _local_file(value: str | PathLike[str]) -> Path: + path = Path(value).absolute() + if path.is_symlink(): + raise ValueError("Selected agent files must not be symlinks") + if not path.is_file(): + raise ValueError("Selected agent files must be regular files") + if path.stat().st_size > _MAX_BYTES: + raise ValueError("An agent file must not exceed 50 MiB") + return path.resolve() + + +def _prepare_selection( + files: Mapping[str, str | PathLike[str]], options: _RequestOptions, defaults: Headers +) -> list[tuple[str, Path]]: + if len(files) > _MAX_FILES: + raise ValueError("Initial agent files must not exceed 50 files") + headers = {key.lower(): value for key, value in defaults.items()} + headers.update({key.lower(): value for key, value in (options["extra_headers"] or {}).items()}) + if len(files) > 1 and "idempotency-key" in headers and not isinstance(headers["idempotency-key"], Omit): + raise ValueError("One Idempotency-Key cannot be reused for multiple file uploads") + selected: list[tuple[str, Path]] = [] + destinations: set[str] = set() + size = 0 + for destination, source in files.items(): + destination = _destination(destination) + if any( + destination == prior or destination.startswith(prior + "/") or prior.startswith(destination + "/") + for prior in destinations + ): + raise ValueError("Selected agent files contain duplicate or conflicting destinations") + destinations.add(destination) + local = _local_file(source) + size += local.stat().st_size + selected.append((destination, local)) + if size > _MAX_BYTES: + raise ValueError("Initial agent files must not exceed 50 MiB in total") + return selected + + +def _matches(parts: tuple[str, ...], pattern: tuple[str, ...]) -> bool: + if not pattern: + return not parts + if pattern[0] == "**": + return _matches(parts, pattern[1:]) or bool(parts) and _matches(parts[1:], pattern) + return bool(parts) and fnmatchcase(parts[0], pattern[0]) and _matches(parts[1:], pattern[1:]) + + +def directory_files(root: str | PathLike[str], destination: str, include: Sequence[str]) -> dict[str, Path]: + directory = Path(root).absolute() + if isinstance(include, str) or not include: + raise ValueError("Select directory files with explicit include patterns") + if directory.is_symlink() or not directory.is_dir(): + raise ValueError("The selected directory must be a directory, not a symlink") + directory = directory.resolve() + # System aliases such as macOS /tmp are allowed above the chosen root. + _destination(destination + "/selected-file") + patterns: list[tuple[str, ...]] = [] + for pattern in include: + if not pattern or PurePosixPath(pattern).is_absolute() or ".." in PurePosixPath(pattern).parts: + raise ValueError("Include patterns must stay inside the selected directory") + patterns.append(PurePosixPath(pattern).parts) + selected: dict[str, Path] = {} + for current, directories, names in walk(directory, followlinks=False): + directories.sort() + for name in sorted([*directories, *names]): + source = Path(current) / name + relative = source.relative_to(directory) + chosen = any(_matches(relative.parts, pattern) for pattern in patterns) + if source.is_symlink(): + if chosen: + raise ValueError("Selected agent files must not be symlinks") + if name in directories: + directories.remove(name) + continue + if chosen and name not in directories: + selected[str(PurePosixPath(destination) / relative.as_posix())] = source + if not selected: + raise ValueError("The include patterns did not select any files") + return selected + + +def prepare(resource: Files, files: Mapping[str, str | PathLike[str]], options: _RequestOptions) -> PreparedAgentFiles: + selected = _prepare_selection(files, options, resource._client.default_headers) + prepared = PreparedAgentFiles() + try: + for destination, source in selected: + uploaded = resource._client.files.create(file=source, purpose="user_data", **options) + prepared.uploaded_file_ids.append(uploaded.id) + prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) + except Exception as error: + raise AgentFilePreparationError(prepared) from error + return prepared + + +async def async_prepare( + resource: AsyncFiles, files: Mapping[str, str | PathLike[str]], options: _RequestOptions +) -> PreparedAgentFiles: + selected = _prepare_selection(files, options, resource._client.default_headers) + prepared = PreparedAgentFiles() + try: + for destination, source in selected: + uploaded = await resource._client.files.create(file=source, purpose="user_data", **options) + prepared.uploaded_file_ids.append(uploaded.id) + prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) + except CancelledError as error: + error.__dict__["prepared"] = prepared + raise + except Exception as error: + raise AgentFilePreparationError(prepared) from error + return prepared + + +def _upload_file(file: FileTypes, path: str) -> str: + destination = _destination(path) + content = file[1] if isinstance(file, tuple) else file + if isinstance(content, PathLike): + _local_file(content) + elif isinstance(content, bytes) and len(content) > _MAX_BYTES: + raise ValueError("An agent file must not exceed 50 MiB") + return destination + + +def upload( + resource: Files, environment_id: str, file: FileTypes, path: str, options: _RequestOptions +) -> StagedAgentFile: + destination = _upload_file(file, path) + uploaded = resource._client.files.create(file=file, purpose="user_data", **options) + try: + staged = resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) + except Exception as error: + raise AgentFileStagingError(uploaded.id) from error + return StagedAgentFile(uploaded_file_id=uploaded.id, file=staged) + + +async def async_upload( + resource: AsyncFiles, environment_id: str, file: FileTypes, path: str, options: _RequestOptions +) -> StagedAgentFile: + destination = _upload_file(file, path) + uploaded = await resource._client.files.create(file=file, purpose="user_data", **options) + try: + staged = await resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) + except CancelledError as error: + error.__dict__["uploaded_file_id"] = uploaded.id + raise + except Exception as error: + raise AgentFileStagingError(uploaded.id) from error + return StagedAgentFile(uploaded_file_id=uploaded.id, file=staged) diff --git a/src/openai/resources/beta/agents/environments/files.py b/src/openai/resources/beta/agents/environments/files.py index b2ad8913bf..e3c03fe63e 100644 --- a/src/openai/resources/beta/agents/environments/files.py +++ b/src/openai/resources/beta/agents/environments/files.py @@ -2,19 +2,29 @@ from __future__ import annotations -from typing import Optional +from os import PathLike +from typing import Mapping, Optional, Sequence from typing_extensions import Literal, overload import httpx2 from ..... import _legacy_response -from ....._types import Body, Omit, Query, Headers, NotGiven, omit, not_given +from ....._types import Body, Omit, Query, Headers, NotGiven, FileTypes, omit, not_given from ....._utils import path_template, required_args, maybe_transform, async_maybe_transform from ....._compat import cached_property from ....._resource import SyncAPIResource, AsyncAPIResource from ....._response import to_streamed_response_wrapper, async_to_streamed_response_wrapper from .....pagination import SyncTokenPage, AsyncTokenPage from ....._base_client import AsyncPaginator, make_request_options +from .....lib.beta.agents._files import ( + StagedAgentFile, + PreparedAgentFiles, + upload, + prepare, + async_upload, + async_prepare, + directory_files, +) from .....types.beta.agents.environments import file_list_params, file_create_params from .....types.beta.agents.environments.environment_file import EnvironmentFile @@ -22,6 +32,40 @@ class Files(SyncAPIResource): + def prepare( + self, + files: Mapping[str, str | PathLike[str]], + *, + extra_headers: Headers | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> PreparedAgentFiles: + """Beta: upload selected files for a new hosted environment.""" + return prepare(self, files, {"extra_headers": extra_headers, "timeout": timeout}) + + def prepare_directory( + self, + root: str | PathLike[str], + *, + destination: str, + include: Sequence[str], + extra_headers: Headers | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> PreparedAgentFiles: + """Beta: prepare an explicitly selected directory snapshot.""" + return self.prepare(directory_files(root, destination, include), extra_headers=extra_headers, timeout=timeout) + + def upload( + self, + environment_id: str, + *, + file: FileTypes, + path: str, + extra_headers: Headers | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> StagedAgentFile: + """Beta: upload and stage one file in an existing environment.""" + return upload(self, environment_id, file, path, {"extra_headers": extra_headers, "timeout": timeout}) + @cached_property def with_raw_response(self) -> FilesWithRawResponse: """ @@ -222,6 +266,44 @@ def list( class AsyncFiles(AsyncAPIResource): + async def prepare( + self, + files: Mapping[str, str | PathLike[str]], + *, + extra_headers: Headers | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> PreparedAgentFiles: + """Beta: upload selected files for a new hosted environment.""" + return await async_prepare(self, files, {"extra_headers": extra_headers, "timeout": timeout}) + + async def prepare_directory( + self, + root: str | PathLike[str], + *, + destination: str, + include: Sequence[str], + extra_headers: Headers | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> PreparedAgentFiles: + """Beta: prepare an explicitly selected directory snapshot.""" + return await self.prepare( + directory_files(root, destination, include), extra_headers=extra_headers, timeout=timeout + ) + + async def upload( + self, + environment_id: str, + *, + file: FileTypes, + path: str, + extra_headers: Headers | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> StagedAgentFile: + """Beta: upload and stage one file in an existing environment.""" + return await async_upload( + self, environment_id, file, path, {"extra_headers": extra_headers, "timeout": timeout} + ) + @cached_property def with_raw_response(self) -> AsyncFilesWithRawResponse: """ diff --git a/src/openai/resources/beta/agents/sessions/artifacts.py b/src/openai/resources/beta/agents/sessions/artifacts.py index 60e6d067ad..f5500de551 100644 --- a/src/openai/resources/beta/agents/sessions/artifacts.py +++ b/src/openai/resources/beta/agents/sessions/artifacts.py @@ -22,6 +22,8 @@ ) from .....pagination import SyncCursorPage, AsyncCursorPage from ....._base_client import AsyncPaginator, make_request_options +from .....lib.beta.agents._result import AgentTurnResult +from .....lib.beta.agents._artifacts import AgentResultArtifacts, AsyncAgentResultArtifacts from .....types.beta.agents.sessions import artifact_list_params from .....types.beta.agents.sessions.session_artifact import SessionArtifact from .....types.beta.agents.sessions.session_artifact_deleted import SessionArtifactDeleted @@ -30,6 +32,10 @@ class Artifacts(SyncAPIResource): + def for_result(self, result: AgentTurnResult) -> AgentResultArtifacts: + """Beta: download artifacts scoped to the exact completed result turn.""" + return AgentResultArtifacts(self, result) + @cached_property def with_raw_response(self) -> ArtifactsWithRawResponse: """ @@ -254,6 +260,10 @@ def content( class AsyncArtifacts(AsyncAPIResource): + def for_result(self, result: AgentTurnResult) -> AsyncAgentResultArtifacts: + """Beta: download artifacts scoped to the exact completed result turn.""" + return AsyncAgentResultArtifacts(self, result) + @cached_property def with_raw_response(self) -> AsyncArtifactsWithRawResponse: """ diff --git a/tests/lib/streaming/agents/test_artifacts.py b/tests/lib/streaming/agents/test_artifacts.py new file mode 100644 index 0000000000..44c0bfddfc --- /dev/null +++ b/tests/lib/streaming/agents/test_artifacts.py @@ -0,0 +1,122 @@ +from __future__ import annotations + +from typing import Any, Iterator, AsyncIterator +from pathlib import Path +from typing_extensions import override + +import httpx2 +import pytest + +from openai import OpenAI, AsyncOpenAI +from openai._compat import model_parse +from openai.lib.beta.agents import AgentTurnResult +from openai.types.beta.agents.sessions.turn import Turn +from tests.lib.streaming.agents.test_streams import Server, sdk as sdk, turn_event + + +class ArtifactBody(httpx2.SyncByteStream, httpx2.AsyncByteStream): + def __init__(self) -> None: + self.closed = False + self.chunks = 0 + + @override + def __iter__(self) -> Iterator[bytes]: + for value in (b"report", b"-", b"done"): + self.chunks += 1 + yield value + + @override + async def __aiter__(self) -> AsyncIterator[bytes]: + for value in self: + yield value + + @override + def close(self) -> None: + self.closed = True + + @override + async def aclose(self) -> None: + self.closed = True + + +def artifact(id: str, turn_id: str = "turn_root") -> dict[str, Any]: + return { + "id": id, + "created_at": 1, + "environment_id": "env_test", + "object": "agent.session.artifact", + "path": "/workspace/outputs/report.md", + "session_id": "session_test", + "size_bytes": 11, + "turn_id": turn_id, + } + + +class ArtifactServer(Server): + def __init__(self) -> None: + super().__init__() + self.artifacts = [artifact("old", "old_turn"), artifact("selected")] + self.content = ArtifactBody() + + @override + def handle(self, request: httpx2.Request) -> httpx2.Response: + self.requests.append(request) + if request.url.path.endswith("/content"): + assert request.url.path.endswith("/selected/content") + return httpx2.Response(200, headers={"content-type": "application/octet-stream"}, stream=self.content) + after = request.url.params.get("after") + offset = next((index + 1 for index, value in enumerate(self.artifacts) if value["id"] == after), 0) + return httpx2.Response( + 200, + json={ + "object": "list", + "data": self.artifacts[offset : offset + 1], + "has_more": offset + 1 < len(self.artifacts), + }, + ) + + +@pytest.fixture +def server() -> ArtifactServer: + return ArtifactServer() + + +def result() -> AgentTurnResult: + return AgentTurnResult(turn=model_parse(Turn, turn_event("completed")["turn"]), messages=[]) + + +async def download(sdk: OpenAI | AsyncOpenAI, path: Path) -> Any: + if isinstance(sdk, AsyncOpenAI): + return await sdk.beta.agents.sessions.artifacts.for_result(result()).download( + "/workspace/outputs/report.md", to=path, extra_headers={"X-Synthetic": "test"}, timeout=7 + ) + return sdk.beta.agents.sessions.artifacts.for_result(result()).download( + "/workspace/outputs/report.md", to=path, extra_headers={"X-Synthetic": "test"}, timeout=7 + ) + + +async def test_exact_turn_artifact_pagination_streaming_and_options( + sdk: OpenAI | AsyncOpenAI, server: ArtifactServer, tmp_path: Path +) -> None: + path = tmp_path / "chosen.md" + selected = await download(sdk, path) + assert selected.id == "selected" + assert path.read_bytes() == b"report-done" + assert len(server.requests) == 3 + assert server.content.chunks == 3 + assert server.content.closed + assert all(request.headers["X-Synthetic"] == "test" for request in server.requests) + assert all(request.extensions["timeout"]["read"] == 7 for request in server.requests) + + +@pytest.mark.parametrize("ambiguous", [False, True]) +async def test_missing_or_ambiguous_artifact_does_not_open_destination( + sdk: OpenAI | AsyncOpenAI, server: ArtifactServer, tmp_path: Path, ambiguous: bool +) -> None: + server.artifacts = [artifact("selected"), artifact("duplicate")] if ambiguous else [artifact("old", "old_turn")] + path = tmp_path / "chosen.md" + path.write_text("keep") + with pytest.raises(ValueError, match="More than one" if ambiguous else "No artifact"): + await download(sdk, path) + assert path.read_text() == "keep" + assert server.content.chunks == 0 diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py new file mode 100644 index 0000000000..82e383507f --- /dev/null +++ b/tests/lib/streaming/agents/test_files.py @@ -0,0 +1,270 @@ +from __future__ import annotations + +import json +from typing import Any +from asyncio import CancelledError +from pathlib import Path +from typing_extensions import override + +import httpx2 +import pytest + +from openai import OpenAI, AsyncOpenAI +from openai.lib.beta.agents import AgentFileStagingError, AgentFilePreparationError +from tests.lib.streaming.agents.test_streams import Server, sdk as sdk + + +class FilesServer(Server): + def __init__(self) -> None: + super().__init__() + self.uploads = 0 + self.fail_upload = 0 + self.fail_stage = False + self.cancel_upload = 0 + self.cancel_stage = False + + @override + def handle(self, request: httpx2.Request) -> httpx2.Response: + self.requests.append(request) + if request.url.path == "/v1/files": + self.uploads += 1 + if self.uploads == self.cancel_upload: + raise CancelledError() + if self.uploads == self.fail_upload: + return httpx2.Response(500, json={"error": {"message": "Synthetic upload failure"}}) + return httpx2.Response( + 200, + json={ + "id": f"file_{self.uploads}", + "bytes": 3, + "created_at": 1, + "filename": "synthetic.txt", + "object": "file", + "purpose": "user_data", + "status": "processed", + }, + ) + assert request.url.path == "/v1/agents/environments/env_test/files" + if self.cancel_stage: + raise CancelledError() + if self.fail_stage: + return httpx2.Response(500, json={"error": {"message": "Synthetic staging failure"}}) + body = json.loads(request.content) + return httpx2.Response( + 200, + json={ + "id": "envfile_test", + "object": "agent.environment.file", + "environment_id": "env_test", + "path": body["path"], + "size_bytes": 3, + "created_at": 1, + }, + ) + + +@pytest.fixture +def server() -> FilesServer: + return FilesServer() + + +async def prepare(sdk: OpenAI | AsyncOpenAI, files: dict[str, Path], **options: Any) -> Any: + if isinstance(sdk, AsyncOpenAI): + return await sdk.beta.agents.environments.files.prepare(files, **options) + return sdk.beta.agents.environments.files.prepare(files, **options) + + +async def test_prepare_uploads_selected_files_and_preserves_options( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + source = tmp_path / "source.txt" + source.write_text("abc") + result = await prepare(sdk, {"/workspace/source.txt": source}, extra_headers={"X-Synthetic": "test"}, timeout=7) + assert result.files == [{"type": "file_id", "file_id": "file_1", "path": "/workspace/source.txt"}] + assert result.uploaded_file_ids == ["file_1"] + assert all(request.headers["X-Synthetic"] == "test" for request in server.requests) + assert b"abc" in server.requests[0].content + + +async def test_preflight_all_files_before_upload( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + source = tmp_path / "source.txt" + source.write_text("abc") + with pytest.raises(ValueError, match="regular files"): + await prepare(sdk, {"/workspace/source.txt": source, "/workspace/missing.txt": tmp_path / "missing"}) + assert server.uploads == 0 + with pytest.raises(ValueError, match="conflicting"): + await prepare(sdk, {"/workspace/source.txt": source, "/workspace/source.txt/child": source}) + assert server.uploads == 0 + with pytest.raises(ValueError, match="Idempotency-Key"): + await prepare( + sdk, {"/workspace/a": source, "/workspace/b": source}, extra_headers={"idempotency-key": "synthetic"} + ) + assert server.uploads == 0 + + +async def test_partial_upload_ownership_is_available( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + source = tmp_path / "source.txt" + source.write_text("abc") + server.fail_upload = 2 + with pytest.raises(AgentFilePreparationError) as caught: + await prepare(sdk, {"/workspace/a": source, "/workspace/b": source}) + assert vars(caught.value)["prepared"].uploaded_file_ids == ["file_1"] + assert caught.value.prepared.files[0]["path"] == "/workspace/a" + assert all(request.method == "POST" for request in server.requests) + + +async def test_directory_selection_and_symlink_preflight( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + (tmp_path / "keep.md").write_text("abc") + (tmp_path / "omit.txt").write_text("secret") + if isinstance(sdk, AsyncOpenAI): + result = await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["*.md"] + ) + else: + result = sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["*.md"] + ) + assert result.files[0]["path"] == "/workspace/docs/keep.md" + assert server.uploads == 1 + (tmp_path / "link.md").symlink_to(tmp_path / "omit.txt") + with pytest.raises(ValueError, match="symlink"): + if isinstance(sdk, AsyncOpenAI): + await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["*.md"] + ) + else: + sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["*.md"] + ) + assert server.uploads == 1 + + +async def test_live_staging_and_failure_ownership(sdk: OpenAI | AsyncOpenAI, server: FilesServer) -> None: + async def stage() -> Any: + if isinstance(sdk, AsyncOpenAI): + return await sdk.beta.agents.environments.files.upload( + "env_test", file=("source.txt", b"abc"), path="/workspace/source.txt" + ) + return sdk.beta.agents.environments.files.upload( + "env_test", file=("source.txt", b"abc"), path="/workspace/source.txt" + ) + + result = await stage() + assert result.uploaded_file_id == "file_1" + assert result.file.path == "/workspace/source.txt" + server.fail_stage = True + with pytest.raises(AgentFileStagingError) as caught: + await stage() + assert vars(caught.value)["uploaded_file_id"] == "file_2" + assert all(request.method == "POST" for request in server.requests) + + +async def test_count_and_aggregate_size_limits_precede_upload( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + source = tmp_path / "large" + with source.open("wb") as content: + content.truncate(26 * 1024 * 1024) + with pytest.raises(ValueError, match="in total"): + await prepare(sdk, {"/workspace/a": source, "/workspace/b": source}) + with pytest.raises(ValueError, match="50 files"): + await prepare(sdk, {f"/workspace/{index}": source for index in range(51)}) + with source.open("wb") as content: + content.truncate(50 * 1024 * 1024 + 1) + with pytest.raises(ValueError, match="50 MiB"): + await prepare(sdk, {"/workspace/a": source}) + assert server.uploads == 0 + + +async def test_async_cancellation_retains_only_observed_uploads( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + if not isinstance(sdk, AsyncOpenAI): + pytest.skip("Async cancellation contract") + source = tmp_path / "source.txt" + source.write_text("abc") + server.cancel_upload = 2 + with pytest.raises(CancelledError) as caught: + await sdk.beta.agents.environments.files.prepare({"/workspace/a": source, "/workspace/b": source}) + assert vars(caught.value)["prepared"].uploaded_file_ids == ["file_1"] + assert all(request.method == "POST" for request in server.requests) + + +async def test_client_default_idempotency_key_is_rejected_for_batch( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + source = tmp_path / "source.txt" + source.write_text("abc") + sdk._custom_headers = {"Idempotency-Key": "synthetic"} + with pytest.raises(ValueError, match="Idempotency-Key"): + await prepare(sdk, {"/workspace/a": source, "/workspace/b": source}) + assert server.uploads == 0 + + +@pytest.mark.parametrize( + "destination", + [ + "/tmp/source", + "/workspace/../source", + "/workspace//source", + "/workspace/./source", + "/workspace/source/", + "/workspace/a\\b", + "/workspace/\x00", + "/workspace/.codex/source", + "/workspace/.managed-agents/source", + "/workspace/.managed-agents-internal/source", + "/workspace/outputs", + ], +) +async def test_invalid_hosted_destination_fails_before_upload( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, destination: str +) -> None: + source = tmp_path / "source.txt" + source.write_text("abc") + with pytest.raises(ValueError): + await prepare(sdk, {destination: source}) + assert server.uploads == 0 + + +async def test_cancelled_staging_preserves_observed_upload_id(sdk: OpenAI | AsyncOpenAI, server: FilesServer) -> None: + if not isinstance(sdk, AsyncOpenAI): + pytest.skip("Async cancellation contract") + server.cancel_stage = True + with pytest.raises(CancelledError) as caught: + await sdk.beta.agents.environments.files.upload( + "env_test", file=("source.txt", b"abc"), path="/workspace/source.txt" + ) + assert vars(caught.value)["uploaded_file_id"] == "file_1" + assert all(request.method == "POST" for request in server.requests) + + +async def test_system_alias_ancestor_is_allowed_but_directory_entries_are_not_followed( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + actual = tmp_path / "actual" + root = actual / "chosen" + root.mkdir(parents=True) + (root / "keep.md").write_text("abc") + outside = tmp_path / "outside" + outside.mkdir() + (outside / "private.md").write_text("private") + (root / "linked").symlink_to(outside, target_is_directory=True) + alias = tmp_path / "system-alias" + alias.symlink_to(actual, target_is_directory=True) + if isinstance(sdk, AsyncOpenAI): + prepared = await sdk.beta.agents.environments.files.prepare_directory( + alias / "chosen", destination="/workspace/docs", include=["**/*.md"] + ) + else: + prepared = sdk.beta.agents.environments.files.prepare_directory( + alias / "chosen", destination="/workspace/docs", include=["**/*.md"] + ) + assert [item["path"] for item in prepared.files] == ["/workspace/docs/keep.md"] + assert server.uploads == 1 From 0beb3b40427ea6b18d927c9964f59c5a57a321b5 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 02:16:10 +0000 Subject: [PATCH 03/18] fix(agents): retain upload sources and standard request options --- src/openai/lib/beta/agents/_artifacts.py | 36 +++- src/openai/lib/beta/agents/_files.py | 156 ++++++++++++------ .../beta/agents/environments/files.py | 52 +++++- tests/lib/streaming/agents/test_artifacts.py | 15 +- tests/lib/streaming/agents/test_files.py | 71 +++++++- 5 files changed, 267 insertions(+), 63 deletions(-) diff --git a/src/openai/lib/beta/agents/_artifacts.py b/src/openai/lib/beta/agents/_artifacts.py index 485f13bd1b..7ca4696e2c 100644 --- a/src/openai/lib/beta/agents/_artifacts.py +++ b/src/openai/lib/beta/agents/_artifacts.py @@ -6,7 +6,7 @@ import httpx2 from ._result import AgentTurnResult -from ...._types import Headers, NotGiven, not_given +from ...._types import Body, Query, Headers, NotGiven, not_given from ....types.beta.agents.sessions.session_artifact import SessionArtifact if TYPE_CHECKING: @@ -27,11 +27,19 @@ def download( *, to: str | PathLike[str], extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> SessionArtifact: """Download one exact turn/path to an explicit local destination.""" selected: SessionArtifact | None = None - for artifact in self._resource.list(self._session_id, extra_headers=extra_headers, timeout=timeout): + for artifact in self._resource.list( + self._session_id, + extra_headers=extra_headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, + ): if artifact.session_id == self._session_id and artifact.turn_id == self._turn_id and artifact.path == path: if selected is not None: raise ValueError("More than one artifact matches this turn and path") @@ -39,7 +47,12 @@ def download( if selected is None: raise ValueError("No artifact matches this turn and path") with self._resource.with_streaming_response.content( - selected.id, session_id=self._session_id, extra_headers=extra_headers, timeout=timeout + selected.id, + session_id=self._session_id, + extra_headers=extra_headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, ) as content: content.stream_to_file(to) return selected @@ -59,11 +72,19 @@ async def download( *, to: str | PathLike[str], extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> SessionArtifact: """Download one exact turn/path to an explicit local destination.""" selected: SessionArtifact | None = None - async for artifact in self._resource.list(self._session_id, extra_headers=extra_headers, timeout=timeout): + async for artifact in self._resource.list( + self._session_id, + extra_headers=extra_headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, + ): if artifact.session_id == self._session_id and artifact.turn_id == self._turn_id and artifact.path == path: if selected is not None: raise ValueError("More than one artifact matches this turn and path") @@ -71,7 +92,12 @@ async def download( if selected is None: raise ValueError("No artifact matches this turn and path") async with self._resource.with_streaming_response.content( - selected.id, session_id=self._session_id, extra_headers=extra_headers, timeout=timeout + selected.id, + session_id=self._session_id, + extra_headers=extra_headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, ) as content: await content.stream_to_file(to) return selected diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index cfe8f05ad0..d7fbe7fc26 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -1,21 +1,34 @@ from __future__ import annotations -from os import PathLike, walk -from typing import TYPE_CHECKING, Mapping, Sequence +from io import IOBase +from os import PathLike, walk, fstat +from stat import S_ISREG +from typing import TYPE_CHECKING, Mapping, BinaryIO, Sequence, Generator, cast from asyncio import CancelledError from fnmatch import fnmatchcase from pathlib import Path, PurePosixPath +from contextlib import ExitStack, contextmanager from dataclasses import field, dataclass +from typing_extensions import TypedDict -from ...._types import Omit, Headers, FileTypes +import httpx2 + +from ...._types import Body, Omit, Query, Headers, NotGiven, FileTypes from ...._exceptions import OpenAIError from ....types.beta.hosted_environment_file_param import HostedEnvironmentFileParam from ....types.beta.agents.environments.environment_file import EnvironmentFile if TYPE_CHECKING: - from ...streaming.agents._streams import _RequestOptions from ....resources.beta.agents.environments.files import Files, AsyncFiles + +class _RequestOptions(TypedDict): + extra_headers: Headers | None + extra_query: Query | None + extra_body: Body | None + timeout: float | httpx2.Timeout | None | NotGiven + + _MAX_FILES = 50 _MAX_BYTES = 50 * 1024 * 1024 @@ -64,6 +77,8 @@ def _destination(value: str) -> str: or value == "/workspace/outputs" ): raise ValueError("Agent file destination uses a reserved environment path") + if len(value) > 4096: + raise ValueError("Agent file destinations must not exceed 4096 characters") return value @@ -75,19 +90,19 @@ def _local_file(value: str | PathLike[str]) -> Path: raise ValueError("Selected agent files must be regular files") if path.stat().st_size > _MAX_BYTES: raise ValueError("An agent file must not exceed 50 MiB") - return path.resolve() + return path def _prepare_selection( - files: Mapping[str, str | PathLike[str]], options: _RequestOptions, defaults: Headers -) -> list[tuple[str, Path]]: + files: Mapping[str, str | PathLike[str]], options: _RequestOptions, defaults: Headers, stack: ExitStack +) -> list[tuple[str, Path, BinaryIO, int]]: if len(files) > _MAX_FILES: raise ValueError("Initial agent files must not exceed 50 files") headers = {key.lower(): value for key, value in defaults.items()} headers.update({key.lower(): value for key, value in (options["extra_headers"] or {}).items()}) if len(files) > 1 and "idempotency-key" in headers and not isinstance(headers["idempotency-key"], Omit): raise ValueError("One Idempotency-Key cannot be reused for multiple file uploads") - selected: list[tuple[str, Path]] = [] + selected: list[tuple[str, Path, BinaryIO, int]] = [] destinations: set[str] = set() size = 0 for destination, source in files.items(): @@ -99,8 +114,9 @@ def _prepare_selection( raise ValueError("Selected agent files contain duplicate or conflicting destinations") destinations.add(destination) local = _local_file(source) - size += local.stat().st_size - selected.append((destination, local)) + handle, length = _open_local(local, stack) + size += length + selected.append((destination, local, handle, length)) if size > _MAX_BYTES: raise ValueError("Initial agent files must not exceed 50 MiB in total") return selected @@ -148,52 +164,99 @@ def directory_files(root: str | PathLike[str], destination: str, include: Sequen return selected +def _open_local(path: Path, stack: ExitStack) -> tuple[BinaryIO, int]: + before = path.lstat() + if not S_ISREG(before.st_mode): + raise ValueError("Selected agent files must be regular files, not symlinks") + handle = stack.enter_context(path.open("rb")) + opened = fstat(handle.fileno()) + if (before.st_dev, before.st_ino, before.st_size) != (opened.st_dev, opened.st_ino, opened.st_size): + raise ValueError("Selected agent file changed during preparation") + if opened.st_size > _MAX_BYTES: + raise ValueError("An agent file must not exceed 50 MiB") + return handle, opened.st_size + + +def _unchanged_size(handle: BinaryIO, length: int) -> None: + if fstat(handle.fileno()).st_size != length: + raise ValueError("Selected agent file changed during preparation") + + def prepare(resource: Files, files: Mapping[str, str | PathLike[str]], options: _RequestOptions) -> PreparedAgentFiles: - selected = _prepare_selection(files, options, resource._client.default_headers) - prepared = PreparedAgentFiles() - try: - for destination, source in selected: - uploaded = resource._client.files.create(file=source, purpose="user_data", **options) - prepared.uploaded_file_ids.append(uploaded.id) - prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) - except Exception as error: - raise AgentFilePreparationError(prepared) from error - return prepared + with ExitStack() as stack: + selected = _prepare_selection(files, options, resource._client.default_headers, stack) + prepared = PreparedAgentFiles() + try: + for destination, source, handle, length in selected: + _unchanged_size(handle, length) + uploaded = resource._client.files.create(file=(source.name, handle), purpose="user_data", **options) + prepared.uploaded_file_ids.append(uploaded.id) + prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) + except Exception as error: + raise AgentFilePreparationError(prepared) from error + return prepared async def async_prepare( resource: AsyncFiles, files: Mapping[str, str | PathLike[str]], options: _RequestOptions ) -> PreparedAgentFiles: - selected = _prepare_selection(files, options, resource._client.default_headers) - prepared = PreparedAgentFiles() - try: - for destination, source in selected: - uploaded = await resource._client.files.create(file=source, purpose="user_data", **options) - prepared.uploaded_file_ids.append(uploaded.id) - prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) - except CancelledError as error: - error.__dict__["prepared"] = prepared - raise - except Exception as error: - raise AgentFilePreparationError(prepared) from error - return prepared - - -def _upload_file(file: FileTypes, path: str) -> str: - destination = _destination(path) + with ExitStack() as stack: + selected = _prepare_selection(files, options, resource._client.default_headers, stack) + prepared = PreparedAgentFiles() + try: + for destination, source, handle, length in selected: + _unchanged_size(handle, length) + uploaded = await resource._client.files.create( + file=(source.name, handle), purpose="user_data", **options + ) + prepared.uploaded_file_ids.append(uploaded.id) + prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) + except CancelledError as error: + error.__dict__["prepared"] = prepared + raise + except Exception as error: + raise AgentFilePreparationError(prepared) from error + return prepared + + +@contextmanager +def _upload_content(file: FileTypes) -> Generator[FileTypes, None, None]: content = file[1] if isinstance(file, tuple) else file - if isinstance(content, PathLike): - _local_file(content) - elif isinstance(content, bytes) and len(content) > _MAX_BYTES: - raise ValueError("An agent file must not exceed 50 MiB") - return destination + with ExitStack() as stack: + if isinstance(content, PathLike): + path = _local_file(content) + handle, length = _open_local(path, stack) + _unchanged_size(handle, length) + yield cast(FileTypes, (file[0], handle, *file[2:])) if isinstance(file, tuple) else (path.name, handle) + return + size = None + if isinstance(content, (bytes, str)): + size = len(content.encode("utf-8") if isinstance(content, str) else content) + elif isinstance(content, IOBase): + try: + if content.seekable(): + position = content.tell() + try: + size = content.seek(0, 2) + finally: + content.seek(position) + else: + metadata = fstat(content.fileno()) + if S_ISREG(metadata.st_mode): + size = metadata.st_size + except (OSError, ValueError): + pass + if size is not None and size > _MAX_BYTES: + raise ValueError("An agent file must not exceed 50 MiB") + yield file def upload( resource: Files, environment_id: str, file: FileTypes, path: str, options: _RequestOptions ) -> StagedAgentFile: - destination = _upload_file(file, path) - uploaded = resource._client.files.create(file=file, purpose="user_data", **options) + destination = _destination(path) + with _upload_content(file) as content: + uploaded = resource._client.files.create(file=content, purpose="user_data", **options) try: staged = resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) except Exception as error: @@ -204,8 +267,9 @@ def upload( async def async_upload( resource: AsyncFiles, environment_id: str, file: FileTypes, path: str, options: _RequestOptions ) -> StagedAgentFile: - destination = _upload_file(file, path) - uploaded = await resource._client.files.create(file=file, purpose="user_data", **options) + destination = _destination(path) + with _upload_content(file) as content: + uploaded = await resource._client.files.create(file=content, purpose="user_data", **options) try: staged = await resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) except CancelledError as error: diff --git a/src/openai/resources/beta/agents/environments/files.py b/src/openai/resources/beta/agents/environments/files.py index e3c03fe63e..0d59fe0a74 100644 --- a/src/openai/resources/beta/agents/environments/files.py +++ b/src/openai/resources/beta/agents/environments/files.py @@ -37,10 +37,16 @@ def prepare( files: Mapping[str, str | PathLike[str]], *, extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> PreparedAgentFiles: """Beta: upload selected files for a new hosted environment.""" - return prepare(self, files, {"extra_headers": extra_headers, "timeout": timeout}) + return prepare( + self, + files, + {"extra_headers": extra_headers, "extra_query": extra_query, "extra_body": extra_body, "timeout": timeout}, + ) def prepare_directory( self, @@ -49,10 +55,18 @@ def prepare_directory( destination: str, include: Sequence[str], extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> PreparedAgentFiles: """Beta: prepare an explicitly selected directory snapshot.""" - return self.prepare(directory_files(root, destination, include), extra_headers=extra_headers, timeout=timeout) + return self.prepare( + directory_files(root, destination, include), + extra_headers=extra_headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, + ) def upload( self, @@ -61,10 +75,18 @@ def upload( file: FileTypes, path: str, extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> StagedAgentFile: """Beta: upload and stage one file in an existing environment.""" - return upload(self, environment_id, file, path, {"extra_headers": extra_headers, "timeout": timeout}) + return upload( + self, + environment_id, + file, + path, + {"extra_headers": extra_headers, "extra_query": extra_query, "extra_body": extra_body, "timeout": timeout}, + ) @cached_property def with_raw_response(self) -> FilesWithRawResponse: @@ -271,10 +293,16 @@ async def prepare( files: Mapping[str, str | PathLike[str]], *, extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> PreparedAgentFiles: """Beta: upload selected files for a new hosted environment.""" - return await async_prepare(self, files, {"extra_headers": extra_headers, "timeout": timeout}) + return await async_prepare( + self, + files, + {"extra_headers": extra_headers, "extra_query": extra_query, "extra_body": extra_body, "timeout": timeout}, + ) async def prepare_directory( self, @@ -283,11 +311,17 @@ async def prepare_directory( destination: str, include: Sequence[str], extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> PreparedAgentFiles: """Beta: prepare an explicitly selected directory snapshot.""" return await self.prepare( - directory_files(root, destination, include), extra_headers=extra_headers, timeout=timeout + directory_files(root, destination, include), + extra_headers=extra_headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, ) async def upload( @@ -297,11 +331,17 @@ async def upload( file: FileTypes, path: str, extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> StagedAgentFile: """Beta: upload and stage one file in an existing environment.""" return await async_upload( - self, environment_id, file, path, {"extra_headers": extra_headers, "timeout": timeout} + self, + environment_id, + file, + path, + {"extra_headers": extra_headers, "extra_query": extra_query, "extra_body": extra_body, "timeout": timeout}, ) @cached_property diff --git a/tests/lib/streaming/agents/test_artifacts.py b/tests/lib/streaming/agents/test_artifacts.py index 44c0bfddfc..9b02e68c5c 100644 --- a/tests/lib/streaming/agents/test_artifacts.py +++ b/tests/lib/streaming/agents/test_artifacts.py @@ -88,10 +88,20 @@ def result() -> AgentTurnResult: async def download(sdk: OpenAI | AsyncOpenAI, path: Path) -> Any: if isinstance(sdk, AsyncOpenAI): return await sdk.beta.agents.sessions.artifacts.for_result(result()).download( - "/workspace/outputs/report.md", to=path, extra_headers={"X-Synthetic": "test"}, timeout=7 + "/workspace/outputs/report.md", + to=path, + extra_headers={"X-Synthetic": "test"}, + extra_query={"synthetic": "query"}, + extra_body={"synthetic": "body"}, + timeout=7, ) return sdk.beta.agents.sessions.artifacts.for_result(result()).download( - "/workspace/outputs/report.md", to=path, extra_headers={"X-Synthetic": "test"}, timeout=7 + "/workspace/outputs/report.md", + to=path, + extra_headers={"X-Synthetic": "test"}, + extra_query={"synthetic": "query"}, + extra_body={"synthetic": "body"}, + timeout=7, ) @@ -106,6 +116,7 @@ async def test_exact_turn_artifact_pagination_streaming_and_options( assert server.content.chunks == 3 assert server.content.closed assert all(request.headers["X-Synthetic"] == "test" for request in server.requests) + assert all(request.url.params["synthetic"] == "query" for request in server.requests) assert all(request.extensions["timeout"]["read"] == 7 for request in server.requests) diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 82e383507f..6ad5271de3 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -1,7 +1,8 @@ from __future__ import annotations import json -from typing import Any +from io import BytesIO +from typing import Any, Callable from asyncio import CancelledError from pathlib import Path from typing_extensions import override @@ -22,12 +23,15 @@ def __init__(self) -> None: self.fail_stage = False self.cancel_upload = 0 self.cancel_stage = False + self.after_upload: Callable[[], None] | None = None @override def handle(self, request: httpx2.Request) -> httpx2.Response: self.requests.append(request) if request.url.path == "/v1/files": self.uploads += 1 + if self.after_upload is not None: + self.after_upload() if self.uploads == self.cancel_upload: raise CancelledError() if self.uploads == self.fail_upload: @@ -79,11 +83,20 @@ async def test_prepare_uploads_selected_files_and_preserves_options( ) -> None: source = tmp_path / "source.txt" source.write_text("abc") - result = await prepare(sdk, {"/workspace/source.txt": source}, extra_headers={"X-Synthetic": "test"}, timeout=7) + result = await prepare( + sdk, + {"/workspace/source.txt": source}, + extra_headers={"X-Synthetic": "test"}, + extra_query={"synthetic": "query"}, + extra_body={"synthetic": "body"}, + timeout=7, + ) assert result.files == [{"type": "file_id", "file_id": "file_1", "path": "/workspace/source.txt"}] assert result.uploaded_file_ids == ["file_1"] assert all(request.headers["X-Synthetic"] == "test" for request in server.requests) assert b"abc" in server.requests[0].content + assert server.requests[0].url.params["synthetic"] == "query" + assert b"body" in server.requests[0].content async def test_preflight_all_files_before_upload( @@ -149,15 +162,25 @@ async def test_live_staging_and_failure_ownership(sdk: OpenAI | AsyncOpenAI, ser async def stage() -> Any: if isinstance(sdk, AsyncOpenAI): return await sdk.beta.agents.environments.files.upload( - "env_test", file=("source.txt", b"abc"), path="/workspace/source.txt" + "env_test", + file=("source.txt", b"abc"), + path="/workspace/source.txt", + extra_query={"synthetic": "query"}, + extra_body={"synthetic": "body"}, ) return sdk.beta.agents.environments.files.upload( - "env_test", file=("source.txt", b"abc"), path="/workspace/source.txt" + "env_test", + file=("source.txt", b"abc"), + path="/workspace/source.txt", + extra_query={"synthetic": "query"}, + extra_body={"synthetic": "body"}, ) result = await stage() assert result.uploaded_file_id == "file_1" assert result.file.path == "/workspace/source.txt" + assert all(request.url.params["synthetic"] == "query" for request in server.requests) + assert all(b"body" in request.content for request in server.requests) server.fail_stage = True with pytest.raises(AgentFileStagingError) as caught: await stage() @@ -221,6 +244,7 @@ async def test_client_default_idempotency_key_is_rejected_for_batch( "/workspace/.managed-agents/source", "/workspace/.managed-agents-internal/source", "/workspace/outputs", + "/workspace/" + "x" * 4096, ], ) async def test_invalid_hosted_destination_fails_before_upload( @@ -268,3 +292,42 @@ async def test_system_alias_ancestor_is_allowed_but_directory_entries_are_not_fo ) assert [item["path"] for item in prepared.files] == ["/workspace/docs/keep.md"] assert server.uploads == 1 + + +async def test_prepared_files_keep_opened_sources_during_batch_upload( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + first = tmp_path / "first.txt" + second = tmp_path / "second.txt" + first.write_text("first-original") + second.write_text("second-original") + + def replace_later_source() -> None: + if server.uploads == 1: + second.unlink() + second.write_text("replacement") + + server.after_upload = replace_later_source + await prepare(sdk, {"/workspace/first": first, "/workspace/second": second}) + assert b"second-original" in server.requests[1].content + assert b"replacement" not in server.requests[1].content + + +@pytest.mark.parametrize("file_backed", [False, True]) +async def test_known_stream_sizes_fail_before_upload_without_moving_position( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, file_backed: bool +) -> None: + content = (tmp_path / "large").open("w+b") if file_backed else BytesIO() + with content: + content.seek(50 * 1024 * 1024) + content.write(b"x") + content.seek(3) + with pytest.raises(ValueError, match="50 MiB"): + if isinstance(sdk, AsyncOpenAI): + await sdk.beta.agents.environments.files.upload( + "env_test", file=("large", content), path="/workspace/large" + ) + else: + sdk.beta.agents.environments.files.upload("env_test", file=("large", content), path="/workspace/large") + assert content.tell() == 3 + assert not server.requests From 6a548dedfa6d842c6aebea9d9fea757bea26e2e9 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 02:19:50 +0000 Subject: [PATCH 04/18] fix(agents): preserve async cleanup and offload directory selection --- src/openai/lib/beta/agents/_files.py | 19 +++-- .../beta/agents/environments/files.py | 3 +- tests/lib/streaming/agents/test_files.py | 82 +++++++++++++++++++ 3 files changed, 98 insertions(+), 6 deletions(-) diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index d7fbe7fc26..4278f4f676 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -4,14 +4,15 @@ from os import PathLike, walk, fstat from stat import S_ISREG from typing import TYPE_CHECKING, Mapping, BinaryIO, Sequence, Generator, cast -from asyncio import CancelledError from fnmatch import fnmatchcase from pathlib import Path, PurePosixPath from contextlib import ExitStack, contextmanager from dataclasses import field, dataclass from typing_extensions import TypedDict +import anyio import httpx2 +from anyio.to_thread import run_sync from ...._types import Body, Omit, Query, Headers, NotGiven, FileTypes from ...._exceptions import OpenAIError @@ -164,6 +165,10 @@ def directory_files(root: str | PathLike[str], destination: str, include: Sequen return selected +async def async_directory_files(root: str | PathLike[str], destination: str, include: Sequence[str]) -> dict[str, Path]: + return await run_sync(directory_files, root, destination, include) + + def _open_local(path: Path, stack: ExitStack) -> tuple[BinaryIO, int]: before = path.lstat() if not S_ISREG(before.st_mode): @@ -211,7 +216,7 @@ async def async_prepare( ) prepared.uploaded_file_ids.append(uploaded.id) prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) - except CancelledError as error: + except anyio.get_cancelled_exc_class() as error: error.__dict__["prepared"] = prepared raise except Exception as error: @@ -230,8 +235,8 @@ def _upload_content(file: FileTypes) -> Generator[FileTypes, None, None]: yield cast(FileTypes, (file[0], handle, *file[2:])) if isinstance(file, tuple) else (path.name, handle) return size = None - if isinstance(content, (bytes, str)): - size = len(content.encode("utf-8") if isinstance(content, str) else content) + if isinstance(content, bytes): + size = len(content) elif isinstance(content, IOBase): try: if content.seekable(): @@ -254,6 +259,8 @@ def _upload_content(file: FileTypes) -> Generator[FileTypes, None, None]: def upload( resource: Files, environment_id: str, file: FileTypes, path: str, options: _RequestOptions ) -> StagedAgentFile: + if not environment_id: + raise ValueError("Expected a non-empty environment_id") destination = _destination(path) with _upload_content(file) as content: uploaded = resource._client.files.create(file=content, purpose="user_data", **options) @@ -267,12 +274,14 @@ def upload( async def async_upload( resource: AsyncFiles, environment_id: str, file: FileTypes, path: str, options: _RequestOptions ) -> StagedAgentFile: + if not environment_id: + raise ValueError("Expected a non-empty environment_id") destination = _destination(path) with _upload_content(file) as content: uploaded = await resource._client.files.create(file=content, purpose="user_data", **options) try: staged = await resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) - except CancelledError as error: + except anyio.get_cancelled_exc_class() as error: error.__dict__["uploaded_file_id"] = uploaded.id raise except Exception as error: diff --git a/src/openai/resources/beta/agents/environments/files.py b/src/openai/resources/beta/agents/environments/files.py index 0d59fe0a74..4c35468df7 100644 --- a/src/openai/resources/beta/agents/environments/files.py +++ b/src/openai/resources/beta/agents/environments/files.py @@ -24,6 +24,7 @@ async_upload, async_prepare, directory_files, + async_directory_files, ) from .....types.beta.agents.environments import file_list_params, file_create_params from .....types.beta.agents.environments.environment_file import EnvironmentFile @@ -317,7 +318,7 @@ async def prepare_directory( ) -> PreparedAgentFiles: """Beta: prepare an explicitly selected directory snapshot.""" return await self.prepare( - directory_files(root, destination, include), + await async_directory_files(root, destination, include), extra_headers=extra_headers, extra_query=extra_query, extra_body=extra_body, diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 6ad5271de3..60dee43e16 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -331,3 +331,85 @@ async def test_known_stream_sizes_fail_before_upload_without_moving_position( sdk.beta.agents.environments.files.upload("env_test", file=("large", content), path="/workspace/large") assert content.tell() == 3 assert not server.requests + + +async def test_empty_environment_id_fails_before_upload(sdk: OpenAI | AsyncOpenAI, server: FilesServer) -> None: + with pytest.raises(ValueError, match="environment_id"): + if isinstance(sdk, AsyncOpenAI): + await sdk.beta.agents.environments.files.upload( + "", file=("source.txt", b"abc"), path="/workspace/source.txt" + ) + else: + sdk.beta.agents.environments.files.upload("", file=("source.txt", b"abc"), path="/workspace/source.txt") + assert not server.requests + + +@pytest.mark.parametrize("staging", [False, True]) +def test_trio_cancellation_preserves_observed_uploads(tmp_path: Path, staging: bool) -> None: + import trio + import anyio + + source = tmp_path / "source.txt" + source.write_text("abc") + server = FilesServer() + captured: list[BaseException] = [] + + async def run() -> None: + with trio.CancelScope() as scope: + + async def handle(request: httpx2.Request) -> httpx2.Response: + should_cancel = request.url.path.endswith("env_test/files") if staging else server.uploads == 1 + if should_cancel: + scope.cancel() + await trio.lowlevel.checkpoint() + return server.handle(request) + + async with AsyncOpenAI( + api_key="synthetic", + base_url="https://sdk-test.example/v1", + max_retries=0, + http_client=httpx2.AsyncClient(transport=httpx2.MockTransport(handle), trust_env=False), + ) as client: + try: + if staging: + await client.beta.agents.environments.files.upload( + "env_test", file=source, path="/workspace/source.txt" + ) + else: + await client.beta.agents.environments.files.prepare( + {"/workspace/a": source, "/workspace/b": source} + ) + except trio.Cancelled as error: + captured.append(error) + raise + + anyio.run(run, backend="trio") + assert len(captured) == 1 + if staging: + assert vars(captured[0])["uploaded_file_id"] == "file_1" + else: + assert vars(captured[0])["prepared"].uploaded_file_ids == ["file_1"] + + +async def test_async_directory_selection_runs_off_event_loop( + sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + import threading + + from openai.lib.beta.agents import _files + + if not isinstance(sdk, AsyncOpenAI): + pytest.skip("Async traversal contract") + (tmp_path / "source.txt").write_text("abc") + current_thread = threading.get_ident() + original = _files.directory_files + + def select(root: Any, destination: str, include: Any) -> Any: + assert threading.get_ident() != current_thread + return original(root, destination, include) + + monkeypatch.setattr(_files, "directory_files", select) + prepared = await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["*.txt"] + ) + assert prepared.uploaded_file_ids == ["file_1"] From c6a4ceaf18a3bf9a6fef7ae501dc809001f4885d Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 02:32:57 +0000 Subject: [PATCH 05/18] fix(agents): preserve selected file identity and async responsiveness --- helpers.md | 4 +- src/openai/lib/beta/agents/_files.py | 81 +++++++++-- tests/lib/streaming/agents/test_files.py | 177 +++++++++++++++++++++++ 3 files changed, 247 insertions(+), 15 deletions(-) diff --git a/helpers.md b/helpers.md index 46ff7c90e5..f103a1b216 100644 --- a/helpers.md +++ b/helpers.md @@ -664,4 +664,6 @@ Files API when ready to delete them. Preparation errors expose partial uploads through `error.prepared`. A batch of multiple uploads cannot share one explicit `Idempotency-Key`. -With `AsyncOpenAI`, await preparation, staging, and artifact downloads. +Use application-controlled local paths and stable source directories. These +convenience helpers are not a filesystem sandbox for untrusted paths or hostile +local writers. With `AsyncOpenAI`, await preparation, staging, and artifact downloads. diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index 4278f4f676..067b58ce72 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -8,7 +8,7 @@ from pathlib import Path, PurePosixPath from contextlib import ExitStack, contextmanager from dataclasses import field, dataclass -from typing_extensions import TypedDict +from typing_extensions import TypedDict, override import anyio import httpx2 @@ -34,6 +34,16 @@ class _RequestOptions(TypedDict): _MAX_BYTES = 50 * 1024 * 1024 +@dataclass(frozen=True) +class _SelectedFile(PathLike[str]): + path: Path + identity: tuple[int, int, int] + + @override + def __fspath__(self) -> str: + return str(self.path) + + @dataclass class PreparedAgentFiles: """Beta: file inputs and upload ownership for explicit Files API cleanup.""" @@ -115,7 +125,7 @@ def _prepare_selection( raise ValueError("Selected agent files contain duplicate or conflicting destinations") destinations.add(destination) local = _local_file(source) - handle, length = _open_local(local, stack) + handle, length = _open_local(local, stack, source.identity if isinstance(source, _SelectedFile) else None) size += length selected.append((destination, local, handle, length)) if size > _MAX_BYTES: @@ -131,7 +141,11 @@ def _matches(parts: tuple[str, ...], pattern: tuple[str, ...]) -> bool: return bool(parts) and fnmatchcase(parts[0], pattern[0]) and _matches(parts[1:], pattern[1:]) -def directory_files(root: str | PathLike[str], destination: str, include: Sequence[str]) -> dict[str, Path]: +def _walk_error(error: OSError) -> None: + raise error + + +def directory_files(root: str | PathLike[str], destination: str, include: Sequence[str]) -> dict[str, PathLike[str]]: directory = Path(root).absolute() if isinstance(include, str) or not include: raise ValueError("Select directory files with explicit include patterns") @@ -139,14 +153,16 @@ def directory_files(root: str | PathLike[str], destination: str, include: Sequen raise ValueError("The selected directory must be a directory, not a symlink") directory = directory.resolve() # System aliases such as macOS /tmp are allowed above the chosen root. - _destination(destination + "/selected-file") + _destination(destination + "/a") # Validate with the shortest possible filename. patterns: list[tuple[str, ...]] = [] for pattern in include: if not pattern or PurePosixPath(pattern).is_absolute() or ".." in PurePosixPath(pattern).parts: raise ValueError("Include patterns must stay inside the selected directory") patterns.append(PurePosixPath(pattern).parts) - selected: dict[str, Path] = {} - for current, directories, names in walk(directory, followlinks=False): + selected: dict[str, PathLike[str]] = {} + for current, directories, names in walk(directory, followlinks=False, onerror=_walk_error): + if not Path(current).resolve().is_relative_to(directory): + raise ValueError("Selected directory traversal left its root") directories.sort() for name in sorted([*directories, *names]): source = Path(current) / name @@ -159,20 +175,30 @@ def directory_files(root: str | PathLike[str], destination: str, include: Sequen directories.remove(name) continue if chosen and name not in directories: - selected[str(PurePosixPath(destination) / relative.as_posix())] = source + canonical = source.resolve() + if not canonical.is_relative_to(directory): + raise ValueError("Selected agent file is outside the chosen directory") + metadata = canonical.lstat() + selected[str(PurePosixPath(destination) / relative.as_posix())] = _SelectedFile( + canonical, (metadata.st_dev, metadata.st_ino, metadata.st_size) + ) if not selected: raise ValueError("The include patterns did not select any files") return selected -async def async_directory_files(root: str | PathLike[str], destination: str, include: Sequence[str]) -> dict[str, Path]: - return await run_sync(directory_files, root, destination, include) +async def async_directory_files( + root: str | PathLike[str], destination: str, include: Sequence[str] +) -> dict[str, PathLike[str]]: + return await run_sync(directory_files, root, destination, include, abandon_on_cancel=True) -def _open_local(path: Path, stack: ExitStack) -> tuple[BinaryIO, int]: +def _open_local(path: Path, stack: ExitStack, expected: tuple[int, int, int] | None = None) -> tuple[BinaryIO, int]: before = path.lstat() if not S_ISREG(before.st_mode): raise ValueError("Selected agent files must be regular files, not symlinks") + if expected is not None and expected != (before.st_dev, before.st_ino, before.st_size): + raise ValueError("Selected agent file changed after directory selection") handle = stack.enter_context(path.open("rb")) opened = fstat(handle.fileno()) if (before.st_dev, before.st_ino, before.st_size) != (opened.st_dev, opened.st_ino, opened.st_size): @@ -199,6 +225,9 @@ def prepare(resource: Files, files: Mapping[str, str | PathLike[str]], options: prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) except Exception as error: raise AgentFilePreparationError(prepared) from error + except BaseException as error: + error.__dict__["prepared"] = prepared + raise return prepared @@ -206,13 +235,13 @@ async def async_prepare( resource: AsyncFiles, files: Mapping[str, str | PathLike[str]], options: _RequestOptions ) -> PreparedAgentFiles: with ExitStack() as stack: - selected = _prepare_selection(files, options, resource._client.default_headers, stack) + selected = await run_sync(_prepare_selection, files, options, resource._client.default_headers, stack) prepared = PreparedAgentFiles() try: for destination, source, handle, length in selected: - _unchanged_size(handle, length) + content = await run_sync(_read_local, handle, length) uploaded = await resource._client.files.create( - file=(source.name, handle), purpose="user_data", **options + file=(source.name, content), purpose="user_data", **options ) prepared.uploaded_file_ids.append(uploaded.id) prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) @@ -224,6 +253,23 @@ async def async_prepare( return prepared +def _read_local(handle: BinaryIO, length: int) -> bytes: + if length > _MAX_BYTES: + raise ValueError("An agent file must not exceed 50 MiB") + _unchanged_size(handle, length) + content = handle.read(length + 1) + if len(content) != length: + raise ValueError("Selected agent file changed during preparation") + return content + + +def _snapshot_upload(file: FileTypes) -> FileTypes: + assert isinstance(file, tuple) and isinstance(file[1], IOBase) + handle = cast(BinaryIO, file[1]) + content = _read_local(handle, fstat(handle.fileno()).st_size) + return cast(FileTypes, (file[0], content, *file[2:])) + + @contextmanager def _upload_content(file: FileTypes) -> Generator[FileTypes, None, None]: content = file[1] if isinstance(file, tuple) else file @@ -250,7 +296,8 @@ def _upload_content(file: FileTypes) -> Generator[FileTypes, None, None]: if S_ISREG(metadata.st_mode): size = metadata.st_size except (OSError, ValueError): - pass + # Some streams expose neither a seekable length nor a file descriptor. + size = None if size is not None and size > _MAX_BYTES: raise ValueError("An agent file must not exceed 50 MiB") yield file @@ -268,6 +315,9 @@ def upload( staged = resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) except Exception as error: raise AgentFileStagingError(uploaded.id) from error + except BaseException as error: + error.__dict__["uploaded_file_id"] = uploaded.id + raise return StagedAgentFile(uploaded_file_id=uploaded.id, file=staged) @@ -277,7 +327,10 @@ async def async_upload( if not environment_id: raise ValueError("Expected a non-empty environment_id") destination = _destination(path) + original = file[1] if isinstance(file, tuple) else file with _upload_content(file) as content: + if isinstance(original, PathLike): + content = await run_sync(_snapshot_upload, content) uploaded = await resource._client.files.create(file=content, purpose="user_data", **options) try: staged = await resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 60dee43e16..1777b6812f 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -413,3 +413,180 @@ def select(root: Any, destination: str, include: Any) -> Any: tmp_path, destination="/workspace/docs", include=["*.txt"] ) assert prepared.uploaded_file_ids == ["file_1"] + + +async def prepare_directory(sdk: OpenAI | AsyncOpenAI, root: Path, destination: str = "/workspace/docs") -> Any: + if isinstance(sdk, AsyncOpenAI): + return await sdk.beta.agents.environments.files.prepare_directory( + root, destination=destination, include=["**/*.txt"] + ) + return sdk.beta.agents.environments.files.prepare_directory(root, destination=destination, include=["**/*.txt"]) + + +async def test_directory_walk_error_is_not_a_partial_success( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + from openai.lib.beta.agents import _files + + (tmp_path / "keep.txt").write_text("abc") + + def failed_walk(root: Any, **options: Any) -> Any: + yield str(root), [], ["keep.txt"] + options["onerror"](PermissionError("synthetic unreadable subtree")) + + monkeypatch.setattr(_files, "walk", failed_walk) + with pytest.raises(PermissionError, match="synthetic"): + await prepare_directory(sdk, tmp_path) + assert not server.requests + + +async def test_cached_directory_entry_cannot_leave_selected_root( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + from openai.lib.beta.agents import _files + + root = tmp_path / "chosen" + child = root / "child" + child.mkdir(parents=True) + outside = tmp_path / "outside" + outside.mkdir() + (outside / "other.txt").write_text("outside") + + def changed_walk(directory: Any, **_options: Any) -> Any: + yield str(directory), ["child"], [] + child.rename(root / "original-child") + child.symlink_to(outside, target_is_directory=True) + yield str(child), [], ["other.txt"] + + monkeypatch.setattr(_files, "walk", changed_walk) + with pytest.raises(ValueError, match="left its root"): + await prepare_directory(sdk, root) + assert not server.requests + + +async def test_directory_selection_keeps_file_identity_until_open( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + from openai.lib.beta.agents import _files + from openai.resources.beta.agents.environments import files as resource_files + + source = tmp_path / "chosen.txt" + source.write_text("abc") + original = _files.directory_files + + def select_then_replace(root: Any, destination: str, include: Any) -> Any: + selected = original(root, destination, include) + source.rename(tmp_path / "original.txt") + source.write_text("xyz") + return selected + + monkeypatch.setattr(_files, "directory_files", select_then_replace) + monkeypatch.setattr(resource_files, "directory_files", select_then_replace) + with pytest.raises(ValueError, match="changed after directory selection"): + await prepare_directory(sdk, tmp_path) + assert not server.requests + + +async def test_directory_destination_uses_actual_filename_length( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + (tmp_path / "a").write_text("abc") + destination = "/workspace/" + "x" * (4094 - len("/workspace/")) + if isinstance(sdk, AsyncOpenAI): + prepared = await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination=destination, include=["a"] + ) + else: + prepared = sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination=destination, include=["a"] + ) + assert len(prepared.files[0]["path"]) == 4096 + assert server.uploads == 1 + + +async def test_async_path_reads_run_off_event_loop( + sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + import threading + + from openai.lib.beta.agents import _files + + if not isinstance(sdk, AsyncOpenAI): + pytest.skip("Async disk read contract") + source = tmp_path / "source.txt" + source.write_text("abc") + current_thread = threading.get_ident() + original = _files._read_local + calls: list[int] = [] + + def read(handle: Any, length: int) -> bytes: + assert threading.get_ident() != current_thread + calls.append(length) + return original(handle, length) + + monkeypatch.setattr(_files, "_read_local", read) + await sdk.beta.agents.environments.files.prepare({"/workspace/source.txt": source}) + await sdk.beta.agents.environments.files.upload("env_test", file=source, path="/workspace/source.txt") + assert calls == [3, 3] + + +@pytest.mark.parametrize("staging", [False, True]) +async def test_sync_interruption_preserves_observed_uploads( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, staging: bool +) -> None: + if isinstance(sdk, AsyncOpenAI): + pytest.skip("Synchronous interruption contract") + source = tmp_path / "source.txt" + source.write_text("abc") + + def interrupt(*_args: Any, **_kwargs: Any) -> Any: + raise KeyboardInterrupt() + + if staging: + monkeypatch.setattr(sdk.beta.agents.environments.files, "create", interrupt) + with pytest.raises(KeyboardInterrupt) as caught: + sdk.beta.agents.environments.files.upload("env_test", file=source, path="/workspace/source.txt") + assert vars(caught.value)["uploaded_file_id"] == "file_1" + else: + + def after_upload() -> None: + if server.uploads == 2: + raise KeyboardInterrupt() + + server.after_upload = after_upload + with pytest.raises(KeyboardInterrupt) as caught: + sdk.beta.agents.environments.files.prepare({"/workspace/a": source, "/workspace/b": source}) + assert vars(caught.value)["prepared"].uploaded_file_ids == ["file_1"] + + +async def test_async_selection_cancellation_does_not_wait_for_scan(monkeypatch: pytest.MonkeyPatch) -> None: + import threading + + import anyio + + from openai.lib.beta.agents import _files + + started = threading.Event() + release = threading.Event() + timer = threading.Timer(2, release.set) + + def scan(*_args: Any) -> Any: + started.set() + release.wait() + return {} + + async def select() -> None: + await _files.async_directory_files("synthetic", "/workspace/docs", ["*"]) + + monkeypatch.setattr(_files, "directory_files", scan) + timer.start() + try: + async with anyio.create_task_group() as tasks: + tasks.start_soon(select) + while not started.is_set(): + await anyio.sleep(0) + tasks.cancel_scope.cancel() + assert not release.is_set() + finally: + release.set() + timer.cancel() From 6f18dc68a6c29d744fa3e48b8478d3e7fa891d5f Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 02:34:57 +0000 Subject: [PATCH 06/18] docs(agents): scope local path-selection assumptions --- helpers.md | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/helpers.md b/helpers.md index f103a1b216..5df8503cdb 100644 --- a/helpers.md +++ b/helpers.md @@ -664,6 +664,6 @@ Files API when ready to delete them. Preparation errors expose partial uploads through `error.prepared`. A batch of multiple uploads cannot share one explicit `Idempotency-Key`. -Use application-controlled local paths and stable source directories. These -convenience helpers are not a filesystem sandbox for untrusted paths or hostile -local writers. With `AsyncOpenAI`, await preparation, staging, and artifact downloads. +Local path/directory uploads are intended for static application-owned files and +stable directories. They do not sandbox untrusted path selection or hostile local +filesystem writers. With `AsyncOpenAI`, await preparation, staging, and artifact downloads. From 6fc7b78e56006af379f514292e2f7a1500cf1695 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 02:43:22 +0000 Subject: [PATCH 07/18] feat(agents): read result artifacts in memory and keep worker cleanup local --- helpers.md | 6 +- src/openai/lib/beta/agents/_artifacts.py | 102 ++++++++++++++----- src/openai/lib/beta/agents/_files.py | 56 +++++----- tests/lib/streaming/agents/test_artifacts.py | 33 ++++++ tests/lib/streaming/agents/test_files.py | 58 +++++++++++ 5 files changed, 207 insertions(+), 48 deletions(-) diff --git a/helpers.md b/helpers.md index 5df8503cdb..e6a0a57a6a 100644 --- a/helpers.md +++ b/helpers.md @@ -651,7 +651,11 @@ with client.beta.agents.sessions.stream( ) as stream: result = stream.get_final_result() -artifact = client.beta.agents.sessions.artifacts.for_result(result).download( +artifacts = client.beta.agents.sessions.artifacts.for_result(result) + +# Read in memory, or stream to an application-owned local path. +report_bytes = artifacts.content("/workspace/outputs/report.md").content +artifact = artifacts.download( "/workspace/outputs/report.md", to=Path("downloaded-report.md") ) ``` diff --git a/src/openai/lib/beta/agents/_artifacts.py b/src/openai/lib/beta/agents/_artifacts.py index 7ca4696e2c..002da07295 100644 --- a/src/openai/lib/beta/agents/_artifacts.py +++ b/src/openai/lib/beta/agents/_artifacts.py @@ -5,8 +5,10 @@ import httpx2 +from ._files import _RequestOptions from ._result import AgentTurnResult from ...._types import Body, Query, Headers, NotGiven, not_given +from ...._legacy_response import HttpxBinaryResponseContent from ....types.beta.agents.sessions.session_artifact import SessionArtifact if TYPE_CHECKING: @@ -21,6 +23,25 @@ def __init__(self, resource: Artifacts, result: AgentTurnResult) -> None: self._session_id = result.session_id self._turn_id = result.turn_id + def content( + self, + path: str, + *, + extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> HttpxBinaryResponseContent: + """Read one exact turn/path using the SDK's native binary response.""" + options: _RequestOptions = { + "extra_headers": extra_headers, + "extra_query": extra_query, + "extra_body": extra_body, + "timeout": timeout, + } + selected = self._find(path, options) + return self._resource.content(selected.id, session_id=self._session_id, **options) + def download( self, path: str, @@ -32,13 +53,29 @@ def download( timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> SessionArtifact: """Download one exact turn/path to an explicit local destination.""" - selected: SessionArtifact | None = None - for artifact in self._resource.list( - self._session_id, + options: _RequestOptions = { + "extra_headers": extra_headers, + "extra_query": extra_query, + "extra_body": extra_body, + "timeout": timeout, + } + selected = self._find(path, options) + with self._resource.with_streaming_response.content( + selected.id, + session_id=self._session_id, extra_headers=extra_headers, extra_query=extra_query, extra_body=extra_body, timeout=timeout, + ) as content: + content.stream_to_file(to) + return selected + + def _find(self, path: str, options: _RequestOptions) -> SessionArtifact: + selected: SessionArtifact | None = None + for artifact in self._resource.list( + self._session_id, + **options, ): if artifact.session_id == self._session_id and artifact.turn_id == self._turn_id and artifact.path == path: if selected is not None: @@ -46,15 +83,6 @@ def download( selected = artifact if selected is None: raise ValueError("No artifact matches this turn and path") - with self._resource.with_streaming_response.content( - selected.id, - session_id=self._session_id, - extra_headers=extra_headers, - extra_query=extra_query, - extra_body=extra_body, - timeout=timeout, - ) as content: - content.stream_to_file(to) return selected @@ -66,6 +94,25 @@ def __init__(self, resource: AsyncArtifacts, result: AgentTurnResult) -> None: self._session_id = result.session_id self._turn_id = result.turn_id + async def content( + self, + path: str, + *, + extra_headers: Headers | None = None, + extra_query: Query | None = None, + extra_body: Body | None = None, + timeout: float | httpx2.Timeout | None | NotGiven = not_given, + ) -> HttpxBinaryResponseContent: + """Read one exact turn/path using the SDK's native binary response.""" + options: _RequestOptions = { + "extra_headers": extra_headers, + "extra_query": extra_query, + "extra_body": extra_body, + "timeout": timeout, + } + selected = await self._find(path, options) + return await self._resource.content(selected.id, session_id=self._session_id, **options) + async def download( self, path: str, @@ -77,13 +124,29 @@ async def download( timeout: float | httpx2.Timeout | None | NotGiven = not_given, ) -> SessionArtifact: """Download one exact turn/path to an explicit local destination.""" - selected: SessionArtifact | None = None - async for artifact in self._resource.list( - self._session_id, + options: _RequestOptions = { + "extra_headers": extra_headers, + "extra_query": extra_query, + "extra_body": extra_body, + "timeout": timeout, + } + selected = await self._find(path, options) + async with self._resource.with_streaming_response.content( + selected.id, + session_id=self._session_id, extra_headers=extra_headers, extra_query=extra_query, extra_body=extra_body, timeout=timeout, + ) as content: + await content.stream_to_file(to) + return selected + + async def _find(self, path: str, options: _RequestOptions) -> SessionArtifact: + selected: SessionArtifact | None = None + async for artifact in self._resource.list( + self._session_id, + **options, ): if artifact.session_id == self._session_id and artifact.turn_id == self._turn_id and artifact.path == path: if selected is not None: @@ -91,13 +154,4 @@ async def download( selected = artifact if selected is None: raise ValueError("No artifact matches this turn and path") - async with self._resource.with_streaming_response.content( - selected.id, - session_id=self._session_id, - extra_headers=extra_headers, - extra_query=extra_query, - extra_body=extra_body, - timeout=timeout, - ) as content: - await content.stream_to_file(to) return selected diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index 067b58ce72..b57c4b31e9 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -231,26 +231,35 @@ def prepare(resource: Files, files: Mapping[str, str | PathLike[str]], options: return prepared +def _snapshot_selection( + files: Mapping[str, str | PathLike[str]], options: _RequestOptions, defaults: Headers +) -> list[tuple[str, str, bytes]]: + # The worker owns every handle it opens, even if its caller is cancelled. + with ExitStack() as stack: + selected = _prepare_selection(files, options, defaults, stack) + return [ + (destination, source.name, _read_local(handle, length)) for destination, source, handle, length in selected + ] + + async def async_prepare( resource: AsyncFiles, files: Mapping[str, str | PathLike[str]], options: _RequestOptions ) -> PreparedAgentFiles: - with ExitStack() as stack: - selected = await run_sync(_prepare_selection, files, options, resource._client.default_headers, stack) - prepared = PreparedAgentFiles() - try: - for destination, source, handle, length in selected: - content = await run_sync(_read_local, handle, length) - uploaded = await resource._client.files.create( - file=(source.name, content), purpose="user_data", **options - ) - prepared.uploaded_file_ids.append(uploaded.id) - prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) - except anyio.get_cancelled_exc_class() as error: - error.__dict__["prepared"] = prepared - raise - except Exception as error: - raise AgentFilePreparationError(prepared) from error - return prepared + selected = await run_sync( + _snapshot_selection, files, options, resource._client.default_headers, abandon_on_cancel=True + ) + prepared = PreparedAgentFiles() + try: + for destination, filename, content in selected: + uploaded = await resource._client.files.create(file=(filename, content), purpose="user_data", **options) + prepared.uploaded_file_ids.append(uploaded.id) + prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) + except anyio.get_cancelled_exc_class() as error: + error.__dict__["prepared"] = prepared + raise + except Exception as error: + raise AgentFilePreparationError(prepared) from error + return prepared def _read_local(handle: BinaryIO, length: int) -> bytes: @@ -264,10 +273,11 @@ def _read_local(handle: BinaryIO, length: int) -> bytes: def _snapshot_upload(file: FileTypes) -> FileTypes: - assert isinstance(file, tuple) and isinstance(file[1], IOBase) - handle = cast(BinaryIO, file[1]) - content = _read_local(handle, fstat(handle.fileno()).st_size) - return cast(FileTypes, (file[0], content, *file[2:])) + with _upload_content(file) as opened: + assert isinstance(opened, tuple) and isinstance(opened[1], IOBase) + handle = cast(BinaryIO, opened[1]) + content = _read_local(handle, fstat(handle.fileno()).st_size) + return cast(FileTypes, (opened[0], content, *opened[2:])) @contextmanager @@ -328,9 +338,9 @@ async def async_upload( raise ValueError("Expected a non-empty environment_id") destination = _destination(path) original = file[1] if isinstance(file, tuple) else file + if isinstance(original, PathLike): + file = await run_sync(_snapshot_upload, file, abandon_on_cancel=True) with _upload_content(file) as content: - if isinstance(original, PathLike): - content = await run_sync(_snapshot_upload, content) uploaded = await resource._client.files.create(file=content, purpose="user_data", **options) try: staged = await resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) diff --git a/tests/lib/streaming/agents/test_artifacts.py b/tests/lib/streaming/agents/test_artifacts.py index 9b02e68c5c..e2df2bc44c 100644 --- a/tests/lib/streaming/agents/test_artifacts.py +++ b/tests/lib/streaming/agents/test_artifacts.py @@ -131,3 +131,36 @@ async def test_missing_or_ambiguous_artifact_does_not_open_destination( await download(sdk, path) assert path.read_text() == "keep" assert server.content.chunks == 0 + + +async def test_native_binary_content_matches_exact_turn_and_preserves_options( + sdk: OpenAI | AsyncOpenAI, server: ArtifactServer +) -> None: + from openai._legacy_response import HttpxBinaryResponseContent + + options: Any = {"extra_headers": {"X-Synthetic": "test"}, "extra_query": {"synthetic": "query"}, "timeout": 7} + content = ( + await sdk.beta.agents.sessions.artifacts.for_result(result()).content("/workspace/outputs/report.md", **options) + if isinstance(sdk, AsyncOpenAI) + else sdk.beta.agents.sessions.artifacts.for_result(result()).content("/workspace/outputs/report.md", **options) + ) + assert isinstance(content, HttpxBinaryResponseContent) + assert content.content == b"report-done" + assert len(server.requests) == 3 + assert all(request.url.params["synthetic"] == "query" for request in server.requests) + assert all(request.headers["X-Synthetic"] == "test" for request in server.requests) + assert all(request.extensions["timeout"]["read"] == 7 for request in server.requests) + assert server.content.closed + + +@pytest.mark.parametrize("ambiguous", [False, True]) +async def test_native_content_rejects_missing_or_ambiguous_artifact( + sdk: OpenAI | AsyncOpenAI, server: ArtifactServer, ambiguous: bool +) -> None: + server.artifacts = [artifact("selected"), artifact("duplicate")] if ambiguous else [artifact("old", "old_turn")] + with pytest.raises(ValueError, match="More than one" if ambiguous else "No artifact"): + if isinstance(sdk, AsyncOpenAI): + await sdk.beta.agents.sessions.artifacts.for_result(result()).content("/workspace/outputs/report.md") + else: + sdk.beta.agents.sessions.artifacts.for_result(result()).content("/workspace/outputs/report.md") + assert server.content.chunks == 0 diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 1777b6812f..fe9392309a 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -590,3 +590,61 @@ async def select() -> None: finally: release.set() timer.cancel() + + +@pytest.mark.parametrize("mode", ["prepare", "directory", "upload"]) +async def test_native_cancellation_keeps_opened_handles_owned_by_worker( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, mode: str +) -> None: + import asyncio + import threading + + from openai.lib.beta.agents import _files + + if not isinstance(sdk, AsyncOpenAI): + pytest.skip("Native async cancellation contract") + source = tmp_path / "source.txt" + source.write_text("abc") + entered = threading.Event() + release = threading.Event() + finished = threading.Event() + handles: list[Any] = [] + original_open = _files._open_local + snapshot_name = "_snapshot_upload" if mode == "upload" else "_snapshot_selection" + original_snapshot = getattr(_files, snapshot_name) + + def delayed_open(*args: Any, **kwargs: Any) -> Any: + entered.set() + release.wait() + handle, size = original_open(*args, **kwargs) + handles.append(handle) + return handle, size + + def snapshot(*args: Any, **kwargs: Any) -> Any: + try: + return original_snapshot(*args, **kwargs) + finally: + finished.set() + + monkeypatch.setattr(_files, "_open_local", delayed_open) + monkeypatch.setattr(_files, snapshot_name, snapshot) + if mode == "upload": + operation = sdk.beta.agents.environments.files.upload("env_test", file=source, path="/workspace/source.txt") + elif mode == "directory": + operation = sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["*.txt"] + ) + else: + operation = sdk.beta.agents.environments.files.prepare({"/workspace/source.txt": source}) + task = asyncio.create_task(operation) + try: + await asyncio.wait_for(asyncio.to_thread(entered.wait), 2) + task.cancel() + with pytest.raises(CancelledError): + await task + finally: + release.set() + await asyncio.wait_for(asyncio.to_thread(finished.wait), 2) + assert len(handles) == 1 + assert handles[0].closed + assert not server.requests From a7cc1966c00749892f2ac37e8029d89d0cff8cf4 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 02:51:59 +0000 Subject: [PATCH 08/18] test(agents): use consistent asyncio imports --- tests/lib/streaming/agents/test_files.py | 13 ++++++------- 1 file changed, 6 insertions(+), 7 deletions(-) diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index fe9392309a..6938be38dd 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -1,9 +1,9 @@ from __future__ import annotations import json +import asyncio from io import BytesIO from typing import Any, Callable -from asyncio import CancelledError from pathlib import Path from typing_extensions import override @@ -33,7 +33,7 @@ def handle(self, request: httpx2.Request) -> httpx2.Response: if self.after_upload is not None: self.after_upload() if self.uploads == self.cancel_upload: - raise CancelledError() + raise asyncio.CancelledError() if self.uploads == self.fail_upload: return httpx2.Response(500, json={"error": {"message": "Synthetic upload failure"}}) return httpx2.Response( @@ -50,7 +50,7 @@ def handle(self, request: httpx2.Request) -> httpx2.Response: ) assert request.url.path == "/v1/agents/environments/env_test/files" if self.cancel_stage: - raise CancelledError() + raise asyncio.CancelledError() if self.fail_stage: return httpx2.Response(500, json={"error": {"message": "Synthetic staging failure"}}) body = json.loads(request.content) @@ -213,7 +213,7 @@ async def test_async_cancellation_retains_only_observed_uploads( source = tmp_path / "source.txt" source.write_text("abc") server.cancel_upload = 2 - with pytest.raises(CancelledError) as caught: + with pytest.raises(asyncio.CancelledError) as caught: await sdk.beta.agents.environments.files.prepare({"/workspace/a": source, "/workspace/b": source}) assert vars(caught.value)["prepared"].uploaded_file_ids == ["file_1"] assert all(request.method == "POST" for request in server.requests) @@ -261,7 +261,7 @@ async def test_cancelled_staging_preserves_observed_upload_id(sdk: OpenAI | Asyn if not isinstance(sdk, AsyncOpenAI): pytest.skip("Async cancellation contract") server.cancel_stage = True - with pytest.raises(CancelledError) as caught: + with pytest.raises(asyncio.CancelledError) as caught: await sdk.beta.agents.environments.files.upload( "env_test", file=("source.txt", b"abc"), path="/workspace/source.txt" ) @@ -596,7 +596,6 @@ async def select() -> None: async def test_native_cancellation_keeps_opened_handles_owned_by_worker( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, mode: str ) -> None: - import asyncio import threading from openai.lib.beta.agents import _files @@ -640,7 +639,7 @@ def snapshot(*args: Any, **kwargs: Any) -> Any: try: await asyncio.wait_for(asyncio.to_thread(entered.wait), 2) task.cancel() - with pytest.raises(CancelledError): + with pytest.raises(asyncio.CancelledError): await task finally: release.set() From baaef5dd24ac400a2a230cbadeaebf89f3b0de39 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 14:46:04 +0000 Subject: [PATCH 09/18] chore(agents): order combined beta helper exports --- src/openai/lib/beta/agents/__init__.py | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/src/openai/lib/beta/agents/__init__.py b/src/openai/lib/beta/agents/__init__.py index 125ae0018d..80ca2a8b4e 100644 --- a/src/openai/lib/beta/agents/__init__.py +++ b/src/openai/lib/beta/agents/__init__.py @@ -1,16 +1,16 @@ """Beta helpers for hosted Agents API tools and turn results.""" -from ._tools import ( - FunctionTool as FunctionTool, - function_tool as function_tool, - pydantic_function_tool as pydantic_function_tool, -) from ._files import ( StagedAgentFile as StagedAgentFile, PreparedAgentFiles as PreparedAgentFiles, AgentFileStagingError as AgentFileStagingError, AgentFilePreparationError as AgentFilePreparationError, ) +from ._tools import ( + FunctionTool as FunctionTool, + function_tool as function_tool, + pydantic_function_tool as pydantic_function_tool, +) from ._output import agent_text_format as agent_text_format from ._result import ( AgentTurnResult as AgentTurnResult, From 62fcfb2ff1d2be77fd5880c67e61f48ea78fadc7 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 19:23:21 +0000 Subject: [PATCH 10/18] fix(agents): preserve interrupted uploads and defer file limits to API --- src/openai/lib/beta/agents/_files.py | 51 +++--------------- tests/lib/streaming/agents/test_files.py | 69 +++++++++++++++--------- 2 files changed, 49 insertions(+), 71 deletions(-) diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index b57c4b31e9..f8bf3cf79f 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -10,7 +10,6 @@ from dataclasses import field, dataclass from typing_extensions import TypedDict, override -import anyio import httpx2 from anyio.to_thread import run_sync @@ -30,10 +29,6 @@ class _RequestOptions(TypedDict): timeout: float | httpx2.Timeout | None | NotGiven -_MAX_FILES = 50 -_MAX_BYTES = 50 * 1024 * 1024 - - @dataclass(frozen=True) class _SelectedFile(PathLike[str]): path: Path @@ -88,8 +83,6 @@ def _destination(value: str) -> str: or value == "/workspace/outputs" ): raise ValueError("Agent file destination uses a reserved environment path") - if len(value) > 4096: - raise ValueError("Agent file destinations must not exceed 4096 characters") return value @@ -99,23 +92,18 @@ def _local_file(value: str | PathLike[str]) -> Path: raise ValueError("Selected agent files must not be symlinks") if not path.is_file(): raise ValueError("Selected agent files must be regular files") - if path.stat().st_size > _MAX_BYTES: - raise ValueError("An agent file must not exceed 50 MiB") return path def _prepare_selection( files: Mapping[str, str | PathLike[str]], options: _RequestOptions, defaults: Headers, stack: ExitStack ) -> list[tuple[str, Path, BinaryIO, int]]: - if len(files) > _MAX_FILES: - raise ValueError("Initial agent files must not exceed 50 files") headers = {key.lower(): value for key, value in defaults.items()} headers.update({key.lower(): value for key, value in (options["extra_headers"] or {}).items()}) if len(files) > 1 and "idempotency-key" in headers and not isinstance(headers["idempotency-key"], Omit): raise ValueError("One Idempotency-Key cannot be reused for multiple file uploads") selected: list[tuple[str, Path, BinaryIO, int]] = [] destinations: set[str] = set() - size = 0 for destination, source in files.items(): destination = _destination(destination) if any( @@ -126,10 +114,7 @@ def _prepare_selection( destinations.add(destination) local = _local_file(source) handle, length = _open_local(local, stack, source.identity if isinstance(source, _SelectedFile) else None) - size += length selected.append((destination, local, handle, length)) - if size > _MAX_BYTES: - raise ValueError("Initial agent files must not exceed 50 MiB in total") return selected @@ -203,8 +188,6 @@ def _open_local(path: Path, stack: ExitStack, expected: tuple[int, int, int] | N opened = fstat(handle.fileno()) if (before.st_dev, before.st_ino, before.st_size) != (opened.st_dev, opened.st_ino, opened.st_size): raise ValueError("Selected agent file changed during preparation") - if opened.st_size > _MAX_BYTES: - raise ValueError("An agent file must not exceed 50 MiB") return handle, opened.st_size @@ -254,17 +237,15 @@ async def async_prepare( uploaded = await resource._client.files.create(file=(filename, content), purpose="user_data", **options) prepared.uploaded_file_ids.append(uploaded.id) prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) - except anyio.get_cancelled_exc_class() as error: - error.__dict__["prepared"] = prepared - raise except Exception as error: raise AgentFilePreparationError(prepared) from error + except BaseException as error: + error.__dict__["prepared"] = prepared + raise return prepared def _read_local(handle: BinaryIO, length: int) -> bytes: - if length > _MAX_BYTES: - raise ValueError("An agent file must not exceed 50 MiB") _unchanged_size(handle, length) content = handle.read(length + 1) if len(content) != length: @@ -290,26 +271,6 @@ def _upload_content(file: FileTypes) -> Generator[FileTypes, None, None]: _unchanged_size(handle, length) yield cast(FileTypes, (file[0], handle, *file[2:])) if isinstance(file, tuple) else (path.name, handle) return - size = None - if isinstance(content, bytes): - size = len(content) - elif isinstance(content, IOBase): - try: - if content.seekable(): - position = content.tell() - try: - size = content.seek(0, 2) - finally: - content.seek(position) - else: - metadata = fstat(content.fileno()) - if S_ISREG(metadata.st_mode): - size = metadata.st_size - except (OSError, ValueError): - # Some streams expose neither a seekable length nor a file descriptor. - size = None - if size is not None and size > _MAX_BYTES: - raise ValueError("An agent file must not exceed 50 MiB") yield file @@ -344,9 +305,9 @@ async def async_upload( uploaded = await resource._client.files.create(file=content, purpose="user_data", **options) try: staged = await resource.create(environment_id, type="file_id", file_id=uploaded.id, path=destination, **options) - except anyio.get_cancelled_exc_class() as error: - error.__dict__["uploaded_file_id"] = uploaded.id - raise except Exception as error: raise AgentFileStagingError(uploaded.id) from error + except BaseException as error: + error.__dict__["uploaded_file_id"] = uploaded.id + raise return StagedAgentFile(uploaded_file_id=uploaded.id, file=staged) diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 6938be38dd..a61e00cb36 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -10,7 +10,7 @@ import httpx2 import pytest -from openai import OpenAI, AsyncOpenAI +from openai import OpenAI, AsyncOpenAI, BadRequestError from openai.lib.beta.agents import AgentFileStagingError, AgentFilePreparationError from tests.lib.streaming.agents.test_streams import Server, sdk as sdk @@ -188,21 +188,31 @@ async def stage() -> Any: assert all(request.method == "POST" for request in server.requests) -async def test_count_and_aggregate_size_limits_precede_upload( - sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +async def test_file_count_is_left_to_the_api(sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path) -> None: + source = tmp_path / "source.txt" + source.write_text("abc") + result = await prepare(sdk, {f"/workspace/{index}": source for index in range(51)}) + assert len(result.files) == server.uploads == 51 + + +@pytest.mark.parametrize("count", [1, 2]) +async def test_large_prepared_files_reach_the_api( + sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, count: int ) -> None: source = tmp_path / "large" with source.open("wb") as content: - content.truncate(26 * 1024 * 1024) - with pytest.raises(ValueError, match="in total"): - await prepare(sdk, {"/workspace/a": source, "/workspace/b": source}) - with pytest.raises(ValueError, match="50 files"): - await prepare(sdk, {f"/workspace/{index}": source for index in range(51)}) - with source.open("wb") as content: - content.truncate(50 * 1024 * 1024 + 1) - with pytest.raises(ValueError, match="50 MiB"): - await prepare(sdk, {"/workspace/a": source}) - assert server.uploads == 0 + content.truncate(51 * 1024 * 1024 // count) + response = httpx2.Response(400, request=httpx2.Request("POST", "https://sdk-test.example/v1/files")) + rejection = BadRequestError("Synthetic server-owned limit", response=response, body=None) + + def reject(**_kwargs: Any) -> Any: + raise rejection + + monkeypatch.setattr(sdk.files, "create", reject) + with pytest.raises(AgentFilePreparationError) as caught: + await prepare(sdk, {f"/workspace/{index}": source for index in range(count)}) + assert caught.value.__cause__ is rejection + assert caught.value.prepared.uploaded_file_ids == [] async def test_async_cancellation_retains_only_observed_uploads( @@ -244,7 +254,6 @@ async def test_client_default_idempotency_key_is_rejected_for_batch( "/workspace/.managed-agents/source", "/workspace/.managed-agents-internal/source", "/workspace/outputs", - "/workspace/" + "x" * 4096, ], ) async def test_invalid_hosted_destination_fails_before_upload( @@ -314,23 +323,30 @@ def replace_later_source() -> None: @pytest.mark.parametrize("file_backed", [False, True]) -async def test_known_stream_sizes_fail_before_upload_without_moving_position( - sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, file_backed: bool +async def test_large_uploads_preserve_api_errors( + sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, file_backed: bool ) -> None: content = (tmp_path / "large").open("w+b") if file_backed else BytesIO() + response = httpx2.Response(400, request=httpx2.Request("POST", "https://sdk-test.example/v1/files")) + rejection = BadRequestError("Synthetic server-owned limit", response=response, body=None) + + def reject(**_kwargs: Any) -> Any: + raise rejection + + monkeypatch.setattr(sdk.files, "create", reject) with content: content.seek(50 * 1024 * 1024) content.write(b"x") content.seek(3) - with pytest.raises(ValueError, match="50 MiB"): + with pytest.raises(BadRequestError) as caught: if isinstance(sdk, AsyncOpenAI): await sdk.beta.agents.environments.files.upload( "env_test", file=("large", content), path="/workspace/large" ) else: sdk.beta.agents.environments.files.upload("env_test", file=("large", content), path="/workspace/large") + assert caught.value is rejection assert content.tell() == 3 - assert not server.requests async def test_empty_environment_id_fails_before_upload(sdk: OpenAI | AsyncOpenAI, server: FilesServer) -> None: @@ -487,11 +503,11 @@ def select_then_replace(root: Any, destination: str, include: Any) -> Any: assert not server.requests -async def test_directory_destination_uses_actual_filename_length( +async def test_directory_preserves_long_destination_for_api_validation( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path ) -> None: (tmp_path / "a").write_text("abc") - destination = "/workspace/" + "x" * (4094 - len("/workspace/")) + destination = "/workspace/" + "x" * 4096 if isinstance(sdk, AsyncOpenAI): prepared = await sdk.beta.agents.environments.files.prepare_directory( tmp_path, destination=destination, include=["a"] @@ -500,7 +516,7 @@ async def test_directory_destination_uses_actual_filename_length( prepared = sdk.beta.agents.environments.files.prepare_directory( tmp_path, destination=destination, include=["a"] ) - assert len(prepared.files[0]["path"]) == 4096 + assert prepared.files[0]["path"] == destination + "/a" assert server.uploads == 1 @@ -531,11 +547,9 @@ def read(handle: Any, length: int) -> bytes: @pytest.mark.parametrize("staging", [False, True]) -async def test_sync_interruption_preserves_observed_uploads( +async def test_interruption_preserves_observed_uploads( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, staging: bool ) -> None: - if isinstance(sdk, AsyncOpenAI): - pytest.skip("Synchronous interruption contract") source = tmp_path / "source.txt" source.write_text("abc") @@ -545,7 +559,10 @@ def interrupt(*_args: Any, **_kwargs: Any) -> Any: if staging: monkeypatch.setattr(sdk.beta.agents.environments.files, "create", interrupt) with pytest.raises(KeyboardInterrupt) as caught: - sdk.beta.agents.environments.files.upload("env_test", file=source, path="/workspace/source.txt") + if isinstance(sdk, AsyncOpenAI): + await sdk.beta.agents.environments.files.upload("env_test", file=source, path="/workspace/source.txt") + else: + sdk.beta.agents.environments.files.upload("env_test", file=source, path="/workspace/source.txt") assert vars(caught.value)["uploaded_file_id"] == "file_1" else: @@ -555,7 +572,7 @@ def after_upload() -> None: server.after_upload = after_upload with pytest.raises(KeyboardInterrupt) as caught: - sdk.beta.agents.environments.files.prepare({"/workspace/a": source, "/workspace/b": source}) + await prepare(sdk, {"/workspace/a": source, "/workspace/b": source}) assert vars(caught.value)["prepared"].uploaded_file_ids == ["file_1"] From 4380393560b079444917d92657f502cb3e41512a Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 20:10:43 +0000 Subject: [PATCH 11/18] fix(agents): validate file destination ancestors without pairwise scans --- src/openai/lib/beta/agents/_files.py | 15 ++++++--------- tests/lib/streaming/agents/test_files.py | 24 ++++++++++++++++++++++++ 2 files changed, 30 insertions(+), 9 deletions(-) diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index f8bf3cf79f..16ae7895c4 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -102,16 +102,13 @@ def _prepare_selection( headers.update({key.lower(): value for key, value in (options["extra_headers"] or {}).items()}) if len(files) > 1 and "idempotency-key" in headers and not isinstance(headers["idempotency-key"], Omit): raise ValueError("One Idempotency-Key cannot be reused for multiple file uploads") - selected: list[tuple[str, Path, BinaryIO, int]] = [] - destinations: set[str] = set() - for destination, source in files.items(): - destination = _destination(destination) - if any( - destination == prior or destination.startswith(prior + "/") or prior.startswith(destination + "/") - for prior in destinations - ): + entries = list(files.items()) + destinations = {_destination(destination) for destination, _ in entries} + for destination in destinations: + if any(str(parent) in destinations for parent in PurePosixPath(destination).parents): raise ValueError("Selected agent files contain duplicate or conflicting destinations") - destinations.add(destination) + selected: list[tuple[str, Path, BinaryIO, int]] = [] + for destination, source in entries: local = _local_file(source) handle, length = _open_local(local, stack, source.identity if isinstance(source, _SelectedFile) else None) selected.append((destination, local, handle, length)) diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index a61e00cb36..2f5a6ec447 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -117,6 +117,30 @@ async def test_preflight_all_files_before_upload( assert server.uploads == 0 +@pytest.mark.parametrize("reverse", [False, True]) +async def test_destination_conflicts_are_checked_before_opening_sources( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, reverse: bool +) -> None: + destinations = ["/workspace/docs", "/workspace/docs/nested/source.txt"] + if reverse: + destinations.reverse() + with pytest.raises(ValueError, match="conflicting"): + await prepare(sdk, {destination: tmp_path / "missing" for destination in destinations}) + assert not server.requests + + +async def test_large_selection_keeps_input_order( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + source = tmp_path / "source.txt" + source.write_text("abc") + destinations = [f"/workspace/group-{index % 17}/file-{index}.txt" for index in reversed(range(257))] + result = await prepare(sdk, dict.fromkeys(destinations, source)) + assert [item["path"] for item in result.files] == destinations + assert [item["file_id"] for item in result.files] == [f"file_{index}" for index in range(1, 258)] + assert server.uploads == 257 + + async def test_partial_upload_ownership_is_available( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path ) -> None: From b45172840ccc657508a4cfdd2ecff1a16356e29e Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 20:11:43 +0000 Subject: [PATCH 12/18] fix(agents): skip directory subtrees outside file selections --- src/openai/lib/beta/agents/_files.py | 17 ++++-- tests/lib/streaming/agents/test_files.py | 72 ++++++++++++++++++++++++ 2 files changed, 85 insertions(+), 4 deletions(-) diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index 16ae7895c4..5e1daa6a7c 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -115,12 +115,17 @@ def _prepare_selection( return selected -def _matches(parts: tuple[str, ...], pattern: tuple[str, ...]) -> bool: +def _matches(parts: tuple[str, ...], pattern: tuple[str, ...], *, prefix: bool = False) -> bool: + # A directory prefix is useful only if the pattern can select a descendant. + if prefix and not parts: + return bool(pattern) if not pattern: return not parts if pattern[0] == "**": - return _matches(parts, pattern[1:]) or bool(parts) and _matches(parts[1:], pattern) - return bool(parts) and fnmatchcase(parts[0], pattern[0]) and _matches(parts[1:], pattern[1:]) + return ( + _matches(parts, pattern[1:], prefix=prefix) or bool(parts) and _matches(parts[1:], pattern, prefix=prefix) + ) + return bool(parts) and fnmatchcase(parts[0], pattern[0]) and _matches(parts[1:], pattern[1:], prefix=prefix) def _walk_error(error: OSError) -> None: @@ -156,7 +161,11 @@ def directory_files(root: str | PathLike[str], destination: str, include: Sequen if name in directories: directories.remove(name) continue - if chosen and name not in directories: + if name in directories: + if not any(_matches(relative.parts, pattern, prefix=True) for pattern in patterns): + directories.remove(name) + continue + if chosen: canonical = source.resolve() if not canonical.is_relative_to(directory): raise ValueError("Selected agent file is outside the chosen directory") diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 2f5a6ec447..6d7d9d3756 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -463,6 +463,78 @@ async def prepare_directory(sdk: OpenAI | AsyncOpenAI, root: Path, destination: return sdk.beta.agents.environments.files.prepare_directory(root, destination=destination, include=["**/*.txt"]) +@pytest.mark.parametrize( + "include,expected", + [ + (["one.txt"], ["one.txt"]), + (["*.txt"], ["one.txt"]), + (["docs/*.txt"], ["docs/one.txt"]), + (["docs/**/*.txt"], ["docs/one.txt", "docs/nested/one.txt"]), + (["docs/*/*.txt"], ["docs/nested/one.txt"]), + ], +) +async def test_directory_selection_skips_unrelated_unreadable_subtrees( + sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, include: list[str], expected: list[str] +) -> None: + import os + + for name in ["one.txt", "docs/one.txt", "docs/nested/one.txt", "private/one.txt"]: + source = tmp_path / name + source.parent.mkdir(parents=True, exist_ok=True) + source.write_text("abc") + scan = os.scandir + visited: list[Path] = [] + + def scandir(path: Any) -> Any: + current = Path(path) + visited.append(current) + if current == tmp_path / "private": + raise PermissionError("Synthetic unreadable unrelated subtree") + return scan(path) + + monkeypatch.setattr(os, "scandir", scandir) + if isinstance(sdk, AsyncOpenAI): + result = await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=include + ) + else: + result = sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=include + ) + assert [file["path"] for file in result.files] == [f"/workspace/docs/{name}" for name in expected] + assert tmp_path / "private" not in visited + if include == ["docs/*.txt"]: + assert tmp_path / "docs/nested" not in visited + + +@pytest.mark.parametrize("include", [["**/*.txt"], ["private/*.txt"]]) +async def test_directory_selection_propagates_needed_subtree_errors( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, include: list[str] +) -> None: + import os + + (tmp_path / "one.txt").write_text("abc") + (tmp_path / "private").mkdir() + scan = os.scandir + + def scandir(path: Any) -> Any: + if Path(path) == tmp_path / "private": + raise PermissionError("Synthetic unreadable selected subtree") + return scan(path) + + monkeypatch.setattr(os, "scandir", scandir) + with pytest.raises(PermissionError, match="selected subtree"): + if isinstance(sdk, AsyncOpenAI): + await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=include + ) + else: + sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=include + ) + assert not server.requests + + async def test_directory_walk_error_is_not_a_partial_success( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: From 234d5c296a16212659f06d87cd1be5ea49fd4973 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 20:16:45 +0000 Subject: [PATCH 13/18] refactor(agents): expand file globs one directory segment at a time --- src/openai/lib/beta/agents/_files.py | 79 ++++++++++++------------ tests/lib/streaming/agents/test_files.py | 67 ++++++++++++-------- 2 files changed, 81 insertions(+), 65 deletions(-) diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index 5e1daa6a7c..c108a4c058 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -1,8 +1,9 @@ from __future__ import annotations +import os from io import IOBase -from os import PathLike, walk, fstat -from stat import S_ISREG +from os import PathLike, fstat +from stat import S_ISDIR, S_ISREG from typing import TYPE_CHECKING, Mapping, BinaryIO, Sequence, Generator, cast from fnmatch import fnmatchcase from pathlib import Path, PurePosixPath @@ -115,21 +116,31 @@ def _prepare_selection( return selected -def _matches(parts: tuple[str, ...], pattern: tuple[str, ...], *, prefix: bool = False) -> bool: - # A directory prefix is useful only if the pattern can select a descendant. - if prefix and not parts: - return bool(pattern) +def _glob_files(current: Path, pattern: tuple[str, ...], root: Path) -> Generator[Path, None, None]: + try: + metadata = current.lstat() + except (FileNotFoundError, NotADirectoryError): + return if not pattern: - return not parts - if pattern[0] == "**": - return ( - _matches(parts, pattern[1:], prefix=prefix) or bool(parts) and _matches(parts[1:], pattern, prefix=prefix) - ) - return bool(parts) and fnmatchcase(parts[0], pattern[0]) and _matches(parts[1:], pattern[1:], prefix=prefix) - - -def _walk_error(error: OSError) -> None: - raise error + yield current + return + if not S_ISDIR(metadata.st_mode): # Do not follow directory symlinks. + return + if not current.resolve().is_relative_to(root): + raise ValueError("Selected directory traversal left its root") + segment, rest = pattern[0], pattern[1:] + if segment == "**": + yield from _glob_files(current, rest, root) + elif not any(character in segment for character in "*?["): + yield from _glob_files(current / segment, rest, root) + return + # Unlike glob/pathlib.glob, propagate errors for directories we need to read. + with os.scandir(current) as entries: + children = sorted(Path(entry.path) for entry in entries if segment == "**" or fnmatchcase(entry.name, segment)) + for child in children: + if segment == "**" and not rest: + yield child + yield from _glob_files(child, pattern if segment == "**" else rest, root) def directory_files(root: str | PathLike[str], destination: str, include: Sequence[str]) -> dict[str, PathLike[str]]: @@ -147,32 +158,20 @@ def directory_files(root: str | PathLike[str], destination: str, include: Sequen raise ValueError("Include patterns must stay inside the selected directory") patterns.append(PurePosixPath(pattern).parts) selected: dict[str, PathLike[str]] = {} - for current, directories, names in walk(directory, followlinks=False, onerror=_walk_error): - if not Path(current).resolve().is_relative_to(directory): - raise ValueError("Selected directory traversal left its root") - directories.sort() - for name in sorted([*directories, *names]): - source = Path(current) / name - relative = source.relative_to(directory) - chosen = any(_matches(relative.parts, pattern) for pattern in patterns) + for pattern in patterns: + for source in _glob_files(directory, pattern, directory): if source.is_symlink(): - if chosen: - raise ValueError("Selected agent files must not be symlinks") - if name in directories: - directories.remove(name) - continue - if name in directories: - if not any(_matches(relative.parts, pattern, prefix=True) for pattern in patterns): - directories.remove(name) + raise ValueError("Selected agent files must not be symlinks") + if source.is_dir(): continue - if chosen: - canonical = source.resolve() - if not canonical.is_relative_to(directory): - raise ValueError("Selected agent file is outside the chosen directory") - metadata = canonical.lstat() - selected[str(PurePosixPath(destination) / relative.as_posix())] = _SelectedFile( - canonical, (metadata.st_dev, metadata.st_ino, metadata.st_size) - ) + canonical = source.resolve() + if not canonical.is_relative_to(directory): + raise ValueError("Selected agent file is outside the chosen directory") + metadata = canonical.lstat() + relative = source.relative_to(directory) + selected[str(PurePosixPath(destination) / relative.as_posix())] = _SelectedFile( + canonical, (metadata.st_dev, metadata.st_ino, metadata.st_size) + ) if not selected: raise ValueError("The include patterns did not select any files") return selected diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 6d7d9d3756..621bd6e552 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -507,6 +507,39 @@ def scandir(path: Any) -> Any: assert tmp_path / "docs/nested" not in visited +async def test_directory_selection_expands_only_wildcard_segments( + sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + import os + + reports = tmp_path / "docs/team-a/reports" + reports.mkdir(parents=True) + (reports / "one.txt").write_text("abc") + (tmp_path / "docs/private").mkdir() + (tmp_path / "docs/team-a/unrelated").mkdir() + scan = os.scandir + visited: list[Path] = [] + + def scandir(path: Any) -> Any: + current = Path(path) + visited.append(current) + if current.name in ("private", "unrelated"): + raise PermissionError("Synthetic unrelated subtree") + return scan(path) + + monkeypatch.setattr(os, "scandir", scandir) + if isinstance(sdk, AsyncOpenAI): + result = await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["docs/team-*/reports/*.txt"] + ) + else: + result = sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["docs/team-*/reports/*.txt"] + ) + assert [item["path"] for item in result.files] == ["/workspace/docs/docs/team-a/reports/one.txt"] + assert visited == [tmp_path / "docs", reports] + + @pytest.mark.parametrize("include", [["**/*.txt"], ["private/*.txt"]]) async def test_directory_selection_propagates_needed_subtree_errors( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, include: list[str] @@ -535,27 +568,10 @@ def scandir(path: Any) -> Any: assert not server.requests -async def test_directory_walk_error_is_not_a_partial_success( - sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch -) -> None: - from openai.lib.beta.agents import _files - - (tmp_path / "keep.txt").write_text("abc") - - def failed_walk(root: Any, **options: Any) -> Any: - yield str(root), [], ["keep.txt"] - options["onerror"](PermissionError("synthetic unreadable subtree")) - - monkeypatch.setattr(_files, "walk", failed_walk) - with pytest.raises(PermissionError, match="synthetic"): - await prepare_directory(sdk, tmp_path) - assert not server.requests - - async def test_cached_directory_entry_cannot_leave_selected_root( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: - from openai.lib.beta.agents import _files + import os root = tmp_path / "chosen" child = root / "child" @@ -563,15 +579,16 @@ async def test_cached_directory_entry_cannot_leave_selected_root( outside = tmp_path / "outside" outside.mkdir() (outside / "other.txt").write_text("outside") + scan = os.scandir - def changed_walk(directory: Any, **_options: Any) -> Any: - yield str(directory), ["child"], [] - child.rename(root / "original-child") - child.symlink_to(outside, target_is_directory=True) - yield str(child), [], ["other.txt"] + def changed_scan(path: Any) -> Any: + if Path(path) == child: + child.rename(root / "original-child") + child.symlink_to(outside, target_is_directory=True) + return scan(path) - monkeypatch.setattr(_files, "walk", changed_walk) - with pytest.raises(ValueError, match="left its root"): + monkeypatch.setattr(os, "scandir", changed_scan) + with pytest.raises(ValueError, match="outside the chosen directory"): await prepare_directory(sdk, root) assert not server.requests From 94688b732440b30ea228e51f356c63a675717bf2 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 20:19:28 +0000 Subject: [PATCH 14/18] refactor(agents): use standard pathlib glob semantics for file selection --- helpers.md | 1 + src/openai/lib/beta/agents/_files.py | 44 ++++----------------- tests/lib/streaming/agents/test_files.py | 50 ++++++++++++------------ 3 files changed, 34 insertions(+), 61 deletions(-) diff --git a/helpers.md b/helpers.md index e6a0a57a6a..159001783f 100644 --- a/helpers.md +++ b/helpers.md @@ -663,6 +663,7 @@ artifact = artifacts.download( Use `prepare_directory("docs", destination="/workspace/docs", include=["**/*.md"])` for a selected directory snapshot, or `files.upload(environment_id, file=Path(...), path="/workspace/source.pdf")` to stage a file in an existing environment. +Directory selection follows `Path.glob` semantics, including skipping unreadable directories. Uploads remain caller-owned: use `prepared.uploaded_file_ids` with the ordinary Files API when ready to delete them. Preparation errors expose partial uploads through `error.prepared`. A batch of multiple uploads cannot share one explicit diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index c108a4c058..e0a1fd592a 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -1,11 +1,9 @@ from __future__ import annotations -import os from io import IOBase from os import PathLike, fstat -from stat import S_ISDIR, S_ISREG +from stat import S_ISREG from typing import TYPE_CHECKING, Mapping, BinaryIO, Sequence, Generator, cast -from fnmatch import fnmatchcase from pathlib import Path, PurePosixPath from contextlib import ExitStack, contextmanager from dataclasses import field, dataclass @@ -116,33 +114,6 @@ def _prepare_selection( return selected -def _glob_files(current: Path, pattern: tuple[str, ...], root: Path) -> Generator[Path, None, None]: - try: - metadata = current.lstat() - except (FileNotFoundError, NotADirectoryError): - return - if not pattern: - yield current - return - if not S_ISDIR(metadata.st_mode): # Do not follow directory symlinks. - return - if not current.resolve().is_relative_to(root): - raise ValueError("Selected directory traversal left its root") - segment, rest = pattern[0], pattern[1:] - if segment == "**": - yield from _glob_files(current, rest, root) - elif not any(character in segment for character in "*?["): - yield from _glob_files(current / segment, rest, root) - return - # Unlike glob/pathlib.glob, propagate errors for directories we need to read. - with os.scandir(current) as entries: - children = sorted(Path(entry.path) for entry in entries if segment == "**" or fnmatchcase(entry.name, segment)) - for child in children: - if segment == "**" and not rest: - yield child - yield from _glob_files(child, pattern if segment == "**" else rest, root) - - def directory_files(root: str | PathLike[str], destination: str, include: Sequence[str]) -> dict[str, PathLike[str]]: directory = Path(root).absolute() if isinstance(include, str) or not include: @@ -152,16 +123,17 @@ def directory_files(root: str | PathLike[str], destination: str, include: Sequen directory = directory.resolve() # System aliases such as macOS /tmp are allowed above the chosen root. _destination(destination + "/a") # Validate with the shortest possible filename. - patterns: list[tuple[str, ...]] = [] for pattern in include: if not pattern or PurePosixPath(pattern).is_absolute() or ".." in PurePosixPath(pattern).parts: raise ValueError("Include patterns must stay inside the selected directory") - patterns.append(PurePosixPath(pattern).parts) selected: dict[str, PathLike[str]] = {} - for pattern in patterns: - for source in _glob_files(directory, pattern, directory): - if source.is_symlink(): - raise ValueError("Selected agent files must not be symlinks") + for pattern in include: + for source in sorted(directory.glob(pattern)): + selected_path = source + while selected_path != directory: + if selected_path.is_symlink(): + raise ValueError("Selected agent files must not be symlinks") + selected_path = selected_path.parent if source.is_dir(): continue canonical = source.resolve() diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 621bd6e552..80b9004c3a 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -469,7 +469,7 @@ async def prepare_directory(sdk: OpenAI | AsyncOpenAI, root: Path, destination: (["one.txt"], ["one.txt"]), (["*.txt"], ["one.txt"]), (["docs/*.txt"], ["docs/one.txt"]), - (["docs/**/*.txt"], ["docs/one.txt", "docs/nested/one.txt"]), + (["docs/**/*.txt"], ["docs/nested/one.txt", "docs/one.txt"]), (["docs/*/*.txt"], ["docs/nested/one.txt"]), ], ) @@ -537,11 +537,12 @@ def scandir(path: Any) -> Any: tmp_path, destination="/workspace/docs", include=["docs/team-*/reports/*.txt"] ) assert [item["path"] for item in result.files] == ["/workspace/docs/docs/team-a/reports/one.txt"] - assert visited == [tmp_path / "docs", reports] + assert tmp_path / "docs/private" not in visited + assert tmp_path / "docs/team-a/unrelated" not in visited -@pytest.mark.parametrize("include", [["**/*.txt"], ["private/*.txt"]]) -async def test_directory_selection_propagates_needed_subtree_errors( +@pytest.mark.parametrize("include", [["**/*.txt"], ["private/*.txt", "one.txt"]]) +async def test_directory_selection_follows_glob_unreadable_subtree_semantics( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, include: list[str] ) -> None: import os @@ -556,39 +557,38 @@ def scandir(path: Any) -> Any: return scan(path) monkeypatch.setattr(os, "scandir", scandir) - with pytest.raises(PermissionError, match="selected subtree"): - if isinstance(sdk, AsyncOpenAI): - await sdk.beta.agents.environments.files.prepare_directory( - tmp_path, destination="/workspace/docs", include=include - ) - else: - sdk.beta.agents.environments.files.prepare_directory( - tmp_path, destination="/workspace/docs", include=include - ) - assert not server.requests + if isinstance(sdk, AsyncOpenAI): + result = await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=include + ) + else: + result = sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=include + ) + assert [item["path"] for item in result.files] == ["/workspace/docs/one.txt"] + assert server.uploads == 1 -async def test_cached_directory_entry_cannot_leave_selected_root( +async def test_selected_directory_cannot_be_replaced_by_symlink( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: - import os - root = tmp_path / "chosen" child = root / "child" child.mkdir(parents=True) + (child / "other.txt").write_text("inside") outside = tmp_path / "outside" outside.mkdir() (outside / "other.txt").write_text("outside") - scan = os.scandir + glob = Path.glob - def changed_scan(path: Any) -> Any: - if Path(path) == child: - child.rename(root / "original-child") - child.symlink_to(outside, target_is_directory=True) - return scan(path) + def replaced_glob(path: Path, pattern: str) -> Any: + selected = list(glob(path, pattern)) + child.rename(root / "original-child") + child.symlink_to(outside, target_is_directory=True) + yield from selected - monkeypatch.setattr(os, "scandir", changed_scan) - with pytest.raises(ValueError, match="outside the chosen directory"): + monkeypatch.setattr(Path, "glob", replaced_glob) + with pytest.raises(ValueError, match="symlink"): await prepare_directory(sdk, root) assert not server.requests From b8653b43fcfafe20ba0a76c7009a2560c612e07d Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 20:21:21 +0000 Subject: [PATCH 15/18] test(agents): consolidate directory selection scenarios --- tests/lib/streaming/agents/test_files.py | 82 ++++-------------------- 1 file changed, 14 insertions(+), 68 deletions(-) diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index 80b9004c3a..da77700dec 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -212,13 +212,6 @@ async def stage() -> Any: assert all(request.method == "POST" for request in server.requests) -async def test_file_count_is_left_to_the_api(sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path) -> None: - source = tmp_path / "source.txt" - source.write_text("abc") - result = await prepare(sdk, {f"/workspace/{index}": source for index in range(51)}) - assert len(result.files) == server.uploads == 51 - - @pytest.mark.parametrize("count", [1, 2]) async def test_large_prepared_files_reach_the_api( sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, count: int @@ -463,87 +456,40 @@ async def prepare_directory(sdk: OpenAI | AsyncOpenAI, root: Path, destination: return sdk.beta.agents.environments.files.prepare_directory(root, destination=destination, include=["**/*.txt"]) -@pytest.mark.parametrize( - "include,expected", - [ - (["one.txt"], ["one.txt"]), - (["*.txt"], ["one.txt"]), - (["docs/*.txt"], ["docs/one.txt"]), - (["docs/**/*.txt"], ["docs/nested/one.txt", "docs/one.txt"]), - (["docs/*/*.txt"], ["docs/nested/one.txt"]), - ], -) -async def test_directory_selection_skips_unrelated_unreadable_subtrees( - sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, include: list[str], expected: list[str] +@pytest.mark.parametrize("include", ["one.txt", "docs/team-*/reports/*.txt"]) +async def test_directory_selects_requested_files_without_unrelated_subtrees( + sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, include: str ) -> None: import os - for name in ["one.txt", "docs/one.txt", "docs/nested/one.txt", "private/one.txt"]: + for name in ["one.txt", "docs/team-a/reports/one.txt"]: source = tmp_path / name source.parent.mkdir(parents=True, exist_ok=True) source.write_text("abc") + for name in ["private", "docs/private", "docs/team-a/private"]: + (tmp_path / name).mkdir() scan = os.scandir - visited: list[Path] = [] def scandir(path: Any) -> Any: - current = Path(path) - visited.append(current) - if current == tmp_path / "private": - raise PermissionError("Synthetic unreadable unrelated subtree") - return scan(path) - - monkeypatch.setattr(os, "scandir", scandir) - if isinstance(sdk, AsyncOpenAI): - result = await sdk.beta.agents.environments.files.prepare_directory( - tmp_path, destination="/workspace/docs", include=include - ) - else: - result = sdk.beta.agents.environments.files.prepare_directory( - tmp_path, destination="/workspace/docs", include=include - ) - assert [file["path"] for file in result.files] == [f"/workspace/docs/{name}" for name in expected] - assert tmp_path / "private" not in visited - if include == ["docs/*.txt"]: - assert tmp_path / "docs/nested" not in visited - - -async def test_directory_selection_expands_only_wildcard_segments( - sdk: OpenAI | AsyncOpenAI, tmp_path: Path, monkeypatch: pytest.MonkeyPatch -) -> None: - import os - - reports = tmp_path / "docs/team-a/reports" - reports.mkdir(parents=True) - (reports / "one.txt").write_text("abc") - (tmp_path / "docs/private").mkdir() - (tmp_path / "docs/team-a/unrelated").mkdir() - scan = os.scandir - visited: list[Path] = [] - - def scandir(path: Any) -> Any: - current = Path(path) - visited.append(current) - if current.name in ("private", "unrelated"): + if Path(path).name == "private": raise PermissionError("Synthetic unrelated subtree") return scan(path) monkeypatch.setattr(os, "scandir", scandir) if isinstance(sdk, AsyncOpenAI): result = await sdk.beta.agents.environments.files.prepare_directory( - tmp_path, destination="/workspace/docs", include=["docs/team-*/reports/*.txt"] + tmp_path, destination="/workspace/docs", include=[include] ) else: result = sdk.beta.agents.environments.files.prepare_directory( - tmp_path, destination="/workspace/docs", include=["docs/team-*/reports/*.txt"] + tmp_path, destination="/workspace/docs", include=[include] ) - assert [item["path"] for item in result.files] == ["/workspace/docs/docs/team-a/reports/one.txt"] - assert tmp_path / "docs/private" not in visited - assert tmp_path / "docs/team-a/unrelated" not in visited + expected = "one.txt" if include == "one.txt" else "docs/team-a/reports/one.txt" + assert [item["path"] for item in result.files] == [f"/workspace/docs/{expected}"] -@pytest.mark.parametrize("include", [["**/*.txt"], ["private/*.txt", "one.txt"]]) async def test_directory_selection_follows_glob_unreadable_subtree_semantics( - sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, include: list[str] + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch ) -> None: import os @@ -559,11 +505,11 @@ def scandir(path: Any) -> Any: monkeypatch.setattr(os, "scandir", scandir) if isinstance(sdk, AsyncOpenAI): result = await sdk.beta.agents.environments.files.prepare_directory( - tmp_path, destination="/workspace/docs", include=include + tmp_path, destination="/workspace/docs", include=["**/*.txt"] ) else: result = sdk.beta.agents.environments.files.prepare_directory( - tmp_path, destination="/workspace/docs", include=include + tmp_path, destination="/workspace/docs", include=["**/*.txt"] ) assert [item["path"] for item in result.files] == ["/workspace/docs/one.txt"] assert server.uploads == 1 From 4e2b103abc35f9606667bb5af2b3d96978c487ed Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 20:29:55 +0000 Subject: [PATCH 16/18] fix(agents): snapshot async uploads one file at a time --- src/openai/lib/beta/agents/_files.py | 18 +++++++----- tests/lib/streaming/agents/test_files.py | 37 +++++++++++++++++++++--- 2 files changed, 44 insertions(+), 11 deletions(-) diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index e0a1fd592a..b746f0f121 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -193,13 +193,15 @@ def prepare(resource: Files, files: Mapping[str, str | PathLike[str]], options: def _snapshot_selection( files: Mapping[str, str | PathLike[str]], options: _RequestOptions, defaults: Headers -) -> list[tuple[str, str, bytes]]: +) -> list[tuple[str, _SelectedFile]]: # The worker owns every handle it opens, even if its caller is cancelled. with ExitStack() as stack: selected = _prepare_selection(files, options, defaults, stack) - return [ - (destination, source.name, _read_local(handle, length)) for destination, source, handle, length in selected - ] + paths: list[tuple[str, _SelectedFile]] = [] + for destination, source, handle, length in selected: + metadata = fstat(handle.fileno()) + paths.append((destination, _SelectedFile(source, (metadata.st_dev, metadata.st_ino, length)))) + return paths async def async_prepare( @@ -210,8 +212,10 @@ async def async_prepare( ) prepared = PreparedAgentFiles() try: - for destination, filename, content in selected: - uploaded = await resource._client.files.create(file=(filename, content), purpose="user_data", **options) + for destination, source in selected: + content = await run_sync(_snapshot_upload, source, abandon_on_cancel=True) + uploaded = await resource._client.files.create(file=content, purpose="user_data", **options) + del content prepared.uploaded_file_ids.append(uploaded.id) prepared.files.append({"type": "file_id", "file_id": uploaded.id, "path": destination}) except Exception as error: @@ -244,7 +248,7 @@ def _upload_content(file: FileTypes) -> Generator[FileTypes, None, None]: with ExitStack() as stack: if isinstance(content, PathLike): path = _local_file(content) - handle, length = _open_local(path, stack) + handle, length = _open_local(path, stack, content.identity if isinstance(content, _SelectedFile) else None) _unchanged_size(handle, length) yield cast(FileTypes, (file[0], handle, *file[2:])) if isinstance(file, tuple) else (path.name, handle) return diff --git a/tests/lib/streaming/agents/test_files.py b/tests/lib/streaming/agents/test_files.py index da77700dec..a27ada2a57 100644 --- a/tests/lib/streaming/agents/test_files.py +++ b/tests/lib/streaming/agents/test_files.py @@ -320,7 +320,7 @@ async def test_system_alias_ancestor_is_allowed_but_directory_entries_are_not_fo assert server.uploads == 1 -async def test_prepared_files_keep_opened_sources_during_batch_upload( +async def test_batch_replacement_never_uploads_changed_sources( sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path ) -> None: first = tmp_path / "first.txt" @@ -334,9 +334,38 @@ def replace_later_source() -> None: second.write_text("replacement") server.after_upload = replace_later_source - await prepare(sdk, {"/workspace/first": first, "/workspace/second": second}) - assert b"second-original" in server.requests[1].content - assert b"replacement" not in server.requests[1].content + if isinstance(sdk, AsyncOpenAI): + with pytest.raises(AgentFilePreparationError) as caught: + await prepare(sdk, {"/workspace/first": first, "/workspace/second": second}) + assert caught.value.prepared.uploaded_file_ids == ["file_1"] + assert isinstance(caught.value.__cause__, ValueError) + assert server.uploads == 1 + else: + await prepare(sdk, {"/workspace/first": first, "/workspace/second": second}) + assert b"second-original" in server.requests[1].content + assert b"replacement" not in server.requests[1].content + + +async def test_async_prepare_reads_each_file_only_when_ready_to_upload( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + from openai.lib.beta.agents import _files + + if not isinstance(sdk, AsyncOpenAI): + pytest.skip("Async buffering contract") + source = tmp_path / "source.txt" + source.write_text("abc") + read = _files._read_local + uploads_at_read: list[int] = [] + + def read_one(handle: Any, length: int) -> bytes: + uploads_at_read.append(server.uploads) + return read(handle, length) + + monkeypatch.setattr(_files, "_read_local", read_one) + result = await sdk.beta.agents.environments.files.prepare({f"/workspace/{index}.txt": source for index in range(3)}) + assert uploads_at_read == [0, 1, 2] + assert result.uploaded_file_ids == ["file_1", "file_2", "file_3"] @pytest.mark.parametrize("file_backed", [False, True]) From 0aff3cc0d098eb68503d5594d7d5a3949d9f91ef Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 20:35:44 +0000 Subject: [PATCH 17/18] docs(agents): note future batch upload opportunity --- src/openai/lib/beta/agents/_files.py | 2 ++ 1 file changed, 2 insertions(+) diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index b746f0f121..67f6c031a6 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -212,6 +212,8 @@ async def async_prepare( ) prepared = PreparedAgentFiles() try: + # Upload sequentially for now. TODO: use batch/archive uploads (e.g. ZIP) + # when the API supports them. for destination, source in selected: content = await run_sync(_snapshot_upload, source, abandon_on_cancel=True) uploaded = await resource._client.files.create(file=content, purpose="user_data", **options) From 51dfacadd53a9a755fb125eb130262a3e07c9301 Mon Sep 17 00:00:00 2001 From: Alex Chang Date: Thu, 1 Oct 2026 20:35:59 +0000 Subject: [PATCH 18/18] docs(agents): clarify application-owned archive orchestration --- src/openai/lib/beta/agents/_files.py | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/src/openai/lib/beta/agents/_files.py b/src/openai/lib/beta/agents/_files.py index 67f6c031a6..15e494edb5 100644 --- a/src/openai/lib/beta/agents/_files.py +++ b/src/openai/lib/beta/agents/_files.py @@ -212,8 +212,9 @@ async def async_prepare( ) prepared = PreparedAgentFiles() try: - # Upload sequentially for now. TODO: use batch/archive uploads (e.g. ZIP) - # when the API supports them. + # TODO: Use API batch uploads when available. Applications can own preflight + # and archive (e.g. ZIP) upload/extraction; that orchestration is outside + # these helpers, which upload sequentially for now. for destination, source in selected: content = await run_sync(_snapshot_upload, source, abandon_on_cancel=True) uploaded = await resource._client.files.create(file=content, purpose="user_data", **options)