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/// Ephemeral vector indexes for caller-owned semantic retrieval.
14///
15/// This capability is independent from [`MemoryStore`]: callers own document
16/// admission, embedding generation, lifecycle, and result fusion.
17pub 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/// SQLite-backed memory store with dual-track Markdown export.
26///
27/// Requires the `sqlite` Cargo feature.
28#[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// ============================================================================
45// Configuration
46// ============================================================================
47
48/// Configuration for relevance scoring
49#[derive(Debug, Clone, Serialize, Deserialize)]
50#[serde(rename_all = "camelCase")]
51pub struct RelevanceConfig {
52    /// Exponential decay half-life in days (default: 30.0)
53    #[serde(default = "RelevanceConfig::default_decay_days")]
54    pub decay_days: f32,
55    /// Weight for importance factor (default: 0.7)
56    #[serde(default = "RelevanceConfig::default_importance_weight")]
57    pub importance_weight: f32,
58    /// Weight for recency factor (default: 0.3)
59    #[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/// Policy controlling automatic pruning of long-term memory.
86#[derive(Debug, Clone, Serialize, Deserialize)]
87#[serde(rename_all = "camelCase")]
88pub struct PrunePolicy {
89    /// Items older than this many days AND below `min_importance_to_keep` are deleted (default: 90).
90    #[serde(default = "PrunePolicy::default_max_age_days")]
91    pub max_age_days: u32,
92    /// Items with importance below this threshold are eligible for age-based deletion (default: 0.5).
93    /// High-importance items are never age-pruned.
94    #[serde(default = "PrunePolicy::default_min_importance_to_keep")]
95    pub min_importance_to_keep: f32,
96    /// Hard cap on total items; when exceeded, lowest-relevance items are removed.
97    /// 0 means unlimited (default: 0).
98    #[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// ============================================================================
122// Memory Item
123// ============================================================================
124
125/// A single memory item
126#[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    /// Stable, case- and whitespace-normalized content fingerprint used by
185    /// default stores to collapse exact durable duplicates.
186    ///
187    /// Very short memories return `None` so generic fragments such as "ok" or
188    /// "done" are not accidentally merged.
189    pub fn content_fingerprint(&self) -> Option<String> {
190        memory_content_fingerprint(&self.content)
191    }
192
193    /// Merge a later duplicate observation into this canonical memory item.
194    ///
195    /// The canonical id is preserved while importance, tags, list-style
196    /// metadata, access stats, and duplicate audit metadata are consolidated.
197    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    /// Calculate relevance score at a given timestamp using the provided config
207    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    /// Calculate relevance score with default config
214    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/// Type of memory
403#[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// ============================================================================
413// Memory Store Trait
414// ============================================================================
415
416#[async_trait::async_trait]
417pub trait MemoryStore: Send + Sync {
418    async fn store(&self, item: MemoryItem) -> anyhow::Result<()>;
419    /// Store a memory and return the canonical item that now represents it.
420    ///
421    /// Default stores may merge durable duplicate content into an existing item
422    /// and return that existing item with updated metadata. Custom backends that
423    /// do not implement deduplication can rely on this default wrapper.
424    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    /// Remove stale or excess items according to `policy`.
442    ///
443    /// Returns the number of items deleted. The default implementation is a
444    /// no-op (returns 0) for backwards compatibility.
445    async fn prune(&self, policy: &PrunePolicy) -> anyhow::Result<usize> {
446        let _ = policy;
447        Ok(0)
448    }
449}
450
451// ============================================================================
452// Shared helpers
453// ============================================================================
454
455/// Score an index entry for sorting (avoids loading full MemoryItem from disk)
456fn 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
560// ============================================================================
561// In-Memory Store
562// ============================================================================
563
564/// In-memory `MemoryStore` implementation.
565///
566/// Useful for testing and ephemeral (non-persistent) use cases.
567pub 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        // Phase 1: remove items that are old and below the importance threshold,
730        // unless they have explicit curation/provenance protection.
731        items.retain(|item| {
732            memory_is_prune_protected(item)
733                || item.importance >= min_importance
734                || item.timestamp >= cutoff
735        });
736
737        // Phase 2: if still over the cap, keep protected memories first and then
738        // the highest-relevance unprotected items.
739        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// ============================================================================
773// File-Based Memory Store
774// ============================================================================
775
776#[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
799/// File-based memory store with atomic writes and in-memory index.
800///
801/// ```text
802/// memory_dir/
803///   index.json
804///   items/{id}.json
805/// ```
806pub 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    /// Rebuild the index from item files on disk (useful for corruption recovery).
922    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        // Phase 1: collect IDs that are old and low-importance, unless the full
1144        // item carries curation/provenance protection.
1145        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        // Phase 2: enforce max_items cap by removing lowest-relevance
1166        // unprotected items. Protected items are hard-exempt, so the store may
1167        // remain above the cap if the user pinned more memories than the cap.
1168        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// ============================================================================
1212// Tests
1213// ============================================================================
1214
1215#[cfg(test)]
1216mod tests {
1217    use super::*;
1218
1219    // MemoryItem tests
1220
1221    #[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    // relevance_score_at tests
1328
1329    #[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        // High importance weight → score dominated by importance
1335        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        // Short decay → recent item still scores well
1344        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(); // 30-day half-life
1358        let score = old_item.relevance_score_at(Utc::now(), &config);
1359        // After 60 days (2 half-lives), decay ≈ exp(-2) ≈ 0.135
1360        // score ≈ 0.5*0.7 + 0.135*0.3 ≈ 0.39
1361        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    // RelevanceConfig tests
1372
1373    #[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    // InMemoryStore tests
1382
1383    #[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    // PrunePolicy / prune() tests
1517
1518    #[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        // A recent item with low importance should NOT be pruned (not old enough)
1582        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}