From c7004bb7f513a1ff3c8f0a70079975b0772e8351 Mon Sep 17 00:00:00 2001 From: aoshen02 Date: Mon, 28 Sep 2026 15:31:03 +0000 Subject: [PATCH 1/4] [Bugfix][Core] Allow AuxOutput reset after abort without running requests Co-authored-by: OpenAI Codex Signed-off-by: aoshen02 --- vllm/v1/core/sched/scheduler.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/vllm/v1/core/sched/scheduler.py b/vllm/v1/core/sched/scheduler.py index 8ebedef4ec31..6926cce52ede 100644 --- a/vllm/v1/core/sched/scheduler.py +++ b/vllm/v1/core/sched/scheduler.py @@ -2771,7 +2771,7 @@ def reset_prefix_cache( is no running requests taking KV cache. """ if reset_running_requests and self.aux_output_connector is not None: - if self._pause_state != PauseState.PAUSED_ALL: + if self.running and self._pause_state != PauseState.PAUSED_ALL: raise RuntimeError( "AuxOutput Connector only supports resetting running requests " "after pause(mode='keep')." From 75052844e505f362e776849a69528bfcc1a8595f Mon Sep 17 00:00:00 2001 From: aoshen02 Date: Tue, 29 Sep 2026 03:56:08 +0000 Subject: [PATCH 2/4] Allow AuxOutput reset independently of pause mode Keep the in-flight output guard and test reset without a keep-mode pause. Co-authored-by: Codex Signed-off-by: aoshen02 --- tests/v1/core/test_scheduler.py | 6 +----- vllm/v1/core/sched/scheduler.py | 18 ++++++++---------- 2 files changed, 9 insertions(+), 15 deletions(-) diff --git a/tests/v1/core/test_scheduler.py b/tests/v1/core/test_scheduler.py index 91c93ed68fda..6a8a9009f32f 100644 --- a/tests/v1/core/test_scheduler.py +++ b/tests/v1/core/test_scheduler.py @@ -1562,11 +1562,7 @@ def test_scheduler_reset_prefix_cache(): assert not scheduler.reset_prefix_cache() scheduler.aux_output_connector.reset.assert_not_called() - with pytest.raises(RuntimeError, match=r"pause\(mode='keep'\)"): - scheduler.reset_prefix_cache(reset_running_requests=True) - - # pause(mode="keep") also waits for scheduled model outputs to drain. - scheduler.set_pause_state(PauseState.PAUSED_ALL) + # Reset requires drained model outputs, not a particular pause mode. with pytest.raises(RuntimeError, match="model output is in flight"): scheduler.reset_prefix_cache(reset_running_requests=True) for request in requests: diff --git a/vllm/v1/core/sched/scheduler.py b/vllm/v1/core/sched/scheduler.py index 86aca8d69e7c..3a9fdf219d29 100644 --- a/vllm/v1/core/sched/scheduler.py +++ b/vllm/v1/core/sched/scheduler.py @@ -2781,16 +2781,14 @@ def reset_prefix_cache( Otherwise, this method will only reset the KV prefix cache when there is no running requests taking KV cache. """ - if reset_running_requests and self.aux_output_connector is not None: - if self.running and self._pause_state != PauseState.PAUSED_ALL: - raise RuntimeError( - "AuxOutput Connector only supports resetting running requests " - "after pause(mode='keep')." - ) - if any(request.num_in_flight_tokens for request in self.requests.values()): - raise RuntimeError( - "AuxOutput Connector cannot reset while model output is in flight." - ) + if ( + reset_running_requests + and self.aux_output_connector is not None + and any(request.num_in_flight_tokens for request in self.requests.values()) + ): + raise RuntimeError( + "AuxOutput Connector cannot reset while model output is in flight." + ) if reset_running_requests: # For logging. timestamp = time.monotonic() From b66a2a07163398c199b82a0d3fd183342f2fe5e1 Mon Sep 17 00:00:00 2001 From: aoshen02 Date: Tue, 29 Sep 2026 07:32:21 +0000 Subject: [PATCH 3/4] Drain queued model outputs before resetting prefix cache Remove the AuxOutput-specific in-flight guard and consume queued outputs at the EngineCoreProc reset boundary without scheduling new batches. Co-authored-by: Codex Signed-off-by: aoshen02 --- tests/v1/core/test_scheduler.py | 4 +-- tests/v1/engine/test_engine_core.py | 46 +++++++++++++++++++++++++++++ vllm/v1/core/sched/scheduler.py | 8 ----- vllm/v1/engine/core.py | 14 +++++++-- 4 files changed, 59 insertions(+), 13 deletions(-) diff --git a/tests/v1/core/test_scheduler.py b/tests/v1/core/test_scheduler.py index 6a8a9009f32f..01c9b3ed07a2 100644 --- a/tests/v1/core/test_scheduler.py +++ b/tests/v1/core/test_scheduler.py @@ -1562,9 +1562,7 @@ def test_scheduler_reset_prefix_cache(): assert not scheduler.reset_prefix_cache() scheduler.aux_output_connector.reset.assert_not_called() - # Reset requires drained model outputs, not a particular pause mode. - with pytest.raises(RuntimeError, match="model output is in flight"): - scheduler.reset_prefix_cache(reset_running_requests=True) + # EngineCore consumes pending model outputs before resetting the scheduler. for request in requests: request.num_in_flight_tokens = 0 diff --git a/tests/v1/engine/test_engine_core.py b/tests/v1/engine/test_engine_core.py index 2d9f7e8632d7..56910b430088 100644 --- a/tests/v1/engine/test_engine_core.py +++ b/tests/v1/engine/test_engine_core.py @@ -6,6 +6,7 @@ import uuid from collections import deque from concurrent.futures import Future, ThreadPoolExecutor +from typing import Any from unittest.mock import MagicMock, PropertyMock, patch import pytest @@ -701,6 +702,51 @@ def _pausable_engine_core_proc() -> EngineCoreProc: return core +@pytest.mark.parametrize("num_batches", [0, 1, 2]) +@pytest.mark.parametrize("reset_running_requests", [False, True]) +def test_reset_prefix_cache_consumes_outputs_before_reset( + num_batches, reset_running_requests +): + """Reset drains queued results in order without scheduling or losing outputs.""" + core = _pausable_engine_core_proc() + core.batch_queue = deque() + core.batch_queue_size = 3 + core.output_queue = MagicMock() + core.log_error_detail = MagicMock() + core.capture_iteration_details = MagicMock() + core._attach_iteration_details = MagicMock() + core._process_aborts_queue = MagicMock() + core.scheduler.has_requests.return_value = True + consumed: list[int] = [] + for index in range(num_batches): + future: Future[Any] = Future() + future.set_result(index) + core.batch_queue.appendleft((future, index, future)) + + def consume(scheduled, output): + assert scheduled == output == len(consumed) + consumed.append(output) + return {0: output} + + def reset(*args): + assert len(consumed) == num_batches + assert not core.batch_queue + assert core.output_queue.put_nowait.call_count == num_batches + return True + + core.scheduler.update_from_output.side_effect = consume + core.scheduler.reset_prefix_cache.side_effect = reset + assert core.reset_prefix_cache(reset_running_requests, reset_connector=True) + core.scheduler.reset_prefix_cache.assert_called_once_with( + reset_running_requests, True + ) + core.scheduler.schedule.assert_not_called() + core.model_executor.execute_model.assert_not_called() + assert [call.args[0] for call in core.output_queue.put_nowait.call_args_list] == [ + (0, index) for index in range(num_batches) + ] + + @pytest.mark.parametrize( "pause_state,has_requests,has_batches", [ diff --git a/vllm/v1/core/sched/scheduler.py b/vllm/v1/core/sched/scheduler.py index 3a9fdf219d29..24ecc4534d6d 100644 --- a/vllm/v1/core/sched/scheduler.py +++ b/vllm/v1/core/sched/scheduler.py @@ -2781,14 +2781,6 @@ def reset_prefix_cache( Otherwise, this method will only reset the KV prefix cache when there is no running requests taking KV cache. """ - if ( - reset_running_requests - and self.aux_output_connector is not None - and any(request.num_in_flight_tokens for request in self.requests.values()) - ): - raise RuntimeError( - "AuxOutput Connector cannot reset while model output is in flight." - ) if reset_running_requests: # For logging. timestamp = time.monotonic() diff --git a/vllm/v1/engine/core.py b/vllm/v1/engine/core.py index 5c09a0059e02..ce315ebaebf8 100644 --- a/vllm/v1/engine/core.py +++ b/vllm/v1/engine/core.py @@ -668,7 +668,7 @@ def post_step(self, model_executed: bool) -> None: self.scheduler.update_draft_token_ids(draft_token_ids) def step_with_batch_queue( - self, + self, *, schedule_new_batch: bool = True ) -> tuple[dict[int, EngineCoreOutputs] | None, bool]: """Schedule and execute batches with the batch queue. Note that if nothing to output in this step, None is returned. @@ -693,7 +693,7 @@ def step_with_batch_queue( model_executed = False deferred_scheduler_output = None - if self.scheduler.has_requests(): + if schedule_new_batch and self.scheduler.has_requests(): scheduler_output = self.scheduler.schedule(self._should_throttle_prefills()) with self.log_error_detail(scheduler_output): exec_future = self.model_executor.execute_model( @@ -1541,6 +1541,16 @@ def _process_engine_step(self) -> bool: return model_executed + def reset_prefix_cache( + self, reset_running_requests: bool = False, reset_connector: bool = False + ) -> bool: + # Consume old outputs before resetting request state or cache generations. + while self.batch_queue: + outputs, _ = self.step_with_batch_queue(schedule_new_batch=False) + for output in outputs.items() if outputs else (): + self.output_queue.put_nowait(output) + return super().reset_prefix_cache(reset_running_requests, reset_connector) + def _notify_idle_state_callbacks(self) -> None: while self._idle_state_callbacks: callback = self._idle_state_callbacks.pop() From 279db25ac45b1d006013f9a7090d3acf1925e639 Mon Sep 17 00:00:00 2001 From: aoshen02 Date: Tue, 29 Sep 2026 07:46:13 +0000 Subject: [PATCH 4/4] Keep AuxOutput reset scoped to completed pauses Remove the added EngineCore queue draining and its tests. Keep removal of the AuxOutput-specific reset guards, with callers responsible for awaiting pause completion before reset. Co-authored-by: Codex Signed-off-by: aoshen02 --- tests/v1/core/test_scheduler.py | 2 +- tests/v1/engine/test_engine_core.py | 46 ----------------------------- vllm/v1/engine/core.py | 14 ++------- 3 files changed, 3 insertions(+), 59 deletions(-) diff --git a/tests/v1/core/test_scheduler.py b/tests/v1/core/test_scheduler.py index 01c9b3ed07a2..d439a197e926 100644 --- a/tests/v1/core/test_scheduler.py +++ b/tests/v1/core/test_scheduler.py @@ -1562,7 +1562,7 @@ def test_scheduler_reset_prefix_cache(): assert not scheduler.reset_prefix_cache() scheduler.aux_output_connector.reset.assert_not_called() - # EngineCore consumes pending model outputs before resetting the scheduler. + # Pause completes pending model outputs before the caller resets the scheduler. for request in requests: request.num_in_flight_tokens = 0 diff --git a/tests/v1/engine/test_engine_core.py b/tests/v1/engine/test_engine_core.py index 56910b430088..2d9f7e8632d7 100644 --- a/tests/v1/engine/test_engine_core.py +++ b/tests/v1/engine/test_engine_core.py @@ -6,7 +6,6 @@ import uuid from collections import deque from concurrent.futures import Future, ThreadPoolExecutor -from typing import Any from unittest.mock import MagicMock, PropertyMock, patch import pytest @@ -702,51 +701,6 @@ def _pausable_engine_core_proc() -> EngineCoreProc: return core -@pytest.mark.parametrize("num_batches", [0, 1, 2]) -@pytest.mark.parametrize("reset_running_requests", [False, True]) -def test_reset_prefix_cache_consumes_outputs_before_reset( - num_batches, reset_running_requests -): - """Reset drains queued results in order without scheduling or losing outputs.""" - core = _pausable_engine_core_proc() - core.batch_queue = deque() - core.batch_queue_size = 3 - core.output_queue = MagicMock() - core.log_error_detail = MagicMock() - core.capture_iteration_details = MagicMock() - core._attach_iteration_details = MagicMock() - core._process_aborts_queue = MagicMock() - core.scheduler.has_requests.return_value = True - consumed: list[int] = [] - for index in range(num_batches): - future: Future[Any] = Future() - future.set_result(index) - core.batch_queue.appendleft((future, index, future)) - - def consume(scheduled, output): - assert scheduled == output == len(consumed) - consumed.append(output) - return {0: output} - - def reset(*args): - assert len(consumed) == num_batches - assert not core.batch_queue - assert core.output_queue.put_nowait.call_count == num_batches - return True - - core.scheduler.update_from_output.side_effect = consume - core.scheduler.reset_prefix_cache.side_effect = reset - assert core.reset_prefix_cache(reset_running_requests, reset_connector=True) - core.scheduler.reset_prefix_cache.assert_called_once_with( - reset_running_requests, True - ) - core.scheduler.schedule.assert_not_called() - core.model_executor.execute_model.assert_not_called() - assert [call.args[0] for call in core.output_queue.put_nowait.call_args_list] == [ - (0, index) for index in range(num_batches) - ] - - @pytest.mark.parametrize( "pause_state,has_requests,has_batches", [ diff --git a/vllm/v1/engine/core.py b/vllm/v1/engine/core.py index ce315ebaebf8..5c09a0059e02 100644 --- a/vllm/v1/engine/core.py +++ b/vllm/v1/engine/core.py @@ -668,7 +668,7 @@ def post_step(self, model_executed: bool) -> None: self.scheduler.update_draft_token_ids(draft_token_ids) def step_with_batch_queue( - self, *, schedule_new_batch: bool = True + self, ) -> tuple[dict[int, EngineCoreOutputs] | None, bool]: """Schedule and execute batches with the batch queue. Note that if nothing to output in this step, None is returned. @@ -693,7 +693,7 @@ def step_with_batch_queue( model_executed = False deferred_scheduler_output = None - if schedule_new_batch and self.scheduler.has_requests(): + if self.scheduler.has_requests(): scheduler_output = self.scheduler.schedule(self._should_throttle_prefills()) with self.log_error_detail(scheduler_output): exec_future = self.model_executor.execute_model( @@ -1541,16 +1541,6 @@ def _process_engine_step(self) -> bool: return model_executed - def reset_prefix_cache( - self, reset_running_requests: bool = False, reset_connector: bool = False - ) -> bool: - # Consume old outputs before resetting request state or cache generations. - while self.batch_queue: - outputs, _ = self.step_with_batch_queue(schedule_new_batch=False) - for output in outputs.items() if outputs else (): - self.output_queue.put_nowait(output) - return super().reset_prefix_cache(reset_running_requests, reset_connector) - def _notify_idle_state_callbacks(self) -> None: while self._idle_state_callbacks: callback = self._idle_state_callbacks.pop()