From 756f01ad7e88b164b2e9c37cea8a5149454e7d6e Mon Sep 17 00:00:00 2001 From: Sohom Chakraborty <16609933+sohom-cs@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:36:22 -0400 Subject: [PATCH] [Bugfix][KV Offload] Keep stepping until a reset_cache flush is delivered, and make reset_cache re-entrant reset_cache() moves every in-flight job id into _current_batch_jobs_to_flush and deliberately leaves it set: workers must wait on those jobs before a post-reset store reuses their CPU chunks. Two gaps: - has_pending_push_work() only checked self._jobs, which reset_cache() just cleared. With no requests left, the engine stopped stepping and the flush set was not delivered until an unrelated request arrived. - A second reset_cache() before the next step hit assert not self._current_batch_jobs_to_flush and took the engine core down. RL loops call reset_prefix_cache(reset_connector=True) every iteration, and EngineCore drains all queued utility calls before it steps again. Count a pending flush set as push work, and merge a second reset's ids into the pending set instead of asserting. Co-authored-by: Claude Signed-off-by: Sohom Chakraborty <16609933+sohom-cs@users.noreply.github.com> --- .../offloading_connector/test_scheduler.py | 36 +++++++++++++++++++ .../kv_connector/v1/offloading/scheduler.py | 13 +++++-- 2 files changed, 46 insertions(+), 3 deletions(-) diff --git a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py index 5b3bf1c8bebe..9d1a7f72d0c9 100644 --- a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py +++ b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py @@ -2478,6 +2478,42 @@ def test_reset_cache(request_runner, async_scheduling: bool): assert group_state.next_stored_chunk_idx == 0 +def test_reset_cache_flush_is_delivered_when_idle_and_reset_is_reentrant( + request_runner, +): + """A reset that discards an in-flight store after the last request + finished must keep the engine stepping until its flush set reaches the + workers, and a second reset before that step must not assert (RL loops + call reset_prefix_cache(reset_connector=True) every iteration, and + EngineCore drains all queued utility calls before it steps again).""" + runner = request_runner( + block_size=4, num_gpu_blocks=100, async_scheduling=False, blocks_per_chunk=1 + ) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + runner.manager.has_pending_work.return_value = False + runner.new_request(token_ids=[0] * 8) + runner.run(decoded_tokens=[EOS_TOKEN_ID], complete_transfers=False) + cs = runner.connector_scheduler + store_job_ids = set(cs._jobs) + assert store_job_ids + assert not runner.scheduler.has_unfinished_requests() + + cs.reset_cache() + cs.reset_cache() + + assert cs._current_batch_jobs_to_flush == store_job_ids + assert cs.has_pending_push_work() + assert runner.scheduler.has_requests() + + scheduler_output = runner.scheduler.schedule() + meta = scheduler_output.kv_connector_metadata + assert isinstance(meta, OffloadingConnectorMetadata) + assert meta.jobs_to_flush == store_job_ids + assert not cs.has_pending_push_work() + + @pytest.mark.parametrize("async_scheduling", [True, False]) def test_reset_cache_finalizes_finished_request_with_pending_store( request_runner, async_scheduling: bool diff --git a/vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py b/vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py index 47aa8654feb8..024330c8e18b 100644 --- a/vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py +++ b/vllm/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py @@ -1835,8 +1835,14 @@ def has_pending_push_work(self) -> bool: While True, build_connector_meta() and update_connector_output() continue to be called even when no requests are scheduled. + A flush set left by reset_cache() counts: it reaches the workers only + through build_connector_meta(). """ - return bool(self._jobs) or self.manager.has_pending_work() + return ( + bool(self._jobs) + or bool(self._current_batch_jobs_to_flush) + or self.manager.has_pending_work() + ) def update_connector_output(self, connector_output: KVConnectorOutput): """Update KVConnector state from worker-side connectors output. @@ -1990,9 +1996,10 @@ def take_events(self) -> Iterable[KVCacheEvent]: def reset_cache(self) -> None: """Reset the offloading manager cache, evicting all stored chunks.""" - # reset_cache cannot be called in the middle of a schedule step + # reset_cache cannot be called in the middle of a schedule step. + # _current_batch_jobs_to_flush may still hold a previous reset's flush + # set if no step ran since; the new ids are merged into it. assert not self._current_batch_load_jobs - assert not self._current_batch_jobs_to_flush assert not self._current_batch_allocated_block_ids # Flush all in-flight jobs