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#[derive(Clone, Copy, Debug)]
111pub enum RuntimeWriteOperation {
112 MergeEntity,
113 MergeNote,
114 UpdateSymmetricEdge,
115}
116
117impl RuntimeWriteOperation {
118 fn operation(self) -> &'static str {
119 match self {
120 Self::MergeEntity => "merge_entity",
121 Self::MergeNote => "merge_note",
122 Self::UpdateSymmetricEdge => "update_edge",
123 }
124 }
125
126 fn fallback_site(self) -> crate::timeout_sink::Site {
127 match self {
128 Self::MergeEntity => crate::timeout_sink::Site::DirectRouteRuntimeMergeEntity,
129 Self::MergeNote => crate::timeout_sink::Site::DirectRouteRuntimeMergeNote,
130 Self::UpdateSymmetricEdge => {
131 crate::timeout_sink::Site::DirectRouteRuntimeUpdateSymmetricEdge
132 }
133 }
134 }
135}
136
137pub(crate) const FALLBACK_WAL_AUTOCHECKPOINT_PAGES: u32 = 4_000;
145
146#[derive(Clone, Copy, Debug, PartialEq, Eq)]
147enum CheckpointOwnership {
148 Unclaimed,
149 Claiming,
150 Claimed,
151}
152
153struct CheckpointOwnershipState {
154 phase: CheckpointOwnership,
155 #[cfg(test)]
156 connection_waiters: usize,
157}
158
159#[cfg(test)]
160struct CheckpointConnectionConfigPause {
161 selected: std::sync::Barrier,
162 resume: std::sync::Barrier,
163}
164
165#[cfg(test)]
166impl CheckpointConnectionConfigPause {
167 fn new() -> Self {
168 Self {
169 selected: std::sync::Barrier::new(2),
170 resume: std::sync::Barrier::new(2),
171 }
172 }
173}
174
175struct CheckpointOwnershipGate {
176 state: Mutex<CheckpointOwnershipState>,
177 changed: Condvar,
178 #[cfg(test)]
179 connection_config_pause: Mutex<Option<Arc<CheckpointConnectionConfigPause>>>,
180 #[cfg(test)]
181 claim_lock_observed: Mutex<Option<std::sync::mpsc::SyncSender<bool>>>,
182}
183
184impl CheckpointOwnershipGate {
185 fn new() -> Self {
186 Self {
187 state: Mutex::new(CheckpointOwnershipState {
188 phase: CheckpointOwnership::Unclaimed,
189 #[cfg(test)]
190 connection_waiters: 0,
191 }),
192 changed: Condvar::new(),
193 #[cfg(test)]
194 connection_config_pause: Mutex::new(None),
195 #[cfg(test)]
196 claim_lock_observed: Mutex::new(None),
197 }
198 }
199
200 fn begin_claim(&self) -> bool {
203 #[cfg(test)]
204 let claim_lock_observed = self.claim_lock_observed.lock().take();
205 #[cfg(test)]
206 let mut state = if let Some(observed) = claim_lock_observed {
207 match self.state.try_lock() {
208 Some(state) => {
209 let _ = observed.send(false);
210 state
211 }
212 None => {
213 let _ = observed.send(true);
214 self.state.lock()
215 }
216 }
217 } else {
218 self.state.lock()
219 };
220 #[cfg(not(test))]
221 let mut state = self.state.lock();
222 loop {
223 match state.phase {
224 CheckpointOwnership::Unclaimed => {
225 state.phase = CheckpointOwnership::Claiming;
226 self.changed.notify_all();
227 return true;
228 }
229 CheckpointOwnership::Claiming => self.changed.wait(&mut state),
230 CheckpointOwnership::Claimed => return false,
231 }
232 }
233 }
234
235 fn finish_claim(&self, succeeded: bool) {
236 let mut state = self.state.lock();
237 debug_assert_eq!(state.phase, CheckpointOwnership::Claiming);
238 state.phase = if succeeded {
239 CheckpointOwnership::Claimed
240 } else {
241 CheckpointOwnership::Unclaimed
242 };
243 self.changed.notify_all();
244 }
245
246 fn settled_state(&self) -> parking_lot::MutexGuard<'_, CheckpointOwnershipState> {
247 let mut state = self.state.lock();
248 while state.phase == CheckpointOwnership::Claiming {
249 #[cfg(test)]
250 {
251 state.connection_waiters += 1;
252 self.changed.notify_all();
253 }
254 self.changed.wait(&mut state);
255 #[cfg(test)]
256 {
257 state.connection_waiters -= 1;
258 self.changed.notify_all();
259 }
260 }
261 state
262 }
263
264 #[cfg(test)]
265 fn wal_autocheckpoint_pages(&self) -> u32 {
266 let state = self.settled_state();
267 match state.phase {
268 CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
269 CheckpointOwnership::Claimed => 0,
270 CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
271 }
272 }
273
274 fn configure_wal_autocheckpoint(&self, conn: &Connection) -> Result<(), SqliteError> {
279 let state = self.settled_state();
280 let pages = match state.phase {
281 CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
282 CheckpointOwnership::Claimed => 0,
283 CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
284 };
285 #[cfg(test)]
286 if let Some(pause) = self.connection_config_pause.lock().take() {
287 pause.selected.wait();
288 pause.resume.wait();
289 }
290 conn.pragma_update(None, "wal_autocheckpoint", pages)?;
291 drop(state);
292 Ok(())
293 }
294}
295
296fn deny_retired_writer(_context: AuthContext<'_>) -> Authorization {
297 Authorization::Deny
298}
299
300pub(crate) const TEST_HARNESS_ENV: &str = "KHIVE_TEST_HARNESS";
301
302#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, Serialize)]
304#[serde(rename_all = "snake_case")]
305pub enum WalCeilingSource {
306 BackendField,
307 Environment,
308 #[default]
309 Default,
310}
311
312#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
315pub struct WalCeilingPolicy {
316 pub bytes: u64,
317 pub source: WalCeilingSource,
318}
319
320impl WalCeilingPolicy {
321 pub fn validate_static(
325 self,
326 file_backed: bool,
327 wal_mode: bool,
328 read_only: bool,
329 ) -> Result<(), SqliteError> {
330 if self.bytes == 0 {
331 return Ok(());
332 }
333 if i64::try_from(self.bytes).is_err() {
334 return Err(SqliteError::WalCeilingOffsetOverflow { bytes: self.bytes });
335 }
336 if read_only {
337 return Ok(());
338 }
339 if !file_backed {
340 return Err(SqliteError::WalCeilingUnsupported {
341 bytes: self.bytes,
342 backend_kind: "in-memory backend",
343 });
344 }
345 if !wal_mode {
346 return Err(SqliteError::WalCeilingUnsupported {
347 bytes: self.bytes,
348 backend_kind: "non-WAL backend",
349 });
350 }
351 Ok(())
352 }
353
354 pub fn effective_bytes(self, read_only: bool) -> u64 {
356 if read_only {
357 0
358 } else {
359 self.bytes
360 }
361 }
362}
363
364#[derive(Clone, Debug)]
366pub struct PoolConfig {
367 pub path: Option<PathBuf>,
369 pub code_map_vfs: Option<String>,
371 #[cfg(any(unix, windows))]
374 pub expected_file_identity: Option<DatabaseFileIdentity>,
375 pub max_readers: usize,
377 pub wal_mode: bool,
379 pub busy_timeout: Duration,
383 pub checkout_timeout: Duration,
395 pub journal_size_limit_bytes: i64,
401 pub read_only: bool,
409 pub wal_ceiling: WalCeilingPolicy,
412 pub write_queue_enabled: Option<bool>,
440 pub write_queue_capacity: usize,
445 pub write_routing_strict: bool,
456 pub write_admission_deadline_ms: u64,
470 pub disk_guard_config: Option<EffectiveDiskGuardConfig>,
473 pub volume_lock_dir: Option<PathBuf>,
477 pub read_tx_max_age: Duration,
487}
488
489const WRITE_ADMISSION_DEADLINE_MS_RANGE: std::ops::RangeInclusive<u64> = 100..=10_000;
491const DEFAULT_WRITE_ADMISSION_DEADLINE_MS: u64 = 2000;
492
493impl Default for PoolConfig {
494 fn default() -> Self {
495 Self {
496 path: None,
497 code_map_vfs: None,
498 #[cfg(any(unix, windows))]
499 expected_file_identity: None,
500 max_readers: std::thread::available_parallelism()
501 .map(|n| n.get())
502 .unwrap_or(1)
503 .clamp(1, DEFAULT_READER_CAP),
504 wal_mode: true,
505 busy_timeout: Duration::from_secs(
506 std::env::var("KHIVE_BUSY_TIMEOUT_SECS")
507 .ok()
508 .and_then(|v| v.parse::<u64>().ok())
509 .unwrap_or(30),
510 ),
511 checkout_timeout: Duration::from_secs(
512 std::env::var("KHIVE_CHECKOUT_TIMEOUT_SECS")
513 .ok()
514 .and_then(|v| v.parse::<u64>().ok())
515 .unwrap_or(5),
516 ),
517 journal_size_limit_bytes: std::env::var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES")
518 .ok()
519 .and_then(|v| v.parse::<i64>().ok())
520 .unwrap_or(DEFAULT_JOURNAL_SIZE_LIMIT_BYTES),
521 read_only: false,
522 wal_ceiling: WalCeilingPolicy::default(),
523 write_queue_enabled: std::env::var_os("KHIVE_WRITE_QUEUE").map(|v| {
528 v.to_str()
529 .is_some_and(|v| v == "1" || v.eq_ignore_ascii_case("true"))
530 }),
531 write_queue_capacity: std::env::var("KHIVE_WRITE_QUEUE_CAPACITY")
532 .ok()
533 .and_then(|v| v.parse::<usize>().ok())
534 .filter(|&n| n > 0)
535 .unwrap_or(DEFAULT_WRITE_QUEUE_CAPACITY),
536 write_routing_strict: std::env::var("KHIVE_WRITE_ROUTING")
537 .map(|v| v.eq_ignore_ascii_case("strict"))
538 .unwrap_or(false),
539 write_admission_deadline_ms: std::env::var("KHIVE_WRITE_ADMISSION_DEADLINE_MS")
540 .ok()
541 .and_then(|v| v.parse::<u64>().ok())
542 .unwrap_or(DEFAULT_WRITE_ADMISSION_DEADLINE_MS),
543 disk_guard_config: None,
544 #[cfg(test)]
545 volume_lock_dir: Some(test_volume_lock_dir()),
546 #[cfg(not(test))]
547 volume_lock_dir: crate::default_volume_lock_dir().ok(),
548 read_tx_max_age: crate::checkpoint::tx_age_thresholds_from_env(
549 Duration::from_secs(30),
550 Duration::from_secs(120),
551 )
552 .1,
553 }
554 }
555}
556
557#[cfg(any(test, feature = "test-support"))]
558impl PoolConfig {
559 pub fn for_test() -> Self {
564 Self {
565 max_readers: 2,
566 volume_lock_dir: Some(test_volume_lock_dir()),
567 ..Self::default()
568 }
569 }
570}
571
572fn refuse_home_data_store_in_tests(config: &PoolConfig) -> Result<(), SqliteError> {
589 if std::env::var(TEST_HARNESS_ENV).as_deref() != Ok("1") {
590 return Ok(());
591 }
592
593 let Some(path) = config.path.as_deref() else {
594 return Ok(());
595 };
596 if path
597 .as_os_str()
598 .as_encoded_bytes()
599 .get(..5)
600 .is_some_and(|prefix| prefix.eq_ignore_ascii_case(b"file:"))
601 {
602 return Err(SqliteError::InvalidData(format!(
603 "test harness refused SQLite URI database path {}; use a filesystem path outside \
604 HOME/.khive (deliberate sessions against a real store run the built binary \
605 directly, outside the Cargo test environment)",
606 path.display()
607 )));
608 }
609
610 let Some(home) = std::env::var_os("HOME") else {
611 return Ok(());
612 };
613 let canonical_path = canonicalize_deepest_existing(path)?;
614 let canonical_home_data_dir =
615 canonicalize_deepest_existing(&PathBuf::from(home).join(".khive"))?;
616 if canonical_path.starts_with(&canonical_home_data_dir) {
617 return Err(SqliteError::InvalidData(format!(
618 "test harness refused to open SQLite database under HOME/.khive: {} \
619 (deliberate sessions against a real store run the built binary directly, \
620 outside the Cargo test environment)",
621 canonical_path.display()
622 )));
623 }
624 Ok(())
625}
626
627fn canonicalize_deepest_existing(path: &Path) -> Result<PathBuf, SqliteError> {
628 let absolute = if path.is_absolute() {
629 path.to_path_buf()
630 } else {
631 std::env::current_dir().map_err(SqliteError::Io)?.join(path)
632 };
633
634 for ancestor in absolute.ancestors() {
635 match fs::canonicalize(ancestor) {
636 Ok(mut canonical) => {
637 let missing = absolute.strip_prefix(ancestor).map_err(|error| {
638 SqliteError::InvalidData(format!(
639 "failed to preserve missing path components for {}: {error}",
640 absolute.display()
641 ))
642 })?;
643 canonical.push(missing);
644 return Ok(canonical);
645 }
646 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
647 Err(error) => {
648 return Err(SqliteError::InvalidData(format!(
649 "failed to canonicalize database path ancestor {}: {error}",
650 ancestor.display()
651 )));
652 }
653 }
654 }
655
656 Err(SqliteError::InvalidData(format!(
657 "database path has no canonicalizable ancestor: {}",
658 absolute.display()
659 )))
660}
661
662fn validate_write_admission_deadline(deadline_ms: u64) -> Result<(), SqliteError> {
667 if WRITE_ADMISSION_DEADLINE_MS_RANGE.contains(&deadline_ms) {
668 return Ok(());
669 }
670 Err(SqliteError::InvalidConfig(format!(
671 "write_admission_deadline_ms must be in [{}, {}] ms, got {deadline_ms}",
672 WRITE_ADMISSION_DEADLINE_MS_RANGE.start(),
673 WRITE_ADMISSION_DEADLINE_MS_RANGE.end()
674 )))
675}
676
677#[derive(Clone, Debug, Default, PartialEq, Eq, serde::Serialize)]
679pub struct SearchMechanismSnapshot {
680 pub dispatches_by_backend_and_kind: BTreeMap<String, BTreeMap<String, u64>>,
684 pub note_candidate_hydration_rows: u64,
688}
689
690pub struct ConnectionPool {
704 writer: Arc<Mutex<Connection>>,
705 retirement_connection: Mutex<Option<Connection>>,
707 #[cfg(any(test, feature = "test-support"))]
708 statement_observer: Arc<crate::statement_observer::StatementObserverHub>,
709 main_pool_generation: OnceLock<u64>,
710 checkpoint_ownership: CheckpointOwnershipGate,
719 pooled_writer_retired: AtomicBool,
724 writer_acquisition_counters: Arc<WriterAcquisitionCounters>,
729 write_admission: Arc<WriteAdmission>,
732 disk_guard_config: Option<EffectiveDiskGuardConfig>,
735 reader_acquisition_counters: ReaderAcquisitionCounters,
740 search_dispatches: Mutex<BTreeMap<String, BTreeMap<String, u64>>>,
743 note_candidate_hydration_rows: AtomicU64,
744 readers: ArrayQueue<Connection>,
745 max_readers: usize,
746 config: PoolConfig,
747 read_only_open_target: Option<PathBuf>,
755 sql_bridge_reader_slots: Arc<Semaphore>,
756 sql_bridge_writer_slots: Arc<Semaphore>,
757 writer_task: OnceLock<Option<WriterTaskHandle>>,
762 writer_task_join: Mutex<Option<tokio::task::JoinHandle<()>>>,
770 writer_task_join_stored: AtomicBool,
775 origin: TxOrigin,
781 identity_path: Option<PathBuf>,
787 #[cfg(any(unix, windows))]
790 opened_file_identity: Option<DatabaseFileIdentity>,
791 opened_database_id: Option<uuid::Uuid>,
794 identity_registration: Option<PoolIdentityRegistration>,
797 #[cfg(test)]
803 writer_task_spawn_count: std::sync::atomic::AtomicUsize,
804}
805
806impl Drop for ConnectionPool {
807 fn drop(&mut self) {
815 while let Some(conn) = self.readers.pop() {
816 drop(conn);
817 }
818 }
819}
820
821enum ReaderLease<'pool> {
822 Pooled(Connection),
823 Shared(parking_lot::MutexGuard<'pool, Connection>),
824}
825
826pub struct ReaderRow<'row, 'statement> {
836 row: &'row rusqlite::Row<'statement>,
837}
838
839impl ReaderRow<'_, '_> {
840 pub fn get<I: rusqlite::RowIndex, T: rusqlite::types::FromSql>(
842 &self,
843 index: I,
844 ) -> rusqlite::Result<T> {
845 self.row.get(index)
846 }
847
848 pub fn get_ref<I: rusqlite::RowIndex>(
850 &self,
851 index: I,
852 ) -> rusqlite::Result<rusqlite::types::ValueRef<'_>> {
853 self.row.get_ref(index)
854 }
855}
856
857struct ReaderQueryInProgress<'a>(&'a Cell<bool>);
859
860impl Drop for ReaderQueryInProgress<'_> {
861 fn drop(&mut self) {
862 self.0.set(false);
863 }
864}
865
866pub(crate) struct ReaderAdmission {
872 slot: tokio::sync::OwnedSemaphorePermit,
873 started: Instant,
874}
875
876pub struct ReaderGuard<'pool> {
879 lease: Option<ReaderLease<'pool>>,
880 admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
884 pool: &'pool ConnectionPool,
885 reusable: Cell<bool>,
886 query_in_progress: Cell<bool>,
887 checked_out_at: Instant,
888 dirty: Cell<bool>,
896 operation: Option<&'static str>,
902}
903
904impl<'pool> ReaderGuard<'pool> {
905 pub(crate) fn conn(&self) -> &Connection {
914 match self
915 .lease
916 .as_ref()
917 .expect("reader guard missing connection")
918 {
919 ReaderLease::Pooled(conn) => conn,
920 ReaderLease::Shared(guard) => guard,
921 }
922 }
923
924 pub fn query_row<T, P, F>(&self, sql: &str, params: P, f: F) -> Result<T, SqliteError>
943 where
944 P: rusqlite::Params,
945 F: FnOnce(&ReaderRow<'_, '_>) -> rusqlite::Result<T>,
946 {
947 crate::sql_bridge::reader_capability_admits(sql).map_err(SqliteError::InvalidData)?;
948 if !self.reusable.get() {
949 return Err(SqliteError::InvalidData(
950 "reader lease is quarantined after failed read cleanup".into(),
951 ));
952 }
953 if self.query_in_progress.replace(true) {
957 return Err(SqliteError::InvalidData(
958 "reader lease is already executing a query".into(),
959 ));
960 }
961 let _in_progress = ReaderQueryInProgress(&self.query_in_progress);
962 self.mark_dirty();
963 crate::read_cancellation::run_borrowed_reader(self, |conn, admission| {
964 conn.query_row(sql, params, |row| {
965 if !admission.admits() {
973 return Err(rusqlite::Error::SqliteFailure(
974 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_INTERRUPT),
975 Some("request stopped before the mapper".into()),
976 ));
977 }
978 f(&ReaderRow { row })
979 })
980 .map_err(|error| {
981 StorageError::driver(StorageCapability::Sql, "reader_guard.query_row", error)
982 })
983 })
984 .map_err(|error| {
985 self.pool.record_reader_query_error(&error);
986 match error {
987 StorageError::Driver {
988 capability,
989 operation,
990 source,
991 } => match source.downcast::<rusqlite::Error>() {
992 Ok(error) => SqliteError::Rusqlite(*error),
993 Err(source) => SqliteError::RequestReadStopped(StorageError::Driver {
994 capability,
995 operation,
996 source,
997 }),
998 },
999 other => SqliteError::RequestReadStopped(other),
1000 }
1001 })
1002 }
1003
1004 pub(crate) fn discard(&self) {
1008 self.reusable.set(false);
1009 }
1010
1011 pub(crate) fn mark_dirty(&self) {
1014 self.dirty.set(true);
1015 }
1016
1017 pub(crate) fn label_operation(&mut self, operation: &'static str) {
1021 self.operation = Some(operation);
1022 }
1023}
1024
1025impl<'pool> Drop for ReaderGuard<'pool> {
1026 fn drop(&mut self) {
1027 let Some(lease) = self.lease.take() else {
1028 return;
1029 };
1030
1031 match lease {
1032 ReaderLease::Pooled(conn) if self.reusable.get() => {
1033 self.pool.return_reader(conn, self.dirty.get())
1034 }
1035 ReaderLease::Pooled(conn) => {
1036 close_connection_quietly(conn);
1037 self.pool.replace_discarded_reader_slot();
1038 }
1039 ReaderLease::Shared(guard) if !self.reusable.get() => {
1040 self.pool.retire_pooled_writer(&guard);
1041 }
1042 ReaderLease::Shared(guard) => {
1043 if self.dirty.get() && !restore_shared_reader_state(&guard, &self.pool.config) {
1051 self.pool.retire_pooled_writer(&guard);
1052 }
1053 }
1054 }
1055
1056 drop(self.admission_slot.take());
1060 self.pool
1061 .reader_acquisition_counters
1062 .record_checkout_completed(self.checked_out_at.elapsed(), self.operation);
1063 }
1064}
1065
1066pub(crate) struct SharedReaderTransactionGuard {
1082 conn: parking_lot::ArcMutexGuard<parking_lot::RawMutex, Connection>,
1083 admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
1084 pool: Arc<ConnectionPool>,
1085 checked_out_at: Instant,
1086 poison: Cell<bool>,
1092}
1093
1094impl SharedReaderTransactionGuard {
1095 pub(crate) fn conn(&self) -> &Connection {
1096 &self.conn
1097 }
1098
1099 pub(crate) fn poison(&self) {
1102 self.poison.set(true);
1103 }
1104}
1105
1106impl Drop for SharedReaderTransactionGuard {
1107 fn drop(&mut self) {
1108 let mut restored = !self.poison.get();
1115 if restored && !self.conn.is_autocommit() {
1116 restored = self.conn.execute_batch("ROLLBACK").is_ok() && self.conn.is_autocommit();
1117 }
1118 if restored {
1119 restored = restore_shared_reader_state(&self.conn, &self.pool.config);
1120 }
1121 if !restored {
1122 self.pool.retire_pooled_writer(&self.conn);
1123 }
1124 drop(self.admission_slot.take());
1125 self.pool
1126 .reader_acquisition_counters
1127 .record_checkout_completed(
1128 self.checked_out_at.elapsed(),
1129 Some("explicit_sql_read_transaction"),
1130 );
1131 }
1132}
1133
1134impl ConnectionPool {
1135 pub(crate) fn checkout_shared_reader_transaction(
1146 self: &Arc<Self>,
1147 should_stop: impl Fn() -> bool,
1148 ) -> Result<Option<SharedReaderTransactionGuard>, SqliteError> {
1149 debug_assert_eq!(
1150 self.max_readers, 0,
1151 "the owned shared-reader-transaction guard exists only for the degraded, \
1152 single-connection backend"
1153 );
1154 self.ensure_pooled_writer_active()?;
1155 let started = Instant::now();
1156 let admission_slot = loop {
1157 if should_stop() {
1158 return Ok(None);
1159 }
1160 match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
1161 Ok(slot) => break slot,
1162 Err(tokio::sync::TryAcquireError::Closed) => {
1163 return Err(SqliteError::InvalidData(
1164 "reader admission semaphore is closed".to_string(),
1165 ));
1166 }
1167 Err(tokio::sync::TryAcquireError::NoPermits) => {}
1168 }
1169 if started.elapsed() >= self.config.checkout_timeout {
1170 self.reader_acquisition_counters.record_checkout_timeout();
1171 return Err(pool_exhausted_error(
1172 self.config.checkout_timeout,
1173 self.max_readers,
1174 ));
1175 }
1176 thread::yield_now();
1177 };
1178
1179 loop {
1180 if should_stop() {
1181 return Ok(None);
1182 }
1183 let remaining = self
1184 .config
1185 .checkout_timeout
1186 .saturating_sub(started.elapsed());
1187 if remaining.is_zero() {
1188 self.reader_acquisition_counters.record_checkout_timeout();
1189 return Err(pool_exhausted_error(
1190 self.config.checkout_timeout,
1191 self.max_readers,
1192 ));
1193 }
1194 if let Some(conn) = self
1195 .writer
1196 .try_lock_arc_for(remaining.min(Duration::from_millis(2)))
1197 {
1198 self.ensure_pooled_writer_active()?;
1199 self.reader_acquisition_counters.record_pooled_checkout();
1200 return Ok(Some(SharedReaderTransactionGuard {
1201 conn,
1202 admission_slot: Some(admission_slot),
1203 pool: Arc::clone(self),
1204 checked_out_at: Instant::now(),
1205 poison: Cell::new(false),
1206 }));
1207 }
1208 }
1209 }
1210}
1211
1212#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1218#[allow(dead_code)] pub(crate) enum StandaloneReaderPurpose {
1220 ExplicitSqlReadTransaction,
1222 BootSchemaProbe,
1225 DiagnosticsIndependentSnapshot,
1228}
1229
1230impl StandaloneReaderPurpose {
1231 fn is_infrastructure(self) -> bool {
1232 matches!(
1233 self,
1234 Self::BootSchemaProbe | Self::DiagnosticsIndependentSnapshot
1235 )
1236 }
1237}
1238
1239#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
1247pub struct ReaderAcquisitionSnapshot {
1248 pub reader_admission_capacity: usize,
1251 pub available_reader_admission_slots: usize,
1254 pub acquisitions: u64,
1257 pub pooled_checkouts: u64,
1259 pub standalone_opens: u64,
1262 pub infrastructure_standalone_opens: u64,
1265 pub checkout_timeouts: u64,
1269 pub busy_timeouts: u64,
1277 pub active_pooled_checkouts: u64,
1279 pub peak_active_pooled_checkouts: u64,
1281 pub completed_pooled_checkouts: u64,
1283 pub max_completed_hold_micros: u64,
1286 pub max_completed_hold_operation: Option<&'static str>,
1291 pub reader_replacement_open_failures: u64,
1297}
1298
1299#[derive(Debug, Default, Clone, Copy)]
1303struct LongestCompletedHold {
1304 micros: u64,
1305 operation: Option<&'static str>,
1306}
1307
1308#[derive(Debug, Default)]
1309struct ReaderAcquisitionCounters {
1310 pooled_checkouts: AtomicU64,
1311 standalone_opens: AtomicU64,
1312 infrastructure_standalone_opens: AtomicU64,
1313 checkout_timeouts: AtomicU64,
1314 busy_timeouts: AtomicU64,
1315 active_pooled_checkouts: AtomicU64,
1316 peak_active_pooled_checkouts: AtomicU64,
1317 completed_pooled_checkouts: AtomicU64,
1318 longest_completed_hold: parking_lot::Mutex<LongestCompletedHold>,
1319 reader_replacement_open_failures: AtomicU64,
1320}
1321
1322impl ReaderAcquisitionCounters {
1323 fn record_pooled_checkout(&self) {
1324 self.pooled_checkouts.fetch_add(1, Ordering::Relaxed);
1325 let active = self
1326 .active_pooled_checkouts
1327 .fetch_add(1, Ordering::Relaxed)
1328 .saturating_add(1);
1329 self.peak_active_pooled_checkouts
1330 .fetch_max(active, Ordering::Relaxed);
1331 }
1332
1333 fn record_checkout_timeout(&self) {
1334 self.checkout_timeouts.fetch_add(1, Ordering::Relaxed);
1335 }
1336
1337 fn record_busy_timeout(&self) {
1338 self.busy_timeouts.fetch_add(1, Ordering::Relaxed);
1339 }
1340
1341 fn record_reader_replacement_open_failure(&self) {
1342 self.reader_replacement_open_failures
1343 .fetch_add(1, Ordering::Relaxed);
1344 }
1345
1346 fn record_standalone_open(&self, purpose: StandaloneReaderPurpose) {
1347 if purpose.is_infrastructure() {
1348 self.infrastructure_standalone_opens
1349 .fetch_add(1, Ordering::Relaxed);
1350 } else {
1351 self.standalone_opens.fetch_add(1, Ordering::Relaxed);
1352 }
1353 }
1354
1355 fn record_checkout_completed(&self, hold: Duration, operation: Option<&'static str>) {
1356 let previous = self.active_pooled_checkouts.fetch_sub(1, Ordering::Relaxed);
1357 debug_assert!(previous > 0, "reader active-checkout counter underflow");
1358 self.completed_pooled_checkouts
1359 .fetch_add(1, Ordering::Relaxed);
1360 let micros = u64::try_from(hold.as_micros()).unwrap_or(u64::MAX);
1361 let mut longest = self.longest_completed_hold.lock();
1365 if micros > longest.micros {
1366 longest.micros = micros;
1367 longest.operation = operation;
1368 }
1369 }
1370
1371 fn snapshot(
1372 &self,
1373 reader_admission_capacity: usize,
1374 available_reader_admission_slots: usize,
1375 ) -> ReaderAcquisitionSnapshot {
1376 let pooled_checkouts = self.pooled_checkouts.load(Ordering::Relaxed);
1377 let standalone_opens = self.standalone_opens.load(Ordering::Relaxed);
1378 let longest_completed_hold = *self.longest_completed_hold.lock();
1379 ReaderAcquisitionSnapshot {
1380 reader_admission_capacity,
1381 available_reader_admission_slots,
1382 acquisitions: pooled_checkouts.saturating_add(standalone_opens),
1383 pooled_checkouts,
1384 standalone_opens,
1385 infrastructure_standalone_opens: self
1386 .infrastructure_standalone_opens
1387 .load(Ordering::Relaxed),
1388 checkout_timeouts: self.checkout_timeouts.load(Ordering::Relaxed),
1389 busy_timeouts: self.busy_timeouts.load(Ordering::Relaxed),
1390 active_pooled_checkouts: self.active_pooled_checkouts.load(Ordering::Relaxed),
1391 peak_active_pooled_checkouts: self.peak_active_pooled_checkouts.load(Ordering::Relaxed),
1392 completed_pooled_checkouts: self.completed_pooled_checkouts.load(Ordering::Relaxed),
1393 max_completed_hold_micros: longest_completed_hold.micros,
1394 max_completed_hold_operation: longest_completed_hold.operation,
1395 reader_replacement_open_failures: self
1396 .reader_replacement_open_failures
1397 .load(Ordering::Relaxed),
1398 }
1399 }
1400}
1401
1402impl ConnectionPool {
1403 pub fn new(config: PoolConfig) -> Result<Self, SqliteError> {
1411 refuse_home_data_store_in_tests(&config)?;
1412 validate_write_admission_deadline(config.write_admission_deadline_ms)?;
1413 if let Some(policy) = config.disk_guard_config {
1414 policy.validate()?;
1415 }
1416 config.wal_ceiling.validate_static(
1417 config.path.is_some(),
1418 config.wal_mode,
1419 config.read_only,
1420 )?;
1421 code_map::validate_pool(&config)?;
1422 if config.path.is_some() && !config.read_only {
1423 let lock_dir =
1424 crate::disk_guard_config::require_volume_lock_dir(config.volume_lock_dir.clone())?;
1425 if !lock_dir.is_absolute() {
1426 return Err(SqliteError::CapacityUnavailable {
1427 phase: CapacityUnavailablePhase::Lock,
1428 message: "configured volume-lock directory is not absolute".to_string(),
1429 });
1430 }
1431 }
1432
1433 let mut config = config;
1437 let inert_memory_queue_request =
1438 config.path.is_none() && config.write_queue_enabled == Some(true);
1439 config.write_queue_enabled =
1440 Some(config.write_queue_enabled.unwrap_or(config.path.is_some()));
1441 if inert_memory_queue_request {
1442 tracing::warn!(
1443 "write queue explicitly requested for an in-memory pool; it is inert because \
1444 in-memory pools cannot host a writer task"
1445 );
1446 }
1447
1448 let (origin, identity_path) = match config.path.as_ref() {
1453 Some(path) if config.code_map_vfs.is_some() => code_map::guarded_identity(path),
1454 Some(path) => {
1455 let (identity, canonical) = mint_db_identity(path)?;
1456 (TxOrigin::Database(identity), Some(canonical))
1457 }
1458 None => (TxOrigin::Memory, None),
1459 };
1460 let (write_admission, disk_guard_config) = if !config.read_only {
1461 if identity_path.as_deref().and_then(Path::parent).is_some() {
1462 let disk_guard = match config.disk_guard_config {
1463 Some(policy) => policy,
1464 None => resolve_disk_guard_config(None, None)?,
1465 };
1466 disk_guard.validate()?;
1467 if disk_guard.legacy_environment_present {
1468 tracing::warn!(
1469 "legacy SQLite reserve setting is deprecated; \
1470 use KHIVE_SQLITE_DISK_RESERVE_BYTES"
1471 );
1472 }
1473 if disk_guard.reserve_bytes == 0 {
1474 tracing::warn!(
1475 "SQLite disk reserve is explicitly zero; \
1476 new logical writes will not be floor-refused"
1477 );
1478 }
1479 (
1480 Arc::new(WriteAdmission::new(
1481 identity_path.clone(),
1482 disk_guard.reserve_bytes,
1483 disk_guard.guard_deadline_ms,
1484 config.volume_lock_dir.clone(),
1485 )?),
1486 Some(disk_guard),
1487 )
1488 } else {
1489 (
1490 Arc::new(WriteAdmission::new(
1491 None,
1492 0,
1493 DEFAULT_DISK_GUARD_DEADLINE_MS,
1494 None,
1495 )?),
1496 None,
1497 )
1498 }
1499 } else {
1500 (
1501 Arc::new(WriteAdmission::new(
1502 None,
1503 0,
1504 DEFAULT_DISK_GUARD_DEADLINE_MS,
1505 None,
1506 )?),
1507 None,
1508 )
1509 };
1510 let read_only_open_target = read_only_open_target(&config, identity_path.as_deref())?;
1511 #[cfg(any(unix, windows))]
1512 let identity_before_open = identity_path
1513 .as_deref()
1514 .map(database_file_identity_if_exists)
1515 .transpose()?
1516 .flatten();
1517 #[cfg(any(unix, windows))]
1518 claimed_file_identity::verify_before_open(&config, identity_before_open)?;
1519 let retirement_connection = Connection::open_in_memory()?;
1520 retirement_connection.authorizer(Some(deny_retired_writer))?;
1521 let mut initialization_lease = None;
1523 let mut writer = open_writer_connection(
1524 &config,
1525 read_only_open_target.as_deref(),
1526 identity_path.as_deref(),
1527 )?;
1528 validate_wal_ceiling_at_open(&writer, &config)?;
1529 writer.busy_timeout(config.busy_timeout)?;
1532 let initial_database_id = if identity_path.is_some() {
1536 read_database_id(&writer)?
1537 } else {
1538 None
1539 };
1540 #[cfg(test)]
1541 if let Some(path) = identity_path.as_deref() {
1542 run_identity_open_hook(
1543 path,
1544 IdentityOpenStage::AfterMainOpenBeforeFirstStat,
1545 Some(&writer),
1546 );
1547 }
1548 #[cfg(any(unix, windows))]
1549 let identity_before_write = identity_path
1550 .as_deref()
1551 .map(database_file_identity)
1552 .transpose()?;
1553 #[cfg(any(unix, windows))]
1554 if identity_before_open.is_some() && identity_before_open != identity_before_write {
1555 return Err(SqliteError::InvalidData(
1556 "database file identity changed while opening the pool".to_string(),
1557 ));
1558 }
1559 #[cfg(any(unix, windows))]
1560 if let Some(path) = identity_path.as_deref() {
1561 let opened = opened_sqlite_file_identity(&writer, path)?;
1562 if identity_before_write != Some(opened) {
1563 return Err(SqliteError::InvalidData(
1564 "database file identity changed while opening the pool".to_string(),
1565 ));
1566 }
1567 }
1568 write_admission.verify_current_volume()?;
1571 let opened_database_id = if identity_path.is_some()
1572 && !config.read_only
1573 && initial_database_id.is_none()
1574 {
1575 match write_admission.check() {
1576 Ok(()) => {
1577 initialization_lease = write_admission.acquire()?;
1578 match initialize_database_id_with_admission(&mut writer, Some(&write_admission))
1579 {
1580 Ok(id) => Some(id),
1581 Err(SqliteError::CapacityFloor { .. }) => None,
1582 Err(error) => return Err(error),
1583 }
1584 }
1585 Err(SqliteError::CapacityFloor { .. }) => None,
1589 Err(error) => return Err(error),
1590 }
1591 } else {
1592 initial_database_id
1593 };
1594 drop(initialization_lease.take());
1595 #[cfg(test)]
1596 if let Some(path) = identity_path.as_deref() {
1597 run_identity_open_hook(
1598 path,
1599 IdentityOpenStage::AfterInitialIdentityWrite,
1600 Some(&writer),
1601 );
1602 }
1603 #[cfg(any(unix, windows))]
1604 let opened_file_identity = identity_path
1605 .as_deref()
1606 .map(|path| opened_sqlite_file_identity(&writer, path))
1607 .transpose()?;
1608 #[cfg(any(unix, windows))]
1609 if identity_before_write != opened_file_identity {
1610 return Err(SqliteError::InvalidData(
1611 "database file identity changed while opening the pool".to_string(),
1612 ));
1613 }
1614 let wal_enabled = configure_writer_connection(&writer, &config)?;
1615 let max_readers = effective_reader_count(&config, wal_enabled);
1616
1617 let readers = ArrayQueue::new(max_readers.max(1));
1618
1619 #[cfg(any(test, feature = "test-support"))]
1620 let statement_observer = crate::statement_observer::StatementObserverHub::new()?;
1621 #[cfg(any(test, feature = "test-support"))]
1622 crate::statement_observer::install(&writer, &statement_observer)?;
1623
1624 let mut pool = Self {
1625 writer: Arc::new(Mutex::new(writer)),
1626 retirement_connection: Mutex::new(Some(retirement_connection)),
1627 #[cfg(any(test, feature = "test-support"))]
1628 statement_observer,
1629 main_pool_generation: OnceLock::new(),
1630 checkpoint_ownership: CheckpointOwnershipGate::new(),
1631 pooled_writer_retired: AtomicBool::new(false),
1632 writer_acquisition_counters: Arc::new(WriterAcquisitionCounters::default()),
1633 write_admission,
1634 disk_guard_config,
1635 reader_acquisition_counters: ReaderAcquisitionCounters::default(),
1636 search_dispatches: Mutex::new(BTreeMap::new()),
1637 note_candidate_hydration_rows: AtomicU64::new(0),
1638 readers,
1639 max_readers,
1640 config,
1641 read_only_open_target,
1642 sql_bridge_reader_slots: Arc::new(Semaphore::new(max_readers.max(1))),
1643 sql_bridge_writer_slots: Arc::new(Semaphore::new(1)),
1644 writer_task: OnceLock::new(),
1645 writer_task_join: Mutex::new(None),
1646 writer_task_join_stored: AtomicBool::new(false),
1647 origin,
1648 identity_path,
1649 #[cfg(any(unix, windows))]
1650 opened_file_identity,
1651 opened_database_id,
1652 identity_registration: None,
1653 #[cfg(test)]
1654 writer_task_spawn_count: std::sync::atomic::AtomicUsize::new(0),
1655 };
1656
1657 for _ in 0..pool.max_readers {
1658 let conn = pool.open_reader_connection()?;
1659 pool.readers
1660 .push(conn)
1661 .expect("reader queue must have capacity during pool initialization");
1662 }
1663
1664 if !pool.config.read_only {
1669 crate::timeout_sink::init(
1670 pool.canonical_path().and_then(Path::parent),
1671 &crate::timeout_sink::db_label(&pool),
1672 );
1673 }
1674
1675 pool.identity_registration = pool.canonical_path().map(PoolIdentityRegistration::new);
1676 Ok(pool)
1677 }
1678
1679 pub fn reader(&self) -> Result<ReaderGuard<'_>, SqliteError> {
1689 self.reader_until(|| false)?.ok_or_else(|| {
1690 SqliteError::InvalidData("uncancelled reader checkout stopped unexpectedly".into())
1691 })
1692 }
1693
1694 pub(crate) fn reader_until<C>(
1700 &self,
1701 should_stop: C,
1702 ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
1703 where
1704 C: Fn() -> bool,
1705 {
1706 let started = Instant::now();
1707 let mut admission_attempt = 0u32;
1708 let admission_slot = loop {
1709 if should_stop() {
1710 return Ok(None);
1711 }
1712 match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
1713 Ok(slot) => break slot,
1714 Err(tokio::sync::TryAcquireError::Closed) => {
1715 return Err(SqliteError::InvalidData(
1716 "reader admission semaphore is closed".to_string(),
1717 ));
1718 }
1719 Err(tokio::sync::TryAcquireError::NoPermits) => {}
1720 }
1721 if started.elapsed() >= self.config.checkout_timeout {
1722 self.reader_acquisition_counters.record_checkout_timeout();
1723 return Err(pool_exhausted_error(
1724 self.config.checkout_timeout,
1725 self.max_readers,
1726 ));
1727 }
1728 match admission_attempt {
1729 0..=7 => {
1730 let spins = 1usize << admission_attempt;
1731 for _ in 0..spins {
1732 std::hint::spin_loop();
1733 }
1734 }
1735 8..=15 => thread::yield_now(),
1736 _ => {
1737 let remaining = self
1738 .config
1739 .checkout_timeout
1740 .saturating_sub(started.elapsed());
1741 let sleep =
1742 Duration::from_micros(50 * (1u64 << (admission_attempt - 16).min(6)));
1743 thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
1744 }
1745 }
1746 admission_attempt = admission_attempt.saturating_add(1);
1747 };
1748
1749 self.reader_with_admission(
1750 ReaderAdmission {
1751 slot: admission_slot,
1752 started,
1753 },
1754 should_stop,
1755 )
1756 }
1757
1758 pub(crate) async fn acquire_reader_admission(
1767 &self,
1768 capability: StorageCapability,
1769 operation: &'static str,
1770 ) -> Result<ReaderAdmission, StorageError> {
1771 let context = khive_storage::capture_request_read_context();
1772 let started = Instant::now();
1773 let stopped = context.stop_reason().is_some();
1774 let outcome: Result<Option<ReaderAdmission>, SqliteError> = if stopped {
1775 Ok(None)
1776 } else {
1777 tokio::select! {
1778 biased;
1779 _ = context.wait_for_stop() => Ok(None),
1780 waited = tokio::time::timeout(
1781 self.config.checkout_timeout,
1782 Arc::clone(&self.sql_bridge_reader_slots).acquire_owned(),
1783 ) => match waited {
1784 Ok(Ok(slot)) => Ok(Some(ReaderAdmission { slot, started })),
1785 Ok(Err(_closed)) => Err(SqliteError::InvalidData(
1786 "reader admission semaphore is closed".to_string(),
1787 )),
1788 Err(_elapsed) => {
1789 self.reader_acquisition_counters.record_checkout_timeout();
1790 Err(pool_exhausted_error(
1791 self.config.checkout_timeout,
1792 self.max_readers,
1793 ))
1794 }
1795 },
1796 }
1797 };
1798 match outcome {
1799 Ok(Some(admission)) => Ok(admission),
1800 Ok(None) => Err(self.reader_checkout_refusal(capability, operation, None)),
1801 Err(error) => Err(self.reader_checkout_refusal(capability, operation, Some(error))),
1802 }
1803 }
1804
1805 pub(crate) fn reader_with_admission<C>(
1812 &self,
1813 admission: ReaderAdmission,
1814 should_stop: C,
1815 ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
1816 where
1817 C: Fn() -> bool,
1818 {
1819 let ReaderAdmission {
1820 slot: admission_slot,
1821 started,
1822 } = admission;
1823
1824 if self.max_readers == 0 {
1825 self.ensure_pooled_writer_active()?;
1826 loop {
1827 if should_stop() {
1828 return Ok(None);
1829 }
1830 let remaining = self
1831 .config
1832 .checkout_timeout
1833 .saturating_sub(started.elapsed());
1834 if remaining.is_zero() {
1835 self.reader_acquisition_counters.record_checkout_timeout();
1836 return Err(pool_exhausted_error(
1837 self.config.checkout_timeout,
1838 self.max_readers,
1839 ));
1840 }
1841 if let Some(guard) = self
1842 .writer
1843 .try_lock_for(remaining.min(Duration::from_millis(2)))
1844 {
1845 self.ensure_pooled_writer_active()?;
1846 self.reader_acquisition_counters.record_pooled_checkout();
1847 return Ok(Some(ReaderGuard {
1848 lease: Some(ReaderLease::Shared(guard)),
1849 admission_slot: Some(admission_slot),
1850 pool: self,
1851 reusable: Cell::new(true),
1852 query_in_progress: Cell::new(false),
1853 checked_out_at: Instant::now(),
1854 dirty: Cell::new(false),
1855 operation: None,
1856 }));
1857 }
1858 }
1859 }
1860
1861 let mut attempt = 0u32;
1862
1863 loop {
1864 if should_stop() {
1865 return Ok(None);
1866 }
1867 if let Some(conn) = self.readers.pop() {
1868 self.reader_acquisition_counters.record_pooled_checkout();
1869 return Ok(Some(ReaderGuard {
1870 lease: Some(ReaderLease::Pooled(conn)),
1871 admission_slot: Some(admission_slot),
1872 pool: self,
1873 reusable: Cell::new(true),
1874 query_in_progress: Cell::new(false),
1875 checked_out_at: Instant::now(),
1876 dirty: Cell::new(false),
1877 operation: None,
1878 }));
1879 }
1880
1881 if started.elapsed() >= self.config.checkout_timeout {
1882 self.reader_acquisition_counters.record_checkout_timeout();
1883 return Err(pool_exhausted_error(
1884 self.config.checkout_timeout,
1885 self.max_readers,
1886 ));
1887 }
1888
1889 match attempt {
1890 0..=7 => {
1891 let spins = 1usize << attempt;
1892 for _ in 0..spins {
1893 std::hint::spin_loop();
1894 }
1895 }
1896 8..=15 => thread::yield_now(),
1897 _ => {
1898 let remaining = self
1899 .config
1900 .checkout_timeout
1901 .saturating_sub(started.elapsed());
1902 let sleep = Duration::from_micros(50 * (1u64 << (attempt - 16).min(6)));
1903 thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
1904 }
1905 }
1906
1907 attempt = attempt.saturating_add(1);
1908 }
1909 }
1910
1911 #[track_caller]
1917 pub fn writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
1918 self.writer_with_checkout_probe(true, false)
1919 }
1920
1921 #[track_caller]
1924 pub(crate) fn writer_for_admitted_operation(&self) -> Result<WriterGuard<'_>, SqliteError> {
1925 self.writer_with_checkout_probe(false, false)
1926 }
1927
1928 fn writer_for_checkpoint_operation(&self) -> Result<WriterGuard<'_>, SqliteError> {
1931 self.writer_with_checkout_probe(false, true)
1932 }
1933
1934 #[track_caller]
1935 fn writer_with_checkout_probe(
1936 &self,
1937 compatibility_probe: bool,
1938 checkpoint_bypass: bool,
1939 ) -> Result<WriterGuard<'_>, SqliteError> {
1940 if !checkpoint_bypass {
1941 self.write_admission.ensure_settled()?;
1942 }
1943 self.ensure_pooled_writer_active()?;
1944 let volume_lease = if checkpoint_bypass {
1945 None
1946 } else {
1947 self.write_admission.acquire()?
1948 };
1949 let Some(guard) = self.writer.try_lock_for(self.config.checkout_timeout) else {
1950 self.writer_acquisition_counters
1951 .pooled_timeouts
1952 .fetch_add(1, Ordering::Relaxed);
1953 let message = format!(
1954 "timed out after {:?} waiting for sqlite writer connection",
1955 self.config.checkout_timeout
1956 );
1957 crate::timeout_sink::emit_timeout(
1958 &crate::timeout_sink::db_label(self),
1959 crate::timeout_sink::Site::PoolAdmission,
1960 &message,
1961 Some(
1962 self.config
1963 .checkout_timeout
1964 .as_millis()
1965 .min(u128::from(u64::MAX)) as u64,
1966 ),
1967 );
1968 return Err(SqliteError::WriterPoolCheckoutTimeout {
1969 timeout: self.config.checkout_timeout,
1970 });
1971 };
1972 self.ensure_pooled_writer_active()?;
1973 #[cfg(any(unix, windows))]
1974 if let Some(path) = self.identity_path.as_deref() {
1975 self.verify_opened_file_identity(path)?;
1976 self.verify_connection_file_identity(&guard, path)?;
1977 }
1978 if compatibility_probe && !checkpoint_bypass {
1982 self.write_admission.check()?;
1983 }
1984 self.writer_acquisition_counters
1985 .pooled_acquisitions
1986 .fetch_add(1, Ordering::Relaxed);
1987 Ok(WriterGuard {
1988 guard,
1989 origin: self.origin(),
1990 pool: self,
1991 admission: self.write_admission.as_ref(),
1992 _volume_lease: volume_lease,
1993 })
1994 }
1995
1996 #[track_caller]
2001 pub fn try_writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
2002 self.writer()
2003 }
2004
2005 #[cfg(test)]
2008 #[track_caller]
2009 fn writer_until<C>(&self, should_stop: C) -> Result<Option<WriterGuard<'_>>, SqliteError>
2010 where
2011 C: Fn() -> bool,
2012 {
2013 self.writer_until_with_checkout_probe(should_stop, true)
2014 }
2015
2016 #[track_caller]
2017 pub(crate) fn writer_until_for_admitted_operation<C>(
2018 &self,
2019 should_stop: C,
2020 ) -> Result<Option<WriterGuard<'_>>, SqliteError>
2021 where
2022 C: Fn() -> bool,
2023 {
2024 self.writer_until_with_checkout_probe(should_stop, false)
2025 }
2026
2027 #[track_caller]
2028 fn writer_until_with_checkout_probe<C>(
2029 &self,
2030 should_stop: C,
2031 compatibility_probe: bool,
2032 ) -> Result<Option<WriterGuard<'_>>, SqliteError>
2033 where
2034 C: Fn() -> bool,
2035 {
2036 self.ensure_pooled_writer_active()?;
2037 if should_stop() {
2038 return Ok(None);
2039 }
2040 let volume_lease = self.write_admission.acquire()?;
2041 let started = Instant::now();
2042 loop {
2043 if should_stop() {
2044 return Ok(None);
2045 }
2046 let remaining = self
2047 .config
2048 .checkout_timeout
2049 .saturating_sub(started.elapsed());
2050 if let Some(guard) = self
2051 .writer
2052 .try_lock_for(remaining.min(Duration::from_millis(2)))
2053 {
2054 if should_stop() {
2057 return Ok(None);
2058 }
2059 self.ensure_pooled_writer_active()?;
2060 #[cfg(any(unix, windows))]
2061 if let Some(path) = self.identity_path.as_deref() {
2062 self.verify_opened_file_identity(path)?;
2063 self.verify_connection_file_identity(&guard, path)?;
2064 }
2065 if compatibility_probe {
2066 self.write_admission.check()?;
2067 }
2068 self.writer_acquisition_counters
2069 .pooled_acquisitions
2070 .fetch_add(1, Ordering::Relaxed);
2071 return Ok(Some(WriterGuard {
2072 guard,
2073 origin: self.origin(),
2074 pool: self,
2075 admission: self.write_admission.as_ref(),
2076 _volume_lease: volume_lease,
2077 }));
2078 }
2079 if started.elapsed() >= self.config.checkout_timeout {
2080 self.writer_acquisition_counters
2081 .pooled_timeouts
2082 .fetch_add(1, Ordering::Relaxed);
2083 let message = format!(
2084 "timed out after {:?} waiting for sqlite writer connection",
2085 self.config.checkout_timeout
2086 );
2087 crate::timeout_sink::emit_timeout(
2088 &crate::timeout_sink::db_label(self),
2089 crate::timeout_sink::Site::PoolAdmission,
2090 &message,
2091 Some(
2092 self.config
2093 .checkout_timeout
2094 .as_millis()
2095 .min(u128::from(u64::MAX)) as u64,
2096 ),
2097 );
2098 return Err(SqliteError::WriterPoolCheckoutTimeout {
2099 timeout: self.config.checkout_timeout,
2100 });
2101 }
2102 }
2103 }
2104
2105 pub fn try_checkpoint_nowait(&self) -> Result<CheckpointGuard<'_>, SqliteError> {
2123 self.ensure_pooled_writer_active()?;
2124 let guard = self.writer.try_lock().ok_or_else(|| {
2125 SqliteError::InvalidData(
2126 "writer connection busy (checkpoint skipped this tick)".to_string(),
2127 )
2128 })?;
2129 self.ensure_pooled_writer_active()?;
2130 Ok(CheckpointGuard { guard })
2131 }
2132
2133 pub(crate) fn retire_pooled_writer(&self, conn: &Connection) {
2134 self.pooled_writer_retired.store(true, Ordering::Release);
2135 if let Err(error) = conn.authorizer(Some(deny_retired_writer)) {
2136 tracing::error!(
2137 %error,
2138 "failed to install the retired pooled-writer quarantine authorizer"
2139 );
2140 }
2141 }
2142
2143 #[cfg(test)]
2146 pub(crate) fn probe_retired_pooled_writer_for_test(&self) -> rusqlite::Result<i64> {
2147 self.writer
2148 .lock()
2149 .query_row("SELECT 1", [], |row| row.get(0))
2150 }
2151
2152 #[cfg(test)]
2156 pub(crate) fn leave_pooled_writer_transaction_open_for_test(&self) -> rusqlite::Result<()> {
2157 self.writer.lock().execute_batch("BEGIN IMMEDIATE")
2158 }
2159
2160 fn ensure_pooled_writer_active(&self) -> Result<(), SqliteError> {
2161 if self.pooled_writer_retired.load(Ordering::Acquire) {
2162 return Err(SqliteError::InvalidData(
2163 "pooled writer connection retired after a terminal transaction fault".to_string(),
2164 ));
2165 }
2166 Ok(())
2167 }
2168
2169 pub fn writer_acquisition_snapshot(&self) -> WriterAcquisitionSnapshot {
2172 self.writer_acquisition_counters
2173 .snapshot(self.write_admission.lease_timeouts())
2174 }
2175
2176 pub fn record_search_dispatch(&self, backend_id: &str, requested_kind: &str) {
2180 let mut dispatches = self.search_dispatches.lock();
2181 let count = dispatches
2182 .entry(backend_id.to_owned())
2183 .or_default()
2184 .entry(requested_kind.to_owned())
2185 .or_default();
2186 *count = count.saturating_add(1);
2187 }
2188
2189 pub fn record_note_candidate_hydration_row(&self) {
2193 let _ = self.note_candidate_hydration_rows.fetch_update(
2194 Ordering::Relaxed,
2195 Ordering::Relaxed,
2196 |current| Some(current.saturating_add(1)),
2197 );
2198 }
2199
2200 pub fn search_mechanism_snapshot(&self) -> SearchMechanismSnapshot {
2202 SearchMechanismSnapshot {
2203 dispatches_by_backend_and_kind: self.search_dispatches.lock().clone(),
2204 note_candidate_hydration_rows: self
2205 .note_candidate_hydration_rows
2206 .load(Ordering::Relaxed),
2207 }
2208 }
2209
2210 pub fn reader_acquisition_snapshot(&self) -> ReaderAcquisitionSnapshot {
2214 self.reader_acquisition_counters.snapshot(
2215 self.max_readers.max(1),
2216 self.sql_bridge_reader_slots.available_permits(),
2217 )
2218 }
2219
2220 pub(crate) fn record_reader_admission_timeout(&self) {
2224 self.reader_acquisition_counters.record_checkout_timeout();
2225 }
2226
2227 pub(crate) fn record_reader_query_error(&self, error: &StorageError) {
2233 if crate::read_cancellation::storage_error_sqlite_code(error)
2234 == Some(rusqlite::ErrorCode::DatabaseBusy)
2235 {
2236 self.reader_acquisition_counters.record_busy_timeout();
2237 }
2238 }
2239
2240 pub(crate) fn writer_acquisition_counters(&self) -> Arc<WriterAcquisitionCounters> {
2242 Arc::clone(&self.writer_acquisition_counters)
2243 }
2244
2245 pub(crate) fn write_admission(&self) -> Arc<WriteAdmission> {
2246 Arc::clone(&self.write_admission)
2247 }
2248
2249 pub fn effective_disk_guard_config(&self) -> Option<EffectiveDiskGuardConfig> {
2250 self.disk_guard_config
2251 }
2252
2253 #[cfg(test)]
2254 pub(crate) fn set_test_write_admission(
2255 &mut self,
2256 floor_bytes: u64,
2257 probe: impl Fn(&Path) -> std::io::Result<u64> + Send + Sync + 'static,
2258 ) {
2259 let database_path = if self.config.read_only {
2260 None
2261 } else {
2262 self.canonical_path().map(Path::to_path_buf)
2263 };
2264 let admission = Arc::new(
2265 WriteAdmission::new(
2266 database_path,
2267 floor_bytes,
2268 DEFAULT_DISK_GUARD_DEADLINE_MS,
2269 self.config.volume_lock_dir.clone(),
2270 )
2271 .expect("test admission volume identity"),
2272 );
2273 admission.set_test_space_probe(probe);
2274 self.write_admission = admission;
2275 if self.canonical_path().is_some() && !self.config.read_only {
2276 self.disk_guard_config = Some(EffectiveDiskGuardConfig {
2277 reserve_bytes: floor_bytes,
2278 guard_deadline_ms: DEFAULT_DISK_GUARD_DEADLINE_MS,
2279 reserve_source: DiskGuardConfigSource::Backend,
2280 deadline_source: DiskGuardConfigSource::Default,
2281 legacy_environment_present: false,
2282 });
2283 }
2284 }
2285
2286 pub fn available_readers(&self) -> usize {
2288 self.readers.len()
2289 }
2290
2291 pub fn max_readers(&self) -> usize {
2293 self.max_readers
2294 }
2295
2296 pub fn config(&self) -> &PoolConfig {
2298 &self.config
2299 }
2300
2301 #[cfg(any(test, feature = "test-support"))]
2311 pub fn observe_test_statement_starts(
2312 &self,
2313 limit: usize,
2314 ) -> Result<crate::statement_observer::StatementStartObservation, SqliteError> {
2315 self.statement_observer.observe(limit)
2316 }
2317
2318 pub fn main_pool_generation(&self) -> u64 {
2322 *self.main_pool_generation.get_or_init(|| {
2323 NEXT_MAIN_POOL_GENERATION
2324 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |next| {
2325 next.checked_add(1)
2326 })
2327 .expect("main pool generation exhausted")
2328 })
2329 }
2330
2331 pub(crate) fn reader_admission_timeout(&self, operation: &'static str) -> StorageError {
2335 StorageError::AdmissionTimeout {
2336 operation: operation.into(),
2337 timeout_ms: u64::try_from(self.config.checkout_timeout.as_millis()).unwrap_or(u64::MAX),
2338 pool_identity: Some(
2339 self.identity_registration
2340 .as_ref()
2341 .map(PoolIdentityRegistration::label)
2342 .unwrap_or_else(|| ":memory:".to_string()),
2343 ),
2344 }
2345 }
2346
2347 pub(crate) fn resolve_reader_checkout<'p>(
2366 &self,
2367 capability: StorageCapability,
2368 operation: &'static str,
2369 outcome: Result<Option<ReaderGuard<'p>>, SqliteError>,
2370 ) -> Result<ReaderGuard<'p>, StorageError> {
2371 match outcome {
2372 Ok(Some(mut guard)) => {
2373 guard.label_operation(operation);
2374 Ok(guard)
2375 }
2376 Ok(None) => Err(self.reader_checkout_refusal(capability, operation, None)),
2377 Err(error) => Err(self.reader_checkout_refusal(capability, operation, Some(error))),
2378 }
2379 }
2380
2381 fn reader_checkout_refusal(
2386 &self,
2387 capability: StorageCapability,
2388 operation: &'static str,
2389 error: Option<SqliteError>,
2390 ) -> StorageError {
2391 let Some(error) = error else {
2392 return StorageError::Timeout {
2393 operation: operation.into(),
2394 };
2395 };
2396 let is_pool_exhausted = matches!(
2397 &error,
2398 SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
2399 if code.code == rusqlite::ErrorCode::DatabaseBusy
2400 );
2401 if is_pool_exhausted {
2402 self.reader_admission_timeout(operation)
2403 } else {
2404 StorageError::driver(capability, operation, error)
2405 }
2406 }
2407
2408 pub(crate) fn sql_bridge_reader_slots(&self) -> Arc<Semaphore> {
2411 Arc::clone(&self.sql_bridge_reader_slots)
2412 }
2413
2414 pub(crate) fn sql_bridge_writer_slots(&self) -> Arc<Semaphore> {
2416 Arc::clone(&self.sql_bridge_writer_slots)
2417 }
2418
2419 pub fn retirement_writer_holds(&self) -> usize {
2424 usize::from(self.writer.is_locked())
2425 + usize::from(self.sql_bridge_writer_slots.available_permits() == 0)
2426 }
2427
2428 pub fn origin(&self) -> TxOrigin {
2434 self.origin.clone()
2435 }
2436
2437 pub fn canonical_path(&self) -> Option<&Path> {
2443 self.identity_path.as_deref()
2444 }
2445
2446 #[cfg(unix)]
2448 pub fn opened_file_identity(&self) -> Option<(u64, u64)> {
2449 self.opened_file_identity
2450 .map(DatabaseFileIdentity::unix_parts)
2451 }
2452
2453 #[cfg(any(unix, windows))]
2456 pub fn opened_file_identity_record(&self) -> Option<DatabaseFileIdentity> {
2457 self.opened_file_identity
2458 }
2459
2460 pub fn database_owner_identity(
2465 &self,
2466 ) -> Result<DatabaseOwnerIdentity, DatabaseOwnerIdentityError> {
2467 if self.identity_path.is_none() {
2468 return Err(DatabaseOwnerIdentityError::InMemory);
2469 }
2470 #[cfg(any(unix, windows))]
2471 {
2472 let durable_id = self
2473 .opened_database_id
2474 .ok_or(DatabaseOwnerIdentityError::DurableIdentityUnavailable)?;
2475 let file_identity = self
2476 .opened_file_identity
2477 .ok_or(DatabaseOwnerIdentityError::PhysicalIdentityUnavailable)?;
2478 Ok(DatabaseOwnerIdentity {
2479 durable_id,
2480 file_identity,
2481 })
2482 }
2483 #[cfg(not(any(unix, windows)))]
2484 {
2485 Err(DatabaseOwnerIdentityError::UnsupportedPlatform)
2486 }
2487 }
2488
2489 pub fn verify_database_owner(
2491 &self,
2492 expected: &DatabaseOwnerIdentity,
2493 ) -> Result<(), DatabaseOwnerIdentityError> {
2494 self.database_owner_identity()?.verify_owner(expected)
2495 }
2496
2497 pub fn write_queue_active(&self) -> bool {
2509 debug_assert!(
2510 self.config.write_queue_enabled.is_some(),
2511 "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
2512 before any write_queue_active read"
2513 );
2514 self.config.write_queue_enabled.unwrap_or(false) && self.config.path.is_some()
2515 }
2516
2517 pub fn writer_task_join_was_stored(&self) -> bool {
2523 self.writer_task_join_stored.load(Ordering::SeqCst)
2524 }
2525
2526 pub fn writer_task_handle(&self) -> Result<Option<WriterTaskHandle>, StorageError> {
2550 debug_assert!(
2558 self.config.write_queue_enabled.is_some(),
2559 "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
2560 before any writer_task_handle read"
2561 );
2562 if !self.config.write_queue_enabled.unwrap_or(false) {
2563 return Ok(None);
2564 }
2565 if let Some(existing) = self.writer_task.get() {
2568 return Ok(existing.clone());
2569 }
2570 if tokio::runtime::Handle::try_current().is_err() {
2574 return Err(StorageError::WriterTaskNoRuntime);
2575 }
2576 Ok(self
2577 .writer_task
2578 .get_or_init(|| {
2579 #[cfg(test)]
2580 self.writer_task_spawn_count
2581 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2582
2583 match crate::writer_task::spawn(self, self.config.write_queue_capacity) {
2584 Ok(handle) => Some(handle),
2585 Err(e) => {
2586 tracing::warn!(
2587 error = %e,
2588 "KHIVE_WRITE_QUEUE=1 but the writer task failed to spawn; \
2589 writes fall back to the pool-mutex path"
2590 );
2591 None
2592 }
2593 }
2594 })
2595 .clone())
2596 }
2597
2598 pub(crate) fn writer_task_for_write(
2608 &self,
2609 cached: Option<&WriterTaskHandle>,
2610 operation: &'static str,
2611 ) -> Result<Option<WriterTaskHandle>, StorageError> {
2612 let handle = match cached {
2613 Some(handle) => Some(handle.clone()),
2614 None => match self.writer_task_handle() {
2615 Ok(handle) => handle,
2616 Err(error) if self.config.write_routing_strict => return Err(error),
2617 Err(_) => None,
2618 },
2619 };
2620
2621 if handle.is_none() && self.config.write_routing_strict {
2622 return Err(StorageError::Pool {
2623 operation: operation.into(),
2624 message: "strict write routing requires a writer-task handle; no handle is \
2625 available, so the direct writer fallback was refused"
2626 .into(),
2627 });
2628 }
2629 Ok(handle)
2630 }
2631
2632 pub(crate) fn record_direct_route(&self, site: crate::timeout_sink::Site) {
2636 if self.write_queue_active() {
2637 crate::timeout_sink::emit_direct_route_violation(
2638 &crate::timeout_sink::db_label(self),
2639 site,
2640 );
2641 }
2642 }
2643
2644 pub fn writer_task_for_runtime_write(
2649 &self,
2650 operation: RuntimeWriteOperation,
2651 ) -> Result<Option<WriterTaskHandle>, StorageError> {
2652 let handle = self.writer_task_for_write(None, operation.operation())?;
2653 if handle.is_none() {
2654 self.record_direct_route(operation.fallback_site());
2655 }
2656 Ok(handle)
2657 }
2658
2659 #[cfg(test)]
2664 pub(crate) fn writer_task_spawn_count(&self) -> usize {
2665 self.writer_task_spawn_count
2666 .load(std::sync::atomic::Ordering::SeqCst)
2667 }
2668
2669 pub(crate) fn set_writer_task_join(&self, join: tokio::task::JoinHandle<()>) {
2682 let first_store = !self.writer_task_join_stored.swap(true, Ordering::SeqCst);
2687 debug_assert!(
2688 first_store,
2689 "writer task JoinHandle stored twice (even counting a taken one); \
2690 the writer_task OnceLock is supposed to make spawn at-most-once per pool"
2691 );
2692 if first_store {
2693 *self.writer_task_join.lock() = Some(join);
2694 }
2695 }
2696
2697 pub fn take_writer_task_join(&self) -> Option<tokio::task::JoinHandle<()>> {
2715 self.writer_task_join.lock().take()
2716 }
2717
2718 fn open_reader_connection(&self) -> Result<Connection, SqliteError> {
2719 let path = self.read_connection_path()?;
2720 #[cfg(any(unix, windows))]
2721 if let Some(identity_path) = self.identity_path.as_deref() {
2722 self.verify_opened_file_identity(identity_path)?;
2723 }
2724 let conn = open_reader_connection(path, &self.config)?;
2725 #[cfg(any(unix, windows))]
2726 if let Some(identity_path) = self.identity_path.as_deref() {
2727 self.verify_connection_file_identity(&conn, identity_path)?;
2728 }
2729 self.verify_opened_database_id(&conn)?;
2730 #[cfg(any(test, feature = "test-support"))]
2731 crate::statement_observer::install(&conn, &self.statement_observer)?;
2732 Ok(conn)
2733 }
2734
2735 fn read_connection_path(&self) -> Result<&Path, SqliteError> {
2736 self.read_only_open_target
2737 .as_deref()
2738 .or(self.identity_path.as_deref())
2739 .ok_or_else(|| {
2740 SqliteError::InvalidData(
2741 "in-memory databases do not support standalone connections".to_string(),
2742 )
2743 })
2744 }
2745
2746 pub fn open_standalone_writer(&self) -> Result<Connection, SqliteError> {
2757 self.write_admission.check()?;
2758 self.open_standalone_writer_for_admitted_operation()
2759 }
2760
2761 pub(crate) fn open_standalone_writer_for_admitted_operation(
2765 &self,
2766 ) -> Result<Connection, SqliteError> {
2767 let conn = self.open_standalone_writer_untracked()?;
2768 self.writer_acquisition_counters
2769 .standalone_acquisitions
2770 .fetch_add(1, Ordering::Relaxed);
2771 Ok(conn)
2772 }
2773
2774 pub(crate) fn open_standalone_writer_untracked(&self) -> Result<Connection, SqliteError> {
2784 let path = self.identity_path.as_deref().ok_or_else(|| {
2785 SqliteError::InvalidData(
2786 "in-memory databases do not support standalone connections".to_string(),
2787 )
2788 })?;
2789
2790 if self.config.read_only {
2791 return Err(SqliteError::InvalidData(
2792 "database is read-only: standalone write connections are not permitted".to_string(),
2793 ));
2794 }
2795
2796 #[cfg(any(unix, windows))]
2800 self.verify_opened_file_identity(path)?;
2801
2802 #[cfg(test)]
2803 run_identity_open_hook(path, IdentityOpenStage::BeforeStandaloneOpen, None);
2804
2805 let conn = claimed_file_identity::open_connection(
2806 &self.config,
2807 path,
2808 OpenFlags::SQLITE_OPEN_READ_WRITE
2809 | OpenFlags::SQLITE_OPEN_NO_MUTEX
2810 | OpenFlags::SQLITE_OPEN_URI,
2811 self.identity_path.as_deref(),
2812 )?;
2813 #[cfg(test)]
2814 run_identity_open_hook(path, IdentityOpenStage::AfterStandaloneOpen, Some(&conn));
2815 #[cfg(any(unix, windows))]
2816 self.verify_connection_file_identity(&conn, path)?;
2817 self.verify_opened_database_id(&conn)?;
2818 #[cfg(feature = "namespace-trigram-proto")]
2819 register_namespace_trigram(&conn)?;
2820 register_writer_clock(&conn)?;
2821 register_rfc3339_key(&conn)?;
2825 conn.busy_timeout(self.config.busy_timeout)?;
2826 if self.config.code_map_vfs.is_none() {
2827 self.checkpoint_ownership
2828 .configure_wal_autocheckpoint(&conn)?;
2829 }
2830 conn.pragma_update(None, "foreign_keys", "ON")?;
2831 conn.pragma_update(None, "synchronous", "NORMAL")?;
2832
2833 let wal_enabled =
2834 self.config.wal_mode && current_journal_mode(&conn)?.eq_ignore_ascii_case("wal");
2835 if wal_enabled {
2836 conn.pragma_update(
2837 None,
2838 "journal_size_limit",
2839 self.config.journal_size_limit_bytes,
2840 )?;
2841 }
2842
2843 #[cfg(any(test, feature = "test-support"))]
2844 crate::statement_observer::install(&conn, &self.statement_observer)?;
2845 Ok(conn)
2846 }
2847
2848 #[cfg(any(unix, windows))]
2849 fn verify_opened_file_identity(&self, path: &Path) -> Result<(), SqliteError> {
2850 let Some(expected) = self.opened_file_identity else {
2851 return Err(SqliteError::InvalidData(
2852 "file-backed pool has no opened database file identity".to_string(),
2853 ));
2854 };
2855 let current = database_file_identity(path).ok();
2856 if current != Some(expected) {
2857 return Err(SqliteError::InvalidData(
2858 "pool database file identity changed since the first open; refusing standalone connection"
2859 .to_string(),
2860 ));
2861 }
2862 Ok(())
2863 }
2864
2865 #[cfg(any(unix, windows))]
2866 fn verify_connection_file_identity(
2867 &self,
2868 conn: &Connection,
2869 path: &Path,
2870 ) -> Result<(), SqliteError> {
2871 let opened = opened_sqlite_file_identity(conn, path)?;
2872 if self.opened_file_identity != Some(opened) {
2873 return Err(SqliteError::InvalidData(
2874 "pool database file identity changed since the first open; refusing standalone connection"
2875 .to_string(),
2876 ));
2877 }
2878 Ok(())
2879 }
2880
2881 fn verify_opened_database_id(&self, conn: &Connection) -> Result<(), SqliteError> {
2882 let Some(expected) = self.opened_database_id else {
2886 return Ok(());
2887 };
2888 if self.identity_path.is_some() && read_database_id(conn)? != Some(expected) {
2889 return Err(SqliteError::InvalidData(
2890 "pool database identity changed since the first open; refusing standalone connection"
2891 .to_string(),
2892 ));
2893 }
2894 Ok(())
2895 }
2896
2897 #[cfg(test)]
2901 pub(crate) fn effective_wal_autocheckpoint_pages(&self) -> u32 {
2902 self.checkpoint_ownership.wal_autocheckpoint_pages()
2903 }
2904
2905 pub fn claim_checkpoint_ownership(&self) -> Result<(), SqliteError> {
2927 if !self.checkpoint_ownership.begin_claim() {
2928 return Ok(());
2929 }
2930 let result = (|| {
2931 if !self.config.read_only {
2932 let writer = self.writer_for_checkpoint_operation()?;
2933 writer.conn().pragma_update(None, "wal_autocheckpoint", 0)?;
2934 }
2935 Ok(())
2936 })();
2937 self.checkpoint_ownership.finish_claim(result.is_ok());
2938 result
2939 }
2940
2941 pub async fn propagate_checkpoint_claim_to_writer_task(&self) -> Result<(), StorageError> {
2950 let Some(handle) = self.writer_task_handle()? else {
2951 return Ok(());
2952 };
2953 handle
2954 .send_top_level(|conn| {
2955 conn.pragma_update(None, "wal_autocheckpoint", 0)
2956 .map_err(|e| StorageError::Pool {
2957 operation: "claim_checkpoint_ownership".into(),
2958 message: e.to_string(),
2959 })
2960 })
2961 .await
2962 }
2963
2964 pub(crate) fn open_standalone_reader(
2973 &self,
2974 purpose: StandaloneReaderPurpose,
2975 ) -> Result<Connection, SqliteError> {
2976 let path = self.read_connection_path()?;
2977
2978 #[cfg(any(unix, windows))]
2979 if let Some(identity_path) = self.identity_path.as_deref() {
2980 self.verify_opened_file_identity(identity_path)?;
2981 }
2982
2983 let conn = claimed_file_identity::open_connection(
2984 &self.config,
2985 path,
2986 OpenFlags::SQLITE_OPEN_READ_ONLY
2987 | OpenFlags::SQLITE_OPEN_NO_MUTEX
2988 | OpenFlags::SQLITE_OPEN_URI,
2989 self.identity_path.as_deref(),
2990 )?;
2991 #[cfg(any(unix, windows))]
2992 if let Some(identity_path) = self.identity_path.as_deref() {
2993 self.verify_connection_file_identity(&conn, identity_path)?;
2994 }
2995 self.verify_opened_database_id(&conn)?;
2996 configure_reader_connection(&conn, &self.config)?;
2997 conn.pragma_update(None, "synchronous", "NORMAL")?;
2998 self.reader_acquisition_counters
2999 .record_standalone_open(purpose);
3000 #[cfg(any(test, feature = "test-support"))]
3001 crate::statement_observer::install(&conn, &self.statement_observer)?;
3002 Ok(conn)
3003 }
3004
3005 fn return_reader(&self, conn: Connection, dirty: bool) {
3006 if self.max_readers == 0 {
3007 return;
3008 }
3009
3010 if reset_reader_connection(&conn, dirty, &self.config)
3011 && reader_connection_is_healthy(&conn)
3012 {
3013 self.enqueue_reader_slot(conn);
3014 return;
3015 }
3016
3017 close_connection_quietly(conn);
3018 self.replace_discarded_reader_slot();
3019 }
3020
3021 fn enqueue_reader_slot(&self, conn: Connection) {
3025 if let Err(conn) = self.readers.push(conn) {
3026 eprintln!("[sqlite-pool] reader pool queue full, discarding replacement connection");
3027 close_connection_quietly(conn);
3028 }
3029 }
3030
3031 fn replace_discarded_reader_slot(&self) {
3038 match self.open_reader_connection() {
3039 Ok(conn) => self.enqueue_reader_slot(conn),
3040 Err(error) => {
3041 self.reader_acquisition_counters
3042 .record_reader_replacement_open_failure();
3043 tracing::warn!(
3044 %error,
3045 "sqlite-pool: reader replacement connection failed to open; the physical \
3046 pool permanently shrinks by one slot below max_readers"
3047 );
3048 }
3049 }
3050 }
3051}
3052
3053const MAX_SYMLINK_DEPTH: u32 = 40;
3058
3059fn mint_db_identity(configured_path: &Path) -> Result<(DbIdentity, PathBuf), SqliteError> {
3092 let absolute = if configured_path.is_absolute() {
3093 configured_path.to_path_buf()
3094 } else {
3095 let cwd = std::env::current_dir().map_err(|e| {
3096 SqliteError::InvalidData(format!(
3097 "cannot mint database identity for {configured_path:?}: failed to resolve the \
3098 process current directory: {e}"
3099 ))
3100 })?;
3101 cwd.join(configured_path)
3102 };
3103
3104 if absolute.exists() {
3105 let canonical = absolute.canonicalize().map_err(|e| {
3106 SqliteError::InvalidData(format!(
3107 "cannot mint database identity: failed to canonicalize existing path \
3108 {absolute:?}: {e}"
3109 ))
3110 })?;
3111 return Ok((
3112 DbIdentity::new(canonical.clone().into_os_string()),
3113 canonical,
3114 ));
3115 }
3116
3117 let resolved_target = resolve_symlink_chain(&absolute)?;
3118 let parent = resolved_target.parent().ok_or_else(|| {
3119 SqliteError::InvalidData(format!(
3120 "cannot mint database identity for {resolved_target:?}: path has no parent \
3121 directory"
3122 ))
3123 })?;
3124 let file_name = resolved_target.file_name().ok_or_else(|| {
3125 SqliteError::InvalidData(format!(
3126 "cannot mint database identity for {resolved_target:?}: path has no file name"
3127 ))
3128 })?;
3129 let canonical_parent = parent.canonicalize().map_err(|e| {
3130 SqliteError::InvalidData(format!(
3131 "cannot mint database identity: parent directory {parent:?} of first-open path \
3132 {resolved_target:?} does not exist or is inaccessible: {e}"
3133 ))
3134 })?;
3135 let mut identity_path = canonical_parent;
3136 identity_path.push(file_name);
3137 Ok((
3138 DbIdentity::new(identity_path.clone().into_os_string()),
3139 identity_path,
3140 ))
3141}
3142
3143#[cfg(any(unix, windows))]
3144fn database_file_identity_if_exists(
3145 path: &Path,
3146) -> Result<Option<DatabaseFileIdentity>, SqliteError> {
3147 match database_file_identity(path) {
3148 Ok(identity) => Ok(Some(identity)),
3149 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
3150 Err(error) => Err(error.into()),
3151 }
3152}
3153
3154#[cfg(any(unix, windows))]
3155pub(crate) fn opened_sqlite_file_identity(
3156 conn: &Connection,
3157 path: &Path,
3158) -> Result<DatabaseFileIdentity, SqliteError> {
3159 #[cfg(unix)]
3160 {
3161 verify_sqlite_opened_file_still_at_path(conn)?;
3162 Ok(database_file_identity(path)?)
3163 }
3164 #[cfg(windows)]
3165 {
3166 let opened = sqlite_opened_file_identity(conn)?;
3167 if database_file_identity(path).ok() != Some(opened) {
3168 return Err(SqliteError::InvalidData(
3169 "database file identity changed while SQLite held the opened file".to_string(),
3170 ));
3171 }
3172 Ok(opened)
3173 }
3174}
3175
3176fn read_database_id(conn: &Connection) -> Result<Option<uuid::Uuid>, SqliteError> {
3180 let table_exists: bool = conn.query_row(
3181 "SELECT count(*) != 0 FROM main.sqlite_master WHERE type = 'table' AND name = ?1",
3182 [DATABASE_ID_TABLE],
3183 |row| row.get(0),
3184 )?;
3185 if !table_exists {
3186 return Ok(None);
3189 }
3190 let id: String = conn.query_row(
3191 &format!("SELECT id FROM main.{DATABASE_ID_TABLE} WHERE singleton = 1"),
3192 [],
3193 |row| row.get(0),
3194 )?;
3195 let id = uuid::Uuid::parse_str(&id).map_err(|error| {
3196 SqliteError::InvalidData(format!("invalid stored database identity: {error}"))
3197 })?;
3198 Ok(Some(id))
3199}
3200
3201#[cfg(test)]
3202fn initialize_database_id(conn: &mut Connection) -> Result<uuid::Uuid, SqliteError> {
3203 initialize_database_id_with_admission(conn, None)
3204}
3205
3206fn initialize_database_id_with_admission(
3207 conn: &mut Connection,
3208 admission: Option<&WriteAdmission>,
3209) -> Result<uuid::Uuid, SqliteError> {
3210 if let Some(id) = read_database_id(conn)? {
3211 return Ok(id);
3212 }
3213 let transaction = conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
3214 if let Some(admission) = admission {
3215 if let Err(error) = admission.check() {
3216 let rollback = transaction.rollback();
3217 return Err(crate::migrations::capacity_refusal_after_rollback(
3218 conn,
3219 rollback,
3220 error,
3221 "database identity bootstrap",
3222 ));
3223 }
3224 }
3225 transaction.execute_batch(&format!(
3226 "CREATE TABLE IF NOT EXISTS main.{DATABASE_ID_TABLE} (\
3227 singleton INTEGER PRIMARY KEY CHECK (singleton = 1), \
3228 id TEXT NOT NULL\
3229 )"
3230 ))?;
3231 transaction.execute(
3232 &format!("INSERT OR IGNORE INTO main.{DATABASE_ID_TABLE} (singleton, id) VALUES (1, ?1)"),
3233 [uuid::Uuid::new_v4().to_string()],
3234 )?;
3235 let id: String = transaction.query_row(
3236 &format!("SELECT id FROM main.{DATABASE_ID_TABLE} WHERE singleton = 1"),
3237 [],
3238 |row| row.get(0),
3239 )?;
3240 let id = uuid::Uuid::parse_str(&id).map_err(|error| {
3241 SqliteError::InvalidData(format!("invalid stored database identity: {error}"))
3242 })?;
3243 transaction.commit()?;
3244 Ok(id)
3245}
3246
3247#[cfg(unix)]
3248fn verify_sqlite_opened_file_still_at_path(conn: &Connection) -> Result<(), SqliteError> {
3249 let mut moved: std::ffi::c_int = 0;
3250 let result = unsafe {
3256 rusqlite::ffi::sqlite3_file_control(
3257 conn.handle(),
3258 c"main".as_ptr(),
3259 rusqlite::ffi::SQLITE_FCNTL_HAS_MOVED,
3260 (&mut moved as *mut std::ffi::c_int).cast(),
3261 )
3262 };
3263 if result != rusqlite::ffi::SQLITE_OK {
3264 return Err(SqliteError::InvalidData(format!(
3265 "cannot verify opened database file identity (SQLite file control {result})"
3266 )));
3267 }
3268 if moved != 0 {
3269 return Err(SqliteError::InvalidData(
3270 "database file identity changed while SQLite held the opened file".to_string(),
3271 ));
3272 }
3273 Ok(())
3274}
3275
3276fn resolve_symlink_chain(path: &Path) -> Result<PathBuf, SqliteError> {
3282 let mut current = path.to_path_buf();
3283 for _ in 0..MAX_SYMLINK_DEPTH {
3284 match fs::symlink_metadata(¤t) {
3285 Ok(meta) if meta.file_type().is_symlink() => {
3286 let target = fs::read_link(¤t).map_err(|e| {
3287 SqliteError::InvalidData(format!(
3288 "cannot mint database identity: failed to read symlink {current:?}: {e}"
3289 ))
3290 })?;
3291 current = if target.is_absolute() {
3292 target
3293 } else {
3294 match current.parent() {
3295 Some(parent) => parent.join(&target),
3296 None => target,
3297 }
3298 };
3299 }
3300 _ => return Ok(current),
3301 }
3302 }
3303 Err(SqliteError::InvalidData(format!(
3304 "cannot mint database identity for {path:?}: symlink chain exceeds \
3305 {MAX_SYMLINK_DEPTH} levels"
3306 )))
3307}
3308
3309fn effective_reader_count(config: &PoolConfig, wal_enabled: bool) -> usize {
3310 if config.path.is_some() && (config.read_only || config.code_map_vfs.is_some()) {
3311 config.max_readers.max(1)
3312 } else if config.path.is_some() && config.wal_mode && wal_enabled {
3313 config.max_readers
3314 } else {
3315 0
3316 }
3317}
3318
3319fn open_writer_connection(
3320 config: &PoolConfig,
3321 read_only_open_target: Option<&Path>,
3322 identity_path: Option<&Path>,
3323) -> Result<Connection, SqliteError> {
3324 claimed_file_identity::open_writer(config, read_only_open_target, identity_path)
3325}
3326
3327fn validate_wal_ceiling_at_open(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
3332 let bytes = config.wal_ceiling.effective_bytes(config.read_only);
3333 if bytes == 0 {
3334 return Ok(());
3335 }
3336 let page_size: i64 = conn.pragma_query_value(None, "page_size", |row| row.get(0))?;
3337 let page_size = u64::try_from(page_size).map_err(|_| {
3338 SqliteError::InvalidData("SQLite reported a negative page size".to_string())
3339 })?;
3340 let minimum_bytes = page_size.checked_add(56).ok_or_else(|| {
3341 SqliteError::InvalidData("SQLite page size overflowed the WAL frame floor".to_string())
3342 })?;
3343 if bytes < minimum_bytes {
3344 return Err(SqliteError::WalCeilingBelowMinimum {
3345 bytes,
3346 page_size,
3347 minimum_bytes,
3348 });
3349 }
3350 Err(SqliteError::WalCapacityUnavailable {
3351 bytes,
3352 capability: "WAL I/O limiter",
3353 })
3354}
3355
3356fn read_only_open_target(
3379 config: &PoolConfig,
3380 physical_path: Option<&Path>,
3381) -> Result<Option<PathBuf>, SqliteError> {
3382 if !config.read_only || config.code_map_vfs.is_some() {
3383 return Ok(None);
3384 }
3385 let Some(path) = physical_path else {
3386 return Ok(None);
3387 };
3388 read_only_wal_open_target_for_path(path).map(Some)
3389}
3390
3391fn read_only_wal_open_target_for_path(path: &Path) -> Result<PathBuf, SqliteError> {
3392 if !sqlite_header_uses_wal(path)? {
3393 return Ok(path.to_path_buf());
3394 }
3395
3396 let shm = sqlite_sidecar_path(path, "-shm");
3397 match fs::metadata(&shm) {
3398 Ok(metadata) if metadata.permissions().readonly() => {
3399 let wal = sqlite_sidecar_path(path, "-wal");
3400 match fs::metadata(&wal) {
3401 Ok(_) => Ok(path.to_path_buf()),
3402 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3403 Err(SqliteError::InvalidData(format!(
3404 "read-only WAL snapshot {} has a shared-memory sidecar {} but no WAL \
3405 sidecar {}; refusing the inconsistent sidecar set before SQLite open",
3406 path.display(),
3407 shm.display(),
3408 wal.display(),
3409 )))
3410 }
3411 Err(error) => Err(SqliteError::Io(error)),
3412 }
3413 }
3414 Ok(_) => Err(SqliteError::InvalidData(format!(
3415 "read-only WAL snapshot {} has a writable WAL shared-memory sidecar {}; close every \
3416 live writer and remove the transient -shm file (or make a genuinely frozen snapshot) \
3417 before inspection",
3418 path.display(),
3419 shm.display(),
3420 ))),
3421 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3422 let wal = sqlite_sidecar_path(path, "-wal");
3423 match fs::metadata(&wal) {
3424 Ok(metadata) if metadata.len() > 0 => Err(SqliteError::InvalidData(format!(
3425 "read-only WAL snapshot {} has a non-empty WAL sidecar {} but no read-only \
3426 shared-memory sidecar {}; refusing before SQLite open because immutable \
3427 mode would omit committed WAL frames and ordinary read-only mode would \
3428 create or mutate -shm; include the frozen read-only -shm beside this \
3429 snapshot, or checkpoint a writable copy before inspection",
3430 path.display(),
3431 wal.display(),
3432 shm.display(),
3433 ))),
3434 Ok(_) => sqlite_immutable_uri(path),
3435 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
3436 sqlite_immutable_uri(path)
3437 }
3438 Err(error) => Err(SqliteError::Io(error)),
3439 }
3440 }
3441 Err(error) => Err(SqliteError::Io(error)),
3442 }
3443}
3444
3445pub(crate) fn open_read_only_snapshot_connection(path: &Path) -> Result<Connection, SqliteError> {
3446 let (_, physical_path) = mint_db_identity(path)?;
3447 let target = read_only_wal_open_target_for_path(&physical_path)?;
3448 let conn = Connection::open_with_flags(&target, reader_open_flags())?;
3449 #[cfg(feature = "namespace-trigram-proto")]
3450 register_namespace_trigram(&conn)?;
3451 Ok(conn)
3452}
3453
3454fn sqlite_header_uses_wal(path: &Path) -> Result<bool, SqliteError> {
3455 let mut file = fs::File::open(path)?;
3456 let mut header = [0_u8; 20];
3457 if let Err(error) = file.read_exact(&mut header) {
3458 if error.kind() == std::io::ErrorKind::UnexpectedEof {
3459 return Ok(false);
3460 }
3461 return Err(SqliteError::Io(error));
3462 }
3463 Ok(&header[..16] == b"SQLite format 3\0" && header[18] == 2 && header[19] == 2)
3464}
3465
3466fn sqlite_sidecar_path(path: &Path, suffix: &str) -> PathBuf {
3467 let mut sidecar = path.as_os_str().to_os_string();
3468 sidecar.push(suffix);
3469 PathBuf::from(sidecar)
3470}
3471
3472fn sqlite_immutable_uri(path: &Path) -> Result<PathBuf, SqliteError> {
3473 let absolute = if path.is_absolute() {
3474 path.to_path_buf()
3475 } else {
3476 std::env::current_dir()?.join(path)
3477 };
3478 let mut uri = String::from("file:");
3479
3480 #[cfg(unix)]
3481 {
3482 use std::os::unix::ffi::OsStrExt as _;
3483 push_sqlite_uri_path(&mut uri, absolute.as_os_str().as_bytes());
3484 }
3485
3486 #[cfg(not(unix))]
3487 {
3488 let path = absolute.to_str().ok_or_else(|| {
3489 SqliteError::InvalidData(format!(
3490 "read-only WAL snapshot path is not representable as a SQLite URI: {}",
3491 absolute.display()
3492 ))
3493 })?;
3494 let normalized = path.replace('\\', "/");
3495 if cfg!(windows) && !normalized.starts_with('/') {
3496 uri.push('/');
3497 }
3498 push_sqlite_uri_path(&mut uri, normalized.as_bytes());
3499 }
3500
3501 uri.push_str("?mode=ro&immutable=1");
3502 Ok(PathBuf::from(uri))
3503}
3504
3505fn push_sqlite_uri_path(uri: &mut String, bytes: &[u8]) {
3506 const HEX: &[u8; 16] = b"0123456789ABCDEF";
3507 for &byte in bytes {
3508 if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~' | b'/') {
3509 uri.push(byte as char);
3510 } else {
3511 uri.push('%');
3512 uri.push(HEX[(byte >> 4) as usize] as char);
3513 uri.push(HEX[(byte & 0x0f) as usize] as char);
3514 }
3515 }
3516}
3517
3518fn open_reader_connection(path: &Path, config: &PoolConfig) -> Result<Connection, SqliteError> {
3519 let conn = claimed_file_identity::open_connection(
3520 config,
3521 path,
3522 reader_open_flags(),
3523 config.path.as_deref(),
3524 )?;
3525 configure_reader_connection(&conn, config)?;
3526 Ok(conn)
3527}
3528
3529fn writer_open_flags() -> OpenFlags {
3530 OpenFlags::SQLITE_OPEN_READ_WRITE
3531 | OpenFlags::SQLITE_OPEN_CREATE
3532 | OpenFlags::SQLITE_OPEN_URI
3533 | OpenFlags::SQLITE_OPEN_NO_MUTEX
3534}
3535
3536fn writer_read_only_open_flags() -> OpenFlags {
3539 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
3540}
3541
3542fn reader_open_flags() -> OpenFlags {
3543 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
3544}
3545
3546#[cfg(feature = "namespace-trigram-proto")]
3547fn register_namespace_trigram(conn: &Connection) -> Result<(), SqliteError> {
3548 crate::namespace_trigram_proto::register(conn).map_err(SqliteError::InvalidData)
3549}
3550
3551fn register_writer_clock(conn: &Connection) -> Result<(), SqliteError> {
3552 conn.create_scalar_function(
3555 "khive_now_micros",
3556 0,
3557 rusqlite::functions::FunctionFlags::SQLITE_UTF8,
3558 |_| Ok(chrono::Utc::now().timestamp_micros()),
3559 )?;
3560 Ok(())
3561}
3562
3563pub(crate) fn rfc3339_instant_key(instant: chrono::DateTime<chrono::Utc>) -> Vec<u8> {
3566 let mut key = Vec::with_capacity(12);
3567 key.extend_from_slice(&((instant.timestamp() as u64) ^ (1_u64 << 63)).to_be_bytes());
3568 key.extend_from_slice(&instant.timestamp_subsec_nanos().to_be_bytes());
3569 key
3570}
3571
3572pub(crate) fn strict_rfc3339_key(text: &str) -> Option<Vec<u8>> {
3576 chrono::DateTime::parse_from_rfc3339(text)
3577 .ok()
3578 .map(|instant| rfc3339_instant_key(instant.with_timezone(&chrono::Utc)))
3579}
3580
3581pub(crate) fn register_rfc3339_key(conn: &Connection) -> rusqlite::Result<()> {
3583 use rusqlite::functions::FunctionFlags;
3584 use rusqlite::types::ValueRef;
3585
3586 conn.create_scalar_function(
3587 "khive_rfc3339_key",
3588 1,
3589 FunctionFlags::SQLITE_UTF8
3590 | FunctionFlags::SQLITE_DETERMINISTIC
3591 | FunctionFlags::SQLITE_INNOCUOUS,
3592 |ctx| {
3593 let text = match ctx.get_raw(0) {
3594 ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
3595 _ => None,
3596 };
3597 let key = text
3598 .and_then(|text| text.parse::<chrono::DateTime<chrono::Utc>>().ok())
3599 .map(rfc3339_instant_key);
3600 Ok(key)
3601 },
3602 )?;
3603 conn.create_scalar_function(
3608 "khive_rfc3339_strict_key",
3609 1,
3610 FunctionFlags::SQLITE_UTF8
3611 | FunctionFlags::SQLITE_DETERMINISTIC
3612 | FunctionFlags::SQLITE_INNOCUOUS,
3613 |ctx| {
3614 let text = match ctx.get_raw(0) {
3615 ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
3616 _ => None,
3617 };
3618 let key = text.and_then(strict_rfc3339_key);
3619 Ok(key)
3620 },
3621 )?;
3622 Ok(())
3623}
3624
3625fn configure_writer_connection(
3626 conn: &Connection,
3627 config: &PoolConfig,
3628) -> Result<bool, SqliteError> {
3629 #[cfg(feature = "namespace-trigram-proto")]
3630 register_namespace_trigram(conn)?;
3631 register_writer_clock(conn)?;
3632 register_rfc3339_key(conn)?;
3633 if config.read_only {
3634 conn.pragma_update(None, "foreign_keys", "ON")?;
3638 conn.busy_timeout(config.busy_timeout)?;
3639 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3640 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3641 conn.pragma_update(None, "temp_store", "MEMORY")?;
3642 conn.pragma_update(None, "query_only", "ON")?;
3643
3644 let wal_enabled =
3645 config.wal_mode && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
3646 return Ok(wal_enabled);
3647 }
3648
3649 let wants_wal = config.path.is_some() && config.wal_mode;
3650
3651 if wants_wal {
3652 conn.pragma_update(None, "journal_mode", "WAL")?;
3653 }
3654 code_map::require_delete_journal(conn, config)?;
3655
3656 conn.pragma_update(None, "synchronous", "NORMAL")?;
3657 conn.pragma_update(None, "foreign_keys", "ON")?;
3658 conn.busy_timeout(config.busy_timeout)?;
3659 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3660 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3661 conn.pragma_update(None, "temp_store", "MEMORY")?;
3662 if config.code_map_vfs.is_none() {
3667 conn.pragma_update(
3668 None,
3669 "wal_autocheckpoint",
3670 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3671 )?;
3672 }
3673
3674 let wal_enabled = wants_wal && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
3675
3676 if wal_enabled {
3677 conn.pragma_update(None, "journal_size_limit", config.journal_size_limit_bytes)?;
3678 }
3679
3680 Ok(wal_enabled)
3681}
3682
3683fn configure_reader_connection(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
3684 #[cfg(feature = "namespace-trigram-proto")]
3685 register_namespace_trigram(conn)?;
3686 register_rfc3339_key(conn)?;
3687 conn.pragma_update(None, "foreign_keys", "ON")?;
3688 conn.busy_timeout(config.busy_timeout)?;
3689 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3690 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
3691 conn.pragma_update(None, "temp_store", "MEMORY")?;
3692 Ok(())
3693}
3694
3695fn current_journal_mode(conn: &Connection) -> Result<String, SqliteError> {
3696 conn.pragma_query_value(None, "journal_mode", |row| row.get::<_, String>(0))
3697 .map(|mode| mode.to_ascii_lowercase())
3698 .map_err(Into::into)
3699}
3700
3701fn reset_reader_connection(conn: &Connection, dirty: bool, config: &PoolConfig) -> bool {
3702 if !conn.is_autocommit() {
3703 match conn.execute_batch("ROLLBACK") {
3704 Ok(()) => {}
3705 Err(rusqlite::Error::SqliteFailure(err, _)) => {
3706 if matches!(
3707 err.code,
3708 rusqlite::ErrorCode::CannotOpen
3709 | rusqlite::ErrorCode::DatabaseCorrupt
3710 | rusqlite::ErrorCode::NotADatabase
3711 | rusqlite::ErrorCode::DiskFull
3712 ) {
3713 return false;
3714 }
3715 }
3716 Err(_) => return false,
3717 }
3718 if !conn.is_autocommit() {
3719 return false;
3720 }
3721 }
3722
3723 if !dirty {
3724 return true;
3725 }
3726
3727 reader_connection_state_is_pristine(conn)
3728 && reader_connection_settings_match_baseline(conn, config, 0)
3729}
3730
3731fn reader_connection_state_is_pristine(conn: &Connection) -> bool {
3746 let has_temp_objects: bool = match conn.query_row(
3747 "SELECT EXISTS(SELECT 1 FROM sqlite_temp_master)",
3748 [],
3749 |row| row.get(0),
3750 ) {
3751 Ok(v) => v,
3752 Err(_) => return false,
3753 };
3754 if has_temp_objects {
3755 return false;
3756 }
3757
3758 let attached_databases: i64 = match conn.query_row(
3759 "SELECT COUNT(*) FROM pragma_database_list WHERE name NOT IN ('main', 'temp')",
3760 [],
3761 |row| row.get(0),
3762 ) {
3763 Ok(v) => v,
3764 Err(_) => return false,
3765 };
3766 attached_databases == 0
3767}
3768
3769fn reader_connection_settings_match_baseline(
3783 conn: &Connection,
3784 config: &PoolConfig,
3785 expected_query_only: i64,
3786) -> bool {
3787 let expected_busy_timeout_ms =
3788 i64::try_from(config.busy_timeout.as_millis()).unwrap_or(i64::MAX);
3789 let expected_cache_size: i64 = CACHE_SIZE_KIB.parse().unwrap_or(-65536);
3790 let checks: [(&str, i64); 8] = [
3791 ("query_only", expected_query_only),
3792 ("writable_schema", 0),
3793 ("foreign_keys", 1),
3794 ("busy_timeout", expected_busy_timeout_ms),
3795 ("cache_size", expected_cache_size),
3796 ("temp_store", 2),
3797 ("read_uncommitted", 0),
3798 ("defer_foreign_keys", 0),
3799 ];
3800 checks.iter().all(|(pragma, expected)| {
3801 conn.pragma_query_value(None, pragma, |row| row.get::<_, i64>(0))
3802 .map(|actual| actual == *expected)
3803 .unwrap_or(false)
3804 })
3805}
3806
3807fn restore_shared_reader_state(conn: &Connection, config: &PoolConfig) -> bool {
3816 if !detach_non_main_databases(conn) {
3817 return false;
3818 }
3819 if !drop_temp_objects(conn) {
3820 return false;
3821 }
3822 let expected_query_only = i64::from(config.read_only);
3823 if reset_observable_settings(conn, config, expected_query_only).is_err() {
3824 return false;
3825 }
3826 reader_connection_state_is_pristine(conn)
3827 && reader_connection_settings_match_baseline(conn, config, expected_query_only)
3828}
3829
3830fn detach_non_main_databases(conn: &Connection) -> bool {
3831 loop {
3832 let name: Option<String> = match conn.query_row(
3833 "SELECT name FROM pragma_database_list WHERE name NOT IN ('main', 'temp') LIMIT 1",
3834 [],
3835 |row| row.get(0),
3836 ) {
3837 Ok(name) => Some(name),
3838 Err(rusqlite::Error::QueryReturnedNoRows) => None,
3839 Err(_) => return false,
3840 };
3841 let Some(name) = name else {
3842 return true;
3843 };
3844 let quoted = format!("\"{}\"", name.replace('"', "\"\""));
3845 if conn
3846 .execute_batch(&format!("DETACH DATABASE {quoted}"))
3847 .is_err()
3848 {
3849 return false;
3850 }
3851 }
3852}
3853
3854fn drop_temp_objects(conn: &Connection) -> bool {
3855 for (kind, ddl_keyword) in [
3859 ("view", "VIEW"),
3860 ("trigger", "TRIGGER"),
3861 ("index", "INDEX"),
3862 ("table", "TABLE"),
3863 ] {
3864 loop {
3865 let name: Option<String> = match conn.query_row(
3866 "SELECT name FROM sqlite_temp_master WHERE type = ?1 LIMIT 1",
3867 [kind],
3868 |row| row.get(0),
3869 ) {
3870 Ok(name) => Some(name),
3871 Err(rusqlite::Error::QueryReturnedNoRows) => None,
3872 Err(_) => return false,
3873 };
3874 let Some(name) = name else {
3875 break;
3876 };
3877 let quoted = format!("\"{}\"", name.replace('"', "\"\""));
3878 if conn
3879 .execute_batch(&format!("DROP {ddl_keyword} IF EXISTS temp.{quoted}"))
3880 .is_err()
3881 {
3882 return false;
3883 }
3884 }
3885 }
3886 true
3887}
3888
3889fn reset_observable_settings(
3890 conn: &Connection,
3891 config: &PoolConfig,
3892 expected_query_only: i64,
3893) -> Result<(), rusqlite::Error> {
3894 conn.pragma_update(None, "query_only", expected_query_only)?;
3895 conn.pragma_update(None, "writable_schema", 0)?;
3896 conn.pragma_update(None, "foreign_keys", "ON")?;
3897 conn.busy_timeout(config.busy_timeout)?;
3898 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3899 conn.pragma_update(None, "temp_store", "MEMORY")?;
3900 conn.pragma_update(None, "read_uncommitted", 0)?;
3901 conn.pragma_update(None, "defer_foreign_keys", 0)?;
3902 Ok(())
3903}
3904
3905fn reader_connection_is_healthy(conn: &Connection) -> bool {
3906 match conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0)) {
3907 Ok(_) => true,
3908 Err(rusqlite::Error::SqliteFailure(err, _)) => !matches!(
3909 err.code,
3910 rusqlite::ErrorCode::CannotOpen
3911 | rusqlite::ErrorCode::NotADatabase
3912 | rusqlite::ErrorCode::DatabaseCorrupt
3913 | rusqlite::ErrorCode::PermissionDenied
3914 | rusqlite::ErrorCode::SystemIoFailure
3915 ),
3916 Err(_) => true,
3917 }
3918}
3919
3920fn close_connection_quietly(conn: Connection) {
3921 match conn.close() {
3922 Ok(()) => {}
3923 Err((conn, _)) => drop(conn),
3924 }
3925}
3926
3927fn pool_exhausted_error(timeout: Duration, max_readers: usize) -> SqliteError {
3928 rusqlite::Error::SqliteFailure(
3929 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
3930 Some(format!(
3931 "Pool exhausted: no reader available after {timeout:?} (max_readers={max_readers})"
3932 )),
3933 )
3934 .into()
3935}
3936
3937#[cfg(test)]
3938#[path = "runtime_write_routing_tests.rs"]
3939mod runtime_write_routing_tests;
3940
3941#[cfg(test)]
3942#[path = "database_owner_identity_pool_tests.rs"]
3943mod database_owner_identity_pool_tests;
3944
3945#[cfg(test)]
3946#[path = "pool_tests.rs"]
3947mod tests;
3948
3949#[cfg(all(test, any(unix, windows)))]
3950#[path = "pool_identity_admission_tests.rs"]
3951mod identity_admission_tests;