Skip to main content

atman_runtime/
message_stream.rs

1//! Incrementally maintained message accumulator.  `window()` returns a
2//! `MessageWindow` anchored at the last compaction summary (zero-copy);
3//! `full_messages()` returns a shared `Arc<Vec<Message>>` of every message.
4
5use std::ops::Deref;
6use std::sync::{Arc, Mutex};
7
8use crate::compaction::is_compaction_summary;
9use crate::event::EventEnvelope;
10use crate::message::Message;
11
12#[derive(Clone)]
13pub struct MessageWindow {
14    messages: Arc<Vec<Message>>,
15    start: usize,
16}
17
18impl MessageWindow {
19    pub fn to_vec(&self) -> Vec<Message> {
20        self.as_slice().to_vec()
21    }
22
23    fn as_slice(&self) -> &[Message] {
24        &self.messages[self.start..]
25    }
26}
27
28impl Deref for MessageWindow {
29    type Target = [Message];
30
31    fn deref(&self) -> &Self::Target {
32        self.as_slice()
33    }
34}
35
36struct Acc {
37    compacted: Vec<(u64, Message)>,
38    full_raw: Vec<(u64, Message)>,
39    replayed: usize,
40    full_cache: Arc<Vec<Message>>,
41    window_cache: MessageWindow,
42}
43
44pub struct MessageStream {
45    events: Arc<Mutex<Vec<EventEnvelope>>>,
46    initial_compacted: Vec<(u64, Message)>,
47    initial_raw: Vec<(u64, Message)>,
48    acc: Mutex<Acc>,
49}
50
51impl MessageStream {
52    pub fn new(events: Arc<Mutex<Vec<EventEnvelope>>>) -> Self {
53        let empty = Arc::new(Vec::new());
54        Self {
55            events,
56            initial_compacted: Vec::new(),
57            initial_raw: Vec::new(),
58            acc: Mutex::new(Acc {
59                compacted: Vec::new(),
60                full_raw: Vec::new(),
61                replayed: 0,
62                full_cache: Arc::clone(&empty),
63                window_cache: MessageWindow {
64                    messages: empty,
65                    start: 0,
66                },
67            }),
68        }
69    }
70
71    pub fn with_initial(
72        events: Arc<Mutex<Vec<EventEnvelope>>>,
73        compacted: Vec<(u64, Message)>,
74        raw: Vec<(u64, Message)>,
75    ) -> Self {
76        let full: Arc<Vec<Message>> = Arc::new(raw.iter().map(|(_, msg)| msg.clone()).collect());
77        let window_messages: Arc<Vec<Message>> =
78            Arc::new(compacted.iter().map(|(_, msg)| msg.clone()).collect());
79        let start = window_messages
80            .iter()
81            .rposition(is_compaction_summary)
82            .unwrap_or(0);
83        let window = MessageWindow {
84            messages: window_messages,
85            start,
86        };
87        Self {
88            events,
89            initial_compacted: compacted.clone(),
90            initial_raw: raw.clone(),
91            acc: Mutex::new(Acc {
92                compacted,
93                full_raw: raw,
94                replayed: 0,
95                full_cache: full,
96                window_cache: window,
97            }),
98        }
99    }
100
101    pub fn full_messages(&self) -> Arc<Vec<Message>> {
102        let events = self.events.lock().expect("events poisoned");
103        let mut acc = self.acc.lock().expect("acc poisoned");
104        self.ensure_fresh_locked(&events, &mut acc);
105        Arc::clone(&acc.full_cache)
106    }
107
108    pub fn window(&self) -> MessageWindow {
109        let events = self.events.lock().expect("events poisoned");
110        let mut acc = self.acc.lock().expect("acc poisoned");
111        self.ensure_fresh_locked(&events, &mut acc);
112        acc.window_cache.clone()
113    }
114
115    fn ensure_fresh_locked(&self, events: &[EventEnvelope], acc: &mut Acc) {
116        if acc.compacted.is_empty() {
117            acc.compacted = self.initial_compacted.clone();
118            acc.full_raw = self.initial_raw.clone();
119        }
120        if acc.replayed >= events.len() {
121            return;
122        }
123        let spawned_flow_ids = crate::projection::message_window::spawned_flow_ids(events);
124        for ev in &events[acc.replayed..] {
125            crate::projection::message_window::apply_envelope_to_messages(
126                ev,
127                &spawned_flow_ids,
128                &mut acc.compacted,
129            );
130            match &ev.event {
131                crate::event::Event::UserMsg {
132                    message,
133                    flow_run_id,
134                    ..
135                }
136                | crate::event::Event::AssistantMsg {
137                    message,
138                    flow_run_id,
139                    ..
140                }
141                | crate::event::Event::ToolResultMsg {
142                    message,
143                    flow_run_id,
144                    ..
145                } if crate::projection::message_window::message_belongs_to_root(
146                    flow_run_id.as_ref(),
147                    &spawned_flow_ids,
148                ) =>
149                {
150                    acc.full_raw.push((ev.seq, message.clone()));
151                }
152                crate::event::Event::SystemMsg { message, .. } => {
153                    acc.full_raw.push((ev.seq, message.clone()));
154                }
155                _ => {}
156            }
157        }
158        acc.replayed = events.len();
159
160        let compacted: Vec<Message> = acc.compacted.iter().map(|(_, msg)| msg).cloned().collect();
161        let start = compacted
162            .iter()
163            .rposition(is_compaction_summary)
164            .unwrap_or(0);
165        let compacted_arc = Arc::new(compacted);
166        acc.window_cache = MessageWindow {
167            messages: Arc::clone(&compacted_arc),
168            start,
169        };
170
171        let raw: Vec<Message> = acc.full_raw.iter().map(|(_, msg)| msg).cloned().collect();
172        acc.full_cache = Arc::new(raw);
173    }
174}
175
176#[cfg(test)]
177mod tests {
178    use super::*;
179    use crate::event::TurnId;
180    use crate::event::{Event, EventEnvelope};
181    use crate::message::{MessageOrigin, MessagePart, MessageRole};
182
183    fn user(text: &str) -> Message {
184        Message {
185            role: MessageRole::User,
186            parts: vec![MessagePart::Text {
187                text: text.to_string(),
188            }],
189            turn_id: TurnId::now(),
190            origin: MessageOrigin::User,
191        }
192    }
193
194    fn assistant(text: &str) -> Message {
195        Message {
196            role: MessageRole::Assistant,
197            parts: vec![MessagePart::Text {
198                text: text.to_string(),
199            }],
200            turn_id: TurnId::now(),
201            origin: MessageOrigin::User,
202        }
203    }
204
205    fn compact_summary(text: &str) -> Message {
206        Message::system_compact_summary(TurnId::now(), text, 0, 1, 2)
207    }
208
209    fn make_msg_event(ty: &str, msg: &Message, _seq: u64) -> Event {
210        match ty {
211            "user_msg" => Event::UserMsg {
212                turn_id: msg.turn_id.clone(),
213                flow_run_id: None,
214                message: msg.clone(),
215            },
216            "assistant_msg" => Event::AssistantMsg {
217                turn_id: msg.turn_id.clone(),
218                flow_run_id: None,
219                message: msg.clone(),
220            },
221            "system_msg" => Event::SystemMsg {
222                turn_id: msg.turn_id.clone(),
223                message: msg.clone(),
224            },
225            _ => unreachable!(),
226        }
227    }
228
229    fn make_context_compact(
230        range_start: u64,
231        range_end: u64,
232        before_tokens: u64,
233        after_tokens: u64,
234        summary_text: &str,
235        replacement_msg_seq: u64,
236    ) -> Event {
237        Event::ContextCompact {
238            session_id: "test".into(),
239            before_tokens,
240            after_tokens,
241            compacted_range_start: range_start,
242            compacted_range_end: range_end,
243            summary_text: Some(summary_text.into()),
244            replacement_msg_seq: Some(replacement_msg_seq),
245        }
246    }
247
248    fn event_envelopes(events: Vec<Event>) -> Arc<Mutex<Vec<EventEnvelope>>> {
249        Arc::new(Mutex::new(
250            events
251                .into_iter()
252                .enumerate()
253                .map(|(i, event)| EventEnvelope::new((i + 1) as u64, event))
254                .collect(),
255        ))
256    }
257
258    #[test]
259    fn full_messages_filters_only_message_events() {
260        let u1 = user("hello");
261        let a1 = assistant("hi there");
262        let events = event_envelopes(vec![
263            make_msg_event("user_msg", &u1, 1),
264            Event::TurnStart {
265                turn_id: TurnId::now(),
266            },
267            make_msg_event("assistant_msg", &a1, 2),
268            Event::LlmCall {
269                model: "m".into(),
270                provider: "p".into(),
271                usage: crate::provider::TokenUsage::default(),
272                wallclock_ms: 0,
273                ttft_ms: None,
274                tokens_per_second: None,
275                status: crate::event::LlmCallStatus::Ok,
276                run_id: None,
277                node_id: None,
278            },
279        ]);
280        let ms = MessageStream::new(events);
281        let msgs = ms.full_messages();
282        assert_eq!(msgs.len(), 2);
283        assert_eq!(msgs[0].text_concat(), "hello");
284        assert_eq!(msgs[1].text_concat(), "hi there");
285    }
286
287    #[test]
288    fn window_no_summary_returns_all() {
289        let events = vec![
290            make_msg_event("user_msg", &user("a"), 1),
291            make_msg_event("assistant_msg", &assistant("b"), 2),
292            make_msg_event("user_msg", &user("c"), 3),
293        ];
294        let ms = MessageStream::new(event_envelopes(events));
295        assert_eq!(ms.window().len(), 3);
296    }
297
298    #[test]
299    fn window_single_summary_starts_from_it() {
300        let s1 = compact_summary("summary 1");
301        let events = vec![
302            make_msg_event("user_msg", &user("old"), 1),
303            make_msg_event("assistant_msg", &assistant("old"), 2),
304            make_msg_event("system_msg", &s1, 3),
305            make_msg_event("user_msg", &user("new"), 4),
306            make_msg_event("assistant_msg", &assistant("new"), 5),
307        ];
308        let ms = MessageStream::new(event_envelopes(events));
309        let w = ms.window();
310        assert_eq!(w.len(), 3);
311        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
312    }
313
314    #[test]
315    fn window_multiple_summaries_uses_last() {
316        let s1 = compact_summary("summary 1");
317        let s2 = compact_summary("summary 2");
318        let events = vec![
319            make_msg_event("system_msg", &s1, 1),
320            make_msg_event("user_msg", &user("m1"), 2),
321            make_msg_event("system_msg", &s2, 3),
322            make_msg_event("user_msg", &user("m2"), 4),
323        ];
324        let ms = MessageStream::new(event_envelopes(events));
325        let w = ms.window();
326        assert_eq!(w.len(), 2);
327        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
328        if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
329            assert_eq!(summary, "summary 2");
330        }
331    }
332
333    #[test]
334    fn window_no_prefix_before_summary() {
335        let s1 = compact_summary("summary");
336        let events = vec![
337            make_msg_event("user_msg", &user("very old"), 1),
338            make_msg_event("assistant_msg", &assistant("very old"), 2),
339            make_msg_event("system_msg", &s1, 3),
340            make_msg_event("user_msg", &user("new"), 4),
341        ];
342        let ms = MessageStream::new(event_envelopes(events));
343        let w = ms.window();
344        assert_eq!(w.len(), 2);
345        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
346        assert_eq!(w[1].text_concat(), "new");
347    }
348
349    #[test]
350    fn window_empty_stream_returns_empty() {
351        let ms = MessageStream::new(event_envelopes(Vec::new()));
352        assert!(ms.window().is_empty());
353    }
354
355    #[test]
356    fn context_compact_replaces_range_with_summary() {
357        let events = vec![
358            make_msg_event("user_msg", &user("old u1"), 1),
359            make_msg_event("assistant_msg", &assistant("old a1"), 2),
360            make_msg_event("user_msg", &user("old u2"), 3),
361            make_msg_event("system_msg", &compact_summary("summary"), 4),
362            make_context_compact(0, 2, 100, 50, "compaction summary text", 4),
363            make_msg_event("user_msg", &user("after compact"), 5),
364        ];
365        let ms = MessageStream::new(event_envelopes(events));
366        let w = ms.window();
367        assert_eq!(w.len(), 2);
368        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
369    }
370
371    #[test]
372    fn multiple_compactions_applied_in_order() {
373        let events = vec![
374            make_msg_event("user_msg", &user("a"), 1),
375            make_msg_event("assistant_msg", &assistant("b"), 2),
376            make_msg_event("system_msg", &compact_summary("s1"), 3),
377            make_context_compact(0, 1, 200, 100, "first summary", 3),
378            make_msg_event("user_msg", &user("c"), 4),
379            make_msg_event("assistant_msg", &assistant("d"), 5),
380            make_msg_event("system_msg", &compact_summary("s2"), 6),
381            make_context_compact(1, 2, 150, 80, "second summary", 6),
382            make_msg_event("user_msg", &user("e"), 7),
383        ];
384        let ms = MessageStream::new(event_envelopes(events));
385        let w = ms.window();
386        assert_eq!(w.len(), 2);
387        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
388        if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
389            assert_eq!(summary, "second summary");
390        }
391    }
392
393    #[test]
394    fn compact_then_user_message_produces_summary_plus_user() {
395        let events = vec![
396            make_msg_event("user_msg", &user("old u1"), 1),
397            make_msg_event("assistant_msg", &assistant("old a1"), 2),
398            make_msg_event("user_msg", &user("old u2"), 3),
399            make_msg_event("system_msg", &compact_summary("compact summary"), 4),
400            make_context_compact(0, 2, 200, 100, "compact summary", 4),
401            make_msg_event("user_msg", &user("new message after compact"), 5),
402        ];
403        let ms = MessageStream::new(event_envelopes(events));
404        let w = ms.window();
405        assert_eq!(w.len(), 2);
406        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
407        assert_eq!(w[1].text_concat(), "new message after compact");
408    }
409
410    #[test]
411    fn no_compaction_window_equals_full_messages() {
412        let events = vec![
413            make_msg_event("user_msg", &user("first"), 1),
414            make_msg_event("assistant_msg", &assistant("second"), 2),
415            make_msg_event("user_msg", &user("third"), 3),
416        ];
417        let ms = MessageStream::new(event_envelopes(events));
418        assert_eq!(ms.full_messages().len(), 3);
419        assert_eq!(ms.window().len(), 3);
420    }
421
422    #[test]
423    fn full_messages_retains_compacted_history() {
424        let events = vec![
425            make_msg_event("user_msg", &user("old u1"), 1),
426            make_msg_event("assistant_msg", &assistant("old a1"), 2),
427            make_msg_event("user_msg", &user("old u2"), 3),
428            make_msg_event("system_msg", &compact_summary("summary"), 4),
429            make_context_compact(0, 2, 200, 100, "summary", 4),
430            make_msg_event("user_msg", &user("after compact"), 5),
431        ];
432        let ms = MessageStream::new(event_envelopes(events));
433        // Window: only compact summary + messages after it
434        let w = ms.window();
435        assert_eq!(w.len(), 2);
436        // Full: all messages including compacted ones
437        let f = ms.full_messages();
438        assert_eq!(f.len(), 5, "full must retain compacted messages");
439        assert_eq!(f[0].text_concat(), "old u1");
440        assert_eq!(f[1].text_concat(), "old a1");
441        assert_eq!(f[2].text_concat(), "old u2");
442    }
443
444    #[test]
445    fn third_compaction_replaces_second_summary() {
446        let events = vec![
447            make_msg_event("user_msg", &user("a"), 1),
448            make_msg_event("assistant_msg", &assistant("b"), 2),
449            make_msg_event("system_msg", &compact_summary("s1"), 3),
450            make_context_compact(0, 1, 100, 50, "s1 text", 3),
451            make_msg_event("user_msg", &user("c"), 4),
452            make_msg_event("system_msg", &compact_summary("s2"), 5),
453            make_context_compact(0, 1, 80, 40, "s2 text", 6),
454            make_msg_event("user_msg", &user("d"), 6),
455            make_msg_event("system_msg", &compact_summary("s3"), 7),
456            make_context_compact(0, 1, 70, 30, "s3 text", 9),
457            make_msg_event("user_msg", &user("final"), 8),
458        ];
459        let ms = MessageStream::new(event_envelopes(events));
460        let w = ms.window();
461        assert_eq!(w.len(), 2);
462        if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
463            assert_eq!(summary, "s3 text");
464        }
465        assert_eq!(w[1].text_concat(), "final");
466    }
467
468    #[test]
469    fn checkpoint_replaces_live_window_and_accepts_following_messages() {
470        let old = user("old");
471        let summary = compact_summary("checkpoint summary");
472        let retained = user("retained current user");
473        let events = event_envelopes(vec![
474            make_msg_event("user_msg", &old, 1),
475            Event::Checkpoint {
476                session_id: "test".into(),
477                messages: vec![summary.clone(), retained.clone()],
478                window_tokens: 10,
479            },
480            make_msg_event("assistant_msg", &assistant("next provider output"), 3),
481        ]);
482        let ms = MessageStream::new(events);
483
484        let window = ms.window();
485        assert_eq!(window.len(), 3);
486        assert!(matches!(
487            window[0].parts[0],
488            MessagePart::CompactSummary { .. }
489        ));
490        assert_eq!(window[1].text_concat(), "retained current user");
491        assert_eq!(window[2].text_concat(), "next provider output");
492        assert!(!window.iter().any(|message| message.text_concat() == "old"));
493    }
494
495    #[test]
496    fn reopened_session_uses_compacted_window_before_new_events() {
497        let initial_compacted = vec![
498            (10, compact_summary("checkpoint summary")),
499            (11, assistant("retained tail")),
500        ];
501        let initial_raw = vec![(1, user("dead user")), (2, assistant("dead assistant"))];
502        let ms = MessageStream::with_initial(
503            Arc::new(Mutex::new(Vec::new())),
504            initial_compacted,
505            initial_raw,
506        );
507
508        let window = ms.window();
509        assert_eq!(window.len(), 2);
510        assert!(matches!(
511            window[0].parts[0],
512            MessagePart::CompactSummary { .. }
513        ));
514        assert_eq!(window[1].text_concat(), "retained tail");
515
516        let full = ms.full_messages();
517        assert_eq!(full.len(), 2);
518        assert_eq!(full[0].text_concat(), "dead user");
519        assert_eq!(full[1].text_concat(), "dead assistant");
520    }
521
522    #[test]
523    fn reopened_session_keeps_initial_messages_after_new_events() {
524        let initial_compacted = vec![
525            (1, compact_summary("compaction summary")),
526            (2, assistant("tail assistant")),
527        ];
528        let initial_raw = vec![
529            (1, compact_summary("compaction summary")),
530            (2, assistant("tail assistant")),
531        ];
532        let events = Arc::new(Mutex::new(Vec::new()));
533        let ms = MessageStream::with_initial(events.clone(), initial_compacted, initial_raw);
534
535        events.lock().unwrap().push(EventEnvelope::new(
536            1,
537            Event::TurnStart {
538                turn_id: TurnId::now(),
539            },
540        ));
541        events.lock().unwrap().push(EventEnvelope::new(
542            2,
543            Event::UserMsg {
544                turn_id: TurnId::now(),
545                flow_run_id: None,
546                message: user("latest user"),
547            },
548        ));
549
550        let w = ms.window();
551        assert_eq!(w.len(), 3);
552        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
553        assert_eq!(w[1].text_concat(), "tail assistant");
554        assert_eq!(w[2].text_concat(), "latest user");
555    }
556
557    /// Regression: after a runtime ContextCompact, full_messages() must still
558    /// contain the pre-compact messages.
559    #[test]
560    fn full_messages_retains_pre_compact_history_after_runtime_compact() {
561        let initial_compacted = vec![
562            (1, compact_summary("prior summary")),
563            (2, user("old user")),
564            (3, assistant("old assistant")),
565        ];
566        let initial_raw = vec![
567            (1, compact_summary("prior summary")),
568            (2, user("old user")),
569            (3, assistant("old assistant")),
570        ];
571        let events = Arc::new(Mutex::new(Vec::new()));
572        let ms = MessageStream::with_initial(events.clone(), initial_compacted, initial_raw);
573
574        events.lock().unwrap().push(EventEnvelope::new(
575            10,
576            Event::UserMsg {
577                turn_id: TurnId::now(),
578                flow_run_id: None,
579                message: user("new user before compact"),
580            },
581        ));
582        events.lock().unwrap().push(EventEnvelope::new(
583            11,
584            Event::AssistantMsg {
585                turn_id: TurnId::now(),
586                flow_run_id: None,
587                message: assistant("new assistant before compact"),
588            },
589        ));
590
591        let before = ms.full_messages();
592        assert_eq!(before.len(), 5, "pre-compact full should have all 5 msgs");
593
594        events.lock().unwrap().push(EventEnvelope::new(
595            12,
596            Event::SystemMsg {
597                turn_id: TurnId::now(),
598                message: compact_summary("runtime summary"),
599            },
600        ));
601        events.lock().unwrap().push(EventEnvelope::new(
602            13,
603            Event::ContextCompact {
604                session_id: "test".into(),
605                before_tokens: 1000,
606                after_tokens: 100,
607                compacted_range_start: 1,
608                compacted_range_end: 2,
609                summary_text: Some("runtime summary".into()),
610                replacement_msg_seq: Some(12),
611            },
612        ));
613
614        events.lock().unwrap().push(EventEnvelope::new(
615            14,
616            Event::UserMsg {
617                turn_id: TurnId::now(),
618                flow_run_id: None,
619                message: user("after compact user"),
620            },
621        ));
622
623        let w = ms.window();
624        assert_eq!(w.len(), 4, "window after compact");
625        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
626
627        let full = ms.full_messages();
628        assert!(
629            full.len() >= 6,
630            "full must retain pre-compact history, got {} msgs: {:?}",
631            full.len(),
632            full.iter().map(|m| m.text_concat()).collect::<Vec<_>>()
633        );
634        let texts: Vec<String> = full.iter().map(|m| m.text_concat()).collect();
635        assert!(
636            texts.iter().any(|t| t.contains("old user")),
637            "full must contain pre-compact 'old user', got: {:?}",
638            texts
639        );
640        assert!(
641            texts.iter().any(|t| t.contains("old assistant")),
642            "full must contain pre-compact 'old assistant', got: {:?}",
643            texts
644        );
645    }
646
647    #[test]
648    fn live_window_keeps_root_tree_once_and_excludes_spawned_tree() {
649        let events = Arc::new(Mutex::new(Vec::new()));
650        let stream = MessageStream::new(Arc::clone(&events));
651        let root = crate::event::FlowRunId::now();
652        let ordinary = crate::event::FlowRunId::now();
653        let spawned = crate::event::FlowRunId::now();
654        let descendant = crate::event::FlowRunId::now();
655        let flow_start = |run_id, parent_run_id, spawned| Event::FlowStart {
656            run_id,
657            flow_name: "test".into(),
658            spawned,
659            parent_run_id,
660            parent_node_id: None,
661        };
662
663        events.lock().unwrap().extend([
664            EventEnvelope::new(1, flow_start(root.clone(), None, false)),
665            EventEnvelope::new(2, flow_start(ordinary.clone(), Some(root.clone()), false)),
666            EventEnvelope::new(3, flow_start(spawned.clone(), Some(root.clone()), true)),
667            EventEnvelope::new(
668                4,
669                flow_start(descendant.clone(), Some(spawned.clone()), false),
670            ),
671            EventEnvelope::new(
672                5,
673                Event::AssistantMsg {
674                    turn_id: TurnId::now(),
675                    flow_run_id: Some(root.clone()),
676                    message: assistant("root one"),
677                },
678            ),
679        ]);
680        assert_eq!(stream.window().len(), 1);
681
682        events.lock().unwrap().extend([
683            EventEnvelope::new(
684                6,
685                Event::AssistantMsg {
686                    turn_id: TurnId::now(),
687                    flow_run_id: Some(ordinary),
688                    message: assistant("ordinary one"),
689                },
690            ),
691            EventEnvelope::new(
692                7,
693                Event::AssistantMsg {
694                    turn_id: TurnId::now(),
695                    flow_run_id: Some(spawned),
696                    message: assistant("spawned one"),
697                },
698            ),
699            EventEnvelope::new(
700                8,
701                Event::AssistantMsg {
702                    turn_id: TurnId::now(),
703                    flow_run_id: Some(descendant),
704                    message: assistant("spawned descendant one"),
705                },
706            ),
707        ]);
708        let second = stream.window();
709        assert_eq!(second.len(), 2);
710        assert_eq!(second[0].text_concat(), "root one");
711        assert_eq!(second[1].text_concat(), "ordinary one");
712
713        events.lock().unwrap().push(EventEnvelope::new(
714            9,
715            Event::AssistantMsg {
716                turn_id: TurnId::now(),
717                flow_run_id: None,
718                message: assistant("durable root"),
719            },
720        ));
721        let third = stream.window();
722        assert_eq!(third.len(), 3);
723        assert_eq!(third[2].text_concat(), "durable root");
724        assert_eq!(stream.window().len(), 3);
725    }
726}