Skip to main content

a3s_memory/
lib.rs

1//! A3S Memory — pluggable memory storage for AI agents.
2//!
3//! Provides the `MemoryStore` trait, `MemoryItem`, `MemoryType`,
4//! configuration types, and the following backend implementations:
5//!
6//! | Backend | Feature | Search |
7//! |---------|---------|--------|
8//! | [`FileMemoryStore`] | always available | substring |
9//! | [`SqliteMemoryStore`] | `sqlite` | BM25 via FTS5 |
10//!
11//! The [`FileMemoryStore`] is the default and requires no additional dependencies.
12
13/// Evidence-backed, namespace-isolated durable memory repository.
14///
15/// This additive V2 API keeps retrieval pure and applies bounded multi-node
16/// changes atomically. The existing [`MemoryStore`] API remains available for
17/// backwards compatibility.
18pub mod repository;
19
20/// Vector indexes for caller-owned semantic retrieval.
21///
22/// This capability is independent from [`MemoryStore`]: callers own document
23/// admission, embedding generation, lifecycle, and result fusion.
24pub 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/// SQLite-backed memory store with dual-track Markdown export.
37///
38/// Requires the `sqlite` Cargo feature.
39#[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// ============================================================================
56// Configuration
57// ============================================================================
58
59/// Configuration for relevance scoring
60#[derive(Debug, Clone, Serialize, Deserialize)]
61#[serde(rename_all = "camelCase")]
62pub struct RelevanceConfig {
63    /// Exponential decay half-life in days (default: 30.0)
64    #[serde(default = "RelevanceConfig::default_decay_days")]
65    pub decay_days: f32,
66    /// Weight for importance factor (default: 0.7)
67    #[serde(default = "RelevanceConfig::default_importance_weight")]
68    pub importance_weight: f32,
69    /// Weight for recency factor (default: 0.3)
70    #[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/// Policy controlling automatic pruning of long-term memory.
97#[derive(Debug, Clone, Serialize, Deserialize)]
98#[serde(rename_all = "camelCase")]
99pub struct PrunePolicy {
100    /// Items older than this many days AND below `min_importance_to_keep` are deleted (default: 90).
101    #[serde(default = "PrunePolicy::default_max_age_days")]
102    pub max_age_days: u32,
103    /// Items with importance below this threshold are eligible for age-based deletion (default: 0.5).
104    /// High-importance items are never age-pruned.
105    #[serde(default = "PrunePolicy::default_min_importance_to_keep")]
106    pub min_importance_to_keep: f32,
107    /// Hard cap on total items; when exceeded, lowest-relevance items are removed.
108    /// 0 means unlimited (default: 0).
109    #[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// ============================================================================
133// Memory Item
134// ============================================================================
135
136/// A single memory item
137#[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    /// Stable, case- and whitespace-normalized content fingerprint used by
196    /// default stores to collapse exact durable duplicates.
197    ///
198    /// Very short memories return `None` so generic fragments such as "ok" or
199    /// "done" are not accidentally merged.
200    pub fn content_fingerprint(&self) -> Option<String> {
201        memory_content_fingerprint(&self.content)
202    }
203
204    /// Merge a later duplicate observation into this canonical memory item.
205    ///
206    /// The canonical id is preserved while importance, tags, list-style
207    /// metadata, access stats, and duplicate audit metadata are consolidated.
208    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    /// Calculate relevance score at a given timestamp using the provided config
218    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    /// Calculate relevance score with default config
225    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/// Type of memory
414#[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// ============================================================================
424// Memory Store Trait
425// ============================================================================
426
427#[async_trait::async_trait]
428pub trait MemoryStore: Send + Sync {
429    async fn store(&self, item: MemoryItem) -> anyhow::Result<()>;
430    /// Store a memory and return the canonical item that now represents it.
431    ///
432    /// Default stores may merge durable duplicate content into an existing item
433    /// and return that existing item with updated metadata. Custom backends that
434    /// do not implement deduplication can rely on this default wrapper.
435    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    /// Remove stale or excess items according to `policy`.
453    ///
454    /// Returns the number of items deleted. The default implementation is a
455    /// no-op (returns 0) for backwards compatibility.
456    async fn prune(&self, policy: &PrunePolicy) -> anyhow::Result<usize> {
457        let _ = policy;
458        Ok(0)
459    }
460}
461
462// ============================================================================
463// Shared helpers
464// ============================================================================
465
466/// Score an index entry for sorting (avoids loading full MemoryItem from disk)
467fn 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
571// ============================================================================
572// In-Memory Store
573// ============================================================================
574
575/// In-memory `MemoryStore` implementation.
576///
577/// Useful for testing and ephemeral (non-persistent) use cases.
578pub 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        // Phase 1: remove items that are old and below the importance threshold,
741        // unless they have explicit curation/provenance protection.
742        items.retain(|item| {
743            memory_is_prune_protected(item)
744                || item.importance >= min_importance
745                || item.timestamp >= cutoff
746        });
747
748        // Phase 2: if still over the cap, keep protected memories first and then
749        // the highest-relevance unprotected items.
750        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// ============================================================================
784// File-Based Memory Store
785// ============================================================================
786
787#[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
810/// File-based memory store with atomic writes and in-memory index.
811///
812/// ```text
813/// memory_dir/
814///   index.json
815///   items/{id}.json
816/// ```
817pub 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    /// Rebuild the index from item files on disk (useful for corruption recovery).
933    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        // Phase 1: collect IDs that are old and low-importance, unless the full
1155        // item carries curation/provenance protection.
1156        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        // Phase 2: enforce max_items cap by removing lowest-relevance
1177        // unprotected items. Protected items are hard-exempt, so the store may
1178        // remain above the cap if the user pinned more memories than the cap.
1179        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// ============================================================================
1223// Tests
1224// ============================================================================
1225
1226#[cfg(test)]
1227mod tests {
1228    use super::*;
1229
1230    // MemoryItem tests
1231
1232    #[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    // relevance_score_at tests
1339
1340    #[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        // High importance weight → score dominated by importance
1346        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        // Short decay → recent item still scores well
1355        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(); // 30-day half-life
1369        let score = old_item.relevance_score_at(Utc::now(), &config);
1370        // After 60 days (2 half-lives), decay ≈ exp(-2) ≈ 0.135
1371        // score ≈ 0.5*0.7 + 0.135*0.3 ≈ 0.39
1372        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    // RelevanceConfig tests
1383
1384    #[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    // InMemoryStore tests
1393
1394    #[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    // PrunePolicy / prune() tests
1528
1529    #[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        // A recent item with low importance should NOT be pruned (not old enough)
1593        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}