1use std::path::Path;
9use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
10use std::sync::Arc;
11
12use rusqlite::OptionalExtension;
13
14use crate::database_owner_identity::{DatabaseOwnerIdentity, DatabaseOwnerIdentityError};
15use crate::error::SqliteError;
16use crate::pool::{ConnectionPool, PoolConfig, WalCeilingPolicy};
17use crate::sql_bridge::SqlBridge;
18use crate::stores::{agents, attachment, blob, entity, event, graph, note, sparse, text, vectors};
19
20mod code_map;
21#[path = "backend/schema_readiness.rs"]
22mod memory_visibility;
23mod pack_schema;
24mod policy_open;
25
26#[cfg(test)]
27#[path = "backend/memory_visibility_tests.rs"]
28mod memory_visibility_tests;
29
30#[cfg(any(unix, windows))]
31mod claimed_file_identity;
32
33fn sqlite_table_exists(conn: &rusqlite::Connection, table: &str) -> Result<bool, SqliteError> {
34 conn.query_row(
35 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
36 rusqlite::params![table],
37 |row| row.get::<_, i64>(0),
38 )
39 .optional()
40 .map(|row| row.is_some())
41 .map_err(SqliteError::Rusqlite)
42}
43
44fn ensure_fts_rowid_map_backfilled(
100 conn: &rusqlite::Connection,
101 table: &str,
102 admission: &crate::pool::WriteAdmission,
103) -> Result<(), SqliteError> {
104 let map = text::rowid_map_table(table);
105 let state = text::rowid_map_state_table(table);
106 let already_backfilled: bool = conn.query_row(
107 &format!("SELECT EXISTS(SELECT 1 FROM {state} WHERE key = 'backfill' AND value = ?1)"),
108 rusqlite::params![text::ROWID_MAP_BACKFILL_COMPLETE],
109 |row| row.get(0),
110 )?;
111 if already_backfilled {
112 return Ok(());
113 }
114 let fts_has_a_row: bool = conn.query_row(
115 &format!("SELECT EXISTS(SELECT 1 FROM {table} LIMIT 1)"),
116 [],
117 |row| row.get(0),
118 )?;
119 if !fts_has_a_row {
120 return Ok(());
121 }
122
123 conn.execute_batch("BEGIN IMMEDIATE")?;
124 if let Err(error) = admission.check() {
125 return Err(crate::migrations::capacity_refusal_after_rollback(
126 conn,
127 conn.execute_batch("ROLLBACK"),
128 error,
129 "FTS rowid-map backfill",
130 ));
131 }
132 let result: Result<(), SqliteError> = (|| {
133 conn.execute_batch(&format!(
134 "DELETE FROM {map} WHERE NOT EXISTS ( \
135 SELECT 1 FROM {table} \
136 WHERE {table}.rowid = {map}.rowid \
137 AND {table}.namespace = {map}.namespace \
138 AND {table}.subject_id = {map}.subject_id \
139 )"
140 ))?;
141 conn.execute_batch(&format!(
142 "INSERT OR REPLACE INTO {map} (namespace, subject_id, rowid) \
143 SELECT namespace, subject_id, rowid FROM {table} \
144 WHERE namespace IS NOT NULL AND subject_id IS NOT NULL \
145 ORDER BY updated_at ASC, rowid ASC"
146 ))?;
147 conn.execute_batch(&format!(
148 "DELETE FROM {table} \
149 WHERE namespace IS NOT NULL AND subject_id IS NOT NULL \
150 AND rowid NOT IN (SELECT rowid FROM {map})"
151 ))?;
152 conn.execute(
153 &format!("INSERT OR REPLACE INTO {state} (key, value) VALUES ('backfill', ?1)"),
154 rusqlite::params![text::ROWID_MAP_BACKFILL_COMPLETE],
155 )?;
156 Ok(())
157 })();
158
159 match result {
160 Ok(()) => {
161 conn.execute_batch("COMMIT")?;
162 Ok(())
163 }
164 Err(e) => {
165 let _ = conn.execute_batch("ROLLBACK");
166 Err(e)
167 }
168 }
169}
170
171fn warn_scan_fallback_once(table: &str) {
177 static WARNED: std::sync::Once = std::sync::Once::new();
178 WARNED.call_once(|| {
179 tracing::warn!(
180 table,
181 "opened a read-only text-search table with no rowid-map sidecar, or with a sidecar \
182 that has never proven a completed backfill (no durable completion marker); a \
183 read-only connection cannot create, backfill, or reconcile the map itself, so this \
184 falls back to pre-map namespace/subject_id scan predicates for get/delete on this \
185 table rather than trusting a map that might be partial"
186 );
187 });
188}
189
190fn validate_vector_model_key(model_key: &str) -> Result<(), SqliteError> {
191 if model_key.is_empty()
192 || !model_key
193 .chars()
194 .all(|c| c.is_ascii_alphanumeric() || c == '_')
195 {
196 return Err(SqliteError::InvalidData(format!(
197 "invalid model_key '{}': must be non-empty and contain only \
198 alphanumeric/underscore characters",
199 model_key
200 )));
201 }
202 Ok(())
203}
204
205fn validate_vector_table_columns(
206 conn: &rusqlite::Connection,
207 table: &str,
208) -> Result<(), SqliteError> {
209 let pragma = format!("PRAGMA table_xinfo({table})");
210 let mut stmt = conn.prepare(&pragma)?;
211 let mut rows = stmt.query([])?;
212 let mut has_field = false;
213 let mut has_embedding_model = false;
214 while let Some(row) = rows.next()? {
215 let name: String = row.get(1)?;
216 if name == "field" {
217 has_field = true;
218 }
219 if name == "embedding_model" {
220 has_embedding_model = true;
221 }
222 }
223 if !has_field || !has_embedding_model {
224 return Err(SqliteError::InvalidData(format!(
225 "vec0 table '{table}' is missing required column(s) (field={has_field}, \
226 embedding_model={has_embedding_model}); this is a pre-v0.2.8 vector schema and is \
227 not supported — recreate the database"
228 )));
229 }
230 Ok(())
231}
232
233#[derive(Clone, Copy, Debug)]
234pub(crate) enum StoreSchemaKind {
235 Entities,
236 Graph,
237 Notes,
238 Events,
239 Agents,
240}
241
242impl StoreSchemaKind {
243 fn initializer(self) -> fn(&rusqlite::Connection) -> Result<(), rusqlite::Error> {
244 match self {
245 Self::Entities => entity::ensure_entities_schema,
246 Self::Graph => graph::ensure_graph_schema,
247 Self::Notes => note::ensure_notes_schema,
248 Self::Events => event::ensure_events_schema,
249 Self::Agents => agents::ensure_agents_schema,
250 }
251 }
252}
253
254#[derive(Default)]
255pub(crate) struct StoreSchemaGate {
256 pub(crate) ready: AtomicBool,
257 #[cfg(test)]
258 pub(crate) attempts: AtomicUsize,
259}
260
261impl StoreSchemaGate {
262 pub(crate) fn ensure(
263 &self,
264 conn: &rusqlite::Connection,
265 ensure: fn(&rusqlite::Connection) -> Result<(), rusqlite::Error>,
266 ) -> Result<(), rusqlite::Error> {
267 if self.ready.load(Ordering::Acquire) {
270 return Ok(());
271 }
272 #[cfg(test)]
273 self.attempts.fetch_add(1, Ordering::Relaxed);
274 ensure(conn)?;
275 self.ready.store(true, Ordering::Release);
276 Ok(())
277 }
278}
279
280pub struct StorageBackend {
286 pool: Arc<ConnectionPool>,
287 is_file_backed: bool,
288 path: Option<std::path::PathBuf>,
289 vector_tables_ready: parking_lot::Mutex<std::collections::HashSet<String>>,
297 notes_seq_repair_runs: AtomicUsize,
304 store_schemas: [Arc<StoreSchemaGate>; 5],
305}
306
307impl StorageBackend {
308 pub fn sqlite(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
322 Self::sqlite_with_pool_config(path, PoolConfig::default(), None)
323 }
324
325 #[cfg(any(test, feature = "test-support"))]
327 pub fn sqlite_for_test(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
328 Self::sqlite_with_pool_config(path, PoolConfig::for_test(), None)
329 }
330
331 #[cfg(any(test, feature = "test-support"))]
333 pub fn sqlite_for_test_with_journal_mode(
334 path: impl AsRef<Path>,
335 wal_mode: bool,
336 busy_timeout: std::time::Duration,
337 ) -> Result<Self, SqliteError> {
338 Self::sqlite_with_pool_config(
339 path,
340 PoolConfig {
341 wal_mode,
342 busy_timeout,
343 write_queue_enabled: Some(true),
344 ..PoolConfig::for_test()
345 },
346 None,
347 )
348 }
349
350 #[cfg(any(test, feature = "test-support"))]
355 pub fn sqlite_for_test_with_journal_mode_in(
356 path: impl AsRef<Path>,
357 wal_mode: bool,
358 busy_timeout: std::time::Duration,
359 volume_lock_dir: std::path::PathBuf,
360 ) -> Result<Self, SqliteError> {
361 Self::sqlite_with_pool_config(
362 path,
363 PoolConfig {
364 wal_mode,
365 busy_timeout,
366 write_queue_enabled: Some(true),
367 volume_lock_dir: Some(volume_lock_dir),
368 ..PoolConfig::for_test()
369 },
370 None,
371 )
372 }
373
374 pub fn sqlite_with_max_readers(
378 path: impl AsRef<Path>,
379 max_readers: Option<usize>,
380 ) -> Result<Self, SqliteError> {
381 Self::sqlite_with_pool_config(path, PoolConfig::default(), max_readers)
382 }
383
384 pub fn sqlite_with_max_readers_and_wal_ceiling(
388 path: impl AsRef<Path>,
389 max_readers: Option<usize>,
390 wal_ceiling: WalCeilingPolicy,
391 ) -> Result<Self, SqliteError> {
392 Self::sqlite_with_pool_config(
393 path,
394 PoolConfig {
395 wal_ceiling,
396 ..PoolConfig::default()
397 },
398 max_readers,
399 )
400 }
401
402 pub(crate) fn sqlite_with_pool_config(
403 path: impl AsRef<Path>,
404 pool_config: PoolConfig,
405 max_readers: Option<usize>,
406 ) -> Result<Self, SqliteError> {
407 crate::extension::ensure_extensions_loaded();
408 let resolved = path.as_ref().to_path_buf();
409 let read_only =
410 std::fs::metadata(&resolved).is_ok_and(|metadata| metadata.permissions().readonly());
411 let mut config = PoolConfig {
412 path: Some(resolved.clone()),
413 read_only,
414 ..pool_config
415 };
416 if let Some(max_readers) = max_readers {
417 config.max_readers = max_readers;
418 }
419 if read_only {
420 config.write_queue_enabled = Some(false);
421 }
422 let pool = ConnectionPool::new(config)?;
423 Ok(Self {
424 pool: Arc::new(pool),
425 is_file_backed: true,
426 path: Some(resolved),
427 vector_tables_ready: Default::default(),
428 notes_seq_repair_runs: AtomicUsize::new(0),
429 store_schemas: std::array::from_fn(|_| Arc::new(StoreSchemaGate::default())),
430 })
431 }
432
433 pub fn sqlite_read_only(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
444 Self::sqlite_read_only_with_pool_config(path, PoolConfig::default(), None)
445 }
446
447 #[cfg(any(test, feature = "test-support"))]
449 pub fn sqlite_read_only_for_test(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
450 Self::sqlite_read_only_with_pool_config(path, PoolConfig::for_test(), None)
451 }
452
453 pub fn sqlite_read_only_with_max_readers(
455 path: impl AsRef<Path>,
456 max_readers: Option<usize>,
457 ) -> Result<Self, SqliteError> {
458 Self::sqlite_read_only_with_pool_config(path, PoolConfig::default(), max_readers)
459 }
460
461 pub fn sqlite_read_only_with_max_readers_and_wal_ceiling(
464 path: impl AsRef<Path>,
465 max_readers: Option<usize>,
466 wal_ceiling: WalCeilingPolicy,
467 ) -> Result<Self, SqliteError> {
468 Self::sqlite_read_only_with_pool_config(
469 path,
470 PoolConfig {
471 wal_ceiling,
472 ..PoolConfig::default()
473 },
474 max_readers,
475 )
476 }
477
478 fn sqlite_read_only_with_pool_config(
479 path: impl AsRef<Path>,
480 pool_config: PoolConfig,
481 max_readers: Option<usize>,
482 ) -> Result<Self, SqliteError> {
483 crate::extension::ensure_extensions_loaded();
484 let resolved = path.as_ref().to_path_buf();
485 let mut config = PoolConfig {
486 path: Some(resolved.clone()),
487 read_only: true,
488 write_queue_enabled: Some(false),
489 ..pool_config
490 };
491 if let Some(max_readers) = max_readers {
492 config.max_readers = max_readers;
493 }
494 let pool = ConnectionPool::new(config)?;
500 Ok(Self {
501 pool: Arc::new(pool),
502 is_file_backed: true,
503 path: Some(resolved),
504 vector_tables_ready: Default::default(),
505 notes_seq_repair_runs: AtomicUsize::new(0),
506 store_schemas: std::array::from_fn(|_| Arc::new(StoreSchemaGate::default())),
507 })
508 }
509
510 pub fn memory() -> Result<Self, SqliteError> {
516 crate::extension::ensure_extensions_loaded();
517 let config = PoolConfig {
518 path: None,
519 ..PoolConfig::default()
520 };
521 let pool = ConnectionPool::new(config)?;
522 Ok(Self {
523 pool: Arc::new(pool),
524 is_file_backed: false,
525 path: None,
526 vector_tables_ready: Default::default(),
527 notes_seq_repair_runs: AtomicUsize::new(0),
528 store_schemas: std::array::from_fn(|_| Arc::new(StoreSchemaGate::default())),
529 })
530 }
531
532 pub fn sql(&self) -> Arc<dyn khive_storage::SqlAccess> {
536 Arc::new(SqlBridge::new(Arc::clone(&self.pool), self.is_file_backed))
537 }
538
539 pub fn apply_schema(
546 &self,
547 plan: &crate::migrations::ServiceSchemaPlan,
548 ) -> Result<(), SqliteError> {
549 let admission = self.pool.write_admission();
550 crate::migrations::apply_schema_plan_with_admission(
551 &mut self.pool.migration_transactions(),
552 plan,
553 &admission,
554 )
555 }
556
557 pub fn apply_pack_ddl_statements(
573 &self,
574 statements: &[&'static str],
575 ) -> Result<(), SqliteError> {
576 self.apply_pack_ddl_statements_with_columns(statements, &[])
577 }
578
579 pub fn apply_pack_ddl_statements_with_columns(
586 &self,
587 statements: &[&'static str],
588 additions: &[khive_types::PackColumnAddition],
589 ) -> Result<(), SqliteError> {
590 let writer = self.pool.writer_for_admitted_operation()?;
591 writer.transaction(|conn| {
592 pack_schema::add_missing_columns(conn, additions)?;
593 for &stmt in statements {
594 conn.execute_batch(stmt)?;
595 }
596 pack_schema::validate_columns(conn, additions)?;
597 Ok(())
598 })
599 }
600
601 pub fn validate_pack_schema_columns(
604 &self,
605 additions: &[khive_types::PackColumnAddition],
606 ) -> Result<(), SqliteError> {
607 if additions.is_empty() {
608 return Ok(());
609 }
610 let reader = self.pool.reader()?;
611 pack_schema::validate_columns(reader.conn(), additions)
612 }
613
614 pub fn prepare_core_schema(&self) -> Result<u32, SqliteError> {
624 if self.is_read_only() {
625 let reader = self.pool.reader()?;
626 crate::migrations::validate_schema_is_current(reader.conn())
627 } else {
628 let latest = crate::migrations::MIGRATIONS
629 .last()
630 .map(|migration| migration.version)
631 .unwrap_or(0);
632 {
633 let reader = self.pool.reader()?;
634 let current = crate::migrations::read_schema_version(reader.conn())?;
635 if current >= latest {
636 return crate::migrations::validate_schema_is_current(reader.conn());
637 }
638 }
639 let owner = crate::stores::blob::acquire_database_gc_owner_for_path_blocking(
640 self.pool.canonical_path().map(Path::to_path_buf),
641 )
642 .map_err(|error| {
643 SqliteError::InvalidData(format!(
644 "failed to acquire database GC owner before schema preparation: {error}"
645 ))
646 })?;
647 self.run_core_migrations(&owner)
648 }
649 }
650
651 pub fn attachment_cutover_status(
653 &self,
654 ) -> Result<crate::migrations::AttachmentCutoverStatus, SqliteError> {
655 let reader = self.pool.reader()?;
656 crate::migrations::attachment_cutover_status(reader.conn())
657 }
658
659 fn require_attachment_cutover_owner(
660 &self,
661 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
662 ) -> Result<(), SqliteError> {
663 let sql = self.sql();
664 let backend_path = sql.database_path();
665 if owner.database_path() != backend_path.as_deref() {
666 return Err(SqliteError::InvalidData(format!(
667 "attachment cutover GC owner targets {:?}, but this backend is {:?}",
668 owner.database_path(),
669 backend_path.as_deref()
670 )));
671 }
672 Ok(())
673 }
674
675 pub fn stage_attachment_cutover(
679 &self,
680 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
681 ) -> Result<(), SqliteError> {
682 self.require_attachment_cutover_owner(owner)?;
683 if self.is_read_only() {
684 return Err(SqliteError::InvalidData(
685 "cannot stage attachment cutover on a read-only backend".into(),
686 ));
687 }
688 let mut writer = self.pool.writer_for_admitted_operation()?;
689 let admission = self.pool.write_admission();
690 crate::migrations::stage_attachment_cutover_with_admission(writer.conn_mut(), &admission)
691 }
692
693 pub fn apply_verified_attachments(
695 &self,
696 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
697 attachments: &[khive_storage::Attachment],
698 ) -> Result<(), SqliteError> {
699 self.require_attachment_cutover_owner(owner)?;
700 if self.is_read_only() {
701 return Err(SqliteError::InvalidData(
702 "cannot apply verified attachments on a read-only backend".into(),
703 ));
704 }
705 let mut writer = self.pool.writer_for_admitted_operation()?;
706 let tx = writer
707 .conn_mut()
708 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
709 let admission = self.pool.write_admission();
710 if let Err(error) = admission.check() {
711 let rollback = tx.rollback();
712 return Err(crate::migrations::capacity_refusal_after_rollback(
713 writer.conn(),
714 rollback,
715 error,
716 "verified attachment publication",
717 ));
718 }
719 for attachment in attachments {
720 attachment
721 .validate()
722 .map_err(|error| SqliteError::InvalidData(error.to_string()))?;
723 crate::migrations::apply_generic_verified_attachment(
724 &tx,
725 &attachment.record_uuid.to_string(),
726 attachment.substrate.as_str(),
727 &attachment.role,
728 &attachment.content_ref,
729 attachment.media_type.as_deref(),
730 attachment.size_bytes,
731 attachment.created_at,
732 )?;
733 }
734 tx.commit()?;
735 Ok(())
736 }
737
738 pub fn finalize_attachment_cutover(
741 &self,
742 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
743 ) -> Result<(), SqliteError> {
744 self.require_attachment_cutover_owner(owner)?;
745 if self.is_read_only() {
746 return Err(SqliteError::InvalidData(
747 "cannot finalize attachment cutover on a read-only backend".into(),
748 ));
749 }
750 let mut writer = self.pool.writer_for_admitted_operation()?;
751 let admission = self.pool.write_admission();
752 crate::migrations::finalize_attachment_cutover_with_admission(writer.conn_mut(), &admission)
753 }
754
755 pub fn entities(&self) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
759 self.entities_for_namespace("local")
760 }
761
762 pub fn entities_for_namespace(
766 &self,
767 namespace: &str,
768 ) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
769 if namespace.trim().is_empty() {
770 return Err(SqliteError::InvalidData(
771 "entities namespace must be non-empty".to_string(),
772 ));
773 }
774 self.ensure_store_schema(StoreSchemaKind::Entities)?;
775
776 Ok(Arc::new(entity::SqlEntityStore::new(
777 Arc::clone(&self.pool),
778 self.is_file_backed,
779 )))
780 }
781
782 pub fn attachments(&self) -> Result<Arc<dyn khive_storage::AttachmentStore>, SqliteError> {
789 Ok(Arc::new(attachment::SqlAttachmentStore::new(
790 Arc::clone(&self.pool),
791 self.is_file_backed,
792 )))
793 }
794
795 pub fn graph(&self) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
800 self.graph_for_namespace("local")
801 }
802
803 pub fn graph_for_namespace(
805 &self,
806 namespace: &str,
807 ) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
808 if namespace.trim().is_empty() {
809 return Err(SqliteError::InvalidData(
810 "graph namespace must be non-empty".to_string(),
811 ));
812 }
813 self.ensure_store_schema(StoreSchemaKind::Graph)?;
814
815 Ok(Arc::new(
816 graph::SqlGraphStore::new_scoped(
817 Arc::clone(&self.pool),
818 self.is_file_backed,
819 namespace.trim().to_string(),
820 )
821 .with_index_repair(crate::stores::index_repair::IndexRepairContext::new(
822 Arc::clone(&self.pool),
823 self.store_schemas.clone(),
824 crate::stores::index_repair::IndexReadKind::Graph,
825 )),
826 ))
827 }
828
829 fn ensure_store_schema(&self, kind: StoreSchemaKind) -> Result<(), SqliteError> {
830 if self.is_read_only()
831 || self.store_schemas[kind as usize]
832 .ready
833 .load(Ordering::Acquire)
834 {
835 return Ok(());
836 }
837 let writer = self.constructor_writer()?;
838 self.ensure_store_schema_with_writer(kind, writer.conn())
839 }
840
841 fn ensure_store_schema_with_writer(
842 &self,
843 kind: StoreSchemaKind,
844 conn: &rusqlite::Connection,
845 ) -> Result<(), SqliteError> {
846 self.store_schemas[kind as usize].ensure(conn, kind.initializer())?;
847 Ok(())
848 }
849
850 fn constructor_writer(
851 &self,
852 ) -> Result<crate::pool::PooledAutocommitWriteUnit<'_>, SqliteError> {
853 let context = khive_storage::capture_request_read_context();
854 let Some(operation) = context.store_acquisition_operation() else {
855 return self.pool.autocommit_write_unit();
856 };
857 let writer = self
858 .pool
859 .writer_until_for_admitted_operation(|| context.blocking_stop_reason().is_some())?
860 .ok_or_else(|| {
861 SqliteError::RequestReadStopped(khive_storage::StorageError::Timeout {
862 operation: operation.into(),
863 })
864 })?;
865 writer.admit_autocommit()
866 }
867
868 pub fn notes(&self) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
872 self.notes_for_namespace("local")
873 }
874
875 pub fn notes_for_namespace(
879 &self,
880 namespace: &str,
881 ) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
882 if namespace.trim().is_empty() {
883 return Err(SqliteError::InvalidData(
884 "notes namespace must be non-empty".to_string(),
885 ));
886 }
887 if !self.is_read_only()
888 && (!self.store_schemas[StoreSchemaKind::Notes as usize]
889 .ready
890 .load(Ordering::Acquire)
891 || self.notes_seq_repair_runs.load(Ordering::Relaxed) == 0)
892 {
893 let writer = self.constructor_writer()?;
894 self.ensure_store_schema_with_writer(StoreSchemaKind::Notes, writer.conn())?;
895
896 if self.notes_seq_repair_runs.load(Ordering::Relaxed) == 0 {
903 note::repair_notes_seq(writer.conn())?;
904 self.notes_seq_repair_runs.fetch_add(1, Ordering::Relaxed);
905 }
906 }
907
908 Ok(Arc::new(
909 note::SqlNoteStore::new(Arc::clone(&self.pool), self.is_file_backed).with_index_repair(
910 crate::stores::index_repair::IndexRepairContext::new(
911 Arc::clone(&self.pool),
912 self.store_schemas.clone(),
913 crate::stores::index_repair::IndexReadKind::Notes,
914 ),
915 ),
916 ))
917 }
918
919 pub fn notes_seq_repair_run_count(&self) -> usize {
924 self.notes_seq_repair_runs.load(Ordering::Relaxed)
925 }
926
927 pub fn events(&self) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
932 self.events_for_namespace("local")
933 }
934
935 pub fn events_for_namespace(
937 &self,
938 namespace: &str,
939 ) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
940 if namespace.trim().is_empty() {
941 return Err(SqliteError::InvalidData(
942 "events namespace must be non-empty".to_string(),
943 ));
944 }
945 self.ensure_store_schema(StoreSchemaKind::Events)?;
946
947 Ok(Arc::new(event::SqlEventStore::new_scoped(
948 Arc::clone(&self.pool),
949 self.is_file_backed,
950 namespace.trim().to_string(),
951 )))
952 }
953
954 pub fn agents(&self) -> Result<Arc<dyn khive_storage::AgentStore>, SqliteError> {
959 self.ensure_store_schema(StoreSchemaKind::Agents)?;
960
961 Ok(Arc::new(agents::SqlAgentStore::new(
962 Arc::clone(&self.pool),
963 self.is_file_backed,
964 )))
965 }
966
967 pub fn vectors(
973 &self,
974 model_key: &str,
975 embedding_model: &str,
976 dimensions: usize,
977 ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
978 self.vectors_for_namespace(model_key, embedding_model, dimensions, "local")
979 }
980
981 pub fn vectors_for_namespace(
991 &self,
992 model_key: &str,
993 embedding_model: &str,
994 dimensions: usize,
995 namespace: &str,
996 ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
997 validate_vector_model_key(model_key)?;
998 if namespace.trim().is_empty() {
999 return Err(SqliteError::InvalidData(
1000 "vector store namespace must be non-empty".to_string(),
1001 ));
1002 }
1003 self.ensure_vector_tables(&[(model_key, dimensions)])?;
1004 Ok(Arc::new(vectors::SqliteVecStore::new(
1005 Arc::clone(&self.pool),
1006 self.is_file_backed,
1007 model_key.to_string(),
1008 embedding_model.to_string(),
1009 dimensions,
1010 namespace.trim().to_string(),
1011 )?))
1012 }
1013
1014 pub fn ensure_vector_tables(&self, models: &[(&str, usize)]) -> Result<(), SqliteError> {
1023 for (model_key, _) in models {
1024 validate_vector_model_key(model_key)?;
1025 }
1026 if models.is_empty() {
1027 return Ok(());
1028 }
1029
1030 crate::extension::ensure_extensions_loaded();
1032
1033 if self.is_read_only() {
1038 let reader = self.pool.reader()?;
1042 for (model_key, _) in models {
1043 let table = format!("vec_{model_key}");
1044 if !sqlite_table_exists(reader.conn(), &table)? {
1045 return Err(SqliteError::InvalidData(format!(
1046 "read-only database has no vector table '{table}'; create and populate it in \
1047 a writable copy before opening the snapshot"
1048 )));
1049 }
1050 validate_vector_table_columns(reader.conn(), &table)?;
1051 }
1052 return Ok(());
1053 }
1054
1055 let pending = self.unprepared_vector_tables(models);
1056 if pending.is_empty() {
1057 return Ok(());
1058 }
1059 let writer = self.constructor_writer()?;
1060 let pending = self.unprepared_vector_tables(&pending);
1063 if pending.is_empty() {
1064 return Ok(());
1065 }
1066
1067 for (model_key, _) in &pending {
1071 let table = format!("vec_{model_key}");
1072 if sqlite_table_exists(writer.conn(), &table)? {
1078 validate_vector_table_columns(writer.conn(), &table)?;
1079 }
1080 }
1081
1082 writer
1091 .conn()
1092 .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
1093
1094 writer
1098 .conn()
1099 .execute_batch(crate::migrations::ANN_WRITE_LOG_DDL)?;
1100 writer
1101 .conn()
1102 .execute_batch(crate::migrations::ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL)?;
1103 writer
1104 .conn()
1105 .execute_batch(crate::migrations::ANN_CONSUMER_PENDING_DDL)?;
1106 for (model_key, dimensions) in &pending {
1108 let ddl = format!(
1109 "CREATE VIRTUAL TABLE IF NOT EXISTS vec_{} USING vec0(\
1110 subject_id TEXT PRIMARY KEY, \
1111 namespace TEXT NOT NULL, \
1112 kind TEXT NOT NULL, \
1113 field TEXT NOT NULL, \
1114 embedding_model TEXT NOT NULL, \
1115 embedding float[{}] distance_metric=cosine\
1116 )",
1117 model_key, dimensions
1118 );
1119 writer.conn().execute_batch(&ddl)?;
1120 }
1121 self.vector_tables_ready.lock().extend(
1124 pending
1125 .iter()
1126 .map(|(model_key, _)| (*model_key).to_string()),
1127 );
1128 Ok(())
1129 }
1130
1131 fn unprepared_vector_tables<'a>(&self, models: &[(&'a str, usize)]) -> Vec<(&'a str, usize)> {
1133 let ready = self.vector_tables_ready.lock();
1134 models
1135 .iter()
1136 .filter(|(model_key, _)| !ready.contains(*model_key))
1137 .copied()
1138 .collect()
1139 }
1140
1141 pub fn register_embedding_model(
1146 &self,
1147 engine_name: &str,
1148 model_id: &str,
1149 key_version: &str,
1150 dimensions: u32,
1151 ) -> Result<(), SqliteError> {
1152 let writer = self.pool.autocommit_write_unit()?;
1153 writer
1154 .conn()
1155 .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
1156
1157 let now = chrono::Utc::now().timestamp_micros();
1158 let canonical_key =
1159 format!("{engine_name}:{model_id}:{key_version}:{dimensions}").into_bytes();
1160 let id = uuid::Uuid::new_v4();
1161 writer.conn().execute(
1162 "INSERT INTO _embedding_models \
1163 (id, engine_name, model_id, key_version, dim, output_dim, status, \
1164 activated_at, superseded_at, superseded_by, canonical_key, created_at) \
1165 VALUES (?1, ?2, ?3, ?4, ?5, NULL, 'active', ?6, NULL, NULL, ?7, ?8) \
1166 ON CONFLICT(canonical_key) DO UPDATE SET \
1167 status = 'active', \
1168 activated_at = COALESCE(_embedding_models.activated_at, excluded.activated_at)",
1169 rusqlite::params![
1170 id.as_bytes().as_slice(),
1171 engine_name,
1172 model_id,
1173 key_version,
1174 dimensions as i64,
1175 now,
1176 canonical_key,
1177 now,
1178 ],
1179 )?;
1180 Ok(())
1181 }
1182
1183 pub fn sparse(
1187 &self,
1188 model_key: &str,
1189 ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
1190 self.sparse_for_namespace(model_key, "local")
1191 }
1192
1193 pub fn sparse_for_namespace(
1197 &self,
1198 model_key: &str,
1199 namespace: &str,
1200 ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
1201 if model_key.is_empty()
1202 || !model_key
1203 .chars()
1204 .all(|c| c.is_ascii_alphanumeric() || c == '_')
1205 {
1206 return Err(SqliteError::InvalidData(format!(
1207 "invalid model_key '{}': must be non-empty and contain only alphanumeric/underscore characters",
1208 model_key
1209 )));
1210 }
1211 if namespace.trim().is_empty() {
1212 return Err(SqliteError::InvalidData(
1213 "sparse store namespace must be non-empty".to_string(),
1214 ));
1215 }
1216
1217 if self.is_read_only() {
1218 let table = format!("sparse_{model_key}");
1219 let reader = self.pool.reader()?;
1220 if !sqlite_table_exists(reader.conn(), &table)? {
1221 return Err(SqliteError::InvalidData(format!(
1222 "read-only database has no sparse table '{table}'; create and populate it in \
1223 a writable copy before opening the snapshot"
1224 )));
1225 }
1226 } else {
1227 let writer = self.pool.autocommit_write_unit()?;
1228 sparse::ensure_sparse_schema(writer.conn(), model_key)
1229 .map_err(SqliteError::Rusqlite)?;
1230 }
1231
1232 Ok(Arc::new(sparse::SqliteSparseStore::new(
1233 Arc::clone(&self.pool),
1234 self.is_file_backed,
1235 model_key.to_string(),
1236 namespace.trim().to_string(),
1237 )?))
1238 }
1239
1240 pub fn text(&self, table_key: &str) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
1247 self.text_with_tokenizer(table_key, "trigram")
1248 }
1249
1250 pub fn text_with_tokenizer(
1258 &self,
1259 table_key: &str,
1260 tokenizer: &str,
1261 ) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
1262 if table_key.is_empty()
1263 || !table_key
1264 .chars()
1265 .all(|c| c.is_ascii_alphanumeric() || c == '_')
1266 {
1267 return Err(SqliteError::InvalidData(format!(
1268 "invalid table_key '{}': must be non-empty and contain only \
1269 alphanumeric/underscore characters",
1270 table_key
1271 )));
1272 }
1273 if table_key.ends_with("_rowids") || table_key.ends_with("_rowids_state") {
1283 return Err(SqliteError::InvalidData(format!(
1284 "invalid table_key '{}': must not end in '_rowids' or '_rowids_state' — those \
1285 suffixes are reserved for a text table's own rowid-map sidecar and its \
1286 completion-marker state table (see text::rowid_map_table, \
1287 text::rowid_map_state_table)",
1288 table_key
1289 )));
1290 }
1291 if tokenizer.is_empty()
1292 || !tokenizer
1293 .chars()
1294 .all(|c| c.is_ascii_alphanumeric() || c == '_')
1295 {
1296 return Err(SqliteError::InvalidData(format!(
1297 "invalid tokenizer '{}': must be non-empty and contain only \
1298 alphanumeric/underscore characters",
1299 tokenizer
1300 )));
1301 }
1302
1303 let ddl = format!(
1304 "CREATE VIRTUAL TABLE IF NOT EXISTS fts_{} USING fts5(\
1305 subject_id UNINDEXED, \
1306 kind UNINDEXED, \
1307 title, \
1308 body, \
1309 tags UNINDEXED, \
1310 namespace UNINDEXED, \
1311 metadata UNINDEXED, \
1312 updated_at UNINDEXED, \
1313 record_kind, \
1314 tokenize = '{}'\
1315 )",
1316 table_key, tokenizer
1317 );
1318 let table = format!("fts_{table_key}");
1319 if self.is_read_only() {
1320 let reader = self.pool.reader()?;
1321 if !sqlite_table_exists(reader.conn(), &table)? {
1322 return Err(SqliteError::InvalidData(format!(
1323 "read-only database has no text-search table '{table}'; create and populate \
1324 it in a writable copy before opening the snapshot"
1325 )));
1326 }
1327 let map = text::rowid_map_table(&table);
1341 let map_exists = sqlite_table_exists(reader.conn(), &map)?;
1342 let marker_present = if map_exists {
1343 let state = text::rowid_map_state_table(&table);
1344 sqlite_table_exists(reader.conn(), &state)?
1345 && reader.conn().query_row(
1346 &format!(
1347 "SELECT EXISTS(SELECT 1 FROM {state} WHERE key = 'backfill' AND value = ?1)"
1348 ),
1349 rusqlite::params![text::ROWID_MAP_BACKFILL_COMPLETE],
1350 |row| row.get(0),
1351 )?
1352 } else {
1353 false
1354 };
1355 if !map_exists || !marker_present {
1356 warn_scan_fallback_once(&table);
1357 return Ok(Arc::new(text::Fts5TextSearch::new_scan_fallback(
1358 Arc::clone(&self.pool),
1359 self.is_file_backed,
1360 table_key.to_string(),
1361 )));
1362 }
1363 } else {
1364 let writer = self.pool.autocommit_write_unit()?;
1365 writer.conn().execute_batch(&ddl)?;
1366 writer.conn().execute_batch(&text::rowid_map_ddl(&table))?;
1367 ensure_fts_rowid_map_backfilled(writer.conn(), &table, &self.pool.write_admission())?;
1368 }
1369
1370 Ok(Arc::new(text::Fts5TextSearch::new(
1371 Arc::clone(&self.pool),
1372 self.is_file_backed,
1373 table_key.to_string(),
1374 )))
1375 }
1376
1377 pub fn blob_store(
1385 &self,
1386 config_root: Option<&Path>,
1387 floor_bytes: Option<u64>,
1388 ) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
1389 let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
1390 let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
1391 Ok(Arc::new(blob::FsBlobStore::new(root, floor)?))
1392 }
1393
1394 pub fn blob_store_read_only(
1399 &self,
1400 config_root: Option<&Path>,
1401 floor_bytes: Option<u64>,
1402 ) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
1403 let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
1404 let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
1405 Ok(Arc::new(blob::FsBlobStore::open_existing(root, floor)?))
1406 }
1407
1408 pub fn is_file_backed(&self) -> bool {
1410 self.is_file_backed
1411 }
1412
1413 pub fn is_read_only(&self) -> bool {
1416 self.pool.config().read_only
1417 }
1418
1419 pub fn data_dir(&self) -> Option<std::path::PathBuf> {
1422 self.path.as_ref()?.parent().map(|p| p.to_path_buf())
1423 }
1424
1425 pub fn ann_root(&self) -> Option<std::path::PathBuf> {
1433 ann_root_for(self.path.as_ref()?)
1434 }
1435
1436 pub fn pool(&self) -> &ConnectionPool {
1438 &self.pool
1439 }
1440
1441 pub fn database_owner_identity(
1444 &self,
1445 ) -> Result<DatabaseOwnerIdentity, DatabaseOwnerIdentityError> {
1446 self.pool.database_owner_identity()
1447 }
1448
1449 pub fn verify_database_owner(
1451 &self,
1452 expected: &DatabaseOwnerIdentity,
1453 ) -> Result<(), DatabaseOwnerIdentityError> {
1454 self.pool.verify_database_owner(expected)
1455 }
1456
1457 pub fn pool_arc(&self) -> Arc<ConnectionPool> {
1459 Arc::clone(&self.pool)
1460 }
1461}
1462
1463fn ann_root_for(path: &std::path::Path) -> Option<std::path::PathBuf> {
1468 let mut file = path.file_name()?.to_os_string();
1469 file.push(".ann");
1470 path.parent().map(|p| p.join(file))
1471}
1472
1473#[cfg(test)]
1474#[path = "backend/store_accessor_tests.rs"]
1475mod store_accessor_tests;
1476
1477#[cfg(test)]
1478#[path = "backend/store_accessor_index_tests.rs"]
1479mod store_accessor_index_tests;
1480
1481#[cfg(test)]
1482mod tests {
1483 use super::*;
1484 use khive_storage::types::{EdgeFilter, SqlStatement, SqlValue};
1485 use khive_storage::{EntityFilter, EventFilter};
1486
1487 include!("backend_read_admission_tests.rs");
1488
1489 #[cfg(unix)]
1490 #[tokio::test]
1491 async fn sqlite_detects_chmod_read_only_snapshot_and_core_reads_succeed() {
1492 use std::os::unix::fs::PermissionsExt;
1493
1494 let dir = tempfile::tempdir().unwrap();
1495 let path = dir.path().join("chmod_snapshot.db");
1496 {
1497 let writable =
1498 StorageBackend::sqlite_for_test(&path).expect("create writable database");
1499 writable
1500 .prepare_core_schema()
1501 .expect("migrate writable snapshot source");
1502 }
1503
1504 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
1505 permissions.set_mode(0o444);
1506 std::fs::set_permissions(&path, permissions).unwrap();
1507 freeze_snapshot_sidecars(&path);
1508
1509 let read_only = StorageBackend::sqlite(&path).expect("auto-detect read-only mode");
1510 assert!(read_only.is_read_only());
1511 assert_eq!(
1512 read_only.pool().config().write_queue_enabled,
1513 Some(false),
1514 "read-only boot must not attempt to spawn a writer task"
1515 );
1516 assert!(read_only
1517 .pool()
1518 .writer_task_handle()
1519 .expect("disabled writer task is a valid configuration")
1520 .is_none());
1521 read_only
1522 .prepare_core_schema()
1523 .expect("current snapshot validates without migration writes");
1524
1525 let entities = read_only.entities().expect("entity store opens read-only");
1526 let graph = read_only.graph().expect("graph store opens read-only");
1527 let notes = read_only.notes().expect("note store opens read-only");
1528 let events = read_only.events().expect("event store opens read-only");
1529 assert_eq!(
1530 entities
1531 .count_entities("local", khive_storage::EntityFilter::default())
1532 .await
1533 .unwrap(),
1534 0
1535 );
1536 assert_eq!(
1537 graph
1538 .count_edges(khive_storage::types::EdgeFilter::default())
1539 .await
1540 .unwrap(),
1541 0
1542 );
1543 assert_eq!(notes.count_notes("local", None).await.unwrap(), 0);
1544 assert_eq!(
1545 events
1546 .count_events(khive_storage::EventFilter::default())
1547 .await
1548 .unwrap(),
1549 0
1550 );
1551 assert_eq!(
1552 read_only.notes_seq_repair_run_count(),
1553 0,
1554 "read-only store acquisition must not run the DML repair"
1555 );
1556 }
1557
1558 #[test]
1559 fn memory_backend_creates_successfully() {
1560 let backend = StorageBackend::memory().expect("memory backend should create");
1561 assert!(!backend.is_file_backed());
1562 }
1563
1564 #[test]
1565 fn file_backend_creates_successfully() {
1566 let dir = tempfile::tempdir().unwrap();
1567 let path = dir.path().join("test.db");
1568 let backend = StorageBackend::sqlite(&path).expect("file backend should create");
1569 assert!(backend.is_file_backed());
1570 assert!(path.exists());
1571 }
1572
1573 #[test]
1574 fn data_dir_returns_none_for_memory_backend() {
1575 let backend = StorageBackend::memory().expect("memory backend");
1576 assert!(backend.data_dir().is_none());
1577 }
1578
1579 #[test]
1580 fn data_dir_returns_parent_dir_for_file_backend() {
1581 let dir = tempfile::tempdir().unwrap();
1582 let path = dir.path().join("data.db");
1583 let backend = StorageBackend::sqlite_for_test(&path).expect("file backend");
1584 let got = backend.data_dir().expect("file backend must return Some");
1585 assert_eq!(got, dir.path());
1586 }
1587
1588 include!("backend/ann_root_tests.rs");
1589
1590 #[tokio::test]
1591 async fn sql_access_memory_roundtrip() {
1592 let backend = StorageBackend::memory().unwrap();
1593 let sql = backend.sql();
1594
1595 let mut writer = sql.writer().await.unwrap();
1596 writer
1597 .execute_script(
1598 "CREATE TABLE test_rt (id TEXT PRIMARY KEY, value INTEGER NOT NULL)".into(),
1599 )
1600 .await
1601 .unwrap();
1602
1603 let affected = writer
1604 .execute(SqlStatement {
1605 sql: "INSERT INTO test_rt (id, value) VALUES (?1, ?2)".into(),
1606 params: vec![SqlValue::Text("row1".into()), SqlValue::Integer(42)],
1607 label: None,
1608 })
1609 .await
1610 .unwrap();
1611 assert_eq!(affected, 1);
1612
1613 let mut reader = sql.reader().await.unwrap();
1614 let row = reader
1615 .query_row(SqlStatement {
1616 sql: "SELECT id, value FROM test_rt WHERE id = ?1".into(),
1617 params: vec![SqlValue::Text("row1".into())],
1618 label: None,
1619 })
1620 .await
1621 .unwrap();
1622
1623 let row = row.expect("should find the inserted row");
1624 assert_eq!(row.columns.len(), 2);
1625 match &row.columns[0].value {
1626 SqlValue::Text(s) => assert_eq!(s, "row1"),
1627 other => panic!("expected Text, got {other:?}"),
1628 }
1629 match &row.columns[1].value {
1630 SqlValue::Integer(v) => assert_eq!(*v, 42),
1631 other => panic!("expected Integer, got {other:?}"),
1632 }
1633 }
1634
1635 #[tokio::test]
1636 async fn sql_access_file_roundtrip() {
1637 let dir = tempfile::tempdir().unwrap();
1638 let path = dir.path().join("test_roundtrip.db");
1639 let backend = StorageBackend::sqlite_for_test(&path).unwrap();
1640 let sql = backend.sql();
1641
1642 let mut writer = sql.writer().await.unwrap();
1643 writer
1644 .execute_script("CREATE TABLE test_f (k TEXT PRIMARY KEY, v TEXT)".into())
1645 .await
1646 .unwrap();
1647 writer
1648 .execute(SqlStatement {
1649 sql: "INSERT INTO test_f (k, v) VALUES (?1, ?2)".into(),
1650 params: vec![
1651 SqlValue::Text("hello".into()),
1652 SqlValue::Text("world".into()),
1653 ],
1654 label: None,
1655 })
1656 .await
1657 .unwrap();
1658
1659 let mut reader = sql.reader().await.unwrap();
1660 let rows = reader
1661 .query_all(SqlStatement {
1662 sql: "SELECT k, v FROM test_f".into(),
1663 params: vec![],
1664 label: None,
1665 })
1666 .await
1667 .unwrap();
1668 assert_eq!(rows.len(), 1);
1669 match &rows[0].columns[1].value {
1670 SqlValue::Text(s) => assert_eq!(s, "world"),
1671 other => panic!("expected Text, got {other:?}"),
1672 }
1673 }
1674
1675 #[test]
1676 fn sqlite_read_only_missing_path_does_not_create_file() {
1677 let dir = tempfile::tempdir().unwrap();
1678 let path = dir.path().join("missing_ro.db");
1679 assert!(!path.exists());
1680
1681 let result = StorageBackend::sqlite_read_only(&path);
1682 assert!(
1683 result.is_err(),
1684 "opening a missing path read-only must fail"
1685 );
1686 assert!(
1687 !path.exists(),
1688 "opening a missing path read-only must not create the file"
1689 );
1690 }
1691
1692 #[test]
1693 fn sqlite_read_only_sparse_store_requires_existing_table_without_writer_acquisition() {
1694 let dir = tempfile::tempdir().unwrap();
1695 let path = dir.path().join("ro_sparse_tables.db");
1696 {
1697 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
1698 writable
1699 .prepare_core_schema()
1700 .expect("migrate snapshot source");
1701 writable
1702 .sparse("present")
1703 .expect("create the optional sparse table while writable");
1704 }
1705 #[cfg(unix)]
1706 freeze_snapshot_sidecars(&path);
1707
1708 let read_only = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
1709 read_only
1710 .prepare_core_schema()
1711 .expect("validate exact current migration ledger");
1712 read_only
1713 .sparse("present")
1714 .expect("an existing sparse table must open read-only");
1715 let missing = match read_only.sparse("missing") {
1716 Ok(_) => panic!("a missing sparse table must fail during store acquisition"),
1717 Err(error) => error,
1718 };
1719 assert!(
1720 missing.to_string().contains("sparse_missing"),
1721 "the diagnostic must name the absent table: {missing}"
1722 );
1723 assert_eq!(
1724 read_only.pool().writer_acquisition_snapshot(),
1725 crate::pool::WriterAcquisitionSnapshot::default(),
1726 "construction, exact-ledger validation, and optional sparse-table inspection must \
1727 use reader connections only"
1728 );
1729 }
1730
1731 #[test]
1732 fn sqlite_read_only_text_store_requires_existing_table_without_writer_acquisition() {
1733 let dir = tempfile::tempdir().unwrap();
1734 let path = dir.path().join("ro_text_tables.db");
1735 {
1736 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
1737 writable
1738 .prepare_core_schema()
1739 .expect("migrate snapshot source");
1740 writable
1741 .text("present")
1742 .expect("create the optional FTS table while writable");
1743 }
1744 #[cfg(unix)]
1745 freeze_snapshot_sidecars(&path);
1746
1747 let read_only = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
1748 read_only
1749 .prepare_core_schema()
1750 .expect("validate exact current migration ledger");
1751 read_only
1752 .text("present")
1753 .expect("an existing FTS table must open read-only");
1754 let missing = match read_only.text("missing") {
1755 Ok(_) => panic!("a missing FTS table must fail during store acquisition"),
1756 Err(error) => error,
1757 };
1758 assert!(
1759 missing.to_string().contains("fts_missing"),
1760 "the diagnostic must name the absent table: {missing}"
1761 );
1762 assert_eq!(
1763 read_only.pool().writer_acquisition_snapshot(),
1764 crate::pool::WriterAcquisitionSnapshot::default(),
1765 "construction, exact-ledger validation, and optional FTS inspection must use reader \
1766 connections only"
1767 );
1768 }
1769
1770 #[cfg(feature = "vectors")]
1771 #[test]
1772 fn sqlite_read_only_vector_store_schema_check_uses_no_writer_acquisition() {
1773 let dir = tempfile::tempdir().unwrap();
1774 let path = dir.path().join("ro_vector_tables.db");
1775 {
1776 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
1777 writable
1778 .prepare_core_schema()
1779 .expect("migrate snapshot source");
1780 writable
1781 .vectors("present", "present", 3)
1782 .expect("create the optional vector table while writable");
1783 }
1784 #[cfg(unix)]
1785 freeze_snapshot_sidecars(&path);
1786
1787 let read_only = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
1788 read_only
1789 .prepare_core_schema()
1790 .expect("validate exact current migration ledger");
1791 read_only
1792 .vectors("present", "present", 3)
1793 .expect("an existing vector table must open read-only");
1794 assert!(
1795 read_only.vectors("missing", "missing", 3).is_err(),
1796 "a missing vector table must fail during store acquisition"
1797 );
1798 assert_eq!(
1799 read_only.pool().writer_acquisition_snapshot(),
1800 crate::pool::WriterAcquisitionSnapshot::default(),
1801 "construction, exact-ledger validation, and optional vector inspection must use \
1802 reader connections only"
1803 );
1804 }
1805
1806 #[tokio::test]
1807 async fn sqlite_read_only_sql_writer_rejects_ddl_and_insert() {
1808 let dir = tempfile::tempdir().unwrap();
1809 let path = dir.path().join("ro_writer.db");
1810
1811 {
1813 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
1814 let sql = writable.sql();
1815 let mut writer = sql.writer().await.unwrap();
1816 writer
1817 .execute_script("CREATE TABLE ro_existing (id INTEGER PRIMARY KEY)".into())
1818 .await
1819 .unwrap();
1820 }
1821 #[cfg(unix)]
1822 freeze_snapshot_sidecars(&path);
1823
1824 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1825 let sql = ro.sql();
1826
1827 let writer_result = sql.writer().await;
1829 assert!(
1830 writer_result.is_err(),
1831 "sql().writer() must be rejected on a read-only backend"
1832 );
1833 }
1834
1835 #[tokio::test]
1836 #[cfg(feature = "vectors")]
1837 async fn vectors_roundtrip_via_public_api() {
1838 let backend = StorageBackend::memory().unwrap();
1839 let store = backend.vectors("test_api", "test_api", 3).unwrap();
1840
1841 let id = uuid::Uuid::new_v4();
1842 store
1843 .insert(
1844 id,
1845 khive_types::SubstrateKind::Entity,
1846 "local",
1847 "content",
1848 vec![vec![1.0, 0.0, 0.0]],
1849 )
1850 .await
1851 .unwrap();
1852
1853 let hits = store
1854 .search(khive_storage::types::VectorSearchRequest {
1855 query_vectors: vec![vec![1.0, 0.0, 0.0]],
1856 top_k: 1,
1857 namespace: None,
1858 kind: None,
1859 embedding_model: None,
1860 filter: None,
1861 backend_hints: None,
1862 })
1863 .await
1864 .unwrap();
1865
1866 assert_eq!(hits.len(), 1);
1867 assert_eq!(hits[0].subject_id, id);
1868 assert!(hits[0].score.to_f64() > 0.99);
1869 }
1870
1871 #[tokio::test]
1872 #[cfg(feature = "vectors")]
1873 async fn vectors_direct_store_leaves_provenance_to_migration() {
1874 let backend = StorageBackend::memory().unwrap();
1875 {
1876 let reader = backend.pool.reader().unwrap();
1877 assert!(!sqlite_table_exists(reader.conn(), "vector_provenance").unwrap());
1878 assert_eq!(
1879 crate::migrations::read_schema_version(reader.conn()).unwrap(),
1880 0
1881 );
1882 }
1883 let store = backend
1884 .vectors("direct_provenance", "direct_provenance", 3)
1885 .unwrap();
1886 let reader = backend.pool.reader().unwrap();
1887 assert!(!sqlite_table_exists(reader.conn(), "vector_provenance").unwrap());
1888 assert!(!sqlite_table_exists(reader.conn(), "_schema_migrations").unwrap());
1889 drop(reader);
1890
1891 let id = uuid::Uuid::new_v4();
1892 store
1893 .insert(
1894 id,
1895 khive_types::SubstrateKind::Entity,
1896 "local",
1897 "content",
1898 vec![vec![1.0, 0.0, 0.0]],
1899 )
1900 .await
1901 .unwrap();
1902 let provenance = store.provenance(id).await.unwrap().unwrap();
1903 assert_eq!(provenance.embedding_model, "direct_provenance");
1904 assert_eq!(provenance.text_fingerprint, None);
1905 assert_eq!(provenance.updated_at, None);
1906 assert!(store.delete(id).await.unwrap());
1907 assert!(store.provenance(id).await.unwrap().is_none());
1908 }
1909
1910 #[tokio::test]
1911 #[cfg(feature = "vectors")]
1912 async fn vectors_after_lazy_notes_and_events_schema_stay_unmigrated() {
1913 let backend = StorageBackend::memory().unwrap();
1914 backend.notes().unwrap();
1915 backend.events().unwrap();
1916 {
1917 let reader = backend.pool.reader().unwrap();
1918 let note_key_columns: u32 = reader
1919 .conn()
1920 .query_row(
1921 "SELECT count(*) FROM pragma_table_xinfo('notes') WHERE name = 'key'",
1922 [],
1923 |row| row.get(0),
1924 )
1925 .unwrap();
1926 assert_eq!(note_key_columns, 1);
1927 assert_eq!(
1928 crate::migrations::read_schema_version(reader.conn()).unwrap(),
1929 0
1930 );
1931 }
1932
1933 let store = backend.vectors("after_notes", "after_notes", 3).unwrap();
1934 let id = uuid::Uuid::new_v4();
1935 store
1936 .insert(
1937 id,
1938 khive_types::SubstrateKind::Entity,
1939 "local",
1940 "content",
1941 vec![vec![1.0, 0.0, 0.0]],
1942 )
1943 .await
1944 .unwrap();
1945 let provenance = store.provenance(id).await.unwrap().unwrap();
1946 assert_eq!(provenance.embedding_model, "after_notes");
1947 assert_eq!(provenance.text_fingerprint, None);
1948 assert_eq!(provenance.updated_at, None);
1949
1950 let reader = backend.pool.reader().unwrap();
1951 assert!(!sqlite_table_exists(reader.conn(), "vector_provenance").unwrap());
1952 assert!(!sqlite_table_exists(reader.conn(), "_schema_migrations").unwrap());
1953 }
1954
1955 #[tokio::test]
1956 #[cfg(feature = "vectors")]
1957 async fn vectors_creates_table_idempotently() {
1958 let backend = StorageBackend::memory().unwrap();
1959
1960 let store1 = backend.vectors("idempotent", "idempotent", 3).unwrap();
1961 let store2 = backend.vectors("idempotent", "idempotent", 3).unwrap();
1962
1963 let id = uuid::Uuid::new_v4();
1964 store1
1965 .insert(
1966 id,
1967 khive_types::SubstrateKind::Entity,
1968 "local",
1969 "content",
1970 vec![vec![1.0, 0.0, 0.0]],
1971 )
1972 .await
1973 .unwrap();
1974
1975 let count = store2.count().await.unwrap();
1976 assert_eq!(count, 1);
1977 }
1978
1979 #[cfg(feature = "vectors")]
1982 fn legacy_vector_table_ddl(model_key: &str) -> String {
1983 format!(
1984 "CREATE VIRTUAL TABLE vec_{model_key} USING vec0(\
1985 subject_id TEXT PRIMARY KEY, namespace TEXT NOT NULL, kind TEXT NOT NULL, \
1986 embedding float[3] distance_metric=cosine)"
1987 )
1988 }
1989
1990 #[cfg(feature = "vectors")]
1991 #[test]
1992 fn repeated_vector_store_fetches_take_the_writer_once_per_model() {
1993 let dir = tempfile::tempdir().unwrap();
1994 let backend = StorageBackend::sqlite_for_test(dir.path().join("vector_once.db")).unwrap();
1995
1996 let before = backend.pool.writer_acquisition_snapshot();
1997 for namespace in ["local", "tenant_a", "tenant_b", "local"] {
1998 backend
1999 .vectors_for_namespace("fetched_once", "fetched-once", 3, namespace)
2000 .expect("vector store");
2001 }
2002 let after = backend.pool.writer_acquisition_snapshot();
2003 assert_eq!(
2004 after.pooled_acquisitions - before.pooled_acquisitions,
2005 1,
2006 "only the first fetch of a model may check out the writer"
2007 );
2008 assert_eq!(
2009 after.writer_task_acquisitions,
2010 before.writer_task_acquisitions
2011 );
2012 assert_eq!(
2013 after.standalone_acquisitions,
2014 before.standalone_acquisitions
2015 );
2016
2017 {
2019 let reader = backend.pool.reader().unwrap();
2020 assert!(sqlite_table_exists(reader.conn(), "vec_fetched_once").unwrap());
2021 assert!(sqlite_table_exists(reader.conn(), "_embedding_models").unwrap());
2022 assert!(sqlite_table_exists(reader.conn(), "ann_write_log").unwrap());
2023 }
2024
2025 let before = backend.pool.writer_acquisition_snapshot();
2027 backend
2028 .vectors("fetched_second", "fetched-second", 3)
2029 .expect("second model");
2030 backend
2031 .vectors("fetched_second", "fetched-second", 3)
2032 .expect("second model again");
2033 let after = backend.pool.writer_acquisition_snapshot();
2034 assert_eq!(
2035 after.pooled_acquisitions - before.pooled_acquisitions,
2036 1,
2037 "a second model is prepared once, independently of the first"
2038 );
2039 let reader = backend.pool.reader().unwrap();
2040 assert!(sqlite_table_exists(reader.conn(), "vec_fetched_second").unwrap());
2041 }
2042
2043 #[cfg(feature = "vectors")]
2044 #[test]
2045 fn ensure_vector_tables_prepares_only_models_that_are_not_ready() {
2046 let dir = tempfile::tempdir().unwrap();
2047 let backend = StorageBackend::sqlite_for_test(dir.path().join("vector_batch.db")).unwrap();
2048 let checkouts = || {
2049 backend
2050 .pool
2051 .writer_acquisition_snapshot()
2052 .pooled_acquisitions
2053 };
2054
2055 let start = checkouts();
2056 backend
2057 .ensure_vector_tables(&[("batch_a", 3), ("batch_b", 3)])
2058 .expect("cold batch");
2059 assert_eq!(checkouts() - start, 1, "a cold batch shares one checkout");
2060
2061 backend
2062 .ensure_vector_tables(&[("batch_b", 3), ("batch_a", 3)])
2063 .expect("ready batch");
2064 assert_eq!(checkouts() - start, 1, "a ready batch takes no writer");
2065
2066 backend
2067 .ensure_vector_tables(&[("batch_a", 3), ("batch_c", 3)])
2068 .expect("partly cold batch");
2069 assert_eq!(
2070 checkouts() - start,
2071 2,
2072 "a batch with one cold model takes the writer once"
2073 );
2074 let reader = backend.pool.reader().unwrap();
2075 assert!(sqlite_table_exists(reader.conn(), "vec_batch_c").unwrap());
2076 }
2077
2078 #[cfg(feature = "vectors")]
2079 #[test]
2080 fn warm_vector_store_fetch_finishes_while_the_pool_writer_is_held() {
2081 let dir = tempfile::tempdir().unwrap();
2082 let backend =
2083 Arc::new(StorageBackend::sqlite_for_test(dir.path().join("vector_warm.db")).unwrap());
2084 backend
2085 .vectors("held_writer", "held-writer", 3)
2086 .expect("cold fetch");
2087
2088 let writer = backend.pool.try_writer().unwrap();
2089 let before = backend.pool.writer_acquisition_snapshot();
2090 let worker_backend = Arc::clone(&backend);
2091 let (finished, result) = std::sync::mpsc::channel();
2092 let worker = std::thread::spawn(move || {
2093 let fetched = worker_backend
2094 .vectors_for_namespace("held_writer", "held-writer", 3, "tenant_a")
2095 .map(|_| ());
2096 finished.send(fetched).unwrap();
2097 });
2098 let while_held = result.recv_timeout(std::time::Duration::from_secs(2));
2099 drop(writer);
2101 worker.join().unwrap();
2102 while_held
2103 .expect("a warm vector store fetch must return while the pool writer is held")
2104 .expect("warm vector store fetch");
2105 assert_eq!(
2106 backend
2107 .pool
2108 .writer_acquisition_snapshot()
2109 .pooled_acquisitions,
2110 before.pooled_acquisitions
2111 );
2112 }
2113
2114 #[cfg(feature = "vectors")]
2115 #[test]
2116 fn failed_vector_table_check_is_retried_and_not_recorded_as_ready() {
2117 let dir = tempfile::tempdir().unwrap();
2118 let backend = StorageBackend::sqlite_for_test(dir.path().join("vector_retry.db")).unwrap();
2119 backend
2120 .pool
2121 .try_writer()
2122 .unwrap()
2123 .conn()
2124 .execute_batch(&legacy_vector_table_ddl("retried"))
2125 .unwrap();
2126
2127 let before = backend.pool.writer_acquisition_snapshot();
2128 for _ in 0..2 {
2129 let error = match backend.vectors("retried", "retried", 3) {
2130 Ok(_) => panic!("a vec0 table without the required columns must be rejected"),
2131 Err(error) => error,
2132 };
2133 assert!(
2134 matches!(&error, SqliteError::InvalidData(message)
2135 if message.contains("vec_retried")
2136 && message.contains("missing required column")),
2137 "unexpected error: {error}"
2138 );
2139 }
2140 let after = backend.pool.writer_acquisition_snapshot();
2141 assert_eq!(
2142 after.pooled_acquisitions - before.pooled_acquisitions,
2143 2,
2144 "a failed check must be repeated by the next fetch"
2145 );
2146
2147 backend
2149 .pool
2150 .try_writer()
2151 .unwrap()
2152 .conn()
2153 .execute_batch("DROP TABLE vec_retried")
2154 .unwrap();
2155 backend
2156 .vectors("retried", "retried", 3)
2157 .expect("fetch after the legacy table is gone");
2158 let before = backend.pool.writer_acquisition_snapshot();
2159 backend
2160 .vectors("retried", "retried", 3)
2161 .expect("warm fetch after the check passed");
2162 assert_eq!(
2163 backend
2164 .pool
2165 .writer_acquisition_snapshot()
2166 .pooled_acquisitions,
2167 before.pooled_acquisitions
2168 );
2169 }
2170
2171 #[cfg(feature = "vectors")]
2172 #[test]
2173 fn vector_table_created_after_open_is_validated_before_use() {
2174 let dir = tempfile::tempdir().unwrap();
2175 let path = dir.path().join("vector_late.db");
2176 let backend = StorageBackend::sqlite_for_test(&path).unwrap();
2177 backend
2178 .vectors("known_model", "known-model", 3)
2179 .expect("prepare one model");
2180
2181 rusqlite::Connection::open(&path)
2183 .unwrap()
2184 .execute_batch(&legacy_vector_table_ddl("late_model"))
2185 .unwrap();
2186
2187 let error = match backend.vectors("late_model", "late-model", 3) {
2188 Ok(_) => panic!("a late legacy vec0 table must be validated before use"),
2189 Err(error) => error,
2190 };
2191 assert!(
2192 matches!(&error, SqliteError::InvalidData(message)
2193 if message.contains("vec_late_model")
2194 && message.contains("missing required column")),
2195 "unexpected error: {error}"
2196 );
2197 }
2198
2199 #[tokio::test]
2200 async fn text_roundtrip_via_public_api() {
2201 let backend = StorageBackend::memory().unwrap();
2202 let store = backend.text("test_api").unwrap();
2203
2204 let id = uuid::Uuid::new_v4();
2205 let doc = khive_storage::types::TextDocument {
2206 subject_id: id,
2207 kind: khive_types::SubstrateKind::Entity,
2208 record_kind: None,
2209 title: Some("Test Title".to_string()),
2210 body: "This is a searchable document about Rust.".to_string(),
2211 tags: vec!["rust".to_string()],
2212 namespace: "test_ns".to_string(),
2213 metadata: None,
2214 updated_at: chrono::Utc::now(),
2215 };
2216 store.upsert_document(doc).await.unwrap();
2217
2218 let hits = store
2219 .search(khive_storage::types::TextSearchRequest {
2220 query: "Rust".to_string(),
2221 mode: khive_storage::types::TextQueryMode::Plain,
2222 filter: Some(khive_storage::types::TextFilter {
2223 namespaces: vec!["test_ns".to_string()],
2224 ..Default::default()
2225 }),
2226 top_k: 1,
2227 snippet_chars: 64,
2228 })
2229 .await
2230 .unwrap();
2231
2232 assert_eq!(hits.len(), 1);
2233 assert_eq!(hits[0].subject_id, id);
2234 assert!(hits[0].score.to_f64() > 0.0);
2235 }
2236
2237 #[tokio::test]
2238 async fn text_creates_table_idempotently() {
2239 let backend = StorageBackend::memory().unwrap();
2240
2241 let store1 = backend.text("idempotent_fts").unwrap();
2242 let store2 = backend.text("idempotent_fts").unwrap();
2243
2244 let id = uuid::Uuid::new_v4();
2245 let doc = khive_storage::types::TextDocument {
2246 subject_id: id,
2247 kind: khive_types::SubstrateKind::Note,
2248 record_kind: None,
2249 title: None,
2250 body: "Hello world.".to_string(),
2251 tags: vec![],
2252 namespace: "test_ns".to_string(),
2253 metadata: None,
2254 updated_at: chrono::Utc::now(),
2255 };
2256 store1.upsert_document(doc).await.unwrap();
2257
2258 let count = store2
2259 .count(khive_storage::types::TextFilter {
2260 namespaces: vec!["test_ns".to_string()],
2261 ..Default::default()
2262 })
2263 .await
2264 .unwrap();
2265 assert_eq!(count, 1);
2266 }
2267
2268 #[tokio::test]
2281 async fn text_repeated_open_after_backfill_does_not_scale_with_row_count() {
2282 use std::sync::atomic::{AtomicU64, Ordering};
2283 use std::sync::Arc;
2284
2285 fn repeated_open_work(backend: &StorageBackend) -> u64 {
2286 let work = Arc::new(AtomicU64::new(0));
2287 let counted = Arc::clone(&work);
2288 {
2289 let writer = backend.pool().writer().unwrap();
2290 writer
2291 .conn()
2292 .progress_handler(
2293 1,
2294 Some(move || {
2295 counted.fetch_add(1, Ordering::Relaxed);
2296 false
2297 }),
2298 )
2299 .unwrap();
2300 }
2301
2302 let result = (0..500).try_for_each(|_| backend.text("hot_path_reopen").map(|_| ()));
2306 backend
2307 .pool()
2308 .writer()
2309 .unwrap()
2310 .conn()
2311 .progress_handler(0, None::<fn() -> bool>)
2312 .unwrap();
2313 result.expect("repeated text-store opens must succeed");
2314 work.load(Ordering::Relaxed)
2315 }
2316
2317 let backend = StorageBackend::memory().unwrap();
2318 let store = backend.text("hot_path_reopen").unwrap();
2319 let body = "the quick brown fox jumps over the lazy dog ".repeat(35);
2320 let mut seeded = 0;
2321 let mut work = Vec::new();
2322 for target_rows in [100, 5_000] {
2323 for _ in seeded..target_rows {
2324 let doc = khive_storage::types::TextDocument {
2325 subject_id: uuid::Uuid::new_v4(),
2326 kind: khive_types::SubstrateKind::Note,
2327 record_kind: Some("memory".to_string()),
2328 title: None,
2329 body: body.clone(),
2330 tags: vec![],
2331 namespace: "test_ns".to_string(),
2332 metadata: None,
2333 updated_at: chrono::Utc::now(),
2334 };
2335 store.upsert_document(doc).await.unwrap();
2336 }
2337 seeded = target_rows;
2338 assert_eq!(
2339 store
2340 .count(khive_storage::types::TextFilter {
2341 namespaces: vec!["test_ns".to_string()],
2342 ..Default::default()
2343 })
2344 .await
2345 .unwrap(),
2346 target_rows,
2347 "the work comparison requires both declared row populations"
2348 );
2349
2350 let _ = backend.text("hot_path_reopen").unwrap();
2353 work.push(repeated_open_work(&backend));
2354 }
2355
2356 let [small, large] = [work[0], work[1]];
2357 assert!(small > 0 && large > 0, "both work meters must be active");
2358 assert!(
2359 large <= small * 2,
2360 "500 repeated backend.text() calls used {small} SQLite VM progress units at \
2361 100 rows and {large} at 5,000 rows; growing the table 50-fold must not \
2362 more than double already-backfilled work (for example via COUNT(*))"
2363 );
2364 }
2365
2366 #[tokio::test]
2372 async fn text_open_after_legacy_seed_backfills_the_map_with_full_parity() {
2373 let backend = StorageBackend::memory().unwrap();
2374 let table_key = "legacy_seed_parity";
2375 let table = format!("fts_{table_key}");
2376 let map = format!("{table}_rowids");
2377
2378 let _ = backend.text(table_key).unwrap();
2381
2382 {
2386 let writer = backend.pool().writer().unwrap();
2387 writer.conn().execute_batch("BEGIN").unwrap();
2388 {
2389 let mut insert = writer
2390 .conn()
2391 .prepare(&format!(
2392 "INSERT INTO {table} \
2393 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, \
2394 record_kind) \
2395 VALUES (?1, 'note', '', 'legacy body', '[]', 'test_ns', NULL, 0, 'memory')"
2396 ))
2397 .unwrap();
2398 for i in 0..500 {
2399 insert
2400 .execute(rusqlite::params![format!("legacy-{i}")])
2401 .unwrap();
2402 }
2403 }
2404 writer.conn().execute_batch("COMMIT").unwrap();
2405 }
2406 {
2407 let writer = backend.pool().writer().unwrap();
2408 let map_count: i64 = writer
2409 .conn()
2410 .query_row(&format!("SELECT COUNT(*) FROM {map}"), [], |row| row.get(0))
2411 .unwrap();
2412 assert_eq!(
2413 map_count, 0,
2414 "the raw-SQL seed must bypass the map, reproducing a genuinely pre-migration db"
2415 );
2416 }
2417
2418 let _ = backend.text(table_key).unwrap();
2421
2422 {
2423 let writer = backend.pool().writer().unwrap();
2424 let mismatched: i64 = writer
2425 .conn()
2426 .query_row(
2427 &format!(
2428 "SELECT \
2429 (SELECT COUNT(*) FROM {table} WHERE rowid NOT IN (SELECT rowid FROM {map})) + \
2430 (SELECT COUNT(*) FROM {map} WHERE rowid NOT IN (SELECT rowid FROM {table}))"
2431 ),
2432 [],
2433 |row| row.get(0),
2434 )
2435 .unwrap();
2436 assert_eq!(
2437 mismatched, 0,
2438 "backfill must give every FTS row exactly one map entry, both directions"
2439 );
2440 let fts_count: i64 = writer
2441 .conn()
2442 .query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| {
2443 row.get(0)
2444 })
2445 .unwrap();
2446 let map_count: i64 = writer
2447 .conn()
2448 .query_row(&format!("SELECT COUNT(*) FROM {map}"), [], |row| row.get(0))
2449 .unwrap();
2450 assert_eq!(fts_count, 500);
2451 assert_eq!(map_count, 500);
2452 }
2453 }
2454
2455 #[tokio::test]
2464 async fn text_open_reconciles_a_partial_map_instead_of_treating_it_as_complete() {
2465 let backend = StorageBackend::memory().unwrap();
2466 let table_key = "partial_map_reconcile";
2467 let table = format!("fts_{table_key}");
2468 let map = format!("{table}_rowids");
2469 let state = format!("{map}_state");
2470
2471 let _ = backend.text(table_key).unwrap();
2474
2475 let a = uuid::Uuid::new_v4();
2476 let b = uuid::Uuid::new_v4();
2477 {
2478 let writer = backend.pool().writer().unwrap();
2479 writer.conn().execute_batch("BEGIN").unwrap();
2480 writer
2481 .conn()
2482 .execute(
2483 &format!(
2484 "INSERT INTO {table} \
2485 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2486 updated_at, record_kind) \
2487 VALUES (1, ?1, 'note', '', 'doc a', '[]', 'test_ns', NULL, 0, 'memory')"
2488 ),
2489 rusqlite::params![a.to_string()],
2490 )
2491 .expect("insert A's fts row");
2492 writer
2493 .conn()
2494 .execute(
2495 &format!(
2496 "INSERT INTO {table} \
2497 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2498 updated_at, record_kind) \
2499 VALUES (2, ?1, 'note', '', 'doc b', '[]', 'test_ns', NULL, 0, 'memory')"
2500 ),
2501 rusqlite::params![b.to_string()],
2502 )
2503 .expect("insert B's fts row");
2504 writer
2506 .conn()
2507 .execute(
2508 &format!(
2509 "INSERT INTO {map} (namespace, subject_id, rowid) VALUES ('test_ns', ?1, 2)"
2510 ),
2511 rusqlite::params![b.to_string()],
2512 )
2513 .expect("insert B's own map entry, leaving A's missing");
2514 writer
2518 .conn()
2519 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2520 .expect("clear the completion marker");
2521 writer.conn().execute_batch("COMMIT").unwrap();
2522 }
2523
2524 let store = backend.text(table_key).unwrap();
2525
2526 let a_mapped: i64 = {
2527 let writer = backend.pool().writer().unwrap();
2528 writer
2529 .conn()
2530 .query_row(
2531 &format!("SELECT COUNT(*) FROM {map} WHERE namespace = 'test_ns' AND subject_id = ?1"),
2532 rusqlite::params![a.to_string()],
2533 |row| row.get(0),
2534 )
2535 .unwrap()
2536 };
2537 assert_eq!(
2538 a_mapped, 1,
2539 "the partial map must be reconciled, not left missing A's entry"
2540 );
2541
2542 let fetched_a = store.get_document("test_ns", a).await.unwrap();
2543 assert!(
2544 fetched_a.is_some(),
2545 "get_document(A) must work once the partial map is reconciled"
2546 );
2547 }
2548
2549 #[tokio::test]
2561 async fn text_open_removes_a_wrong_key_map_row_before_backfilling_the_right_one() {
2562 let backend = StorageBackend::memory().unwrap();
2563 let table_key = "wrong_key_map_row";
2564 let table = format!("fts_{table_key}");
2565 let map = format!("{table}_rowids");
2566 let state = format!("{map}_state");
2567
2568 let _ = backend.text(table_key).unwrap();
2569
2570 let a = uuid::Uuid::new_v4();
2571 let b = uuid::Uuid::new_v4();
2572 {
2573 let writer = backend.pool().writer().unwrap();
2574 writer.conn().execute_batch("BEGIN").unwrap();
2575 writer
2576 .conn()
2577 .execute(
2578 &format!(
2579 "INSERT INTO {table} \
2580 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2581 updated_at, record_kind) \
2582 VALUES (7, ?1, 'note', '', 'doc b', '[]', 'test_ns', NULL, 0, 'memory')"
2583 ),
2584 rusqlite::params![b.to_string()],
2585 )
2586 .expect("insert B's live fts row at rowid 7");
2587 writer
2588 .conn()
2589 .execute(
2590 &format!(
2591 "INSERT INTO {map} (namespace, subject_id, rowid) VALUES ('test_ns', ?1, 7)"
2592 ),
2593 rusqlite::params![a.to_string()],
2594 )
2595 .expect("insert A's stale map row still pointing at rowid 7");
2596 writer
2597 .conn()
2598 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2599 .expect("clear the completion marker");
2600 writer.conn().execute_batch("COMMIT").unwrap();
2601 }
2602
2603 let store = backend.text(table_key).unwrap();
2604
2605 let map_rows: Vec<(String, i64)> = {
2606 let writer = backend.pool().writer().unwrap();
2607 let mut stmt = writer
2608 .conn()
2609 .prepare(&format!(
2610 "SELECT subject_id, rowid FROM {map} ORDER BY rowid"
2611 ))
2612 .unwrap();
2613 let rows = stmt
2614 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
2615 .unwrap()
2616 .collect::<Result<Vec<_>, _>>()
2617 .unwrap();
2618 rows
2619 };
2620 assert_eq!(
2621 map_rows,
2622 vec![(b.to_string(), 7)],
2623 "the stale (A, 7) map row must be removed and replaced by the correct (B, 7) row, \
2624 not left alongside it"
2625 );
2626
2627 assert!(
2628 store.get_document("test_ns", a).await.unwrap().is_none(),
2629 "A's stale map entry is gone, so get_document(A) must find nothing"
2630 );
2631 let fetched_b = store.get_document("test_ns", b).await.unwrap();
2632 assert!(
2633 fetched_b.is_some(),
2634 "get_document(B) must find the live row now correctly mapped"
2635 );
2636 assert_eq!(fetched_b.unwrap().body, "doc b");
2637 }
2638
2639 #[tokio::test]
2646 async fn text_open_sweeps_the_duplicate_that_lost_the_survivor_race() {
2647 let backend = StorageBackend::memory().unwrap();
2648 let table_key = "duplicate_loser_swept";
2649 let table = format!("fts_{table_key}");
2650 let map = format!("{table}_rowids");
2651 let state = format!("{map}_state");
2652
2653 let _ = backend.text(table_key).unwrap();
2654
2655 let dup = uuid::Uuid::new_v4();
2656 {
2657 let writer = backend.pool().writer().unwrap();
2658 writer.conn().execute_batch("BEGIN").unwrap();
2659 writer
2661 .conn()
2662 .execute(
2663 &format!(
2664 "INSERT INTO {table} \
2665 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2666 updated_at, record_kind) \
2667 VALUES (10, ?1, 'note', '', 'older body', '[]', 'test_ns', NULL, 1, \
2668 'memory')"
2669 ),
2670 rusqlite::params![dup.to_string()],
2671 )
2672 .expect("insert the older/losing duplicate");
2673 writer
2675 .conn()
2676 .execute(
2677 &format!(
2678 "INSERT INTO {table} \
2679 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2680 updated_at, record_kind) \
2681 VALUES (20, ?1, 'note', '', 'newer body', '[]', 'test_ns', NULL, 5, \
2682 'memory')"
2683 ),
2684 rusqlite::params![dup.to_string()],
2685 )
2686 .expect("insert the newer/surviving duplicate");
2687 writer
2688 .conn()
2689 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2690 .expect("clear the completion marker");
2691 writer.conn().execute_batch("COMMIT").unwrap();
2692 }
2693
2694 let _ = backend.text(table_key).unwrap();
2695
2696 let writer = backend.pool().writer().unwrap();
2697 let map_rows: Vec<i64> = writer
2698 .conn()
2699 .prepare(&format!(
2700 "SELECT rowid FROM {map} WHERE namespace = 'test_ns' AND subject_id = ?1"
2701 ))
2702 .unwrap()
2703 .query_map(rusqlite::params![dup.to_string()], |row| row.get(0))
2704 .unwrap()
2705 .collect::<Result<Vec<_>, _>>()
2706 .unwrap();
2707 assert_eq!(
2708 map_rows,
2709 vec![20],
2710 "exactly one map row must survive, at the newer (by updated_at) rowid"
2711 );
2712
2713 let fts_rowids: Vec<i64> = writer
2714 .conn()
2715 .prepare(&format!("SELECT rowid FROM {table} ORDER BY rowid"))
2716 .unwrap()
2717 .query_map([], |row| row.get(0))
2718 .unwrap()
2719 .collect::<Result<Vec<_>, _>>()
2720 .unwrap();
2721 assert_eq!(
2722 fts_rowids,
2723 vec![20],
2724 "the losing duplicate (rowid 10) must be deleted from the FTS table itself, not just \
2725 left out of the map as an unmapped live row"
2726 );
2727 }
2728
2729 #[tokio::test]
2736 async fn text_open_after_marker_written_does_not_rescan_even_a_corrupted_map() {
2737 let backend = StorageBackend::memory().unwrap();
2738 let table_key = "marker_no_rescan";
2739 let table = format!("fts_{table_key}");
2740 let map = format!("{table}_rowids");
2741 let state = format!("{map}_state");
2742
2743 let store = backend.text(table_key).unwrap();
2744 store
2745 .upsert_document(khive_storage::types::TextDocument {
2746 subject_id: uuid::Uuid::new_v4(),
2747 kind: khive_types::SubstrateKind::Note,
2748 record_kind: Some("memory".to_string()),
2749 title: None,
2750 body: "seed".to_string(),
2751 tags: vec![],
2752 namespace: "test_ns".to_string(),
2753 metadata: None,
2754 updated_at: chrono::Utc::now(),
2755 })
2756 .await
2757 .unwrap();
2758
2759 let _ = backend.text(table_key).unwrap();
2766
2767 let marked: bool = {
2768 let writer = backend.pool().writer().unwrap();
2769 writer
2770 .conn()
2771 .query_row(
2772 &format!(
2773 "SELECT EXISTS(SELECT 1 FROM {state} WHERE key = 'backfill' AND value = 'complete')"
2774 ),
2775 [],
2776 |row| row.get(0),
2777 )
2778 .unwrap()
2779 };
2780 assert!(
2781 marked,
2782 "a completion marker must exist once the table has held a row"
2783 );
2784
2785 {
2786 let writer = backend.pool().writer().unwrap();
2787 writer
2788 .conn()
2789 .execute(&format!("DELETE FROM {map}"), [])
2790 .expect("corrupt the map by deleting its row directly");
2791 }
2792
2793 let _ = backend.text(table_key).unwrap();
2794
2795 let map_count: i64 = {
2796 let writer = backend.pool().writer().unwrap();
2797 writer
2798 .conn()
2799 .query_row(&format!("SELECT COUNT(*) FROM {map}"), [], |row| row.get(0))
2800 .unwrap()
2801 };
2802 assert_eq!(
2803 map_count, 0,
2804 "a marker-complete table must not be re-scanned on open, even to reconcile a map \
2805 an external actor emptied out from under it"
2806 );
2807 }
2808
2809 #[tokio::test]
2813 async fn text_open_writable_legacy_backfill_excludes_null_key_rows() {
2814 let backend = StorageBackend::memory().unwrap();
2815 let table_key = "legacy_null_key";
2816 let table = format!("fts_{table_key}");
2817 let map = format!("{table}_rowids");
2818 let state = format!("{map}_state");
2819
2820 let _ = backend.text(table_key).unwrap();
2821 {
2822 let writer = backend.pool().writer().unwrap();
2823 writer.conn().execute_batch("BEGIN").unwrap();
2824 writer
2825 .conn()
2826 .execute(
2827 &format!(
2828 "INSERT INTO {table} \
2829 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2830 updated_at, record_kind) \
2831 VALUES (1, NULL, 'note', '', 'null-key body', '[]', NULL, NULL, 0, '')"
2832 ),
2833 [],
2834 )
2835 .expect("insert legacy NULL-key fts row");
2836 writer
2837 .conn()
2838 .execute(
2839 &format!(
2840 "INSERT INTO {table} \
2841 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2842 updated_at, record_kind) \
2843 VALUES (2, 'legacy-1', 'note', '', 'normal body', '[]', 'test_ns', NULL, \
2844 0, 'memory')"
2845 ),
2846 [],
2847 )
2848 .expect("insert legacy normal-key fts row");
2849 writer
2850 .conn()
2851 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2852 .expect("clear the completion marker written for the then-empty table");
2853 writer.conn().execute_batch("COMMIT").unwrap();
2854 }
2855
2856 let _ = backend.text(table_key).unwrap();
2858
2859 let writer = backend.pool().writer().unwrap();
2860 let map_count: i64 = writer
2861 .conn()
2862 .query_row(&format!("SELECT COUNT(*) FROM {map}"), [], |row| row.get(0))
2863 .unwrap();
2864 assert_eq!(map_count, 1, "only the non-NULL-key row may be mapped");
2865 let mapped_subject: String = writer
2866 .conn()
2867 .query_row(&format!("SELECT subject_id FROM {map}"), [], |row| {
2868 row.get(0)
2869 })
2870 .unwrap();
2871 assert_eq!(mapped_subject, "legacy-1");
2872 }
2873
2874 #[tokio::test]
2879 async fn text_open_writable_legacy_backfill_survivor_is_chosen_by_updated_at_not_rowid() {
2880 let backend = StorageBackend::memory().unwrap();
2881 let table_key = "legacy_updated_at_survivor";
2882 let table = format!("fts_{table_key}");
2883 let map = format!("{table}_rowids");
2884 let state = format!("{map}_state");
2885
2886 let _ = backend.text(table_key).unwrap();
2887 {
2888 let writer = backend.pool().writer().unwrap();
2889 writer.conn().execute_batch("BEGIN").unwrap();
2890 writer
2892 .conn()
2893 .execute(
2894 &format!(
2895 "INSERT INTO {table} \
2896 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2897 updated_at, record_kind) \
2898 VALUES (100, 'dup', 'note', '', 'newer body', '[]', 'test_ns', NULL, \
2899 500, 'memory')"
2900 ),
2901 [],
2902 )
2903 .expect("insert newer-but-lower-rowid fts row");
2904 writer
2906 .conn()
2907 .execute(
2908 &format!(
2909 "INSERT INTO {table} \
2910 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2911 updated_at, record_kind) \
2912 VALUES (200, 'dup', 'note', '', 'older body', '[]', 'test_ns', NULL, \
2913 100, 'memory')"
2914 ),
2915 [],
2916 )
2917 .expect("insert older-but-higher-rowid fts row");
2918 writer
2919 .conn()
2920 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2921 .expect("clear the completion marker written for the then-empty table");
2922 writer.conn().execute_batch("COMMIT").unwrap();
2923 }
2924
2925 let _ = backend.text(table_key).unwrap();
2926
2927 let writer = backend.pool().writer().unwrap();
2928 let mapped_rowid: i64 = writer
2929 .conn()
2930 .query_row(
2931 &format!(
2932 "SELECT rowid FROM {map} WHERE namespace = 'test_ns' AND subject_id = 'dup'"
2933 ),
2934 [],
2935 |row| row.get(0),
2936 )
2937 .expect("read dup's map entry");
2938 assert_eq!(
2939 mapped_rowid, 100,
2940 "the newer document (by updated_at) must survive even though its rowid is lower"
2941 );
2942 }
2943
2944 include!("backend/key_validation_tests.rs");
2945
2946 #[tokio::test]
2947 async fn sqlite_read_only_graph_store_rejects_upsert_edge() {
2948 use khive_storage::types::Edge;
2949 use khive_types::EdgeRelation;
2950
2951 let dir = tempfile::tempdir().unwrap();
2952 let path = dir.path().join("ro_graph.db");
2953
2954 {
2956 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
2957 writable.graph().unwrap();
2958 }
2959 #[cfg(unix)]
2960 freeze_snapshot_sidecars(&path);
2961
2962 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
2963 let store = match ro.graph() {
2964 Ok(store) => store,
2965 Err(_) => return,
2968 };
2969
2970 let now = chrono::Utc::now();
2971 let edge = Edge {
2972 id: uuid::Uuid::new_v4().into(),
2973 namespace: "local".to_string(),
2974 source_id: uuid::Uuid::new_v4(),
2975 target_id: uuid::Uuid::new_v4(),
2976 relation: EdgeRelation::Extends,
2977 weight: 0.8,
2978 created_at: now,
2979 updated_at: now,
2980 deleted_at: None,
2981 metadata: None,
2982 target_backend: None,
2983 };
2984
2985 let result = store.upsert_edge(edge).await;
2986 assert!(
2987 result.is_err(),
2988 "upsert_edge on a read-only backend must reject, not silently no-op"
2989 );
2990 }
2991
2992 #[tokio::test]
2993 async fn sqlite_read_only_event_store_rejects_append_event() {
2994 use khive_types::{EventKind, EventOutcome, SubstrateKind};
2995
2996 let dir = tempfile::tempdir().unwrap();
2997 let path = dir.path().join("ro_events.db");
2998
2999 {
3000 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
3001 writable.events().unwrap();
3002 }
3003 #[cfg(unix)]
3004 freeze_snapshot_sidecars(&path);
3005
3006 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
3007 let store = match ro.events() {
3008 Ok(store) => store,
3009 Err(_) => return,
3010 };
3011
3012 let event = khive_storage::event::Event::new(
3013 "local",
3014 "test.verb",
3015 EventKind::Audit,
3016 SubstrateKind::Entity,
3017 "test-actor",
3018 )
3019 .with_outcome(EventOutcome::Success);
3020
3021 let result = store.append_event(event).await;
3022 assert!(
3023 result.is_err(),
3024 "append_event on a read-only backend must reject, not silently no-op"
3025 );
3026 }
3027
3028 #[tokio::test]
3029 async fn sqlite_read_only_text_store_rejects_upsert_document() {
3030 use khive_storage::types::TextDocument;
3031 use khive_types::SubstrateKind;
3032
3033 let dir = tempfile::tempdir().unwrap();
3034 let path = dir.path().join("ro_text.db");
3035
3036 {
3037 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
3038 writable.text("ro_test").unwrap();
3039 }
3040 #[cfg(unix)]
3041 freeze_snapshot_sidecars(&path);
3042
3043 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
3044 let store = match ro.text("ro_test") {
3045 Ok(store) => store,
3046 Err(_) => return,
3047 };
3048
3049 let doc = TextDocument {
3050 subject_id: uuid::Uuid::new_v4(),
3051 kind: SubstrateKind::Entity,
3052 record_kind: None,
3053 title: Some("Title".to_string()),
3054 body: "Body text.".to_string(),
3055 tags: vec![],
3056 namespace: "local".to_string(),
3057 metadata: None,
3058 updated_at: chrono::Utc::now(),
3059 };
3060
3061 let result = store.upsert_document(doc).await;
3062 assert!(
3063 result.is_err(),
3064 "upsert_document on a read-only backend must reject, not silently no-op"
3065 );
3066 }
3067
3068 #[tokio::test]
3075 async fn sqlite_read_only_text_store_without_rowid_map_falls_back_to_scan_predicates() {
3076 let dir = tempfile::tempdir().unwrap();
3077 let path = dir.path().join("ro_text_no_map.db");
3078
3079 let id = uuid::Uuid::new_v4();
3080 {
3081 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
3082 let writer = writable.pool().try_writer().unwrap();
3083 writer
3084 .conn()
3085 .execute_batch(
3086 "CREATE VIRTUAL TABLE IF NOT EXISTS fts_ro_no_map USING fts5(\
3087 subject_id UNINDEXED, kind UNINDEXED, title, body, tags UNINDEXED, \
3088 namespace UNINDEXED, metadata UNINDEXED, updated_at UNINDEXED, \
3089 record_kind, tokenize = 'trigram')",
3090 )
3091 .unwrap();
3092 writer
3093 .conn()
3094 .execute(
3095 "INSERT INTO fts_ro_no_map \
3096 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, \
3097 record_kind) \
3098 VALUES (?1, 'note', '', 'legacy body', '[]', 'local', NULL, 0, NULL)",
3099 rusqlite::params![id.to_string()],
3100 )
3101 .unwrap();
3102 }
3103 #[cfg(unix)]
3104 freeze_snapshot_sidecars(&path);
3105
3106 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
3107 let store = ro
3108 .text("ro_no_map")
3109 .expect("a read-only FTS table with no sidecar map must still open successfully");
3110
3111 let fetched = store
3112 .get_document("local", id)
3113 .await
3114 .expect("scan-fallback get_document must not error against a missing map table");
3115 assert!(
3116 fetched.is_some(),
3117 "scan-fallback get_document must still find the legacy row"
3118 );
3119 assert_eq!(fetched.unwrap().subject_id, id);
3120 }
3121
3122 #[tokio::test]
3132 async fn sqlite_read_only_text_store_with_unmarked_map_falls_back_to_scan_predicates() {
3133 let dir = tempfile::tempdir().unwrap();
3134 let path = dir.path().join("ro_text_unmarked_map.db");
3135
3136 let a = uuid::Uuid::new_v4();
3137 let b = uuid::Uuid::new_v4();
3138 {
3139 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
3140 let writer = writable.pool().try_writer().unwrap();
3141 writer
3142 .conn()
3143 .execute_batch(
3144 "CREATE VIRTUAL TABLE IF NOT EXISTS fts_ro_unmarked USING fts5(\
3145 subject_id UNINDEXED, kind UNINDEXED, title, body, tags UNINDEXED, \
3146 namespace UNINDEXED, metadata UNINDEXED, updated_at UNINDEXED, \
3147 record_kind, tokenize = 'trigram')",
3148 )
3149 .unwrap();
3150 writer
3151 .conn()
3152 .execute_batch(&text::rowid_map_ddl("fts_ro_unmarked"))
3153 .unwrap();
3154 writer
3155 .conn()
3156 .execute(
3157 "INSERT INTO fts_ro_unmarked \
3158 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
3159 updated_at, record_kind) \
3160 VALUES (1, ?1, 'note', '', 'doc a', '[]', 'local', NULL, 0, NULL)",
3161 rusqlite::params![a.to_string()],
3162 )
3163 .unwrap();
3164 writer
3165 .conn()
3166 .execute(
3167 "INSERT INTO fts_ro_unmarked \
3168 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
3169 updated_at, record_kind) \
3170 VALUES (2, ?1, 'note', '', 'doc b', '[]', 'local', NULL, 0, NULL)",
3171 rusqlite::params![b.to_string()],
3172 )
3173 .unwrap();
3174 writer
3177 .conn()
3178 .execute(
3179 "INSERT INTO fts_ro_unmarked_rowids (namespace, subject_id, rowid) \
3180 VALUES ('local', ?1, 2)",
3181 rusqlite::params![b.to_string()],
3182 )
3183 .unwrap();
3184 }
3185 #[cfg(unix)]
3186 freeze_snapshot_sidecars(&path);
3187
3188 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
3189 let store = ro
3190 .text("ro_unmarked")
3191 .expect("a read-only FTS table with an unmarked map must still open successfully");
3192
3193 let fetched_a = store
3194 .get_document("local", a)
3195 .await
3196 .expect("scan-fallback get_document must not error against an unmarked map");
3197 assert!(
3198 fetched_a.is_some(),
3199 "A has no map entry, so only the scan fallback (not a map join) can find it -- \
3200 proving the unmarked map was not trusted"
3201 );
3202 assert_eq!(fetched_a.unwrap().body, "doc a");
3203 }
3204
3205 #[tokio::test]
3206 async fn blob_store_roundtrip_via_public_api() {
3207 let dir = tempfile::tempdir().unwrap();
3208 let path = dir.path().join("blob_backend.db");
3209 let backend = StorageBackend::sqlite_for_test(&path).unwrap();
3210
3211 let store = backend.blob_store(None, Some(0)).unwrap();
3215 let bytes = b"backend-level blob roundtrip".to_vec();
3216 let content_ref = store.put(bytes.clone()).await.unwrap();
3217 assert_eq!(
3218 store
3219 .get_bounded_verified(&content_ref, bytes.len() as u64)
3220 .await
3221 .unwrap(),
3222 bytes
3223 );
3224 }
3225
3226 #[test]
3227 fn blob_store_defaults_root_beside_db_file() {
3228 let dir = tempfile::tempdir().unwrap();
3229 let path = dir.path().join("blob_default.db");
3230 let backend = StorageBackend::sqlite_for_test(&path).unwrap();
3231
3232 let _store = backend.blob_store(None, None).unwrap();
3236 assert!(
3237 dir.path().join("blobs").is_dir(),
3238 "default root must be created beside the database file"
3239 );
3240 }
3241
3242 #[test]
3243 fn blob_store_errors_for_in_memory_backend_with_no_override() {
3244 let backend = StorageBackend::memory().unwrap();
3245 assert!(backend.blob_store(None, None).is_err());
3246 }
3247
3248 #[test]
3249 fn blob_store_accepts_explicit_root_for_in_memory_backend() {
3250 let dir = tempfile::tempdir().unwrap();
3251 let backend = StorageBackend::memory().unwrap();
3252 let store = backend.blob_store(Some(dir.path()), None);
3253 assert!(store.is_ok());
3254 }
3255
3256 include!("backend/migration_tests.rs");
3257
3258 #[test]
3259 fn pack_ddl_plan_rolls_back_all_statements_on_failure() {
3260 let backend = StorageBackend::memory().unwrap();
3261 let error = backend
3262 .apply_pack_ddl_statements(&[
3263 "CREATE TABLE IF NOT EXISTS pack_schema_first (id INTEGER PRIMARY KEY)",
3264 "CREATE INDEX IF NOT EXISTS pack_schema_second ON pack_schema_missing(id)",
3265 ])
3266 .unwrap_err();
3267
3268 assert!(
3269 error.to_string().contains("pack_schema_missing"),
3270 "schema-plan error must retain the failing SQLite diagnostic: {error}"
3271 );
3272
3273 let reader = backend.pool().reader().unwrap();
3274 let visible_objects: i64 = reader
3275 .conn()
3276 .query_row(
3277 "SELECT COUNT(*) FROM sqlite_master \
3278 WHERE name IN ('pack_schema_first', 'pack_schema_second')",
3279 [],
3280 |row| row.get(0),
3281 )
3282 .unwrap();
3283 assert_eq!(visible_objects, 0);
3284 }
3285
3286 #[test]
3287 fn pack_ddl_plan_applies_all_statements_idempotently() {
3288 const PLAN: &[&str] = &[
3289 "CREATE TABLE IF NOT EXISTS pack_schema_success (id INTEGER PRIMARY KEY, value TEXT)",
3290 "CREATE INDEX IF NOT EXISTS pack_schema_success_value_idx \
3291 ON pack_schema_success(value)",
3292 ];
3293
3294 let backend = StorageBackend::memory().unwrap();
3295 backend.apply_pack_ddl_statements(PLAN).unwrap();
3296 backend.apply_pack_ddl_statements(PLAN).unwrap();
3297
3298 let reader = backend.pool().reader().unwrap();
3299 let visible_objects: i64 = reader
3300 .conn()
3301 .query_row(
3302 "SELECT COUNT(*) FROM sqlite_master \
3303 WHERE name IN ('pack_schema_success', 'pack_schema_success_value_idx')",
3304 [],
3305 |row| row.get(0),
3306 )
3307 .unwrap();
3308 assert_eq!(visible_objects, 2);
3309 }
3310
3311 fn issue_1029_pool(write_queue_enabled: bool) -> (tempfile::TempDir, StorageBackend) {
3320 let dir = tempfile::tempdir().unwrap();
3321 let path = dir.path().join("issue_1029.db");
3322 let config = crate::pool::PoolConfig {
3323 path: Some(path.clone()),
3324 busy_timeout: std::time::Duration::from_millis(200),
3325 write_queue_enabled: Some(write_queue_enabled),
3326 ..crate::pool::PoolConfig::for_test()
3327 };
3328 let pool = ConnectionPool::new(config).expect("fresh tenant-shaped pool should open");
3329 let backend = StorageBackend {
3330 pool: Arc::new(pool),
3331 is_file_backed: true,
3332 path: Some(path),
3333 vector_tables_ready: Default::default(),
3334 notes_seq_repair_runs: AtomicUsize::new(0),
3335 store_schemas: std::array::from_fn(|_| Arc::new(StoreSchemaGate::default())),
3336 };
3337 (dir, backend)
3338 }
3339
3340 async fn issue_1029_create_entity_shaped_sequence(
3341 backend: &StorageBackend,
3342 ) -> Result<(), String> {
3343 let entities = backend
3344 .entities_for_namespace("tenant_ns")
3345 .map_err(|e| format!("entities_for_namespace: {e}"))?;
3346 let entity = khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029Repro");
3347 let entity_id = entity.id;
3348 entities
3349 .upsert_entity(entity)
3350 .await
3351 .map_err(|e| format!("upsert_entity: {e}"))?;
3352
3353 let text = backend.text("entities").map_err(|e| format!("text: {e}"))?;
3354 let doc = khive_storage::types::TextDocument {
3355 subject_id: entity_id,
3356 kind: khive_types::SubstrateKind::Entity,
3357 record_kind: None,
3358 title: Some("Issue1029Repro".to_string()),
3359 body: "issue 1029 repro body".to_string(),
3360 tags: vec![],
3361 namespace: "tenant_ns".to_string(),
3362 metadata: None,
3363 updated_at: chrono::Utc::now(),
3364 };
3365 text.upsert_document(doc)
3366 .await
3367 .map_err(|e| format!("fts_upsert: {e}"))
3368 }
3369
3370 #[tokio::test]
3376 async fn issue_1029_create_entity_shaped_sequence_write_queue_off() {
3377 let (_dir, backend) = issue_1029_pool(false);
3378 let result = issue_1029_create_entity_shaped_sequence(&backend).await;
3379 assert!(
3380 result.is_ok(),
3381 "khive#1029 repro (KHIVE_WRITE_QUEUE off): fts_upsert step failed: {:?}",
3382 result.err()
3383 );
3384 }
3385
3386 #[tokio::test]
3392 async fn issue_1029_create_entity_shaped_sequence_write_queue_on() {
3393 let (_dir, backend) = issue_1029_pool(true);
3394 let result = issue_1029_create_entity_shaped_sequence(&backend).await;
3395 assert!(
3396 result.is_ok(),
3397 "khive#1029 repro (KHIVE_WRITE_QUEUE=1): fts_upsert step failed: {:?}",
3398 result.err()
3399 );
3400 }
3401
3402 #[tokio::test]
3410 async fn issue_1029_two_pools_same_file_write_queue_on() {
3411 let dir = tempfile::tempdir().unwrap();
3412 let path = dir.path().join("issue_1029_two_pools.db");
3413
3414 let cfg = |p: std::path::PathBuf| crate::pool::PoolConfig {
3415 path: Some(p),
3416 busy_timeout: std::time::Duration::from_millis(200),
3417 write_queue_enabled: Some(true),
3418 ..crate::pool::PoolConfig::for_test()
3419 };
3420
3421 let pool_a = ConnectionPool::new(cfg(path.clone())).expect("pool A should open");
3422 let backend_a = StorageBackend {
3423 pool: Arc::new(pool_a),
3424 is_file_backed: true,
3425 path: Some(path.clone()),
3426 vector_tables_ready: Default::default(),
3427 notes_seq_repair_runs: AtomicUsize::new(0),
3428 store_schemas: std::array::from_fn(|_| Arc::new(StoreSchemaGate::default())),
3429 };
3430 let pool_b = ConnectionPool::new(cfg(path.clone())).expect("pool B should open");
3431 let backend_b = StorageBackend {
3432 pool: Arc::new(pool_b),
3433 is_file_backed: true,
3434 path: Some(path),
3435 vector_tables_ready: Default::default(),
3436 notes_seq_repair_runs: AtomicUsize::new(0),
3437 store_schemas: std::array::from_fn(|_| Arc::new(StoreSchemaGate::default())),
3438 };
3439
3440 let entities = backend_a
3441 .entities_for_namespace("tenant_ns")
3442 .expect("entities_for_namespace on pool A");
3443 let entity =
3444 khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029TwoPools");
3445 let entity_id = entity.id;
3446 entities
3447 .upsert_entity(entity)
3448 .await
3449 .expect("pool A entity upsert should succeed");
3450
3451 let text = backend_b.text("entities").expect("text on pool B");
3452 let doc = khive_storage::types::TextDocument {
3453 subject_id: entity_id,
3454 kind: khive_types::SubstrateKind::Entity,
3455 record_kind: None,
3456 title: Some("Issue1029TwoPools".to_string()),
3457 body: "issue 1029 two-pool repro body".to_string(),
3458 tags: vec![],
3459 namespace: "tenant_ns".to_string(),
3460 metadata: None,
3461 updated_at: chrono::Utc::now(),
3462 };
3463 let result = text.upsert_document(doc).await;
3464 assert!(
3465 result.is_ok(),
3466 "khive#1029 two-pool repro: fts_upsert on an independent pool for the \
3467 same tenant DB file failed: {:?}",
3468 result.err()
3469 );
3470 }
3471
3472 struct StarvationCaptureSubscriber {
3475 events: Arc<std::sync::Mutex<Vec<std::collections::BTreeMap<String, String>>>>,
3476 }
3477
3478 impl tracing::Subscriber for StarvationCaptureSubscriber {
3479 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
3480 true
3481 }
3482 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
3483 tracing::span::Id::from_u64(1)
3484 }
3485 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
3486 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
3487 fn event(&self, event: &tracing::Event<'_>) {
3488 #[derive(Default)]
3489 struct FieldVisitor(std::collections::BTreeMap<String, String>);
3490 impl tracing::field::Visit for FieldVisitor {
3491 fn record_debug(
3492 &mut self,
3493 field: &tracing::field::Field,
3494 value: &dyn std::fmt::Debug,
3495 ) {
3496 self.0
3497 .insert(field.name().to_string(), format!("{value:?}"));
3498 }
3499 }
3500 let mut visitor = FieldVisitor::default();
3501 event.record(&mut visitor);
3502 self.events.lock().unwrap().push(visitor.0);
3503 }
3504 fn enter(&self, _: &tracing::span::Id) {}
3505 fn exit(&self, _: &tracing::span::Id) {}
3506 }
3507
3508 #[tokio::test]
3521 #[serial_test::serial(tx_registry)]
3522 async fn issue_1029_starvation_warn_reports_registered_transactions() {
3523 let (_dir, backend) = issue_1029_pool(false);
3524 let text = backend.text("entities").expect("text store");
3527
3528 let holder = backend
3532 .pool
3533 .open_standalone_writer()
3534 .expect("holder connection");
3535 holder
3536 .execute_batch("BEGIN IMMEDIATE")
3537 .expect("holder BEGIN IMMEDIATE");
3538 let fixture =
3539 khive_storage::tx_registry::register(Some("issue_1029_fixture_tx".to_string()));
3540
3541 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
3542 let subscriber = StarvationCaptureSubscriber {
3543 events: Arc::clone(&events),
3544 };
3545 let guard = tracing::subscriber::set_default(subscriber);
3546
3547 let doc = khive_storage::types::TextDocument {
3548 subject_id: uuid::Uuid::new_v4(),
3549 kind: khive_types::SubstrateKind::Entity,
3550 record_kind: None,
3551 title: Some("Issue1029Starved".to_string()),
3552 body: "issue 1029 starvation diagnostic body".to_string(),
3553 tags: vec![],
3554 namespace: "tenant_ns".to_string(),
3555 metadata: None,
3556 updated_at: chrono::Utc::now(),
3557 };
3558 let result = text.upsert_document(doc).await;
3559
3560 drop(guard);
3561 drop(fixture);
3562 holder
3563 .execute_batch("ROLLBACK")
3564 .expect("holder ROLLBACK releases the lock");
3565
3566 assert!(
3567 result.is_err(),
3568 "upsert_document must starve while another connection holds the write lock"
3569 );
3570
3571 let events = events.lock().unwrap();
3572 let warn = events
3573 .iter()
3574 .find(|fields| {
3575 fields
3576 .get("message")
3577 .is_some_and(|m| m.contains("text write starved"))
3578 })
3579 .unwrap_or_else(|| panic!("expected a starvation WARN, captured events: {events:?}"));
3580 assert!(
3581 warn.get("op").is_some_and(|op| op.contains("fts_upsert")),
3582 "WARN must name the starved operation, got: {warn:?}"
3583 );
3584 assert!(
3585 warn.get("open_txs")
3586 .is_some_and(|txs| txs.contains("issue_1029_fixture_tx")),
3587 "WARN must list the registered holder label, got: {warn:?}"
3588 );
3589 let count: usize = warn
3590 .get("open_tx_count")
3591 .expect("WARN must carry open_tx_count")
3592 .parse()
3593 .expect("open_tx_count must be numeric");
3594 assert!(
3595 count >= 1,
3596 "open_tx_count must count the fixture, got {count}"
3597 );
3598 }
3599}
3600
3601#[cfg(test)]
3602#[path = "backend_admission_tests.rs"]
3603mod admission_tests;