1use std::collections::HashSet;
9use std::fs::File;
10use std::io::{BufRead, BufReader, Read};
11use std::path::{Path, PathBuf};
12
13use anyhow::{Context, Result, anyhow, bail};
14use flate2::read::MultiGzDecoder;
15use serde::Deserialize;
16use serde_json::Value;
17
18#[derive(Debug, Clone, Deserialize)]
19pub(crate) struct RequestTraceRecord {
20 pub(crate) schema: TraceSchema,
21 pub(crate) event_type: String,
22 pub(crate) event_time_unix_ms: u64,
23 #[serde(default)]
24 pub(crate) agent_context: Option<AgentContextFields>,
25 #[serde(default)]
26 pub(crate) request: Option<RequestTraceRequestMetrics>,
27 #[serde(default)]
28 pub(crate) tool: Option<RequestTraceToolEventMetrics>,
29}
30
31#[derive(Debug, Clone, Copy, Deserialize, PartialEq, Eq)]
32pub(crate) enum TraceSchema {
33 #[serde(rename = "dynamo.request.trace.v1")]
34 RequestV1,
35}
36
37#[derive(Debug, Clone, Deserialize)]
38pub(crate) struct AgentContextFields {
39 pub(crate) session_id: String,
40 #[serde(default)]
41 pub(crate) parent_session_id: Option<String>,
42}
43
44#[derive(Debug, Clone, Deserialize)]
45pub(crate) struct RequestTraceRequestMetrics {
46 pub(crate) request_id: String,
47 #[serde(default)]
48 pub(crate) output_tokens: Option<u64>,
49 #[serde(default)]
50 pub(crate) request_received_ms: Option<u64>,
51 #[serde(default)]
52 pub(crate) total_time_ms: Option<f64>,
53 #[serde(default)]
54 pub(crate) replay: Option<RequestTraceReplayMetrics>,
55}
56
57#[derive(Debug, Clone, Deserialize)]
58pub(crate) struct RequestTraceReplayMetrics {
59 pub(crate) trace_block_size: usize,
60 pub(crate) input_length: usize,
61 pub(crate) input_sequence_hashes: Vec<u64>,
62}
63
64#[derive(Debug, Clone, Deserialize)]
65pub(crate) struct RequestTraceToolEventMetrics {
66 pub(crate) tool_call_id: String,
67 pub(crate) tool_class: String,
68 #[serde(default)]
69 pub(crate) claude: Option<ClaudeToolReplayMetrics>,
70 #[serde(default)]
71 pub(crate) started_at_unix_ms: Option<u64>,
72 #[serde(default)]
73 pub(crate) ended_at_unix_ms: Option<u64>,
74 #[serde(default)]
75 pub(crate) duration_ms: Option<f64>,
76 #[serde(default)]
77 pub(crate) status: Option<String>,
78 #[serde(default)]
79 pub(crate) output_bytes: Option<u64>,
80 #[serde(default)]
81 pub(crate) output_tokens: Option<u64>,
82 #[serde(default)]
83 pub(crate) error_type: Option<String>,
84}
85
86#[derive(Debug, Clone, Deserialize)]
92pub(crate) struct ClaudeToolReplayMetrics {
93 pub(crate) source_request_id: String,
95 #[serde(default)]
97 pub(crate) consumer_request_id: Option<String>,
98 #[serde(default)]
100 pub(crate) child_session_id: Option<String>,
101 pub(crate) execution_mode: String,
103}
104
105#[derive(Debug, Clone)]
106pub struct RequestEntry {
107 pub(crate) start_ms: i64,
108 pub(crate) end_ms: i64,
109 pub(crate) agent_context: Option<AgentContextFields>,
110 pub(crate) request: RequestTraceRequestMetrics,
111 pub(crate) replay: RequestTraceReplayMetrics,
112}
113
114#[derive(Debug, Clone)]
115pub struct ToolEntry {
116 pub(crate) session_id: String,
117 pub(crate) start_ms: i64,
118 pub(crate) end_ms: i64,
119 pub(crate) tool_call_id: String,
120 pub(crate) tool_class: String,
121 pub(crate) claude: Option<ClaudeToolReplayMetrics>,
122 pub(crate) status: String,
123 pub(crate) duration_ms: f64,
124 pub(crate) output_bytes: Option<u64>,
125 pub(crate) output_tokens: Option<u64>,
126 pub(crate) error_type: Option<String>,
127}
128
129#[derive(Debug, Default)]
130pub struct LoadedAgentTrace {
131 pub requests: Vec<RequestEntry>,
132 pub tools: Vec<ToolEntry>,
133}
134
135#[derive(Debug, Clone, Copy, PartialEq, Eq)]
136pub enum RequestTraceMode {
137 Standard,
138 Agentic,
139}
140
141impl LoadedAgentTrace {
142 pub fn mode(&self) -> Result<RequestTraceMode> {
143 if self.requests.is_empty() {
144 bail!("Dynamo request trace contains no requests");
145 }
146 let contextual_requests = self
147 .requests
148 .iter()
149 .filter(|request| request.agent_context.is_some())
150 .count();
151 match contextual_requests {
152 0 => Ok(RequestTraceMode::Standard),
153 count if count == self.requests.len() => Ok(RequestTraceMode::Agentic),
154 _ => bail!("Dynamo request trace cannot mix requests with and without agent_context"),
155 }
156 }
157
158 pub fn ensure_agentic_compatible(&self) -> Result<()> {
159 if self
160 .requests
161 .iter()
162 .any(|request| request.agent_context.is_none())
163 {
164 bail!("agentic lowering requires agent_context on every request");
165 }
166 Ok(())
167 }
168}
169
170pub fn load_request_trace_records(paths: &[PathBuf]) -> Result<LoadedAgentTrace> {
173 let mut loaded = LoadedAgentTrace::default();
174 let mut request_ids = HashSet::new();
175
176 for path in paths {
177 let reader = open_trace_reader(path)?;
178 for (line_index, line) in reader.lines().enumerate() {
179 let line = line
180 .with_context(|| format!("failed to read {}:{}", path.display(), line_index + 1))?;
181 if line.trim().is_empty() {
182 continue;
183 }
184 let Some(record) = parse_trace_record(&line).with_context(|| {
185 format!("failed to parse {}:{}", path.display(), line_index + 1)
186 })?
187 else {
188 continue;
189 };
190 let _schema = record.schema;
191 if !matches!(
192 record.event_type.as_str(),
193 "request_end" | "tool_start" | "tool_end" | "tool_error"
194 ) {
195 bail!(
196 "request trace schema only supports request_end/tool_* events, got {} at {}:{}",
197 record.event_type,
198 path.display(),
199 line_index + 1
200 );
201 }
202 if record.event_type == "request_end" {
203 let entry = request_entry(record).with_context(|| {
204 format!(
205 "invalid request_end at {}:{}",
206 path.display(),
207 line_index + 1
208 )
209 })?;
210 if !request_ids.insert(entry.request.request_id.clone()) {
211 bail!(
212 "duplicate request_id {} at {}:{}",
213 entry.request.request_id,
214 path.display(),
215 line_index + 1
216 );
217 }
218 loaded.requests.push(entry);
219 } else if matches!(record.event_type.as_str(), "tool_end" | "tool_error") {
220 let terminal_event = record.event_type.clone();
221 if let Some(tool) = tool_entry(record, terminal_event) {
222 loaded.tools.push(tool);
223 }
224 }
225 }
226 }
227
228 if loaded.requests.is_empty() {
229 bail!("no request_end records with replay fields found");
230 }
231
232 Ok(loaded)
233}
234
235fn open_trace_reader(path: &Path) -> Result<Box<dyn BufRead>> {
236 let file = File::open(path).with_context(|| format!("failed to open {}", path.display()))?;
237 let reader: Box<dyn Read> = if path.extension().and_then(|ext| ext.to_str()) == Some("gz") {
238 Box::new(MultiGzDecoder::new(file))
239 } else {
240 Box::new(file)
241 };
242 Ok(Box::new(BufReader::new(reader)))
243}
244
245fn parse_trace_record(line: &str) -> Result<Option<RequestTraceRecord>> {
246 let value: Value = serde_json::from_str(line)?;
247 let event = value.get("event").unwrap_or(&value);
248 if !event.is_object() {
249 return Ok(None);
250 }
251 Ok(Some(serde_json::from_value(event.clone())?))
252}
253
254fn request_entry(record: RequestTraceRecord) -> Result<RequestEntry> {
255 let request = record
256 .request
257 .ok_or_else(|| anyhow!("request_end record is missing request payload"))?;
258 let replay = request
259 .replay
260 .clone()
261 .ok_or_else(|| anyhow!("request payload is missing replay metrics"))?;
262 if replay.trace_block_size == 0 {
263 bail!("request replay trace_block_size must be greater than 0");
264 }
265 let expected_hashes = replay.input_length.div_ceil(replay.trace_block_size);
266 if replay.input_sequence_hashes.len() != expected_hashes {
267 bail!(
268 "input_length {} with trace_block_size {} requires exactly {} replay hashes, got {}",
269 replay.input_length,
270 replay.trace_block_size,
271 expected_hashes,
272 replay.input_sequence_hashes.len()
273 );
274 }
275 if request.request_received_ms.is_none() {
276 bail!("request trace is missing request_received_ms");
277 }
278 if request.output_tokens.is_none() {
279 bail!("request trace is missing output_tokens");
280 }
281
282 let (start_ms, end_ms) = request_times(record.event_time_unix_ms, &request);
283 Ok(RequestEntry {
284 start_ms,
285 end_ms,
286 agent_context: record.agent_context,
287 request,
288 replay,
289 })
290}
291
292fn tool_entry(record: RequestTraceRecord, terminal_event: String) -> Option<ToolEntry> {
293 let context = record.agent_context?;
294 let tool = record.tool?;
295 let end_ms = tool
296 .ended_at_unix_ms
297 .map(saturating_i64)
298 .unwrap_or_else(|| saturating_i64(record.event_time_unix_ms));
299 let start_ms = tool
300 .started_at_unix_ms
301 .map(saturating_i64)
302 .or_else(|| {
303 tool.duration_ms
304 .map(|duration_ms| end_ms.saturating_sub(duration_ms.max(0.0).round() as i64))
305 })
306 .unwrap_or(end_ms);
307 if end_ms < start_ms {
308 return None;
309 }
310 let duration_ms = tool
311 .duration_ms
312 .unwrap_or_else(|| (end_ms - start_ms).max(0) as f64);
313 let status = tool.status.unwrap_or_else(|| {
314 if terminal_event == "tool_error" {
315 "error".to_string()
316 } else {
317 "succeeded".to_string()
318 }
319 });
320 Some(ToolEntry {
321 session_id: context.session_id,
322 start_ms,
323 end_ms,
324 tool_call_id: tool.tool_call_id,
325 tool_class: tool.tool_class,
326 claude: tool.claude,
327 status,
328 duration_ms,
329 output_bytes: tool.output_bytes,
330 output_tokens: tool.output_tokens,
331 error_type: tool.error_type,
332 })
333}
334
335pub(crate) fn request_times(
336 event_time_unix_ms: u64,
337 request: &RequestTraceRequestMetrics,
338) -> (i64, i64) {
339 let total_ms = request
340 .total_time_ms
341 .map(|value| value.max(0.0).round() as u64)
342 .unwrap_or_else(|| {
343 event_time_unix_ms
344 .saturating_sub(request.request_received_ms.unwrap_or(event_time_unix_ms))
345 });
346 let end_ms = request
347 .request_received_ms
348 .map(|start| start.saturating_add(total_ms))
349 .unwrap_or(event_time_unix_ms);
350 let start_ms = request
351 .request_received_ms
352 .unwrap_or_else(|| event_time_unix_ms.saturating_sub(total_ms));
353 (saturating_i64(start_ms), saturating_i64(end_ms))
354}
355
356pub(crate) fn saturating_i64(value: u64) -> i64 {
357 value.min(i64::MAX as u64) as i64
358}
359
360#[cfg(test)]
361mod tests {
362 use std::io::Write;
363
364 use tempfile::NamedTempFile;
365
366 use super::*;
367
368 #[test]
369 fn request_times_uses_event_time_when_total_duration_is_missing() {
370 let request = RequestTraceRequestMetrics {
371 request_id: "req".to_string(),
372 output_tokens: Some(1),
373 request_received_ms: Some(1_000),
374 total_time_ms: None,
375 replay: None,
376 };
377
378 assert_eq!(request_times(1_250, &request), (1_000, 1_250));
379 }
380
381 #[test]
382 fn loads_context_free_request_trace() {
383 let mut file = NamedTempFile::new().unwrap();
384 writeln!(
385 file,
386 r#"{{"schema":"dynamo.request.trace.v1","event_type":"request_end","event_time_unix_ms":1100,"request":{{"request_id":"req-1","request_received_ms":1000,"output_tokens":4,"replay":{{"trace_block_size":2,"input_length":3,"input_sequence_hashes":[11,22]}}}}}}"#
387 )
388 .unwrap();
389
390 let loaded = load_request_trace_records(&[file.path().to_path_buf()]).unwrap();
391 assert_eq!(loaded.requests.len(), 1);
392 assert!(loaded.requests[0].agent_context.is_none());
393 assert_eq!(loaded.requests[0].start_ms, 1_000);
394 assert_eq!(loaded.requests[0].end_ms, 1_100);
395 }
396
397 #[test]
398 fn rejects_unknown_trace_schema() {
399 let mut file = NamedTempFile::new().unwrap();
400 writeln!(
401 file,
402 r#"{{"schema":"dynamo.request.trace.v2","event_type":"request_end","event_time_unix_ms":1100,"request":{{"request_id":"req-1","request_received_ms":1000,"output_tokens":4,"replay":{{"trace_block_size":2,"input_length":2,"input_sequence_hashes":[11]}}}}}}"#
403 )
404 .unwrap();
405
406 let error = load_request_trace_records(&[file.path().to_path_buf()]).unwrap_err();
407 assert!(error.to_string().contains("failed to parse"));
408 assert!(format!("{error:#}").contains("unknown variant `dynamo.request.trace.v2`"));
409 }
410
411 #[test]
412 fn request_trace_requires_arrival_time_and_output_length() {
413 let mut file = NamedTempFile::new().unwrap();
414 writeln!(
415 file,
416 r#"{{"schema":"dynamo.request.trace.v1","event_type":"request_end","event_time_unix_ms":1100,"request":{{"request_id":"req-1","replay":{{"trace_block_size":2,"input_length":2,"input_sequence_hashes":[11]}}}}}}"#
417 )
418 .unwrap();
419
420 let error = load_request_trace_records(&[file.path().to_path_buf()]).unwrap_err();
421 assert!(format!("{error:#}").contains("request trace is missing request_received_ms"));
422 }
423
424 #[test]
425 fn rejects_duplicate_request_ids() {
426 let mut file = NamedTempFile::new().unwrap();
427 for _ in 0..2 {
428 writeln!(
429 file,
430 r#"{{"schema":"dynamo.request.trace.v1","event_type":"request_end","event_time_unix_ms":1100,"request":{{"request_id":"req-1","request_received_ms":1000,"output_tokens":4,"replay":{{"trace_block_size":2,"input_length":2,"input_sequence_hashes":[11]}}}}}}"#
431 )
432 .unwrap();
433 }
434
435 let error = load_request_trace_records(&[file.path().to_path_buf()]).unwrap_err();
436 assert!(error.to_string().contains("duplicate request_id req-1"));
437 }
438
439 #[test]
440 fn rejects_extra_replay_hashes() {
441 let mut file = NamedTempFile::new().unwrap();
442 writeln!(
443 file,
444 r#"{{"schema":"dynamo.request.trace.v1","event_type":"request_end","event_time_unix_ms":1100,"request":{{"request_id":"req-1","request_received_ms":1000,"output_tokens":4,"replay":{{"trace_block_size":2,"input_length":2,"input_sequence_hashes":[11,22]}}}}}}"#
445 )
446 .unwrap();
447
448 let error = load_request_trace_records(&[file.path().to_path_buf()]).unwrap_err();
449 assert!(format!("{error:#}").contains("requires exactly 1 replay hashes, got 2"));
450 }
451
452 #[test]
453 fn loads_agentic_request_trace_rows_and_tool_events() {
454 let mut file = NamedTempFile::new().unwrap();
455 writeln!(
456 file,
457 r#"{{"schema":"dynamo.request.trace.v1","event_type":"request_end","event_time_unix_ms":1100,"event_source":"dynamo","agent_context":{{"session_id":"root"}},"request":{{"request_id":"req-1","model":"test","request_received_ms":1000,"output_tokens":4,"replay":{{"trace_block_size":2,"input_length":3,"input_sequence_hashes":[11,22]}}}}}}"#
458 )
459 .unwrap();
460 writeln!(
461 file,
462 r#"{{"schema":"dynamo.request.trace.v1","event_type":"tool_end","event_time_unix_ms":1200,"event_source":"harness","agent_context":{{"session_id":"root"}},"tool":{{"tool_call_id":"tool-1","tool_class":"search","claude":{{"source_request_id":"req-1","consumer_request_id":"req-2","child_session_id":"child","execution_mode":"background"}},"started_at_unix_ms":1110,"ended_at_unix_ms":1200,"status":"succeeded","duration_ms":90}}}}"#
463 )
464 .unwrap();
465
466 let loaded = load_request_trace_records(&[file.path().to_path_buf()]).unwrap();
467 loaded.ensure_agentic_compatible().unwrap();
468 assert_eq!(loaded.requests.len(), 1);
469 assert_eq!(
470 loaded.requests[0]
471 .agent_context
472 .as_ref()
473 .expect("agent context")
474 .session_id,
475 "root"
476 );
477 assert_eq!(loaded.tools.len(), 1);
478 assert_eq!(loaded.tools[0].tool_call_id, "tool-1");
479 assert_eq!(loaded.tools[0].tool_class, "search");
480 let claude = loaded.tools[0].claude.as_ref().expect("Claude metadata");
481 assert_eq!(claude.source_request_id, "req-1");
482 assert_eq!(claude.consumer_request_id.as_deref(), Some("req-2"));
483 assert_eq!(claude.child_session_id.as_deref(), Some("child"));
484 assert_eq!(claude.execution_mode, "background");
485 }
486}