1use std::collections::HashMap;
2use std::path::Path;
3
4use crate::event;
5use crate::event_log::reader::{parse_json_lines, read_event_envelopes};
6use crate::message::{Message, MessagePart};
7use crate::nodegraph;
8use crate::provider;
9use crate::session::SessionOpenError;
10use serde_json;
11
12#[derive(Debug, Clone)]
13pub enum TranscriptEntry {
14 Message {
15 message: Message,
16 flow_run_id: Option<String>,
17 },
18 CompactionSummary {
19 range_start: usize,
20 range_end: usize,
21 compacted_count: usize,
22 before_tokens: u64,
23 after_tokens: u64,
24 summary: String,
25 ts: Option<chrono::DateTime<chrono::Utc>>,
26 },
27 DiffPreview {
28 title: String,
29 old_content: Option<String>,
30 new_content: Option<String>,
31 unified_diff: Option<String>,
32 },
33 FlowGraph {
34 run_id: String,
35 flow_name: String,
36 graph: nodegraph::FlowGraph,
37 ts: Option<chrono::DateTime<chrono::Utc>>,
38 },
39 FlowStart {
40 run_id: String,
41 flow_name: String,
42 parent_run_id: Option<String>,
43 parent_node_id: Option<String>,
44 spawned: bool,
45 ts: Option<chrono::DateTime<chrono::Utc>>,
46 },
47 FlowNodeStart {
48 run_id: String,
49 node_id: String,
50 kind: nodegraph::NodeKind,
51 label: String,
52 parent_node_id: Option<String>,
53 ts: Option<chrono::DateTime<chrono::Utc>>,
54 },
55 FlowNodeEnd {
56 run_id: String,
57 node_id: String,
58 status: event::FlowNodeStatus,
59 output_preview: Option<String>,
60 ts: Option<chrono::DateTime<chrono::Utc>>,
61 },
62 ToolNode {
63 run_id: String,
64 parent_node_id: String,
65 tool_use_id: String,
66 tool_name: String,
67 args_preview: String,
68 ts: Option<chrono::DateTime<chrono::Utc>>,
69 },
70 FlowDone {
71 run_id: String,
72 ok: bool,
73 cancelled: bool,
74 ts: Option<chrono::DateTime<chrono::Utc>>,
75 },
76 LlmCall {
77 model: String,
78 usage: provider::TokenUsage,
79 wallclock_ms: u64,
80 ttft_ms: Option<u64>,
81 tokens_per_second: Option<f64>,
82 run_id: Option<event::FlowRunId>,
83 node_id: Option<String>,
84 ts: Option<chrono::DateTime<chrono::Utc>>,
85 },
86 TerminalFinalState {
87 handle: String,
88 screen: crate::tools::term::TerminalScreen,
89 },
90 MermaidDiagram {
91 source: String,
92 },
93}
94
95pub fn replay_messages_from(path: &Path) -> Result<Vec<Message>, SessionOpenError> {
96 Ok(replay_messages_with_seq(path)?
97 .into_iter()
98 .map(|(_, msg)| msg)
99 .collect())
100}
101
102pub fn replay_messages_with_seq(path: &Path) -> Result<Vec<(u64, Message)>, SessionOpenError> {
103 let envelopes = read_event_envelopes(path)?;
104 Ok(envelopes.as_slice().to_messages_with_seq())
105}
106
107pub fn replay_all_messages_with_seq(path: &Path) -> Result<Vec<(u64, Message)>, SessionOpenError> {
108 let envelopes = read_event_envelopes(path)?;
109 let spawned_flow_ids = spawned_flow_ids(&envelopes);
110 Ok(envelopes
111 .iter()
112 .filter_map(|env| match &env.event {
113 crate::event::Event::UserMsg {
114 message,
115 flow_run_id,
116 ..
117 }
118 | crate::event::Event::AssistantMsg {
119 message,
120 flow_run_id,
121 ..
122 }
123 | crate::event::Event::ToolResultMsg {
124 message,
125 flow_run_id,
126 ..
127 } if message_belongs_to_root(flow_run_id.as_ref(), &spawned_flow_ids) => {
128 Some((env.seq, message.clone()))
129 }
130 crate::event::Event::SystemMsg { message, .. } => Some((env.seq, message.clone())),
131 _ => None,
132 })
133 .collect())
134}
135
136#[derive(Debug, Clone)]
137pub struct AttachmentPatch {
138 part_index: usize,
139 file_basename: String,
140 reason: String,
141}
142
143pub fn parse_ts(v: &serde_json::Value) -> Option<chrono::DateTime<chrono::Utc>> {
144 v.get("ts")?
145 .as_str()
146 .and_then(|s| chrono::DateTime::parse_from_rfc3339(s).ok())
147 .map(|dt| dt.with_timezone(&chrono::Utc))
148}
149
150pub fn parse_context_compact_event(v: &serde_json::Value) -> Option<CompactReplayEvent> {
151 if v["type"].as_str() != Some("context_compact") {
152 return None;
153 }
154 Some(CompactReplayEvent {
155 range_start: v["compacted_range_start"].as_u64().unwrap_or(0) as usize,
156 range_end: v["compacted_range_end"].as_u64().unwrap_or(0) as usize,
157 replacement_msg_seq: v["replacement_msg_seq"].as_u64(),
158 })
159}
160
161#[derive(Debug, Clone)]
162pub struct CompactReplayEvent {
163 range_start: usize,
164 range_end: usize,
165 replacement_msg_seq: Option<u64>,
166}
167
168pub fn collect_attachment_patches(
169 values: &[serde_json::Value],
170) -> HashMap<u64, Vec<AttachmentPatch>> {
171 let mut map: HashMap<u64, Vec<AttachmentPatch>> = HashMap::new();
172 for v in values {
173 if v["type"].as_str() == Some("attachment_degraded") {
174 let Some(msg_seq) = v["message_seq"].as_u64() else {
175 continue;
176 };
177 let Some(part_index) = v["part_index"].as_u64() else {
178 continue;
179 };
180 let file_basename = v["file_basename"].as_str().unwrap_or("").to_string();
181 let reason = v["reason"].as_str().unwrap_or("degraded").to_string();
182 map.entry(msg_seq).or_default().push(AttachmentPatch {
183 part_index: part_index as usize,
184 file_basename,
185 reason,
186 });
187 }
188 }
189 map
190}
191
192pub fn apply_attachment_patches(msg: &mut Message, patches: &[AttachmentPatch]) {
193 for p in patches {
194 if let Some(part) = msg.parts.get_mut(p.part_index) {
195 *part = MessagePart::Text {
196 text: format!(
197 "[attachment unavailable: {} — {}]",
198 p.file_basename, p.reason
199 ),
200 };
201 }
202 }
203}
204
205pub fn replay_transcript_from(path: &Path) -> Result<Vec<TranscriptEntry>, SessionOpenError> {
206 let text = match std::fs::read_to_string(path) {
207 Ok(t) => t,
208 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
209 Err(e) => {
210 return Err(SessionOpenError::Replay {
211 path: path.to_path_buf(),
212 source: e,
213 });
214 }
215 };
216 let values = parse_json_lines(&text);
217 let mut flow_parents = std::collections::HashMap::new();
218 let mut spawned_flow_ids = std::collections::HashSet::new();
219 for value in &values {
220 if value["type"].as_str() != Some("flow_start") {
221 continue;
222 }
223 let Some(raw_run_id) = value["run_id"].as_str() else {
224 continue;
225 };
226 let Ok(run_id) = uuid::Uuid::parse_str(raw_run_id) else {
227 continue;
228 };
229 let run_id = crate::event::FlowRunId(run_id);
230 let parent = value["parent_run_id"]
231 .as_str()
232 .and_then(|raw| uuid::Uuid::parse_str(raw).ok())
233 .map(crate::event::FlowRunId);
234 flow_parents.insert(run_id.clone(), parent);
235 if value["spawned"].as_bool().unwrap_or(false) {
236 spawned_flow_ids.insert(run_id);
237 }
238 }
239 loop {
240 let descendants: Vec<_> = flow_parents
241 .iter()
242 .filter_map(|(run_id, parent)| {
243 (!spawned_flow_ids.contains(run_id)
244 && parent
245 .as_ref()
246 .is_some_and(|parent| spawned_flow_ids.contains(parent)))
247 .then_some(run_id.clone())
248 })
249 .collect();
250 if descendants.is_empty() {
251 break;
252 }
253 spawned_flow_ids.extend(descendants);
254 }
255 let known_flow_ids = flow_parents
256 .into_keys()
257 .collect::<std::collections::HashSet<_>>();
258 let patches = collect_attachment_patches(&values);
259 let mut out = Vec::new();
260 let mut msg_indices: Vec<usize> = Vec::new();
261 let mut msg_seqs: Vec<u64> = Vec::new();
262 for v in &values {
263 let ty = v["type"].as_str().unwrap_or("");
264 match ty {
265 "user_msg" | "assistant_msg" | "tool_result_msg" | "system_msg" => {
266 if let Some(m) = v.get("message")
267 && let Ok(mut msg) = serde_json::from_value::<Message>(m.clone())
268 {
269 let seq = v["seq"].as_u64().unwrap_or(0);
270 if let Some(ps) = patches.get(&seq) {
271 apply_attachment_patches(&mut msg, ps);
272 }
273 let flow_run_id = v["flow_run_id"].as_str().and_then(|raw| {
274 let run_id = uuid::Uuid::parse_str(raw).ok()?;
275 let run_id = crate::event::FlowRunId(run_id);
276 if known_flow_ids.contains(&run_id) && !spawned_flow_ids.contains(&run_id) {
277 None
278 } else {
279 Some(raw.to_string())
280 }
281 });
282 msg_indices.push(out.len());
283 msg_seqs.push(seq);
284 out.push(TranscriptEntry::Message {
285 message: msg,
286 flow_run_id,
287 });
288 }
289 }
290 "context_compact" => {
291 let Some(event) = parse_context_compact_event(v) else {
292 continue;
293 };
294 if event.range_start > event.range_end || event.range_end >= msg_indices.len() {
295 continue;
296 }
297 let Some(replacement_seq) = event.replacement_msg_seq else {
298 continue;
299 };
300 let Some(replacement_pos) = msg_seqs.iter().position(|seq| *seq == replacement_seq)
301 else {
302 continue;
303 };
304 let replacement_out_idx = msg_indices[replacement_pos];
305 let replacement_entry = out.remove(replacement_out_idx);
306 let removed_out_start = msg_indices[event.range_start];
307 let removed_count = event.range_end - event.range_start + 1;
308 for _ in 0..removed_count {
309 out.remove(removed_out_start);
310 }
311 msg_indices.drain(event.range_start..=event.range_end);
312 msg_seqs.drain(event.range_start..=event.range_end);
313 out.insert(removed_out_start, replacement_entry);
314 msg_indices.insert(event.range_start, removed_out_start);
315 msg_seqs.insert(event.range_start, replacement_seq);
316 for (i, ordinal_out_idx) in msg_indices.iter_mut().enumerate() {
317 if i > event.range_start {
318 *ordinal_out_idx =
319 ordinal_out_idx.saturating_sub(removed_count.saturating_sub(1));
320 }
321 }
322 }
323 "compaction_summary" => {
324 out.push(TranscriptEntry::CompactionSummary {
325 range_start: v["range_start"].as_u64().unwrap_or(0) as usize,
326 range_end: v["range_end"].as_u64().unwrap_or(0) as usize,
327 compacted_count: v["compacted_count"].as_u64().unwrap_or(0) as usize,
328 before_tokens: v["before_tokens"].as_u64().unwrap_or(0),
329 after_tokens: v["after_tokens"].as_u64().unwrap_or(0),
330 summary: v["summary"].as_str().unwrap_or("").to_string(),
331 ts: parse_ts(v),
332 });
333 }
334 "diff_preview" => {
335 out.push(TranscriptEntry::DiffPreview {
336 title: v["title"].as_str().unwrap_or("").to_string(),
337 old_content: v["old_content"].as_str().map(String::from),
338 new_content: v["new_content"].as_str().map(String::from),
339 unified_diff: v["unified_diff"].as_str().map(String::from),
340 });
341 }
342 "flow_graph" => {
343 let run_id = v["run_id"].as_str().unwrap_or("").to_string();
344 let flow_name = v
345 .get("graph")
346 .and_then(|g| g["flow_name"].as_str())
347 .unwrap_or("")
348 .to_string();
349 let ts = parse_ts(v);
350 if let Some(g) = v.get("graph")
351 && let Ok(graph) = serde_json::from_value::<nodegraph::FlowGraph>(g.clone())
352 {
353 out.push(TranscriptEntry::FlowGraph {
354 run_id,
355 flow_name,
356 graph,
357 ts,
358 });
359 }
360 }
361 "flow_start" => {
362 let run_id = v["run_id"].as_str().unwrap_or("").to_string();
363 let flow_name = v["flow_name"].as_str().unwrap_or("").to_string();
364 let parent_run_id = v["parent_run_id"].as_str().map(String::from);
365 let parent_node_id = v["parent_node_id"].as_str().map(String::from);
366 let spawned = v.get("spawned").and_then(|s| s.as_bool()).unwrap_or(false);
367 let ts = parse_ts(v);
368 out.push(TranscriptEntry::FlowStart {
369 run_id,
370 flow_name,
371 parent_run_id,
372 parent_node_id,
373 spawned,
374 ts,
375 });
376 }
377 "flow_node_start" => {
378 let run_id = v["run_id"].as_str().unwrap_or("").to_string();
379 let node_id = v["node_id"].as_str().unwrap_or("").to_string();
380 let label = v["label"].as_str().unwrap_or(&node_id).to_string();
381 let parent_node_id = v["parent_node_id"].as_str().map(String::from);
382 let kind = v
383 .get("kind")
384 .and_then(|k| serde_json::from_value(k.clone()).ok())
385 .unwrap_or(nodegraph::NodeKind::UserConfirm);
386 let ts = parse_ts(v);
387 out.push(TranscriptEntry::FlowNodeStart {
388 run_id,
389 node_id,
390 kind,
391 label,
392 parent_node_id,
393 ts,
394 });
395 }
396 "flow_node_end" => {
397 let run_id = v["run_id"].as_str().unwrap_or("").to_string();
398 let node_id = v["node_id"].as_str().unwrap_or("").to_string();
399 let status: event::FlowNodeStatus = v
400 .get("status")
401 .and_then(|s| serde_json::from_value(s.clone()).ok())
402 .unwrap_or(event::FlowNodeStatus::Ok);
403 let output_preview = v["output_preview"].as_str().map(String::from);
404 let ts = parse_ts(v);
405 out.push(TranscriptEntry::FlowNodeEnd {
406 run_id,
407 node_id,
408 status,
409 output_preview,
410 ts,
411 });
412 }
413 "tool_node" => {
414 let run_id = v["run_id"].as_str().unwrap_or("").to_string();
415 let parent_node_id = v["parent_node_id"].as_str().unwrap_or("").to_string();
416 let tool_use_id = v["tool_use_id"].as_str().unwrap_or("").to_string();
417 let tool_name = v["tool_name"].as_str().unwrap_or("").to_string();
418 let args_preview = v["args_preview"].as_str().unwrap_or("").to_string();
419 let ts = parse_ts(v);
420 out.push(TranscriptEntry::ToolNode {
421 run_id,
422 parent_node_id,
423 tool_use_id,
424 tool_name,
425 args_preview,
426 ts,
427 });
428 }
429 "flow_end" => {
430 let run_id = v["run_id"].as_str().unwrap_or("").to_string();
431 let ok = v["status"]["kind"].as_str() == Some("ok");
432 let cancelled = v["status"]["kind"].as_str() == Some("cancelled");
433 let ts = parse_ts(v);
434 out.push(TranscriptEntry::FlowDone {
435 run_id,
436 ok,
437 cancelled,
438 ts,
439 });
440 }
441 "llm_call" => {
442 let model = v["model"].as_str().unwrap_or("").to_string();
443 let usage: provider::TokenUsage = v
444 .get("usage")
445 .and_then(|u| serde_json::from_value(u.clone()).ok())
446 .unwrap_or_default();
447 let wallclock_ms = v["wallclock_ms"].as_u64().unwrap_or(0);
448 let ttft_ms = v["ttft_ms"].as_u64();
449 let tokens_per_second = v["tokens_per_second"].as_f64();
450 let run_id = v["run_id"]
451 .as_str()
452 .and_then(|s| uuid::Uuid::parse_str(s).ok())
453 .map(event::FlowRunId);
454 let node_id = v["node_id"].as_str().map(String::from);
455 let ts = parse_ts(v);
456 out.push(TranscriptEntry::LlmCall {
457 model,
458 usage,
459 wallclock_ms,
460 ttft_ms,
461 tokens_per_second,
462 run_id,
463 node_id,
464 ts,
465 });
466 }
467 "terminal_final_state" => {
468 let handle = v["handle"].as_str().unwrap_or("").to_string();
469 if let Some(screen) = v.get("screen")
470 && let Ok(screen) =
471 serde_json::from_value::<crate::tools::term::TerminalScreen>(screen.clone())
472 {
473 out.push(TranscriptEntry::TerminalFinalState { handle, screen });
474 }
475 }
476 "mermaid_diagram" => {
477 if let Some(source) = v.get("source").and_then(|s| s.as_str()) {
478 out.push(TranscriptEntry::MermaidDiagram {
479 source: source.to_string(),
480 });
481 }
482 }
483 _ => {}
484 }
485 }
486 Ok(out)
487}
488
489pub trait MessageProjection {
490 fn to_messages(&self) -> Vec<Message>;
491 fn to_messages_with_seq(&self) -> Vec<(u64, Message)>;
492}
493
494impl MessageProjection for [crate::event::EventEnvelope] {
495 fn to_messages(&self) -> Vec<Message> {
496 self.to_messages_with_seq()
497 .into_iter()
498 .map(|(_, msg)| msg)
499 .collect()
500 }
501
502 fn to_messages_with_seq(&self) -> Vec<(u64, Message)> {
503 let spawned_flow_ids = spawned_flow_ids(self);
504 let mut acc: Vec<(u64, Message)> = Vec::new();
505 for env in self {
506 apply_envelope_to_messages(env, &spawned_flow_ids, &mut acc);
507 }
508 acc
509 }
510}
511
512pub(crate) fn spawned_flow_ids(
513 envelopes: &[crate::event::EventEnvelope],
514) -> std::collections::HashSet<crate::event::FlowRunId> {
515 let mut parents = std::collections::HashMap::new();
516 let mut spawned = std::collections::HashSet::new();
517 for env in envelopes {
518 if let crate::event::Event::FlowStart {
519 run_id,
520 parent_run_id,
521 spawned: is_spawned,
522 ..
523 } = &env.event
524 {
525 parents.insert(run_id.clone(), parent_run_id.clone());
526 if *is_spawned {
527 spawned.insert(run_id.clone());
528 }
529 }
530 }
531 loop {
532 let descendants: Vec<_> = parents
533 .iter()
534 .filter_map(|(run_id, parent)| {
535 (!spawned.contains(run_id)
536 && parent
537 .as_ref()
538 .is_some_and(|parent| spawned.contains(parent)))
539 .then_some(run_id.clone())
540 })
541 .collect();
542 if descendants.is_empty() {
543 break;
544 }
545 spawned.extend(descendants);
546 }
547 spawned
548}
549
550pub(crate) fn message_belongs_to_root(
551 flow_run_id: Option<&crate::event::FlowRunId>,
552 spawned_flow_ids: &std::collections::HashSet<crate::event::FlowRunId>,
553) -> bool {
554 flow_run_id.is_none_or(|run_id| !spawned_flow_ids.contains(run_id))
555}
556
557pub(crate) fn apply_envelope_to_messages(
558 env: &crate::event::EventEnvelope,
559 spawned_flow_ids: &std::collections::HashSet<crate::event::FlowRunId>,
560 acc: &mut Vec<(u64, Message)>,
561) {
562 match &env.event {
563 crate::event::Event::UserMsg {
564 message,
565 flow_run_id,
566 ..
567 }
568 | crate::event::Event::AssistantMsg {
569 message,
570 flow_run_id,
571 ..
572 }
573 | crate::event::Event::ToolResultMsg {
574 message,
575 flow_run_id,
576 ..
577 } if message_belongs_to_root(flow_run_id.as_ref(), spawned_flow_ids) => {
578 acc.push((env.seq, message.clone()));
579 }
580 crate::event::Event::SystemMsg { message, .. } => {
581 acc.push((env.seq, message.clone()));
582 }
583 crate::event::Event::ContextCompact {
584 compacted_range_start,
585 compacted_range_end,
586 replacement_msg_seq,
587 summary_text,
588 after_tokens,
589 before_tokens,
590 ..
591 } => {
592 let range_start = *compacted_range_start as usize;
593 let range_end = *compacted_range_end as usize;
594 if range_start > range_end || range_end >= acc.len() {
595 return;
596 }
597 let Some(rep_seq) = replacement_msg_seq else {
598 return;
599 };
600 let Some(rep_idx) = acc.iter().position(|(s, _)| *s == *rep_seq) else {
601 return;
602 };
603 if *after_tokens >= *before_tokens {
604 return;
605 }
606 let replacement = acc.remove(rep_idx);
607 let removed_count = range_end - range_start + 1;
608 for _ in 0..removed_count {
609 acc.remove(range_start);
610 }
611 let insertion_idx = range_start.min(acc.len());
612 if let Some(summary) = summary_text {
613 acc.insert(
614 insertion_idx,
615 (
616 *rep_seq,
617 Message::system_compact_summary(
618 crate::event::TurnId::now(),
619 summary.clone(),
620 range_start as u64,
621 range_end as u64,
622 removed_count,
623 ),
624 ),
625 );
626 } else {
627 acc.insert(insertion_idx, replacement);
628 }
629 }
630 crate::event::Event::Checkpoint { messages, .. } => {
631 acc.clear();
632 acc.extend(
633 messages
634 .iter()
635 .cloned()
636 .enumerate()
637 .map(|(index, message)| (u64::MAX.saturating_sub(index as u64), message)),
638 );
639 }
640 crate::event::Event::AttachmentDegraded {
641 message_seq,
642 part_index,
643 file_basename,
644 reason,
645 ..
646 } => {
647 if let Some((_, msg)) = acc.iter_mut().find(|(s, _)| *s == *message_seq) {
648 if let Some(part) = msg.parts.get_mut(*part_index) {
649 *part = MessagePart::Text {
650 text: format!("[attachment unavailable: {} — {}]", file_basename, reason),
651 };
652 }
653 }
654 }
655 _ => {}
656 }
657}
658
659#[cfg(test)]
660mod tests {
661 use crate::event::{Event, EventEnvelope, FlowRunId, TurnId};
662 use crate::message::{Message, MessageOrigin, MessagePart, MessageRole};
663 use uuid::Uuid;
664
665 fn message(role: MessageRole, text: &str) -> Message {
666 Message {
667 role,
668 parts: vec![MessagePart::Text {
669 text: text.to_string(),
670 }],
671 turn_id: TurnId::now(),
672 origin: MessageOrigin::User,
673 }
674 }
675
676 fn flow_start(run_id: FlowRunId, parent_run_id: Option<FlowRunId>, spawned: bool) -> Event {
677 Event::FlowStart {
678 run_id,
679 flow_name: "test".into(),
680 spawned,
681 parent_run_id,
682 parent_node_id: None,
683 }
684 }
685
686 #[test]
687 fn replay_excludes_subagent_messages() {
688 let root = FlowRunId(Uuid::now_v7());
689 let child = FlowRunId(Uuid::now_v7());
690 let envelopes = vec![
691 EventEnvelope::new(1, flow_start(root.clone(), None, false)),
692 EventEnvelope::new(2, flow_start(child.clone(), Some(root), true)),
693 EventEnvelope::new(
694 3,
695 Event::AssistantMsg {
696 turn_id: TurnId::now(),
697 flow_run_id: Some(child),
698 message: message(MessageRole::Assistant, "child"),
699 },
700 ),
701 ];
702 assert!(super::MessageProjection::to_messages(envelopes.as_slice()).is_empty());
703 }
704
705 #[test]
706 fn replay_excludes_ordinary_subflow_execution_messages() {
707 let root = FlowRunId(Uuid::now_v7());
708 let child = FlowRunId(Uuid::now_v7());
709 let envelopes = vec![
710 EventEnvelope::new(1, flow_start(root.clone(), None, false)),
711 EventEnvelope::new(2, flow_start(child.clone(), Some(root), false)),
712 EventEnvelope::new(
713 3,
714 Event::AssistantMsg {
715 turn_id: TurnId::now(),
716 flow_run_id: Some(child),
717 message: message(MessageRole::Assistant, "ordinary child"),
718 },
719 ),
720 ];
721 let messages = super::MessageProjection::to_messages(envelopes.as_slice());
722 assert_eq!(messages.len(), 1);
723 assert_eq!(messages[0].text_concat(), "ordinary child");
724 }
725
726 #[test]
727 fn replay_excludes_descendants_of_spawned_flows() {
728 let root = FlowRunId(Uuid::now_v7());
729 let spawned = FlowRunId(Uuid::now_v7());
730 let descendant = FlowRunId(Uuid::now_v7());
731 let envelopes = vec![
732 EventEnvelope::new(1, flow_start(root.clone(), None, false)),
733 EventEnvelope::new(2, flow_start(spawned.clone(), Some(root), true)),
734 EventEnvelope::new(3, flow_start(descendant.clone(), Some(spawned), false)),
735 EventEnvelope::new(
736 4,
737 Event::AssistantMsg {
738 turn_id: TurnId::now(),
739 flow_run_id: Some(descendant),
740 message: message(MessageRole::Assistant, "spawned descendant"),
741 },
742 ),
743 ];
744 assert!(super::MessageProjection::to_messages(envelopes.as_slice()).is_empty());
745 }
746
747 #[test]
748 fn replay_keeps_unknown_nonspawned_message() {
749 let orphan = FlowRunId(Uuid::now_v7());
750 let envelopes = vec![EventEnvelope::new(
751 1,
752 Event::AssistantMsg {
753 turn_id: TurnId::now(),
754 flow_run_id: Some(orphan),
755 message: message(MessageRole::Assistant, "orphan"),
756 },
757 )];
758 let messages = super::MessageProjection::to_messages(envelopes.as_slice());
759 assert_eq!(messages.len(), 1);
760 assert_eq!(messages[0].text_concat(), "orphan");
761 }
762}