mirror of
https://github.com/openai/codex.git
synced 2026-09-13 11:47:17 +00:00
@@ -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;
|
||||
|
||||
@@ -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<Session>,
|
||||
turn_context: &Arc<TurnContext>,
|
||||
request: McpToolCallRequest<'_>,
|
||||
) -> Result<CallToolResult, String> {
|
||||
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<serde_json::Value>,
|
||||
request_meta: Option<serde_json::Value>,
|
||||
server_origin: Option<&'a str>,
|
||||
connector_id: Option<String>,
|
||||
connector_name: Option<String>,
|
||||
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) {
|
||||
|
||||
@@ -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<String>,
|
||||
}
|
||||
|
||||
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
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
@@ -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<Vec<u8>> = 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(())
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user