From c5683358238664de1d81b0d55a7c15ca28b3a6b4 Mon Sep 17 00:00:00 2001 From: Moein Khazraee Date: Fri, 24 Jul 2026 00:55:15 -0700 Subject: [PATCH 1/4] Add partial submit_load results for Secondary tiers Signed-off-by: Moein Khazraee --- .../tiering/test_tiering_offloading.py | 42 +++++++++++++++++++ vllm/v1/kv_offload/tiering/base.py | 8 +++- vllm/v1/kv_offload/tiering/manager.py | 31 +++++++++++--- 3 files changed, 75 insertions(+), 6 deletions(-) diff --git a/tests/v1/kv_offload/tiering/test_tiering_offloading.py b/tests/v1/kv_offload/tiering/test_tiering_offloading.py index b19a270efb46..64b7106ce1ad 100644 --- a/tests/v1/kv_offload/tiering/test_tiering_offloading.py +++ b/tests/v1/kv_offload/tiering/test_tiering_offloading.py @@ -403,6 +403,48 @@ 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], + ), + ( + (), + [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: + self.secondary_tier1.completed_jobs.append( + JobResult( + job_id=job_metadata.job_id, + success=False, + successful_keys=tuple(blocks[i] for i in successful_indices), + ) + ) + + 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() diff --git a/vllm/v1/kv_offload/tiering/base.py b/vllm/v1/kv_offload/tiering/base.py index 3577b4fa0ed5..c53a22267c56 100644 --- a/vllm/v1/kv_offload/tiering/base.py +++ b/vllm/v1/kv_offload/tiering/base.py @@ -54,10 +54,16 @@ class JobMetadata: @dataclass class JobResult: - """Result of an async transfer job (successful or failed).""" + """Result of an async transfer job. + + For load jobs, ``successful_keys`` identifies completed keys when + ``success`` is False. An empty collection preserves the legacy + all-or-nothing failure behavior. + """ job_id: JobId success: bool + successful_keys: Collection[OffloadKey] = () class ParentManager(ABC): diff --git a/vllm/v1/kv_offload/tiering/manager.py b/vllm/v1/kv_offload/tiering/manager.py index d1d3c4211590..f4b9cf69af50 100644 --- a/vllm/v1/kv_offload/tiering/manager.py +++ b/vllm/v1/kv_offload/tiering/manager.py @@ -266,11 +266,32 @@ 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, - ) + if completed_job.success: + self.primary_tier.complete_write( + job_metadata.keys, + job_metadata.req_context, + True, + ) + else: + successful_keys = completed_job.successful_keys + failed_keys = set(job_metadata.keys) + assert failed_keys.issuperset(successful_keys), ( + f"Finished job_id {job_id} from tier #{i}" + f" ({tier.tier_type}) reported unknown successful keys" + ) + failed_keys.difference_update(successful_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, + ) else: # primary→secondary transfer completed. # Decrement ref_cnt on primary blocks. From 1807c8b3cbff772db771b07074a54385787a194a Mon Sep 17 00:00:00 2001 From: Moein Khazraee Date: Thu, 30 Jul 2026 12:16:31 -0700 Subject: [PATCH 2/4] Addressed reviewer comments. Signed-off-by: Moein Khazraee --- .../tiering/test_tiering_offloading.py | 9 ++- vllm/v1/kv_offload/tiering/base.py | 11 ++-- vllm/v1/kv_offload/tiering/manager.py | 66 +++++++++++-------- 3 files changed, 51 insertions(+), 35 deletions(-) diff --git a/tests/v1/kv_offload/tiering/test_tiering_offloading.py b/tests/v1/kv_offload/tiering/test_tiering_offloading.py index 64b7106ce1ad..78be7b5f4228 100644 --- a/tests/v1/kv_offload/tiering/test_tiering_offloading.py +++ b/tests/v1/kv_offload/tiering/test_tiering_offloading.py @@ -411,7 +411,7 @@ def test_promotion_from_secondary(self, manager_setup): [LookupResult.HIT, LookupResult.MISS, LookupResult.HIT], ), ( - (), + None, [LookupResult.MISS, LookupResult.MISS, LookupResult.MISS], ), ], @@ -425,11 +425,16 @@ def test_failed_promotion_keeps_only_successful_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=tuple(blocks[i] for i in successful_indices), + successful_keys=successful_keys, ) ) diff --git a/vllm/v1/kv_offload/tiering/base.py b/vllm/v1/kv_offload/tiering/base.py index c53a22267c56..f82037eb6b99 100644 --- a/vllm/v1/kv_offload/tiering/base.py +++ b/vllm/v1/kv_offload/tiering/base.py @@ -54,16 +54,13 @@ class JobMetadata: @dataclass class JobResult: - """Result of an async transfer job. - - For load jobs, ``successful_keys`` identifies completed keys when - ``success`` is False. An empty collection preserves the legacy - all-or-nothing failure behavior. - """ + """Result of an async transfer job.""" job_id: JobId success: bool - successful_keys: Collection[OffloadKey] = () + # Only applicable to promotion jobs. On partial failure, identifies the + # keys that were successfully loaded. + successful_keys: Collection[OffloadKey] | None = None class ParentManager(ABC): diff --git a/vllm/v1/kv_offload/tiering/manager.py b/vllm/v1/kv_offload/tiering/manager.py index f4b9cf69af50..9ceb76f0ff6f 100644 --- a/vllm/v1/kv_offload/tiering/manager.py +++ b/vllm/v1/kv_offload/tiering/manager.py @@ -49,6 +49,7 @@ from vllm.v1.kv_offload.tiering.base import ( JobId, JobMetadata, + JobResult, ParentManager, SecondaryTierManager, TieringOffloadingMetrics, @@ -243,6 +244,44 @@ 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: + if completed_job.success: + self.primary_tier.complete_write( + job_metadata.keys, + job_metadata.req_context, + True, + ) + return + + successful_keys = completed_job.successful_keys + if not successful_keys: + self.primary_tier.complete_write( + job_metadata.keys, + job_metadata.req_context, + False, + ) + return + + failed_keys = set(job_metadata.keys) + assert failed_keys.issuperset(successful_keys), ( + f"Finished promotion job_id {completed_job.job_id} " + "reported unknown successful keys" + ) + failed_keys.difference_update(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, + ) + def _process_finished_jobs(self): """ Unconditionally poll all secondary tiers for completed jobs. @@ -266,32 +305,7 @@ def _process_finished_jobs(self): if job_metadata.is_promotion: # secondary→primary transfer (promotion) completed. # Make blocks available in primary tier. - if completed_job.success: - self.primary_tier.complete_write( - job_metadata.keys, - job_metadata.req_context, - True, - ) - else: - successful_keys = completed_job.successful_keys - failed_keys = set(job_metadata.keys) - assert failed_keys.issuperset(successful_keys), ( - f"Finished job_id {job_id} from tier #{i}" - f" ({tier.tier_type}) reported unknown successful keys" - ) - failed_keys.difference_update(successful_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, - ) + self._complete_promotion(job_metadata, completed_job) else: # primary→secondary transfer completed. # Decrement ref_cnt on primary blocks. From 9cee7b01d8509fe87549c1a6c95352a3c82f1955 Mon Sep 17 00:00:00 2001 From: Moein Khazraee <33970824+mkhazraee@users.noreply.github.com> Date: Mon, 3 Aug 2026 11:31:55 -0700 Subject: [PATCH 3/4] Apply clarification comment suggestions from code review Co-authored-by: Or Ozeri Signed-off-by: Moein Khazraee --- vllm/v1/kv_offload/tiering/base.py | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/vllm/v1/kv_offload/tiering/base.py b/vllm/v1/kv_offload/tiering/base.py index f82037eb6b99..2075b556fbaf 100644 --- a/vllm/v1/kv_offload/tiering/base.py +++ b/vllm/v1/kv_offload/tiering/base.py @@ -57,9 +57,11 @@ class JobResult: """Result of an async transfer job.""" job_id: JobId + # True if all keys succeeded; False if all or some failed. success: bool # Only applicable to promotion jobs. On partial failure, identifies the - # keys that were successfully loaded. + # 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 From f07f0f6e0625c6c4ef6881de6d8351a72cde3306 Mon Sep 17 00:00:00 2001 From: Moein Khazraee Date: Mon, 3 Aug 2026 12:05:43 -0700 Subject: [PATCH 4/4] Style clean up suggested by Varun Signed-off-by: Moein Khazraee --- vllm/v1/kv_offload/tiering/manager.py | 39 ++++++++++++--------------- 1 file changed, 17 insertions(+), 22 deletions(-) diff --git a/vllm/v1/kv_offload/tiering/manager.py b/vllm/v1/kv_offload/tiering/manager.py index 9ceb76f0ff6f..5b8f92a45845 100644 --- a/vllm/v1/kv_offload/tiering/manager.py +++ b/vllm/v1/kv_offload/tiering/manager.py @@ -247,34 +247,29 @@ def _maybe_process_finished_jobs(self): 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: - self.primary_tier.complete_write( - job_metadata.keys, - job_metadata.req_context, - True, + 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" ) - return + failed_keys_set.difference_update(successful_keys) + failed_keys = failed_keys_set + else: + successful_keys = () + failed_keys = job_metadata.keys - successful_keys = completed_job.successful_keys - if not successful_keys: + if successful_keys: self.primary_tier.complete_write( - job_metadata.keys, + successful_keys, job_metadata.req_context, - False, + True, ) - return - - failed_keys = set(job_metadata.keys) - assert failed_keys.issuperset(successful_keys), ( - f"Finished promotion job_id {completed_job.job_id} " - "reported unknown successful keys" - ) - failed_keys.difference_update(successful_keys) - self.primary_tier.complete_write( - successful_keys, - job_metadata.req_context, - True, - ) if failed_keys: self.primary_tier.complete_write( failed_keys,