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