Skip to main content

dynamo_data_gen/request_trace/
agentic.rs

1// SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
2// SPDX-License-Identifier: Apache-2.0
3
4//! Agentic lowering: infer the workflow DAG and attribute tool spans to the LLM
5//! row that consumed them.
6
7use std::collections::{HashMap, HashSet, VecDeque};
8
9use crate::{AgenticMooncakeRow, AgenticToolEvent, RollingHashIdMapper};
10use anyhow::{Context, Result, anyhow, bail};
11
12use super::load::{LoadedAgentTrace, RequestEntry, ToolEntry};
13
14/// Streams agentic Mooncake-compatible rows into the replay builder.
15///
16/// This is an in-memory compatibility layer; it does not write a Mooncake trace.
17pub fn lower_agentic_mooncake_rows<F>(mut loaded: LoadedAgentTrace, mut emit: F) -> Result<usize>
18where
19    F: FnMut(usize, AgenticMooncakeRow) -> Result<()>,
20{
21    loaded.ensure_agentic_compatible()?;
22    let global_start_ms = loaded
23        .requests
24        .iter()
25        .map(|request| request.start_ms)
26        .min()
27        .ok_or_else(|| anyhow!("no request records to convert"))?;
28    let trace_block_size = loaded.requests[0].replay.trace_block_size;
29    for request in &loaded.requests {
30        if request.replay.trace_block_size != trace_block_size {
31            bail!(
32                "mixed replay trace_block_size values are not supported: {} and {}",
33                trace_block_size,
34                request.replay.trace_block_size
35            );
36        }
37    }
38
39    loaded.requests.sort_by(|left, right| {
40        (left.start_ms, left.end_ms, &left.request.request_id).cmp(&(
41            right.start_ms,
42            right.end_ms,
43            &right.request.request_id,
44        ))
45    });
46
47    let mut id_to_index = HashMap::new();
48    for (idx, request) in loaded.requests.iter().enumerate() {
49        if id_to_index
50            .insert(request.request.request_id.clone(), idx)
51            .is_some()
52        {
53            bail!("duplicate request_id {}", request.request.request_id);
54        }
55    }
56
57    let mut session_to_indices: HashMap<String, Vec<usize>> = HashMap::new();
58    let mut parent_by_session: HashMap<String, String> = HashMap::new();
59    for (idx, request) in loaded.requests.iter().enumerate() {
60        let session_id = session_id_for(request);
61        session_to_indices
62            .entry(session_id.clone())
63            .or_default()
64            .push(idx);
65        if let Some(parent) = request
66            .agent_context
67            .as_ref()
68            .and_then(|context| context.parent_session_id.as_deref())
69            .map(str::to_string)
70        {
71            match parent_by_session.get(&session_id) {
72                Some(existing) if existing != &parent => {
73                    bail!(
74                        "session {} has conflicting parent_session_id values: {} and {}",
75                        session_id,
76                        existing,
77                        parent
78                    );
79                }
80                Some(_) => {}
81                None => {
82                    parent_by_session.insert(session_id, parent);
83                }
84            }
85        }
86    }
87    for indices in session_to_indices.values_mut() {
88        indices.sort_by_key(|idx| {
89            let request = &loaded.requests[*idx];
90            (
91                request.start_ms,
92                request.end_ms,
93                request.request.request_id.clone(),
94            )
95        });
96    }
97
98    let mut explicit_tool_by_child = HashMap::new();
99    let mut background_sessions = HashSet::new();
100    for tool in &loaded.tools {
101        let Some(claude) = tool.claude.as_ref() else {
102            continue;
103        };
104        if !matches!(claude.execution_mode.as_str(), "blocking" | "background") {
105            bail!(
106                "tool {} has unsupported execution_mode {}",
107                tool.tool_call_id,
108                claude.execution_mode
109            );
110        }
111        for request_id in [
112            Some(claude.source_request_id.as_str()),
113            claude.consumer_request_id.as_deref(),
114        ]
115        .into_iter()
116        .flatten()
117        {
118            let Some(request_idx) = id_to_index.get(request_id) else {
119                bail!(
120                    "tool {} references unknown request_id {}",
121                    tool.tool_call_id,
122                    request_id
123                );
124            };
125            if session_id_for(&loaded.requests[*request_idx]) != tool.session_id {
126                bail!(
127                    "tool {} request {} belongs to a different session",
128                    tool.tool_call_id,
129                    request_id
130                );
131            }
132        }
133        let Some(child_session_id) = claude.child_session_id.as_deref() else {
134            continue;
135        };
136        if !session_to_indices.contains_key(child_session_id) {
137            continue;
138        }
139        if explicit_tool_by_child
140            .insert(child_session_id.to_string(), tool)
141            .is_some()
142        {
143            bail!("multiple tool events reference child session {child_session_id}");
144        }
145        if claude.execution_mode == "background" {
146            background_sessions.insert(child_session_id.to_string());
147        }
148    }
149
150    let mut wait_for: Vec<Vec<String>> = vec![Vec::new(); loaded.requests.len()];
151    let mut branches: Vec<Vec<String>> = vec![Vec::new(); loaded.requests.len()];
152    let mut prefix_reset = vec![false; loaded.requests.len()];
153    let mut previous_request_start_ms = vec![None; loaded.requests.len()];
154
155    for indices in session_to_indices.values() {
156        for (pos, idx) in indices.iter().copied().enumerate() {
157            prefix_reset[idx] = pos == 0;
158            if pos > 0 {
159                let previous_request = &loaded.requests[indices[pos - 1]];
160                let previous = &previous_request.request.request_id;
161                push_unique(&mut wait_for[idx], previous.clone());
162                previous_request_start_ms[idx] = Some(previous_request.start_ms);
163            }
164        }
165    }
166
167    for (session_id, parent_id) in &parent_by_session {
168        let Some(child_indices) = session_to_indices.get(session_id) else {
169            continue;
170        };
171        let Some(parent_indices) = session_to_indices.get(parent_id) else {
172            continue;
173        };
174        let first_child_idx = child_indices[0];
175        let last_finishing_child_idx = *child_indices
176            .iter()
177            .max_by(|left, right| {
178                let left_request = &loaded.requests[**left];
179                let right_request = &loaded.requests[**right];
180                (
181                    left_request.end_ms,
182                    left_request.start_ms,
183                    &left_request.request.request_id,
184                )
185                    .cmp(&(
186                        right_request.end_ms,
187                        right_request.start_ms,
188                        &right_request.request.request_id,
189                    ))
190            })
191            .expect("child session is non-empty");
192        if let Some(tool) = explicit_tool_by_child.get(session_id) {
193            let claude = tool
194                .claude
195                .as_ref()
196                .expect("explicit child tool has Claude metadata");
197            let source_request_id = claude.source_request_id.as_str();
198            let parent_spawn_idx = id_to_index[source_request_id];
199            if !parent_indices.contains(&parent_spawn_idx) {
200                bail!(
201                    "tool {} source request {} is not in parent session {}",
202                    tool.tool_call_id,
203                    source_request_id,
204                    parent_id
205                );
206            }
207            let parent_request_id = loaded.requests[parent_spawn_idx].request.request_id.clone();
208            push_unique(&mut wait_for[first_child_idx], parent_request_id);
209            let child_request_id = loaded.requests[first_child_idx].request.request_id.clone();
210            push_unique(&mut branches[parent_spawn_idx], child_request_id);
211            if let Some(consumer_request_id) = claude.consumer_request_id.as_deref() {
212                let parent_join_idx = id_to_index[consumer_request_id];
213                if !parent_indices.contains(&parent_join_idx) {
214                    bail!(
215                        "tool {} consumer request {} is not in parent session {}",
216                        tool.tool_call_id,
217                        consumer_request_id,
218                        parent_id
219                    );
220                }
221                let child_request_id = loaded.requests[last_finishing_child_idx]
222                    .request
223                    .request_id
224                    .clone();
225                push_unique(&mut wait_for[parent_join_idx], child_request_id);
226            }
227            continue;
228        }
229
230        let child_start_ms = loaded.requests[first_child_idx].start_ms;
231        let child_end_ms = loaded.requests[last_finishing_child_idx].end_ms;
232        if let Some(parent_spawn_idx) =
233            latest_request_starting_before(&loaded.requests, parent_indices, child_start_ms)
234        {
235            let parent_request_id = loaded.requests[parent_spawn_idx].request.request_id.clone();
236            push_unique(&mut wait_for[first_child_idx], parent_request_id);
237            let child_request_id = loaded.requests[first_child_idx].request.request_id.clone();
238            push_unique(&mut branches[parent_spawn_idx], child_request_id);
239        }
240        if let Some(parent_join_idx) =
241            first_request_starting_after(&loaded.requests, parent_indices, child_end_ms)
242        {
243            let child_request_id = loaded.requests[last_finishing_child_idx]
244                .request
245                .request_id
246                .clone();
247            push_unique(&mut wait_for[parent_join_idx], child_request_id);
248        }
249    }
250    validate_dependency_dag(&loaded.requests, &wait_for, &id_to_index)?;
251
252    let mut tools_by_session: HashMap<String, Vec<ToolEntry>> = HashMap::new();
253    for tool in loaded.tools {
254        tools_by_session
255            .entry(tool.session_id.clone())
256            .or_default()
257            .push(tool);
258    }
259    for tools in tools_by_session.values_mut() {
260        tools.sort_by_key(|tool| (tool.start_ms, tool.end_ms));
261    }
262
263    let mut mapper = RollingHashIdMapper::new(trace_block_size);
264    for (idx, request) in loaded.requests.iter().enumerate() {
265        let hash_ids = mapper.ids_for_sequence_hashes(&request.replay.input_sequence_hashes);
266        let output_length = request.request.output_tokens.ok_or_else(|| {
267            anyhow!(
268                "request {} is missing output length",
269                request.request.request_id
270            )
271        })?;
272        let session_id = session_id_for(request);
273        let dep_end_ms = wait_for[idx]
274            .iter()
275            .filter_map(|dependency| id_to_index.get(dependency))
276            .map(|dep_idx| loaded.requests[*dep_idx].end_ms)
277            .max();
278        let (delay, tool_wait_ms, tool_events) = if let Some(dep_end_ms) = dep_end_ms {
279            let observed_gap_ms = request.start_ms.saturating_sub(dep_end_ms).max(0) as f64;
280            let tool_event_start_ms = previous_request_start_ms[idx].unwrap_or(dep_end_ms);
281            let (raw_tool_wait_ms, contributing) = tools_by_session
282                .get(&session_id)
283                .map(|tools| {
284                    collect_tools_in_window(
285                        tools,
286                        &request.request.request_id,
287                        tool_event_start_ms,
288                        dep_end_ms,
289                        request.start_ms,
290                    )
291                })
292                .unwrap_or_else(|| (0.0, Vec::new()));
293            let tool_wait_ms = raw_tool_wait_ms.min(observed_gap_ms);
294            let non_tool_wait_ms = (observed_gap_ms - tool_wait_ms).max(0.0);
295            let events = contributing
296                .into_iter()
297                .map(tool_entry_to_event)
298                .collect::<Vec<_>>();
299            (
300                Some(non_tool_wait_ms),
301                (tool_wait_ms > 0.0).then_some(tool_wait_ms),
302                events,
303            )
304        } else {
305            (None, None, Vec::new())
306        };
307
308        emit(
309            trace_block_size,
310            AgenticMooncakeRow {
311                request_id: request.request.request_id.clone(),
312                session_id: Some(session_id.clone()),
313                input_length: Some(request.replay.input_length),
314                output_length: Some(
315                    usize::try_from(output_length)
316                        .context("output length does not fit in usize")?,
317                ),
318                hash_ids: Some(hash_ids),
319                request_kind: Some(
320                    if background_sessions.contains(&session_id) {
321                        "background_agent"
322                    } else if parent_by_session.contains_key(&session_id) {
323                        "agent"
324                    } else {
325                        "foreground"
326                    }
327                    .to_string(),
328                ),
329                timestamp: Some((request.start_ms - global_start_ms) as f64),
330                delay,
331                wait_for: std::mem::take(&mut wait_for[idx]),
332                branches: std::mem::take(&mut branches[idx]),
333                prefix_reset: Some(prefix_reset[idx]),
334                tool_wait_ms,
335                tool_events,
336                ..Default::default()
337            },
338        )?;
339    }
340
341    Ok(trace_block_size)
342}
343
344fn session_id_for(request: &RequestEntry) -> String {
345    request
346        .agent_context
347        .as_ref()
348        .map(|context| context.session_id.clone())
349        .unwrap_or_else(|| request.request.request_id.clone())
350}
351
352fn latest_request_starting_before(
353    requests: &[RequestEntry],
354    indices: &[usize],
355    timestamp_ms: i64,
356) -> Option<usize> {
357    indices
358        .iter()
359        .copied()
360        .filter(|idx| requests[*idx].start_ms <= timestamp_ms)
361        .max_by_key(|idx| requests[*idx].start_ms)
362}
363
364fn first_request_starting_after(
365    requests: &[RequestEntry],
366    indices: &[usize],
367    timestamp_ms: i64,
368) -> Option<usize> {
369    indices
370        .iter()
371        .copied()
372        .filter(|idx| requests[*idx].start_ms >= timestamp_ms)
373        .min_by_key(|idx| requests[*idx].start_ms)
374}
375
376/// Return tools completed since the previous request started, while computing
377/// wait time only from their overlap with `[wait_start_ms, end_ms]`.
378fn collect_tools_in_window<'a>(
379    tools: &'a [ToolEntry],
380    request_id: &str,
381    event_start_ms: i64,
382    wait_start_ms: i64,
383    end_ms: i64,
384) -> (f64, Vec<&'a ToolEntry>) {
385    let mut contributing: Vec<&ToolEntry> = Vec::new();
386    let mut intervals = Vec::new();
387    for tool in tools {
388        let claude = tool.claude.as_ref();
389        if let Some(consumer_request_id) =
390            claude.and_then(|metadata| metadata.consumer_request_id.as_deref())
391        {
392            if consumer_request_id != request_id {
393                continue;
394            }
395        } else if claude.is_some_and(|metadata| metadata.execution_mode == "background")
396            || tool.end_ms <= event_start_ms
397            || tool.end_ms > end_ms
398        {
399            continue;
400        }
401        contributing.push(tool);
402        let clipped_start = tool.start_ms.max(wait_start_ms);
403        let clipped_end = tool.end_ms.min(end_ms);
404        if clipped_end > clipped_start {
405            intervals.push((clipped_start, clipped_end));
406        }
407    }
408    intervals.sort_unstable();
409
410    let mut total = 0_i64;
411    let mut current: Option<(i64, i64)> = None;
412    for (start, end) in intervals {
413        match current {
414            None => current = Some((start, end)),
415            Some((current_start, current_end)) if start <= current_end => {
416                current = Some((current_start, current_end.max(end)));
417            }
418            Some((current_start, current_end)) => {
419                total += current_end - current_start;
420                current = Some((start, end));
421            }
422        }
423    }
424    if let Some((current_start, current_end)) = current {
425        total += current_end - current_start;
426    }
427    (total as f64, contributing)
428}
429
430fn tool_entry_to_event(entry: &ToolEntry) -> AgenticToolEvent {
431    AgenticToolEvent {
432        tool_call_id: entry.tool_call_id.clone(),
433        tool_class: entry.tool_class.clone(),
434        started_at_unix_ms: entry.start_ms.max(0) as u64,
435        ended_at_unix_ms: entry.end_ms.max(0) as u64,
436        duration_ms: entry.duration_ms,
437        status: entry.status.clone(),
438        output_bytes: entry.output_bytes,
439        output_tokens: entry.output_tokens,
440        error_type: entry.error_type.clone(),
441    }
442}
443
444fn push_unique(values: &mut Vec<String>, value: String) {
445    if !values.iter().any(|existing| existing == &value) {
446        values.push(value);
447    }
448}
449
450fn validate_dependency_dag(
451    requests: &[RequestEntry],
452    wait_for: &[Vec<String>],
453    id_to_index: &HashMap<String, usize>,
454) -> Result<()> {
455    let mut indegree = wait_for.iter().map(Vec::len).collect::<Vec<_>>();
456    let mut dependents = vec![Vec::new(); requests.len()];
457    for (request_idx, dependencies) in wait_for.iter().enumerate() {
458        for dependency in dependencies {
459            let dependency_idx = id_to_index.get(dependency).ok_or_else(|| {
460                anyhow!(
461                    "request {} depends on unknown request {}",
462                    requests[request_idx].request.request_id,
463                    dependency
464                )
465            })?;
466            dependents[*dependency_idx].push(request_idx);
467        }
468    }
469
470    let mut ready = indegree
471        .iter()
472        .enumerate()
473        .filter_map(|(idx, count)| (*count == 0).then_some(idx))
474        .collect::<VecDeque<_>>();
475    let mut visited = 0;
476    while let Some(idx) = ready.pop_front() {
477        visited += 1;
478        for dependent in &dependents[idx] {
479            indegree[*dependent] -= 1;
480            if indegree[*dependent] == 0 {
481                ready.push_back(*dependent);
482            }
483        }
484    }
485    if visited != requests.len() {
486        bail!("agentic request dependencies contain a cycle");
487    }
488    Ok(())
489}
490
491#[cfg(test)]
492mod tests {
493    use super::*;
494    use crate::request_trace::load::{
495        AgentContextFields, ClaudeToolReplayMetrics, RequestEntry, RequestTraceReplayMetrics,
496        RequestTraceRequestMetrics, ToolEntry,
497    };
498
499    fn request(
500        request_id: &str,
501        start_ms: i64,
502        end_ms: i64,
503        sequence_hashes: Vec<u64>,
504    ) -> RequestEntry {
505        RequestEntry {
506            start_ms,
507            end_ms,
508            agent_context: None,
509            request: RequestTraceRequestMetrics {
510                request_id: request_id.to_string(),
511                output_tokens: Some(5),
512                request_received_ms: Some(start_ms as u64),
513                total_time_ms: Some((end_ms - start_ms) as f64),
514                ..Default::default()
515            },
516            replay: RequestTraceReplayMetrics {
517                trace_block_size: 2,
518                input_length: sequence_hashes.len() * 2,
519                input_sequence_hashes: sequence_hashes,
520            },
521        }
522    }
523
524    fn contextual_request(
525        request_id: &str,
526        session_id: &str,
527        parent_session_id: Option<&str>,
528        start_ms: i64,
529        end_ms: i64,
530        sequence_hashes: Vec<u64>,
531    ) -> RequestEntry {
532        let mut entry = request(request_id, start_ms, end_ms, sequence_hashes);
533        entry.agent_context = Some(AgentContextFields {
534            session_id: session_id.to_string(),
535            parent_session_id: parent_session_id.map(str::to_string),
536        });
537        entry
538    }
539
540    fn tool(
541        session_id: &str,
542        tool_call_id: &str,
543        tool_class: &str,
544        start_ms: i64,
545        end_ms: i64,
546    ) -> ToolEntry {
547        ToolEntry {
548            session_id: session_id.to_string(),
549            start_ms,
550            end_ms,
551            tool_call_id: tool_call_id.to_string(),
552            tool_class: tool_class.to_string(),
553            claude: None,
554            status: "succeeded".to_string(),
555            duration_ms: (end_ms - start_ms).max(0) as f64,
556            output_bytes: None,
557            output_tokens: None,
558            error_type: None,
559        }
560    }
561
562    fn lower_rows(loaded: LoadedAgentTrace) -> Result<Vec<AgenticMooncakeRow>> {
563        let mut rows = Vec::with_capacity(loaded.requests.len());
564        lower_agentic_mooncake_rows(loaded, |_, row| {
565            rows.push(row);
566            Ok(())
567        })?;
568        Ok(rows)
569    }
570
571    #[test]
572    fn agentic_lowering_builds_sequential_waits_and_tool_wait_components() {
573        let loaded = LoadedAgentTrace {
574            requests: vec![
575                contextual_request("r1", "root", None, 1_000, 1_100, vec![11]),
576                contextual_request("r2", "root", None, 1_300, 1_400, vec![11, 22]),
577            ],
578            tools: vec![tool("root", "call-1", "ls", 1_150, 1_250)],
579        };
580
581        let rows = lower_rows(loaded).unwrap();
582
583        assert_eq!(rows.len(), 2);
584        assert!(rows[0].wait_for.is_empty());
585        assert_eq!(rows[0].prefix_reset, Some(true));
586        assert!(rows[0].tool_events.is_empty());
587        assert_eq!(rows[1].wait_for, vec!["r1"]);
588        assert_eq!(rows[1].delay, Some(100.0));
589        assert_eq!(rows[1].tool_wait_ms, Some(100.0));
590        assert_eq!(rows[1].dependency_delay_ms(), 200.0);
591        assert_eq!(rows[1].tool_events.len(), 1);
592        assert_eq!(rows[1].tool_events[0].tool_class, "ls");
593        assert_eq!(rows[1].tool_events[0].tool_call_id, "call-1");
594        assert_eq!(rows[1].tool_events[0].started_at_unix_ms, 1_150);
595        assert_eq!(rows[1].tool_events[0].ended_at_unix_ms, 1_250);
596    }
597
598    #[test]
599    fn agentic_lowering_attaches_parallel_tool_events_with_union_wait() {
600        let loaded = LoadedAgentTrace {
601            requests: vec![
602                contextual_request("r1", "root", None, 1_000, 1_100, vec![11]),
603                contextual_request("r2", "root", None, 1_400, 1_500, vec![11, 22]),
604            ],
605            // Two tools that overlap heavily: union is 200ms (1_100..1_300),
606            // naive sum would be 350ms.
607            tools: vec![
608                tool("root", "call-1", "read", 1_100, 1_300),
609                tool("root", "call-2", "read", 1_150, 1_250),
610                tool("root", "call-3", "find", 1_200, 1_250),
611            ],
612        };
613
614        let rows = lower_rows(loaded).unwrap();
615
616        assert_eq!(rows[1].tool_wait_ms, Some(200.0));
617        assert_eq!(rows[1].tool_events.len(), 3);
618        let classes: Vec<_> = rows[1]
619            .tool_events
620            .iter()
621            .map(|event| event.tool_class.as_str())
622            .collect();
623        assert!(classes.contains(&"read"));
624        assert!(classes.contains(&"find"));
625    }
626
627    #[test]
628    fn agentic_lowering_adds_subagent_launch_and_join_dependencies() {
629        let loaded = LoadedAgentTrace {
630            requests: vec![
631                contextual_request("parent-1", "root", None, 1_000, 1_100, vec![11]),
632                contextual_request("child-1", "child", Some("root"), 1_200, 1_300, vec![33]),
633                contextual_request("parent-2", "root", None, 1_500, 1_600, vec![11, 22]),
634            ],
635            tools: Vec::new(),
636        };
637
638        let rows = lower_rows(loaded).unwrap();
639        let by_id = rows
640            .iter()
641            .map(|row| (row.request_id.as_str(), row))
642            .collect::<HashMap<_, _>>();
643
644        assert_eq!(by_id["child-1"].wait_for, vec!["parent-1"]);
645        assert_eq!(by_id["parent-1"].branches, vec!["child-1"]);
646        assert_eq!(by_id["parent-2"].wait_for, vec!["parent-1", "child-1"]);
647        assert_eq!(by_id["parent-2"].delay, Some(200.0));
648    }
649
650    #[test]
651    fn explicit_background_agent_causality_allows_parent_work_until_join() {
652        let mut agent_tool = tool("root", "agent-call", "Agent", 1_100, 1_800);
653        agent_tool.claude = Some(ClaudeToolReplayMetrics {
654            source_request_id: "parent-1".to_string(),
655            consumer_request_id: Some("parent-3".to_string()),
656            child_session_id: Some("child".to_string()),
657            execution_mode: "background".to_string(),
658        });
659        let loaded = LoadedAgentTrace {
660            requests: vec![
661                contextual_request("parent-1", "root", None, 1_000, 1_100, vec![11]),
662                contextual_request("child-1", "child", Some("root"), 1_200, 1_700, vec![33]),
663                contextual_request("parent-2", "root", None, 1_300, 1_400, vec![11, 22]),
664                contextual_request("parent-3", "root", None, 1_850, 1_950, vec![11, 22, 44]),
665            ],
666            tools: vec![agent_tool],
667        };
668
669        let rows = lower_rows(loaded).unwrap();
670        let by_id = rows
671            .iter()
672            .map(|row| (row.request_id.as_str(), row))
673            .collect::<HashMap<_, _>>();
674
675        assert_eq!(by_id["parent-1"].branches, vec!["child-1"]);
676        assert_eq!(by_id["child-1"].wait_for, vec!["parent-1"]);
677        assert_eq!(
678            by_id["child-1"].request_kind.as_deref(),
679            Some("background_agent")
680        );
681        assert_eq!(by_id["parent-2"].wait_for, vec!["parent-1"]);
682        assert_eq!(by_id["parent-3"].wait_for, vec!["parent-2", "child-1"]);
683        assert_eq!(by_id["parent-3"].tool_wait_ms, Some(100.0));
684        assert_eq!(by_id["parent-3"].delay, Some(50.0));
685        assert_eq!(by_id["parent-3"].tool_events.len(), 1);
686    }
687
688    #[test]
689    fn explicit_causality_rejects_cycles() {
690        let mut agent_tool = tool("root", "agent-call", "Agent", 1_100, 1_200);
691        agent_tool.claude = Some(ClaudeToolReplayMetrics {
692            source_request_id: "parent-2".to_string(),
693            consumer_request_id: Some("parent-1".to_string()),
694            child_session_id: Some("child".to_string()),
695            execution_mode: "background".to_string(),
696        });
697        let loaded = LoadedAgentTrace {
698            requests: vec![
699                contextual_request("parent-1", "root", None, 1_000, 1_100, vec![11]),
700                contextual_request("child-1", "child", Some("root"), 1_200, 1_300, vec![33]),
701                contextual_request("parent-2", "root", None, 1_400, 1_500, vec![11, 22]),
702            ],
703            tools: vec![agent_tool],
704        };
705
706        let err = lower_rows(loaded).unwrap_err();
707        assert!(err.to_string().contains("dependencies contain a cycle"));
708    }
709
710    #[test]
711    fn missing_child_trace_replays_as_external_background_tool() {
712        let mut agent_tool = tool("root", "agent-call", "Agent", 1_100, 1_250);
713        agent_tool.claude = Some(ClaudeToolReplayMetrics {
714            source_request_id: "parent-1".to_string(),
715            consumer_request_id: Some("parent-2".to_string()),
716            child_session_id: Some("missing-child".to_string()),
717            execution_mode: "background".to_string(),
718        });
719        let rows = lower_rows(LoadedAgentTrace {
720            requests: vec![
721                contextual_request("parent-1", "root", None, 1_000, 1_100, vec![11]),
722                contextual_request("parent-2", "root", None, 1_300, 1_400, vec![11, 22]),
723            ],
724            tools: vec![agent_tool],
725        })
726        .unwrap();
727
728        assert!(rows[0].branches.is_empty());
729        assert_eq!(rows[1].wait_for, vec!["parent-1"]);
730        assert_eq!(rows[1].tool_wait_ms, Some(150.0));
731        assert_eq!(rows[1].delay, Some(50.0));
732        assert_eq!(rows[1].tool_events.len(), 1);
733    }
734
735    #[test]
736    fn agentic_lowering_rejects_conflicting_session_parents() {
737        let loaded = LoadedAgentTrace {
738            requests: vec![
739                contextual_request("child-1", "child", Some("root-a"), 1_000, 1_100, vec![11]),
740                contextual_request("child-2", "child", Some("root-b"), 1_200, 1_300, vec![22]),
741            ],
742            tools: Vec::new(),
743        };
744
745        let err = lower_rows(loaded).unwrap_err();
746        assert!(err.to_string().contains("conflicting parent_session_id"));
747    }
748
749    #[test]
750    fn agentic_lowering_joins_on_last_finishing_child_request() {
751        let loaded = LoadedAgentTrace {
752            requests: vec![
753                contextual_request("parent-1", "root", None, 1_000, 1_100, vec![11]),
754                contextual_request("child-slow", "child", Some("root"), 1_200, 1_900, vec![33]),
755                contextual_request("child-fast", "child", Some("root"), 1_300, 1_400, vec![44]),
756                contextual_request("parent-2", "root", None, 1_500, 1_600, vec![11, 22]),
757                contextual_request("parent-3", "root", None, 2_000, 2_100, vec![11, 22, 33]),
758            ],
759            tools: Vec::new(),
760        };
761
762        let rows = lower_rows(loaded).unwrap();
763        let by_id = rows
764            .iter()
765            .map(|row| (row.request_id.as_str(), row))
766            .collect::<HashMap<_, _>>();
767
768        assert!(!by_id["parent-2"].wait_for.contains(&"child-fast".into()));
769        assert!(by_id["parent-3"].wait_for.contains(&"child-slow".into()));
770    }
771}