fix(grpc): drain in-flight requests on SIGTERM for graceful scale-down - #1172
Kangyan-Zhou wants to merge 7 commits into
Conversation
When a Kubernetes pod running the SGLang gRPC servicer is scaled down,
the old signal handler only set `stop_event` and immediately invoked
`request_manager.shutdown()`, which cancels tasks and pushes
`{"error": "Server shutting down"}` into every in-flight request's
out_queue. In-flight prefill/decode streams were aborted mid-generation.
Mirror the HTTP `TokenizerManager.sigterm_watchdog` pattern:
- Add `GrpcRequestManager.begin_drain()` (non-destructive; flips
`gracefully_exit` so the health servicer reports NOT_SERVING and K8s
stops routing new traffic).
- Add `SGLangSchedulerServicer.begin_drain()` as a thin forwarder.
- In `serve_grpc`, call `servicer.begin_drain()` from the signal handler
and insert a drain loop that polls `rid_to_state` every 5 s until
empty, excluding `state.finished` entries (which linger 5 s via the
existing `cleanup()` task). Honour `SGL_FORCE_SHUTDOWN` as an escape
hatch, matching HTTP mode.
- Fix `GrpcRequestManager.handle_loop` to run `while True:` instead of
`while not gracefully_exit:` so scheduler outputs keep flowing to
in-flight streaming requests during drain; the loop now exits only
via task cancellation from `shutdown()` or an unrecoverable ZMQ
error. HTTP's `TokenizerManager.handle_loop` uses the same pattern.
No CLI or default changes: operators continue to tune drain time via
K8s `terminationGracePeriodSeconds`, with SIGKILL as the ultimate
backstop (same as HTTP mode).
Unit tests cover `begin_drain` (non-destructive with populated
`rid_to_state`/`asyncio_tasks`, idempotent) and the health-flag wiring.
Signed-off-by: Kangyan Zhou <zky314343421@gmail.com>
The root pytest.ini only picks up e2e_test/; each Python package is expected to carry its own [tool.pytest.ini_options] (see bindings/python/pyproject.toml, clients/python/pyproject.toml). Add the same wiring to grpc_servicer so `pytest grpc_servicer/` discovers the graceful-shutdown unit tests, and declare pytest + pytest-asyncio as dev dependencies. Signed-off-by: Kangyan Zhou <zky314343421@gmail.com>
The critical fix in 2d4d2186 (handle_loop must be `while True:`, not `while not self.gracefully_exit:`) had no test coverage. If someone reverts it, none of the existing tests would fail because they don't exercise handle_loop. Add a source-inspection test that asserts the invariant directly. It's not a behavioural test — but a reverted gate would produce a silent runtime regression (stalled streams on drain), and a grep-style guard is cheap, explicit about why the invariant matters, and caught at import time rather than at production shutdown. Signed-off-by: Kangyan Zhou <zky314343421@gmail.com>
The watchdog body was `while not self.gracefully_exit: await asyncio.sleep(1)` — it sleeps until the flag flips, then returns with no side effects. Its docstring claimed parity with `TokenizerManager.sigterm_watchdog`, but the HTTP version actually drains `rid_to_state` and kills the process tree. With the drain now implemented in `serve_grpc` (server.py:291), the stub is dead wiring and misleading to readers. Also clarify the intent of the SIGTERM/SIGQUIT registrations in auto_create_handle_loop: SIGTERM here is a startup-window fallback that serve_grpc overrides; SIGQUIT stays owned here for scheduler-crash forwarding. Signed-off-by: Kangyan Zhou <zky314343421@gmail.com>
📝 WalkthroughWalkthroughThis change implements a graceful drain mechanism for the gRPC servicer that transitions the system into a draining state upon shutdown, prevents new requests from being enqueued, and waits for in-flight requests to complete before final cleanup. The handle loop now runs unconditionally to forward scheduler outputs during drain mode. Changes
Sequence DiagramsequenceDiagram
participant Signal as Signal Handler
participant Servicer as SGLangServicer
participant ReqMgr as RequestManager
participant Server as Server (shutdown)
participant Client as Client (in-flight RPC)
Signal->>Servicer: SIGTERM/SIGINT
Servicer->>ReqMgr: begin_drain()
activate ReqMgr
ReqMgr->>ReqMgr: set gracefully_exit = True
deactivate ReqMgr
Servicer->>Servicer: set health to NOT_SERVING
Note over Servicer: New RPCs blocked
Client->>Servicer: Generate/Embed RPC
Servicer-->>Client: UNAVAILABLE (Server shutting down)
Note over ReqMgr: In-flight requests continue
ReqMgr->>ReqMgr: forward scheduler outputs
Client->>ReqMgr: await completion
Server->>Server: drain loop
loop Check remaining requests
Server->>ReqMgr: inspect rid_to_state
alt requests remain
Server->>Server: sleep 5s
else all finished
Server->>Server: proceed to shutdown
end
end
ReqMgr->>Server: handle_loop_task completes
Server->>ReqMgr: shutdown()
Estimated Code Review Effort🎯 3 (Moderate) | ⏱️ ~25 minutes Possibly Related PRs
Suggested Reviewers
Poem
🚥 Pre-merge checks | ✅ 4 | ❌ 1❌ Failed checks (1 warning)
✅ Passed checks (4 passed)
✏️ Tip: You can configure your own custom pre-merge checks in the settings. ✨ Finishing Touches🧪 Generate unit tests (beta)
Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out. Comment |
There was a problem hiding this comment.
Code Review
This pull request implements a graceful shutdown and drain mechanism for the gRPC servicer, ensuring that in-flight requests can complete before the server terminates. Key changes include modifying the request manager's event loop to remain active during the drain phase, updating signal handlers to trigger a drain state, and adding a monitoring loop that waits for active requests to finish. Review feedback highlights a potential hang if background tasks fail unexpectedly and suggests optimizations for logging verbosity and the responsiveness of the drain loop.
| while True: | ||
| if get_bool_env_var("SGL_FORCE_SHUTDOWN"): | ||
| logger.warning("SGL_FORCE_SHUTDOWN set; skipping drain") | ||
| break | ||
|
|
||
| # Finished requests linger in rid_to_state for 5 s via the | ||
| # cleanup() task in _handle_batch_output; exclude them so | ||
| # the drain loop exits promptly once real work is done. | ||
| remaining_rids = [ | ||
| rid | ||
| for rid, state in servicer.request_manager.rid_to_state.items() | ||
| if not state.finished | ||
| ] | ||
| remain_num_req = len(remaining_rids) | ||
| if remain_num_req == 0: | ||
| logger.info("Drain complete; no in-flight requests") | ||
| break | ||
|
|
||
| logger.info( | ||
| "Gracefully exiting... Remaining number of requests %d. Remaining requests %s", | ||
| remain_num_req, | ||
| remaining_rids, | ||
| ) | ||
| await asyncio.sleep(5) |
There was a problem hiding this comment.
There is a potential hang in this drain loop if the background handle_loop task in GrpcRequestManager terminates unexpectedly (for example, due to a ZMQError as seen in request_manager.py line 524). If handle_loop stops, scheduler outputs will no longer be processed, and in-flight requests will never be marked as finished. This would cause the drain loop to wait indefinitely until the pod is forcefully killed by the orchestrator (e.g., Kubernetes SIGKILL after terminationGracePeriodSeconds). Consider adding a mechanism to detect if the background processing task is still alive during the drain phase.
| logger.info( | ||
| "Gracefully exiting... Remaining number of requests %d. Remaining requests %s", | ||
| remain_num_req, | ||
| remaining_rids, | ||
| ) |
There was a problem hiding this comment.
Logging the full list of remaining_rids can be extremely verbose in high-throughput scenarios where many requests are in-flight during a scale-down. This could lead to log flooding and performance degradation. It is better to truncate the list in the log message.
| logger.info( | |
| "Gracefully exiting... Remaining number of requests %d. Remaining requests %s", | |
| remain_num_req, | |
| remaining_rids, | |
| ) | |
| logger.info( | |
| "Gracefully exiting... Remaining number of requests %d. Remaining requests %s", | |
| remain_num_req, | |
| remaining_rids[:10] + (["..."] if remain_num_req > 10 else []), | |
| ) |
| remain_num_req, | ||
| remaining_rids, | ||
| ) | ||
| await asyncio.sleep(5) |
There was a problem hiding this comment.
A 5-second sleep interval for the drain loop is relatively long. It can delay the final termination of the pod by several seconds even after all in-flight requests have finished. Reducing this to 1 second would make the shutdown process more responsive without significantly increasing CPU overhead.
| await asyncio.sleep(5) | |
| await asyncio.sleep(1) |
|
Warning You have reached your daily quota limit. Please wait up to 24 hours and I will start processing your requests again! |
There was a problem hiding this comment.
Actionable comments posted: 2
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/server.py (1)
277-319:⚠️ Potential issue | 🟠 MajorAdd admission guards to reject new
Generate/EmbedRPCs during graceful shutdown.Line 282 calls
begin_drain()which setsgracefully_exit, butGenerate(line 207) andEmbed(line 266) lack admission checks. New HTTP/2 connections can still start these RPCs during the drain window (beforeserver.stop()at line 327), queuing work inrid_to_stateand extending the drain loop beyond the intended in-flight-request window.The
HealthCheckRPC already has the pattern (line 319): add equivalentif self.request_manager.gracefully_exit:guards to bothGenerateandEmbedentry points to abort withgrpc.StatusCode.UNAVAILABLE.🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@grpc_servicer/smg_grpc_servicer/sglang/server.py` around lines 277 - 319, Generate and Embed need the same admission-check used by HealthCheck to reject new RPCs once shutdown begins: in the Generate and Embed methods, at the very start (before allocating request state or queuing work into request_manager.rid_to_state), check self.request_manager.gracefully_exit and if true abort immediately with grpc.StatusCode.UNAVAILABLE (mirroring the HealthCheck implementation) so begin_drain() actually prevents new work from being accepted during the drain window; update both Generate and Embed entry paths to perform this guard and return the same error response used by HealthCheck.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Inline comments:
In `@grpc_servicer/pyproject.toml`:
- Around line 35-38: The dev extra in pyproject.toml currently lists only pytest
packages but tests import smg_grpc_servicer.sglang which does a module-level
import sglang, so add sglang to the dev extra (e.g., include "sglang" with an
appropriate version spec) so pip install -e .[dev] pulls it in, or alternatively
update the repository README / CONTRIBUTING to state that developers must
install with .[dev,sglang]; modify the dev extras entry in pyproject.toml (the
"dev" list) or the docs accordingly to resolve ImportError during test
collection.
In `@grpc_servicer/tests/sglang/test_graceful_shutdown.py`:
- Around line 58-75: Replace the fragile string checks in
test_handle_loop_not_gated_on_gracefully_exit with an AST-based assertion: parse
inspect.getsource(GrpcRequestManager.handle_loop) into an AST, find the
top-level ast.While node(s) in the function body, and assert at least one while
loop has a condition that is a literal True (ast.Constant True or
ast.NameConstant True) and that none of the top-level while conditions represent
a negated or false condition that references self.gracefully_exit (e.g.,
ast.UnaryOp(Not)/ast.Compare against False or direct Name/Attribute comparing to
False); update the test to reference the GrpcRequestManager.handle_loop function
name when locating the AST node so comments/docstrings or incidental text cannot
bypass the check.
---
Outside diff comments:
In `@grpc_servicer/smg_grpc_servicer/sglang/server.py`:
- Around line 277-319: Generate and Embed need the same admission-check used by
HealthCheck to reject new RPCs once shutdown begins: in the Generate and Embed
methods, at the very start (before allocating request state or queuing work into
request_manager.rid_to_state), check self.request_manager.gracefully_exit and if
true abort immediately with grpc.StatusCode.UNAVAILABLE (mirroring the
HealthCheck implementation) so begin_drain() actually prevents new work from
being accepted during the drain window; update both Generate and Embed entry
paths to perform this guard and return the same error response used by
HealthCheck.
🪄 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: c6fec908-55e1-4fd8-b21b-9068cbc80036
📒 Files selected for processing (6)
grpc_servicer/pyproject.tomlgrpc_servicer/smg_grpc_servicer/sglang/request_manager.pygrpc_servicer/smg_grpc_servicer/sglang/server.pygrpc_servicer/smg_grpc_servicer/sglang/servicer.pygrpc_servicer/tests/conftest.pygrpc_servicer/tests/sglang/test_graceful_shutdown.py
| dev = [ | ||
| "pytest>=8", | ||
| "pytest-asyncio>=0.24", | ||
| ] |
There was a problem hiding this comment.
🧩 Analysis chain
🏁 Script executed:
#!/bin/bash
# Description: Show dev/sglang extras and SGLang imports used by the new tests.
sed -n '/^\[project.optional-dependencies\]/,/^\[project.urls\]/p' grpc_servicer/pyproject.toml
rg -n -C2 'from smg_grpc_servicer\.sglang|from sglang|import sglang' \
grpc_servicer/tests grpc_servicer/smg_grpc_servicer/sglangRepository: lightseekorg/smg
Length of output: 9178
🏁 Script executed:
#!/bin/bash
# Check pytest configuration in pyproject.toml and conftest.py
echo "=== Pytest configuration in pyproject.toml ==="
sed -n '/\[tool.pytest/,/^\[/p' grpc_servicer/pyproject.toml | head -50
echo -e "\n=== Check for conftest.py in tests directory ==="
find grpc_servicer/tests -name "conftest.py" -exec cat {} \;
echo -e "\n=== Check for skipif or pytest markers in sglang tests ==="
head -30 grpc_servicer/tests/sglang/test_graceful_shutdown.pyRepository: lightseekorg/smg
Length of output: 1880
🏁 Script executed:
#!/bin/bash
# Check CI configuration to see how tests are actually invoked
echo "=== Looking for CI/workflow config files ==="
find . -maxdepth 3 -type f \( -name "*.yml" -o -name "*.yaml" \) -path "*/.github/workflows/*" | head -10
echo -e "\n=== Check for tox.ini or similar test config ==="
find grpc_servicer -maxdepth 2 -type f \( -name "tox.ini" -o -name "pytest.ini" -o -name "setup.cfg" \)
echo -e "\n=== Check pyproject.toml for test/dev instructions or test-related sections ==="
grep -A5 -B2 "test\|dev\|pytest" grpc_servicer/pyproject.toml | head -60Repository: lightseekorg/smg
Length of output: 1177
🏁 Script executed:
#!/bin/bash
# Search for CI workflows that run grpc_servicer tests
echo "=== Search for grpc_servicer test runs in CI ==="
grep -r "grpc_servicer" .github/workflows/ | grep -E "test|pytest|install" | head -20
echo -e "\n=== Check main test/CI workflows for patterns ==="
for f in .github/workflows/test*.yml .github/workflows/*test*.yml .github/workflows/ci*.yml; do
if [ -f "$f" ]; then
echo "=== File: $f ==="
grep -A5 -B5 "grpc.servicer\|pytest\|\[dev\]" "$f" | head -40
fi
done
echo -e "\n=== Check for any Makefile or shell scripts in grpc_servicer ==="
find grpc_servicer -maxdepth 2 -type f \( -name "Makefile" -o -name "*.sh" \) 2>/dev/null
echo -e "\n=== Check README or dev docs in grpc_servicer ==="
find grpc_servicer -maxdepth 2 -name "README*" -o -name "DEVELOPMENT*" 2>/dev/nullRepository: lightseekorg/smg
Length of output: 1907
🏁 Script executed:
#!/bin/bash
# Check for any imports of sglang conditional on availability
echo "=== Check conftest.py for sglang import guards ==="
find grpc_servicer -name "conftest.py" -exec cat {} \;
echo -e "\n=== Check if any test file has importorskip or pytest.mark.skipif for sglang ==="
grep -r "importorskip\|skipif\|sglang" grpc_servicer/tests/ --include="*.py" | grep -E "skip|import"Repository: lightseekorg/smg
Length of output: 724
🏁 Script executed:
#!/bin/bash
# Check DEVELOPMENT.md for test setup instructions
echo "=== grpc_servicer/DEVELOPMENT.md ==="
cat grpc_servicer/DEVELOPMENT.md
echo -e "\n=== grpc_servicer/README.md (first 100 lines) ==="
head -100 grpc_servicer/README.mdRepository: lightseekorg/smg
Length of output: 3898
🏁 Script executed:
#!/bin/bash
# Look for any CI job that specifically runs grpc_servicer tests
echo "=== Search for unit-tests job definition ==="
grep -A 50 "unit-tests:" .github/workflows/pr-test-rust.yml | head -80
echo -e "\n=== Check if there's a separate grpc_servicer test workflow ==="
grep -l "grpc.servicer\|grpc_servicer" .github/workflows/*.ymlRepository: lightseekorg/smg
Length of output: 2437
Add sglang dependency to the dev extra, or document the required .[dev,sglang] installation.
The new unit test file imports smg_grpc_servicer.sglang.*, which has module-level import sglang statements. Installing pip install -e .[dev] provides pytest but not sglang, causing ImportError during test collection. Since the dev extra is positioned for development/testing, it should either include sglang or be clearly documented as requiring .[dev,sglang].
Proposed fix if dev is intended to run the full test suite
dev = [
"pytest>=8",
"pytest-asyncio>=0.24",
+ "sglang>=0.5.10",
]📝 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.
| dev = [ | |
| "pytest>=8", | |
| "pytest-asyncio>=0.24", | |
| ] | |
| dev = [ | |
| "pytest>=8", | |
| "pytest-asyncio>=0.24", | |
| "sglang>=0.5.10", | |
| ] |
🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/pyproject.toml` around lines 35 - 38, The dev extra in
pyproject.toml currently lists only pytest packages but tests import
smg_grpc_servicer.sglang which does a module-level import sglang, so add sglang
to the dev extra (e.g., include "sglang" with an appropriate version spec) so
pip install -e .[dev] pulls it in, or alternatively update the repository README
/ CONTRIBUTING to state that developers must install with .[dev,sglang]; modify
the dev extras entry in pyproject.toml (the "dev" list) or the docs accordingly
to resolve ImportError during test collection.
| def test_handle_loop_not_gated_on_gracefully_exit(): | ||
| """Regression guard for the drain invariant. | ||
|
|
||
| handle_loop must run `while True:`, not `while not self.gracefully_exit:`. | ||
| If it gates on the flag, begin_drain() exits the loop on the next ZMQ | ||
| recv, stalling in-flight streaming requests and defeating the drain. | ||
| See docs/superpowers/specs/2026-04-16-grpc-graceful-shutdown-design.md. | ||
| """ | ||
| src = inspect.getsource(GrpcRequestManager.handle_loop) | ||
|
|
||
| assert "while True:" in src, ( | ||
| "handle_loop must run `while True:` so scheduler outputs keep " | ||
| "flowing to in-flight streams during drain" | ||
| ) | ||
| assert "while not self.gracefully_exit" not in src, ( | ||
| "handle_loop must not gate on gracefully_exit; doing so stalls " | ||
| "streaming requests after begin_drain() is called" | ||
| ) |
There was a problem hiding this comment.
🧹 Nitpick | 🔵 Trivial
Make the regression guard inspect the actual AST loop condition.
The current string checks can miss regressions such as while self.gracefully_exit is False: and can also be fooled by comments/docstrings. Check the top-level ast.While condition instead.
More robust source-inspection guard
+import ast
import inspect
+import textwrap
from unittest.mock import MagicMock def test_handle_loop_not_gated_on_gracefully_exit():
@@
- src = inspect.getsource(GrpcRequestManager.handle_loop)
-
- assert "while True:" in src, (
+ src = textwrap.dedent(inspect.getsource(GrpcRequestManager.handle_loop))
+ tree = ast.parse(src)
+ func = tree.body[0]
+ top_level_while = next(
+ (node for node in func.body if isinstance(node, ast.While)),
+ None,
+ )
+
+ assert top_level_while is not None
+ assert isinstance(top_level_while.test, ast.Constant)
+ assert top_level_while.test.value is True, (
"handle_loop must run `while True:` so scheduler outputs keep "
"flowing to in-flight streams during drain"
)
- assert "while not self.gracefully_exit" not in src, (
- "handle_loop must not gate on gracefully_exit; doing so stalls "
- "streaming requests after begin_drain() is called"
- )🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed.
In `@grpc_servicer/tests/sglang/test_graceful_shutdown.py` around lines 58 - 75,
Replace the fragile string checks in
test_handle_loop_not_gated_on_gracefully_exit with an AST-based assertion: parse
inspect.getsource(GrpcRequestManager.handle_loop) into an AST, find the
top-level ast.While node(s) in the function body, and assert at least one while
loop has a condition that is a literal True (ast.Constant True or
ast.NameConstant True) and that none of the top-level while conditions represent
a negated or false condition that references self.gracefully_exit (e.g.,
ast.UnaryOp(Not)/ast.Compare against False or direct Name/Attribute comparing to
False); update the test to reference the GrpcRequestManager.handle_loop function
name when locating the AST node so comments/docstrings or incidental text cannot
bypass the check.
Remove the test file + pytest config. Validation of the fix came from the end-to-end K8s manual test (GLM-5-FP8 PD-disagg decoder, streaming request + kubectl delete), which is the behaviour this PR actually protects. The code-level unit tests were thin (flag-set + source inspection) and added test infrastructure for no load-bearing signal. Signed-off-by: Kangyan Zhou <zky314343421@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: ac0f8cb08c
ℹ️ 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".
| # timeout; K8s terminationGracePeriodSeconds SIGKILL is the backstop. | ||
| # rid_to_state is only mutated from the event-loop thread, so a | ||
| # list() snapshot is safe without additional guarding. | ||
| while True: |
There was a problem hiding this comment.
Stop accepting new RPCs before waiting for drain completion
The new drain loop runs before the server is stopped, so the process still accepts Generate/Embed RPCs while waiting for rid_to_state to reach zero. In environments with persistent gRPC clients (or direct pod traffic), new requests can continue to enter during drain and keep remain_num_req non-zero, causing termination to stall until external SIGKILL and defeating the graceful-shutdown goal. Move listener shutdown/rejection earlier (or explicitly reject inference RPCs once draining starts) so the drain set is bounded.
Useful? React with 👍 / 👎.
| if self.gracefully_exit: | ||
| logger.debug(f"ZMQ recv interrupted during shutdown: {e}") | ||
| break | ||
| logger.error(f"ZMQ error in handle loop: {e}\n{get_exception_traceback()}") | ||
| else: | ||
| logger.error(f"ZMQ error in handle loop: {e}\n{get_exception_traceback()}") | ||
| break |
There was a problem hiding this comment.
Avoid treating drain mode as terminal ZMQ shutdown
begin_drain() now sets gracefully_exit before real shutdown, but this branch still treats any ZMQError under that flag as shutdown and exits handle_loop. During a long drain window, a transient socket error would stop forwarding scheduler outputs to request queues, which can leave in-flight streams stuck and prevent the drain loop from ever completing cleanly. Drain and terminal-shutdown states should be separated for this error path.
Useful? React with 👍 / 👎.
Addresses PR smg-project#1172 review feedback: - Detect a dead handle_loop task in the drain loop and break out instead of blocking until SIGKILL. If the scheduler or ZMQ channel dies during drain, scheduler outputs can never reach rid_to_state so in-flight rids would otherwise linger forever. Expose the task via a new GrpcRequestManager.handle_loop_task attribute. - Reject new Generate/Embed RPCs with UNAVAILABLE once gracefully_exit is set, mirroring the existing HealthCheck pattern. K8s already removes the pod from Service Endpoints on NOT_SERVING, but persistent clients and direct pod traffic could otherwise keep feeding work into rid_to_state and stall the drain loop indefinitely. - Truncate remaining_rids in the drain log to the first 10 entries (plus an ellipsis when more follow) to avoid log flooding under high concurrency. Signed-off-by: Kangyan Zhou <zky314343421@gmail.com>
There was a problem hiding this comment.
Caution
Some comments are outside the diff and can’t be posted inline due to platform limitations.
⚠️ Outside diff range comments (2)
grpc_servicer/smg_grpc_servicer/sglang/server.py (1)
279-286:⚠️ Potential issue | 🟠 MajorCoordinate drain with the temporary SIGTERM handler and warmup thread.
This handler only covers signals received after it is registered. A SIGTERM handled earlier by
GrpcSignalHandler.sigterm_handleronly setsrequest_manager.gracefully_exit, sostop_eventis never set. Also, the warmup thread is not told that drain started; it can later callset_serving()or hit the newUNAVAILABLEGenerate/Embed path and kill the process from warmup error handling. Add a shared shutdown event and consume any pre-existinggracefully_exitafter installing this handler.Proposed direction
+ shutdown_started = threading.Event() + # Start warmup in a separate thread warmup_thread = threading.Thread( target=_wait_and_warmup_grpc, - args=(server_args, health_servicer), + args=(server_args, health_servicer, shutdown_started), ) warmup_thread.start() @@ def signal_handler(): logger.info("Received shutdown signal") + shutdown_started.set() # Flip health to NOT_SERVING and mark the request manager as # draining so K8s stops routing new traffic to this pod and # in-flight requests are allowed to finish. servicer.begin_drain() stop_event.set() @@ for sig in (signal.SIGTERM, signal.SIGINT): loop.add_signal_handler(sig, signal_handler) + + # Close the window where GrpcRequestManager's temporary SIGTERM handler + # may already have marked the manager as draining before this handler + # was installed. + if request_manager.gracefully_exit: + signal_handler()Also gate warmup before setting SERVING and before fatal warmup failure handling:
def _wait_and_warmup_grpc( server_args: ServerArgs, health_servicer: SGLangHealthServicer | None = None, + shutdown_started: threading.Event | None = None, ): """Wait for gRPC server to be ready and execute warmup.""" + if shutdown_started is not None and shutdown_started.is_set(): + return + if not server_args.skip_server_warmup: - if not _execute_grpc_server_warmup(server_args): + if not _execute_grpc_server_warmup(server_args, shutdown_started): return @@ # Mark health service as SERVING after warmup completes + if shutdown_started is not None and shutdown_started.is_set(): + return if health_servicer: health_servicer.set_serving()🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@grpc_servicer/smg_grpc_servicer/sglang/server.py` around lines 279 - 286, The shutdown flow doesn't coordinate the existing GrpcSignalHandler.sigterm_handler and the new local signal handler/warmup thread: sigterm may have set request_manager.gracefully_exit earlier but never sets stop_event or notifies the warmup thread, allowing warmup to later call set_serving() or abort fatally. Fix by introducing and using a shared shutdown event (e.g., shutdown_event) that both the local handler and warmup thread observe; update the local signal handler (the closure passed to loop.add_signal_handler) to set shutdown_event and call servicer.begin_drain() and stop_event.set(); after registering the handler, immediately check request_manager.gracefully_exit and if true set shutdown_event/stop_event and call servicer.begin_drain() so pre-existing SIGTERM is consumed; finally gate warmup (in the warmup thread code that calls set_serving() and handles fatal warmup errors) to check shutdown_event before setting SERVING or taking fatal-exit paths so warmup aborts cleanly instead of killing the process.grpc_servicer/smg_grpc_servicer/sglang/request_manager.py (1)
492-532:⚠️ Potential issue | 🟠 MajorKeep the handle loop alive after non-ZMQ processing errors during drain.
Line 531 still breaks on any generic exception once
gracefully_exitis true. That means one recoverable output-processing error during drain can completehandle_loop_task;server.pythen treats forwarding as dead and proceeds to destructive shutdown, aborting unrelated in-flight streams. Only cancellation or unrecoverable ZMQ/socket teardown should stop forwarding.Proposed fix
except Exception as e: logger.error(f"Handle loop error: {e}\n{get_exception_traceback()}") - if self.gracefully_exit: - break + # Keep forwarding scheduler outputs during drain; socket teardown + # is handled by the ZMQError path and task cancellation is not + # swallowed here. + continue🤖 Prompt for AI Agents
Verify each finding against the current code and only fix it if needed. In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py` around lines 492 - 532, The handle loop currently breaks on any generic Exception (in the except Exception as e: block) which lets a recoverable processing error during drain terminate handle_loop_task and trigger shutdown; instead, change the except Exception handler in the handle_loop (the loop that calls self.recv_from_scheduler.recv_pyobj() and dispatches to _handle_batch_output, _handle_embedding_output, _handle_health_check_output, _handle_abort_req, GetLoadsReqOutput/get_loads_communicator) so that it logs the error and continues the loop (i.e., do not break when self.gracefully_exit is True); only exit the loop on explicit cancellation (asyncio.CancelledError) or unrecoverable ZMQ/socket errors (zmq.error.ZMQError as already handled) so forwarding stays alive for recoverable processing failures.
🤖 Prompt for all review comments with AI agents
Verify each finding against the current code and only fix it if needed.
Outside diff comments:
In `@grpc_servicer/smg_grpc_servicer/sglang/request_manager.py`:
- Around line 492-532: The handle loop currently breaks on any generic Exception
(in the except Exception as e: block) which lets a recoverable processing error
during drain terminate handle_loop_task and trigger shutdown; instead, change
the except Exception handler in the handle_loop (the loop that calls
self.recv_from_scheduler.recv_pyobj() and dispatches to _handle_batch_output,
_handle_embedding_output, _handle_health_check_output, _handle_abort_req,
GetLoadsReqOutput/get_loads_communicator) so that it logs the error and
continues the loop (i.e., do not break when self.gracefully_exit is True); only
exit the loop on explicit cancellation (asyncio.CancelledError) or unrecoverable
ZMQ/socket errors (zmq.error.ZMQError as already handled) so forwarding stays
alive for recoverable processing failures.
In `@grpc_servicer/smg_grpc_servicer/sglang/server.py`:
- Around line 279-286: The shutdown flow doesn't coordinate the existing
GrpcSignalHandler.sigterm_handler and the new local signal handler/warmup
thread: sigterm may have set request_manager.gracefully_exit earlier but never
sets stop_event or notifies the warmup thread, allowing warmup to later call
set_serving() or abort fatally. Fix by introducing and using a shared shutdown
event (e.g., shutdown_event) that both the local handler and warmup thread
observe; update the local signal handler (the closure passed to
loop.add_signal_handler) to set shutdown_event and call servicer.begin_drain()
and stop_event.set(); after registering the handler, immediately check
request_manager.gracefully_exit and if true set shutdown_event/stop_event and
call servicer.begin_drain() so pre-existing SIGTERM is consumed; finally gate
warmup (in the warmup thread code that calls set_serving() and handles fatal
warmup errors) to check shutdown_event before setting SERVING or taking
fatal-exit paths so warmup aborts cleanly instead of killing the process.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Organization UI
Review profile: ASSERTIVE
Plan: Pro
Run ID: 941fbc32-aa70-4ba2-a8c8-202a195ec206
📒 Files selected for processing (3)
grpc_servicer/smg_grpc_servicer/sglang/request_manager.pygrpc_servicer/smg_grpc_servicer/sglang/server.pygrpc_servicer/smg_grpc_servicer/sglang/servicer.py
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 90917f7e98
ℹ️ 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 self.request_manager.gracefully_exit: | ||
| await context.abort( | ||
| grpc.StatusCode.UNAVAILABLE, | ||
| "Server is shutting down", | ||
| ) |
There was a problem hiding this comment.
Exempt internal warmup RPCs from drain-time rejection
The new gracefully_exit guard aborts every Generate/Embed call, including the server’s own warmup calls, so if SIGTERM arrives before warmup finishes the warmup thread can receive UNAVAILABLE and follow its existing fatal path in server.py (_execute_grpc_server_warmup catches the RPC error and calls kill_process_tree). In that startup/shutdown overlap, the process is force-killed instead of completing the intended drain, which can still drop in-flight direct-to-pod requests.
Useful? React with 👍 / 👎.
Signed-off-by: Kangyan Zhou <zky314343421@gmail.com>
There was a problem hiding this comment.
💡 Codex Review
Here are some automated review suggestions for this pull request.
Reviewed commit: 95ed234a16
ℹ️ 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".
| remaining_rids = [ | ||
| rid | ||
| for rid, state in servicer.request_manager.rid_to_state.items() | ||
| if not state.finished | ||
| ] |
There was a problem hiding this comment.
Include active RPC handlers in drain completion check
The drain loop decides shutdown readiness solely from request_manager.rid_to_state entries marked unfinished, but that map does not cover every still-active RPC handler. In particular, generate_request can be mid-flight with no active RID between its n>1 prefix phase and launching phase-2 generators, so this check can read zero and proceed to servicer.shutdown() while a client request is still being processed, causing that request to be aborted despite graceful-drain mode. Consider tracking handler-level in-flight RPCs (or a broader active-request counter) instead of relying only on RID state.
Useful? React with 👍 / 👎.
| # Flip health to NOT_SERVING and mark the request manager as | ||
| # draining so K8s stops routing new traffic to this pod and | ||
| # in-flight requests are allowed to finish. | ||
| servicer.begin_drain() |
There was a problem hiding this comment.
Prevent warmup from restoring SERVING after drain begins
Calling begin_drain() here marks health as NOT_SERVING, but the warmup thread later calls health_servicer.set_serving() unconditionally in _wait_and_warmup_grpc. If SIGTERM arrives during warmup and drain lasts for in-flight requests, readiness can flip back to SERVING during drain, so Kubernetes may resume routing traffic to a pod that is intentionally rejecting inference RPCs with UNAVAILABLE. Drain mode should make health status sticky (or warmup should skip set_serving once gracefully_exit is true).
Useful? React with 👍 / 👎.
|
This pull request has been automatically marked as stale because it has not had any activity within 14 days. It will be automatically closed if no further activity occurs within 16 days. Leave a comment if you feel this pull request should remain open. Thank you! |
|
This pull request has been automatically closed due to inactivity. Please feel free to reopen if you intend to continue working on it. Thank you! |
Summary
request_manager.shutdown()eagerly (which cancels tasks and pushes{"error": "Server shutting down"}into every request's out_queue).TokenizerManager.sigterm_watchdogpattern: a newbegin_drain()onGrpcRequestManagerandSGLangSchedulerServicerthat only flipsgracefully_exit(+ health NOT_SERVING), and a drain loop inserve_grpcthat pollsrid_to_stateevery 5 s (excluding already-finished entries) before the destructive teardown runs.SGL_FORCE_SHUTDOWNescape hatch honoured, same as HTTP.GrpcRequestManager.handle_loop: changewhile not self.gracefully_exit:towhile True:so scheduler outputs keep flowing to in-flight streaming requests during drain (the old gate stalled streams the momentbegin_drain()flipped the flag). HTTP'shandle_loopuses the samewhile True:for the same reason; the loop now exits only via asyncio task cancellation fromshutdown()or an unrecoverableZMQError.terminationGracePeriodSeconds, with SIGKILL as the ultimate backstop (same as HTTP mode).Test plan
pre-commit run --files …clean across all touched filesManual K8s validation — PD-disagg decoder under load on
prod-sci-us-central1-1with GLM-5-FP8, TP=8, PP=1. Same streaming request, samekubectl deleteat T+15 s:v0.5.10-fixes(pre-fix baseline){"error":"Stream error: Server shutting down"}after 862 SSE lines.v0.5.10-graceful-shutdown-test(this PR applied on top of the baseline)finish_reason:"stop"after 4418 SSE lines.Confirms the drain loop lets in-flight streaming requests complete cleanly before the gRPC servicer tears down on K8s scale-down.
Summary by CodeRabbit
Release Notes
SGL_FORCE_SHUTDOWNenvironment variable to skip graceful draining when needed.