mod code_map;
#[path = "pool/identity_registry.rs"]
mod identity_registry;
#[cfg(any(test, feature = "test-support"))]
#[path = "pool/test_volume_lock_dir.rs"]
mod test_volume_lock_dir;
#[path = "pool/write_units.rs"]
mod write_units;
#[path = "pool/writer_acquisition.rs"]
mod writer_acquisition;
#[cfg(test)]
use identity_registry::pool_identity_suffix;
use identity_registry::PoolIdentityRegistration;
#[cfg(any(test, feature = "test-support"))]
use test_volume_lock_dir::test_volume_lock_dir;
#[cfg(test)]
use write_units::STARTUP_SPACE_PROBE;
pub use write_units::{CheckpointGuard, CheckpointResult, WriterAcquisitionSnapshot, WriterGuard};
pub(crate) use write_units::{
PooledAutocommitWriteUnit, PooledTransactionWriteUnit, StandaloneTransactionWriteUnit,
WriteAdmission, WriterAcquisitionCounters,
};
use crossbeam_queue::ArrayQueue;
use parking_lot::{Condvar, Mutex};
use rusqlite::hooks::{AuthContext, Authorization};
use rusqlite::{Connection, OpenFlags};
use serde::Serialize;
use sha2::{Digest, Sha256};
use std::cell::Cell;
use std::collections::{BTreeMap, HashMap};
use std::fs;
use std::io::Read as _;
use std::ops::{Deref, DerefMut};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::{Arc, OnceLock};
use std::thread;
use std::time::{Duration, Instant};
use tokio::sync::Semaphore;
use crate::database_owner_identity::{DatabaseOwnerIdentity, DatabaseOwnerIdentityError};
use crate::disk_guard::{LeaseRefusal, VolumeIdentity, VolumeLease};
#[cfg(test)]
use crate::disk_guard_config::DiskGuardConfigSource;
use crate::disk_guard_config::{
resolve_disk_guard_config, EffectiveDiskGuardConfig, DEFAULT_DISK_GUARD_DEADLINE_MS,
};
use crate::error::SqliteError;
#[cfg(windows)]
use crate::file_identity::sqlite_opened_file_identity;
#[cfg(any(unix, windows))]
use crate::file_identity::{database_file_identity, DatabaseFileIdentity};
use crate::writer_task::{execute_wrapped_transaction, WriterTaskHandle};
use khive_storage::error::StorageError;
use khive_storage::tx_registry::{DbIdentity, TxOrigin};
use khive_storage::CapacityUnavailablePhase;
use khive_storage::StorageCapability;
mod claimed_file_identity;
#[cfg(unix)]
mod claimed_file_observer;
#[cfg(unix)]
pub use claimed_file_observer::initialize as initialize_claimed_file_observer;
#[cfg(all(test, unix))]
#[path = "pool/claimed_file_identity_tests.rs"]
mod claimed_file_identity_tests;
pub(crate) const CACHE_SIZE_KIB: &str = "-65536";
const MMAP_SIZE_BYTES: &str = "1073741824";
const DEFAULT_READER_CAP: usize = 8;
const DEFAULT_JOURNAL_SIZE_LIMIT_BYTES: i64 = 67_108_864; const DEFAULT_WRITE_QUEUE_CAPACITY: usize = 256;
const DATABASE_ID_TABLE: &str = "_khive_database_identity";
static NEXT_MAIN_POOL_GENERATION: AtomicU64 = AtomicU64::new(1);
#[cfg(test)]
#[derive(Clone, Copy, PartialEq, Eq)]
enum IdentityOpenStage {
AfterMainOpenBeforeFirstStat,
AfterInitialIdentityWrite,
BeforeStandaloneOpen,
AfterStandaloneOpen,
}
#[cfg(test)]
type IdentityOpenHook = Box<dyn Fn(&Path, IdentityOpenStage, Option<&Connection>)>;
#[cfg(test)]
thread_local! {
static IDENTITY_OPEN_HOOK: std::cell::RefCell<Option<IdentityOpenHook>> =
const { std::cell::RefCell::new(None) };
}
#[cfg(test)]
fn run_identity_open_hook(path: &Path, stage: IdentityOpenStage, conn: Option<&Connection>) {
IDENTITY_OPEN_HOOK.with(|hook| {
if let Some(hook) = hook.borrow().as_ref() {
hook(path, stage, conn);
}
});
}
#[path = "pool/admission.rs"]
mod admission;
use admission::CheckpointOwnershipGate;
pub use admission::RuntimeWriteOperation;
pub(crate) use admission::FALLBACK_WAL_AUTOCHECKPOINT_PAGES;
#[cfg(test)]
use admission::{CheckpointConnectionConfigPause, CheckpointOwnership};
fn deny_retired_writer(_context: AuthContext<'_>) -> Authorization {
Authorization::Deny
}
pub(crate) const TEST_HARNESS_ENV: &str = "KHIVE_TEST_HARNESS";
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum WalCeilingSource {
BackendField,
Environment,
#[default]
Default,
}
#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
pub struct WalCeilingPolicy {
pub bytes: u64,
pub source: WalCeilingSource,
}
impl WalCeilingPolicy {
pub fn validate_static(
self,
file_backed: bool,
wal_mode: bool,
read_only: bool,
) -> Result<(), SqliteError> {
if self.bytes == 0 {
return Ok(());
}
if i64::try_from(self.bytes).is_err() {
return Err(SqliteError::WalCeilingOffsetOverflow { bytes: self.bytes });
}
if read_only {
return Ok(());
}
if !file_backed {
return Err(SqliteError::WalCeilingUnsupported {
bytes: self.bytes,
backend_kind: "in-memory backend",
});
}
if !wal_mode {
return Err(SqliteError::WalCeilingUnsupported {
bytes: self.bytes,
backend_kind: "non-WAL backend",
});
}
Ok(())
}
pub fn effective_bytes(self, read_only: bool) -> u64 {
if read_only {
0
} else {
self.bytes
}
}
}
#[derive(Clone, Debug)]
pub struct PoolConfig {
pub path: Option<PathBuf>,
pub code_map_vfs: Option<String>,
#[cfg(any(unix, windows))]
pub expected_file_identity: Option<DatabaseFileIdentity>,
pub max_readers: usize,
pub wal_mode: bool,
pub busy_timeout: Duration,
pub checkout_timeout: Duration,
pub journal_size_limit_bytes: i64,
pub read_only: bool,
pub wal_ceiling: WalCeilingPolicy,
pub write_queue_enabled: Option<bool>,
pub write_queue_capacity: usize,
pub write_routing_strict: bool,
pub write_admission_deadline_ms: u64,
pub disk_guard_config: Option<EffectiveDiskGuardConfig>,
pub volume_lock_dir: Option<PathBuf>,
pub read_tx_max_age: Duration,
}
const WRITE_ADMISSION_DEADLINE_MS_RANGE: std::ops::RangeInclusive<u64> = 100..=10_000;
const DEFAULT_WRITE_ADMISSION_DEADLINE_MS: u64 = 2000;
impl Default for PoolConfig {
fn default() -> Self {
Self {
path: None,
code_map_vfs: None,
#[cfg(any(unix, windows))]
expected_file_identity: None,
max_readers: std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(1)
.clamp(1, DEFAULT_READER_CAP),
wal_mode: true,
busy_timeout: Duration::from_secs(crate::env::env_parse_or(
"KHIVE_BUSY_TIMEOUT_SECS",
30,
)),
checkout_timeout: Duration::from_secs(crate::env::env_parse_or(
"KHIVE_CHECKOUT_TIMEOUT_SECS",
5,
)),
journal_size_limit_bytes: crate::env::env_parse_or(
"KHIVE_JOURNAL_SIZE_LIMIT_BYTES",
DEFAULT_JOURNAL_SIZE_LIMIT_BYTES,
),
read_only: false,
wal_ceiling: WalCeilingPolicy::default(),
write_queue_enabled: std::env::var_os("KHIVE_WRITE_QUEUE").map(|v| {
v.to_str()
.is_some_and(|v| v == "1" || v.eq_ignore_ascii_case("true"))
}),
write_queue_capacity: std::env::var("KHIVE_WRITE_QUEUE_CAPACITY")
.ok()
.and_then(|v| v.parse::<usize>().ok())
.filter(|&n| n > 0)
.unwrap_or(DEFAULT_WRITE_QUEUE_CAPACITY),
write_routing_strict: std::env::var("KHIVE_WRITE_ROUTING")
.map(|v| v.eq_ignore_ascii_case("strict"))
.unwrap_or(false),
write_admission_deadline_ms: crate::env::env_parse_or(
"KHIVE_WRITE_ADMISSION_DEADLINE_MS",
DEFAULT_WRITE_ADMISSION_DEADLINE_MS,
),
disk_guard_config: None,
#[cfg(test)]
volume_lock_dir: Some(test_volume_lock_dir()),
#[cfg(not(test))]
volume_lock_dir: crate::default_volume_lock_dir().ok(),
read_tx_max_age: crate::checkpoint::tx_age_thresholds_from_env(
Duration::from_secs(30),
Duration::from_secs(120),
)
.1,
}
}
}
#[cfg(any(test, feature = "test-support"))]
impl PoolConfig {
pub fn for_test() -> Self {
Self {
max_readers: 2,
volume_lock_dir: Some(test_volume_lock_dir()),
..Self::default()
}
}
}
fn refuse_home_data_store_in_tests(config: &PoolConfig) -> Result<(), SqliteError> {
if std::env::var(TEST_HARNESS_ENV).as_deref() != Ok("1") {
return Ok(());
}
let Some(path) = config.path.as_deref() else {
return Ok(());
};
if path
.as_os_str()
.as_encoded_bytes()
.get(..5)
.is_some_and(|prefix| prefix.eq_ignore_ascii_case(b"file:"))
{
return Err(SqliteError::InvalidData(format!(
"test harness refused SQLite URI database path {}; use a filesystem path outside \
HOME/.khive (deliberate sessions against a real store run the built binary \
directly, outside the Cargo test environment)",
path.display()
)));
}
let Some(home) = std::env::var_os("HOME") else {
return Ok(());
};
let canonical_path = canonicalize_deepest_existing(path)?;
let canonical_home_data_dir =
canonicalize_deepest_existing(&PathBuf::from(home).join(".khive"))?;
if canonical_path.starts_with(&canonical_home_data_dir) {
return Err(SqliteError::InvalidData(format!(
"test harness refused to open SQLite database under HOME/.khive: {} \
(deliberate sessions against a real store run the built binary directly, \
outside the Cargo test environment)",
canonical_path.display()
)));
}
Ok(())
}
fn canonicalize_deepest_existing(path: &Path) -> Result<PathBuf, SqliteError> {
let absolute = if path.is_absolute() {
path.to_path_buf()
} else {
std::env::current_dir().map_err(SqliteError::Io)?.join(path)
};
for ancestor in absolute.ancestors() {
match fs::canonicalize(ancestor) {
Ok(mut canonical) => {
let missing = absolute.strip_prefix(ancestor).map_err(|error| {
SqliteError::InvalidData(format!(
"failed to preserve missing path components for {}: {error}",
absolute.display()
))
})?;
canonical.push(missing);
return Ok(canonical);
}
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
Err(error) => {
return Err(SqliteError::InvalidData(format!(
"failed to canonicalize database path ancestor {}: {error}",
ancestor.display()
)));
}
}
}
Err(SqliteError::InvalidData(format!(
"database path has no canonicalizable ancestor: {}",
absolute.display()
)))
}
fn validate_write_admission_deadline(deadline_ms: u64) -> Result<(), SqliteError> {
if WRITE_ADMISSION_DEADLINE_MS_RANGE.contains(&deadline_ms) {
return Ok(());
}
Err(SqliteError::InvalidConfig(format!(
"write_admission_deadline_ms must be in [{}, {}] ms, got {deadline_ms}",
WRITE_ADMISSION_DEADLINE_MS_RANGE.start(),
WRITE_ADMISSION_DEADLINE_MS_RANGE.end()
)))
}
#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize)]
pub struct SearchMechanismSnapshot {
pub dispatches_by_backend_and_kind: BTreeMap<String, BTreeMap<String, u64>>,
pub note_candidate_hydration_rows: u64,
}
pub struct ConnectionPool {
writer: Arc<Mutex<Connection>>,
retirement_connection: Mutex<Option<Connection>>,
#[cfg(any(test, feature = "test-support"))]
statement_observer: Arc<crate::statement_observer::StatementObserverHub>,
main_pool_generation: OnceLock<u64>,
checkpoint_ownership: CheckpointOwnershipGate,
pooled_writer_retired: AtomicBool,
writer_acquisition_counters: Arc<WriterAcquisitionCounters>,
write_admission: Arc<WriteAdmission>,
disk_guard_config: Option<EffectiveDiskGuardConfig>,
reader_acquisition_counters: ReaderAcquisitionCounters,
search_dispatches: Mutex<BTreeMap<String, BTreeMap<String, u64>>>,
note_candidate_hydration_rows: AtomicU64,
readers: ArrayQueue<Connection>,
max_readers: usize,
config: PoolConfig,
read_only_open_target: Option<PathBuf>,
sql_bridge_reader_slots: Arc<Semaphore>,
sql_bridge_writer_slots: Arc<Semaphore>,
writer_task: OnceLock<Option<WriterTaskHandle>>,
writer_task_join: Mutex<Option<tokio::task::JoinHandle<()>>>,
writer_task_join_stored: AtomicBool,
origin: TxOrigin,
identity_path: Option<PathBuf>,
#[cfg(any(unix, windows))]
opened_file_identity: Option<DatabaseFileIdentity>,
opened_database_id: Option<uuid::Uuid>,
identity_registration: Option<PoolIdentityRegistration>,
#[cfg(test)]
writer_task_spawn_count: std::sync::atomic::AtomicUsize,
}
impl Drop for ConnectionPool {
fn drop(&mut self) {
while let Some(conn) = self.readers.pop() {
drop(conn);
}
}
}
enum ReaderLease<'pool> {
Pooled(Connection),
Shared(parking_lot::MutexGuard<'pool, Connection>),
}
pub struct ReaderRow<'row, 'statement> {
row: &'row rusqlite::Row<'statement>,
}
impl ReaderRow<'_, '_> {
pub fn get<I: rusqlite::RowIndex, T: rusqlite::types::FromSql>(
&self,
index: I,
) -> rusqlite::Result<T> {
self.row.get(index)
}
pub fn get_ref<I: rusqlite::RowIndex>(
&self,
index: I,
) -> rusqlite::Result<rusqlite::types::ValueRef<'_>> {
self.row.get_ref(index)
}
}
struct ReaderQueryInProgress<'a>(&'a Cell<bool>);
impl Drop for ReaderQueryInProgress<'_> {
fn drop(&mut self) {
self.0.set(false);
}
}
pub(crate) struct ReaderAdmission {
slot: tokio::sync::OwnedSemaphorePermit,
started: Instant,
}
pub struct ReaderGuard<'pool> {
lease: Option<ReaderLease<'pool>>,
admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
pool: &'pool ConnectionPool,
reusable: Cell<bool>,
query_in_progress: Cell<bool>,
checked_out_at: Instant,
dirty: Cell<bool>,
operation: Option<&'static str>,
}
impl<'pool> ReaderGuard<'pool> {
pub(crate) fn conn(&self) -> &Connection {
match self
.lease
.as_ref()
.expect("reader guard missing connection")
{
ReaderLease::Pooled(conn) => conn,
ReaderLease::Shared(guard) => guard,
}
}
pub fn query_row<T, P, F>(&self, sql: &str, params: P, f: F) -> Result<T, SqliteError>
where
P: rusqlite::Params,
F: FnOnce(&ReaderRow<'_, '_>) -> rusqlite::Result<T>,
{
crate::sql_bridge::reader_capability_admits(sql).map_err(SqliteError::InvalidData)?;
if !self.reusable.get() {
return Err(SqliteError::InvalidData(
"reader lease is quarantined after failed read cleanup".into(),
));
}
if self.query_in_progress.replace(true) {
return Err(SqliteError::InvalidData(
"reader lease is already executing a query".into(),
));
}
let _in_progress = ReaderQueryInProgress(&self.query_in_progress);
self.mark_dirty();
crate::read_cancellation::run_borrowed_reader(self, |conn, admission| {
conn.query_row(sql, params, |row| {
if !admission.admits() {
return Err(rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_INTERRUPT),
Some("request stopped before the mapper".into()),
));
}
f(&ReaderRow { row })
})
.map_err(|error| {
StorageError::driver(StorageCapability::Sql, "reader_guard.query_row", error)
})
})
.map_err(|error| {
self.pool.record_reader_query_error(&error);
match error {
StorageError::Driver {
capability,
operation,
source,
} => match source.downcast::<rusqlite::Error>() {
Ok(error) => SqliteError::Rusqlite(*error),
Err(source) => SqliteError::RequestReadStopped(StorageError::Driver {
capability,
operation,
source,
}),
},
other => SqliteError::RequestReadStopped(other),
}
})
}
pub(crate) fn discard(&self) {
self.reusable.set(false);
}
pub(crate) fn mark_dirty(&self) {
self.dirty.set(true);
}
pub(crate) fn label_operation(&mut self, operation: &'static str) {
self.operation = Some(operation);
}
}
impl<'pool> Drop for ReaderGuard<'pool> {
fn drop(&mut self) {
let Some(lease) = self.lease.take() else {
return;
};
match lease {
ReaderLease::Pooled(conn) if self.reusable.get() => {
self.pool.return_reader(conn, self.dirty.get())
}
ReaderLease::Pooled(conn) => {
close_connection_quietly(conn);
self.pool.replace_discarded_reader_slot();
}
ReaderLease::Shared(guard) if !self.reusable.get() => {
self.pool.retire_pooled_writer(&guard);
}
ReaderLease::Shared(guard) => {
if self.dirty.get() && !restore_shared_reader_state(&guard, &self.pool.config) {
self.pool.retire_pooled_writer(&guard);
}
}
}
drop(self.admission_slot.take());
self.pool
.reader_acquisition_counters
.record_checkout_completed(self.checked_out_at.elapsed(), self.operation);
}
}
pub(crate) struct SharedReaderTransactionGuard {
conn: parking_lot::ArcMutexGuard<parking_lot::RawMutex, Connection>,
admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
pool: Arc<ConnectionPool>,
checked_out_at: Instant,
poison: Cell<bool>,
}
impl SharedReaderTransactionGuard {
pub(crate) fn conn(&self) -> &Connection {
&self.conn
}
pub(crate) fn poison(&self) {
self.poison.set(true);
}
}
impl Drop for SharedReaderTransactionGuard {
fn drop(&mut self) {
let mut restored = !self.poison.get();
if restored && !self.conn.is_autocommit() {
restored = self.conn.execute_batch("ROLLBACK").is_ok() && self.conn.is_autocommit();
}
if restored {
restored = restore_shared_reader_state(&self.conn, &self.pool.config);
}
if !restored {
self.pool.retire_pooled_writer(&self.conn);
}
drop(self.admission_slot.take());
self.pool
.reader_acquisition_counters
.record_checkout_completed(
self.checked_out_at.elapsed(),
Some("explicit_sql_read_transaction"),
);
}
}
impl ConnectionPool {
pub(crate) fn checkout_shared_reader_transaction(
self: &Arc<Self>,
should_stop: impl Fn() -> bool,
) -> Result<Option<SharedReaderTransactionGuard>, SqliteError> {
debug_assert_eq!(
self.max_readers, 0,
"the owned shared-reader-transaction guard exists only for the degraded, \
single-connection backend"
);
self.ensure_pooled_writer_active()?;
let started = Instant::now();
let admission_slot = loop {
if should_stop() {
return Ok(None);
}
match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
Ok(slot) => break slot,
Err(tokio::sync::TryAcquireError::Closed) => {
return Err(SqliteError::InvalidData(
"reader admission semaphore is closed".to_string(),
));
}
Err(tokio::sync::TryAcquireError::NoPermits) => {}
}
if started.elapsed() >= self.config.checkout_timeout {
self.reader_acquisition_counters.record_checkout_timeout();
return Err(pool_exhausted_error(
self.config.checkout_timeout,
self.max_readers,
));
}
thread::yield_now();
};
loop {
if should_stop() {
return Ok(None);
}
let remaining = self
.config
.checkout_timeout
.saturating_sub(started.elapsed());
if remaining.is_zero() {
self.reader_acquisition_counters.record_checkout_timeout();
return Err(pool_exhausted_error(
self.config.checkout_timeout,
self.max_readers,
));
}
if let Some(conn) = self
.writer
.try_lock_arc_for(remaining.min(Duration::from_millis(2)))
{
self.ensure_pooled_writer_active()?;
self.reader_acquisition_counters.record_pooled_checkout();
return Ok(Some(SharedReaderTransactionGuard {
conn,
admission_slot: Some(admission_slot),
pool: Arc::clone(self),
checked_out_at: Instant::now(),
poison: Cell::new(false),
}));
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[allow(dead_code)] pub(crate) enum StandaloneReaderPurpose {
ExplicitSqlReadTransaction,
BootSchemaProbe,
DiagnosticsIndependentSnapshot,
}
impl StandaloneReaderPurpose {
fn is_infrastructure(self) -> bool {
matches!(
self,
Self::BootSchemaProbe | Self::DiagnosticsIndependentSnapshot
)
}
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct ReaderAcquisitionSnapshot {
pub reader_admission_capacity: usize,
pub available_reader_admission_slots: usize,
pub acquisitions: u64,
pub pooled_checkouts: u64,
pub standalone_opens: u64,
pub infrastructure_standalone_opens: u64,
pub checkout_timeouts: u64,
pub busy_timeouts: u64,
pub active_pooled_checkouts: u64,
pub peak_active_pooled_checkouts: u64,
pub completed_pooled_checkouts: u64,
pub max_completed_hold_micros: u64,
pub max_completed_hold_operation: Option<&'static str>,
pub reader_replacement_open_failures: u64,
}
#[derive(Debug, Default, Clone, Copy)]
struct LongestCompletedHold {
micros: u64,
operation: Option<&'static str>,
}
#[derive(Debug, Default)]
struct ReaderAcquisitionCounters {
pooled_checkouts: AtomicU64,
standalone_opens: AtomicU64,
infrastructure_standalone_opens: AtomicU64,
checkout_timeouts: AtomicU64,
busy_timeouts: AtomicU64,
active_pooled_checkouts: AtomicU64,
peak_active_pooled_checkouts: AtomicU64,
completed_pooled_checkouts: AtomicU64,
longest_completed_hold: parking_lot::Mutex<LongestCompletedHold>,
reader_replacement_open_failures: AtomicU64,
}
impl ReaderAcquisitionCounters {
fn record_pooled_checkout(&self) {
self.pooled_checkouts.fetch_add(1, Ordering::Relaxed);
let active = self
.active_pooled_checkouts
.fetch_add(1, Ordering::Relaxed)
.saturating_add(1);
self.peak_active_pooled_checkouts
.fetch_max(active, Ordering::Relaxed);
}
fn record_checkout_timeout(&self) {
self.checkout_timeouts.fetch_add(1, Ordering::Relaxed);
}
fn record_busy_timeout(&self) {
self.busy_timeouts.fetch_add(1, Ordering::Relaxed);
}
fn record_reader_replacement_open_failure(&self) {
self.reader_replacement_open_failures
.fetch_add(1, Ordering::Relaxed);
}
fn record_standalone_open(&self, purpose: StandaloneReaderPurpose) {
if purpose.is_infrastructure() {
self.infrastructure_standalone_opens
.fetch_add(1, Ordering::Relaxed);
} else {
self.standalone_opens.fetch_add(1, Ordering::Relaxed);
}
}
fn record_checkout_completed(&self, hold: Duration, operation: Option<&'static str>) {
let previous = self.active_pooled_checkouts.fetch_sub(1, Ordering::Relaxed);
debug_assert!(previous > 0, "reader active-checkout counter underflow");
self.completed_pooled_checkouts
.fetch_add(1, Ordering::Relaxed);
let micros = u64::try_from(hold.as_micros()).unwrap_or(u64::MAX);
let mut longest = self.longest_completed_hold.lock();
if micros > longest.micros {
longest.micros = micros;
longest.operation = operation;
}
}
fn snapshot(
&self,
reader_admission_capacity: usize,
available_reader_admission_slots: usize,
) -> ReaderAcquisitionSnapshot {
let pooled_checkouts = self.pooled_checkouts.load(Ordering::Relaxed);
let standalone_opens = self.standalone_opens.load(Ordering::Relaxed);
let longest_completed_hold = *self.longest_completed_hold.lock();
ReaderAcquisitionSnapshot {
reader_admission_capacity,
available_reader_admission_slots,
acquisitions: pooled_checkouts.saturating_add(standalone_opens),
pooled_checkouts,
standalone_opens,
infrastructure_standalone_opens: self
.infrastructure_standalone_opens
.load(Ordering::Relaxed),
checkout_timeouts: self.checkout_timeouts.load(Ordering::Relaxed),
busy_timeouts: self.busy_timeouts.load(Ordering::Relaxed),
active_pooled_checkouts: self.active_pooled_checkouts.load(Ordering::Relaxed),
peak_active_pooled_checkouts: self.peak_active_pooled_checkouts.load(Ordering::Relaxed),
completed_pooled_checkouts: self.completed_pooled_checkouts.load(Ordering::Relaxed),
max_completed_hold_micros: longest_completed_hold.micros,
max_completed_hold_operation: longest_completed_hold.operation,
reader_replacement_open_failures: self
.reader_replacement_open_failures
.load(Ordering::Relaxed),
}
}
}
impl ConnectionPool {
pub fn new(config: PoolConfig) -> Result<Self, SqliteError> {
refuse_home_data_store_in_tests(&config)?;
validate_write_admission_deadline(config.write_admission_deadline_ms)?;
if let Some(policy) = config.disk_guard_config {
policy.validate()?;
}
config.wal_ceiling.validate_static(
config.path.is_some(),
config.wal_mode,
config.read_only,
)?;
code_map::validate_pool(&config)?;
if config.path.is_some() && !config.read_only {
let lock_dir =
crate::disk_guard_config::require_volume_lock_dir(config.volume_lock_dir.clone())?;
if !lock_dir.is_absolute() {
return Err(SqliteError::CapacityUnavailable {
phase: CapacityUnavailablePhase::Lock,
message: "configured volume-lock directory is not absolute".to_string(),
});
}
}
let mut config = config;
let inert_memory_queue_request =
config.path.is_none() && config.write_queue_enabled == Some(true);
config.write_queue_enabled =
Some(config.write_queue_enabled.unwrap_or(config.path.is_some()));
if inert_memory_queue_request {
tracing::warn!(
"write queue explicitly requested for an in-memory pool; it is inert because \
in-memory pools cannot host a writer task"
);
}
let (origin, identity_path) = match config.path.as_ref() {
Some(path) if config.code_map_vfs.is_some() => code_map::guarded_identity(path),
Some(path) => {
let (identity, canonical) = mint_db_identity(path)?;
(TxOrigin::Database(identity), Some(canonical))
}
None => (TxOrigin::Memory, None),
};
let (write_admission, disk_guard_config) = if !config.read_only {
if identity_path.as_deref().and_then(Path::parent).is_some() {
let disk_guard = match config.disk_guard_config {
Some(policy) => policy,
None => resolve_disk_guard_config(None, None)?,
};
disk_guard.validate()?;
if disk_guard.legacy_environment_present {
tracing::warn!(
"legacy SQLite reserve setting is deprecated; \
use KHIVE_SQLITE_DISK_RESERVE_BYTES"
);
}
if disk_guard.reserve_bytes == 0 {
tracing::warn!(
"SQLite disk reserve is explicitly zero; \
new logical writes will not be floor-refused"
);
}
(
Arc::new(WriteAdmission::new(
identity_path.clone(),
disk_guard.reserve_bytes,
disk_guard.guard_deadline_ms,
config.volume_lock_dir.clone(),
)?),
Some(disk_guard),
)
} else {
(
Arc::new(WriteAdmission::new(
None,
0,
DEFAULT_DISK_GUARD_DEADLINE_MS,
None,
)?),
None,
)
}
} else {
(
Arc::new(WriteAdmission::new(
None,
0,
DEFAULT_DISK_GUARD_DEADLINE_MS,
None,
)?),
None,
)
};
let read_only_open_target = read_only_open_target(&config, identity_path.as_deref())?;
#[cfg(any(unix, windows))]
let identity_before_open = identity_path
.as_deref()
.map(database_file_identity_if_exists)
.transpose()?
.flatten();
#[cfg(any(unix, windows))]
claimed_file_identity::verify_before_open(&config, identity_before_open)?;
let retirement_connection = Connection::open_in_memory()?;
retirement_connection.authorizer(Some(deny_retired_writer))?;
let mut initialization_lease = None;
let mut writer = open_writer_connection(
&config,
read_only_open_target.as_deref(),
identity_path.as_deref(),
)?;
validate_wal_ceiling_at_open(&writer, &config)?;
writer.busy_timeout(config.busy_timeout)?;
let initial_database_id = if identity_path.is_some() {
read_database_id(&writer)?
} else {
None
};
#[cfg(test)]
if let Some(path) = identity_path.as_deref() {
run_identity_open_hook(
path,
IdentityOpenStage::AfterMainOpenBeforeFirstStat,
Some(&writer),
);
}
#[cfg(any(unix, windows))]
let identity_before_write = identity_path
.as_deref()
.map(database_file_identity)
.transpose()?;
#[cfg(any(unix, windows))]
if identity_before_open.is_some() && identity_before_open != identity_before_write {
return Err(SqliteError::InvalidData(
"database file identity changed while opening the pool".to_string(),
));
}
#[cfg(any(unix, windows))]
if let Some(path) = identity_path.as_deref() {
let opened = opened_sqlite_file_identity(&writer, path)?;
if identity_before_write != Some(opened) {
return Err(SqliteError::InvalidData(
"database file identity changed while opening the pool".to_string(),
));
}
}
write_admission.verify_current_volume()?;
let opened_database_id = if identity_path.is_some()
&& !config.read_only
&& initial_database_id.is_none()
{
match write_admission.check() {
Ok(()) => {
initialization_lease = write_admission.acquire()?;
match initialize_database_id_with_admission(&mut writer, Some(&write_admission))
{
Ok(id) => Some(id),
Err(SqliteError::CapacityFloor { .. }) => None,
Err(error) => return Err(error),
}
}
Err(SqliteError::CapacityFloor { .. }) => None,
Err(error) => return Err(error),
}
} else {
initial_database_id
};
drop(initialization_lease.take());
#[cfg(test)]
if let Some(path) = identity_path.as_deref() {
run_identity_open_hook(
path,
IdentityOpenStage::AfterInitialIdentityWrite,
Some(&writer),
);
}
#[cfg(any(unix, windows))]
let opened_file_identity = identity_path
.as_deref()
.map(|path| opened_sqlite_file_identity(&writer, path))
.transpose()?;
#[cfg(any(unix, windows))]
if identity_before_write != opened_file_identity {
return Err(SqliteError::InvalidData(
"database file identity changed while opening the pool".to_string(),
));
}
let wal_enabled = configure_writer_connection(&writer, &config)?;
let max_readers = effective_reader_count(&config, wal_enabled);
let readers = ArrayQueue::new(max_readers.max(1));
#[cfg(any(test, feature = "test-support"))]
let statement_observer = crate::statement_observer::StatementObserverHub::new()?;
#[cfg(any(test, feature = "test-support"))]
crate::statement_observer::install(&writer, &statement_observer)?;
let mut pool = Self {
writer: Arc::new(Mutex::new(writer)),
retirement_connection: Mutex::new(Some(retirement_connection)),
#[cfg(any(test, feature = "test-support"))]
statement_observer,
main_pool_generation: OnceLock::new(),
checkpoint_ownership: CheckpointOwnershipGate::new(),
pooled_writer_retired: AtomicBool::new(false),
writer_acquisition_counters: Arc::new(WriterAcquisitionCounters::default()),
write_admission,
disk_guard_config,
reader_acquisition_counters: ReaderAcquisitionCounters::default(),
search_dispatches: Mutex::new(BTreeMap::new()),
note_candidate_hydration_rows: AtomicU64::new(0),
readers,
max_readers,
config,
read_only_open_target,
sql_bridge_reader_slots: Arc::new(Semaphore::new(max_readers.max(1))),
sql_bridge_writer_slots: Arc::new(Semaphore::new(1)),
writer_task: OnceLock::new(),
writer_task_join: Mutex::new(None),
writer_task_join_stored: AtomicBool::new(false),
origin,
identity_path,
#[cfg(any(unix, windows))]
opened_file_identity,
opened_database_id,
identity_registration: None,
#[cfg(test)]
writer_task_spawn_count: std::sync::atomic::AtomicUsize::new(0),
};
for _ in 0..pool.max_readers {
let conn = pool.open_reader_connection()?;
pool.readers
.push(conn)
.expect("reader queue must have capacity during pool initialization");
}
if !pool.config.read_only {
crate::timeout_sink::init(
pool.canonical_path().and_then(Path::parent),
&crate::timeout_sink::db_label(&pool),
);
}
pool.identity_registration = pool.canonical_path().map(PoolIdentityRegistration::new);
Ok(pool)
}
pub fn reader(&self) -> Result<ReaderGuard<'_>, SqliteError> {
self.reader_until(|| false)?.ok_or_else(|| {
SqliteError::InvalidData("uncancelled reader checkout stopped unexpectedly".into())
})
}
pub(crate) fn reader_until<C>(
&self,
should_stop: C,
) -> Result<Option<ReaderGuard<'_>>, SqliteError>
where
C: Fn() -> bool,
{
let started = Instant::now();
let mut admission_attempt = 0u32;
let admission_slot = loop {
if should_stop() {
return Ok(None);
}
match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
Ok(slot) => break slot,
Err(tokio::sync::TryAcquireError::Closed) => {
return Err(SqliteError::InvalidData(
"reader admission semaphore is closed".to_string(),
));
}
Err(tokio::sync::TryAcquireError::NoPermits) => {}
}
if started.elapsed() >= self.config.checkout_timeout {
self.reader_acquisition_counters.record_checkout_timeout();
return Err(pool_exhausted_error(
self.config.checkout_timeout,
self.max_readers,
));
}
match admission_attempt {
0..=7 => {
let spins = 1usize << admission_attempt;
for _ in 0..spins {
std::hint::spin_loop();
}
}
8..=15 => thread::yield_now(),
_ => {
let remaining = self
.config
.checkout_timeout
.saturating_sub(started.elapsed());
let sleep =
Duration::from_micros(50 * (1u64 << (admission_attempt - 16).min(6)));
thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
}
}
admission_attempt = admission_attempt.saturating_add(1);
};
self.reader_with_admission(
ReaderAdmission {
slot: admission_slot,
started,
},
should_stop,
)
}
pub(crate) async fn acquire_reader_admission(
&self,
capability: StorageCapability,
operation: &'static str,
) -> Result<ReaderAdmission, StorageError> {
let context = khive_storage::capture_request_read_context();
let started = Instant::now();
let stopped = context.stop_reason().is_some();
let outcome: Result<Option<ReaderAdmission>, SqliteError> = if stopped {
Ok(None)
} else {
tokio::select! {
biased;
_ = context.wait_for_stop() => Ok(None),
waited = tokio::time::timeout(
self.config.checkout_timeout,
Arc::clone(&self.sql_bridge_reader_slots).acquire_owned(),
) => match waited {
Ok(Ok(slot)) => Ok(Some(ReaderAdmission { slot, started })),
Ok(Err(_closed)) => Err(SqliteError::InvalidData(
"reader admission semaphore is closed".to_string(),
)),
Err(_elapsed) => {
self.reader_acquisition_counters.record_checkout_timeout();
Err(pool_exhausted_error(
self.config.checkout_timeout,
self.max_readers,
))
}
},
}
};
match outcome {
Ok(Some(admission)) => Ok(admission),
Ok(None) => Err(self.reader_checkout_refusal(capability, operation, None)),
Err(error) => Err(self.reader_checkout_refusal(capability, operation, Some(error))),
}
}
pub(crate) fn reader_with_admission<C>(
&self,
admission: ReaderAdmission,
should_stop: C,
) -> Result<Option<ReaderGuard<'_>>, SqliteError>
where
C: Fn() -> bool,
{
let ReaderAdmission {
slot: admission_slot,
started,
} = admission;
if self.max_readers == 0 {
self.ensure_pooled_writer_active()?;
loop {
if should_stop() {
return Ok(None);
}
let remaining = self
.config
.checkout_timeout
.saturating_sub(started.elapsed());
if remaining.is_zero() {
self.reader_acquisition_counters.record_checkout_timeout();
return Err(pool_exhausted_error(
self.config.checkout_timeout,
self.max_readers,
));
}
if let Some(guard) = self
.writer
.try_lock_for(remaining.min(Duration::from_millis(2)))
{
self.ensure_pooled_writer_active()?;
self.reader_acquisition_counters.record_pooled_checkout();
return Ok(Some(ReaderGuard {
lease: Some(ReaderLease::Shared(guard)),
admission_slot: Some(admission_slot),
pool: self,
reusable: Cell::new(true),
query_in_progress: Cell::new(false),
checked_out_at: Instant::now(),
dirty: Cell::new(false),
operation: None,
}));
}
}
}
let mut attempt = 0u32;
loop {
if should_stop() {
return Ok(None);
}
if let Some(conn) = self.readers.pop() {
self.reader_acquisition_counters.record_pooled_checkout();
return Ok(Some(ReaderGuard {
lease: Some(ReaderLease::Pooled(conn)),
admission_slot: Some(admission_slot),
pool: self,
reusable: Cell::new(true),
query_in_progress: Cell::new(false),
checked_out_at: Instant::now(),
dirty: Cell::new(false),
operation: None,
}));
}
if started.elapsed() >= self.config.checkout_timeout {
self.reader_acquisition_counters.record_checkout_timeout();
return Err(pool_exhausted_error(
self.config.checkout_timeout,
self.max_readers,
));
}
match attempt {
0..=7 => {
let spins = 1usize << attempt;
for _ in 0..spins {
std::hint::spin_loop();
}
}
8..=15 => thread::yield_now(),
_ => {
let remaining = self
.config
.checkout_timeout
.saturating_sub(started.elapsed());
let sleep = Duration::from_micros(50 * (1u64 << (attempt - 16).min(6)));
thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
}
}
attempt = attempt.saturating_add(1);
}
}
#[track_caller]
pub fn writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
self.writer_with_checkout_probe(true, false)
}
#[track_caller]
pub(crate) fn writer_for_admitted_operation(&self) -> Result<WriterGuard<'_>, SqliteError> {
self.writer_with_checkout_probe(false, false)
}
fn writer_for_checkpoint_operation(&self) -> Result<WriterGuard<'_>, SqliteError> {
self.writer_with_checkout_probe(false, true)
}
#[track_caller]
fn writer_with_checkout_probe(
&self,
compatibility_probe: bool,
checkpoint_bypass: bool,
) -> Result<WriterGuard<'_>, SqliteError> {
if !checkpoint_bypass {
self.write_admission.ensure_settled()?;
}
self.ensure_pooled_writer_active()?;
let volume_lease = if checkpoint_bypass {
None
} else {
self.write_admission.acquire()?
};
let Some(guard) = self.writer.try_lock_for(self.config.checkout_timeout) else {
self.writer_acquisition_counters
.pooled_timeouts
.fetch_add(1, Ordering::Relaxed);
let message = format!(
"timed out after {:?} waiting for sqlite writer connection",
self.config.checkout_timeout
);
crate::timeout_sink::emit_timeout(
&crate::timeout_sink::db_label(self),
crate::timeout_sink::Site::PoolAdmission,
&message,
Some(
self.config
.checkout_timeout
.as_millis()
.min(u128::from(u64::MAX)) as u64,
),
);
return Err(SqliteError::WriterPoolCheckoutTimeout {
timeout: self.config.checkout_timeout,
});
};
self.ensure_pooled_writer_active()?;
#[cfg(any(unix, windows))]
if let Some(path) = self.identity_path.as_deref() {
self.verify_opened_file_identity(path)?;
self.verify_connection_file_identity(&guard, path)?;
}
if compatibility_probe && !checkpoint_bypass {
self.write_admission.check()?;
}
self.writer_acquisition_counters
.pooled_acquisitions
.fetch_add(1, Ordering::Relaxed);
Ok(WriterGuard {
guard,
origin: self.origin(),
pool: self,
admission: self.write_admission.as_ref(),
_volume_lease: volume_lease,
})
}
#[track_caller]
pub fn try_writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
self.writer()
}
#[cfg(test)]
#[track_caller]
fn writer_until<C>(&self, should_stop: C) -> Result<Option<WriterGuard<'_>>, SqliteError>
where
C: Fn() -> bool,
{
self.writer_until_with_checkout_probe(should_stop, true)
}
#[track_caller]
pub(crate) fn writer_until_for_admitted_operation<C>(
&self,
should_stop: C,
) -> Result<Option<WriterGuard<'_>>, SqliteError>
where
C: Fn() -> bool,
{
self.writer_until_with_checkout_probe(should_stop, false)
}
#[track_caller]
fn writer_until_with_checkout_probe<C>(
&self,
should_stop: C,
compatibility_probe: bool,
) -> Result<Option<WriterGuard<'_>>, SqliteError>
where
C: Fn() -> bool,
{
self.ensure_pooled_writer_active()?;
if should_stop() {
return Ok(None);
}
let volume_lease = self.write_admission.acquire()?;
let started = Instant::now();
loop {
if should_stop() {
return Ok(None);
}
let remaining = self
.config
.checkout_timeout
.saturating_sub(started.elapsed());
if let Some(guard) = self
.writer
.try_lock_for(remaining.min(Duration::from_millis(2)))
{
if should_stop() {
return Ok(None);
}
self.ensure_pooled_writer_active()?;
#[cfg(any(unix, windows))]
if let Some(path) = self.identity_path.as_deref() {
self.verify_opened_file_identity(path)?;
self.verify_connection_file_identity(&guard, path)?;
}
if compatibility_probe {
self.write_admission.check()?;
}
self.writer_acquisition_counters
.pooled_acquisitions
.fetch_add(1, Ordering::Relaxed);
return Ok(Some(WriterGuard {
guard,
origin: self.origin(),
pool: self,
admission: self.write_admission.as_ref(),
_volume_lease: volume_lease,
}));
}
if started.elapsed() >= self.config.checkout_timeout {
self.writer_acquisition_counters
.pooled_timeouts
.fetch_add(1, Ordering::Relaxed);
let message = format!(
"timed out after {:?} waiting for sqlite writer connection",
self.config.checkout_timeout
);
crate::timeout_sink::emit_timeout(
&crate::timeout_sink::db_label(self),
crate::timeout_sink::Site::PoolAdmission,
&message,
Some(
self.config
.checkout_timeout
.as_millis()
.min(u128::from(u64::MAX)) as u64,
),
);
return Err(SqliteError::WriterPoolCheckoutTimeout {
timeout: self.config.checkout_timeout,
});
}
}
}
pub fn try_checkpoint_nowait(&self) -> Result<CheckpointGuard<'_>, SqliteError> {
self.ensure_pooled_writer_active()?;
let guard = self.writer.try_lock().ok_or_else(|| {
SqliteError::InvalidData(
"writer connection busy (checkpoint skipped this tick)".to_string(),
)
})?;
self.ensure_pooled_writer_active()?;
Ok(CheckpointGuard { guard })
}
pub(crate) fn retire_pooled_writer(&self, conn: &Connection) {
self.pooled_writer_retired.store(true, Ordering::Release);
if let Err(error) = conn.authorizer(Some(deny_retired_writer)) {
tracing::error!(
%error,
"failed to install the retired pooled-writer quarantine authorizer"
);
}
}
#[cfg(test)]
pub(crate) fn probe_retired_pooled_writer_for_test(&self) -> rusqlite::Result<i64> {
self.writer
.lock()
.query_row("SELECT 1", [], |row| row.get(0))
}
#[cfg(test)]
pub(crate) fn leave_pooled_writer_transaction_open_for_test(&self) -> rusqlite::Result<()> {
self.writer.lock().execute_batch("BEGIN IMMEDIATE")
}
fn ensure_pooled_writer_active(&self) -> Result<(), SqliteError> {
if self.pooled_writer_retired.load(Ordering::Acquire) {
return Err(SqliteError::InvalidData(
"pooled writer connection retired after a terminal transaction fault".to_string(),
));
}
Ok(())
}
pub fn writer_acquisition_snapshot(&self) -> WriterAcquisitionSnapshot {
self.writer_acquisition_counters
.snapshot(self.write_admission.lease_timeouts())
}
pub fn record_search_dispatch(&self, backend_id: &str, requested_kind: &str) {
let mut dispatches = self.search_dispatches.lock();
let count = dispatches
.entry(backend_id.to_owned())
.or_default()
.entry(requested_kind.to_owned())
.or_default();
*count = count.saturating_add(1);
}
pub fn record_note_candidate_hydration_row(&self) {
let _ = self.note_candidate_hydration_rows.fetch_update(
Ordering::Relaxed,
Ordering::Relaxed,
|current| Some(current.saturating_add(1)),
);
}
pub fn search_mechanism_snapshot(&self) -> SearchMechanismSnapshot {
SearchMechanismSnapshot {
dispatches_by_backend_and_kind: self.search_dispatches.lock().clone(),
note_candidate_hydration_rows: self
.note_candidate_hydration_rows
.load(Ordering::Relaxed),
}
}
pub fn reader_acquisition_snapshot(&self) -> ReaderAcquisitionSnapshot {
self.reader_acquisition_counters.snapshot(
self.max_readers.max(1),
self.sql_bridge_reader_slots.available_permits(),
)
}
pub(crate) fn record_reader_admission_timeout(&self) {
self.reader_acquisition_counters.record_checkout_timeout();
}
pub(crate) fn record_reader_query_error(&self, error: &StorageError) {
if crate::read_cancellation::storage_error_sqlite_code(error)
== Some(rusqlite::ErrorCode::DatabaseBusy)
{
self.reader_acquisition_counters.record_busy_timeout();
}
}
pub(crate) fn writer_acquisition_counters(&self) -> Arc<WriterAcquisitionCounters> {
Arc::clone(&self.writer_acquisition_counters)
}
pub(crate) fn write_admission(&self) -> Arc<WriteAdmission> {
Arc::clone(&self.write_admission)
}
pub fn effective_disk_guard_config(&self) -> Option<EffectiveDiskGuardConfig> {
self.disk_guard_config
}
#[cfg(test)]
pub(crate) fn set_test_write_admission(
&mut self,
floor_bytes: u64,
probe: impl Fn(&Path) -> std::io::Result<u64> + Send + Sync + 'static,
) {
let database_path = if self.config.read_only {
None
} else {
self.canonical_path().map(Path::to_path_buf)
};
let admission = Arc::new(
WriteAdmission::new(
database_path,
floor_bytes,
DEFAULT_DISK_GUARD_DEADLINE_MS,
self.config.volume_lock_dir.clone(),
)
.expect("test admission volume identity"),
);
admission.set_test_space_probe(probe);
self.write_admission = admission;
if self.canonical_path().is_some() && !self.config.read_only {
self.disk_guard_config = Some(EffectiveDiskGuardConfig {
reserve_bytes: floor_bytes,
guard_deadline_ms: DEFAULT_DISK_GUARD_DEADLINE_MS,
reserve_source: DiskGuardConfigSource::Backend,
deadline_source: DiskGuardConfigSource::Default,
legacy_environment_present: false,
});
}
}
pub fn available_readers(&self) -> usize {
self.readers.len()
}
pub fn max_readers(&self) -> usize {
self.max_readers
}
pub fn config(&self) -> &PoolConfig {
&self.config
}
#[cfg(any(test, feature = "test-support"))]
pub fn observe_test_statement_starts(
&self,
limit: usize,
) -> Result<crate::statement_observer::StatementStartObservation, SqliteError> {
self.statement_observer.observe(limit)
}
pub fn main_pool_generation(&self) -> u64 {
*self.main_pool_generation.get_or_init(|| {
NEXT_MAIN_POOL_GENERATION
.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |next| {
next.checked_add(1)
})
.expect("main pool generation exhausted")
})
}
pub(crate) fn reader_admission_timeout(&self, operation: &'static str) -> StorageError {
StorageError::AdmissionTimeout {
operation: operation.into(),
timeout_ms: u64::try_from(self.config.checkout_timeout.as_millis()).unwrap_or(u64::MAX),
pool_identity: Some(
self.identity_registration
.as_ref()
.map(PoolIdentityRegistration::label)
.unwrap_or_else(|| ":memory:".to_string()),
),
}
}
pub(crate) fn resolve_reader_checkout<'p>(
&self,
capability: StorageCapability,
operation: &'static str,
outcome: Result<Option<ReaderGuard<'p>>, SqliteError>,
) -> Result<ReaderGuard<'p>, StorageError> {
match outcome {
Ok(Some(mut guard)) => {
guard.label_operation(operation);
Ok(guard)
}
Ok(None) => Err(self.reader_checkout_refusal(capability, operation, None)),
Err(error) => Err(self.reader_checkout_refusal(capability, operation, Some(error))),
}
}
fn reader_checkout_refusal(
&self,
capability: StorageCapability,
operation: &'static str,
error: Option<SqliteError>,
) -> StorageError {
let Some(error) = error else {
return StorageError::Timeout {
operation: operation.into(),
};
};
let is_pool_exhausted = matches!(
&error,
SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
if code.code == rusqlite::ErrorCode::DatabaseBusy
);
if is_pool_exhausted {
self.reader_admission_timeout(operation)
} else {
StorageError::driver(capability, operation, error)
}
}
pub(crate) fn sql_bridge_reader_slots(&self) -> Arc<Semaphore> {
Arc::clone(&self.sql_bridge_reader_slots)
}
pub(crate) fn sql_bridge_writer_slots(&self) -> Arc<Semaphore> {
Arc::clone(&self.sql_bridge_writer_slots)
}
pub fn retirement_writer_holds(&self) -> usize {
usize::from(self.writer.is_locked())
+ usize::from(self.sql_bridge_writer_slots.available_permits() == 0)
}
pub fn origin(&self) -> TxOrigin {
self.origin.clone()
}
pub fn canonical_path(&self) -> Option<&Path> {
self.identity_path.as_deref()
}
#[cfg(unix)]
pub fn opened_file_identity(&self) -> Option<(u64, u64)> {
self.opened_file_identity
.map(DatabaseFileIdentity::unix_parts)
}
#[cfg(any(unix, windows))]
pub fn opened_file_identity_record(&self) -> Option<DatabaseFileIdentity> {
self.opened_file_identity
}
pub fn database_owner_identity(
&self,
) -> Result<DatabaseOwnerIdentity, DatabaseOwnerIdentityError> {
if self.identity_path.is_none() {
return Err(DatabaseOwnerIdentityError::InMemory);
}
#[cfg(any(unix, windows))]
{
let durable_id = self
.opened_database_id
.ok_or(DatabaseOwnerIdentityError::DurableIdentityUnavailable)?;
let file_identity = self
.opened_file_identity
.ok_or(DatabaseOwnerIdentityError::PhysicalIdentityUnavailable)?;
Ok(DatabaseOwnerIdentity {
durable_id,
file_identity,
})
}
#[cfg(not(any(unix, windows)))]
{
Err(DatabaseOwnerIdentityError::UnsupportedPlatform)
}
}
pub fn verify_database_owner(
&self,
expected: &DatabaseOwnerIdentity,
) -> Result<(), DatabaseOwnerIdentityError> {
self.database_owner_identity()?.verify_owner(expected)
}
pub fn write_queue_active(&self) -> bool {
debug_assert!(
self.config.write_queue_enabled.is_some(),
"write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
before any write_queue_active read"
);
self.config.write_queue_enabled.unwrap_or(false) && self.config.path.is_some()
}
pub fn writer_task_join_was_stored(&self) -> bool {
self.writer_task_join_stored.load(Ordering::SeqCst)
}
pub fn writer_task_handle(&self) -> Result<Option<WriterTaskHandle>, StorageError> {
debug_assert!(
self.config.write_queue_enabled.is_some(),
"write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
before any writer_task_handle read"
);
if !self.config.write_queue_enabled.unwrap_or(false) {
return Ok(None);
}
if let Some(existing) = self.writer_task.get() {
return Ok(existing.clone());
}
if tokio::runtime::Handle::try_current().is_err() {
return Err(StorageError::WriterTaskNoRuntime);
}
Ok(self
.writer_task
.get_or_init(|| {
#[cfg(test)]
self.writer_task_spawn_count
.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
match crate::writer_task::spawn(self, self.config.write_queue_capacity) {
Ok(handle) => Some(handle),
Err(e) => {
tracing::warn!(
error = %e,
"KHIVE_WRITE_QUEUE=1 but the writer task failed to spawn; \
writes fall back to the pool-mutex path"
);
None
}
}
})
.clone())
}
pub(crate) fn writer_task_for_write(
&self,
cached: Option<&WriterTaskHandle>,
operation: &'static str,
) -> Result<Option<WriterTaskHandle>, StorageError> {
let handle = match cached {
Some(handle) => Some(handle.clone()),
None => match self.writer_task_handle() {
Ok(handle) => handle,
Err(error) if self.config.write_routing_strict => return Err(error),
Err(_) => None,
},
};
if handle.is_none() && self.config.write_routing_strict {
return Err(StorageError::Pool {
operation: operation.into(),
message: "strict write routing requires a writer-task handle; no handle is \
available, so the direct writer fallback was refused"
.into(),
});
}
Ok(handle)
}
pub(crate) fn record_direct_route(&self, site: crate::timeout_sink::Site) {
if self.write_queue_active() {
crate::timeout_sink::emit_direct_route_violation(
&crate::timeout_sink::db_label(self),
site,
);
}
}
pub fn writer_task_for_runtime_write(
&self,
operation: RuntimeWriteOperation,
) -> Result<Option<WriterTaskHandle>, StorageError> {
let handle = self.writer_task_for_write(None, operation.operation())?;
if handle.is_none() {
self.record_direct_route(operation.fallback_site());
}
Ok(handle)
}
#[cfg(test)]
pub(crate) fn writer_task_spawn_count(&self) -> usize {
self.writer_task_spawn_count
.load(std::sync::atomic::Ordering::SeqCst)
}
pub(crate) fn set_writer_task_join(&self, join: tokio::task::JoinHandle<()>) {
let first_store = !self.writer_task_join_stored.swap(true, Ordering::SeqCst);
debug_assert!(
first_store,
"writer task JoinHandle stored twice (even counting a taken one); \
the writer_task OnceLock is supposed to make spawn at-most-once per pool"
);
if first_store {
*self.writer_task_join.lock() = Some(join);
}
}
pub fn take_writer_task_join(&self) -> Option<tokio::task::JoinHandle<()>> {
self.writer_task_join.lock().take()
}
fn open_reader_connection(&self) -> Result<Connection, SqliteError> {
let path = self.read_connection_path()?;
#[cfg(any(unix, windows))]
if let Some(identity_path) = self.identity_path.as_deref() {
self.verify_opened_file_identity(identity_path)?;
}
let conn = open_reader_connection(path, &self.config)?;
#[cfg(any(unix, windows))]
if let Some(identity_path) = self.identity_path.as_deref() {
self.verify_connection_file_identity(&conn, identity_path)?;
}
self.verify_opened_database_id(&conn)?;
#[cfg(any(test, feature = "test-support"))]
crate::statement_observer::install(&conn, &self.statement_observer)?;
Ok(conn)
}
fn read_connection_path(&self) -> Result<&Path, SqliteError> {
self.read_only_open_target
.as_deref()
.or(self.identity_path.as_deref())
.ok_or_else(|| {
SqliteError::InvalidData(
"in-memory databases do not support standalone connections".to_string(),
)
})
}
pub fn open_standalone_writer(&self) -> Result<Connection, SqliteError> {
self.write_admission.check()?;
self.open_standalone_writer_for_admitted_operation()
}
pub(crate) fn open_standalone_writer_for_admitted_operation(
&self,
) -> Result<Connection, SqliteError> {
let conn = self.open_standalone_writer_untracked()?;
self.writer_acquisition_counters
.standalone_acquisitions
.fetch_add(1, Ordering::Relaxed);
Ok(conn)
}
pub(crate) fn open_standalone_writer_untracked(&self) -> Result<Connection, SqliteError> {
let path = self.identity_path.as_deref().ok_or_else(|| {
SqliteError::InvalidData(
"in-memory databases do not support standalone connections".to_string(),
)
})?;
if self.config.read_only {
return Err(SqliteError::InvalidData(
"database is read-only: standalone write connections are not permitted".to_string(),
));
}
#[cfg(any(unix, windows))]
self.verify_opened_file_identity(path)?;
#[cfg(test)]
run_identity_open_hook(path, IdentityOpenStage::BeforeStandaloneOpen, None);
let conn = claimed_file_identity::open_connection(
&self.config,
path,
OpenFlags::SQLITE_OPEN_READ_WRITE
| OpenFlags::SQLITE_OPEN_NO_MUTEX
| OpenFlags::SQLITE_OPEN_URI,
self.identity_path.as_deref(),
)?;
#[cfg(test)]
run_identity_open_hook(path, IdentityOpenStage::AfterStandaloneOpen, Some(&conn));
#[cfg(any(unix, windows))]
self.verify_connection_file_identity(&conn, path)?;
self.verify_opened_database_id(&conn)?;
#[cfg(feature = "namespace-trigram-proto")]
register_namespace_trigram(&conn)?;
register_writer_clock(&conn)?;
register_rfc3339_key(&conn)?;
conn.busy_timeout(self.config.busy_timeout)?;
if self.config.code_map_vfs.is_none() {
self.checkpoint_ownership
.configure_wal_autocheckpoint(&conn)?;
}
conn.pragma_update(None, "foreign_keys", "ON")?;
conn.pragma_update(None, "synchronous", "NORMAL")?;
let wal_enabled =
self.config.wal_mode && current_journal_mode(&conn)?.eq_ignore_ascii_case("wal");
if wal_enabled {
conn.pragma_update(
None,
"journal_size_limit",
self.config.journal_size_limit_bytes,
)?;
}
#[cfg(any(test, feature = "test-support"))]
crate::statement_observer::install(&conn, &self.statement_observer)?;
Ok(conn)
}
#[cfg(any(unix, windows))]
fn verify_opened_file_identity(&self, path: &Path) -> Result<(), SqliteError> {
let Some(expected) = self.opened_file_identity else {
return Err(SqliteError::InvalidData(
"file-backed pool has no opened database file identity".to_string(),
));
};
let current = database_file_identity(path).ok();
if current != Some(expected) {
return Err(SqliteError::InvalidData(
"pool database file identity changed since the first open; refusing standalone connection"
.to_string(),
));
}
Ok(())
}
#[cfg(any(unix, windows))]
fn verify_connection_file_identity(
&self,
conn: &Connection,
path: &Path,
) -> Result<(), SqliteError> {
let opened = opened_sqlite_file_identity(conn, path)?;
if self.opened_file_identity != Some(opened) {
return Err(SqliteError::InvalidData(
"pool database file identity changed since the first open; refusing standalone connection"
.to_string(),
));
}
Ok(())
}
fn verify_opened_database_id(&self, conn: &Connection) -> Result<(), SqliteError> {
let Some(expected) = self.opened_database_id else {
return Ok(());
};
if self.identity_path.is_some() && read_database_id(conn)? != Some(expected) {
return Err(SqliteError::InvalidData(
"pool database identity changed since the first open; refusing standalone connection"
.to_string(),
));
}
Ok(())
}
#[cfg(test)]
pub(crate) fn effective_wal_autocheckpoint_pages(&self) -> u32 {
self.checkpoint_ownership.wal_autocheckpoint_pages()
}
pub fn claim_checkpoint_ownership(&self) -> Result<(), SqliteError> {
if !self.checkpoint_ownership.begin_claim() {
return Ok(());
}
let result = (|| {
if !self.config.read_only {
let writer = self.writer_for_checkpoint_operation()?;
writer.conn().pragma_update(None, "wal_autocheckpoint", 0)?;
}
Ok(())
})();
self.checkpoint_ownership.finish_claim(result.is_ok());
result
}
pub async fn propagate_checkpoint_claim_to_writer_task(&self) -> Result<(), StorageError> {
let Some(handle) = self.writer_task_handle()? else {
return Ok(());
};
handle
.send_top_level(|conn| {
conn.pragma_update(None, "wal_autocheckpoint", 0)
.map_err(|e| StorageError::Pool {
operation: "claim_checkpoint_ownership".into(),
message: e.to_string(),
})
})
.await
}
pub(crate) fn open_standalone_reader(
&self,
purpose: StandaloneReaderPurpose,
) -> Result<Connection, SqliteError> {
let path = self.read_connection_path()?;
#[cfg(any(unix, windows))]
if let Some(identity_path) = self.identity_path.as_deref() {
self.verify_opened_file_identity(identity_path)?;
}
let conn = claimed_file_identity::open_connection(
&self.config,
path,
OpenFlags::SQLITE_OPEN_READ_ONLY
| OpenFlags::SQLITE_OPEN_NO_MUTEX
| OpenFlags::SQLITE_OPEN_URI,
self.identity_path.as_deref(),
)?;
#[cfg(any(unix, windows))]
if let Some(identity_path) = self.identity_path.as_deref() {
self.verify_connection_file_identity(&conn, identity_path)?;
}
self.verify_opened_database_id(&conn)?;
configure_reader_connection(&conn, &self.config)?;
conn.pragma_update(None, "synchronous", "NORMAL")?;
self.reader_acquisition_counters
.record_standalone_open(purpose);
#[cfg(any(test, feature = "test-support"))]
crate::statement_observer::install(&conn, &self.statement_observer)?;
Ok(conn)
}
fn return_reader(&self, conn: Connection, dirty: bool) {
if self.max_readers == 0 {
return;
}
if reset_reader_connection(&conn, dirty, &self.config)
&& reader_connection_is_healthy(&conn)
{
self.enqueue_reader_slot(conn);
return;
}
close_connection_quietly(conn);
self.replace_discarded_reader_slot();
}
fn enqueue_reader_slot(&self, conn: Connection) {
if let Err(conn) = self.readers.push(conn) {
eprintln!("[sqlite-pool] reader pool queue full, discarding replacement connection");
close_connection_quietly(conn);
}
}
fn replace_discarded_reader_slot(&self) {
match self.open_reader_connection() {
Ok(conn) => self.enqueue_reader_slot(conn),
Err(error) => {
self.reader_acquisition_counters
.record_reader_replacement_open_failure();
tracing::warn!(
%error,
"sqlite-pool: reader replacement connection failed to open; the physical \
pool permanently shrinks by one slot below max_readers"
);
}
}
}
}
const MAX_SYMLINK_DEPTH: u32 = 40;
fn mint_db_identity(configured_path: &Path) -> Result<(DbIdentity, PathBuf), SqliteError> {
let absolute = if configured_path.is_absolute() {
configured_path.to_path_buf()
} else {
let cwd = std::env::current_dir().map_err(|e| {
SqliteError::InvalidData(format!(
"cannot mint database identity for {configured_path:?}: failed to resolve the \
process current directory: {e}"
))
})?;
cwd.join(configured_path)
};
if absolute.exists() {
let canonical = absolute.canonicalize().map_err(|e| {
SqliteError::InvalidData(format!(
"cannot mint database identity: failed to canonicalize existing path \
{absolute:?}: {e}"
))
})?;
return Ok((
DbIdentity::new(canonical.clone().into_os_string()),
canonical,
));
}
let resolved_target = resolve_symlink_chain(&absolute)?;
let parent = resolved_target.parent().ok_or_else(|| {
SqliteError::InvalidData(format!(
"cannot mint database identity for {resolved_target:?}: path has no parent \
directory"
))
})?;
let file_name = resolved_target.file_name().ok_or_else(|| {
SqliteError::InvalidData(format!(
"cannot mint database identity for {resolved_target:?}: path has no file name"
))
})?;
let canonical_parent = parent.canonicalize().map_err(|e| {
SqliteError::InvalidData(format!(
"cannot mint database identity: parent directory {parent:?} of first-open path \
{resolved_target:?} does not exist or is inaccessible: {e}"
))
})?;
let mut identity_path = canonical_parent;
identity_path.push(file_name);
Ok((
DbIdentity::new(identity_path.clone().into_os_string()),
identity_path,
))
}
#[cfg(any(unix, windows))]
fn database_file_identity_if_exists(
path: &Path,
) -> Result<Option<DatabaseFileIdentity>, SqliteError> {
match database_file_identity(path) {
Ok(identity) => Ok(Some(identity)),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(error) => Err(error.into()),
}
}
#[cfg(any(unix, windows))]
pub(crate) fn opened_sqlite_file_identity(
conn: &Connection,
path: &Path,
) -> Result<DatabaseFileIdentity, SqliteError> {
#[cfg(unix)]
{
verify_sqlite_opened_file_still_at_path(conn)?;
Ok(database_file_identity(path)?)
}
#[cfg(windows)]
{
let opened = sqlite_opened_file_identity(conn)?;
if database_file_identity(path).ok() != Some(opened) {
return Err(SqliteError::InvalidData(
"database file identity changed while SQLite held the opened file".to_string(),
));
}
Ok(opened)
}
}
fn read_database_id(conn: &Connection) -> Result<Option<uuid::Uuid>, SqliteError> {
let table_exists: bool = conn.query_row(
"SELECT count(*) != 0 FROM main.sqlite_master WHERE type = 'table' AND name = ?1",
[DATABASE_ID_TABLE],
|row| row.get(0),
)?;
if !table_exists {
return Ok(None);
}
let id: String = conn.query_row(
&format!("SELECT id FROM main.{DATABASE_ID_TABLE} WHERE singleton = 1"),
[],
|row| row.get(0),
)?;
let id = uuid::Uuid::parse_str(&id).map_err(|error| {
SqliteError::InvalidData(format!("invalid stored database identity: {error}"))
})?;
Ok(Some(id))
}
#[cfg(test)]
fn initialize_database_id(conn: &mut Connection) -> Result<uuid::Uuid, SqliteError> {
initialize_database_id_with_admission(conn, None)
}
fn initialize_database_id_with_admission(
conn: &mut Connection,
admission: Option<&WriteAdmission>,
) -> Result<uuid::Uuid, SqliteError> {
if let Some(id) = read_database_id(conn)? {
return Ok(id);
}
let transaction = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
if let Some(admission) = admission {
if let Err(error) = admission.check() {
let rollback = transaction.rollback();
return Err(crate::migrations::capacity_refusal_after_rollback(
conn,
rollback,
error,
"database identity bootstrap",
));
}
}
transaction.execute_batch(&format!(
"CREATE TABLE IF NOT EXISTS main.{DATABASE_ID_TABLE} (\
singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
id TEXT NOT NULL\
)"
))?;
transaction.execute(
&format!("INSERT OR IGNORE INTO main.{DATABASE_ID_TABLE} (singleton, id) VALUES (1, ?1)"),
[uuid::Uuid::new_v4().to_string()],
)?;
let id: String = transaction.query_row(
&format!("SELECT id FROM main.{DATABASE_ID_TABLE} WHERE singleton = 1"),
[],
|row| row.get(0),
)?;
let id = uuid::Uuid::parse_str(&id).map_err(|error| {
SqliteError::InvalidData(format!("invalid stored database identity: {error}"))
})?;
transaction.commit()?;
Ok(id)
}
#[cfg(unix)]
fn verify_sqlite_opened_file_still_at_path(conn: &Connection) -> Result<(), SqliteError> {
let mut moved: std::ffi::c_int = 0;
let result = unsafe {
rusqlite::ffi::sqlite3_file_control(
conn.handle(),
c"main".as_ptr(),
rusqlite::ffi::SQLITE_FCNTL_HAS_MOVED,
(&mut moved as *mut std::ffi::c_int).cast(),
)
};
if result != rusqlite::ffi::SQLITE_OK {
return Err(SqliteError::InvalidData(format!(
"cannot verify opened database file identity (SQLite file control {result})"
)));
}
if moved != 0 {
return Err(SqliteError::InvalidData(
"database file identity changed while SQLite held the opened file".to_string(),
));
}
Ok(())
}
fn resolve_symlink_chain(path: &Path) -> Result<PathBuf, SqliteError> {
let mut current = path.to_path_buf();
for _ in 0..MAX_SYMLINK_DEPTH {
match fs::symlink_metadata(¤t) {
Ok(meta) if meta.file_type().is_symlink() => {
let target = fs::read_link(¤t).map_err(|e| {
SqliteError::InvalidData(format!(
"cannot mint database identity: failed to read symlink {current:?}: {e}"
))
})?;
current = if target.is_absolute() {
target
} else {
match current.parent() {
Some(parent) => parent.join(&target),
None => target,
}
};
}
_ => return Ok(current),
}
}
Err(SqliteError::InvalidData(format!(
"cannot mint database identity for {path:?}: symlink chain exceeds \
{MAX_SYMLINK_DEPTH} levels"
)))
}
fn effective_reader_count(config: &PoolConfig, wal_enabled: bool) -> usize {
if config.path.is_some() && (config.read_only || config.code_map_vfs.is_some()) {
config.max_readers.max(1)
} else if config.path.is_some() && config.wal_mode && wal_enabled {
config.max_readers
} else {
0
}
}
fn open_writer_connection(
config: &PoolConfig,
read_only_open_target: Option<&Path>,
identity_path: Option<&Path>,
) -> Result<Connection, SqliteError> {
claimed_file_identity::open_writer(config, read_only_open_target, identity_path)
}
fn validate_wal_ceiling_at_open(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
let bytes = config.wal_ceiling.effective_bytes(config.read_only);
if bytes == 0 {
return Ok(());
}
let page_size: i64 = conn.pragma_query_value(None, "page_size", |row| row.get(0))?;
let page_size = u64::try_from(page_size).map_err(|_| {
SqliteError::InvalidData("SQLite reported a negative page size".to_string())
})?;
let minimum_bytes = page_size.checked_add(56).ok_or_else(|| {
SqliteError::InvalidData("SQLite page size overflowed the WAL frame floor".to_string())
})?;
if bytes < minimum_bytes {
return Err(SqliteError::WalCeilingBelowMinimum {
bytes,
page_size,
minimum_bytes,
});
}
Err(SqliteError::WalCapacityUnavailable {
bytes,
capability: "WAL I/O limiter",
})
}
fn read_only_open_target(
config: &PoolConfig,
physical_path: Option<&Path>,
) -> Result<Option<PathBuf>, SqliteError> {
if !config.read_only || config.code_map_vfs.is_some() {
return Ok(None);
}
let Some(path) = physical_path else {
return Ok(None);
};
read_only_wal_open_target_for_path(path).map(Some)
}
fn read_only_wal_open_target_for_path(path: &Path) -> Result<PathBuf, SqliteError> {
if !sqlite_header_uses_wal(path)? {
return Ok(path.to_path_buf());
}
let shm = sqlite_sidecar_path(path, "-shm");
match fs::metadata(&shm) {
Ok(metadata) if metadata.permissions().readonly() => {
let wal = sqlite_sidecar_path(path, "-wal");
match fs::metadata(&wal) {
Ok(_) => Ok(path.to_path_buf()),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
Err(SqliteError::InvalidData(format!(
"read-only WAL snapshot {} has a shared-memory sidecar {} but no WAL \
sidecar {}; refusing the inconsistent sidecar set before SQLite open",
path.display(),
shm.display(),
wal.display(),
)))
}
Err(error) => Err(SqliteError::Io(error)),
}
}
Ok(_) => Err(SqliteError::InvalidData(format!(
"read-only WAL snapshot {} has a writable WAL shared-memory sidecar {}; close every \
live writer and remove the transient -shm file (or make a genuinely frozen snapshot) \
before inspection",
path.display(),
shm.display(),
))),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
let wal = sqlite_sidecar_path(path, "-wal");
match fs::metadata(&wal) {
Ok(metadata) if metadata.len() > 0 => Err(SqliteError::InvalidData(format!(
"read-only WAL snapshot {} has a non-empty WAL sidecar {} but no read-only \
shared-memory sidecar {}; refusing before SQLite open because immutable \
mode would omit committed WAL frames and ordinary read-only mode would \
create or mutate -shm; include the frozen read-only -shm beside this \
snapshot, or checkpoint a writable copy before inspection",
path.display(),
wal.display(),
shm.display(),
))),
Ok(_) => sqlite_immutable_uri(path),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
sqlite_immutable_uri(path)
}
Err(error) => Err(SqliteError::Io(error)),
}
}
Err(error) => Err(SqliteError::Io(error)),
}
}
pub(crate) fn open_read_only_snapshot_connection(path: &Path) -> Result<Connection, SqliteError> {
let (_, physical_path) = mint_db_identity(path)?;
let target = read_only_wal_open_target_for_path(&physical_path)?;
let conn = Connection::open_with_flags(&target, reader_open_flags())?;
#[cfg(feature = "namespace-trigram-proto")]
register_namespace_trigram(&conn)?;
Ok(conn)
}
fn sqlite_header_uses_wal(path: &Path) -> Result<bool, SqliteError> {
let mut file = fs::File::open(path)?;
let mut header = [0_u8; 20];
if let Err(error) = file.read_exact(&mut header) {
if error.kind() == std::io::ErrorKind::UnexpectedEof {
return Ok(false);
}
return Err(SqliteError::Io(error));
}
Ok(&header[..16] == b"SQLite format 3\0" && header[18] == 2 && header[19] == 2)
}
fn sqlite_sidecar_path(path: &Path, suffix: &str) -> PathBuf {
let mut sidecar = path.as_os_str().to_os_string();
sidecar.push(suffix);
PathBuf::from(sidecar)
}
fn sqlite_immutable_uri(path: &Path) -> Result<PathBuf, SqliteError> {
let absolute = if path.is_absolute() {
path.to_path_buf()
} else {
std::env::current_dir()?.join(path)
};
let mut uri = String::from("file:");
#[cfg(unix)]
{
use std::os::unix::ffi::OsStrExt as _;
push_sqlite_uri_path(&mut uri, absolute.as_os_str().as_bytes());
}
#[cfg(not(unix))]
{
let path = absolute.to_str().ok_or_else(|| {
SqliteError::InvalidData(format!(
"read-only WAL snapshot path is not representable as a SQLite URI: {}",
absolute.display()
))
})?;
let normalized = path.replace('\\', "/");
if cfg!(windows) && !normalized.starts_with('/') {
uri.push('/');
}
push_sqlite_uri_path(&mut uri, normalized.as_bytes());
}
uri.push_str("?mode=ro&immutable=1");
Ok(PathBuf::from(uri))
}
fn push_sqlite_uri_path(uri: &mut String, bytes: &[u8]) {
const HEX: &[u8; 16] = b"0123456789ABCDEF";
for &byte in bytes {
if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~' | b'/') {
uri.push(byte as char);
} else {
uri.push('%');
uri.push(HEX[(byte >> 4) as usize] as char);
uri.push(HEX[(byte & 0x0f) as usize] as char);
}
}
}
fn open_reader_connection(path: &Path, config: &PoolConfig) -> Result<Connection, SqliteError> {
let conn = claimed_file_identity::open_connection(
config,
path,
reader_open_flags(),
config.path.as_deref(),
)?;
configure_reader_connection(&conn, config)?;
Ok(conn)
}
fn writer_open_flags() -> OpenFlags {
OpenFlags::SQLITE_OPEN_READ_WRITE
| OpenFlags::SQLITE_OPEN_CREATE
| OpenFlags::SQLITE_OPEN_URI
| OpenFlags::SQLITE_OPEN_NO_MUTEX
}
fn writer_read_only_open_flags() -> OpenFlags {
OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
}
fn reader_open_flags() -> OpenFlags {
OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
}
#[cfg(feature = "namespace-trigram-proto")]
fn register_namespace_trigram(conn: &Connection) -> Result<(), SqliteError> {
crate::namespace_trigram_proto::register(conn).map_err(SqliteError::InvalidData)
}
fn register_writer_clock(conn: &Connection) -> Result<(), SqliteError> {
conn.create_scalar_function(
"khive_now_micros",
0,
rusqlite::functions::FunctionFlags::SQLITE_UTF8,
|_| Ok(chrono::Utc::now().timestamp_micros()),
)?;
Ok(())
}
pub(crate) fn rfc3339_instant_key(instant: chrono::DateTime<chrono::Utc>) -> Vec<u8> {
let mut key = Vec::with_capacity(12);
key.extend_from_slice(&((instant.timestamp() as u64) ^ (1_u64 << 63)).to_be_bytes());
key.extend_from_slice(&instant.timestamp_subsec_nanos().to_be_bytes());
key
}
pub(crate) fn strict_rfc3339_key(text: &str) -> Option<Vec<u8>> {
chrono::DateTime::parse_from_rfc3339(text)
.ok()
.map(|instant| rfc3339_instant_key(instant.with_timezone(&chrono::Utc)))
}
pub(crate) fn register_rfc3339_key(conn: &Connection) -> rusqlite::Result<()> {
use rusqlite::functions::FunctionFlags;
use rusqlite::types::ValueRef;
conn.create_scalar_function(
"khive_rfc3339_key",
1,
FunctionFlags::SQLITE_UTF8
| FunctionFlags::SQLITE_DETERMINISTIC
| FunctionFlags::SQLITE_INNOCUOUS,
|ctx| {
let text = match ctx.get_raw(0) {
ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
_ => None,
};
let key = text
.and_then(|text| text.parse::<chrono::DateTime<chrono::Utc>>().ok())
.map(rfc3339_instant_key);
Ok(key)
},
)?;
conn.create_scalar_function(
"khive_rfc3339_strict_key",
1,
FunctionFlags::SQLITE_UTF8
| FunctionFlags::SQLITE_DETERMINISTIC
| FunctionFlags::SQLITE_INNOCUOUS,
|ctx| {
let text = match ctx.get_raw(0) {
ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
_ => None,
};
let key = text.and_then(strict_rfc3339_key);
Ok(key)
},
)?;
Ok(())
}
fn configure_writer_connection(
conn: &Connection,
config: &PoolConfig,
) -> Result<bool, SqliteError> {
#[cfg(feature = "namespace-trigram-proto")]
register_namespace_trigram(conn)?;
register_writer_clock(conn)?;
register_rfc3339_key(conn)?;
if config.read_only {
conn.pragma_update(None, "foreign_keys", "ON")?;
conn.busy_timeout(config.busy_timeout)?;
conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
conn.pragma_update(None, "temp_store", "MEMORY")?;
conn.pragma_update(None, "query_only", "ON")?;
let wal_enabled =
config.wal_mode && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
return Ok(wal_enabled);
}
let wants_wal = config.path.is_some() && config.wal_mode;
if wants_wal {
conn.pragma_update(None, "journal_mode", "WAL")?;
}
code_map::require_delete_journal(conn, config)?;
conn.pragma_update(None, "synchronous", "NORMAL")?;
conn.pragma_update(None, "foreign_keys", "ON")?;
conn.busy_timeout(config.busy_timeout)?;
conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
conn.pragma_update(None, "temp_store", "MEMORY")?;
if config.code_map_vfs.is_none() {
conn.pragma_update(
None,
"wal_autocheckpoint",
FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
)?;
}
let wal_enabled = wants_wal && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
if wal_enabled {
conn.pragma_update(None, "journal_size_limit", config.journal_size_limit_bytes)?;
}
Ok(wal_enabled)
}
fn configure_reader_connection(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
#[cfg(feature = "namespace-trigram-proto")]
register_namespace_trigram(conn)?;
register_rfc3339_key(conn)?;
conn.pragma_update(None, "foreign_keys", "ON")?;
conn.busy_timeout(config.busy_timeout)?;
conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
conn.pragma_update(None, "temp_store", "MEMORY")?;
Ok(())
}
fn current_journal_mode(conn: &Connection) -> Result<String, SqliteError> {
conn.pragma_query_value(None, "journal_mode", |row| row.get::<_, String>(0))
.map(|mode| mode.to_ascii_lowercase())
.map_err(Into::into)
}
fn reset_reader_connection(conn: &Connection, dirty: bool, config: &PoolConfig) -> bool {
if !conn.is_autocommit() {
match conn.execute_batch("ROLLBACK") {
Ok(()) => {}
Err(rusqlite::Error::SqliteFailure(err, _)) => {
if matches!(
err.code,
rusqlite::ErrorCode::CannotOpen
| rusqlite::ErrorCode::DatabaseCorrupt
| rusqlite::ErrorCode::NotADatabase
| rusqlite::ErrorCode::DiskFull
) {
return false;
}
}
Err(_) => return false,
}
if !conn.is_autocommit() {
return false;
}
}
if !dirty {
return true;
}
reader_connection_state_is_pristine(conn)
&& reader_connection_settings_match_baseline(conn, config, 0)
}
fn reader_connection_state_is_pristine(conn: &Connection) -> bool {
let has_temp_objects: bool = match conn.query_row(
"SELECT EXISTS(SELECT 1 FROM sqlite_temp_master)",
[],
|row| row.get(0),
) {
Ok(v) => v,
Err(_) => return false,
};
if has_temp_objects {
return false;
}
let attached_databases: i64 = match conn.query_row(
"SELECT COUNT(*) FROM pragma_database_list WHERE name NOT IN ('main', 'temp')",
[],
|row| row.get(0),
) {
Ok(v) => v,
Err(_) => return false,
};
attached_databases == 0
}
fn reader_connection_settings_match_baseline(
conn: &Connection,
config: &PoolConfig,
expected_query_only: i64,
) -> bool {
let expected_busy_timeout_ms =
i64::try_from(config.busy_timeout.as_millis()).unwrap_or(i64::MAX);
let expected_cache_size: i64 = CACHE_SIZE_KIB.parse().unwrap_or(-65536);
let checks: [(&str, i64); 8] = [
("query_only", expected_query_only),
("writable_schema", 0),
("foreign_keys", 1),
("busy_timeout", expected_busy_timeout_ms),
("cache_size", expected_cache_size),
("temp_store", 2),
("read_uncommitted", 0),
("defer_foreign_keys", 0),
];
checks.iter().all(|(pragma, expected)| {
conn.pragma_query_value(None, pragma, |row| row.get::<_, i64>(0))
.map(|actual| actual == *expected)
.unwrap_or(false)
})
}
fn restore_shared_reader_state(conn: &Connection, config: &PoolConfig) -> bool {
if !detach_non_main_databases(conn) {
return false;
}
if !drop_temp_objects(conn) {
return false;
}
let expected_query_only = i64::from(config.read_only);
if reset_observable_settings(conn, config, expected_query_only).is_err() {
return false;
}
reader_connection_state_is_pristine(conn)
&& reader_connection_settings_match_baseline(conn, config, expected_query_only)
}
fn detach_non_main_databases(conn: &Connection) -> bool {
loop {
let name: Option<String> = match conn.query_row(
"SELECT name FROM pragma_database_list WHERE name NOT IN ('main', 'temp') LIMIT 1",
[],
|row| row.get(0),
) {
Ok(name) => Some(name),
Err(rusqlite::Error::QueryReturnedNoRows) => None,
Err(_) => return false,
};
let Some(name) = name else {
return true;
};
let quoted = format!("\"{}\"", name.replace('"', "\"\""));
if conn
.execute_batch(&format!("DETACH DATABASE {quoted}"))
.is_err()
{
return false;
}
}
}
fn drop_temp_objects(conn: &Connection) -> bool {
for (kind, ddl_keyword) in [
("view", "VIEW"),
("trigger", "TRIGGER"),
("index", "INDEX"),
("table", "TABLE"),
] {
loop {
let name: Option<String> = match conn.query_row(
"SELECT name FROM sqlite_temp_master WHERE type = ?1 LIMIT 1",
[kind],
|row| row.get(0),
) {
Ok(name) => Some(name),
Err(rusqlite::Error::QueryReturnedNoRows) => None,
Err(_) => return false,
};
let Some(name) = name else {
break;
};
let quoted = format!("\"{}\"", name.replace('"', "\"\""));
if conn
.execute_batch(&format!("DROP {ddl_keyword} IF EXISTS temp.{quoted}"))
.is_err()
{
return false;
}
}
}
true
}
fn reset_observable_settings(
conn: &Connection,
config: &PoolConfig,
expected_query_only: i64,
) -> Result<(), rusqlite::Error> {
conn.pragma_update(None, "query_only", expected_query_only)?;
conn.pragma_update(None, "writable_schema", 0)?;
conn.pragma_update(None, "foreign_keys", "ON")?;
conn.busy_timeout(config.busy_timeout)?;
conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
conn.pragma_update(None, "temp_store", "MEMORY")?;
conn.pragma_update(None, "read_uncommitted", 0)?;
conn.pragma_update(None, "defer_foreign_keys", 0)?;
Ok(())
}
fn reader_connection_is_healthy(conn: &Connection) -> bool {
match conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0)) {
Ok(_) => true,
Err(rusqlite::Error::SqliteFailure(err, _)) => !matches!(
err.code,
rusqlite::ErrorCode::CannotOpen
| rusqlite::ErrorCode::NotADatabase
| rusqlite::ErrorCode::DatabaseCorrupt
| rusqlite::ErrorCode::PermissionDenied
| rusqlite::ErrorCode::SystemIoFailure
),
Err(_) => true,
}
}
fn close_connection_quietly(conn: Connection) {
match conn.close() {
Ok(()) => {}
Err((conn, _)) => drop(conn),
}
}
fn pool_exhausted_error(timeout: Duration, max_readers: usize) -> SqliteError {
rusqlite::Error::SqliteFailure(
rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
Some(format!(
"Pool exhausted: no reader available after {timeout:?} (max_readers={max_readers})"
)),
)
.into()
}
#[cfg(test)]
#[path = "runtime_write_routing_tests.rs"]
mod runtime_write_routing_tests;
#[cfg(test)]
#[path = "database_owner_identity_pool_tests.rs"]
mod database_owner_identity_pool_tests;
#[cfg(test)]
#[path = "pool_tests.rs"]
mod tests;
#[cfg(all(test, any(unix, windows)))]
#[path = "pool_identity_admission_tests.rs"]
mod identity_admission_tests;