Skip to main content

mnemara_store_file/
lib.rs

1use async_trait::async_trait;
2use mnemara_core::{
3    ArchiveReceipt, ArchiveRequest, BatchUpsertRequest, ChangefeedEvent, ChangefeedEventKind,
4    ChangefeedReport, ChangefeedRequest, CompactionReport, CompactionRequest,
5    ConfiguredRecallScorer, DeleteReceipt, DeleteRequest, EngineConfig, Error, ExportRequest,
6    GraphInspectionReport, GraphInspectionRequest, ImportFailure, ImportMode, ImportReport,
7    ImportRequest, IntegrityCheckReport, IntegrityCheckRequest, LineageLink, LineageRelationKind,
8    MaintenanceStats, MemoryHistoricalState, MemoryQualityState, MemoryRecord, MemoryRecordKind,
9    MemoryScope, MemoryStore, MemoryTrustLevel, NamespaceStats, PlannedRecallCandidate,
10    PortableRecord, PortableStorePackage, RecallExplanation, RecallHistoricalMode, RecallHit,
11    RecallPlanner, RecallPlanningProfile, RecallPlanningTrace, RecallQuery, RecallResult,
12    RecallScorer, RecallTemporalOrder, RecallTraceCandidate, RecoverReceipt, RecoverRequest,
13    RepairReport, RepairRequest, Result, SemanticEmbedder, SnapshotManifest, StoreStatsReport,
14    StoreStatsRequest, SuppressReceipt, SuppressRequest, SynthesisProposal, SynthesisReport,
15    SynthesisRequest, TimeTravelRecallRequest, UpsertReceipt, UpsertRequest,
16    build_graph_inspection_report,
17};
18use serde::{Deserialize, Serialize};
19use std::collections::{BTreeMap, BTreeSet, HashMap};
20use std::fmt;
21use std::fs;
22use std::hash::{Hash, Hasher};
23use std::path::{Path, PathBuf};
24use std::sync::Arc;
25use std::time::{SystemTime, UNIX_EPOCH};
26
27#[derive(Clone)]
28pub struct FileStoreConfig {
29    pub data_dir: PathBuf,
30    pub engine_config: EngineConfig,
31    pub shared_embedder: Option<SharedEmbedderConfig>,
32}
33
34#[derive(Clone)]
35pub struct SharedEmbedderConfig {
36    pub embedder: Arc<dyn SemanticEmbedder>,
37    pub provider_note: String,
38}
39
40impl SharedEmbedderConfig {
41    fn new(embedder: Arc<dyn SemanticEmbedder>, provider_note: impl Into<String>) -> Self {
42        Self {
43            embedder,
44            provider_note: provider_note.into(),
45        }
46    }
47}
48
49impl fmt::Debug for SharedEmbedderConfig {
50    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
51        formatter
52            .debug_struct("SharedEmbedderConfig")
53            .field("provider_note", &self.provider_note)
54            .finish()
55    }
56}
57
58impl fmt::Debug for FileStoreConfig {
59    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
60        formatter
61            .debug_struct("FileStoreConfig")
62            .field("data_dir", &self.data_dir)
63            .field("engine_config", &self.engine_config)
64            .field("shared_embedder", &self.shared_embedder)
65            .finish()
66    }
67}
68
69impl FileStoreConfig {
70    pub fn new(data_dir: impl AsRef<Path>) -> Self {
71        Self {
72            data_dir: data_dir.as_ref().to_path_buf(),
73            engine_config: EngineConfig::default(),
74            shared_embedder: None,
75        }
76    }
77
78    pub fn with_engine_config(mut self, engine_config: EngineConfig) -> Self {
79        self.engine_config = engine_config;
80        self
81    }
82
83    pub fn with_shared_embedder(
84        mut self,
85        embedder: Arc<dyn SemanticEmbedder>,
86        provider_note: impl Into<String>,
87    ) -> Self {
88        self.shared_embedder = Some(SharedEmbedderConfig::new(embedder, provider_note));
89        self
90    }
91
92    fn recall_planner(&self) -> RecallPlanner {
93        if let Some(shared_embedder) = &self.shared_embedder {
94            RecallPlanner::with_shared_embedder(
95                self.engine_config.recall_planning_profile,
96                self.engine_config.graph_expansion_max_hops,
97                self.engine_config.recall_scorer_kind,
98                self.engine_config.recall_scoring_profile,
99                self.engine_config.recall_policy_profile,
100                Arc::clone(&shared_embedder.embedder),
101                shared_embedder.provider_note.clone(),
102            )
103        } else {
104            RecallPlanner::from_engine_config(&self.engine_config)
105        }
106    }
107}
108
109#[derive(Debug)]
110pub struct FileMemoryStore {
111    config: FileStoreConfig,
112}
113
114#[derive(Debug, Clone, Serialize, Deserialize)]
115struct StoredRecord {
116    record: MemoryRecord,
117    idempotency_key: Option<String>,
118}
119
120#[derive(Debug, Clone, Serialize, Deserialize)]
121struct VersionedStoredRecord {
122    record_id: String,
123    observed_at_unix_ms: u64,
124    record: Option<MemoryRecord>,
125    idempotency_key: Option<String>,
126}
127
128#[derive(Debug, Clone)]
129struct IdempotencyMapping {
130    scoped_key: String,
131    record_id: String,
132}
133
134type ScopeKeyParts = (
135    String,
136    String,
137    String,
138    Option<String>,
139    Option<String>,
140    String,
141);
142
143#[derive(Debug, Clone, Copy, Default)]
144struct IntegritySummary {
145    scanned_records: u64,
146    scanned_idempotency_keys: u64,
147    stale_idempotency_keys: u64,
148    missing_idempotency_keys: u64,
149    duplicate_active_records: u64,
150}
151
152#[derive(Debug, Clone, Copy, Default)]
153struct RelativeTemporalBounds {
154    after_unix_ms: Option<u64>,
155    before_unix_ms: Option<u64>,
156}
157
158const PORTABLE_PACKAGE_VERSION: u32 = 1;
159
160impl FileMemoryStore {
161    pub fn open(config: FileStoreConfig) -> Result<Self> {
162        fs::create_dir_all(Self::records_dir(&config.data_dir)).map_err(|err| {
163            Error::Backend(format!(
164                "failed to create records dir {}: {err}",
165                Self::records_dir(&config.data_dir).display()
166            ))
167        })?;
168        fs::create_dir_all(Self::idempotency_dir(&config.data_dir)).map_err(|err| {
169            Error::Backend(format!(
170                "failed to create idempotency dir {}: {err}",
171                Self::idempotency_dir(&config.data_dir).display()
172            ))
173        })?;
174        fs::create_dir_all(Self::versions_dir(&config.data_dir)).map_err(|err| {
175            Error::Backend(format!(
176                "failed to create versions dir {}: {err}",
177                Self::versions_dir(&config.data_dir).display()
178            ))
179        })?;
180        fs::create_dir_all(Self::changefeed_dir(&config.data_dir)).map_err(|err| {
181            Error::Backend(format!(
182                "failed to create changefeed dir {}: {err}",
183                Self::changefeed_dir(&config.data_dir).display()
184            ))
185        })?;
186        Ok(Self { config })
187    }
188
189    fn records_dir(data_dir: &Path) -> PathBuf {
190        data_dir.join("records")
191    }
192
193    fn idempotency_dir(data_dir: &Path) -> PathBuf {
194        data_dir.join("idempotency")
195    }
196
197    fn versions_dir(data_dir: &Path) -> PathBuf {
198        data_dir.join("versions")
199    }
200
201    fn changefeed_dir(data_dir: &Path) -> PathBuf {
202        data_dir.join("changefeed")
203    }
204
205    fn record_path(&self, record_id: &str) -> PathBuf {
206        Self::records_dir(&self.config.data_dir).join(format!("{}.json", hex_key(record_id)))
207    }
208
209    fn version_path(&self, record_id: &str, observed_at_unix_ms: u64) -> PathBuf {
210        Self::versions_dir(&self.config.data_dir).join(format!(
211            "{}-{:020}.json",
212            hex_key(record_id),
213            observed_at_unix_ms
214        ))
215    }
216
217    fn changefeed_path(&self, sequence: u64) -> PathBuf {
218        Self::changefeed_dir(&self.config.data_dir).join(format!("{sequence:020}.json"))
219    }
220
221    fn idempotency_scope_key(scope: &MemoryScope, key: &str) -> String {
222        format!(
223            "{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}",
224            scope.tenant_id,
225            scope.namespace,
226            scope.actor_id,
227            scope.conversation_id.as_deref().unwrap_or(""),
228            scope.session_id.as_deref().unwrap_or(""),
229            key
230        )
231    }
232
233    fn idempotency_path(&self, scope: &MemoryScope, key: &str) -> PathBuf {
234        let scoped = Self::idempotency_scope_key(scope, key);
235        Self::idempotency_dir(&self.config.data_dir).join(format!("{}.txt", hex_key(&scoped)))
236    }
237
238    fn validate_record(&self, record: &MemoryRecord) -> Result<()> {
239        record.validate()
240    }
241
242    fn validate_upsert_request(&self, request: &UpsertRequest) -> Result<()> {
243        self.validate_record(&request.record)?;
244        if request.idempotency_key.is_none()
245            && self
246                .config
247                .engine_config
248                .ingestion
249                .idempotent_writes_required
250        {
251            return Err(Error::InvalidRequest(
252                "idempotency_key is required by the current ingestion policy".to_string(),
253            ));
254        }
255        if self.config.engine_config.ingestion.require_source_labels
256            && request.record.scope.labels.is_empty()
257        {
258            return Err(Error::InvalidRequest(
259                "at least one source label is required by the current ingestion policy".to_string(),
260            ));
261        }
262        Ok(())
263    }
264
265    fn validate_delete_request(&self, request: &DeleteRequest) -> Result<()> {
266        if request.tenant_id.trim().is_empty() {
267            return Err(Error::InvalidRequest(
268                "delete tenant_id is required".to_string(),
269            ));
270        }
271        if request.namespace.trim().is_empty() {
272            return Err(Error::InvalidRequest(
273                "delete namespace is required".to_string(),
274            ));
275        }
276        if request.record_id.trim().is_empty() {
277            return Err(Error::InvalidRequest(
278                "delete record_id is required".to_string(),
279            ));
280        }
281        if request.audit_reason.trim().is_empty() {
282            return Err(Error::InvalidRequest(
283                "delete audit_reason is required".to_string(),
284            ));
285        }
286        Ok(())
287    }
288
289    fn validate_archive_request(&self, request: &ArchiveRequest) -> Result<()> {
290        Self::validate_lifecycle_request(
291            "archive",
292            &request.tenant_id,
293            &request.namespace,
294            &request.record_id,
295            &request.audit_reason,
296        )
297    }
298
299    fn validate_suppress_request(&self, request: &SuppressRequest) -> Result<()> {
300        Self::validate_lifecycle_request(
301            "suppress",
302            &request.tenant_id,
303            &request.namespace,
304            &request.record_id,
305            &request.audit_reason,
306        )
307    }
308
309    fn validate_recover_request(&self, request: &RecoverRequest) -> Result<()> {
310        Self::validate_lifecycle_request(
311            "recover",
312            &request.tenant_id,
313            &request.namespace,
314            &request.record_id,
315            &request.audit_reason,
316        )?;
317        match request.quality_state {
318            MemoryQualityState::Active | MemoryQualityState::Verified => {}
319            _ => {
320                return Err(Error::InvalidRequest(
321                    "recover quality_state must be Active or Verified".to_string(),
322                ));
323            }
324        }
325        if matches!(
326            request.historical_state,
327            Some(MemoryHistoricalState::Superseded)
328        ) {
329            return Err(Error::InvalidRequest(
330                "recover historical_state cannot be Superseded".to_string(),
331            ));
332        }
333        Ok(())
334    }
335
336    fn validate_lifecycle_request(
337        action: &str,
338        tenant_id: &str,
339        namespace: &str,
340        record_id: &str,
341        audit_reason: &str,
342    ) -> Result<()> {
343        if tenant_id.trim().is_empty() {
344            return Err(Error::InvalidRequest(format!(
345                "{action} tenant_id is required"
346            )));
347        }
348        if namespace.trim().is_empty() {
349            return Err(Error::InvalidRequest(format!(
350                "{action} namespace is required"
351            )));
352        }
353        if record_id.trim().is_empty() {
354            return Err(Error::InvalidRequest(format!(
355                "{action} record_id is required"
356            )));
357        }
358        if audit_reason.trim().is_empty() {
359            return Err(Error::InvalidRequest(format!(
360                "{action} audit_reason is required"
361            )));
362        }
363        Ok(())
364    }
365
366    fn validate_record_scope(
367        stored: &StoredRecord,
368        tenant_id: &str,
369        namespace: &str,
370    ) -> Result<()> {
371        if stored.record.scope.tenant_id != tenant_id {
372            return Err(Error::InvalidRequest(format!(
373                "record {} does not belong to tenant {}",
374                stored.record.id, tenant_id
375            )));
376        }
377        if stored.record.scope.namespace != namespace {
378            return Err(Error::InvalidRequest(format!(
379                "record {} does not belong to namespace {}",
380                stored.record.id, namespace
381            )));
382        }
383        Ok(())
384    }
385
386    fn validate_import_request(
387        &self,
388        request: &ImportRequest,
389    ) -> (u64, bool, Vec<ImportFailure>, Vec<PortableRecord>) {
390        let mut validated_records = 0u64;
391        let mut failures = Vec::new();
392        let mut entries = Vec::with_capacity(request.package.records.len());
393
394        if request.package.package_version != PORTABLE_PACKAGE_VERSION {
395            failures.push(ImportFailure {
396                record_id: None,
397                reason: format!(
398                    "unsupported portable package version {}; expected {}",
399                    request.package.package_version, PORTABLE_PACKAGE_VERSION
400                ),
401            });
402        }
403
404        if request.package.manifest.record_count != request.package.records.len() as u64 {
405            failures.push(ImportFailure {
406                record_id: None,
407                reason: format!(
408                    "portable package manifest record_count={} does not match payload records={}",
409                    request.package.manifest.record_count,
410                    request.package.records.len()
411                ),
412            });
413        }
414
415        for entry in &request.package.records {
416            match self.validate_record(&entry.record) {
417                Ok(()) => {
418                    validated_records += 1;
419                    entries.push(entry.clone());
420                }
421                Err(error) => failures.push(ImportFailure {
422                    record_id: Some(entry.record.id.clone()),
423                    reason: error.to_string(),
424                }),
425            }
426        }
427
428        (validated_records, failures.is_empty(), failures, entries)
429    }
430
431    fn is_pinned(record: &MemoryRecord) -> bool {
432        record.scope.trust_level == MemoryTrustLevel::Pinned
433    }
434
435    fn retention_exempt(&self, record: &MemoryRecord) -> bool {
436        self.config.engine_config.retention.pinned_records_exempt && Self::is_pinned(record)
437    }
438
439    fn now_unix_ms() -> Result<u64> {
440        SystemTime::now()
441            .duration_since(UNIX_EPOCH)
442            .map_err(|err| Error::Backend(format!("system clock error: {err}")))
443            .map(|value| value.as_millis() as u64)
444    }
445
446    fn load_record(&self, record_id: &str) -> Result<Option<StoredRecord>> {
447        let path = self.record_path(record_id);
448        if !path.exists() {
449            return Ok(None);
450        }
451        let raw = fs::read(&path).map_err(|err| {
452            Error::Backend(format!("failed to read record {}: {err}", path.display()))
453        })?;
454        let stored = serde_json::from_slice::<StoredRecord>(&raw).map_err(|err| {
455            Error::Backend(format!("failed to decode record {}: {err}", path.display()))
456        })?;
457        Ok(Some(stored))
458    }
459
460    fn persist_record(&self, stored: &StoredRecord) -> Result<()> {
461        let path = self.record_path(&stored.record.id);
462        let encoded = serde_json::to_vec(stored)
463            .map_err(|err| Error::Backend(format!("failed to encode record: {err}")))?;
464        fs::write(&path, encoded).map_err(|err| {
465            Error::Backend(format!("failed to write record {}: {err}", path.display()))
466        })?;
467        Ok(())
468    }
469
470    fn persist_record_version(
471        &self,
472        record_id: &str,
473        observed_at_unix_ms: u64,
474        stored: Option<&StoredRecord>,
475    ) -> Result<()> {
476        let version = VersionedStoredRecord {
477            record_id: record_id.to_string(),
478            observed_at_unix_ms,
479            record: stored.map(|value| value.record.clone()),
480            idempotency_key: stored.and_then(|value| value.idempotency_key.clone()),
481        };
482        let path = self.version_path(record_id, observed_at_unix_ms);
483        let encoded = serde_json::to_vec(&version)
484            .map_err(|err| Error::Backend(format!("failed to encode record version: {err}")))?;
485        fs::write(&path, encoded).map_err(|err| {
486            Error::Backend(format!(
487                "failed to write record version {}: {err}",
488                path.display()
489            ))
490        })?;
491        Ok(())
492    }
493
494    fn next_changefeed_sequence(&self) -> Result<u64> {
495        let mut max_sequence = 0u64;
496        for entry in fs::read_dir(Self::changefeed_dir(&self.config.data_dir)).map_err(|err| {
497            Error::Backend(format!(
498                "failed to read changefeed dir {}: {err}",
499                Self::changefeed_dir(&self.config.data_dir).display()
500            ))
501        })? {
502            let entry = entry.map_err(|err| {
503                Error::Backend(format!("failed to iterate changefeed dir: {err}"))
504            })?;
505            let path = entry.path();
506            let Some(stem) = path.file_stem().and_then(|value| value.to_str()) else {
507                continue;
508            };
509            if let Ok(sequence) = stem.parse::<u64>() {
510                max_sequence = max_sequence.max(sequence);
511            }
512        }
513        Ok(max_sequence.saturating_add(1))
514    }
515
516    #[allow(clippy::too_many_arguments)]
517    fn append_changefeed_event(
518        &self,
519        kind: ChangefeedEventKind,
520        record: Option<&MemoryRecord>,
521        tenant_id: &str,
522        namespace: &str,
523        record_id: Option<&str>,
524        summary: Option<String>,
525        occurred_at_unix_ms: u64,
526    ) -> Result<()> {
527        let sequence = self.next_changefeed_sequence()?;
528        let event = ChangefeedEvent {
529            sequence,
530            event_id: format!("change:{}:{sequence}", self.backend_kind()),
531            kind,
532            tenant_id: tenant_id.to_string(),
533            namespace: namespace.to_string(),
534            record_id: record_id.map(ToOwned::to_owned),
535            occurred_at_unix_ms,
536            summary,
537            record: record.cloned(),
538        };
539        let path = self.changefeed_path(sequence);
540        let encoded = serde_json::to_vec(&event)
541            .map_err(|err| Error::Backend(format!("failed to encode changefeed event: {err}")))?;
542        fs::write(&path, encoded).map_err(|err| {
543            Error::Backend(format!(
544                "failed to write changefeed event {}: {err}",
545                path.display()
546            ))
547        })?;
548        Ok(())
549    }
550
551    fn iterate_record_versions(&self) -> Result<Vec<VersionedStoredRecord>> {
552        let mut versions = Vec::new();
553        let dir = Self::versions_dir(&self.config.data_dir);
554        for entry in fs::read_dir(&dir).map_err(|err| {
555            Error::Backend(format!(
556                "failed to read versions dir {}: {err}",
557                dir.display()
558            ))
559        })? {
560            let entry = entry
561                .map_err(|err| Error::Backend(format!("failed to iterate versions dir: {err}")))?;
562            let path = entry.path();
563            if path.extension().and_then(|ext| ext.to_str()) != Some("json") {
564                continue;
565            }
566            let raw = fs::read(&path).map_err(|err| {
567                Error::Backend(format!(
568                    "failed to read record version {}: {err}",
569                    path.display()
570                ))
571            })?;
572            let version = serde_json::from_slice::<VersionedStoredRecord>(&raw).map_err(|err| {
573                Error::Backend(format!(
574                    "failed to decode record version {}: {err}",
575                    path.display()
576                ))
577            })?;
578            versions.push(version);
579        }
580        Ok(versions)
581    }
582
583    fn records_as_of(&self, as_of_unix_ms: u64) -> Result<Vec<StoredRecord>> {
584        let mut latest = BTreeMap::<String, VersionedStoredRecord>::new();
585        for version in self.iterate_record_versions()? {
586            if version.observed_at_unix_ms > as_of_unix_ms {
587                continue;
588            }
589            let replace = latest
590                .get(&version.record_id)
591                .is_none_or(|existing| existing.observed_at_unix_ms <= version.observed_at_unix_ms);
592            if replace {
593                latest.insert(version.record_id.clone(), version);
594            }
595        }
596
597        Ok(latest
598            .into_values()
599            .filter_map(|version| {
600                version.record.map(|record| StoredRecord {
601                    record,
602                    idempotency_key: version.idempotency_key,
603                })
604            })
605            .collect())
606    }
607
608    fn persist_imported_record(&self, stored: &StoredRecord) -> Result<()> {
609        self.persist_record(stored)?;
610        if let Some(idempotency_key) = &stored.idempotency_key {
611            fs::write(
612                self.idempotency_path(&stored.record.scope, idempotency_key),
613                stored.record.id.as_bytes(),
614            )
615            .map_err(|err| Error::Backend(format!("failed to write idempotency mapping: {err}")))?;
616        }
617        Ok(())
618    }
619
620    fn remove_record(&self, record_id: &str) -> Result<bool> {
621        let path = self.record_path(record_id);
622        if !path.exists() {
623            return Ok(false);
624        }
625        fs::remove_file(&path).map_err(|err| {
626            Error::Backend(format!("failed to remove record {}: {err}", path.display()))
627        })?;
628        Ok(true)
629    }
630
631    fn remove_idempotency_mapping(&self, stored: &StoredRecord) -> Result<()> {
632        if let Some(idempotency_key) = &stored.idempotency_key {
633            let path = self.idempotency_path(&stored.record.scope, idempotency_key);
634            if path.exists() {
635                fs::remove_file(&path).map_err(|err| {
636                    Error::Backend(format!(
637                        "failed to remove idempotency mapping {}: {err}",
638                        path.display()
639                    ))
640                })?;
641            }
642        }
643        Ok(())
644    }
645
646    fn clear_all_records(&self) -> Result<()> {
647        for stored in self.iterate_records()? {
648            self.remove_record(&stored.record.id)?;
649            self.remove_idempotency_mapping(&stored)?;
650        }
651        Ok(())
652    }
653
654    fn iterate_records(&self) -> Result<Vec<StoredRecord>> {
655        let mut records = Vec::new();
656        let dir = Self::records_dir(&self.config.data_dir);
657        for entry in fs::read_dir(&dir).map_err(|err| {
658            Error::Backend(format!(
659                "failed to read records dir {}: {err}",
660                dir.display()
661            ))
662        })? {
663            let entry = entry.map_err(|err| {
664                Error::Backend(format!(
665                    "failed to iterate records dir {}: {err}",
666                    dir.display()
667                ))
668            })?;
669            let path = entry.path();
670            if path.extension().and_then(|ext| ext.to_str()) != Some("json") {
671                continue;
672            }
673            let raw = fs::read(&path).map_err(|err| {
674                Error::Backend(format!("failed to read record {}: {err}", path.display()))
675            })?;
676            let stored = serde_json::from_slice::<StoredRecord>(&raw).map_err(|err| {
677                Error::Backend(format!("failed to decode record {}: {err}", path.display()))
678            })?;
679            records.push(stored);
680        }
681        Ok(records)
682    }
683
684    fn iterate_idempotency_mappings(&self) -> Result<Vec<IdempotencyMapping>> {
685        let mut mappings = Vec::new();
686        for entry in fs::read_dir(Self::idempotency_dir(&self.config.data_dir)).map_err(|err| {
687            Error::Backend(format!(
688                "failed to read idempotency dir {}: {err}",
689                Self::idempotency_dir(&self.config.data_dir).display()
690            ))
691        })? {
692            let entry = entry.map_err(|err| {
693                Error::Backend(format!("failed to iterate idempotency dir: {err}"))
694            })?;
695            let path = entry.path();
696            if path.extension().and_then(|value| value.to_str()) != Some("txt") {
697                continue;
698            }
699            let Some(stem) = path.file_stem().and_then(|value| value.to_str()) else {
700                continue;
701            };
702            let Some(scoped_key) = unhex_key(stem) else {
703                continue;
704            };
705            let record_id = fs::read_to_string(&path).map_err(|err| {
706                Error::Backend(format!(
707                    "failed to read idempotency mapping {}: {err}",
708                    path.display()
709                ))
710            })?;
711            mappings.push(IdempotencyMapping {
712                scoped_key,
713                record_id,
714            });
715        }
716        Ok(mappings)
717    }
718
719    fn parse_scope_key(scoped_key: &str) -> Option<ScopeKeyParts> {
720        let parts = scoped_key.split('\u{1f}').collect::<Vec<_>>();
721        if parts.len() != 6 {
722            return None;
723        }
724        Some((
725            parts[0].to_string(),
726            parts[1].to_string(),
727            parts[2].to_string(),
728            (!parts[3].is_empty()).then(|| parts[3].to_string()),
729            (!parts[4].is_empty()).then(|| parts[4].to_string()),
730            parts[5].to_string(),
731        ))
732    }
733
734    fn scope_matches_filters(
735        tenant_id: &str,
736        namespace: &str,
737        tenant_filter: Option<&str>,
738        namespace_filter: Option<&str>,
739    ) -> bool {
740        tenant_filter.is_none_or(|expected| tenant_id == expected)
741            && namespace_filter.is_none_or(|expected| namespace == expected)
742    }
743
744    fn build_integrity_summary(
745        &self,
746        tenant_filter: Option<&str>,
747        namespace_filter: Option<&str>,
748    ) -> Result<IntegritySummary> {
749        let records = self.iterate_records()?;
750        let mappings = self.iterate_idempotency_mappings()?;
751        let filtered_records = records
752            .iter()
753            .filter(|stored| {
754                Self::scope_matches_filters(
755                    &stored.record.scope.tenant_id,
756                    &stored.record.scope.namespace,
757                    tenant_filter,
758                    namespace_filter,
759                )
760            })
761            .collect::<Vec<_>>();
762
763        let mut mapping_lookup = HashMap::new();
764        let mut stale_idempotency_keys = 0u64;
765        let mut scanned_idempotency_keys = 0u64;
766        for mapping in &mappings {
767            let Some((tenant_id, namespace, _, _, _, idempotency_key)) =
768                Self::parse_scope_key(&mapping.scoped_key)
769            else {
770                stale_idempotency_keys += 1;
771                scanned_idempotency_keys += 1;
772                continue;
773            };
774            if !Self::scope_matches_filters(&tenant_id, &namespace, tenant_filter, namespace_filter)
775            {
776                continue;
777            }
778            scanned_idempotency_keys += 1;
779            let Some(stored) = self.load_record(&mapping.record_id)? else {
780                stale_idempotency_keys += 1;
781                continue;
782            };
783            if stored.record.scope.tenant_id != tenant_id
784                || stored.record.scope.namespace != namespace
785                || stored.idempotency_key.as_deref() != Some(idempotency_key.as_str())
786                || Self::idempotency_scope_key(&stored.record.scope, &idempotency_key)
787                    != mapping.scoped_key
788            {
789                stale_idempotency_keys += 1;
790                continue;
791            }
792            mapping_lookup.insert(mapping.scoped_key.clone(), mapping.record_id.clone());
793        }
794
795        let mut duplicate_groups = HashMap::<String, usize>::new();
796        let mut missing_idempotency_keys = 0u64;
797        let mut duplicate_active_records = 0u64;
798        for stored in &filtered_records {
799            if let Some(idempotency_key) = &stored.idempotency_key {
800                let scoped_key = Self::idempotency_scope_key(&stored.record.scope, idempotency_key);
801                if mapping_lookup.get(&scoped_key) != Some(&stored.record.id) {
802                    missing_idempotency_keys += 1;
803                }
804            }
805
806            if !matches!(
807                stored.record.quality_state,
808                MemoryQualityState::Archived
809                    | MemoryQualityState::Deleted
810                    | MemoryQualityState::Suppressed
811            ) {
812                *duplicate_groups
813                    .entry(Self::dedup_signature(&stored.record))
814                    .or_default() += 1;
815            }
816        }
817
818        for group_size in duplicate_groups.into_values() {
819            if group_size > 1 {
820                duplicate_active_records += (group_size - 1) as u64;
821            }
822        }
823
824        Ok(IntegritySummary {
825            scanned_records: filtered_records.len() as u64,
826            scanned_idempotency_keys,
827            stale_idempotency_keys,
828            missing_idempotency_keys,
829            duplicate_active_records,
830        })
831    }
832
833    fn build_stats_report(&self, request: &StoreStatsRequest) -> Result<StoreStatsReport> {
834        let records = self.iterate_records()?;
835        let tenant_filter = request.tenant_id.as_deref();
836        let namespace_filter = request.namespace.as_deref();
837        let filtered_records = records
838            .iter()
839            .filter(|stored| {
840                Self::scope_matches_filters(
841                    &stored.record.scope.tenant_id,
842                    &stored.record.scope.namespace,
843                    tenant_filter,
844                    namespace_filter,
845                )
846            })
847            .collect::<Vec<_>>();
848        let now_unix_ms = Self::now_unix_ms()?;
849        let integrity = self.build_integrity_summary(tenant_filter, namespace_filter)?;
850        let mut namespace_map = BTreeMap::<(String, String), NamespaceStats>::new();
851        let mut duplicate_groups = HashMap::<String, usize>::new();
852        let mut tombstoned_records = 0u64;
853        let mut expired_records = 0u64;
854        let mut historical_records = 0u64;
855        let mut superseded_records = 0u64;
856        let mut lineage_links = 0u64;
857
858        for stored in &filtered_records {
859            let key = (
860                stored.record.scope.tenant_id.clone(),
861                stored.record.scope.namespace.clone(),
862            );
863            let entry = namespace_map
864                .entry(key.clone())
865                .or_insert_with(|| NamespaceStats {
866                    tenant_id: key.0.clone(),
867                    namespace: key.1.clone(),
868                    active_records: 0,
869                    archived_records: 0,
870                    deleted_records: 0,
871                    suppressed_records: 0,
872                    pinned_records: 0,
873                });
874            match stored.record.quality_state {
875                MemoryQualityState::Archived => entry.archived_records += 1,
876                MemoryQualityState::Deleted => {
877                    entry.deleted_records += 1;
878                    tombstoned_records += 1;
879                }
880                MemoryQualityState::Suppressed => entry.suppressed_records += 1,
881                _ => entry.active_records += 1,
882            }
883            if stored.record.scope.trust_level == MemoryTrustLevel::Pinned {
884                entry.pinned_records += 1;
885            }
886            if stored
887                .record
888                .expires_at_unix_ms
889                .is_some_and(|value| value <= now_unix_ms)
890            {
891                expired_records += 1;
892            }
893            if matches!(
894                stored.record.historical_state,
895                MemoryHistoricalState::Historical
896            ) {
897                historical_records += 1;
898            }
899            if matches!(
900                stored.record.historical_state,
901                MemoryHistoricalState::Superseded
902            ) {
903                superseded_records += 1;
904            }
905            lineage_links += stored.record.lineage.len() as u64;
906            if !matches!(
907                stored.record.quality_state,
908                MemoryQualityState::Archived
909                    | MemoryQualityState::Deleted
910                    | MemoryQualityState::Suppressed
911            ) {
912                *duplicate_groups
913                    .entry(Self::dedup_signature(&stored.record))
914                    .or_default() += 1;
915            }
916        }
917
918        let mut duplicate_candidate_groups = 0u64;
919        let mut duplicate_candidate_records = 0u64;
920        for group_size in duplicate_groups.into_values() {
921            if group_size > 1 {
922                duplicate_candidate_groups += 1;
923                duplicate_candidate_records += (group_size - 1) as u64;
924            }
925        }
926
927        Ok(StoreStatsReport {
928            generated_at_unix_ms: now_unix_ms,
929            total_records: filtered_records.len() as u64,
930            storage_bytes: dir_size(&self.config.data_dir)?,
931            namespaces: namespace_map.into_values().collect(),
932            maintenance: MaintenanceStats {
933                duplicate_candidate_groups,
934                duplicate_candidate_records,
935                tombstoned_records,
936                expired_records,
937                stale_idempotency_keys: integrity.stale_idempotency_keys,
938                historical_records,
939                superseded_records,
940                lineage_links,
941            },
942            engine: self.config.engine_config.tuning_info(),
943        })
944    }
945
946    fn matches_scope(candidate: &MemoryScope, query: &MemoryScope) -> bool {
947        candidate.tenant_id == query.tenant_id
948            && candidate.namespace == query.namespace
949            && candidate.actor_id == query.actor_id
950            && (query.conversation_id.is_none()
951                || candidate.conversation_id == query.conversation_id)
952            && (query.session_id.is_none() || candidate.session_id == query.session_id)
953    }
954
955    fn record_passes_filters(
956        record: &MemoryRecord,
957        query: &RecallQuery,
958        relative_bounds: RelativeTemporalBounds,
959    ) -> bool {
960        if !Self::matches_scope(&record.scope, &query.scope) {
961            return false;
962        }
963
964        if let Some(expires_at_unix_ms) = record.expires_at_unix_ms
965            && expires_at_unix_ms <= Self::now_unix_ms().unwrap_or(u64::MAX)
966        {
967            return false;
968        }
969
970        if let Some(min_score) = query.filters.min_importance_score
971            && record.importance_score < min_score
972        {
973            return false;
974        }
975
976        if let Some(source) = &query.filters.source
977            && &record.scope.source != source
978        {
979            return false;
980        }
981
982        if let Some(from_unix_ms) = query.filters.from_unix_ms
983            && record.updated_at_unix_ms < from_unix_ms
984        {
985            return false;
986        }
987
988        if let Some(to_unix_ms) = query.filters.to_unix_ms
989            && record.updated_at_unix_ms > to_unix_ms
990        {
991            return false;
992        }
993
994        if let Some(after_unix_ms) = relative_bounds.after_unix_ms
995            && Self::record_temporal_anchor(record) <= after_unix_ms
996        {
997            return false;
998        }
999
1000        if let Some(before_unix_ms) = relative_bounds.before_unix_ms
1001            && Self::record_temporal_anchor(record) >= before_unix_ms
1002        {
1003            return false;
1004        }
1005
1006        if !query.filters.trust_levels.is_empty()
1007            && !query
1008                .filters
1009                .trust_levels
1010                .contains(&record.scope.trust_level)
1011        {
1012            return false;
1013        }
1014
1015        if !query.filters.required_labels.is_empty()
1016            && !query.filters.required_labels.iter().all(|label| {
1017                record
1018                    .scope
1019                    .labels
1020                    .iter()
1021                    .any(|candidate| candidate == label)
1022            })
1023        {
1024            return false;
1025        }
1026
1027        if !query.filters.kinds.is_empty() && !query.filters.kinds.contains(&record.kind) {
1028            return false;
1029        }
1030
1031        if let Some(episode_id) = &query.filters.episode_id
1032            && record.episode.as_ref().map(|episode| &episode.episode_id) != Some(episode_id)
1033        {
1034            return false;
1035        }
1036
1037        if !query.filters.continuity_states.is_empty()
1038            && !record.episode.as_ref().is_some_and(|episode| {
1039                query
1040                    .filters
1041                    .continuity_states
1042                    .contains(&episode.continuity_state)
1043            })
1044        {
1045            return false;
1046        }
1047
1048        if query.filters.unresolved_only
1049            && !record
1050                .episode
1051                .as_ref()
1052                .is_some_and(|episode| episode.continuity_state.is_unresolved())
1053        {
1054            return false;
1055        }
1056
1057        if let Some(lineage_record_id) = &query.filters.lineage_record_id
1058            && record.id != *lineage_record_id
1059            && !record
1060                .lineage
1061                .iter()
1062                .any(|link| &link.record_id == lineage_record_id)
1063        {
1064            return false;
1065        }
1066
1067        if !query.filters.boundary_labels.is_empty()
1068            && !record.episode.as_ref().is_some_and(|episode| {
1069                episode.boundary_label.as_ref().is_some_and(|label| {
1070                    query
1071                        .filters
1072                        .boundary_labels
1073                        .iter()
1074                        .any(|expected| expected == label)
1075                })
1076            })
1077        {
1078            return false;
1079        }
1080
1081        if let Some(recurrence_key) = &query.filters.recurrence_key
1082            && record
1083                .episode
1084                .as_ref()
1085                .and_then(|episode| episode.recurrence_key.as_ref())
1086                != Some(recurrence_key)
1087        {
1088            return false;
1089        }
1090
1091        if !query.filters.conflict_states.is_empty()
1092            && !record
1093                .conflict
1094                .as_ref()
1095                .is_some_and(|conflict| query.filters.conflict_states.contains(&conflict.state))
1096        {
1097            return false;
1098        }
1099
1100        if !query.filters.resolution_kinds.is_empty()
1101            && !record.conflict.as_ref().is_some_and(|conflict| {
1102                query
1103                    .filters
1104                    .resolution_kinds
1105                    .contains(&conflict.resolution)
1106            })
1107        {
1108            return false;
1109        }
1110
1111        if query.filters.unresolved_conflicts_only
1112            && !record.conflict.as_ref().is_some_and(|conflict| {
1113                matches!(
1114                    conflict.state,
1115                    mnemara_core::ConflictReviewState::PotentialConflict
1116                        | mnemara_core::ConflictReviewState::UnderReview
1117                )
1118            })
1119        {
1120            return false;
1121        }
1122
1123        match query.filters.historical_mode {
1124            RecallHistoricalMode::CurrentOnly => {
1125                if !matches!(record.historical_state, MemoryHistoricalState::Current) {
1126                    return false;
1127                }
1128            }
1129            RecallHistoricalMode::HistoricalOnly => {
1130                if matches!(record.historical_state, MemoryHistoricalState::Current) {
1131                    return false;
1132                }
1133            }
1134            RecallHistoricalMode::IncludeHistorical => {}
1135        }
1136
1137        if !query.filters.states.is_empty() {
1138            if !query.filters.states.contains(&record.quality_state) {
1139                return false;
1140            }
1141        } else {
1142            match record.quality_state {
1143                MemoryQualityState::Archived if !query.filters.include_archived => return false,
1144                MemoryQualityState::Deleted | MemoryQualityState::Suppressed => return false,
1145                _ => {}
1146            }
1147        }
1148
1149        true
1150    }
1151
1152    fn relative_temporal_bounds(
1153        records: &[StoredRecord],
1154        query: &RecallQuery,
1155    ) -> Result<RelativeTemporalBounds> {
1156        let mut bounds = RelativeTemporalBounds::default();
1157        if let Some(after_record_id) = &query.filters.after_record_id {
1158            let Some(anchor) = records
1159                .iter()
1160                .find(|stored| {
1161                    stored.record.id == *after_record_id
1162                        && Self::matches_scope(&stored.record.scope, &query.scope)
1163                })
1164                .map(|stored| Self::record_temporal_anchor(&stored.record))
1165            else {
1166                return Err(Error::InvalidRequest(format!(
1167                    "after_record_id '{after_record_id}' was not found in recall scope"
1168                )));
1169            };
1170            bounds.after_unix_ms = Some(anchor);
1171        }
1172        if let Some(before_record_id) = &query.filters.before_record_id {
1173            let Some(anchor) = records
1174                .iter()
1175                .find(|stored| {
1176                    stored.record.id == *before_record_id
1177                        && Self::matches_scope(&stored.record.scope, &query.scope)
1178                })
1179                .map(|stored| Self::record_temporal_anchor(&stored.record))
1180            else {
1181                return Err(Error::InvalidRequest(format!(
1182                    "before_record_id '{before_record_id}' was not found in recall scope"
1183                )));
1184            };
1185            bounds.before_unix_ms = Some(anchor);
1186        }
1187        if let (Some(after), Some(before)) = (bounds.after_unix_ms, bounds.before_unix_ms)
1188            && after >= before
1189        {
1190            return Err(Error::InvalidRequest(
1191                "after_record_id must refer to an earlier record than before_record_id".to_string(),
1192            ));
1193        }
1194        Ok(bounds)
1195    }
1196
1197    fn approximate_tokens(record: &MemoryRecord) -> usize {
1198        let content_tokens = record.content.split_whitespace().count();
1199        let summary_tokens = record
1200            .summary
1201            .as_deref()
1202            .map(|summary| summary.split_whitespace().count())
1203            .unwrap_or(0);
1204        content_tokens + summary_tokens
1205    }
1206
1207    fn record_temporal_anchor(record: &MemoryRecord) -> u64 {
1208        record
1209            .episode
1210            .as_ref()
1211            .and_then(|episode| episode.last_active_unix_ms.or(episode.started_at_unix_ms))
1212            .unwrap_or(record.updated_at_unix_ms)
1213    }
1214
1215    fn selected_channels_for_hit(hit: &RecallHit, empty_query: bool) -> Vec<String> {
1216        let mut selected_channels = if empty_query {
1217            vec!["temporal".to_string(), "policy".to_string()]
1218        } else {
1219            vec!["lexical".to_string(), "policy".to_string()]
1220        };
1221        if hit.breakdown.semantic > 0.0 {
1222            selected_channels.push("semantic".to_string());
1223        }
1224        if hit.breakdown.metadata > 0.0 {
1225            selected_channels.push("metadata".to_string());
1226        }
1227        if hit.breakdown.episodic > 0.0 {
1228            selected_channels.push("episodic".to_string());
1229        }
1230        if hit.breakdown.salience > 0.0 {
1231            selected_channels.push("salience".to_string());
1232        }
1233        if hit.breakdown.curation > 0.0 {
1234            selected_channels.push("curation".to_string());
1235        }
1236        if hit.record.conflict.is_some() {
1237            selected_channels.push("conflict".to_string());
1238        }
1239        selected_channels.sort();
1240        selected_channels.dedup();
1241        selected_channels
1242    }
1243
1244    fn planning_profile_note(profile: RecallPlanningProfile) -> &'static str {
1245        match profile {
1246            RecallPlanningProfile::FastPath => "planning_profile=fast_path",
1247            RecallPlanningProfile::ContinuityAware => "planning_profile=continuity_aware",
1248        }
1249    }
1250
1251    fn dedup_signature(record: &MemoryRecord) -> String {
1252        format!(
1253            "{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}",
1254            record.scope.tenant_id,
1255            record.scope.namespace,
1256            record.scope.actor_id,
1257            record.kind as u8,
1258            record.content.trim().to_ascii_lowercase(),
1259            record
1260                .summary
1261                .clone()
1262                .unwrap_or_default()
1263                .trim()
1264                .to_ascii_lowercase()
1265        )
1266    }
1267
1268    fn summary_record_id(signature: &str) -> String {
1269        let mut hasher = std::collections::hash_map::DefaultHasher::new();
1270        signature.hash(&mut hasher);
1271        format!("compacted-summary-{:016x}", hasher.finish())
1272    }
1273
1274    fn compaction_summary_record(
1275        group: &[StoredRecord],
1276        signature: &str,
1277        now_unix_ms: u64,
1278    ) -> StoredRecord {
1279        let canonical = &group[0].record;
1280        let representative_summary = canonical
1281            .summary
1282            .clone()
1283            .filter(|value| !value.trim().is_empty())
1284            .unwrap_or_else(|| canonical.content.clone());
1285        let cluster_size = group.len();
1286        let max_importance_score = group
1287            .iter()
1288            .map(|stored| stored.record.importance_score)
1289            .fold(canonical.importance_score, f32::max);
1290
1291        let mut metadata = BTreeMap::new();
1292        metadata.insert(
1293            "compaction_reason".to_string(),
1294            "duplicate_cluster_rollup".to_string(),
1295        );
1296        metadata.insert(
1297            "compaction_cluster_size".to_string(),
1298            cluster_size.to_string(),
1299        );
1300        metadata.insert("representative_record_id".to_string(), canonical.id.clone());
1301
1302        let mut labels = canonical.scope.labels.clone();
1303        if !labels.iter().any(|label| label == "compacted") {
1304            labels.push("compacted".to_string());
1305        }
1306
1307        StoredRecord {
1308            record: MemoryRecord {
1309                id: Self::summary_record_id(signature),
1310                scope: MemoryScope {
1311                    tenant_id: canonical.scope.tenant_id.clone(),
1312                    namespace: canonical.scope.namespace.clone(),
1313                    actor_id: canonical.scope.actor_id.clone(),
1314                    conversation_id: canonical.scope.conversation_id.clone(),
1315                    session_id: canonical.scope.session_id.clone(),
1316                    source: canonical.scope.source.clone(),
1317                    labels,
1318                    trust_level: canonical.scope.trust_level,
1319                },
1320                kind: mnemara_core::MemoryRecordKind::Summary,
1321                content: format!(
1322                    "Compacted {} related records into a durable summary. Representative memory: {}",
1323                    cluster_size, representative_summary
1324                ),
1325                summary: Some(format!(
1326                    "{} related records: {}",
1327                    cluster_size, representative_summary
1328                )),
1329                source_id: None,
1330                metadata,
1331                quality_state: if matches!(canonical.quality_state, MemoryQualityState::Verified) {
1332                    MemoryQualityState::Verified
1333                } else {
1334                    MemoryQualityState::Active
1335                },
1336                created_at_unix_ms: now_unix_ms,
1337                updated_at_unix_ms: now_unix_ms,
1338                expires_at_unix_ms: None,
1339                importance_score: max_importance_score,
1340                artifact: canonical.artifact.clone(),
1341                episode: canonical.episode.clone(),
1342                historical_state: MemoryHistoricalState::Current,
1343                lineage: group
1344                    .iter()
1345                    .map(|stored| LineageLink {
1346                        record_id: stored.record.id.clone(),
1347                        relation: LineageRelationKind::ConsolidatedFrom,
1348                        confidence: 1.0,
1349                    })
1350                    .collect(),
1351                conflict: canonical.conflict.clone(),
1352            },
1353            idempotency_key: None,
1354        }
1355    }
1356
1357    fn cold_tiering_candidates(
1358        &self,
1359        tenant_id: &str,
1360        namespace: Option<&str>,
1361        now_unix_ms: u64,
1362    ) -> Result<Vec<StoredRecord>> {
1363        let cold_archive_after_days = self.config.engine_config.compaction.cold_archive_after_days;
1364        if cold_archive_after_days == 0 {
1365            return Ok(Vec::new());
1366        }
1367        let archive_threshold_ms =
1368            u64::from(cold_archive_after_days).saturating_mul(24 * 60 * 60 * 1_000);
1369        let max_importance = f32::from(
1370            self.config
1371                .engine_config
1372                .compaction
1373                .cold_archive_importance_threshold_per_mille,
1374        ) / 1000.0;
1375
1376        Ok(self
1377            .iterate_records()?
1378            .into_iter()
1379            .filter(|stored| stored.record.scope.tenant_id == tenant_id)
1380            .filter(|stored| namespace.is_none_or(|value| stored.record.scope.namespace == value))
1381            .filter(|stored| {
1382                matches!(
1383                    stored.record.quality_state,
1384                    MemoryQualityState::Draft
1385                        | MemoryQualityState::Active
1386                        | MemoryQualityState::Verified
1387                )
1388            })
1389            .filter(|stored| {
1390                stored.record.scope.trust_level != mnemara_core::MemoryTrustLevel::Pinned
1391            })
1392            .filter(|stored| {
1393                now_unix_ms.saturating_sub(stored.record.updated_at_unix_ms) > archive_threshold_ms
1394                    && stored.record.importance_score <= max_importance
1395            })
1396            .collect())
1397    }
1398
1399    fn synthesis_group_key(record: &MemoryRecord) -> String {
1400        format!(
1401            "{}\u{1f}{}\u{1f}{}\u{1f}{}\u{1f}{}",
1402            record.scope.namespace,
1403            record.scope.actor_id,
1404            record.scope.conversation_id.as_deref().unwrap_or(""),
1405            record.scope.session_id.as_deref().unwrap_or(""),
1406            record.scope.source
1407        )
1408    }
1409
1410    fn synthesis_record_id(source_ids: &[String]) -> String {
1411        let mut hasher = std::collections::hash_map::DefaultHasher::new();
1412        for source_id in source_ids {
1413            source_id.hash(&mut hasher);
1414        }
1415        format!("synthesis-proposal-{:016x}", hasher.finish())
1416    }
1417
1418    fn synthesis_content(group: &[StoredRecord]) -> String {
1419        let snippets = group
1420            .iter()
1421            .take(5)
1422            .map(|stored| {
1423                stored
1424                    .record
1425                    .summary
1426                    .clone()
1427                    .filter(|value| !value.trim().is_empty())
1428                    .unwrap_or_else(|| stored.record.content.clone())
1429            })
1430            .map(|value| {
1431                let value = value.trim();
1432                let truncated = value.chars().take(160).collect::<String>();
1433                if truncated.len() < value.len() {
1434                    format!("{truncated}...")
1435                } else {
1436                    truncated
1437                }
1438            })
1439            .collect::<Vec<_>>();
1440        format!(
1441            "Synthesized {} source memories into a reviewable summary: {}",
1442            group.len(),
1443            snippets.join("; ")
1444        )
1445    }
1446
1447    fn synthesis_proposal_from_group(
1448        group: &[StoredRecord],
1449        reason: &str,
1450        now_unix_ms: u64,
1451    ) -> SynthesisProposal {
1452        let canonical = &group[0].record;
1453        let source_record_ids = group
1454            .iter()
1455            .map(|stored| stored.record.id.clone())
1456            .collect::<Vec<_>>();
1457        let max_importance_score = group
1458            .iter()
1459            .map(|stored| stored.record.importance_score)
1460            .fold(canonical.importance_score, f32::max);
1461        let mut labels = canonical.scope.labels.clone();
1462        for label in ["synthesized", "review-required"] {
1463            if !labels.iter().any(|existing| existing == label) {
1464                labels.push(label.to_string());
1465            }
1466        }
1467        let content = Self::synthesis_content(group);
1468        let mut metadata = BTreeMap::new();
1469        metadata.insert("synthesis_reason".to_string(), reason.to_string());
1470        metadata.insert(
1471            "synthesis_strategy".to_string(),
1472            "deterministic_scope_rollup".to_string(),
1473        );
1474        metadata.insert(
1475            "synthesis_source_count".to_string(),
1476            group.len().to_string(),
1477        );
1478        metadata.insert("review_state".to_string(), "proposed".to_string());
1479        let confidence = (0.55 + (group.len().min(9) as f32 * 0.05)).min(0.95);
1480        let proposed_record = MemoryRecord {
1481            id: Self::synthesis_record_id(&source_record_ids),
1482            scope: MemoryScope {
1483                tenant_id: canonical.scope.tenant_id.clone(),
1484                namespace: canonical.scope.namespace.clone(),
1485                actor_id: canonical.scope.actor_id.clone(),
1486                conversation_id: canonical.scope.conversation_id.clone(),
1487                session_id: canonical.scope.session_id.clone(),
1488                source: canonical.scope.source.clone(),
1489                labels,
1490                trust_level: MemoryTrustLevel::Derived,
1491            },
1492            kind: MemoryRecordKind::Summary,
1493            content: content.clone(),
1494            summary: Some(content),
1495            source_id: None,
1496            metadata: metadata.clone(),
1497            quality_state: MemoryQualityState::Draft,
1498            created_at_unix_ms: now_unix_ms,
1499            updated_at_unix_ms: now_unix_ms,
1500            expires_at_unix_ms: None,
1501            importance_score: max_importance_score,
1502            artifact: None,
1503            episode: canonical.episode.clone(),
1504            historical_state: MemoryHistoricalState::Current,
1505            lineage: source_record_ids
1506                .iter()
1507                .map(|record_id| LineageLink {
1508                    record_id: record_id.clone(),
1509                    relation: LineageRelationKind::DerivedFrom,
1510                    confidence,
1511                })
1512                .collect(),
1513            conflict: None,
1514        };
1515        SynthesisProposal {
1516            proposed_record,
1517            source_record_ids,
1518            confidence,
1519            rationale: "records share tenant, namespace, actor, conversation, session, and source"
1520                .to_string(),
1521            metadata,
1522        }
1523    }
1524
1525    fn build_explanations(
1526        &self,
1527        scorer: ConfiguredRecallScorer,
1528        planning_profile: RecallPlanningProfile,
1529        query: &RecallQuery,
1530        planned: &[PlannedRecallCandidate],
1531        selected_record_ids: &[String],
1532        trace_id: &str,
1533    ) -> (Vec<RecallHit>, Option<RecallExplanation>) {
1534        let selected_set = selected_record_ids.iter().cloned().collect::<BTreeSet<_>>();
1535        let hits = planned
1536            .iter()
1537            .filter(|candidate| selected_set.contains(&candidate.hit.record.id))
1538            .map(|candidate| {
1539                let mut enriched = candidate.hit.clone();
1540                if query.include_explanation {
1541                    let selected_channels = Self::selected_channels_for_hit(
1542                        &candidate.hit,
1543                        query.query_text.trim().is_empty(),
1544                    );
1545                    enriched.explanation = Some(RecallExplanation {
1546                        selected_channels,
1547                        policy_notes: vec![if query.query_text.trim().is_empty() {
1548                            "recent_scope_scan".to_string()
1549                        } else {
1550                            "initial_file_backend_scoring".to_string()
1551                        }],
1552                        trace_id: Some(trace_id.to_string()),
1553                        planning_trace: None,
1554                        planning_profile: Some(planning_profile),
1555                        policy_profile: Some(scorer.policy_profile()),
1556                        scorer_kind: Some(scorer.scorer_kind()),
1557                        scoring_profile: Some(scorer.scoring_profile()),
1558                    });
1559                    if let Some(explanation) = enriched.explanation.as_mut() {
1560                        explanation
1561                            .policy_notes
1562                            .push(scorer.profile_note().to_string());
1563                        explanation
1564                            .policy_notes
1565                            .push(scorer.policy_profile_note().to_string());
1566                        explanation
1567                            .policy_notes
1568                            .push(Self::planning_profile_note(planning_profile).to_string());
1569                        if let Some(note) = scorer.embedding_note() {
1570                            explanation.policy_notes.push(note.to_string());
1571                        }
1572                        if query.filters.episode_id.is_some() {
1573                            explanation
1574                                .policy_notes
1575                                .push("episode_filter_applied".to_string());
1576                        }
1577                        if query.filters.unresolved_only {
1578                            explanation
1579                                .policy_notes
1580                                .push("unresolved_only_filter_applied".to_string());
1581                        }
1582                        explanation.policy_notes.push(format!(
1583                            "matched_terms={}",
1584                            candidate.matched_terms.join(",")
1585                        ));
1586                    }
1587                }
1588                enriched
1589            })
1590            .collect::<Vec<_>>();
1591
1592        let planning_trace = query.include_explanation.then(|| RecallPlanningTrace {
1593            trace_id: trace_id.to_string(),
1594            token_budget_applied: query.token_budget.is_some(),
1595            candidates: planned
1596                .iter()
1597                .enumerate()
1598                .map(|(index, candidate)| {
1599                    let selected = selected_set.contains(&candidate.hit.record.id);
1600                    let selection_rank = selected_record_ids
1601                        .iter()
1602                        .position(|record_id| record_id == &candidate.hit.record.id)
1603                        .map(|position| position as u32 + 1);
1604                    let selected_channels = Self::selected_channels_for_hit(
1605                        &candidate.hit,
1606                        query.query_text.trim().is_empty(),
1607                    );
1608
1609                    let mut filter_reasons = Vec::new();
1610                    if selected {
1611                        filter_reasons.push("retained".to_string());
1612                    } else {
1613                        if index >= query.max_items {
1614                            filter_reasons.push("max_items_exhausted".to_string());
1615                        }
1616                        if query.token_budget.is_some() {
1617                            filter_reasons.push("token_budget_exhausted".to_string());
1618                        }
1619                    }
1620
1621                    RecallTraceCandidate {
1622                        record_id: candidate.hit.record.id.clone(),
1623                        kind: candidate.hit.record.kind,
1624                        selected,
1625                        planner_stage: candidate.planner_stage,
1626                        candidate_sources: candidate.candidate_sources.clone(),
1627                        relation_reasons: candidate.relation_reasons.clone(),
1628                        selection_rank,
1629                        matched_terms: candidate.matched_terms.clone(),
1630                        selected_channels,
1631                        filter_reasons,
1632                        decision_reason: if selected {
1633                            "selected_by_rank".to_string()
1634                        } else if query.token_budget.is_some() {
1635                            "excluded_by_rank_or_budget".to_string()
1636                        } else {
1637                            "excluded_by_rank".to_string()
1638                        },
1639                        breakdown: candidate.hit.breakdown.clone(),
1640                    }
1641                })
1642                .collect(),
1643        });
1644
1645        let explanation = query.include_explanation.then(|| {
1646            let mut selected_channels = if query.query_text.trim().is_empty() {
1647                vec!["temporal".to_string(), "policy".to_string()]
1648            } else {
1649                vec!["lexical".to_string(), "policy".to_string()]
1650            };
1651            for channel in [
1652                "semantic", "metadata", "episodic", "salience", "curation", "conflict",
1653            ] {
1654                let present = planned.iter().any(|candidate| match channel {
1655                    "semantic" => candidate.hit.breakdown.semantic > 0.0,
1656                    "metadata" => candidate.hit.breakdown.metadata > 0.0,
1657                    "episodic" => candidate.hit.breakdown.episodic > 0.0,
1658                    "salience" => candidate.hit.breakdown.salience > 0.0,
1659                    "curation" => candidate.hit.breakdown.curation > 0.0,
1660                    "conflict" => candidate.hit.record.conflict.is_some(),
1661                    _ => false,
1662                });
1663                if present && !selected_channels.iter().any(|existing| existing == channel) {
1664                    selected_channels.push(channel.to_string());
1665                }
1666            }
1667            let mut policy_notes = vec![if query.query_text.trim().is_empty() {
1668                "recent_scope_scan".to_string()
1669            } else {
1670                "initial_file_backend_scoring".to_string()
1671            }];
1672            policy_notes.push(scorer.profile_note().to_string());
1673            policy_notes.push(scorer.policy_profile_note().to_string());
1674            policy_notes.push(Self::planning_profile_note(planning_profile).to_string());
1675            if let Some(note) = scorer.embedding_note() {
1676                policy_notes.push(note.to_string());
1677            }
1678            if query.filters.episode_id.is_some() {
1679                policy_notes.push("episode_filter_applied".to_string());
1680            }
1681            if query.filters.unresolved_only {
1682                policy_notes.push("unresolved_only_filter_applied".to_string());
1683            }
1684            if query.filters.before_record_id.is_some() || query.filters.after_record_id.is_some() {
1685                policy_notes.push("relative_temporal_filter_applied".to_string());
1686            }
1687            if !query.filters.boundary_labels.is_empty() || query.filters.recurrence_key.is_some() {
1688                policy_notes.push("episodic_boundary_filter_applied".to_string());
1689            }
1690            if !query.filters.conflict_states.is_empty()
1691                || !query.filters.resolution_kinds.is_empty()
1692                || query.filters.unresolved_conflicts_only
1693            {
1694                policy_notes.push("conflict_review_filter_applied".to_string());
1695            }
1696            RecallExplanation {
1697                selected_channels,
1698                policy_notes,
1699                trace_id: Some(trace_id.to_string()),
1700                planning_trace,
1701                planning_profile: Some(planning_profile),
1702                policy_profile: Some(scorer.policy_profile()),
1703                scorer_kind: Some(scorer.scorer_kind()),
1704                scoring_profile: Some(scorer.scoring_profile()),
1705            }
1706        });
1707
1708        (hits, explanation)
1709    }
1710
1711    fn recall_from_stored_records(
1712        &self,
1713        query: RecallQuery,
1714        stored_records: Vec<StoredRecord>,
1715    ) -> Result<RecallResult> {
1716        let empty_query = query.query_text.trim().is_empty();
1717        let planner = self.config.recall_planner();
1718        let scorer = planner.scorer();
1719        let planning_profile = planner.effective_profile(&query);
1720        let relative_bounds = Self::relative_temporal_bounds(&stored_records, &query)?;
1721        let records = stored_records
1722            .into_iter()
1723            .filter(|stored| Self::record_passes_filters(&stored.record, &query, relative_bounds))
1724            .map(|stored| stored.record)
1725            .collect::<Vec<_>>();
1726        let mut scored = planner.plan(&records, &query);
1727        match query.filters.temporal_order {
1728            RecallTemporalOrder::Relevance if empty_query => {
1729                scored.sort_by(|left, right| {
1730                    Self::record_temporal_anchor(&right.hit.record)
1731                        .cmp(&Self::record_temporal_anchor(&left.hit.record))
1732                        .then_with(|| {
1733                            right
1734                                .hit
1735                                .record
1736                                .importance_score
1737                                .total_cmp(&left.hit.record.importance_score)
1738                        })
1739                        .then_with(|| left.hit.record.id.cmp(&right.hit.record.id))
1740                });
1741            }
1742            RecallTemporalOrder::Relevance => {
1743                scored.sort_by(|left, right| {
1744                    right
1745                        .hit
1746                        .breakdown
1747                        .total
1748                        .total_cmp(&left.hit.breakdown.total)
1749                        .then_with(|| left.hit.record.id.cmp(&right.hit.record.id))
1750                });
1751            }
1752            RecallTemporalOrder::ChronologicalAsc => {
1753                scored.sort_by(|left, right| {
1754                    Self::record_temporal_anchor(&left.hit.record)
1755                        .cmp(&Self::record_temporal_anchor(&right.hit.record))
1756                        .then_with(|| {
1757                            right
1758                                .hit
1759                                .breakdown
1760                                .total
1761                                .total_cmp(&left.hit.breakdown.total)
1762                        })
1763                        .then_with(|| left.hit.record.id.cmp(&right.hit.record.id))
1764                });
1765            }
1766            RecallTemporalOrder::ChronologicalDesc => {
1767                scored.sort_by(|left, right| {
1768                    Self::record_temporal_anchor(&right.hit.record)
1769                        .cmp(&Self::record_temporal_anchor(&left.hit.record))
1770                        .then_with(|| {
1771                            right
1772                                .hit
1773                                .breakdown
1774                                .total
1775                                .total_cmp(&left.hit.breakdown.total)
1776                        })
1777                        .then_with(|| left.hit.record.id.cmp(&right.hit.record.id))
1778                });
1779            }
1780        }
1781
1782        let examined = scored.len();
1783        let mut selected_ids = Vec::with_capacity(query.max_items);
1784        let mut remaining_budget = query.token_budget.unwrap_or(usize::MAX);
1785        for candidate in &scored {
1786            if selected_ids.len() >= query.max_items {
1787                break;
1788            }
1789            let estimated_tokens = Self::approximate_tokens(&candidate.hit.record);
1790            if selected_ids.is_empty() || estimated_tokens <= remaining_budget {
1791                remaining_budget = remaining_budget.saturating_sub(estimated_tokens);
1792                selected_ids.push(candidate.hit.record.id.clone());
1793            }
1794        }
1795
1796        let trace_id = format!(
1797            "recall:{}:{}:{}",
1798            query.scope.tenant_id, query.scope.namespace, examined
1799        );
1800        let (hits, explanation) = self.build_explanations(
1801            scorer,
1802            planning_profile,
1803            &query,
1804            &scored,
1805            &selected_ids,
1806            &trace_id,
1807        );
1808        Ok(RecallResult {
1809            hits,
1810            total_candidates_examined: examined,
1811            explanation,
1812        })
1813    }
1814
1815    fn apply_retention_for_namespace(&self, tenant_id: &str, namespace: &str) -> Result<()> {
1816        let now_unix_ms = Self::now_unix_ms()?;
1817        let retention = &self.config.engine_config.retention;
1818        let ttl_window_ms = u64::from(retention.ttl_days).saturating_mul(24 * 60 * 60 * 1_000);
1819        let archive_window_ms =
1820            u64::from(retention.archive_after_days).saturating_mul(24 * 60 * 60 * 1_000);
1821
1822        let mut namespace_records = self
1823            .iterate_records()?
1824            .into_iter()
1825            .filter(|stored| {
1826                stored.record.scope.tenant_id == tenant_id
1827                    && stored.record.scope.namespace == namespace
1828            })
1829            .collect::<Vec<_>>();
1830
1831        for stored in &mut namespace_records {
1832            if self.retention_exempt(&stored.record) {
1833                continue;
1834            }
1835
1836            let expired_by_explicit_deadline = stored
1837                .record
1838                .expires_at_unix_ms
1839                .is_some_and(|expires_at| expires_at <= now_unix_ms);
1840            let expired_by_ttl = ttl_window_ms > 0
1841                && now_unix_ms.saturating_sub(stored.record.created_at_unix_ms) > ttl_window_ms;
1842
1843            if expired_by_explicit_deadline || expired_by_ttl {
1844                self.remove_record(&stored.record.id)?;
1845                self.remove_idempotency_mapping(stored)?;
1846                continue;
1847            }
1848
1849            let should_archive_by_age = archive_window_ms > 0
1850                && now_unix_ms.saturating_sub(stored.record.created_at_unix_ms) > archive_window_ms
1851                && matches!(
1852                    stored.record.quality_state,
1853                    MemoryQualityState::Draft
1854                        | MemoryQualityState::Active
1855                        | MemoryQualityState::Verified
1856                );
1857            if should_archive_by_age && stored.record.quality_state != MemoryQualityState::Archived
1858            {
1859                stored.record.quality_state = MemoryQualityState::Archived;
1860                stored.record.historical_state = MemoryHistoricalState::Historical;
1861                stored.record.updated_at_unix_ms = now_unix_ms;
1862                self.persist_record(stored)?;
1863            }
1864        }
1865
1866        if retention.max_records_per_namespace > 0 {
1867            let mut active = self
1868                .iterate_records()?
1869                .into_iter()
1870                .filter(|stored| {
1871                    stored.record.scope.tenant_id == tenant_id
1872                        && stored.record.scope.namespace == namespace
1873                        && !self.retention_exempt(&stored.record)
1874                        && matches!(
1875                            stored.record.quality_state,
1876                            MemoryQualityState::Draft
1877                                | MemoryQualityState::Active
1878                                | MemoryQualityState::Verified
1879                        )
1880                })
1881                .collect::<Vec<_>>();
1882            if active.len() > retention.max_records_per_namespace {
1883                active.sort_by(|left, right| {
1884                    left.record
1885                        .updated_at_unix_ms
1886                        .cmp(&right.record.updated_at_unix_ms)
1887                        .then_with(|| {
1888                            left.record
1889                                .importance_score
1890                                .total_cmp(&right.record.importance_score)
1891                        })
1892                        .then_with(|| left.record.id.cmp(&right.record.id))
1893                });
1894                let archive_count = active.len() - retention.max_records_per_namespace;
1895                for stored in active.iter_mut().take(archive_count) {
1896                    stored.record.quality_state = MemoryQualityState::Archived;
1897                    stored.record.historical_state = MemoryHistoricalState::Historical;
1898                    stored.record.updated_at_unix_ms = now_unix_ms;
1899                    self.persist_record(stored)?;
1900                }
1901            }
1902        }
1903        Ok(())
1904    }
1905}
1906
1907#[async_trait]
1908impl MemoryStore for FileMemoryStore {
1909    fn backend_kind(&self) -> &'static str {
1910        "file"
1911    }
1912
1913    async fn upsert(&self, request: UpsertRequest) -> Result<UpsertReceipt> {
1914        self.validate_upsert_request(&request)?;
1915
1916        if let Some(idempotency_key) = &request.idempotency_key {
1917            let path = self.idempotency_path(&request.record.scope, idempotency_key);
1918            if path.exists() {
1919                let existing_record_id = fs::read_to_string(&path).map_err(|err| {
1920                    Error::Backend(format!(
1921                        "failed to read idempotency mapping {}: {err}",
1922                        path.display()
1923                    ))
1924                })?;
1925                if existing_record_id != request.record.id {
1926                    return Err(Error::Conflict(format!(
1927                        "idempotency key already belongs to record {}",
1928                        existing_record_id
1929                    )));
1930                }
1931                if self.load_record(&existing_record_id)?.is_some() {
1932                    return Ok(UpsertReceipt {
1933                        record_id: existing_record_id,
1934                        deduplicated: true,
1935                        summary_refreshed: false,
1936                    });
1937                }
1938                fs::remove_file(&path).map_err(|err| {
1939                    Error::Backend(format!(
1940                        "failed to clear stale idempotency mapping {}: {err}",
1941                        path.display()
1942                    ))
1943                })?;
1944            }
1945        }
1946
1947        let key = request.record.id.clone();
1948        let tenant_id = request.record.scope.tenant_id.clone();
1949        let namespace = request.record.scope.namespace.clone();
1950        let observed_at_unix_ms = request.record.updated_at_unix_ms;
1951        let existing = self.load_record(&key)?;
1952        let deduplicated = existing.is_some();
1953        let stored = StoredRecord {
1954            record: request.record,
1955            idempotency_key: request.idempotency_key,
1956        };
1957        if let Some(existing) = existing
1958            && existing.idempotency_key != stored.idempotency_key
1959        {
1960            self.remove_idempotency_mapping(&existing)?;
1961        }
1962        self.persist_record(&stored)?;
1963        if let Some(idempotency_key) = &stored.idempotency_key {
1964            fs::write(
1965                self.idempotency_path(&stored.record.scope, idempotency_key),
1966                key.as_bytes(),
1967            )
1968            .map_err(|err| Error::Backend(format!("failed to write idempotency mapping: {err}")))?;
1969        }
1970        self.persist_record_version(&key, observed_at_unix_ms, Some(&stored))?;
1971        self.append_changefeed_event(
1972            ChangefeedEventKind::Upserted,
1973            Some(&stored.record),
1974            &tenant_id,
1975            &namespace,
1976            Some(&key),
1977            Some(if deduplicated {
1978                "record replaced".to_string()
1979            } else {
1980                "record inserted".to_string()
1981            }),
1982            observed_at_unix_ms,
1983        )?;
1984        self.apply_retention_for_namespace(&tenant_id, &namespace)?;
1985        Ok(UpsertReceipt {
1986            record_id: key,
1987            deduplicated,
1988            summary_refreshed: false,
1989        })
1990    }
1991
1992    async fn batch_upsert(&self, request: BatchUpsertRequest) -> Result<Vec<UpsertReceipt>> {
1993        if request.requests.len() > self.config.engine_config.max_batch_size {
1994            return Err(Error::InvalidRequest(format!(
1995                "batch size {} exceeds configured max_batch_size {}",
1996                request.requests.len(),
1997                self.config.engine_config.max_batch_size
1998            )));
1999        }
2000        for request in &request.requests {
2001            self.validate_upsert_request(request)?;
2002        }
2003        let mut receipts = Vec::with_capacity(request.requests.len());
2004        for request in request.requests {
2005            receipts.push(self.upsert(request).await?);
2006        }
2007        Ok(receipts)
2008    }
2009
2010    async fn recall(&self, query: RecallQuery) -> Result<RecallResult> {
2011        self.recall_from_stored_records(query, self.iterate_records()?)
2012    }
2013
2014    async fn recall_as_of(&self, request: TimeTravelRecallRequest) -> Result<RecallResult> {
2015        self.recall_from_stored_records(request.query, self.records_as_of(request.as_of_unix_ms)?)
2016    }
2017
2018    async fn compact(&self, request: CompactionRequest) -> Result<CompactionReport> {
2019        if request.tenant_id.trim().is_empty() {
2020            return Err(Error::InvalidRequest(
2021                "compaction tenant_id is required".to_string(),
2022            ));
2023        }
2024
2025        let mut groups: HashMap<String, Vec<StoredRecord>> = HashMap::new();
2026        for stored in self.iterate_records()? {
2027            if stored.record.scope.tenant_id != request.tenant_id {
2028                continue;
2029            }
2030            if let Some(namespace) = &request.namespace
2031                && stored.record.scope.namespace != *namespace
2032            {
2033                continue;
2034            }
2035            if matches!(
2036                stored.record.quality_state,
2037                MemoryQualityState::Archived
2038                    | MemoryQualityState::Deleted
2039                    | MemoryQualityState::Suppressed
2040            ) {
2041                continue;
2042            }
2043            groups
2044                .entry(Self::dedup_signature(&stored.record))
2045                .or_default()
2046                .push(stored);
2047        }
2048
2049        let mut deduplicated_records = 0u64;
2050        let mut archived_records = 0u64;
2051        let mut summarized_clusters = 0u64;
2052        let mut superseded_records = 0u64;
2053        let mut lineage_links_created = 0u64;
2054        let now_unix_ms = Self::now_unix_ms()?;
2055        for group in groups.values_mut() {
2056            if group.len() < 2 {
2057                continue;
2058            }
2059            group.sort_by(|left, right| {
2060                right
2061                    .record
2062                    .updated_at_unix_ms
2063                    .cmp(&left.record.updated_at_unix_ms)
2064                    .then_with(|| {
2065                        right
2066                            .record
2067                            .importance_score
2068                            .total_cmp(&left.record.importance_score)
2069                    })
2070                    .then_with(|| left.record.id.cmp(&right.record.id))
2071            });
2072            let signature = Self::dedup_signature(&group[0].record);
2073            if self
2074                .config
2075                .engine_config
2076                .compaction
2077                .summarize_after_record_count
2078                > 0
2079                && group.len()
2080                    >= self
2081                        .config
2082                        .engine_config
2083                        .compaction
2084                        .summarize_after_record_count
2085            {
2086                summarized_clusters += 1;
2087                lineage_links_created += group.len() as u64;
2088                if !request.dry_run {
2089                    let summary = Self::compaction_summary_record(group, &signature, now_unix_ms);
2090                    self.persist_record(&summary)?;
2091                }
2092            }
2093            for duplicate in group.iter_mut().skip(1) {
2094                deduplicated_records += 1;
2095                archived_records += 1;
2096                superseded_records += 1;
2097                if request.dry_run {
2098                    continue;
2099                }
2100                duplicate.record.quality_state = MemoryQualityState::Archived;
2101                duplicate.record.historical_state = MemoryHistoricalState::Superseded;
2102                duplicate.record.lineage.push(LineageLink {
2103                    record_id: Self::summary_record_id(&signature),
2104                    relation: LineageRelationKind::SupersededBy,
2105                    confidence: 1.0,
2106                });
2107                lineage_links_created += 1;
2108                duplicate.record.updated_at_unix_ms = Self::now_unix_ms()?;
2109                self.persist_record(duplicate)?;
2110            }
2111        }
2112
2113        for mut candidate in self.cold_tiering_candidates(
2114            &request.tenant_id,
2115            request.namespace.as_deref(),
2116            now_unix_ms,
2117        )? {
2118            archived_records += 1;
2119            if request.dry_run {
2120                continue;
2121            }
2122            candidate.record.quality_state = MemoryQualityState::Archived;
2123            candidate.record.historical_state = MemoryHistoricalState::Historical;
2124            candidate.record.updated_at_unix_ms = now_unix_ms;
2125            self.persist_record(&candidate)?;
2126        }
2127
2128        Ok(CompactionReport {
2129            deduplicated_records,
2130            archived_records,
2131            summarized_clusters,
2132            pruned_graph_edges: 0,
2133            superseded_records,
2134            lineage_links_created,
2135            dry_run: request.dry_run,
2136        })
2137    }
2138
2139    async fn synthesize(&self, request: SynthesisRequest) -> Result<SynthesisReport> {
2140        if request.tenant_id.trim().is_empty() {
2141            return Err(Error::InvalidRequest(
2142                "synthesis tenant_id is required".to_string(),
2143            ));
2144        }
2145        if request.min_source_records == 0 {
2146            return Err(Error::InvalidRequest(
2147                "synthesis min_source_records must be greater than zero".to_string(),
2148            ));
2149        }
2150        if request.max_source_records < request.min_source_records {
2151            return Err(Error::InvalidRequest(
2152                "synthesis max_source_records must be greater than or equal to min_source_records"
2153                    .to_string(),
2154            ));
2155        }
2156
2157        let mut scanned_records = 0u64;
2158        let mut groups: HashMap<String, Vec<StoredRecord>> = HashMap::new();
2159        for stored in self.iterate_records()? {
2160            scanned_records += 1;
2161            let record = &stored.record;
2162            if record.scope.tenant_id != request.tenant_id {
2163                continue;
2164            }
2165            if let Some(namespace) = &request.namespace
2166                && record.scope.namespace != *namespace
2167            {
2168                continue;
2169            }
2170            if let Some(actor_id) = &request.actor_id
2171                && record.scope.actor_id != *actor_id
2172            {
2173                continue;
2174            }
2175            if let Some(conversation_id) = &request.conversation_id
2176                && record.scope.conversation_id.as_deref() != Some(conversation_id.as_str())
2177            {
2178                continue;
2179            }
2180            if let Some(session_id) = &request.session_id
2181                && record.scope.session_id.as_deref() != Some(session_id.as_str())
2182            {
2183                continue;
2184            }
2185            if request
2186                .from_unix_ms
2187                .is_some_and(|from| record.updated_at_unix_ms < from)
2188                || request
2189                    .to_unix_ms
2190                    .is_some_and(|to| record.updated_at_unix_ms > to)
2191            {
2192                continue;
2193            }
2194            if record.kind == MemoryRecordKind::Summary
2195                || matches!(
2196                    record.quality_state,
2197                    MemoryQualityState::Archived
2198                        | MemoryQualityState::Deleted
2199                        | MemoryQualityState::Suppressed
2200                )
2201            {
2202                continue;
2203            }
2204            groups
2205                .entry(Self::synthesis_group_key(record))
2206                .or_default()
2207                .push(stored);
2208        }
2209
2210        let now_unix_ms = Self::now_unix_ms()?;
2211        let mut eligible_records = 0u64;
2212        let mut proposals = Vec::new();
2213        for group in groups.values_mut() {
2214            if group.len() < request.min_source_records {
2215                continue;
2216            }
2217            group.sort_by(|left, right| {
2218                right
2219                    .record
2220                    .updated_at_unix_ms
2221                    .cmp(&left.record.updated_at_unix_ms)
2222                    .then_with(|| {
2223                        right
2224                            .record
2225                            .importance_score
2226                            .total_cmp(&left.record.importance_score)
2227                    })
2228                    .then_with(|| left.record.id.cmp(&right.record.id))
2229            });
2230            group.truncate(request.max_source_records);
2231            eligible_records += group.len() as u64;
2232            proposals.push(Self::synthesis_proposal_from_group(
2233                group,
2234                &request.reason,
2235                now_unix_ms,
2236            ));
2237        }
2238        proposals.sort_by(|left, right| {
2239            right
2240                .proposed_record
2241                .importance_score
2242                .total_cmp(&left.proposed_record.importance_score)
2243                .then_with(|| left.proposed_record.id.cmp(&right.proposed_record.id))
2244        });
2245        if request.max_proposals > 0 {
2246            proposals.truncate(request.max_proposals);
2247        }
2248
2249        let mut persisted_records = 0u64;
2250        let lineage_links_created = proposals
2251            .iter()
2252            .map(|proposal| proposal.source_record_ids.len() as u64)
2253            .sum();
2254        if !request.dry_run {
2255            for proposal in &proposals {
2256                let stored = StoredRecord {
2257                    record: proposal.proposed_record.clone(),
2258                    idempotency_key: None,
2259                };
2260                self.persist_record(&stored)?;
2261                self.persist_record_version(
2262                    &stored.record.id,
2263                    stored.record.updated_at_unix_ms,
2264                    Some(&stored),
2265                )?;
2266                self.append_changefeed_event(
2267                    ChangefeedEventKind::Upserted,
2268                    Some(&stored.record),
2269                    &stored.record.scope.tenant_id,
2270                    &stored.record.scope.namespace,
2271                    Some(&stored.record.id),
2272                    Some("synthesis proposal persisted".to_string()),
2273                    stored.record.updated_at_unix_ms,
2274                )?;
2275                persisted_records += 1;
2276            }
2277        }
2278
2279        Ok(SynthesisReport {
2280            dry_run: request.dry_run,
2281            scanned_records,
2282            eligible_records,
2283            proposed_records: proposals.len() as u64,
2284            persisted_records,
2285            lineage_links_created,
2286            proposals,
2287        })
2288    }
2289
2290    async fn delete(&self, request: DeleteRequest) -> Result<DeleteReceipt> {
2291        self.validate_delete_request(&request)?;
2292        let Some(stored) = self.load_record(&request.record_id)? else {
2293            return Ok(DeleteReceipt {
2294                record_id: request.record_id,
2295                tombstoned: false,
2296                hard_deleted: false,
2297            });
2298        };
2299        if stored.record.scope.tenant_id != request.tenant_id {
2300            return Err(Error::InvalidRequest(format!(
2301                "record {} does not belong to tenant {}",
2302                request.record_id, request.tenant_id
2303            )));
2304        }
2305        if stored.record.scope.namespace != request.namespace {
2306            return Err(Error::InvalidRequest(format!(
2307                "record {} does not belong to namespace {}",
2308                request.record_id, request.namespace
2309            )));
2310        }
2311
2312        if request.hard_delete {
2313            let now_unix_ms = Self::now_unix_ms()?;
2314            self.remove_record(&request.record_id)?;
2315            self.remove_idempotency_mapping(&stored)?;
2316            self.persist_record_version(&request.record_id, now_unix_ms, None)?;
2317            self.append_changefeed_event(
2318                ChangefeedEventKind::Deleted,
2319                None,
2320                &request.tenant_id,
2321                &request.namespace,
2322                Some(&request.record_id),
2323                Some("record hard deleted".to_string()),
2324                now_unix_ms,
2325            )?;
2326        } else {
2327            let mut tombstone = stored;
2328            tombstone.record.quality_state = MemoryQualityState::Deleted;
2329            tombstone.record.updated_at_unix_ms = Self::now_unix_ms()?;
2330            self.persist_record(&tombstone)?;
2331            self.persist_record_version(
2332                &request.record_id,
2333                tombstone.record.updated_at_unix_ms,
2334                Some(&tombstone),
2335            )?;
2336            self.append_changefeed_event(
2337                ChangefeedEventKind::Deleted,
2338                Some(&tombstone.record),
2339                &request.tenant_id,
2340                &request.namespace,
2341                Some(&request.record_id),
2342                Some("record tombstoned".to_string()),
2343                tombstone.record.updated_at_unix_ms,
2344            )?;
2345        }
2346
2347        Ok(DeleteReceipt {
2348            record_id: request.record_id,
2349            tombstoned: !request.hard_delete,
2350            hard_deleted: request.hard_delete,
2351        })
2352    }
2353
2354    async fn archive(&self, request: ArchiveRequest) -> Result<ArchiveReceipt> {
2355        self.validate_archive_request(&request)?;
2356        let Some(mut stored) = self.load_record(&request.record_id)? else {
2357            return Err(Error::InvalidRequest(format!(
2358                "record {} was not found",
2359                request.record_id
2360            )));
2361        };
2362        Self::validate_record_scope(&stored, &request.tenant_id, &request.namespace)?;
2363
2364        let previous_quality_state = stored.record.quality_state;
2365        let previous_historical_state = stored.record.historical_state;
2366        let changed = previous_quality_state != MemoryQualityState::Archived
2367            || previous_historical_state == MemoryHistoricalState::Current;
2368        let historical_state = match previous_historical_state {
2369            MemoryHistoricalState::Current => MemoryHistoricalState::Historical,
2370            other => other,
2371        };
2372
2373        if changed && !request.dry_run {
2374            stored.record.quality_state = MemoryQualityState::Archived;
2375            stored.record.historical_state = historical_state;
2376            stored.record.updated_at_unix_ms = Self::now_unix_ms()?;
2377            self.persist_record(&stored)?;
2378            self.persist_record_version(
2379                &request.record_id,
2380                stored.record.updated_at_unix_ms,
2381                Some(&stored),
2382            )?;
2383            self.append_changefeed_event(
2384                ChangefeedEventKind::Archived,
2385                Some(&stored.record),
2386                &request.tenant_id,
2387                &request.namespace,
2388                Some(&request.record_id),
2389                Some("record archived".to_string()),
2390                stored.record.updated_at_unix_ms,
2391            )?;
2392        }
2393
2394        Ok(ArchiveReceipt {
2395            record_id: request.record_id,
2396            previous_quality_state,
2397            previous_historical_state,
2398            quality_state: MemoryQualityState::Archived,
2399            historical_state,
2400            changed,
2401            dry_run: request.dry_run,
2402        })
2403    }
2404
2405    async fn suppress(&self, request: SuppressRequest) -> Result<SuppressReceipt> {
2406        self.validate_suppress_request(&request)?;
2407        let Some(mut stored) = self.load_record(&request.record_id)? else {
2408            return Err(Error::InvalidRequest(format!(
2409                "record {} was not found",
2410                request.record_id
2411            )));
2412        };
2413        Self::validate_record_scope(&stored, &request.tenant_id, &request.namespace)?;
2414
2415        let previous_quality_state = stored.record.quality_state;
2416        let previous_historical_state = stored.record.historical_state;
2417        let changed = previous_quality_state != MemoryQualityState::Suppressed;
2418
2419        if changed && !request.dry_run {
2420            stored.record.quality_state = MemoryQualityState::Suppressed;
2421            stored.record.updated_at_unix_ms = Self::now_unix_ms()?;
2422            self.persist_record(&stored)?;
2423            self.persist_record_version(
2424                &request.record_id,
2425                stored.record.updated_at_unix_ms,
2426                Some(&stored),
2427            )?;
2428            self.append_changefeed_event(
2429                ChangefeedEventKind::Suppressed,
2430                Some(&stored.record),
2431                &request.tenant_id,
2432                &request.namespace,
2433                Some(&request.record_id),
2434                Some("record suppressed".to_string()),
2435                stored.record.updated_at_unix_ms,
2436            )?;
2437        }
2438
2439        Ok(SuppressReceipt {
2440            record_id: request.record_id,
2441            previous_quality_state,
2442            previous_historical_state,
2443            quality_state: MemoryQualityState::Suppressed,
2444            historical_state: previous_historical_state,
2445            changed,
2446            dry_run: request.dry_run,
2447        })
2448    }
2449
2450    async fn recover(&self, request: RecoverRequest) -> Result<RecoverReceipt> {
2451        self.validate_recover_request(&request)?;
2452        let Some(mut stored) = self.load_record(&request.record_id)? else {
2453            return Err(Error::InvalidRequest(format!(
2454                "record {} was not found",
2455                request.record_id
2456            )));
2457        };
2458        Self::validate_record_scope(&stored, &request.tenant_id, &request.namespace)?;
2459
2460        let previous_quality_state = stored.record.quality_state;
2461        let previous_historical_state = stored.record.historical_state;
2462        let historical_state = request
2463            .historical_state
2464            .unwrap_or(MemoryHistoricalState::Current);
2465        let changed = previous_quality_state != request.quality_state
2466            || previous_historical_state != historical_state;
2467
2468        if changed && !request.dry_run {
2469            stored.record.quality_state = request.quality_state;
2470            stored.record.historical_state = historical_state;
2471            stored.record.updated_at_unix_ms = Self::now_unix_ms()?;
2472            self.persist_record(&stored)?;
2473            self.persist_record_version(
2474                &request.record_id,
2475                stored.record.updated_at_unix_ms,
2476                Some(&stored),
2477            )?;
2478            self.append_changefeed_event(
2479                ChangefeedEventKind::Recovered,
2480                Some(&stored.record),
2481                &request.tenant_id,
2482                &request.namespace,
2483                Some(&request.record_id),
2484                Some("record recovered".to_string()),
2485                stored.record.updated_at_unix_ms,
2486            )?;
2487        }
2488
2489        Ok(RecoverReceipt {
2490            record_id: request.record_id,
2491            previous_quality_state,
2492            previous_historical_state,
2493            quality_state: request.quality_state,
2494            historical_state,
2495            changed,
2496            dry_run: request.dry_run,
2497        })
2498    }
2499
2500    async fn snapshot(&self) -> Result<SnapshotManifest> {
2501        let records = self.iterate_records()?;
2502        let namespaces = records
2503            .iter()
2504            .map(|stored| stored.record.scope.namespace.clone())
2505            .collect::<BTreeSet<_>>()
2506            .into_iter()
2507            .collect::<Vec<_>>();
2508        let created_at_unix_ms = Self::now_unix_ms()?;
2509        let storage_bytes = dir_size(&self.config.data_dir)?;
2510
2511        Ok(SnapshotManifest {
2512            snapshot_id: format!("snapshot-{created_at_unix_ms}"),
2513            created_at_unix_ms,
2514            namespaces,
2515            record_count: records.len() as u64,
2516            storage_bytes,
2517            engine: self.config.engine_config.tuning_info(),
2518        })
2519    }
2520
2521    async fn stats(&self, request: StoreStatsRequest) -> Result<StoreStatsReport> {
2522        self.build_stats_report(&request)
2523    }
2524
2525    async fn inspect_graph(
2526        &self,
2527        request: GraphInspectionRequest,
2528    ) -> Result<GraphInspectionReport> {
2529        let records = self
2530            .iterate_records()?
2531            .into_iter()
2532            .map(|stored| stored.record)
2533            .collect::<Vec<_>>();
2534        Ok(build_graph_inspection_report(
2535            &records,
2536            &request,
2537            Self::now_unix_ms()?,
2538        ))
2539    }
2540
2541    async fn changefeed(&self, request: ChangefeedRequest) -> Result<ChangefeedReport> {
2542        let limit = request.limit.unwrap_or(100).max(1);
2543        let after_sequence = request.after_sequence.unwrap_or(0);
2544        let mut events = Vec::new();
2545        let mut truncated = false;
2546
2547        for entry in fs::read_dir(Self::changefeed_dir(&self.config.data_dir)).map_err(|err| {
2548            Error::Backend(format!(
2549                "failed to read changefeed dir {}: {err}",
2550                Self::changefeed_dir(&self.config.data_dir).display()
2551            ))
2552        })? {
2553            let entry = entry.map_err(|err| {
2554                Error::Backend(format!("failed to iterate changefeed dir: {err}"))
2555            })?;
2556            let path = entry.path();
2557            if path.extension().and_then(|ext| ext.to_str()) != Some("json") {
2558                continue;
2559            }
2560            let raw = fs::read(&path).map_err(|err| {
2561                Error::Backend(format!(
2562                    "failed to read changefeed event {}: {err}",
2563                    path.display()
2564                ))
2565            })?;
2566            let event = serde_json::from_slice::<ChangefeedEvent>(&raw).map_err(|err| {
2567                Error::Backend(format!(
2568                    "failed to decode changefeed event {}: {err}",
2569                    path.display()
2570                ))
2571            })?;
2572            if event.sequence <= after_sequence
2573                || request
2574                    .tenant_id
2575                    .as_ref()
2576                    .is_some_and(|tenant_id| &event.tenant_id != tenant_id)
2577                || request
2578                    .namespace
2579                    .as_ref()
2580                    .is_some_and(|namespace| &event.namespace != namespace)
2581            {
2582                continue;
2583            }
2584            events.push(event);
2585        }
2586
2587        events.sort_by_key(|event| event.sequence);
2588        if events.len() > limit {
2589            events.truncate(limit);
2590            truncated = true;
2591        }
2592        let last_sequence = events.last().map(|event| event.sequence);
2593        Ok(ChangefeedReport {
2594            events,
2595            last_sequence,
2596            truncated,
2597        })
2598    }
2599
2600    async fn integrity_check(
2601        &self,
2602        request: IntegrityCheckRequest,
2603    ) -> Result<IntegrityCheckReport> {
2604        let summary = self
2605            .build_integrity_summary(request.tenant_id.as_deref(), request.namespace.as_deref())?;
2606        Ok(IntegrityCheckReport {
2607            generated_at_unix_ms: Self::now_unix_ms()?,
2608            healthy: summary.stale_idempotency_keys == 0
2609                && summary.missing_idempotency_keys == 0
2610                && summary.duplicate_active_records == 0,
2611            scanned_records: summary.scanned_records,
2612            scanned_idempotency_keys: summary.scanned_idempotency_keys,
2613            stale_idempotency_keys: summary.stale_idempotency_keys,
2614            missing_idempotency_keys: summary.missing_idempotency_keys,
2615            duplicate_active_records: summary.duplicate_active_records,
2616        })
2617    }
2618
2619    async fn repair(&self, request: RepairRequest) -> Result<RepairReport> {
2620        if request.reason.trim().is_empty() {
2621            return Err(Error::InvalidRequest(
2622                "repair reason is required".to_string(),
2623            ));
2624        }
2625        if !request.remove_stale_idempotency_keys && !request.rebuild_missing_idempotency_keys {
2626            return Err(Error::InvalidRequest(
2627                "repair requires at least one enabled action".to_string(),
2628            ));
2629        }
2630
2631        let tenant_filter = request.tenant_id.as_deref();
2632        let namespace_filter = request.namespace.as_deref();
2633        let summary = self.build_integrity_summary(tenant_filter, namespace_filter)?;
2634        let records = self.iterate_records()?;
2635        let mappings = self.iterate_idempotency_mappings()?;
2636        let mut removed_stale_idempotency_keys = 0u64;
2637        let mut rebuilt_missing_idempotency_keys = 0u64;
2638
2639        if request.remove_stale_idempotency_keys {
2640            for mapping in &mappings {
2641                let Some((tenant_id, namespace, _, _, _, idempotency_key)) =
2642                    Self::parse_scope_key(&mapping.scoped_key)
2643                else {
2644                    continue;
2645                };
2646                if !Self::scope_matches_filters(
2647                    &tenant_id,
2648                    &namespace,
2649                    tenant_filter,
2650                    namespace_filter,
2651                ) {
2652                    continue;
2653                }
2654                let stale = match self.load_record(&mapping.record_id)? {
2655                    Some(stored) => {
2656                        stored.record.scope.tenant_id != tenant_id
2657                            || stored.record.scope.namespace != namespace
2658                            || stored.idempotency_key.as_deref() != Some(idempotency_key.as_str())
2659                            || Self::idempotency_scope_key(&stored.record.scope, &idempotency_key)
2660                                != mapping.scoped_key
2661                    }
2662                    None => true,
2663                };
2664                if stale {
2665                    removed_stale_idempotency_keys += 1;
2666                    if !request.dry_run {
2667                        let path = Self::idempotency_dir(&self.config.data_dir)
2668                            .join(format!("{}.txt", hex_key(&mapping.scoped_key)));
2669                        if path.exists() {
2670                            fs::remove_file(&path).map_err(|err| {
2671                                Error::Backend(format!(
2672                                    "failed to remove stale idempotency mapping {}: {err}",
2673                                    path.display()
2674                                ))
2675                            })?;
2676                        }
2677                    }
2678                }
2679            }
2680        }
2681
2682        if request.rebuild_missing_idempotency_keys {
2683            let existing = self.iterate_idempotency_mappings()?;
2684            let existing_lookup = existing
2685                .into_iter()
2686                .map(|mapping| (mapping.scoped_key, mapping.record_id))
2687                .collect::<HashMap<_, _>>();
2688
2689            for stored in &records {
2690                if !Self::scope_matches_filters(
2691                    &stored.record.scope.tenant_id,
2692                    &stored.record.scope.namespace,
2693                    tenant_filter,
2694                    namespace_filter,
2695                ) {
2696                    continue;
2697                }
2698                let Some(idempotency_key) = &stored.idempotency_key else {
2699                    continue;
2700                };
2701                let scoped_key = Self::idempotency_scope_key(&stored.record.scope, idempotency_key);
2702                if existing_lookup.get(&scoped_key) == Some(&stored.record.id) {
2703                    continue;
2704                }
2705                rebuilt_missing_idempotency_keys += 1;
2706                if !request.dry_run {
2707                    fs::write(
2708                        Self::idempotency_dir(&self.config.data_dir)
2709                            .join(format!("{}.txt", hex_key(&scoped_key))),
2710                        stored.record.id.as_bytes(),
2711                    )
2712                    .map_err(|err| {
2713                        Error::Backend(format!("failed to rebuild idempotency mapping: {err}"))
2714                    })?;
2715                }
2716            }
2717        }
2718
2719        let stale_after = if request.remove_stale_idempotency_keys {
2720            0
2721        } else {
2722            summary.stale_idempotency_keys
2723        };
2724        let missing_after = if request.rebuild_missing_idempotency_keys {
2725            0
2726        } else {
2727            summary.missing_idempotency_keys
2728        };
2729
2730        Ok(RepairReport {
2731            dry_run: request.dry_run,
2732            scanned_records: summary.scanned_records,
2733            scanned_idempotency_keys: summary.scanned_idempotency_keys,
2734            removed_stale_idempotency_keys,
2735            rebuilt_missing_idempotency_keys,
2736            healthy_after: stale_after == 0
2737                && missing_after == 0
2738                && summary.duplicate_active_records == 0,
2739        })
2740    }
2741
2742    async fn export(&self, request: ExportRequest) -> Result<PortableStorePackage> {
2743        let exported_at_unix_ms = Self::now_unix_ms()?;
2744        let mut namespaces = BTreeSet::new();
2745        let mut records = Vec::new();
2746        for stored in self.iterate_records()? {
2747            if request
2748                .tenant_id
2749                .as_deref()
2750                .is_some_and(|tenant_id| stored.record.scope.tenant_id != tenant_id)
2751            {
2752                continue;
2753            }
2754            if request
2755                .namespace
2756                .as_deref()
2757                .is_some_and(|namespace| stored.record.scope.namespace != namespace)
2758            {
2759                continue;
2760            }
2761            if !request.include_archived
2762                && stored.record.quality_state == MemoryQualityState::Archived
2763            {
2764                continue;
2765            }
2766            namespaces.insert(format!(
2767                "{}:{}",
2768                stored.record.scope.tenant_id, stored.record.scope.namespace
2769            ));
2770            records.push(PortableRecord {
2771                record: stored.record,
2772                idempotency_key: stored.idempotency_key,
2773            });
2774        }
2775
2776        let storage_bytes = records
2777            .iter()
2778            .map(|entry| {
2779                entry.record.content.len()
2780                    + entry.record.summary.as_deref().map(str::len).unwrap_or(0)
2781            })
2782            .sum::<usize>() as u64;
2783
2784        Ok(PortableStorePackage {
2785            package_version: PORTABLE_PACKAGE_VERSION,
2786            exported_at_unix_ms,
2787            manifest: SnapshotManifest {
2788                snapshot_id: format!("portable-export-{exported_at_unix_ms}"),
2789                created_at_unix_ms: exported_at_unix_ms,
2790                namespaces: namespaces.into_iter().collect(),
2791                record_count: records.len() as u64,
2792                storage_bytes,
2793                engine: self.config.engine_config.tuning_info(),
2794            },
2795            records,
2796        })
2797    }
2798
2799    async fn import(&self, request: ImportRequest) -> Result<ImportReport> {
2800        let snapshot_id = request.package.manifest.snapshot_id.clone();
2801        let package_version = request.package.package_version;
2802        let (validated_records, compatible_package, failed_records, entries) =
2803            self.validate_import_request(&request);
2804        let apply_changes = compatible_package
2805            && failed_records.is_empty()
2806            && !request.dry_run
2807            && !matches!(request.mode, ImportMode::Validate);
2808        let mut imported_records = 0u64;
2809        let mut skipped_records = 0u64;
2810
2811        if apply_changes && matches!(request.mode, ImportMode::Replace) {
2812            self.clear_all_records()?;
2813        }
2814
2815        for entry in entries {
2816            if matches!(request.mode, ImportMode::Merge)
2817                && self.load_record(&entry.record.id)?.is_some()
2818            {
2819                skipped_records += 1;
2820                continue;
2821            }
2822            if apply_changes {
2823                self.persist_imported_record(&StoredRecord {
2824                    record: entry.record,
2825                    idempotency_key: entry.idempotency_key,
2826                })?;
2827            }
2828            imported_records += 1;
2829        }
2830
2831        Ok(ImportReport {
2832            mode: request.mode,
2833            dry_run: request.dry_run,
2834            applied: apply_changes,
2835            compatible_package,
2836            package_version,
2837            validated_records,
2838            imported_records,
2839            skipped_records,
2840            replaced_existing: matches!(request.mode, ImportMode::Replace),
2841            snapshot_id,
2842            failed_records,
2843        })
2844    }
2845}
2846
2847fn dir_size(path: &Path) -> Result<u64> {
2848    let mut total = 0u64;
2849    for entry in fs::read_dir(path)
2850        .map_err(|err| Error::Backend(format!("failed to read dir {}: {err}", path.display())))?
2851    {
2852        let entry = entry.map_err(|err| {
2853            Error::Backend(format!("failed to iterate dir {}: {err}", path.display()))
2854        })?;
2855        let entry_path = entry.path();
2856        if entry_path.is_dir() {
2857            total = total.saturating_add(dir_size(&entry_path)?);
2858        } else {
2859            total = total.saturating_add(
2860                entry
2861                    .metadata()
2862                    .map_err(|err| {
2863                        Error::Backend(format!(
2864                            "failed to stat file {}: {err}",
2865                            entry_path.display()
2866                        ))
2867                    })?
2868                    .len(),
2869            );
2870        }
2871    }
2872    Ok(total)
2873}
2874
2875fn hex_key(input: &str) -> String {
2876    let mut output = String::with_capacity(input.len() * 2);
2877    for byte in input.as_bytes() {
2878        output.push(nibble_to_hex(byte >> 4));
2879        output.push(nibble_to_hex(byte & 0x0f));
2880    }
2881    output
2882}
2883
2884fn nibble_to_hex(value: u8) -> char {
2885    match value {
2886        0..=9 => (b'0' + value) as char,
2887        10..=15 => (b'a' + (value - 10)) as char,
2888        _ => '0',
2889    }
2890}
2891
2892fn unhex_key(input: &str) -> Option<String> {
2893    if input.len() % 2 != 0 {
2894        return None;
2895    }
2896
2897    let mut bytes = Vec::with_capacity(input.len() / 2);
2898    let chars = input.as_bytes().chunks_exact(2);
2899    for chunk in chars {
2900        let high = hex_to_nibble(chunk[0] as char)?;
2901        let low = hex_to_nibble(chunk[1] as char)?;
2902        bytes.push((high << 4) | low);
2903    }
2904
2905    String::from_utf8(bytes).ok()
2906}
2907
2908fn hex_to_nibble(value: char) -> Option<u8> {
2909    match value {
2910        '0'..='9' => Some(value as u8 - b'0'),
2911        'a'..='f' => Some(value as u8 - b'a' + 10),
2912        'A'..='F' => Some(value as u8 - b'A' + 10),
2913        _ => None,
2914    }
2915}