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}