Skip to main content

dynamo_data_gen/request_trace/
load.rs

1// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Source-format records and JSONL/gz loader.
5//!
6//! Local subset of Dynamo request trace rows consumed by replay.
7
8use 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/// Claude-only evidence used to disambiguate an offline replay DAG.
87///
88/// A completed session reveals the future request that consumed a tool result;
89/// live ZMQ tool producers do not have that information. Missing metadata is
90/// expected and leaves agentic lowering on its timestamp-based fallback.
91#[derive(Debug, Clone, Deserialize)]
92pub(crate) struct ClaudeToolReplayMetrics {
93    /// Request that emitted the tool call.
94    pub(crate) source_request_id: String,
95    /// Later request that consumed the terminal result.
96    #[serde(default)]
97    pub(crate) consumer_request_id: Option<String>,
98    /// Child agent session launched by the tool, when present in the export.
99    #[serde(default)]
100    pub(crate) child_session_id: Option<String>,
101    /// Whether parent execution blocked or continued in the background.
102    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
170/// Records other than `request_end` / `tool_end` / `tool_error` are skipped.
171/// Errors if no `request_end` rows were found.
172pub 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}