Skip to main content

mnemo_core/query/
recall.rs

1use std::collections::HashSet;
2
3use serde::{Deserialize, Serialize};
4use uuid::Uuid;
5
6use crate::error::Result;
7use crate::hash::compute_content_hash;
8use crate::model::event::{AgentEvent, EventType};
9use crate::model::memory::{MemoryRecord, MemoryType, Scope};
10use crate::query::MnemoEngine;
11use crate::storage::MemoryFilter;
12#[allow(unused_imports)]
13use base64::Engine as _;
14
15#[derive(Debug, Clone, Default, Serialize, Deserialize)]
16pub struct TemporalRange {
17    pub after: Option<String>,
18    pub before: Option<String>,
19}
20
21impl TemporalRange {
22    pub fn new() -> Self {
23        Self::default()
24    }
25}
26
27#[derive(Debug, Clone, Serialize, Deserialize)]
28pub struct RecallRequest {
29    pub query: String,
30    pub agent_id: Option<String>,
31    pub limit: Option<usize>,
32    pub memory_type: Option<MemoryType>,
33    pub memory_types: Option<Vec<MemoryType>>,
34    pub scope: Option<Scope>,
35    pub min_importance: Option<f32>,
36    pub tags: Option<Vec<String>>,
37    pub org_id: Option<String>,
38    pub strategy: Option<String>,
39    pub temporal_range: Option<TemporalRange>,
40    pub recency_half_life_hours: Option<f64>,
41    pub hybrid_weights: Option<Vec<f32>>,
42    pub rrf_k: Option<f32>,
43    pub as_of: Option<String>,
44    /// When set, each `ScoredMemory` is augmented with a `score_breakdown`
45    /// that reports the per-signal score contributions (vector, bm25, graph,
46    /// recency) and final RRF rank.
47    pub explain: Option<bool>,
48    /// v0.4.0-rc3 (Task B1) — when `Some(true)` AND the engine has a
49    /// [`ProvenanceSigner`](crate::provenance::ProvenanceSigner)
50    /// attached, the response carries a [`ReadProvenance`](crate::provenance::ReadProvenance)
51    /// HMAC receipt over the recalled records. Default `None` keeps
52    /// the recall hot-path overhead at zero for callers that don't
53    /// need verifiable receipts.
54    pub with_provenance: Option<bool>,
55    /// v0.4.4 — typed retrieval mode. When `Some`, takes precedence
56    /// over the legacy `strategy` field (which stays in place for
57    /// backwards compatibility). When `None`, the engine falls back
58    /// to parsing `strategy` exactly as in v0.4.3. See
59    /// [`crate::retrieval::RetrievalMode`].
60    #[serde(default, skip_serializing_if = "Option::is_none")]
61    pub mode: Option<crate::retrieval::RetrievalMode>,
62    /// v0.4.7 — opt-in current-fact resolver. When `Some`, the
63    /// engine runs a post-processor over the standard recall result
64    /// set that groups candidates by `cfg.fact_key` and keeps the
65    /// most-recent write per group. See
66    /// [`crate::query::current_fact_resolver`] for the contract +
67    /// the MINTEval arXiv:2605.18565 anchor. Default `None` keeps
68    /// the read path unchanged.
69    #[serde(default, skip_serializing_if = "Option::is_none")]
70    pub current_fact_resolver:
71        Option<crate::query::current_fact_resolver::CurrentFactResolverConfig>,
72    /// v0.4.8 — opt-in orientation cache. When `Some` AND the
73    /// engine has an
74    /// [`OrientationCacheStore`][crate::query::orientation_cache::OrientationCacheStore]
75    /// attached, the engine maintains a per-namespace, constant-token
76    /// "context map" updated from each recall hit, and returns a
77    /// bounded rendering in
78    /// [`RecallResponse::orientation_cache`]. PEEK-anchored
79    /// (arXiv:2605.19932). Default `None` keeps the read path
80    /// unchanged.
81    #[serde(default, skip_serializing_if = "Option::is_none")]
82    pub orientation_cache: Option<crate::query::orientation_cache::OrientationCacheConfig>,
83    /// v0.4.12 — opt-in cost-aware evidence budget. When `Some`, the
84    /// engine runs the [`crate::query::evidence`] selector over the
85    /// ranked candidate set and returns the smallest prefix that
86    /// clears the configured sufficiency bar (capped by
87    /// `max_evidence`). Purely subtractive — it never reorders the
88    /// retrieval's top-k. Default `None` keeps the read path unchanged
89    /// (front-loaded top-`limit`).
90    #[serde(default, skip_serializing_if = "Option::is_none")]
91    pub evidence_budget: Option<crate::query::evidence::EvidenceBudget>,
92    /// EMBER (arXiv:2606.05894) — opt-in budgeted evidence retention.
93    /// When `Some(budget)`, the engine builds a
94    /// [`RetentionReport`](crate::query::retained::RetentionReport) that
95    /// packs the recalled hits into at most `budget` retained tokens as
96    /// verbatim *evidence capsules* (excerpt + retrieval key), ranked by
97    /// a `recency × hit-rate` recoverability heuristic, and returns it in
98    /// [`RecallResponse::retained_evidence`]. Purely **additive** — the
99    /// `memories` list is unchanged, so the default read path is
100    /// unaffected. See [`crate::query::retained`].
101    #[serde(default, skip_serializing_if = "Option::is_none")]
102    pub retained_token_budget: Option<usize>,
103    /// v0.4.15 — domain-scoped recall predicate (MASDR-RAG,
104    /// arXiv:2606.11350). When set (or when
105    /// [`mode`](Self::mode) is [`RetrievalMode::DomainScoped`][crate::retrieval::RetrievalMode::DomainScoped]),
106    /// the candidate set is restricted to the metadata-defined
107    /// sub-corpus described by this [`DomainScope`][crate::retrieval::DomainScope]
108    /// *before* the dense similarity step, countering vector-search
109    /// dilution at scale. Default `None` keeps the read path unchanged.
110    #[serde(default, skip_serializing_if = "Option::is_none")]
111    pub domain_scope: Option<crate::retrieval::DomainScope>,
112}
113
114impl RecallRequest {
115    pub fn new(query: String) -> Self {
116        Self {
117            query,
118            agent_id: None,
119            limit: None,
120            memory_type: None,
121            memory_types: None,
122            scope: None,
123            min_importance: None,
124            tags: None,
125            org_id: None,
126            strategy: None,
127            temporal_range: None,
128            recency_half_life_hours: None,
129            hybrid_weights: None,
130            rrf_k: None,
131            as_of: None,
132            explain: None,
133            with_provenance: None,
134            mode: None,
135            current_fact_resolver: None,
136            orientation_cache: None,
137            evidence_budget: None,
138            retained_token_budget: None,
139            domain_scope: None,
140        }
141    }
142}
143
144/// v0.4.7 — one entry of the supersession chain returned when the
145/// current-fact resolver is enabled with
146/// [`CurrentFactResolverConfig::include_supersession_chain`][crate::query::current_fact_resolver::CurrentFactResolverConfig::include_supersession_chain]
147/// set to `true`. Carries the prior fact version's id + the
148/// timestamps so an auditor can reconstruct the timeline.
149#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
150pub struct SupersededRecord {
151    pub id: Uuid,
152    pub fact_id: String,
153    pub superseded_by: Uuid,
154    /// Timestamp of the winning current record.
155    pub superseded_at: String,
156    /// Timestamp of the older record being marked superseded.
157    pub prior_updated_at: String,
158}
159
160/// Per-signal score contributions for a single recall hit.
161///
162/// Emitted when `RecallRequest.explain = Some(true)`. Each field is the
163/// raw signal score used as input to reciprocal-rank fusion (0 when the
164/// memory didn't appear in that list).
165#[derive(Debug, Clone, Default, Serialize, Deserialize)]
166pub struct ScoreBreakdown {
167    pub vector: f32,
168    pub bm25: f32,
169    pub graph: f32,
170    pub recency: f32,
171    /// 0-based position of the memory in the fused ranking.
172    pub rrf_rank: u32,
173}
174
175#[non_exhaustive]
176#[derive(Debug, Clone, Serialize, Deserialize)]
177pub struct RecallResponse {
178    pub memories: Vec<ScoredMemory>,
179    pub total: usize,
180    /// HMAC receipt over the recalled records — present iff the
181    /// caller set `RecallRequest.with_provenance = Some(true)` AND
182    /// the engine has a `ProvenanceSigner` attached.
183    /// See [`crate::provenance`].
184    #[serde(skip_serializing_if = "Option::is_none", default)]
185    pub provenance: Option<crate::provenance::ReadProvenance>,
186    /// v0.4.7 — older fact-versions dropped by the current-fact
187    /// resolver, in newest-superseded → oldest order. Present iff
188    /// the caller set
189    /// [`CurrentFactResolverConfig::include_supersession_chain`][crate::query::current_fact_resolver::CurrentFactResolverConfig::include_supersession_chain]
190    /// to `true` AND the resolver actually dropped any candidates.
191    #[serde(skip_serializing_if = "Option::is_none", default)]
192    pub superseded: Option<Vec<SupersededRecord>>,
193    /// v0.4.8 — bounded, namespace-scoped orientation map rendered
194    /// after the recall ran. Present iff the caller set
195    /// [`RecallRequest::orientation_cache`] AND the engine has an
196    /// [`OrientationCacheStore`][crate::query::orientation_cache::OrientationCacheStore]
197    /// attached AND the config did not set `include_in_response =
198    /// false`. PEEK-anchored (arXiv:2605.19932).
199    #[serde(skip_serializing_if = "Option::is_none", default)]
200    pub orientation_cache: Option<crate::query::orientation_cache::RenderedContextMap>,
201    /// v0.4.12 — diagnostics from the cost-aware evidence budget.
202    /// Present iff the caller set [`RecallRequest::evidence_budget`].
203    /// Reports the scorer used, how many candidates were examined vs
204    /// returned, the cumulative sufficiency score, and whether
205    /// early-stop / the cap fired. See [`crate::query::evidence`].
206    #[serde(skip_serializing_if = "Option::is_none", default)]
207    pub evidence_selection: Option<crate::query::evidence::EvidenceSelectionReport>,
208    /// EMBER (arXiv:2606.05894) — budgeted evidence-retention view.
209    /// Present iff the caller set
210    /// [`RecallRequest::retained_token_budget`]. Carries verbatim
211    /// evidence capsules (excerpt + retrieval key) packed under the
212    /// requested token cap, ranked by recoverability. Additive: the
213    /// `memories` list above is unchanged. See [`crate::query::retained`].
214    #[serde(skip_serializing_if = "Option::is_none", default)]
215    pub retained_evidence: Option<crate::query::retained::RetentionReport>,
216    /// v0.5.1 — active-reconstruction belief-state node (MRAgent,
217    /// arXiv:2606.06036). Present iff the caller selected the
218    /// `reconstruct` strategy ([`RetrievalMode::Reconstruct`][crate::retrieval::RetrievalMode::Reconstruct]).
219    /// Carries a deterministic summary synthesised from the retrieved
220    /// candidates plus the linked/causal context gathered by walking the
221    /// memory graph. Additive: `memories` is exactly the top-k the default
222    /// hybrid (`auto`) path returns, so the raw read path is unchanged.
223    #[serde(skip_serializing_if = "Option::is_none", default)]
224    pub reconstruction: Option<ReconstructedBelief>,
225}
226
227impl RecallResponse {
228    pub fn new(memories: Vec<ScoredMemory>, total: usize) -> Self {
229        Self {
230            memories,
231            total,
232            provenance: None,
233            superseded: None,
234            orientation_cache: None,
235            evidence_selection: None,
236            retained_evidence: None,
237            reconstruction: None,
238        }
239    }
240}
241
242/// v0.5.1 — a reconstructed belief-state node (MRAgent, arXiv:2606.06036).
243///
244/// Produced by the `reconstruct` recall strategy. Rather than returning
245/// the top-k hits alone, the strategy walks the memory graph from those
246/// hits to gather linked/causal context and synthesises a deterministic
247/// summary the caller receives ALONGSIDE the raw `memories`. The synthesis
248/// is rule-based (no LLM), so the same inputs always yield the same node —
249/// it is an honest substrate for A/B-ing reconstruction vs. retrieval on
250/// your own data, not a claim that retrieval is wrong.
251#[non_exhaustive]
252#[derive(Debug, Clone, Serialize, Deserialize)]
253pub struct ReconstructedBelief {
254    /// The cue (query) the belief was reconstructed for.
255    pub cue: String,
256    /// Deterministic summary: direct evidence (the retrieved hits) followed
257    /// by the linked/causal context gathered from the memory graph.
258    pub summary: String,
259    /// Ids of the retrieved candidates that seeded the reconstruction.
260    pub source_ids: Vec<Uuid>,
261    /// Ids of graph-linked memories pulled in as causal/linked context
262    /// (not present in `source_ids`).
263    pub linked_context_ids: Vec<Uuid>,
264    /// Mean retrieval score of the source hits — a coarse confidence proxy.
265    pub confidence: f32,
266}
267
268#[non_exhaustive]
269#[derive(Debug, Clone, Serialize, Deserialize)]
270pub struct ScoredMemory {
271    pub id: Uuid,
272    pub content: String,
273    pub agent_id: String,
274    pub memory_type: MemoryType,
275    pub scope: Scope,
276    pub importance: f32,
277    pub tags: Vec<String>,
278    pub metadata: serde_json::Value,
279    pub score: f32,
280    pub access_count: u64,
281    pub created_at: String,
282    pub updated_at: String,
283    #[serde(skip_serializing_if = "Option::is_none")]
284    pub score_breakdown: Option<ScoreBreakdown>,
285}
286
287impl From<(MemoryRecord, f32)> for ScoredMemory {
288    fn from((record, score): (MemoryRecord, f32)) -> Self {
289        Self {
290            id: record.id,
291            content: record.content,
292            agent_id: record.agent_id,
293            memory_type: record.memory_type,
294            scope: record.scope,
295            importance: record.importance,
296            tags: record.tags,
297            metadata: record.metadata,
298            score,
299            access_count: record.access_count,
300            created_at: record.created_at,
301            updated_at: record.updated_at,
302            score_breakdown: None,
303        }
304    }
305}
306
307/// Get a memory by ID, checking cache first then falling back to storage.
308async fn get_memory_cached(engine: &MnemoEngine, id: Uuid) -> Result<Option<MemoryRecord>> {
309    if let Some(ref cache) = engine.cache
310        && let Some(record) = cache.get(id)
311    {
312        return Ok(Some(record));
313    }
314    let result = engine.storage.get_memory(id).await?;
315    if let Some(ref record) = result
316        && let Some(ref cache) = engine.cache
317    {
318        cache.put(record.clone());
319    }
320    Ok(result)
321}
322
323pub async fn execute(engine: &MnemoEngine, request: RecallRequest) -> Result<RecallResponse> {
324    let limit = request.limit.unwrap_or(10).min(100);
325    let agent_id = request
326        .agent_id
327        .clone()
328        .unwrap_or_else(|| engine.default_agent_id.clone());
329    super::validate_agent_id(&agent_id)?;
330
331    // Determine strategy. v0.4.4: prefer the typed
332    // `mode: Option<RetrievalMode>` field when set; fall back to the
333    // legacy `strategy: Option<String>` field otherwise. Backwards
334    // compatible — SDKs that only marshal `strategy` continue to work.
335    let strategy = if let Some(ref mode) = request.mode {
336        mode.to_strategy_str()
337    } else if request
338        .domain_scope
339        .as_ref()
340        .map(|s| !s.is_empty())
341        .unwrap_or(false)
342    {
343        // v0.4.15 — a domain_scope predicate selects domain-scoped recall
344        // even when the caller didn't set the typed mode (ergonomic for
345        // SDKs that only marshal a `scope` kwarg).
346        "domain_scoped"
347    } else {
348        request.strategy.as_deref().unwrap_or("auto")
349    };
350
351    // v0.5.13 — fail loud, never silent-empty. Semantic and the semantic legs
352    // of hybrid/auto/graph/domain_scoped all depend on a real query vector. The
353    // no-op embedder returns an all-zero vector, which would make these paths
354    // silently return an empty or meaningless result set. Refuse with a typed
355    // error instead. Purely lexical (BM25) and exact/metadata recall need no
356    // embedder and are unaffected.
357    let needs_semantic = matches!(
358        strategy,
359        "semantic" | "hybrid" | "auto" | "graph" | "domain_scoped"
360    );
361    if needs_semantic && !engine.embedding.is_semantic_capable() {
362        return Err(crate::error::Error::EmbedderNotConfigured {
363            requested: strategy.to_string(),
364            backend: engine.storage.backend_name().to_string(),
365        });
366    }
367
368    // Compute query embedding (needed for semantic/hybrid/auto)
369    let query_embedding = engine.embedding.embed(&request.query).await?;
370
371    // Pre-compute accessible memory IDs for permission-safe ANN pre-filtering
372    let accessible_ids: HashSet<Uuid> = engine
373        .storage
374        .list_accessible_memory_ids(&agent_id, super::MAX_BATCH_QUERY_LIMIT)
375        .await?
376        .into_iter()
377        .collect();
378    let perm_filter = |id: Uuid| accessible_ids.contains(&id);
379
380    let mut scored_memories: Vec<(MemoryRecord, f32)> = Vec::new();
381    let mut breakdowns: std::collections::HashMap<Uuid, ScoreBreakdown> =
382        std::collections::HashMap::new();
383
384    match strategy {
385        "lexical" => {
386            // BM25-only path
387            if let Some(ref ft) = engine.full_text {
388                let bm25_results = ft.search(&request.query, limit * 3)?;
389                for (id, score) in bm25_results {
390                    if let Some(record) = get_memory_cached(engine, id).await?
391                        && passes_filters(&record, &request, &agent_id, engine).await
392                    {
393                        scored_memories.push((record, score));
394                    }
395                }
396            }
397        }
398        "semantic" => {
399            // Vector-only path with permission pre-filtering
400            let search_results =
401                engine
402                    .index
403                    .filtered_search(&query_embedding, limit * 3, &perm_filter)?;
404            for (id, distance) in search_results {
405                if let Some(record) = get_memory_cached(engine, id).await?
406                    && passes_filters(&record, &request, &agent_id, engine).await
407                {
408                    let score = 1.0 - distance;
409                    scored_memories.push((record, score));
410                }
411            }
412        }
413        "domain_scoped" => {
414            // v0.4.15 — domain-scoped recall (MASDR-RAG, arXiv:2606.11350).
415            // Restrict the candidate universe to the metadata-defined
416            // sub-corpus BEFORE the dense similarity step, so off-domain
417            // (but semantically similar) records can never enter the
418            // top-k. Then a single vector pass over the sub-corpus.
419            //
420            // The sub-corpus id-set is resolved from storage by the
421            // `DomainScope` predicate and composed with the permission
422            // filter, so the ANN sees only (accessible ∩ in-domain) ids.
423            let domain_ids: Option<HashSet<Uuid>> = match request.domain_scope.as_ref() {
424                Some(scope) if !scope.is_empty() => {
425                    // Coarse narrowing on org_id at the storage layer, then
426                    // exact predicate matching (namespace / doc_class / tags).
427                    let coarse = MemoryFilter {
428                        agent_id: None,
429                        memory_type: None,
430                        scope: None,
431                        tags: None,
432                        min_importance: None,
433                        org_id: scope.org_id.clone(),
434                        thread_id: None,
435                        include_deleted: false,
436                    };
437                    let records = engine
438                        .storage
439                        .list_memories(&coarse, super::MAX_BATCH_QUERY_LIMIT, 0)
440                        .await?;
441                    Some(
442                        records
443                            .iter()
444                            .filter(|r| scope.matches(r))
445                            .map(|r| r.id)
446                            .collect(),
447                    )
448                }
449                // DomainScoped selected without a predicate degrades to a
450                // plain vector pass (no extra restriction).
451                _ => None,
452            };
453
454            let domain_filter = |id: Uuid| {
455                perm_filter(id) && domain_ids.as_ref().map(|d| d.contains(&id)).unwrap_or(true)
456            };
457            let search_results =
458                engine
459                    .index
460                    .filtered_search(&query_embedding, limit * 3, &domain_filter)?;
461            for (id, distance) in search_results {
462                if let Some(record) = get_memory_cached(engine, id).await?
463                    && passes_filters(&record, &request, &agent_id, engine).await
464                {
465                    let score = 1.0 - distance;
466                    scored_memories.push((record, score));
467                }
468            }
469        }
470        "graph" => {
471            // Seed from vector results with permission pre-filtering, then expand via graph relations
472            let search_results =
473                engine
474                    .index
475                    .filtered_search(&query_embedding, limit * 3, &perm_filter)?;
476            let mut seeds: Vec<(Uuid, f32)> = Vec::new();
477            for (id, distance) in &search_results {
478                if let Some(record) = get_memory_cached(engine, *id).await?
479                    && passes_filters(&record, &request, &agent_id, engine).await
480                {
481                    seeds.push((*id, 1.0 - distance));
482                }
483            }
484
485            // Collect graph-expanded results with configurable multi-hop traversal
486            let max_hops = 2;
487            let mut seen: HashSet<Uuid> = seeds.iter().map(|(id, _)| *id).collect();
488            let mut graph_ranked: Vec<(Uuid, f32)> = Vec::new();
489
490            // Seeds get score 1.0
491            for &(id, _) in &seeds {
492                graph_ranked.push((id, 1.0));
493            }
494
495            // Multi-hop expansion with exponential decay
496            let mut frontier: Vec<Uuid> = seeds.iter().map(|(id, _)| *id).collect();
497            let mut decay = 0.5_f32;
498            for _hop in 0..max_hops {
499                let mut next_frontier: Vec<Uuid> = Vec::new();
500                for &id in &frontier {
501                    let from_rels = engine.storage.get_relations_from(id).await?;
502                    let to_rels = engine.storage.get_relations_to(id).await?;
503                    for rel in from_rels.iter().chain(to_rels.iter()) {
504                        let related_id = if rel.source_id == id {
505                            rel.target_id
506                        } else {
507                            rel.source_id
508                        };
509                        if seen.insert(related_id)
510                            && let Some(record) = get_memory_cached(engine, related_id).await?
511                            && passes_filters(&record, &request, &agent_id, engine).await
512                        {
513                            graph_ranked.push((related_id, decay));
514                            next_frontier.push(related_id);
515                        }
516                    }
517                }
518                frontier = next_frontier;
519                decay *= 0.5;
520            }
521
522            // Use RRF fusion with vector + graph lists
523            let mut v_sorted: Vec<(Uuid, f32)> = seeds.clone();
524            v_sorted.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
525            graph_ranked.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
526
527            let ranked_lists = vec![v_sorted, graph_ranked];
528            let rrf_k = request.rrf_k.unwrap_or(60.0);
529            let fused = if let Some(ref weights) = request.hybrid_weights {
530                crate::query::retrieval::weighted_reciprocal_rank_fusion(
531                    &ranked_lists,
532                    rrf_k,
533                    weights,
534                )
535            } else {
536                crate::query::retrieval::reciprocal_rank_fusion(&ranked_lists, rrf_k)
537            };
538
539            for (id, score) in fused {
540                if let Some(record) = get_memory_cached(engine, id).await?
541                    && passes_filters(&record, &request, &agent_id, engine).await
542                {
543                    scored_memories.push((record, score));
544                }
545            }
546        }
547        "exact" => {
548            // Filter-based exact matching, no embedding needed
549            // When as_of is set, include deleted records so the as_of filter can evaluate them
550            let filter = MemoryFilter {
551                agent_id: Some(agent_id.clone()),
552                memory_type: request.memory_type,
553                scope: request.scope,
554                tags: request.tags.clone(),
555                min_importance: request.min_importance,
556                org_id: request.org_id.clone(),
557                thread_id: None,
558                include_deleted: request.as_of.is_some(),
559            };
560            let memories = engine.storage.list_memories(&filter, limit, 0).await?;
561            for record in memories {
562                if passes_filters(&record, &request, &agent_id, engine).await {
563                    scored_memories.push((record, 1.0));
564                }
565            }
566        }
567        _ => {
568            // "auto" or "hybrid" — use hybrid if full_text available, else semantic
569            let vector_results =
570                engine
571                    .index
572                    .filtered_search(&query_embedding, limit * 3, &perm_filter)?;
573            let mut vector_ranked: Vec<(Uuid, f32)> = Vec::new();
574            for (id, distance) in vector_results {
575                vector_ranked.push((id, 1.0 - distance));
576            }
577
578            if let Some(ref ft) = engine.full_text {
579                // Hybrid: RRF fusion of vector + BM25 + recency
580                let bm25_results = ft.search(&request.query, limit * 3)?;
581
582                // Build recency-scored list from vector candidates
583                let mut recency_ranked: Vec<(Uuid, f32)> = Vec::new();
584                for &(id, _) in &vector_ranked {
585                    if let Some(record) = get_memory_cached(engine, id).await? {
586                        let r_score = crate::query::retrieval::recency_score(
587                            &record.created_at,
588                            request.recency_half_life_hours.unwrap_or(168.0),
589                        );
590                        recency_ranked.push((id, r_score));
591                    }
592                }
593                // Also add BM25 candidates to recency
594                for &(id, _) in &bm25_results {
595                    if !recency_ranked.iter().any(|(rid, _)| *rid == id)
596                        && let Some(record) = get_memory_cached(engine, id).await?
597                    {
598                        let r_score = crate::query::retrieval::recency_score(
599                            &record.created_at,
600                            request.recency_half_life_hours.unwrap_or(168.0),
601                        );
602                        recency_ranked.push((id, r_score));
603                    }
604                }
605
606                // Sort each list by score descending
607                let mut v_sorted = vector_ranked.clone();
608                v_sorted.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
609                let mut b_sorted = bm25_results;
610                b_sorted.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
611                recency_ranked
612                    .sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
613
614                // Graph expansion signal: from top-10 vector results, multi-hop expansion
615                let max_hops = 2;
616                let mut graph_ranked: Vec<(Uuid, f32)> = Vec::new();
617                let top_seeds: Vec<Uuid> =
618                    vector_ranked.iter().take(10).map(|(id, _)| *id).collect();
619                let mut graph_seen: HashSet<Uuid> = top_seeds.iter().copied().collect();
620                for &seed_id in &top_seeds {
621                    graph_ranked.push((seed_id, 1.0));
622                }
623                let mut frontier: Vec<Uuid> = top_seeds;
624                let mut decay = 0.5_f32;
625                for _hop in 0..max_hops {
626                    let mut next_frontier: Vec<Uuid> = Vec::new();
627                    for &fid in &frontier {
628                        match engine.storage.get_relations_from(fid).await {
629                            Ok(from_rels) => {
630                                for rel in &from_rels {
631                                    if graph_seen.insert(rel.target_id) {
632                                        graph_ranked.push((rel.target_id, decay));
633                                        next_frontier.push(rel.target_id);
634                                    }
635                                }
636                            }
637                            Err(e) => {
638                                tracing::warn!(memory_id = %fid, error = %e, "graph expansion: failed to get outgoing relations");
639                            }
640                        }
641                        match engine.storage.get_relations_to(fid).await {
642                            Ok(to_rels) => {
643                                for rel in &to_rels {
644                                    if graph_seen.insert(rel.source_id) {
645                                        graph_ranked.push((rel.source_id, decay));
646                                        next_frontier.push(rel.source_id);
647                                    }
648                                }
649                            }
650                            Err(e) => {
651                                tracing::warn!(memory_id = %fid, error = %e, "graph expansion: failed to get incoming relations");
652                            }
653                        }
654                    }
655                    frontier = next_frontier;
656                    decay *= 0.5;
657                }
658                graph_ranked
659                    .sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
660
661                // Capture per-signal score maps before moving the ranked lists
662                // into the fusion call, so `explain=true` can surface each
663                // signal's contribution in the response.
664                let explain = request.explain.unwrap_or(false);
665                type SignalMap = std::collections::HashMap<Uuid, f32>;
666                let (vector_map, bm25_map, recency_map, graph_map): (
667                    SignalMap,
668                    SignalMap,
669                    SignalMap,
670                    SignalMap,
671                ) = if explain {
672                    (
673                        v_sorted.iter().copied().collect(),
674                        b_sorted.iter().copied().collect(),
675                        recency_ranked.iter().copied().collect(),
676                        graph_ranked.iter().copied().collect(),
677                    )
678                } else {
679                    Default::default()
680                };
681
682                let ranked_lists = vec![v_sorted, b_sorted, recency_ranked, graph_ranked];
683                let rrf_k = request.rrf_k.unwrap_or(60.0);
684                let fused = if let Some(ref weights) = request.hybrid_weights {
685                    crate::query::retrieval::weighted_reciprocal_rank_fusion(
686                        &ranked_lists,
687                        rrf_k,
688                        weights,
689                    )
690                } else {
691                    crate::query::retrieval::reciprocal_rank_fusion(&ranked_lists, rrf_k)
692                };
693
694                for (rank, (id, score)) in fused.into_iter().enumerate() {
695                    if let Some(record) = get_memory_cached(engine, id).await?
696                        && passes_filters(&record, &request, &agent_id, engine).await
697                    {
698                        scored_memories.push((record, score));
699                        if explain {
700                            breakdowns.insert(
701                                id,
702                                ScoreBreakdown {
703                                    vector: vector_map.get(&id).copied().unwrap_or(0.0),
704                                    bm25: bm25_map.get(&id).copied().unwrap_or(0.0),
705                                    graph: graph_map.get(&id).copied().unwrap_or(0.0),
706                                    recency: recency_map.get(&id).copied().unwrap_or(0.0),
707                                    rrf_rank: rank as u32,
708                                },
709                            );
710                        }
711                    }
712                }
713            } else {
714                // Fallback to semantic-only
715                for (id, score) in vector_ranked {
716                    if let Some(record) = get_memory_cached(engine, id).await?
717                        && passes_filters(&record, &request, &agent_id, engine).await
718                    {
719                        scored_memories.push((record, score));
720                    }
721                }
722            }
723        }
724    }
725
726    // Sort by score descending
727    scored_memories.sort_by(|a, b| b.1.partial_cmp(&a.1).unwrap_or(std::cmp::Ordering::Equal));
728    scored_memories.truncate(limit);
729
730    // v0.4.12 — opt-in cost-aware evidence budget. Runs only when the
731    // caller set `request.evidence_budget`. The selector operates on
732    // the already-ranked list and returns the smallest prefix that
733    // clears the sufficiency bar (capped by `max_evidence`); it never
734    // reorders, so the top-k cosine/RRF ordering is preserved. Applied
735    // BEFORE `touch_memory` so we do not mark-accessed evidence the
736    // budget trimmed away (cost-aware on the write side too). See
737    // [`crate::query::evidence`].
738    let evidence_selection = if let Some(ref budget) = request.evidence_budget {
739        let cosine_default = crate::query::evidence::CosineScorer;
740        let scorer: &dyn crate::query::evidence::EvidenceScorer =
741            match (budget.scorer, engine.evidence_scorer.as_ref()) {
742                (crate::query::evidence::ScorerKind::Delta, Some(s)) => s.as_ref(),
743                _ => &cosine_default,
744            };
745        // Pass the query embedding only when it is non-degenerate
746        // (NoopEmbedding yields all-zero vectors, for which cosine is
747        // undefined and the scorer should fall back to retrieval score).
748        let q_emb: Option<&[f32]> = if query_embedding.iter().any(|v| *v != 0.0) {
749            Some(query_embedding.as_slice())
750        } else {
751            None
752        };
753        let candidates: Vec<crate::query::evidence::EvidenceCandidate<'_>> = scored_memories
754            .iter()
755            .map(|(r, score)| crate::query::evidence::EvidenceCandidate {
756                content: &r.content,
757                embedding: r.embedding.as_deref(),
758                retrieval_score: *score,
759            })
760            .collect();
761        let selection = crate::query::evidence::select_within_budget(
762            &candidates,
763            budget,
764            scorer,
765            &request.query,
766            q_emb,
767        );
768        let keep = selection.keep;
769        drop(candidates);
770        scored_memories.truncate(keep);
771        Some(selection.report)
772    } else {
773        None
774    };
775
776    let _total_pre_resolver = scored_memories.len();
777
778    // Touch accessed memories
779    for (record, _) in &scored_memories {
780        if let Err(e) = engine.storage.touch_memory(record.id).await {
781            tracing::warn!(memory_id = %record.id, error = %e, "failed to update access timestamp");
782        }
783    }
784
785    // Decrypt content if encryption is configured
786    if let Some(ref enc) = engine.encryption {
787        for (record, _) in &mut scored_memories {
788            match base64::engine::general_purpose::STANDARD.decode(&record.content) {
789                Ok(encrypted_bytes) => match enc.decrypt(&encrypted_bytes) {
790                    Ok(decrypted) => match String::from_utf8(decrypted) {
791                        Ok(plaintext) => record.content = plaintext,
792                        Err(e) => {
793                            tracing::error!(memory_id = %record.id, error = %e, "decrypted content is not valid UTF-8");
794                            record.content = "[content unavailable: decryption error]".to_string();
795                        }
796                    },
797                    Err(e) => {
798                        tracing::error!(memory_id = %record.id, error = %e, "failed to decrypt memory content");
799                        record.content = "[content unavailable: decryption error]".to_string();
800                    }
801                },
802                Err(e) => {
803                    tracing::error!(memory_id = %record.id, error = %e, "failed to decode encrypted content");
804                    record.content = "[content unavailable: decryption error]".to_string();
805                }
806            }
807        }
808    }
809
810    // Keep the underlying records around if the caller asked for a
811    // provenance receipt (Task B1) — the HMAC chain needs the
812    // content_hash + prev_hash off each record before they get
813    // collapsed into ScoredMemory.
814    let provenance_records: Option<Vec<MemoryRecord>> =
815        if request.with_provenance == Some(true) && engine.provenance_signer.is_some() {
816            Some(scored_memories.iter().map(|(r, _)| r.clone()).collect())
817        } else {
818            None
819        };
820
821    let memories: Vec<ScoredMemory> = scored_memories
822        .into_iter()
823        .map(|(record, score)| {
824            let id = record.id;
825            let mut scored = ScoredMemory::from((record, score));
826            if let Some(breakdown) = breakdowns.remove(&id) {
827                scored.score_breakdown = Some(breakdown);
828            }
829            scored
830        })
831        .collect();
832
833    // v0.4.7 — opt-in current-fact resolver post-process. Runs only
834    // when the caller set `request.current_fact_resolver`. The
835    // resolver groups by `cfg.fact_key`, keeps the most-recent
836    // write per group, and (optionally) returns the older versions
837    // as a supersession chain. See
838    // [`crate::query::current_fact_resolver`] for the MINTEval
839    // arXiv:2605.18565 anchor + the contract.
840    let (memories, superseded_chain) = if let Some(ref cfg) = request.current_fact_resolver {
841        let out = crate::query::current_fact_resolver::resolve(cfg, memories);
842        let chain = if cfg.include_supersession_chain && !out.superseded.is_empty() {
843            Some(out.superseded)
844        } else {
845            None
846        };
847        (out.kept, chain)
848    } else {
849        (memories, None)
850    };
851    let total = memories.len();
852
853    // v0.5.1 — active reconstruction (MRAgent, arXiv:2606.06036). When the
854    // caller selected the `reconstruct` strategy, walk the memory graph from
855    // the retrieved hits to gather linked/causal context and synthesise a
856    // deterministic belief-state node returned ALONGSIDE the raw hits. The
857    // `memories` list above is untouched, so this is purely additive.
858    let reconstruction = if strategy == "reconstruct" {
859        Some(reconstruct_belief(engine, &request, &agent_id, &memories).await)
860    } else {
861        None
862    };
863
864    // v0.4.8 — opt-in orientation cache. Runs only when the caller
865    // set `request.orientation_cache` AND the engine has an
866    // `OrientationCacheStore` attached. Per-namespace map is
867    // updated from the hits + a bounded rendering is returned. See
868    // [`crate::query::orientation_cache`] for the PEEK
869    // arXiv:2605.19932 anchor + the contract.
870    let orientation_rendered = match (
871        request.orientation_cache.as_ref(),
872        engine.orientation_cache_store.as_ref(),
873    ) {
874        (Some(cfg), Some(store)) => {
875            let ns = crate::query::orientation_cache::resolve_namespace(
876                cfg,
877                &agent_id,
878                request.org_id.as_deref(),
879            );
880            let rendered =
881                crate::query::orientation_cache::update_and_render(store, cfg, &ns, &memories);
882            if cfg.include_in_response {
883                Some(rendered)
884            } else {
885                None
886            }
887        }
888        _ => None,
889    };
890
891    // Emit MemoryRead event with hash chain linking (fire-and-forget)
892    let now = chrono::Utc::now().to_rfc3339();
893    let event_content_hash = compute_content_hash(&request.query, &agent_id, &now);
894    let prev_event_hash = match engine.storage.get_latest_event_hash(&agent_id, None).await {
895        Ok(hash) => hash,
896        Err(e) => {
897            tracing::warn!(error = %e, "failed to get latest event hash, starting new chain segment");
898            None
899        }
900    };
901    let event_prev_hash = Some(crate::hash::compute_chain_hash(
902        &event_content_hash,
903        prev_event_hash.as_deref(),
904    ));
905    let mut event = AgentEvent {
906        id: Uuid::now_v7(),
907        agent_id: agent_id.clone(),
908        thread_id: None,
909        run_id: None,
910        parent_event_id: None,
911        event_type: EventType::MemoryRead,
912        payload: serde_json::json!({
913            "query": request.query,
914            "results": total,
915            "strategy": strategy,
916        }),
917        trace_id: None,
918        span_id: None,
919        model: None,
920        tokens_input: None,
921        tokens_output: None,
922        latency_ms: None,
923        cost_usd: None,
924        timestamp: now.clone(),
925        logical_clock: 0,
926        content_hash: event_content_hash,
927        prev_hash: event_prev_hash,
928        embedding: None,
929    };
930    // Optionally embed the event payload
931    if engine.embed_events
932        && let Ok(emb) = engine.embedding.embed(&event.payload.to_string()).await
933    {
934        event.embedding = Some(emb);
935    }
936    if let Err(e) = engine.storage.insert_event(&event).await {
937        tracing::error!(event_id = %event.id, error = %e, "failed to insert audit event");
938    }
939
940    // v0.4.0-rc3 (B1) — sign a ReadProvenance over the recalled
941    // records when the caller opted in. Failures are non-fatal:
942    // missing signer or HMAC error degrades to "no provenance" so the
943    // recall still returns. The caller can detect by `provenance.is_none()`.
944    let provenance = if let (Some(records), Some(signer)) =
945        (provenance_records, engine.provenance_signer.as_ref())
946    {
947        match signer.sign(&agent_id, &request.query, &records) {
948            Ok(p) => Some(p),
949            Err(e) => {
950                tracing::warn!(error = %e, "failed to sign read provenance; degrading to no-provenance response");
951                None
952            }
953        }
954    } else {
955        None
956    };
957
958    // EMBER (arXiv:2606.05894) — opt-in budgeted evidence retention.
959    // Runs only when the caller set `request.retained_token_budget`.
960    // Builds verbatim evidence capsules (excerpt + retrieval key) packed
961    // under the token cap, ranked by `recency × hit-rate` recoverability.
962    // Computed from the FINAL `memories` (post current-fact resolver,
963    // decrypted) and returned ALONGSIDE them — `memories` is not
964    // modified, so the default read path is unaffected. See
965    // [`crate::query::retained`].
966    let retained_evidence = request.retained_token_budget.map(|budget| {
967        let retain_now = chrono::Utc::now();
968        let candidates: Vec<crate::query::retained::RetentionCandidate<'_>> = memories
969            .iter()
970            .map(|m| {
971                let age_hours = chrono::DateTime::parse_from_rfc3339(&m.updated_at)
972                    .or_else(|_| chrono::DateTime::parse_from_rfc3339(&m.created_at))
973                    .map(|ts| {
974                        (retain_now - ts.with_timezone(&chrono::Utc)).num_seconds() as f64 / 3600.0
975                    })
976                    .unwrap_or(0.0);
977                crate::query::retained::RetentionCandidate {
978                    id: m.id,
979                    content: &m.content,
980                    access_count: m.access_count,
981                    age_hours,
982                    retrieval_score: m.score,
983                }
984            })
985            .collect();
986        crate::query::retained::retain_within_budget(
987            &candidates,
988            budget,
989            crate::query::retained::DEFAULT_EXCERPT_TOKENS,
990        )
991    });
992
993    Ok(RecallResponse {
994        memories,
995        total,
996        provenance,
997        superseded: superseded_chain,
998        orientation_cache: orientation_rendered,
999        evidence_selection,
1000        retained_evidence,
1001        reconstruction,
1002    })
1003}
1004
1005/// v0.5.1 — synthesise a [`ReconstructedBelief`] from the retrieved hits
1006/// (MRAgent, arXiv:2606.06036). Walks one hop of memory-graph relations
1007/// outward from each hit to gather linked/causal context, then renders a
1008/// deterministic, rule-based summary (no LLM). Used only by the
1009/// `reconstruct` strategy; the raw `memories` are left unchanged.
1010async fn reconstruct_belief(
1011    engine: &MnemoEngine,
1012    request: &RecallRequest,
1013    agent_id: &str,
1014    memories: &[ScoredMemory],
1015) -> ReconstructedBelief {
1016    let cue = request.query.clone();
1017    if memories.is_empty() {
1018        return ReconstructedBelief {
1019            cue: cue.clone(),
1020            summary: format!("No memories matched the cue \"{cue}\"."),
1021            source_ids: Vec::new(),
1022            linked_context_ids: Vec::new(),
1023            confidence: 0.0,
1024        };
1025    }
1026
1027    let source_ids: Vec<Uuid> = memories.iter().map(|m| m.id).collect();
1028    let mut seen: HashSet<Uuid> = source_ids.iter().copied().collect();
1029
1030    // Walk one hop of relations outward from each hit to gather
1031    // linked/causal context. Deterministic order: hits in rank order, and
1032    // within a hit, outgoing relations before incoming.
1033    let mut linked: Vec<(Uuid, String)> = Vec::new();
1034    for m in memories {
1035        let from_rels = engine
1036            .storage
1037            .get_relations_from(m.id)
1038            .await
1039            .unwrap_or_default();
1040        let to_rels = engine
1041            .storage
1042            .get_relations_to(m.id)
1043            .await
1044            .unwrap_or_default();
1045        for rel in from_rels.iter().chain(to_rels.iter()) {
1046            let linked_id = if rel.source_id == m.id {
1047                rel.target_id
1048            } else {
1049                rel.source_id
1050            };
1051            if seen.insert(linked_id)
1052                && let Ok(Some(mut rec)) = engine.storage.get_memory(linked_id).await
1053                && passes_filters(&rec, request, agent_id, engine).await
1054            {
1055                decrypt_record_content(engine, &mut rec);
1056                linked.push((linked_id, rec.content));
1057            }
1058        }
1059    }
1060
1061    // Deterministic, rule-based belief summary (no LLM).
1062    let mut summary = format!("Reconstructed belief for cue \"{cue}\":\n\nDirect evidence:\n");
1063    for (i, m) in memories.iter().enumerate() {
1064        summary.push_str(&format!("{}. {}\n", i + 1, excerpt(&m.content, 200)));
1065    }
1066    if linked.is_empty() {
1067        summary.push_str("\n(No linked context found in the memory graph.)\n");
1068    } else {
1069        summary.push_str("\nLinked context (from graph relations):\n");
1070        for (_, content) in &linked {
1071            summary.push_str(&format!("- {}\n", excerpt(content, 160)));
1072        }
1073    }
1074
1075    let confidence = memories.iter().map(|m| m.score).sum::<f32>() / memories.len() as f32;
1076
1077    ReconstructedBelief {
1078        cue,
1079        summary,
1080        source_ids,
1081        linked_context_ids: linked.into_iter().map(|(id, _)| id).collect(),
1082        confidence,
1083    }
1084}
1085
1086/// First non-empty line of `content`, truncated to `max` chars (char-safe).
1087fn excerpt(content: &str, max: usize) -> String {
1088    let line = content.lines().find(|l| !l.trim().is_empty()).unwrap_or("");
1089    let trimmed = line.trim();
1090    if trimmed.chars().count() <= max {
1091        trimmed.to_string()
1092    } else {
1093        let mut out: String = trimmed.chars().take(max).collect();
1094        out.push('…');
1095        out
1096    }
1097}
1098
1099/// Decrypt a record's content in place if engine-level encryption is on.
1100/// Mirrors the read-path decryption in [`execute`]; used by
1101/// [`reconstruct_belief`] for graph-linked records fetched after the main
1102/// decrypt loop.
1103fn decrypt_record_content(engine: &MnemoEngine, record: &mut MemoryRecord) {
1104    if let Some(ref enc) = engine.encryption {
1105        if let Ok(bytes) = base64::engine::general_purpose::STANDARD.decode(&record.content)
1106            && let Ok(plain) = enc.decrypt(&bytes)
1107            && let Ok(text) = String::from_utf8(plain)
1108        {
1109            record.content = text;
1110        } else {
1111            record.content = "[content unavailable: decryption error]".to_string();
1112        }
1113    }
1114}
1115
1116async fn passes_filters(
1117    record: &MemoryRecord,
1118    request: &RecallRequest,
1119    agent_id: &str,
1120    engine: &MnemoEngine,
1121) -> bool {
1122    // Experience-tier plan records (DocTrace, arXiv:2606.10921) are never
1123    // surfaced by ordinary recall — they are replayed only via
1124    // `recall_plan`. Skip them unless the caller explicitly asks for the
1125    // reserved tag.
1126    if record
1127        .tags
1128        .iter()
1129        .any(|t| t == crate::query::experience::EXPERIENCE_PLAN_TAG)
1130        && !request
1131            .tags
1132            .as_ref()
1133            .map(|ts| {
1134                ts.iter()
1135                    .any(|t| t == crate::query::experience::EXPERIENCE_PLAN_TAG)
1136            })
1137            .unwrap_or(false)
1138    {
1139        return false;
1140    }
1141
1142    // Skip deleted (unless as_of is set — the as_of filter handles deleted records)
1143    if request.as_of.is_none() && record.is_deleted() {
1144        return false;
1145    }
1146
1147    // Skip expired
1148    if let Some(ref expires_at) = record.expires_at
1149        && let Ok(exp) = chrono::DateTime::parse_from_rfc3339(expires_at)
1150        && exp < chrono::Utc::now()
1151    {
1152        return false;
1153    }
1154
1155    // Skip quarantined
1156    if record.quarantined {
1157        return false;
1158    }
1159
1160    // Scope filter (explicit request scope filter, separate from visibility below)
1161    if let Some(ref s) = request.scope
1162        && record.scope != *s
1163    {
1164        return false;
1165    }
1166
1167    // Type filter: memory_types (multi) takes precedence over memory_type (single)
1168    if let Some(ref mts) = request.memory_types {
1169        if !mts.contains(&record.memory_type) {
1170            return false;
1171        }
1172    } else if let Some(ref mt) = request.memory_type
1173        && record.memory_type != *mt
1174    {
1175        return false;
1176    }
1177
1178    // Importance filter
1179    if let Some(min_imp) = request.min_importance
1180        && record.importance < min_imp
1181    {
1182        return false;
1183    }
1184
1185    // Tags filter
1186    if let Some(ref req_tags) = request.tags
1187        && !req_tags.iter().any(|t| record.tags.contains(t))
1188    {
1189        return false;
1190    }
1191
1192    // Temporal range filter (parse to DateTime for correct comparison)
1193    if let Some(ref tr) = request.temporal_range {
1194        if let Some(ref after) = tr.after
1195            && let (Ok(after_dt), Ok(record_dt)) = (
1196                chrono::DateTime::parse_from_rfc3339(after),
1197                chrono::DateTime::parse_from_rfc3339(&record.created_at),
1198            )
1199            && record_dt < after_dt
1200        {
1201            return false;
1202        }
1203        if let Some(ref before) = tr.before
1204            && let (Ok(before_dt), Ok(record_dt)) = (
1205                chrono::DateTime::parse_from_rfc3339(before),
1206                chrono::DateTime::parse_from_rfc3339(&record.created_at),
1207            )
1208            && record_dt > before_dt
1209        {
1210            return false;
1211        }
1212    }
1213
1214    // Point-in-time as_of filter: show memory state at time T
1215    if let Some(ref as_of) = request.as_of {
1216        if let (Ok(as_of_dt), Ok(record_dt)) = (
1217            chrono::DateTime::parse_from_rfc3339(as_of),
1218            chrono::DateTime::parse_from_rfc3339(&record.created_at),
1219        ) && record_dt > as_of_dt
1220        {
1221            // Exclude memories created after as_of
1222            return false;
1223        }
1224        // Exclude memories already deleted at as_of
1225        if let Some(ref deleted_at) = record.deleted_at
1226            && let (Ok(del_dt), Ok(as_of_dt)) = (
1227                chrono::DateTime::parse_from_rfc3339(deleted_at),
1228                chrono::DateTime::parse_from_rfc3339(as_of),
1229            )
1230            && del_dt <= as_of_dt
1231        {
1232            return false;
1233        }
1234    }
1235
1236    // Scope-based visibility
1237    match record.scope {
1238        Scope::Public | Scope::Global => true,
1239        Scope::Shared => {
1240            record.agent_id == agent_id
1241                || engine
1242                    .storage
1243                    .check_permission(
1244                        record.id,
1245                        agent_id,
1246                        crate::model::acl::Permission::Read,
1247                    )
1248                    .await
1249                    .unwrap_or_else(|e| {
1250                        tracing::warn!(memory_id = %record.id, error = %e, "permission check failed, denying access");
1251                        false
1252                    })
1253        }
1254        Scope::Private => record.agent_id == agent_id,
1255    }
1256}