Repository navigation
feat(grpc): add FlushCache and profile RPCs with worker-abstracted admin ops - #1655
Conversation
📝 WalkthroughWalkthroughAdds scheduler admin operations (FlushCache, StartProfile, StopProfile): protobuf contracts, client helpers with local timeouts and trace injection, servicer handlers that dispatch communicator requests and aggregate per-rank results, GrpcRequestManager communicator routing API, gateway HTTP handlers and WorkerManager fan-out, worker-level admin implementations, and E2E tests validating both gRPC and HTTP modes. ChangesAdmin Operations Feature
Estimated code review effort🎯 4 (Complex) | ⏱️ ~60 minutes Possibly related PRs
Suggested reviewers
Poem
🚥 Pre-merge checks | ✅ 5✅ Passed checks (5 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches📝 Generate docstrings
🧪 Generate unit tests (beta)
Comment |
There was a problem hiding this comment.
Code Review
This pull request introduces gateway admin operations for cache flushing and profiling, enabling parallel fan-out to both HTTP and gRPC workers. Key changes include adding new protobuf messages and RPCs, implementing shared admin operations in the gRPC client, and updating the worker manager and gRPC servicer to support these operations. Feedback on the changes highlights several critical issues: a potential panic in the Rust client when converting non-finite float timeouts to durations, a similar float sanitization issue in the Python servicer, a logic bug where explicit profiling parameters are incorrectly overridden by environment variables, and a robust matching issue for targeted workers in data-parallel deployments.
Important
The consumer version of Gemini Code Assist on GitHub is being sunset. Starting June 18, 2026, new organization installations will be blocked, and all code review activity will officially cease on July 17, 2026.
For more details on the timeline and next steps, please review the Help Documentation.
| let deadline = std::time::Duration::from_secs_f32(timeout_s.max(0.0)) | ||
| + $crate::FLUSH_RPC_DEADLINE_MARGIN; |
There was a problem hiding this comment.
Using std::time::Duration::from_secs_f32 directly on timeout_s.max(0.0) can cause a panic if timeout_s is NaN or Infinity (since from_secs_f32 panics on non-finite values). To prevent potential denial-of-service panics, sanitize timeout_s to ensure it is finite before converting it to a Duration.
| let deadline = std::time::Duration::from_secs_f32(timeout_s.max(0.0)) | |
| + $crate::FLUSH_RPC_DEADLINE_MARGIN; | |
| let secs = if timeout_s.is_finite() { timeout_s.max(0.0) } else { 0.0 }; | |
| let deadline = std::time::Duration::from_secs_f32(secs) | |
| + $crate::FLUSH_RPC_DEADLINE_MARGIN; |
References
- Do not introduce panics in code that interacts with external systems if the upstream server does not handle the error. Instead, handle the error gracefully or propagate it appropriately.
| ) -> common_pb2.FlushCacheResponse: | ||
| """Flush the KV cache on all scheduler processes.""" | ||
| logger.debug("Receive flush cache request (timeout_s=%.1f)", request.timeout_s) | ||
| timeout_s = request.timeout_s |
There was a problem hiding this comment.
If request.timeout_s is NaN or negative, it can cause unexpected behavior in max(30.0, timeout_s + 10.0) and downstream timeout handling. Sanitize timeout_s to ensure it is a valid, non-negative number.
| timeout_s = request.timeout_s | |
| timeout_s = request.timeout_s if (request.timeout_s == request.timeout_s and request.timeout_s >= 0.0) else 0.0 |
| with_stack = request.with_stack if request.HasField("with_stack") else None | ||
| with_stack = (with_stack is not False) and get_bool_env_var( | ||
| "SGLANG_PROFILE_WITH_STACK", "true" | ||
| ) | ||
| record_shapes = request.record_shapes if request.HasField("record_shapes") else None | ||
| record_shapes = (record_shapes is not False) and get_bool_env_var( | ||
| "SGLANG_PROFILE_RECORD_SHAPES", "true" | ||
| ) |
There was a problem hiding this comment.
The current logic for resolving with_stack and record_shapes has a bug where an explicit user request of True is overridden by the environment variable if the environment variable is set to False (because (with_stack is not False) and get_bool_env_var(...) evaluates to the env var value when with_stack is True). Explicit request parameters should take precedence over environment variable defaults. Simplifying this to check HasField and fallback to the env var only when unset resolves this issue and improves readability.
with_stack = request.with_stack if request.HasField("with_stack") else get_bool_env_var(
"SGLANG_PROFILE_WITH_STACK", "true"
)
record_shapes = request.record_shapes if request.HasField("record_shapes") else get_bool_env_var(
"SGLANG_PROFILE_RECORD_SHAPES", "true"
)| match url_filter { | ||
| Some(url) => workers.into_iter().filter(|w| w.url() == url).collect(), | ||
| None => workers, | ||
| } |
There was a problem hiding this comment.
In data-parallel (DP) deployments, workers are registered with a DP rank suffix (e.g., @0 or @1) appended to their URL. If a user targets a specific worker using its base URL (without the suffix), w.url() == url will fail to match. Checking both w.url() and w.base_url() ensures robust matching for targeted admin operations in DP deployments.
| match url_filter { | |
| Some(url) => workers.into_iter().filter(|w| w.url() == url).collect(), | |
| None => workers, | |
| } | |
| match url_filter { | |
| Some(url) => workers | |
| .into_iter() | |
| .filter(|w| w.url() == url || w.base_url() == url) | |
| .collect(), | |
| None => workers, | |
| } |
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 60c40d5e32
ℹ️ 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".
| let workers = Self::admin_target_workers(worker_registry, worker_url); | ||
| let total_workers = workers.len(); |
There was a problem hiding this comment.
Deduplicate DP-aware gRPC workers before profiling
In a dp_aware gRPC deployment, create_worker registers one logical worker per DP rank (URLs suffixed with @rank), but each SGLang StartProfile RPC already fans out to all scheduler DP ranks via the servicer's communicator. Calling admin_target_workers here therefore sends the same profile-start request once per logical rank to the same backend, so multi-DP gRPC workers can report partial failures such as “already profiling” even though a single RPC would have profiled every rank. The start/stop/flush admin paths should target each gRPC backend once (or otherwise avoid re-fanning-out over logical DP workers).
Useful? React with 👍 / 👎.
|
👋 The PR description doesn't fully follow
Please update the PR description so reviewers have the context they need. |
| # Env-var overrides mirror sglang's TokenizerCommunicatorMixin defaults. | ||
| with_stack = request.with_stack if request.HasField("with_stack") else None | ||
| with_stack = (with_stack is not False) and get_bool_env_var( | ||
| "SGLANG_PROFILE_WITH_STACK", "true" | ||
| ) |
There was a problem hiding this comment.
🟡 Nit: The env-var logic here acts as a mask (AND), not a fallback default. If a caller explicitly sends with_stack=True in the gRPC request, an operator who sets SGLANG_PROFILE_WITH_STACK=false will silently override it to False:
# request.with_stack = True → with_stack = True
# (True is not False) and get_bool_env_var(..., "false") → False
Same issue applies to record_shapes below.
If the intent is "use request value when set, otherwise fall back to env var":
| # Env-var overrides mirror sglang's TokenizerCommunicatorMixin defaults. | |
| with_stack = request.with_stack if request.HasField("with_stack") else None | |
| with_stack = (with_stack is not False) and get_bool_env_var( | |
| "SGLANG_PROFILE_WITH_STACK", "true" | |
| ) | |
| if request.HasField("with_stack"): | |
| with_stack = request.with_stack | |
| else: | |
| with_stack = get_bool_env_var("SGLANG_PROFILE_WITH_STACK", "true") |
If the AND-masking is intentional (operator can force-disable regardless of client request), a brief comment would help future readers understand the contract.
| if not results: | ||
| return False, "No response from scheduler" | ||
| failures = [r for r in results if not r.success] | ||
| if failures: | ||
| return False, " | ".join(r.message or "failed" for r in failures) |
There was a problem hiding this comment.
🟡 Nit: Missing blank line — PEP 8 (E302) expects two blank lines between top-level definitions. _aggregate_communicator_results ends on the line just before SAMPLING_DEFAULT_KEYS.
| if not results: | |
| return False, "No response from scheduler" | |
| failures = [r for r in results if not r.success] | |
| if failures: | |
| return False, " | ".join(r.message or "failed" for r in failures) | |
| SAMPLING_DEFAULT_KEYS = ( |
There was a problem hiding this comment.
Well-structured PR — the worker-trait abstraction for admin ops is a clean pattern that follows the established dual-path dispatch. Two minor nits posted inline (env-var masking semantics on profile options, PEP 8 spacing). No blocking issues found.
Summary: 0 🔴 Important · 2 🟡 Nit · 0 🟣 Pre-existing
There was a problem hiding this comment.
Actionable comments posted: 5
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (1)
grpc_servicer/smg_grpc_servicer/sglang/request_manager.py (1)
59-108:⚠️ Potential issue | 🟠 MajorEnsure
_GrpcCommunicatoruses a single asyncio event loop for both waiters andhandle_recvIn
grpc_servicer/smg_grpc_servicer/sglang/request_manager.py(59-108 and 1001-1030),_GrpcCommunicator.__call__creates theasyncio.Lock/asyncio.Eventon the caller’s running loop, buthandle_recv()completes that event fromhandle_loop()running on the loop chosen byauto_create_handle_loop()(get_or_create_event_loop()). Ifsend_communicator_req()is ever called from a different event loop/thread than the one that startedhandle_loop(),Event.set()will interact with Futures from another loop, riskingFuture attached to a different loopand/or hangs.Start
handle_loop()on the caller’sasyncio.get_running_loop()(insidesend_communicator_req), or route/restrictsend_communicator_req()calls so waiter and completer always run on the same event loop (or explicitly reject cross-loop use).🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py` around lines 59 - 108, The problem is cross-event-loop use: _GrpcCommunicator.__call__ (which creates asyncio.Event/Lock) can run on a different loop than handle_loop()/handle_recv (started via auto_create_handle_loop()/get_or_create_event_loop), causing "Future attached to a different loop" errors; fix by ensuring the handle loop is started on the caller's running loop or disallow cross-loop calls. Concretely, in send_communicator_req (the caller that spawns handle_loop/get_or_create_event_loop) obtain asyncio.get_running_loop() and start handle_loop on that loop (or pass that loop explicitly into get_or_create_event_loop/_GrpcCommunicator) so that __call__ (which creates self._lock/self._result_event) and handle_recv both use the same loop; alternatively add an explicit runtime check in _GrpcCommunicator.__call__ to raise if called from a different loop than the communicator's handle_loop, referencing methods/classes: _GrpcCommunicator, __call__, send_communicator_req, handle_loop, get_or_create_event_loop, handle_recv, and the attributes _result_event/_lock.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/protocols/src/worker.rs`:
- Around line 1008-1024: The ProfileOptions struct's profile_by_stage is a plain
bool so missing values are serialized as false and override backend defaults;
change the profile_by_stage field to Option<bool> (i.e., pub profile_by_stage:
Option<bool>) and add serde(skip_serializing_if = "Option::is_none") so unset
values remain None and are omitted during serialization (this preserves backend
defaults when StartProfileRequest is flattened).
In `@e2e_test/router/test_admin_ops.py`:
- Around line 90-99: Add a positive test alongside
test_profile_url_filter_without_match_returns_404 that POSTs to
f"{gateway.base_url}/start_profile" with a real worker URL (one from
setup_backend/gateway that will resolve) and assert resp.status_code == 200 and
resp.json()["workers_profiled"] == 1; run this assertion for both dispatch modes
(e.g., parametrize the test over the two modes or loop modes) to ensure
targeted-dispatch behavior is locked in, referencing the existing
test_profile_url_filter_without_match_returns_404 and the setup_backend fixture
to locate the gateway and worker URL.
In `@grpc_servicer/smg_grpc_servicer/sglang/servicer.py`:
- Around line 627-633: Before constructing FlushCacheReqInput and calling
request_manager.send_communicator_req in the flush-cache request handler (the
block that calls logger.debug and uses request.timeout_s, comm_timeout, and
FlushCacheReqInput), validate timeout_s with math.isfinite(timeout_s) and
timeout_s >= 0; if the check fails, abort the gRPC call with an INVALID_ARGUMENT
status (use the gRPC context abort / raise with grpc.StatusCode.INVALID_ARGUMENT
and a clear message) so negative or non-finite values are rejected at the
boundary instead of being forwarded to the scheduler.
In `@model_gateway/src/server.rs`:
- Around line 583-605: The handlers start_profile and stop_profile must stop
using Option<Json<_>> because an empty body with Content-Type: application/json
will cause extractor rejection; instead accept the raw request body (e.g.,
axum::body::Bytes or TypedBody) in start_profile and stop_profile, check if
bytes.is_empty() and if so use
StartProfileRequest::default()/StopProfileRequest::default(), otherwise attempt
serde_json::from_slice to deserialize into the respective request struct and
return a 400 on parse errors; then call WorkerManager::start_profile_all and
WorkerManager::stop_profile_all with the resolved request.options and
url.as_deref() as before.
In `@model_gateway/src/worker/worker.rs`:
- Around line 516-527: The code currently serializes the full ProfileOptions
with serde_json::to_value in the start_profile HTTP path, which forwards
gateway-only fields to the backend; instead, build and serialize the
backend-facing subset (the same subset used by the gRPC path) or explicitly
remove router-only fields before calling serde_json::to_value. Update the
start_profile HTTP branch (where ProfileOptions is converted and admin_http_post
is called) to map/convert ProfileOptions into the backend request struct (or
filter out fields like URL targeting) and serialize that sanitized struct for
the admin_http_post body so HTTP and gRPC behavior remain consistent.
---
Outside diff comments:
In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py`:
- Around line 59-108: The problem is cross-event-loop use:
_GrpcCommunicator.__call__ (which creates asyncio.Event/Lock) can run on a
different loop than handle_loop()/handle_recv (started via
auto_create_handle_loop()/get_or_create_event_loop), causing "Future attached to
a different loop" errors; fix by ensuring the handle loop is started on the
caller's running loop or disallow cross-loop calls. Concretely, in
send_communicator_req (the caller that spawns
handle_loop/get_or_create_event_loop) obtain asyncio.get_running_loop() and
start handle_loop on that loop (or pass that loop explicitly into
get_or_create_event_loop/_GrpcCommunicator) so that __call__ (which creates
self._lock/self._result_event) and handle_recv both use the same loop;
alternatively add an explicit runtime check in _GrpcCommunicator.__call__ to
raise if called from a different loop than the communicator's handle_loop,
referencing methods/classes: _GrpcCommunicator, __call__, send_communicator_req,
handle_loop, get_or_create_event_loop, handle_recv, and the attributes
_result_event/_lock.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 6710b8df-cd3d-456b-96d9-4dc19c91b8a8
📒 Files selected for processing (14)
crates/grpc_client/proto/common.protocrates/grpc_client/proto/sglang_scheduler.protocrates/grpc_client/src/lib.rscrates/grpc_client/src/sglang_scheduler.rscrates/protocols/src/worker.rse2e_test/router/test_admin_ops.pygrpc_servicer/smg_grpc_servicer/sglang/request_manager.pygrpc_servicer/smg_grpc_servicer/sglang/server.pygrpc_servicer/smg_grpc_servicer/sglang/servicer.pymodel_gateway/src/routers/grpc/client.rsmodel_gateway/src/server.rsmodel_gateway/src/worker/error.rsmodel_gateway/src/worker/manager.rsmodel_gateway/src/worker/worker.rs
| #[derive(Debug, Clone, Default, Serialize, Deserialize)] | ||
| #[serde(default)] | ||
| pub struct ProfileOptions { | ||
| #[serde(skip_serializing_if = "Option::is_none")] | ||
| pub output_dir: Option<String>, | ||
| #[serde(skip_serializing_if = "Option::is_none")] | ||
| pub start_step: Option<i32>, | ||
| #[serde(skip_serializing_if = "Option::is_none")] | ||
| pub num_steps: Option<i32>, | ||
| #[serde(skip_serializing_if = "Option::is_none")] | ||
| pub activities: Option<Vec<String>>, | ||
| #[serde(skip_serializing_if = "Option::is_none")] | ||
| pub with_stack: Option<bool>, | ||
| #[serde(skip_serializing_if = "Option::is_none")] | ||
| pub record_shapes: Option<bool>, | ||
| pub profile_by_stage: bool, | ||
| } |
There was a problem hiding this comment.
profile_by_stage cannot preserve omission.
ProfileOptions is documented as letting unset fields fall back to backend defaults, but profile_by_stage is a bare bool. With #[serde(default)] plus the flattened StartProfileRequest, payloads like {} or {"url":"..."} are always re-serialized downstream as "profile_by_stage": false, so the gateway overrides the backend default instead of preserving “unset”.
Suggested fix
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
#[serde(default)]
pub struct ProfileOptions {
#[serde(skip_serializing_if = "Option::is_none")]
pub output_dir: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub start_step: Option<i32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub num_steps: Option<i32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub activities: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub with_stack: Option<bool>,
#[serde(skip_serializing_if = "Option::is_none")]
pub record_shapes: Option<bool>,
- pub profile_by_stage: bool,
+ #[serde(skip_serializing_if = "Option::is_none")]
+ pub profile_by_stage: Option<bool>,
}Based on PR context, unset profile parameters are supposed to defer to backend defaults.
📝 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.
| #[derive(Debug, Clone, Default, Serialize, Deserialize)] | |
| #[serde(default)] | |
| pub struct ProfileOptions { | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub output_dir: Option<String>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub start_step: Option<i32>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub num_steps: Option<i32>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub activities: Option<Vec<String>>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub with_stack: Option<bool>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub record_shapes: Option<bool>, | |
| pub profile_by_stage: bool, | |
| } | |
| #[derive(Debug, Clone, Default, Serialize, Deserialize)] | |
| #[serde(default)] | |
| pub struct ProfileOptions { | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub output_dir: Option<String>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub start_step: Option<i32>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub num_steps: Option<i32>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub activities: Option<Vec<String>>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub with_stack: Option<bool>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub record_shapes: Option<bool>, | |
| #[serde(skip_serializing_if = "Option::is_none")] | |
| pub profile_by_stage: Option<bool>, | |
| } |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@crates/protocols/src/worker.rs` around lines 1008 - 1024, The ProfileOptions
struct's profile_by_stage is a plain bool so missing values are serialized as
false and override backend defaults; change the profile_by_stage field to
Option<bool> (i.e., pub profile_by_stage: Option<bool>) and add
serde(skip_serializing_if = "Option::is_none") so unset values remain None and
are omitted during serialization (this preserves backend defaults when
StartProfileRequest is flattened).
| def test_profile_url_filter_without_match_returns_404(self, setup_backend): | ||
| _backend, _model, _client, gateway = setup_backend | ||
|
|
||
| resp = httpx.post( | ||
| f"{gateway.base_url}/start_profile", | ||
| json={"url": "http://nonexistent:9999"}, | ||
| timeout=30.0, | ||
| ) | ||
| assert resp.status_code == 404, resp.text | ||
|
|
There was a problem hiding this comment.
🛠️ Refactor suggestion | 🟠 Major | ⚡ Quick win
Add a positive URL-filter match case for profiling.
The suite validates only the 404 no-match path. Please add one test where url matches a real worker and assert 200 plus workers_profiled == 1 for both modes to lock in targeted-dispatch behavior.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@e2e_test/router/test_admin_ops.py` around lines 90 - 99, Add a positive test
alongside test_profile_url_filter_without_match_returns_404 that POSTs to
f"{gateway.base_url}/start_profile" with a real worker URL (one from
setup_backend/gateway that will resolve) and assert resp.status_code == 200 and
resp.json()["workers_profiled"] == 1; run this assertion for both dispatch modes
(e.g., parametrize the test over the two modes or loop modes) to ensure
targeted-dispatch behavior is locked in, referencing the existing
test_profile_url_filter_without_match_returns_404 and the setup_backend fixture
to locate the gateway and worker URL.
| logger.debug("Receive flush cache request (timeout_s=%.1f)", request.timeout_s) | ||
| timeout_s = request.timeout_s | ||
| comm_timeout = max(30.0, timeout_s + 10.0) | ||
| try: | ||
| results = await self.request_manager.send_communicator_req( | ||
| FlushCacheReqInput(timeout_s=timeout_s), | ||
| "flush_cache_communicator", |
There was a problem hiding this comment.
Reject invalid timeout_s before sending FlushCacheReqInput.
The proto contract only defines 0 and positive wait times, but this forwards request.timeout_s verbatim. Negative or non-finite values currently fall through to the scheduler and can surface later as backend-specific failures instead of a clean INVALID_ARGUMENT at the gRPC boundary.
A small guard here for math.isfinite(timeout_s) and timeout_s >= 0 would keep the contract deterministic.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@grpc_servicer/smg_grpc_servicer/sglang/servicer.py` around lines 627 - 633,
Before constructing FlushCacheReqInput and calling
request_manager.send_communicator_req in the flush-cache request handler (the
block that calls logger.debug and uses request.timeout_s, comm_timeout, and
FlushCacheReqInput), validate timeout_s with math.isfinite(timeout_s) and
timeout_s >= 0; if the check fails, abort the gRPC call with an INVALID_ARGUMENT
status (use the gRPC context abort / raise with grpc.StatusCode.INVALID_ARGUMENT
and a clear message) so negative or non-finite values are rejected at the
boundary instead of being forwarded to the scheduler.
| async fn start_profile( | ||
| State(state): State<Arc<AppState>>, | ||
| body: Option<Json<StartProfileRequest>>, | ||
| ) -> Response { | ||
| let body = body.map_or_else(StartProfileRequest::default, |Json(body)| body); | ||
| WorkerManager::start_profile_all( | ||
| &state.context.worker_registry, | ||
| &body.options, | ||
| body.url.as_deref(), | ||
| ) | ||
| .await | ||
| .into_response() | ||
| } | ||
|
|
||
| async fn stop_profile( | ||
| State(state): State<Arc<AppState>>, | ||
| body: Option<Json<StopProfileRequest>>, | ||
| ) -> Response { | ||
| let body = body.map_or_else(StopProfileRequest::default, |Json(body)| body); | ||
| WorkerManager::stop_profile_all(&state.context.worker_registry, body.url.as_deref()) | ||
| .await | ||
| .into_response() | ||
| } |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
set -euo pipefail
# Show axum version
if [ -f Cargo.toml ]; then
rg -n "axum\s*=" Cargo.toml || true
fi
rg -n "axum::" model_gateway/Cargo.toml model_gateway/src -S || true
# Inspect the handlers and surrounding route setup
sed -n '520,700p' model_gateway/src/server.rs
# If route registration is elsewhere near reported lines
sed -n '880,980p' model_gateway/src/server.rsRepository: lightseekorg/smg
Length of output: 23633
🌐 Web query:
axum Option<Json<T>> OptionalFromRequest empty body Content-Type application/json
💡 Result:
In axum, when using Option<Json<T>> as an extractor, the behavior is governed by the OptionalFromRequest trait [1][2]. If a request has no Content-Type header and an empty body, the Json extractor will return None [1][2][3]. However, if the request includes a Content-Type: application/json header but has an empty or malformed body, the Json extractor will typically return an error (rejection) rather than None [1][2]. This is because the implementation of OptionalFromRequest for Json is designed to return None only when the request does not specify the JSON content type [3]. If the content type header is present, axum attempts to parse the body as JSON; if that parsing fails (due to an empty or invalid body), the extraction fails with a JsonRejection [3][4]. If you require different behavior, such as treating an empty body as None even when the Content-Type header is present, you may need to implement a custom extractor or manually handle the Result<Json<T>, JsonRejection> and convert errors as needed [1][5]. Top results: [1], [2], [3], [4]
Citations:
- 1: Optional Json doesn't work in axum 0.8.1 tokio-rs/axum#3152
- 2: Implement
OptionalFromRequestfor the Json extractor tokio-rs/axum#3142 - 3: https://tg-rs.github.io/carapax/src/axum/json.rs.html
- 4: https://tg-rs.github.io/carapax/axum/struct.Json.html
- 5: https://readmedium.com/how-to-make-an-api-endpoint-accept-post-requests-with-an-optional-json-body-in-axum-23a22c23d582
Handle empty request bodies explicitly for /start_profile and /stop_profile (avoid Option<Json<_>> reliance).
In model_gateway/src/server.rs (lines 583-605), both handlers default via Option<Json<_>>, but axum’s Option<Json<T>> yields None only when the request does not include Content-Type: application/json. If the client sends Content-Type: application/json with a zero-length/empty body, JSON extraction rejects (so the default isn’t used). Deserialize conditionally (e.g., accept raw bytes/optional body and treat empty as default, or use a small custom extractor) instead of relying on Option<Json<_>>.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@model_gateway/src/server.rs` around lines 583 - 605, The handlers
start_profile and stop_profile must stop using Option<Json<_>> because an empty
body with Content-Type: application/json will cause extractor rejection; instead
accept the raw request body (e.g., axum::body::Bytes or TypedBody) in
start_profile and stop_profile, check if bytes.is_empty() and if so use
StartProfileRequest::default()/StopProfileRequest::default(), otherwise attempt
serde_json::from_slice to deserialize into the respective request struct and
return a 400 on parse errors; then call WorkerManager::start_profile_all and
WorkerManager::stop_profile_all with the resolved request.options and
url.as_deref() as before.
| let body = | ||
| serde_json::to_value(options).map_err(|e| WorkerError::OperationFailed { | ||
| url: self.url().to_string(), | ||
| operation: "start_profile".to_string(), | ||
| reason: format!("failed to serialize profile options: {e}"), | ||
| })?; | ||
| admin_http_post( | ||
| self.http_client(), | ||
| self.endpoint_url("/start_profile"), | ||
| self.api_key(), | ||
| Some(body), | ||
| "start_profile", |
There was a problem hiding this comment.
Do not forward gateway-only profile filter fields to backend HTTP /start_profile.
serde_json::to_value(options) serializes the full ProfileOptions, while the gRPC path sends a backend-facing subset. This makes HTTP/gRPC behavior diverge and can break HTTP workers when router-only fields (for example URL targeting) are present.
Suggested fix
- let body =
- serde_json::to_value(options).map_err(|e| WorkerError::OperationFailed {
- url: self.url().to_string(),
- operation: "start_profile".to_string(),
- reason: format!("failed to serialize profile options: {e}"),
- })?;
+ // Keep HTTP payload aligned with backend-facing gRPC request fields.
+ let body = serde_json::json!({
+ "output_dir": options.output_dir,
+ "start_step": options.start_step,
+ "num_steps": options.num_steps,
+ "activities": options.activities,
+ "with_stack": options.with_stack,
+ "record_shapes": options.record_shapes,
+ "profile_by_stage": options.profile_by_stage,
+ });📝 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.
| let body = | |
| serde_json::to_value(options).map_err(|e| WorkerError::OperationFailed { | |
| url: self.url().to_string(), | |
| operation: "start_profile".to_string(), | |
| reason: format!("failed to serialize profile options: {e}"), | |
| })?; | |
| admin_http_post( | |
| self.http_client(), | |
| self.endpoint_url("/start_profile"), | |
| self.api_key(), | |
| Some(body), | |
| "start_profile", | |
| // Keep HTTP payload aligned with backend-facing gRPC request fields. | |
| let body = serde_json::json!({ | |
| "output_dir": options.output_dir, | |
| "start_step": options.start_step, | |
| "num_steps": options.num_steps, | |
| "activities": options.activities, | |
| "with_stack": options.with_stack, | |
| "record_shapes": options.record_shapes, | |
| "profile_by_stage": options.profile_by_stage, | |
| }); | |
| admin_http_post( | |
| self.http_client(), | |
| self.endpoint_url("/start_profile"), | |
| self.api_key(), | |
| Some(body), | |
| "start_profile", |
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@model_gateway/src/worker/worker.rs` around lines 516 - 527, The code
currently serializes the full ProfileOptions with serde_json::to_value in the
start_profile HTTP path, which forwards gateway-only fields to the backend;
instead, build and serialize the backend-facing subset (the same subset used by
the gRPC path) or explicitly remove router-only fields before calling
serde_json::to_value. Update the start_profile HTTP branch (where ProfileOptions
is converted and admin_http_post is called) to map/convert ProfileOptions into
the backend request struct (or filter out fields like URL targeting) and
serialize that sanitized struct for the admin_http_post body so HTTP and gRPC
behavior remain consistent.
…min ops Add KV-cache flush and profiler start/stop support for gRPC-mode engines, dispatched per worker behind the Worker trait: - proto: FlushCache/StartProfile/StopProfile messages in common.proto (shared across engines, same pattern as GetTokenizer); RPCs on SglangScheduler - grpc_client: impl_admin_ops! macro shared by engine clients with local deadlines and trace injection; wired for sglang - Worker trait: flush_cache/start_profile/stop_profile default impls dispatch on connection mode (HTTP endpoint vs gRPC RPC), mirroring check_health_async - WorkerManager: flush_cache_all becomes one uniform fan-out, so gRPC workers are no longer silently skipped; start_profile_all and stop_profile_all support optional worker-URL targeting for PD mode - server: /start_profile and /stop_profile admin routes - grpc_servicer (sglang): FlushCache/StartProfile/StopProfile handlers, communicator plumbing, and the send_communicator_req + serve_grpc(on_request_manager_ready) surface that sglang's gRPC entrypoint requires since sgl-project/sglang#22500. The communicator response socket only binds tokenizer_ipc_name when recv_from_scheduler has not already bound it (--skip-tokenizer-init mode, #1591) Supersedes #1088 Signed-off-by: Simo Lin <25425177+slin1237@users.noreply.github.com>
60c40d5 to
0f467fd
Compare
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 0f467fd61f
ℹ️ 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".
| let workers = worker_registry.get_all(); | ||
| let total_workers = workers.len(); | ||
| match url_filter { | ||
| Some(url) => workers.into_iter().filter(|w| w.url() == url).collect(), |
There was a problem hiding this comment.
Accept DP base URLs when targeting profiles
In DP-aware deployments the registered worker URL is internally rewritten to {base}@{rank} while endpoint_url() sends admin calls to the stripped base_url, so a caller that targets the real backend URL (for example grpc://prefill:port from a profiling workflow) will hit this exact w.url() == url filter and get a 404 even though the worker exists. Match base_url() as well (and deduplicate ranks as needed) so single-worker profiling works for DP-aware workers using the externally known URL.
Useful? React with 👍 / 👎.
| f"{gateway.base_url}/start_profile", | ||
| content=b"{not json", | ||
| headers={"content-type": "application/json"}, | ||
| timeout=30.0, |
There was a problem hiding this comment.
🟡 Nit: This assertion expects 400 but will likely get 422. The start_profile handler uses Option<Json<StartProfileRequest>> — when Content-Type: application/json is present, axum's OptionalFromRequest impl tries to parse the body. Malformed JSON triggers JsonRejection::JsonSyntaxError, which defaults to 422 Unprocessable Entity, not 400. There's no custom rejection handler on the admin routes to remap it.
Either change the assertion to 422, or (if 400 is the desired contract) switch the handler from Option<Json<_>> to a custom extractor / raw Bytes with explicit parse + error mapping — which would also fix the related Option<Json<_>> empty-body issue that was flagged elsewhere.
There was a problem hiding this comment.
Actionable comments posted: 1
♻️ Duplicate comments (2)
grpc_servicer/smg_grpc_servicer/sglang/servicer.py (1)
629-635:⚠️ Potential issue | 🟡 Minor | ⚡ Quick winValidate
timeout_sat the RPC boundary before forwarding.On Line 630 and Line 634,
timeout_sis forwarded without checking finiteness/non-negativity. Invalid values (negative/NaN/inf) should be rejected asINVALID_ARGUMENTinstead of reaching scheduler-specific failure paths.Suggested patch
+import math @@ logger.debug("Receive flush cache request (timeout_s=%.1f)", request.timeout_s) timeout_s = request.timeout_s + if not math.isfinite(timeout_s) or timeout_s < 0: + await context.abort( + grpc.StatusCode.INVALID_ARGUMENT, + "timeout_s must be a finite value >= 0", + ) comm_timeout = max(30.0, timeout_s + 10.0)🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@grpc_servicer/smg_grpc_servicer/sglang/servicer.py` around lines 629 - 635, Validate request.timeout_s at the RPC boundary before using it: check request.timeout_s (assigned to timeout_s) for finiteness (not NaN or +/-inf) and non-negativity, and if invalid return an INVALID_ARGUMENT gRPC error immediately instead of forwarding; then proceed to compute comm_timeout and call request_manager.send_communicator_req with FlushCacheReqInput(timeout_s=timeout_s) and "flush_cache_communicator" only when the value is valid.crates/protocols/src/worker.rs (1)
1007-1024:⚠️ Potential issue | 🟠 Major | 🏗️ Heavy lift
profile_by_stagecurrently cannot preserve “unset” behavior.Line [1007] documents fallback-to-backend-default semantics for unset options, but
profile_by_stageis a non-optionalboolon Line [1023], so omission is coerced tofalseand forwarded explicitly. That makes the contract non-preserving for this field (especially through flattened request serialization).Please either make this field optional end-to-end (protocol + worker mapping + gRPC contract) or narrow the contract/docs to state this field is always explicit.
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/protocols/src/worker.rs` around lines 1007 - 1024, The ProfileOptions struct’s profile_by_stage is a plain bool so omission becomes false and breaks the “unset -> backend default” contract; change profile_by_stage to an Option<bool> (e.g., pub profile_by_stage: Option<bool>) and add serde(skip_serializing_if = "Option::is_none") so it can be omitted, then propagate this optional type through the worker/gRPC mapping code that reads ProfileOptions (update any conversion functions, request builders, and protobuf serialization/deserialization that reference ProfileOptions.profile_by_stage) to treat None as “use backend default” while keeping existing true/false semantics when Some(value).
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@grpc_servicer/smg_grpc_servicer/sglang/servicer.py`:
- Around line 647-651: The exception handling for the FlushCache path currently
logs the error then calls context.set_code/set_details and returns a
FlushCacheResponse; change this to follow servicer convention: after logging
(use logger.error(...) and get_exception_traceback()), await
context.abort(grpc.StatusCode.INTERNAL, str(e)) instead of
context.set_code/set_details and returning common_pb2.FlushCacheResponse,
removing the manual return; apply the same change to the other INTERNAL error
block in this file (the other try/except that currently uses
context.set_code/set_details and returns a payload) so both handlers uniformly
use await context.abort(grpc.StatusCode.INTERNAL, str(e)).
---
Duplicate comments:
In `@crates/protocols/src/worker.rs`:
- Around line 1007-1024: The ProfileOptions struct’s profile_by_stage is a plain
bool so omission becomes false and breaks the “unset -> backend default”
contract; change profile_by_stage to an Option<bool> (e.g., pub
profile_by_stage: Option<bool>) and add serde(skip_serializing_if =
"Option::is_none") so it can be omitted, then propagate this optional type
through the worker/gRPC mapping code that reads ProfileOptions (update any
conversion functions, request builders, and protobuf
serialization/deserialization that reference ProfileOptions.profile_by_stage) to
treat None as “use backend default” while keeping existing true/false semantics
when Some(value).
In `@grpc_servicer/smg_grpc_servicer/sglang/servicer.py`:
- Around line 629-635: Validate request.timeout_s at the RPC boundary before
using it: check request.timeout_s (assigned to timeout_s) for finiteness (not
NaN or +/-inf) and non-negativity, and if invalid return an INVALID_ARGUMENT
gRPC error immediately instead of forwarding; then proceed to compute
comm_timeout and call request_manager.send_communicator_req with
FlushCacheReqInput(timeout_s=timeout_s) and "flush_cache_communicator" only when
the value is valid.
🪄 Autofix (Beta)
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 0a04daf2-0cc4-4668-abce-5486c1678d10
📒 Files selected for processing (14)
crates/grpc_client/proto/common.protocrates/grpc_client/proto/sglang_scheduler.protocrates/grpc_client/src/lib.rscrates/grpc_client/src/sglang_scheduler.rscrates/protocols/src/worker.rse2e_test/router/test_admin_ops.pygrpc_servicer/smg_grpc_servicer/sglang/request_manager.pygrpc_servicer/smg_grpc_servicer/sglang/server.pygrpc_servicer/smg_grpc_servicer/sglang/servicer.pymodel_gateway/src/routers/grpc/client.rsmodel_gateway/src/server.rsmodel_gateway/src/worker/error.rsmodel_gateway/src/worker/manager.rsmodel_gateway/src/worker/worker.rs
| except Exception as e: | ||
| logger.error(f"FlushCache failed: {e}\n{get_exception_traceback()}") | ||
| context.set_code(grpc.StatusCode.INTERNAL) | ||
| context.set_details(f"Flush cache failed: {e}") | ||
| return common_pb2.FlushCacheResponse(success=False, message=f"Flush cache failed: {e}") |
There was a problem hiding this comment.
Use context.abort(..., str(e)) for INTERNAL errors in servicer handlers.
Line 649-651 and Line 719-721 currently use set_code/set_details and return payloads for INTERNAL paths. Please align with the servicer convention and abort directly with str(e) so error propagation is consistent across implementations.
Based on learnings: repository-wide servicer convention is to log then await context.abort(grpc.StatusCode.INTERNAL, str(e)) without substituting messages.
Suggested patch
except Exception as e:
logger.error(f"FlushCache failed: {e}\n{get_exception_traceback()}")
- context.set_code(grpc.StatusCode.INTERNAL)
- context.set_details(f"Flush cache failed: {e}")
- return common_pb2.FlushCacheResponse(success=False, message=f"Flush cache failed: {e}")
+ await context.abort(grpc.StatusCode.INTERNAL, str(e))
@@
except Exception as e:
logger.error(f"{op_name} failed: {e}\n{get_exception_traceback()}")
- context.set_code(grpc.StatusCode.INTERNAL)
- context.set_details(f"{op_name} failed: {e}")
- return common_pb2.ProfileResponse(success=False, message=f"{op_name} failed: {e}")
+ await context.abort(grpc.StatusCode.INTERNAL, str(e))Also applies to: 718-721
🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
In `@grpc_servicer/smg_grpc_servicer/sglang/servicer.py` around lines 647 - 651,
The exception handling for the FlushCache path currently logs the error then
calls context.set_code/set_details and returns a FlushCacheResponse; change this
to follow servicer convention: after logging (use logger.error(...) and
get_exception_traceback()), await context.abort(grpc.StatusCode.INTERNAL,
str(e)) instead of context.set_code/set_details and returning
common_pb2.FlushCacheResponse, removing the manual return; apply the same change
to the other INTERNAL error block in this file (the other try/except that
currently uses context.set_code/set_details and returns a payload) so both
handlers uniformly use await context.abort(grpc.StatusCode.INTERNAL, str(e)).
Source: Learnings
Summary
Adds KV-cache flush and profiling (start/stop) support for engines behind SMG, dispatched per worker behind the
Workertrait. Supersedes #1088 — thanks @Kangyan-Zhou for the groundwork; the Python transport layer from that PR is kept (its API is now a public contract relied on by merged sgl-project/sglang#22500), while the router side is restructured so admin ops are abstracted behind the worker like every other engine-facing operation.What changed
common.protoFlushCacheRequest/Response,StartProfileRequest,StopProfileRequest,ProfileResponse— shared across engines (same pattern asGetTokenizer)sglang_scheduler.protoFlushCache,StartProfile,StopProfileRPCscrates/grpc_clientimpl_admin_ops!macro shared by engine clients (local deadline + trace injection on every RPC), wired for sglangGrpcClientdispatcherUnimplemented(theget_loadsgating pattern)Workertraitflush_cache/start_profile/stop_profiledefault impls dispatch on connection mode — HTTP workers get the engine's native endpoint, gRPC workers get the RPC (mirrorscheck_health_async)WorkerManagerflush_cache_allis one uniform fan-out — gRPC workers are no longer silently skipped; newstart_profile_all/stop_profile_allwith optional worker-URL targetingserver.rs/start_profile,/stop_profileadmin routes;/flush_cacheresponse now also reportstotal_grpc_workersgrpc_servicer/sglangFlushCache/StartProfile/StopProfilehandlers; communicator plumbing;send_communicator_req+serve_grpc(on_request_manager_ready=...)— the surface sglang's gRPC entrypoint requires since sgl-project/sglang#22500Why restructure relative to #1088
WorkerManager::flush_cache_all; every future admin op would duplicate that split. TheWorkertrait already owns dual-path dispatch (check_health_async), so admin ops follow the same shape — the tokenspeed follow-up needs only proto + servicer + client wiring, zero router changes.recv_from_scheduleralready bindstokenizer_ipc_nameunder--skip-tokenizer-init; feat(grpc): add FlushCache RPC and profiling support for gRPC mode #1088's unconditional second bind of the same endpoint would fail at startup. The communicator-response socket now binds only when the endpoint is free, and both receive loops share one dispatch helper.{"url": ...}for PD mode (bench_serving --profile-prefill-url-style workflows).GrpcClientmatch including theMlx/TokenSpeedvariants added since feat(grpc): add FlushCache RPC and profiling support for gRPC mode #1088.Compatibility notes
serve_grpc(..., on_request_manager_ready=...)andrequest_manager.send_communicator_req(...)(merged in [Observability] Add HTTP sidecar endpoints and FlushCache gRPC RPC for gRPC mode sgl-project/sglang#22500). sglang gRPC mode is broken againstsmg-grpc-servicermain until this lands.FlushCacheReqInput(timeout_s=...)exists at sglang HEAD (Optional[float]).smg-grpc-proto/smg-grpc-servicerversion bumps are intentionally left to a separate PR.Follow-up (separate PR)
TokenSpeed: add the 3 RPCs to
tokenspeed_scheduler.proto(reusing the common messages),impl_admin_ops!()in its Rust client, servicer handlers calling the existingAsyncLLM.flush_cache/start_profile/stop_profile, flip the dispatcher arms, and extend the e2e marker toengine("sglang", "tokenspeed"). No changes needed in the tokenspeed repo.Test plan
cargo +nightly fmt --allcleancargo clippy --all-targets --all-features -- -D warningscleancargo test --workspace— 3502 passed, 0 failed (includes new fan-out unit tests against loopback admin stubs)e2e_test/router/test_admin_ops.py: flush + profile round-trips through the gateway for bothgrpcandhttpbackends, URL-filter 404, malformed-JSON 400Summary by CodeRabbit
New Features
UX/Behavior
Tests