mirror of
https://github.com/pchuan98/codex.git
synced 2026-07-01 00:31:56 +08:00
[codex] Trace exec-server JSON-RPC requests (#27466)
## Why Exec-server JSON-RPC calls can cross local and remote transports, but trace context stopped at the RPC boundary. That made client and server work difficult to correlate when diagnosing latency or failures. ## What changed - Propagate the current W3C trace context on outbound JSON-RPC requests. - Parent inbound request spans from received trace context. - Record the received JSON-RPC method on server spans and keep each span open through response enqueue. - Add only the OTEL dependencies required by the exec-server crate. ## Stack Review and land this stack in order: 1. #27466 — trace exec-server JSON-RPC requests **(this PR)** 2. #27467 — record bounded connection, request, and process lifecycle metrics 3. #27470 — observe remote registration and Noise rendezvous lifecycle ## Validation - `just test -p codex-exec-server --lib` (153 passed) - `just bazel-lock-check` - `just fix -p codex-exec-server`
This commit is contained in:
committed by
GitHub
Unverified
parent
4907f0c2c3
commit
74dcce594d
@@ -1,6 +1,7 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use tokio::sync::mpsc;
|
||||
use tracing::Instrument;
|
||||
use tracing::debug;
|
||||
use tracing::warn;
|
||||
|
||||
@@ -101,30 +102,40 @@ async fn run_connection(
|
||||
JsonRpcConnectionEvent::Message(message) => match message {
|
||||
codex_exec_server_protocol::JSONRPCMessage::Request(request) => {
|
||||
if let Some(route) = router.request_route(request.method.as_str()) {
|
||||
let request_span = request_span(request.method.as_str(), &request);
|
||||
let message = tokio::select! {
|
||||
message = route(Arc::clone(&handler), request) => message,
|
||||
message = route(Arc::clone(&handler), request).instrument(request_span.clone()) => message,
|
||||
_ = disconnected_rx.changed() => {
|
||||
request_span.record("result", "disconnected");
|
||||
debug!("exec-server transport disconnected while handling request");
|
||||
break;
|
||||
}
|
||||
};
|
||||
let result = request_result(&message);
|
||||
if let Some(message) = message
|
||||
&& outgoing_tx.send(message).await.is_err()
|
||||
{
|
||||
request_span.record("result", "disconnected");
|
||||
break;
|
||||
}
|
||||
} else if outgoing_tx
|
||||
.send(RpcServerOutboundMessage::Error {
|
||||
request_id: request.id,
|
||||
error: method_not_found(format!(
|
||||
"exec-server stub does not implement `{}` yet",
|
||||
request.method
|
||||
)),
|
||||
})
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
break;
|
||||
request_span.record("result", result);
|
||||
} else {
|
||||
let request_span = request_span("unknown", &request);
|
||||
if outgoing_tx
|
||||
.send(RpcServerOutboundMessage::Error {
|
||||
request_id: request.id,
|
||||
error: method_not_found(format!(
|
||||
"exec-server stub does not implement `{}` yet",
|
||||
request.method
|
||||
)),
|
||||
})
|
||||
.await
|
||||
.is_err()
|
||||
{
|
||||
request_span.record("result", "disconnected");
|
||||
break;
|
||||
}
|
||||
request_span.record("result", "error");
|
||||
}
|
||||
}
|
||||
codex_exec_server_protocol::JSONRPCMessage::Notification(notification) => {
|
||||
@@ -184,6 +195,36 @@ async fn run_connection(
|
||||
let _ = outbound_task.await;
|
||||
}
|
||||
|
||||
fn request_span(
|
||||
span_name: &str,
|
||||
request: &codex_exec_server_protocol::JSONRPCRequest,
|
||||
) -> tracing::Span {
|
||||
let method = request.method.as_str();
|
||||
let span = tracing::info_span!(
|
||||
"codex.exec_server.request",
|
||||
otel.kind = "server",
|
||||
otel.name = span_name,
|
||||
method,
|
||||
result = tracing::field::Empty,
|
||||
);
|
||||
if let Some(trace) = &request.trace
|
||||
&& !codex_otel::set_parent_from_w3c_trace_context(&span, trace)
|
||||
{
|
||||
warn!(method, "ignoring invalid inbound exec-server trace carrier");
|
||||
}
|
||||
span
|
||||
}
|
||||
|
||||
fn request_result(message: &Option<RpcServerOutboundMessage>) -> &'static str {
|
||||
match message {
|
||||
Some(RpcServerOutboundMessage::Error { .. }) => "error",
|
||||
Some(
|
||||
RpcServerOutboundMessage::Response { .. } | RpcServerOutboundMessage::Notification(_),
|
||||
)
|
||||
| None => "success",
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::collections::HashMap;
|
||||
@@ -196,6 +237,11 @@ mod tests {
|
||||
use codex_exec_server_protocol::JSONRPCResponse;
|
||||
use codex_exec_server_protocol::RequestId;
|
||||
use codex_utils_path_uri::PathUri;
|
||||
use opentelemetry::trace::SpanId;
|
||||
use opentelemetry::trace::TraceId;
|
||||
use opentelemetry::trace::TracerProvider as _;
|
||||
use opentelemetry_sdk::trace::InMemorySpanExporter;
|
||||
use opentelemetry_sdk::trace::SdkTracerProvider;
|
||||
use pretty_assertions::assert_eq;
|
||||
use serde::Serialize;
|
||||
use serde::de::DeserializeOwned;
|
||||
@@ -207,7 +253,10 @@ mod tests {
|
||||
use tokio::io::duplex;
|
||||
use tokio::task::JoinHandle;
|
||||
use tokio::time::timeout;
|
||||
use tracing_subscriber::filter::filter_fn;
|
||||
use tracing_subscriber::prelude::*;
|
||||
|
||||
use super::request_span;
|
||||
use super::run_connection;
|
||||
use crate::ExecServerRuntimePaths;
|
||||
use crate::ProcessId;
|
||||
@@ -228,6 +277,57 @@ mod tests {
|
||||
use crate::protocol::TerminateResponse;
|
||||
use crate::server::session_registry::SessionRegistry;
|
||||
|
||||
#[test]
|
||||
fn request_span_uses_bounded_name_wire_method_and_inbound_trace_parent() {
|
||||
let span_exporter = InMemorySpanExporter::default();
|
||||
let tracer_provider = SdkTracerProvider::builder()
|
||||
.with_simple_exporter(span_exporter.clone())
|
||||
.build();
|
||||
let tracer = tracer_provider.tracer("exec-server-test");
|
||||
let subscriber = tracing_subscriber::registry().with(
|
||||
tracing_opentelemetry::layer()
|
||||
.with_tracer(tracer)
|
||||
.with_filter(filter_fn(codex_otel::OtelProvider::trace_export_filter)),
|
||||
);
|
||||
let trace_id = TraceId::from_hex("00000000000000000000000000000001").expect("trace id");
|
||||
let parent_span_id = SpanId::from_hex("0000000000000002").expect("span id");
|
||||
let trace = codex_protocol::protocol::W3cTraceContext {
|
||||
traceparent: Some(format!("00-{trace_id}-{parent_span_id}-01")),
|
||||
tracestate: None,
|
||||
};
|
||||
|
||||
let method = "custom/method";
|
||||
tracing::subscriber::with_default(subscriber, || {
|
||||
tracing::callsite::rebuild_interest_cache();
|
||||
let request = JSONRPCRequest {
|
||||
id: RequestId::Integer(1),
|
||||
method: method.to_string(),
|
||||
params: None,
|
||||
trace: Some(trace),
|
||||
};
|
||||
let request_span = request_span("unknown", &request);
|
||||
request_span.in_scope(|| {});
|
||||
drop(request_span);
|
||||
});
|
||||
|
||||
tracer_provider.force_flush().expect("flush traces");
|
||||
let spans = span_exporter.get_finished_spans().expect("span export");
|
||||
let request_span = spans
|
||||
.iter()
|
||||
.find(|span| span.name.as_ref() == "unknown")
|
||||
.expect("unknown method span");
|
||||
assert_eq!(
|
||||
request_span
|
||||
.attributes
|
||||
.iter()
|
||||
.find(|attribute| attribute.key.as_str() == "method")
|
||||
.map(|attribute| attribute.value.clone()),
|
||||
Some(opentelemetry::Value::String(method.into()))
|
||||
);
|
||||
assert_eq!(request_span.span_context.trace_id(), trace_id);
|
||||
assert_eq!(request_span.parent_span_id, parent_span_id);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn connection_accepts_pipelined_scalar_requests() {
|
||||
let registry = SessionRegistry::new();
|
||||
@@ -382,6 +482,7 @@ mod tests {
|
||||
id: RequestId::Integer(id),
|
||||
method: method.to_string(),
|
||||
params: Some(serde_json::to_value(params).expect("serialize params")),
|
||||
trace: None,
|
||||
}),
|
||||
)
|
||||
.await;
|
||||
|
||||
@@ -74,6 +74,7 @@ async fn stdio_listen_transport_serves_initialize() {
|
||||
})
|
||||
.expect("initialize params should serialize"),
|
||||
),
|
||||
trace: None,
|
||||
});
|
||||
write_jsonrpc_line(&mut client_writer, &initialize).await;
|
||||
|
||||
|
||||
Reference in New Issue
Block a user