Skip to main content

vv_agent/memory/manager/
mod.rs

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}