zkr 0.2.2

Evidence-backed temporal memory for personal agents
Documentation
use super::*;
use rusqlite::{Transaction, TransactionBehavior};

const SCHEMA_VERSION: i64 = 8;
pub(super) const CLAIM_TIME_INTERVAL_ERROR: &str = "invalid claim half-open time interval";

pub(super) fn migrate(connection: &mut Connection) -> Result<()> {
    connection.execute_batch("PRAGMA foreign_keys = ON;")?;
    let version = connection.query_row("PRAGMA user_version", [], |row| row.get::<_, i64>(0))?;
    if !(0..=SCHEMA_VERSION).contains(&version) {
        return Err(Error::Invalid(format!(
            "database schema version {version} is newer than supported version {SCHEMA_VERSION}"
        )));
    }
    if version == SCHEMA_VERSION {
        return Ok(());
    }
    let transaction = connection.transaction_with_behavior(TransactionBehavior::Immediate)?;
    if version < 1 {
        migrate_v1(&transaction)?;
    }
    repair_v1_shape(&transaction)?;
    if version < 2 {
        migrate_v2(&transaction)?;
    }
    if version < 3 {
        migrate_v3(&transaction)?;
    }
    if version < 4 {
        migrate_v4(&transaction)?;
    }
    if version < 5 {
        migrate_v5(&transaction)?;
    }
    if version < 6 {
        migrate_v6(&transaction)?;
    }
    if version < 7 {
        migrate_v7(&transaction)?;
    }
    if version < 8 {
        migrate_v8(&transaction)?;
    }
    set_version(&transaction, SCHEMA_VERSION)?;
    transaction.commit()?;
    Ok(())
}

fn migrate_v1(transaction: &Transaction<'_>) -> Result<()> {
    transaction.execute_batch(
        "
        CREATE TABLE IF NOT EXISTS sources(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, revision INTEGER NOT NULL, kind TEXT NOT NULL, content TEXT NOT NULL, captured_at INTEGER NOT NULL, recorded_at INTEGER NOT NULL, deleted_at INTEGER);
        CREATE INDEX IF NOT EXISTS sources_scope ON sources(tenant_id, person_id, id);
        CREATE VIRTUAL TABLE IF NOT EXISTS source_fts USING fts5(source_id UNINDEXED, tenant_id UNINDEXED, person_id UNINDEXED, content, tokenize='unicode61');
        CREATE TABLE IF NOT EXISTS evidence(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, source_id TEXT NOT NULL REFERENCES sources(id), source_revision INTEGER NOT NULL, quote TEXT NOT NULL, recorded_at INTEGER NOT NULL, deleted_at INTEGER);
        CREATE INDEX IF NOT EXISTS evidence_scope ON evidence(tenant_id, person_id, source_id);
        CREATE TABLE IF NOT EXISTS claims(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, subject TEXT NOT NULL, predicate TEXT NOT NULL, value TEXT NOT NULL, valid_from INTEGER NOT NULL, valid_until INTEGER, recorded_from INTEGER NOT NULL, recorded_until INTEGER, status TEXT NOT NULL);
        CREATE INDEX IF NOT EXISTS claims_scope ON claims(tenant_id, person_id, status);
        CREATE TABLE IF NOT EXISTS claim_evidence(tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, claim_id TEXT NOT NULL REFERENCES claims(id), evidence_id TEXT NOT NULL REFERENCES evidence(id), relation TEXT NOT NULL, confidence_basis_points INTEGER NOT NULL, PRIMARY KEY(tenant_id, person_id, claim_id, evidence_id));
        CREATE TABLE IF NOT EXISTS daily_reviews(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, day TEXT NOT NULL, summary TEXT NOT NULL, evidence_ids TEXT NOT NULL, recorded_at INTEGER NOT NULL);
        CREATE INDEX IF NOT EXISTS reviews_scope ON daily_reviews(tenant_id, person_id, day);
        CREATE TABLE IF NOT EXISTS embeddings(tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, target_kind TEXT NOT NULL, target_id TEXT NOT NULL, model TEXT NOT NULL, version TEXT NOT NULL, dimension INTEGER NOT NULL, input_hash TEXT NOT NULL, normalization TEXT NOT NULL, distance TEXT NOT NULL, vector TEXT NOT NULL, PRIMARY KEY(tenant_id, person_id, target_kind, target_id, model, version));
        CREATE INDEX IF NOT EXISTS embeddings_scope ON embeddings(tenant_id, person_id, target_kind, target_id);
        ",
    )?;
    Ok(())
}

fn repair_v1_shape(transaction: &Transaction<'_>) -> Result<()> {
    ensure_column(transaction, "sources", "ingestion_key", "TEXT")?;
    ensure_column(
        transaction,
        "embeddings",
        "target_revision",
        "INTEGER NOT NULL DEFAULT 0",
    )?;
    ensure_column(
        transaction,
        "embeddings",
        "created_at",
        "INTEGER NOT NULL DEFAULT 0",
    )?;
    transaction.execute_batch(
        "CREATE UNIQUE INDEX IF NOT EXISTS sources_ingestion_key ON sources(tenant_id, person_id, ingestion_key) WHERE ingestion_key IS NOT NULL;",
    )?;
    Ok(())
}

fn migrate_v2(transaction: &Transaction<'_>) -> Result<()> {
    transaction.execute_batch(
        "
        CREATE TABLE IF NOT EXISTS evidence_locators(tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, evidence_id TEXT NOT NULL REFERENCES evidence(id), device_id TEXT NOT NULL, provider TEXT NOT NULL, stream_id TEXT NOT NULL, segment_id TEXT NOT NULL, start_ms INTEGER NOT NULL, end_ms INTEGER NOT NULL, PRIMARY KEY(tenant_id, person_id, evidence_id));
        ",
    )?;
    Ok(())
}

fn migrate_v3(transaction: &Transaction<'_>) -> Result<()> {
    transaction.execute_batch(
        "
        CREATE INDEX IF NOT EXISTS sources_scope ON sources(tenant_id, person_id, id);
        CREATE INDEX IF NOT EXISTS evidence_scope ON evidence(tenant_id, person_id, source_id);
        CREATE INDEX IF NOT EXISTS claims_scope ON claims(tenant_id, person_id, status);
        CREATE INDEX IF NOT EXISTS reviews_scope ON daily_reviews(tenant_id, person_id, day);
        CREATE INDEX IF NOT EXISTS embeddings_scope ON embeddings(tenant_id, person_id, target_kind, target_id);
        ",
    )?;
    Ok(())
}

fn migrate_v4(transaction: &Transaction<'_>) -> Result<()> {
    validate_scope(transaction, false, false)?;
    install_core_scope_triggers(transaction)?;
    Ok(())
}

fn migrate_v5(transaction: &Transaction<'_>) -> Result<()> {
    ensure_column(
        transaction,
        "claims",
        "kind",
        "TEXT NOT NULL DEFAULT 'fact'",
    )?;
    transaction.execute_batch(
        "
        CREATE TABLE IF NOT EXISTS profile_entries(id TEXT PRIMARY KEY, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, key TEXT NOT NULL, value TEXT NOT NULL, stability TEXT NOT NULL, claim_id TEXT NOT NULL REFERENCES claims(id), recorded_at INTEGER NOT NULL);
        CREATE INDEX IF NOT EXISTS profile_entries_scope ON profile_entries(tenant_id, person_id, recorded_at);
        ",
    )?;
    install_profile_triggers(transaction)?;
    Ok(())
}

fn migrate_v6(transaction: &Transaction<'_>) -> Result<()> {
    ensure_column(transaction, "sources", "origin_evidence_id", "TEXT")?;
    ensure_column(transaction, "sources", "origin_claim_id", "TEXT")?;
    transaction.execute(
        "UPDATE sources SET origin_evidence_id = (SELECT MIN(e.id) FROM evidence e WHERE e.source_id = sources.id AND e.tenant_id = sources.tenant_id AND e.person_id = sources.person_id) WHERE origin_evidence_id IS NULL AND (SELECT COUNT(*) FROM evidence e WHERE e.source_id = sources.id AND e.tenant_id = sources.tenant_id AND e.person_id = sources.person_id) = 1",
        [],
    )?;
    transaction.execute(
        r#"UPDATE sources SET origin_claim_id = (SELECT MIN(ce.claim_id) FROM claim_evidence ce WHERE ce.evidence_id = sources.origin_evidence_id AND ce.tenant_id = sources.tenant_id AND ce.person_id = sources.person_id AND ce.relation = '"supports"') WHERE origin_claim_id IS NULL AND origin_evidence_id IS NOT NULL AND (SELECT COUNT(*) FROM claim_evidence ce WHERE ce.evidence_id = sources.origin_evidence_id AND ce.tenant_id = sources.tenant_id AND ce.person_id = sources.person_id AND ce.relation = '"supports"') = 1"#,
        [],
    )?;
    validate_scope(transaction, false, true)?;
    transaction.execute_batch(
        "DROP TRIGGER IF EXISTS profile_entry_scope_insert;
         DROP TRIGGER IF EXISTS profile_entry_scope_update;",
    )?;
    validate_profile_references(transaction)?;
    transaction.execute(
        "UPDATE profile_entries SET key = (SELECT c.predicate FROM claims c WHERE c.id = profile_entries.claim_id AND c.tenant_id = profile_entries.tenant_id AND c.person_id = profile_entries.person_id AND c.kind = 'profile_fact'), value = (SELECT c.value FROM claims c WHERE c.id = profile_entries.claim_id AND c.tenant_id = profile_entries.tenant_id AND c.person_id = profile_entries.person_id AND c.kind = 'profile_fact') WHERE EXISTS (SELECT 1 FROM claims c WHERE c.id = profile_entries.claim_id AND c.tenant_id = profile_entries.tenant_id AND c.person_id = profile_entries.person_id AND c.kind = 'profile_fact')",
        [],
    )?;
    transaction.execute(
        "DELETE FROM profile_entries AS old WHERE EXISTS (SELECT 1 FROM profile_entries newer WHERE newer.tenant_id = old.tenant_id AND newer.person_id = old.person_id AND newer.key = old.key AND (newer.recorded_at > old.recorded_at OR (newer.recorded_at = old.recorded_at AND newer.id > old.id)))",
        [],
    )?;
    validate_scope(transaction, true, true)?;
    transaction.execute_batch(
        "
        CREATE UNIQUE INDEX IF NOT EXISTS profile_entries_current ON profile_entries(tenant_id, person_id, key);
        CREATE TRIGGER IF NOT EXISTS source_origin_scope_insert BEFORE INSERT ON sources FOR EACH ROW WHEN (NEW.origin_evidence_id IS NOT NULL AND NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.origin_evidence_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) OR (NEW.origin_claim_id IS NOT NULL AND NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.origin_claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) BEGIN SELECT RAISE(ABORT, 'source origin scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS source_origin_scope_update BEFORE UPDATE OF origin_evidence_id, origin_claim_id, tenant_id, person_id ON sources FOR EACH ROW WHEN (NEW.origin_evidence_id IS NOT NULL AND NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.origin_evidence_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) OR (NEW.origin_claim_id IS NOT NULL AND NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.origin_claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) BEGIN SELECT RAISE(ABORT, 'source origin scope mismatch'); END;
        ",
    )?;
    install_profile_triggers(transaction)?;
    Ok(())
}

fn migrate_v7(transaction: &Transaction<'_>) -> Result<()> {
    let invalid = transaction.query_row(
        "SELECT EXISTS(SELECT 1 FROM claims WHERE (valid_until IS NOT NULL AND valid_until <= valid_from) OR (recorded_until IS NOT NULL AND recorded_until <= recorded_from))",
        [],
        |row| row.get::<_, bool>(0),
    )?;
    if invalid {
        return Err(Error::Invalid(
            "legacy claim has an invalid half-open time interval".to_owned(),
        ));
    }
    transaction.execute_batch(&format!(
        "CREATE TRIGGER IF NOT EXISTS claim_time_interval_insert BEFORE INSERT ON claims FOR EACH ROW WHEN (NEW.valid_until IS NOT NULL AND NEW.valid_until <= NEW.valid_from) OR (NEW.recorded_until IS NOT NULL AND NEW.recorded_until <= NEW.recorded_from) BEGIN SELECT RAISE(ABORT, '{CLAIM_TIME_INTERVAL_ERROR}'); END;
         CREATE TRIGGER IF NOT EXISTS claim_time_interval_update BEFORE UPDATE OF valid_from, valid_until, recorded_from, recorded_until ON claims FOR EACH ROW WHEN (NEW.valid_until IS NOT NULL AND NEW.valid_until <= NEW.valid_from) OR (NEW.recorded_until IS NOT NULL AND NEW.recorded_until <= NEW.recorded_from) BEGIN SELECT RAISE(ABORT, '{CLAIM_TIME_INTERVAL_ERROR}'); END;",
    ))?;
    Ok(())
}

fn migrate_v8(transaction: &Transaction<'_>) -> Result<()> {
    transaction.execute_batch(
        "CREATE TABLE IF NOT EXISTS corrections(tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, superseded_claim_id TEXT NOT NULL REFERENCES claims(id), claim_id TEXT NOT NULL REFERENCES claims(id), source_id TEXT NOT NULL REFERENCES sources(id), evidence_id TEXT NOT NULL REFERENCES evidence(id), valid_at INTEGER NOT NULL, recorded_at INTEGER NOT NULL, PRIMARY KEY(tenant_id, person_id, superseded_claim_id, claim_id));
         CREATE TABLE IF NOT EXISTS memory_commits(sequence INTEGER PRIMARY KEY AUTOINCREMENT, tenant_id TEXT NOT NULL, person_id TEXT NOT NULL, recorded_at INTEGER NOT NULL);
         CREATE INDEX IF NOT EXISTS memory_commits_scope ON memory_commits(tenant_id, person_id, sequence);
         CREATE TABLE IF NOT EXISTS memory_export_events(commit_sequence INTEGER NOT NULL REFERENCES memory_commits(sequence) ON DELETE CASCADE, event_index INTEGER NOT NULL, payload TEXT NOT NULL CHECK(json_valid(payload)), PRIMARY KEY(commit_sequence, event_index));",
    )?;
    transaction.execute_batch(
        "CREATE TRIGGER IF NOT EXISTS correction_scope_insert BEFORE INSERT ON corrections FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.superseded_claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) OR NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) OR NOT EXISTS (SELECT 1 FROM sources WHERE id = NEW.source_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id AND origin_claim_id = NEW.claim_id AND origin_evidence_id = NEW.evidence_id) OR NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.evidence_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id AND source_id = NEW.source_id) BEGIN SELECT RAISE(ABORT, 'correction scope mismatch'); END;
         CREATE TRIGGER IF NOT EXISTS correction_scope_update BEFORE UPDATE OF tenant_id, person_id, superseded_claim_id, claim_id, source_id, evidence_id ON corrections FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.superseded_claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) OR NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) OR NOT EXISTS (SELECT 1 FROM sources WHERE id = NEW.source_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id AND origin_claim_id = NEW.claim_id AND origin_evidence_id = NEW.evidence_id) OR NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.evidence_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id AND source_id = NEW.source_id) BEGIN SELECT RAISE(ABORT, 'correction scope mismatch'); END;",
    )?;
    let empty: bool = transaction.query_row(
        "SELECT NOT EXISTS(SELECT 1 FROM memory_commits)",
        [],
        |row| row.get(0),
    )?;
    if empty {
        transaction.execute_batch(
            "INSERT INTO memory_commits(tenant_id, person_id, recorded_at)
             SELECT tenant_id, person_id, MAX(recorded_at) FROM (
               SELECT tenant_id, person_id, recorded_at FROM sources
               UNION ALL SELECT tenant_id, person_id, recorded_at FROM evidence
               UNION ALL SELECT tenant_id, person_id, recorded_from FROM claims
               UNION ALL SELECT tenant_id, person_id, recorded_at FROM profile_entries
               UNION ALL SELECT tenant_id, person_id, recorded_at FROM daily_reviews
             ) GROUP BY tenant_id, person_id;
             INSERT INTO memory_export_events(commit_sequence, event_index, payload)
             SELECT commits.sequence, ROW_NUMBER() OVER (PARTITION BY snapshots.tenant_id, snapshots.person_id ORDER BY record_kind, record_id) - 1, payload
             FROM (
               SELECT s.tenant_id, s.person_id, 'source' record_kind, s.id record_id,
                 json_object('kind','source','record',json_object('source',json_object('id',s.id,'tenant_id',s.tenant_id,'person_id',s.person_id,'revision',s.revision,'kind',json(s.kind),'content',s.content,'captured_at',s.captured_at,'recorded_at',s.recorded_at,'deleted_at',s.deleted_at),'ingestion_key',s.ingestion_key,'origin_evidence_id',s.origin_evidence_id,'origin_claim_id',s.origin_claim_id)) payload FROM sources s
               UNION ALL SELECT e.tenant_id,e.person_id,'evidence',e.id,
                 json_object('kind','evidence','record',json_object('evidence',json_object('id',e.id,'tenant_id',e.tenant_id,'person_id',e.person_id,'source_id',e.source_id,'source_revision',e.source_revision,'quote',e.quote,'byte_range',NULL,'recorded_at',e.recorded_at),'locator',CASE WHEN l.evidence_id IS NULL THEN NULL ELSE json_object('device_id',l.device_id,'provider',l.provider,'stream_id',l.stream_id,'segment_id',l.segment_id,'start_ms',l.start_ms,'end_ms',l.end_ms) END,'deleted_at',e.deleted_at)) FROM evidence e LEFT JOIN evidence_locators l ON l.evidence_id=e.id AND l.tenant_id=e.tenant_id AND l.person_id=e.person_id
               UNION ALL SELECT c.tenant_id,c.person_id,'claim',c.id,
                 json_object('kind','claim','record',json_object('id',c.id,'tenant_id',c.tenant_id,'person_id',c.person_id,'subject',c.subject,'predicate',c.predicate,'value',c.value,'kind',c.kind,'valid_time',json_object('from',c.valid_from,'until',c.valid_until),'recorded_time',json_object('from',c.recorded_from,'until',c.recorded_until),'status',c.status)) FROM claims c
               UNION ALL SELECT ce.tenant_id,ce.person_id,'claim_evidence',ce.claim_id||':'||ce.evidence_id,
                 json_object('kind','claim_evidence','record',json_object('tenant_id',ce.tenant_id,'person_id',ce.person_id,'claim_id',ce.claim_id,'evidence_id',ce.evidence_id,'relation',json(ce.relation),'confidence_basis_points',ce.confidence_basis_points)) FROM claim_evidence ce
               UNION ALL SELECT p.tenant_id,p.person_id,'profile',p.id,
                 json_object('kind','profile','record',json_object('id',p.id,'tenant_id',p.tenant_id,'person_id',p.person_id,'key',p.key,'value',p.value,'stability',json(p.stability),'claim_id',p.claim_id,'recorded_at',p.recorded_at)) FROM profile_entries p
               UNION ALL SELECT r.tenant_id,r.person_id,'review',r.id,
                 json_object('kind','daily_review','record',json_object('id',r.id,'tenant_id',r.tenant_id,'person_id',r.person_id,'day',r.day,'summary',r.summary,'evidence_ids',json(r.evidence_ids),'recorded_at',r.recorded_at)) FROM daily_reviews r
             ) snapshots JOIN memory_commits commits USING(tenant_id, person_id);",
        )?;
        let oversized: bool = transaction.query_row(
            "SELECT EXISTS(SELECT 1 FROM memory_export_events WHERE length(CAST(payload AS BLOB)) > ?1)",
            [MAX_EXPORT_RECORD_BYTES as i64],
            |row| row.get(0),
        )?;
        if oversized {
            return Err(Error::Invalid(format!(
                "legacy authoritative record exceeds the {MAX_EXPORT_RECORD_BYTES}-byte export compatibility limit"
            )));
        }
    }
    Ok(())
}

fn ensure_column(
    transaction: &Transaction<'_>,
    table: &str,
    column: &str,
    definition: &str,
) -> Result<()> {
    let exists = transaction.query_row(
        &format!("SELECT EXISTS(SELECT 1 FROM pragma_table_info('{table}') WHERE name = ?1)"),
        [column],
        |row| row.get::<_, bool>(0),
    )?;
    if !exists {
        transaction.execute(
            &format!("ALTER TABLE {table} ADD COLUMN {column} {definition}"),
            [],
        )?;
    }
    Ok(())
}

fn set_version(transaction: &Transaction<'_>, version: i64) -> Result<()> {
    transaction.execute_batch(&format!("PRAGMA user_version = {version};"))?;
    Ok(())
}

fn validate_scope(
    transaction: &Transaction<'_>,
    require_profile_content: bool,
    require_origins: bool,
) -> Result<()> {
    let mut checks = vec![
        (
            "evidence",
            "SELECT EXISTS(SELECT 1 FROM evidence e LEFT JOIN sources s ON s.id = e.source_id AND s.tenant_id = e.tenant_id AND s.person_id = e.person_id WHERE s.id IS NULL)",
        ),
        (
            "evidence locator",
            "SELECT EXISTS(SELECT 1 FROM evidence_locators l LEFT JOIN evidence e ON e.id = l.evidence_id AND e.tenant_id = l.tenant_id AND e.person_id = l.person_id WHERE e.id IS NULL)",
        ),
        (
            "claim evidence",
            "SELECT EXISTS(SELECT 1 FROM claim_evidence ce LEFT JOIN claims c ON c.id = ce.claim_id AND c.tenant_id = ce.tenant_id AND c.person_id = ce.person_id LEFT JOIN evidence e ON e.id = ce.evidence_id AND e.tenant_id = ce.tenant_id AND e.person_id = ce.person_id WHERE c.id IS NULL OR e.id IS NULL)",
        ),
        (
            "embedding",
            "SELECT EXISTS(SELECT 1 FROM embeddings p WHERE p.target_kind NOT IN ('source', 'evidence', 'claim') OR (p.target_kind = 'source' AND NOT EXISTS (SELECT 1 FROM sources s WHERE s.id = p.target_id AND s.tenant_id = p.tenant_id AND s.person_id = p.person_id)) OR (p.target_kind = 'evidence' AND NOT EXISTS (SELECT 1 FROM evidence e WHERE e.id = p.target_id AND e.tenant_id = p.tenant_id AND e.person_id = p.person_id)) OR (p.target_kind = 'claim' AND NOT EXISTS (SELECT 1 FROM claims c WHERE c.id = p.target_id AND c.tenant_id = p.tenant_id AND c.person_id = p.person_id)))",
        ),
        (
            "daily review",
            "SELECT EXISTS(SELECT 1 FROM daily_reviews r LEFT JOIN json_each(CASE WHEN json_valid(r.evidence_ids) THEN r.evidence_ids ELSE '[]' END) citation LEFT JOIN evidence e ON e.id = citation.value AND e.tenant_id = r.tenant_id AND e.person_id = r.person_id WHERE json_valid(r.evidence_ids) = 0 OR e.id IS NULL)",
        ),
    ];
    if require_origins {
        checks.push((
            "source origin",
            "SELECT EXISTS(SELECT 1 FROM sources s LEFT JOIN evidence e ON e.id = s.origin_evidence_id AND e.tenant_id = s.tenant_id AND e.person_id = s.person_id LEFT JOIN claims c ON c.id = s.origin_claim_id AND c.tenant_id = s.tenant_id AND c.person_id = s.person_id WHERE (s.origin_evidence_id IS NOT NULL AND e.id IS NULL) OR (s.origin_claim_id IS NOT NULL AND c.id IS NULL))",
        ));
    }
    if require_profile_content {
        checks.push((
            "profile entry",
            "SELECT EXISTS(SELECT 1 FROM profile_entries p LEFT JOIN claims c ON c.id = p.claim_id AND c.tenant_id = p.tenant_id AND c.person_id = p.person_id WHERE c.id IS NULL OR c.kind != 'profile_fact' OR c.predicate != p.key OR c.value != p.value)",
        ));
    }
    for (name, query) in checks {
        if transaction.query_row(query, [], |row| row.get::<_, bool>(0))? {
            return Err(Error::Invalid(format!(
                "legacy {name} is inconsistent with schema invariants"
            )));
        }
    }
    Ok(())
}

fn validate_profile_references(transaction: &Transaction<'_>) -> Result<()> {
    let inconsistent = transaction.query_row(
        "SELECT EXISTS(SELECT 1 FROM profile_entries p LEFT JOIN claims c ON c.id = p.claim_id AND c.tenant_id = p.tenant_id AND c.person_id = p.person_id WHERE c.id IS NULL OR c.kind != 'profile_fact')",
        [],
        |row| row.get::<_, bool>(0),
    )?;
    if inconsistent {
        return Err(Error::Invalid(
            "legacy profile entry is inconsistent with schema invariants".to_owned(),
        ));
    }
    Ok(())
}

fn install_core_scope_triggers(transaction: &Transaction<'_>) -> Result<()> {
    transaction.execute_batch(
        "
        CREATE TRIGGER IF NOT EXISTS evidence_scope_insert BEFORE INSERT ON evidence FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM sources WHERE id = NEW.source_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) BEGIN SELECT RAISE(ABORT, 'evidence source scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS evidence_scope_update BEFORE UPDATE OF source_id, tenant_id, person_id ON evidence FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM sources WHERE id = NEW.source_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) BEGIN SELECT RAISE(ABORT, 'evidence source scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS locator_scope_insert BEFORE INSERT ON evidence_locators FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.evidence_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) BEGIN SELECT RAISE(ABORT, 'evidence locator scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS locator_scope_update BEFORE UPDATE OF evidence_id, tenant_id, person_id ON evidence_locators FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.evidence_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) BEGIN SELECT RAISE(ABORT, 'evidence locator scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS claim_evidence_scope_insert BEFORE INSERT ON claim_evidence FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) OR NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.evidence_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) BEGIN SELECT RAISE(ABORT, 'claim evidence scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS claim_evidence_scope_update BEFORE UPDATE OF claim_id, evidence_id, tenant_id, person_id ON claim_evidence FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) OR NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.evidence_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id) BEGIN SELECT RAISE(ABORT, 'claim evidence scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS embedding_scope_insert BEFORE INSERT ON embeddings FOR EACH ROW WHEN NEW.target_kind NOT IN ('source', 'evidence', 'claim') OR (NEW.target_kind = 'source' AND NOT EXISTS (SELECT 1 FROM sources WHERE id = NEW.target_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) OR (NEW.target_kind = 'evidence' AND NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.target_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) OR (NEW.target_kind = 'claim' AND NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.target_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) BEGIN SELECT RAISE(ABORT, 'embedding target scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS embedding_scope_update BEFORE UPDATE OF target_kind, target_id, tenant_id, person_id ON embeddings FOR EACH ROW WHEN NEW.target_kind NOT IN ('source', 'evidence', 'claim') OR (NEW.target_kind = 'source' AND NOT EXISTS (SELECT 1 FROM sources WHERE id = NEW.target_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) OR (NEW.target_kind = 'evidence' AND NOT EXISTS (SELECT 1 FROM evidence WHERE id = NEW.target_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) OR (NEW.target_kind = 'claim' AND NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.target_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id)) BEGIN SELECT RAISE(ABORT, 'embedding target scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS review_scope_insert BEFORE INSERT ON daily_reviews FOR EACH ROW WHEN json_valid(NEW.evidence_ids) = 0 OR EXISTS (SELECT 1 FROM json_each(NEW.evidence_ids) citation LEFT JOIN evidence e ON e.id = citation.value AND e.tenant_id = NEW.tenant_id AND e.person_id = NEW.person_id WHERE e.id IS NULL) BEGIN SELECT RAISE(ABORT, 'daily review evidence scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS review_scope_update BEFORE UPDATE OF tenant_id, person_id, evidence_ids ON daily_reviews FOR EACH ROW WHEN json_valid(NEW.evidence_ids) = 0 OR EXISTS (SELECT 1 FROM json_each(NEW.evidence_ids) citation LEFT JOIN evidence e ON e.id = citation.value AND e.tenant_id = NEW.tenant_id AND e.person_id = NEW.person_id WHERE e.id IS NULL) BEGIN SELECT RAISE(ABORT, 'daily review evidence scope mismatch'); END;
        ",
    )?;
    Ok(())
}

fn install_profile_triggers(transaction: &Transaction<'_>) -> Result<()> {
    transaction.execute_batch(
        "
        CREATE TRIGGER IF NOT EXISTS profile_entry_scope_insert BEFORE INSERT ON profile_entries FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id AND kind = 'profile_fact' AND predicate = NEW.key AND value = NEW.value) BEGIN SELECT RAISE(ABORT, 'profile entry claim scope mismatch'); END;
        CREATE TRIGGER IF NOT EXISTS profile_entry_scope_update BEFORE UPDATE OF key, value, claim_id, tenant_id, person_id ON profile_entries FOR EACH ROW WHEN NOT EXISTS (SELECT 1 FROM claims WHERE id = NEW.claim_id AND tenant_id = NEW.tenant_id AND person_id = NEW.person_id AND kind = 'profile_fact' AND predicate = NEW.key AND value = NEW.value) BEGIN SELECT RAISE(ABORT, 'profile entry claim scope mismatch'); END;
        ",
    )?;
    Ok(())
}