Skip to content

vllm.v1.simple_kv_offload.manager

Scheduler-side manager for SimpleCPUOffloadConnector.

Classes:

BlockStoreMeta dataclass

Per-block metadata snapshot for BlockStored event emission.

Captured at store-prep time (when the Request is available) and carried through TransferMeta to async completion, where the Request is gone.

Source code in vllm/v1/simple_kv_offload/manager.py
@dataclass
class BlockStoreMeta:
    """Per-block metadata snapshot for BlockStored event emission.

    Captured at store-prep time (when the Request is available) and carried
    through TransferMeta to async completion, where the Request is gone.
    """

    token_ids: list[int]
    parent_block_hash: ExternalBlockHash | None
    lora_id: int | None
    lora_name: str | None
    extra_keys: tuple[Any, ...] | None
    block_size: int | None = None

BoundaryStoreStats dataclass

Cumulative counters for boundary hand-off stores.

Every branch that declines a hand-off is silent by design: the boundary is simply not offloaded and a later request misses. Without counters the only symptom is a lower cache hit rate with nothing in the logs, so each decline reason is counted and exposed through SimpleCPUOffloadConnector.get_boundary_store_stats(). Reset by reset(), which carries any not-yet-drained deltas; not otherwise cleared.

Source code in vllm/v1/simple_kv_offload/manager.py
@dataclass
class BoundaryStoreStats:
    """Cumulative counters for boundary hand-off stores.

    Every branch that declines a hand-off is silent by design: the boundary is
    simply not offloaded and a later request misses. Without counters the only
    symptom is a lower cache hit rate with nothing in the logs, so each decline
    reason is counted and exposed through
    ``SimpleCPUOffloadConnector.get_boundary_store_stats()``. Reset by
    ``reset()``, which carries any not-yet-drained deltas; not otherwise
    cleared.
    """

    published: int = 0
    stored: int = 0
    dropped_cpu_full: int = 0
    dropped_request_gone: int = 0
    skipped_already_cached: int = 0
    skipped_in_flight: int = 0
    dropped_null_block: int = 0
    dropped_not_hashed: int = 0

SimpleCPUOffloadScheduler

Scheduler-side manager for CPU offloading.

Methods:

Source code in vllm/v1/simple_kv_offload/manager.py
 141
 142
 143
 144
 145
 146
 147
 148
 149
 150
 151
 152
 153
 154
 155
 156
 157
 158
 159
 160
 161
 162
 163
 164
 165
 166
 167
 168
 169
 170
 171
 172
 173
 174
 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
1007
1008
1009
1010
1011
1012
1013
1014
1015
1016
1017
1018
1019
1020
1021
1022
1023
1024
1025
1026
1027
1028
1029
1030
1031
1032
1033
1034
1035
1036
1037
1038
1039
1040
1041
1042
1043
1044
1045
1046
1047
1048
1049
1050
1051
1052
1053
1054
1055
1056
1057
1058
1059
1060
1061
1062
1063
1064
1065
1066
1067
1068
1069
1070
1071
1072
1073
1074
1075
1076
1077
1078
1079
1080
1081
1082
1083
1084
1085
1086
1087
1088
1089
1090
1091
1092
1093
1094
1095
1096
1097
1098
1099
1100
1101
1102
1103
1104
1105
1106
1107
1108
1109
1110
1111
1112
1113
1114
1115
1116
1117
1118
1119
1120
1121
1122
1123
1124
1125
1126
1127
1128
1129
1130
1131
1132
1133
1134
1135
1136
1137
1138
1139
1140
1141
1142
1143
1144
1145
1146
1147
1148
1149
1150
1151
1152
1153
1154
1155
1156
1157
1158
1159
1160
1161
1162
1163
1164
1165
1166
1167
1168
1169
1170
1171
1172
1173
1174
1175
1176
1177
1178
1179
1180
1181
1182
1183
1184
1185
1186
1187
1188
1189
1190
1191
1192
1193
1194
1195
1196
1197
1198
1199
1200
1201
1202
1203
1204
1205
1206
1207
1208
1209
1210
1211
1212
1213
1214
1215
1216
1217
1218
1219
1220
1221
1222
1223
1224
1225
1226
1227
1228
1229
1230
1231
1232
1233
1234
1235
1236
1237
1238
1239
1240
1241
1242
1243
1244
1245
1246
1247
1248
1249
1250
1251
1252
1253
1254
1255
1256
1257
1258
1259
1260
1261
1262
1263
1264
1265
1266
1267
1268
1269
1270
1271
1272
1273
1274
1275
1276
1277
1278
1279
1280
1281
1282
1283
1284
1285
1286
1287
1288
1289
1290
1291
1292
1293
1294
1295
1296
1297
1298
1299
1300
1301
1302
1303
1304
1305
1306
1307
1308
1309
1310
1311
1312
1313
1314
1315
1316
1317
1318
1319
1320
1321
1322
1323
1324
1325
1326
1327
1328
1329
1330
1331
1332
1333
1334
1335
1336
1337
1338
1339
1340
1341
1342
1343
1344
1345
1346
1347
1348
1349
1350
1351
1352
1353
1354
1355
1356
1357
1358
1359
1360
1361
1362
1363
1364
1365
1366
1367
1368
1369
1370
1371
1372
1373
1374
1375
1376
1377
1378
1379
1380
1381
1382
1383
1384
1385
1386
1387
1388
1389
1390
1391
1392
1393
1394
1395
1396
1397
1398
1399
1400
1401
1402
1403
1404
1405
1406
1407
1408
1409
1410
1411
1412
1413
1414
1415
1416
1417
1418
1419
1420
1421
1422
1423
1424
1425
1426
1427
1428
1429
1430
1431
1432
1433
1434
1435
1436
1437
1438
1439
1440
1441
1442
1443
1444
1445
class SimpleCPUOffloadScheduler:
    """Scheduler-side manager for CPU offloading."""

    def __init__(
        self,
        vllm_config: VllmConfig,
        kv_cache_config: "KVCacheConfig | None",
        cpu_capacity_bytes: int,
        scheduler_block_size: int,
        hash_block_size: int,
        lazy_offload: bool = False,
        disk_capacity_bytes: int = 0,
        use_page_cache: bool = False,
    ):
        self.vllm_config = vllm_config
        self.kv_cache_config = kv_cache_config
        # When disk mode is active, the offload pool size is disk-based.
        offload_capacity = (
            disk_capacity_bytes if disk_capacity_bytes > 0 else cpu_capacity_bytes
        )
        self.enable_kv_cache_events = (
            vllm_config.kv_events_config is not None
            and vllm_config.kv_events_config.enable_kv_cache_events
        )
        dcp_world_size = vllm_config.parallel_config.decode_context_parallel_size
        self.cp_world_size = dcp_world_size
        self.block_size = scheduler_block_size
        self.hash_block_size = hash_block_size
        assert self.block_size % self.hash_block_size == 0
        # Derive a CPU KVCacheConfig from the GPU config and build a coordinator
        assert kv_cache_config is not None
        self.cpu_kv_cache_config = self._derive_cpu_config(
            kv_cache_config, offload_capacity
        )
        self.num_cpu_blocks = self.cpu_kv_cache_config.num_blocks
        self.prefix_cacheable_group_ids = (
            self.cpu_kv_cache_config.prefix_cacheable_group_ids
        )
        self.kv_event_medium = MEDIUM_STORAGE if disk_capacity_bytes > 0 else MEDIUM_CPU
        # Must track kv_event_medium.
        self._info_labelvalues: tuple[str, ...] = (
            "disk" if disk_capacity_bytes > 0 else "cpu",
            str(use_page_cache).lower(),
            str(lazy_offload).lower(),
            str(self.num_cpu_blocks),
        )
        # Find the full attention kv group for prefix cache matching.
        self.fa_gidx = -1
        for g_idx, g in enumerate(self.cpu_kv_cache_config.kv_cache_groups):
            if isinstance(g.kv_cache_spec, FullAttentionSpec):
                self.fa_gidx = g_idx
                break
        assert 0 <= self.fa_gidx < len(self.cpu_kv_cache_config.kv_cache_groups)
        logger.info(
            "SimpleCPUOffloadScheduler: Allocating %d offload blocks "
            "(%.2f GB, mode=%s, backend=%s)",
            self.num_cpu_blocks,
            offload_capacity / (1024**3),
            "lazy" if lazy_offload else "eager",
            "disk" if disk_capacity_bytes > 0 else "cpu",
        )

        spec_config = vllm_config.speculative_config
        use_eagle_block_drop = (
            spec_config is not None and spec_config.use_eagle_block_drop()
        )
        self.cpu_coordinator: KVCacheCoordinator = get_kv_cache_coordinator(
            kv_cache_config=self.cpu_kv_cache_config,
            max_model_len=vllm_config.model_config.max_model_len,
            max_in_flight_tokens=vllm_config.max_in_flight_tokens,
            use_eagle=use_eagle_block_drop,
            enable_caching=True,
            enable_kv_cache_events=self.enable_kv_cache_events,
            dcp_world_size=dcp_world_size,
            pcp_world_size=1,
            scheduler_block_size=self.block_size,
            hash_block_size=self.hash_block_size,
            allow_partial_hash_hits=not lazy_offload,
        )
        self.group_block_sizes = self.cpu_coordinator.group_block_sizes
        # FA group's own resolved block_size; divides scheduler_block_size (the
        # LCM) but is NOT assumed to equal it.
        self.fa_block_size: int = self.group_block_sizes[self.fa_gidx]
        assert self.block_size % self.fa_block_size == 0
        self.cpu_block_pool: BlockPool = self.cpu_coordinator.block_pool
        # GPU block pool reference - bound after scheduler builds kv_cache_manager
        self._gpu_block_pool: BlockPool | None = None

        # Load metadata
        self._reqs_to_load: dict[str, LoadRequestState] = {}
        # Inverse map: load_event_idx -> req_ids. Keyed by load_event_idx because
        # the worker reports completions by event index, not request id.
        self._load_event_to_reqs: dict[int, list[str]] = {}

        # Pending (cpu_hit_blocks, hit_length, num_computed_tokens) tuples from
        # find_longest_cache_hit, kept pinned via touch() while awaiting
        # update_state_after_alloc(). ``num_computed_tokens`` is the local
        # prefix the hit was resolved against; the load path needs it to place
        # external blocks and must not re-derive it from block metadata.
        self._pending_cpu_hits: dict[
            str, tuple[tuple[list[KVCacheBlock], ...], int, int]
        ] = {}

        # Store metadata
        self._lazy_mode = lazy_offload
        # Lazy mode: use a cursor to track the last scanned block in the GPU free queue.
        self._cursor: KVCacheBlock | None = None
        if self._lazy_mode:
            self._target_free = self._estimate_lazy_target_blocks(
                kv_cache_config,
                vllm_config.scheduler_config.max_num_batched_tokens,
                self.cp_world_size,
            )
        else:
            self._target_free = 0
        self._store_event_to_blocks: dict[int, TransferMeta] = {}
        self._abandoned_store_event_to_blocks: dict[int, TransferMeta] = {}
        # Eager mode only
        self._reqs_to_store: dict[str, StoreRequestState] = {}
        self._store_event_to_reqs: dict[int, list[str]] = {}
        self._in_flight_store_gpu_blocks: set[int] = set()
        self._pending_finished_stores: list[TransferMeta] = []
        self._abandoned_reqs_to_load: dict[str, LoadRequestState] = {}

        # Event counters
        self._load_event_counter: int = 0
        self._store_event_counter: int = 0
        self.boundary_store_stats = BoundaryStoreStats()
        # Interval stats state drained by get_stats()
        self._boundary_stats_snapshot = BoundaryStoreStats()
        self._interval_load_blocks_completed = 0

        # For TP/PP: track partial store completions across steps.
        # Events must be reported by all world_size workers before considered complete.
        self._expected_worker_count = vllm_config.parallel_config.world_size
        self._store_event_pending_counts: dict[int, int] = {}

    @staticmethod
    def _derive_cpu_config(
        gpu_config: "KVCacheConfig", cpu_capacity_bytes: int
    ) -> "KVCacheConfig":
        """Derive a CPU KVCacheConfig from the GPU config.
        Same kv_cache_groups, num_blocks scaled by CPU/GPU memory ratio."""
        # Import here to avoid potential circular imports
        from vllm.v1.kv_cache_interface import KVCacheTensor

        assert len(gpu_config.kv_cache_tensors) > 0

        # Every KVCacheTensor describes placement within the same backing allocation,
        # so its size is the total GPU KV cache size.
        gpu_total_bytes = gpu_config.kv_cache_tensors[0].size
        num_gpu_blocks = gpu_config.num_blocks
        num_cpu_blocks = max(1, num_gpu_blocks * cpu_capacity_bytes // gpu_total_bytes)
        # Create CPU kv_cache_tensors mirroring GPU by scaling size proportionally.
        cpu_tensors = [
            KVCacheTensor(
                size=t.size // num_gpu_blocks * num_cpu_blocks,
                layers=list(t.layers),
                layer_stride=t.layer_stride,
                block_stride=t.block_stride,
                offset=t.offset,
            )
            for t in gpu_config.kv_cache_tensors
        ]

        return replace(
            gpu_config,
            num_blocks=num_cpu_blocks,
            kv_cache_tensors=cpu_tensors,
        )

    @staticmethod
    def _estimate_lazy_target_blocks(
        kv_cache_config: "KVCacheConfig",
        max_num_batched_tokens: int,
        cp_world_size: int = 1,
    ) -> int:
        """GPU blocks to keep available (free/offloaded) per step in lazy mode."""
        WATERMARK_RATIO = 1.0  # Reserve larger space to avoid running out of GPU blocks
        target = 0
        for g in kv_cache_config.prefix_cacheable_groups:
            spec = g.kv_cache_spec
            block_size = resolve_dcp_kv_block_size(spec, cp_world_size)
            if isinstance(spec, MambaSpec):
                target += 2
            elif isinstance(spec, SlidingWindowSpec):
                target += cdiv(spec.sliding_window, block_size) + 1
            else:
                target += cdiv(max_num_batched_tokens, block_size)
        return int(target * (1 + WATERMARK_RATIO))

    def bind_gpu_block_pool(self, gpu_block_pool: BlockPool) -> None:
        """Bind GPU block pool so that we can touch blocks during stores.
        Called by Scheduler after kv_cache_manager is ready."""
        self._gpu_block_pool = gpu_block_pool

    def get_num_new_matched_tokens(
        self, request: "Request", num_computed_tokens: int
    ) -> tuple[int | None, bool]:
        """Return (num_new_tokens, is_async) from consecutive CPU cache hits."""
        # Pins found CPU blocks so they survive LRU eviction until
        # update_state_after_alloc() consumes them. Any pin from an earlier
        # call on the same request (e.g. retry after a failed allocate_slots)
        # is dropped first.
        if stale := self._pending_cpu_hits.pop(request.request_id, None):
            self._free_pending_cpu_hit(stale)

        if request.skip_reading_prefix_cache:
            return 0, False

        if num_computed_tokens % self.block_size != 0:
            # Transfers are whole-block copies, so an external suffix cannot
            # start in the middle of a destination block. When the GPU-side
            # coordinator also serves fine-grained hits, the local prefix can
            # land on a hash boundary that is not a scheduler-block boundary;
            # such a request keeps its local hit and skips the external one.
            logger.warning_once(
                "SimpleCPU external lookup requires scheduler-block-aligned "
                "local tokens, got %d tokens with scheduler block size %d.",
                num_computed_tokens,
                self.block_size,
            )
            return 0, False

        num_skipped_hashes = num_computed_tokens // self.hash_block_size
        remaining_hashes = request.block_hashes[num_skipped_hashes:]

        if not remaining_hashes:
            return 0, False
        # Must recompute at least the last token, matching the logic in
        # kv_cache_manager.get_computed_blocks().
        max_hit_len = request.num_tokens - 1 - num_computed_tokens
        if max_hit_len <= 0:
            return 0, False
        cpu_hit_blocks, hit_length, _ = self.cpu_coordinator.find_longest_cache_hit(
            remaining_hashes, max_hit_len
        )

        if hit_length > 0:
            pin_blocks = [
                blk for grp in cpu_hit_blocks for blk in grp if not blk.is_null
            ]
            self.cpu_block_pool.touch(pin_blocks)
            self._pending_cpu_hits[request.request_id] = (
                cpu_hit_blocks,
                hit_length,
                num_computed_tokens,
            )
            return hit_length, True
        return 0, False

    # TODO(yifan): this API now only matches the suffix part of the prefix cache. A more
    # general API should scan blocks in both GPU and CPU block pool in a single pass.
    def update_state_after_alloc(
        self,
        request: "Request",
        blocks: "KVCacheBlocks",
        num_external_tokens: int,
    ) -> None:
        req_id = request.request_id
        block_ids_by_group = blocks.get_block_ids()
        num_groups = len(block_ids_by_group)

        # Store tracking (eager mode only). Register the request;
        # block IDs are accumulated from scheduler_output in
        # _prepare_eager_store_specs via yield_req_data.
        if not self._lazy_mode and req_id not in self._reqs_to_store:
            self._reqs_to_store[req_id] = StoreRequestState(
                request=request,
                block_ids=tuple([] for _ in range(num_groups)),
                num_stored_blocks=[0] * num_groups,
            )

        # Pop the CPU hit cached by get_num_new_matched_tokens(). The
        # found blocks were pinned there to survive LRU eviction in the window
        # between get_num_new_matched_tokens() and this matching call.
        pending = self._pending_cpu_hits.pop(req_id, None)

        if num_external_tokens == 0:
            if pending is not None:
                logger.warning(
                    "SimpleCPUOffloadScheduler: update_state_after_alloc "
                    "called for req_id=%s with no external tokens but "
                    "get_num_new_matched_tokens() unexpectedly recorded "
                    "a pending CPU hit; releasing the stale pin.",
                    req_id,
                )
                self._free_pending_cpu_hit(pending)
            return

        if pending is None:
            logger.warning(
                "SimpleCPUOffloadScheduler: update_state_after_alloc called "
                "for req_id=%s with num_external_tokens=%d but no pending "
                "CPU hit from get_num_new_matched_tokens(); skipping load.",
                req_id,
                num_external_tokens,
            )
            return

        cpu_hit_blocks_full, _, num_computed_tokens = pending

        assert num_computed_tokens % self.block_size == 0, (
            "SimpleCPU external loads must start at a scheduler-block boundary"
        )
        # Fine-grained hits are hash-block aligned, not group-block aligned, so
        # the per-group block count below rounds up instead of dividing exactly.
        # Hash alignment is still required: it is what makes the accepted
        # external suffix start on a boundary every group can address.
        assert num_external_tokens % self.hash_block_size == 0, (
            f"num_external_tokens={num_external_tokens} is not aligned to "
            f"hash_block_size={self.hash_block_size}"
        )

        # The scheduler may have accepted fewer blocks than
        # get_num_new_matched_tokens() reported.
        # (e.g. due to token budget in test_partial_gpu_prefix_plus_cpu_load).
        # Take only the leading N blocks per group matching num_external_tokens;
        # the rest will be released along with the temp pin below.
        cpu_hit_blocks: list[list[KVCacheBlock]] = []
        for g in range(num_groups):
            if g not in self.prefix_cacheable_group_ids:
                cpu_hit_blocks.append([])
                continue
            g_block_size = self.group_block_sizes[g]
            n_take_g = cdiv(num_external_tokens, g_block_size)
            cpu_hit_blocks.append(cpu_hit_blocks_full[g][:n_take_g])

        gpu_block_ids: list[int] = []
        cpu_block_ids: list[int] = []
        cpu_blocks_to_touch: list[KVCacheBlock] = []

        for g in range(num_groups):
            cpu_blocks_g = cpu_hit_blocks[g]
            n_ext_g = len(cpu_blocks_g)
            if n_ext_g == 0:
                continue

            g_block_size = self.group_block_sizes[g]
            gpu_ext_start = num_computed_tokens // g_block_size
            group_gpu_ids = block_ids_by_group[g]

            for i, cpu_blk in enumerate(cpu_blocks_g):
                # Skip null blocks (e.g. sliding window or mamba padding).
                if cpu_blk.is_null:
                    continue
                gpu_block_ids.append(group_gpu_ids[gpu_ext_start + i])
                cpu_block_ids.append(cpu_blk.block_id)
                cpu_blocks_to_touch.append(cpu_blk)

        # Touch CPU blocks to prevent eviction during async load.
        self.cpu_block_pool.touch(cpu_blocks_to_touch)
        # Release the temporary pin held since get_num_new_matched_tokens().
        self._free_pending_cpu_hit(pending)

        # Touch GPU blocks to prevent freeing during async load
        assert self._gpu_block_pool is not None
        self._gpu_block_pool.touch(
            [self._gpu_block_pool.blocks[bid] for bid in gpu_block_ids]
        )

        assert self._reqs_to_load.get(req_id) is None
        self._reqs_to_load[req_id] = LoadRequestState(
            request=request, transfer_meta=TransferMeta(gpu_block_ids, cpu_block_ids)
        )

    def build_connector_meta(
        self,
        scheduler_output: SchedulerOutput,
    ) -> SimpleCPUOffloadMetadata:
        # --- Stores ---
        store_event = -1
        store_gpu, store_cpu, store_req_ids, store_meta = self.prepare_store_specs(
            scheduler_output
        )
        for transfer in self._pending_finished_stores:
            store_gpu.extend(transfer.gpu_block_ids)
            store_cpu.extend(transfer.cpu_block_ids)
            if store_meta is not None:
                assert transfer.block_meta is not None
                store_meta.extend(transfer.block_meta)
        self._pending_finished_stores.clear()

        if store_gpu:
            store_event = self._store_event_counter
            self._store_event_counter += 1
            self._store_event_to_blocks[store_event] = TransferMeta(
                store_gpu, store_cpu, store_meta
            )
            if store_req_ids:  # For eager mode only, track req->blocks mapping
                self._store_event_to_reqs[store_event] = store_req_ids
                for req_id in store_req_ids:
                    store_state = self._reqs_to_store.get(req_id)
                    if store_state is not None:
                        store_state.store_events.add(store_event)

        # --- Loads ---
        load_event = -1
        load_gpu: list[int] = []
        load_cpu: list[int] = []
        load_req_ids: list[str] = []
        for req_id, load_state in self._reqs_to_load.items():
            if load_state.load_event is not None:
                continue
            assert load_state.transfer_meta is not None
            load_gpu.extend(load_state.transfer_meta.gpu_block_ids)
            load_cpu.extend(load_state.transfer_meta.cpu_block_ids)
            load_req_ids.append(req_id)
        if load_req_ids:
            load_event = self._load_event_counter
            self._load_event_counter += 1
            for req_id in load_req_ids:
                self._reqs_to_load[req_id].load_event = load_event
            self._load_event_to_reqs[load_event] = load_req_ids

        result = SimpleCPUOffloadMetadata(
            load_event=load_event,
            load_gpu_blocks=load_gpu,
            load_cpu_blocks=load_cpu,
            load_event_to_reqs={
                event_idx: list(req_ids)
                for event_idx, req_ids in self._load_event_to_reqs.items()
            },
            store_event=store_event,
            store_gpu_blocks=store_gpu,
            store_cpu_blocks=store_cpu,
            need_flush=bool(scheduler_output.preempted_req_ids),
        )
        return result

    def prepare_store_specs(
        self, scheduler_output: SchedulerOutput
    ) -> tuple[
        list[int],
        list[int],
        list[str],
        list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None,
    ]:
        """Prepare store specs for the store event."""
        if self._lazy_mode:
            return self._prepare_lazy_store_specs()
        else:
            return self._prepare_eager_store_specs(scheduler_output)

    def _prepare_lazy_store_specs(
        self,
    ) -> tuple[
        list[int],
        list[int],
        list[str],
        list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None,
    ]:
        """Single-pass cursor walk: offload cached GPU blocks near eviction.

        Walks the GPU free queue from the cursor, counting blocks that are
        free-or-offloaded (safe for the allocator to evict). Stops when
        target_free blocks are covered or CPU capacity is reached.
        """
        gpu_pool = self._gpu_block_pool
        if gpu_pool is None or self._target_free <= 0:
            return [], [], [], None

        free_queue = gpu_pool.free_block_queue
        cpu_pool = self.cpu_block_pool
        num_cpu_free = cpu_pool.get_num_free_blocks()

        # Validate cursor: stale if block was removed from free queue.
        if self._cursor is not None and self._cursor.ref_cnt > 0:
            self._cursor = None

        gpu_ids: list[int] = []
        last_visited = self._cursor

        for covered, node in enumerate(free_queue.iter_blocks_after(self._cursor)):
            if covered >= self._target_free or len(gpu_ids) >= num_cpu_free:
                break

            last_visited = node
            bhash = node.block_hash

            if (
                bhash is not None
                and not node.is_null
                and cpu_pool.cached_block_hash_to_block.get_one_block(bhash) is None
            ):
                gpu_ids.append(node.block_id)

        self._cursor = last_visited

        # Batch-allocate CPU blocks.
        if gpu_ids:
            cpu_blocks = cpu_pool.get_new_blocks(len(gpu_ids))
            cpu_ids = [blk.block_id for blk in cpu_blocks]
            # Touch GPU blocks to prevent eviction during async copy.
            gpu_pool.touch([gpu_pool.blocks[bid] for bid in gpu_ids])
        else:
            cpu_ids = []

        return gpu_ids, cpu_ids, [], None

    def _prepare_eager_store_specs(
        self, scheduler_output: SchedulerOutput
    ) -> tuple[
        list[int],
        list[int],
        list[str],
        list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None,
    ]:
        """Identify newly computed blocks to offload from scheduler requests.

        Only considers blocks whose KV data has been **confirmed computed** by
        the GPU. Blocks from the current step are stored on a later step, or by
        the finish-time flush if the request completes first.

        Returns:
            (gpu_block_ids, cpu_block_ids, req_ids, block_meta) for the store
            event. ``block_meta`` is None when kv cache events are disabled.

        """
        merged_gpu_block_ids: list[int] = []
        merged_cpu_block_ids: list[int] = []
        req_ids: list[str] = []
        merged_block_meta: list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None = (
            [] if self.enable_kv_cache_events else None
        )

        gpu_block_pool = self._gpu_block_pool
        if gpu_block_pool is None:
            return [], [], [], merged_block_meta
        cpu_block_pool = self.cpu_block_pool
        num_groups = len(self.cpu_kv_cache_config.kv_cache_groups)
        # Dedup against blocks already scheduled.
        in_flight = self._in_flight_store_gpu_blocks
        num_free = cpu_block_pool.get_num_free_blocks()

        block_state = scheduler_output.kv_connector_block_state
        boundary_offloads = (
            block_state.boundary_state_offloads if block_state is not None else {}
        )
        stats = self.boundary_store_stats
        preempted_req_ids = scheduler_output.preempted_req_ids or ()
        for req_id, entries in boundary_offloads.items():
            # The scheduler drains handoffs without filtering by request
            # liveness. A request that finished or was preempted in this step
            # is giving up its blocks, which may already have been reallocated
            # to another request; reading them would publish unrelated KV under
            # a valid hash. Drop conservatively, as the mooncake store does.
            store_state = self._reqs_to_store.get(req_id)
            if (
                store_state is None
                or store_state.finished
                or req_id in preempted_req_ids
                or req_id in scheduler_output.finished_req_ids
            ):
                stats.published += len(entries)
                stats.dropped_request_gone += len(entries)
                continue
            scheduled_for_req = False
            # ``boundary_tokens`` is not needed here: the cache key comes from
            # the handed-off block itself, not from the offered boundary.
            for _, gpu_block_id, _ in entries:
                stats.published += 1
                gpu_block = gpu_block_pool.blocks[gpu_block_id]
                admission = self._classify_store_candidate(gpu_block)
                if admission is _StoreAdmission.NULL_BLOCK:
                    stats.dropped_null_block += 1
                    continue
                if admission is _StoreAdmission.NOT_HASHED:
                    stats.dropped_not_hashed += 1
                    continue
                if admission is _StoreAdmission.IN_FLIGHT:
                    stats.skipped_in_flight += 1
                    continue
                if admission is _StoreAdmission.ALREADY_CACHED:
                    stats.skipped_already_cached += 1
                    continue
                if num_free <= 0:
                    stats.dropped_cpu_full += 1
                    logger.warning_once(
                        "SimpleCPU dropped a boundary-state handoff because "
                        "the CPU block pool is full; this boundary will not be retried."
                    )
                    continue
                cpu_block = cpu_block_pool.get_new_blocks(1)[0]
                merged_gpu_block_ids.append(gpu_block_id)
                merged_cpu_block_ids.append(cpu_block.block_id)
                if merged_block_meta is not None:
                    # Keep the metadata list index-parallel with the block ids.
                    # The hand-off block's own hashes are resolved at completion,
                    # so there is nothing to capture here; emit the event without
                    # a payload, as the size-mismatch fallback already does.
                    merged_block_meta.append({})
                in_flight.add(gpu_block_id)
                gpu_block_pool.touch([gpu_block])
                num_free -= 1
                stats.stored += 1
                scheduled_for_req = True
            if scheduled_for_req:
                req_ids.append(req_id)

        for req_id, new_block_id_groups, preempted in yield_req_data(scheduler_output):
            state = self._reqs_to_store.get(req_id)
            if state is None or state.finished:
                continue

            if preempted:
                state.block_ids = tuple([] for _ in range(num_groups))
                state.num_stored_blocks = [0] * num_groups
            if new_block_id_groups:
                for g in range(min(num_groups, len(new_block_id_groups))):
                    if new_block_id_groups[g] is not None:
                        # Accumulate new block IDs.
                        state.block_ids[g].extend(new_block_id_groups[g])

            num_new_tokens = scheduler_output.num_scheduled_tokens.get(req_id, 0)
            if num_new_tokens == 0:
                continue

            block_ids_by_group = state.block_ids
            if not block_ids_by_group:
                continue

            gpu_block_ids, advanced_per_group, block_meta = (
                self._select_eager_blocks_to_store(state, block_ids_by_group)
            )

            # Batch allocate the CPU destinations.
            n_to_alloc = len(gpu_block_ids)
            if n_to_alloc > 0:
                cpu_blocks_alloc = cpu_block_pool.get_new_blocks(n_to_alloc)
                cpu_block_ids = [blk.block_id for blk in cpu_blocks_alloc]
            else:
                cpu_block_ids = []

            if cpu_block_ids:
                req_ids.append(req_id)
                merged_gpu_block_ids.extend(gpu_block_ids)
                merged_cpu_block_ids.extend(cpu_block_ids)
                in_flight.update(gpu_block_ids)
                if merged_block_meta is not None:
                    assert block_meta is not None
                    merged_block_meta.extend(block_meta)

                # Touch GPU blocks to prevent freeing during async copy
                gpu_block_pool.touch(
                    [gpu_block_pool.blocks[bid] for bid in gpu_block_ids]
                )

                logger.debug(
                    "Request %s: Scheduling store of %d blocks to CPU (%d groups)",
                    req_id,
                    len(cpu_block_ids),
                    num_groups,
                )

            # Advance per-group cursors (includes cached hits + newly stored)
            for g in range(num_groups):
                state.num_stored_blocks[g] += advanced_per_group[g]

        # A request can contribute both boundary and positional blocks to the
        # same event. Completion tracking is set-based, so keep one request ID.
        req_ids = list(dict.fromkeys(req_ids))
        return merged_gpu_block_ids, merged_cpu_block_ids, req_ids, merged_block_meta

    def _cached_gpu_block(
        self, resolved_hashes: "BlockHashList", block_idx: int, group_id: int
    ) -> "KVCacheBlock | None":
        """Return the GPU block still cached under this group's block hash."""
        if block_idx >= len(resolved_hashes):
            return None
        assert self._gpu_block_pool is not None
        blocks = self._gpu_block_pool.get_cached_block(
            resolved_hashes[block_idx], [group_id]
        )
        return blocks[0] if blocks else None

    def _select_eager_blocks_to_store(
        self,
        state: StoreRequestState,
        block_ids_by_group: tuple[list[int], ...],
    ) -> tuple[
        list[int],
        list[int],
        list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None,
    ]:
        """Return confirmed eager blocks, cursor advances, and event metadata."""
        assert self._gpu_block_pool is not None
        gpu_block_ids: list[int] = []
        advanced_per_group = [0] * len(self.cpu_kv_cache_config.kv_cache_groups)
        block_meta: list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None = (
            [] if self.enable_kv_cache_events else None
        )
        request = state.request
        confirmed_tokens = request.num_computed_tokens - request.num_output_placeholders
        # Truncate to the granularity a lookup can actually land on. With
        # fine-grained hits that is hash_block_size; otherwise hits only land on
        # the scheduler block (the group LCM). Using the LCM unconditionally
        # would drop whole blocks that sit between the last LCM boundary and a
        # reachable fine-grained boundary, so the other groups would hold that
        # boundary and this one would not, and the joint hybrid lookup would
        # reconcile to zero.
        store_alignment = (
            self.hash_block_size
            if self.cpu_coordinator.enable_partial_hash_hits
            else self.block_size
        )
        aligned_tokens = confirmed_tokens // store_alignment * store_alignment
        num_free = self.cpu_block_pool.get_num_free_blocks()

        for g, group_gpu_ids in enumerate(block_ids_by_group):
            if g not in self.prefix_cacheable_group_ids:
                continue
            if len(gpu_block_ids) >= num_free:
                break
            group_manager = self.cpu_coordinator.single_type_managers[g]
            if not group_manager.has_positionally_stable_blocks:
                continue
            # FIXME (yifan): handle CPU cache eviction, where
            # num_stored_blocks can be stale and omit evicted blocks in
            # the middle of the request.
            group_size = self.group_block_sizes[g]
            ready = min(len(group_gpu_ids), aligned_tokens // group_size)
            resolved_hashes = resolve_block_hashes(
                request.block_hashes, self.hash_block_size, group_size
            )
            curr_mm_idx = 0
            secondary_mm_idx = 0
            start = state.num_stored_blocks[g]
            for i, gpu_block_id in enumerate(group_gpu_ids[start:ready], start=start):
                gpu_block = self._gpu_block_pool.blocks[gpu_block_id]
                if gpu_block.is_null:
                    # Sliding-window groups null pages that left the window
                    # before the connector sees the block table, but the
                    # retained prefix-cache tail stays hashed in the GPU free
                    # queue; store it from there.
                    recovered = self._cached_gpu_block(resolved_hashes, i, g)
                    if recovered is None:
                        advanced_per_group[g] += 1
                        continue
                    gpu_block = recovered
                    gpu_block_id = gpu_block.block_id
                if (
                    self._classify_store_candidate(gpu_block)
                    is not _StoreAdmission.READY
                ):
                    advanced_per_group[g] += 1
                    continue
                if len(gpu_block_ids) >= num_free:
                    break
                primary_block_hash = gpu_block.block_hash
                assert primary_block_hash is not None
                gpu_block_ids.append(gpu_block_id)
                advanced_per_group[g] += 1
                if block_meta is not None:
                    token_start = i * group_size
                    token_end = token_start + group_size
                    parent_hash = (
                        None
                        if i == 0
                        else maybe_convert_block_hash(resolved_hashes[i - 1])
                    )
                    lora_req = request.lora_request
                    extra_keys, curr_mm_idx = generate_block_hash_extra_keys(
                        request, token_start, token_end, curr_mm_idx
                    )
                    meta_by_hash = {
                        primary_block_hash: BlockStoreMeta(
                            token_ids=list(
                                request.all_token_ids[token_start:token_end]
                            ),
                            parent_block_hash=parent_hash,
                            lora_id=lora_req.adapter_id if lora_req else None,
                            lora_name=lora_req.name if lora_req else None,
                            extra_keys=extra_keys,
                        )
                    }
                    first_hash_idx = token_start // self.hash_block_size
                    last_hash_idx = token_end // self.hash_block_size
                    for hash_idx in range(first_hash_idx, last_hash_idx):
                        block_hash = make_block_hash_with_group_id(
                            request.block_hashes[hash_idx], g
                        )
                        if block_hash is None or block_hash == primary_block_hash:
                            continue
                        hash_start = hash_idx * self.hash_block_size
                        hash_end = hash_start + self.hash_block_size
                        secondary_extra_keys, secondary_mm_idx = (
                            generate_block_hash_extra_keys(
                                request,
                                hash_start,
                                hash_end,
                                secondary_mm_idx,
                            )
                        )
                        meta_by_hash[block_hash] = BlockStoreMeta(
                            token_ids=list(request.all_token_ids[hash_start:hash_end]),
                            parent_block_hash=(
                                None
                                if hash_idx == 0
                                else maybe_convert_block_hash(
                                    request.block_hashes[hash_idx - 1]
                                )
                            ),
                            lora_id=lora_req.adapter_id if lora_req else None,
                            lora_name=lora_req.name if lora_req else None,
                            extra_keys=secondary_extra_keys,
                            block_size=self.hash_block_size,
                        )
                    block_meta.append(meta_by_hash)

        return gpu_block_ids, advanced_per_group, block_meta

    def _classify_store_candidate(self, gpu_block: "KVCacheBlock") -> _StoreAdmission:
        if gpu_block.is_null:
            return _StoreAdmission.NULL_BLOCK
        if gpu_block.block_hash is None:
            return _StoreAdmission.NOT_HASHED
        if gpu_block.block_id in self._in_flight_store_gpu_blocks:
            return _StoreAdmission.IN_FLIGHT
        if (
            self.cpu_block_pool.cached_block_hash_to_block.get_one_block(
                gpu_block.block_hash
            )
            is not None
        ):
            return _StoreAdmission.ALREADY_CACHED
        return _StoreAdmission.READY

    def get_boundary_store_stats(self) -> BoundaryStoreStats:
        return replace(self.boundary_store_stats)

    def get_stats(self) -> SimpleCPUOffloadStats:
        """Drain per-step stats for the connector's stats hooks."""
        stats = SimpleCPUOffloadStats()

        current = self.boundary_store_stats
        for outcome, field_name in OUTCOME_TO_FIELD.items():
            delta = getattr(current, field_name) - getattr(
                self._boundary_stats_snapshot, field_name
            )
            if delta > 0:
                stats.increase_counter(MetricName.SAVE_OUTCOMES, delta, (outcome,))
        self._boundary_stats_snapshot = replace(current)

        if self._interval_load_blocks_completed:
            stats.increase_counter(
                MetricName.LOAD_BLOCKS,
                self._interval_load_blocks_completed,
            )
            self._interval_load_blocks_completed = 0

        stats.set_gauge(
            MetricName.USED_BLOCKS,
            self.num_cpu_blocks - self.cpu_block_pool.get_num_free_blocks(),
        )
        pending = sum(
            len(t.cpu_block_ids)
            for t in (
                *self._store_event_to_blocks.values(),
                *self._pending_finished_stores,
                *self._abandoned_store_event_to_blocks.values(),
            )
        )
        stats.set_gauge(MetricName.PENDING_STORE_BLOCKS, pending)
        stats.set_gauge(MetricName.INFO, 1, self._info_labelvalues)
        return stats

    def update_connector_output(self, connector_output: KVConnectorOutput) -> None:
        """Handle async transfer completions from worker.

        Load completions arrive via finished_recving (real req_ids).
        Store completions arrive via kv_connector_worker_meta as
        per-event worker counts. We accumulate across steps and process
        a store event only when all workers have reported completion.
        """
        # --- Load completions ---
        for req_id in list(connector_output.finished_recving or []):
            completed_blocks = self._cleanup_load_request(req_id)
            if completed_blocks:
                self._interval_load_blocks_completed += completed_blocks

        # --- Store completions ---
        meta = connector_output.kv_connector_worker_meta
        if not isinstance(meta, SimpleCPUOffloadWorkerMetadata):
            return
        for event_idx, count in meta.completed_store_events.items():
            total = self._store_event_pending_counts.get(event_idx, 0) + count
            if total >= self._expected_worker_count:
                self._store_event_pending_counts.pop(event_idx, None)
                self._process_store_event(event_idx)
            else:
                self._store_event_pending_counts[event_idx] = total

    def _process_store_event(self, event_idx: int) -> None:
        """Process a fully-completed store event."""
        transfer = self._store_event_to_blocks.pop(event_idx, None)
        if transfer is None:
            transfer = self._abandoned_store_event_to_blocks.pop(event_idx, None)
            if transfer is None:
                return  # guard stale events from before a reset() call
            self._release_transfer_refs(transfer)
            return

        if not self._lazy_mode:
            self._in_flight_store_gpu_blocks.difference_update(transfer.gpu_block_ids)

        self._process_store_completion(
            transfer.gpu_block_ids,
            transfer.cpu_block_ids,
            transfer.block_meta,
        )
        logger.debug(
            "Store event %d completed: cached %d blocks to CPU",
            event_idx,
            len(transfer.cpu_block_ids),
        )

        # Eager only: update per-req state
        if not self._lazy_mode:
            for req_id in self._store_event_to_reqs.pop(event_idx, []):
                state = self._reqs_to_store.get(req_id)
                if state is None:
                    continue
                state.store_events.discard(event_idx)
                if state.finished and not state.store_events:
                    self._cleanup_store_request(req_id)

    def _process_store_completion(
        self,
        gpu_block_ids: list[int],
        cpu_block_ids: list[int],
        block_meta: list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None = None,
    ) -> None:
        """Register copied blocks in the CPU prefix cache and release refs."""
        assert len(cpu_block_ids) == len(gpu_block_ids)
        assert block_meta is None or len(block_meta) == len(gpu_block_ids)

        cpu_blocks = [self.cpu_block_pool.blocks[bid] for bid in cpu_block_ids]

        assert self._gpu_block_pool is not None
        for i, (gpu_block_id, cpu_block) in enumerate(zip(gpu_block_ids, cpu_blocks)):
            gpu_block = self._gpu_block_pool.blocks[gpu_block_id]
            primary_hash = gpu_block.block_hash
            assert primary_hash is not None
            self.cpu_block_pool._insert_block_hash(
                primary_hash,
                cpu_block,
                num_tokens=gpu_block.block_hash_num_tokens,
            )
            secondary_hashes = self._gpu_block_pool.cached_block_hashes_by_block.get(
                gpu_block_id, ()
            )
            for block_hash in secondary_hashes:
                self.cpu_block_pool._insert_block_hash(
                    block_hash, cpu_block, num_tokens=None
                )
            if self.enable_kv_cache_events:
                meta_by_hash = block_meta[i] if block_meta is not None else {}
                primary_group_idx = get_group_id(primary_hash)
                primary_block_size = self.group_block_sizes[primary_group_idx]
                events_to_emit: list[
                    tuple[BlockHashWithGroupId, BlockStoreMeta | None, int]
                ] = [
                    (
                        primary_hash,
                        meta_by_hash.get(primary_hash),
                        primary_block_size,
                    )
                ]
                for secondary_hash in secondary_hashes:
                    secondary_meta = meta_by_hash.get(secondary_hash)
                    events_to_emit.append(
                        (
                            secondary_hash,
                            secondary_meta,
                            secondary_meta.block_size
                            if secondary_meta is not None
                            and secondary_meta.block_size is not None
                            else 0,
                        )
                    )

                for block_hash, meta, event_block_size in events_to_emit:
                    group_idx = get_group_id(block_hash)
                    spec = self.cpu_kv_cache_config.kv_cache_groups[
                        group_idx
                    ].kv_cache_spec
                    if meta is not None and len(meta.token_ids) != event_block_size:
                        # token_ids were sliced with a different g_block_size
                        # (e.g. Mamba+DCP where capture uses spec.block_size *
                        # cp_world_size but the event emits spec.block_size * 1).
                        # Emit empty metadata until #49962 fixes the sizing.
                        token_ids = []
                        parent_block_hash = None
                        extra_keys = None
                    else:
                        token_ids = meta.token_ids if meta else []
                        parent_block_hash = meta.parent_block_hash if meta else None
                        extra_keys = meta.extra_keys if meta else None
                    self.cpu_block_pool.kv_event_queue.append(
                        BlockStored(
                            block_hashes=[
                                maybe_convert_block_hash(get_block_hash(block_hash))
                            ],
                            parent_block_hash=parent_block_hash,
                            token_ids=token_ids,
                            block_size=event_block_size,
                            lora_id=meta.lora_id if meta else None,
                            medium=self.kv_event_medium,
                            lora_name=meta.lora_name if meta else None,
                            extra_keys=to_event_extra_keys(extra_keys and [extra_keys]),
                            group_idx=group_idx,
                            kv_cache_spec_kind=get_kv_cache_spec_kind(spec).value,
                            kv_cache_spec_sliding_window=(
                                get_kv_cache_spec_sliding_window(spec)
                            ),
                            locality="LOCAL",
                        )
                    )

        # Free CPU and GPU blocks' ref counts to turn them into prefix cache
        self.cpu_block_pool.free_blocks(cpu_blocks)
        self._gpu_block_pool.free_blocks(
            self._gpu_block_pool.blocks[bid] for bid in gpu_block_ids
        )

    def _release_transfer_refs(self, transfer: TransferMeta) -> None:
        """Release transfer refs without making copied data cacheable."""
        cpu_blocks = [self.cpu_block_pool.blocks[bid] for bid in transfer.cpu_block_ids]
        self.cpu_block_pool.free_blocks(cpu_blocks)
        assert self._gpu_block_pool is not None
        self._gpu_block_pool.free_blocks(
            self._gpu_block_pool.blocks[bid] for bid in transfer.gpu_block_ids
        )

    def has_pending_stores(self) -> bool:
        """Return True if a store transfer is queued or in flight."""
        return bool(
            self._pending_finished_stores
            or self._store_event_to_blocks
            or self._abandoned_store_event_to_blocks
        )

    def request_finished(
        self,
        request: "Request",
        block_ids: list[int],
    ) -> tuple[bool, dict[str, Any] | None]:
        """Always returns (False, None). GPU blocks are protected by ref_cnt,
        so the scheduler can free blocks immediately."""
        req_id = request.request_id

        # Release any temp CPU hit pin from get_num_new_matched_tokens()
        # if request is canceled or preempted before update_state_after_alloc()
        pending = self._pending_cpu_hits.pop(req_id, None)
        if pending is not None:
            self._free_pending_cpu_hit(pending)

        # Handle load: defer cleanup if load is in-flight
        load_state = self._reqs_to_load.get(req_id)
        if load_state is not None:
            if load_state.load_event is not None:
                load_state.finished = True  # Defer: load in-flight
            else:
                self._cleanup_load_request(req_id)

        # Handle store (eager mode only): defer cleanup if stores in-flight
        if not self._lazy_mode:
            store_state = self._reqs_to_store.get(req_id)
            if store_state is not None:
                if store_state.store_events:
                    store_state.finished = True  # Defer: stores in-flight
                else:
                    self._cleanup_store_request(req_id)

        return False, None

    def request_finished_all_groups(
        self,
        request: "Request",
        block_ids: tuple[list[int], ...],
    ) -> tuple[bool, dict[str, Any] | None]:
        self._queue_finished_eager_store(request, block_ids)
        return self.request_finished(request, block_ids=[])

    def _queue_finished_eager_store(
        self, request: "Request", block_ids: tuple[list[int], ...]
    ) -> None:
        """Queue confirmed eager blocks omitted after a request enters decode."""
        state = self._reqs_to_store.get(request.request_id)
        gpu_pool = self._gpu_block_pool
        if state is None or gpu_pool is None:
            return
        gpu_ids, _, block_meta = self._select_eager_blocks_to_store(state, block_ids)
        partial_tail = self._find_fa_partial_tail_source(request, block_ids)
        if (
            partial_tail is not None
            # Decode tokens can fill the boundary block, in which case the
            # positional scan above already selected it. Appending again would
            # allocate two CPU blocks for one GPU block.
            and partial_tail not in gpu_ids
            and len(gpu_ids) < self.cpu_block_pool.get_num_free_blocks()
        ):
            gpu_ids.append(partial_tail)
            if block_meta is not None:
                # The block's own hashes are registered at completion, so no
                # capture is needed here; emit the event without a payload, as
                # the size-mismatch fallback above already does.
                block_meta.append({})
        if not gpu_ids:
            return
        cpu_blocks = self.cpu_block_pool.get_new_blocks(len(gpu_ids))
        self._pending_finished_stores.append(
            TransferMeta(gpu_ids, [b.block_id for b in cpu_blocks], block_meta)
        )
        self._in_flight_store_gpu_blocks.update(gpu_ids)
        gpu_pool.touch([gpu_pool.blocks[bid] for bid in gpu_ids])

    def _find_fa_partial_tail_source(
        self, request: "Request", block_ids: tuple[list[int], ...]
    ) -> int | None:
        """Locate the full-attention prompt-tail block for a fine-grained hit.

        Mirrors ``FullAttentionManager._cache_partial_tail_block``: only the
        final prompt hash boundary is eligible, and boundaries that land on a
        physical block edge are already covered by the positional scan.

        Attention block tables are append-only, so the block is located
        positionally, as the mooncake store and offloading connector also do;
        no new core hand-off is required. The boundary key is registered at
        completion along with the block's other hashes, so nothing has to be
        captured here.
        """
        if not self.cpu_coordinator.enable_partial_hash_hits:
            return None
        assert self._gpu_block_pool is not None
        boundary_tokens = (
            request.num_prompt_tokens // self.hash_block_size * self.hash_block_size
        )
        if boundary_tokens == 0 or boundary_tokens > request.num_computed_tokens:
            return None
        if boundary_tokens % self.fa_block_size == 0:
            return None
        block_idx = boundary_tokens // self.fa_block_size
        fa_block_ids = block_ids[self.fa_gidx]
        if block_idx >= len(fa_block_ids):
            return None
        gpu_block_id = fa_block_ids[block_idx]
        gpu_block = self._gpu_block_pool.blocks[gpu_block_id]
        if gpu_block.is_null or gpu_block_id in self._in_flight_store_gpu_blocks:
            return None
        hash_idx = boundary_tokens // self.hash_block_size - 1
        if hash_idx >= len(request.block_hashes):
            return None
        block_hash = make_block_hash_with_group_id(
            request.block_hashes[hash_idx], self.fa_gidx
        )
        # Only worth copying when the boundary key is actually registered on
        # the source block; completion replays the block's own hashes.
        if (
            gpu_block.block_hash != block_hash
            and block_hash
            not in self._gpu_block_pool.cached_block_hashes_by_block.get(
                gpu_block_id, ()
            )
        ):
            return None
        if self.cpu_block_pool.cached_block_hash_to_block.get_one_block(block_hash):
            return None
        return gpu_block_id

    def _free_pending_cpu_hit(
        self, pending: tuple[tuple[list["KVCacheBlock"], ...], int, int]
    ) -> None:
        """Release the temporary CPU block pin taken in get_num_new_matched_tokens()."""
        cpu_hit_blocks, _hit_length, _num_computed_tokens = pending
        blocks_to_free = [
            blk for grp in cpu_hit_blocks for blk in grp if not blk.is_null
        ]
        if blocks_to_free:
            self.cpu_block_pool.free_blocks(blocks_to_free)

    def _cleanup_load_request(self, req_id: str) -> int:
        """Release all load resources for a request.

        Shared between request_finished() and update_connector_output() paths.
        Removes the request from _reqs_to_load, cleans up event mappings,
        and frees CPU/GPU touch refs.

        Returns the number of blocks in the load if it had been issued to
        the worker, else 0 (never-issued loads are not counted as completed
        by the caller).
        """
        state = self._reqs_to_load.pop(req_id, None)
        if state is None:
            state = self._abandoned_reqs_to_load.pop(req_id, None)
        if state is None:
            return 0
        # Remove from load event mapping (only this req, not whole event)
        if state.load_event is not None:
            reqs = self._load_event_to_reqs.get(state.load_event)
            if reqs is not None:
                with contextlib.suppress(ValueError):
                    reqs.remove(req_id)
                if not reqs:
                    self._load_event_to_reqs.pop(state.load_event, None)

        if state.transfer_meta is not None:
            # Free CPU touch refs
            self.cpu_block_pool.free_blocks(
                self.cpu_block_pool.blocks[bid]
                for bid in state.transfer_meta.cpu_block_ids
            )
            # Free GPU touch refs
            assert self._gpu_block_pool is not None
            self._gpu_block_pool.free_blocks(
                self._gpu_block_pool.blocks[bid]
                for bid in state.transfer_meta.gpu_block_ids
            )

        if state.load_event is None:
            return 0
        return len(state.transfer_meta.gpu_block_ids)

    def _cleanup_store_request(self, req_id: str) -> None:
        """Release store metadata for a request.

        Metadata-only cleanup but no block freeing. Job completion handles
        block caching and GPU ref freeing via _process_store_completion().
        """
        state = self._reqs_to_store.pop(req_id, None)
        if state is None:
            return
        for event_idx in list(state.store_events):
            if (reqs := self._store_event_to_reqs.get(event_idx)) is not None:
                with contextlib.suppress(ValueError):
                    reqs.remove(req_id)
                if not reqs:
                    self._store_event_to_reqs.pop(event_idx, None)
        state.store_events.clear()

    def take_events(self) -> Iterable[KVCacheEvent]:
        events = self.cpu_block_pool.take_events()
        for event in events:
            if isinstance(event, BlockRemoved):
                event.medium = self.kv_event_medium
                event.locality = "LOCAL"
        return events

    def reset(self) -> bool:
        """Abandon pending transfers and reset the CPU cache when safe.

        Worker-side DMA may still be using blocks after reset is requested.
        Keep those block refs pinned until the existing completion path reports
        the transfer finished, then release refs without caching abandoned
        store results.
        """
        self._abandoned_store_event_to_blocks.update(self._store_event_to_blocks)
        for transfer in self._pending_finished_stores:
            self._release_transfer_refs(transfer)
        self._pending_finished_stores.clear()

        self._store_event_to_blocks.clear()
        self._in_flight_store_gpu_blocks.clear()

        # Loads that have not been sent to the worker cannot have running DMA.
        # In-flight loads stay pinned and are cleaned up on completion.
        for req_id in list(self._reqs_to_load):
            state = self._reqs_to_load.pop(req_id)
            if state.load_event is None:
                self._reqs_to_load[req_id] = state
                self._cleanup_load_request(req_id)
            else:
                self._abandoned_reqs_to_load[req_id] = state

        self._reqs_to_store.clear()
        self._store_event_to_reqs.clear()
        self._store_event_pending_counts = {
            event_idx: count
            for event_idx, count in self._store_event_pending_counts.items()
            if event_idx in self._abandoned_store_event_to_blocks
        }
        self._cursor = None
        # Seed the fresh counters with outcomes since the last drain so a
        # mid-interval reset() does not drop them from the next report.
        undrained = BoundaryStoreStats(
            **{
                f.name: max(
                    0,
                    getattr(self.boundary_store_stats, f.name)
                    - getattr(self._boundary_stats_snapshot, f.name),
                )
                for f in fields(BoundaryStoreStats)
            }
        )
        self.boundary_store_stats = undrained
        self._boundary_stats_snapshot = BoundaryStoreStats()
        # NOTE: _load_event_counter / _store_event_counter are not
        # reset as they are monotonic and must stay ahead of the workers
        # high-water marks to avoid event index collisions

        if self._abandoned_store_event_to_blocks or self._abandoned_reqs_to_load:
            return False

        return self.cpu_block_pool.reset_prefix_cache()

_cached_gpu_block(resolved_hashes, block_idx, group_id)

Return the GPU block still cached under this group's block hash.

Source code in vllm/v1/simple_kv_offload/manager.py
def _cached_gpu_block(
    self, resolved_hashes: "BlockHashList", block_idx: int, group_id: int
) -> "KVCacheBlock | None":
    """Return the GPU block still cached under this group's block hash."""
    if block_idx >= len(resolved_hashes):
        return None
    assert self._gpu_block_pool is not None
    blocks = self._gpu_block_pool.get_cached_block(
        resolved_hashes[block_idx], [group_id]
    )
    return blocks[0] if blocks else None

_cleanup_load_request(req_id)

Release all load resources for a request.

Shared between request_finished() and update_connector_output() paths. Removes the request from _reqs_to_load, cleans up event mappings, and frees CPU/GPU touch refs.

Returns the number of blocks in the load if it had been issued to the worker, else 0 (never-issued loads are not counted as completed by the caller).

Source code in vllm/v1/simple_kv_offload/manager.py
def _cleanup_load_request(self, req_id: str) -> int:
    """Release all load resources for a request.

    Shared between request_finished() and update_connector_output() paths.
    Removes the request from _reqs_to_load, cleans up event mappings,
    and frees CPU/GPU touch refs.

    Returns the number of blocks in the load if it had been issued to
    the worker, else 0 (never-issued loads are not counted as completed
    by the caller).
    """
    state = self._reqs_to_load.pop(req_id, None)
    if state is None:
        state = self._abandoned_reqs_to_load.pop(req_id, None)
    if state is None:
        return 0
    # Remove from load event mapping (only this req, not whole event)
    if state.load_event is not None:
        reqs = self._load_event_to_reqs.get(state.load_event)
        if reqs is not None:
            with contextlib.suppress(ValueError):
                reqs.remove(req_id)
            if not reqs:
                self._load_event_to_reqs.pop(state.load_event, None)

    if state.transfer_meta is not None:
        # Free CPU touch refs
        self.cpu_block_pool.free_blocks(
            self.cpu_block_pool.blocks[bid]
            for bid in state.transfer_meta.cpu_block_ids
        )
        # Free GPU touch refs
        assert self._gpu_block_pool is not None
        self._gpu_block_pool.free_blocks(
            self._gpu_block_pool.blocks[bid]
            for bid in state.transfer_meta.gpu_block_ids
        )

    if state.load_event is None:
        return 0
    return len(state.transfer_meta.gpu_block_ids)

_cleanup_store_request(req_id)

Release store metadata for a request.

Metadata-only cleanup but no block freeing. Job completion handles block caching and GPU ref freeing via _process_store_completion().

Source code in vllm/v1/simple_kv_offload/manager.py
def _cleanup_store_request(self, req_id: str) -> None:
    """Release store metadata for a request.

    Metadata-only cleanup but no block freeing. Job completion handles
    block caching and GPU ref freeing via _process_store_completion().
    """
    state = self._reqs_to_store.pop(req_id, None)
    if state is None:
        return
    for event_idx in list(state.store_events):
        if (reqs := self._store_event_to_reqs.get(event_idx)) is not None:
            with contextlib.suppress(ValueError):
                reqs.remove(req_id)
            if not reqs:
                self._store_event_to_reqs.pop(event_idx, None)
    state.store_events.clear()

_derive_cpu_config(gpu_config, cpu_capacity_bytes) staticmethod

Derive a CPU KVCacheConfig from the GPU config. Same kv_cache_groups, num_blocks scaled by CPU/GPU memory ratio.

Source code in vllm/v1/simple_kv_offload/manager.py
@staticmethod
def _derive_cpu_config(
    gpu_config: "KVCacheConfig", cpu_capacity_bytes: int
) -> "KVCacheConfig":
    """Derive a CPU KVCacheConfig from the GPU config.
    Same kv_cache_groups, num_blocks scaled by CPU/GPU memory ratio."""
    # Import here to avoid potential circular imports
    from vllm.v1.kv_cache_interface import KVCacheTensor

    assert len(gpu_config.kv_cache_tensors) > 0

    # Every KVCacheTensor describes placement within the same backing allocation,
    # so its size is the total GPU KV cache size.
    gpu_total_bytes = gpu_config.kv_cache_tensors[0].size
    num_gpu_blocks = gpu_config.num_blocks
    num_cpu_blocks = max(1, num_gpu_blocks * cpu_capacity_bytes // gpu_total_bytes)
    # Create CPU kv_cache_tensors mirroring GPU by scaling size proportionally.
    cpu_tensors = [
        KVCacheTensor(
            size=t.size // num_gpu_blocks * num_cpu_blocks,
            layers=list(t.layers),
            layer_stride=t.layer_stride,
            block_stride=t.block_stride,
            offset=t.offset,
        )
        for t in gpu_config.kv_cache_tensors
    ]

    return replace(
        gpu_config,
        num_blocks=num_cpu_blocks,
        kv_cache_tensors=cpu_tensors,
    )

_estimate_lazy_target_blocks(kv_cache_config, max_num_batched_tokens, cp_world_size=1) staticmethod

GPU blocks to keep available (free/offloaded) per step in lazy mode.

Source code in vllm/v1/simple_kv_offload/manager.py
@staticmethod
def _estimate_lazy_target_blocks(
    kv_cache_config: "KVCacheConfig",
    max_num_batched_tokens: int,
    cp_world_size: int = 1,
) -> int:
    """GPU blocks to keep available (free/offloaded) per step in lazy mode."""
    WATERMARK_RATIO = 1.0  # Reserve larger space to avoid running out of GPU blocks
    target = 0
    for g in kv_cache_config.prefix_cacheable_groups:
        spec = g.kv_cache_spec
        block_size = resolve_dcp_kv_block_size(spec, cp_world_size)
        if isinstance(spec, MambaSpec):
            target += 2
        elif isinstance(spec, SlidingWindowSpec):
            target += cdiv(spec.sliding_window, block_size) + 1
        else:
            target += cdiv(max_num_batched_tokens, block_size)
    return int(target * (1 + WATERMARK_RATIO))

_find_fa_partial_tail_source(request, block_ids)

Locate the full-attention prompt-tail block for a fine-grained hit.

Mirrors FullAttentionManager._cache_partial_tail_block: only the final prompt hash boundary is eligible, and boundaries that land on a physical block edge are already covered by the positional scan.

Attention block tables are append-only, so the block is located positionally, as the mooncake store and offloading connector also do; no new core hand-off is required. The boundary key is registered at completion along with the block's other hashes, so nothing has to be captured here.

Source code in vllm/v1/simple_kv_offload/manager.py
def _find_fa_partial_tail_source(
    self, request: "Request", block_ids: tuple[list[int], ...]
) -> int | None:
    """Locate the full-attention prompt-tail block for a fine-grained hit.

    Mirrors ``FullAttentionManager._cache_partial_tail_block``: only the
    final prompt hash boundary is eligible, and boundaries that land on a
    physical block edge are already covered by the positional scan.

    Attention block tables are append-only, so the block is located
    positionally, as the mooncake store and offloading connector also do;
    no new core hand-off is required. The boundary key is registered at
    completion along with the block's other hashes, so nothing has to be
    captured here.
    """
    if not self.cpu_coordinator.enable_partial_hash_hits:
        return None
    assert self._gpu_block_pool is not None
    boundary_tokens = (
        request.num_prompt_tokens // self.hash_block_size * self.hash_block_size
    )
    if boundary_tokens == 0 or boundary_tokens > request.num_computed_tokens:
        return None
    if boundary_tokens % self.fa_block_size == 0:
        return None
    block_idx = boundary_tokens // self.fa_block_size
    fa_block_ids = block_ids[self.fa_gidx]
    if block_idx >= len(fa_block_ids):
        return None
    gpu_block_id = fa_block_ids[block_idx]
    gpu_block = self._gpu_block_pool.blocks[gpu_block_id]
    if gpu_block.is_null or gpu_block_id in self._in_flight_store_gpu_blocks:
        return None
    hash_idx = boundary_tokens // self.hash_block_size - 1
    if hash_idx >= len(request.block_hashes):
        return None
    block_hash = make_block_hash_with_group_id(
        request.block_hashes[hash_idx], self.fa_gidx
    )
    # Only worth copying when the boundary key is actually registered on
    # the source block; completion replays the block's own hashes.
    if (
        gpu_block.block_hash != block_hash
        and block_hash
        not in self._gpu_block_pool.cached_block_hashes_by_block.get(
            gpu_block_id, ()
        )
    ):
        return None
    if self.cpu_block_pool.cached_block_hash_to_block.get_one_block(block_hash):
        return None
    return gpu_block_id

_free_pending_cpu_hit(pending)

Release the temporary CPU block pin taken in get_num_new_matched_tokens().

Source code in vllm/v1/simple_kv_offload/manager.py
def _free_pending_cpu_hit(
    self, pending: tuple[tuple[list["KVCacheBlock"], ...], int, int]
) -> None:
    """Release the temporary CPU block pin taken in get_num_new_matched_tokens()."""
    cpu_hit_blocks, _hit_length, _num_computed_tokens = pending
    blocks_to_free = [
        blk for grp in cpu_hit_blocks for blk in grp if not blk.is_null
    ]
    if blocks_to_free:
        self.cpu_block_pool.free_blocks(blocks_to_free)

_prepare_eager_store_specs(scheduler_output)

Identify newly computed blocks to offload from scheduler requests.

Only considers blocks whose KV data has been confirmed computed by the GPU. Blocks from the current step are stored on a later step, or by the finish-time flush if the request completes first.

Returns:

  • list[int] –

    (gpu_block_ids, cpu_block_ids, req_ids, block_meta) for the store

  • list[int] –

    event. block_meta is None when kv cache events are disabled.

Source code in vllm/v1/simple_kv_offload/manager.py
def _prepare_eager_store_specs(
    self, scheduler_output: SchedulerOutput
) -> tuple[
    list[int],
    list[int],
    list[str],
    list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None,
]:
    """Identify newly computed blocks to offload from scheduler requests.

    Only considers blocks whose KV data has been **confirmed computed** by
    the GPU. Blocks from the current step are stored on a later step, or by
    the finish-time flush if the request completes first.

    Returns:
        (gpu_block_ids, cpu_block_ids, req_ids, block_meta) for the store
        event. ``block_meta`` is None when kv cache events are disabled.

    """
    merged_gpu_block_ids: list[int] = []
    merged_cpu_block_ids: list[int] = []
    req_ids: list[str] = []
    merged_block_meta: list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None = (
        [] if self.enable_kv_cache_events else None
    )

    gpu_block_pool = self._gpu_block_pool
    if gpu_block_pool is None:
        return [], [], [], merged_block_meta
    cpu_block_pool = self.cpu_block_pool
    num_groups = len(self.cpu_kv_cache_config.kv_cache_groups)
    # Dedup against blocks already scheduled.
    in_flight = self._in_flight_store_gpu_blocks
    num_free = cpu_block_pool.get_num_free_blocks()

    block_state = scheduler_output.kv_connector_block_state
    boundary_offloads = (
        block_state.boundary_state_offloads if block_state is not None else {}
    )
    stats = self.boundary_store_stats
    preempted_req_ids = scheduler_output.preempted_req_ids or ()
    for req_id, entries in boundary_offloads.items():
        # The scheduler drains handoffs without filtering by request
        # liveness. A request that finished or was preempted in this step
        # is giving up its blocks, which may already have been reallocated
        # to another request; reading them would publish unrelated KV under
        # a valid hash. Drop conservatively, as the mooncake store does.
        store_state = self._reqs_to_store.get(req_id)
        if (
            store_state is None
            or store_state.finished
            or req_id in preempted_req_ids
            or req_id in scheduler_output.finished_req_ids
        ):
            stats.published += len(entries)
            stats.dropped_request_gone += len(entries)
            continue
        scheduled_for_req = False
        # ``boundary_tokens`` is not needed here: the cache key comes from
        # the handed-off block itself, not from the offered boundary.
        for _, gpu_block_id, _ in entries:
            stats.published += 1
            gpu_block = gpu_block_pool.blocks[gpu_block_id]
            admission = self._classify_store_candidate(gpu_block)
            if admission is _StoreAdmission.NULL_BLOCK:
                stats.dropped_null_block += 1
                continue
            if admission is _StoreAdmission.NOT_HASHED:
                stats.dropped_not_hashed += 1
                continue
            if admission is _StoreAdmission.IN_FLIGHT:
                stats.skipped_in_flight += 1
                continue
            if admission is _StoreAdmission.ALREADY_CACHED:
                stats.skipped_already_cached += 1
                continue
            if num_free <= 0:
                stats.dropped_cpu_full += 1
                logger.warning_once(
                    "SimpleCPU dropped a boundary-state handoff because "
                    "the CPU block pool is full; this boundary will not be retried."
                )
                continue
            cpu_block = cpu_block_pool.get_new_blocks(1)[0]
            merged_gpu_block_ids.append(gpu_block_id)
            merged_cpu_block_ids.append(cpu_block.block_id)
            if merged_block_meta is not None:
                # Keep the metadata list index-parallel with the block ids.
                # The hand-off block's own hashes are resolved at completion,
                # so there is nothing to capture here; emit the event without
                # a payload, as the size-mismatch fallback already does.
                merged_block_meta.append({})
            in_flight.add(gpu_block_id)
            gpu_block_pool.touch([gpu_block])
            num_free -= 1
            stats.stored += 1
            scheduled_for_req = True
        if scheduled_for_req:
            req_ids.append(req_id)

    for req_id, new_block_id_groups, preempted in yield_req_data(scheduler_output):
        state = self._reqs_to_store.get(req_id)
        if state is None or state.finished:
            continue

        if preempted:
            state.block_ids = tuple([] for _ in range(num_groups))
            state.num_stored_blocks = [0] * num_groups
        if new_block_id_groups:
            for g in range(min(num_groups, len(new_block_id_groups))):
                if new_block_id_groups[g] is not None:
                    # Accumulate new block IDs.
                    state.block_ids[g].extend(new_block_id_groups[g])

        num_new_tokens = scheduler_output.num_scheduled_tokens.get(req_id, 0)
        if num_new_tokens == 0:
            continue

        block_ids_by_group = state.block_ids
        if not block_ids_by_group:
            continue

        gpu_block_ids, advanced_per_group, block_meta = (
            self._select_eager_blocks_to_store(state, block_ids_by_group)
        )

        # Batch allocate the CPU destinations.
        n_to_alloc = len(gpu_block_ids)
        if n_to_alloc > 0:
            cpu_blocks_alloc = cpu_block_pool.get_new_blocks(n_to_alloc)
            cpu_block_ids = [blk.block_id for blk in cpu_blocks_alloc]
        else:
            cpu_block_ids = []

        if cpu_block_ids:
            req_ids.append(req_id)
            merged_gpu_block_ids.extend(gpu_block_ids)
            merged_cpu_block_ids.extend(cpu_block_ids)
            in_flight.update(gpu_block_ids)
            if merged_block_meta is not None:
                assert block_meta is not None
                merged_block_meta.extend(block_meta)

            # Touch GPU blocks to prevent freeing during async copy
            gpu_block_pool.touch(
                [gpu_block_pool.blocks[bid] for bid in gpu_block_ids]
            )

            logger.debug(
                "Request %s: Scheduling store of %d blocks to CPU (%d groups)",
                req_id,
                len(cpu_block_ids),
                num_groups,
            )

        # Advance per-group cursors (includes cached hits + newly stored)
        for g in range(num_groups):
            state.num_stored_blocks[g] += advanced_per_group[g]

    # A request can contribute both boundary and positional blocks to the
    # same event. Completion tracking is set-based, so keep one request ID.
    req_ids = list(dict.fromkeys(req_ids))
    return merged_gpu_block_ids, merged_cpu_block_ids, req_ids, merged_block_meta

_prepare_lazy_store_specs()

Single-pass cursor walk: offload cached GPU blocks near eviction.

Walks the GPU free queue from the cursor, counting blocks that are free-or-offloaded (safe for the allocator to evict). Stops when target_free blocks are covered or CPU capacity is reached.

Source code in vllm/v1/simple_kv_offload/manager.py
def _prepare_lazy_store_specs(
    self,
) -> tuple[
    list[int],
    list[int],
    list[str],
    list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None,
]:
    """Single-pass cursor walk: offload cached GPU blocks near eviction.

    Walks the GPU free queue from the cursor, counting blocks that are
    free-or-offloaded (safe for the allocator to evict). Stops when
    target_free blocks are covered or CPU capacity is reached.
    """
    gpu_pool = self._gpu_block_pool
    if gpu_pool is None or self._target_free <= 0:
        return [], [], [], None

    free_queue = gpu_pool.free_block_queue
    cpu_pool = self.cpu_block_pool
    num_cpu_free = cpu_pool.get_num_free_blocks()

    # Validate cursor: stale if block was removed from free queue.
    if self._cursor is not None and self._cursor.ref_cnt > 0:
        self._cursor = None

    gpu_ids: list[int] = []
    last_visited = self._cursor

    for covered, node in enumerate(free_queue.iter_blocks_after(self._cursor)):
        if covered >= self._target_free or len(gpu_ids) >= num_cpu_free:
            break

        last_visited = node
        bhash = node.block_hash

        if (
            bhash is not None
            and not node.is_null
            and cpu_pool.cached_block_hash_to_block.get_one_block(bhash) is None
        ):
            gpu_ids.append(node.block_id)

    self._cursor = last_visited

    # Batch-allocate CPU blocks.
    if gpu_ids:
        cpu_blocks = cpu_pool.get_new_blocks(len(gpu_ids))
        cpu_ids = [blk.block_id for blk in cpu_blocks]
        # Touch GPU blocks to prevent eviction during async copy.
        gpu_pool.touch([gpu_pool.blocks[bid] for bid in gpu_ids])
    else:
        cpu_ids = []

    return gpu_ids, cpu_ids, [], None

_process_store_completion(gpu_block_ids, cpu_block_ids, block_meta=None)

Register copied blocks in the CPU prefix cache and release refs.

Source code in vllm/v1/simple_kv_offload/manager.py
def _process_store_completion(
    self,
    gpu_block_ids: list[int],
    cpu_block_ids: list[int],
    block_meta: list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None = None,
) -> None:
    """Register copied blocks in the CPU prefix cache and release refs."""
    assert len(cpu_block_ids) == len(gpu_block_ids)
    assert block_meta is None or len(block_meta) == len(gpu_block_ids)

    cpu_blocks = [self.cpu_block_pool.blocks[bid] for bid in cpu_block_ids]

    assert self._gpu_block_pool is not None
    for i, (gpu_block_id, cpu_block) in enumerate(zip(gpu_block_ids, cpu_blocks)):
        gpu_block = self._gpu_block_pool.blocks[gpu_block_id]
        primary_hash = gpu_block.block_hash
        assert primary_hash is not None
        self.cpu_block_pool._insert_block_hash(
            primary_hash,
            cpu_block,
            num_tokens=gpu_block.block_hash_num_tokens,
        )
        secondary_hashes = self._gpu_block_pool.cached_block_hashes_by_block.get(
            gpu_block_id, ()
        )
        for block_hash in secondary_hashes:
            self.cpu_block_pool._insert_block_hash(
                block_hash, cpu_block, num_tokens=None
            )
        if self.enable_kv_cache_events:
            meta_by_hash = block_meta[i] if block_meta is not None else {}
            primary_group_idx = get_group_id(primary_hash)
            primary_block_size = self.group_block_sizes[primary_group_idx]
            events_to_emit: list[
                tuple[BlockHashWithGroupId, BlockStoreMeta | None, int]
            ] = [
                (
                    primary_hash,
                    meta_by_hash.get(primary_hash),
                    primary_block_size,
                )
            ]
            for secondary_hash in secondary_hashes:
                secondary_meta = meta_by_hash.get(secondary_hash)
                events_to_emit.append(
                    (
                        secondary_hash,
                        secondary_meta,
                        secondary_meta.block_size
                        if secondary_meta is not None
                        and secondary_meta.block_size is not None
                        else 0,
                    )
                )

            for block_hash, meta, event_block_size in events_to_emit:
                group_idx = get_group_id(block_hash)
                spec = self.cpu_kv_cache_config.kv_cache_groups[
                    group_idx
                ].kv_cache_spec
                if meta is not None and len(meta.token_ids) != event_block_size:
                    # token_ids were sliced with a different g_block_size
                    # (e.g. Mamba+DCP where capture uses spec.block_size *
                    # cp_world_size but the event emits spec.block_size * 1).
                    # Emit empty metadata until #49962 fixes the sizing.
                    token_ids = []
                    parent_block_hash = None
                    extra_keys = None
                else:
                    token_ids = meta.token_ids if meta else []
                    parent_block_hash = meta.parent_block_hash if meta else None
                    extra_keys = meta.extra_keys if meta else None
                self.cpu_block_pool.kv_event_queue.append(
                    BlockStored(
                        block_hashes=[
                            maybe_convert_block_hash(get_block_hash(block_hash))
                        ],
                        parent_block_hash=parent_block_hash,
                        token_ids=token_ids,
                        block_size=event_block_size,
                        lora_id=meta.lora_id if meta else None,
                        medium=self.kv_event_medium,
                        lora_name=meta.lora_name if meta else None,
                        extra_keys=to_event_extra_keys(extra_keys and [extra_keys]),
                        group_idx=group_idx,
                        kv_cache_spec_kind=get_kv_cache_spec_kind(spec).value,
                        kv_cache_spec_sliding_window=(
                            get_kv_cache_spec_sliding_window(spec)
                        ),
                        locality="LOCAL",
                    )
                )

    # Free CPU and GPU blocks' ref counts to turn them into prefix cache
    self.cpu_block_pool.free_blocks(cpu_blocks)
    self._gpu_block_pool.free_blocks(
        self._gpu_block_pool.blocks[bid] for bid in gpu_block_ids
    )

_process_store_event(event_idx)

Process a fully-completed store event.

Source code in vllm/v1/simple_kv_offload/manager.py
def _process_store_event(self, event_idx: int) -> None:
    """Process a fully-completed store event."""
    transfer = self._store_event_to_blocks.pop(event_idx, None)
    if transfer is None:
        transfer = self._abandoned_store_event_to_blocks.pop(event_idx, None)
        if transfer is None:
            return  # guard stale events from before a reset() call
        self._release_transfer_refs(transfer)
        return

    if not self._lazy_mode:
        self._in_flight_store_gpu_blocks.difference_update(transfer.gpu_block_ids)

    self._process_store_completion(
        transfer.gpu_block_ids,
        transfer.cpu_block_ids,
        transfer.block_meta,
    )
    logger.debug(
        "Store event %d completed: cached %d blocks to CPU",
        event_idx,
        len(transfer.cpu_block_ids),
    )

    # Eager only: update per-req state
    if not self._lazy_mode:
        for req_id in self._store_event_to_reqs.pop(event_idx, []):
            state = self._reqs_to_store.get(req_id)
            if state is None:
                continue
            state.store_events.discard(event_idx)
            if state.finished and not state.store_events:
                self._cleanup_store_request(req_id)

_queue_finished_eager_store(request, block_ids)

Queue confirmed eager blocks omitted after a request enters decode.

Source code in vllm/v1/simple_kv_offload/manager.py
def _queue_finished_eager_store(
    self, request: "Request", block_ids: tuple[list[int], ...]
) -> None:
    """Queue confirmed eager blocks omitted after a request enters decode."""
    state = self._reqs_to_store.get(request.request_id)
    gpu_pool = self._gpu_block_pool
    if state is None or gpu_pool is None:
        return
    gpu_ids, _, block_meta = self._select_eager_blocks_to_store(state, block_ids)
    partial_tail = self._find_fa_partial_tail_source(request, block_ids)
    if (
        partial_tail is not None
        # Decode tokens can fill the boundary block, in which case the
        # positional scan above already selected it. Appending again would
        # allocate two CPU blocks for one GPU block.
        and partial_tail not in gpu_ids
        and len(gpu_ids) < self.cpu_block_pool.get_num_free_blocks()
    ):
        gpu_ids.append(partial_tail)
        if block_meta is not None:
            # The block's own hashes are registered at completion, so no
            # capture is needed here; emit the event without a payload, as
            # the size-mismatch fallback above already does.
            block_meta.append({})
    if not gpu_ids:
        return
    cpu_blocks = self.cpu_block_pool.get_new_blocks(len(gpu_ids))
    self._pending_finished_stores.append(
        TransferMeta(gpu_ids, [b.block_id for b in cpu_blocks], block_meta)
    )
    self._in_flight_store_gpu_blocks.update(gpu_ids)
    gpu_pool.touch([gpu_pool.blocks[bid] for bid in gpu_ids])

_release_transfer_refs(transfer)

Release transfer refs without making copied data cacheable.

Source code in vllm/v1/simple_kv_offload/manager.py
def _release_transfer_refs(self, transfer: TransferMeta) -> None:
    """Release transfer refs without making copied data cacheable."""
    cpu_blocks = [self.cpu_block_pool.blocks[bid] for bid in transfer.cpu_block_ids]
    self.cpu_block_pool.free_blocks(cpu_blocks)
    assert self._gpu_block_pool is not None
    self._gpu_block_pool.free_blocks(
        self._gpu_block_pool.blocks[bid] for bid in transfer.gpu_block_ids
    )

_select_eager_blocks_to_store(state, block_ids_by_group)

Return confirmed eager blocks, cursor advances, and event metadata.

Source code in vllm/v1/simple_kv_offload/manager.py
def _select_eager_blocks_to_store(
    self,
    state: StoreRequestState,
    block_ids_by_group: tuple[list[int], ...],
) -> tuple[
    list[int],
    list[int],
    list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None,
]:
    """Return confirmed eager blocks, cursor advances, and event metadata."""
    assert self._gpu_block_pool is not None
    gpu_block_ids: list[int] = []
    advanced_per_group = [0] * len(self.cpu_kv_cache_config.kv_cache_groups)
    block_meta: list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None = (
        [] if self.enable_kv_cache_events else None
    )
    request = state.request
    confirmed_tokens = request.num_computed_tokens - request.num_output_placeholders
    # Truncate to the granularity a lookup can actually land on. With
    # fine-grained hits that is hash_block_size; otherwise hits only land on
    # the scheduler block (the group LCM). Using the LCM unconditionally
    # would drop whole blocks that sit between the last LCM boundary and a
    # reachable fine-grained boundary, so the other groups would hold that
    # boundary and this one would not, and the joint hybrid lookup would
    # reconcile to zero.
    store_alignment = (
        self.hash_block_size
        if self.cpu_coordinator.enable_partial_hash_hits
        else self.block_size
    )
    aligned_tokens = confirmed_tokens // store_alignment * store_alignment
    num_free = self.cpu_block_pool.get_num_free_blocks()

    for g, group_gpu_ids in enumerate(block_ids_by_group):
        if g not in self.prefix_cacheable_group_ids:
            continue
        if len(gpu_block_ids) >= num_free:
            break
        group_manager = self.cpu_coordinator.single_type_managers[g]
        if not group_manager.has_positionally_stable_blocks:
            continue
        # FIXME (yifan): handle CPU cache eviction, where
        # num_stored_blocks can be stale and omit evicted blocks in
        # the middle of the request.
        group_size = self.group_block_sizes[g]
        ready = min(len(group_gpu_ids), aligned_tokens // group_size)
        resolved_hashes = resolve_block_hashes(
            request.block_hashes, self.hash_block_size, group_size
        )
        curr_mm_idx = 0
        secondary_mm_idx = 0
        start = state.num_stored_blocks[g]
        for i, gpu_block_id in enumerate(group_gpu_ids[start:ready], start=start):
            gpu_block = self._gpu_block_pool.blocks[gpu_block_id]
            if gpu_block.is_null:
                # Sliding-window groups null pages that left the window
                # before the connector sees the block table, but the
                # retained prefix-cache tail stays hashed in the GPU free
                # queue; store it from there.
                recovered = self._cached_gpu_block(resolved_hashes, i, g)
                if recovered is None:
                    advanced_per_group[g] += 1
                    continue
                gpu_block = recovered
                gpu_block_id = gpu_block.block_id
            if (
                self._classify_store_candidate(gpu_block)
                is not _StoreAdmission.READY
            ):
                advanced_per_group[g] += 1
                continue
            if len(gpu_block_ids) >= num_free:
                break
            primary_block_hash = gpu_block.block_hash
            assert primary_block_hash is not None
            gpu_block_ids.append(gpu_block_id)
            advanced_per_group[g] += 1
            if block_meta is not None:
                token_start = i * group_size
                token_end = token_start + group_size
                parent_hash = (
                    None
                    if i == 0
                    else maybe_convert_block_hash(resolved_hashes[i - 1])
                )
                lora_req = request.lora_request
                extra_keys, curr_mm_idx = generate_block_hash_extra_keys(
                    request, token_start, token_end, curr_mm_idx
                )
                meta_by_hash = {
                    primary_block_hash: BlockStoreMeta(
                        token_ids=list(
                            request.all_token_ids[token_start:token_end]
                        ),
                        parent_block_hash=parent_hash,
                        lora_id=lora_req.adapter_id if lora_req else None,
                        lora_name=lora_req.name if lora_req else None,
                        extra_keys=extra_keys,
                    )
                }
                first_hash_idx = token_start // self.hash_block_size
                last_hash_idx = token_end // self.hash_block_size
                for hash_idx in range(first_hash_idx, last_hash_idx):
                    block_hash = make_block_hash_with_group_id(
                        request.block_hashes[hash_idx], g
                    )
                    if block_hash is None or block_hash == primary_block_hash:
                        continue
                    hash_start = hash_idx * self.hash_block_size
                    hash_end = hash_start + self.hash_block_size
                    secondary_extra_keys, secondary_mm_idx = (
                        generate_block_hash_extra_keys(
                            request,
                            hash_start,
                            hash_end,
                            secondary_mm_idx,
                        )
                    )
                    meta_by_hash[block_hash] = BlockStoreMeta(
                        token_ids=list(request.all_token_ids[hash_start:hash_end]),
                        parent_block_hash=(
                            None
                            if hash_idx == 0
                            else maybe_convert_block_hash(
                                request.block_hashes[hash_idx - 1]
                            )
                        ),
                        lora_id=lora_req.adapter_id if lora_req else None,
                        lora_name=lora_req.name if lora_req else None,
                        extra_keys=secondary_extra_keys,
                        block_size=self.hash_block_size,
                    )
                block_meta.append(meta_by_hash)

    return gpu_block_ids, advanced_per_group, block_meta

bind_gpu_block_pool(gpu_block_pool)

Bind GPU block pool so that we can touch blocks during stores. Called by Scheduler after kv_cache_manager is ready.

Source code in vllm/v1/simple_kv_offload/manager.py
def bind_gpu_block_pool(self, gpu_block_pool: BlockPool) -> None:
    """Bind GPU block pool so that we can touch blocks during stores.
    Called by Scheduler after kv_cache_manager is ready."""
    self._gpu_block_pool = gpu_block_pool

get_num_new_matched_tokens(request, num_computed_tokens)

Return (num_new_tokens, is_async) from consecutive CPU cache hits.

Source code in vllm/v1/simple_kv_offload/manager.py
def get_num_new_matched_tokens(
    self, request: "Request", num_computed_tokens: int
) -> tuple[int | None, bool]:
    """Return (num_new_tokens, is_async) from consecutive CPU cache hits."""
    # Pins found CPU blocks so they survive LRU eviction until
    # update_state_after_alloc() consumes them. Any pin from an earlier
    # call on the same request (e.g. retry after a failed allocate_slots)
    # is dropped first.
    if stale := self._pending_cpu_hits.pop(request.request_id, None):
        self._free_pending_cpu_hit(stale)

    if request.skip_reading_prefix_cache:
        return 0, False

    if num_computed_tokens % self.block_size != 0:
        # Transfers are whole-block copies, so an external suffix cannot
        # start in the middle of a destination block. When the GPU-side
        # coordinator also serves fine-grained hits, the local prefix can
        # land on a hash boundary that is not a scheduler-block boundary;
        # such a request keeps its local hit and skips the external one.
        logger.warning_once(
            "SimpleCPU external lookup requires scheduler-block-aligned "
            "local tokens, got %d tokens with scheduler block size %d.",
            num_computed_tokens,
            self.block_size,
        )
        return 0, False

    num_skipped_hashes = num_computed_tokens // self.hash_block_size
    remaining_hashes = request.block_hashes[num_skipped_hashes:]

    if not remaining_hashes:
        return 0, False
    # Must recompute at least the last token, matching the logic in
    # kv_cache_manager.get_computed_blocks().
    max_hit_len = request.num_tokens - 1 - num_computed_tokens
    if max_hit_len <= 0:
        return 0, False
    cpu_hit_blocks, hit_length, _ = self.cpu_coordinator.find_longest_cache_hit(
        remaining_hashes, max_hit_len
    )

    if hit_length > 0:
        pin_blocks = [
            blk for grp in cpu_hit_blocks for blk in grp if not blk.is_null
        ]
        self.cpu_block_pool.touch(pin_blocks)
        self._pending_cpu_hits[request.request_id] = (
            cpu_hit_blocks,
            hit_length,
            num_computed_tokens,
        )
        return hit_length, True
    return 0, False

get_stats()

Drain per-step stats for the connector's stats hooks.

Source code in vllm/v1/simple_kv_offload/manager.py
def get_stats(self) -> SimpleCPUOffloadStats:
    """Drain per-step stats for the connector's stats hooks."""
    stats = SimpleCPUOffloadStats()

    current = self.boundary_store_stats
    for outcome, field_name in OUTCOME_TO_FIELD.items():
        delta = getattr(current, field_name) - getattr(
            self._boundary_stats_snapshot, field_name
        )
        if delta > 0:
            stats.increase_counter(MetricName.SAVE_OUTCOMES, delta, (outcome,))
    self._boundary_stats_snapshot = replace(current)

    if self._interval_load_blocks_completed:
        stats.increase_counter(
            MetricName.LOAD_BLOCKS,
            self._interval_load_blocks_completed,
        )
        self._interval_load_blocks_completed = 0

    stats.set_gauge(
        MetricName.USED_BLOCKS,
        self.num_cpu_blocks - self.cpu_block_pool.get_num_free_blocks(),
    )
    pending = sum(
        len(t.cpu_block_ids)
        for t in (
            *self._store_event_to_blocks.values(),
            *self._pending_finished_stores,
            *self._abandoned_store_event_to_blocks.values(),
        )
    )
    stats.set_gauge(MetricName.PENDING_STORE_BLOCKS, pending)
    stats.set_gauge(MetricName.INFO, 1, self._info_labelvalues)
    return stats

has_pending_stores()

Return True if a store transfer is queued or in flight.

Source code in vllm/v1/simple_kv_offload/manager.py
def has_pending_stores(self) -> bool:
    """Return True if a store transfer is queued or in flight."""
    return bool(
        self._pending_finished_stores
        or self._store_event_to_blocks
        or self._abandoned_store_event_to_blocks
    )

prepare_store_specs(scheduler_output)

Prepare store specs for the store event.

Source code in vllm/v1/simple_kv_offload/manager.py
def prepare_store_specs(
    self, scheduler_output: SchedulerOutput
) -> tuple[
    list[int],
    list[int],
    list[str],
    list[dict[BlockHashWithGroupId, BlockStoreMeta]] | None,
]:
    """Prepare store specs for the store event."""
    if self._lazy_mode:
        return self._prepare_lazy_store_specs()
    else:
        return self._prepare_eager_store_specs(scheduler_output)

request_finished(request, block_ids)

Always returns (False, None). GPU blocks are protected by ref_cnt, so the scheduler can free blocks immediately.

Source code in vllm/v1/simple_kv_offload/manager.py
def request_finished(
    self,
    request: "Request",
    block_ids: list[int],
) -> tuple[bool, dict[str, Any] | None]:
    """Always returns (False, None). GPU blocks are protected by ref_cnt,
    so the scheduler can free blocks immediately."""
    req_id = request.request_id

    # Release any temp CPU hit pin from get_num_new_matched_tokens()
    # if request is canceled or preempted before update_state_after_alloc()
    pending = self._pending_cpu_hits.pop(req_id, None)
    if pending is not None:
        self._free_pending_cpu_hit(pending)

    # Handle load: defer cleanup if load is in-flight
    load_state = self._reqs_to_load.get(req_id)
    if load_state is not None:
        if load_state.load_event is not None:
            load_state.finished = True  # Defer: load in-flight
        else:
            self._cleanup_load_request(req_id)

    # Handle store (eager mode only): defer cleanup if stores in-flight
    if not self._lazy_mode:
        store_state = self._reqs_to_store.get(req_id)
        if store_state is not None:
            if store_state.store_events:
                store_state.finished = True  # Defer: stores in-flight
            else:
                self._cleanup_store_request(req_id)

    return False, None

reset()

Abandon pending transfers and reset the CPU cache when safe.

Worker-side DMA may still be using blocks after reset is requested. Keep those block refs pinned until the existing completion path reports the transfer finished, then release refs without caching abandoned store results.

Source code in vllm/v1/simple_kv_offload/manager.py
def reset(self) -> bool:
    """Abandon pending transfers and reset the CPU cache when safe.

    Worker-side DMA may still be using blocks after reset is requested.
    Keep those block refs pinned until the existing completion path reports
    the transfer finished, then release refs without caching abandoned
    store results.
    """
    self._abandoned_store_event_to_blocks.update(self._store_event_to_blocks)
    for transfer in self._pending_finished_stores:
        self._release_transfer_refs(transfer)
    self._pending_finished_stores.clear()

    self._store_event_to_blocks.clear()
    self._in_flight_store_gpu_blocks.clear()

    # Loads that have not been sent to the worker cannot have running DMA.
    # In-flight loads stay pinned and are cleaned up on completion.
    for req_id in list(self._reqs_to_load):
        state = self._reqs_to_load.pop(req_id)
        if state.load_event is None:
            self._reqs_to_load[req_id] = state
            self._cleanup_load_request(req_id)
        else:
            self._abandoned_reqs_to_load[req_id] = state

    self._reqs_to_store.clear()
    self._store_event_to_reqs.clear()
    self._store_event_pending_counts = {
        event_idx: count
        for event_idx, count in self._store_event_pending_counts.items()
        if event_idx in self._abandoned_store_event_to_blocks
    }
    self._cursor = None
    # Seed the fresh counters with outcomes since the last drain so a
    # mid-interval reset() does not drop them from the next report.
    undrained = BoundaryStoreStats(
        **{
            f.name: max(
                0,
                getattr(self.boundary_store_stats, f.name)
                - getattr(self._boundary_stats_snapshot, f.name),
            )
            for f in fields(BoundaryStoreStats)
        }
    )
    self.boundary_store_stats = undrained
    self._boundary_stats_snapshot = BoundaryStoreStats()
    # NOTE: _load_event_counter / _store_event_counter are not
    # reset as they are monotonic and must stay ahead of the workers
    # high-water marks to avoid event index collisions

    if self._abandoned_store_event_to_blocks or self._abandoned_reqs_to_load:
        return False

    return self.cpu_block_pool.reset_prefix_cache()

update_connector_output(connector_output)

Handle async transfer completions from worker.

Load completions arrive via finished_recving (real req_ids). Store completions arrive via kv_connector_worker_meta as per-event worker counts. We accumulate across steps and process a store event only when all workers have reported completion.

Source code in vllm/v1/simple_kv_offload/manager.py
def update_connector_output(self, connector_output: KVConnectorOutput) -> None:
    """Handle async transfer completions from worker.

    Load completions arrive via finished_recving (real req_ids).
    Store completions arrive via kv_connector_worker_meta as
    per-event worker counts. We accumulate across steps and process
    a store event only when all workers have reported completion.
    """
    # --- Load completions ---
    for req_id in list(connector_output.finished_recving or []):
        completed_blocks = self._cleanup_load_request(req_id)
        if completed_blocks:
            self._interval_load_blocks_completed += completed_blocks

    # --- Store completions ---
    meta = connector_output.kv_connector_worker_meta
    if not isinstance(meta, SimpleCPUOffloadWorkerMetadata):
        return
    for event_idx, count in meta.completed_store_events.items():
        total = self._store_event_pending_counts.get(event_idx, 0) + count
        if total >= self._expected_worker_count:
            self._store_event_pending_counts.pop(event_idx, None)
            self._process_store_event(event_idx)
        else:
            self._store_event_pending_counts[event_idx] = total