From e3de61bd8f25852630f1b577765ff57a8f8951c6 Mon Sep 17 00:00:00 2001 From: Ketor Date: Sat, 12 Sep 2026 00:14:18 +0800 Subject: [PATCH] test(vllm): verify preemption source lifetime and consumer recovery --- CHANGELOG.md | 10 ++ VERSION | 2 +- integration/common/pyproject.toml | 2 +- integration/lmcache/pyproject.toml | 4 +- integration/vllm/README.md | 11 ++ integration/vllm/pyproject.toml | 4 +- integration/vllm/tests/test_preempt_fence.py | 167 ++++++++++++++++-- .../tests/test_scheduler_full_block_ids.py | 103 +++++++++-- 8 files changed, 261 insertions(+), 42 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 9e871a9..69efe7c 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,6 +2,16 @@ ## Unreleased +### v2.26.3 — vLLM acceptance coverage + +- Replace the forwarding-only preemption test with a source-lifetime regression: + block reuse must wait for the old SAVE, which must retain the original bytes. +- Cover consumer MRV1 synchronous/asynchronous recovery, prior-step/same-step + preemption, continued prefill, decode and completion without duplicate LOAD + or unauthorized SAVE. +- Clarify that compatible object layouts do not validate old cache contents. + This release changes tests and documentation, not connector runtime behavior. + ### vLLM preemption safety - Fence cancelled receive work and in-flight SAVE sources through vLLM's diff --git a/VERSION b/VERSION index cd74b3e..90f3e72 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -2.26.2 \ No newline at end of file +2.26.3 \ No newline at end of file diff --git a/integration/common/pyproject.toml b/integration/common/pyproject.toml index cc326cd..7a9c5e5 100644 --- a/integration/common/pyproject.toml +++ b/integration/common/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "dfkv-common" -version = "2.26.2" +version = "2.26.3" description = "Canonical namespace and pool-key schema shared by dfkv connectors" requires-python = ">=3.9" diff --git a/integration/lmcache/pyproject.toml b/integration/lmcache/pyproject.toml index fbbf9ac..19335d7 100644 --- a/integration/lmcache/pyproject.toml +++ b/integration/lmcache/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "dfkv-connector" -version = "2.26.2" +version = "2.26.3" description = "LMCache RemoteConnector for the dfkv KV cache (ctypes over libdfkv.so)" readme = "README.md" requires-python = ">=3.9" @@ -14,7 +14,7 @@ authors = [{ name = "Wine93", email = "wine93.info@gmail.com" }] # runtime by path (DFKV_LIB / remote_storage_plugin.dfkv.lib). No bundled .so, # no CPython extension, so the wheel is platform-independent. dependencies = [ - "dfkv-common==2.26.2", + "dfkv-common==2.26.3", "lmcache", "torch", ] diff --git a/integration/vllm/README.md b/integration/vllm/README.md index a51383e..c7be46b 100644 --- a/integration/vllm/README.md +++ b/integration/vllm/README.md @@ -274,6 +274,17 @@ SAVE queue entries also carry a worker-local generation. Reusing a request ID cannot revive an old queued SAVE or let its cancellation decrement the new generation's completion counter. +The engine must call `handle_preemptions` before model execution can reuse +blocks. `start_load_kv` is not a substitute: the engine can defer it until after +forward when there are no synchronous loads. Consumer-only cached resumes +still require LOAD metadata even though they never authorize SAVE. + +The v2.26.2 lifecycle fix preserves the v2.26.1 object format; it cannot detect +or repair equal-length KV objects corrupted by an earlier block-reuse race. +When upgrading a deployment that could have encountered that race, use a fresh +`model_revision` consistently across writers and readers to isolate old objects. +Keep the old namespace for normal retention/eviction rather than clearing it. + Different namespace/key bytes are a cold miss. Byte-layout changes not captured by the effective cache-spec identity still require a source-controlled layout-ID bump and coordinated writer/reader deployment. diff --git a/integration/vllm/pyproject.toml b/integration/vllm/pyproject.toml index 62de3a9..4d0dd28 100644 --- a/integration/vllm/pyproject.toml +++ b/integration/vllm/pyproject.toml @@ -4,12 +4,12 @@ build-backend = "setuptools.build_meta" [project] name = "dfkv-vllm" -version = "2.26.2" +version = "2.26.3" description = "Direct vLLM KVConnectorBase_V1 connector for dfkv (GPUDirect RDMA, no LMCache)" requires-python = ">=3.12" # vllm + torch are provided by the runtime image; not pinned here. The connector # itself is pure Python (ctypes over libdfkv.so), so there is no native build. -dependencies = ["dfkv-common==2.26.2"] +dependencies = ["dfkv-common==2.26.3"] # Telemetry is opt-in: the OTel SDK is only needed when DFKV_METRICS_ENABLED / # DFKV_TRACE_ENABLED is set. Without this extra the connector stays dependency- diff --git a/integration/vllm/tests/test_preempt_fence.py b/integration/vllm/tests/test_preempt_fence.py index da95690..393b4a9 100644 --- a/integration/vllm/tests/test_preempt_fence.py +++ b/integration/vllm/tests/test_preempt_fence.py @@ -8,14 +8,27 @@ dfkv_vllm on PYTHONPATH, e.g. pip install -e integration/vllm) """ +import ctypes import threading import types import unittest +from unittest.mock import patch try: + import torch + from vllm.v1.core.kv_cache_utils import BlockHash + from vllm.v1.kv_cache_interface import FullAttentionSpec, KVCacheGroupSpec + from dfkv_vllm.connector import DfkvStoreConnector - from dfkv_vllm.data import DfkvStoreConnectorMetadata - from dfkv_vllm.worker import KVCacheStoreSendingThread + from dfkv_vllm.coordinator import DfkvStoreCoordinator + from dfkv_vllm.data import ( + ChunkedTokenDatabase, + DfkvStoreConnectorMetadata, + KeyMetadata, + PoolKey, + ReqMeta, + ) + from dfkv_vllm.worker import DfkvStoreWorker, KVCacheStoreSendingThread HAVE_VLLM = True except ImportError: # pragma: no cover - vllm not installed @@ -62,27 +75,143 @@ def _mk_thread(self, coord) -> "KVCacheStoreSendingThread": self._threads.append(t) return t - def test_connector_fences_preemption_before_starting_loads(self): - calls: list[str] = [] + def test_pre_forward_preserves_preempted_save_source_until_put_finishes(self): + # Exercise the installed engine hook when available, without importing + # or constructing GPUModelRunner. Older engines can still exercise the + # connector's pre-forward contract, but that is NOT engine-order proof. + try: + from vllm.v1.worker.gpu.kv_connector import ActiveKVConnector + except ImportError: + ActiveKVConnector = None + + block_size = 16 + original = bytes([17]) * block_size + replacement = bytes([99]) * block_size + pool = ctypes.create_string_buffer(2 * block_size) + source = ctypes.addressof(pool) + block_size + ctypes.memmove(source, original, block_size) + key_metadata = KeyMetadata( + model_name="preempt-source-lifetime", dp_size=1, dp_rank=-1, + tp_size=1, tp_rank=0, pcp_size=1, pcp_rank=0, + dcp_size=1, dcp_rank=0, pp_size=1, pp_rank=0, + ) + database = ChunkedTokenDatabase(key_metadata, block_size=block_size) + database.set_seg_layout([ + (ctypes.addressof(pool), block_size, block_size), + ]) + spec = FullAttentionSpec( + block_size=block_size, num_kv_heads=1, head_size=1, + head_size_v=0, dtype=torch.uint8, + ) + coordinator = DfkvStoreCoordinator( + [KVCacheGroupSpec(["kv"], spec)], block_size, block_size, + ) + put_entered = threading.Event() + release_put = threading.Event() + fence_or_step_finished = threading.Event() + overwritten = threading.Event() + saved: dict[bytes, bytes] = {} + step_errors: list[BaseException] = [] - class Worker: - def handle_preemptions(self, metadata): - calls.append("preemption") + class DelayedMemoryClient: + # Substitute only the native boundary: descriptor construction, + # queue/dequeue, cancellation and in-flight tracking stay real. + def batch_exist(self, keys): + return [0] * len(keys) - def start_load_kv(self, metadata): - calls.append("load") + def batch_put_sg(self, keys, pointers, sizes): + put_entered.set() + if not release_put.wait(10): + raise TimeoutError("test did not release the pending SAVE") + for key, ptrs, lengths in zip(keys, pointers, sizes, strict=True): + saved[key] = b"".join( + ctypes.string_at(ptr, length) + for ptr, length in zip(ptrs, lengths, strict=True) + ) + return [0] * len(keys) + sender = KVCacheStoreSendingThread( + client=DelayedMemoryClient(), coord=coordinator, + token_databases=[database], block_size=block_size, + tp_rank=0, stripe_idx=0, stripe_step=1, kv_role="kv_producer", + ready_event=threading.Event(), + ) + worker = object.__new__(DfkvStoreWorker) + worker.kv_send_thread = sender + worker.kv_recv_thread = None connector = object.__new__(DfkvStoreConnector) - connector.connector_worker = Worker() - connector._begin_call = lambda: None - connector._finish_call = lambda: None - metadata = DfkvStoreConnectorMetadata(set(), {"resumed"}) - connector._get_connector_metadata = lambda: metadata - - connector.handle_preemptions(metadata) - self.assertEqual(calls, ["preemption"]) - connector.start_load_kv(None) - self.assertEqual(calls, ["preemption", "load"]) + connector.connector_worker = worker + connector._shutdown_condition = threading.Condition() + connector._shutdown = False + connector._inflight_calls = 0 + metadata = DfkvStoreConnectorMetadata(set(), {"preempted"}) + if ActiveKVConnector is not None: + runner_hook = object.__new__(ActiveKVConnector) + runner_hook.kv_connector = connector + runner_hook._disabled = False + runner_hook._pending_load_start = False + scheduler_output = types.SimpleNamespace( + kv_connector_metadata=metadata, has_sync_kv_loads=False, + ) + + def next_model_step(): + try: + if ActiveKVConnector is not None: + runner_hook.pre_forward(scheduler_output) + else: + connector.handle_preemptions(metadata) + # Model reuse of the same physical block immediately after + # the real pre-forward hook admits the next model step. + ctypes.memmove(source, replacement, block_size) + overwritten.set() + except BaseException as exc: + step_errors.append(exc) + finally: + fence_or_step_finished.set() + + # Observe the real condition wait, not an arbitrary sleep. Acquiring + # its lock below guarantees the hook has actually blocked, or finished + # early (the v2.26.1 missing-hook bug), before inspecting source bytes. + condition_wait = sender._active_cv.wait + + def observed_wait(timeout=None): + fence_or_step_finished.set() + return condition_wait(timeout) + + step = threading.Thread(target=next_model_step, daemon=True) + request = ReqMeta( + req_id="preempted", token_len_chunk=block_size, + block_ids=([1],), block_hashes=[BlockHash(b"x" * 32)], + ) + with patch.object(sender._active_cv, "wait", observed_wait): + try: + sender.start() + self.assertTrue(sender.ready_event.wait(5)) + sender.add_stored_request(request) + self.assertTrue(sender.add_request(request)) + self.assertTrue(put_entered.wait(5), "SAVE never reached native PUT") + step.start() + self.assertTrue(fence_or_step_finished.wait(5)) + with sender._active_cv: + self.assertEqual(step_errors, []) + self.assertFalse( + overwritten.is_set(), + "pre-forward admitted block reuse while SAVE was reading", + ) + self.assertEqual(ctypes.string_at(source, block_size), original) + release_put.set() + step.join(5) + self.assertFalse(step.is_alive(), "preemption fence did not release") + self.assertEqual(step_errors, []) + self.assertTrue(overwritten.is_set()) + key = PoolKey(key_metadata, request.block_hashes[0].hex()).to_bytes() + self.assertEqual(saved, {key: original}) + self.assertEqual(ctypes.string_at(source, block_size), replacement) + finally: + release_put.set() + sender.stop(cancel_pending=True) + if step.ident is not None: + step.join(5) def test_wait_blocks_while_put_executes_and_returns_after(self): diff --git a/integration/vllm/tests/test_scheduler_full_block_ids.py b/integration/vllm/tests/test_scheduler_full_block_ids.py index e7614f2..f5d510e 100644 --- a/integration/vllm/tests/test_scheduler_full_block_ids.py +++ b/integration/vllm/tests/test_scheduler_full_block_ids.py @@ -91,41 +91,110 @@ def test_cache_bypass_discards_stale_external_admission(): assert scheduler.load_specs == {"unrelated": other_spec} -def test_consumer_cached_resume_emits_load_without_save(): +@pytest.mark.parametrize("same_step_preemption", [False, True]) +@pytest.mark.parametrize("async_pending", [False, True]) +def test_consumer_cached_resume_emits_load_without_save( + async_pending, same_step_preemption, +): scheduler = object.__new__(DfkvStoreScheduler) scheduler.kv_role = "kv_consumer" scheduler.client = MagicMock() - scheduler.load_specs = {"resumed": LoadSpec(0, 4, True)} + scheduler.client.lookup.return_value = 4 + scheduler.lookup_async = False + scheduler.load_async = async_pending + scheduler.load_specs = {} scheduler._request_trackers = {} - scheduler._preempted_req_ids = {"resumed"} - scheduler._unfinished_request_ids = {"resumed"} + scheduler._preempted_req_ids = set() + scheduler._unfinished_request_ids = set() + scheduler._unfinished_requests = {} scheduler._allocated_req_ids = set() scheduler._block_size = 4 request = SimpleNamespace( - request_id="resumed", block_hashes=[], all_token_ids=list(range(8)), - num_computed_tokens=4, + request_id="resumed", block_hashes=[], all_token_ids=list(range(16)), + num_tokens=16, num_computed_tokens=0, + ) + old_blocks = ([1, 2, 3, 4], [5, 6, 7, 8]) + scheduler.update_state_after_alloc( + request, SimpleNamespace(get_block_ids=lambda: old_blocks), 0, ) - blocks = ([10, 11],) - scheduler._unfinished_requests = {"resumed": (request, blocks)} step = SimpleNamespace( finished_req_ids=set(), preempted_req_ids=set(), - scheduled_new_reqs=[], - scheduled_cached_reqs=SimpleNamespace( - req_ids=["resumed"], new_block_ids=[([11],)], - num_computed_tokens=[4], - ), - num_scheduled_tokens={"resumed": 4}, + scheduled_new_reqs=[SimpleNamespace( + req_id="resumed", num_computed_tokens=0, block_ids=old_blocks, + prefill_token_ids=None, prompt_token_ids=list(range(8)), + )], + scheduled_cached_reqs=SimpleNamespace(req_ids=[]), + num_scheduled_tokens={"resumed": 8}, ) + assert scheduler.build_connector_meta(step).requests == [] - metadata = scheduler.build_connector_meta(step) + step.preempted_req_ids = {"resumed"} + step.scheduled_new_reqs = [] + step.num_scheduled_tokens = {} + if not same_step_preemption: + assert scheduler.build_connector_meta(step).requests == [] + step.preempted_req_ids = set() + + # Lookup and allocation admit a real LOAD; the resumed MRV1 request only + # carries this step's delta, not the complete destination block table. + assert scheduler.get_num_new_matched_tokens(request, 0) == (4, async_pending) + blocks = ([10, 11], [20, 21]) + scheduler.update_state_after_alloc( + request, SimpleNamespace(get_block_ids=lambda: blocks), 4, + ) + request.num_computed_tokens = 4 + resumed_cached = SimpleNamespace( + req_ids=["resumed"], new_block_ids=[([11], [21])], + num_computed_tokens=[4], + ) + if not async_pending: + step.scheduled_cached_reqs = resumed_cached + step.num_scheduled_tokens = {"resumed": 4} + metadata = scheduler.build_connector_meta(step) assert len(metadata.requests) == 1 restored = metadata.requests[0] + assert restored.req_id == "resumed" assert restored.load_spec is not None and restored.load_spec.can_load + assert restored.load_spec.vllm_cached_tokens == 0 + assert restored.load_spec.kvpool_cached_tokens == 4 assert restored.can_save is False assert restored.block_ids == blocks - assert "resumed" not in scheduler.load_specs - assert "resumed" not in scheduler._preempted_req_ids + + if async_pending: + # Waiting another tick must not submit the already issued LOAD again. + step.preempted_req_ids = set() + assert scheduler.build_connector_meta(step).requests == [] + # After completion MRV1 can resume without allocating another block. + resumed_cached.new_block_ids = [None] + step.scheduled_cached_reqs = resumed_cached + step.num_scheduled_tokens = {"resumed": 4} + assert scheduler.build_connector_meta(step).requests == [] + + # Finish re-prefilling prompt + generated history, then decode. Consumers + # must neither reload the prefix nor SAVE newly computed chunk/decode KV. + for computed, count, delta in ( + (8, 4, ([12], [22])), + (12, 4, ([13], [23])), + (16, 1, ([14], [24])), + (17, 1, None), + ): + if computed >= request.num_tokens: + request.all_token_ids.append(computed) + request.num_tokens += 1 + request.num_computed_tokens = computed + step.preempted_req_ids = set() + step.scheduled_cached_reqs = SimpleNamespace( + req_ids=["resumed"], new_block_ids=[delta], + num_computed_tokens=[computed], + ) + step.num_scheduled_tokens = {"resumed": count} + assert scheduler.build_connector_meta(step).requests == [] + + step.finished_req_ids = {"resumed"} + step.scheduled_cached_reqs = SimpleNamespace(req_ids=[]) + step.num_scheduled_tokens = {} + assert scheduler.build_connector_meta(step).requests == [] @pytest.mark.parametrize("cached_resume", [False, True])