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