Skip to main content

vv_agent/memory/manager/
mod.rs

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}