1use 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 let w = ms.window();
435 assert_eq!(w.len(), 2);
436 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 #[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}