mirror of
https://github.com/openai/codex.git
synced 2026-09-13 11:47:17 +00:00
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
This commit is contained in:
@@ -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<T>(
|
||||
session: &Arc<Session>,
|
||||
turn: &TurnContext,
|
||||
call_id: &str,
|
||||
invocation: McpInvocation,
|
||||
operation: impl Future<Output = Result<T, FunctionCallError>>,
|
||||
) -> Result<Box<dyn ToolOutput>, 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<String>) -> Option<String> {
|
||||
input.and_then(|value| {
|
||||
let trimmed = value.trim().to_string();
|
||||
|
||||
@@ -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<ListResourceTemplatesPayload, FunctionCallError> = 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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<ListResourcesPayload, FunctionCallError> = 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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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<ReadResourcePayload, FunctionCallError> = 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
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user