Skip to main content

dynamo_mocker/loadgen/
types.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::LocalBlockHash;
5use dynamo_kv_router::protocols::{
6    BlockHashOptions, ExternalSequenceBlockHash, WorkerId, compute_block_hash_for_seq,
7    compute_seq_hash_for_block,
8};
9use dynamo_tokens::SequenceHash;
10use uuid::Uuid;
11
12use super::trace::synthesize_validated_trace_tokens;
13use crate::common::protocols::DirectRequest;
14
15pub const OUTPUT_REPLAY_ID_ANNOTATION_KEY: &str = "output_replay_id";
16pub const OUTPUT_REPLAY_CONSUMER_RUNTIME_KEY: &str = "output_replay_consumer";
17
18pub fn output_replay_id_annotation(replay_key: &str) -> String {
19    format!("{OUTPUT_REPLAY_ID_ANNOTATION_KEY}:{replay_key}")
20}
21
22pub fn effective_replay_key(
23    request_id: Option<&str>,
24    session_id: Option<&str>,
25    turn_index: usize,
26    line_index: usize,
27) -> String {
28    if let Some(request_id) = request_id.map(str::trim).filter(|value| !value.is_empty()) {
29        return request_id.to_string();
30    }
31    if let Some(session_id) = session_id.map(str::trim).filter(|value| !value.is_empty()) {
32        return format!("{session_id}:{turn_index}");
33    }
34    format!("line:{line_index}")
35}
36
37#[derive(Debug, Clone, PartialEq)]
38pub struct Trace {
39    pub block_size: usize,
40    pub sessions: Vec<SessionTrace>,
41}
42
43#[derive(Debug, Clone, PartialEq)]
44pub struct AgenticTrace {
45    pub block_size: usize,
46    pub turns: Vec<AgenticTurnTrace>,
47}
48
49#[derive(Debug, Clone, PartialEq)]
50pub enum DynamoRequestTrace {
51    Standard(Trace),
52    Agentic(AgenticTrace),
53}
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56pub enum TraceFileFormat {
57    Mooncake,
58    /// Mooncake-shaped rows where follow-up turns contain new input deltas.
59    /// Offline replay accumulates each generated output and the next input delta
60    /// per session before computing engine block hashes. Use this only for delta
61    /// traces: it expands compact session turns into cumulative prompts and can
62    /// use much more memory than `Mooncake`.
63    MooncakeDelta,
64    /// Mooncake request/cache rows plus explicit request-level workflow
65    /// dependencies. Each row dispatches after `wait_for` completions plus its
66    /// authored delay/tool wait.
67    AgenticMooncake,
68    AppliedComputeAgentic,
69    Dynamo,
70}
71
72impl TraceFileFormat {
73    pub fn as_str(self) -> &'static str {
74        match self {
75            Self::Mooncake => "mooncake",
76            Self::MooncakeDelta => "mooncake-delta",
77            Self::AgenticMooncake => "agentic_mooncake",
78            Self::AppliedComputeAgentic => "applied_compute_agentic",
79            Self::Dynamo => "dynamo",
80        }
81    }
82}
83
84#[derive(Debug, Clone, PartialEq)]
85pub struct SessionTrace {
86    pub session_id: String,
87    pub first_arrival_timestamp_ms: Option<f64>,
88    pub turns: Vec<TurnTrace>,
89}
90
91#[derive(Debug, Clone, Default, PartialEq)]
92pub struct TurnTrace {
93    pub input_length: usize,
94    pub max_output_tokens: usize,
95    pub output_token_ids: Option<Vec<u32>>,
96    pub replay_key: Option<String>,
97    pub hash_ids: Vec<u32>,
98    pub delay_after_previous_ms: f64,
99    pub priority: i32,
100    pub strict_priority: u32,
101    pub policy_class: Option<String>,
102}
103
104#[derive(Debug, Clone, Default, PartialEq)]
105pub struct AgenticTurnTrace {
106    pub request_id: String,
107    pub session_id: String,
108    pub input_length: usize,
109    pub max_output_tokens: usize,
110    pub output_token_ids: Option<Vec<u32>>,
111    pub replay_key: Option<String>,
112    pub hash_ids: Vec<u32>,
113    pub first_ready_timestamp_ms: Option<f64>,
114    pub delay_after_dependencies_ms: f64,
115    pub priority: i32,
116    pub strict_priority: u32,
117    pub policy_class: Option<String>,
118    pub wait_for: Vec<String>,
119    pub prefix_reset: bool,
120}
121
122#[derive(Debug, Clone)]
123pub struct LengthSpec {
124    pub mean: usize,
125    pub stddev: f64,
126}
127
128#[derive(Debug, Clone)]
129pub enum ArrivalSpec {
130    Burst,
131    ConstantQps { qps: f64 },
132    PoissonQps { qps: f64 },
133    GammaQps { qps: f64, smoothness: f64 },
134}
135
136#[derive(Debug, Clone)]
137pub enum DelaySpec {
138    None,
139    ConstantMs(f64),
140    ExponentialMs { mean_ms: f64 },
141}
142
143#[derive(Debug, Clone)]
144pub struct SyntheticTraceSpec {
145    pub block_size: usize,
146    pub num_sessions: usize,
147    pub turns_per_session: usize,
148    pub input_tokens: LengthSpec,
149    pub output_tokens: LengthSpec,
150    pub shared_prefix_ratio: f64,
151    pub num_prefix_groups: usize,
152    pub first_turn_arrivals: ArrivalSpec,
153    pub inter_turn_delays: DelaySpec,
154    pub seed: u64,
155    pub arrival_seed: u64,
156}
157
158#[derive(Debug, Clone, Copy)]
159pub enum SequenceHashMode {
160    Raw,
161    Cumulative,
162}
163
164#[derive(Debug, Clone, Copy)]
165pub enum SessionPartitionSpec {
166    Random { num_partitions: usize, seed: u64 },
167    RoundRobin { num_partitions: usize },
168}
169
170#[derive(Debug, Clone)]
171pub struct RouterSequence {
172    pub worker_id: WorkerId,
173    pub local_hashes: Vec<LocalBlockHash>,
174    pub external_hashes: Vec<ExternalSequenceBlockHash>,
175}
176
177#[derive(Debug, Clone, PartialEq, Eq)]
178pub struct ReplayRequestHashes {
179    pub local_block_hashes: Vec<LocalBlockHash>,
180    pub sequence_hashes: Vec<SequenceHash>,
181}
182
183impl ReplayRequestHashes {
184    pub(crate) fn from_tokens(tokens: &[u32], engine_block_size: u32) -> Self {
185        let local_block_hashes =
186            compute_block_hash_for_seq(tokens, engine_block_size, BlockHashOptions::default());
187        let sequence_hashes = compute_seq_hash_for_block(&local_block_hashes);
188
189        Self {
190            local_block_hashes,
191            sequence_hashes,
192        }
193    }
194}
195
196#[derive(Debug, Clone)]
197pub struct ReadyTurn {
198    pub request_uuid: Uuid,
199    pub session_id: String,
200    pub turn_index: usize,
201    pub emit_session_metadata: bool,
202    pub replay_key: Option<String>,
203    pub scheduled_ready_at_ms: f64,
204    pub replay_hashes: Option<ReplayRequestHashes>,
205    pub request: DirectRequest,
206}
207
208/// A request whose prompt may still be represented by one hash id per trace
209/// block. Offline replay keeps this compact form while an aggregated or
210/// prefill router queues the request and materializes tokens only when a
211/// worker admits it.
212#[derive(Debug)]
213pub(crate) enum ReplayRequestPayload {
214    Materialized(DirectRequest),
215    Deferred {
216        request_metadata: DirectRequest,
217        input_length: usize,
218        hash_ids: Vec<u32>,
219        trace_block_size: usize,
220    },
221}
222
223impl ReplayRequestPayload {
224    pub(crate) fn materialized(request: DirectRequest) -> Self {
225        Self::Materialized(request)
226    }
227
228    pub(super) fn deferred(
229        request_metadata: DirectRequest,
230        input_length: usize,
231        hash_ids: Vec<u32>,
232        trace_block_size: usize,
233    ) -> Self {
234        debug_assert!(request_metadata.tokens.is_empty());
235        Self::Deferred {
236            request_metadata,
237            input_length,
238            hash_ids,
239            trace_block_size,
240        }
241    }
242
243    pub(crate) fn input_length(&self) -> usize {
244        match self {
245            Self::Materialized(request) => request.tokens.len(),
246            Self::Deferred { input_length, .. } => *input_length,
247        }
248    }
249
250    pub(crate) fn metadata(&self) -> &DirectRequest {
251        match self {
252            Self::Materialized(request) => request,
253            Self::Deferred {
254                request_metadata, ..
255            } => request_metadata,
256        }
257    }
258
259    pub(crate) fn metadata_mut(&mut self) -> &mut DirectRequest {
260        match self {
261            Self::Materialized(request) => request,
262            Self::Deferred {
263                request_metadata, ..
264            } => request_metadata,
265        }
266    }
267
268    pub(crate) fn materialized_tokens(&self) -> Option<&[u32]> {
269        match self {
270            Self::Materialized(request) => Some(&request.tokens),
271            Self::Deferred { .. } => None,
272        }
273    }
274
275    pub(crate) fn materialized_request(&self) -> Option<&DirectRequest> {
276        match self {
277            Self::Materialized(request) => Some(request),
278            Self::Deferred { .. } => None,
279        }
280    }
281
282    pub(crate) fn prompt_tokens(&self) -> Vec<u32> {
283        match self {
284            Self::Materialized(request) => request.tokens.clone(),
285            Self::Deferred {
286                input_length,
287                hash_ids,
288                trace_block_size,
289                ..
290            } => synthesize_validated_trace_tokens(*input_length, hash_ids, *trace_block_size),
291        }
292    }
293
294    pub(crate) fn into_direct_request(self) -> DirectRequest {
295        match self {
296            Self::Materialized(request) => request,
297            Self::Deferred {
298                mut request_metadata,
299                input_length,
300                hash_ids,
301                trace_block_size,
302            } => {
303                request_metadata.tokens =
304                    synthesize_validated_trace_tokens(input_length, &hash_ids, trace_block_size);
305                request_metadata
306            }
307        }
308    }
309
310    pub(crate) fn materialize(&mut self) -> Option<&DirectRequest> {
311        if matches!(self, Self::Deferred { .. }) {
312            let payload = std::mem::replace(self, Self::Materialized(DirectRequest::default()));
313            *self = Self::Materialized(payload.into_direct_request());
314        }
315        self.materialized_request()
316    }
317}
318
319#[derive(Debug)]
320pub(crate) struct CompactReadyTurn {
321    pub(crate) request_uuid: Uuid,
322    pub(crate) session_id: String,
323    pub(crate) turn_index: usize,
324    pub(crate) replay_key: Option<String>,
325    pub(crate) scheduled_ready_at_ms: f64,
326    pub(crate) replay_hashes: Option<ReplayRequestHashes>,
327    pub(crate) emit_session_metadata: bool,
328    pub(crate) request: ReplayRequestPayload,
329}
330
331impl CompactReadyTurn {
332    pub(crate) fn into_ready_turn(self) -> ReadyTurn {
333        ReadyTurn {
334            request_uuid: self.request_uuid,
335            session_id: self.session_id,
336            turn_index: self.turn_index,
337            emit_session_metadata: self.emit_session_metadata,
338            replay_key: self.replay_key,
339            scheduled_ready_at_ms: self.scheduled_ready_at_ms,
340            replay_hashes: self.replay_hashes,
341            request: self.request.into_direct_request(),
342        }
343    }
344}