Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
14 changes: 14 additions & 0 deletions crates/protocols/src/worker.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1145,6 +1145,20 @@ impl WorkerLoadResponse {
self.loads.iter().map(|l| l.num_used_tokens as i64).sum()
}

/// Total in-flight requests (running + waiting) summed across all DP ranks.
///
/// This is the worker's *global* request backlog as reported by the backend
/// itself — unlike a gateway's local in-flight counter, it counts requests
/// dispatched by every gateway, so it is invariant to the number of gateway
/// replicas. The `cache_aware` policy uses it to keep imbalance detection
/// and shortest-queue selection consistent across a multi-gateway deployment.
pub fn total_running_waiting_reqs(&self) -> i64 {
self.loads
.iter()
.map(|l| i64::from(l.num_running_reqs) + i64::from(l.num_waiting_reqs))
.sum()
}

/// Total queued (waiting, uncached) tokens summed across all DP ranks.
pub fn total_waiting_uncached_tokens(&self) -> i64 {
self.loads
Expand Down
14 changes: 6 additions & 8 deletions model_gateway/src/app_context.rs
Original file line number Diff line number Diff line change
Expand Up @@ -576,7 +576,7 @@ impl AppContextBuilder {
.client
.as_ref()
.ok_or_else(|| "client must be set before load monitor".to_string())?;
self.worker_monitor = Some(Arc::new(WorkerMonitor::new(
let worker_monitor = Arc::new(WorkerMonitor::new(
self.worker_registry
.as_ref()
.ok_or_else(|| "worker_registry must be set before load monitor".to_string())?
Expand All @@ -588,7 +588,11 @@ impl AppContextBuilder {
client.clone(),
config.load_monitor_interval_secs,
config.engine_metrics,
)));
));
if let Some(ref registry) = self.policy_registry {
registry.set_load_receiver(Some(worker_monitor.subscribe()));
}
self.worker_monitor = Some(worker_monitor);
Ok(self)
}

Expand Down Expand Up @@ -660,12 +664,6 @@ impl AppContextBuilder {
// and any other existing cache-aware policies.
if let Some(ref registry) = self.policy_registry {
registry.set_kv_event_monitor(Some(Arc::clone(&monitor)));
// Wire the backend load snapshot so cache-aware policies can use
// the KV-usage imbalance trigger. `with_worker_monitor` ran
// earlier in the build chain, so this is already set.
if let Some(ref worker_monitor) = self.worker_monitor {
registry.set_load_receiver(Some(worker_monitor.subscribe()));
}
}

self.kv_event_monitor = Some(monitor);
Expand Down
Loading
Loading