[model-gateway] Add streaming metrics for harmony gRPC router (#15147)
This commit is contained in:
@@ -756,6 +756,7 @@ pub mod smg_labels {
|
|||||||
pub const BACKEND_REGULAR: &str = "regular";
|
pub const BACKEND_REGULAR: &str = "regular";
|
||||||
pub const BACKEND_PD: &str = "pd";
|
pub const BACKEND_PD: &str = "pd";
|
||||||
pub const BACKEND_EXTERNAL: &str = "external";
|
pub const BACKEND_EXTERNAL: &str = "external";
|
||||||
|
pub const BACKEND_HARMONY: &str = "harmony";
|
||||||
|
|
||||||
// Connection modes
|
// Connection modes
|
||||||
pub const CONNECTION_HTTP: &str = "http";
|
pub const CONNECTION_HTTP: &str = "http";
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ use std::{
|
|||||||
collections::{hash_map::Entry::Vacant, HashMap},
|
collections::{hash_map::Entry::Vacant, HashMap},
|
||||||
io,
|
io,
|
||||||
sync::Arc,
|
sync::Arc,
|
||||||
|
time::Instant,
|
||||||
};
|
};
|
||||||
|
|
||||||
use axum::{body::Body, http::StatusCode, response::Response};
|
use axum::{body::Body, http::StatusCode, response::Response};
|
||||||
@@ -19,6 +20,7 @@ use super::{
|
|||||||
};
|
};
|
||||||
use crate::{
|
use crate::{
|
||||||
grpc_client::sglang_proto::generate_complete::MatchedStop::{MatchedStopStr, MatchedTokenId},
|
grpc_client::sglang_proto::generate_complete::MatchedStop::{MatchedStopStr, MatchedTokenId},
|
||||||
|
observability::metrics::{smg_labels, SmgMetrics, StreamingMetricsParams},
|
||||||
protocols::{
|
protocols::{
|
||||||
chat::{
|
chat::{
|
||||||
ChatCompletionRequest, ChatCompletionStreamResponse, ChatMessageDelta, ChatStreamChoice,
|
ChatCompletionRequest, ChatCompletionStreamResponse, ChatMessageDelta, ChatStreamChoice,
|
||||||
@@ -183,6 +185,10 @@ impl HarmonyStreamingProcessor {
|
|||||||
original_request: Arc<ChatCompletionRequest>,
|
original_request: Arc<ChatCompletionRequest>,
|
||||||
tx: &mpsc::UnboundedSender<Result<Bytes, io::Error>>,
|
tx: &mpsc::UnboundedSender<Result<Bytes, io::Error>>,
|
||||||
) -> Result<(), String> {
|
) -> Result<(), String> {
|
||||||
|
// Timing for metrics
|
||||||
|
let start_time = Instant::now();
|
||||||
|
let mut first_token_time: Option<Instant> = None;
|
||||||
|
|
||||||
// Per-index state management (for n>1 support)
|
// Per-index state management (for n>1 support)
|
||||||
let mut parsers: HashMap<u32, HarmonyParserAdapter> = HashMap::new();
|
let mut parsers: HashMap<u32, HarmonyParserAdapter> = HashMap::new();
|
||||||
let mut is_firsts: HashMap<u32, bool> = HashMap::new();
|
let mut is_firsts: HashMap<u32, bool> = HashMap::new();
|
||||||
@@ -202,6 +208,11 @@ impl HarmonyStreamingProcessor {
|
|||||||
let chunk = chunk_wrapper.as_sglang();
|
let chunk = chunk_wrapper.as_sglang();
|
||||||
let index = chunk.index;
|
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
|
// Initialize parser for this index if needed
|
||||||
if let Vacant(e) = parsers.entry(index) {
|
if let Vacant(e) = parsers.entry(index) {
|
||||||
e.insert(
|
e.insert(
|
||||||
@@ -285,11 +296,12 @@ impl HarmonyStreamingProcessor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Emit final usage if requested
|
// Compute totals once for both usage chunk and metrics
|
||||||
if let Some(true) = stream_options.as_ref().and_then(|so| so.include_usage) {
|
|
||||||
let total_prompt: u32 = prompt_tokens.values().sum();
|
let total_prompt: u32 = prompt_tokens.values().sum();
|
||||||
let total_completion: u32 = completion_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) {
|
||||||
Self::emit_usage_chunk(
|
Self::emit_usage_chunk(
|
||||||
total_prompt,
|
total_prompt,
|
||||||
total_completion,
|
total_completion,
|
||||||
@@ -302,6 +314,18 @@ impl HarmonyStreamingProcessor {
|
|||||||
// Mark stream as completed successfully to prevent abort on drop
|
// Mark stream as completed successfully to prevent abort on drop
|
||||||
grpc_stream.mark_completed();
|
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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -313,6 +337,10 @@ impl HarmonyStreamingProcessor {
|
|||||||
original_request: Arc<ChatCompletionRequest>,
|
original_request: Arc<ChatCompletionRequest>,
|
||||||
tx: &mpsc::UnboundedSender<Result<Bytes, io::Error>>,
|
tx: &mpsc::UnboundedSender<Result<Bytes, io::Error>>,
|
||||||
) -> Result<(), String> {
|
) -> Result<(), String> {
|
||||||
|
// Timing for metrics
|
||||||
|
let start_time = Instant::now();
|
||||||
|
let mut first_token_time: Option<Instant> = None;
|
||||||
|
|
||||||
// Phase 1: Process prefill stream (collect metadata)
|
// Phase 1: Process prefill stream (collect metadata)
|
||||||
let mut prompt_tokens: HashMap<u32, u32> = HashMap::new();
|
let mut prompt_tokens: HashMap<u32, u32> = HashMap::new();
|
||||||
|
|
||||||
@@ -342,6 +370,11 @@ impl HarmonyStreamingProcessor {
|
|||||||
let chunk = chunk_wrapper.as_sglang();
|
let chunk = chunk_wrapper.as_sglang();
|
||||||
let index = chunk.index;
|
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
|
// Initialize parser for this index if needed
|
||||||
if let Vacant(e) = parsers.entry(index) {
|
if let Vacant(e) = parsers.entry(index) {
|
||||||
e.insert(
|
e.insert(
|
||||||
@@ -425,11 +458,12 @@ impl HarmonyStreamingProcessor {
|
|||||||
// This ensures that if client disconnects during decode, BOTH streams send abort
|
// This ensures that if client disconnects during decode, BOTH streams send abort
|
||||||
prefill_stream.mark_completed();
|
prefill_stream.mark_completed();
|
||||||
|
|
||||||
// Emit final usage if requested
|
// Compute totals once for both usage chunk and metrics
|
||||||
if let Some(true) = stream_options.as_ref().and_then(|so| so.include_usage) {
|
|
||||||
let total_prompt: u32 = prompt_tokens.values().sum();
|
let total_prompt: u32 = prompt_tokens.values().sum();
|
||||||
let total_completion: u32 = completion_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) {
|
||||||
Self::emit_usage_chunk(
|
Self::emit_usage_chunk(
|
||||||
total_prompt,
|
total_prompt,
|
||||||
total_completion,
|
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(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user