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