1use std::path::Path;
9use std::sync::atomic::{AtomicUsize, Ordering};
10use std::sync::Arc;
11
12use rusqlite::OptionalExtension;
13
14use crate::error::SqliteError;
15use crate::pool::{ConnectionPool, PoolConfig};
16use crate::sql_bridge::SqlBridge;
17use crate::stores::{agents, attachment, blob, entity, event, graph, note, sparse, text, vectors};
18
19fn sqlite_table_exists(conn: &rusqlite::Connection, table: &str) -> Result<bool, SqliteError> {
20 conn.query_row(
21 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
22 rusqlite::params![table],
23 |row| row.get::<_, i64>(0),
24 )
25 .optional()
26 .map(|row| row.is_some())
27 .map_err(SqliteError::Rusqlite)
28}
29
30fn validate_vector_table_columns(
31 conn: &rusqlite::Connection,
32 table: &str,
33) -> Result<(), SqliteError> {
34 let pragma = format!("PRAGMA table_xinfo({table})");
35 let mut stmt = conn.prepare(&pragma)?;
36 let mut rows = stmt.query([])?;
37 let mut has_field = false;
38 let mut has_embedding_model = false;
39 while let Some(row) = rows.next()? {
40 let name: String = row.get(1)?;
41 if name == "field" {
42 has_field = true;
43 }
44 if name == "embedding_model" {
45 has_embedding_model = true;
46 }
47 }
48 if !has_field || !has_embedding_model {
49 return Err(SqliteError::InvalidData(format!(
50 "vec0 table '{table}' is missing required column(s) (field={has_field}, \
51 embedding_model={has_embedding_model}); this is a pre-v0.2.8 vector schema and is \
52 not supported — recreate the database"
53 )));
54 }
55 Ok(())
56}
57
58pub struct StorageBackend {
60 pool: Arc<ConnectionPool>,
61 is_file_backed: bool,
62 path: Option<std::path::PathBuf>,
63 notes_seq_repair_runs: AtomicUsize,
70}
71
72impl StorageBackend {
73 pub fn sqlite(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
81 crate::extension::ensure_extensions_loaded();
82 let resolved = path.as_ref().to_path_buf();
83 let read_only =
84 std::fs::metadata(&resolved).is_ok_and(|metadata| metadata.permissions().readonly());
85 let mut config = PoolConfig {
86 path: Some(resolved.clone()),
87 read_only,
88 ..PoolConfig::default()
89 };
90 if read_only {
91 config.write_queue_enabled = Some(false);
92 }
93 let pool = ConnectionPool::new(config)?;
94 Ok(Self {
95 pool: Arc::new(pool),
96 is_file_backed: true,
97 path: Some(resolved),
98 notes_seq_repair_runs: AtomicUsize::new(0),
99 })
100 }
101
102 pub fn sqlite_read_only(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
113 crate::extension::ensure_extensions_loaded();
114 let resolved = path.as_ref().to_path_buf();
115 let config = PoolConfig {
116 path: Some(resolved.clone()),
117 read_only: true,
118 write_queue_enabled: Some(false),
119 ..PoolConfig::default()
120 };
121 let pool = ConnectionPool::new(config)?;
127 Ok(Self {
128 pool: Arc::new(pool),
129 is_file_backed: true,
130 path: Some(resolved),
131 notes_seq_repair_runs: AtomicUsize::new(0),
132 })
133 }
134
135 pub fn memory() -> Result<Self, SqliteError> {
141 crate::extension::ensure_extensions_loaded();
142 let config = PoolConfig {
143 path: None,
144 ..PoolConfig::default()
145 };
146 let pool = ConnectionPool::new(config)?;
147 Ok(Self {
148 pool: Arc::new(pool),
149 is_file_backed: false,
150 path: None,
151 notes_seq_repair_runs: AtomicUsize::new(0),
152 })
153 }
154
155 pub fn sql(&self) -> Arc<dyn khive_storage::SqlAccess> {
159 Arc::new(SqlBridge::new(Arc::clone(&self.pool), self.is_file_backed))
160 }
161
162 pub fn apply_schema(
168 &self,
169 plan: &crate::migrations::ServiceSchemaPlan,
170 ) -> Result<(), SqliteError> {
171 let writer = self.pool.try_writer()?;
172 crate::migrations::apply_schema_plan(writer.conn(), plan)
173 }
174
175 pub fn apply_pack_ddl_statements(
192 &self,
193 statements: &[&'static str],
194 ) -> Result<(), SqliteError> {
195 let writer = self.pool.try_writer()?;
196 writer.transaction(|conn| {
197 for &stmt in statements {
198 conn.execute_batch(stmt)?;
199 }
200 Ok(())
201 })
202 }
203
204 pub fn prepare_core_schema(&self) -> Result<u32, SqliteError> {
214 if self.is_read_only() {
215 let reader = self.pool.reader()?;
216 crate::migrations::validate_schema_is_current(reader.conn())
217 } else {
218 let latest = crate::migrations::MIGRATIONS
219 .last()
220 .map(|migration| migration.version)
221 .unwrap_or(0);
222 {
223 let reader = self.pool.reader()?;
224 let current = crate::migrations::read_schema_version(reader.conn())?;
225 if current >= latest {
226 return crate::migrations::validate_schema_is_current(reader.conn());
227 }
228 }
229 let owner = crate::stores::blob::acquire_database_gc_owner_for_path_blocking(
230 self.pool.canonical_path().map(Path::to_path_buf),
231 )
232 .map_err(|error| {
233 SqliteError::InvalidData(format!(
234 "failed to acquire database GC owner before schema preparation: {error}"
235 ))
236 })?;
237 let mut writer = self.pool.try_writer()?;
238 crate::migrations::run_migrations_with_database_gc_owner(writer.conn_mut(), &owner)
239 }
240 }
241
242 pub fn attachment_cutover_status(
244 &self,
245 ) -> Result<crate::migrations::AttachmentCutoverStatus, SqliteError> {
246 if self.is_read_only() {
247 let reader = self.pool.reader()?;
248 crate::migrations::attachment_cutover_status(reader.conn())
249 } else {
250 let writer = self.pool.try_writer()?;
251 crate::migrations::attachment_cutover_status(writer.conn())
252 }
253 }
254
255 fn require_attachment_cutover_owner(
256 &self,
257 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
258 ) -> Result<(), SqliteError> {
259 let sql = self.sql();
260 let backend_path = sql.database_path();
261 if owner.database_path() != backend_path.as_deref() {
262 return Err(SqliteError::InvalidData(format!(
263 "attachment cutover GC owner targets {:?}, but this backend is {:?}",
264 owner.database_path(),
265 backend_path.as_deref()
266 )));
267 }
268 Ok(())
269 }
270
271 pub fn stage_attachment_cutover(
275 &self,
276 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
277 ) -> Result<(), SqliteError> {
278 self.require_attachment_cutover_owner(owner)?;
279 if self.is_read_only() {
280 return Err(SqliteError::InvalidData(
281 "cannot stage attachment cutover on a read-only backend".into(),
282 ));
283 }
284 let mut writer = self.pool.try_writer()?;
285 crate::migrations::stage_attachment_cutover(writer.conn_mut())
286 }
287
288 pub fn apply_verified_attachments(
290 &self,
291 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
292 attachments: &[khive_storage::Attachment],
293 ) -> Result<(), SqliteError> {
294 self.require_attachment_cutover_owner(owner)?;
295 if self.is_read_only() {
296 return Err(SqliteError::InvalidData(
297 "cannot apply verified attachments on a read-only backend".into(),
298 ));
299 }
300 let mut writer = self.pool.try_writer()?;
301 let tx = writer
302 .conn_mut()
303 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
304 for attachment in attachments {
305 attachment
306 .validate()
307 .map_err(|error| SqliteError::InvalidData(error.to_string()))?;
308 crate::migrations::apply_generic_verified_attachment(
309 &tx,
310 &attachment.record_uuid.to_string(),
311 attachment.substrate.as_str(),
312 &attachment.role,
313 &attachment.content_ref,
314 attachment.media_type.as_deref(),
315 attachment.size_bytes,
316 attachment.created_at,
317 )?;
318 }
319 tx.commit()?;
320 Ok(())
321 }
322
323 pub fn finalize_attachment_cutover(
326 &self,
327 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
328 ) -> Result<(), SqliteError> {
329 self.require_attachment_cutover_owner(owner)?;
330 if self.is_read_only() {
331 return Err(SqliteError::InvalidData(
332 "cannot finalize attachment cutover on a read-only backend".into(),
333 ));
334 }
335 let mut writer = self.pool.try_writer()?;
336 crate::migrations::finalize_attachment_cutover(writer.conn_mut())
337 }
338
339 pub fn entities(&self) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
343 self.entities_for_namespace("local")
344 }
345
346 pub fn entities_for_namespace(
350 &self,
351 namespace: &str,
352 ) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
353 if namespace.trim().is_empty() {
354 return Err(SqliteError::InvalidData(
355 "entities namespace must be non-empty".to_string(),
356 ));
357 }
358 if !self.is_read_only() {
359 let writer = self.pool.try_writer()?;
360 entity::ensure_entities_schema(writer.conn())?;
361 }
362
363 Ok(Arc::new(entity::SqlEntityStore::new(
364 Arc::clone(&self.pool),
365 self.is_file_backed,
366 )))
367 }
368
369 pub fn attachments(&self) -> Result<Arc<dyn khive_storage::AttachmentStore>, SqliteError> {
376 Ok(Arc::new(attachment::SqlAttachmentStore::new(
377 Arc::clone(&self.pool),
378 self.is_file_backed,
379 )))
380 }
381
382 pub fn graph(&self) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
387 self.graph_for_namespace("local")
388 }
389
390 pub fn graph_for_namespace(
392 &self,
393 namespace: &str,
394 ) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
395 if namespace.trim().is_empty() {
396 return Err(SqliteError::InvalidData(
397 "graph namespace must be non-empty".to_string(),
398 ));
399 }
400 if !self.is_read_only() {
401 let writer = self.pool.try_writer()?;
402 graph::ensure_graph_schema(writer.conn())?;
403 }
404
405 Ok(Arc::new(graph::SqlGraphStore::new_scoped(
406 Arc::clone(&self.pool),
407 self.is_file_backed,
408 namespace.trim().to_string(),
409 )))
410 }
411
412 pub fn notes(&self) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
416 self.notes_for_namespace("local")
417 }
418
419 pub fn notes_for_namespace(
423 &self,
424 namespace: &str,
425 ) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
426 if namespace.trim().is_empty() {
427 return Err(SqliteError::InvalidData(
428 "notes namespace must be non-empty".to_string(),
429 ));
430 }
431 if !self.is_read_only() {
432 let writer = self.pool.try_writer()?;
433 note::ensure_notes_schema(writer.conn())?;
434
435 if self.notes_seq_repair_runs.load(Ordering::Relaxed) == 0 {
442 note::repair_notes_seq(writer.conn())?;
443 self.notes_seq_repair_runs.fetch_add(1, Ordering::Relaxed);
444 }
445 }
446
447 Ok(Arc::new(note::SqlNoteStore::new(
448 Arc::clone(&self.pool),
449 self.is_file_backed,
450 )))
451 }
452
453 pub fn notes_seq_repair_run_count(&self) -> usize {
458 self.notes_seq_repair_runs.load(Ordering::Relaxed)
459 }
460
461 pub fn events(&self) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
466 self.events_for_namespace("local")
467 }
468
469 pub fn events_for_namespace(
471 &self,
472 namespace: &str,
473 ) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
474 if namespace.trim().is_empty() {
475 return Err(SqliteError::InvalidData(
476 "events namespace must be non-empty".to_string(),
477 ));
478 }
479 if !self.is_read_only() {
480 let writer = self.pool.try_writer()?;
481 event::ensure_events_schema(writer.conn())?;
482 }
483
484 Ok(Arc::new(event::SqlEventStore::new_scoped(
485 Arc::clone(&self.pool),
486 self.is_file_backed,
487 namespace.trim().to_string(),
488 )))
489 }
490
491 pub fn agents(&self) -> Result<Arc<dyn khive_storage::AgentStore>, SqliteError> {
496 if !self.is_read_only() {
497 let writer = self.pool.try_writer()?;
498 agents::ensure_agents_schema(writer.conn())?;
499 }
500
501 Ok(Arc::new(agents::SqlAgentStore::new(
502 Arc::clone(&self.pool),
503 self.is_file_backed,
504 )))
505 }
506
507 pub fn vectors(
513 &self,
514 model_key: &str,
515 embedding_model: &str,
516 dimensions: usize,
517 ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
518 self.vectors_for_namespace(model_key, embedding_model, dimensions, "local")
519 }
520
521 pub fn vectors_for_namespace(
531 &self,
532 model_key: &str,
533 embedding_model: &str,
534 dimensions: usize,
535 namespace: &str,
536 ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
537 if model_key.is_empty()
538 || !model_key
539 .chars()
540 .all(|c| c.is_ascii_alphanumeric() || c == '_')
541 {
542 return Err(SqliteError::InvalidData(format!(
543 "invalid model_key '{}': must be non-empty and contain only \
544 alphanumeric/underscore characters",
545 model_key
546 )));
547 }
548 if namespace.trim().is_empty() {
549 return Err(SqliteError::InvalidData(
550 "vector store namespace must be non-empty".to_string(),
551 ));
552 }
553
554 crate::extension::ensure_extensions_loaded();
556
557 let table = format!("vec_{}", model_key);
558
559 if self.is_read_only() {
560 let reader = self.pool.reader()?;
564 if !sqlite_table_exists(reader.conn(), &table)? {
565 return Err(SqliteError::InvalidData(format!(
566 "read-only database has no vector table '{table}'; create and populate it in \
567 a writable copy before opening the snapshot"
568 )));
569 }
570 validate_vector_table_columns(reader.conn(), &table)?;
571 drop(reader);
572 return Ok(Arc::new(vectors::SqliteVecStore::new(
573 Arc::clone(&self.pool),
574 self.is_file_backed,
575 model_key.to_string(),
576 embedding_model.to_string(),
577 dimensions,
578 namespace.trim().to_string(),
579 )?));
580 }
581
582 let writer = self.pool.try_writer()?;
583
584 let table_exists = sqlite_table_exists(writer.conn(), &table)?;
591
592 if table_exists {
593 validate_vector_table_columns(writer.conn(), &table)?;
599 }
600
601 writer
610 .conn()
611 .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
612
613 writer
617 .conn()
618 .execute_batch(crate::migrations::ANN_WRITE_LOG_DDL)?;
619 writer
620 .conn()
621 .execute_batch(crate::migrations::ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL)?;
622 writer
623 .conn()
624 .execute_batch(crate::migrations::ANN_CONSUMER_PENDING_DDL)?;
625
626 let ddl = format!(
629 "CREATE VIRTUAL TABLE IF NOT EXISTS vec_{} USING vec0(\
630 subject_id TEXT PRIMARY KEY, \
631 namespace TEXT NOT NULL, \
632 kind TEXT NOT NULL, \
633 field TEXT NOT NULL, \
634 embedding_model TEXT NOT NULL, \
635 embedding float[{}] distance_metric=cosine\
636 )",
637 model_key, dimensions
638 );
639 writer.conn().execute_batch(&ddl)?;
640
641 Ok(Arc::new(vectors::SqliteVecStore::new(
642 Arc::clone(&self.pool),
643 self.is_file_backed,
644 model_key.to_string(),
645 embedding_model.to_string(),
646 dimensions,
647 namespace.trim().to_string(),
648 )?))
649 }
650
651 pub fn register_embedding_model(
656 &self,
657 engine_name: &str,
658 model_id: &str,
659 key_version: &str,
660 dimensions: u32,
661 ) -> Result<(), SqliteError> {
662 let writer = self.pool.try_writer()?;
663 writer
664 .conn()
665 .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
666
667 let now = chrono::Utc::now().timestamp_micros();
668 let canonical_key =
669 format!("{engine_name}:{model_id}:{key_version}:{dimensions}").into_bytes();
670 let id = uuid::Uuid::new_v4();
671 writer.conn().execute(
672 "INSERT INTO _embedding_models \
673 (id, engine_name, model_id, key_version, dim, output_dim, status, \
674 activated_at, superseded_at, superseded_by, canonical_key, created_at) \
675 VALUES (?1, ?2, ?3, ?4, ?5, NULL, 'active', ?6, NULL, NULL, ?7, ?8) \
676 ON CONFLICT(canonical_key) DO UPDATE SET \
677 status = 'active', \
678 activated_at = COALESCE(_embedding_models.activated_at, excluded.activated_at)",
679 rusqlite::params![
680 id.as_bytes().as_slice(),
681 engine_name,
682 model_id,
683 key_version,
684 dimensions as i64,
685 now,
686 canonical_key,
687 now,
688 ],
689 )?;
690 Ok(())
691 }
692
693 pub fn sparse(
697 &self,
698 model_key: &str,
699 ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
700 self.sparse_for_namespace(model_key, "local")
701 }
702
703 pub fn sparse_for_namespace(
707 &self,
708 model_key: &str,
709 namespace: &str,
710 ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
711 if model_key.is_empty()
712 || !model_key
713 .chars()
714 .all(|c| c.is_ascii_alphanumeric() || c == '_')
715 {
716 return Err(SqliteError::InvalidData(format!(
717 "invalid model_key '{}': must be non-empty and contain only alphanumeric/underscore characters",
718 model_key
719 )));
720 }
721 if namespace.trim().is_empty() {
722 return Err(SqliteError::InvalidData(
723 "sparse store namespace must be non-empty".to_string(),
724 ));
725 }
726
727 if self.is_read_only() {
728 let table = format!("sparse_{model_key}");
729 let reader = self.pool.reader()?;
730 if !sqlite_table_exists(reader.conn(), &table)? {
731 return Err(SqliteError::InvalidData(format!(
732 "read-only database has no sparse table '{table}'; create and populate it in \
733 a writable copy before opening the snapshot"
734 )));
735 }
736 } else {
737 let writer = self.pool.try_writer()?;
738 sparse::ensure_sparse_schema(writer.conn(), model_key)
739 .map_err(SqliteError::Rusqlite)?;
740 }
741
742 Ok(Arc::new(sparse::SqliteSparseStore::new(
743 Arc::clone(&self.pool),
744 self.is_file_backed,
745 model_key.to_string(),
746 namespace.trim().to_string(),
747 )?))
748 }
749
750 pub fn text(&self, table_key: &str) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
757 self.text_with_tokenizer(table_key, "trigram")
758 }
759
760 pub fn text_with_tokenizer(
768 &self,
769 table_key: &str,
770 tokenizer: &str,
771 ) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
772 if table_key.is_empty()
773 || !table_key
774 .chars()
775 .all(|c| c.is_ascii_alphanumeric() || c == '_')
776 {
777 return Err(SqliteError::InvalidData(format!(
778 "invalid table_key '{}': must be non-empty and contain only \
779 alphanumeric/underscore characters",
780 table_key
781 )));
782 }
783 if tokenizer.is_empty()
784 || !tokenizer
785 .chars()
786 .all(|c| c.is_ascii_alphanumeric() || c == '_')
787 {
788 return Err(SqliteError::InvalidData(format!(
789 "invalid tokenizer '{}': must be non-empty and contain only \
790 alphanumeric/underscore characters",
791 tokenizer
792 )));
793 }
794
795 let ddl = format!(
796 "CREATE VIRTUAL TABLE IF NOT EXISTS fts_{} USING fts5(\
797 subject_id UNINDEXED, \
798 kind UNINDEXED, \
799 title, \
800 body, \
801 tags UNINDEXED, \
802 namespace UNINDEXED, \
803 metadata UNINDEXED, \
804 updated_at UNINDEXED, \
805 tokenize = '{}'\
806 )",
807 table_key, tokenizer
808 );
809 if self.is_read_only() {
810 let table = format!("fts_{table_key}");
811 let reader = self.pool.reader()?;
812 if !sqlite_table_exists(reader.conn(), &table)? {
813 return Err(SqliteError::InvalidData(format!(
814 "read-only database has no text-search table '{table}'; create and populate \
815 it in a writable copy before opening the snapshot"
816 )));
817 }
818 } else {
819 let writer = self.pool.try_writer()?;
820 writer.conn().execute_batch(&ddl)?;
821 }
822
823 Ok(Arc::new(text::Fts5TextSearch::new(
824 Arc::clone(&self.pool),
825 self.is_file_backed,
826 table_key.to_string(),
827 )))
828 }
829
830 pub fn blob_store(
838 &self,
839 config_root: Option<&Path>,
840 floor_bytes: Option<u64>,
841 ) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
842 let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
843 let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
844 Ok(Arc::new(blob::FsBlobStore::new(root, floor)?))
845 }
846
847 pub fn blob_store_read_only(
852 &self,
853 config_root: Option<&Path>,
854 floor_bytes: Option<u64>,
855 ) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
856 let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
857 let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
858 Ok(Arc::new(blob::FsBlobStore::open_existing(root, floor)?))
859 }
860
861 pub fn is_file_backed(&self) -> bool {
863 self.is_file_backed
864 }
865
866 pub fn is_read_only(&self) -> bool {
869 self.pool.config().read_only
870 }
871
872 pub fn data_dir(&self) -> Option<std::path::PathBuf> {
875 self.path.as_ref()?.parent().map(|p| p.to_path_buf())
876 }
877
878 pub fn ann_root(&self) -> Option<std::path::PathBuf> {
886 ann_root_for(self.path.as_ref()?)
887 }
888
889 pub fn pool(&self) -> &ConnectionPool {
891 &self.pool
892 }
893
894 pub fn pool_arc(&self) -> Arc<ConnectionPool> {
896 Arc::clone(&self.pool)
897 }
898}
899
900fn ann_root_for(path: &std::path::Path) -> Option<std::path::PathBuf> {
905 let mut file = path.file_name()?.to_os_string();
906 file.push(".ann");
907 path.parent().map(|p| p.join(file))
908}
909
910#[cfg(test)]
911mod tests {
912 use super::*;
913 use khive_storage::types::{SqlStatement, SqlValue};
914
915 #[cfg(unix)]
921 fn freeze_snapshot_sidecars(path: &std::path::Path) {
922 use std::os::unix::fs::PermissionsExt;
923 for suffix in ["-wal", "-shm"] {
924 let mut name = path.file_name().expect("db file name").to_os_string();
925 name.push(suffix);
926 let sidecar = path.parent().expect("db parent dir").join(name);
927 if sidecar.exists() {
928 let mut permissions = std::fs::metadata(&sidecar)
929 .expect("sidecar metadata")
930 .permissions();
931 permissions.set_mode(0o444);
932 std::fs::set_permissions(&sidecar, permissions).expect("freeze sidecar");
933 }
934 }
935 }
936
937 #[cfg(unix)]
938 #[tokio::test]
939 async fn sqlite_detects_chmod_read_only_snapshot_and_core_reads_succeed() {
940 use std::os::unix::fs::PermissionsExt;
941
942 let dir = tempfile::tempdir().unwrap();
943 let path = dir.path().join("chmod_snapshot.db");
944 {
945 let writable = StorageBackend::sqlite(&path).expect("create writable database");
946 writable
947 .prepare_core_schema()
948 .expect("migrate writable snapshot source");
949 }
950
951 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
952 permissions.set_mode(0o444);
953 std::fs::set_permissions(&path, permissions).unwrap();
954 freeze_snapshot_sidecars(&path);
955
956 let read_only = StorageBackend::sqlite(&path).expect("auto-detect read-only mode");
957 assert!(read_only.is_read_only());
958 assert_eq!(
959 read_only.pool().config().write_queue_enabled,
960 Some(false),
961 "read-only boot must not attempt to spawn a writer task"
962 );
963 assert!(read_only
964 .pool()
965 .writer_task_handle()
966 .expect("disabled writer task is a valid configuration")
967 .is_none());
968 read_only
969 .prepare_core_schema()
970 .expect("current snapshot validates without migration writes");
971
972 let entities = read_only.entities().expect("entity store opens read-only");
973 let graph = read_only.graph().expect("graph store opens read-only");
974 let notes = read_only.notes().expect("note store opens read-only");
975 let events = read_only.events().expect("event store opens read-only");
976 assert_eq!(
977 entities
978 .count_entities("local", khive_storage::EntityFilter::default())
979 .await
980 .unwrap(),
981 0
982 );
983 assert_eq!(
984 graph
985 .count_edges(khive_storage::types::EdgeFilter::default())
986 .await
987 .unwrap(),
988 0
989 );
990 assert_eq!(notes.count_notes("local", None).await.unwrap(), 0);
991 assert_eq!(
992 events
993 .count_events(khive_storage::EventFilter::default())
994 .await
995 .unwrap(),
996 0
997 );
998 assert_eq!(
999 read_only.notes_seq_repair_run_count(),
1000 0,
1001 "read-only store acquisition must not run the DML repair"
1002 );
1003 }
1004
1005 #[test]
1006 fn memory_backend_creates_successfully() {
1007 let backend = StorageBackend::memory().expect("memory backend should create");
1008 assert!(!backend.is_file_backed());
1009 }
1010
1011 #[test]
1012 fn file_backend_creates_successfully() {
1013 let dir = tempfile::tempdir().unwrap();
1014 let path = dir.path().join("test.db");
1015 let backend = StorageBackend::sqlite(&path).expect("file backend should create");
1016 assert!(backend.is_file_backed());
1017 assert!(path.exists());
1018 }
1019
1020 #[test]
1021 fn data_dir_returns_none_for_memory_backend() {
1022 let backend = StorageBackend::memory().expect("memory backend");
1023 assert!(backend.data_dir().is_none());
1024 }
1025
1026 #[test]
1027 fn data_dir_returns_parent_dir_for_file_backend() {
1028 let dir = tempfile::tempdir().unwrap();
1029 let path = dir.path().join("data.db");
1030 let backend = StorageBackend::sqlite(&path).expect("file backend");
1031 let got = backend.data_dir().expect("file backend must return Some");
1032 assert_eq!(got, dir.path());
1033 }
1034
1035 #[test]
1036 fn ann_root_is_database_scoped_sibling_dir() {
1037 let dir = tempfile::tempdir().unwrap();
1038 let path = dir.path().join("data.db");
1039 let backend = StorageBackend::sqlite(&path).expect("file backend");
1040 let got = backend.ann_root().expect("file backend must return Some");
1041 assert_eq!(got, dir.path().join("data.db.ann"));
1042 assert!(StorageBackend::memory().unwrap().ann_root().is_none());
1043 }
1044
1045 #[cfg(unix)]
1051 #[test]
1052 fn ann_root_distinct_for_non_utf8_filenames() {
1053 use std::os::unix::ffi::OsStrExt;
1054 let path_a = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xff.db"));
1055 let path_b = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xfe.db"));
1056 let root_a = ann_root_for(&path_a).expect("Some for a file path");
1057 let root_b = ann_root_for(&path_b).expect("Some for a file path");
1058 assert_ne!(
1059 root_a, root_b,
1060 "distinct database files must map to distinct ANN roots"
1061 );
1062 }
1063
1064 #[tokio::test]
1065 async fn sql_access_memory_roundtrip() {
1066 let backend = StorageBackend::memory().unwrap();
1067 let sql = backend.sql();
1068
1069 let mut writer = sql.writer().await.unwrap();
1070 writer
1071 .execute_script(
1072 "CREATE TABLE test_rt (id TEXT PRIMARY KEY, value INTEGER NOT NULL)".into(),
1073 )
1074 .await
1075 .unwrap();
1076
1077 let affected = writer
1078 .execute(SqlStatement {
1079 sql: "INSERT INTO test_rt (id, value) VALUES (?1, ?2)".into(),
1080 params: vec![SqlValue::Text("row1".into()), SqlValue::Integer(42)],
1081 label: None,
1082 })
1083 .await
1084 .unwrap();
1085 assert_eq!(affected, 1);
1086
1087 let mut reader = sql.reader().await.unwrap();
1088 let row = reader
1089 .query_row(SqlStatement {
1090 sql: "SELECT id, value FROM test_rt WHERE id = ?1".into(),
1091 params: vec![SqlValue::Text("row1".into())],
1092 label: None,
1093 })
1094 .await
1095 .unwrap();
1096
1097 let row = row.expect("should find the inserted row");
1098 assert_eq!(row.columns.len(), 2);
1099 match &row.columns[0].value {
1100 SqlValue::Text(s) => assert_eq!(s, "row1"),
1101 other => panic!("expected Text, got {other:?}"),
1102 }
1103 match &row.columns[1].value {
1104 SqlValue::Integer(v) => assert_eq!(*v, 42),
1105 other => panic!("expected Integer, got {other:?}"),
1106 }
1107 }
1108
1109 #[tokio::test]
1110 async fn sql_access_file_roundtrip() {
1111 let dir = tempfile::tempdir().unwrap();
1112 let path = dir.path().join("test_roundtrip.db");
1113 let backend = StorageBackend::sqlite(&path).unwrap();
1114 let sql = backend.sql();
1115
1116 let mut writer = sql.writer().await.unwrap();
1117 writer
1118 .execute_script("CREATE TABLE test_f (k TEXT PRIMARY KEY, v TEXT)".into())
1119 .await
1120 .unwrap();
1121 writer
1122 .execute(SqlStatement {
1123 sql: "INSERT INTO test_f (k, v) VALUES (?1, ?2)".into(),
1124 params: vec![
1125 SqlValue::Text("hello".into()),
1126 SqlValue::Text("world".into()),
1127 ],
1128 label: None,
1129 })
1130 .await
1131 .unwrap();
1132
1133 let mut reader = sql.reader().await.unwrap();
1134 let rows = reader
1135 .query_all(SqlStatement {
1136 sql: "SELECT k, v FROM test_f".into(),
1137 params: vec![],
1138 label: None,
1139 })
1140 .await
1141 .unwrap();
1142 assert_eq!(rows.len(), 1);
1143 match &rows[0].columns[1].value {
1144 SqlValue::Text(s) => assert_eq!(s, "world"),
1145 other => panic!("expected Text, got {other:?}"),
1146 }
1147 }
1148
1149 #[test]
1150 fn sqlite_read_only_missing_path_does_not_create_file() {
1151 let dir = tempfile::tempdir().unwrap();
1152 let path = dir.path().join("missing_ro.db");
1153 assert!(!path.exists());
1154
1155 let result = StorageBackend::sqlite_read_only(&path);
1156 assert!(
1157 result.is_err(),
1158 "opening a missing path read-only must fail"
1159 );
1160 assert!(
1161 !path.exists(),
1162 "opening a missing path read-only must not create the file"
1163 );
1164 }
1165
1166 #[test]
1167 fn sqlite_read_only_sparse_store_requires_existing_table_without_writer_acquisition() {
1168 let dir = tempfile::tempdir().unwrap();
1169 let path = dir.path().join("ro_sparse_tables.db");
1170 {
1171 let writable = StorageBackend::sqlite(&path).unwrap();
1172 writable
1173 .prepare_core_schema()
1174 .expect("migrate snapshot source");
1175 writable
1176 .sparse("present")
1177 .expect("create the optional sparse table while writable");
1178 }
1179 #[cfg(unix)]
1180 freeze_snapshot_sidecars(&path);
1181
1182 let read_only = StorageBackend::sqlite_read_only(&path).unwrap();
1183 read_only
1184 .prepare_core_schema()
1185 .expect("validate exact current migration ledger");
1186 read_only
1187 .sparse("present")
1188 .expect("an existing sparse table must open read-only");
1189 let missing = match read_only.sparse("missing") {
1190 Ok(_) => panic!("a missing sparse table must fail during store acquisition"),
1191 Err(error) => error,
1192 };
1193 assert!(
1194 missing.to_string().contains("sparse_missing"),
1195 "the diagnostic must name the absent table: {missing}"
1196 );
1197 assert_eq!(
1198 read_only.pool().writer_acquisition_snapshot(),
1199 crate::pool::WriterAcquisitionSnapshot::default(),
1200 "construction, exact-ledger validation, and optional sparse-table inspection must \
1201 use reader connections only"
1202 );
1203 }
1204
1205 #[test]
1206 fn sqlite_read_only_text_store_requires_existing_table_without_writer_acquisition() {
1207 let dir = tempfile::tempdir().unwrap();
1208 let path = dir.path().join("ro_text_tables.db");
1209 {
1210 let writable = StorageBackend::sqlite(&path).unwrap();
1211 writable
1212 .prepare_core_schema()
1213 .expect("migrate snapshot source");
1214 writable
1215 .text("present")
1216 .expect("create the optional FTS table while writable");
1217 }
1218 #[cfg(unix)]
1219 freeze_snapshot_sidecars(&path);
1220
1221 let read_only = StorageBackend::sqlite_read_only(&path).unwrap();
1222 read_only
1223 .prepare_core_schema()
1224 .expect("validate exact current migration ledger");
1225 read_only
1226 .text("present")
1227 .expect("an existing FTS table must open read-only");
1228 let missing = match read_only.text("missing") {
1229 Ok(_) => panic!("a missing FTS table must fail during store acquisition"),
1230 Err(error) => error,
1231 };
1232 assert!(
1233 missing.to_string().contains("fts_missing"),
1234 "the diagnostic must name the absent table: {missing}"
1235 );
1236 assert_eq!(
1237 read_only.pool().writer_acquisition_snapshot(),
1238 crate::pool::WriterAcquisitionSnapshot::default(),
1239 "construction, exact-ledger validation, and optional FTS inspection must use reader \
1240 connections only"
1241 );
1242 }
1243
1244 #[cfg(feature = "vectors")]
1245 #[test]
1246 fn sqlite_read_only_vector_store_schema_check_uses_no_writer_acquisition() {
1247 let dir = tempfile::tempdir().unwrap();
1248 let path = dir.path().join("ro_vector_tables.db");
1249 {
1250 let writable = StorageBackend::sqlite(&path).unwrap();
1251 writable
1252 .prepare_core_schema()
1253 .expect("migrate snapshot source");
1254 writable
1255 .vectors("present", "present", 3)
1256 .expect("create the optional vector table while writable");
1257 }
1258 #[cfg(unix)]
1259 freeze_snapshot_sidecars(&path);
1260
1261 let read_only = StorageBackend::sqlite_read_only(&path).unwrap();
1262 read_only
1263 .prepare_core_schema()
1264 .expect("validate exact current migration ledger");
1265 read_only
1266 .vectors("present", "present", 3)
1267 .expect("an existing vector table must open read-only");
1268 assert!(
1269 read_only.vectors("missing", "missing", 3).is_err(),
1270 "a missing vector table must fail during store acquisition"
1271 );
1272 assert_eq!(
1273 read_only.pool().writer_acquisition_snapshot(),
1274 crate::pool::WriterAcquisitionSnapshot::default(),
1275 "construction, exact-ledger validation, and optional vector inspection must use \
1276 reader connections only"
1277 );
1278 }
1279
1280 #[tokio::test]
1281 async fn sqlite_read_only_sql_writer_rejects_ddl_and_insert() {
1282 let dir = tempfile::tempdir().unwrap();
1283 let path = dir.path().join("ro_writer.db");
1284
1285 {
1287 let writable = StorageBackend::sqlite(&path).unwrap();
1288 let sql = writable.sql();
1289 let mut writer = sql.writer().await.unwrap();
1290 writer
1291 .execute_script("CREATE TABLE ro_existing (id INTEGER PRIMARY KEY)".into())
1292 .await
1293 .unwrap();
1294 }
1295 #[cfg(unix)]
1296 freeze_snapshot_sidecars(&path);
1297
1298 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1299 let sql = ro.sql();
1300
1301 let writer_result = sql.writer().await;
1303 assert!(
1304 writer_result.is_err(),
1305 "sql().writer() must be rejected on a read-only backend"
1306 );
1307 }
1308
1309 #[tokio::test]
1310 #[cfg(feature = "vectors")]
1311 async fn vectors_roundtrip_via_public_api() {
1312 let backend = StorageBackend::memory().unwrap();
1313 let store = backend.vectors("test_api", "test_api", 3).unwrap();
1314
1315 let id = uuid::Uuid::new_v4();
1316 store
1317 .insert(
1318 id,
1319 khive_types::SubstrateKind::Entity,
1320 "local",
1321 "content",
1322 vec![vec![1.0, 0.0, 0.0]],
1323 )
1324 .await
1325 .unwrap();
1326
1327 let hits = store
1328 .search(khive_storage::types::VectorSearchRequest {
1329 query_vectors: vec![vec![1.0, 0.0, 0.0]],
1330 top_k: 1,
1331 namespace: None,
1332 kind: None,
1333 embedding_model: None,
1334 filter: None,
1335 backend_hints: None,
1336 })
1337 .await
1338 .unwrap();
1339
1340 assert_eq!(hits.len(), 1);
1341 assert_eq!(hits[0].subject_id, id);
1342 assert!(hits[0].score.to_f64() > 0.99);
1343 }
1344
1345 #[tokio::test]
1346 #[cfg(feature = "vectors")]
1347 async fn vectors_creates_table_idempotently() {
1348 let backend = StorageBackend::memory().unwrap();
1349
1350 let store1 = backend.vectors("idempotent", "idempotent", 3).unwrap();
1351 let store2 = backend.vectors("idempotent", "idempotent", 3).unwrap();
1352
1353 let id = uuid::Uuid::new_v4();
1354 store1
1355 .insert(
1356 id,
1357 khive_types::SubstrateKind::Entity,
1358 "local",
1359 "content",
1360 vec![vec![1.0, 0.0, 0.0]],
1361 )
1362 .await
1363 .unwrap();
1364
1365 let count = store2.count().await.unwrap();
1366 assert_eq!(count, 1);
1367 }
1368
1369 #[tokio::test]
1370 async fn text_roundtrip_via_public_api() {
1371 let backend = StorageBackend::memory().unwrap();
1372 let store = backend.text("test_api").unwrap();
1373
1374 let id = uuid::Uuid::new_v4();
1375 let doc = khive_storage::types::TextDocument {
1376 subject_id: id,
1377 kind: khive_types::SubstrateKind::Entity,
1378 title: Some("Test Title".to_string()),
1379 body: "This is a searchable document about Rust.".to_string(),
1380 tags: vec!["rust".to_string()],
1381 namespace: "test_ns".to_string(),
1382 metadata: None,
1383 updated_at: chrono::Utc::now(),
1384 };
1385 store.upsert_document(doc).await.unwrap();
1386
1387 let hits = store
1388 .search(khive_storage::types::TextSearchRequest {
1389 query: "Rust".to_string(),
1390 mode: khive_storage::types::TextQueryMode::Plain,
1391 filter: Some(khive_storage::types::TextFilter {
1392 namespaces: vec!["test_ns".to_string()],
1393 ..Default::default()
1394 }),
1395 top_k: 1,
1396 snippet_chars: 64,
1397 })
1398 .await
1399 .unwrap();
1400
1401 assert_eq!(hits.len(), 1);
1402 assert_eq!(hits[0].subject_id, id);
1403 assert!(hits[0].score.to_f64() > 0.0);
1404 }
1405
1406 #[tokio::test]
1407 async fn text_creates_table_idempotently() {
1408 let backend = StorageBackend::memory().unwrap();
1409
1410 let store1 = backend.text("idempotent_fts").unwrap();
1411 let store2 = backend.text("idempotent_fts").unwrap();
1412
1413 let id = uuid::Uuid::new_v4();
1414 let doc = khive_storage::types::TextDocument {
1415 subject_id: id,
1416 kind: khive_types::SubstrateKind::Note,
1417 title: None,
1418 body: "Hello world.".to_string(),
1419 tags: vec![],
1420 namespace: "test_ns".to_string(),
1421 metadata: None,
1422 updated_at: chrono::Utc::now(),
1423 };
1424 store1.upsert_document(doc).await.unwrap();
1425
1426 let count = store2
1427 .count(khive_storage::types::TextFilter {
1428 namespaces: vec!["test_ns".to_string()],
1429 ..Default::default()
1430 })
1431 .await
1432 .unwrap();
1433 assert_eq!(count, 1);
1434 }
1435
1436 #[test]
1437 fn invalid_model_key_rejected() {
1438 let backend = StorageBackend::memory().unwrap();
1439 assert!(backend.vectors("bad key!", "bad key!", 3).is_err());
1440 assert!(backend.vectors("", "", 3).is_err());
1441 }
1442
1443 #[test]
1444 fn invalid_table_key_rejected() {
1445 let backend = StorageBackend::memory().unwrap();
1446 assert!(backend.text("bad key!").is_err());
1447 assert!(backend.text("").is_err());
1448 }
1449
1450 #[tokio::test]
1451 async fn sqlite_read_only_graph_store_rejects_upsert_edge() {
1452 use khive_storage::types::Edge;
1453 use khive_types::EdgeRelation;
1454
1455 let dir = tempfile::tempdir().unwrap();
1456 let path = dir.path().join("ro_graph.db");
1457
1458 {
1460 let writable = StorageBackend::sqlite(&path).unwrap();
1461 writable.graph().unwrap();
1462 }
1463 #[cfg(unix)]
1464 freeze_snapshot_sidecars(&path);
1465
1466 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1467 let store = match ro.graph() {
1468 Ok(store) => store,
1469 Err(_) => return,
1472 };
1473
1474 let now = chrono::Utc::now();
1475 let edge = Edge {
1476 id: uuid::Uuid::new_v4().into(),
1477 namespace: "local".to_string(),
1478 source_id: uuid::Uuid::new_v4(),
1479 target_id: uuid::Uuid::new_v4(),
1480 relation: EdgeRelation::Extends,
1481 weight: 0.8,
1482 created_at: now,
1483 updated_at: now,
1484 deleted_at: None,
1485 metadata: None,
1486 target_backend: None,
1487 };
1488
1489 let result = store.upsert_edge(edge).await;
1490 assert!(
1491 result.is_err(),
1492 "upsert_edge on a read-only backend must reject, not silently no-op"
1493 );
1494 }
1495
1496 #[tokio::test]
1497 async fn sqlite_read_only_event_store_rejects_append_event() {
1498 use khive_types::{EventKind, EventOutcome, SubstrateKind};
1499
1500 let dir = tempfile::tempdir().unwrap();
1501 let path = dir.path().join("ro_events.db");
1502
1503 {
1504 let writable = StorageBackend::sqlite(&path).unwrap();
1505 writable.events().unwrap();
1506 }
1507 #[cfg(unix)]
1508 freeze_snapshot_sidecars(&path);
1509
1510 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1511 let store = match ro.events() {
1512 Ok(store) => store,
1513 Err(_) => return,
1514 };
1515
1516 let event = khive_storage::event::Event::new(
1517 "local",
1518 "test.verb",
1519 EventKind::Audit,
1520 SubstrateKind::Entity,
1521 "test-actor",
1522 )
1523 .with_outcome(EventOutcome::Success);
1524
1525 let result = store.append_event(event).await;
1526 assert!(
1527 result.is_err(),
1528 "append_event on a read-only backend must reject, not silently no-op"
1529 );
1530 }
1531
1532 #[tokio::test]
1533 async fn sqlite_read_only_text_store_rejects_upsert_document() {
1534 use khive_storage::types::TextDocument;
1535 use khive_types::SubstrateKind;
1536
1537 let dir = tempfile::tempdir().unwrap();
1538 let path = dir.path().join("ro_text.db");
1539
1540 {
1541 let writable = StorageBackend::sqlite(&path).unwrap();
1542 writable.text("ro_test").unwrap();
1543 }
1544 #[cfg(unix)]
1545 freeze_snapshot_sidecars(&path);
1546
1547 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1548 let store = match ro.text("ro_test") {
1549 Ok(store) => store,
1550 Err(_) => return,
1551 };
1552
1553 let doc = TextDocument {
1554 subject_id: uuid::Uuid::new_v4(),
1555 kind: SubstrateKind::Entity,
1556 title: Some("Title".to_string()),
1557 body: "Body text.".to_string(),
1558 tags: vec![],
1559 namespace: "local".to_string(),
1560 metadata: None,
1561 updated_at: chrono::Utc::now(),
1562 };
1563
1564 let result = store.upsert_document(doc).await;
1565 assert!(
1566 result.is_err(),
1567 "upsert_document on a read-only backend must reject, not silently no-op"
1568 );
1569 }
1570
1571 #[tokio::test]
1572 async fn blob_store_roundtrip_via_public_api() {
1573 let dir = tempfile::tempdir().unwrap();
1574 let path = dir.path().join("blob_backend.db");
1575 let backend = StorageBackend::sqlite(&path).unwrap();
1576
1577 let store = backend.blob_store(None, Some(0)).unwrap();
1581 let bytes = b"backend-level blob roundtrip".to_vec();
1582 let content_ref = store.put(bytes.clone()).await.unwrap();
1583 assert_eq!(
1584 store
1585 .get_bounded_verified(&content_ref, bytes.len() as u64)
1586 .await
1587 .unwrap(),
1588 bytes
1589 );
1590 }
1591
1592 #[test]
1593 fn blob_store_defaults_root_beside_db_file() {
1594 let dir = tempfile::tempdir().unwrap();
1595 let path = dir.path().join("blob_default.db");
1596 let backend = StorageBackend::sqlite(&path).unwrap();
1597
1598 let _store = backend.blob_store(None, None).unwrap();
1602 assert!(
1603 dir.path().join("blobs").is_dir(),
1604 "default root must be created beside the database file"
1605 );
1606 }
1607
1608 #[test]
1609 fn blob_store_errors_for_in_memory_backend_with_no_override() {
1610 let backend = StorageBackend::memory().unwrap();
1611 assert!(backend.blob_store(None, None).is_err());
1612 }
1613
1614 #[test]
1615 fn blob_store_accepts_explicit_root_for_in_memory_backend() {
1616 let dir = tempfile::tempdir().unwrap();
1617 let backend = StorageBackend::memory().unwrap();
1618 let store = backend.blob_store(Some(dir.path()), None);
1619 assert!(store.is_ok());
1620 }
1621
1622 #[test]
1623 fn apply_schema_runs_migrations_idempotently() {
1624 static MIGRATIONS: &[crate::migrations::Migration] = &[crate::migrations::Migration {
1625 id: "001_init",
1626 up_sql: "CREATE TABLE IF NOT EXISTS schema_test (id TEXT PRIMARY KEY);",
1627 down_sql: None,
1628 is_already_applied: None,
1629 }];
1630 let plan = crate::migrations::ServiceSchemaPlan {
1631 service: "schema_test_svc",
1632 sqlite: MIGRATIONS,
1633 postgres: &[],
1634 };
1635
1636 let backend = StorageBackend::memory().unwrap();
1637 backend.apply_schema(&plan).unwrap();
1638 backend.apply_schema(&plan).unwrap();
1639
1640 let reader = backend.pool().reader().unwrap();
1641 let count: i64 = reader
1642 .conn()
1643 .query_row(
1644 "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='schema_test'",
1645 [],
1646 |row| row.get(0),
1647 )
1648 .unwrap();
1649 assert_eq!(count, 1);
1650 }
1651
1652 #[test]
1653 fn pack_ddl_plan_rolls_back_all_statements_on_failure() {
1654 let backend = StorageBackend::memory().unwrap();
1655 let error = backend
1656 .apply_pack_ddl_statements(&[
1657 "CREATE TABLE IF NOT EXISTS pack_schema_first (id INTEGER PRIMARY KEY)",
1658 "CREATE INDEX IF NOT EXISTS pack_schema_second ON pack_schema_missing(id)",
1659 ])
1660 .unwrap_err();
1661
1662 assert!(
1663 error.to_string().contains("pack_schema_missing"),
1664 "schema-plan error must retain the failing SQLite diagnostic: {error}"
1665 );
1666
1667 let reader = backend.pool().reader().unwrap();
1668 let visible_objects: i64 = reader
1669 .conn()
1670 .query_row(
1671 "SELECT COUNT(*) FROM sqlite_master \
1672 WHERE name IN ('pack_schema_first', 'pack_schema_second')",
1673 [],
1674 |row| row.get(0),
1675 )
1676 .unwrap();
1677 assert_eq!(visible_objects, 0);
1678 }
1679
1680 #[test]
1681 fn pack_ddl_plan_applies_all_statements_idempotently() {
1682 const PLAN: &[&str] = &[
1683 "CREATE TABLE IF NOT EXISTS pack_schema_success (id INTEGER PRIMARY KEY, value TEXT)",
1684 "CREATE INDEX IF NOT EXISTS pack_schema_success_value_idx \
1685 ON pack_schema_success(value)",
1686 ];
1687
1688 let backend = StorageBackend::memory().unwrap();
1689 backend.apply_pack_ddl_statements(PLAN).unwrap();
1690 backend.apply_pack_ddl_statements(PLAN).unwrap();
1691
1692 let reader = backend.pool().reader().unwrap();
1693 let visible_objects: i64 = reader
1694 .conn()
1695 .query_row(
1696 "SELECT COUNT(*) FROM sqlite_master \
1697 WHERE name IN ('pack_schema_success', 'pack_schema_success_value_idx')",
1698 [],
1699 |row| row.get(0),
1700 )
1701 .unwrap();
1702 assert_eq!(visible_objects, 2);
1703 }
1704
1705 fn issue_1029_pool(write_queue_enabled: bool) -> (tempfile::TempDir, StorageBackend) {
1714 let dir = tempfile::tempdir().unwrap();
1715 let path = dir.path().join("issue_1029.db");
1716 let config = crate::pool::PoolConfig {
1717 path: Some(path.clone()),
1718 busy_timeout: std::time::Duration::from_millis(200),
1719 write_queue_enabled: Some(write_queue_enabled),
1720 ..crate::pool::PoolConfig::default()
1721 };
1722 let pool = ConnectionPool::new(config).expect("fresh tenant-shaped pool should open");
1723 let backend = StorageBackend {
1724 pool: Arc::new(pool),
1725 is_file_backed: true,
1726 path: Some(path),
1727 notes_seq_repair_runs: AtomicUsize::new(0),
1728 };
1729 (dir, backend)
1730 }
1731
1732 async fn issue_1029_create_entity_shaped_sequence(
1733 backend: &StorageBackend,
1734 ) -> Result<(), String> {
1735 let entities = backend
1736 .entities_for_namespace("tenant_ns")
1737 .map_err(|e| format!("entities_for_namespace: {e}"))?;
1738 let entity = khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029Repro");
1739 let entity_id = entity.id;
1740 entities
1741 .upsert_entity(entity)
1742 .await
1743 .map_err(|e| format!("upsert_entity: {e}"))?;
1744
1745 let text = backend.text("entities").map_err(|e| format!("text: {e}"))?;
1746 let doc = khive_storage::types::TextDocument {
1747 subject_id: entity_id,
1748 kind: khive_types::SubstrateKind::Entity,
1749 title: Some("Issue1029Repro".to_string()),
1750 body: "issue 1029 repro body".to_string(),
1751 tags: vec![],
1752 namespace: "tenant_ns".to_string(),
1753 metadata: None,
1754 updated_at: chrono::Utc::now(),
1755 };
1756 text.upsert_document(doc)
1757 .await
1758 .map_err(|e| format!("fts_upsert: {e}"))
1759 }
1760
1761 #[tokio::test]
1767 async fn issue_1029_create_entity_shaped_sequence_write_queue_off() {
1768 let (_dir, backend) = issue_1029_pool(false);
1769 let result = issue_1029_create_entity_shaped_sequence(&backend).await;
1770 assert!(
1771 result.is_ok(),
1772 "khive#1029 repro (KHIVE_WRITE_QUEUE off): fts_upsert step failed: {:?}",
1773 result.err()
1774 );
1775 }
1776
1777 #[tokio::test]
1783 async fn issue_1029_create_entity_shaped_sequence_write_queue_on() {
1784 let (_dir, backend) = issue_1029_pool(true);
1785 let result = issue_1029_create_entity_shaped_sequence(&backend).await;
1786 assert!(
1787 result.is_ok(),
1788 "khive#1029 repro (KHIVE_WRITE_QUEUE=1): fts_upsert step failed: {:?}",
1789 result.err()
1790 );
1791 }
1792
1793 #[tokio::test]
1801 async fn issue_1029_two_pools_same_file_write_queue_on() {
1802 let dir = tempfile::tempdir().unwrap();
1803 let path = dir.path().join("issue_1029_two_pools.db");
1804
1805 let cfg = |p: std::path::PathBuf| crate::pool::PoolConfig {
1806 path: Some(p),
1807 busy_timeout: std::time::Duration::from_millis(200),
1808 write_queue_enabled: Some(true),
1809 ..crate::pool::PoolConfig::default()
1810 };
1811
1812 let pool_a = ConnectionPool::new(cfg(path.clone())).expect("pool A should open");
1813 let backend_a = StorageBackend {
1814 pool: Arc::new(pool_a),
1815 is_file_backed: true,
1816 path: Some(path.clone()),
1817 notes_seq_repair_runs: AtomicUsize::new(0),
1818 };
1819 let pool_b = ConnectionPool::new(cfg(path.clone())).expect("pool B should open");
1820 let backend_b = StorageBackend {
1821 pool: Arc::new(pool_b),
1822 is_file_backed: true,
1823 path: Some(path),
1824 notes_seq_repair_runs: AtomicUsize::new(0),
1825 };
1826
1827 let entities = backend_a
1828 .entities_for_namespace("tenant_ns")
1829 .expect("entities_for_namespace on pool A");
1830 let entity =
1831 khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029TwoPools");
1832 let entity_id = entity.id;
1833 entities
1834 .upsert_entity(entity)
1835 .await
1836 .expect("pool A entity upsert should succeed");
1837
1838 let text = backend_b.text("entities").expect("text on pool B");
1839 let doc = khive_storage::types::TextDocument {
1840 subject_id: entity_id,
1841 kind: khive_types::SubstrateKind::Entity,
1842 title: Some("Issue1029TwoPools".to_string()),
1843 body: "issue 1029 two-pool repro body".to_string(),
1844 tags: vec![],
1845 namespace: "tenant_ns".to_string(),
1846 metadata: None,
1847 updated_at: chrono::Utc::now(),
1848 };
1849 let result = text.upsert_document(doc).await;
1850 assert!(
1851 result.is_ok(),
1852 "khive#1029 two-pool repro: fts_upsert on an independent pool for the \
1853 same tenant DB file failed: {:?}",
1854 result.err()
1855 );
1856 }
1857
1858 struct StarvationCaptureSubscriber {
1861 events: Arc<std::sync::Mutex<Vec<std::collections::BTreeMap<String, String>>>>,
1862 }
1863
1864 impl tracing::Subscriber for StarvationCaptureSubscriber {
1865 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
1866 true
1867 }
1868 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
1869 tracing::span::Id::from_u64(1)
1870 }
1871 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
1872 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
1873 fn event(&self, event: &tracing::Event<'_>) {
1874 #[derive(Default)]
1875 struct FieldVisitor(std::collections::BTreeMap<String, String>);
1876 impl tracing::field::Visit for FieldVisitor {
1877 fn record_debug(
1878 &mut self,
1879 field: &tracing::field::Field,
1880 value: &dyn std::fmt::Debug,
1881 ) {
1882 self.0
1883 .insert(field.name().to_string(), format!("{value:?}"));
1884 }
1885 }
1886 let mut visitor = FieldVisitor::default();
1887 event.record(&mut visitor);
1888 self.events.lock().unwrap().push(visitor.0);
1889 }
1890 fn enter(&self, _: &tracing::span::Id) {}
1891 fn exit(&self, _: &tracing::span::Id) {}
1892 }
1893
1894 #[tokio::test]
1907 #[serial_test::serial(tx_registry)]
1908 async fn issue_1029_starvation_warn_reports_registered_transactions() {
1909 let (_dir, backend) = issue_1029_pool(false);
1910 let text = backend.text("entities").expect("text store");
1913
1914 let holder = backend
1918 .pool
1919 .open_standalone_writer()
1920 .expect("holder connection");
1921 holder
1922 .execute_batch("BEGIN IMMEDIATE")
1923 .expect("holder BEGIN IMMEDIATE");
1924 let fixture =
1925 khive_storage::tx_registry::register(Some("issue_1029_fixture_tx".to_string()));
1926
1927 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
1928 let subscriber = StarvationCaptureSubscriber {
1929 events: Arc::clone(&events),
1930 };
1931 let guard = tracing::subscriber::set_default(subscriber);
1932
1933 let doc = khive_storage::types::TextDocument {
1934 subject_id: uuid::Uuid::new_v4(),
1935 kind: khive_types::SubstrateKind::Entity,
1936 title: Some("Issue1029Starved".to_string()),
1937 body: "issue 1029 starvation diagnostic body".to_string(),
1938 tags: vec![],
1939 namespace: "tenant_ns".to_string(),
1940 metadata: None,
1941 updated_at: chrono::Utc::now(),
1942 };
1943 let result = text.upsert_document(doc).await;
1944
1945 drop(guard);
1946 drop(fixture);
1947 holder
1948 .execute_batch("ROLLBACK")
1949 .expect("holder ROLLBACK releases the lock");
1950
1951 assert!(
1952 result.is_err(),
1953 "upsert_document must starve while another connection holds the write lock"
1954 );
1955
1956 let events = events.lock().unwrap();
1957 let warn = events
1958 .iter()
1959 .find(|fields| {
1960 fields
1961 .get("message")
1962 .is_some_and(|m| m.contains("text write starved"))
1963 })
1964 .unwrap_or_else(|| panic!("expected a starvation WARN, captured events: {events:?}"));
1965 assert!(
1966 warn.get("op").is_some_and(|op| op.contains("fts_upsert")),
1967 "WARN must name the starved operation, got: {warn:?}"
1968 );
1969 assert!(
1970 warn.get("open_txs")
1971 .is_some_and(|txs| txs.contains("issue_1029_fixture_tx")),
1972 "WARN must list the registered holder label, got: {warn:?}"
1973 );
1974 let count: usize = warn
1975 .get("open_tx_count")
1976 .expect("WARN must carry open_tx_count")
1977 .parse()
1978 .expect("open_tx_count must be numeric");
1979 assert!(
1980 count >= 1,
1981 "open_tx_count must count the fixture, got {count}"
1982 );
1983 }
1984}