Skip to content
Open
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
Original file line number Diff line number Diff line change
Expand Up @@ -3151,7 +3151,12 @@ def get_index_k_buffer(
return full_view[:, 0]

def get_num_available_tokens(
self, *, token_num_upper_bound: int, batch_size: int = 1, max_num_draft_tokens: int = 0
self,
*,
token_num_upper_bound: int,
batch_size: int = 1,
max_num_draft_tokens: int = 0,
max_beam_width: int = 1,
) -> int:
"""Clamp ``token_num_upper_bound`` to the allocatable token capacity.

Expand All @@ -3163,6 +3168,9 @@ def get_num_available_tokens(
``max_num_tokens``) stay consistent because a helix context forward
replicates all tokens on every rank, so both bounds constrain the
same request-length variable.

``max_beam_width`` is accepted for interface parity with the V1
manager; V2 only supports a beam width of 1.
"""
extra_tokens = self.num_extra_kv_tokens + max_num_draft_tokens
# Token num upper bound is the maximum number of tokens that can be allocated in the kv cache manager.
Expand Down
13 changes: 11 additions & 2 deletions tensorrt_llm/_torch/pyexecutor/model_engine.py
Original file line number Diff line number Diff line change
Expand Up @@ -3335,16 +3335,25 @@ def free_warmup_requests() -> None:
available_tokens = kv_cache_manager.get_num_available_tokens(
token_num_upper_bound=max_seq_len,
batch_size=batch_size,
max_num_draft_tokens=_kv_draft)
max_num_draft_tokens=_kv_draft,
max_beam_width=self.max_beam_width)

# Also consider draft KV cache capacity when it exists
if draft_kv_cache_manager is not None:
draft_available_tokens = draft_kv_cache_manager.get_num_available_tokens(
batch_size=batch_size,
token_num_upper_bound=max_seq_len,
max_num_draft_tokens=_kv_draft)
max_num_draft_tokens=_kv_draft,
max_beam_width=self.max_beam_width)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
available_tokens = min(available_tokens, draft_available_tokens)

# With beam search the free blocks may not cover even one token per
# beam. add_dummy_requests cannot always catch this (it does not size
# VSWA pools per beam), so skip this batch size here.
if self.max_beam_width > 1 and available_tokens < 1:
free_warmup_requests()
return None

token_num = max(
ENC_DEC_CUDA_GRAPH_DUMMY_TOKEN_NUM if is_enc_dec else 1,
min(
Expand Down
69 changes: 66 additions & 3 deletions tensorrt_llm/_torch/pyexecutor/resource_manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -1347,6 +1347,25 @@ def add_dummy_requests(
_populate_dummy_mrope_config(req, token_num, is_gen)
requests.append(req)

# Beam search allocates most blocks once per beam, so the single-block
# check above does not guarantee the dummy requests fit. Skip padding
# instead of failing inside the block manager. VSWA pools are sized
# per window, which this full-attention count does not model.
if beam_width > 1 and batch_request_infos and not self.is_vswa:
num_appended_tokens = self.num_extra_kv_tokens + (_kv_draft
if is_gen else 0)
num_required_blocks = sum(
self._get_num_blocks_for_dummy_request(
token_num, num_appended_tokens, beam_width)
for _, token_num, _ in batch_request_infos)
if num_required_blocks > available_blocks:
logger.debug(
f"[add_dummy_requests] {len(batch_request_infos)} dummy "
f"requests with beam_width={beam_width} need "
f"{num_required_blocks} blocks, only {available_blocks} "
f"free; skipping.")
return None

try:
# Use add_sequence_batch for all dummy requests, then add extra tokens.
# This must happen before is_gen state modifications below, which may
Expand Down Expand Up @@ -1839,18 +1858,62 @@ def get_num_kv_blocks(self, num_tokens: int) -> int:
def get_num_available_tokens(self,
token_num_upper_bound: int,
max_num_draft_tokens: int = 0,
max_beam_width: int = 1,
**kwargs) -> int:
"""Return a token count such that one sequence of any length up to it
fits in the free blocks.

Args:
token_num_upper_bound: Upper bound on the returned token count.
max_num_draft_tokens: Draft tokens appended after the sequence.
max_beam_width: Beam width of the sequence. With beam search, only
blocks fully covered by the prompt are shared among beams; the
rest are allocated once per beam.
"""
free_blocks = self.get_num_free_blocks()
result = min(
token_num_upper_bound, free_blocks * self.tokens_per_block -
self.num_extra_kv_tokens - max_num_draft_tokens)
num_appended_tokens = self.num_extra_kv_tokens + max_num_draft_tokens
if max_beam_width > 1 and self.kv_cache_type != CacheTypeCpp.CROSS:
# Block usage is not monotonic in the sequence length (a
# block-aligned prompt shares all of its blocks), so bound it by
# the worst case: a partially filled last prompt block followed by
# the appended tokens, all allocated per beam.
max_blocks_per_beam = math.ceil(
(self.tokens_per_block - 1 + num_appended_tokens) /
self.tokens_per_block)
num_shared_blocks = free_blocks - max_beam_width * max_blocks_per_beam
capacity = (num_shared_blocks + 1) * self.tokens_per_block - 1
Comment thread
coderabbitai[bot] marked this conversation as resolved.
else:
Comment thread
coderabbitai[bot] marked this conversation as resolved.
capacity = free_blocks * self.tokens_per_block - num_appended_tokens
result = min(token_num_upper_bound, capacity)
logger.debug(
f"[get_num_available_tokens] free_blocks={free_blocks}, "
f"tokens_per_block={self.tokens_per_block}, "
f"num_extra_kv_tokens={self.num_extra_kv_tokens}, "
f"max_beam_width={max_beam_width}, "
f"token_num_upper_bound={token_num_upper_bound}, result={result}")
return result

def _get_num_blocks_for_dummy_request(self, token_num: int,
num_appended_tokens: int,
beam_width: int) -> int:
"""Number of blocks ``add_dummy_requests`` allocates for one sequence
of ``token_num`` prompt tokens followed by ``num_appended_tokens``
tokens added one at a time.

Blocks fully covered by the prompt are shared among beams (for cross
KV, the partial last prompt block is shared too); every other block is
allocated once per beam.
"""
num_blocks = math.ceil(
(token_num + num_appended_tokens) / self.tokens_per_block)
if beam_width == 1:
return num_blocks
if self.kv_cache_type == CacheTypeCpp.CROSS:
num_shared_blocks = math.ceil(token_num / self.tokens_per_block)
else:
num_shared_blocks = token_num // self.tokens_per_block
return num_shared_blocks + beam_width * (num_blocks - num_shared_blocks)

def get_buffers(self,
layer_idx: int,
kv_layout: str = "NHD") -> Optional[torch.Tensor]:
Expand Down
114 changes: 113 additions & 1 deletion tests/unittest/_torch/executor/test_pytorch_model_engine_warmup.py
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,11 @@
from tensorrt_llm._torch.pyexecutor.engine.runners.no_kv_cache import NoKVCacheRunner
from tensorrt_llm._torch.pyexecutor.kv_cache.kv_cache_manager_v2 import KVCacheManagerV2
from tensorrt_llm._torch.pyexecutor.model_engine import PyTorchModelEngine
from tensorrt_llm._torch.pyexecutor.resource_manager import ResourceManager, ResourceManagerType
from tensorrt_llm._torch.pyexecutor.resource_manager import (
KVCacheManager,
ResourceManager,
ResourceManagerType,
)
from tensorrt_llm._torch.speculative.utils import update_draft_len
from tensorrt_llm.llmapi import CudaGraphConfig, KvCacheConfig
from tensorrt_llm.llmapi.llm_args import (
Expand Down Expand Up @@ -200,6 +204,114 @@ def test_warmup_builders_resynchronize_stale_draft_length(
)


_BEAM_WIDTH = 4
_BEAM_BATCH_SIZE = 4
_BEAM_MAX_SEQ_LEN = 32
_BEAM_TOKENS_PER_BLOCK = 8


def _build_beam_search_engine(max_tokens: int) -> PyTorchModelEngine:
llm_args = TorchLlmArgs(
model="dummy",
max_batch_size=_BEAM_BATCH_SIZE,
max_num_tokens=_BEAM_BATCH_SIZE * _BEAM_MAX_SEQ_LEN,
max_seq_len=_BEAM_MAX_SEQ_LEN,
max_beam_width=_BEAM_WIDTH,
kv_cache_config=KvCacheConfig(max_tokens=max_tokens, enable_block_reuse=False),
cuda_graph_config=CudaGraphConfig(batch_sizes=[1, 2, 4]),
)
return _DummyModelEngine(llm_args, torch.half)


def _create_beam_search_kv_cache_manager(
engine: PyTorchModelEngine, max_tokens: int
) -> KVCacheManager:
return KVCacheManager(
KvCacheConfig(max_tokens=max_tokens, enable_block_reuse=False),
tensorrt_llm.bindings.internal.batch_manager.CacheType.SELF,
num_layers=1,
num_kv_heads=engine.model.config.num_key_value_heads,
head_dim=engine.model.config.head_dim,
tokens_per_block=_BEAM_TOKENS_PER_BLOCK,
max_seq_len=_BEAM_MAX_SEQ_LEN,
max_batch_size=_BEAM_BATCH_SIZE,
max_beam_width=_BEAM_WIDTH,
mapping=Mapping(world_size=1, tp_size=1, rank=0),
dtype=tensorrt_llm.bindings.DataType.HALF,
)


def test_cuda_graph_warmup_request_fits_beam_search_kv_blocks() -> None:
"""The max-length warmup request must fit the blocks that beam search
allocates per beam instead of failing inside the block manager."""
# 16 blocks. The three one-token warmup requests take one block per beam
# (12), so the max-length request gets 4 blocks: one tail block per beam.
engine = _build_beam_search_engine(max_tokens=128)
kv_cache_manager = _create_beam_search_kv_cache_manager(engine, max_tokens=128)
resource_manager = ResourceManager({ResourceManagerType.KV_CACHE_MANAGER: kv_cache_manager})
try:
total_free = kv_cache_manager.get_num_free_blocks()
assert total_free == 16
warmup_request = engine._create_cuda_graph_warmup_request(
resource_manager, _BEAM_BATCH_SIZE, 0
)
with engine._release_batch_context(warmup_request, resource_manager) as batch:
assert batch is not None
assert len(batch.generation_requests) == _BEAM_BATCH_SIZE
assert kv_cache_manager.get_num_free_blocks() == 0
assert kv_cache_manager.get_num_free_blocks() == total_free
finally:
kv_cache_manager.shutdown()


def test_cuda_graph_warmup_request_passes_beam_width_to_kv_cache_managers() -> None:
"""Both the target and the draft KV cache capacity must be computed for
the configured beam width."""
engine = _build_beam_search_engine(max_tokens=1024)
target = _create_beam_search_kv_cache_manager(engine, max_tokens=1024)
draft = _create_beam_search_kv_cache_manager(engine, max_tokens=1024)
target.get_num_available_tokens = Mock(wraps=target.get_num_available_tokens)
draft.get_num_available_tokens = Mock(wraps=draft.get_num_available_tokens)
engine._get_draft_kv_cache_manager = lambda resource_manager: draft
resource_manager = ResourceManager({ResourceManagerType.KV_CACHE_MANAGER: target})
try:
managers = (target, draft)
total_free = [manager.get_num_free_blocks() for manager in managers]
warmup_request = engine._create_cuda_graph_warmup_request(
resource_manager, _BEAM_BATCH_SIZE, 0
)
with engine._release_batch_context(warmup_request, resource_manager) as batch:
assert batch is not None
for manager, free_blocks in zip(managers, total_free):
manager.get_num_available_tokens.assert_called_once()
kwargs = manager.get_num_available_tokens.call_args.kwargs
assert kwargs["max_beam_width"] == _BEAM_WIDTH
assert kwargs["batch_size"] == _BEAM_BATCH_SIZE
assert manager.get_num_free_blocks() == free_blocks
finally:
target.shutdown()
draft.shutdown()


def test_cuda_graph_warmup_request_skips_beam_search_batch_below_one_token() -> None:
"""When the capacity reported for beam search is below one token, the
batch size is skipped and the warmup requests already added are freed."""
# Stands in for a manager whose pool cannot hold one more token per beam
# but whose add_dummy_requests cannot detect it (VSWA pools).
engine = _build_beam_search_engine(max_tokens=1024)
kv_cache_manager = _create_beam_search_kv_cache_manager(engine, max_tokens=1024)
kv_cache_manager.get_num_available_tokens = Mock(return_value=0)
resource_manager = ResourceManager({ResourceManagerType.KV_CACHE_MANAGER: kv_cache_manager})
try:
total_free = kv_cache_manager.get_num_free_blocks()
assert (
engine._create_cuda_graph_warmup_request(resource_manager, _BEAM_BATCH_SIZE, 0) is None
)
assert kv_cache_manager.get_num_free_blocks() == total_free
finally:
kv_cache_manager.shutdown()


class _Tracker:
"""Records method-call order via mock side_effects."""

Expand Down
Loading