Skip to main content

atman_runtime/projection/
message_window.rs

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}