mirror of
https://github.com/whit3rabbit/anyllm-proxy.git
synced 2026-09-22 08:00:51 +00:00
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) <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.6
parent
6a02944cc9
commit
b9a76bf154
@@ -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();
|
||||
|
||||
@@ -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<ConcurrencyPermit>,
|
||||
) -> 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(
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user