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