From 9df628f38c84f99f03db67e8cb1f249b6106f89d Mon Sep 17 00:00:00 2001 From: nicholasclark-openai Date: Wed, 25 Mar 2026 11:28:17 -0700 Subject: [PATCH] Simplify MCP span diff Co-authored-by: Codex --- codex-rs/core/src/lib.rs | 1 - codex-rs/core/src/mcp_tool_call.rs | 171 ++++++++++++++--------- codex-rs/core/src/network_trace.rs | 87 ------------ codex-rs/core/tests/suite/mod.rs | 2 - codex-rs/core/tests/suite/otel_apps.rs | 104 -------------- codex-rs/core/tests/suite/rmcp_client.rs | 3 +- 6 files changed, 108 insertions(+), 260 deletions(-) delete mode 100644 codex-rs/core/src/network_trace.rs delete mode 100644 codex-rs/core/tests/suite/otel_apps.rs diff --git a/codex-rs/core/src/lib.rs b/codex-rs/core/src/lib.rs index 4dec339bb6..71f4328aed 100644 --- a/codex-rs/core/src/lib.rs +++ b/codex-rs/core/src/lib.rs @@ -53,7 +53,6 @@ mod mcp_tool_approval_templates; pub mod models_manager; mod network_policy_decision; pub mod network_proxy_loader; -mod network_trace; mod original_image_detail; mod packages; pub use mcp_connection_manager::MCP_SANDBOX_STATE_CAPABILITY; diff --git a/codex-rs/core/src/mcp_tool_call.rs b/codex-rs/core/src/mcp_tool_call.rs index 3484dfedef..fc4f2f61a5 100644 --- a/codex-rs/core/src/mcp_tool_call.rs +++ b/codex-rs/core/src/mcp_tool_call.rs @@ -51,6 +51,9 @@ use std::path::Path; use std::sync::Arc; use toml_edit::value; use tracing::Instrument; +use tracing::Span; +use tracing::field::Empty; +use url::Url; /// Handles the specified tool call dispatches the appropriate /// `McpToolCallBegin` and `McpToolCallEnd` events to the `Session`. @@ -159,21 +162,39 @@ pub(crate) async fn handle_mcp_tool_call( maybe_mark_thread_memory_mode_polluted(sess.as_ref(), turn_context.as_ref()).await; let start = Instant::now(); - let result = execute_mcp_tool_call( - &sess, - turn_context, - McpToolCallRequest { - server: &server, + let result = async { + sess.call_tool( + &server, + &tool_name, + arguments_value.clone(), + request_meta.clone(), + ) + .await + .map_err(|e| format!("tool call error: {e:?}")) + } + .instrument(mcp_tool_call_span( + sess.as_ref(), + turn_context.as_ref(), + McpToolCallSpanFields { + server_name: &server, tool_name: &tool_name, call_id: &call_id, - arguments_value: arguments_value.clone(), - request_meta: request_meta.clone(), server_origin: server_origin.as_deref(), - connector_id: connector_id.clone(), - connector_name: connector_name.clone(), + connector_id: connector_id.as_deref(), + connector_name: connector_name.as_deref(), }, - ) + )) .await; + let result = sanitize_mcp_tool_result_for_model( + turn_context + .model_info + .input_modalities + .contains(&InputModality::Image), + result, + ); + if let Err(error) = &result { + tracing::warn!("MCP tool call error: {error:?}"); + } let tool_call_end_event = EventMsg::McpToolCallEnd(McpToolCallEndEvent { call_id: call_id.clone(), invocation, @@ -246,21 +267,34 @@ pub(crate) async fn handle_mcp_tool_call( let start = Instant::now(); // Perform the tool call. - let result = execute_mcp_tool_call( - &sess, - turn_context, - McpToolCallRequest { - server: &server, + let result = async { + sess.call_tool(&server, &tool_name, arguments_value.clone(), request_meta) + .await + .map_err(|e| format!("tool call error: {e:?}")) + } + .instrument(mcp_tool_call_span( + sess.as_ref(), + turn_context.as_ref(), + McpToolCallSpanFields { + server_name: &server, tool_name: &tool_name, call_id: &call_id, - arguments_value: arguments_value.clone(), - request_meta, server_origin: server_origin.as_deref(), - connector_id, - connector_name, + connector_id: connector_id.as_deref(), + connector_name: connector_name.as_deref(), }, - ) + )) .await; + let result = sanitize_mcp_tool_result_for_model( + turn_context + .model_info + .input_modalities + .contains(&InputModality::Image), + result, + ); + if let Err(error) = &result { + tracing::warn!("MCP tool call error: {error:?}"); + } let tool_call_end_event = EventMsg::McpToolCallEnd(McpToolCallEndEvent { call_id: call_id.clone(), invocation, @@ -284,57 +318,64 @@ pub(crate) async fn handle_mcp_tool_call( CallToolResult::from_result(result) } -async fn execute_mcp_tool_call( - sess: &Arc, - turn_context: &Arc, - request: McpToolCallRequest<'_>, -) -> Result { - let tool_call_span = crate::network_trace::mcp_tool_call_span( - sess.as_ref(), - turn_context.as_ref(), - crate::network_trace::McpToolCallTrace { - server_name: request.server, - tool_name: request.tool_name, - call_id: request.call_id, - server_origin: request.server_origin, - connector_id: request.connector_id.as_deref(), - connector_name: request.connector_name.as_deref(), - }, +fn mcp_tool_call_span( + session: &Session, + turn_context: &TurnContext, + fields: McpToolCallSpanFields<'_>, +) -> Span { + let transport = match fields.server_origin { + Some("stdio") => "stdio", + Some(_) => "streamable_http", + None => "", + }; + let span = tracing::info_span!( + "mcp.tools.call", + otel.kind = "client", + rpc.system = "jsonrpc", + rpc.method = "tools/call", + mcp.server.name = fields.server_name, + mcp.server.origin = fields.server_origin.unwrap_or(""), + mcp.transport = transport, + mcp.connector.id = fields.connector_id.unwrap_or(""), + mcp.connector.name = fields.connector_name.unwrap_or(""), + tool.name = fields.tool_name, + tool.call_id = fields.call_id, + conversation.id = Empty, + session.id = Empty, + turn.id = Empty, + server.address = Empty, + server.port = Empty, ); - let result = async { - sess.call_tool( - request.server, - request.tool_name, - request.arguments_value, - request.request_meta, - ) - .await - .map_err(|e| format!("tool call error: {e:?}")) - } - .instrument(tool_call_span) - .await; - let result = sanitize_mcp_tool_result_for_model( - turn_context - .model_info - .input_modalities - .contains(&InputModality::Image), - result, - ); - if let Err(error) = &result { - tracing::warn!("MCP tool call error: {error:?}"); - } - result + let conversation_id = session.conversation_id.to_string(); + span.record("conversation.id", conversation_id.as_str()); + span.record("session.id", conversation_id.as_str()); + span.record("turn.id", turn_context.sub_id.as_str()); + record_server_fields(&span, fields.server_origin); + span } -struct McpToolCallRequest<'a> { - server: &'a str, +struct McpToolCallSpanFields<'a> { + server_name: &'a str, tool_name: &'a str, call_id: &'a str, - arguments_value: Option, - request_meta: Option, server_origin: Option<&'a str>, - connector_id: Option, - connector_name: Option, + connector_id: Option<&'a str>, + connector_name: Option<&'a str>, +} + +fn record_server_fields(span: &Span, url: Option<&str>) { + let Some(url) = url else { + return; + }; + let Ok(parsed) = Url::parse(url) else { + return; + }; + if let Some(host) = parsed.host_str() { + span.record("server.address", host); + } + if let Some(port) = parsed.port_or_known_default() { + span.record("server.port", port as i64); + } } async fn maybe_mark_thread_memory_mode_polluted(sess: &Session, turn_context: &TurnContext) { diff --git a/codex-rs/core/src/network_trace.rs b/codex-rs/core/src/network_trace.rs deleted file mode 100644 index cdfa15075c..0000000000 --- a/codex-rs/core/src/network_trace.rs +++ /dev/null @@ -1,87 +0,0 @@ -use tracing::Span; -use tracing::field::Empty; -use url::Url; - -use crate::codex::Session; -use crate::codex::TurnContext; - -struct CorrelationFields { - conversation_id: String, - session_id: String, - turn_id: Option, -} - -impl CorrelationFields { - fn from_turn_context(session: &Session, turn_context: &TurnContext) -> Self { - Self { - conversation_id: session.conversation_id.to_string(), - session_id: session.conversation_id.to_string(), - turn_id: Some(turn_context.sub_id.clone()), - } - } - - fn record_on(&self, span: &Span) { - span.record("conversation.id", self.conversation_id.as_str()); - span.record("session.id", self.session_id.as_str()); - if let Some(turn_id) = self.turn_id.as_deref() { - span.record("turn.id", turn_id); - } - } -} - -fn record_server_fields(span: &Span, url: Option<&str>) { - let Some(url) = url else { - return; - }; - let Ok(parsed) = Url::parse(url) else { - return; - }; - if let Some(host) = parsed.host_str() { - span.record("server.address", host); - } - if let Some(port) = parsed.port_or_known_default() { - span.record("server.port", port as i64); - } -} - -pub(crate) struct McpToolCallTrace<'a> { - pub server_name: &'a str, - pub tool_name: &'a str, - pub call_id: &'a str, - pub server_origin: Option<&'a str>, - pub connector_id: Option<&'a str>, - pub connector_name: Option<&'a str>, -} - -pub(crate) fn mcp_tool_call_span( - session: &Session, - turn_context: &TurnContext, - trace: McpToolCallTrace<'_>, -) -> Span { - let transport = match trace.server_origin { - Some("stdio") => "stdio", - Some(_) => "streamable_http", - None => "", - }; - let span = tracing::info_span!( - "mcp.tools.call", - otel.kind = "client", - rpc.system = "jsonrpc", - rpc.method = "tools/call", - mcp.server.name = trace.server_name, - mcp.server.origin = trace.server_origin.unwrap_or(""), - mcp.transport = transport, - mcp.connector.id = trace.connector_id.unwrap_or(""), - mcp.connector.name = trace.connector_name.unwrap_or(""), - tool.name = trace.tool_name, - tool.call_id = trace.call_id, - conversation.id = Empty, - session.id = Empty, - turn.id = Empty, - server.address = Empty, - server.port = Empty, - ); - CorrelationFields::from_turn_context(session, turn_context).record_on(&span); - record_server_fields(&span, trace.server_origin); - span -} diff --git a/codex-rs/core/tests/suite/mod.rs b/codex-rs/core/tests/suite/mod.rs index 8e72d9efc1..f4891d58c5 100644 --- a/codex-rs/core/tests/suite/mod.rs +++ b/codex-rs/core/tests/suite/mod.rs @@ -92,8 +92,6 @@ mod model_visible_layout; mod models_cache_ttl; mod models_etag_responses; mod otel; -#[cfg(not(target_os = "windows"))] -mod otel_apps; mod pending_input; mod permissions_messages; mod personality; diff --git a/codex-rs/core/tests/suite/otel_apps.rs b/codex-rs/core/tests/suite/otel_apps.rs deleted file mode 100644 index 75e45dac14..0000000000 --- a/codex-rs/core/tests/suite/otel_apps.rs +++ /dev/null @@ -1,104 +0,0 @@ -#![allow(clippy::unwrap_used, clippy::expect_used)] - -use anyhow::Result; -use codex_core::CodexAuth; -use codex_features::Feature; -use codex_protocol::protocol::EventMsg; -use codex_protocol::protocol::Op; -use codex_protocol::user_input::UserInput; -use core_test_support::apps_test_server::AppsTestServer; -use core_test_support::responses::ev_assistant_message; -use core_test_support::responses::ev_completed; -use core_test_support::responses::ev_function_call; -use core_test_support::responses::ev_response_created; -use core_test_support::responses::mount_sse_once; -use core_test_support::responses::sse; -use core_test_support::responses::start_mock_server; -use core_test_support::skip_if_no_network; -use core_test_support::test_codex::test_codex; -use core_test_support::wait_for_event; -use std::sync::Mutex; -use tracing::Level; -use tracing_subscriber::fmt::format::FmtSpan; -use tracing_test::internal::MockWriter; - -#[tokio::test(flavor = "current_thread")] -async fn codex_apps_mcp_span_records_connector_metadata() -> Result<()> { - skip_if_no_network!(Ok(())); - - let buffer: &'static Mutex> = Box::leak(Box::new(Mutex::new(Vec::new()))); - let subscriber = tracing_subscriber::fmt() - .with_level(true) - .with_ansi(false) - .with_max_level(Level::TRACE) - .with_span_events(FmtSpan::FULL) - .with_writer(MockWriter::new(buffer)) - .finish(); - let _guard = tracing::subscriber::set_default(subscriber); - - let server = start_mock_server().await; - let apps_server = AppsTestServer::mount(&server).await?; - - mount_sse_once( - &server, - sse(vec![ - ev_response_created("resp-1"), - ev_function_call( - "calendar-call-1", - "mcp__codex_apps__calendar_create_event", - r#"{"title":"Lunch","starts_at":"2026-03-10T12:00:00Z"}"#, - ), - ev_completed("resp-1"), - ]), - ) - .await; - mount_sse_once( - &server, - sse(vec![ - ev_assistant_message("msg-1", "calendar tool completed successfully."), - ev_completed("resp-2"), - ]), - ) - .await; - - let fixture = test_codex() - .with_auth(CodexAuth::create_dummy_chatgpt_auth_for_testing()) - .with_config(move |config| { - config - .features - .enable(Feature::Apps) - .expect("test config should allow feature update"); - config.chatgpt_base_url = apps_server.chatgpt_base_url; - config.model = Some("gpt-5-codex".to_string()); - }) - .build(&server) - .await?; - - fixture - .codex - .submit(Op::UserInput { - items: vec![UserInput::Text { - text: "create a calendar event".into(), - text_elements: Vec::new(), - }], - final_output_json_schema: None, - }) - .await?; - - wait_for_event(&fixture.codex, |event| { - matches!(event, EventMsg::TurnComplete(_)) - }) - .await; - - let logs = String::from_utf8(buffer.lock().unwrap().clone()).unwrap(); - assert!( - logs.contains("mcp.tools.call{otel.kind=\"client\"") - && logs.contains("mcp.server.name=\"codex_apps\"") - && logs.contains("mcp.connector.id=\"calendar\"") - && logs.contains("mcp.connector.name=\"Calendar\"") - && logs.contains("tool.name=\"calendar_create_event\""), - "missing connector metadata on mcp.tools.call span\nlogs:\n{logs}" - ); - - Ok(()) -} diff --git a/codex-rs/core/tests/suite/rmcp_client.rs b/codex-rs/core/tests/suite/rmcp_client.rs index b7e1b84a44..9a6122bb3a 100644 --- a/codex-rs/core/tests/suite/rmcp_client.rs +++ b/codex-rs/core/tests/suite/rmcp_client.rs @@ -207,6 +207,7 @@ async fn stdio_server_round_trip() -> anyhow::Result<()> { assert!( logs.contains("turn{otel.name=\"session_task.turn\"") && logs.contains("mcp.tools.call{otel.kind=\"client\"") + && logs.contains("mcp.client.operation{otel.kind=\"client\"") && logs.contains("rpc.system=\"jsonrpc\"") && logs.contains("rpc.method=\"tools/call\"") && logs.contains("mcp.server.name=\"rmcp\"") @@ -214,7 +215,7 @@ async fn stdio_server_round_trip() -> anyhow::Result<()> { && logs.contains("tool.name=\"echo\"") && logs.contains("tool.call_id=\"call-123\"") && logs.contains("turn.id="), - "missing mcp.tools.call span nested under session_task.turn\nlogs:\n{logs}" + "missing MCP tracing spans nested under session_task.turn\nlogs:\n{logs}" ); server.verify().await;