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