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