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 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 fn ensure_vectors_table(&self, dim: usize) -> Result<(), DbError> {
520 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 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 pub fn embedding_dim(&self) -> Result<Option<usize>, DbError> {
1984 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 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 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 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}