From 2fa932f54668e114926b4c84368e32c4f16b5216 Mon Sep 17 00:00:00 2001 From: Alex Date: Mon, 15 Jun 2026 17:41:15 +0800 Subject: [PATCH 1/4] test(kv_connector): add request_finished fence population tests Add unit tests to verify that request_finished correctly populates the pending store fence index when in-flight store jobs exist. This prevents data corruption by tracking pending non-sliding-window blocks until offloading completes. Also verifies proper cleanup when no pending stores remain. Signed-off-by: Alex --- .../offloading_connector/test_scheduler.py | 371 +++++++++++++++++- 1 file changed, 367 insertions(+), 4 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 973fcc63e31d..2ea9e1be7196 100644 --- a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py +++ b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py @@ -950,8 +950,8 @@ def setup(r, max_offload_tokens): token_ids=[0] * offloaded_block_size * 3, kv_transfer_params={"max_offload_tokens": max_offload_tokens}, ) - r.manager.prepare_store.side_effect = ( - lambda keys, req_context: generate_store_output(keys) + r.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) ) # Pending offloads drain via non-blocking stepping, not a flush, so no @@ -1049,8 +1049,8 @@ def test_offload_prompt_only(request_runner, async_scheduling: bool): extra_config_overrides={"offload_prompt_only": True}, ) - runner.manager.prepare_store.side_effect = ( - lambda keys, req_context: generate_store_output(keys) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) ) runner.new_request(token_ids=[0] * offloaded_block_size * num_prompt_blocks) @@ -2053,3 +2053,366 @@ def test_full_attn_store_then_load(self, request_runner, async_scheduling: bool) (1, 1), ), ) + + +# --------------------------------------------------------------------------- +# Tests for request_finished fence population with in-flight pending stores. +# --------------------------------------------------------------------------- + + +def test_request_finished_with_pending_stores_populates_fence(request_runner): + """When a request finishes with in-flight store jobs, the fence index + (_block_id_to_pending_jobs) is correctly populated with the store jobs' + non_sliding_window_block_ids. + + This prevents data corruption when a subsequent request reuses the same + GPU blocks before the store completes. + """ + block_size = 4 + block_size_factor = 1 + offloaded_block_size = block_size * block_size_factor + + runner = request_runner( + block_size=block_size, + num_gpu_blocks=100, + async_scheduling=False, + block_size_factor=block_size_factor, + ) + + # Create a request with 3 blocks of data. + runner.new_request(token_ids=[0] * offloaded_block_size * 3) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + + # Schedule to create store jobs (but don't complete transfers). + runner.scheduler.schedule() + runner._update_gpu_blocks() + + # Verify store job was created. + assert len(runner.connector_scheduler._jobs) > 0 + job_id = next(iter(runner.connector_scheduler._jobs)) + job_status = runner.connector_scheduler._jobs[job_id] + assert job_status.is_store + non_sw_block_ids = job_status.non_sliding_window_block_ids or [] + assert len(non_sw_block_ids) > 0 + + # Fence should be empty before request_finished + # (non-sliding-window blocks are only registered at request finish). + assert runner.connector_scheduler._block_id_to_pending_jobs == {} + + # Finish the request — triggers request_finished which populates the fence. + req_id = str(runner.req_id) + runner.scheduler.finish_requests({req_id}, RequestStatus.FINISHED_STOPPED) + + # Verify fence index is populated with the job's non-SW block IDs. + for bid in non_sw_block_ids: + assert bid in runner.connector_scheduler._block_id_to_pending_jobs + assert job_id in runner.connector_scheduler._block_id_to_pending_jobs[bid] + + # req_status should still exist because the job is still in-flight. + assert req_id in runner.connector_scheduler._req_status + + +def test_request_finished_no_pending_stores_cleanup(request_runner): + """When all store jobs complete before the request finishes, + request_finished cleans up req_status (no fence entries needed).""" + block_size = 4 + block_size_factor = 1 + offloaded_block_size = block_size * block_size_factor + + runner = request_runner( + block_size=block_size, + num_gpu_blocks=100, + async_scheduling=False, + block_size_factor=block_size_factor, + ) + + runner.new_request(token_ids=[0] * offloaded_block_size * 3) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + + # Run with complete_transfers=True — all stores finish before request ends. + runner.run( + decoded_tokens=[EOS_TOKEN_ID], + expected_stored=(0, 1, 2), + expected_flushed=(0, 1, 2), + ) + + # All jobs completed → fence index empty. + assert runner.connector_scheduler._block_id_to_pending_jobs == {} + assert len(runner.connector_scheduler._jobs) == 0 + + # req_status should already be cleaned up: + # update_connector_output deletes it when the last job completes + # for a finished request. + req_id = str(runner.req_id) + assert req_id not in runner.connector_scheduler._req_status + + +def test_request_finished_sliding_window_blocks_not_in_fence(request_runner): + """Sliding window blocks are registered in the fence at store creation + time, NOT at request_finished. request_finished only registers + non_sliding_window_block_ids. + + For a pure sliding window group, non_sliding_window_block_ids is empty, + so request_finished adds zero new fence entries. + """ + block_size = 4 + sliding_window = 8 # 2 blocks + + kv_cache_groups = [ + KVCacheGroupSpec( + ["layer0"], + SlidingWindowSpec( + block_size=block_size, + num_kv_heads=1, + head_size=1, + dtype=torch.float32, + sliding_window=sliding_window, + ), + ), + ] + + runner = request_runner( + block_size=block_size, + num_gpu_blocks=100, + async_scheduling=False, + kv_cache_groups=kv_cache_groups, + ) + + # 3 blocks of prompt (12 tokens), window = 2 blocks. + runner.new_request(token_ids=[0] * block_size * 3) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + + # Schedule to create store jobs. + runner.scheduler.schedule() + runner._update_gpu_blocks() + + # Verify store job was created. + assert len(runner.connector_scheduler._jobs) > 0 + job_id = next(iter(runner.connector_scheduler._jobs)) + job_status = runner.connector_scheduler._jobs[job_id] + assert job_status.is_store + + # For sliding window groups, all block IDs go into + # sliding_window_block_ids; non_sliding_window_block_ids is empty. + sw_block_ids = set(job_status.sliding_window_block_ids or []) + non_sw_block_ids = set(job_status.non_sliding_window_block_ids or []) + assert len(sw_block_ids) > 0 + assert len(non_sw_block_ids) == 0 + + # Sliding window blocks should already be in the fence + # (registered at store creation time in _build_store_jobs). + for bid in sw_block_ids: + assert bid in runner.connector_scheduler._block_id_to_pending_jobs + assert job_id in runner.connector_scheduler._block_id_to_pending_jobs[bid] + + fence_before = dict(runner.connector_scheduler._block_id_to_pending_jobs) + + # Finish the request — request_finished should NOT add new entries + # because non_sliding_window_block_ids is empty for SW groups. + req_id = str(runner.req_id) + runner.scheduler.finish_requests({req_id}, RequestStatus.FINISHED_STOPPED) + + fence_after = runner.connector_scheduler._block_id_to_pending_jobs + + # Fence unchanged: no new entries from request_finished. + assert fence_before == fence_after + + +def test_multiple_in_flight_stores_all_registered_in_fence(request_runner): + """When a request finishes with multiple in-flight store jobs, + ALL jobs' non_sliding_window_block_ids are registered in the fence. + """ + block_size = 4 + block_size_factor = 1 + offloaded_block_size = block_size * block_size_factor + + runner = request_runner( + block_size=block_size, + num_gpu_blocks=100, + async_scheduling=False, + block_size_factor=block_size_factor, + ) + + # Create a request with 1 block. + runner.new_request(token_ids=[0] * offloaded_block_size) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + + # Schedule to create the first store job. + runner.scheduler.schedule() + runner._update_gpu_blocks() + + assert len(runner.connector_scheduler._jobs) == 1 + job_id_1 = next(iter(runner.connector_scheduler._jobs)) + job_status_1 = runner.connector_scheduler._jobs[job_id_1] + non_sw_1 = list(job_status_1.non_sliding_window_block_ids or []) + assert len(non_sw_1) > 0 + + # Manually inject a second in-flight store job to simulate a request + # that produced multiple store jobs across decode steps. + req_id = str(runner.req_id) + req_status = runner.connector_scheduler._req_status[req_id] + fake_block_id = 99 + fake_job_id = 999 + runner.connector_scheduler._jobs[fake_job_id] = TransferJobStatus( + req_id=req_id, + pending_count=runner.connector_scheduler.config.num_workers, + keys=set(), + is_store=True, + non_sliding_window_block_ids=[fake_block_id], + ) + req_status.transfer_jobs.add(fake_job_id) + + assert len(req_status.transfer_jobs) == 2 + + # Finish the request — both jobs should be registered in the fence. + runner.scheduler.finish_requests({req_id}, RequestStatus.FINISHED_STOPPED) + + fence = runner.connector_scheduler._block_id_to_pending_jobs + for bid in non_sw_1: + assert bid in fence + assert job_id_1 in fence[bid] + assert fake_block_id in fence + assert fake_job_id in fence[fake_block_id] + + +def test_fence_full_lifecycle_populate_and_cleanup(request_runner): + """Verifies the complete fence lifecycle using the standard runner flow: + 1. Store job created (in-flight) + 2. Request finishes → request_finished populates the fence + 3. build_connector_meta detects all requests finished → flushes jobs + 4. Flush completes the job + 5. update_connector_output cleans up the fence and req_status + + The fence population is verified by + test_request_finished_with_pending_stores_populates_fence. + This test verifies the cleanup after the full lifecycle. + """ + block_size = 4 + block_size_factor = 1 + offloaded_block_size = block_size * block_size_factor + + runner = request_runner( + block_size=block_size, + num_gpu_blocks=100, + async_scheduling=False, + block_size_factor=block_size_factor, + ) + + runner.new_request(token_ids=[0] * offloaded_block_size) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + + # Run with complete_transfers=False: creates store job, finishes request, + # populates fence, flushes job (which completes it), cleans up fence. + runner.run( + decoded_tokens=[EOS_TOKEN_ID], + complete_transfers=False, + expected_stored=(0,), + expected_flushed=(0,), + ) + + # After the full lifecycle, fence should be empty. + assert runner.connector_scheduler._block_id_to_pending_jobs == {} + # Job should be removed. + assert len(runner.connector_scheduler._jobs) == 0 + # req_status should be removed. + req_id = str(runner.req_id) + assert req_id not in runner.connector_scheduler._req_status + + +def test_request_finished_mixed_full_attn_and_sliding_window( + request_runner, +): + """With both FullAttention and SlidingWindow groups, a single store job + has both non_sliding_window_block_ids and sliding_window_block_ids. + + request_finished only registers non-SW blocks in the fence. + SW blocks were already registered at store creation time. + """ + block_size = 4 + sliding_window = 8 # 2 blocks + + kv_cache_groups = [ + KVCacheGroupSpec( + ["layer0"], + FullAttentionSpec( + block_size=block_size, + num_kv_heads=1, + head_size=1, + dtype=torch.float32, + ), + ), + KVCacheGroupSpec( + ["layer1"], + SlidingWindowSpec( + block_size=block_size, + num_kv_heads=1, + head_size=1, + dtype=torch.float32, + sliding_window=sliding_window, + ), + ), + ] + + runner = request_runner( + block_size=block_size, + num_gpu_blocks=100, + async_scheduling=False, + kv_cache_groups=kv_cache_groups, + ) + + # 3 blocks of prompt (12 tokens) — enough for both groups. + runner.new_request(token_ids=[0] * block_size * 3) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + + runner.scheduler.schedule() + runner._update_gpu_blocks() + + assert len(runner.connector_scheduler._jobs) > 0 + job_id = next(iter(runner.connector_scheduler._jobs)) + job_status = runner.connector_scheduler._jobs[job_id] + assert job_status.is_store + + # Both types of block IDs should be present. + non_sw_block_ids = set(job_status.non_sliding_window_block_ids or []) + sw_block_ids = set(job_status.sliding_window_block_ids or []) + assert len(non_sw_block_ids) > 0 + assert len(sw_block_ids) > 0 + + # SW blocks should already be in the fence + # (registered at store creation time). + for bid in sw_block_ids: + assert bid in runner.connector_scheduler._block_id_to_pending_jobs + assert job_id in (runner.connector_scheduler._block_id_to_pending_jobs[bid]) + + # Snapshot the SW fence entries before request_finished. + sw_fence_entries = { + bid: set(runner.connector_scheduler._block_id_to_pending_jobs[bid]) + for bid in sw_block_ids + } + + # Finish the request — only non-SW blocks should be newly registered. + req_id = str(runner.req_id) + runner.scheduler.finish_requests({req_id}, RequestStatus.FINISHED_STOPPED) + + fence = runner.connector_scheduler._block_id_to_pending_jobs + + # Non-SW blocks are now in the fence. + for bid in non_sw_block_ids: + assert bid in fence + assert job_id in fence[bid] + + # SW fence entries are unchanged (not duplicated). + for bid, expected_jobs in sw_fence_entries.items(): + assert fence[bid] == expected_jobs From 9fffb4b879953d41268fdb5b93fdd9ec8dc75101 Mon Sep 17 00:00:00 2001 From: Alex Date: Tue, 16 Jun 2026 10:49:51 +0800 Subject: [PATCH 2/4] test(kv_connector): improve fence test quality with observable behavior assertions - Remove redundant tests (no_pending_stores_cleanup, sliding_window_blocks_not_in_fence) - Rewrite test_multiple_in_flight_stores to use natural multi-job scenario instead of fake job injection, asserting expected_flushed over internal state - Remove unused TransferJobStatus import Signed-off-by: AlexHuang Signed-off-by: Alex --- .../offloading_connector/test_scheduler.py | 173 +++--------------- 1 file changed, 30 insertions(+), 143 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 2ea9e1be7196..c674fe0c63d7 100644 --- a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py +++ b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py @@ -2114,9 +2114,15 @@ def test_request_finished_with_pending_stores_populates_fence(request_runner): assert req_id in runner.connector_scheduler._req_status -def test_request_finished_no_pending_stores_cleanup(request_runner): - """When all store jobs complete before the request finishes, - request_finished cleans up req_status (no fence entries needed).""" +def test_multiple_in_flight_stores_all_flushed(request_runner): + """When a request finishes with multiple in-flight store jobs, + ALL jobs are flushed via the fence mechanism (observable behavior). + + Uses two runner.run() calls to naturally produce two store jobs: + - First run: 4 decoded tokens fills blocks 0-1 → job_0 created for both + - Second run: 4 more tokens + EOS fills block 2 → job_1 created, + request finishes, "all finished" flush fires → both jobs flushed + """ block_size = 4 block_size_factor = 1 offloaded_block_size = block_size * block_size_factor @@ -2128,159 +2134,40 @@ def test_request_finished_no_pending_stores_cleanup(request_runner): block_size_factor=block_size_factor, ) - runner.new_request(token_ids=[0] * offloaded_block_size * 3) + # 4 prompt tokens → 1 GPU block (block 0) + runner.new_request(token_ids=[0] * offloaded_block_size) runner.manager.prepare_store.side_effect = lambda keys, req_context: ( generate_store_output(keys) ) - # Run with complete_transfers=True — all stores finish before request ends. + # First run: 4 decoded tokens → blocks 0-1 full → job_0 created for both. + # complete_transfers=False keeps job_0 in-flight (not completed yet). runner.run( - decoded_tokens=[EOS_TOKEN_ID], - expected_stored=(0, 1, 2), - expected_flushed=(0, 1, 2), - ) - - # All jobs completed → fence index empty. - assert runner.connector_scheduler._block_id_to_pending_jobs == {} - assert len(runner.connector_scheduler._jobs) == 0 - - # req_status should already be cleaned up: - # update_connector_output deletes it when the last job completes - # for a finished request. - req_id = str(runner.req_id) - assert req_id not in runner.connector_scheduler._req_status - - -def test_request_finished_sliding_window_blocks_not_in_fence(request_runner): - """Sliding window blocks are registered in the fence at store creation - time, NOT at request_finished. request_finished only registers - non_sliding_window_block_ids. - - For a pure sliding window group, non_sliding_window_block_ids is empty, - so request_finished adds zero new fence entries. - """ - block_size = 4 - sliding_window = 8 # 2 blocks - - kv_cache_groups = [ - KVCacheGroupSpec( - ["layer0"], - SlidingWindowSpec( - block_size=block_size, - num_kv_heads=1, - head_size=1, - dtype=torch.float32, - sliding_window=sliding_window, - ), - ), - ] - - runner = request_runner( - block_size=block_size, - num_gpu_blocks=100, - async_scheduling=False, - kv_cache_groups=kv_cache_groups, - ) - - # 3 blocks of prompt (12 tokens), window = 2 blocks. - runner.new_request(token_ids=[0] * block_size * 3) - runner.manager.prepare_store.side_effect = lambda keys, req_context: ( - generate_store_output(keys) + decoded_tokens=[0] * offloaded_block_size, + complete_transfers=False, + expected_stored=(), + expected_flushed=(), ) - # Schedule to create store jobs. - runner.scheduler.schedule() - runner._update_gpu_blocks() - - # Verify store job was created. - assert len(runner.connector_scheduler._jobs) > 0 - job_id = next(iter(runner.connector_scheduler._jobs)) - job_status = runner.connector_scheduler._jobs[job_id] - assert job_status.is_store - - # For sliding window groups, all block IDs go into - # sliding_window_block_ids; non_sliding_window_block_ids is empty. - sw_block_ids = set(job_status.sliding_window_block_ids or []) - non_sw_block_ids = set(job_status.non_sliding_window_block_ids or []) - assert len(sw_block_ids) > 0 - assert len(non_sw_block_ids) == 0 - - # Sliding window blocks should already be in the fence - # (registered at store creation time in _build_store_jobs). - for bid in sw_block_ids: - assert bid in runner.connector_scheduler._block_id_to_pending_jobs - assert job_id in runner.connector_scheduler._block_id_to_pending_jobs[bid] - - fence_before = dict(runner.connector_scheduler._block_id_to_pending_jobs) - - # Finish the request — request_finished should NOT add new entries - # because non_sliding_window_block_ids is empty for SW groups. - req_id = str(runner.req_id) - runner.scheduler.finish_requests({req_id}, RequestStatus.FINISHED_STOPPED) - - fence_after = runner.connector_scheduler._block_id_to_pending_jobs - - # Fence unchanged: no new entries from request_finished. - assert fence_before == fence_after - - -def test_multiple_in_flight_stores_all_registered_in_fence(request_runner): - """When a request finishes with multiple in-flight store jobs, - ALL jobs' non_sliding_window_block_ids are registered in the fence. - """ - block_size = 4 - block_size_factor = 1 - offloaded_block_size = block_size * block_size_factor - - runner = request_runner( - block_size=block_size, - num_gpu_blocks=100, - async_scheduling=False, - block_size_factor=block_size_factor, - ) + # Verify job_0 is in-flight. + assert len(runner.connector_scheduler._jobs) == 1 - # Create a request with 1 block. - runner.new_request(token_ids=[0] * offloaded_block_size) + # Second run: 4 more tokens + EOS → block 2 full → job_1 created. + # Request finishes → "all finished" flush fires → both jobs flushed. + # Re-set prepare_store after run() calls manager.reset_mock(). runner.manager.prepare_store.side_effect = lambda keys, req_context: ( generate_store_output(keys) ) - - # Schedule to create the first store job. - runner.scheduler.schedule() - runner._update_gpu_blocks() - - assert len(runner.connector_scheduler._jobs) == 1 - job_id_1 = next(iter(runner.connector_scheduler._jobs)) - job_status_1 = runner.connector_scheduler._jobs[job_id_1] - non_sw_1 = list(job_status_1.non_sliding_window_block_ids or []) - assert len(non_sw_1) > 0 - - # Manually inject a second in-flight store job to simulate a request - # that produced multiple store jobs across decode steps. - req_id = str(runner.req_id) - req_status = runner.connector_scheduler._req_status[req_id] - fake_block_id = 99 - fake_job_id = 999 - runner.connector_scheduler._jobs[fake_job_id] = TransferJobStatus( - req_id=req_id, - pending_count=runner.connector_scheduler.config.num_workers, - keys=set(), - is_store=True, - non_sliding_window_block_ids=[fake_block_id], + runner.run( + decoded_tokens=[0] * offloaded_block_size + [EOS_TOKEN_ID], + complete_transfers=False, + expected_stored=(0, 1, 2), + expected_flushed=(0, 1, 2), ) - req_status.transfer_jobs.add(fake_job_id) - - assert len(req_status.transfer_jobs) == 2 - # Finish the request — both jobs should be registered in the fence. - runner.scheduler.finish_requests({req_id}, RequestStatus.FINISHED_STOPPED) - - fence = runner.connector_scheduler._block_id_to_pending_jobs - for bid in non_sw_1: - assert bid in fence - assert job_id_1 in fence[bid] - assert fake_block_id in fence - assert fake_job_id in fence[fake_block_id] + # Post-condition: fence cleaned up after full lifecycle. + assert runner.connector_scheduler._block_id_to_pending_jobs == {} + assert len(runner.connector_scheduler._jobs) == 0 def test_fence_full_lifecycle_populate_and_cleanup(request_runner): From 1df1902efe1c172d8f6966d8cbad43873b8c0bc9 Mon Sep 17 00:00:00 2001 From: Alex Date: Tue, 16 Jun 2026 14:59:01 +0800 Subject: [PATCH 3/4] test(kv_connector): address reviewer suggestions for fence tests - Add post_step_fn callback to runner.run() for capturing mid-lifecycle state - Convert test_request_finished_with_pending_stores_populates_fence to use runner.run() - Remove redundant side_effect reset in test_multiple_in_flight_stores_all_flushed - Convert test_request_finished_mixed_full_attn_and_sliding_window to use runner.run() - Add cleanup verification to mixed architecture test Signed-off-by: AlexHuang Signed-off-by: Alex --- .../offloading_connector/test_scheduler.py | 132 ++++++++++-------- .../unit/offloading_connector/utils.py | 17 ++- 2 files changed, 84 insertions(+), 65 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 c674fe0c63d7..a645fea0f165 100644 --- a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py +++ b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py @@ -2085,33 +2085,38 @@ def test_request_finished_with_pending_stores_populates_fence(request_runner): generate_store_output(keys) ) - # Schedule to create store jobs (but don't complete transfers). - runner.scheduler.schedule() - runner._update_gpu_blocks() - - # Verify store job was created. - assert len(runner.connector_scheduler._jobs) > 0 - job_id = next(iter(runner.connector_scheduler._jobs)) - job_status = runner.connector_scheduler._jobs[job_id] - assert job_status.is_store - non_sw_block_ids = job_status.non_sliding_window_block_ids or [] - assert len(non_sw_block_ids) > 0 - - # Fence should be empty before request_finished - # (non-sliding-window blocks are only registered at request finish). - assert runner.connector_scheduler._block_id_to_pending_jobs == {} + # Capture fence state at each step to verify it was populated. + fence_snapshots: list[dict] = [] + job_block_ids: set[int] = set() - # Finish the request — triggers request_finished which populates the fence. - req_id = str(runner.req_id) - runner.scheduler.finish_requests({req_id}, RequestStatus.FINISHED_STOPPED) + def capture_fence(): + fence_snapshots.append( + dict(runner.connector_scheduler._block_id_to_pending_jobs) + ) + for js in runner.connector_scheduler._jobs.values(): + if js.is_store: + job_block_ids.update(js.non_sliding_window_block_ids or []) - # Verify fence index is populated with the job's non-SW block IDs. - for bid in non_sw_block_ids: - assert bid in runner.connector_scheduler._block_id_to_pending_jobs - assert job_id in runner.connector_scheduler._block_id_to_pending_jobs[bid] + # Run the full lifecycle: create job → finish request → flush → cleanup. + runner.run( + decoded_tokens=[EOS_TOKEN_ID], + complete_transfers=False, + expected_stored=(0, 1, 2), + expected_flushed=(0, 1, 2), + post_step_fn=capture_fence, + ) + + # Verify fence was populated at some point during the run. + assert len(job_block_ids) > 0, "No store job was created" + populated_fence = next((f for f in fence_snapshots if len(f) > 0), None) + assert populated_fence is not None, "Fence was never populated" - # req_status should still exist because the job is still in-flight. - assert req_id in runner.connector_scheduler._req_status + # Verify fence contained the job's non-SW block IDs. + for bid in job_block_ids: + assert bid in populated_fence, f"Block {bid} not in fence: {populated_fence}" + + # Verify fence is empty after full lifecycle (cleanup happened). + assert runner.connector_scheduler._block_id_to_pending_jobs == {} def test_multiple_in_flight_stores_all_flushed(request_runner): @@ -2154,10 +2159,8 @@ def test_multiple_in_flight_stores_all_flushed(request_runner): # Second run: 4 more tokens + EOS → block 2 full → job_1 created. # Request finishes → "all finished" flush fires → both jobs flushed. - # Re-set prepare_store after run() calls manager.reset_mock(). - runner.manager.prepare_store.side_effect = lambda keys, req_context: ( - generate_store_output(keys) - ) + # Note: prepare_store.side_effect persists across run() calls because + # manager.reset_mock() only resets call history, not side_effect. runner.run( decoded_tokens=[0] * offloaded_block_size + [EOS_TOKEN_ID], complete_transfers=False, @@ -2263,43 +2266,48 @@ def test_request_finished_mixed_full_attn_and_sliding_window( generate_store_output(keys) ) - runner.scheduler.schedule() - runner._update_gpu_blocks() + # Capture fence state and job block IDs at each step. + fence_snapshots: list[dict] = [] + sw_block_ids: set[int] = set() + non_sw_block_ids: set[int] = set() - assert len(runner.connector_scheduler._jobs) > 0 - job_id = next(iter(runner.connector_scheduler._jobs)) - job_status = runner.connector_scheduler._jobs[job_id] - assert job_status.is_store - - # Both types of block IDs should be present. - non_sw_block_ids = set(job_status.non_sliding_window_block_ids or []) - sw_block_ids = set(job_status.sliding_window_block_ids or []) - assert len(non_sw_block_ids) > 0 - assert len(sw_block_ids) > 0 - - # SW blocks should already be in the fence - # (registered at store creation time). - for bid in sw_block_ids: - assert bid in runner.connector_scheduler._block_id_to_pending_jobs - assert job_id in (runner.connector_scheduler._block_id_to_pending_jobs[bid]) + def capture_fence(): + fence_snapshots.append( + dict(runner.connector_scheduler._block_id_to_pending_jobs) + ) + for js in runner.connector_scheduler._jobs.values(): + if js.is_store: + sw_block_ids.update(js.sliding_window_block_ids or []) + non_sw_block_ids.update(js.non_sliding_window_block_ids or []) - # Snapshot the SW fence entries before request_finished. - sw_fence_entries = { - bid: set(runner.connector_scheduler._block_id_to_pending_jobs[bid]) - for bid in sw_block_ids - } + # Run the full lifecycle: create job → finish request → flush → cleanup. + runner.run( + decoded_tokens=[EOS_TOKEN_ID], + complete_transfers=False, + expected_stored=(0, 1, 2), + expected_flushed=(0, 1, 2), + post_step_fn=capture_fence, + ) - # Finish the request — only non-SW blocks should be newly registered. - req_id = str(runner.req_id) - runner.scheduler.finish_requests({req_id}, RequestStatus.FINISHED_STOPPED) + # Verify job had both SW and non-SW blocks. + assert len(sw_block_ids) > 0, "No SW blocks in store job" + assert len(non_sw_block_ids) > 0, "No non-SW blocks in store job" - fence = runner.connector_scheduler._block_id_to_pending_jobs + # Find the fence snapshot where both SW and non-SW blocks were present. + # SW blocks should appear at creation time, non-SW at request_finished. + populated_fence = None + for fence in fence_snapshots: + has_sw = all(bid in fence for bid in sw_block_ids) + has_non_sw = all(bid in fence for bid in non_sw_block_ids) + if has_sw and has_non_sw: + populated_fence = fence + break - # Non-SW blocks are now in the fence. - for bid in non_sw_block_ids: - assert bid in fence - assert job_id in fence[bid] + assert populated_fence is not None, ( + f"Fence never contained both SW {sw_block_ids} and " + f"non-SW {non_sw_block_ids} blocks. Snapshots: {fence_snapshots}" + ) - # SW fence entries are unchanged (not duplicated). - for bid, expected_jobs in sw_fence_entries.items(): - assert fence[bid] == expected_jobs + # Verify fence is empty after full lifecycle (cleanup happened). + assert runner.connector_scheduler._block_id_to_pending_jobs == {} + assert len(runner.connector_scheduler._jobs) == 0 diff --git a/tests/v1/kv_connector/unit/offloading_connector/utils.py b/tests/v1/kv_connector/unit/offloading_connector/utils.py index f6a354ebd43b..446453191468 100644 --- a/tests/v1/kv_connector/unit/offloading_connector/utils.py +++ b/tests/v1/kv_connector/unit/offloading_connector/utils.py @@ -1,6 +1,6 @@ # SPDX-License-Identifier: Apache-2.0 # SPDX-FileCopyrightText: Copyright contributors to the vLLM project -from collections.abc import Iterable, Iterator +from collections.abc import Callable, Iterable, Iterator from dataclasses import dataclass from typing import Any from unittest.mock import MagicMock @@ -430,7 +430,12 @@ def _update_gpu_blocks(self): for block_idx, block in enumerate(blocks): self.gpu_blocks[block.block_id] = GPUBlock(group_idx, block_idx) - def _run(self, decoded_tokens: list[int], complete_transfers: bool): + def _run( + self, + decoded_tokens: list[int], + complete_transfers: bool, + post_step_fn: Callable[[], None] | None = None, + ): """ Runs multiple engine (scheduler + worker) steps. Assumes a single request is running. @@ -438,6 +443,8 @@ def _run(self, decoded_tokens: list[int], complete_transfers: bool): Args: decoded_tokens: the tokens to yield at each step. complete_transfers: complete transfers immediately + post_step_fn: optional callback invoked after each step's + update_from_output(), before the next schedule(). """ tokens_iter = iter(decoded_tokens) @@ -500,6 +507,9 @@ def _run(self, decoded_tokens: list[int], complete_transfers: bool): else: self.scheduler.update_from_output(scheduler_output, model_runner_output) + if post_step_fn is not None: + post_step_fn() + if ( prev_token_id == EOS_TOKEN_ID and prev_token_id != token_id @@ -545,6 +555,7 @@ def run( expected_stored: tuple[int | tuple[int, int], ...] = (), expected_loaded: tuple[int | tuple[int, int], ...] = (), expected_flushed: tuple[int | tuple[int, int], ...] = (), + post_step_fn: Callable[[], None] | None = None, ): """ Runs multiple engine (scheduler + worker) steps. @@ -570,7 +581,7 @@ def run( expected_flushed_gpu_blocks = self._to_gpu_blocks(expected_flushed) self.manager.reset_mock() - self._run(decoded_tokens, complete_transfers) + self._run(decoded_tokens, complete_transfers, post_step_fn=post_step_fn) loaded_gpu_blocks: set[GPUBlock] = set() for transfer in self.completed_loads: From 06f5d3d375d10f9ffa97ccc509d480a948e8016a Mon Sep 17 00:00:00 2001 From: Alex Date: Wed, 17 Jun 2026 18:03:07 +0800 Subject: [PATCH 4/4] test(kv_connector): adapt fence tests for non-blocking drain (#45595) - Adapt tests to use block-reuse fence flush instead of all-finished flush - Add test_multiple_in_flight_stores_all_flushed_by_fence (multi-job scenario) - Add req_status cleanup verification to populates_fence test - Strengthen existing fence tests with post_step_fn snapshot verification - Remove test_fence_full_lifecycle (overlaps with populates_fence) Signed-off-by: Alex --- .../offloading_connector/test_scheduler.py | 174 ++++++++++-------- 1 file changed, 99 insertions(+), 75 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 a645fea0f165..7d4940183c58 100644 --- a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py +++ b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py @@ -835,9 +835,27 @@ def test_fence_at_update_state_after_alloc(request_runner): runner.manager.prepare_store.side_effect = lambda keys, req_context: ( generate_store_output(keys) ) - runner.run(decoded_tokens=[EOS_TOKEN_ID], complete_transfers=False) + + # Capture fence snapshots to verify block 0 is registered. + fence_snapshots: list[dict] = [] + + def capture_fence(): + fence_snapshots.append( + dict(runner.connector_scheduler._block_id_to_pending_jobs) + ) + + runner.run( + decoded_tokens=[EOS_TOKEN_ID], + complete_transfers=False, + post_step_fn=capture_fence, + ) assert runner.connector_scheduler._block_id_to_pending_jobs + # Verify fence was populated with the store job's block IDs. + populated_fence = next((f for f in fence_snapshots if f), None) + assert populated_fence is not None, "Fence was never populated" + assert len(populated_fence) > 0, "Fence is empty" + runner.scheduler.reset_prefix_cache() runner.new_request(token_ids=[0] * 4) runner.connector_scheduler._maximal_prefix_lookup = lambda key, req_context: 1 @@ -868,9 +886,27 @@ def test_fence_at_build_store_jobs(request_runner): runner.manager.prepare_store.side_effect = lambda keys, req_context: ( generate_store_output(keys) ) - runner.run(decoded_tokens=[EOS_TOKEN_ID], complete_transfers=False) + + # Capture fence snapshots to verify block 0 is registered. + fence_snapshots: list[dict] = [] + + def capture_fence(): + fence_snapshots.append( + dict(runner.connector_scheduler._block_id_to_pending_jobs) + ) + + runner.run( + decoded_tokens=[EOS_TOKEN_ID], + complete_transfers=False, + post_step_fn=capture_fence, + ) assert runner.connector_scheduler._block_id_to_pending_jobs + # Verify fence was populated with the store job's block IDs. + populated_fence = next((f for f in fence_snapshots if f), None) + assert populated_fence is not None, "Fence was never populated" + assert len(populated_fence) > 0, "Fence is empty" + runner.scheduler.reset_prefix_cache() runner.new_request(token_ids=[1] * 4) runner.connector_scheduler._maximal_prefix_lookup = lambda key, req_context: 0 @@ -2072,15 +2108,17 @@ def test_request_finished_with_pending_stores_populates_fence(request_runner): block_size_factor = 1 offloaded_block_size = block_size * block_size_factor + # Use 2 GPU blocks so the second run reuses the same blocks, + # triggering a fence-based flush of the in-flight job from run 1. runner = request_runner( block_size=block_size, - num_gpu_blocks=100, + num_gpu_blocks=2, async_scheduling=False, block_size_factor=block_size_factor, ) - # Create a request with 3 blocks of data. - runner.new_request(token_ids=[0] * offloaded_block_size * 3) + # 4 prompt tokens → 1 GPU block (block 0) + runner.new_request(token_ids=[0] * offloaded_block_size) runner.manager.prepare_store.side_effect = lambda keys, req_context: ( generate_store_output(keys) ) @@ -2097,12 +2135,11 @@ def capture_fence(): if js.is_store: job_block_ids.update(js.non_sliding_window_block_ids or []) - # Run the full lifecycle: create job → finish request → flush → cleanup. + # Run 1: create store job, finish request, populate fence. + # With non-blocking drain (#45595), the job stays in-flight. runner.run( decoded_tokens=[EOS_TOKEN_ID], complete_transfers=False, - expected_stored=(0, 1, 2), - expected_flushed=(0, 1, 2), post_step_fn=capture_fence, ) @@ -2115,108 +2152,83 @@ def capture_fence(): for bid in job_block_ids: assert bid in populated_fence, f"Block {bid} not in fence: {populated_fence}" + # Run 2: block reuse triggers fence-based flush → cleanup. + runner.scheduler.reset_prefix_cache() + runner.new_request(token_ids=[0] * offloaded_block_size) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + runner.run( + decoded_tokens=[EOS_TOKEN_ID], + expected_stored=(0,), + expected_flushed=(0,), + ) + # Verify fence is empty after full lifecycle (cleanup happened). assert runner.connector_scheduler._block_id_to_pending_jobs == {} + # req_status should be removed. + req_id = str(runner.req_id) + assert req_id not in runner.connector_scheduler._req_status -def test_multiple_in_flight_stores_all_flushed(request_runner): +def test_multiple_in_flight_stores_all_flushed_by_fence(request_runner): """When a request finishes with multiple in-flight store jobs, - ALL jobs are flushed via the fence mechanism (observable behavior). + ALL jobs are flushed when a new request reuses their blocks. - Uses two runner.run() calls to naturally produce two store jobs: - - First run: 4 decoded tokens fills blocks 0-1 → job_0 created for both - - Second run: 4 more tokens + EOS fills block 2 → job_1 created, - request finishes, "all finished" flush fires → both jobs flushed + Uses three runner.run() calls: + - Run 1: decode fills a block → job_0 created + - Run 2: decode fills another block + EOS → job_1 created, request finishes + - Run 3: block reuse → both jobs flushed via fence """ block_size = 4 block_size_factor = 1 offloaded_block_size = block_size * block_size_factor + # 4 GPU blocks: block 0 is null, blocks 1-3 are usable. runner = request_runner( block_size=block_size, - num_gpu_blocks=100, + num_gpu_blocks=4, async_scheduling=False, block_size_factor=block_size_factor, ) - # 4 prompt tokens → 1 GPU block (block 0) + # Prompt: 4 tokens → block 1 runner.new_request(token_ids=[0] * offloaded_block_size) runner.manager.prepare_store.side_effect = lambda keys, req_context: ( generate_store_output(keys) ) - # First run: 4 decoded tokens → blocks 0-1 full → job_0 created for both. - # complete_transfers=False keeps job_0 in-flight (not completed yet). + # Run 1: 4 decoded tokens → block 2 full → job_0 created for block 1. runner.run( decoded_tokens=[0] * offloaded_block_size, complete_transfers=False, - expected_stored=(), - expected_flushed=(), ) + assert len(runner.connector_scheduler._jobs) >= 1 - # Verify job_0 is in-flight. - assert len(runner.connector_scheduler._jobs) == 1 - - # Second run: 4 more tokens + EOS → block 2 full → job_1 created. - # Request finishes → "all finished" flush fires → both jobs flushed. - # Note: prepare_store.side_effect persists across run() calls because - # manager.reset_mock() only resets call history, not side_effect. + # Run 2: 4 more tokens + EOS → block 3 full → more jobs created. + # Request finishes → all jobs registered in fence. runner.run( decoded_tokens=[0] * offloaded_block_size + [EOS_TOKEN_ID], complete_transfers=False, - expected_stored=(0, 1, 2), - expected_flushed=(0, 1, 2), - ) - - # Post-condition: fence cleaned up after full lifecycle. - assert runner.connector_scheduler._block_id_to_pending_jobs == {} - assert len(runner.connector_scheduler._jobs) == 0 - - -def test_fence_full_lifecycle_populate_and_cleanup(request_runner): - """Verifies the complete fence lifecycle using the standard runner flow: - 1. Store job created (in-flight) - 2. Request finishes → request_finished populates the fence - 3. build_connector_meta detects all requests finished → flushes jobs - 4. Flush completes the job - 5. update_connector_output cleans up the fence and req_status - - The fence population is verified by - test_request_finished_with_pending_stores_populates_fence. - This test verifies the cleanup after the full lifecycle. - """ - block_size = 4 - block_size_factor = 1 - offloaded_block_size = block_size * block_size_factor - - runner = request_runner( - block_size=block_size, - num_gpu_blocks=100, - async_scheduling=False, - block_size_factor=block_size_factor, ) + num_jobs = len(runner.connector_scheduler._jobs) + assert num_jobs >= 2, f"Expected multiple in-flight jobs, got {num_jobs}" - runner.new_request(token_ids=[0] * offloaded_block_size) + # Run 3: block reuse → fence flushes both jobs. + runner.scheduler.reset_prefix_cache() + runner.new_request(token_ids=[0] * offloaded_block_size * 3) runner.manager.prepare_store.side_effect = lambda keys, req_context: ( generate_store_output(keys) ) - - # Run with complete_transfers=False: creates store job, finishes request, - # populates fence, flushes job (which completes it), cleans up fence. runner.run( decoded_tokens=[EOS_TOKEN_ID], - complete_transfers=False, - expected_stored=(0,), - expected_flushed=(0,), + expected_stored=(0, 1, 2), + expected_flushed=(0, 1, 2), ) - # After the full lifecycle, fence should be empty. + # Post-condition: fence cleaned up, all jobs gone. assert runner.connector_scheduler._block_id_to_pending_jobs == {} - # Job should be removed. assert len(runner.connector_scheduler._jobs) == 0 - # req_status should be removed. - req_id = str(runner.req_id) - assert req_id not in runner.connector_scheduler._req_status def test_request_finished_mixed_full_attn_and_sliding_window( @@ -2253,15 +2265,17 @@ def test_request_finished_mixed_full_attn_and_sliding_window( ), ] + # Use 4 GPU blocks (2 per group) so run 2 reuses the same blocks, + # triggering a fence-based flush. runner = request_runner( block_size=block_size, - num_gpu_blocks=100, + num_gpu_blocks=4, async_scheduling=False, kv_cache_groups=kv_cache_groups, ) - # 3 blocks of prompt (12 tokens) — enough for both groups. - runner.new_request(token_ids=[0] * block_size * 3) + # 1 block of prompt (4 tokens) — 1 block per group. + runner.new_request(token_ids=[0] * block_size) runner.manager.prepare_store.side_effect = lambda keys, req_context: ( generate_store_output(keys) ) @@ -2280,12 +2294,10 @@ def capture_fence(): sw_block_ids.update(js.sliding_window_block_ids or []) non_sw_block_ids.update(js.non_sliding_window_block_ids or []) - # Run the full lifecycle: create job → finish request → flush → cleanup. + # Run 1: create store job, finish request, populate fence. runner.run( decoded_tokens=[EOS_TOKEN_ID], complete_transfers=False, - expected_stored=(0, 1, 2), - expected_flushed=(0, 1, 2), post_step_fn=capture_fence, ) @@ -2308,6 +2320,18 @@ def capture_fence(): f"non-SW {non_sw_block_ids} blocks. Snapshots: {fence_snapshots}" ) + # Run 2: block reuse triggers fence-based flush of the old job. + runner.scheduler.reset_prefix_cache() + runner.new_request(token_ids=[0] * block_size) + runner.manager.prepare_store.side_effect = lambda keys, req_context: ( + generate_store_output(keys) + ) + runner.run( + decoded_tokens=[EOS_TOKEN_ID], + expected_stored=((0, 0), (1, 0)), + expected_flushed=((1, 0),), + ) + # Verify fence is empty after full lifecycle (cleanup happened). assert runner.connector_scheduler._block_id_to_pending_jobs == {} assert len(runner.connector_scheduler._jobs) == 0