From 47b446ef70b1f074b885d901c12996eb78690f84 Mon Sep 17 00:00:00 2001 From: David <12414531+DavidBellamy@users.noreply.github.com> Date: Fri, 7 Aug 2026 02:26:17 -0700 Subject: [PATCH 1/3] fix(cache): include scheduler occupancy in owner pressure Signed-off-by: David <12414531+DavidBellamy@users.noreply.github.com> --- model_gateway/src/policies/cache_aware.rs | 334 +++++++++++++++++++++- 1 file changed, 326 insertions(+), 8 deletions(-) diff --git a/model_gateway/src/policies/cache_aware.rs b/model_gateway/src/policies/cache_aware.rs index 5e6b598fa..a60a77146 100644 --- a/model_gateway/src/policies/cache_aware.rs +++ b/model_gateway/src/policies/cache_aware.rs @@ -38,9 +38,10 @@ ------------------------------------------- Restricts cached-owner candidates to workers within 10 percentage points of the least-pressured owner and below 90% pressure. Non-owner candidates use - the equivalent fleet-scoped guard for cold fallback. Pressure is the larger - of KV token usage and utilization; waiting requests break ties. Missing or - stale telemetry fails open to existing owners. + the equivalent fleet-scoped guard for cold fallback. Pressure is the largest + of KV token usage, utilization, and scheduler occupancy when every engine + reports a running-request cap for every DP rank. Missing, partial, or stale + telemetry fails open to existing owners. Configuration Parameters: ------------------------ @@ -346,7 +347,14 @@ impl CacheAwarePolicy { } let loads = self.engine_loads.read(); - let mut pressure_by_index = HashMap::with_capacity(healthy_indices.len()); + struct RawWorkerPressure { + idx: usize, + backend_pressure: f64, + scheduler_pressure: Option, + waiting_requests: i64, + } + + let mut raw_pressure = Vec::with_capacity(healthy_indices.len()); // Degrade the whole decision to legacy request-count routing unless all // candidates have comparable, fresh engine telemetry. @@ -357,27 +365,66 @@ impl CacheAwarePolicy { return None; } - let pressure = load + let backend_pressure = load .response .loads .iter() .map(|rank| rank.token_usage.max(rank.utilization)) .fold(0.0_f64, f64::max); - if !pressure.is_finite() || pressure < 0.0 { + if !backend_pressure.is_finite() || backend_pressure < 0.0 { return None; } + let running_requests: i64 = load + .response + .loads + .iter() + .map(|rank| i64::from(rank.num_running_reqs.max(0))) + .sum(); let waiting_requests = load .response .loads .iter() .map(|rank| i64::from(rank.num_waiting_reqs.max(0))) .sum(); + let reported_max_running = load.response.loads.iter().try_fold(0_i64, |total, rank| { + let rank_cap = i64::from(rank.max_running_requests); + (rank_cap > 0).then(|| total.saturating_add(rank_cap)) + }); + let scheduler_pressure = reported_max_running.map(|max_running| { + (running_requests.saturating_add(waiting_requests) as f64 / max_running as f64) + .clamp(0.0, 1.0) + }); - pressure_by_index.insert( + raw_pressure.push(RawWorkerPressure { idx, + backend_pressure, + scheduler_pressure, + waiting_requests, + }); + } + + // Scheduler counts and caps must be comparable across the entire + // candidate set. In particular, SGLang's Prometheus-only fallback can + // repeat request gauges for every TP rank while omitting the cap. A + // missing cap therefore disables scheduler occupancy for this decision + // instead of making that worker look artificially saturated. + let use_scheduler_pressure = raw_pressure + .iter() + .all(|load| load.scheduler_pressure.is_some()); + let mut pressure_by_index = HashMap::with_capacity(raw_pressure.len()); + for load in raw_pressure { + let pressure = if use_scheduler_pressure { + load.backend_pressure + .max(load.scheduler_pressure.unwrap_or_default()) + } else { + load.backend_pressure + }; + + pressure_by_index.insert( + load.idx, WorkerPressure { pressure, - waiting_requests, + waiting_requests: load.waiting_requests, }, ); } @@ -1995,6 +2042,49 @@ mod tests { } } + fn engine_load_with_scheduler( + token_usage: f64, + utilization: f64, + running: i32, + waiting: i32, + max_running: i32, + ) -> WorkerLoadResponse { + WorkerLoadResponse { + loads: vec![SchedulerLoadSnapshot { + token_usage, + utilization, + num_running_reqs: running, + num_waiting_reqs: waiting, + max_running_requests: max_running, + ..Default::default() + }], + dp_rank_count: 1, + ..Default::default() + } + } + + fn engine_load_with_scheduler_ranks(ranks: &[(i32, i32, i32)]) -> WorkerLoadResponse { + WorkerLoadResponse { + loads: ranks + .iter() + .enumerate() + .map( + |(dp_rank, &(running, waiting, max_running))| SchedulerLoadSnapshot { + dp_rank: i32::try_from(dp_rank).unwrap(), + token_usage: 0.10, + utilization: 0.10, + num_running_reqs: running, + num_waiting_reqs: waiting, + max_running_requests: max_running, + ..Default::default() + }, + ) + .collect(), + dp_rank_count: i32::try_from(ranks.len()).unwrap(), + ..Default::default() + } + } + fn two_workers() -> Vec> { vec![ Arc::new( @@ -2012,6 +2102,24 @@ mod tests { ] } + fn two_workers_with_max_running(max_running: u16) -> Vec> { + ["http://w1:8000", "http://w2:8000"] + .into_iter() + .map(|url| { + Arc::new( + BasicWorkerBuilder::new(url) + .worker_type(WorkerType::Regular) + .labels(HashMap::from([( + "max_running_requests".to_string(), + max_running.to_string(), + )])) + .health_config(no_health_check()) + .build(), + ) as Arc + }) + .collect() + } + fn engine_aware_policy() -> CacheAwarePolicy { CacheAwarePolicy::with_config(CacheAwareConfig { engine_load: true, @@ -2034,6 +2142,28 @@ mod tests { .insert_text(text, workers[0].url()); } + fn scheduler_pressure_policy() -> CacheAwarePolicy { + CacheAwarePolicy::with_config(CacheAwareConfig { + engine_load: true, + eviction_interval_secs: 0, + max_cached_owners_per_prefix: 8, + cache_owner_spill_cooldown_secs: 60, + ..Default::default() + }) + } + + fn select_text(policy: &CacheAwarePolicy, workers: &[Arc], text: &str) -> usize { + policy + .select_worker( + workers, + &SelectWorkerInfo { + request_text: Some(text), + ..Default::default() + }, + ) + .unwrap() + } + #[test] fn test_engine_load_keeps_cache_affinity_when_pressure_is_close() { let policy = engine_aware_policy(); @@ -2080,6 +2210,194 @@ mod tests { assert_eq!(selected, 1); } + #[test] + fn scheduler_pressure_keeps_affinity_below_watermark() { + let policy = scheduler_pressure_policy(); + let workers = two_workers(); + prime_worker_one_affinity(&policy, &workers, "shared prefix"); + policy.update_loads(&HashMap::from([ + ( + "http://w1:8000".to_string(), + engine_load_with_scheduler(0.10, 0.10, 32, 1, 37), + ), + ( + "http://w2:8000".to_string(), + engine_load_with_scheduler(0.10, 0.10, 0, 0, 37), + ), + ])); + + let selected = select_text(&policy, &workers, "shared prefix"); + + assert_eq!(selected, 0); + assert!(policy.replication_state.is_empty()); + } + + #[test] + fn scheduler_pressure_spills_above_watermark() { + let policy = scheduler_pressure_policy(); + let workers = two_workers(); + prime_worker_one_affinity(&policy, &workers, "shared prefix"); + policy.update_loads(&HashMap::from([ + ( + "http://w1:8000".to_string(), + engine_load_with_scheduler(0.10, 0.10, 33, 1, 37), + ), + ( + "http://w2:8000".to_string(), + engine_load_with_scheduler(0.10, 0.10, 0, 0, 37), + ), + ])); + + let selected = select_text(&policy, &workers, "shared prefix"); + + assert_eq!(selected, 1); + assert_eq!(policy.replication_state.len(), 1); + } + + #[test] + fn scheduler_pressure_sums_complete_rank_caps() { + let policy = scheduler_pressure_policy(); + let workers = two_workers(); + prime_worker_one_affinity(&policy, &workers, "shared prefix"); + policy.update_loads(&HashMap::from([ + ( + "http://w1:8000".to_string(), + engine_load_with_scheduler_ranks(&[(17, 0, 37), (17, 0, 37)]), + ), + ( + "http://w2:8000".to_string(), + engine_load_with_scheduler_ranks(&[(0, 0, 37), (0, 0, 37)]), + ), + ])); + + let selected = select_text(&policy, &workers, "shared prefix"); + + assert_eq!(selected, 0); + assert!(policy.replication_state.is_empty()); + } + + #[test] + fn scheduler_pressure_ignores_metadata_cap_when_rank_caps_are_missing() { + let policy = scheduler_pressure_policy(); + let workers = two_workers_with_max_running(37); + prime_worker_one_affinity(&policy, &workers, "shared prefix"); + policy.update_loads(&HashMap::from([ + ( + "http://w1:8000".to_string(), + engine_load_with_scheduler_ranks(&[(17, 0, 0), (17, 0, 0)]), + ), + ( + "http://w2:8000".to_string(), + engine_load_with_scheduler_ranks(&[(0, 0, 0), (0, 0, 0)]), + ), + ])); + + let selected = select_text(&policy, &workers, "shared prefix"); + + assert_eq!(selected, 0); + assert!(policy.replication_state.is_empty()); + } + + #[test] + fn scheduler_pressure_ignores_partial_rank_caps() { + let policy = scheduler_pressure_policy(); + let workers = two_workers(); + prime_worker_one_affinity(&policy, &workers, "shared prefix"); + policy.update_loads(&HashMap::from([ + ( + "http://w1:8000".to_string(), + engine_load_with_scheduler_ranks(&[(17, 0, 37), (17, 0, 0)]), + ), + ( + "http://w2:8000".to_string(), + engine_load_with_scheduler_ranks(&[(0, 0, 37), (0, 0, 37)]), + ), + ])); + + let selected = select_text(&policy, &workers, "shared prefix"); + + assert_eq!(selected, 0); + assert!(policy.replication_state.is_empty()); + } + + #[test] + fn scheduler_pressure_ignores_tp_duplicated_counts_without_snapshot_cap() { + let policy = scheduler_pressure_policy(); + let workers = two_workers_with_max_running(37); + prime_worker_one_affinity(&policy, &workers, "shared prefix"); + policy.update_loads(&HashMap::from([ + ( + "http://w1:8000".to_string(), + engine_load_with_scheduler(0.10, 0.10, 176, 32, 0), + ), + ( + "http://w2:8000".to_string(), + engine_load_with_scheduler(0.10, 0.10, 0, 0, 0), + ), + ])); + + let selected = select_text(&policy, &workers, "shared prefix"); + + assert_eq!(selected, 0); + assert!(policy.replication_state.is_empty()); + } + + #[test] + fn scheduler_pressure_uses_an_existing_suitable_owner_before_spilling() { + let policy = scheduler_pressure_policy(); + let workers = make_workers(&["http://w1:8000", "http://w2:8000", "http://w3:8000"]); + policy.init_workers(&workers); + let model_id = normalize_model_key(workers[0].model_id()); + let tree = policy.string_trees.get(model_id).unwrap().value().clone(); + tree.insert_text("hot prefix", workers[0].url()); + tree.insert_text("hot prefix", workers[1].url()); + policy.update_loads(&HashMap::from([ + ( + "http://w1:8000".to_string(), + engine_load_with_scheduler(0.20, 0.20, 37, 116, 37), + ), + ( + "http://w2:8000".to_string(), + engine_load_with_scheduler(0.20, 0.20, 20, 0, 37), + ), + ( + "http://w3:8000".to_string(), + engine_load_with_scheduler(0.20, 0.20, 0, 0, 37), + ), + ])); + + let selected = select_text(&policy, &workers, "hot prefix"); + + assert_eq!(selected, 1); + assert!(policy.replication_state.is_empty()); + assert_eq!(tree.match_prefix_with_counts("hot prefix").tenants.len(), 2); + } + + #[test] + fn scheduler_pressure_skips_a_saturated_spill_candidate() { + let policy = scheduler_pressure_policy(); + let workers = make_workers(&["http://w1:8000", "http://w2:8000", "http://w3:8000"]); + prime_worker_one_affinity(&policy, &workers, "hot prefix"); + policy.update_loads(&HashMap::from([ + ( + "http://w1:8000".to_string(), + engine_load_with_scheduler(0.20, 0.20, 37, 116, 37), + ), + ( + "http://w2:8000".to_string(), + engine_load_with_scheduler(0.20, 0.20, 37, 0, 37), + ), + ( + "http://w3:8000".to_string(), + engine_load_with_scheduler(0.20, 0.20, 0, 0, 37), + ), + ])); + + let selected = select_text(&policy, &workers, "hot prefix"); + + assert_eq!(selected, 2); + } + #[test] fn test_engine_load_falls_back_when_telemetry_is_partial() { let policy = engine_aware_policy(); From b828e2a3fb2f17e1c0d9c6792cab5ba536041b67 Mon Sep 17 00:00:00 2001 From: David <12414531+DavidBellamy@users.noreply.github.com> Date: Fri, 7 Aug 2026 16:26:01 -0700 Subject: [PATCH 2/3] fix(cache): surface legacy scheduler capacity Signed-off-by: David <12414531+DavidBellamy@users.noreply.github.com> --- model_gateway/src/worker/monitor.rs | 79 +++++++++++++++++++++++++++-- 1 file changed, 74 insertions(+), 5 deletions(-) diff --git a/model_gateway/src/worker/monitor.rs b/model_gateway/src/worker/monitor.rs index 14a9aa86a..b57310f0e 100644 --- a/model_gateway/src/worker/monitor.rs +++ b/model_gateway/src/worker/monitor.rs @@ -593,7 +593,12 @@ impl WorkerMonitor { .find(|p| m.has(&format!("{p}token_usage")))?; if let Some(loads) = standard_loads { - return Self::combine_sglang_standard_loads(loads, &m, prefix); + return Self::combine_sglang_standard_loads( + loads, + &m, + prefix, + worker.max_running_requests().map_or(0, i32::from), + ); } Some(Self::single_rank(SchedulerLoadSnapshot { @@ -623,11 +628,18 @@ impl WorkerMonitor { /// Prometheus values are averaged across emitted samples. Assigning the /// average to every DP snapshot preserves the aggregate throughput when /// metrics are per-DP, while avoiding multiplication when a TP-only - /// engine repeats the same sample on every tensor rank. + /// engine repeats the same sample on every tensor rank. Older SGLang + /// versions omit `max_running_requests` from this canonical endpoint, so + /// divide the instance-wide capacity discovered from `/server_info` across + /// a complete set of missing per-DP caps. Mixed reported/missing caps stay + /// incomplete so scheduler-pressure routing fails open. The Prometheus-only + /// fallback remains unchanged and does not receive this capacity, because + /// its request gauges may be repeated by TP rank. fn combine_sglang_standard_loads( loads: Vec, metrics: &PromScrape, prefix: &str, + fallback_max_running_requests: i32, ) -> Option { if loads.is_empty() { return None; @@ -642,6 +654,15 @@ impl WorkerMonitor { let cache_hit_rate = metrics.mean(&format!("{prefix}cache_hit_rate")); let utilization = metrics.mean(&format!("{prefix}utilization")); let dp_rank_count = i32::try_from(loads.len()).unwrap_or(i32::MAX); + let all_caps_missing = loads.iter().all(|load| load.max_running_requests <= 0); + let fallback_per_dp = if all_caps_missing { + fallback_max_running_requests + .max(0) + .checked_div(dp_rank_count) + .unwrap_or(0) + } else { + 0 + }; let loads = loads .into_iter() .map(|load| { @@ -658,7 +679,11 @@ impl WorkerMonitor { gen_throughput, cache_hit_rate, utilization, - max_running_requests: load.max_running_requests.max(0), + max_running_requests: if load.max_running_requests > 0 { + load.max_running_requests + } else { + fallback_per_dp + }, ..Default::default() } }) @@ -1439,13 +1464,15 @@ sglang:gen_throughput{model_name="glm",tp_rank="7"} 1147 assert_eq!(metrics.sum("sglang:num_queue_reqs"), 32.0); let response = - WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:").unwrap(); + WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:", 37) + .unwrap(); assert_eq!(response.dp_rank_count, 1); assert_eq!(response.loads[0].num_running_reqs, 28); assert_eq!(response.loads[0].num_waiting_reqs, 6); assert_eq!(response.loads[0].num_total_reqs, 34); assert_eq!(response.loads[0].num_waiting_uncached_tokens, 151477); assert_eq!(response.loads[0].num_used_tokens, 1420277); + assert_eq!(response.loads[0].max_running_requests, 37); assert!((response.effective_token_usage() - 0.98).abs() < 1e-12); assert_eq!(response.total_gen_throughput(), 1147.0); } @@ -1464,7 +1491,8 @@ sglang:gen_throughput{model_name="glm",tp_rank="7"} 1147 ); let response = - WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:").unwrap(); + WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:", 37) + .unwrap(); assert_eq!(response.dp_rank_count, 2); assert_eq!( response @@ -1490,10 +1518,51 @@ sglang:gen_throughput{model_name="glm",tp_rank="7"} 1147 .sum::(), 7 ); + assert!(response + .loads + .iter() + .all(|load| load.max_running_requests == 18)); + assert_eq!( + response + .loads + .iter() + .map(|load| load.max_running_requests) + .sum::(), + 36 + ); assert!((response.effective_token_usage() - 0.3).abs() < f64::EPSILON); assert!((response.total_gen_throughput() - 30.0).abs() < f64::EPSILON); } + #[test] + fn sglang_standard_load_prefers_reported_capacity_over_metadata_fallback() { + let standard: Vec = serde_json::from_str( + r#"[{"dp_rank":0,"num_reqs":34,"num_waiting_reqs":6,"max_running_requests":45}]"#, + ) + .unwrap(); + let metrics = PromScrape::parse("sglang:token_usage 0.5\n"); + + let response = + WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:", 37) + .unwrap(); + assert_eq!(response.loads[0].max_running_requests, 45); + } + + #[test] + fn sglang_standard_load_keeps_partial_rank_capacity_incomplete() { + let standard: Vec = serde_json::from_str( + r#"[{"dp_rank":0,"num_reqs":34,"max_running_requests":45},{"dp_rank":1,"num_reqs":7}]"#, + ) + .unwrap(); + let metrics = PromScrape::parse("sglang:token_usage 0.5\n"); + + let response = + WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:", 90) + .unwrap(); + assert_eq!(response.loads[0].max_running_requests, 45); + assert_eq!(response.loads[1].max_running_requests, 0); + } + #[test] fn sglang_v054_underscore_prefix_is_recognized() { // SGLang v0.5.4+ renamed the metric prefix `sglang:` -> `sglang_`. From 11c6a2f63933bd51421b69a8e2b200c954ec57a9 Mon Sep 17 00:00:00 2001 From: David <12414531+DavidBellamy@users.noreply.github.com> Date: Fri, 7 Aug 2026 16:38:38 -0700 Subject: [PATCH 3/3] Revert "fix(cache): surface legacy scheduler capacity" This reverts commit b828e2a3fb2f17e1c0d9c6792cab5ba536041b67. --- model_gateway/src/worker/monitor.rs | 79 ++--------------------------- 1 file changed, 5 insertions(+), 74 deletions(-) diff --git a/model_gateway/src/worker/monitor.rs b/model_gateway/src/worker/monitor.rs index b57310f0e..14a9aa86a 100644 --- a/model_gateway/src/worker/monitor.rs +++ b/model_gateway/src/worker/monitor.rs @@ -593,12 +593,7 @@ impl WorkerMonitor { .find(|p| m.has(&format!("{p}token_usage")))?; if let Some(loads) = standard_loads { - return Self::combine_sglang_standard_loads( - loads, - &m, - prefix, - worker.max_running_requests().map_or(0, i32::from), - ); + return Self::combine_sglang_standard_loads(loads, &m, prefix); } Some(Self::single_rank(SchedulerLoadSnapshot { @@ -628,18 +623,11 @@ impl WorkerMonitor { /// Prometheus values are averaged across emitted samples. Assigning the /// average to every DP snapshot preserves the aggregate throughput when /// metrics are per-DP, while avoiding multiplication when a TP-only - /// engine repeats the same sample on every tensor rank. Older SGLang - /// versions omit `max_running_requests` from this canonical endpoint, so - /// divide the instance-wide capacity discovered from `/server_info` across - /// a complete set of missing per-DP caps. Mixed reported/missing caps stay - /// incomplete so scheduler-pressure routing fails open. The Prometheus-only - /// fallback remains unchanged and does not receive this capacity, because - /// its request gauges may be repeated by TP rank. + /// engine repeats the same sample on every tensor rank. fn combine_sglang_standard_loads( loads: Vec, metrics: &PromScrape, prefix: &str, - fallback_max_running_requests: i32, ) -> Option { if loads.is_empty() { return None; @@ -654,15 +642,6 @@ impl WorkerMonitor { let cache_hit_rate = metrics.mean(&format!("{prefix}cache_hit_rate")); let utilization = metrics.mean(&format!("{prefix}utilization")); let dp_rank_count = i32::try_from(loads.len()).unwrap_or(i32::MAX); - let all_caps_missing = loads.iter().all(|load| load.max_running_requests <= 0); - let fallback_per_dp = if all_caps_missing { - fallback_max_running_requests - .max(0) - .checked_div(dp_rank_count) - .unwrap_or(0) - } else { - 0 - }; let loads = loads .into_iter() .map(|load| { @@ -679,11 +658,7 @@ impl WorkerMonitor { gen_throughput, cache_hit_rate, utilization, - max_running_requests: if load.max_running_requests > 0 { - load.max_running_requests - } else { - fallback_per_dp - }, + max_running_requests: load.max_running_requests.max(0), ..Default::default() } }) @@ -1464,15 +1439,13 @@ sglang:gen_throughput{model_name="glm",tp_rank="7"} 1147 assert_eq!(metrics.sum("sglang:num_queue_reqs"), 32.0); let response = - WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:", 37) - .unwrap(); + WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:").unwrap(); assert_eq!(response.dp_rank_count, 1); assert_eq!(response.loads[0].num_running_reqs, 28); assert_eq!(response.loads[0].num_waiting_reqs, 6); assert_eq!(response.loads[0].num_total_reqs, 34); assert_eq!(response.loads[0].num_waiting_uncached_tokens, 151477); assert_eq!(response.loads[0].num_used_tokens, 1420277); - assert_eq!(response.loads[0].max_running_requests, 37); assert!((response.effective_token_usage() - 0.98).abs() < 1e-12); assert_eq!(response.total_gen_throughput(), 1147.0); } @@ -1491,8 +1464,7 @@ sglang:gen_throughput{model_name="glm",tp_rank="7"} 1147 ); let response = - WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:", 37) - .unwrap(); + WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:").unwrap(); assert_eq!(response.dp_rank_count, 2); assert_eq!( response @@ -1518,51 +1490,10 @@ sglang:gen_throughput{model_name="glm",tp_rank="7"} 1147 .sum::(), 7 ); - assert!(response - .loads - .iter() - .all(|load| load.max_running_requests == 18)); - assert_eq!( - response - .loads - .iter() - .map(|load| load.max_running_requests) - .sum::(), - 36 - ); assert!((response.effective_token_usage() - 0.3).abs() < f64::EPSILON); assert!((response.total_gen_throughput() - 30.0).abs() < f64::EPSILON); } - #[test] - fn sglang_standard_load_prefers_reported_capacity_over_metadata_fallback() { - let standard: Vec = serde_json::from_str( - r#"[{"dp_rank":0,"num_reqs":34,"num_waiting_reqs":6,"max_running_requests":45}]"#, - ) - .unwrap(); - let metrics = PromScrape::parse("sglang:token_usage 0.5\n"); - - let response = - WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:", 37) - .unwrap(); - assert_eq!(response.loads[0].max_running_requests, 45); - } - - #[test] - fn sglang_standard_load_keeps_partial_rank_capacity_incomplete() { - let standard: Vec = serde_json::from_str( - r#"[{"dp_rank":0,"num_reqs":34,"max_running_requests":45},{"dp_rank":1,"num_reqs":7}]"#, - ) - .unwrap(); - let metrics = PromScrape::parse("sglang:token_usage 0.5\n"); - - let response = - WorkerMonitor::combine_sglang_standard_loads(standard, &metrics, "sglang:", 90) - .unwrap(); - assert_eq!(response.loads[0].max_running_requests, 45); - assert_eq!(response.loads[1].max_running_requests, 0); - } - #[test] fn sglang_v054_underscore_prefix_is_recognized() { // SGLang v0.5.4+ renamed the metric prefix `sglang:` -> `sglang_`.