Skip to content
Closed
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
47 changes: 4 additions & 43 deletions components/src/dynamo/vllm/cache_info.py
Original file line number Diff line number Diff line change
@@ -1,20 +1,10 @@
# SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0

import logging
from typing import Any

from vllm.config import VllmConfig
from vllm.v1.engine.async_llm import AsyncLLM

logger = logging.getLogger(__name__)

DYNAMO_KV_EVENT_BLOCK_SIZE_KEY = "dynamo_kv_event_block_size"
MAIN_ATTENTION_KV_CACHE_KINDS = {
"full_attention",
"mla_attention",
"sink_full_attention",
}


def get_configured_kv_event_block_size(vllm_config: VllmConfig) -> int:
Expand All @@ -26,43 +16,14 @@ def get_configured_kv_event_block_size(vllm_config: VllmConfig) -> int:
)


def select_main_attention_block_size(
group_metadata: list[dict[str, Any]],
fallback_block_size: int,
) -> int:
"""Select the main-attention KV block size from engine cache-group metadata."""
if not group_metadata:
return fallback_block_size

for group in group_metadata:
if group.get("kind") in MAIN_ATTENTION_KV_CACHE_KINDS:
return group.get("block_size", fallback_block_size)

return fallback_block_size


async def configure_kv_event_block_size(
engine: AsyncLLM,
vllm_config: VllmConfig,
) -> int:
"""Fetch engine cache-group metadata and cache the KV event block size on vLLM config."""
fallback_block_size = vllm_config.cache_config.block_size
try:
group_metadata = await engine.engine_core.call_utility_async(
"get_kv_cache_group_metadata"
)
except Exception as e:
logger.warning(
"Failed to fetch KV cache group metadata; falling back to "
"vLLM cache_config.block_size: %s",
e,
)
kv_event_block_size = fallback_block_size
else:
kv_event_block_size = select_main_attention_block_size(
group_metadata,
fallback_block_size,
)
"""Cache the engine's effective attention block size on vLLM config."""
kv_event_block_size = engine.vllm_config.cache_config.effective_attention_block_size
if kv_event_block_size is None:
kv_event_block_size = vllm_config.cache_config.block_size

if vllm_config.additional_config is None:
vllm_config.additional_config = {}
Expand Down
36 changes: 36 additions & 0 deletions components/src/dynamo/vllm/tests/test_vllm_cache_info.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,36 @@
# SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
# SPDX-License-Identifier: Apache-2.0

from types import SimpleNamespace

import pytest

from dynamo.vllm.cache_info import (
DYNAMO_KV_EVENT_BLOCK_SIZE_KEY,
configure_kv_event_block_size,
)

pytestmark = [
pytest.mark.unit,
pytest.mark.vllm,
pytest.mark.gpu_0,
pytest.mark.pre_merge,
]


@pytest.mark.asyncio
@pytest.mark.parametrize(
"effective_size,expected", [(16, 16), (32, 32), (1056, 1056), (None, 16)]
)
async def test_configure_uses_effective_attention_block_size(effective_size, expected):
config = SimpleNamespace(
cache_config=SimpleNamespace(
block_size=16, effective_attention_block_size=effective_size
),
additional_config=None,
)
engine = SimpleNamespace(vllm_config=config)

assert await configure_kv_event_block_size(engine, config) == expected
assert config.additional_config[DYNAMO_KV_EVENT_BLOCK_SIZE_KEY] == expected
assert config.cache_config.block_size == 16
44 changes: 44 additions & 0 deletions tests/router/common.py
Original file line number Diff line number Diff line change
Expand Up @@ -148,6 +148,50 @@ async def _assert_overlap_scores(
########################################################


def _test_full_block_event_delivery(
engine_workers, model_name: str, block_size: int, request_plane: str
):
"""Require a worker's full KV blocks to reach the router's index."""
with managed_runtime(request_plane=request_plane) as runtime:
endpoint = runtime.endpoint(
f"{engine_workers.namespace}.{engine_workers.component_name}.generate"
)

async def run_test():
router = _create_kv_router_with_timeout(
router_factory=lambda: KvRouter(
endpoint=endpoint,
block_size=block_size,
kv_router_config=KvRouterConfig(use_kv_events=True),
),
num_workers=1,
engine_workers=engine_workers,
)
worker_ids = await wait_for_workers_ready(endpoint, router, 1, model_name)
await send_request_via_python_kv_router(
kv_python_router=router,
model_name=model_name,
token_ids=list(range(1, 3 * block_size + 1)),
stop_conditions={"ignore_eos": True, "max_tokens": 2},
worker_id=worker_ids[0],
)
deadline = time.monotonic() + 10
while time.monotonic() < deadline:
events = json.loads(await router.dump_events())
stored_blocks = sum(
len(event["event"]["data"].get("stored", {}).get("blocks", []))
for event in events
)
if stored_blocks >= 3:
return
await asyncio.sleep(0.1)
raise AssertionError(
f"Expected three full KV blocks in the index, got {events}"
)

asyncio.run(run_test())


def _test_router_basic(
engine_workers,
block_size: int,
Expand Down
49 changes: 46 additions & 3 deletions tests/router/test_router_e2e_with_vllm.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@
import aiohttp
import pytest

from tests.router.common import _test_full_block_event_delivery
from tests.router.e2e_harness import (
ManagedEngineProcessMixin,
run_basic_router_test,
Expand All @@ -30,7 +31,7 @@
wait_for_indexer_workers_active,
)
from tests.utils.constants import DynamoPortRange
from tests.utils.gpu_args import build_gpu_mem_args
from tests.utils.gpu_args import build_gpu_mem_args, map_cuda_visible_devices
from tests.utils.managed_process import ManagedProcess
from tests.utils.port_utils import (
allocate_contiguous_ports,
Expand Down Expand Up @@ -225,6 +226,7 @@ def __init__(
num_gpu_blocks_override = vllm_args.get("num_gpu_blocks_override")
max_model_len = vllm_args.get("max_model_len")
enforce_eager = vllm_args.get("enforce_eager", False)
tensor_parallel_size = vllm_args.get("tensor_parallel_size", 1)

self.model_name = model
self.block_size = vllm_args.get("block_size", BLOCK_SIZE)
Expand All @@ -251,10 +253,18 @@ def __init__(
)
)
else:
# No DP; worker sees one GPU
gpu_device = str(gpu_start_index + worker_idx)
worker_start_gpu = gpu_start_index + worker_idx * tensor_parallel_size
gpu_device = map_cuda_visible_devices(
range(worker_start_gpu, worker_start_gpu + tensor_parallel_size),
os.environ.get("CUDA_VISIBLE_DEVICES"),
)

command = ["python3", "-m", "dynamo.vllm", "--model", model]
for name in ("tensor_parallel_size", "decode_context_parallel_size"):
if name in vllm_args:
command.extend(
[f"--{name.replace('_', '-')}", str(vllm_args[name])]
)

if "block_size" in vllm_args:
command.extend(["--block-size", str(vllm_args["block_size"])])
Expand Down Expand Up @@ -584,6 +594,39 @@ def launch_indexer(self):
init_delay_reason = "initialize NIXL before starting next worker"


@pytest.mark.gpu_4
@pytest.mark.h100
@pytest.mark.nightly
@pytest.mark.model("Qwen/Qwen2.5-3B-Instruct")
@pytest.mark.profiled_vram_gib(6.9)
@pytest.mark.requested_vllm_kv_cache_bytes(268_435_456)
@pytest.mark.timeout(600)
@pytest.mark.parametrize("request_plane", ["tcp"], indirect=True)
def test_vllm_dcp2_event_indexing(
request,
runtime_services_dynamic_ports,
predownload_models,
set_ucx_tls_no_mm,
request_plane,
):
model_name = "Qwen/Qwen2.5-3B-Instruct"
with VLLMProcess(
request,
vllm_args={
"model": model_name,
"block_size": 16,
"tensor_parallel_size": 4,
"decode_context_parallel_size": 2,
"kv_cache_memory_bytes": 268_435_456,
"max_model_len": 256,
"enforce_eager": True,
},
num_workers=1,
request_plane=request_plane,
) as engine_workers:
_test_full_block_event_delivery(engine_workers, model_name, 32, request_plane)


@pytest.mark.pre_merge
@pytest.mark.gpu_1
@pytest.mark.profiled_vram_gib(6.9) # actual profiled peak with kv-bytes
Expand Down
Loading