Skip to content

vllm.v1.kv_offload.tiering.manager

TieringOffloadingManager: Multi-tier KV cache offloading orchestrator.

This manager coordinates between a CPU primary tier (with direct GPU access) and zero or more secondary tiers (Storage, Network, etc.) to provide hierarchical KV cache offloading.

Key Design Principles: 1. Always offload to all tiers — When a chunk is stored to the primary tier, it is cascaded to ALL secondary tiers 2. Primary tier is the gateway — Secondary tiers cannot access GPU memory directly; all data flows through the CPU primary tier 3. Staged promotion — Chunks in secondary tiers must be promoted to the primary tier before GPU can access them 4. Transparent retry mechanism — Return None from lookup() to signal "data is being promoted, try later" 5. ref_cnt as eviction protection — primary.prepare_read() increments ref_cnt, protecting chunks from eviction until complete_read() is called

Classes:

CPUPrimaryTierOffloadingManager

Bases: CPUOffloadingManager

CPUOffloadingManager with a primary/secondary transfer interface.

The inherited prepare_store/complete_store/prepare_load/complete_load are the GPU-facing OffloadingManager interface. These aliases expose the same operations from the secondary tier perspective, where read/write refers to secondary accessing primary. This avoids confusion when reading TieringOffloadingManager code (e.g. calling prepare_load inside a cascade/store path would be misleading).

Methods:

  • get_kv_memoryview –

    Return the memoryview over the primary tier's KV cache buffer.

  • prepare_read –

    Pin chunks for a CPU-to-secondary transfer.

Source code in vllm/v1/kv_offload/tiering/manager.py
class CPUPrimaryTierOffloadingManager(CPUOffloadingManager):
    """CPUOffloadingManager with a primary/secondary transfer interface.

    The inherited prepare_store/complete_store/prepare_load/complete_load are the
    GPU-facing OffloadingManager interface. These aliases expose the same operations
    from the secondary tier perspective, where read/write refers to secondary
    accessing primary. This avoids confusion when reading TieringOffloadingManager
    code (e.g. calling prepare_load inside a cascade/store path would be misleading).
    """

    def __init__(
        self,
        num_chunks: int,
        mmap_region: SharedOffloadRegion,
        cache_policy: str = "lru",
        cache_policy_module_path: str | None = None,
        enable_events: bool = False,
    ):
        super().__init__(
            num_chunks=num_chunks,
            cache_policy=cache_policy,
            cache_policy_module_path=cache_policy_module_path,
            enable_events=enable_events,
        )
        self._mmap_region = mmap_region
        # read/write is for CPU<->secondary transfers,
        # load/store is for CPU<->GPU transfers.
        # These aliases avoid calling prepare_load inside a store path.
        self.complete_read = self.complete_load
        self.prepare_write = self.prepare_store
        self.complete_write = self.complete_store

        self._kv_memoryview = mmap_region.create_kv_memoryview()

    def prepare_read(
        self, keys: Collection[OffloadKey], req_context: ReqContext
    ) -> LoadStoreSpec:
        """Pin chunks for a CPU-to-secondary transfer.

        Cascade reads are implementation details of tiering, not additional
        request accesses, so they must not alter request-scoped recency.
        """
        return self._prepare_load(keys, req_context, record_access=False)

    def get_kv_memoryview(self) -> memoryview:
        """Return the memoryview over the primary tier's KV cache buffer.

        The view has shape (num_chunks, row_stride_bytes) and is backed by the
        SharedOffloadRegion mmap.  Secondary tiers address chunk *c* as
        ``view[c]``.
        """
        return self._kv_memoryview

    @override
    def shutdown(self) -> None:
        super().shutdown()
        self._kv_memoryview.release()
        self._mmap_region.cleanup()

get_kv_memoryview()

Return the memoryview over the primary tier's KV cache buffer.

The view has shape (num_chunks, row_stride_bytes) and is backed by the SharedOffloadRegion mmap. Secondary tiers address chunk c as view[c].

Source code in vllm/v1/kv_offload/tiering/manager.py
def get_kv_memoryview(self) -> memoryview:
    """Return the memoryview over the primary tier's KV cache buffer.

    The view has shape (num_chunks, row_stride_bytes) and is backed by the
    SharedOffloadRegion mmap.  Secondary tiers address chunk *c* as
    ``view[c]``.
    """
    return self._kv_memoryview

prepare_read(keys, req_context)

Pin chunks for a CPU-to-secondary transfer.

Cascade reads are implementation details of tiering, not additional request accesses, so they must not alter request-scoped recency.

Source code in vllm/v1/kv_offload/tiering/manager.py
def prepare_read(
    self, keys: Collection[OffloadKey], req_context: ReqContext
) -> LoadStoreSpec:
    """Pin chunks for a CPU-to-secondary transfer.

    Cascade reads are implementation details of tiering, not additional
    request accesses, so they must not alter request-scoped recency.
    """
    return self._prepare_load(keys, req_context, record_access=False)

PendingPromotion dataclass

Accumulator for chunks awaiting submit_load() for one (tier, request).

Source code in vllm/v1/kv_offload/tiering/manager.py
@dataclass
class PendingPromotion:
    """Accumulator for chunks awaiting submit_load() for one (tier, request)."""

    req_context: ReqContext
    keys: list[OffloadKey] = field(default_factory=list)
    chunk_ids: list[int] = field(default_factory=list)

TieringOffloadingManager

Bases: OffloadingManager

Orchestrates multi-tier KV cache offloading.

This manager coordinates between a CPU primary tier (with direct GPU access) and zero or more secondary tiers (Storage, Network, etc.) to provide hierarchical KV cache offloading.

Key internal state
  • Minimal state tracking; relies on secondary tiers to report completion via get_finished_jobs()
  • Secondary tiers return JobResult objects containing all necessary information
  • job_id_counter: monotonically increasing counter for job IDs

Methods:

  • __init__ –

    Initialize the TieringOffloadingManager.

  • complete_load –

    Mark chunks as done loading from primary tier to GPU.

  • complete_store –

    Mark chunks as done storing from GPU to primary tier.

  • create_store_job –

    Pin chunks in the primary tier and create a tracked store job.

  • lookup –

    Check whether a single chunk is offloaded and ready.

  • on_new_request –

    Query each secondary tier for its offload policy preference.

  • on_schedule_end –

    End-of-schedule hook: process finished jobs, flush deferred

  • prepare_load –

    Prepare chunks to be loaded from primary tier to GPU.

  • prepare_store –

    Prepare chunks to be stored from GPU to primary tier.

  • reset_cache –

    Reset transfer bookkeeping and primary-tier cache.

  • shutdown –

    Shut down secondary tiers before releasing primary resources.

  • take_events –

    Yield events owned by the primary and secondary tiers.

  • touch –

    Mark chunks as recently used in all tiers.

Source code in vllm/v1/kv_offload/tiering/manager.py
 175
 176
 177
 178
 179
 180
 181
 182
 183
 184
 185
 186
 187
 188
 189
 190
 191
 192
 193
 194
 195
 196
 197
 198
 199
 200
 201
 202
 203
 204
 205
 206
 207
 208
 209
 210
 211
 212
 213
 214
 215
 216
 217
 218
 219
 220
 221
 222
 223
 224
 225
 226
 227
 228
 229
 230
 231
 232
 233
 234
 235
 236
 237
 238
 239
 240
 241
 242
 243
 244
 245
 246
 247
 248
 249
 250
 251
 252
 253
 254
 255
 256
 257
 258
 259
 260
 261
 262
 263
 264
 265
 266
 267
 268
 269
 270
 271
 272
 273
 274
 275
 276
 277
 278
 279
 280
 281
 282
 283
 284
 285
 286
 287
 288
 289
 290
 291
 292
 293
 294
 295
 296
 297
 298
 299
 300
 301
 302
 303
 304
 305
 306
 307
 308
 309
 310
 311
 312
 313
 314
 315
 316
 317
 318
 319
 320
 321
 322
 323
 324
 325
 326
 327
 328
 329
 330
 331
 332
 333
 334
 335
 336
 337
 338
 339
 340
 341
 342
 343
 344
 345
 346
 347
 348
 349
 350
 351
 352
 353
 354
 355
 356
 357
 358
 359
 360
 361
 362
 363
 364
 365
 366
 367
 368
 369
 370
 371
 372
 373
 374
 375
 376
 377
 378
 379
 380
 381
 382
 383
 384
 385
 386
 387
 388
 389
 390
 391
 392
 393
 394
 395
 396
 397
 398
 399
 400
 401
 402
 403
 404
 405
 406
 407
 408
 409
 410
 411
 412
 413
 414
 415
 416
 417
 418
 419
 420
 421
 422
 423
 424
 425
 426
 427
 428
 429
 430
 431
 432
 433
 434
 435
 436
 437
 438
 439
 440
 441
 442
 443
 444
 445
 446
 447
 448
 449
 450
 451
 452
 453
 454
 455
 456
 457
 458
 459
 460
 461
 462
 463
 464
 465
 466
 467
 468
 469
 470
 471
 472
 473
 474
 475
 476
 477
 478
 479
 480
 481
 482
 483
 484
 485
 486
 487
 488
 489
 490
 491
 492
 493
 494
 495
 496
 497
 498
 499
 500
 501
 502
 503
 504
 505
 506
 507
 508
 509
 510
 511
 512
 513
 514
 515
 516
 517
 518
 519
 520
 521
 522
 523
 524
 525
 526
 527
 528
 529
 530
 531
 532
 533
 534
 535
 536
 537
 538
 539
 540
 541
 542
 543
 544
 545
 546
 547
 548
 549
 550
 551
 552
 553
 554
 555
 556
 557
 558
 559
 560
 561
 562
 563
 564
 565
 566
 567
 568
 569
 570
 571
 572
 573
 574
 575
 576
 577
 578
 579
 580
 581
 582
 583
 584
 585
 586
 587
 588
 589
 590
 591
 592
 593
 594
 595
 596
 597
 598
 599
 600
 601
 602
 603
 604
 605
 606
 607
 608
 609
 610
 611
 612
 613
 614
 615
 616
 617
 618
 619
 620
 621
 622
 623
 624
 625
 626
 627
 628
 629
 630
 631
 632
 633
 634
 635
 636
 637
 638
 639
 640
 641
 642
 643
 644
 645
 646
 647
 648
 649
 650
 651
 652
 653
 654
 655
 656
 657
 658
 659
 660
 661
 662
 663
 664
 665
 666
 667
 668
 669
 670
 671
 672
 673
 674
 675
 676
 677
 678
 679
 680
 681
 682
 683
 684
 685
 686
 687
 688
 689
 690
 691
 692
 693
 694
 695
 696
 697
 698
 699
 700
 701
 702
 703
 704
 705
 706
 707
 708
 709
 710
 711
 712
 713
 714
 715
 716
 717
 718
 719
 720
 721
 722
 723
 724
 725
 726
 727
 728
 729
 730
 731
 732
 733
 734
 735
 736
 737
 738
 739
 740
 741
 742
 743
 744
 745
 746
 747
 748
 749
 750
 751
 752
 753
 754
 755
 756
 757
 758
 759
 760
 761
 762
 763
 764
 765
 766
 767
 768
 769
 770
 771
 772
 773
 774
 775
 776
 777
 778
 779
 780
 781
 782
 783
 784
 785
 786
 787
 788
 789
 790
 791
 792
 793
 794
 795
 796
 797
 798
 799
 800
 801
 802
 803
 804
 805
 806
 807
 808
 809
 810
 811
 812
 813
 814
 815
 816
 817
 818
 819
 820
 821
 822
 823
 824
 825
 826
 827
 828
 829
 830
 831
 832
 833
 834
 835
 836
 837
 838
 839
 840
 841
 842
 843
 844
 845
 846
 847
 848
 849
 850
 851
 852
 853
 854
 855
 856
 857
 858
 859
 860
 861
 862
 863
 864
 865
 866
 867
 868
 869
 870
 871
 872
 873
 874
 875
 876
 877
 878
 879
 880
 881
 882
 883
 884
 885
 886
 887
 888
 889
 890
 891
 892
 893
 894
 895
 896
 897
 898
 899
 900
 901
 902
 903
 904
 905
 906
 907
 908
 909
 910
 911
 912
 913
 914
 915
 916
 917
 918
 919
 920
 921
 922
 923
 924
 925
 926
 927
 928
 929
 930
 931
 932
 933
 934
 935
 936
 937
 938
 939
 940
 941
 942
 943
 944
 945
 946
 947
 948
 949
 950
 951
 952
 953
 954
 955
 956
 957
 958
 959
 960
 961
 962
 963
 964
 965
 966
 967
 968
 969
 970
 971
 972
 973
 974
 975
 976
 977
 978
 979
 980
 981
 982
 983
 984
 985
 986
 987
 988
 989
 990
 991
 992
 993
 994
 995
 996
 997
 998
 999
1000
1001
1002
1003
1004
1005
1006
class TieringOffloadingManager(OffloadingManager):
    """Orchestrates multi-tier KV cache offloading.

    This manager coordinates between a CPU primary tier (with direct GPU access)
    and zero or more secondary tiers (Storage, Network, etc.) to provide
    hierarchical KV cache offloading.

    Key internal state:
      - Minimal state tracking; relies on secondary tiers to report completion
        via get_finished_jobs()
      - Secondary tiers return JobResult objects containing all necessary
        information
      - job_id_counter: monotonically increasing counter for job IDs
    """

    def __init__(
        self,
        primary_tier: CPUPrimaryTierOffloadingManager,
        secondary_tiers: list[SecondaryTierManager] | None = None,
    ):
        """Initialize the TieringOffloadingManager.

        Args:
            primary_tier: The primary tier manager (CPU-based).
            secondary_tiers: List of secondary tier managers (e.g., Storage,
                            Network). Can be None or empty list.

        """
        self.primary_tier: CPUPrimaryTierOffloadingManager = primary_tier
        self.secondary_tiers = secondary_tiers or []

        self._job_id_counter: int = 0
        # Job tracking: maps job_id to metadata for all in-flight transfers.
        # TransferJob.is_promotion distinguishes direction:
        #   True:  secondary → primary (promotion)
        #   False: primary → secondary (cascade)
        self._jobs: dict[JobId, JobMetadata] = {}
        primary_view = self.primary_tier.get_kv_memoryview()
        assert primary_view.strides is not None
        self._metrics = TieringMetricsTracker(
            tier_types=[tier.tier_type for tier in self.secondary_tiers],
            num_primary_chunks=self.primary_tier._num_chunks,
            primary_chunk_size=primary_view.strides[0],
        )

        # Pending promotion requests accumulated during lookup() calls; flushed
        # as one batched submit_load() per (tier, request) in on_schedule_end().
        # Outer key: tier index. Inner key: req_context.req_id — the same ReqContext
        # object is reused for all chunk lookups of a given request per engine step.
        self._pending_load_submissions: dict[int, dict[str, PendingPromotion]] = {}

        # Gate for once-per-step execution of _maybe_process_finished_jobs().
        # Reset at the end of each step in on_schedule_end().
        self._processed_jobs_this_step: bool = False

        # Per-request state for prepared GPU->primary stores and finalization.
        # Secondary tiers are finalized only after pending primary stores reach
        # complete_store(), since complete_store() can still submit cascades.
        self._req_state: dict[str, RequestState] = {}

        # Cached ParentManager wrappers for each secondary tier.
        self._tier_parents: dict[SecondaryTierManager, _SecondaryTierFacingParent] = {
            tier: _SecondaryTierFacingParent(self, tier_idx)
            for tier_idx, tier in enumerate(self.secondary_tiers)
        }

        self._tier_index: dict[SecondaryTierManager, int] = {
            tier: i for i, tier in enumerate(self.secondary_tiers)
        }

    @property
    def _transfer_jobs(self) -> dict[JobId, JobMetadata]:
        return self._jobs

    def _next_job_id(self) -> JobId:
        """Generate a unique job ID for async transfer tracking."""
        job_id = self._job_id_counter
        self._job_id_counter += 1
        return job_id

    def _register_job(self, transfer_job: TransferJob, tier_idx: int) -> None:
        job_metadata = JobMetadata(transfer_job, tier_idx)
        self._jobs[transfer_job.job_id] = job_metadata
        self._metrics.on_job_registered(job_metadata)

    def _pop_job(self, job_id: JobId) -> JobMetadata | None:
        return self._jobs.pop(job_id, None)

    def _maybe_process_finished_jobs(self):
        """Poll secondary tiers for completed jobs (at most once per step).

        Guarded by _processed_jobs_this_step: the first call in an engine step
        does the actual polling; subsequent calls are no-ops. The flag is reset
        in on_schedule_end() at the end of each step.
        """
        if self._processed_jobs_this_step:
            return
        self._processed_jobs_this_step = True
        self._process_finished_jobs()

    def _complete_promotion(
        self, job_metadata: JobMetadata, completed_job: JobResult
    ) -> None:
        transfer_job = job_metadata.transfer_job
        successful_keys = completed_job.successful_keys
        failed_keys: Collection[OffloadKey]
        if completed_job.success:
            successful_keys = transfer_job.keys
            failed_keys = ()
        elif successful_keys:
            failed_keys_set = set(transfer_job.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 = transfer_job.keys

        if successful_keys:
            self.primary_tier.complete_write(
                successful_keys,
                transfer_job.req_context,
                True,
            )
        if failed_keys:
            self.primary_tier.complete_write(
                failed_keys,
                transfer_job.req_context,
                False,
            )

    def _process_finished_jobs(self):
        """Unconditionally poll all secondary tiers for completed jobs.

        This method:
        1. Calls get_finished_jobs() on each secondary tier
        2. For completed stores (primary→secondary): calls primary.complete_read()
           to decrement ref_cnt
        3. For completed loads (secondary→primary): calls primary.complete_write()
           to make chunks available
        """
        for i, tier in enumerate(self.secondary_tiers):
            for completed_job in tier.get_finished_jobs():
                job_id = completed_job.job_id
                job_metadata = self._pop_job(job_id)
                assert job_metadata is not None, (
                    f"Finished job_id {job_id} from tier #{i}"
                    f" ({tier.tier_type}) not in _jobs"
                )
                assert job_metadata.tier_idx == i, (
                    f"Finished job_id {job_id} reported by tier #{i}"
                    f" but belongs to tier #{job_metadata.tier_idx}"
                )
                transfer_job = job_metadata.transfer_job
                self._metrics.on_job_finished(job_metadata, completed_job)

                if transfer_job.is_promotion:
                    # secondary→primary transfer (promotion) completed.
                    # Make chunks available in primary tier.
                    self._complete_promotion(job_metadata, completed_job)
                else:
                    # primary→secondary transfer completed.
                    # Decrement ref_cnt on primary chunks.
                    self.primary_tier.complete_read(
                        transfer_job.keys, transfer_job.req_context
                    )
                    if completed_job.success:
                        self._update_backpressure(tier, job_metadata, completed_job)

    def _should_store_to_tier(
        self, tier: SecondaryTierManager, num_blocks: int
    ) -> bool:
        detector = tier.bp_detector
        if detector is None:
            return True
        return detector.should_store(num_blocks)

    def _update_backpressure(
        self,
        tier: SecondaryTierManager,
        job_metadata: JobMetadata,
        completed_job: JobResult,
    ) -> None:
        detector = tier.bp_detector
        if detector is None:
            return
        was_under_pressure = detector.is_under_pressure()
        tj = job_metadata.transfer_job
        num_bytes = (
            completed_job.transfer_bytes
            if completed_job.transfer_bytes is not None
            else len(tj.keys) * tier.block_size_bytes
        )
        detector.update(tj.submit_time, num_bytes)
        if detector.is_under_pressure() != was_under_pressure:
            tier_idx = self._tier_index[tier]
            logger.info(
                "Tier #%d (%s) back-pressure %s (stats=%s)",
                tier_idx,
                tier.tier_type,
                "activated" if detector.is_under_pressure() else "cleared",
                detector.stats,
            )

    @override
    def lookup(
        self,
        key: OffloadKey,
        req_context: ReqContext,
        *,
        exclude_tier_idx: int | None = None,
    ) -> LookupResult:
        """Check whether a single chunk is offloaded and ready.

        Algorithm:
            1. Process any completed async jobs first.
            2. Query primary tier — short-circuit on hit or in-flight.
            3. On primary miss, query secondary tiers — stop on first
               hit and initiate promotion.

        Args:
            key: Chunk hash to look up.
            req_context: Per-request context.
            exclude_tier_idx: Skip this tier index during the lookup.

        Returns:
            HIT       — chunk is ready in the primary tier.
            HIT_PENDING — chunk found but not yet readable (write
                        in-flight on the primary tier).
            RETRY     — promotion started or a secondary tier is busy.
            MISS      — chunk not found in any tier, or primary is full
                        and cannot accept a promotion.

        """
        # Poll first so a promotion that finished since the last call is
        # already reflected as HIT (not stale HIT_PENDING/MISS) below, and
        # so chunks freed by cascade or promotion completions are evictable
        # in time for a promotion this lookup may initiate.
        self._maybe_process_finished_jobs()

        start_time = time.monotonic()
        primary_hit = self.primary_tier.lookup(key, req_context)
        lookup_duration = time.monotonic() - start_time
        self._metrics.on_lookup(
            req_context,
            key,
            self._metrics.primary_tier_label,
            primary_hit,
            lookup_duration,
        )
        if primary_hit is LookupResult.HIT:
            return LookupResult.HIT
        if primary_hit is LookupResult.HIT_PENDING:
            return LookupResult.HIT_PENDING

        any_retry = False
        for i, tier in enumerate(self.secondary_tiers):
            if i == exclude_tier_idx:
                continue
            if not req_context.load_tier_filter.allows(tier.medium, tier.locality):
                continue
            labelvalues = self._metrics.tier_label(i)
            start_time = time.monotonic()
            result = tier.lookup(key, req_context)
            lookup_duration = time.monotonic() - start_time
            if result is LookupResult.HIT:
                self._metrics.on_lookup(
                    req_context,
                    key,
                    labelvalues,
                    result,
                    lookup_duration,
                )
                promoted = self._initiate_promotion(i, key, req_context)
                return LookupResult.MISS if not promoted else LookupResult.HIT_PENDING
            if result is LookupResult.RETRY:
                any_retry = True
            self._metrics.on_lookup(
                req_context,
                key,
                labelvalues,
                result,
                lookup_duration,
            )

        if any_retry:
            return LookupResult.RETRY
        return LookupResult.MISS

    def _initiate_promotion(
        self,
        tier_idx: int,
        key: OffloadKey,
        req_context: ReqContext,
    ) -> bool:
        """Queue a chunk for promotion from a secondary tier to the primary tier.

        Allocates space in the primary tier immediately (sets ref_cnt=-1 so
        subsequent lookups within the same step see the slot as in-flight),
        then defers the actual submit_load() call to _flush_pending_promotions()
        so all chunks queued during one engine step are submitted as a single
        batched job.

        Args:
            tier_idx: The secondary tier index to promote from
            key: Chunk to promote
            req_context: Per-request context forwarded to primary.prepare_write().

        Returns:
            True if promotion was initiated, False if primary tier is full.

        """
        # Allocate space in primary tier for promoted chunk.
        # Must happen immediately so primary.lookup() returns None (in-flight)
        # for this key on any subsequent lookup() call within the same step,
        # preventing duplicate promotion attempts.
        primary_write_result = self.primary_tier.prepare_write([key], req_context)

        if primary_write_result is None:
            # Primary tier is full; caller should treat the chunk as unavailable
            # rather than retrying indefinitely.
            self._metrics.on_promotion_allocation_failure()
            return False

        store_spec = primary_write_result.store_spec
        assert isinstance(store_spec, CPULoadStoreSpec)
        # Defer submit_load to on_schedule_end(). Group by (tier, request) so
        # each request's chunks are submitted as one batched job per tier.
        tier_pending = self._pending_load_submissions.setdefault(tier_idx, {})
        ctx_id = req_context.req_id
        if ctx_id not in tier_pending:
            tier_pending[ctx_id] = PendingPromotion(
                keys=[], chunk_ids=[], req_context=req_context
            )
        entry = tier_pending[ctx_id]
        entry.keys.extend(primary_write_result.keys_to_store)
        entry.chunk_ids.extend(store_spec.chunk_ids)
        return True

    def _flush_pending_promotions(self) -> None:
        """Submit one batched submit_load() per (tier, request).

        Called from on_schedule_end() at the end of each scheduler step,
        flushing all promotion requests deferred during lookup().
        """
        if not self._pending_load_submissions:
            return

        for tier_idx, pending_by_ctx in self._pending_load_submissions.items():
            tier = self.secondary_tiers[tier_idx]
            for entry in pending_by_ctx.values():
                job_id = self._next_job_id()
                job_metadata = TransferJob(
                    job_id=job_id,
                    keys=entry.keys,
                    chunk_ids=np.array(entry.chunk_ids, dtype=np.int32),
                    is_promotion=True,
                    req_context=entry.req_context,
                )
                self._register_job(job_metadata, tier_idx)
                tier.submit_load(job_metadata)

        self._pending_load_submissions.clear()

    @override
    def prepare_load(
        self, keys: Collection[OffloadKey], req_context: ReqContext
    ) -> LoadStoreSpec:
        """Prepare chunks to be loaded from primary tier to GPU.

        Callers only pass keys already confirmed HIT by lookup() earlier this
        step.

        This increments ref_cnt on the chunks in the primary tier, protecting
        them from eviction during the transfer.

        Args:
            keys: Chunks to prepare for loading.
            req_context: Per-request context.

        Returns:
            LoadStoreSpec for reading from primary tier.

        """
        return self.primary_tier.prepare_load(keys, req_context)

    @override
    def touch(self, keys: Collection[OffloadKey], req_context: ReqContext):
        """Mark chunks as recently used in all tiers.

        Args:
            keys: Chunks to mark as recently used.
            req_context: Per-request context.

        """
        self.primary_tier.touch(keys, req_context)
        for tier in self.secondary_tiers:
            tier.touch(keys, req_context)

    @override
    def complete_load(self, keys: Collection[OffloadKey], req_context: ReqContext):
        """Mark chunks as done loading from primary tier to GPU.

        This decrements ref_cnt on the chunks in the primary tier, allowing
        them to be evicted again.

        Args:
            keys: Chunks that finished loading.
            req_context: Per-request context.

        """
        self.primary_tier.complete_load(keys, req_context)

    @override
    def prepare_store(
        self, keys: Collection[OffloadKey], req_context: ReqContext
    ) -> PrepareStoreOutput | None:
        """Prepare chunks to be stored from GPU to primary tier.

        CRITICAL: This method calls _maybe_process_finished_jobs() FIRST to ensure
        that any completed async transfers have their ref_cnt decremented
        before the primary tier makes eviction decisions.

        For request-level tiers, chunks already present in the primary tier
        are immediately cascaded via submit_store().

        Args:
            keys: Chunks to prepare for storing.
            req_context: Per-request context.

        Returns:
            PrepareStoreOutput describing where to store chunks and what was
            evicted, or None if store cannot proceed.

        """
        # Step 1: Poll for completed async jobs FIRST
        # _process_finished_jobs() handles two kinds of completions here:
        #  - Cascade completions (store to a secondary tier, either a local
        #    cascade or a store job created for a remote requester via
        #    create_store_job()): decrements ref_cnt on the primary chunks
        #    that were read, making them evictable again once ref_cnt hits 0.
        #  - Promotion completions (secondary->primary loads): sets a
        #    not-yet-ready chunk's ref_cnt from -1 to 0 via complete_write(),
        #    making it evictable for the first time.
        # Both must be accounted for before the eviction decision below.
        self._maybe_process_finished_jobs()

        # Step 2: Store to primary tier (new chunks only).
        # Cascading of these newly-stored chunks to ALL secondary tiers
        # happens later in complete_store(), after the GPU→Primary transfer
        # completes.
        primary_result = self.primary_tier.prepare_store(keys, req_context)

        if primary_result is None:
            return None

        if primary_result.keys_to_store:
            state = self._req_state[req_context.req_id]
            state.pending_primary_stores += 1

        # Step 3: For request-level tiers, cascade chunks already in primary
        request_level_tiers = self._req_state[req_context.req_id].request_level_tiers
        if request_level_tiers:
            keys_to_store_set = set(primary_result.keys_to_store)
            keys_already_in_primary = tuple(
                k for k in keys if k not in keys_to_store_set
            )
            if keys_already_in_primary:
                self._cascade_existing_chunks_to_request_level_tiers(
                    keys_already_in_primary, req_context, request_level_tiers
                )

        return primary_result

    def _cascade_existing_chunks_to_request_level_tiers(
        self,
        keys: Sequence[OffloadKey],
        req_context: ReqContext,
        request_level_tiers: set[int],
    ) -> None:
        """For tiers that requested request-level policy, submit_store() for
        chunks that are already present in the primary tier.

        A key whose primary write is still in flight (HIT_PENDING) cannot be
        dropped: prepare_store already excluded it as present, and the
        scheduler advances past its chunk, so no path offers it again. Park it
        instead. MISS keys are dropped, since nothing is there to read.

        The primary tier resolves every key it holds, so RETRY cannot reach
        here. Parking on it would have no guarantee of ever draining, which is
        what makes parking HIT_PENDING safe, so it is rejected rather than
        guessed at.
        """
        state = self._req_state[req_context.req_id]
        ready_keys = []
        for key in keys:
            result = self.primary_tier.lookup(key, req_context)
            if result is LookupResult.HIT:
                ready_keys.append(key)
            elif result is LookupResult.HIT_PENDING:
                state.pending_cascade_keys.append(key)
            else:
                assert result is LookupResult.MISS, (
                    f"primary tier returned {result} for a cascade key"
                )
        if not ready_keys:
            return

        for tier_idx in request_level_tiers:
            tier = self.secondary_tiers[tier_idx]
            if not self._should_store_to_tier(tier, len(ready_keys)):
                continue
            job_metadata = self.create_store_job(ready_keys, req_context, tier_idx)
            tier.submit_store(job_metadata)

    def _flush_pending_cascades(self) -> None:
        """Retry request-level cascades parked on an in-flight primary write.

        A parked key always resolves, to HIT or to MISS, so the set drains and
        a request cannot be held from finalization forever.
        """
        for req_id, state in list(self._req_state.items()):
            if not state.pending_cascade_keys:
                continue
            assert state.request_level_tiers
            keys, state.pending_cascade_keys = state.pending_cascade_keys, []
            self._cascade_existing_chunks_to_request_level_tiers(
                keys, state.req_context, state.request_level_tiers
            )
            self._maybe_finalize_request(req_id)

    @override
    def complete_store(
        self,
        keys: Collection[OffloadKey],
        req_context: ReqContext,
        success: bool = True,
    ) -> None:
        """Mark chunks as done storing from GPU to primary tier.

        This is where secondary tier cascading happens — after chunks are
        confirmed to be in the primary tier, they are cascaded to ALL
        secondary tiers.

        For each secondary tier:
        1. Call primary.prepare_read() to get LoadStoreSpec AND increment
           ref_cnt (protecting chunks during async transfer)
        2. Call tier.submit_store() to start async transfer: primary→secondary
        3. Track the job in _store_jobs dictionary

        Args:
            keys: Chunks that finished storing.
            success: Whether the GPU→primary transfer succeeded.
            req_context: Per-request context forwarded to primary.prepare_read().

        """
        # Step 1: Complete store in primary tier (makes chunks loadable)
        self.primary_tier.complete_store(keys, req_context, success)

        if success:
            # Step 2: Cascade to ALL secondary tiers
            # For each secondary tier, call primary.prepare_read() to get the
            # LoadStoreSpec AND to increment ref_cnt (protecting chunks from
            # eviction during the async transfer). One prepare_read() call per
            # secondary tier.
            for tier_idx, tier in enumerate(self.secondary_tiers):
                if not self._should_store_to_tier(tier, len(keys)):
                    continue
                job_metadata = self.create_store_job(keys, req_context, tier_idx)
                tier.submit_store(job_metadata)

        # Note: The async transfers are now in flight. Their completion is
        # tracked via get_finished_jobs() / _maybe_process_finished_jobs().
        req_id = req_context.req_id
        state = self._req_state[req_id]
        assert state.pending_primary_stores > 0
        state.pending_primary_stores -= 1
        self._maybe_finalize_request(req_id)

    def create_store_job(
        self,
        keys: Collection[OffloadKey],
        req_context: ReqContext,
        tier_idx: int = 0,
    ) -> TransferJob:
        """Pin chunks in the primary tier and create a tracked store job.

        Calls prepare_read() to increment ref_cnt (protecting chunks
        from eviction during the async transfer), allocates a job ID,
        and registers the job in _jobs.

        The caller is responsible for the actual data transfer and
        reporting completion via get_finished_jobs().
        """
        primary_chunks_spec = self.primary_tier.prepare_read(keys, req_context)
        assert isinstance(primary_chunks_spec, CPULoadStoreSpec)
        job_id = self._next_job_id()
        job_metadata = TransferJob(
            job_id=job_id,
            keys=keys,
            chunk_ids=primary_chunks_spec.chunk_ids,
            is_promotion=False,
            req_context=req_context,
        )
        self._register_job(job_metadata, tier_idx)
        return job_metadata

    @override
    def on_new_request(
        self,
        req_context: ReqContext,
        *,
        exclude_tier_idx: int | None = None,
    ) -> RequestOffloadingContext:
        """Query each secondary tier for its offload policy preference.

        Returns REQUEST_LEVEL if ANY secondary tier wants request-level.
        Only stores REQUEST_LEVEL tier decisions for use in prepare_store.
        """
        state = RequestState(req_context=req_context)
        self._metrics.on_new_request(req_context)
        for tier_idx, tier in enumerate(self.secondary_tiers):
            if tier_idx == exclude_tier_idx:
                continue
            tier_ctx = tier.on_new_request(req_context)
            if tier_ctx.policy == OffloadPolicy.REQUEST_LEVEL:
                if state.request_level_tiers is None:
                    state.request_level_tiers = set()
                state.request_level_tiers.add(tier_idx)
        self._req_state[req_context.req_id] = state

        policy = (
            OffloadPolicy.REQUEST_LEVEL
            if state.request_level_tiers
            else OffloadPolicy.CHUNK_LEVEL
        )
        return RequestOffloadingContext(policy=policy)

    @override
    def on_request_finished(
        self,
        req_context: ReqContext,
        *,
        exclude_tier_idx: int | None = None,
    ) -> None:
        self.primary_tier.on_request_finished(req_context)
        state = self._req_state[req_context.req_id]
        state.is_finished = True
        self._maybe_finalize_request(req_context.req_id, exclude_tier_idx)

    def _maybe_finalize_request(
        self,
        req_id: str,
        exclude_tier_idx: int | None = None,
    ) -> None:
        """Finalize secondary tiers once no more cascades can be submitted.

        Their finalization is delayed until pending GPU->primary stores
        finish, since those callbacks may still submit secondary stores.
        """
        state = self._req_state[req_id]
        if not state.is_finished:
            return
        if state.pending_primary_stores != 0:
            return
        if state.pending_cascade_keys:
            return

        for tier_idx, tier in enumerate(self.secondary_tiers):
            if tier_idx == exclude_tier_idx:
                continue
            tier.on_request_finished(state.req_context)
        self._metrics.on_request_finished(state.req_context)
        del self._req_state[req_id]

    @override
    def on_schedule_end(self, context: ScheduleEndContext) -> None:
        """End-of-schedule hook: process finished jobs, flush deferred
        promotions, and reset the per-step gate.

        Called once per scheduler step from
        OffloadingConnectorScheduler.build_connector_meta().
        """
        # Catch-all poll: guarantees jobs are processed even on steps where
        # lookup()/prepare_store() were never called (e.g. no requests
        # scheduled but a tier still has_pending_work()).
        self._maybe_process_finished_jobs()

        for tier in self.secondary_tiers:
            tier.serve_external_requests(self._tier_parents[tier])

        # Reset the per-step gate AFTER serve_external_requests so that
        # lookup() calls within it skip redundant _process_finished_jobs().
        self._processed_jobs_this_step = False

        self._flush_pending_promotions()
        self._flush_pending_cascades()
        for tier in self.secondary_tiers:
            tier.on_schedule_end(context)

        for req_id in context.new_req_ids:
            state = self._req_state.get(req_id)
            if state is None:
                continue
            self._metrics.on_request_allocated(state.req_context)

    @override
    def has_pending_work(self) -> bool:
        # In-flight primary<->secondary transfers (pending promotions are
        # translated to transfer jobs in on_schedule_end), plus any work the
        # secondary tiers themselves still have outstanding.
        return (
            bool(self._jobs)
            or any(state.pending_cascade_keys for state in self._req_state.values())
            or any(tier.has_pending_work() for tier in self.secondary_tiers)
        )

    @override
    def take_events(self) -> Iterable[OffloadingEvent]:
        """Yield events owned by the primary and secondary tiers.

        Yields:
            New OffloadingEvents collected by each tier since the last call.

        """
        yield from self.primary_tier.take_events()
        for tier in self.secondary_tiers:
            yield from tier.take_events()

    @override
    def reset_cache(self) -> None:
        """Reset transfer bookkeeping and primary-tier cache.

        Called during sleep, weight update, or resume. Each secondary tier
        drains its in-flight transfers via drain_jobs() so no tier I/O is
        touching primary memory before the primary tier is reset. A stuck
        tier will block here visibly — preferable to silent corruption
        from reusing primary slots while a transfer is mid-copy.

        Secondary tiers are intentionally not reset: persistent stores
        (FS, network) keep their data across resets. Active request state is
        retained so those requests can continue after the reset; finished
        requests are finalized and removed.
        """
        for tier in self.secondary_tiers:
            tier.drain_jobs()
        # All tier I/O has stopped; consume their completion notifications
        # so manager bookkeeping is consistent before the primary reset.
        self._process_finished_jobs()
        assert not self._jobs

        # Deferred promotion submissions reserve primary slots that the
        # reset below invalidates; their submit_load() has not yet been
        # called so no tier I/O is touching that memory.
        self._pending_load_submissions.clear()
        self._metrics.assert_idle()

        finished_req_ids = []
        for req_id, state in self._req_state.items():
            state.pending_primary_stores = 0
            state.pending_cascade_keys.clear()
            if not state.is_finished:
                continue
            for tier in self.secondary_tiers:
                tier.on_request_finished(state.req_context)
            self._metrics.on_request_finished(state.req_context)
            finished_req_ids.append(req_id)

        self.primary_tier.reset_cache()

        for req_id in finished_req_ids:
            del self._req_state[req_id]
        self._processed_jobs_this_step = False

        for tier in self.secondary_tiers:
            if tier.bp_detector is not None:
                tier.bp_detector.reset()

    @override
    def get_stats(self) -> OffloadingConnectorStats | None:
        stats = self.primary_tier.get_stats()

        if stats is not None and stats.is_empty():
            stats = None

        self._metrics.record_backpressure(self.secondary_tiers)
        metrics_stats = self._metrics.take_stats()
        if metrics_stats is not None:
            if stats is None:
                stats = metrics_stats
            else:
                stats.aggregate(metrics_stats)

        for tier in self.secondary_tiers:
            tier_stats = tier.get_stats()
            if tier_stats is None or tier_stats.is_empty():
                continue
            if stats is None:
                stats = tier_stats
            else:
                stats.aggregate(tier_stats)

        return stats

    @override
    def shutdown(self) -> None:
        """Shut down secondary tiers before releasing primary resources.

        Every secondary tier is given a shutdown attempt. If any shutdown
        fails, preserve the primary mmap because a failed tier may still use it.
        """
        shutdown_error: Exception | None = None
        for tier_idx, tier in enumerate(self.secondary_tiers):
            try:
                tier.shutdown()
            except Exception as exc:
                shutdown_error = exc
                logger.exception(
                    "Failed to shut down secondary tier #%d "
                    "(tier_type=%s, impl_class=%s)",
                    tier_idx,
                    tier.tier_type,
                    type(tier).__name__,
                )

        if shutdown_error is not None:
            raise shutdown_error

        self.primary_tier.shutdown()

__init__(primary_tier, secondary_tiers=None)

Initialize the TieringOffloadingManager.

Parameters:

Source code in vllm/v1/kv_offload/tiering/manager.py
def __init__(
    self,
    primary_tier: CPUPrimaryTierOffloadingManager,
    secondary_tiers: list[SecondaryTierManager] | None = None,
):
    """Initialize the TieringOffloadingManager.

    Args:
        primary_tier: The primary tier manager (CPU-based).
        secondary_tiers: List of secondary tier managers (e.g., Storage,
                        Network). Can be None or empty list.

    """
    self.primary_tier: CPUPrimaryTierOffloadingManager = primary_tier
    self.secondary_tiers = secondary_tiers or []

    self._job_id_counter: int = 0
    # Job tracking: maps job_id to metadata for all in-flight transfers.
    # TransferJob.is_promotion distinguishes direction:
    #   True:  secondary → primary (promotion)
    #   False: primary → secondary (cascade)
    self._jobs: dict[JobId, JobMetadata] = {}
    primary_view = self.primary_tier.get_kv_memoryview()
    assert primary_view.strides is not None
    self._metrics = TieringMetricsTracker(
        tier_types=[tier.tier_type for tier in self.secondary_tiers],
        num_primary_chunks=self.primary_tier._num_chunks,
        primary_chunk_size=primary_view.strides[0],
    )

    # Pending promotion requests accumulated during lookup() calls; flushed
    # as one batched submit_load() per (tier, request) in on_schedule_end().
    # Outer key: tier index. Inner key: req_context.req_id — the same ReqContext
    # object is reused for all chunk lookups of a given request per engine step.
    self._pending_load_submissions: dict[int, dict[str, PendingPromotion]] = {}

    # Gate for once-per-step execution of _maybe_process_finished_jobs().
    # Reset at the end of each step in on_schedule_end().
    self._processed_jobs_this_step: bool = False

    # Per-request state for prepared GPU->primary stores and finalization.
    # Secondary tiers are finalized only after pending primary stores reach
    # complete_store(), since complete_store() can still submit cascades.
    self._req_state: dict[str, RequestState] = {}

    # Cached ParentManager wrappers for each secondary tier.
    self._tier_parents: dict[SecondaryTierManager, _SecondaryTierFacingParent] = {
        tier: _SecondaryTierFacingParent(self, tier_idx)
        for tier_idx, tier in enumerate(self.secondary_tiers)
    }

    self._tier_index: dict[SecondaryTierManager, int] = {
        tier: i for i, tier in enumerate(self.secondary_tiers)
    }

_cascade_existing_chunks_to_request_level_tiers(keys, req_context, request_level_tiers)

For tiers that requested request-level policy, submit_store() for chunks that are already present in the primary tier.

A key whose primary write is still in flight (HIT_PENDING) cannot be dropped: prepare_store already excluded it as present, and the scheduler advances past its chunk, so no path offers it again. Park it instead. MISS keys are dropped, since nothing is there to read.

The primary tier resolves every key it holds, so RETRY cannot reach here. Parking on it would have no guarantee of ever draining, which is what makes parking HIT_PENDING safe, so it is rejected rather than guessed at.

Source code in vllm/v1/kv_offload/tiering/manager.py
def _cascade_existing_chunks_to_request_level_tiers(
    self,
    keys: Sequence[OffloadKey],
    req_context: ReqContext,
    request_level_tiers: set[int],
) -> None:
    """For tiers that requested request-level policy, submit_store() for
    chunks that are already present in the primary tier.

    A key whose primary write is still in flight (HIT_PENDING) cannot be
    dropped: prepare_store already excluded it as present, and the
    scheduler advances past its chunk, so no path offers it again. Park it
    instead. MISS keys are dropped, since nothing is there to read.

    The primary tier resolves every key it holds, so RETRY cannot reach
    here. Parking on it would have no guarantee of ever draining, which is
    what makes parking HIT_PENDING safe, so it is rejected rather than
    guessed at.
    """
    state = self._req_state[req_context.req_id]
    ready_keys = []
    for key in keys:
        result = self.primary_tier.lookup(key, req_context)
        if result is LookupResult.HIT:
            ready_keys.append(key)
        elif result is LookupResult.HIT_PENDING:
            state.pending_cascade_keys.append(key)
        else:
            assert result is LookupResult.MISS, (
                f"primary tier returned {result} for a cascade key"
            )
    if not ready_keys:
        return

    for tier_idx in request_level_tiers:
        tier = self.secondary_tiers[tier_idx]
        if not self._should_store_to_tier(tier, len(ready_keys)):
            continue
        job_metadata = self.create_store_job(ready_keys, req_context, tier_idx)
        tier.submit_store(job_metadata)

_flush_pending_cascades()

Retry request-level cascades parked on an in-flight primary write.

A parked key always resolves, to HIT or to MISS, so the set drains and a request cannot be held from finalization forever.

Source code in vllm/v1/kv_offload/tiering/manager.py
def _flush_pending_cascades(self) -> None:
    """Retry request-level cascades parked on an in-flight primary write.

    A parked key always resolves, to HIT or to MISS, so the set drains and
    a request cannot be held from finalization forever.
    """
    for req_id, state in list(self._req_state.items()):
        if not state.pending_cascade_keys:
            continue
        assert state.request_level_tiers
        keys, state.pending_cascade_keys = state.pending_cascade_keys, []
        self._cascade_existing_chunks_to_request_level_tiers(
            keys, state.req_context, state.request_level_tiers
        )
        self._maybe_finalize_request(req_id)

_flush_pending_promotions()

Submit one batched submit_load() per (tier, request).

Called from on_schedule_end() at the end of each scheduler step, flushing all promotion requests deferred during lookup().

Source code in vllm/v1/kv_offload/tiering/manager.py
def _flush_pending_promotions(self) -> None:
    """Submit one batched submit_load() per (tier, request).

    Called from on_schedule_end() at the end of each scheduler step,
    flushing all promotion requests deferred during lookup().
    """
    if not self._pending_load_submissions:
        return

    for tier_idx, pending_by_ctx in self._pending_load_submissions.items():
        tier = self.secondary_tiers[tier_idx]
        for entry in pending_by_ctx.values():
            job_id = self._next_job_id()
            job_metadata = TransferJob(
                job_id=job_id,
                keys=entry.keys,
                chunk_ids=np.array(entry.chunk_ids, dtype=np.int32),
                is_promotion=True,
                req_context=entry.req_context,
            )
            self._register_job(job_metadata, tier_idx)
            tier.submit_load(job_metadata)

    self._pending_load_submissions.clear()

_initiate_promotion(tier_idx, key, req_context)

Queue a chunk for promotion from a secondary tier to the primary tier.

Allocates space in the primary tier immediately (sets ref_cnt=-1 so subsequent lookups within the same step see the slot as in-flight), then defers the actual submit_load() call to _flush_pending_promotions() so all chunks queued during one engine step are submitted as a single batched job.

Parameters:

  • tier_idx

    (int) –

    The secondary tier index to promote from

  • key

    (OffloadKey) –

    Chunk to promote

  • req_context

    (ReqContext) –

    Per-request context forwarded to primary.prepare_write().

Returns:

  • bool –

    True if promotion was initiated, False if primary tier is full.

Source code in vllm/v1/kv_offload/tiering/manager.py
def _initiate_promotion(
    self,
    tier_idx: int,
    key: OffloadKey,
    req_context: ReqContext,
) -> bool:
    """Queue a chunk for promotion from a secondary tier to the primary tier.

    Allocates space in the primary tier immediately (sets ref_cnt=-1 so
    subsequent lookups within the same step see the slot as in-flight),
    then defers the actual submit_load() call to _flush_pending_promotions()
    so all chunks queued during one engine step are submitted as a single
    batched job.

    Args:
        tier_idx: The secondary tier index to promote from
        key: Chunk to promote
        req_context: Per-request context forwarded to primary.prepare_write().

    Returns:
        True if promotion was initiated, False if primary tier is full.

    """
    # Allocate space in primary tier for promoted chunk.
    # Must happen immediately so primary.lookup() returns None (in-flight)
    # for this key on any subsequent lookup() call within the same step,
    # preventing duplicate promotion attempts.
    primary_write_result = self.primary_tier.prepare_write([key], req_context)

    if primary_write_result is None:
        # Primary tier is full; caller should treat the chunk as unavailable
        # rather than retrying indefinitely.
        self._metrics.on_promotion_allocation_failure()
        return False

    store_spec = primary_write_result.store_spec
    assert isinstance(store_spec, CPULoadStoreSpec)
    # Defer submit_load to on_schedule_end(). Group by (tier, request) so
    # each request's chunks are submitted as one batched job per tier.
    tier_pending = self._pending_load_submissions.setdefault(tier_idx, {})
    ctx_id = req_context.req_id
    if ctx_id not in tier_pending:
        tier_pending[ctx_id] = PendingPromotion(
            keys=[], chunk_ids=[], req_context=req_context
        )
    entry = tier_pending[ctx_id]
    entry.keys.extend(primary_write_result.keys_to_store)
    entry.chunk_ids.extend(store_spec.chunk_ids)
    return True

_maybe_finalize_request(req_id, exclude_tier_idx=None)

Finalize secondary tiers once no more cascades can be submitted.

Their finalization is delayed until pending GPU->primary stores finish, since those callbacks may still submit secondary stores.

Source code in vllm/v1/kv_offload/tiering/manager.py
def _maybe_finalize_request(
    self,
    req_id: str,
    exclude_tier_idx: int | None = None,
) -> None:
    """Finalize secondary tiers once no more cascades can be submitted.

    Their finalization is delayed until pending GPU->primary stores
    finish, since those callbacks may still submit secondary stores.
    """
    state = self._req_state[req_id]
    if not state.is_finished:
        return
    if state.pending_primary_stores != 0:
        return
    if state.pending_cascade_keys:
        return

    for tier_idx, tier in enumerate(self.secondary_tiers):
        if tier_idx == exclude_tier_idx:
            continue
        tier.on_request_finished(state.req_context)
    self._metrics.on_request_finished(state.req_context)
    del self._req_state[req_id]

_maybe_process_finished_jobs()

Poll secondary tiers for completed jobs (at most once per step).

Guarded by _processed_jobs_this_step: the first call in an engine step does the actual polling; subsequent calls are no-ops. The flag is reset in on_schedule_end() at the end of each step.

Source code in vllm/v1/kv_offload/tiering/manager.py
def _maybe_process_finished_jobs(self):
    """Poll secondary tiers for completed jobs (at most once per step).

    Guarded by _processed_jobs_this_step: the first call in an engine step
    does the actual polling; subsequent calls are no-ops. The flag is reset
    in on_schedule_end() at the end of each step.
    """
    if self._processed_jobs_this_step:
        return
    self._processed_jobs_this_step = True
    self._process_finished_jobs()

_next_job_id()

Generate a unique job ID for async transfer tracking.

Source code in vllm/v1/kv_offload/tiering/manager.py
def _next_job_id(self) -> JobId:
    """Generate a unique job ID for async transfer tracking."""
    job_id = self._job_id_counter
    self._job_id_counter += 1
    return job_id

_process_finished_jobs()

Unconditionally poll all secondary tiers for completed jobs.

This method: 1. Calls get_finished_jobs() on each secondary tier 2. For completed stores (primary→secondary): calls primary.complete_read() to decrement ref_cnt 3. For completed loads (secondary→primary): calls primary.complete_write() to make chunks available

Source code in vllm/v1/kv_offload/tiering/manager.py
def _process_finished_jobs(self):
    """Unconditionally poll all secondary tiers for completed jobs.

    This method:
    1. Calls get_finished_jobs() on each secondary tier
    2. For completed stores (primary→secondary): calls primary.complete_read()
       to decrement ref_cnt
    3. For completed loads (secondary→primary): calls primary.complete_write()
       to make chunks available
    """
    for i, tier in enumerate(self.secondary_tiers):
        for completed_job in tier.get_finished_jobs():
            job_id = completed_job.job_id
            job_metadata = self._pop_job(job_id)
            assert job_metadata is not None, (
                f"Finished job_id {job_id} from tier #{i}"
                f" ({tier.tier_type}) not in _jobs"
            )
            assert job_metadata.tier_idx == i, (
                f"Finished job_id {job_id} reported by tier #{i}"
                f" but belongs to tier #{job_metadata.tier_idx}"
            )
            transfer_job = job_metadata.transfer_job
            self._metrics.on_job_finished(job_metadata, completed_job)

            if transfer_job.is_promotion:
                # secondary→primary transfer (promotion) completed.
                # Make chunks available in primary tier.
                self._complete_promotion(job_metadata, completed_job)
            else:
                # primary→secondary transfer completed.
                # Decrement ref_cnt on primary chunks.
                self.primary_tier.complete_read(
                    transfer_job.keys, transfer_job.req_context
                )
                if completed_job.success:
                    self._update_backpressure(tier, job_metadata, completed_job)

complete_load(keys, req_context)

Mark chunks as done loading from primary tier to GPU.

This decrements ref_cnt on the chunks in the primary tier, allowing them to be evicted again.

Parameters:

  • keys

    (Collection[OffloadKey]) –

    Chunks that finished loading.

  • req_context

    (ReqContext) –

    Per-request context.

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def complete_load(self, keys: Collection[OffloadKey], req_context: ReqContext):
    """Mark chunks as done loading from primary tier to GPU.

    This decrements ref_cnt on the chunks in the primary tier, allowing
    them to be evicted again.

    Args:
        keys: Chunks that finished loading.
        req_context: Per-request context.

    """
    self.primary_tier.complete_load(keys, req_context)

complete_store(keys, req_context, success=True)

Mark chunks as done storing from GPU to primary tier.

This is where secondary tier cascading happens — after chunks are confirmed to be in the primary tier, they are cascaded to ALL secondary tiers.

For each secondary tier: 1. Call primary.prepare_read() to get LoadStoreSpec AND increment ref_cnt (protecting chunks during async transfer) 2. Call tier.submit_store() to start async transfer: primary→secondary 3. Track the job in _store_jobs dictionary

Parameters:

  • keys

    (Collection[OffloadKey]) –

    Chunks that finished storing.

  • success

    (bool, default: True ) –

    Whether the GPU→primary transfer succeeded.

  • req_context

    (ReqContext) –

    Per-request context forwarded to primary.prepare_read().

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def complete_store(
    self,
    keys: Collection[OffloadKey],
    req_context: ReqContext,
    success: bool = True,
) -> None:
    """Mark chunks as done storing from GPU to primary tier.

    This is where secondary tier cascading happens — after chunks are
    confirmed to be in the primary tier, they are cascaded to ALL
    secondary tiers.

    For each secondary tier:
    1. Call primary.prepare_read() to get LoadStoreSpec AND increment
       ref_cnt (protecting chunks during async transfer)
    2. Call tier.submit_store() to start async transfer: primary→secondary
    3. Track the job in _store_jobs dictionary

    Args:
        keys: Chunks that finished storing.
        success: Whether the GPU→primary transfer succeeded.
        req_context: Per-request context forwarded to primary.prepare_read().

    """
    # Step 1: Complete store in primary tier (makes chunks loadable)
    self.primary_tier.complete_store(keys, req_context, success)

    if success:
        # Step 2: Cascade to ALL secondary tiers
        # For each secondary tier, call primary.prepare_read() to get the
        # LoadStoreSpec AND to increment ref_cnt (protecting chunks from
        # eviction during the async transfer). One prepare_read() call per
        # secondary tier.
        for tier_idx, tier in enumerate(self.secondary_tiers):
            if not self._should_store_to_tier(tier, len(keys)):
                continue
            job_metadata = self.create_store_job(keys, req_context, tier_idx)
            tier.submit_store(job_metadata)

    # Note: The async transfers are now in flight. Their completion is
    # tracked via get_finished_jobs() / _maybe_process_finished_jobs().
    req_id = req_context.req_id
    state = self._req_state[req_id]
    assert state.pending_primary_stores > 0
    state.pending_primary_stores -= 1
    self._maybe_finalize_request(req_id)

create_store_job(keys, req_context, tier_idx=0)

Pin chunks in the primary tier and create a tracked store job.

Calls prepare_read() to increment ref_cnt (protecting chunks from eviction during the async transfer), allocates a job ID, and registers the job in _jobs.

The caller is responsible for the actual data transfer and reporting completion via get_finished_jobs().

Source code in vllm/v1/kv_offload/tiering/manager.py
def create_store_job(
    self,
    keys: Collection[OffloadKey],
    req_context: ReqContext,
    tier_idx: int = 0,
) -> TransferJob:
    """Pin chunks in the primary tier and create a tracked store job.

    Calls prepare_read() to increment ref_cnt (protecting chunks
    from eviction during the async transfer), allocates a job ID,
    and registers the job in _jobs.

    The caller is responsible for the actual data transfer and
    reporting completion via get_finished_jobs().
    """
    primary_chunks_spec = self.primary_tier.prepare_read(keys, req_context)
    assert isinstance(primary_chunks_spec, CPULoadStoreSpec)
    job_id = self._next_job_id()
    job_metadata = TransferJob(
        job_id=job_id,
        keys=keys,
        chunk_ids=primary_chunks_spec.chunk_ids,
        is_promotion=False,
        req_context=req_context,
    )
    self._register_job(job_metadata, tier_idx)
    return job_metadata

lookup(key, req_context, *, exclude_tier_idx=None)

Check whether a single chunk is offloaded and ready.

Algorithm
  1. Process any completed async jobs first.
  2. Query primary tier — short-circuit on hit or in-flight.
  3. On primary miss, query secondary tiers — stop on first hit and initiate promotion.

Parameters:

  • key

    (OffloadKey) –

    Chunk hash to look up.

  • req_context

    (ReqContext) –

    Per-request context.

  • exclude_tier_idx

    (int | None, default: None ) –

    Skip this tier index during the lookup.

Returns:

  • LookupResult –

    HIT — chunk is ready in the primary tier.

  • LookupResult –

    HIT_PENDING — chunk found but not yet readable (write in-flight on the primary tier).

  • LookupResult –

    RETRY — promotion started or a secondary tier is busy.

  • LookupResult –

    MISS — chunk not found in any tier, or primary is full and cannot accept a promotion.

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def lookup(
    self,
    key: OffloadKey,
    req_context: ReqContext,
    *,
    exclude_tier_idx: int | None = None,
) -> LookupResult:
    """Check whether a single chunk is offloaded and ready.

    Algorithm:
        1. Process any completed async jobs first.
        2. Query primary tier — short-circuit on hit or in-flight.
        3. On primary miss, query secondary tiers — stop on first
           hit and initiate promotion.

    Args:
        key: Chunk hash to look up.
        req_context: Per-request context.
        exclude_tier_idx: Skip this tier index during the lookup.

    Returns:
        HIT       — chunk is ready in the primary tier.
        HIT_PENDING — chunk found but not yet readable (write
                    in-flight on the primary tier).
        RETRY     — promotion started or a secondary tier is busy.
        MISS      — chunk not found in any tier, or primary is full
                    and cannot accept a promotion.

    """
    # Poll first so a promotion that finished since the last call is
    # already reflected as HIT (not stale HIT_PENDING/MISS) below, and
    # so chunks freed by cascade or promotion completions are evictable
    # in time for a promotion this lookup may initiate.
    self._maybe_process_finished_jobs()

    start_time = time.monotonic()
    primary_hit = self.primary_tier.lookup(key, req_context)
    lookup_duration = time.monotonic() - start_time
    self._metrics.on_lookup(
        req_context,
        key,
        self._metrics.primary_tier_label,
        primary_hit,
        lookup_duration,
    )
    if primary_hit is LookupResult.HIT:
        return LookupResult.HIT
    if primary_hit is LookupResult.HIT_PENDING:
        return LookupResult.HIT_PENDING

    any_retry = False
    for i, tier in enumerate(self.secondary_tiers):
        if i == exclude_tier_idx:
            continue
        if not req_context.load_tier_filter.allows(tier.medium, tier.locality):
            continue
        labelvalues = self._metrics.tier_label(i)
        start_time = time.monotonic()
        result = tier.lookup(key, req_context)
        lookup_duration = time.monotonic() - start_time
        if result is LookupResult.HIT:
            self._metrics.on_lookup(
                req_context,
                key,
                labelvalues,
                result,
                lookup_duration,
            )
            promoted = self._initiate_promotion(i, key, req_context)
            return LookupResult.MISS if not promoted else LookupResult.HIT_PENDING
        if result is LookupResult.RETRY:
            any_retry = True
        self._metrics.on_lookup(
            req_context,
            key,
            labelvalues,
            result,
            lookup_duration,
        )

    if any_retry:
        return LookupResult.RETRY
    return LookupResult.MISS

on_new_request(req_context, *, exclude_tier_idx=None)

Query each secondary tier for its offload policy preference.

Returns REQUEST_LEVEL if ANY secondary tier wants request-level. Only stores REQUEST_LEVEL tier decisions for use in prepare_store.

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def on_new_request(
    self,
    req_context: ReqContext,
    *,
    exclude_tier_idx: int | None = None,
) -> RequestOffloadingContext:
    """Query each secondary tier for its offload policy preference.

    Returns REQUEST_LEVEL if ANY secondary tier wants request-level.
    Only stores REQUEST_LEVEL tier decisions for use in prepare_store.
    """
    state = RequestState(req_context=req_context)
    self._metrics.on_new_request(req_context)
    for tier_idx, tier in enumerate(self.secondary_tiers):
        if tier_idx == exclude_tier_idx:
            continue
        tier_ctx = tier.on_new_request(req_context)
        if tier_ctx.policy == OffloadPolicy.REQUEST_LEVEL:
            if state.request_level_tiers is None:
                state.request_level_tiers = set()
            state.request_level_tiers.add(tier_idx)
    self._req_state[req_context.req_id] = state

    policy = (
        OffloadPolicy.REQUEST_LEVEL
        if state.request_level_tiers
        else OffloadPolicy.CHUNK_LEVEL
    )
    return RequestOffloadingContext(policy=policy)

on_schedule_end(context)

End-of-schedule hook: process finished jobs, flush deferred promotions, and reset the per-step gate.

Called once per scheduler step from OffloadingConnectorScheduler.build_connector_meta().

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def on_schedule_end(self, context: ScheduleEndContext) -> None:
    """End-of-schedule hook: process finished jobs, flush deferred
    promotions, and reset the per-step gate.

    Called once per scheduler step from
    OffloadingConnectorScheduler.build_connector_meta().
    """
    # Catch-all poll: guarantees jobs are processed even on steps where
    # lookup()/prepare_store() were never called (e.g. no requests
    # scheduled but a tier still has_pending_work()).
    self._maybe_process_finished_jobs()

    for tier in self.secondary_tiers:
        tier.serve_external_requests(self._tier_parents[tier])

    # Reset the per-step gate AFTER serve_external_requests so that
    # lookup() calls within it skip redundant _process_finished_jobs().
    self._processed_jobs_this_step = False

    self._flush_pending_promotions()
    self._flush_pending_cascades()
    for tier in self.secondary_tiers:
        tier.on_schedule_end(context)

    for req_id in context.new_req_ids:
        state = self._req_state.get(req_id)
        if state is None:
            continue
        self._metrics.on_request_allocated(state.req_context)

prepare_load(keys, req_context)

Prepare chunks to be loaded from primary tier to GPU.

Callers only pass keys already confirmed HIT by lookup() earlier this step.

This increments ref_cnt on the chunks in the primary tier, protecting them from eviction during the transfer.

Parameters:

  • keys

    (Collection[OffloadKey]) –

    Chunks to prepare for loading.

  • req_context

    (ReqContext) –

    Per-request context.

Returns:

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def prepare_load(
    self, keys: Collection[OffloadKey], req_context: ReqContext
) -> LoadStoreSpec:
    """Prepare chunks to be loaded from primary tier to GPU.

    Callers only pass keys already confirmed HIT by lookup() earlier this
    step.

    This increments ref_cnt on the chunks in the primary tier, protecting
    them from eviction during the transfer.

    Args:
        keys: Chunks to prepare for loading.
        req_context: Per-request context.

    Returns:
        LoadStoreSpec for reading from primary tier.

    """
    return self.primary_tier.prepare_load(keys, req_context)

prepare_store(keys, req_context)

Prepare chunks to be stored from GPU to primary tier.

CRITICAL: This method calls _maybe_process_finished_jobs() FIRST to ensure that any completed async transfers have their ref_cnt decremented before the primary tier makes eviction decisions.

For request-level tiers, chunks already present in the primary tier are immediately cascaded via submit_store().

Parameters:

  • keys

    (Collection[OffloadKey]) –

    Chunks to prepare for storing.

  • req_context

    (ReqContext) –

    Per-request context.

Returns:

  • PrepareStoreOutput | None –

    PrepareStoreOutput describing where to store chunks and what was

  • PrepareStoreOutput | None –

    evicted, or None if store cannot proceed.

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def prepare_store(
    self, keys: Collection[OffloadKey], req_context: ReqContext
) -> PrepareStoreOutput | None:
    """Prepare chunks to be stored from GPU to primary tier.

    CRITICAL: This method calls _maybe_process_finished_jobs() FIRST to ensure
    that any completed async transfers have their ref_cnt decremented
    before the primary tier makes eviction decisions.

    For request-level tiers, chunks already present in the primary tier
    are immediately cascaded via submit_store().

    Args:
        keys: Chunks to prepare for storing.
        req_context: Per-request context.

    Returns:
        PrepareStoreOutput describing where to store chunks and what was
        evicted, or None if store cannot proceed.

    """
    # Step 1: Poll for completed async jobs FIRST
    # _process_finished_jobs() handles two kinds of completions here:
    #  - Cascade completions (store to a secondary tier, either a local
    #    cascade or a store job created for a remote requester via
    #    create_store_job()): decrements ref_cnt on the primary chunks
    #    that were read, making them evictable again once ref_cnt hits 0.
    #  - Promotion completions (secondary->primary loads): sets a
    #    not-yet-ready chunk's ref_cnt from -1 to 0 via complete_write(),
    #    making it evictable for the first time.
    # Both must be accounted for before the eviction decision below.
    self._maybe_process_finished_jobs()

    # Step 2: Store to primary tier (new chunks only).
    # Cascading of these newly-stored chunks to ALL secondary tiers
    # happens later in complete_store(), after the GPU→Primary transfer
    # completes.
    primary_result = self.primary_tier.prepare_store(keys, req_context)

    if primary_result is None:
        return None

    if primary_result.keys_to_store:
        state = self._req_state[req_context.req_id]
        state.pending_primary_stores += 1

    # Step 3: For request-level tiers, cascade chunks already in primary
    request_level_tiers = self._req_state[req_context.req_id].request_level_tiers
    if request_level_tiers:
        keys_to_store_set = set(primary_result.keys_to_store)
        keys_already_in_primary = tuple(
            k for k in keys if k not in keys_to_store_set
        )
        if keys_already_in_primary:
            self._cascade_existing_chunks_to_request_level_tiers(
                keys_already_in_primary, req_context, request_level_tiers
            )

    return primary_result

reset_cache()

Reset transfer bookkeeping and primary-tier cache.

Called during sleep, weight update, or resume. Each secondary tier drains its in-flight transfers via drain_jobs() so no tier I/O is touching primary memory before the primary tier is reset. A stuck tier will block here visibly — preferable to silent corruption from reusing primary slots while a transfer is mid-copy.

Secondary tiers are intentionally not reset: persistent stores (FS, network) keep their data across resets. Active request state is retained so those requests can continue after the reset; finished requests are finalized and removed.

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def reset_cache(self) -> None:
    """Reset transfer bookkeeping and primary-tier cache.

    Called during sleep, weight update, or resume. Each secondary tier
    drains its in-flight transfers via drain_jobs() so no tier I/O is
    touching primary memory before the primary tier is reset. A stuck
    tier will block here visibly — preferable to silent corruption
    from reusing primary slots while a transfer is mid-copy.

    Secondary tiers are intentionally not reset: persistent stores
    (FS, network) keep their data across resets. Active request state is
    retained so those requests can continue after the reset; finished
    requests are finalized and removed.
    """
    for tier in self.secondary_tiers:
        tier.drain_jobs()
    # All tier I/O has stopped; consume their completion notifications
    # so manager bookkeeping is consistent before the primary reset.
    self._process_finished_jobs()
    assert not self._jobs

    # Deferred promotion submissions reserve primary slots that the
    # reset below invalidates; their submit_load() has not yet been
    # called so no tier I/O is touching that memory.
    self._pending_load_submissions.clear()
    self._metrics.assert_idle()

    finished_req_ids = []
    for req_id, state in self._req_state.items():
        state.pending_primary_stores = 0
        state.pending_cascade_keys.clear()
        if not state.is_finished:
            continue
        for tier in self.secondary_tiers:
            tier.on_request_finished(state.req_context)
        self._metrics.on_request_finished(state.req_context)
        finished_req_ids.append(req_id)

    self.primary_tier.reset_cache()

    for req_id in finished_req_ids:
        del self._req_state[req_id]
    self._processed_jobs_this_step = False

    for tier in self.secondary_tiers:
        if tier.bp_detector is not None:
            tier.bp_detector.reset()

shutdown()

Shut down secondary tiers before releasing primary resources.

Every secondary tier is given a shutdown attempt. If any shutdown fails, preserve the primary mmap because a failed tier may still use it.

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def shutdown(self) -> None:
    """Shut down secondary tiers before releasing primary resources.

    Every secondary tier is given a shutdown attempt. If any shutdown
    fails, preserve the primary mmap because a failed tier may still use it.
    """
    shutdown_error: Exception | None = None
    for tier_idx, tier in enumerate(self.secondary_tiers):
        try:
            tier.shutdown()
        except Exception as exc:
            shutdown_error = exc
            logger.exception(
                "Failed to shut down secondary tier #%d "
                "(tier_type=%s, impl_class=%s)",
                tier_idx,
                tier.tier_type,
                type(tier).__name__,
            )

    if shutdown_error is not None:
        raise shutdown_error

    self.primary_tier.shutdown()

take_events()

Yield events owned by the primary and secondary tiers.

Yields:

  • Iterable[OffloadingEvent] –

    New OffloadingEvents collected by each tier since the last call.

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def take_events(self) -> Iterable[OffloadingEvent]:
    """Yield events owned by the primary and secondary tiers.

    Yields:
        New OffloadingEvents collected by each tier since the last call.

    """
    yield from self.primary_tier.take_events()
    for tier in self.secondary_tiers:
        yield from tier.take_events()

touch(keys, req_context)

Mark chunks as recently used in all tiers.

Parameters:

  • keys

    (Collection[OffloadKey]) –

    Chunks to mark as recently used.

  • req_context

    (ReqContext) –

    Per-request context.

Source code in vllm/v1/kv_offload/tiering/manager.py
@override
def touch(self, keys: Collection[OffloadKey], req_context: ReqContext):
    """Mark chunks as recently used in all tiers.

    Args:
        keys: Chunks to mark as recently used.
        req_context: Per-request context.

    """
    self.primary_tier.touch(keys, req_context)
    for tier in self.secondary_tiers:
        tier.touch(keys, req_context)

_SecondaryTierFacingParent

Bases: ParentManager

Wrapper that implements ParentManager by delegating to the TieringOffloadingManager with exclude_tier_idx set to the origin tier.

Source code in vllm/v1/kv_offload/tiering/manager.py
class _SecondaryTierFacingParent(ParentManager):
    """Wrapper that implements ParentManager by delegating to the
    TieringOffloadingManager with exclude_tier_idx set to the origin tier."""

    __slots__ = ("_m", "_origin_idx")

    def __init__(
        self,
        manager: "TieringOffloadingManager",
        tier_idx: int,
    ):
        self._m = manager
        self._origin_idx = tier_idx

    def on_new_request(self, req_context: ReqContext) -> RequestOffloadingContext:
        return self._m.on_new_request(req_context, exclude_tier_idx=self._origin_idx)

    def lookup(self, key: OffloadKey, req_context: ReqContext) -> LookupResult:
        return self._m.lookup(key, req_context, exclude_tier_idx=self._origin_idx)

    def create_store_job(
        self, keys: Collection[OffloadKey], req_context: ReqContext
    ) -> TransferJob:
        return self._m.create_store_job(keys, req_context, self._origin_idx)

    def on_request_finished(self, req_context: ReqContext) -> None:
        return self._m.on_request_finished(
            req_context, exclude_tier_idx=self._origin_idx
        )