1#[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
144pub struct G1Manager {
149 backend: G1ManagerBackend,
150}
151
152impl G1Manager {
153 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 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 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 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 #[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 #[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}