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