diff --git a/codex-rs/exec-server/src/client.rs b/codex-rs/exec-server/src/client.rs index 2e2f4ce450e1..94992bf587e5 100644 --- a/codex-rs/exec-server/src/client.rs +++ b/codex-rs/exec-server/src/client.rs @@ -1310,6 +1310,7 @@ mod tests { } #[tokio::test(flavor = "current_thread")] + #[serial_test::serial(exec_server_tracing)] async fn process_start_propagates_caller_trace_context_across_background_task() { let (client_stdin, server_reader) = duplex(1 << 20); let (mut server_writer, client_stdout) = duplex(1 << 20); diff --git a/codex-rs/exec-server/src/remote.rs b/codex-rs/exec-server/src/remote.rs index 8870e0efe98c..ef200885a76b 100644 --- a/codex-rs/exec-server/src/remote.rs +++ b/codex-rs/exec-server/src/remote.rs @@ -711,6 +711,7 @@ mod tests { } #[tokio::test(flavor = "current_thread")] + #[serial_test::serial(exec_server_tracing)] async fn register_environment_posts_with_auth_provider_headers() { let provider = SdkTracerProvider::builder().build(); let tracer = provider.tracer("exec-server-test"); diff --git a/codex-rs/exec-server/src/rpc.rs b/codex-rs/exec-server/src/rpc.rs index 9671f25d3933..2875efcdfb92 100644 --- a/codex-rs/exec-server/src/rpc.rs +++ b/codex-rs/exec-server/src/rpc.rs @@ -363,23 +363,19 @@ impl RpcClient { }) } - #[tracing::instrument( - name = "codex.exec_server.request", - level = "info", - skip_all, - fields( - otel.kind = "client", - otel.name = method, - method, - ) - )] + fn allocate_request_id(&self) -> RequestId { + RequestId::Integer(self.next_request_id.fetch_add(1, Ordering::SeqCst)) + } + pub(crate) async fn call(&self, method: &str, params: &P) -> Result where P: Serialize, T: DeserializeOwned, { let _call_slot = self.acquire_regular_call_slot()?; - self.call_inner(method, params, RpcCallTimeout::None).await + let request_id = self.allocate_request_id(); + self.call_inner(request_id, method, params, RpcCallTimeout::None) + .await } pub(crate) async fn call_with_timeout( @@ -393,20 +389,16 @@ impl RpcClient { T: DeserializeOwned, { let _call_slot = self.acquire_regular_call_slot()?; - self.call_inner(method, params, RpcCallTimeout::After(call_timeout)) - .await - } - - #[tracing::instrument( - name = "codex.exec_server.request", - level = "info", - skip_all, - fields( - otel.kind = "client", - otel.name = method, + let request_id = self.allocate_request_id(); + self.call_inner( + request_id, method, + params, + RpcCallTimeout::After(call_timeout), ) - )] + .await + } + pub(crate) async fn call_for_cleanup( &self, method: &str, @@ -426,11 +418,27 @@ impl RpcClient { } }, }; - self.call_inner(method, params, RpcCallTimeout::None).await + let request_id = self.allocate_request_id(); + self.call_inner(request_id, method, params, RpcCallTimeout::None) + .await } + #[tracing::instrument( + name = "codex.exec_server.request", + level = "info", + skip_all, + fields( + otel.kind = "client", + otel.name = method, + rpc.system = "jsonrpc", + rpc.method = method, + rpc.request_id = %request_id, + method, + ) + )] async fn call_inner( &self, + request_id: RequestId, method: &str, params: &P, call_timeout: RpcCallTimeout, @@ -439,7 +447,6 @@ impl RpcClient { P: Serialize, T: DeserializeOwned, { - let request_id = RequestId::Integer(self.next_request_id.fetch_add(1, Ordering::SeqCst)); let (response_tx, response_rx) = oneshot::channel(); { let mut pending = self.pending.lock().await; @@ -967,6 +974,7 @@ mod tests { } #[tokio::test(flavor = "current_thread")] + #[serial_test::serial(exec_server_tracing)] async fn rpc_client_propagates_current_trace_context() { let span_exporter = InMemorySpanExporter::default(); let tracer_provider = SdkTracerProvider::builder() diff --git a/codex-rs/exec-server/src/server.rs b/codex-rs/exec-server/src/server.rs index 7988e0870b6b..39f80c3acf33 100644 --- a/codex-rs/exec-server/src/server.rs +++ b/codex-rs/exec-server/src/server.rs @@ -47,6 +47,7 @@ mod tests { use crate::ExecServerTelemetry; #[tokio::test] + #[serial_test::serial(exec_server_tracing)] async fn telemetry_entrypoint_emits_root_span() { let exporter = InMemorySpanExporter::default(); let provider = SdkTracerProvider::builder() diff --git a/codex-rs/exec-server/src/server/processor.rs b/codex-rs/exec-server/src/server/processor.rs index 63f6bd45d9a6..15343754c9ee 100644 --- a/codex-rs/exec-server/src/server/processor.rs +++ b/codex-rs/exec-server/src/server/processor.rs @@ -129,7 +129,7 @@ async fn run_connection( codex_exec_server_protocol::JSONRPCMessage::Request(request) => { let request_started_at = Instant::now(); if let Some((method, route)) = router.request_route(request.method.as_str()) { - let request_span = request_span(method, &request); + let request_span = request_span(method, &request, transport); let message = tokio::select! { message = route(Arc::clone(&handler), request).instrument(request_span.clone()) => message, _ = disconnected_rx.changed() => { @@ -160,7 +160,7 @@ async fn run_connection( drop(request_span); } else { let method = "unknown"; - let request_span = request_span(method, &request); + let request_span = request_span(method, &request, transport); if outgoing_tx .send(RpcServerOutboundMessage::Error { request_id: request.id, @@ -244,12 +244,17 @@ async fn run_connection( fn request_span( span_name: &str, request: &codex_exec_server_protocol::JSONRPCRequest, + transport: ConnectionTransport, ) -> tracing::Span { let method = request.method.as_str(); let span = tracing::info_span!( "codex.exec_server.request", otel.kind = "server", otel.name = span_name, + rpc.system = "jsonrpc", + rpc.method = method, + rpc.transport = transport.metric_tag(), + rpc.request_id = %request.id, method, result = tracing::field::Empty, ); @@ -324,6 +329,7 @@ mod tests { use crate::server::session_registry::SessionRegistry; #[test] + #[serial_test::serial(exec_server_tracing)] fn request_span_uses_bounded_name_wire_method_and_inbound_trace_parent() { let span_exporter = InMemorySpanExporter::default(); let tracer_provider = SdkTracerProvider::builder() @@ -351,7 +357,11 @@ mod tests { params: None, trace: Some(trace), }; - let request_span = request_span("unknown", &request); + let request_span = request_span( + "unknown", + &request, + crate::telemetry::ConnectionTransport::Stdio, + ); request_span.in_scope(|| {}); drop(request_span); }); diff --git a/codex-rs/exec-server/src/telemetry.rs b/codex-rs/exec-server/src/telemetry.rs index 650faacb2c63..6e83bb1d5de5 100644 --- a/codex-rs/exec-server/src/telemetry.rs +++ b/codex-rs/exec-server/src/telemetry.rs @@ -52,7 +52,7 @@ pub(crate) enum ConnectionTransport { } impl ConnectionTransport { - fn metric_tag(self) -> &'static str { + pub(crate) fn metric_tag(self) -> &'static str { match self { Self::Relay => "relay", Self::Stdio => "stdio", diff --git a/codex-rs/exec-server/src/trace_context_tests.rs b/codex-rs/exec-server/src/trace_context_tests.rs index 92a820afc43a..91f3edd8ff72 100644 --- a/codex-rs/exec-server/src/trace_context_tests.rs +++ b/codex-rs/exec-server/src/trace_context_tests.rs @@ -5,6 +5,7 @@ use tracing_subscriber::prelude::*; use super::current_trace_context_headers; #[test] +#[serial_test::serial(exec_server_tracing)] fn creates_traceparent_header_from_current_span() { let provider = SdkTracerProvider::builder().build(); let tracer = provider.tracer("exec-server-test");