Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 10 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion VERSION
Original file line number Diff line number Diff line change
@@ -1 +1 @@
2.26.2
2.26.3
2 changes: 1 addition & 1 deletion integration/common/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"

Expand Down
4 changes: 2 additions & 2 deletions integration/lmcache/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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",
]
Expand Down
11 changes: 11 additions & 0 deletions integration/vllm/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
4 changes: 2 additions & 2 deletions integration/vllm/pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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-
Expand Down
167 changes: 148 additions & 19 deletions integration/vllm/tests/test_preempt_fence.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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):
Expand Down
Loading
Loading