1use 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 MooncakeDelta,
64 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#[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}