Skip to main content

aisimulate_core/replay/loadgen/
dynamo.rs

1// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Dynamo request-trace-v1 loading without a dependency on Dynamo crates.
5
6use std::collections::{HashMap, HashSet};
7use std::fs::File;
8use std::io::{BufRead, BufReader, Read};
9use std::path::PathBuf;
10
11use anyhow::{Context, Result, anyhow, bail, ensure};
12use flate2::read::MultiGzDecoder;
13use serde::Deserialize;
14
15use super::trace::assign_dependency_component_play_ids;
16use super::{
17    AGENTIC_MOONCAKE_SCHEMA, AGENTIC_MOONCAKE_VERSION, AgenticDependency,
18    AgenticDependencyRelation, AgenticDependencyTrigger, AgenticHashIdScope, AgenticMooncakeHeader,
19    AgenticMooncakeRow, AgenticSourceProvenance, AgenticTrace, MooncakeRow, Trace,
20};
21
22#[derive(Debug, Clone, PartialEq)]
23pub enum DynamoRequestTrace {
24    Standard(Trace),
25    Agentic(AgenticTrace),
26}
27
28#[derive(Debug, Clone, Deserialize)]
29struct Record {
30    schema: String,
31    event_type: String,
32    event_time_unix_ms: u64,
33    #[serde(default)]
34    agent_context: Option<AgentContext>,
35    #[serde(default)]
36    request: Option<RequestMetrics>,
37    #[serde(default)]
38    tool: Option<ToolMetrics>,
39}
40
41#[derive(Debug, Clone, Deserialize)]
42struct AgentContext {
43    session_id: String,
44    #[serde(default)]
45    parent_session_id: Option<String>,
46}
47
48#[derive(Debug, Clone, Deserialize)]
49struct RequestMetrics {
50    request_id: String,
51    #[serde(default)]
52    model: Option<String>,
53    #[serde(default)]
54    output_tokens: Option<u64>,
55    #[serde(default)]
56    request_received_ms: Option<u64>,
57    #[serde(default)]
58    total_time_ms: Option<f64>,
59    replay: ReplayMetrics,
60}
61
62#[derive(Debug, Clone, Deserialize)]
63struct ReplayMetrics {
64    trace_block_size: usize,
65    input_length: usize,
66    input_sequence_hashes: Vec<u64>,
67}
68
69#[derive(Debug, Clone, Deserialize)]
70struct ToolMetrics {
71    tool_call_id: String,
72    tool_class: String,
73    #[serde(default)]
74    claude: Option<ClaudeToolMetrics>,
75    #[serde(default)]
76    started_at_unix_ms: Option<u64>,
77    #[serde(default)]
78    ended_at_unix_ms: Option<u64>,
79    #[serde(default)]
80    duration_ms: Option<f64>,
81}
82
83#[derive(Debug, Clone, Deserialize)]
84struct ClaudeToolMetrics {
85    source_request_id: String,
86    #[serde(default)]
87    consumer_request_id: Option<String>,
88    #[serde(default)]
89    child_session_id: Option<String>,
90    execution_mode: String,
91}
92
93#[derive(Debug, Clone)]
94struct RequestEntry {
95    start_ms: i64,
96    end_ms: i64,
97    agent_context: Option<AgentContext>,
98    request: RequestMetrics,
99}
100
101#[derive(Debug, Clone)]
102struct ToolEntry {
103    session_id: String,
104    tool_call_id: String,
105    tool_class: String,
106    claude: Option<ClaudeToolMetrics>,
107}
108
109#[derive(Debug, Default)]
110struct LoadedEntries {
111    requests: Vec<RequestEntry>,
112    tools: Vec<ToolEntry>,
113}
114
115impl DynamoRequestTrace {
116    pub fn from_request_trace_files(
117        paths: &[PathBuf],
118        expected_block_size: Option<usize>,
119    ) -> Result<Self> {
120        ensure!(!paths.is_empty(), "Dynamo trace requires at least one path");
121        let loaded = load_entries(paths)?;
122        let mut entries = loaded.requests;
123        let contextual = entries
124            .iter()
125            .filter(|entry| entry.agent_context.is_some())
126            .count();
127        if contextual != 0 && contextual != entries.len() {
128            bail!("Dynamo request trace cannot mix requests with and without agent_context");
129        }
130        let block_size = entries[0].request.replay.trace_block_size;
131        ensure!(block_size > 0, "embedded trace block size must be positive");
132        if entries
133            .iter()
134            .any(|entry| entry.request.replay.trace_block_size != block_size)
135        {
136            bail!("mixed replay trace_block_size values are not supported");
137        }
138        if let Some(expected) = expected_block_size {
139            ensure!(
140                expected == block_size,
141                "trace_block_size {expected} does not match embedded Dynamo request trace block size {block_size}"
142            );
143        }
144        entries.sort_by(|left, right| {
145            (left.start_ms, left.end_ms, &left.request.request_id).cmp(&(
146                right.start_ms,
147                right.end_ms,
148                &right.request.request_id,
149            ))
150        });
151        if contextual == 0 {
152            lower_standard(entries, block_size).map(Self::Standard)
153        } else {
154            lower_agentic(entries, loaded.tools, block_size).map(Self::Agentic)
155        }
156    }
157}
158
159fn open_reader(path: &PathBuf) -> Result<Box<dyn BufRead>> {
160    let file = File::open(path).with_context(|| format!("failed to open {}", path.display()))?;
161    let reader: Box<dyn Read> = if path.extension().and_then(|value| value.to_str()) == Some("gz") {
162        Box::new(MultiGzDecoder::new(file))
163    } else {
164        Box::new(file)
165    };
166    Ok(Box::new(BufReader::new(reader)))
167}
168
169fn parse_record(line: &str) -> Result<Option<Record>> {
170    let value: serde_json::Value = serde_json::from_str(line)?;
171    let Some(object) = value.as_object() else {
172        return Ok(None);
173    };
174    let event = object.get("event").unwrap_or(&value);
175    let event_type = event
176        .get("event_type")
177        .and_then(serde_json::Value::as_str)
178        .or_else(|| object.get("event_type").and_then(serde_json::Value::as_str));
179    if matches!(event_type, Some("request_payload" | "tool_start")) {
180        return Ok(None);
181    }
182    ensure!(
183        matches!(event_type, Some("request_end" | "tool_end" | "tool_error")),
184        "request trace supports request_end, terminal tool events, request_payload, and tool_start; got {event_type:?}"
185    );
186    let record: Record = serde_json::from_value(event.clone())?;
187    ensure!(
188        record.schema == "dynamo.request.trace.v1",
189        "unsupported Dynamo request trace schema {:?}",
190        record.schema
191    );
192    ensure!(
193        Some(record.event_type.as_str()) == event_type,
194        "record event_type changed while decoding"
195    );
196    Ok(Some(record))
197}
198
199fn load_entries(paths: &[PathBuf]) -> Result<LoadedEntries> {
200    let mut loaded = LoadedEntries::default();
201    let mut request_ids = HashSet::new();
202    for path in paths {
203        for (line_index, line) in open_reader(path)?.lines().enumerate() {
204            let line = line
205                .with_context(|| format!("failed to read {}:{}", path.display(), line_index + 1))?;
206            if line.trim().is_empty() {
207                continue;
208            }
209            let Some(record) = parse_record(&line).with_context(|| {
210                format!("failed to parse {}:{}", path.display(), line_index + 1)
211            })?
212            else {
213                continue;
214            };
215            if record.event_type == "request_end" {
216                let request = record
217                    .request
218                    .context("request_end is missing request metrics")?;
219                ensure!(
220                    !request.request_id.trim().is_empty(),
221                    "request_id must be nonempty"
222                );
223                ensure!(
224                    request_ids.insert(request.request_id.clone()),
225                    "duplicate request_id {:?}",
226                    request.request_id
227                );
228                let (start_ms, end_ms) = request_times(record.event_time_unix_ms, &request)?;
229                loaded.requests.push(RequestEntry {
230                    start_ms,
231                    end_ms,
232                    agent_context: record.agent_context,
233                    request,
234                });
235            } else if let Some(tool) = tool_entry(record)? {
236                loaded.tools.push(tool);
237            }
238        }
239    }
240    ensure!(
241        !loaded.requests.is_empty(),
242        "Dynamo trace contains no request_end records"
243    );
244    Ok(loaded)
245}
246
247fn request_times(event_time_unix_ms: u64, request: &RequestMetrics) -> Result<(i64, i64)> {
248    let total_ms = request
249        .total_time_ms
250        .map(|value| {
251            ensure!(
252                value.is_finite() && value >= 0.0,
253                "request duration must be finite and nonnegative"
254            );
255            Ok(value.round() as u64)
256        })
257        .transpose()?;
258    let end_ms = match (request.request_received_ms, total_ms) {
259        (Some(start), Some(duration)) => start.saturating_add(duration),
260        _ => event_time_unix_ms,
261    };
262    let start_ms = request
263        .request_received_ms
264        .unwrap_or_else(|| event_time_unix_ms.saturating_sub(total_ms.unwrap_or(0)));
265    Ok((saturating_i64(start_ms), saturating_i64(end_ms)))
266}
267
268fn tool_entry(record: Record) -> Result<Option<ToolEntry>> {
269    let Some(context) = record.agent_context else {
270        return Ok(None);
271    };
272    let Some(tool) = record.tool else {
273        return Ok(None);
274    };
275    ensure!(
276        !context.session_id.trim().is_empty(),
277        "tool session_id must be nonempty"
278    );
279    ensure!(
280        !tool.tool_call_id.trim().is_empty(),
281        "tool_call_id must be nonempty"
282    );
283    ensure!(
284        !tool.tool_class.trim().is_empty(),
285        "tool_class must be nonempty"
286    );
287    if let Some(duration) = tool.duration_ms {
288        ensure!(
289            duration.is_finite() && duration >= 0.0,
290            "tool duration must be finite and nonnegative"
291        );
292    }
293    let end_ms = saturating_i64(tool.ended_at_unix_ms.unwrap_or(record.event_time_unix_ms));
294    let start_ms = tool
295        .started_at_unix_ms
296        .map(saturating_i64)
297        .or_else(|| {
298            tool.duration_ms
299                .map(|duration| end_ms.saturating_sub(duration.round() as i64))
300        })
301        .unwrap_or(end_ms);
302    ensure!(end_ms >= start_ms, "tool end time precedes start time");
303    Ok(Some(ToolEntry {
304        session_id: context.session_id,
305        tool_call_id: tool.tool_call_id,
306        tool_class: tool.tool_class,
307        claude: tool.claude,
308    }))
309}
310
311fn saturating_i64(value: u64) -> i64 {
312    value.min(i64::MAX as u64) as i64
313}
314
315fn lower_standard(entries: Vec<RequestEntry>, block_size: usize) -> Result<Trace> {
316    let first_start = entries
317        .iter()
318        .map(|entry| entry.start_ms)
319        .min()
320        .ok_or_else(|| anyhow!("Dynamo trace contains no requests"))?;
321    let rows = entries
322        .into_iter()
323        .map(|entry| -> Result<MooncakeRow> {
324            Ok(MooncakeRow {
325                request_id: Some(entry.request.request_id),
326                input_length: Some(entry.request.replay.input_length),
327                output_length: Some(
328                    usize::try_from(
329                        entry
330                            .request
331                            .output_tokens
332                            .context("missing output_tokens")?,
333                    )
334                    .context("output_tokens does not fit usize")?,
335                ),
336                hash_ids: Some(entry.request.replay.input_sequence_hashes),
337                timestamp: Some((entry.start_ms - first_start) as f64),
338                ..Default::default()
339            })
340        })
341        .collect::<Result<Vec<_>>>()?;
342    Trace::from_mooncake_rows(rows, block_size)
343}
344
345fn lower_agentic(
346    entries: Vec<RequestEntry>,
347    tools: Vec<ToolEntry>,
348    block_size: usize,
349) -> Result<AgenticTrace> {
350    let first_start = entries
351        .iter()
352        .map(|entry| entry.start_ms)
353        .min()
354        .ok_or_else(|| anyhow!("Dynamo trace contains no requests"))?;
355    let id_to_index = entries
356        .iter()
357        .enumerate()
358        .map(|(index, entry)| (entry.request.request_id.clone(), index))
359        .collect::<HashMap<_, _>>();
360    let mut by_session: HashMap<String, Vec<usize>> = HashMap::new();
361    let mut parent_by_session: HashMap<String, String> = HashMap::new();
362    for (index, entry) in entries.iter().enumerate() {
363        let context = entry
364            .agent_context
365            .as_ref()
366            .context("agentic request is missing agent_context")?;
367        ensure!(
368            !context.session_id.trim().is_empty(),
369            "session_id must be nonempty"
370        );
371        by_session
372            .entry(context.session_id.clone())
373            .or_default()
374            .push(index);
375        if let Some(parent) = context.parent_session_id.as_ref() {
376            match parent_by_session.get(&context.session_id) {
377                Some(existing) if existing != parent => bail!(
378                    "session {:?} has conflicting parent_session_id values {:?} and {:?}",
379                    context.session_id,
380                    existing,
381                    parent
382                ),
383                Some(_) => {}
384                None => {
385                    parent_by_session.insert(context.session_id.clone(), parent.clone());
386                }
387            }
388        }
389    }
390    for indices in by_session.values_mut() {
391        indices.sort_by_key(|index| {
392            let entry = &entries[*index];
393            (
394                entry.start_ms,
395                entry.end_ms,
396                entry.request.request_id.clone(),
397            )
398        });
399    }
400    let mut dependencies = vec![Vec::<AgenticDependency>::new(); entries.len()];
401    for indices in by_session.values() {
402        for pair in indices.windows(2) {
403            push_dependency(
404                &mut dependencies[pair[1]],
405                dependency_between(
406                    &entries,
407                    pair[0],
408                    pair[1],
409                    AgenticDependencyTrigger::Completion,
410                    AgenticDependencyRelation::Sequence,
411                ),
412            );
413        }
414    }
415
416    let mut explicit_tool_by_child: HashMap<String, &ToolEntry> = HashMap::new();
417    for tool in &tools {
418        let Some(claude) = tool.claude.as_ref() else {
419            continue;
420        };
421        ensure!(
422            matches!(claude.execution_mode.as_str(), "blocking" | "background"),
423            "tool {:?} ({}) has unsupported execution_mode {:?}",
424            tool.tool_call_id,
425            tool.tool_class,
426            claude.execution_mode
427        );
428        for request_id in [
429            Some(claude.source_request_id.as_str()),
430            claude.consumer_request_id.as_deref(),
431        ]
432        .into_iter()
433        .flatten()
434        {
435            let request_index = id_to_index.get(request_id).with_context(|| {
436                format!(
437                    "tool {:?} references unknown request_id {:?}",
438                    tool.tool_call_id, request_id
439                )
440            })?;
441            let request_session = &entries[*request_index]
442                .agent_context
443                .as_ref()
444                .expect("validated agent context")
445                .session_id;
446            ensure!(
447                request_session == &tool.session_id,
448                "tool {:?} request {:?} belongs to session {:?}, expected {:?}",
449                tool.tool_call_id,
450                request_id,
451                request_session,
452                tool.session_id
453            );
454        }
455        let Some(child_session) = claude.child_session_id.as_ref() else {
456            continue;
457        };
458        if !by_session.contains_key(child_session) {
459            continue;
460        }
461        ensure!(
462            explicit_tool_by_child
463                .insert(child_session.clone(), tool)
464                .is_none(),
465            "multiple tool events reference child session {:?}",
466            child_session
467        );
468    }
469
470    for (child_session, parent_session) in &parent_by_session {
471        let child_indices = by_session
472            .get(child_session)
473            .expect("child session must have requests");
474        let parent_indices = by_session.get(parent_session).with_context(|| {
475            format!(
476                "child session {:?} references unknown parent session {:?}",
477                child_session, parent_session
478            )
479        })?;
480        let first_child = child_indices[0];
481        let last_child = *child_indices
482            .iter()
483            .max_by_key(|index| {
484                let entry = &entries[**index];
485                (entry.end_ms, entry.start_ms, &entry.request.request_id)
486            })
487            .expect("child session must be nonempty");
488
489        if let Some(tool) = explicit_tool_by_child.get(child_session) {
490            let claude = tool.claude.as_ref().expect("explicit tool has metadata");
491            let parent_spawn = id_to_index[&claude.source_request_id];
492            ensure!(
493                parent_indices.contains(&parent_spawn),
494                "tool {:?} source request {:?} is not in parent session {:?}",
495                tool.tool_call_id,
496                claude.source_request_id,
497                parent_session
498            );
499            push_dependency(
500                &mut dependencies[first_child],
501                dependency_between(
502                    &entries,
503                    parent_spawn,
504                    first_child,
505                    AgenticDependencyTrigger::Dispatch,
506                    AgenticDependencyRelation::Spawn,
507                ),
508            );
509            if let Some(consumer) = claude.consumer_request_id.as_ref() {
510                let parent_join = id_to_index[consumer];
511                ensure!(
512                    parent_indices.contains(&parent_join),
513                    "tool {:?} consumer request {:?} is not in parent session {:?}",
514                    tool.tool_call_id,
515                    consumer,
516                    parent_session
517                );
518                push_dependency(
519                    &mut dependencies[parent_join],
520                    dependency_between(
521                        &entries,
522                        last_child,
523                        parent_join,
524                        AgenticDependencyTrigger::Completion,
525                        AgenticDependencyRelation::Join,
526                    ),
527                );
528            }
529            continue;
530        }
531
532        if let Some(parent_spawn) = parent_indices
533            .iter()
534            .copied()
535            .filter(|index| entries[*index].start_ms <= entries[first_child].start_ms)
536            .max_by_key(|index| entries[*index].start_ms)
537        {
538            push_dependency(
539                &mut dependencies[first_child],
540                dependency_between(
541                    &entries,
542                    parent_spawn,
543                    first_child,
544                    AgenticDependencyTrigger::Dispatch,
545                    AgenticDependencyRelation::Spawn,
546                ),
547            );
548        }
549        if let Some(parent_join) = parent_indices
550            .iter()
551            .copied()
552            .filter(|index| entries[*index].start_ms >= entries[last_child].end_ms)
553            .min_by_key(|index| entries[*index].start_ms)
554        {
555            push_dependency(
556                &mut dependencies[parent_join],
557                dependency_between(
558                    &entries,
559                    last_child,
560                    parent_join,
561                    AgenticDependencyTrigger::Completion,
562                    AgenticDependencyRelation::Join,
563                ),
564            );
565        }
566    }
567
568    let mut rows = entries
569        .iter()
570        .enumerate()
571        .map(|(index, entry)| -> Result<AgenticMooncakeRow> {
572            let context = entry
573                .agent_context
574                .as_ref()
575                .expect("validated agent context");
576            let request_dependencies = dependencies[index].clone();
577            Ok(AgenticMooncakeRow {
578                request_id: entry.request.request_id.clone(),
579                play_id: "dynamo-request-trace".to_string(),
580                session_id: context.session_id.clone(),
581                model: entry
582                    .request
583                    .model
584                    .clone()
585                    .unwrap_or_else(|| "unknown".to_string()),
586                input_length: Some(entry.request.replay.input_length),
587                output_length: Some(
588                    usize::try_from(
589                        entry
590                            .request
591                            .output_tokens
592                            .context("missing output_tokens")?,
593                    )
594                    .context("output_tokens does not fit usize")?,
595                ),
596                output_token_ids: None,
597                hash_ids: Some(entry.request.replay.input_sequence_hashes.clone()),
598                not_before_ms: if request_dependencies.is_empty() {
599                    (entry.start_ms - first_start) as f64
600                } else {
601                    0.0
602                },
603                priority: None,
604                strict_priority: None,
605                policy_class: None,
606                dependencies: request_dependencies,
607            })
608        })
609        .collect::<Result<Vec<_>>>()?;
610    assign_dependency_component_play_ids(&mut rows, "dynamo-play");
611    AgenticTrace::from_agentic_mooncake_rows(
612        AgenticMooncakeHeader {
613            schema: AGENTIC_MOONCAKE_SCHEMA.to_string(),
614            version: AGENTIC_MOONCAKE_VERSION,
615            block_size,
616            hash_id_scope: AgenticHashIdScope::Local,
617            source: AgenticSourceProvenance {
618                format: "dynamo.request.trace.v1".to_string(),
619                digest: format!("requests:{};tools:{}", rows.len(), tools.len()),
620            },
621        },
622        rows,
623    )
624}
625
626fn dependency_between(
627    entries: &[RequestEntry],
628    source: usize,
629    target: usize,
630    trigger: AgenticDependencyTrigger,
631    relation: AgenticDependencyRelation,
632) -> AgenticDependency {
633    let source_time = match trigger {
634        AgenticDependencyTrigger::Dispatch => entries[source].start_ms,
635        AgenticDependencyTrigger::Completion => entries[source].end_ms,
636    };
637    AgenticDependency {
638        request_id: entries[source].request.request_id.clone(),
639        trigger,
640        delay_ms: entries[target].start_ms.saturating_sub(source_time) as f64,
641        relation,
642    }
643}
644
645fn push_dependency(dependencies: &mut Vec<AgenticDependency>, dependency: AgenticDependency) {
646    if dependencies.iter().any(|existing| {
647        existing.request_id == dependency.request_id && existing.trigger == dependency.trigger
648    }) {
649        return;
650    }
651    dependencies.push(dependency);
652}
653
654#[cfg(test)]
655mod tests {
656    use std::io::Write;
657
658    use serde_json::json;
659    use tempfile::NamedTempFile;
660
661    use super::*;
662
663    fn request(id: &str, start_ms: u64, session_id: Option<&str>) -> serde_json::Value {
664        let mut value = json!({
665            "schema": "dynamo.request.trace.v1",
666            "event_type": "request_end",
667            "event_time_unix_ms": start_ms + 10,
668            "request": {
669                "request_id": id,
670                "output_tokens": 2,
671                "request_received_ms": start_ms,
672                "total_time_ms": 10,
673                "replay": {
674                    "trace_block_size": 4,
675                    "input_length": 4,
676                    "input_sequence_hashes": [11]
677                }
678            }
679        });
680        if let Some(session_id) = session_id {
681            value["agent_context"] = json!({"session_id": session_id});
682        }
683        value
684    }
685
686    fn child_request(
687        id: &str,
688        start_ms: u64,
689        session_id: &str,
690        parent_session_id: &str,
691    ) -> serde_json::Value {
692        let mut value = request(id, start_ms, Some(session_id));
693        value["agent_context"]["parent_session_id"] = json!(parent_session_id);
694        value
695    }
696
697    fn child_tool(
698        source_request_id: &str,
699        consumer_request_id: Option<&str>,
700        child_session_id: &str,
701        execution_mode: &str,
702    ) -> serde_json::Value {
703        json!({
704            "schema": "dynamo.request.trace.v1",
705            "event_type": "tool_end",
706            "event_time_unix_ms": 115,
707            "agent_context": {"session_id": "parent"},
708            "tool": {
709                "tool_call_id": "tool-1",
710                "tool_class": "agent",
711                "started_at_unix_ms": 110,
712                "ended_at_unix_ms": 115,
713                "claude": {
714                    "source_request_id": source_request_id,
715                    "consumer_request_id": consumer_request_id,
716                    "child_session_id": child_session_id,
717                    "execution_mode": execution_mode
718                }
719            }
720        })
721    }
722
723    fn dependencies<'a>(trace: &'a AgenticTrace, request_id: &str) -> &'a [AgenticDependency] {
724        trace
725            .nodes()
726            .iter()
727            .find(|node| node.request_id() == request_id)
728            .expect("request must exist")
729            .dependencies()
730    }
731
732    fn trace_file(rows: &[serde_json::Value]) -> NamedTempFile {
733        let mut file = NamedTempFile::new().unwrap();
734        for row in rows {
735            writeln!(file, "{}", serde_json::to_string(row).unwrap()).unwrap();
736        }
737        file
738    }
739
740    #[test]
741    fn loads_standard_multi_file_trace() {
742        let first = trace_file(&[request("a", 100, None)]);
743        let second = trace_file(&[request("b", 120, None)]);
744        let loaded = DynamoRequestTrace::from_request_trace_files(
745            &[first.path().to_path_buf(), second.path().to_path_buf()],
746            Some(4),
747        )
748        .unwrap();
749        let DynamoRequestTrace::Standard(trace) = loaded else {
750            panic!("expected standard trace");
751        };
752        assert_eq!(trace.sessions.len(), 2);
753        assert_eq!(trace.sessions[0].first_arrival_timestamp_ms, Some(0.0));
754        assert_eq!(trace.sessions[1].first_arrival_timestamp_ms, Some(20.0));
755    }
756
757    #[test]
758    fn loads_agentic_trace_as_dependency_graph() {
759        let file = trace_file(&[
760            request("a", 100, Some("session")),
761            request("b", 120, Some("session")),
762        ]);
763        let loaded =
764            DynamoRequestTrace::from_request_trace_files(&[file.path().to_path_buf()], Some(4))
765                .unwrap();
766        let DynamoRequestTrace::Agentic(trace) = loaded else {
767            panic!("expected agentic trace");
768        };
769        assert_eq!(trace.node_count(), 2);
770        assert_eq!(trace.play_count(), 1);
771    }
772
773    #[test]
774    fn missing_duration_uses_request_end_event_time() {
775        let mut first = request("a", 100, Some("session"));
776        first["event_time_unix_ms"] = json!(150);
777        first["request"]
778            .as_object_mut()
779            .unwrap()
780            .remove("total_time_ms");
781        let file = trace_file(&[first, request("b", 160, Some("session"))]);
782
783        let DynamoRequestTrace::Agentic(trace) =
784            DynamoRequestTrace::from_request_trace_files(&[file.path().to_path_buf()], Some(4))
785                .unwrap()
786        else {
787            panic!("expected agentic trace");
788        };
789
790        let dependency = &dependencies(&trace, "b")[0];
791        assert_eq!(dependency.request_id, "a");
792        assert_eq!(dependency.delay_ms, 10.0);
793    }
794
795    #[test]
796    fn blocking_child_preserves_parent_sequence_spawn_and_join() {
797        let file = trace_file(&[
798            request("parent-source", 100, Some("parent")),
799            child_tool(
800                "parent-source",
801                Some("parent-consumer"),
802                "child",
803                "blocking",
804            ),
805            child_request("child-first", 120, "child", "parent"),
806            child_request("child-last", 150, "child", "parent"),
807            request("parent-consumer", 200, Some("parent")),
808        ]);
809
810        let DynamoRequestTrace::Agentic(trace) =
811            DynamoRequestTrace::from_request_trace_files(&[file.path().to_path_buf()], Some(4))
812                .unwrap()
813        else {
814            panic!("expected agentic trace");
815        };
816
817        assert!(
818            dependencies(&trace, "child-first")
819                .iter()
820                .any(|dependency| {
821                    dependency.request_id == "parent-source"
822                        && dependency.trigger == AgenticDependencyTrigger::Dispatch
823                        && dependency.relation == AgenticDependencyRelation::Spawn
824                })
825        );
826        let consumer = dependencies(&trace, "parent-consumer");
827        assert!(consumer.iter().any(|dependency| {
828            dependency.request_id == "parent-source"
829                && dependency.relation == AgenticDependencyRelation::Sequence
830        }));
831        assert!(consumer.iter().any(|dependency| {
832            dependency.request_id == "child-last"
833                && dependency.trigger == AgenticDependencyTrigger::Completion
834                && dependency.relation == AgenticDependencyRelation::Join
835        }));
836    }
837
838    #[test]
839    fn background_child_launches_without_implicit_parent_join() {
840        let file = trace_file(&[
841            request("parent-source", 100, Some("parent")),
842            child_tool("parent-source", None, "child", "background"),
843            child_request("child", 120, "child", "parent"),
844            request("parent-next", 130, Some("parent")),
845        ]);
846
847        let DynamoRequestTrace::Agentic(trace) =
848            DynamoRequestTrace::from_request_trace_files(&[file.path().to_path_buf()], Some(4))
849                .unwrap()
850        else {
851            panic!("expected agentic trace");
852        };
853
854        assert!(dependencies(&trace, "child").iter().any(|dependency| {
855            dependency.request_id == "parent-source"
856                && dependency.trigger == AgenticDependencyTrigger::Dispatch
857                && dependency.relation == AgenticDependencyRelation::Spawn
858        }));
859        assert_eq!(dependencies(&trace, "parent-next").len(), 1);
860        assert_eq!(
861            dependencies(&trace, "parent-next")[0].request_id,
862            "parent-source"
863        );
864    }
865
866    #[test]
867    fn timestamp_fallback_infers_spawn_and_last_child_join() {
868        let file = trace_file(&[
869            request("parent-source", 100, Some("parent")),
870            child_request("child-first", 120, "child", "parent"),
871            child_request("child-last", 150, "child", "parent"),
872            request("parent-consumer", 200, Some("parent")),
873        ]);
874
875        let DynamoRequestTrace::Agentic(trace) =
876            DynamoRequestTrace::from_request_trace_files(&[file.path().to_path_buf()], Some(4))
877                .unwrap()
878        else {
879            panic!("expected agentic trace");
880        };
881
882        assert!(
883            dependencies(&trace, "child-first")
884                .iter()
885                .any(|dependency| {
886                    dependency.request_id == "parent-source"
887                        && dependency.relation == AgenticDependencyRelation::Spawn
888                })
889        );
890        assert!(
891            dependencies(&trace, "parent-consumer")
892                .iter()
893                .any(|dependency| {
894                    dependency.request_id == "child-last"
895                        && dependency.relation == AgenticDependencyRelation::Join
896                })
897        );
898    }
899
900    #[test]
901    fn independent_agent_sessions_become_independent_plays() {
902        let file = trace_file(&[
903            request("a", 100, Some("session-a")),
904            request("b", 120, Some("session-b")),
905        ]);
906        let loaded =
907            DynamoRequestTrace::from_request_trace_files(&[file.path().to_path_buf()], Some(4))
908                .unwrap();
909        let DynamoRequestTrace::Agentic(trace) = loaded else {
910            panic!("expected agentic trace");
911        };
912        assert_eq!(trace.node_count(), 2);
913        assert_eq!(trace.play_count(), 2);
914    }
915}