1use std::path::Path;
8use std::sync::atomic::{AtomicUsize, Ordering};
9use std::sync::Arc;
10
11use rusqlite::OptionalExtension;
12
13use crate::error::SqliteError;
14use crate::pool::{ConnectionPool, PoolConfig};
15use crate::sql_bridge::SqlBridge;
16use crate::stores::{blob, entity, event, graph, note, sparse, text, vectors};
17
18pub struct StorageBackend {
20 pool: Arc<ConnectionPool>,
21 is_file_backed: bool,
22 path: Option<std::path::PathBuf>,
23 notes_seq_repair_runs: AtomicUsize,
30}
31
32impl StorageBackend {
33 pub fn sqlite(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
39 crate::extension::ensure_extensions_loaded();
40 let resolved = path.as_ref().to_path_buf();
41 let config = PoolConfig {
42 path: Some(resolved.clone()),
43 ..PoolConfig::default()
44 };
45 let pool = ConnectionPool::new(config)?;
46 Ok(Self {
47 pool: Arc::new(pool),
48 is_file_backed: true,
49 path: Some(resolved),
50 notes_seq_repair_runs: AtomicUsize::new(0),
51 })
52 }
53
54 pub fn sqlite_read_only(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
64 crate::extension::ensure_extensions_loaded();
65 let resolved = path.as_ref().to_path_buf();
66 let config = PoolConfig {
67 path: Some(resolved.clone()),
68 read_only: true,
69 ..PoolConfig::default()
70 };
71 let pool = ConnectionPool::new(config)?;
77 Ok(Self {
78 pool: Arc::new(pool),
79 is_file_backed: true,
80 path: Some(resolved),
81 notes_seq_repair_runs: AtomicUsize::new(0),
82 })
83 }
84
85 pub fn memory() -> Result<Self, SqliteError> {
91 crate::extension::ensure_extensions_loaded();
92 let config = PoolConfig {
93 path: None,
94 ..PoolConfig::default()
95 };
96 let pool = ConnectionPool::new(config)?;
97 Ok(Self {
98 pool: Arc::new(pool),
99 is_file_backed: false,
100 path: None,
101 notes_seq_repair_runs: AtomicUsize::new(0),
102 })
103 }
104
105 pub fn sql(&self) -> Arc<dyn khive_storage::SqlAccess> {
109 Arc::new(SqlBridge::new(Arc::clone(&self.pool), self.is_file_backed))
110 }
111
112 pub fn apply_schema(
118 &self,
119 plan: &crate::migrations::ServiceSchemaPlan,
120 ) -> Result<(), SqliteError> {
121 let writer = self.pool.try_writer()?;
122 crate::migrations::apply_schema_plan(writer.conn(), plan)
123 }
124
125 pub fn apply_pack_ddl_statements(
142 &self,
143 statements: &[&'static str],
144 ) -> Result<(), SqliteError> {
145 let writer = self.pool.try_writer()?;
146 writer.transaction(|conn| {
147 for &stmt in statements {
148 conn.execute_batch(stmt)?;
149 }
150 Ok(())
151 })
152 }
153
154 pub fn entities(&self) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
158 self.entities_for_namespace("local")
159 }
160
161 pub fn entities_for_namespace(
165 &self,
166 namespace: &str,
167 ) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
168 if namespace.trim().is_empty() {
169 return Err(SqliteError::InvalidData(
170 "entities namespace must be non-empty".to_string(),
171 ));
172 }
173 let writer = self.pool.try_writer()?;
174 entity::ensure_entities_schema(writer.conn())?;
175
176 Ok(Arc::new(entity::SqlEntityStore::new(
177 Arc::clone(&self.pool),
178 self.is_file_backed,
179 )))
180 }
181
182 pub fn graph(&self) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
187 self.graph_for_namespace("local")
188 }
189
190 pub fn graph_for_namespace(
192 &self,
193 namespace: &str,
194 ) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
195 if namespace.trim().is_empty() {
196 return Err(SqliteError::InvalidData(
197 "graph namespace must be non-empty".to_string(),
198 ));
199 }
200 let writer = self.pool.try_writer()?;
201 graph::ensure_graph_schema(writer.conn())?;
202
203 Ok(Arc::new(graph::SqlGraphStore::new_scoped(
204 Arc::clone(&self.pool),
205 self.is_file_backed,
206 namespace.trim().to_string(),
207 )))
208 }
209
210 pub fn notes(&self) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
214 self.notes_for_namespace("local")
215 }
216
217 pub fn notes_for_namespace(
221 &self,
222 namespace: &str,
223 ) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
224 if namespace.trim().is_empty() {
225 return Err(SqliteError::InvalidData(
226 "notes namespace must be non-empty".to_string(),
227 ));
228 }
229 let writer = self.pool.try_writer()?;
230 note::ensure_notes_schema(writer.conn())?;
231
232 if self.notes_seq_repair_runs.load(Ordering::Relaxed) == 0 {
239 note::repair_notes_seq(writer.conn())?;
240 self.notes_seq_repair_runs.fetch_add(1, Ordering::Relaxed);
241 }
242
243 Ok(Arc::new(note::SqlNoteStore::new(
244 Arc::clone(&self.pool),
245 self.is_file_backed,
246 )))
247 }
248
249 pub fn notes_seq_repair_run_count(&self) -> usize {
254 self.notes_seq_repair_runs.load(Ordering::Relaxed)
255 }
256
257 pub fn events(&self) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
262 self.events_for_namespace("local")
263 }
264
265 pub fn events_for_namespace(
267 &self,
268 namespace: &str,
269 ) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
270 if namespace.trim().is_empty() {
271 return Err(SqliteError::InvalidData(
272 "events namespace must be non-empty".to_string(),
273 ));
274 }
275 let writer = self.pool.try_writer()?;
276 event::ensure_events_schema(writer.conn())?;
277
278 Ok(Arc::new(event::SqlEventStore::new_scoped(
279 Arc::clone(&self.pool),
280 self.is_file_backed,
281 namespace.trim().to_string(),
282 )))
283 }
284
285 pub fn vectors(
291 &self,
292 model_key: &str,
293 embedding_model: &str,
294 dimensions: usize,
295 ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
296 self.vectors_for_namespace(model_key, embedding_model, dimensions, "local")
297 }
298
299 pub fn vectors_for_namespace(
309 &self,
310 model_key: &str,
311 embedding_model: &str,
312 dimensions: usize,
313 namespace: &str,
314 ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
315 if model_key.is_empty()
316 || !model_key
317 .chars()
318 .all(|c| c.is_ascii_alphanumeric() || c == '_')
319 {
320 return Err(SqliteError::InvalidData(format!(
321 "invalid model_key '{}': must be non-empty and contain only \
322 alphanumeric/underscore characters",
323 model_key
324 )));
325 }
326 if namespace.trim().is_empty() {
327 return Err(SqliteError::InvalidData(
328 "vector store namespace must be non-empty".to_string(),
329 ));
330 }
331
332 crate::extension::ensure_extensions_loaded();
334
335 let table = format!("vec_{}", model_key);
336 let writer = self.pool.try_writer()?;
337
338 let table_exists: bool = writer
345 .conn()
346 .query_row(
347 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
348 rusqlite::params![&table],
349 |row| row.get::<_, i64>(0),
350 )
351 .optional()
352 .map_err(SqliteError::Rusqlite)?
353 .is_some();
354
355 if table_exists {
356 let pragma = format!("PRAGMA table_xinfo({})", table);
362 let mut stmt = writer.conn().prepare(&pragma)?;
363 let mut rows = stmt.query([])?;
364 let mut has_field = false;
365 let mut has_embedding_model = false;
366 while let Some(row) = rows.next()? {
367 let name: String = row.get(1)?;
368 if name == "field" {
369 has_field = true;
370 }
371 if name == "embedding_model" {
372 has_embedding_model = true;
373 }
374 }
375 if !has_field || !has_embedding_model {
376 return Err(SqliteError::InvalidData(format!(
377 "vec0 table '{}' is missing required column(s) (field={}, \
378 embedding_model={}); this is a pre-v0.2.8 vector schema and is \
379 not supported — recreate the database",
380 table, has_field, has_embedding_model,
381 )));
382 }
383 }
384
385 writer
394 .conn()
395 .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
396
397 writer
401 .conn()
402 .execute_batch(crate::migrations::ANN_WRITE_LOG_DDL)?;
403 writer
404 .conn()
405 .execute_batch(crate::migrations::ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL)?;
406
407 let ddl = format!(
410 "CREATE VIRTUAL TABLE IF NOT EXISTS vec_{} USING vec0(\
411 subject_id TEXT PRIMARY KEY, \
412 namespace TEXT NOT NULL, \
413 kind TEXT NOT NULL, \
414 field TEXT NOT NULL, \
415 embedding_model TEXT NOT NULL, \
416 embedding float[{}] distance_metric=cosine\
417 )",
418 model_key, dimensions
419 );
420 writer.conn().execute_batch(&ddl)?;
421
422 Ok(Arc::new(vectors::SqliteVecStore::new(
423 Arc::clone(&self.pool),
424 self.is_file_backed,
425 model_key.to_string(),
426 embedding_model.to_string(),
427 dimensions,
428 namespace.trim().to_string(),
429 )?))
430 }
431
432 pub fn register_embedding_model(
437 &self,
438 engine_name: &str,
439 model_id: &str,
440 key_version: &str,
441 dimensions: u32,
442 ) -> Result<(), SqliteError> {
443 let writer = self.pool.try_writer()?;
444 writer
445 .conn()
446 .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
447
448 let now = chrono::Utc::now().timestamp_micros();
449 let canonical_key =
450 format!("{engine_name}:{model_id}:{key_version}:{dimensions}").into_bytes();
451 let id = uuid::Uuid::new_v4();
452 writer.conn().execute(
453 "INSERT INTO _embedding_models \
454 (id, engine_name, model_id, key_version, dim, output_dim, status, \
455 activated_at, superseded_at, superseded_by, canonical_key, created_at) \
456 VALUES (?1, ?2, ?3, ?4, ?5, NULL, 'active', ?6, NULL, NULL, ?7, ?8) \
457 ON CONFLICT(canonical_key) DO UPDATE SET \
458 status = 'active', \
459 activated_at = COALESCE(_embedding_models.activated_at, excluded.activated_at)",
460 rusqlite::params![
461 id.as_bytes().as_slice(),
462 engine_name,
463 model_id,
464 key_version,
465 dimensions as i64,
466 now,
467 canonical_key,
468 now,
469 ],
470 )?;
471 Ok(())
472 }
473
474 pub fn sparse(
478 &self,
479 model_key: &str,
480 ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
481 self.sparse_for_namespace(model_key, "local")
482 }
483
484 pub fn sparse_for_namespace(
488 &self,
489 model_key: &str,
490 namespace: &str,
491 ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
492 if model_key.is_empty()
493 || !model_key
494 .chars()
495 .all(|c| c.is_ascii_alphanumeric() || c == '_')
496 {
497 return Err(SqliteError::InvalidData(format!(
498 "invalid model_key '{}': must be non-empty and contain only alphanumeric/underscore characters",
499 model_key
500 )));
501 }
502 if namespace.trim().is_empty() {
503 return Err(SqliteError::InvalidData(
504 "sparse store namespace must be non-empty".to_string(),
505 ));
506 }
507
508 let writer = self.pool.try_writer()?;
509 sparse::ensure_sparse_schema(writer.conn(), model_key).map_err(SqliteError::Rusqlite)?;
510
511 Ok(Arc::new(sparse::SqliteSparseStore::new(
512 Arc::clone(&self.pool),
513 self.is_file_backed,
514 model_key.to_string(),
515 namespace.trim().to_string(),
516 )?))
517 }
518
519 pub fn text(&self, table_key: &str) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
526 self.text_with_tokenizer(table_key, "trigram")
527 }
528
529 pub fn text_with_tokenizer(
537 &self,
538 table_key: &str,
539 tokenizer: &str,
540 ) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
541 if table_key.is_empty()
542 || !table_key
543 .chars()
544 .all(|c| c.is_ascii_alphanumeric() || c == '_')
545 {
546 return Err(SqliteError::InvalidData(format!(
547 "invalid table_key '{}': must be non-empty and contain only \
548 alphanumeric/underscore characters",
549 table_key
550 )));
551 }
552 if tokenizer.is_empty()
553 || !tokenizer
554 .chars()
555 .all(|c| c.is_ascii_alphanumeric() || c == '_')
556 {
557 return Err(SqliteError::InvalidData(format!(
558 "invalid tokenizer '{}': must be non-empty and contain only \
559 alphanumeric/underscore characters",
560 tokenizer
561 )));
562 }
563
564 let ddl = format!(
565 "CREATE VIRTUAL TABLE IF NOT EXISTS fts_{} USING fts5(\
566 subject_id UNINDEXED, \
567 kind UNINDEXED, \
568 title, \
569 body, \
570 tags UNINDEXED, \
571 namespace UNINDEXED, \
572 metadata UNINDEXED, \
573 updated_at UNINDEXED, \
574 tokenize = '{}'\
575 )",
576 table_key, tokenizer
577 );
578 let writer = self.pool.try_writer()?;
579 writer.conn().execute_batch(&ddl)?;
580
581 Ok(Arc::new(text::Fts5TextSearch::new(
582 Arc::clone(&self.pool),
583 self.is_file_backed,
584 table_key.to_string(),
585 )))
586 }
587
588 pub fn blob_store(
596 &self,
597 config_root: Option<&Path>,
598 floor_bytes: Option<u64>,
599 ) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
600 let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
601 let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
602 Ok(Arc::new(blob::FsBlobStore::new(root, floor)?))
603 }
604
605 pub fn is_file_backed(&self) -> bool {
607 self.is_file_backed
608 }
609
610 pub fn data_dir(&self) -> Option<std::path::PathBuf> {
613 self.path.as_ref()?.parent().map(|p| p.to_path_buf())
614 }
615
616 pub fn ann_root(&self) -> Option<std::path::PathBuf> {
624 ann_root_for(self.path.as_ref()?)
625 }
626
627 pub fn pool(&self) -> &ConnectionPool {
629 &self.pool
630 }
631
632 pub fn pool_arc(&self) -> Arc<ConnectionPool> {
634 Arc::clone(&self.pool)
635 }
636}
637
638fn ann_root_for(path: &std::path::Path) -> Option<std::path::PathBuf> {
643 let mut file = path.file_name()?.to_os_string();
644 file.push(".ann");
645 path.parent().map(|p| p.join(file))
646}
647
648#[cfg(test)]
649mod tests {
650 use super::*;
651 use khive_storage::types::{SqlStatement, SqlValue};
652
653 #[test]
654 fn memory_backend_creates_successfully() {
655 let backend = StorageBackend::memory().expect("memory backend should create");
656 assert!(!backend.is_file_backed());
657 }
658
659 #[test]
660 fn file_backend_creates_successfully() {
661 let dir = tempfile::tempdir().unwrap();
662 let path = dir.path().join("test.db");
663 let backend = StorageBackend::sqlite(&path).expect("file backend should create");
664 assert!(backend.is_file_backed());
665 assert!(path.exists());
666 }
667
668 #[test]
669 fn data_dir_returns_none_for_memory_backend() {
670 let backend = StorageBackend::memory().expect("memory backend");
671 assert!(backend.data_dir().is_none());
672 }
673
674 #[test]
675 fn data_dir_returns_parent_dir_for_file_backend() {
676 let dir = tempfile::tempdir().unwrap();
677 let path = dir.path().join("data.db");
678 let backend = StorageBackend::sqlite(&path).expect("file backend");
679 let got = backend.data_dir().expect("file backend must return Some");
680 assert_eq!(got, dir.path());
681 }
682
683 #[test]
684 fn ann_root_is_database_scoped_sibling_dir() {
685 let dir = tempfile::tempdir().unwrap();
686 let path = dir.path().join("data.db");
687 let backend = StorageBackend::sqlite(&path).expect("file backend");
688 let got = backend.ann_root().expect("file backend must return Some");
689 assert_eq!(got, dir.path().join("data.db.ann"));
690 assert!(StorageBackend::memory().unwrap().ann_root().is_none());
691 }
692
693 #[cfg(unix)]
699 #[test]
700 fn ann_root_distinct_for_non_utf8_filenames() {
701 use std::os::unix::ffi::OsStrExt;
702 let path_a = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xff.db"));
703 let path_b = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xfe.db"));
704 let root_a = ann_root_for(&path_a).expect("Some for a file path");
705 let root_b = ann_root_for(&path_b).expect("Some for a file path");
706 assert_ne!(
707 root_a, root_b,
708 "distinct database files must map to distinct ANN roots"
709 );
710 }
711
712 #[tokio::test]
713 async fn sql_access_memory_roundtrip() {
714 let backend = StorageBackend::memory().unwrap();
715 let sql = backend.sql();
716
717 let mut writer = sql.writer().await.unwrap();
718 writer
719 .execute_script(
720 "CREATE TABLE test_rt (id TEXT PRIMARY KEY, value INTEGER NOT NULL)".into(),
721 )
722 .await
723 .unwrap();
724
725 let affected = writer
726 .execute(SqlStatement {
727 sql: "INSERT INTO test_rt (id, value) VALUES (?1, ?2)".into(),
728 params: vec![SqlValue::Text("row1".into()), SqlValue::Integer(42)],
729 label: None,
730 })
731 .await
732 .unwrap();
733 assert_eq!(affected, 1);
734
735 let mut reader = sql.reader().await.unwrap();
736 let row = reader
737 .query_row(SqlStatement {
738 sql: "SELECT id, value FROM test_rt WHERE id = ?1".into(),
739 params: vec![SqlValue::Text("row1".into())],
740 label: None,
741 })
742 .await
743 .unwrap();
744
745 let row = row.expect("should find the inserted row");
746 assert_eq!(row.columns.len(), 2);
747 match &row.columns[0].value {
748 SqlValue::Text(s) => assert_eq!(s, "row1"),
749 other => panic!("expected Text, got {other:?}"),
750 }
751 match &row.columns[1].value {
752 SqlValue::Integer(v) => assert_eq!(*v, 42),
753 other => panic!("expected Integer, got {other:?}"),
754 }
755 }
756
757 #[tokio::test]
758 async fn sql_access_file_roundtrip() {
759 let dir = tempfile::tempdir().unwrap();
760 let path = dir.path().join("test_roundtrip.db");
761 let backend = StorageBackend::sqlite(&path).unwrap();
762 let sql = backend.sql();
763
764 let mut writer = sql.writer().await.unwrap();
765 writer
766 .execute_script("CREATE TABLE test_f (k TEXT PRIMARY KEY, v TEXT)".into())
767 .await
768 .unwrap();
769 writer
770 .execute(SqlStatement {
771 sql: "INSERT INTO test_f (k, v) VALUES (?1, ?2)".into(),
772 params: vec![
773 SqlValue::Text("hello".into()),
774 SqlValue::Text("world".into()),
775 ],
776 label: None,
777 })
778 .await
779 .unwrap();
780
781 let mut reader = sql.reader().await.unwrap();
782 let rows = reader
783 .query_all(SqlStatement {
784 sql: "SELECT k, v FROM test_f".into(),
785 params: vec![],
786 label: None,
787 })
788 .await
789 .unwrap();
790 assert_eq!(rows.len(), 1);
791 match &rows[0].columns[1].value {
792 SqlValue::Text(s) => assert_eq!(s, "world"),
793 other => panic!("expected Text, got {other:?}"),
794 }
795 }
796
797 #[test]
798 fn sqlite_read_only_missing_path_does_not_create_file() {
799 let dir = tempfile::tempdir().unwrap();
800 let path = dir.path().join("missing_ro.db");
801 assert!(!path.exists());
802
803 let result = StorageBackend::sqlite_read_only(&path);
804 assert!(
805 result.is_err(),
806 "opening a missing path read-only must fail"
807 );
808 assert!(
809 !path.exists(),
810 "opening a missing path read-only must not create the file"
811 );
812 }
813
814 #[tokio::test]
815 async fn sqlite_read_only_sql_writer_rejects_ddl_and_insert() {
816 let dir = tempfile::tempdir().unwrap();
817 let path = dir.path().join("ro_writer.db");
818
819 {
821 let writable = StorageBackend::sqlite(&path).unwrap();
822 let sql = writable.sql();
823 let mut writer = sql.writer().await.unwrap();
824 writer
825 .execute_script("CREATE TABLE ro_existing (id INTEGER PRIMARY KEY)".into())
826 .await
827 .unwrap();
828 }
829
830 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
831 let sql = ro.sql();
832
833 let writer_result = sql.writer().await;
835 assert!(
836 writer_result.is_err(),
837 "sql().writer() must be rejected on a read-only backend"
838 );
839 }
840
841 #[tokio::test]
842 #[cfg(feature = "vectors")]
843 async fn vectors_roundtrip_via_public_api() {
844 let backend = StorageBackend::memory().unwrap();
845 let store = backend.vectors("test_api", "test_api", 3).unwrap();
846
847 let id = uuid::Uuid::new_v4();
848 store
849 .insert(
850 id,
851 khive_types::SubstrateKind::Entity,
852 "local",
853 "content",
854 vec![vec![1.0, 0.0, 0.0]],
855 )
856 .await
857 .unwrap();
858
859 let hits = store
860 .search(khive_storage::types::VectorSearchRequest {
861 query_vectors: vec![vec![1.0, 0.0, 0.0]],
862 top_k: 1,
863 namespace: None,
864 kind: None,
865 embedding_model: None,
866 filter: None,
867 backend_hints: None,
868 })
869 .await
870 .unwrap();
871
872 assert_eq!(hits.len(), 1);
873 assert_eq!(hits[0].subject_id, id);
874 assert!(hits[0].score.to_f64() > 0.99);
875 }
876
877 #[tokio::test]
878 #[cfg(feature = "vectors")]
879 async fn vectors_creates_table_idempotently() {
880 let backend = StorageBackend::memory().unwrap();
881
882 let store1 = backend.vectors("idempotent", "idempotent", 3).unwrap();
883 let store2 = backend.vectors("idempotent", "idempotent", 3).unwrap();
884
885 let id = uuid::Uuid::new_v4();
886 store1
887 .insert(
888 id,
889 khive_types::SubstrateKind::Entity,
890 "local",
891 "content",
892 vec![vec![1.0, 0.0, 0.0]],
893 )
894 .await
895 .unwrap();
896
897 let count = store2.count().await.unwrap();
898 assert_eq!(count, 1);
899 }
900
901 #[tokio::test]
902 async fn text_roundtrip_via_public_api() {
903 let backend = StorageBackend::memory().unwrap();
904 let store = backend.text("test_api").unwrap();
905
906 let id = uuid::Uuid::new_v4();
907 let doc = khive_storage::types::TextDocument {
908 subject_id: id,
909 kind: khive_types::SubstrateKind::Entity,
910 title: Some("Test Title".to_string()),
911 body: "This is a searchable document about Rust.".to_string(),
912 tags: vec!["rust".to_string()],
913 namespace: "test_ns".to_string(),
914 metadata: None,
915 updated_at: chrono::Utc::now(),
916 };
917 store.upsert_document(doc).await.unwrap();
918
919 let hits = store
920 .search(khive_storage::types::TextSearchRequest {
921 query: "Rust".to_string(),
922 mode: khive_storage::types::TextQueryMode::Plain,
923 filter: Some(khive_storage::types::TextFilter {
924 namespaces: vec!["test_ns".to_string()],
925 ..Default::default()
926 }),
927 top_k: 1,
928 snippet_chars: 64,
929 })
930 .await
931 .unwrap();
932
933 assert_eq!(hits.len(), 1);
934 assert_eq!(hits[0].subject_id, id);
935 assert!(hits[0].score.to_f64() > 0.0);
936 }
937
938 #[tokio::test]
939 async fn text_creates_table_idempotently() {
940 let backend = StorageBackend::memory().unwrap();
941
942 let store1 = backend.text("idempotent_fts").unwrap();
943 let store2 = backend.text("idempotent_fts").unwrap();
944
945 let id = uuid::Uuid::new_v4();
946 let doc = khive_storage::types::TextDocument {
947 subject_id: id,
948 kind: khive_types::SubstrateKind::Note,
949 title: None,
950 body: "Hello world.".to_string(),
951 tags: vec![],
952 namespace: "test_ns".to_string(),
953 metadata: None,
954 updated_at: chrono::Utc::now(),
955 };
956 store1.upsert_document(doc).await.unwrap();
957
958 let count = store2
959 .count(khive_storage::types::TextFilter {
960 namespaces: vec!["test_ns".to_string()],
961 ..Default::default()
962 })
963 .await
964 .unwrap();
965 assert_eq!(count, 1);
966 }
967
968 #[test]
969 fn invalid_model_key_rejected() {
970 let backend = StorageBackend::memory().unwrap();
971 assert!(backend.vectors("bad key!", "bad key!", 3).is_err());
972 assert!(backend.vectors("", "", 3).is_err());
973 }
974
975 #[test]
976 fn invalid_table_key_rejected() {
977 let backend = StorageBackend::memory().unwrap();
978 assert!(backend.text("bad key!").is_err());
979 assert!(backend.text("").is_err());
980 }
981
982 #[tokio::test]
983 async fn sqlite_read_only_graph_store_rejects_upsert_edge() {
984 use khive_storage::types::Edge;
985 use khive_types::EdgeRelation;
986
987 let dir = tempfile::tempdir().unwrap();
988 let path = dir.path().join("ro_graph.db");
989
990 {
992 let writable = StorageBackend::sqlite(&path).unwrap();
993 writable.graph().unwrap();
994 }
995
996 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
997 let store = match ro.graph() {
998 Ok(store) => store,
999 Err(_) => return,
1002 };
1003
1004 let now = chrono::Utc::now();
1005 let edge = Edge {
1006 id: uuid::Uuid::new_v4().into(),
1007 namespace: "local".to_string(),
1008 source_id: uuid::Uuid::new_v4(),
1009 target_id: uuid::Uuid::new_v4(),
1010 relation: EdgeRelation::Extends,
1011 weight: 0.8,
1012 created_at: now,
1013 updated_at: now,
1014 deleted_at: None,
1015 metadata: None,
1016 target_backend: None,
1017 };
1018
1019 let result = store.upsert_edge(edge).await;
1020 assert!(
1021 result.is_err(),
1022 "upsert_edge on a read-only backend must reject, not silently no-op"
1023 );
1024 }
1025
1026 #[tokio::test]
1027 async fn sqlite_read_only_event_store_rejects_append_event() {
1028 use khive_types::{EventKind, EventOutcome, SubstrateKind};
1029
1030 let dir = tempfile::tempdir().unwrap();
1031 let path = dir.path().join("ro_events.db");
1032
1033 {
1034 let writable = StorageBackend::sqlite(&path).unwrap();
1035 writable.events().unwrap();
1036 }
1037
1038 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1039 let store = match ro.events() {
1040 Ok(store) => store,
1041 Err(_) => return,
1042 };
1043
1044 let event = khive_storage::event::Event::new(
1045 "local",
1046 "test.verb",
1047 EventKind::Audit,
1048 SubstrateKind::Entity,
1049 "test-actor",
1050 )
1051 .with_outcome(EventOutcome::Success);
1052
1053 let result = store.append_event(event).await;
1054 assert!(
1055 result.is_err(),
1056 "append_event on a read-only backend must reject, not silently no-op"
1057 );
1058 }
1059
1060 #[tokio::test]
1061 async fn sqlite_read_only_text_store_rejects_upsert_document() {
1062 use khive_storage::types::TextDocument;
1063 use khive_types::SubstrateKind;
1064
1065 let dir = tempfile::tempdir().unwrap();
1066 let path = dir.path().join("ro_text.db");
1067
1068 {
1069 let writable = StorageBackend::sqlite(&path).unwrap();
1070 writable.text("ro_test").unwrap();
1071 }
1072
1073 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1074 let store = match ro.text("ro_test") {
1075 Ok(store) => store,
1076 Err(_) => return,
1077 };
1078
1079 let doc = TextDocument {
1080 subject_id: uuid::Uuid::new_v4(),
1081 kind: SubstrateKind::Entity,
1082 title: Some("Title".to_string()),
1083 body: "Body text.".to_string(),
1084 tags: vec![],
1085 namespace: "local".to_string(),
1086 metadata: None,
1087 updated_at: chrono::Utc::now(),
1088 };
1089
1090 let result = store.upsert_document(doc).await;
1091 assert!(
1092 result.is_err(),
1093 "upsert_document on a read-only backend must reject, not silently no-op"
1094 );
1095 }
1096
1097 #[tokio::test]
1098 async fn blob_store_roundtrip_via_public_api() {
1099 let dir = tempfile::tempdir().unwrap();
1100 let path = dir.path().join("blob_backend.db");
1101 let backend = StorageBackend::sqlite(&path).unwrap();
1102
1103 let store = backend.blob_store(None, Some(0)).unwrap();
1107 let bytes = b"backend-level blob roundtrip".to_vec();
1108 let content_ref = store.put(bytes.clone()).await.unwrap();
1109 assert_eq!(store.get(&content_ref).await.unwrap(), bytes);
1110 }
1111
1112 #[test]
1113 fn blob_store_defaults_root_beside_db_file() {
1114 let dir = tempfile::tempdir().unwrap();
1115 let path = dir.path().join("blob_default.db");
1116 let backend = StorageBackend::sqlite(&path).unwrap();
1117
1118 let _store = backend.blob_store(None, None).unwrap();
1122 assert!(
1123 dir.path().join("blobs").is_dir(),
1124 "default root must be created beside the database file"
1125 );
1126 }
1127
1128 #[test]
1129 fn blob_store_errors_for_in_memory_backend_with_no_override() {
1130 let backend = StorageBackend::memory().unwrap();
1131 assert!(backend.blob_store(None, None).is_err());
1132 }
1133
1134 #[test]
1135 fn blob_store_accepts_explicit_root_for_in_memory_backend() {
1136 let dir = tempfile::tempdir().unwrap();
1137 let backend = StorageBackend::memory().unwrap();
1138 let store = backend.blob_store(Some(dir.path()), None);
1139 assert!(store.is_ok());
1140 }
1141
1142 #[test]
1143 fn apply_schema_runs_migrations_idempotently() {
1144 static MIGRATIONS: &[crate::migrations::Migration] = &[crate::migrations::Migration {
1145 id: "001_init",
1146 up_sql: "CREATE TABLE IF NOT EXISTS schema_test (id TEXT PRIMARY KEY);",
1147 down_sql: None,
1148 is_already_applied: None,
1149 }];
1150 let plan = crate::migrations::ServiceSchemaPlan {
1151 service: "schema_test_svc",
1152 sqlite: MIGRATIONS,
1153 postgres: &[],
1154 };
1155
1156 let backend = StorageBackend::memory().unwrap();
1157 backend.apply_schema(&plan).unwrap();
1158 backend.apply_schema(&plan).unwrap();
1159
1160 let reader = backend.pool().reader().unwrap();
1161 let count: i64 = reader
1162 .conn()
1163 .query_row(
1164 "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='schema_test'",
1165 [],
1166 |row| row.get(0),
1167 )
1168 .unwrap();
1169 assert_eq!(count, 1);
1170 }
1171
1172 #[test]
1173 fn pack_ddl_plan_rolls_back_all_statements_on_failure() {
1174 let backend = StorageBackend::memory().unwrap();
1175 let error = backend
1176 .apply_pack_ddl_statements(&[
1177 "CREATE TABLE IF NOT EXISTS pack_schema_first (id INTEGER PRIMARY KEY)",
1178 "CREATE INDEX IF NOT EXISTS pack_schema_second ON pack_schema_missing(id)",
1179 ])
1180 .unwrap_err();
1181
1182 assert!(
1183 error.to_string().contains("pack_schema_missing"),
1184 "schema-plan error must retain the failing SQLite diagnostic: {error}"
1185 );
1186
1187 let reader = backend.pool().reader().unwrap();
1188 let visible_objects: i64 = reader
1189 .conn()
1190 .query_row(
1191 "SELECT COUNT(*) FROM sqlite_master \
1192 WHERE name IN ('pack_schema_first', 'pack_schema_second')",
1193 [],
1194 |row| row.get(0),
1195 )
1196 .unwrap();
1197 assert_eq!(visible_objects, 0);
1198 }
1199
1200 #[test]
1201 fn pack_ddl_plan_applies_all_statements_idempotently() {
1202 const PLAN: &[&str] = &[
1203 "CREATE TABLE IF NOT EXISTS pack_schema_success (id INTEGER PRIMARY KEY, value TEXT)",
1204 "CREATE INDEX IF NOT EXISTS pack_schema_success_value_idx \
1205 ON pack_schema_success(value)",
1206 ];
1207
1208 let backend = StorageBackend::memory().unwrap();
1209 backend.apply_pack_ddl_statements(PLAN).unwrap();
1210 backend.apply_pack_ddl_statements(PLAN).unwrap();
1211
1212 let reader = backend.pool().reader().unwrap();
1213 let visible_objects: i64 = reader
1214 .conn()
1215 .query_row(
1216 "SELECT COUNT(*) FROM sqlite_master \
1217 WHERE name IN ('pack_schema_success', 'pack_schema_success_value_idx')",
1218 [],
1219 |row| row.get(0),
1220 )
1221 .unwrap();
1222 assert_eq!(visible_objects, 2);
1223 }
1224
1225 fn issue_1029_pool(write_queue_enabled: bool) -> (tempfile::TempDir, StorageBackend) {
1234 let dir = tempfile::tempdir().unwrap();
1235 let path = dir.path().join("issue_1029.db");
1236 let config = crate::pool::PoolConfig {
1237 path: Some(path.clone()),
1238 busy_timeout: std::time::Duration::from_millis(200),
1239 write_queue_enabled,
1240 ..crate::pool::PoolConfig::default()
1241 };
1242 let pool = ConnectionPool::new(config).expect("fresh tenant-shaped pool should open");
1243 let backend = StorageBackend {
1244 pool: Arc::new(pool),
1245 is_file_backed: true,
1246 path: Some(path),
1247 notes_seq_repair_runs: AtomicUsize::new(0),
1248 };
1249 (dir, backend)
1250 }
1251
1252 async fn issue_1029_create_entity_shaped_sequence(
1253 backend: &StorageBackend,
1254 ) -> Result<(), String> {
1255 let entities = backend
1256 .entities_for_namespace("tenant_ns")
1257 .map_err(|e| format!("entities_for_namespace: {e}"))?;
1258 let entity = khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029Repro");
1259 let entity_id = entity.id;
1260 entities
1261 .upsert_entity(entity)
1262 .await
1263 .map_err(|e| format!("upsert_entity: {e}"))?;
1264
1265 let text = backend.text("entities").map_err(|e| format!("text: {e}"))?;
1266 let doc = khive_storage::types::TextDocument {
1267 subject_id: entity_id,
1268 kind: khive_types::SubstrateKind::Entity,
1269 title: Some("Issue1029Repro".to_string()),
1270 body: "issue 1029 repro body".to_string(),
1271 tags: vec![],
1272 namespace: "tenant_ns".to_string(),
1273 metadata: None,
1274 updated_at: chrono::Utc::now(),
1275 };
1276 text.upsert_document(doc)
1277 .await
1278 .map_err(|e| format!("fts_upsert: {e}"))
1279 }
1280
1281 #[tokio::test]
1287 async fn issue_1029_create_entity_shaped_sequence_write_queue_off() {
1288 let (_dir, backend) = issue_1029_pool(false);
1289 let result = issue_1029_create_entity_shaped_sequence(&backend).await;
1290 assert!(
1291 result.is_ok(),
1292 "khive#1029 repro (KHIVE_WRITE_QUEUE off): fts_upsert step failed: {:?}",
1293 result.err()
1294 );
1295 }
1296
1297 #[tokio::test]
1303 async fn issue_1029_create_entity_shaped_sequence_write_queue_on() {
1304 let (_dir, backend) = issue_1029_pool(true);
1305 let result = issue_1029_create_entity_shaped_sequence(&backend).await;
1306 assert!(
1307 result.is_ok(),
1308 "khive#1029 repro (KHIVE_WRITE_QUEUE=1): fts_upsert step failed: {:?}",
1309 result.err()
1310 );
1311 }
1312
1313 #[tokio::test]
1321 async fn issue_1029_two_pools_same_file_write_queue_on() {
1322 let dir = tempfile::tempdir().unwrap();
1323 let path = dir.path().join("issue_1029_two_pools.db");
1324
1325 let cfg = |p: std::path::PathBuf| crate::pool::PoolConfig {
1326 path: Some(p),
1327 busy_timeout: std::time::Duration::from_millis(200),
1328 write_queue_enabled: true,
1329 ..crate::pool::PoolConfig::default()
1330 };
1331
1332 let pool_a = ConnectionPool::new(cfg(path.clone())).expect("pool A should open");
1333 let backend_a = StorageBackend {
1334 pool: Arc::new(pool_a),
1335 is_file_backed: true,
1336 path: Some(path.clone()),
1337 notes_seq_repair_runs: AtomicUsize::new(0),
1338 };
1339 let pool_b = ConnectionPool::new(cfg(path.clone())).expect("pool B should open");
1340 let backend_b = StorageBackend {
1341 pool: Arc::new(pool_b),
1342 is_file_backed: true,
1343 path: Some(path),
1344 notes_seq_repair_runs: AtomicUsize::new(0),
1345 };
1346
1347 let entities = backend_a
1348 .entities_for_namespace("tenant_ns")
1349 .expect("entities_for_namespace on pool A");
1350 let entity =
1351 khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029TwoPools");
1352 let entity_id = entity.id;
1353 entities
1354 .upsert_entity(entity)
1355 .await
1356 .expect("pool A entity upsert should succeed");
1357
1358 let text = backend_b.text("entities").expect("text on pool B");
1359 let doc = khive_storage::types::TextDocument {
1360 subject_id: entity_id,
1361 kind: khive_types::SubstrateKind::Entity,
1362 title: Some("Issue1029TwoPools".to_string()),
1363 body: "issue 1029 two-pool repro body".to_string(),
1364 tags: vec![],
1365 namespace: "tenant_ns".to_string(),
1366 metadata: None,
1367 updated_at: chrono::Utc::now(),
1368 };
1369 let result = text.upsert_document(doc).await;
1370 assert!(
1371 result.is_ok(),
1372 "khive#1029 two-pool repro: fts_upsert on an independent pool for the \
1373 same tenant DB file failed: {:?}",
1374 result.err()
1375 );
1376 }
1377
1378 struct StarvationCaptureSubscriber {
1381 events: Arc<std::sync::Mutex<Vec<std::collections::BTreeMap<String, String>>>>,
1382 }
1383
1384 impl tracing::Subscriber for StarvationCaptureSubscriber {
1385 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
1386 true
1387 }
1388 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
1389 tracing::span::Id::from_u64(1)
1390 }
1391 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
1392 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
1393 fn event(&self, event: &tracing::Event<'_>) {
1394 #[derive(Default)]
1395 struct FieldVisitor(std::collections::BTreeMap<String, String>);
1396 impl tracing::field::Visit for FieldVisitor {
1397 fn record_debug(
1398 &mut self,
1399 field: &tracing::field::Field,
1400 value: &dyn std::fmt::Debug,
1401 ) {
1402 self.0
1403 .insert(field.name().to_string(), format!("{value:?}"));
1404 }
1405 }
1406 let mut visitor = FieldVisitor::default();
1407 event.record(&mut visitor);
1408 self.events.lock().unwrap().push(visitor.0);
1409 }
1410 fn enter(&self, _: &tracing::span::Id) {}
1411 fn exit(&self, _: &tracing::span::Id) {}
1412 }
1413
1414 #[tokio::test]
1427 #[serial_test::serial(tx_registry)]
1428 async fn issue_1029_starvation_warn_reports_registered_transactions() {
1429 let (_dir, backend) = issue_1029_pool(false);
1430 let text = backend.text("entities").expect("text store");
1433
1434 let holder = backend
1438 .pool
1439 .open_standalone_writer()
1440 .expect("holder connection");
1441 holder
1442 .execute_batch("BEGIN IMMEDIATE")
1443 .expect("holder BEGIN IMMEDIATE");
1444 let fixture =
1445 khive_storage::tx_registry::register(Some("issue_1029_fixture_tx".to_string()));
1446
1447 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
1448 let subscriber = StarvationCaptureSubscriber {
1449 events: Arc::clone(&events),
1450 };
1451 let guard = tracing::subscriber::set_default(subscriber);
1452
1453 let doc = khive_storage::types::TextDocument {
1454 subject_id: uuid::Uuid::new_v4(),
1455 kind: khive_types::SubstrateKind::Entity,
1456 title: Some("Issue1029Starved".to_string()),
1457 body: "issue 1029 starvation diagnostic body".to_string(),
1458 tags: vec![],
1459 namespace: "tenant_ns".to_string(),
1460 metadata: None,
1461 updated_at: chrono::Utc::now(),
1462 };
1463 let result = text.upsert_document(doc).await;
1464
1465 drop(guard);
1466 drop(fixture);
1467 holder
1468 .execute_batch("ROLLBACK")
1469 .expect("holder ROLLBACK releases the lock");
1470
1471 assert!(
1472 result.is_err(),
1473 "upsert_document must starve while another connection holds the write lock"
1474 );
1475
1476 let events = events.lock().unwrap();
1477 let warn = events
1478 .iter()
1479 .find(|fields| {
1480 fields
1481 .get("message")
1482 .is_some_and(|m| m.contains("text write starved"))
1483 })
1484 .unwrap_or_else(|| panic!("expected a starvation WARN, captured events: {events:?}"));
1485 assert!(
1486 warn.get("op").is_some_and(|op| op.contains("fts_upsert")),
1487 "WARN must name the starved operation, got: {warn:?}"
1488 );
1489 assert!(
1490 warn.get("open_txs")
1491 .is_some_and(|txs| txs.contains("issue_1029_fixture_tx")),
1492 "WARN must list the registered holder label, got: {warn:?}"
1493 );
1494 let count: usize = warn
1495 .get("open_tx_count")
1496 .expect("WARN must carry open_tx_count")
1497 .parse()
1498 .expect("open_tx_count must be numeric");
1499 assert!(
1500 count >= 1,
1501 "open_tx_count must count the fixture, got {count}"
1502 );
1503 }
1504}