use std::path::Path;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use rusqlite::OptionalExtension;
use crate::error::SqliteError;
use crate::pool::{ConnectionPool, PoolConfig};
use crate::sql_bridge::SqlBridge;
use crate::stores::{agents, attachment, blob, entity, event, graph, note, sparse, text, vectors};
fn sqlite_table_exists(conn: &rusqlite::Connection, table: &str) -> Result<bool, SqliteError> {
conn.query_row(
"SELECT 1 FROM sqlite_master WHERE type='table' AND name=?1",
rusqlite::params![table],
|row| row.get::<_, i64>(0),
)
.optional()
.map(|row| row.is_some())
.map_err(SqliteError::Rusqlite)
}
fn validate_vector_table_columns(
conn: &rusqlite::Connection,
table: &str,
) -> Result<(), SqliteError> {
let pragma = format!("PRAGMA table_xinfo({table})");
let mut stmt = conn.prepare(&pragma)?;
let mut rows = stmt.query([])?;
let mut has_field = false;
let mut has_embedding_model = false;
while let Some(row) = rows.next()? {
let name: String = row.get(1)?;
if name == "field" {
has_field = true;
}
if name == "embedding_model" {
has_embedding_model = true;
}
}
if !has_field || !has_embedding_model {
return Err(SqliteError::InvalidData(format!(
"vec0 table '{table}' is missing required column(s) (field={has_field}, \
embedding_model={has_embedding_model}); this is a pre-v0.2.8 vector schema and is \
not supported — recreate the database"
)));
}
Ok(())
}
pub struct StorageBackend {
pool: Arc<ConnectionPool>,
is_file_backed: bool,
path: Option<std::path::PathBuf>,
notes_seq_repair_runs: AtomicUsize,
}
impl StorageBackend {
pub fn sqlite(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
crate::extension::ensure_extensions_loaded();
let resolved = path.as_ref().to_path_buf();
let read_only =
std::fs::metadata(&resolved).is_ok_and(|metadata| metadata.permissions().readonly());
let mut config = PoolConfig {
path: Some(resolved.clone()),
read_only,
..PoolConfig::default()
};
if read_only {
config.write_queue_enabled = Some(false);
}
let pool = ConnectionPool::new(config)?;
Ok(Self {
pool: Arc::new(pool),
is_file_backed: true,
path: Some(resolved),
notes_seq_repair_runs: AtomicUsize::new(0),
})
}
pub fn sqlite_read_only(path: impl AsRef<Path>) -> Result<Self, SqliteError> {
crate::extension::ensure_extensions_loaded();
let resolved = path.as_ref().to_path_buf();
let config = PoolConfig {
path: Some(resolved.clone()),
read_only: true,
write_queue_enabled: Some(false),
..PoolConfig::default()
};
let pool = ConnectionPool::new(config)?;
Ok(Self {
pool: Arc::new(pool),
is_file_backed: true,
path: Some(resolved),
notes_seq_repair_runs: AtomicUsize::new(0),
})
}
pub fn memory() -> Result<Self, SqliteError> {
crate::extension::ensure_extensions_loaded();
let config = PoolConfig {
path: None,
..PoolConfig::default()
};
let pool = ConnectionPool::new(config)?;
Ok(Self {
pool: Arc::new(pool),
is_file_backed: false,
path: None,
notes_seq_repair_runs: AtomicUsize::new(0),
})
}
pub fn sql(&self) -> Arc<dyn khive_storage::SqlAccess> {
Arc::new(SqlBridge::new(Arc::clone(&self.pool), self.is_file_backed))
}
pub fn apply_schema(
&self,
plan: &crate::migrations::ServiceSchemaPlan,
) -> Result<(), SqliteError> {
let writer = self.pool.try_writer()?;
crate::migrations::apply_schema_plan(writer.conn(), plan)
}
pub fn apply_pack_ddl_statements(
&self,
statements: &[&'static str],
) -> Result<(), SqliteError> {
let writer = self.pool.try_writer()?;
writer.transaction(|conn| {
for &stmt in statements {
conn.execute_batch(stmt)?;
}
Ok(())
})
}
pub fn prepare_core_schema(&self) -> Result<u32, SqliteError> {
if self.is_read_only() {
let reader = self.pool.reader()?;
crate::migrations::validate_schema_is_current(reader.conn())
} else {
let latest = crate::migrations::MIGRATIONS
.last()
.map(|migration| migration.version)
.unwrap_or(0);
{
let reader = self.pool.reader()?;
let current = crate::migrations::read_schema_version(reader.conn())?;
if current >= latest {
return crate::migrations::validate_schema_is_current(reader.conn());
}
}
let owner = crate::stores::blob::acquire_database_gc_owner_for_path_blocking(
self.pool.canonical_path().map(Path::to_path_buf),
)
.map_err(|error| {
SqliteError::InvalidData(format!(
"failed to acquire database GC owner before schema preparation: {error}"
))
})?;
let mut writer = self.pool.try_writer()?;
crate::migrations::run_migrations_with_database_gc_owner(writer.conn_mut(), &owner)
}
}
pub fn attachment_cutover_status(
&self,
) -> Result<crate::migrations::AttachmentCutoverStatus, SqliteError> {
if self.is_read_only() {
let reader = self.pool.reader()?;
crate::migrations::attachment_cutover_status(reader.conn())
} else {
let writer = self.pool.try_writer()?;
crate::migrations::attachment_cutover_status(writer.conn())
}
}
fn require_attachment_cutover_owner(
&self,
owner: &crate::stores::blob::DatabaseGcOwnerGuard,
) -> Result<(), SqliteError> {
let sql = self.sql();
let backend_path = sql.database_path();
if owner.database_path() != backend_path.as_deref() {
return Err(SqliteError::InvalidData(format!(
"attachment cutover GC owner targets {:?}, but this backend is {:?}",
owner.database_path(),
backend_path.as_deref()
)));
}
Ok(())
}
pub fn stage_attachment_cutover(
&self,
owner: &crate::stores::blob::DatabaseGcOwnerGuard,
) -> Result<(), SqliteError> {
self.require_attachment_cutover_owner(owner)?;
if self.is_read_only() {
return Err(SqliteError::InvalidData(
"cannot stage attachment cutover on a read-only backend".into(),
));
}
let mut writer = self.pool.try_writer()?;
crate::migrations::stage_attachment_cutover(writer.conn_mut())
}
pub fn apply_verified_attachments(
&self,
owner: &crate::stores::blob::DatabaseGcOwnerGuard,
attachments: &[khive_storage::Attachment],
) -> Result<(), SqliteError> {
self.require_attachment_cutover_owner(owner)?;
if self.is_read_only() {
return Err(SqliteError::InvalidData(
"cannot apply verified attachments on a read-only backend".into(),
));
}
let mut writer = self.pool.try_writer()?;
let tx = writer
.conn_mut()
.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
for attachment in attachments {
attachment
.validate()
.map_err(|error| SqliteError::InvalidData(error.to_string()))?;
crate::migrations::apply_generic_verified_attachment(
&tx,
&attachment.record_uuid.to_string(),
attachment.substrate.as_str(),
&attachment.role,
&attachment.content_ref,
attachment.media_type.as_deref(),
attachment.size_bytes,
attachment.created_at,
)?;
}
tx.commit()?;
Ok(())
}
pub fn finalize_attachment_cutover(
&self,
owner: &crate::stores::blob::DatabaseGcOwnerGuard,
) -> Result<(), SqliteError> {
self.require_attachment_cutover_owner(owner)?;
if self.is_read_only() {
return Err(SqliteError::InvalidData(
"cannot finalize attachment cutover on a read-only backend".into(),
));
}
let mut writer = self.pool.try_writer()?;
crate::migrations::finalize_attachment_cutover(writer.conn_mut())
}
pub fn entities(&self) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
self.entities_for_namespace("local")
}
pub fn entities_for_namespace(
&self,
namespace: &str,
) -> Result<Arc<dyn khive_storage::EntityStore>, SqliteError> {
if namespace.trim().is_empty() {
return Err(SqliteError::InvalidData(
"entities namespace must be non-empty".to_string(),
));
}
if !self.is_read_only() {
let writer = self.pool.try_writer()?;
entity::ensure_entities_schema(writer.conn())?;
}
Ok(Arc::new(entity::SqlEntityStore::new(
Arc::clone(&self.pool),
self.is_file_backed,
)))
}
pub fn attachments(&self) -> Result<Arc<dyn khive_storage::AttachmentStore>, SqliteError> {
Ok(Arc::new(attachment::SqlAttachmentStore::new(
Arc::clone(&self.pool),
self.is_file_backed,
)))
}
pub fn graph(&self) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
self.graph_for_namespace("local")
}
pub fn graph_for_namespace(
&self,
namespace: &str,
) -> Result<Arc<dyn khive_storage::GraphStore>, SqliteError> {
if namespace.trim().is_empty() {
return Err(SqliteError::InvalidData(
"graph namespace must be non-empty".to_string(),
));
}
if !self.is_read_only() {
let writer = self.pool.try_writer()?;
graph::ensure_graph_schema(writer.conn())?;
}
Ok(Arc::new(graph::SqlGraphStore::new_scoped(
Arc::clone(&self.pool),
self.is_file_backed,
namespace.trim().to_string(),
)))
}
pub fn notes(&self) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
self.notes_for_namespace("local")
}
pub fn notes_for_namespace(
&self,
namespace: &str,
) -> Result<Arc<dyn khive_storage::NoteStore>, SqliteError> {
if namespace.trim().is_empty() {
return Err(SqliteError::InvalidData(
"notes namespace must be non-empty".to_string(),
));
}
if !self.is_read_only() {
let writer = self.pool.try_writer()?;
note::ensure_notes_schema(writer.conn())?;
if self.notes_seq_repair_runs.load(Ordering::Relaxed) == 0 {
note::repair_notes_seq(writer.conn())?;
self.notes_seq_repair_runs.fetch_add(1, Ordering::Relaxed);
}
}
Ok(Arc::new(note::SqlNoteStore::new(
Arc::clone(&self.pool),
self.is_file_backed,
)))
}
pub fn notes_seq_repair_run_count(&self) -> usize {
self.notes_seq_repair_runs.load(Ordering::Relaxed)
}
pub fn events(&self) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
self.events_for_namespace("local")
}
pub fn events_for_namespace(
&self,
namespace: &str,
) -> Result<Arc<dyn khive_storage::EventStore>, SqliteError> {
if namespace.trim().is_empty() {
return Err(SqliteError::InvalidData(
"events namespace must be non-empty".to_string(),
));
}
if !self.is_read_only() {
let writer = self.pool.try_writer()?;
event::ensure_events_schema(writer.conn())?;
}
Ok(Arc::new(event::SqlEventStore::new_scoped(
Arc::clone(&self.pool),
self.is_file_backed,
namespace.trim().to_string(),
)))
}
pub fn agents(&self) -> Result<Arc<dyn khive_storage::AgentStore>, SqliteError> {
if !self.is_read_only() {
let writer = self.pool.try_writer()?;
agents::ensure_agents_schema(writer.conn())?;
}
Ok(Arc::new(agents::SqlAgentStore::new(
Arc::clone(&self.pool),
self.is_file_backed,
)))
}
pub fn vectors(
&self,
model_key: &str,
embedding_model: &str,
dimensions: usize,
) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
self.vectors_for_namespace(model_key, embedding_model, dimensions, "local")
}
pub fn vectors_for_namespace(
&self,
model_key: &str,
embedding_model: &str,
dimensions: usize,
namespace: &str,
) -> Result<Arc<dyn khive_storage::VectorStore>, SqliteError> {
if model_key.is_empty()
|| !model_key
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_')
{
return Err(SqliteError::InvalidData(format!(
"invalid model_key '{}': must be non-empty and contain only \
alphanumeric/underscore characters",
model_key
)));
}
if namespace.trim().is_empty() {
return Err(SqliteError::InvalidData(
"vector store namespace must be non-empty".to_string(),
));
}
crate::extension::ensure_extensions_loaded();
let table = format!("vec_{}", model_key);
if self.is_read_only() {
let reader = self.pool.reader()?;
if !sqlite_table_exists(reader.conn(), &table)? {
return Err(SqliteError::InvalidData(format!(
"read-only database has no vector table '{table}'; create and populate it in \
a writable copy before opening the snapshot"
)));
}
validate_vector_table_columns(reader.conn(), &table)?;
drop(reader);
return Ok(Arc::new(vectors::SqliteVecStore::new(
Arc::clone(&self.pool),
self.is_file_backed,
model_key.to_string(),
embedding_model.to_string(),
dimensions,
namespace.trim().to_string(),
)?));
}
let writer = self.pool.try_writer()?;
let table_exists = sqlite_table_exists(writer.conn(), &table)?;
if table_exists {
validate_vector_table_columns(writer.conn(), &table)?;
}
writer
.conn()
.execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
writer
.conn()
.execute_batch(crate::migrations::ANN_WRITE_LOG_DDL)?;
writer
.conn()
.execute_batch(crate::migrations::ANN_WRITE_LOG_MODEL_SEQ_INDEX_DDL)?;
writer
.conn()
.execute_batch(crate::migrations::ANN_CONSUMER_PENDING_DDL)?;
let ddl = format!(
"CREATE VIRTUAL TABLE IF NOT EXISTS vec_{} USING vec0(\
subject_id TEXT PRIMARY KEY, \
namespace TEXT NOT NULL, \
kind TEXT NOT NULL, \
field TEXT NOT NULL, \
embedding_model TEXT NOT NULL, \
embedding float[{}] distance_metric=cosine\
)",
model_key, dimensions
);
writer.conn().execute_batch(&ddl)?;
Ok(Arc::new(vectors::SqliteVecStore::new(
Arc::clone(&self.pool),
self.is_file_backed,
model_key.to_string(),
embedding_model.to_string(),
dimensions,
namespace.trim().to_string(),
)?))
}
pub fn register_embedding_model(
&self,
engine_name: &str,
model_id: &str,
key_version: &str,
dimensions: u32,
) -> Result<(), SqliteError> {
let writer = self.pool.try_writer()?;
writer
.conn()
.execute_batch(crate::migrations::EMBEDDING_MODELS_DDL)?;
let now = chrono::Utc::now().timestamp_micros();
let canonical_key =
format!("{engine_name}:{model_id}:{key_version}:{dimensions}").into_bytes();
let id = uuid::Uuid::new_v4();
writer.conn().execute(
"INSERT INTO _embedding_models \
(id, engine_name, model_id, key_version, dim, output_dim, status, \
activated_at, superseded_at, superseded_by, canonical_key, created_at) \
VALUES (?1, ?2, ?3, ?4, ?5, NULL, 'active', ?6, NULL, NULL, ?7, ?8) \
ON CONFLICT(canonical_key) DO UPDATE SET \
status = 'active', \
activated_at = COALESCE(_embedding_models.activated_at, excluded.activated_at)",
rusqlite::params![
id.as_bytes().as_slice(),
engine_name,
model_id,
key_version,
dimensions as i64,
now,
canonical_key,
now,
],
)?;
Ok(())
}
pub fn sparse(
&self,
model_key: &str,
) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
self.sparse_for_namespace(model_key, "local")
}
pub fn sparse_for_namespace(
&self,
model_key: &str,
namespace: &str,
) -> Result<Arc<dyn khive_storage::SparseStore>, SqliteError> {
if model_key.is_empty()
|| !model_key
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_')
{
return Err(SqliteError::InvalidData(format!(
"invalid model_key '{}': must be non-empty and contain only alphanumeric/underscore characters",
model_key
)));
}
if namespace.trim().is_empty() {
return Err(SqliteError::InvalidData(
"sparse store namespace must be non-empty".to_string(),
));
}
if self.is_read_only() {
let table = format!("sparse_{model_key}");
let reader = self.pool.reader()?;
if !sqlite_table_exists(reader.conn(), &table)? {
return Err(SqliteError::InvalidData(format!(
"read-only database has no sparse table '{table}'; create and populate it in \
a writable copy before opening the snapshot"
)));
}
} else {
let writer = self.pool.try_writer()?;
sparse::ensure_sparse_schema(writer.conn(), model_key)
.map_err(SqliteError::Rusqlite)?;
}
Ok(Arc::new(sparse::SqliteSparseStore::new(
Arc::clone(&self.pool),
self.is_file_backed,
model_key.to_string(),
namespace.trim().to_string(),
)?))
}
pub fn text(&self, table_key: &str) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
self.text_with_tokenizer(table_key, "trigram")
}
pub fn text_with_tokenizer(
&self,
table_key: &str,
tokenizer: &str,
) -> Result<Arc<dyn khive_storage::TextSearch>, SqliteError> {
if table_key.is_empty()
|| !table_key
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_')
{
return Err(SqliteError::InvalidData(format!(
"invalid table_key '{}': must be non-empty and contain only \
alphanumeric/underscore characters",
table_key
)));
}
if tokenizer.is_empty()
|| !tokenizer
.chars()
.all(|c| c.is_ascii_alphanumeric() || c == '_')
{
return Err(SqliteError::InvalidData(format!(
"invalid tokenizer '{}': must be non-empty and contain only \
alphanumeric/underscore characters",
tokenizer
)));
}
let ddl = format!(
"CREATE VIRTUAL TABLE IF NOT EXISTS fts_{} USING fts5(\
subject_id UNINDEXED, \
kind UNINDEXED, \
title, \
body, \
tags UNINDEXED, \
namespace UNINDEXED, \
metadata UNINDEXED, \
updated_at UNINDEXED, \
tokenize = '{}'\
)",
table_key, tokenizer
);
if self.is_read_only() {
let table = format!("fts_{table_key}");
let reader = self.pool.reader()?;
if !sqlite_table_exists(reader.conn(), &table)? {
return Err(SqliteError::InvalidData(format!(
"read-only database has no text-search table '{table}'; create and populate \
it in a writable copy before opening the snapshot"
)));
}
} else {
let writer = self.pool.try_writer()?;
writer.conn().execute_batch(&ddl)?;
}
Ok(Arc::new(text::Fts5TextSearch::new(
Arc::clone(&self.pool),
self.is_file_backed,
table_key.to_string(),
)))
}
pub fn blob_store(
&self,
config_root: Option<&Path>,
floor_bytes: Option<u64>,
) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
Ok(Arc::new(blob::FsBlobStore::new(root, floor)?))
}
pub fn blob_store_read_only(
&self,
config_root: Option<&Path>,
floor_bytes: Option<u64>,
) -> Result<Arc<dyn khive_storage::BlobStore>, SqliteError> {
let root = blob::resolve_blob_root(self.data_dir().as_deref(), config_root)?;
let floor = floor_bytes.unwrap_or(blob::FsBlobStore::DEFAULT_FLOOR_BYTES);
Ok(Arc::new(blob::FsBlobStore::open_existing(root, floor)?))
}
pub fn is_file_backed(&self) -> bool {
self.is_file_backed
}
pub fn is_read_only(&self) -> bool {
self.pool.config().read_only
}
pub fn data_dir(&self) -> Option<std::path::PathBuf> {
self.path.as_ref()?.parent().map(|p| p.to_path_buf())
}
pub fn ann_root(&self) -> Option<std::path::PathBuf> {
ann_root_for(self.path.as_ref()?)
}
pub fn pool(&self) -> &ConnectionPool {
&self.pool
}
pub fn pool_arc(&self) -> Arc<ConnectionPool> {
Arc::clone(&self.pool)
}
}
fn ann_root_for(path: &std::path::Path) -> Option<std::path::PathBuf> {
let mut file = path.file_name()?.to_os_string();
file.push(".ann");
path.parent().map(|p| p.join(file))
}
#[cfg(test)]
mod tests {
use super::*;
use khive_storage::types::{SqlStatement, SqlValue};
#[cfg(unix)]
fn freeze_snapshot_sidecars(path: &std::path::Path) {
use std::os::unix::fs::PermissionsExt;
for suffix in ["-wal", "-shm"] {
let mut name = path.file_name().expect("db file name").to_os_string();
name.push(suffix);
let sidecar = path.parent().expect("db parent dir").join(name);
if sidecar.exists() {
let mut permissions = std::fs::metadata(&sidecar)
.expect("sidecar metadata")
.permissions();
permissions.set_mode(0o444);
std::fs::set_permissions(&sidecar, permissions).expect("freeze sidecar");
}
}
}
#[cfg(unix)]
#[tokio::test]
async fn sqlite_detects_chmod_read_only_snapshot_and_core_reads_succeed() {
use std::os::unix::fs::PermissionsExt;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("chmod_snapshot.db");
{
let writable = StorageBackend::sqlite(&path).expect("create writable database");
writable
.prepare_core_schema()
.expect("migrate writable snapshot source");
}
let mut permissions = std::fs::metadata(&path).unwrap().permissions();
permissions.set_mode(0o444);
std::fs::set_permissions(&path, permissions).unwrap();
freeze_snapshot_sidecars(&path);
let read_only = StorageBackend::sqlite(&path).expect("auto-detect read-only mode");
assert!(read_only.is_read_only());
assert_eq!(
read_only.pool().config().write_queue_enabled,
Some(false),
"read-only boot must not attempt to spawn a writer task"
);
assert!(read_only
.pool()
.writer_task_handle()
.expect("disabled writer task is a valid configuration")
.is_none());
read_only
.prepare_core_schema()
.expect("current snapshot validates without migration writes");
let entities = read_only.entities().expect("entity store opens read-only");
let graph = read_only.graph().expect("graph store opens read-only");
let notes = read_only.notes().expect("note store opens read-only");
let events = read_only.events().expect("event store opens read-only");
assert_eq!(
entities
.count_entities("local", khive_storage::EntityFilter::default())
.await
.unwrap(),
0
);
assert_eq!(
graph
.count_edges(khive_storage::types::EdgeFilter::default())
.await
.unwrap(),
0
);
assert_eq!(notes.count_notes("local", None).await.unwrap(), 0);
assert_eq!(
events
.count_events(khive_storage::EventFilter::default())
.await
.unwrap(),
0
);
assert_eq!(
read_only.notes_seq_repair_run_count(),
0,
"read-only store acquisition must not run the DML repair"
);
}
#[test]
fn memory_backend_creates_successfully() {
let backend = StorageBackend::memory().expect("memory backend should create");
assert!(!backend.is_file_backed());
}
#[test]
fn file_backend_creates_successfully() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("test.db");
let backend = StorageBackend::sqlite(&path).expect("file backend should create");
assert!(backend.is_file_backed());
assert!(path.exists());
}
#[test]
fn data_dir_returns_none_for_memory_backend() {
let backend = StorageBackend::memory().expect("memory backend");
assert!(backend.data_dir().is_none());
}
#[test]
fn data_dir_returns_parent_dir_for_file_backend() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("data.db");
let backend = StorageBackend::sqlite(&path).expect("file backend");
let got = backend.data_dir().expect("file backend must return Some");
assert_eq!(got, dir.path());
}
#[test]
fn ann_root_is_database_scoped_sibling_dir() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("data.db");
let backend = StorageBackend::sqlite(&path).expect("file backend");
let got = backend.ann_root().expect("file backend must return Some");
assert_eq!(got, dir.path().join("data.db.ann"));
assert!(StorageBackend::memory().unwrap().ann_root().is_none());
}
#[cfg(unix)]
#[test]
fn ann_root_distinct_for_non_utf8_filenames() {
use std::os::unix::ffi::OsStrExt;
let path_a = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xff.db"));
let path_b = std::path::Path::new("/data").join(std::ffi::OsStr::from_bytes(b"\xfe.db"));
let root_a = ann_root_for(&path_a).expect("Some for a file path");
let root_b = ann_root_for(&path_b).expect("Some for a file path");
assert_ne!(
root_a, root_b,
"distinct database files must map to distinct ANN roots"
);
}
#[tokio::test]
async fn sql_access_memory_roundtrip() {
let backend = StorageBackend::memory().unwrap();
let sql = backend.sql();
let mut writer = sql.writer().await.unwrap();
writer
.execute_script(
"CREATE TABLE test_rt (id TEXT PRIMARY KEY, value INTEGER NOT NULL)".into(),
)
.await
.unwrap();
let affected = writer
.execute(SqlStatement {
sql: "INSERT INTO test_rt (id, value) VALUES (?1, ?2)".into(),
params: vec![SqlValue::Text("row1".into()), SqlValue::Integer(42)],
label: None,
})
.await
.unwrap();
assert_eq!(affected, 1);
let mut reader = sql.reader().await.unwrap();
let row = reader
.query_row(SqlStatement {
sql: "SELECT id, value FROM test_rt WHERE id = ?1".into(),
params: vec![SqlValue::Text("row1".into())],
label: None,
})
.await
.unwrap();
let row = row.expect("should find the inserted row");
assert_eq!(row.columns.len(), 2);
match &row.columns[0].value {
SqlValue::Text(s) => assert_eq!(s, "row1"),
other => panic!("expected Text, got {other:?}"),
}
match &row.columns[1].value {
SqlValue::Integer(v) => assert_eq!(*v, 42),
other => panic!("expected Integer, got {other:?}"),
}
}
#[tokio::test]
async fn sql_access_file_roundtrip() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("test_roundtrip.db");
let backend = StorageBackend::sqlite(&path).unwrap();
let sql = backend.sql();
let mut writer = sql.writer().await.unwrap();
writer
.execute_script("CREATE TABLE test_f (k TEXT PRIMARY KEY, v TEXT)".into())
.await
.unwrap();
writer
.execute(SqlStatement {
sql: "INSERT INTO test_f (k, v) VALUES (?1, ?2)".into(),
params: vec![
SqlValue::Text("hello".into()),
SqlValue::Text("world".into()),
],
label: None,
})
.await
.unwrap();
let mut reader = sql.reader().await.unwrap();
let rows = reader
.query_all(SqlStatement {
sql: "SELECT k, v FROM test_f".into(),
params: vec![],
label: None,
})
.await
.unwrap();
assert_eq!(rows.len(), 1);
match &rows[0].columns[1].value {
SqlValue::Text(s) => assert_eq!(s, "world"),
other => panic!("expected Text, got {other:?}"),
}
}
#[test]
fn sqlite_read_only_missing_path_does_not_create_file() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("missing_ro.db");
assert!(!path.exists());
let result = StorageBackend::sqlite_read_only(&path);
assert!(
result.is_err(),
"opening a missing path read-only must fail"
);
assert!(
!path.exists(),
"opening a missing path read-only must not create the file"
);
}
#[test]
fn sqlite_read_only_sparse_store_requires_existing_table_without_writer_acquisition() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("ro_sparse_tables.db");
{
let writable = StorageBackend::sqlite(&path).unwrap();
writable
.prepare_core_schema()
.expect("migrate snapshot source");
writable
.sparse("present")
.expect("create the optional sparse table while writable");
}
#[cfg(unix)]
freeze_snapshot_sidecars(&path);
let read_only = StorageBackend::sqlite_read_only(&path).unwrap();
read_only
.prepare_core_schema()
.expect("validate exact current migration ledger");
read_only
.sparse("present")
.expect("an existing sparse table must open read-only");
let missing = match read_only.sparse("missing") {
Ok(_) => panic!("a missing sparse table must fail during store acquisition"),
Err(error) => error,
};
assert!(
missing.to_string().contains("sparse_missing"),
"the diagnostic must name the absent table: {missing}"
);
assert_eq!(
read_only.pool().writer_acquisition_snapshot(),
crate::pool::WriterAcquisitionSnapshot::default(),
"construction, exact-ledger validation, and optional sparse-table inspection must \
use reader connections only"
);
}
#[test]
fn sqlite_read_only_text_store_requires_existing_table_without_writer_acquisition() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("ro_text_tables.db");
{
let writable = StorageBackend::sqlite(&path).unwrap();
writable
.prepare_core_schema()
.expect("migrate snapshot source");
writable
.text("present")
.expect("create the optional FTS table while writable");
}
#[cfg(unix)]
freeze_snapshot_sidecars(&path);
let read_only = StorageBackend::sqlite_read_only(&path).unwrap();
read_only
.prepare_core_schema()
.expect("validate exact current migration ledger");
read_only
.text("present")
.expect("an existing FTS table must open read-only");
let missing = match read_only.text("missing") {
Ok(_) => panic!("a missing FTS table must fail during store acquisition"),
Err(error) => error,
};
assert!(
missing.to_string().contains("fts_missing"),
"the diagnostic must name the absent table: {missing}"
);
assert_eq!(
read_only.pool().writer_acquisition_snapshot(),
crate::pool::WriterAcquisitionSnapshot::default(),
"construction, exact-ledger validation, and optional FTS inspection must use reader \
connections only"
);
}
#[cfg(feature = "vectors")]
#[test]
fn sqlite_read_only_vector_store_schema_check_uses_no_writer_acquisition() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("ro_vector_tables.db");
{
let writable = StorageBackend::sqlite(&path).unwrap();
writable
.prepare_core_schema()
.expect("migrate snapshot source");
writable
.vectors("present", "present", 3)
.expect("create the optional vector table while writable");
}
#[cfg(unix)]
freeze_snapshot_sidecars(&path);
let read_only = StorageBackend::sqlite_read_only(&path).unwrap();
read_only
.prepare_core_schema()
.expect("validate exact current migration ledger");
read_only
.vectors("present", "present", 3)
.expect("an existing vector table must open read-only");
assert!(
read_only.vectors("missing", "missing", 3).is_err(),
"a missing vector table must fail during store acquisition"
);
assert_eq!(
read_only.pool().writer_acquisition_snapshot(),
crate::pool::WriterAcquisitionSnapshot::default(),
"construction, exact-ledger validation, and optional vector inspection must use \
reader connections only"
);
}
#[tokio::test]
async fn sqlite_read_only_sql_writer_rejects_ddl_and_insert() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("ro_writer.db");
{
let writable = StorageBackend::sqlite(&path).unwrap();
let sql = writable.sql();
let mut writer = sql.writer().await.unwrap();
writer
.execute_script("CREATE TABLE ro_existing (id INTEGER PRIMARY KEY)".into())
.await
.unwrap();
}
#[cfg(unix)]
freeze_snapshot_sidecars(&path);
let ro = StorageBackend::sqlite_read_only(&path).unwrap();
let sql = ro.sql();
let writer_result = sql.writer().await;
assert!(
writer_result.is_err(),
"sql().writer() must be rejected on a read-only backend"
);
}
#[tokio::test]
#[cfg(feature = "vectors")]
async fn vectors_roundtrip_via_public_api() {
let backend = StorageBackend::memory().unwrap();
let store = backend.vectors("test_api", "test_api", 3).unwrap();
let id = uuid::Uuid::new_v4();
store
.insert(
id,
khive_types::SubstrateKind::Entity,
"local",
"content",
vec![vec![1.0, 0.0, 0.0]],
)
.await
.unwrap();
let hits = store
.search(khive_storage::types::VectorSearchRequest {
query_vectors: vec![vec![1.0, 0.0, 0.0]],
top_k: 1,
namespace: None,
kind: None,
embedding_model: None,
filter: None,
backend_hints: None,
})
.await
.unwrap();
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].subject_id, id);
assert!(hits[0].score.to_f64() > 0.99);
}
#[tokio::test]
#[cfg(feature = "vectors")]
async fn vectors_creates_table_idempotently() {
let backend = StorageBackend::memory().unwrap();
let store1 = backend.vectors("idempotent", "idempotent", 3).unwrap();
let store2 = backend.vectors("idempotent", "idempotent", 3).unwrap();
let id = uuid::Uuid::new_v4();
store1
.insert(
id,
khive_types::SubstrateKind::Entity,
"local",
"content",
vec![vec![1.0, 0.0, 0.0]],
)
.await
.unwrap();
let count = store2.count().await.unwrap();
assert_eq!(count, 1);
}
#[tokio::test]
async fn text_roundtrip_via_public_api() {
let backend = StorageBackend::memory().unwrap();
let store = backend.text("test_api").unwrap();
let id = uuid::Uuid::new_v4();
let doc = khive_storage::types::TextDocument {
subject_id: id,
kind: khive_types::SubstrateKind::Entity,
title: Some("Test Title".to_string()),
body: "This is a searchable document about Rust.".to_string(),
tags: vec!["rust".to_string()],
namespace: "test_ns".to_string(),
metadata: None,
updated_at: chrono::Utc::now(),
};
store.upsert_document(doc).await.unwrap();
let hits = store
.search(khive_storage::types::TextSearchRequest {
query: "Rust".to_string(),
mode: khive_storage::types::TextQueryMode::Plain,
filter: Some(khive_storage::types::TextFilter {
namespaces: vec!["test_ns".to_string()],
..Default::default()
}),
top_k: 1,
snippet_chars: 64,
})
.await
.unwrap();
assert_eq!(hits.len(), 1);
assert_eq!(hits[0].subject_id, id);
assert!(hits[0].score.to_f64() > 0.0);
}
#[tokio::test]
async fn text_creates_table_idempotently() {
let backend = StorageBackend::memory().unwrap();
let store1 = backend.text("idempotent_fts").unwrap();
let store2 = backend.text("idempotent_fts").unwrap();
let id = uuid::Uuid::new_v4();
let doc = khive_storage::types::TextDocument {
subject_id: id,
kind: khive_types::SubstrateKind::Note,
title: None,
body: "Hello world.".to_string(),
tags: vec![],
namespace: "test_ns".to_string(),
metadata: None,
updated_at: chrono::Utc::now(),
};
store1.upsert_document(doc).await.unwrap();
let count = store2
.count(khive_storage::types::TextFilter {
namespaces: vec!["test_ns".to_string()],
..Default::default()
})
.await
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn invalid_model_key_rejected() {
let backend = StorageBackend::memory().unwrap();
assert!(backend.vectors("bad key!", "bad key!", 3).is_err());
assert!(backend.vectors("", "", 3).is_err());
}
#[test]
fn invalid_table_key_rejected() {
let backend = StorageBackend::memory().unwrap();
assert!(backend.text("bad key!").is_err());
assert!(backend.text("").is_err());
}
#[tokio::test]
async fn sqlite_read_only_graph_store_rejects_upsert_edge() {
use khive_storage::types::Edge;
use khive_types::EdgeRelation;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("ro_graph.db");
{
let writable = StorageBackend::sqlite(&path).unwrap();
writable.graph().unwrap();
}
#[cfg(unix)]
freeze_snapshot_sidecars(&path);
let ro = StorageBackend::sqlite_read_only(&path).unwrap();
let store = match ro.graph() {
Ok(store) => store,
Err(_) => return,
};
let now = chrono::Utc::now();
let edge = Edge {
id: uuid::Uuid::new_v4().into(),
namespace: "local".to_string(),
source_id: uuid::Uuid::new_v4(),
target_id: uuid::Uuid::new_v4(),
relation: EdgeRelation::Extends,
weight: 0.8,
created_at: now,
updated_at: now,
deleted_at: None,
metadata: None,
target_backend: None,
};
let result = store.upsert_edge(edge).await;
assert!(
result.is_err(),
"upsert_edge on a read-only backend must reject, not silently no-op"
);
}
#[tokio::test]
async fn sqlite_read_only_event_store_rejects_append_event() {
use khive_types::{EventKind, EventOutcome, SubstrateKind};
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("ro_events.db");
{
let writable = StorageBackend::sqlite(&path).unwrap();
writable.events().unwrap();
}
#[cfg(unix)]
freeze_snapshot_sidecars(&path);
let ro = StorageBackend::sqlite_read_only(&path).unwrap();
let store = match ro.events() {
Ok(store) => store,
Err(_) => return,
};
let event = khive_storage::event::Event::new(
"local",
"test.verb",
EventKind::Audit,
SubstrateKind::Entity,
"test-actor",
)
.with_outcome(EventOutcome::Success);
let result = store.append_event(event).await;
assert!(
result.is_err(),
"append_event on a read-only backend must reject, not silently no-op"
);
}
#[tokio::test]
async fn sqlite_read_only_text_store_rejects_upsert_document() {
use khive_storage::types::TextDocument;
use khive_types::SubstrateKind;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("ro_text.db");
{
let writable = StorageBackend::sqlite(&path).unwrap();
writable.text("ro_test").unwrap();
}
#[cfg(unix)]
freeze_snapshot_sidecars(&path);
let ro = StorageBackend::sqlite_read_only(&path).unwrap();
let store = match ro.text("ro_test") {
Ok(store) => store,
Err(_) => return,
};
let doc = TextDocument {
subject_id: uuid::Uuid::new_v4(),
kind: SubstrateKind::Entity,
title: Some("Title".to_string()),
body: "Body text.".to_string(),
tags: vec![],
namespace: "local".to_string(),
metadata: None,
updated_at: chrono::Utc::now(),
};
let result = store.upsert_document(doc).await;
assert!(
result.is_err(),
"upsert_document on a read-only backend must reject, not silently no-op"
);
}
#[tokio::test]
async fn blob_store_roundtrip_via_public_api() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("blob_backend.db");
let backend = StorageBackend::sqlite(&path).unwrap();
let store = backend.blob_store(None, Some(0)).unwrap();
let bytes = b"backend-level blob roundtrip".to_vec();
let content_ref = store.put(bytes.clone()).await.unwrap();
assert_eq!(
store
.get_bounded_verified(&content_ref, bytes.len() as u64)
.await
.unwrap(),
bytes
);
}
#[test]
fn blob_store_defaults_root_beside_db_file() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("blob_default.db");
let backend = StorageBackend::sqlite(&path).unwrap();
let _store = backend.blob_store(None, None).unwrap();
assert!(
dir.path().join("blobs").is_dir(),
"default root must be created beside the database file"
);
}
#[test]
fn blob_store_errors_for_in_memory_backend_with_no_override() {
let backend = StorageBackend::memory().unwrap();
assert!(backend.blob_store(None, None).is_err());
}
#[test]
fn blob_store_accepts_explicit_root_for_in_memory_backend() {
let dir = tempfile::tempdir().unwrap();
let backend = StorageBackend::memory().unwrap();
let store = backend.blob_store(Some(dir.path()), None);
assert!(store.is_ok());
}
#[test]
fn apply_schema_runs_migrations_idempotently() {
static MIGRATIONS: &[crate::migrations::Migration] = &[crate::migrations::Migration {
id: "001_init",
up_sql: "CREATE TABLE IF NOT EXISTS schema_test (id TEXT PRIMARY KEY);",
down_sql: None,
is_already_applied: None,
}];
let plan = crate::migrations::ServiceSchemaPlan {
service: "schema_test_svc",
sqlite: MIGRATIONS,
postgres: &[],
};
let backend = StorageBackend::memory().unwrap();
backend.apply_schema(&plan).unwrap();
backend.apply_schema(&plan).unwrap();
let reader = backend.pool().reader().unwrap();
let count: i64 = reader
.conn()
.query_row(
"SELECT COUNT(*) FROM sqlite_master WHERE type='table' AND name='schema_test'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn pack_ddl_plan_rolls_back_all_statements_on_failure() {
let backend = StorageBackend::memory().unwrap();
let error = backend
.apply_pack_ddl_statements(&[
"CREATE TABLE IF NOT EXISTS pack_schema_first (id INTEGER PRIMARY KEY)",
"CREATE INDEX IF NOT EXISTS pack_schema_second ON pack_schema_missing(id)",
])
.unwrap_err();
assert!(
error.to_string().contains("pack_schema_missing"),
"schema-plan error must retain the failing SQLite diagnostic: {error}"
);
let reader = backend.pool().reader().unwrap();
let visible_objects: i64 = reader
.conn()
.query_row(
"SELECT COUNT(*) FROM sqlite_master \
WHERE name IN ('pack_schema_first', 'pack_schema_second')",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(visible_objects, 0);
}
#[test]
fn pack_ddl_plan_applies_all_statements_idempotently() {
const PLAN: &[&str] = &[
"CREATE TABLE IF NOT EXISTS pack_schema_success (id INTEGER PRIMARY KEY, value TEXT)",
"CREATE INDEX IF NOT EXISTS pack_schema_success_value_idx \
ON pack_schema_success(value)",
];
let backend = StorageBackend::memory().unwrap();
backend.apply_pack_ddl_statements(PLAN).unwrap();
backend.apply_pack_ddl_statements(PLAN).unwrap();
let reader = backend.pool().reader().unwrap();
let visible_objects: i64 = reader
.conn()
.query_row(
"SELECT COUNT(*) FROM sqlite_master \
WHERE name IN ('pack_schema_success', 'pack_schema_success_value_idx')",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(visible_objects, 2);
}
fn issue_1029_pool(write_queue_enabled: bool) -> (tempfile::TempDir, StorageBackend) {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("issue_1029.db");
let config = crate::pool::PoolConfig {
path: Some(path.clone()),
busy_timeout: std::time::Duration::from_millis(200),
write_queue_enabled: Some(write_queue_enabled),
..crate::pool::PoolConfig::default()
};
let pool = ConnectionPool::new(config).expect("fresh tenant-shaped pool should open");
let backend = StorageBackend {
pool: Arc::new(pool),
is_file_backed: true,
path: Some(path),
notes_seq_repair_runs: AtomicUsize::new(0),
};
(dir, backend)
}
async fn issue_1029_create_entity_shaped_sequence(
backend: &StorageBackend,
) -> Result<(), String> {
let entities = backend
.entities_for_namespace("tenant_ns")
.map_err(|e| format!("entities_for_namespace: {e}"))?;
let entity = khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029Repro");
let entity_id = entity.id;
entities
.upsert_entity(entity)
.await
.map_err(|e| format!("upsert_entity: {e}"))?;
let text = backend.text("entities").map_err(|e| format!("text: {e}"))?;
let doc = khive_storage::types::TextDocument {
subject_id: entity_id,
kind: khive_types::SubstrateKind::Entity,
title: Some("Issue1029Repro".to_string()),
body: "issue 1029 repro body".to_string(),
tags: vec![],
namespace: "tenant_ns".to_string(),
metadata: None,
updated_at: chrono::Utc::now(),
};
text.upsert_document(doc)
.await
.map_err(|e| format!("fts_upsert: {e}"))
}
#[tokio::test]
async fn issue_1029_create_entity_shaped_sequence_write_queue_off() {
let (_dir, backend) = issue_1029_pool(false);
let result = issue_1029_create_entity_shaped_sequence(&backend).await;
assert!(
result.is_ok(),
"khive#1029 repro (KHIVE_WRITE_QUEUE off): fts_upsert step failed: {:?}",
result.err()
);
}
#[tokio::test]
async fn issue_1029_create_entity_shaped_sequence_write_queue_on() {
let (_dir, backend) = issue_1029_pool(true);
let result = issue_1029_create_entity_shaped_sequence(&backend).await;
assert!(
result.is_ok(),
"khive#1029 repro (KHIVE_WRITE_QUEUE=1): fts_upsert step failed: {:?}",
result.err()
);
}
#[tokio::test]
async fn issue_1029_two_pools_same_file_write_queue_on() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("issue_1029_two_pools.db");
let cfg = |p: std::path::PathBuf| crate::pool::PoolConfig {
path: Some(p),
busy_timeout: std::time::Duration::from_millis(200),
write_queue_enabled: Some(true),
..crate::pool::PoolConfig::default()
};
let pool_a = ConnectionPool::new(cfg(path.clone())).expect("pool A should open");
let backend_a = StorageBackend {
pool: Arc::new(pool_a),
is_file_backed: true,
path: Some(path.clone()),
notes_seq_repair_runs: AtomicUsize::new(0),
};
let pool_b = ConnectionPool::new(cfg(path.clone())).expect("pool B should open");
let backend_b = StorageBackend {
pool: Arc::new(pool_b),
is_file_backed: true,
path: Some(path),
notes_seq_repair_runs: AtomicUsize::new(0),
};
let entities = backend_a
.entities_for_namespace("tenant_ns")
.expect("entities_for_namespace on pool A");
let entity =
khive_storage::entity::Entity::new("tenant_ns", "concept", "Issue1029TwoPools");
let entity_id = entity.id;
entities
.upsert_entity(entity)
.await
.expect("pool A entity upsert should succeed");
let text = backend_b.text("entities").expect("text on pool B");
let doc = khive_storage::types::TextDocument {
subject_id: entity_id,
kind: khive_types::SubstrateKind::Entity,
title: Some("Issue1029TwoPools".to_string()),
body: "issue 1029 two-pool repro body".to_string(),
tags: vec![],
namespace: "tenant_ns".to_string(),
metadata: None,
updated_at: chrono::Utc::now(),
};
let result = text.upsert_document(doc).await;
assert!(
result.is_ok(),
"khive#1029 two-pool repro: fts_upsert on an independent pool for the \
same tenant DB file failed: {:?}",
result.err()
);
}
struct StarvationCaptureSubscriber {
events: Arc<std::sync::Mutex<Vec<std::collections::BTreeMap<String, String>>>>,
}
impl tracing::Subscriber for StarvationCaptureSubscriber {
fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
true
}
fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
tracing::span::Id::from_u64(1)
}
fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
fn event(&self, event: &tracing::Event<'_>) {
#[derive(Default)]
struct FieldVisitor(std::collections::BTreeMap<String, String>);
impl tracing::field::Visit for FieldVisitor {
fn record_debug(
&mut self,
field: &tracing::field::Field,
value: &dyn std::fmt::Debug,
) {
self.0
.insert(field.name().to_string(), format!("{value:?}"));
}
}
let mut visitor = FieldVisitor::default();
event.record(&mut visitor);
self.events.lock().unwrap().push(visitor.0);
}
fn enter(&self, _: &tracing::span::Id) {}
fn exit(&self, _: &tracing::span::Id) {}
}
#[tokio::test]
#[serial_test::serial(tx_registry)]
async fn issue_1029_starvation_warn_reports_registered_transactions() {
let (_dir, backend) = issue_1029_pool(false);
let text = backend.text("entities").expect("text store");
let holder = backend
.pool
.open_standalone_writer()
.expect("holder connection");
holder
.execute_batch("BEGIN IMMEDIATE")
.expect("holder BEGIN IMMEDIATE");
let fixture =
khive_storage::tx_registry::register(Some("issue_1029_fixture_tx".to_string()));
let events = Arc::new(std::sync::Mutex::new(Vec::new()));
let subscriber = StarvationCaptureSubscriber {
events: Arc::clone(&events),
};
let guard = tracing::subscriber::set_default(subscriber);
let doc = khive_storage::types::TextDocument {
subject_id: uuid::Uuid::new_v4(),
kind: khive_types::SubstrateKind::Entity,
title: Some("Issue1029Starved".to_string()),
body: "issue 1029 starvation diagnostic body".to_string(),
tags: vec![],
namespace: "tenant_ns".to_string(),
metadata: None,
updated_at: chrono::Utc::now(),
};
let result = text.upsert_document(doc).await;
drop(guard);
drop(fixture);
holder
.execute_batch("ROLLBACK")
.expect("holder ROLLBACK releases the lock");
assert!(
result.is_err(),
"upsert_document must starve while another connection holds the write lock"
);
let events = events.lock().unwrap();
let warn = events
.iter()
.find(|fields| {
fields
.get("message")
.is_some_and(|m| m.contains("text write starved"))
})
.unwrap_or_else(|| panic!("expected a starvation WARN, captured events: {events:?}"));
assert!(
warn.get("op").is_some_and(|op| op.contains("fts_upsert")),
"WARN must name the starved operation, got: {warn:?}"
);
assert!(
warn.get("open_txs")
.is_some_and(|txs| txs.contains("issue_1029_fixture_tx")),
"WARN must list the registered holder label, got: {warn:?}"
);
let count: usize = warn
.get("open_tx_count")
.expect("WARN must carry open_tx_count")
.parse()
.expect("open_tx_count must be numeric");
assert!(
count >= 1,
"open_tx_count must count the fixture, got {count}"
);
}
}