use std::ops::Deref;
use std::sync::{Arc, Mutex};
use crate::compaction::is_compaction_summary;
use crate::event::EventEnvelope;
use crate::message::Message;
#[derive(Clone)]
pub struct MessageWindow {
messages: Arc<Vec<Message>>,
start: usize,
}
impl MessageWindow {
pub fn to_vec(&self) -> Vec<Message> {
self.as_slice().to_vec()
}
fn as_slice(&self) -> &[Message] {
&self.messages[self.start..]
}
}
impl Deref for MessageWindow {
type Target = [Message];
fn deref(&self) -> &Self::Target {
self.as_slice()
}
}
struct Acc {
compacted: Vec<(u64, Message)>,
full_raw: Vec<(u64, Message)>,
replayed: usize,
full_cache: Arc<Vec<Message>>,
window_cache: MessageWindow,
}
pub struct MessageStream {
events: Arc<Mutex<Vec<EventEnvelope>>>,
initial_compacted: Vec<(u64, Message)>,
initial_raw: Vec<(u64, Message)>,
acc: Mutex<Acc>,
}
impl MessageStream {
pub fn new(events: Arc<Mutex<Vec<EventEnvelope>>>) -> Self {
let empty = Arc::new(Vec::new());
Self {
events,
initial_compacted: Vec::new(),
initial_raw: Vec::new(),
acc: Mutex::new(Acc {
compacted: Vec::new(),
full_raw: Vec::new(),
replayed: 0,
full_cache: Arc::clone(&empty),
window_cache: MessageWindow {
messages: empty,
start: 0,
},
}),
}
}
pub fn with_initial(
events: Arc<Mutex<Vec<EventEnvelope>>>,
compacted: Vec<(u64, Message)>,
raw: Vec<(u64, Message)>,
) -> Self {
let full: Arc<Vec<Message>> = Arc::new(raw.iter().map(|(_, msg)| msg.clone()).collect());
let window_messages: Arc<Vec<Message>> =
Arc::new(compacted.iter().map(|(_, msg)| msg.clone()).collect());
let start = window_messages
.iter()
.rposition(is_compaction_summary)
.unwrap_or(0);
let window = MessageWindow {
messages: window_messages,
start,
};
Self {
events,
initial_compacted: compacted.clone(),
initial_raw: raw.clone(),
acc: Mutex::new(Acc {
compacted,
full_raw: raw,
replayed: 0,
full_cache: full,
window_cache: window,
}),
}
}
pub fn full_messages(&self) -> Arc<Vec<Message>> {
let events = self.events.lock().expect("events poisoned");
let mut acc = self.acc.lock().expect("acc poisoned");
self.ensure_fresh_locked(&events, &mut acc);
Arc::clone(&acc.full_cache)
}
pub fn window(&self) -> MessageWindow {
let events = self.events.lock().expect("events poisoned");
let mut acc = self.acc.lock().expect("acc poisoned");
self.ensure_fresh_locked(&events, &mut acc);
acc.window_cache.clone()
}
fn ensure_fresh_locked(&self, events: &[EventEnvelope], acc: &mut Acc) {
if acc.compacted.is_empty() {
acc.compacted = self.initial_compacted.clone();
acc.full_raw = self.initial_raw.clone();
}
if acc.replayed >= events.len() {
return;
}
let spawned_flow_ids = crate::projection::message_window::spawned_flow_ids(events);
for ev in &events[acc.replayed..] {
crate::projection::message_window::apply_envelope_to_messages(
ev,
&spawned_flow_ids,
&mut acc.compacted,
);
match &ev.event {
crate::event::Event::UserMsg {
message,
flow_run_id,
..
}
| crate::event::Event::AssistantMsg {
message,
flow_run_id,
..
}
| crate::event::Event::ToolResultMsg {
message,
flow_run_id,
..
} if crate::projection::message_window::message_belongs_to_root(
flow_run_id.as_ref(),
&spawned_flow_ids,
) =>
{
acc.full_raw.push((ev.seq, message.clone()));
}
crate::event::Event::SystemMsg { message, .. } => {
acc.full_raw.push((ev.seq, message.clone()));
}
_ => {}
}
}
acc.replayed = events.len();
let compacted: Vec<Message> = acc.compacted.iter().map(|(_, msg)| msg).cloned().collect();
let start = compacted
.iter()
.rposition(is_compaction_summary)
.unwrap_or(0);
let compacted_arc = Arc::new(compacted);
acc.window_cache = MessageWindow {
messages: Arc::clone(&compacted_arc),
start,
};
let raw: Vec<Message> = acc.full_raw.iter().map(|(_, msg)| msg).cloned().collect();
acc.full_cache = Arc::new(raw);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::event::TurnId;
use crate::event::{Event, EventEnvelope};
use crate::message::{MessageOrigin, MessagePart, MessageRole};
fn user(text: &str) -> Message {
Message {
role: MessageRole::User,
parts: vec![MessagePart::Text {
text: text.to_string(),
}],
turn_id: TurnId::now(),
origin: MessageOrigin::User,
}
}
fn assistant(text: &str) -> Message {
Message {
role: MessageRole::Assistant,
parts: vec![MessagePart::Text {
text: text.to_string(),
}],
turn_id: TurnId::now(),
origin: MessageOrigin::User,
}
}
fn compact_summary(text: &str) -> Message {
Message::system_compact_summary(TurnId::now(), text, 0, 1, 2)
}
fn make_msg_event(ty: &str, msg: &Message, _seq: u64) -> Event {
match ty {
"user_msg" => Event::UserMsg {
turn_id: msg.turn_id.clone(),
flow_run_id: None,
message: msg.clone(),
},
"assistant_msg" => Event::AssistantMsg {
turn_id: msg.turn_id.clone(),
flow_run_id: None,
message: msg.clone(),
},
"system_msg" => Event::SystemMsg {
turn_id: msg.turn_id.clone(),
message: msg.clone(),
},
_ => unreachable!(),
}
}
fn make_context_compact(
range_start: u64,
range_end: u64,
before_tokens: u64,
after_tokens: u64,
summary_text: &str,
replacement_msg_seq: u64,
) -> Event {
Event::ContextCompact {
session_id: "test".into(),
before_tokens,
after_tokens,
compacted_range_start: range_start,
compacted_range_end: range_end,
summary_text: Some(summary_text.into()),
replacement_msg_seq: Some(replacement_msg_seq),
}
}
fn event_envelopes(events: Vec<Event>) -> Arc<Mutex<Vec<EventEnvelope>>> {
Arc::new(Mutex::new(
events
.into_iter()
.enumerate()
.map(|(i, event)| EventEnvelope::new((i + 1) as u64, event))
.collect(),
))
}
#[test]
fn full_messages_filters_only_message_events() {
let u1 = user("hello");
let a1 = assistant("hi there");
let events = event_envelopes(vec![
make_msg_event("user_msg", &u1, 1),
Event::TurnStart {
turn_id: TurnId::now(),
},
make_msg_event("assistant_msg", &a1, 2),
Event::LlmCall {
model: "m".into(),
provider: "p".into(),
usage: crate::provider::TokenUsage::default(),
wallclock_ms: 0,
ttft_ms: None,
tokens_per_second: None,
status: crate::event::LlmCallStatus::Ok,
run_id: None,
node_id: None,
},
]);
let ms = MessageStream::new(events);
let msgs = ms.full_messages();
assert_eq!(msgs.len(), 2);
assert_eq!(msgs[0].text_concat(), "hello");
assert_eq!(msgs[1].text_concat(), "hi there");
}
#[test]
fn window_no_summary_returns_all() {
let events = vec![
make_msg_event("user_msg", &user("a"), 1),
make_msg_event("assistant_msg", &assistant("b"), 2),
make_msg_event("user_msg", &user("c"), 3),
];
let ms = MessageStream::new(event_envelopes(events));
assert_eq!(ms.window().len(), 3);
}
#[test]
fn window_single_summary_starts_from_it() {
let s1 = compact_summary("summary 1");
let events = vec![
make_msg_event("user_msg", &user("old"), 1),
make_msg_event("assistant_msg", &assistant("old"), 2),
make_msg_event("system_msg", &s1, 3),
make_msg_event("user_msg", &user("new"), 4),
make_msg_event("assistant_msg", &assistant("new"), 5),
];
let ms = MessageStream::new(event_envelopes(events));
let w = ms.window();
assert_eq!(w.len(), 3);
assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
}
#[test]
fn window_multiple_summaries_uses_last() {
let s1 = compact_summary("summary 1");
let s2 = compact_summary("summary 2");
let events = vec![
make_msg_event("system_msg", &s1, 1),
make_msg_event("user_msg", &user("m1"), 2),
make_msg_event("system_msg", &s2, 3),
make_msg_event("user_msg", &user("m2"), 4),
];
let ms = MessageStream::new(event_envelopes(events));
let w = ms.window();
assert_eq!(w.len(), 2);
assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
assert_eq!(summary, "summary 2");
}
}
#[test]
fn window_no_prefix_before_summary() {
let s1 = compact_summary("summary");
let events = vec![
make_msg_event("user_msg", &user("very old"), 1),
make_msg_event("assistant_msg", &assistant("very old"), 2),
make_msg_event("system_msg", &s1, 3),
make_msg_event("user_msg", &user("new"), 4),
];
let ms = MessageStream::new(event_envelopes(events));
let w = ms.window();
assert_eq!(w.len(), 2);
assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
assert_eq!(w[1].text_concat(), "new");
}
#[test]
fn window_empty_stream_returns_empty() {
let ms = MessageStream::new(event_envelopes(Vec::new()));
assert!(ms.window().is_empty());
}
#[test]
fn context_compact_replaces_range_with_summary() {
let events = vec![
make_msg_event("user_msg", &user("old u1"), 1),
make_msg_event("assistant_msg", &assistant("old a1"), 2),
make_msg_event("user_msg", &user("old u2"), 3),
make_msg_event("system_msg", &compact_summary("summary"), 4),
make_context_compact(0, 2, 100, 50, "compaction summary text", 4),
make_msg_event("user_msg", &user("after compact"), 5),
];
let ms = MessageStream::new(event_envelopes(events));
let w = ms.window();
assert_eq!(w.len(), 2);
assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
}
#[test]
fn multiple_compactions_applied_in_order() {
let events = vec![
make_msg_event("user_msg", &user("a"), 1),
make_msg_event("assistant_msg", &assistant("b"), 2),
make_msg_event("system_msg", &compact_summary("s1"), 3),
make_context_compact(0, 1, 200, 100, "first summary", 3),
make_msg_event("user_msg", &user("c"), 4),
make_msg_event("assistant_msg", &assistant("d"), 5),
make_msg_event("system_msg", &compact_summary("s2"), 6),
make_context_compact(1, 2, 150, 80, "second summary", 6),
make_msg_event("user_msg", &user("e"), 7),
];
let ms = MessageStream::new(event_envelopes(events));
let w = ms.window();
assert_eq!(w.len(), 2);
assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
assert_eq!(summary, "second summary");
}
}
#[test]
fn compact_then_user_message_produces_summary_plus_user() {
let events = vec![
make_msg_event("user_msg", &user("old u1"), 1),
make_msg_event("assistant_msg", &assistant("old a1"), 2),
make_msg_event("user_msg", &user("old u2"), 3),
make_msg_event("system_msg", &compact_summary("compact summary"), 4),
make_context_compact(0, 2, 200, 100, "compact summary", 4),
make_msg_event("user_msg", &user("new message after compact"), 5),
];
let ms = MessageStream::new(event_envelopes(events));
let w = ms.window();
assert_eq!(w.len(), 2);
assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
assert_eq!(w[1].text_concat(), "new message after compact");
}
#[test]
fn no_compaction_window_equals_full_messages() {
let events = vec![
make_msg_event("user_msg", &user("first"), 1),
make_msg_event("assistant_msg", &assistant("second"), 2),
make_msg_event("user_msg", &user("third"), 3),
];
let ms = MessageStream::new(event_envelopes(events));
assert_eq!(ms.full_messages().len(), 3);
assert_eq!(ms.window().len(), 3);
}
#[test]
fn full_messages_retains_compacted_history() {
let events = vec![
make_msg_event("user_msg", &user("old u1"), 1),
make_msg_event("assistant_msg", &assistant("old a1"), 2),
make_msg_event("user_msg", &user("old u2"), 3),
make_msg_event("system_msg", &compact_summary("summary"), 4),
make_context_compact(0, 2, 200, 100, "summary", 4),
make_msg_event("user_msg", &user("after compact"), 5),
];
let ms = MessageStream::new(event_envelopes(events));
let w = ms.window();
assert_eq!(w.len(), 2);
let f = ms.full_messages();
assert_eq!(f.len(), 5, "full must retain compacted messages");
assert_eq!(f[0].text_concat(), "old u1");
assert_eq!(f[1].text_concat(), "old a1");
assert_eq!(f[2].text_concat(), "old u2");
}
#[test]
fn third_compaction_replaces_second_summary() {
let events = vec![
make_msg_event("user_msg", &user("a"), 1),
make_msg_event("assistant_msg", &assistant("b"), 2),
make_msg_event("system_msg", &compact_summary("s1"), 3),
make_context_compact(0, 1, 100, 50, "s1 text", 3),
make_msg_event("user_msg", &user("c"), 4),
make_msg_event("system_msg", &compact_summary("s2"), 5),
make_context_compact(0, 1, 80, 40, "s2 text", 6),
make_msg_event("user_msg", &user("d"), 6),
make_msg_event("system_msg", &compact_summary("s3"), 7),
make_context_compact(0, 1, 70, 30, "s3 text", 9),
make_msg_event("user_msg", &user("final"), 8),
];
let ms = MessageStream::new(event_envelopes(events));
let w = ms.window();
assert_eq!(w.len(), 2);
if let MessagePart::CompactSummary { summary, .. } = &w[0].parts[0] {
assert_eq!(summary, "s3 text");
}
assert_eq!(w[1].text_concat(), "final");
}
#[test]
fn checkpoint_replaces_live_window_and_accepts_following_messages() {
let old = user("old");
let summary = compact_summary("checkpoint summary");
let retained = user("retained current user");
let events = event_envelopes(vec![
make_msg_event("user_msg", &old, 1),
Event::Checkpoint {
session_id: "test".into(),
messages: vec![summary.clone(), retained.clone()],
window_tokens: 10,
},
make_msg_event("assistant_msg", &assistant("next provider output"), 3),
]);
let ms = MessageStream::new(events);
let window = ms.window();
assert_eq!(window.len(), 3);
assert!(matches!(
window[0].parts[0],
MessagePart::CompactSummary { .. }
));
assert_eq!(window[1].text_concat(), "retained current user");
assert_eq!(window[2].text_concat(), "next provider output");
assert!(!window.iter().any(|message| message.text_concat() == "old"));
}
#[test]
fn reopened_session_uses_compacted_window_before_new_events() {
let initial_compacted = vec![
(10, compact_summary("checkpoint summary")),
(11, assistant("retained tail")),
];
let initial_raw = vec![(1, user("dead user")), (2, assistant("dead assistant"))];
let ms = MessageStream::with_initial(
Arc::new(Mutex::new(Vec::new())),
initial_compacted,
initial_raw,
);
let window = ms.window();
assert_eq!(window.len(), 2);
assert!(matches!(
window[0].parts[0],
MessagePart::CompactSummary { .. }
));
assert_eq!(window[1].text_concat(), "retained tail");
let full = ms.full_messages();
assert_eq!(full.len(), 2);
assert_eq!(full[0].text_concat(), "dead user");
assert_eq!(full[1].text_concat(), "dead assistant");
}
#[test]
fn reopened_session_keeps_initial_messages_after_new_events() {
let initial_compacted = vec![
(1, compact_summary("compaction summary")),
(2, assistant("tail assistant")),
];
let initial_raw = vec![
(1, compact_summary("compaction summary")),
(2, assistant("tail assistant")),
];
let events = Arc::new(Mutex::new(Vec::new()));
let ms = MessageStream::with_initial(events.clone(), initial_compacted, initial_raw);
events.lock().unwrap().push(EventEnvelope::new(
1,
Event::TurnStart {
turn_id: TurnId::now(),
},
));
events.lock().unwrap().push(EventEnvelope::new(
2,
Event::UserMsg {
turn_id: TurnId::now(),
flow_run_id: None,
message: user("latest user"),
},
));
let w = ms.window();
assert_eq!(w.len(), 3);
assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
assert_eq!(w[1].text_concat(), "tail assistant");
assert_eq!(w[2].text_concat(), "latest user");
}
#[test]
fn full_messages_retains_pre_compact_history_after_runtime_compact() {
let initial_compacted = vec![
(1, compact_summary("prior summary")),
(2, user("old user")),
(3, assistant("old assistant")),
];
let initial_raw = vec![
(1, compact_summary("prior summary")),
(2, user("old user")),
(3, assistant("old assistant")),
];
let events = Arc::new(Mutex::new(Vec::new()));
let ms = MessageStream::with_initial(events.clone(), initial_compacted, initial_raw);
events.lock().unwrap().push(EventEnvelope::new(
10,
Event::UserMsg {
turn_id: TurnId::now(),
flow_run_id: None,
message: user("new user before compact"),
},
));
events.lock().unwrap().push(EventEnvelope::new(
11,
Event::AssistantMsg {
turn_id: TurnId::now(),
flow_run_id: None,
message: assistant("new assistant before compact"),
},
));
let before = ms.full_messages();
assert_eq!(before.len(), 5, "pre-compact full should have all 5 msgs");
events.lock().unwrap().push(EventEnvelope::new(
12,
Event::SystemMsg {
turn_id: TurnId::now(),
message: compact_summary("runtime summary"),
},
));
events.lock().unwrap().push(EventEnvelope::new(
13,
Event::ContextCompact {
session_id: "test".into(),
before_tokens: 1000,
after_tokens: 100,
compacted_range_start: 1,
compacted_range_end: 2,
summary_text: Some("runtime summary".into()),
replacement_msg_seq: Some(12),
},
));
events.lock().unwrap().push(EventEnvelope::new(
14,
Event::UserMsg {
turn_id: TurnId::now(),
flow_run_id: None,
message: user("after compact user"),
},
));
let w = ms.window();
assert_eq!(w.len(), 4, "window after compact");
assert!(matches!(w[0].parts[0], MessagePart::CompactSummary { .. }));
let full = ms.full_messages();
assert!(
full.len() >= 6,
"full must retain pre-compact history, got {} msgs: {:?}",
full.len(),
full.iter().map(|m| m.text_concat()).collect::<Vec<_>>()
);
let texts: Vec<String> = full.iter().map(|m| m.text_concat()).collect();
assert!(
texts.iter().any(|t| t.contains("old user")),
"full must contain pre-compact 'old user', got: {:?}",
texts
);
assert!(
texts.iter().any(|t| t.contains("old assistant")),
"full must contain pre-compact 'old assistant', got: {:?}",
texts
);
}
#[test]
fn live_window_keeps_root_tree_once_and_excludes_spawned_tree() {
let events = Arc::new(Mutex::new(Vec::new()));
let stream = MessageStream::new(Arc::clone(&events));
let root = crate::event::FlowRunId::now();
let ordinary = crate::event::FlowRunId::now();
let spawned = crate::event::FlowRunId::now();
let descendant = crate::event::FlowRunId::now();
let flow_start = |run_id, parent_run_id, spawned| Event::FlowStart {
run_id,
flow_name: "test".into(),
spawned,
parent_run_id,
parent_node_id: None,
};
events.lock().unwrap().extend([
EventEnvelope::new(1, flow_start(root.clone(), None, false)),
EventEnvelope::new(2, flow_start(ordinary.clone(), Some(root.clone()), false)),
EventEnvelope::new(3, flow_start(spawned.clone(), Some(root.clone()), true)),
EventEnvelope::new(
4,
flow_start(descendant.clone(), Some(spawned.clone()), false),
),
EventEnvelope::new(
5,
Event::AssistantMsg {
turn_id: TurnId::now(),
flow_run_id: Some(root.clone()),
message: assistant("root one"),
},
),
]);
assert_eq!(stream.window().len(), 1);
events.lock().unwrap().extend([
EventEnvelope::new(
6,
Event::AssistantMsg {
turn_id: TurnId::now(),
flow_run_id: Some(ordinary),
message: assistant("ordinary one"),
},
),
EventEnvelope::new(
7,
Event::AssistantMsg {
turn_id: TurnId::now(),
flow_run_id: Some(spawned),
message: assistant("spawned one"),
},
),
EventEnvelope::new(
8,
Event::AssistantMsg {
turn_id: TurnId::now(),
flow_run_id: Some(descendant),
message: assistant("spawned descendant one"),
},
),
]);
let second = stream.window();
assert_eq!(second.len(), 2);
assert_eq!(second[0].text_concat(), "root one");
assert_eq!(second[1].text_concat(), "ordinary one");
events.lock().unwrap().push(EventEnvelope::new(
9,
Event::AssistantMsg {
turn_id: TurnId::now(),
flow_run_id: None,
message: assistant("durable root"),
},
));
let third = stream.window();
assert_eq!(third.len(), 3);
assert_eq!(third[2].text_concat(), "durable root");
assert_eq!(stream.window().len(), 3);
}
}