From 372ae11e5df30a9b7e197a579bd9806e655bcb68 Mon Sep 17 00:00:00 2001 From: Roy Han Date: Tue, 5 May 2026 14:24:25 -0700 Subject: [PATCH] add deferred image content reads --- .../schema/json/ClientRequest.json | 7 + .../schema/json/v2/ThreadResumeParams.json | 7 + .../schema/typescript/v2/LargeContentMode.ts | 5 + .../schema/typescript/v2/index.ts | 1 + .../src/protocol/common.rs | 8 + .../app-server-protocol/src/protocol/v2.rs | 91 ++++++--- codex-rs/app-server/README.md | 52 ++++- codex-rs/app-server/src/message_processor.rs | 3 + codex-rs/app-server/src/request_processors.rs | 4 + .../request_processors/thread_lifecycle.rs | 1 + .../request_processors/thread_processor.rs | 185 +++++++++++++----- .../thread_processor_tests.rs | 1 + codex-rs/app-server/src/thread_state.rs | 2 + .../app-server/tests/common/mcp_process.rs | 12 +- .../app-server/tests/suite/v2/thread_read.rs | 149 ++++++++++++++ .../tests/suite/v2/thread_shell_command.rs | 1 + 16 files changed, 456 insertions(+), 73 deletions(-) create mode 100644 codex-rs/app-server-protocol/schema/typescript/v2/LargeContentMode.ts diff --git a/codex-rs/app-server-protocol/schema/json/ClientRequest.json b/codex-rs/app-server-protocol/schema/json/ClientRequest.json index 9ea9893f5b..153bd9bf37 100644 --- a/codex-rs/app-server-protocol/schema/json/ClientRequest.json +++ b/codex-rs/app-server-protocol/schema/json/ClientRequest.json @@ -1478,6 +1478,13 @@ ], "type": "object" }, + "LargeContentMode": { + "enum": [ + "inline", + "deferred" + ], + "type": "string" + }, "ListMcpServerStatusParams": { "properties": { "cursor": { diff --git a/codex-rs/app-server-protocol/schema/json/v2/ThreadResumeParams.json b/codex-rs/app-server-protocol/schema/json/v2/ThreadResumeParams.json index 9fe5c7f47f..3b24408510 100644 --- a/codex-rs/app-server-protocol/schema/json/v2/ThreadResumeParams.json +++ b/codex-rs/app-server-protocol/schema/json/v2/ThreadResumeParams.json @@ -215,6 +215,13 @@ ], "type": "string" }, + "LargeContentMode": { + "enum": [ + "inline", + "deferred" + ], + "type": "string" + }, "LocalShellAction": { "oneOf": [ { diff --git a/codex-rs/app-server-protocol/schema/typescript/v2/LargeContentMode.ts b/codex-rs/app-server-protocol/schema/typescript/v2/LargeContentMode.ts new file mode 100644 index 0000000000..8d0bb63df4 --- /dev/null +++ b/codex-rs/app-server-protocol/schema/typescript/v2/LargeContentMode.ts @@ -0,0 +1,5 @@ +// GENERATED CODE! DO NOT MODIFY BY HAND! + +// This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. + +export type LargeContentMode = "inline" | "deferred"; diff --git a/codex-rs/app-server-protocol/schema/typescript/v2/index.ts b/codex-rs/app-server-protocol/schema/typescript/v2/index.ts index 129ab6d745..868b2f8698 100644 --- a/codex-rs/app-server-protocol/schema/typescript/v2/index.ts +++ b/codex-rs/app-server-protocol/schema/typescript/v2/index.ts @@ -177,6 +177,7 @@ export type { ItemCompletedNotification } from "./ItemCompletedNotification"; export type { ItemGuardianApprovalReviewCompletedNotification } from "./ItemGuardianApprovalReviewCompletedNotification"; export type { ItemGuardianApprovalReviewStartedNotification } from "./ItemGuardianApprovalReviewStartedNotification"; export type { ItemStartedNotification } from "./ItemStartedNotification"; +export type { LargeContentMode } from "./LargeContentMode"; export type { ListMcpServerStatusParams } from "./ListMcpServerStatusParams"; export type { ListMcpServerStatusResponse } from "./ListMcpServerStatusResponse"; export type { LoginAccountParams } from "./LoginAccountParams"; diff --git a/codex-rs/app-server-protocol/src/protocol/common.rs b/codex-rs/app-server-protocol/src/protocol/common.rs index 4176b8e2da..1e62da8565 100644 --- a/codex-rs/app-server-protocol/src/protocol/common.rs +++ b/codex-rs/app-server-protocol/src/protocol/common.rs @@ -584,6 +584,13 @@ client_request_definitions! { serialization: None, response: v2::ThreadTurnsItemsListResponse, }, + #[experimental("thread/item/content/read")] + ThreadItemContentRead => "thread/item/content/read" { + params: v2::ThreadItemContentReadParams, + // Explicitly concurrent: this primarily reads append-only rollout storage. + serialization: None, + response: v2::ThreadItemContentReadResponse, + }, /// Append raw Responses API items to the thread history without starting a user turn. ThreadInjectItems => "thread/inject_items" { params: v2::ThreadInjectItemsParams, @@ -1844,6 +1851,7 @@ mod tests { cursor: None, limit: None, sort_direction: None, + large_content: None, }, }; assert_eq!(thread_turns_list.serialization_scope(), None); diff --git a/codex-rs/app-server-protocol/src/protocol/v2.rs b/codex-rs/app-server-protocol/src/protocol/v2.rs index da62956465..72053be016 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2.rs @@ -3980,6 +3980,11 @@ pub struct ThreadResumeParams { #[experimental("thread/resume.excludeTurns")] #[serde(default, skip_serializing_if = "std::ops::Not::not")] pub exclude_turns: bool, + /// Controls whether large item payloads are embedded in returned turns or + /// replaced with deferred-content metadata. + #[experimental("thread/resume.largeContent")] + #[ts(optional = nullable)] + pub large_content: Option, /// Deprecated and ignored by app-server. Kept only so older clients can /// continue sending the field while rollout persistence always uses the /// limited history policy. @@ -4677,6 +4682,10 @@ pub struct ThreadTurnsListParams { /// Optional turn pagination direction; defaults to descending. #[ts(optional = nullable)] pub sort_direction: Option, + /// Controls whether large item payloads are embedded in returned turns or + /// replaced with deferred-content metadata. + #[ts(optional = nullable)] + pub large_content: Option, } #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] @@ -4694,6 +4703,15 @@ pub struct ThreadTurnsListResponse { pub backwards_cursor: Option, } +#[derive(Serialize, Deserialize, Debug, Clone, Copy, PartialEq, Eq, JsonSchema, TS, Default)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub enum LargeContentMode { + #[default] + Inline, + Deferred, +} + #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] #[serde(rename_all = "camelCase")] #[ts(export_to = "v2/")] @@ -4709,6 +4727,10 @@ pub struct ThreadTurnsItemsListParams { /// Optional item pagination direction; defaults to ascending. #[ts(optional = nullable)] pub sort_direction: Option, + /// Controls whether large item payloads are embedded in returned items or + /// replaced with deferred-content metadata. + #[ts(optional = nullable)] + pub large_content: Option, } #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] @@ -4726,6 +4748,25 @@ pub struct ThreadTurnsItemsListResponse { pub backwards_cursor: Option, } +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ThreadItemContentReadParams { + pub thread_id: String, + pub turn_id: String, + pub item_id: String, + pub content_id: String, +} + +#[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] +#[serde(rename_all = "camelCase")] +#[ts(export_to = "v2/")] +pub struct ThreadItemContentReadResponse { + pub mime_type: String, + pub data_base64: String, + pub byte_length: u64, +} + #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] #[serde(rename_all = "camelCase")] #[ts(export_to = "v2/")] @@ -9108,6 +9149,31 @@ mod tests { assert_eq!(decoded, response); } + #[test] + fn image_generation_defaults_missing_content_for_legacy_payloads() { + let item: ThreadItem = serde_json::from_value(json!({ + "type": "imageGeneration", + "id": "ig_123", + "status": "completed", + "revisedPrompt": null, + "result": "Zm9v", + "savedPath": null + })) + .expect("legacy image generation item should deserialize"); + + assert_eq!( + item, + ThreadItem::ImageGeneration { + id: "ig_123".to_string(), + status: "completed".to_string(), + revised_prompt: None, + content: ImageGenerationContent::default(), + result: "Zm9v".to_string(), + saved_path: None, + } + ); + } + #[test] fn fs_read_file_params_round_trip() { let params = FsReadFileParams { @@ -12037,29 +12103,4 @@ mod tests { "unexpected error: {err}" ); } - - #[test] - fn image_generation_defaults_missing_content_for_legacy_payloads() { - let item: ThreadItem = serde_json::from_value(json!({ - "type": "imageGeneration", - "id": "ig_123", - "status": "completed", - "revisedPrompt": null, - "result": "Zm9v", - "savedPath": null - })) - .expect("legacy image generation item should deserialize"); - - assert_eq!( - item, - ThreadItem::ImageGeneration { - id: "ig_123".to_string(), - status: "completed".to_string(), - revised_prompt: None, - content: ImageGenerationContent::default(), - result: "Zm9v".to_string(), - saved_path: None, - } - ); - } } diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index 18a9d80821..7c8d853201 100644 --- a/codex-rs/app-server/README.md +++ b/codex-rs/app-server/README.md @@ -143,14 +143,15 @@ Example with notification opt-out: ## API Overview - `thread/start` — create a new thread; emits `thread/started` (including the current `thread.status`) and auto-subscribes you to turn/item events for that thread. When the request includes a `cwd` and the resolved sandbox is `workspace-write` or full access, app-server also marks that project as trusted in the user `config.toml`. Pass `sessionStartSource: "clear"` when starting a replacement thread after clearing the current session so `SessionStart` hooks receive `source: "clear"` instead of the default `"startup"`. For permissions, prefer experimental `permissions` profile selection; the legacy `sandbox` shorthand is still accepted but cannot be combined with `permissions`. Experimental `environments` selects the sticky execution environments for turns on the thread; omit it to use the server default, pass `[]` to disable environments, or pass explicit environment ids with per-environment `cwd`. -- `thread/resume` — reopen an existing thread by id so subsequent `turn/start` calls append to it. Accepts the same permission override rules as `thread/start`. +- `thread/resume` — reopen an existing thread by id so subsequent `turn/start` calls append to it. Accepts the same permission override rules as `thread/start`. Experimental clients can pass `largeContent: "deferred"` to replace large item payloads such as generated-image bytes with metadata placeholders. - `thread/fork` — fork an existing thread into a new thread id by copying the stored history; if the source thread is currently mid-turn, the fork records the same interruption marker as `turn/interrupt` instead of inheriting an unmarked partial turn suffix. The returned `thread.forkedFromId` points at the source thread when known. Accepts `ephemeral: true` for an in-memory temporary fork, emits `thread/started` (including the current `thread.status`), and auto-subscribes you to turn/item events for the new thread. Experimental clients can pass `excludeTurns: true` when they plan to page fork history via `thread/turns/list` instead of receiving the full turn array immediately. Accepts the same permission override rules as `thread/start`. - `thread/start`, `thread/resume`, and `thread/fork` responses include the legacy `sandbox` compatibility projection. Experimental clients can read response `permissionProfile` for the exact active runtime permissions and `activePermissionProfile` for the named or implicit built-in profile identity/provenance when known. - `thread/list` — page through stored rollouts; supports cursor-based pagination and optional `modelProviders`, `sourceKinds`, `archived`, `cwd`, and `searchTerm` filters. Each returned `thread` includes `status` (`ThreadStatus`), defaulting to `notLoaded` when the thread is not currently loaded. - `thread/loaded/list` — list the thread ids currently loaded in memory. - `thread/read` — read a stored thread by id without resuming it; optionally include turns via `includeTurns`. The returned `thread` includes `status` (`ThreadStatus`), defaulting to `notLoaded` when the thread is not currently loaded. -- `thread/turns/list` — experimental; page through a stored thread’s turn history without resuming it; supports cursor-based pagination with `sortDirection`, `nextCursor`, and `backwardsCursor`. -- `thread/turns/items/list` — experimental; page through one stored turn’s items with the same cursor model. +- `thread/turns/list` — experimental; page through a stored thread’s turn history without resuming it; supports cursor-based pagination with `sortDirection`, `nextCursor`, and `backwardsCursor`, plus `largeContent: "deferred"` for generated-image payloads. +- `thread/turns/items/list` — experimental; page through one stored turn’s items with the same cursor model and optional `largeContent: "deferred"` handling. +- `thread/item/content/read` — experimental; fetch deferred large-content bytes for one stored item. - `thread/metadata/update` — patch stored thread metadata in sqlite; currently supports updating persisted `gitInfo` fields and returns the refreshed `thread`. - `thread/memoryMode/set` — experimental; set a thread’s persisted memory eligibility to `"enabled"` or `"disabled"` for either a loaded thread or a stored rollout; returns `{}` on success. - `memory/reset` — experimental; clear the current `CODEX_HOME/memories` directory and reset persisted memory stage data in sqlite while preserving existing thread memory modes; returns `{}` on success. @@ -443,6 +444,51 @@ Every returned `Turn` includes `itemsView`, which tells clients whether the `ite } } ``` +### Example: Defer generated-image content (experimental) + +Pass `largeContent: "deferred"` when listing a turn's items to keep generated-image bytes out of the page response. Deferred `imageGeneration` items retain the legacy `result` field as an empty string and expose a structured `content` descriptor that can be loaded separately. For a just-completed live turn, callers may need to retry a content read briefly while rollout persistence catches up with the live turn view. + +```json +{ "method": "thread/turns/items/list", "id": 25, "params": { + "threadId": "thr_123", + "turnId": "turn_456", + "limit": 100, + "sortDirection": "asc", + "largeContent": "deferred" +} } +{ "id": 25, "result": { + "data": [{ + "type": "imageGeneration", + "id": "ig_789", + "status": "completed", + "revisedPrompt": null, + "content": { + "type": "deferred", + "contentId": "result", + "mimeType": "image/png", + "byteLength": 1234567, + "width": null, + "height": null + }, + "result": "", + "savedPath": null + }], + "nextCursor": null, + "backwardsCursor": "newer-items-cursor-or-null" +} } +{ "method": "thread/item/content/read", "id": 26, "params": { + "threadId": "thr_123", + "turnId": "turn_456", + "itemId": "ig_789", + "contentId": "result" +} } +{ "id": 26, "result": { + "mimeType": "image/png", + "dataBase64": "...", + "byteLength": 1234567 +} } +``` + ### Example: Update stored thread metadata Use `thread/metadata/update` to patch sqlite-backed metadata for a thread without resuming it. Today this supports persisted `gitInfo`; omitted fields are left unchanged, while explicit `null` clears a stored value. diff --git a/codex-rs/app-server/src/message_processor.rs b/codex-rs/app-server/src/message_processor.rs index 1790563b2b..60fe5d17e2 100644 --- a/codex-rs/app-server/src/message_processor.rs +++ b/codex-rs/app-server/src/message_processor.rs @@ -1047,6 +1047,9 @@ impl MessageProcessor { ClientRequest::ThreadTurnsItemsList { params, .. } => { self.thread_processor.thread_turns_items_list(params).await } + ClientRequest::ThreadItemContentRead { params, .. } => { + self.thread_processor.thread_item_content_read(params).await + } ClientRequest::ThreadShellCommand { params, .. } => { self.thread_processor .thread_shell_command(&request_id, params) diff --git a/codex-rs/app-server/src/request_processors.rs b/codex-rs/app-server/src/request_processors.rs index 75da930f3b..7cb5a0a104 100644 --- a/codex-rs/app-server/src/request_processors.rs +++ b/codex-rs/app-server/src/request_processors.rs @@ -70,9 +70,11 @@ use codex_app_server_protocol::GitInfo as ApiGitInfo; use codex_app_server_protocol::HookMetadata; use codex_app_server_protocol::HooksListParams; use codex_app_server_protocol::HooksListResponse; +use codex_app_server_protocol::ImageGenerationContent; use codex_app_server_protocol::InitializeParams; use codex_app_server_protocol::InitializeResponse; use codex_app_server_protocol::JSONRPCErrorError; +use codex_app_server_protocol::LargeContentMode; use codex_app_server_protocol::ListMcpServerStatusParams; use codex_app_server_protocol::ListMcpServerStatusResponse; use codex_app_server_protocol::LoginAccountParams; @@ -173,6 +175,8 @@ use codex_app_server_protocol::ThreadIncrementElicitationResponse; use codex_app_server_protocol::ThreadInjectItemsParams; use codex_app_server_protocol::ThreadInjectItemsResponse; use codex_app_server_protocol::ThreadItem; +use codex_app_server_protocol::ThreadItemContentReadParams; +use codex_app_server_protocol::ThreadItemContentReadResponse; use codex_app_server_protocol::ThreadListCwdFilter; use codex_app_server_protocol::ThreadListParams; use codex_app_server_protocol::ThreadListResponse; diff --git a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs index 4a677d91ab..d542ad6ea0 100644 --- a/codex-rs/app-server/src/request_processors/thread_lifecycle.rs +++ b/codex-rs/app-server/src/request_processors/thread_lifecycle.rs @@ -538,6 +538,7 @@ pub(super) async fn handle_pending_thread_resume_request( active_turn.as_ref(), ); } + super::thread_processor::apply_large_content_mode_to_thread(&mut thread, pending.large_content); let thread_status = thread_watch_manager .loaded_status_for_thread(&thread.id) diff --git a/codex-rs/app-server/src/request_processors/thread_processor.rs b/codex-rs/app-server/src/request_processors/thread_processor.rs index 3dd3d3426d..779acfa245 100644 --- a/codex-rs/app-server/src/request_processors/thread_processor.rs +++ b/codex-rs/app-server/src/request_processors/thread_processor.rs @@ -546,6 +546,15 @@ impl ThreadRequestProcessor { .map(|response| Some(response.into())) } + pub(crate) async fn thread_item_content_read( + &self, + params: ThreadItemContentReadParams, + ) -> Result, JSONRPCErrorError> { + self.thread_item_content_read_response_inner(params) + .await + .map(|response| Some(response.into())) + } + pub(crate) async fn thread_shell_command( &self, request_id: &ConnectionRequestId, @@ -2034,6 +2043,7 @@ impl ThreadRequestProcessor { cursor, limit, sort_direction, + large_content, } = params; let thread_uuid = ThreadId::from_string(&thread_id) @@ -2048,22 +2058,8 @@ impl ThreadRequestProcessor { // every request. Rollback and compaction events can change earlier turns, so // the server has to rebuild the full turn list until turn metadata is indexed // separately. - let loaded_thread = self.thread_manager.get_thread(thread_uuid).await.ok(); - let has_live_running_thread = match loaded_thread.as_ref() { - Some(thread) => matches!(thread.agent_status().await, AgentStatus::Running), - None => false, - }; - let active_turn = if loaded_thread.is_some() { - // Persisted history may not yet include the currently running turn. The - // app-server listener has already projected live turn events into ThreadState, - // so merge that in-memory snapshot before paginating. - let thread_state = self.thread_state_manager.thread_state(thread_uuid).await; - let state = thread_state.lock().await; - state.active_turn_snapshot() - } else { - None - }; - let turns = reconstruct_thread_turns_for_turns_list( + let (has_live_running_thread, active_turn) = self.live_turn_read_context(thread_uuid).await; + let mut turns = reconstruct_thread_turns_for_turns_list( &items, self.thread_watch_manager .loaded_status_for_thread(&thread_uuid.to_string()) @@ -2071,6 +2067,7 @@ impl ThreadRequestProcessor { has_live_running_thread, active_turn, ); + apply_large_content_mode_to_turns(&mut turns, large_content.unwrap_or_default()); let page = paginate_thread_turns( turns, cursor.as_deref(), @@ -2094,6 +2091,7 @@ impl ThreadRequestProcessor { cursor, limit, sort_direction, + large_content, } = params; let thread_uuid = ThreadId::from_string(&thread_id) @@ -2102,18 +2100,7 @@ impl ThreadRequestProcessor { .load_thread_turns_list_history(thread_uuid) .await .map_err(thread_read_view_error)?; - let loaded_thread = self.thread_manager.get_thread(thread_uuid).await.ok(); - let has_live_running_thread = match loaded_thread.as_ref() { - Some(thread) => matches!(thread.agent_status().await, AgentStatus::Running), - None => false, - }; - let active_turn = if loaded_thread.is_some() { - let thread_state = self.thread_state_manager.thread_state(thread_uuid).await; - let state = thread_state.lock().await; - state.active_turn_snapshot() - } else { - None - }; + let (has_live_running_thread, active_turn) = self.live_turn_read_context(thread_uuid).await; let turns = reconstruct_thread_turns_for_turns_list( &items, self.thread_watch_manager @@ -2122,11 +2109,12 @@ impl ThreadRequestProcessor { has_live_running_thread, active_turn, ); - let turn_items = turns + let mut turn_items = turns .into_iter() .find(|turn| turn.id == turn_id) .ok_or_else(|| invalid_request(format!("turn not found: {turn_id}")))? .items; + apply_large_content_mode_to_items(&mut turn_items, large_content.unwrap_or_default()); let page = paginate_thread_items( turn_items, cursor.as_deref(), @@ -2140,6 +2128,82 @@ impl ThreadRequestProcessor { }) } + async fn thread_item_content_read_response_inner( + &self, + params: ThreadItemContentReadParams, + ) -> Result { + let ThreadItemContentReadParams { + thread_id, + turn_id, + item_id, + content_id, + } = params; + + let thread_uuid = ThreadId::from_string(&thread_id) + .map_err(|err| invalid_request(format!("invalid thread id: {err}")))?; + let items = self + .load_thread_turns_list_history(thread_uuid) + .await + .map_err(thread_read_view_error)?; + let (has_live_running_thread, active_turn) = self.live_turn_read_context(thread_uuid).await; + let turns = reconstruct_thread_turns_for_turns_list( + &items, + self.thread_watch_manager + .loaded_status_for_thread(&thread_uuid.to_string()) + .await, + has_live_running_thread, + active_turn, + ); + let turn = turns + .into_iter() + .find(|turn| turn.id == turn_id) + .ok_or_else(|| invalid_request(format!("turn not found: {turn_id}")))?; + let item = turn + .items + .into_iter() + .find(|item| item.id() == item_id) + .ok_or_else(|| invalid_request(format!("item not found: {item_id}")))?; + let ThreadItem::ImageGeneration { + content, result, .. + } = item + else { + return Err(invalid_request(format!( + "item {item_id} does not contain deferred content" + ))); + }; + if content_id != IMAGE_GENERATION_RESULT_CONTENT_ID { + return Err(invalid_request(format!("content not found: {content_id}"))); + } + let ImageGenerationContent::Inline { + mime_type, + byte_length, + .. + } = content + else { + return Err(internal_error(format!( + "image generation item {item_id} was unexpectedly already deferred" + ))); + }; + Ok(ThreadItemContentReadResponse { + mime_type, + data_base64: result, + byte_length, + }) + } + + async fn live_turn_read_context(&self, thread_id: ThreadId) -> (bool, Option) { + let Ok(thread) = self.thread_manager.get_thread(thread_id).await else { + return (false, None); + }; + let has_live_running_thread = matches!(thread.agent_status().await, AgentStatus::Running); + // Persisted history may not yet include the currently running turn. The + // app-server listener has already projected live turn events into ThreadState, + // so merge that in-memory snapshot when serving history reads. + let thread_state = self.thread_state_manager.thread_state(thread_id).await; + let state = thread_state.lock().await; + (has_live_running_thread, state.active_turn_snapshot()) + } + async fn load_thread_turns_list_history( &self, thread_id: ThreadId, @@ -2334,6 +2398,7 @@ impl ThreadRequestProcessor { developer_instructions, personality, exclude_turns, + large_content, persist_extended_history: _persist_extended_history, } = params; let include_turns = !exclude_turns; @@ -2458,6 +2523,7 @@ impl ThreadRequestProcessor { return Ok(()); } }; + apply_large_content_mode_to_thread(&mut thread, large_content.unwrap_or_default()); self.thread_watch_manager .upsert_thread(thread.clone()) @@ -2691,6 +2757,7 @@ impl ThreadRequestProcessor { emit_thread_goal_update, thread_goal_state_db, include_turns: !params.exclude_turns, + large_content: params.large_content.unwrap_or_default(), }), ); if listener_command_tx.send(command).is_err() { @@ -3325,6 +3392,7 @@ const THREAD_TURNS_DEFAULT_LIMIT: usize = 25; const THREAD_TURNS_MAX_LIMIT: usize = 100; const THREAD_ITEMS_DEFAULT_LIMIT: usize = 100; const THREAD_ITEMS_MAX_LIMIT: usize = 500; +const IMAGE_GENERATION_RESULT_CONTENT_ID: &str = "result"; fn thread_backwards_cursor_for_sort_key( thread: &StoredThread, @@ -3488,11 +3556,9 @@ fn paginate_thread_items( .as_ref() .and_then(|anchor| items.iter().position(|item| item.id() == anchor.item_id)); if anchor.is_some() && anchor_index.is_none() { - return Err(JSONRPCErrorError { - code: INVALID_REQUEST_ERROR_CODE, - message: "invalid cursor: anchor item is no longer present".to_string(), - data: None, - }); + return Err(invalid_request( + "invalid cursor: anchor item is no longer present", + )); } let mut keyed_items: Vec<_> = items.into_iter().enumerate().collect(); @@ -3555,19 +3621,50 @@ fn serialize_thread_items_cursor( item_id: item_id.to_string(), include_anchor, }) - .map_err(|err| JSONRPCErrorError { - code: INTERNAL_ERROR_CODE, - message: format!("failed to serialize cursor: {err}"), - data: None, - }) + .map_err(|err| internal_error(format!("failed to serialize cursor: {err}"))) } fn parse_thread_items_cursor(cursor: &str) -> Result { - serde_json::from_str(cursor).map_err(|_| JSONRPCErrorError { - code: INVALID_REQUEST_ERROR_CODE, - message: format!("invalid cursor: {cursor}"), - data: None, - }) + serde_json::from_str(cursor).map_err(|_| invalid_request(format!("invalid cursor: {cursor}"))) +} + +pub(super) fn apply_large_content_mode_to_thread(thread: &mut Thread, mode: LargeContentMode) { + apply_large_content_mode_to_turns(&mut thread.turns, mode); +} + +fn apply_large_content_mode_to_turns(turns: &mut [Turn], mode: LargeContentMode) { + for turn in turns { + apply_large_content_mode_to_items(&mut turn.items, mode); + } +} + +fn apply_large_content_mode_to_items(items: &mut [ThreadItem], mode: LargeContentMode) { + if matches!(mode, LargeContentMode::Inline) { + return; + } + + for item in items { + if let ThreadItem::ImageGeneration { + content, result, .. + } = item + && let ImageGenerationContent::Inline { + mime_type, + byte_length, + width, + height, + .. + } = content + { + *content = ImageGenerationContent::Deferred { + content_id: IMAGE_GENERATION_RESULT_CONTENT_ID.to_string(), + mime_type: mime_type.clone(), + byte_length: *byte_length, + width: *width, + height: *height, + }; + result.clear(); + } + } } fn reconstruct_thread_turns_for_turns_list( diff --git a/codex-rs/app-server/src/request_processors/thread_processor_tests.rs b/codex-rs/app-server/src/request_processors/thread_processor_tests.rs index 4f3e476e48..b713b987cc 100644 --- a/codex-rs/app-server/src/request_processors/thread_processor_tests.rs +++ b/codex-rs/app-server/src/request_processors/thread_processor_tests.rs @@ -525,6 +525,7 @@ mod thread_processor_behavior_tests { developer_instructions: None, personality: None, exclude_turns: false, + large_content: None, persist_extended_history: false, }; let config_snapshot = ThreadConfigSnapshot { diff --git a/codex-rs/app-server/src/thread_state.rs b/codex-rs/app-server/src/thread_state.rs index dddbcf483b..b4dac56939 100644 --- a/codex-rs/app-server/src/thread_state.rs +++ b/codex-rs/app-server/src/thread_state.rs @@ -1,5 +1,6 @@ use crate::outgoing_message::ConnectionId; use crate::outgoing_message::ConnectionRequestId; +use codex_app_server_protocol::LargeContentMode; use codex_app_server_protocol::RequestId; use codex_app_server_protocol::ThreadGoal; use codex_app_server_protocol::ThreadHistoryBuilder; @@ -33,6 +34,7 @@ pub(crate) struct PendingThreadResumeRequest { pub(crate) emit_thread_goal_update: bool, pub(crate) thread_goal_state_db: Option, pub(crate) include_turns: bool, + pub(crate) large_content: LargeContentMode, } // ThreadListenerCommand is used to perform operations in the context of the thread listener, for serialization purposes. diff --git a/codex-rs/app-server/tests/common/mcp_process.rs b/codex-rs/app-server/tests/common/mcp_process.rs index 3d508b618f..4a9877af2a 100644 --- a/codex-rs/app-server/tests/common/mcp_process.rs +++ b/codex-rs/app-server/tests/common/mcp_process.rs @@ -74,7 +74,7 @@ use codex_app_server_protocol::ThreadArchiveParams; use codex_app_server_protocol::ThreadCompactStartParams; use codex_app_server_protocol::ThreadForkParams; use codex_app_server_protocol::ThreadInjectItemsParams; -use codex_app_server_protocol::ThreadTurnsItemsListParams; +use codex_app_server_protocol::ThreadItemContentReadParams; use codex_app_server_protocol::ThreadListParams; use codex_app_server_protocol::ThreadLoadedListParams; use codex_app_server_protocol::ThreadMemoryModeSetParams; @@ -90,6 +90,7 @@ use codex_app_server_protocol::ThreadRollbackParams; use codex_app_server_protocol::ThreadSetNameParams; use codex_app_server_protocol::ThreadShellCommandParams; use codex_app_server_protocol::ThreadStartParams; +use codex_app_server_protocol::ThreadTurnsItemsListParams; use codex_app_server_protocol::ThreadTurnsListParams; use codex_app_server_protocol::ThreadUnarchiveParams; use codex_app_server_protocol::ThreadUnsubscribeParams; @@ -532,6 +533,15 @@ impl McpProcess { self.send_request("thread/turns/items/list", params).await } + /// Send a `thread/item/content/read` JSON-RPC request. + pub async fn send_thread_item_content_read_request( + &mut self, + params: ThreadItemContentReadParams, + ) -> anyhow::Result { + let params = Some(serde_json::to_value(params)?); + self.send_request("thread/item/content/read", params).await + } + /// Send a `model/list` JSON-RPC request. pub async fn send_list_models_request( &mut self, diff --git a/codex-rs/app-server/tests/suite/v2/thread_read.rs b/codex-rs/app-server/tests/suite/v2/thread_read.rs index 0dc616dc86..5f0ab9e95f 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_read.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_read.rs @@ -9,16 +9,20 @@ use codex_app_server::in_process; use codex_app_server::in_process::InProcessStartArgs; use codex_app_server_protocol::ClientInfo; use codex_app_server_protocol::ClientRequest; +use codex_app_server_protocol::ImageGenerationContent; use codex_app_server_protocol::InitializeCapabilities; use codex_app_server_protocol::InitializeParams; use codex_app_server_protocol::JSONRPCError; use codex_app_server_protocol::JSONRPCResponse; +use codex_app_server_protocol::LargeContentMode; use codex_app_server_protocol::RequestId; use codex_app_server_protocol::SessionSource; use codex_app_server_protocol::SortDirection; use codex_app_server_protocol::ThreadForkParams; use codex_app_server_protocol::ThreadForkResponse; use codex_app_server_protocol::ThreadItem; +use codex_app_server_protocol::ThreadItemContentReadParams; +use codex_app_server_protocol::ThreadItemContentReadResponse; use codex_app_server_protocol::ThreadListParams; use codex_app_server_protocol::ThreadListResponse; use codex_app_server_protocol::ThreadNameUpdatedNotification; @@ -31,6 +35,8 @@ use codex_app_server_protocol::ThreadSetNameResponse; use codex_app_server_protocol::ThreadStartParams; use codex_app_server_protocol::ThreadStartResponse; use codex_app_server_protocol::ThreadStatus; +use codex_app_server_protocol::ThreadTurnsItemsListParams; +use codex_app_server_protocol::ThreadTurnsItemsListResponse; use codex_app_server_protocol::ThreadTurnsListParams; use codex_app_server_protocol::ThreadTurnsListResponse; use codex_app_server_protocol::TurnItemsView; @@ -223,6 +229,7 @@ async fn thread_turns_list_can_page_backward_and_forward() -> Result<()> { cursor: None, limit: Some(2), sort_direction: Some(SortDirection::Desc), + large_content: None, }) .await?; let read_resp: JSONRPCResponse = timeout( @@ -249,6 +256,7 @@ async fn thread_turns_list_can_page_backward_and_forward() -> Result<()> { cursor: Some(next_cursor), limit: Some(10), sort_direction: Some(SortDirection::Desc), + large_content: None, }) .await?; let read_resp: JSONRPCResponse = timeout( @@ -267,6 +275,7 @@ async fn thread_turns_list_can_page_backward_and_forward() -> Result<()> { cursor: Some(backwards_cursor), limit: Some(10), sort_direction: Some(SortDirection::Asc), + large_content: None, }) .await?; let read_resp: JSONRPCResponse = timeout( @@ -280,6 +289,119 @@ async fn thread_turns_list_can_page_backward_and_forward() -> Result<()> { Ok(()) } +#[tokio::test] +async fn thread_turn_items_list_defers_image_generation_content_and_reads_it_back() -> Result<()> { + let server = create_mock_responses_server_repeating_assistant("Done").await; + let codex_home = TempDir::new()?; + create_config_toml(codex_home.path(), &server.uri())?; + + let filename_ts = "2025-01-05T12-00-00"; + let conversation_id = create_fake_rollout_with_text_elements( + codex_home.path(), + filename_ts, + "2025-01-05T12:00:00Z", + "make an image", + vec![], + Some("mock_provider"), + /*git_info*/ None, + )?; + let rollout_path = rollout_path(codex_home.path(), filename_ts, &conversation_id); + append_image_generation_end( + rollout_path.as_path(), + "2025-01-05T12:00:01Z", + "ig_123", + "Zm9v", + )?; + + let mut mcp = McpProcess::new(codex_home.path()).await?; + timeout(DEFAULT_READ_TIMEOUT, mcp.initialize()).await??; + + let turns_list_id = mcp + .send_thread_turns_list_request(ThreadTurnsListParams { + thread_id: conversation_id.clone(), + cursor: None, + limit: Some(10), + sort_direction: Some(SortDirection::Asc), + large_content: Some(LargeContentMode::Deferred), + }) + .await?; + let turns_list_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(turns_list_id)), + ) + .await??; + let ThreadTurnsListResponse { data, .. } = + to_response::(turns_list_resp)?; + let turn = data.first().expect("expected one turn"); + let turn_id = turn.id.clone(); + let deferred_item = turn + .items + .iter() + .find(|item| item.id() == "ig_123") + .expect("expected deferred image generation item"); + assert_eq!( + deferred_item, + &ThreadItem::ImageGeneration { + id: "ig_123".to_string(), + status: "completed".to_string(), + revised_prompt: None, + content: ImageGenerationContent::Deferred { + content_id: "result".to_string(), + mime_type: "image/png".to_string(), + byte_length: 3, + width: None, + height: None, + }, + result: String::new(), + saved_path: None, + } + ); + + let items_list_id = mcp + .send_thread_turns_items_list_request(ThreadTurnsItemsListParams { + thread_id: conversation_id.clone(), + turn_id: turn_id.clone(), + cursor: None, + limit: Some(10), + sort_direction: Some(SortDirection::Asc), + large_content: Some(LargeContentMode::Deferred), + }) + .await?; + let items_list_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(items_list_id)), + ) + .await??; + let ThreadTurnsItemsListResponse { data, .. } = + to_response::(items_list_resp)?; + assert!(data.contains(deferred_item)); + + let content_read_id = mcp + .send_thread_item_content_read_request(ThreadItemContentReadParams { + thread_id: conversation_id, + turn_id, + item_id: "ig_123".to_string(), + content_id: "result".to_string(), + }) + .await?; + let content_read_resp: JSONRPCResponse = timeout( + DEFAULT_READ_TIMEOUT, + mcp.read_stream_until_response_message(RequestId::Integer(content_read_id)), + ) + .await??; + let content = to_response::(content_read_resp)?; + assert_eq!( + content, + ThreadItemContentReadResponse { + mime_type: "image/png".to_string(), + data_base64: "Zm9v".to_string(), + byte_length: 3, + } + ); + + Ok(()) +} + #[tokio::test] async fn thread_turns_list_reads_store_history_without_rollout_path() -> Result<()> { let codex_home = TempDir::new()?; @@ -334,6 +456,7 @@ async fn thread_turns_list_reads_store_history_without_rollout_path() -> Result< cursor: None, limit: Some(10), sort_direction: Some(SortDirection::Asc), + large_content: None, }, }) .await? @@ -583,6 +706,7 @@ async fn thread_turns_list_rejects_cursor_when_anchor_turn_is_rolled_back() -> R cursor: None, limit: Some(2), sort_direction: Some(SortDirection::Desc), + large_content: None, }) .await?; let read_resp: JSONRPCResponse = timeout( @@ -607,6 +731,7 @@ async fn thread_turns_list_rejects_cursor_when_anchor_turn_is_rolled_back() -> R cursor: Some(backwards_cursor), limit: Some(10), sort_direction: Some(SortDirection::Asc), + large_content: None, }) .await?; let read_err: JSONRPCError = timeout( @@ -963,6 +1088,7 @@ async fn thread_turns_list_rejects_unmaterialized_loaded_thread() -> Result<()> cursor: None, limit: None, sort_direction: None, + large_content: None, }) .await?; let read_err: JSONRPCError = timeout( @@ -1068,6 +1194,29 @@ fn append_user_message(path: &Path, timestamp: &str, text: &str) -> std::io::Res ) } +fn append_image_generation_end( + path: &Path, + timestamp: &str, + call_id: &str, + result: &str, +) -> std::io::Result<()> { + let mut file = std::fs::OpenOptions::new().append(true).open(path)?; + writeln!( + file, + "{}", + json!({ + "timestamp": timestamp, + "type":"event_msg", + "payload": { + "type":"image_generation_end", + "call_id": call_id, + "status": "completed", + "result": result + } + }) + ) +} + fn append_thread_rollback(path: &Path, timestamp: &str, num_turns: u32) -> std::io::Result<()> { let mut file = std::fs::OpenOptions::new().append(true).open(path)?; writeln!( diff --git a/codex-rs/app-server/tests/suite/v2/thread_shell_command.rs b/codex-rs/app-server/tests/suite/v2/thread_shell_command.rs index eebc4077df..27830a4b28 100644 --- a/codex-rs/app-server/tests/suite/v2/thread_shell_command.rs +++ b/codex-rs/app-server/tests/suite/v2/thread_shell_command.rs @@ -151,6 +151,7 @@ async fn thread_shell_command_history_responses_exclude_persisted_command_execut cursor: None, limit: None, sort_direction: Some(SortDirection::Asc), + large_content: None, }) .await?; let turns_list_resp: JSONRPCResponse = timeout(