diff --git a/MODULE.bazel.lock b/MODULE.bazel.lock index 76850d7bd3..ea7a984f3a 100644 --- a/MODULE.bazel.lock +++ b/MODULE.bazel.lock @@ -1119,6 +1119,7 @@ "group_0.13.0": "{\"dependencies\":[{\"default_features\":false,\"name\":\"ff\",\"req\":\"^0.13\"},{\"name\":\"memuse\",\"optional\":true,\"req\":\"^0.2\"},{\"default_features\":false,\"name\":\"rand\",\"optional\":true,\"req\":\"^0.8\"},{\"default_features\":false,\"name\":\"rand_core\",\"req\":\"^0.6\"},{\"name\":\"rand_xorshift\",\"optional\":true,\"req\":\"^0.3\"},{\"default_features\":false,\"name\":\"subtle\",\"req\":\"^2.2.1\"}],\"features\":{\"alloc\":[],\"default\":[\"alloc\"],\"tests\":[\"alloc\",\"rand\",\"rand_xorshift\"],\"wnaf-memuse\":[\"alloc\",\"memuse\"]}}", "gzip-header_1.0.0": "{\"dependencies\":[{\"name\":\"crc32fast\",\"req\":\"^1.2.1\"}],\"features\":{}}", "h2_0.4.13": "{\"dependencies\":[{\"name\":\"atomic-waker\",\"req\":\"^1.0.0\"},{\"name\":\"bytes\",\"req\":\"^1\"},{\"default_features\":false,\"kind\":\"dev\",\"name\":\"env_logger\",\"req\":\"^0.10\"},{\"name\":\"fnv\",\"req\":\"^1.0.5\"},{\"default_features\":false,\"name\":\"futures-core\",\"req\":\"^0.3\"},{\"default_features\":false,\"name\":\"futures-sink\",\"req\":\"^0.3\"},{\"kind\":\"dev\",\"name\":\"hex\",\"req\":\"^0.4.3\"},{\"name\":\"http\",\"req\":\"^1\"},{\"features\":[\"std\"],\"name\":\"indexmap\",\"req\":\"^2\"},{\"default_features\":false,\"kind\":\"dev\",\"name\":\"quickcheck\",\"req\":\"^1.0.3\"},{\"kind\":\"dev\",\"name\":\"rand\",\"req\":\"^0.8.4\"},{\"kind\":\"dev\",\"name\":\"serde\",\"req\":\"^1.0.0\"},{\"kind\":\"dev\",\"name\":\"serde_json\",\"req\":\"^1.0.0\"},{\"name\":\"slab\",\"req\":\"^0.4.2\"},{\"features\":[\"io-util\"],\"name\":\"tokio\",\"req\":\"^1\"},{\"features\":[\"rt-multi-thread\",\"macros\",\"sync\",\"net\"],\"kind\":\"dev\",\"name\":\"tokio\",\"req\":\"^1\"},{\"kind\":\"dev\",\"name\":\"tokio-rustls\",\"req\":\"^0.26\"},{\"features\":[\"codec\",\"io\"],\"name\":\"tokio-util\",\"req\":\"^0.7.1\"},{\"default_features\":false,\"features\":[\"std\"],\"name\":\"tracing\",\"req\":\"^0.1.35\"},{\"kind\":\"dev\",\"name\":\"walkdir\",\"req\":\"^2.3.2\"},{\"kind\":\"dev\",\"name\":\"webpki-roots\",\"req\":\"^1\"}],\"features\":{\"stream\":[],\"unstable\":[]}}", + "h2_0.4.16": "{\"dependencies\":[{\"name\":\"atomic-waker\",\"req\":\"^1.0.0\"},{\"name\":\"bytes\",\"req\":\"^1\"},{\"default_features\":false,\"kind\":\"dev\",\"name\":\"env_logger\",\"req\":\"^0.10\"},{\"name\":\"fnv\",\"req\":\"^1.0.5\"},{\"default_features\":false,\"name\":\"futures-core\",\"req\":\"^0.3\"},{\"default_features\":false,\"name\":\"futures-sink\",\"req\":\"^0.3\"},{\"kind\":\"dev\",\"name\":\"hex\",\"req\":\"^0.4.3\"},{\"name\":\"http\",\"req\":\"^1.1\"},{\"features\":[\"std\"],\"name\":\"indexmap\",\"req\":\"^2\"},{\"default_features\":false,\"kind\":\"dev\",\"name\":\"quickcheck\",\"req\":\"^1.0.3\"},{\"kind\":\"dev\",\"name\":\"rand\",\"req\":\"^0.8.4\"},{\"kind\":\"dev\",\"name\":\"serde\",\"req\":\"^1.0.0\"},{\"kind\":\"dev\",\"name\":\"serde_json\",\"req\":\"^1.0.0\"},{\"name\":\"slab\",\"req\":\"^0.4.2\"},{\"features\":[\"io-util\"],\"name\":\"tokio\",\"req\":\"^1\"},{\"features\":[\"rt-multi-thread\",\"macros\",\"sync\",\"net\"],\"kind\":\"dev\",\"name\":\"tokio\",\"req\":\"^1\"},{\"kind\":\"dev\",\"name\":\"tokio-rustls\",\"req\":\"^0.26\"},{\"features\":[\"codec\",\"io\"],\"name\":\"tokio-util\",\"req\":\"^0.7.1\"},{\"default_features\":false,\"features\":[\"std\"],\"name\":\"tracing\",\"req\":\"^0.1.35\"},{\"kind\":\"dev\",\"name\":\"walkdir\",\"req\":\"^2.3.2\"},{\"kind\":\"dev\",\"name\":\"webpki-roots\",\"req\":\"^1\"}],\"features\":{\"stream\":[],\"unstable\":[]}}", "half_2.7.1": "{\"dependencies\":[{\"features\":[\"derive\"],\"name\":\"arbitrary\",\"optional\":true,\"req\":\"^1.4.1\"},{\"default_features\":false,\"features\":[\"derive\"],\"name\":\"bytemuck\",\"optional\":true,\"req\":\"^1.4.1\"},{\"name\":\"cfg-if\",\"req\":\"^1.0.0\"},{\"kind\":\"dev\",\"name\":\"criterion\",\"req\":\"^0.5\"},{\"name\":\"crunchy\",\"req\":\"^0.2.2\",\"target\":\"cfg(target_arch = \\\"spirv\\\")\"},{\"kind\":\"dev\",\"name\":\"crunchy\",\"req\":\"^0.2.2\"},{\"default_features\":false,\"features\":[\"libm\"],\"name\":\"num-traits\",\"optional\":true,\"req\":\"^0.2.16\"},{\"kind\":\"dev\",\"name\":\"quickcheck\",\"req\":\"^1.0\"},{\"kind\":\"dev\",\"name\":\"quickcheck_macros\",\"req\":\"^1.0\"},{\"default_features\":false,\"features\":[\"thread_rng\"],\"name\":\"rand\",\"optional\":true,\"req\":\"^0.9.0\"},{\"kind\":\"dev\",\"name\":\"rand\",\"req\":\"^0.9.0\"},{\"default_features\":false,\"name\":\"rand_distr\",\"optional\":true,\"req\":\"^0.5.0\"},{\"name\":\"rkyv\",\"optional\":true,\"req\":\"^0.8.0\"},{\"default_features\":false,\"features\":[\"derive\"],\"name\":\"serde\",\"optional\":true,\"req\":\"^1.0\"},{\"default_features\":false,\"features\":[\"derive\",\"simd\"],\"name\":\"zerocopy\",\"req\":\"^0.8.26\"}],\"features\":{\"alloc\":[],\"default\":[\"std\"],\"nightly\":[],\"rand_distr\":[\"dep:rand\",\"dep:rand_distr\"],\"std\":[\"alloc\"],\"use-intrinsics\":[],\"zerocopy\":[]}}", "hash32_0.2.1": "{\"dependencies\":[{\"default_features\":false,\"name\":\"byteorder\",\"req\":\"^1.2.2\"},{\"kind\":\"dev\",\"name\":\"hash32-derive\",\"req\":\"^0.1.0\"}],\"features\":{}}", "hash32_0.3.1": "{\"dependencies\":[{\"default_features\":false,\"name\":\"byteorder\",\"req\":\"^1.2.2\"}],\"features\":{}}", diff --git a/codex-rs/Cargo.lock b/codex-rs/Cargo.lock index 0b3b4f8dfd..eb794ef447 100644 --- a/codex-rs/Cargo.lock +++ b/codex-rs/Cargo.lock @@ -7687,9 +7687,9 @@ dependencies = [ [[package]] name = "h2" -version = "0.4.13" +version = "0.4.16" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "2f44da3a8150a6703ed5d34e164b875fd14c2cdab9af1252a9a1020bde2bdc54" +checksum = "a9f37a958b41b3b19ee2707c06439c0e9e547e847223eb791ecb0cb821c65e27" dependencies = [ "atomic-waker", "bytes", diff --git a/codex-rs/app-server-protocol/schema/json/ClientRequest.json b/codex-rs/app-server-protocol/schema/json/ClientRequest.json index d44523b425..830672e5ae 100644 --- a/codex-rs/app-server-protocol/schema/json/ClientRequest.json +++ b/codex-rs/app-server-protocol/schema/json/ClientRequest.json @@ -1912,6 +1912,13 @@ }, "McpResourceReadParams": { "properties": { + "originCallId": { + "description": "Originating MCP tool call used to select the resource's app.", + "type": [ + "string", + "null" + ] + }, "server": { "type": "string" }, diff --git a/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.schemas.json b/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.schemas.json index bb6b8f1755..cc99f737b8 100644 --- a/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.schemas.json +++ b/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.schemas.json @@ -13295,6 +13295,13 @@ "McpResourceReadParams": { "$schema": "http://json-schema.org/draft-07/schema#", "properties": { + "originCallId": { + "description": "Originating MCP tool call used to select the resource's app.", + "type": [ + "string", + "null" + ] + }, "server": { "type": "string" }, @@ -13323,6 +13330,13 @@ "$ref": "#/definitions/v2/ResourceContent" }, "type": "array" + }, + "originCallId": { + "description": "Originating call when the server applied app-specific resource scoping.", + "type": [ + "string", + "null" + ] } }, "required": [ diff --git a/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.v2.schemas.json b/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.v2.schemas.json index 59415b86bd..d2f9019e8b 100644 --- a/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.v2.schemas.json +++ b/codex-rs/app-server-protocol/schema/json/codex_app_server_protocol.v2.schemas.json @@ -9525,6 +9525,13 @@ "McpResourceReadParams": { "$schema": "http://json-schema.org/draft-07/schema#", "properties": { + "originCallId": { + "description": "Originating MCP tool call used to select the resource's app.", + "type": [ + "string", + "null" + ] + }, "server": { "type": "string" }, @@ -9553,6 +9560,13 @@ "$ref": "#/definitions/ResourceContent" }, "type": "array" + }, + "originCallId": { + "description": "Originating call when the server applied app-specific resource scoping.", + "type": [ + "string", + "null" + ] } }, "required": [ diff --git a/codex-rs/app-server-protocol/schema/json/v2/McpResourceReadParams.json b/codex-rs/app-server-protocol/schema/json/v2/McpResourceReadParams.json index 2fe58155d2..26435c73e6 100644 --- a/codex-rs/app-server-protocol/schema/json/v2/McpResourceReadParams.json +++ b/codex-rs/app-server-protocol/schema/json/v2/McpResourceReadParams.json @@ -1,6 +1,13 @@ { "$schema": "http://json-schema.org/draft-07/schema#", "properties": { + "originCallId": { + "description": "Originating MCP tool call used to select the resource's app.", + "type": [ + "string", + "null" + ] + }, "server": { "type": "string" }, diff --git a/codex-rs/app-server-protocol/schema/json/v2/McpResourceReadResponse.json b/codex-rs/app-server-protocol/schema/json/v2/McpResourceReadResponse.json index b1a4012344..99060ef873 100644 --- a/codex-rs/app-server-protocol/schema/json/v2/McpResourceReadResponse.json +++ b/codex-rs/app-server-protocol/schema/json/v2/McpResourceReadResponse.json @@ -59,6 +59,13 @@ "$ref": "#/definitions/ResourceContent" }, "type": "array" + }, + "originCallId": { + "description": "Originating call when the server applied app-specific resource scoping.", + "type": [ + "string", + "null" + ] } }, "required": [ diff --git a/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-experimental.json.zst b/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-experimental.json.zst index 24ec8911c7..7175b91988 100644 Binary files a/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-experimental.json.zst and b/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-experimental.json.zst differ diff --git a/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst b/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst index 668a62f3af..0b1cefd57a 100644 Binary files a/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst and b/codex-rs/app-server-protocol/schema/precomputed/app-server-exports-stable.json.zst differ diff --git a/codex-rs/app-server-protocol/schema/typescript/v2/McpResourceReadParams.ts b/codex-rs/app-server-protocol/schema/typescript/v2/McpResourceReadParams.ts index c48795f27e..c09a3ab2a4 100644 --- a/codex-rs/app-server-protocol/schema/typescript/v2/McpResourceReadParams.ts +++ b/codex-rs/app-server-protocol/schema/typescript/v2/McpResourceReadParams.ts @@ -2,4 +2,8 @@ // This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. -export type McpResourceReadParams = { threadId?: string | null, server: string, uri: string, }; +export type McpResourceReadParams = { threadId?: string | null, +/** + * Originating MCP tool call used to select the resource's app. + */ +originCallId?: string | null, server: string, uri: string, }; diff --git a/codex-rs/app-server-protocol/schema/typescript/v2/McpResourceReadResponse.ts b/codex-rs/app-server-protocol/schema/typescript/v2/McpResourceReadResponse.ts index 2af1dbcd09..2b1194d150 100644 --- a/codex-rs/app-server-protocol/schema/typescript/v2/McpResourceReadResponse.ts +++ b/codex-rs/app-server-protocol/schema/typescript/v2/McpResourceReadResponse.ts @@ -3,4 +3,8 @@ // This file was generated by [ts-rs](https://github.com/Aleph-Alpha/ts-rs). Do not edit this file manually. import type { ResourceContent } from "../ResourceContent"; -export type McpResourceReadResponse = { contents: Array, }; +export type McpResourceReadResponse = { contents: Array, +/** + * Originating call when the server applied app-specific resource scoping. + */ +originCallId: string | null, }; diff --git a/codex-rs/app-server-protocol/src/protocol/common.rs b/codex-rs/app-server-protocol/src/protocol/common.rs index cbde00fce3..a73455111f 100644 --- a/codex-rs/app-server-protocol/src/protocol/common.rs +++ b/codex-rs/app-server-protocol/src/protocol/common.rs @@ -2328,6 +2328,7 @@ mod tests { request_id: request_id(), params: v2::McpResourceReadParams { thread_id: Some("thread-1".to_string()), + origin_call_id: None, server: "server-a".to_string(), uri: "file:///tmp/resource".to_string(), }, @@ -2515,6 +2516,7 @@ mod tests { request_id: request_id(), params: v2::McpResourceReadParams { thread_id: None, + origin_call_id: None, server: "server-a".to_string(), uri: "file:///tmp/resource".to_string(), }, diff --git a/codex-rs/app-server-protocol/src/protocol/v2/mcp.rs b/codex-rs/app-server-protocol/src/protocol/v2/mcp.rs index 2f806c0c28..2196a3b96a 100644 --- a/codex-rs/app-server-protocol/src/protocol/v2/mcp.rs +++ b/codex-rs/app-server-protocol/src/protocol/v2/mcp.rs @@ -85,6 +85,9 @@ pub struct ListMcpServerStatusResponse { pub struct McpResourceReadParams { #[ts(optional = nullable)] pub thread_id: Option, + /// Originating MCP tool call used to select the resource's app. + #[ts(optional = nullable)] + pub origin_call_id: Option, pub server: String, pub uri: String, } @@ -94,6 +97,8 @@ pub struct McpResourceReadParams { #[ts(export_to = "v2/")] pub struct McpResourceReadResponse { pub contents: Vec, + /// Originating call when the server applied app-specific resource scoping. + pub origin_call_id: Option, } #[derive(Serialize, Deserialize, Debug, Clone, PartialEq, JsonSchema, TS)] diff --git a/codex-rs/app-server/README.md b/codex-rs/app-server/README.md index 7b63b45475..40479fcd23 100644 --- a/codex-rs/app-server/README.md +++ b/codex-rs/app-server/README.md @@ -278,7 +278,7 @@ Example with notification opt-out: - `tool/requestUserInput` — prompt the user with 1–3 short questions for a tool call and return their answers (experimental). - `config/mcpServer/reload` — reload MCP server config from disk and queue a refresh for loaded threads (applied on each thread's next active turn); returns `{}`. Use this after editing `config.toml` without restarting the server. - `mcpServerStatus/list` — enumerate configured MCP servers with their tools, auth status, server info, owning `pluginId` (`null` for servers not contributed by a plugin), plus resources/resource templates for `full` detail; supports optional `threadId` and cursor+limit pagination. If `threadId` is omitted, the server reads from the latest global config directly. If `detail` is omitted, the server defaults to `full`. An `unknown` auth status means OAuth support could not be determined; `unsupported` means OAuth is known not to be supported. -- `mcpServer/resource/read` — read a resource from a configured MCP server by optional `threadId`, `server`, and `uri`, returning text/blob resource `contents`. If `threadId` is omitted, the server reads from the latest MCP config directly. +- `mcpServer/resource/read` — read a resource from a configured MCP server by optional `threadId`, `server`, and `uri`, returning text/blob resource `contents`. Pass `originCallId` with `threadId` to scope a Codex app widget to the app and account of the completed tool call that produced it; successful scoped reads return the same `originCallId`. If `threadId` is omitted, the server reads from the latest MCP config directly. - `mcpServer/tool/call` — call a tool on a thread's configured MCP server by `threadId`, `server`, `tool`, optional `arguments`, and optional `_meta`, returning the MCP tool result. - `windowsSandbox/setupStart` — start Windows sandbox setup for the selected mode (`elevated` or `unelevated`); accepts an optional absolute `cwd` to target setup for a specific workspace, returns `{ started: true }` immediately, and later emits `windowsSandbox/setupCompleted`. - `feedback/upload` — submit a feedback report (classification + optional reason/logs, conversation_id, and optional `extraLogFiles` attachments array); returns the tracking thread id. diff --git a/codex-rs/app-server/src/request_processors/mcp_processor.rs b/codex-rs/app-server/src/request_processors/mcp_processor.rs index 69c2c285ef..eea21cfd1d 100644 --- a/codex-rs/app-server/src/request_processors/mcp_processor.rs +++ b/codex-rs/app-server/src/request_processors/mcp_processor.rs @@ -422,6 +422,7 @@ impl McpRequestProcessor { let outgoing = Arc::clone(&self.outgoing); let McpResourceReadParams { thread_id, + origin_call_id, server, uri, } = params; @@ -431,12 +432,22 @@ impl McpRequestProcessor { let request_id = request_id.clone(); tokio::spawn(async move { - let result = thread.read_mcp_resource(&server, &uri).await; - Self::send_mcp_resource_read_response(outgoing, request_id, result).await; + let origin_call_id = + origin_call_id.filter(|_| server == codex_mcp::CODEX_APPS_MCP_SERVER_NAME); + let result = match origin_call_id.as_deref() { + Some(call_id) => thread.read_mcp_resource_for_call(call_id, &uri).await, + None => thread.read_mcp_resource(&server, &uri).await, + }; + Self::send_mcp_resource_read_response(outgoing, request_id, result, origin_call_id) + .await; }); return Ok(()); } + if origin_call_id.is_some() { + return Err(invalid_request("originCallId requires threadId")); + } + let config = self.load_latest_config(/*fallback_cwd*/ None).await?; let mcp_manager = self.thread_manager.mcp_manager(); let mcp_config = mcp_manager.runtime_config(&config).await; @@ -463,7 +474,10 @@ impl McpRequestProcessor { ) .await .and_then(|result| serde_json::to_value(result).map_err(anyhow::Error::from)); - Self::send_mcp_resource_read_response(outgoing, request_id, result).await; + Self::send_mcp_resource_read_response( + outgoing, request_id, result, /*origin_call_id*/ None, + ) + .await; }); Ok(()) } @@ -472,6 +486,7 @@ impl McpRequestProcessor { outgoing: Arc, request_id: ConnectionRequestId, result: anyhow::Result, + origin_call_id: Option, ) { let result = result .map_err(|error| internal_error(format!("{error:#}"))) @@ -481,6 +496,10 @@ impl McpRequestProcessor { "failed to deserialize MCP resource read response: {error}" )) }) + }) + .map(|mut response| { + response.origin_call_id = origin_call_id; + response }); outgoing.send_result(request_id, result).await; } diff --git a/codex-rs/app-server/tests/suite/v2/mcp_resource.rs b/codex-rs/app-server/tests/suite/v2/mcp_resource.rs index 7286ee5c55..2860ee994d 100644 --- a/codex-rs/app-server/tests/suite/v2/mcp_resource.rs +++ b/codex-rs/app-server/tests/suite/v2/mcp_resource.rs @@ -1,4 +1,5 @@ use std::sync::Arc; +use std::sync::atomic::AtomicBool; use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::time::Duration; @@ -36,11 +37,14 @@ use core_test_support::responses; use pretty_assertions::assert_eq; use rmcp::handler::server::ServerHandler; use rmcp::model::BooleanSchema; +use rmcp::model::CallToolRequestParams; +use rmcp::model::CallToolResult; use rmcp::model::ElicitRequestParams; use rmcp::model::ElicitResult; use rmcp::model::ElicitationAction; use rmcp::model::ElicitationSchema; use rmcp::model::ListResourcesResult; +use rmcp::model::ListToolsResult; use rmcp::model::MetaObject; use rmcp::model::PaginatedRequestParams; use rmcp::model::PrimitiveSchemaDefinition; @@ -62,8 +66,9 @@ use tokio::net::TcpListener; use tokio::task::JoinHandle; use tokio::time::timeout; -const DEFAULT_READ_TIMEOUT: Duration = Duration::from_secs(10); +pub(super) const DEFAULT_READ_TIMEOUT: Duration = Duration::from_secs(10); const TEST_RESOURCE_URI: &str = "test://codex/resource"; +pub(super) const TEST_WIDGET_RESOURCE_URI: &str = "ui://widget/checkout-session.html"; const TEST_BLOB_RESOURCE_URI: &str = "test://codex/resource.bin"; const TEST_RESOURCE_BLOB: &str = "YmluYXJ5LXJlc291cmNl"; const TEST_RESOURCE_TEXT: &str = "Resource body from the MCP server."; @@ -116,6 +121,7 @@ async fn mcp_resource_read_returns_resource_contents() -> Result<()> { request_id, params: McpResourceReadParams { thread_id: Some(thread.id), + origin_call_id: None, server: "codex_apps".to_string(), uri: TEST_RESOURCE_URI.to_string(), }, @@ -588,6 +594,7 @@ apps = true request_id, params: McpResourceReadParams { thread_id: None, + origin_call_id: None, server: "codex_apps".to_string(), uri: TEST_RESOURCE_URI.to_string(), }, @@ -599,6 +606,7 @@ apps = true request_id, params: McpResourceReadParams { thread_id: None, + origin_call_id: None, server: "codex_apps".to_string(), uri: TEST_ELICITATION_RESOURCE_URI.to_string(), }, @@ -613,6 +621,7 @@ apps = true text: TEST_ELICITATION_RESOURCE_TEXT.to_string(), meta: None, }], + origin_call_id: None, } ); @@ -665,6 +674,7 @@ async fn mcp_resource_read_returns_error_for_unknown_thread() -> Result<()> { request_id: RequestId::Integer(1), params: McpResourceReadParams { thread_id: Some("00000000-0000-4000-8000-000000000000".to_string()), + origin_call_id: None, server: "codex_apps".to_string(), uri: TEST_RESOURCE_URI.to_string(), }, @@ -684,7 +694,7 @@ async fn mcp_resource_read_returns_error_for_unknown_thread() -> Result<()> { Ok(()) } -async fn start_resource_test_app_server( +pub(super) async fn start_resource_test_app_server( apps_server_url: &str, responses_server_uri: &str, environment: ResourceTestEnvironment, @@ -734,12 +744,12 @@ async fn start_resource_test_app_server_with_extra_config( Ok((codex_home, mcp)) } -enum ResourceTestEnvironment { +pub(super) enum ResourceTestEnvironment { Auto, Local, } -async fn start_resource_apps_mcp_server() +pub(super) async fn start_resource_apps_mcp_server() -> Result<(String, Arc, JoinHandle<()>)> { let listener = TcpListener::bind("127.0.0.1:0").await?; let addr = listener.local_addr()?; @@ -780,14 +790,17 @@ fn expected_resource_read_response() -> McpResourceReadResponse { meta: None, }, ], + origin_call_id: None, } } #[derive(Debug, Default)] -struct ResourceAppsMcpCalls { +pub(super) struct ResourceAppsMcpCalls { list_resources: AtomicUsize, main_prompt_reads: AtomicUsize, reference_reads: AtomicUsize, + pub(super) tools_enabled: AtomicBool, + pub(super) best_buy_app_only: AtomicBool, } impl ResourceAppsMcpCalls { @@ -814,8 +827,70 @@ struct ResourceAppsMcpServer { impl ServerHandler for ResourceAppsMcpServer { fn get_info(&self) -> ServerInfo { - ServerInfo::new(ServerCapabilities::builder().enable_resources().build()) - .with_protocol_version(ProtocolVersion::V_2025_06_18) + let mut info = ServerInfo::new(ServerCapabilities::builder().enable_resources().build()) + .with_protocol_version(ProtocolVersion::V_2025_06_18); + if self.calls.tools_enabled.load(Ordering::Relaxed) { + info.capabilities.tools = Some(Default::default()); + } + info + } + + async fn list_tools( + &self, + _request: Option, + _context: RequestContext, + ) -> Result { + let tools = ["best_buy", "walmart"] + .into_iter() + .map(|connector_id| { + let mut ui = json!({ "resourceUri": TEST_WIDGET_RESOURCE_URI }); + if connector_id == "best_buy" + && self.calls.best_buy_app_only.load(Ordering::Relaxed) + { + ui["visibility"] = json!(["app"]); + } + serde_json::from_value(json!({ + "name": format!("{connector_id}_product_search"), + "description": "Search products.", + "inputSchema": { "type": "object" }, + "annotations": { "readOnlyHint": true }, + "_meta": { + "connector_id": connector_id, + "connector_name": connector_id, + "link_id": format!("link_{connector_id}"), + "ui": ui, + "openai/outputTemplate": TEST_WIDGET_RESOURCE_URI, + "_codex_apps": { + "resource_uri": format!( + "/{connector_id}/link_{connector_id}/{connector_id}_product_search" + ), + "contains_mcp_source": true, + }, + }, + })) + }) + .collect::>>() + .map_err(|error| rmcp::ErrorData::internal_error(error.to_string(), None))?; + Ok(ListToolsResult::with_all_items(tools)) + } + + async fn call_tool( + &self, + request: CallToolRequestParams, + _context: RequestContext, + ) -> Result { + if request + .arguments + .as_ref() + .and_then(|arguments| arguments.get("query")) + == Some(&json!("fail")) + { + return Ok( + CallToolResult::structured_error(json!({ "error": "search failed" })).into(), + ); + } + + Ok(CallToolResult::structured(json!({ "products": [] })).into()) } async fn list_resources( @@ -868,6 +943,36 @@ impl ServerHandler for ResourceAppsMcpServer { context: RequestContext, ) -> Result { let uri = request.uri; + if uri == TEST_WIDGET_RESOURCE_URI { + let request_meta = context + .meta + .0 + .0 + .get("x-codex-turn-metadata") + .and_then(|metadata| metadata.get("mcp_request_meta")); + let connector_id = request_meta + .and_then(|metadata| metadata.pointer("/selected_connector_ids/0")) + .and_then(serde_json::Value::as_str) + .ok_or_else(|| rmcp::ErrorData::invalid_params("missing app scope", None))?; + let expected_link_id = format!("link_{connector_id}"); + if request_meta + .and_then(|metadata| metadata.get("link_id")) + .and_then(serde_json::Value::as_str) + != Some(expected_link_id.as_str()) + { + return Err(rmcp::ErrorData::invalid_params("wrong account scope", None)); + } + + return Ok( + ReadResourceResult::new(vec![ResourceContents::TextResourceContents { + uri, + mime_type: Some("text/html".to_string()), + text: format!("{connector_id}"), + meta: None, + }]) + .into(), + ); + } if uri == TEST_ELICITATION_RESOURCE_URI { let requested_schema = ElicitationSchema::builder() .required_property( diff --git a/codex-rs/app-server/tests/suite/v2/mcp_resource_origin.rs b/codex-rs/app-server/tests/suite/v2/mcp_resource_origin.rs new file mode 100644 index 0000000000..0168b97b49 --- /dev/null +++ b/codex-rs/app-server/tests/suite/v2/mcp_resource_origin.rs @@ -0,0 +1,241 @@ +use super::mcp_resource::DEFAULT_READ_TIMEOUT; +use super::mcp_resource::ResourceTestEnvironment; +use super::mcp_resource::TEST_WIDGET_RESOURCE_URI; +use super::mcp_resource::start_resource_apps_mcp_server; +use super::mcp_resource::start_resource_test_app_server; +use anyhow::Result; +use app_test_support::TestAppServer; +use codex_app_server_protocol::ClientRequest; +use codex_app_server_protocol::McpResourceContent; +use codex_app_server_protocol::McpResourceReadParams; +use codex_app_server_protocol::McpResourceReadResponse; +use codex_app_server_protocol::RequestId; +use codex_app_server_protocol::ThreadHistoryMode; +use codex_app_server_protocol::ThreadResumeParams; +use codex_app_server_protocol::ThreadResumeResponse; +use codex_app_server_protocol::ThreadStartParams; +use codex_app_server_protocol::ThreadStartResponse; +use codex_app_server_protocol::TurnStartParams; +use codex_app_server_protocol::TurnStartResponse; +use codex_app_server_protocol::UserInput; +use core_test_support::responses; +use pretty_assertions::assert_eq; +use serde_json::json; +use std::sync::atomic::Ordering; +use tokio::time::timeout; + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn widget_reads_survive_history_modes_restarts_and_app_only_visibility() -> Result<()> { + let responses_server = responses::start_mock_server().await; + let (apps_server_url, calls, apps_server_handle) = start_resource_apps_mcp_server().await?; + calls.tools_enabled.store(true, Ordering::Relaxed); + let (codex_home, mut app_server) = start_resource_test_app_server( + &apps_server_url, + &responses_server.uri(), + ResourceTestEnvironment::Auto, + ) + .await?; + + let mut tool_events = vec![responses::ev_response_created("widget-tools")]; + tool_events.extend( + [ + ("best-buy-call", "best_buy", "lamps", "link_best_buy"), + ("walmart-call", "walmart", "lamps", "link_walmart"), + ("failed-call", "walmart", "fail", "link_walmart"), + ("ambiguous-account-call", "walmart", "lamps", "link_other"), + ] + .into_iter() + .map(|(call_id, app, query, link_id)| { + responses::ev_function_call_with_namespace( + call_id, + &format!("mcp__codex_apps__{app}"), + "_product_search", + &json!({ "query": query, "link_id": link_id }).to_string(), + ) + }), + ); + tool_events.push(responses::ev_completed("widget-tools")); + let response = [ + responses::sse(tool_events), + responses::sse(vec![responses::ev_completed("widget-done")]), + ]; + let mut model_responses = std::iter::repeat_n(response, 4) + .flatten() + .collect::>(); + model_responses.push(responses::sse(vec![responses::ev_completed( + "app-only-visibility", + )])); + let response_mock = responses::mount_sse_sequence(&responses_server, model_responses).await; + + let mut persistent_thread_ids = Vec::new(); + for (history_mode, ephemeral) in [ + (ThreadHistoryMode::Legacy, false), + (ThreadHistoryMode::Paginated, false), + (ThreadHistoryMode::Legacy, true), + (ThreadHistoryMode::Paginated, true), + ] { + let ThreadStartResponse { thread, .. } = app_server + .start_thread(ThreadStartParams { + history_mode: Some(history_mode), + ephemeral: ephemeral.then_some(true), + ..Default::default() + }) + .await?; + let turn_id = app_server + .send_turn_start_request(TurnStartParams { + thread_id: thread.id.clone(), + input: vec![UserInput::Text { + text: "Find a lamp.".to_string(), + text_elements: Vec::new(), + }], + ..Default::default() + }) + .await?; + let _: TurnStartResponse = + timeout(DEFAULT_READ_TIMEOUT, app_server.read_response(turn_id)).await??; + timeout( + DEFAULT_READ_TIMEOUT, + app_server.read_stream_until_notification_message("turn/completed"), + ) + .await??; + + for (call_id, connector_id) in [("best-buy-call", "best_buy"), ("walmart-call", "walmart")] + { + let response = read_widget(&mut app_server, &thread.id, call_id).await?; + assert_eq!( + response, + McpResourceReadResponse { + contents: vec![McpResourceContent::Text { + uri: TEST_WIDGET_RESOURCE_URI.to_string(), + mime_type: Some("text/html".to_string()), + text: format!("{connector_id}"), + meta: None, + }], + origin_call_id: Some(call_id.to_string()), + } + ); + } + + for (thread_id, call_id, uri, expected_error) in [ + ( + Some(thread.id.clone()), + "failed-call", + TEST_WIDGET_RESOURCE_URI, + "was not found", + ), + ( + Some(thread.id.clone()), + "walmart-call", + "ui://widget/wrong.html", + "does not match", + ), + ( + Some(thread.id.clone()), + "ambiguous-account-call", + TEST_WIDGET_RESOURCE_URI, + "ambiguous account", + ), + ( + None, + "walmart-call", + TEST_WIDGET_RESOURCE_URI, + "requires threadId", + ), + ] { + let request_id = app_server + .send_mcp_resource_read_request(McpResourceReadParams { + thread_id, + origin_call_id: Some(call_id.to_string()), + server: "codex_apps".to_string(), + uri: uri.to_string(), + }) + .await?; + let error = timeout( + DEFAULT_READ_TIMEOUT, + app_server.read_stream_until_error_message(RequestId::Integer(request_id)), + ) + .await??; + assert!( + error.error.message.contains(expected_error), + "expected {expected_error:?}, got: {error:?}" + ); + } + if !ephemeral { + persistent_thread_ids.push(thread.id); + } + } + + timeout(DEFAULT_READ_TIMEOUT, app_server.shutdown_gracefully()).await??; + calls.best_buy_app_only.store(true, Ordering::Relaxed); + let mut restarted = TestAppServer::builder() + .with_codex_home(codex_home.path()) + .build_initialized() + .await?; + let model_visibility_thread_id = persistent_thread_ids[0].clone(); + for thread_id in persistent_thread_ids { + let resume_id = restarted + .send_thread_resume_request(ThreadResumeParams { + thread_id, + ..Default::default() + }) + .await?; + let ThreadResumeResponse { thread, .. } = + timeout(DEFAULT_READ_TIMEOUT, restarted.read_response(resume_id)).await??; + for call_id in ["walmart-call", "best-buy-call"] { + let response = read_widget(&mut restarted, &thread.id, call_id).await?; + assert_eq!(response.origin_call_id.as_deref(), Some(call_id)); + } + } + + let turn_id = restarted + .send_turn_start_request(TurnStartParams { + thread_id: model_visibility_thread_id, + input: vec![UserInput::Text { + text: "Find another lamp.".to_string(), + text_elements: Vec::new(), + }], + ..Default::default() + }) + .await?; + let _: TurnStartResponse = + timeout(DEFAULT_READ_TIMEOUT, restarted.read_response(turn_id)).await??; + timeout( + DEFAULT_READ_TIMEOUT, + restarted.read_stream_until_notification_message("turn/completed"), + ) + .await??; + let requests = response_mock.requests(); + let model_request = requests.last().expect("post-resume model request"); + assert!( + model_request + .tool_by_name("mcp__codex_apps__best_buy", "_product_search") + .is_none() + ); + assert!( + model_request + .tool_by_name("mcp__codex_apps__walmart", "_product_search") + .is_some() + ); + + apps_server_handle.abort(); + let _ = apps_server_handle.await; + Ok(()) +} + +async fn read_widget( + app_server: &mut TestAppServer, + thread_id: &str, + call_id: &str, +) -> Result { + app_server + .request(|request_id| ClientRequest::McpResourceRead { + request_id, + params: McpResourceReadParams { + thread_id: Some(thread_id.to_string()), + origin_call_id: Some(call_id.to_string()), + server: "codex_apps".to_string(), + uri: TEST_WIDGET_RESOURCE_URI.to_string(), + }, + }) + .await +} diff --git a/codex-rs/app-server/tests/suite/v2/mod.rs b/codex-rs/app-server/tests/suite/v2/mod.rs index 373eb271de..75a7952eec 100644 --- a/codex-rs/app-server/tests/suite/v2/mod.rs +++ b/codex-rs/app-server/tests/suite/v2/mod.rs @@ -42,6 +42,7 @@ mod marketplace_add; mod marketplace_remove; mod marketplace_upgrade; mod mcp_resource; +mod mcp_resource_origin; mod mcp_server_elicitation; mod mcp_server_status; mod mcp_tool; diff --git a/codex-rs/codex-mcp/src/binding.rs b/codex-rs/codex-mcp/src/binding.rs index fc48427399..c1e9373ef7 100644 --- a/codex-rs/codex-mcp/src/binding.rs +++ b/codex-rs/codex-mcp/src/binding.rs @@ -75,15 +75,23 @@ impl McpBinding { self.plugins_available } - /// Returns the frozen catalog captured for this binding. + /// Returns the frozen model-visible catalog captured for this binding. pub fn tools(&self) -> &[ToolInfo] { &self.tools } - /// Binds a call to the exact client and metadata advertised by this binding. + /// Returns permitted tool metadata, including app-only tools. + pub fn tool_info(&self, server: &str, tool: &str) -> Option<&ToolInfo> { + self.calls + .get(&(server.to_string(), tool.to_string())) + .map(PreparedMcpCall::tool_info) + } + + /// Binds a model-visible call to the exact client and metadata in this binding. pub fn prepare_call(&self, server: &str, tool: &str) -> Option { self.calls .get(&(server.to_string(), tool.to_string())) + .filter(|call| crate::tool_is_model_visible(call.tool_info())) .cloned() } diff --git a/codex-rs/codex-mcp/src/connection_manager/tool_catalog.rs b/codex-rs/codex-mcp/src/connection_manager/tool_catalog.rs index 9e41137701..4075bd616a 100644 --- a/codex-rs/codex-mcp/src/connection_manager/tool_catalog.rs +++ b/codex-rs/codex-mcp/src/connection_manager/tool_catalog.rs @@ -288,11 +288,11 @@ impl McpConnectionSet { let mut tools = Vec::with_capacity(listed_tools.len()); let mut calls = std::collections::HashMap::with_capacity(listed_tools.len()); for tool_info in listed_tools { - if !crate::tool_is_model_visible(&tool_info) { - continue; - } + let model_visible = crate::tool_is_model_visible(&tool_info); let Some(client) = clients.client(&tool_info.server_name) else { - tools.push(tool_info); + if model_visible { + tools.push(tool_info); + } continue; }; let Some(call) = self.prepare_call(&tool_info, client, Arc::clone(&config), *revision) @@ -311,7 +311,9 @@ impl McpConnectionSet { ), call, ); - tools.push(tool_info); + if model_visible { + tools.push(tool_info); + } } McpBinding::new( Arc::clone(self), diff --git a/codex-rs/codex-mcp/src/lib.rs b/codex-rs/codex-mcp/src/lib.rs index 0e6fad1fe1..f47b926d14 100644 --- a/codex-rs/codex-mcp/src/lib.rs +++ b/codex-rs/codex-mcp/src/lib.rs @@ -102,6 +102,7 @@ mod openai_docs_source_attribution; mod pagination; mod plugin_config; mod resource_client; +mod resource_origin; pub(crate) mod rmcp_client; pub(crate) mod runtime; pub(crate) mod server; diff --git a/codex-rs/codex-mcp/src/resource_origin.rs b/codex-rs/codex-mcp/src/resource_origin.rs new file mode 100644 index 0000000000..0a21c4c63e --- /dev/null +++ b/codex-rs/codex-mcp/src/resource_origin.rs @@ -0,0 +1,250 @@ +//! Thread-owned, bounded provenance for app-hosted widget resources. + +use std::collections::VecDeque; + +use anyhow::Context; +use codex_connectors::AppToolPolicyEvaluator; +use codex_connectors::AppToolPolicyInput; +use codex_protocol::ThreadId; +use codex_protocol::items::McpToolCallStatus; +use codex_protocol::items::TurnItem; +use codex_protocol::protocol::EventMsg; +use rmcp::model::ReadResourceRequestParams; +use rmcp::model::ReadResourceResult; + +use crate::CODEX_APPS_MCP_SERVER_NAME; +use crate::McpBinding; + +const MAX_ORIGINS: usize = 64; +const MAX_ORIGIN_BYTES: usize = 1024; + +#[derive(Default)] +pub(crate) struct ResourceOrigins { + origins: VecDeque, + turns: VecDeque, + current_turn_id: Option, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub(crate) struct ResourceOrigin { + call_id: String, + turn_id: Option, + tool: String, + connector_id: String, + link_id: Option, + uri: String, + ambiguous_account: bool, +} + +impl ResourceOrigins { + pub(crate) fn observe(&mut self, event: &EventMsg) { + match event { + EventMsg::TurnStarted(event) if event.turn_id.len() <= MAX_ORIGIN_BYTES => { + self.current_turn_id = Some(event.turn_id.clone()); + if self.turns.back() != Some(&event.turn_id) { + self.turns.push_back(event.turn_id.clone()); + if self.turns.len() > MAX_ORIGINS { + self.turns.pop_front(); + } + } + } + EventMsg::ItemCompleted(event) => { + let TurnItem::McpToolCall(item) = &event.item else { + return; + }; + if item.status == McpToolCallStatus::Completed { + self.remember( + &item.id, + Some(&event.turn_id), + &item.server, + &item.tool, + &item.arguments, + item.connector_id.as_deref(), + item.link_id.as_deref(), + item.mcp_app_resource_uri.as_deref(), + ); + } + } + EventMsg::McpToolCallEnd(event) + if event + .result + .as_ref() + .is_ok_and(|result| !result.is_error.unwrap_or(false)) => + { + let turn_id = self.current_turn_id.clone(); + self.remember( + &event.call_id, + turn_id.as_deref(), + &event.invocation.server, + &event.invocation.tool, + event + .invocation + .arguments + .as_ref() + .unwrap_or(&serde_json::Value::Null), + event.connector_id.as_deref(), + event.link_id.as_deref(), + event.mcp_app_resource_uri.as_deref(), + ); + } + EventMsg::ThreadRolledBack(event) => { + for _ in 0..event.num_turns { + let Some(turn_id) = self.turns.pop_back() else { + *self = Self::default(); + return; + }; + self.origins.retain(|origin| { + origin.turn_id.is_some() + && origin.turn_id.as_deref() != Some(turn_id.as_str()) + }); + } + self.current_turn_id = self.turns.back().cloned(); + } + _ => {} + } + } + + pub(crate) fn find(&self, call_id: &str) -> anyhow::Result { + self.origins + .iter() + .rev() + .find(|origin| origin.call_id == call_id) + .cloned() + .context("originating MCP tool call was not found or did not complete successfully") + } + + #[expect( + clippy::too_many_arguments, + reason = "only bounded provenance fields are retained" + )] + fn remember( + &mut self, + call_id: &str, + turn_id: Option<&str>, + server: &str, + tool: &str, + arguments: &serde_json::Value, + connector_id: Option<&str>, + link_id: Option<&str>, + uri: Option<&str>, + ) { + if server != CODEX_APPS_MCP_SERVER_NAME { + return; + } + let Some(connector_id) = connector_id.filter(|value| !value.trim().is_empty()) else { + return; + }; + let Some(uri) = uri else { + return; + }; + let link_id = link_id.filter(|value| !value.trim().is_empty()); + let ambiguous_account = arguments + .get("link_id") + .and_then(serde_json::Value::as_str) + .filter(|value| !value.trim().is_empty()) + .is_some_and(|argument_link_id| Some(argument_link_id) != link_id); + let origin = ResourceOrigin { + call_id: call_id.to_owned(), + turn_id: turn_id.map(str::to_owned), + tool: tool.to_owned(), + connector_id: connector_id.to_owned(), + link_id: link_id.map(str::to_owned), + uri: uri.to_owned(), + ambiguous_account, + }; + if origin.byte_len() > MAX_ORIGIN_BYTES { + return; + } + + if let Some(index) = self + .origins + .iter() + .position(|existing| existing.call_id == origin.call_id) + { + self.origins.remove(index); + } + if self.origins.len() >= MAX_ORIGINS { + self.origins.pop_front(); + } + self.origins.push_back(origin); + } +} + +impl ResourceOrigin { + fn byte_len(&self) -> usize { + self.call_id.len() + + self.turn_id.as_ref().map_or(0, String::len) + + self.tool.len() + + self.connector_id.len() + + self.link_id.as_ref().map_or(0, String::len) + + self.uri.len() + } + + pub(crate) async fn read( + &self, + binding: &McpBinding, + thread_id: ThreadId, + uri: &str, + ) -> anyhow::Result { + if self.uri != uri { + anyhow::bail!("originating MCP tool call does not match the requested resource"); + } + if self.ambiguous_account { + anyhow::bail!("originating MCP tool call has ambiguous account selection"); + } + + let tool_info = binding + .tool_info(CODEX_APPS_MCP_SERVER_NAME, &self.tool) + .context("originating MCP tool is unavailable")?; + if tool_info.connector_id.as_deref() != Some(self.connector_id.as_str()) { + anyhow::bail!("originating MCP tool connector does not match its app context"); + } + let tool_meta = tool_info.tool.meta.as_ref().map(|meta| &meta.0); + let current_link_id = tool_meta + .and_then(|meta| meta.get("link_id")) + .and_then(serde_json::Value::as_str) + .filter(|value| !value.trim().is_empty()); + if current_link_id != self.link_id.as_deref() { + anyhow::bail!("originating MCP tool link does not match its app context"); + } + if self.link_id.is_none() + && tool_meta + .and_then(|meta| meta.get("_codex_apps")) + .and_then(|meta| meta.get("requires_explicit_link_id")) + .and_then(serde_json::Value::as_bool) + == Some(true) + { + anyhow::bail!("originating MCP tool requires an explicit account link"); + } + + let annotations = tool_info.tool.annotations.as_ref(); + if !AppToolPolicyEvaluator::new(&binding.config().config_layer_stack) + .policy(AppToolPolicyInput { + connector_id: Some(&self.connector_id), + tool_name: tool_info.tool.name.as_ref(), + tool_title: tool_info.tool.title.as_deref(), + destructive_hint: annotations.and_then(|value| value.destructive_hint), + open_world_hint: annotations.and_then(|value| value.open_world_hint), + }) + .enabled + { + anyhow::bail!("originating MCP tool is disabled by app configuration"); + } + + let meta = serde_json::from_value(serde_json::json!({ + "threadId": thread_id, + "x-codex-turn-metadata": { + "mcp_request_meta": { + "selected_connector_ids": [&self.connector_id], + "link_id": &self.link_id, + } + }, + }))?; + binding + .read_resource( + CODEX_APPS_MCP_SERVER_NAME, + ReadResourceRequestParams::new(uri).with_meta(meta), + ) + .await + } +} diff --git a/codex-rs/codex-mcp/src/runtime.rs b/codex-rs/codex-mcp/src/runtime.rs index 8fc1a3a945..0aeff678ea 100644 --- a/codex-rs/codex-mcp/src/runtime.rs +++ b/codex-rs/codex-mcp/src/runtime.rs @@ -24,11 +24,13 @@ use codex_exec_server::HttpClient; use codex_exec_server::RouteAwareHttpClient; use codex_login::AuthManager; use codex_login::CodexAuth; +use codex_protocol::ThreadId; use codex_protocol::capabilities::SelectedCapabilityRoot; use codex_protocol::mcp::CallToolResult; use codex_protocol::mcp::ClientMcpExtensions; use codex_protocol::models::PermissionProfile; use codex_protocol::protocol::Event; +use codex_protocol::protocol::EventMsg; use codex_rmcp_client::ElicitationResponse; use codex_rmcp_client::with_http_headers_helper; use codex_utils_path_uri::PathUri; @@ -46,6 +48,7 @@ use crate::connection_manager::McpConnectionSet; use crate::elicitation::ElicitationLifecycle; use crate::elicitation::ElicitationRequestRouter; use crate::elicitation::ElicitationReviewerHandle; +use crate::resource_origin::ResourceOrigins; use crate::server::EffectiveMcpServer; use crate::tool_catalog_cache::McpToolCatalogCache; use crate::tools::ToolInfo; @@ -88,6 +91,7 @@ pub struct McpRuntime { current: ArcSwap, reconnect_pending: AtomicBool, elicitation_router: ElicitationRequestRouter, + resource_origins: Mutex, } struct PublishedMcpRuntime { @@ -171,9 +175,47 @@ impl McpRuntime { }), reconnect_pending: AtomicBool::new(false), elicitation_router: ElicitationRequestRouter::default(), + resource_origins: Mutex::default(), } } + /// Updates this thread's bounded resource provenance from a live or restored event. + pub fn observe_event(&self, event: &EventMsg) { + if !matches!( + event, + EventMsg::TurnStarted(_) + | EventMsg::ItemCompleted(_) + | EventMsg::McpToolCallEnd(_) + | EventMsg::ThreadRolledBack(_) + ) { + return; + } + self.resource_origins + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .observe(event); + } + + /// Reads a widget through the current binding of the app tool that produced it. + pub async fn read_resource_for_call( + &self, + thread_id: ThreadId, + call_id: &str, + uri: &str, + ) -> anyhow::Result { + let origin = self + .resource_origins + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + .find(call_id)?; + let binding = self + .current_binding_for_call(crate::CODEX_APPS_MCP_SERVER_NAME) + .await + .ok_or_else(|| anyhow::anyhow!("codex_apps MCP server is unavailable"))?; + + origin.read(&binding, thread_id, uri).await + } + pub async fn new(input: McpRuntimeInput) -> Self { let runtime = Self::empty(input.config.prefix_mcp_tool_names); runtime.replace(input).await; diff --git a/codex-rs/core/src/codex_thread.rs b/codex-rs/core/src/codex_thread.rs index cdc8c01320..f03e00767e 100644 --- a/codex-rs/core/src/codex_thread.rs +++ b/codex-rs/core/src/codex_thread.rs @@ -747,6 +747,23 @@ impl CodexThread { Ok(serde_json::to_value(result)?) } + /// Reads an app resource using the current authority of its originating tool call. + pub async fn read_mcp_resource_for_call( + &self, + call_id: &str, + uri: &str, + ) -> anyhow::Result { + self.session.refresh_mcp_if_dirty().await; + let result = self + .session + .services + .mcp_runtime + .read_resource_for_call(self.session.thread_id, call_id, uri) + .await?; + + Ok(serde_json::to_value(result)?) + } + pub async fn call_mcp_tool( &self, server: &str, diff --git a/codex-rs/core/src/session/mod.rs b/codex-rs/core/src/session/mod.rs index 340f471ef7..777b1c882c 100644 --- a/codex-rs/core/src/session/mod.rs +++ b/codex-rs/core/src/session/mod.rs @@ -2129,6 +2129,7 @@ impl Session { } async fn send_event_raw_with_persistence(&self, event: Event, persist: bool) { + self.services.mcp_runtime.observe_event(&event.msg); // Persist the event into rollout storage; the store applies its persistence policy. if persist { let rollout_items = vec![RolloutItem::EventMsg(event.msg.clone())]; diff --git a/codex-rs/core/src/session/session.rs b/codex-rs/core/src/session/session.rs index c82d0a0315..f622dc0cf1 100644 --- a/codex-rs/core/src/session/session.rs +++ b/codex-rs/core/src/session/session.rs @@ -1223,6 +1223,11 @@ impl Session { let mcp_runtime = Arc::new(McpRuntime::empty( mcp_projection.config.prefix_mcp_tool_names, )); + for item in initial_history.get_rollout_items() { + if let RolloutItem::EventMsg(event) = item { + mcp_runtime.observe_event(event); + } + } let session_extension_data = codex_extension_api::ExtensionData::new(session_id.to_string()); session_extension_data.insert(analytics_events_client.clone());