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