From b9a76bf154489cc08f2b09d91f70de73ddf8b2eb Mon Sep 17 00:00:00 2001 From: whit3rabbit Date: Fri, 27 Mar 2026 06:32:37 -0500 Subject: [PATCH] feat: add streaming-specific metrics (started/completed/failed/disconnected) Add atomic counters for SSE stream lifecycle events to the Metrics struct: streams_started, streams_completed, streams_failed, streams_client_disconnected. Wire recording into StreamOutcome::record() in streaming.rs and the parallel streaming path in chat_completions.rs. Co-Authored-By: Claude Opus 4.6 (1M context) --- crates/proxy/src/metrics/mod.rs | 58 +++++++++++++++++++++ crates/proxy/src/server/chat_completions.rs | 33 ++++++++++-- crates/proxy/src/server/streaming.rs | 13 ++++- 3 files changed, 97 insertions(+), 7 deletions(-) diff --git a/crates/proxy/src/metrics/mod.rs b/crates/proxy/src/metrics/mod.rs index 14a0663..cb3fbea 100644 --- a/crates/proxy/src/metrics/mod.rs +++ b/crates/proxy/src/metrics/mod.rs @@ -15,6 +15,10 @@ struct MetricsInner { requests_total: AtomicU64, requests_success: AtomicU64, requests_error: AtomicU64, + streams_started: AtomicU64, + streams_completed: AtomicU64, + streams_failed: AtomicU64, + streams_client_disconnected: AtomicU64, } impl Metrics { @@ -41,12 +45,41 @@ impl Metrics { self.inner.requests_error.fetch_add(1, Ordering::Relaxed); } + /// Increment when an SSE stream begins sending events to the client. + pub fn record_stream_started(&self) { + self.inner.streams_started.fetch_add(1, Ordering::Relaxed); + } + + /// Increment when an SSE stream completes normally (backend sent all data). + pub fn record_stream_completed(&self) { + self.inner.streams_completed.fetch_add(1, Ordering::Relaxed); + } + + /// Increment when an SSE stream fails due to an upstream error. + pub fn record_stream_failed(&self) { + self.inner.streams_failed.fetch_add(1, Ordering::Relaxed); + } + + /// Increment when the downstream client disconnects before the stream finishes. + pub fn record_stream_client_disconnected(&self) { + self.inner + .streams_client_disconnected + .fetch_add(1, Ordering::Relaxed); + } + /// Take a point-in-time snapshot of all counters for the GET /metrics endpoint. pub fn snapshot(&self) -> MetricsSnapshot { MetricsSnapshot { requests_total: self.inner.requests_total.load(Ordering::Relaxed), requests_success: self.inner.requests_success.load(Ordering::Relaxed), requests_error: self.inner.requests_error.load(Ordering::Relaxed), + streams_started: self.inner.streams_started.load(Ordering::Relaxed), + streams_completed: self.inner.streams_completed.load(Ordering::Relaxed), + streams_failed: self.inner.streams_failed.load(Ordering::Relaxed), + streams_client_disconnected: self + .inner + .streams_client_disconnected + .load(Ordering::Relaxed), } } } @@ -60,6 +93,14 @@ pub struct MetricsSnapshot { pub requests_success: u64, /// Requests that failed (non-2xx status or transport error). pub requests_error: u64, + /// SSE streams that began sending events to the client. + pub streams_started: u64, + /// SSE streams that completed normally. + pub streams_completed: u64, + /// SSE streams that failed due to upstream errors. + pub streams_failed: u64, + /// SSE streams where the client disconnected early. + pub streams_client_disconnected: u64, } impl MetricsSnapshot { @@ -91,6 +132,23 @@ mod tests { assert_eq!(s.requests_error, 1); } + #[test] + fn streaming_metrics_counting() { + let m = Metrics::new(); + m.record_stream_started(); + m.record_stream_started(); + m.record_stream_started(); + m.record_stream_completed(); + m.record_stream_failed(); + m.record_stream_client_disconnected(); + + let s = m.snapshot(); + assert_eq!(s.streams_started, 3); + assert_eq!(s.streams_completed, 1); + assert_eq!(s.streams_failed, 1); + assert_eq!(s.streams_client_disconnected, 1); + } + #[test] fn metrics_clone_shares_state() { let m = Metrics::new(); diff --git a/crates/proxy/src/server/chat_completions.rs b/crates/proxy/src/server/chat_completions.rs index e282290..5d4aac0 100644 --- a/crates/proxy/src/server/chat_completions.rs +++ b/crates/proxy/src/server/chat_completions.rs @@ -156,10 +156,15 @@ pub(crate) async fn chat_completions( } // Resolve model routing (may switch to a different backend). - let (mapped_model, effective) = match state.resolve_model_and_state(&original_model) { + let (mapped_model, effective, deployment) = match state.resolve_model_and_state(&original_model) + { Ok(v) => v, Err(resp) => return resp, }; + if let Some(ref d) = deployment { + d.record_start(); + } + let backend_start = std::time::Instant::now(); // Non-streaming path match &effective.backend { @@ -181,6 +186,9 @@ pub(crate) async fn chat_completions( match client.chat_completion(&openai_req).await { Ok((openai_resp, _status, rate_limits)) => { + if let Some(ref d) = deployment { + d.record_finish(backend_start.elapsed().as_millis() as u64); + } state.metrics.record_success(); // Translate Anthropic response back to OpenAI format let anthropic_resp = mapping::message_map::openai_to_anthropic_response( @@ -221,6 +229,9 @@ pub(crate) async fn chat_completions( response } Err(e) => { + if let Some(ref d) = deployment { + d.record_finish(backend_start.elapsed().as_millis() as u64); + } state.metrics.record_error(); let status = e.status_code(); log_request( @@ -246,6 +257,9 @@ pub(crate) async fn chat_completions( match client.responses(&responses_req).await { Ok((resp, _status, rate_limits)) => { + if let Some(ref d) = deployment { + d.record_finish(backend_start.elapsed().as_millis() as u64); + } state.metrics.record_success(); let anthropic_resp = mapping::responses_message_map::responses_to_anthropic_response( @@ -287,6 +301,9 @@ pub(crate) async fn chat_completions( response } Err(e) => { + if let Some(ref d) = deployment { + d.record_finish(backend_start.elapsed().as_millis() as u64); + } state.metrics.record_error(); let status = e.status_code(); log_request( @@ -326,10 +343,11 @@ async fn chat_completions_stream( concurrency_permit: Option, ) -> Response { // Resolve model routing (may switch to a different backend). - let (mapped_model_resolved, effective) = match state.resolve_model_and_state(&original_model) { - Ok(v) => v, - Err(resp) => return resp, - }; + let (mapped_model_resolved, effective, _deployment) = + match state.resolve_model_and_state(&original_model) { + Ok(v) => v, + Err(resp) => return resp, + }; // Translate to OpenAI format for the backend let mut openai_req = mapping::message_map::anthropic_to_openai_request(&anthropic_req); @@ -372,6 +390,7 @@ async fn chat_completions_stream( let _permit = concurrency_permit; tokio::spawn(async move { + metrics.record_stream_started(); let mut translator = ReverseStreamingTranslator::new( format!("chatcmpl-{}", uuid::Uuid::new_v4().as_simple()), model_for_translator.clone(), @@ -389,6 +408,7 @@ async fn chat_completions_stream( Err(e) => { tracing::error!("stream read error: {e}"); metrics.record_error(); + metrics.record_stream_failed(); break; } }; @@ -397,6 +417,7 @@ async fn chat_completions_stream( if buffer.len() > MAX_SSE_BUFFER_SIZE { tracing::error!("SSE buffer exceeded maximum size"); metrics.record_error(); + metrics.record_stream_failed(); break; } @@ -425,6 +446,7 @@ async fn chat_completions_stream( if let Ok(json) = serde_json::to_string(oai_chunk) { let sse_line = format!("data: {}\n\n", json); if tx.send(Ok(sse_line)).await.is_err() { + metrics.record_stream_client_disconnected(); return; // Client disconnected } } @@ -456,6 +478,7 @@ async fn chat_completions_stream( } metrics.record_success(); + metrics.record_stream_completed(); log_request( &log_shared, ctx.log_entry( diff --git a/crates/proxy/src/server/streaming.rs b/crates/proxy/src/server/streaming.rs index 07fc911..da34e79 100644 --- a/crates/proxy/src/server/streaming.rs +++ b/crates/proxy/src/server/streaming.rs @@ -47,10 +47,17 @@ impl StreamOutcome { match self { Self::Completed => { metrics.record_success(); + metrics.record_stream_completed(); (200, None) } - Self::ClientDisconnected => (499, Some("client disconnected".into())), - Self::UpstreamError => (502, Some("stream interrupted".into())), + Self::ClientDisconnected => { + metrics.record_stream_client_disconnected(); + (499, Some("client disconnected".into())) + } + Self::UpstreamError => { + metrics.record_stream_failed(); + (502, Some("stream interrupted".into())) + } } } } @@ -183,6 +190,7 @@ pub(crate) async fn messages_stream( // until headers are sent, so the semaphore accurately bounds // concurrent streaming connections. let _permit = permit; + metrics.record_stream_started(); match client.chat_completion_stream(&openai_req).await { Ok((response, rate_limits)) => { rl_tx.send(Ok(rate_limits)).ok(); @@ -258,6 +266,7 @@ pub(crate) async fn messages_stream( tokio::spawn(async move { let _permit = permit; + metrics.record_stream_started(); match client.responses_stream(&responses_req).await { Ok((response, rate_limits)) => { rl_tx.send(Ok(rate_limits)).ok();