From ec0ffaacc8a6437681b5743cfdcc7841ada6b4d7 Mon Sep 17 00:00:00 2001 From: Bugen Zhao Date: Thu, 2 Jul 2026 22:55:25 +0800 Subject: [PATCH] [Rust Frontend] Improve scheduler stats logging parity (#47435) Signed-off-by: Bugen Zhao --- rust/src/engine-core-client/src/client.rs | 8 + rust/src/engine-core-client/src/metrics.rs | 13 + .../engine-core-client/src/protocol/stats.rs | 4 +- rust/src/llm/src/lib.rs | 2 +- rust/src/llm/src/log_stats.rs | 403 ++++++++++++++++-- rust/src/metrics/src/scheduler.rs | 121 +++++- 6 files changed, 521 insertions(+), 30 deletions(-) diff --git a/rust/src/engine-core-client/src/client.rs b/rust/src/engine-core-client/src/client.rs index 42680706a32..babc8179329 100644 --- a/rust/src/engine-core-client/src/client.rs +++ b/rust/src/engine-core-client/src/client.rs @@ -372,6 +372,14 @@ impl EngineCoreClient { self.engines.len() } + /// Return the engine-side indices connected to this client. + pub fn engine_indices(&self) -> Vec { + self.engines + .iter() + .map(|engine| engine.engine_id.engine_index().expect("engine id must encode as u16")) + .collect() + } + /// Return the engine identities of all engines connected to this client. pub fn engine_identities(&self) -> Vec<&[u8]> { self.engines.iter().map(|engine| &*engine.engine_id).collect() diff --git a/rust/src/engine-core-client/src/metrics.rs b/rust/src/engine-core-client/src/metrics.rs index a939b02f654..e129fdd6a27 100644 --- a/rust/src/engine-core-client/src/metrics.rs +++ b/rust/src/engine-core-client/src/metrics.rs @@ -84,6 +84,10 @@ pub(crate) fn record_scheduler_stats( .spec_decode_num_accepted_tokens .get_or_create(&labels) .inc_by(spec_decoding_stats.num_accepted_tokens); + metrics.log_stats.get_or_create(&labels).observe_spec_decode( + spec_decoding_stats.num_drafts, + &spec_decoding_stats.num_accepted_tokens_per_pos, + ); for (position, accepted_tokens) in spec_decoding_stats.num_accepted_tokens_per_pos.iter().copied().enumerate() @@ -119,6 +123,15 @@ pub(crate) fn record_scheduler_stats( .inc_by(perf_stats.num_write_bytes_per_gpu); } + if let Some(cudagraph_stats) = &stats.cudagraph_stats { + metrics.log_stats.get_or_create(&labels).observe_cudagraph( + cudagraph_stats.num_unpadded_tokens, + cudagraph_stats.num_padded_tokens, + cudagraph_stats.num_paddings, + &cudagraph_stats.runtime_mode, + ); + } + // Sampled KV-cache residency histograms. if !stats.kv_cache_eviction_events.is_empty() { let kv_block_lifetime_seconds = metrics.kv_block_lifetime_seconds.get_or_create(&labels); diff --git a/rust/src/engine-core-client/src/protocol/stats.rs b/rust/src/engine-core-client/src/protocol/stats.rs index 9f35f8301e4..1fc9606bcdf 100644 --- a/rust/src/engine-core-client/src/protocol/stats.rs +++ b/rust/src/engine-core-client/src/protocol/stats.rs @@ -141,7 +141,7 @@ pub struct PerfStats { /// Original Python definition: /// #[derive(Debug, Clone, Default, PartialEq, Serialize, Deserialize)] -pub struct CudagraphStat { +pub struct CudagraphStats { /// Number of real tokens in the captured batch before padding. pub num_unpadded_tokens: u64, /// Number of padded tokens in the captured batch. @@ -182,7 +182,7 @@ pub struct SchedulerStats { /// Connector-specific KV transfer stats, kept opaque for now. pub kv_connector_stats: Option>, /// CUDA graph runtime stats when graph metrics are enabled. - pub cudagraph_stats: Option, + pub cudagraph_stats: Option, /// Estimated MFU/performance stats, when enabled. pub perf_stats: Option, } diff --git a/rust/src/llm/src/lib.rs b/rust/src/llm/src/lib.rs index 9adfc737b63..ce15a970c63 100644 --- a/rust/src/llm/src/lib.rs +++ b/rust/src/llm/src/lib.rs @@ -51,7 +51,7 @@ impl Llm { if enabled { let stats_logger = StatsLogger::start( self.client.model_name().to_string(), - self.client.engine_count(), + self.client.engine_indices(), ); self.stats_logger = Some(stats_logger); } else { diff --git a/rust/src/llm/src/log_stats.rs b/rust/src/llm/src/log_stats.rs index 7d3a149731b..81c9093a108 100644 --- a/rust/src/llm/src/log_stats.rs +++ b/rust/src/llm/src/log_stats.rs @@ -4,10 +4,12 @@ use std::time::{Duration, Instant}; use tokio_util::task::AbortOnDropHandle; use tracing::{debug, info}; use vllm_metrics::{ - EngineLabels, F64Gauge, METRICS, PromptTokenSourceLabels, U64Counter, U64Gauge, + EngineLabels, F64Gauge, METRICS, PromptTokenSourceLabels, SchedulerLogStatsAccumulator, + SchedulerLogStatsInterval, U64Counter, U64Gauge, WaitingReasonLabels, }; const LOG_STATS_INTERVAL: Duration = Duration::from_secs(10); +const WAITING_REASON_DEFERRED: &str = "deferred"; /// Cached, cloned metric handles for one engine. Each clone shares the same /// underlying `Arc` as the prometheus `Family` entry, so reads go @@ -18,20 +20,58 @@ struct EngineMetrics { generation_tokens: U64Counter, prefix_cache_queries: U64Counter, prefix_cache_hits: U64Counter, + external_prefix_cache_queries: U64Counter, + external_prefix_cache_hits: U64Counter, + num_preemptions: U64Counter, + spec_decode_num_drafts: U64Counter, + spec_decode_num_draft_tokens: U64Counter, + spec_decode_num_accepted_tokens: U64Counter, + estimated_flops_per_gpu: U64Counter, + estimated_read_bytes_per_gpu: U64Counter, + estimated_write_bytes_per_gpu: U64Counter, + log_stats: SchedulerLogStatsAccumulator, // Gauges for instantaneous scheduler state. scheduler_running: U64Gauge, scheduler_waiting: U64Gauge, + scheduler_deferred: U64Gauge, kv_cache_usage: F64Gauge, } /// Accumulated snapshot values from the last logging interval, used to compute /// deltas. +#[derive(Default)] struct CounterSnapshot { prompt_tokens: u64, generation_tokens: u64, prefix_cache_queries: u64, prefix_cache_hits: u64, + external_prefix_cache_queries: u64, + external_prefix_cache_hits: u64, + num_preemptions: u64, + spec_decode_num_drafts: u64, + spec_decode_num_draft_tokens: u64, + spec_decode_num_accepted_tokens: u64, + estimated_flops_per_gpu: u64, + estimated_read_bytes_per_gpu: u64, + estimated_write_bytes_per_gpu: u64, +} + +/// Derived spec-decoding values for one logging interval. +struct SpecDecodingLogStats { + mean_acceptance_length: f64, + accepted_throughput: f64, + draft_throughput: f64, + accepted_tokens: u64, + draft_tokens: u64, + per_position_acceptance_rates: Vec, + draft_acceptance_rate: f64, +} + +/// Derived MFU values for one logging interval. +struct MfuLogStats { + tflops_per_gpu: f64, + gbps_per_gpu: f64, } /// Periodic stats logger that mirrors Python vLLM's `LoggingStatLogger`. @@ -46,18 +86,20 @@ pub(crate) struct StatsLogger { impl StatsLogger { /// Start the background stats logging task. - pub(crate) fn start(model_name: String, engine_count: usize) -> Self { + pub(crate) fn start(model_name: String, engine_indices: Vec) -> Self { let task = AbortOnDropHandle::new(tokio::spawn(async move { - run_stats_logger(model_name, engine_count).await; + run_stats_logger(model_name, engine_indices).await; })); Self { _task: task } } } /// Resolve and clone all metric handles once so the hot path is lock-free. -fn resolve_engine_metrics(model_name: &str, engine_count: usize) -> Vec { +fn resolve_engine_metrics(model_name: &str, engine_indices: &[u32]) -> Vec { let m = &METRICS; - (0..engine_count as u32) + engine_indices + .iter() + .copied() .map(|engine| { let el = EngineLabels { model_name: model_name.to_string(), @@ -68,6 +110,11 @@ fn resolve_engine_metrics(model_name: &str, engine_count: usize) -> Vec Vec) { + let engines = resolve_engine_metrics(&model_name, &engine_indices); let mut interval = tokio::time::interval(LOG_STATS_INTERVAL); interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay); @@ -109,6 +191,7 @@ async fn run_stats_logger(model_name: String, engine_count: usize) { } let curr = read_counters(&engines); + let raw_log_stats = drain_scheduler_log_stats(&engines); let prompt_throughput = curr.prompt_tokens.wrapping_sub(prev.prompt_tokens) as f64 / elapsed; @@ -121,17 +204,37 @@ async fn run_stats_logger(model_name: String, engine_count: usize) { && last_prompt_throughput == 0.0 && last_generation_throughput == 0.0; + /// Emit one stats line at DEBUG while idle and INFO while active. + macro_rules! log_stats_line { + ($($arg:tt)*) => { + if is_idle { + debug!($($arg)*); + } else { + info!($($arg)*); + } + }; + } + // Read scheduler gauges (aggregate across engines). let (num_running, num_waiting, kv_cache_usage) = read_scheduler_gauges(&engines); + let num_deferred = read_deferred_waiting(&engines); + let delta_preemptions = curr.num_preemptions.wrapping_sub(prev.num_preemptions); // Compute prefix cache hit rate over this interval. let delta_queries = curr.prefix_cache_queries.wrapping_sub(prev.prefix_cache_queries); - let prefix_cache_hit_rate = if delta_queries > 0 { - let delta_hits = curr.prefix_cache_hits.wrapping_sub(prev.prefix_cache_hits); - delta_hits as f64 / delta_queries as f64 * 100.0 - } else { - 0.0 - }; + let delta_hits = curr.prefix_cache_hits.wrapping_sub(prev.prefix_cache_hits); + let prefix_cache_hit_rate = cache_hit_rate(delta_hits, delta_queries); + + let delta_external_queries = curr + .external_prefix_cache_queries + .wrapping_sub(prev.external_prefix_cache_queries); + let delta_external_hits = + curr.external_prefix_cache_hits.wrapping_sub(prev.external_prefix_cache_hits); + let external_prefix_cache_hit_rate = + cache_hit_rate(delta_external_hits, delta_external_queries); + let spec_decoding_log_stats = + spec_decoding_log_stats(&curr, &prev, elapsed, &raw_log_stats); + let mfu_log_stats = mfu_log_stats(&curr, &prev, elapsed, engines.len()); // Build the log line. msg.clear(); @@ -140,17 +243,70 @@ async fn run_stats_logger(model_name: String, engine_count: usize) { "Avg prompt tput: {prompt_throughput:.1} toks/s, \ Avg generation tput: {generation_throughput:.1} toks/s, \ Reqs Running: {num_running}, \ - Waiting: {num_waiting}, \ - GPU KV cache used: {:.1}%, \ + Waiting: {num_waiting}" + ) + .unwrap(); + if num_deferred > 0 { + write!(msg, ", Deferred: {num_deferred} reqs").unwrap(); + } + if delta_preemptions > 0 { + write!(msg, ", Preemptions: {delta_preemptions}").unwrap(); + } + write!( + msg, + ", GPU KV cache used: {:.1}%, \ Prefix cache hit rate: {prefix_cache_hit_rate:.1}%", kv_cache_usage * 100.0, ) .unwrap(); + if delta_external_queries > 0 { + write!( + msg, + ", External prefix cache hit rate: {external_prefix_cache_hit_rate:.1}%" + ) + .unwrap(); + } - if is_idle { - debug!("{msg}"); - } else { - info!("{msg}"); + log_stats_line!("{msg}"); + + if let Some(spec_stats) = spec_decoding_log_stats { + msg.clear(); + write!( + msg, + "SpecDecoding metrics: \ + Mean acceptance length: {:.2}, \ + Accepted throughput: {:.2} tokens/s, \ + Drafted throughput: {:.2} tokens/s, \ + Accepted: {} tokens, \ + Drafted: {} tokens", + spec_stats.mean_acceptance_length, + spec_stats.accepted_throughput, + spec_stats.draft_throughput, + spec_stats.accepted_tokens, + spec_stats.draft_tokens, + ) + .unwrap(); + if !spec_stats.per_position_acceptance_rates.is_empty() { + msg.push_str(", Per-position acceptance rate: "); + format_position_rates(&mut msg, &spec_stats.per_position_acceptance_rates); + } + write!( + msg, + ", Avg Draft acceptance rate: {:.1}%", + spec_stats.draft_acceptance_rate, + ) + .unwrap(); + log_stats_line!("{msg}"); + } + + // TODO: Decide on best way to surface CUDAGraph interval samples. + + if let Some(mfu_stats) = mfu_log_stats { + log_stats_line!( + "MFU: {:.1} TF/s/GPU {:.1} GB/s/GPU", + mfu_stats.tflops_per_gpu, + mfu_stats.gbps_per_gpu, + ); } last_prompt_throughput = prompt_throughput; @@ -162,17 +318,21 @@ async fn run_stats_logger(model_name: String, engine_count: usize) { /// Read the current cumulative counter values for throughput computation. fn read_counters(engines: &[EngineMetrics]) -> CounterSnapshot { - let mut snap = CounterSnapshot { - prompt_tokens: 0, - generation_tokens: 0, - prefix_cache_queries: 0, - prefix_cache_hits: 0, - }; + let mut snap = CounterSnapshot::default(); for e in engines { snap.prompt_tokens += e.prompt_tokens_computed.get(); snap.generation_tokens += e.generation_tokens.get(); snap.prefix_cache_queries += e.prefix_cache_queries.get(); snap.prefix_cache_hits += e.prefix_cache_hits.get(); + snap.external_prefix_cache_queries += e.external_prefix_cache_queries.get(); + snap.external_prefix_cache_hits += e.external_prefix_cache_hits.get(); + snap.num_preemptions += e.num_preemptions.get(); + snap.spec_decode_num_drafts += e.spec_decode_num_drafts.get(); + snap.spec_decode_num_draft_tokens += e.spec_decode_num_draft_tokens.get(); + snap.spec_decode_num_accepted_tokens += e.spec_decode_num_accepted_tokens.get(); + snap.estimated_flops_per_gpu += e.estimated_flops_per_gpu.get(); + snap.estimated_read_bytes_per_gpu += e.estimated_read_bytes_per_gpu.get(); + snap.estimated_write_bytes_per_gpu += e.estimated_write_bytes_per_gpu.get(); } snap } @@ -197,3 +357,196 @@ fn read_scheduler_gauges(engines: &[EngineMetrics]) -> (u64, u64, f64) { (num_running, num_waiting, kv_cache_usage) } + +/// Read deferred waiting requests across all engines. +fn read_deferred_waiting(engines: &[EngineMetrics]) -> u64 { + engines.iter().map(|e| e.scheduler_deferred.get()).sum() +} + +/// Return the cache hit rate as a percentage for a counter delta. +fn cache_hit_rate(hits: u64, queries: u64) -> f64 { + if queries > 0 { + hits as f64 / queries as f64 * 100.0 + } else { + 0.0 + } +} + +/// Compute aggregate spec-decoding stats for one logging interval. +fn spec_decoding_log_stats( + curr: &CounterSnapshot, + prev: &CounterSnapshot, + elapsed: f64, + raw_log_stats: &SchedulerLogStatsInterval, +) -> Option { + let num_drafts = curr.spec_decode_num_drafts.wrapping_sub(prev.spec_decode_num_drafts); + if num_drafts == 0 { + return None; + } + + let draft_tokens = curr + .spec_decode_num_draft_tokens + .wrapping_sub(prev.spec_decode_num_draft_tokens); + let accepted_tokens = curr + .spec_decode_num_accepted_tokens + .wrapping_sub(prev.spec_decode_num_accepted_tokens); + + let (accepted_throughput, draft_throughput) = if elapsed > 0.0 { + ( + accepted_tokens as f64 / elapsed, + draft_tokens as f64 / elapsed, + ) + } else { + (0.0, 0.0) + }; + let draft_acceptance_rate = if draft_tokens > 0 { + accepted_tokens as f64 / draft_tokens as f64 * 100.0 + } else { + f64::NAN + }; + let per_position_acceptance_rates = if raw_log_stats.spec_num_drafts > 0 { + raw_log_stats + .spec_accepted_tokens_per_pos + .iter() + .map(|accepted_tokens| *accepted_tokens as f64 / raw_log_stats.spec_num_drafts as f64) + .collect() + } else { + Vec::new() + }; + + Some(SpecDecodingLogStats { + mean_acceptance_length: 1.0 + accepted_tokens as f64 / num_drafts as f64, + accepted_throughput, + draft_throughput, + accepted_tokens, + draft_tokens, + per_position_acceptance_rates, + draft_acceptance_rate, + }) +} + +/// Compute average per-GPU MFU rates for one logging interval. +fn mfu_log_stats( + curr: &CounterSnapshot, + prev: &CounterSnapshot, + elapsed: f64, + engine_count: usize, +) -> Option { + let flops = curr.estimated_flops_per_gpu.wrapping_sub(prev.estimated_flops_per_gpu); + let read_bytes = curr + .estimated_read_bytes_per_gpu + .wrapping_sub(prev.estimated_read_bytes_per_gpu); + let write_bytes = curr + .estimated_write_bytes_per_gpu + .wrapping_sub(prev.estimated_write_bytes_per_gpu); + + if flops == 0 && read_bytes == 0 && write_bytes == 0 { + return None; + } + + let denominator = elapsed * engine_count.max(1) as f64; + let (tflops_per_gpu, gbps_per_gpu) = if denominator > 0.0 { + ( + flops as f64 / denominator / 1e12, + (read_bytes as f64 + write_bytes as f64) / denominator / 1e9, + ) + } else { + (0.0, 0.0) + }; + + Some(MfuLogStats { + tflops_per_gpu, + gbps_per_gpu, + }) +} + +/// Drain raw scheduler DTO stats for the configured model and engines. +fn drain_scheduler_log_stats(engines: &[EngineMetrics]) -> SchedulerLogStatsInterval { + let mut interval = SchedulerLogStatsInterval::default(); + for engine in engines { + interval.merge(engine.log_stats.drain()); + } + interval +} + +/// Append spec-decoding per-position acceptance rates like Python's logger. +fn format_position_rates(output: &mut String, rates: &[f64]) { + for (position, rate) in rates.iter().enumerate() { + if position > 0 { + output.push_str(", "); + } + write!(output, "{rate:.3}").unwrap(); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn cache_hit_rate_returns_percent_for_non_empty_queries() { + assert_eq!(cache_hit_rate(25, 100), 25.0); + assert_eq!(cache_hit_rate(0, 0), 0.0); + } + + #[test] + fn spec_decoding_log_stats_uses_interval_deltas() { + let raw_log_stats = SchedulerLogStatsInterval { + spec_num_drafts: 4, + spec_accepted_tokens_per_pos: vec![4, 2, 1], + ..Default::default() + }; + let prev = CounterSnapshot { + spec_decode_num_drafts: 10, + spec_decode_num_draft_tokens: 100, + spec_decode_num_accepted_tokens: 40, + ..Default::default() + }; + let curr = CounterSnapshot { + spec_decode_num_drafts: 14, + spec_decode_num_draft_tokens: 120, + spec_decode_num_accepted_tokens: 52, + ..Default::default() + }; + + let stats = spec_decoding_log_stats(&curr, &prev, 2.0, &raw_log_stats).unwrap(); + + assert_eq!(stats.mean_acceptance_length, 4.0); + assert_eq!(stats.accepted_throughput, 6.0); + assert_eq!(stats.draft_throughput, 10.0); + assert_eq!(stats.accepted_tokens, 12); + assert_eq!(stats.draft_tokens, 20); + assert_eq!(stats.per_position_acceptance_rates, vec![1.0, 0.5, 0.25]); + assert_eq!(stats.draft_acceptance_rate, 60.0); + } + + #[test] + fn mfu_log_stats_averages_per_gpu_across_engines() { + let prev = CounterSnapshot { + estimated_flops_per_gpu: 10, + estimated_read_bytes_per_gpu: 10, + estimated_write_bytes_per_gpu: 10, + ..Default::default() + }; + let curr = CounterSnapshot { + estimated_flops_per_gpu: 4_000_000_000_010, + estimated_read_bytes_per_gpu: 2_000_000_010, + estimated_write_bytes_per_gpu: 2_000_000_010, + ..Default::default() + }; + + let stats = mfu_log_stats(&curr, &prev, 2.0, 2).unwrap(); + + assert_eq!(stats.tflops_per_gpu, 1.0); + assert_eq!(stats.gbps_per_gpu, 1.0); + } + + #[test] + fn format_position_rates_uses_three_decimal_places() { + let mut output = String::new(); + + format_position_rates(&mut output, &[1.0, 0.5, 0.25]); + + assert_eq!(output, "1.000, 0.500, 0.250"); + } +} diff --git a/rust/src/metrics/src/scheduler.rs b/rust/src/metrics/src/scheduler.rs index ec5f8d4e9f3..c38cc291dc7 100644 --- a/rust/src/metrics/src/scheduler.rs +++ b/rust/src/metrics/src/scheduler.rs @@ -1,4 +1,5 @@ -use std::collections::BTreeSet; +use std::collections::{BTreeMap, BTreeSet}; +use std::sync::{Arc, Mutex}; use itertools::Itertools as _; use prometheus_client::encoding::{EncodeLabelSet, EncodeLabelValue, LabelValueEncoder}; @@ -62,6 +63,90 @@ pub struct LoraInfoLabels { pub waiting_lora_adapters: LoraAdapterNames, } +/// CUDA graph sample key used for periodic text-log aggregation. +#[derive(Clone, Debug, PartialEq, Eq, PartialOrd, Ord)] +pub struct CudagraphLogKey { + pub num_unpadded_tokens: u64, + pub num_padded_tokens: u64, + pub num_paddings: u64, + pub runtime_mode: String, +} + +/// Raw scheduler stats accumulated for one periodic text-log interval. +#[derive(Default)] +pub struct SchedulerLogStatsInterval { + pub spec_num_drafts: u64, + pub spec_accepted_tokens_per_pos: Vec, + pub cudagraph_counts: BTreeMap, +} + +impl SchedulerLogStatsInterval { + /// Merge another drained interval into this one. + pub fn merge(&mut self, other: Self) { + self.spec_num_drafts += other.spec_num_drafts; + + if self.spec_accepted_tokens_per_pos.len() < other.spec_accepted_tokens_per_pos.len() { + self.spec_accepted_tokens_per_pos + .resize(other.spec_accepted_tokens_per_pos.len(), 0); + } + for (position, accepted_tokens) in + other.spec_accepted_tokens_per_pos.into_iter().enumerate() + { + self.spec_accepted_tokens_per_pos[position] += accepted_tokens; + } + + for (key, count) in other.cudagraph_counts { + *self.cudagraph_counts.entry(key).or_default() += count; + } + } +} + +/// Internal, non-Prometheus accumulator for periodic text logs that need raw +/// scheduler DTOs. +#[derive(Clone, Default)] +pub struct SchedulerLogStatsAccumulator { + inner: Arc>, +} + +impl SchedulerLogStatsAccumulator { + /// Observe spec-decoding fields needed for per-position text-log rates. + pub fn observe_spec_decode(&self, num_drafts: u64, accepted_tokens_per_pos: &[u64]) { + let mut inner = self.inner.lock().expect("scheduler log stats accumulator poisoned"); + inner.spec_num_drafts += num_drafts; + + if inner.spec_accepted_tokens_per_pos.len() < accepted_tokens_per_pos.len() { + inner.spec_accepted_tokens_per_pos.resize(accepted_tokens_per_pos.len(), 0); + } + for (position, accepted_tokens) in accepted_tokens_per_pos.iter().copied().enumerate() { + inner.spec_accepted_tokens_per_pos[position] += accepted_tokens; + } + } + + /// Observe one CUDA graph runtime sample for the interval table. + pub fn observe_cudagraph( + &self, + num_unpadded_tokens: u64, + num_padded_tokens: u64, + num_paddings: u64, + runtime_mode: &str, + ) { + let mut inner = self.inner.lock().expect("scheduler log stats accumulator poisoned"); + let key = CudagraphLogKey { + num_unpadded_tokens, + num_padded_tokens, + num_paddings, + runtime_mode: runtime_mode.to_string(), + }; + *inner.cudagraph_counts.entry(key).or_default() += 1; + } + + /// Drain and reset the current text-log interval. + pub fn drain(&self) -> SchedulerLogStatsInterval { + let mut inner = self.inner.lock().expect("scheduler log stats accumulator poisoned"); + std::mem::take(&mut *inner) + } +} + /// Scheduler/batch-scoped Prometheus families exported from `SchedulerStats`. pub struct SchedulerMetrics { // Scheduler state gauges. @@ -95,6 +180,9 @@ pub struct SchedulerMetrics { pub kv_block_lifetime_seconds: HistogramFamily, pub kv_block_idle_before_evict_seconds: HistogramFamily, pub kv_block_reuse_gap_seconds: HistogramFamily, + + /// Non-Prometheus interval accumulators for periodic text-log helpers. + pub log_stats: Family, } impl SchedulerMetrics { @@ -265,13 +353,14 @@ impl SchedulerMetrics { kv_block_lifetime_seconds, kv_block_idle_before_evict_seconds, kv_block_reuse_gap_seconds, + log_stats: Family::default(), } } } #[cfg(test)] mod tests { - use crate::{EngineLabels, Metrics}; + use crate::{CudagraphLogKey, EngineLabels, Metrics, SchedulerLogStatsAccumulator}; #[test] fn perf_counters_render_with_a_single_total_suffix() { @@ -301,4 +390,32 @@ mod tests { assert!(!rendered.contains("vllm:estimated_read_bytes_per_gpu_total_total")); assert!(!rendered.contains("vllm:estimated_write_bytes_per_gpu_total_total")); } + + #[test] + fn log_stats_accumulator_drains_interval_data() { + let accumulator = SchedulerLogStatsAccumulator::default(); + + accumulator.observe_spec_decode(2, &[1, 2]); + accumulator.observe_spec_decode(3, &[3, 4, 5]); + accumulator.observe_cudagraph(8, 16, 8, "FULL"); + accumulator.observe_cudagraph(8, 16, 8, "FULL"); + + let interval = accumulator.drain(); + + assert_eq!(interval.spec_num_drafts, 5); + assert_eq!(interval.spec_accepted_tokens_per_pos, vec![4, 6, 5]); + assert_eq!( + interval + .cudagraph_counts + .get(&CudagraphLogKey { + num_unpadded_tokens: 8, + num_padded_tokens: 16, + num_paddings: 8, + runtime_mode: "FULL".to_string(), + }) + .copied(), + Some(2) + ); + assert_eq!(accumulator.drain().spec_num_drafts, 0); + } }