Skip to main content

dynamo_mocker/replay/
artifacts.rs

1// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4use dynamo_kv_router::protocols::{KvCacheEvent, StorageTier};
5use uuid::Uuid;
6
7use crate::common::protocols::OutputSignal;
8use crate::loadgen::ReplayRequestHashes;
9
10/// Build the minimal Native G1 artifact that exercises the producer→indexer
11/// parent-ordering contract.
12///
13/// This is test-only plumbing shared by the mocker unit regression and
14/// downstream indexer parity coverage.
15#[cfg(any(test, feature = "test-support"))]
16#[doc(hidden)]
17pub fn native_g1_parent_chain_artifact(block_size: usize) -> ReplayWorkerArtifacts {
18    use crate::common::protocols::{G1Backend, KvEventPublishers};
19    use crate::common::sequence::ActiveSequence;
20    use crate::kv_manager::{G1Acquire, G1Manager};
21    use crate::scheduler::capture_router_event_sink;
22
23    assert!(block_size >= 2, "block size must be at least 2");
24    let computed_after = block_size
25        .checked_mul(3)
26        .expect("ordering regression token count overflow");
27    let prompt_len = computed_after - 2;
28    let prompt_len_u32 =
29        u32::try_from(prompt_len).expect("ordering regression prompt length must fit in u32");
30    let mut tokens = (0..prompt_len_u32).collect::<Vec<_>>();
31    let mut sequence = ActiveSequence::new(tokens.clone(), 2, Some(block_size), true, false);
32    let owner = Uuid::from_u128(1);
33    let (events, sink) = capture_router_event_sink(1);
34    let publishers = KvEventPublishers::new(Some(sink), None);
35    let mut manager = G1Manager::new_with_backend(3, block_size, publishers, 0, G1Backend::Native);
36
37    let creation = sequence
38        .take_creation_signal()
39        .expect("three-block sequence must allocate G1 blocks");
40    assert!(matches!(
41        manager.process_for_request(owner, &creation, 0),
42        G1Acquire::Ready(3)
43    ));
44
45    for token in [prompt_len_u32, prompt_len_u32 + 1] {
46        assert!(sequence.push(token).is_none());
47        tokens.push(token);
48    }
49    manager.finalize_computed_prefix(owner, 0, computed_after, &mut sequence);
50
51    let kv_events = events
52        .drain()
53        .into_iter()
54        .enumerate()
55        .map(|(ordinal, event)| ReplayTimedKvEvent {
56            event: event.event,
57            storage_tier: event.storage_tier,
58            timestamp_us: ordinal as u64,
59        })
60        .collect::<Vec<_>>();
61    let request_timestamp = kv_events.len() as u64;
62
63    ReplayWorkerArtifacts {
64        requests: vec![ReplayTimedRequest {
65            uuid: owner,
66            timestamp_us: request_timestamp,
67            scheduled_ready_at_ms: request_timestamp as f64 / 1000.0,
68            input_length: tokens.len(),
69            output_length: 0,
70            replay_hashes: ReplayRequestHashes::from_tokens(
71                &tokens,
72                u32::try_from(block_size).expect("block size must fit in u32"),
73            ),
74        }],
75        output_signals: Vec::new(),
76        kv_events,
77    }
78}
79
80#[derive(Debug, Clone)]
81pub struct ReplayTimedRequest {
82    pub uuid: Uuid,
83    pub timestamp_us: u64,
84    pub scheduled_ready_at_ms: f64,
85    pub input_length: usize,
86    pub output_length: usize,
87    pub replay_hashes: ReplayRequestHashes,
88}
89
90#[derive(Debug, Clone)]
91pub struct ReplayTimedOutputSignal {
92    pub signal: OutputSignal,
93    pub timestamp_us: u64,
94}
95
96#[derive(Debug, Clone)]
97pub struct ReplayTimedKvEvent {
98    pub event: KvCacheEvent,
99    pub storage_tier: StorageTier,
100    pub timestamp_us: u64,
101}
102
103#[derive(Debug, Clone, Default)]
104pub struct ReplayWorkerArtifacts {
105    pub requests: Vec<ReplayTimedRequest>,
106    pub output_signals: Vec<ReplayTimedOutputSignal>,
107    pub kv_events: Vec<ReplayTimedKvEvent>,
108}