Skip to main content

meerkat_mobkit/memory/
sqlite_store.rs

1//! Bundled per-realm SQLite agent-memory store (§7.3).
2//!
3//! One database per realm at `<root>/<pct-encoded-realm>.sqlite3` — the same
4//! directory and encoding scheme the markdown store uses (deliberately NOT
5//! `<persistent_state>/memory/`, which belongs to meerkat's session semantic
6//! memory). WAL journaling, busy-timeout, plain B-tree lookups only: the
7//! bright-line ratchet (§12) forbids retrieval-index machinery here, and
8//! recall quality is the LLM Selector's job, not the store's.
9//!
10//! Every write path — including `remember`/`forget` and the markdown import
11//! — flows through the staged-batch validator and a single-transaction
12//! apply with one audit row per op (§8.5 crash semantics).
13
14use 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
44/// Per-scope retention floors (§7.3): exceeded floors WARN the steward via
45/// tracing; deterministic code never evicts.
46pub const DEFAULT_SCOPE_FLOOR_RECORDS: usize = 4_000;
47pub const DEFAULT_SCOPE_FLOOR_BYTES: usize = 32 * 1024 * 1024;
48
49/// Staged-but-uncommitted batches older than this are garbage-collected on
50/// realm open — a dead producer leaves a token that is never applied.
51const 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/// Bundled SQLite store. Cheap to clone; connections are cached per realm
167/// and shared across clones.
168#[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    /// §10.1 write-seam enforcement: consulted for every LLM-authored
175    /// create/supersede across ALL write paths (direct and staged commits),
176    /// so taint/posture quarantine holds for any caller — the Recorder
177    /// tool, staged batches, and future stages alike. Shared across clones
178    /// so wiring the gate once covers every handle.
179    llm_write_gate: Arc<Mutex<Option<Arc<dyn LlmWriteGate>>>>,
180    /// §10.2 P3 extension: evidence-ref resolvability for `agent_verified`
181    /// retiers. Optional like the write gate — the wiring that enables the
182    /// steward installs it; absent, the P2 claim-presence rule stands
183    /// alone. Shared across clones.
184    evidence_resolver: Arc<Mutex<Option<Arc<dyn EvidenceRefResolver>>>>,
185    /// §9.3 timeline sink for quarantined-write events. Shared across
186    /// clones; absent, the tracing warn is the only surface.
187    event_sink: Arc<Mutex<Option<Arc<dyn crate::memory::events::MemoryEventSink>>>>,
188}
189
190/// §10.2 P3: whether an [`EvidenceRef`] resolves against the persistent
191/// session store (session exists; a cited range lies within the persisted
192/// transcript). The semantic endorsement half of an `agent_verified` retier
193/// is the dream's judgment (recorded in the op rationale); this is the
194/// mechanical half.
195pub 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    /// Install the §10.1 LLM write gate. Wiring installs it at startup,
220    /// before any member can dispatch a write.
221    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    /// Install the §10.1 gate only when none is present. The classic-mob
229    /// builder path uses this so it never clobbers a taint-tracking gate an
230    /// embedder (or the gateway) installed before handing the store over.
231    /// Returns whether this call installed the gate.
232    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    /// Install the §10.2 evidence-ref resolver. The steward wiring installs
245    /// it at startup; from then on every staged retier to `agent_verified`
246    /// must cite evidence that resolves against the session store.
247    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    /// Wire the §9.3 timeline sink for quarantined-write events.
255    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    /// Wire the §9.3 sink only when none is present — the classic-mob
263    /// builder path uses this so an embedder-installed sink survives.
264    /// Returns whether this call installed the sink.
265    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    /// Same directory + percent-encoding scheme as
309    /// `MarkdownAgentMemoryStore::path_for`, one database per realm.
310    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        // WAL survives crashes without blocking readers; query_row because
327        // the pragma returns the new mode.
328        conn.query_row("PRAGMA journal_mode=WAL", [], |_| Ok(()))
329            .map_err(sql_err)?;
330        conn.execute_batch(SCHEMA_SQL).map_err(sql_err)?;
331        // Column migrations for stores created before the columns joined
332        // SCHEMA_SQL (CREATE TABLE IF NOT EXISTS never alters).
333        if ensure_column(
334            &conn,
335            "records",
336            "ever_quarantined",
337            "INTEGER NOT NULL DEFAULT 0",
338        )? {
339            // Backfill the durable §10.2 marker: currently-quarantined rows
340            // directly; tombstoned rows through their audit trail (the
341            // tombstone apply nulls status_detail, so the audit row's
342            // `"quarantined":"<reason>"` is the only remaining evidence
343            // that a row once landed quarantined).
344            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            // Conservative backfill (mirrors ever_quarantined above): the
359            // propose-time taint fact for pre-migration proposals lived only
360            // in the in-memory SessionTaintTracker and is unrecoverable, so
361            // still-live proposals route through the operator-gated
362            // promotion path instead of reading as clean. Terminal statuses
363            // (accepted/rejected) are never re-verdicted and stay untouched.
364            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        // Stage GC spares tokens referenced by a still-pending gated
373        // promotion (§10.2) — the operator's decision window outranks the
374        // dead-producer sweep; deny/timeout resolution discards them.
375        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    /// One-shot migration (§7.3): un-imported markdown files for this realm
388    /// are imported through the staged-commit path (ids and timestamps
389    /// preserved; kind=fact, trust=agent_observed, identity scope, agent
390    /// author with empty evidence) and renamed to `<file>.imported` —
391    /// user-inspectable data is never deleted.
392    ///
393    /// §7.3 invites hand edits, so content problems must never make the
394    /// realm store unopenable: an invalid record is skipped loudly (warn +
395    /// count in the import audit row) and the rest of the file imports; a
396    /// file that fails wholesale (bad identity stem, over the size cap,
397    /// residual batch-validation failure) is warned about, set aside as
398    /// `<file>.import-failed`, and the remaining files continue. Only real
399    /// I/O errors propagate into the open.
400    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        // Skip ids already present (idempotence if a rename previously
464        // failed) and dedup ids within the file (hand-edits happen).
465        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            // Pre-validate each record with the same deterministic checks
485            // the staged validator applies, so one bad hand-edited record
486            // skips loudly instead of failing the whole batch.
487            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            // Gate deliberately absent: the import migrates records the
543            // markdown store already accepted; it is not a new LLM write.
544            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            // Deterministic write guard (§7.3): an exact content-hash
606            // duplicate short-circuits to the existing id — no new row.
607            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                // RPC/SDK writes are application-principal writes (§7.2);
627                // the P1 Recorder threads real agent authorship.
628                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    /// Shared tombstone path for the wire `forget` (Application principal)
687    /// and the Recorder's `forget_authored` (Agent principal).
688    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    /// §8.2 Recorder create: agent-authored, gate-enforced, dedup-guarded.
807    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            // Deterministic write guard (§7.3): an exact content-hash
827            // duplicate short-circuits to the existing active record.
828            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                    // §10.2: LLM writes enter at the ceiling; the staged
862                    // validator rejects anything higher.
863                    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                        // Counted apart from ambient injection (§9.2): a
935                        // pull on purpose is a much stronger usefulness
936                        // signal than a push that may have been ignored.
937                        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    /// Newest-first injection-ledger rows for a realm (§9.2). Read surface
1029    /// for the steward's usage audit and the console Memory panel.
1030    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        // §10.4 secret hygiene: proposals bypass the staged validator (the
1049        // row is not a record yet), so the write-seam refusal is applied
1050        // here directly.
1051        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        // §10.1: capture the quarantine decision AT PROPOSE TIME. The taint
1063        // tracker is in-memory and session-sticky; re-deriving when the
1064        // steward dreams would both under-quarantine (tracker restart,
1065        // reset boundary, eviction) and over-quarantine (identity tainted
1066        // later by an unrelated ingestion). The persisted fact makes the
1067        // steward's accept downgrade deterministic shell law.
1068        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// ---- provider trait implementations ----
1183
1184#[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// ---- steward read/write surface (§8.5) ----
1358
1359/// Per-scope store overview row for the dream's orient phase (§8.5) and
1360/// the P3b console Memory panel.
1361#[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/// One pending (or held) mob/operator-scope proposal awaiting a dream
1372/// verdict (§8.5 promotion).
1373#[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    /// §10.1 propose-time taint fact: `Some(reason)` when the write gate
1382    /// would have quarantined this author at propose time. A plain steward
1383    /// "accept" on a tainted proposal downgrades to an operator gate.
1384    pub taint: Option<String>,
1385}
1386
1387/// One retired identity awaiting an exit-interview harvest (§8.5).
1388#[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/// One gated quarantine-promotion (§10.2): the staged batch commits only on
1397/// gating approval.
1398#[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    /// The per-scope retention floors this store warns against (§7.3);
1412    /// rendered into the dream's orient overview as floor pressure.
1413    pub fn scope_floors(&self) -> (usize, usize) {
1414        (self.scope_floor_records, self.scope_floor_bytes)
1415    }
1416
1417    /// Per-scope counts and byte totals for a realm — the orient phase's
1418    /// one cheap aggregate.
1419    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    /// Pending/held proposals, oldest first (§8.5 promotion queue).
1479    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    /// Record a dream verdict on a proposal: `accepted`, `rejected`, or
1540    /// `held` (held stays in the pending queue for the next dream).
1541    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    /// Quarantined records, newest first — the dream's review queue (§8.5).
1576    /// The steward is the one stage that reads these bodies wholesale; the
1577    /// caller renders them defanged.
1578    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    /// Records by id, any status — the gather phase's bounded body fetch.
1608    /// Missing ids are skipped (the model may cite stale ids).
1609    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    /// Most recently updated records in a realm, any scope, active or
1632    /// quarantined — the gather phase filters (e.g. recent distillates by
1633    /// author) host-side.
1634    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    /// Record a retired identity for the next dream's exit-interview
1664    /// harvest (§8.5). Idempotent per (identity, retired_at_ms).
1665    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    /// Pending exit-interview harvests, oldest first.
1693    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    /// Mark one exit-interview harvest done.
1730    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    /// Record a gated quarantine-promotion: gating `pending_id` → staged
1754    /// batch token (§10.2). Only a gating approval commits the token.
1755    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    /// Look up a still-pending gated promotion by its gating pending id.
1787    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    /// All still-pending gated promotions (dream-start reconciliation).
1823    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    /// Resolve a gated promotion: `committed`, `denied`, or `expired`.
1863    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    /// Re-key a gated promotion after a gating escalation minted a
1893    /// successor pending entry.
1894    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    /// Discard a staged-but-uncommitted batch (denied/expired gated
1918    /// promotions; §8.5 crash semantics keep this safe — an unapplied stage
1919    /// row is never visible).
1920    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// ---- console Memory panel read surface (§9.3, P3b) ----
1934
1935/// One page of panel records: strictly-descending `(updated_at_ms,
1936/// memory_id)` keyset pagination.
1937#[derive(Debug, Clone, PartialEq, Eq)]
1938pub struct PanelRecordsPage {
1939    pub records: Vec<super::records::MemoryRecord>,
1940    /// Pass back as `cursor` to continue; `None` when exhausted.
1941    pub next_cursor: Option<(u64, String)>,
1942}
1943
1944/// One steward dream run reconstructed from its audit rows (every committed
1945/// op records its `Steward { run_id }` author in the audit `detail`).
1946#[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    /// op kind → count (create/supersede/tombstone/retier/set_rank).
1953    pub op_kinds: std::collections::BTreeMap<String, u64>,
1954    /// Ops that landed quarantined at the write seam.
1955    pub quarantined_ops: u64,
1956    /// Bounded sample of touched record ids, newest first.
1957    pub memory_ids: Vec<String>,
1958    /// Bounded sample of op rationales, newest first.
1959    pub rationales: Vec<String>,
1960}
1961
1962/// Bounds for [`SqliteAgentMemoryStore::dream_history`]: audit rows scanned
1963/// per call and per-run sample sizes. The panel is a summary surface, not a
1964/// full audit export.
1965const 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    /// Realms with a store file on disk (panel realm picker).
1971    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    /// One record by id, any status.
1977    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    /// Panel record listing: optional scope/status filters, newest-updated
1992    /// first, keyset cursor. Any status is visible here — the panel is an
1993    /// inspection surface and renders status explicitly.
1994    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    /// Supersede lineage around one record, oldest first: ancestors via the
2072    /// `supersedes` pointer, the record itself, then committed successors
2073    /// via the `Superseded { by }` status link. When the tip has no
2074    /// committed successor, records *claiming* to supersede it (e.g. a
2075    /// quarantined supersede that left the prior active, §10.1) are
2076    /// appended without recursing — claims are visible but never extend
2077    /// the walk. Bounded by `max_len`, cycle-safe.
2078    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                            // Trailing claimants: visible, not walked.
2131                            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    /// Newest-first injection-ledger rows for one record (panel usage view).
2160    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    /// Dream-run summaries reconstructed from steward audit rows, newest
2211    /// run first. Bounded scan ([`DREAM_HISTORY_SCAN_ROWS`]); runs older
2212    /// than the scan window fall off the panel, which is acceptable for a
2213    /// history summary surface.
2214    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    /// Recent tombstones for the Distiller's pre-injected "never re-create
2311    /// these" list (§8.4). The mechanical backstop for exact recreation is
2312    /// the staged validator's content-hash check; this list closes the
2313    /// paraphrase gap at the prompt level.
2314    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/// Body fetch for selector-chosen ids (§8.3): a plain by-id read over the
2364/// composed scopes, wire-compat projected, returned in `ids` order. Only
2365/// active records in the requested scopes qualify — the selector judged a
2366/// manifest of exactly those.
2367#[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                        // The full MemoryRecord is in hand before projection
2418                        // strips it — carry scope + trust so injected bodies
2419                        // render their §7.2 labels.
2420                        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
2438// ---- blocking internals ----
2439
2440async 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
2448/// §10.2 P3 validator extension, enforced at the store seam (stage and
2449/// commit): every `Retier` to `agent_verified` requires the target record's
2450/// verification claim to cite at least one `EvidenceRef` that resolves
2451/// against the session store. No resolver wired ⇒ the P2 claim-presence
2452/// rule stands alone (wiring that enables the steward installs one).
2453fn 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            // Unknown record — validate_batch already rejects this.
2484            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
2505/// Validates (against the live transaction) and applies a batch atomically:
2506/// one SQLite transaction, one audit row per op (§8.5).
2507///
2508/// `gate` is the §10.1 LLM write gate: consulted once per batch (the
2509/// quarantine decision is a property of the author's session/posture and of
2510/// the batch's cited evidence, not of individual ops — a batch with any
2511/// tainted evidence quarantines wholesale, conservative direction) and
2512/// applied to every create/supersede in the batch. `None` only for the
2513/// markdown import, which migrates already-accepted records rather than
2514/// writing new LLM output.
2515fn 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            // §8.3 / §7.1: the superseding record inherits the prior's rank
2664            // until the next dream re-ranks. rank_set_at_ms stays NULL so
2665            // the successor also remains in the manifest's recent slice —
2666            // a fresh correction is selector-visible on the next assembly.
2667            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                // A quarantined supersede must not retire the active prior:
2691                // otherwise a tainted session could silently blank a good
2692                // record by "updating" it. The quarantined successor keeps
2693                // its `supersedes` lineage edge; the steward resolves the
2694                // fork at review (promote → prior superseded; tombstone →
2695                // lineage unchanged).
2696                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            // Rank is steward metadata: updated_at_ms is deliberately NOT
2724            // bumped, or every re-rank would flood the manifest's
2725            // "updated since last rank" recent slice.
2726            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    // §10.1: the gate's verdict lands as row status. Quarantined records are
2761    // write-only — every read surface filters on status_kind = 'active'.
2762    let (status_kind, status_detail) = match quarantine {
2763        Some(reason) => ("quarantined", Some(reason)),
2764        None => ("active", None),
2765    };
2766    // §10.2 durable taint: set when landing quarantined, inherited from any
2767    // direct ancestor (derivation source or superseded prior) that carries
2768    // it or currently sits quarantined. Materialized transitively at each
2769    // insert, so one level suffices; the validator's chain walk remains the
2770    // enforcement.
2771    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
2813/// One-level ancestor check backing the materialized `ever_quarantined`
2814/// inheritance in [`insert_record`].
2815fn 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
2839/// Validator view over a live connection/transaction. Rows in a realm DB
2840/// are realm-homogeneous by construction, so the view carries the realm to
2841/// reconstruct full scopes for the validator's realm-confinement checks.
2842struct 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
2926// ---- row mapping ----
2927
2928struct 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
3074/// Wire-compat projection: MemoryRecord → AgentMemoryRecord keeps
3075/// memory_id/title/body/tags/timestamps (§7.3 — recall stays
3076/// wire-compatible).
3077fn 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
3088/// §8.3 WorkingSet(k): top-K ranked (steward ordering) ∪ recent/unranked
3089/// slice (unranked, or updated since their last rank), newest first, the
3090/// union capped at 2*k. Full: every active record, ranked first.
3091fn 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
3174/// §7.3 retention floors: warn (never evict) when a scope outgrows its
3175/// record-count or byte floor — retention pressure is a dream input, not a
3176/// FIFO.
3177fn 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
3209/// Pure floor check, unit-tested separately from the tracing side effect.
3210fn 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
3225/// Markdown-import failure split: content problems are contained (skip the
3226/// file, keep the store open); I/O problems propagate into the open.
3227enum MarkdownImportError {
3228    Content(String),
3229    Io(AgentMemoryError),
3230}
3231
3232/// One summary audit row per markdown-import file with skips or a wholesale
3233/// failure: the durable, operator-visible counterpart of the tracing warns.
3234fn 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
3262/// Add `column` to `table` when absent. Returns true when the column was
3263/// just added (the caller's cue to run a one-time backfill).
3264fn 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        // A deliberate non-LLM re-add of the same content passes the
3460        // tombstone-recreation guard (which targets LLM authors, §8.4) and
3461        // mints a fresh id.
3462        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        // Steward ranks the record, then the RPC update path supersedes it.
3480        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        // Only the successor is recallable (memory never argues with
3505        // itself), and it inherited the steward rank.
3506        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        // Chain is preserved on the row.
3521        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        // Rank the first three; ranking does not count as an update, so the
3553        // ranked records leave the recent/unranked slice.
3554        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        // top-2 ranked = ids[0], ids[1]; recent slice = the two unranked
3575        // (ids[4], ids[3] newest-first); union capped at 4.
3576        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        // A ranked record updated after its rank re-enters the recent slice
3588        // via supersede (rank inheritance keeps it selector-visible).
3589        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        // The producer "dies": no commit. Nothing is visible, in this
3630        // instance or a fresh one over the same directory.
3631        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        // Commit applies the batch and burns the token.
3646        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    /// Stores created before the `ever_quarantined`/`taint` columns existed
3655    /// migrate on open, with the conservative backfill: currently-quarantined
3656    /// rows flagged directly, tombstoned rows flagged through their audit
3657    /// trail (the tombstone apply nulled the `quarantined` status_detail).
3658    #[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        // Any realm operation opens the connection and runs the migration;
3735        // pending_proposals also exercises the proposals `taint` migration.
3736        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    /// The proposals `taint` migration conservatively marks still-live
3759    /// (pending/held) proposals tainted: the propose-time taint fact lived
3760    /// only in the in-memory tracker and is unrecoverable after the restart
3761    /// that accompanies the upgrade, so a plain steward accept downgrades
3762    /// to the operator gate instead of clean-accepting a possibly-tainted
3763    /// pre-migration proposal.
3764    #[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        // Terminal statuses are never re-verdicted: the backfill leaves them
3822        // alone.
3823        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        // A stage referenced by a still-PENDING gated promotion: the
3862        // operator's decision window outranks the dead-producer sweep.
3863        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        // A stage referenced by a RESOLVED promotion must NOT be exempt —
3880        // this pins the `status = 'pending'` filter in the GC query.
3881        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        // Age every stage row past the 24h GC horizon, then reopen.
3901        {
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        // Agent author above the LLM ceiling.
3936        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        // Operator tier is never staged-assignable, for any author.
3961        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        // Seed an untrusted record, merge it into a "fresh" consolidated
3992        // record, then try to retier the merge product upward.
3993        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        // Import audit trail exists (one audit row per imported record).
4098        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        // Reopening does not re-import (file renamed) and keeps counts.
4111        let reopened = SqliteAgentMemoryStore::open(dir.path())?;
4112        assert_eq!(reopened.recall(recall_all(id, "family")).await?.len(), 2);
4113        Ok(())
4114    }
4115
4116    /// Store-seam gate stand-in: quarantines LLM writes whose evidence
4117    /// cites the tainted session.
4118    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        // Agent write from a tainted session lands quarantined, carrying a
4155        // verification claim (so the later retier passes the claim check
4156        // and only the taint ceiling can stop it).
4157        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        // Steward release: create a copy derived from the origin, tombstone
4176        // the origin (exactly the dream's release group).
4177        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        // The durable taint marker persisted through the release: both the
4220        // tombstoned origin and the copy (inherited via derived_from).
4221        {
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        // §10.2 "capped forever": retiering the released copy to
4244        // agent_verified must be rejected even though the quarantined
4245        // origin is now tombstoned.
4246        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        // Hand-edits happen (§7.3 invites them): append one record with an
4275        // oversized title and one carrying a secret. Both must skip loudly;
4276        // the valid record must still import; the open must succeed.
4277        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        // A file whose stem is not an agent identity (whitespace never
4293        // validates) fails wholesale: set aside as .import-failed, never
4294        // taking the realm store down.
4295        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        // The skips are counted in import audit rows.
4314        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        // Reopen: no re-import attempts, store stays healthy.
4346        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        // Tainted at propose time: the fact is persisted on the row.
4366        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        // Clean propose: no taint.
4370        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        // §10.4: the proposal seam refuses secrets with the class named.
4400        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        // The wire remember path flows through the staged validator's
4420        // §10.4 chokepoint.
4421        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        // Clean writes pass.
4434        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    // ---- §10.1 write gate ----
4603
4604    /// Gate that quarantines every LLM-authored write (the
4605    /// `llm_writes = "quarantined"` posture / a permanently tainted session).
4606    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        // Quarantined records are write-only: recall and manifest (the
4649        // coordinator's two read surfaces) must never return them.
4650        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        // Non-LLM principals are not gated: the RPC remember path
4666        // (Application author) lands active through the same gate.
4667        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        // Clean evidence: lands Active, tier-ceilinged at AgentObserved.
4702        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        // Tainted evidence range: session-tainted ⇒ the write quarantines,
4723        // for the Distiller author (not just Agent authors).
4724        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        // Reset boundary: quarantines without any content taint (§8.4).
4747        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        // Active records never appear; a since_ms in the future filters out.
4779        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        // A tainted "update" must not blank the good record.
4813        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        // The verification is a CLAIM in provenance; the tier stays at the
4871        // LLM ceiling (§10.2).
4872        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        // Recall sees it (identity scope, active).
4893        let records = store.recall(recall_all(id, "family")).await?;
4894        assert_eq!(records.len(), 1);
4895
4896        // forget_authored tombstones it with agent authorship.
4897        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        // An agent may only supersede within its OWN identity scope: the
4915        // staged validator rejects the batch even when the caller lies
4916        // about the scope (single-lineage supersede stays with the record's
4917        // own writers, §8.2).
4918        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        // Keyset pagination: strictly-descending (updated_at_ms, id) with
4955        // no row repeated or skipped across pages.
4956        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        // Status filter.
4990        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        // The same chain comes back oldest-first from every entry point.
5023        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        // Bounded.
5037        let bounded = store.supersede_chain("family", &root.memory_id, 2).await?;
5038        assert_eq!(bounded.len(), 2);
5039        Ok(())
5040    }
5041}