From b6cc897fea1cc8ee4b04de5f75edca8f54b79f50 Mon Sep 17 00:00:00 2001 From: fzyzcjy <5236035+fzyzcjy@users.noreply.github.com> Date: Wed, 15 Jul 2026 14:27:48 +0800 Subject: [PATCH] Fix non-existent abort mode in Scheduler.pause_generation and inline retract_all (#30673) --- python/sglang/srt/managers/schedule_batch.py | 17 +--- python/sglang/srt/managers/scheduler.py | 16 +++- .../test_scheduler_pause_generation.py | 91 +++++++++++++------ 3 files changed, 75 insertions(+), 49 deletions(-) diff --git a/python/sglang/srt/managers/schedule_batch.py b/python/sglang/srt/managers/schedule_batch.py index 8336955d1..83548ae0b 100755 --- a/python/sglang/srt/managers/schedule_batch.py +++ b/python/sglang/srt/managers/schedule_batch.py @@ -1724,8 +1724,7 @@ def retract_all( tree_cache: BasePrefixCache, hisparse_coordinator: Optional[HiSparseCoordinator], offload_kv: bool = True, -) -> List[Req]: - retracted_reqs = reqs +) -> None: for idx in range(len(reqs)): release_req( req=reqs[idx], @@ -1737,7 +1736,6 @@ def retract_all( hisparse_coordinator=hisparse_coordinator, offload_kv=offload_kv, ) - return retracted_reqs def compute_extend_logprob_start_len( @@ -2571,19 +2569,6 @@ class ScheduleBatch(ScheduleBatchDisaggregationDecodeMixin): evict_from_tree_cache(self.tree_cache, num_tokens) return self.token_to_kv_pool_allocator.available_size() >= num_tokens - def retract_all(self, server_args: ServerArgs, offload_kv: bool = True): - retracted_reqs = retract_all( - reqs=self.reqs, - server_args=server_args, - req_to_token_pool=self.req_to_token_pool, - token_to_kv_pool_allocator=self.token_to_kv_pool_allocator, - tree_cache=self.tree_cache, - hisparse_coordinator=self.hisparse_coordinator, - offload_kv=offload_kv, - ) - self.reqs = [] - return retracted_reqs - def retract_decode( self, server_args: ServerArgs ) -> Tuple[List[Req], float, List[Req]]: diff --git a/python/sglang/srt/managers/scheduler.py b/python/sglang/srt/managers/scheduler.py index 7eeeeb5b9..45c3c8c57 100644 --- a/python/sglang/srt/managers/scheduler.py +++ b/python/sglang/srt/managers/scheduler.py @@ -166,6 +166,7 @@ from sglang.srt.managers.schedule_batch import ( NextBatchPlan, Req, ScheduleBatch, + retract_all, ) from sglang.srt.managers.schedule_policy import ( AddReqResult, @@ -4064,6 +4065,7 @@ class Scheduler( raise NotImplementedError() def pause_generation(self, recv_req: PauseGenerationReqInput): + assert recv_req.mode in ("in_place", "retract") self._engine_paused = True if recv_req.mode == "in_place": @@ -4102,16 +4104,24 @@ class Scheduler( self.last_batch = None self.cur_batch_for_debug = None - if recv_req.mode == "retract" and not self.running_batch.is_empty(): + if not self.running_batch.is_empty(): self.running_batch.filter_batch() if len(self.running_batch.reqs) != 0: # Decode-side retract always rebootstraps (recomputes the KV from # the prefill), so skip the device->host KV offload that release_req # would otherwise do; the offloaded copy would be immediately # discarded. Non-decode modes ignore offload_kv (they never offload). - retracted_reqs = self.running_batch.retract_all( - self.server_args, offload_kv=False + retracted_reqs = self.running_batch.reqs + retract_all( + reqs=retracted_reqs, + server_args=self.server_args, + req_to_token_pool=self.running_batch.req_to_token_pool, + token_to_kv_pool_allocator=self.running_batch.token_to_kv_pool_allocator, + tree_cache=self.running_batch.tree_cache, + hisparse_coordinator=self.running_batch.hisparse_coordinator, + offload_kv=False, ) + self.running_batch.reqs = [] for req in retracted_reqs: if self.disaggregation_mode == DisaggregationMode.DECODE: if req.output_ids: diff --git a/test/registered/unit/managers/test_scheduler_pause_generation.py b/test/registered/unit/managers/test_scheduler_pause_generation.py index 596714afd..eac3c1feb 100644 --- a/test/registered/unit/managers/test_scheduler_pause_generation.py +++ b/test/registered/unit/managers/test_scheduler_pause_generation.py @@ -1,7 +1,7 @@ import unittest from collections import deque from types import SimpleNamespace -from unittest.mock import MagicMock +from unittest.mock import MagicMock, patch from sglang.test.ci.ci_register import register_cpu_ci from sglang.test.test_utils import maybe_stub_sgl_kernel @@ -98,14 +98,28 @@ class TestSchedulerPauseGeneration(unittest.TestCase): last_batch.filter_batch.assert_not_called() scheduler.running_batch.merge_batch.assert_not_called() - def test_abort_clears_state(self): - """abort mode should clear last_batch and cur_batch_for_debug.""" + def test_abort_mode_rejected_at_scheduler(self): + """abort mode must be rejected by the scheduler-side assert.""" + scheduler = self._new_scheduler() + + with self.assertRaises(AssertionError): + scheduler.pause_generation(PauseGenerationReqInput(mode="abort")) + + def test_default_mode_rejected_at_scheduler(self): + """bare PauseGenerationReqInput defaults to abort and must be rejected.""" + scheduler = self._new_scheduler() + + with self.assertRaises(AssertionError): + scheduler.pause_generation(PauseGenerationReqInput()) + + def test_retract_clears_last_batch_state(self): + """retract mode should clear last_batch and cur_batch_for_debug.""" scheduler = self._new_scheduler() scheduler.last_batch = MagicMock() scheduler.last_batch.forward_mode.is_extend.return_value = False scheduler.cur_batch_for_debug = MagicMock() - scheduler.pause_generation(PauseGenerationReqInput(mode="abort")) + scheduler.pause_generation(PauseGenerationReqInput(mode="retract")) self.assertTrue(scheduler._engine_paused) self.assertIsNone(scheduler.last_batch) @@ -121,24 +135,56 @@ class TestSchedulerPauseGeneration(unittest.TestCase): scheduler.waiting_queue = [] scheduler._add_request_to_queue = MagicMock() - retracted = [MagicMock(), MagicMock()] - scheduler.running_batch.retract_all.return_value = retracted scheduler.running_batch.filter_batch = MagicMock() scheduler.server_args = MagicMock() + reqs_before = scheduler.running_batch.reqs + + with patch("sglang.srt.managers.scheduler.retract_all") as mock_retract_all: + scheduler.pause_generation(PauseGenerationReqInput(mode="retract")) + + self.assertTrue(scheduler._engine_paused) + mock_retract_all.assert_called_once() + self.assertIs(mock_retract_all.call_args.kwargs["reqs"], reqs_before) + self.assertEqual(scheduler.running_batch.reqs, []) + self.assertEqual(scheduler._add_request_to_queue.call_count, 2) + self.assertEqual( + [call.args[0] for call in scheduler._add_request_to_queue.call_args_list], + reqs_before, + ) + self.assertIsNone(scheduler.chunked_req) + + def test_retract_empty_running_batch_requeues_nothing(self): + """retract with empty running_batch must not release or requeue any request.""" + scheduler = self._new_scheduler() + scheduler.waiting_queue = [] + original_reqs = scheduler.running_batch.reqs scheduler.pause_generation(PauseGenerationReqInput(mode="retract")) self.assertTrue(scheduler._engine_paused) - scheduler.running_batch.retract_all.assert_called_once() - self.assertEqual(scheduler._add_request_to_queue.call_count, 2) - self.assertIsNone(scheduler.chunked_req) + self.assertEqual(len(scheduler.waiting_queue), 0) + self.assertIs(scheduler.running_batch.reqs, original_reqs) + + def test_retract_drains_overlap_queue(self): + """retract with overlap enabled should drain the result_queue.""" + scheduler = self._new_scheduler() + scheduler.enable_overlap = True + mock_batch = MagicMock() + mock_batch.forward_mode.is_extend.return_value = False + scheduler.last_batch = mock_batch + scheduler.result_queue = deque([(MagicMock(), MagicMock())]) + scheduler.process_batch_result = MagicMock() + + scheduler.pause_generation(PauseGenerationReqInput(mode="retract")) + + scheduler.process_batch_result.assert_called_once() + self.assertEqual(len(scheduler.result_queue), 0) def test_pd_decode_retract_requeues_for_rebootstrap(self): """PD decode retract should rebootstrap instead of resuming stale CPU KV.""" scheduler = self._new_scheduler() scheduler.disaggregation_mode = DisaggregationMode.DECODE scheduler.last_batch = None - scheduler.running_batch.reqs = [MagicMock()] scheduler.running_batch.is_empty.return_value = False scheduler._add_request_to_queue = MagicMock() scheduler.disagg_decode_prealloc_queue = MagicMock() @@ -147,11 +193,12 @@ class TestSchedulerPauseGeneration(unittest.TestCase): output_ids=[10, 11, 12], time_stats=MagicMock(), ) - scheduler.running_batch.retract_all.return_value = [req] + scheduler.running_batch.reqs = [req] scheduler.running_batch.filter_batch = MagicMock() scheduler.server_args = MagicMock() - scheduler.pause_generation(PauseGenerationReqInput(mode="retract")) + with patch("sglang.srt.managers.scheduler.retract_all") as mock_retract_all: + scheduler.pause_generation(PauseGenerationReqInput(mode="retract")) scheduler._add_request_to_queue.assert_not_called() scheduler.disagg_decode_prealloc_queue.hold_rebootstrap.assert_called_once_with( @@ -162,9 +209,8 @@ class TestSchedulerPauseGeneration(unittest.TestCase): self.assertTrue(req.pd_rebootstrap_in_progress) # Rebootstrap recomputes the KV from the prefill, so the retract must skip # the device->host KV offload rather than offload-then-delete it. - scheduler.running_batch.retract_all.assert_called_once_with( - scheduler.server_args, offload_kv=False - ) + mock_retract_all.assert_called_once() + self.assertEqual(mock_retract_all.call_args.kwargs["offload_kv"], False) def test_pd_decode_continue_releases_held_rebootstrap(self): """continue_generation must enqueue staged rebootstrap reqs on resume.""" @@ -180,21 +226,6 @@ class TestSchedulerPauseGeneration(unittest.TestCase): scheduler.disagg_decode_prealloc_queue.enqueue_held_rebootstrap.assert_called_once_with() self.assertFalse(scheduler._engine_paused) - def test_abort_drains_overlap_queue(self): - """abort with overlap enabled should drain the result_queue.""" - scheduler = self._new_scheduler() - scheduler.enable_overlap = True - mock_batch = MagicMock() - mock_batch.forward_mode.is_extend.return_value = False - scheduler.last_batch = mock_batch - scheduler.result_queue = deque([(MagicMock(), MagicMock())]) - scheduler.process_batch_result = MagicMock() - - scheduler.pause_generation(PauseGenerationReqInput(mode="abort")) - - scheduler.process_batch_result.assert_called_once() - self.assertEqual(len(scheduler.result_queue), 0) - if __name__ == "__main__": unittest.main()