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;
11use std::io::BufReader;
12use std::io::Read;
13use std::path::Path;
14use std::path::PathBuf;
15
16use anyhow::Context;
17use anyhow::Result;
18use anyhow::anyhow;
19use anyhow::bail;
20use flate2::read::MultiGzDecoder;
21use serde::de::IgnoredAny;
22use serde::{Deserialize, Serialize};
23
24#[derive(Debug, Clone, Deserialize)]
25pub(crate) struct RequestTraceRecord {
26    pub(crate) schema: TraceSchema,
27    pub(crate) event_type: String,
28    pub(crate) event_time_unix_ms: u64,
29    #[serde(default)]
30    pub(crate) agent_context: Option<AgentContextFields>,
31    #[serde(default)]
32    pub(crate) request: Option<RequestTraceRequestMetrics>,
33    #[serde(default)]
34    pub(crate) tool: Option<RequestTraceToolEventMetrics>,
35}
36
37#[derive(Debug, Clone, Copy, Deserialize, PartialEq, Eq)]
38pub(crate) enum TraceSchema {
39    #[serde(rename = "dynamo.request.trace.v1")]
40    RequestV1,
41}
42
43#[derive(Debug, Clone, Default, Deserialize)]
44pub(crate) struct AgentContextFields {
45    pub(crate) session_id: String,
46    #[serde(default)]
47    pub(crate) parent_session_id: Option<String>,
48}
49
50#[derive(Debug, Clone, Default, Deserialize)]
51pub(crate) struct RequestTraceRequestMetrics {
52    pub(crate) request_id: String,
53    #[serde(default)]
54    pub(crate) x_request_id: Option<String>,
55    #[serde(default)]
56    pub(crate) model: Option<String>,
57    #[serde(default)]
58    pub(crate) input_tokens: Option<u64>,
59    #[serde(default)]
60    pub(crate) output_tokens: Option<u64>,
61    #[serde(default)]
62    pub(crate) cached_tokens: Option<u64>,
63    #[serde(default)]
64    pub(crate) request_received_ms: Option<u64>,
65    #[serde(default)]
66    pub(crate) prefill_wait_time_ms: Option<f64>,
67    #[serde(default)]
68    pub(crate) prefill_time_ms: Option<f64>,
69    #[serde(default)]
70    pub(crate) ttft_ms: Option<f64>,
71    #[serde(default)]
72    pub(crate) total_time_ms: Option<f64>,
73    #[serde(default)]
74    pub(crate) avg_itl_ms: Option<f64>,
75    #[serde(default)]
76    pub(crate) kv_hit_rate: Option<f64>,
77    #[serde(default)]
78    pub(crate) kv_transfer_estimated_latency_ms: Option<f64>,
79    #[serde(default)]
80    pub(crate) queue_depth: Option<u64>,
81    #[serde(default)]
82    pub(crate) worker: Option<RequestTraceWorkerInfo>,
83    #[serde(default)]
84    pub(crate) replay: Option<RequestTraceReplayMetrics>,
85    #[serde(default)]
86    pub(crate) finish_reason_metadata: Option<RequestTraceFinishReasonMetadata>,
87}
88
89#[derive(Debug, Clone, Deserialize, Serialize)]
90pub(crate) struct RequestTraceWorkerInfo {
91    #[serde(default, skip_serializing_if = "Option::is_none")]
92    pub(crate) prefill_worker_id: Option<u64>,
93    #[serde(default, skip_serializing_if = "Option::is_none")]
94    pub(crate) prefill_dp_rank: Option<u32>,
95    #[serde(default, skip_serializing_if = "Option::is_none")]
96    pub(crate) decode_worker_id: Option<u64>,
97    #[serde(default, skip_serializing_if = "Option::is_none")]
98    pub(crate) decode_dp_rank: Option<u32>,
99}
100
101#[derive(Debug, Clone, Deserialize, Serialize)]
102pub(crate) struct RequestTraceFinishReasonMetadata {
103    #[serde(default, skip_serializing_if = "Option::is_none")]
104    pub(crate) finish_reason: Option<String>,
105    #[serde(default, skip_serializing_if = "Option::is_none")]
106    pub(crate) backend_finish_reason: Option<String>,
107    #[serde(default, skip_serializing_if = "Option::is_none")]
108    pub(crate) stop_reason: Option<serde_json::Value>,
109    #[serde(default, skip_serializing_if = "Vec::is_empty")]
110    pub(crate) tool_calls: Vec<RequestTraceToolCall>,
111}
112
113#[derive(Debug, Clone, Deserialize, Serialize)]
114pub(crate) struct RequestTraceToolCall {
115    pub(crate) choice_index: u32,
116    pub(crate) tool_call_index: u32,
117    #[serde(default, skip_serializing_if = "Option::is_none")]
118    pub(crate) id: Option<String>,
119    #[serde(default, skip_serializing_if = "Option::is_none")]
120    pub(crate) name: Option<String>,
121}
122
123#[derive(Debug, Clone, Deserialize)]
124pub(crate) struct RequestTraceReplayMetrics {
125    pub(crate) trace_block_size: usize,
126    pub(crate) input_length: usize,
127    pub(crate) input_sequence_hashes: Vec<u64>,
128}
129
130#[derive(Debug, Clone, Deserialize)]
131pub(crate) struct RequestTraceToolEventMetrics {
132    pub(crate) tool_call_id: String,
133    pub(crate) tool_class: String,
134    #[serde(default)]
135    pub(crate) claude: Option<ClaudeToolReplayMetrics>,
136    #[serde(default)]
137    pub(crate) started_at_unix_ms: Option<u64>,
138    #[serde(default)]
139    pub(crate) ended_at_unix_ms: Option<u64>,
140    #[serde(default)]
141    pub(crate) duration_ms: Option<f64>,
142    #[serde(default)]
143    pub(crate) status: Option<String>,
144    #[serde(default)]
145    pub(crate) output_bytes: Option<u64>,
146    #[serde(default)]
147    pub(crate) output_tokens: Option<u64>,
148    #[serde(default)]
149    pub(crate) error_type: Option<String>,
150}
151
152/// Claude-only evidence used to disambiguate an offline replay DAG.
153///
154/// A completed session reveals the future request that consumed a tool result;
155/// live ZMQ tool producers do not have that information. Missing metadata is
156/// expected and leaves agentic lowering on its timestamp-based fallback.
157#[derive(Debug, Clone, Deserialize)]
158pub(crate) struct ClaudeToolReplayMetrics {
159    /// Request that emitted the tool call.
160    pub(crate) source_request_id: String,
161    /// Later request that consumed the terminal result.
162    #[serde(default)]
163    pub(crate) consumer_request_id: Option<String>,
164    /// Child agent session launched by the tool, when present in the export.
165    #[serde(default)]
166    pub(crate) child_session_id: Option<String>,
167    /// Whether parent execution blocked or continued in the background.
168    pub(crate) execution_mode: String,
169}
170
171#[derive(Debug, Deserialize)]
172#[serde(untagged)]
173enum JsonLineEnvelope {
174    Object(TraceRecordEnvelope),
175    Other(IgnoredAny),
176}
177
178#[derive(Debug, Deserialize)]
179struct TraceRecordEnvelope {
180    #[serde(default)]
181    event_type: Option<String>,
182    #[serde(default)]
183    event: Option<TraceEventEnvelope>,
184}
185
186#[derive(Debug, Deserialize)]
187#[serde(untagged)]
188enum TraceEventEnvelope {
189    Object(TraceEventFields),
190    Other(IgnoredAny),
191}
192
193#[derive(Debug, Deserialize)]
194struct TraceEventFields {
195    #[serde(default)]
196    event_type: Option<String>,
197}
198
199#[derive(Debug, Deserialize)]
200struct WrappedTraceRecord {
201    event: RequestTraceRecord,
202}
203
204#[derive(Debug, Clone)]
205pub struct RequestEntry {
206    pub(crate) start_ms: i64,
207    pub(crate) end_ms: i64,
208    pub(crate) agent_context: Option<AgentContextFields>,
209    pub(crate) request: RequestTraceRequestMetrics,
210    pub(crate) replay: RequestTraceReplayMetrics,
211}
212
213#[derive(Debug, Clone)]
214pub struct ToolEntry {
215    pub(crate) session_id: String,
216    pub(crate) start_ms: i64,
217    pub(crate) end_ms: i64,
218    pub(crate) tool_call_id: String,
219    pub(crate) tool_class: String,
220    pub(crate) claude: Option<ClaudeToolReplayMetrics>,
221    pub(crate) status: String,
222    pub(crate) duration_ms: f64,
223    pub(crate) output_bytes: Option<u64>,
224    pub(crate) output_tokens: Option<u64>,
225    pub(crate) error_type: Option<String>,
226}
227
228#[derive(Debug, Default)]
229pub struct LoadedAgentTrace {
230    pub requests: Vec<RequestEntry>,
231    pub tools: Vec<ToolEntry>,
232}
233
234#[derive(Debug, Clone, Copy, PartialEq, Eq)]
235pub enum RequestTraceMode {
236    Standard,
237    Agentic,
238}
239
240impl LoadedAgentTrace {
241    pub fn mode(&self) -> Result<RequestTraceMode> {
242        if self.requests.is_empty() {
243            bail!("Dynamo request trace contains no requests");
244        }
245        let contextual_requests = self
246            .requests
247            .iter()
248            .filter(|request| request.agent_context.is_some())
249            .count();
250        match contextual_requests {
251            0 => Ok(RequestTraceMode::Standard),
252            count if count == self.requests.len() => Ok(RequestTraceMode::Agentic),
253            _ => bail!("Dynamo request trace cannot mix requests with and without agent_context"),
254        }
255    }
256
257    pub fn ensure_agentic_compatible(&self) -> Result<()> {
258        if self
259            .requests
260            .iter()
261            .any(|request| request.agent_context.is_none())
262        {
263            bail!("agentic lowering requires agent_context on every request");
264        }
265        Ok(())
266    }
267}
268
269/// `request_payload` records are skipped; replay consumes `request_end` and
270/// terminal tool events only. Errors if no `request_end` rows were found.
271pub fn load_request_trace_records(paths: &[PathBuf]) -> Result<LoadedAgentTrace> {
272    let mut loaded = LoadedAgentTrace::default();
273    let mut request_ids = HashSet::new();
274
275    for path in paths {
276        let reader = open_trace_reader(path)?;
277        for (line_index, line) in reader.lines().enumerate() {
278            let line = line
279                .with_context(|| format!("failed to read {}:{}", path.display(), line_index + 1))?;
280            if line.trim().is_empty() {
281                continue;
282            }
283            let Some(record) = parse_trace_record(&line).with_context(|| {
284                format!("failed to parse {}:{}", path.display(), line_index + 1)
285            })?
286            else {
287                continue;
288            };
289            let _schema = record.schema;
290            if record.event_type == "request_payload" {
291                continue;
292            }
293            if !matches!(
294                record.event_type.as_str(),
295                "request_end" | "tool_start" | "tool_end" | "tool_error"
296            ) {
297                bail!(
298                    "request trace schema only supports request_end/tool_* and request_payload events, got {} at {}:{}",
299                    record.event_type,
300                    path.display(),
301                    line_index + 1
302                );
303            }
304            if record.event_type == "request_end" {
305                let entry = request_entry(record).with_context(|| {
306                    format!(
307                        "invalid request_end at {}:{}",
308                        path.display(),
309                        line_index + 1
310                    )
311                })?;
312                if !request_ids.insert(entry.request.request_id.clone()) {
313                    bail!(
314                        "duplicate request_id {} at {}:{}",
315                        entry.request.request_id,
316                        path.display(),
317                        line_index + 1
318                    );
319                }
320                loaded.requests.push(entry);
321            } else if matches!(record.event_type.as_str(), "tool_end" | "tool_error") {
322                let terminal_event = record.event_type.clone();
323                if let Some(tool) = tool_entry(record, terminal_event) {
324                    loaded.tools.push(tool);
325                }
326            }
327        }
328    }
329
330    if loaded.requests.is_empty() {
331        bail!("no request_end records with replay fields found");
332    }
333
334    Ok(loaded)
335}
336
337fn open_trace_reader(path: &Path) -> Result<Box<dyn BufRead>> {
338    let file = File::open(path).with_context(|| format!("failed to open {}", path.display()))?;
339    let reader: Box<dyn Read> = if path.extension().and_then(|ext| ext.to_str()) == Some("gz") {
340        Box::new(MultiGzDecoder::new(file))
341    } else {
342        Box::new(file)
343    };
344    Ok(Box::new(BufReader::new(reader)))
345}
346
347fn parse_trace_record(line: &str) -> Result<Option<RequestTraceRecord>> {
348    let envelope = match serde_json::from_str::<JsonLineEnvelope>(line)? {
349        JsonLineEnvelope::Object(envelope) => envelope,
350        JsonLineEnvelope::Other(_) => return Ok(None),
351    };
352
353    match envelope.event {
354        Some(TraceEventEnvelope::Object(event)) => {
355            if event.event_type.as_deref() == Some("request_payload") {
356                return Ok(None);
357            }
358            Ok(Some(
359                serde_json::from_str::<WrappedTraceRecord>(line)?.event,
360            ))
361        }
362        Some(TraceEventEnvelope::Other(_)) => Ok(None),
363        None => {
364            if envelope.event_type.as_deref() == Some("request_payload") {
365                return Ok(None);
366            }
367            Ok(Some(serde_json::from_str(line)?))
368        }
369    }
370}
371
372fn request_entry(record: RequestTraceRecord) -> Result<RequestEntry> {
373    let request = record
374        .request
375        .ok_or_else(|| anyhow!("request_end record is missing request payload"))?;
376    let replay = request
377        .replay
378        .clone()
379        .ok_or_else(|| anyhow!("request payload is missing replay metrics"))?;
380    if replay.trace_block_size == 0 {
381        bail!("request replay trace_block_size must be greater than 0");
382    }
383    let expected_hashes = replay.input_length.div_ceil(replay.trace_block_size);
384    if replay.input_sequence_hashes.len() != expected_hashes {
385        bail!(
386            "input_length {} with trace_block_size {} requires exactly {} replay hashes, got {}",
387            replay.input_length,
388            replay.trace_block_size,
389            expected_hashes,
390            replay.input_sequence_hashes.len()
391        );
392    }
393    if request.request_received_ms.is_none() {
394        bail!("request trace is missing request_received_ms");
395    }
396    if request.output_tokens.is_none() {
397        bail!("request trace is missing output_tokens");
398    }
399
400    let (start_ms, end_ms) = request_times(record.event_time_unix_ms, &request);
401    Ok(RequestEntry {
402        start_ms,
403        end_ms,
404        agent_context: record.agent_context,
405        request,
406        replay,
407    })
408}
409
410fn tool_entry(record: RequestTraceRecord, terminal_event: String) -> Option<ToolEntry> {
411    let context = record.agent_context?;
412    let tool = record.tool?;
413    let end_ms = tool
414        .ended_at_unix_ms
415        .map(saturating_i64)
416        .unwrap_or_else(|| saturating_i64(record.event_time_unix_ms));
417    let start_ms = tool
418        .started_at_unix_ms
419        .map(saturating_i64)
420        .or_else(|| {
421            tool.duration_ms
422                .map(|duration_ms| end_ms.saturating_sub(duration_ms.max(0.0).round() as i64))
423        })
424        .unwrap_or(end_ms);
425    if end_ms < start_ms {
426        return None;
427    }
428    let duration_ms = tool
429        .duration_ms
430        .unwrap_or_else(|| (end_ms - start_ms).max(0) as f64);
431    let status = tool.status.unwrap_or_else(|| {
432        if terminal_event == "tool_error" {
433            "error".to_string()
434        } else {
435            "succeeded".to_string()
436        }
437    });
438    Some(ToolEntry {
439        session_id: context.session_id,
440        start_ms,
441        end_ms,
442        tool_call_id: tool.tool_call_id,
443        tool_class: tool.tool_class,
444        claude: tool.claude,
445        status,
446        duration_ms,
447        output_bytes: tool.output_bytes,
448        output_tokens: tool.output_tokens,
449        error_type: tool.error_type,
450    })
451}
452
453pub(crate) fn request_times(
454    event_time_unix_ms: u64,
455    request: &RequestTraceRequestMetrics,
456) -> (i64, i64) {
457    let total_ms = request
458        .total_time_ms
459        .map(|value| value.max(0.0).round() as u64)
460        .unwrap_or_else(|| {
461            event_time_unix_ms
462                .saturating_sub(request.request_received_ms.unwrap_or(event_time_unix_ms))
463        });
464    let end_ms = request
465        .request_received_ms
466        .map(|start| start.saturating_add(total_ms))
467        .unwrap_or(event_time_unix_ms);
468    let start_ms = request
469        .request_received_ms
470        .unwrap_or_else(|| event_time_unix_ms.saturating_sub(total_ms));
471    (saturating_i64(start_ms), saturating_i64(end_ms))
472}
473
474pub(crate) fn saturating_i64(value: u64) -> i64 {
475    value.min(i64::MAX as u64) as i64
476}
477
478#[cfg(test)]
479mod tests {
480    use std::io::Write;
481
482    use tempfile::NamedTempFile;
483
484    use super::*;
485
486    #[test]
487    fn request_times_uses_event_time_when_total_duration_is_missing() {
488        let request = RequestTraceRequestMetrics {
489            request_id: "req".to_string(),
490            output_tokens: Some(1),
491            request_received_ms: Some(1_000),
492            ..Default::default()
493        };
494
495        assert_eq!(request_times(1_250, &request), (1_000, 1_250));
496    }
497
498    #[test]
499    fn loads_context_free_request_trace() {
500        let mut file = NamedTempFile::new().unwrap();
501        writeln!(
502            file,
503            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]}}}}}}"#
504        )
505        .unwrap();
506
507        let loaded = load_request_trace_records(&[file.path().to_path_buf()]).unwrap();
508        assert_eq!(loaded.requests.len(), 1);
509        assert!(loaded.requests[0].agent_context.is_none());
510        assert_eq!(loaded.requests[0].start_ms, 1_000);
511        assert_eq!(loaded.requests[0].end_ms, 1_100);
512    }
513
514    #[test]
515    fn loads_wrapped_request_trace_record() {
516        let mut file = NamedTempFile::new().unwrap();
517        writeln!(
518            file,
519            r#"{{"timestamp":1,"event":{{"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]}}}}}}}}"#
520        )
521        .unwrap();
522
523        let loaded = load_request_trace_records(&[file.path().to_path_buf()]).unwrap();
524        assert_eq!(loaded.requests.len(), 1);
525        assert_eq!(loaded.requests[0].request.request_id, "req-1");
526    }
527
528    #[test]
529    fn skips_request_payload_records() {
530        let mut file = NamedTempFile::new().unwrap();
531        writeln!(
532            file,
533            r#"{{"schema":"dynamo.request.trace.v1","event_type":"request_payload","event_time_unix_ms":1050,"payload":{{"request_id":"req-1","endpoint":"openai.chat_completion","model":"test","request":{{"model":"test","messages":[{{"role":"user","content":"hi"}}]}},"payload_complete":true}}}}"#
534        )
535        .unwrap();
536        let large_payload = "x".repeat(4096);
537        writeln!(
538            file,
539            "{}",
540            serde_json::json!({
541                "timestamp": 1051,
542                "event": {
543                    "schema": "dynamo.request.trace.v1",
544                    "event_type": "request_payload",
545                    "event_time_unix_ms": 1051,
546                    "payload": {
547                        "request_id": "req-1",
548                        "endpoint": "openai.chat_completion",
549                        "model": "test",
550                        "request": {
551                            "model": "test",
552                            "messages": [{
553                                "role": "user",
554                                "content": large_payload.clone(),
555                            }],
556                        },
557                        "response": {
558                            "choices": [{
559                                "message": {
560                                    "role": "assistant",
561                                    "content": large_payload,
562                                },
563                            }],
564                        },
565                        "payload_complete": true,
566                    },
567                },
568            })
569        )
570        .unwrap();
571        writeln!(
572            file,
573            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]}}}}}}"#
574        )
575        .unwrap();
576
577        let loaded = load_request_trace_records(&[file.path().to_path_buf()]).unwrap();
578        assert_eq!(loaded.requests.len(), 1);
579        assert_eq!(loaded.requests[0].request.request_id, "req-1");
580    }
581
582    #[test]
583    fn rejects_unknown_trace_schema() {
584        let mut file = NamedTempFile::new().unwrap();
585        writeln!(
586            file,
587            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]}}}}}}"#
588        )
589        .unwrap();
590
591        let error = load_request_trace_records(&[file.path().to_path_buf()]).unwrap_err();
592        assert!(error.to_string().contains("failed to parse"));
593        assert!(format!("{error:#}").contains("unknown variant `dynamo.request.trace.v2`"));
594    }
595
596    #[test]
597    fn request_trace_requires_arrival_time_and_output_length() {
598        let mut file = NamedTempFile::new().unwrap();
599        writeln!(
600            file,
601            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]}}}}}}"#
602        )
603        .unwrap();
604
605        let error = load_request_trace_records(&[file.path().to_path_buf()]).unwrap_err();
606        assert!(format!("{error:#}").contains("request trace is missing request_received_ms"));
607    }
608
609    #[test]
610    fn rejects_duplicate_request_ids() {
611        let mut file = NamedTempFile::new().unwrap();
612        for _ in 0..2 {
613            writeln!(
614                file,
615                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]}}}}}}"#
616            )
617            .unwrap();
618        }
619
620        let error = load_request_trace_records(&[file.path().to_path_buf()]).unwrap_err();
621        assert!(error.to_string().contains("duplicate request_id req-1"));
622    }
623
624    #[test]
625    fn rejects_extra_replay_hashes() {
626        let mut file = NamedTempFile::new().unwrap();
627        writeln!(
628            file,
629            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]}}}}}}"#
630        )
631        .unwrap();
632
633        let error = load_request_trace_records(&[file.path().to_path_buf()]).unwrap_err();
634        assert!(format!("{error:#}").contains("requires exactly 1 replay hashes, got 2"));
635    }
636
637    #[test]
638    fn loads_agentic_request_trace_rows_and_tool_events() {
639        let mut file = NamedTempFile::new().unwrap();
640        writeln!(
641            file,
642            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]}}}}}}"#
643        )
644        .unwrap();
645        writeln!(
646            file,
647            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}}}}"#
648        )
649        .unwrap();
650
651        let loaded = load_request_trace_records(&[file.path().to_path_buf()]).unwrap();
652        loaded.ensure_agentic_compatible().unwrap();
653        assert_eq!(loaded.requests.len(), 1);
654        assert_eq!(
655            loaded.requests[0]
656                .agent_context
657                .as_ref()
658                .expect("agent context")
659                .session_id,
660            "root"
661        );
662        assert_eq!(loaded.tools.len(), 1);
663        assert_eq!(loaded.tools[0].tool_call_id, "tool-1");
664        assert_eq!(loaded.tools[0].tool_class, "search");
665        let claude = loaded.tools[0].claude.as_ref().expect("Claude metadata");
666        assert_eq!(claude.source_request_id, "req-1");
667        assert_eq!(claude.consumer_request_id.as_deref(), Some("req-2"));
668        assert_eq!(claude.child_session_id.as_deref(), Some("child"));
669        assert_eq!(claude.execution_mode, "background");
670    }
671}