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 = sanitized;
168 let mut message_length =
169 self.calculate_effective_length(&working_messages, total_tokens, recent_tool_call_ids);
170 if let Some(session_memory) = self.session_memory.as_mut() {
171 let text_messages = working_messages
172 .iter()
173 .filter(|message| {
174 !matches!(message.role, MessageRole::System | MessageRole::Tool)
175 && !message.content.trim().is_empty()
176 })
177 .count();
178 let runtime_callback = self.runtime_callbacks.session_memory.as_ref();
179 let should_extract = match runtime_callback {
180 Some(_) => session_memory
181 .should_extract_with_runtime_callback(message_length, text_messages),
182 None => session_memory.should_extract(message_length, text_messages),
183 };
184 if should_extract {
185 let _ = match runtime_callback {
186 Some(callback) => session_memory.extract_with_runtime_callback(
187 &working_messages,
188 cycle_index as i32,
189 message_length,
190 callback,
191 self.runtime_callbacks.session_memory_diagnostic.as_ref(),
192 )?,
193 None => session_memory.extract(
194 &working_messages,
195 cycle_index as i32,
196 message_length,
197 ),
198 };
199 }
200 }
201 if !force && self.should_preemptive_microcompact(message_length) {
202 let (microcompacted, cleared) =
203 self.microcompact_messages(&working_messages, cycle_index);
204 if cleared > 0 {
205 working_messages = microcompacted;
206 mode = mode.max(MemoryCompactMode::Micro);
207 changed = true;
208 message_length = self.calculate_effective_length(&working_messages, None, None);
209 }
210 }
211 if !force && message_length <= self.autocompact_threshold() {
212 let (warned, warning_inserted) =
213 self.maybe_append_memory_warning(&working_messages, message_length);
214 if warning_inserted {
215 mode = mode.max(MemoryCompactMode::Structural);
216 changed = true;
217 }
218 return Ok(MemoryCompactionOutcome::new(
219 messages, warned, mode, changed,
220 ));
221 }
222 let mut summary_source = working_messages;
223 if !force {
224 let (image_compacted, image_changed) =
225 compact_processed_image_messages(&summary_source);
226 let (artifact_compacted, artifact_changed) =
227 self.compact_large_tool_results(&image_compacted, artifact_cycle_index);
228 if (image_changed || artifact_changed)
229 && count_messages_tokens(&artifact_compacted, &self.config.model)
230 <= self.autocompact_threshold()
231 {
232 return Ok(MemoryCompactionOutcome::new(
233 messages,
234 artifact_compacted,
235 mode.max(MemoryCompactMode::Structural),
236 true,
237 ));
238 }
239 if image_changed || artifact_changed {
240 mode = mode.max(MemoryCompactMode::Structural);
241 summary_source = artifact_compacted;
242 }
243 }
244 let (compacted, summary_changed) = self.compress_memory(
245 &summary_source,
246 artifact_cycle_index,
247 self.runtime_callbacks.memory_compaction.as_ref(),
248 )?;
249 if summary_changed {
250 mode = mode.max(MemoryCompactMode::Summary);
251 let post_compaction_tokens = count_messages_tokens(&compacted, &self.config.model);
252 if let Some(session_memory) = self.session_memory.as_mut() {
253 session_memory.on_compaction(Some(post_compaction_tokens));
254 }
255 }
256 Ok(MemoryCompactionOutcome::new(
257 messages,
258 compacted,
259 mode,
260 changed || summary_changed,
261 ))
262 }
263
264 fn remove_previous_summary(&self, messages: &[Message]) -> Vec<Message> {
265 messages
266 .iter()
267 .filter(|message| {
268 !(message.role == MessageRole::System
269 && message.name.as_deref() == Some(MEMORY_SUMMARY_NAME))
270 })
271 .cloned()
272 .collect()
273 }
274}
275
276impl MemoryCompactionOutcome {
277 fn new(
278 original: &[Message],
279 messages: Vec<Message>,
280 mode: MemoryCompactMode,
281 changed: bool,
282 ) -> Self {
283 let content_changed = messages != original;
284 let mode = if !content_changed {
285 MemoryCompactMode::None
286 } else if mode == MemoryCompactMode::None {
287 MemoryCompactMode::Structural
288 } else {
289 mode
290 };
291 Self {
292 messages,
293 changed,
294 mode,
295 }
296 }
297
298 fn into_tuple(self) -> (Vec<Message>, bool) {
299 (self.messages, self.changed)
300 }
301}