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 pub explain: Option<bool>,
48 pub with_provenance: Option<bool>,
55 #[serde(default, skip_serializing_if = "Option::is_none")]
61 pub mode: Option<crate::retrieval::RetrievalMode>,
62 #[serde(default, skip_serializing_if = "Option::is_none")]
70 pub current_fact_resolver:
71 Option<crate::query::current_fact_resolver::CurrentFactResolverConfig>,
72 #[serde(default, skip_serializing_if = "Option::is_none")]
82 pub orientation_cache: Option<crate::query::orientation_cache::OrientationCacheConfig>,
83 #[serde(default, skip_serializing_if = "Option::is_none")]
91 pub evidence_budget: Option<crate::query::evidence::EvidenceBudget>,
92 #[serde(default, skip_serializing_if = "Option::is_none")]
102 pub retained_token_budget: Option<usize>,
103 #[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#[derive(Debug, Clone, Serialize, Deserialize, PartialEq)]
150pub struct SupersededRecord {
151 pub id: Uuid,
152 pub fact_id: String,
153 pub superseded_by: Uuid,
154 pub superseded_at: String,
156 pub prior_updated_at: String,
158}
159
160#[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 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 #[serde(skip_serializing_if = "Option::is_none", default)]
185 pub provenance: Option<crate::provenance::ReadProvenance>,
186 #[serde(skip_serializing_if = "Option::is_none", default)]
192 pub superseded: Option<Vec<SupersededRecord>>,
193 #[serde(skip_serializing_if = "Option::is_none", default)]
200 pub orientation_cache: Option<crate::query::orientation_cache::RenderedContextMap>,
201 #[serde(skip_serializing_if = "Option::is_none", default)]
207 pub evidence_selection: Option<crate::query::evidence::EvidenceSelectionReport>,
208 #[serde(skip_serializing_if = "Option::is_none", default)]
215 pub retained_evidence: Option<crate::query::retained::RetentionReport>,
216 #[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#[non_exhaustive]
252#[derive(Debug, Clone, Serialize, Deserialize)]
253pub struct ReconstructedBelief {
254 pub cue: String,
256 pub summary: String,
259 pub source_ids: Vec<Uuid>,
261 pub linked_context_ids: Vec<Uuid>,
264 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
307async 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 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 "domain_scoped"
347 } else {
348 request.strategy.as_deref().unwrap_or("auto")
349 };
350
351 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 let query_embedding = engine.embedding.embed(&request.query).await?;
370
371 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 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 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 let domain_ids: Option<HashSet<Uuid>> = match request.domain_scope.as_ref() {
424 Some(scope) if !scope.is_empty() => {
425 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 _ => 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 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 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 for &(id, _) in &seeds {
492 graph_ranked.push((id, 1.0));
493 }
494
495 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 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 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 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 let bm25_results = ft.search(&request.query, limit * 3)?;
581
582 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 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 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 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 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 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 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 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 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 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 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 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 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 let reconstruction = if strategy == "reconstruct" {
859 Some(reconstruct_belief(engine, &request, &agent_id, &memories).await)
860 } else {
861 None
862 };
863
864 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 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 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 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 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
1005async 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 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 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
1086fn 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
1099fn 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 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 if request.as_of.is_none() && record.is_deleted() {
1144 return false;
1145 }
1146
1147 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 if record.quarantined {
1157 return false;
1158 }
1159
1160 if let Some(ref s) = request.scope
1162 && record.scope != *s
1163 {
1164 return false;
1165 }
1166
1167 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 if let Some(min_imp) = request.min_importance
1180 && record.importance < min_imp
1181 {
1182 return false;
1183 }
1184
1185 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 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 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 return false;
1223 }
1224 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 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}