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
19mod pack_schema;
20
21fn sqlite_table_exists(conn: &rusqlite::Connection, table: &str) -> Result<bool, SqliteError> {
22 conn.query_row(
23 "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
24 rusqlite::params![table],
25 |row| row.get::<_, i64>(0),
26 )
27 .optional()
28 .map(|row| row.is_some())
29 .map_err(SqliteError::Rusqlite)
30}
31
32fn ensure_fts_rowid_map_backfilled(
88 conn: &rusqlite::Connection,
89 table: &str,
90) -> Result<(), SqliteError> {
91 let map = text::rowid_map_table(table);
92 let state = text::rowid_map_state_table(table);
93 let already_backfilled: bool = conn.query_row(
94 &format!("SELECT EXISTS(SELECT 1 FROM {state} WHERE key = 'backfill' AND value = ?1)"),
95 rusqlite::params![text::ROWID_MAP_BACKFILL_COMPLETE],
96 |row| row.get(0),
97 )?;
98 if already_backfilled {
99 return Ok(());
100 }
101 let fts_has_a_row: bool = conn.query_row(
102 &format!("SELECT EXISTS(SELECT 1 FROM {table} LIMIT 1)"),
103 [],
104 |row| row.get(0),
105 )?;
106 if !fts_has_a_row {
107 return Ok(());
108 }
109
110 conn.execute_batch("BEGIN IMMEDIATE")?;
111 let result: Result<(), SqliteError> = (|| {
112 conn.execute_batch(&format!(
113 "DELETE FROM {map} WHERE NOT EXISTS ( \
114 SELECT 1 FROM {table} \
115 WHERE {table}.rowid = {map}.rowid \
116 AND {table}.namespace = {map}.namespace \
117 AND {table}.subject_id = {map}.subject_id \
118 )"
119 ))?;
120 conn.execute_batch(&format!(
121 "INSERT OR REPLACE INTO {map} (namespace, subject_id, rowid) \
122 SELECT namespace, subject_id, rowid FROM {table} \
123 WHERE namespace IS NOT NULL AND subject_id IS NOT NULL \
124 ORDER BY updated_at ASC, rowid ASC"
125 ))?;
126 conn.execute_batch(&format!(
127 "DELETE FROM {table} \
128 WHERE namespace IS NOT NULL AND subject_id IS NOT NULL \
129 AND rowid NOT IN (SELECT rowid FROM {map})"
130 ))?;
131 conn.execute(
132 &format!("INSERT OR REPLACE INTO {state} (key, value) VALUES ('backfill', ?1)"),
133 rusqlite::params![text::ROWID_MAP_BACKFILL_COMPLETE],
134 )?;
135 Ok(())
136 })();
137
138 match result {
139 Ok(()) => {
140 conn.execute_batch("COMMIT")?;
141 Ok(())
142 }
143 Err(e) => {
144 let _ = conn.execute_batch("ROLLBACK");
145 Err(e)
146 }
147 }
148}
149
150fn warn_scan_fallback_once(table: &str) {
156 static WARNED: std::sync::Once = std::sync::Once::new();
157 WARNED.call_once(|| {
158 tracing::warn!(
159 table,
160 "opened a read-only text-search table with no rowid-map sidecar, or with a sidecar \
161 that has never proven a completed backfill (no durable completion marker); a \
162 read-only connection cannot create, backfill, or reconcile the map itself, so this \
163 falls back to pre-map namespace/subject_id scan predicates for get/delete on this \
164 table rather than trusting a map that might be partial"
165 );
166 });
167}
168
169fn validate_vector_model_key(model_key: &str) -> Result<(), SqliteError> {
170 if model_key.is_empty()
171 || !model_key
172 .chars()
173 .all(|c| c.is_ascii_alphanumeric() || c == '_')
174 {
175 return Err(SqliteError::InvalidData(format!(
176 "invalid model_key '{}': must be non-empty and contain only \
177 alphanumeric/underscore characters",
178 model_key
179 )));
180 }
181 Ok(())
182}
183
184fn validate_vector_table_columns(
185 conn: &rusqlite::Connection,
186 table: &str,
187) -> Result<(), SqliteError> {
188 let pragma = format!("PRAGMA table_xinfo({table})");
189 let mut stmt = conn.prepare(&pragma)?;
190 let mut rows = stmt.query([])?;
191 let mut has_field = false;
192 let mut has_embedding_model = false;
193 while let Some(row) = rows.next()? {
194 let name: String = row.get(1)?;
195 if name == "field" {
196 has_field = true;
197 }
198 if name == "embedding_model" {
199 has_embedding_model = true;
200 }
201 }
202 if !has_field || !has_embedding_model {
203 return Err(SqliteError::InvalidData(format!(
204 "vec0 table '{table}' is missing required column(s) (field={has_field}, \
205 embedding_model={has_embedding_model}); this is a pre-v0.2.8 vector schema and is \
206 not supported — recreate the database"
207 )));
208 }
209 Ok(())
210}
211
212pub struct StorageBackend {
214 pool: Arc<ConnectionPool>,
215 is_file_backed: bool,
216 path: Option<std::path::PathBuf>,
217 notes_seq_repair_runs: AtomicUsize,
224}
225
226impl StorageBackend {
227 pub fn sqlite(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
235 Self::sqlite_with_pool_config(path, PoolConfig::default(), None)
236 }
237
238 #[cfg(any(test, feature = "test-support"))]
240 pub fn sqlite_for_test(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
241 Self::sqlite_with_pool_config(path, PoolConfig::for_test(), None)
242 }
243
244 pub fn sqlite_with_max_readers(
247 path: impl AsRef<Path>,
248 max_readers: Option<usize>,
249 ) -> Result<Self, SqliteError> {
250 Self::sqlite_with_pool_config(path, PoolConfig::default(), max_readers)
251 }
252
253 fn sqlite_with_pool_config(
254 path: impl AsRef<Path>,
255 pool_config: PoolConfig,
256 max_readers: Option<usize>,
257 ) -> Result<Self, SqliteError> {
258 crate::extension::ensure_extensions_loaded();
259 let resolved = path.as_ref().to_path_buf();
260 let read_only =
261 std::fs::metadata(&resolved).is_ok_and(|metadata| metadata.permissions().readonly());
262 let mut config = PoolConfig {
263 path: Some(resolved.clone()),
264 read_only,
265 ..pool_config
266 };
267 if let Some(max_readers) = max_readers {
268 config.max_readers = max_readers;
269 }
270 if read_only {
271 config.write_queue_enabled = Some(false);
272 }
273 let pool = ConnectionPool::new(config)?;
274 Ok(Self {
275 pool: Arc::new(pool),
276 is_file_backed: true,
277 path: Some(resolved),
278 notes_seq_repair_runs: AtomicUsize::new(0),
279 })
280 }
281
282 pub fn sqlite_read_only(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
293 Self::sqlite_read_only_with_pool_config(path, PoolConfig::default(), None)
294 }
295
296 #[cfg(any(test, feature = "test-support"))]
298 pub fn sqlite_read_only_for_test(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
299 Self::sqlite_read_only_with_pool_config(path, PoolConfig::for_test(), None)
300 }
301
302 pub fn sqlite_read_only_with_max_readers(
304 path: impl AsRef<Path>,
305 max_readers: Option<usize>,
306 ) -> Result<Self, SqliteError> {
307 Self::sqlite_read_only_with_pool_config(path, PoolConfig::default(), max_readers)
308 }
309
310 fn sqlite_read_only_with_pool_config(
311 path: impl AsRef<Path>,
312 pool_config: PoolConfig,
313 max_readers: Option<usize>,
314 ) -> Result<Self, SqliteError> {
315 crate::extension::ensure_extensions_loaded();
316 let resolved = path.as_ref().to_path_buf();
317 let mut config = PoolConfig {
318 path: Some(resolved.clone()),
319 read_only: true,
320 write_queue_enabled: Some(false),
321 ..pool_config
322 };
323 if let Some(max_readers) = max_readers {
324 config.max_readers = max_readers;
325 }
326 let pool = ConnectionPool::new(config)?;
332 Ok(Self {
333 pool: Arc::new(pool),
334 is_file_backed: true,
335 path: Some(resolved),
336 notes_seq_repair_runs: AtomicUsize::new(0),
337 })
338 }
339
340 pub fn memory() -> Result<Self, SqliteError> {
346 crate::extension::ensure_extensions_loaded();
347 let config = PoolConfig {
348 path: None,
349 ..PoolConfig::default()
350 };
351 let pool = ConnectionPool::new(config)?;
352 Ok(Self {
353 pool: Arc::new(pool),
354 is_file_backed: false,
355 path: None,
356 notes_seq_repair_runs: AtomicUsize::new(0),
357 })
358 }
359
360 pub fn sql(&self) -> Arc<dyn khive_storage::SqlAccess> {
364 Arc::new(SqlBridge::new(Arc::clone(&self.pool), self.is_file_backed))
365 }
366
367 pub fn apply_schema(
374 &self,
375 plan: &crate::migrations::ServiceSchemaPlan,
376 ) -> Result<(), SqliteError> {
377 let writer = self.pool.try_writer()?;
378 crate::migrations::apply_schema_plan(writer.conn(), plan)
379 }
380
381 pub fn apply_pack_ddl_statements(
397 &self,
398 statements: &[&'static str],
399 ) -> Result<(), SqliteError> {
400 self.apply_pack_ddl_statements_with_columns(statements, &[])
401 }
402
403 pub fn apply_pack_ddl_statements_with_columns(
410 &self,
411 statements: &[&'static str],
412 additions: &[khive_types::PackColumnAddition],
413 ) -> Result<(), SqliteError> {
414 let writer = self.pool.try_writer()?;
415 writer.transaction(|conn| {
416 pack_schema::add_missing_columns(conn, additions)?;
417 for &stmt in statements {
418 conn.execute_batch(stmt)?;
419 }
420 pack_schema::validate_columns(conn, additions)?;
421 Ok(())
422 })
423 }
424
425 pub fn validate_pack_schema_columns(
428 &self,
429 additions: &[khive_types::PackColumnAddition],
430 ) -> Result<(), SqliteError> {
431 if additions.is_empty() {
432 return Ok(());
433 }
434 let reader = self.pool.reader()?;
435 pack_schema::validate_columns(reader.conn(), additions)
436 }
437
438 pub fn prepare_core_schema(&self) -> Result<u32, SqliteError> {
448 if self.is_read_only() {
449 let reader = self.pool.reader()?;
450 crate::migrations::validate_schema_is_current(reader.conn())
451 } else {
452 let latest = crate::migrations::MIGRATIONS
453 .last()
454 .map(|migration| migration.version)
455 .unwrap_or(0);
456 {
457 let reader = self.pool.reader()?;
458 let current = crate::migrations::read_schema_version(reader.conn())?;
459 if current >= latest {
460 return crate::migrations::validate_schema_is_current(reader.conn());
461 }
462 }
463 let owner = crate::stores::blob::acquire_database_gc_owner_for_path_blocking(
464 self.pool.canonical_path().map(Path::to_path_buf),
465 )
466 .map_err(|error| {
467 SqliteError::InvalidData(format!(
468 "failed to acquire database GC owner before schema preparation: {error}"
469 ))
470 })?;
471 let mut writer = self.pool.try_writer()?;
472 crate::migrations::run_migrations_with_database_gc_owner(writer.conn_mut(), &owner)
473 }
474 }
475
476 pub fn schema_version(&self) -> Result<u32, SqliteError> {
483 if self.is_read_only() {
484 let reader = self.pool.reader()?;
485 crate::migrations::read_schema_version(reader.conn())
486 } else {
487 let writer = self.pool.try_writer()?;
488 crate::migrations::read_schema_version(writer.conn())
489 }
490 }
491
492 pub fn attachment_cutover_status(
494 &self,
495 ) -> Result<crate::migrations::AttachmentCutoverStatus, SqliteError> {
496 if self.is_read_only() {
497 let reader = self.pool.reader()?;
498 crate::migrations::attachment_cutover_status(reader.conn())
499 } else {
500 let writer = self.pool.try_writer()?;
501 crate::migrations::attachment_cutover_status(writer.conn())
502 }
503 }
504
505 fn require_attachment_cutover_owner(
506 &self,
507 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
508 ) -> Result<(), SqliteError> {
509 let sql = self.sql();
510 let backend_path = sql.database_path();
511 if owner.database_path() != backend_path.as_deref() {
512 return Err(SqliteError::InvalidData(format!(
513 "attachment cutover GC owner targets {:?}, but this backend is {:?}",
514 owner.database_path(),
515 backend_path.as_deref()
516 )));
517 }
518 Ok(())
519 }
520
521 pub fn stage_attachment_cutover(
525 &self,
526 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
527 ) -> Result<(), SqliteError> {
528 self.require_attachment_cutover_owner(owner)?;
529 if self.is_read_only() {
530 return Err(SqliteError::InvalidData(
531 "cannot stage attachment cutover on a read-only backend".into(),
532 ));
533 }
534 let mut writer = self.pool.try_writer()?;
535 crate::migrations::stage_attachment_cutover(writer.conn_mut())
536 }
537
538 pub fn apply_verified_attachments(
540 &self,
541 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
542 attachments: &[khive_storage::Attachment],
543 ) -> Result<(), SqliteError> {
544 self.require_attachment_cutover_owner(owner)?;
545 if self.is_read_only() {
546 return Err(SqliteError::InvalidData(
547 "cannot apply verified attachments on a read-only backend".into(),
548 ));
549 }
550 let mut writer = self.pool.try_writer()?;
551 let tx = writer
552 .conn_mut()
553 .transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
554 for attachment in attachments {
555 attachment
556 .validate()
557 .map_err(|error| SqliteError::InvalidData(error.to_string()))?;
558 crate::migrations::apply_generic_verified_attachment(
559 &tx,
560 &attachment.record_uuid.to_string(),
561 attachment.substrate.as_str(),
562 &attachment.role,
563 &attachment.content_ref,
564 attachment.media_type.as_deref(),
565 attachment.size_bytes,
566 attachment.created_at,
567 )?;
568 }
569 tx.commit()?;
570 Ok(())
571 }
572
573 pub fn finalize_attachment_cutover(
576 &self,
577 owner: &crate::stores::blob::DatabaseGcOwnerGuard,
578 ) -> Result<(), SqliteError> {
579 self.require_attachment_cutover_owner(owner)?;
580 if self.is_read_only() {
581 return Err(SqliteError::InvalidData(
582 "cannot finalize attachment cutover on a read-only backend".into(),
583 ));
584 }
585 let mut writer = self.pool.try_writer()?;
586 crate::migrations::finalize_attachment_cutover(writer.conn_mut())
587 }
588
589 pub fn entities(&self) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
593 self.entities_for_namespace("local")
594 }
595
596 pub fn entities_for_namespace(
600 &self,
601 namespace: &str,
602 ) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
603 if namespace.trim().is_empty() {
604 return Err(SqliteError::InvalidData(
605 "entities namespace must be non-empty".to_string(),
606 ));
607 }
608 if !self.is_read_only() {
609 let writer = self.pool.try_writer()?;
610 entity::ensure_entities_schema(writer.conn())?;
611 }
612
613 Ok(Arc::new(entity::SqlEntityStore::new(
614 Arc::clone(&self.pool),
615 self.is_file_backed,
616 )))
617 }
618
619 pub fn attachments(&self) -> Result<Arc<dyn khive_storage::AttachmentStore>, SqliteError> {
626 Ok(Arc::new(attachment::SqlAttachmentStore::new(
627 Arc::clone(&self.pool),
628 self.is_file_backed,
629 )))
630 }
631
632 pub fn graph(&self) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
637 self.graph_for_namespace("local")
638 }
639
640 pub fn graph_for_namespace(
642 &self,
643 namespace: &str,
644 ) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
645 if namespace.trim().is_empty() {
646 return Err(SqliteError::InvalidData(
647 "graph namespace must be non-empty".to_string(),
648 ));
649 }
650 if !self.is_read_only() {
651 let writer = self.pool.try_writer()?;
652 graph::ensure_graph_schema(writer.conn())?;
653 }
654
655 Ok(Arc::new(graph::SqlGraphStore::new_scoped(
656 Arc::clone(&self.pool),
657 self.is_file_backed,
658 namespace.trim().to_string(),
659 )))
660 }
661
662 fn constructor_writer(&self) -> Result<crate::pool::WriterGuard<'_>, SqliteError> {
663 let context = khive_storage::capture_request_read_context();
664 let Some(operation) = context.store_acquisition_operation() else {
665 return self.pool.try_writer();
666 };
667 self.pool
668 .writer_until(|| context.blocking_stop_reason().is_some())?
669 .ok_or_else(|| {
670 SqliteError::RequestReadStopped(khive_storage::StorageError::Timeout {
671 operation: operation.into(),
672 })
673 })
674 }
675
676 pub fn notes(&self) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
680 self.notes_for_namespace("local")
681 }
682
683 pub fn notes_for_namespace(
687 &self,
688 namespace: &str,
689 ) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
690 if namespace.trim().is_empty() {
691 return Err(SqliteError::InvalidData(
692 "notes namespace must be non-empty".to_string(),
693 ));
694 }
695 if !self.is_read_only() {
696 let writer = self.constructor_writer()?;
697 note::ensure_notes_schema(writer.conn())?;
698
699 if self.notes_seq_repair_runs.load(Ordering::Relaxed) == 0 {
706 note::repair_notes_seq(writer.conn())?;
707 self.notes_seq_repair_runs.fetch_add(1, Ordering::Relaxed);
708 }
709 }
710
711 Ok(Arc::new(note::SqlNoteStore::new(
712 Arc::clone(&self.pool),
713 self.is_file_backed,
714 )))
715 }
716
717 pub fn notes_seq_repair_run_count(&self) -> usize {
722 self.notes_seq_repair_runs.load(Ordering::Relaxed)
723 }
724
725 pub fn events(&self) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
730 self.events_for_namespace("local")
731 }
732
733 pub fn events_for_namespace(
735 &self,
736 namespace: &str,
737 ) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
738 if namespace.trim().is_empty() {
739 return Err(SqliteError::InvalidData(
740 "events namespace must be non-empty".to_string(),
741 ));
742 }
743 if !self.is_read_only() {
744 let writer = self.constructor_writer()?;
745 event::ensure_events_schema(writer.conn())?;
746 }
747
748 Ok(Arc::new(event::SqlEventStore::new_scoped(
749 Arc::clone(&self.pool),
750 self.is_file_backed,
751 namespace.trim().to_string(),
752 )))
753 }
754
755 pub fn agents(&self) -> Result<Arc<dyn khive_storage::AgentStore>, SqliteError> {
760 if !self.is_read_only() {
761 let writer = self.pool.try_writer()?;
762 agents::ensure_agents_schema(writer.conn())?;
763 }
764
765 Ok(Arc::new(agents::SqlAgentStore::new(
766 Arc::clone(&self.pool),
767 self.is_file_backed,
768 )))
769 }
770
771 pub fn vectors(
777 &self,
778 model_key: &str,
779 embedding_model: &str,
780 dimensions: usize,
781 ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
782 self.vectors_for_namespace(model_key, embedding_model, dimensions, "local")
783 }
784
785 pub fn vectors_for_namespace(
795 &self,
796 model_key: &str,
797 embedding_model: &str,
798 dimensions: usize,
799 namespace: &str,
800 ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
801 validate_vector_model_key(model_key)?;
802 if namespace.trim().is_empty() {
803 return Err(SqliteError::InvalidData(
804 "vector store namespace must be non-empty".to_string(),
805 ));
806 }
807 self.ensure_vector_tables(&[(model_key, dimensions)])?;
808 Ok(Arc::new(vectors::SqliteVecStore::new(
809 Arc::clone(&self.pool),
810 self.is_file_backed,
811 model_key.to_string(),
812 embedding_model.to_string(),
813 dimensions,
814 namespace.trim().to_string(),
815 )?))
816 }
817
818 pub fn ensure_vector_tables(&self, models: &[(&str, usize)]) -> Result<(), SqliteError> {
821 for (model_key, _) in models {
822 validate_vector_model_key(model_key)?;
823 }
824 if models.is_empty() {
825 return Ok(());
826 }
827
828 crate::extension::ensure_extensions_loaded();
830
831 if self.is_read_only() {
832 let reader = self.pool.reader()?;
836 for (model_key, _) in models {
837 let table = format!("vec_{model_key}");
838 if !sqlite_table_exists(reader.conn(), &table)? {
839 return Err(SqliteError::InvalidData(format!(
840 "read-only database has no vector table '{table}'; create and populate it in \
841 a writable copy before opening the snapshot"
842 )));
843 }
844 validate_vector_table_columns(reader.conn(), &table)?;
845 }
846 return Ok(());
847 }
848
849 let writer = self.constructor_writer()?;
850
851 for (model_key, _) in models {
855 let table = format!("vec_{model_key}");
856 if sqlite_table_exists(writer.conn(), &table)? {
862 validate_vector_table_columns(writer.conn(), &table)?;
863 }
864 }
865
866 writer
875 .conn()
876 .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
877
878 writer
882 .conn()
883 .execute_batch(crate::migrations::ANN_WRITE_LOG_DDL)?;
884 writer
885 .conn()
886 .execute_batch(crate::migrations::ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL)?;
887 writer
888 .conn()
889 .execute_batch(crate::migrations::ANN_CONSUMER_PENDING_DDL)?;
890
891 for (model_key, dimensions) in models {
893 let ddl = format!(
894 "CREATE VIRTUAL TABLE IF NOT EXISTS vec_{} USING vec0(\
895 subject_id TEXT PRIMARY KEY, \
896 namespace TEXT NOT NULL, \
897 kind TEXT NOT NULL, \
898 field TEXT NOT NULL, \
899 embedding_model TEXT NOT NULL, \
900 embedding float[{}] distance_metric=cosine\
901 )",
902 model_key, dimensions
903 );
904 writer.conn().execute_batch(&ddl)?;
905 }
906 Ok(())
907 }
908
909 pub fn register_embedding_model(
914 &self,
915 engine_name: &str,
916 model_id: &str,
917 key_version: &str,
918 dimensions: u32,
919 ) -> Result<(), SqliteError> {
920 let writer = self.pool.try_writer()?;
921 writer
922 .conn()
923 .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
924
925 let now = chrono::Utc::now().timestamp_micros();
926 let canonical_key =
927 format!("{engine_name}:{model_id}:{key_version}:{dimensions}").into_bytes();
928 let id = uuid::Uuid::new_v4();
929 writer.conn().execute(
930 "INSERT INTO _embedding_models \
931 (id, engine_name, model_id, key_version, dim, output_dim, status, \
932 activated_at, superseded_at, superseded_by, canonical_key, created_at) \
933 VALUES (?1, ?2, ?3, ?4, ?5, NULL, 'active', ?6, NULL, NULL, ?7, ?8) \
934 ON CONFLICT(canonical_key) DO UPDATE SET \
935 status = 'active', \
936 activated_at = COALESCE(_embedding_models.activated_at, excluded.activated_at)",
937 rusqlite::params![
938 id.as_bytes().as_slice(),
939 engine_name,
940 model_id,
941 key_version,
942 dimensions as i64,
943 now,
944 canonical_key,
945 now,
946 ],
947 )?;
948 Ok(())
949 }
950
951 pub fn sparse(
955 &self,
956 model_key: &str,
957 ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
958 self.sparse_for_namespace(model_key, "local")
959 }
960
961 pub fn sparse_for_namespace(
965 &self,
966 model_key: &str,
967 namespace: &str,
968 ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
969 if model_key.is_empty()
970 || !model_key
971 .chars()
972 .all(|c| c.is_ascii_alphanumeric() || c == '_')
973 {
974 return Err(SqliteError::InvalidData(format!(
975 "invalid model_key '{}': must be non-empty and contain only alphanumeric/underscore characters",
976 model_key
977 )));
978 }
979 if namespace.trim().is_empty() {
980 return Err(SqliteError::InvalidData(
981 "sparse store namespace must be non-empty".to_string(),
982 ));
983 }
984
985 if self.is_read_only() {
986 let table = format!("sparse_{model_key}");
987 let reader = self.pool.reader()?;
988 if !sqlite_table_exists(reader.conn(), &table)? {
989 return Err(SqliteError::InvalidData(format!(
990 "read-only database has no sparse table '{table}'; create and populate it in \
991 a writable copy before opening the snapshot"
992 )));
993 }
994 } else {
995 let writer = self.pool.try_writer()?;
996 sparse::ensure_sparse_schema(writer.conn(), model_key)
997 .map_err(SqliteError::Rusqlite)?;
998 }
999
1000 Ok(Arc::new(sparse::SqliteSparseStore::new(
1001 Arc::clone(&self.pool),
1002 self.is_file_backed,
1003 model_key.to_string(),
1004 namespace.trim().to_string(),
1005 )?))
1006 }
1007
1008 pub fn text(&self, table_key: &str) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
1015 self.text_with_tokenizer(table_key, "trigram")
1016 }
1017
1018 pub fn text_with_tokenizer(
1026 &self,
1027 table_key: &str,
1028 tokenizer: &str,
1029 ) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
1030 if table_key.is_empty()
1031 || !table_key
1032 .chars()
1033 .all(|c| c.is_ascii_alphanumeric() || c == '_')
1034 {
1035 return Err(SqliteError::InvalidData(format!(
1036 "invalid table_key '{}': must be non-empty and contain only \
1037 alphanumeric/underscore characters",
1038 table_key
1039 )));
1040 }
1041 if table_key.ends_with("_rowids") || table_key.ends_with("_rowids_state") {
1051 return Err(SqliteError::InvalidData(format!(
1052 "invalid table_key '{}': must not end in '_rowids' or '_rowids_state' — those \
1053 suffixes are reserved for a text table's own rowid-map sidecar and its \
1054 completion-marker state table (see text::rowid_map_table, \
1055 text::rowid_map_state_table)",
1056 table_key
1057 )));
1058 }
1059 if tokenizer.is_empty()
1060 || !tokenizer
1061 .chars()
1062 .all(|c| c.is_ascii_alphanumeric() || c == '_')
1063 {
1064 return Err(SqliteError::InvalidData(format!(
1065 "invalid tokenizer '{}': must be non-empty and contain only \
1066 alphanumeric/underscore characters",
1067 tokenizer
1068 )));
1069 }
1070
1071 let ddl = format!(
1072 "CREATE VIRTUAL TABLE IF NOT EXISTS fts_{} USING fts5(\
1073 subject_id UNINDEXED, \
1074 kind UNINDEXED, \
1075 title, \
1076 body, \
1077 tags UNINDEXED, \
1078 namespace UNINDEXED, \
1079 metadata UNINDEXED, \
1080 updated_at UNINDEXED, \
1081 record_kind, \
1082 tokenize = '{}'\
1083 )",
1084 table_key, tokenizer
1085 );
1086 let table = format!("fts_{table_key}");
1087 if self.is_read_only() {
1088 let reader = self.pool.reader()?;
1089 if !sqlite_table_exists(reader.conn(), &table)? {
1090 return Err(SqliteError::InvalidData(format!(
1091 "read-only database has no text-search table '{table}'; create and populate \
1092 it in a writable copy before opening the snapshot"
1093 )));
1094 }
1095 let map = text::rowid_map_table(&table);
1109 let map_exists = sqlite_table_exists(reader.conn(), &map)?;
1110 let marker_present = if map_exists {
1111 let state = text::rowid_map_state_table(&table);
1112 sqlite_table_exists(reader.conn(), &state)?
1113 && reader.conn().query_row(
1114 &format!(
1115 "SELECT EXISTS(SELECT 1 FROM {state} WHERE key = 'backfill' AND value = ?1)"
1116 ),
1117 rusqlite::params![text::ROWID_MAP_BACKFILL_COMPLETE],
1118 |row| row.get(0),
1119 )?
1120 } else {
1121 false
1122 };
1123 if !map_exists || !marker_present {
1124 warn_scan_fallback_once(&table);
1125 return Ok(Arc::new(text::Fts5TextSearch::new_scan_fallback(
1126 Arc::clone(&self.pool),
1127 self.is_file_backed,
1128 table_key.to_string(),
1129 )));
1130 }
1131 } else {
1132 let writer = self.pool.try_writer()?;
1133 writer.conn().execute_batch(&ddl)?;
1134 writer.conn().execute_batch(&text::rowid_map_ddl(&table))?;
1135 ensure_fts_rowid_map_backfilled(writer.conn(), &table)?;
1136 }
1137
1138 Ok(Arc::new(text::Fts5TextSearch::new(
1139 Arc::clone(&self.pool),
1140 self.is_file_backed,
1141 table_key.to_string(),
1142 )))
1143 }
1144
1145 pub fn blob_store(
1153 &self,
1154 config_root: Option<&Path>,
1155 floor_bytes: Option<u64>,
1156 ) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
1157 let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
1158 let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
1159 Ok(Arc::new(blob::FsBlobStore::new(root, floor)?))
1160 }
1161
1162 pub fn blob_store_read_only(
1167 &self,
1168 config_root: Option<&Path>,
1169 floor_bytes: Option<u64>,
1170 ) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
1171 let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
1172 let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
1173 Ok(Arc::new(blob::FsBlobStore::open_existing(root, floor)?))
1174 }
1175
1176 pub fn is_file_backed(&self) -> bool {
1178 self.is_file_backed
1179 }
1180
1181 pub fn is_read_only(&self) -> bool {
1184 self.pool.config().read_only
1185 }
1186
1187 pub fn data_dir(&self) -> Option<std::path::PathBuf> {
1190 self.path.as_ref()?.parent().map(|p| p.to_path_buf())
1191 }
1192
1193 pub fn ann_root(&self) -> Option<std::path::PathBuf> {
1201 ann_root_for(self.path.as_ref()?)
1202 }
1203
1204 pub fn pool(&self) -> &ConnectionPool {
1206 &self.pool
1207 }
1208
1209 pub fn pool_arc(&self) -> Arc<ConnectionPool> {
1211 Arc::clone(&self.pool)
1212 }
1213}
1214
1215fn ann_root_for(path: &std::path::Path) -> Option<std::path::PathBuf> {
1220 let mut file = path.file_name()?.to_os_string();
1221 file.push(".ann");
1222 path.parent().map(|p| p.join(file))
1223}
1224
1225#[cfg(test)]
1226mod tests {
1227 use super::*;
1228 use khive_storage::types::{EdgeFilter, SqlStatement, SqlValue};
1229 use khive_storage::{EntityFilter, EventFilter};
1230
1231 #[tokio::test]
1232 async fn ordinary_store_accessors_ignore_request_read_cancellation() {
1233 let backend = StorageBackend::memory().unwrap();
1234 let (_sender, receiver) = tokio::sync::watch::channel(true);
1235 khive_storage::scope_request_read_cancellation(receiver, async {
1236 backend.notes().expect("ordinary notes accessor");
1237 backend.events().expect("ordinary events accessor");
1238 #[cfg(feature = "vectors")]
1239 backend
1240 .vectors("ordinary_store", "ordinary-store", 8)
1241 .expect("ordinary vectors accessor");
1242 })
1243 .await;
1244 }
1245
1246 #[tokio::test]
1247 async fn admitted_store_constructor_finishes_ddl_after_cancellation() {
1248 let backend = StorageBackend::memory().unwrap();
1249 let (sender, receiver) = tokio::sync::watch::channel(false);
1250 let fired = Arc::new(std::sync::atomic::AtomicBool::new(false));
1251 {
1252 let writer = backend.pool.writer().unwrap();
1253 let fired = fired.clone();
1254 writer
1255 .conn()
1256 .authorizer(Some(move |_: rusqlite::hooks::AuthContext<'_>| {
1257 fired.store(true, Ordering::SeqCst);
1258 sender.send_replace(true);
1259 rusqlite::hooks::Authorization::Allow
1260 }))
1261 .unwrap();
1262 }
1263 let result = khive_storage::scope_request_read_cancellation(receiver, async {
1264 khive_storage::capture_request_read_context()
1265 .scope_store_acquisition("admitted_notes_store", || backend.notes())
1266 })
1267 .await;
1268 let writer = backend.pool.writer().unwrap();
1269 writer
1270 .conn()
1271 .authorizer(
1272 None::<fn(rusqlite::hooks::AuthContext<'_>) -> rusqlite::hooks::Authorization>,
1273 )
1274 .unwrap();
1275 result.expect("request cancellation must not interrupt admitted constructor DDL");
1276 assert!(
1277 fired.load(Ordering::SeqCst),
1278 "cancellation must fire inside actual SQLite work"
1279 );
1280 assert_eq!(
1281 backend.notes_seq_repair_run_count(),
1282 1,
1283 "constructor must finish schema repair"
1284 );
1285 assert!(sqlite_table_exists(writer.conn(), "notes_seq").unwrap());
1286 }
1287
1288 #[cfg(unix)]
1289 use khive_storage::test_support::freeze_snapshot_sidecars;
1290
1291 #[tokio::test]
1292 async fn hot_path_guard_g2_file_backed_read_suite_uses_only_pooled_readers() {
1293 let dir = tempfile::tempdir().unwrap();
1294 let backend = StorageBackend::sqlite_for_test(dir.path().join("hot_path_g2.db")).unwrap();
1295 backend.prepare_core_schema().unwrap();
1296
1297 let entities = backend.entities().unwrap();
1300 let notes = backend.notes().unwrap();
1301 let graph = backend.graph().unwrap();
1302 let events = backend.events().unwrap();
1303 let text = backend
1304 .text_with_tokenizer("hot_path_g2", "unicode61")
1305 .unwrap();
1306 let agents = backend.agents().unwrap();
1307 let attachments = backend.attachments().unwrap();
1308 let sparse = backend.sparse("hot_path_g2").unwrap();
1309 #[cfg(feature = "vectors")]
1310 let vectors = backend.vectors("hot_path_g2", "test-model", 2).unwrap();
1311 let sql = backend.sql();
1312
1313 let before = backend.pool().reader_acquisition_snapshot();
1314 assert_eq!(
1315 entities
1316 .count_entities("local", EntityFilter::default())
1317 .await
1318 .unwrap(),
1319 0
1320 );
1321 assert_eq!(notes.count_notes("local", None).await.unwrap(), 0);
1322 assert_eq!(graph.count_edges(EdgeFilter::default()).await.unwrap(), 0);
1323 assert_eq!(
1324 events.count_events(EventFilter::default()).await.unwrap(),
1325 0
1326 );
1327 assert!(text
1328 .get_document("local", uuid::Uuid::new_v4())
1329 .await
1330 .unwrap()
1331 .is_none());
1332 assert!(agents.get("no-such-agent").await.unwrap().is_none());
1333 assert!(attachments
1334 .get_attachment(uuid::Uuid::new_v4(), "primary")
1335 .await
1336 .unwrap()
1337 .is_none());
1338 assert_eq!(sparse.count().await.unwrap(), 0);
1339 #[cfg(feature = "vectors")]
1340 assert_eq!(vectors.count().await.unwrap(), 0);
1341
1342 let mut raw = sql.reader().await.unwrap();
1343 assert!(matches!(
1344 raw.query_scalar(SqlStatement {
1345 sql: "SELECT 1".into(),
1346 params: Vec::new(),
1347 label: None,
1348 })
1349 .await
1350 .unwrap(),
1351 Some(SqlValue::Integer(1))
1352 ));
1353
1354 let after = backend.pool().reader_acquisition_snapshot();
1355 let expected_pooled_delta = 9 + u64::from(cfg!(feature = "vectors"));
1356 assert_eq!(
1357 after.pooled_checkouts - before.pooled_checkouts,
1358 expected_pooled_delta,
1359 "each ordinary file-backed read must check out exactly one pooled reader"
1360 );
1361 assert_eq!(
1362 after.standalone_opens, before.standalone_opens,
1363 "ADR-166 G2: ordinary file-backed read verbs must not open standalone readers"
1364 );
1365 assert_eq!(after.active_pooled_checkouts, 0);
1366 assert_eq!(
1367 after.completed_pooled_checkouts - before.completed_pooled_checkouts,
1368 expected_pooled_delta
1369 );
1370 }
1371
1372 #[cfg(unix)]
1373 #[tokio::test]
1374 async fn sqlite_detects_chmod_read_only_snapshot_and_core_reads_succeed() {
1375 use std::os::unix::fs::PermissionsExt;
1376
1377 let dir = tempfile::tempdir().unwrap();
1378 let path = dir.path().join("chmod_snapshot.db");
1379 {
1380 let writable =
1381 StorageBackend::sqlite_for_test(&path).expect("create writable database");
1382 writable
1383 .prepare_core_schema()
1384 .expect("migrate writable snapshot source");
1385 }
1386
1387 let mut permissions = std::fs::metadata(&path).unwrap().permissions();
1388 permissions.set_mode(0o444);
1389 std::fs::set_permissions(&path, permissions).unwrap();
1390 freeze_snapshot_sidecars(&path);
1391
1392 let read_only = StorageBackend::sqlite(&path).expect("auto-detect read-only mode");
1393 assert!(read_only.is_read_only());
1394 assert_eq!(
1395 read_only.pool().config().write_queue_enabled,
1396 Some(false),
1397 "read-only boot must not attempt to spawn a writer task"
1398 );
1399 assert!(read_only
1400 .pool()
1401 .writer_task_handle()
1402 .expect("disabled writer task is a valid configuration")
1403 .is_none());
1404 read_only
1405 .prepare_core_schema()
1406 .expect("current snapshot validates without migration writes");
1407
1408 let entities = read_only.entities().expect("entity store opens read-only");
1409 let graph = read_only.graph().expect("graph store opens read-only");
1410 let notes = read_only.notes().expect("note store opens read-only");
1411 let events = read_only.events().expect("event store opens read-only");
1412 assert_eq!(
1413 entities
1414 .count_entities("local", khive_storage::EntityFilter::default())
1415 .await
1416 .unwrap(),
1417 0
1418 );
1419 assert_eq!(
1420 graph
1421 .count_edges(khive_storage::types::EdgeFilter::default())
1422 .await
1423 .unwrap(),
1424 0
1425 );
1426 assert_eq!(notes.count_notes("local", None).await.unwrap(), 0);
1427 assert_eq!(
1428 events
1429 .count_events(khive_storage::EventFilter::default())
1430 .await
1431 .unwrap(),
1432 0
1433 );
1434 assert_eq!(
1435 read_only.notes_seq_repair_run_count(),
1436 0,
1437 "read-only store acquisition must not run the DML repair"
1438 );
1439 }
1440
1441 #[test]
1442 fn memory_backend_creates_successfully() {
1443 let backend = StorageBackend::memory().expect("memory backend should create");
1444 assert!(!backend.is_file_backed());
1445 }
1446
1447 #[test]
1448 fn file_backend_creates_successfully() {
1449 let dir = tempfile::tempdir().unwrap();
1450 let path = dir.path().join("test.db");
1451 let backend = StorageBackend::sqlite(&path).expect("file backend should create");
1452 assert!(backend.is_file_backed());
1453 assert!(path.exists());
1454 }
1455
1456 #[test]
1457 fn data_dir_returns_none_for_memory_backend() {
1458 let backend = StorageBackend::memory().expect("memory backend");
1459 assert!(backend.data_dir().is_none());
1460 }
1461
1462 #[test]
1463 fn data_dir_returns_parent_dir_for_file_backend() {
1464 let dir = tempfile::tempdir().unwrap();
1465 let path = dir.path().join("data.db");
1466 let backend = StorageBackend::sqlite_for_test(&path).expect("file backend");
1467 let got = backend.data_dir().expect("file backend must return Some");
1468 assert_eq!(got, dir.path());
1469 }
1470
1471 #[test]
1472 fn ann_root_is_database_scoped_sibling_dir() {
1473 let dir = tempfile::tempdir().unwrap();
1474 let path = dir.path().join("data.db");
1475 let backend = StorageBackend::sqlite_for_test(&path).expect("file backend");
1476 let got = backend.ann_root().expect("file backend must return Some");
1477 assert_eq!(got, dir.path().join("data.db.ann"));
1478 assert!(StorageBackend::memory().unwrap().ann_root().is_none());
1479 }
1480
1481 #[cfg(unix)]
1487 #[test]
1488 fn ann_root_distinct_for_non_utf8_filenames() {
1489 use std::os::unix::ffi::OsStrExt;
1490 let path_a = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xff.db"));
1491 let path_b = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xfe.db"));
1492 let root_a = ann_root_for(&path_a).expect("Some for a file path");
1493 let root_b = ann_root_for(&path_b).expect("Some for a file path");
1494 assert_ne!(
1495 root_a, root_b,
1496 "distinct database files must map to distinct ANN roots"
1497 );
1498 }
1499
1500 #[tokio::test]
1501 async fn sql_access_memory_roundtrip() {
1502 let backend = StorageBackend::memory().unwrap();
1503 let sql = backend.sql();
1504
1505 let mut writer = sql.writer().await.unwrap();
1506 writer
1507 .execute_script(
1508 "CREATE TABLE test_rt (id TEXT PRIMARY KEY, value INTEGER NOT NULL)".into(),
1509 )
1510 .await
1511 .unwrap();
1512
1513 let affected = writer
1514 .execute(SqlStatement {
1515 sql: "INSERT INTO test_rt (id, value) VALUES (?1, ?2)".into(),
1516 params: vec![SqlValue::Text("row1".into()), SqlValue::Integer(42)],
1517 label: None,
1518 })
1519 .await
1520 .unwrap();
1521 assert_eq!(affected, 1);
1522
1523 let mut reader = sql.reader().await.unwrap();
1524 let row = reader
1525 .query_row(SqlStatement {
1526 sql: "SELECT id, value FROM test_rt WHERE id = ?1".into(),
1527 params: vec![SqlValue::Text("row1".into())],
1528 label: None,
1529 })
1530 .await
1531 .unwrap();
1532
1533 let row = row.expect("should find the inserted row");
1534 assert_eq!(row.columns.len(), 2);
1535 match &row.columns[0].value {
1536 SqlValue::Text(s) => assert_eq!(s, "row1"),
1537 other => panic!("expected Text, got {other:?}"),
1538 }
1539 match &row.columns[1].value {
1540 SqlValue::Integer(v) => assert_eq!(*v, 42),
1541 other => panic!("expected Integer, got {other:?}"),
1542 }
1543 }
1544
1545 #[tokio::test]
1546 async fn sql_access_file_roundtrip() {
1547 let dir = tempfile::tempdir().unwrap();
1548 let path = dir.path().join("test_roundtrip.db");
1549 let backend = StorageBackend::sqlite_for_test(&path).unwrap();
1550 let sql = backend.sql();
1551
1552 let mut writer = sql.writer().await.unwrap();
1553 writer
1554 .execute_script("CREATE TABLE test_f (k TEXT PRIMARY KEY, v TEXT)".into())
1555 .await
1556 .unwrap();
1557 writer
1558 .execute(SqlStatement {
1559 sql: "INSERT INTO test_f (k, v) VALUES (?1, ?2)".into(),
1560 params: vec![
1561 SqlValue::Text("hello".into()),
1562 SqlValue::Text("world".into()),
1563 ],
1564 label: None,
1565 })
1566 .await
1567 .unwrap();
1568
1569 let mut reader = sql.reader().await.unwrap();
1570 let rows = reader
1571 .query_all(SqlStatement {
1572 sql: "SELECT k, v FROM test_f".into(),
1573 params: vec![],
1574 label: None,
1575 })
1576 .await
1577 .unwrap();
1578 assert_eq!(rows.len(), 1);
1579 match &rows[0].columns[1].value {
1580 SqlValue::Text(s) => assert_eq!(s, "world"),
1581 other => panic!("expected Text, got {other:?}"),
1582 }
1583 }
1584
1585 #[test]
1586 fn sqlite_read_only_missing_path_does_not_create_file() {
1587 let dir = tempfile::tempdir().unwrap();
1588 let path = dir.path().join("missing_ro.db");
1589 assert!(!path.exists());
1590
1591 let result = StorageBackend::sqlite_read_only(&path);
1592 assert!(
1593 result.is_err(),
1594 "opening a missing path read-only must fail"
1595 );
1596 assert!(
1597 !path.exists(),
1598 "opening a missing path read-only must not create the file"
1599 );
1600 }
1601
1602 #[test]
1603 fn sqlite_read_only_sparse_store_requires_existing_table_without_writer_acquisition() {
1604 let dir = tempfile::tempdir().unwrap();
1605 let path = dir.path().join("ro_sparse_tables.db");
1606 {
1607 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
1608 writable
1609 .prepare_core_schema()
1610 .expect("migrate snapshot source");
1611 writable
1612 .sparse("present")
1613 .expect("create the optional sparse table while writable");
1614 }
1615 #[cfg(unix)]
1616 freeze_snapshot_sidecars(&path);
1617
1618 let read_only = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
1619 read_only
1620 .prepare_core_schema()
1621 .expect("validate exact current migration ledger");
1622 read_only
1623 .sparse("present")
1624 .expect("an existing sparse table must open read-only");
1625 let missing = match read_only.sparse("missing") {
1626 Ok(_) => panic!("a missing sparse table must fail during store acquisition"),
1627 Err(error) => error,
1628 };
1629 assert!(
1630 missing.to_string().contains("sparse_missing"),
1631 "the diagnostic must name the absent table: {missing}"
1632 );
1633 assert_eq!(
1634 read_only.pool().writer_acquisition_snapshot(),
1635 crate::pool::WriterAcquisitionSnapshot::default(),
1636 "construction, exact-ledger validation, and optional sparse-table inspection must \
1637 use reader connections only"
1638 );
1639 }
1640
1641 #[test]
1642 fn sqlite_read_only_text_store_requires_existing_table_without_writer_acquisition() {
1643 let dir = tempfile::tempdir().unwrap();
1644 let path = dir.path().join("ro_text_tables.db");
1645 {
1646 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
1647 writable
1648 .prepare_core_schema()
1649 .expect("migrate snapshot source");
1650 writable
1651 .text("present")
1652 .expect("create the optional FTS table while writable");
1653 }
1654 #[cfg(unix)]
1655 freeze_snapshot_sidecars(&path);
1656
1657 let read_only = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
1658 read_only
1659 .prepare_core_schema()
1660 .expect("validate exact current migration ledger");
1661 read_only
1662 .text("present")
1663 .expect("an existing FTS table must open read-only");
1664 let missing = match read_only.text("missing") {
1665 Ok(_) => panic!("a missing FTS table must fail during store acquisition"),
1666 Err(error) => error,
1667 };
1668 assert!(
1669 missing.to_string().contains("fts_missing"),
1670 "the diagnostic must name the absent table: {missing}"
1671 );
1672 assert_eq!(
1673 read_only.pool().writer_acquisition_snapshot(),
1674 crate::pool::WriterAcquisitionSnapshot::default(),
1675 "construction, exact-ledger validation, and optional FTS inspection must use reader \
1676 connections only"
1677 );
1678 }
1679
1680 #[cfg(feature = "vectors")]
1681 #[test]
1682 fn sqlite_read_only_vector_store_schema_check_uses_no_writer_acquisition() {
1683 let dir = tempfile::tempdir().unwrap();
1684 let path = dir.path().join("ro_vector_tables.db");
1685 {
1686 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
1687 writable
1688 .prepare_core_schema()
1689 .expect("migrate snapshot source");
1690 writable
1691 .vectors("present", "present", 3)
1692 .expect("create the optional vector table while writable");
1693 }
1694 #[cfg(unix)]
1695 freeze_snapshot_sidecars(&path);
1696
1697 let read_only = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
1698 read_only
1699 .prepare_core_schema()
1700 .expect("validate exact current migration ledger");
1701 read_only
1702 .vectors("present", "present", 3)
1703 .expect("an existing vector table must open read-only");
1704 assert!(
1705 read_only.vectors("missing", "missing", 3).is_err(),
1706 "a missing vector table must fail during store acquisition"
1707 );
1708 assert_eq!(
1709 read_only.pool().writer_acquisition_snapshot(),
1710 crate::pool::WriterAcquisitionSnapshot::default(),
1711 "construction, exact-ledger validation, and optional vector inspection must use \
1712 reader connections only"
1713 );
1714 }
1715
1716 #[tokio::test]
1717 async fn sqlite_read_only_sql_writer_rejects_ddl_and_insert() {
1718 let dir = tempfile::tempdir().unwrap();
1719 let path = dir.path().join("ro_writer.db");
1720
1721 {
1723 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
1724 let sql = writable.sql();
1725 let mut writer = sql.writer().await.unwrap();
1726 writer
1727 .execute_script("CREATE TABLE ro_existing (id INTEGER PRIMARY KEY)".into())
1728 .await
1729 .unwrap();
1730 }
1731 #[cfg(unix)]
1732 freeze_snapshot_sidecars(&path);
1733
1734 let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1735 let sql = ro.sql();
1736
1737 let writer_result = sql.writer().await;
1739 assert!(
1740 writer_result.is_err(),
1741 "sql().writer() must be rejected on a read-only backend"
1742 );
1743 }
1744
1745 #[tokio::test]
1746 #[cfg(feature = "vectors")]
1747 async fn vectors_roundtrip_via_public_api() {
1748 let backend = StorageBackend::memory().unwrap();
1749 let store = backend.vectors("test_api", "test_api", 3).unwrap();
1750
1751 let id = uuid::Uuid::new_v4();
1752 store
1753 .insert(
1754 id,
1755 khive_types::SubstrateKind::Entity,
1756 "local",
1757 "content",
1758 vec![vec![1.0, 0.0, 0.0]],
1759 )
1760 .await
1761 .unwrap();
1762
1763 let hits = store
1764 .search(khive_storage::types::VectorSearchRequest {
1765 query_vectors: vec![vec![1.0, 0.0, 0.0]],
1766 top_k: 1,
1767 namespace: None,
1768 kind: None,
1769 embedding_model: None,
1770 filter: None,
1771 backend_hints: None,
1772 })
1773 .await
1774 .unwrap();
1775
1776 assert_eq!(hits.len(), 1);
1777 assert_eq!(hits[0].subject_id, id);
1778 assert!(hits[0].score.to_f64() > 0.99);
1779 }
1780
1781 #[tokio::test]
1782 #[cfg(feature = "vectors")]
1783 async fn vectors_creates_table_idempotently() {
1784 let backend = StorageBackend::memory().unwrap();
1785
1786 let store1 = backend.vectors("idempotent", "idempotent", 3).unwrap();
1787 let store2 = backend.vectors("idempotent", "idempotent", 3).unwrap();
1788
1789 let id = uuid::Uuid::new_v4();
1790 store1
1791 .insert(
1792 id,
1793 khive_types::SubstrateKind::Entity,
1794 "local",
1795 "content",
1796 vec![vec![1.0, 0.0, 0.0]],
1797 )
1798 .await
1799 .unwrap();
1800
1801 let count = store2.count().await.unwrap();
1802 assert_eq!(count, 1);
1803 }
1804
1805 #[tokio::test]
1806 async fn text_roundtrip_via_public_api() {
1807 let backend = StorageBackend::memory().unwrap();
1808 let store = backend.text("test_api").unwrap();
1809
1810 let id = uuid::Uuid::new_v4();
1811 let doc = khive_storage::types::TextDocument {
1812 subject_id: id,
1813 kind: khive_types::SubstrateKind::Entity,
1814 record_kind: None,
1815 title: Some("Test Title".to_string()),
1816 body: "This is a searchable document about Rust.".to_string(),
1817 tags: vec!["rust".to_string()],
1818 namespace: "test_ns".to_string(),
1819 metadata: None,
1820 updated_at: chrono::Utc::now(),
1821 };
1822 store.upsert_document(doc).await.unwrap();
1823
1824 let hits = store
1825 .search(khive_storage::types::TextSearchRequest {
1826 query: "Rust".to_string(),
1827 mode: khive_storage::types::TextQueryMode::Plain,
1828 filter: Some(khive_storage::types::TextFilter {
1829 namespaces: vec!["test_ns".to_string()],
1830 ..Default::default()
1831 }),
1832 top_k: 1,
1833 snippet_chars: 64,
1834 })
1835 .await
1836 .unwrap();
1837
1838 assert_eq!(hits.len(), 1);
1839 assert_eq!(hits[0].subject_id, id);
1840 assert!(hits[0].score.to_f64() > 0.0);
1841 }
1842
1843 #[tokio::test]
1844 async fn text_creates_table_idempotently() {
1845 let backend = StorageBackend::memory().unwrap();
1846
1847 let store1 = backend.text("idempotent_fts").unwrap();
1848 let store2 = backend.text("idempotent_fts").unwrap();
1849
1850 let id = uuid::Uuid::new_v4();
1851 let doc = khive_storage::types::TextDocument {
1852 subject_id: id,
1853 kind: khive_types::SubstrateKind::Note,
1854 record_kind: None,
1855 title: None,
1856 body: "Hello world.".to_string(),
1857 tags: vec![],
1858 namespace: "test_ns".to_string(),
1859 metadata: None,
1860 updated_at: chrono::Utc::now(),
1861 };
1862 store1.upsert_document(doc).await.unwrap();
1863
1864 let count = store2
1865 .count(khive_storage::types::TextFilter {
1866 namespaces: vec!["test_ns".to_string()],
1867 ..Default::default()
1868 })
1869 .await
1870 .unwrap();
1871 assert_eq!(count, 1);
1872 }
1873
1874 #[tokio::test]
1887 async fn text_repeated_open_after_backfill_does_not_scale_with_row_count() {
1888 use std::sync::atomic::{AtomicU64, Ordering};
1889 use std::sync::Arc;
1890
1891 fn repeated_open_work(backend: &StorageBackend) -> u64 {
1892 let work = Arc::new(AtomicU64::new(0));
1893 let counted = Arc::clone(&work);
1894 {
1895 let writer = backend.pool().writer().unwrap();
1896 writer
1897 .conn()
1898 .progress_handler(
1899 1,
1900 Some(move || {
1901 counted.fetch_add(1, Ordering::Relaxed);
1902 false
1903 }),
1904 )
1905 .unwrap();
1906 }
1907
1908 let result = (0..500).try_for_each(|_| backend.text("hot_path_reopen").map(|_| ()));
1912 backend
1913 .pool()
1914 .writer()
1915 .unwrap()
1916 .conn()
1917 .progress_handler(0, None::<fn() -> bool>)
1918 .unwrap();
1919 result.expect("repeated text-store opens must succeed");
1920 work.load(Ordering::Relaxed)
1921 }
1922
1923 let backend = StorageBackend::memory().unwrap();
1924 let store = backend.text("hot_path_reopen").unwrap();
1925 let body = "the quick brown fox jumps over the lazy dog ".repeat(35);
1926 let mut seeded = 0;
1927 let mut work = Vec::new();
1928 for target_rows in [100, 5_000] {
1929 for _ in seeded..target_rows {
1930 let doc = khive_storage::types::TextDocument {
1931 subject_id: uuid::Uuid::new_v4(),
1932 kind: khive_types::SubstrateKind::Note,
1933 record_kind: Some("memory".to_string()),
1934 title: None,
1935 body: body.clone(),
1936 tags: vec![],
1937 namespace: "test_ns".to_string(),
1938 metadata: None,
1939 updated_at: chrono::Utc::now(),
1940 };
1941 store.upsert_document(doc).await.unwrap();
1942 }
1943 seeded = target_rows;
1944 assert_eq!(
1945 store
1946 .count(khive_storage::types::TextFilter {
1947 namespaces: vec!["test_ns".to_string()],
1948 ..Default::default()
1949 })
1950 .await
1951 .unwrap(),
1952 target_rows,
1953 "the work comparison requires both declared row populations"
1954 );
1955
1956 let _ = backend.text("hot_path_reopen").unwrap();
1959 work.push(repeated_open_work(&backend));
1960 }
1961
1962 let [small, large] = [work[0], work[1]];
1963 assert!(small > 0 && large > 0, "both work meters must be active");
1964 assert!(
1965 large <= small * 2,
1966 "500 repeated backend.text() calls used {small} SQLite VM progress units at \
1967 100 rows and {large} at 5,000 rows; growing the table 50-fold must not \
1968 more than double already-backfilled work (for example via COUNT(*))"
1969 );
1970 }
1971
1972 #[tokio::test]
1978 async fn text_open_after_legacy_seed_backfills_the_map_with_full_parity() {
1979 let backend = StorageBackend::memory().unwrap();
1980 let table_key = "legacy_seed_parity";
1981 let table = format!("fts_{table_key}");
1982 let map = format!("{table}_rowids");
1983
1984 let _ = backend.text(table_key).unwrap();
1987
1988 {
1992 let writer = backend.pool().writer().unwrap();
1993 writer.conn().execute_batch("BEGIN").unwrap();
1994 {
1995 let mut insert = writer
1996 .conn()
1997 .prepare(&format!(
1998 "INSERT INTO {table} \
1999 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, \
2000 record_kind) \
2001 VALUES (?1, 'note', '', 'legacy body', '[]', 'test_ns', NULL, 0, 'memory')"
2002 ))
2003 .unwrap();
2004 for i in 0..500 {
2005 insert
2006 .execute(rusqlite::params![format!("legacy-{i}")])
2007 .unwrap();
2008 }
2009 }
2010 writer.conn().execute_batch("COMMIT").unwrap();
2011 }
2012 {
2013 let writer = backend.pool().writer().unwrap();
2014 let map_count: i64 = writer
2015 .conn()
2016 .query_row(&format!("SELECT COUNT(*) FROM {map}"), [], |row| row.get(0))
2017 .unwrap();
2018 assert_eq!(
2019 map_count, 0,
2020 "the raw-SQL seed must bypass the map, reproducing a genuinely pre-migration db"
2021 );
2022 }
2023
2024 let _ = backend.text(table_key).unwrap();
2027
2028 {
2029 let writer = backend.pool().writer().unwrap();
2030 let mismatched: i64 = writer
2031 .conn()
2032 .query_row(
2033 &format!(
2034 "SELECT \
2035 (SELECT COUNT(*) FROM {table} WHERE rowid NOT IN (SELECT rowid FROM {map})) + \
2036 (SELECT COUNT(*) FROM {map} WHERE rowid NOT IN (SELECT rowid FROM {table}))"
2037 ),
2038 [],
2039 |row| row.get(0),
2040 )
2041 .unwrap();
2042 assert_eq!(
2043 mismatched, 0,
2044 "backfill must give every FTS row exactly one map entry, both directions"
2045 );
2046 let fts_count: i64 = writer
2047 .conn()
2048 .query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| {
2049 row.get(0)
2050 })
2051 .unwrap();
2052 let map_count: i64 = writer
2053 .conn()
2054 .query_row(&format!("SELECT COUNT(*) FROM {map}"), [], |row| row.get(0))
2055 .unwrap();
2056 assert_eq!(fts_count, 500);
2057 assert_eq!(map_count, 500);
2058 }
2059 }
2060
2061 #[tokio::test]
2070 async fn text_open_reconciles_a_partial_map_instead_of_treating_it_as_complete() {
2071 let backend = StorageBackend::memory().unwrap();
2072 let table_key = "partial_map_reconcile";
2073 let table = format!("fts_{table_key}");
2074 let map = format!("{table}_rowids");
2075 let state = format!("{map}_state");
2076
2077 let _ = backend.text(table_key).unwrap();
2080
2081 let a = uuid::Uuid::new_v4();
2082 let b = uuid::Uuid::new_v4();
2083 {
2084 let writer = backend.pool().writer().unwrap();
2085 writer.conn().execute_batch("BEGIN").unwrap();
2086 writer
2087 .conn()
2088 .execute(
2089 &format!(
2090 "INSERT INTO {table} \
2091 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2092 updated_at, record_kind) \
2093 VALUES (1, ?1, 'note', '', 'doc a', '[]', 'test_ns', NULL, 0, 'memory')"
2094 ),
2095 rusqlite::params![a.to_string()],
2096 )
2097 .expect("insert A's fts row");
2098 writer
2099 .conn()
2100 .execute(
2101 &format!(
2102 "INSERT INTO {table} \
2103 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2104 updated_at, record_kind) \
2105 VALUES (2, ?1, 'note', '', 'doc b', '[]', 'test_ns', NULL, 0, 'memory')"
2106 ),
2107 rusqlite::params![b.to_string()],
2108 )
2109 .expect("insert B's fts row");
2110 writer
2112 .conn()
2113 .execute(
2114 &format!(
2115 "INSERT INTO {map} (namespace, subject_id, rowid) VALUES ('test_ns', ?1, 2)"
2116 ),
2117 rusqlite::params![b.to_string()],
2118 )
2119 .expect("insert B's own map entry, leaving A's missing");
2120 writer
2124 .conn()
2125 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2126 .expect("clear the completion marker");
2127 writer.conn().execute_batch("COMMIT").unwrap();
2128 }
2129
2130 let store = backend.text(table_key).unwrap();
2131
2132 let a_mapped: i64 = {
2133 let writer = backend.pool().writer().unwrap();
2134 writer
2135 .conn()
2136 .query_row(
2137 &format!("SELECT COUNT(*) FROM {map} WHERE namespace = 'test_ns' AND subject_id = ?1"),
2138 rusqlite::params![a.to_string()],
2139 |row| row.get(0),
2140 )
2141 .unwrap()
2142 };
2143 assert_eq!(
2144 a_mapped, 1,
2145 "the partial map must be reconciled, not left missing A's entry"
2146 );
2147
2148 let fetched_a = store.get_document("test_ns", a).await.unwrap();
2149 assert!(
2150 fetched_a.is_some(),
2151 "get_document(A) must work once the partial map is reconciled"
2152 );
2153 }
2154
2155 #[tokio::test]
2167 async fn text_open_removes_a_wrong_key_map_row_before_backfilling_the_right_one() {
2168 let backend = StorageBackend::memory().unwrap();
2169 let table_key = "wrong_key_map_row";
2170 let table = format!("fts_{table_key}");
2171 let map = format!("{table}_rowids");
2172 let state = format!("{map}_state");
2173
2174 let _ = backend.text(table_key).unwrap();
2175
2176 let a = uuid::Uuid::new_v4();
2177 let b = uuid::Uuid::new_v4();
2178 {
2179 let writer = backend.pool().writer().unwrap();
2180 writer.conn().execute_batch("BEGIN").unwrap();
2181 writer
2182 .conn()
2183 .execute(
2184 &format!(
2185 "INSERT INTO {table} \
2186 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2187 updated_at, record_kind) \
2188 VALUES (7, ?1, 'note', '', 'doc b', '[]', 'test_ns', NULL, 0, 'memory')"
2189 ),
2190 rusqlite::params![b.to_string()],
2191 )
2192 .expect("insert B's live fts row at rowid 7");
2193 writer
2194 .conn()
2195 .execute(
2196 &format!(
2197 "INSERT INTO {map} (namespace, subject_id, rowid) VALUES ('test_ns', ?1, 7)"
2198 ),
2199 rusqlite::params![a.to_string()],
2200 )
2201 .expect("insert A's stale map row still pointing at rowid 7");
2202 writer
2203 .conn()
2204 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2205 .expect("clear the completion marker");
2206 writer.conn().execute_batch("COMMIT").unwrap();
2207 }
2208
2209 let store = backend.text(table_key).unwrap();
2210
2211 let map_rows: Vec<(String, i64)> = {
2212 let writer = backend.pool().writer().unwrap();
2213 let mut stmt = writer
2214 .conn()
2215 .prepare(&format!(
2216 "SELECT subject_id, rowid FROM {map} ORDER BY rowid"
2217 ))
2218 .unwrap();
2219 let rows = stmt
2220 .query_map([], |row| Ok((row.get(0)?, row.get(1)?)))
2221 .unwrap()
2222 .collect::<Result<Vec<_>, _>>()
2223 .unwrap();
2224 rows
2225 };
2226 assert_eq!(
2227 map_rows,
2228 vec![(b.to_string(), 7)],
2229 "the stale (A, 7) map row must be removed and replaced by the correct (B, 7) row, \
2230 not left alongside it"
2231 );
2232
2233 assert!(
2234 store.get_document("test_ns", a).await.unwrap().is_none(),
2235 "A's stale map entry is gone, so get_document(A) must find nothing"
2236 );
2237 let fetched_b = store.get_document("test_ns", b).await.unwrap();
2238 assert!(
2239 fetched_b.is_some(),
2240 "get_document(B) must find the live row now correctly mapped"
2241 );
2242 assert_eq!(fetched_b.unwrap().body, "doc b");
2243 }
2244
2245 #[tokio::test]
2252 async fn text_open_sweeps_the_duplicate_that_lost_the_survivor_race() {
2253 let backend = StorageBackend::memory().unwrap();
2254 let table_key = "duplicate_loser_swept";
2255 let table = format!("fts_{table_key}");
2256 let map = format!("{table}_rowids");
2257 let state = format!("{map}_state");
2258
2259 let _ = backend.text(table_key).unwrap();
2260
2261 let dup = uuid::Uuid::new_v4();
2262 {
2263 let writer = backend.pool().writer().unwrap();
2264 writer.conn().execute_batch("BEGIN").unwrap();
2265 writer
2267 .conn()
2268 .execute(
2269 &format!(
2270 "INSERT INTO {table} \
2271 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2272 updated_at, record_kind) \
2273 VALUES (10, ?1, 'note', '', 'older body', '[]', 'test_ns', NULL, 1, \
2274 'memory')"
2275 ),
2276 rusqlite::params![dup.to_string()],
2277 )
2278 .expect("insert the older/losing duplicate");
2279 writer
2281 .conn()
2282 .execute(
2283 &format!(
2284 "INSERT INTO {table} \
2285 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2286 updated_at, record_kind) \
2287 VALUES (20, ?1, 'note', '', 'newer body', '[]', 'test_ns', NULL, 5, \
2288 'memory')"
2289 ),
2290 rusqlite::params![dup.to_string()],
2291 )
2292 .expect("insert the newer/surviving duplicate");
2293 writer
2294 .conn()
2295 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2296 .expect("clear the completion marker");
2297 writer.conn().execute_batch("COMMIT").unwrap();
2298 }
2299
2300 let _ = backend.text(table_key).unwrap();
2301
2302 let writer = backend.pool().writer().unwrap();
2303 let map_rows: Vec<i64> = writer
2304 .conn()
2305 .prepare(&format!(
2306 "SELECT rowid FROM {map} WHERE namespace = 'test_ns' AND subject_id = ?1"
2307 ))
2308 .unwrap()
2309 .query_map(rusqlite::params![dup.to_string()], |row| row.get(0))
2310 .unwrap()
2311 .collect::<Result<Vec<_>, _>>()
2312 .unwrap();
2313 assert_eq!(
2314 map_rows,
2315 vec![20],
2316 "exactly one map row must survive, at the newer (by updated_at) rowid"
2317 );
2318
2319 let fts_rowids: Vec<i64> = writer
2320 .conn()
2321 .prepare(&format!("SELECT rowid FROM {table} ORDER BY rowid"))
2322 .unwrap()
2323 .query_map([], |row| row.get(0))
2324 .unwrap()
2325 .collect::<Result<Vec<_>, _>>()
2326 .unwrap();
2327 assert_eq!(
2328 fts_rowids,
2329 vec![20],
2330 "the losing duplicate (rowid 10) must be deleted from the FTS table itself, not just \
2331 left out of the map as an unmapped live row"
2332 );
2333 }
2334
2335 #[tokio::test]
2342 async fn text_open_after_marker_written_does_not_rescan_even_a_corrupted_map() {
2343 let backend = StorageBackend::memory().unwrap();
2344 let table_key = "marker_no_rescan";
2345 let table = format!("fts_{table_key}");
2346 let map = format!("{table}_rowids");
2347 let state = format!("{map}_state");
2348
2349 let store = backend.text(table_key).unwrap();
2350 store
2351 .upsert_document(khive_storage::types::TextDocument {
2352 subject_id: uuid::Uuid::new_v4(),
2353 kind: khive_types::SubstrateKind::Note,
2354 record_kind: Some("memory".to_string()),
2355 title: None,
2356 body: "seed".to_string(),
2357 tags: vec![],
2358 namespace: "test_ns".to_string(),
2359 metadata: None,
2360 updated_at: chrono::Utc::now(),
2361 })
2362 .await
2363 .unwrap();
2364
2365 let _ = backend.text(table_key).unwrap();
2372
2373 let marked: bool = {
2374 let writer = backend.pool().writer().unwrap();
2375 writer
2376 .conn()
2377 .query_row(
2378 &format!(
2379 "SELECT EXISTS(SELECT 1 FROM {state} WHERE key = 'backfill' AND value = 'complete')"
2380 ),
2381 [],
2382 |row| row.get(0),
2383 )
2384 .unwrap()
2385 };
2386 assert!(
2387 marked,
2388 "a completion marker must exist once the table has held a row"
2389 );
2390
2391 {
2392 let writer = backend.pool().writer().unwrap();
2393 writer
2394 .conn()
2395 .execute(&format!("DELETE FROM {map}"), [])
2396 .expect("corrupt the map by deleting its row directly");
2397 }
2398
2399 let _ = backend.text(table_key).unwrap();
2400
2401 let map_count: i64 = {
2402 let writer = backend.pool().writer().unwrap();
2403 writer
2404 .conn()
2405 .query_row(&format!("SELECT COUNT(*) FROM {map}"), [], |row| row.get(0))
2406 .unwrap()
2407 };
2408 assert_eq!(
2409 map_count, 0,
2410 "a marker-complete table must not be re-scanned on open, even to reconcile a map \
2411 an external actor emptied out from under it"
2412 );
2413 }
2414
2415 #[tokio::test]
2419 async fn text_open_writable_legacy_backfill_excludes_null_key_rows() {
2420 let backend = StorageBackend::memory().unwrap();
2421 let table_key = "legacy_null_key";
2422 let table = format!("fts_{table_key}");
2423 let map = format!("{table}_rowids");
2424 let state = format!("{map}_state");
2425
2426 let _ = backend.text(table_key).unwrap();
2427 {
2428 let writer = backend.pool().writer().unwrap();
2429 writer.conn().execute_batch("BEGIN").unwrap();
2430 writer
2431 .conn()
2432 .execute(
2433 &format!(
2434 "INSERT INTO {table} \
2435 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2436 updated_at, record_kind) \
2437 VALUES (1, NULL, 'note', '', 'null-key body', '[]', NULL, NULL, 0, '')"
2438 ),
2439 [],
2440 )
2441 .expect("insert legacy NULL-key fts row");
2442 writer
2443 .conn()
2444 .execute(
2445 &format!(
2446 "INSERT INTO {table} \
2447 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2448 updated_at, record_kind) \
2449 VALUES (2, 'legacy-1', 'note', '', 'normal body', '[]', 'test_ns', NULL, \
2450 0, 'memory')"
2451 ),
2452 [],
2453 )
2454 .expect("insert legacy normal-key fts row");
2455 writer
2456 .conn()
2457 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2458 .expect("clear the completion marker written for the then-empty table");
2459 writer.conn().execute_batch("COMMIT").unwrap();
2460 }
2461
2462 let _ = backend.text(table_key).unwrap();
2464
2465 let writer = backend.pool().writer().unwrap();
2466 let map_count: i64 = writer
2467 .conn()
2468 .query_row(&format!("SELECT COUNT(*) FROM {map}"), [], |row| row.get(0))
2469 .unwrap();
2470 assert_eq!(map_count, 1, "only the non-NULL-key row may be mapped");
2471 let mapped_subject: String = writer
2472 .conn()
2473 .query_row(&format!("SELECT subject_id FROM {map}"), [], |row| {
2474 row.get(0)
2475 })
2476 .unwrap();
2477 assert_eq!(mapped_subject, "legacy-1");
2478 }
2479
2480 #[tokio::test]
2485 async fn text_open_writable_legacy_backfill_survivor_is_chosen_by_updated_at_not_rowid() {
2486 let backend = StorageBackend::memory().unwrap();
2487 let table_key = "legacy_updated_at_survivor";
2488 let table = format!("fts_{table_key}");
2489 let map = format!("{table}_rowids");
2490 let state = format!("{map}_state");
2491
2492 let _ = backend.text(table_key).unwrap();
2493 {
2494 let writer = backend.pool().writer().unwrap();
2495 writer.conn().execute_batch("BEGIN").unwrap();
2496 writer
2498 .conn()
2499 .execute(
2500 &format!(
2501 "INSERT INTO {table} \
2502 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2503 updated_at, record_kind) \
2504 VALUES (100, 'dup', 'note', '', 'newer body', '[]', 'test_ns', NULL, \
2505 500, 'memory')"
2506 ),
2507 [],
2508 )
2509 .expect("insert newer-but-lower-rowid fts row");
2510 writer
2512 .conn()
2513 .execute(
2514 &format!(
2515 "INSERT INTO {table} \
2516 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2517 updated_at, record_kind) \
2518 VALUES (200, 'dup', 'note', '', 'older body', '[]', 'test_ns', NULL, \
2519 100, 'memory')"
2520 ),
2521 [],
2522 )
2523 .expect("insert older-but-higher-rowid fts row");
2524 writer
2525 .conn()
2526 .execute(&format!("DELETE FROM {state} WHERE key = 'backfill'"), [])
2527 .expect("clear the completion marker written for the then-empty table");
2528 writer.conn().execute_batch("COMMIT").unwrap();
2529 }
2530
2531 let _ = backend.text(table_key).unwrap();
2532
2533 let writer = backend.pool().writer().unwrap();
2534 let mapped_rowid: i64 = writer
2535 .conn()
2536 .query_row(
2537 &format!(
2538 "SELECT rowid FROM {map} WHERE namespace = 'test_ns' AND subject_id = 'dup'"
2539 ),
2540 [],
2541 |row| row.get(0),
2542 )
2543 .expect("read dup's map entry");
2544 assert_eq!(
2545 mapped_rowid, 100,
2546 "the newer document (by updated_at) must survive even though its rowid is lower"
2547 );
2548 }
2549
2550 #[test]
2551 fn invalid_model_key_rejected() {
2552 let backend = StorageBackend::memory().unwrap();
2553 assert!(backend.vectors("bad key!", "bad key!", 3).is_err());
2554 assert!(backend.vectors("", "", 3).is_err());
2555 }
2556
2557 #[test]
2558 fn invalid_table_key_rejected() {
2559 let backend = StorageBackend::memory().unwrap();
2560 assert!(backend.text("bad key!").is_err());
2561 assert!(backend.text("").is_err());
2562 }
2563
2564 #[test]
2570 fn table_key_ending_in_rowids_suffix_rejected() {
2571 let backend = StorageBackend::memory().unwrap();
2572 assert!(backend.text("entities_rowids").is_err());
2573 assert!(backend.text("notes_rowids").is_err());
2574 assert!(backend.text("anything_rowids").is_err());
2575 }
2576
2577 #[test]
2580 fn table_key_containing_but_not_ending_in_rowids_suffix_accepted() {
2581 let backend = StorageBackend::memory().unwrap();
2582 assert!(backend.text("rowids_but_not_at_the_end").is_ok());
2583 }
2584
2585 #[test]
2592 fn table_key_ending_in_rowids_state_suffix_rejected() {
2593 let backend = StorageBackend::memory().unwrap();
2594 assert!(backend.text("entities_rowids_state").is_err());
2595 assert!(backend.text("notes_rowids_state").is_err());
2596 assert!(backend.text("anything_rowids_state").is_err());
2597 }
2598
2599 #[test]
2602 fn table_key_containing_but_not_ending_in_rowids_state_suffix_accepted() {
2603 let backend = StorageBackend::memory().unwrap();
2604 assert!(backend.text("rowids_state_but_not_at_the_end").is_ok());
2605 }
2606
2607 #[tokio::test]
2608 async fn sqlite_read_only_graph_store_rejects_upsert_edge() {
2609 use khive_storage::types::Edge;
2610 use khive_types::EdgeRelation;
2611
2612 let dir = tempfile::tempdir().unwrap();
2613 let path = dir.path().join("ro_graph.db");
2614
2615 {
2617 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
2618 writable.graph().unwrap();
2619 }
2620 #[cfg(unix)]
2621 freeze_snapshot_sidecars(&path);
2622
2623 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
2624 let store = match ro.graph() {
2625 Ok(store) => store,
2626 Err(_) => return,
2629 };
2630
2631 let now = chrono::Utc::now();
2632 let edge = Edge {
2633 id: uuid::Uuid::new_v4().into(),
2634 namespace: "local".to_string(),
2635 source_id: uuid::Uuid::new_v4(),
2636 target_id: uuid::Uuid::new_v4(),
2637 relation: EdgeRelation::Extends,
2638 weight: 0.8,
2639 created_at: now,
2640 updated_at: now,
2641 deleted_at: None,
2642 metadata: None,
2643 target_backend: None,
2644 };
2645
2646 let result = store.upsert_edge(edge).await;
2647 assert!(
2648 result.is_err(),
2649 "upsert_edge on a read-only backend must reject, not silently no-op"
2650 );
2651 }
2652
2653 #[tokio::test]
2654 async fn sqlite_read_only_event_store_rejects_append_event() {
2655 use khive_types::{EventKind, EventOutcome, SubstrateKind};
2656
2657 let dir = tempfile::tempdir().unwrap();
2658 let path = dir.path().join("ro_events.db");
2659
2660 {
2661 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
2662 writable.events().unwrap();
2663 }
2664 #[cfg(unix)]
2665 freeze_snapshot_sidecars(&path);
2666
2667 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
2668 let store = match ro.events() {
2669 Ok(store) => store,
2670 Err(_) => return,
2671 };
2672
2673 let event = khive_storage::event::Event::new(
2674 "local",
2675 "test.verb",
2676 EventKind::Audit,
2677 SubstrateKind::Entity,
2678 "test-actor",
2679 )
2680 .with_outcome(EventOutcome::Success);
2681
2682 let result = store.append_event(event).await;
2683 assert!(
2684 result.is_err(),
2685 "append_event on a read-only backend must reject, not silently no-op"
2686 );
2687 }
2688
2689 #[tokio::test]
2690 async fn sqlite_read_only_text_store_rejects_upsert_document() {
2691 use khive_storage::types::TextDocument;
2692 use khive_types::SubstrateKind;
2693
2694 let dir = tempfile::tempdir().unwrap();
2695 let path = dir.path().join("ro_text.db");
2696
2697 {
2698 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
2699 writable.text("ro_test").unwrap();
2700 }
2701 #[cfg(unix)]
2702 freeze_snapshot_sidecars(&path);
2703
2704 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
2705 let store = match ro.text("ro_test") {
2706 Ok(store) => store,
2707 Err(_) => return,
2708 };
2709
2710 let doc = TextDocument {
2711 subject_id: uuid::Uuid::new_v4(),
2712 kind: SubstrateKind::Entity,
2713 record_kind: None,
2714 title: Some("Title".to_string()),
2715 body: "Body text.".to_string(),
2716 tags: vec![],
2717 namespace: "local".to_string(),
2718 metadata: None,
2719 updated_at: chrono::Utc::now(),
2720 };
2721
2722 let result = store.upsert_document(doc).await;
2723 assert!(
2724 result.is_err(),
2725 "upsert_document on a read-only backend must reject, not silently no-op"
2726 );
2727 }
2728
2729 #[tokio::test]
2736 async fn sqlite_read_only_text_store_without_rowid_map_falls_back_to_scan_predicates() {
2737 let dir = tempfile::tempdir().unwrap();
2738 let path = dir.path().join("ro_text_no_map.db");
2739
2740 let id = uuid::Uuid::new_v4();
2741 {
2742 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
2743 let writer = writable.pool().try_writer().unwrap();
2744 writer
2745 .conn()
2746 .execute_batch(
2747 "CREATE VIRTUAL TABLE IF NOT EXISTS fts_ro_no_map USING fts5(\
2748 subject_id UNINDEXED, kind UNINDEXED, title, body, tags UNINDEXED, \
2749 namespace UNINDEXED, metadata UNINDEXED, updated_at UNINDEXED, \
2750 record_kind, tokenize = 'trigram')",
2751 )
2752 .unwrap();
2753 writer
2754 .conn()
2755 .execute(
2756 "INSERT INTO fts_ro_no_map \
2757 (subject_id, kind, title, body, tags, namespace, metadata, updated_at, \
2758 record_kind) \
2759 VALUES (?1, 'note', '', 'legacy body', '[]', 'local', NULL, 0, NULL)",
2760 rusqlite::params![id.to_string()],
2761 )
2762 .unwrap();
2763 }
2764 #[cfg(unix)]
2765 freeze_snapshot_sidecars(&path);
2766
2767 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
2768 let store = ro
2769 .text("ro_no_map")
2770 .expect("a read-only FTS table with no sidecar map must still open successfully");
2771
2772 let fetched = store
2773 .get_document("local", id)
2774 .await
2775 .expect("scan-fallback get_document must not error against a missing map table");
2776 assert!(
2777 fetched.is_some(),
2778 "scan-fallback get_document must still find the legacy row"
2779 );
2780 assert_eq!(fetched.unwrap().subject_id, id);
2781 }
2782
2783 #[tokio::test]
2793 async fn sqlite_read_only_text_store_with_unmarked_map_falls_back_to_scan_predicates() {
2794 let dir = tempfile::tempdir().unwrap();
2795 let path = dir.path().join("ro_text_unmarked_map.db");
2796
2797 let a = uuid::Uuid::new_v4();
2798 let b = uuid::Uuid::new_v4();
2799 {
2800 let writable = StorageBackend::sqlite_for_test(&path).unwrap();
2801 let writer = writable.pool().try_writer().unwrap();
2802 writer
2803 .conn()
2804 .execute_batch(
2805 "CREATE VIRTUAL TABLE IF NOT EXISTS fts_ro_unmarked USING fts5(\
2806 subject_id UNINDEXED, kind UNINDEXED, title, body, tags UNINDEXED, \
2807 namespace UNINDEXED, metadata UNINDEXED, updated_at UNINDEXED, \
2808 record_kind, tokenize = 'trigram')",
2809 )
2810 .unwrap();
2811 writer
2812 .conn()
2813 .execute_batch(&text::rowid_map_ddl("fts_ro_unmarked"))
2814 .unwrap();
2815 writer
2816 .conn()
2817 .execute(
2818 "INSERT INTO fts_ro_unmarked \
2819 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2820 updated_at, record_kind) \
2821 VALUES (1, ?1, 'note', '', 'doc a', '[]', 'local', NULL, 0, NULL)",
2822 rusqlite::params![a.to_string()],
2823 )
2824 .unwrap();
2825 writer
2826 .conn()
2827 .execute(
2828 "INSERT INTO fts_ro_unmarked \
2829 (rowid, subject_id, kind, title, body, tags, namespace, metadata, \
2830 updated_at, record_kind) \
2831 VALUES (2, ?1, 'note', '', 'doc b', '[]', 'local', NULL, 0, NULL)",
2832 rusqlite::params![b.to_string()],
2833 )
2834 .unwrap();
2835 writer
2838 .conn()
2839 .execute(
2840 "INSERT INTO fts_ro_unmarked_rowids (namespace, subject_id, rowid) \
2841 VALUES ('local', ?1, 2)",
2842 rusqlite::params![b.to_string()],
2843 )
2844 .unwrap();
2845 }
2846 #[cfg(unix)]
2847 freeze_snapshot_sidecars(&path);
2848
2849 let ro = StorageBackend::sqlite_read_only_for_test(&path).unwrap();
2850 let store = ro
2851 .text("ro_unmarked")
2852 .expect("a read-only FTS table with an unmarked map must still open successfully");
2853
2854 let fetched_a = store
2855 .get_document("local", a)
2856 .await
2857 .expect("scan-fallback get_document must not error against an unmarked map");
2858 assert!(
2859 fetched_a.is_some(),
2860 "A has no map entry, so only the scan fallback (not a map join) can find it -- \
2861 proving the unmarked map was not trusted"
2862 );
2863 assert_eq!(fetched_a.unwrap().body, "doc a");
2864 }
2865
2866 #[tokio::test]
2867 async fn blob_store_roundtrip_via_public_api() {
2868 let dir = tempfile::tempdir().unwrap();
2869 let path = dir.path().join("blob_backend.db");
2870 let backend = StorageBackend::sqlite_for_test(&path).unwrap();
2871
2872 let store = backend.blob_store(None, Some(0)).unwrap();
2876 let bytes = b"backend-level blob roundtrip".to_vec();
2877 let content_ref = store.put(bytes.clone()).await.unwrap();
2878 assert_eq!(
2879 store
2880 .get_bounded_verified(&content_ref, bytes.len() as u64)
2881 .await
2882 .unwrap(),
2883 bytes
2884 );
2885 }
2886
2887 #[test]
2888 fn blob_store_defaults_root_beside_db_file() {
2889 let dir = tempfile::tempdir().unwrap();
2890 let path = dir.path().join("blob_default.db");
2891 let backend = StorageBackend::sqlite_for_test(&path).unwrap();
2892
2893 let _store = backend.blob_store(None, None).unwrap();
2897 assert!(
2898 dir.path().join("blobs").is_dir(),
2899 "default root must be created beside the database file"
2900 );
2901 }
2902
2903 #[test]
2904 fn blob_store_errors_for_in_memory_backend_with_no_override() {
2905 let backend = StorageBackend::memory().unwrap();
2906 assert!(backend.blob_store(None, None).is_err());
2907 }
2908
2909 #[test]
2910 fn blob_store_accepts_explicit_root_for_in_memory_backend() {
2911 let dir = tempfile::tempdir().unwrap();
2912 let backend = StorageBackend::memory().unwrap();
2913 let store = backend.blob_store(Some(dir.path()), None);
2914 assert!(store.is_ok());
2915 }
2916
2917 #[test]
2918 fn apply_schema_runs_migrations_idempotently() {
2919 static MIGRATIONS: &[crate::migrations::Migration] = &[crate::migrations::Migration {
2920 id: "001_init",
2921 up_sql: "CREATE TABLE IF NOT EXISTS schema_test (id TEXT PRIMARY KEY);",
2922 down_sql: None,
2923 is_already_applied: None,
2924 }];
2925 let plan = crate::migrations::ServiceSchemaPlan {
2926 service: "schema_test_svc",
2927 sqlite: MIGRATIONS,
2928 postgres: &[],
2929 };
2930
2931 let backend = StorageBackend::memory().unwrap();
2932 backend.apply_schema(&plan).unwrap();
2933 backend.apply_schema(&plan).unwrap();
2934
2935 let reader = backend.pool().reader().unwrap();
2936 let count: i64 = reader
2937 .conn()
2938 .query_row(
2939 "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='schema_test'",
2940 [],
2941 |row| row.get(0),
2942 )
2943 .unwrap();
2944 assert_eq!(count, 1);
2945 }
2946
2947 #[test]
2948 fn pack_ddl_plan_rolls_back_all_statements_on_failure() {
2949 let backend = StorageBackend::memory().unwrap();
2950 let error = backend
2951 .apply_pack_ddl_statements(&[
2952 "CREATE TABLE IF NOT EXISTS pack_schema_first (id INTEGER PRIMARY KEY)",
2953 "CREATE INDEX IF NOT EXISTS pack_schema_second ON pack_schema_missing(id)",
2954 ])
2955 .unwrap_err();
2956
2957 assert!(
2958 error.to_string().contains("pack_schema_missing"),
2959 "schema-plan error must retain the failing SQLite diagnostic: {error}"
2960 );
2961
2962 let reader = backend.pool().reader().unwrap();
2963 let visible_objects: i64 = reader
2964 .conn()
2965 .query_row(
2966 "SELECT COUNT(*) FROM sqlite_master \
2967 WHERE name IN ('pack_schema_first', 'pack_schema_second')",
2968 [],
2969 |row| row.get(0),
2970 )
2971 .unwrap();
2972 assert_eq!(visible_objects, 0);
2973 }
2974
2975 #[test]
2976 fn pack_ddl_plan_applies_all_statements_idempotently() {
2977 const PLAN: &[&str] = &[
2978 "CREATE TABLE IF NOT EXISTS pack_schema_success (id INTEGER PRIMARY KEY, value TEXT)",
2979 "CREATE INDEX IF NOT EXISTS pack_schema_success_value_idx \
2980 ON pack_schema_success(value)",
2981 ];
2982
2983 let backend = StorageBackend::memory().unwrap();
2984 backend.apply_pack_ddl_statements(PLAN).unwrap();
2985 backend.apply_pack_ddl_statements(PLAN).unwrap();
2986
2987 let reader = backend.pool().reader().unwrap();
2988 let visible_objects: i64 = reader
2989 .conn()
2990 .query_row(
2991 "SELECT COUNT(*) FROM sqlite_master \
2992 WHERE name IN ('pack_schema_success', 'pack_schema_success_value_idx')",
2993 [],
2994 |row| row.get(0),
2995 )
2996 .unwrap();
2997 assert_eq!(visible_objects, 2);
2998 }
2999
3000 fn issue_1029_pool(write_queue_enabled: bool) -> (tempfile::TempDir, StorageBackend) {
3009 let dir = tempfile::tempdir().unwrap();
3010 let path = dir.path().join("issue_1029.db");
3011 let config = crate::pool::PoolConfig {
3012 path: Some(path.clone()),
3013 busy_timeout: std::time::Duration::from_millis(200),
3014 write_queue_enabled: Some(write_queue_enabled),
3015 ..crate::pool::PoolConfig::for_test()
3016 };
3017 let pool = ConnectionPool::new(config).expect("fresh tenant-shaped pool should open");
3018 let backend = StorageBackend {
3019 pool: Arc::new(pool),
3020 is_file_backed: true,
3021 path: Some(path),
3022 notes_seq_repair_runs: AtomicUsize::new(0),
3023 };
3024 (dir, backend)
3025 }
3026
3027 async fn issue_1029_create_entity_shaped_sequence(
3028 backend: &StorageBackend,
3029 ) -> Result<(), String> {
3030 let entities = backend
3031 .entities_for_namespace("tenant_ns")
3032 .map_err(|e| format!("entities_for_namespace: {e}"))?;
3033 let entity = khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029Repro");
3034 let entity_id = entity.id;
3035 entities
3036 .upsert_entity(entity)
3037 .await
3038 .map_err(|e| format!("upsert_entity: {e}"))?;
3039
3040 let text = backend.text("entities").map_err(|e| format!("text: {e}"))?;
3041 let doc = khive_storage::types::TextDocument {
3042 subject_id: entity_id,
3043 kind: khive_types::SubstrateKind::Entity,
3044 record_kind: None,
3045 title: Some("Issue1029Repro".to_string()),
3046 body: "issue 1029 repro body".to_string(),
3047 tags: vec![],
3048 namespace: "tenant_ns".to_string(),
3049 metadata: None,
3050 updated_at: chrono::Utc::now(),
3051 };
3052 text.upsert_document(doc)
3053 .await
3054 .map_err(|e| format!("fts_upsert: {e}"))
3055 }
3056
3057 #[tokio::test]
3063 async fn issue_1029_create_entity_shaped_sequence_write_queue_off() {
3064 let (_dir, backend) = issue_1029_pool(false);
3065 let result = issue_1029_create_entity_shaped_sequence(&backend).await;
3066 assert!(
3067 result.is_ok(),
3068 "khive#1029 repro (KHIVE_WRITE_QUEUE off): fts_upsert step failed: {:?}",
3069 result.err()
3070 );
3071 }
3072
3073 #[tokio::test]
3079 async fn issue_1029_create_entity_shaped_sequence_write_queue_on() {
3080 let (_dir, backend) = issue_1029_pool(true);
3081 let result = issue_1029_create_entity_shaped_sequence(&backend).await;
3082 assert!(
3083 result.is_ok(),
3084 "khive#1029 repro (KHIVE_WRITE_QUEUE=1): fts_upsert step failed: {:?}",
3085 result.err()
3086 );
3087 }
3088
3089 #[tokio::test]
3097 async fn issue_1029_two_pools_same_file_write_queue_on() {
3098 let dir = tempfile::tempdir().unwrap();
3099 let path = dir.path().join("issue_1029_two_pools.db");
3100
3101 let cfg = |p: std::path::PathBuf| crate::pool::PoolConfig {
3102 path: Some(p),
3103 busy_timeout: std::time::Duration::from_millis(200),
3104 write_queue_enabled: Some(true),
3105 ..crate::pool::PoolConfig::for_test()
3106 };
3107
3108 let pool_a = ConnectionPool::new(cfg(path.clone())).expect("pool A should open");
3109 let backend_a = StorageBackend {
3110 pool: Arc::new(pool_a),
3111 is_file_backed: true,
3112 path: Some(path.clone()),
3113 notes_seq_repair_runs: AtomicUsize::new(0),
3114 };
3115 let pool_b = ConnectionPool::new(cfg(path.clone())).expect("pool B should open");
3116 let backend_b = StorageBackend {
3117 pool: Arc::new(pool_b),
3118 is_file_backed: true,
3119 path: Some(path),
3120 notes_seq_repair_runs: AtomicUsize::new(0),
3121 };
3122
3123 let entities = backend_a
3124 .entities_for_namespace("tenant_ns")
3125 .expect("entities_for_namespace on pool A");
3126 let entity =
3127 khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029TwoPools");
3128 let entity_id = entity.id;
3129 entities
3130 .upsert_entity(entity)
3131 .await
3132 .expect("pool A entity upsert should succeed");
3133
3134 let text = backend_b.text("entities").expect("text on pool B");
3135 let doc = khive_storage::types::TextDocument {
3136 subject_id: entity_id,
3137 kind: khive_types::SubstrateKind::Entity,
3138 record_kind: None,
3139 title: Some("Issue1029TwoPools".to_string()),
3140 body: "issue 1029 two-pool repro body".to_string(),
3141 tags: vec![],
3142 namespace: "tenant_ns".to_string(),
3143 metadata: None,
3144 updated_at: chrono::Utc::now(),
3145 };
3146 let result = text.upsert_document(doc).await;
3147 assert!(
3148 result.is_ok(),
3149 "khive#1029 two-pool repro: fts_upsert on an independent pool for the \
3150 same tenant DB file failed: {:?}",
3151 result.err()
3152 );
3153 }
3154
3155 struct StarvationCaptureSubscriber {
3158 events: Arc<std::sync::Mutex<Vec<std::collections::BTreeMap<String, String>>>>,
3159 }
3160
3161 impl tracing::Subscriber for StarvationCaptureSubscriber {
3162 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
3163 true
3164 }
3165 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
3166 tracing::span::Id::from_u64(1)
3167 }
3168 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
3169 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
3170 fn event(&self, event: &tracing::Event<'_>) {
3171 #[derive(Default)]
3172 struct FieldVisitor(std::collections::BTreeMap<String, String>);
3173 impl tracing::field::Visit for FieldVisitor {
3174 fn record_debug(
3175 &mut self,
3176 field: &tracing::field::Field,
3177 value: &dyn std::fmt::Debug,
3178 ) {
3179 self.0
3180 .insert(field.name().to_string(), format!("{value:?}"));
3181 }
3182 }
3183 let mut visitor = FieldVisitor::default();
3184 event.record(&mut visitor);
3185 self.events.lock().unwrap().push(visitor.0);
3186 }
3187 fn enter(&self, _: &tracing::span::Id) {}
3188 fn exit(&self, _: &tracing::span::Id) {}
3189 }
3190
3191 #[tokio::test]
3204 #[serial_test::serial(tx_registry)]
3205 async fn issue_1029_starvation_warn_reports_registered_transactions() {
3206 let (_dir, backend) = issue_1029_pool(false);
3207 let text = backend.text("entities").expect("text store");
3210
3211 let holder = backend
3215 .pool
3216 .open_standalone_writer()
3217 .expect("holder connection");
3218 holder
3219 .execute_batch("BEGIN IMMEDIATE")
3220 .expect("holder BEGIN IMMEDIATE");
3221 let fixture =
3222 khive_storage::tx_registry::register(Some("issue_1029_fixture_tx".to_string()));
3223
3224 let events = Arc::new(std::sync::Mutex::new(Vec::new()));
3225 let subscriber = StarvationCaptureSubscriber {
3226 events: Arc::clone(&events),
3227 };
3228 let guard = tracing::subscriber::set_default(subscriber);
3229
3230 let doc = khive_storage::types::TextDocument {
3231 subject_id: uuid::Uuid::new_v4(),
3232 kind: khive_types::SubstrateKind::Entity,
3233 record_kind: None,
3234 title: Some("Issue1029Starved".to_string()),
3235 body: "issue 1029 starvation diagnostic body".to_string(),
3236 tags: vec![],
3237 namespace: "tenant_ns".to_string(),
3238 metadata: None,
3239 updated_at: chrono::Utc::now(),
3240 };
3241 let result = text.upsert_document(doc).await;
3242
3243 drop(guard);
3244 drop(fixture);
3245 holder
3246 .execute_batch("ROLLBACK")
3247 .expect("holder ROLLBACK releases the lock");
3248
3249 assert!(
3250 result.is_err(),
3251 "upsert_document must starve while another connection holds the write lock"
3252 );
3253
3254 let events = events.lock().unwrap();
3255 let warn = events
3256 .iter()
3257 .find(|fields| {
3258 fields
3259 .get("message")
3260 .is_some_and(|m| m.contains("text write starved"))
3261 })
3262 .unwrap_or_else(|| panic!("expected a starvation WARN, captured events: {events:?}"));
3263 assert!(
3264 warn.get("op").is_some_and(|op| op.contains("fts_upsert")),
3265 "WARN must name the starved operation, got: {warn:?}"
3266 );
3267 assert!(
3268 warn.get("open_txs")
3269 .is_some_and(|txs| txs.contains("issue_1029_fixture_tx")),
3270 "WARN must list the registered holder label, got: {warn:?}"
3271 );
3272 let count: usize = warn
3273 .get("open_tx_count")
3274 .expect("WARN must carry open_tx_count")
3275 .parse()
3276 .expect("open_tx_count must be numeric");
3277 assert!(
3278 count >= 1,
3279 "open_tx_count must count the fixture, got {count}"
3280 );
3281 }
3282}