From 1b594980f3eb6998554097eed87e98a3ecb7ebfb Mon Sep 17 00:00:00 2001 From: jif Date: Mon, 3 Aug 2026 09:44:32 +0000 Subject: [PATCH] Deduplicate MCP resource operation handling (#36716) ## What changed - Add a shared runner for MCP resource operation lifecycle events, output serialization, truncation, timing, and error handling. - Use it for listing resources, listing resource templates, and reading resources. GitOrigin-RevId: 84cae2e01a096d5b8ed97ea1cb462f01fe2ed1f9 --- .../core/src/tools/handlers/mcp_resource.rs | 51 ++++++++++++++ .../list_mcp_resource_templates.rs | 67 ++----------------- .../mcp_resource/list_mcp_resources.rs | 67 ++----------------- .../mcp_resource/read_mcp_resource.rs | 67 ++----------------- 4 files changed, 63 insertions(+), 189 deletions(-) diff --git a/codex-rs/core/src/tools/handlers/mcp_resource.rs b/codex-rs/core/src/tools/handlers/mcp_resource.rs index 9196566a14..0cd5857302 100644 --- a/codex-rs/core/src/tools/handlers/mcp_resource.rs +++ b/codex-rs/core/src/tools/handlers/mcp_resource.rs @@ -1,6 +1,8 @@ use std::collections::HashMap; +use std::future::Future; use std::sync::Arc; use std::time::Duration; +use std::time::Instant; use codex_mcp::CODEX_APPS_MCP_SERVER_NAME; use codex_protocol::items::McpToolCallError; @@ -8,6 +10,7 @@ use codex_protocol::items::McpToolCallItem; use codex_protocol::items::McpToolCallStatus; use codex_protocol::items::TurnItem; use codex_protocol::mcp::CallToolResult; +use codex_protocol::models::function_call_output_content_items_to_text; use codex_protocol::protocol::TruncationPolicy; use codex_utils_output_truncation::truncate_text; use rmcp::model::ListResourceTemplatesResult; @@ -24,6 +27,8 @@ use crate::function_tool::FunctionCallError; use crate::session::session::Session; use crate::session::turn_context::TurnContext; use crate::tools::context::FunctionToolOutput; +use crate::tools::context::ToolOutput; +use crate::tools::context::boxed_tool_output; use codex_protocol::protocol::McpInvocation; mod list_mcp_resource_templates; @@ -280,6 +285,52 @@ async fn emit_tool_call_end( session.emit_turn_item_completed(turn, item).await; } +async fn run_resource_operation( + session: &Arc, + turn: &TurnContext, + call_id: &str, + invocation: McpInvocation, + operation: impl Future>, +) -> Result, FunctionCallError> +where + T: Serialize, +{ + emit_tool_call_begin(session, turn, call_id, invocation.clone()).await; + let start = Instant::now(); + let result = operation.await.and_then(|payload| { + serialize_function_output(payload, turn.model_info.truncation_policy.into()) + }); + + match result { + Ok(output) => { + let content = + function_call_output_content_items_to_text(&output.body).unwrap_or_default(); + emit_tool_call_end( + session, + turn, + call_id, + invocation, + start.elapsed(), + Ok(call_tool_result_from_content(&content, output.success)), + ) + .await; + Ok(boxed_tool_output(output)) + } + Err(error) => { + emit_tool_call_end( + session, + turn, + call_id, + invocation, + start.elapsed(), + Err(error.to_string()), + ) + .await; + Err(error) + } + } +} + fn normalize_optional_string(input: Option) -> Option { input.and_then(|value| { let trimmed = value.trim().to_string(); diff --git a/codex-rs/core/src/tools/handlers/mcp_resource/list_mcp_resource_templates.rs b/codex-rs/core/src/tools/handlers/mcp_resource/list_mcp_resource_templates.rs index eb28ce3e5a..94b0132376 100644 --- a/codex-rs/core/src/tools/handlers/mcp_resource/list_mcp_resource_templates.rs +++ b/codex-rs/core/src/tools/handlers/mcp_resource/list_mcp_resource_templates.rs @@ -1,13 +1,9 @@ -use std::time::Instant; - use crate::function_tool::FunctionCallError; use crate::tools::context::ToolInvocation; use crate::tools::context::ToolPayload; -use crate::tools::context::boxed_tool_output; use crate::tools::handlers::mcp_resource_spec::create_list_mcp_resource_templates_tool; use crate::tools::registry::CoreToolRuntime; use crate::tools::registry::ToolExecutor; -use codex_protocol::models::function_call_output_content_items_to_text; use codex_protocol::protocol::McpInvocation; use codex_tools::ToolName; use codex_tools::ToolSpec; @@ -16,15 +12,12 @@ use rmcp::model::PaginatedRequestParams; use super::ListResourceTemplatesArgs; use super::ListResourceTemplatesPayload; -use super::call_tool_result_from_content; -use super::emit_tool_call_begin; -use super::emit_tool_call_end; use super::ensure_model_can_access_mcp_server; use super::model_can_access_mcp_server; use super::normalize_optional_string; use super::parse_args_with_default; use super::parse_arguments; -use super::serialize_function_output; +use super::run_resource_operation; pub struct ListMcpResourceTemplatesHandler; @@ -82,10 +75,7 @@ impl ListMcpResourceTemplatesHandler { arguments: arguments.clone(), }; - emit_tool_call_begin(&session, turn.as_ref(), &call_id, invocation.clone()).await; - let start = Instant::now(); - - let payload_result: Result = async { + run_resource_operation(&session, turn.as_ref(), &call_id, invocation, async { if let Some(server_name) = server.clone() { ensure_model_can_access_mcp_server(turn.as_ref(), &server_name)?; let params = cursor @@ -117,57 +107,8 @@ impl ListMcpResourceTemplatesHandler { .await; Ok(ListResourceTemplatesPayload::from_all_servers(templates)) } - } - .await; - let truncation_policy = turn.model_info.truncation_policy.into(); - - match payload_result { - Ok(payload) => match serialize_function_output(payload, truncation_policy) { - Ok(output) => { - let content = function_call_output_content_items_to_text(&output.body) - .unwrap_or_default(); - let duration = start.elapsed(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Ok(call_tool_result_from_content(&content, output.success)), - ) - .await; - Ok(boxed_tool_output(output)) - } - Err(err) => { - let duration = start.elapsed(); - let message = err.to_string(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Err(message.clone()), - ) - .await; - Err(err) - } - }, - Err(err) => { - let duration = start.elapsed(); - let message = err.to_string(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Err(message.clone()), - ) - .await; - Err(err) - } - } + }) + .await } } diff --git a/codex-rs/core/src/tools/handlers/mcp_resource/list_mcp_resources.rs b/codex-rs/core/src/tools/handlers/mcp_resource/list_mcp_resources.rs index 464837bc0a..a966fbf4d2 100644 --- a/codex-rs/core/src/tools/handlers/mcp_resource/list_mcp_resources.rs +++ b/codex-rs/core/src/tools/handlers/mcp_resource/list_mcp_resources.rs @@ -1,13 +1,9 @@ -use std::time::Instant; - use crate::function_tool::FunctionCallError; use crate::tools::context::ToolInvocation; use crate::tools::context::ToolPayload; -use crate::tools::context::boxed_tool_output; use crate::tools::handlers::mcp_resource_spec::create_list_mcp_resources_tool; use crate::tools::registry::CoreToolRuntime; use crate::tools::registry::ToolExecutor; -use codex_protocol::models::function_call_output_content_items_to_text; use codex_protocol::protocol::McpInvocation; use codex_tools::ToolName; use codex_tools::ToolSpec; @@ -16,15 +12,12 @@ use rmcp::model::PaginatedRequestParams; use super::ListResourcesArgs; use super::ListResourcesPayload; -use super::call_tool_result_from_content; -use super::emit_tool_call_begin; -use super::emit_tool_call_end; use super::ensure_model_can_access_mcp_server; use super::model_can_access_mcp_server; use super::normalize_optional_string; use super::parse_args_with_default; use super::parse_arguments; -use super::serialize_function_output; +use super::run_resource_operation; pub struct ListMcpResourcesHandler; @@ -82,10 +75,7 @@ impl ListMcpResourcesHandler { arguments: arguments.clone(), }; - emit_tool_call_begin(&session, turn.as_ref(), &call_id, invocation.clone()).await; - let start = Instant::now(); - - let payload_result: Result = async { + run_resource_operation(&session, turn.as_ref(), &call_id, invocation, async { if let Some(server_name) = server.clone() { ensure_model_can_access_mcp_server(turn.as_ref(), &server_name)?; let params = cursor @@ -115,57 +105,8 @@ impl ListMcpResourcesHandler { .await; Ok(ListResourcesPayload::from_all_servers(resources)) } - } - .await; - let truncation_policy = turn.model_info.truncation_policy.into(); - - match payload_result { - Ok(payload) => match serialize_function_output(payload, truncation_policy) { - Ok(output) => { - let content = function_call_output_content_items_to_text(&output.body) - .unwrap_or_default(); - let duration = start.elapsed(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Ok(call_tool_result_from_content(&content, output.success)), - ) - .await; - Ok(boxed_tool_output(output)) - } - Err(err) => { - let duration = start.elapsed(); - let message = err.to_string(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Err(message.clone()), - ) - .await; - Err(err) - } - }, - Err(err) => { - let duration = start.elapsed(); - let message = err.to_string(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Err(message.clone()), - ) - .await; - Err(err) - } - } + }) + .await } } diff --git a/codex-rs/core/src/tools/handlers/mcp_resource/read_mcp_resource.rs b/codex-rs/core/src/tools/handlers/mcp_resource/read_mcp_resource.rs index 2a83e5961f..3049f391c9 100644 --- a/codex-rs/core/src/tools/handlers/mcp_resource/read_mcp_resource.rs +++ b/codex-rs/core/src/tools/handlers/mcp_resource/read_mcp_resource.rs @@ -1,13 +1,9 @@ -use std::time::Instant; - use crate::function_tool::FunctionCallError; use crate::tools::context::ToolInvocation; use crate::tools::context::ToolPayload; -use crate::tools::context::boxed_tool_output; use crate::tools::handlers::mcp_resource_spec::create_read_mcp_resource_tool; use crate::tools::registry::CoreToolRuntime; use crate::tools::registry::ToolExecutor; -use codex_protocol::models::function_call_output_content_items_to_text; use codex_protocol::protocol::McpInvocation; use codex_tools::ToolName; use codex_tools::ToolSpec; @@ -16,14 +12,11 @@ use rmcp::model::ReadResourceRequestParams; use super::ReadResourceArgs; use super::ReadResourcePayload; -use super::call_tool_result_from_content; -use super::emit_tool_call_begin; -use super::emit_tool_call_end; use super::ensure_model_can_access_mcp_server; use super::normalize_required_string; use super::parse_args; use super::parse_arguments; -use super::serialize_function_output; +use super::run_resource_operation; pub struct ReadMcpResourceHandler; @@ -81,10 +74,7 @@ impl ReadMcpResourceHandler { arguments: arguments.clone(), }; - emit_tool_call_begin(&session, turn.as_ref(), &call_id, invocation.clone()).await; - let start = Instant::now(); - - let payload_result: Result = async { + run_resource_operation(&session, turn.as_ref(), &call_id, invocation, async { ensure_model_can_access_mcp_server(turn.as_ref(), &server)?; let result = mcp .read_resource(&server, ReadResourceRequestParams::new(uri.clone())) @@ -98,57 +88,8 @@ impl ReadMcpResourceHandler { uri, result, }) - } - .await; - let truncation_policy = turn.model_info.truncation_policy.into(); - - match payload_result { - Ok(payload) => match serialize_function_output(payload, truncation_policy) { - Ok(output) => { - let content = function_call_output_content_items_to_text(&output.body) - .unwrap_or_default(); - let duration = start.elapsed(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Ok(call_tool_result_from_content(&content, output.success)), - ) - .await; - Ok(boxed_tool_output(output)) - } - Err(err) => { - let duration = start.elapsed(); - let message = err.to_string(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Err(message.clone()), - ) - .await; - Err(err) - } - }, - Err(err) => { - let duration = start.elapsed(); - let message = err.to_string(); - emit_tool_call_end( - &session, - turn.as_ref(), - &call_id, - invocation, - duration, - Err(message.clone()), - ) - .await; - Err(err) - } - } + }) + .await } }