1use 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}