diff --git a/codex-rs/exec-server/src/connection.rs b/codex-rs/exec-server/src/connection.rs index e832a529c8..eb3ffb41d6 100644 --- a/codex-rs/exec-server/src/connection.rs +++ b/codex-rs/exec-server/src/connection.rs @@ -40,8 +40,6 @@ pub(crate) const WEBSOCKET_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(30 #[derive(Debug)] pub(crate) enum JsonRpcConnectionEvent { Message(JSONRPCMessage), - // Constructed by the stacked Noise transport change. - #[allow(dead_code)] TracedMessage { message: JSONRPCMessage, trace: Option, diff --git a/codex-rs/exec-server/src/noise_relay/harness.rs b/codex-rs/exec-server/src/noise_relay/harness.rs index 9969a3e130..f0f99184ed 100644 --- a/codex-rs/exec-server/src/noise_relay/harness.rs +++ b/codex-rs/exec-server/src/noise_relay/harness.rs @@ -7,6 +7,9 @@ //! records; inbound records are reordered before decryption and reassembly. use futures::FutureExt; +use std::time::Instant; + +use codex_exec_server_protocol::JSONRPCMessage; use futures::Sink; use futures::SinkExt; use futures::Stream; @@ -15,7 +18,9 @@ use tokio::sync::mpsc; use tokio::sync::watch; use tokio_tungstenite::tungstenite::Message; use tracing::Instrument; +use tracing::Span; use tracing::debug; +use tracing::field; use tracing::info; use tracing::warn; use uuid::Uuid; @@ -36,6 +41,7 @@ use crate::noise_relay::message_framing::NOISE_RECORD_PLAINTEXT_LEN; use crate::noise_relay::message_framing::frame_jsonrpc_message; use crate::noise_relay::ordered_ciphertext::OrderedCiphertextFrames; use crate::noise_relay::take_next_sequence; +use crate::noise_relay::trace_context::NoiseTraceContext; use crate::relay::RelayFrameBodyKind; use crate::relay::decode_relay_message_frame; use crate::relay::encode_relay_message_frame; @@ -66,6 +72,13 @@ const NOISE_RELAY_RESET_DISCONNECT_REASON: &str = "Noise relay stream reset"; // Give a Pong already queued behind data a bounded chance to reach the reader. const MAX_FRAMES_DRAINED_AFTER_PONG_DEADLINE: usize = 32; +struct PendingOutboundMessage { + framed: Vec, + offset: usize, + started_at: Instant, + span: Span, +} + /// Adapt one harness rendezvous websocket into an authenticated JSON-RPC connection. /// /// The returned connection is not usable until the background task completes @@ -264,6 +277,7 @@ where let mut next_outbound_seq = 0u32; let mut inbound_ciphertexts = OrderedCiphertextFrames::default(); let mut inbound_decoder = JsonRpcMessageDecoder::default(); + let mut trace_context = NoiseTraceContext::default(); let mut keepalive = tokio::time::interval_at( tokio::time::Instant::now() + WEBSOCKET_KEEPALIVE_INTERVAL, WEBSOCKET_KEEPALIVE_INTERVAL, @@ -275,7 +289,7 @@ where // Keep one framed message as a cursor. Sending one Noise record per loop // creates a scheduling point for keepalive and inbound control frames // without splitting the WebSocket reader and writer. - let mut pending_outbound: Option<(Vec, usize)> = None; + let mut pending_outbound: Option = None; let mut force_incoming = false; let mut frames_drained_after_pong_deadline = 0usize; 'relay: loop { @@ -344,12 +358,30 @@ where let Some(message) = maybe_message else { break; }; - pending_outbound = Some(match frame_jsonrpc_message(&message) { - Ok(framed) => (framed, 0), + trace_context.observe_request(&message); + record_outbound_message_dequeued(&message); + let span = outbound_message_span(&message); + let started_at = Instant::now(); + let frame_span = + tracing::info_span!(parent: &span, "exec_server.noise.frame_request"); + let framed = match frame_span.in_scope(|| frame_jsonrpc_message(&message)) { + Ok(framed) => framed, Err(error) => { warn!("failed to frame JSON-RPC payload for Noise relay: {error}"); break; } + }; + span.record("framed_bytes", framed.len()); + span.record("records", framed.len().div_ceil(NOISE_RECORD_PLAINTEXT_LEN)); + span.record( + "frame_complete_ms", + started_at.elapsed().as_secs_f64() * 1_000.0, + ); + pending_outbound = Some(PendingOutboundMessage { + framed, + offset: 0, + started_at, + span, }); } _ = std::future::ready(()), if pending_outbound.is_some() && !force_incoming && !pong_deadline_expired => { @@ -360,19 +392,40 @@ where break 'relay; } }; - let (ciphertext, next_offset, message_complete) = { - let Some((framed, offset)) = pending_outbound.as_ref() else { + let (ciphertext, next_offset, message_complete, send_span) = { + let Some(pending) = pending_outbound.as_ref() else { continue; }; - let next_offset = (*offset + NOISE_RECORD_PLAINTEXT_LEN).min(framed.len()); - let ciphertext = match transport.encrypt(&framed[*offset..next_offset]) { + let next_offset = (pending.offset + NOISE_RECORD_PLAINTEXT_LEN) + .min(pending.framed.len()); + let encrypt_span = tracing::info_span!( + parent: &pending.span, + "exec_server.noise.encrypt_record", + ); + let ciphertext = match encrypt_span.in_scope(|| { + transport.encrypt(&pending.framed[pending.offset..next_offset]) + }) { Ok(ciphertext) => ciphertext, Err(error) => { warn!("failed to encrypt JSON-RPC payload for Noise relay: {error}"); break 'relay; } }; - (ciphertext, next_offset, next_offset == framed.len()) + pending.span.record( + "encrypt_complete_ms", + pending.started_at.elapsed().as_secs_f64() * 1_000.0, + ); + let send_span = tracing::info_span!( + parent: &pending.span, + "exec_server.noise.websocket_send", + seq, + ); + ( + ciphertext, + next_offset, + next_offset == pending.framed.len(), + send_span, + ) }; let frame = RelayMessageFrame::data(stream_id.clone(), seq, ciphertext); // A Pong can arrive after the readiness check while this write owns the @@ -384,15 +437,21 @@ where Message::Binary(encode_relay_message_frame(&frame).into()), pong_watchdog.write_deadline(tokio::time::Instant::now()), ) + .instrument(send_span) .await { warn!("failed to write Noise relay websocket: {error}"); break 'relay; } if message_complete { - pending_outbound = None; - } else if let Some((_framed, offset)) = pending_outbound.as_mut() { - *offset = next_offset; + if let Some(pending) = pending_outbound.take() { + pending.span.record( + "send_complete_ms", + pending.started_at.elapsed().as_secs_f64() * 1_000.0, + ); + } + } else if let Some(pending) = pending_outbound.as_mut() { + pending.offset = next_offset; } } _ = &mut pong_deadline, if pong_watchdog.deadline().is_some() && !force_incoming => { @@ -427,6 +486,7 @@ where } match incoming_message { Ok(Message::Binary(payload)) => { + let frame_received_at = Instant::now(); let frame = match decode_relay_message_frame(payload.as_ref()) { Ok(frame) => frame, Err(error) => { @@ -450,9 +510,14 @@ where &mut inbound_ciphertexts, &mut transport, &mut inbound_decoder, - data, - pong_watchdog.write_deadline(tokio::time::Instant::now()), + InboundRelayFrame { + data, + delivery_deadline: pong_watchdog + .write_deadline(tokio::time::Instant::now()), + received_at: frame_received_at, + }, &incoming_tx, + &mut trace_context, ) .await { @@ -561,46 +626,150 @@ where Ok(()) } +fn record_outbound_message_dequeued(message: &JSONRPCMessage) { + let (message_kind, method, trace) = match message { + JSONRPCMessage::Request(request) => { + ("request", request.method.as_str(), request.trace.as_ref()) + } + JSONRPCMessage::Notification(notification) => { + ("notification", notification.method.as_str(), None) + } + JSONRPCMessage::Response(_) => ("response", "", None), + JSONRPCMessage::Error(_) => ("error", "", None), + }; + let span = tracing::info_span!( + "exec_server.noise.harness_outbound_dequeued", + message_kind, + method, + ); + if let Some(trace) = trace { + let _ = codex_otel::set_parent_from_w3c_trace_context(&span, trace); + } + span.in_scope(|| {}); +} + +fn outbound_message_span(message: &JSONRPCMessage) -> Span { + let (message_kind, method, trace) = match message { + JSONRPCMessage::Request(request) => { + ("request", request.method.as_str(), request.trace.as_ref()) + } + JSONRPCMessage::Notification(notification) => { + ("notification", notification.method.as_str(), None) + } + JSONRPCMessage::Response(_) => ("response", "", None), + JSONRPCMessage::Error(_) => ("error", "", None), + }; + let span = tracing::info_span!( + "exec_server.noise.harness_outbound", + message_kind, + method, + framed_bytes = field::Empty, + records = field::Empty, + frame_complete_ms = field::Empty, + encrypt_complete_ms = field::Empty, + send_complete_ms = field::Empty, + ); + if let Some(trace) = trace { + let _ = codex_otel::set_parent_from_w3c_trace_context(&span, trace); + } + span +} + /// Order and decrypt one relay frame, then emit any complete JSON-RPC messages. /// Relay records and JSON-RPC messages do not share boundaries, so reassembly /// happens after decryption. +struct InboundRelayFrame { + data: RelayData, + delivery_deadline: tokio::time::Instant, + received_at: Instant, +} + async fn receive_data( inbound_ciphertexts: &mut OrderedCiphertextFrames, transport: &mut NoiseTransport, decoder: &mut JsonRpcMessageDecoder, - data: RelayData, - delivery_deadline: tokio::time::Instant, + frame: InboundRelayFrame, incoming_tx: &mpsc::Sender, + trace_context: &mut NoiseTraceContext, ) -> Result<(), ExecServerError> { + let InboundRelayFrame { + data, + delivery_deadline, + received_at: frame_received_at, + } = frame; + let frame_decode_elapsed = frame_received_at.elapsed(); + let ciphertext_bytes = data.payload.len(); // Ordering must happen before decryption because Noise transport nonces are // implicit. A future or duplicate ciphertext passed directly to Clatter // would desynchronize the channel. - for ciphertext in inbound_ciphertexts.push(data.seq, data.payload)? { + let reorder_started_at = Instant::now(); + let ciphertexts = inbound_ciphertexts.push(data.seq, data.payload)?; + let reorder_elapsed = reorder_started_at.elapsed(); + for ciphertext in ciphertexts { + let decrypt_started_at = Instant::now(); let plaintext = transport.decrypt(&ciphertext).map_err(|error| { ExecServerError::Protocol(format!("Noise relay decryption failed: {error}")) })?; + let decrypt_elapsed = decrypt_started_at.elapsed(); // The authenticated byte stream can carry partial or multiple JSON-RPC // messages; emit only complete, successfully parsed messages. - for message in decoder.push(&plaintext)? { - send_incoming_event( + let decode_started_at = Instant::now(); + let messages = decoder.push(&plaintext)?; + let decode_elapsed = decode_started_at.elapsed(); + for message in messages { + let trace = trace_context.return_trace(&message); + let inbound_span = inbound_message_span( + &message, + trace.as_ref(), + ciphertext_bytes, + frame_decode_elapsed, + reorder_elapsed, + decrypt_elapsed, + decode_elapsed, + ); + let delivery_started_at = Instant::now(); + let send_result = send_incoming_event( incoming_tx, - JsonRpcConnectionEvent::Message(message), + || JsonRpcConnectionEvent::TracedMessage { + message, + trace, + queued_at: Instant::now(), + }, delivery_deadline, ) - .await?; + .instrument(inbound_span.clone()) + .await; + inbound_span.record( + "delivery_ms", + delivery_started_at.elapsed().as_secs_f64() * 1_000.0, + ); + inbound_span.record( + "receive_complete_ms", + frame_received_at.elapsed().as_secs_f64() * 1_000.0, + ); + send_result?; } } Ok(()) } -async fn send_incoming_event( +async fn send_incoming_event( incoming_tx: &mpsc::Sender, - event: JsonRpcConnectionEvent, + event: F, deadline: tokio::time::Instant, -) -> Result<(), ExecServerError> { - match tokio::time::timeout_at(deadline, incoming_tx.send(event)).await { - Ok(Ok(())) => Ok(()), +) -> Result<(), ExecServerError> +where + F: FnOnce() -> JsonRpcConnectionEvent, +{ + let reserve = incoming_tx.reserve().instrument(tracing::info_span!( + "exec_server.noise.harness_inbound_delivery" + )); + match tokio::time::timeout_at(deadline, reserve).await { + Ok(Ok(permit)) => { + permit.send(event()); + Ok(()) + } Ok(Err(_)) => Err(ExecServerError::Closed), Err(_) => { warn!( @@ -612,6 +781,41 @@ async fn send_incoming_event( } } +fn inbound_message_span( + message: &JSONRPCMessage, + trace: Option<&codex_protocol::protocol::W3cTraceContext>, + ciphertext_bytes: usize, + frame_decode_elapsed: std::time::Duration, + reorder_elapsed: std::time::Duration, + decrypt_elapsed: std::time::Duration, + decode_elapsed: std::time::Duration, +) -> Span { + let (message_kind, method) = match message { + JSONRPCMessage::Request(request) => ("request", request.method.as_str()), + JSONRPCMessage::Notification(notification) => { + ("notification", notification.method.as_str()) + } + JSONRPCMessage::Response(_) => ("response", ""), + JSONRPCMessage::Error(_) => ("error", ""), + }; + let span = tracing::info_span!( + "exec_server.noise.harness_inbound", + message_kind, + method, + ciphertext_bytes, + frame_decode_ms = frame_decode_elapsed.as_secs_f64() * 1_000.0, + reorder_ms = reorder_elapsed.as_secs_f64() * 1_000.0, + decrypt_ms = decrypt_elapsed.as_secs_f64() * 1_000.0, + decode_ms = decode_elapsed.as_secs_f64() * 1_000.0, + delivery_ms = field::Empty, + receive_complete_ms = field::Empty, + ); + if let Some(trace) = trace { + let _ = codex_otel::set_parent_from_w3c_trace_context(&span, trace); + } + span +} + fn send_malformed(incoming_tx: &mpsc::Sender, reason: String) { let _ = incoming_tx.try_send(JsonRpcConnectionEvent::MalformedMessage { reason }); } diff --git a/codex-rs/exec-server/src/noise_relay/harness_tests.rs b/codex-rs/exec-server/src/noise_relay/harness_tests.rs index 348351e61c..7b6cbabbbf 100644 --- a/codex-rs/exec-server/src/noise_relay/harness_tests.rs +++ b/codex-rs/exec-server/src/noise_relay/harness_tests.rs @@ -226,7 +226,7 @@ async fn application_event_delivery_is_bounded() -> Result<()> { let result = send_incoming_event( &incoming_tx, - JsonRpcConnectionEvent::MalformedMessage { + || JsonRpcConnectionEvent::MalformedMessage { reason: "blocked event".to_string(), }, Instant::now() + Duration::from_millis(10),