From 7964543521b34a2961ce692cf816ed97a9cc1f03 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 10:48:48 -0700 Subject: [PATCH 01/18] exec-server: stream files in chunks --- codex-rs/exec-server/README.md | 3 + codex-rs/exec-server/src/client.rs | 132 +++++++++++++++ codex-rs/exec-server/src/file_read.rs | 140 +++++++++++++++ codex-rs/exec-server/src/lib.rs | 2 + codex-rs/exec-server/src/local_file_system.rs | 42 +++++ codex-rs/exec-server/src/protocol.rs | 41 +++++ .../src/server/file_system_handler.rs | 49 ++++++ codex-rs/exec-server/src/server/handler.rs | 31 ++++ codex-rs/exec-server/src/server/registry.rs | 24 +++ codex-rs/exec-server/tests/file_stream.rs | 159 ++++++++++++++++++ 10 files changed, 623 insertions(+) create mode 100644 codex-rs/exec-server/src/file_read.rs create mode 100644 codex-rs/exec-server/tests/file_stream.rs diff --git a/codex-rs/exec-server/README.md b/codex-rs/exec-server/README.md index 54b00ea92ec6..eff62ee62b0c 100644 --- a/codex-rs/exec-server/README.md +++ b/codex-rs/exec-server/README.md @@ -343,6 +343,8 @@ invalid or unavailable paths. For compatibility, requests also accept native absolute path strings and normalize them to `file:` URIs: - `fs/readFile` +- `fs/open`, `fs/readBlock`, and `fs/close` (internal transport for + `ExecServerClient::stream`) - `fs/writeFile` - `fs/createDirectory` - `fs/getMetadata` @@ -381,6 +383,7 @@ The crate exports: - `ExecServerClient` - `ExecServerError` +- `FileReadStream` - `ExecServerClientConnectOptions` - `RemoteExecServerConnectArgs` - protocol request/response structs for process and filesystem RPCs diff --git a/codex-rs/exec-server/src/client.rs b/codex-rs/exec-server/src/client.rs index 978c502b7d94..54f7709e2415 100644 --- a/codex-rs/exec-server/src/client.rs +++ b/codex-rs/exec-server/src/client.rs @@ -1,15 +1,22 @@ use std::collections::BTreeMap; use std::collections::HashMap; +use std::pin::Pin; use std::sync::Arc; use std::sync::Mutex as StdMutex; use std::sync::OnceLock; use std::sync::atomic::AtomicU64; +use std::task::Context; +use std::task::Poll; use std::time::Duration; use arc_swap::ArcSwap; +use bytes::Bytes; use codex_app_server_protocol::JSONRPCNotification; use futures::FutureExt; +use futures::Stream; +use futures::StreamExt; use futures::future::BoxFuture; +use futures::stream::BoxStream; use serde_json::Value; use tokio::sync::Mutex; use tokio::sync::Semaphore; @@ -45,21 +52,30 @@ use crate::protocol::ExecOutputDeltaNotification; use crate::protocol::ExecParams; use crate::protocol::ExecResponse; use crate::protocol::FS_CANONICALIZE_METHOD; +use crate::protocol::FS_CLOSE_METHOD; use crate::protocol::FS_COPY_METHOD; use crate::protocol::FS_CREATE_DIRECTORY_METHOD; use crate::protocol::FS_GET_METADATA_METHOD; +use crate::protocol::FS_OPEN_METHOD; +use crate::protocol::FS_READ_BLOCK_METHOD; use crate::protocol::FS_READ_DIRECTORY_METHOD; use crate::protocol::FS_READ_FILE_METHOD; use crate::protocol::FS_REMOVE_METHOD; use crate::protocol::FS_WRITE_FILE_METHOD; use crate::protocol::FsCanonicalizeParams; use crate::protocol::FsCanonicalizeResponse; +use crate::protocol::FsCloseParams; +use crate::protocol::FsCloseResponse; use crate::protocol::FsCopyParams; use crate::protocol::FsCopyResponse; use crate::protocol::FsCreateDirectoryParams; use crate::protocol::FsCreateDirectoryResponse; use crate::protocol::FsGetMetadataParams; use crate::protocol::FsGetMetadataResponse; +use crate::protocol::FsOpenParams; +use crate::protocol::FsOpenResponse; +use crate::protocol::FsReadBlockParams; +use crate::protocol::FsReadBlockResponse; use crate::protocol::FsReadDirectoryParams; use crate::protocol::FsReadDirectoryResponse; use crate::protocol::FsReadFileParams; @@ -95,6 +111,26 @@ const INITIALIZE_TIMEOUT: Duration = Duration::from_secs(10); const PROCESS_EVENT_CHANNEL_CAPACITY: usize = 256; const PROCESS_EVENT_RETAINED_BYTES: usize = 1024 * 1024; +/// Stream of immutable file blocks returned by [`ExecServerClient::stream`]. +pub struct FileReadStream { + inner: BoxStream<'static, Result>, +} + +impl Stream for FileReadStream { + type Item = Result; + + fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + self.inner.as_mut().poll_next(cx) + } +} + +struct FileReadRegistration { + client: ExecServerClient, + handle_id: String, + runtime: Option, + active: bool, +} + impl Default for ExecServerClientConnectOptions { fn default() -> Self { Self { @@ -429,6 +465,83 @@ impl ExecServerClient { self.call(FS_READ_FILE_METHOD, ¶ms).await } + /// Opens an unsandboxed file and returns a demand-driven stream of 1 MiB blocks. + pub async fn stream( + &self, + params: FsReadFileParams, + ) -> Result { + let response = self + .fs_open(FsOpenParams { + path: params.path, + sandbox: params.sandbox, + }) + .await?; + let registration = FileReadRegistration { + client: self.clone(), + handle_id: response.handle_id, + runtime: tokio::runtime::Handle::try_current().ok(), + active: true, + }; + Ok(FileReadStream { + inner: futures::stream::try_unfold(Some((registration, 0_u64)), |state| async move { + let Some((mut registration, offset)) = state else { + return Ok(None); + }; + let response = registration + .client + .fs_read_block(FsReadBlockParams { + handle_id: registration.handle_id.clone(), + offset, + len: crate::file_read::FILE_READ_BLOCK_SIZE, + }) + .await?; + let chunk = Bytes::from(response.chunk.into_inner()); + if chunk.len() > crate::file_read::FILE_READ_BLOCK_SIZE { + return Err(ExecServerError::Protocol(format!( + "{FS_READ_BLOCK_METHOD} returned {} bytes, maximum is {}", + chunk.len(), + crate::file_read::FILE_READ_BLOCK_SIZE + ))); + } + if response.eof { + registration.active = false; + return if chunk.is_empty() { + Ok(None) + } else { + Ok(Some((chunk, None))) + }; + } + if chunk.is_empty() { + return Err(ExecServerError::Protocol(format!( + "{FS_READ_BLOCK_METHOD} returned an empty non-terminal block" + ))); + } + let next_offset = offset.checked_add(chunk.len() as u64).ok_or_else(|| { + ExecServerError::Protocol(format!( + "{FS_READ_BLOCK_METHOD} offset overflowed after {offset} bytes" + )) + })?; + Ok(Some((chunk, Some((registration, next_offset))))) + }) + .boxed(), + }) + } + + async fn fs_open(&self, params: FsOpenParams) -> Result { + self.call(FS_OPEN_METHOD, ¶ms).await + } + + async fn fs_read_block( + &self, + params: FsReadBlockParams, + ) -> Result { + self.call(FS_READ_BLOCK_METHOD, ¶ms).await + } + + async fn fs_close(&self, params: FsCloseParams) -> Result { + self.call(FS_CLOSE_METHOD, ¶ms).await + } + pub async fn fs_write_file( &self, params: FsWriteFileParams, @@ -601,6 +714,25 @@ impl ExecServerClient { } } +impl Drop for FileReadRegistration { + fn drop(&mut self) { + if !self.active { + return; + } + let client = self.client.clone(); + let handle_id = self.handle_id.clone(); + let runtime = self + .runtime + .clone() + .or_else(|| tokio::runtime::Handle::try_current().ok()); + if let Some(runtime) = runtime { + runtime.spawn(async move { + let _ = client.fs_close(FsCloseParams { handle_id }).await; + }); + } + } +} + impl From for ExecServerError { fn from(value: RpcCallError) -> Self { match value { diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs new file mode 100644 index 000000000000..245e5d2b523e --- /dev/null +++ b/codex-rs/exec-server/src/file_read.rs @@ -0,0 +1,140 @@ +use std::collections::HashMap; +use std::io; +use std::io::SeekFrom; +use std::sync::Arc; + +use tokio::io::AsyncReadExt; +use tokio::io::AsyncSeekExt; +use tokio::sync::Mutex; +use tokio::sync::mpsc; +use tokio::sync::oneshot; +use uuid::Uuid; + +pub(crate) const FILE_READ_BLOCK_SIZE: usize = 1024 * 1024; +const MAX_OPEN_FILE_READS: usize = 32; + +#[derive(Debug, Eq, PartialEq)] +pub(crate) struct FileReadBlock { + pub(crate) bytes: Vec, + pub(crate) eof: bool, +} + +#[derive(Clone, Default)] +pub(crate) struct FileReadHandleManager { + handles: Arc>>>, +} + +struct FileReadRequest { + offset: u64, + len: usize, + response: oneshot::Sender>, +} + +impl FileReadHandleManager { + pub(crate) async fn open(&self, file: tokio::fs::File) -> io::Result { + let mut handles = self.handles.lock().await; + if handles.len() >= MAX_OPEN_FILE_READS { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!("at most {MAX_OPEN_FILE_READS} file reads may be open per connection"), + )); + } + let handle_id = Uuid::new_v4().to_string(); + let (sender, receiver) = mpsc::channel(1); + handles.insert(handle_id.clone(), sender); + tokio::spawn(serve_file_reads(file, receiver)); + Ok(handle_id) + } + + pub(crate) async fn read_block( + &self, + handle_id: &str, + offset: u64, + len: usize, + ) -> io::Result { + validate_read_block_len(len)?; + let sender = { + let handles = self.handles.lock().await; + handles + .get(handle_id) + .cloned() + .ok_or_else(|| unknown_handle_error(handle_id))? + }; + let (response, result) = oneshot::channel(); + if sender + .send(FileReadRequest { + offset, + len, + response, + }) + .await + .is_err() + { + self.close(handle_id).await; + return Err(unknown_handle_error(handle_id)); + } + let result = result.await.map_err(|_| { + io::Error::new( + io::ErrorKind::BrokenPipe, + format!("file read handle `{handle_id}` stopped unexpectedly"), + ) + })?; + if result.is_err() || result.as_ref().is_ok_and(|block| block.eof) { + self.close(handle_id).await; + } + result + } + + pub(crate) async fn close(&self, handle_id: &str) { + self.handles.lock().await.remove(handle_id); + } + + pub(crate) async fn close_all(&self) { + self.handles.lock().await.clear(); + } +} + +async fn serve_file_reads( + mut file: tokio::fs::File, + mut requests: mpsc::Receiver, +) { + while let Some(request) = requests.recv().await { + let result = read_block(&mut file, request.offset, request.len).await; + let should_stop = result.is_err(); + let _ = request.response.send(result); + if should_stop { + break; + } + } +} + +async fn read_block( + file: &mut tokio::fs::File, + offset: u64, + len: usize, +) -> io::Result { + file.seek(SeekFrom::Start(offset)).await?; + let mut bytes = Vec::with_capacity(len); + file.take(len as u64).read_to_end(&mut bytes).await?; + Ok(FileReadBlock { + eof: bytes.len() < len, + bytes, + }) +} + +fn validate_read_block_len(len: usize) -> io::Result<()> { + if !(1..=FILE_READ_BLOCK_SIZE).contains(&len) { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!("file read block length must be between 1 and {FILE_READ_BLOCK_SIZE}"), + )); + } + Ok(()) +} + +fn unknown_handle_error(handle_id: &str) -> io::Error { + io::Error::new( + io::ErrorKind::NotFound, + format!("unknown file read handle `{handle_id}`"), + ) +} diff --git a/codex-rs/exec-server/src/lib.rs b/codex-rs/exec-server/src/lib.rs index d6a9aacb5114..bb3364dd9a93 100644 --- a/codex-rs/exec-server/src/lib.rs +++ b/codex-rs/exec-server/src/lib.rs @@ -5,6 +5,7 @@ mod connection; mod environment; mod environment_provider; mod environment_toml; +mod file_read; mod fs_helper; mod fs_helper_main; mod fs_sandbox; @@ -25,6 +26,7 @@ mod server; pub use client::ExecServerClient; pub use client::ExecServerError; +pub use client::FileReadStream; pub use client::http_client::HttpResponseBodyStream; pub use client::http_client::ReqwestHttpClient; pub use client_api::ExecServerClientConnectOptions; diff --git a/codex-rs/exec-server/src/local_file_system.rs b/codex-rs/exec-server/src/local_file_system.rs index 34b3213ecb92..98730f4c5531 100644 --- a/codex-rs/exec-server/src/local_file_system.rs +++ b/codex-rs/exec-server/src/local_file_system.rs @@ -79,6 +79,20 @@ impl LocalFileSystem { } impl LocalFileSystem { + pub(crate) async fn open_file_for_read( + &self, + path: &PathUri, + sandbox: Option<&FileSystemSandboxContext>, + ) -> FileSystemResult { + if sandbox.is_some_and(FileSystemSandboxContext::should_run_in_sandbox) { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + "streaming file reads do not support platform sandboxing", + )); + } + self.unsandboxed.open_file_for_read(path, sandbox).await + } + async fn canonicalize( &self, path: &PathUri, @@ -239,6 +253,17 @@ impl ExecutorFileSystem for LocalFileSystem { } impl UnsandboxedFileSystem { + async fn open_file_for_read( + &self, + path: &PathUri, + sandbox: Option<&FileSystemSandboxContext>, + ) -> FileSystemResult { + reject_platform_sandbox_context(sandbox)?; + self.file_system + .open_file_for_read(path, /*sandbox*/ None) + .await + } + async fn canonicalize( &self, path: &PathUri, @@ -414,6 +439,23 @@ impl ExecutorFileSystem for UnsandboxedFileSystem { } impl DirectFileSystem { + async fn open_file_for_read( + &self, + path: &PathUri, + sandbox: Option<&FileSystemSandboxContext>, + ) -> FileSystemResult { + reject_sandbox_context(sandbox)?; + let path = path.to_abs_path()?; + let file = tokio::fs::File::open(path.as_path()).await?; + if !file.metadata().await?.is_file() { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!("path `{}` is not a file", path.display()), + )); + } + Ok(file) + } + async fn canonicalize( &self, path: &PathUri, diff --git a/codex-rs/exec-server/src/protocol.rs b/codex-rs/exec-server/src/protocol.rs index ee116ef08e2b..d840898b8339 100644 --- a/codex-rs/exec-server/src/protocol.rs +++ b/codex-rs/exec-server/src/protocol.rs @@ -21,6 +21,9 @@ pub const EXEC_EXITED_METHOD: &str = "process/exited"; pub const EXEC_CLOSED_METHOD: &str = "process/closed"; pub const ENVIRONMENT_INFO_METHOD: &str = "environment/info"; pub const FS_READ_FILE_METHOD: &str = "fs/readFile"; +pub(crate) const FS_OPEN_METHOD: &str = "fs/open"; +pub(crate) const FS_READ_BLOCK_METHOD: &str = "fs/readBlock"; +pub(crate) const FS_CLOSE_METHOD: &str = "fs/close"; pub const FS_WRITE_FILE_METHOD: &str = "fs/writeFile"; pub const FS_CREATE_DIRECTORY_METHOD: &str = "fs/createDirectory"; pub const FS_GET_METADATA_METHOD: &str = "fs/getMetadata"; @@ -210,6 +213,44 @@ pub struct FsReadFileResponse { pub data_base64: String, } +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct FsOpenParams { + pub path: PathUri, + pub sandbox: Option, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct FsOpenResponse { + pub handle_id: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct FsReadBlockParams { + pub handle_id: String, + pub offset: u64, + pub len: usize, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct FsReadBlockResponse { + pub chunk: ByteChunk, + pub eof: bool, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct FsCloseParams { + pub handle_id: String, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "camelCase")] +pub(crate) struct FsCloseResponse {} + #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub struct FsWriteFileParams { diff --git a/codex-rs/exec-server/src/server/file_system_handler.rs b/codex-rs/exec-server/src/server/file_system_handler.rs index 080d4829d089..f8704d63c763 100644 --- a/codex-rs/exec-server/src/server/file_system_handler.rs +++ b/codex-rs/exec-server/src/server/file_system_handler.rs @@ -9,16 +9,23 @@ use crate::CreateDirectoryOptions; use crate::ExecServerRuntimePaths; use crate::ExecutorFileSystem; use crate::RemoveOptions; +use crate::file_read::FileReadHandleManager; use crate::local_file_system::LocalFileSystem; use crate::protocol::FS_WRITE_FILE_METHOD; use crate::protocol::FsCanonicalizeParams; use crate::protocol::FsCanonicalizeResponse; +use crate::protocol::FsCloseParams; +use crate::protocol::FsCloseResponse; use crate::protocol::FsCopyParams; use crate::protocol::FsCopyResponse; use crate::protocol::FsCreateDirectoryParams; use crate::protocol::FsCreateDirectoryResponse; use crate::protocol::FsGetMetadataParams; use crate::protocol::FsGetMetadataResponse; +use crate::protocol::FsOpenParams; +use crate::protocol::FsOpenResponse; +use crate::protocol::FsReadBlockParams; +use crate::protocol::FsReadBlockResponse; use crate::protocol::FsReadDirectoryEntry; use crate::protocol::FsReadDirectoryParams; use crate::protocol::FsReadDirectoryResponse; @@ -35,15 +42,57 @@ use crate::rpc::not_found; #[derive(Clone)] pub(crate) struct FileSystemHandler { file_system: LocalFileSystem, + file_reads: FileReadHandleManager, } impl FileSystemHandler { pub(crate) fn new(runtime_paths: ExecServerRuntimePaths) -> Self { Self { file_system: LocalFileSystem::with_runtime_paths(runtime_paths), + file_reads: FileReadHandleManager::default(), } } + pub(crate) async fn shutdown(&self) { + self.file_reads.close_all().await; + } + + pub(crate) async fn open( + &self, + params: FsOpenParams, + ) -> Result { + let file = self + .file_system + .open_file_for_read(¶ms.path, params.sandbox.as_ref()) + .await + .map_err(map_fs_error)?; + let handle_id = self.file_reads.open(file).await.map_err(map_fs_error)?; + Ok(FsOpenResponse { handle_id }) + } + + pub(crate) async fn read_block( + &self, + params: FsReadBlockParams, + ) -> Result { + let block = self + .file_reads + .read_block(¶ms.handle_id, params.offset, params.len) + .await + .map_err(map_fs_error)?; + Ok(FsReadBlockResponse { + chunk: block.bytes.into(), + eof: block.eof, + }) + } + + pub(crate) async fn close( + &self, + params: FsCloseParams, + ) -> Result { + self.file_reads.close(¶ms.handle_id).await; + Ok(FsCloseResponse {}) + } + pub(crate) async fn read_file( &self, params: FsReadFileParams, diff --git a/codex-rs/exec-server/src/server/handler.rs b/codex-rs/exec-server/src/server/handler.rs index 73e8f22684e6..e0ede94c7742 100644 --- a/codex-rs/exec-server/src/server/handler.rs +++ b/codex-rs/exec-server/src/server/handler.rs @@ -19,12 +19,18 @@ use crate::protocol::ExecParams; use crate::protocol::ExecResponse; use crate::protocol::FsCanonicalizeParams; use crate::protocol::FsCanonicalizeResponse; +use crate::protocol::FsCloseParams; +use crate::protocol::FsCloseResponse; use crate::protocol::FsCopyParams; use crate::protocol::FsCopyResponse; use crate::protocol::FsCreateDirectoryParams; use crate::protocol::FsCreateDirectoryResponse; use crate::protocol::FsGetMetadataParams; use crate::protocol::FsGetMetadataResponse; +use crate::protocol::FsOpenParams; +use crate::protocol::FsOpenResponse; +use crate::protocol::FsReadBlockParams; +use crate::protocol::FsReadBlockResponse; use crate::protocol::FsReadDirectoryParams; use crate::protocol::FsReadDirectoryResponse; use crate::protocol::FsReadFileParams; @@ -87,6 +93,7 @@ impl ExecServerHandler { self.background_task_shutdown.cancel(); self.background_tasks.close(); self.background_tasks.wait().await; + self.file_system.shutdown().await; if let Some(session) = self.session() { session.detach().await; } @@ -234,6 +241,30 @@ impl ExecServerHandler { self.file_system.read_file(params).await } + pub(crate) async fn fs_open( + &self, + params: FsOpenParams, + ) -> Result { + self.require_initialized_for("filesystem")?; + self.file_system.open(params).await + } + + pub(crate) async fn fs_read_block( + &self, + params: FsReadBlockParams, + ) -> Result { + self.require_initialized_for("filesystem")?; + self.file_system.read_block(params).await + } + + pub(crate) async fn fs_close( + &self, + params: FsCloseParams, + ) -> Result { + self.require_initialized_for("filesystem")?; + self.file_system.close(params).await + } + pub(crate) async fn fs_write_file( &self, params: FsWriteFileParams, diff --git a/codex-rs/exec-server/src/server/registry.rs b/codex-rs/exec-server/src/server/registry.rs index 4ba7ce865195..8f48aeaf99b7 100644 --- a/codex-rs/exec-server/src/server/registry.rs +++ b/codex-rs/exec-server/src/server/registry.rs @@ -8,17 +8,23 @@ use crate::protocol::EXEC_TERMINATE_METHOD; use crate::protocol::EXEC_WRITE_METHOD; use crate::protocol::ExecParams; use crate::protocol::FS_CANONICALIZE_METHOD; +use crate::protocol::FS_CLOSE_METHOD; use crate::protocol::FS_COPY_METHOD; use crate::protocol::FS_CREATE_DIRECTORY_METHOD; use crate::protocol::FS_GET_METADATA_METHOD; +use crate::protocol::FS_OPEN_METHOD; +use crate::protocol::FS_READ_BLOCK_METHOD; use crate::protocol::FS_READ_DIRECTORY_METHOD; use crate::protocol::FS_READ_FILE_METHOD; use crate::protocol::FS_REMOVE_METHOD; use crate::protocol::FS_WRITE_FILE_METHOD; use crate::protocol::FsCanonicalizeParams; +use crate::protocol::FsCloseParams; use crate::protocol::FsCopyParams; use crate::protocol::FsCreateDirectoryParams; use crate::protocol::FsGetMetadataParams; +use crate::protocol::FsOpenParams; +use crate::protocol::FsReadBlockParams; use crate::protocol::FsReadDirectoryParams; use crate::protocol::FsReadFileParams; use crate::protocol::FsRemoveParams; @@ -93,6 +99,24 @@ pub(crate) fn build_router() -> RpcRouter { handler.fs_read_file(params).await }, ); + router.request( + FS_OPEN_METHOD, + |handler: Arc, params: FsOpenParams| async move { + handler.fs_open(params).await + }, + ); + router.request( + FS_READ_BLOCK_METHOD, + |handler: Arc, params: FsReadBlockParams| async move { + handler.fs_read_block(params).await + }, + ); + router.request( + FS_CLOSE_METHOD, + |handler: Arc, params: FsCloseParams| async move { + handler.fs_close(params).await + }, + ); router.request( FS_WRITE_FILE_METHOD, |handler: Arc, params: FsWriteFileParams| async move { diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs new file mode 100644 index 000000000000..44e014f1e7f8 --- /dev/null +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -0,0 +1,159 @@ +mod common; + +use anyhow::Result; +use codex_exec_server::ExecServerClient; +use codex_exec_server::FileSystemSandboxContext; +use codex_exec_server::FsReadFileParams; +use codex_exec_server::RemoteExecServerConnectArgs; +use codex_protocol::models::PermissionProfile; +use codex_protocol::permissions::FileSystemAccessMode; +use codex_protocol::permissions::FileSystemPath; +use codex_protocol::permissions::FileSystemSandboxEntry; +use codex_protocol::permissions::FileSystemSandboxPolicy; +use codex_protocol::permissions::NetworkSandboxPolicy; +use codex_utils_absolute_path::AbsolutePathBuf; +use codex_utils_path_uri::PathUri; +use futures::TryStreamExt; +use pretty_assertions::assert_eq; +use tempfile::TempDir; + +use crate::common::exec_server::exec_server; + +const BLOCK_SIZE: usize = 1024 * 1024; + +#[tokio::test] +async fn stream_reads_file_in_one_mib_blocks() -> Result<()> { + let server = exec_server().await?; + let client = connect_client(server.websocket_url()).await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("blocks.bin"); + let contents = (0..BLOCK_SIZE * 2 + 17) + .map(|index| (index % 251) as u8) + .collect::>(); + std::fs::write(&path, &contents)?; + + let chunks = client + .stream(FsReadFileParams { + path: PathUri::from_path(path)?, + sandbox: None, + }) + .await? + .try_collect::>() + .await?; + + assert_eq!( + chunks.iter().map(bytes::Bytes::len).collect::>(), + vec![BLOCK_SIZE, BLOCK_SIZE, 17] + ); + assert_eq!( + chunks + .iter() + .flat_map(|chunk| chunk.iter().copied()) + .collect::>(), + contents + ); + Ok(()) +} + +#[tokio::test] +async fn stream_stops_after_an_exact_block_boundary() -> Result<()> { + let server = exec_server().await?; + let client = connect_client(server.websocket_url()).await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("exact-blocks.bin"); + std::fs::write(&path, vec![b'x'; BLOCK_SIZE * 2])?; + + let chunks = client + .stream(FsReadFileParams { + path: PathUri::from_path(path)?, + sandbox: None, + }) + .await? + .try_collect::>() + .await?; + + assert_eq!( + chunks.iter().map(bytes::Bytes::len).collect::>(), + vec![BLOCK_SIZE, BLOCK_SIZE] + ); + Ok(()) +} + +#[tokio::test] +async fn stream_rejects_platform_sandbox() -> Result<()> { + let server = exec_server().await?; + let client = connect_client(server.websocket_url()).await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("sandboxed.txt"); + std::fs::write(&path, "sandboxed hello")?; + + let result = client + .stream(FsReadFileParams { + path: PathUri::from_path(&path)?, + sandbox: Some(read_only_sandbox(tmp.path().to_path_buf())), + }) + .await; + + let Err(error) = result else { + panic!("sandboxed stream should be rejected"); + }; + assert_eq!( + error.to_string(), + "exec-server rejected request (-32600): streaming file reads do not support platform sandboxing" + ); + Ok(()) +} + +#[cfg(unix)] +#[tokio::test] +async fn stream_keeps_reading_the_open_file_after_path_replacement() -> Result<()> { + let server = exec_server().await?; + let client = connect_client(server.websocket_url()).await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("replaceable.bin"); + std::fs::write(&path, vec![b'a'; BLOCK_SIZE + 1])?; + let mut stream = client + .stream(FsReadFileParams { + path: PathUri::from_path(&path)?, + sandbox: None, + }) + .await?; + + assert_eq!( + stream.try_next().await?, + Some(bytes::Bytes::from(vec![b'a'; BLOCK_SIZE])) + ); + let replacement = tmp.path().join("replacement.bin"); + std::fs::write(&replacement, vec![b'b'; BLOCK_SIZE + 1])?; + std::fs::remove_file(&path)?; + std::fs::rename(replacement, &path)?; + + assert_eq!( + stream.try_next().await?, + Some(bytes::Bytes::from_static(b"a")) + ); + assert_eq!(stream.try_next().await?, None); + Ok(()) +} + +async fn connect_client(websocket_url: &str) -> Result { + Ok( + ExecServerClient::connect_websocket(RemoteExecServerConnectArgs::new( + websocket_url.to_string(), + "file-stream-test".to_string(), + )) + .await?, + ) +} + +fn read_only_sandbox(path: std::path::PathBuf) -> FileSystemSandboxContext { + let path = AbsolutePathBuf::from_absolute_path(&path) + .unwrap_or_else(|err| panic!("sandbox path should be absolute: {err}")); + FileSystemSandboxContext::from_permission_profile(PermissionProfile::from_runtime_permissions( + &FileSystemSandboxPolicy::restricted(vec![FileSystemSandboxEntry { + path: FileSystemPath::Path { path }, + access: FileSystemAccessMode::Read, + }]), + NetworkSandboxPolicy::Restricted, + )) +} From 8108e89897df4e74a73e5c1a6dfdeef88c3cc6dc Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 11:01:34 -0700 Subject: [PATCH 02/18] exec-server: reuse request loop for file reads --- codex-rs/exec-server/src/file_read.rs | 62 +++++---------------------- 1 file changed, 11 insertions(+), 51 deletions(-) diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs index 245e5d2b523e..b56cf240e210 100644 --- a/codex-rs/exec-server/src/file_read.rs +++ b/codex-rs/exec-server/src/file_read.rs @@ -6,8 +6,6 @@ use std::sync::Arc; use tokio::io::AsyncReadExt; use tokio::io::AsyncSeekExt; use tokio::sync::Mutex; -use tokio::sync::mpsc; -use tokio::sync::oneshot; use uuid::Uuid; pub(crate) const FILE_READ_BLOCK_SIZE: usize = 1024 * 1024; @@ -21,13 +19,7 @@ pub(crate) struct FileReadBlock { #[derive(Clone, Default)] pub(crate) struct FileReadHandleManager { - handles: Arc>>>, -} - -struct FileReadRequest { - offset: u64, - len: usize, - response: oneshot::Sender>, + handles: Arc>>, } impl FileReadHandleManager { @@ -40,9 +32,7 @@ impl FileReadHandleManager { )); } let handle_id = Uuid::new_v4().to_string(); - let (sender, receiver) = mpsc::channel(1); - handles.insert(handle_id.clone(), sender); - tokio::spawn(serve_file_reads(file, receiver)); + handles.insert(handle_id.clone(), file); Ok(handle_id) } @@ -53,34 +43,18 @@ impl FileReadHandleManager { len: usize, ) -> io::Result { validate_read_block_len(len)?; - let sender = { - let handles = self.handles.lock().await; + let mut file = { + let mut handles = self.handles.lock().await; handles - .get(handle_id) - .cloned() + .remove(handle_id) .ok_or_else(|| unknown_handle_error(handle_id))? }; - let (response, result) = oneshot::channel(); - if sender - .send(FileReadRequest { - offset, - len, - response, - }) - .await - .is_err() - { - self.close(handle_id).await; - return Err(unknown_handle_error(handle_id)); - } - let result = result.await.map_err(|_| { - io::Error::new( - io::ErrorKind::BrokenPipe, - format!("file read handle `{handle_id}` stopped unexpectedly"), - ) - })?; - if result.is_err() || result.as_ref().is_ok_and(|block| block.eof) { - self.close(handle_id).await; + let result = read_block(&mut file, offset, len).await; + if result.as_ref().is_ok_and(|block| !block.eof) { + self.handles + .lock() + .await + .insert(handle_id.to_string(), file); } result } @@ -94,20 +68,6 @@ impl FileReadHandleManager { } } -async fn serve_file_reads( - mut file: tokio::fs::File, - mut requests: mpsc::Receiver, -) { - while let Some(request) = requests.recv().await { - let result = read_block(&mut file, request.offset, request.len).await; - let should_stop = result.is_err(); - let _ = request.response.send(result); - if should_stop { - break; - } - } -} - async fn read_block( file: &mut tokio::fs::File, offset: u64, From ae377302de3a7d3551bcd8b3c7e31f146bcfa369 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 11:05:07 -0700 Subject: [PATCH 03/18] exec-server: test file read handles --- codex-rs/exec-server/src/file_read.rs | 4 ++ codex-rs/exec-server/src/file_read_tests.rs | 75 +++++++++++++++++++++ 2 files changed, 79 insertions(+) create mode 100644 codex-rs/exec-server/src/file_read_tests.rs diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs index b56cf240e210..6f53807b9a3b 100644 --- a/codex-rs/exec-server/src/file_read.rs +++ b/codex-rs/exec-server/src/file_read.rs @@ -98,3 +98,7 @@ fn unknown_handle_error(handle_id: &str) -> io::Error { format!("unknown file read handle `{handle_id}`"), ) } + +#[cfg(test)] +#[path = "file_read_tests.rs"] +mod tests; diff --git a/codex-rs/exec-server/src/file_read_tests.rs b/codex-rs/exec-server/src/file_read_tests.rs new file mode 100644 index 000000000000..fb3f62fac9f7 --- /dev/null +++ b/codex-rs/exec-server/src/file_read_tests.rs @@ -0,0 +1,75 @@ +use std::io; + +use anyhow::Result; +use pretty_assertions::assert_eq; +use tempfile::TempDir; + +use super::FileReadBlock; +use super::FileReadHandleManager; +use super::MAX_OPEN_FILE_READS; + +#[tokio::test] +async fn reads_blocks_at_non_sequential_offsets() -> Result<()> { + let tmp = TempDir::new()?; + let path = tmp.path().join("non-sequential.bin"); + std::fs::write(&path, b"0123456789")?; + let manager = FileReadHandleManager::default(); + let handle_id = manager.open(tokio::fs::File::open(&path).await?).await?; + + assert_eq!( + manager + .read_block(&handle_id, /*offset*/ 6, /*len*/ 3) + .await?, + FileReadBlock { + bytes: b"678".to_vec(), + eof: false, + } + ); + assert_eq!( + manager + .read_block(&handle_id, /*offset*/ 1, /*len*/ 2) + .await?, + FileReadBlock { + bytes: b"12".to_vec(), + eof: false, + } + ); + assert_eq!( + manager + .read_block(&handle_id, /*offset*/ 8, /*len*/ 4) + .await?, + FileReadBlock { + bytes: b"89".to_vec(), + eof: true, + } + ); + Ok(()) +} + +#[tokio::test] +async fn limits_open_files_and_releases_capacity_on_close() -> Result<()> { + let tmp = TempDir::new()?; + let path = tmp.path().join("limited.bin"); + std::fs::write(&path, b"limited")?; + let manager = FileReadHandleManager::default(); + let mut handles = Vec::with_capacity(MAX_OPEN_FILE_READS); + for _ in 0..MAX_OPEN_FILE_READS { + handles.push(manager.open(tokio::fs::File::open(&path).await?).await?); + } + + let error = manager + .open(tokio::fs::File::open(&path).await?) + .await + .expect_err("opening beyond the limit should fail"); + assert_eq!( + (error.kind(), error.to_string()), + ( + io::ErrorKind::InvalidInput, + format!("at most {MAX_OPEN_FILE_READS} file reads may be open per connection"), + ) + ); + + manager.close(&handles[0]).await; + manager.open(tokio::fs::File::open(&path).await?).await?; + Ok(()) +} From 741b2b9213b060213beac26b7964baaf70f9438f Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 11:08:39 -0700 Subject: [PATCH 04/18] exec-server: test file reads over rpc --- codex-rs/exec-server/src/file_read.rs | 4 - codex-rs/exec-server/src/file_read_tests.rs | 75 -------- codex-rs/exec-server/tests/file_stream.rs | 188 ++++++++++++++++++++ 3 files changed, 188 insertions(+), 79 deletions(-) delete mode 100644 codex-rs/exec-server/src/file_read_tests.rs diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs index 6f53807b9a3b..b56cf240e210 100644 --- a/codex-rs/exec-server/src/file_read.rs +++ b/codex-rs/exec-server/src/file_read.rs @@ -98,7 +98,3 @@ fn unknown_handle_error(handle_id: &str) -> io::Error { format!("unknown file read handle `{handle_id}`"), ) } - -#[cfg(test)] -#[path = "file_read_tests.rs"] -mod tests; diff --git a/codex-rs/exec-server/src/file_read_tests.rs b/codex-rs/exec-server/src/file_read_tests.rs deleted file mode 100644 index fb3f62fac9f7..000000000000 --- a/codex-rs/exec-server/src/file_read_tests.rs +++ /dev/null @@ -1,75 +0,0 @@ -use std::io; - -use anyhow::Result; -use pretty_assertions::assert_eq; -use tempfile::TempDir; - -use super::FileReadBlock; -use super::FileReadHandleManager; -use super::MAX_OPEN_FILE_READS; - -#[tokio::test] -async fn reads_blocks_at_non_sequential_offsets() -> Result<()> { - let tmp = TempDir::new()?; - let path = tmp.path().join("non-sequential.bin"); - std::fs::write(&path, b"0123456789")?; - let manager = FileReadHandleManager::default(); - let handle_id = manager.open(tokio::fs::File::open(&path).await?).await?; - - assert_eq!( - manager - .read_block(&handle_id, /*offset*/ 6, /*len*/ 3) - .await?, - FileReadBlock { - bytes: b"678".to_vec(), - eof: false, - } - ); - assert_eq!( - manager - .read_block(&handle_id, /*offset*/ 1, /*len*/ 2) - .await?, - FileReadBlock { - bytes: b"12".to_vec(), - eof: false, - } - ); - assert_eq!( - manager - .read_block(&handle_id, /*offset*/ 8, /*len*/ 4) - .await?, - FileReadBlock { - bytes: b"89".to_vec(), - eof: true, - } - ); - Ok(()) -} - -#[tokio::test] -async fn limits_open_files_and_releases_capacity_on_close() -> Result<()> { - let tmp = TempDir::new()?; - let path = tmp.path().join("limited.bin"); - std::fs::write(&path, b"limited")?; - let manager = FileReadHandleManager::default(); - let mut handles = Vec::with_capacity(MAX_OPEN_FILE_READS); - for _ in 0..MAX_OPEN_FILE_READS { - handles.push(manager.open(tokio::fs::File::open(&path).await?).await?); - } - - let error = manager - .open(tokio::fs::File::open(&path).await?) - .await - .expect_err("opening beyond the limit should fail"); - assert_eq!( - (error.kind(), error.to_string()), - ( - io::ErrorKind::InvalidInput, - format!("at most {MAX_OPEN_FILE_READS} file reads may be open per connection"), - ) - ); - - manager.close(&handles[0]).await; - manager.open(tokio::fs::File::open(&path).await?).await?; - Ok(()) -} diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index 44e014f1e7f8..dadd188912ca 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -1,9 +1,14 @@ mod common; use anyhow::Result; +use base64::Engine as _; +use codex_app_server_protocol::JSONRPCError; +use codex_app_server_protocol::JSONRPCMessage; +use codex_app_server_protocol::JSONRPCResponse; use codex_exec_server::ExecServerClient; use codex_exec_server::FileSystemSandboxContext; use codex_exec_server::FsReadFileParams; +use codex_exec_server::InitializeParams; use codex_exec_server::RemoteExecServerConnectArgs; use codex_protocol::models::PermissionProfile; use codex_protocol::permissions::FileSystemAccessMode; @@ -15,11 +20,27 @@ use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; use futures::TryStreamExt; use pretty_assertions::assert_eq; +use serde::Deserialize; +use serde::de::DeserializeOwned; use tempfile::TempDir; +use crate::common::exec_server::ExecServerHarness; use crate::common::exec_server::exec_server; const BLOCK_SIZE: usize = 1024 * 1024; +const OPEN_FILE_LIMIT: usize = 32; + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "camelCase")] +struct OpenFileResponse { + handle_id: String, +} + +#[derive(Debug, Deserialize)] +struct ReadBlockResponse { + chunk: String, + eof: bool, +} #[tokio::test] async fn stream_reads_file_in_one_mib_blocks() -> Result<()> { @@ -136,6 +157,173 @@ async fn stream_keeps_reading_the_open_file_after_path_replacement() -> Result<( Ok(()) } +#[tokio::test] +async fn read_block_supports_non_sequential_offsets_and_lengths() -> Result<()> { + let mut server = exec_server().await?; + initialize_exec_server(&mut server).await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("non-sequential.bin"); + std::fs::write(&path, b"0123456789")?; + let open: OpenFileResponse = rpc_call( + &mut server, + "fs/open", + serde_json::json!({ + "path": PathUri::from_path(path)?, + "sandbox": null, + }), + ) + .await?; + + assert_eq!( + read_block( + &mut server, + &open.handle_id, + /*offset*/ 6, + /*len*/ 3 + ) + .await?, + (b"678".to_vec(), false) + ); + assert_eq!( + read_block( + &mut server, + &open.handle_id, + /*offset*/ 1, + /*len*/ 2 + ) + .await?, + (b"12".to_vec(), false) + ); + assert_eq!( + read_block( + &mut server, + &open.handle_id, + /*offset*/ 8, + /*len*/ 4 + ) + .await?, + (b"89".to_vec(), true) + ); + server.shutdown().await?; + Ok(()) +} + +#[tokio::test] +async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> Result<()> { + let mut server = exec_server().await?; + initialize_exec_server(&mut server).await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("limited.bin"); + std::fs::write(&path, b"limited")?; + let path = PathUri::from_path(path)?; + let mut handles = Vec::with_capacity(OPEN_FILE_LIMIT); + for _ in 0..OPEN_FILE_LIMIT { + let open: OpenFileResponse = rpc_call( + &mut server, + "fs/open", + serde_json::json!({ "path": path, "sandbox": null }), + ) + .await?; + handles.push(open.handle_id); + } + + let response = rpc_message( + &mut server, + "fs/open", + serde_json::json!({ "path": path, "sandbox": null }), + ) + .await?; + let JSONRPCMessage::Error(JSONRPCError { error, .. }) = response else { + anyhow::bail!("expected opening beyond the limit to fail, got {response:?}"); + }; + assert_eq!( + (error.code, error.message), + ( + -32600, + format!("at most {OPEN_FILE_LIMIT} file reads may be open per connection"), + ) + ); + + let _: serde_json::Value = rpc_call( + &mut server, + "fs/close", + serde_json::json!({ "handleId": handles[0] }), + ) + .await?; + let _: OpenFileResponse = rpc_call( + &mut server, + "fs/open", + serde_json::json!({ "path": path, "sandbox": null }), + ) + .await?; + server.shutdown().await?; + Ok(()) +} + +async fn initialize_exec_server(server: &mut ExecServerHarness) -> Result<()> { + let _: serde_json::Value = rpc_call( + server, + "initialize", + serde_json::to_value(InitializeParams { + client_name: "file-stream-protocol-test".to_string(), + resume_session_id: None, + })?, + ) + .await?; + server + .send_notification("initialized", serde_json::json!({})) + .await?; + Ok(()) +} + +async fn read_block( + server: &mut ExecServerHarness, + handle_id: &str, + offset: u64, + len: usize, +) -> Result<(Vec, bool)> { + let response: ReadBlockResponse = rpc_call( + server, + "fs/readBlock", + serde_json::json!({ "handleId": handle_id, "offset": offset, "len": len }), + ) + .await?; + Ok(( + base64::engine::general_purpose::STANDARD.decode(response.chunk)?, + response.eof, + )) +} + +async fn rpc_call( + server: &mut ExecServerHarness, + method: &str, + params: serde_json::Value, +) -> Result +where + T: DeserializeOwned, +{ + let response = rpc_message(server, method, params).await?; + let JSONRPCMessage::Response(JSONRPCResponse { result, .. }) = response else { + anyhow::bail!("expected successful `{method}` response, got {response:?}"); + }; + Ok(serde_json::from_value(result)?) +} + +async fn rpc_message( + server: &mut ExecServerHarness, + method: &str, + params: serde_json::Value, +) -> Result { + let request_id = server.send_request(method, params).await?; + server + .wait_for_event(|event| match event { + JSONRPCMessage::Response(JSONRPCResponse { id, .. }) + | JSONRPCMessage::Error(JSONRPCError { id, .. }) => id == &request_id, + JSONRPCMessage::Request(_) | JSONRPCMessage::Notification(_) => false, + }) + .await +} + async fn connect_client(websocket_url: &str) -> Result { Ok( ExecServerClient::connect_websocket(RemoteExecServerConnectArgs::new( From f2c9b00242865e7eedc036f26b9c1891aeb3088c Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 11:18:05 -0700 Subject: [PATCH 05/18] exec-server: allow more open file reads --- codex-rs/exec-server/src/file_read.rs | 2 +- codex-rs/exec-server/tests/file_stream.rs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs index b56cf240e210..bc15a9adf98a 100644 --- a/codex-rs/exec-server/src/file_read.rs +++ b/codex-rs/exec-server/src/file_read.rs @@ -9,7 +9,7 @@ use tokio::sync::Mutex; use uuid::Uuid; pub(crate) const FILE_READ_BLOCK_SIZE: usize = 1024 * 1024; -const MAX_OPEN_FILE_READS: usize = 32; +const MAX_OPEN_FILE_READS: usize = 128; #[derive(Debug, Eq, PartialEq)] pub(crate) struct FileReadBlock { diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index dadd188912ca..e091f3fa1532 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -28,7 +28,7 @@ use crate::common::exec_server::ExecServerHarness; use crate::common::exec_server::exec_server; const BLOCK_SIZE: usize = 1024 * 1024; -const OPEN_FILE_LIMIT: usize = 32; +const OPEN_FILE_LIMIT: usize = 128; #[derive(Debug, Deserialize)] #[serde(rename_all = "camelCase")] From 66c112502af10af4f488b01f1b1fb503e2001da5 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 11:23:23 -0700 Subject: [PATCH 06/18] exec-server: use positional file reads --- codex-rs/exec-server/src/file_read.rs | 64 ++++++++++++++++++--------- 1 file changed, 42 insertions(+), 22 deletions(-) diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs index bc15a9adf98a..4e45f053306d 100644 --- a/codex-rs/exec-server/src/file_read.rs +++ b/codex-rs/exec-server/src/file_read.rs @@ -1,10 +1,8 @@ use std::collections::HashMap; +use std::fs::File; use std::io; -use std::io::SeekFrom; use std::sync::Arc; -use tokio::io::AsyncReadExt; -use tokio::io::AsyncSeekExt; use tokio::sync::Mutex; use uuid::Uuid; @@ -19,11 +17,12 @@ pub(crate) struct FileReadBlock { #[derive(Clone, Default)] pub(crate) struct FileReadHandleManager { - handles: Arc>>, + handles: Arc>>>, } impl FileReadHandleManager { pub(crate) async fn open(&self, file: tokio::fs::File) -> io::Result { + let file = Arc::new(file.into_std().await); let mut handles = self.handles.lock().await; if handles.len() >= MAX_OPEN_FILE_READS { return Err(io::Error::new( @@ -43,18 +42,22 @@ impl FileReadHandleManager { len: usize, ) -> io::Result { validate_read_block_len(len)?; - let mut file = { - let mut handles = self.handles.lock().await; + let file = { + let handles = self.handles.lock().await; handles - .remove(handle_id) + .get(handle_id) + .cloned() .ok_or_else(|| unknown_handle_error(handle_id))? }; - let result = read_block(&mut file, offset, len).await; - if result.as_ref().is_ok_and(|block| !block.eof) { - self.handles - .lock() - .await - .insert(handle_id.to_string(), file); + let result = + match tokio::task::spawn_blocking(move || read_block_at(&file, offset, len)).await { + Ok(result) => result, + Err(error) => Err(io::Error::other(format!( + "file read task stopped unexpectedly: {error}" + ))), + }; + if result.is_err() || result.as_ref().is_ok_and(|block| block.eof) { + self.close(handle_id).await; } result } @@ -68,20 +71,37 @@ impl FileReadHandleManager { } } -async fn read_block( - file: &mut tokio::fs::File, - offset: u64, - len: usize, -) -> io::Result { - file.seek(SeekFrom::Start(offset)).await?; - let mut bytes = Vec::with_capacity(len); - file.take(len as u64).read_to_end(&mut bytes).await?; +fn read_block_at(file: &File, offset: u64, len: usize) -> io::Result { + let mut bytes = vec![0; len]; + let mut bytes_read = 0; + while bytes_read < len { + let read_offset = offset.checked_add(bytes_read as u64).ok_or_else(|| { + io::Error::new(io::ErrorKind::InvalidInput, "file read offset overflowed") + })?; + match read_file_at(file, &mut bytes[bytes_read..], read_offset) { + Ok(0) => break, + Ok(read) => bytes_read += read, + Err(error) if error.kind() == io::ErrorKind::Interrupted => {} + Err(error) => return Err(error), + } + } + bytes.truncate(bytes_read); Ok(FileReadBlock { - eof: bytes.len() < len, + eof: bytes_read < len, bytes, }) } +#[cfg(unix)] +fn read_file_at(file: &File, bytes: &mut [u8], offset: u64) -> io::Result { + std::os::unix::fs::FileExt::read_at(file, bytes, offset) +} + +#[cfg(windows)] +fn read_file_at(file: &File, bytes: &mut [u8], offset: u64) -> io::Result { + std::os::windows::fs::FileExt::seek_read(file, bytes, offset) +} + fn validate_read_block_len(len: usize) -> io::Result<()> { if !(1..=FILE_READ_BLOCK_SIZE).contains(&len) { return Err(io::Error::new( From 6fb792556141aaf015e4dee59a8ced1def0ea3a4 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 11:39:21 -0700 Subject: [PATCH 07/18] exec-server: assign file handles on the client --- codex-rs/exec-server/src/client.rs | 15 +++++++------- codex-rs/exec-server/src/file_read.rs | 14 ++++++++++--- codex-rs/exec-server/src/protocol.rs | 1 + .../src/server/file_system_handler.rs | 6 +++++- codex-rs/exec-server/tests/file_stream.rs | 20 ++++++++++++++++--- 5 files changed, 42 insertions(+), 14 deletions(-) diff --git a/codex-rs/exec-server/src/client.rs b/codex-rs/exec-server/src/client.rs index 54f7709e2415..d14a18839425 100644 --- a/codex-rs/exec-server/src/client.rs +++ b/codex-rs/exec-server/src/client.rs @@ -25,6 +25,7 @@ use tokio::sync::watch; use tokio::time::timeout; use tracing::debug; +use uuid::Uuid; use crate::ProcessId; use crate::client_api::ExecServerClientConnectOptions; @@ -470,18 +471,18 @@ impl ExecServerClient { &self, params: FsReadFileParams, ) -> Result { - let response = self - .fs_open(FsOpenParams { - path: params.path, - sandbox: params.sandbox, - }) - .await?; let registration = FileReadRegistration { client: self.clone(), - handle_id: response.handle_id, + handle_id: Uuid::new_v4().to_string(), runtime: tokio::runtime::Handle::try_current().ok(), active: true, }; + self.fs_open(FsOpenParams { + handle_id: registration.handle_id.clone(), + path: params.path, + sandbox: params.sandbox, + }) + .await?; Ok(FileReadStream { inner: futures::stream::try_unfold(Some((registration, 0_u64)), |state| async move { let Some((mut registration, offset)) = state else { diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs index 4e45f053306d..26a4afb84ea0 100644 --- a/codex-rs/exec-server/src/file_read.rs +++ b/codex-rs/exec-server/src/file_read.rs @@ -4,7 +4,6 @@ use std::io; use std::sync::Arc; use tokio::sync::Mutex; -use uuid::Uuid; pub(crate) const FILE_READ_BLOCK_SIZE: usize = 1024 * 1024; const MAX_OPEN_FILE_READS: usize = 128; @@ -21,16 +20,25 @@ pub(crate) struct FileReadHandleManager { } impl FileReadHandleManager { - pub(crate) async fn open(&self, file: tokio::fs::File) -> io::Result { + pub(crate) async fn open( + &self, + handle_id: String, + file: tokio::fs::File, + ) -> io::Result { let file = Arc::new(file.into_std().await); let mut handles = self.handles.lock().await; + if handles.contains_key(&handle_id) { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!("file read handle `{handle_id}` already exists"), + )); + } if handles.len() >= MAX_OPEN_FILE_READS { return Err(io::Error::new( io::ErrorKind::InvalidInput, format!("at most {MAX_OPEN_FILE_READS} file reads may be open per connection"), )); } - let handle_id = Uuid::new_v4().to_string(); handles.insert(handle_id.clone(), file); Ok(handle_id) } diff --git a/codex-rs/exec-server/src/protocol.rs b/codex-rs/exec-server/src/protocol.rs index d840898b8339..8ebd7e5eaf11 100644 --- a/codex-rs/exec-server/src/protocol.rs +++ b/codex-rs/exec-server/src/protocol.rs @@ -216,6 +216,7 @@ pub struct FsReadFileResponse { #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] pub(crate) struct FsOpenParams { + pub handle_id: String, pub path: PathUri, pub sandbox: Option, } diff --git a/codex-rs/exec-server/src/server/file_system_handler.rs b/codex-rs/exec-server/src/server/file_system_handler.rs index f8704d63c763..603d50735533 100644 --- a/codex-rs/exec-server/src/server/file_system_handler.rs +++ b/codex-rs/exec-server/src/server/file_system_handler.rs @@ -66,7 +66,11 @@ impl FileSystemHandler { .open_file_for_read(¶ms.path, params.sandbox.as_ref()) .await .map_err(map_fs_error)?; - let handle_id = self.file_reads.open(file).await.map_err(map_fs_error)?; + let handle_id = self + .file_reads + .open(params.handle_id, file) + .await + .map_err(map_fs_error)?; Ok(FsOpenResponse { handle_id }) } diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index e091f3fa1532..422fe6d7a776 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -23,6 +23,7 @@ use pretty_assertions::assert_eq; use serde::Deserialize; use serde::de::DeserializeOwned; use tempfile::TempDir; +use uuid::Uuid; use crate::common::exec_server::ExecServerHarness; use crate::common::exec_server::exec_server; @@ -168,6 +169,7 @@ async fn read_block_supports_non_sequential_offsets_and_lengths() -> Result<()> &mut server, "fs/open", serde_json::json!({ + "handleId": Uuid::new_v4().to_string(), "path": PathUri::from_path(path)?, "sandbox": null, }), @@ -221,7 +223,11 @@ async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> let open: OpenFileResponse = rpc_call( &mut server, "fs/open", - serde_json::json!({ "path": path, "sandbox": null }), + serde_json::json!({ + "handleId": Uuid::new_v4().to_string(), + "path": path, + "sandbox": null, + }), ) .await?; handles.push(open.handle_id); @@ -230,7 +236,11 @@ async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> let response = rpc_message( &mut server, "fs/open", - serde_json::json!({ "path": path, "sandbox": null }), + serde_json::json!({ + "handleId": Uuid::new_v4().to_string(), + "path": path, + "sandbox": null, + }), ) .await?; let JSONRPCMessage::Error(JSONRPCError { error, .. }) = response else { @@ -253,7 +263,11 @@ async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> let _: OpenFileResponse = rpc_call( &mut server, "fs/open", - serde_json::json!({ "path": path, "sandbox": null }), + serde_json::json!({ + "handleId": Uuid::new_v4().to_string(), + "path": path, + "sandbox": null, + }), ) .await?; server.shutdown().await?; From 5ab7aba16a2fc8c995c45dda9741568940ee0540 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 11:51:34 -0700 Subject: [PATCH 08/18] exec-server: reject special files before streaming --- codex-rs/exec-server/src/local_file_system.rs | 5 +-- codex-rs/exec-server/tests/file_stream.rs | 40 +++++++++++++++++++ 2 files changed, 42 insertions(+), 3 deletions(-) diff --git a/codex-rs/exec-server/src/local_file_system.rs b/codex-rs/exec-server/src/local_file_system.rs index 98730f4c5531..ce1d3761db3c 100644 --- a/codex-rs/exec-server/src/local_file_system.rs +++ b/codex-rs/exec-server/src/local_file_system.rs @@ -446,14 +446,13 @@ impl DirectFileSystem { ) -> FileSystemResult { reject_sandbox_context(sandbox)?; let path = path.to_abs_path()?; - let file = tokio::fs::File::open(path.as_path()).await?; - if !file.metadata().await?.is_file() { + if !tokio::fs::metadata(path.as_path()).await?.is_file() { return Err(io::Error::new( io::ErrorKind::InvalidInput, format!("path `{}` is not a file", path.display()), )); } - Ok(file) + tokio::fs::File::open(path.as_path()).await } async fn canonicalize( diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index 422fe6d7a776..50cfba99a6f1 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -22,7 +22,9 @@ use futures::TryStreamExt; use pretty_assertions::assert_eq; use serde::Deserialize; use serde::de::DeserializeOwned; +use std::time::Duration; use tempfile::TempDir; +use tokio::time::timeout; use uuid::Uuid; use crate::common::exec_server::ExecServerHarness; @@ -126,6 +128,44 @@ async fn stream_rejects_platform_sandbox() -> Result<()> { Ok(()) } +#[cfg(unix)] +#[tokio::test] +async fn stream_rejects_fifo_without_waiting_for_a_writer() -> Result<()> { + let server = exec_server().await?; + let client = connect_client(server.websocket_url()).await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("named-pipe"); + let output = std::process::Command::new("mkfifo").arg(&path).output()?; + if !output.status.success() { + anyhow::bail!( + "mkfifo failed: stdout={} stderr={}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + } + + let result = timeout( + Duration::from_secs(1), + client.stream(FsReadFileParams { + path: PathUri::from_path(&path)?, + sandbox: None, + }), + ) + .await + .expect("opening a FIFO should not wait for a writer"); + let Err(error) = result else { + panic!("streaming a FIFO should be rejected"); + }; + assert_eq!( + error.to_string(), + format!( + "exec-server rejected request (-32600): path `{}` is not a file", + path.display() + ) + ); + Ok(()) +} + #[cfg(unix)] #[tokio::test] async fn stream_keeps_reading_the_open_file_after_path_replacement() -> Result<()> { From 8a5bf73af916f7eeb3de873d753bd92a60ebe15f Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 12:25:42 -0700 Subject: [PATCH 09/18] exec-server: expose file streams through filesystem --- codex-rs/Cargo.lock | 2 + codex-rs/exec-server/README.md | 2 +- codex-rs/exec-server/src/client.rs | 6 +- codex-rs/exec-server/src/file_read.rs | 6 +- codex-rs/exec-server/src/lib.rs | 2 + codex-rs/exec-server/src/local_file_system.rs | 75 +++++++++++++++++++ .../exec-server/src/remote_file_system.rs | 32 ++++++++ codex-rs/exec-server/tests/file_stream.rs | 34 --------- .../exec-server/tests/file_system/shared.rs | 48 ++++++++++++ codex-rs/file-system/Cargo.toml | 2 + codex-rs/file-system/src/lib.rs | 59 +++++++++++++++ 11 files changed, 227 insertions(+), 41 deletions(-) diff --git a/codex-rs/Cargo.lock b/codex-rs/Cargo.lock index 547a05774a71..3ae1695a7637 100644 --- a/codex-rs/Cargo.lock +++ b/codex-rs/Cargo.lock @@ -3000,8 +3000,10 @@ dependencies = [ name = "codex-file-system" version = "0.0.0" dependencies = [ + "bytes", "codex-protocol", "codex-utils-path-uri", + "futures", "serde", ] diff --git a/codex-rs/exec-server/README.md b/codex-rs/exec-server/README.md index eff62ee62b0c..c845b5be45e8 100644 --- a/codex-rs/exec-server/README.md +++ b/codex-rs/exec-server/README.md @@ -344,7 +344,7 @@ absolute path strings and normalize them to `file:` URIs: - `fs/readFile` - `fs/open`, `fs/readBlock`, and `fs/close` (internal transport for - `ExecServerClient::stream`) + `ExecServerClient::stream` and `ExecutorFileSystem::read_file_stream`) - `fs/writeFile` - `fs/createDirectory` - `fs/getMetadata` diff --git a/codex-rs/exec-server/src/client.rs b/codex-rs/exec-server/src/client.rs index d14a18839425..9df7dce3dd7c 100644 --- a/codex-rs/exec-server/src/client.rs +++ b/codex-rs/exec-server/src/client.rs @@ -493,15 +493,15 @@ impl ExecServerClient { .fs_read_block(FsReadBlockParams { handle_id: registration.handle_id.clone(), offset, - len: crate::file_read::FILE_READ_BLOCK_SIZE, + len: crate::FILE_READ_CHUNK_SIZE, }) .await?; let chunk = Bytes::from(response.chunk.into_inner()); - if chunk.len() > crate::file_read::FILE_READ_BLOCK_SIZE { + if chunk.len() > crate::FILE_READ_CHUNK_SIZE { return Err(ExecServerError::Protocol(format!( "{FS_READ_BLOCK_METHOD} returned {} bytes, maximum is {}", chunk.len(), - crate::file_read::FILE_READ_BLOCK_SIZE + crate::FILE_READ_CHUNK_SIZE ))); } if response.eof { diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs index 26a4afb84ea0..53ea46278f68 100644 --- a/codex-rs/exec-server/src/file_read.rs +++ b/codex-rs/exec-server/src/file_read.rs @@ -3,9 +3,9 @@ use std::fs::File; use std::io; use std::sync::Arc; +use codex_file_system::FILE_READ_CHUNK_SIZE; use tokio::sync::Mutex; -pub(crate) const FILE_READ_BLOCK_SIZE: usize = 1024 * 1024; const MAX_OPEN_FILE_READS: usize = 128; #[derive(Debug, Eq, PartialEq)] @@ -111,10 +111,10 @@ fn read_file_at(file: &File, bytes: &mut [u8], offset: u64) -> io::Result } fn validate_read_block_len(len: usize) -> io::Result<()> { - if !(1..=FILE_READ_BLOCK_SIZE).contains(&len) { + if !(1..=FILE_READ_CHUNK_SIZE).contains(&len) { return Err(io::Error::new( io::ErrorKind::InvalidInput, - format!("file read block length must be between 1 and {FILE_READ_BLOCK_SIZE}"), + format!("file read block length must be between 1 and {FILE_READ_CHUNK_SIZE}"), )); } Ok(()) diff --git a/codex-rs/exec-server/src/lib.rs b/codex-rs/exec-server/src/lib.rs index bb3364dd9a93..95b7767876b9 100644 --- a/codex-rs/exec-server/src/lib.rs +++ b/codex-rs/exec-server/src/lib.rs @@ -36,7 +36,9 @@ pub use codex_file_system::CopyOptions; pub use codex_file_system::CreateDirectoryOptions; pub use codex_file_system::ExecutorFileSystem; pub use codex_file_system::ExecutorFileSystemFuture; +pub use codex_file_system::FILE_READ_CHUNK_SIZE; pub use codex_file_system::FileMetadata; +pub use codex_file_system::FileSystemReadStream; pub use codex_file_system::FileSystemResult; pub use codex_file_system::FileSystemSandboxContext; pub use codex_file_system::ReadDirectoryEntry; diff --git a/codex-rs/exec-server/src/local_file_system.rs b/codex-rs/exec-server/src/local_file_system.rs index ce1d3761db3c..753a531fec88 100644 --- a/codex-rs/exec-server/src/local_file_system.rs +++ b/codex-rs/exec-server/src/local_file_system.rs @@ -1,3 +1,4 @@ +use bytes::Bytes; use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; use std::path::Path; @@ -7,13 +8,16 @@ use std::sync::LazyLock; use std::time::SystemTime; use std::time::UNIX_EPOCH; use tokio::io; +use tokio::io::AsyncReadExt; use crate::CopyOptions; use crate::CreateDirectoryOptions; use crate::ExecServerRuntimePaths; use crate::ExecutorFileSystem; use crate::ExecutorFileSystemFuture; +use crate::FILE_READ_CHUNK_SIZE; use crate::FileMetadata; +use crate::FileSystemReadStream; use crate::FileSystemResult; use crate::FileSystemSandboxContext; use crate::ReadDirectoryEntry; @@ -111,6 +115,15 @@ impl LocalFileSystem { file_system.read_file(path, sandbox).await } + async fn read_file_stream( + &self, + path: &PathUri, + sandbox: Option<&FileSystemSandboxContext>, + ) -> FileSystemResult { + let (file_system, sandbox) = self.file_system_for(sandbox)?; + file_system.read_file_stream(path, sandbox).await + } + async fn write_file( &self, path: &PathUri, @@ -190,6 +203,14 @@ impl ExecutorFileSystem for LocalFileSystem { Box::pin(LocalFileSystem::read_file(self, path, sandbox)) } + fn read_file_stream<'a>( + &'a self, + path: &'a PathUri, + sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(LocalFileSystem::read_file_stream(self, path, sandbox)) + } + fn write_file<'a>( &'a self, path: &'a PathUri, @@ -282,6 +303,17 @@ impl UnsandboxedFileSystem { self.file_system.read_file(path, /*sandbox*/ None).await } + async fn read_file_stream( + &self, + path: &PathUri, + sandbox: Option<&FileSystemSandboxContext>, + ) -> FileSystemResult { + reject_platform_sandbox_context(sandbox)?; + self.file_system + .read_file_stream(path, /*sandbox*/ None) + .await + } + async fn write_file( &self, path: &PathUri, @@ -374,6 +406,14 @@ impl ExecutorFileSystem for UnsandboxedFileSystem { Box::pin(UnsandboxedFileSystem::read_file(self, path, sandbox)) } + fn read_file_stream<'a>( + &'a self, + path: &'a PathUri, + sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(UnsandboxedFileSystem::read_file_stream(self, path, sandbox)) + } + fn write_file<'a>( &'a self, path: &'a PathUri, @@ -484,6 +524,33 @@ impl DirectFileSystem { tokio::fs::read(path.as_path()).await } + async fn read_file_stream( + &self, + path: &PathUri, + sandbox: Option<&FileSystemSandboxContext>, + ) -> FileSystemResult { + let file = self.open_file_for_read(path, sandbox).await?; + Ok(FileSystemReadStream::new(futures::stream::try_unfold( + file, + |mut file| async move { + let mut bytes = vec![0; FILE_READ_CHUNK_SIZE]; + let mut bytes_read = 0; + while bytes_read < bytes.len() { + let read = file.read(&mut bytes[bytes_read..]).await?; + if read == 0 { + break; + } + bytes_read += read; + } + if bytes_read == 0 { + return Ok(None); + } + bytes.truncate(bytes_read); + Ok(Some((Bytes::from(bytes), file))) + }, + ))) + } + async fn write_file( &self, path: &PathUri, @@ -650,6 +717,14 @@ impl ExecutorFileSystem for DirectFileSystem { Box::pin(DirectFileSystem::read_file(self, path, sandbox)) } + fn read_file_stream<'a>( + &'a self, + path: &'a PathUri, + sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(DirectFileSystem::read_file_stream(self, path, sandbox)) + } + fn write_file<'a>( &'a self, path: &'a PathUri, diff --git a/codex-rs/exec-server/src/remote_file_system.rs b/codex-rs/exec-server/src/remote_file_system.rs index 7c04f267eea2..9c335d37f21f 100644 --- a/codex-rs/exec-server/src/remote_file_system.rs +++ b/codex-rs/exec-server/src/remote_file_system.rs @@ -1,6 +1,7 @@ use base64::Engine as _; use base64::engine::general_purpose::STANDARD; use codex_utils_path_uri::PathUri; +use futures::TryStreamExt; use tokio::io; use tracing::trace; @@ -10,6 +11,7 @@ use crate::ExecServerError; use crate::ExecutorFileSystem; use crate::ExecutorFileSystemFuture; use crate::FileMetadata; +use crate::FileSystemReadStream; use crate::FileSystemResult; use crate::FileSystemSandboxContext; use crate::ReadDirectoryEntry; @@ -76,6 +78,28 @@ impl RemoteFileSystem { }) } + async fn read_file_stream( + &self, + path: &PathUri, + sandbox: Option<&FileSystemSandboxContext>, + ) -> FileSystemResult { + if sandbox.is_some_and(FileSystemSandboxContext::should_run_in_sandbox) { + return Ok(FileSystemReadStream::from_bytes( + self.read_file(path, sandbox).await?, + )); + } + trace!("remote fs read_file_stream"); + let client = self.client.get().await.map_err(map_remote_error)?; + let stream = client + .stream(FsReadFileParams { + path: path.clone(), + sandbox: remote_sandbox_context(sandbox), + }) + .await + .map_err(map_remote_error)?; + Ok(FileSystemReadStream::new(stream.map_err(map_remote_error))) + } + async fn write_file( &self, path: &PathUri, @@ -222,6 +246,14 @@ impl ExecutorFileSystem for RemoteFileSystem { Box::pin(RemoteFileSystem::read_file(self, path, sandbox)) } + fn read_file_stream<'a>( + &'a self, + path: &'a PathUri, + sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(RemoteFileSystem::read_file_stream(self, path, sandbox)) + } + fn write_file<'a>( &'a self, path: &'a PathUri, diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index 50cfba99a6f1..7753c4a4c17f 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -45,40 +45,6 @@ struct ReadBlockResponse { eof: bool, } -#[tokio::test] -async fn stream_reads_file_in_one_mib_blocks() -> Result<()> { - let server = exec_server().await?; - let client = connect_client(server.websocket_url()).await?; - let tmp = TempDir::new()?; - let path = tmp.path().join("blocks.bin"); - let contents = (0..BLOCK_SIZE * 2 + 17) - .map(|index| (index % 251) as u8) - .collect::>(); - std::fs::write(&path, &contents)?; - - let chunks = client - .stream(FsReadFileParams { - path: PathUri::from_path(path)?, - sandbox: None, - }) - .await? - .try_collect::>() - .await?; - - assert_eq!( - chunks.iter().map(bytes::Bytes::len).collect::>(), - vec![BLOCK_SIZE, BLOCK_SIZE, 17] - ); - assert_eq!( - chunks - .iter() - .flat_map(|chunk| chunk.iter().copied()) - .collect::>(), - contents - ); - Ok(()) -} - #[tokio::test] async fn stream_stops_after_an_exact_block_boundary() -> Result<()> { let server = exec_server().await?; diff --git a/codex-rs/exec-server/tests/file_system/shared.rs b/codex-rs/exec-server/tests/file_system/shared.rs index 27bdb22d85dc..08a04502c8ec 100644 --- a/codex-rs/exec-server/tests/file_system/shared.rs +++ b/codex-rs/exec-server/tests/file_system/shared.rs @@ -2,6 +2,7 @@ use anyhow::Context; use anyhow::Result; use codex_exec_server::CopyOptions; use codex_exec_server::CreateDirectoryOptions; +use codex_exec_server::FILE_READ_CHUNK_SIZE; use codex_exec_server::FileMetadata; use codex_exec_server::ReadDirectoryEntry; use codex_exec_server::RemoveOptions; @@ -11,6 +12,7 @@ use codex_protocol::models::PermissionProfile; use codex_sandboxing::policy_transforms::effective_file_system_sandbox_policy; use codex_sandboxing::policy_transforms::effective_network_sandbox_policy; use codex_utils_path_uri::PathUri; +use futures::TryStreamExt; use pretty_assertions::assert_eq; use std::path::Path; use tempfile::TempDir; @@ -193,6 +195,44 @@ async fn file_system_read_file_returns_bytes( Ok(()) } +#[test_case(FileSystemImplementation::Local ; "local")] +#[test_case(FileSystemImplementation::Remote ; "remote")] +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn file_system_read_file_stream_returns_one_mib_chunks( + implementation: FileSystemImplementation, +) -> Result<()> { + let context = create_file_system_context(implementation).await?; + let file_system = context.file_system; + + let tmp = TempDir::new()?; + let file_path = tmp.path().join("blocks.bin"); + let contents = (0..FILE_READ_CHUNK_SIZE * 2 + 17) + .map(|index| (index % 251) as u8) + .collect::>(); + std::fs::write(&file_path, &contents)?; + + let chunks = file_system + .read_file_stream(&PathUri::from_path(file_path)?, /*sandbox*/ None) + .await + .with_context(|| format!("mode={implementation}"))? + .try_collect::>() + .await?; + + assert_eq!( + chunks.iter().map(bytes::Bytes::len).collect::>(), + vec![FILE_READ_CHUNK_SIZE, FILE_READ_CHUNK_SIZE, 17] + ); + assert_eq!( + chunks + .iter() + .flat_map(|chunk| chunk.iter().copied()) + .collect::>(), + contents + ); + + Ok(()) +} + #[test_case(FileSystemImplementation::Local ; "local")] #[test_case(FileSystemImplementation::Remote ; "remote")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] @@ -447,6 +487,14 @@ async fn file_system_sandboxed_metadata_and_read_allow_readable_root( .with_context(|| format!("mode={implementation}"))?; assert_eq!(contents, b"sandboxed hello"); + let chunks = file_system + .read_file_stream(&PathUri::from_path(file_path)?, Some(&sandbox)) + .await + .with_context(|| format!("mode={implementation}"))? + .try_collect::>() + .await?; + assert_eq!(chunks, vec![bytes::Bytes::from_static(b"sandboxed hello")]); + Ok(()) } diff --git a/codex-rs/file-system/Cargo.toml b/codex-rs/file-system/Cargo.toml index d1beefbc8b08..4248e174a088 100644 --- a/codex-rs/file-system/Cargo.toml +++ b/codex-rs/file-system/Cargo.toml @@ -8,8 +8,10 @@ license.workspace = true workspace = true [dependencies] +bytes = { workspace = true } codex-protocol = { workspace = true } codex-utils-path-uri = { workspace = true } +futures = { workspace = true } serde = { workspace = true, features = ["derive"] } [lib] diff --git a/codex-rs/file-system/src/lib.rs b/codex-rs/file-system/src/lib.rs index 429c2ead50a1..5b8d22bcc5ec 100644 --- a/codex-rs/file-system/src/lib.rs +++ b/codex-rs/file-system/src/lib.rs @@ -1,3 +1,4 @@ +use bytes::Bytes; use codex_protocol::config_types::WindowsSandboxLevel; use codex_protocol::models::PermissionProfile; use codex_protocol::models::SandboxEnforcement; @@ -8,10 +9,16 @@ use codex_protocol::permissions::FileSystemSpecialPath; use codex_protocol::permissions::NetworkSandboxPolicy; use codex_protocol::protocol::SandboxPolicy; use codex_utils_path_uri::PathUri; +use futures::Stream; use std::future::Future; use std::io; use std::path::Path; use std::pin::Pin; +use std::task::Context; +use std::task::Poll; + +/// Maximum chunk size returned by [`ExecutorFileSystem::read_file_stream`]. +pub const FILE_READ_CHUNK_SIZE: usize = 1024 * 1024; #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct CreateDirectoryOptions { @@ -139,6 +146,42 @@ pub type FileSystemResult = io::Result; pub type ExecutorFileSystemFuture<'a, T> = Pin> + Send + 'a>>; +/// Stream of immutable chunks read from an [`ExecutorFileSystem`]. +pub struct FileSystemReadStream { + inner: Pin> + Send + 'static>>, +} + +impl FileSystemReadStream { + /// Wraps a filesystem byte stream. + pub fn new(stream: impl Stream> + Send + 'static) -> Self { + Self { + inner: Box::pin(stream), + } + } + + /// Splits an already-buffered file into standard-sized chunks. + pub fn from_bytes(bytes: Vec) -> Self { + Self::new(futures::stream::try_unfold( + Bytes::from(bytes), + |mut remaining| async move { + if remaining.is_empty() { + return Ok(None); + } + let chunk = remaining.split_to(remaining.len().min(FILE_READ_CHUNK_SIZE)); + Ok(Some((chunk, remaining))) + }, + )) + } +} + +impl Stream for FileSystemReadStream { + type Item = FileSystemResult; + + fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + self.inner.as_mut().poll_next(cx) + } +} + /// Abstract filesystem access used by components that may operate locally or via /// a remote environment. pub trait ExecutorFileSystem: Send + Sync { @@ -155,6 +198,22 @@ pub trait ExecutorFileSystem: Send + Sync { sandbox: Option<&'a FileSystemSandboxContext>, ) -> ExecutorFileSystemFuture<'a, Vec>; + /// Reads a file as a stream of chunks no larger than [`FILE_READ_CHUNK_SIZE`]. + /// + /// Implementations that cannot keep an open file handle may buffer the file + /// before returning the stream. + fn read_file_stream<'a>( + &'a self, + path: &'a PathUri, + sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(async move { + Ok(FileSystemReadStream::from_bytes( + self.read_file(path, sandbox).await?, + )) + }) + } + /// Reads a file and decodes it as UTF-8 text. fn read_file_text<'a>( &'a self, From 1ca386e3466c4a5c8264322eace8972ca79df443 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 12:43:59 -0700 Subject: [PATCH 10/18] filesystem: require explicit stream support --- codex-rs/config/src/loader/tests.rs | 14 +++++++++++ codex-rs/core-plugins/src/provider_tests.rs | 9 +++++++ codex-rs/core/src/agents_md_tests.rs | 14 +++++++++++ .../exec-server/src/remote_file_system.rs | 5 ++-- .../exec-server/src/sandboxed_file_system.rs | 14 +++++++++++ .../exec-server/tests/file_system/shared.rs | 8 ------ .../mcp/src/executor_plugin/provider_tests.rs | 9 +++++++ .../tests/executor_file_system_authority.rs | 14 +++++++++++ codex-rs/file-system/src/lib.rs | 25 +------------------ 9 files changed, 78 insertions(+), 34 deletions(-) diff --git a/codex-rs/config/src/loader/tests.rs b/codex-rs/config/src/loader/tests.rs index 543c307c3548..068e29aa6ec9 100644 --- a/codex-rs/config/src/loader/tests.rs +++ b/codex-rs/config/src/loader/tests.rs @@ -3,6 +3,7 @@ use codex_file_system::CopyOptions; use codex_file_system::CreateDirectoryOptions; use codex_file_system::ExecutorFileSystemFuture; use codex_file_system::FileMetadata; +use codex_file_system::FileSystemReadStream; use codex_file_system::FileSystemSandboxContext; use codex_file_system::ReadDirectoryEntry; use codex_file_system::RemoveOptions; @@ -36,6 +37,19 @@ impl ExecutorFileSystem for TestFileSystem { }) } + fn read_file_stream<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(async { + Err(std::io::Error::new( + std::io::ErrorKind::Unsupported, + "test filesystem does not support streaming reads", + )) + }) + } + fn write_file<'a>( &'a self, _path: &'a PathUri, diff --git a/codex-rs/core-plugins/src/provider_tests.rs b/codex-rs/core-plugins/src/provider_tests.rs index a6a7e2fa7463..3f92a47f0b9f 100644 --- a/codex-rs/core-plugins/src/provider_tests.rs +++ b/codex-rs/core-plugins/src/provider_tests.rs @@ -8,6 +8,7 @@ use codex_exec_server::EnvironmentManager; use codex_exec_server::ExecutorFileSystem; use codex_exec_server::ExecutorFileSystemFuture; use codex_exec_server::FileMetadata; +use codex_exec_server::FileSystemReadStream; use codex_exec_server::FileSystemResult; use codex_exec_server::FileSystemSandboxContext; use codex_exec_server::LOCAL_ENVIRONMENT_ID; @@ -89,6 +90,14 @@ impl ExecutorFileSystem for SyntheticPluginFileSystem { }) } + fn read_file_stream<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(async { Self::unsupported() }) + } + fn write_file<'a>( &'a self, _path: &'a PathUri, diff --git a/codex-rs/core/src/agents_md_tests.rs b/codex-rs/core/src/agents_md_tests.rs index 8d0987ca4dcd..aec3dae3a2e3 100644 --- a/codex-rs/core/src/agents_md_tests.rs +++ b/codex-rs/core/src/agents_md_tests.rs @@ -11,6 +11,7 @@ use codex_exec_server::CreateDirectoryOptions; use codex_exec_server::Environment; use codex_exec_server::ExecutorFileSystemFuture; use codex_exec_server::FileMetadata; +use codex_exec_server::FileSystemReadStream; use codex_exec_server::FileSystemSandboxContext; use codex_exec_server::LOCAL_FS; use codex_exec_server::ReadDirectoryEntry; @@ -141,6 +142,19 @@ impl ExecutorFileSystem for FailingFileSystem { Box::pin(FailingFileSystem::read_file(self, path, sandbox)) } + fn read_file_stream<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(async { + Err(io::Error::new( + io::ErrorKind::Unsupported, + "failing filesystem does not support streaming reads", + )) + }) + } + fn write_file<'a>( &'a self, path: &'a PathUri, diff --git a/codex-rs/exec-server/src/remote_file_system.rs b/codex-rs/exec-server/src/remote_file_system.rs index 9c335d37f21f..49fbd2475fcb 100644 --- a/codex-rs/exec-server/src/remote_file_system.rs +++ b/codex-rs/exec-server/src/remote_file_system.rs @@ -84,8 +84,9 @@ impl RemoteFileSystem { sandbox: Option<&FileSystemSandboxContext>, ) -> FileSystemResult { if sandbox.is_some_and(FileSystemSandboxContext::should_run_in_sandbox) { - return Ok(FileSystemReadStream::from_bytes( - self.read_file(path, sandbox).await?, + return Err(io::Error::new( + io::ErrorKind::Unsupported, + "streaming file reads do not support platform sandboxing", )); } trace!("remote fs read_file_stream"); diff --git a/codex-rs/exec-server/src/sandboxed_file_system.rs b/codex-rs/exec-server/src/sandboxed_file_system.rs index f1ed02f76c9b..d0beb57a43a8 100644 --- a/codex-rs/exec-server/src/sandboxed_file_system.rs +++ b/codex-rs/exec-server/src/sandboxed_file_system.rs @@ -10,6 +10,7 @@ use crate::ExecServerRuntimePaths; use crate::ExecutorFileSystem; use crate::ExecutorFileSystemFuture; use crate::FileMetadata; +use crate::FileSystemReadStream; use crate::FileSystemResult; use crate::FileSystemSandboxContext; use crate::ReadDirectoryEntry; @@ -265,6 +266,19 @@ impl ExecutorFileSystem for SandboxedFileSystem { Box::pin(SandboxedFileSystem::read_file(self, path, sandbox)) } + fn read_file_stream<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(async { + Err(io::Error::new( + io::ErrorKind::Unsupported, + "streaming file reads do not support platform sandboxing", + )) + }) + } + fn write_file<'a>( &'a self, path: &'a PathUri, diff --git a/codex-rs/exec-server/tests/file_system/shared.rs b/codex-rs/exec-server/tests/file_system/shared.rs index 08a04502c8ec..4b48f1eed243 100644 --- a/codex-rs/exec-server/tests/file_system/shared.rs +++ b/codex-rs/exec-server/tests/file_system/shared.rs @@ -487,14 +487,6 @@ async fn file_system_sandboxed_metadata_and_read_allow_readable_root( .with_context(|| format!("mode={implementation}"))?; assert_eq!(contents, b"sandboxed hello"); - let chunks = file_system - .read_file_stream(&PathUri::from_path(file_path)?, Some(&sandbox)) - .await - .with_context(|| format!("mode={implementation}"))? - .try_collect::>() - .await?; - assert_eq!(chunks, vec![bytes::Bytes::from_static(b"sandboxed hello")]); - Ok(()) } diff --git a/codex-rs/ext/mcp/src/executor_plugin/provider_tests.rs b/codex-rs/ext/mcp/src/executor_plugin/provider_tests.rs index dcaac40b698b..3bc9ef579624 100644 --- a/codex-rs/ext/mcp/src/executor_plugin/provider_tests.rs +++ b/codex-rs/ext/mcp/src/executor_plugin/provider_tests.rs @@ -8,6 +8,7 @@ use codex_exec_server::CreateDirectoryOptions; use codex_exec_server::ExecutorFileSystem; use codex_exec_server::ExecutorFileSystemFuture; use codex_exec_server::FileMetadata; +use codex_exec_server::FileSystemReadStream; use codex_exec_server::FileSystemResult; use codex_exec_server::FileSystemSandboxContext; use codex_exec_server::ReadDirectoryEntry; @@ -73,6 +74,14 @@ impl ExecutorFileSystem for SyntheticExecutorFileSystem { }) } + fn read_file_stream<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(async { Self::unsupported() }) + } + fn write_file<'a>( &'a self, _path: &'a PathUri, diff --git a/codex-rs/ext/skills/tests/executor_file_system_authority.rs b/codex-rs/ext/skills/tests/executor_file_system_authority.rs index 02443351fd90..f3d17e35255b 100644 --- a/codex-rs/ext/skills/tests/executor_file_system_authority.rs +++ b/codex-rs/ext/skills/tests/executor_file_system_authority.rs @@ -12,6 +12,7 @@ use codex_exec_server::EnvironmentManager; use codex_exec_server::ExecutorFileSystem; use codex_exec_server::ExecutorFileSystemFuture; use codex_exec_server::FileMetadata; +use codex_exec_server::FileSystemReadStream; use codex_exec_server::FileSystemSandboxContext; use codex_exec_server::ReadDirectoryEntry; use codex_exec_server::RemoveOptions; @@ -109,6 +110,19 @@ impl ExecutorFileSystem for SyntheticFileSystem { Box::pin(SyntheticFileSystem::read_file(self, path)) } + fn read_file_stream<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(async { + Err(io::Error::new( + io::ErrorKind::Unsupported, + "synthetic filesystem does not support streaming reads", + )) + }) + } + fn write_file<'a>( &'a self, _path: &'a PathUri, diff --git a/codex-rs/file-system/src/lib.rs b/codex-rs/file-system/src/lib.rs index 5b8d22bcc5ec..2472efa94a6d 100644 --- a/codex-rs/file-system/src/lib.rs +++ b/codex-rs/file-system/src/lib.rs @@ -158,20 +158,6 @@ impl FileSystemReadStream { inner: Box::pin(stream), } } - - /// Splits an already-buffered file into standard-sized chunks. - pub fn from_bytes(bytes: Vec) -> Self { - Self::new(futures::stream::try_unfold( - Bytes::from(bytes), - |mut remaining| async move { - if remaining.is_empty() { - return Ok(None); - } - let chunk = remaining.split_to(remaining.len().min(FILE_READ_CHUNK_SIZE)); - Ok(Some((chunk, remaining))) - }, - )) - } } impl Stream for FileSystemReadStream { @@ -199,20 +185,11 @@ pub trait ExecutorFileSystem: Send + Sync { ) -> ExecutorFileSystemFuture<'a, Vec>; /// Reads a file as a stream of chunks no larger than [`FILE_READ_CHUNK_SIZE`]. - /// - /// Implementations that cannot keep an open file handle may buffer the file - /// before returning the stream. fn read_file_stream<'a>( &'a self, path: &'a PathUri, sandbox: Option<&'a FileSystemSandboxContext>, - ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { - Box::pin(async move { - Ok(FileSystemReadStream::from_bytes( - self.read_file(path, sandbox).await?, - )) - }) - } + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream>; /// Reads a file and decodes it as UTF-8 text. fn read_file_text<'a>( From 79533b0508248cad55e81a0132e06d2d52285655 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 12:45:39 -0700 Subject: [PATCH 11/18] exec-server: keep file handles open after eof --- codex-rs/exec-server/src/client.rs | 11 ++++++- codex-rs/exec-server/src/file_read.rs | 2 +- codex-rs/exec-server/tests/file_stream.rs | 40 +++++++++++++++++++++++ 3 files changed, 51 insertions(+), 2 deletions(-) diff --git a/codex-rs/exec-server/src/client.rs b/codex-rs/exec-server/src/client.rs index 9df7dce3dd7c..9211679939a2 100644 --- a/codex-rs/exec-server/src/client.rs +++ b/codex-rs/exec-server/src/client.rs @@ -505,7 +505,16 @@ impl ExecServerClient { ))); } if response.eof { - registration.active = false; + if registration + .client + .fs_close(FsCloseParams { + handle_id: registration.handle_id.clone(), + }) + .await + .is_ok() + { + registration.active = false; + } return if chunk.is_empty() { Ok(None) } else { diff --git a/codex-rs/exec-server/src/file_read.rs b/codex-rs/exec-server/src/file_read.rs index 53ea46278f68..8cc84b67300d 100644 --- a/codex-rs/exec-server/src/file_read.rs +++ b/codex-rs/exec-server/src/file_read.rs @@ -64,7 +64,7 @@ impl FileReadHandleManager { "file read task stopped unexpectedly: {error}" ))), }; - if result.is_err() || result.as_ref().is_ok_and(|block| block.eof) { + if result.is_err() { self.close(handle_id).await; } result diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index 7753c4a4c17f..2e068e8f5d6e 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -69,6 +69,30 @@ async fn stream_stops_after_an_exact_block_boundary() -> Result<()> { Ok(()) } +#[tokio::test] +async fn completed_streams_release_handle_capacity() -> Result<()> { + let server = exec_server().await?; + let client = connect_client(server.websocket_url()).await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("repeated.txt"); + std::fs::write(&path, b"repeated")?; + let path = PathUri::from_path(path)?; + + for _ in 0..=OPEN_FILE_LIMIT { + let chunks = client + .stream(FsReadFileParams { + path: path.clone(), + sandbox: None, + }) + .await? + .try_collect::>() + .await?; + assert_eq!(chunks, vec![bytes::Bytes::from_static(b"repeated")]); + } + + Ok(()) +} + #[tokio::test] async fn stream_rejects_platform_sandbox() -> Result<()> { let server = exec_server().await?; @@ -212,6 +236,22 @@ async fn read_block_supports_non_sequential_offsets_and_lengths() -> Result<()> .await?, (b"89".to_vec(), true) ); + assert_eq!( + read_block( + &mut server, + &open.handle_id, + /*offset*/ 0, + /*len*/ 2 + ) + .await?, + (b"01".to_vec(), false) + ); + let _: serde_json::Value = rpc_call( + &mut server, + "fs/close", + serde_json::json!({ "handleId": open.handle_id }), + ) + .await?; server.shutdown().await?; Ok(()) } From 25f4daf9a05134818ecec022482a9964b749d447 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 12:56:14 -0700 Subject: [PATCH 12/18] exec-server: expose file streams only through filesystem --- codex-rs/exec-server/README.md | 3 +- codex-rs/exec-server/src/client.rs | 128 ++---------------- codex-rs/exec-server/src/lib.rs | 1 - .../exec-server/src/remote_file_stream.rs | 121 +++++++++++++++++ .../exec-server/src/remote_file_system.rs | 13 +- codex-rs/exec-server/tests/file_stream.rs | 71 ++++------ 6 files changed, 160 insertions(+), 177 deletions(-) create mode 100644 codex-rs/exec-server/src/remote_file_stream.rs diff --git a/codex-rs/exec-server/README.md b/codex-rs/exec-server/README.md index c845b5be45e8..1d251cab9eec 100644 --- a/codex-rs/exec-server/README.md +++ b/codex-rs/exec-server/README.md @@ -344,7 +344,7 @@ absolute path strings and normalize them to `file:` URIs: - `fs/readFile` - `fs/open`, `fs/readBlock`, and `fs/close` (internal transport for - `ExecServerClient::stream` and `ExecutorFileSystem::read_file_stream`) + `ExecutorFileSystem::read_file_stream`) - `fs/writeFile` - `fs/createDirectory` - `fs/getMetadata` @@ -383,7 +383,6 @@ The crate exports: - `ExecServerClient` - `ExecServerError` -- `FileReadStream` - `ExecServerClientConnectOptions` - `RemoteExecServerConnectArgs` - protocol request/response structs for process and filesystem RPCs diff --git a/codex-rs/exec-server/src/client.rs b/codex-rs/exec-server/src/client.rs index 9211679939a2..e6e1f6884f31 100644 --- a/codex-rs/exec-server/src/client.rs +++ b/codex-rs/exec-server/src/client.rs @@ -1,22 +1,15 @@ use std::collections::BTreeMap; use std::collections::HashMap; -use std::pin::Pin; use std::sync::Arc; use std::sync::Mutex as StdMutex; use std::sync::OnceLock; use std::sync::atomic::AtomicU64; -use std::task::Context; -use std::task::Poll; use std::time::Duration; use arc_swap::ArcSwap; -use bytes::Bytes; use codex_app_server_protocol::JSONRPCNotification; use futures::FutureExt; -use futures::Stream; -use futures::StreamExt; use futures::future::BoxFuture; -use futures::stream::BoxStream; use serde_json::Value; use tokio::sync::Mutex; use tokio::sync::Semaphore; @@ -25,7 +18,6 @@ use tokio::sync::watch; use tokio::time::timeout; use tracing::debug; -use uuid::Uuid; use crate::ProcessId; use crate::client_api::ExecServerClientConnectOptions; @@ -112,26 +104,6 @@ const INITIALIZE_TIMEOUT: Duration = Duration::from_secs(10); const PROCESS_EVENT_CHANNEL_CAPACITY: usize = 256; const PROCESS_EVENT_RETAINED_BYTES: usize = 1024 * 1024; -/// Stream of immutable file blocks returned by [`ExecServerClient::stream`]. -pub struct FileReadStream { - inner: BoxStream<'static, Result>, -} - -impl Stream for FileReadStream { - type Item = Result; - - fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - self.inner.as_mut().poll_next(cx) - } -} - -struct FileReadRegistration { - client: ExecServerClient, - handle_id: String, - runtime: Option, - active: bool, -} - impl Default for ExecServerClientConnectOptions { fn default() -> Self { Self { @@ -466,89 +438,24 @@ impl ExecServerClient { self.call(FS_READ_FILE_METHOD, ¶ms).await } - /// Opens an unsandboxed file and returns a demand-driven stream of 1 MiB blocks. - pub async fn stream( + pub(crate) async fn fs_open( &self, - params: FsReadFileParams, - ) -> Result { - let registration = FileReadRegistration { - client: self.clone(), - handle_id: Uuid::new_v4().to_string(), - runtime: tokio::runtime::Handle::try_current().ok(), - active: true, - }; - self.fs_open(FsOpenParams { - handle_id: registration.handle_id.clone(), - path: params.path, - sandbox: params.sandbox, - }) - .await?; - Ok(FileReadStream { - inner: futures::stream::try_unfold(Some((registration, 0_u64)), |state| async move { - let Some((mut registration, offset)) = state else { - return Ok(None); - }; - let response = registration - .client - .fs_read_block(FsReadBlockParams { - handle_id: registration.handle_id.clone(), - offset, - len: crate::FILE_READ_CHUNK_SIZE, - }) - .await?; - let chunk = Bytes::from(response.chunk.into_inner()); - if chunk.len() > crate::FILE_READ_CHUNK_SIZE { - return Err(ExecServerError::Protocol(format!( - "{FS_READ_BLOCK_METHOD} returned {} bytes, maximum is {}", - chunk.len(), - crate::FILE_READ_CHUNK_SIZE - ))); - } - if response.eof { - if registration - .client - .fs_close(FsCloseParams { - handle_id: registration.handle_id.clone(), - }) - .await - .is_ok() - { - registration.active = false; - } - return if chunk.is_empty() { - Ok(None) - } else { - Ok(Some((chunk, None))) - }; - } - if chunk.is_empty() { - return Err(ExecServerError::Protocol(format!( - "{FS_READ_BLOCK_METHOD} returned an empty non-terminal block" - ))); - } - let next_offset = offset.checked_add(chunk.len() as u64).ok_or_else(|| { - ExecServerError::Protocol(format!( - "{FS_READ_BLOCK_METHOD} offset overflowed after {offset} bytes" - )) - })?; - Ok(Some((chunk, Some((registration, next_offset))))) - }) - .boxed(), - }) - } - - async fn fs_open(&self, params: FsOpenParams) -> Result { + params: FsOpenParams, + ) -> Result { self.call(FS_OPEN_METHOD, ¶ms).await } - async fn fs_read_block( + pub(crate) async fn fs_read_block( &self, params: FsReadBlockParams, ) -> Result { self.call(FS_READ_BLOCK_METHOD, ¶ms).await } - async fn fs_close(&self, params: FsCloseParams) -> Result { + pub(crate) async fn fs_close( + &self, + params: FsCloseParams, + ) -> Result { self.call(FS_CLOSE_METHOD, ¶ms).await } @@ -724,25 +631,6 @@ impl ExecServerClient { } } -impl Drop for FileReadRegistration { - fn drop(&mut self) { - if !self.active { - return; - } - let client = self.client.clone(); - let handle_id = self.handle_id.clone(); - let runtime = self - .runtime - .clone() - .or_else(|| tokio::runtime::Handle::try_current().ok()); - if let Some(runtime) = runtime { - runtime.spawn(async move { - let _ = client.fs_close(FsCloseParams { handle_id }).await; - }); - } - } -} - impl From for ExecServerError { fn from(value: RpcCallError) -> Self { match value { diff --git a/codex-rs/exec-server/src/lib.rs b/codex-rs/exec-server/src/lib.rs index 95b7767876b9..4b4d8c99030e 100644 --- a/codex-rs/exec-server/src/lib.rs +++ b/codex-rs/exec-server/src/lib.rs @@ -26,7 +26,6 @@ mod server; pub use client::ExecServerClient; pub use client::ExecServerError; -pub use client::FileReadStream; pub use client::http_client::HttpResponseBodyStream; pub use client::http_client::ReqwestHttpClient; pub use client_api::ExecServerClientConnectOptions; diff --git a/codex-rs/exec-server/src/remote_file_stream.rs b/codex-rs/exec-server/src/remote_file_stream.rs new file mode 100644 index 000000000000..2a22dcf2ae66 --- /dev/null +++ b/codex-rs/exec-server/src/remote_file_stream.rs @@ -0,0 +1,121 @@ +use bytes::Bytes; +use codex_utils_path_uri::PathUri; +use tokio::io; +use uuid::Uuid; + +use super::map_remote_error; +use crate::ExecServerClient; +use crate::FILE_READ_CHUNK_SIZE; +use crate::FileSystemReadStream; +use crate::FileSystemResult; +use crate::FileSystemSandboxContext; +use crate::protocol::FS_READ_BLOCK_METHOD; +use crate::protocol::FsCloseParams; +use crate::protocol::FsOpenParams; +use crate::protocol::FsReadBlockParams; + +struct FileReadRegistration { + client: ExecServerClient, + handle_id: String, + runtime: Option, + active: bool, +} + +pub(super) async fn open( + client: ExecServerClient, + path: PathUri, + sandbox: Option, +) -> FileSystemResult { + let registration = FileReadRegistration { + client, + handle_id: Uuid::new_v4().to_string(), + runtime: tokio::runtime::Handle::try_current().ok(), + active: true, + }; + registration + .client + .fs_open(FsOpenParams { + handle_id: registration.handle_id.clone(), + path, + sandbox, + }) + .await + .map_err(map_remote_error)?; + Ok(FileSystemReadStream::new(futures::stream::try_unfold( + Some((registration, 0_u64)), + |state| async move { + let Some((mut registration, offset)) = state else { + return Ok(None); + }; + let response = registration + .client + .fs_read_block(FsReadBlockParams { + handle_id: registration.handle_id.clone(), + offset, + len: FILE_READ_CHUNK_SIZE, + }) + .await + .map_err(map_remote_error)?; + let chunk = Bytes::from(response.chunk.into_inner()); + if chunk.len() > FILE_READ_CHUNK_SIZE { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!( + "{FS_READ_BLOCK_METHOD} returned {} bytes, maximum is {}", + chunk.len(), + FILE_READ_CHUNK_SIZE + ), + )); + } + if response.eof { + if registration + .client + .fs_close(FsCloseParams { + handle_id: registration.handle_id.clone(), + }) + .await + .is_ok() + { + registration.active = false; + } + return if chunk.is_empty() { + Ok(None) + } else { + Ok(Some((chunk, None))) + }; + } + if chunk.is_empty() { + return Err(io::Error::new( + io::ErrorKind::InvalidData, + format!("{FS_READ_BLOCK_METHOD} returned an empty non-terminal block"), + )); + } + let next_offset = offset.checked_add(chunk.len() as u64).ok_or_else(|| { + io::Error::new( + io::ErrorKind::InvalidData, + format!("{FS_READ_BLOCK_METHOD} offset overflowed after {offset} bytes"), + ) + })?; + Ok(Some((chunk, Some((registration, next_offset))))) + }, + ))) +} + +impl Drop for FileReadRegistration { + fn drop(&mut self) { + if !self.active { + return; + } + let client = self.client.clone(); + let handle_id = self.handle_id.clone(); + let runtime = self + .runtime + .clone() + .or_else(|| tokio::runtime::Handle::try_current().ok()); + if let Some(runtime) = runtime { + runtime.spawn(async move { + let _ = client.fs_close(FsCloseParams { handle_id }).await; + }); + } + } +} diff --git a/codex-rs/exec-server/src/remote_file_system.rs b/codex-rs/exec-server/src/remote_file_system.rs index 49fbd2475fcb..f9136e146108 100644 --- a/codex-rs/exec-server/src/remote_file_system.rs +++ b/codex-rs/exec-server/src/remote_file_system.rs @@ -1,7 +1,6 @@ use base64::Engine as _; use base64::engine::general_purpose::STANDARD; use codex_utils_path_uri::PathUri; -use futures::TryStreamExt; use tokio::io; use tracing::trace; @@ -29,6 +28,9 @@ use crate::protocol::FsWriteFileParams; const INVALID_REQUEST_ERROR_CODE: i64 = -32600; const NOT_FOUND_ERROR_CODE: i64 = -32004; +#[path = "remote_file_stream.rs"] +mod file_stream; + pub(crate) struct RemoteFileSystem { client: LazyRemoteExecServerClient, } @@ -91,14 +93,7 @@ impl RemoteFileSystem { } trace!("remote fs read_file_stream"); let client = self.client.get().await.map_err(map_remote_error)?; - let stream = client - .stream(FsReadFileParams { - path: path.clone(), - sandbox: remote_sandbox_context(sandbox), - }) - .await - .map_err(map_remote_error)?; - Ok(FileSystemReadStream::new(stream.map_err(map_remote_error))) + file_stream::open(client, path.clone(), remote_sandbox_context(sandbox)).await } async fn write_file( diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index 2e068e8f5d6e..fd0b7f8d9e3f 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -5,11 +5,10 @@ use base64::Engine as _; use codex_app_server_protocol::JSONRPCError; use codex_app_server_protocol::JSONRPCMessage; use codex_app_server_protocol::JSONRPCResponse; -use codex_exec_server::ExecServerClient; +use codex_exec_server::Environment; +use codex_exec_server::ExecutorFileSystem; use codex_exec_server::FileSystemSandboxContext; -use codex_exec_server::FsReadFileParams; use codex_exec_server::InitializeParams; -use codex_exec_server::RemoteExecServerConnectArgs; use codex_protocol::models::PermissionProfile; use codex_protocol::permissions::FileSystemAccessMode; use codex_protocol::permissions::FileSystemPath; @@ -22,6 +21,7 @@ use futures::TryStreamExt; use pretty_assertions::assert_eq; use serde::Deserialize; use serde::de::DeserializeOwned; +use std::sync::Arc; use std::time::Duration; use tempfile::TempDir; use tokio::time::timeout; @@ -48,16 +48,13 @@ struct ReadBlockResponse { #[tokio::test] async fn stream_stops_after_an_exact_block_boundary() -> Result<()> { let server = exec_server().await?; - let client = connect_client(server.websocket_url()).await?; + let file_system = connect_file_system(server.websocket_url())?; let tmp = TempDir::new()?; let path = tmp.path().join("exact-blocks.bin"); std::fs::write(&path, vec![b'x'; BLOCK_SIZE * 2])?; - let chunks = client - .stream(FsReadFileParams { - path: PathUri::from_path(path)?, - sandbox: None, - }) + let chunks = file_system + .read_file_stream(&PathUri::from_path(path)?, /*sandbox*/ None) .await? .try_collect::>() .await?; @@ -72,18 +69,15 @@ async fn stream_stops_after_an_exact_block_boundary() -> Result<()> { #[tokio::test] async fn completed_streams_release_handle_capacity() -> Result<()> { let server = exec_server().await?; - let client = connect_client(server.websocket_url()).await?; + let file_system = connect_file_system(server.websocket_url())?; let tmp = TempDir::new()?; let path = tmp.path().join("repeated.txt"); std::fs::write(&path, b"repeated")?; let path = PathUri::from_path(path)?; for _ in 0..=OPEN_FILE_LIMIT { - let chunks = client - .stream(FsReadFileParams { - path: path.clone(), - sandbox: None, - }) + let chunks = file_system + .read_file_stream(&path, /*sandbox*/ None) .await? .try_collect::>() .await?; @@ -96,24 +90,25 @@ async fn completed_streams_release_handle_capacity() -> Result<()> { #[tokio::test] async fn stream_rejects_platform_sandbox() -> Result<()> { let server = exec_server().await?; - let client = connect_client(server.websocket_url()).await?; + let file_system = connect_file_system(server.websocket_url())?; let tmp = TempDir::new()?; let path = tmp.path().join("sandboxed.txt"); std::fs::write(&path, "sandboxed hello")?; - let result = client - .stream(FsReadFileParams { - path: PathUri::from_path(&path)?, - sandbox: Some(read_only_sandbox(tmp.path().to_path_buf())), - }) + let result = file_system + .read_file_stream( + &PathUri::from_path(&path)?, + Some(&read_only_sandbox(tmp.path().to_path_buf())), + ) .await; let Err(error) = result else { panic!("sandboxed stream should be rejected"); }; + assert_eq!(error.kind(), std::io::ErrorKind::Unsupported); assert_eq!( error.to_string(), - "exec-server rejected request (-32600): streaming file reads do not support platform sandboxing" + "streaming file reads do not support platform sandboxing" ); Ok(()) } @@ -122,7 +117,7 @@ async fn stream_rejects_platform_sandbox() -> Result<()> { #[tokio::test] async fn stream_rejects_fifo_without_waiting_for_a_writer() -> Result<()> { let server = exec_server().await?; - let client = connect_client(server.websocket_url()).await?; + let file_system = connect_file_system(server.websocket_url())?; let tmp = TempDir::new()?; let path = tmp.path().join("named-pipe"); let output = std::process::Command::new("mkfifo").arg(&path).output()?; @@ -136,10 +131,7 @@ async fn stream_rejects_fifo_without_waiting_for_a_writer() -> Result<()> { let result = timeout( Duration::from_secs(1), - client.stream(FsReadFileParams { - path: PathUri::from_path(&path)?, - sandbox: None, - }), + file_system.read_file_stream(&PathUri::from_path(&path)?, /*sandbox*/ None), ) .await .expect("opening a FIFO should not wait for a writer"); @@ -148,10 +140,7 @@ async fn stream_rejects_fifo_without_waiting_for_a_writer() -> Result<()> { }; assert_eq!( error.to_string(), - format!( - "exec-server rejected request (-32600): path `{}` is not a file", - path.display() - ) + format!("path `{}` is not a file", path.display()) ); Ok(()) } @@ -160,15 +149,12 @@ async fn stream_rejects_fifo_without_waiting_for_a_writer() -> Result<()> { #[tokio::test] async fn stream_keeps_reading_the_open_file_after_path_replacement() -> Result<()> { let server = exec_server().await?; - let client = connect_client(server.websocket_url()).await?; + let file_system = connect_file_system(server.websocket_url())?; let tmp = TempDir::new()?; let path = tmp.path().join("replaceable.bin"); std::fs::write(&path, vec![b'a'; BLOCK_SIZE + 1])?; - let mut stream = client - .stream(FsReadFileParams { - path: PathUri::from_path(&path)?, - sandbox: None, - }) + let mut stream = file_system + .read_file_stream(&PathUri::from_path(&path)?, /*sandbox*/ None) .await?; assert_eq!( @@ -384,14 +370,9 @@ async fn rpc_message( .await } -async fn connect_client(websocket_url: &str) -> Result { - Ok( - ExecServerClient::connect_websocket(RemoteExecServerConnectArgs::new( - websocket_url.to_string(), - "file-stream-test".to_string(), - )) - .await?, - ) +fn connect_file_system(websocket_url: &str) -> Result> { + let environment = Environment::create_for_tests(Some(websocket_url.to_string()))?; + Ok(environment.get_filesystem()) } fn read_only_sandbox(path: std::path::PathBuf) -> FileSystemSandboxContext { From f4d3c13f98ba18609e52526d55ace16c821fc67a Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 13:10:36 -0700 Subject: [PATCH 13/18] codex: expose file read protocol methods (#28354) --- codex-rs/exec-server/src/client.rs | 9 +- codex-rs/exec-server/src/lib.rs | 7 ++ codex-rs/exec-server/src/protocol.rs | 12 +- codex-rs/exec-server/tests/file_stream.rs | 135 +++++++++------------- 4 files changed, 71 insertions(+), 92 deletions(-) diff --git a/codex-rs/exec-server/src/client.rs b/codex-rs/exec-server/src/client.rs index e6e1f6884f31..883e81871b5a 100644 --- a/codex-rs/exec-server/src/client.rs +++ b/codex-rs/exec-server/src/client.rs @@ -438,21 +438,18 @@ impl ExecServerClient { self.call(FS_READ_FILE_METHOD, ¶ms).await } - pub(crate) async fn fs_open( - &self, - params: FsOpenParams, - ) -> Result { + pub async fn fs_open(&self, params: FsOpenParams) -> Result { self.call(FS_OPEN_METHOD, ¶ms).await } - pub(crate) async fn fs_read_block( + pub async fn fs_read_block( &self, params: FsReadBlockParams, ) -> Result { self.call(FS_READ_BLOCK_METHOD, ¶ms).await } - pub(crate) async fn fs_close( + pub async fn fs_close( &self, params: FsCloseParams, ) -> Result { diff --git a/codex-rs/exec-server/src/lib.rs b/codex-rs/exec-server/src/lib.rs index 4b4d8c99030e..9d79b7b6c3e9 100644 --- a/codex-rs/exec-server/src/lib.rs +++ b/codex-rs/exec-server/src/lib.rs @@ -62,6 +62,7 @@ pub use process::ExecProcessEventReceiver; pub use process::ExecProcessFuture; pub use process::StartedExecProcess; pub use process_id::ProcessId; +pub use protocol::ByteChunk; pub use protocol::EnvironmentInfo; pub use protocol::ExecClosedNotification; pub use protocol::ExecEnvPolicy; @@ -72,12 +73,18 @@ pub use protocol::ExecParams; pub use protocol::ExecResponse; pub use protocol::FsCanonicalizeParams; pub use protocol::FsCanonicalizeResponse; +pub use protocol::FsCloseParams; +pub use protocol::FsCloseResponse; pub use protocol::FsCopyParams; pub use protocol::FsCopyResponse; pub use protocol::FsCreateDirectoryParams; pub use protocol::FsCreateDirectoryResponse; pub use protocol::FsGetMetadataParams; pub use protocol::FsGetMetadataResponse; +pub use protocol::FsOpenParams; +pub use protocol::FsOpenResponse; +pub use protocol::FsReadBlockParams; +pub use protocol::FsReadBlockResponse; pub use protocol::FsReadDirectoryEntry; pub use protocol::FsReadDirectoryParams; pub use protocol::FsReadDirectoryResponse; diff --git a/codex-rs/exec-server/src/protocol.rs b/codex-rs/exec-server/src/protocol.rs index 8ebd7e5eaf11..2017c7f0ff21 100644 --- a/codex-rs/exec-server/src/protocol.rs +++ b/codex-rs/exec-server/src/protocol.rs @@ -215,7 +215,7 @@ pub struct FsReadFileResponse { #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] -pub(crate) struct FsOpenParams { +pub struct FsOpenParams { pub handle_id: String, pub path: PathUri, pub sandbox: Option, @@ -223,13 +223,13 @@ pub(crate) struct FsOpenParams { #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] -pub(crate) struct FsOpenResponse { +pub struct FsOpenResponse { pub handle_id: String, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] -pub(crate) struct FsReadBlockParams { +pub struct FsReadBlockParams { pub handle_id: String, pub offset: u64, pub len: usize, @@ -237,20 +237,20 @@ pub(crate) struct FsReadBlockParams { #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] -pub(crate) struct FsReadBlockResponse { +pub struct FsReadBlockResponse { pub chunk: ByteChunk, pub eof: bool, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] -pub(crate) struct FsCloseParams { +pub struct FsCloseParams { pub handle_id: String, } #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] -pub(crate) struct FsCloseResponse {} +pub struct FsCloseResponse {} #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[serde(rename_all = "camelCase")] diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index fd0b7f8d9e3f..1506f9e7523f 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -1,14 +1,19 @@ mod common; use anyhow::Result; -use base64::Engine as _; use codex_app_server_protocol::JSONRPCError; use codex_app_server_protocol::JSONRPCMessage; use codex_app_server_protocol::JSONRPCResponse; use codex_exec_server::Environment; +use codex_exec_server::ExecServerClient; use codex_exec_server::ExecutorFileSystem; use codex_exec_server::FileSystemSandboxContext; +use codex_exec_server::FsCloseParams; +use codex_exec_server::FsOpenParams; +use codex_exec_server::FsReadBlockParams; +use codex_exec_server::FsReadBlockResponse; use codex_exec_server::InitializeParams; +use codex_exec_server::RemoteExecServerConnectArgs; use codex_protocol::models::PermissionProfile; use codex_protocol::permissions::FileSystemAccessMode; use codex_protocol::permissions::FileSystemPath; @@ -39,12 +44,6 @@ struct OpenFileResponse { handle_id: String, } -#[derive(Debug, Deserialize)] -struct ReadBlockResponse { - chunk: String, - eof: bool, -} - #[tokio::test] async fn stream_stops_after_an_exact_block_boundary() -> Result<()> { let server = exec_server().await?; @@ -177,67 +176,61 @@ async fn stream_keeps_reading_the_open_file_after_path_replacement() -> Result<( #[tokio::test] async fn read_block_supports_non_sequential_offsets_and_lengths() -> Result<()> { let mut server = exec_server().await?; - initialize_exec_server(&mut server).await?; + let client = ExecServerClient::connect_websocket(RemoteExecServerConnectArgs::new( + server.websocket_url().to_string(), + "file-stream-protocol-test".to_string(), + )) + .await?; let tmp = TempDir::new()?; let path = tmp.path().join("non-sequential.bin"); std::fs::write(&path, b"0123456789")?; - let open: OpenFileResponse = rpc_call( - &mut server, - "fs/open", - serde_json::json!({ - "handleId": Uuid::new_v4().to_string(), - "path": PathUri::from_path(path)?, - "sandbox": null, - }), - ) - .await?; + let open = client + .fs_open(FsOpenParams { + handle_id: Uuid::new_v4().to_string(), + path: PathUri::from_path(path)?, + sandbox: None, + }) + .await?; + let mut blocks = Vec::new(); + for (offset, len) in [(6, 3), (1, 2), (8, 4), (0, 2)] { + blocks.push( + client + .fs_read_block(FsReadBlockParams { + handle_id: open.handle_id.clone(), + offset, + len, + }) + .await?, + ); + } assert_eq!( - read_block( - &mut server, - &open.handle_id, - /*offset*/ 6, - /*len*/ 3 - ) - .await?, - (b"678".to_vec(), false) - ); - assert_eq!( - read_block( - &mut server, - &open.handle_id, - /*offset*/ 1, - /*len*/ 2 - ) - .await?, - (b"12".to_vec(), false) - ); - assert_eq!( - read_block( - &mut server, - &open.handle_id, - /*offset*/ 8, - /*len*/ 4 - ) - .await?, - (b"89".to_vec(), true) - ); - assert_eq!( - read_block( - &mut server, - &open.handle_id, - /*offset*/ 0, - /*len*/ 2 - ) - .await?, - (b"01".to_vec(), false) + blocks, + vec![ + FsReadBlockResponse { + chunk: b"678".to_vec().into(), + eof: false, + }, + FsReadBlockResponse { + chunk: b"12".to_vec().into(), + eof: false, + }, + FsReadBlockResponse { + chunk: b"89".to_vec().into(), + eof: true, + }, + FsReadBlockResponse { + chunk: b"01".to_vec().into(), + eof: false, + }, + ] ); - let _: serde_json::Value = rpc_call( - &mut server, - "fs/close", - serde_json::json!({ "handleId": open.handle_id }), - ) - .await?; + client + .fs_close(FsCloseParams { + handle_id: open.handle_id, + }) + .await?; + drop(client); server.shutdown().await?; Ok(()) } @@ -322,24 +315,6 @@ async fn initialize_exec_server(server: &mut ExecServerHarness) -> Result<()> { Ok(()) } -async fn read_block( - server: &mut ExecServerHarness, - handle_id: &str, - offset: u64, - len: usize, -) -> Result<(Vec, bool)> { - let response: ReadBlockResponse = rpc_call( - server, - "fs/readBlock", - serde_json::json!({ "handleId": handle_id, "offset": offset, "len": len }), - ) - .await?; - Ok(( - base64::engine::general_purpose::STANDARD.decode(response.chunk)?, - response.eof, - )) -} - async fn rpc_call( server: &mut ExecServerHarness, method: &str, From 701738caec389854b5c78c000dab895fb115248f Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 13:26:21 -0700 Subject: [PATCH 14/18] codex: simplify file read streaming (#28354) --- codex-rs/exec-server/Cargo.toml | 2 +- codex-rs/exec-server/src/local_file_system.rs | 22 +-- codex-rs/exec-server/tests/file_stream.rs | 136 +++++------------- .../exec-server/tests/file_system/shared.rs | 9 +- 4 files changed, 46 insertions(+), 123 deletions(-) diff --git a/codex-rs/exec-server/Cargo.toml b/codex-rs/exec-server/Cargo.toml index 687a1ee59c7b..1879400f048f 100644 --- a/codex-rs/exec-server/Cargo.toml +++ b/codex-rs/exec-server/Cargo.toml @@ -44,7 +44,7 @@ tokio = { workspace = true, features = [ "sync", "time", ] } -tokio-util = { workspace = true, features = ["rt"] } +tokio-util = { workspace = true, features = ["io", "rt"] } tokio-tungstenite = { workspace = true } tracing = { workspace = true } uuid = { workspace = true, features = ["v4"] } diff --git a/codex-rs/exec-server/src/local_file_system.rs b/codex-rs/exec-server/src/local_file_system.rs index 753a531fec88..a7711c79424d 100644 --- a/codex-rs/exec-server/src/local_file_system.rs +++ b/codex-rs/exec-server/src/local_file_system.rs @@ -1,4 +1,3 @@ -use bytes::Bytes; use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; use std::path::Path; @@ -8,7 +7,7 @@ use std::sync::LazyLock; use std::time::SystemTime; use std::time::UNIX_EPOCH; use tokio::io; -use tokio::io::AsyncReadExt; +use tokio_util::io::ReaderStream; use crate::CopyOptions; use crate::CreateDirectoryOptions; @@ -530,24 +529,9 @@ impl DirectFileSystem { sandbox: Option<&FileSystemSandboxContext>, ) -> FileSystemResult { let file = self.open_file_for_read(path, sandbox).await?; - Ok(FileSystemReadStream::new(futures::stream::try_unfold( + Ok(FileSystemReadStream::new(ReaderStream::with_capacity( file, - |mut file| async move { - let mut bytes = vec![0; FILE_READ_CHUNK_SIZE]; - let mut bytes_read = 0; - while bytes_read < bytes.len() { - let read = file.read(&mut bytes[bytes_read..]).await?; - if read == 0 { - break; - } - bytes_read += read; - } - if bytes_read == 0 { - return Ok(None); - } - bytes.truncate(bytes_read); - Ok(Some((Bytes::from(bytes), file))) - }, + FILE_READ_CHUNK_SIZE, ))) } diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index 1506f9e7523f..4d3f889594dd 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -1,18 +1,15 @@ mod common; use anyhow::Result; -use codex_app_server_protocol::JSONRPCError; -use codex_app_server_protocol::JSONRPCMessage; -use codex_app_server_protocol::JSONRPCResponse; use codex_exec_server::Environment; use codex_exec_server::ExecServerClient; +use codex_exec_server::ExecServerError; use codex_exec_server::ExecutorFileSystem; use codex_exec_server::FileSystemSandboxContext; use codex_exec_server::FsCloseParams; use codex_exec_server::FsOpenParams; use codex_exec_server::FsReadBlockParams; use codex_exec_server::FsReadBlockResponse; -use codex_exec_server::InitializeParams; use codex_exec_server::RemoteExecServerConnectArgs; use codex_protocol::models::PermissionProfile; use codex_protocol::permissions::FileSystemAccessMode; @@ -24,26 +21,17 @@ use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; use futures::TryStreamExt; use pretty_assertions::assert_eq; -use serde::Deserialize; -use serde::de::DeserializeOwned; use std::sync::Arc; use std::time::Duration; use tempfile::TempDir; use tokio::time::timeout; use uuid::Uuid; -use crate::common::exec_server::ExecServerHarness; use crate::common::exec_server::exec_server; const BLOCK_SIZE: usize = 1024 * 1024; const OPEN_FILE_LIMIT: usize = 128; -#[derive(Debug, Deserialize)] -#[serde(rename_all = "camelCase")] -struct OpenFileResponse { - handle_id: String, -} - #[tokio::test] async fn stream_stops_after_an_exact_block_boundary() -> Result<()> { let server = exec_server().await?; @@ -238,111 +226,61 @@ async fn read_block_supports_non_sequential_offsets_and_lengths() -> Result<()> #[tokio::test] async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> Result<()> { let mut server = exec_server().await?; - initialize_exec_server(&mut server).await?; + let client = ExecServerClient::connect_websocket(RemoteExecServerConnectArgs::new( + server.websocket_url().to_string(), + "file-stream-protocol-test".to_string(), + )) + .await?; let tmp = TempDir::new()?; let path = tmp.path().join("limited.bin"); std::fs::write(&path, b"limited")?; let path = PathUri::from_path(path)?; let mut handles = Vec::with_capacity(OPEN_FILE_LIMIT); for _ in 0..OPEN_FILE_LIMIT { - let open: OpenFileResponse = rpc_call( - &mut server, - "fs/open", - serde_json::json!({ - "handleId": Uuid::new_v4().to_string(), - "path": path, - "sandbox": null, - }), - ) - .await?; + let open = client + .fs_open(FsOpenParams { + handle_id: Uuid::new_v4().to_string(), + path: path.clone(), + sandbox: None, + }) + .await?; handles.push(open.handle_id); } - let response = rpc_message( - &mut server, - "fs/open", - serde_json::json!({ - "handleId": Uuid::new_v4().to_string(), - "path": path, - "sandbox": null, - }), - ) - .await?; - let JSONRPCMessage::Error(JSONRPCError { error, .. }) = response else { - anyhow::bail!("expected opening beyond the limit to fail, got {response:?}"); + let error = client + .fs_open(FsOpenParams { + handle_id: Uuid::new_v4().to_string(), + path: path.clone(), + sandbox: None, + }) + .await + .expect_err("opening beyond the limit should fail"); + let ExecServerError::Server { code, message } = error else { + anyhow::bail!("expected server error, got {error:?}"); }; assert_eq!( - (error.code, error.message), + (code, message), ( -32600, format!("at most {OPEN_FILE_LIMIT} file reads may be open per connection"), ) ); - let _: serde_json::Value = rpc_call( - &mut server, - "fs/close", - serde_json::json!({ "handleId": handles[0] }), - ) - .await?; - let _: OpenFileResponse = rpc_call( - &mut server, - "fs/open", - serde_json::json!({ - "handleId": Uuid::new_v4().to_string(), - "path": path, - "sandbox": null, - }), - ) - .await?; - server.shutdown().await?; - Ok(()) -} - -async fn initialize_exec_server(server: &mut ExecServerHarness) -> Result<()> { - let _: serde_json::Value = rpc_call( - server, - "initialize", - serde_json::to_value(InitializeParams { - client_name: "file-stream-protocol-test".to_string(), - resume_session_id: None, - })?, - ) - .await?; - server - .send_notification("initialized", serde_json::json!({})) + client + .fs_close(FsCloseParams { + handle_id: handles.remove(0), + }) .await?; - Ok(()) -} - -async fn rpc_call( - server: &mut ExecServerHarness, - method: &str, - params: serde_json::Value, -) -> Result -where - T: DeserializeOwned, -{ - let response = rpc_message(server, method, params).await?; - let JSONRPCMessage::Response(JSONRPCResponse { result, .. }) = response else { - anyhow::bail!("expected successful `{method}` response, got {response:?}"); - }; - Ok(serde_json::from_value(result)?) -} - -async fn rpc_message( - server: &mut ExecServerHarness, - method: &str, - params: serde_json::Value, -) -> Result { - let request_id = server.send_request(method, params).await?; - server - .wait_for_event(|event| match event { - JSONRPCMessage::Response(JSONRPCResponse { id, .. }) - | JSONRPCMessage::Error(JSONRPCError { id, .. }) => id == &request_id, - JSONRPCMessage::Request(_) | JSONRPCMessage::Notification(_) => false, + client + .fs_open(FsOpenParams { + handle_id: Uuid::new_v4().to_string(), + path, + sandbox: None, }) - .await + .await?; + drop(client); + server.shutdown().await?; + Ok(()) } fn connect_file_system(websocket_url: &str) -> Result> { diff --git a/codex-rs/exec-server/tests/file_system/shared.rs b/codex-rs/exec-server/tests/file_system/shared.rs index 4b48f1eed243..97e4f5a51f61 100644 --- a/codex-rs/exec-server/tests/file_system/shared.rs +++ b/codex-rs/exec-server/tests/file_system/shared.rs @@ -198,7 +198,7 @@ async fn file_system_read_file_returns_bytes( #[test_case(FileSystemImplementation::Local ; "local")] #[test_case(FileSystemImplementation::Remote ; "remote")] #[tokio::test(flavor = "multi_thread", worker_threads = 2)] -async fn file_system_read_file_stream_returns_one_mib_chunks( +async fn file_system_read_file_stream_returns_bounded_chunks( implementation: FileSystemImplementation, ) -> Result<()> { let context = create_file_system_context(implementation).await?; @@ -218,9 +218,10 @@ async fn file_system_read_file_stream_returns_one_mib_chunks( .try_collect::>() .await?; - assert_eq!( - chunks.iter().map(bytes::Bytes::len).collect::>(), - vec![FILE_READ_CHUNK_SIZE, FILE_READ_CHUNK_SIZE, 17] + assert!( + chunks + .iter() + .all(|chunk| !chunk.is_empty() && chunk.len() <= FILE_READ_CHUNK_SIZE) ); assert_eq!( chunks From ca49b47d8a5593db4f5a9e6fe88ba12534f3b64d Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Mon, 15 Jun 2026 13:38:57 -0700 Subject: [PATCH 15/18] codex: fix CI failure on PR #28354 --- codex-rs/exec-server/tests/file_stream.rs | 2 ++ 1 file changed, 2 insertions(+) diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index 4d3f889594dd..d708d5610869 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -22,8 +22,10 @@ use codex_utils_path_uri::PathUri; use futures::TryStreamExt; use pretty_assertions::assert_eq; use std::sync::Arc; +#[cfg(unix)] use std::time::Duration; use tempfile::TempDir; +#[cfg(unix)] use tokio::time::timeout; use uuid::Uuid; From e43701d6bc8472608ad81bc9fa6b51374184b8fd Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Tue, 16 Jun 2026 08:24:38 -0700 Subject: [PATCH 16/18] codex: bound file read handle IDs (#28354) --- .../exec-server/src/remote_file_stream.rs | 2 +- .../src/server/file_system_handler.rs | 14 +++++++ codex-rs/exec-server/tests/file_stream.rs | 42 +++++++++++++++++-- 3 files changed, 53 insertions(+), 5 deletions(-) diff --git a/codex-rs/exec-server/src/remote_file_stream.rs b/codex-rs/exec-server/src/remote_file_stream.rs index 2a22dcf2ae66..107dea51d7d5 100644 --- a/codex-rs/exec-server/src/remote_file_stream.rs +++ b/codex-rs/exec-server/src/remote_file_stream.rs @@ -28,7 +28,7 @@ pub(super) async fn open( ) -> FileSystemResult { let registration = FileReadRegistration { client, - handle_id: Uuid::new_v4().to_string(), + handle_id: Uuid::new_v4().simple().to_string(), runtime: tokio::runtime::Handle::try_current().ok(), active: true, }; diff --git a/codex-rs/exec-server/src/server/file_system_handler.rs b/codex-rs/exec-server/src/server/file_system_handler.rs index 603d50735533..ddf17f21cd93 100644 --- a/codex-rs/exec-server/src/server/file_system_handler.rs +++ b/codex-rs/exec-server/src/server/file_system_handler.rs @@ -39,6 +39,8 @@ use crate::rpc::internal_error; use crate::rpc::invalid_request; use crate::rpc::not_found; +const MAX_FILE_READ_HANDLE_ID_BYTES: usize = 32; + #[derive(Clone)] pub(crate) struct FileSystemHandler { file_system: LocalFileSystem, @@ -61,6 +63,7 @@ impl FileSystemHandler { &self, params: FsOpenParams, ) -> Result { + validate_file_read_handle_id(¶ms.handle_id)?; let file = self .file_system .open_file_for_read(¶ms.path, params.sandbox.as_ref()) @@ -78,6 +81,7 @@ impl FileSystemHandler { &self, params: FsReadBlockParams, ) -> Result { + validate_file_read_handle_id(¶ms.handle_id)?; let block = self .file_reads .read_block(¶ms.handle_id, params.offset, params.len) @@ -93,6 +97,7 @@ impl FileSystemHandler { &self, params: FsCloseParams, ) -> Result { + validate_file_read_handle_id(¶ms.handle_id)?; self.file_reads.close(¶ms.handle_id).await; Ok(FsCloseResponse {}) } @@ -229,6 +234,15 @@ impl FileSystemHandler { } } +fn validate_file_read_handle_id(handle_id: &str) -> Result<(), JSONRPCErrorError> { + if handle_id.len() > MAX_FILE_READ_HANDLE_ID_BYTES { + return Err(invalid_request(format!( + "file read handle ID must not exceed {MAX_FILE_READ_HANDLE_ID_BYTES} bytes" + ))); + } + Ok(()) +} + fn map_fs_error(err: io::Error) -> JSONRPCErrorError { match err.kind() { io::ErrorKind::NotFound => not_found(err.to_string()), diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index d708d5610869..79e318f66a8a 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -176,7 +176,7 @@ async fn read_block_supports_non_sequential_offsets_and_lengths() -> Result<()> std::fs::write(&path, b"0123456789")?; let open = client .fs_open(FsOpenParams { - handle_id: Uuid::new_v4().to_string(), + handle_id: Uuid::new_v4().simple().to_string(), path: PathUri::from_path(path)?, sandbox: None, }) @@ -241,7 +241,7 @@ async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> for _ in 0..OPEN_FILE_LIMIT { let open = client .fs_open(FsOpenParams { - handle_id: Uuid::new_v4().to_string(), + handle_id: Uuid::new_v4().simple().to_string(), path: path.clone(), sandbox: None, }) @@ -251,7 +251,7 @@ async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> let error = client .fs_open(FsOpenParams { - handle_id: Uuid::new_v4().to_string(), + handle_id: Uuid::new_v4().simple().to_string(), path: path.clone(), sandbox: None, }) @@ -275,7 +275,7 @@ async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> .await?; client .fs_open(FsOpenParams { - handle_id: Uuid::new_v4().to_string(), + handle_id: Uuid::new_v4().simple().to_string(), path, sandbox: None, }) @@ -285,6 +285,40 @@ async fn open_enforces_the_per_connection_limit_and_close_releases_capacity() -> Ok(()) } +#[tokio::test] +async fn open_rejects_handle_ids_longer_than_32_bytes() -> Result<()> { + let server = exec_server().await?; + let client = ExecServerClient::connect_websocket(RemoteExecServerConnectArgs::new( + server.websocket_url().to_string(), + "file-stream-protocol-test".to_string(), + )) + .await?; + let tmp = TempDir::new()?; + let path = tmp.path().join("handle-id-limit.bin"); + std::fs::write(&path, b"limited")?; + + let error = client + .fs_open(FsOpenParams { + handle_id: "x".repeat(33), + path: PathUri::from_path(path)?, + sandbox: None, + }) + .await + .expect_err("oversized handle ID should fail"); + + let ExecServerError::Server { code, message } = error else { + anyhow::bail!("expected server error, got {error:?}"); + }; + assert_eq!( + (code, message), + ( + -32600, + "file read handle ID must not exceed 32 bytes".to_string(), + ) + ); + Ok(()) +} + fn connect_file_system(websocket_url: &str) -> Result> { let environment = Environment::create_for_tests(Some(websocket_url.to_string()))?; Ok(environment.get_filesystem()) From 699b25fa76904e82ab7a251de56df3925ea0f4b6 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Tue, 16 Jun 2026 08:58:36 -0700 Subject: [PATCH 17/18] codex: validate opened file handles (#28354) --- codex-rs/Cargo.lock | 2 + codex-rs/exec-server/Cargo.toml | 9 +++ codex-rs/exec-server/src/lib.rs | 1 + codex-rs/exec-server/src/local_file_system.rs | 36 +++++---- codex-rs/exec-server/src/regular_file.rs | 48 ++++++++++++ codex-rs/exec-server/tests/file_stream.rs | 77 ++++++++++++++++--- 6 files changed, 149 insertions(+), 24 deletions(-) create mode 100644 codex-rs/exec-server/src/regular_file.rs diff --git a/codex-rs/Cargo.lock b/codex-rs/Cargo.lock index 3ae1695a7637..6a1a75c96ddf 100644 --- a/codex-rs/Cargo.lock +++ b/codex-rs/Cargo.lock @@ -2851,6 +2851,7 @@ dependencies = [ "ctor 0.6.3", "futures", "http 1.4.0", + "libc", "pretty_assertions", "prost 0.14.3", "reqwest 0.12.28", @@ -2866,6 +2867,7 @@ dependencies = [ "toml 0.9.11+spec-1.1.0", "tracing", "uuid", + "windows-sys 0.52.0", "wiremock", ] diff --git a/codex-rs/exec-server/Cargo.toml b/codex-rs/exec-server/Cargo.toml index 1879400f048f..a9b626fcfcd7 100644 --- a/codex-rs/exec-server/Cargo.toml +++ b/codex-rs/exec-server/Cargo.toml @@ -49,6 +49,15 @@ tokio-tungstenite = { workspace = true } tracing = { workspace = true } uuid = { workspace = true, features = ["v4"] } +[target.'cfg(unix)'.dependencies] +libc = { workspace = true } + +[target.'cfg(windows)'.dependencies] +windows-sys = { version = "0.52", features = [ + "Win32_Foundation", + "Win32_Storage_FileSystem", +] } + [dev-dependencies] anyhow = { workspace = true } codex-test-binary-support = { workspace = true } diff --git a/codex-rs/exec-server/src/lib.rs b/codex-rs/exec-server/src/lib.rs index 9d79b7b6c3e9..edae7e93508c 100644 --- a/codex-rs/exec-server/src/lib.rs +++ b/codex-rs/exec-server/src/lib.rs @@ -14,6 +14,7 @@ mod local_process; mod process; mod process_id; mod protocol; +mod regular_file; mod relay; mod relay_proto; mod remote; diff --git a/codex-rs/exec-server/src/local_file_system.rs b/codex-rs/exec-server/src/local_file_system.rs index a7711c79424d..3129606d6b27 100644 --- a/codex-rs/exec-server/src/local_file_system.rs +++ b/codex-rs/exec-server/src/local_file_system.rs @@ -7,6 +7,7 @@ use std::sync::LazyLock; use std::time::SystemTime; use std::time::UNIX_EPOCH; use tokio::io; +use tokio::io::AsyncReadExt; use tokio_util::io::ReaderStream; use crate::CopyOptions; @@ -21,10 +22,18 @@ use crate::FileSystemResult; use crate::FileSystemSandboxContext; use crate::ReadDirectoryEntry; use crate::RemoveOptions; +use crate::regular_file; use crate::sandboxed_file_system::SandboxedFileSystem; const MAX_READ_FILE_BYTES: u64 = 512 * 1024 * 1024; +fn file_too_large_error() -> io::Error { + io::Error::new( + io::ErrorKind::InvalidInput, + format!("file is too large to read: limit is {MAX_READ_FILE_BYTES} bytes"), + ) +} + pub static LOCAL_FS: LazyLock> = LazyLock::new(|| -> Arc { Arc::new(LocalFileSystem::unsandboxed()) }); @@ -485,13 +494,7 @@ impl DirectFileSystem { ) -> FileSystemResult { reject_sandbox_context(sandbox)?; let path = path.to_abs_path()?; - if !tokio::fs::metadata(path.as_path()).await?.is_file() { - return Err(io::Error::new( - io::ErrorKind::InvalidInput, - format!("path `{}` is not a file", path.display()), - )); - } - tokio::fs::File::open(path.as_path()).await + regular_file::open(path.as_path()).await } async fn canonicalize( @@ -511,16 +514,19 @@ impl DirectFileSystem { path: &PathUri, sandbox: Option<&FileSystemSandboxContext>, ) -> FileSystemResult> { - reject_sandbox_context(sandbox)?; - let path = path.to_abs_path()?; - let metadata = tokio::fs::metadata(path.as_path()).await?; + let file = self.open_file_for_read(path, sandbox).await?; + let metadata = file.metadata().await?; if metadata.len() > MAX_READ_FILE_BYTES { - return Err(io::Error::new( - io::ErrorKind::InvalidInput, - format!("file is too large to read: limit is {MAX_READ_FILE_BYTES} bytes"), - )); + return Err(file_too_large_error()); + } + let mut bytes = Vec::with_capacity(metadata.len() as usize); + file.take(MAX_READ_FILE_BYTES + 1) + .read_to_end(&mut bytes) + .await?; + if bytes.len() as u64 > MAX_READ_FILE_BYTES { + return Err(file_too_large_error()); } - tokio::fs::read(path.as_path()).await + Ok(bytes) } async fn read_file_stream( diff --git a/codex-rs/exec-server/src/regular_file.rs b/codex-rs/exec-server/src/regular_file.rs new file mode 100644 index 000000000000..55540faddfb1 --- /dev/null +++ b/codex-rs/exec-server/src/regular_file.rs @@ -0,0 +1,48 @@ +use std::io; +use std::path::Path; + +pub(crate) async fn open(path: &Path) -> io::Result { + let mut options = tokio::fs::OpenOptions::new(); + options.read(true); + configure_open(&mut options); + + let file = options.open(path).await?; + if !is_disk_file(&file) || !file.metadata().await?.is_file() { + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!("path `{}` is not a file", path.display()), + )); + } + Ok(file) +} + +#[cfg(unix)] +fn configure_open(options: &mut tokio::fs::OpenOptions) { + options.custom_flags(libc::O_NONBLOCK); +} + +#[cfg(windows)] +fn configure_open(options: &mut tokio::fs::OpenOptions) { + use windows_sys::Win32::Storage::FileSystem::SECURITY_IDENTIFICATION; + + options.security_qos_flags(SECURITY_IDENTIFICATION); +} + +#[cfg(not(any(unix, windows)))] +fn configure_open(_options: &mut tokio::fs::OpenOptions) {} + +#[cfg(windows)] +fn is_disk_file(file: &tokio::fs::File) -> bool { + use std::os::windows::io::AsRawHandle; + use windows_sys::Win32::Foundation::HANDLE; + use windows_sys::Win32::Storage::FileSystem::FILE_TYPE_DISK; + use windows_sys::Win32::Storage::FileSystem::GetFileType; + + // SAFETY: `file` owns this handle for the duration of the call. + unsafe { GetFileType(file.as_raw_handle() as HANDLE) == FILE_TYPE_DISK } +} + +#[cfg(not(windows))] +fn is_disk_file(_file: &tokio::fs::File) -> bool { + true +} diff --git a/codex-rs/exec-server/tests/file_stream.rs b/codex-rs/exec-server/tests/file_stream.rs index 79e318f66a8a..e30e8c468902 100644 --- a/codex-rs/exec-server/tests/file_stream.rs +++ b/codex-rs/exec-server/tests/file_stream.rs @@ -22,10 +22,12 @@ use codex_utils_path_uri::PathUri; use futures::TryStreamExt; use pretty_assertions::assert_eq; use std::sync::Arc; -#[cfg(unix)] +#[cfg(any(unix, windows))] use std::time::Duration; use tempfile::TempDir; -#[cfg(unix)] +#[cfg(windows)] +use tokio::net::windows::named_pipe::ServerOptions; +#[cfg(any(unix, windows))] use tokio::time::timeout; use uuid::Uuid; @@ -104,7 +106,7 @@ async fn stream_rejects_platform_sandbox() -> Result<()> { #[cfg(unix)] #[tokio::test] -async fn stream_rejects_fifo_without_waiting_for_a_writer() -> Result<()> { +async fn file_reads_reject_fifo_without_waiting_for_a_writer() -> Result<()> { let server = exec_server().await?; let file_system = connect_file_system(server.websocket_url())?; let tmp = TempDir::new()?; @@ -118,18 +120,75 @@ async fn stream_rejects_fifo_without_waiting_for_a_writer() -> Result<()> { ); } - let result = timeout( + let path_uri = PathUri::from_path(&path)?; + let read_error = timeout( Duration::from_secs(1), - file_system.read_file_stream(&PathUri::from_path(&path)?, /*sandbox*/ None), + file_system.read_file(&path_uri, /*sandbox*/ None), ) .await - .expect("opening a FIFO should not wait for a writer"); - let Err(error) = result else { + .expect("reading a FIFO should not wait for a writer") + .expect_err("reading a FIFO should be rejected"); + let stream_result = timeout( + Duration::from_secs(1), + file_system.read_file_stream(&path_uri, /*sandbox*/ None), + ) + .await + .expect("streaming a FIFO should not wait for a writer"); + let Err(stream_error) = stream_result else { panic!("streaming a FIFO should be rejected"); }; + let expected = format!("path `{}` is not a file", path.display()); assert_eq!( - error.to_string(), - format!("path `{}` is not a file", path.display()) + (read_error.to_string(), stream_error.to_string()), + (expected.clone(), expected) + ); + Ok(()) +} + +#[cfg(windows)] +#[tokio::test] +async fn file_reads_reject_named_pipes() -> Result<()> { + let server = exec_server().await?; + let file_system = connect_file_system(server.websocket_url())?; + + let read_path = format!(r"\\.\pipe\codex-fs-read-{}", Uuid::new_v4()); + let _read_pipe = ServerOptions::new() + .first_pipe_instance(true) + .create(&read_path)?; + let read_error = timeout( + Duration::from_secs(1), + file_system.read_file( + &PathUri::from_path(std::path::Path::new(&read_path))?, + /*sandbox*/ None, + ), + ) + .await + .expect("reading a named pipe should not hang") + .expect_err("reading a named pipe should be rejected"); + + let stream_path = format!(r"\\.\pipe\codex-fs-stream-{}", Uuid::new_v4()); + let _stream_pipe = ServerOptions::new() + .first_pipe_instance(true) + .create(&stream_path)?; + let stream_result = timeout( + Duration::from_secs(1), + file_system.read_file_stream( + &PathUri::from_path(std::path::Path::new(&stream_path))?, + /*sandbox*/ None, + ), + ) + .await + .expect("streaming a named pipe should not hang"); + let Err(stream_error) = stream_result else { + panic!("streaming a named pipe should be rejected"); + }; + + assert_eq!( + (read_error.kind(), stream_error.kind()), + ( + std::io::ErrorKind::InvalidInput, + std::io::ErrorKind::InvalidInput, + ) ); Ok(()) } From ee57b6d7220d6572b73a037ef40d7603076dbe97 Mon Sep 17 00:00:00 2001 From: pakrym-oai Date: Tue, 16 Jun 2026 09:28:52 -0700 Subject: [PATCH 18/18] codex: fix CI failure on PR #28354 --- .../tests/suite/v2/external_agent_config.rs | 69 ++++++++----------- 1 file changed, 30 insertions(+), 39 deletions(-) diff --git a/codex-rs/app-server/tests/suite/v2/external_agent_config.rs b/codex-rs/app-server/tests/suite/v2/external_agent_config.rs index b75d703c4805..927fab760709 100644 --- a/codex-rs/app-server/tests/suite/v2/external_agent_config.rs +++ b/codex-rs/app-server/tests/suite/v2/external_agent_config.rs @@ -718,24 +718,14 @@ async fn external_agent_config_import_returns_before_background_session_import_f let session_path = session_dir.join("session.jsonl"); std::fs::create_dir_all(&project_root)?; std::fs::create_dir_all(&session_dir)?; - std::fs::write( - &session_path, - serde_json::json!({ - "type": "user", - "cwd": &project_root, - "timestamp": &recent_timestamp, - "message": { "content": "first request" }, - }) - .to_string(), - )?; - - let project_config_dir = project_root.join(".codex"); - std::fs::create_dir_all(&project_config_dir)?; - let project_config = project_config_dir.join("config.toml"); - let status = std::process::Command::new("mkfifo") - .arg(&project_config) - .status()?; - assert!(status.success()); + let session_contents = serde_json::json!({ + "type": "user", + "cwd": &project_root, + "timestamp": &recent_timestamp, + "message": { "content": "first request" }, + }) + .to_string(); + std::fs::write(&session_path, &session_contents)?; let home_dir = codex_home.path().display().to_string(); let mut mcp = @@ -758,6 +748,12 @@ async fn external_agent_config_import_returns_before_background_session_import_f assert_eq!(detected.items.len(), 1); let detected_items = detected.items; + std::fs::remove_file(&session_path)?; + let status = std::process::Command::new("mkfifo") + .arg(&session_path) + .status()?; + assert!(status.success()); + let request_id = mcp .send_raw_request( "externalAgentConfig/import", @@ -796,28 +792,23 @@ async fn external_agent_config_import_returns_before_background_session_import_f let response: ExternalAgentConfigImportResponse = to_response(response)?; assert_eq!(response, ExternalAgentConfigImportResponse {}); - let writer = tokio::spawn(async move { - let mut file = tokio::fs::OpenOptions::new() - .write(true) - .open(&project_config) - .await?; - file.write_all(b"\n").await - }); - timeout(DEFAULT_TIMEOUT, writer).await???; - - let notification = timeout( - DEFAULT_TIMEOUT, - mcp.read_stream_until_notification_message("externalAgentConfig/import/completed"), - ) - .await??; - assert_eq!(notification.method, "externalAgentConfig/import/completed"); + for _ in 0..2 { + timeout(DEFAULT_TIMEOUT, async { + let mut file = tokio::fs::OpenOptions::new() + .write(true) + .open(&session_path) + .await?; + file.write_all(session_contents.as_bytes()).await + }) + .await??; - let notification = timeout( - DEFAULT_TIMEOUT, - mcp.read_stream_until_notification_message("externalAgentConfig/import/completed"), - ) - .await??; - assert_eq!(notification.method, "externalAgentConfig/import/completed"); + let notification = timeout( + DEFAULT_TIMEOUT, + mcp.read_stream_until_notification_message("externalAgentConfig/import/completed"), + ) + .await??; + assert_eq!(notification.method, "externalAgentConfig/import/completed"); + } let request_id = mcp .send_thread_list_request(ThreadListParams {