Repository navigation
feat(load-monitor): add dedicated config, /v1/loads endpoint, and gRPC support - #552
Conversation
…C support The load monitor previously reused worker_startup_check_interval_secs for its polling interval, only supported HTTP workers (gRPC returned -1), and used the deprecated /get_load endpoint. - Add load_monitor_interval_secs config field (default: 10s) to RouterConfig with builder method and > 0 validation - Add --load-monitor-interval CLI arg in main.rs (help_heading: Load Monitoring) - Update app_context.rs to use config.load_monitor_interval_secs - Add get_loads() to SglangSchedulerClient wrapping the GetLoads gRPC RPC - Add get_loads() to GrpcClient wrapper (SGLang only, sums num_used_tokens) - Rename parse_load_response -> fetch_http_load, take &Worker instead of separate url/api_key params - Update HTTP endpoint from /get_load to /v1/loads?include=core with new response format (loads[].num_used_tokens) - Add fetch_grpc_load() using Worker::get_grpc_client() for gRPC workers - Replace is_http bool + if/else with match on ConnectionMode enum - Add load_monitor_interval to Python RouterArgs dataclass + CLI arg - Add load_monitor_interval to PyO3 bindings (struct, constructor, builder) Signed-off-by: Simon Lin <simon@lightseek.ai> Signed-off-by: Simo Lin <linsimo.mark@gmail.com>
|
Note Currently processing new changes in this PR. This may take a few minutes, please wait... 📒 Files selected for processing (2)
✏️ Tip: You can disable in-progress messages and the fortune message in your review settings. 📝 WalkthroughWalkthroughThis pull request introduces a configurable Changes
Sequence Diagram(s)sequenceDiagram
participant LM as LoadMonitor
participant WM as WorkerManager
participant HTTP as HTTP Backend
participant GRPC as gRPC Backend
participant Worker as Worker Instance
LM->>LM: Sleep for load_monitor_interval secs
LM->>WM: get_all_worker_loads()
alt Connection Mode: HTTP
WM->>HTTP: GET /v1/loads?include=core
HTTP-->>WM: {loads: [{num_used_tokens: X}, ...]}
WM->>WM: Sum num_used_tokens
else Connection Mode: gRPC
WM->>Worker: get_grpc_client()
Worker-->>WM: GrpcClient(Sglang)
WM->>GRPC: get_loads(include: ["core"])
GRPC-->>WM: GetLoadsResponse{loads: [...]}
WM->>WM: Sum num_used_tokens from loads
end
WM-->>LM: WorkerLoadInfo{load: aggregated_tokens}
LM->>LM: Update metrics & schedule accordingly
Estimated code review effort🎯 3 (Moderate) | ⏱️ ~25 minutes Possibly related PRs
Suggested labels
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 2 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (2 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 capabilities by introducing a dedicated configuration for the polling interval, updating the HTTP load fetching mechanism to a more modern endpoint, and crucially, adding support for gRPC-based load data collection. These changes ensure that load data is accurately and consistently gathered across both HTTP and gRPC worker types, improving the overall reliability and functionality of the load balancing system. 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
|
There was a problem hiding this comment.
Code Review
This pull request introduces a dedicated configuration for the load monitor, migrates to the new /v1/loads endpoint for HTTP load fetching, and adds support for fetching load from gRPC workers. The changes are well-implemented and address the issues described. I've suggested a couple of minor refactorings in model_gateway/src/core/worker_manager.rs to improve code readability by reducing nesting in the error handling logic for fetch_http_load and fetch_grpc_load, aligning with best practices for match statement refactoring.
There was a problem hiding this comment.
Actionable comments posted: 3
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
bindings/python/src/lib.rs (1)
691-699:⚠️ Potential issue | 🟠 MajorInserting the new parameter mid-signature creates a potential backward-incompatibility risk for external positional callers.
Although the codebase exclusively uses keyword argument unpacking (
_Router(**args_dict)), this PyO3 binding is a public API. Placingload_monitor_intervalbeforecache_threshold(lines 698–699) could break external code using positional arguments. Recommend either appending this parameter at the end of the signature or enforcing keyword-only usage with a documented migration path.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@bindings/python/src/lib.rs` around lines 691 - 699, The new parameter load_monitor_interval was inserted before cache_threshold in the pyo3(signature = ...) for the router constructor, which can break external positional callers; fix by moving load_monitor_interval to the end of the signature (after cache_threshold) so existing positional argument order stays stable, or alternatively make the constructor keyword-only (so callers must use names) and document the change; update the pyo3(signature = ...) attribute that wraps the _Router constructor (the signature tuple containing worker_urls, policy, host, port, worker_startup_timeout_secs, worker_startup_check_interval, load_monitor_interval, cache_threshold) to implement your chosen approach and update any docstrings/comments accordingly.
🤖 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/config/types.rs`:
- Line 24: The new required field load_monitor_interval_secs should get a serde
default to preserve backward compatibility; add the attribute #[serde(default =
"default_load_monitor_interval_secs")] to the load_monitor_interval_secs field
so deserialization uses the existing Default impl value (10) when the key is
absent, and ensure a function named default_load_monitor_interval_secs exists
and returns the u64 default used by the Default impl.
In `@model_gateway/src/core/worker_manager.rs`:
- Around line 218-232: fetch_http_load currently swallows all
HTTP/transport/JSON errors and returns -1, making triage hard; update the match
handling around req.send().await and subsequent branches to log/record
diagnostics (e.g., using the same observability mechanism used by the gRPC path)
including the request context, HTTP status code when r.status() is not
successful, response body or text when available, and any JSON deserialization
error from r.json::<Value>().await, and still return -1 on failure; locate the
logic in fetch_http_load where req and r are used and add structured
logging/error reporting for the Err case of req.send().await, the non-success
status branch, and the Err case of r.json::<Value>().await so failures carry
actionable details.
In `@model_gateway/src/routers/grpc/client.rs`:
- Around line 159-170: The current get_loads implementation accumulates per-rank
num_used_tokens into an i32 which can overflow; update the Sglang branch to sum
into a 64-bit accumulator instead (e.g., map resp.loads.iter().map(|l|
l.num_used_tokens as i64).sum::<i64>() or use checked_add while iterating) to
avoid silent overflow, then convert the final i64 total to isize using a
fallible conversion (TryFrom) and return an Err if the conversion would
overflow; modify the symbols get_loads, Sglang, resp.loads, and num_used_tokens
accordingly so aggregation uses i64 and safe conversion to isize.
---
Outside diff comments:
In `@bindings/python/src/lib.rs`:
- Around line 691-699: The new parameter load_monitor_interval was inserted
before cache_threshold in the pyo3(signature = ...) for the router constructor,
which can break external positional callers; fix by moving load_monitor_interval
to the end of the signature (after cache_threshold) so existing positional
argument order stays stable, or alternatively make the constructor keyword-only
(so callers must use names) and document the change; update the pyo3(signature =
...) attribute that wraps the _Router constructor (the signature tuple
containing worker_urls, policy, host, port, worker_startup_timeout_secs,
worker_startup_check_interval, load_monitor_interval, cache_threshold) to
implement your chosen approach and update any docstrings/comments accordingly.
ℹ️ Review info
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
📒 Files selected for processing (10)
bindings/python/src/lib.rsbindings/python/src/smg/router_args.pygrpc_client/src/sglang_scheduler.rsmodel_gateway/src/app_context.rsmodel_gateway/src/config/builder.rsmodel_gateway/src/config/types.rsmodel_gateway/src/config/validation.rsmodel_gateway/src/core/worker_manager.rsmodel_gateway/src/main.rsmodel_gateway/src/routers/grpc/client.rs
| /// Get the total token load from the backend. | ||
| /// Only supported for SGLang backends. Returns summed num_used_tokens across all DP ranks. | ||
| pub async fn get_loads(&self) -> Result<isize, Box<dyn std::error::Error + Send + Sync>> { | ||
| match self { | ||
| Self::Sglang(client) => { | ||
| let resp = client.get_loads(vec!["core".to_string()]).await?; | ||
| let total: i32 = resp.loads.iter().map(|l| l.num_used_tokens).sum(); | ||
| Ok(total as isize) | ||
| } | ||
| _ => Err("GetLoads RPC not supported for this backend".into()), | ||
| } | ||
| } |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Search for num_used_tokens usage and get_loads calls
echo "=== Searching for num_used_tokens usage ==="
rg -n "num_used_tokens" --type rust --type proto
echo -e "\n=== Searching for get_loads calls ==="
rg -n "get_loads\s*\(" --type rust
echo -e "\n=== Looking for proto definitions ==="
find . -name "*.proto" | head -20Repository: lightseekorg/smg
Length of output: 1127
🏁 Script executed:
#!/bin/bash
# Read proto file to understand field types
echo "=== Content of sglang_scheduler.proto ==="
cat ./grpc_client/proto/sglang_scheduler.proto
echo -e "\n=== Check worker_manager.rs usage of get_loads ==="
sed -n '235,250p' model_gateway/src/core/worker_manager.rs
echo -e "\n=== Check get_loads in server.rs ==="
sed -n '425,435p' model_gateway/src/server.rsRepository: lightseekorg/smg
Length of output: 14953
🏁 Script executed:
#!/bin/bash
# Check test cases and typical DP rank counts
echo "=== Test for get_loads ==="
cat model_gateway/tests/api/api_endpoints_test.rs | sed -n '1220,1270p'
echo -e "\n=== Look for DP rank usage patterns ==="
rg -n "dp_rank\|data_parallel" --type rust -A 2 -B 2 | head -60
echo -e "\n=== Check if there's any overflow handling elsewhere ==="
rg -n "checked_add\|overflow\|saturating" --type rust | head -20Repository: lightseekorg/smg
Length of output: 1644
🏁 Script executed:
#!/bin/bash
# Read the actual client.rs get_loads implementation in full context
echo "=== Full client.rs get_loads implementation ==="
sed -n '159,175p' model_gateway/src/routers/grpc/client.rs
# Check how loads are typically used downstream
echo -e "\n=== Full fetch_grpc_load context in worker_manager.rs ==="
sed -n '230,260p' model_gateway/src/core/worker_manager.rs
# Look for any load balancing logic that depends on isize values
echo -e "\n=== Search for load balancing with get_loads results ==="
rg -n "fetch_grpc_load\|get_loads" model_gateway/src/core/worker_manager.rs -A 5 -B 2
# Check response types and JSON serialization
echo -e "\n=== Check how loads are returned in HTTP response ==="
rg -n "get_all_worker_loads\|GetLoadsRequest" model_gateway/src --type rust -A 10 | head -50Repository: lightseekorg/smg
Length of output: 1986
🏁 Script executed:
#!/bin/bash
# Check if there are multiple DP ranks being handled
echo "=== Check GetLoadsRequest and DP rank filtering ==="
rg -n "GetLoadsRequest\|dp_rank" grpc_client/src/sglang_scheduler.rs -A 3 -B 3
# Look at how the response loads are iterated
echo -e "\n=== Check sglang_scheduler.rs get_loads implementation ==="
sed -n '285,310p' grpc_client/src/sglang_scheduler.rs
# Check what i32::MAX is and typical token counts
echo -e "\n=== Search for token count values in code/tests ==="
rg -n "num_used_tokens\|num_tokens\|token.*count" model_gateway/tests --type rust -A 2 | head -40
# Look for max_total_num_tokens to understand scale
echo -e "\n=== Check token capacity and max values ==="
rg -n "max_total_num_tokens\|max.*token\|token.*capacity" --type rust -B 1 -A 1 | head -30Repository: lightseekorg/smg
Length of output: 1061
🏁 Script executed:
#!/bin/bash
# Look for typical deployment configurations and scaling patterns
echo "=== Check for configuration docs and deployment info ==="
find . -name "*.md" -o -name "*.yaml" -o -name "*.toml" | xargs grep -l "dp_rank\|data_parallel\|distributed" 2>/dev/null | head -5
# Check if there are any saturating or checked operations used elsewhere for token counts
echo -e "\n=== Look at how similar aggregations are done in the codebase ==="
rg -n "\.sum\(\)" --type rust -B 3 -A 1 | head -40
# Check if i32::MAX could realistically be exceeded
echo -e "\n=== Look for large token count values ==="
rg -n "100000\|1000000\|1e6\|1e7\|1e8" --type rust | head -20
# Check proto message definitions for other numeric aggregations
echo -e "\n=== Check AggregateMetrics in proto ==="
sed -n '600,620p' grpc_client/proto/sglang_scheduler.protoRepository: lightseekorg/smg
Length of output: 3073
🏁 Script executed:
#!/bin/bash
# Read the AggregateMetrics section and understand what it contains
echo "=== Reading AggregateMetrics from proto (lines 595-615) ==="
sed -n '595,615p' grpc_client/proto/sglang_scheduler.proto
# Check if there are any examples or documentation about typical DP rank counts
echo -e "\n=== Search for DP rank documentation/examples ==="
rg -n "dp_rank\|data.parallel\|distributed" docs/ --type md -A 2 -B 2 | head -50
# Look at the full GetLoadsResponse structure
echo -e "\n=== Full GetLoadsResponse structure ==="
sed -n '516,545p' grpc_client/proto/sglang_scheduler.proto
# Check actual test values to see realistic token counts
echo -e "\n=== Look at mock worker responses ==="
rg -n "mock.*load\|num_used_tokens\|SchedulerLoad" model_gateway/tests --type rust -B 2 -A 2 | head -40Repository: lightseekorg/smg
Length of output: 996
🏁 Script executed:
#!/bin/bash
# Look for proto import statements and see what the actual field types are
echo "=== Check proto imports and ensure int32 is standard ==="
head -20 grpc_client/proto/sglang_scheduler.proto
# Verify the exact aggregation logic - does it sum ALL ranks or just request "core"?
echo -e "\n=== Check the vec![\"core\".to_string()] parameter ==="
rg -n "vec!\[\"core\"" model_gateway/src/routers/grpc/client.rs -B 3 -A 3
# Let me verify if "core" means single rank or all ranks
echo -e "\n=== Check if include parameter filters to single rank or includes all ==="
grep -n "include\|GetLoadsRequest" grpc_client/proto/sglang_scheduler.proto -A 5 | head -20
# Check actual test implementation to see what response looks like
echo -e "\n=== Check the mock implementation or test helper ==="
rg -n "GetLoadsResponse\|SchedulerLoad" model_gateway/tests --type rust -B 2 -A 3Repository: lightseekorg/smg
Length of output: 1801
🏁 Script executed:
#!/bin/bash
# Let me verify the actual usage pattern more directly
echo "=== Check if GetLoadsResponse.loads is plural (multiple entries) ==="
sed -n '504,520p' grpc_client/proto/sglang_scheduler.proto
# Confirm the comment about "Returns summed num_used_tokens across all DP ranks"
echo -e "\n=== Re-read the full docstring and implementation ==="
sed -n '159,170p' model_gateway/src/routers/grpc/client.rs
# Check if there's any documentation about typical DP rank counts/token values
echo -e "\n=== Look for deployment guides mentioning DP ranks ==="
find ./docs -name "*.md" -exec grep -l "distributed\|scale\|rank" {} \; | head -5
# Quick verification: is resp.loads a Vec or single item?
echo -e "\n=== Verify GetLoadsResponse proto fields ==="
sed -n '500,512p' grpc_client/proto/sglang_scheduler.protoRepository: lightseekorg/smg
Length of output: 2192
🏁 Script executed:
#!/bin/bash
# Final verification: confirm that large token counts are realistic
echo "=== Check max_total_num_tokens to understand token capacity scale ==="
sed -n '516,535p' grpc_client/proto/sglang_scheduler.proto
echo -e "\n=== Verify that casting to isize preserves the corrupted value ==="
python3 << 'PYEOF'
# Simulate the overflow
i32_max = 2**31 - 1 # 2,147,483,647
print(f"i32::MAX = {i32_max}")
# Scenario: 8 DP ranks with 400M tokens each
total = 8 * 400_000_000
print(f"8 ranks × 400M tokens = {total}")
# Overflow wrapping behavior in i32
overflowed = total % (2**32)
if overflowed > i32_max:
overflowed = overflowed - 2**32
print(f"After i32 overflow: {overflowed}")
print(f"This corrupted value affects scheduling: smaller value -> more work assigned")
PYEOFRepository: lightseekorg/smg
Length of output: 824
Prevent overflow in get_loads token aggregation.
Line 165 sums per-rank loads into i32 and line 166 casts directly to isize; with multiple DP ranks (common in distributed deployments), the total can exceed i32::MAX, causing silent overflow and corrupting scheduling decisions.
Consider using i64 with checked arithmetic:
💡 Proposed fix
pub async fn get_loads(&self) -> Result<isize, Box<dyn std::error::Error + Send + Sync>> {
match self {
Self::Sglang(client) => {
let resp = client.get_loads(vec!["core".to_string()]).await?;
- let total: i32 = resp.loads.iter().map(|l| l.num_used_tokens).sum();
- Ok(total as isize)
+ let total = resp.loads.iter().try_fold(0_i64, |acc, l| {
+ let tokens = i64::from(l.num_used_tokens);
+ if tokens < 0 {
+ return Err("GetLoads returned negative num_used_tokens".into());
+ }
+ acc.checked_add(tokens).ok_or("GetLoads total overflowed i64".into())
+ })?;
+ let total =
+ isize::try_from(total).map_err(|_| "GetLoads total does not fit in isize".into())?;
+ Ok(total)
}
_ => Err("GetLoads RPC not supported for this backend".into()),
}
}📝 Committable suggestion
‼️ IMPORTANT
Carefully review the code before committing. Ensure that it accurately replaces the highlighted code, contains no missing lines, and has no issues with indentation. Thoroughly test & benchmark the code to ensure it meets the requirements.
| /// Get the total token load from the backend. | |
| /// Only supported for SGLang backends. Returns summed num_used_tokens across all DP ranks. | |
| pub async fn get_loads(&self) -> Result<isize, Box<dyn std::error::Error + Send + Sync>> { | |
| match self { | |
| Self::Sglang(client) => { | |
| let resp = client.get_loads(vec!["core".to_string()]).await?; | |
| let total: i32 = resp.loads.iter().map(|l| l.num_used_tokens).sum(); | |
| Ok(total as isize) | |
| } | |
| _ => Err("GetLoads RPC not supported for this backend".into()), | |
| } | |
| } | |
| /// Get the total token load from the backend. | |
| /// Only supported for SGLang backends. Returns summed num_used_tokens across all DP ranks. | |
| pub async fn get_loads(&self) -> Result<isize, Box<dyn std::error::Error + Send + Sync>> { | |
| match self { | |
| Self::Sglang(client) => { | |
| let resp = client.get_loads(vec!["core".to_string()]).await?; | |
| let total = resp.loads.iter().try_fold(0_i64, |acc, l| { | |
| let tokens = i64::from(l.num_used_tokens); | |
| if tokens < 0 { | |
| return Err("GetLoads returned negative num_used_tokens".into()); | |
| } | |
| acc.checked_add(tokens).ok_or("GetLoads total overflowed i64".into()) | |
| })?; | |
| let total = | |
| isize::try_from(total).map_err(|_| "GetLoads total does not fit in isize".into())?; | |
| Ok(total) | |
| } | |
| _ => Err("GetLoads RPC not supported for this backend".into()), | |
| } | |
| } |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@model_gateway/src/routers/grpc/client.rs` around lines 159 - 170, The current
get_loads implementation accumulates per-rank num_used_tokens into an i32 which
can overflow; update the Sglang branch to sum into a 64-bit accumulator instead
(e.g., map resp.loads.iter().map(|l| l.num_used_tokens as i64).sum::<i64>() or
use checked_add while iterating) to avoid silent overflow, then convert the
final i64 total to isize using a fallible conversion (TryFrom) and return an Err
if the conversion would overflow; modify the symbols get_loads, Sglang,
resp.loads, and num_used_tokens accordingly so aggregation uses i64 and safe
conversion to isize.
- Flatten nested match in fetch_http_load and fetch_grpc_load with early returns, making the happy path more linear - Add #[serde(default = "default_load_monitor_interval_secs")] so existing config files without the field deserialize with default 10 Signed-off-by: Simon Lin <simon@lightseek.ai> 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: 8be2d4b0f0
ℹ️ 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".
| } | ||
| }; | ||
|
|
||
| match grpc_client.get_loads().await { |
There was a problem hiding this comment.
Bound GetLoads RPC with timeout in load monitor
The new gRPC load path awaits grpc_client.get_loads().await with no deadline, while the HTTP path uses REQUEST_TIMEOUT. If a gRPC worker becomes unresponsive but keeps the connection open, get_all_worker_loads can hang inside future::join_all, which blocks /get_loads and prevents PowerOfTwo policies from receiving any load updates from other healthy workers. Please wrap this RPC in tokio::time::timeout (or configure a per-request gRPC timeout) so a stuck backend degrades to -1 instead of stalling the entire load collection pass.
Useful? React with 👍 / 👎.
…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>
…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>
Description
Problem
The load monitor had several issues:
worker_startup_check_interval_secs(30s) instead of having its own dedicated polling interval-1(no load data collected)/get_loadreturning[{num_tokens: N}]instead of the new/v1/loadsendpoint (sglang#16976) andGetLoadsgRPC RPC (sglang#17087)Solution
load_monitor_interval_secsconfig (default 10s) with CLI arg--load-monitor-intervalGET /v1/loads?include=corewith new response format (loads[].num_used_tokens)GetLoadsRPC through the existing per-workerGrpcClientmatchonConnectionModeenum instead of boolean flag for HTTP/gRPC dispatchChanges
model_gateway/src/config/types.rs— Addload_monitor_interval_secs: u64field (default: 10)model_gateway/src/config/builder.rs— Addload_monitor_interval_secs()builder methodmodel_gateway/src/config/validation.rs— Add> 0validationmodel_gateway/src/main.rs— Add--load-monitor-intervalCLI argmodel_gateway/src/app_context.rs— Useconfig.load_monitor_interval_secsinstead of startup check intervalgrpc_client/src/sglang_scheduler.rs— Addget_loads()method wrapping theGetLoadsRPCmodel_gateway/src/routers/grpc/client.rs— Addget_loads()toGrpcClient(SGLang only, sumsnum_used_tokens)model_gateway/src/core/worker_manager.rs— Rename tofetch_http_load(takes&Worker), update to/v1/loads?include=core, addfetch_grpc_load(), match onConnectionModebindings/python/src/smg/router_args.py— Addload_monitor_intervalfield +--load-monitor-intervalCLI argbindings/python/src/lib.rs— Addload_monitor_intervalto PyO3 struct, constructor, and builder wiringTest Plan
cargo build— compiles cleanlycargo test -p smg— 16 passedcargo test -p smg-grpc-client— 28 passedcargo clippy -p smg -p smg-grpc-client— no warningspytest bindings/python/tests/— 213 passedChecklist
cargo +nightly fmtpassescargo clippy --all-targets --all-features -- -D warningspassesSummary by CodeRabbit
New Features
Configuration