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
47 changes: 47 additions & 0 deletions tests/v1/kv_offload/tiering/test_tiering_offloading.py
Original file line number Diff line number Diff line change
Expand Up @@ -403,6 +403,53 @@ def test_promotion_from_secondary(self, manager_setup):
# Next lookup should succeed
assert count_hits(self.manager, blocks) == 3

@pytest.mark.parametrize(
("successful_indices", "expected_results"),
[
(
(0, 2),
[LookupResult.HIT, LookupResult.MISS, LookupResult.HIT],
),
(
None,
[LookupResult.MISS, LookupResult.MISS, LookupResult.MISS],
),
],
ids=["partial", "legacy-full-failure"],
)
def test_failed_promotion_keeps_only_successful_blocks(
self, manager_setup, successful_indices, expected_results
):
blocks = to_keys(range(3))
for block in blocks:
self.secondary_tier1.blocks[block] = True

def submit_partial(job_metadata: JobMetadata) -> None:
successful_keys = (
None
if successful_indices is None
else tuple(blocks[i] for i in successful_indices)
)
self.secondary_tier1.completed_jobs.append(
JobResult(
job_id=job_metadata.job_id,
success=False,
successful_keys=successful_keys,
)
)

self.secondary_tier1.submit_load = submit_partial

for block in blocks:
assert self.manager.lookup(block, _CTX) is LookupResult.RETRY

self._simulate_on_schedule_end()
self._simulate_on_schedule_end()

assert [
self.primary_tier.lookup(block, _CTX) for block in blocks
] == expected_results

def test_lookup_reports_sync_delay_for_resolved_lookups(self, manager_setup):
"""Resolved lookups report one sync delay sample on allocation."""
self._start_request()
Expand Down
7 changes: 6 additions & 1 deletion vllm/v1/kv_offload/tiering/base.py
Original file line number Diff line number Diff line change
Expand Up @@ -54,10 +54,15 @@ class JobMetadata:

@dataclass
class JobResult:
"""Result of an async transfer job (successful or failed)."""
"""Result of an async transfer job."""

job_id: JobId
# True if all keys succeeded; False if all or some failed.
success: bool
Comment thread
mkhazraee marked this conversation as resolved.
# Only applicable to promotion jobs. On partial failure, identifies the
# keys that were successfully loaded. None means all keys share the fate
# indicated by `success`. Must be a subset of the job's original keys.
successful_keys: Collection[OffloadKey] | None = None


class ParentManager(ABC):
Expand Down
40 changes: 35 additions & 5 deletions vllm/v1/kv_offload/tiering/manager.py
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@
from vllm.v1.kv_offload.tiering.base import (
JobId,
JobMetadata,
JobResult,
ParentManager,
SecondaryTierManager,
TieringOffloadingMetrics,
Expand Down Expand Up @@ -243,6 +244,39 @@ def _maybe_process_finished_jobs(self):
self._processed_jobs_this_step = True
self._process_finished_jobs()

def _complete_promotion(
self, job_metadata: JobMetadata, completed_job: JobResult
) -> None:
successful_keys = completed_job.successful_keys
failed_keys: Collection[OffloadKey]
if completed_job.success:
successful_keys = job_metadata.keys
failed_keys = ()
elif successful_keys:
failed_keys_set = set(job_metadata.keys)
assert failed_keys_set.issuperset(successful_keys), (
f"Finished promotion job_id {completed_job.job_id} "
"reported unknown successful keys"
)
failed_keys_set.difference_update(successful_keys)
failed_keys = failed_keys_set
else:
successful_keys = ()
failed_keys = job_metadata.keys

if successful_keys:
self.primary_tier.complete_write(
successful_keys,
job_metadata.req_context,
True,
)
if failed_keys:
self.primary_tier.complete_write(
failed_keys,
job_metadata.req_context,
False,
)
Comment thread
mkhazraee marked this conversation as resolved.

def _process_finished_jobs(self):
"""
Unconditionally poll all secondary tiers for completed jobs.
Expand All @@ -266,11 +300,7 @@ def _process_finished_jobs(self):
if job_metadata.is_promotion:
# secondary→primary transfer (promotion) completed.
# Make blocks available in primary tier.
self.primary_tier.complete_write(
job_metadata.keys,
job_metadata.req_context,
completed_job.success,
)
self._complete_promotion(job_metadata, completed_job)
else:
# primary→secondary transfer completed.
# Decrement ref_cnt on primary blocks.
Expand Down
Loading