1use std::collections::BTreeSet;
2
3mod compaction;
4mod config;
5mod emergency;
6mod helpers;
7mod limits;
8mod microcompact;
9mod normalization;
10mod prompts;
11mod session_context;
12mod warnings;
13
14use crate::events::{MemoryCompactMode, ReservedOutputSource};
15use crate::memory::message_sanitizer::filter_empty_assistant_messages;
16use crate::memory::session::SessionMemory;
17use crate::memory::token_utils::count_messages_tokens;
18use crate::memory::{RuntimeMemoryCallbackError, RuntimeMemoryCallbacks};
19use crate::types::{Message, MessageRole};
20
21pub use config::{MemoryManagerConfig, SummaryCallback};
22
23use helpers::compact_processed_image_messages;
24
25const MEMORY_SUMMARY_NAME: &str = "memory_summary";
26
27#[derive(Debug, Clone)]
28pub struct MemoryManager {
29 pub config: MemoryManagerConfig,
30 session_memory: Option<SessionMemory>,
31 model_max_output_tokens: Option<u64>,
32 reserved_output_source: ReservedOutputSource,
33 runtime_callbacks: RuntimeMemoryCallbacks,
34}
35
36#[derive(Debug)]
37pub(crate) struct MemoryCompactionOutcome {
38 pub(crate) messages: Vec<Message>,
39 pub(crate) changed: bool,
40 pub(crate) mode: MemoryCompactMode,
41}
42
43impl MemoryManager {
44 pub fn new(mut config: MemoryManagerConfig) -> Self {
45 let session_memory = config.session_memory.take();
46 Self {
47 config,
48 session_memory,
49 model_max_output_tokens: None,
50 reserved_output_source: ReservedOutputSource::FrameworkFallback,
51 runtime_callbacks: RuntimeMemoryCallbacks::default(),
52 }
53 }
54
55 pub(crate) fn with_capacity_observation(
56 mut self,
57 model_max_output_tokens: Option<u64>,
58 reserved_output_source: ReservedOutputSource,
59 ) -> Self {
60 self.model_max_output_tokens = model_max_output_tokens;
61 self.reserved_output_source = reserved_output_source;
62 self
63 }
64
65 pub(crate) fn with_runtime_callbacks(mut self, callbacks: RuntimeMemoryCallbacks) -> Self {
66 self.runtime_callbacks = callbacks;
67 self
68 }
69
70 pub(crate) fn model_max_output_tokens(&self) -> Option<u64> {
71 self.model_max_output_tokens
72 }
73
74 pub(crate) fn reserved_output_source(&self) -> ReservedOutputSource {
75 self.reserved_output_source
76 }
77
78 pub fn compact(&mut self, messages: &[Message], force: bool) -> (Vec<Message>, bool) {
79 self.compact_for_cycle_with_usage_inner(messages, 0, None, force, None, None)
80 .expect("public memory compaction has no runtime callback control flow")
81 .into_tuple()
82 }
83
84 pub fn compact_for_cycle(
85 &mut self,
86 messages: &[Message],
87 cycle_index: u32,
88 force: bool,
89 ) -> (Vec<Message>, bool) {
90 self.compact_for_cycle_with_usage_inner(
91 messages,
92 cycle_index,
93 Some(cycle_index),
94 force,
95 None,
96 None,
97 )
98 .expect("public memory compaction has no runtime callback control flow")
99 .into_tuple()
100 }
101
102 pub fn compact_for_cycle_with_usage(
103 &mut self,
104 messages: &[Message],
105 cycle_index: u32,
106 force: bool,
107 total_tokens: Option<u64>,
108 recent_tool_call_ids: Option<&BTreeSet<String>>,
109 ) -> (Vec<Message>, bool) {
110 self.compact_for_cycle_with_usage_inner(
111 messages,
112 cycle_index,
113 Some(cycle_index),
114 force,
115 total_tokens,
116 recent_tool_call_ids,
117 )
118 .expect("public memory compaction has no runtime callback control flow")
119 .into_tuple()
120 }
121
122 pub(crate) fn compact_for_cycle_with_usage_observed(
123 &mut self,
124 messages: &[Message],
125 cycle_index: u32,
126 force: bool,
127 total_tokens: Option<u64>,
128 recent_tool_call_ids: Option<&BTreeSet<String>>,
129 ) -> Result<MemoryCompactionOutcome, RuntimeMemoryCallbackError> {
130 self.compact_for_cycle_with_usage_inner(
131 messages,
132 cycle_index,
133 Some(cycle_index),
134 force,
135 total_tokens,
136 recent_tool_call_ids,
137 )
138 }
139
140 fn compact_for_cycle_with_usage_inner(
141 &mut self,
142 messages: &[Message],
143 cycle_index: u32,
144 artifact_cycle_index: Option<u32>,
145 force: bool,
146 total_tokens: Option<u64>,
147 recent_tool_call_ids: Option<&BTreeSet<String>>,
148 ) -> Result<MemoryCompactionOutcome, RuntimeMemoryCallbackError> {
149 if messages.is_empty() {
150 return Ok(MemoryCompactionOutcome::new(
151 messages,
152 Vec::new(),
153 MemoryCompactMode::None,
154 false,
155 ));
156 }
157
158 let cleaned = self.remove_previous_summary(messages);
159 let sanitized = filter_empty_assistant_messages(&cleaned);
160 let changed_by_sanitize = sanitized != messages;
161 let mut changed = changed_by_sanitize;
162 let mut mode = if changed_by_sanitize {
163 MemoryCompactMode::Structural
164 } else {
165 MemoryCompactMode::None
166 };
167 let mut working_messages = self.apply_session_memory_context(&sanitized);
168 if working_messages != sanitized {
169 mode = mode.max(MemoryCompactMode::Structural);
170 }
171 let mut message_length =
172 self.calculate_effective_length(&working_messages, total_tokens, recent_tool_call_ids);
173 if let Some(session_memory) = self.session_memory.as_mut() {
174 let text_messages = working_messages
175 .iter()
176 .filter(|message| {
177 !matches!(message.role, MessageRole::System | MessageRole::Tool)
178 && !message.content.trim().is_empty()
179 })
180 .count();
181 let runtime_callback = self.runtime_callbacks.session_memory.as_ref();
182 let should_extract = match runtime_callback {
183 Some(_) => session_memory
184 .should_extract_with_runtime_callback(message_length, text_messages),
185 None => session_memory.should_extract(message_length, text_messages),
186 };
187 let extracted = if should_extract {
188 match runtime_callback {
189 Some(callback) => session_memory.extract_with_runtime_callback(
190 &working_messages,
191 cycle_index as i32,
192 message_length,
193 callback,
194 self.runtime_callbacks.session_memory_diagnostic.as_ref(),
195 )?,
196 None => session_memory.extract(
197 &working_messages,
198 cycle_index as i32,
199 message_length,
200 ),
201 }
202 } else {
203 0
204 };
205 if extracted > 0 {
206 let before_refresh = working_messages;
207 working_messages = self.apply_session_memory_context(&sanitized);
208 if working_messages != before_refresh {
209 mode = mode.max(MemoryCompactMode::Structural);
210 }
211 message_length =
212 self.calculate_effective_length(&working_messages, None, recent_tool_call_ids);
213 }
214 }
215 if !force && self.should_preemptive_microcompact(message_length) {
216 let (microcompacted, cleared) =
217 self.microcompact_messages(&working_messages, cycle_index);
218 if cleared > 0 {
219 working_messages = microcompacted;
220 mode = mode.max(MemoryCompactMode::Micro);
221 changed = true;
222 message_length = self.calculate_effective_length(&working_messages, None, None);
223 }
224 }
225 if !force && message_length <= self.autocompact_threshold() {
226 let (warned, warning_inserted) =
227 self.maybe_append_memory_warning(&working_messages, message_length);
228 if warning_inserted {
229 mode = mode.max(MemoryCompactMode::Structural);
230 changed = true;
231 }
232 return Ok(MemoryCompactionOutcome::new(
233 messages, warned, mode, changed,
234 ));
235 }
236 let mut summary_source = self.strip_session_memory_context(&working_messages);
237 if summary_source != working_messages {
238 mode = mode.max(MemoryCompactMode::Structural);
239 }
240 if !force {
241 let (image_compacted, image_changed) =
242 compact_processed_image_messages(&summary_source);
243 let (artifact_compacted, artifact_changed) =
244 self.compact_large_tool_results(&image_compacted, artifact_cycle_index);
245 if (image_changed || artifact_changed)
246 && count_messages_tokens(&artifact_compacted, &self.config.model)
247 <= self.autocompact_threshold()
248 {
249 return Ok(MemoryCompactionOutcome::new(
250 messages,
251 artifact_compacted,
252 mode.max(MemoryCompactMode::Structural),
253 true,
254 ));
255 }
256 if image_changed || artifact_changed {
257 mode = mode.max(MemoryCompactMode::Structural);
258 summary_source = artifact_compacted;
259 }
260 }
261 let (compacted, summary_changed) = self.compress_memory(
262 &summary_source,
263 artifact_cycle_index,
264 self.runtime_callbacks.memory_compaction.as_ref(),
265 )?;
266 if summary_changed {
267 mode = mode.max(MemoryCompactMode::Summary);
268 let post_compaction_tokens = count_messages_tokens(
269 &self.apply_session_memory_context(&compacted),
270 &self.config.model,
271 );
272 if let Some(session_memory) = self.session_memory.as_mut() {
273 session_memory.on_compaction(Some(post_compaction_tokens));
274 }
275 }
276 Ok(MemoryCompactionOutcome::new(
277 messages,
278 compacted,
279 mode,
280 changed || summary_changed,
281 ))
282 }
283
284 fn remove_previous_summary(&self, messages: &[Message]) -> Vec<Message> {
285 messages
286 .iter()
287 .filter(|message| {
288 !(message.role == MessageRole::System
289 && message.name.as_deref() == Some(MEMORY_SUMMARY_NAME))
290 })
291 .cloned()
292 .collect()
293 }
294}
295
296impl MemoryCompactionOutcome {
297 fn new(
298 original: &[Message],
299 messages: Vec<Message>,
300 mode: MemoryCompactMode,
301 changed: bool,
302 ) -> Self {
303 let content_changed = messages != original;
304 let mode = if !content_changed {
305 MemoryCompactMode::None
306 } else if mode == MemoryCompactMode::None {
307 MemoryCompactMode::Structural
308 } else {
309 mode
310 };
311 Self {
312 messages,
313 changed,
314 mode,
315 }
316 }
317
318 fn into_tuple(self) -> (Vec<Message>, bool) {
319 (self.messages, self.changed)
320 }
321}