1pub mod repository;
19
20pub mod vector;
25
26pub use vector::{
27 InMemoryVectorIndex, VectorBudgetResource, VectorIndex, VectorIndexDescriptor,
28 VectorIndexError, VectorIndexObservation, VectorIndexStatus, VectorMetric,
29 VectorMutationConsistency, VectorNormalization, VectorRecord, VectorResult, VectorRevision,
30 VectorSearchHit, VectorSearchRequest, VectorSearchResult,
31};
32
33#[cfg(feature = "sqlite")]
34pub use vector::SqliteVectorIndex;
35
36#[cfg(feature = "sqlite")]
40pub mod sqlite;
41
42#[cfg(feature = "sqlite")]
43pub use sqlite::SqliteMemoryStore;
44
45use anyhow::Context as _;
46use chrono::{DateTime, Utc};
47use serde::{Deserialize, Serialize};
48use std::collections::HashMap;
49use std::sync::OnceLock;
50use tokio::sync::RwLock;
51
52const MIN_DEDUPE_FINGERPRINT_CHARS: usize = 24;
53const PRUNE_PROTECTED_ACCESS_COUNT: u32 = 3;
54
55#[derive(Debug, Clone, Serialize, Deserialize)]
61#[serde(rename_all = "camelCase")]
62pub struct RelevanceConfig {
63 #[serde(default = "RelevanceConfig::default_decay_days")]
65 pub decay_days: f32,
66 #[serde(default = "RelevanceConfig::default_importance_weight")]
68 pub importance_weight: f32,
69 #[serde(default = "RelevanceConfig::default_recency_weight")]
71 pub recency_weight: f32,
72}
73
74impl RelevanceConfig {
75 fn default_decay_days() -> f32 {
76 30.0
77 }
78 fn default_importance_weight() -> f32 {
79 0.7
80 }
81 fn default_recency_weight() -> f32 {
82 0.3
83 }
84}
85
86impl Default for RelevanceConfig {
87 fn default() -> Self {
88 Self {
89 decay_days: 30.0,
90 importance_weight: 0.7,
91 recency_weight: 0.3,
92 }
93 }
94}
95
96#[derive(Debug, Clone, Serialize, Deserialize)]
98#[serde(rename_all = "camelCase")]
99pub struct PrunePolicy {
100 #[serde(default = "PrunePolicy::default_max_age_days")]
102 pub max_age_days: u32,
103 #[serde(default = "PrunePolicy::default_min_importance_to_keep")]
106 pub min_importance_to_keep: f32,
107 #[serde(default)]
110 pub max_items: usize,
111}
112
113impl PrunePolicy {
114 fn default_max_age_days() -> u32 {
115 90
116 }
117 fn default_min_importance_to_keep() -> f32 {
118 0.5
119 }
120}
121
122impl Default for PrunePolicy {
123 fn default() -> Self {
124 Self {
125 max_age_days: 90,
126 min_importance_to_keep: 0.5,
127 max_items: 0,
128 }
129 }
130}
131
132#[derive(Debug, Clone, Serialize, Deserialize)]
138pub struct MemoryItem {
139 pub id: String,
140 pub content: String,
141 pub timestamp: DateTime<Utc>,
142 pub importance: f32,
143 pub tags: Vec<String>,
144 pub memory_type: MemoryType,
145 pub metadata: HashMap<String, String>,
146 pub access_count: u32,
147 pub last_accessed: Option<DateTime<Utc>>,
148 #[serde(skip)]
149 pub content_lower: String,
150}
151
152impl MemoryItem {
153 pub fn new(content: impl Into<String>) -> Self {
154 let content = content.into();
155 let content_lower = content.to_lowercase();
156 Self {
157 id: uuid::Uuid::new_v4().to_string(),
158 content,
159 timestamp: Utc::now(),
160 importance: 0.5,
161 tags: Vec::new(),
162 memory_type: MemoryType::Episodic,
163 metadata: HashMap::new(),
164 access_count: 0,
165 last_accessed: None,
166 content_lower,
167 }
168 }
169
170 pub fn with_importance(mut self, importance: f32) -> Self {
171 self.importance = importance.clamp(0.0, 1.0);
172 self
173 }
174
175 pub fn with_tags(mut self, tags: Vec<String>) -> Self {
176 self.tags = tags;
177 self
178 }
179
180 pub fn with_tag(mut self, tag: impl Into<String>) -> Self {
181 self.tags.push(tag.into());
182 self
183 }
184
185 pub fn with_type(mut self, memory_type: MemoryType) -> Self {
186 self.memory_type = memory_type;
187 self
188 }
189
190 pub fn with_metadata(mut self, key: impl Into<String>, value: impl Into<String>) -> Self {
191 self.metadata.insert(key.into(), value.into());
192 self
193 }
194
195 pub fn content_fingerprint(&self) -> Option<String> {
201 memory_content_fingerprint(&self.content)
202 }
203
204 pub fn merge_duplicate(self, incoming: MemoryItem) -> MemoryItem {
209 merge_duplicate_memory_item(self, incoming)
210 }
211
212 pub fn record_access(&mut self) {
213 self.access_count += 1;
214 self.last_accessed = Some(Utc::now());
215 }
216
217 pub fn relevance_score_at(&self, now: DateTime<Utc>, config: &RelevanceConfig) -> f32 {
219 let age_days = (now - self.timestamp).num_seconds() as f32 / 86400.0;
220 let decay = (-age_days / config.decay_days).exp();
221 self.importance * config.importance_weight + decay * config.recency_weight
222 }
223
224 pub fn relevance_score(&self) -> f32 {
226 self.relevance_score_at(Utc::now(), &RelevanceConfig::default())
227 }
228}
229
230fn normalize_item_for_store(mut item: MemoryItem) -> MemoryItem {
231 item.content_lower = item.content.to_lowercase();
232 item
233}
234
235fn memory_content_fingerprint(content: &str) -> Option<String> {
236 let fingerprint = content
237 .split_whitespace()
238 .collect::<Vec<_>>()
239 .join(" ")
240 .to_lowercase();
241 if fingerprint.chars().count() < MIN_DEDUPE_FINGERPRINT_CHARS {
242 None
243 } else {
244 Some(fingerprint)
245 }
246}
247
248fn memories_are_store_duplicates(existing: &MemoryItem, incoming: &MemoryItem) -> bool {
249 if existing.id == incoming.id {
250 return true;
251 }
252 existing.content_fingerprint().is_some()
253 && existing.content_fingerprint() == incoming.content_fingerprint()
254}
255
256fn memory_index_entry_is_duplicate(entry: &IndexEntry, incoming: &MemoryItem) -> bool {
257 if entry.id == incoming.id {
258 return true;
259 }
260 incoming.content_fingerprint().is_some()
261 && memory_content_fingerprint(&entry.content_lower) == incoming.content_fingerprint()
262}
263
264fn merge_duplicate_memory_item(existing: MemoryItem, incoming: MemoryItem) -> MemoryItem {
265 let incoming = normalize_item_for_store(incoming);
266 let mut merged = existing.clone();
267
268 if should_replace_duplicate_content(&existing, &incoming) {
269 merged.content = incoming.content.clone();
270 merged.content_lower = incoming.content_lower.clone();
271 }
272 merged.importance = existing.importance.max(incoming.importance);
273 merged.timestamp = existing.timestamp.max(incoming.timestamp);
274 merged.memory_type = stronger_memory_type(existing.memory_type, incoming.memory_type);
275 merged.access_count = existing.access_count.max(incoming.access_count);
276 merged.last_accessed = max_optional_datetime(existing.last_accessed, incoming.last_accessed);
277
278 merge_tags(&mut merged.tags, &incoming.tags);
279 merge_metadata(&mut merged.metadata, &incoming.metadata);
280 record_duplicate_metadata(&mut merged.metadata, &incoming.id);
281
282 normalize_item_for_store(merged)
283}
284
285fn should_replace_duplicate_content(existing: &MemoryItem, incoming: &MemoryItem) -> bool {
286 incoming.importance > existing.importance
287 || (incoming.importance == existing.importance
288 && incoming.content.chars().count() > existing.content.chars().count())
289}
290
291fn stronger_memory_type(existing: MemoryType, incoming: MemoryType) -> MemoryType {
292 if memory_type_strength(incoming) > memory_type_strength(existing) {
293 incoming
294 } else {
295 existing
296 }
297}
298
299fn memory_type_strength(memory_type: MemoryType) -> u8 {
300 match memory_type {
301 MemoryType::Procedural | MemoryType::Semantic => 3,
302 MemoryType::Working => 2,
303 MemoryType::Episodic => 1,
304 }
305}
306
307fn max_optional_datetime(
308 left: Option<DateTime<Utc>>,
309 right: Option<DateTime<Utc>>,
310) -> Option<DateTime<Utc>> {
311 match (left, right) {
312 (Some(left), Some(right)) => Some(left.max(right)),
313 (Some(left), None) => Some(left),
314 (None, Some(right)) => Some(right),
315 (None, None) => None,
316 }
317}
318
319fn merge_tags(existing: &mut Vec<String>, incoming: &[String]) {
320 for tag in incoming {
321 if !existing.contains(tag) {
322 existing.push(tag.clone());
323 }
324 }
325}
326
327fn merge_metadata(existing: &mut HashMap<String, String>, incoming: &HashMap<String, String>) {
328 for (key, value) in incoming {
329 if value.trim().is_empty() {
330 continue;
331 }
332 match existing.get_mut(key) {
333 Some(current) if current == value => {}
334 Some(current) if is_list_metadata_key(key) => {
335 *current = merge_metadata_list(current, value);
336 }
337 Some(_) => {}
338 None => {
339 existing.insert(key.clone(), value.clone());
340 }
341 }
342 }
343}
344
345fn is_list_metadata_key(key: &str) -> bool {
346 matches!(
347 key,
348 "supersedes" | "conflicts_with" | "tools" | "aliases" | "entity_aliases"
349 )
350}
351
352fn merge_metadata_list(existing: &str, incoming: &str) -> String {
353 let mut values = Vec::new();
354 for raw in existing.split(',').chain(incoming.split(',')) {
355 let value = raw.trim();
356 if !value.is_empty() && !values.iter().any(|seen| seen == value) {
357 values.push(value.to_string());
358 }
359 }
360 values.join(",")
361}
362
363fn record_duplicate_metadata(metadata: &mut HashMap<String, String>, incoming_id: &str) {
364 let count = metadata
365 .get("duplicate_count")
366 .and_then(|value| value.parse::<u32>().ok())
367 .unwrap_or(0)
368 + 1;
369 metadata.insert("duplicate_count".to_string(), count.to_string());
370 metadata.insert("last_duplicate_at".to_string(), Utc::now().to_rfc3339());
371 if !incoming_id.trim().is_empty() {
372 let duplicate_ids = metadata
373 .get("duplicate_ids")
374 .map(|existing| merge_metadata_list(existing, incoming_id))
375 .unwrap_or_else(|| incoming_id.to_string());
376 metadata.insert("duplicate_ids".to_string(), duplicate_ids);
377 }
378}
379
380fn memory_is_prune_protected(item: &MemoryItem) -> bool {
381 item.access_count >= PRUNE_PROTECTED_ACCESS_COUNT
382 || item.tags.iter().any(|tag| {
383 matches!(
384 tag.as_str(),
385 "keep" | "pinned" | "protected" | "consolidated" | "conflict"
386 )
387 })
388 || metadata_truthy(&item.metadata, "keep")
389 || metadata_truthy(&item.metadata, "pinned")
390 || metadata_truthy(&item.metadata, "protected")
391 || metadata_nonempty(&item.metadata, "supersedes")
392 || metadata_nonempty(&item.metadata, "conflicts_with")
393}
394
395fn metadata_truthy(metadata: &HashMap<String, String>, key: &str) -> bool {
396 metadata
397 .get(key)
398 .map(|value| {
399 matches!(
400 value.trim().to_ascii_lowercase().as_str(),
401 "1" | "true" | "yes" | "keep" | "pinned" | "protected"
402 )
403 })
404 .unwrap_or(false)
405}
406
407fn metadata_nonempty(metadata: &HashMap<String, String>, key: &str) -> bool {
408 metadata
409 .get(key)
410 .is_some_and(|value| !value.trim().is_empty())
411}
412
413#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
415#[serde(rename_all = "snake_case")]
416pub enum MemoryType {
417 Episodic,
418 Semantic,
419 Procedural,
420 Working,
421}
422
423#[async_trait::async_trait]
428pub trait MemoryStore: Send + Sync {
429 async fn store(&self, item: MemoryItem) -> anyhow::Result<()>;
430 async fn store_and_return(&self, item: MemoryItem) -> anyhow::Result<MemoryItem> {
436 self.store(item.clone()).await?;
437 Ok(item)
438 }
439 async fn retrieve(&self, id: &str) -> anyhow::Result<Option<MemoryItem>>;
440 async fn search(&self, query: &str, limit: usize) -> anyhow::Result<Vec<MemoryItem>>;
441 async fn search_by_tags(
442 &self,
443 tags: &[String],
444 limit: usize,
445 ) -> anyhow::Result<Vec<MemoryItem>>;
446 async fn get_recent(&self, limit: usize) -> anyhow::Result<Vec<MemoryItem>>;
447 async fn get_important(&self, threshold: f32, limit: usize) -> anyhow::Result<Vec<MemoryItem>>;
448 async fn delete(&self, id: &str) -> anyhow::Result<()>;
449 async fn clear(&self) -> anyhow::Result<()>;
450 async fn count(&self) -> anyhow::Result<usize>;
451
452 async fn prune(&self, policy: &PrunePolicy) -> anyhow::Result<usize> {
457 let _ = policy;
458 Ok(0)
459 }
460}
461
462fn index_score(entry: &IndexEntry, now: DateTime<Utc>, config: &RelevanceConfig) -> f32 {
468 let age_days = (now - entry.timestamp).num_seconds() as f32 / 86400.0;
469 let decay = (-age_days / config.decay_days).exp();
470 entry.importance * config.importance_weight + decay * config.recency_weight
471}
472
473fn sort_by_relevance(items: &mut [MemoryItem]) {
474 let now = Utc::now();
475 let config = RelevanceConfig::default();
476 items.sort_by(|a, b| {
477 b.relevance_score_at(now, &config)
478 .partial_cmp(&a.relevance_score_at(now, &config))
479 .unwrap_or(std::cmp::Ordering::Equal)
480 });
481}
482
483fn memory_type_to_query_key(memory_type: MemoryType) -> &'static str {
484 match memory_type {
485 MemoryType::Episodic => "episodic",
486 MemoryType::Semantic => "semantic",
487 MemoryType::Procedural => "procedural",
488 MemoryType::Working => "working",
489 }
490}
491
492fn query_terms(query: &str) -> Vec<String> {
493 let mut terms: Vec<String> = query
494 .to_lowercase()
495 .split(|ch: char| {
496 !(ch.is_alphanumeric() || matches!(ch, '/' | '\\' | '_' | '-' | '.' | ':' | '@'))
497 })
498 .map(str::trim)
499 .filter(|term| term.chars().count() >= 2)
500 .map(ToOwned::to_owned)
501 .collect();
502 terms.sort();
503 terms.dedup();
504 terms
505}
506
507fn lexical_match_score(
508 content_lower: &str,
509 tags: &[String],
510 memory_type: MemoryType,
511 query_lower: &str,
512 terms: &[String],
513) -> Option<f32> {
514 if query_lower.trim().is_empty() {
515 return Some(0.0);
516 }
517
518 let mut score = 0.0;
519 let mut matched_terms = 0usize;
520 if !query_lower.is_empty() && content_lower.contains(query_lower) {
521 score += 1.25;
522 }
523
524 let memory_type = memory_type_to_query_key(memory_type);
525 for term in terms {
526 let mut matched = false;
527 if content_lower.contains(term) {
528 score += 0.35;
529 matched = true;
530 }
531 if tags.iter().any(|tag| tag.to_lowercase().contains(term)) {
532 score += 0.55;
533 matched = true;
534 }
535 if memory_type.contains(term) {
536 score += 0.20;
537 matched = true;
538 }
539 if matched {
540 matched_terms += 1;
541 }
542 }
543
544 if score <= 0.0 {
545 return None;
546 }
547
548 if !terms.is_empty() {
549 score += matched_terms as f32 / terms.len() as f32;
550 }
551 Some(score)
552}
553
554fn index_search_score(
555 entry: &IndexEntry,
556 now: DateTime<Utc>,
557 config: &RelevanceConfig,
558 query_lower: &str,
559 terms: &[String],
560) -> Option<f32> {
561 let lexical = lexical_match_score(
562 &entry.content_lower,
563 &entry.tags,
564 entry.memory_type,
565 query_lower,
566 terms,
567 )?;
568 Some(index_score(entry, now, config) + lexical)
569}
570
571pub struct InMemoryStore {
579 items: RwLock<Vec<MemoryItem>>,
580}
581
582impl Default for InMemoryStore {
583 fn default() -> Self {
584 Self::new()
585 }
586}
587
588impl InMemoryStore {
589 pub fn new() -> Self {
590 Self {
591 items: RwLock::new(Vec::new()),
592 }
593 }
594}
595
596#[async_trait::async_trait]
597impl MemoryStore for InMemoryStore {
598 async fn store(&self, item: MemoryItem) -> anyhow::Result<()> {
599 self.store_and_return(item).await.map(|_| ())
600 }
601
602 async fn store_and_return(&self, item: MemoryItem) -> anyhow::Result<MemoryItem> {
603 let item = normalize_item_for_store(item);
604 let mut items = self.items.write().await;
605 if let Some(pos) = items.iter().position(|i| i.id == item.id) {
606 items[pos] = item.clone();
607 return Ok(item);
608 }
609
610 if let Some(pos) = items
611 .iter()
612 .position(|existing| memories_are_store_duplicates(existing, &item))
613 {
614 let merged = merge_duplicate_memory_item(items[pos].clone(), item);
615 items[pos] = merged.clone();
616 return Ok(merged);
617 }
618
619 items.push(item.clone());
620 Ok(item)
621 }
622
623 async fn retrieve(&self, id: &str) -> anyhow::Result<Option<MemoryItem>> {
624 let mut items = self.items.write().await;
625 let Some(item) = items.iter_mut().find(|i| i.id == id) else {
626 return Ok(None);
627 };
628 item.record_access();
629 Ok(Some(item.clone()))
630 }
631
632 async fn search(&self, query: &str, limit: usize) -> anyhow::Result<Vec<MemoryItem>> {
633 let query_lower = query.to_lowercase();
634 let config = RelevanceConfig::default();
635 let now = Utc::now();
636 let terms = query_terms(&query_lower);
637 let mut items = self.items.write().await;
638 let mut scored: Vec<(usize, f32)> = items
639 .iter()
640 .enumerate()
641 .filter_map(|(idx, item)| {
642 let lexical = lexical_match_score(
643 &item.content_lower,
644 &item.tags,
645 item.memory_type,
646 &query_lower,
647 &terms,
648 )?;
649 Some((idx, item.relevance_score_at(now, &config) + lexical))
650 })
651 .collect();
652 scored.sort_by(|a, b| {
653 b.1.partial_cmp(&a.1)
654 .unwrap_or(std::cmp::Ordering::Equal)
655 .then_with(|| items[a.0].timestamp.cmp(&items[b.0].timestamp))
656 });
657 let ids: Vec<usize> = scored.into_iter().take(limit).map(|(idx, _)| idx).collect();
658 let mut matches = Vec::with_capacity(ids.len());
659 for idx in ids {
660 items[idx].record_access();
661 matches.push(items[idx].clone());
662 }
663 Ok(matches)
664 }
665
666 async fn search_by_tags(
667 &self,
668 tags: &[String],
669 limit: usize,
670 ) -> anyhow::Result<Vec<MemoryItem>> {
671 let config = RelevanceConfig::default();
672 let now = Utc::now();
673 let mut items = self.items.write().await;
674 let mut scored: Vec<(usize, f32)> = items
675 .iter()
676 .enumerate()
677 .filter(|(_, item)| tags.iter().any(|t| item.tags.contains(t)))
678 .map(|(idx, item)| (idx, item.relevance_score_at(now, &config)))
679 .collect();
680 scored.sort_by(|a, b| {
681 b.1.partial_cmp(&a.1)
682 .unwrap_or(std::cmp::Ordering::Equal)
683 .then_with(|| items[a.0].timestamp.cmp(&items[b.0].timestamp))
684 });
685 let ids: Vec<usize> = scored.into_iter().take(limit).map(|(idx, _)| idx).collect();
686 let mut matches = Vec::with_capacity(ids.len());
687 for idx in ids {
688 items[idx].record_access();
689 matches.push(items[idx].clone());
690 }
691 Ok(matches)
692 }
693
694 async fn get_recent(&self, limit: usize) -> anyhow::Result<Vec<MemoryItem>> {
695 let items = self.items.read().await;
696 let mut sorted: Vec<MemoryItem> = items.iter().cloned().collect();
697 sorted.sort_by_key(|item| std::cmp::Reverse(item.timestamp));
698 sorted.truncate(limit);
699 Ok(sorted)
700 }
701
702 async fn get_important(&self, threshold: f32, limit: usize) -> anyhow::Result<Vec<MemoryItem>> {
703 let items = self.items.read().await;
704 let mut matches: Vec<MemoryItem> = items
705 .iter()
706 .filter(|i| i.importance >= threshold)
707 .cloned()
708 .collect();
709 matches.sort_by(|a, b| {
710 b.importance
711 .partial_cmp(&a.importance)
712 .unwrap_or(std::cmp::Ordering::Equal)
713 });
714 matches.truncate(limit);
715 Ok(matches)
716 }
717
718 async fn delete(&self, id: &str) -> anyhow::Result<()> {
719 self.items.write().await.retain(|i| i.id != id);
720 Ok(())
721 }
722
723 async fn clear(&self) -> anyhow::Result<()> {
724 self.items.write().await.clear();
725 Ok(())
726 }
727
728 async fn count(&self) -> anyhow::Result<usize> {
729 Ok(self.items.read().await.len())
730 }
731
732 async fn prune(&self, policy: &PrunePolicy) -> anyhow::Result<usize> {
733 let now = Utc::now();
734 let cutoff = now - chrono::Duration::days(policy.max_age_days as i64);
735 let min_importance = policy.min_importance_to_keep;
736
737 let mut items = self.items.write().await;
738 let before = items.len();
739
740 items.retain(|item| {
743 memory_is_prune_protected(item)
744 || item.importance >= min_importance
745 || item.timestamp >= cutoff
746 });
747
748 if policy.max_items > 0 && items.len() > policy.max_items {
751 let config = RelevanceConfig::default();
752 let protected_count = items
753 .iter()
754 .filter(|item| memory_is_prune_protected(item))
755 .count();
756 let unprotected_to_keep = policy.max_items.saturating_sub(protected_count);
757 let mut unprotected_seen = 0usize;
758 items.sort_by(|a, b| {
759 memory_is_prune_protected(b)
760 .cmp(&memory_is_prune_protected(a))
761 .then_with(|| {
762 b.relevance_score_at(now, &config)
763 .partial_cmp(&a.relevance_score_at(now, &config))
764 .unwrap_or(std::cmp::Ordering::Equal)
765 })
766 });
767 items.retain(|item| {
768 if memory_is_prune_protected(item) {
769 true
770 } else if unprotected_seen < unprotected_to_keep {
771 unprotected_seen += 1;
772 true
773 } else {
774 false
775 }
776 });
777 }
778
779 Ok(before - items.len())
780 }
781}
782
783#[derive(Debug, Clone, Serialize, Deserialize)]
788struct IndexEntry {
789 id: String,
790 content_lower: String,
791 tags: Vec<String>,
792 importance: f32,
793 timestamp: DateTime<Utc>,
794 memory_type: MemoryType,
795}
796
797impl From<&MemoryItem> for IndexEntry {
798 fn from(item: &MemoryItem) -> Self {
799 Self {
800 id: item.id.clone(),
801 content_lower: item.content.to_lowercase(),
802 tags: item.tags.clone(),
803 importance: item.importance,
804 timestamp: item.timestamp,
805 memory_type: item.memory_type,
806 }
807 }
808}
809
810pub struct FileMemoryStore {
818 items_dir: std::path::PathBuf,
819 index_path: std::path::PathBuf,
820 index: RwLock<Vec<IndexEntry>>,
821}
822
823static FILE_MEMORY_INDEX_LOCK: OnceLock<tokio::sync::Mutex<()>> = OnceLock::new();
824
825fn file_memory_index_lock() -> &'static tokio::sync::Mutex<()> {
826 FILE_MEMORY_INDEX_LOCK.get_or_init(|| tokio::sync::Mutex::new(()))
827}
828
829impl FileMemoryStore {
830 pub async fn new(dir: impl AsRef<std::path::Path>) -> anyhow::Result<Self> {
831 let dir = dir.as_ref().to_path_buf();
832 let items_dir = dir.join("items");
833 let index_path = dir.join("index.json");
834
835 tokio::fs::create_dir_all(&items_dir)
836 .await
837 .with_context(|| {
838 format!("Failed to create memory directory: {}", items_dir.display())
839 })?;
840
841 let index = if index_path.exists() {
842 let data = tokio::fs::read_to_string(&index_path)
843 .await
844 .with_context(|| {
845 format!("Failed to read memory index: {}", index_path.display())
846 })?;
847 serde_json::from_str(&data).unwrap_or_default()
848 } else {
849 Vec::new()
850 };
851
852 Ok(Self {
853 items_dir,
854 index_path,
855 index: RwLock::new(index),
856 })
857 }
858
859 fn safe_id(id: &str) -> String {
860 id.replace(['/', '\\'], "_").replace("..", "_")
861 }
862
863 fn item_path(&self, id: &str) -> std::path::PathBuf {
864 self.items_dir.join(format!("{}.json", Self::safe_id(id)))
865 }
866
867 async fn read_index_from_disk(&self) -> anyhow::Result<Vec<IndexEntry>> {
868 if !self.index_path.exists() {
869 return Ok(Vec::new());
870 }
871 let data = tokio::fs::read_to_string(&self.index_path)
872 .await
873 .with_context(|| {
874 format!("Failed to read memory index: {}", self.index_path.display())
875 })?;
876 Ok(serde_json::from_str(&data).unwrap_or_default())
877 }
878
879 async fn current_index(&self) -> Vec<IndexEntry> {
880 match self.read_index_from_disk().await {
881 Ok(index) => {
882 *self.index.write().await = index.clone();
883 index
884 }
885 Err(_) => self.index.read().await.clone(),
886 }
887 }
888
889 async fn write_index_entries(&self, index: &[IndexEntry]) -> anyhow::Result<()> {
890 let json = serde_json::to_string(index).context("Failed to serialize memory index")?;
891 let tmp = self
892 .index_path
893 .with_extension(format!("json.{}.tmp", uuid::Uuid::new_v4()));
894 tokio::fs::write(&tmp, json.as_bytes())
895 .await
896 .context("Failed to write memory index temp file")?;
897 tokio::fs::rename(&tmp, &self.index_path)
898 .await
899 .context("Failed to rename memory index")?;
900 Ok(())
901 }
902
903 async fn save_index(&self) -> anyhow::Result<()> {
904 let index = self.index.read().await.clone();
905 self.write_index_entries(&index).await
906 }
907
908 async fn save_item(&self, item: &MemoryItem) -> anyhow::Result<()> {
909 let path = self.item_path(&item.id);
910 let json = serde_json::to_string_pretty(item)
911 .with_context(|| format!("Failed to serialize memory item: {}", item.id))?;
912 let tmp = path.with_extension("json.tmp");
913 tokio::fs::write(&tmp, json.as_bytes())
914 .await
915 .with_context(|| format!("Failed to write memory item: {}", item.id))?;
916 tokio::fs::rename(&tmp, &path)
917 .await
918 .with_context(|| format!("Failed to rename memory item: {}", item.id))?;
919 Ok(())
920 }
921
922 async fn load_item_without_access(&self, id: &str) -> anyhow::Result<Option<MemoryItem>> {
923 let path = self.item_path(id);
924 if !path.exists() {
925 return Ok(None);
926 }
927 let data = tokio::fs::read_to_string(&path).await?;
928 let item: MemoryItem = serde_json::from_str(&data)?;
929 Ok(Some(normalize_item_for_store(item)))
930 }
931
932 pub async fn rebuild_index(&self) -> anyhow::Result<usize> {
934 let _guard = file_memory_index_lock().lock().await;
935 let mut entries = tokio::fs::read_dir(&self.items_dir).await?;
936 let mut new_index = Vec::new();
937 while let Some(entry) = entries.next_entry().await? {
938 let path = entry.path();
939 if path.extension().is_some_and(|ext| ext == "json") {
940 if let Ok(data) = tokio::fs::read_to_string(&path).await {
941 if let Ok(item) = serde_json::from_str::<MemoryItem>(&data) {
942 new_index.push(IndexEntry::from(&item));
943 }
944 }
945 }
946 }
947 let count = new_index.len();
948 self.write_index_entries(&new_index).await?;
949 *self.index.write().await = new_index;
950 Ok(count)
951 }
952}
953
954#[async_trait::async_trait]
955impl MemoryStore for FileMemoryStore {
956 async fn store(&self, item: MemoryItem) -> anyhow::Result<()> {
957 self.store_and_return(item).await.map(|_| ())
958 }
959
960 async fn store_and_return(&self, item: MemoryItem) -> anyhow::Result<MemoryItem> {
961 let _guard = file_memory_index_lock().lock().await;
962 let mut item = normalize_item_for_store(item);
963 item.id = Self::safe_id(&item.id);
964 let mut index = self.read_index_from_disk().await.unwrap_or_default();
965
966 let duplicate_id = index
967 .iter()
968 .find(|entry| memory_index_entry_is_duplicate(entry, &item))
969 .map(|entry| entry.id.clone());
970 if let Some(duplicate_id) = duplicate_id {
971 if duplicate_id != item.id {
972 if let Some(existing) = self.load_item_without_access(&duplicate_id).await? {
973 if memories_are_store_duplicates(&existing, &item) {
974 let merged = merge_duplicate_memory_item(existing, item.clone());
975 self.save_item(&merged).await?;
976 if item.id != merged.id {
977 let stale_path = self.item_path(&item.id);
978 if stale_path.exists() {
979 let _ = tokio::fs::remove_file(stale_path).await;
980 }
981 }
982 index.retain(|entry| entry.id != item.id && entry.id != merged.id);
983 index.push(IndexEntry::from(&merged));
984 self.write_index_entries(&index).await?;
985 *self.index.write().await = index;
986 return Ok(merged);
987 }
988 }
989 }
990 }
991
992 self.save_item(&item).await?;
993 let entry = IndexEntry::from(&item);
994 if let Some(pos) = index.iter().position(|e| e.id == item.id) {
995 index[pos] = entry;
996 } else {
997 index.push(entry);
998 }
999 self.write_index_entries(&index).await?;
1000 *self.index.write().await = index;
1001 Ok(item)
1002 }
1003
1004 async fn retrieve(&self, id: &str) -> anyhow::Result<Option<MemoryItem>> {
1005 let path = self.item_path(id);
1006 if !path.exists() {
1007 return Ok(None);
1008 }
1009 let data = tokio::fs::read_to_string(&path).await?;
1010 let mut item: MemoryItem = serde_json::from_str(&data)?;
1011 item.content_lower = item.content.to_lowercase();
1012 item.record_access();
1013 self.save_item(&item).await?;
1014 Ok(Some(item))
1015 }
1016
1017 async fn search(&self, query: &str, limit: usize) -> anyhow::Result<Vec<MemoryItem>> {
1018 let query_lower = query.to_lowercase();
1019 let index = self.current_index().await;
1020 let now = Utc::now();
1021 let config = RelevanceConfig::default();
1022 let terms = query_terms(&query_lower);
1023 let mut matches: Vec<(&IndexEntry, f32)> = index
1024 .iter()
1025 .filter_map(|e| {
1026 Some((
1027 e,
1028 index_search_score(e, now, &config, &query_lower, &terms)?,
1029 ))
1030 })
1031 .collect();
1032 matches.sort_by(|a, b| {
1033 b.1.partial_cmp(&a.1)
1034 .unwrap_or(std::cmp::Ordering::Equal)
1035 .then_with(|| b.0.timestamp.cmp(&a.0.timestamp))
1036 });
1037 let ids: Vec<String> = matches
1038 .iter()
1039 .take(limit)
1040 .map(|(e, _)| e.id.clone())
1041 .collect();
1042 let mut items = Vec::with_capacity(ids.len());
1043 for id in ids {
1044 if let Some(item) = self.retrieve(&id).await? {
1045 items.push(item);
1046 }
1047 }
1048 Ok(items)
1049 }
1050
1051 async fn search_by_tags(
1052 &self,
1053 tags: &[String],
1054 limit: usize,
1055 ) -> anyhow::Result<Vec<MemoryItem>> {
1056 let index = self.current_index().await;
1057 let now = Utc::now();
1058 let config = RelevanceConfig::default();
1059 let mut matches: Vec<&IndexEntry> = index
1060 .iter()
1061 .filter(|e| tags.iter().any(|t| e.tags.contains(t)))
1062 .collect();
1063 matches.sort_by(|a, b| {
1064 index_score(a, now, &config)
1065 .partial_cmp(&index_score(b, now, &config))
1066 .unwrap_or(std::cmp::Ordering::Equal)
1067 .reverse()
1068 });
1069 let ids: Vec<String> = matches.iter().take(limit).map(|e| e.id.clone()).collect();
1070 let mut items = Vec::with_capacity(ids.len());
1071 for id in ids {
1072 if let Some(item) = self.retrieve(&id).await? {
1073 items.push(item);
1074 }
1075 }
1076 sort_by_relevance(&mut items);
1077 Ok(items)
1078 }
1079
1080 async fn get_recent(&self, limit: usize) -> anyhow::Result<Vec<MemoryItem>> {
1081 let index = self.current_index().await;
1082 let mut sorted: Vec<&IndexEntry> = index.iter().collect();
1083 sorted.sort_by_key(|entry| std::cmp::Reverse(entry.timestamp));
1084 let ids: Vec<String> = sorted.iter().take(limit).map(|e| e.id.clone()).collect();
1085 let mut items = Vec::with_capacity(ids.len());
1086 for id in ids {
1087 if let Some(item) = self.retrieve(&id).await? {
1088 items.push(item);
1089 }
1090 }
1091 items.sort_by_key(|item| std::cmp::Reverse(item.timestamp));
1092 Ok(items)
1093 }
1094
1095 async fn get_important(&self, threshold: f32, limit: usize) -> anyhow::Result<Vec<MemoryItem>> {
1096 let index = self.current_index().await;
1097 let mut matches: Vec<&IndexEntry> =
1098 index.iter().filter(|e| e.importance >= threshold).collect();
1099 matches.sort_by(|a, b| {
1100 b.importance
1101 .partial_cmp(&a.importance)
1102 .unwrap_or(std::cmp::Ordering::Equal)
1103 });
1104 let ids: Vec<String> = matches.iter().take(limit).map(|e| e.id.clone()).collect();
1105 let mut items = Vec::with_capacity(ids.len());
1106 for id in ids {
1107 if let Some(item) = self.retrieve(&id).await? {
1108 items.push(item);
1109 }
1110 }
1111 items.sort_by(|a, b| {
1112 b.importance
1113 .partial_cmp(&a.importance)
1114 .unwrap_or(std::cmp::Ordering::Equal)
1115 });
1116 Ok(items)
1117 }
1118
1119 async fn delete(&self, id: &str) -> anyhow::Result<()> {
1120 let _guard = file_memory_index_lock().lock().await;
1121 let path = self.item_path(id);
1122 if path.exists() {
1123 tokio::fs::remove_file(&path).await?;
1124 }
1125 let mut index = self.read_index_from_disk().await.unwrap_or_default();
1126 index.retain(|e| e.id != id);
1127 self.write_index_entries(&index).await?;
1128 *self.index.write().await = index;
1129 Ok(())
1130 }
1131
1132 async fn clear(&self) -> anyhow::Result<()> {
1133 let _guard = file_memory_index_lock().lock().await;
1134 let mut entries = tokio::fs::read_dir(&self.items_dir).await?;
1135 while let Some(entry) = entries.next_entry().await? {
1136 let path = entry.path();
1137 if path.extension().is_some_and(|ext| ext == "json") {
1138 let _ = tokio::fs::remove_file(&path).await;
1139 }
1140 }
1141 self.index.write().await.clear();
1142 self.save_index().await
1143 }
1144
1145 async fn count(&self) -> anyhow::Result<usize> {
1146 Ok(self.current_index().await.len())
1147 }
1148
1149 async fn prune(&self, policy: &PrunePolicy) -> anyhow::Result<usize> {
1150 let now = Utc::now();
1151 let cutoff = now - chrono::Duration::days(policy.max_age_days as i64);
1152 let min_importance = policy.min_importance_to_keep;
1153
1154 let mut items = Vec::new();
1157 for entry in self.current_index().await {
1158 if let Some(item) = self.load_item_without_access(&entry.id).await? {
1159 items.push(item);
1160 }
1161 }
1162 let phase1_ids: Vec<String> = items
1163 .iter()
1164 .filter(|item| {
1165 !memory_is_prune_protected(item)
1166 && item.importance < min_importance
1167 && item.timestamp < cutoff
1168 })
1169 .map(|item| item.id.clone())
1170 .collect();
1171 let mut deleted = phase1_ids.len();
1172 for id in &phase1_ids {
1173 self.delete(id).await?;
1174 }
1175
1176 if policy.max_items > 0 {
1180 let config = RelevanceConfig::default();
1181 let phase2_ids: Vec<String> = {
1182 let mut remaining = Vec::new();
1183 for entry in self.current_index().await {
1184 if let Some(item) = self.load_item_without_access(&entry.id).await? {
1185 remaining.push(item);
1186 }
1187 }
1188 if remaining.len() <= policy.max_items {
1189 Vec::new()
1190 } else {
1191 let protected_count = remaining
1192 .iter()
1193 .filter(|item| memory_is_prune_protected(item))
1194 .count();
1195 let unprotected_to_keep = policy.max_items.saturating_sub(protected_count);
1196 let mut unprotected: Vec<MemoryItem> = remaining
1197 .into_iter()
1198 .filter(|item| !memory_is_prune_protected(item))
1199 .collect();
1200 unprotected.sort_by(|a, b| {
1201 b.relevance_score_at(now, &config)
1202 .partial_cmp(&a.relevance_score_at(now, &config))
1203 .unwrap_or(std::cmp::Ordering::Equal)
1204 });
1205 unprotected
1206 .into_iter()
1207 .skip(unprotected_to_keep)
1208 .map(|item| item.id)
1209 .collect()
1210 }
1211 };
1212 deleted += phase2_ids.len();
1213 for id in &phase2_ids {
1214 self.delete(id).await?;
1215 }
1216 }
1217
1218 Ok(deleted)
1219 }
1220}
1221
1222#[cfg(test)]
1227mod tests {
1228 use super::*;
1229
1230 #[test]
1233 fn test_memory_item_creation() {
1234 let item = MemoryItem::new("Test memory")
1235 .with_importance(0.8)
1236 .with_tag("test")
1237 .with_type(MemoryType::Semantic);
1238 assert_eq!(item.content, "Test memory");
1239 assert_eq!(item.importance, 0.8);
1240 assert_eq!(item.tags, vec!["test"]);
1241 assert_eq!(item.memory_type, MemoryType::Semantic);
1242 }
1243
1244 #[test]
1245 fn test_memory_item_importance_clamped() {
1246 assert_eq!(MemoryItem::new("x").with_importance(1.5).importance, 1.0);
1247 assert_eq!(MemoryItem::new("x").with_importance(-0.5).importance, 0.0);
1248 }
1249
1250 #[test]
1251 fn test_memory_item_record_access() {
1252 let mut item = MemoryItem::new("test");
1253 assert_eq!(item.access_count, 0);
1254 item.record_access();
1255 assert_eq!(item.access_count, 1);
1256 assert!(item.last_accessed.is_some());
1257 }
1258
1259 #[test]
1260 fn test_memory_item_merge_duplicate_preserves_canonical_id() {
1261 let existing = MemoryItem::new("Run focused memory store tests after parser changes.")
1262 .with_importance(0.4)
1263 .with_tag("memory")
1264 .with_metadata("source", "workflow");
1265 let existing_id = existing.id.clone();
1266 let incoming =
1267 MemoryItem::new("Run focused memory store regression tests after parser changes.")
1268 .with_importance(0.9)
1269 .with_tag("tests")
1270 .with_metadata("supersedes", "old-memory");
1271
1272 let merged = existing.merge_duplicate(incoming);
1273
1274 assert_eq!(merged.id, existing_id);
1275 assert!(merged.content.contains("regression tests"));
1276 assert_eq!(merged.importance, 0.9);
1277 assert!(merged.tags.contains(&"memory".to_string()));
1278 assert!(merged.tags.contains(&"tests".to_string()));
1279 assert_eq!(
1280 merged.metadata.get("duplicate_count").map(String::as_str),
1281 Some("1")
1282 );
1283 assert_eq!(
1284 merged.metadata.get("supersedes").map(String::as_str),
1285 Some("old-memory")
1286 );
1287 }
1288
1289 #[test]
1290 fn test_content_fingerprint_normalizes_only_case_and_whitespace() {
1291 let canonical = MemoryItem::new("Use C++ for the native client implementation.");
1292 let equivalent = MemoryItem::new(" use C++ for the native client implementation. ");
1293 let distinct = MemoryItem::new("Use C for the native client implementation.");
1294
1295 assert_eq!(
1296 canonical.content_fingerprint(),
1297 equivalent.content_fingerprint()
1298 );
1299 assert_ne!(
1300 canonical.content_fingerprint(),
1301 distinct.content_fingerprint()
1302 );
1303 }
1304
1305 #[test]
1306 fn test_memory_item_default_type_is_episodic() {
1307 assert_eq!(MemoryItem::new("test").memory_type, MemoryType::Episodic);
1308 }
1309
1310 #[test]
1311 fn test_memory_item_all_types() {
1312 assert_eq!(
1313 MemoryItem::new("e")
1314 .with_type(MemoryType::Episodic)
1315 .memory_type,
1316 MemoryType::Episodic
1317 );
1318 assert_eq!(
1319 MemoryItem::new("s")
1320 .with_type(MemoryType::Semantic)
1321 .memory_type,
1322 MemoryType::Semantic
1323 );
1324 assert_eq!(
1325 MemoryItem::new("p")
1326 .with_type(MemoryType::Procedural)
1327 .memory_type,
1328 MemoryType::Procedural
1329 );
1330 assert_eq!(
1331 MemoryItem::new("w")
1332 .with_type(MemoryType::Working)
1333 .memory_type,
1334 MemoryType::Working
1335 );
1336 }
1337
1338 #[test]
1341 fn test_relevance_score_uses_config() {
1342 let item = MemoryItem::new("test").with_importance(1.0);
1343 let now = Utc::now();
1344
1345 let config_importance = RelevanceConfig {
1347 decay_days: 30.0,
1348 importance_weight: 0.9,
1349 recency_weight: 0.1,
1350 };
1351 let score = item.relevance_score_at(now, &config_importance);
1352 assert!(score > 0.95, "score was {score}");
1353
1354 let config_fast_decay = RelevanceConfig {
1356 decay_days: 1.0,
1357 importance_weight: 0.7,
1358 recency_weight: 0.3,
1359 };
1360 let score2 = item.relevance_score_at(now, &config_fast_decay);
1361 assert!(score2 > 0.9, "score was {score2}");
1362 }
1363
1364 #[test]
1365 fn test_relevance_score_decays_with_age() {
1366 let mut old_item = MemoryItem::new("old").with_importance(0.5);
1367 old_item.timestamp = Utc::now() - chrono::Duration::days(60);
1368 let config = RelevanceConfig::default(); let score = old_item.relevance_score_at(Utc::now(), &config);
1370 assert!(score < 0.45, "score was {score}");
1373 }
1374
1375 #[test]
1376 fn test_relevance_score_default_uses_default_config() {
1377 let item = MemoryItem::new("test").with_importance(0.9);
1378 let score = item.relevance_score();
1379 assert!(score > 0.6);
1380 }
1381
1382 #[test]
1385 fn test_relevance_config_defaults() {
1386 let c = RelevanceConfig::default();
1387 assert_eq!(c.decay_days, 30.0);
1388 assert_eq!(c.importance_weight, 0.7);
1389 assert_eq!(c.recency_weight, 0.3);
1390 }
1391
1392 #[tokio::test]
1395 async fn test_in_memory_store_retrieve() {
1396 let store = InMemoryStore::new();
1397 let item = MemoryItem::new("hello").with_tag("test");
1398 store.store(item.clone()).await.unwrap();
1399 let r = store.retrieve(&item.id).await.unwrap();
1400 assert!(r.is_some());
1401 assert_eq!(r.unwrap().content, "hello");
1402 let r = store.retrieve(&item.id).await.unwrap().unwrap();
1403 assert_eq!(r.access_count, 2);
1404 }
1405
1406 #[tokio::test]
1407 async fn test_in_memory_store_retrieve_nonexistent() {
1408 let store = InMemoryStore::new();
1409 assert!(store.retrieve("nope").await.unwrap().is_none());
1410 }
1411
1412 #[tokio::test]
1413 async fn test_in_memory_store_upsert() {
1414 let store = InMemoryStore::new();
1415 let mut item = MemoryItem::new("original");
1416 let id = item.id.clone();
1417 store.store(item.clone()).await.unwrap();
1418 item.content = "updated".to_string();
1419 item.content_lower = "updated".to_string();
1420 store.store(item).await.unwrap();
1421 assert_eq!(store.count().await.unwrap(), 1);
1422 assert_eq!(
1423 store.retrieve(&id).await.unwrap().unwrap().content,
1424 "updated"
1425 );
1426 }
1427
1428 #[tokio::test]
1429 async fn test_in_memory_store_search_and_tags() {
1430 let store = InMemoryStore::new();
1431 store
1432 .store(MemoryItem::new("create file").with_tag("file"))
1433 .await
1434 .unwrap();
1435 store
1436 .store(MemoryItem::new("delete file").with_tag("file"))
1437 .await
1438 .unwrap();
1439 store
1440 .store(MemoryItem::new("create dir").with_tag("dir"))
1441 .await
1442 .unwrap();
1443 assert_eq!(store.search("create", 10).await.unwrap().len(), 2);
1444 assert_eq!(
1445 store
1446 .search_by_tags(&["file".to_string()], 10)
1447 .await
1448 .unwrap()
1449 .len(),
1450 2
1451 );
1452 }
1453
1454 #[tokio::test]
1455 async fn test_in_memory_store_search_relevance_order() {
1456 let store = InMemoryStore::new();
1457 store
1458 .store(MemoryItem::new("rust tip").with_importance(0.3))
1459 .await
1460 .unwrap();
1461 store
1462 .store(MemoryItem::new("rust trick").with_importance(0.9))
1463 .await
1464 .unwrap();
1465 let results = store.search("rust", 10).await.unwrap();
1466 assert_eq!(results.len(), 2);
1467 assert!(results[0].importance >= results[1].importance);
1468 }
1469
1470 #[tokio::test]
1471 async fn test_in_memory_store_delete_and_clear() {
1472 let store = InMemoryStore::new();
1473 let item = MemoryItem::new("to delete");
1474 let id = item.id.clone();
1475 store.store(item).await.unwrap();
1476 store.delete(&id).await.unwrap();
1477 assert_eq!(store.count().await.unwrap(), 0);
1478
1479 for i in 0..3 {
1480 store
1481 .store(MemoryItem::new(format!("item {i}")))
1482 .await
1483 .unwrap();
1484 }
1485 store.clear().await.unwrap();
1486 assert_eq!(store.count().await.unwrap(), 0);
1487 }
1488
1489 #[tokio::test]
1490 async fn test_in_memory_store_get_recent() {
1491 let store = InMemoryStore::new();
1492 for i in 0..5 {
1493 let mut item = MemoryItem::new(format!("item {i}"));
1494 item.timestamp = Utc::now() + chrono::Duration::seconds(i as i64);
1495 store.store(item).await.unwrap();
1496 }
1497 let recent = store.get_recent(3).await.unwrap();
1498 assert_eq!(recent.len(), 3);
1499 assert!(recent[0].timestamp >= recent[1].timestamp);
1500 }
1501
1502 #[tokio::test]
1503 async fn test_in_memory_store_get_important() {
1504 let store = InMemoryStore::new();
1505 store
1506 .store(MemoryItem::new("low").with_importance(0.2))
1507 .await
1508 .unwrap();
1509 store
1510 .store(MemoryItem::new("high").with_importance(0.9))
1511 .await
1512 .unwrap();
1513 store
1514 .store(MemoryItem::new("medium").with_importance(0.5))
1515 .await
1516 .unwrap();
1517 let results = store.get_important(0.7, 10).await.unwrap();
1518 assert_eq!(results.len(), 1);
1519 assert_eq!(results[0].content, "high");
1520 }
1521
1522 #[test]
1523 fn test_in_memory_store_default() {
1524 let _store: InMemoryStore = InMemoryStore::default();
1525 }
1526
1527 #[test]
1530 fn test_prune_policy_defaults() {
1531 let p = PrunePolicy::default();
1532 assert_eq!(p.max_age_days, 90);
1533 assert_eq!(p.min_importance_to_keep, 0.5);
1534 assert_eq!(p.max_items, 0);
1535 }
1536
1537 #[tokio::test]
1538 async fn test_prune_removes_old_low_importance() {
1539 let store = InMemoryStore::new();
1540 let mut old_item = MemoryItem::new("stale memory").with_importance(0.2);
1541 old_item.timestamp = Utc::now() - chrono::Duration::days(100);
1542 store.store(old_item).await.unwrap();
1543
1544 let policy = PrunePolicy {
1545 max_age_days: 90,
1546 min_importance_to_keep: 0.5,
1547 max_items: 0,
1548 };
1549 let deleted = store.prune(&policy).await.unwrap();
1550 assert_eq!(deleted, 1);
1551 assert_eq!(store.count().await.unwrap(), 0);
1552 }
1553
1554 #[tokio::test]
1555 async fn test_prune_keeps_high_importance() {
1556 let store = InMemoryStore::new();
1557 let mut old_item = MemoryItem::new("important memory").with_importance(0.9);
1558 old_item.timestamp = Utc::now() - chrono::Duration::days(100);
1559 store.store(old_item).await.unwrap();
1560
1561 let policy = PrunePolicy {
1562 max_age_days: 90,
1563 min_importance_to_keep: 0.5,
1564 max_items: 0,
1565 };
1566 let deleted = store.prune(&policy).await.unwrap();
1567 assert_eq!(deleted, 0);
1568 assert_eq!(store.count().await.unwrap(), 1);
1569 }
1570
1571 #[tokio::test]
1572 async fn test_prune_max_items() {
1573 let store = InMemoryStore::new();
1574 for i in 0..10 {
1575 store
1576 .store(MemoryItem::new(format!("item {i}")).with_importance(i as f32 * 0.1))
1577 .await
1578 .unwrap();
1579 }
1580 let policy = PrunePolicy {
1581 max_age_days: 9999,
1582 min_importance_to_keep: 0.0,
1583 max_items: 5,
1584 };
1585 let deleted = store.prune(&policy).await.unwrap();
1586 assert_eq!(deleted, 5);
1587 assert_eq!(store.count().await.unwrap(), 5);
1588 }
1589
1590 #[tokio::test]
1591 async fn test_prune_keeps_recent_low_importance() {
1592 let store = InMemoryStore::new();
1594 store
1595 .store(MemoryItem::new("fresh").with_importance(0.1))
1596 .await
1597 .unwrap();
1598
1599 let policy = PrunePolicy {
1600 max_age_days: 90,
1601 min_importance_to_keep: 0.5,
1602 max_items: 0,
1603 };
1604 let deleted = store.prune(&policy).await.unwrap();
1605 assert_eq!(deleted, 0);
1606 assert_eq!(store.count().await.unwrap(), 1);
1607 }
1608}
1609
1610#[cfg(test)]
1611mod file_memory_store_tests {
1612 use super::*;
1613 use tempfile::TempDir;
1614
1615 async fn setup() -> (TempDir, FileMemoryStore) {
1616 let dir = TempDir::new().unwrap();
1617 let store = FileMemoryStore::new(dir.path()).await.unwrap();
1618 (dir, store)
1619 }
1620
1621 #[tokio::test]
1622 async fn test_store_and_retrieve() {
1623 let (_dir, store) = setup().await;
1624 let item = MemoryItem::new("hello world");
1625 let id = item.id.clone();
1626 store.store(item).await.unwrap();
1627 let r = store.retrieve(&id).await.unwrap().unwrap();
1628 assert_eq!(r.content, "hello world");
1629 }
1630
1631 #[tokio::test]
1632 async fn test_retrieve_nonexistent() {
1633 let (_dir, store) = setup().await;
1634 assert!(store.retrieve("nonexistent").await.unwrap().is_none());
1635 }
1636
1637 #[tokio::test]
1638 async fn test_search_by_content() {
1639 let (_dir, store) = setup().await;
1640 store
1641 .store(MemoryItem::new("rust programming"))
1642 .await
1643 .unwrap();
1644 store
1645 .store(MemoryItem::new("python scripting"))
1646 .await
1647 .unwrap();
1648 store
1649 .store(MemoryItem::new("rust async patterns"))
1650 .await
1651 .unwrap();
1652 let results = store.search("rust", 10).await.unwrap();
1653 assert_eq!(results.len(), 2);
1654 }
1655
1656 #[tokio::test]
1657 async fn test_search_matches_non_contiguous_terms() {
1658 let (_dir, store) = setup().await;
1659 store
1660 .store(MemoryItem::new(
1661 "Success: release preflight\nTools: bash\nResult: provider verification passed",
1662 ))
1663 .await
1664 .unwrap();
1665 let results = store.search("release provider check", 10).await.unwrap();
1666 assert_eq!(results.len(), 1);
1667 assert!(results[0].content.contains("release preflight"));
1668 }
1669
1670 #[tokio::test]
1671 async fn test_retrieve_records_access_on_disk() {
1672 let (dir, store) = setup().await;
1673 let item = MemoryItem::new("access me");
1674 let id = item.id.clone();
1675 store.store(item).await.unwrap();
1676
1677 let item = store.retrieve(&id).await.unwrap().unwrap();
1678 assert_eq!(item.access_count, 1);
1679
1680 let reopened = FileMemoryStore::new(dir.path()).await.unwrap();
1681 let item = reopened.retrieve(&id).await.unwrap().unwrap();
1682 assert_eq!(item.access_count, 2);
1683 assert!(item.last_accessed.is_some());
1684 }
1685
1686 #[tokio::test]
1687 async fn test_search_limit() {
1688 let (_dir, store) = setup().await;
1689 for i in 0..10 {
1690 store
1691 .store(MemoryItem::new(format!("item {i}")))
1692 .await
1693 .unwrap();
1694 }
1695 assert_eq!(store.search("item", 3).await.unwrap().len(), 3);
1696 }
1697
1698 #[tokio::test]
1699 async fn test_search_by_tags() {
1700 let (_dir, store) = setup().await;
1701 store
1702 .store(MemoryItem::new("one").with_tags(vec!["rust".into(), "async".into()]))
1703 .await
1704 .unwrap();
1705 store
1706 .store(MemoryItem::new("two").with_tags(vec!["python".into()]))
1707 .await
1708 .unwrap();
1709 store
1710 .store(MemoryItem::new("three").with_tags(vec!["rust".into()]))
1711 .await
1712 .unwrap();
1713 assert_eq!(
1714 store
1715 .search_by_tags(&["rust".to_string()], 10)
1716 .await
1717 .unwrap()
1718 .len(),
1719 2
1720 );
1721 }
1722
1723 #[tokio::test]
1724 async fn test_get_recent_ordered() {
1725 let (_dir, store) = setup().await;
1726 for i in 0..5 {
1727 let mut item = MemoryItem::new(format!("item {i}"));
1728 item.timestamp = Utc::now() + chrono::Duration::seconds(i as i64);
1729 store.store(item).await.unwrap();
1730 }
1731 let results = store.get_recent(3).await.unwrap();
1732 assert_eq!(results.len(), 3);
1733 assert!(results[0].timestamp >= results[1].timestamp);
1734 }
1735
1736 #[tokio::test]
1737 async fn test_get_important() {
1738 let (_dir, store) = setup().await;
1739 store
1740 .store(MemoryItem::new("low").with_importance(0.1))
1741 .await
1742 .unwrap();
1743 store
1744 .store(MemoryItem::new("high").with_importance(0.9))
1745 .await
1746 .unwrap();
1747 store
1748 .store(MemoryItem::new("medium").with_importance(0.5))
1749 .await
1750 .unwrap();
1751 let results = store.get_important(0.0, 2).await.unwrap();
1752 assert_eq!(results.len(), 2);
1753 assert!(results[0].importance >= results[1].importance);
1754 }
1755
1756 #[tokio::test]
1757 async fn test_delete() {
1758 let (_dir, store) = setup().await;
1759 let item = MemoryItem::new("to delete");
1760 let id = item.id.clone();
1761 store.store(item).await.unwrap();
1762 store.delete(&id).await.unwrap();
1763 assert_eq!(store.count().await.unwrap(), 0);
1764 assert!(store.retrieve(&id).await.unwrap().is_none());
1765 }
1766
1767 #[tokio::test]
1768 async fn test_delete_nonexistent() {
1769 let (_dir, store) = setup().await;
1770 store.delete("nonexistent").await.unwrap();
1771 }
1772
1773 #[tokio::test]
1774 async fn test_clear() {
1775 let (_dir, store) = setup().await;
1776 for i in 0..5 {
1777 store
1778 .store(MemoryItem::new(format!("item {i}")))
1779 .await
1780 .unwrap();
1781 }
1782 store.clear().await.unwrap();
1783 assert_eq!(store.count().await.unwrap(), 0);
1784 }
1785
1786 #[tokio::test]
1787 async fn test_persistence_across_instances() {
1788 let dir = TempDir::new().unwrap();
1789 {
1790 let store = FileMemoryStore::new(dir.path()).await.unwrap();
1791 store
1792 .store(MemoryItem::new("persistent data").with_tags(vec!["test".into()]))
1793 .await
1794 .unwrap();
1795 }
1796 {
1797 let store = FileMemoryStore::new(dir.path()).await.unwrap();
1798 assert_eq!(store.count().await.unwrap(), 1);
1799 assert_eq!(store.search("persistent", 10).await.unwrap().len(), 1);
1800 }
1801 }
1802
1803 #[tokio::test]
1804 async fn test_stale_instances_merge_index_on_store() {
1805 let dir = TempDir::new().unwrap();
1806 let store_a = FileMemoryStore::new(dir.path()).await.unwrap();
1807 let store_b = FileMemoryStore::new(dir.path()).await.unwrap();
1808
1809 store_a
1810 .store(MemoryItem::new("alpha stale merge"))
1811 .await
1812 .unwrap();
1813 store_b
1814 .store(MemoryItem::new("beta stale merge"))
1815 .await
1816 .unwrap();
1817
1818 let reopened = FileMemoryStore::new(dir.path()).await.unwrap();
1819 assert_eq!(reopened.count().await.unwrap(), 2);
1820 assert_eq!(reopened.search("stale merge", 10).await.unwrap().len(), 2);
1821 }
1822
1823 #[tokio::test]
1824 async fn test_rebuild_index() {
1825 let dir = TempDir::new().unwrap();
1826 {
1827 let store = FileMemoryStore::new(dir.path()).await.unwrap();
1828 store.store(MemoryItem::new("alpha")).await.unwrap();
1829 store.store(MemoryItem::new("beta")).await.unwrap();
1830 }
1831 tokio::fs::remove_file(dir.path().join("index.json"))
1832 .await
1833 .unwrap();
1834 {
1835 let store = FileMemoryStore::new(dir.path()).await.unwrap();
1836 assert_eq!(store.count().await.unwrap(), 0);
1837 store.rebuild_index().await.unwrap();
1838 assert_eq!(store.count().await.unwrap(), 2);
1839 }
1840 }
1841
1842 #[tokio::test]
1843 async fn test_path_traversal_prevention() {
1844 let (_dir, store) = setup().await;
1845 let mut item = MemoryItem::new("sneaky");
1846 item.id = "../../../etc/passwd".to_string();
1847 store.store(item).await.unwrap();
1848 let results = store.search("sneaky", 10).await.unwrap();
1849 assert_eq!(results.len(), 1);
1850 assert!(!results[0].id.contains('/'));
1851 assert!(!results[0].id.contains(".."));
1852 }
1853
1854 #[tokio::test]
1855 async fn test_importance_threshold() {
1856 let (_dir, store) = setup().await;
1857 store
1858 .store(MemoryItem::new("low").with_importance(0.2))
1859 .await
1860 .unwrap();
1861 store
1862 .store(MemoryItem::new("high").with_importance(0.8))
1863 .await
1864 .unwrap();
1865 let results = store.get_important(0.5, 10).await.unwrap();
1866 assert_eq!(results.len(), 1);
1867 assert_eq!(results[0].content, "high");
1868 }
1869
1870 #[tokio::test]
1871 async fn test_file_prune_removes_old_low_importance() {
1872 let (_dir, store) = setup().await;
1873 let mut old_item = MemoryItem::new("stale").with_importance(0.2);
1874 old_item.timestamp = Utc::now() - chrono::Duration::days(100);
1875 store.store(old_item).await.unwrap();
1876
1877 let policy = PrunePolicy {
1878 max_age_days: 90,
1879 min_importance_to_keep: 0.5,
1880 max_items: 0,
1881 };
1882 let deleted = store.prune(&policy).await.unwrap();
1883 assert_eq!(deleted, 1);
1884 assert_eq!(store.count().await.unwrap(), 0);
1885 }
1886
1887 #[tokio::test]
1888 async fn test_file_prune_keeps_high_importance() {
1889 let (_dir, store) = setup().await;
1890 let mut old_item = MemoryItem::new("important").with_importance(0.9);
1891 old_item.timestamp = Utc::now() - chrono::Duration::days(100);
1892 store.store(old_item).await.unwrap();
1893
1894 let policy = PrunePolicy {
1895 max_age_days: 90,
1896 min_importance_to_keep: 0.5,
1897 max_items: 0,
1898 };
1899 let deleted = store.prune(&policy).await.unwrap();
1900 assert_eq!(deleted, 0);
1901 assert_eq!(store.count().await.unwrap(), 1);
1902 }
1903
1904 #[tokio::test]
1905 async fn test_file_prune_max_items() {
1906 let (_dir, store) = setup().await;
1907 for i in 0..10 {
1908 store
1909 .store(MemoryItem::new(format!("item {i}")).with_importance(i as f32 * 0.1))
1910 .await
1911 .unwrap();
1912 }
1913 let policy = PrunePolicy {
1914 max_age_days: 9999,
1915 min_importance_to_keep: 0.0,
1916 max_items: 5,
1917 };
1918 let deleted = store.prune(&policy).await.unwrap();
1919 assert_eq!(deleted, 5);
1920 assert_eq!(store.count().await.unwrap(), 5);
1921 }
1922}