Skip to main content

dynamo_mocker/kv_manager/
g1_manager.rs

1// SPDX-FileCopyrightText: Copyright (c) 2024-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Temporary scheduler-facing facade for selecting a G1 implementation.
5//!
6//! This module is not a third KV manager. It forwards each operation to either
7//! KVBM G1 or vLLM G1, selected when the manager is constructed. Remove this
8//! facade when the KVBM G1 implementation is removed.
9
10#[cfg(feature = "kvbm-offload")]
11use std::sync::{Arc, Mutex};
12
13#[cfg(feature = "kvbm-offload")]
14use dynamo_tokens::PositionalLineageHash;
15use dynamo_tokens::blocks::UniqueBlock;
16#[cfg(feature = "kvbm-offload")]
17use kvbm_logical::ImmutableBlock;
18use uuid::Uuid;
19
20#[cfg(feature = "kvbm-offload")]
21use crate::common::protocols::G1;
22use crate::common::protocols::{
23    G1Backend, KvEventPublishers, MockerEvictionBackend, MoveBlock, PrefillCost,
24};
25use crate::common::sequence::ActiveSequence;
26#[cfg(feature = "kvbm-offload")]
27use crate::kvbm_offload::{MockOffloadEngine, OffloadId};
28
29use super::G1Acquire;
30use super::kvbm_backend;
31#[cfg(feature = "kvbm-offload")]
32use super::kvbm_backend::{BatchSwapInOutcome, SwapInRegistrationBlock, SwapInRegistrationOutcome};
33use super::vllm_backend::{
34    NativeDecodeBlockReservation, NativeDestinationReservation, VllmAcquire, VllmBlockLayout,
35    VllmKvManager,
36};
37
38enum G1ManagerBackend {
39    Kvbm(kvbm_backend::KvManager),
40    Native(VllmKvManager),
41}
42
43fn validate_plh_alignment(blocks: &[UniqueBlock], plh_count: usize) {
44    let full_blocks = blocks
45        .iter()
46        .filter(|block| matches!(block, UniqueBlock::FullBlock(_)))
47        .count();
48    assert_eq!(plh_count, full_blocks, "PLHs must align with full blocks");
49}
50
51fn into_g1_acquire<T>(outcome: VllmAcquire<T>) -> G1Acquire<T> {
52    match outcome {
53        VllmAcquire::Ready(value) => G1Acquire::Ready(value),
54        VllmAcquire::CapacityExhausted => G1Acquire::CapacityExhausted,
55    }
56}
57
58fn process_vllm_event(
59    manager: &mut VllmKvManager,
60    owner: Uuid,
61    event: &MoveBlock,
62    reusable_prefix_blocks: usize,
63) -> VllmAcquire<usize> {
64    match event {
65        MoveBlock::Use(blocks, local_hashes, plhs, token_ids, parent) => {
66            validate_plh_alignment(blocks, plhs.len());
67            manager.use_for_request(
68                owner,
69                blocks,
70                local_hashes,
71                token_ids.as_deref(),
72                parent.as_ref(),
73                reusable_prefix_blocks,
74            )
75        }
76        MoveBlock::Deref(blocks) => {
77            assert_eq!(reusable_prefix_blocks, 0);
78            manager.deref_for_request(owner, blocks);
79            VllmAcquire::Ready(1)
80        }
81        MoveBlock::Promote(uuid, hash, parent_hash, local_hash, _plh, token_ids) => {
82            assert_eq!(reusable_prefix_blocks, 0);
83            manager.promote_for_request(
84                owner,
85                *uuid,
86                *hash,
87                *parent_hash,
88                *local_hash,
89                token_ids.clone(),
90            );
91            VllmAcquire::Ready(1)
92        }
93    }
94}
95
96pub(crate) struct DecodeBlockReservation {
97    inner: DecodeBlockReservationBackend,
98}
99
100enum DecodeBlockReservationBackend {
101    Kvbm(kvbm_backend::DecodeBlockReservation),
102    Native(NativeDecodeBlockReservation),
103}
104
105impl DecodeBlockReservation {
106    pub(crate) fn len(&self) -> usize {
107        match &self.inner {
108            DecodeBlockReservationBackend::Kvbm(reservation) => reservation.len(),
109            DecodeBlockReservationBackend::Native(reservation) => reservation.len(),
110        }
111    }
112}
113
114pub(crate) struct DestinationReservation {
115    inner: DestinationReservationBackend,
116}
117
118enum DestinationReservationBackend {
119    Kvbm(kvbm_backend::VllmDestinationReservation),
120    Native(NativeDestinationReservation),
121}
122
123impl DestinationReservation {
124    pub(crate) fn transferable_prompt_tokens(&self, block_size: usize) -> usize {
125        match &self.inner {
126            DestinationReservationBackend::Kvbm(reservation) => {
127                reservation.transferable_prompt_tokens(block_size)
128            }
129            DestinationReservationBackend::Native(reservation) => {
130                reservation.transferable_prompt_tokens(block_size)
131            }
132        }
133    }
134
135    #[cfg(test)]
136    pub(crate) fn len(&self) -> usize {
137        match &self.inner {
138            DestinationReservationBackend::Kvbm(reservation) => reservation.len(),
139            DestinationReservationBackend::Native(reservation) => reservation.len(),
140        }
141    }
142}
143
144/// Temporary facade over the G1 implementation selected by [`G1Backend`].
145///
146/// This is not a separate manager implementation and will be removed with the
147/// KVBM G1 backend.
148pub struct G1Manager {
149    backend: G1ManagerBackend,
150}
151
152impl G1Manager {
153    /// Constructs the default KVBM G1 backend.
154    pub fn new_with_event_sink(
155        max_capacity: usize,
156        block_size: usize,
157        kv_event_publishers: KvEventPublishers,
158        dp_rank: u32,
159    ) -> Self {
160        Self::new_with_backend(
161            max_capacity,
162            block_size,
163            kv_event_publishers,
164            dp_rank,
165            G1Backend::Kvbm,
166        )
167    }
168
169    pub fn new_with_backend(
170        max_capacity: usize,
171        block_size: usize,
172        kv_event_publishers: KvEventPublishers,
173        dp_rank: u32,
174        backend: G1Backend,
175    ) -> Self {
176        Self::new_with_backend_and_caching(
177            max_capacity,
178            block_size,
179            kv_event_publishers,
180            dp_rank,
181            backend,
182            true,
183        )
184    }
185
186    pub fn new_with_backend_and_caching(
187        max_capacity: usize,
188        block_size: usize,
189        kv_event_publishers: KvEventPublishers,
190        dp_rank: u32,
191        backend: G1Backend,
192        enable_prefix_caching: bool,
193    ) -> Self {
194        let backend = match backend {
195            G1Backend::Kvbm => {
196                G1ManagerBackend::Kvbm(kvbm_backend::KvManager::new_with_event_sink(
197                    max_capacity,
198                    block_size,
199                    kv_event_publishers,
200                    dp_rank,
201                ))
202            }
203            G1Backend::Native => G1ManagerBackend::Native(VllmKvManager::new_with_event_sink(
204                max_capacity,
205                block_size,
206                enable_prefix_caching,
207                kv_event_publishers,
208                dp_rank,
209            )),
210        };
211        Self { backend }
212    }
213
214    /// Make newly allocated full blocks prefix-cache-visible once the current
215    /// scheduling decision reaches their complete token range. KVBM retains
216    /// its historical eager lifecycle; the native backend mirrors vLLM's
217    /// per-request `allocate_slots()` / `cache_blocks()` boundary.
218    pub(crate) fn finalize_computed_prefix(
219        &mut self,
220        owner: Uuid,
221        computed_before: usize,
222        computed_after: usize,
223        sequence: &mut ActiveSequence,
224    ) {
225        if let G1ManagerBackend::Native(manager) = &mut self.backend {
226            // A single scheduling decision can complete several prompt blocks
227            // and the mutable tail together. Publish/cache the earlier full
228            // blocks first so the promoted tail never references a parent that
229            // is not cache-visible yet.
230            manager.finalize_computed_prefix(owner, computed_before, computed_after);
231            if let Some(promote) = sequence.promote_computed_tail(computed_after) {
232                assert!(
233                    matches!(
234                        process_vllm_event(manager, owner, &promote, 0),
235                        VllmAcquire::Ready(_)
236                    ),
237                    "computed-tail promotion must be infallible"
238                );
239            }
240        }
241    }
242
243    pub fn new_with_eviction_backend(
244        max_capacity: usize,
245        block_size: usize,
246        kv_event_publishers: KvEventPublishers,
247        dp_rank: u32,
248        eviction_backend: MockerEvictionBackend,
249    ) -> Self {
250        Self {
251            backend: G1ManagerBackend::Kvbm(kvbm_backend::KvManager::new_with_eviction_backend(
252                max_capacity,
253                block_size,
254                kv_event_publishers,
255                dp_rank,
256                eviction_backend,
257            )),
258        }
259    }
260
261    /// Owner-aware production entrypoint. KVBM deliberately ignores `owner`;
262    /// native G1 uses it to address the exact physical request block table.
263    pub(crate) fn process_for_request(
264        &mut self,
265        owner: Uuid,
266        event: &MoveBlock,
267        reusable_prefix_blocks: usize,
268    ) -> G1Acquire<usize> {
269        match &mut self.backend {
270            G1ManagerBackend::Kvbm(manager) => manager.process(event),
271            G1ManagerBackend::Native(manager) => into_g1_acquire(process_vllm_event(
272                manager,
273                owner,
274                event,
275                reusable_prefix_blocks,
276            )),
277        }
278    }
279
280    /// Compatibility entrypoint for KVBM-focused unit tests. Native callers
281    /// must provide request ownership through [`Self::process_for_request`].
282    #[cfg(test)]
283    pub(crate) fn process(&mut self, event: &MoveBlock) -> G1Acquire<usize> {
284        match &mut self.backend {
285            G1ManagerBackend::Kvbm(manager) => manager.process(event),
286            G1ManagerBackend::Native(_) => {
287                panic!("native G1 operations require a request owner")
288            }
289        }
290    }
291
292    pub(crate) fn reserve_decode_blocks(
293        &mut self,
294        count: usize,
295    ) -> G1Acquire<DecodeBlockReservation> {
296        match &mut self.backend {
297            G1ManagerBackend::Kvbm(manager) => {
298                manager
299                    .reserve_decode_blocks(count)
300                    .map(|inner| DecodeBlockReservation {
301                        inner: DecodeBlockReservationBackend::Kvbm(inner),
302                    })
303            }
304            G1ManagerBackend::Native(manager) => {
305                into_g1_acquire(manager.reserve_decode_blocks(count)).map(|inner| {
306                    DecodeBlockReservation {
307                        inner: DecodeBlockReservationBackend::Native(inner),
308                    }
309                })
310            }
311        }
312    }
313
314    pub(crate) fn process_decode_signal_for_request(
315        &mut self,
316        owner: Uuid,
317        event: &MoveBlock,
318        reservation: &mut DecodeBlockReservation,
319    ) {
320        match (&mut self.backend, &mut reservation.inner) {
321            (G1ManagerBackend::Kvbm(manager), DecodeBlockReservationBackend::Kvbm(inner)) => {
322                manager.process_decode_signal(event, inner);
323            }
324            (G1ManagerBackend::Native(manager), DecodeBlockReservationBackend::Native(inner)) => {
325                match event {
326                    MoveBlock::Use(blocks, local_hashes, plhs, token_ids, parent) => {
327                        validate_plh_alignment(blocks, plhs.len());
328                        manager.use_decode_reservation_for_request(
329                            owner,
330                            blocks,
331                            local_hashes,
332                            token_ids.as_deref(),
333                            parent.as_ref(),
334                            inner,
335                        );
336                    }
337                    _ => assert!(
338                        matches!(
339                            process_vllm_event(manager, owner, event, 0),
340                            VllmAcquire::Ready(_)
341                        ),
342                        "non-Use decode signal must be infallible"
343                    ),
344                }
345            }
346            _ => panic!("decode reservation belongs to a different G1 backend"),
347        }
348    }
349
350    pub(crate) fn release_decode_reservation(&mut self, reservation: DecodeBlockReservation) {
351        match (&mut self.backend, reservation.inner) {
352            (G1ManagerBackend::Kvbm(_), DecodeBlockReservationBackend::Kvbm(inner)) => drop(inner),
353            (G1ManagerBackend::Native(manager), DecodeBlockReservationBackend::Native(inner)) => {
354                manager.release_decode_reservation(inner);
355            }
356            _ => panic!("decode reservation belongs to a different G1 backend"),
357        }
358    }
359
360    pub(crate) fn reserve_destination_at(
361        &mut self,
362        owner: Uuid,
363        sequence: &ActiveSequence,
364        eviction_now_ms: Option<f64>,
365    ) -> G1Acquire<DestinationReservation> {
366        match &mut self.backend {
367            G1ManagerBackend::Kvbm(manager) => manager
368                .reserve_destination_at(sequence, eviction_now_ms)
369                .map(|inner| DestinationReservation {
370                    inner: DestinationReservationBackend::Kvbm(inner),
371                }),
372            G1ManagerBackend::Native(manager) => {
373                let layout = match sequence.prepare_allocation(sequence.num_input_tokens()) {
374                    None => None,
375                    Some(MoveBlock::Use(blocks, local_hashes, plhs, token_ids, parent)) => {
376                        validate_plh_alignment(&blocks, plhs.len());
377                        Some(VllmBlockLayout::new(
378                            blocks,
379                            local_hashes,
380                            token_ids,
381                            parent,
382                        ))
383                    }
384                    Some(_) => panic!("destination allocation must be a Use signal"),
385                };
386                into_g1_acquire(manager.reserve_destination_at(owner, layout, eviction_now_ms)).map(
387                    |inner| DestinationReservation {
388                        inner: DestinationReservationBackend::Native(inner),
389                    },
390                )
391            }
392        }
393    }
394
395    pub(crate) fn activate_destination(&mut self, reservation: DestinationReservation) {
396        match (&mut self.backend, reservation.inner) {
397            (G1ManagerBackend::Kvbm(manager), DestinationReservationBackend::Kvbm(inner)) => {
398                manager.activate_destination(inner);
399            }
400            (G1ManagerBackend::Native(manager), DestinationReservationBackend::Native(inner)) => {
401                manager.activate_destination(inner);
402            }
403            _ => panic!("destination reservation belongs to a different G1 backend"),
404        }
405    }
406
407    pub(crate) fn cancel_destination(&mut self, reservation: DestinationReservation) {
408        match (&mut self.backend, reservation.inner) {
409            (G1ManagerBackend::Kvbm(_), DestinationReservationBackend::Kvbm(inner)) => drop(inner),
410            (G1ManagerBackend::Native(manager), DestinationReservationBackend::Native(inner)) => {
411                manager.cancel_destination(inner);
412            }
413            _ => panic!("destination reservation belongs to a different G1 backend"),
414        }
415    }
416
417    pub fn num_active_blocks(&self) -> usize {
418        match &self.backend {
419            G1ManagerBackend::Kvbm(manager) => manager.num_active_blocks(),
420            G1ManagerBackend::Native(manager) => manager.num_active_blocks(),
421        }
422    }
423
424    pub fn num_active_block_refs(&self) -> usize {
425        match &self.backend {
426            G1ManagerBackend::Kvbm(manager) => manager.num_active_block_refs(),
427            G1ManagerBackend::Native(manager) => manager.num_active_block_refs(),
428        }
429    }
430
431    pub fn num_inactive_blocks(&self) -> usize {
432        match &self.backend {
433            G1ManagerBackend::Kvbm(manager) => manager.num_inactive_blocks(),
434            G1ManagerBackend::Native(manager) => manager.num_inactive_blocks(),
435        }
436    }
437
438    pub fn get_active_perc(&self) -> f64 {
439        match &self.backend {
440            G1ManagerBackend::Kvbm(manager) => manager.get_active_perc(),
441            G1ManagerBackend::Native(manager) => manager.get_active_perc(),
442        }
443    }
444
445    pub fn max_capacity(&self) -> usize {
446        match &self.backend {
447            G1ManagerBackend::Kvbm(manager) => manager.max_capacity(),
448            G1ManagerBackend::Native(manager) => manager.max_capacity(),
449        }
450    }
451
452    pub fn block_size(&self) -> usize {
453        match &self.backend {
454            G1ManagerBackend::Kvbm(manager) => manager.block_size(),
455            G1ManagerBackend::Native(manager) => manager.block_size(),
456        }
457    }
458
459    pub fn dp_rank(&self) -> u32 {
460        match &self.backend {
461            G1ManagerBackend::Kvbm(manager) => manager.dp_rank(),
462            G1ManagerBackend::Native(manager) => manager.dp_rank(),
463        }
464    }
465
466    #[cfg(test)]
467    pub(crate) fn request_block_count(&self, owner: Uuid, sequence: &ActiveSequence) -> usize {
468        match &self.backend {
469            G1ManagerBackend::Kvbm(manager) => manager.active_block_ids(sequence).len(),
470            G1ManagerBackend::Native(manager) => manager.request_block_count(owner),
471        }
472    }
473
474    pub fn get_prefill_cost(&self, sequence: &ActiveSequence) -> PrefillCost {
475        match &self.backend {
476            G1ManagerBackend::Kvbm(manager) => manager.get_prefill_cost(sequence),
477            G1ManagerBackend::Native(manager) => manager.get_prefill_cost(sequence),
478        }
479    }
480
481    #[cfg(feature = "kvbm-offload")]
482    pub fn attach_new_offload_engine(
483        &mut self,
484        engine: MockOffloadEngine,
485    ) -> Arc<Mutex<MockOffloadEngine>> {
486        match &mut self.backend {
487            G1ManagerBackend::Kvbm(manager) => manager.attach_new_offload_engine(engine),
488            G1ManagerBackend::Native(_) => {
489                panic!("legacy kvbm-offload cannot be attached to native G1")
490            }
491        }
492    }
493
494    /// Returns whether the legacy KVBM backend has a lower-tier offload
495    /// engine. Native G1 always returns `false`, so native requests cannot
496    /// enter the KVBM swap-in lifecycle.
497    #[cfg(feature = "kvbm-offload")]
498    pub fn has_offload_engine(&self) -> bool {
499        match &self.backend {
500            G1ManagerBackend::Kvbm(manager) => manager.has_offload_engine(),
501            G1ManagerBackend::Native(_) => false,
502        }
503    }
504
505    #[cfg(feature = "kvbm-offload")]
506    pub fn tick_offload_engine(&mut self, now_ms: f64) {
507        if let G1ManagerBackend::Kvbm(manager) = &mut self.backend {
508            manager.tick_offload_engine(now_ms);
509        }
510    }
511
512    #[cfg(feature = "kvbm-offload")]
513    pub fn earliest_offload_deadline(&self) -> Option<f64> {
514        match &self.backend {
515            G1ManagerBackend::Kvbm(manager) => manager.earliest_offload_deadline(),
516            G1ManagerBackend::Native(_) => None,
517        }
518    }
519
520    pub(crate) fn refresh_offload_dependency(
521        &self,
522        dependency: super::OffloadDependency,
523    ) -> Option<super::OffloadDependency> {
524        match &self.backend {
525            G1ManagerBackend::Kvbm(manager) => manager.refresh_offload_dependency(dependency),
526            G1ManagerBackend::Native(_) => None,
527        }
528    }
529
530    #[cfg(feature = "kvbm-offload")]
531    pub fn try_batch_swap_in(
532        &mut self,
533        remaining_plhs: &[PositionalLineageHash],
534        prefix_pins: Vec<ImmutableBlock<G1>>,
535        now_ms: Option<f64>,
536    ) -> BatchSwapInOutcome {
537        match &mut self.backend {
538            G1ManagerBackend::Kvbm(manager) => {
539                manager.try_batch_swap_in(remaining_plhs, prefix_pins, now_ms)
540            }
541            G1ManagerBackend::Native(_) => {
542                drop(prefix_pins);
543                BatchSwapInOutcome::NoHits
544            }
545        }
546    }
547
548    #[cfg(feature = "kvbm-offload")]
549    pub(crate) fn cancel_swap_in(&mut self, id: OffloadId) -> bool {
550        match &mut self.backend {
551            G1ManagerBackend::Kvbm(manager) => manager.cancel_swap_in(id),
552            G1ManagerBackend::Native(_) => false,
553        }
554    }
555
556    #[cfg(feature = "kvbm-offload")]
557    pub(crate) fn register_completed_swap_in(
558        &mut self,
559        id: OffloadId,
560        entries: Vec<SwapInRegistrationBlock>,
561        parent_hash: Option<u64>,
562    ) -> SwapInRegistrationOutcome {
563        match &mut self.backend {
564            G1ManagerBackend::Kvbm(manager) => {
565                manager.register_completed_swap_in(id, entries, parent_hash)
566            }
567            G1ManagerBackend::Native(_) => {
568                panic!("native G1 cannot complete a legacy KVBM swap-in")
569            }
570        }
571    }
572
573    #[cfg(feature = "kvbm-offload")]
574    pub(crate) fn try_pin_g1_prefix(
575        &mut self,
576        prefix_plhs: &[PositionalLineageHash],
577    ) -> Option<Vec<ImmutableBlock<G1>>> {
578        match &mut self.backend {
579            G1ManagerBackend::Kvbm(manager) => manager.try_pin_g1_prefix(prefix_plhs),
580            G1ManagerBackend::Native(_) => None,
581        }
582    }
583}
584
585#[cfg(test)]
586mod tests {
587    use dynamo_kv_router::protocols::KvCacheEventData;
588
589    use super::*;
590
591    #[test]
592    #[should_panic(expected = "PLHs must align with full blocks")]
593    fn vllm_boundary_rejects_misaligned_lineage_metadata() {
594        let mut manager =
595            G1Manager::new_with_backend(2, 4, KvEventPublishers::default(), 0, G1Backend::Native);
596        let event = MoveBlock::Use(
597            vec![UniqueBlock::FullBlock(7)],
598            vec![107],
599            Vec::new(),
600            None,
601            None,
602        );
603
604        let _ = manager.process_for_request(Uuid::from_u128(1), &event, 0);
605    }
606
607    #[test]
608    fn native_finalization_publishes_parent_before_promoted_tail() {
609        let stored = crate::replay::native_g1_parent_chain_artifact(4)
610            .kv_events
611            .into_iter()
612            .map(|event| match event.event.data {
613                KvCacheEventData::Stored(stored) => stored,
614                other => panic!("expected Stored event, got {other:?}"),
615            })
616            .collect::<Vec<_>>();
617        assert_eq!(stored.len(), 2);
618        assert_eq!(stored[0].blocks.len(), 2);
619        assert_eq!(stored[0].parent_hash, None);
620        assert_eq!(stored[1].blocks.len(), 1);
621        assert_eq!(
622            stored[1].parent_hash,
623            Some(stored[0].blocks[1].block_hash),
624            "the promoted tail must be published after its parent block"
625        );
626    }
627}