diff --git a/helpers.md b/helpers.md index 54902af655..159001783f 100644 --- a/helpers.md +++ b/helpers.md @@ -633,3 +633,42 @@ 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() + +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") +) +``` + +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 +`Idempotency-Key`. + +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. diff --git a/src/openai/lib/beta/agents/__init__.py b/src/openai/lib/beta/agents/__init__.py index 6b004a47e7..80ca2a8b4e 100644 --- a/src/openai/lib/beta/agents/__init__.py +++ b/src/openai/lib/beta/agents/__init__.py @@ -1,5 +1,11 @@ """Beta helpers for hosted Agents API tools and turn results.""" +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, @@ -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..002da07295 --- /dev/null +++ b/src/openai/lib/beta/agents/_artifacts.py @@ -0,0 +1,157 @@ +from __future__ import annotations + +from os import PathLike +from typing import TYPE_CHECKING + +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: + 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 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, + *, + 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.""" + 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: + 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") + 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 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, + *, + 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.""" + 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: + 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") + 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..15e494edb5 --- /dev/null +++ b/src/openai/lib/beta/agents/_files.py @@ -0,0 +1,297 @@ +from __future__ import annotations + +from io import IOBase +from os import PathLike, fstat +from stat import S_ISREG +from typing import TYPE_CHECKING, Mapping, BinaryIO, Sequence, Generator, cast +from pathlib import Path, PurePosixPath +from contextlib import ExitStack, contextmanager +from dataclasses import field, dataclass +from typing_extensions import TypedDict, override + +import httpx2 +from anyio.to_thread import run_sync + +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 ....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 + + +@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.""" + + 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") + return path + + +def _prepare_selection( + files: Mapping[str, str | PathLike[str]], options: _RequestOptions, defaults: Headers, stack: ExitStack +) -> list[tuple[str, Path, BinaryIO, int]]: + 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") + 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") + 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)) + return selected + + +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") + 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 + "/a") # Validate with the shortest possible filename. + 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") + selected: dict[str, PathLike[str]] = {} + 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() + 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 + + +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, 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): + raise ValueError("Selected agent file changed during preparation") + 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: + 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 + except BaseException as error: + error.__dict__["prepared"] = prepared + raise + return prepared + + +def _snapshot_selection( + files: Mapping[str, str | PathLike[str]], options: _RequestOptions, defaults: Headers +) -> 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) + 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( + resource: AsyncFiles, files: Mapping[str, str | PathLike[str]], options: _RequestOptions +) -> PreparedAgentFiles: + selected = await run_sync( + _snapshot_selection, files, options, resource._client.default_headers, abandon_on_cancel=True + ) + prepared = PreparedAgentFiles() + try: + # 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) + 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: + raise AgentFilePreparationError(prepared) from error + except BaseException as error: + error.__dict__["prepared"] = prepared + raise + return prepared + + +def _read_local(handle: BinaryIO, length: int) -> bytes: + _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: + 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 +def _upload_content(file: FileTypes) -> Generator[FileTypes, None, None]: + content = file[1] if isinstance(file, tuple) else file + with ExitStack() as stack: + if isinstance(content, PathLike): + path = _local_file(content) + 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 + yield file + + +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) + 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 + except BaseException as error: + error.__dict__["uploaded_file_id"] = uploaded.id + raise + 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: + if not environment_id: + 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: + 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 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/src/openai/resources/beta/agents/environments/files.py b/src/openai/resources/beta/agents/environments/files.py index b2ad8913bf..4c35468df7 100644 --- a/src/openai/resources/beta/agents/environments/files.py +++ b/src/openai/resources/beta/agents/environments/files.py @@ -2,19 +2,30 @@ 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, + async_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 +33,62 @@ class Files(SyncAPIResource): + def prepare( + self, + 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, "extra_query": extra_query, "extra_body": extra_body, "timeout": timeout}, + ) + + def prepare_directory( + self, + root: str | PathLike[str], + *, + 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, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, + ) + + def upload( + self, + environment_id: str, + *, + 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, "extra_query": extra_query, "extra_body": extra_body, "timeout": timeout}, + ) + @cached_property def with_raw_response(self) -> FilesWithRawResponse: """ @@ -222,6 +289,62 @@ def list( class AsyncFiles(AsyncAPIResource): + async def prepare( + self, + 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, "extra_query": extra_query, "extra_body": extra_body, "timeout": timeout}, + ) + + async def prepare_directory( + self, + root: str | PathLike[str], + *, + 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( + await async_directory_files(root, destination, include), + extra_headers=extra_headers, + extra_query=extra_query, + extra_body=extra_body, + timeout=timeout, + ) + + async def upload( + self, + environment_id: str, + *, + 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, "extra_query": extra_query, "extra_body": extra_body, "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..e2df2bc44c --- /dev/null +++ b/tests/lib/streaming/agents/test_artifacts.py @@ -0,0 +1,166 @@ +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"}, + 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"}, + extra_query={"synthetic": "query"}, + extra_body={"synthetic": "body"}, + 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.url.params["synthetic"] == "query" 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 + + +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 new file mode 100644 index 0000000000..a27ada2a57 --- /dev/null +++ b/tests/lib/streaming/agents/test_files.py @@ -0,0 +1,754 @@ +from __future__ import annotations + +import json +import asyncio +from io import BytesIO +from typing import Any, Callable +from pathlib import Path +from typing_extensions import override + +import httpx2 +import pytest + +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 + + +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 + 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 asyncio.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 asyncio.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"}, + 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( + 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 + + +@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: + 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", + 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", + 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() + assert vars(caught.value)["uploaded_file_id"] == "file_2" + assert all(request.method == "POST" for request in server.requests) + + +@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(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( + 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(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) + + +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(asyncio.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 + + +async def test_batch_replacement_never_uploads_changed_sources( + 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 + 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]) +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(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 + + +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"] + + +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"]) + + +@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/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 + + def scandir(path: Any) -> Any: + 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=[include] + ) + else: + result = sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=[include] + ) + 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}"] + + +async def test_directory_selection_follows_glob_unreadable_subtree_semantics( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> 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) + if isinstance(sdk, AsyncOpenAI): + result = await sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["**/*.txt"] + ) + else: + result = sdk.beta.agents.environments.files.prepare_directory( + tmp_path, destination="/workspace/docs", include=["**/*.txt"] + ) + assert [item["path"] for item in result.files] == ["/workspace/docs/one.txt"] + assert server.uploads == 1 + + +async def test_selected_directory_cannot_be_replaced_by_symlink( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch +) -> None: + 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") + glob = Path.glob + + 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(Path, "glob", replaced_glob) + with pytest.raises(ValueError, match="symlink"): + 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_preserves_long_destination_for_api_validation( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path +) -> None: + (tmp_path / "a").write_text("abc") + destination = "/workspace/" + "x" * 4096 + 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 prepared.files[0]["path"] == destination + "/a" + 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_interruption_preserves_observed_uploads( + sdk: OpenAI | AsyncOpenAI, server: FilesServer, tmp_path: Path, monkeypatch: pytest.MonkeyPatch, staging: bool +) -> None: + 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: + 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: + + def after_upload() -> None: + if server.uploads == 2: + raise KeyboardInterrupt() + + server.after_upload = after_upload + with pytest.raises(KeyboardInterrupt) as caught: + await prepare(sdk, {"/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() + + +@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 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(asyncio.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