From bd9c3a47d6d48be7442ec8a280a046c8a9903e9f Mon Sep 17 00:00:00 2001 From: Simo Lin Date: Sun, 14 Dec 2025 20:54:43 -0800 Subject: [PATCH] [model-gateway] Add streaming metrics for harmony gRPC router (#15147) --- .../src/observability/metrics.rs | 1 + .../src/routers/grpc/harmony/streaming.rs | 58 +++++++++++++++++-- 2 files changed, 53 insertions(+), 6 deletions(-) diff --git a/sgl-model-gateway/src/observability/metrics.rs b/sgl-model-gateway/src/observability/metrics.rs index 4fb170994..b33b3dff1 100644 --- a/sgl-model-gateway/src/observability/metrics.rs +++ b/sgl-model-gateway/src/observability/metrics.rs @@ -756,6 +756,7 @@ pub mod smg_labels { pub const BACKEND_REGULAR: &str = "regular"; pub const BACKEND_PD: &str = "pd"; pub const BACKEND_EXTERNAL: &str = "external"; + pub const BACKEND_HARMONY: &str = "harmony"; // Connection modes pub const CONNECTION_HTTP: &str = "http"; diff --git a/sgl-model-gateway/src/routers/grpc/harmony/streaming.rs b/sgl-model-gateway/src/routers/grpc/harmony/streaming.rs index afcf7dd54..174ae9d14 100644 --- a/sgl-model-gateway/src/routers/grpc/harmony/streaming.rs +++ b/sgl-model-gateway/src/routers/grpc/harmony/streaming.rs @@ -4,6 +4,7 @@ use std::{ collections::{hash_map::Entry::Vacant, HashMap}, io, sync::Arc, + time::Instant, }; use axum::{body::Body, http::StatusCode, response::Response}; @@ -19,6 +20,7 @@ use super::{ }; use crate::{ grpc_client::sglang_proto::generate_complete::MatchedStop::{MatchedStopStr, MatchedTokenId}, + observability::metrics::{smg_labels, SmgMetrics, StreamingMetricsParams}, protocols::{ chat::{ ChatCompletionRequest, ChatCompletionStreamResponse, ChatMessageDelta, ChatStreamChoice, @@ -183,6 +185,10 @@ impl HarmonyStreamingProcessor { original_request: Arc, tx: &mpsc::UnboundedSender>, ) -> Result<(), String> { + // Timing for metrics + let start_time = Instant::now(); + let mut first_token_time: Option = None; + // Per-index state management (for n>1 support) let mut parsers: HashMap = HashMap::new(); let mut is_firsts: HashMap = HashMap::new(); @@ -202,6 +208,11 @@ impl HarmonyStreamingProcessor { let chunk = chunk_wrapper.as_sglang(); let index = chunk.index; + // Track first token time for TTFT metric + if first_token_time.is_none() { + first_token_time = Some(Instant::now()); + } + // Initialize parser for this index if needed if let Vacant(e) = parsers.entry(index) { e.insert( @@ -285,11 +296,12 @@ impl HarmonyStreamingProcessor { } } + // Compute totals once for both usage chunk and metrics + let total_prompt: u32 = prompt_tokens.values().sum(); + let total_completion: u32 = completion_tokens.values().sum(); + // Emit final usage if requested if let Some(true) = stream_options.as_ref().and_then(|so| so.include_usage) { - let total_prompt: u32 = prompt_tokens.values().sum(); - let total_completion: u32 = completion_tokens.values().sum(); - Self::emit_usage_chunk( total_prompt, total_completion, @@ -302,6 +314,18 @@ impl HarmonyStreamingProcessor { // Mark stream as completed successfully to prevent abort on drop grpc_stream.mark_completed(); + // Record streaming metrics + SmgMetrics::record_streaming_metrics(StreamingMetricsParams { + router_type: smg_labels::ROUTER_GRPC, + backend_type: smg_labels::BACKEND_HARMONY, + model_id: &original_request.model, + endpoint: smg_labels::ENDPOINT_CHAT, + ttft: first_token_time.map(|t| t.duration_since(start_time)), + generation_duration: start_time.elapsed(), + input_tokens: Some(total_prompt as u64), + output_tokens: total_completion as u64, + }); + Ok(()) } @@ -313,6 +337,10 @@ impl HarmonyStreamingProcessor { original_request: Arc, tx: &mpsc::UnboundedSender>, ) -> Result<(), String> { + // Timing for metrics + let start_time = Instant::now(); + let mut first_token_time: Option = None; + // Phase 1: Process prefill stream (collect metadata) let mut prompt_tokens: HashMap = HashMap::new(); @@ -342,6 +370,11 @@ impl HarmonyStreamingProcessor { let chunk = chunk_wrapper.as_sglang(); let index = chunk.index; + // Track first token time for TTFT metric + if first_token_time.is_none() { + first_token_time = Some(Instant::now()); + } + // Initialize parser for this index if needed if let Vacant(e) = parsers.entry(index) { e.insert( @@ -425,11 +458,12 @@ impl HarmonyStreamingProcessor { // This ensures that if client disconnects during decode, BOTH streams send abort prefill_stream.mark_completed(); + // Compute totals once for both usage chunk and metrics + let total_prompt: u32 = prompt_tokens.values().sum(); + let total_completion: u32 = completion_tokens.values().sum(); + // Emit final usage if requested if let Some(true) = stream_options.as_ref().and_then(|so| so.include_usage) { - let total_prompt: u32 = prompt_tokens.values().sum(); - let total_completion: u32 = completion_tokens.values().sum(); - Self::emit_usage_chunk( total_prompt, total_completion, @@ -439,6 +473,18 @@ impl HarmonyStreamingProcessor { )?; } + // Record streaming metrics + SmgMetrics::record_streaming_metrics(StreamingMetricsParams { + router_type: smg_labels::ROUTER_GRPC, + backend_type: smg_labels::BACKEND_HARMONY, + model_id: &original_request.model, + endpoint: smg_labels::ENDPOINT_CHAT, + ttft: first_token_time.map(|t| t.duration_since(start_time)), + generation_duration: start_time.elapsed(), + input_tokens: Some(total_prompt as u64), + output_tokens: total_completion as u64, + }); + Ok(()) }