From f04cb75025a6f4a6d07c70be42ed7b4e62f10cc5 Mon Sep 17 00:00:00 2001 From: Kangyan Zhou Date: Thu, 16 Apr 2026 16:20:57 -0700 Subject: [PATCH 1/7] fix(grpc): drain in-flight requests on SIGTERM for graceful scale-down 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 --- .../sglang/request_manager.py | 27 +++++-- .../smg_grpc_servicer/sglang/server.py | 36 ++++++++- .../smg_grpc_servicer/sglang/servicer.py | 11 +++ .../tests/sglang/test_graceful_shutdown.py | 75 +++++++++++++++++++ 4 files changed, 141 insertions(+), 8 deletions(-) create mode 100644 grpc_servicer/tests/sglang/test_graceful_shutdown.py diff --git a/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py b/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py index 8814f7f702..401b1ab476 100644 --- a/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py +++ b/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py @@ -479,8 +479,13 @@ async def handle_loop(self): """ Main event loop - processes outputs from scheduler. Mimics TokenizerManager's handle_loop. + + Runs until the task is cancelled (by shutdown()) or an unrecoverable + ZMQ error. It must not gate on ``gracefully_exit`` — during drain + the health flag is set but scheduler outputs must still be forwarded + so in-flight streaming requests can finish their token sequences. """ - while not self.gracefully_exit: + while True: try: # Receive from scheduler recv_obj = await self.recv_from_scheduler.recv_pyobj() @@ -508,16 +513,14 @@ async def handle_loop(self): logger.warning(f"Unknown output type: {type(recv_obj)}") except zmq.error.Again: - # Timeout, check if we should exit - if self.gracefully_exit: - break + # Timeout on non-blocking recv; keep polling. continue except zmq.error.ZMQError as e: - # Socket closed or other ZMQ error - exit cleanly if shutting down + # Socket closed or unrecoverable ZMQ error. 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 except Exception as e: logger.error(f"Handle loop error: {e}\n{get_exception_traceback()}") @@ -817,6 +820,16 @@ def record_request_for_crash_dump(self, obj): } ) + def begin_drain(self) -> None: + """Mark the manager as draining. + + Idempotent and non-destructive: flips gracefully_exit so the health + servicer reports NOT_SERVING and the server-side drain loop can + begin, but does not cancel tasks or enqueue shutdown errors on + in-flight requests. Safe to call from a synchronous signal handler. + """ + self.gracefully_exit = True + async def shutdown(self): """Gracefully shutdown the request manager.""" logger.info("Shutting down GrpcRequestManager") diff --git a/grpc_servicer/smg_grpc_servicer/sglang/server.py b/grpc_servicer/smg_grpc_servicer/sglang/server.py index 99230711bb..29124c9c71 100644 --- a/grpc_servicer/smg_grpc_servicer/sglang/server.py +++ b/grpc_servicer/smg_grpc_servicer/sglang/server.py @@ -21,7 +21,7 @@ from sglang.srt.disaggregation.utils import FAKE_BOOTSTRAP_HOST, DisaggregationMode from sglang.srt.managers.disagg_service import start_disagg_service from sglang.srt.server_args import ServerArgs -from sglang.srt.utils import kill_process_tree +from sglang.srt.utils import get_bool_env_var, kill_process_tree from sglang.utils import get_exception_traceback from smg_grpc_proto import sglang_scheduler_pb2, sglang_scheduler_pb2_grpc @@ -276,6 +276,10 @@ def _cert_config_fetcher(): def signal_handler(): logger.info("Received shutdown signal") + # 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): @@ -284,6 +288,36 @@ def signal_handler(): try: await stop_event.wait() finally: + # Drain phase: wait for in-flight requests to finish naturally. + # Mirrors tokenizer_manager.sigterm_watchdog. No in-process + # 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: + 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) + logger.info("Shutting down gRPC server") # Shutdown request manager first - this closes ZMQ sockets and stops background tasks diff --git a/grpc_servicer/smg_grpc_servicer/sglang/servicer.py b/grpc_servicer/smg_grpc_servicer/sglang/servicer.py index aa2ee916e4..9ce71b56bc 100644 --- a/grpc_servicer/smg_grpc_servicer/sglang/servicer.py +++ b/grpc_servicer/smg_grpc_servicer/sglang/servicer.py @@ -1098,6 +1098,17 @@ def _create_completion_response( ), ) + def begin_drain(self) -> None: + """Mark the service as draining so health flips to NOT_SERVING. + + Non-destructive: does not cancel in-flight requests. Safe to call + from a synchronous signal handler. Must be followed by shutdown() + (after the drain loop completes) for final cleanup. + """ + if self.health_servicer: + self.health_servicer.set_not_serving() + self.request_manager.begin_drain() + async def shutdown(self): """Shutdown the service.""" logger.info("Shutting down gRPC service") diff --git a/grpc_servicer/tests/sglang/test_graceful_shutdown.py b/grpc_servicer/tests/sglang/test_graceful_shutdown.py new file mode 100644 index 0000000000..3a8b0ca7df --- /dev/null +++ b/grpc_servicer/tests/sglang/test_graceful_shutdown.py @@ -0,0 +1,75 @@ +"""Unit tests for graceful shutdown wiring in the SGLang gRPC servicer. + +Covers the drain-signal plumbing added for K8s scale-down: the request +manager's begin_drain() method and the health servicer's response to the +gracefully_exit flag. +""" + +from unittest.mock import MagicMock + +import pytest +from grpc_health.v1 import health_pb2 +from smg_grpc_servicer.sglang.health_servicer import SGLangHealthServicer +from smg_grpc_servicer.sglang.request_manager import GrpcRequestManager + +_SENTINEL_RID = "rid-under-drain" +_SENTINEL_TASK = "fake-task" + + +def _make_bare_request_manager() -> GrpcRequestManager: + """Return a GrpcRequestManager with only the attributes begin_drain touches. + + The real __init__ constructs ZMQ sockets, a scheduler channel, etc., none + of which begin_drain depends on. object.__new__ + explicit attribute + assignment gives us an isolated unit under test. We seed rid_to_state + and asyncio_tasks with sentinels so the non-destructive assertion is + meaningful (an empty-state fixture would pass trivially). + """ + mgr = object.__new__(GrpcRequestManager) + mgr.gracefully_exit = False + mgr.rid_to_state = {_SENTINEL_RID: MagicMock(finished=False)} + mgr.asyncio_tasks = {_SENTINEL_TASK} + return mgr + + +def test_begin_drain_sets_flag(): + mgr = _make_bare_request_manager() + + mgr.begin_drain() + + assert mgr.gracefully_exit is True + # begin_drain must not cancel tasks or evict in-flight request state. + assert _SENTINEL_RID in mgr.rid_to_state + assert mgr.asyncio_tasks == {_SENTINEL_TASK} + + +def test_begin_drain_idempotent(): + mgr = _make_bare_request_manager() + + mgr.begin_drain() + mgr.begin_drain() + + assert mgr.gracefully_exit is True + assert _SENTINEL_RID in mgr.rid_to_state + assert mgr.asyncio_tasks == {_SENTINEL_TASK} + + +@pytest.mark.asyncio +async def test_health_reflects_drain_flag(): + request_manager = MagicMock() + request_manager.gracefully_exit = False + + health = SGLangHealthServicer(request_manager=request_manager, scheduler_info={}) + health.set_serving() + + ctx = MagicMock() + req = health_pb2.HealthCheckRequest(service="") + + # Sanity check: SERVING before drain starts. + resp = await health.Check(req, ctx) + assert resp.status == health_pb2.HealthCheckResponse.SERVING + + # After begin_drain flips the flag, health must flip to NOT_SERVING. + request_manager.gracefully_exit = True + resp = await health.Check(req, ctx) + assert resp.status == health_pb2.HealthCheckResponse.NOT_SERVING From d6aed74f5fcb0df239ae09815d7cd60185a029a3 Mon Sep 17 00:00:00 2001 From: Kangyan Zhou Date: Thu, 16 Apr 2026 16:23:49 -0700 Subject: [PATCH 2/7] test(grpc-servicer): add pytest config so unit tests are discoverable 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 --- grpc_servicer/pyproject.toml | 14 ++++++++++++++ grpc_servicer/tests/conftest.py | 10 ++++++++++ 2 files changed, 24 insertions(+) create mode 100644 grpc_servicer/tests/conftest.py diff --git a/grpc_servicer/pyproject.toml b/grpc_servicer/pyproject.toml index 14172f2807..ad6dd4a3c1 100644 --- a/grpc_servicer/pyproject.toml +++ b/grpc_servicer/pyproject.toml @@ -32,6 +32,10 @@ classifiers = [ [project.optional-dependencies] vllm = ["vllm>=0.19.0"] sglang = ["sglang>=0.5.10"] +dev = [ + "pytest>=8", + "pytest-asyncio>=0.24", +] [project.urls] Homepage = "https://github.com/lightseekorg/smg" @@ -40,3 +44,13 @@ Repository = "https://github.com/lightseekorg/smg" [tool.setuptools.packages.find] where = ["."] include = ["smg_grpc_servicer*"] + +[tool.pytest.ini_options] +testpaths = ["tests"] +python_files = ["test_*.py"] +python_classes = ["Test*"] +python_functions = ["test_*"] +asyncio_mode = "strict" +markers = [ + "unit: mark test as a unit test (no GPU or scheduler required)", +] diff --git a/grpc_servicer/tests/conftest.py b/grpc_servicer/tests/conftest.py new file mode 100644 index 0000000000..12c7c1a3d3 --- /dev/null +++ b/grpc_servicer/tests/conftest.py @@ -0,0 +1,10 @@ +"""Pytest configuration for smg-grpc-servicer unit tests. + +These tests run without GPU resources or a live scheduler process. +""" + + +def pytest_configure(config): + config.addinivalue_line( + "markers", "unit: mark test as a unit test (no GPU or scheduler required)" + ) From db30688d37b7f32c96ba2a00171cf92c0bf54d14 Mon Sep 17 00:00:00 2001 From: Kangyan Zhou Date: Thu, 16 Apr 2026 16:25:19 -0700 Subject: [PATCH 3/7] test(grpc): add regression guard for handle_loop drain invariant MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- .../tests/sglang/test_graceful_shutdown.py | 21 +++++++++++++++++++ 1 file changed, 21 insertions(+) diff --git a/grpc_servicer/tests/sglang/test_graceful_shutdown.py b/grpc_servicer/tests/sglang/test_graceful_shutdown.py index 3a8b0ca7df..a5e6baaab4 100644 --- a/grpc_servicer/tests/sglang/test_graceful_shutdown.py +++ b/grpc_servicer/tests/sglang/test_graceful_shutdown.py @@ -5,6 +5,7 @@ gracefully_exit flag. """ +import inspect from unittest.mock import MagicMock import pytest @@ -54,6 +55,26 @@ def test_begin_drain_idempotent(): assert mgr.asyncio_tasks == {_SENTINEL_TASK} +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" + ) + + @pytest.mark.asyncio async def test_health_reflects_drain_flag(): request_manager = MagicMock() From 263fd2c508890f9f71259ce04a0bb5e0c5772281 Mon Sep 17 00:00:00 2001 From: Kangyan Zhou Date: Thu, 16 Apr 2026 16:28:07 -0700 Subject: [PATCH 4/7] refactor(grpc): remove stub sigterm_watchdog task MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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 --- .../smg_grpc_servicer/sglang/request_manager.py | 13 ++++--------- 1 file changed, 4 insertions(+), 9 deletions(-) diff --git a/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py b/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py index 401b1ab476..26b9cb97c3 100644 --- a/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py +++ b/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py @@ -921,20 +921,15 @@ def auto_create_handle_loop(self): self.event_loop = loop # We only add signal handler when the tokenizer manager is in the main thread - # due to the CPython limitation. + # due to the CPython limitation. The SIGTERM handler here is a startup-window + # fallback; once serve_grpc runs, it overrides SIGTERM with the drain-aware + # signal_handler in server.py. SIGQUIT stays owned here to forward scheduler + # crashes to a process-tree kill. if threading.current_thread() is threading.main_thread(): signal_handler = GrpcSignalHandler(self) loop.add_signal_handler(signal.SIGTERM, signal_handler.sigterm_handler) - # Update the signal handler for the process. It overrides the sigquit handler in the launch phase. loop.add_signal_handler(signal.SIGQUIT, signal_handler.running_phase_sigquit_handler) - self.asyncio_tasks.add(loop.create_task(print_exception_wrapper(self.sigterm_watchdog))) - - async def sigterm_watchdog(self): - """Watchdog to handle SIGTERM gracefully, matching TokenizerManager pattern.""" - while not self.gracefully_exit: - await asyncio.sleep(1.0) - def _req_stats_init( self, obj: TokenizedGenerateReqInput | TokenizedEmbeddingReqInput, From ac0f8cb08cdd8f211e95a8b5103f40508dd9f4cf Mon Sep 17 00:00:00 2001 From: Kangyan Zhou Date: Mon, 20 Apr 2026 23:41:03 -0700 Subject: [PATCH 5/7] test(grpc): drop unit tests for graceful shutdown 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 --- grpc_servicer/pyproject.toml | 14 --- grpc_servicer/tests/conftest.py | 10 -- .../tests/sglang/test_graceful_shutdown.py | 96 ------------------- 3 files changed, 120 deletions(-) delete mode 100644 grpc_servicer/tests/conftest.py delete mode 100644 grpc_servicer/tests/sglang/test_graceful_shutdown.py diff --git a/grpc_servicer/pyproject.toml b/grpc_servicer/pyproject.toml index ad6dd4a3c1..14172f2807 100644 --- a/grpc_servicer/pyproject.toml +++ b/grpc_servicer/pyproject.toml @@ -32,10 +32,6 @@ classifiers = [ [project.optional-dependencies] vllm = ["vllm>=0.19.0"] sglang = ["sglang>=0.5.10"] -dev = [ - "pytest>=8", - "pytest-asyncio>=0.24", -] [project.urls] Homepage = "https://github.com/lightseekorg/smg" @@ -44,13 +40,3 @@ Repository = "https://github.com/lightseekorg/smg" [tool.setuptools.packages.find] where = ["."] include = ["smg_grpc_servicer*"] - -[tool.pytest.ini_options] -testpaths = ["tests"] -python_files = ["test_*.py"] -python_classes = ["Test*"] -python_functions = ["test_*"] -asyncio_mode = "strict" -markers = [ - "unit: mark test as a unit test (no GPU or scheduler required)", -] diff --git a/grpc_servicer/tests/conftest.py b/grpc_servicer/tests/conftest.py deleted file mode 100644 index 12c7c1a3d3..0000000000 --- a/grpc_servicer/tests/conftest.py +++ /dev/null @@ -1,10 +0,0 @@ -"""Pytest configuration for smg-grpc-servicer unit tests. - -These tests run without GPU resources or a live scheduler process. -""" - - -def pytest_configure(config): - config.addinivalue_line( - "markers", "unit: mark test as a unit test (no GPU or scheduler required)" - ) diff --git a/grpc_servicer/tests/sglang/test_graceful_shutdown.py b/grpc_servicer/tests/sglang/test_graceful_shutdown.py deleted file mode 100644 index a5e6baaab4..0000000000 --- a/grpc_servicer/tests/sglang/test_graceful_shutdown.py +++ /dev/null @@ -1,96 +0,0 @@ -"""Unit tests for graceful shutdown wiring in the SGLang gRPC servicer. - -Covers the drain-signal plumbing added for K8s scale-down: the request -manager's begin_drain() method and the health servicer's response to the -gracefully_exit flag. -""" - -import inspect -from unittest.mock import MagicMock - -import pytest -from grpc_health.v1 import health_pb2 -from smg_grpc_servicer.sglang.health_servicer import SGLangHealthServicer -from smg_grpc_servicer.sglang.request_manager import GrpcRequestManager - -_SENTINEL_RID = "rid-under-drain" -_SENTINEL_TASK = "fake-task" - - -def _make_bare_request_manager() -> GrpcRequestManager: - """Return a GrpcRequestManager with only the attributes begin_drain touches. - - The real __init__ constructs ZMQ sockets, a scheduler channel, etc., none - of which begin_drain depends on. object.__new__ + explicit attribute - assignment gives us an isolated unit under test. We seed rid_to_state - and asyncio_tasks with sentinels so the non-destructive assertion is - meaningful (an empty-state fixture would pass trivially). - """ - mgr = object.__new__(GrpcRequestManager) - mgr.gracefully_exit = False - mgr.rid_to_state = {_SENTINEL_RID: MagicMock(finished=False)} - mgr.asyncio_tasks = {_SENTINEL_TASK} - return mgr - - -def test_begin_drain_sets_flag(): - mgr = _make_bare_request_manager() - - mgr.begin_drain() - - assert mgr.gracefully_exit is True - # begin_drain must not cancel tasks or evict in-flight request state. - assert _SENTINEL_RID in mgr.rid_to_state - assert mgr.asyncio_tasks == {_SENTINEL_TASK} - - -def test_begin_drain_idempotent(): - mgr = _make_bare_request_manager() - - mgr.begin_drain() - mgr.begin_drain() - - assert mgr.gracefully_exit is True - assert _SENTINEL_RID in mgr.rid_to_state - assert mgr.asyncio_tasks == {_SENTINEL_TASK} - - -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" - ) - - -@pytest.mark.asyncio -async def test_health_reflects_drain_flag(): - request_manager = MagicMock() - request_manager.gracefully_exit = False - - health = SGLangHealthServicer(request_manager=request_manager, scheduler_info={}) - health.set_serving() - - ctx = MagicMock() - req = health_pb2.HealthCheckRequest(service="") - - # Sanity check: SERVING before drain starts. - resp = await health.Check(req, ctx) - assert resp.status == health_pb2.HealthCheckResponse.SERVING - - # After begin_drain flips the flag, health must flip to NOT_SERVING. - request_manager.gracefully_exit = True - resp = await health.Check(req, ctx) - assert resp.status == health_pb2.HealthCheckResponse.NOT_SERVING From 90917f7e986a96b9c5f84b5461b528b54806a869 Mon Sep 17 00:00:00 2001 From: Kangyan Zhou Date: Tue, 21 Apr 2026 18:33:25 -0700 Subject: [PATCH 6/7] fix(grpc): harden graceful drain against hangs and new-RPC leaks Addresses PR #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 --- .../sglang/request_manager.py | 7 ++++++- .../smg_grpc_servicer/sglang/server.py | 19 ++++++++++++++++++- .../smg_grpc_servicer/sglang/servicer.py | 17 +++++++++++++++++ 3 files changed, 41 insertions(+), 2 deletions(-) diff --git a/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py b/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py index 26b9cb97c3..d7a449bc66 100644 --- a/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py +++ b/grpc_servicer/smg_grpc_servicer/sglang/request_manager.py @@ -195,6 +195,10 @@ def __init__( # State Management (from TokenizerManager) self.rid_to_state: dict[str, GrpcReqState] = {} self.asyncio_tasks: set = set() + # Separate handle_loop ref so the drain loop can detect if scheduler + # output forwarding has died — without it, rid_to_state would never + # drain and the pod would wait for SIGKILL. + self.handle_loop_task: asyncio.Task | None = None self.gracefully_exit = False self.no_create_loop = False self.event_loop = None @@ -916,7 +920,8 @@ def auto_create_handle_loop(self): self.no_create_loop = True loop = get_or_create_event_loop() - self.asyncio_tasks.add(loop.create_task(print_exception_wrapper(self.handle_loop))) + self.handle_loop_task = loop.create_task(print_exception_wrapper(self.handle_loop)) + self.asyncio_tasks.add(self.handle_loop_task) self.event_loop = loop diff --git a/grpc_servicer/smg_grpc_servicer/sglang/server.py b/grpc_servicer/smg_grpc_servicer/sglang/server.py index 29124c9c71..c9a3da6ccc 100644 --- a/grpc_servicer/smg_grpc_servicer/sglang/server.py +++ b/grpc_servicer/smg_grpc_servicer/sglang/server.py @@ -298,6 +298,17 @@ def signal_handler(): logger.warning("SGL_FORCE_SHUTDOWN set; skipping drain") break + # If handle_loop has died (e.g. fatal ZMQError, scheduler + # crash), scheduler outputs will never be forwarded and + # in-flight rids can never be marked finished — abort the + # drain instead of blocking until SIGKILL. + handle_loop_task = servicer.request_manager.handle_loop_task + if handle_loop_task is not None and handle_loop_task.done(): + logger.warning( + "handle_loop task has terminated; aborting drain and proceeding to shutdown" + ) + 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. @@ -311,10 +322,16 @@ def signal_handler(): logger.info("Drain complete; no in-flight requests") break + # Truncate rid list in logs: under high load there can be + # thousands of in-flight requests and the full list would + # flood the log every poll interval. + log_rids = remaining_rids[:10] + if remain_num_req > 10: + log_rids = log_rids + ["..."] logger.info( "Gracefully exiting... Remaining number of requests %d. Remaining requests %s", remain_num_req, - remaining_rids, + log_rids, ) await asyncio.sleep(5) diff --git a/grpc_servicer/smg_grpc_servicer/sglang/servicer.py b/grpc_servicer/smg_grpc_servicer/sglang/servicer.py index 9ce71b56bc..647a2be155 100644 --- a/grpc_servicer/smg_grpc_servicer/sglang/servicer.py +++ b/grpc_servicer/smg_grpc_servicer/sglang/servicer.py @@ -212,6 +212,16 @@ async def Generate( """Handle generation requests with streaming responses.""" logger.info(f"Receive generation request: {request.request_id}") + # Reject new RPCs once drain has started. K8s removes the pod from + # Service Endpoints when health flips to NOT_SERVING, but persistent + # clients or direct pod traffic could otherwise keep feeding work + # into rid_to_state and stall the drain loop indefinitely. + if self.request_manager.gracefully_exit: + await context.abort( + grpc.StatusCode.UNAVAILABLE, + "Server is shutting down", + ) + try: # Convert gRPC request to internal format tokenized_req = self._convert_generate_request(request) @@ -271,6 +281,13 @@ async def Embed( """Handle embedding requests.""" logger.info(f"Receive embedding request: {request.request_id}") + # Reject new RPCs once drain has started (same rationale as Generate). + if self.request_manager.gracefully_exit: + await context.abort( + grpc.StatusCode.UNAVAILABLE, + "Server is shutting down", + ) + try: tokenized_req = self._convert_embed_request(request) From 95ed234a16e93fff4baa2f2d752fa595bdfb8038 Mon Sep 17 00:00:00 2001 From: Kangyan Zhou Date: Wed, 22 Apr 2026 23:10:59 -0700 Subject: [PATCH 7/7] chore: retrigger CI Signed-off-by: Kangyan Zhou