1use crate::config::{EmbeddingConfig, MemoryLimits, PoolConfig};
4use crate::error::MemoryError;
5use crate::quantize::unpack_quantized;
6#[cfg(any(
7 feature = "turbo-quant-codec",
8 feature = "fib-quant-codec",
9 feature = "per-dim-codec"
10))]
11use crate::types::{DerivedVectorArtifactGenerationV1, VectorArtifactBuildReceiptV1};
12use crate::types::{EpisodeOutcome, Role, VectorSearchReceiptV1, VerificationStatus};
13use chrono::{DateTime, Utc};
14use rusqlite::{params, Connection, OpenFlags, OptionalExtension};
15use serde::{Deserialize, Serialize};
16use stack_ids::ContentDigest;
17#[cfg(any(
18 feature = "turbo-quant-codec",
19 feature = "fib-quant-codec",
20 feature = "per-dim-codec"
21))]
22use stack_ids::DigestBuilder;
23use std::path::Path;
24
25const MIGRATION_V1: &str = r#"
27-- CONVERSATIONS
28CREATE TABLE sessions (
29 id TEXT PRIMARY KEY,
30 channel TEXT NOT NULL DEFAULT 'repl',
31 created_at TEXT NOT NULL DEFAULT (datetime('now')),
32 updated_at TEXT NOT NULL DEFAULT (datetime('now')),
33 metadata TEXT
34);
35
36CREATE INDEX idx_sessions_updated ON sessions(updated_at DESC);
37
38CREATE TABLE messages (
39 id INTEGER PRIMARY KEY AUTOINCREMENT,
40 session_id TEXT NOT NULL REFERENCES sessions(id) ON DELETE CASCADE,
41 role TEXT NOT NULL CHECK (role IN ('system', 'user', 'assistant', 'tool')),
42 content TEXT NOT NULL,
43 token_count INTEGER,
44 created_at TEXT NOT NULL DEFAULT (datetime('now')),
45 metadata TEXT
46);
47
48CREATE INDEX idx_messages_session ON messages(session_id, created_at ASC);
49CREATE INDEX idx_messages_created ON messages(created_at DESC);
50
51-- KNOWLEDGE (Facts)
52CREATE TABLE facts (
53 id TEXT PRIMARY KEY,
54 namespace TEXT NOT NULL DEFAULT 'general',
55 content TEXT NOT NULL,
56 source TEXT,
57 embedding BLOB,
58 created_at TEXT NOT NULL DEFAULT (datetime('now')),
59 updated_at TEXT NOT NULL DEFAULT (datetime('now')),
60 metadata TEXT
61);
62
63CREATE INDEX idx_facts_namespace ON facts(namespace);
64CREATE INDEX idx_facts_updated ON facts(updated_at DESC);
65
66CREATE TABLE facts_rowid_map (
67 rowid INTEGER PRIMARY KEY AUTOINCREMENT,
68 fact_id TEXT NOT NULL UNIQUE REFERENCES facts(id) ON DELETE CASCADE
69);
70
71CREATE VIRTUAL TABLE facts_fts USING fts5(
72 content,
73 content='',
74 content_rowid='rowid',
75 tokenize='porter unicode61'
76);
77
78-- DOCUMENTS (Chunked content)
79CREATE TABLE documents (
80 id TEXT PRIMARY KEY,
81 title TEXT NOT NULL,
82 source_path TEXT,
83 namespace TEXT NOT NULL DEFAULT 'general',
84 created_at TEXT NOT NULL DEFAULT (datetime('now')),
85 metadata TEXT
86);
87
88CREATE TABLE chunks (
89 id TEXT PRIMARY KEY,
90 document_id TEXT NOT NULL REFERENCES documents(id) ON DELETE CASCADE,
91 chunk_index INTEGER NOT NULL,
92 content TEXT NOT NULL,
93 token_count INTEGER,
94 embedding BLOB,
95 created_at TEXT NOT NULL DEFAULT (datetime('now'))
96);
97
98CREATE INDEX idx_chunks_document ON chunks(document_id, chunk_index ASC);
99
100CREATE TABLE chunks_rowid_map (
101 rowid INTEGER PRIMARY KEY AUTOINCREMENT,
102 chunk_id TEXT NOT NULL UNIQUE REFERENCES chunks(id) ON DELETE CASCADE
103);
104
105CREATE VIRTUAL TABLE chunks_fts USING fts5(
106 content,
107 content='',
108 content_rowid='rowid',
109 tokenize='porter unicode61'
110);
111
112-- EMBEDDING METADATA
113CREATE TABLE embedding_metadata (
114 id INTEGER PRIMARY KEY CHECK (id = 1),
115 model_name TEXT NOT NULL,
116 dimensions INTEGER NOT NULL,
117 updated_at TEXT NOT NULL DEFAULT (datetime('now'))
118);
119"#;
120
121const MIGRATION_V2: &str = r#"
123ALTER TABLE messages ADD COLUMN embedding BLOB;
124
125CREATE TABLE messages_rowid_map (
126 rowid INTEGER PRIMARY KEY AUTOINCREMENT,
127 message_id INTEGER NOT NULL UNIQUE REFERENCES messages(id) ON DELETE CASCADE
128);
129
130CREATE VIRTUAL TABLE messages_fts USING fts5(
131 content,
132 content='',
133 content_rowid='rowid',
134 tokenize='porter unicode61'
135);
136"#;
137
138const MIGRATION_V3: &str = r#"
140ALTER TABLE embedding_metadata ADD COLUMN embeddings_dirty INTEGER NOT NULL DEFAULT 0;
141"#;
142
143const MIGRATION_V4: &str = r#"
145CREATE TABLE IF NOT EXISTS hnsw_metadata (
146 key TEXT PRIMARY KEY,
147 value TEXT NOT NULL
148);
149"#;
150
151const MIGRATION_V5: &str = r#"
153ALTER TABLE facts ADD COLUMN embedding_q8 BLOB;
154ALTER TABLE chunks ADD COLUMN embedding_q8 BLOB;
155ALTER TABLE messages ADD COLUMN embedding_q8 BLOB;
156
157CREATE TABLE IF NOT EXISTS hnsw_keymap (
158 node_id INTEGER PRIMARY KEY,
159 item_key TEXT NOT NULL UNIQUE,
160 deleted INTEGER NOT NULL DEFAULT 0
161);
162
163CREATE INDEX idx_hnsw_keymap_key ON hnsw_keymap(item_key);
164"#;
165
166const MIGRATION_V6: &str = r#"
168CREATE TABLE IF NOT EXISTS episodes (
169 document_id TEXT PRIMARY KEY REFERENCES documents(id) ON DELETE CASCADE,
170 cause_ids TEXT NOT NULL,
171 effect_type TEXT NOT NULL,
172 outcome TEXT NOT NULL DEFAULT 'pending',
173 confidence REAL NOT NULL DEFAULT 0.0,
174 verification_status TEXT NOT NULL DEFAULT '{"status":"unverified"}',
175 experiment_id TEXT,
176 created_at TEXT NOT NULL DEFAULT (datetime('now'))
177);
178
179CREATE INDEX IF NOT EXISTS idx_episodes_effect_type ON episodes(effect_type);
180CREATE INDEX IF NOT EXISTS idx_episodes_outcome ON episodes(outcome);
181CREATE INDEX IF NOT EXISTS idx_episodes_experiment_id ON episodes(experiment_id);
182"#;
183
184const MIGRATION_V7: &str = r#"
186ALTER TABLE episodes ADD COLUMN updated_at TEXT NOT NULL DEFAULT (datetime('now'));
187ALTER TABLE episodes ADD COLUMN search_text TEXT NOT NULL DEFAULT '';
188ALTER TABLE episodes ADD COLUMN embedding BLOB;
189ALTER TABLE episodes ADD COLUMN embedding_q8 BLOB;
190
191CREATE TABLE IF NOT EXISTS episodes_rowid_map (
192 rowid INTEGER PRIMARY KEY AUTOINCREMENT,
193 document_id TEXT NOT NULL UNIQUE REFERENCES episodes(document_id) ON DELETE CASCADE
194);
195
196CREATE VIRTUAL TABLE episodes_fts USING fts5(
197 content,
198 content='',
199 content_rowid='rowid',
200 tokenize='porter unicode61'
201);
202
203CREATE TABLE IF NOT EXISTS pending_index_ops (
204 item_key TEXT PRIMARY KEY,
205 entity_type TEXT NOT NULL,
206 op_kind TEXT NOT NULL CHECK (op_kind IN ('upsert', 'delete')),
207 attempt_count INTEGER NOT NULL DEFAULT 0,
208 last_error TEXT,
209 updated_at TEXT NOT NULL DEFAULT (datetime('now'))
210);
211
212INSERT OR IGNORE INTO hnsw_metadata (key, value) VALUES ('sidecar_dirty', '0');
213
214UPDATE episodes
215SET search_text = TRIM(
216 COALESCE(effect_type, '') || ' ' ||
217 COALESCE(outcome, '') || ' ' ||
218 COALESCE(experiment_id, '') || ' ' ||
219 COALESCE(cause_ids, '')
220)
221WHERE search_text = '';
222
223INSERT OR IGNORE INTO episodes_rowid_map (document_id)
224SELECT document_id FROM episodes;
225
226INSERT INTO episodes_fts (rowid, content)
227SELECT rm.rowid, e.search_text
228FROM episodes_rowid_map rm
229JOIN episodes e ON e.document_id = rm.document_id;
230"#;
231
232const MIGRATION_V8: &str = r#"
234ALTER TABLE episodes ADD COLUMN trace_id TEXT;
235"#;
236
237const MIGRATION_V9: &str = "";
245
246const MIGRATION_V18: &str = r#"
248CREATE TABLE IF NOT EXISTS search_receipts (
249 receipt_id TEXT PRIMARY KEY,
250 schema_version TEXT NOT NULL,
251 evaluation_time TEXT NOT NULL,
252 search_profile TEXT NOT NULL,
253 candidate_backend TEXT NOT NULL,
254 approximate INTEGER NOT NULL CHECK (approximate IN (0, 1)),
255 exact_rerank INTEGER NOT NULL CHECK (exact_rerank IN (0, 1)),
256 fallback TEXT,
257 requested_candidates INTEGER NOT NULL CHECK (requested_candidates >= 0),
258 returned_candidates INTEGER NOT NULL CHECK (returned_candidates >= 0),
259 post_filter_candidates INTEGER NOT NULL CHECK (post_filter_candidates >= 0),
260 result_ids_json TEXT NOT NULL,
261 receipt_json TEXT NOT NULL,
262 receipt_digest TEXT NOT NULL,
263 created_at TEXT NOT NULL DEFAULT (datetime('now'))
264);
265
266CREATE INDEX IF NOT EXISTS idx_search_receipts_created
267ON search_receipts(created_at DESC);
268
269CREATE INDEX IF NOT EXISTS idx_search_receipts_backend
270ON search_receipts(candidate_backend);
271"#;
272
273const MIGRATION_V19: &str = r#"
275CREATE TABLE IF NOT EXISTS derived_vector_artifacts (
276 item_key TEXT NOT NULL,
277 codec_family TEXT NOT NULL,
278 codec_profile_digest TEXT NOT NULL,
279 source_embedding_digest TEXT NOT NULL,
280 encoded_digest TEXT NOT NULL,
281 artifact_digest TEXT NOT NULL,
282 encoding TEXT NOT NULL,
283 dim INTEGER NOT NULL,
284 encoded BLOB NOT NULL,
285 created_at TEXT NOT NULL DEFAULT (datetime('now')),
286 status TEXT NOT NULL DEFAULT 'active',
287 PRIMARY KEY (item_key, codec_family, codec_profile_digest)
288);
289
290CREATE INDEX IF NOT EXISTS idx_derived_vector_artifacts_profile
291ON derived_vector_artifacts(codec_family, codec_profile_digest, status);
292
293CREATE INDEX IF NOT EXISTS idx_derived_vector_artifacts_source_digest
294ON derived_vector_artifacts(source_embedding_digest);
295"#;
296
297const MIGRATION_V20: &str = r#"
299-- Procedural migration; see run_migration_v20.
300"#;
301
302const MIGRATION_V21: &str = r#"
304CREATE TABLE IF NOT EXISTS derived_vector_artifact_generations (
305 generation_id TEXT PRIMARY KEY,
306 schema_version TEXT NOT NULL,
307 codec_family TEXT NOT NULL,
308 codec_profile_digest TEXT NOT NULL,
309 source_snapshot_digest TEXT NOT NULL,
310 source_row_count INTEGER NOT NULL,
311 artifact_count INTEGER NOT NULL,
312 source_tables_json TEXT NOT NULL,
313 dim INTEGER NOT NULL,
314 encoding TEXT NOT NULL,
315 created_at TEXT NOT NULL,
316 build_receipt_id TEXT,
317 artifact_manifest_digest TEXT NOT NULL,
318 status TEXT NOT NULL CHECK (status IN ('active', 'superseded', 'invalidated', 'failed')),
319 degradations_json TEXT NOT NULL DEFAULT '[]'
320);
321
322CREATE INDEX IF NOT EXISTS idx_derived_vector_generations_profile
323ON derived_vector_artifact_generations(codec_family, codec_profile_digest, status, created_at DESC);
324"#;
325
326const MIGRATION_V23: &str = r#"
329ALTER TABLE derived_vector_artifacts ADD COLUMN codec_governance_receipt_id TEXT;
330ALTER TABLE derived_vector_artifacts ADD COLUMN codec_profile TEXT;
331ALTER TABLE derived_vector_artifacts ADD COLUMN degradation_budget REAL;
332ALTER TABLE derived_vector_artifacts ADD COLUMN raw_source_artifact_id TEXT;
333"#;
334
335const MIGRATION_V24: &str = r#"
337ALTER TABLE derived_vector_artifact_generations ADD COLUMN config_json TEXT NOT NULL DEFAULT '{}';
338"#;
339
340const MIGRATION_V22: &str = r#"
343ALTER TABLE episodes ADD COLUMN valid_time TEXT;
344ALTER TABLE episodes ADD COLUMN recorded_time TEXT NOT NULL DEFAULT (datetime('now'));
345ALTER TABLE episodes ADD COLUMN superseded_by TEXT;
346ALTER TABLE episodes ADD COLUMN fact_digest TEXT;
347CREATE INDEX IF NOT EXISTS idx_episodes_recorded ON episodes(recorded_time ASC);
348CREATE INDEX IF NOT EXISTS idx_episodes_valid ON episodes(valid_time);
349CREATE INDEX IF NOT EXISTS idx_episodes_superseded ON episodes(superseded_by) WHERE superseded_by IS NOT NULL;
350UPDATE episodes SET recorded_time = updated_at WHERE recorded_time IS NULL OR recorded_time = '';
351"#;
352
353#[allow(deprecated)]
355const MIGRATIONS: &[(u32, &str)] = &[
356 (1, MIGRATION_V1),
357 (2, MIGRATION_V2),
358 (3, MIGRATION_V3),
359 (4, MIGRATION_V4),
360 (5, MIGRATION_V5),
361 (6, MIGRATION_V6),
362 (7, MIGRATION_V7),
363 (8, MIGRATION_V8),
364 (9, MIGRATION_V9),
365 (10, crate::projection_import::MIGRATION_V10),
366 (11, crate::projection_storage::MIGRATION_V11),
367 (12, crate::projection_storage::MIGRATION_V12),
368 (13, crate::projection_storage::MIGRATION_V13),
369 (14, crate::projection_storage::MIGRATION_V14),
370 (15, crate::projection_storage::MIGRATION_V15),
371 (16, crate::projection_storage::MIGRATION_V16),
372 (17, crate::projection_storage::MIGRATION_V17),
373 (18, MIGRATION_V18),
374 (19, MIGRATION_V19),
375 (20, MIGRATION_V20),
376 (21, MIGRATION_V21),
377 (22, MIGRATION_V22),
378 (23, MIGRATION_V23),
379 (24, MIGRATION_V24),
380];
381
382pub const MAX_SCHEMA_VERSION: u32 = 24;
384
385fn run_migration_v9(conn: &Connection) -> Result<(), MemoryError> {
387 let episodes_exist: bool = conn
389 .query_row(
390 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type='table' AND name='episodes'",
391 [],
392 |row| row.get(0),
393 )
394 .map_err(|e| MemoryError::MigrationFailed {
395 version: 9,
396 reason: format!("existence check failed: {e}"),
397 })?;
398
399 if !episodes_exist {
400 conn.execute_batch(
402 "CREATE TABLE IF NOT EXISTS episode_causes (
403 episode_id TEXT NOT NULL,
404 cause_node_id TEXT NOT NULL,
405 ordinal INTEGER NOT NULL DEFAULT 0,
406 PRIMARY KEY (episode_id, cause_node_id)
407 );
408 CREATE INDEX IF NOT EXISTS idx_episode_causes_cause ON episode_causes(cause_node_id);",
409 )?;
410 return Ok(());
411 }
412
413 conn.execute_batch("PRAGMA foreign_keys = OFF;")?;
415
416 conn.execute_batch(
417 "CREATE TABLE episodes_new (
418 episode_id TEXT PRIMARY KEY,
419 document_id TEXT NOT NULL REFERENCES documents(id) ON DELETE CASCADE,
420 cause_ids TEXT NOT NULL,
421 effect_type TEXT NOT NULL,
422 outcome TEXT NOT NULL DEFAULT 'pending',
423 confidence REAL NOT NULL DEFAULT 0.0,
424 verification_status TEXT NOT NULL DEFAULT '{\"status\":\"unverified\"}',
425 experiment_id TEXT,
426 created_at TEXT NOT NULL DEFAULT (datetime('now')),
427 updated_at TEXT NOT NULL DEFAULT (datetime('now')),
428 search_text TEXT NOT NULL DEFAULT '',
429 embedding BLOB,
430 embedding_q8 BLOB,
431 trace_id TEXT
432 )",
433 )?;
434
435 conn.execute_batch(
437 "INSERT INTO episodes_new
438 (episode_id, document_id, cause_ids, effect_type, outcome, confidence,
439 verification_status, experiment_id, created_at, updated_at,
440 search_text, embedding, embedding_q8, trace_id)
441 SELECT
442 document_id || '-ep0',
443 document_id, cause_ids, effect_type, outcome, confidence,
444 verification_status, experiment_id, created_at, updated_at,
445 search_text, embedding, embedding_q8, trace_id
446 FROM episodes",
447 )?;
448
449 conn.execute_batch("DROP TABLE episodes")?;
450 conn.execute_batch("ALTER TABLE episodes_new RENAME TO episodes")?;
451
452 conn.execute_batch(
453 "CREATE INDEX idx_episodes_document_id ON episodes(document_id);
454 CREATE INDEX idx_episodes_effect_type ON episodes(effect_type);
455 CREATE INDEX idx_episodes_outcome ON episodes(outcome);
456 CREATE INDEX idx_episodes_experiment_id ON episodes(experiment_id);",
457 )?;
458
459 conn.execute_batch(
461 "DROP TABLE IF EXISTS episodes_rowid_map;
462 CREATE TABLE episodes_rowid_map (
463 rowid INTEGER PRIMARY KEY AUTOINCREMENT,
464 episode_id TEXT NOT NULL UNIQUE,
465 document_id TEXT
466 );
467 INSERT INTO episodes_rowid_map (episode_id, document_id)
468 SELECT episode_id, document_id FROM episodes;",
469 )?;
470
471 conn.execute_batch(
473 "DROP TABLE IF EXISTS episodes_fts;
474 CREATE VIRTUAL TABLE episodes_fts USING fts5(
475 content,
476 content='',
477 content_rowid='rowid',
478 tokenize='porter unicode61'
479 );
480 INSERT INTO episodes_fts (rowid, content)
481 SELECT rm.rowid, e.search_text
482 FROM episodes_rowid_map rm
483 JOIN episodes e ON e.episode_id = rm.episode_id;",
484 )?;
485
486 conn.execute_batch(
488 "CREATE TABLE IF NOT EXISTS episode_causes (
489 episode_id TEXT NOT NULL,
490 cause_node_id TEXT NOT NULL,
491 ordinal INTEGER NOT NULL DEFAULT 0,
492 PRIMARY KEY (episode_id, cause_node_id)
493 );
494 CREATE INDEX IF NOT EXISTS idx_episode_causes_cause ON episode_causes(cause_node_id);",
495 )?;
496
497 conn.execute_batch(
499 "INSERT OR IGNORE INTO episode_causes (episode_id, cause_node_id, ordinal)
500 SELECT e.episode_id, je.value, CAST(je.key AS INTEGER)
501 FROM episodes e, json_each(e.cause_ids) je;",
502 )?;
503
504 conn.execute_batch("PRAGMA foreign_keys = ON;")?;
505
506 Ok(())
507}
508
509#[derive(Debug, Clone, Copy, PartialEq, Eq)]
511pub enum VerifyMode {
512 Quick,
514 Full,
516}
517
518#[derive(Debug, Clone)]
520pub struct IntegrityReport {
521 pub ok: bool,
522 pub schema_version: u32,
523 pub fact_count: usize,
524 pub chunk_count: usize,
525 pub message_count: usize,
526 pub facts_missing_embeddings: usize,
527 pub chunks_missing_embeddings: usize,
528 pub issues: Vec<String>,
529}
530
531#[derive(Debug, Clone, Copy, PartialEq, Eq)]
533pub enum ReconcileAction {
534 ReportOnly,
535 RebuildFts,
536 ReEmbed,
537}
538
539#[derive(Debug, Clone, Copy, PartialEq, Eq)]
541pub(crate) enum IndexOpKind {
542 Upsert,
543 Delete,
544}
545
546impl IndexOpKind {
547 pub(crate) fn as_str(self) -> &'static str {
548 match self {
549 Self::Upsert => "upsert",
550 Self::Delete => "delete",
551 }
552 }
553
554 fn parse(raw: &str, item_key: &str) -> Result<Self, MemoryError> {
555 match raw {
556 "upsert" => Ok(Self::Upsert),
557 "delete" => Ok(Self::Delete),
558 other => Err(MemoryError::CorruptData {
559 table: "pending_index_ops",
560 row_id: item_key.to_string(),
561 detail: format!("invalid op_kind '{other}'"),
562 }),
563 }
564 }
565}
566
567#[derive(Debug, Clone)]
569pub(crate) struct PendingIndexOp {
570 pub item_key: String,
571 pub entity_type: String,
572 pub op_kind: IndexOpKind,
573 pub attempt_count: u32,
574 pub last_error: Option<String>,
575}
576
577pub fn with_transaction<F, T>(conn: &Connection, f: F) -> Result<T, MemoryError>
579where
580 F: FnOnce(&rusqlite::Transaction<'_>) -> Result<T, MemoryError>,
581{
582 let tx = conn.unchecked_transaction()?;
583 let result = f(&tx)?;
584 tx.commit()?;
585 Ok(result)
586}
587
588pub fn open_database(
590 path: &Path,
591 pool: &PoolConfig,
592 limits: &MemoryLimits,
593) -> Result<Connection, MemoryError> {
594 open_database_internal(path, pool, limits.max_db_size_bytes, true)
595}
596
597pub fn open_database_connection(
599 path: &Path,
600 pool: &PoolConfig,
601 limits: &MemoryLimits,
602) -> Result<Connection, MemoryError> {
603 open_database_internal(path, pool, limits.max_db_size_bytes, false)
604}
605
606pub(crate) fn open_database_internal(
607 path: &Path,
608 pool: &PoolConfig,
609 max_db_size_bytes: u64,
610 run_schema_migrations: bool,
611) -> Result<Connection, MemoryError> {
612 create_parent_dirs(path)?;
613 let conn = Connection::open(path)?;
614 configure_connection(&conn, path, pool, max_db_size_bytes, false)?;
615 if run_schema_migrations {
616 run_migrations(&conn)?;
617 }
618 Ok(conn)
619}
620
621pub(crate) fn open_pool_member_connection(
622 path: &Path,
623 pool: &PoolConfig,
624 limits: &MemoryLimits,
625 query_only: bool,
626) -> Result<Connection, MemoryError> {
627 create_parent_dirs(path)?;
628 let flags = OpenFlags::SQLITE_OPEN_READ_WRITE | OpenFlags::SQLITE_OPEN_CREATE;
629 let conn = Connection::open_with_flags(path, flags)?;
630 configure_connection(&conn, path, pool, limits.max_db_size_bytes, query_only)?;
631 Ok(conn)
632}
633
634fn create_parent_dirs(path: &Path) -> Result<(), MemoryError> {
635 if let Some(parent) = path.parent() {
636 if !parent.as_os_str().is_empty() {
637 std::fs::create_dir_all(parent).map_err(|e| {
638 MemoryError::StorageError(format!(
639 "failed to create database directory {}: {}",
640 parent.display(),
641 e
642 ))
643 })?;
644 }
645 }
646 Ok(())
647}
648
649fn configure_connection(
650 conn: &Connection,
651 path: &Path,
652 pool: &PoolConfig,
653 max_db_size_bytes: u64,
654 query_only: bool,
655) -> Result<(), MemoryError> {
656 let journal_mode = if pool.enable_wal { "WAL" } else { "DELETE" };
657 conn.execute_batch(&format!(
658 "PRAGMA journal_mode = {};
659 PRAGMA foreign_keys = ON;
660 PRAGMA busy_timeout = {};
661 PRAGMA synchronous = NORMAL;
662 PRAGMA temp_store = MEMORY;
663 PRAGMA wal_autocheckpoint = {};",
664 journal_mode, pool.busy_timeout_ms, pool.wal_autocheckpoint,
665 ))?;
666
667 if query_only {
668 conn.execute_batch("PRAGMA query_only = ON;")?;
669 }
670
671 let actual_journal_mode: String =
672 conn.query_row("PRAGMA journal_mode", [], |row| row.get(0))?;
673 let expected_journal_mode = if pool.enable_wal { "wal" } else { "delete" };
674 if actual_journal_mode.to_lowercase() != expected_journal_mode {
675 return Err(MemoryError::StorageError(format!(
676 "SQLite journal mode mismatch for {}: requested {}, got {}",
677 path.display(),
678 expected_journal_mode,
679 actual_journal_mode
680 )));
681 }
682
683 if max_db_size_bytes > 0 {
684 let page_size: u64 = conn.query_row("PRAGMA page_size", [], |row| row.get(0))?;
685 let max_page_count = max_db_size_bytes.div_ceil(page_size);
686
687 const MAX_SQLITE_PAGE_COUNT: u64 = 1_073_741_823; const MIN_SQLITE_PAGE_COUNT: u64 = 1;
690 if !(MIN_SQLITE_PAGE_COUNT..=MAX_SQLITE_PAGE_COUNT).contains(&max_page_count) {
691 return Err(MemoryError::StorageError(format!(
692 "Invalid max_page_count {}: must be between {} and {}",
693 max_page_count, MIN_SQLITE_PAGE_COUNT, MAX_SQLITE_PAGE_COUNT
694 )));
695 }
696
697 let actual_max_page_count: u64 = conn.query_row(
698 &format!("PRAGMA max_page_count = {}", max_page_count),
699 [],
700 |row| row.get(0),
701 )?;
702 let page_count: u64 = conn.query_row("PRAGMA page_count", [], |row| row.get(0))?;
703
704 if page_count > actual_max_page_count {
705 return Err(MemoryError::DatabaseSizeLimitExceeded {
706 current: page_count.saturating_mul(page_size),
707 limit: max_db_size_bytes,
708 });
709 }
710 }
711
712 let foreign_keys_enabled: bool = conn.query_row("PRAGMA foreign_keys", [], |row| row.get(0))?;
714 if !foreign_keys_enabled {
715 return Err(MemoryError::StorageError(
716 "PRAGMA foreign_keys failed to enable after configuration".to_string(),
717 ));
718 }
719
720 Ok(())
721}
722
723pub fn run_migrations(conn: &Connection) -> Result<(), MemoryError> {
725 let user_version: u32 = conn
726 .query_row("PRAGMA user_version", [], |row| row.get(0))
727 .map_err(|e| MemoryError::MigrationFailed {
728 version: 0,
729 reason: format!("failed to read PRAGMA user_version: {e}"),
730 })?;
731
732 if user_version > MAX_SCHEMA_VERSION {
733 return Err(MemoryError::SchemaAhead {
734 found: user_version,
735 supported: MAX_SCHEMA_VERSION,
736 });
737 }
738
739 conn.execute_batch(
740 "CREATE TABLE IF NOT EXISTS _schema_version (
741 version INTEGER PRIMARY KEY,
742 applied_at TEXT NOT NULL DEFAULT (datetime('now'))
743 );",
744 )?;
745
746 for &(version, sql) in MIGRATIONS {
747 let current_version: u32 = conn
748 .query_row(
749 "SELECT COALESCE(MAX(version), 0) FROM _schema_version",
750 [],
751 |row| row.get(0),
752 )
753 .unwrap_or(0);
754
755 if current_version >= version {
756 continue;
757 }
758
759 with_transaction(conn, |tx| {
760 match version {
761 9 => run_migration_v9(tx).map_err(|e| MemoryError::MigrationFailed {
762 version,
763 reason: e.to_string(),
764 })?,
765 16 => run_migration_v16(tx).map_err(|e| MemoryError::MigrationFailed {
766 version,
767 reason: e.to_string(),
768 })?,
769 17 => run_migration_v17(tx).map_err(|e| MemoryError::MigrationFailed {
770 version,
771 reason: e.to_string(),
772 })?,
773 20 => run_migration_v20(tx).map_err(|e| MemoryError::MigrationFailed {
774 version,
775 reason: e.to_string(),
776 })?,
777 21 => run_migration_v21(tx).map_err(|e| MemoryError::MigrationFailed {
778 version,
779 reason: e.to_string(),
780 })?,
781 _ => tx
782 .execute_batch(sql)
783 .map_err(|e| MemoryError::MigrationFailed {
784 version,
785 reason: e.to_string(),
786 })?,
787 }
788 tx.execute(
789 "INSERT INTO _schema_version (version) VALUES (?1)",
790 params![version],
791 )
792 .map_err(|e| MemoryError::MigrationFailed {
793 version,
794 reason: e.to_string(),
795 })?;
796 Ok(())
797 })?;
798
799 tracing::info!("Applied migration V{}", version);
800 }
801
802 let final_version: u32 = conn
803 .query_row(
804 "SELECT COALESCE(MAX(version), 0) FROM _schema_version",
805 [],
806 |row| row.get(0),
807 )
808 .unwrap_or(0);
809 conn.execute_batch(&format!("PRAGMA user_version = {};", final_version))?;
810
811 Ok(())
812}
813
814fn run_migration_v16(conn: &Connection) -> Result<(), rusqlite::Error> {
815 add_column_if_missing(conn, "projection_import_log", "kernel_payload_json", "TEXT")?;
816 add_column_if_missing(
817 conn,
818 "projection_import_failures",
819 "kernel_payload_json",
820 "TEXT",
821 )?;
822 Ok(())
823}
824
825fn run_migration_v17(conn: &Connection) -> Result<(), rusqlite::Error> {
826 add_column_if_missing(conn, "projection_import_log", "episode_bundle_id", "TEXT")?;
827 add_column_if_missing(conn, "projection_import_log", "episode_bundle_json", "TEXT")?;
828 add_column_if_missing(
829 conn,
830 "projection_import_log",
831 "execution_context_json",
832 "TEXT",
833 )?;
834 add_column_if_missing(
835 conn,
836 "projection_import_failures",
837 "episode_bundle_id",
838 "TEXT",
839 )?;
840 add_column_if_missing(
841 conn,
842 "projection_import_failures",
843 "episode_bundle_json",
844 "TEXT",
845 )?;
846 add_column_if_missing(
847 conn,
848 "projection_import_failures",
849 "execution_context_json",
850 "TEXT",
851 )?;
852 Ok(())
853}
854
855fn run_migration_v20(conn: &Connection) -> Result<(), rusqlite::Error> {
856 add_column_if_missing(conn, "derived_vector_artifacts", "encoded_digest", "TEXT")?;
857 conn.execute(
858 "UPDATE derived_vector_artifacts
859 SET encoded_digest = artifact_digest
860 WHERE encoded_digest IS NULL OR encoded_digest = ''",
861 [],
862 )?;
863 add_column_if_missing(
864 conn,
865 "derived_vector_artifacts",
866 "encoding",
867 "TEXT NOT NULL DEFAULT 'turbo_code_wire_v1'",
868 )?;
869 add_column_if_missing(
870 conn,
871 "derived_vector_artifacts",
872 "dim",
873 "INTEGER NOT NULL DEFAULT 0",
874 )?;
875 add_column_if_missing(
876 conn,
877 "derived_vector_artifacts",
878 "status",
879 "TEXT NOT NULL DEFAULT 'active'",
880 )?;
881 conn.execute_batch(
882 "CREATE INDEX IF NOT EXISTS idx_derived_vector_artifacts_profile
883 ON derived_vector_artifacts(codec_family, codec_profile_digest, status);
884 CREATE INDEX IF NOT EXISTS idx_derived_vector_artifacts_source_digest
885 ON derived_vector_artifacts(source_embedding_digest);",
886 )?;
887 Ok(())
888}
889
890fn run_migration_v21(conn: &Connection) -> Result<(), rusqlite::Error> {
891 conn.execute_batch(MIGRATION_V21)?;
892 add_column_if_missing(conn, "derived_vector_artifacts", "generation_id", "TEXT")?;
893 conn.execute_batch(
894 "CREATE INDEX IF NOT EXISTS idx_derived_vector_artifacts_generation
895 ON derived_vector_artifacts(generation_id, status);",
896 )?;
897 Ok(())
898}
899
900const SEARCH_RECEIPT_SCHEMA_VERSION: &str = "vector_search_receipt_v1";
901
902#[derive(Debug, Serialize, Deserialize)]
903struct StoredVectorSearchReceiptV1 {
904 #[serde(default = "default_search_receipt_schema_version")]
905 schema_version: String,
906 receipt_id: String,
907 evaluation_time: DateTime<Utc>,
908 #[serde(default)]
909 receipt_digest: Option<String>,
910 #[serde(default)]
911 trace_id: Option<String>,
912 #[serde(default)]
913 attempt_family_id: Option<String>,
914 #[serde(default)]
915 attempt_id: Option<String>,
916 #[serde(default)]
917 replay_of: Option<String>,
918 query_embedding_digest: Option<String>,
919 #[serde(default)]
920 query_text_digest: Option<String>,
921 #[serde(default)]
922 query_input_digest: Option<String>,
923 #[serde(default)]
924 filter_digest: Option<String>,
925 #[serde(default)]
926 redaction_state: Option<String>,
927 #[serde(default)]
928 budget_id: Option<String>,
929 #[serde(default)]
930 deadline_at: Option<DateTime<Utc>>,
931 search_profile: String,
932 candidate_backend: String,
933 codec_family: Option<String>,
934 codec_profile_digest: Option<String>,
935 #[serde(default)]
936 artifact_profile_digest: Option<String>,
937 #[serde(default)]
938 artifact_count: Option<u64>,
939 #[serde(default)]
940 artifact_corruption_count: Option<u64>,
941 #[serde(default)]
942 artifact_missing_count: Option<u64>,
943 #[serde(default)]
944 vector_artifact_manifest_digest: Option<String>,
945 #[serde(default)]
946 artifact_generation_id: Option<String>,
947 #[serde(default)]
948 approximate_scanned_count: Option<u64>,
949 #[serde(default)]
950 approximate_returned_count: Option<u64>,
951 #[serde(default)]
952 raw_rows_loaded_count: Option<u64>,
953 #[serde(default)]
954 filter_strategy: Option<String>,
955 #[serde(default)]
956 vector_artifact_count: Option<u64>,
957 #[serde(default)]
958 vector_artifact_missing_count: Option<u64>,
959 #[serde(default)]
960 vector_artifact_stale_count: Option<u64>,
961 #[serde(default)]
962 exact_rerank_count: Option<u64>,
963 #[serde(default)]
964 approximate_candidate_count: Option<u64>,
965 #[serde(default)]
966 fallback_reason: Option<String>,
967 approximate: bool,
968 requested_candidates: u64,
969 returned_candidates: u64,
970 post_filter_candidates: u64,
971 fallback: Option<String>,
972 exact_rerank: bool,
973 result_ids: Vec<String>,
974 degradations: Vec<String>,
975}
976
977fn default_search_receipt_schema_version() -> String {
978 SEARCH_RECEIPT_SCHEMA_VERSION.to_string()
979}
980
981fn b3_digest(bytes: &[u8]) -> String {
982 format!("blake3:{}", ContentDigest::compute(bytes).hex())
983}
984
985#[cfg(any(
987 feature = "turbo-quant-codec",
988 feature = "fib-quant-codec",
989 feature = "per-dim-codec"
990))]
991#[derive(Debug, Clone)]
992pub(crate) struct DerivedVectorArtifactRow {
993 pub item_key: String,
994 pub generation_id: Option<String>,
995 pub codec_family: String,
996 pub codec_profile_digest: String,
997 pub source_embedding_digest: String,
998 pub encoded_digest: String,
999 pub encoding: String,
1000 pub dim: usize,
1001 pub status: String,
1002 pub encoded: Vec<u8>,
1003 pub codec_governance_receipt_id: Option<String>,
1005 pub codec_profile: Option<String>,
1006 pub degradation_budget: Option<f64>,
1007 pub raw_source_artifact_id: Option<String>,
1008}
1009
1010#[cfg(any(
1012 feature = "turbo-quant-codec",
1013 feature = "fib-quant-codec",
1014 feature = "per-dim-codec"
1015))]
1016#[derive(Debug, Clone)]
1017#[allow(dead_code)]
1018pub(crate) struct DerivedVectorArtifactGenerationRow {
1019 pub generation_id: String,
1020 pub codec_family: String,
1021 pub codec_profile_digest: String,
1022 pub source_snapshot_digest: String,
1023 pub source_row_count: usize,
1024 pub artifact_count: usize,
1025 pub dim: usize,
1026 pub encoding: String,
1027 pub artifact_manifest_digest: String,
1028 pub status: String,
1029}
1030
1031#[cfg(any(
1033 feature = "turbo-quant-codec",
1034 feature = "fib-quant-codec",
1035 feature = "per-dim-codec"
1036))]
1037pub(crate) fn source_embedding_digest(
1038 blob: &[u8],
1039 expected_dim: usize,
1040) -> Result<String, MemoryError> {
1041 validate_vector_blob_len(blob, expected_dim)?;
1042 let mut builder = DigestBuilder::new();
1043 builder
1044 .update_str("semantic-memory.source_embedding.v1")
1045 .separator()
1046 .update(&(expected_dim as u64).to_le_bytes())
1047 .separator()
1048 .update(blob);
1049 Ok(format!("blake3:{}", builder.finalize().hex()))
1050}
1051
1052#[cfg(any(
1053 feature = "turbo-quant-codec",
1054 feature = "fib-quant-codec",
1055 feature = "per-dim-codec"
1056))]
1057fn source_snapshot_digest(rows: &[DerivedVectorArtifactRow], dim: usize) -> String {
1058 let mut entries = rows
1059 .iter()
1060 .map(|row| (row.item_key.as_str(), row.source_embedding_digest.as_str()))
1061 .collect::<Vec<_>>();
1062 entries.sort_unstable();
1063
1064 let mut builder = DigestBuilder::new();
1065 builder
1066 .update_str("semantic-memory.vector_source_snapshot.v1")
1067 .separator()
1068 .update(&(dim as u64).to_le_bytes())
1069 .separator();
1070 for (item_key, source_embedding_digest) in entries {
1071 builder
1072 .update_str(item_key)
1073 .separator()
1074 .update_str(source_embedding_digest)
1075 .separator();
1076 }
1077 format!("blake3:{}", builder.finalize().hex())
1078}
1079
1080#[cfg(any(
1081 feature = "turbo-quant-codec",
1082 feature = "fib-quant-codec",
1083 feature = "per-dim-codec"
1084))]
1085pub(crate) fn current_source_snapshot_digest(
1086 conn: &Connection,
1087 dim: usize,
1088) -> Result<(String, usize), MemoryError> {
1089 let mut stmt = conn.prepare(
1090 "SELECT 'fact:' || id AS item_key, embedding FROM facts WHERE embedding IS NOT NULL
1091 UNION ALL
1092 SELECT 'chunk:' || id AS item_key, embedding FROM chunks WHERE embedding IS NOT NULL
1093 UNION ALL
1094 SELECT 'msg:' || id AS item_key, embedding FROM messages WHERE embedding IS NOT NULL
1095 UNION ALL
1096 SELECT 'episode:' || episode_id AS item_key, embedding FROM episodes WHERE embedding IS NOT NULL",
1097 )?;
1098 let rows = stmt.query_map([], |row| {
1099 Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?))
1100 })?;
1101
1102 let mut entries = Vec::new();
1103 for row in rows {
1104 let (item_key, blob) = row?;
1105 entries.push((item_key, source_embedding_digest(&blob, dim)?));
1106 }
1107 entries.sort_unstable();
1108
1109 let mut builder = DigestBuilder::new();
1110 builder
1111 .update_str("semantic-memory.vector_source_snapshot.v1")
1112 .separator()
1113 .update(&(dim as u64).to_le_bytes())
1114 .separator();
1115 for (item_key, source_embedding_digest) in &entries {
1116 builder
1117 .update_str(item_key)
1118 .separator()
1119 .update_str(source_embedding_digest)
1120 .separator();
1121 }
1122 Ok((
1123 format!("blake3:{}", builder.finalize().hex()),
1124 entries.len(),
1125 ))
1126}
1127
1128#[cfg(any(
1129 feature = "turbo-quant-codec",
1130 feature = "fib-quant-codec",
1131 feature = "per-dim-codec"
1132))]
1133fn derived_artifact_manifest_digest(rows: &[DerivedVectorArtifactRow]) -> String {
1134 let mut entries = rows
1135 .iter()
1136 .map(|row| {
1137 (
1138 row.item_key.as_str(),
1139 row.source_embedding_digest.as_str(),
1140 row.encoded_digest.as_str(),
1141 )
1142 })
1143 .collect::<Vec<_>>();
1144 entries.sort_unstable();
1145
1146 let mut builder = DigestBuilder::new();
1147 builder
1148 .update_str("semantic-memory.vector_artifact_manifest.v1")
1149 .separator();
1150 for (item_key, source_embedding_digest, encoded_digest) in entries {
1151 builder
1152 .update_str(item_key)
1153 .separator()
1154 .update_str(source_embedding_digest)
1155 .separator()
1156 .update_str(encoded_digest)
1157 .separator();
1158 }
1159 format!("blake3:{}", builder.finalize().hex())
1160}
1161
1162#[cfg(any(
1163 feature = "turbo-quant-codec",
1164 feature = "fib-quant-codec",
1165 feature = "per-dim-codec"
1166))]
1167pub(crate) fn upsert_derived_vector_artifact(
1168 conn: &Connection,
1169 row: &DerivedVectorArtifactRow,
1170) -> Result<(), MemoryError> {
1171 conn.execute(
1172 "INSERT OR REPLACE INTO derived_vector_artifacts
1173 (item_key, generation_id, codec_family, codec_profile_digest, source_embedding_digest,
1174 encoded_digest, artifact_digest, encoding, dim, encoded, created_at, status,
1175 codec_governance_receipt_id, codec_profile, degradation_budget, raw_source_artifact_id)
1176 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?6, ?7, ?8, ?9, datetime('now'), ?10, ?11, ?12, ?13, ?14)",
1177 params![
1178 row.item_key,
1179 row.generation_id.as_deref(),
1180 row.codec_family,
1181 row.codec_profile_digest,
1182 row.source_embedding_digest,
1183 row.encoded_digest,
1184 row.encoding,
1185 i64::try_from(row.dim)
1186 .map_err(|err| MemoryError::Other(format!("artifact dim overflow: {err}")))?,
1187 row.encoded,
1188 row.status,
1189 row.codec_governance_receipt_id.as_deref(),
1190 row.codec_profile.as_deref(),
1191 row.degradation_budget,
1192 row.raw_source_artifact_id.as_deref(),
1193 ],
1194 )?;
1195 Ok(())
1196}
1197
1198pub fn delete_derived_vector_artifact(
1199 conn: &Connection,
1200 item_key: &str,
1201) -> Result<(), MemoryError> {
1202 conn.execute(
1203 "DELETE FROM derived_vector_artifacts WHERE item_key = ?1",
1204 params![item_key],
1205 )?;
1206 Ok(())
1207}
1208
1209pub fn invalidate_derived_vector_artifact(
1210 conn: &Connection,
1211 item_key: &str,
1212) -> Result<(), MemoryError> {
1213 conn.execute(
1214 "UPDATE derived_vector_artifacts
1215 SET status = 'invalidated'
1216 WHERE item_key = ?1 AND status = 'active'",
1217 params![item_key],
1218 )?;
1219 conn.execute(
1220 "UPDATE derived_vector_artifact_generations
1221 SET status = 'invalidated'
1222 WHERE status = 'active'",
1223 [],
1224 )?;
1225 Ok(())
1226}
1227
1228#[cfg(any(
1229 feature = "turbo-quant-codec",
1230 feature = "fib-quant-codec",
1231 feature = "per-dim-codec"
1232))]
1233#[allow(dead_code)]
1234pub(crate) fn load_derived_vector_artifacts_by_profile(
1235 conn: &Connection,
1236 codec_family: &str,
1237 codec_profile_digest: &str,
1238) -> Result<Vec<DerivedVectorArtifactRow>, MemoryError> {
1239 let mut stmt = conn.prepare(
1240 "SELECT item_key, generation_id, codec_family, codec_profile_digest, source_embedding_digest,
1241 encoded_digest, encoding, dim, status, encoded,
1242 codec_governance_receipt_id, codec_profile, degradation_budget, raw_source_artifact_id
1243 FROM derived_vector_artifacts
1244 WHERE codec_family = ?1 AND codec_profile_digest = ?2 AND status = 'active'",
1245 )?;
1246 let rows = stmt.query_map(params![codec_family, codec_profile_digest], |row| {
1247 let dim_i64: i64 = row.get(7)?;
1248 Ok(DerivedVectorArtifactRow {
1249 item_key: row.get(0)?,
1250 generation_id: row.get(1)?,
1251 codec_family: row.get(2)?,
1252 codec_profile_digest: row.get(3)?,
1253 source_embedding_digest: row.get(4)?,
1254 encoded_digest: row.get(5)?,
1255 encoding: row.get(6)?,
1256 dim: usize::try_from(dim_i64).map_err(|err| {
1257 rusqlite::Error::FromSqlConversionFailure(
1258 7,
1259 rusqlite::types::Type::Integer,
1260 Box::new(err),
1261 )
1262 })?,
1263 status: row.get(8)?,
1264 encoded: row.get(9)?,
1265 codec_governance_receipt_id: row.get(10)?,
1266 codec_profile: row.get(11)?,
1267 degradation_budget: row.get(12)?,
1268 raw_source_artifact_id: row.get(13)?,
1269 })
1270 })?;
1271
1272 let mut artifacts = Vec::new();
1273 for row in rows {
1274 artifacts.push(row?);
1275 }
1276 Ok(artifacts)
1277}
1278
1279#[cfg(any(
1280 feature = "turbo-quant-codec",
1281 feature = "fib-quant-codec",
1282 feature = "per-dim-codec"
1283))]
1284pub(crate) fn load_derived_vector_artifacts_by_generation(
1285 conn: &Connection,
1286 generation_id: &str,
1287) -> Result<Vec<DerivedVectorArtifactRow>, MemoryError> {
1288 let mut stmt = conn.prepare(
1289 "SELECT item_key, generation_id, codec_family, codec_profile_digest, source_embedding_digest,
1290 encoded_digest, encoding, dim, status, encoded,
1291 codec_governance_receipt_id, codec_profile, degradation_budget, raw_source_artifact_id
1292 FROM derived_vector_artifacts
1293 WHERE generation_id = ?1 AND status = 'active'",
1294 )?;
1295 let rows = stmt.query_map(params![generation_id], |row| {
1296 let dim_i64: i64 = row.get(7)?;
1297 Ok(DerivedVectorArtifactRow {
1298 item_key: row.get(0)?,
1299 generation_id: row.get(1)?,
1300 codec_family: row.get(2)?,
1301 codec_profile_digest: row.get(3)?,
1302 source_embedding_digest: row.get(4)?,
1303 encoded_digest: row.get(5)?,
1304 encoding: row.get(6)?,
1305 dim: usize::try_from(dim_i64).map_err(|err| {
1306 rusqlite::Error::FromSqlConversionFailure(
1307 7,
1308 rusqlite::types::Type::Integer,
1309 Box::new(err),
1310 )
1311 })?,
1312 status: row.get(8)?,
1313 encoded: row.get(9)?,
1314 codec_governance_receipt_id: row.get(10)?,
1315 codec_profile: row.get(11)?,
1316 degradation_budget: row.get(12)?,
1317 raw_source_artifact_id: row.get(13)?,
1318 })
1319 })?;
1320
1321 let mut artifacts = Vec::new();
1322 for row in rows {
1323 artifacts.push(row?);
1324 }
1325 Ok(artifacts)
1326}
1327
1328#[cfg(any(
1329 feature = "turbo-quant-codec",
1330 feature = "fib-quant-codec",
1331 feature = "per-dim-codec"
1332))]
1333pub(crate) fn current_derived_vector_generation(
1334 conn: &Connection,
1335 codec_family: &str,
1336 codec_profile_digest: &str,
1337) -> Result<Option<DerivedVectorArtifactGenerationRow>, MemoryError> {
1338 conn.query_row(
1339 "SELECT generation_id, codec_family, codec_profile_digest, source_snapshot_digest,
1340 source_row_count, artifact_count, dim, encoding, artifact_manifest_digest, status
1341 FROM derived_vector_artifact_generations
1342 WHERE codec_family = ?1 AND codec_profile_digest = ?2 AND status = 'active'
1343 ORDER BY created_at DESC
1344 LIMIT 1",
1345 params![codec_family, codec_profile_digest],
1346 |row| {
1347 let source_row_count: i64 = row.get(4)?;
1348 let artifact_count: i64 = row.get(5)?;
1349 let dim: i64 = row.get(6)?;
1350 Ok(DerivedVectorArtifactGenerationRow {
1351 generation_id: row.get(0)?,
1352 codec_family: row.get(1)?,
1353 codec_profile_digest: row.get(2)?,
1354 source_snapshot_digest: row.get(3)?,
1355 source_row_count: usize::try_from(source_row_count).map_err(|err| {
1356 rusqlite::Error::FromSqlConversionFailure(
1357 4,
1358 rusqlite::types::Type::Integer,
1359 Box::new(err),
1360 )
1361 })?,
1362 artifact_count: usize::try_from(artifact_count).map_err(|err| {
1363 rusqlite::Error::FromSqlConversionFailure(
1364 5,
1365 rusqlite::types::Type::Integer,
1366 Box::new(err),
1367 )
1368 })?,
1369 dim: usize::try_from(dim).map_err(|err| {
1370 rusqlite::Error::FromSqlConversionFailure(
1371 6,
1372 rusqlite::types::Type::Integer,
1373 Box::new(err),
1374 )
1375 })?,
1376 encoding: row.get(7)?,
1377 artifact_manifest_digest: row.get(8)?,
1378 status: row.get(9)?,
1379 })
1380 },
1381 )
1382 .optional()
1383 .map_err(MemoryError::from)
1384}
1385
1386pub fn count_derived_vector_artifacts(
1387 conn: &Connection,
1388 codec_family: &str,
1389 codec_profile_digest: &str,
1390) -> Result<usize, MemoryError> {
1391 let count: i64 = conn.query_row(
1392 "SELECT COUNT(*) FROM derived_vector_artifacts
1393 WHERE codec_family = ?1 AND codec_profile_digest = ?2 AND status = 'active'",
1394 params![codec_family, codec_profile_digest],
1395 |row| row.get(0),
1396 )?;
1397 usize::try_from(count)
1398 .map_err(|err| MemoryError::Other(format!("derived artifact count overflow: {err}")))
1399}
1400
1401#[cfg(feature = "turbo-quant-codec")]
1402pub(crate) fn rebuild_turbo_quant_artifacts(
1403 conn: &Connection,
1404 dim: usize,
1405 bits: u8,
1406 projections: usize,
1407 seed: u64,
1408) -> Result<VectorArtifactBuildReceiptV1, MemoryError> {
1409 use crate::vector_codec::{TurboQuantCodec, VectorCodec};
1410
1411 let started = std::time::Instant::now();
1412 let codec = TurboQuantCodec::new(dim, bits, projections, seed)?;
1413 let codebook_json = "{}".to_string();
1414 let codec_profile_digest = codec.profile().digest();
1415 let generation_id = uuid::Uuid::new_v4().to_string();
1416 let mut source_row_count = 0usize;
1417 let mut artifact_count = 0usize;
1418 let mut skipped_row_count = 0usize;
1419 let mut degradations = Vec::new();
1420
1421 let mut stmt = conn.prepare(
1422 "SELECT 'fact:' || id AS item_key, embedding FROM facts WHERE embedding IS NOT NULL
1423 UNION ALL
1424 SELECT 'chunk:' || id AS item_key, embedding FROM chunks WHERE embedding IS NOT NULL
1425 UNION ALL
1426 SELECT 'msg:' || id AS item_key, embedding FROM messages WHERE embedding IS NOT NULL
1427 UNION ALL
1428 SELECT 'episode:' || episode_id AS item_key, embedding FROM episodes WHERE embedding IS NOT NULL",
1429 )?;
1430 let rows = stmt.query_map([], |row| {
1431 Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?))
1432 })?;
1433
1434 let mut pending = Vec::new();
1435 for row in rows {
1436 let (item_key, blob) = row?;
1437 source_row_count += 1;
1438 let embedding = match decode_f32_le(&blob, dim) {
1439 Ok(embedding) => embedding,
1440 Err(err) => {
1441 skipped_row_count += 1;
1442 degradations.push(format!(
1443 "skipped {item_key}: invalid authoritative embedding: {err}"
1444 ));
1445 continue;
1446 }
1447 };
1448 let artifact = match codec.encode(&embedding) {
1449 Ok(artifact) => artifact,
1450 Err(err) => {
1451 skipped_row_count += 1;
1452 degradations.push(format!("skipped {item_key}: encode failed: {err}"));
1453 continue;
1454 }
1455 };
1456 pending.push(DerivedVectorArtifactRow {
1457 item_key,
1458 generation_id: Some(generation_id.clone()),
1459 codec_family: "turbo_quant".to_string(),
1460 codec_profile_digest: codec_profile_digest.clone(),
1461 source_embedding_digest: source_embedding_digest(&blob, dim)?,
1462 encoded_digest: artifact.artifact_digest,
1463 encoding: "turbo_code_wire_v1".to_string(),
1464 dim,
1465 status: "active".to_string(),
1466 encoded: artifact.encoded,
1467 codec_governance_receipt_id: None,
1470 codec_profile: None,
1471 degradation_budget: None,
1472 raw_source_artifact_id: None,
1473 });
1474 }
1475 drop(stmt);
1476
1477 let build_receipt_id = uuid::Uuid::new_v4().to_string();
1478 let source_snapshot_digest = source_snapshot_digest(&pending, dim);
1479 let artifact_manifest_digest = derived_artifact_manifest_digest(&pending);
1480 let source_tables = vec![
1481 "facts".to_string(),
1482 "chunks".to_string(),
1483 "messages".to_string(),
1484 "episodes".to_string(),
1485 ];
1486 let generation_manifest = DerivedVectorArtifactGenerationV1 {
1487 schema_version: "derived_vector_artifact_generation_v1".to_string(),
1488 generation_id: generation_id.clone(),
1489 codec_family: "turbo_quant".to_string(),
1490 codec_profile_digest: codec_profile_digest.clone(),
1491 source_snapshot_digest: source_snapshot_digest.clone(),
1492 source_row_count,
1493 artifact_count: pending.len(),
1494 source_tables,
1495 dim,
1496 encoding: "turbo_code_wire_v1".to_string(),
1497 created_at: Utc::now(),
1498 build_receipt_id: Some(build_receipt_id.clone()),
1499 artifact_manifest_digest: artifact_manifest_digest.clone(),
1500 status: if skipped_row_count == 0 {
1501 "active".to_string()
1502 } else {
1503 "failed".to_string()
1504 },
1505 degradations: degradations.clone(),
1506 };
1507
1508 with_transaction(conn, |tx| {
1509 tx.execute(
1510 "UPDATE derived_vector_artifact_generations
1511 SET status = 'superseded'
1512 WHERE codec_family = ?1 AND codec_profile_digest = ?2 AND status = 'active'",
1513 params!["turbo_quant", &codec_profile_digest],
1514 )?;
1515 tx.execute(
1516 "DELETE FROM derived_vector_artifacts
1517 WHERE codec_family = ?1 AND codec_profile_digest = ?2",
1518 params!["turbo_quant", &codec_profile_digest],
1519 )?;
1520 tx.execute(
1521 "INSERT INTO derived_vector_artifact_generations
1522 (generation_id, schema_version, codec_family, codec_profile_digest,
1523 source_snapshot_digest, source_row_count, artifact_count, source_tables_json,
1524 dim, encoding, created_at, build_receipt_id, artifact_manifest_digest,
1525 status, degradations_json, config_json)
1526 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)",
1527 params![
1528 generation_manifest.generation_id,
1529 generation_manifest.schema_version,
1530 generation_manifest.codec_family,
1531 generation_manifest.codec_profile_digest,
1532 generation_manifest.source_snapshot_digest,
1533 i64::try_from(generation_manifest.source_row_count).map_err(|err| {
1534 MemoryError::Other(format!("source row count overflow: {err}"))
1535 })?,
1536 i64::try_from(generation_manifest.artifact_count).map_err(|err| {
1537 MemoryError::Other(format!("artifact count overflow: {err}"))
1538 })?,
1539 serde_json::to_string(&generation_manifest.source_tables)
1540 .map_err(|err| MemoryError::Other(err.to_string()))?,
1541 i64::try_from(generation_manifest.dim)
1542 .map_err(|err| MemoryError::Other(format!("artifact dim overflow: {err}")))?,
1543 generation_manifest.encoding,
1544 generation_manifest.created_at.to_rfc3339(),
1545 generation_manifest.build_receipt_id,
1546 generation_manifest.artifact_manifest_digest,
1547 generation_manifest.status,
1548 serde_json::to_string(&generation_manifest.degradations)
1549 .map_err(|err| MemoryError::Other(err.to_string()))?,
1550 codebook_json,
1551 ],
1552 )?;
1553 for row in &pending {
1554 upsert_derived_vector_artifact(tx, row)?;
1555 artifact_count += 1;
1556 }
1557 Ok(())
1558 })?;
1559
1560 Ok(VectorArtifactBuildReceiptV1 {
1561 schema_version: "vector_artifact_build_receipt_v1".to_string(),
1562 codec_family: "turbo_quant".to_string(),
1563 codec_profile_digest,
1564 source_row_count,
1565 artifact_count,
1566 generation_id: Some(generation_id),
1567 source_snapshot_digest: Some(source_snapshot_digest),
1568 artifact_manifest_digest: Some(artifact_manifest_digest),
1569 build_receipt_id: Some(build_receipt_id),
1570 skipped_row_count,
1571 elapsed_ms: started.elapsed().as_millis(),
1572 created_at: Utc::now(),
1573 degradations,
1574 })
1575}
1576
1577#[cfg(feature = "per-dim-codec")]
1578pub(crate) fn rebuild_per_dim_artifacts(
1579 conn: &Connection,
1580 dim: usize,
1581 bits: u8,
1582) -> Result<VectorArtifactBuildReceiptV1, MemoryError> {
1583 use crate::vector_codec::{PerDimCodec, VectorCodec};
1584
1585 let started = std::time::Instant::now();
1586 let codec = PerDimCodec::new(dim, bits)?;
1587 let codebook_json = codec.config_json()?;
1588 let codec_profile_digest = codec.profile().digest();
1589 let generation_id = uuid::Uuid::new_v4().to_string();
1590 let mut source_row_count = 0usize;
1591 let mut artifact_count = 0usize;
1592 let mut skipped_row_count = 0usize;
1593 let mut degradations = Vec::new();
1594
1595 let mut stmt = conn.prepare(
1596 "SELECT 'fact:' || id AS item_key, embedding FROM facts WHERE embedding IS NOT NULL
1597 UNION ALL
1598 SELECT 'chunk:' || id AS item_key, embedding FROM chunks WHERE embedding IS NOT NULL
1599 UNION ALL
1600 SELECT 'msg:' || id AS item_key, embedding FROM messages WHERE embedding IS NOT NULL
1601 UNION ALL
1602 SELECT 'episode:' || episode_id AS item_key, embedding FROM episodes WHERE embedding IS NOT NULL",
1603 )?;
1604 let rows = stmt.query_map([], |row| {
1605 Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?))
1606 })?;
1607
1608 let mut pending = Vec::new();
1609 for row in rows {
1610 let (item_key, blob) = row?;
1611 source_row_count += 1;
1612 let embedding = match decode_f32_le(&blob, dim) {
1613 Ok(embedding) => embedding,
1614 Err(err) => {
1615 skipped_row_count += 1;
1616 degradations.push(format!(
1617 "skipped {item_key}: invalid authoritative embedding: {err}"
1618 ));
1619 continue;
1620 }
1621 };
1622 let artifact = match codec.encode(&embedding) {
1623 Ok(artifact) => artifact,
1624 Err(err) => {
1625 skipped_row_count += 1;
1626 degradations.push(format!("skipped {item_key}: encode failed: {err}"));
1627 continue;
1628 }
1629 };
1630 pending.push(DerivedVectorArtifactRow {
1631 item_key,
1632 generation_id: Some(generation_id.clone()),
1633 codec_family: "per_dim".to_string(),
1634 codec_profile_digest: codec_profile_digest.clone(),
1635 source_embedding_digest: source_embedding_digest(&blob, dim)?,
1636 encoded_digest: artifact.artifact_digest,
1637 encoding: "per_dim_wire_v1".to_string(),
1638 dim,
1639 status: "active".to_string(),
1640 encoded: artifact.encoded,
1641 codec_governance_receipt_id: None,
1644 codec_profile: None,
1645 degradation_budget: None,
1646 raw_source_artifact_id: None,
1647 });
1648 }
1649 drop(stmt);
1650
1651 let build_receipt_id = uuid::Uuid::new_v4().to_string();
1652 let source_snapshot_digest = source_snapshot_digest(&pending, dim);
1653 let artifact_manifest_digest = derived_artifact_manifest_digest(&pending);
1654 let source_tables = vec![
1655 "facts".to_string(),
1656 "chunks".to_string(),
1657 "messages".to_string(),
1658 "episodes".to_string(),
1659 ];
1660 let generation_manifest = DerivedVectorArtifactGenerationV1 {
1661 schema_version: "derived_vector_artifact_generation_v1".to_string(),
1662 generation_id: generation_id.clone(),
1663 codec_family: "per_dim".to_string(),
1664 codec_profile_digest: codec_profile_digest.clone(),
1665 source_snapshot_digest: source_snapshot_digest.clone(),
1666 source_row_count,
1667 artifact_count: pending.len(),
1668 source_tables,
1669 dim,
1670 encoding: "per_dim_wire_v1".to_string(),
1671 created_at: Utc::now(),
1672 build_receipt_id: Some(build_receipt_id.clone()),
1673 artifact_manifest_digest: artifact_manifest_digest.clone(),
1674 status: if skipped_row_count == 0 {
1675 "active".to_string()
1676 } else {
1677 "failed".to_string()
1678 },
1679 degradations: degradations.clone(),
1680 };
1681
1682 with_transaction(conn, |tx| {
1683 tx.execute(
1684 "UPDATE derived_vector_artifact_generations
1685 SET status = 'superseded'
1686 WHERE codec_family = ?1 AND codec_profile_digest = ?2 AND status = 'active'",
1687 params!["per_dim", &codec_profile_digest],
1688 )?;
1689 tx.execute(
1690 "DELETE FROM derived_vector_artifacts
1691 WHERE codec_family = ?1 AND codec_profile_digest = ?2",
1692 params!["per_dim", &codec_profile_digest],
1693 )?;
1694 tx.execute(
1695 "INSERT INTO derived_vector_artifact_generations
1696 (generation_id, schema_version, codec_family, codec_profile_digest,
1697 source_snapshot_digest, source_row_count, artifact_count, source_tables_json,
1698 dim, encoding, created_at, build_receipt_id, artifact_manifest_digest,
1699 status, degradations_json, config_json)
1700 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)",
1701 params![
1702 generation_manifest.generation_id,
1703 generation_manifest.schema_version,
1704 generation_manifest.codec_family,
1705 generation_manifest.codec_profile_digest,
1706 generation_manifest.source_snapshot_digest,
1707 i64::try_from(generation_manifest.source_row_count).map_err(|err| {
1708 MemoryError::Other(format!("source row count overflow: {err}"))
1709 })?,
1710 i64::try_from(generation_manifest.artifact_count).map_err(|err| {
1711 MemoryError::Other(format!("artifact count overflow: {err}"))
1712 })?,
1713 serde_json::to_string(&generation_manifest.source_tables)
1714 .map_err(|err| MemoryError::Other(err.to_string()))?,
1715 i64::try_from(generation_manifest.dim)
1716 .map_err(|err| MemoryError::Other(format!("artifact dim overflow: {err}")))?,
1717 generation_manifest.encoding,
1718 generation_manifest.created_at.to_rfc3339(),
1719 generation_manifest.build_receipt_id,
1720 generation_manifest.artifact_manifest_digest,
1721 generation_manifest.status,
1722 serde_json::to_string(&generation_manifest.degradations)
1723 .map_err(|err| MemoryError::Other(err.to_string()))?,
1724 codebook_json,
1725 ],
1726 )?;
1727 for row in &pending {
1728 upsert_derived_vector_artifact(tx, row)?;
1729 artifact_count += 1;
1730 }
1731 Ok(())
1732 })?;
1733
1734 Ok(VectorArtifactBuildReceiptV1 {
1735 schema_version: "vector_artifact_build_receipt_v1".to_string(),
1736 codec_family: "per_dim".to_string(),
1737 codec_profile_digest,
1738 source_row_count,
1739 artifact_count,
1740 generation_id: Some(generation_id),
1741 source_snapshot_digest: Some(source_snapshot_digest),
1742 artifact_manifest_digest: Some(artifact_manifest_digest),
1743 build_receipt_id: Some(build_receipt_id),
1744 skipped_row_count,
1745 elapsed_ms: started.elapsed().as_millis(),
1746 created_at: Utc::now(),
1747 degradations,
1748 })
1749}
1750
1751#[cfg(feature = "fib-quant-codec")]
1752pub(crate) fn rebuild_fib_quant_artifacts(
1753 conn: &Connection,
1754 dim: usize,
1755 block_count: usize,
1756) -> Result<VectorArtifactBuildReceiptV1, MemoryError> {
1757 use crate::vector_codec::{FibQuantCodec, VectorCodec};
1758
1759 let started = std::time::Instant::now();
1760 let codec = FibQuantCodec::new(dim, block_count, 8)?;
1761 let codebook_json = serde_json::to_string(codec.codebook())
1762 .map_err(|e| MemoryError::QuantizationError(format!("serialize codebook: {e}")))?;
1763 let codec_profile_digest = codec.profile().digest();
1764 let generation_id = uuid::Uuid::new_v4().to_string();
1765 let mut source_row_count = 0usize;
1766 let mut artifact_count = 0usize;
1767 let mut skipped_row_count = 0usize;
1768 let mut degradations = Vec::new();
1769
1770 let mut stmt = conn.prepare(
1771 "SELECT 'fact:' || id AS item_key, embedding FROM facts WHERE embedding IS NOT NULL
1772 UNION ALL
1773 SELECT 'chunk:' || id AS item_key, embedding FROM chunks WHERE embedding IS NOT NULL
1774 UNION ALL
1775 SELECT 'msg:' || id AS item_key, embedding FROM messages WHERE embedding IS NOT NULL
1776 UNION ALL
1777 SELECT 'episode:' || episode_id AS item_key, embedding FROM episodes WHERE embedding IS NOT NULL",
1778 )?;
1779 let rows = stmt.query_map([], |row| {
1780 Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?))
1781 })?;
1782
1783 let mut pending = Vec::new();
1784 for row in rows {
1785 let (item_key, blob) = row?;
1786 source_row_count += 1;
1787 let embedding = match decode_f32_le(&blob, dim) {
1788 Ok(embedding) => embedding,
1789 Err(err) => {
1790 skipped_row_count += 1;
1791 degradations.push(format!(
1792 "skipped {item_key}: invalid authoritative embedding: {err}"
1793 ));
1794 continue;
1795 }
1796 };
1797 let artifact = match codec.encode(&embedding) {
1798 Ok(artifact) => artifact,
1799 Err(err) => {
1800 skipped_row_count += 1;
1801 degradations.push(format!("skipped {item_key}: encode failed: {err}"));
1802 continue;
1803 }
1804 };
1805 pending.push(DerivedVectorArtifactRow {
1806 item_key,
1807 generation_id: Some(generation_id.clone()),
1808 codec_family: "fib_quant".to_string(),
1809 codec_profile_digest: codec_profile_digest.clone(),
1810 source_embedding_digest: source_embedding_digest(&blob, dim)?,
1811 encoded_digest: artifact.artifact_digest,
1812 encoding: "fib_code_wire_v1".to_string(),
1813 dim,
1814 status: "active".to_string(),
1815 encoded: artifact.encoded,
1816 codec_governance_receipt_id: None,
1819 codec_profile: None,
1820 degradation_budget: None,
1821 raw_source_artifact_id: None,
1822 });
1823 }
1824 drop(stmt);
1825
1826 let build_receipt_id = uuid::Uuid::new_v4().to_string();
1827 let source_snapshot_digest = source_snapshot_digest(&pending, dim);
1828 let artifact_manifest_digest = derived_artifact_manifest_digest(&pending);
1829 let source_tables = vec![
1830 "facts".to_string(),
1831 "chunks".to_string(),
1832 "messages".to_string(),
1833 "episodes".to_string(),
1834 ];
1835 let generation_manifest = DerivedVectorArtifactGenerationV1 {
1836 schema_version: "derived_vector_artifact_generation_v1".to_string(),
1837 generation_id: generation_id.clone(),
1838 codec_family: "fib_quant".to_string(),
1839 codec_profile_digest: codec_profile_digest.clone(),
1840 source_snapshot_digest: source_snapshot_digest.clone(),
1841 source_row_count,
1842 artifact_count: pending.len(),
1843 source_tables,
1844 dim,
1845 encoding: "fib_code_wire_v1".to_string(),
1846 created_at: Utc::now(),
1847 build_receipt_id: Some(build_receipt_id.clone()),
1848 artifact_manifest_digest: artifact_manifest_digest.clone(),
1849 status: if skipped_row_count == 0 {
1850 "active".to_string()
1851 } else {
1852 "failed".to_string()
1853 },
1854 degradations: degradations.clone(),
1855 };
1856
1857 with_transaction(conn, |tx| {
1858 tx.execute(
1859 "UPDATE derived_vector_artifact_generations
1860 SET status = 'superseded'
1861 WHERE codec_family = ?1 AND codec_profile_digest = ?2 AND status = 'active'",
1862 params!["fib_quant", &codec_profile_digest],
1863 )?;
1864 tx.execute(
1865 "DELETE FROM derived_vector_artifacts
1866 WHERE codec_family = ?1 AND codec_profile_digest = ?2",
1867 params!["fib_quant", &codec_profile_digest],
1868 )?;
1869 tx.execute(
1870 "INSERT INTO derived_vector_artifact_generations
1871 (generation_id, schema_version, codec_family, codec_profile_digest,
1872 source_snapshot_digest, source_row_count, artifact_count, source_tables_json,
1873 dim, encoding, created_at, build_receipt_id, artifact_manifest_digest,
1874 status, degradations_json, config_json)
1875 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14, ?15, ?16)",
1876 params![
1877 generation_manifest.generation_id,
1878 generation_manifest.schema_version,
1879 generation_manifest.codec_family,
1880 generation_manifest.codec_profile_digest,
1881 generation_manifest.source_snapshot_digest,
1882 i64::try_from(generation_manifest.source_row_count).map_err(|err| {
1883 MemoryError::Other(format!("source row count overflow: {err}"))
1884 })?,
1885 i64::try_from(generation_manifest.artifact_count).map_err(|err| {
1886 MemoryError::Other(format!("artifact count overflow: {err}"))
1887 })?,
1888 serde_json::to_string(&generation_manifest.source_tables)
1889 .map_err(|err| MemoryError::Other(err.to_string()))?,
1890 i64::try_from(generation_manifest.dim)
1891 .map_err(|err| MemoryError::Other(format!("artifact dim overflow: {err}")))?,
1892 generation_manifest.encoding,
1893 generation_manifest.created_at.to_rfc3339(),
1894 generation_manifest.build_receipt_id,
1895 generation_manifest.artifact_manifest_digest,
1896 generation_manifest.status,
1897 serde_json::to_string(&generation_manifest.degradations)
1898 .map_err(|err| MemoryError::Other(err.to_string()))?,
1899 codebook_json,
1900 ],
1901 )?;
1902 for row in &pending {
1903 upsert_derived_vector_artifact(tx, row)?;
1904 artifact_count += 1;
1905 }
1906 Ok(())
1907 })?;
1908
1909 Ok(VectorArtifactBuildReceiptV1 {
1910 schema_version: "vector_artifact_build_receipt_v1".to_string(),
1911 codec_family: "fib_quant".to_string(),
1912 codec_profile_digest,
1913 source_row_count,
1914 artifact_count,
1915 generation_id: Some(generation_id),
1916 source_snapshot_digest: Some(source_snapshot_digest),
1917 artifact_manifest_digest: Some(artifact_manifest_digest),
1918 build_receipt_id: Some(build_receipt_id),
1919 skipped_row_count,
1920 elapsed_ms: started.elapsed().as_millis(),
1921 created_at: Utc::now(),
1922 degradations,
1923 })
1924}
1925
1926#[cfg(feature = "fib-quant-codec")]
1927pub(crate) fn load_fib_codebook(
1928 conn: &Connection,
1929 profile_digest: &str,
1930) -> Result<fib_quant::FibCodebookV1, MemoryError> {
1931 let json: String = conn.query_row(
1932 "SELECT config_json FROM derived_vector_artifact_generations WHERE codec_family='fib_quant' AND codec_profile_digest=?1 AND status='active' ORDER BY created_at DESC LIMIT 1",
1933 [profile_digest],
1934 |row| row.get(0),
1935 )?;
1936 serde_json::from_str(&json)
1937 .map_err(|e| MemoryError::QuantizationError(format!("deserialize codebook: {e}")))
1938}
1939
1940fn receipt_count_to_u64(value: usize, field: &'static str) -> Result<u64, MemoryError> {
1941 u64::try_from(value).map_err(|err| MemoryError::Other(format!("{field} is too large: {err}")))
1942}
1943
1944fn receipt_count_to_i64(value: u64, field: &'static str) -> Result<i64, MemoryError> {
1945 i64::try_from(value).map_err(|err| MemoryError::Other(format!("{field} is too large: {err}")))
1946}
1947
1948fn receipt_count_to_usize(
1949 value: u64,
1950 receipt_id: &str,
1951 field: &'static str,
1952) -> Result<usize, MemoryError> {
1953 usize::try_from(value).map_err(|err| MemoryError::CorruptData {
1954 table: "search_receipts",
1955 row_id: receipt_id.to_string(),
1956 detail: format!("{field} does not fit this platform: {err}"),
1957 })
1958}
1959
1960fn stored_search_receipt(
1961 receipt: &VectorSearchReceiptV1,
1962) -> Result<StoredVectorSearchReceiptV1, MemoryError> {
1963 Ok(StoredVectorSearchReceiptV1 {
1964 schema_version: SEARCH_RECEIPT_SCHEMA_VERSION.to_string(),
1965 receipt_id: receipt.receipt_id.clone(),
1966 evaluation_time: receipt.evaluation_time,
1967 receipt_digest: receipt.receipt_digest.clone(),
1968 trace_id: receipt.trace_id.clone(),
1969 attempt_family_id: receipt.attempt_family_id.clone(),
1970 attempt_id: receipt.attempt_id.clone(),
1971 replay_of: receipt.replay_of.clone(),
1972 query_embedding_digest: receipt.query_embedding_digest.clone(),
1973 query_text_digest: receipt.query_text_digest.clone(),
1974 query_input_digest: receipt.query_input_digest.clone(),
1975 filter_digest: receipt.filter_digest.clone(),
1976 redaction_state: receipt.redaction_state.clone(),
1977 budget_id: receipt.budget_id.clone(),
1978 deadline_at: receipt.deadline_at,
1979 search_profile: receipt.search_profile.clone(),
1980 candidate_backend: receipt.candidate_backend.clone(),
1981 codec_family: receipt.codec_family.clone(),
1982 codec_profile_digest: receipt.codec_profile_digest.clone(),
1983 artifact_profile_digest: receipt.artifact_profile_digest.clone(),
1984 artifact_count: receipt
1985 .artifact_count
1986 .map(|value| receipt_count_to_u64(value, "artifact_count"))
1987 .transpose()?,
1988 artifact_corruption_count: receipt
1989 .artifact_corruption_count
1990 .map(|value| receipt_count_to_u64(value, "artifact_corruption_count"))
1991 .transpose()?,
1992 artifact_missing_count: receipt
1993 .artifact_missing_count
1994 .map(|value| receipt_count_to_u64(value, "artifact_missing_count"))
1995 .transpose()?,
1996 vector_artifact_manifest_digest: receipt.vector_artifact_manifest_digest.clone(),
1997 artifact_generation_id: receipt.artifact_generation_id.clone(),
1998 approximate_scanned_count: receipt
1999 .approximate_scanned_count
2000 .map(|value| receipt_count_to_u64(value, "approximate_scanned_count"))
2001 .transpose()?,
2002 approximate_returned_count: receipt
2003 .approximate_returned_count
2004 .map(|value| receipt_count_to_u64(value, "approximate_returned_count"))
2005 .transpose()?,
2006 raw_rows_loaded_count: receipt
2007 .raw_rows_loaded_count
2008 .map(|value| receipt_count_to_u64(value, "raw_rows_loaded_count"))
2009 .transpose()?,
2010 filter_strategy: receipt.filter_strategy.clone(),
2011 vector_artifact_count: receipt
2012 .vector_artifact_count
2013 .map(|value| receipt_count_to_u64(value, "vector_artifact_count"))
2014 .transpose()?,
2015 vector_artifact_missing_count: receipt
2016 .vector_artifact_missing_count
2017 .map(|value| receipt_count_to_u64(value, "vector_artifact_missing_count"))
2018 .transpose()?,
2019 vector_artifact_stale_count: receipt
2020 .vector_artifact_stale_count
2021 .map(|value| receipt_count_to_u64(value, "vector_artifact_stale_count"))
2022 .transpose()?,
2023 exact_rerank_count: receipt
2024 .exact_rerank_count
2025 .map(|value| receipt_count_to_u64(value, "exact_rerank_count"))
2026 .transpose()?,
2027 approximate_candidate_count: receipt
2028 .approximate_candidate_count
2029 .map(|value| receipt_count_to_u64(value, "approximate_candidate_count"))
2030 .transpose()?,
2031 fallback_reason: receipt.fallback_reason.clone(),
2032 approximate: receipt.approximate,
2033 requested_candidates: receipt_count_to_u64(
2034 receipt.requested_candidates,
2035 "requested_candidates",
2036 )?,
2037 returned_candidates: receipt_count_to_u64(
2038 receipt.returned_candidates,
2039 "returned_candidates",
2040 )?,
2041 post_filter_candidates: receipt_count_to_u64(
2042 receipt.post_filter_candidates,
2043 "post_filter_candidates",
2044 )?,
2045 fallback: receipt.fallback.clone(),
2046 exact_rerank: receipt.exact_rerank,
2047 result_ids: receipt.result_ids.clone(),
2048 degradations: receipt.degradations.clone(),
2049 })
2050}
2051
2052fn search_receipt_from_stored(
2053 stored: StoredVectorSearchReceiptV1,
2054) -> Result<VectorSearchReceiptV1, MemoryError> {
2055 if stored.schema_version != SEARCH_RECEIPT_SCHEMA_VERSION {
2056 return Err(MemoryError::CorruptData {
2057 table: "search_receipts",
2058 row_id: stored.receipt_id,
2059 detail: format!(
2060 "unsupported receipt schema version '{}'",
2061 stored.schema_version
2062 ),
2063 });
2064 }
2065
2066 Ok(VectorSearchReceiptV1 {
2067 schema_version: stored.schema_version.clone(),
2068 receipt_digest: stored.receipt_digest,
2069 receipt_id: stored.receipt_id.clone(),
2070 evaluation_time: stored.evaluation_time,
2071 trace_id: stored.trace_id,
2072 attempt_family_id: stored.attempt_family_id,
2073 attempt_id: stored.attempt_id,
2074 replay_of: stored.replay_of,
2075 query_embedding_digest: stored.query_embedding_digest,
2076 query_text_digest: stored.query_text_digest,
2077 query_input_digest: stored.query_input_digest,
2078 filter_digest: stored.filter_digest,
2079 redaction_state: stored.redaction_state,
2080 budget_id: stored.budget_id,
2081 deadline_at: stored.deadline_at,
2082 search_profile: stored.search_profile,
2083 candidate_backend: stored.candidate_backend,
2084 codec_family: stored.codec_family,
2085 codec_profile_digest: stored.codec_profile_digest,
2086 artifact_profile_digest: stored.artifact_profile_digest,
2087 artifact_count: stored
2088 .artifact_count
2089 .map(|value| receipt_count_to_usize(value, &stored.receipt_id, "artifact_count"))
2090 .transpose()?,
2091 artifact_corruption_count: stored
2092 .artifact_corruption_count
2093 .map(|value| {
2094 receipt_count_to_usize(value, &stored.receipt_id, "artifact_corruption_count")
2095 })
2096 .transpose()?,
2097 artifact_missing_count: stored
2098 .artifact_missing_count
2099 .map(|value| {
2100 receipt_count_to_usize(value, &stored.receipt_id, "artifact_missing_count")
2101 })
2102 .transpose()?,
2103 vector_artifact_manifest_digest: stored.vector_artifact_manifest_digest,
2104 artifact_generation_id: stored.artifact_generation_id,
2105 approximate_scanned_count: stored
2106 .approximate_scanned_count
2107 .map(|value| {
2108 receipt_count_to_usize(value, &stored.receipt_id, "approximate_scanned_count")
2109 })
2110 .transpose()?,
2111 approximate_returned_count: stored
2112 .approximate_returned_count
2113 .map(|value| {
2114 receipt_count_to_usize(value, &stored.receipt_id, "approximate_returned_count")
2115 })
2116 .transpose()?,
2117 raw_rows_loaded_count: stored
2118 .raw_rows_loaded_count
2119 .map(|value| receipt_count_to_usize(value, &stored.receipt_id, "raw_rows_loaded_count"))
2120 .transpose()?,
2121 filter_strategy: stored.filter_strategy,
2122 vector_artifact_count: stored
2123 .vector_artifact_count
2124 .map(|value| receipt_count_to_usize(value, &stored.receipt_id, "vector_artifact_count"))
2125 .transpose()?,
2126 vector_artifact_missing_count: stored
2127 .vector_artifact_missing_count
2128 .map(|value| {
2129 receipt_count_to_usize(value, &stored.receipt_id, "vector_artifact_missing_count")
2130 })
2131 .transpose()?,
2132 vector_artifact_stale_count: stored
2133 .vector_artifact_stale_count
2134 .map(|value| {
2135 receipt_count_to_usize(value, &stored.receipt_id, "vector_artifact_stale_count")
2136 })
2137 .transpose()?,
2138 exact_rerank_count: stored
2139 .exact_rerank_count
2140 .map(|value| receipt_count_to_usize(value, &stored.receipt_id, "exact_rerank_count"))
2141 .transpose()?,
2142 approximate_candidate_count: stored
2143 .approximate_candidate_count
2144 .map(|value| {
2145 receipt_count_to_usize(value, &stored.receipt_id, "approximate_candidate_count")
2146 })
2147 .transpose()?,
2148 fallback_reason: stored.fallback_reason,
2149 approximate: stored.approximate,
2150 requested_candidates: receipt_count_to_usize(
2151 stored.requested_candidates,
2152 &stored.receipt_id,
2153 "requested_candidates",
2154 )?,
2155 returned_candidates: receipt_count_to_usize(
2156 stored.returned_candidates,
2157 &stored.receipt_id,
2158 "returned_candidates",
2159 )?,
2160 post_filter_candidates: receipt_count_to_usize(
2161 stored.post_filter_candidates,
2162 &stored.receipt_id,
2163 "post_filter_candidates",
2164 )?,
2165 fallback: stored.fallback,
2166 exact_rerank: stored.exact_rerank,
2167 result_ids: stored.result_ids,
2168 degradations: stored.degradations,
2169 })
2170}
2171
2172pub fn store_search_receipt(
2177 conn: &Connection,
2178 receipt: &VectorSearchReceiptV1,
2179) -> Result<(), MemoryError> {
2180 let stored = stored_search_receipt(receipt)?;
2181 let receipt_json = serde_json::to_string(&stored)
2182 .map_err(|err| MemoryError::Other(format!("failed to serialize search receipt: {err}")))?;
2183 let receipt_digest = b3_digest(receipt_json.as_bytes());
2184
2185 let existing_digest: Option<String> = conn
2186 .query_row(
2187 "SELECT receipt_digest FROM search_receipts WHERE receipt_id = ?1",
2188 params![&stored.receipt_id],
2189 |row| row.get(0),
2190 )
2191 .optional()?;
2192 if let Some(existing_digest) = existing_digest {
2193 if existing_digest == receipt_digest {
2194 return Ok(());
2195 }
2196 return Err(MemoryError::SearchReceiptConflict {
2197 receipt_id: stored.receipt_id,
2198 });
2199 }
2200
2201 let result_ids_json = serde_json::to_string(&stored.result_ids).map_err(|err| {
2202 MemoryError::Other(format!(
2203 "failed to serialize search receipt result IDs: {err}"
2204 ))
2205 })?;
2206 conn.execute(
2207 "INSERT INTO search_receipts (
2208 receipt_id,
2209 schema_version,
2210 evaluation_time,
2211 search_profile,
2212 candidate_backend,
2213 approximate,
2214 exact_rerank,
2215 fallback,
2216 requested_candidates,
2217 returned_candidates,
2218 post_filter_candidates,
2219 result_ids_json,
2220 receipt_json,
2221 receipt_digest
2222 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13, ?14)",
2223 params![
2224 &stored.receipt_id,
2225 SEARCH_RECEIPT_SCHEMA_VERSION,
2226 stored.evaluation_time.to_rfc3339(),
2227 &stored.search_profile,
2228 &stored.candidate_backend,
2229 if stored.approximate { 1_i64 } else { 0_i64 },
2230 if stored.exact_rerank { 1_i64 } else { 0_i64 },
2231 &stored.fallback,
2232 receipt_count_to_i64(stored.requested_candidates, "requested_candidates")?,
2233 receipt_count_to_i64(stored.returned_candidates, "returned_candidates")?,
2234 receipt_count_to_i64(stored.post_filter_candidates, "post_filter_candidates")?,
2235 &result_ids_json,
2236 &receipt_json,
2237 &receipt_digest,
2238 ],
2239 )?;
2240 Ok(())
2241}
2242
2243pub fn get_search_receipt(
2245 conn: &Connection,
2246 receipt_id: &str,
2247) -> Result<Option<VectorSearchReceiptV1>, MemoryError> {
2248 let row: Option<(String, String, String)> = conn
2249 .query_row(
2250 "SELECT schema_version, receipt_json, receipt_digest
2251 FROM search_receipts
2252 WHERE receipt_id = ?1",
2253 params![receipt_id],
2254 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
2255 )
2256 .optional()?;
2257
2258 let Some((schema_version, receipt_json, receipt_digest)) = row else {
2259 return Ok(None);
2260 };
2261 if schema_version != SEARCH_RECEIPT_SCHEMA_VERSION {
2262 return Err(MemoryError::CorruptData {
2263 table: "search_receipts",
2264 row_id: receipt_id.to_string(),
2265 detail: format!("unsupported receipt schema version '{schema_version}'"),
2266 });
2267 }
2268
2269 let stored: StoredVectorSearchReceiptV1 =
2270 serde_json::from_str(&receipt_json).map_err(|err| MemoryError::CorruptData {
2271 table: "search_receipts",
2272 row_id: receipt_id.to_string(),
2273 detail: format!("invalid receipt JSON: {err}"),
2274 })?;
2275 let mut receipt = search_receipt_from_stored(stored)?;
2276 receipt.receipt_digest = Some(receipt_digest);
2277 Ok(Some(receipt))
2278}
2279
2280fn add_column_if_missing(
2281 conn: &Connection,
2282 table: &str,
2283 column: &str,
2284 column_sql: &str,
2285) -> Result<(), rusqlite::Error> {
2286 let pragma = format!("PRAGMA table_info({table})");
2287 let mut stmt = conn.prepare(&pragma)?;
2288 let exists = stmt
2289 .query_map([], |row| row.get::<_, String>(1))?
2290 .collect::<Result<Vec<_>, _>>()?
2291 .into_iter()
2292 .any(|name| name == column);
2293
2294 if !exists {
2295 conn.execute(
2296 &format!("ALTER TABLE {table} ADD COLUMN {column} {column_sql}"),
2297 [],
2298 )?;
2299 }
2300
2301 Ok(())
2302}
2303
2304pub fn check_embedding_metadata(
2306 conn: &Connection,
2307 config: &EmbeddingConfig,
2308) -> Result<(), MemoryError> {
2309 let existing: Option<(String, usize)> = conn
2311 .query_row(
2312 "SELECT model_name, dimensions FROM embedding_metadata WHERE id = 1",
2313 [],
2314 |row| Ok((row.get(0)?, row.get(1)?)),
2315 )
2316 .ok();
2317
2318 match existing {
2319 Some((model, dims)) => {
2320 if model != config.model || dims != config.dimensions {
2321 tracing::warn!(
2322 stored_model = %model,
2323 stored_dims = dims,
2324 configured_model = %config.model,
2325 configured_dims = config.dimensions,
2326 "Embedding model changed. Existing embeddings are stale."
2327 );
2328 conn.execute(
2329 "UPDATE embedding_metadata
2330 SET model_name = ?1,
2331 dimensions = ?2,
2332 embeddings_dirty = 1,
2333 updated_at = datetime('now')
2334 WHERE id = 1",
2335 params![config.model, config.dimensions],
2336 )?;
2337 }
2338 }
2339 None => {
2340 conn.execute(
2341 "INSERT INTO embedding_metadata (id, model_name, dimensions) VALUES (1, ?1, ?2)",
2342 params![config.model, config.dimensions],
2343 )?;
2344 }
2345 }
2346
2347 Ok(())
2348}
2349
2350pub fn embedding_to_bytes(embedding: &[f32]) -> Vec<u8> {
2352 encode_f32_le(embedding)
2353}
2354
2355pub fn encode_f32_le(values: &[f32]) -> Vec<u8> {
2357 let mut bytes = Vec::with_capacity(values.len() * 4);
2358 for value in values {
2359 bytes.extend_from_slice(&value.to_le_bytes());
2360 }
2361 bytes
2362}
2363
2364pub(crate) fn validate_embedding(values: &[f32], expected_dim: usize) -> Result<(), MemoryError> {
2366 if values.len() != expected_dim {
2367 return Err(MemoryError::EmbeddingDimensionMismatch {
2368 expected: expected_dim,
2369 actual: values.len(),
2370 });
2371 }
2372 if let Some((index, _)) = values
2373 .iter()
2374 .enumerate()
2375 .find(|(_, value)| !value.is_finite())
2376 {
2377 return Err(MemoryError::NonFiniteEmbeddingValue { index });
2378 }
2379 Ok(())
2380}
2381
2382pub(crate) fn validate_embedding_batch(
2384 values: &[Vec<f32>],
2385 requested: usize,
2386 expected_dim: usize,
2387) -> Result<(), MemoryError> {
2388 if values.len() != requested {
2389 return Err(MemoryError::EmbeddingBatchCountMismatch {
2390 requested,
2391 returned: values.len(),
2392 });
2393 }
2394 for embedding in values {
2395 validate_embedding(embedding, expected_dim)?;
2396 }
2397 Ok(())
2398}
2399
2400pub(crate) fn validate_vector_blob_len(
2402 bytes: &[u8],
2403 expected_dim: usize,
2404) -> Result<(), MemoryError> {
2405 let expected_bytes = expected_dim
2406 .checked_mul(4)
2407 .ok_or_else(|| MemoryError::InvalidConfig {
2408 field: "embedding.dimensions",
2409 reason: "dimension byte length overflow".to_string(),
2410 })?;
2411 if bytes.len() != expected_bytes {
2412 return Err(MemoryError::VectorBlobLengthMismatch {
2413 expected_bytes,
2414 actual_bytes: bytes.len(),
2415 });
2416 }
2417 Ok(())
2418}
2419
2420#[allow(clippy::manual_is_multiple_of)]
2422pub fn decode_f32_le(bytes: &[u8], expected_dim: usize) -> Result<Vec<f32>, MemoryError> {
2423 validate_vector_blob_len(bytes, expected_dim)?;
2424 decode_f32_le_unchecked_dim(bytes)
2425}
2426
2427#[allow(clippy::manual_is_multiple_of)]
2429pub fn bytes_to_embedding(bytes: &[u8]) -> Result<Vec<f32>, MemoryError> {
2430 if bytes.len() % 4 != 0 {
2431 return Err(MemoryError::InvalidEmbedding {
2432 expected_bytes: bytes.len() - (bytes.len() % 4),
2433 actual_bytes: bytes.len(),
2434 });
2435 }
2436
2437 decode_f32_le_unchecked_dim(bytes)
2438}
2439
2440fn decode_f32_le_unchecked_dim(bytes: &[u8]) -> Result<Vec<f32>, MemoryError> {
2441 let mut embedding = Vec::with_capacity(bytes.len() / 4);
2442 for (index, chunk) in bytes.chunks_exact(4).enumerate() {
2443 let value = f32::from_le_bytes([chunk[0], chunk[1], chunk[2], chunk[3]]);
2444 if !value.is_finite() {
2445 return Err(MemoryError::NonFiniteEmbeddingValue { index });
2446 }
2447 embedding.push(value);
2448 }
2449 Ok(embedding)
2450}
2451
2452pub fn is_embeddings_dirty(conn: &Connection) -> Result<bool, MemoryError> {
2453 let dirty: i32 = conn
2454 .query_row(
2455 "SELECT COALESCE(embeddings_dirty, 0) FROM embedding_metadata WHERE id = 1",
2456 [],
2457 |row| row.get(0),
2458 )
2459 .unwrap_or(0);
2460 Ok(dirty != 0)
2461}
2462
2463pub fn clear_embeddings_dirty(conn: &Connection) -> Result<(), MemoryError> {
2464 conn.execute(
2465 "UPDATE embedding_metadata SET embeddings_dirty = 0 WHERE id = 1",
2466 [],
2467 )?;
2468 Ok(())
2469}
2470
2471#[cfg(feature = "hnsw")]
2472pub(crate) fn queue_pending_index_op(
2473 tx: &rusqlite::Transaction<'_>,
2474 item_key: &str,
2475 entity_type: &str,
2476 op_kind: IndexOpKind,
2477) -> Result<(), MemoryError> {
2478 tx.execute(
2479 "INSERT INTO pending_index_ops (item_key, entity_type, op_kind, attempt_count, last_error, updated_at)
2480 VALUES (?1, ?2, ?3, 0, NULL, datetime('now'))
2481 ON CONFLICT(item_key) DO UPDATE SET
2482 entity_type = excluded.entity_type,
2483 op_kind = excluded.op_kind,
2484 attempt_count = 0,
2485 last_error = NULL,
2486 updated_at = datetime('now')",
2487 params![item_key, entity_type, op_kind.as_str()],
2488 )?;
2489 mark_sidecar_dirty(tx)?;
2490 Ok(())
2491}
2492
2493#[cfg(feature = "hnsw")]
2494pub(crate) use IndexOpKind as PendingIndexOpKind;
2495
2496#[cfg(feature = "hnsw")]
2497pub(crate) fn enqueue_pending_index_op(
2498 tx: &rusqlite::Transaction<'_>,
2499 item_key: &str,
2500 entity_type: &str,
2501 op_kind: PendingIndexOpKind,
2502) -> Result<(), MemoryError> {
2503 queue_pending_index_op(tx, item_key, entity_type, op_kind)
2504}
2505
2506pub(crate) fn list_pending_index_ops(
2507 conn: &Connection,
2508) -> Result<Vec<PendingIndexOp>, MemoryError> {
2509 let table_exists: bool = conn
2510 .query_row(
2511 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type='table' AND name='pending_index_ops'",
2512 [],
2513 |row| row.get(0),
2514 )
2515 .unwrap_or(false);
2516 if !table_exists {
2517 return Ok(Vec::new());
2518 }
2519
2520 let mut stmt = conn.prepare(
2521 "SELECT item_key, entity_type, op_kind, attempt_count, last_error
2522 FROM pending_index_ops
2523 ORDER BY updated_at ASC, item_key ASC",
2524 )?;
2525 let rows = stmt
2526 .query_map([], |row| {
2527 let item_key: String = row.get(0)?;
2528 let op_kind: String = row.get(2)?;
2529 Ok(PendingIndexOp {
2530 item_key: item_key.clone(),
2531 entity_type: row.get(1)?,
2532 op_kind: IndexOpKind::parse(&op_kind, &item_key).map_err(|e| {
2533 rusqlite::Error::FromSqlConversionFailure(
2534 2,
2535 rusqlite::types::Type::Text,
2536 Box::new(e),
2537 )
2538 })?,
2539 attempt_count: row.get::<_, i64>(3)? as u32,
2540 last_error: row.get(4)?,
2541 })
2542 })?
2543 .collect::<Result<Vec<_>, _>>()?;
2544 Ok(rows)
2545}
2546
2547#[cfg(feature = "hnsw")]
2548pub(crate) fn pending_index_op_count(conn: &Connection) -> Result<usize, MemoryError> {
2549 let table_exists: bool = conn
2550 .query_row(
2551 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type='table' AND name='pending_index_ops'",
2552 [],
2553 |row| row.get(0),
2554 )
2555 .unwrap_or(false);
2556 if !table_exists {
2557 return Ok(0);
2558 }
2559
2560 let count: i64 = conn.query_row("SELECT COUNT(*) FROM pending_index_ops", [], |row| {
2561 row.get(0)
2562 })?;
2563 Ok(count as usize)
2564}
2565
2566#[cfg(feature = "hnsw")]
2567pub(crate) fn mark_pending_index_ops_failed(
2568 conn: &Connection,
2569 item_keys: &[String],
2570 error: &str,
2571) -> Result<(), MemoryError> {
2572 with_transaction(conn, |tx| {
2573 for item_key in item_keys {
2574 tx.execute(
2575 "UPDATE pending_index_ops
2576 SET attempt_count = attempt_count + 1,
2577 last_error = ?1,
2578 updated_at = datetime('now')
2579 WHERE item_key = ?2",
2580 params![error, item_key],
2581 )?;
2582 }
2583 Ok(())
2584 })
2585}
2586
2587#[cfg(feature = "hnsw")]
2588pub(crate) fn clear_pending_index_ops(
2589 conn: &Connection,
2590 item_keys: &[String],
2591) -> Result<(), MemoryError> {
2592 with_transaction(conn, |tx| {
2593 for item_key in item_keys {
2594 tx.execute(
2595 "DELETE FROM pending_index_ops WHERE item_key = ?1",
2596 params![item_key],
2597 )?;
2598 }
2599 Ok(())
2600 })
2601}
2602
2603#[cfg(feature = "hnsw")]
2604pub(crate) fn clear_all_pending_index_ops(conn: &Connection) -> Result<(), MemoryError> {
2605 conn.execute("DELETE FROM pending_index_ops", [])?;
2606 Ok(())
2607}
2608
2609#[cfg(feature = "hnsw")]
2610pub(crate) fn load_embedding_for_index_key(
2611 conn: &Connection,
2612 item_key: &str,
2613) -> Result<Option<Vec<f32>>, MemoryError> {
2614 let Some((domain, raw_id)) = item_key.split_once(':') else {
2615 return Err(MemoryError::InvalidKey(item_key.to_string()));
2616 };
2617
2618 let blob_result: Result<Option<Vec<u8>>, rusqlite::Error> = match domain {
2619 "fact" => conn.query_row(
2620 "SELECT embedding FROM facts WHERE id = ?1",
2621 params![raw_id],
2622 |row| row.get(0),
2623 ),
2624 "chunk" => conn.query_row(
2625 "SELECT embedding FROM chunks WHERE id = ?1",
2626 params![raw_id],
2627 |row| row.get(0),
2628 ),
2629 "msg" => {
2630 let message_id = raw_id
2631 .parse::<i64>()
2632 .map_err(|e| MemoryError::InvalidKey(format!("{}: {e}", item_key)))?;
2633 conn.query_row(
2634 "SELECT embedding FROM messages WHERE id = ?1",
2635 params![message_id],
2636 |row| row.get(0),
2637 )
2638 }
2639 "episode" => conn.query_row(
2640 "SELECT embedding FROM episodes WHERE episode_id = ?1",
2641 params![raw_id],
2642 |row| row.get(0),
2643 ),
2644 _ => return Err(MemoryError::InvalidKey(item_key.to_string())),
2645 };
2646
2647 let blob = match blob_result {
2648 Ok(blob) => blob,
2649 Err(rusqlite::Error::QueryReturnedNoRows) => None,
2650 Err(err) => return Err(err.into()),
2651 };
2652
2653 blob.map(|bytes| bytes_to_embedding(&bytes)).transpose()
2654}
2655
2656#[cfg(feature = "hnsw")]
2657fn mark_sidecar_dirty(tx: &rusqlite::Transaction<'_>) -> Result<(), MemoryError> {
2658 tx.execute(
2659 "INSERT INTO hnsw_metadata (key, value) VALUES ('sidecar_dirty', '1')
2660 ON CONFLICT(key) DO UPDATE SET value = '1'",
2661 [],
2662 )?;
2663 Ok(())
2664}
2665
2666#[cfg(feature = "hnsw")]
2667pub(crate) fn is_sidecar_dirty(conn: &Connection) -> Result<bool, MemoryError> {
2668 let dirty: Option<String> = conn
2670 .query_row(
2671 "SELECT value FROM hnsw_metadata WHERE key = 'sidecar_dirty'",
2672 [],
2673 |row| row.get(0),
2674 )
2675 .ok();
2676 Ok(matches!(dirty.as_deref(), Some("1")))
2677}
2678
2679#[cfg(feature = "hnsw")]
2680pub(crate) fn set_sidecar_dirty(conn: &Connection, dirty: bool) -> Result<(), MemoryError> {
2681 conn.execute(
2682 "INSERT INTO hnsw_metadata (key, value) VALUES ('sidecar_dirty', ?1)
2683 ON CONFLICT(key) DO UPDATE SET value = excluded.value",
2684 params![if dirty { "1" } else { "0" }],
2685 )?;
2686 Ok(())
2687}
2688
2689pub(crate) fn parse_optional_json(
2690 table: &'static str,
2691 row_id: &str,
2692 field: &'static str,
2693 raw: Option<&str>,
2694) -> Result<Option<serde_json::Value>, MemoryError> {
2695 match raw {
2696 Some(raw) => serde_json::from_str(raw)
2697 .map(Some)
2698 .map_err(|e| MemoryError::CorruptData {
2699 table,
2700 row_id: row_id.to_string(),
2701 detail: format!("invalid {field}: {e}"),
2702 }),
2703 None => Ok(None),
2704 }
2705}
2706
2707pub(crate) fn parse_string_list_json(
2708 table: &'static str,
2709 row_id: &str,
2710 field: &'static str,
2711 raw: &str,
2712) -> Result<Vec<String>, MemoryError> {
2713 serde_json::from_str(raw).map_err(|e| MemoryError::CorruptData {
2714 table,
2715 row_id: row_id.to_string(),
2716 detail: format!("invalid {field}: {e}"),
2717 })
2718}
2719
2720pub(crate) fn parse_role(
2721 table: &'static str,
2722 row_id: &str,
2723 raw: &str,
2724) -> Result<Role, MemoryError> {
2725 Role::from_str_value(raw).ok_or_else(|| MemoryError::CorruptData {
2726 table,
2727 row_id: row_id.to_string(),
2728 detail: format!("invalid role '{raw}'"),
2729 })
2730}
2731
2732pub(crate) fn parse_episode_outcome(
2733 row_id: &str,
2734 raw: &str,
2735) -> Result<EpisodeOutcome, MemoryError> {
2736 EpisodeOutcome::from_str_value(raw).ok_or_else(|| MemoryError::CorruptData {
2737 table: "episodes",
2738 row_id: row_id.to_string(),
2739 detail: format!("invalid outcome '{raw}'"),
2740 })
2741}
2742
2743pub(crate) fn parse_verification_status(
2744 row_id: &str,
2745 raw: &str,
2746) -> Result<VerificationStatus, MemoryError> {
2747 serde_json::from_str(raw).map_err(|e| MemoryError::CorruptData {
2748 table: "episodes",
2749 row_id: row_id.to_string(),
2750 detail: format!("invalid verification_status: {e}"),
2751 })
2752}
2753
2754pub fn verify_integrity_sync(
2756 conn: &Connection,
2757 mode: VerifyMode,
2758) -> Result<IntegrityReport, MemoryError> {
2759 let mut issues = Vec::new();
2760
2761 let schema_version: u32 = conn
2762 .query_row("PRAGMA user_version", [], |row| row.get(0))
2763 .unwrap_or_else(|e| {
2764 issues.push(format!("failed to read schema version: {e}"));
2765 0
2766 });
2767 if schema_version > MAX_SCHEMA_VERSION {
2768 issues.push(format!(
2769 "schema version {} is ahead of supported {}",
2770 schema_version, MAX_SCHEMA_VERSION
2771 ));
2772 }
2773
2774 let fact_count: usize = conn
2775 .query_row("SELECT COUNT(*) FROM facts", [], |row| row.get(0))
2776 .unwrap_or_else(|e| {
2777 issues.push(format!("failed to count facts: {e}"));
2778 0
2779 });
2780 let chunk_count: usize = conn
2781 .query_row("SELECT COUNT(*) FROM chunks", [], |row| row.get(0))
2782 .unwrap_or_else(|e| {
2783 issues.push(format!("failed to count chunks: {e}"));
2784 0
2785 });
2786 let message_count: usize = conn
2787 .query_row("SELECT COUNT(*) FROM messages", [], |row| row.get(0))
2788 .unwrap_or_else(|e| {
2789 issues.push(format!("failed to count messages: {e}"));
2790 0
2791 });
2792 let episode_count: usize = conn
2793 .query_row("SELECT COUNT(*) FROM episodes", [], |row| row.get(0))
2794 .unwrap_or_else(|e| {
2795 issues.push(format!("failed to count episodes: {e}"));
2796 0
2797 });
2798
2799 let facts_missing_embeddings: usize = conn
2800 .query_row(
2801 "SELECT COUNT(*) FROM facts WHERE embedding IS NULL",
2802 [],
2803 |row| row.get(0),
2804 )
2805 .unwrap_or_else(|e| {
2806 issues.push(format!("failed to count facts missing embeddings: {e}"));
2807 0
2808 });
2809 let chunks_missing_embeddings: usize = conn
2810 .query_row(
2811 "SELECT COUNT(*) FROM chunks WHERE embedding IS NULL",
2812 [],
2813 |row| row.get(0),
2814 )
2815 .unwrap_or_else(|e| {
2816 issues.push(format!("failed to count chunks missing embeddings: {e}"));
2817 0
2818 });
2819 let episodes_missing_embeddings: usize = conn
2820 .query_row(
2821 "SELECT COUNT(*) FROM episodes WHERE embedding IS NULL",
2822 [],
2823 |row| row.get(0),
2824 )
2825 .unwrap_or_else(|e| {
2826 issues.push(format!("failed to count episodes missing embeddings: {e}"));
2827 0
2828 });
2829
2830 if facts_missing_embeddings > 0 {
2831 issues.push(format!(
2832 "{} facts missing embeddings",
2833 facts_missing_embeddings
2834 ));
2835 }
2836 if chunks_missing_embeddings > 0 {
2837 issues.push(format!(
2838 "{} chunks missing embeddings",
2839 chunks_missing_embeddings
2840 ));
2841 }
2842 if episodes_missing_embeddings > 0 {
2843 issues.push(format!(
2844 "{} episodes missing embeddings",
2845 episodes_missing_embeddings
2846 ));
2847 }
2848
2849 let pending_ops = list_pending_index_ops(conn).unwrap_or_default();
2850 if !pending_ops.is_empty() {
2851 issues.push(format!(
2852 "{} pending HNSW sidecar ops queued in SQLite",
2853 pending_ops.len()
2854 ));
2855 for op in pending_ops.iter().take(5) {
2856 let op_kind = op.op_kind.as_str();
2857 let detail = match &op.last_error {
2858 Some(last_error) => format!(
2859 "{} {} {} (attempts: {}, last_error: {})",
2860 op.entity_type,
2861 op.op_kind.as_str(),
2862 op.item_key,
2863 op.attempt_count,
2864 last_error
2865 ),
2866 None => format!(
2867 "{} {} {} (attempts: {})",
2868 op.entity_type, op_kind, op.item_key, op.attempt_count
2869 ),
2870 };
2871 issues.push(format!("pending sidecar op: {detail}"));
2872 }
2873 }
2874
2875 if mode == VerifyMode::Full {
2876 let dims: usize = conn
2877 .query_row(
2878 "SELECT dimensions FROM embedding_metadata WHERE id = 1",
2879 [],
2880 |row| row.get(0),
2881 )
2882 .unwrap_or_else(|e| {
2883 issues.push(format!("failed to read embedding dimensions: {e}"));
2884 0
2885 });
2886
2887 verify_fts_drift(conn, "facts", "facts_rowid_map", fact_count, &mut issues);
2888 verify_fts_drift(conn, "chunks", "chunks_rowid_map", chunk_count, &mut issues);
2889 verify_fts_drift(
2890 conn,
2891 "messages",
2892 "messages_rowid_map",
2893 message_count,
2894 &mut issues,
2895 );
2896 verify_fts_drift(
2897 conn,
2898 "episodes",
2899 "episodes_rowid_map",
2900 episode_count,
2901 &mut issues,
2902 );
2903
2904 verify_blob_table(conn, "facts", "id", "embedding", dims, &mut issues)?;
2905 verify_blob_table(conn, "chunks", "id", "embedding", dims, &mut issues)?;
2906 verify_blob_table(conn, "messages", "id", "embedding", dims, &mut issues)?;
2907 verify_blob_table(
2908 conn,
2909 "episodes",
2910 "episode_id",
2911 "embedding",
2912 dims,
2913 &mut issues,
2914 )?;
2915
2916 verify_quantized_table(conn, "facts", "id", dims, &mut issues)?;
2917 verify_quantized_table(conn, "chunks", "id", dims, &mut issues)?;
2918 verify_quantized_table(conn, "messages", "id", dims, &mut issues)?;
2919 verify_quantized_table(conn, "episodes", "episode_id", dims, &mut issues)?;
2920
2921 verify_session_rows(conn, &mut issues)?;
2922 verify_message_rows(conn, &mut issues)?;
2923 verify_fact_rows(conn, &mut issues)?;
2924 verify_document_rows(conn, &mut issues)?;
2925 verify_episode_rows(conn, &mut issues)?;
2926
2927 let integrity_check: String = conn
2928 .query_row("PRAGMA integrity_check", [], |row| row.get(0))
2929 .unwrap_or_else(|_| "error".to_string());
2930 if integrity_check != "ok" {
2931 issues.push(format!("SQLite integrity_check: {}", integrity_check));
2932 }
2933 }
2934
2935 Ok(IntegrityReport {
2936 ok: issues.is_empty(),
2937 schema_version,
2938 fact_count,
2939 chunk_count,
2940 message_count,
2941 facts_missing_embeddings,
2942 chunks_missing_embeddings,
2943 issues,
2944 })
2945}
2946
2947pub fn reconcile_fts(conn: &Connection) -> Result<(), MemoryError> {
2949 with_transaction(conn, |tx| {
2950 tx.execute_batch("DROP TABLE IF EXISTS facts_fts")?;
2951 tx.execute_batch("DELETE FROM facts_rowid_map")?;
2952 tx.execute_batch(
2953 "CREATE VIRTUAL TABLE facts_fts USING fts5(
2954 content,
2955 content='',
2956 content_rowid='rowid',
2957 tokenize='porter unicode61'
2958 )",
2959 )?;
2960 tx.execute_batch("INSERT INTO facts_rowid_map (fact_id) SELECT id FROM facts")?;
2961 tx.execute_batch(
2962 "INSERT INTO facts_fts (rowid, content)
2963 SELECT rm.rowid, f.content
2964 FROM facts_rowid_map rm
2965 JOIN facts f ON f.id = rm.fact_id",
2966 )?;
2967
2968 tx.execute_batch("DROP TABLE IF EXISTS chunks_fts")?;
2969 tx.execute_batch("DELETE FROM chunks_rowid_map")?;
2970 tx.execute_batch(
2971 "CREATE VIRTUAL TABLE chunks_fts USING fts5(
2972 content,
2973 content='',
2974 content_rowid='rowid',
2975 tokenize='porter unicode61'
2976 )",
2977 )?;
2978 tx.execute_batch("INSERT INTO chunks_rowid_map (chunk_id) SELECT id FROM chunks")?;
2979 tx.execute_batch(
2980 "INSERT INTO chunks_fts (rowid, content)
2981 SELECT rm.rowid, c.content
2982 FROM chunks_rowid_map rm
2983 JOIN chunks c ON c.id = rm.chunk_id",
2984 )?;
2985
2986 tx.execute_batch("DROP TABLE IF EXISTS messages_fts")?;
2987 tx.execute_batch("DELETE FROM messages_rowid_map")?;
2988 tx.execute_batch(
2989 "CREATE VIRTUAL TABLE messages_fts USING fts5(
2990 content,
2991 content='',
2992 content_rowid='rowid',
2993 tokenize='porter unicode61'
2994 )",
2995 )?;
2996 tx.execute_batch("INSERT INTO messages_rowid_map (message_id) SELECT id FROM messages")?;
2997 tx.execute_batch(
2998 "INSERT INTO messages_fts (rowid, content)
2999 SELECT rm.rowid, m.content
3000 FROM messages_rowid_map rm
3001 JOIN messages m ON m.id = rm.message_id",
3002 )?;
3003
3004 tx.execute_batch("DROP TABLE IF EXISTS episodes_fts")?;
3005 tx.execute_batch("DELETE FROM episodes_rowid_map")?;
3006 tx.execute_batch(
3007 "CREATE VIRTUAL TABLE episodes_fts USING fts5(
3008 content,
3009 content='',
3010 content_rowid='rowid',
3011 tokenize='porter unicode61'
3012 )",
3013 )?;
3014 tx.execute_batch(
3015 "INSERT INTO episodes_rowid_map (episode_id, document_id) SELECT episode_id, document_id FROM episodes",
3016 )?;
3017 tx.execute_batch(
3018 "INSERT INTO episodes_fts (rowid, content)
3019 SELECT rm.rowid, e.search_text
3020 FROM episodes_rowid_map rm
3021 JOIN episodes e ON e.episode_id = rm.episode_id",
3022 )?;
3023
3024 Ok(())
3025 })?;
3026
3027 tracing::info!("FTS indexes reconciled");
3028 Ok(())
3029}
3030
3031fn verify_fts_drift(
3032 conn: &Connection,
3033 label: &str,
3034 map_table: &str,
3035 source_count: usize,
3036 issues: &mut Vec<String>,
3037) {
3038 let table_exists: bool = conn
3039 .query_row(
3040 "SELECT COUNT(*) > 0 FROM sqlite_master WHERE type='table' AND name = ?1",
3041 params![map_table],
3042 |row| row.get(0),
3043 )
3044 .unwrap_or(false);
3045 if !table_exists {
3046 if source_count > 0 {
3047 issues.push(format!("{} rows exist but {} is missing", label, map_table));
3048 }
3049 return;
3050 }
3051
3052 let sql = format!("SELECT COUNT(*) FROM {}", map_table);
3053 let indexed_count: usize = conn.query_row(&sql, [], |row| row.get(0)).unwrap_or(0);
3054 if indexed_count != source_count {
3055 issues.push(format!(
3056 "FTS {} index drift: {} rows in map vs {} source rows",
3057 label, indexed_count, source_count
3058 ));
3059 }
3060}
3061
3062fn verify_blob_table(
3063 conn: &Connection,
3064 table: &'static str,
3065 id_column: &'static str,
3066 blob_column: &'static str,
3067 expected_dims: usize,
3068 issues: &mut Vec<String>,
3069) -> Result<(), MemoryError> {
3070 if expected_dims == 0 {
3071 return Ok(());
3072 }
3073
3074 let sql = format!(
3075 "SELECT CAST({id_column} AS TEXT), {blob_column} FROM {table} WHERE {blob_column} IS NOT NULL"
3076 );
3077 let mut stmt = conn.prepare(&sql)?;
3078 let rows = stmt.query_map([], |row| {
3079 Ok((row.get::<_, String>(0)?, row.get::<_, Vec<u8>>(1)?))
3080 })?;
3081
3082 for row in rows {
3083 let (row_id, blob) = row?;
3084 match bytes_to_embedding(&blob) {
3085 Ok(embedding) if embedding.len() != expected_dims => issues.push(format!(
3086 "{}({}) has embedding dimension {} but expected {}",
3087 table,
3088 row_id,
3089 embedding.len(),
3090 expected_dims
3091 )),
3092 Ok(_) => {}
3093 Err(err) => issues.push(format!(
3094 "{}({}) invalid embedding blob: {}",
3095 table, row_id, err
3096 )),
3097 }
3098 }
3099
3100 Ok(())
3101}
3102
3103fn verify_quantized_table(
3104 conn: &Connection,
3105 table: &'static str,
3106 id_column: &'static str,
3107 expected_dims: usize,
3108 issues: &mut Vec<String>,
3109) -> Result<(), MemoryError> {
3110 if expected_dims == 0 {
3111 return Ok(());
3112 }
3113
3114 let sql = format!(
3115 "SELECT CAST({id_column} AS TEXT), embedding_q8 FROM {table} WHERE embedding IS NOT NULL"
3116 );
3117 let mut stmt = conn.prepare(&sql)?;
3118 let rows = stmt.query_map([], |row| {
3119 Ok((row.get::<_, String>(0)?, row.get::<_, Option<Vec<u8>>>(1)?))
3120 })?;
3121
3122 for row in rows {
3123 let (row_id, blob) = row?;
3124 match blob {
3125 Some(blob) => {
3126 if let Err(err) = unpack_quantized(&blob, expected_dims) {
3127 issues.push(format!(
3128 "{}({}) invalid quantized embedding: {}",
3129 table, row_id, err
3130 ));
3131 }
3132 }
3133 None => issues.push(format!("{}({}) missing quantized embedding", table, row_id)),
3134 }
3135 }
3136
3137 Ok(())
3138}
3139
3140fn verify_session_rows(conn: &Connection, issues: &mut Vec<String>) -> Result<(), MemoryError> {
3141 let mut stmt = conn.prepare("SELECT id, metadata FROM sessions WHERE metadata IS NOT NULL")?;
3142 let rows = stmt.query_map([], |row| {
3143 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3144 })?;
3145 for row in rows {
3146 let (id, metadata) = row?;
3147 if let Err(err) = parse_optional_json("sessions", &id, "metadata", Some(&metadata)) {
3148 issues.push(err.to_string());
3149 }
3150 }
3151 Ok(())
3152}
3153
3154fn verify_message_rows(conn: &Connection, issues: &mut Vec<String>) -> Result<(), MemoryError> {
3155 let mut stmt = conn.prepare("SELECT id, role, metadata FROM messages")?;
3156 let rows = stmt.query_map([], |row| {
3157 Ok((
3158 row.get::<_, i64>(0)?,
3159 row.get::<_, String>(1)?,
3160 row.get::<_, Option<String>>(2)?,
3161 ))
3162 })?;
3163 for row in rows {
3164 let (id, role, metadata) = row?;
3165 let row_id = id.to_string();
3166 if let Err(err) = parse_role("messages", &row_id, &role) {
3167 issues.push(err.to_string());
3168 }
3169 if let Err(err) = parse_optional_json("messages", &row_id, "metadata", metadata.as_deref())
3170 {
3171 issues.push(err.to_string());
3172 }
3173 }
3174 Ok(())
3175}
3176
3177fn verify_fact_rows(conn: &Connection, issues: &mut Vec<String>) -> Result<(), MemoryError> {
3178 let mut stmt = conn.prepare("SELECT id, metadata FROM facts WHERE metadata IS NOT NULL")?;
3179 let rows = stmt.query_map([], |row| {
3180 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3181 })?;
3182 for row in rows {
3183 let (id, metadata) = row?;
3184 if let Err(err) = parse_optional_json("facts", &id, "metadata", Some(&metadata)) {
3185 issues.push(err.to_string());
3186 }
3187 }
3188 Ok(())
3189}
3190
3191fn verify_document_rows(conn: &Connection, issues: &mut Vec<String>) -> Result<(), MemoryError> {
3192 let mut stmt = conn.prepare("SELECT id, metadata FROM documents WHERE metadata IS NOT NULL")?;
3193 let rows = stmt.query_map([], |row| {
3194 Ok((row.get::<_, String>(0)?, row.get::<_, String>(1)?))
3195 })?;
3196 for row in rows {
3197 let (id, metadata) = row?;
3198 if let Err(err) = parse_optional_json("documents", &id, "metadata", Some(&metadata)) {
3199 issues.push(err.to_string());
3200 }
3201 }
3202 Ok(())
3203}
3204
3205fn verify_episode_rows(conn: &Connection, issues: &mut Vec<String>) -> Result<(), MemoryError> {
3206 let mut stmt = conn.prepare(
3207 "SELECT episode_id, cause_ids, outcome, verification_status
3208 FROM episodes",
3209 )?;
3210 let rows = stmt.query_map([], |row| {
3211 Ok((
3212 row.get::<_, String>(0)?,
3213 row.get::<_, String>(1)?,
3214 row.get::<_, String>(2)?,
3215 row.get::<_, String>(3)?,
3216 ))
3217 })?;
3218 for row in rows {
3219 let (episode_id, cause_ids, outcome, verification_status) = row?;
3220 if let Err(err) = parse_string_list_json("episodes", &episode_id, "cause_ids", &cause_ids) {
3221 issues.push(err.to_string());
3222 }
3223 if let Err(err) = parse_episode_outcome(&episode_id, &outcome) {
3224 issues.push(err.to_string());
3225 }
3226 if let Err(err) = parse_verification_status(&episode_id, &verification_status) {
3227 issues.push(err.to_string());
3228 }
3229 }
3230 Ok(())
3231}