1mod code_map;
3#[path = "pool/identity_registry.rs"]
4mod identity_registry;
5#[cfg(any(test, feature = "test-support"))]
6#[path = "pool/test_volume_lock_dir.rs"]
7mod test_volume_lock_dir;
8#[path = "pool/write_units.rs"]
9mod write_units;
10#[path = "pool/writer_acquisition.rs"]
11mod writer_acquisition;
12#[cfg(test)]
13use identity_registry::pool_identity_suffix;
14use identity_registry::PoolIdentityRegistration;
15#[cfg(any(test, feature = "test-support"))]
16use test_volume_lock_dir::test_volume_lock_dir;
17#[cfg(test)]
18use write_units::STARTUP_SPACE_PROBE;
19pub use write_units::{CheckpointGuard, CheckpointResult, WriterAcquisitionSnapshot, WriterGuard};
20pub(crate) use write_units::{
21 PooledAutocommitWriteUnit, PooledTransactionWriteUnit, StandaloneTransactionWriteUnit,
22 WriteAdmission, WriterAcquisitionCounters,
23};
24
25use crossbeam_queue::ArrayQueue;
26use parking_lot::{Condvar, Mutex};
27use rusqlite::hooks::{AuthContext, Authorization};
28use rusqlite::{Connection, OpenFlags};
29use serde::Serialize;
30use sha2::{Digest, Sha256};
31use std::cell::Cell;
32use std::collections::{BTreeMap, HashMap};
33use std::fs;
34use std::io::Read as _;
35use std::ops::{Deref, DerefMut};
36use std::path::{Path, PathBuf};
37use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
38use std::sync::{Arc, OnceLock};
39use std::thread;
40use std::time::{Duration, Instant};
41use tokio::sync::Semaphore;
42
43use crate::database_owner_identity::{DatabaseOwnerIdentity, DatabaseOwnerIdentityError};
44use crate::disk_guard::{LeaseRefusal, VolumeIdentity, VolumeLease};
45#[cfg(test)]
46use crate::disk_guard_config::DiskGuardConfigSource;
47use crate::disk_guard_config::{
48 resolve_disk_guard_config, EffectiveDiskGuardConfig, DEFAULT_DISK_GUARD_DEADLINE_MS,
49};
50use crate::error::SqliteError;
51#[cfg(windows)]
52use crate::file_identity::sqlite_opened_file_identity;
53#[cfg(any(unix, windows))]
54use crate::file_identity::{database_file_identity, DatabaseFileIdentity};
55use crate::writer_task::{execute_wrapped_transaction, WriterTaskHandle};
56use khive_storage::error::StorageError;
57use khive_storage::tx_registry::{DbIdentity, TxOrigin};
58use khive_storage::CapacityUnavailablePhase;
59use khive_storage::StorageCapability;
60
61mod claimed_file_identity;
62
63#[cfg(unix)]
64mod claimed_file_observer;
65
66#[cfg(unix)]
67pub use claimed_file_observer::initialize as initialize_claimed_file_observer;
68
69#[cfg(all(test, unix))]
70#[path = "pool/claimed_file_identity_tests.rs"]
71mod claimed_file_identity_tests;
72
73pub(crate) const CACHE_SIZE_KIB: &str = "-65536";
74const MMAP_SIZE_BYTES: &str = "1073741824";
75const DEFAULT_READER_CAP: usize = 8;
76
77const DEFAULT_JOURNAL_SIZE_LIMIT_BYTES: i64 = 67_108_864; const DEFAULT_WRITE_QUEUE_CAPACITY: usize = 256;
79const DATABASE_ID_TABLE: &str = "_khive_database_identity";
80static NEXT_MAIN_POOL_GENERATION: AtomicU64 = AtomicU64::new(1);
81
82#[cfg(test)]
83#[derive(Clone, Copy, PartialEq, Eq)]
84enum IdentityOpenStage {
85 AfterMainOpenBeforeFirstStat,
86 AfterInitialIdentityWrite,
87 BeforeStandaloneOpen,
88 AfterStandaloneOpen,
89}
90
91#[cfg(test)]
92type IdentityOpenHook = Box<dyn Fn(&Path, IdentityOpenStage, Option<&Connection>)>;
93
94#[cfg(test)]
95thread_local! {
96 static IDENTITY_OPEN_HOOK: std::cell::RefCell<Option<IdentityOpenHook>> =
97 const { std::cell::RefCell::new(None) };
98}
99
100#[cfg(test)]
101fn run_identity_open_hook(path: &Path, stage: IdentityOpenStage, conn: Option<&Connection>) {
102 IDENTITY_OPEN_HOOK.with(|hook| {
103 if let Some(hook) = hook.borrow().as_ref() {
104 hook(path, stage, conn);
105 }
106 });
107}
108
109#[path = "pool/admission.rs"]
110mod admission;
111use admission::CheckpointOwnershipGate;
112pub use admission::RuntimeWriteOperation;
113pub(crate) use admission::FALLBACK_WAL_AUTOCHECKPOINT_PAGES;
114#[cfg(test)]
115use admission::{CheckpointConnectionConfigPause, CheckpointOwnership};
116
117fn deny_retired_writer(_context: AuthContext<'_>) -> Authorization {
118 Authorization::Deny
119}
120
121pub(crate) const TEST_HARNESS_ENV: &str = "KHIVE_TEST_HARNESS";
122
123#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)]
125#[serde(rename_all = "snake_case")]
126pub enum WalCeilingSource {
127 BackendField,
128 Environment,
129 #[default]
130 Default,
131}
132
133#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
136pub struct WalCeilingPolicy {
137 pub bytes: u64,
138 pub source: WalCeilingSource,
139}
140
141impl WalCeilingPolicy {
142 pub fn validate_static(
146 self,
147 file_backed: bool,
148 wal_mode: bool,
149 read_only: bool,
150 ) -> Result<(), SqliteError> {
151 if self.bytes == 0 {
152 return Ok(());
153 }
154 if i64::try_from(self.bytes).is_err() {
155 return Err(SqliteError::WalCeilingOffsetOverflow { bytes: self.bytes });
156 }
157 if read_only {
158 return Ok(());
159 }
160 if !file_backed {
161 return Err(SqliteError::WalCeilingUnsupported {
162 bytes: self.bytes,
163 backend_kind: "in-memory backend",
164 });
165 }
166 if !wal_mode {
167 return Err(SqliteError::WalCeilingUnsupported {
168 bytes: self.bytes,
169 backend_kind: "non-WAL backend",
170 });
171 }
172 Ok(())
173 }
174
175 pub fn effective_bytes(self, read_only: bool) -> u64 {
177 if read_only {
178 0
179 } else {
180 self.bytes
181 }
182 }
183}
184
185#[derive(Clone, Debug)]
187pub struct PoolConfig {
188 pub path: Option<PathBuf>,
190 pub code_map_vfs: Option<String>,
192 #[cfg(any(unix, windows))]
195 pub expected_file_identity: Option<DatabaseFileIdentity>,
196 pub max_readers: usize,
198 pub wal_mode: bool,
200 pub busy_timeout: Duration,
204 pub checkout_timeout: Duration,
216 pub journal_size_limit_bytes: i64,
222 pub read_only: bool,
230 pub wal_ceiling: WalCeilingPolicy,
233 pub write_queue_enabled: Option<bool>,
261 pub write_queue_capacity: usize,
266 pub write_routing_strict: bool,
277 pub write_admission_deadline_ms: u64,
291 pub disk_guard_config: Option<EffectiveDiskGuardConfig>,
294 pub volume_lock_dir: Option<PathBuf>,
298 pub read_tx_max_age: Duration,
308}
309
310const WRITE_ADMISSION_DEADLINE_MS_RANGE: std::ops::RangeInclusive<u64> = 100..=10_000;
312const DEFAULT_WRITE_ADMISSION_DEADLINE_MS: u64 = 2000;
313
314impl Default for PoolConfig {
315 fn default() -> Self {
316 Self {
317 path: None,
318 code_map_vfs: None,
319 #[cfg(any(unix, windows))]
320 expected_file_identity: None,
321 max_readers: std::thread::available_parallelism()
322 .map(|n| n.get())
323 .unwrap_or(1)
324 .clamp(1, DEFAULT_READER_CAP),
325 wal_mode: true,
326 busy_timeout: Duration::from_secs(crate::env::env_parse_or(
327 "KHIVE_BUSY_TIMEOUT_SECS",
328 30,
329 )),
330 checkout_timeout: Duration::from_secs(crate::env::env_parse_or(
331 "KHIVE_CHECKOUT_TIMEOUT_SECS",
332 5,
333 )),
334 journal_size_limit_bytes: crate::env::env_parse_or(
335 "KHIVE_JOURNAL_SIZE_LIMIT_BYTES",
336 DEFAULT_JOURNAL_SIZE_LIMIT_BYTES,
337 ),
338 read_only: false,
339 wal_ceiling: WalCeilingPolicy::default(),
340 write_queue_enabled: std::env::var_os("KHIVE_WRITE_QUEUE").map(|v| {
345 v.to_str()
346 .is_some_and(|v| v == "1" || v.eq_ignore_ascii_case("true"))
347 }),
348 write_queue_capacity: std::env::var("KHIVE_WRITE_QUEUE_CAPACITY")
349 .ok()
350 .and_then(|v| v.parse::<usize>().ok())
351 .filter(|&n| n > 0)
352 .unwrap_or(DEFAULT_WRITE_QUEUE_CAPACITY),
353 write_routing_strict: std::env::var("KHIVE_WRITE_ROUTING")
354 .map(|v| v.eq_ignore_ascii_case("strict"))
355 .unwrap_or(false),
356 write_admission_deadline_ms: crate::env::env_parse_or(
357 "KHIVE_WRITE_ADMISSION_DEADLINE_MS",
358 DEFAULT_WRITE_ADMISSION_DEADLINE_MS,
359 ),
360 disk_guard_config: None,
361 #[cfg(test)]
362 volume_lock_dir: Some(test_volume_lock_dir()),
363 #[cfg(not(test))]
364 volume_lock_dir: crate::default_volume_lock_dir().ok(),
365 read_tx_max_age: crate::checkpoint::tx_age_thresholds_from_env(
366 Duration::from_secs(30),
367 Duration::from_secs(120),
368 )
369 .1,
370 }
371 }
372}
373
374#[cfg(any(test, feature = "test-support"))]
375impl PoolConfig {
376 pub fn for_test() -> Self {
381 Self {
382 max_readers: 2,
383 volume_lock_dir: Some(test_volume_lock_dir()),
384 ..Self::default()
385 }
386 }
387}
388
389fn refuse_home_data_store_in_tests(config: &PoolConfig) -> Result<(), SqliteError> {
406 if std::env::var(TEST_HARNESS_ENV).as_deref() != Ok("1") {
407 return Ok(());
408 }
409
410 let Some(path) = config.path.as_deref() else {
411 return Ok(());
412 };
413 if path
414 .as_os_str()
415 .as_encoded_bytes()
416 .get(..5)
417 .is_some_and(|prefix| prefix.eq_ignore_ascii_case(b"file:"))
418 {
419 return Err(SqliteError::InvalidData(format!(
420 "test harness refused SQLite URI database path {}; use a filesystem path outside \
421 HOME/.khive (deliberate sessions against a real store run the built binary \
422 directly, outside the Cargo test environment)",
423 path.display()
424 )));
425 }
426
427 let Some(home) = std::env::var_os("HOME") else {
428 return Ok(());
429 };
430 let canonical_path = canonicalize_deepest_existing(path)?;
431 let canonical_home_data_dir =
432 canonicalize_deepest_existing(&PathBuf::from(home).join(".khive"))?;
433 if canonical_path.starts_with(&canonical_home_data_dir) {
434 return Err(SqliteError::InvalidData(format!(
435 "test harness refused to open SQLite database under HOME/.khive: {} \
436 (deliberate sessions against a real store run the built binary directly, \
437 outside the Cargo test environment)",
438 canonical_path.display()
439 )));
440 }
441 Ok(())
442}
443
444fn canonicalize_deepest_existing(path: &Path) -> Result<PathBuf, SqliteError> {
445 let absolute = if path.is_absolute() {
446 path.to_path_buf()
447 } else {
448 std::env::current_dir().map_err(SqliteError::Io)?.join(path)
449 };
450
451 for ancestor in absolute.ancestors() {
452 match fs::canonicalize(ancestor) {
453 Ok(mut canonical) => {
454 let missing = absolute.strip_prefix(ancestor).map_err(|error| {
455 SqliteError::InvalidData(format!(
456 "failed to preserve missing path components for {}: {error}",
457 absolute.display()
458 ))
459 })?;
460 canonical.push(missing);
461 return Ok(canonical);
462 }
463 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
464 Err(error) => {
465 return Err(SqliteError::InvalidData(format!(
466 "failed to canonicalize database path ancestor {}: {error}",
467 ancestor.display()
468 )));
469 }
470 }
471 }
472
473 Err(SqliteError::InvalidData(format!(
474 "database path has no canonicalizable ancestor: {}",
475 absolute.display()
476 )))
477}
478
479fn validate_write_admission_deadline(deadline_ms: u64) -> Result<(), SqliteError> {
484 if WRITE_ADMISSION_DEADLINE_MS_RANGE.contains(&deadline_ms) {
485 return Ok(());
486 }
487 Err(SqliteError::InvalidConfig(format!(
488 "write_admission_deadline_ms must be in [{}, {}] ms, got {deadline_ms}",
489 WRITE_ADMISSION_DEADLINE_MS_RANGE.start(),
490 WRITE_ADMISSION_DEADLINE_MS_RANGE.end()
491 )))
492}
493
494#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize)]
496pub struct SearchMechanismSnapshot {
497 pub dispatches_by_backend_and_kind: BTreeMap<String, BTreeMap<String, u64>>,
501 pub note_candidate_hydration_rows: u64,
505}
506
507pub struct ConnectionPool {
521 writer: Arc<Mutex<Connection>>,
522 retirement_connection: Mutex<Option<Connection>>,
524 #[cfg(any(test, feature = "test-support"))]
525 statement_observer: Arc<crate::statement_observer::StatementObserverHub>,
526 main_pool_generation: OnceLock<u64>,
527 checkpoint_ownership: CheckpointOwnershipGate,
536 pooled_writer_retired: AtomicBool,
541 writer_acquisition_counters: Arc<WriterAcquisitionCounters>,
546 write_admission: Arc<WriteAdmission>,
549 disk_guard_config: Option<EffectiveDiskGuardConfig>,
552 reader_acquisition_counters: ReaderAcquisitionCounters,
557 search_dispatches: Mutex<BTreeMap<String, BTreeMap<String, u64>>>,
560 note_candidate_hydration_rows: AtomicU64,
561 readers: ArrayQueue<Connection>,
562 max_readers: usize,
563 config: PoolConfig,
564 read_only_open_target: Option<PathBuf>,
572 sql_bridge_reader_slots: Arc<Semaphore>,
573 sql_bridge_writer_slots: Arc<Semaphore>,
574 writer_task: OnceLock<Option<WriterTaskHandle>>,
579 writer_task_join: Mutex<Option<tokio::task::JoinHandle<()>>>,
587 writer_task_join_stored: AtomicBool,
592 origin: TxOrigin,
598 identity_path: Option<PathBuf>,
604 #[cfg(any(unix, windows))]
607 opened_file_identity: Option<DatabaseFileIdentity>,
608 opened_database_id: Option<uuid::Uuid>,
611 identity_registration: Option<PoolIdentityRegistration>,
614 #[cfg(test)]
620 writer_task_spawn_count: std::sync::atomic::AtomicUsize,
621}
622
623impl Drop for ConnectionPool {
624 fn drop(&mut self) {
632 while let Some(conn) = self.readers.pop() {
633 drop(conn);
634 }
635 }
636}
637
638enum ReaderLease<'pool> {
639 Pooled(Connection),
640 Shared(parking_lot::MutexGuard<'pool, Connection>),
641}
642
643pub struct ReaderRow<'row, 'statement> {
653 row: &'row rusqlite::Row<'statement>,
654}
655
656impl ReaderRow<'_, '_> {
657 pub fn get<I: rusqlite::RowIndex, T: rusqlite::types::FromSql>(
659 &self,
660 index: I,
661 ) -> rusqlite::Result<T> {
662 self.row.get(index)
663 }
664
665 pub fn get_ref<I: rusqlite::RowIndex>(
667 &self,
668 index: I,
669 ) -> rusqlite::Result<rusqlite::types::ValueRef<'_>> {
670 self.row.get_ref(index)
671 }
672}
673
674struct ReaderQueryInProgress<'a>(&'a Cell<bool>);
676
677impl Drop for ReaderQueryInProgress<'_> {
678 fn drop(&mut self) {
679 self.0.set(false);
680 }
681}
682
683pub(crate) struct ReaderAdmission {
689 slot: tokio::sync::OwnedSemaphorePermit,
690 started: Instant,
691}
692
693pub struct ReaderGuard<'pool> {
696 lease: Option<ReaderLease<'pool>>,
697 admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
701 pool: &'pool ConnectionPool,
702 reusable: Cell<bool>,
703 query_in_progress: Cell<bool>,
704 checked_out_at: Instant,
705 dirty: Cell<bool>,
713 operation: Option<&'static str>,
719}
720
721impl<'pool> ReaderGuard<'pool> {
722 pub(crate) fn conn(&self) -> &Connection {
731 match self
732 .lease
733 .as_ref()
734 .expect("reader guard missing connection")
735 {
736 ReaderLease::Pooled(conn) => conn,
737 ReaderLease::Shared(guard) => guard,
738 }
739 }
740
741 pub fn query_row<T, P, F>(&self, sql: &str, params: P, f: F) -> Result<T, SqliteError>
760 where
761 P: rusqlite::Params,
762 F: FnOnce(&ReaderRow<'_, '_>) -> rusqlite::Result<T>,
763 {
764 crate::sql_bridge::reader_capability_admits(sql).map_err(SqliteError::InvalidData)?;
765 if !self.reusable.get() {
766 return Err(SqliteError::InvalidData(
767 "reader lease is quarantined after failed read cleanup".into(),
768 ));
769 }
770 if self.query_in_progress.replace(true) {
774 return Err(SqliteError::InvalidData(
775 "reader lease is already executing a query".into(),
776 ));
777 }
778 let _in_progress = ReaderQueryInProgress(&self.query_in_progress);
779 self.mark_dirty();
780 crate::read_cancellation::run_borrowed_reader(self, |conn, admission| {
781 conn.query_row(sql, params, |row| {
782 if !admission.admits() {
790 return Err(rusqlite::Error::SqliteFailure(
791 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_INTERRUPT),
792 Some("request stopped before the mapper".into()),
793 ));
794 }
795 f(&ReaderRow { row })
796 })
797 .map_err(|error| {
798 StorageError::driver(StorageCapability::Sql, "reader_guard.query_row", error)
799 })
800 })
801 .map_err(|error| {
802 self.pool.record_reader_query_error(&error);
803 match error {
804 StorageError::Driver {
805 capability,
806 operation,
807 source,
808 } => match source.downcast::<rusqlite::Error>() {
809 Ok(error) => SqliteError::Rusqlite(*error),
810 Err(source) => SqliteError::RequestReadStopped(StorageError::Driver {
811 capability,
812 operation,
813 source,
814 }),
815 },
816 other => SqliteError::RequestReadStopped(other),
817 }
818 })
819 }
820
821 pub(crate) fn discard(&self) {
825 self.reusable.set(false);
826 }
827
828 pub(crate) fn mark_dirty(&self) {
831 self.dirty.set(true);
832 }
833
834 pub(crate) fn label_operation(&mut self, operation: &'static str) {
838 self.operation = Some(operation);
839 }
840}
841
842impl<'pool> Drop for ReaderGuard<'pool> {
843 fn drop(&mut self) {
844 let Some(lease) = self.lease.take() else {
845 return;
846 };
847
848 match lease {
849 ReaderLease::Pooled(conn) if self.reusable.get() => {
850 self.pool.return_reader(conn, self.dirty.get())
851 }
852 ReaderLease::Pooled(conn) => {
853 close_connection_quietly(conn);
854 self.pool.replace_discarded_reader_slot();
855 }
856 ReaderLease::Shared(guard) if !self.reusable.get() => {
857 self.pool.retire_pooled_writer(&guard);
858 }
859 ReaderLease::Shared(guard) => {
860 if self.dirty.get() && !restore_shared_reader_state(&guard, &self.pool.config) {
868 self.pool.retire_pooled_writer(&guard);
869 }
870 }
871 }
872
873 drop(self.admission_slot.take());
877 self.pool
878 .reader_acquisition_counters
879 .record_checkout_completed(self.checked_out_at.elapsed(), self.operation);
880 }
881}
882
883pub(crate) struct SharedReaderTransactionGuard {
899 conn: parking_lot::ArcMutexGuard<parking_lot::RawMutex, Connection>,
900 admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
901 pool: Arc<ConnectionPool>,
902 checked_out_at: Instant,
903 poison: Cell<bool>,
909}
910
911impl SharedReaderTransactionGuard {
912 pub(crate) fn conn(&self) -> &Connection {
913 &self.conn
914 }
915
916 pub(crate) fn poison(&self) {
919 self.poison.set(true);
920 }
921}
922
923impl Drop for SharedReaderTransactionGuard {
924 fn drop(&mut self) {
925 let mut restored = !self.poison.get();
932 if restored && !self.conn.is_autocommit() {
933 restored = self.conn.execute_batch("ROLLBACK").is_ok() && self.conn.is_autocommit();
934 }
935 if restored {
936 restored = restore_shared_reader_state(&self.conn, &self.pool.config);
937 }
938 if !restored {
939 self.pool.retire_pooled_writer(&self.conn);
940 }
941 drop(self.admission_slot.take());
942 self.pool
943 .reader_acquisition_counters
944 .record_checkout_completed(
945 self.checked_out_at.elapsed(),
946 Some("explicit_sql_read_transaction"),
947 );
948 }
949}
950
951impl ConnectionPool {
952 pub(crate) fn checkout_shared_reader_transaction(
963 self: &Arc<Self>,
964 should_stop: impl Fn() -> bool,
965 ) -> Result<Option<SharedReaderTransactionGuard>, SqliteError> {
966 debug_assert_eq!(
967 self.max_readers, 0,
968 "the owned shared-reader-transaction guard exists only for the degraded, \
969 single-connection backend"
970 );
971 self.ensure_pooled_writer_active()?;
972 let started = Instant::now();
973 let admission_slot = loop {
974 if should_stop() {
975 return Ok(None);
976 }
977 match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
978 Ok(slot) => break slot,
979 Err(tokio::sync::TryAcquireError::Closed) => {
980 return Err(SqliteError::InvalidData(
981 "reader admission semaphore is closed".to_string(),
982 ));
983 }
984 Err(tokio::sync::TryAcquireError::NoPermits) => {}
985 }
986 if started.elapsed() >= self.config.checkout_timeout {
987 self.reader_acquisition_counters.record_checkout_timeout();
988 return Err(pool_exhausted_error(
989 self.config.checkout_timeout,
990 self.max_readers,
991 ));
992 }
993 thread::yield_now();
994 };
995
996 loop {
997 if should_stop() {
998 return Ok(None);
999 }
1000 let remaining = self
1001 .config
1002 .checkout_timeout
1003 .saturating_sub(started.elapsed());
1004 if remaining.is_zero() {
1005 self.reader_acquisition_counters.record_checkout_timeout();
1006 return Err(pool_exhausted_error(
1007 self.config.checkout_timeout,
1008 self.max_readers,
1009 ));
1010 }
1011 if let Some(conn) = self
1012 .writer
1013 .try_lock_arc_for(remaining.min(Duration::from_millis(2)))
1014 {
1015 self.ensure_pooled_writer_active()?;
1016 self.reader_acquisition_counters.record_pooled_checkout();
1017 return Ok(Some(SharedReaderTransactionGuard {
1018 conn,
1019 admission_slot: Some(admission_slot),
1020 pool: Arc::clone(self),
1021 checked_out_at: Instant::now(),
1022 poison: Cell::new(false),
1023 }));
1024 }
1025 }
1026 }
1027}
1028
1029#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1035#[allow(dead_code)] pub(crate) enum StandaloneReaderPurpose {
1037 ExplicitSqlReadTransaction,
1039 BootSchemaProbe,
1042 DiagnosticsIndependentSnapshot,
1045}
1046
1047impl StandaloneReaderPurpose {
1048 fn is_infrastructure(self) -> bool {
1049 matches!(
1050 self,
1051 Self::BootSchemaProbe | Self::DiagnosticsIndependentSnapshot
1052 )
1053 }
1054}
1055
1056#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
1064pub struct ReaderAcquisitionSnapshot {
1065 pub reader_admission_capacity: usize,
1068 pub available_reader_admission_slots: usize,
1071 pub acquisitions: u64,
1074 pub pooled_checkouts: u64,
1076 pub standalone_opens: u64,
1079 pub infrastructure_standalone_opens: u64,
1082 pub checkout_timeouts: u64,
1086 pub busy_timeouts: u64,
1094 pub active_pooled_checkouts: u64,
1096 pub peak_active_pooled_checkouts: u64,
1098 pub completed_pooled_checkouts: u64,
1100 pub max_completed_hold_micros: u64,
1103 pub max_completed_hold_operation: Option<&'static str>,
1108 pub reader_replacement_open_failures: u64,
1114}
1115
1116#[derive(Debug, Default, Clone, Copy)]
1120struct LongestCompletedHold {
1121 micros: u64,
1122 operation: Option<&'static str>,
1123}
1124
1125#[derive(Debug, Default)]
1126struct ReaderAcquisitionCounters {
1127 pooled_checkouts: AtomicU64,
1128 standalone_opens: AtomicU64,
1129 infrastructure_standalone_opens: AtomicU64,
1130 checkout_timeouts: AtomicU64,
1131 busy_timeouts: AtomicU64,
1132 active_pooled_checkouts: AtomicU64,
1133 peak_active_pooled_checkouts: AtomicU64,
1134 completed_pooled_checkouts: AtomicU64,
1135 longest_completed_hold: parking_lot::Mutex<LongestCompletedHold>,
1136 reader_replacement_open_failures: AtomicU64,
1137}
1138
1139impl ReaderAcquisitionCounters {
1140 fn record_pooled_checkout(&self) {
1141 self.pooled_checkouts.fetch_add(1, Ordering::Relaxed);
1142 let active = self
1143 .active_pooled_checkouts
1144 .fetch_add(1, Ordering::Relaxed)
1145 .saturating_add(1);
1146 self.peak_active_pooled_checkouts
1147 .fetch_max(active, Ordering::Relaxed);
1148 }
1149
1150 fn record_checkout_timeout(&self) {
1151 self.checkout_timeouts.fetch_add(1, Ordering::Relaxed);
1152 }
1153
1154 fn record_busy_timeout(&self) {
1155 self.busy_timeouts.fetch_add(1, Ordering::Relaxed);
1156 }
1157
1158 fn record_reader_replacement_open_failure(&self) {
1159 self.reader_replacement_open_failures
1160 .fetch_add(1, Ordering::Relaxed);
1161 }
1162
1163 fn record_standalone_open(&self, purpose: StandaloneReaderPurpose) {
1164 if purpose.is_infrastructure() {
1165 self.infrastructure_standalone_opens
1166 .fetch_add(1, Ordering::Relaxed);
1167 } else {
1168 self.standalone_opens.fetch_add(1, Ordering::Relaxed);
1169 }
1170 }
1171
1172 fn record_checkout_completed(&self, hold: Duration, operation: Option<&'static str>) {
1173 let previous = self.active_pooled_checkouts.fetch_sub(1, Ordering::Relaxed);
1174 debug_assert!(previous > 0, "reader active-checkout counter underflow");
1175 self.completed_pooled_checkouts
1176 .fetch_add(1, Ordering::Relaxed);
1177 let micros = u64::try_from(hold.as_micros()).unwrap_or(u64::MAX);
1178 let mut longest = self.longest_completed_hold.lock();
1182 if micros > longest.micros {
1183 longest.micros = micros;
1184 longest.operation = operation;
1185 }
1186 }
1187
1188 fn snapshot(
1189 &self,
1190 reader_admission_capacity: usize,
1191 available_reader_admission_slots: usize,
1192 ) -> ReaderAcquisitionSnapshot {
1193 let pooled_checkouts = self.pooled_checkouts.load(Ordering::Relaxed);
1194 let standalone_opens = self.standalone_opens.load(Ordering::Relaxed);
1195 let longest_completed_hold = *self.longest_completed_hold.lock();
1196 ReaderAcquisitionSnapshot {
1197 reader_admission_capacity,
1198 available_reader_admission_slots,
1199 acquisitions: pooled_checkouts.saturating_add(standalone_opens),
1200 pooled_checkouts,
1201 standalone_opens,
1202 infrastructure_standalone_opens: self
1203 .infrastructure_standalone_opens
1204 .load(Ordering::Relaxed),
1205 checkout_timeouts: self.checkout_timeouts.load(Ordering::Relaxed),
1206 busy_timeouts: self.busy_timeouts.load(Ordering::Relaxed),
1207 active_pooled_checkouts: self.active_pooled_checkouts.load(Ordering::Relaxed),
1208 peak_active_pooled_checkouts: self.peak_active_pooled_checkouts.load(Ordering::Relaxed),
1209 completed_pooled_checkouts: self.completed_pooled_checkouts.load(Ordering::Relaxed),
1210 max_completed_hold_micros: longest_completed_hold.micros,
1211 max_completed_hold_operation: longest_completed_hold.operation,
1212 reader_replacement_open_failures: self
1213 .reader_replacement_open_failures
1214 .load(Ordering::Relaxed),
1215 }
1216 }
1217}
1218
1219impl ConnectionPool {
1220 pub fn new(config: PoolConfig) -> Result<Self, SqliteError> {
1228 refuse_home_data_store_in_tests(&config)?;
1229 validate_write_admission_deadline(config.write_admission_deadline_ms)?;
1230 if let Some(policy) = config.disk_guard_config {
1231 policy.validate()?;
1232 }
1233 config.wal_ceiling.validate_static(
1234 config.path.is_some(),
1235 config.wal_mode,
1236 config.read_only,
1237 )?;
1238 code_map::validate_pool(&config)?;
1239 if config.path.is_some() && !config.read_only {
1240 let lock_dir =
1241 crate::disk_guard_config::require_volume_lock_dir(config.volume_lock_dir.clone())?;
1242 if !lock_dir.is_absolute() {
1243 return Err(SqliteError::CapacityUnavailable {
1244 phase: CapacityUnavailablePhase::Lock,
1245 message: "configured volume-lock directory is not absolute".to_string(),
1246 });
1247 }
1248 }
1249
1250 let mut config = config;
1254 let inert_memory_queue_request =
1255 config.path.is_none() && config.write_queue_enabled == Some(true);
1256 config.write_queue_enabled =
1257 Some(config.write_queue_enabled.unwrap_or(config.path.is_some()));
1258 if inert_memory_queue_request {
1259 tracing::warn!(
1260 "write queue explicitly requested for an in-memory pool; it is inert because \
1261 in-memory pools cannot host a writer task"
1262 );
1263 }
1264
1265 let (origin, identity_path) = match config.path.as_ref() {
1270 Some(path) if config.code_map_vfs.is_some() => code_map::guarded_identity(path),
1271 Some(path) => {
1272 let (identity, canonical) = mint_db_identity(path)?;
1273 (TxOrigin::Database(identity), Some(canonical))
1274 }
1275 None => (TxOrigin::Memory, None),
1276 };
1277 let (write_admission, disk_guard_config) = if !config.read_only {
1278 if identity_path.as_deref().and_then(Path::parent).is_some() {
1279 let disk_guard = match config.disk_guard_config {
1280 Some(policy) => policy,
1281 None => resolve_disk_guard_config(None, None)?,
1282 };
1283 disk_guard.validate()?;
1284 if disk_guard.legacy_environment_present {
1285 tracing::warn!(
1286 "legacy SQLite reserve setting is deprecated; \
1287 use KHIVE_SQLITE_DISK_RESERVE_BYTES"
1288 );
1289 }
1290 if disk_guard.reserve_bytes == 0 {
1291 tracing::warn!(
1292 "SQLite disk reserve is explicitly zero; \
1293 new logical writes will not be floor-refused"
1294 );
1295 }
1296 (
1297 Arc::new(WriteAdmission::new(
1298 identity_path.clone(),
1299 disk_guard.reserve_bytes,
1300 disk_guard.guard_deadline_ms,
1301 config.volume_lock_dir.clone(),
1302 )?),
1303 Some(disk_guard),
1304 )
1305 } else {
1306 (
1307 Arc::new(WriteAdmission::new(
1308 None,
1309 0,
1310 DEFAULT_DISK_GUARD_DEADLINE_MS,
1311 None,
1312 )?),
1313 None,
1314 )
1315 }
1316 } else {
1317 (
1318 Arc::new(WriteAdmission::new(
1319 None,
1320 0,
1321 DEFAULT_DISK_GUARD_DEADLINE_MS,
1322 None,
1323 )?),
1324 None,
1325 )
1326 };
1327 let read_only_open_target = read_only_open_target(&config, identity_path.as_deref())?;
1328 #[cfg(any(unix, windows))]
1329 let identity_before_open = identity_path
1330 .as_deref()
1331 .map(database_file_identity_if_exists)
1332 .transpose()?
1333 .flatten();
1334 #[cfg(any(unix, windows))]
1335 claimed_file_identity::verify_before_open(&config, identity_before_open)?;
1336 let retirement_connection = Connection::open_in_memory()?;
1337 retirement_connection.authorizer(Some(deny_retired_writer))?;
1338 let mut initialization_lease = None;
1340 let mut writer = open_writer_connection(
1341 &config,
1342 read_only_open_target.as_deref(),
1343 identity_path.as_deref(),
1344 )?;
1345 validate_wal_ceiling_at_open(&writer, &config)?;
1346 writer.busy_timeout(config.busy_timeout)?;
1349 let initial_database_id = if identity_path.is_some() {
1353 read_database_id(&writer)?
1354 } else {
1355 None
1356 };
1357 #[cfg(test)]
1358 if let Some(path) = identity_path.as_deref() {
1359 run_identity_open_hook(
1360 path,
1361 IdentityOpenStage::AfterMainOpenBeforeFirstStat,
1362 Some(&writer),
1363 );
1364 }
1365 #[cfg(any(unix, windows))]
1366 let identity_before_write = identity_path
1367 .as_deref()
1368 .map(database_file_identity)
1369 .transpose()?;
1370 #[cfg(any(unix, windows))]
1371 if identity_before_open.is_some() && identity_before_open != identity_before_write {
1372 return Err(SqliteError::InvalidData(
1373 "database file identity changed while opening the pool".to_string(),
1374 ));
1375 }
1376 #[cfg(any(unix, windows))]
1377 if let Some(path) = identity_path.as_deref() {
1378 let opened = opened_sqlite_file_identity(&writer, path)?;
1379 if identity_before_write != Some(opened) {
1380 return Err(SqliteError::InvalidData(
1381 "database file identity changed while opening the pool".to_string(),
1382 ));
1383 }
1384 }
1385 write_admission.verify_current_volume()?;
1388 let opened_database_id = if identity_path.is_some()
1389 && !config.read_only
1390 && initial_database_id.is_none()
1391 {
1392 match write_admission.check() {
1393 Ok(()) => {
1394 initialization_lease = write_admission.acquire()?;
1395 match initialize_database_id_with_admission(&mut writer, Some(&write_admission))
1396 {
1397 Ok(id) => Some(id),
1398 Err(SqliteError::CapacityFloor { .. }) => None,
1399 Err(error) => return Err(error),
1400 }
1401 }
1402 Err(SqliteError::CapacityFloor { .. }) => None,
1406 Err(error) => return Err(error),
1407 }
1408 } else {
1409 initial_database_id
1410 };
1411 drop(initialization_lease.take());
1412 #[cfg(test)]
1413 if let Some(path) = identity_path.as_deref() {
1414 run_identity_open_hook(
1415 path,
1416 IdentityOpenStage::AfterInitialIdentityWrite,
1417 Some(&writer),
1418 );
1419 }
1420 #[cfg(any(unix, windows))]
1421 let opened_file_identity = identity_path
1422 .as_deref()
1423 .map(|path| opened_sqlite_file_identity(&writer, path))
1424 .transpose()?;
1425 #[cfg(any(unix, windows))]
1426 if identity_before_write != opened_file_identity {
1427 return Err(SqliteError::InvalidData(
1428 "database file identity changed while opening the pool".to_string(),
1429 ));
1430 }
1431 let wal_enabled = configure_writer_connection(&writer, &config)?;
1432 let max_readers = effective_reader_count(&config, wal_enabled);
1433
1434 let readers = ArrayQueue::new(max_readers.max(1));
1435
1436 #[cfg(any(test, feature = "test-support"))]
1437 let statement_observer = crate::statement_observer::StatementObserverHub::new()?;
1438 #[cfg(any(test, feature = "test-support"))]
1439 crate::statement_observer::install(&writer, &statement_observer)?;
1440
1441 let mut pool = Self {
1442 writer: Arc::new(Mutex::new(writer)),
1443 retirement_connection: Mutex::new(Some(retirement_connection)),
1444 #[cfg(any(test, feature = "test-support"))]
1445 statement_observer,
1446 main_pool_generation: OnceLock::new(),
1447 checkpoint_ownership: CheckpointOwnershipGate::new(),
1448 pooled_writer_retired: AtomicBool::new(false),
1449 writer_acquisition_counters: Arc::new(WriterAcquisitionCounters::default()),
1450 write_admission,
1451 disk_guard_config,
1452 reader_acquisition_counters: ReaderAcquisitionCounters::default(),
1453 search_dispatches: Mutex::new(BTreeMap::new()),
1454 note_candidate_hydration_rows: AtomicU64::new(0),
1455 readers,
1456 max_readers,
1457 config,
1458 read_only_open_target,
1459 sql_bridge_reader_slots: Arc::new(Semaphore::new(max_readers.max(1))),
1460 sql_bridge_writer_slots: Arc::new(Semaphore::new(1)),
1461 writer_task: OnceLock::new(),
1462 writer_task_join: Mutex::new(None),
1463 writer_task_join_stored: AtomicBool::new(false),
1464 origin,
1465 identity_path,
1466 #[cfg(any(unix, windows))]
1467 opened_file_identity,
1468 opened_database_id,
1469 identity_registration: None,
1470 #[cfg(test)]
1471 writer_task_spawn_count: std::sync::atomic::AtomicUsize::new(0),
1472 };
1473
1474 for _ in 0..pool.max_readers {
1475 let conn = pool.open_reader_connection()?;
1476 pool.readers
1477 .push(conn)
1478 .expect("reader queue must have capacity during pool initialization");
1479 }
1480
1481 if !pool.config.read_only {
1486 crate::timeout_sink::init(
1487 pool.canonical_path().and_then(Path::parent),
1488 &crate::timeout_sink::db_label(&pool),
1489 );
1490 }
1491
1492 pool.identity_registration = pool.canonical_path().map(PoolIdentityRegistration::new);
1493 Ok(pool)
1494 }
1495
1496 pub fn reader(&self) -> Result<ReaderGuard<'_>, SqliteError> {
1506 self.reader_until(|| false)?.ok_or_else(|| {
1507 SqliteError::InvalidData("uncancelled reader checkout stopped unexpectedly".into())
1508 })
1509 }
1510
1511 pub(crate) fn reader_until<C>(
1517 &self,
1518 should_stop: C,
1519 ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
1520 where
1521 C: Fn() -> bool,
1522 {
1523 let started = Instant::now();
1524 let mut admission_attempt = 0u32;
1525 let admission_slot = loop {
1526 if should_stop() {
1527 return Ok(None);
1528 }
1529 match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
1530 Ok(slot) => break slot,
1531 Err(tokio::sync::TryAcquireError::Closed) => {
1532 return Err(SqliteError::InvalidData(
1533 "reader admission semaphore is closed".to_string(),
1534 ));
1535 }
1536 Err(tokio::sync::TryAcquireError::NoPermits) => {}
1537 }
1538 if started.elapsed() >= self.config.checkout_timeout {
1539 self.reader_acquisition_counters.record_checkout_timeout();
1540 return Err(pool_exhausted_error(
1541 self.config.checkout_timeout,
1542 self.max_readers,
1543 ));
1544 }
1545 match admission_attempt {
1546 0..=7 => {
1547 let spins = 1usize << admission_attempt;
1548 for _ in 0..spins {
1549 std::hint::spin_loop();
1550 }
1551 }
1552 8..=15 => thread::yield_now(),
1553 _ => {
1554 let remaining = self
1555 .config
1556 .checkout_timeout
1557 .saturating_sub(started.elapsed());
1558 let sleep =
1559 Duration::from_micros(50 * (1u64 << (admission_attempt - 16).min(6)));
1560 thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
1561 }
1562 }
1563 admission_attempt = admission_attempt.saturating_add(1);
1564 };
1565
1566 self.reader_with_admission(
1567 ReaderAdmission {
1568 slot: admission_slot,
1569 started,
1570 },
1571 should_stop,
1572 )
1573 }
1574
1575 pub(crate) async fn acquire_reader_admission(
1584 &self,
1585 capability: StorageCapability,
1586 operation: &'static str,
1587 ) -> Result<ReaderAdmission, StorageError> {
1588 let context = khive_storage::capture_request_read_context();
1589 let started = Instant::now();
1590 let stopped = context.stop_reason().is_some();
1591 let outcome: Result<Option<ReaderAdmission>, SqliteError> = if stopped {
1592 Ok(None)
1593 } else {
1594 tokio::select! {
1595 biased;
1596 _ = context.wait_for_stop() => Ok(None),
1597 waited = tokio::time::timeout(
1598 self.config.checkout_timeout,
1599 Arc::clone(&self.sql_bridge_reader_slots).acquire_owned(),
1600 ) => match waited {
1601 Ok(Ok(slot)) => Ok(Some(ReaderAdmission { slot, started })),
1602 Ok(Err(_closed)) => Err(SqliteError::InvalidData(
1603 "reader admission semaphore is closed".to_string(),
1604 )),
1605 Err(_elapsed) => {
1606 self.reader_acquisition_counters.record_checkout_timeout();
1607 Err(pool_exhausted_error(
1608 self.config.checkout_timeout,
1609 self.max_readers,
1610 ))
1611 }
1612 },
1613 }
1614 };
1615 match outcome {
1616 Ok(Some(admission)) => Ok(admission),
1617 Ok(None) => Err(self.reader_checkout_refusal(capability, operation, None)),
1618 Err(error) => Err(self.reader_checkout_refusal(capability, operation, Some(error))),
1619 }
1620 }
1621
1622 pub(crate) fn reader_with_admission<C>(
1629 &self,
1630 admission: ReaderAdmission,
1631 should_stop: C,
1632 ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
1633 where
1634 C: Fn() -> bool,
1635 {
1636 let ReaderAdmission {
1637 slot: admission_slot,
1638 started,
1639 } = admission;
1640
1641 if self.max_readers == 0 {
1642 self.ensure_pooled_writer_active()?;
1643 loop {
1644 if should_stop() {
1645 return Ok(None);
1646 }
1647 let remaining = self
1648 .config
1649 .checkout_timeout
1650 .saturating_sub(started.elapsed());
1651 if remaining.is_zero() {
1652 self.reader_acquisition_counters.record_checkout_timeout();
1653 return Err(pool_exhausted_error(
1654 self.config.checkout_timeout,
1655 self.max_readers,
1656 ));
1657 }
1658 if let Some(guard) = self
1659 .writer
1660 .try_lock_for(remaining.min(Duration::from_millis(2)))
1661 {
1662 self.ensure_pooled_writer_active()?;
1663 self.reader_acquisition_counters.record_pooled_checkout();
1664 return Ok(Some(ReaderGuard {
1665 lease: Some(ReaderLease::Shared(guard)),
1666 admission_slot: Some(admission_slot),
1667 pool: self,
1668 reusable: Cell::new(true),
1669 query_in_progress: Cell::new(false),
1670 checked_out_at: Instant::now(),
1671 dirty: Cell::new(false),
1672 operation: None,
1673 }));
1674 }
1675 }
1676 }
1677
1678 let mut attempt = 0u32;
1679
1680 loop {
1681 if should_stop() {
1682 return Ok(None);
1683 }
1684 if let Some(conn) = self.readers.pop() {
1685 self.reader_acquisition_counters.record_pooled_checkout();
1686 return Ok(Some(ReaderGuard {
1687 lease: Some(ReaderLease::Pooled(conn)),
1688 admission_slot: Some(admission_slot),
1689 pool: self,
1690 reusable: Cell::new(true),
1691 query_in_progress: Cell::new(false),
1692 checked_out_at: Instant::now(),
1693 dirty: Cell::new(false),
1694 operation: None,
1695 }));
1696 }
1697
1698 if started.elapsed() >= self.config.checkout_timeout {
1699 self.reader_acquisition_counters.record_checkout_timeout();
1700 return Err(pool_exhausted_error(
1701 self.config.checkout_timeout,
1702 self.max_readers,
1703 ));
1704 }
1705
1706 match attempt {
1707 0..=7 => {
1708 let spins = 1usize << attempt;
1709 for _ in 0..spins {
1710 std::hint::spin_loop();
1711 }
1712 }
1713 8..=15 => thread::yield_now(),
1714 _ => {
1715 let remaining = self
1716 .config
1717 .checkout_timeout
1718 .saturating_sub(started.elapsed());
1719 let sleep = Duration::from_micros(50 * (1u64 << (attempt - 16).min(6)));
1720 thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
1721 }
1722 }
1723
1724 attempt = attempt.saturating_add(1);
1725 }
1726 }
1727
1728 #[track_caller]
1734 pub fn writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
1735 self.writer_with_checkout_probe(true, false)
1736 }
1737
1738 #[track_caller]
1741 pub(crate) fn writer_for_admitted_operation(&self) -> Result<WriterGuard<'_>, SqliteError> {
1742 self.writer_with_checkout_probe(false, false)
1743 }
1744
1745 fn writer_for_checkpoint_operation(&self) -> Result<WriterGuard<'_>, SqliteError> {
1748 self.writer_with_checkout_probe(false, true)
1749 }
1750
1751 #[track_caller]
1752 fn writer_with_checkout_probe(
1753 &self,
1754 compatibility_probe: bool,
1755 checkpoint_bypass: bool,
1756 ) -> Result<WriterGuard<'_>, SqliteError> {
1757 if !checkpoint_bypass {
1758 self.write_admission.ensure_settled()?;
1759 }
1760 self.ensure_pooled_writer_active()?;
1761 let volume_lease = if checkpoint_bypass {
1762 None
1763 } else {
1764 self.write_admission.acquire()?
1765 };
1766 let Some(guard) = self.writer.try_lock_for(self.config.checkout_timeout) else {
1767 self.writer_acquisition_counters
1768 .pooled_timeouts
1769 .fetch_add(1, Ordering::Relaxed);
1770 let message = format!(
1771 "timed out after {:?} waiting for sqlite writer connection",
1772 self.config.checkout_timeout
1773 );
1774 crate::timeout_sink::emit_timeout(
1775 &crate::timeout_sink::db_label(self),
1776 crate::timeout_sink::Site::PoolAdmission,
1777 &message,
1778 Some(
1779 self.config
1780 .checkout_timeout
1781 .as_millis()
1782 .min(u128::from(u64::MAX)) as u64,
1783 ),
1784 );
1785 return Err(SqliteError::WriterPoolCheckoutTimeout {
1786 timeout: self.config.checkout_timeout,
1787 });
1788 };
1789 self.ensure_pooled_writer_active()?;
1790 #[cfg(any(unix, windows))]
1791 if let Some(path) = self.identity_path.as_deref() {
1792 self.verify_opened_file_identity(path)?;
1793 self.verify_connection_file_identity(&guard, path)?;
1794 }
1795 if compatibility_probe && !checkpoint_bypass {
1799 self.write_admission.check()?;
1800 }
1801 self.writer_acquisition_counters
1802 .pooled_acquisitions
1803 .fetch_add(1, Ordering::Relaxed);
1804 Ok(WriterGuard {
1805 guard,
1806 origin: self.origin(),
1807 pool: self,
1808 admission: self.write_admission.as_ref(),
1809 _volume_lease: volume_lease,
1810 })
1811 }
1812
1813 #[track_caller]
1818 pub fn try_writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
1819 self.writer()
1820 }
1821
1822 #[cfg(test)]
1825 #[track_caller]
1826 fn writer_until<C>(&self, should_stop: C) -> Result<Option<WriterGuard<'_>>, SqliteError>
1827 where
1828 C: Fn() -> bool,
1829 {
1830 self.writer_until_with_checkout_probe(should_stop, true)
1831 }
1832
1833 #[track_caller]
1834 pub(crate) fn writer_until_for_admitted_operation<C>(
1835 &self,
1836 should_stop: C,
1837 ) -> Result<Option<WriterGuard<'_>>, SqliteError>
1838 where
1839 C: Fn() -> bool,
1840 {
1841 self.writer_until_with_checkout_probe(should_stop, false)
1842 }
1843
1844 #[track_caller]
1845 fn writer_until_with_checkout_probe<C>(
1846 &self,
1847 should_stop: C,
1848 compatibility_probe: bool,
1849 ) -> Result<Option<WriterGuard<'_>>, SqliteError>
1850 where
1851 C: Fn() -> bool,
1852 {
1853 self.ensure_pooled_writer_active()?;
1854 if should_stop() {
1855 return Ok(None);
1856 }
1857 let volume_lease = self.write_admission.acquire()?;
1858 let started = Instant::now();
1859 loop {
1860 if should_stop() {
1861 return Ok(None);
1862 }
1863 let remaining = self
1864 .config
1865 .checkout_timeout
1866 .saturating_sub(started.elapsed());
1867 if let Some(guard) = self
1868 .writer
1869 .try_lock_for(remaining.min(Duration::from_millis(2)))
1870 {
1871 if should_stop() {
1874 return Ok(None);
1875 }
1876 self.ensure_pooled_writer_active()?;
1877 #[cfg(any(unix, windows))]
1878 if let Some(path) = self.identity_path.as_deref() {
1879 self.verify_opened_file_identity(path)?;
1880 self.verify_connection_file_identity(&guard, path)?;
1881 }
1882 if compatibility_probe {
1883 self.write_admission.check()?;
1884 }
1885 self.writer_acquisition_counters
1886 .pooled_acquisitions
1887 .fetch_add(1, Ordering::Relaxed);
1888 return Ok(Some(WriterGuard {
1889 guard,
1890 origin: self.origin(),
1891 pool: self,
1892 admission: self.write_admission.as_ref(),
1893 _volume_lease: volume_lease,
1894 }));
1895 }
1896 if started.elapsed() >= self.config.checkout_timeout {
1897 self.writer_acquisition_counters
1898 .pooled_timeouts
1899 .fetch_add(1, Ordering::Relaxed);
1900 let message = format!(
1901 "timed out after {:?} waiting for sqlite writer connection",
1902 self.config.checkout_timeout
1903 );
1904 crate::timeout_sink::emit_timeout(
1905 &crate::timeout_sink::db_label(self),
1906 crate::timeout_sink::Site::PoolAdmission,
1907 &message,
1908 Some(
1909 self.config
1910 .checkout_timeout
1911 .as_millis()
1912 .min(u128::from(u64::MAX)) as u64,
1913 ),
1914 );
1915 return Err(SqliteError::WriterPoolCheckoutTimeout {
1916 timeout: self.config.checkout_timeout,
1917 });
1918 }
1919 }
1920 }
1921
1922 pub fn try_checkpoint_nowait(&self) -> Result<CheckpointGuard<'_>, SqliteError> {
1940 self.ensure_pooled_writer_active()?;
1941 let guard = self.writer.try_lock().ok_or_else(|| {
1942 SqliteError::InvalidData(
1943 "writer connection busy (checkpoint skipped this tick)".to_string(),
1944 )
1945 })?;
1946 self.ensure_pooled_writer_active()?;
1947 Ok(CheckpointGuard { guard })
1948 }
1949
1950 pub(crate) fn retire_pooled_writer(&self, conn: &Connection) {
1951 self.pooled_writer_retired.store(true, Ordering::Release);
1952 if let Err(error) = conn.authorizer(Some(deny_retired_writer)) {
1953 tracing::error!(
1954 %error,
1955 "failed to install the retired pooled-writer quarantine authorizer"
1956 );
1957 }
1958 }
1959
1960 #[cfg(test)]
1963 pub(crate) fn probe_retired_pooled_writer_for_test(&self) -> rusqlite::Result<i64> {
1964 self.writer
1965 .lock()
1966 .query_row("SELECT 1", [], |row| row.get(0))
1967 }
1968
1969 #[cfg(test)]
1973 pub(crate) fn leave_pooled_writer_transaction_open_for_test(&self) -> rusqlite::Result<()> {
1974 self.writer.lock().execute_batch("BEGIN IMMEDIATE")
1975 }
1976
1977 fn ensure_pooled_writer_active(&self) -> Result<(), SqliteError> {
1978 if self.pooled_writer_retired.load(Ordering::Acquire) {
1979 return Err(SqliteError::InvalidData(
1980 "pooled writer connection retired after a terminal transaction fault".to_string(),
1981 ));
1982 }
1983 Ok(())
1984 }
1985
1986 pub fn writer_acquisition_snapshot(&self) -> WriterAcquisitionSnapshot {
1989 self.writer_acquisition_counters
1990 .snapshot(self.write_admission.lease_timeouts())
1991 }
1992
1993 pub fn record_search_dispatch(&self, backend_id: &str, requested_kind: &str) {
1997 let mut dispatches = self.search_dispatches.lock();
1998 let count = dispatches
1999 .entry(backend_id.to_owned())
2000 .or_default()
2001 .entry(requested_kind.to_owned())
2002 .or_default();
2003 *count = count.saturating_add(1);
2004 }
2005
2006 pub fn record_note_candidate_hydration_row(&self) {
2010 let _ = self.note_candidate_hydration_rows.fetch_update(
2011 Ordering::Relaxed,
2012 Ordering::Relaxed,
2013 |current| Some(current.saturating_add(1)),
2014 );
2015 }
2016
2017 pub fn search_mechanism_snapshot(&self) -> SearchMechanismSnapshot {
2019 SearchMechanismSnapshot {
2020 dispatches_by_backend_and_kind: self.search_dispatches.lock().clone(),
2021 note_candidate_hydration_rows: self
2022 .note_candidate_hydration_rows
2023 .load(Ordering::Relaxed),
2024 }
2025 }
2026
2027 pub fn reader_acquisition_snapshot(&self) -> ReaderAcquisitionSnapshot {
2031 self.reader_acquisition_counters.snapshot(
2032 self.max_readers.max(1),
2033 self.sql_bridge_reader_slots.available_permits(),
2034 )
2035 }
2036
2037 pub(crate) fn record_reader_admission_timeout(&self) {
2041 self.reader_acquisition_counters.record_checkout_timeout();
2042 }
2043
2044 pub(crate) fn record_reader_query_error(&self, error: &StorageError) {
2050 if crate::read_cancellation::storage_error_sqlite_code(error)
2051 == Some(rusqlite::ErrorCode::DatabaseBusy)
2052 {
2053 self.reader_acquisition_counters.record_busy_timeout();
2054 }
2055 }
2056
2057 pub(crate) fn writer_acquisition_counters(&self) -> Arc<WriterAcquisitionCounters> {
2059 Arc::clone(&self.writer_acquisition_counters)
2060 }
2061
2062 pub(crate) fn write_admission(&self) -> Arc<WriteAdmission> {
2063 Arc::clone(&self.write_admission)
2064 }
2065
2066 pub fn effective_disk_guard_config(&self) -> Option<EffectiveDiskGuardConfig> {
2067 self.disk_guard_config
2068 }
2069
2070 #[cfg(test)]
2071 pub(crate) fn set_test_write_admission(
2072 &mut self,
2073 floor_bytes: u64,
2074 probe: impl Fn(&Path) -> std::io::Result<u64> + Send + Sync + 'static,
2075 ) {
2076 let database_path = if self.config.read_only {
2077 None
2078 } else {
2079 self.canonical_path().map(Path::to_path_buf)
2080 };
2081 let admission = Arc::new(
2082 WriteAdmission::new(
2083 database_path,
2084 floor_bytes,
2085 DEFAULT_DISK_GUARD_DEADLINE_MS,
2086 self.config.volume_lock_dir.clone(),
2087 )
2088 .expect("test admission volume identity"),
2089 );
2090 admission.set_test_space_probe(probe);
2091 self.write_admission = admission;
2092 if self.canonical_path().is_some() && !self.config.read_only {
2093 self.disk_guard_config = Some(EffectiveDiskGuardConfig {
2094 reserve_bytes: floor_bytes,
2095 guard_deadline_ms: DEFAULT_DISK_GUARD_DEADLINE_MS,
2096 reserve_source: DiskGuardConfigSource::Backend,
2097 deadline_source: DiskGuardConfigSource::Default,
2098 legacy_environment_present: false,
2099 });
2100 }
2101 }
2102
2103 pub fn available_readers(&self) -> usize {
2105 self.readers.len()
2106 }
2107
2108 pub fn max_readers(&self) -> usize {
2110 self.max_readers
2111 }
2112
2113 pub fn config(&self) -> &PoolConfig {
2115 &self.config
2116 }
2117
2118 #[cfg(any(test, feature = "test-support"))]
2128 pub fn observe_test_statement_starts(
2129 &self,
2130 limit: usize,
2131 ) -> Result<crate::statement_observer::StatementStartObservation, SqliteError> {
2132 self.statement_observer.observe(limit)
2133 }
2134
2135 pub fn main_pool_generation(&self) -> u64 {
2139 *self.main_pool_generation.get_or_init(|| {
2140 NEXT_MAIN_POOL_GENERATION
2141 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |next| {
2142 next.checked_add(1)
2143 })
2144 .expect("main pool generation exhausted")
2145 })
2146 }
2147
2148 pub(crate) fn reader_admission_timeout(&self, operation: &'static str) -> StorageError {
2152 StorageError::AdmissionTimeout {
2153 operation: operation.into(),
2154 timeout_ms: u64::try_from(self.config.checkout_timeout.as_millis()).unwrap_or(u64::MAX),
2155 pool_identity: Some(
2156 self.identity_registration
2157 .as_ref()
2158 .map(PoolIdentityRegistration::label)
2159 .unwrap_or_else(|| ":memory:".to_string()),
2160 ),
2161 }
2162 }
2163
2164 pub(crate) fn resolve_reader_checkout<'p>(
2183 &self,
2184 capability: StorageCapability,
2185 operation: &'static str,
2186 outcome: Result<Option<ReaderGuard<'p>>, SqliteError>,
2187 ) -> Result<ReaderGuard<'p>, StorageError> {
2188 match outcome {
2189 Ok(Some(mut guard)) => {
2190 guard.label_operation(operation);
2191 Ok(guard)
2192 }
2193 Ok(None) => Err(self.reader_checkout_refusal(capability, operation, None)),
2194 Err(error) => Err(self.reader_checkout_refusal(capability, operation, Some(error))),
2195 }
2196 }
2197
2198 fn reader_checkout_refusal(
2203 &self,
2204 capability: StorageCapability,
2205 operation: &'static str,
2206 error: Option<SqliteError>,
2207 ) -> StorageError {
2208 let Some(error) = error else {
2209 return StorageError::Timeout {
2210 operation: operation.into(),
2211 };
2212 };
2213 let is_pool_exhausted = matches!(
2214 &error,
2215 SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
2216 if code.code == rusqlite::ErrorCode::DatabaseBusy
2217 );
2218 if is_pool_exhausted {
2219 self.reader_admission_timeout(operation)
2220 } else {
2221 StorageError::driver(capability, operation, error)
2222 }
2223 }
2224
2225 pub(crate) fn sql_bridge_reader_slots(&self) -> Arc<Semaphore> {
2228 Arc::clone(&self.sql_bridge_reader_slots)
2229 }
2230
2231 pub(crate) fn sql_bridge_writer_slots(&self) -> Arc<Semaphore> {
2233 Arc::clone(&self.sql_bridge_writer_slots)
2234 }
2235
2236 pub fn retirement_writer_holds(&self) -> usize {
2241 usize::from(self.writer.is_locked())
2242 + usize::from(self.sql_bridge_writer_slots.available_permits() == 0)
2243 }
2244
2245 pub fn origin(&self) -> TxOrigin {
2251 self.origin.clone()
2252 }
2253
2254 pub fn canonical_path(&self) -> Option<&Path> {
2260 self.identity_path.as_deref()
2261 }
2262
2263 #[cfg(unix)]
2265 pub fn opened_file_identity(&self) -> Option<(u64, u64)> {
2266 self.opened_file_identity
2267 .map(DatabaseFileIdentity::unix_parts)
2268 }
2269
2270 #[cfg(any(unix, windows))]
2273 pub fn opened_file_identity_record(&self) -> Option<DatabaseFileIdentity> {
2274 self.opened_file_identity
2275 }
2276
2277 pub fn database_owner_identity(
2282 &self,
2283 ) -> Result<DatabaseOwnerIdentity, DatabaseOwnerIdentityError> {
2284 if self.identity_path.is_none() {
2285 return Err(DatabaseOwnerIdentityError::InMemory);
2286 }
2287 #[cfg(any(unix, windows))]
2288 {
2289 let durable_id = self
2290 .opened_database_id
2291 .ok_or(DatabaseOwnerIdentityError::DurableIdentityUnavailable)?;
2292 let file_identity = self
2293 .opened_file_identity
2294 .ok_or(DatabaseOwnerIdentityError::PhysicalIdentityUnavailable)?;
2295 Ok(DatabaseOwnerIdentity {
2296 durable_id,
2297 file_identity,
2298 })
2299 }
2300 #[cfg(not(any(unix, windows)))]
2301 {
2302 Err(DatabaseOwnerIdentityError::UnsupportedPlatform)
2303 }
2304 }
2305
2306 pub fn verify_database_owner(
2308 &self,
2309 expected: &DatabaseOwnerIdentity,
2310 ) -> Result<(), DatabaseOwnerIdentityError> {
2311 self.database_owner_identity()?.verify_owner(expected)
2312 }
2313
2314 pub fn write_queue_active(&self) -> bool {
2326 debug_assert!(
2327 self.config.write_queue_enabled.is_some(),
2328 "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
2329 before any write_queue_active read"
2330 );
2331 self.config.write_queue_enabled.unwrap_or(false) && self.config.path.is_some()
2332 }
2333
2334 pub fn writer_task_join_was_stored(&self) -> bool {
2340 self.writer_task_join_stored.load(Ordering::SeqCst)
2341 }
2342
2343 pub fn writer_task_handle(&self) -> Result<Option<WriterTaskHandle>, StorageError> {
2367 debug_assert!(
2375 self.config.write_queue_enabled.is_some(),
2376 "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
2377 before any writer_task_handle read"
2378 );
2379 if !self.config.write_queue_enabled.unwrap_or(false) {
2380 return Ok(None);
2381 }
2382 if let Some(existing) = self.writer_task.get() {
2385 return Ok(existing.clone());
2386 }
2387 if tokio::runtime::Handle::try_current().is_err() {
2391 return Err(StorageError::WriterTaskNoRuntime);
2392 }
2393 Ok(self
2394 .writer_task
2395 .get_or_init(|| {
2396 #[cfg(test)]
2397 self.writer_task_spawn_count
2398 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2399
2400 match crate::writer_task::spawn(self, self.config.write_queue_capacity) {
2401 Ok(handle) => Some(handle),
2402 Err(e) => {
2403 tracing::warn!(
2404 error = %e,
2405 "KHIVE_WRITE_QUEUE=1 but the writer task failed to spawn; \
2406 writes fall back to the pool-mutex path"
2407 );
2408 None
2409 }
2410 }
2411 })
2412 .clone())
2413 }
2414
2415 pub(crate) fn writer_task_for_write(
2425 &self,
2426 cached: Option<&WriterTaskHandle>,
2427 operation: &'static str,
2428 ) -> Result<Option<WriterTaskHandle>, StorageError> {
2429 let handle = match cached {
2430 Some(handle) => Some(handle.clone()),
2431 None => match self.writer_task_handle() {
2432 Ok(handle) => handle,
2433 Err(error) if self.config.write_routing_strict => return Err(error),
2434 Err(_) => None,
2435 },
2436 };
2437
2438 if handle.is_none() && self.config.write_routing_strict {
2439 return Err(StorageError::Pool {
2440 operation: operation.into(),
2441 message: "strict write routing requires a writer-task handle; no handle is \
2442 available, so the direct writer fallback was refused"
2443 .into(),
2444 });
2445 }
2446 Ok(handle)
2447 }
2448
2449 pub(crate) fn record_direct_route(&self, site: crate::timeout_sink::Site) {
2453 if self.write_queue_active() {
2454 crate::timeout_sink::emit_direct_route_violation(
2455 &crate::timeout_sink::db_label(self),
2456 site,
2457 );
2458 }
2459 }
2460
2461 pub fn writer_task_for_runtime_write(
2466 &self,
2467 operation: RuntimeWriteOperation,
2468 ) -> Result<Option<WriterTaskHandle>, StorageError> {
2469 let handle = self.writer_task_for_write(None, operation.operation())?;
2470 if handle.is_none() {
2471 self.record_direct_route(operation.fallback_site());
2472 }
2473 Ok(handle)
2474 }
2475
2476 #[cfg(test)]
2481 pub(crate) fn writer_task_spawn_count(&self) -> usize {
2482 self.writer_task_spawn_count
2483 .load(std::sync::atomic::Ordering::SeqCst)
2484 }
2485
2486 pub(crate) fn set_writer_task_join(&self, join: tokio::task::JoinHandle<()>) {
2499 let first_store = !self.writer_task_join_stored.swap(true, Ordering::SeqCst);
2504 debug_assert!(
2505 first_store,
2506 "writer task JoinHandle stored twice (even counting a taken one); \
2507 the writer_task OnceLock is supposed to make spawn at-most-once per pool"
2508 );
2509 if first_store {
2510 *self.writer_task_join.lock() = Some(join);
2511 }
2512 }
2513
2514 pub fn take_writer_task_join(&self) -> Option<tokio::task::JoinHandle<()>> {
2532 self.writer_task_join.lock().take()
2533 }
2534
2535 fn open_reader_connection(&self) -> Result<Connection, SqliteError> {
2536 let path = self.read_connection_path()?;
2537 #[cfg(any(unix, windows))]
2538 if let Some(identity_path) = self.identity_path.as_deref() {
2539 self.verify_opened_file_identity(identity_path)?;
2540 }
2541 let conn = open_reader_connection(path, &self.config)?;
2542 #[cfg(any(unix, windows))]
2543 if let Some(identity_path) = self.identity_path.as_deref() {
2544 self.verify_connection_file_identity(&conn, identity_path)?;
2545 }
2546 self.verify_opened_database_id(&conn)?;
2547 #[cfg(any(test, feature = "test-support"))]
2548 crate::statement_observer::install(&conn, &self.statement_observer)?;
2549 Ok(conn)
2550 }
2551
2552 fn read_connection_path(&self) -> Result<&Path, SqliteError> {
2553 self.read_only_open_target
2554 .as_deref()
2555 .or(self.identity_path.as_deref())
2556 .ok_or_else(|| {
2557 SqliteError::InvalidData(
2558 "in-memory databases do not support standalone connections".to_string(),
2559 )
2560 })
2561 }
2562
2563 pub fn open_standalone_writer(&self) -> Result<Connection, SqliteError> {
2574 self.write_admission.check()?;
2575 self.open_standalone_writer_for_admitted_operation()
2576 }
2577
2578 pub(crate) fn open_standalone_writer_for_admitted_operation(
2582 &self,
2583 ) -> Result<Connection, SqliteError> {
2584 let conn = self.open_standalone_writer_untracked()?;
2585 self.writer_acquisition_counters
2586 .standalone_acquisitions
2587 .fetch_add(1, Ordering::Relaxed);
2588 Ok(conn)
2589 }
2590
2591 pub(crate) fn open_standalone_writer_untracked(&self) -> Result<Connection, SqliteError> {
2601 let path = self.identity_path.as_deref().ok_or_else(|| {
2602 SqliteError::InvalidData(
2603 "in-memory databases do not support standalone connections".to_string(),
2604 )
2605 })?;
2606
2607 if self.config.read_only {
2608 return Err(SqliteError::InvalidData(
2609 "database is read-only: standalone write connections are not permitted".to_string(),
2610 ));
2611 }
2612
2613 #[cfg(any(unix, windows))]
2617 self.verify_opened_file_identity(path)?;
2618
2619 #[cfg(test)]
2620 run_identity_open_hook(path, IdentityOpenStage::BeforeStandaloneOpen, None);
2621
2622 let conn = claimed_file_identity::open_connection(
2623 &self.config,
2624 path,
2625 OpenFlags::SQLITE_OPEN_READ_WRITE
2626 | OpenFlags::SQLITE_OPEN_NO_MUTEX
2627 | OpenFlags::SQLITE_OPEN_URI,
2628 self.identity_path.as_deref(),
2629 )?;
2630 #[cfg(test)]
2631 run_identity_open_hook(path, IdentityOpenStage::AfterStandaloneOpen, Some(&conn));
2632 #[cfg(any(unix, windows))]
2633 self.verify_connection_file_identity(&conn, path)?;
2634 self.verify_opened_database_id(&conn)?;
2635 #[cfg(feature = "namespace-trigram-proto")]
2636 register_namespace_trigram(&conn)?;
2637 register_writer_clock(&conn)?;
2638 register_rfc3339_key(&conn)?;
2642 conn.busy_timeout(self.config.busy_timeout)?;
2643 if self.config.code_map_vfs.is_none() {
2644 self.checkpoint_ownership
2645 .configure_wal_autocheckpoint(&conn)?;
2646 }
2647 conn.pragma_update(None, "foreign_keys", "ON")?;
2648 conn.pragma_update(None, "synchronous", "NORMAL")?;
2649
2650 let wal_enabled =
2651 self.config.wal_mode && current_journal_mode(&conn)?.eq_ignore_ascii_case("wal");
2652 if wal_enabled {
2653 conn.pragma_update(
2654 None,
2655 "journal_size_limit",
2656 self.config.journal_size_limit_bytes,
2657 )?;
2658 }
2659
2660 #[cfg(any(test, feature = "test-support"))]
2661 crate::statement_observer::install(&conn, &self.statement_observer)?;
2662 Ok(conn)
2663 }
2664
2665 #[cfg(any(unix, windows))]
2666 fn verify_opened_file_identity(&self, path: &Path) -> Result<(), SqliteError> {
2667 let Some(expected) = self.opened_file_identity else {
2668 return Err(SqliteError::InvalidData(
2669 "file-backed pool has no opened database file identity".to_string(),
2670 ));
2671 };
2672 let current = database_file_identity(path).ok();
2673 if current != Some(expected) {
2674 return Err(SqliteError::InvalidData(
2675 "pool database file identity changed since the first open; refusing standalone connection"
2676 .to_string(),
2677 ));
2678 }
2679 Ok(())
2680 }
2681
2682 #[cfg(any(unix, windows))]
2683 fn verify_connection_file_identity(
2684 &self,
2685 conn: &Connection,
2686 path: &Path,
2687 ) -> Result<(), SqliteError> {
2688 let opened = opened_sqlite_file_identity(conn, path)?;
2689 if self.opened_file_identity != Some(opened) {
2690 return Err(SqliteError::InvalidData(
2691 "pool database file identity changed since the first open; refusing standalone connection"
2692 .to_string(),
2693 ));
2694 }
2695 Ok(())
2696 }
2697
2698 fn verify_opened_database_id(&self, conn: &Connection) -> Result<(), SqliteError> {
2699 let Some(expected) = self.opened_database_id else {
2703 return Ok(());
2704 };
2705 if self.identity_path.is_some() && read_database_id(conn)? != Some(expected) {
2706 return Err(SqliteError::InvalidData(
2707 "pool database identity changed since the first open; refusing standalone connection"
2708 .to_string(),
2709 ));
2710 }
2711 Ok(())
2712 }
2713
2714 #[cfg(test)]
2718 pub(crate) fn effective_wal_autocheckpoint_pages(&self) -> u32 {
2719 self.checkpoint_ownership.wal_autocheckpoint_pages()
2720 }
2721
2722 pub fn claim_checkpoint_ownership(&self) -> Result<(), SqliteError> {
2744 if !self.checkpoint_ownership.begin_claim() {
2745 return Ok(());
2746 }
2747 let result = (|| {
2748 if !self.config.read_only {
2749 let writer = self.writer_for_checkpoint_operation()?;
2750 writer.conn().pragma_update(None, "wal_autocheckpoint", 0)?;
2751 }
2752 Ok(())
2753 })();
2754 self.checkpoint_ownership.finish_claim(result.is_ok());
2755 result
2756 }
2757
2758 pub async fn propagate_checkpoint_claim_to_writer_task(&self) -> Result<(), StorageError> {
2767 let Some(handle) = self.writer_task_handle()? else {
2768 return Ok(());
2769 };
2770 handle
2771 .send_top_level(|conn| {
2772 conn.pragma_update(None, "wal_autocheckpoint", 0)
2773 .map_err(|e| StorageError::Pool {
2774 operation: "claim_checkpoint_ownership".into(),
2775 message: e.to_string(),
2776 })
2777 })
2778 .await
2779 }
2780
2781 pub(crate) fn open_standalone_reader(
2790 &self,
2791 purpose: StandaloneReaderPurpose,
2792 ) -> Result<Connection, SqliteError> {
2793 let path = self.read_connection_path()?;
2794
2795 #[cfg(any(unix, windows))]
2796 if let Some(identity_path) = self.identity_path.as_deref() {
2797 self.verify_opened_file_identity(identity_path)?;
2798 }
2799
2800 let conn = claimed_file_identity::open_connection(
2801 &self.config,
2802 path,
2803 OpenFlags::SQLITE_OPEN_READ_ONLY
2804 | OpenFlags::SQLITE_OPEN_NO_MUTEX
2805 | OpenFlags::SQLITE_OPEN_URI,
2806 self.identity_path.as_deref(),
2807 )?;
2808 #[cfg(any(unix, windows))]
2809 if let Some(identity_path) = self.identity_path.as_deref() {
2810 self.verify_connection_file_identity(&conn, identity_path)?;
2811 }
2812 self.verify_opened_database_id(&conn)?;
2813 configure_reader_connection(&conn, &self.config)?;
2814 conn.pragma_update(None, "synchronous", "NORMAL")?;
2815 self.reader_acquisition_counters
2816 .record_standalone_open(purpose);
2817 #[cfg(any(test, feature = "test-support"))]
2818 crate::statement_observer::install(&conn, &self.statement_observer)?;
2819 Ok(conn)
2820 }
2821
2822 fn return_reader(&self, conn: Connection, dirty: bool) {
2823 if self.max_readers == 0 {
2824 return;
2825 }
2826
2827 if reset_reader_connection(&conn, dirty, &self.config)
2828 && reader_connection_is_healthy(&conn)
2829 {
2830 self.enqueue_reader_slot(conn);
2831 return;
2832 }
2833
2834 close_connection_quietly(conn);
2835 self.replace_discarded_reader_slot();
2836 }
2837
2838 fn enqueue_reader_slot(&self, conn: Connection) {
2842 if let Err(conn) = self.readers.push(conn) {
2843 eprintln!("[sqlite-pool] reader pool queue full, discarding replacement connection");
2844 close_connection_quietly(conn);
2845 }
2846 }
2847
2848 fn replace_discarded_reader_slot(&self) {
2855 match self.open_reader_connection() {
2856 Ok(conn) => self.enqueue_reader_slot(conn),
2857 Err(error) => {
2858 self.reader_acquisition_counters
2859 .record_reader_replacement_open_failure();
2860 tracing::warn!(
2861 %error,
2862 "sqlite-pool: reader replacement connection failed to open; the physical \
2863 pool permanently shrinks by one slot below max_readers"
2864 );
2865 }
2866 }
2867 }
2868}
2869
2870const MAX_SYMLINK_DEPTH: u32 = 40;
2875
2876fn mint_db_identity(configured_path: &Path) -> Result<(DbIdentity, PathBuf), SqliteError> {
2909 let absolute = if configured_path.is_absolute() {
2910 configured_path.to_path_buf()
2911 } else {
2912 let cwd = std::env::current_dir().map_err(|e| {
2913 SqliteError::InvalidData(format!(
2914 "cannot mint database identity for {configured_path:?}: failed to resolve the \
2915 process current directory: {e}"
2916 ))
2917 })?;
2918 cwd.join(configured_path)
2919 };
2920
2921 if absolute.exists() {
2922 let canonical = absolute.canonicalize().map_err(|e| {
2923 SqliteError::InvalidData(format!(
2924 "cannot mint database identity: failed to canonicalize existing path \
2925 {absolute:?}: {e}"
2926 ))
2927 })?;
2928 return Ok((
2929 DbIdentity::new(canonical.clone().into_os_string()),
2930 canonical,
2931 ));
2932 }
2933
2934 let resolved_target = resolve_symlink_chain(&absolute)?;
2935 let parent = resolved_target.parent().ok_or_else(|| {
2936 SqliteError::InvalidData(format!(
2937 "cannot mint database identity for {resolved_target:?}: path has no parent \
2938 directory"
2939 ))
2940 })?;
2941 let file_name = resolved_target.file_name().ok_or_else(|| {
2942 SqliteError::InvalidData(format!(
2943 "cannot mint database identity for {resolved_target:?}: path has no file name"
2944 ))
2945 })?;
2946 let canonical_parent = parent.canonicalize().map_err(|e| {
2947 SqliteError::InvalidData(format!(
2948 "cannot mint database identity: parent directory {parent:?} of first-open path \
2949 {resolved_target:?} does not exist or is inaccessible: {e}"
2950 ))
2951 })?;
2952 let mut identity_path = canonical_parent;
2953 identity_path.push(file_name);
2954 Ok((
2955 DbIdentity::new(identity_path.clone().into_os_string()),
2956 identity_path,
2957 ))
2958}
2959
2960#[cfg(any(unix, windows))]
2961fn database_file_identity_if_exists(
2962 path: &Path,
2963) -> Result<Option<DatabaseFileIdentity>, SqliteError> {
2964 match database_file_identity(path) {
2965 Ok(identity) => Ok(Some(identity)),
2966 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
2967 Err(error) => Err(error.into()),
2968 }
2969}
2970
2971#[cfg(any(unix, windows))]
2972pub(crate) fn opened_sqlite_file_identity(
2973 conn: &Connection,
2974 path: &Path,
2975) -> Result<DatabaseFileIdentity, SqliteError> {
2976 #[cfg(unix)]
2977 {
2978 verify_sqlite_opened_file_still_at_path(conn)?;
2979 Ok(database_file_identity(path)?)
2980 }
2981 #[cfg(windows)]
2982 {
2983 let opened = sqlite_opened_file_identity(conn)?;
2984 if database_file_identity(path).ok() != Some(opened) {
2985 return Err(SqliteError::InvalidData(
2986 "database file identity changed while SQLite held the opened file".to_string(),
2987 ));
2988 }
2989 Ok(opened)
2990 }
2991}
2992
2993fn read_database_id(conn: &Connection) -> Result<Option<uuid::Uuid>, SqliteError> {
2997 let table_exists: bool = conn.query_row(
2998 "SELECT count(*) != 0 FROM main.sqlite_master WHERE type = 'table' AND name = ?1",
2999 [DATABASE_ID_TABLE],
3000 |row| row.get(0),
3001 )?;
3002 if !table_exists {
3003 return Ok(None);
3006 }
3007 let id: String = conn.query_row(
3008 &format!("SELECT id FROM main.{DATABASE_ID_TABLE} WHERE singleton = 1"),
3009 [],
3010 |row| row.get(0),
3011 )?;
3012 let id = uuid::Uuid::parse_str(&id).map_err(|error| {
3013 SqliteError::InvalidData(format!("invalid stored database identity: {error}"))
3014 })?;
3015 Ok(Some(id))
3016}
3017
3018#[cfg(test)]
3019fn initialize_database_id(conn: &mut Connection) -> Result<uuid::Uuid, SqliteError> {
3020 initialize_database_id_with_admission(conn, None)
3021}
3022
3023fn initialize_database_id_with_admission(
3024 conn: &mut Connection,
3025 admission: Option<&WriteAdmission>,
3026) -> Result<uuid::Uuid, SqliteError> {
3027 if let Some(id) = read_database_id(conn)? {
3028 return Ok(id);
3029 }
3030 let transaction = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
3031 if let Some(admission) = admission {
3032 if let Err(error) = admission.check() {
3033 let rollback = transaction.rollback();
3034 return Err(crate::migrations::capacity_refusal_after_rollback(
3035 conn,
3036 rollback,
3037 error,
3038 "database identity bootstrap",
3039 ));
3040 }
3041 }
3042 transaction.execute_batch(&format!(
3043 "CREATE TABLE IF NOT EXISTS main.{DATABASE_ID_TABLE} (\
3044 singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
3045 id TEXT NOT NULL\
3046 )"
3047 ))?;
3048 transaction.execute(
3049 &format!("INSERT OR IGNORE INTO main.{DATABASE_ID_TABLE} (singleton, id) VALUES (1, ?1)"),
3050 [uuid::Uuid::new_v4().to_string()],
3051 )?;
3052 let id: String = transaction.query_row(
3053 &format!("SELECT id FROM main.{DATABASE_ID_TABLE} WHERE singleton = 1"),
3054 [],
3055 |row| row.get(0),
3056 )?;
3057 let id = uuid::Uuid::parse_str(&id).map_err(|error| {
3058 SqliteError::InvalidData(format!("invalid stored database identity: {error}"))
3059 })?;
3060 transaction.commit()?;
3061 Ok(id)
3062}
3063
3064#[cfg(unix)]
3065fn verify_sqlite_opened_file_still_at_path(conn: &Connection) -> Result<(), SqliteError> {
3066 let mut moved: std::ffi::c_int = 0;
3067 let result = unsafe {
3073 rusqlite::ffi::sqlite3_file_control(
3074 conn.handle(),
3075 c"main".as_ptr(),
3076 rusqlite::ffi::SQLITE_FCNTL_HAS_MOVED,
3077 (&mut moved as *mut std::ffi::c_int).cast(),
3078 )
3079 };
3080 if result != rusqlite::ffi::SQLITE_OK {
3081 return Err(SqliteError::InvalidData(format!(
3082 "cannot verify opened database file identity (SQLite file control {result})"
3083 )));
3084 }
3085 if moved != 0 {
3086 return Err(SqliteError::InvalidData(
3087 "database file identity changed while SQLite held the opened file".to_string(),
3088 ));
3089 }
3090 Ok(())
3091}
3092
3093fn resolve_symlink_chain(path: &Path) -> Result<PathBuf, SqliteError> {
3099 let mut current = path.to_path_buf();
3100 for _ in 0..MAX_SYMLINK_DEPTH {
3101 match fs::symlink_metadata(¤t) {
3102 Ok(meta) if meta.file_type().is_symlink() => {
3103 let target = fs::read_link(¤t).map_err(|e| {
3104 SqliteError::InvalidData(format!(
3105 "cannot mint database identity: failed to read symlink {current:?}: {e}"
3106 ))
3107 })?;
3108 current = if target.is_absolute() {
3109 target
3110 } else {
3111 match current.parent() {
3112 Some(parent) => parent.join(&target),
3113 None => target,
3114 }
3115 };
3116 }
3117 _ => return Ok(current),
3118 }
3119 }
3120 Err(SqliteError::InvalidData(format!(
3121 "cannot mint database identity for {path:?}: symlink chain exceeds \
3122 {MAX_SYMLINK_DEPTH} levels"
3123 )))
3124}
3125
3126fn effective_reader_count(config: &PoolConfig, wal_enabled: bool) -> usize {
3127 if config.path.is_some() && (config.read_only || config.code_map_vfs.is_some()) {
3128 config.max_readers.max(1)
3129 } else if config.path.is_some() && config.wal_mode && wal_enabled {
3130 config.max_readers
3131 } else {
3132 0
3133 }
3134}
3135
3136fn open_writer_connection(
3137 config: &PoolConfig,
3138 read_only_open_target: Option<&Path>,
3139 identity_path: Option<&Path>,
3140) -> Result<Connection, SqliteError> {
3141 claimed_file_identity::open_writer(config, read_only_open_target, identity_path)
3142}
3143
3144fn validate_wal_ceiling_at_open(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
3149 let bytes = config.wal_ceiling.effective_bytes(config.read_only);
3150 if bytes == 0 {
3151 return Ok(());
3152 }
3153 let page_size: i64 = conn.pragma_query_value(None, "page_size", |row| row.get(0))?;
3154 let page_size = u64::try_from(page_size).map_err(|_| {
3155 SqliteError::InvalidData("SQLite reported a negative page size".to_string())
3156 })?;
3157 let minimum_bytes = page_size.checked_add(56).ok_or_else(|| {
3158 SqliteError::InvalidData("SQLite page size overflowed the WAL frame floor".to_string())
3159 })?;
3160 if bytes < minimum_bytes {
3161 return Err(SqliteError::WalCeilingBelowMinimum {
3162 bytes,
3163 page_size,
3164 minimum_bytes,
3165 });
3166 }
3167 Err(SqliteError::WalCapacityUnavailable {
3168 bytes,
3169 capability: "WAL I/O limiter",
3170 })
3171}
3172
3173fn read_only_open_target(
3196 config: &PoolConfig,
3197 physical_path: Option<&Path>,
3198) -> Result<Option<PathBuf>, SqliteError> {
3199 if !config.read_only || config.code_map_vfs.is_some() {
3200 return Ok(None);
3201 }
3202 let Some(path) = physical_path else {
3203 return Ok(None);
3204 };
3205 read_only_wal_open_target_for_path(path).map(Some)
3206}
3207
3208fn read_only_wal_open_target_for_path(path: &Path) -> Result<PathBuf, SqliteError> {
3209 if !sqlite_header_uses_wal(path)? {
3210 return Ok(path.to_path_buf());
3211 }
3212
3213 let shm = sqlite_sidecar_path(path, "-shm");
3214 match fs::metadata(&shm) {
3215 Ok(metadata) if metadata.permissions().readonly() => {
3216 let wal = sqlite_sidecar_path(path, "-wal");
3217 match fs::metadata(&wal) {
3218 Ok(_) => Ok(path.to_path_buf()),
3219 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3220 Err(SqliteError::InvalidData(format!(
3221 "read-only WAL snapshot {} has a shared-memory sidecar {} but no WAL \
3222 sidecar {}; refusing the inconsistent sidecar set before SQLite open",
3223 path.display(),
3224 shm.display(),
3225 wal.display(),
3226 )))
3227 }
3228 Err(error) => Err(SqliteError::Io(error)),
3229 }
3230 }
3231 Ok(_) => Err(SqliteError::InvalidData(format!(
3232 "read-only WAL snapshot {} has a writable WAL shared-memory sidecar {}; close every \
3233 live writer and remove the transient -shm file (or make a genuinely frozen snapshot) \
3234 before inspection",
3235 path.display(),
3236 shm.display(),
3237 ))),
3238 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3239 let wal = sqlite_sidecar_path(path, "-wal");
3240 match fs::metadata(&wal) {
3241 Ok(metadata) if metadata.len() > 0 => Err(SqliteError::InvalidData(format!(
3242 "read-only WAL snapshot {} has a non-empty WAL sidecar {} but no read-only \
3243 shared-memory sidecar {}; refusing before SQLite open because immutable \
3244 mode would omit committed WAL frames and ordinary read-only mode would \
3245 create or mutate -shm; include the frozen read-only -shm beside this \
3246 snapshot, or checkpoint a writable copy before inspection",
3247 path.display(),
3248 wal.display(),
3249 shm.display(),
3250 ))),
3251 Ok(_) => sqlite_immutable_uri(path),
3252 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3253 sqlite_immutable_uri(path)
3254 }
3255 Err(error) => Err(SqliteError::Io(error)),
3256 }
3257 }
3258 Err(error) => Err(SqliteError::Io(error)),
3259 }
3260}
3261
3262pub(crate) fn open_read_only_snapshot_connection(path: &Path) -> Result<Connection, SqliteError> {
3263 let (_, physical_path) = mint_db_identity(path)?;
3264 let target = read_only_wal_open_target_for_path(&physical_path)?;
3265 let conn = Connection::open_with_flags(&target, reader_open_flags())?;
3266 #[cfg(feature = "namespace-trigram-proto")]
3267 register_namespace_trigram(&conn)?;
3268 Ok(conn)
3269}
3270
3271fn sqlite_header_uses_wal(path: &Path) -> Result<bool, SqliteError> {
3272 let mut file = fs::File::open(path)?;
3273 let mut header = [0_u8; 20];
3274 if let Err(error) = file.read_exact(&mut header) {
3275 if error.kind() == std::io::ErrorKind::UnexpectedEof {
3276 return Ok(false);
3277 }
3278 return Err(SqliteError::Io(error));
3279 }
3280 Ok(&header[..16] == b"SQLite format 3\0" && header[18] == 2 && header[19] == 2)
3281}
3282
3283fn sqlite_sidecar_path(path: &Path, suffix: &str) -> PathBuf {
3284 let mut sidecar = path.as_os_str().to_os_string();
3285 sidecar.push(suffix);
3286 PathBuf::from(sidecar)
3287}
3288
3289fn sqlite_immutable_uri(path: &Path) -> Result<PathBuf, SqliteError> {
3290 let absolute = if path.is_absolute() {
3291 path.to_path_buf()
3292 } else {
3293 std::env::current_dir()?.join(path)
3294 };
3295 let mut uri = String::from("file:");
3296
3297 #[cfg(unix)]
3298 {
3299 use std::os::unix::ffi::OsStrExt as _;
3300 push_sqlite_uri_path(&mut uri, absolute.as_os_str().as_bytes());
3301 }
3302
3303 #[cfg(not(unix))]
3304 {
3305 let path = absolute.to_str().ok_or_else(|| {
3306 SqliteError::InvalidData(format!(
3307 "read-only WAL snapshot path is not representable as a SQLite URI: {}",
3308 absolute.display()
3309 ))
3310 })?;
3311 let normalized = path.replace('\\', "/");
3312 if cfg!(windows) && !normalized.starts_with('/') {
3313 uri.push('/');
3314 }
3315 push_sqlite_uri_path(&mut uri, normalized.as_bytes());
3316 }
3317
3318 uri.push_str("?mode=ro&immutable=1");
3319 Ok(PathBuf::from(uri))
3320}
3321
3322fn push_sqlite_uri_path(uri: &mut String, bytes: &[u8]) {
3323 const HEX: &[u8; 16] = b"0123456789ABCDEF";
3324 for &byte in bytes {
3325 if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~' | b'/') {
3326 uri.push(byte as char);
3327 } else {
3328 uri.push('%');
3329 uri.push(HEX[(byte >> 4) as usize] as char);
3330 uri.push(HEX[(byte & 0x0f) as usize] as char);
3331 }
3332 }
3333}
3334
3335fn open_reader_connection(path: &Path, config: &PoolConfig) -> Result<Connection, SqliteError> {
3336 let conn = claimed_file_identity::open_connection(
3337 config,
3338 path,
3339 reader_open_flags(),
3340 config.path.as_deref(),
3341 )?;
3342 configure_reader_connection(&conn, config)?;
3343 Ok(conn)
3344}
3345
3346fn writer_open_flags() -> OpenFlags {
3347 OpenFlags::SQLITE_OPEN_READ_WRITE
3348 | OpenFlags::SQLITE_OPEN_CREATE
3349 | OpenFlags::SQLITE_OPEN_URI
3350 | OpenFlags::SQLITE_OPEN_NO_MUTEX
3351}
3352
3353fn writer_read_only_open_flags() -> OpenFlags {
3356 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
3357}
3358
3359fn reader_open_flags() -> OpenFlags {
3360 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
3361}
3362
3363#[cfg(feature = "namespace-trigram-proto")]
3364fn register_namespace_trigram(conn: &Connection) -> Result<(), SqliteError> {
3365 crate::namespace_trigram_proto::register(conn).map_err(SqliteError::InvalidData)
3366}
3367
3368fn register_writer_clock(conn: &Connection) -> Result<(), SqliteError> {
3369 conn.create_scalar_function(
3372 "khive_now_micros",
3373 0,
3374 rusqlite::functions::FunctionFlags::SQLITE_UTF8,
3375 |_| Ok(chrono::Utc::now().timestamp_micros()),
3376 )?;
3377 Ok(())
3378}
3379
3380pub(crate) fn rfc3339_instant_key(instant: chrono::DateTime<chrono::Utc>) -> Vec<u8> {
3383 let mut key = Vec::with_capacity(12);
3384 key.extend_from_slice(&((instant.timestamp() as u64) ^ (1_u64 << 63)).to_be_bytes());
3385 key.extend_from_slice(&instant.timestamp_subsec_nanos().to_be_bytes());
3386 key
3387}
3388
3389pub(crate) fn strict_rfc3339_key(text: &str) -> Option<Vec<u8>> {
3393 chrono::DateTime::parse_from_rfc3339(text)
3394 .ok()
3395 .map(|instant| rfc3339_instant_key(instant.with_timezone(&chrono::Utc)))
3396}
3397
3398pub(crate) fn register_rfc3339_key(conn: &Connection) -> rusqlite::Result<()> {
3400 use rusqlite::functions::FunctionFlags;
3401 use rusqlite::types::ValueRef;
3402
3403 conn.create_scalar_function(
3404 "khive_rfc3339_key",
3405 1,
3406 FunctionFlags::SQLITE_UTF8
3407 | FunctionFlags::SQLITE_DETERMINISTIC
3408 | FunctionFlags::SQLITE_INNOCUOUS,
3409 |ctx| {
3410 let text = match ctx.get_raw(0) {
3411 ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
3412 _ => None,
3413 };
3414 let key = text
3415 .and_then(|text| text.parse::<chrono::DateTime<chrono::Utc>>().ok())
3416 .map(rfc3339_instant_key);
3417 Ok(key)
3418 },
3419 )?;
3420 conn.create_scalar_function(
3425 "khive_rfc3339_strict_key",
3426 1,
3427 FunctionFlags::SQLITE_UTF8
3428 | FunctionFlags::SQLITE_DETERMINISTIC
3429 | FunctionFlags::SQLITE_INNOCUOUS,
3430 |ctx| {
3431 let text = match ctx.get_raw(0) {
3432 ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
3433 _ => None,
3434 };
3435 let key = text.and_then(strict_rfc3339_key);
3436 Ok(key)
3437 },
3438 )?;
3439 Ok(())
3440}
3441
3442fn configure_writer_connection(
3443 conn: &Connection,
3444 config: &PoolConfig,
3445) -> Result<bool, SqliteError> {
3446 #[cfg(feature = "namespace-trigram-proto")]
3447 register_namespace_trigram(conn)?;
3448 register_writer_clock(conn)?;
3449 register_rfc3339_key(conn)?;
3450 if config.read_only {
3451 conn.pragma_update(None, "foreign_keys", "ON")?;
3455 conn.busy_timeout(config.busy_timeout)?;
3456 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3457 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3458 conn.pragma_update(None, "temp_store", "MEMORY")?;
3459 conn.pragma_update(None, "query_only", "ON")?;
3460
3461 let wal_enabled =
3462 config.wal_mode && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
3463 return Ok(wal_enabled);
3464 }
3465
3466 let wants_wal = config.path.is_some() && config.wal_mode;
3467
3468 if wants_wal {
3469 conn.pragma_update(None, "journal_mode", "WAL")?;
3470 }
3471 code_map::require_delete_journal(conn, config)?;
3472
3473 conn.pragma_update(None, "synchronous", "NORMAL")?;
3474 conn.pragma_update(None, "foreign_keys", "ON")?;
3475 conn.busy_timeout(config.busy_timeout)?;
3476 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3477 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3478 conn.pragma_update(None, "temp_store", "MEMORY")?;
3479 if config.code_map_vfs.is_none() {
3484 conn.pragma_update(
3485 None,
3486 "wal_autocheckpoint",
3487 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3488 )?;
3489 }
3490
3491 let wal_enabled = wants_wal && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
3492
3493 if wal_enabled {
3494 conn.pragma_update(None, "journal_size_limit", config.journal_size_limit_bytes)?;
3495 }
3496
3497 Ok(wal_enabled)
3498}
3499
3500fn configure_reader_connection(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
3501 #[cfg(feature = "namespace-trigram-proto")]
3502 register_namespace_trigram(conn)?;
3503 register_rfc3339_key(conn)?;
3504 conn.pragma_update(None, "foreign_keys", "ON")?;
3505 conn.busy_timeout(config.busy_timeout)?;
3506 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3507 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3508 conn.pragma_update(None, "temp_store", "MEMORY")?;
3509 Ok(())
3510}
3511
3512fn current_journal_mode(conn: &Connection) -> Result<String, SqliteError> {
3513 conn.pragma_query_value(None, "journal_mode", |row| row.get::<_, String>(0))
3514 .map(|mode| mode.to_ascii_lowercase())
3515 .map_err(Into::into)
3516}
3517
3518fn reset_reader_connection(conn: &Connection, dirty: bool, config: &PoolConfig) -> bool {
3519 if !conn.is_autocommit() {
3520 match conn.execute_batch("ROLLBACK") {
3521 Ok(()) => {}
3522 Err(rusqlite::Error::SqliteFailure(err, _)) => {
3523 if matches!(
3524 err.code,
3525 rusqlite::ErrorCode::CannotOpen
3526 | rusqlite::ErrorCode::DatabaseCorrupt
3527 | rusqlite::ErrorCode::NotADatabase
3528 | rusqlite::ErrorCode::DiskFull
3529 ) {
3530 return false;
3531 }
3532 }
3533 Err(_) => return false,
3534 }
3535 if !conn.is_autocommit() {
3536 return false;
3537 }
3538 }
3539
3540 if !dirty {
3541 return true;
3542 }
3543
3544 reader_connection_state_is_pristine(conn)
3545 && reader_connection_settings_match_baseline(conn, config, 0)
3546}
3547
3548fn reader_connection_state_is_pristine(conn: &Connection) -> bool {
3563 let has_temp_objects: bool = match conn.query_row(
3564 "SELECT EXISTS(SELECT 1 FROM sqlite_temp_master)",
3565 [],
3566 |row| row.get(0),
3567 ) {
3568 Ok(v) => v,
3569 Err(_) => return false,
3570 };
3571 if has_temp_objects {
3572 return false;
3573 }
3574
3575 let attached_databases: i64 = match conn.query_row(
3576 "SELECT COUNT(*) FROM pragma_database_list WHERE name NOT IN ('main', 'temp')",
3577 [],
3578 |row| row.get(0),
3579 ) {
3580 Ok(v) => v,
3581 Err(_) => return false,
3582 };
3583 attached_databases == 0
3584}
3585
3586fn reader_connection_settings_match_baseline(
3600 conn: &Connection,
3601 config: &PoolConfig,
3602 expected_query_only: i64,
3603) -> bool {
3604 let expected_busy_timeout_ms =
3605 i64::try_from(config.busy_timeout.as_millis()).unwrap_or(i64::MAX);
3606 let expected_cache_size: i64 = CACHE_SIZE_KIB.parse().unwrap_or(-65536);
3607 let checks: [(&str, i64); 8] = [
3608 ("query_only", expected_query_only),
3609 ("writable_schema", 0),
3610 ("foreign_keys", 1),
3611 ("busy_timeout", expected_busy_timeout_ms),
3612 ("cache_size", expected_cache_size),
3613 ("temp_store", 2),
3614 ("read_uncommitted", 0),
3615 ("defer_foreign_keys", 0),
3616 ];
3617 checks.iter().all(|(pragma, expected)| {
3618 conn.pragma_query_value(None, pragma, |row| row.get::<_, i64>(0))
3619 .map(|actual| actual == *expected)
3620 .unwrap_or(false)
3621 })
3622}
3623
3624fn restore_shared_reader_state(conn: &Connection, config: &PoolConfig) -> bool {
3633 if !detach_non_main_databases(conn) {
3634 return false;
3635 }
3636 if !drop_temp_objects(conn) {
3637 return false;
3638 }
3639 let expected_query_only = i64::from(config.read_only);
3640 if reset_observable_settings(conn, config, expected_query_only).is_err() {
3641 return false;
3642 }
3643 reader_connection_state_is_pristine(conn)
3644 && reader_connection_settings_match_baseline(conn, config, expected_query_only)
3645}
3646
3647fn detach_non_main_databases(conn: &Connection) -> bool {
3648 loop {
3649 let name: Option<String> = match conn.query_row(
3650 "SELECT name FROM pragma_database_list WHERE name NOT IN ('main', 'temp') LIMIT 1",
3651 [],
3652 |row| row.get(0),
3653 ) {
3654 Ok(name) => Some(name),
3655 Err(rusqlite::Error::QueryReturnedNoRows) => None,
3656 Err(_) => return false,
3657 };
3658 let Some(name) = name else {
3659 return true;
3660 };
3661 let quoted = format!("\"{}\"", name.replace('"', "\"\""));
3662 if conn
3663 .execute_batch(&format!("DETACH DATABASE {quoted}"))
3664 .is_err()
3665 {
3666 return false;
3667 }
3668 }
3669}
3670
3671fn drop_temp_objects(conn: &Connection) -> bool {
3672 for (kind, ddl_keyword) in [
3676 ("view", "VIEW"),
3677 ("trigger", "TRIGGER"),
3678 ("index", "INDEX"),
3679 ("table", "TABLE"),
3680 ] {
3681 loop {
3682 let name: Option<String> = match conn.query_row(
3683 "SELECT name FROM sqlite_temp_master WHERE type = ?1 LIMIT 1",
3684 [kind],
3685 |row| row.get(0),
3686 ) {
3687 Ok(name) => Some(name),
3688 Err(rusqlite::Error::QueryReturnedNoRows) => None,
3689 Err(_) => return false,
3690 };
3691 let Some(name) = name else {
3692 break;
3693 };
3694 let quoted = format!("\"{}\"", name.replace('"', "\"\""));
3695 if conn
3696 .execute_batch(&format!("DROP {ddl_keyword} IF EXISTS temp.{quoted}"))
3697 .is_err()
3698 {
3699 return false;
3700 }
3701 }
3702 }
3703 true
3704}
3705
3706fn reset_observable_settings(
3707 conn: &Connection,
3708 config: &PoolConfig,
3709 expected_query_only: i64,
3710) -> Result<(), rusqlite::Error> {
3711 conn.pragma_update(None, "query_only", expected_query_only)?;
3712 conn.pragma_update(None, "writable_schema", 0)?;
3713 conn.pragma_update(None, "foreign_keys", "ON")?;
3714 conn.busy_timeout(config.busy_timeout)?;
3715 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3716 conn.pragma_update(None, "temp_store", "MEMORY")?;
3717 conn.pragma_update(None, "read_uncommitted", 0)?;
3718 conn.pragma_update(None, "defer_foreign_keys", 0)?;
3719 Ok(())
3720}
3721
3722fn reader_connection_is_healthy(conn: &Connection) -> bool {
3723 match conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0)) {
3724 Ok(_) => true,
3725 Err(rusqlite::Error::SqliteFailure(err, _)) => !matches!(
3726 err.code,
3727 rusqlite::ErrorCode::CannotOpen
3728 | rusqlite::ErrorCode::NotADatabase
3729 | rusqlite::ErrorCode::DatabaseCorrupt
3730 | rusqlite::ErrorCode::PermissionDenied
3731 | rusqlite::ErrorCode::SystemIoFailure
3732 ),
3733 Err(_) => true,
3734 }
3735}
3736
3737fn close_connection_quietly(conn: Connection) {
3738 match conn.close() {
3739 Ok(()) => {}
3740 Err((conn, _)) => drop(conn),
3741 }
3742}
3743
3744fn pool_exhausted_error(timeout: Duration, max_readers: usize) -> SqliteError {
3745 rusqlite::Error::SqliteFailure(
3746 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
3747 Some(format!(
3748 "Pool exhausted: no reader available after {timeout:?} (max_readers={max_readers})"
3749 )),
3750 )
3751 .into()
3752}
3753
3754#[cfg(test)]
3755#[path = "runtime_write_routing_tests.rs"]
3756mod runtime_write_routing_tests;
3757
3758#[cfg(test)]
3759#[path = "database_owner_identity_pool_tests.rs"]
3760mod database_owner_identity_pool_tests;
3761
3762#[cfg(test)]
3763#[path = "pool_tests.rs"]
3764mod tests;
3765
3766#[cfg(all(test, any(unix, windows)))]
3767#[path = "pool_identity_admission_tests.rs"]
3768mod identity_admission_tests;