From 1a1971ec6ca184555c716753f9d863df1d87febf Mon Sep 17 00:00:00 2001 From: AlpinDale Date: Fri, 24 Jul 2026 09:50:22 +0000 Subject: [PATCH 1/4] [sync] [Bugfix][KV Offloading] Handle queued request aborts without allocated KV blocks (#49146) Upstream-vLLM: 94ed0bf4e023255d8a0c98da962b6c01ee7dee9e Co-authored-by: Chauncey --- .sync/vllm-sha | 2 +- .../kv_connector/v1/offloading/scheduler.py | 32 +++++++++++++----- .../offloading_connector/test_scheduler.py | 33 +++++++++++++++++++ 3 files changed, 57 insertions(+), 10 deletions(-) diff --git a/.sync/vllm-sha b/.sync/vllm-sha index f30a3bbcde..a20f036d73 100644 --- a/.sync/vllm-sha +++ b/.sync/vllm-sha @@ -1 +1 @@ -1940c8441eb39d7db88764b8e700d37e50573a84 +94ed0bf4e023255d8a0c98da962b6c01ee7dee9e diff --git a/aphrodite/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py b/aphrodite/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py index a61753e93e..0ee7971ae2 100644 --- a/aphrodite/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py +++ b/aphrodite/distributed/kv_transfer/kv_connector/v1/offloading/scheduler.py @@ -291,16 +291,21 @@ def update_block_id_groups(self, new_block_id_groups: tuple[list[int], ...] | No for group_state, new_blocks in zip(self.group_states, new_block_id_groups): group_state.block_ids.extend(new_blocks) - def storable_chunks(self, group_config: "GroupOffloadConfig", num_offloadable_tokens: int) -> int: - """Number of leading offloaded blocks eligible for store. + def storable_chunks( + self, + group_config: "GroupOffloadConfig", + group_state: RequestGroupState, + num_offloadable_tokens: int, + ) -> int: + """Number of allocated leading offloaded chunks eligible for store. - For eagle/MTP groups the volatile trailing block of the offloadable + For eagle/MTP groups the volatile trailing chunk of the offloadable range is excluded while decoding: the draft-layer KV of the last accepted position may be rewritten after spec-token rejection. During - prefill the trailing block is stable (the draft input for a chunk's + prefill the trailing chunk is stable (the draft input for a chunk's last position is the next prompt token), so it is stored immediately. The exclusion must be applied consistently everywhere - ``next_stored_chunk_idx`` is derived: otherwise the trailing block of + ``next_stored_chunk_idx`` is derived: otherwise the trailing chunk of each step is skipped on collection but jumped over by ``next_stored_chunk_idx``, so it is never re-considered and a permanent hole breaks prefix-reuse lookup. @@ -309,7 +314,8 @@ def storable_chunks(self, group_config: "GroupOffloadConfig", num_offloadable_to is_decoding = num_offloadable_tokens > self.req.num_prompt_tokens if group_config.is_eagle_group and is_decoding: num_blocks = max(0, num_blocks - 1) - return num_blocks + num_allocated_chunks = len(group_state.block_ids) // self.config.blocks_per_chunk + return min(num_blocks, num_allocated_chunks) def advance_stored_idx(self, num_offloadable_tokens: int) -> None: # max(): at the prefill->decode transition of a block-aligned prompt, @@ -318,7 +324,7 @@ def advance_stored_idx(self, num_offloadable_tokens: int) -> None: for group_config, group_state in zip(self.config.kv_group_configs, self.group_states): group_state.next_stored_chunk_idx = max( group_state.next_stored_chunk_idx, - self.storable_chunks(group_config, num_offloadable_tokens), + self.storable_chunks(group_config, group_state, num_offloadable_tokens), ) def update_num_hit_chunks(self, num_cached_tokens: int) -> None: @@ -909,7 +915,11 @@ def _build_store_jobs( # or unreachable by the load path's alignment constraints. new_offload_keys: list[OffloadKey] = [] for group_config, group_state in zip(self.config.kv_group_configs, req_status.group_states): - num_blocks = req_status.storable_chunks(group_config, num_offloadable_tokens) + num_blocks = req_status.storable_chunks( + group_config, + group_state, + num_offloadable_tokens, + ) start_block_idx = group_state.next_stored_chunk_idx if num_blocks <= start_block_idx: @@ -973,7 +983,11 @@ def _build_store_jobs( non_sliding_window_block_ids: list[int] = [] for group_config, group_state in zip(self.config.kv_group_configs, req_status.group_states): is_sliding_window = group_config.sliding_window_size_in_chunks is not None - num_blocks = req_status.storable_chunks(group_config, num_offloadable_tokens) + num_blocks = req_status.storable_chunks( + group_config, + group_state, + num_offloadable_tokens, + ) start_block_idx = group_state.next_stored_chunk_idx block_ids = group_state.block_ids num_group_blocks = 0 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 8ed973e193..54c8ade93d 100644 --- a/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py +++ b/tests/v1/kv_connector/unit/offloading_connector/test_scheduler.py @@ -6,6 +6,9 @@ import pytest import torch +from aphrodite.distributed.kv_transfer.kv_connector.v1.offloading.common import ( + OffloadingConnectorMetadata, +) from aphrodite.distributed.kv_transfer.kv_connector.v1.offloading.metrics import ( OffloadingConnectorStats, _ConnectorMetricName, @@ -97,6 +100,36 @@ def test_last_block_offloaded_at_request_finish(request_runner, async_scheduling ) +@pytest.mark.parametrize("async_scheduling", [True, False]) +def test_abort_queued_request_does_not_build_store_job(request_runner, async_scheduling: bool): + """Aborting a never-scheduled request must not store unallocated KV.""" + block_size = 4 + runner = request_runner( + block_size=block_size, + num_gpu_blocks=8, + async_scheduling=async_scheduling, + ) + + runner.new_request(token_ids=[0] * (block_size * 4)) + runner.scheduler.schedule() + + runner.new_request(token_ids=[1] * (block_size * 4)) + queued_req_id = str(runner.req_id) + assert any(request.request_id == queued_req_id for request in runner.scheduler.waiting) + + runner.scheduler.finish_requests(queued_req_id, RequestStatus.FINISHED_ABORTED) + req_status = runner.connector_scheduler._req_status[queued_req_id] + assert all(group_state.offload_keys for group_state in req_status.group_states) + assert all(not group_state.block_ids for group_state in req_status.group_states) + + scheduler_output = runner.scheduler.schedule() + + metadata = scheduler_output.kv_connector_metadata + assert isinstance(metadata, OffloadingConnectorMetadata) + assert all(job.req_id != queued_req_id for job in metadata.store_jobs.values()) + assert queued_req_id not in runner.connector_scheduler._req_status + + def test_scheduler_reports_lookup_sync_delay(request_runner): runner = request_runner( block_size=4, From 9bcb77636f97c6a7f0e49c0f3aa1278320bd6873 Mon Sep 17 00:00:00 2001 From: AlpinDale Date: Fri, 24 Jul 2026 09:53:39 +0000 Subject: [PATCH 2/4] [sync] [MRV2] Add encoder cache profiling implementation (#47985) Upstream-vLLM: ea0e9c8f2e4b037e6fa4d914c82bf178c9f65f55 Co-authored-by: Isotr0py --- .sync/vllm-sha | 2 +- aphrodite/multimodal/encoder_budget.py | 25 +++++++++++- aphrodite/v1/worker/gpu/mm/encoder_cache.py | 2 +- aphrodite/v1/worker/gpu/mm/encoder_runner.py | 43 ++++++++++++++++++++ aphrodite/v1/worker/gpu/model_runner.py | 19 +++++++++ 5 files changed, 88 insertions(+), 3 deletions(-) diff --git a/.sync/vllm-sha b/.sync/vllm-sha index a20f036d73..537b145a54 100644 --- a/.sync/vllm-sha +++ b/.sync/vllm-sha @@ -1 +1 @@ -94ed0bf4e023255d8a0c98da962b6c01ee7dee9e +ea0e9c8f2e4b037e6fa4d914c82bf178c9f65f55 diff --git a/aphrodite/multimodal/encoder_budget.py b/aphrodite/multimodal/encoder_budget.py index e9fc8743bb..655072278d 100644 --- a/aphrodite/multimodal/encoder_budget.py +++ b/aphrodite/multimodal/encoder_budget.py @@ -4,6 +4,7 @@ from aphrodite.config import AphroditeConfig, ModelConfig from aphrodite.logger import init_logger +from aphrodite.multimodal.inputs import MultiModalKwargsItem from aphrodite.multimodal.processing import BaseMultiModalProcessor from aphrodite.multimodal.registry import MultiModalRegistry from aphrodite.utils.torch_utils import set_default_torch_num_threads @@ -48,6 +49,8 @@ def __init__( self, aphrodite_config: AphroditeConfig, mm_registry: MultiModalRegistry, + *, + enable_cache: bool = True, ) -> None: super().__init__() @@ -58,7 +61,7 @@ def __init__( self.max_num_reqs = scheduler_config.max_num_seqs with set_default_torch_num_threads(): # Avoid hang during startup - cache = mm_registry.processor_only_cache_from_config(aphrodite_config) + cache = mm_registry.processor_only_cache_from_config(aphrodite_config) if enable_cache else None processor = mm_registry.create_processor(model_config, cache=cache) self.cache = cache @@ -185,3 +188,23 @@ def get_encoder_budget(self) -> int: def reset_cache(self) -> None: if self.cache is not None: self.cache.clear_cache() + + +def get_dummy_encoder_profile_inputs( + mm_registry: MultiModalRegistry, + budget: MultiModalBudget, +) -> list[tuple[str, MultiModalKwargsItem]]: + if budget.get_encoder_budget() <= 0 or not budget.mm_max_toks_per_item: + return [] + + modality = budget.get_modality_with_max_tokens() + max_items_per_batch = budget.mm_max_items_per_batch[modality] + dummy_mm_inputs = mm_registry.get_dummy_mm_inputs( + budget.model_config, + mm_counts={modality: 1}, + processor=budget.processor, + ) + dummy_mm_item = dummy_mm_inputs["mm_kwargs"][modality][0] + assert dummy_mm_item is not None, "Dummy item should be generated" + + return [(modality, dummy_mm_item)] * max_items_per_batch diff --git a/aphrodite/v1/worker/gpu/mm/encoder_cache.py b/aphrodite/v1/worker/gpu/mm/encoder_cache.py index bb349c25a4..f48e659a1b 100644 --- a/aphrodite/v1/worker/gpu/mm/encoder_cache.py +++ b/aphrodite/v1/worker/gpu/mm/encoder_cache.py @@ -26,7 +26,7 @@ def reset_mm_cache(self) -> None: Clear the multi-modal cache that was used during profiling, but no longer needed during inference. """ - # TODO: Implement MM budget for encoder dummy run + # NOTE: v2 encoder cache profiling skips the multi-modal cache pass def reset_encoder_cache(self) -> None: diff --git a/aphrodite/v1/worker/gpu/mm/encoder_runner.py b/aphrodite/v1/worker/gpu/mm/encoder_runner.py index ef31c8728b..7454343c12 100644 --- a/aphrodite/v1/worker/gpu/mm/encoder_runner.py +++ b/aphrodite/v1/worker/gpu/mm/encoder_runner.py @@ -3,7 +3,9 @@ import numpy as np import torch +from aphrodite.logger import init_logger from aphrodite.model_executor.models.interfaces import SupportsMultiModal, supports_realtime +from aphrodite.multimodal.encoder_budget import MultiModalBudget from aphrodite.multimodal.inputs import MultiModalKwargsItem from aphrodite.multimodal.utils import ( get_mm_features_in_window, @@ -13,6 +15,8 @@ from aphrodite.v1.worker.gpu.mm.encoder_cache import EncoderCache from aphrodite.v1.worker.utils import sanity_check_mm_encoder_outputs +logger = init_logger(__name__) + class EncoderRunner: def __init__( @@ -50,6 +54,45 @@ def prepare_mm_inputs( return mm_hashes, mm_kwargs + @torch.inference_mode() + def profile_encoder_cache( + self, + dummy_mm_inputs: list[tuple[str, MultiModalKwargsItem]], + budget: MultiModalBudget, + ) -> None: + """Profile multimodal encoder and temporary encoder cache memory.""" + if (encoder_budget := budget.get_encoder_budget()) <= 0: + return + + if not budget.mm_max_toks_per_item: + logger.info( + "Skipping encoder profiling for embedding-only mode " + "(all modality limits=0 with enable_mm_embeds=True).", + ) + return + + assert dummy_mm_inputs, "Dummy inputs should be generated for encoder profiling" + dummy_modality = dummy_mm_inputs[0][0] + max_mm_items_per_batch = len(dummy_mm_inputs) + + logger.info_once( + "Encoder cache will be initialized with a budget of %s tokens, " + "and profiled with %s %s items of the maximum feature size.", + encoder_budget, + max_mm_items_per_batch, + dummy_modality, + ) + + dummy_encoder_outputs = self.execute_mm_encoder(dummy_mm_inputs) + + sanity_check_mm_encoder_outputs( + dummy_encoder_outputs, + expected_num_items=max_mm_items_per_batch, + ) + self.encoder_cache.encoder_outputs.update( + (f"tmp_{i}", output) for i, output in enumerate(dummy_encoder_outputs) + ) + @torch.inference_mode() def execute_mm_encoder(self, mm_kwargs: list[tuple[str, MultiModalKwargsItem]]) -> list[torch.Tensor]: encoder_outputs: list[torch.Tensor] = [] diff --git a/aphrodite/v1/worker/gpu/model_runner.py b/aphrodite/v1/worker/gpu/model_runner.py index d6058673ab..98909527d5 100644 --- a/aphrodite/v1/worker/gpu/model_runner.py +++ b/aphrodite/v1/worker/gpu/model_runner.py @@ -43,6 +43,10 @@ ) from aphrodite.model_executor.model_loader import get_model_loader from aphrodite.multimodal import MULTIMODAL_REGISTRY +from aphrodite.multimodal.encoder_budget import ( + MultiModalBudget, + get_dummy_encoder_profile_inputs, +) from aphrodite.sequence import IntermediateTensors from aphrodite.tasks import SupportedTask from aphrodite.utils.math_utils import cdiv @@ -641,6 +645,20 @@ def _dummy_pooler_run(self, hidden_states: torch.Tensor) -> None: @torch.inference_mode() def profile_run(self) -> None: + if self.supports_mm_inputs and self.is_first_pp_rank: + mm_config = self.model_config.multimodal_config + if mm_config is not None and not mm_config.skip_mm_profiling: + mm_budget = MultiModalBudget( + self.aphrodite_config, + self.mm_registry, + enable_cache=False, + ) + dummy_mm_inputs = get_dummy_encoder_profile_inputs( + self.mm_registry, + mm_budget, + ) + self.model_state.encoder_runner.profile_encoder_cache(dummy_mm_inputs, mm_budget) + hidden_states, sample_hidden_states = self._dummy_run(self.max_num_tokens, skip_attn=True, is_profile=True) # Only run sampler/pooler on last PP rank (non-last ranks return None). @@ -653,6 +671,7 @@ def profile_run(self) -> None: torch.accelerator.synchronize() del hidden_states, sample_hidden_states + self.reset_encoder_cache() gc.collect() def post_kv_cache_wake_up(self) -> None: From d8926d2b33d1fd47c99d7d696c791201e80352d4 Mon Sep 17 00:00:00 2001 From: AlpinDale Date: Fri, 24 Jul 2026 09:55:28 +0000 Subject: [PATCH 3/4] [sync] [CI][NIXL] Isolate concurrent engine internal ports (#49129) Upstream-vLLM: 6bcda970fde272b06cff0d8f397224dbac96d7ad Co-authored-by: Andreas Karatzas --- .sync/vllm-sha | 2 +- .../nixl_integration/run_accuracy_test.sh | 13 +++++++++++++ 2 files changed, 14 insertions(+), 1 deletion(-) diff --git a/.sync/vllm-sha b/.sync/vllm-sha index 537b145a54..2f7a9fe3d9 100644 --- a/.sync/vllm-sha +++ b/.sync/vllm-sha @@ -1 +1 @@ -ea0e9c8f2e4b037e6fa4d914c82bf178c9f65f55 +6bcda970fde272b06cff0d8f397224dbac96d7ad diff --git a/tests/v1/kv_connector/nixl_integration/run_accuracy_test.sh b/tests/v1/kv_connector/nixl_integration/run_accuracy_test.sh index ea283b7100..26866e351d 100755 --- a/tests/v1/kv_connector/nixl_integration/run_accuracy_test.sh +++ b/tests/v1/kv_connector/nixl_integration/run_accuracy_test.sh @@ -83,6 +83,11 @@ DECODE_BLOCK_SIZE=${DECODE_BLOCK_SIZE:-128} ENFORCE_EAGER=${ENFORCE_EAGER:-1} # Comma-separated extra args for aphrodite serve (e.g. --max-model-len,2048) APHRODITE_SERVE_EXTRA_ARGS=${APHRODITE_SERVE_EXTRA_ARGS:-} +# Pin concurrent prefiller and non-DP decoder engines to separate internal +# port windows. DP decoder ranks retain their existing internal port selection. +PREFILLER_INTERNAL_PORT_BASE=${PREFILLER_INTERNAL_PORT_BASE:-20000} +DECODER_INTERNAL_PORT_BASE=${DECODER_INTERNAL_PORT_BASE:-30000} +INTERNAL_PORT_STRIDE=${INTERNAL_PORT_STRIDE:-100} # Resolve the repository root from the script location instead of `.git`. # The ROCm CI image copies `/aphrodite-workspace` without the Git metadata, so @@ -154,12 +159,14 @@ run_tests_for_model() { PORT=$((8100 + i)) # Calculate side channel port. Avoid clash with with TP workers. SIDE_CHANNEL_PORT=$((5559 + i)) + INTERNAL_PORT=$((PREFILLER_INTERNAL_PORT_BASE + i * INTERNAL_PORT_STRIDE)) echo "Starting prefill instance $i on GPU $GPU_ID, port $PORT" # Build the command with or without model-specific args BASE_CMD="CUDA_VISIBLE_DEVICES=$GPU_ID \ APHRODITE_KV_CACHE_LAYOUT='HND' \ + APHRODITE_PORT=$INTERNAL_PORT \ UCX_NET_DEVICES=all \ APHRODITE_NIXL_SIDE_CHANNEL_PORT=$SIDE_CHANNEL_PORT \ aphrodite serve $model_name \ @@ -208,12 +215,18 @@ run_tests_for_model() { PORT=$((8200 + i)) # Calculate side channel port SIDE_CHANNEL_PORT=$((5659 + i * $DECODER_TP_SIZE)) + INTERNAL_PORT=$((DECODER_INTERNAL_PORT_BASE + i * INTERNAL_PORT_STRIDE)) + DECODER_INTERNAL_PORT_ENV= + if [[ -z "${DP_EP:-}" ]]; then + DECODER_INTERNAL_PORT_ENV="APHRODITE_PORT=$INTERNAL_PORT" + fi echo "Starting decode instance $i on GPU $GPU_ID, port $PORT" # Build the command with or without model-specific args BASE_CMD="CUDA_VISIBLE_DEVICES=$GPU_ID \ APHRODITE_KV_CACHE_LAYOUT=$DECODER_KV_LAYOUT \ + $DECODER_INTERNAL_PORT_ENV \ UCX_NET_DEVICES=all \ APHRODITE_NIXL_SIDE_CHANNEL_PORT=$SIDE_CHANNEL_PORT \ aphrodite serve $model_name \ From c6dbcdcef3fc51384a50ecfa9eb35240d272f8da Mon Sep 17 00:00:00 2001 From: AlpinDale Date: Fri, 24 Jul 2026 09:57:08 +0000 Subject: [PATCH 4/4] [sync] Update BGE-M3 token expectations for leading spaces (#49269) Upstream-vLLM: d9aa35161d569608aed98465db0b0535b1431c28 Co-authored-by: aoshen02 --- .sync/vllm-sha | 2 +- .../test_bge_m3_sparse_io_processor_plugins.py | 12 ++++++------ 2 files changed, 7 insertions(+), 7 deletions(-) diff --git a/.sync/vllm-sha b/.sync/vllm-sha index 2f7a9fe3d9..14df7b2aac 100644 --- a/.sync/vllm-sha +++ b/.sync/vllm-sha @@ -1 +1 @@ -6bcda970fde272b06cff0d8f397224dbac96d7ad +d9aa35161d569608aed98465db0b0535b1431c28 diff --git a/tests/plugins_tests/test_bge_m3_sparse_io_processor_plugins.py b/tests/plugins_tests/test_bge_m3_sparse_io_processor_plugins.py index a22edfcd12..88582647f7 100644 --- a/tests/plugins_tests/test_bge_m3_sparse_io_processor_plugins.py +++ b/tests/plugins_tests/test_bge_m3_sparse_io_processor_plugins.py @@ -43,12 +43,12 @@ def _check_dense_embedding(data, index=0): def _check_sparse_embedding(data, check_tokens=False): expected_weights = [ {"token_id": 32, "weight": 0.0552978515625, "token": "?"}, - {"token_id": 70, "weight": 0.09808349609375, "token": "the"}, - {"token_id": 83, "weight": 0.08154296875, "token": "is"}, - {"token_id": 111, "weight": 0.11810302734375, "token": "of"}, - {"token_id": 4865, "weight": 0.1171875, "token": "What"}, - {"token_id": 9942, "weight": 0.292236328125, "token": "France"}, - {"token_id": 10323, "weight": 0.2802734375, "token": "capital"}, + {"token_id": 70, "weight": 0.09808349609375, "token": " the"}, + {"token_id": 83, "weight": 0.08154296875, "token": " is"}, + {"token_id": 111, "weight": 0.11810302734375, "token": " of"}, + {"token_id": 4865, "weight": 0.1171875, "token": " What"}, + {"token_id": 9942, "weight": 0.292236328125, "token": " France"}, + {"token_id": 10323, "weight": 0.2802734375, "token": " capital"}, ] expected_embed = {x["token_id"]: x for x in expected_weights}