feat(load-monitor): per-group polling with rich load data - #554
Conversation
📝 WalkthroughWalkthroughReworks load monitoring into per-worker-group polling via a new Changes
Sequence DiagramsequenceDiagram
participant Register as Worker Registration
participant LoadMonitor
participant GroupLoop as Group Monitor Loop
participant GrpcClient as gRPC Client
participant Policy as PowerOfTwo Policy
Register->>LoadMonitor: on_group_added(WorkerGroupKey, interval?)
activate LoadMonitor
LoadMonitor->>GroupLoop: spawn per-group polling task
activate GroupLoop
loop Poll cycle
GroupLoop->>GrpcClient: fetch_grpc_load(worker)
activate GrpcClient
GrpcClient-->>GroupLoop: WorkerLoadResponse
deactivate GrpcClient
GroupLoop->>GroupLoop: aggregate per-group loads (HashMap<String, WorkerLoadResponse>)
GroupLoop->>Policy: update_loads(HashMap<String, WorkerLoadResponse>)
Policy-->>GroupLoop: (cache updated)
GroupLoop->>LoadMonitor: tx.send_modify(merged loads)
end
Register->>LoadMonitor: on_group_removed(WorkerGroupKey, worker_urls)
LoadMonitor->>GroupLoop: abort/stop handle
deactivate GroupLoop
deactivate LoadMonitor
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
Suggested reviewers
🚥 Pre-merge checks | ✅ 3✅ Passed checks (3 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches
🧪 Generate unit tests (beta)
Comment |
Summary of ChangesHello, I'm Gemini Code Assist1! I'm currently reviewing this pull request and will post my feedback shortly. In the meantime, here's a summary to help you and other reviewers quickly get up to speed! This pull request significantly enhances the load monitoring and balancing capabilities by moving from a single global polling loop with scalar load values to a more granular, per-group polling system. This change addresses previous limitations by providing richer, per-DP-rank load data, enabling more intelligent routing decisions based on normalized token usage, and allowing for flexible, per-worker polling intervals. The overall impact is a more efficient and fair distribution of requests across workers, especially in heterogeneous environments. Highlights
Changelog
Activity
Using Gemini Code AssistThe full guide for Gemini Code Assist can be found on our documentation page, here are some quick tips. Invoking Gemini You can request assistance from Gemini at any point by creating a comment using either
Customization To customize Gemini Code Assist for GitHub experience, repository maintainers can create a configuration file and/or provide a custom code review style guide (such as PEP-8 for Python) by creating and adding files to a Limitations & Feedback Gemini Code Assist may make mistakes. Please leave feedback on any instances where its feedback is incorrect or counter productive. You can react with 👍 and 👎 on @gemini-code-assist comments. If you're interested in giving your feedback about your experience with Gemini Code Assist for Github and other Google products, sign up here. You can also get AI-powered code generation, chat, as well as code reviews directly in the IDE at no cost with the Gemini Code Assist IDE Extension. Footnotes
|
|
Hi @slin1237, this PR has merge conflicts that must be resolved before it can be merged. Please rebase your branch: git fetch origin main
git rebase origin/main
# resolve any conflicts, then:
git push --force-with-lease |
…sage routing Replace the single-loop load monitor with per-group polling and upgrade from scalar isize load values to full WorkerLoadResponse data. What changed: - protocols/src/worker.rs: Add SchedulerLoadSnapshot, WorkerLoadResponse, and WorkerGroupKey types. Add load_monitor_interval_secs to WorkerSpec. Update WorkerLoadInfo with optional details field. - grpc_client/src/sglang_scheduler.rs: Add From<proto> impls to convert gRPC SchedulerLoad and GetLoadsResponse to protocol types. - model_gateway/src/routers/grpc/client.rs: Update get_loads() return type from isize to WorkerLoadResponse. - model_gateway/src/policies/mod.rs: Change update_loads trait signature to HashMap<String, WorkerLoadResponse>. - model_gateway/src/policies/power_of_two.rs: Compare by effective_token_usage() ratio (0.0-1.0) instead of raw token counts. Fallback to request counts when either worker lacks cached load data. - model_gateway/src/core/worker_manager.rs: Rewrite LoadMonitor with per-group polling loops keyed by (model_id, worker_type, connection_mode). Groups start/stop automatically via on_group_added/on_group_removed. Simplify fetch_http_load to use serde deserialization directly. - model_gateway/src/core/steps/worker/shared/register.rs: Notify LoadMonitor on worker group creation with per-worker interval override. - model_gateway/src/core/steps/worker/local/remove_from_worker_registry.rs: Notify LoadMonitor when a worker group becomes empty. - model_gateway/src/server.rs: Remove load_monitor.start() call since groups start dynamically on worker registration. - model_gateway/src/core/mod.rs: Re-export WorkerGroupKey from protocols. Why: The previous load monitor used a single polling loop for all workers and stored only a scalar sum of num_used_tokens. This prevented per-group interval control and lost rich per-DP-rank metrics (token_usage ratio, throughput, cache hit rate, utilization). Workers with different capacities could not be compared fairly using raw token counts. How: Each unique (model_id, worker_type, connection_mode) group gets its own tokio task polling at a configurable interval (per-worker override via WorkerSpec.load_monitor_interval_secs, falling back to global default). Groups merge results into a shared watch channel using send_modify for atomic updates without cross-group clobbering. PowerOfTwo now compares workers by effective_token_usage() (average across DP ranks) when both have cached data, degrading both to request counts when either is missing to prevent incompatible metric comparison. Refs: #552 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
5760379 to
9b5a802
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 57603796a5
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| for policy in &power_of_two_policies { | ||
| policy.update_loads(&group_loads); | ||
| } |
There was a problem hiding this comment.
Pass full load map when updating PowerOfTwo caches
This call updates each policy with only one group's group_loads, but PowerOfTwoPolicy::update_loads replaces its entire cache via clone_from (see model_gateway/src/policies/power_of_two.rs), so every group tick wipes metrics from other groups. In contexts where one PowerOfTwo instance serves multiple groups (for example prefill/decode groups for the same model), routing for the non-last-updated group loses token-usage data and degrades to request-count behavior unexpectedly.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Same root cause as the coderabbit finding. Fixed in d9d106c by switching clone_from to extend.
| if handles.contains_key(&key) { | ||
| debug!("Load monitor group already running: {key}"); | ||
| return; |
There was a problem hiding this comment.
Restart existing group monitor when interval changes
Returning early when a group handle already exists prevents interval reconfiguration after the first worker is registered. Since registration computes load_monitor_interval_secs from newly added workers, adding a worker with a smaller interval to an existing group cannot lower that group's poll cadence, so the implementation does not honor the documented "minimum interval wins" behavior until process restart.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Pushing back. Changing interval at runtime requires removing and re-adding the worker (which triggers group lifecycle). Hot-reload for intervals is extra complexity for an edge case.
There was a problem hiding this comment.
Code Review
This pull request significantly enhances load monitoring by introducing per-group polling and using richer load data for more intelligent routing decisions, transitioning to dynamic, per-group monitoring loops and improving the PowerOfTwo policy. However, several critical issues have been identified: a zero-value interval in the load monitor can lead to a CPU-exhaustion Denial of Service; there's a memory leak in the cleanup logic for removed worker groups; a logic error in the PowerOfTwoPolicy cache update mechanism can disable high-fidelity load balancing; and an inefficiency in finding the minimum polling interval for new worker groups. These issues require attention to ensure reliability, performance, and security.
There was a problem hiding this comment.
Actionable comments posted: 4
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/core/worker_manager.rs`:
- Around line 234-236: The gRPC fetch path must mirror the HTTP path's
empty-load semantics: add the same check that returns None when
response.loads.is_empty() to the gRPC handling code so empty payloads are
treated as failures consistently; update the gRPC fetch/response handling (where
response.loads is inspected) to return None for empty loads, and apply the same
change to the other occurrence noted around the second block (256-260) so both
paths use identical empty-load logic.
- Around line 487-514: The current merge leaves stale entries for group members
that failed this tick; inside the tx.send_modify closure clear any existing
entries for this group's worker URLs before inserting the new successful loads
so old snapshots are removed. Concretely, in the tx.send_modify(|map| { ... })
closure iterate over the group's worker identifiers (the same URLs in the
workers collection used earlier) and remove them from map (e.g., map.remove(url)
or map.retain(...) for those keys), then map.extend(group_loads); keep
references to group_loads, workers and group_key as in the surrounding code and
ensure removal happens atomically before the extend.
- Around line 322-331: on_group_added() constructs a tokio::time::interval using
per-worker WorkerSpec.load_monitor_interval_secs which can be 0 and will panic;
before creating the timer, validate or clamp the effective interval (use the
WorkerManager.default_interval as fallback) so any value <= 0 becomes
Duration::from_secs(1) (or at least 1s), then use that Duration when calling
tokio::time::interval(); apply the same guard where per-worker overrides are
read/used (e.g., when computing effective_interval from
WorkerSpec.load_monitor_interval_secs and default_interval) to ensure no
zero-second Duration is ever passed to tokio::time::interval().
In `@model_gateway/src/policies/power_of_two.rs`:
- Around line 116-119: The update_loads method currently calls
cached.clone_from(loads), which replaces the entire cached_loads HashMap with a
single group's snapshot and drops other groups' entries; instead, acquire the
write lock on cached_loads and merge the incoming loads into the existing cache
(e.g., iterate over loads and upsert each (key, value) into cached), updating or
inserting per-key entries (WorkerLoadResponse) rather than wholesale replacing
the map so concurrent group polls don't erase other groups' data.
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
📒 Files selected for processing (10)
grpc_client/src/sglang_scheduler.rsmodel_gateway/src/core/mod.rsmodel_gateway/src/core/steps/worker/local/remove_from_worker_registry.rsmodel_gateway/src/core/steps/worker/shared/register.rsmodel_gateway/src/core/worker_manager.rsmodel_gateway/src/policies/mod.rsmodel_gateway/src/policies/power_of_two.rsmodel_gateway/src/routers/grpc/client.rsmodel_gateway/src/server.rsprotocols/src/worker.rs
There was a problem hiding this comment.
♻️ Duplicate comments (2)
model_gateway/src/core/worker_manager.rs (2)
256-261:⚠️ Potential issue | 🟠 MajorgRPC path should also treat empty loads as failure for consistency.
HTTP
fetch_http_loadreturnsNonewhenresponse.loads.is_empty()(lines 234–236), butfetch_grpc_loadreturnsSome(load)unconditionally. This creates transport-dependent routing behavior for the same logical state.Suggested fix
match grpc_client.get_loads().await { - Ok(load) => Some(load), + Ok(load) if !load.loads.is_empty() => Some(load), + Ok(_) => None, Err(e) => { debug!("gRPC GetLoads failed for {}: {e}", worker.url()); None } }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/core/worker_manager.rs` around lines 256 - 261, fetch_grpc_load (the code calling grpc_client.get_loads()) currently returns Some(load) even when the returned load has no entries; make it consistent with fetch_http_load by treating empty loads as failure: after Ok(load) check whether load.loads.is_empty() (or the equivalent empty-check on the returned struct) and if empty log a debug message (similar to the existing Err branch) and return None, otherwise return Some(load).
510-513:⚠️ Potential issue | 🟠 MajorStale load entries persist for workers that fail within a polling cycle.
map.extend(group_loads)inserts successful fetches but doesn't remove entries for workers that fail this tick. If a worker succeeds at tick N and fails at tick N+1, its stale load data remains in the map until the worker is removed from the registry, potentially causing routing decisions based on outdated information.Industry load balancers (AWS NLB, Kubernetes) address this by removing unhealthy endpoints from rotation after consecutive failures, rather than masking them with cached data. Consider clearing this group's entries before extending:
Suggested fix
+ let group_urls: Vec<String> = + workers.iter().map(|w| w.url().to_string()).collect(); tx.send_modify(|map| { + for url in &group_urls { + map.remove(url); + } map.extend(group_loads); });🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/core/worker_manager.rs` around lines 510 - 513, The current tx.send_modify closure uses map.extend(group_loads) which adds successful fetches but leaves stale entries for workers that failed this tick; inside the same atomic closure (tx.send_modify(|map| { ... })), remove or replace the existing entries for this group before extending — e.g., iterate map keys and remove those whose group id matches this group's id (or clear the group's sub-map) and then map.extend(group_loads) so only current successful workers remain; apply this change in the code that calls tx.send_modify (referenced symbols: tx.send_modify, map.extend, group_loads, and the group identifier used in the map keys/registry).
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Duplicate comments:
In `@model_gateway/src/core/worker_manager.rs`:
- Around line 256-261: fetch_grpc_load (the code calling
grpc_client.get_loads()) currently returns Some(load) even when the returned
load has no entries; make it consistent with fetch_http_load by treating empty
loads as failure: after Ok(load) check whether load.loads.is_empty() (or the
equivalent empty-check on the returned struct) and if empty log a debug message
(similar to the existing Err branch) and return None, otherwise return
Some(load).
- Around line 510-513: The current tx.send_modify closure uses
map.extend(group_loads) which adds successful fetches but leaves stale entries
for workers that failed this tick; inside the same atomic closure
(tx.send_modify(|map| { ... })), remove or replace the existing entries for this
group before extending — e.g., iterate map keys and remove those whose group id
matches this group's id (or clear the group's sub-map) and then
map.extend(group_loads) so only current successful workers remain; apply this
change in the code that calls tx.send_modify (referenced symbols:
tx.send_modify, map.extend, group_loads, and the group identifier used in the
map keys/registry).
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
📒 Files selected for processing (2)
model_gateway/src/core/steps/worker/local/remove_from_worker_registry.rsmodel_gateway/src/core/worker_manager.rs
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 6a4534b04a
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| if group_loads.is_empty() { | ||
| debug!("No loads fetched for group {group_key}"); | ||
| continue; | ||
| } |
There was a problem hiding this comment.
Clear group cache when polling returns no load data
If all workers in a group fail a poll cycle (for example, temporary timeouts or backend restart), group_loads is empty and this early continue skips policy.update_loads. That leaves the previous successful snapshot cached in PowerOfTwoPolicy, so routing keeps comparing stale token_usage values instead of degrading to request-count fallback when load data is missing. This makes transient monitor outages bias routing with outdated metrics until a later successful poll.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Pushing back. Keeping last-known-good data during transient failures is intentional. The watch channel now clears and re-inserts per tick (d9d106c), so the subscriber view stays current. The policy cache retaining last-known-good is a deliberate design choice — degrading all routing to request-count fallback when one worker briefly times out is worse.
| .filter_map(|w| w.metadata().spec.load_monitor_interval_secs) | ||
| .min(); |
There was a problem hiding this comment.
Include default interval when computing group minimum
This minimum only considers explicit per-worker overrides, so workers with load_monitor_interval_secs = None are excluded from the calculation. In a mixed group (e.g., one worker with Some(30) and another with None while global is 10), this passes Some(30) to on_group_added, so the group polls slower than the global fallback and violates the documented “minimum effective interval wins” behavior.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Pushing back. Per-worker interval is an explicit override. In a mixed group [Some(30), None] with global=10, we pass Some(30) — the user set 30s intentionally. Workers without an override defer to whatever the group picks. The global default only kicks in when no worker specifies an interval.
What changed: - worker_manager.rs: Clamp load monitor interval to minimum 1 second, preventing a tight-loop DoS when load_monitor_interval_secs is 0. - worker_manager.rs: Change on_group_removed() to accept worker_urls parameter from the caller instead of querying the registry (which is already empty by the time this method is called). - remove_from_worker_registry.rs: Collect worker URLs per group before removal and pass them to on_group_removed() for proper cleanup. Why: - Duration::from_secs(0) creates a timer that ticks immediately and repeatedly, consuming 100% CPU. The interval is user-configurable via WorkerSpec, so a malicious or misconfigured value must be clamped. - on_group_removed() was querying get_workers_filtered() after workers had already been removed from the registry, so stale load entries were never cleaned up from the shared watch channel map. Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
6a4534b to
952cf9d
Compare
There was a problem hiding this comment.
Actionable comments posted: 1
♻️ Duplicate comments (2)
model_gateway/src/core/worker_manager.rs (2)
256-260:⚠️ Potential issue | 🟠 MajorAlign gRPC empty-load handling with the HTTP path.
Line 257 accepts empty payloads as success, while HTTP treats empty
loadsas failure (Line 234-236). This creates transport-dependent load semantics.Suggested fix
- match grpc_client.get_loads().await { - Ok(load) => Some(load), + match grpc_client.get_loads().await { + Ok(load) if !load.loads.is_empty() => Some(load), + Ok(_) => None, Err(e) => { debug!("gRPC GetLoads failed for {}: {e}", worker.url()); None } }🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/core/worker_manager.rs` around lines 256 - 260, The gRPC path currently treats an Ok(empty_load) as success while the HTTP path returns failure for empty loads; modify the grpc_client.get_loads() Ok branch in worker_manager.rs to mirror HTTP behavior by checking the returned load's payload (e.g., inspect load.loads or the appropriate field on the response) and return None when it's empty, otherwise return Some(load); keep the existing debug/error logging (include worker.url() as before) so empty responses are treated consistently across transports.
494-497:⚠️ Potential issue | 🟠 MajorPurge current-group entries before merge to prevent stale snapshots.
Line 512 only extends successful loads, and Line 496 skips updates when all fetches fail. Failed members retain old snapshots indefinitely.
Suggested fix
- if group_loads.is_empty() { - debug!("No loads fetched for group {group_key}"); - continue; - } + let group_urls: Vec<String> = + workers.iter().map(|w| w.url().to_string()).collect(); + + // Always clear current group keys first, then reinsert successful loads. + tx.send_modify(|map| { + for url in &group_urls { + map.remove(url); + } + map.extend(group_loads.clone()); + }); + + if group_loads.is_empty() { + debug!("No loads fetched for group {group_key}"); + continue; + } @@ - tx.send_modify(|map| { - map.extend(group_loads); - }); + // already merged aboveAlso applies to: 510-513
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/core/worker_manager.rs` around lines 494 - 497, Currently when group_loads.is_empty() the code continues without clearing existing snapshots, leaving stale entries; before the early continue and before you extend/merge successful loads you must explicitly purge the current-group entries for group_key (e.g., remove or clear the group_key bucket in the snapshot/store/map) so failed fetches don't leave old snapshots; in short, call the appropriate clear/remove for the current-group storage for group_key prior to checking group_loads.is_empty() and likewise clear it immediately before you extend/merge the new successful loads instead of only appending.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/core/steps/worker/local/remove_from_worker_registry.rs`:
- Around line 111-124: The current cleanup only invokes
app_context.load_monitor.on_group_removed when pool_size == 0, so when a subset
of workers is removed from a non-empty group their URLs are never pruned; update
remove_from_worker_registry logic to always compute the removed URLs for the
group (using urls_by_group and the WorkerGroupKey built from model_id,
worker_type and connection_mode) and call the load monitor to prune those
specific URLs even when pool_size > 0 (e.g., invoke an existing or new load
monitor method to remove specific worker URLs such as
on_workers_removed/on_urls_removed with the key and removed_urls instead of only
calling on_group_removed for empty groups).
---
Duplicate comments:
In `@model_gateway/src/core/worker_manager.rs`:
- Around line 256-260: The gRPC path currently treats an Ok(empty_load) as
success while the HTTP path returns failure for empty loads; modify the
grpc_client.get_loads() Ok branch in worker_manager.rs to mirror HTTP behavior
by checking the returned load's payload (e.g., inspect load.loads or the
appropriate field on the response) and return None when it's empty, otherwise
return Some(load); keep the existing debug/error logging (include worker.url()
as before) so empty responses are treated consistently across transports.
- Around line 494-497: Currently when group_loads.is_empty() the code continues
without clearing existing snapshots, leaving stale entries; before the early
continue and before you extend/merge successful loads you must explicitly purge
the current-group entries for group_key (e.g., remove or clear the group_key
bucket in the snapshot/store/map) so failed fetches don't leave old snapshots;
in short, call the appropriate clear/remove for the current-group storage for
group_key prior to checking group_loads.is_empty() and likewise clear it
immediately before you extend/merge the new successful loads instead of only
appending.
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
📒 Files selected for processing (2)
model_gateway/src/core/steps/worker/local/remove_from_worker_registry.rsmodel_gateway/src/core/worker_manager.rs
| // If the group is now empty, stop its load monitor and clean up entries | ||
| if pool_size == 0 { | ||
| if let Some(ref load_monitor) = app_context.load_monitor { | ||
| let key = WorkerGroupKey { | ||
| model_id: model_id.clone(), | ||
| worker_type: *worker_type, | ||
| connection_mode: *connection_mode, | ||
| }; | ||
| let removed_urls = urls_by_group | ||
| .get(&(*worker_type, *connection_mode, model_id.clone())) | ||
| .map(|v| v.as_slice()) | ||
| .unwrap_or_default(); | ||
| load_monitor.on_group_removed(&key, removed_urls).await; | ||
| } |
There was a problem hiding this comment.
Cleanup is skipped for removed workers when the group still has members.
Line 112 gates cleanup behind pool_size == 0. When a subset is removed, those URLs are not pruned here, so stale load entries can persist for active groups.
Suggested direction
- // If the group is now empty, stop its load monitor and clean up entries
- if pool_size == 0 {
- if let Some(ref load_monitor) = app_context.load_monitor {
+ if let Some(ref load_monitor) = app_context.load_monitor {
+ let removed_urls = urls_by_group
+ .get(&(*worker_type, *connection_mode, model_id.clone()))
+ .map(|v| v.as_slice())
+ .unwrap_or_default();
+
+ // New API: prune removed URLs even when group remains active.
+ load_monitor.on_workers_removed(removed_urls).await;
+
+ if pool_size == 0 {
let key = WorkerGroupKey {
model_id: model_id.clone(),
worker_type: *worker_type,
connection_mode: *connection_mode,
};
- let removed_urls = urls_by_group
- .get(&(*worker_type, *connection_mode, model_id.clone()))
- .map(|v| v.as_slice())
- .unwrap_or_default();
load_monitor.on_group_removed(&key, removed_urls).await;
}
}🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/core/steps/worker/local/remove_from_worker_registry.rs`
around lines 111 - 124, The current cleanup only invokes
app_context.load_monitor.on_group_removed when pool_size == 0, so when a subset
of workers is removed from a non-empty group their URLs are never pruned; update
remove_from_worker_registry logic to always compute the removed URLs for the
group (using urls_by_group and the WorkerGroupKey built from model_id,
worker_type and connection_mode) and call the load monitor to prune those
specific URLs even when pool_size > 0 (e.g., invoke an existing or new load
monitor method to remove specific worker URLs such as
on_workers_removed/on_urls_removed with the key and removed_urls instead of only
calling on_group_removed for empty groups).
There was a problem hiding this comment.
Pushing back. Stale entries for removed workers are cleaned up on the next poll tick — group_monitor_loop clears all group URLs before inserting (d9d106c), and removed workers wont appear in get_workers_filtered, so their URLs wont be re-inserted. Cleanup window is bounded to one interval period.
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Find the group_monitor_loop implementation and the clearing logic
# Look for the section that clears group URLs before inserting new load data
# First, let's see the group_monitor_loop function
ast-grep --pattern $'async fn group_monitor_loop($$$) {
$$$
}'
# Also search for the commit message mentioned by user
git log --oneline --all | grep -i "d9d106c5\|clear\|stale"
# Look for send_modify usage which should show the merge logic
rg -n -A10 -B5 "send_modify" model_gateway/src/core/worker_manager.rs
Length of output: 9054
@slin1237, you're absolutely right. I can see in lines 514-520 of model_gateway/src/core/worker_manager.rs that group_monitor_loop collects all_group_urls from the current get_workers_filtered result, then uses send_modify to remove all those URLs before extending with fresh group_loads. Since removed workers won't appear in the filtered list, their URLs are automatically pruned on the next tick. The comment at line 512 confirms this design. My concern was incorrect—apologies for the noise.
🧠 Learnings used
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: protocols/src/worker.rs:338-343
Timestamp: 2026-02-21T02:36:06.251Z
Learning: Repo: lightseekorg/smg PR: 489
File: protocols/src/worker.rs (impl From<Vec<ModelCard>> for WorkerModels)
Learning: For the single-element case (models.len() == 1), keep the defensive Option guard with a non-panicking fallback to WorkerModels::Wildcard instead of using expect()/unwrap(). This is intentional to avoid introducing panic paths in production code for PR `#489` and similar clippy/lint-only efforts.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/policies/bucket.rs:116-137
Timestamp: 2026-02-21T02:31:03.274Z
Learning: For clippy-focused PRs in this repository, avoid behavior changes (e.g., load-balancing logic). Defer such changes to a separate PR/issue.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 399
File: protocols/src/interactions.rs:505-509
Timestamp: 2026-02-19T03:08:50.192Z
Learning: In code reviews for Rust projects using the validator crate (v0.20.0), ensure that custom validation functions for numeric primitive types (e.g., f32, i32, u32, i16, etc.) accept the value by value, not by reference. Example: fn validate(value: f32) { ... }. The validator derive macro has a hardcoded list of numeric types that are passed by value, while all other types are passed by reference. Apply this guideline whenever validating numeric fields to align with the derive macro behavior.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/core/token_bucket.rs:58-63
Timestamp: 2026-02-21T02:30:51.443Z
Learning: For lint-only/Clippy enforcement PRs in this repository, avoid introducing behavioral changes (e.g., new input validation or logic changes). Treat such PRs as non-functional changes and plan a separate follow-up issue/PR for hardening or behavior changes. This applies broadly to Rust files across the repo; during review, focus on lint/style corrections and clearly note any intentional exceptions.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: protocols/src/responses.rs:928-931
Timestamp: 2026-02-21T02:36:00.882Z
Learning: In Rust code across the repository, use the marker INVARIANT: to document assumptions in safe code. Reserve SAFETY: for explaining why unsafe blocks are sound. This improves clarity of invariants and safety reasoning. Example reference: protocols/src/responses.rs near validate_tool_choice_with_tools().
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: mesh/src/sync.rs:83-83
Timestamp: 2026-02-21T02:37:01.416Z
Learning: General Rust formatting rule: format! with implicit captures only supports simple identifiers, not full expressions like {state.model_id}. For cases where you want to interpolate a field or expression, bind the value first and interpolate the binding, e.g., let model_id = &state.model_id; and then use format!("policy:{}", model_id). In the specific file mesh/src/sync.rs, prefer format!("policy:{}", state.model_id) or bind to a local variable if you need named interpolation, to keep clarity and avoid unintended captures.
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 952cf9df12
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| && w.worker_type() == worker_type | ||
| && w.connection_mode() == connection_mode | ||
| }) | ||
| .filter_map(|w| w.metadata().spec.load_monitor_interval_secs) |
There was a problem hiding this comment.
Propagate per-worker monitor interval into worker metadata
This code derives the group polling interval from w.metadata().spec.load_monitor_interval_secs, but the worker creation paths still build workers via BasicWorkerBuilder::new(...) without ever copying WorkerSpec.load_monitor_interval_secs into metadata (e.g., CreateLocalWorkerStep / CreateExternalWorkersStep), so this filter always sees None and on_group_added always falls back to the global default. In practice, any user-specified per-worker override is silently ignored, so the new per-worker interval feature does not function.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Valid observation. The creation steps use individual builder setters and dont pass load_monitor_interval_secs through. The field exists on WorkerSpec and the monitoring code reads it, but the creation pipeline needs a builder setter or from_spec() usage. Will address in a follow-up.
…oad handling - power_of_two.rs: change update_loads from clone_from to extend so per-group updates don't wipe other groups' cached load data - worker_manager.rs group_monitor_loop: clear all group worker URLs from the watch channel before inserting fresh loads, so workers that failed this tick don't retain stale snapshots - worker_manager.rs fetch_grpc_load: add empty-load guard matching HTTP path behavior — return None when loads vec is empty Refs: #554 Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: d9d106c514
ℹ️ About Codex in GitHub
Codex has been enabled to automatically review pull requests in this repo. Reviews are triggered when you
- Open a pull request for review
- Mark a draft as ready
- Comment "@codex review".
If Codex has suggestions, it will comment; otherwise it will react with 👍.
When you sign up for Codex through ChatGPT, Codex can also answer questions or update the PR, like "@codex address that feedback".
| fn update_loads(&self, loads: &HashMap<String, WorkerLoadResponse>) { | ||
| if let Ok(mut cached) = self.cached_loads.write() { | ||
| cached.clone_from(loads); | ||
| cached.extend(loads.iter().map(|(k, v)| (k.clone(), v.clone()))); |
There was a problem hiding this comment.
Evict stale worker loads when poll results are partial
group_monitor_loop builds group_loads from only successful fetches, but PowerOfTwoPolicy::update_loads only extends the cache and never removes workers that are absent in the latest tick. If one worker in a compared pair times out (or is briefly unreachable), its old WorkerLoadResponse remains in cached_loads, so select_worker still uses token_usage instead of falling back to request counts, routing on stale utilization data.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Same reasoning as my reply to 2863781483 — keeping last-known-good data is intentional. If one worker times out, its stale ratio is better than degrading both workers to request-count fallback. Staleness is bounded to one poll interval.
| ); | ||
|
|
||
| // If the group is now empty, stop its load monitor and clean up entries | ||
| if pool_size == 0 { |
There was a problem hiding this comment.
Clean removed worker loads even when group is non-empty
This cleanup path only runs when pool_size == 0, so removing a worker from a still-active group skips on_group_removed. The monitor loop later removes only URLs that still exist in that group, so departed URLs are never purged from the shared load snapshot for long-lived groups, causing stale-entry growth and allowing old load data to resurface if that URL is re-registered before the next successful poll.
Useful? React with 👍 / 👎.
There was a problem hiding this comment.
Duplicate of the earlier comment. Already addressed — the monitor loop clears and re-inserts on each tick (d9d106c), so departed URLs are cleaned up within one interval period.
There was a problem hiding this comment.
Actionable comments posted: 2
♻️ Duplicate comments (1)
model_gateway/src/core/worker_manager.rs (1)
495-498:⚠️ Potential issue | 🟠 MajorDo not skip stale-entry cleanup when a group has zero successful fetches.
Line 495’s early
continuebypasses the removal block in Line 511-Line 520. If every worker in the group fails in a tick, old snapshots for that group remain in the shared map.Suggested fix
- if group_loads.is_empty() { - debug!("No loads fetched for group {group_key}"); - continue; - } + if group_loads.is_empty() { + debug!("No loads fetched for group {group_key}"); + } debug!( "Fetched loads from {}/{} workers in group {group_key}", group_loads.len(), workers.len() );Also applies to: 511-520
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@model_gateway/src/core/worker_manager.rs` around lines 495 - 498, The early continue when group_loads.is_empty() skips the stale-entry cleanup, so change control flow to always run the removal block for the group even if group_loads is empty: do not return/continue before the stale-snapshot removal that uses group_key and the shared snapshots map; instead, move or call the cleanup logic (the block that removes old snapshots for the group) before any continue/return or guard, or invert the condition so the cleanup always executes and only the subsequent per-load processing is skipped when group_loads.is_empty().
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@model_gateway/src/core/worker_manager.rs`:
- Around line 466-486: The current code spawns an unbounded set of worker load
fetch futures and awaits them with future::join_all, which can overload large
groups; change it to cap per-group concurrency using the same MAX_CONCURRENT
used for fan_out: convert the per-worker mapping into a futures stream (e.g.,
futures::stream::iter or StreamExt::map) and use
buffer_unordered(MAX_CONCURRENT) to run at most MAX_CONCURRENT in-flight
fetches, then collect the results; keep the existing match on
group_key.connection_mode and calls to WorkerManager::fetch_http_load and
WorkerManager::fetch_grpc_load and preserve the (worker.url().to_string(),
response) tuple shape.
In `@model_gateway/src/policies/power_of_two.rs`:
- Around line 116-119: update_loads currently only upserts entries into the
policy cache (cached_loads) and never removes workers that are no longer
reported, causing stale entries to persist; change update_loads to treat the
provided loads as the authoritative set for the group by removing any keys from
cached_loads that are not present in the incoming loads map (i.e., compute the
difference between cached_loads keys and loads.keys and remove the missing ones
before/after extend), and then call this updated update_loads from
LoadMonitor::group_monitor_loop using the same all_group_urls used for watch-map
cleanup so cache pruning is aligned with group membership.
---
Duplicate comments:
In `@model_gateway/src/core/worker_manager.rs`:
- Around line 495-498: The early continue when group_loads.is_empty() skips the
stale-entry cleanup, so change control flow to always run the removal block for
the group even if group_loads is empty: do not return/continue before the
stale-snapshot removal that uses group_key and the shared snapshots map;
instead, move or call the cleanup logic (the block that removes old snapshots
for the group) before any continue/return or guard, or invert the condition so
the cleanup always executes and only the subsequent per-load processing is
skipped when group_loads.is_empty().
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
📒 Files selected for processing (2)
model_gateway/src/core/worker_manager.rsmodel_gateway/src/policies/power_of_two.rs
| let futures: Vec<_> = workers | ||
| .iter() | ||
| .map(|worker| { | ||
| let client = client.clone(); | ||
| let worker = Arc::clone(worker); | ||
| let connection_mode = group_key.connection_mode; | ||
|
|
||
| async move { | ||
| let response = match connection_mode { | ||
| ConnectionMode::Http => { | ||
| WorkerManager::fetch_http_load(&client, &worker).await | ||
| } | ||
| ConnectionMode::Grpc => WorkerManager::fetch_grpc_load(&worker).await, | ||
| }; | ||
| (worker.url().to_string(), response) | ||
| } | ||
| }) | ||
| .collect(); | ||
|
|
||
| let results = future::join_all(futures).await; | ||
|
|
There was a problem hiding this comment.
Cap per-group load fetch concurrency instead of join_all over all workers.
This path is currently unbounded, so large groups can create too many in-flight calls in one tick. Reuse MAX_CONCURRENT here (same as fan_out) to avoid burst pressure.
Suggested fix
- let futures: Vec<_> = workers
- .iter()
- .map(|worker| {
- let client = client.clone();
- let worker = Arc::clone(worker);
- let connection_mode = group_key.connection_mode;
-
- async move {
- let response = match connection_mode {
- ConnectionMode::Http => {
- WorkerManager::fetch_http_load(&client, &worker).await
- }
- ConnectionMode::Grpc => WorkerManager::fetch_grpc_load(&worker).await,
- };
- (worker.url().to_string(), response)
- }
- })
- .collect();
-
- let results = future::join_all(futures).await;
+ let results: Vec<_> = stream::iter(workers.iter().map(|worker| {
+ let client = client.clone();
+ let worker = Arc::clone(worker);
+ let connection_mode = group_key.connection_mode;
+ async move {
+ let response = match connection_mode {
+ ConnectionMode::Http => WorkerManager::fetch_http_load(&client, &worker).await,
+ ConnectionMode::Grpc => WorkerManager::fetch_grpc_load(&worker).await,
+ };
+ (worker.url().to_string(), response)
+ }
+ }))
+ .buffer_unordered(MAX_CONCURRENT)
+ .collect()
+ .await;🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/core/worker_manager.rs` around lines 466 - 486, The current
code spawns an unbounded set of worker load fetch futures and awaits them with
future::join_all, which can overload large groups; change it to cap per-group
concurrency using the same MAX_CONCURRENT used for fan_out: convert the
per-worker mapping into a futures stream (e.g., futures::stream::iter or
StreamExt::map) and use buffer_unordered(MAX_CONCURRENT) to run at most
MAX_CONCURRENT in-flight fetches, then collect the results; keep the existing
match on group_key.connection_mode and calls to WorkerManager::fetch_http_load
and WorkerManager::fetch_grpc_load and preserve the (worker.url().to_string(),
response) tuple shape.
There was a problem hiding this comment.
Pushing back. A single worker group has a handful of workers (typically 1-8). join_all on 8 futures is fine. Buffered concurrency would only matter for groups with 100+ workers, which is not a realistic deployment pattern.
There was a problem hiding this comment.
@slin1237, that makes sense—for 1-8 workers per group, join_all is perfectly reasonable and keeps the code simpler. I'll defer to your domain knowledge on typical deployment patterns here.
🧠 Learnings used
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: protocols/src/worker.rs:338-343
Timestamp: 2026-02-21T02:36:06.251Z
Learning: Repo: lightseekorg/smg PR: 489
File: protocols/src/worker.rs (impl From<Vec<ModelCard>> for WorkerModels)
Learning: For the single-element case (models.len() == 1), keep the defensive Option guard with a non-panicking fallback to WorkerModels::Wildcard instead of using expect()/unwrap(). This is intentional to avoid introducing panic paths in production code for PR `#489` and similar clippy/lint-only efforts.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/benches/wasm_middleware_latency.rs:88-91
Timestamp: 2026-02-21T02:37:02.009Z
Learning: Repo: lightseekorg/smg — For clippy-only/enforcement PRs (e.g., PR `#489`), even micro-optimizations (like replacing an async closure with std::future::ready in benches such as model_gateway/benches/wasm_middleware_latency.rs) should be deferred to a follow-up PR rather than included inline.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/policies/bucket.rs:116-137
Timestamp: 2026-02-21T02:31:03.274Z
Learning: For clippy-focused PRs in this repository, avoid behavior changes (e.g., load-balancing logic). Defer such changes to a separate PR/issue.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 399
File: protocols/src/interactions.rs:505-509
Timestamp: 2026-02-19T03:08:50.192Z
Learning: In code reviews for Rust projects using the validator crate (v0.20.0), ensure that custom validation functions for numeric primitive types (e.g., f32, i32, u32, i16, etc.) accept the value by value, not by reference. Example: fn validate(value: f32) { ... }. The validator derive macro has a hardcoded list of numeric types that are passed by value, while all other types are passed by reference. Apply this guideline whenever validating numeric fields to align with the derive macro behavior.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/core/token_bucket.rs:58-63
Timestamp: 2026-02-21T02:30:51.443Z
Learning: For lint-only/Clippy enforcement PRs in this repository, avoid introducing behavioral changes (e.g., new input validation or logic changes). Treat such PRs as non-functional changes and plan a separate follow-up issue/PR for hardening or behavior changes. This applies broadly to Rust files across the repo; during review, focus on lint/style corrections and clearly note any intentional exceptions.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: protocols/src/responses.rs:928-931
Timestamp: 2026-02-21T02:36:00.882Z
Learning: In Rust code across the repository, use the marker INVARIANT: to document assumptions in safe code. Reserve SAFETY: for explaining why unsafe blocks are sound. This improves clarity of invariants and safety reasoning. Example reference: protocols/src/responses.rs near validate_tool_choice_with_tools().
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: mesh/src/sync.rs:83-83
Timestamp: 2026-02-21T02:37:01.416Z
Learning: General Rust formatting rule: format! with implicit captures only supports simple identifiers, not full expressions like {state.model_id}. For cases where you want to interpolate a field or expression, bind the value first and interpolate the binding, e.g., let model_id = &state.model_id; and then use format!("policy:{}", model_id). In the specific file mesh/src/sync.rs, prefer format!("policy:{}", state.model_id) or bind to a local variable if you need named interpolation, to keep clarity and avoid unintended captures.
| fn update_loads(&self, loads: &HashMap<String, WorkerLoadResponse>) { | ||
| if let Ok(mut cached) = self.cached_loads.write() { | ||
| cached.clone_from(loads); | ||
| cached.extend(loads.iter().map(|(k, v)| (k.clone(), v.clone()))); | ||
| } |
There was a problem hiding this comment.
update_loads upsert-only behavior leaves stale entries in policy cache.
With per-group polling, this method receives only successful samples. cached.extend(...) never removes missing workers, so stale token metrics can persist and prevent the intended request-count fallback on partial failures.
Suggested direction
+impl PowerOfTwoPolicy {
+ pub fn reconcile_group_loads(
+ &self,
+ group_urls: &[String],
+ fresh_loads: &HashMap<String, WorkerLoadResponse>,
+ ) {
+ if let Ok(mut cached) = self.cached_loads.write() {
+ for url in group_urls {
+ cached.remove(url);
+ }
+ cached.extend(fresh_loads.iter().map(|(k, v)| (k.clone(), v.clone())));
+ }
+ }
+}Then call this from LoadMonitor::group_monitor_loop using the same all_group_urls used for watch-map cleanup.
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/policies/power_of_two.rs` around lines 116 - 119,
update_loads currently only upserts entries into the policy cache (cached_loads)
and never removes workers that are no longer reported, causing stale entries to
persist; change update_loads to treat the provided loads as the authoritative
set for the group by removing any keys from cached_loads that are not present in
the incoming loads map (i.e., compute the difference between cached_loads keys
and loads.keys and remove the missing ones before/after extend), and then call
this updated update_loads from LoadMonitor::group_monitor_loop using the same
all_group_urls used for watch-map cleanup so cache pruning is aligned with group
membership.
There was a problem hiding this comment.
Same as 2864073561. Intentional design — policy cache retains last-known-good data. The watch channel (subscriber-facing) properly clears stale entries each tick. The policy cache keeping stale data is preferable to degrading all routing to request-count fallback during transient failures.
There was a problem hiding this comment.
@slin1237, understood — I appreciate the clarification on the design intent. Retaining last-known-good token data in the policy cache during transient failures is a sensible trade-off to maintain routing quality, and it's good to know the watch channel handles the authoritative cleanup each tick.
✏️ Learnings added
Learnt from: slin1237
Repo: lightseekorg/smg PR: 554
File: model_gateway/src/policies/power_of_two.rs:116-119
Timestamp: 2026-02-27T12:44:59.424Z
Learning: In model_gateway/src/policies/power_of_two.rs, the PowerOfTwoPolicy's cached_loads field intentionally uses extend-only behavior (never removes stale entries) to retain last-known-good load data. This avoids degrading all routing to request-count fallback during transient worker failures. The watch channel (subscriber-facing) properly clears stale entries each tick, but the policy cache keeping stale data is preferable for routing decisions.
Note: Learnings are effective only in the context of similar code segments. To apply general coding guidelines, please configure review instructions. You can manage existing learnings in the app.
🧠 Learnings used
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/policies/bucket.rs:116-137
Timestamp: 2026-02-21T02:31:03.274Z
Learning: For clippy-focused PRs in this repository, avoid behavior changes (e.g., load-balancing logic). Defer such changes to a separate PR/issue.
Learnt from: XinyueZhang369
Repo: lightseekorg/smg PR: 399
File: protocols/src/interactions.rs:505-509
Timestamp: 2026-02-19T03:08:50.192Z
Learning: In code reviews for Rust projects using the validator crate (v0.20.0), ensure that custom validation functions for numeric primitive types (e.g., f32, i32, u32, i16, etc.) accept the value by value, not by reference. Example: fn validate(value: f32) { ... }. The validator derive macro has a hardcoded list of numeric types that are passed by value, while all other types are passed by reference. Apply this guideline whenever validating numeric fields to align with the derive macro behavior.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: model_gateway/src/core/token_bucket.rs:58-63
Timestamp: 2026-02-21T02:30:51.443Z
Learning: For lint-only/Clippy enforcement PRs in this repository, avoid introducing behavioral changes (e.g., new input validation or logic changes). Treat such PRs as non-functional changes and plan a separate follow-up issue/PR for hardening or behavior changes. This applies broadly to Rust files across the repo; during review, focus on lint/style corrections and clearly note any intentional exceptions.
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: protocols/src/responses.rs:928-931
Timestamp: 2026-02-21T02:36:00.882Z
Learning: In Rust code across the repository, use the marker INVARIANT: to document assumptions in safe code. Reserve SAFETY: for explaining why unsafe blocks are sound. This improves clarity of invariants and safety reasoning. Example reference: protocols/src/responses.rs near validate_tool_choice_with_tools().
Learnt from: slin1237
Repo: lightseekorg/smg PR: 489
File: mesh/src/sync.rs:83-83
Timestamp: 2026-02-21T02:37:01.416Z
Learning: General Rust formatting rule: format! with implicit captures only supports simple identifiers, not full expressions like {state.model_id}. For cases where you want to interpolate a field or expression, bind the value first and interpolate the binding, e.g., let model_id = &state.model_id; and then use format!("policy:{}", model_id). In the specific file mesh/src/sync.rs, prefer format!("policy:{}", state.model_id) or bind to a local variable if you need named interpolation, to keep clarity and avoid unintended captures.
Summary
Follow-up to #552. Upgrades the load monitor from a single global polling loop with scalar load values to per-group polling with full
WorkerLoadResponsedata andtoken_usageratio-based routing.(model_id, worker_type, connection_mode)group, started/stopped automatically on worker registration/removalSchedulerLoadSnapshotandWorkerLoadResponsetypes capturing all per-DP-rank metrics (token_usage, throughput, cache_hit_rate, utilization, etc.)effective_token_usage()ratio (0.0–1.0) instead of raw token counts, with fallback to request counts when load data is missingWorkerSpec.load_monitor_interval_secsoverride, falling back to globalload_monitor_interval_secsWhat changed
protocols/src/worker.rsSchedulerLoadSnapshot,WorkerLoadResponse,WorkerGroupKeytypes; addload_monitor_interval_secstoWorkerSpec; adddetailsfield toWorkerLoadInfogrpc_client/src/sglang_scheduler.rsFrom<proto>impls forSchedulerLoad→SchedulerLoadSnapshotandGetLoadsResponse→WorkerLoadResponsemodel_gateway/src/routers/grpc/client.rsget_loads()returnsWorkerLoadResponseinstead ofisizemodel_gateway/src/policies/mod.rsupdate_loadstrait signature usesHashMap<String, WorkerLoadResponse>model_gateway/src/policies/power_of_two.rseffective_token_usage()ratio; fallback to request counts when either worker lacks datamodel_gateway/src/core/worker_manager.rsLoadMonitorwith per-group polling,WorkerGroupKey, serde-basedfetch_http_loadmodel_gateway/src/core/steps/worker/shared/register.rsLoadMonitor.on_group_added()after worker registrationmodel_gateway/src/core/steps/worker/local/remove_from_worker_registry.rsLoadMonitor.on_group_removed()when group becomes emptymodel_gateway/src/server.rsload_monitor.start()(groups start dynamically)Why
The previous load monitor had three limitations:
isize(sum ofnum_used_tokens), losing rich per-DP-rank metricstoken_usageratio (0.0–1.0) is the proper normalized metricHow
Each unique
(model_id, worker_type, connection_mode)group gets its own tokio task polling at a configurable interval. Groups merge results into a shared watch channel usingsend_modify()for atomic updates without cross-group clobbering. When workers in the same group specify different intervals, the minimum (fastest) wins. PowerOfTwo degrades both workers to request counts when either is missing cached load data, preventing incompatible metric comparison (the "mixed metric bug").Test plan
cargo build— all crates compilecargo test -p openai-protocol— 29 passedcargo test -p smg-grpc-client— 28 passedcargo test -p smg— 93 + 16 integration tests passedcargo clippy -p smg -p openai-protocol— no warningspytest bindings/python/tests/— 106 passedRefs: #552
Summary by CodeRabbit
New Features
Improvements