Skip to main content

mempal_store_sqlite/
lib.rs

1use std::collections::{BTreeSet, VecDeque};
2use std::fs;
3use std::path::{Path, PathBuf};
4use std::sync::OnceLock;
5use std::time::{SystemTime, UNIX_EPOCH};
6
7use mempal_agent_memory::anchor;
8use mempal_agent_memory::types::{
9    AnchorKind, ChunkNeighbors, Drawer, ExplicitTunnel, KnowledgeCard, KnowledgeCardEvent,
10    KnowledgeCardFilter, KnowledgeEventType, KnowledgeEvidenceLink, KnowledgeEvidenceRole,
11    KnowledgeStatus, KnowledgeTier, MemoryDomain, MemoryKind, NeighborChunk, Provenance,
12    ReindexSource, RuntimeAdoptionEvent, RuntimeAdoptionFilter, RuntimeAdoptionSignal,
13    RuntimeAdoptionTrack, SourceType, TaxonomyEntry, Triple, TripleStats, TunnelEndpoint,
14    TunnelFollowResult,
15};
16use mempal_search_core::build_fts_match_query;
17use rusqlite::{Connection, OptionalExtension, Row, params};
18use serde_json::Value;
19use sha2::{Digest, Sha256};
20use thiserror::Error;
21
22pub const CURRENT_SCHEMA_VERSION: u32 = 9;
23const DRAWER_SELECT_COLUMNS: &str = r#"
24    id,
25    content,
26    wing,
27    room,
28    source_file,
29    source_type,
30    added_at,
31    chunk_index,
32    normalize_version,
33    COALESCE(importance, 0) as importance,
34    memory_kind,
35    domain,
36    field,
37    anchor_kind,
38    anchor_id,
39    parent_anchor_id,
40    provenance,
41    statement,
42    tier,
43    status,
44    supporting_refs,
45    counterexample_refs,
46    teaching_refs,
47    verification_refs,
48    scope_constraints,
49    trigger_hints
50"#;
51
52const V1_SCHEMA_SQL: &str = r#"
53PRAGMA foreign_keys = ON;
54
55CREATE TABLE IF NOT EXISTS drawers (
56    id TEXT PRIMARY KEY,
57    content TEXT NOT NULL,
58    wing TEXT NOT NULL,
59    room TEXT,
60    source_file TEXT,
61    source_type TEXT NOT NULL CHECK(source_type IN ('project', 'conversation', 'manual')),
62    added_at TEXT NOT NULL,
63    chunk_index INTEGER
64);
65
66-- drawer_vectors is created lazily by insert_vector() with the actual
67-- embedding dimension from the configured embedder. This avoids hardcoding
68-- a dimension that may not match the model in use.
69
70CREATE TABLE IF NOT EXISTS triples (
71    id TEXT PRIMARY KEY,
72    subject TEXT NOT NULL,
73    predicate TEXT NOT NULL,
74    object TEXT NOT NULL,
75    valid_from TEXT,
76    valid_to TEXT,
77    confidence REAL DEFAULT 1.0,
78    source_drawer TEXT REFERENCES drawers(id)
79);
80
81CREATE TABLE IF NOT EXISTS taxonomy (
82    wing TEXT NOT NULL,
83    room TEXT NOT NULL DEFAULT '',
84    display_name TEXT,
85    keywords TEXT,
86    PRIMARY KEY (wing, room)
87);
88
89CREATE INDEX IF NOT EXISTS idx_drawers_wing ON drawers(wing);
90CREATE INDEX IF NOT EXISTS idx_drawers_wing_room ON drawers(wing, room);
91CREATE INDEX IF NOT EXISTS idx_triples_subject ON triples(subject);
92CREATE INDEX IF NOT EXISTS idx_triples_object ON triples(object);
93"#;
94
95static SQLITE_VEC_AUTO_EXTENSION: OnceLock<Result<(), String>> = OnceLock::new();
96
97fn current_timestamp() -> String {
98    match SystemTime::now().duration_since(UNIX_EPOCH) {
99        Ok(duration) => duration.as_secs().to_string(),
100        Err(_) => "0".to_string(),
101    }
102}
103
104fn build_tunnel_id(left: &TunnelEndpoint, right: &TunnelEndpoint) -> String {
105    let mut endpoints = [tunnel_endpoint_key(left), tunnel_endpoint_key(right)];
106    endpoints.sort();
107
108    let mut hasher = Sha256::new();
109    for component in [
110        endpoints[0].0.as_str(),
111        endpoints[0].1.as_str(),
112        endpoints[1].0.as_str(),
113        endpoints[1].1.as_str(),
114    ] {
115        hasher.update([0]);
116        hasher.update(component.as_bytes());
117    }
118    let digest = format!("{:x}", hasher.finalize());
119    format!("tunnel_{}", &digest[..16])
120}
121
122fn tunnel_endpoint_key(endpoint: &TunnelEndpoint) -> (String, String) {
123    (
124        endpoint.wing.trim().to_string(),
125        endpoint
126            .room
127            .as_deref()
128            .map(str::trim)
129            .filter(|room| !room.is_empty())
130            .unwrap_or("")
131            .to_string(),
132    )
133}
134
135fn format_tunnel_endpoint(endpoint: &TunnelEndpoint) -> String {
136    match endpoint.room.as_deref() {
137        Some(room) if !room.is_empty() => format!("{}:{room}", endpoint.wing),
138        _ => endpoint.wing.clone(),
139    }
140}
141
142#[derive(Debug, Error)]
143pub enum DbError {
144    #[error("failed to create database directory for {path}")]
145    CreateDir {
146        path: PathBuf,
147        #[source]
148        source: std::io::Error,
149    },
150    #[error("failed to read database metadata for {path}")]
151    Metadata {
152        path: PathBuf,
153        #[source]
154        source: std::io::Error,
155    },
156    #[error(transparent)]
157    Sqlite(#[from] rusqlite::Error),
158    #[error("failed to parse taxonomy keywords JSON")]
159    Json(#[from] serde_json::Error),
160    #[error("invalid source_type stored in database: {0}")]
161    InvalidSourceType(String),
162    #[error("invalid {kind} stored in database: {value}")]
163    InvalidEnumValue { kind: &'static str, value: String },
164    #[error("invalid drawer metadata: {0}")]
165    InvalidDrawerMetadata(String),
166    #[error(
167        "source {source_file} in wing {wing} is protected by {references} knowledge references across {referenced_drawers} drawers"
168    )]
169    SourceProtectedByKnowledgeReferences {
170        source_file: String,
171        wing: String,
172        referenced_drawers: u64,
173        references: u64,
174    },
175    #[error("invalid tunnel: {0}")]
176    InvalidTunnel(String),
177    #[error("failed to register sqlite-vec auto extension: {0}")]
178    RegisterVec(String),
179    #[error("database schema version {current} is newer than supported version {supported}")]
180    UnsupportedSchemaVersion { current: u32, supported: u32 },
181}
182
183pub struct Database {
184    conn: Connection,
185    path: PathBuf,
186}
187
188#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
189pub struct KnowledgeReferenceSummary {
190    pub referenced_drawers: u64,
191    pub references: u64,
192}
193
194impl Database {
195    pub fn open(path: &Path) -> Result<Self, DbError> {
196        if let Some(parent) = path
197            .parent()
198            .filter(|parent| !parent.as_os_str().is_empty())
199        {
200            fs::create_dir_all(parent).map_err(|source| DbError::CreateDir {
201                path: parent.to_path_buf(),
202                source,
203            })?;
204        }
205
206        register_sqlite_vec()?;
207
208        let conn = Connection::open(path)?;
209        conn.execute_batch("PRAGMA foreign_keys = ON;")?;
210        apply_migrations(&conn)?;
211
212        Ok(Self {
213            conn,
214            path: path.to_path_buf(),
215        })
216    }
217
218    pub fn conn(&self) -> &Connection {
219        &self.conn
220    }
221
222    pub fn path(&self) -> &Path {
223        &self.path
224    }
225
226    pub fn insert_drawer(&self, drawer: &Drawer) -> Result<bool, DbError> {
227        anchor::validate_anchor_domain(&drawer.domain, &drawer.anchor_kind)
228            .map_err(|message| DbError::InvalidDrawerMetadata(message.to_string()))?;
229
230        let inserted = self.conn.execute(
231            r#"
232            INSERT INTO drawers (
233                id,
234                content,
235                wing,
236                room,
237                source_file,
238                source_type,
239                added_at,
240                chunk_index,
241                normalize_version,
242                importance,
243                memory_kind,
244                domain,
245                field,
246                anchor_kind,
247                anchor_id,
248                parent_anchor_id,
249                provenance,
250                statement,
251                tier,
252                status,
253                supporting_refs,
254                counterexample_refs,
255                teaching_refs,
256                verification_refs,
257                scope_constraints,
258                trigger_hints
259            )
260            VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16, ?17, ?18, ?19, ?20, ?21, ?22, ?23, ?24, ?25, ?26)
261            ON CONFLICT(id) DO NOTHING
262            "#,
263            params![
264                drawer.id.as_str(),
265                drawer.content.as_str(),
266                drawer.wing.as_str(),
267                drawer.room.as_deref(),
268                drawer.source_file.as_deref(),
269                source_type_as_str(&drawer.source_type),
270                drawer.added_at.as_str(),
271                drawer.chunk_index,
272                i64::from(drawer.normalize_version),
273                drawer.importance,
274                memory_kind_as_str(&drawer.memory_kind),
275                memory_domain_as_str(&drawer.domain),
276                drawer.field.as_str(),
277                anchor_kind_as_str(&drawer.anchor_kind),
278                drawer.anchor_id.as_str(),
279                drawer.parent_anchor_id.as_deref(),
280                drawer.provenance.as_ref().map(provenance_as_str),
281                drawer.statement.as_deref(),
282                drawer.tier.as_ref().map(knowledge_tier_as_str),
283                drawer.status.as_ref().map(knowledge_status_as_str),
284                encode_json(&drawer.supporting_refs)?,
285                encode_json(&drawer.counterexample_refs)?,
286                encode_json(&drawer.teaching_refs)?,
287                encode_json(&drawer.verification_refs)?,
288                drawer.scope_constraints.as_deref(),
289                encode_optional_json(drawer.trigger_hints.as_ref())?,
290            ],
291        )?;
292
293        Ok(inserted == 1)
294    }
295
296    pub fn taxonomy_entries(&self) -> Result<Vec<TaxonomyEntry>, DbError> {
297        let mut statement = self.conn.prepare(
298            "SELECT wing, room, display_name, keywords FROM taxonomy ORDER BY wing, room",
299        )?;
300        let rows = statement.query_map([], |row| {
301            Ok((
302                row.get::<_, String>(0)?,
303                row.get::<_, String>(1)?,
304                row.get::<_, Option<String>>(2)?,
305                row.get::<_, Option<String>>(3)?,
306            ))
307        })?;
308
309        let mut entries = Vec::new();
310        for row in rows {
311            let (wing, room, display_name, keywords_json) = row?;
312            let keywords = parse_keywords(keywords_json.as_deref())?;
313            entries.push(TaxonomyEntry {
314                wing,
315                room,
316                display_name,
317                keywords,
318            });
319        }
320
321        Ok(entries)
322    }
323
324    pub fn upsert_taxonomy_entry(&self, entry: &TaxonomyEntry) -> Result<(), DbError> {
325        let keywords = serde_json::to_string(&entry.keywords)?;
326        self.conn.execute(
327            r#"
328            INSERT INTO taxonomy (wing, room, display_name, keywords)
329            VALUES (?1, ?2, ?3, ?4)
330            ON CONFLICT(wing, room) DO UPDATE SET
331                display_name = excluded.display_name,
332                keywords = excluded.keywords
333            "#,
334            (
335                entry.wing.as_str(),
336                entry.room.as_str(),
337                entry.display_name.as_deref(),
338                keywords.as_str(),
339            ),
340        )?;
341
342        Ok(())
343    }
344
345    /// Returns top drawers sorted by importance (descending), then recency.
346    pub fn top_drawers(&self, limit: usize) -> Result<Vec<Drawer>, DbError> {
347        let limit = i64::try_from(limit)
348            .map_err(|_| rusqlite::Error::InvalidParameterName("limit".to_string()))?;
349        let mut statement = self.conn.prepare(&format!(
350            r#"
351            SELECT {DRAWER_SELECT_COLUMNS}
352            FROM drawers
353            WHERE deleted_at IS NULL
354            ORDER BY importance DESC, CAST(added_at AS INTEGER) DESC, id DESC
355            LIMIT ?1
356            "#,
357        ))?;
358        let rows = statement.query_map([limit], |row| {
359            drawer_from_row(row).map_err(row_decode_error)
360        })?;
361
362        let mut drawers = Vec::new();
363        for row in rows {
364            drawers.push(row?);
365        }
366
367        Ok(drawers)
368    }
369
370    pub fn drawer_exists(&self, drawer_id: &str) -> Result<bool, DbError> {
371        let exists = self.conn.query_row(
372            "SELECT EXISTS(SELECT 1 FROM drawers WHERE id = ?1 AND deleted_at IS NULL)",
373            [drawer_id],
374            |row| row.get::<_, i64>(0),
375        )?;
376        Ok(exists == 1)
377    }
378
379    pub fn insert_vector(&self, drawer_id: &str, vector: &[f32]) -> Result<(), DbError> {
380        self.ensure_vectors_table(vector.len())?;
381        let vector_json = serde_json::to_string(vector)?;
382        self.conn.execute(
383            "INSERT INTO drawer_vectors (id, embedding) VALUES (?1, vec_f32(?2))",
384            (drawer_id, vector_json.as_str()),
385        )?;
386        Ok(())
387    }
388
389    pub fn upsert_drawer_and_replace_vector(
390        &self,
391        drawer: &Drawer,
392        vector: &[f32],
393    ) -> Result<(), DbError> {
394        anchor::validate_anchor_domain(&drawer.domain, &drawer.anchor_kind)
395            .map_err(|message| DbError::InvalidDrawerMetadata(message.to_string()))?;
396        self.ensure_vectors_table(vector.len())?;
397
398        let existing = self
399            .conn
400            .query_row(
401                "SELECT rowid, content FROM drawers WHERE id = ?1 AND deleted_at IS NULL",
402                [drawer.id.as_str()],
403                |row| Ok((row.get::<_, i64>(0)?, row.get::<_, String>(1)?)),
404            )
405            .optional()?;
406
407        if existing.is_none() {
408            if self.insert_drawer(drawer)? {
409                return self.insert_vector(&drawer.id, vector);
410            }
411            return Ok(());
412        }
413
414        let (rowid, old_content) = existing.expect("checked Some");
415        let vector_json = serde_json::to_string(vector)?;
416
417        self.conn.execute_batch("BEGIN IMMEDIATE;")?;
418        let result = (|| -> Result<(), DbError> {
419            if self.table_exists("drawers_fts")? {
420                self.conn.execute(
421                    "INSERT INTO drawers_fts(drawers_fts, rowid, content) VALUES ('delete', ?1, ?2)",
422                    params![rowid, old_content],
423                )?;
424            }
425
426            self.conn.execute(
427                r#"
428                UPDATE drawers
429                SET content = ?2,
430                    wing = ?3,
431                    room = ?4,
432                    source_file = ?5,
433                    source_type = ?6,
434                    added_at = ?7,
435                    chunk_index = ?8,
436                    normalize_version = ?9,
437                    importance = ?10,
438                    memory_kind = ?11,
439                    domain = ?12,
440                    field = ?13,
441                    anchor_kind = ?14,
442                    anchor_id = ?15,
443                    parent_anchor_id = ?16,
444                    provenance = ?17,
445                    statement = ?18,
446                    tier = ?19,
447                    status = ?20,
448                    supporting_refs = ?21,
449                    counterexample_refs = ?22,
450                    teaching_refs = ?23,
451                    verification_refs = ?24,
452                    scope_constraints = ?25,
453                    trigger_hints = ?26
454                WHERE id = ?1 AND deleted_at IS NULL
455                "#,
456                params![
457                    drawer.id.as_str(),
458                    drawer.content.as_str(),
459                    drawer.wing.as_str(),
460                    drawer.room.as_deref(),
461                    drawer.source_file.as_deref(),
462                    source_type_as_str(&drawer.source_type),
463                    drawer.added_at.as_str(),
464                    drawer.chunk_index,
465                    i64::from(drawer.normalize_version),
466                    drawer.importance,
467                    memory_kind_as_str(&drawer.memory_kind),
468                    memory_domain_as_str(&drawer.domain),
469                    drawer.field.as_str(),
470                    anchor_kind_as_str(&drawer.anchor_kind),
471                    drawer.anchor_id.as_str(),
472                    drawer.parent_anchor_id.as_deref(),
473                    drawer.provenance.as_ref().map(provenance_as_str),
474                    drawer.statement.as_deref(),
475                    drawer.tier.as_ref().map(knowledge_tier_as_str),
476                    drawer.status.as_ref().map(knowledge_status_as_str),
477                    encode_json(&drawer.supporting_refs)?,
478                    encode_json(&drawer.counterexample_refs)?,
479                    encode_json(&drawer.teaching_refs)?,
480                    encode_json(&drawer.verification_refs)?,
481                    drawer.scope_constraints.as_deref(),
482                    encode_optional_json(drawer.trigger_hints.as_ref())?,
483                ],
484            )?;
485
486            if self.table_exists("drawers_fts")? {
487                self.conn.execute(
488                    "INSERT INTO drawers_fts(rowid, content) VALUES (?1, ?2)",
489                    params![rowid, drawer.content.as_str()],
490                )?;
491            }
492
493            self.conn.execute(
494                "DELETE FROM drawer_vectors WHERE id = ?1",
495                [drawer.id.as_str()],
496            )?;
497            self.conn.execute(
498                "INSERT INTO drawer_vectors (id, embedding) VALUES (?1, vec_f32(?2))",
499                params![drawer.id.as_str(), vector_json.as_str()],
500            )?;
501
502            Ok(())
503        })();
504
505        match result {
506            Ok(()) => {
507                self.conn.execute_batch("COMMIT;")?;
508                Ok(())
509            }
510            Err(error) => {
511                let _ = self.conn.execute_batch("ROLLBACK;");
512                Err(error)
513            }
514        }
515    }
516
517    /// Ensure drawer_vectors table exists with the right dimension.
518    /// Creates it on first call; errors on dimension mismatch.
519    fn ensure_vectors_table(&self, dim: usize) -> Result<(), DbError> {
520        // Check if table exists
521        let exists: bool = self
522            .conn
523            .query_row(
524                "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type='table' AND name='drawer_vectors')",
525                [],
526                |row| row.get(0),
527            )?;
528
529        if !exists {
530            self.conn.execute_batch(&format!(
531                "CREATE VIRTUAL TABLE IF NOT EXISTS drawer_vectors USING vec0(id TEXT PRIMARY KEY, embedding FLOAT[{dim}]);"
532            ))?;
533        }
534        Ok(())
535    }
536
537    pub fn drawer_count(&self) -> Result<i64, DbError> {
538        Ok(self.conn.query_row(
539            "SELECT COUNT(*) FROM drawers WHERE deleted_at IS NULL",
540            [],
541            |row| row.get(0),
542        )?)
543    }
544
545    pub fn stale_drawer_count(&self, current_normalize_version: u32) -> Result<i64, DbError> {
546        Ok(self.conn.query_row(
547            "SELECT COUNT(*) FROM drawers WHERE deleted_at IS NULL AND normalize_version < ?1",
548            [i64::from(current_normalize_version)],
549            |row| row.get(0),
550        )?)
551    }
552
553    pub fn drawer_count_by_normalize_version(&self) -> Result<Vec<(u32, i64)>, DbError> {
554        let mut statement = self.conn.prepare(
555            r#"
556            SELECT normalize_version, COUNT(*)
557            FROM drawers
558            WHERE deleted_at IS NULL
559            GROUP BY normalize_version
560            ORDER BY normalize_version
561            "#,
562        )?;
563        let rows = statement
564            .query_map([], |row| Ok((row.get::<_, u32>(0)?, row.get::<_, i64>(1)?)))?
565            .collect::<std::result::Result<Vec<_>, _>>()?;
566        Ok(rows)
567    }
568
569    pub fn diary_rollup_days(&self) -> Result<u32, DbError> {
570        let count = self.conn.query_row(
571            r#"
572            SELECT COUNT(DISTINCT substr(source_file, length(source_file) - 9, 10))
573            FROM drawers
574            WHERE deleted_at IS NULL
575              AND wing = 'agent-diary'
576              AND source_file LIKE 'agent-diary://rollup/%'
577            "#,
578            [],
579            |row| row.get::<_, i64>(0),
580        )?;
581        Ok(count as u32)
582    }
583
584    pub fn reindex_sources_stale(
585        &self,
586        current_normalize_version: u32,
587    ) -> Result<Vec<ReindexSource>, DbError> {
588        let mut statement = self.conn.prepare(
589            r#"
590            SELECT source_file, wing, room, COUNT(*)
591            FROM drawers
592            WHERE deleted_at IS NULL AND normalize_version < ?1
593            GROUP BY source_file, wing, room
594            ORDER BY source_file, wing, room
595            "#,
596        )?;
597        let rows = statement
598            .query_map(
599                [i64::from(current_normalize_version)],
600                reindex_source_from_row,
601            )?
602            .collect::<std::result::Result<Vec<_>, _>>()?;
603        Ok(rows)
604    }
605
606    pub fn reindex_sources_force(&self) -> Result<Vec<ReindexSource>, DbError> {
607        let mut statement = self.conn.prepare(
608            r#"
609            SELECT source_file, wing, room, COUNT(*)
610            FROM drawers
611            WHERE deleted_at IS NULL
612            GROUP BY source_file, wing, room
613            ORDER BY source_file, wing, room
614            "#,
615        )?;
616        let rows = statement
617            .query_map([], reindex_source_from_row)?
618            .collect::<std::result::Result<Vec<_>, _>>()?;
619        Ok(rows)
620    }
621
622    /// Count active evidence drawers for one physical source that are protected
623    /// by Stage-1 knowledge reference arrays or Phase-2 evidence links.
624    ///
625    /// Source replacement must not hard-delete these drawers without an
626    /// explicit old-to-new reference mapping. P117 therefore uses this as a
627    /// conservative preflight and skips the whole source when the summary is
628    /// non-empty.
629    pub fn source_knowledge_reference_summary(
630        &self,
631        source_file: &str,
632        wing: &str,
633    ) -> Result<KnowledgeReferenceSummary, DbError> {
634        let (referenced_drawers, references) = self.conn.query_row(
635            r#"
636            WITH stage1_refs(evidence_id) AS (
637                SELECT ref.value
638                FROM drawers AS knowledge, json_each(knowledge.supporting_refs) AS ref
639                WHERE knowledge.deleted_at IS NULL AND knowledge.memory_kind = 'knowledge'
640                UNION ALL
641                SELECT ref.value
642                FROM drawers AS knowledge, json_each(knowledge.counterexample_refs) AS ref
643                WHERE knowledge.deleted_at IS NULL AND knowledge.memory_kind = 'knowledge'
644                UNION ALL
645                SELECT ref.value
646                FROM drawers AS knowledge, json_each(knowledge.teaching_refs) AS ref
647                WHERE knowledge.deleted_at IS NULL AND knowledge.memory_kind = 'knowledge'
648                UNION ALL
649                SELECT ref.value
650                FROM drawers AS knowledge, json_each(knowledge.verification_refs) AS ref
651                WHERE knowledge.deleted_at IS NULL AND knowledge.memory_kind = 'knowledge'
652            ),
653            all_refs(evidence_id) AS (
654                SELECT evidence_id FROM stage1_refs
655                UNION ALL
656                SELECT evidence_drawer_id FROM knowledge_evidence_links
657            )
658            SELECT COUNT(DISTINCT evidence.id), COUNT(*)
659            FROM all_refs
660            JOIN drawers AS evidence ON evidence.id = all_refs.evidence_id
661            WHERE evidence.deleted_at IS NULL
662              AND evidence.source_file = ?1
663              AND evidence.wing = ?2
664            "#,
665            (source_file, wing),
666            |row| Ok((row.get::<_, i64>(0)?, row.get::<_, i64>(1)?)),
667        )?;
668        Ok(KnowledgeReferenceSummary {
669            referenced_drawers: referenced_drawers as u64,
670            references: references as u64,
671        })
672    }
673
674    /// Hard-delete the active drawers for a source scoped to a specific room
675    /// (NULL room matches NULL room), removing their FTS and vector rows too.
676    ///
677    /// Use this when the caller knows the exact room the drawers live in.
678    /// For reindex, prefer [`replace_active_source_drawers_across_rooms`]:
679    /// re-ingesting a physical source may re-route it to a different room, and
680    /// a room-scoped delete would miss the stale drawers in the old room and
681    /// leave duplicates behind.
682    /// Run `f` inside one IMMEDIATE transaction: commit on Ok, roll back on
683    /// Err (P117). Callables must use `_in_txn` method variants — nesting a
684    /// method that opens its own transaction will fail.
685    pub fn with_immediate_transaction<T, E>(
686        &self,
687        f: impl FnOnce(&Self) -> Result<T, E>,
688    ) -> Result<T, E>
689    where
690        E: From<DbError>,
691    {
692        self.conn
693            .execute_batch("BEGIN IMMEDIATE;")
694            .map_err(DbError::from)
695            .map_err(E::from)?;
696        match f(self) {
697            Ok(value) => {
698                if let Err(error) = self
699                    .conn
700                    .execute_batch("COMMIT;")
701                    .map_err(DbError::from)
702                    .map_err(E::from)
703                {
704                    let _ = self.conn.execute_batch("ROLLBACK;");
705                    return Err(error);
706                }
707                Ok(value)
708            }
709            Err(error) => {
710                let _ = self.conn.execute_batch("ROLLBACK;");
711                Err(error)
712            }
713        }
714    }
715
716    /// Non-transactional variant of [`Self::replace_active_source_drawers`]
717    /// for use inside [`Self::with_immediate_transaction`] (P117).
718    pub fn replace_active_source_drawers_in_txn(
719        &self,
720        source_file: &str,
721        wing: &str,
722        room: Option<&str>,
723    ) -> Result<u64, DbError> {
724        self.ensure_source_is_not_knowledge_referenced(source_file, wing)?;
725        let rows = self.active_source_rows_in_room(source_file, wing, room)?;
726        self.delete_source_drawer_rows_in_txn(&rows)
727    }
728
729    /// Non-transactional variant of
730    /// [`Self::replace_active_source_drawers_across_rooms`] for use inside
731    /// [`Self::with_immediate_transaction`] (P117).
732    pub fn replace_active_source_drawers_across_rooms_in_txn(
733        &self,
734        source_file: &str,
735        wing: &str,
736    ) -> Result<u64, DbError> {
737        self.ensure_source_is_not_knowledge_referenced(source_file, wing)?;
738        let rows = self.active_source_rows_all_rooms(source_file, wing)?;
739        self.delete_source_drawer_rows_in_txn(&rows)
740    }
741
742    pub fn replace_active_source_drawers(
743        &self,
744        source_file: &str,
745        wing: &str,
746        room: Option<&str>,
747    ) -> Result<u64, DbError> {
748        self.with_immediate_transaction(|db| {
749            db.replace_active_source_drawers_in_txn(source_file, wing, room)
750        })
751    }
752
753    fn active_source_rows_in_room(
754        &self,
755        source_file: &str,
756        wing: &str,
757        room: Option<&str>,
758    ) -> Result<Vec<(i64, String, String)>, DbError> {
759        let mut statement = self.conn.prepare(
760            r#"
761            SELECT rowid, id, content
762            FROM drawers
763            WHERE deleted_at IS NULL
764              AND source_file = ?1
765              AND wing = ?2
766              AND ((?3 IS NULL AND room IS NULL) OR room = ?3)
767            ORDER BY rowid
768            "#,
769        )?;
770        let rows = statement
771            .query_map((source_file, wing, room), |row| {
772                Ok((
773                    row.get::<_, i64>(0)?,
774                    row.get::<_, String>(1)?,
775                    row.get::<_, String>(2)?,
776                ))
777            })?
778            .collect::<std::result::Result<Vec<_>, _>>()?;
779        Ok(rows)
780    }
781
782    /// Hard-delete the active drawers for a source across ALL rooms in a wing.
783    ///
784    /// This is the correct replace semantics for reindex: a physical source
785    /// file maps to one logical source, so re-indexing it should keep only the
786    /// freshly produced drawers regardless of which room each previous version
787    /// was routed into. Without this, a source that auto-routes to a new room
788    /// on reindex leaves its old-room drawers behind as duplicates.
789    pub fn replace_active_source_drawers_across_rooms(
790        &self,
791        source_file: &str,
792        wing: &str,
793    ) -> Result<u64, DbError> {
794        self.with_immediate_transaction(|db| {
795            db.replace_active_source_drawers_across_rooms_in_txn(source_file, wing)
796        })
797    }
798
799    fn active_source_rows_all_rooms(
800        &self,
801        source_file: &str,
802        wing: &str,
803    ) -> Result<Vec<(i64, String, String)>, DbError> {
804        let mut statement = self.conn.prepare(
805            r#"
806            SELECT rowid, id, content
807            FROM drawers
808            WHERE deleted_at IS NULL
809              AND source_file = ?1
810              AND wing = ?2
811            ORDER BY rowid
812            "#,
813        )?;
814        let rows = statement
815            .query_map((source_file, wing), |row| {
816                Ok((
817                    row.get::<_, i64>(0)?,
818                    row.get::<_, String>(1)?,
819                    row.get::<_, String>(2)?,
820                ))
821            })?
822            .collect::<std::result::Result<Vec<_>, _>>()?;
823        Ok(rows)
824    }
825
826    /// Delete body without transaction control — the caller owns the
827    /// transaction (P117).
828    fn delete_source_drawer_rows_in_txn(
829        &self,
830        rows: &[(i64, String, String)],
831    ) -> Result<u64, DbError> {
832        if rows.is_empty() {
833            return Ok(0);
834        }
835
836        let fts_exists = self.table_exists("drawers_fts")?;
837        let vectors_exist = self.table_exists("drawer_vectors")?;
838
839        for (rowid, id, content) in rows {
840            if fts_exists {
841                self.conn.execute(
842                    "INSERT INTO drawers_fts(drawers_fts, rowid, content) VALUES ('delete', ?1, ?2)",
843                    params![rowid, content],
844                )?;
845            }
846            if vectors_exist {
847                self.conn
848                    .execute("DELETE FROM drawer_vectors WHERE id = ?1", [id])?;
849            }
850            // triples.source_drawer is a FK to drawers(id) (RESTRICT). Drop
851            // the dangling provenance link before the hard delete, otherwise
852            // deleting a drawer referenced by a KG triple fails with a
853            // FOREIGN KEY constraint error. The triple (a KG fact) is kept;
854            // only its stale source pointer is cleared.
855            self.conn.execute(
856                "UPDATE triples SET source_drawer = NULL WHERE source_drawer = ?1",
857                [id],
858            )?;
859            self.conn
860                .execute("DELETE FROM drawers WHERE rowid = ?1", [rowid])?;
861        }
862
863        Ok(rows.len() as u64)
864    }
865
866    fn ensure_source_is_not_knowledge_referenced(
867        &self,
868        source_file: &str,
869        wing: &str,
870    ) -> Result<(), DbError> {
871        let summary = self.source_knowledge_reference_summary(source_file, wing)?;
872        if summary.referenced_drawers == 0 {
873            return Ok(());
874        }
875        Err(DbError::SourceProtectedByKnowledgeReferences {
876            source_file: source_file.to_string(),
877            wing: wing.to_string(),
878            referenced_drawers: summary.referenced_drawers,
879            references: summary.references,
880        })
881    }
882
883    fn table_exists(&self, table_name: &str) -> Result<bool, DbError> {
884        let exists = self.conn.query_row(
885            "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type='table' AND name = ?1)",
886            [table_name],
887            |row| row.get::<_, i64>(0),
888        )?;
889        Ok(exists == 1)
890    }
891
892    pub fn taxonomy_count(&self) -> Result<i64, DbError> {
893        Ok(self
894            .conn
895            .query_row("SELECT COUNT(*) FROM taxonomy", [], |row| row.get(0))?)
896    }
897
898    pub fn scope_counts(&self) -> Result<Vec<(String, Option<String>, i64)>, DbError> {
899        let mut statement = self.conn.prepare(
900            r#"
901            SELECT wing, room, COUNT(*)
902            FROM drawers
903            WHERE deleted_at IS NULL
904            GROUP BY wing, room
905            ORDER BY wing, room
906            "#,
907        )?;
908        let rows = statement
909            .query_map([], |row| {
910                Ok((
911                    row.get::<_, String>(0)?,
912                    row.get::<_, Option<String>>(1)?,
913                    row.get::<_, i64>(2)?,
914                ))
915            })?
916            .collect::<std::result::Result<Vec<_>, _>>()?;
917        Ok(rows)
918    }
919
920    /// Per-field counts for the P106 distill signal: active evidence drawers and
921    /// active promoted-or-canonical knowledge drawers, grouped by `field`.
922    /// Read-only; performs no writes.
923    pub fn distill_field_counts(&self) -> Result<Vec<(String, i64, i64)>, DbError> {
924        let mut statement = self.conn.prepare(
925            r#"
926            SELECT field,
927                   SUM(CASE WHEN memory_kind = 'evidence' THEN 1 ELSE 0 END) AS evidence_count,
928                   SUM(CASE WHEN memory_kind = 'knowledge'
929                             AND status IN ('promoted', 'canonical') THEN 1 ELSE 0 END) AS promoted_count
930            FROM drawers
931            WHERE deleted_at IS NULL
932            GROUP BY field
933            "#,
934        )?;
935        let rows = statement
936            .query_map([], |row| {
937                Ok((
938                    row.get::<_, String>(0)?,
939                    row.get::<_, i64>(1)?,
940                    row.get::<_, i64>(2)?,
941                ))
942            })?
943            .collect::<std::result::Result<Vec<_>, _>>()?;
944        Ok(rows)
945    }
946
947    /// Up to `limit` active evidence drawer ids for a `field`, ordered by rowid
948    /// for deterministic sampling. Read-only; performs no writes.
949    pub fn sample_evidence_drawer_ids(
950        &self,
951        field: &str,
952        limit: usize,
953    ) -> Result<Vec<String>, DbError> {
954        let mut statement = self.conn.prepare(
955            r#"
956            SELECT id
957            FROM drawers
958            WHERE deleted_at IS NULL AND memory_kind = 'evidence' AND field = ?1
959            ORDER BY rowid
960            LIMIT ?2
961            "#,
962        )?;
963        let rows = statement
964            .query_map(params![field, limit as i64], |row| row.get::<_, String>(0))?
965            .collect::<std::result::Result<Vec<_>, _>>()?;
966        Ok(rows)
967    }
968
969    pub fn get_drawer(&self, drawer_id: &str) -> Result<Option<Drawer>, DbError> {
970        let mut statement = self.conn.prepare(&format!(
971            r#"
972            SELECT {DRAWER_SELECT_COLUMNS}
973            FROM drawers
974            WHERE id = ?1 AND deleted_at IS NULL
975            "#,
976        ))?;
977        let mut rows = statement.query_map([drawer_id], |row| {
978            drawer_from_row(row).map_err(row_decode_error)
979        })?;
980
981        match rows.next() {
982            Some(row) => Ok(Some(row?)),
983            None => Ok(None),
984        }
985    }
986
987    pub fn update_knowledge_lifecycle(
988        &self,
989        drawer_id: &str,
990        status: &KnowledgeStatus,
991        verification_refs: &[String],
992        counterexample_refs: &[String],
993    ) -> Result<bool, DbError> {
994        let affected = self.conn.execute(
995            r#"
996            UPDATE drawers
997            SET status = ?2,
998                verification_refs = ?3,
999                counterexample_refs = ?4
1000            WHERE id = ?1
1001              AND deleted_at IS NULL
1002              AND memory_kind = 'knowledge'
1003            "#,
1004            params![
1005                drawer_id,
1006                knowledge_status_as_str(status),
1007                encode_json(verification_refs)?,
1008                encode_json(counterexample_refs)?,
1009            ],
1010        )?;
1011        Ok(affected > 0)
1012    }
1013
1014    pub fn update_knowledge_anchor(
1015        &self,
1016        drawer_id: &str,
1017        anchor_kind: &AnchorKind,
1018        anchor_id: &str,
1019        parent_anchor_id: Option<&str>,
1020    ) -> Result<bool, DbError> {
1021        let affected = self.conn.execute(
1022            r#"
1023            UPDATE drawers
1024            SET anchor_kind = ?2,
1025                anchor_id = ?3,
1026                parent_anchor_id = ?4
1027            WHERE id = ?1
1028              AND deleted_at IS NULL
1029              AND memory_kind = 'knowledge'
1030            "#,
1031            params![
1032                drawer_id,
1033                anchor_kind_as_str(anchor_kind),
1034                anchor_id,
1035                parent_anchor_id,
1036            ],
1037        )?;
1038        Ok(affected > 0)
1039    }
1040
1041    pub fn insert_knowledge_card(&self, card: &KnowledgeCard) -> Result<(), DbError> {
1042        anchor::validate_anchor_domain(&card.domain, &card.anchor_kind)
1043            .map_err(|message| DbError::InvalidDrawerMetadata(message.to_string()))?;
1044
1045        self.conn.execute(
1046            r#"
1047            INSERT INTO knowledge_cards (
1048                id,
1049                statement,
1050                content,
1051                tier,
1052                status,
1053                domain,
1054                field,
1055                anchor_kind,
1056                anchor_id,
1057                parent_anchor_id,
1058                scope_constraints,
1059                trigger_hints,
1060                created_at,
1061                updated_at
1062            )
1063            VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)
1064            "#,
1065            params![
1066                card.id.as_str(),
1067                card.statement.as_str(),
1068                card.content.as_str(),
1069                knowledge_tier_as_str(&card.tier),
1070                knowledge_status_as_str(&card.status),
1071                memory_domain_as_str(&card.domain),
1072                card.field.as_str(),
1073                anchor_kind_as_str(&card.anchor_kind),
1074                card.anchor_id.as_str(),
1075                card.parent_anchor_id.as_deref(),
1076                card.scope_constraints.as_deref(),
1077                encode_optional_json(card.trigger_hints.as_ref())?,
1078                card.created_at.as_str(),
1079                card.updated_at.as_str(),
1080            ],
1081        )?;
1082        Ok(())
1083    }
1084
1085    pub fn get_knowledge_card(&self, card_id: &str) -> Result<Option<KnowledgeCard>, DbError> {
1086        let mut statement = self.conn.prepare(
1087            r#"
1088            SELECT
1089                id,
1090                statement,
1091                content,
1092                tier,
1093                status,
1094                domain,
1095                field,
1096                anchor_kind,
1097                anchor_id,
1098                parent_anchor_id,
1099                scope_constraints,
1100                trigger_hints,
1101                created_at,
1102                updated_at
1103            FROM knowledge_cards
1104            WHERE id = ?1
1105            "#,
1106        )?;
1107        let mut rows = statement.query_map([card_id], |row| {
1108            knowledge_card_from_row(row).map_err(row_decode_error)
1109        })?;
1110
1111        match rows.next() {
1112            Some(row) => Ok(Some(row?)),
1113            None => Ok(None),
1114        }
1115    }
1116
1117    pub fn list_knowledge_cards(
1118        &self,
1119        filter: &KnowledgeCardFilter,
1120    ) -> Result<Vec<KnowledgeCard>, DbError> {
1121        let tier = filter.tier.as_ref().map(knowledge_tier_as_str);
1122        let status = filter.status.as_ref().map(knowledge_status_as_str);
1123        let domain = filter.domain.as_ref().map(memory_domain_as_str);
1124        let anchor_kind = filter.anchor_kind.as_ref().map(anchor_kind_as_str);
1125
1126        let mut statement = self.conn.prepare(
1127            r#"
1128            SELECT
1129                id,
1130                statement,
1131                content,
1132                tier,
1133                status,
1134                domain,
1135                field,
1136                anchor_kind,
1137                anchor_id,
1138                parent_anchor_id,
1139                scope_constraints,
1140                trigger_hints,
1141                created_at,
1142                updated_at
1143            FROM knowledge_cards
1144            WHERE (?1 IS NULL OR tier = ?1)
1145              AND (?2 IS NULL OR status = ?2)
1146              AND (?3 IS NULL OR domain = ?3)
1147              AND (?4 IS NULL OR field = ?4)
1148              AND (?5 IS NULL OR anchor_kind = ?5)
1149              AND (?6 IS NULL OR anchor_id = ?6)
1150            ORDER BY tier, status, id
1151            "#,
1152        )?;
1153        let rows = statement
1154            .query_map(
1155                params![
1156                    tier,
1157                    status,
1158                    domain,
1159                    filter.field.as_deref(),
1160                    anchor_kind,
1161                    filter.anchor_id.as_deref(),
1162                ],
1163                |row| knowledge_card_from_row(row).map_err(row_decode_error),
1164            )?
1165            .collect::<std::result::Result<Vec<_>, _>>()?;
1166        Ok(rows)
1167    }
1168
1169    pub fn knowledge_card_count(&self) -> Result<i64, DbError> {
1170        self.conn
1171            .query_row("SELECT COUNT(*) FROM knowledge_cards", [], |row| row.get(0))
1172            .map_err(Into::into)
1173    }
1174
1175    pub fn insert_runtime_adoption_event(
1176        &self,
1177        event: &RuntimeAdoptionEvent,
1178    ) -> Result<(), DbError> {
1179        self.conn.execute(
1180            r#"
1181            INSERT INTO runtime_adoption_events (
1182                id,
1183                track,
1184                signal,
1185                feature,
1186                query,
1187                context_hash,
1188                card_id,
1189                evaluator_id,
1190                research_report_id,
1191                note,
1192                metadata,
1193                created_at
1194            )
1195            VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12)
1196            "#,
1197            params![
1198                event.id.as_str(),
1199                runtime_adoption_track_as_str(&event.track),
1200                runtime_adoption_signal_as_str(&event.signal),
1201                event.feature.as_str(),
1202                event.query.as_deref(),
1203                event.context_hash.as_deref(),
1204                event.card_id.as_deref(),
1205                event.evaluator_id.as_deref(),
1206                event.research_report_id.as_deref(),
1207                event.note.as_deref(),
1208                encode_optional_json(event.metadata.as_ref())?,
1209                event.created_at.as_str(),
1210            ],
1211        )?;
1212        Ok(())
1213    }
1214
1215    pub fn list_runtime_adoption_events(
1216        &self,
1217        filter: &RuntimeAdoptionFilter,
1218        limit: usize,
1219    ) -> Result<Vec<RuntimeAdoptionEvent>, DbError> {
1220        let track = filter.track.as_ref().map(runtime_adoption_track_as_str);
1221        let limit =
1222            i64::try_from(limit).map_err(|_| DbError::InvalidSourceType("limit".to_string()))?;
1223        let mut statement = self.conn.prepare(
1224            r#"
1225            SELECT
1226                id,
1227                track,
1228                signal,
1229                feature,
1230                query,
1231                context_hash,
1232                card_id,
1233                evaluator_id,
1234                research_report_id,
1235                note,
1236                metadata,
1237                created_at
1238            FROM runtime_adoption_events
1239            WHERE (?1 IS NULL OR track = ?1)
1240              AND (?2 IS NULL OR feature = ?2)
1241            ORDER BY created_at DESC, id DESC
1242            LIMIT ?3
1243            "#,
1244        )?;
1245        let rows = statement
1246            .query_map(params![track, filter.feature.as_deref(), limit], |row| {
1247                runtime_adoption_event_from_row(row).map_err(row_decode_error)
1248            })?
1249            .collect::<std::result::Result<Vec<_>, _>>()?;
1250        Ok(rows)
1251    }
1252
1253    pub fn list_knowledge_drawers_for_card_backfill(
1254        &self,
1255        filter: &KnowledgeCardFilter,
1256    ) -> Result<Vec<Drawer>, DbError> {
1257        let tier = filter.tier.as_ref().map(knowledge_tier_as_str);
1258        let status = filter.status.as_ref().map(knowledge_status_as_str);
1259        let domain = filter.domain.as_ref().map(memory_domain_as_str);
1260        let anchor_kind = filter.anchor_kind.as_ref().map(anchor_kind_as_str);
1261
1262        let mut statement = self.conn.prepare(&format!(
1263            r#"
1264            SELECT {DRAWER_SELECT_COLUMNS}
1265            FROM drawers
1266            WHERE deleted_at IS NULL
1267              AND memory_kind = 'knowledge'
1268              AND (?1 IS NULL OR tier = ?1)
1269              AND (?2 IS NULL OR status = ?2)
1270              AND (?3 IS NULL OR domain = ?3)
1271              AND (?4 IS NULL OR field = ?4)
1272              AND (?5 IS NULL OR anchor_kind = ?5)
1273              AND (?6 IS NULL OR anchor_id = ?6)
1274            ORDER BY id
1275            "#,
1276        ))?;
1277        let rows = statement
1278            .query_map(
1279                params![
1280                    tier,
1281                    status,
1282                    domain,
1283                    filter.field.as_deref(),
1284                    anchor_kind,
1285                    filter.anchor_id.as_deref(),
1286                ],
1287                |row| drawer_from_row(row).map_err(row_decode_error),
1288            )?
1289            .collect::<std::result::Result<Vec<_>, _>>()?;
1290        Ok(rows)
1291    }
1292
1293    pub fn update_knowledge_card(&self, card: &KnowledgeCard) -> Result<bool, DbError> {
1294        anchor::validate_anchor_domain(&card.domain, &card.anchor_kind)
1295            .map_err(|message| DbError::InvalidDrawerMetadata(message.to_string()))?;
1296
1297        let affected = self.conn.execute(
1298            r#"
1299            UPDATE knowledge_cards
1300            SET statement = ?2,
1301                content = ?3,
1302                tier = ?4,
1303                status = ?5,
1304                domain = ?6,
1305                field = ?7,
1306                anchor_kind = ?8,
1307                anchor_id = ?9,
1308                parent_anchor_id = ?10,
1309                scope_constraints = ?11,
1310                trigger_hints = ?12,
1311                updated_at = ?13
1312            WHERE id = ?1
1313            "#,
1314            params![
1315                card.id.as_str(),
1316                card.statement.as_str(),
1317                card.content.as_str(),
1318                knowledge_tier_as_str(&card.tier),
1319                knowledge_status_as_str(&card.status),
1320                memory_domain_as_str(&card.domain),
1321                card.field.as_str(),
1322                anchor_kind_as_str(&card.anchor_kind),
1323                card.anchor_id.as_str(),
1324                card.parent_anchor_id.as_deref(),
1325                card.scope_constraints.as_deref(),
1326                encode_optional_json(card.trigger_hints.as_ref())?,
1327                card.updated_at.as_str(),
1328            ],
1329        )?;
1330        Ok(affected > 0)
1331    }
1332
1333    pub fn insert_knowledge_evidence_link(
1334        &self,
1335        link: &KnowledgeEvidenceLink,
1336    ) -> Result<(), DbError> {
1337        let evidence = self.get_drawer(&link.evidence_drawer_id)?.ok_or_else(|| {
1338            DbError::InvalidDrawerMetadata(format!(
1339                "evidence drawer {} does not exist",
1340                link.evidence_drawer_id
1341            ))
1342        })?;
1343        if evidence.memory_kind != MemoryKind::Evidence {
1344            return Err(DbError::InvalidDrawerMetadata(format!(
1345                "evidence link target {} must be an evidence drawer",
1346                link.evidence_drawer_id
1347            )));
1348        }
1349
1350        self.conn.execute(
1351            r#"
1352            INSERT INTO knowledge_evidence_links (
1353                id,
1354                card_id,
1355                evidence_drawer_id,
1356                role,
1357                note,
1358                created_at
1359            )
1360            VALUES (?1, ?2, ?3, ?4, ?5, ?6)
1361            "#,
1362            params![
1363                link.id.as_str(),
1364                link.card_id.as_str(),
1365                link.evidence_drawer_id.as_str(),
1366                knowledge_evidence_role_as_str(&link.role),
1367                link.note.as_deref(),
1368                link.created_at.as_str(),
1369            ],
1370        )?;
1371        Ok(())
1372    }
1373
1374    pub fn knowledge_evidence_links(
1375        &self,
1376        card_id: &str,
1377    ) -> Result<Vec<KnowledgeEvidenceLink>, DbError> {
1378        let mut statement = self.conn.prepare(
1379            r#"
1380            SELECT id, card_id, evidence_drawer_id, role, note, created_at
1381            FROM knowledge_evidence_links
1382            WHERE card_id = ?1
1383            ORDER BY created_at, id
1384            "#,
1385        )?;
1386        let rows = statement
1387            .query_map([card_id], |row| {
1388                knowledge_evidence_link_from_row(row).map_err(row_decode_error)
1389            })?
1390            .collect::<std::result::Result<Vec<_>, _>>()?;
1391        Ok(rows)
1392    }
1393
1394    pub fn knowledge_evidence_links_for_drawer(
1395        &self,
1396        evidence_drawer_id: &str,
1397    ) -> Result<Vec<KnowledgeEvidenceLink>, DbError> {
1398        let mut statement = self.conn.prepare(
1399            r#"
1400            SELECT id, card_id, evidence_drawer_id, role, note, created_at
1401            FROM knowledge_evidence_links
1402            WHERE evidence_drawer_id = ?1
1403            ORDER BY created_at, id
1404            "#,
1405        )?;
1406        let rows = statement
1407            .query_map([evidence_drawer_id], |row| {
1408                knowledge_evidence_link_from_row(row).map_err(row_decode_error)
1409            })?
1410            .collect::<std::result::Result<Vec<_>, _>>()?;
1411        Ok(rows)
1412    }
1413
1414    pub fn append_knowledge_event(&self, event: &KnowledgeCardEvent) -> Result<(), DbError> {
1415        self.conn.execute(
1416            r#"
1417            INSERT INTO knowledge_events (
1418                id,
1419                card_id,
1420                event_type,
1421                from_status,
1422                to_status,
1423                reason,
1424                actor,
1425                metadata,
1426                created_at
1427            )
1428            VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)
1429            "#,
1430            params![
1431                event.id.as_str(),
1432                event.card_id.as_str(),
1433                knowledge_event_type_as_str(&event.event_type),
1434                event.from_status.as_ref().map(knowledge_status_as_str),
1435                event.to_status.as_ref().map(knowledge_status_as_str),
1436                event.reason.as_str(),
1437                event.actor.as_deref(),
1438                encode_optional_json(event.metadata.as_ref())?,
1439                event.created_at.as_str(),
1440            ],
1441        )?;
1442        Ok(())
1443    }
1444
1445    pub fn knowledge_events(&self, card_id: &str) -> Result<Vec<KnowledgeCardEvent>, DbError> {
1446        let mut statement = self.conn.prepare(
1447            r#"
1448            SELECT
1449                id,
1450                card_id,
1451                event_type,
1452                from_status,
1453                to_status,
1454                reason,
1455                actor,
1456                metadata,
1457                created_at
1458            FROM knowledge_events
1459            WHERE card_id = ?1
1460            ORDER BY created_at, id
1461            "#,
1462        )?;
1463        let rows = statement
1464            .query_map([card_id], |row| {
1465                knowledge_event_from_row(row).map_err(row_decode_error)
1466            })?
1467            .collect::<std::result::Result<Vec<_>, _>>()?;
1468        Ok(rows)
1469    }
1470
1471    pub fn neighbor_chunks(
1472        &self,
1473        source_file: &str,
1474        wing: &str,
1475        room: Option<&str>,
1476        chunk_index: i64,
1477    ) -> Result<ChunkNeighbors, DbError> {
1478        let prev_index = chunk_index - 1;
1479        let next_index = chunk_index + 1;
1480        let sql = r#"
1481            SELECT id, content, chunk_index
1482            FROM drawers
1483            WHERE deleted_at IS NULL
1484              AND source_file = ?1
1485              AND wing = ?2
1486              AND ((?3 IS NULL AND room IS NULL) OR (?3 IS NOT NULL AND room = ?3))
1487              AND chunk_index IN (?4, ?5)
1488            ORDER BY chunk_index, id
1489            "#;
1490        let mut statement = self.conn.prepare(sql)?;
1491        let mut rows = statement.query(params![source_file, wing, room, prev_index, next_index])?;
1492        let mut neighbors = ChunkNeighbors {
1493            prev: None,
1494            next: None,
1495        };
1496
1497        while let Some(row) = rows.next()? {
1498            let row_index = row.get::<_, i64>(2)?;
1499            let Ok(chunk_index) = u32::try_from(row_index) else {
1500                continue;
1501            };
1502            let chunk = NeighborChunk {
1503                drawer_id: row.get(0)?,
1504                content: row.get(1)?,
1505                chunk_index,
1506            };
1507            if row_index == prev_index && neighbors.prev.is_none() {
1508                neighbors.prev = Some(chunk);
1509            } else if row_index == next_index && neighbors.next.is_none() {
1510                neighbors.next = Some(chunk);
1511            }
1512        }
1513
1514        Ok(neighbors)
1515    }
1516
1517    pub fn soft_delete_drawer(&self, drawer_id: &str) -> Result<bool, DbError> {
1518        let timestamp = current_timestamp();
1519        let affected = self.conn.execute(
1520            "UPDATE drawers SET deleted_at = ?1 WHERE id = ?2 AND deleted_at IS NULL",
1521            params![timestamp, drawer_id],
1522        )?;
1523        Ok(affected > 0)
1524    }
1525
1526    pub fn purge_deleted(&self, before: Option<&str>) -> Result<u64, DbError> {
1527        // First collect IDs to purge, then delete from both tables
1528        let ids: Vec<String> = if let Some(before) = before {
1529            let mut stmt = self.conn.prepare(
1530                "SELECT id FROM drawers WHERE deleted_at IS NOT NULL AND deleted_at < ?1",
1531            )?;
1532            stmt.query_map([before], |row| row.get(0))?
1533                .collect::<std::result::Result<Vec<_>, _>>()?
1534        } else {
1535            let mut stmt = self
1536                .conn
1537                .prepare("SELECT id FROM drawers WHERE deleted_at IS NOT NULL")?;
1538            stmt.query_map([], |row| row.get(0))?
1539                .collect::<std::result::Result<Vec<_>, _>>()?
1540        };
1541
1542        if ids.is_empty() {
1543            return Ok(0);
1544        }
1545
1546        // Check if drawer_vectors table exists (lazy-created)
1547        let vectors_exist: bool = self.conn.query_row(
1548            "SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type='table' AND name='drawer_vectors')",
1549            [],
1550            |row| row.get(0),
1551        )?;
1552
1553        // Wrap the purge in a single transaction. Clearing the
1554        // triples.source_drawer FK and deleting the drawer must be atomic:
1555        // another RESTRICT FK (e.g. knowledge_evidence_links.evidence_drawer_id)
1556        // can block `DELETE FROM drawers`, and without a transaction the prior
1557        // `UPDATE triples SET source_drawer = NULL` would have already committed
1558        // — silently dropping provenance for a drawer that was not purged.
1559        self.conn.execute_batch("BEGIN IMMEDIATE;")?;
1560        let result = (|| -> Result<u64, DbError> {
1561            for id in &ids {
1562                if vectors_exist {
1563                    self.conn
1564                        .execute("DELETE FROM drawer_vectors WHERE id = ?1", [id])?;
1565                }
1566                // Clear the triples.source_drawer FK (RESTRICT) before the hard
1567                // delete so purging a soft-deleted drawer referenced by a KG
1568                // triple does not fail with a FOREIGN KEY constraint error.
1569                self.conn.execute(
1570                    "UPDATE triples SET source_drawer = NULL WHERE source_drawer = ?1",
1571                    [id],
1572                )?;
1573                self.conn
1574                    .execute("DELETE FROM drawers WHERE id = ?1", [id])?;
1575            }
1576            Ok(ids.len() as u64)
1577        })();
1578
1579        match result {
1580            Ok(count) => {
1581                self.conn.execute_batch("COMMIT;")?;
1582                Ok(count)
1583            }
1584            Err(error) => {
1585                let _ = self.conn.execute_batch("ROLLBACK;");
1586                Err(error)
1587            }
1588        }
1589    }
1590
1591    pub fn deleted_drawer_count(&self) -> Result<i64, DbError> {
1592        Ok(self.conn.query_row(
1593            "SELECT COUNT(*) FROM drawers WHERE deleted_at IS NOT NULL",
1594            [],
1595            |row| row.get(0),
1596        )?)
1597    }
1598
1599    // --- FTS5 BM25 search ---
1600
1601    pub fn search_fts(
1602        &self,
1603        query: &str,
1604        wing: Option<&str>,
1605        room: Option<&str>,
1606        limit: usize,
1607    ) -> Result<Vec<(String, f64)>, DbError> {
1608        let Some(match_query) = build_fts_match_query(query) else {
1609            return Ok(Vec::new());
1610        };
1611        let limit =
1612            i64::try_from(limit).map_err(|_| DbError::InvalidSourceType("limit".to_string()))?;
1613        let mut stmt = self.conn.prepare(
1614            r#"
1615            SELECT d.id, fts.rank
1616            FROM drawers_fts fts
1617            JOIN drawers d ON d.rowid = fts.rowid
1618            WHERE drawers_fts MATCH ?1
1619              AND d.deleted_at IS NULL
1620              AND (?2 IS NULL OR d.wing = ?2)
1621              AND (?3 IS NULL OR d.room = ?3)
1622            ORDER BY fts.rank
1623            LIMIT ?4
1624            "#,
1625        )?;
1626        let rows = stmt
1627            .query_map((match_query.as_str(), wing, room, limit), |row| {
1628                Ok((row.get::<_, String>(0)?, row.get::<_, f64>(1)?))
1629            })?
1630            .collect::<std::result::Result<Vec<_>, _>>()?;
1631        Ok(rows)
1632    }
1633
1634    // --- Triples (Knowledge Graph) ---
1635
1636    pub fn insert_triple(&self, triple: &Triple) -> Result<(), DbError> {
1637        self.conn.execute(
1638            r#"
1639            INSERT OR REPLACE INTO triples (id, subject, predicate, object, valid_from, valid_to, confidence, source_drawer)
1640            VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)
1641            "#,
1642            params![
1643                triple.id,
1644                triple.subject,
1645                triple.predicate,
1646                triple.object,
1647                triple.valid_from,
1648                triple.valid_to,
1649                triple.confidence,
1650                triple.source_drawer,
1651            ],
1652        )?;
1653        Ok(())
1654    }
1655
1656    pub fn query_triples(
1657        &self,
1658        subject: Option<&str>,
1659        predicate: Option<&str>,
1660        object: Option<&str>,
1661        active_only: bool,
1662    ) -> Result<Vec<Triple>, DbError> {
1663        let active_clause = if active_only {
1664            "AND (valid_to IS NULL OR valid_to > strftime('%s', 'now'))"
1665        } else {
1666            ""
1667        };
1668        let sql = format!(
1669            r#"
1670            SELECT id, subject, predicate, object, valid_from, valid_to, confidence, source_drawer
1671            FROM triples
1672            WHERE (?1 IS NULL OR subject = ?1)
1673              AND (?2 IS NULL OR predicate = ?2)
1674              AND (?3 IS NULL OR object = ?3)
1675              {active_clause}
1676            ORDER BY confidence DESC, id
1677            "#
1678        );
1679        let mut stmt = self.conn.prepare(&sql)?;
1680        let rows = stmt
1681            .query_map((subject, predicate, object), |row| {
1682                Ok(Triple {
1683                    id: row.get(0)?,
1684                    subject: row.get(1)?,
1685                    predicate: row.get(2)?,
1686                    object: row.get(3)?,
1687                    valid_from: row.get(4)?,
1688                    valid_to: row.get(5)?,
1689                    confidence: row.get(6)?,
1690                    source_drawer: row.get(7)?,
1691                })
1692            })?
1693            .collect::<std::result::Result<Vec<_>, _>>()?;
1694        Ok(rows)
1695    }
1696
1697    pub fn invalidate_triple(&self, triple_id: &str) -> Result<bool, DbError> {
1698        let timestamp = current_timestamp();
1699        let affected = self.conn.execute(
1700            "UPDATE triples SET valid_to = ?1 WHERE id = ?2 AND valid_to IS NULL",
1701            params![timestamp, triple_id],
1702        )?;
1703        Ok(affected > 0)
1704    }
1705
1706    pub fn triple_count(&self) -> Result<i64, DbError> {
1707        Ok(self
1708            .conn
1709            .query_row("SELECT COUNT(*) FROM triples", [], |row| row.get(0))?)
1710    }
1711
1712    pub fn timeline_for_entity(&self, entity: &str) -> Result<Vec<Triple>, DbError> {
1713        let mut stmt = self.conn.prepare(
1714            r#"
1715            SELECT id, subject, predicate, object, valid_from, valid_to, confidence, source_drawer
1716            FROM triples
1717            WHERE subject = ?1 OR object = ?1
1718            ORDER BY COALESCE(valid_from, '0') ASC, id ASC
1719            "#,
1720        )?;
1721        let rows = stmt
1722            .query_map([entity], |row| {
1723                Ok(Triple {
1724                    id: row.get(0)?,
1725                    subject: row.get(1)?,
1726                    predicate: row.get(2)?,
1727                    object: row.get(3)?,
1728                    valid_from: row.get(4)?,
1729                    valid_to: row.get(5)?,
1730                    confidence: row.get(6)?,
1731                    source_drawer: row.get(7)?,
1732                })
1733            })?
1734            .collect::<std::result::Result<Vec<_>, _>>()?;
1735        Ok(rows)
1736    }
1737
1738    pub fn triple_stats(&self) -> Result<TripleStats, DbError> {
1739        let total: i64 = self
1740            .conn
1741            .query_row("SELECT COUNT(*) FROM triples", [], |row| row.get(0))?;
1742        let active: i64 = self.conn.query_row(
1743            "SELECT COUNT(*) FROM triples WHERE valid_to IS NULL",
1744            [],
1745            |row| row.get(0),
1746        )?;
1747        let expired = total - active;
1748        let entities: i64 = self.conn.query_row(
1749            r#"
1750            SELECT COUNT(DISTINCT entity) FROM (
1751                SELECT subject AS entity FROM triples
1752                UNION
1753                SELECT object AS entity FROM triples
1754            )
1755            "#,
1756            [],
1757            |row| row.get(0),
1758        )?;
1759        let mut top_predicates_stmt = self.conn.prepare(
1760            "SELECT predicate, COUNT(*) as cnt FROM triples GROUP BY predicate ORDER BY cnt DESC LIMIT 5",
1761        )?;
1762        let top_predicates = top_predicates_stmt
1763            .query_map([], |row| {
1764                Ok((row.get::<_, String>(0)?, row.get::<_, i64>(1)?))
1765            })?
1766            .collect::<std::result::Result<Vec<_>, _>>()?;
1767        Ok(TripleStats {
1768            total,
1769            active,
1770            expired,
1771            entities,
1772            top_predicates,
1773        })
1774    }
1775
1776    // --- Tunnels (cross-Wing discovery) ---
1777
1778    pub fn find_tunnels(&self) -> Result<Vec<(String, Vec<String>)>, DbError> {
1779        let mut stmt = self.conn.prepare(
1780            r#"
1781            SELECT room, GROUP_CONCAT(DISTINCT wing) as wings
1782            FROM drawers
1783            WHERE deleted_at IS NULL AND room IS NOT NULL AND room != ''
1784            GROUP BY room
1785            HAVING COUNT(DISTINCT wing) > 1
1786            ORDER BY room
1787            "#,
1788        )?;
1789        let rows = stmt
1790            .query_map([], |row| {
1791                let room: String = row.get(0)?;
1792                let wings_csv: String = row.get(1)?;
1793                Ok((room, wings_csv))
1794            })?
1795            .collect::<std::result::Result<Vec<_>, _>>()?;
1796        Ok(rows
1797            .into_iter()
1798            .map(|(room, wings_csv)| {
1799                let wings = wings_csv.split(',').map(ToOwned::to_owned).collect();
1800                (room, wings)
1801            })
1802            .collect())
1803    }
1804
1805    pub fn create_tunnel(
1806        &self,
1807        left: &TunnelEndpoint,
1808        right: &TunnelEndpoint,
1809        label: &str,
1810        created_by: Option<&str>,
1811    ) -> Result<ExplicitTunnel, DbError> {
1812        let left = normalize_tunnel_endpoint(left)?;
1813        let right = normalize_tunnel_endpoint(right)?;
1814        let label = label.trim();
1815        if label.is_empty() {
1816            return Err(DbError::InvalidTunnel("label is required".to_string()));
1817        }
1818        if left == right {
1819            return Err(DbError::InvalidTunnel(
1820                "self-link is not allowed".to_string(),
1821            ));
1822        }
1823
1824        let id = build_tunnel_id(&left, &right);
1825        let created_at = current_timestamp();
1826        let created_by = created_by
1827            .map(str::trim)
1828            .filter(|value| !value.is_empty())
1829            .map(ToOwned::to_owned);
1830        self.conn.execute(
1831            r#"
1832            INSERT INTO tunnels (
1833                id, left_wing, left_room, right_wing, right_room,
1834                label, created_at, created_by, deleted_at
1835            )
1836            VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, NULL)
1837            ON CONFLICT(id) DO UPDATE SET
1838                label = CASE
1839                    WHEN tunnels.deleted_at IS NOT NULL THEN excluded.label
1840                    ELSE tunnels.label
1841                END,
1842                created_at = CASE
1843                    WHEN tunnels.deleted_at IS NOT NULL THEN excluded.created_at
1844                    ELSE tunnels.created_at
1845                END,
1846                created_by = CASE
1847                    WHEN tunnels.deleted_at IS NOT NULL THEN excluded.created_by
1848                    ELSE tunnels.created_by
1849                END,
1850                deleted_at = NULL
1851            "#,
1852            params![
1853                id, left.wing, left.room, right.wing, right.room, label, created_at, created_by,
1854            ],
1855        )?;
1856
1857        self.get_explicit_tunnel(&id)?
1858            .ok_or_else(|| DbError::InvalidTunnel(format!("failed to create tunnel {id}")))
1859    }
1860
1861    pub fn list_explicit_tunnels(
1862        &self,
1863        wing: Option<&str>,
1864    ) -> Result<Vec<ExplicitTunnel>, DbError> {
1865        let wing = wing.map(str::trim).filter(|value| !value.is_empty());
1866        let mut statement = self.conn.prepare(
1867            r#"
1868            SELECT id, left_wing, left_room, right_wing, right_room,
1869                   label, created_at, created_by, deleted_at
1870            FROM tunnels
1871            WHERE deleted_at IS NULL
1872              AND (?1 IS NULL OR left_wing = ?1 OR right_wing = ?1)
1873            ORDER BY left_wing, left_room, right_wing, right_room, id
1874            "#,
1875        )?;
1876        let rows = statement
1877            .query_map([wing], explicit_tunnel_from_row)?
1878            .collect::<std::result::Result<Vec<_>, _>>()?;
1879        Ok(rows)
1880    }
1881
1882    pub fn delete_explicit_tunnel(&self, tunnel_id: &str) -> Result<bool, DbError> {
1883        let timestamp = current_timestamp();
1884        let affected = self.conn.execute(
1885            "UPDATE tunnels SET deleted_at = ?1 WHERE id = ?2 AND deleted_at IS NULL",
1886            params![timestamp, tunnel_id],
1887        )?;
1888        Ok(affected > 0)
1889    }
1890
1891    pub fn follow_explicit_tunnels(
1892        &self,
1893        from: &TunnelEndpoint,
1894        max_hops: u8,
1895    ) -> Result<Vec<TunnelFollowResult>, DbError> {
1896        if !(1..=2).contains(&max_hops) {
1897            return Err(DbError::InvalidTunnel(
1898                "max_hops must be 1 or 2".to_string(),
1899            ));
1900        }
1901
1902        let from = normalize_tunnel_endpoint(from)?;
1903        let tunnels = self.list_explicit_tunnels(None)?;
1904        let mut visited = BTreeSet::from([from.clone()]);
1905        let mut queue = VecDeque::from([(from, 0_u8)]);
1906        let mut results = Vec::new();
1907
1908        while let Some((current, hop)) = queue.pop_front() {
1909            if hop >= max_hops {
1910                continue;
1911            }
1912            let next_hop = hop + 1;
1913            for tunnel in &tunnels {
1914                let neighbor = if tunnel.left == current {
1915                    Some(tunnel.right.clone())
1916                } else if tunnel.right == current {
1917                    Some(tunnel.left.clone())
1918                } else {
1919                    None
1920                };
1921                let Some(neighbor) = neighbor else {
1922                    continue;
1923                };
1924                if !visited.insert(neighbor.clone()) {
1925                    continue;
1926                }
1927                results.push(TunnelFollowResult {
1928                    endpoint: neighbor.clone(),
1929                    via_tunnel_id: tunnel.id.clone(),
1930                    hop: next_hop,
1931                });
1932                queue.push_back((neighbor, next_hop));
1933            }
1934        }
1935
1936        results.sort_by(|left, right| {
1937            left.hop
1938                .cmp(&right.hop)
1939                .then_with(|| left.endpoint.cmp(&right.endpoint))
1940                .then_with(|| left.via_tunnel_id.cmp(&right.via_tunnel_id))
1941        });
1942        Ok(results)
1943    }
1944
1945    pub fn explicit_tunnel_hints(
1946        &self,
1947        wing: &str,
1948        room: Option<&str>,
1949    ) -> Result<Vec<String>, DbError> {
1950        let endpoint = TunnelEndpoint {
1951            wing: wing.to_string(),
1952            room: room.map(ToOwned::to_owned),
1953        };
1954        let hints = self
1955            .follow_explicit_tunnels(&endpoint, 1)?
1956            .into_iter()
1957            .map(|result| format_tunnel_endpoint(&result.endpoint))
1958            .collect::<BTreeSet<_>>()
1959            .into_iter()
1960            .collect();
1961        Ok(hints)
1962    }
1963
1964    fn get_explicit_tunnel(&self, tunnel_id: &str) -> Result<Option<ExplicitTunnel>, DbError> {
1965        let mut statement = self.conn.prepare(
1966            r#"
1967            SELECT id, left_wing, left_room, right_wing, right_room,
1968                   label, created_at, created_by, deleted_at
1969            FROM tunnels
1970            WHERE id = ?1 AND deleted_at IS NULL
1971            "#,
1972        )?;
1973        let mut rows = statement.query_map([tunnel_id], explicit_tunnel_from_row)?;
1974        match rows.next() {
1975            Some(row) => Ok(Some(row?)),
1976            None => Ok(None),
1977        }
1978    }
1979
1980    // --- Embedding dimension management ---
1981
1982    /// Returns the current embedding dimension from the vec0 table, or None if the table is empty.
1983    pub fn embedding_dim(&self) -> Result<Option<usize>, DbError> {
1984        // sqlite-vec stores dimension in table schema; probe by checking a row
1985        let result: std::result::Result<i64, _> = self.conn.query_row(
1986            "SELECT vec_length(embedding) FROM drawer_vectors LIMIT 1",
1987            [],
1988            |row| row.get(0),
1989        );
1990        match result {
1991            Ok(dim) => Ok(Some(dim as usize)),
1992            Err(rusqlite::Error::QueryReturnedNoRows) => Ok(None),
1993            Err(e) => Err(DbError::Sqlite(e)),
1994        }
1995    }
1996
1997    /// Drop and recreate the drawer_vectors table with the specified dimension.
1998    /// All existing vectors are lost — caller must re-embed after this.
1999    pub fn recreate_vectors_table(&self, dim: usize) -> Result<(), DbError> {
2000        self.conn.execute_batch(&format!(
2001            r#"
2002            DROP TABLE IF EXISTS drawer_vectors;
2003            CREATE VIRTUAL TABLE drawer_vectors USING vec0(
2004                id TEXT PRIMARY KEY,
2005                embedding FLOAT[{dim}]
2006            );
2007            "#
2008        ))?;
2009        Ok(())
2010    }
2011
2012    /// Returns all active (non-deleted) drawer IDs and their content for re-embedding.
2013    pub fn all_active_drawers(&self) -> Result<Vec<(String, String)>, DbError> {
2014        let mut stmt = self
2015            .conn
2016            .prepare("SELECT id, content FROM drawers WHERE deleted_at IS NULL ORDER BY id")?;
2017        let rows = stmt
2018            .query_map([], |row| {
2019                Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
2020            })?
2021            .collect::<std::result::Result<Vec<_>, _>>()?;
2022        Ok(rows)
2023    }
2024
2025    pub fn database_size_bytes(&self) -> Result<u64, DbError> {
2026        fs::metadata(&self.path)
2027            .map(|metadata| metadata.len())
2028            .map_err(|source| DbError::Metadata {
2029                path: self.path.clone(),
2030                source,
2031            })
2032    }
2033
2034    pub fn schema_version(&self) -> Result<u32, DbError> {
2035        read_user_version(&self.conn)
2036    }
2037}
2038
2039fn apply_migrations(conn: &Connection) -> Result<(), DbError> {
2040    let current_version = read_user_version(conn)?;
2041    if current_version > CURRENT_SCHEMA_VERSION {
2042        return Err(DbError::UnsupportedSchemaVersion {
2043            current: current_version,
2044            supported: CURRENT_SCHEMA_VERSION,
2045        });
2046    }
2047
2048    for migration in migrations()
2049        .iter()
2050        .filter(|migration| migration.version > current_version)
2051    {
2052        apply_migration_atomic(conn, migration)?;
2053    }
2054
2055    Ok(())
2056}
2057
2058fn apply_migration_atomic(conn: &Connection, migration: &Migration) -> Result<(), DbError> {
2059    conn.execute_batch("BEGIN IMMEDIATE;")?;
2060    if let Err(error) = (|| -> Result<(), DbError> {
2061        conn.execute_batch(migration.sql)?;
2062        set_user_version(conn, migration.version)?;
2063        conn.execute_batch("COMMIT;")?;
2064        Ok(())
2065    })() {
2066        let _ = conn.execute_batch("ROLLBACK;");
2067        return Err(error);
2068    }
2069    Ok(())
2070}
2071
2072fn read_user_version(conn: &Connection) -> Result<u32, DbError> {
2073    let version = conn.query_row("PRAGMA user_version", [], |row| row.get::<_, u32>(0))?;
2074    Ok(version)
2075}
2076
2077fn set_user_version(conn: &Connection, version: u32) -> Result<(), DbError> {
2078    conn.execute_batch(&format!("PRAGMA user_version = {version};"))?;
2079    Ok(())
2080}
2081
2082fn normalize_tunnel_endpoint(endpoint: &TunnelEndpoint) -> Result<TunnelEndpoint, DbError> {
2083    let wing = endpoint.wing.trim();
2084    if wing.is_empty() {
2085        return Err(DbError::InvalidTunnel(
2086            "endpoint wing is required".to_string(),
2087        ));
2088    }
2089    let room = endpoint
2090        .room
2091        .as_deref()
2092        .map(str::trim)
2093        .filter(|room| !room.is_empty())
2094        .map(ToOwned::to_owned);
2095    Ok(TunnelEndpoint {
2096        wing: wing.to_string(),
2097        room,
2098    })
2099}
2100
2101fn explicit_tunnel_from_row(row: &Row<'_>) -> rusqlite::Result<ExplicitTunnel> {
2102    Ok(ExplicitTunnel {
2103        id: row.get(0)?,
2104        left: TunnelEndpoint {
2105            wing: row.get(1)?,
2106            room: row.get(2)?,
2107        },
2108        right: TunnelEndpoint {
2109            wing: row.get(3)?,
2110            room: row.get(4)?,
2111        },
2112        label: row.get(5)?,
2113        created_at: row.get(6)?,
2114        created_by: row.get(7)?,
2115        deleted_at: row.get(8)?,
2116    })
2117}
2118
2119fn reindex_source_from_row(row: &Row<'_>) -> rusqlite::Result<ReindexSource> {
2120    let drawer_count = row.get::<_, i64>(3)?;
2121    Ok(ReindexSource {
2122        source_file: row.get(0)?,
2123        wing: row.get(1)?,
2124        room: row.get(2)?,
2125        drawer_count: drawer_count as u64,
2126    })
2127}
2128
2129const V2_MIGRATION_SQL: &str = r#"
2130ALTER TABLE drawers ADD COLUMN deleted_at TEXT;
2131CREATE INDEX IF NOT EXISTS idx_drawers_deleted_at ON drawers(deleted_at);
2132"#;
2133
2134const V3_MIGRATION_SQL: &str = r#"
2135CREATE VIRTUAL TABLE IF NOT EXISTS drawers_fts USING fts5(
2136    content,
2137    content='drawers',
2138    content_rowid='rowid'
2139);
2140
2141-- Populate FTS from existing drawers (excluding soft-deleted)
2142INSERT INTO drawers_fts(rowid, content)
2143    SELECT rowid, content FROM drawers WHERE deleted_at IS NULL;
2144
2145-- Keep FTS in sync: INSERT trigger
2146CREATE TRIGGER IF NOT EXISTS drawers_ai AFTER INSERT ON drawers BEGIN
2147    INSERT INTO drawers_fts(rowid, content) VALUES (new.rowid, new.content);
2148END;
2149
2150-- Keep FTS in sync: soft-delete (UPDATE deleted_at) removes from FTS
2151CREATE TRIGGER IF NOT EXISTS drawers_au_softdelete AFTER UPDATE OF deleted_at ON drawers
2152    WHEN new.deleted_at IS NOT NULL AND old.deleted_at IS NULL BEGIN
2153    INSERT INTO drawers_fts(drawers_fts, rowid, content) VALUES ('delete', old.rowid, old.content);
2154END;
2155
2156-- No DELETE trigger on drawers — soft-deleted rows are already removed from FTS
2157-- by the UPDATE trigger above. Physical DELETE (purge) skips FTS because the
2158-- entry is already gone.
2159"#;
2160
2161const V4_MIGRATION_SQL: &str = r#"
2162ALTER TABLE drawers ADD COLUMN importance INTEGER DEFAULT 0;
2163"#;
2164
2165const V5_MIGRATION_SQL: &str = r#"
2166ALTER TABLE drawers ADD COLUMN memory_kind TEXT NOT NULL CHECK(memory_kind IN ('evidence', 'knowledge')) DEFAULT 'evidence';
2167ALTER TABLE drawers ADD COLUMN domain TEXT NOT NULL CHECK(domain IN ('project', 'agent', 'skill', 'global')) DEFAULT 'project';
2168ALTER TABLE drawers ADD COLUMN field TEXT NOT NULL DEFAULT 'general';
2169ALTER TABLE drawers ADD COLUMN anchor_kind TEXT NOT NULL CHECK(anchor_kind IN ('global', 'repo', 'worktree')) DEFAULT 'repo';
2170ALTER TABLE drawers ADD COLUMN anchor_id TEXT NOT NULL DEFAULT 'repo://legacy';
2171ALTER TABLE drawers ADD COLUMN parent_anchor_id TEXT;
2172ALTER TABLE drawers ADD COLUMN provenance TEXT CHECK(provenance IN ('runtime', 'research', 'human'));
2173ALTER TABLE drawers ADD COLUMN statement TEXT;
2174ALTER TABLE drawers ADD COLUMN tier TEXT CHECK(tier IN ('qi', 'shu', 'dao_ren', 'dao_tian'));
2175ALTER TABLE drawers ADD COLUMN status TEXT CHECK(status IN ('candidate', 'promoted', 'canonical', 'demoted', 'retired'));
2176ALTER TABLE drawers ADD COLUMN supporting_refs TEXT NOT NULL DEFAULT '[]';
2177ALTER TABLE drawers ADD COLUMN counterexample_refs TEXT NOT NULL DEFAULT '[]';
2178ALTER TABLE drawers ADD COLUMN teaching_refs TEXT NOT NULL DEFAULT '[]';
2179ALTER TABLE drawers ADD COLUMN verification_refs TEXT NOT NULL DEFAULT '[]';
2180ALTER TABLE drawers ADD COLUMN scope_constraints TEXT;
2181ALTER TABLE drawers ADD COLUMN trigger_hints TEXT;
2182
2183UPDATE drawers
2184SET memory_kind = 'evidence',
2185    domain = 'project',
2186    field = 'general',
2187    anchor_kind = 'repo',
2188    anchor_id = 'repo://legacy',
2189    parent_anchor_id = NULL,
2190    provenance = CASE source_type
2191        WHEN 'project' THEN 'research'
2192        WHEN 'conversation' THEN 'human'
2193        WHEN 'manual' THEN 'human'
2194        ELSE NULL
2195    END
2196WHERE memory_kind = 'evidence'
2197  AND domain = 'project'
2198  AND field = 'general'
2199  AND anchor_kind = 'repo'
2200  AND anchor_id = 'repo://legacy'
2201  AND parent_anchor_id IS NULL
2202  AND provenance IS NULL;
2203"#;
2204
2205const V6_MIGRATION_SQL: &str = r#"
2206CREATE TABLE IF NOT EXISTS tunnels (
2207    id TEXT PRIMARY KEY,
2208    left_wing TEXT NOT NULL,
2209    left_room TEXT,
2210    right_wing TEXT NOT NULL,
2211    right_room TEXT,
2212    label TEXT NOT NULL,
2213    created_at TEXT NOT NULL,
2214    created_by TEXT,
2215    deleted_at TEXT
2216);
2217
2218CREATE INDEX IF NOT EXISTS idx_tunnels_left
2219    ON tunnels(left_wing, left_room)
2220    WHERE deleted_at IS NULL;
2221
2222CREATE INDEX IF NOT EXISTS idx_tunnels_right
2223    ON tunnels(right_wing, right_room)
2224    WHERE deleted_at IS NULL;
2225"#;
2226
2227const V7_MIGRATION_SQL: &str = r#"
2228ALTER TABLE drawers ADD COLUMN normalize_version INTEGER NOT NULL DEFAULT 1;
2229
2230CREATE INDEX IF NOT EXISTS idx_drawers_normalize_version
2231    ON drawers(normalize_version)
2232    WHERE deleted_at IS NULL;
2233"#;
2234
2235const V8_MIGRATION_SQL: &str = r#"
2236CREATE TABLE IF NOT EXISTS knowledge_cards (
2237    id TEXT PRIMARY KEY,
2238    statement TEXT NOT NULL,
2239    content TEXT NOT NULL,
2240    tier TEXT NOT NULL CHECK(tier IN ('qi', 'shu', 'dao_ren', 'dao_tian')),
2241    status TEXT NOT NULL CHECK(status IN ('candidate', 'promoted', 'canonical', 'demoted', 'retired')),
2242    domain TEXT NOT NULL CHECK(domain IN ('project', 'agent', 'skill', 'global')),
2243    field TEXT NOT NULL DEFAULT 'general',
2244    anchor_kind TEXT NOT NULL CHECK(anchor_kind IN ('global', 'repo', 'worktree')),
2245    anchor_id TEXT NOT NULL,
2246    parent_anchor_id TEXT,
2247    scope_constraints TEXT,
2248    trigger_hints TEXT,
2249    created_at TEXT NOT NULL,
2250    updated_at TEXT NOT NULL
2251);
2252
2253CREATE TABLE IF NOT EXISTS knowledge_evidence_links (
2254    id TEXT PRIMARY KEY,
2255    card_id TEXT NOT NULL,
2256    evidence_drawer_id TEXT NOT NULL,
2257    role TEXT NOT NULL CHECK(role IN ('supporting', 'verification', 'counterexample', 'teaching')),
2258    note TEXT,
2259    created_at TEXT NOT NULL,
2260    UNIQUE(card_id, evidence_drawer_id, role),
2261    FOREIGN KEY(card_id) REFERENCES knowledge_cards(id) ON DELETE RESTRICT,
2262    FOREIGN KEY(evidence_drawer_id) REFERENCES drawers(id) ON DELETE RESTRICT
2263);
2264
2265CREATE TABLE IF NOT EXISTS knowledge_events (
2266    id TEXT PRIMARY KEY,
2267    card_id TEXT NOT NULL,
2268    event_type TEXT NOT NULL CHECK(event_type IN ('created', 'promoted', 'demoted', 'retired', 'linked', 'unlinked', 'updated', 'published_anchor')),
2269    from_status TEXT,
2270    to_status TEXT,
2271    reason TEXT NOT NULL,
2272    actor TEXT,
2273    metadata TEXT,
2274    created_at TEXT NOT NULL,
2275    FOREIGN KEY(card_id) REFERENCES knowledge_cards(id) ON DELETE RESTRICT
2276);
2277
2278CREATE INDEX IF NOT EXISTS idx_knowledge_cards_tier_status
2279    ON knowledge_cards(tier, status);
2280
2281CREATE INDEX IF NOT EXISTS idx_knowledge_cards_domain_field
2282    ON knowledge_cards(domain, field);
2283
2284CREATE INDEX IF NOT EXISTS idx_knowledge_cards_anchor
2285    ON knowledge_cards(anchor_kind, anchor_id);
2286
2287CREATE INDEX IF NOT EXISTS idx_knowledge_evidence_links_card
2288    ON knowledge_evidence_links(card_id);
2289
2290CREATE INDEX IF NOT EXISTS idx_knowledge_evidence_links_evidence
2291    ON knowledge_evidence_links(evidence_drawer_id);
2292
2293CREATE INDEX IF NOT EXISTS idx_knowledge_events_card_created_at
2294    ON knowledge_events(card_id, created_at);
2295
2296CREATE TRIGGER IF NOT EXISTS knowledge_events_no_update
2297BEFORE UPDATE ON knowledge_events
2298BEGIN
2299    SELECT RAISE(ABORT, 'knowledge_events are append-only');
2300END;
2301
2302CREATE TRIGGER IF NOT EXISTS knowledge_events_no_delete
2303BEFORE DELETE ON knowledge_events
2304BEGIN
2305    SELECT RAISE(ABORT, 'knowledge_events are append-only');
2306END;
2307"#;
2308
2309const V9_MIGRATION_SQL: &str = r#"
2310CREATE TABLE IF NOT EXISTS runtime_adoption_events (
2311    id TEXT PRIMARY KEY,
2312    track TEXT NOT NULL CHECK(track IN ('runtime_adoption', 'card_context', 'card_embedding', 'evaluator', 'research_adapter')),
2313    signal TEXT NOT NULL CHECK(signal IN ('used', 'accepted', 'rejected', 'miss', 'rollback', 'contradiction', 'neutral')),
2314    feature TEXT NOT NULL,
2315    query TEXT,
2316    context_hash TEXT,
2317    card_id TEXT,
2318    evaluator_id TEXT,
2319    research_report_id TEXT,
2320    note TEXT,
2321    metadata TEXT,
2322    created_at TEXT NOT NULL,
2323    FOREIGN KEY(card_id) REFERENCES knowledge_cards(id) ON DELETE SET NULL
2324);
2325
2326CREATE INDEX IF NOT EXISTS idx_runtime_adoption_events_track_created_at
2327    ON runtime_adoption_events(track, created_at);
2328
2329CREATE INDEX IF NOT EXISTS idx_runtime_adoption_events_feature
2330    ON runtime_adoption_events(feature);
2331
2332CREATE INDEX IF NOT EXISTS idx_runtime_adoption_events_signal
2333    ON runtime_adoption_events(signal);
2334"#;
2335
2336fn migrations() -> &'static [Migration] {
2337    static MIGRATIONS: &[Migration] = &[
2338        Migration {
2339            version: 1,
2340            sql: V1_SCHEMA_SQL,
2341        },
2342        Migration {
2343            version: 2,
2344            sql: V2_MIGRATION_SQL,
2345        },
2346        Migration {
2347            version: 3,
2348            sql: V3_MIGRATION_SQL,
2349        },
2350        Migration {
2351            version: 4,
2352            sql: V4_MIGRATION_SQL,
2353        },
2354        Migration {
2355            version: 5,
2356            sql: V5_MIGRATION_SQL,
2357        },
2358        Migration {
2359            version: 6,
2360            sql: V6_MIGRATION_SQL,
2361        },
2362        Migration {
2363            version: 7,
2364            sql: V7_MIGRATION_SQL,
2365        },
2366        Migration {
2367            version: 8,
2368            sql: V8_MIGRATION_SQL,
2369        },
2370        Migration {
2371            version: 9,
2372            sql: V9_MIGRATION_SQL,
2373        },
2374    ];
2375    MIGRATIONS
2376}
2377
2378struct Migration {
2379    version: u32,
2380    sql: &'static str,
2381}
2382
2383fn register_sqlite_vec() -> Result<(), DbError> {
2384    SQLITE_VEC_AUTO_EXTENSION
2385        .get_or_init(|| unsafe {
2386            // sqlite-vec exposes a standard SQLite extension init symbol; auto-registration
2387            // makes vec0 available on every subsequently opened connection in this process.
2388            let init: rusqlite::auto_extension::RawAutoExtension =
2389                std::mem::transmute::<*const (), rusqlite::auto_extension::RawAutoExtension>(
2390                    sqlite_vec::sqlite3_vec_init as *const (),
2391                );
2392
2393            rusqlite::auto_extension::register_auto_extension(init)
2394                .map_err(|error| error.to_string())
2395        })
2396        .as_ref()
2397        .map(|_| ())
2398        .map_err(|message| DbError::RegisterVec(message.clone()))
2399}
2400
2401fn source_type_as_str(source_type: &SourceType) -> &'static str {
2402    match source_type {
2403        SourceType::Project => "project",
2404        SourceType::Conversation => "conversation",
2405        SourceType::Manual => "manual",
2406    }
2407}
2408
2409fn source_type_from_str(source_type: &str) -> Result<SourceType, DbError> {
2410    match source_type {
2411        "project" => Ok(SourceType::Project),
2412        "conversation" => Ok(SourceType::Conversation),
2413        "manual" => Ok(SourceType::Manual),
2414        other => Err(DbError::InvalidSourceType(other.to_string())),
2415    }
2416}
2417
2418fn memory_kind_as_str(memory_kind: &MemoryKind) -> &'static str {
2419    match memory_kind {
2420        MemoryKind::Evidence => "evidence",
2421        MemoryKind::Knowledge => "knowledge",
2422    }
2423}
2424
2425fn memory_kind_from_str(memory_kind: &str) -> Result<MemoryKind, DbError> {
2426    match memory_kind {
2427        "evidence" => Ok(MemoryKind::Evidence),
2428        "knowledge" => Ok(MemoryKind::Knowledge),
2429        other => Err(DbError::InvalidEnumValue {
2430            kind: "memory_kind",
2431            value: other.to_string(),
2432        }),
2433    }
2434}
2435
2436fn memory_domain_as_str(domain: &MemoryDomain) -> &'static str {
2437    match domain {
2438        MemoryDomain::Project => "project",
2439        MemoryDomain::Agent => "agent",
2440        MemoryDomain::Skill => "skill",
2441        MemoryDomain::Global => "global",
2442    }
2443}
2444
2445fn memory_domain_from_str(domain: &str) -> Result<MemoryDomain, DbError> {
2446    match domain {
2447        "project" => Ok(MemoryDomain::Project),
2448        "agent" => Ok(MemoryDomain::Agent),
2449        "skill" => Ok(MemoryDomain::Skill),
2450        "global" => Ok(MemoryDomain::Global),
2451        other => Err(DbError::InvalidEnumValue {
2452            kind: "domain",
2453            value: other.to_string(),
2454        }),
2455    }
2456}
2457
2458fn anchor_kind_as_str(anchor_kind: &AnchorKind) -> &'static str {
2459    match anchor_kind {
2460        AnchorKind::Global => "global",
2461        AnchorKind::Repo => "repo",
2462        AnchorKind::Worktree => "worktree",
2463    }
2464}
2465
2466fn anchor_kind_from_str(anchor_kind: &str) -> Result<AnchorKind, DbError> {
2467    match anchor_kind {
2468        "global" => Ok(AnchorKind::Global),
2469        "repo" => Ok(AnchorKind::Repo),
2470        "worktree" => Ok(AnchorKind::Worktree),
2471        other => Err(DbError::InvalidEnumValue {
2472            kind: "anchor_kind",
2473            value: other.to_string(),
2474        }),
2475    }
2476}
2477
2478fn provenance_as_str(provenance: &Provenance) -> &'static str {
2479    match provenance {
2480        Provenance::Runtime => "runtime",
2481        Provenance::Research => "research",
2482        Provenance::Human => "human",
2483    }
2484}
2485
2486fn provenance_from_str(provenance: &str) -> Result<Provenance, DbError> {
2487    match provenance {
2488        "runtime" => Ok(Provenance::Runtime),
2489        "research" => Ok(Provenance::Research),
2490        "human" => Ok(Provenance::Human),
2491        other => Err(DbError::InvalidEnumValue {
2492            kind: "provenance",
2493            value: other.to_string(),
2494        }),
2495    }
2496}
2497
2498fn knowledge_tier_as_str(tier: &KnowledgeTier) -> &'static str {
2499    match tier {
2500        KnowledgeTier::Qi => "qi",
2501        KnowledgeTier::Shu => "shu",
2502        KnowledgeTier::DaoRen => "dao_ren",
2503        KnowledgeTier::DaoTian => "dao_tian",
2504    }
2505}
2506
2507fn knowledge_tier_from_str(tier: &str) -> Result<KnowledgeTier, DbError> {
2508    match tier {
2509        "qi" => Ok(KnowledgeTier::Qi),
2510        "shu" => Ok(KnowledgeTier::Shu),
2511        "dao_ren" => Ok(KnowledgeTier::DaoRen),
2512        "dao_tian" => Ok(KnowledgeTier::DaoTian),
2513        other => Err(DbError::InvalidEnumValue {
2514            kind: "tier",
2515            value: other.to_string(),
2516        }),
2517    }
2518}
2519
2520fn knowledge_status_as_str(status: &KnowledgeStatus) -> &'static str {
2521    match status {
2522        KnowledgeStatus::Candidate => "candidate",
2523        KnowledgeStatus::Promoted => "promoted",
2524        KnowledgeStatus::Canonical => "canonical",
2525        KnowledgeStatus::Demoted => "demoted",
2526        KnowledgeStatus::Retired => "retired",
2527    }
2528}
2529
2530fn knowledge_status_from_str(status: &str) -> Result<KnowledgeStatus, DbError> {
2531    match status {
2532        "candidate" => Ok(KnowledgeStatus::Candidate),
2533        "promoted" => Ok(KnowledgeStatus::Promoted),
2534        "canonical" => Ok(KnowledgeStatus::Canonical),
2535        "demoted" => Ok(KnowledgeStatus::Demoted),
2536        "retired" => Ok(KnowledgeStatus::Retired),
2537        other => Err(DbError::InvalidEnumValue {
2538            kind: "status",
2539            value: other.to_string(),
2540        }),
2541    }
2542}
2543
2544fn knowledge_evidence_role_as_str(role: &KnowledgeEvidenceRole) -> &'static str {
2545    match role {
2546        KnowledgeEvidenceRole::Supporting => "supporting",
2547        KnowledgeEvidenceRole::Verification => "verification",
2548        KnowledgeEvidenceRole::Counterexample => "counterexample",
2549        KnowledgeEvidenceRole::Teaching => "teaching",
2550    }
2551}
2552
2553fn knowledge_evidence_role_from_str(role: &str) -> Result<KnowledgeEvidenceRole, DbError> {
2554    match role {
2555        "supporting" => Ok(KnowledgeEvidenceRole::Supporting),
2556        "verification" => Ok(KnowledgeEvidenceRole::Verification),
2557        "counterexample" => Ok(KnowledgeEvidenceRole::Counterexample),
2558        "teaching" => Ok(KnowledgeEvidenceRole::Teaching),
2559        other => Err(DbError::InvalidEnumValue {
2560            kind: "knowledge_evidence_role",
2561            value: other.to_string(),
2562        }),
2563    }
2564}
2565
2566fn knowledge_event_type_as_str(event_type: &KnowledgeEventType) -> &'static str {
2567    match event_type {
2568        KnowledgeEventType::Created => "created",
2569        KnowledgeEventType::Promoted => "promoted",
2570        KnowledgeEventType::Demoted => "demoted",
2571        KnowledgeEventType::Retired => "retired",
2572        KnowledgeEventType::Linked => "linked",
2573        KnowledgeEventType::Unlinked => "unlinked",
2574        KnowledgeEventType::Updated => "updated",
2575        KnowledgeEventType::PublishedAnchor => "published_anchor",
2576    }
2577}
2578
2579fn knowledge_event_type_from_str(event_type: &str) -> Result<KnowledgeEventType, DbError> {
2580    match event_type {
2581        "created" => Ok(KnowledgeEventType::Created),
2582        "promoted" => Ok(KnowledgeEventType::Promoted),
2583        "demoted" => Ok(KnowledgeEventType::Demoted),
2584        "retired" => Ok(KnowledgeEventType::Retired),
2585        "linked" => Ok(KnowledgeEventType::Linked),
2586        "unlinked" => Ok(KnowledgeEventType::Unlinked),
2587        "updated" => Ok(KnowledgeEventType::Updated),
2588        "published_anchor" => Ok(KnowledgeEventType::PublishedAnchor),
2589        other => Err(DbError::InvalidEnumValue {
2590            kind: "knowledge_event_type",
2591            value: other.to_string(),
2592        }),
2593    }
2594}
2595
2596fn runtime_adoption_track_as_str(track: &RuntimeAdoptionTrack) -> &'static str {
2597    match track {
2598        RuntimeAdoptionTrack::RuntimeAdoption => "runtime_adoption",
2599        RuntimeAdoptionTrack::CardContext => "card_context",
2600        RuntimeAdoptionTrack::CardEmbedding => "card_embedding",
2601        RuntimeAdoptionTrack::Evaluator => "evaluator",
2602        RuntimeAdoptionTrack::ResearchAdapter => "research_adapter",
2603    }
2604}
2605
2606fn runtime_adoption_track_from_str(track: &str) -> Result<RuntimeAdoptionTrack, DbError> {
2607    match track {
2608        "runtime_adoption" => Ok(RuntimeAdoptionTrack::RuntimeAdoption),
2609        "card_context" => Ok(RuntimeAdoptionTrack::CardContext),
2610        "card_embedding" => Ok(RuntimeAdoptionTrack::CardEmbedding),
2611        "evaluator" => Ok(RuntimeAdoptionTrack::Evaluator),
2612        "research_adapter" => Ok(RuntimeAdoptionTrack::ResearchAdapter),
2613        other => Err(DbError::InvalidEnumValue {
2614            kind: "runtime_adoption_track",
2615            value: other.to_string(),
2616        }),
2617    }
2618}
2619
2620fn runtime_adoption_signal_as_str(signal: &RuntimeAdoptionSignal) -> &'static str {
2621    match signal {
2622        RuntimeAdoptionSignal::Used => "used",
2623        RuntimeAdoptionSignal::Accepted => "accepted",
2624        RuntimeAdoptionSignal::Rejected => "rejected",
2625        RuntimeAdoptionSignal::Miss => "miss",
2626        RuntimeAdoptionSignal::Rollback => "rollback",
2627        RuntimeAdoptionSignal::Contradiction => "contradiction",
2628        RuntimeAdoptionSignal::Neutral => "neutral",
2629    }
2630}
2631
2632fn runtime_adoption_signal_from_str(signal: &str) -> Result<RuntimeAdoptionSignal, DbError> {
2633    match signal {
2634        "used" => Ok(RuntimeAdoptionSignal::Used),
2635        "accepted" => Ok(RuntimeAdoptionSignal::Accepted),
2636        "rejected" => Ok(RuntimeAdoptionSignal::Rejected),
2637        "miss" => Ok(RuntimeAdoptionSignal::Miss),
2638        "rollback" => Ok(RuntimeAdoptionSignal::Rollback),
2639        "contradiction" => Ok(RuntimeAdoptionSignal::Contradiction),
2640        "neutral" => Ok(RuntimeAdoptionSignal::Neutral),
2641        other => Err(DbError::InvalidEnumValue {
2642            kind: "runtime_adoption_signal",
2643            value: other.to_string(),
2644        }),
2645    }
2646}
2647
2648fn encode_json<T: serde::Serialize + ?Sized>(value: &T) -> Result<String, DbError> {
2649    Ok(serde_json::to_string(value)?)
2650}
2651
2652fn encode_optional_json<T: serde::Serialize>(value: Option<&T>) -> Result<Option<String>, DbError> {
2653    value.map(encode_json).transpose()
2654}
2655
2656fn parse_string_list(raw: Option<&str>) -> Result<Vec<String>, DbError> {
2657    let Some(raw) = raw else {
2658        return Ok(Vec::new());
2659    };
2660    Ok(serde_json::from_str::<Vec<String>>(raw)?)
2661}
2662
2663fn parse_optional_json<T>(raw: Option<&str>) -> Result<Option<T>, DbError>
2664where
2665    T: serde::de::DeserializeOwned,
2666{
2667    raw.map(serde_json::from_str)
2668        .transpose()
2669        .map_err(DbError::from)
2670}
2671
2672fn row_decode_error(error: DbError) -> rusqlite::Error {
2673    rusqlite::Error::FromSqlConversionFailure(0, rusqlite::types::Type::Text, Box::new(error))
2674}
2675
2676fn knowledge_card_from_row(row: &Row<'_>) -> Result<KnowledgeCard, DbError> {
2677    let tier = knowledge_tier_from_str(&row.get::<_, String>(3)?)?;
2678    let status = knowledge_status_from_str(&row.get::<_, String>(4)?)?;
2679    let domain = memory_domain_from_str(&row.get::<_, String>(5)?)?;
2680    let anchor_kind = anchor_kind_from_str(&row.get::<_, String>(7)?)?;
2681    let trigger_hints = parse_optional_json(row.get::<_, Option<String>>(11)?.as_deref())?;
2682
2683    anchor::validate_anchor_domain(&domain, &anchor_kind)
2684        .map_err(|message| DbError::InvalidDrawerMetadata(message.to_string()))?;
2685
2686    Ok(KnowledgeCard {
2687        id: row.get(0)?,
2688        statement: row.get(1)?,
2689        content: row.get(2)?,
2690        tier,
2691        status,
2692        domain,
2693        field: row.get(6)?,
2694        anchor_kind,
2695        anchor_id: row.get(8)?,
2696        parent_anchor_id: row.get(9)?,
2697        scope_constraints: row.get(10)?,
2698        trigger_hints,
2699        created_at: row.get(12)?,
2700        updated_at: row.get(13)?,
2701    })
2702}
2703
2704fn knowledge_evidence_link_from_row(row: &Row<'_>) -> Result<KnowledgeEvidenceLink, DbError> {
2705    Ok(KnowledgeEvidenceLink {
2706        id: row.get(0)?,
2707        card_id: row.get(1)?,
2708        evidence_drawer_id: row.get(2)?,
2709        role: knowledge_evidence_role_from_str(&row.get::<_, String>(3)?)?,
2710        note: row.get(4)?,
2711        created_at: row.get(5)?,
2712    })
2713}
2714
2715fn knowledge_event_from_row(row: &Row<'_>) -> Result<KnowledgeCardEvent, DbError> {
2716    let from_status = row
2717        .get::<_, Option<String>>(3)?
2718        .as_deref()
2719        .map(knowledge_status_from_str)
2720        .transpose()?;
2721    let to_status = row
2722        .get::<_, Option<String>>(4)?
2723        .as_deref()
2724        .map(knowledge_status_from_str)
2725        .transpose()?;
2726    let metadata = parse_optional_json(row.get::<_, Option<String>>(7)?.as_deref())?;
2727
2728    Ok(KnowledgeCardEvent {
2729        id: row.get(0)?,
2730        card_id: row.get(1)?,
2731        event_type: knowledge_event_type_from_str(&row.get::<_, String>(2)?)?,
2732        from_status,
2733        to_status,
2734        reason: row.get(5)?,
2735        actor: row.get(6)?,
2736        metadata,
2737        created_at: row.get(8)?,
2738    })
2739}
2740
2741fn runtime_adoption_event_from_row(row: &Row<'_>) -> Result<RuntimeAdoptionEvent, DbError> {
2742    let metadata = parse_optional_json(row.get::<_, Option<String>>(10)?.as_deref())?;
2743    Ok(RuntimeAdoptionEvent {
2744        id: row.get(0)?,
2745        track: runtime_adoption_track_from_str(&row.get::<_, String>(1)?)?,
2746        signal: runtime_adoption_signal_from_str(&row.get::<_, String>(2)?)?,
2747        feature: row.get(3)?,
2748        query: row.get(4)?,
2749        context_hash: row.get(5)?,
2750        card_id: row.get(6)?,
2751        evaluator_id: row.get(7)?,
2752        research_report_id: row.get(8)?,
2753        note: row.get(9)?,
2754        metadata,
2755        created_at: row.get(11)?,
2756    })
2757}
2758
2759fn drawer_from_row(row: &Row<'_>) -> Result<Drawer, DbError> {
2760    let source_type = source_type_from_str(&row.get::<_, String>(5)?)?;
2761    let memory_kind = memory_kind_from_str(&row.get::<_, String>(10)?)?;
2762    let domain = memory_domain_from_str(&row.get::<_, String>(11)?)?;
2763    let field = row.get::<_, String>(12)?;
2764    let anchor_kind = anchor_kind_from_str(&row.get::<_, String>(13)?)?;
2765    let anchor_id = row.get::<_, String>(14)?;
2766    let parent_anchor_id = row.get::<_, Option<String>>(15)?;
2767    let provenance = row
2768        .get::<_, Option<String>>(16)?
2769        .as_deref()
2770        .map(provenance_from_str)
2771        .transpose()?;
2772    let statement = row.get::<_, Option<String>>(17)?;
2773    let tier = row
2774        .get::<_, Option<String>>(18)?
2775        .as_deref()
2776        .map(knowledge_tier_from_str)
2777        .transpose()?;
2778    let status = row
2779        .get::<_, Option<String>>(19)?
2780        .as_deref()
2781        .map(knowledge_status_from_str)
2782        .transpose()?;
2783    let supporting_refs = parse_string_list(row.get::<_, Option<String>>(20)?.as_deref())?;
2784    let counterexample_refs = parse_string_list(row.get::<_, Option<String>>(21)?.as_deref())?;
2785    let teaching_refs = parse_string_list(row.get::<_, Option<String>>(22)?.as_deref())?;
2786    let verification_refs = parse_string_list(row.get::<_, Option<String>>(23)?.as_deref())?;
2787    let scope_constraints = row.get::<_, Option<String>>(24)?;
2788    let trigger_hints = parse_optional_json(row.get::<_, Option<String>>(25)?.as_deref())?;
2789
2790    anchor::validate_anchor_domain(&domain, &anchor_kind)
2791        .map_err(|message| DbError::InvalidDrawerMetadata(message.to_string()))?;
2792
2793    Ok(Drawer {
2794        id: row.get(0)?,
2795        content: row.get(1)?,
2796        wing: row.get(2)?,
2797        room: row.get(3)?,
2798        source_file: row.get(4)?,
2799        source_type,
2800        added_at: row.get(6)?,
2801        chunk_index: row.get(7)?,
2802        normalize_version: row.get(8)?,
2803        importance: row.get(9)?,
2804        memory_kind,
2805        domain,
2806        field,
2807        anchor_kind,
2808        anchor_id,
2809        parent_anchor_id,
2810        provenance,
2811        statement,
2812        tier,
2813        status,
2814        supporting_refs,
2815        counterexample_refs,
2816        teaching_refs,
2817        verification_refs,
2818        scope_constraints,
2819        trigger_hints,
2820    })
2821}
2822
2823fn parse_keywords(raw: Option<&str>) -> Result<Vec<String>, DbError> {
2824    let Some(raw) = raw else {
2825        return Ok(Vec::new());
2826    };
2827
2828    let value: Value = serde_json::from_str(raw)?;
2829    let keywords = value
2830        .as_array()
2831        .into_iter()
2832        .flatten()
2833        .filter_map(|item| item.as_str())
2834        .map(ToOwned::to_owned)
2835        .collect();
2836
2837    Ok(keywords)
2838}
2839
2840#[cfg(test)]
2841mod tests {
2842    use super::*;
2843
2844    #[test]
2845    fn test_atomic_migration_rolls_back_partial_schema_changes() {
2846        let conn = Connection::open_in_memory().expect("open in-memory");
2847        conn.execute_batch(
2848            r#"
2849            CREATE TABLE drawers (
2850                id TEXT PRIMARY KEY,
2851                content TEXT NOT NULL
2852            );
2853            PRAGMA user_version = 4;
2854            "#,
2855        )
2856        .expect("create base schema");
2857
2858        let migration = Migration {
2859            version: 5,
2860            sql: r#"
2861            ALTER TABLE drawers ADD COLUMN memory_kind TEXT;
2862            ALTER TABLE missing_table ADD COLUMN nope TEXT;
2863            "#,
2864        };
2865
2866        let error = apply_migration_atomic(&conn, &migration).expect_err("migration should fail");
2867        assert!(
2868            matches!(error, DbError::Sqlite(_)),
2869            "unexpected error: {error:?}"
2870        );
2871        assert_eq!(read_user_version(&conn).expect("user_version"), 4);
2872
2873        let mut stmt = conn
2874            .prepare("PRAGMA table_info(drawers)")
2875            .expect("table_info");
2876        let columns = stmt
2877            .query_map([], |row| row.get::<_, String>(1))
2878            .expect("query columns")
2879            .collect::<std::result::Result<Vec<_>, _>>()
2880            .expect("collect columns");
2881
2882        assert!(
2883            !columns.iter().any(|column| column == "memory_kind"),
2884            "failed migration must not leave partial columns behind"
2885        );
2886    }
2887}