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