1use std::collections::BTreeMap;
2use std::collections::BTreeSet;
3use std::fmt;
4use std::sync::Arc;
5
6mod compaction;
7mod config;
8mod emergency;
9mod helpers;
10mod limits;
11mod microcompact;
12mod normalization;
13mod prompts;
14mod session_context;
15mod warnings;
16
17use crate::events::{MemoryCompactMode, ReservedOutputSource};
18use crate::memory::message_sanitizer::filter_empty_assistant_messages;
19use crate::memory::session::SessionMemory;
20use crate::memory::token_utils::count_messages_tokens;
21use crate::memory::{RuntimeMemoryCallbackError, RuntimeMemoryCallbacks};
22use crate::tools::ToolResultRetention;
23use crate::types::{Message, MessageRole};
24use crate::workspace::WorkspaceBackend;
25
26pub use config::{MemoryManagerConfig, SummaryCallback};
27
28use helpers::compact_processed_image_messages;
29
30const MEMORY_SUMMARY_NAME: &str = "memory_summary";
31
32#[derive(Clone)]
33pub struct MemoryManager {
34 pub config: MemoryManagerConfig,
35 session_memory: Option<SessionMemory>,
36 model_max_output_tokens: Option<u64>,
37 reserved_output_source: ReservedOutputSource,
38 runtime_callbacks: RuntimeMemoryCallbacks,
39 workspace_backend: Option<Arc<dyn WorkspaceBackend>>,
40 artifact_namespace: String,
41 tool_result_retentions: BTreeMap<String, ToolResultRetention>,
42 recovery_tool_available: bool,
43}
44
45#[derive(Debug)]
46pub(crate) struct MemoryCompactionOutcome {
47 pub(crate) messages: Vec<Message>,
48 pub(crate) changed: bool,
49 pub(crate) mode: MemoryCompactMode,
50 pub(crate) archived_count: usize,
51 pub(crate) reclaimed_tokens: u64,
52 pub(crate) artifact_failure_count: usize,
53}
54
55struct CompactionRequest<'a> {
56 cycle_index: u32,
57 artifact_cycle_index: Option<u32>,
58 force: bool,
59 total_tokens: Option<u64>,
60 recent_tool_call_ids: Option<&'a BTreeSet<String>>,
61 microcompaction_plan: Option<microcompact::MicrocompactionPlan>,
62}
63
64impl fmt::Debug for MemoryManager {
65 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
66 formatter
67 .debug_struct("MemoryManager")
68 .field("config", &self.config)
69 .field("has_session_memory", &self.session_memory.is_some())
70 .field("model_max_output_tokens", &self.model_max_output_tokens)
71 .field("reserved_output_source", &self.reserved_output_source)
72 .field("has_workspace_backend", &self.workspace_backend.is_some())
73 .field("artifact_namespace", &self.artifact_namespace)
74 .field("tool_result_retentions", &self.tool_result_retentions)
75 .field("recovery_tool_available", &self.recovery_tool_available)
76 .finish()
77 }
78}
79
80impl MemoryManager {
81 pub fn new(mut config: MemoryManagerConfig) -> Self {
82 config
83 .microcompaction_policy
84 .validate()
85 .expect("MemoryManagerConfig has an invalid microcompaction policy");
86 let session_memory = config.session_memory.take();
87 Self {
88 config,
89 session_memory,
90 model_max_output_tokens: None,
91 reserved_output_source: ReservedOutputSource::FrameworkFallback,
92 runtime_callbacks: RuntimeMemoryCallbacks::default(),
93 workspace_backend: None,
94 artifact_namespace: "memory".to_string(),
95 tool_result_retentions: BTreeMap::new(),
96 recovery_tool_available: false,
97 }
98 }
99
100 pub fn with_workspace_backend(mut self, backend: Arc<dyn WorkspaceBackend>) -> Self {
101 self.workspace_backend = Some(backend);
102 self
103 }
104
105 pub fn with_recovery_tool_available(mut self, available: bool) -> Self {
106 self.recovery_tool_available = available;
107 self
108 }
109
110 pub(crate) fn with_archive_context(
111 mut self,
112 artifact_namespace: impl Into<String>,
113 tool_result_retentions: BTreeMap<String, ToolResultRetention>,
114 recovery_tool_available: bool,
115 ) -> Self {
116 self.artifact_namespace = artifact_namespace.into();
117 self.tool_result_retentions = tool_result_retentions;
118 self.recovery_tool_available = recovery_tool_available;
119 self
120 }
121
122 pub(crate) fn with_capacity_observation(
123 mut self,
124 model_max_output_tokens: Option<u64>,
125 reserved_output_source: ReservedOutputSource,
126 ) -> Self {
127 self.model_max_output_tokens = model_max_output_tokens;
128 self.reserved_output_source = reserved_output_source;
129 self
130 }
131
132 pub(crate) fn with_runtime_callbacks(mut self, callbacks: RuntimeMemoryCallbacks) -> Self {
133 self.runtime_callbacks = callbacks;
134 self
135 }
136
137 pub(crate) fn model_max_output_tokens(&self) -> Option<u64> {
138 self.model_max_output_tokens
139 }
140
141 pub(crate) fn reserved_output_source(&self) -> ReservedOutputSource {
142 self.reserved_output_source
143 }
144
145 pub fn compact(&mut self, messages: &[Message], force: bool) -> (Vec<Message>, bool) {
146 self.compact_for_cycle_with_usage_inner(
147 messages,
148 CompactionRequest {
149 cycle_index: 0,
150 artifact_cycle_index: None,
151 force,
152 total_tokens: None,
153 recent_tool_call_ids: None,
154 microcompaction_plan: None,
155 },
156 )
157 .expect("public memory compaction has no runtime callback control flow")
158 .into_tuple()
159 }
160
161 pub fn compact_for_cycle(
162 &mut self,
163 messages: &[Message],
164 cycle_index: u32,
165 force: bool,
166 ) -> (Vec<Message>, bool) {
167 self.compact_for_cycle_with_usage_inner(
168 messages,
169 CompactionRequest {
170 cycle_index,
171 artifact_cycle_index: Some(cycle_index),
172 force,
173 total_tokens: None,
174 recent_tool_call_ids: None,
175 microcompaction_plan: None,
176 },
177 )
178 .expect("public memory compaction has no runtime callback control flow")
179 .into_tuple()
180 }
181
182 pub fn compact_for_cycle_with_usage(
183 &mut self,
184 messages: &[Message],
185 cycle_index: u32,
186 force: bool,
187 total_tokens: Option<u64>,
188 recent_tool_call_ids: Option<&BTreeSet<String>>,
189 ) -> (Vec<Message>, bool) {
190 self.compact_for_cycle_with_usage_inner(
191 messages,
192 CompactionRequest {
193 cycle_index,
194 artifact_cycle_index: Some(cycle_index),
195 force,
196 total_tokens,
197 recent_tool_call_ids,
198 microcompaction_plan: None,
199 },
200 )
201 .expect("public memory compaction has no runtime callback control flow")
202 .into_tuple()
203 }
204
205 pub(crate) fn compact_for_cycle_with_usage_observed(
206 &mut self,
207 messages: &[Message],
208 cycle_index: u32,
209 force: bool,
210 total_tokens: Option<u64>,
211 recent_tool_call_ids: Option<&BTreeSet<String>>,
212 microcompaction_plan: Option<microcompact::MicrocompactionPlan>,
213 ) -> Result<MemoryCompactionOutcome, RuntimeMemoryCallbackError> {
214 self.compact_for_cycle_with_usage_inner(
215 messages,
216 CompactionRequest {
217 cycle_index,
218 artifact_cycle_index: Some(cycle_index),
219 force,
220 total_tokens,
221 recent_tool_call_ids,
222 microcompaction_plan,
223 },
224 )
225 }
226
227 fn compact_for_cycle_with_usage_inner(
228 &mut self,
229 messages: &[Message],
230 request: CompactionRequest<'_>,
231 ) -> Result<MemoryCompactionOutcome, RuntimeMemoryCallbackError> {
232 let CompactionRequest {
233 cycle_index,
234 artifact_cycle_index,
235 force,
236 total_tokens,
237 recent_tool_call_ids,
238 microcompaction_plan,
239 } = request;
240 if messages.is_empty() {
241 return Ok(MemoryCompactionOutcome::new(
242 messages,
243 Vec::new(),
244 MemoryCompactMode::None,
245 false,
246 microcompact::MicrocompactionApplication::default(),
247 ));
248 }
249
250 let cleaned = self.remove_previous_summary(messages);
251 let sanitized = filter_empty_assistant_messages(&cleaned);
252 let changed_by_sanitize = sanitized != messages;
253 let mut changed = changed_by_sanitize;
254 let mut mode = if changed_by_sanitize {
255 MemoryCompactMode::Structural
256 } else {
257 MemoryCompactMode::None
258 };
259 let mut working_messages = sanitized;
260 let mut message_length =
261 self.calculate_effective_length(&working_messages, total_tokens, recent_tool_call_ids);
262 if let Some(session_memory) = self.session_memory.as_mut() {
263 let text_messages = working_messages
264 .iter()
265 .filter(|message| {
266 !matches!(message.role, MessageRole::System | MessageRole::Tool)
267 && !message.content.trim().is_empty()
268 })
269 .count();
270 let runtime_callback = self.runtime_callbacks.session_memory.as_ref();
271 let should_extract = match runtime_callback {
272 Some(_) => session_memory
273 .should_extract_with_runtime_callback(message_length, text_messages),
274 None => session_memory.should_extract(message_length, text_messages),
275 };
276 if should_extract {
277 let _ = match runtime_callback {
278 Some(callback) => session_memory.extract_with_runtime_callback(
279 &working_messages,
280 cycle_index as i32,
281 message_length,
282 callback,
283 self.runtime_callbacks.session_memory_diagnostic.as_ref(),
284 )?,
285 None => session_memory.extract(
286 &working_messages,
287 cycle_index as i32,
288 message_length,
289 ),
290 };
291 }
292 }
293 let mut microcompaction = microcompact::MicrocompactionApplication::default();
294 if !force && self.should_preemptive_microcompact(message_length) {
295 let plan = microcompaction_plan.unwrap_or_else(|| {
296 self.plan_microcompaction(&working_messages, cycle_index, message_length)
297 });
298 microcompaction = self.apply_microcompaction(&working_messages, &plan);
299 if microcompaction.archived_count > 0 {
300 working_messages = microcompaction.messages.clone();
301 mode = mode.max(MemoryCompactMode::Micro);
302 changed = true;
303 message_length = message_length.saturating_sub(microcompaction.reclaimed_tokens);
304 }
305 }
306 if !force && message_length <= self.autocompact_threshold() {
307 let (warned, warning_inserted) =
308 self.maybe_append_memory_warning(&working_messages, message_length);
309 if warning_inserted {
310 mode = mode.max(MemoryCompactMode::Structural);
311 changed = true;
312 }
313 return Ok(MemoryCompactionOutcome::new(
314 messages,
315 warned,
316 mode,
317 changed,
318 microcompaction,
319 ));
320 }
321 let mut summary_source = working_messages;
322 if !force {
323 let before_structural_tokens =
324 count_messages_tokens(&summary_source, &self.config.model);
325 let (image_compacted, image_changed) =
326 compact_processed_image_messages(&summary_source);
327 let (artifact_compacted, artifact_changed) =
328 self.compact_large_tool_results(&image_compacted, artifact_cycle_index);
329 let after_structural_tokens =
330 count_messages_tokens(&artifact_compacted, &self.config.model);
331 message_length = if after_structural_tokens >= before_structural_tokens {
332 message_length.saturating_add(
333 after_structural_tokens.saturating_sub(before_structural_tokens),
334 )
335 } else {
336 message_length.saturating_sub(
337 before_structural_tokens.saturating_sub(after_structural_tokens),
338 )
339 };
340 if (image_changed || artifact_changed) && message_length <= self.autocompact_threshold()
341 {
342 return Ok(MemoryCompactionOutcome::new(
343 messages,
344 artifact_compacted,
345 mode.max(MemoryCompactMode::Structural),
346 true,
347 microcompaction,
348 ));
349 }
350 if image_changed || artifact_changed {
351 mode = mode.max(MemoryCompactMode::Structural);
352 summary_source = artifact_compacted;
353 }
354 }
355 let (compacted, summary_changed) = self.compress_memory(
356 &summary_source,
357 artifact_cycle_index,
358 self.runtime_callbacks.memory_compaction.as_ref(),
359 )?;
360 if summary_changed {
361 mode = mode.max(MemoryCompactMode::Summary);
362 let post_compaction_tokens = count_messages_tokens(&compacted, &self.config.model);
363 if let Some(session_memory) = self.session_memory.as_mut() {
364 session_memory.on_compaction(Some(post_compaction_tokens));
365 }
366 }
367 Ok(MemoryCompactionOutcome::new(
368 messages,
369 compacted,
370 mode,
371 changed || summary_changed,
372 microcompaction,
373 ))
374 }
375
376 fn remove_previous_summary(&self, messages: &[Message]) -> Vec<Message> {
377 messages
378 .iter()
379 .filter(|message| {
380 !(message.role == MessageRole::System
381 && message.name.as_deref() == Some(MEMORY_SUMMARY_NAME))
382 })
383 .cloned()
384 .collect()
385 }
386}
387
388impl MemoryCompactionOutcome {
389 fn new(
390 original: &[Message],
391 messages: Vec<Message>,
392 mode: MemoryCompactMode,
393 changed: bool,
394 microcompaction: microcompact::MicrocompactionApplication,
395 ) -> Self {
396 let content_changed = messages != original;
397 let mode = if !content_changed {
398 MemoryCompactMode::None
399 } else if mode == MemoryCompactMode::None {
400 MemoryCompactMode::Structural
401 } else {
402 mode
403 };
404 Self {
405 messages,
406 changed,
407 mode,
408 archived_count: microcompaction.archived_count,
409 reclaimed_tokens: microcompaction.reclaimed_tokens,
410 artifact_failure_count: microcompaction.artifact_failure_count,
411 }
412 }
413
414 fn into_tuple(self) -> (Vec<Message>, bool) {
415 (self.messages, self.changed)
416 }
417}