mirror of
https://github.com/openai/codex.git
synced 2026-09-13 11:47:17 +00:00
Keep top MCP span PR core-only
Co-authored-by: Codex <noreply@openai.com>
This commit is contained in:
@@ -207,7 +207,6 @@ 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\"")
|
||||
|
||||
@@ -64,8 +64,6 @@ use tokio::io::BufReader;
|
||||
use tokio::process::Command;
|
||||
use tokio::sync::Mutex;
|
||||
use tokio::time;
|
||||
use tracing::Instrument;
|
||||
use tracing::field::Empty;
|
||||
use tracing::info;
|
||||
use tracing::warn;
|
||||
|
||||
@@ -1054,93 +1052,41 @@ impl RmcpClient {
|
||||
Fut: std::future::Future<Output = std::result::Result<T, rmcp::service::ServiceError>>,
|
||||
{
|
||||
let service = self.service().await?;
|
||||
match Self::run_service_operation_once(
|
||||
Arc::clone(&service),
|
||||
label,
|
||||
timeout,
|
||||
self.service_operation_span(label),
|
||||
&operation,
|
||||
)
|
||||
.await
|
||||
match Self::run_service_operation_once(Arc::clone(&service), label, timeout, &operation)
|
||||
.await
|
||||
{
|
||||
Ok(result) => Ok(result),
|
||||
Err(error) if Self::is_session_expired_404(&error) => {
|
||||
self.reinitialize_after_session_expiry(&service).await?;
|
||||
let recovered_service = self.service().await?;
|
||||
Self::run_service_operation_once(
|
||||
recovered_service,
|
||||
label,
|
||||
timeout,
|
||||
self.service_operation_span(label),
|
||||
&operation,
|
||||
)
|
||||
.await
|
||||
.map_err(Into::into)
|
||||
Self::run_service_operation_once(recovered_service, label, timeout, &operation)
|
||||
.await
|
||||
.map_err(Into::into)
|
||||
}
|
||||
Err(error) => Err(error.into()),
|
||||
}
|
||||
}
|
||||
|
||||
fn service_operation_span(&self, label: &str) -> tracing::Span {
|
||||
let span = tracing::info_span!(
|
||||
"mcp.client.operation",
|
||||
otel.kind = "client",
|
||||
rpc.system = "jsonrpc",
|
||||
rpc.method = label,
|
||||
mcp.transport = Empty,
|
||||
mcp.server.name = Empty,
|
||||
server.address = Empty,
|
||||
server.port = Empty,
|
||||
);
|
||||
|
||||
match &self.transport_recipe {
|
||||
TransportRecipe::Stdio { .. } => {
|
||||
span.record("mcp.transport", "stdio");
|
||||
}
|
||||
TransportRecipe::StreamableHttp {
|
||||
server_name, url, ..
|
||||
} => {
|
||||
span.record("mcp.transport", "streamable_http");
|
||||
span.record("mcp.server.name", server_name.as_str());
|
||||
if let Ok(parsed_url) = reqwest::Url::parse(url) {
|
||||
if let Some(host) = parsed_url.host_str() {
|
||||
span.record("server.address", host);
|
||||
}
|
||||
if let Some(port) = parsed_url.port_or_known_default() {
|
||||
span.record("server.port", port as i64);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
span
|
||||
}
|
||||
|
||||
async fn run_service_operation_once<T, F, Fut>(
|
||||
service: Arc<RunningService<RoleClient, LoggingClientHandler>>,
|
||||
label: &str,
|
||||
timeout: Option<Duration>,
|
||||
operation_span: tracing::Span,
|
||||
operation: &F,
|
||||
) -> std::result::Result<T, ClientOperationError>
|
||||
where
|
||||
F: Fn(Arc<RunningService<RoleClient, LoggingClientHandler>>) -> Fut,
|
||||
Fut: std::future::Future<Output = std::result::Result<T, rmcp::service::ServiceError>>,
|
||||
{
|
||||
async move {
|
||||
match timeout {
|
||||
Some(duration) => time::timeout(duration, operation(service))
|
||||
.await
|
||||
.map_err(|_| ClientOperationError::Timeout {
|
||||
label: label.to_string(),
|
||||
duration,
|
||||
})?
|
||||
.map_err(ClientOperationError::from),
|
||||
None => operation(service).await.map_err(ClientOperationError::from),
|
||||
}
|
||||
match timeout {
|
||||
Some(duration) => time::timeout(duration, operation(service))
|
||||
.await
|
||||
.map_err(|_| ClientOperationError::Timeout {
|
||||
label: label.to_string(),
|
||||
duration,
|
||||
})?
|
||||
.map_err(ClientOperationError::from),
|
||||
None => operation(service).await.map_err(ClientOperationError::from),
|
||||
}
|
||||
.instrument(operation_span)
|
||||
.await
|
||||
}
|
||||
|
||||
fn is_session_expired_404(error: &ClientOperationError) -> bool {
|
||||
|
||||
Reference in New Issue
Block a user