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::collections::{HashMap, HashSet, VecDeque};
6use std::ops::Deref;
7use std::sync::{Arc, Mutex};
8
9use crate::compaction::is_compaction_summary;
10use crate::event::EventEnvelope;
11use crate::message::Message;
12
13#[derive(Clone)]
14pub struct MessageWindow {
15    messages: Arc<Vec<Message>>,
16    start: usize,
17}
18
19impl MessageWindow {
20    pub fn to_vec(&self) -> Vec<Message> {
21        self.as_slice().to_vec()
22    }
23
24    fn as_slice(&self) -> &[Message] {
25        &self.messages[self.start..]
26    }
27}
28
29impl Deref for MessageWindow {
30    type Target = [Message];
31
32    fn deref(&self) -> &Self::Target {
33        self.as_slice()
34    }
35}
36
37struct Acc {
38    compacted: Vec<(u64, Message)>,
39    compacted_positions: HashMap<u64, usize>,
40    full_raw: Vec<(u64, Message)>,
41    full_positions: HashMap<u64, usize>,
42    replayed: usize,
43    projection_revision: u64,
44    ownership: FlowOwnership,
45    full_cache: Arc<Vec<Message>>,
46    window_cache: MessageWindow,
47}
48
49#[derive(Default)]
50struct FlowOwnership {
51    children: HashMap<crate::event::FlowRunId, HashSet<crate::event::FlowRunId>>,
52    spawned: HashSet<crate::event::FlowRunId>,
53}
54
55impl FlowOwnership {
56    fn observe(&mut self, event: &crate::event::Event) {
57        let crate::event::Event::FlowStart {
58            run_id,
59            parent_run_id,
60            spawned,
61            ..
62        } = event
63        else {
64            return;
65        };
66        if let Some(parent) = parent_run_id {
67            self.children
68                .entry(parent.clone())
69                .or_default()
70                .insert(run_id.clone());
71        }
72        if !*spawned
73            && !parent_run_id
74                .as_ref()
75                .is_some_and(|parent| self.spawned.contains(parent))
76        {
77            return;
78        }
79        let mut queue = VecDeque::from([run_id.clone()]);
80        while let Some(parent) = queue.pop_front() {
81            if !self.spawned.insert(parent.clone()) {
82                continue;
83            }
84            if let Some(children) = self.children.get(&parent) {
85                queue.extend(children.iter().cloned());
86            }
87        }
88    }
89}
90
91pub struct MessageStream {
92    events: Arc<Mutex<Vec<EventEnvelope>>>,
93    acc: Mutex<Acc>,
94}
95
96impl MessageStream {
97    pub fn new(events: Arc<Mutex<Vec<EventEnvelope>>>) -> Self {
98        let empty = Arc::new(Vec::new());
99        Self {
100            events,
101            acc: Mutex::new(Acc {
102                compacted: Vec::new(),
103                compacted_positions: HashMap::new(),
104                full_raw: Vec::new(),
105                full_positions: HashMap::new(),
106                replayed: 0,
107                projection_revision: 0,
108                ownership: FlowOwnership::default(),
109                full_cache: Arc::clone(&empty),
110                window_cache: MessageWindow {
111                    messages: empty,
112                    start: 0,
113                },
114            }),
115        }
116    }
117
118    pub fn with_initial(
119        events: Arc<Mutex<Vec<EventEnvelope>>>,
120        compacted: Vec<(u64, Message)>,
121        raw: Vec<(u64, Message)>,
122    ) -> Self {
123        let full: Arc<Vec<Message>> = Arc::new(raw.iter().map(|(_, msg)| msg.clone()).collect());
124        let window_messages: Arc<Vec<Message>> =
125            Arc::new(compacted.iter().map(|(_, msg)| msg.clone()).collect());
126        let start = window_messages
127            .iter()
128            .rposition(is_compaction_summary)
129            .unwrap_or(0);
130        let window = MessageWindow {
131            messages: window_messages,
132            start,
133        };
134        let projection_revision = u64::from(!compacted.is_empty() || !raw.is_empty());
135        let compacted_positions = crate::projection::message_window::message_positions(&compacted);
136        let full_positions = crate::projection::message_window::message_positions(&raw);
137        Self {
138            events,
139            acc: Mutex::new(Acc {
140                compacted,
141                compacted_positions,
142                full_raw: raw,
143                full_positions,
144                replayed: 0,
145                projection_revision,
146                ownership: FlowOwnership::default(),
147                full_cache: full,
148                window_cache: window,
149            }),
150        }
151    }
152
153    pub fn full_messages(&self) -> Arc<Vec<Message>> {
154        let events = self.events.lock().expect("events poisoned");
155        let mut acc = self.acc.lock().expect("acc poisoned");
156        self.ensure_fresh_locked(&events, &mut acc);
157        Arc::clone(&acc.full_cache)
158    }
159
160    pub fn window(&self) -> MessageWindow {
161        let events = self.events.lock().expect("events poisoned");
162        let mut acc = self.acc.lock().expect("acc poisoned");
163        self.ensure_fresh_locked(&events, &mut acc);
164        acc.window_cache.clone()
165    }
166
167    fn ensure_fresh_locked(&self, events: &[EventEnvelope], acc: &mut Acc) {
168        if acc.replayed >= events.len() {
169            return;
170        }
171        let mut compacted_changed = false;
172        let mut full_changed = false;
173        for event in &events[acc.replayed..] {
174            acc.ownership.observe(&event.event);
175        }
176        for ev in &events[acc.replayed..] {
177            compacted_changed |= crate::projection::message_window::apply_envelope_to_messages(
178                ev,
179                &acc.ownership.spawned,
180                &mut acc.compacted,
181                &mut acc.compacted_positions,
182            );
183            match &ev.event {
184                crate::event::Event::UserMsg {
185                    message,
186                    flow_run_id,
187                    ..
188                }
189                | crate::event::Event::AssistantMsg {
190                    message,
191                    flow_run_id,
192                    ..
193                }
194                | crate::event::Event::ToolResultMsg {
195                    message,
196                    flow_run_id,
197                    ..
198                }
199                | crate::event::Event::SystemMsg {
200                    message,
201                    flow_run_id,
202                    ..
203                } if crate::projection::message_window::message_belongs_to_root(
204                    flow_run_id.as_ref(),
205                    &acc.ownership.spawned,
206                ) =>
207                {
208                    acc.full_positions.insert(ev.seq, acc.full_raw.len());
209                    acc.full_raw.push((ev.seq, message.clone()));
210                    full_changed = true;
211                }
212                crate::event::Event::AttachmentDegraded { .. } => {
213                    full_changed |= crate::projection::message_window::apply_envelope_to_messages(
214                        ev,
215                        &acc.ownership.spawned,
216                        &mut acc.full_raw,
217                        &mut acc.full_positions,
218                    );
219                }
220                _ => {}
221            }
222        }
223        acc.replayed = events.len();
224
225        if compacted_changed {
226            let compacted = acc
227                .compacted
228                .iter()
229                .map(|(_, message)| message.clone())
230                .collect::<Vec<_>>();
231            let start = compacted
232                .iter()
233                .rposition(is_compaction_summary)
234                .unwrap_or(0);
235            acc.window_cache = MessageWindow {
236                messages: Arc::new(compacted),
237                start,
238            };
239        }
240        if full_changed {
241            acc.full_cache = Arc::new(
242                acc.full_raw
243                    .iter()
244                    .map(|(_, message)| message.clone())
245                    .collect(),
246            );
247        }
248        if compacted_changed || full_changed {
249            acc.projection_revision = acc.projection_revision.saturating_add(1);
250        }
251    }
252}
253
254#[cfg(test)]
255mod tests {
256    use super::*;
257    use crate::event::TurnId;
258    use crate::event::{Event, EventEnvelope};
259    use crate::message::{MessageOrigin, MessagePart, MessageRole};
260
261    fn user(text: &str) -> Message {
262        Message {
263            role: MessageRole::User,
264            parts: vec![MessagePart::Text {
265                text: text.to_string(),
266            }],
267            turn_id: TurnId::now(),
268            origin: MessageOrigin::User,
269        }
270    }
271
272    fn assistant(text: &str) -> Message {
273        Message {
274            role: MessageRole::Assistant,
275            parts: vec![MessagePart::Text {
276                text: text.to_string(),
277            }],
278            turn_id: TurnId::now(),
279            origin: MessageOrigin::User,
280        }
281    }
282
283    fn compact_summary(text: &str) -> Message {
284        Message::system_compact_summary(TurnId::now(), text, 0, 1, 2)
285    }
286
287    fn make_msg_event(ty: &str, msg: &Message, _seq: u64) -> Event {
288        match ty {
289            "user_msg" => Event::UserMsg {
290                turn_id: msg.turn_id.clone(),
291                flow_run_id: None,
292                message: msg.clone(),
293            },
294            "assistant_msg" => Event::AssistantMsg {
295                turn_id: msg.turn_id.clone(),
296                flow_run_id: None,
297                message: msg.clone(),
298            },
299            "system_msg" => Event::SystemMsg {
300                turn_id: msg.turn_id.clone(),
301                flow_run_id: None,
302                message: msg.clone(),
303            },
304            _ => unreachable!(),
305        }
306    }
307
308    fn make_context_compact(
309        range_start: u64,
310        range_end: u64,
311        before_tokens: u64,
312        after_tokens: u64,
313        summary_text: &str,
314        replacement_msg_seq: u64,
315    ) -> Event {
316        Event::ContextCompact {
317            session_id: "test".into(),
318            flow_run_id: None,
319            before_tokens,
320            after_tokens,
321            compacted_range_start: range_start,
322            compacted_range_end: range_end,
323            summary_text: Some(summary_text.into()),
324            replacement_msg_seq: Some(replacement_msg_seq),
325        }
326    }
327
328    fn event_envelopes(events: Vec<Event>) -> Arc<Mutex<Vec<EventEnvelope>>> {
329        Arc::new(Mutex::new(
330            events
331                .into_iter()
332                .enumerate()
333                .map(|(i, event)| EventEnvelope::new((i + 1) as u64, event))
334                .collect(),
335        ))
336    }
337
338    #[test]
339    fn full_messages_filters_only_message_events() {
340        let u1 = user("hello");
341        let a1 = assistant("hi there");
342        let events = event_envelopes(vec![
343            make_msg_event("user_msg", &u1, 1),
344            Event::TurnStart {
345                turn_id: TurnId::now(),
346            },
347            make_msg_event("assistant_msg", &a1, 2),
348            Event::LlmCall {
349                model: "m".into(),
350                provider: "p".into(),
351                context_plan_id: None,
352                context_epoch: None,
353                context_tokens: None,
354                usage_source: None,
355                context_call_purpose: None,
356                context_call_identity: None,
357                context_cache: None,
358                assistant_tool_batch_width: None,
359                usage: crate::provider::TokenUsage::default(),
360                wallclock_ms: 0,
361                ttft_ms: None,
362                tokens_per_second: None,
363                status: crate::event::LlmCallStatus::Ok,
364                run_id: None,
365                node_id: None,
366            },
367        ]);
368        let ms = MessageStream::new(events);
369        let msgs = ms.full_messages();
370        assert_eq!(msgs.len(), 2);
371        assert_eq!(msgs[0].text_concat(), "hello");
372        assert_eq!(msgs[1].text_concat(), "hi there");
373    }
374
375    #[test]
376    fn non_message_events_preserve_message_cache_identity() {
377        let events = Arc::new(Mutex::new(Vec::new()));
378        let stream = MessageStream::with_initial(
379            Arc::clone(&events),
380            vec![(1, user("window"))],
381            vec![(1, user("full"))],
382        );
383        let full_before = stream.full_messages();
384        let window_before = stream.window();
385        let revision_before = stream.acc.lock().unwrap().projection_revision;
386        events.lock().unwrap().extend([
387            EventEnvelope::new(
388                2,
389                Event::TurnStart {
390                    turn_id: TurnId::now(),
391                },
392            ),
393            EventEnvelope::new(
394                3,
395                Event::FlowStart {
396                    run_id: crate::event::FlowRunId::now(),
397                    flow_name: "root".into(),
398                    parent_run_id: None,
399                    parent_node_id: None,
400                    spawned: false,
401                },
402            ),
403            EventEnvelope::new(
404                4,
405                Event::LlmCall {
406                    model: "model".into(),
407                    provider: "provider".into(),
408                    context_plan_id: None,
409                    context_epoch: None,
410                    context_tokens: None,
411                    usage_source: None,
412                    context_call_purpose: None,
413                    context_call_identity: None,
414                    context_cache: None,
415                    assistant_tool_batch_width: None,
416                    usage: crate::provider::TokenUsage::default(),
417                    wallclock_ms: 0,
418                    ttft_ms: None,
419                    tokens_per_second: None,
420                    status: crate::event::LlmCallStatus::Ok,
421                    run_id: None,
422                    node_id: None,
423                },
424            ),
425        ]);
426
427        let full_after = stream.full_messages();
428        let window_after = stream.window();
429        let acc = stream.acc.lock().unwrap();
430
431        assert!(Arc::ptr_eq(&full_before, &full_after));
432        assert!(Arc::ptr_eq(&window_before.messages, &window_after.messages));
433        assert_eq!(acc.projection_revision, revision_before);
434        assert_eq!(acc.replayed, 3);
435    }
436
437    #[test]
438    fn compaction_rebuilds_window_without_cloning_full_history() {
439        let summary = compact_summary("summary");
440        let compacted = vec![
441            (1, user("old")),
442            (2, assistant("old")),
443            (3, summary.clone()),
444        ];
445        let raw = compacted.clone();
446        let events = Arc::new(Mutex::new(Vec::new()));
447        let stream = MessageStream::with_initial(Arc::clone(&events), compacted, raw);
448        let full_before = stream.full_messages();
449        let window_before = stream.window();
450        events.lock().unwrap().push(EventEnvelope::new(
451            4,
452            make_context_compact(0, 1, 100, 10, "summary", 3),
453        ));
454
455        let full_after = stream.full_messages();
456        let window_after = stream.window();
457
458        assert!(Arc::ptr_eq(&full_before, &full_after));
459        assert!(!Arc::ptr_eq(
460            &window_before.messages,
461            &window_after.messages
462        ));
463        assert_eq!(window_after.len(), 1);
464        assert!(matches!(
465            window_after[0].parts[0],
466            MessagePart::CompactSummary { .. }
467        ));
468    }
469
470    #[test]
471    fn attachment_degradation_updates_both_message_views() {
472        let message = user("attachment");
473        let events = Arc::new(Mutex::new(Vec::new()));
474        let stream = MessageStream::with_initial(
475            Arc::clone(&events),
476            vec![(5, message.clone())],
477            vec![(5, message)],
478        );
479        let full_before = stream.full_messages();
480        let window_before = stream.window();
481        events.lock().unwrap().push(EventEnvelope::new(
482            6,
483            Event::AttachmentDegraded {
484                turn_id: None,
485                flow_run_id: None,
486                message_seq: 5,
487                part_index: 0,
488                file_basename: "image.png".into(),
489                reason: "unreadable".into(),
490            },
491        ));
492
493        let full_after = stream.full_messages();
494        let window_after = stream.window();
495
496        assert!(!Arc::ptr_eq(&full_before, &full_after));
497        assert!(!Arc::ptr_eq(
498            &window_before.messages,
499            &window_after.messages
500        ));
501        assert_eq!(full_after[0].text_concat(), window_after[0].text_concat());
502        assert!(full_after[0].text_concat().contains("image.png"));
503    }
504
505    #[test]
506    fn window_no_summary_returns_all() {
507        let events = vec![
508            make_msg_event("user_msg", &user("a"), 1),
509            make_msg_event("assistant_msg", &assistant("b"), 2),
510            make_msg_event("user_msg", &user("c"), 3),
511        ];
512        let ms = MessageStream::new(event_envelopes(events));
513        assert_eq!(ms.window().len(), 3);
514    }
515
516    #[test]
517    fn window_single_summary_starts_from_it() {
518        let s1 = compact_summary("summary 1");
519        let events = vec![
520            make_msg_event("user_msg", &user("old"), 1),
521            make_msg_event("assistant_msg", &assistant("old"), 2),
522            make_msg_event("system_msg", &s1, 3),
523            make_msg_event("user_msg", &user("new"), 4),
524            make_msg_event("assistant_msg", &assistant("new"), 5),
525        ];
526        let ms = MessageStream::new(event_envelopes(events));
527        let w = ms.window();
528        assert_eq!(w.len(), 3);
529        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
530    }
531
532    #[test]
533    fn window_multiple_summaries_uses_last() {
534        let s1 = compact_summary("summary 1");
535        let s2 = compact_summary("summary 2");
536        let events = vec![
537            make_msg_event("system_msg", &s1, 1),
538            make_msg_event("user_msg", &user("m1"), 2),
539            make_msg_event("system_msg", &s2, 3),
540            make_msg_event("user_msg", &user("m2"), 4),
541        ];
542        let ms = MessageStream::new(event_envelopes(events));
543        let w = ms.window();
544        assert_eq!(w.len(), 2);
545        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
546        if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
547            assert_eq!(summary, "summary 2");
548        }
549    }
550
551    #[test]
552    fn window_no_prefix_before_summary() {
553        let s1 = compact_summary("summary");
554        let events = vec![
555            make_msg_event("user_msg", &user("very old"), 1),
556            make_msg_event("assistant_msg", &assistant("very old"), 2),
557            make_msg_event("system_msg", &s1, 3),
558            make_msg_event("user_msg", &user("new"), 4),
559        ];
560        let ms = MessageStream::new(event_envelopes(events));
561        let w = ms.window();
562        assert_eq!(w.len(), 2);
563        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
564        assert_eq!(w[1].text_concat(), "new");
565    }
566
567    #[test]
568    fn window_empty_stream_returns_empty() {
569        let ms = MessageStream::new(event_envelopes(Vec::new()));
570        assert!(ms.window().is_empty());
571    }
572
573    #[test]
574    fn context_compact_replaces_range_with_summary() {
575        let events = vec![
576            make_msg_event("user_msg", &user("old u1"), 1),
577            make_msg_event("assistant_msg", &assistant("old a1"), 2),
578            make_msg_event("user_msg", &user("old u2"), 3),
579            make_msg_event("system_msg", &compact_summary("summary"), 4),
580            make_context_compact(0, 2, 100, 50, "compaction summary text", 4),
581            make_msg_event("user_msg", &user("after compact"), 5),
582        ];
583        let ms = MessageStream::new(event_envelopes(events));
584        let w = ms.window();
585        assert_eq!(w.len(), 2);
586        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
587    }
588
589    #[test]
590    fn multiple_compactions_applied_in_order() {
591        let events = vec![
592            make_msg_event("user_msg", &user("a"), 1),
593            make_msg_event("assistant_msg", &assistant("b"), 2),
594            make_msg_event("system_msg", &compact_summary("s1"), 3),
595            make_context_compact(0, 1, 200, 100, "first summary", 3),
596            make_msg_event("user_msg", &user("c"), 4),
597            make_msg_event("assistant_msg", &assistant("d"), 5),
598            make_msg_event("system_msg", &compact_summary("s2"), 6),
599            make_context_compact(1, 2, 150, 80, "second summary", 7),
600            make_msg_event("user_msg", &user("e"), 7),
601        ];
602        let ms = MessageStream::new(event_envelopes(events));
603        let w = ms.window();
604        assert_eq!(w.len(), 2);
605        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
606        if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
607            assert_eq!(summary, "second summary");
608        }
609    }
610
611    #[test]
612    fn compact_then_user_message_produces_summary_plus_user() {
613        let events = vec![
614            make_msg_event("user_msg", &user("old u1"), 1),
615            make_msg_event("assistant_msg", &assistant("old a1"), 2),
616            make_msg_event("user_msg", &user("old u2"), 3),
617            make_msg_event("system_msg", &compact_summary("compact summary"), 4),
618            make_context_compact(0, 2, 200, 100, "compact summary", 4),
619            make_msg_event("user_msg", &user("new message after compact"), 5),
620        ];
621        let ms = MessageStream::new(event_envelopes(events));
622        let w = ms.window();
623        assert_eq!(w.len(), 2);
624        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
625        assert_eq!(w[1].text_concat(), "new message after compact");
626    }
627
628    #[test]
629    fn no_compaction_window_equals_full_messages() {
630        let events = vec![
631            make_msg_event("user_msg", &user("first"), 1),
632            make_msg_event("assistant_msg", &assistant("second"), 2),
633            make_msg_event("user_msg", &user("third"), 3),
634        ];
635        let ms = MessageStream::new(event_envelopes(events));
636        assert_eq!(ms.full_messages().len(), 3);
637        assert_eq!(ms.window().len(), 3);
638    }
639
640    #[test]
641    fn full_messages_retains_compacted_history() {
642        let events = vec![
643            make_msg_event("user_msg", &user("old u1"), 1),
644            make_msg_event("assistant_msg", &assistant("old a1"), 2),
645            make_msg_event("user_msg", &user("old u2"), 3),
646            make_msg_event("system_msg", &compact_summary("summary"), 4),
647            make_context_compact(0, 2, 200, 100, "summary", 4),
648            make_msg_event("user_msg", &user("after compact"), 5),
649        ];
650        let ms = MessageStream::new(event_envelopes(events));
651        // Window: only compact summary + messages after it
652        let w = ms.window();
653        assert_eq!(w.len(), 2);
654        // Full: all messages including compacted ones
655        let f = ms.full_messages();
656        assert_eq!(f.len(), 5, "full must retain compacted messages");
657        assert_eq!(f[0].text_concat(), "old u1");
658        assert_eq!(f[1].text_concat(), "old a1");
659        assert_eq!(f[2].text_concat(), "old u2");
660    }
661
662    #[test]
663    fn third_compaction_replaces_second_summary() {
664        let events = vec![
665            make_msg_event("user_msg", &user("a"), 1),
666            make_msg_event("assistant_msg", &assistant("b"), 2),
667            make_msg_event("system_msg", &compact_summary("s1"), 3),
668            make_context_compact(0, 1, 100, 50, "s1 text", 3),
669            make_msg_event("user_msg", &user("c"), 4),
670            make_msg_event("system_msg", &compact_summary("s2"), 5),
671            make_context_compact(0, 1, 80, 40, "s2 text", 6),
672            make_msg_event("user_msg", &user("d"), 6),
673            make_msg_event("system_msg", &compact_summary("s3"), 7),
674            make_context_compact(0, 1, 70, 30, "s3 text", 9),
675            make_msg_event("user_msg", &user("final"), 8),
676        ];
677        let ms = MessageStream::new(event_envelopes(events));
678        let w = ms.window();
679        assert_eq!(w.len(), 2);
680        if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
681            assert_eq!(summary, "s3 text");
682        }
683        assert_eq!(w[1].text_concat(), "final");
684    }
685
686    #[test]
687    fn checkpoint_replaces_live_window_and_accepts_following_messages() {
688        let old = user("old");
689        let summary = compact_summary("checkpoint summary");
690        let retained = user("retained current user");
691        let events = event_envelopes(vec![
692            make_msg_event("user_msg", &old, 1),
693            Event::Checkpoint {
694                session_id: "test".into(),
695                flow_run_id: None,
696                messages: vec![summary.clone(), retained.clone()],
697                window_tokens: 10,
698            },
699            make_msg_event("assistant_msg", &assistant("next provider output"), 3),
700        ]);
701        let ms = MessageStream::new(events);
702
703        let window = ms.window();
704        assert_eq!(window.len(), 3);
705        assert!(matches!(
706            window[0].parts[0],
707            MessagePart::CompactSummary { .. }
708        ));
709        assert_eq!(window[1].text_concat(), "retained current user");
710        assert_eq!(window[2].text_concat(), "next provider output");
711        assert!(!window.iter().any(|message| message.text_concat() == "old"));
712    }
713
714    #[test]
715    fn reopened_session_uses_compacted_window_before_new_events() {
716        let initial_compacted = vec![
717            (10, compact_summary("checkpoint summary")),
718            (11, assistant("retained tail")),
719        ];
720        let initial_raw = vec![(1, user("dead user")), (2, assistant("dead assistant"))];
721        let ms = MessageStream::with_initial(
722            Arc::new(Mutex::new(Vec::new())),
723            initial_compacted,
724            initial_raw,
725        );
726
727        let window = ms.window();
728        assert_eq!(window.len(), 2);
729        assert!(matches!(
730            window[0].parts[0],
731            MessagePart::CompactSummary { .. }
732        ));
733        assert_eq!(window[1].text_concat(), "retained tail");
734
735        let full = ms.full_messages();
736        assert_eq!(full.len(), 2);
737        assert_eq!(full[0].text_concat(), "dead user");
738        assert_eq!(full[1].text_concat(), "dead assistant");
739    }
740
741    #[test]
742    fn reopened_session_keeps_initial_messages_after_new_events() {
743        let initial_compacted = vec![
744            (1, compact_summary("compaction summary")),
745            (2, assistant("tail assistant")),
746        ];
747        let initial_raw = vec![
748            (1, compact_summary("compaction summary")),
749            (2, assistant("tail assistant")),
750        ];
751        let events = Arc::new(Mutex::new(Vec::new()));
752        let ms = MessageStream::with_initial(events.clone(), initial_compacted, initial_raw);
753
754        events.lock().unwrap().push(EventEnvelope::new(
755            1,
756            Event::TurnStart {
757                turn_id: TurnId::now(),
758            },
759        ));
760        events.lock().unwrap().push(EventEnvelope::new(
761            2,
762            Event::UserMsg {
763                turn_id: TurnId::now(),
764                flow_run_id: None,
765                message: user("latest user"),
766            },
767        ));
768
769        let w = ms.window();
770        assert_eq!(w.len(), 3);
771        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
772        assert_eq!(w[1].text_concat(), "tail assistant");
773        assert_eq!(w[2].text_concat(), "latest user");
774    }
775
776    /// Regression: after a runtime ContextCompact, full_messages() must still
777    /// contain the pre-compact messages.
778    #[test]
779    fn full_messages_retains_pre_compact_history_after_runtime_compact() {
780        let initial_compacted = vec![
781            (1, compact_summary("prior summary")),
782            (2, user("old user")),
783            (3, assistant("old assistant")),
784        ];
785        let initial_raw = vec![
786            (1, compact_summary("prior summary")),
787            (2, user("old user")),
788            (3, assistant("old assistant")),
789        ];
790        let events = Arc::new(Mutex::new(Vec::new()));
791        let ms = MessageStream::with_initial(events.clone(), initial_compacted, initial_raw);
792
793        events.lock().unwrap().push(EventEnvelope::new(
794            10,
795            Event::UserMsg {
796                turn_id: TurnId::now(),
797                flow_run_id: None,
798                message: user("new user before compact"),
799            },
800        ));
801        events.lock().unwrap().push(EventEnvelope::new(
802            11,
803            Event::AssistantMsg {
804                turn_id: TurnId::now(),
805                flow_run_id: None,
806                message: assistant("new assistant before compact"),
807            },
808        ));
809
810        let before = ms.full_messages();
811        assert_eq!(before.len(), 5, "pre-compact full should have all 5 msgs");
812
813        events.lock().unwrap().push(EventEnvelope::new(
814            12,
815            Event::SystemMsg {
816                turn_id: TurnId::now(),
817                flow_run_id: None,
818                message: compact_summary("runtime summary"),
819            },
820        ));
821        events.lock().unwrap().push(EventEnvelope::new(
822            13,
823            Event::ContextCompact {
824                session_id: "test".into(),
825                flow_run_id: None,
826                before_tokens: 1000,
827                after_tokens: 100,
828                compacted_range_start: 1,
829                compacted_range_end: 2,
830                summary_text: Some("runtime summary".into()),
831                replacement_msg_seq: Some(12),
832            },
833        ));
834
835        events.lock().unwrap().push(EventEnvelope::new(
836            14,
837            Event::UserMsg {
838                turn_id: TurnId::now(),
839                flow_run_id: None,
840                message: user("after compact user"),
841            },
842        ));
843
844        let w = ms.window();
845        assert_eq!(w.len(), 4, "window after compact");
846        assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
847
848        let full = ms.full_messages();
849        assert!(
850            full.len() >= 6,
851            "full must retain pre-compact history, got {} msgs: {:?}",
852            full.len(),
853            full.iter().map(|m| m.text_concat()).collect::<Vec<_>>()
854        );
855        let texts: Vec<String> = full.iter().map(|m| m.text_concat()).collect();
856        assert!(
857            texts.iter().any(|t| t.contains("old user")),
858            "full must contain pre-compact 'old user', got: {:?}",
859            texts
860        );
861        assert!(
862            texts.iter().any(|t| t.contains("old assistant")),
863            "full must contain pre-compact 'old assistant', got: {:?}",
864            texts
865        );
866    }
867
868    #[test]
869    fn spawned_compaction_and_checkpoint_do_not_mutate_root_window() {
870        let child_run_id = crate::event::FlowRunId::now();
871        let events = event_envelopes(vec![
872            make_msg_event("user_msg", &user("root user"), 1),
873            make_msg_event("assistant_msg", &assistant("root assistant"), 2),
874            Event::FlowStart {
875                run_id: child_run_id.clone(),
876                flow_name: "child".into(),
877                parent_run_id: None,
878                parent_node_id: None,
879                spawned: true,
880            },
881            Event::SystemMsg {
882                turn_id: TurnId::now(),
883                flow_run_id: Some(child_run_id.clone()),
884                message: compact_summary("child summary"),
885            },
886            Event::ContextCompact {
887                session_id: "test".into(),
888                flow_run_id: Some(child_run_id.clone()),
889                before_tokens: 100,
890                after_tokens: 10,
891                compacted_range_start: 0,
892                compacted_range_end: 1,
893                summary_text: Some("child summary".into()),
894                replacement_msg_seq: Some(4),
895            },
896            Event::Checkpoint {
897                session_id: "test".into(),
898                flow_run_id: Some(child_run_id),
899                messages: vec![compact_summary("child checkpoint")],
900                window_tokens: 10,
901            },
902        ]);
903        let stream = MessageStream::new(events);
904
905        assert_eq!(
906            stream
907                .window()
908                .iter()
909                .map(Message::text_concat)
910                .collect::<Vec<_>>(),
911            ["root user", "root assistant"]
912        );
913        assert_eq!(stream.full_messages().len(), 2);
914    }
915
916    #[test]
917    fn live_window_keeps_root_tree_once_and_excludes_spawned_tree() {
918        let events = Arc::new(Mutex::new(Vec::new()));
919        let stream = MessageStream::new(Arc::clone(&events));
920        let root = crate::event::FlowRunId::now();
921        let ordinary = crate::event::FlowRunId::now();
922        let spawned = crate::event::FlowRunId::now();
923        let descendant = crate::event::FlowRunId::now();
924        let flow_start = |run_id, parent_run_id, spawned| Event::FlowStart {
925            run_id,
926            flow_name: "test".into(),
927            spawned,
928            parent_run_id,
929            parent_node_id: None,
930        };
931
932        events.lock().unwrap().extend([
933            EventEnvelope::new(1, flow_start(root.clone(), None, false)),
934            EventEnvelope::new(2, flow_start(ordinary.clone(), Some(root.clone()), false)),
935            EventEnvelope::new(3, flow_start(spawned.clone(), Some(root.clone()), true)),
936            EventEnvelope::new(
937                4,
938                flow_start(descendant.clone(), Some(spawned.clone()), false),
939            ),
940            EventEnvelope::new(
941                5,
942                Event::AssistantMsg {
943                    turn_id: TurnId::now(),
944                    flow_run_id: Some(root.clone()),
945                    message: assistant("root one"),
946                },
947            ),
948        ]);
949        assert_eq!(stream.window().len(), 1);
950
951        events.lock().unwrap().extend([
952            EventEnvelope::new(
953                6,
954                Event::AssistantMsg {
955                    turn_id: TurnId::now(),
956                    flow_run_id: Some(ordinary),
957                    message: assistant("ordinary one"),
958                },
959            ),
960            EventEnvelope::new(
961                7,
962                Event::AssistantMsg {
963                    turn_id: TurnId::now(),
964                    flow_run_id: Some(spawned.clone()),
965                    message: assistant("spawned one"),
966                },
967            ),
968            EventEnvelope::new(
969                8,
970                Event::AssistantMsg {
971                    turn_id: TurnId::now(),
972                    flow_run_id: Some(descendant),
973                    message: assistant("spawned descendant one"),
974                },
975            ),
976        ]);
977        let second = stream.window();
978        assert_eq!(second.len(), 2);
979        assert_eq!(second[0].text_concat(), "root one");
980        assert_eq!(second[1].text_concat(), "ordinary one");
981
982        events.lock().unwrap().push(EventEnvelope::new(
983            9,
984            Event::AssistantMsg {
985                turn_id: TurnId::now(),
986                flow_run_id: None,
987                message: assistant("durable root"),
988            },
989        ));
990        let third = stream.window();
991        assert_eq!(third.len(), 3);
992        assert_eq!(third[2].text_concat(), "durable root");
993        assert_eq!(stream.window().len(), 3);
994
995        events.lock().unwrap().extend([
996            EventEnvelope::new(
997                10,
998                Event::SystemMsg {
999                    turn_id: TurnId::now(),
1000                    flow_run_id: Some(root),
1001                    message: Message::system_text(TurnId::now(), "root system"),
1002                },
1003            ),
1004            EventEnvelope::new(
1005                11,
1006                Event::SystemMsg {
1007                    turn_id: TurnId::now(),
1008                    flow_run_id: Some(spawned),
1009                    message: Message::system_text(TurnId::now(), "spawned system"),
1010                },
1011            ),
1012            EventEnvelope::new(
1013                12,
1014                Event::SystemMsg {
1015                    turn_id: TurnId::now(),
1016                    flow_run_id: None,
1017                    message: Message::system_text(TurnId::now(), "durable system"),
1018                },
1019            ),
1020        ]);
1021        let fourth = stream.window();
1022        assert_eq!(fourth.len(), 5);
1023        assert!(
1024            fourth
1025                .iter()
1026                .any(|message| message.text_concat() == "root system")
1027        );
1028        assert!(
1029            fourth
1030                .iter()
1031                .all(|message| message.text_concat() != "spawned system")
1032        );
1033        assert!(
1034            fourth
1035                .iter()
1036                .any(|message| message.text_concat() == "durable system")
1037        );
1038    }
1039
1040    #[test]
1041    fn batch_refresh_classifies_messages_after_collecting_flow_ownership() {
1042        let events = Arc::new(Mutex::new(Vec::new()));
1043        let stream = MessageStream::new(Arc::clone(&events));
1044        let spawned = crate::event::FlowRunId::now();
1045        events.lock().unwrap().extend([
1046            EventEnvelope::new(
1047                1,
1048                Event::AssistantMsg {
1049                    turn_id: TurnId::now(),
1050                    flow_run_id: Some(spawned.clone()),
1051                    message: assistant("spawned output"),
1052                },
1053            ),
1054            EventEnvelope::new(
1055                2,
1056                Event::FlowStart {
1057                    run_id: spawned,
1058                    flow_name: "spawned".into(),
1059                    parent_run_id: None,
1060                    parent_node_id: None,
1061                    spawned: true,
1062                },
1063            ),
1064        ]);
1065
1066        assert!(stream.window().is_empty());
1067        assert!(stream.full_messages().is_empty());
1068    }
1069}