1use std::collections::HashMap;
15use std::fs;
16use std::path::{Path, PathBuf};
17use std::sync::atomic::{AtomicU64, Ordering};
18use std::sync::{Arc, Mutex};
19use std::time::{SystemTime, UNIX_EPOCH};
20
21use async_trait::async_trait;
22use rusqlite::{Connection, OptionalExtension, params};
23
24use crate::identity_first::AgentIdentity;
25use crate::identity_first::agent_memory::{
26 AgentMemoryError, AgentMemoryForgetResult, AgentMemoryProvider, AgentMemoryRecallRequest,
27 AgentMemoryRecord, AuthoredWriteReceipt, NewAgentMemory, compact_whitespace,
28 decode_path_segment, encode_path_segment, new_memory_id, normalize_tags, read_markdown_records,
29 select_recall_records,
30};
31use crate::memory::taint::LlmWriteGate;
32
33use super::records::{
34 InjectionLogEntry, InjectionSurface, ManifestTier, MemoryAuthor, MemoryId, MemoryKind,
35 MemoryProvenance, MemoryScope, NewMemoryRecord, ProposalId, RecordMeta, RecordStatus,
36 TrustTier, UsageEvent, UsageStats, age_days, content_hash, validate_record_fields,
37};
38use super::staged::{
39 CommitReceipt, DEFAULT_TOMBSTONE_RECREATE_WINDOW_MS, StageToken, StagedBatchKind,
40 StagedBatchView, StagedMemoryStore, StagedMutationBatch, StagedOp, StagedRecordView,
41 validate_batch,
42};
43
44pub const DEFAULT_SCOPE_FLOOR_RECORDS: usize = 4_000;
47pub const DEFAULT_SCOPE_FLOOR_BYTES: usize = 32 * 1024 * 1024;
48
49const STAGE_GC_MAX_AGE_MS: u64 = 24 * 60 * 60 * 1000;
52
53const SQLITE_BUSY_TIMEOUT_MS: u64 = 5_000;
54
55const SCHEMA_SQL: &str = "
56CREATE TABLE IF NOT EXISTS records (
57 memory_id TEXT PRIMARY KEY,
58 scope_kind TEXT NOT NULL,
59 scope_key TEXT NOT NULL,
60 kind TEXT NOT NULL,
61 title TEXT NOT NULL,
62 description TEXT NOT NULL DEFAULT '',
63 body TEXT NOT NULL,
64 tags TEXT NOT NULL DEFAULT '[]',
65 provenance TEXT NOT NULL,
66 trust TEXT NOT NULL,
67 status_kind TEXT NOT NULL,
68 status_detail TEXT,
69 supersedes TEXT,
70 derived_from TEXT NOT NULL DEFAULT '[]',
71 working_set_rank INTEGER,
72 rank_set_at_ms INTEGER,
73 content_hash TEXT NOT NULL,
74 created_at_ms INTEGER NOT NULL,
75 updated_at_ms INTEGER NOT NULL,
76 usage_stats TEXT NOT NULL DEFAULT '{}',
77 tombstoned_at_ms INTEGER,
78 -- §10.2 durable taint marker: 1 when the record landed quarantined or
79 -- descends from a record that did. Survives the tombstone that a
80 -- quarantine release applies to the origin (which erases the
81 -- `quarantined` status), so the transitive ceiling holds forever.
82 ever_quarantined INTEGER NOT NULL DEFAULT 0
83);
84CREATE INDEX IF NOT EXISTS records_scope_idx
85 ON records(scope_kind, scope_key, status_kind);
86CREATE INDEX IF NOT EXISTS records_scope_hash_idx
87 ON records(scope_kind, scope_key, content_hash);
88
89CREATE TABLE IF NOT EXISTS proposals (
90 proposal_id TEXT PRIMARY KEY,
91 scope_kind TEXT NOT NULL,
92 scope_key TEXT NOT NULL,
93 record TEXT NOT NULL,
94 author TEXT NOT NULL,
95 status TEXT NOT NULL DEFAULT 'pending',
96 created_at_ms INTEGER NOT NULL,
97 -- §10.1: quarantine decision captured AT PROPOSE TIME (the taint
98 -- tracker is in-memory and session-sticky; re-deriving at dream time
99 -- both under- and over-quarantines). NULL = clean at propose time.
100 taint TEXT
101);
102
103CREATE TABLE IF NOT EXISTS audit (
104 audit_id INTEGER PRIMARY KEY AUTOINCREMENT,
105 stage_token TEXT NOT NULL,
106 op_index INTEGER NOT NULL,
107 op_kind TEXT NOT NULL,
108 memory_id TEXT,
109 detail TEXT NOT NULL,
110 applied_at_ms INTEGER NOT NULL
111);
112
113CREATE TABLE IF NOT EXISTS stage (
114 token TEXT PRIMARY KEY,
115 batch TEXT NOT NULL,
116 created_at_ms INTEGER NOT NULL
117);
118
119-- Injection ledger (§9.2): plain telemetry appends, deliberately outside
120-- the staged-batch path — rows here are observations about delivery, not
121-- record mutations. session_key is NULL for build-time assembly, where the
122-- session does not exist yet.
123CREATE TABLE IF NOT EXISTS injections (
124 injection_id INTEGER PRIMARY KEY AUTOINCREMENT,
125 record_id TEXT NOT NULL,
126 identity TEXT NOT NULL,
127 session_key TEXT,
128 surface TEXT NOT NULL,
129 at_ms INTEGER NOT NULL
130);
131CREATE INDEX IF NOT EXISTS injections_record_idx
132 ON injections(record_id, at_ms);
133
134-- Exit-interview queue (§8.5): identities recorded by the retire/delete
135-- hooks; the next dream harvests each pending row and marks it done.
136CREATE TABLE IF NOT EXISTS pending_harvests (
137 identity TEXT NOT NULL,
138 session_key TEXT,
139 cause TEXT NOT NULL,
140 retired_at_ms INTEGER NOT NULL,
141 status TEXT NOT NULL DEFAULT 'pending',
142 PRIMARY KEY (identity, retired_at_ms)
143);
144
145-- Quarantine-promotions awaiting operator approval through the gating
146-- flow (§10.2): gating pending_id → staged batch token. Only a gating
147-- approval commits the token; deny/timeout discards it.
148CREATE TABLE IF NOT EXISTS pending_promotions (
149 pending_id TEXT PRIMARY KEY,
150 stage_token TEXT NOT NULL,
151 record_id TEXT NOT NULL,
152 scope_kind TEXT NOT NULL,
153 scope_key TEXT NOT NULL,
154 rationale TEXT,
155 status TEXT NOT NULL DEFAULT 'pending',
156 created_at_ms INTEGER NOT NULL,
157 resolved_at_ms INTEGER
158);
159
160CREATE TABLE IF NOT EXISTS dream_runs (
161 run_id TEXT PRIMARY KEY,
162 partition_label TEXT NOT NULL DEFAULT 'realm',
163 started_at_ms INTEGER NOT NULL,
164 completed_at_ms INTEGER NOT NULL,
165 ops_committed INTEGER NOT NULL,
166 detail TEXT NOT NULL
167);
168CREATE INDEX IF NOT EXISTS dream_runs_completed
169 ON dream_runs(completed_at_ms DESC);
170
171CREATE TABLE IF NOT EXISTS dream_audit_verdicts (
172 run_id TEXT NOT NULL,
173 record_id TEXT NOT NULL,
174 verdict TEXT NOT NULL,
175 rationale TEXT NOT NULL,
176 created_at_ms INTEGER NOT NULL,
177 resolved_at_ms INTEGER,
178 resolution TEXT,
179 PRIMARY KEY (run_id, record_id)
180);
181CREATE INDEX IF NOT EXISTS dream_audit_verdicts_open
182 ON dream_audit_verdicts(record_id, resolved_at_ms);
183";
184
185const RECORD_COLUMNS: &str = "memory_id, scope_kind, scope_key, kind, title, description, body, \
186 tags, provenance, trust, status_kind, status_detail, supersedes, derived_from, \
187 working_set_rank, rank_set_at_ms, content_hash, created_at_ms, updated_at_ms, \
188 usage_stats, tombstoned_at_ms";
189
190#[derive(Clone)]
193pub struct SqliteAgentMemoryStore {
194 root: PathBuf,
195 scope_floor_records: usize,
196 scope_floor_bytes: usize,
197 connections: Arc<Mutex<HashMap<String, Arc<Mutex<Connection>>>>>,
198 llm_write_gate: Arc<Mutex<Option<Arc<dyn LlmWriteGate>>>>,
204 evidence_resolver: Arc<Mutex<Option<Arc<dyn EvidenceRefResolver>>>>,
209 event_sink: Arc<Mutex<Option<Arc<dyn crate::memory::events::MemoryEventSink>>>>,
212}
213
214pub trait EvidenceRefResolver: Send + Sync {
220 fn resolves(&self, evidence: &crate::memory::records::EvidenceRef) -> Result<(), String>;
221}
222
223impl SqliteAgentMemoryStore {
224 pub fn open(root: impl Into<PathBuf>) -> Result<Self, AgentMemoryError> {
225 let root = root.into();
226 if root.as_os_str().is_empty() {
227 return Err(AgentMemoryError::InvalidConfig(
228 "agent memory root path must not be empty".to_string(),
229 ));
230 }
231 fs::create_dir_all(&root).map_err(|err| AgentMemoryError::Io(err.to_string()))?;
232 Ok(Self {
233 root,
234 scope_floor_records: DEFAULT_SCOPE_FLOOR_RECORDS,
235 scope_floor_bytes: DEFAULT_SCOPE_FLOOR_BYTES,
236 connections: Arc::new(Mutex::new(HashMap::new())),
237 llm_write_gate: Arc::new(Mutex::new(None)),
238 evidence_resolver: Arc::new(Mutex::new(None)),
239 event_sink: Arc::new(Mutex::new(None)),
240 })
241 }
242
243 pub fn set_llm_write_gate(&self, gate: Arc<dyn LlmWriteGate>) {
246 *self
247 .llm_write_gate
248 .lock()
249 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(gate);
250 }
251
252 pub fn set_llm_write_gate_if_absent(&self, gate: Arc<dyn LlmWriteGate>) -> bool {
257 let mut guard = self
258 .llm_write_gate
259 .lock()
260 .unwrap_or_else(std::sync::PoisonError::into_inner);
261 if guard.is_some() {
262 return false;
263 }
264 *guard = Some(gate);
265 true
266 }
267
268 pub fn set_evidence_resolver(&self, resolver: Arc<dyn EvidenceRefResolver>) {
272 *self
273 .evidence_resolver
274 .lock()
275 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(resolver);
276 }
277
278 pub fn set_event_sink(&self, sink: Arc<dyn crate::memory::events::MemoryEventSink>) {
280 *self
281 .event_sink
282 .lock()
283 .unwrap_or_else(std::sync::PoisonError::into_inner) = Some(sink);
284 }
285
286 pub fn set_event_sink_if_absent(
290 &self,
291 sink: Arc<dyn crate::memory::events::MemoryEventSink>,
292 ) -> bool {
293 let mut guard = self
294 .event_sink
295 .lock()
296 .unwrap_or_else(std::sync::PoisonError::into_inner);
297 if guard.is_some() {
298 return false;
299 }
300 *guard = Some(sink);
301 true
302 }
303
304 fn gate(&self) -> Option<Arc<dyn LlmWriteGate>> {
305 self.llm_write_gate
306 .lock()
307 .unwrap_or_else(std::sync::PoisonError::into_inner)
308 .clone()
309 }
310
311 fn resolver(&self) -> Option<Arc<dyn EvidenceRefResolver>> {
312 self.evidence_resolver
313 .lock()
314 .unwrap_or_else(std::sync::PoisonError::into_inner)
315 .clone()
316 }
317
318 fn events(&self) -> Option<Arc<dyn crate::memory::events::MemoryEventSink>> {
319 self.event_sink
320 .lock()
321 .unwrap_or_else(std::sync::PoisonError::into_inner)
322 .clone()
323 }
324
325 #[cfg(test)]
326 fn with_scope_floors(mut self, records: usize, bytes: usize) -> Self {
327 self.scope_floor_records = records;
328 self.scope_floor_bytes = bytes;
329 self
330 }
331
332 pub fn path_for_realm(&self, realm: &str) -> PathBuf {
335 self.root
336 .join(format!("{}.sqlite3", encode_path_segment(realm)))
337 }
338
339 fn realm_connection(&self, realm: &str) -> Result<Arc<Mutex<Connection>>, AgentMemoryError> {
340 let mut connections = self
341 .connections
342 .lock()
343 .unwrap_or_else(std::sync::PoisonError::into_inner);
344 if let Some(existing) = connections.get(realm) {
345 return Ok(existing.clone());
346 }
347 let mut conn = Connection::open(self.path_for_realm(realm)).map_err(sql_err)?;
348 conn.busy_timeout(std::time::Duration::from_millis(SQLITE_BUSY_TIMEOUT_MS))
349 .map_err(sql_err)?;
350 conn.query_row("PRAGMA journal_mode=WAL", [], |_| Ok(()))
353 .map_err(sql_err)?;
354 conn.execute_batch(SCHEMA_SQL).map_err(sql_err)?;
355 if ensure_column(
358 &conn,
359 "records",
360 "ever_quarantined",
361 "INTEGER NOT NULL DEFAULT 0",
362 )? {
363 conn.execute(
369 "UPDATE records SET ever_quarantined = 1 WHERE status_kind = 'quarantined'",
370 [],
371 )
372 .map_err(sql_err)?;
373 conn.execute(
374 "UPDATE records SET ever_quarantined = 1 WHERE status_kind = 'tombstoned' \
375 AND memory_id IN (SELECT memory_id FROM audit \
376 WHERE detail LIKE '%\"quarantined\":\"%')",
377 [],
378 )
379 .map_err(sql_err)?;
380 }
381 if ensure_column(&conn, "proposals", "taint", "TEXT")? {
382 conn.execute(
389 "UPDATE proposals SET taint = 'pre-migration proposal: propose-time \
390 taint fact unrecoverable' WHERE status IN ('pending', 'held')",
391 [],
392 )
393 .map_err(sql_err)?;
394 }
395 let now = now_ms();
396 conn.execute(
400 "DELETE FROM stage WHERE created_at_ms < ?1 AND token NOT IN \
401 (SELECT stage_token FROM pending_promotions WHERE status = 'pending')",
402 params![(now.saturating_sub(STAGE_GC_MAX_AGE_MS)) as i64],
403 )
404 .map_err(sql_err)?;
405 self.import_markdown_realm(&mut conn, realm)?;
406 let shared = Arc::new(Mutex::new(conn));
407 connections.insert(realm.to_string(), shared.clone());
408 Ok(shared)
409 }
410
411 fn import_markdown_realm(
425 &self,
426 conn: &mut Connection,
427 realm: &str,
428 ) -> Result<(), AgentMemoryError> {
429 let realm_dir = self.root.join(encode_path_segment(realm));
430 if !realm_dir.is_dir() {
431 return Ok(());
432 }
433 let entries =
434 fs::read_dir(&realm_dir).map_err(|err| AgentMemoryError::Io(err.to_string()))?;
435 let mut files: Vec<PathBuf> = entries
436 .filter_map(|entry| entry.ok().map(|e| e.path()))
437 .filter(|path| path.extension().is_some_and(|ext| ext == "md"))
438 .collect();
439 files.sort();
440 for path in files {
441 match self.import_markdown_file(conn, realm, &path) {
442 Ok(()) => {}
443 Err(MarkdownImportError::Content(reason)) => {
444 tracing::warn!(
445 file = %path.display(),
446 reason,
447 "agent memory markdown import: file failed and was set aside as \
448 .import-failed (fix and rename back to .md to retry); the realm \
449 store stays open"
450 );
451 record_import_audit(conn, &path, 0, 1, std::slice::from_ref(&reason))?;
452 let mut failed_name = path.as_os_str().to_owned();
453 failed_name.push(".import-failed");
454 fs::rename(&path, PathBuf::from(failed_name))
455 .map_err(|err| AgentMemoryError::Io(err.to_string()))?;
456 }
457 Err(MarkdownImportError::Io(err)) => return Err(err),
458 }
459 }
460 Ok(())
461 }
462
463 fn import_markdown_file(
464 &self,
465 conn: &mut Connection,
466 realm: &str,
467 path: &Path,
468 ) -> Result<(), MarkdownImportError> {
469 let Some(stem) = path.file_stem().and_then(|stem| stem.to_str()) else {
470 return Ok(());
471 };
472 let identity_str = decode_path_segment(stem);
473 let identity = AgentIdentity::parse(&identity_str).map_err(|err| {
474 MarkdownImportError::Content(format!(
475 "'{}' does not decode to an agent identity: {err}",
476 path.display()
477 ))
478 })?;
479 let records = read_markdown_records(path).map_err(|err| match err {
480 AgentMemoryError::Io(_) => MarkdownImportError::Io(err),
481 other => MarkdownImportError::Content(other.to_string()),
482 })?;
483 let scope = MemoryScope::Identity {
484 realm: realm.to_string(),
485 identity: identity.as_str().to_string(),
486 };
487 let mut seen = std::collections::HashSet::new();
490 let mut ops = Vec::new();
491 let mut skip_reasons: Vec<String> = Vec::new();
492 for record in records {
493 if !seen.insert(record.memory_id.clone()) {
494 continue;
495 }
496 let exists: Option<i64> = conn
497 .query_row(
498 "SELECT 1 FROM records WHERE memory_id = ?1",
499 params![record.memory_id],
500 |row| row.get(0),
501 )
502 .optional()
503 .map_err(sql_err)
504 .map_err(MarkdownImportError::Io)?;
505 if exists.is_some() {
506 continue;
507 }
508 let mut skip = |record_id: &str, reason: String| {
512 tracing::warn!(
513 file = %path.display(),
514 memory_id = record_id,
515 reason,
516 "agent memory markdown import: record skipped"
517 );
518 skip_reasons.push(format!("{record_id}: {reason}"));
519 };
520 if let Err(reason) = validate_record_fields(&record.title, "", &record.body) {
521 skip(&record.memory_id, reason);
522 continue;
523 }
524 if let Some(class) = crate::memory::secrets::detect_record_secret(
525 &record.title,
526 "",
527 &record.body,
528 &record.tags,
529 ) {
530 skip(
531 &record.memory_id,
532 format!("matches the '{class}' secret pattern class (§10.4)"),
533 );
534 continue;
535 }
536 ops.push(StagedOp::Create {
537 id: Some(record.memory_id),
538 scope: scope.clone(),
539 record: NewMemoryRecord {
540 kind: MemoryKind::Fact,
541 title: record.title,
542 description: String::new(),
543 body: record.body,
544 tags: record.tags,
545 evidence: Vec::new(),
546 verification: None,
547 },
548 trust: TrustTier::AgentObserved,
549 derived_from: Vec::new(),
550 rationale: Some("markdown import".to_string()),
551 created_at_ms: Some(record.created_at_ms),
552 updated_at_ms: Some(record.updated_at_ms),
553 });
554 }
555 let imported = ops.len();
556 if !ops.is_empty() {
557 let batch = StagedMutationBatch {
558 kind: StagedBatchKind::FreshWrite,
559 realm: realm.to_string(),
560 author: MemoryAuthor::Agent {
561 identity: identity.as_str().to_string(),
562 },
563 ops,
564 };
565 let token = mint_token("import");
566 apply_batch_tx(conn, &batch, None, None, &token, now_ms()).map_err(|err| {
569 MarkdownImportError::Content(format!("batch validation failed: {err}"))
570 })?;
571 }
572 if !skip_reasons.is_empty() {
573 record_import_audit(conn, path, imported, skip_reasons.len(), &skip_reasons)
574 .map_err(MarkdownImportError::Io)?;
575 }
576 let mut imported_name = path.as_os_str().to_owned();
577 imported_name.push(".imported");
578 fs::rename(path, PathBuf::from(imported_name))
579 .map_err(|err| MarkdownImportError::Io(AgentMemoryError::Io(err.to_string())))?;
580 Ok(())
581 }
582
583 fn with_realm_conn<T>(
584 &self,
585 realm: &str,
586 f: impl FnOnce(&mut Connection) -> Result<T, AgentMemoryError>,
587 ) -> Result<T, AgentMemoryError> {
588 let conn = self.realm_connection(realm)?;
589 let mut guard = conn
590 .lock()
591 .unwrap_or_else(std::sync::PoisonError::into_inner);
592 f(&mut guard)
593 }
594
595 fn recall_blocking(
596 &self,
597 request: AgentMemoryRecallRequest,
598 ) -> Result<Vec<AgentMemoryRecord>, AgentMemoryError> {
599 let scope = MemoryScope::Identity {
600 realm: request.realm.clone(),
601 identity: request.identity.as_str().to_string(),
602 };
603 let records =
604 self.with_realm_conn(&request.realm, |conn| active_scope_records(conn, &scope))?;
605 let projected = records.into_iter().map(project_record).collect();
606 Ok(select_recall_records(projected, &request))
607 }
608
609 fn remember_blocking(
610 &self,
611 realm: &str,
612 identity: &AgentIdentity,
613 memory: NewAgentMemory,
614 ) -> Result<AgentMemoryRecord, AgentMemoryError> {
615 let title = compact_whitespace(&memory.title);
616 let body = memory.body.trim().to_string();
617 validate_record_fields(&title, "", &body).map_err(AgentMemoryError::InvalidRecord)?;
618 let tags = normalize_tags(memory.tags)?;
619 let scope = MemoryScope::Identity {
620 realm: realm.to_string(),
621 identity: identity.as_str().to_string(),
622 };
623 let hash = content_hash(&title, &body);
624 let floor_records = self.scope_floor_records;
625 let floor_bytes = self.scope_floor_bytes;
626 let gate = self.gate();
627 let events = self.events();
628 self.with_realm_conn(realm, |conn| {
629 let existing: Option<MemoryRecordRow> = conn
632 .query_row(
633 &format!(
634 "SELECT {RECORD_COLUMNS} FROM records \
635 WHERE scope_kind = ?1 AND scope_key = ?2 AND content_hash = ?3 \
636 AND status_kind = 'active' \
637 ORDER BY created_at_ms ASC LIMIT 1"
638 ),
639 params![scope.kind_str(), scope.key(), hash],
640 row_to_record_row,
641 )
642 .optional()
643 .map_err(sql_err)?;
644 if let Some(row) = existing {
645 return Ok(project_record(row.into_record(scope.realm())?));
646 }
647 let batch = StagedMutationBatch {
648 kind: StagedBatchKind::FreshWrite,
649 realm: realm.to_string(),
650 author: MemoryAuthor::Application,
653 ops: vec![StagedOp::Create {
654 id: None,
655 scope: scope.clone(),
656 record: NewMemoryRecord {
657 kind: MemoryKind::Fact,
658 title,
659 description: String::new(),
660 body,
661 tags: tags.clone(),
662 evidence: Vec::new(),
663 verification: None,
664 },
665 trust: TrustTier::AgentObserved,
666 derived_from: Vec::new(),
667 rationale: None,
668 created_at_ms: None,
669 updated_at_ms: None,
670 }],
671 };
672 let receipt = apply_batch_tx(
673 conn,
674 &batch,
675 gate.as_deref(),
676 events.as_deref(),
677 &mint_token("direct"),
678 now_ms(),
679 )?;
680 warn_if_scope_floors_exceeded(conn, &scope, floor_records, floor_bytes)?;
681 let memory_id = receipt.memory_ids.first().cloned().ok_or_else(|| {
682 AgentMemoryError::Io("remember commit returned no record id".to_string())
683 })?;
684 let record = load_record(conn, scope.realm(), &memory_id)?.ok_or_else(|| {
685 AgentMemoryError::Io("remembered record vanished mid-commit".to_string())
686 })?;
687 Ok(project_record(record))
688 })
689 }
690
691 fn forget_blocking(
692 &self,
693 realm: &str,
694 identity: &AgentIdentity,
695 memory_id: &str,
696 ) -> Result<AgentMemoryForgetResult, AgentMemoryError> {
697 let memory_id = memory_id.trim().to_string();
698 if memory_id.is_empty() {
699 return Err(AgentMemoryError::InvalidRecord(
700 "memory_id must not be empty".to_string(),
701 ));
702 }
703 let scope = MemoryScope::Identity {
704 realm: realm.to_string(),
705 identity: identity.as_str().to_string(),
706 };
707 self.forget_in_scope_blocking(&scope, &memory_id, MemoryAuthor::Application)
708 }
709
710 fn forget_in_scope_blocking(
713 &self,
714 scope: &MemoryScope,
715 memory_id: &str,
716 author: MemoryAuthor,
717 ) -> Result<AgentMemoryForgetResult, AgentMemoryError> {
718 let memory_id = memory_id.to_string();
719 let gate = self.gate();
720 let events = self.events();
721 self.with_realm_conn(scope.realm(), |conn| {
722 let record = load_record(conn, scope.realm(), &memory_id)?;
723 let deletable = record.is_some_and(|record| {
724 record.scope == *scope && record.status != RecordStatus::Tombstoned
725 });
726 if !deletable {
727 return Ok(AgentMemoryForgetResult {
728 memory_id,
729 deleted: false,
730 });
731 }
732 let batch = StagedMutationBatch {
733 kind: StagedBatchKind::FreshWrite,
734 realm: scope.realm().to_string(),
735 author,
736 ops: vec![StagedOp::Tombstone {
737 id: memory_id.clone(),
738 rationale: None,
739 }],
740 };
741 apply_batch_tx(
742 conn,
743 &batch,
744 gate.as_deref(),
745 events.as_deref(),
746 &mint_token("direct"),
747 now_ms(),
748 )?;
749 Ok(AgentMemoryForgetResult {
750 memory_id,
751 deleted: true,
752 })
753 })
754 }
755
756 fn supersede_blocking(
757 &self,
758 scope: &MemoryScope,
759 prior: &str,
760 record: NewMemoryRecord,
761 ) -> Result<MemoryId, AgentMemoryError> {
762 self.supersede_with_author_blocking(scope, prior, record, MemoryAuthor::Application)
763 .map(|receipt| receipt.memory_id)
764 }
765
766 fn supersede_with_author_blocking(
767 &self,
768 scope: &MemoryScope,
769 prior: &str,
770 record: NewMemoryRecord,
771 author: MemoryAuthor,
772 ) -> Result<AuthoredWriteReceipt, AgentMemoryError> {
773 let title = compact_whitespace(&record.title);
774 let body = record.body.trim().to_string();
775 validate_record_fields(&title, &record.description, &body)
776 .map_err(AgentMemoryError::InvalidRecord)?;
777 let tags = normalize_tags(record.tags)?;
778 let realm = scope.realm().to_string();
779 let expected_scope = scope.clone();
780 let gate = self.gate();
781 let events = self.events();
782 self.with_realm_conn(&realm, |conn| {
783 let existing = load_record(conn, &realm, prior)?.ok_or_else(|| {
784 AgentMemoryError::InvalidRecord(format!("record '{prior}' does not exist"))
785 })?;
786 if existing.scope != expected_scope {
787 return Err(AgentMemoryError::InvalidRecord(format!(
788 "record '{prior}' does not belong to the requested scope"
789 )));
790 }
791 let batch = StagedMutationBatch {
792 kind: StagedBatchKind::FreshWrite,
793 realm: realm.clone(),
794 author,
795 ops: vec![StagedOp::Supersede {
796 id: None,
797 prior: prior.to_string(),
798 record: NewMemoryRecord {
799 title,
800 body,
801 tags,
802 ..record
803 },
804 trust: TrustTier::AgentObserved,
805 derived_from: Vec::new(),
806 rationale: None,
807 }],
808 };
809 let receipt = apply_batch_tx(
810 conn,
811 &batch,
812 gate.as_deref(),
813 events.as_deref(),
814 &mint_token("direct"),
815 now_ms(),
816 )?;
817 let memory_id = receipt.memory_ids.first().cloned().ok_or_else(|| {
818 AgentMemoryError::Io("supersede commit returned no record id".to_string())
819 })?;
820 let record = load_record(conn, &realm, &memory_id)?.ok_or_else(|| {
821 AgentMemoryError::Io("superseding record vanished mid-commit".to_string())
822 })?;
823 Ok(AuthoredWriteReceipt {
824 memory_id,
825 status: record.status,
826 })
827 })
828 }
829
830 fn remember_authored_blocking(
832 &self,
833 scope: &MemoryScope,
834 record: NewMemoryRecord,
835 author: MemoryAuthor,
836 ) -> Result<AuthoredWriteReceipt, AgentMemoryError> {
837 let title = compact_whitespace(&record.title);
838 let body = record.body.trim().to_string();
839 validate_record_fields(&title, &record.description, &body)
840 .map_err(AgentMemoryError::InvalidRecord)?;
841 let tags = normalize_tags(record.tags)?;
842 let hash = content_hash(&title, &body);
843 let realm = scope.realm().to_string();
844 let scope = scope.clone();
845 let floor_records = self.scope_floor_records;
846 let floor_bytes = self.scope_floor_bytes;
847 let gate = self.gate();
848 let events = self.events();
849 self.with_realm_conn(&realm, |conn| {
850 let existing: Option<MemoryRecordRow> = conn
853 .query_row(
854 &format!(
855 "SELECT {RECORD_COLUMNS} FROM records \
856 WHERE scope_kind = ?1 AND scope_key = ?2 AND content_hash = ?3 \
857 AND status_kind = 'active' \
858 ORDER BY created_at_ms ASC LIMIT 1"
859 ),
860 params![scope.kind_str(), scope.key(), hash],
861 row_to_record_row,
862 )
863 .optional()
864 .map_err(sql_err)?;
865 if let Some(row) = existing {
866 let record = row.into_record(scope.realm())?;
867 return Ok(AuthoredWriteReceipt {
868 memory_id: record.id,
869 status: record.status,
870 });
871 }
872 let batch = StagedMutationBatch {
873 kind: StagedBatchKind::FreshWrite,
874 realm: realm.clone(),
875 author,
876 ops: vec![StagedOp::Create {
877 id: None,
878 scope: scope.clone(),
879 record: NewMemoryRecord {
880 title,
881 body,
882 tags,
883 ..record
884 },
885 trust: TrustTier::AgentObserved,
888 derived_from: Vec::new(),
889 rationale: None,
890 created_at_ms: None,
891 updated_at_ms: None,
892 }],
893 };
894 let receipt = apply_batch_tx(
895 conn,
896 &batch,
897 gate.as_deref(),
898 events.as_deref(),
899 &mint_token("direct"),
900 now_ms(),
901 )?;
902 warn_if_scope_floors_exceeded(conn, &scope, floor_records, floor_bytes)?;
903 let memory_id = receipt.memory_ids.first().cloned().ok_or_else(|| {
904 AgentMemoryError::Io("remember commit returned no record id".to_string())
905 })?;
906 let record = load_record(conn, scope.realm(), &memory_id)?.ok_or_else(|| {
907 AgentMemoryError::Io("remembered record vanished mid-commit".to_string())
908 })?;
909 Ok(AuthoredWriteReceipt {
910 memory_id,
911 status: record.status,
912 })
913 })
914 }
915
916 fn manifest_blocking(
917 &self,
918 scopes: &[MemoryScope],
919 tier: ManifestTier,
920 ) -> Result<Vec<RecordMeta>, AgentMemoryError> {
921 let now = now_ms();
922 let mut out = Vec::new();
923 for scope in scopes {
924 let metas =
925 self.with_realm_conn(scope.realm(), |conn| scope_manifest(conn, scope, tier, now))?;
926 out.extend(metas);
927 }
928 Ok(out)
929 }
930
931 fn mark_usage_blocking(
932 &self,
933 ids: &[MemoryId],
934 event: UsageEvent,
935 ) -> Result<(), AgentMemoryError> {
936 let now = now_ms();
937 for realm in self.known_realms()? {
938 self.with_realm_conn(&realm, |conn| {
939 for id in ids {
940 let usage_json: Option<String> = conn
941 .query_row(
942 "SELECT usage_stats FROM records WHERE memory_id = ?1",
943 params![id],
944 |row| row.get(0),
945 )
946 .optional()
947 .map_err(sql_err)?;
948 let Some(usage_json) = usage_json else {
949 continue;
950 };
951 let mut usage: UsageStats =
952 serde_json::from_str(&usage_json).unwrap_or_default();
953 match event {
954 UsageEvent::Injected => {
955 usage.injected_count += 1;
956 usage.last_injected_at_ms = Some(now);
957 }
958 UsageEvent::ExplicitRecall => {
962 usage.explicit_recall_count += 1;
963 usage.last_recalled_at_ms = Some(now);
964 }
965 UsageEvent::JudgedUseful => {
966 usage.judged_useful_count += 1;
967 usage.last_useful_at_ms = Some(now);
968 }
969 }
970 conn.execute(
971 "UPDATE records SET usage_stats = ?1 WHERE memory_id = ?2",
972 params![json_string(&usage)?, id],
973 )
974 .map_err(sql_err)?;
975 }
976 Ok(())
977 })?;
978 }
979 Ok(())
980 }
981
982 fn log_injections_blocking(
983 &self,
984 realm: &str,
985 entries: &[InjectionLogEntry],
986 ) -> Result<(), AgentMemoryError> {
987 if entries.is_empty() {
988 return Ok(());
989 }
990 self.with_realm_conn(realm, |conn| {
991 let mut stmt = conn
992 .prepare(
993 "INSERT INTO injections (record_id, identity, session_key, surface, at_ms) \
994 VALUES (?1, ?2, ?3, ?4, ?5)",
995 )
996 .map_err(sql_err)?;
997 for entry in entries {
998 stmt.execute(params![
999 entry.record_id,
1000 entry.identity,
1001 entry.session_key,
1002 entry.surface.as_str(),
1003 entry.at_ms as i64,
1004 ])
1005 .map_err(sql_err)?;
1006 }
1007 Ok(())
1008 })
1009 }
1010
1011 fn injection_log_blocking(
1012 &self,
1013 realm: &str,
1014 limit: usize,
1015 ) -> Result<Vec<InjectionLogEntry>, AgentMemoryError> {
1016 self.with_realm_conn(realm, |conn| {
1017 let mut stmt = conn
1018 .prepare(
1019 "SELECT record_id, identity, session_key, surface, at_ms FROM injections \
1020 ORDER BY injection_id DESC LIMIT ?1",
1021 )
1022 .map_err(sql_err)?;
1023 let rows = stmt
1024 .query_map(params![limit as i64], |row| {
1025 Ok((
1026 row.get::<_, String>(0)?,
1027 row.get::<_, String>(1)?,
1028 row.get::<_, Option<String>>(2)?,
1029 row.get::<_, String>(3)?,
1030 row.get::<_, i64>(4)?,
1031 ))
1032 })
1033 .map_err(sql_err)?;
1034 let mut entries = Vec::new();
1035 for row in rows {
1036 let (record_id, identity, session_key, surface, at_ms) = row.map_err(sql_err)?;
1037 let surface = InjectionSurface::parse(&surface).ok_or_else(|| {
1038 AgentMemoryError::Parse(format!("unknown injection surface '{surface}'"))
1039 })?;
1040 entries.push(InjectionLogEntry {
1041 record_id,
1042 identity,
1043 session_key,
1044 surface,
1045 at_ms: at_ms as u64,
1046 });
1047 }
1048 Ok(entries)
1049 })
1050 }
1051
1052 pub async fn injection_log(
1055 &self,
1056 realm: &str,
1057 limit: usize,
1058 ) -> Result<Vec<InjectionLogEntry>, AgentMemoryError> {
1059 let store = self.clone();
1060 let realm = realm.to_string();
1061 run_blocking(move || store.injection_log_blocking(&realm, limit)).await
1062 }
1063
1064 fn propose_blocking(
1065 &self,
1066 scope: &MemoryScope,
1067 record: NewMemoryRecord,
1068 author: MemoryAuthor,
1069 ) -> Result<ProposalId, AgentMemoryError> {
1070 validate_record_fields(&record.title, &record.description, &record.body)
1071 .map_err(AgentMemoryError::InvalidRecord)?;
1072 if let Some(class) = crate::memory::secrets::detect_record_secret(
1076 &record.title,
1077 &record.description,
1078 &record.body,
1079 &record.tags,
1080 ) {
1081 return Err(AgentMemoryError::InvalidRecord(
1082 crate::memory::staged::StagedBatchError::SecretDetected { op_index: 0, class }
1083 .to_string(),
1084 ));
1085 }
1086 let taint = self.gate().and_then(|gate| {
1093 gate.quarantine_reason(&author, StagedBatchKind::FreshWrite, &record.evidence)
1094 });
1095 if let Some(reason) = taint.as_deref() {
1096 tracing::warn!(
1097 realm = scope.realm(),
1098 author = ?author,
1099 reason,
1100 "agent memory: proposal from tainted context recorded as tainted; a plain \
1101 steward accept will downgrade to an operator gate"
1102 );
1103 }
1104 let proposal_id = mint_token("prop");
1105 self.with_realm_conn(scope.realm(), |conn| {
1106 conn.execute(
1107 "INSERT INTO proposals (proposal_id, scope_kind, scope_key, record, author, \
1108 status, created_at_ms, taint) VALUES (?1, ?2, ?3, ?4, ?5, 'pending', ?6, ?7)",
1109 params![
1110 proposal_id,
1111 scope.kind_str(),
1112 scope.key(),
1113 json_string(&record)?,
1114 json_string(&author)?,
1115 now_ms() as i64,
1116 taint,
1117 ],
1118 )
1119 .map_err(sql_err)?;
1120 Ok(())
1121 })?;
1122 Ok(proposal_id)
1123 }
1124
1125 fn stage_blocking(&self, batch: StagedMutationBatch) -> Result<StageToken, AgentMemoryError> {
1126 let realm = batch.realm.clone();
1127 let resolver = self.resolver();
1128 self.with_realm_conn(&realm, |conn| {
1129 {
1130 let view = ConnBatchView {
1131 conn,
1132 realm: &batch.realm,
1133 };
1134 validate_batch(
1135 &batch,
1136 &view,
1137 DEFAULT_TOMBSTONE_RECREATE_WINDOW_MS,
1138 now_ms(),
1139 )
1140 .map_err(|err| AgentMemoryError::InvalidRecord(err.to_string()))?;
1141 }
1142 check_verified_retier_evidence(conn, &batch, resolver.as_deref())?;
1143 let token = mint_token("stage");
1144 conn.execute(
1145 "INSERT INTO stage (token, batch, created_at_ms) VALUES (?1, ?2, ?3)",
1146 params![token, json_string(&batch)?, now_ms() as i64],
1147 )
1148 .map_err(sql_err)?;
1149 Ok(StageToken {
1150 realm: realm.clone(),
1151 token,
1152 })
1153 })
1154 }
1155
1156 fn commit_blocking(&self, token: StageToken) -> Result<CommitReceipt, AgentMemoryError> {
1157 let gate = self.gate();
1158 let resolver = self.resolver();
1159 let events = self.events();
1160 self.with_realm_conn(&token.realm, |conn| {
1161 let batch_json: Option<String> = conn
1162 .query_row(
1163 "SELECT batch FROM stage WHERE token = ?1",
1164 params![token.token],
1165 |row| row.get(0),
1166 )
1167 .optional()
1168 .map_err(sql_err)?;
1169 let Some(batch_json) = batch_json else {
1170 return Err(AgentMemoryError::InvalidRecord(format!(
1171 "unknown or expired stage token '{}'",
1172 token.token
1173 )));
1174 };
1175 let batch: StagedMutationBatch = serde_json::from_str(&batch_json)
1176 .map_err(|err| AgentMemoryError::Parse(err.to_string()))?;
1177 check_verified_retier_evidence(conn, &batch, resolver.as_deref())?;
1178 apply_batch_tx(
1179 conn,
1180 &batch,
1181 gate.as_deref(),
1182 events.as_deref(),
1183 &token.token,
1184 now_ms(),
1185 )
1186 })
1187 }
1188
1189 fn known_realms(&self) -> Result<Vec<String>, AgentMemoryError> {
1190 let entries =
1191 fs::read_dir(&self.root).map_err(|err| AgentMemoryError::Io(err.to_string()))?;
1192 let mut realms = Vec::new();
1193 for entry in entries.filter_map(Result::ok) {
1194 let path = entry.path();
1195 if path.extension().is_some_and(|ext| ext == "sqlite3")
1196 && let Some(stem) = path.file_stem().and_then(|stem| stem.to_str())
1197 {
1198 realms.push(decode_path_segment(stem));
1199 }
1200 }
1201 realms.sort();
1202 Ok(realms)
1203 }
1204}
1205
1206#[async_trait]
1209impl AgentMemoryProvider for SqliteAgentMemoryStore {
1210 async fn recall(
1211 &self,
1212 request: AgentMemoryRecallRequest,
1213 ) -> Result<Vec<AgentMemoryRecord>, AgentMemoryError> {
1214 let store = self.clone();
1215 run_blocking(move || store.recall_blocking(request)).await
1216 }
1217
1218 fn supports_remember(&self) -> bool {
1219 true
1220 }
1221
1222 async fn remember(
1223 &self,
1224 realm: &str,
1225 identity: &AgentIdentity,
1226 memory: NewAgentMemory,
1227 ) -> Result<AgentMemoryRecord, AgentMemoryError> {
1228 let store = self.clone();
1229 let realm = realm.to_string();
1230 let identity = identity.clone();
1231 run_blocking(move || store.remember_blocking(&realm, &identity, memory)).await
1232 }
1233
1234 fn supports_forget(&self) -> bool {
1235 true
1236 }
1237
1238 async fn forget(
1239 &self,
1240 realm: &str,
1241 identity: &AgentIdentity,
1242 memory_id: &str,
1243 ) -> Result<AgentMemoryForgetResult, AgentMemoryError> {
1244 let store = self.clone();
1245 let realm = realm.to_string();
1246 let identity = identity.clone();
1247 let memory_id = memory_id.to_string();
1248 run_blocking(move || store.forget_blocking(&realm, &identity, &memory_id)).await
1249 }
1250
1251 async fn manifest(
1252 &self,
1253 scopes: &[MemoryScope],
1254 tier: ManifestTier,
1255 ) -> Result<Vec<RecordMeta>, AgentMemoryError> {
1256 let store = self.clone();
1257 let scopes = scopes.to_vec();
1258 run_blocking(move || store.manifest_blocking(&scopes, tier)).await
1259 }
1260
1261 fn supports_manifest(&self) -> bool {
1262 true
1263 }
1264
1265 async fn supersede(
1266 &self,
1267 scope: &MemoryScope,
1268 prior: &str,
1269 record: NewMemoryRecord,
1270 ) -> Result<MemoryId, AgentMemoryError> {
1271 let store = self.clone();
1272 let scope = scope.clone();
1273 let prior = prior.to_string();
1274 run_blocking(move || store.supersede_blocking(&scope, &prior, record)).await
1275 }
1276
1277 fn supports_supersede(&self) -> bool {
1278 true
1279 }
1280
1281 async fn mark_usage(
1282 &self,
1283 ids: &[MemoryId],
1284 event: UsageEvent,
1285 ) -> Result<(), AgentMemoryError> {
1286 let store = self.clone();
1287 let ids = ids.to_vec();
1288 run_blocking(move || store.mark_usage_blocking(&ids, event)).await
1289 }
1290
1291 async fn log_injections(
1292 &self,
1293 realm: &str,
1294 entries: &[InjectionLogEntry],
1295 ) -> Result<(), AgentMemoryError> {
1296 let store = self.clone();
1297 let realm = realm.to_string();
1298 let entries = entries.to_vec();
1299 run_blocking(move || store.log_injections_blocking(&realm, &entries)).await
1300 }
1301
1302 async fn propose(
1303 &self,
1304 scope: &MemoryScope,
1305 record: NewMemoryRecord,
1306 author: MemoryAuthor,
1307 ) -> Result<ProposalId, AgentMemoryError> {
1308 let store = self.clone();
1309 let scope = scope.clone();
1310 run_blocking(move || store.propose_blocking(&scope, record, author)).await
1311 }
1312
1313 fn supports_propose(&self) -> bool {
1314 true
1315 }
1316
1317 async fn remember_authored(
1318 &self,
1319 scope: &MemoryScope,
1320 record: NewMemoryRecord,
1321 author: MemoryAuthor,
1322 ) -> Result<AuthoredWriteReceipt, AgentMemoryError> {
1323 let store = self.clone();
1324 let scope = scope.clone();
1325 run_blocking(move || store.remember_authored_blocking(&scope, record, author)).await
1326 }
1327
1328 async fn supersede_authored(
1329 &self,
1330 scope: &MemoryScope,
1331 prior: &str,
1332 record: NewMemoryRecord,
1333 author: MemoryAuthor,
1334 ) -> Result<AuthoredWriteReceipt, AgentMemoryError> {
1335 let store = self.clone();
1336 let scope = scope.clone();
1337 let prior = prior.to_string();
1338 run_blocking(move || store.supersede_with_author_blocking(&scope, &prior, record, author))
1339 .await
1340 }
1341
1342 async fn forget_authored(
1343 &self,
1344 scope: &MemoryScope,
1345 memory_id: &str,
1346 author: MemoryAuthor,
1347 ) -> Result<AgentMemoryForgetResult, AgentMemoryError> {
1348 let memory_id = memory_id.trim().to_string();
1349 if memory_id.is_empty() {
1350 return Err(AgentMemoryError::InvalidRecord(
1351 "memory_id must not be empty".to_string(),
1352 ));
1353 }
1354 let store = self.clone();
1355 let scope = scope.clone();
1356 run_blocking(move || store.forget_in_scope_blocking(&scope, &memory_id, author)).await
1357 }
1358
1359 fn supports_authored_writes(&self) -> bool {
1360 true
1361 }
1362
1363 fn as_sqlite_store(&self) -> Option<&SqliteAgentMemoryStore> {
1364 Some(self)
1365 }
1366}
1367
1368#[async_trait]
1369impl StagedMemoryStore for SqliteAgentMemoryStore {
1370 async fn stage(&self, batch: StagedMutationBatch) -> Result<StageToken, AgentMemoryError> {
1371 let store = self.clone();
1372 run_blocking(move || store.stage_blocking(batch)).await
1373 }
1374
1375 async fn commit(&self, token: StageToken) -> Result<CommitReceipt, AgentMemoryError> {
1376 let store = self.clone();
1377 run_blocking(move || store.commit_blocking(token)).await
1378 }
1379}
1380
1381#[derive(Debug, Clone, PartialEq, Eq)]
1386pub struct ScopeOverview {
1387 pub scope: MemoryScope,
1388 pub active: u64,
1389 pub quarantined: u64,
1390 pub superseded: u64,
1391 pub tombstoned: u64,
1392 pub body_bytes: u64,
1393}
1394
1395#[derive(Debug, Clone, PartialEq)]
1398pub struct PendingProposal {
1399 pub proposal_id: ProposalId,
1400 pub scope: MemoryScope,
1401 pub record: NewMemoryRecord,
1402 pub author: MemoryAuthor,
1403 pub status: String,
1404 pub created_at_ms: u64,
1405 pub taint: Option<String>,
1409}
1410
1411#[derive(Debug, Clone, PartialEq, Eq)]
1413pub struct PendingHarvest {
1414 pub identity: String,
1415 pub session_key: Option<String>,
1416 pub cause: String,
1417 pub retired_at_ms: u64,
1418}
1419
1420#[derive(Debug, Clone, PartialEq, Eq)]
1423pub struct PendingPromotion {
1424 pub pending_id: String,
1425 pub stage_token: String,
1426 pub record_id: MemoryId,
1427 pub scope_kind: String,
1428 pub scope_key: String,
1429 pub rationale: Option<String>,
1430 pub status: String,
1431 pub created_at_ms: u64,
1432}
1433
1434#[derive(Debug, Clone, PartialEq, Eq)]
1438pub struct PersistedDreamRun {
1439 pub run_id: String,
1440 pub partition_label: String,
1441 pub started_at_ms: u64,
1442 pub completed_at_ms: u64,
1443 pub ops_committed: u64,
1444 pub detail: String,
1446}
1447
1448#[derive(Debug, Clone, PartialEq, Eq)]
1452pub struct DreamAuditVerdict {
1453 pub run_id: String,
1454 pub record_id: String,
1455 pub verdict: String,
1456 pub rationale: String,
1457 pub created_at_ms: u64,
1458 pub resolved_at_ms: Option<u64>,
1459 pub resolution: Option<String>,
1460}
1461
1462impl SqliteAgentMemoryStore {
1463 pub fn scope_floors(&self) -> (usize, usize) {
1466 (self.scope_floor_records, self.scope_floor_bytes)
1467 }
1468
1469 pub async fn scope_overview(
1472 &self,
1473 realm: &str,
1474 ) -> Result<Vec<ScopeOverview>, AgentMemoryError> {
1475 let store = self.clone();
1476 let realm = realm.to_string();
1477 run_blocking(move || {
1478 store.with_realm_conn(&realm, |conn| {
1479 let mut stmt = conn
1480 .prepare(
1481 "SELECT scope_kind, scope_key, status_kind, COUNT(*), \
1482 COALESCE(SUM(LENGTH(body)), 0) FROM records \
1483 GROUP BY scope_kind, scope_key, status_kind",
1484 )
1485 .map_err(sql_err)?;
1486 let rows = stmt
1487 .query_map([], |row| {
1488 Ok((
1489 row.get::<_, String>(0)?,
1490 row.get::<_, String>(1)?,
1491 row.get::<_, String>(2)?,
1492 row.get::<_, i64>(3)?,
1493 row.get::<_, i64>(4)?,
1494 ))
1495 })
1496 .map_err(sql_err)?;
1497 let mut by_scope: HashMap<(String, String), ScopeOverview> = HashMap::new();
1498 for row in rows {
1499 let (scope_kind, scope_key, status_kind, count, bytes) =
1500 row.map_err(sql_err)?;
1501 let scope = scope_from_parts(&scope_kind, &scope_key, &realm)?;
1502 let entry =
1503 by_scope
1504 .entry((scope_kind, scope_key))
1505 .or_insert_with(|| ScopeOverview {
1506 scope,
1507 active: 0,
1508 quarantined: 0,
1509 superseded: 0,
1510 tombstoned: 0,
1511 body_bytes: 0,
1512 });
1513 match status_kind.as_str() {
1514 "active" => entry.active = count as u64,
1515 "quarantined" => entry.quarantined = count as u64,
1516 "superseded" => entry.superseded = count as u64,
1517 "tombstoned" => entry.tombstoned = count as u64,
1518 _ => {}
1519 }
1520 entry.body_bytes += bytes as u64;
1521 }
1522 let mut overview: Vec<ScopeOverview> = by_scope.into_values().collect();
1523 overview.sort_by(|a, b| a.scope.cmp(&b.scope));
1524 Ok(overview)
1525 })
1526 })
1527 .await
1528 }
1529
1530 pub async fn pending_proposals(
1532 &self,
1533 realm: &str,
1534 limit: usize,
1535 ) -> Result<Vec<PendingProposal>, AgentMemoryError> {
1536 let store = self.clone();
1537 let realm = realm.to_string();
1538 run_blocking(move || {
1539 store.with_realm_conn(&realm, |conn| {
1540 let mut stmt = conn
1541 .prepare(
1542 "SELECT proposal_id, scope_kind, scope_key, record, author, status, \
1543 created_at_ms, taint FROM proposals WHERE status IN ('pending', 'held') \
1544 ORDER BY created_at_ms ASC LIMIT ?1",
1545 )
1546 .map_err(sql_err)?;
1547 let rows = stmt
1548 .query_map(params![limit as i64], |row| {
1549 Ok((
1550 row.get::<_, String>(0)?,
1551 row.get::<_, String>(1)?,
1552 row.get::<_, String>(2)?,
1553 row.get::<_, String>(3)?,
1554 row.get::<_, String>(4)?,
1555 row.get::<_, String>(5)?,
1556 row.get::<_, i64>(6)?,
1557 row.get::<_, Option<String>>(7)?,
1558 ))
1559 })
1560 .map_err(sql_err)?;
1561 let mut proposals = Vec::new();
1562 for row in rows {
1563 let (
1564 proposal_id,
1565 scope_kind,
1566 scope_key,
1567 record,
1568 author,
1569 status,
1570 created,
1571 taint,
1572 ) = row.map_err(sql_err)?;
1573 proposals.push(PendingProposal {
1574 proposal_id,
1575 scope: scope_from_parts(&scope_kind, &scope_key, &realm)?,
1576 record: serde_json::from_str(&record)
1577 .map_err(|err| AgentMemoryError::Parse(err.to_string()))?,
1578 author: serde_json::from_str(&author)
1579 .map_err(|err| AgentMemoryError::Parse(err.to_string()))?,
1580 status,
1581 created_at_ms: created as u64,
1582 taint,
1583 });
1584 }
1585 Ok(proposals)
1586 })
1587 })
1588 .await
1589 }
1590
1591 pub async fn set_proposal_status(
1594 &self,
1595 realm: &str,
1596 proposal_id: &str,
1597 status: &str,
1598 ) -> Result<(), AgentMemoryError> {
1599 if !matches!(status, "accepted" | "rejected" | "held" | "pending") {
1600 return Err(AgentMemoryError::InvalidRecord(format!(
1601 "unknown proposal status '{status}'"
1602 )));
1603 }
1604 let store = self.clone();
1605 let realm = realm.to_string();
1606 let proposal_id = proposal_id.to_string();
1607 let status = status.to_string();
1608 run_blocking(move || {
1609 store.with_realm_conn(&realm, |conn| {
1610 let updated = conn
1611 .execute(
1612 "UPDATE proposals SET status = ?1 WHERE proposal_id = ?2",
1613 params![status, proposal_id],
1614 )
1615 .map_err(sql_err)?;
1616 if updated == 0 {
1617 return Err(AgentMemoryError::InvalidRecord(format!(
1618 "unknown proposal '{proposal_id}'"
1619 )));
1620 }
1621 Ok(())
1622 })
1623 })
1624 .await
1625 }
1626
1627 pub async fn quarantined_records(
1631 &self,
1632 realm: &str,
1633 limit: usize,
1634 ) -> Result<Vec<super::records::MemoryRecord>, AgentMemoryError> {
1635 let store = self.clone();
1636 let realm = realm.to_string();
1637 run_blocking(move || {
1638 store.with_realm_conn(&realm, |conn| {
1639 let mut stmt = conn
1640 .prepare(&format!(
1641 "SELECT {RECORD_COLUMNS} FROM records \
1642 WHERE status_kind = 'quarantined' \
1643 ORDER BY created_at_ms DESC LIMIT ?1"
1644 ))
1645 .map_err(sql_err)?;
1646 let rows = stmt
1647 .query_map(params![limit as i64], row_to_record_row)
1648 .map_err(sql_err)?;
1649 let mut records = Vec::new();
1650 for row in rows {
1651 records.push(row.map_err(sql_err)?.into_record(&realm)?);
1652 }
1653 Ok(records)
1654 })
1655 })
1656 .await
1657 }
1658
1659 pub async fn records_by_ids(
1662 &self,
1663 realm: &str,
1664 ids: &[String],
1665 ) -> Result<Vec<super::records::MemoryRecord>, AgentMemoryError> {
1666 let store = self.clone();
1667 let realm = realm.to_string();
1668 let ids = ids.to_vec();
1669 run_blocking(move || {
1670 store.with_realm_conn(&realm, |conn| {
1671 let mut records = Vec::new();
1672 for id in &ids {
1673 if let Some(record) = load_record(conn, &realm, id)? {
1674 records.push(record);
1675 }
1676 }
1677 Ok(records)
1678 })
1679 })
1680 .await
1681 }
1682
1683 pub async fn recent_records(
1687 &self,
1688 realm: &str,
1689 limit: usize,
1690 ) -> Result<Vec<super::records::MemoryRecord>, AgentMemoryError> {
1691 let store = self.clone();
1692 let realm = realm.to_string();
1693 run_blocking(move || {
1694 store.with_realm_conn(&realm, |conn| {
1695 let mut stmt = conn
1696 .prepare(&format!(
1697 "SELECT {RECORD_COLUMNS} FROM records \
1698 WHERE status_kind IN ('active', 'quarantined') \
1699 ORDER BY updated_at_ms DESC LIMIT ?1"
1700 ))
1701 .map_err(sql_err)?;
1702 let rows = stmt
1703 .query_map(params![limit as i64], row_to_record_row)
1704 .map_err(sql_err)?;
1705 let mut records = Vec::new();
1706 for row in rows {
1707 records.push(row.map_err(sql_err)?.into_record(&realm)?);
1708 }
1709 Ok(records)
1710 })
1711 })
1712 .await
1713 }
1714
1715 pub async fn record_pending_harvest(
1718 &self,
1719 realm: &str,
1720 identity: &str,
1721 session_key: Option<&str>,
1722 cause: &str,
1723 ) -> Result<(), AgentMemoryError> {
1724 let store = self.clone();
1725 let realm = realm.to_string();
1726 let identity = identity.to_string();
1727 let session_key = session_key.map(str::to_string);
1728 let cause = cause.to_string();
1729 run_blocking(move || {
1730 store.with_realm_conn(&realm, |conn| {
1731 conn.execute(
1732 "INSERT OR IGNORE INTO pending_harvests \
1733 (identity, session_key, cause, retired_at_ms, status) \
1734 VALUES (?1, ?2, ?3, ?4, 'pending')",
1735 params![identity, session_key, cause, now_ms() as i64],
1736 )
1737 .map_err(sql_err)?;
1738 Ok(())
1739 })
1740 })
1741 .await
1742 }
1743
1744 pub async fn pending_harvests(
1746 &self,
1747 realm: &str,
1748 limit: usize,
1749 ) -> Result<Vec<PendingHarvest>, AgentMemoryError> {
1750 let store = self.clone();
1751 let realm = realm.to_string();
1752 run_blocking(move || {
1753 store.with_realm_conn(&realm, |conn| {
1754 let mut stmt = conn
1755 .prepare(
1756 "SELECT identity, session_key, cause, retired_at_ms FROM \
1757 pending_harvests WHERE status = 'pending' \
1758 ORDER BY retired_at_ms ASC LIMIT ?1",
1759 )
1760 .map_err(sql_err)?;
1761 let rows = stmt
1762 .query_map(params![limit as i64], |row| {
1763 Ok(PendingHarvest {
1764 identity: row.get(0)?,
1765 session_key: row.get(1)?,
1766 cause: row.get(2)?,
1767 retired_at_ms: row.get::<_, i64>(3)? as u64,
1768 })
1769 })
1770 .map_err(sql_err)?;
1771 let mut harvests = Vec::new();
1772 for row in rows {
1773 harvests.push(row.map_err(sql_err)?);
1774 }
1775 Ok(harvests)
1776 })
1777 })
1778 .await
1779 }
1780
1781 pub async fn mark_harvest_complete(
1783 &self,
1784 realm: &str,
1785 identity: &str,
1786 retired_at_ms: u64,
1787 ) -> Result<(), AgentMemoryError> {
1788 let store = self.clone();
1789 let realm = realm.to_string();
1790 let identity = identity.to_string();
1791 run_blocking(move || {
1792 store.with_realm_conn(&realm, |conn| {
1793 conn.execute(
1794 "UPDATE pending_harvests SET status = 'harvested' \
1795 WHERE identity = ?1 AND retired_at_ms = ?2",
1796 params![identity, retired_at_ms as i64],
1797 )
1798 .map_err(sql_err)?;
1799 Ok(())
1800 })
1801 })
1802 .await
1803 }
1804
1805 pub async fn record_pending_promotion(
1808 &self,
1809 realm: &str,
1810 promotion: PendingPromotion,
1811 ) -> Result<(), AgentMemoryError> {
1812 let store = self.clone();
1813 let realm = realm.to_string();
1814 run_blocking(move || {
1815 store.with_realm_conn(&realm, |conn| {
1816 conn.execute(
1817 "INSERT INTO pending_promotions (pending_id, stage_token, record_id, \
1818 scope_kind, scope_key, rationale, status, created_at_ms) \
1819 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
1820 params![
1821 promotion.pending_id,
1822 promotion.stage_token,
1823 promotion.record_id,
1824 promotion.scope_kind,
1825 promotion.scope_key,
1826 promotion.rationale,
1827 promotion.status,
1828 promotion.created_at_ms as i64,
1829 ],
1830 )
1831 .map_err(sql_err)?;
1832 Ok(())
1833 })
1834 })
1835 .await
1836 }
1837
1838 pub async fn pending_promotion_by_id(
1840 &self,
1841 realm: &str,
1842 pending_id: &str,
1843 ) -> Result<Option<PendingPromotion>, AgentMemoryError> {
1844 let store = self.clone();
1845 let realm = realm.to_string();
1846 let pending_id = pending_id.to_string();
1847 run_blocking(move || {
1848 store.with_realm_conn(&realm, |conn| {
1849 conn.query_row(
1850 "SELECT pending_id, stage_token, record_id, scope_kind, scope_key, \
1851 rationale, status, created_at_ms FROM pending_promotions \
1852 WHERE pending_id = ?1 AND status = 'pending'",
1853 params![pending_id],
1854 |row| {
1855 Ok(PendingPromotion {
1856 pending_id: row.get(0)?,
1857 stage_token: row.get(1)?,
1858 record_id: row.get(2)?,
1859 scope_kind: row.get(3)?,
1860 scope_key: row.get(4)?,
1861 rationale: row.get(5)?,
1862 status: row.get(6)?,
1863 created_at_ms: row.get::<_, i64>(7)? as u64,
1864 })
1865 },
1866 )
1867 .optional()
1868 .map_err(sql_err)
1869 })
1870 })
1871 .await
1872 }
1873
1874 pub async fn pending_promotions(
1876 &self,
1877 realm: &str,
1878 ) -> Result<Vec<PendingPromotion>, AgentMemoryError> {
1879 let store = self.clone();
1880 let realm = realm.to_string();
1881 run_blocking(move || {
1882 store.with_realm_conn(&realm, |conn| {
1883 let mut stmt = conn
1884 .prepare(
1885 "SELECT pending_id, stage_token, record_id, scope_kind, scope_key, \
1886 rationale, status, created_at_ms FROM pending_promotions \
1887 WHERE status = 'pending' ORDER BY created_at_ms ASC",
1888 )
1889 .map_err(sql_err)?;
1890 let rows = stmt
1891 .query_map([], |row| {
1892 Ok(PendingPromotion {
1893 pending_id: row.get(0)?,
1894 stage_token: row.get(1)?,
1895 record_id: row.get(2)?,
1896 scope_kind: row.get(3)?,
1897 scope_key: row.get(4)?,
1898 rationale: row.get(5)?,
1899 status: row.get(6)?,
1900 created_at_ms: row.get::<_, i64>(7)? as u64,
1901 })
1902 })
1903 .map_err(sql_err)?;
1904 let mut promotions = Vec::new();
1905 for row in rows {
1906 promotions.push(row.map_err(sql_err)?);
1907 }
1908 Ok(promotions)
1909 })
1910 })
1911 .await
1912 }
1913
1914 pub async fn resolve_pending_promotion(
1916 &self,
1917 realm: &str,
1918 pending_id: &str,
1919 status: &str,
1920 ) -> Result<(), AgentMemoryError> {
1921 if !matches!(status, "committed" | "denied" | "expired") {
1922 return Err(AgentMemoryError::InvalidRecord(format!(
1923 "unknown promotion resolution '{status}'"
1924 )));
1925 }
1926 let store = self.clone();
1927 let realm = realm.to_string();
1928 let pending_id = pending_id.to_string();
1929 let status = status.to_string();
1930 run_blocking(move || {
1931 store.with_realm_conn(&realm, |conn| {
1932 conn.execute(
1933 "UPDATE pending_promotions SET status = ?1, resolved_at_ms = ?2 \
1934 WHERE pending_id = ?3",
1935 params![status, now_ms() as i64, pending_id],
1936 )
1937 .map_err(sql_err)?;
1938 Ok(())
1939 })
1940 })
1941 .await
1942 }
1943
1944 pub async fn rekey_pending_promotion(
1947 &self,
1948 realm: &str,
1949 old_pending_id: &str,
1950 new_pending_id: &str,
1951 ) -> Result<(), AgentMemoryError> {
1952 let store = self.clone();
1953 let realm = realm.to_string();
1954 let old_pending_id = old_pending_id.to_string();
1955 let new_pending_id = new_pending_id.to_string();
1956 run_blocking(move || {
1957 store.with_realm_conn(&realm, |conn| {
1958 conn.execute(
1959 "UPDATE pending_promotions SET pending_id = ?1 WHERE pending_id = ?2",
1960 params![new_pending_id, old_pending_id],
1961 )
1962 .map_err(sql_err)?;
1963 Ok(())
1964 })
1965 })
1966 .await
1967 }
1968
1969 pub async fn discard_stage(&self, token: StageToken) -> Result<(), AgentMemoryError> {
1973 let store = self.clone();
1974 run_blocking(move || {
1975 store.with_realm_conn(&token.realm, |conn| {
1976 conn.execute("DELETE FROM stage WHERE token = ?1", params![token.token])
1977 .map_err(sql_err)?;
1978 Ok(())
1979 })
1980 })
1981 .await
1982 }
1983}
1984
1985#[derive(Debug, Clone, PartialEq, Eq)]
1990pub struct PanelRecordsPage {
1991 pub records: Vec<super::records::MemoryRecord>,
1992 pub next_cursor: Option<(u64, String)>,
1994}
1995
1996#[derive(Debug, Clone, Default, PartialEq, Eq)]
1999pub struct DreamRunAudit {
2000 pub run_id: String,
2001 pub first_op_at_ms: u64,
2002 pub last_op_at_ms: u64,
2003 pub ops: u64,
2004 pub op_kinds: std::collections::BTreeMap<String, u64>,
2006 pub quarantined_ops: u64,
2008 pub memory_ids: Vec<String>,
2010 pub rationales: Vec<String>,
2012}
2013
2014const DREAM_HISTORY_SCAN_ROWS: usize = 5_000;
2018const DREAM_HISTORY_ID_SAMPLE: usize = 12;
2019const DREAM_HISTORY_RATIONALE_SAMPLE: usize = 6;
2020
2021impl SqliteAgentMemoryStore {
2022 pub async fn panel_realms(&self) -> Result<Vec<String>, AgentMemoryError> {
2024 let store = self.clone();
2025 run_blocking(move || store.known_realms()).await
2026 }
2027
2028 pub async fn record_by_id(
2030 &self,
2031 realm: &str,
2032 memory_id: &str,
2033 ) -> Result<Option<super::records::MemoryRecord>, AgentMemoryError> {
2034 let store = self.clone();
2035 let realm = realm.to_string();
2036 let memory_id = memory_id.to_string();
2037 run_blocking(move || {
2038 store.with_realm_conn(&realm, |conn| load_record(conn, &realm, &memory_id))
2039 })
2040 .await
2041 }
2042
2043 pub async fn records_page(
2047 &self,
2048 realm: &str,
2049 scope_kind: Option<&str>,
2050 scope_key: Option<&str>,
2051 status_kind: Option<&str>,
2052 limit: usize,
2053 cursor: Option<(u64, String)>,
2054 ) -> Result<PanelRecordsPage, AgentMemoryError> {
2055 let store = self.clone();
2056 let realm = realm.to_string();
2057 let scope_kind = scope_kind.map(str::to_string);
2058 let scope_key = scope_key.map(str::to_string);
2059 let status_kind = status_kind.map(str::to_string);
2060 let limit = limit.max(1);
2061 run_blocking(move || {
2062 store.with_realm_conn(&realm, |conn| {
2063 let mut clauses: Vec<String> = Vec::new();
2064 let mut values: Vec<rusqlite::types::Value> = Vec::new();
2065 if let Some(kind) = &scope_kind {
2066 values.push(kind.clone().into());
2067 clauses.push(format!("scope_kind = ?{}", values.len()));
2068 }
2069 if let Some(key) = &scope_key {
2070 values.push(key.clone().into());
2071 clauses.push(format!("scope_key = ?{}", values.len()));
2072 }
2073 if let Some(status) = &status_kind {
2074 values.push(status.clone().into());
2075 clauses.push(format!("status_kind = ?{}", values.len()));
2076 }
2077 if let Some((after_ms, after_id)) = &cursor {
2078 values.push((*after_ms as i64).into());
2079 let ms_slot = values.len();
2080 values.push(after_id.clone().into());
2081 let id_slot = values.len();
2082 clauses.push(format!(
2083 "(updated_at_ms < ?{ms_slot} OR (updated_at_ms = ?{ms_slot} \
2084 AND memory_id < ?{id_slot}))"
2085 ));
2086 }
2087 let where_sql = if clauses.is_empty() {
2088 String::new()
2089 } else {
2090 format!("WHERE {}", clauses.join(" AND "))
2091 };
2092 values.push(((limit + 1) as i64).into());
2093 let sql = format!(
2094 "SELECT {RECORD_COLUMNS} FROM records {where_sql} \
2095 ORDER BY updated_at_ms DESC, memory_id DESC LIMIT ?{}",
2096 values.len()
2097 );
2098 let mut stmt = conn.prepare(&sql).map_err(sql_err)?;
2099 let rows = stmt
2100 .query_map(rusqlite::params_from_iter(values), row_to_record_row)
2101 .map_err(sql_err)?;
2102 let mut records = Vec::new();
2103 for row in rows {
2104 records.push(row.map_err(sql_err)?.into_record(&realm)?);
2105 }
2106 let next_cursor = if records.len() > limit {
2107 records.truncate(limit);
2108 records
2109 .last()
2110 .map(|record| (record.updated_at_ms, record.id.clone()))
2111 } else {
2112 None
2113 };
2114 Ok(PanelRecordsPage {
2115 records,
2116 next_cursor,
2117 })
2118 })
2119 })
2120 .await
2121 }
2122
2123 pub async fn supersede_chain(
2131 &self,
2132 realm: &str,
2133 memory_id: &str,
2134 max_len: usize,
2135 ) -> Result<Vec<super::records::MemoryRecord>, AgentMemoryError> {
2136 let store = self.clone();
2137 let realm = realm.to_string();
2138 let memory_id = memory_id.to_string();
2139 let max_len = max_len.max(1);
2140 run_blocking(move || {
2141 store.with_realm_conn(&realm, |conn| {
2142 let Some(origin) = load_record(conn, &realm, &memory_id)? else {
2143 return Ok(Vec::new());
2144 };
2145 let mut seen: std::collections::BTreeSet<String> =
2146 std::collections::BTreeSet::from([origin.id.clone()]);
2147 let mut ancestors: Vec<super::records::MemoryRecord> = Vec::new();
2148 let mut parent_id = origin.supersedes.clone();
2149 while let Some(id) = parent_id {
2150 if ancestors.len() + 1 >= max_len || !seen.insert(id.clone()) {
2151 break;
2152 }
2153 let Some(parent) = load_record(conn, &realm, &id)? else {
2154 break;
2155 };
2156 parent_id = parent.supersedes.clone();
2157 ancestors.push(parent);
2158 }
2159 ancestors.reverse();
2160 let mut chain = ancestors;
2161 chain.push(origin);
2162 loop {
2163 if chain.len() >= max_len {
2164 return Ok(chain);
2165 }
2166 let tip = chain.last().unwrap_or_else(|| unreachable!());
2167 let successor_id = match &tip.status {
2168 super::records::RecordStatus::Superseded { by } => Some(by.clone()),
2169 _ => None,
2170 };
2171 match successor_id {
2172 Some(id) => {
2173 if !seen.insert(id.clone()) {
2174 return Ok(chain);
2175 }
2176 let Some(successor) = load_record(conn, &realm, &id)? else {
2177 return Ok(chain);
2178 };
2179 chain.push(successor);
2180 }
2181 None => {
2182 let tip_id = tip.id.clone();
2184 let mut stmt = conn
2185 .prepare(&format!(
2186 "SELECT {RECORD_COLUMNS} FROM records \
2187 WHERE supersedes = ?1 ORDER BY created_at_ms ASC"
2188 ))
2189 .map_err(sql_err)?;
2190 let rows = stmt
2191 .query_map(params![tip_id], row_to_record_row)
2192 .map_err(sql_err)?;
2193 for row in rows {
2194 if chain.len() >= max_len {
2195 break;
2196 }
2197 let claimant = row.map_err(sql_err)?.into_record(&realm)?;
2198 if seen.insert(claimant.id.clone()) {
2199 chain.push(claimant);
2200 }
2201 }
2202 return Ok(chain);
2203 }
2204 }
2205 }
2206 })
2207 })
2208 .await
2209 }
2210
2211 pub async fn injection_log_for_record(
2213 &self,
2214 realm: &str,
2215 record_id: &str,
2216 limit: usize,
2217 ) -> Result<Vec<InjectionLogEntry>, AgentMemoryError> {
2218 let store = self.clone();
2219 let realm = realm.to_string();
2220 let record_id = record_id.to_string();
2221 run_blocking(move || {
2222 store.with_realm_conn(&realm, |conn| {
2223 let mut stmt = conn
2224 .prepare(
2225 "SELECT record_id, identity, session_key, surface, at_ms \
2226 FROM injections WHERE record_id = ?1 \
2227 ORDER BY at_ms DESC, injection_id DESC LIMIT ?2",
2228 )
2229 .map_err(sql_err)?;
2230 let rows = stmt
2231 .query_map(params![record_id, limit as i64], |row| {
2232 Ok((
2233 row.get::<_, String>(0)?,
2234 row.get::<_, String>(1)?,
2235 row.get::<_, Option<String>>(2)?,
2236 row.get::<_, String>(3)?,
2237 row.get::<_, i64>(4)?,
2238 ))
2239 })
2240 .map_err(sql_err)?;
2241 let mut entries = Vec::new();
2242 for row in rows {
2243 let (record_id, identity, session_key, surface, at_ms) =
2244 row.map_err(sql_err)?;
2245 let surface = InjectionSurface::parse(&surface).ok_or_else(|| {
2246 AgentMemoryError::Parse(format!("unknown injection surface '{surface}'"))
2247 })?;
2248 entries.push(InjectionLogEntry {
2249 record_id,
2250 identity,
2251 session_key,
2252 surface,
2253 at_ms: at_ms as u64,
2254 });
2255 }
2256 Ok(entries)
2257 })
2258 })
2259 .await
2260 }
2261
2262 pub async fn save_dream_run(
2268 &self,
2269 realm: &str,
2270 run: PersistedDreamRun,
2271 ) -> Result<(), AgentMemoryError> {
2272 let store = self.clone();
2273 let realm = realm.to_string();
2274 run_blocking(move || {
2275 store.with_realm_conn(&realm, |conn| {
2276 conn.execute(
2277 "INSERT OR REPLACE INTO dream_runs (run_id, partition_label, started_at_ms, completed_at_ms, ops_committed, detail) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
2278 rusqlite::params![
2279 run.run_id,
2280 run.partition_label,
2281 run.started_at_ms,
2282 run.completed_at_ms,
2283 run.ops_committed,
2284 run.detail,
2285 ],
2286 )
2287 .map_err(sql_err)?;
2288 Ok(())
2289 })
2290 })
2291 .await
2292 }
2293
2294 pub async fn dream_runs(
2296 &self,
2297 realm: &str,
2298 limit: usize,
2299 ) -> Result<Vec<PersistedDreamRun>, AgentMemoryError> {
2300 let store = self.clone();
2301 let realm = realm.to_string();
2302 let limit = limit.max(1);
2303 run_blocking(move || {
2304 store.with_realm_conn(&realm, |conn| {
2305 let mut stmt = conn
2306 .prepare(
2307 "SELECT run_id, partition_label, started_at_ms, completed_at_ms, ops_committed, detail FROM dream_runs ORDER BY completed_at_ms DESC, run_id DESC LIMIT ?1",
2308 )
2309 .map_err(sql_err)?;
2310 let rows = stmt
2311 .query_map([limit as i64], |row| {
2312 Ok(PersistedDreamRun {
2313 run_id: row.get(0)?,
2314 partition_label: row.get(1)?,
2315 started_at_ms: row.get(2)?,
2316 completed_at_ms: row.get(3)?,
2317 ops_committed: row.get(4)?,
2318 detail: row.get(5)?,
2319 })
2320 })
2321 .map_err(sql_err)?
2322 .collect::<Result<Vec<_>, _>>()
2323 .map_err(sql_err)?;
2324 Ok(rows)
2325 })
2326 })
2327 .await
2328 }
2329
2330 pub async fn save_dream_audit_verdicts(
2334 &self,
2335 realm: &str,
2336 run_id: &str,
2337 verdicts: Vec<(String, String, String)>,
2338 ) -> Result<(), AgentMemoryError> {
2339 if verdicts.is_empty() {
2340 return Ok(());
2341 }
2342 let store = self.clone();
2343 let realm = realm.to_string();
2344 let run_id = run_id.to_string();
2345 let now = now_ms();
2346 run_blocking(move || {
2347 store.with_realm_conn(&realm, |conn| {
2348 for (record_id, verdict, rationale) in &verdicts {
2349 conn.execute(
2350 "INSERT OR REPLACE INTO dream_audit_verdicts (run_id, record_id, verdict, rationale, created_at_ms) VALUES (?1, ?2, ?3, ?4, ?5)",
2351 rusqlite::params![run_id, record_id, verdict, rationale, now],
2352 )
2353 .map_err(sql_err)?;
2354 }
2355 Ok(())
2356 })
2357 })
2358 .await
2359 }
2360
2361 pub async fn open_dream_audit_verdicts(
2364 &self,
2365 realm: &str,
2366 limit: usize,
2367 ) -> Result<Vec<DreamAuditVerdict>, AgentMemoryError> {
2368 let store = self.clone();
2369 let realm = realm.to_string();
2370 let limit = limit.max(1);
2371 run_blocking(move || {
2372 store.with_realm_conn(&realm, |conn| {
2373 let mut stmt = conn
2374 .prepare(
2375 "SELECT run_id, record_id, verdict, rationale, created_at_ms, resolved_at_ms, resolution FROM dream_audit_verdicts WHERE resolved_at_ms IS NULL ORDER BY created_at_ms DESC, record_id ASC LIMIT ?1",
2376 )
2377 .map_err(sql_err)?;
2378 let rows = stmt
2379 .query_map([limit as i64], |row| {
2380 Ok(DreamAuditVerdict {
2381 run_id: row.get(0)?,
2382 record_id: row.get(1)?,
2383 verdict: row.get(2)?,
2384 rationale: row.get(3)?,
2385 created_at_ms: row.get(4)?,
2386 resolved_at_ms: row.get(5)?,
2387 resolution: row.get(6)?,
2388 })
2389 })
2390 .map_err(sql_err)?
2391 .collect::<Result<Vec<_>, _>>()
2392 .map_err(sql_err)?;
2393 Ok(rows)
2394 })
2395 })
2396 .await
2397 }
2398
2399 pub async fn resolve_dream_audit_verdicts(
2402 &self,
2403 realm: &str,
2404 record_id: &str,
2405 resolution: &str,
2406 ) -> Result<usize, AgentMemoryError> {
2407 let store = self.clone();
2408 let realm = realm.to_string();
2409 let record_id = record_id.to_string();
2410 let resolution = resolution.to_string();
2411 let now = now_ms();
2412 run_blocking(move || {
2413 store.with_realm_conn(&realm, |conn| {
2414 let changed = conn
2415 .execute(
2416 "UPDATE dream_audit_verdicts SET resolved_at_ms = ?1, resolution = ?2 WHERE record_id = ?3 AND resolved_at_ms IS NULL",
2417 rusqlite::params![now, resolution, record_id],
2418 )
2419 .map_err(sql_err)?;
2420 Ok(changed)
2421 })
2422 })
2423 .await
2424 }
2425
2426 pub async fn dream_history(
2427 &self,
2428 realm: &str,
2429 max_runs: usize,
2430 ) -> Result<Vec<DreamRunAudit>, AgentMemoryError> {
2431 let store = self.clone();
2432 let realm = realm.to_string();
2433 let max_runs = max_runs.max(1);
2434 run_blocking(move || {
2435 store.with_realm_conn(&realm, |conn| {
2436 let mut stmt = conn
2437 .prepare(
2438 "SELECT op_kind, memory_id, detail, applied_at_ms FROM audit \
2439 ORDER BY applied_at_ms DESC, audit_id DESC LIMIT ?1",
2440 )
2441 .map_err(sql_err)?;
2442 let rows = stmt
2443 .query_map(params![DREAM_HISTORY_SCAN_ROWS as i64], |row| {
2444 Ok((
2445 row.get::<_, String>(0)?,
2446 row.get::<_, Option<String>>(1)?,
2447 row.get::<_, String>(2)?,
2448 row.get::<_, i64>(3)?,
2449 ))
2450 })
2451 .map_err(sql_err)?;
2452 let mut order: Vec<String> = Vec::new();
2453 let mut runs: HashMap<String, DreamRunAudit> = HashMap::new();
2454 for row in rows {
2455 let (op_kind, memory_id, detail, applied_at_ms) = row.map_err(sql_err)?;
2456 let detail: serde_json::Value =
2457 serde_json::from_str(&detail).unwrap_or_default();
2458 let author = detail.get("author");
2459 let is_steward = author
2460 .and_then(|author| author.get("author"))
2461 .and_then(serde_json::Value::as_str)
2462 == Some("steward");
2463 if !is_steward {
2464 continue;
2465 }
2466 let Some(run_id) = author
2467 .and_then(|author| author.get("run_id"))
2468 .and_then(serde_json::Value::as_str)
2469 else {
2470 continue;
2471 };
2472 if !runs.contains_key(run_id) {
2473 if runs.len() >= max_runs {
2474 continue;
2475 }
2476 order.push(run_id.to_string());
2477 }
2478 let run = runs
2479 .entry(run_id.to_string())
2480 .or_insert_with(|| DreamRunAudit {
2481 run_id: run_id.to_string(),
2482 first_op_at_ms: applied_at_ms as u64,
2483 last_op_at_ms: applied_at_ms as u64,
2484 ..DreamRunAudit::default()
2485 });
2486 run.ops += 1;
2487 run.first_op_at_ms = run.first_op_at_ms.min(applied_at_ms as u64);
2488 run.last_op_at_ms = run.last_op_at_ms.max(applied_at_ms as u64);
2489 *run.op_kinds.entry(op_kind).or_insert(0) += 1;
2490 if !detail
2491 .get("quarantined")
2492 .map(serde_json::Value::is_null)
2493 .unwrap_or(true)
2494 {
2495 run.quarantined_ops += 1;
2496 }
2497 if let Some(memory_id) = memory_id
2498 && run.memory_ids.len() < DREAM_HISTORY_ID_SAMPLE
2499 {
2500 run.memory_ids.push(memory_id);
2501 }
2502 if let Some(rationale) =
2503 detail.get("rationale").and_then(serde_json::Value::as_str)
2504 && !rationale.is_empty()
2505 && run.rationales.len() < DREAM_HISTORY_RATIONALE_SAMPLE
2506 {
2507 run.rationales.push(rationale.to_string());
2508 }
2509 }
2510 Ok(order
2511 .into_iter()
2512 .filter_map(|run_id| runs.remove(&run_id))
2513 .collect())
2514 })
2515 })
2516 .await
2517 }
2518}
2519
2520#[async_trait]
2521impl crate::memory::distiller::TombstoneSource for SqliteAgentMemoryStore {
2522 async fn recent_tombstones(
2527 &self,
2528 scope: &MemoryScope,
2529 since_ms: u64,
2530 limit: usize,
2531 ) -> Result<Vec<crate::memory::distiller::TombstoneMeta>, AgentMemoryError> {
2532 let store = self.clone();
2533 let scope = scope.clone();
2534 run_blocking(move || {
2535 store.with_realm_conn(scope.realm(), |conn| {
2536 let mut statement = conn
2537 .prepare(
2538 "SELECT title, kind, tombstoned_at_ms FROM records \
2539 WHERE scope_kind = ?1 AND scope_key = ?2 \
2540 AND status_kind = 'tombstoned' AND tombstoned_at_ms >= ?3 \
2541 ORDER BY tombstoned_at_ms DESC LIMIT ?4",
2542 )
2543 .map_err(sql_err)?;
2544 let rows = statement
2545 .query_map(
2546 params![scope.kind_str(), scope.key(), since_ms as i64, limit as i64],
2547 |row| {
2548 Ok((
2549 row.get::<_, String>(0)?,
2550 row.get::<_, String>(1)?,
2551 row.get::<_, i64>(2)?,
2552 ))
2553 },
2554 )
2555 .map_err(sql_err)?;
2556 let mut tombstones = Vec::new();
2557 for row in rows {
2558 let (title, kind, tombstoned_at_ms) = row.map_err(sql_err)?;
2559 let kind = MemoryKind::parse(&kind).ok_or_else(|| {
2560 AgentMemoryError::Parse(format!("unknown record kind '{kind}'"))
2561 })?;
2562 tombstones.push(crate::memory::distiller::TombstoneMeta {
2563 title,
2564 kind,
2565 tombstoned_at_ms: tombstoned_at_ms as u64,
2566 });
2567 }
2568 Ok(tombstones)
2569 })
2570 })
2571 .await
2572 }
2573}
2574
2575#[async_trait]
2580impl crate::memory::selector::SelectedRecordFetch for SqliteAgentMemoryStore {
2581 async fn fetch_records(
2582 &self,
2583 scopes: &[MemoryScope],
2584 ids: &[String],
2585 ) -> Result<Vec<AgentMemoryRecord>, AgentMemoryError> {
2586 let store = self.clone();
2587 let scopes = scopes.to_vec();
2588 let ids = ids.to_vec();
2589 run_blocking(move || {
2590 let mut records = Vec::new();
2591 for id in &ids {
2592 for scope in &scopes {
2593 let found = store.with_realm_conn(scope.realm(), |conn| {
2594 load_record(conn, scope.realm(), id)
2595 })?;
2596 if let Some(record) = found
2597 && record.scope == *scope
2598 && matches!(record.status, RecordStatus::Active)
2599 {
2600 records.push(project_record(record));
2601 break;
2602 }
2603 }
2604 }
2605 Ok(records)
2606 })
2607 .await
2608 }
2609
2610 async fn fetch_records_annotated(
2611 &self,
2612 scopes: &[MemoryScope],
2613 ids: &[String],
2614 ) -> Result<Vec<crate::memory::selector::AnnotatedRecord>, AgentMemoryError> {
2615 let store = self.clone();
2616 let scopes = scopes.to_vec();
2617 let ids = ids.to_vec();
2618 run_blocking(move || {
2619 let mut records = Vec::new();
2620 for id in &ids {
2621 for scope in &scopes {
2622 let found = store.with_realm_conn(scope.realm(), |conn| {
2623 load_record(conn, scope.realm(), id)
2624 })?;
2625 if let Some(record) = found
2626 && record.scope == *scope
2627 && matches!(record.status, RecordStatus::Active)
2628 {
2629 let provenance = Some(crate::memory::selector::RecordProvenance {
2633 scope: record.scope.clone(),
2634 trust: record.trust,
2635 });
2636 records.push(crate::memory::selector::AnnotatedRecord {
2637 record: project_record(record),
2638 provenance,
2639 });
2640 break;
2641 }
2642 }
2643 }
2644 Ok(records)
2645 })
2646 .await
2647 }
2648}
2649
2650async fn run_blocking<T: Send + 'static>(
2653 f: impl FnOnce() -> Result<T, AgentMemoryError> + Send + 'static,
2654) -> Result<T, AgentMemoryError> {
2655 tokio::task::spawn_blocking(f)
2656 .await
2657 .map_err(|err| AgentMemoryError::Io(format!("agent memory task failed: {err}")))?
2658}
2659
2660fn check_verified_retier_evidence(
2666 conn: &Connection,
2667 batch: &StagedMutationBatch,
2668 resolver: Option<&dyn EvidenceRefResolver>,
2669) -> Result<(), AgentMemoryError> {
2670 let Some(resolver) = resolver else {
2671 return Ok(());
2672 };
2673 for (op_index, op) in batch.ops.iter().enumerate() {
2674 let StagedOp::Retier { id, trust, .. } = op else {
2675 continue;
2676 };
2677 if *trust != TrustTier::AgentVerified {
2678 continue;
2679 }
2680 let reject = |reason: String| {
2681 AgentMemoryError::InvalidRecord(
2682 super::staged::StagedBatchError::UnresolvableEvidence { op_index, reason }
2683 .to_string(),
2684 )
2685 };
2686 let provenance: Option<String> = conn
2687 .query_row(
2688 "SELECT provenance FROM records WHERE memory_id = ?1",
2689 params![id],
2690 |row| row.get(0),
2691 )
2692 .optional()
2693 .map_err(sql_err)?;
2694 let Some(provenance) = provenance else {
2695 continue;
2697 };
2698 let provenance: MemoryProvenance = serde_json::from_str(&provenance)
2699 .map_err(|err| AgentMemoryError::Parse(err.to_string()))?;
2700 let evidence = provenance
2701 .verification
2702 .as_ref()
2703 .map(|claim| claim.evidence.as_slice())
2704 .unwrap_or(&[]);
2705 if evidence.is_empty() {
2706 return Err(reject(
2707 "verification claim cites no evidence refs".to_string(),
2708 ));
2709 }
2710 for reference in evidence {
2711 resolver.resolves(reference).map_err(reject)?;
2712 }
2713 }
2714 Ok(())
2715}
2716
2717fn apply_batch_tx(
2728 conn: &mut Connection,
2729 batch: &StagedMutationBatch,
2730 gate: Option<&dyn LlmWriteGate>,
2731 events: Option<&dyn crate::memory::events::MemoryEventSink>,
2732 token: &str,
2733 now: u64,
2734) -> Result<CommitReceipt, AgentMemoryError> {
2735 let evidence: Vec<crate::memory::records::EvidenceRef> = batch
2736 .ops
2737 .iter()
2738 .flat_map(|op| match op {
2739 StagedOp::Create { record, .. } | StagedOp::Supersede { record, .. } => {
2740 record.evidence.clone()
2741 }
2742 _ => Vec::new(),
2743 })
2744 .collect();
2745 let quarantine =
2746 gate.and_then(|gate| gate.quarantine_reason(&batch.author, batch.kind, &evidence));
2747 if let Some(reason) = quarantine.as_deref() {
2748 tracing::warn!(
2749 realm = %batch.realm,
2750 author = ?batch.author,
2751 reason,
2752 "agent memory: LLM-authored write landing quarantined (write-only until review)"
2753 );
2754 if let Some(events) = events {
2755 events.emit(
2756 crate::memory::events::MemoryTimelineEvent::QuarantinedWrite {
2757 realm: batch.realm.clone(),
2758 author: format!("{:?}", batch.author),
2759 reason: reason.to_string(),
2760 },
2761 );
2762 }
2763 }
2764 let tx = conn.transaction().map_err(sql_err)?;
2765 {
2766 let view = ConnBatchView {
2767 conn: &tx,
2768 realm: &batch.realm,
2769 };
2770 validate_batch(batch, &view, DEFAULT_TOMBSTONE_RECREATE_WINDOW_MS, now)
2771 .map_err(|err| AgentMemoryError::InvalidRecord(err.to_string()))?;
2772 }
2773 let mut memory_ids = Vec::with_capacity(batch.ops.len());
2774 for (op_index, op) in batch.ops.iter().enumerate() {
2775 let memory_id = apply_op(&tx, batch, op, quarantine.as_deref(), now)?;
2776 let detail = serde_json::json!({
2777 "op": op.kind_str(),
2778 "author": batch.author,
2779 "rationale": op_rationale(op),
2780 "quarantined": quarantine,
2781 });
2782 tx.execute(
2783 "INSERT INTO audit (stage_token, op_index, op_kind, memory_id, detail, \
2784 applied_at_ms) VALUES (?1, ?2, ?3, ?4, ?5, ?6)",
2785 params![
2786 token,
2787 op_index as i64,
2788 op.kind_str(),
2789 memory_id,
2790 detail.to_string(),
2791 now as i64,
2792 ],
2793 )
2794 .map_err(sql_err)?;
2795 memory_ids.push(memory_id);
2796 }
2797 tx.execute("DELETE FROM stage WHERE token = ?1", params![token])
2798 .map_err(sql_err)?;
2799 tx.commit().map_err(sql_err)?;
2800 Ok(CommitReceipt {
2801 token: token.to_string(),
2802 applied_ops: batch.ops.len(),
2803 memory_ids,
2804 })
2805}
2806
2807fn op_rationale(op: &StagedOp) -> Option<String> {
2808 match op {
2809 StagedOp::Create { rationale, .. }
2810 | StagedOp::Supersede { rationale, .. }
2811 | StagedOp::Tombstone { rationale, .. }
2812 | StagedOp::Retier { rationale, .. } => rationale.clone(),
2813 StagedOp::SetRank { .. } => None,
2814 }
2815}
2816
2817fn apply_op(
2818 conn: &Connection,
2819 batch: &StagedMutationBatch,
2820 op: &StagedOp,
2821 quarantine: Option<&str>,
2822 now: u64,
2823) -> Result<MemoryId, AgentMemoryError> {
2824 match op {
2825 StagedOp::Create {
2826 id,
2827 scope,
2828 record,
2829 trust,
2830 derived_from,
2831 created_at_ms,
2832 updated_at_ms,
2833 ..
2834 } => {
2835 let memory_id = id
2836 .clone()
2837 .unwrap_or_else(|| new_memory_id(&record.title, &record.body));
2838 insert_record(
2839 conn,
2840 &memory_id,
2841 scope,
2842 record,
2843 *trust,
2844 &batch.author,
2845 derived_from,
2846 None,
2847 None,
2848 None,
2849 quarantine,
2850 created_at_ms.unwrap_or(now),
2851 updated_at_ms.unwrap_or(now),
2852 )?;
2853 Ok(memory_id)
2854 }
2855 StagedOp::Supersede {
2856 id,
2857 prior,
2858 record,
2859 trust,
2860 derived_from,
2861 ..
2862 } => {
2863 let prior_row: (String, String, Option<i64>) = conn
2864 .query_row(
2865 "SELECT scope_kind, scope_key, working_set_rank \
2866 FROM records WHERE memory_id = ?1",
2867 params![prior],
2868 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
2869 )
2870 .map_err(sql_err)?;
2871 let scope = scope_from_parts(&prior_row.0, &prior_row.1, &batch.realm)?;
2872 let memory_id = id
2873 .clone()
2874 .unwrap_or_else(|| new_memory_id(&record.title, &record.body));
2875 insert_record(
2880 conn,
2881 &memory_id,
2882 &scope,
2883 record,
2884 *trust,
2885 &batch.author,
2886 derived_from,
2887 Some(prior.clone()),
2888 prior_row.2,
2889 None,
2890 quarantine,
2891 now,
2892 now,
2893 )?;
2894 if quarantine.is_none() {
2895 conn.execute(
2896 "UPDATE records SET status_kind = 'superseded', status_detail = ?1, \
2897 updated_at_ms = ?2 WHERE memory_id = ?3",
2898 params![memory_id, now as i64, prior],
2899 )
2900 .map_err(sql_err)?;
2901 } else {
2902 tracing::warn!(
2909 prior,
2910 successor = %memory_id,
2911 "agent memory: quarantined supersede leaves the prior record active \
2912 pending review"
2913 );
2914 }
2915 Ok(memory_id)
2916 }
2917 StagedOp::Tombstone { id, .. } => {
2918 conn.execute(
2919 "UPDATE records SET status_kind = 'tombstoned', status_detail = NULL, \
2920 tombstoned_at_ms = ?1, updated_at_ms = ?1 WHERE memory_id = ?2",
2921 params![now as i64, id],
2922 )
2923 .map_err(sql_err)?;
2924 Ok(id.clone())
2925 }
2926 StagedOp::Retier { id, trust, .. } => {
2927 conn.execute(
2928 "UPDATE records SET trust = ?1, updated_at_ms = ?2 WHERE memory_id = ?3",
2929 params![trust.as_str(), now as i64, id],
2930 )
2931 .map_err(sql_err)?;
2932 Ok(id.clone())
2933 }
2934 StagedOp::SetRank { id, rank } => {
2935 conn.execute(
2939 "UPDATE records SET working_set_rank = ?1, rank_set_at_ms = ?2 \
2940 WHERE memory_id = ?3",
2941 params![rank.map(|r| r as i64), now as i64, id],
2942 )
2943 .map_err(sql_err)?;
2944 Ok(id.clone())
2945 }
2946 }
2947}
2948
2949#[allow(clippy::too_many_arguments)]
2950fn insert_record(
2951 conn: &Connection,
2952 memory_id: &str,
2953 scope: &MemoryScope,
2954 record: &NewMemoryRecord,
2955 trust: TrustTier,
2956 author: &MemoryAuthor,
2957 derived_from: &[MemoryId],
2958 supersedes: Option<MemoryId>,
2959 working_set_rank: Option<i64>,
2960 rank_set_at_ms: Option<i64>,
2961 quarantine: Option<&str>,
2962 created_at_ms: u64,
2963 updated_at_ms: u64,
2964) -> Result<(), AgentMemoryError> {
2965 let tags = normalize_tags(record.tags.clone())?;
2966 let provenance = MemoryProvenance {
2967 evidence: record.evidence.clone(),
2968 author: author.clone(),
2969 profile: None,
2970 verification: record.verification.clone(),
2971 };
2972 let (status_kind, status_detail) = match quarantine {
2975 Some(reason) => ("quarantined", Some(reason)),
2976 None => ("active", None),
2977 };
2978 let ever_quarantined = quarantine.is_some() || {
2984 let mut ancestors: Vec<&str> = derived_from.iter().map(String::as_str).collect();
2985 if let Some(prior) = supersedes.as_deref() {
2986 ancestors.push(prior);
2987 }
2988 ancestors_reach_quarantine(conn, &ancestors)?
2989 };
2990 conn.execute(
2991 "INSERT INTO records (memory_id, scope_kind, scope_key, kind, title, description, \
2992 body, tags, provenance, trust, status_kind, status_detail, supersedes, derived_from, \
2993 working_set_rank, rank_set_at_ms, content_hash, created_at_ms, updated_at_ms, \
2994 usage_stats, tombstoned_at_ms, ever_quarantined) \
2995 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, \
2996 ?16, ?17, ?18, ?19, ?20, NULL, ?21)",
2997 params![
2998 memory_id,
2999 scope.kind_str(),
3000 scope.key(),
3001 record.kind.as_str(),
3002 record.title,
3003 record.description,
3004 record.body,
3005 json_string(&tags)?,
3006 json_string(&provenance)?,
3007 trust.as_str(),
3008 status_kind,
3009 status_detail,
3010 supersedes,
3011 json_string(&derived_from.to_vec())?,
3012 working_set_rank,
3013 rank_set_at_ms,
3014 content_hash(&record.title, &record.body),
3015 created_at_ms as i64,
3016 updated_at_ms as i64,
3017 json_string(&UsageStats::default())?,
3018 ever_quarantined,
3019 ],
3020 )
3021 .map_err(sql_err)?;
3022 Ok(())
3023}
3024
3025fn ancestors_reach_quarantine(
3028 conn: &Connection,
3029 ancestors: &[&str],
3030) -> Result<bool, AgentMemoryError> {
3031 if ancestors.is_empty() {
3032 return Ok(false);
3033 }
3034 let placeholders = (1..=ancestors.len())
3035 .map(|slot| format!("?{slot}"))
3036 .collect::<Vec<_>>()
3037 .join(", ");
3038 let sql = format!(
3039 "SELECT 1 FROM records WHERE memory_id IN ({placeholders}) \
3040 AND (ever_quarantined = 1 OR status_kind = 'quarantined') LIMIT 1"
3041 );
3042 let hit: Option<i64> = conn
3043 .query_row(&sql, rusqlite::params_from_iter(ancestors.iter()), |row| {
3044 row.get(0)
3045 })
3046 .optional()
3047 .map_err(sql_err)?;
3048 Ok(hit.is_some())
3049}
3050
3051struct ConnBatchView<'a> {
3055 conn: &'a Connection,
3056 realm: &'a str,
3057}
3058
3059impl StagedBatchView for ConnBatchView<'_> {
3060 fn record(&self, id: &str) -> Option<StagedRecordView> {
3061 self.conn
3062 .query_row(
3063 "SELECT scope_kind, scope_key, trust, status_kind, status_detail, supersedes, \
3064 derived_from, content_hash, provenance, ever_quarantined \
3065 FROM records WHERE memory_id = ?1",
3066 params![id],
3067 |row| {
3068 let scope_kind: String = row.get(0)?;
3069 let scope_key: String = row.get(1)?;
3070 let trust: String = row.get(2)?;
3071 let status_kind: String = row.get(3)?;
3072 let status_detail: Option<String> = row.get(4)?;
3073 let supersedes: Option<String> = row.get(5)?;
3074 let derived_from: String = row.get(6)?;
3075 let hash: String = row.get(7)?;
3076 let provenance: String = row.get(8)?;
3077 let ever_quarantined: bool = row.get(9)?;
3078 Ok((
3079 scope_kind,
3080 scope_key,
3081 trust,
3082 status_kind,
3083 status_detail,
3084 supersedes,
3085 derived_from,
3086 hash,
3087 provenance,
3088 ever_quarantined,
3089 ))
3090 },
3091 )
3092 .optional()
3093 .ok()
3094 .flatten()
3095 .and_then(
3096 |(
3097 scope_kind,
3098 scope_key,
3099 trust,
3100 status_kind,
3101 status_detail,
3102 supersedes,
3103 derived_from,
3104 hash,
3105 provenance,
3106 ever_quarantined,
3107 )| {
3108 let scope = scope_from_parts(&scope_kind, &scope_key, self.realm).ok()?;
3109 let provenance: MemoryProvenance = serde_json::from_str(&provenance).ok()?;
3110 Some(StagedRecordView {
3111 scope,
3112 trust: TrustTier::parse(&trust)?,
3113 status: status_from_parts(&status_kind, status_detail),
3114 supersedes,
3115 derived_from: serde_json::from_str(&derived_from).unwrap_or_default(),
3116 content_hash: hash,
3117 has_verification: provenance.verification.is_some(),
3118 ever_quarantined,
3119 })
3120 },
3121 )
3122 }
3123
3124 fn tombstoned_at_ms(&self, scope: &MemoryScope, hash: &str) -> Option<u64> {
3125 self.conn
3126 .query_row(
3127 "SELECT MAX(tombstoned_at_ms) FROM records WHERE scope_kind = ?1 \
3128 AND scope_key = ?2 AND content_hash = ?3 AND status_kind = 'tombstoned'",
3129 params![scope.kind_str(), scope.key(), hash],
3130 |row| row.get::<_, Option<i64>>(0),
3131 )
3132 .ok()
3133 .flatten()
3134 .map(|ms| ms as u64)
3135 }
3136}
3137
3138struct MemoryRecordRow {
3141 memory_id: String,
3142 scope_kind: String,
3143 scope_key: String,
3144 kind: String,
3145 title: String,
3146 description: String,
3147 body: String,
3148 tags: String,
3149 provenance: String,
3150 trust: String,
3151 status_kind: String,
3152 status_detail: Option<String>,
3153 supersedes: Option<String>,
3154 derived_from: String,
3155 working_set_rank: Option<i64>,
3156 created_at_ms: i64,
3157 updated_at_ms: i64,
3158 usage_stats: String,
3159}
3160
3161fn row_to_record_row(row: &rusqlite::Row<'_>) -> rusqlite::Result<MemoryRecordRow> {
3162 Ok(MemoryRecordRow {
3163 memory_id: row.get(0)?,
3164 scope_kind: row.get(1)?,
3165 scope_key: row.get(2)?,
3166 kind: row.get(3)?,
3167 title: row.get(4)?,
3168 description: row.get(5)?,
3169 body: row.get(6)?,
3170 tags: row.get(7)?,
3171 provenance: row.get(8)?,
3172 trust: row.get(9)?,
3173 status_kind: row.get(10)?,
3174 status_detail: row.get(11)?,
3175 supersedes: row.get(12)?,
3176 derived_from: row.get(13)?,
3177 working_set_rank: row.get(14)?,
3178 created_at_ms: row.get(17)?,
3179 updated_at_ms: row.get(18)?,
3180 usage_stats: row.get(19)?,
3181 })
3182}
3183
3184impl MemoryRecordRow {
3185 fn into_record(self, realm: &str) -> Result<super::records::MemoryRecord, AgentMemoryError> {
3186 let scope = scope_from_parts(&self.scope_kind, &self.scope_key, realm)?;
3187 let provenance: MemoryProvenance = serde_json::from_str(&self.provenance)
3188 .map_err(|err| AgentMemoryError::Parse(err.to_string()))?;
3189 Ok(super::records::MemoryRecord {
3190 id: self.memory_id,
3191 scope,
3192 kind: MemoryKind::parse(&self.kind).ok_or_else(|| {
3193 AgentMemoryError::Parse(format!("unknown record kind '{}'", self.kind))
3194 })?,
3195 title: self.title,
3196 description: self.description,
3197 body: self.body,
3198 tags: serde_json::from_str(&self.tags).unwrap_or_default(),
3199 provenance,
3200 trust: TrustTier::parse(&self.trust).ok_or_else(|| {
3201 AgentMemoryError::Parse(format!("unknown trust tier '{}'", self.trust))
3202 })?,
3203 status: status_from_parts(&self.status_kind, self.status_detail),
3204 supersedes: self.supersedes,
3205 derived_from: serde_json::from_str(&self.derived_from).unwrap_or_default(),
3206 working_set_rank: self.working_set_rank.map(|rank| rank as u32),
3207 created_at_ms: self.created_at_ms as u64,
3208 updated_at_ms: self.updated_at_ms as u64,
3209 usage: serde_json::from_str(&self.usage_stats).unwrap_or_default(),
3210 })
3211 }
3212}
3213
3214fn status_from_parts(kind: &str, detail: Option<String>) -> RecordStatus {
3215 match kind {
3216 "superseded" => RecordStatus::Superseded {
3217 by: detail.unwrap_or_default(),
3218 },
3219 "quarantined" => RecordStatus::Quarantined {
3220 reason: detail.unwrap_or_default(),
3221 },
3222 "tombstoned" => RecordStatus::Tombstoned,
3223 _ => RecordStatus::Active,
3224 }
3225}
3226
3227fn scope_from_parts(kind: &str, key: &str, realm: &str) -> Result<MemoryScope, AgentMemoryError> {
3228 match kind {
3229 "identity" => Ok(MemoryScope::Identity {
3230 realm: realm.to_string(),
3231 identity: key.to_string(),
3232 }),
3233 "mob" => Ok(MemoryScope::Mob {
3234 realm: realm.to_string(),
3235 mob: key.to_string(),
3236 }),
3237 "operator" => Ok(MemoryScope::Operator {
3238 realm: realm.to_string(),
3239 operator: key.to_string(),
3240 }),
3241 "realm" => Ok(MemoryScope::Realm {
3242 realm: realm.to_string(),
3243 }),
3244 other => Err(AgentMemoryError::Parse(format!(
3245 "unknown scope kind '{other}'"
3246 ))),
3247 }
3248}
3249
3250fn active_scope_records(
3251 conn: &Connection,
3252 scope: &MemoryScope,
3253) -> Result<Vec<super::records::MemoryRecord>, AgentMemoryError> {
3254 let mut stmt = conn
3255 .prepare(&format!(
3256 "SELECT {RECORD_COLUMNS} FROM records WHERE scope_kind = ?1 AND scope_key = ?2 \
3257 AND status_kind = 'active'"
3258 ))
3259 .map_err(sql_err)?;
3260 let rows = stmt
3261 .query_map(params![scope.kind_str(), scope.key()], row_to_record_row)
3262 .map_err(sql_err)?;
3263 let mut records = Vec::new();
3264 for row in rows {
3265 records.push(row.map_err(sql_err)?.into_record(scope.realm())?);
3266 }
3267 Ok(records)
3268}
3269
3270fn load_record(
3271 conn: &Connection,
3272 realm: &str,
3273 memory_id: &str,
3274) -> Result<Option<super::records::MemoryRecord>, AgentMemoryError> {
3275 let row = conn
3276 .query_row(
3277 &format!("SELECT {RECORD_COLUMNS} FROM records WHERE memory_id = ?1"),
3278 params![memory_id],
3279 row_to_record_row,
3280 )
3281 .optional()
3282 .map_err(sql_err)?;
3283 row.map(|row| row.into_record(realm)).transpose()
3284}
3285
3286fn project_record(record: super::records::MemoryRecord) -> AgentMemoryRecord {
3290 AgentMemoryRecord {
3291 memory_id: record.id,
3292 title: record.title,
3293 body: record.body,
3294 tags: record.tags,
3295 created_at_ms: record.created_at_ms,
3296 updated_at_ms: record.updated_at_ms,
3297 }
3298}
3299
3300fn scope_manifest(
3304 conn: &Connection,
3305 scope: &MemoryScope,
3306 tier: ManifestTier,
3307 now: u64,
3308) -> Result<Vec<RecordMeta>, AgentMemoryError> {
3309 let to_meta = |row: &rusqlite::Row<'_>| -> rusqlite::Result<RecordMeta> {
3310 let kind: String = row.get(1)?;
3311 let updated_at: i64 = row.get(4)?;
3312 let rank: Option<i64> = row.get(5)?;
3313 Ok(RecordMeta {
3314 id: row.get(0)?,
3315 kind: MemoryKind::parse(&kind).unwrap_or(MemoryKind::Fact),
3316 title: row.get(2)?,
3317 description: row.get(3)?,
3318 age_days: age_days(updated_at as u64, now),
3319 rank: rank.map(|rank| rank as u32),
3320 })
3321 };
3322 const META_COLUMNS: &str =
3323 "memory_id, kind, title, description, updated_at_ms, working_set_rank";
3324 match tier {
3325 ManifestTier::Full => {
3326 let mut stmt = conn
3327 .prepare(&format!(
3328 "SELECT {META_COLUMNS} FROM records \
3329 WHERE scope_kind = ?1 AND scope_key = ?2 AND status_kind = 'active' \
3330 ORDER BY (working_set_rank IS NULL) ASC, working_set_rank ASC, \
3331 updated_at_ms DESC, created_at_ms DESC, rowid DESC"
3332 ))
3333 .map_err(sql_err)?;
3334 let rows = stmt
3335 .query_map(params![scope.kind_str(), scope.key()], to_meta)
3336 .map_err(sql_err)?;
3337 rows.collect::<Result<Vec<_>, _>>().map_err(sql_err)
3338 }
3339 ManifestTier::WorkingSet(k) => {
3340 let mut stmt = conn
3341 .prepare(&format!(
3342 "SELECT {META_COLUMNS} FROM records \
3343 WHERE scope_kind = ?1 AND scope_key = ?2 AND status_kind = 'active' \
3344 AND working_set_rank IS NOT NULL \
3345 ORDER BY working_set_rank ASC, updated_at_ms DESC, rowid DESC LIMIT ?3"
3346 ))
3347 .map_err(sql_err)?;
3348 let ranked = stmt
3349 .query_map(params![scope.kind_str(), scope.key(), k as i64], to_meta)
3350 .map_err(sql_err)?
3351 .collect::<Result<Vec<_>, _>>()
3352 .map_err(sql_err)?;
3353 let mut stmt = conn
3354 .prepare(&format!(
3355 "SELECT {META_COLUMNS} FROM records \
3356 WHERE scope_kind = ?1 AND scope_key = ?2 AND status_kind = 'active' \
3357 AND (working_set_rank IS NULL \
3358 OR updated_at_ms > COALESCE(rank_set_at_ms, 0)) \
3359 ORDER BY updated_at_ms DESC, created_at_ms DESC, rowid DESC LIMIT ?3"
3360 ))
3361 .map_err(sql_err)?;
3362 let recent = stmt
3363 .query_map(
3364 params![scope.kind_str(), scope.key(), (2 * k) as i64],
3365 to_meta,
3366 )
3367 .map_err(sql_err)?
3368 .collect::<Result<Vec<_>, _>>()
3369 .map_err(sql_err)?;
3370 let cap = 2 * k;
3371 let mut seen = std::collections::HashSet::new();
3372 let mut union = Vec::new();
3373 for meta in ranked.into_iter().chain(recent) {
3374 if union.len() >= cap {
3375 break;
3376 }
3377 if seen.insert(meta.id.clone()) {
3378 union.push(meta);
3379 }
3380 }
3381 Ok(union)
3382 }
3383 }
3384}
3385
3386fn warn_if_scope_floors_exceeded(
3390 conn: &Connection,
3391 scope: &MemoryScope,
3392 floor_records: usize,
3393 floor_bytes: usize,
3394) -> Result<(), AgentMemoryError> {
3395 let (count, bytes): (i64, Option<i64>) = conn
3396 .query_row(
3397 "SELECT COUNT(*), SUM(LENGTH(title) + LENGTH(description) + LENGTH(body)) \
3398 FROM records WHERE scope_kind = ?1 AND scope_key = ?2 \
3399 AND status_kind != 'tombstoned'",
3400 params![scope.kind_str(), scope.key()],
3401 |row| Ok((row.get(0)?, row.get(1)?)),
3402 )
3403 .map_err(sql_err)?;
3404 if let Some(reason) = scope_floor_warning(
3405 count as usize,
3406 bytes.unwrap_or(0) as usize,
3407 floor_records,
3408 floor_bytes,
3409 ) {
3410 tracing::warn!(
3411 realm = scope.realm(),
3412 scope_kind = scope.kind_str(),
3413 scope_key = scope.key(),
3414 "agent memory scope exceeds retention floor ({reason}); steward consolidation \
3415 needed — records are never evicted automatically"
3416 );
3417 }
3418 Ok(())
3419}
3420
3421fn scope_floor_warning(
3423 count: usize,
3424 bytes: usize,
3425 floor_records: usize,
3426 floor_bytes: usize,
3427) -> Option<String> {
3428 if count > floor_records {
3429 return Some(format!("{count} records > floor {floor_records}"));
3430 }
3431 if bytes > floor_bytes {
3432 return Some(format!("{bytes} bytes > floor {floor_bytes}"));
3433 }
3434 None
3435}
3436
3437enum MarkdownImportError {
3440 Content(String),
3441 Io(AgentMemoryError),
3442}
3443
3444fn record_import_audit(
3447 conn: &Connection,
3448 file: &Path,
3449 imported: usize,
3450 skipped: usize,
3451 reasons: &[String],
3452) -> Result<(), AgentMemoryError> {
3453 const MAX_AUDITED_REASONS: usize = 8;
3454 let detail = serde_json::json!({
3455 "op": "markdown_import",
3456 "file": file.display().to_string(),
3457 "imported": imported,
3458 "skipped": skipped,
3459 "skip_reasons": reasons.iter().take(MAX_AUDITED_REASONS).collect::<Vec<_>>(),
3460 });
3461 conn.execute(
3462 "INSERT INTO audit (stage_token, op_index, op_kind, memory_id, detail, applied_at_ms) \
3463 VALUES (?1, 0, 'import_summary', NULL, ?2, ?3)",
3464 params![
3465 mint_token("import-audit"),
3466 detail.to_string(),
3467 now_ms() as i64,
3468 ],
3469 )
3470 .map_err(sql_err)?;
3471 Ok(())
3472}
3473
3474fn ensure_column(
3477 conn: &Connection,
3478 table: &str,
3479 column: &str,
3480 ddl: &str,
3481) -> Result<bool, AgentMemoryError> {
3482 let mut stmt = conn
3483 .prepare(&format!("PRAGMA table_info({table})"))
3484 .map_err(sql_err)?;
3485 let existing: Vec<String> = stmt
3486 .query_map([], |row| row.get::<_, String>(1))
3487 .map_err(sql_err)?
3488 .collect::<Result<_, _>>()
3489 .map_err(sql_err)?;
3490 if existing.iter().any(|name| name == column) {
3491 return Ok(false);
3492 }
3493 conn.execute(
3494 &format!("ALTER TABLE {table} ADD COLUMN {column} {ddl}"),
3495 [],
3496 )
3497 .map_err(sql_err)?;
3498 Ok(true)
3499}
3500
3501fn json_string<T: serde::Serialize>(value: &T) -> Result<String, AgentMemoryError> {
3502 serde_json::to_string(value).map_err(|err| AgentMemoryError::Parse(err.to_string()))
3503}
3504
3505fn sql_err(err: rusqlite::Error) -> AgentMemoryError {
3506 AgentMemoryError::Io(err.to_string())
3507}
3508
3509fn now_ms() -> u64 {
3510 SystemTime::now()
3511 .duration_since(UNIX_EPOCH)
3512 .map(|duration| duration.as_millis() as u64)
3513 .unwrap_or(0)
3514}
3515
3516fn mint_token(prefix: &str) -> String {
3517 static NEXT_TOKEN_SEQ: AtomicU64 = AtomicU64::new(0);
3518 let seq = NEXT_TOKEN_SEQ.fetch_add(1, Ordering::Relaxed);
3519 let nanos = SystemTime::now()
3520 .duration_since(UNIX_EPOCH)
3521 .map(|duration| duration.as_nanos())
3522 .unwrap_or(0);
3523 format!("{prefix}-{nanos}-{:x}-{seq:x}", std::process::id())
3524}
3525
3526#[cfg(test)]
3527#[allow(
3528 clippy::await_holding_lock,
3529 clippy::cloned_ref_to_slice_refs,
3530 clippy::expect_used,
3531 clippy::let_and_return,
3532 clippy::panic,
3533 clippy::unnecessary_to_owned
3534)]
3535mod tests {
3536 use super::*;
3537 use crate::identity_first::agent_memory::{AgentMemorySelection, MarkdownAgentMemoryStore};
3538 use std::error::Error;
3539
3540 fn identity() -> Result<AgentIdentity, Box<dyn Error>> {
3541 AgentIdentity::parse("identity:luka").map_err(|err| {
3542 std::io::Error::other(format!("test identity should parse: {err}")).into()
3543 })
3544 }
3545
3546 fn identity_scope(realm: &str) -> Result<MemoryScope, Box<dyn Error>> {
3547 Ok(MemoryScope::Identity {
3548 realm: realm.to_string(),
3549 identity: identity()?.as_str().to_string(),
3550 })
3551 }
3552
3553 fn new_memory(title: &str, body: &str) -> NewAgentMemory {
3554 NewAgentMemory {
3555 title: title.to_string(),
3556 body: body.to_string(),
3557 tags: Vec::new(),
3558 }
3559 }
3560
3561 fn recall_all(identity: AgentIdentity, realm: &str) -> AgentMemoryRecallRequest {
3562 AgentMemoryRecallRequest {
3563 identity,
3564 realm: realm.to_string(),
3565 query_text: None,
3566 query_terms: Vec::new(),
3567 selection: AgentMemorySelection::Always,
3568 max_entries: 64,
3569 }
3570 }
3571
3572 fn payload(title: &str, body: &str) -> NewMemoryRecord {
3573 NewMemoryRecord {
3574 kind: MemoryKind::Fact,
3575 title: title.to_string(),
3576 description: String::new(),
3577 body: body.to_string(),
3578 tags: Vec::new(),
3579 evidence: Vec::new(),
3580 verification: None,
3581 }
3582 }
3583
3584 #[tokio::test]
3585 async fn remember_dedups_exact_content_hash() -> Result<(), Box<dyn Error>> {
3586 let dir = tempfile::tempdir()?;
3587 let store = SqliteAgentMemoryStore::open(dir.path())?;
3588 let id = identity()?;
3589
3590 let first = store
3591 .remember("family", &id, new_memory("Same fact", "Same body"))
3592 .await?;
3593 let second = store
3594 .remember("family", &id, new_memory("Same fact", "Same body"))
3595 .await?;
3596 let third = store
3597 .remember("family", &id, new_memory("Other fact", "Other body"))
3598 .await?;
3599
3600 assert_eq!(
3601 first.memory_id, second.memory_id,
3602 "dedup must return the existing id"
3603 );
3604 assert_ne!(first.memory_id, third.memory_id);
3605 let records = store.recall(recall_all(id, "family")).await?;
3606 assert_eq!(records.len(), 2, "duplicate remember must not add a row");
3607 Ok(())
3608 }
3609
3610 #[tokio::test]
3611 async fn recall_scores_contextually_like_markdown_store() -> Result<(), Box<dyn Error>> {
3612 let dir = tempfile::tempdir()?;
3613 let store = SqliteAgentMemoryStore::open(dir.path())?;
3614 let id = identity()?;
3615 store
3616 .remember(
3617 "default",
3618 &id,
3619 NewAgentMemory {
3620 title: "Passport location".to_string(),
3621 body: "The passport is in the blue travel folder.".to_string(),
3622 tags: vec!["travel".to_string()],
3623 },
3624 )
3625 .await?;
3626 store
3627 .remember(
3628 "default",
3629 &id,
3630 new_memory("Unrelated", "Rust release checklist."),
3631 )
3632 .await?;
3633
3634 let matches = store
3635 .recall(AgentMemoryRecallRequest {
3636 identity: id,
3637 realm: "default".to_string(),
3638 query_text: Some("where did I put the passport".to_string()),
3639 query_terms: vec!["passport".to_string()],
3640 selection: AgentMemorySelection::Contextual,
3641 max_entries: 8,
3642 })
3643 .await?;
3644
3645 assert_eq!(matches.len(), 1);
3646 assert_eq!(matches[0].title, "Passport location");
3647 Ok(())
3648 }
3649
3650 #[tokio::test]
3651 async fn forget_tombstones_and_allows_deliberate_readd() -> Result<(), Box<dyn Error>> {
3652 let dir = tempfile::tempdir()?;
3653 let store = SqliteAgentMemoryStore::open(dir.path())?;
3654 let id = identity()?;
3655 let record = store
3656 .remember("family", &id, new_memory("Fact", "Body"))
3657 .await?;
3658
3659 let deleted = store.forget("family", &id, &record.memory_id).await?;
3660 assert!(deleted.deleted);
3661 assert!(
3662 store
3663 .recall(recall_all(id.clone(), "family"))
3664 .await?
3665 .is_empty()
3666 );
3667
3668 let again = store.forget("family", &id, &record.memory_id).await?;
3669 assert!(!again.deleted, "tombstoned record must not delete twice");
3670
3671 let readded = store
3675 .remember("family", &id, new_memory("Fact", "Body"))
3676 .await?;
3677 assert_ne!(readded.memory_id, record.memory_id);
3678 Ok(())
3679 }
3680
3681 #[tokio::test]
3682 async fn supersede_chains_and_inherits_rank() -> Result<(), Box<dyn Error>> {
3683 let dir = tempfile::tempdir()?;
3684 let store = SqliteAgentMemoryStore::open(dir.path())?;
3685 let id = identity()?;
3686 let scope = identity_scope("family")?;
3687 let prior = store
3688 .remember("family", &id, new_memory("DB host", "Use db-old.example."))
3689 .await?;
3690
3691 let token = store
3693 .stage(StagedMutationBatch {
3694 kind: StagedBatchKind::FreshWrite,
3695 realm: "family".to_string(),
3696 author: MemoryAuthor::Steward {
3697 run_id: "dream-1".to_string(),
3698 },
3699 ops: vec![StagedOp::SetRank {
3700 id: prior.memory_id.clone(),
3701 rank: Some(1),
3702 }],
3703 })
3704 .await?;
3705 store.commit(token).await?;
3706
3707 let new_id = store
3708 .supersede(
3709 &scope,
3710 &prior.memory_id,
3711 payload("DB host", "Use db-new.example."),
3712 )
3713 .await?;
3714 assert_ne!(new_id, prior.memory_id);
3715
3716 let records = store.recall(recall_all(id, "family")).await?;
3719 assert_eq!(records.len(), 1);
3720 assert_eq!(records[0].memory_id, new_id);
3721 assert!(records[0].body.contains("db-new"));
3722
3723 let manifest = store.manifest(&[scope.clone()], ManifestTier::Full).await?;
3724 assert_eq!(manifest.len(), 1);
3725 assert_eq!(manifest[0].id, new_id);
3726 assert_eq!(
3727 manifest[0].rank,
3728 Some(1),
3729 "supersede inherits the prior's rank"
3730 );
3731
3732 let conn = store.realm_connection("family")?;
3734 let guard = conn
3735 .lock()
3736 .unwrap_or_else(std::sync::PoisonError::into_inner);
3737 let (status_kind, by): (String, Option<String>) = guard.query_row(
3738 "SELECT status_kind, status_detail FROM records WHERE memory_id = ?1",
3739 params![prior.memory_id],
3740 |row| Ok((row.get(0)?, row.get(1)?)),
3741 )?;
3742 assert_eq!(status_kind, "superseded");
3743 assert_eq!(by.as_deref(), Some(new_id.as_str()));
3744 Ok(())
3745 }
3746
3747 #[tokio::test]
3748 async fn manifest_working_set_unions_ranked_and_recent() -> Result<(), Box<dyn Error>> {
3749 let dir = tempfile::tempdir()?;
3750 let store = SqliteAgentMemoryStore::open(dir.path())?;
3751 let id = identity()?;
3752 let scope = identity_scope("family")?;
3753 let mut ids = Vec::new();
3754 for i in 0..5 {
3755 let record = store
3756 .remember(
3757 "family",
3758 &id,
3759 new_memory(&format!("Fact {i}"), &format!("Body {i}")),
3760 )
3761 .await?;
3762 ids.push(record.memory_id);
3763 }
3764 let token = store
3767 .stage(StagedMutationBatch {
3768 kind: StagedBatchKind::FreshWrite,
3769 realm: "family".to_string(),
3770 author: MemoryAuthor::Steward {
3771 run_id: "dream-1".to_string(),
3772 },
3773 ops: (0..3)
3774 .map(|i| StagedOp::SetRank {
3775 id: ids[i].clone(),
3776 rank: Some(i as u32 + 1),
3777 })
3778 .collect(),
3779 })
3780 .await?;
3781 store.commit(token).await?;
3782
3783 let metas = store
3784 .manifest(&[scope.clone()], ManifestTier::WorkingSet(2))
3785 .await?;
3786 assert_eq!(metas.len(), 4);
3789 assert_eq!(metas[0].id, ids[0]);
3790 assert_eq!(metas[0].rank, Some(1));
3791 assert_eq!(metas[1].id, ids[1]);
3792 assert_eq!(metas[2].id, ids[4], "unranked slice is newest-first");
3793 assert_eq!(metas[3].id, ids[3]);
3794 assert!(
3795 !metas.iter().any(|meta| meta.id == ids[2]),
3796 "rank 3 is outside top-K and, being ranked and un-updated, outside the recent slice"
3797 );
3798
3799 let successor = store
3802 .supersede(&scope, &ids[2], payload("Fact 2", "Corrected body 2"))
3803 .await?;
3804 let metas = store
3805 .manifest(&[scope], ManifestTier::WorkingSet(2))
3806 .await?;
3807 assert!(
3808 metas.iter().any(|meta| meta.id == successor),
3809 "freshly superseded record must be selector-visible before the next dream: {metas:#?}"
3810 );
3811 Ok(())
3812 }
3813
3814 #[tokio::test]
3815 async fn staged_batch_without_commit_leaves_store_unchanged() -> Result<(), Box<dyn Error>> {
3816 let dir = tempfile::tempdir()?;
3817 let store = SqliteAgentMemoryStore::open(dir.path())?;
3818 let id = identity()?;
3819 let scope = identity_scope("family")?;
3820
3821 let token = store
3822 .stage(StagedMutationBatch {
3823 kind: StagedBatchKind::FreshWrite,
3824 realm: "family".to_string(),
3825 author: MemoryAuthor::Steward {
3826 run_id: "dream-crash".to_string(),
3827 },
3828 ops: vec![StagedOp::Create {
3829 id: None,
3830 scope: scope.clone(),
3831 record: payload("Staged fact", "Never committed"),
3832 trust: TrustTier::AgentObserved,
3833 derived_from: Vec::new(),
3834 rationale: None,
3835 created_at_ms: None,
3836 updated_at_ms: None,
3837 }],
3838 })
3839 .await?;
3840
3841 assert!(
3844 store
3845 .recall(recall_all(id.clone(), "family"))
3846 .await?
3847 .is_empty()
3848 );
3849 let reopened = SqliteAgentMemoryStore::open(dir.path())?;
3850 assert!(
3851 reopened
3852 .recall(recall_all(id.clone(), "family"))
3853 .await?
3854 .is_empty()
3855 );
3856
3857 let receipt = store.commit(token.clone()).await?;
3859 assert_eq!(receipt.applied_ops, 1);
3860 assert_eq!(store.recall(recall_all(id, "family")).await?.len(), 1);
3861 let replay = store.commit(token).await;
3862 assert!(matches!(replay, Err(AgentMemoryError::InvalidRecord(_))));
3863 Ok(())
3864 }
3865
3866 #[tokio::test]
3871 async fn ever_quarantined_migration_backfills_old_stores() -> Result<(), Box<dyn Error>> {
3872 let dir = tempfile::tempdir()?;
3873 let db_path = {
3874 let store = SqliteAgentMemoryStore::open(dir.path())?;
3875 store.path_for_realm("family")
3876 };
3877 {
3878 let conn = Connection::open(&db_path)?;
3879 conn.execute_batch(
3880 "CREATE TABLE records (
3881 memory_id TEXT PRIMARY KEY,
3882 scope_kind TEXT NOT NULL,
3883 scope_key TEXT NOT NULL,
3884 kind TEXT NOT NULL,
3885 title TEXT NOT NULL,
3886 description TEXT NOT NULL DEFAULT '',
3887 body TEXT NOT NULL,
3888 tags TEXT NOT NULL DEFAULT '[]',
3889 provenance TEXT NOT NULL,
3890 trust TEXT NOT NULL,
3891 status_kind TEXT NOT NULL,
3892 status_detail TEXT,
3893 supersedes TEXT,
3894 derived_from TEXT NOT NULL DEFAULT '[]',
3895 working_set_rank INTEGER,
3896 rank_set_at_ms INTEGER,
3897 content_hash TEXT NOT NULL,
3898 created_at_ms INTEGER NOT NULL,
3899 updated_at_ms INTEGER NOT NULL,
3900 usage_stats TEXT NOT NULL DEFAULT '{}',
3901 tombstoned_at_ms INTEGER
3902 );
3903 CREATE TABLE proposals (
3904 proposal_id TEXT PRIMARY KEY,
3905 scope_kind TEXT NOT NULL,
3906 scope_key TEXT NOT NULL,
3907 record TEXT NOT NULL,
3908 author TEXT NOT NULL,
3909 status TEXT NOT NULL DEFAULT 'pending',
3910 created_at_ms INTEGER NOT NULL
3911 );
3912 CREATE TABLE audit (
3913 audit_id INTEGER PRIMARY KEY AUTOINCREMENT,
3914 stage_token TEXT NOT NULL,
3915 op_index INTEGER NOT NULL,
3916 op_kind TEXT NOT NULL,
3917 memory_id TEXT,
3918 detail TEXT NOT NULL,
3919 applied_at_ms INTEGER NOT NULL
3920 );",
3921 )?;
3922 let provenance = "{\"author\":{\"author\":\"application\"}}";
3923 let insert = |id: &str, status_kind: &str, detail: Option<&str>| {
3924 conn.execute(
3925 "INSERT INTO records (memory_id, scope_kind, scope_key, kind, title, \
3926 description, body, tags, provenance, trust, status_kind, status_detail, \
3927 supersedes, derived_from, content_hash, created_at_ms, updated_at_ms, \
3928 usage_stats) VALUES (?1, 'identity', 'identity:luka', 'fact', ?1, '', \
3929 'body', '[]', ?2, 'agent_observed', ?3, ?4, NULL, '[]', ?1, 1, 1, '{}')",
3930 params![id, provenance, status_kind, detail],
3931 )
3932 };
3933 insert("mem-clean", "active", None)?;
3934 insert("mem-quarantined", "quarantined", Some("tainted session"))?;
3935 insert("mem-tombstoned-was-quarantined", "tombstoned", None)?;
3936 insert("mem-tombstoned-clean", "tombstoned", None)?;
3937 conn.execute(
3938 "INSERT INTO audit (stage_token, op_index, op_kind, memory_id, detail, \
3939 applied_at_ms) VALUES ('direct-1', 0, 'create', \
3940 'mem-tombstoned-was-quarantined', \
3941 '{\"op\":\"create\",\"quarantined\":\"llm_writes=quarantined policy\"}', 1)",
3942 params![],
3943 )?;
3944 }
3945 let store = SqliteAgentMemoryStore::open(dir.path())?;
3946 assert!(store.pending_proposals("family", 4).await?.is_empty());
3949 let conn = store.realm_connection("family")?;
3950 let guard = conn
3951 .lock()
3952 .unwrap_or_else(std::sync::PoisonError::into_inner);
3953 let flag = |id: &str| -> Result<bool, Box<dyn Error>> {
3954 Ok(guard.query_row(
3955 "SELECT ever_quarantined FROM records WHERE memory_id = ?1",
3956 params![id],
3957 |row| row.get(0),
3958 )?)
3959 };
3960 assert!(!flag("mem-clean")?);
3961 assert!(flag("mem-quarantined")?);
3962 assert!(flag("mem-tombstoned-was-quarantined")?);
3963 assert!(
3964 !flag("mem-tombstoned-clean")?,
3965 "ordinary tombstones must not be poisoned by the backfill"
3966 );
3967 Ok(())
3968 }
3969
3970 #[tokio::test]
3977 async fn proposal_taint_migration_marks_live_proposals_tainted() -> Result<(), Box<dyn Error>> {
3978 let dir = tempfile::tempdir()?;
3979 let db_path = {
3980 let store = SqliteAgentMemoryStore::open(dir.path())?;
3981 store.path_for_realm("family")
3982 };
3983 {
3984 let conn = Connection::open(&db_path)?;
3985 conn.execute_batch(
3986 "CREATE TABLE proposals (
3987 proposal_id TEXT PRIMARY KEY,
3988 scope_kind TEXT NOT NULL,
3989 scope_key TEXT NOT NULL,
3990 record TEXT NOT NULL,
3991 author TEXT NOT NULL,
3992 status TEXT NOT NULL DEFAULT 'pending',
3993 created_at_ms INTEGER NOT NULL
3994 );",
3995 )?;
3996 let record = serde_json::to_string(&NewMemoryRecord {
3997 kind: MemoryKind::Fact,
3998 title: "Shared gotcha".to_string(),
3999 description: String::new(),
4000 body: "proposed before the taint column existed".to_string(),
4001 tags: Vec::new(),
4002 evidence: Vec::new(),
4003 verification: None,
4004 })?;
4005 let author = serde_json::to_string(&MemoryAuthor::Agent {
4006 identity: "identity:luka".to_string(),
4007 })?;
4008 let insert = |id: &str, status: &str| {
4009 conn.execute(
4010 "INSERT INTO proposals (proposal_id, scope_kind, scope_key, record, \
4011 author, status, created_at_ms) VALUES (?1, 'mob', 'mob:home', ?2, ?3, \
4012 ?4, 1)",
4013 params![id, record, author, status],
4014 )
4015 };
4016 insert("prop-pending", "pending")?;
4017 insert("prop-held", "held")?;
4018 insert("prop-accepted", "accepted")?;
4019 insert("prop-rejected", "rejected")?;
4020 }
4021 let store = SqliteAgentMemoryStore::open(dir.path())?;
4022 let proposals = store.pending_proposals("family", 8).await?;
4023 assert_eq!(proposals.len(), 2, "{proposals:?}");
4024 for proposal in &proposals {
4025 let taint = proposal.taint.as_deref().unwrap_or_else(|| {
4026 panic!(
4027 "live pre-migration proposal '{}' must be conservatively tainted",
4028 proposal.proposal_id
4029 )
4030 });
4031 assert!(taint.contains("pre-migration"), "{taint}");
4032 }
4033 let conn = store.realm_connection("family")?;
4036 let guard = conn
4037 .lock()
4038 .unwrap_or_else(std::sync::PoisonError::into_inner);
4039 for id in ["prop-accepted", "prop-rejected"] {
4040 let taint: Option<String> = guard.query_row(
4041 "SELECT taint FROM proposals WHERE proposal_id = ?1",
4042 params![id],
4043 |row| row.get(0),
4044 )?;
4045 assert!(
4046 taint.is_none(),
4047 "terminal proposal '{id}' must stay untouched: {taint:?}"
4048 );
4049 }
4050 Ok(())
4051 }
4052
4053 #[tokio::test]
4054 async fn stale_stage_tokens_gc_on_open() -> Result<(), Box<dyn Error>> {
4055 let dir = tempfile::tempdir()?;
4056 let store = SqliteAgentMemoryStore::open(dir.path())?;
4057 let stage_create = |title: &str| StagedMutationBatch {
4058 kind: StagedBatchKind::FreshWrite,
4059 realm: "family".to_string(),
4060 author: MemoryAuthor::Application,
4061 ops: vec![StagedOp::Create {
4062 id: None,
4063 scope: identity_scope("family").expect("scope"),
4064 record: payload(title, &format!("{title} body")),
4065 trust: TrustTier::AgentObserved,
4066 derived_from: Vec::new(),
4067 rationale: None,
4068 created_at_ms: None,
4069 updated_at_ms: None,
4070 }],
4071 };
4072 let ungated = store.stage(stage_create("Stale")).await?;
4073 let pending_gated = store.stage(stage_create("Gated pending")).await?;
4076 store
4077 .record_pending_promotion(
4078 "family",
4079 PendingPromotion {
4080 pending_id: "gate-pending".to_string(),
4081 stage_token: pending_gated.token.clone(),
4082 record_id: "mem-src-1".to_string(),
4083 scope_kind: "mob".to_string(),
4084 scope_key: "mob:home".to_string(),
4085 rationale: None,
4086 status: "pending".to_string(),
4087 created_at_ms: now_ms(),
4088 },
4089 )
4090 .await?;
4091 let resolved_gated = store.stage(stage_create("Gated resolved")).await?;
4094 store
4095 .record_pending_promotion(
4096 "family",
4097 PendingPromotion {
4098 pending_id: "gate-resolved".to_string(),
4099 stage_token: resolved_gated.token.clone(),
4100 record_id: "mem-src-2".to_string(),
4101 scope_kind: "mob".to_string(),
4102 scope_key: "mob:home".to_string(),
4103 rationale: None,
4104 status: "pending".to_string(),
4105 created_at_ms: now_ms(),
4106 },
4107 )
4108 .await?;
4109 store
4110 .resolve_pending_promotion("family", "gate-resolved", "denied")
4111 .await?;
4112 {
4114 let conn = store.realm_connection("family")?;
4115 let guard = conn
4116 .lock()
4117 .unwrap_or_else(std::sync::PoisonError::into_inner);
4118 guard.execute(
4119 "UPDATE stage SET created_at_ms = created_at_ms - ?1",
4120 params![(STAGE_GC_MAX_AGE_MS + 60_000) as i64],
4121 )?;
4122 }
4123 let reopened = SqliteAgentMemoryStore::open(dir.path())?;
4124 let result = reopened.commit(ungated).await;
4125 assert!(
4126 matches!(result, Err(AgentMemoryError::InvalidRecord(_))),
4127 "aged-out ungated stage token must be garbage-collected on open"
4128 );
4129 let result = reopened.commit(resolved_gated).await;
4130 assert!(
4131 matches!(result, Err(AgentMemoryError::InvalidRecord(_))),
4132 "a stage referenced only by a RESOLVED promotion must still be collected"
4133 );
4134 let receipt = reopened.commit(pending_gated).await.map_err(|err| {
4135 format!("a stage referenced by a pending gated promotion must survive GC: {err}")
4136 })?;
4137 assert_eq!(receipt.applied_ops, 1);
4138 Ok(())
4139 }
4140
4141 #[tokio::test]
4142 async fn stage_rejects_lattice_violations() -> Result<(), Box<dyn Error>> {
4143 let dir = tempfile::tempdir()?;
4144 let store = SqliteAgentMemoryStore::open(dir.path())?;
4145 let scope = identity_scope("family")?;
4146
4147 let above_ceiling = store
4149 .stage(StagedMutationBatch {
4150 kind: StagedBatchKind::FreshWrite,
4151 realm: "family".to_string(),
4152 author: MemoryAuthor::Agent {
4153 identity: identity()?.as_str().to_string(),
4154 },
4155 ops: vec![StagedOp::Create {
4156 id: None,
4157 scope: scope.clone(),
4158 record: payload("Fact", "Body"),
4159 trust: TrustTier::AgentVerified,
4160 derived_from: Vec::new(),
4161 rationale: None,
4162 created_at_ms: None,
4163 updated_at_ms: None,
4164 }],
4165 })
4166 .await;
4167 assert!(matches!(
4168 above_ceiling,
4169 Err(AgentMemoryError::InvalidRecord(_))
4170 ));
4171
4172 let operator_tier = store
4174 .stage(StagedMutationBatch {
4175 kind: StagedBatchKind::FreshWrite,
4176 realm: "family".to_string(),
4177 author: MemoryAuthor::Operator,
4178 ops: vec![StagedOp::Create {
4179 id: None,
4180 scope,
4181 record: payload("Fact", "Body"),
4182 trust: TrustTier::Operator,
4183 derived_from: Vec::new(),
4184 rationale: None,
4185 created_at_ms: None,
4186 updated_at_ms: None,
4187 }],
4188 })
4189 .await;
4190 assert!(matches!(
4191 operator_tier,
4192 Err(AgentMemoryError::InvalidRecord(_))
4193 ));
4194 Ok(())
4195 }
4196
4197 #[tokio::test]
4198 async fn transitive_taint_blocks_laundering_through_store() -> Result<(), Box<dyn Error>> {
4199 let dir = tempfile::tempdir()?;
4200 let store = SqliteAgentMemoryStore::open(dir.path())?;
4201 let scope = identity_scope("family")?;
4202
4203 let seed = store
4206 .stage(StagedMutationBatch {
4207 kind: StagedBatchKind::FreshWrite,
4208 realm: "family".to_string(),
4209 author: MemoryAuthor::Steward {
4210 run_id: "dream-1".to_string(),
4211 },
4212 ops: vec![
4213 StagedOp::Create {
4214 id: Some("mem-tainted".to_string()),
4215 scope: scope.clone(),
4216 record: payload("Web claim", "Untrusted web content"),
4217 trust: TrustTier::Untrusted,
4218 derived_from: Vec::new(),
4219 rationale: None,
4220 created_at_ms: None,
4221 updated_at_ms: None,
4222 },
4223 StagedOp::Create {
4224 id: Some("mem-merged".to_string()),
4225 scope: scope.clone(),
4226 record: {
4227 let mut merged = payload("Consolidated", "Merged content");
4228 merged.verification = Some(super::super::records::VerificationClaim {
4229 checked: "claims verification".to_string(),
4230 evidence: Vec::new(),
4231 });
4232 merged
4233 },
4234 trust: TrustTier::AgentObserved,
4235 derived_from: vec!["mem-tainted".to_string()],
4236 rationale: Some("consolidation".to_string()),
4237 created_at_ms: None,
4238 updated_at_ms: None,
4239 },
4240 ],
4241 })
4242 .await?;
4243 store.commit(seed).await?;
4244
4245 let launder = store
4246 .stage(StagedMutationBatch {
4247 kind: StagedBatchKind::FreshWrite,
4248 realm: "family".to_string(),
4249 author: MemoryAuthor::Steward {
4250 run_id: "dream-2".to_string(),
4251 },
4252 ops: vec![StagedOp::Retier {
4253 id: "mem-merged".to_string(),
4254 trust: TrustTier::AgentVerified,
4255 rationale: Some("launder attempt".to_string()),
4256 }],
4257 })
4258 .await;
4259 let err = match launder {
4260 Err(AgentMemoryError::InvalidRecord(message)) => message,
4261 other => return Err(format!("laundering must be rejected, got {other:?}").into()),
4262 };
4263 assert!(err.contains("untrusted/quarantined"), "{err}");
4264 Ok(())
4265 }
4266
4267 #[tokio::test]
4268 async fn markdown_import_preserves_ids_and_renames_file() -> Result<(), Box<dyn Error>> {
4269 let dir = tempfile::tempdir()?;
4270 let id = identity()?;
4271 let markdown = MarkdownAgentMemoryStore::open(dir.path())?;
4272 let first = markdown.remember(
4273 "family",
4274 &id,
4275 new_memory("Imported fact", "Body one with detail."),
4276 )?;
4277 let second = markdown.remember(
4278 "family",
4279 &id,
4280 NewAgentMemory {
4281 title: "Second fact".to_string(),
4282 body: "Body two with detail.".to_string(),
4283 tags: vec!["travel".to_string()],
4284 },
4285 )?;
4286 let md_path = markdown.path_for("family", &id);
4287
4288 let store = SqliteAgentMemoryStore::open(dir.path())?;
4289 let records = store.recall(recall_all(id.clone(), "family")).await?;
4290 let mut got: Vec<&str> = records.iter().map(|r| r.memory_id.as_str()).collect();
4291 got.sort_unstable();
4292 let mut want = [first.memory_id.as_str(), second.memory_id.as_str()];
4293 want.sort_unstable();
4294 assert_eq!(got, want, "import must preserve memory ids");
4295 let imported = records
4296 .iter()
4297 .find(|record| record.memory_id == second.memory_id)
4298 .ok_or("second record imported")?;
4299 assert_eq!(imported.tags, vec!["travel"]);
4300 assert_eq!(imported.created_at_ms, second.created_at_ms);
4301
4302 assert!(
4303 !md_path.exists(),
4304 "markdown file must be renamed after import"
4305 );
4306 let renamed = md_path.with_extension("md.imported");
4307 assert!(renamed.exists(), "markdown file must survive as .imported");
4308
4309 let conn = store.realm_connection("family")?;
4311 let guard = conn
4312 .lock()
4313 .unwrap_or_else(std::sync::PoisonError::into_inner);
4314 let audits: i64 = guard.query_row(
4315 "SELECT COUNT(*) FROM audit WHERE stage_token LIKE 'import-%'",
4316 [],
4317 |row| row.get(0),
4318 )?;
4319 assert_eq!(audits, 2);
4320 drop(guard);
4321
4322 let reopened = SqliteAgentMemoryStore::open(dir.path())?;
4324 assert_eq!(reopened.recall(recall_all(id, "family")).await?.len(), 2);
4325 Ok(())
4326 }
4327
4328 struct TaintedSessionGate;
4331
4332 impl crate::memory::taint::LlmWriteGate for TaintedSessionGate {
4333 fn quarantine_reason(
4334 &self,
4335 author: &MemoryAuthor,
4336 _kind: StagedBatchKind,
4337 evidence: &[crate::memory::records::EvidenceRef],
4338 ) -> Option<String> {
4339 if !author.is_llm() {
4340 return None;
4341 }
4342 evidence
4343 .iter()
4344 .any(|reference| reference.session_id == "tainted-sess")
4345 .then(|| "evidence cites a tainted session".to_string())
4346 }
4347 }
4348
4349 fn tainted_evidence() -> Vec<crate::memory::records::EvidenceRef> {
4350 vec![crate::memory::records::EvidenceRef {
4351 session_id: "tainted-sess".to_string(),
4352 generation: 0,
4353 revision: None,
4354 range: None,
4355 }]
4356 }
4357
4358 #[tokio::test]
4359 async fn release_then_retier_of_formerly_quarantined_origin_rejected()
4360 -> Result<(), Box<dyn Error>> {
4361 let dir = tempfile::tempdir()?;
4362 let store = SqliteAgentMemoryStore::open(dir.path())?;
4363 store.set_llm_write_gate(std::sync::Arc::new(TaintedSessionGate));
4364 let scope = identity_scope("family")?;
4365
4366 let mut record = payload("Quarantined origin", "possibly poisoned content");
4370 record.evidence = tainted_evidence();
4371 record.verification = Some(crate::memory::records::VerificationClaim {
4372 checked: "claims to have checked".to_string(),
4373 evidence: Vec::new(),
4374 });
4375 let receipt = store
4376 .remember_authored(
4377 &scope,
4378 record,
4379 MemoryAuthor::Agent {
4380 identity: identity()?.as_str().to_string(),
4381 },
4382 )
4383 .await?;
4384 assert!(matches!(receipt.status, RecordStatus::Quarantined { .. }));
4385 let origin = receipt.memory_id;
4386
4387 let mut copy_payload =
4390 payload("Quarantined origin", "possibly poisoned content (released)");
4391 copy_payload.verification = Some(crate::memory::records::VerificationClaim {
4392 checked: "claims to have checked".to_string(),
4393 evidence: Vec::new(),
4394 });
4395 let release = StagedMutationBatch {
4396 kind: StagedBatchKind::FreshWrite,
4397 realm: "family".to_string(),
4398 author: MemoryAuthor::Steward {
4399 run_id: "dream-1".to_string(),
4400 },
4401 ops: vec![
4402 StagedOp::Create {
4403 id: Some("mem-released-copy".to_string()),
4404 scope: scope.clone(),
4405 record: copy_payload,
4406 trust: TrustTier::AgentObserved,
4407 derived_from: vec![origin.clone()],
4408 rationale: Some("quarantine release".to_string()),
4409 created_at_ms: None,
4410 updated_at_ms: None,
4411 },
4412 StagedOp::Tombstone {
4413 id: origin.clone(),
4414 rationale: Some("superseded by quarantine release".to_string()),
4415 },
4416 ],
4417 };
4418 let token = store.stage(release).await?;
4419 store.commit(token).await?;
4420 let released = store
4421 .record_by_id("family", "mem-released-copy")
4422 .await?
4423 .ok_or("released copy exists")?;
4424 assert_eq!(released.status, RecordStatus::Active);
4425 let origin_record = store
4426 .record_by_id("family", &origin)
4427 .await?
4428 .ok_or("origin exists")?;
4429 assert_eq!(origin_record.status, RecordStatus::Tombstoned);
4430
4431 {
4434 let conn = store.realm_connection("family")?;
4435 let guard = conn
4436 .lock()
4437 .unwrap_or_else(std::sync::PoisonError::into_inner);
4438 let flags: Vec<(String, bool)> = {
4439 let mut stmt = guard.prepare(
4440 "SELECT memory_id, ever_quarantined FROM records ORDER BY memory_id",
4441 )?;
4442 let rows = stmt
4443 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))?
4444 .collect::<Result<Vec<_>, _>>()?;
4445 rows
4446 };
4447 for (memory_id, flag) in &flags {
4448 assert!(
4449 flag,
4450 "'{memory_id}' must carry ever_quarantined after the release"
4451 );
4452 }
4453 }
4454
4455 let retier = StagedMutationBatch {
4459 kind: StagedBatchKind::FreshWrite,
4460 realm: "family".to_string(),
4461 author: MemoryAuthor::Steward {
4462 run_id: "dream-2".to_string(),
4463 },
4464 ops: vec![StagedOp::Retier {
4465 id: "mem-released-copy".to_string(),
4466 trust: TrustTier::AgentVerified,
4467 rationale: Some("post-release launder attempt".to_string()),
4468 }],
4469 };
4470 let err = store.stage(retier).await.expect_err("ceiling must hold");
4471 assert!(
4472 err.to_string().contains("provenance chain reaches"),
4473 "{err}"
4474 );
4475 Ok(())
4476 }
4477
4478 #[tokio::test]
4479 async fn markdown_import_skips_bad_records_and_files_loudly() -> Result<(), Box<dyn Error>> {
4480 let dir = tempfile::tempdir()?;
4481 let id = identity()?;
4482 let markdown = MarkdownAgentMemoryStore::open(dir.path())?;
4483 let valid = markdown.remember("family", &id, new_memory("Valid fact", "Valid body."))?;
4484 let md_path = markdown.path_for("family", &id);
4485
4486 let oversized_title = "T".repeat(crate::memory::records::MAX_RECORD_TITLE_BYTES + 10);
4490 let mut content = fs::read_to_string(&md_path)?;
4491 content.push_str(&format!(
4492 "## {oversized_title}\n<!-- mobkit-agent-memory \
4493 {{\"memory_id\":\"mem-bad-title\",\"tags\":[],\"created_at_ms\":1,\
4494 \"updated_at_ms\":1}} -->\nSome body.\n<!-- /mobkit-agent-memory -->\n\n"
4495 ));
4496 content.push_str(
4497 "## Leaked credential\n<!-- mobkit-agent-memory \
4498 {\"memory_id\":\"mem-secret\",\"tags\":[],\"created_at_ms\":2,\
4499 \"updated_at_ms\":2} -->\nthe key was AKIAIOSFODNN7EXAMPLE\n\
4500 <!-- /mobkit-agent-memory -->\n\n",
4501 );
4502 fs::write(&md_path, content)?;
4503
4504 let junk_path = dir.path().join("family").join("not an identity.md");
4508 fs::write(&junk_path, "## Orphan\nnot a memory file\n")?;
4509
4510 let store = SqliteAgentMemoryStore::open(dir.path())?;
4511 let records = store.recall(recall_all(id.clone(), "family")).await?;
4512 assert_eq!(
4513 records
4514 .iter()
4515 .map(|record| record.memory_id.as_str())
4516 .collect::<Vec<_>>(),
4517 vec![valid.memory_id.as_str()],
4518 "only the valid record imports"
4519 );
4520 assert!(!md_path.exists(), "identity file renamed after import");
4521 assert!(md_path.with_extension("md.imported").exists());
4522 assert!(!junk_path.exists(), "junk file set aside");
4523 assert!(junk_path.with_extension("md.import-failed").exists());
4524
4525 let conn = store.realm_connection("family")?;
4527 let guard = conn
4528 .lock()
4529 .unwrap_or_else(std::sync::PoisonError::into_inner);
4530 let summaries: Vec<String> = {
4531 let mut stmt =
4532 guard.prepare("SELECT detail FROM audit WHERE op_kind = 'import_summary'")?;
4533 let rows = stmt
4534 .query_map([], |row| row.get(0))?
4535 .collect::<Result<Vec<_>, _>>()?;
4536 rows
4537 };
4538 assert_eq!(summaries.len(), 2, "one summary per skipping/failing file");
4539 let identity_summary = summaries
4540 .iter()
4541 .find(|detail| detail.contains("mem-bad-title"))
4542 .ok_or("identity-file summary present")?;
4543 assert!(
4544 identity_summary.contains("\"skipped\":2"),
4545 "{identity_summary}"
4546 );
4547 assert!(
4548 identity_summary.contains("secret pattern class"),
4549 "{identity_summary}"
4550 );
4551 assert!(
4552 !identity_summary.contains("AKIAIOSFODNN7EXAMPLE"),
4553 "audit must not echo the secret: {identity_summary}"
4554 );
4555 drop(guard);
4556
4557 let reopened = SqliteAgentMemoryStore::open(dir.path())?;
4559 assert_eq!(reopened.recall(recall_all(id, "family")).await?.len(), 1);
4560 Ok(())
4561 }
4562
4563 #[tokio::test]
4564 async fn propose_captures_taint_at_propose_time_and_refuses_secrets()
4565 -> Result<(), Box<dyn Error>> {
4566 let dir = tempfile::tempdir()?;
4567 let store = SqliteAgentMemoryStore::open(dir.path())?;
4568 store.set_llm_write_gate(std::sync::Arc::new(TaintedSessionGate));
4569 let mob = MemoryScope::Mob {
4570 realm: "family".to_string(),
4571 mob: "mob:home".to_string(),
4572 };
4573 let author = MemoryAuthor::Agent {
4574 identity: identity()?.as_str().to_string(),
4575 };
4576
4577 let mut tainted = payload("Shared gotcha", "from a poisoned session");
4579 tainted.evidence = tainted_evidence();
4580 let tainted_id = store.propose(&mob, tainted, author.clone()).await?;
4581 let clean_id = store
4583 .propose(
4584 &mob,
4585 payload("Clean gotcha", "from a clean session"),
4586 author.clone(),
4587 )
4588 .await?;
4589 let proposals = store.pending_proposals("family", 8).await?;
4590 let by_id: std::collections::HashMap<&str, &PendingProposal> = proposals
4591 .iter()
4592 .map(|proposal| (proposal.proposal_id.as_str(), proposal))
4593 .collect();
4594 let tainted_row = by_id.get(tainted_id.as_str()).ok_or("tainted present")?;
4595 assert!(
4596 tainted_row
4597 .taint
4598 .as_deref()
4599 .is_some_and(|reason| reason.contains("tainted")),
4600 "{:?}",
4601 tainted_row.taint
4602 );
4603 assert!(
4604 by_id
4605 .get(clean_id.as_str())
4606 .ok_or("clean present")?
4607 .taint
4608 .is_none()
4609 );
4610
4611 let err = store
4613 .propose(
4614 &mob,
4615 payload("Creds", "api_key = \"zXy1aB2cD3eF4gH5iJ6k\""),
4616 author,
4617 )
4618 .await
4619 .expect_err("secret-bearing proposal refused");
4620 let message = err.to_string();
4621 assert!(message.contains("credential-assignment"), "{message}");
4622 assert!(!message.contains("zXy1aB2cD3eF4gH5iJ6k"), "{message}");
4623 Ok(())
4624 }
4625
4626 #[tokio::test]
4627 async fn secret_bearing_writes_refused_at_store_seam() -> Result<(), Box<dyn Error>> {
4628 let dir = tempfile::tempdir()?;
4629 let store = SqliteAgentMemoryStore::open(dir.path())?;
4630 let id = identity()?;
4631 let err = store
4634 .remember(
4635 "family",
4636 &id,
4637 new_memory("AWS key", "found AKIAIOSFODNN7EXAMPLE in the logs"),
4638 )
4639 .await
4640 .expect_err("secret-bearing remember refused");
4641 let message = err.to_string();
4642 assert!(message.contains("aws-access-key-id"), "{message}");
4643 assert!(!message.contains("AKIAIOSFODNN7EXAMPLE"), "{message}");
4644
4645 store
4647 .remember(
4648 "family",
4649 &id,
4650 new_memory(
4651 "Key location",
4652 "The AWS key lives in the vault, path infra/aws.",
4653 ),
4654 )
4655 .await?;
4656 Ok(())
4657 }
4658
4659 #[tokio::test]
4660 async fn scope_floors_warn_but_never_evict() -> Result<(), Box<dyn Error>> {
4661 let dir = tempfile::tempdir()?;
4662 let store = SqliteAgentMemoryStore::open(dir.path())?.with_scope_floors(2, usize::MAX);
4663 let id = identity()?;
4664 for i in 0..4 {
4665 store
4666 .remember(
4667 "family",
4668 &id,
4669 new_memory(&format!("Fact {i}"), &format!("Body {i}")),
4670 )
4671 .await?;
4672 }
4673 let records = store.recall(recall_all(id, "family")).await?;
4674 assert_eq!(
4675 records.len(),
4676 4,
4677 "floors warn the steward; deterministic code never evicts"
4678 );
4679 Ok(())
4680 }
4681
4682 #[test]
4683 fn floor_warning_fires_above_either_floor() {
4684 assert!(scope_floor_warning(5, 0, 4, 100).is_some());
4685 assert!(scope_floor_warning(0, 101, 4, 100).is_some());
4686 assert!(scope_floor_warning(4, 100, 4, 100).is_none());
4687 }
4688
4689 #[tokio::test]
4690 async fn mark_usage_updates_counters() -> Result<(), Box<dyn Error>> {
4691 let dir = tempfile::tempdir()?;
4692 let store = SqliteAgentMemoryStore::open(dir.path())?;
4693 let id = identity()?;
4694 let record = store
4695 .remember("family", &id, new_memory("Fact", "Body"))
4696 .await?;
4697
4698 store
4699 .mark_usage(&[record.memory_id.clone()], UsageEvent::Injected)
4700 .await?;
4701 store
4702 .mark_usage(&[record.memory_id.clone()], UsageEvent::ExplicitRecall)
4703 .await?;
4704 store
4705 .mark_usage(&[record.memory_id.clone()], UsageEvent::ExplicitRecall)
4706 .await?;
4707 store
4708 .mark_usage(&[record.memory_id.clone()], UsageEvent::JudgedUseful)
4709 .await?;
4710
4711 let conn = store.realm_connection("family")?;
4712 let guard = conn
4713 .lock()
4714 .unwrap_or_else(std::sync::PoisonError::into_inner);
4715 let usage_json: String = guard.query_row(
4716 "SELECT usage_stats FROM records WHERE memory_id = ?1",
4717 params![record.memory_id],
4718 |row| row.get(0),
4719 )?;
4720 let usage: UsageStats = serde_json::from_str(&usage_json)?;
4721 assert_eq!(usage.injected_count, 1, "ambient injections only");
4722 assert_eq!(usage.explicit_recall_count, 2, "explicit pulls only");
4723 assert_eq!(usage.judged_useful_count, 1);
4724 assert!(usage.last_injected_at_ms.is_some());
4725 assert!(usage.last_recalled_at_ms.is_some());
4726 Ok(())
4727 }
4728
4729 #[tokio::test]
4730 async fn injection_ledger_appends_and_reads_newest_first() -> Result<(), Box<dyn Error>> {
4731 let dir = tempfile::tempdir()?;
4732 let store = SqliteAgentMemoryStore::open(dir.path())?;
4733 let id = identity()?;
4734 let record = store
4735 .remember("family", &id, new_memory("Fact", "Body"))
4736 .await?;
4737
4738 let build_entry = InjectionLogEntry {
4739 record_id: record.memory_id.clone(),
4740 identity: id.as_str().to_string(),
4741 session_key: None,
4742 surface: InjectionSurface::Build,
4743 at_ms: 100,
4744 };
4745 let turn_entry = InjectionLogEntry {
4746 record_id: record.memory_id.clone(),
4747 identity: id.as_str().to_string(),
4748 session_key: Some("session-1".to_string()),
4749 surface: InjectionSurface::Turn,
4750 at_ms: 200,
4751 };
4752 AgentMemoryProvider::log_injections(&store, "family", &[build_entry.clone()]).await?;
4753 AgentMemoryProvider::log_injections(&store, "family", &[turn_entry.clone()]).await?;
4754
4755 let entries = store.injection_log("family", 16).await?;
4756 assert_eq!(entries, vec![turn_entry, build_entry]);
4757
4758 let limited = store.injection_log("family", 1).await?;
4759 assert_eq!(limited.len(), 1);
4760 assert_eq!(limited[0].surface, InjectionSurface::Turn);
4761
4762 let other_realm = store.injection_log("other", 16).await?;
4763 assert!(other_realm.is_empty(), "ledger rows are realm-scoped");
4764 Ok(())
4765 }
4766
4767 #[tokio::test]
4768 async fn propose_queues_for_steward() -> Result<(), Box<dyn Error>> {
4769 let dir = tempfile::tempdir()?;
4770 let store = SqliteAgentMemoryStore::open(dir.path())?;
4771 let scope = MemoryScope::Mob {
4772 realm: "family".to_string(),
4773 mob: "mob:home".to_string(),
4774 };
4775 let proposal_id = store
4776 .propose(
4777 &scope,
4778 payload("Shared fact", "For the mob store"),
4779 MemoryAuthor::Agent {
4780 identity: identity()?.as_str().to_string(),
4781 },
4782 )
4783 .await?;
4784 assert!(proposal_id.starts_with("prop-"));
4785
4786 let conn = store.realm_connection("family")?;
4787 let guard = conn
4788 .lock()
4789 .unwrap_or_else(std::sync::PoisonError::into_inner);
4790 let (status, scope_kind): (String, String) = guard.query_row(
4791 "SELECT status, scope_kind FROM proposals WHERE proposal_id = ?1",
4792 params![proposal_id],
4793 |row| Ok((row.get(0)?, row.get(1)?)),
4794 )?;
4795 assert_eq!(status, "pending");
4796 assert_eq!(scope_kind, "mob");
4797
4798 let author_json: String = guard.query_row(
4799 "SELECT author FROM proposals WHERE proposal_id = ?1",
4800 params![proposal_id],
4801 |row| row.get(0),
4802 )?;
4803 let author: MemoryAuthor = serde_json::from_str(&author_json)?;
4804 assert_eq!(
4805 author,
4806 MemoryAuthor::Agent {
4807 identity: identity()?.as_str().to_string()
4808 },
4809 "proposals carry real authorship (§8.2)"
4810 );
4811 Ok(())
4812 }
4813
4814 struct AlwaysQuarantine;
4819
4820 impl LlmWriteGate for AlwaysQuarantine {
4821 fn quarantine_reason(
4822 &self,
4823 author: &MemoryAuthor,
4824 _kind: StagedBatchKind,
4825 _evidence: &[crate::memory::records::EvidenceRef],
4826 ) -> Option<String> {
4827 author
4828 .is_llm()
4829 .then(|| "session tainted by web tool 'web_search'".to_string())
4830 }
4831 }
4832
4833 fn agent_author() -> Result<MemoryAuthor, Box<dyn Error>> {
4834 Ok(MemoryAuthor::Agent {
4835 identity: identity()?.as_str().to_string(),
4836 })
4837 }
4838
4839 #[tokio::test]
4840 async fn gated_agent_write_lands_quarantined_and_stays_unreadable() -> Result<(), Box<dyn Error>>
4841 {
4842 let dir = tempfile::tempdir()?;
4843 let store = SqliteAgentMemoryStore::open(dir.path())?;
4844 store.set_llm_write_gate(Arc::new(AlwaysQuarantine));
4845 let id = identity()?;
4846 let scope = identity_scope("family")?;
4847
4848 let receipt = store
4849 .remember_authored(
4850 &scope,
4851 payload("Poisoned", "Attacker fact"),
4852 agent_author()?,
4853 )
4854 .await?;
4855 let RecordStatus::Quarantined { reason } = &receipt.status else {
4856 return Err(format!("expected quarantined status, got {:?}", receipt.status).into());
4857 };
4858 assert!(reason.contains("session tainted"), "{reason}");
4859
4860 assert!(
4863 store
4864 .recall(recall_all(id.clone(), "family"))
4865 .await?
4866 .is_empty(),
4867 "quarantined bodies must never reach recall"
4868 );
4869 assert!(
4870 store
4871 .manifest(&[scope.clone()], ManifestTier::Full)
4872 .await?
4873 .is_empty(),
4874 "quarantined records must never reach the manifest"
4875 );
4876
4877 let record = store
4880 .remember("family", &id, new_memory("App fact", "App body"))
4881 .await?;
4882 let records = store.recall(recall_all(id, "family")).await?;
4883 assert_eq!(records.len(), 1);
4884 assert_eq!(records[0].memory_id, record.memory_id);
4885 Ok(())
4886 }
4887
4888 #[tokio::test]
4889 async fn distiller_write_law_holds_at_the_store_seam() -> Result<(), Box<dyn Error>> {
4890 use crate::memory::taint::{ContentTrustConfig, SessionTaintTracker, TaintLlmWriteGate};
4891
4892 let dir = tempfile::tempdir()?;
4893 let store = SqliteAgentMemoryStore::open(dir.path())?;
4894 let tracker = SessionTaintTracker::new(ContentTrustConfig::default());
4895 store.set_llm_write_gate(Arc::new(TaintLlmWriteGate::new(
4896 Some(tracker.clone()),
4897 crate::identity_first::agent_memory::AgentMemoryLlmWrites::Observed,
4898 )));
4899 let scope = identity_scope("family")?;
4900 let author = MemoryAuthor::Distiller {
4901 run_id: "run-1".to_string(),
4902 };
4903 let with_evidence = |title: &str, session: &str| NewMemoryRecord {
4904 evidence: vec![crate::memory::records::EvidenceRef {
4905 session_id: session.to_string(),
4906 generation: 1,
4907 revision: None,
4908 range: Some((0, 3)),
4909 }],
4910 ..payload(title, "Distilled body")
4911 };
4912
4913 let receipt = store
4915 .remember_authored(
4916 &scope,
4917 with_evidence("Clean fact", "sess-clean"),
4918 author.clone(),
4919 )
4920 .await?;
4921 assert_eq!(receipt.status, RecordStatus::Active);
4922 let record = store.with_realm_conn(&"family".to_string(), |conn| {
4923 load_record(conn, "family", &receipt.memory_id)?
4924 .ok_or_else(|| AgentMemoryError::Io("record missing".to_string()))
4925 })?;
4926 assert_eq!(record.trust, TrustTier::AgentObserved);
4927 assert!(matches!(
4928 record.provenance.author,
4929 MemoryAuthor::Distiller { .. }
4930 ));
4931 assert_eq!(record.provenance.evidence.len(), 1);
4932 assert_eq!(record.provenance.evidence[0].range, Some((0, 3)));
4933
4934 tracker.note_current_session("identity:someone", "sess-dirty");
4937 tracker.observe_agent_event(
4938 "identity:someone",
4939 &meerkat_core::event::AgentEvent::ToolResultReceived {
4940 id: "t".to_string(),
4941 name: "web_fetch".to_string(),
4942 content: vec![],
4943 is_error: false,
4944 },
4945 );
4946 let receipt = store
4947 .remember_authored(
4948 &scope,
4949 with_evidence("Tainted fact", "sess-dirty"),
4950 author.clone(),
4951 )
4952 .await?;
4953 let RecordStatus::Quarantined { reason } = &receipt.status else {
4954 return Err(format!("expected quarantine, got {:?}", receipt.status).into());
4955 };
4956 assert!(reason.contains("evidence session tainted"), "{reason}");
4957
4958 tracker.mark_reset_boundary("sess-reset");
4960 let receipt = store
4961 .remember_authored(&scope, with_evidence("Reset fact", "sess-reset"), author)
4962 .await?;
4963 let RecordStatus::Quarantined { reason } = &receipt.status else {
4964 return Err(format!("expected quarantine, got {:?}", receipt.status).into());
4965 };
4966 assert!(reason.contains("reset boundary"), "{reason}");
4967 Ok(())
4968 }
4969
4970 #[tokio::test]
4971 async fn recent_tombstones_lists_scope_tombstones_newest_first() -> Result<(), Box<dyn Error>> {
4972 use crate::memory::distiller::TombstoneSource;
4973
4974 let dir = tempfile::tempdir()?;
4975 let store = SqliteAgentMemoryStore::open(dir.path())?;
4976 let id = identity()?;
4977 let scope = identity_scope("family")?;
4978 let kept = store
4979 .remember("family", &id, new_memory("Kept fact", "Body"))
4980 .await?;
4981 let dropped = store
4982 .remember("family", &id, new_memory("Phone number", "Body 2"))
4983 .await?;
4984 store.forget("family", &id, &dropped.memory_id).await?;
4985
4986 let tombstones = store.recent_tombstones(&scope, 0, 10).await?;
4987 assert_eq!(tombstones.len(), 1);
4988 assert_eq!(tombstones[0].title, "Phone number");
4989 assert!(tombstones[0].tombstoned_at_ms > 0);
4990 assert!(!tombstones.iter().any(|t| t.title == "Kept fact"));
4992 let future = tombstones[0].tombstoned_at_ms + 1;
4993 assert!(
4994 store
4995 .recent_tombstones(&scope, future, 10)
4996 .await?
4997 .is_empty()
4998 );
4999 let _ = kept;
5000 Ok(())
5001 }
5002
5003 #[tokio::test]
5004 async fn quarantined_supersede_leaves_prior_active() -> Result<(), Box<dyn Error>> {
5005 let dir = tempfile::tempdir()?;
5006 let store = SqliteAgentMemoryStore::open(dir.path())?;
5007 let id = identity()?;
5008 let scope = identity_scope("family")?;
5009 let prior = store
5010 .remember("family", &id, new_memory("DB host", "Use db-good.example."))
5011 .await?;
5012
5013 store.set_llm_write_gate(Arc::new(AlwaysQuarantine));
5014 let receipt = store
5015 .supersede_authored(
5016 &scope,
5017 &prior.memory_id,
5018 payload("DB host", "Use db-evil.example."),
5019 agent_author()?,
5020 )
5021 .await?;
5022 assert!(matches!(receipt.status, RecordStatus::Quarantined { .. }));
5023
5024 let records = store.recall(recall_all(id, "family")).await?;
5026 assert_eq!(records.len(), 1);
5027 assert_eq!(records[0].memory_id, prior.memory_id);
5028 assert!(records[0].body.contains("db-good"));
5029 Ok(())
5030 }
5031
5032 #[tokio::test]
5033 async fn gate_covers_staged_commits_not_just_direct_writes() -> Result<(), Box<dyn Error>> {
5034 let dir = tempfile::tempdir()?;
5035 let store = SqliteAgentMemoryStore::open(dir.path())?;
5036 store.set_llm_write_gate(Arc::new(AlwaysQuarantine));
5037 let id = identity()?;
5038 let scope = identity_scope("family")?;
5039
5040 let token = store
5041 .stage(StagedMutationBatch {
5042 kind: StagedBatchKind::FreshWrite,
5043 realm: "family".to_string(),
5044 author: agent_author()?,
5045 ops: vec![StagedOp::Create {
5046 id: None,
5047 scope,
5048 record: payload("Staged fact", "Via staged path"),
5049 trust: TrustTier::AgentObserved,
5050 derived_from: Vec::new(),
5051 rationale: None,
5052 created_at_ms: None,
5053 updated_at_ms: None,
5054 }],
5055 })
5056 .await?;
5057 store.commit(token).await?;
5058 assert!(
5059 store.recall(recall_all(id, "family")).await?.is_empty(),
5060 "the write gate must hold at the store seam for staged commits too"
5061 );
5062 Ok(())
5063 }
5064
5065 #[tokio::test]
5066 async fn ungated_authored_write_lands_active_with_agent_author() -> Result<(), Box<dyn Error>> {
5067 let dir = tempfile::tempdir()?;
5068 let store = SqliteAgentMemoryStore::open(dir.path())?;
5069 let id = identity()?;
5070 let scope = identity_scope("family")?;
5071
5072 let mut record = payload("Observed fact", "Seen in session");
5073 record.verification = Some(super::super::records::VerificationClaim {
5074 checked: "ran the smoke test and watched it pass".to_string(),
5075 evidence: Vec::new(),
5076 });
5077 let receipt = store
5078 .remember_authored(&scope, record, agent_author()?)
5079 .await?;
5080 assert_eq!(receipt.status, RecordStatus::Active);
5081
5082 let conn = store.realm_connection("family")?;
5085 let guard = conn
5086 .lock()
5087 .unwrap_or_else(std::sync::PoisonError::into_inner);
5088 let (trust, provenance_json): (String, String) = guard.query_row(
5089 "SELECT trust, provenance FROM records WHERE memory_id = ?1",
5090 params![receipt.memory_id],
5091 |row| Ok((row.get(0)?, row.get(1)?)),
5092 )?;
5093 assert_eq!(trust, "agent_observed");
5094 let provenance: MemoryProvenance = serde_json::from_str(&provenance_json)?;
5095 assert_eq!(provenance.author, agent_author()?);
5096 assert!(
5097 provenance
5098 .verification
5099 .as_ref()
5100 .is_some_and(|claim| claim.checked.contains("smoke test"))
5101 );
5102 drop(guard);
5103
5104 let records = store.recall(recall_all(id, "family")).await?;
5106 assert_eq!(records.len(), 1);
5107
5108 let scope = identity_scope("family")?;
5110 let result = store
5111 .forget_authored(&scope, &receipt.memory_id, agent_author()?)
5112 .await?;
5113 assert!(result.deleted);
5114 Ok(())
5115 }
5116
5117 #[tokio::test]
5118 async fn authored_update_rejects_cross_identity_scope() -> Result<(), Box<dyn Error>> {
5119 let dir = tempfile::tempdir()?;
5120 let store = SqliteAgentMemoryStore::open(dir.path())?;
5121 let id = identity()?;
5122 let prior = store
5123 .remember("family", &id, new_memory("Fact", "Body"))
5124 .await?;
5125
5126 let other_scope = MemoryScope::Identity {
5131 realm: "family".to_string(),
5132 identity: "identity:other".to_string(),
5133 };
5134 let cross = store
5135 .supersede_authored(
5136 &other_scope,
5137 &prior.memory_id,
5138 payload("Fact", "Hijacked body"),
5139 MemoryAuthor::Agent {
5140 identity: "identity:other".to_string(),
5141 },
5142 )
5143 .await;
5144 assert!(
5145 matches!(cross, Err(AgentMemoryError::InvalidRecord(_))),
5146 "cross-identity update must be rejected, got {cross:?}"
5147 );
5148 Ok(())
5149 }
5150
5151 #[tokio::test]
5152 async fn panel_records_page_paginates_and_filters() -> Result<(), Box<dyn Error>> {
5153 let dir = tempfile::tempdir()?;
5154 let store = SqliteAgentMemoryStore::open(dir.path())?;
5155 let scope = identity_scope("family")?;
5156 for index in 0..5 {
5157 store
5158 .remember_authored(
5159 &scope,
5160 payload(&format!("Fact {index}"), &format!("Body {index}")),
5161 MemoryAuthor::Operator,
5162 )
5163 .await?;
5164 }
5165
5166 let first = store
5169 .records_page("family", Some("identity"), None, None, 2, None)
5170 .await?;
5171 assert_eq!(first.records.len(), 2);
5172 let cursor = first.next_cursor.clone().expect("more pages");
5173 let second = store
5174 .records_page("family", Some("identity"), None, None, 2, Some(cursor))
5175 .await?;
5176 assert_eq!(second.records.len(), 2);
5177 let third_cursor = second.next_cursor.clone().expect("one more page");
5178 let third = store
5179 .records_page(
5180 "family",
5181 Some("identity"),
5182 None,
5183 None,
5184 2,
5185 Some(third_cursor),
5186 )
5187 .await?;
5188 assert_eq!(third.records.len(), 1);
5189 assert_eq!(third.next_cursor, None);
5190 let mut seen: Vec<String> = first
5191 .records
5192 .iter()
5193 .chain(second.records.iter())
5194 .chain(third.records.iter())
5195 .map(|record| record.id.clone())
5196 .collect();
5197 let total = seen.len();
5198 seen.dedup();
5199 assert_eq!(total, 5, "pages cover every record exactly once");
5200
5201 let quarantined = store
5203 .records_page("family", None, None, Some("quarantined"), 10, None)
5204 .await?;
5205 assert!(quarantined.records.is_empty());
5206 Ok(())
5207 }
5208
5209 #[tokio::test]
5210 async fn panel_supersede_chain_walks_both_directions() -> Result<(), Box<dyn Error>> {
5211 let dir = tempfile::tempdir()?;
5212 let store = SqliteAgentMemoryStore::open(dir.path())?;
5213 let scope = identity_scope("family")?;
5214 let root = store
5215 .remember_authored(&scope, payload("Fact", "v1"), MemoryAuthor::Operator)
5216 .await?;
5217 let mid = store
5218 .supersede_authored(
5219 &scope,
5220 &root.memory_id,
5221 payload("Fact", "v2"),
5222 MemoryAuthor::Operator,
5223 )
5224 .await?;
5225 let tip = store
5226 .supersede_authored(
5227 &scope,
5228 &mid.memory_id,
5229 payload("Fact", "v3"),
5230 MemoryAuthor::Operator,
5231 )
5232 .await?;
5233
5234 for entry in [&root.memory_id, &mid.memory_id, &tip.memory_id] {
5236 let chain = store.supersede_chain("family", entry, 16).await?;
5237 let ids: Vec<&str> = chain.iter().map(|record| record.id.as_str()).collect();
5238 assert_eq!(
5239 ids,
5240 [
5241 root.memory_id.as_str(),
5242 mid.memory_id.as_str(),
5243 tip.memory_id.as_str()
5244 ],
5245 "chain from {entry}"
5246 );
5247 }
5248 let bounded = store.supersede_chain("family", &root.memory_id, 2).await?;
5250 assert_eq!(bounded.len(), 2);
5251 Ok(())
5252 }
5253}