Skip to main content

khive_db/
backend.rs

1//! Concrete storage backend providing capability traits.
2//!
3//! `StorageBackend` owns a `ConnectionPool` and provides factory methods for all
4//! capability traits (`GraphStore`, `NoteStore`, `EventStore`, `VectorStore`,
5//! `TextSearch`, `SqlAccess`). File-backed for production; in-memory for tests.
6
7use std::path::Path;
8use std::sync::atomic::{AtomicUsize, Ordering};
9use std::sync::Arc;
10
11use rusqlite::OptionalExtension;
12
13use crate::error::SqliteError;
14use crate::pool::{ConnectionPool, PoolConfig};
15use crate::sql_bridge::SqlBridge;
16use crate::stores::{blob, entity, event, graph, note, sparse, text, vectors};
17
18/// Concrete storage backend providing capability traits.
19pub struct StorageBackend {
20    pool: Arc<ConnectionPool>,
21    is_file_backed: bool,
22    path: Option<std::path::PathBuf>,
23    /// How many times the lazy `notes_seq` anti-join repair has actually
24    /// executed against this backend's pool. Gates `notes_for_namespace` so
25    /// the repair (a full `notes` scan) runs at most once per backend for
26    /// the process's lifetime instead of on every store acquisition (khive
27    /// #827). Also exposed via
28    /// `notes_seq_repair_run_count` for regression tests.
29    notes_seq_repair_runs: AtomicUsize,
30}
31
32impl StorageBackend {
33    /// File-backed SQLite database.
34    ///
35    /// Opens (or creates) the database at `path`. The underlying pool provides
36    /// 1 writer + N readers in WAL mode for concurrent access.
37    /// No schema is applied — call `apply_schema()` for each service.
38    pub fn sqlite(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
39        crate::extension::ensure_extensions_loaded();
40        let resolved = path.as_ref().to_path_buf();
41        let config = PoolConfig {
42            path: Some(resolved.clone()),
43            ..PoolConfig::default()
44        };
45        let pool = ConnectionPool::new(config)?;
46        Ok(Self {
47            pool: Arc::new(pool),
48            is_file_backed: true,
49            path: Some(resolved),
50            notes_seq_repair_runs: AtomicUsize::new(0),
51        })
52    }
53
54    /// File-backed SQLite database opened read-only.
55    ///
56    /// Opens the database at `path` and sets `PRAGMA query_only = ON` on the
57    /// writer connection so that any write attempt (INSERT/UPDATE/DELETE) returns
58    /// an error. Reader connections are opened with `SQLITE_OPEN_READ_ONLY` by the
59    /// pool; this PRAGMA extends that protection to the writer slot.
60    ///
61    /// The database file must already exist — unlike `sqlite()` this constructor
62    /// does not create a new file.
63    pub fn sqlite_read_only(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
64        crate::extension::ensure_extensions_loaded();
65        let resolved = path.as_ref().to_path_buf();
66        let config = PoolConfig {
67            path: Some(resolved.clone()),
68            read_only: true,
69            ..PoolConfig::default()
70        };
71        // `ConnectionPool::new` opens the writer slot with `SQLITE_OPEN_READ_ONLY`
72        // (no `SQLITE_OPEN_CREATE`) and sets `PRAGMA query_only = ON` on it, so a
73        // missing path is rejected instead of created, and any write attempt is
74        // rejected at the SQLite level regardless of which code path reaches the
75        // writer.
76        let pool = ConnectionPool::new(config)?;
77        Ok(Self {
78            pool: Arc::new(pool),
79            is_file_backed: true,
80            path: Some(resolved),
81            notes_seq_repair_runs: AtomicUsize::new(0),
82        })
83    }
84
85    /// In-memory SQLite database (for tests).
86    ///
87    /// All data is lost when the backend is dropped. The pool degrades to
88    /// single-connection mode since in-memory databases cannot be shared
89    /// across multiple connections.
90    pub fn memory() -> Result<Self, SqliteError> {
91        crate::extension::ensure_extensions_loaded();
92        let config = PoolConfig {
93            path: None,
94            ..PoolConfig::default()
95        };
96        let pool = ConnectionPool::new(config)?;
97        Ok(Self {
98            pool: Arc::new(pool),
99            is_file_backed: false,
100            path: None,
101            notes_seq_repair_runs: AtomicUsize::new(0),
102        })
103    }
104
105    /// Get the SQL access capability.
106    ///
107    /// Returns an `Arc<dyn SqlAccess>` suitable for passing to services.
108    pub fn sql(&self) -> Arc<dyn khive_storage::SqlAccess> {
109        Arc::new(SqlBridge::new(Arc::clone(&self.pool), self.is_file_backed))
110    }
111
112    /// Apply a service's schema plan (run migrations).
113    ///
114    /// Each migration in the plan's `sqlite` list is applied idempotently.
115    /// Already-applied migrations are skipped. The `_schema_versions` table
116    /// tracks which migrations have been run.
117    pub fn apply_schema(
118        &self,
119        plan: &crate::migrations::ServiceSchemaPlan,
120    ) -> Result<(), SqliteError> {
121        let writer = self.pool.try_writer()?;
122        crate::migrations::apply_schema_plan(writer.conn(), plan)
123    }
124
125    /// Apply pack-auxiliary DDL statements.
126    ///
127    /// Executes the full plan in one transaction, applying each DDL statement
128    /// idempotently via `execute_batch`. Each statement MUST be self-contained
129    /// and use `CREATE TABLE IF NOT EXISTS` (or equivalent idempotent DDL) so
130    /// that calling this method more than once does not fail.
131    ///
132    /// Pack auxiliary tables are NOT tracked in `_schema_versions` — they are
133    /// non-versioned. Use `apply_schema` with a `ServiceSchemaPlan` when version
134    /// tracking is needed.
135    ///
136    /// This method is lower-level than `PackRuntime::schema_plan()` — the
137    /// runtime bootstrap calls `pack.schema_plan().statements` and passes the
138    /// slice here. The `SchemaPlan` type lives in `khive-runtime` (above this
139    /// crate in the dep chain); this method accepts a plain `&[&'static str]`
140    /// to avoid a circular dependency.
141    pub fn apply_pack_ddl_statements(
142        &self,
143        statements: &[&'static str],
144    ) -> Result<(), SqliteError> {
145        let writer = self.pool.try_writer()?;
146        writer.transaction(|conn| {
147            for &stmt in statements {
148                conn.execute_batch(stmt)?;
149            }
150            Ok(())
151        })
152    }
153
154    /// Get an EntityStore. Applies the entities DDL if not already present.
155    ///
156    /// Idempotent — safe to call multiple times.
157    pub fn entities(&self) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
158        self.entities_for_namespace("local")
159    }
160
161    /// Get an EntityStore. The namespace parameter is validated (non-empty) and
162    /// the entities schema is applied, but the store itself is unscoped — namespace
163    /// is the caller's responsibility on each query/delete call.
164    pub fn entities_for_namespace(
165        &self,
166        namespace: &str,
167    ) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
168        if namespace.trim().is_empty() {
169            return Err(SqliteError::InvalidData(
170                "entities namespace must be non-empty".to_string(),
171            ));
172        }
173        let writer = self.pool.try_writer()?;
174        entity::ensure_entities_schema(writer.conn())?;
175
176        Ok(Arc::new(entity::SqlEntityStore::new(
177            Arc::clone(&self.pool),
178            self.is_file_backed,
179        )))
180    }
181
182    /// Get a GraphStore for the default namespace.
183    ///
184    /// Creates the `graph_edges` table (with indexes) if it does not already
185    /// exist. Idempotent — safe to call multiple times.
186    pub fn graph(&self) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
187        self.graph_for_namespace("local")
188    }
189
190    /// Get a GraphStore scoped to a namespace.
191    pub fn graph_for_namespace(
192        &self,
193        namespace: &str,
194    ) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
195        if namespace.trim().is_empty() {
196            return Err(SqliteError::InvalidData(
197                "graph namespace must be non-empty".to_string(),
198            ));
199        }
200        let writer = self.pool.try_writer()?;
201        graph::ensure_graph_schema(writer.conn())?;
202
203        Ok(Arc::new(graph::SqlGraphStore::new_scoped(
204            Arc::clone(&self.pool),
205            self.is_file_backed,
206            namespace.trim().to_string(),
207        )))
208    }
209
210    /// Get a NoteStore. Applies the notes DDL if not already present.
211    ///
212    /// Idempotent — safe to call multiple times.
213    pub fn notes(&self) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
214        self.notes_for_namespace("local")
215    }
216
217    /// Get a NoteStore. The namespace parameter is validated (non-empty) and
218    /// the notes schema is applied, but the store itself is unscoped — namespace
219    /// is the caller's responsibility on each query/delete call.
220    pub fn notes_for_namespace(
221        &self,
222        namespace: &str,
223    ) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
224        if namespace.trim().is_empty() {
225            return Err(SqliteError::InvalidData(
226                "notes namespace must be non-empty".to_string(),
227            ));
228        }
229        let writer = self.pool.try_writer()?;
230        note::ensure_notes_schema(writer.conn())?;
231
232        // The anti-join repair is a full `notes` scan -- gate it to run at
233        // most once per backend/pool. `try_writer()` blocks for exclusive
234        // access to the single writer connection for this whole function,
235        // so this load-then-run-then-store is race-free: no other caller on
236        // this pool can observe or advance `notes_seq_repair_runs` while we
237        // hold the writer guard (khive #827).
238        if self.notes_seq_repair_runs.load(Ordering::Relaxed) == 0 {
239            note::repair_notes_seq(writer.conn())?;
240            self.notes_seq_repair_runs.fetch_add(1, Ordering::Relaxed);
241        }
242
243        Ok(Arc::new(note::SqlNoteStore::new(
244            Arc::clone(&self.pool),
245            self.is_file_backed,
246        )))
247    }
248
249    /// How many times the lazy `notes_seq` anti-join repair has actually
250    /// executed against this backend's pool. Exposed for regression tests
251    /// asserting the repair runs at most once per backend for the process's
252    /// lifetime, not once per `notes_for_namespace` call (khive #827).
253    pub fn notes_seq_repair_run_count(&self) -> usize {
254        self.notes_seq_repair_runs.load(Ordering::Relaxed)
255    }
256
257    /// Get an EventStore for the default namespace.
258    ///
259    /// Creates the `events` table (with indexes) if it does not already exist.
260    /// Idempotent — safe to call multiple times.
261    pub fn events(&self) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
262        self.events_for_namespace("local")
263    }
264
265    /// Get an EventStore scoped to a namespace.
266    pub fn events_for_namespace(
267        &self,
268        namespace: &str,
269    ) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
270        if namespace.trim().is_empty() {
271            return Err(SqliteError::InvalidData(
272                "events namespace must be non-empty".to_string(),
273            ));
274        }
275        let writer = self.pool.try_writer()?;
276        event::ensure_events_schema(writer.conn())?;
277
278        Ok(Arc::new(event::SqlEventStore::new_scoped(
279            Arc::clone(&self.pool),
280            self.is_file_backed,
281            namespace.trim().to_string(),
282        )))
283    }
284
285    /// Get a VectorStore for a specific embedding model, scoped to the default namespace.
286    ///
287    /// Creates the vec0 virtual table if it does not already exist. The `model_key`
288    /// must contain only ASCII alphanumeric/underscore characters. The `embedding_model`
289    /// is the canonical display name stored in each vector row.
290    pub fn vectors(
291        &self,
292        model_key: &str,
293        embedding_model: &str,
294        dimensions: usize,
295    ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
296        self.vectors_for_namespace(model_key, embedding_model, dimensions, "local")
297    }
298
299    /// Get a VectorStore for a specific embedding model with a default namespace.
300    ///
301    /// Creates the vec0 virtual table if it does not already exist. The `namespace`
302    /// is a default for trait methods that lack a per-call namespace parameter
303    /// (count, delete, info). Access control is enforced at the runtime layer.
304    ///
305    /// The `model_key` must contain only ASCII alphanumeric/underscore characters.
306    /// The `embedding_model` is the canonical display name stored in the `embedding_model`
307    /// column of each vector row (e.g. `"all-minilm-l6-v2"`).
308    pub fn vectors_for_namespace(
309        &self,
310        model_key: &str,
311        embedding_model: &str,
312        dimensions: usize,
313        namespace: &str,
314    ) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
315        if model_key.is_empty()
316            || !model_key
317                .chars()
318                .all(|c| c.is_ascii_alphanumeric() || c == '_')
319        {
320            return Err(SqliteError::InvalidData(format!(
321                "invalid model_key '{}': must be non-empty and contain only \
322                 alphanumeric/underscore characters",
323                model_key
324            )));
325        }
326        if namespace.trim().is_empty() {
327            return Err(SqliteError::InvalidData(
328                "vector store namespace must be non-empty".to_string(),
329            ));
330        }
331
332        // Ensure sqlite-vec is registered before creating vec0 tables.
333        crate::extension::ensure_extensions_loaded();
334
335        let table = format!("vec_{}", model_key);
336        let writer = self.pool.try_writer()?;
337
338        // Detect old-schema vec0 tables that predate the `field` column.
339        // vec0 virtual tables do not support ALTER TABLE, so we must drop and recreate
340        // the table if it exists without the `field` column. Vector data is a cache —
341        // callers can re-embed from the source record after the table is rebuilt.
342        // Use pragma_table_info to check columns directly; substring matching on the
343        // CREATE DDL is fragile (a model_key containing "field" would false-match).
344        let table_exists: bool = writer
345            .conn()
346            .query_row(
347                "SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
348                rusqlite::params![&table],
349                |row| row.get::<_, i64>(0),
350            )
351            .optional()
352            .map_err(SqliteError::Rusqlite)?
353            .is_some();
354
355        if table_exists {
356            // V17 migration (vector_embedding_model_tag_preserving_rebuild) adds
357            // `field` and `embedding_model` to all pre-existing vec0 tables at
358            // migration time.  If this table still lacks either column post-migration
359            // that indicates the database was not migrated — return a hard error
360            // rather than silently dropping data.
361            let pragma = format!("PRAGMA table_xinfo({})", table);
362            let mut stmt = writer.conn().prepare(&pragma)?;
363            let mut rows = stmt.query([])?;
364            let mut has_field = false;
365            let mut has_embedding_model = false;
366            while let Some(row) = rows.next()? {
367                let name: String = row.get(1)?;
368                if name == "field" {
369                    has_field = true;
370                }
371                if name == "embedding_model" {
372                    has_embedding_model = true;
373                }
374            }
375            if !has_field || !has_embedding_model {
376                return Err(SqliteError::InvalidData(format!(
377                    "vec0 table '{}' is missing required column(s) (field={}, \
378                     embedding_model={}); this is a pre-v0.2.8 vector schema and is \
379                     not supported — recreate the database",
380                    table, has_field, has_embedding_model,
381                )));
382            }
383        }
384
385        // Ensure the _embedding_models registry table exists.
386        // This is a no-op when the table already exists. Running it here ensures
387        // the registry is present for any caller that opens a vector store without
388        // first calling run_migrations() (e.g., tests that create stores directly).
389        // Production callers are expected to call run_migrations() at startup, which
390        // creates the registry via V14; this is a belt-and-suspenders fallback.
391        // Schema is defined in `migrations::EMBEDDING_MODELS_DDL` (single source of
392        // truth) to prevent the two copies from silently drifting.
393        writer
394            .conn()
395            .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
396
397        // Same guarantee for the ANN write log: vector write paths append to it
398        // in the same transaction as the vec0 mutation, so it must exist in any
399        // database that hosts vec_* tables.
400        writer
401            .conn()
402            .execute_batch(crate::migrations::ANN_WRITE_LOG_DDL)?;
403        writer
404            .conn()
405            .execute_batch(crate::migrations::ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL)?;
406
407        // Create the vec0 virtual table. Idempotent on fresh databases and after the
408        // old-schema rebuild above.
409        let ddl = format!(
410            "CREATE VIRTUAL TABLE IF NOT EXISTS vec_{} USING vec0(\
411             subject_id TEXT PRIMARY KEY, \
412             namespace TEXT NOT NULL, \
413             kind TEXT NOT NULL, \
414             field TEXT NOT NULL, \
415             embedding_model TEXT NOT NULL, \
416             embedding float[{}] distance_metric=cosine\
417             )",
418            model_key, dimensions
419        );
420        writer.conn().execute_batch(&ddl)?;
421
422        Ok(Arc::new(vectors::SqliteVecStore::new(
423            Arc::clone(&self.pool),
424            self.is_file_backed,
425            model_key.to_string(),
426            embedding_model.to_string(),
427            dimensions,
428            namespace.trim().to_string(),
429        )?))
430    }
431
432    /// Register an embedding model in the `_embedding_models` registry table.
433    ///
434    /// Idempotent: if a row with the same `canonical_key` already exists, updates its
435    /// status back to `'active'` without changing other fields.
436    pub fn register_embedding_model(
437        &self,
438        engine_name: &str,
439        model_id: &str,
440        key_version: &str,
441        dimensions: u32,
442    ) -> Result<(), SqliteError> {
443        let writer = self.pool.try_writer()?;
444        writer
445            .conn()
446            .execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
447
448        let now = chrono::Utc::now().timestamp_micros();
449        let canonical_key =
450            format!("{engine_name}:{model_id}:{key_version}:{dimensions}").into_bytes();
451        let id = uuid::Uuid::new_v4();
452        writer.conn().execute(
453            "INSERT INTO _embedding_models \
454             (id, engine_name, model_id, key_version, dim, output_dim, status, \
455              activated_at, superseded_at, superseded_by, canonical_key, created_at) \
456             VALUES (?1, ?2, ?3, ?4, ?5, NULL, 'active', ?6, NULL, NULL, ?7, ?8) \
457             ON CONFLICT(canonical_key) DO UPDATE SET \
458                status = 'active', \
459                activated_at = COALESCE(_embedding_models.activated_at, excluded.activated_at)",
460            rusqlite::params![
461                id.as_bytes().as_slice(),
462                engine_name,
463                model_id,
464                key_version,
465                dimensions as i64,
466                now,
467                canonical_key,
468                now,
469            ],
470        )?;
471        Ok(())
472    }
473
474    /// Get a SparseStore for a specific model key, scoped to the default namespace.
475    ///
476    /// Creates the sparse table if it does not already exist.
477    pub fn sparse(
478        &self,
479        model_key: &str,
480    ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
481        self.sparse_for_namespace(model_key, "local")
482    }
483
484    /// Get a SparseStore for a specific model key with an explicit default namespace.
485    ///
486    /// The `model_key` must contain only ASCII alphanumeric/underscore characters.
487    pub fn sparse_for_namespace(
488        &self,
489        model_key: &str,
490        namespace: &str,
491    ) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
492        if model_key.is_empty()
493            || !model_key
494                .chars()
495                .all(|c| c.is_ascii_alphanumeric() || c == '_')
496        {
497            return Err(SqliteError::InvalidData(format!(
498                "invalid model_key '{}': must be non-empty and contain only alphanumeric/underscore characters",
499                model_key
500            )));
501        }
502        if namespace.trim().is_empty() {
503            return Err(SqliteError::InvalidData(
504                "sparse store namespace must be non-empty".to_string(),
505            ));
506        }
507
508        let writer = self.pool.try_writer()?;
509        sparse::ensure_sparse_schema(writer.conn(), model_key).map_err(SqliteError::Rusqlite)?;
510
511        Ok(Arc::new(sparse::SqliteSparseStore::new(
512            Arc::clone(&self.pool),
513            self.is_file_backed,
514            model_key.to_string(),
515            namespace.trim().to_string(),
516        )?))
517    }
518
519    /// Get a TextSearch for a specific table key.
520    ///
521    /// Creates the FTS5 virtual table if it does not already exist. Uses the
522    /// `trigram` tokenizer by default (CJK-safe).
523    ///
524    /// The `table_key` must contain only ASCII alphanumeric/underscore characters.
525    pub fn text(&self, table_key: &str) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
526        self.text_with_tokenizer(table_key, "trigram")
527    }
528
529    /// Get a TextSearch with an explicit FTS5 tokenizer.
530    ///
531    /// Use when you need a tokenizer other than the default `trigram` — for
532    /// example `unicode61` for Latin-only corpora.
533    ///
534    /// Both `table_key` and `tokenizer` must contain only ASCII
535    /// alphanumeric/underscore characters.
536    pub fn text_with_tokenizer(
537        &self,
538        table_key: &str,
539        tokenizer: &str,
540    ) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
541        if table_key.is_empty()
542            || !table_key
543                .chars()
544                .all(|c| c.is_ascii_alphanumeric() || c == '_')
545        {
546            return Err(SqliteError::InvalidData(format!(
547                "invalid table_key '{}': must be non-empty and contain only \
548                 alphanumeric/underscore characters",
549                table_key
550            )));
551        }
552        if tokenizer.is_empty()
553            || !tokenizer
554                .chars()
555                .all(|c| c.is_ascii_alphanumeric() || c == '_')
556        {
557            return Err(SqliteError::InvalidData(format!(
558                "invalid tokenizer '{}': must be non-empty and contain only \
559                 alphanumeric/underscore characters",
560                tokenizer
561            )));
562        }
563
564        let ddl = format!(
565            "CREATE VIRTUAL TABLE IF NOT EXISTS fts_{} USING fts5(\
566             subject_id UNINDEXED, \
567             kind UNINDEXED, \
568             title, \
569             body, \
570             tags UNINDEXED, \
571             namespace UNINDEXED, \
572             metadata UNINDEXED, \
573             updated_at UNINDEXED, \
574             tokenize = '{}'\
575             )",
576            table_key, tokenizer
577        );
578        let writer = self.pool.try_writer()?;
579        writer.conn().execute_batch(&ddl)?;
580
581        Ok(Arc::new(text::Fts5TextSearch::new(
582            Arc::clone(&self.pool),
583            self.is_file_backed,
584            table_key.to_string(),
585        )))
586    }
587
588    /// Get a `BlobStore` rooted per khive#292's precedence chain:
589    /// `KHIVE_BLOB_ROOT` env var > `config_root` (a caller-resolved
590    /// `khive.toml` override — `khive-db` has no TOML parser of its own) >
591    /// beside this backend's database directory. `floor_bytes` overrides the
592    /// default 100 GB fail-closed free-space floor (`None` keeps the
593    /// default). Errors if none of the three roots apply — e.g. an in-memory
594    /// backend with no override and no env var has nowhere to default to.
595    pub fn blob_store(
596        &self,
597        config_root: Option<&Path>,
598        floor_bytes: Option<u64>,
599    ) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
600        let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
601        let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
602        Ok(Arc::new(blob::FsBlobStore::new(root, floor)?))
603    }
604
605    /// Is this a file-backed backend?
606    pub fn is_file_backed(&self) -> bool {
607        self.is_file_backed
608    }
609
610    /// Return the directory containing the backend's database file, or `None`
611    /// for an in-memory backend.
612    pub fn data_dir(&self) -> Option<std::path::PathBuf> {
613        self.path.as_ref()?.parent().map(|p| p.to_path_buf())
614    }
615
616    /// Root directory for this database's ANN segment tree, or `None` for an
617    /// in-memory backend. Derived from the database file name itself
618    /// (`<db-file>.ann/` beside the file), so two databases sharing a parent
619    /// directory can never adopt each other's segments or UUID maps. The
620    /// suffix is appended at the `OsString` byte level — a lossy UTF-8
621    /// conversion would collapse distinct non-UTF-8 filenames into one
622    /// replacement-character root, breaking exactly that isolation.
623    pub fn ann_root(&self) -> Option<std::path::PathBuf> {
624        ann_root_for(self.path.as_ref()?)
625    }
626
627    /// Access the underlying pool (escape hatch).
628    pub fn pool(&self) -> &ConnectionPool {
629        &self.pool
630    }
631
632    /// Clone the underlying pool Arc.
633    pub fn pool_arc(&self) -> Arc<ConnectionPool> {
634        Arc::clone(&self.pool)
635    }
636}
637
638/// `<db-file>.ann` sibling of a database file, appended at the `OsString`
639/// byte level: a lossy UTF-8 conversion would collapse distinct non-UTF-8
640/// filenames into one replacement-character root, breaking the per-database
641/// segment isolation that `ann_root` exists to guarantee.
642fn ann_root_for(path: &std::path::Path) -> Option<std::path::PathBuf> {
643    let mut file = path.file_name()?.to_os_string();
644    file.push(".ann");
645    path.parent().map(|p| p.join(file))
646}
647
648#[cfg(test)]
649mod tests {
650    use super::*;
651    use khive_storage::types::{SqlStatement, SqlValue};
652
653    #[test]
654    fn memory_backend_creates_successfully() {
655        let backend = StorageBackend::memory().expect("memory backend should create");
656        assert!(!backend.is_file_backed());
657    }
658
659    #[test]
660    fn file_backend_creates_successfully() {
661        let dir = tempfile::tempdir().unwrap();
662        let path = dir.path().join("test.db");
663        let backend = StorageBackend::sqlite(&path).expect("file backend should create");
664        assert!(backend.is_file_backed());
665        assert!(path.exists());
666    }
667
668    #[test]
669    fn data_dir_returns_none_for_memory_backend() {
670        let backend = StorageBackend::memory().expect("memory backend");
671        assert!(backend.data_dir().is_none());
672    }
673
674    #[test]
675    fn data_dir_returns_parent_dir_for_file_backend() {
676        let dir = tempfile::tempdir().unwrap();
677        let path = dir.path().join("data.db");
678        let backend = StorageBackend::sqlite(&path).expect("file backend");
679        let got = backend.data_dir().expect("file backend must return Some");
680        assert_eq!(got, dir.path());
681    }
682
683    #[test]
684    fn ann_root_is_database_scoped_sibling_dir() {
685        let dir = tempfile::tempdir().unwrap();
686        let path = dir.path().join("data.db");
687        let backend = StorageBackend::sqlite(&path).expect("file backend");
688        let got = backend.ann_root().expect("file backend must return Some");
689        assert_eq!(got, dir.path().join("data.db.ann"));
690        assert!(StorageBackend::memory().unwrap().ann_root().is_none());
691    }
692
693    /// Two distinct non-UTF-8 database filenames must never share an ANN
694    /// root: a lossy UTF-8 conversion collapses both to the replacement
695    /// character, letting one database adopt the other's segments. Exercised
696    /// on the path derivation directly — APFS (macOS CI) refuses to create
697    /// files with non-UTF-8 names, so a real backend cannot be opened there.
698    #[cfg(unix)]
699    #[test]
700    fn ann_root_distinct_for_non_utf8_filenames() {
701        use std::os::unix::ffi::OsStrExt;
702        let path_a = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xff.db"));
703        let path_b = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xfe.db"));
704        let root_a = ann_root_for(&path_a).expect("Some for a file path");
705        let root_b = ann_root_for(&path_b).expect("Some for a file path");
706        assert_ne!(
707            root_a, root_b,
708            "distinct database files must map to distinct ANN roots"
709        );
710    }
711
712    #[tokio::test]
713    async fn sql_access_memory_roundtrip() {
714        let backend = StorageBackend::memory().unwrap();
715        let sql = backend.sql();
716
717        let mut writer = sql.writer().await.unwrap();
718        writer
719            .execute_script(
720                "CREATE TABLE test_rt (id TEXT PRIMARY KEY, value INTEGER NOT NULL)".into(),
721            )
722            .await
723            .unwrap();
724
725        let affected = writer
726            .execute(SqlStatement {
727                sql: "INSERT INTO test_rt (id, value) VALUES (?1, ?2)".into(),
728                params: vec![SqlValue::Text("row1".into()), SqlValue::Integer(42)],
729                label: None,
730            })
731            .await
732            .unwrap();
733        assert_eq!(affected, 1);
734
735        let mut reader = sql.reader().await.unwrap();
736        let row = reader
737            .query_row(SqlStatement {
738                sql: "SELECT id, value FROM test_rt WHERE id = ?1".into(),
739                params: vec![SqlValue::Text("row1".into())],
740                label: None,
741            })
742            .await
743            .unwrap();
744
745        let row = row.expect("should find the inserted row");
746        assert_eq!(row.columns.len(), 2);
747        match &row.columns[0].value {
748            SqlValue::Text(s) => assert_eq!(s, "row1"),
749            other => panic!("expected Text, got {other:?}"),
750        }
751        match &row.columns[1].value {
752            SqlValue::Integer(v) => assert_eq!(*v, 42),
753            other => panic!("expected Integer, got {other:?}"),
754        }
755    }
756
757    #[tokio::test]
758    async fn sql_access_file_roundtrip() {
759        let dir = tempfile::tempdir().unwrap();
760        let path = dir.path().join("test_roundtrip.db");
761        let backend = StorageBackend::sqlite(&path).unwrap();
762        let sql = backend.sql();
763
764        let mut writer = sql.writer().await.unwrap();
765        writer
766            .execute_script("CREATE TABLE test_f (k TEXT PRIMARY KEY, v TEXT)".into())
767            .await
768            .unwrap();
769        writer
770            .execute(SqlStatement {
771                sql: "INSERT INTO test_f (k, v) VALUES (?1, ?2)".into(),
772                params: vec![
773                    SqlValue::Text("hello".into()),
774                    SqlValue::Text("world".into()),
775                ],
776                label: None,
777            })
778            .await
779            .unwrap();
780
781        let mut reader = sql.reader().await.unwrap();
782        let rows = reader
783            .query_all(SqlStatement {
784                sql: "SELECT k, v FROM test_f".into(),
785                params: vec![],
786                label: None,
787            })
788            .await
789            .unwrap();
790        assert_eq!(rows.len(), 1);
791        match &rows[0].columns[1].value {
792            SqlValue::Text(s) => assert_eq!(s, "world"),
793            other => panic!("expected Text, got {other:?}"),
794        }
795    }
796
797    #[test]
798    fn sqlite_read_only_missing_path_does_not_create_file() {
799        let dir = tempfile::tempdir().unwrap();
800        let path = dir.path().join("missing_ro.db");
801        assert!(!path.exists());
802
803        let result = StorageBackend::sqlite_read_only(&path);
804        assert!(
805            result.is_err(),
806            "opening a missing path read-only must fail"
807        );
808        assert!(
809            !path.exists(),
810            "opening a missing path read-only must not create the file"
811        );
812    }
813
814    #[tokio::test]
815    async fn sqlite_read_only_sql_writer_rejects_ddl_and_insert() {
816        let dir = tempfile::tempdir().unwrap();
817        let path = dir.path().join("ro_writer.db");
818
819        // Create the database and a table while writable.
820        {
821            let writable = StorageBackend::sqlite(&path).unwrap();
822            let sql = writable.sql();
823            let mut writer = sql.writer().await.unwrap();
824            writer
825                .execute_script("CREATE TABLE ro_existing (id INTEGER PRIMARY KEY)".into())
826                .await
827                .unwrap();
828        }
829
830        let ro = StorageBackend::sqlite_read_only(&path).unwrap();
831        let sql = ro.sql();
832
833        // Writer acquisition itself must fail for a read-only backend.
834        let writer_result = sql.writer().await;
835        assert!(
836            writer_result.is_err(),
837            "sql().writer() must be rejected on a read-only backend"
838        );
839    }
840
841    #[tokio::test]
842    #[cfg(feature = "vectors")]
843    async fn vectors_roundtrip_via_public_api() {
844        let backend = StorageBackend::memory().unwrap();
845        let store = backend.vectors("test_api", "test_api", 3).unwrap();
846
847        let id = uuid::Uuid::new_v4();
848        store
849            .insert(
850                id,
851                khive_types::SubstrateKind::Entity,
852                "local",
853                "content",
854                vec![vec![1.0, 0.0, 0.0]],
855            )
856            .await
857            .unwrap();
858
859        let hits = store
860            .search(khive_storage::types::VectorSearchRequest {
861                query_vectors: vec![vec![1.0, 0.0, 0.0]],
862                top_k: 1,
863                namespace: None,
864                kind: None,
865                embedding_model: None,
866                filter: None,
867                backend_hints: None,
868            })
869            .await
870            .unwrap();
871
872        assert_eq!(hits.len(), 1);
873        assert_eq!(hits[0].subject_id, id);
874        assert!(hits[0].score.to_f64() > 0.99);
875    }
876
877    #[tokio::test]
878    #[cfg(feature = "vectors")]
879    async fn vectors_creates_table_idempotently() {
880        let backend = StorageBackend::memory().unwrap();
881
882        let store1 = backend.vectors("idempotent", "idempotent", 3).unwrap();
883        let store2 = backend.vectors("idempotent", "idempotent", 3).unwrap();
884
885        let id = uuid::Uuid::new_v4();
886        store1
887            .insert(
888                id,
889                khive_types::SubstrateKind::Entity,
890                "local",
891                "content",
892                vec![vec![1.0, 0.0, 0.0]],
893            )
894            .await
895            .unwrap();
896
897        let count = store2.count().await.unwrap();
898        assert_eq!(count, 1);
899    }
900
901    #[tokio::test]
902    async fn text_roundtrip_via_public_api() {
903        let backend = StorageBackend::memory().unwrap();
904        let store = backend.text("test_api").unwrap();
905
906        let id = uuid::Uuid::new_v4();
907        let doc = khive_storage::types::TextDocument {
908            subject_id: id,
909            kind: khive_types::SubstrateKind::Entity,
910            title: Some("Test Title".to_string()),
911            body: "This is a searchable document about Rust.".to_string(),
912            tags: vec!["rust".to_string()],
913            namespace: "test_ns".to_string(),
914            metadata: None,
915            updated_at: chrono::Utc::now(),
916        };
917        store.upsert_document(doc).await.unwrap();
918
919        let hits = store
920            .search(khive_storage::types::TextSearchRequest {
921                query: "Rust".to_string(),
922                mode: khive_storage::types::TextQueryMode::Plain,
923                filter: Some(khive_storage::types::TextFilter {
924                    namespaces: vec!["test_ns".to_string()],
925                    ..Default::default()
926                }),
927                top_k: 1,
928                snippet_chars: 64,
929            })
930            .await
931            .unwrap();
932
933        assert_eq!(hits.len(), 1);
934        assert_eq!(hits[0].subject_id, id);
935        assert!(hits[0].score.to_f64() > 0.0);
936    }
937
938    #[tokio::test]
939    async fn text_creates_table_idempotently() {
940        let backend = StorageBackend::memory().unwrap();
941
942        let store1 = backend.text("idempotent_fts").unwrap();
943        let store2 = backend.text("idempotent_fts").unwrap();
944
945        let id = uuid::Uuid::new_v4();
946        let doc = khive_storage::types::TextDocument {
947            subject_id: id,
948            kind: khive_types::SubstrateKind::Note,
949            title: None,
950            body: "Hello world.".to_string(),
951            tags: vec![],
952            namespace: "test_ns".to_string(),
953            metadata: None,
954            updated_at: chrono::Utc::now(),
955        };
956        store1.upsert_document(doc).await.unwrap();
957
958        let count = store2
959            .count(khive_storage::types::TextFilter {
960                namespaces: vec!["test_ns".to_string()],
961                ..Default::default()
962            })
963            .await
964            .unwrap();
965        assert_eq!(count, 1);
966    }
967
968    #[test]
969    fn invalid_model_key_rejected() {
970        let backend = StorageBackend::memory().unwrap();
971        assert!(backend.vectors("bad key!", "bad key!", 3).is_err());
972        assert!(backend.vectors("", "", 3).is_err());
973    }
974
975    #[test]
976    fn invalid_table_key_rejected() {
977        let backend = StorageBackend::memory().unwrap();
978        assert!(backend.text("bad key!").is_err());
979        assert!(backend.text("").is_err());
980    }
981
982    #[tokio::test]
983    async fn sqlite_read_only_graph_store_rejects_upsert_edge() {
984        use khive_storage::types::Edge;
985        use khive_types::EdgeRelation;
986
987        let dir = tempfile::tempdir().unwrap();
988        let path = dir.path().join("ro_graph.db");
989
990        // Create the database and the graph schema while writable.
991        {
992            let writable = StorageBackend::sqlite(&path).unwrap();
993            writable.graph().unwrap();
994        }
995
996        let ro = StorageBackend::sqlite_read_only(&path).unwrap();
997        let store = match ro.graph() {
998            Ok(store) => store,
999            // Failing to even open the store on a read-only backend is an
1000            // acceptable rejection — the write path never becomes reachable.
1001            Err(_) => return,
1002        };
1003
1004        let now = chrono::Utc::now();
1005        let edge = Edge {
1006            id: uuid::Uuid::new_v4().into(),
1007            namespace: "local".to_string(),
1008            source_id: uuid::Uuid::new_v4(),
1009            target_id: uuid::Uuid::new_v4(),
1010            relation: EdgeRelation::Extends,
1011            weight: 0.8,
1012            created_at: now,
1013            updated_at: now,
1014            deleted_at: None,
1015            metadata: None,
1016            target_backend: None,
1017        };
1018
1019        let result = store.upsert_edge(edge).await;
1020        assert!(
1021            result.is_err(),
1022            "upsert_edge on a read-only backend must reject, not silently no-op"
1023        );
1024    }
1025
1026    #[tokio::test]
1027    async fn sqlite_read_only_event_store_rejects_append_event() {
1028        use khive_types::{EventKind, EventOutcome, SubstrateKind};
1029
1030        let dir = tempfile::tempdir().unwrap();
1031        let path = dir.path().join("ro_events.db");
1032
1033        {
1034            let writable = StorageBackend::sqlite(&path).unwrap();
1035            writable.events().unwrap();
1036        }
1037
1038        let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1039        let store = match ro.events() {
1040            Ok(store) => store,
1041            Err(_) => return,
1042        };
1043
1044        let event = khive_storage::event::Event::new(
1045            "local",
1046            "test.verb",
1047            EventKind::Audit,
1048            SubstrateKind::Entity,
1049            "test-actor",
1050        )
1051        .with_outcome(EventOutcome::Success);
1052
1053        let result = store.append_event(event).await;
1054        assert!(
1055            result.is_err(),
1056            "append_event on a read-only backend must reject, not silently no-op"
1057        );
1058    }
1059
1060    #[tokio::test]
1061    async fn sqlite_read_only_text_store_rejects_upsert_document() {
1062        use khive_storage::types::TextDocument;
1063        use khive_types::SubstrateKind;
1064
1065        let dir = tempfile::tempdir().unwrap();
1066        let path = dir.path().join("ro_text.db");
1067
1068        {
1069            let writable = StorageBackend::sqlite(&path).unwrap();
1070            writable.text("ro_test").unwrap();
1071        }
1072
1073        let ro = StorageBackend::sqlite_read_only(&path).unwrap();
1074        let store = match ro.text("ro_test") {
1075            Ok(store) => store,
1076            Err(_) => return,
1077        };
1078
1079        let doc = TextDocument {
1080            subject_id: uuid::Uuid::new_v4(),
1081            kind: SubstrateKind::Entity,
1082            title: Some("Title".to_string()),
1083            body: "Body text.".to_string(),
1084            tags: vec![],
1085            namespace: "local".to_string(),
1086            metadata: None,
1087            updated_at: chrono::Utc::now(),
1088        };
1089
1090        let result = store.upsert_document(doc).await;
1091        assert!(
1092            result.is_err(),
1093            "upsert_document on a read-only backend must reject, not silently no-op"
1094        );
1095    }
1096
1097    #[tokio::test]
1098    async fn blob_store_roundtrip_via_public_api() {
1099        let dir = tempfile::tempdir().unwrap();
1100        let path = dir.path().join("blob_backend.db");
1101        let backend = StorageBackend::sqlite(&path).unwrap();
1102
1103        // Explicit floor_bytes=0, not the default 100GB — the free space on
1104        // whatever volume runs this test is not this test's concern (and a
1105        // dev machine or CI runner legitimately may not clear 100GB free).
1106        let store = backend.blob_store(None, Some(0)).unwrap();
1107        let bytes = b"backend-level blob roundtrip".to_vec();
1108        let content_ref = store.put(bytes.clone()).await.unwrap();
1109        assert_eq!(store.get(&content_ref).await.unwrap(), bytes);
1110    }
1111
1112    #[test]
1113    fn blob_store_defaults_root_beside_db_file() {
1114        let dir = tempfile::tempdir().unwrap();
1115        let path = dir.path().join("blob_default.db");
1116        let backend = StorageBackend::sqlite(&path).unwrap();
1117
1118        // `blob_store` creates the root directory eagerly (`FsBlobStore::new`),
1119        // so its existence at the expected default path is directly
1120        // observable without reaching into the trait object.
1121        let _store = backend.blob_store(None, None).unwrap();
1122        assert!(
1123            dir.path().join("blobs").is_dir(),
1124            "default root must be created beside the database file"
1125        );
1126    }
1127
1128    #[test]
1129    fn blob_store_errors_for_in_memory_backend_with_no_override() {
1130        let backend = StorageBackend::memory().unwrap();
1131        assert!(backend.blob_store(None, None).is_err());
1132    }
1133
1134    #[test]
1135    fn blob_store_accepts_explicit_root_for_in_memory_backend() {
1136        let dir = tempfile::tempdir().unwrap();
1137        let backend = StorageBackend::memory().unwrap();
1138        let store = backend.blob_store(Some(dir.path()), None);
1139        assert!(store.is_ok());
1140    }
1141
1142    #[test]
1143    fn apply_schema_runs_migrations_idempotently() {
1144        static MIGRATIONS: &[crate::migrations::Migration] = &[crate::migrations::Migration {
1145            id: "001_init",
1146            up_sql: "CREATE TABLE IF NOT EXISTS schema_test (id TEXT PRIMARY KEY);",
1147            down_sql: None,
1148            is_already_applied: None,
1149        }];
1150        let plan = crate::migrations::ServiceSchemaPlan {
1151            service: "schema_test_svc",
1152            sqlite: MIGRATIONS,
1153            postgres: &[],
1154        };
1155
1156        let backend = StorageBackend::memory().unwrap();
1157        backend.apply_schema(&plan).unwrap();
1158        backend.apply_schema(&plan).unwrap();
1159
1160        let reader = backend.pool().reader().unwrap();
1161        let count: i64 = reader
1162            .conn()
1163            .query_row(
1164                "SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='schema_test'",
1165                [],
1166                |row| row.get(0),
1167            )
1168            .unwrap();
1169        assert_eq!(count, 1);
1170    }
1171
1172    #[test]
1173    fn pack_ddl_plan_rolls_back_all_statements_on_failure() {
1174        let backend = StorageBackend::memory().unwrap();
1175        let error = backend
1176            .apply_pack_ddl_statements(&[
1177                "CREATE TABLE IF NOT EXISTS pack_schema_first (id INTEGER PRIMARY KEY)",
1178                "CREATE INDEX IF NOT EXISTS pack_schema_second ON pack_schema_missing(id)",
1179            ])
1180            .unwrap_err();
1181
1182        assert!(
1183            error.to_string().contains("pack_schema_missing"),
1184            "schema-plan error must retain the failing SQLite diagnostic: {error}"
1185        );
1186
1187        let reader = backend.pool().reader().unwrap();
1188        let visible_objects: i64 = reader
1189            .conn()
1190            .query_row(
1191                "SELECT COUNT(*) FROM sqlite_master \
1192                 WHERE name IN ('pack_schema_first', 'pack_schema_second')",
1193                [],
1194                |row| row.get(0),
1195            )
1196            .unwrap();
1197        assert_eq!(visible_objects, 0);
1198    }
1199
1200    #[test]
1201    fn pack_ddl_plan_applies_all_statements_idempotently() {
1202        const PLAN: &[&str] = &[
1203            "CREATE TABLE IF NOT EXISTS pack_schema_success (id INTEGER PRIMARY KEY, value TEXT)",
1204            "CREATE INDEX IF NOT EXISTS pack_schema_success_value_idx \
1205             ON pack_schema_success(value)",
1206        ];
1207
1208        let backend = StorageBackend::memory().unwrap();
1209        backend.apply_pack_ddl_statements(PLAN).unwrap();
1210        backend.apply_pack_ddl_statements(PLAN).unwrap();
1211
1212        let reader = backend.pool().reader().unwrap();
1213        let visible_objects: i64 = reader
1214            .conn()
1215            .query_row(
1216                "SELECT COUNT(*) FROM sqlite_master \
1217                 WHERE name IN ('pack_schema_success', 'pack_schema_success_value_idx')",
1218                [],
1219                |row| row.get(0),
1220            )
1221            .unwrap();
1222        assert_eq!(visible_objects, 2);
1223    }
1224
1225    /// khive#1029 repro: a `create_entity`-shaped write sequence (entity
1226    /// upsert, then FTS `upsert_document` on the SAME file-backed DB, SAME
1227    /// `StorageBackend`/pool) against a fresh tenant DB file, with a short
1228    /// `busy_timeout` so a genuine lock hang fails fast instead of burning
1229    /// 30s. Runs with `write_queue_enabled: false` — the legacy pool-mutex /
1230    /// standalone-connection path (`KHIVE_WRITE_QUEUE` unset/0 in the
1231    /// hosted symptom report is one of the two configs to check; see the
1232    /// `_write_queue_enabled` sibling below for the flag-on config).
1233    fn issue_1029_pool(write_queue_enabled: bool) -> (tempfile::TempDir, StorageBackend) {
1234        let dir = tempfile::tempdir().unwrap();
1235        let path = dir.path().join("issue_1029.db");
1236        let config = crate::pool::PoolConfig {
1237            path: Some(path.clone()),
1238            busy_timeout: std::time::Duration::from_millis(200),
1239            write_queue_enabled,
1240            ..crate::pool::PoolConfig::default()
1241        };
1242        let pool = ConnectionPool::new(config).expect("fresh tenant-shaped pool should open");
1243        let backend = StorageBackend {
1244            pool: Arc::new(pool),
1245            is_file_backed: true,
1246            path: Some(path),
1247            notes_seq_repair_runs: AtomicUsize::new(0),
1248        };
1249        (dir, backend)
1250    }
1251
1252    async fn issue_1029_create_entity_shaped_sequence(
1253        backend: &StorageBackend,
1254    ) -> Result<(), String> {
1255        let entities = backend
1256            .entities_for_namespace("tenant_ns")
1257            .map_err(|e| format!("entities_for_namespace: {e}"))?;
1258        let entity = khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029Repro");
1259        let entity_id = entity.id;
1260        entities
1261            .upsert_entity(entity)
1262            .await
1263            .map_err(|e| format!("upsert_entity: {e}"))?;
1264
1265        let text = backend.text("entities").map_err(|e| format!("text: {e}"))?;
1266        let doc = khive_storage::types::TextDocument {
1267            subject_id: entity_id,
1268            kind: khive_types::SubstrateKind::Entity,
1269            title: Some("Issue1029Repro".to_string()),
1270            body: "issue 1029 repro body".to_string(),
1271            tags: vec![],
1272            namespace: "tenant_ns".to_string(),
1273            metadata: None,
1274            updated_at: chrono::Utc::now(),
1275        };
1276        text.upsert_document(doc)
1277            .await
1278            .map_err(|e| format!("fts_upsert: {e}"))
1279    }
1280
1281    /// khive#1029 H1/H2 control: `KHIVE_WRITE_QUEUE` unset (legacy pool-mutex
1282    /// / standalone-connection path for both stores, sharing ONE
1283    /// `ConnectionPool` via ONE `StorageBackend` — the topology this test
1284    /// exists to confirm or kill as the lock source, isolated from any
1285    /// multi-pool or multi-backend wiring question).
1286    #[tokio::test]
1287    async fn issue_1029_create_entity_shaped_sequence_write_queue_off() {
1288        let (_dir, backend) = issue_1029_pool(false);
1289        let result = issue_1029_create_entity_shaped_sequence(&backend).await;
1290        assert!(
1291            result.is_ok(),
1292            "khive#1029 repro (KHIVE_WRITE_QUEUE off): fts_upsert step failed: {:?}",
1293            result.err()
1294        );
1295    }
1296
1297    /// khive#1029 H1 direct test: `KHIVE_WRITE_QUEUE=1`, single shared
1298    /// `ConnectionPool`/`StorageBackend` (so the pool-wide `WriterTask` is
1299    /// shared by construction) — isolates whether the WriterTask's
1300    /// transaction lifecycle itself (not a multi-pool topology) is the lock
1301    /// source.
1302    #[tokio::test]
1303    async fn issue_1029_create_entity_shaped_sequence_write_queue_on() {
1304        let (_dir, backend) = issue_1029_pool(true);
1305        let result = issue_1029_create_entity_shaped_sequence(&backend).await;
1306        assert!(
1307            result.is_ok(),
1308            "khive#1029 repro (KHIVE_WRITE_QUEUE=1): fts_upsert step failed: {:?}",
1309            result.err()
1310        );
1311    }
1312
1313    /// khive#1029 H2 direct test: TWO independent `ConnectionPool`s (hence
1314    /// two independent writer connections / two independent `WriterTask`
1315    /// `OnceLock`s) opened against the SAME tenant DB file — the shape a
1316    /// per-store (rather than per-backend) pool construction would produce.
1317    /// Entity writes go through pool A, the FTS write through pool B, each
1318    /// with `write_queue_enabled: true` so each independently spawns its own
1319    /// WriterTask on first access.
1320    #[tokio::test]
1321    async fn issue_1029_two_pools_same_file_write_queue_on() {
1322        let dir = tempfile::tempdir().unwrap();
1323        let path = dir.path().join("issue_1029_two_pools.db");
1324
1325        let cfg = |p: std::path::PathBuf| crate::pool::PoolConfig {
1326            path: Some(p),
1327            busy_timeout: std::time::Duration::from_millis(200),
1328            write_queue_enabled: true,
1329            ..crate::pool::PoolConfig::default()
1330        };
1331
1332        let pool_a = ConnectionPool::new(cfg(path.clone())).expect("pool A should open");
1333        let backend_a = StorageBackend {
1334            pool: Arc::new(pool_a),
1335            is_file_backed: true,
1336            path: Some(path.clone()),
1337            notes_seq_repair_runs: AtomicUsize::new(0),
1338        };
1339        let pool_b = ConnectionPool::new(cfg(path.clone())).expect("pool B should open");
1340        let backend_b = StorageBackend {
1341            pool: Arc::new(pool_b),
1342            is_file_backed: true,
1343            path: Some(path),
1344            notes_seq_repair_runs: AtomicUsize::new(0),
1345        };
1346
1347        let entities = backend_a
1348            .entities_for_namespace("tenant_ns")
1349            .expect("entities_for_namespace on pool A");
1350        let entity =
1351            khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029TwoPools");
1352        let entity_id = entity.id;
1353        entities
1354            .upsert_entity(entity)
1355            .await
1356            .expect("pool A entity upsert should succeed");
1357
1358        let text = backend_b.text("entities").expect("text on pool B");
1359        let doc = khive_storage::types::TextDocument {
1360            subject_id: entity_id,
1361            kind: khive_types::SubstrateKind::Entity,
1362            title: Some("Issue1029TwoPools".to_string()),
1363            body: "issue 1029 two-pool repro body".to_string(),
1364            tags: vec![],
1365            namespace: "tenant_ns".to_string(),
1366            metadata: None,
1367            updated_at: chrono::Utc::now(),
1368        };
1369        let result = text.upsert_document(doc).await;
1370        assert!(
1371            result.is_ok(),
1372            "khive#1029 two-pool repro: fts_upsert on an independent pool for the \
1373             same tenant DB file failed: {:?}",
1374            result.err()
1375        );
1376    }
1377
1378    /// Minimal thread-local capture subscriber for asserting emitted events —
1379    /// mirrors the capture subscriber in `checkpoint.rs`'s tick tests.
1380    struct StarvationCaptureSubscriber {
1381        events: Arc<std::sync::Mutex<Vec<std::collections::BTreeMap<String, String>>>>,
1382    }
1383
1384    impl tracing::Subscriber for StarvationCaptureSubscriber {
1385        fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
1386            true
1387        }
1388        fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
1389            tracing::span::Id::from_u64(1)
1390        }
1391        fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
1392        fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
1393        fn event(&self, event: &tracing::Event<'_>) {
1394            #[derive(Default)]
1395            struct FieldVisitor(std::collections::BTreeMap<String, String>);
1396            impl tracing::field::Visit for FieldVisitor {
1397                fn record_debug(
1398                    &mut self,
1399                    field: &tracing::field::Field,
1400                    value: &dyn std::fmt::Debug,
1401                ) {
1402                    self.0
1403                        .insert(field.name().to_string(), format!("{value:?}"));
1404                }
1405            }
1406            let mut visitor = FieldVisitor::default();
1407            event.record(&mut visitor);
1408            self.events.lock().unwrap().push(visitor.0);
1409        }
1410        fn enter(&self, _: &tracing::span::Id) {}
1411        fn exit(&self, _: &tracing::span::Id) {}
1412    }
1413
1414    /// Regression coverage for the lock-starvation diagnostic itself: when a
1415    /// text write starves on the SQLite write lock, `with_writer_unmanaged`
1416    /// must emit the WARN carrying the `tx_registry` snapshot — operation
1417    /// name, open-transaction count, and the registered labels.
1418    ///
1419    /// `#[serial(tx_registry)]`: the registry is a process-wide singleton
1420    /// shared across this test binary; this group serializes every test that
1421    /// registers fixture entries or asserts snapshot contents (see
1422    /// `checkpoint.rs`, `pool.rs`, `sql_bridge.rs`). The assertion checks the
1423    /// fixture label is PRESENT rather than the snapshot being exactly one
1424    /// entry, so unrelated short-lived production registrations elsewhere in
1425    /// the binary cannot flake it.
1426    #[tokio::test]
1427    #[serial_test::serial(tx_registry)]
1428    async fn issue_1029_starvation_warn_reports_registered_transactions() {
1429        let (_dir, backend) = issue_1029_pool(false);
1430        // Create the store (and its FTS DDL) BEFORE the lock is held, so the
1431        // starvation happens inside `upsert_document` itself.
1432        let text = backend.text("entities").expect("text store");
1433
1434        // Hold a genuine SQLite write lock on a separate standalone writer
1435        // connection, with a registered fixture transaction the diagnostic
1436        // must surface.
1437        let holder = backend
1438            .pool
1439            .open_standalone_writer()
1440            .expect("holder connection");
1441        holder
1442            .execute_batch("BEGIN IMMEDIATE")
1443            .expect("holder BEGIN IMMEDIATE");
1444        let fixture =
1445            khive_storage::tx_registry::register(Some("issue_1029_fixture_tx".to_string()));
1446
1447        let events = Arc::new(std::sync::Mutex::new(Vec::new()));
1448        let subscriber = StarvationCaptureSubscriber {
1449            events: Arc::clone(&events),
1450        };
1451        let guard = tracing::subscriber::set_default(subscriber);
1452
1453        let doc = khive_storage::types::TextDocument {
1454            subject_id: uuid::Uuid::new_v4(),
1455            kind: khive_types::SubstrateKind::Entity,
1456            title: Some("Issue1029Starved".to_string()),
1457            body: "issue 1029 starvation diagnostic body".to_string(),
1458            tags: vec![],
1459            namespace: "tenant_ns".to_string(),
1460            metadata: None,
1461            updated_at: chrono::Utc::now(),
1462        };
1463        let result = text.upsert_document(doc).await;
1464
1465        drop(guard);
1466        drop(fixture);
1467        holder
1468            .execute_batch("ROLLBACK")
1469            .expect("holder ROLLBACK releases the lock");
1470
1471        assert!(
1472            result.is_err(),
1473            "upsert_document must starve while another connection holds the write lock"
1474        );
1475
1476        let events = events.lock().unwrap();
1477        let warn = events
1478            .iter()
1479            .find(|fields| {
1480                fields
1481                    .get("message")
1482                    .is_some_and(|m| m.contains("text write starved"))
1483            })
1484            .unwrap_or_else(|| panic!("expected a starvation WARN, captured events: {events:?}"));
1485        assert!(
1486            warn.get("op").is_some_and(|op| op.contains("fts_upsert")),
1487            "WARN must name the starved operation, got: {warn:?}"
1488        );
1489        assert!(
1490            warn.get("open_txs")
1491                .is_some_and(|txs| txs.contains("issue_1029_fixture_tx")),
1492            "WARN must list the registered holder label, got: {warn:?}"
1493        );
1494        let count: usize = warn
1495            .get("open_tx_count")
1496            .expect("WARN must carry open_tx_count")
1497            .parse()
1498            .expect("open_tx_count must be numeric");
1499        assert!(
1500            count >= 1,
1501            "open_tx_count must count the fixture, got {count}"
1502        );
1503    }
1504}