1use crossbeam_queue::ArrayQueue;
3use parking_lot::{Condvar, Mutex};
4use rusqlite::hooks::{AuthContext, Authorization};
5use rusqlite::{Connection, OpenFlags};
6use sha2::{Digest, Sha256};
7use std::cell::Cell;
8use std::collections::HashMap;
9use std::fs;
10use std::io::Read as _;
11use std::ops::{Deref, DerefMut};
12use std::path::{Path, PathBuf};
13use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
14use std::sync::{Arc, OnceLock};
15use std::thread;
16use std::time::{Duration, Instant};
17use tokio::sync::Semaphore;
18
19use crate::error::SqliteError;
20use crate::writer_task::WriterTaskHandle;
21use khive_storage::error::StorageError;
22use khive_storage::tx_registry::{DbIdentity, TxOrigin};
23use khive_storage::StorageCapability;
24
25const CACHE_SIZE_KIB: &str = "-65536";
26const MMAP_SIZE_BYTES: &str = "1073741824";
27const DEFAULT_READER_CAP: usize = 8;
28
29const DEFAULT_JOURNAL_SIZE_LIMIT_BYTES: i64 = 67_108_864; const DEFAULT_WRITE_QUEUE_CAPACITY: usize = 256;
31static NEXT_MAIN_POOL_GENERATION: AtomicU64 = AtomicU64::new(1);
32
33struct OpenPoolIdentity {
34 count: usize,
35 basename: String,
36 suffix: String,
37}
38
39#[derive(Default)]
42struct PoolIdentityRegistry {
43 paths: HashMap<PathBuf, OpenPoolIdentity>,
44}
45
46fn pool_identity_registry() -> &'static Mutex<PoolIdentityRegistry> {
47 static REGISTRY: OnceLock<Mutex<PoolIdentityRegistry>> = OnceLock::new();
48 REGISTRY.get_or_init(|| Mutex::new(PoolIdentityRegistry::default()))
49}
50
51fn pool_identity_suffix(path: &Path) -> String {
55 #[cfg(unix)]
56 let bytes = {
57 use std::os::unix::ffi::OsStrExt;
58 path.as_os_str().as_bytes().to_vec()
59 };
60 #[cfg(windows)]
61 let bytes = {
62 use std::os::windows::ffi::OsStrExt;
63 path.as_os_str()
64 .encode_wide()
65 .flat_map(u16::to_le_bytes)
66 .collect::<Vec<_>>()
67 };
68 #[cfg(not(any(unix, windows)))]
69 let bytes = path.to_string_lossy().as_bytes().to_vec();
70 let digest = Sha256::digest(&bytes);
71 format!(
72 "{:02x}{:02x}{:02x}{:02x}",
73 digest[0], digest[1], digest[2], digest[3]
74 )
75}
76
77struct PoolIdentityRegistration(PathBuf);
78
79impl PoolIdentityRegistration {
80 fn new(path: &Path) -> Self {
81 let mut registry = pool_identity_registry().lock();
82 if let Some(entry) = registry.paths.get_mut(path) {
83 entry.count += 1;
84 } else {
85 let basename = path
86 .file_name()
87 .unwrap_or_default()
88 .to_string_lossy()
89 .into_owned();
90 let suffix = pool_identity_suffix(path);
91 registry.paths.insert(
92 path.to_path_buf(),
93 OpenPoolIdentity {
94 count: 1,
95 basename,
96 suffix,
97 },
98 );
99 }
100 Self(path.to_path_buf())
101 }
102
103 fn label(&self) -> String {
104 let registry = pool_identity_registry().lock();
105 let entry = ®istry.paths[&self.0];
106 let collides = registry
107 .paths
108 .iter()
109 .any(|(path, other)| path != &self.0 && other.basename == entry.basename);
110 if collides {
111 format!("{}#{}", entry.basename, entry.suffix)
112 } else {
113 entry.basename.clone()
114 }
115 }
116}
117
118impl Drop for PoolIdentityRegistration {
119 fn drop(&mut self) {
120 let mut registry = pool_identity_registry().lock();
121 if let Some(entry) = registry.paths.get_mut(&self.0) {
122 entry.count -= 1;
123 if entry.count == 0 {
124 registry.paths.remove(&self.0);
125 }
126 }
127 }
128}
129
130#[derive(Clone, Copy, Debug)]
132pub enum RuntimeWriteOperation {
133 MergeEntity,
134 MergeNote,
135 UpdateSymmetricEdge,
136}
137
138impl RuntimeWriteOperation {
139 fn operation(self) -> &'static str {
140 match self {
141 Self::MergeEntity => "merge_entity",
142 Self::MergeNote => "merge_note",
143 Self::UpdateSymmetricEdge => "update_edge",
144 }
145 }
146
147 fn fallback_site(self) -> crate::timeout_sink::Site {
148 match self {
149 Self::MergeEntity => crate::timeout_sink::Site::DirectRouteRuntimeMergeEntity,
150 Self::MergeNote => crate::timeout_sink::Site::DirectRouteRuntimeMergeNote,
151 Self::UpdateSymmetricEdge => {
152 crate::timeout_sink::Site::DirectRouteRuntimeUpdateSymmetricEdge
153 }
154 }
155 }
156}
157
158pub(crate) const FALLBACK_WAL_AUTOCHECKPOINT_PAGES: u32 = 4_000;
166
167#[derive(Clone, Copy, Debug, PartialEq, Eq)]
168enum CheckpointOwnership {
169 Unclaimed,
170 Claiming,
171 Claimed,
172}
173
174struct CheckpointOwnershipState {
175 phase: CheckpointOwnership,
176 #[cfg(test)]
177 connection_waiters: usize,
178}
179
180#[cfg(test)]
181struct CheckpointConnectionConfigPause {
182 selected: std::sync::Barrier,
183 resume: std::sync::Barrier,
184}
185
186#[cfg(test)]
187impl CheckpointConnectionConfigPause {
188 fn new() -> Self {
189 Self {
190 selected: std::sync::Barrier::new(2),
191 resume: std::sync::Barrier::new(2),
192 }
193 }
194}
195
196struct CheckpointOwnershipGate {
197 state: Mutex<CheckpointOwnershipState>,
198 changed: Condvar,
199 #[cfg(test)]
200 connection_config_pause: Mutex<Option<Arc<CheckpointConnectionConfigPause>>>,
201 #[cfg(test)]
202 claim_lock_observed: Mutex<Option<std::sync::mpsc::SyncSender<bool>>>,
203}
204
205impl CheckpointOwnershipGate {
206 fn new() -> Self {
207 Self {
208 state: Mutex::new(CheckpointOwnershipState {
209 phase: CheckpointOwnership::Unclaimed,
210 #[cfg(test)]
211 connection_waiters: 0,
212 }),
213 changed: Condvar::new(),
214 #[cfg(test)]
215 connection_config_pause: Mutex::new(None),
216 #[cfg(test)]
217 claim_lock_observed: Mutex::new(None),
218 }
219 }
220
221 fn begin_claim(&self) -> bool {
224 #[cfg(test)]
225 let claim_lock_observed = self.claim_lock_observed.lock().take();
226 #[cfg(test)]
227 let mut state = if let Some(observed) = claim_lock_observed {
228 match self.state.try_lock() {
229 Some(state) => {
230 let _ = observed.send(false);
231 state
232 }
233 None => {
234 let _ = observed.send(true);
235 self.state.lock()
236 }
237 }
238 } else {
239 self.state.lock()
240 };
241 #[cfg(not(test))]
242 let mut state = self.state.lock();
243 loop {
244 match state.phase {
245 CheckpointOwnership::Unclaimed => {
246 state.phase = CheckpointOwnership::Claiming;
247 self.changed.notify_all();
248 return true;
249 }
250 CheckpointOwnership::Claiming => self.changed.wait(&mut state),
251 CheckpointOwnership::Claimed => return false,
252 }
253 }
254 }
255
256 fn finish_claim(&self, succeeded: bool) {
257 let mut state = self.state.lock();
258 debug_assert_eq!(state.phase, CheckpointOwnership::Claiming);
259 state.phase = if succeeded {
260 CheckpointOwnership::Claimed
261 } else {
262 CheckpointOwnership::Unclaimed
263 };
264 self.changed.notify_all();
265 }
266
267 fn settled_state(&self) -> parking_lot::MutexGuard<'_, CheckpointOwnershipState> {
268 let mut state = self.state.lock();
269 while state.phase == CheckpointOwnership::Claiming {
270 #[cfg(test)]
271 {
272 state.connection_waiters += 1;
273 self.changed.notify_all();
274 }
275 self.changed.wait(&mut state);
276 #[cfg(test)]
277 {
278 state.connection_waiters -= 1;
279 self.changed.notify_all();
280 }
281 }
282 state
283 }
284
285 #[cfg(test)]
286 fn wal_autocheckpoint_pages(&self) -> u32 {
287 let state = self.settled_state();
288 match state.phase {
289 CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
290 CheckpointOwnership::Claimed => 0,
291 CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
292 }
293 }
294
295 fn configure_wal_autocheckpoint(&self, conn: &Connection) -> Result<(), SqliteError> {
300 let state = self.settled_state();
301 let pages = match state.phase {
302 CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
303 CheckpointOwnership::Claimed => 0,
304 CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
305 };
306 #[cfg(test)]
307 if let Some(pause) = self.connection_config_pause.lock().take() {
308 pause.selected.wait();
309 pause.resume.wait();
310 }
311 conn.pragma_update(None, "wal_autocheckpoint", pages)?;
312 drop(state);
313 Ok(())
314 }
315}
316
317fn deny_retired_writer(_context: AuthContext<'_>) -> Authorization {
318 Authorization::Deny
319}
320
321pub(crate) const TEST_HARNESS_ENV: &str = "KHIVE_TEST_HARNESS";
322
323#[derive(Clone, Debug)]
325pub struct PoolConfig {
326 pub path: Option<PathBuf>,
328 pub max_readers: usize,
330 pub wal_mode: bool,
332 pub busy_timeout: Duration,
336 pub checkout_timeout: Duration,
340 pub journal_size_limit_bytes: i64,
346 pub read_only: bool,
354 pub write_queue_enabled: Option<bool>,
382 pub write_queue_capacity: usize,
387 pub write_routing_strict: bool,
398 pub write_admission_deadline_ms: u64,
412 pub read_tx_max_age: Duration,
422}
423
424const WRITE_ADMISSION_DEADLINE_MS_RANGE: std::ops::RangeInclusive<u64> = 100..=10_000;
426const DEFAULT_WRITE_ADMISSION_DEADLINE_MS: u64 = 2000;
427
428impl Default for PoolConfig {
429 fn default() -> Self {
430 Self {
431 path: None,
432 max_readers: std::thread::available_parallelism()
433 .map(|n| n.get())
434 .unwrap_or(1)
435 .clamp(1, DEFAULT_READER_CAP),
436 wal_mode: true,
437 busy_timeout: Duration::from_secs(
438 std::env::var("KHIVE_BUSY_TIMEOUT_SECS")
439 .ok()
440 .and_then(|v| v.parse::<u64>().ok())
441 .unwrap_or(30),
442 ),
443 checkout_timeout: Duration::from_secs(
444 std::env::var("KHIVE_CHECKOUT_TIMEOUT_SECS")
445 .ok()
446 .and_then(|v| v.parse::<u64>().ok())
447 .unwrap_or(5),
448 ),
449 journal_size_limit_bytes: std::env::var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES")
450 .ok()
451 .and_then(|v| v.parse::<i64>().ok())
452 .unwrap_or(DEFAULT_JOURNAL_SIZE_LIMIT_BYTES),
453 read_only: false,
454 write_queue_enabled: std::env::var_os("KHIVE_WRITE_QUEUE").map(|v| {
459 v.to_str()
460 .is_some_and(|v| v == "1" || v.eq_ignore_ascii_case("true"))
461 }),
462 write_queue_capacity: std::env::var("KHIVE_WRITE_QUEUE_CAPACITY")
463 .ok()
464 .and_then(|v| v.parse::<usize>().ok())
465 .filter(|&n| n > 0)
466 .unwrap_or(DEFAULT_WRITE_QUEUE_CAPACITY),
467 write_routing_strict: std::env::var("KHIVE_WRITE_ROUTING")
468 .map(|v| v.eq_ignore_ascii_case("strict"))
469 .unwrap_or(false),
470 write_admission_deadline_ms: std::env::var("KHIVE_WRITE_ADMISSION_DEADLINE_MS")
471 .ok()
472 .and_then(|v| v.parse::<u64>().ok())
473 .unwrap_or(DEFAULT_WRITE_ADMISSION_DEADLINE_MS),
474 read_tx_max_age: crate::checkpoint::tx_age_thresholds_from_env(
475 Duration::from_secs(30),
476 Duration::from_secs(120),
477 )
478 .1,
479 }
480 }
481}
482
483#[cfg(any(test, feature = "test-support"))]
484impl PoolConfig {
485 pub fn for_test() -> Self {
490 Self {
491 max_readers: 2,
492 ..Self::default()
493 }
494 }
495}
496
497fn refuse_home_data_store_in_tests(config: &PoolConfig) -> Result<(), SqliteError> {
514 if std::env::var(TEST_HARNESS_ENV).as_deref() != Ok("1") {
515 return Ok(());
516 }
517
518 let Some(path) = config.path.as_deref() else {
519 return Ok(());
520 };
521 if path
522 .as_os_str()
523 .as_encoded_bytes()
524 .get(..5)
525 .is_some_and(|prefix| prefix.eq_ignore_ascii_case(b"file:"))
526 {
527 return Err(SqliteError::InvalidData(format!(
528 "test harness refused SQLite URI database path {}; use a filesystem path outside \
529 HOME/.khive (deliberate sessions against a real store run the built binary \
530 directly, outside the Cargo test environment)",
531 path.display()
532 )));
533 }
534
535 let Some(home) = std::env::var_os("HOME") else {
536 return Ok(());
537 };
538 let canonical_path = canonicalize_deepest_existing(path)?;
539 let canonical_home_data_dir =
540 canonicalize_deepest_existing(&PathBuf::from(home).join(".khive"))?;
541 if canonical_path.starts_with(&canonical_home_data_dir) {
542 return Err(SqliteError::InvalidData(format!(
543 "test harness refused to open SQLite database under HOME/.khive: {} \
544 (deliberate sessions against a real store run the built binary directly, \
545 outside the Cargo test environment)",
546 canonical_path.display()
547 )));
548 }
549 Ok(())
550}
551
552fn canonicalize_deepest_existing(path: &Path) -> Result<PathBuf, SqliteError> {
553 let absolute = if path.is_absolute() {
554 path.to_path_buf()
555 } else {
556 std::env::current_dir().map_err(SqliteError::Io)?.join(path)
557 };
558
559 for ancestor in absolute.ancestors() {
560 match fs::canonicalize(ancestor) {
561 Ok(mut canonical) => {
562 let missing = absolute.strip_prefix(ancestor).map_err(|error| {
563 SqliteError::InvalidData(format!(
564 "failed to preserve missing path components for {}: {error}",
565 absolute.display()
566 ))
567 })?;
568 canonical.push(missing);
569 return Ok(canonical);
570 }
571 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
572 Err(error) => {
573 return Err(SqliteError::InvalidData(format!(
574 "failed to canonicalize database path ancestor {}: {error}",
575 ancestor.display()
576 )));
577 }
578 }
579 }
580
581 Err(SqliteError::InvalidData(format!(
582 "database path has no canonicalizable ancestor: {}",
583 absolute.display()
584 )))
585}
586
587fn validate_write_admission_deadline(deadline_ms: u64) -> Result<(), SqliteError> {
592 if WRITE_ADMISSION_DEADLINE_MS_RANGE.contains(&deadline_ms) {
593 return Ok(());
594 }
595 Err(SqliteError::InvalidConfig(format!(
596 "write_admission_deadline_ms must be in [{}, {}] ms, got {deadline_ms}",
597 WRITE_ADMISSION_DEADLINE_MS_RANGE.start(),
598 WRITE_ADMISSION_DEADLINE_MS_RANGE.end()
599 )))
600}
601
602pub struct ConnectionPool {
616 writer: Arc<Mutex<Connection>>,
617 main_pool_generation: OnceLock<u64>,
618 checkpoint_ownership: CheckpointOwnershipGate,
627 pooled_writer_retired: AtomicBool,
632 writer_acquisition_counters: Arc<WriterAcquisitionCounters>,
637 reader_acquisition_counters: ReaderAcquisitionCounters,
642 readers: ArrayQueue<Connection>,
643 max_readers: usize,
644 config: PoolConfig,
645 read_only_open_target: Option<PathBuf>,
653 sql_bridge_reader_slots: Arc<Semaphore>,
654 sql_bridge_writer_slots: Arc<Semaphore>,
655 writer_task: OnceLock<Option<WriterTaskHandle>>,
660 writer_task_join: Mutex<Option<tokio::task::JoinHandle<()>>>,
668 writer_task_join_stored: AtomicBool,
673 origin: TxOrigin,
679 identity_path: Option<PathBuf>,
685 identity_registration: Option<PoolIdentityRegistration>,
688 #[cfg(test)]
694 writer_task_spawn_count: std::sync::atomic::AtomicUsize,
695}
696
697impl Drop for ConnectionPool {
698 fn drop(&mut self) {
706 while let Some(conn) = self.readers.pop() {
707 drop(conn);
708 }
709 }
710}
711
712enum ReaderLease<'pool> {
713 Pooled(Connection),
714 Shared(parking_lot::MutexGuard<'pool, Connection>),
715}
716
717pub struct ReaderRow<'row, 'statement> {
727 row: &'row rusqlite::Row<'statement>,
728}
729
730impl ReaderRow<'_, '_> {
731 pub fn get<I: rusqlite::RowIndex, T: rusqlite::types::FromSql>(
733 &self,
734 index: I,
735 ) -> rusqlite::Result<T> {
736 self.row.get(index)
737 }
738
739 pub fn get_ref<I: rusqlite::RowIndex>(
741 &self,
742 index: I,
743 ) -> rusqlite::Result<rusqlite::types::ValueRef<'_>> {
744 self.row.get_ref(index)
745 }
746}
747
748struct ReaderQueryInProgress<'a>(&'a Cell<bool>);
750
751impl Drop for ReaderQueryInProgress<'_> {
752 fn drop(&mut self) {
753 self.0.set(false);
754 }
755}
756
757pub struct ReaderGuard<'pool> {
760 lease: Option<ReaderLease<'pool>>,
761 admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
765 pool: &'pool ConnectionPool,
766 reusable: Cell<bool>,
767 query_in_progress: Cell<bool>,
768 checked_out_at: Instant,
769 dirty: Cell<bool>,
777 operation: Option<&'static str>,
783}
784
785impl<'pool> ReaderGuard<'pool> {
786 pub(crate) fn conn(&self) -> &Connection {
795 match self
796 .lease
797 .as_ref()
798 .expect("reader guard missing connection")
799 {
800 ReaderLease::Pooled(conn) => conn,
801 ReaderLease::Shared(guard) => guard,
802 }
803 }
804
805 pub fn query_row<T, P, F>(&self, sql: &str, params: P, f: F) -> Result<T, SqliteError>
824 where
825 P: rusqlite::Params,
826 F: FnOnce(&ReaderRow<'_, '_>) -> rusqlite::Result<T>,
827 {
828 crate::sql_bridge::reader_capability_admits(sql).map_err(SqliteError::InvalidData)?;
829 if !self.reusable.get() {
830 return Err(SqliteError::InvalidData(
831 "reader lease is quarantined after failed read cleanup".into(),
832 ));
833 }
834 if self.query_in_progress.replace(true) {
838 return Err(SqliteError::InvalidData(
839 "reader lease is already executing a query".into(),
840 ));
841 }
842 let _in_progress = ReaderQueryInProgress(&self.query_in_progress);
843 self.mark_dirty();
844 crate::read_cancellation::run_borrowed_reader(self, |conn, admission| {
845 conn.query_row(sql, params, |row| {
846 if !admission.admits() {
854 return Err(rusqlite::Error::SqliteFailure(
855 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_INTERRUPT),
856 Some("request stopped before the mapper".into()),
857 ));
858 }
859 f(&ReaderRow { row })
860 })
861 .map_err(|error| {
862 StorageError::driver(StorageCapability::Sql, "reader_guard.query_row", error)
863 })
864 })
865 .map_err(|error| match error {
866 StorageError::Driver {
867 capability,
868 operation,
869 source,
870 } => match source.downcast::<rusqlite::Error>() {
871 Ok(error) => SqliteError::Rusqlite(*error),
872 Err(source) => SqliteError::RequestReadStopped(StorageError::Driver {
873 capability,
874 operation,
875 source,
876 }),
877 },
878 other => SqliteError::RequestReadStopped(other),
879 })
880 }
881
882 pub(crate) fn discard(&self) {
886 self.reusable.set(false);
887 }
888
889 pub(crate) fn mark_dirty(&self) {
892 self.dirty.set(true);
893 }
894
895 pub(crate) fn label_operation(&mut self, operation: &'static str) {
899 self.operation = Some(operation);
900 }
901}
902
903impl<'pool> Drop for ReaderGuard<'pool> {
904 fn drop(&mut self) {
905 let Some(lease) = self.lease.take() else {
906 return;
907 };
908
909 match lease {
910 ReaderLease::Pooled(conn) if self.reusable.get() => {
911 self.pool.return_reader(conn, self.dirty.get())
912 }
913 ReaderLease::Pooled(conn) => {
914 close_connection_quietly(conn);
915 self.pool.replace_discarded_reader_slot();
916 }
917 ReaderLease::Shared(guard) if !self.reusable.get() => {
918 self.pool.retire_pooled_writer(&guard);
919 }
920 ReaderLease::Shared(guard) => {
921 if self.dirty.get() && !restore_shared_reader_state(&guard, &self.pool.config) {
929 self.pool.retire_pooled_writer(&guard);
930 }
931 }
932 }
933
934 drop(self.admission_slot.take());
938 self.pool
939 .reader_acquisition_counters
940 .record_checkout_completed(self.checked_out_at.elapsed(), self.operation);
941 }
942}
943
944pub(crate) struct SharedReaderTransactionGuard {
960 conn: parking_lot::ArcMutexGuard<parking_lot::RawMutex, Connection>,
961 admission_slot: Option<tokio::sync::OwnedSemaphorePermit>,
962 pool: Arc<ConnectionPool>,
963 checked_out_at: Instant,
964 poison: Cell<bool>,
970}
971
972impl SharedReaderTransactionGuard {
973 pub(crate) fn conn(&self) -> &Connection {
974 &self.conn
975 }
976
977 pub(crate) fn poison(&self) {
980 self.poison.set(true);
981 }
982}
983
984impl Drop for SharedReaderTransactionGuard {
985 fn drop(&mut self) {
986 let mut restored = !self.poison.get();
993 if restored && !self.conn.is_autocommit() {
994 restored = self.conn.execute_batch("ROLLBACK").is_ok() && self.conn.is_autocommit();
995 }
996 if restored {
997 restored = restore_shared_reader_state(&self.conn, &self.pool.config);
998 }
999 if !restored {
1000 self.pool.retire_pooled_writer(&self.conn);
1001 }
1002 drop(self.admission_slot.take());
1003 self.pool
1004 .reader_acquisition_counters
1005 .record_checkout_completed(
1006 self.checked_out_at.elapsed(),
1007 Some("explicit_sql_read_transaction"),
1008 );
1009 }
1010}
1011
1012impl ConnectionPool {
1013 pub(crate) fn checkout_shared_reader_transaction(
1024 self: &Arc<Self>,
1025 should_stop: impl Fn() -> bool,
1026 ) -> Result<Option<SharedReaderTransactionGuard>, SqliteError> {
1027 debug_assert_eq!(
1028 self.max_readers, 0,
1029 "the owned shared-reader-transaction guard exists only for the degraded, \
1030 single-connection backend"
1031 );
1032 self.ensure_pooled_writer_active()?;
1033 let started = Instant::now();
1034 let admission_slot = loop {
1035 if should_stop() {
1036 return Ok(None);
1037 }
1038 match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
1039 Ok(slot) => break slot,
1040 Err(tokio::sync::TryAcquireError::Closed) => {
1041 return Err(SqliteError::InvalidData(
1042 "reader admission semaphore is closed".to_string(),
1043 ));
1044 }
1045 Err(tokio::sync::TryAcquireError::NoPermits) => {}
1046 }
1047 if started.elapsed() >= self.config.checkout_timeout {
1048 self.reader_acquisition_counters.record_checkout_timeout();
1049 return Err(pool_exhausted_error(
1050 self.config.checkout_timeout,
1051 self.max_readers,
1052 ));
1053 }
1054 thread::yield_now();
1055 };
1056
1057 loop {
1058 if should_stop() {
1059 return Ok(None);
1060 }
1061 let remaining = self
1062 .config
1063 .checkout_timeout
1064 .saturating_sub(started.elapsed());
1065 if remaining.is_zero() {
1066 self.reader_acquisition_counters.record_checkout_timeout();
1067 return Err(pool_exhausted_error(
1068 self.config.checkout_timeout,
1069 self.max_readers,
1070 ));
1071 }
1072 if let Some(conn) = self
1073 .writer
1074 .try_lock_arc_for(remaining.min(Duration::from_millis(2)))
1075 {
1076 self.ensure_pooled_writer_active()?;
1077 self.reader_acquisition_counters.record_pooled_checkout();
1078 return Ok(Some(SharedReaderTransactionGuard {
1079 conn,
1080 admission_slot: Some(admission_slot),
1081 pool: Arc::clone(self),
1082 checked_out_at: Instant::now(),
1083 poison: Cell::new(false),
1084 }));
1085 }
1086 }
1087 }
1088}
1089
1090#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1096#[allow(dead_code)] pub(crate) enum StandaloneReaderPurpose {
1098 ExplicitSqlReadTransaction,
1100 BootSchemaProbe,
1103 DiagnosticsIndependentSnapshot,
1106}
1107
1108impl StandaloneReaderPurpose {
1109 fn is_infrastructure(self) -> bool {
1110 matches!(
1111 self,
1112 Self::BootSchemaProbe | Self::DiagnosticsIndependentSnapshot
1113 )
1114 }
1115}
1116
1117#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
1125pub struct ReaderAcquisitionSnapshot {
1126 pub reader_admission_capacity: usize,
1129 pub available_reader_admission_slots: usize,
1132 pub acquisitions: u64,
1135 pub pooled_checkouts: u64,
1137 pub standalone_opens: u64,
1140 pub infrastructure_standalone_opens: u64,
1143 pub checkout_timeouts: u64,
1147 pub active_pooled_checkouts: u64,
1149 pub peak_active_pooled_checkouts: u64,
1151 pub completed_pooled_checkouts: u64,
1153 pub max_completed_hold_micros: u64,
1156 pub max_completed_hold_operation: Option<&'static str>,
1161 pub reader_replacement_open_failures: u64,
1167}
1168
1169#[derive(Debug, Default, Clone, Copy)]
1173struct LongestCompletedHold {
1174 micros: u64,
1175 operation: Option<&'static str>,
1176}
1177
1178#[derive(Debug, Default)]
1179struct ReaderAcquisitionCounters {
1180 pooled_checkouts: AtomicU64,
1181 standalone_opens: AtomicU64,
1182 infrastructure_standalone_opens: AtomicU64,
1183 checkout_timeouts: AtomicU64,
1184 active_pooled_checkouts: AtomicU64,
1185 peak_active_pooled_checkouts: AtomicU64,
1186 completed_pooled_checkouts: AtomicU64,
1187 longest_completed_hold: parking_lot::Mutex<LongestCompletedHold>,
1188 reader_replacement_open_failures: AtomicU64,
1189}
1190
1191impl ReaderAcquisitionCounters {
1192 fn record_pooled_checkout(&self) {
1193 self.pooled_checkouts.fetch_add(1, Ordering::Relaxed);
1194 let active = self
1195 .active_pooled_checkouts
1196 .fetch_add(1, Ordering::Relaxed)
1197 .saturating_add(1);
1198 self.peak_active_pooled_checkouts
1199 .fetch_max(active, Ordering::Relaxed);
1200 }
1201
1202 fn record_checkout_timeout(&self) {
1203 self.checkout_timeouts.fetch_add(1, Ordering::Relaxed);
1204 }
1205
1206 fn record_reader_replacement_open_failure(&self) {
1207 self.reader_replacement_open_failures
1208 .fetch_add(1, Ordering::Relaxed);
1209 }
1210
1211 fn record_standalone_open(&self, purpose: StandaloneReaderPurpose) {
1212 if purpose.is_infrastructure() {
1213 self.infrastructure_standalone_opens
1214 .fetch_add(1, Ordering::Relaxed);
1215 } else {
1216 self.standalone_opens.fetch_add(1, Ordering::Relaxed);
1217 }
1218 }
1219
1220 fn record_checkout_completed(&self, hold: Duration, operation: Option<&'static str>) {
1221 let previous = self.active_pooled_checkouts.fetch_sub(1, Ordering::Relaxed);
1222 debug_assert!(previous > 0, "reader active-checkout counter underflow");
1223 self.completed_pooled_checkouts
1224 .fetch_add(1, Ordering::Relaxed);
1225 let micros = u64::try_from(hold.as_micros()).unwrap_or(u64::MAX);
1226 let mut longest = self.longest_completed_hold.lock();
1230 if micros > longest.micros {
1231 longest.micros = micros;
1232 longest.operation = operation;
1233 }
1234 }
1235
1236 fn snapshot(
1237 &self,
1238 reader_admission_capacity: usize,
1239 available_reader_admission_slots: usize,
1240 ) -> ReaderAcquisitionSnapshot {
1241 let pooled_checkouts = self.pooled_checkouts.load(Ordering::Relaxed);
1242 let standalone_opens = self.standalone_opens.load(Ordering::Relaxed);
1243 let longest_completed_hold = *self.longest_completed_hold.lock();
1244 ReaderAcquisitionSnapshot {
1245 reader_admission_capacity,
1246 available_reader_admission_slots,
1247 acquisitions: pooled_checkouts.saturating_add(standalone_opens),
1248 pooled_checkouts,
1249 standalone_opens,
1250 infrastructure_standalone_opens: self
1251 .infrastructure_standalone_opens
1252 .load(Ordering::Relaxed),
1253 checkout_timeouts: self.checkout_timeouts.load(Ordering::Relaxed),
1254 active_pooled_checkouts: self.active_pooled_checkouts.load(Ordering::Relaxed),
1255 peak_active_pooled_checkouts: self.peak_active_pooled_checkouts.load(Ordering::Relaxed),
1256 completed_pooled_checkouts: self.completed_pooled_checkouts.load(Ordering::Relaxed),
1257 max_completed_hold_micros: longest_completed_hold.micros,
1258 max_completed_hold_operation: longest_completed_hold.operation,
1259 reader_replacement_open_failures: self
1260 .reader_replacement_open_failures
1261 .load(Ordering::Relaxed),
1262 }
1263 }
1264}
1265
1266pub struct WriterGuard<'pool> {
1269 guard: parking_lot::MutexGuard<'pool, Connection>,
1270 origin: TxOrigin,
1274}
1275
1276#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
1285pub struct WriterAcquisitionSnapshot {
1286 pub acquisitions: u64,
1289 pub pooled_acquisitions: u64,
1291 pub standalone_acquisitions: u64,
1293 pub writer_task_acquisitions: u64,
1296 pub timeouts: u64,
1298 pub writer_task_begin_busy: u64,
1303 pub writer_task_begin_busy_absorbed: u64,
1308 pub writer_task_begin_errors: u64,
1311 pub writer_task_request_failures: u64,
1315 pub writer_task_side_effects_unknown: u64,
1320}
1321
1322#[derive(Debug, Default)]
1326pub(crate) struct WriterAcquisitionCounters {
1327 pooled_acquisitions: AtomicU64,
1328 standalone_acquisitions: AtomicU64,
1329 writer_task_acquisitions: AtomicU64,
1330 pooled_timeouts: AtomicU64,
1331 writer_task_begin_busy: AtomicU64,
1332 writer_task_begin_busy_absorbed: AtomicU64,
1333 writer_task_begin_errors: AtomicU64,
1334 writer_task_request_failures: AtomicU64,
1335 writer_task_side_effects_unknown: AtomicU64,
1336}
1337
1338impl WriterAcquisitionCounters {
1339 pub(crate) fn record_writer_task_acquisition(&self) {
1340 self.writer_task_acquisitions
1341 .fetch_add(1, Ordering::Relaxed);
1342 }
1343
1344 pub(crate) fn record_writer_task_begin_busy(&self) {
1350 self.writer_task_begin_busy.fetch_add(1, Ordering::Relaxed);
1351 }
1352
1353 pub(crate) fn record_writer_task_begin_busy_absorbed(&self) {
1359 self.writer_task_begin_busy_absorbed
1360 .fetch_add(1, Ordering::Relaxed);
1361 }
1362
1363 pub(crate) fn record_writer_task_begin_error(&self) {
1367 self.writer_task_begin_errors
1368 .fetch_add(1, Ordering::Relaxed);
1369 }
1370
1371 pub(crate) fn record_writer_task_request_failure(&self) {
1375 self.writer_task_request_failures
1376 .fetch_add(1, Ordering::Relaxed);
1377 }
1378
1379 pub(crate) fn record_writer_task_side_effects_unknown(&self) {
1384 self.writer_task_side_effects_unknown
1385 .fetch_add(1, Ordering::Relaxed);
1386 }
1387
1388 fn snapshot(&self) -> WriterAcquisitionSnapshot {
1389 let pooled_acquisitions = self.pooled_acquisitions.load(Ordering::Relaxed);
1390 let standalone_acquisitions = self.standalone_acquisitions.load(Ordering::Relaxed);
1391 let writer_task_acquisitions = self.writer_task_acquisitions.load(Ordering::Relaxed);
1392 WriterAcquisitionSnapshot {
1393 acquisitions: pooled_acquisitions
1394 .saturating_add(standalone_acquisitions)
1395 .saturating_add(writer_task_acquisitions),
1396 pooled_acquisitions,
1397 standalone_acquisitions,
1398 writer_task_acquisitions,
1399 timeouts: self.pooled_timeouts.load(Ordering::Relaxed),
1400 writer_task_begin_busy: self.writer_task_begin_busy.load(Ordering::Relaxed),
1401 writer_task_begin_busy_absorbed: self
1402 .writer_task_begin_busy_absorbed
1403 .load(Ordering::Relaxed),
1404 writer_task_begin_errors: self.writer_task_begin_errors.load(Ordering::Relaxed),
1405 writer_task_request_failures: self.writer_task_request_failures.load(Ordering::Relaxed),
1406 writer_task_side_effects_unknown: self
1407 .writer_task_side_effects_unknown
1408 .load(Ordering::Relaxed),
1409 }
1410 }
1411}
1412
1413impl<'pool> WriterGuard<'pool> {
1414 pub fn conn(&self) -> &Connection {
1416 &self.guard
1417 }
1418
1419 pub fn conn_mut(&mut self) -> &mut Connection {
1421 &mut self.guard
1422 }
1423
1424 pub fn transaction<F, R>(&self, f: F) -> Result<R, SqliteError>
1427 where
1428 F: FnOnce(&Connection) -> Result<R, SqliteError>,
1429 {
1430 self.guard.execute_batch("BEGIN IMMEDIATE")?;
1431 let _tx_handle = khive_storage::tx_registry::register_scoped(
1432 Some("writer_guard_tx".to_string()),
1433 self.origin.clone(),
1434 );
1435
1436 match f(&self.guard) {
1437 Ok(result) => {
1438 if let Err(err) = self.guard.execute_batch("COMMIT") {
1439 let _ = self.guard.execute_batch("ROLLBACK");
1440 return Err(err.into());
1441 }
1442 Ok(result)
1443 }
1444 Err(err) => {
1445 let _ = self.guard.execute_batch("ROLLBACK");
1446 Err(err)
1447 }
1448 }
1449 }
1450}
1451
1452impl<'pool> Deref for WriterGuard<'pool> {
1453 type Target = Connection;
1454
1455 fn deref(&self) -> &Self::Target {
1456 self.conn()
1457 }
1458}
1459
1460impl<'pool> DerefMut for WriterGuard<'pool> {
1461 fn deref_mut(&mut self) -> &mut Self::Target {
1462 self.conn_mut()
1463 }
1464}
1465
1466impl ConnectionPool {
1467 pub fn new(config: PoolConfig) -> Result<Self, SqliteError> {
1475 refuse_home_data_store_in_tests(&config)?;
1476 validate_write_admission_deadline(config.write_admission_deadline_ms)?;
1477
1478 let mut config = config;
1482 let inert_memory_queue_request =
1483 config.path.is_none() && config.write_queue_enabled == Some(true);
1484 config.write_queue_enabled =
1485 Some(config.write_queue_enabled.unwrap_or(config.path.is_some()));
1486 if inert_memory_queue_request {
1487 tracing::warn!(
1488 "write queue explicitly requested for an in-memory pool; it is inert because \
1489 in-memory pools cannot host a writer task"
1490 );
1491 }
1492
1493 let (origin, identity_path) = match config.path.as_ref() {
1498 Some(path) => {
1499 let (identity, canonical) = mint_db_identity(path)?;
1500 (TxOrigin::Database(identity), Some(canonical))
1501 }
1502 None => (TxOrigin::Memory, None),
1503 };
1504 let read_only_open_target = read_only_open_target(&config, identity_path.as_deref())?;
1505 let writer = open_writer_connection(&config, read_only_open_target.as_deref())?;
1506 let wal_enabled = configure_writer_connection(&writer, &config)?;
1507 let max_readers = effective_reader_count(&config, wal_enabled);
1508
1509 let readers = ArrayQueue::new(max_readers.max(1));
1510
1511 let mut pool = Self {
1512 writer: Arc::new(Mutex::new(writer)),
1513 main_pool_generation: OnceLock::new(),
1514 checkpoint_ownership: CheckpointOwnershipGate::new(),
1515 pooled_writer_retired: AtomicBool::new(false),
1516 writer_acquisition_counters: Arc::new(WriterAcquisitionCounters::default()),
1517 reader_acquisition_counters: ReaderAcquisitionCounters::default(),
1518 readers,
1519 max_readers,
1520 config,
1521 read_only_open_target,
1522 sql_bridge_reader_slots: Arc::new(Semaphore::new(max_readers.max(1))),
1523 sql_bridge_writer_slots: Arc::new(Semaphore::new(1)),
1524 writer_task: OnceLock::new(),
1525 writer_task_join: Mutex::new(None),
1526 writer_task_join_stored: AtomicBool::new(false),
1527 origin,
1528 identity_path,
1529 identity_registration: None,
1530 #[cfg(test)]
1531 writer_task_spawn_count: std::sync::atomic::AtomicUsize::new(0),
1532 };
1533
1534 for _ in 0..pool.max_readers {
1535 let conn = pool.open_reader_connection()?;
1536 pool.readers
1537 .push(conn)
1538 .expect("reader queue must have capacity during pool initialization");
1539 }
1540
1541 if !pool.config.read_only {
1546 crate::timeout_sink::init(
1547 pool.canonical_path().and_then(Path::parent),
1548 &crate::timeout_sink::db_label(&pool),
1549 );
1550 }
1551
1552 pool.identity_registration = pool.canonical_path().map(PoolIdentityRegistration::new);
1553 Ok(pool)
1554 }
1555
1556 pub fn reader(&self) -> Result<ReaderGuard<'_>, SqliteError> {
1566 self.reader_until(|| false)?.ok_or_else(|| {
1567 SqliteError::InvalidData("uncancelled reader checkout stopped unexpectedly".into())
1568 })
1569 }
1570
1571 pub(crate) fn reader_until<C>(
1577 &self,
1578 should_stop: C,
1579 ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
1580 where
1581 C: Fn() -> bool,
1582 {
1583 let started = Instant::now();
1584 let mut admission_attempt = 0u32;
1585 let admission_slot = loop {
1586 if should_stop() {
1587 return Ok(None);
1588 }
1589 match Arc::clone(&self.sql_bridge_reader_slots).try_acquire_owned() {
1590 Ok(slot) => break slot,
1591 Err(tokio::sync::TryAcquireError::Closed) => {
1592 return Err(SqliteError::InvalidData(
1593 "reader admission semaphore is closed".to_string(),
1594 ));
1595 }
1596 Err(tokio::sync::TryAcquireError::NoPermits) => {}
1597 }
1598 if started.elapsed() >= self.config.checkout_timeout {
1599 self.reader_acquisition_counters.record_checkout_timeout();
1600 return Err(pool_exhausted_error(
1601 self.config.checkout_timeout,
1602 self.max_readers,
1603 ));
1604 }
1605 match admission_attempt {
1606 0..=7 => {
1607 let spins = 1usize << admission_attempt;
1608 for _ in 0..spins {
1609 std::hint::spin_loop();
1610 }
1611 }
1612 8..=15 => thread::yield_now(),
1613 _ => {
1614 let remaining = self
1615 .config
1616 .checkout_timeout
1617 .saturating_sub(started.elapsed());
1618 let sleep =
1619 Duration::from_micros(50 * (1u64 << (admission_attempt - 16).min(6)));
1620 thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
1621 }
1622 }
1623 admission_attempt = admission_attempt.saturating_add(1);
1624 };
1625
1626 if self.max_readers == 0 {
1627 self.ensure_pooled_writer_active()?;
1628 loop {
1629 if should_stop() {
1630 return Ok(None);
1631 }
1632 let remaining = self
1633 .config
1634 .checkout_timeout
1635 .saturating_sub(started.elapsed());
1636 if remaining.is_zero() {
1637 self.reader_acquisition_counters.record_checkout_timeout();
1638 return Err(pool_exhausted_error(
1639 self.config.checkout_timeout,
1640 self.max_readers,
1641 ));
1642 }
1643 if let Some(guard) = self
1644 .writer
1645 .try_lock_for(remaining.min(Duration::from_millis(2)))
1646 {
1647 self.ensure_pooled_writer_active()?;
1648 self.reader_acquisition_counters.record_pooled_checkout();
1649 return Ok(Some(ReaderGuard {
1650 lease: Some(ReaderLease::Shared(guard)),
1651 admission_slot: Some(admission_slot),
1652 pool: self,
1653 reusable: Cell::new(true),
1654 query_in_progress: Cell::new(false),
1655 checked_out_at: Instant::now(),
1656 dirty: Cell::new(false),
1657 operation: None,
1658 }));
1659 }
1660 }
1661 }
1662
1663 let mut attempt = 0u32;
1664
1665 loop {
1666 if should_stop() {
1667 return Ok(None);
1668 }
1669 if let Some(conn) = self.readers.pop() {
1670 self.reader_acquisition_counters.record_pooled_checkout();
1671 return Ok(Some(ReaderGuard {
1672 lease: Some(ReaderLease::Pooled(conn)),
1673 admission_slot: Some(admission_slot),
1674 pool: self,
1675 reusable: Cell::new(true),
1676 query_in_progress: Cell::new(false),
1677 checked_out_at: Instant::now(),
1678 dirty: Cell::new(false),
1679 operation: None,
1680 }));
1681 }
1682
1683 if started.elapsed() >= self.config.checkout_timeout {
1684 self.reader_acquisition_counters.record_checkout_timeout();
1685 return Err(pool_exhausted_error(
1686 self.config.checkout_timeout,
1687 self.max_readers,
1688 ));
1689 }
1690
1691 match attempt {
1692 0..=7 => {
1693 let spins = 1usize << attempt;
1694 for _ in 0..spins {
1695 std::hint::spin_loop();
1696 }
1697 }
1698 8..=15 => thread::yield_now(),
1699 _ => {
1700 let remaining = self
1701 .config
1702 .checkout_timeout
1703 .saturating_sub(started.elapsed());
1704 let sleep = Duration::from_micros(50 * (1u64 << (attempt - 16).min(6)));
1705 thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
1706 }
1707 }
1708
1709 attempt = attempt.saturating_add(1);
1710 }
1711 }
1712
1713 pub fn writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
1719 self.ensure_pooled_writer_active()?;
1720 let Some(guard) = self.writer.try_lock_for(self.config.checkout_timeout) else {
1721 self.writer_acquisition_counters
1722 .pooled_timeouts
1723 .fetch_add(1, Ordering::Relaxed);
1724 let message = format!(
1725 "timed out after {:?} waiting for sqlite writer connection",
1726 self.config.checkout_timeout
1727 );
1728 crate::timeout_sink::emit_timeout(
1729 &crate::timeout_sink::db_label(self),
1730 crate::timeout_sink::Site::PoolAdmission,
1731 &message,
1732 Some(
1733 self.config
1734 .checkout_timeout
1735 .as_millis()
1736 .min(u128::from(u64::MAX)) as u64,
1737 ),
1738 );
1739 return Err(SqliteError::WriterPoolCheckoutTimeout {
1740 timeout: self.config.checkout_timeout,
1741 });
1742 };
1743 self.ensure_pooled_writer_active()?;
1744 self.writer_acquisition_counters
1745 .pooled_acquisitions
1746 .fetch_add(1, Ordering::Relaxed);
1747 Ok(WriterGuard {
1748 guard,
1749 origin: self.origin(),
1750 })
1751 }
1752
1753 pub fn try_writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
1758 self.writer()
1759 }
1760
1761 pub(crate) fn writer_until<C>(
1762 &self,
1763 should_stop: C,
1764 ) -> Result<Option<WriterGuard<'_>>, SqliteError>
1765 where
1766 C: Fn() -> bool,
1767 {
1768 self.ensure_pooled_writer_active()?;
1769 let started = Instant::now();
1770 loop {
1771 if should_stop() {
1772 return Ok(None);
1773 }
1774 let remaining = self
1775 .config
1776 .checkout_timeout
1777 .saturating_sub(started.elapsed());
1778 if let Some(guard) = self
1779 .writer
1780 .try_lock_for(remaining.min(Duration::from_millis(2)))
1781 {
1782 if should_stop() {
1785 return Ok(None);
1786 }
1787 self.ensure_pooled_writer_active()?;
1788 self.writer_acquisition_counters
1789 .pooled_acquisitions
1790 .fetch_add(1, Ordering::Relaxed);
1791 return Ok(Some(WriterGuard {
1792 guard,
1793 origin: self.origin(),
1794 }));
1795 }
1796 if started.elapsed() >= self.config.checkout_timeout {
1797 self.writer_acquisition_counters
1798 .pooled_timeouts
1799 .fetch_add(1, Ordering::Relaxed);
1800 let message = format!(
1801 "timed out after {:?} waiting for sqlite writer connection",
1802 self.config.checkout_timeout
1803 );
1804 crate::timeout_sink::emit_timeout(
1805 &crate::timeout_sink::db_label(self),
1806 crate::timeout_sink::Site::PoolAdmission,
1807 &message,
1808 Some(
1809 self.config
1810 .checkout_timeout
1811 .as_millis()
1812 .min(u128::from(u64::MAX)) as u64,
1813 ),
1814 );
1815 return Err(SqliteError::WriterPoolCheckoutTimeout {
1816 timeout: self.config.checkout_timeout,
1817 });
1818 }
1819 }
1820 }
1821
1822 pub fn try_writer_nowait(&self) -> Result<WriterGuard<'_>, SqliteError> {
1831 self.ensure_pooled_writer_active()?;
1832 let guard = self.writer.try_lock().ok_or_else(|| {
1833 SqliteError::InvalidData(
1834 "writer connection busy (checkpoint skipped this tick)".to_string(),
1835 )
1836 })?;
1837 self.ensure_pooled_writer_active()?;
1838 Ok(WriterGuard {
1839 guard,
1840 origin: self.origin(),
1841 })
1842 }
1843
1844 pub(crate) fn retire_pooled_writer(&self, conn: &Connection) {
1845 self.pooled_writer_retired.store(true, Ordering::Release);
1846 if let Err(error) = conn.authorizer(Some(deny_retired_writer)) {
1847 tracing::error!(
1848 %error,
1849 "failed to install the retired pooled-writer quarantine authorizer"
1850 );
1851 }
1852 }
1853
1854 fn ensure_pooled_writer_active(&self) -> Result<(), SqliteError> {
1855 if self.pooled_writer_retired.load(Ordering::Acquire) {
1856 return Err(SqliteError::InvalidData(
1857 "pooled writer connection retired after a terminal transaction fault".to_string(),
1858 ));
1859 }
1860 Ok(())
1861 }
1862
1863 pub fn writer_acquisition_snapshot(&self) -> WriterAcquisitionSnapshot {
1866 self.writer_acquisition_counters.snapshot()
1867 }
1868
1869 pub fn reader_acquisition_snapshot(&self) -> ReaderAcquisitionSnapshot {
1873 self.reader_acquisition_counters.snapshot(
1874 self.max_readers.max(1),
1875 self.sql_bridge_reader_slots.available_permits(),
1876 )
1877 }
1878
1879 pub(crate) fn record_reader_admission_timeout(&self) {
1883 self.reader_acquisition_counters.record_checkout_timeout();
1884 }
1885
1886 pub(crate) fn writer_acquisition_counters(&self) -> Arc<WriterAcquisitionCounters> {
1888 Arc::clone(&self.writer_acquisition_counters)
1889 }
1890
1891 pub fn available_readers(&self) -> usize {
1893 self.readers.len()
1894 }
1895
1896 pub fn max_readers(&self) -> usize {
1898 self.max_readers
1899 }
1900
1901 pub fn config(&self) -> &PoolConfig {
1903 &self.config
1904 }
1905
1906 pub fn main_pool_generation(&self) -> u64 {
1910 *self.main_pool_generation.get_or_init(|| {
1911 NEXT_MAIN_POOL_GENERATION
1912 .fetch_update(Ordering::Relaxed, Ordering::Relaxed, |next| {
1913 next.checked_add(1)
1914 })
1915 .expect("main pool generation exhausted")
1916 })
1917 }
1918
1919 pub(crate) fn reader_admission_timeout(&self, operation: &'static str) -> StorageError {
1923 StorageError::AdmissionTimeout {
1924 operation: operation.into(),
1925 timeout_ms: u64::try_from(self.config.checkout_timeout.as_millis()).unwrap_or(u64::MAX),
1926 pool_identity: Some(
1927 self.identity_registration
1928 .as_ref()
1929 .map(PoolIdentityRegistration::label)
1930 .unwrap_or_else(|| ":memory:".to_string()),
1931 ),
1932 }
1933 }
1934
1935 pub(crate) fn resolve_reader_checkout<'p>(
1954 &self,
1955 capability: StorageCapability,
1956 operation: &'static str,
1957 outcome: Result<Option<ReaderGuard<'p>>, SqliteError>,
1958 ) -> Result<ReaderGuard<'p>, StorageError> {
1959 match outcome {
1960 Ok(Some(mut guard)) => {
1961 guard.label_operation(operation);
1962 Ok(guard)
1963 }
1964 Ok(None) => Err(StorageError::Timeout {
1965 operation: operation.into(),
1966 }),
1967 Err(error) => {
1968 let is_pool_exhausted = matches!(
1969 &error,
1970 SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
1971 if code.code == rusqlite::ErrorCode::DatabaseBusy
1972 );
1973 if is_pool_exhausted {
1974 Err(self.reader_admission_timeout(operation))
1975 } else {
1976 Err(StorageError::driver(capability, operation, error))
1977 }
1978 }
1979 }
1980 }
1981
1982 pub(crate) fn sql_bridge_reader_slots(&self) -> Arc<Semaphore> {
1985 Arc::clone(&self.sql_bridge_reader_slots)
1986 }
1987
1988 pub(crate) fn sql_bridge_writer_slots(&self) -> Arc<Semaphore> {
1990 Arc::clone(&self.sql_bridge_writer_slots)
1991 }
1992
1993 pub fn origin(&self) -> TxOrigin {
1999 self.origin.clone()
2000 }
2001
2002 pub fn canonical_path(&self) -> Option<&Path> {
2008 self.identity_path.as_deref()
2009 }
2010
2011 pub fn write_queue_active(&self) -> bool {
2023 debug_assert!(
2024 self.config.write_queue_enabled.is_some(),
2025 "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
2026 before any write_queue_active read"
2027 );
2028 self.config.write_queue_enabled.unwrap_or(false) && self.config.path.is_some()
2029 }
2030
2031 pub fn writer_task_join_was_stored(&self) -> bool {
2037 self.writer_task_join_stored.load(Ordering::SeqCst)
2038 }
2039
2040 pub fn writer_task_handle(&self) -> Result<Option<WriterTaskHandle>, StorageError> {
2064 debug_assert!(
2072 self.config.write_queue_enabled.is_some(),
2073 "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
2074 before any writer_task_handle read"
2075 );
2076 if !self.config.write_queue_enabled.unwrap_or(false) {
2077 return Ok(None);
2078 }
2079 if let Some(existing) = self.writer_task.get() {
2082 return Ok(existing.clone());
2083 }
2084 if tokio::runtime::Handle::try_current().is_err() {
2088 return Err(StorageError::WriterTaskNoRuntime);
2089 }
2090 Ok(self
2091 .writer_task
2092 .get_or_init(|| {
2093 #[cfg(test)]
2094 self.writer_task_spawn_count
2095 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
2096
2097 match crate::writer_task::spawn(self, self.config.write_queue_capacity) {
2098 Ok(handle) => Some(handle),
2099 Err(e) => {
2100 tracing::warn!(
2101 error = %e,
2102 "KHIVE_WRITE_QUEUE=1 but the writer task failed to spawn; \
2103 writes fall back to the pool-mutex path"
2104 );
2105 None
2106 }
2107 }
2108 })
2109 .clone())
2110 }
2111
2112 pub(crate) fn writer_task_for_write(
2122 &self,
2123 cached: Option<&WriterTaskHandle>,
2124 operation: &'static str,
2125 ) -> Result<Option<WriterTaskHandle>, StorageError> {
2126 let handle = match cached {
2127 Some(handle) => Some(handle.clone()),
2128 None => match self.writer_task_handle() {
2129 Ok(handle) => handle,
2130 Err(error) if self.config.write_routing_strict => return Err(error),
2131 Err(_) => None,
2132 },
2133 };
2134
2135 if handle.is_none() && self.config.write_routing_strict {
2136 return Err(StorageError::Pool {
2137 operation: operation.into(),
2138 message: "strict write routing requires a writer-task handle; no handle is \
2139 available, so the direct writer fallback was refused"
2140 .into(),
2141 });
2142 }
2143 Ok(handle)
2144 }
2145
2146 pub(crate) fn record_direct_route(&self, site: crate::timeout_sink::Site) {
2150 if self.write_queue_active() {
2151 crate::timeout_sink::emit_direct_route_violation(
2152 &crate::timeout_sink::db_label(self),
2153 site,
2154 );
2155 }
2156 }
2157
2158 pub fn writer_task_for_runtime_write(
2163 &self,
2164 operation: RuntimeWriteOperation,
2165 ) -> Result<Option<WriterTaskHandle>, StorageError> {
2166 let handle = self.writer_task_for_write(None, operation.operation())?;
2167 if handle.is_none() {
2168 self.record_direct_route(operation.fallback_site());
2169 }
2170 Ok(handle)
2171 }
2172
2173 #[cfg(test)]
2178 pub(crate) fn writer_task_spawn_count(&self) -> usize {
2179 self.writer_task_spawn_count
2180 .load(std::sync::atomic::Ordering::SeqCst)
2181 }
2182
2183 pub(crate) fn set_writer_task_join(&self, join: tokio::task::JoinHandle<()>) {
2196 let first_store = !self.writer_task_join_stored.swap(true, Ordering::SeqCst);
2201 debug_assert!(
2202 first_store,
2203 "writer task JoinHandle stored twice (even counting a taken one); \
2204 the writer_task OnceLock is supposed to make spawn at-most-once per pool"
2205 );
2206 if first_store {
2207 *self.writer_task_join.lock() = Some(join);
2208 }
2209 }
2210
2211 pub fn take_writer_task_join(&self) -> Option<tokio::task::JoinHandle<()>> {
2229 self.writer_task_join.lock().take()
2230 }
2231
2232 pub fn legacy_conn(&self) -> Arc<Mutex<Connection>> {
2237 Arc::clone(&self.writer)
2238 }
2239
2240 fn open_reader_connection(&self) -> Result<Connection, SqliteError> {
2241 let path = self.read_connection_path()?;
2242 open_reader_connection(path, &self.config)
2243 }
2244
2245 fn read_connection_path(&self) -> Result<&Path, SqliteError> {
2246 self.read_only_open_target
2247 .as_deref()
2248 .or(self.config.path.as_deref())
2249 .ok_or_else(|| {
2250 SqliteError::InvalidData(
2251 "in-memory databases do not support standalone connections".to_string(),
2252 )
2253 })
2254 }
2255
2256 pub fn open_standalone_writer(&self) -> Result<Connection, SqliteError> {
2267 let conn = self.open_standalone_writer_untracked()?;
2268 self.writer_acquisition_counters
2269 .standalone_acquisitions
2270 .fetch_add(1, Ordering::Relaxed);
2271 Ok(conn)
2272 }
2273
2274 pub(crate) fn open_standalone_writer_untracked(&self) -> Result<Connection, SqliteError> {
2284 let path = self.config.path.as_ref().ok_or_else(|| {
2285 SqliteError::InvalidData(
2286 "in-memory databases do not support standalone connections".to_string(),
2287 )
2288 })?;
2289
2290 if self.config.read_only {
2291 return Err(SqliteError::InvalidData(
2292 "database is read-only: standalone write connections are not permitted".to_string(),
2293 ));
2294 }
2295
2296 let conn = Connection::open_with_flags(
2297 path,
2298 OpenFlags::SQLITE_OPEN_READ_WRITE
2299 | OpenFlags::SQLITE_OPEN_NO_MUTEX
2300 | OpenFlags::SQLITE_OPEN_URI,
2301 )?;
2302 register_writer_clock(&conn)?;
2303 conn.busy_timeout(self.config.busy_timeout)?;
2304 self.checkpoint_ownership
2305 .configure_wal_autocheckpoint(&conn)?;
2306 conn.pragma_update(None, "foreign_keys", "ON")?;
2307 conn.pragma_update(None, "synchronous", "NORMAL")?;
2308
2309 let wal_enabled =
2310 self.config.wal_mode && current_journal_mode(&conn)?.eq_ignore_ascii_case("wal");
2311 if wal_enabled {
2312 conn.pragma_update(
2313 None,
2314 "journal_size_limit",
2315 self.config.journal_size_limit_bytes,
2316 )?;
2317 }
2318
2319 Ok(conn)
2320 }
2321
2322 #[cfg(test)]
2326 pub(crate) fn effective_wal_autocheckpoint_pages(&self) -> u32 {
2327 self.checkpoint_ownership.wal_autocheckpoint_pages()
2328 }
2329
2330 pub fn claim_checkpoint_ownership(&self) -> Result<(), SqliteError> {
2352 if !self.checkpoint_ownership.begin_claim() {
2353 return Ok(());
2354 }
2355 let result = (|| {
2356 if !self.config.read_only {
2357 let writer = self.writer()?;
2358 writer.conn().pragma_update(None, "wal_autocheckpoint", 0)?;
2359 }
2360 Ok(())
2361 })();
2362 self.checkpoint_ownership.finish_claim(result.is_ok());
2363 result
2364 }
2365
2366 pub async fn propagate_checkpoint_claim_to_writer_task(&self) -> Result<(), StorageError> {
2375 let Some(handle) = self.writer_task_handle()? else {
2376 return Ok(());
2377 };
2378 handle
2379 .send_top_level(|conn| {
2380 conn.pragma_update(None, "wal_autocheckpoint", 0)
2381 .map_err(|e| StorageError::Pool {
2382 operation: "claim_checkpoint_ownership".into(),
2383 message: e.to_string(),
2384 })
2385 })
2386 .await
2387 }
2388
2389 pub(crate) fn open_standalone_reader(
2398 &self,
2399 purpose: StandaloneReaderPurpose,
2400 ) -> Result<Connection, SqliteError> {
2401 let path = self.read_connection_path()?;
2402
2403 let conn = Connection::open_with_flags(
2404 path,
2405 OpenFlags::SQLITE_OPEN_READ_ONLY
2406 | OpenFlags::SQLITE_OPEN_NO_MUTEX
2407 | OpenFlags::SQLITE_OPEN_URI,
2408 )?;
2409 configure_reader_connection(&conn, &self.config)?;
2410 conn.pragma_update(None, "synchronous", "NORMAL")?;
2411 self.reader_acquisition_counters
2412 .record_standalone_open(purpose);
2413 Ok(conn)
2414 }
2415
2416 fn return_reader(&self, conn: Connection, dirty: bool) {
2417 if self.max_readers == 0 {
2418 return;
2419 }
2420
2421 if reset_reader_connection(&conn, dirty, &self.config)
2422 && reader_connection_is_healthy(&conn)
2423 {
2424 self.enqueue_reader_slot(conn);
2425 return;
2426 }
2427
2428 close_connection_quietly(conn);
2429 self.replace_discarded_reader_slot();
2430 }
2431
2432 fn enqueue_reader_slot(&self, conn: Connection) {
2436 if let Err(conn) = self.readers.push(conn) {
2437 eprintln!("[sqlite-pool] reader pool queue full, discarding replacement connection");
2438 close_connection_quietly(conn);
2439 }
2440 }
2441
2442 fn replace_discarded_reader_slot(&self) {
2449 match self.open_reader_connection() {
2450 Ok(conn) => self.enqueue_reader_slot(conn),
2451 Err(error) => {
2452 self.reader_acquisition_counters
2453 .record_reader_replacement_open_failure();
2454 tracing::warn!(
2455 %error,
2456 "sqlite-pool: reader replacement connection failed to open; the physical \
2457 pool permanently shrinks by one slot below max_readers"
2458 );
2459 }
2460 }
2461 }
2462}
2463
2464const MAX_SYMLINK_DEPTH: u32 = 40;
2469
2470fn mint_db_identity(configured_path: &Path) -> Result<(DbIdentity, PathBuf), SqliteError> {
2503 let absolute = if configured_path.is_absolute() {
2504 configured_path.to_path_buf()
2505 } else {
2506 let cwd = std::env::current_dir().map_err(|e| {
2507 SqliteError::InvalidData(format!(
2508 "cannot mint database identity for {configured_path:?}: failed to resolve the \
2509 process current directory: {e}"
2510 ))
2511 })?;
2512 cwd.join(configured_path)
2513 };
2514
2515 if absolute.exists() {
2516 let canonical = absolute.canonicalize().map_err(|e| {
2517 SqliteError::InvalidData(format!(
2518 "cannot mint database identity: failed to canonicalize existing path \
2519 {absolute:?}: {e}"
2520 ))
2521 })?;
2522 return Ok((
2523 DbIdentity::new(canonical.clone().into_os_string()),
2524 canonical,
2525 ));
2526 }
2527
2528 let resolved_target = resolve_symlink_chain(&absolute)?;
2529 let parent = resolved_target.parent().ok_or_else(|| {
2530 SqliteError::InvalidData(format!(
2531 "cannot mint database identity for {resolved_target:?}: path has no parent \
2532 directory"
2533 ))
2534 })?;
2535 let file_name = resolved_target.file_name().ok_or_else(|| {
2536 SqliteError::InvalidData(format!(
2537 "cannot mint database identity for {resolved_target:?}: path has no file name"
2538 ))
2539 })?;
2540 let canonical_parent = parent.canonicalize().map_err(|e| {
2541 SqliteError::InvalidData(format!(
2542 "cannot mint database identity: parent directory {parent:?} of first-open path \
2543 {resolved_target:?} does not exist or is inaccessible: {e}"
2544 ))
2545 })?;
2546 let mut identity_path = canonical_parent;
2547 identity_path.push(file_name);
2548 Ok((
2549 DbIdentity::new(identity_path.clone().into_os_string()),
2550 identity_path,
2551 ))
2552}
2553
2554fn resolve_symlink_chain(path: &Path) -> Result<PathBuf, SqliteError> {
2560 let mut current = path.to_path_buf();
2561 for _ in 0..MAX_SYMLINK_DEPTH {
2562 match fs::symlink_metadata(¤t) {
2563 Ok(meta) if meta.file_type().is_symlink() => {
2564 let target = fs::read_link(¤t).map_err(|e| {
2565 SqliteError::InvalidData(format!(
2566 "cannot mint database identity: failed to read symlink {current:?}: {e}"
2567 ))
2568 })?;
2569 current = if target.is_absolute() {
2570 target
2571 } else {
2572 match current.parent() {
2573 Some(parent) => parent.join(&target),
2574 None => target,
2575 }
2576 };
2577 }
2578 _ => return Ok(current),
2579 }
2580 }
2581 Err(SqliteError::InvalidData(format!(
2582 "cannot mint database identity for {path:?}: symlink chain exceeds \
2583 {MAX_SYMLINK_DEPTH} levels"
2584 )))
2585}
2586
2587fn effective_reader_count(config: &PoolConfig, wal_enabled: bool) -> usize {
2588 if config.path.is_some() && config.read_only {
2589 config.max_readers.max(1)
2590 } else if config.path.is_some() && config.wal_mode && wal_enabled {
2591 config.max_readers
2592 } else {
2593 0
2594 }
2595}
2596
2597fn open_writer_connection(
2598 config: &PoolConfig,
2599 read_only_open_target: Option<&Path>,
2600) -> Result<Connection, SqliteError> {
2601 match config.path.as_ref() {
2602 Some(path) => {
2603 let flags = if config.read_only {
2604 writer_read_only_open_flags()
2605 } else {
2606 writer_open_flags()
2607 };
2608 let target = if config.read_only {
2609 read_only_open_target.ok_or_else(|| {
2610 SqliteError::InvalidData(
2611 "file-backed read-only pool has no canonical open target".to_string(),
2612 )
2613 })?
2614 } else {
2615 path
2616 };
2617 Connection::open_with_flags(target, flags).map_err(Into::into)
2618 }
2619 None => Connection::open_in_memory().map_err(Into::into),
2620 }
2621}
2622
2623fn read_only_open_target(
2646 config: &PoolConfig,
2647 physical_path: Option<&Path>,
2648) -> Result<Option<PathBuf>, SqliteError> {
2649 if !config.read_only {
2650 return Ok(None);
2651 }
2652 let Some(path) = physical_path else {
2653 return Ok(None);
2654 };
2655 read_only_wal_open_target_for_path(path).map(Some)
2656}
2657
2658fn read_only_wal_open_target_for_path(path: &Path) -> Result<PathBuf, SqliteError> {
2659 if !sqlite_header_uses_wal(path)? {
2660 return Ok(path.to_path_buf());
2661 }
2662
2663 let shm = sqlite_sidecar_path(path, "-shm");
2664 match fs::metadata(&shm) {
2665 Ok(metadata) if metadata.permissions().readonly() => {
2666 let wal = sqlite_sidecar_path(path, "-wal");
2667 match fs::metadata(&wal) {
2668 Ok(_) => Ok(path.to_path_buf()),
2669 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
2670 Err(SqliteError::InvalidData(format!(
2671 "read-only WAL snapshot {} has a shared-memory sidecar {} but no WAL \
2672 sidecar {}; refusing the inconsistent sidecar set before SQLite open",
2673 path.display(),
2674 shm.display(),
2675 wal.display(),
2676 )))
2677 }
2678 Err(error) => Err(SqliteError::Io(error)),
2679 }
2680 }
2681 Ok(_) => Err(SqliteError::InvalidData(format!(
2682 "read-only WAL snapshot {} has a writable WAL shared-memory sidecar {}; close every \
2683 live writer and remove the transient -shm file (or make a genuinely frozen snapshot) \
2684 before inspection",
2685 path.display(),
2686 shm.display(),
2687 ))),
2688 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
2689 let wal = sqlite_sidecar_path(path, "-wal");
2690 match fs::metadata(&wal) {
2691 Ok(metadata) if metadata.len() > 0 => Err(SqliteError::InvalidData(format!(
2692 "read-only WAL snapshot {} has a non-empty WAL sidecar {} but no read-only \
2693 shared-memory sidecar {}; refusing before SQLite open because immutable \
2694 mode would omit committed WAL frames and ordinary read-only mode would \
2695 create or mutate -shm; include the frozen read-only -shm beside this \
2696 snapshot, or checkpoint a writable copy before inspection",
2697 path.display(),
2698 wal.display(),
2699 shm.display(),
2700 ))),
2701 Ok(_) => sqlite_immutable_uri(path),
2702 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
2703 sqlite_immutable_uri(path)
2704 }
2705 Err(error) => Err(SqliteError::Io(error)),
2706 }
2707 }
2708 Err(error) => Err(SqliteError::Io(error)),
2709 }
2710}
2711
2712pub(crate) fn open_read_only_snapshot_connection(path: &Path) -> Result<Connection, SqliteError> {
2713 let (_, physical_path) = mint_db_identity(path)?;
2714 let target = read_only_wal_open_target_for_path(&physical_path)?;
2715 Connection::open_with_flags(&target, reader_open_flags()).map_err(Into::into)
2716}
2717
2718fn sqlite_header_uses_wal(path: &Path) -> Result<bool, SqliteError> {
2719 let mut file = fs::File::open(path)?;
2720 let mut header = [0_u8; 20];
2721 if let Err(error) = file.read_exact(&mut header) {
2722 if error.kind() == std::io::ErrorKind::UnexpectedEof {
2723 return Ok(false);
2724 }
2725 return Err(SqliteError::Io(error));
2726 }
2727 Ok(&header[..16] == b"SQLite format 3\0" && header[18] == 2 && header[19] == 2)
2728}
2729
2730fn sqlite_sidecar_path(path: &Path, suffix: &str) -> PathBuf {
2731 let mut sidecar = path.as_os_str().to_os_string();
2732 sidecar.push(suffix);
2733 PathBuf::from(sidecar)
2734}
2735
2736fn sqlite_immutable_uri(path: &Path) -> Result<PathBuf, SqliteError> {
2737 let absolute = if path.is_absolute() {
2738 path.to_path_buf()
2739 } else {
2740 std::env::current_dir()?.join(path)
2741 };
2742 let mut uri = String::from("file:");
2743
2744 #[cfg(unix)]
2745 {
2746 use std::os::unix::ffi::OsStrExt as _;
2747 push_sqlite_uri_path(&mut uri, absolute.as_os_str().as_bytes());
2748 }
2749
2750 #[cfg(not(unix))]
2751 {
2752 let path = absolute.to_str().ok_or_else(|| {
2753 SqliteError::InvalidData(format!(
2754 "read-only WAL snapshot path is not representable as a SQLite URI: {}",
2755 absolute.display()
2756 ))
2757 })?;
2758 let normalized = path.replace('\\', "/");
2759 if cfg!(windows) && !normalized.starts_with('/') {
2760 uri.push('/');
2761 }
2762 push_sqlite_uri_path(&mut uri, normalized.as_bytes());
2763 }
2764
2765 uri.push_str("?mode=ro&immutable=1");
2766 Ok(PathBuf::from(uri))
2767}
2768
2769fn push_sqlite_uri_path(uri: &mut String, bytes: &[u8]) {
2770 const HEX: &[u8; 16] = b"0123456789ABCDEF";
2771 for &byte in bytes {
2772 if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~' | b'/') {
2773 uri.push(byte as char);
2774 } else {
2775 uri.push('%');
2776 uri.push(HEX[(byte >> 4) as usize] as char);
2777 uri.push(HEX[(byte & 0x0f) as usize] as char);
2778 }
2779 }
2780}
2781
2782fn open_reader_connection(path: &Path, config: &PoolConfig) -> Result<Connection, SqliteError> {
2783 let conn = Connection::open_with_flags(path, reader_open_flags())?;
2784 configure_reader_connection(&conn, config)?;
2785 Ok(conn)
2786}
2787
2788fn writer_open_flags() -> OpenFlags {
2789 OpenFlags::SQLITE_OPEN_READ_WRITE
2790 | OpenFlags::SQLITE_OPEN_CREATE
2791 | OpenFlags::SQLITE_OPEN_URI
2792 | OpenFlags::SQLITE_OPEN_NO_MUTEX
2793}
2794
2795fn writer_read_only_open_flags() -> OpenFlags {
2798 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
2799}
2800
2801fn reader_open_flags() -> OpenFlags {
2802 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
2803}
2804
2805fn register_writer_clock(conn: &Connection) -> Result<(), SqliteError> {
2806 conn.create_scalar_function(
2809 "khive_now_micros",
2810 0,
2811 rusqlite::functions::FunctionFlags::SQLITE_UTF8,
2812 |_| Ok(chrono::Utc::now().timestamp_micros()),
2813 )?;
2814 Ok(())
2815}
2816
2817pub(crate) fn rfc3339_instant_key(instant: chrono::DateTime<chrono::Utc>) -> Vec<u8> {
2820 let mut key = Vec::with_capacity(12);
2821 key.extend_from_slice(&((instant.timestamp() as u64) ^ (1_u64 << 63)).to_be_bytes());
2822 key.extend_from_slice(&instant.timestamp_subsec_nanos().to_be_bytes());
2823 key
2824}
2825
2826fn register_rfc3339_key(conn: &Connection) -> Result<(), SqliteError> {
2827 use rusqlite::functions::FunctionFlags;
2828 use rusqlite::types::ValueRef;
2829
2830 conn.create_scalar_function(
2831 "khive_rfc3339_key",
2832 1,
2833 FunctionFlags::SQLITE_UTF8
2834 | FunctionFlags::SQLITE_DETERMINISTIC
2835 | FunctionFlags::SQLITE_INNOCUOUS,
2836 |ctx| {
2837 let text = match ctx.get_raw(0) {
2838 ValueRef::Text(bytes) => std::str::from_utf8(bytes).ok(),
2839 _ => None,
2840 };
2841 let key = text
2842 .and_then(|text| text.parse::<chrono::DateTime<chrono::Utc>>().ok())
2843 .map(rfc3339_instant_key);
2844 Ok(key)
2845 },
2846 )?;
2847 Ok(())
2848}
2849
2850fn configure_writer_connection(
2851 conn: &Connection,
2852 config: &PoolConfig,
2853) -> Result<bool, SqliteError> {
2854 register_writer_clock(conn)?;
2855 register_rfc3339_key(conn)?;
2856 if config.read_only {
2857 conn.pragma_update(None, "foreign_keys", "ON")?;
2861 conn.busy_timeout(config.busy_timeout)?;
2862 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
2863 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
2864 conn.pragma_update(None, "temp_store", "MEMORY")?;
2865 conn.pragma_update(None, "query_only", "ON")?;
2866
2867 let wal_enabled =
2868 config.wal_mode && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
2869 return Ok(wal_enabled);
2870 }
2871
2872 let wants_wal = config.path.is_some() && config.wal_mode;
2873
2874 if wants_wal {
2875 conn.pragma_update(None, "journal_mode", "WAL")?;
2876 }
2877
2878 conn.pragma_update(None, "synchronous", "NORMAL")?;
2879 conn.pragma_update(None, "foreign_keys", "ON")?;
2880 conn.busy_timeout(config.busy_timeout)?;
2881 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
2882 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
2883 conn.pragma_update(None, "temp_store", "MEMORY")?;
2884 conn.pragma_update(
2889 None,
2890 "wal_autocheckpoint",
2891 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
2892 )?;
2893
2894 let wal_enabled = wants_wal && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
2895
2896 if wal_enabled {
2897 conn.pragma_update(None, "journal_size_limit", config.journal_size_limit_bytes)?;
2898 }
2899
2900 Ok(wal_enabled)
2901}
2902
2903fn configure_reader_connection(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
2904 register_rfc3339_key(conn)?;
2905 conn.pragma_update(None, "foreign_keys", "ON")?;
2906 conn.busy_timeout(config.busy_timeout)?;
2907 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
2908 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
2909 conn.pragma_update(None, "temp_store", "MEMORY")?;
2910 Ok(())
2911}
2912
2913fn current_journal_mode(conn: &Connection) -> Result<String, SqliteError> {
2914 conn.pragma_query_value(None, "journal_mode", |row| row.get::<_, String>(0))
2915 .map(|mode| mode.to_ascii_lowercase())
2916 .map_err(Into::into)
2917}
2918
2919fn reset_reader_connection(conn: &Connection, dirty: bool, config: &PoolConfig) -> bool {
2920 if !conn.is_autocommit() {
2921 match conn.execute_batch("ROLLBACK") {
2922 Ok(()) => {}
2923 Err(rusqlite::Error::SqliteFailure(err, _)) => {
2924 if matches!(
2925 err.code,
2926 rusqlite::ErrorCode::CannotOpen
2927 | rusqlite::ErrorCode::DatabaseCorrupt
2928 | rusqlite::ErrorCode::NotADatabase
2929 | rusqlite::ErrorCode::DiskFull
2930 ) {
2931 return false;
2932 }
2933 }
2934 Err(_) => return false,
2935 }
2936 if !conn.is_autocommit() {
2937 return false;
2938 }
2939 }
2940
2941 if !dirty {
2942 return true;
2943 }
2944
2945 reader_connection_state_is_pristine(conn)
2946 && reader_connection_settings_match_baseline(conn, config, 0)
2947}
2948
2949fn reader_connection_state_is_pristine(conn: &Connection) -> bool {
2964 let has_temp_objects: bool = match conn.query_row(
2965 "SELECT EXISTS(SELECT 1 FROM sqlite_temp_master)",
2966 [],
2967 |row| row.get(0),
2968 ) {
2969 Ok(v) => v,
2970 Err(_) => return false,
2971 };
2972 if has_temp_objects {
2973 return false;
2974 }
2975
2976 let attached_databases: i64 = match conn.query_row(
2977 "SELECT COUNT(*) FROM pragma_database_list WHERE name NOT IN ('main', 'temp')",
2978 [],
2979 |row| row.get(0),
2980 ) {
2981 Ok(v) => v,
2982 Err(_) => return false,
2983 };
2984 attached_databases == 0
2985}
2986
2987fn reader_connection_settings_match_baseline(
3001 conn: &Connection,
3002 config: &PoolConfig,
3003 expected_query_only: i64,
3004) -> bool {
3005 let expected_busy_timeout_ms =
3006 i64::try_from(config.busy_timeout.as_millis()).unwrap_or(i64::MAX);
3007 let expected_cache_size: i64 = CACHE_SIZE_KIB.parse().unwrap_or(-65536);
3008 let checks: [(&str, i64); 8] = [
3009 ("query_only", expected_query_only),
3010 ("writable_schema", 0),
3011 ("foreign_keys", 1),
3012 ("busy_timeout", expected_busy_timeout_ms),
3013 ("cache_size", expected_cache_size),
3014 ("temp_store", 2),
3015 ("read_uncommitted", 0),
3016 ("defer_foreign_keys", 0),
3017 ];
3018 checks.iter().all(|(pragma, expected)| {
3019 conn.pragma_query_value(None, pragma, |row| row.get::<_, i64>(0))
3020 .map(|actual| actual == *expected)
3021 .unwrap_or(false)
3022 })
3023}
3024
3025fn restore_shared_reader_state(conn: &Connection, config: &PoolConfig) -> bool {
3034 if !detach_non_main_databases(conn) {
3035 return false;
3036 }
3037 if !drop_temp_objects(conn) {
3038 return false;
3039 }
3040 let expected_query_only = i64::from(config.read_only);
3041 if reset_observable_settings(conn, config, expected_query_only).is_err() {
3042 return false;
3043 }
3044 reader_connection_state_is_pristine(conn)
3045 && reader_connection_settings_match_baseline(conn, config, expected_query_only)
3046}
3047
3048fn detach_non_main_databases(conn: &Connection) -> bool {
3049 loop {
3050 let name: Option<String> = match conn.query_row(
3051 "SELECT name FROM pragma_database_list WHERE name NOT IN ('main', 'temp') LIMIT 1",
3052 [],
3053 |row| row.get(0),
3054 ) {
3055 Ok(name) => Some(name),
3056 Err(rusqlite::Error::QueryReturnedNoRows) => None,
3057 Err(_) => return false,
3058 };
3059 let Some(name) = name else {
3060 return true;
3061 };
3062 let quoted = format!("\"{}\"", name.replace('"', "\"\""));
3063 if conn
3064 .execute_batch(&format!("DETACH DATABASE {quoted}"))
3065 .is_err()
3066 {
3067 return false;
3068 }
3069 }
3070}
3071
3072fn drop_temp_objects(conn: &Connection) -> bool {
3073 for (kind, ddl_keyword) in [
3077 ("view", "VIEW"),
3078 ("trigger", "TRIGGER"),
3079 ("index", "INDEX"),
3080 ("table", "TABLE"),
3081 ] {
3082 loop {
3083 let name: Option<String> = match conn.query_row(
3084 "SELECT name FROM sqlite_temp_master WHERE type = ?1 LIMIT 1",
3085 [kind],
3086 |row| row.get(0),
3087 ) {
3088 Ok(name) => Some(name),
3089 Err(rusqlite::Error::QueryReturnedNoRows) => None,
3090 Err(_) => return false,
3091 };
3092 let Some(name) = name else {
3093 break;
3094 };
3095 let quoted = format!("\"{}\"", name.replace('"', "\"\""));
3096 if conn
3097 .execute_batch(&format!("DROP {ddl_keyword} IF EXISTS temp.{quoted}"))
3098 .is_err()
3099 {
3100 return false;
3101 }
3102 }
3103 }
3104 true
3105}
3106
3107fn reset_observable_settings(
3108 conn: &Connection,
3109 config: &PoolConfig,
3110 expected_query_only: i64,
3111) -> Result<(), rusqlite::Error> {
3112 conn.pragma_update(None, "query_only", expected_query_only)?;
3113 conn.pragma_update(None, "writable_schema", 0)?;
3114 conn.pragma_update(None, "foreign_keys", "ON")?;
3115 conn.busy_timeout(config.busy_timeout)?;
3116 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
3117 conn.pragma_update(None, "temp_store", "MEMORY")?;
3118 conn.pragma_update(None, "read_uncommitted", 0)?;
3119 conn.pragma_update(None, "defer_foreign_keys", 0)?;
3120 Ok(())
3121}
3122
3123fn reader_connection_is_healthy(conn: &Connection) -> bool {
3124 match conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0)) {
3125 Ok(_) => true,
3126 Err(rusqlite::Error::SqliteFailure(err, _)) => !matches!(
3127 err.code,
3128 rusqlite::ErrorCode::CannotOpen
3129 | rusqlite::ErrorCode::NotADatabase
3130 | rusqlite::ErrorCode::DatabaseCorrupt
3131 | rusqlite::ErrorCode::PermissionDenied
3132 | rusqlite::ErrorCode::SystemIoFailure
3133 ),
3134 Err(_) => true,
3135 }
3136}
3137
3138fn close_connection_quietly(conn: Connection) {
3139 match conn.close() {
3140 Ok(()) => {}
3141 Err((conn, _)) => drop(conn),
3142 }
3143}
3144
3145fn pool_exhausted_error(timeout: Duration, max_readers: usize) -> SqliteError {
3146 rusqlite::Error::SqliteFailure(
3147 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
3148 Some(format!(
3149 "Pool exhausted: no reader available after {timeout:?} (max_readers={max_readers})"
3150 )),
3151 )
3152 .into()
3153}
3154
3155#[cfg(test)]
3156#[path = "runtime_write_routing_tests.rs"]
3157mod runtime_write_routing_tests;
3158
3159#[cfg(test)]
3160mod tests {
3161 use super::*;
3162 use serial_test::serial;
3163
3164 #[test]
3165 fn constructor_writer_cancels_after_entering_the_wait_without_pool_timeout() {
3166 let pool = ConnectionPool::new(PoolConfig {
3167 path: None,
3168 ..PoolConfig::default()
3169 })
3170 .unwrap();
3171 let held = pool.writer().unwrap();
3172 let before = pool.writer_acquisition_snapshot();
3173 let checks = Cell::new(0);
3174 let stopped = pool
3175 .writer_until(|| {
3176 checks.set(checks.get() + 1);
3177 checks.get() == 2
3178 })
3179 .unwrap();
3180 assert!(
3181 stopped.is_none(),
3182 "second predicate check must stop an in-flight wait"
3183 );
3184 assert_eq!(checks.get(), 2);
3185 assert_eq!(pool.writer_acquisition_snapshot(), before);
3186 drop(held);
3187 }
3188
3189 #[tokio::test]
3190 async fn constructor_writer_observes_absolute_blocking_deadline() {
3191 let pool = ConnectionPool::new(PoolConfig {
3192 path: None,
3193 checkout_timeout: Duration::from_secs(5),
3194 ..PoolConfig::default()
3195 })
3196 .unwrap();
3197 let held = pool.writer().unwrap();
3198 let context =
3199 khive_storage::scope_request_read_deadline(Duration::from_millis(20), async {
3200 khive_storage::capture_request_read_context()
3201 })
3202 .await;
3203 let before = pool.writer_acquisition_snapshot();
3204 let started = Instant::now();
3205 let stopped = pool
3206 .writer_until(|| context.blocking_stop_reason().is_some())
3207 .unwrap();
3208 assert!(stopped.is_none());
3209 assert!(
3210 started.elapsed() < Duration::from_secs(1),
3211 "request deadline must beat pool timeout"
3212 );
3213 assert_eq!(pool.writer_acquisition_snapshot(), before);
3214 drop(held);
3215 }
3216
3217 #[test]
3218 fn constructor_writer_preserves_uncancelled_checkout_timeout() {
3219 let pool = ConnectionPool::new(PoolConfig {
3220 path: None,
3221 checkout_timeout: Duration::from_millis(5),
3222 ..PoolConfig::default()
3223 })
3224 .unwrap();
3225 let held = pool.writer().unwrap();
3226 let before = pool.writer_acquisition_snapshot();
3227 let result = pool.writer_until(|| false);
3228 assert!(
3229 matches!(result, Err(SqliteError::WriterPoolCheckoutTimeout { timeout }) if timeout == Duration::from_millis(5))
3230 );
3231 let after = pool.writer_acquisition_snapshot();
3232 assert_eq!(after.timeouts, before.timeouts + 1);
3233 assert_eq!(after.pooled_acquisitions, before.pooled_acquisitions);
3234 drop(held);
3235 }
3236
3237 struct WarningCapture {
3238 messages: Arc<std::sync::Mutex<Vec<String>>>,
3239 }
3240
3241 impl tracing::Subscriber for WarningCapture {
3242 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
3243 true
3244 }
3245
3246 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
3247 tracing::span::Id::from_u64(1)
3248 }
3249
3250 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
3251
3252 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
3253
3254 fn event(&self, event: &tracing::Event<'_>) {
3255 struct Visitor(Option<String>);
3256
3257 impl tracing::field::Visit for Visitor {
3258 fn record_debug(
3259 &mut self,
3260 field: &tracing::field::Field,
3261 value: &dyn std::fmt::Debug,
3262 ) {
3263 if field.name() == "message" {
3264 self.0 = Some(format!("{value:?}"));
3265 }
3266 }
3267 }
3268
3269 let mut visitor = Visitor(None);
3270 event.record(&mut visitor);
3271 if let Some(message) = visitor.0 {
3272 self.messages.lock().unwrap().push(message);
3273 }
3274 }
3275
3276 fn enter(&self, _: &tracing::span::Id) {}
3277
3278 fn exit(&self, _: &tracing::span::Id) {}
3279 }
3280
3281 struct CwdGuard {
3286 original: PathBuf,
3287 }
3288
3289 impl CwdGuard {
3290 fn enter(dir: &Path) -> Self {
3291 let original = std::env::current_dir().unwrap();
3292 std::env::set_current_dir(dir).unwrap();
3293 Self { original }
3294 }
3295 }
3296
3297 impl Drop for CwdGuard {
3298 fn drop(&mut self) {
3299 let _ = std::env::set_current_dir(&self.original);
3300 }
3301 }
3302
3303 const POOL_ENV_VARS: [&str; 7] = [
3304 "KHIVE_BUSY_TIMEOUT_SECS",
3305 "KHIVE_CHECKOUT_TIMEOUT_SECS",
3306 "KHIVE_WAL_AUTOCHECKPOINT_PAGES",
3307 "KHIVE_JOURNAL_SIZE_LIMIT_BYTES",
3308 "KHIVE_WRITE_QUEUE",
3309 "KHIVE_WRITE_QUEUE_CAPACITY",
3310 "KHIVE_WRITE_ROUTING",
3311 ];
3312
3313 struct PoolEnvGuard {
3314 saved: Vec<(&'static str, Option<std::ffi::OsString>)>,
3315 }
3316
3317 impl PoolEnvGuard {
3318 fn capture() -> Self {
3319 Self {
3320 saved: POOL_ENV_VARS
3321 .into_iter()
3322 .map(|key| (key, std::env::var_os(key)))
3323 .collect(),
3324 }
3325 }
3326 }
3327
3328 impl Drop for PoolEnvGuard {
3329 fn drop(&mut self) {
3330 for (key, value) in &self.saved {
3331 match value {
3332 Some(value) => std::env::set_var(key, value),
3333 None => std::env::remove_var(key),
3334 }
3335 }
3336 }
3337 }
3338
3339 fn clear_pool_env() -> PoolEnvGuard {
3340 let guard = PoolEnvGuard::capture();
3341 for var in POOL_ENV_VARS {
3342 std::env::remove_var(var);
3343 }
3344 guard
3345 }
3346
3347 fn wal_autocheckpoint_pages(conn: &Connection) -> u32 {
3348 conn.pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
3349 .expect("read PRAGMA wal_autocheckpoint")
3350 }
3351
3352 fn journal_size_limit_bytes(conn: &Connection) -> i64 {
3353 conn.pragma_query_value(None, "journal_size_limit", |row| row.get(0))
3354 .expect("read PRAGMA journal_size_limit")
3355 }
3356
3357 #[test]
3358 fn read_only_rollback_journal_pool_keeps_a_dedicated_reader() {
3359 let dir = tempfile::tempdir().unwrap();
3360 let path = dir.path().join("read_only_delete_journal.db");
3361 {
3362 let conn = Connection::open(&path).unwrap();
3363 conn.execute_batch("CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY);")
3364 .unwrap();
3365 let mode: String = conn
3366 .pragma_query_value(None, "journal_mode", |row| row.get(0))
3367 .unwrap();
3368 assert_eq!(mode.to_ascii_lowercase(), "delete");
3369 }
3370
3371 let pool = ConnectionPool::new(PoolConfig {
3372 path: Some(path),
3373 read_only: true,
3374 write_queue_enabled: Some(false),
3375 ..PoolConfig::for_test()
3376 })
3377 .unwrap();
3378
3379 assert!(
3380 pool.max_readers() > 0,
3381 "a read-only rollback-journal snapshot must use a genuine read-only reader, not \
3382 alias reader() onto the query-only writer slot"
3383 );
3384 let reader = pool.reader().expect("dedicated read-only reader checkout");
3385 let count: i64 = reader
3386 .conn()
3387 .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
3388 .unwrap();
3389 assert_eq!(count, 0);
3390 drop(reader);
3391 assert_eq!(
3392 pool.writer_acquisition_snapshot(),
3393 WriterAcquisitionSnapshot::default(),
3394 "constructing and reading a rollback-journal snapshot must never acquire the writer"
3395 );
3396 }
3397
3398 fn sqlite_sidecar(path: &Path, suffix: &str) -> PathBuf {
3399 let mut sidecar = path.as_os_str().to_os_string();
3400 sidecar.push(suffix);
3401 PathBuf::from(sidecar)
3402 }
3403
3404 fn directory_entries(path: &Path) -> Vec<std::ffi::OsString> {
3405 let mut entries = std::fs::read_dir(path)
3406 .unwrap()
3407 .map(|entry| entry.unwrap().file_name())
3408 .collect::<Vec<_>>();
3409 entries.sort();
3410 entries
3411 }
3412
3413 #[test]
3418 fn read_only_persistent_wal_without_shm_is_refused_without_mutation() {
3419 let dir = tempfile::tempdir().unwrap();
3420 let source = dir.path().join("wal-source.db");
3421 let snapshot = dir.path().join("snapshot ?#%.db");
3422 let source_wal = sqlite_sidecar(&source, "-wal");
3423 let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
3424 let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
3425
3426 let source_conn = Connection::open(&source).unwrap();
3427 let mode: String = source_conn
3428 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
3429 .unwrap();
3430 assert_eq!(mode.to_ascii_lowercase(), "wal");
3431 source_conn
3432 .pragma_update(None, "wal_autocheckpoint", 0)
3433 .unwrap();
3434 source_conn
3435 .execute_batch(
3436 "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
3437 INSERT INTO snapshot_row(body) VALUES ('committed-only-in-wal');",
3438 )
3439 .unwrap();
3440 assert!(source_wal.exists(), "fixture must retain a WAL sidecar");
3441
3442 std::fs::copy(&source, &snapshot).unwrap();
3443 std::fs::copy(&source_wal, &snapshot_wal).unwrap();
3444 assert!(
3445 !snapshot_shm.exists(),
3446 "fixture intentionally omits the transient shared-memory index"
3447 );
3448
3449 let main_before = std::fs::read(&snapshot).unwrap();
3450 let wal_before = std::fs::read(&snapshot_wal).unwrap();
3451 let entries_before = directory_entries(dir.path());
3452
3453 let error = match ConnectionPool::new(PoolConfig {
3454 path: Some(snapshot.clone()),
3455 read_only: true,
3456 write_queue_enabled: Some(false),
3457 ..PoolConfig::for_test()
3458 }) {
3459 Ok(_) => panic!("a non-empty WAL without its frozen -shm must fail closed"),
3460 Err(error) => error,
3461 };
3462 assert!(
3463 error
3464 .to_string()
3465 .contains("would omit committed WAL frames"),
3466 "diagnostic must explain why neither unsafe open mode is allowed: {error}"
3467 );
3468
3469 assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
3470 assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
3471 assert_eq!(directory_entries(dir.path()), entries_before);
3472 assert!(
3473 !snapshot_shm.exists(),
3474 "read-only admission and every reader must keep the source free of -shm"
3475 );
3476
3477 drop(source_conn);
3478 }
3479
3480 #[test]
3485 fn read_only_persistent_wal_with_read_only_shm_reads_without_mutation() {
3486 let dir = tempfile::tempdir().unwrap();
3487 let source = dir.path().join("wal-source.db");
3488 let snapshot = dir.path().join("frozen-wal-snapshot.db");
3489 let source_wal = sqlite_sidecar(&source, "-wal");
3490 let source_shm = sqlite_sidecar(&source, "-shm");
3491 let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
3492 let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
3493
3494 let source_conn = Connection::open(&source).unwrap();
3495 let mode: String = source_conn
3496 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
3497 .unwrap();
3498 assert_eq!(mode.to_ascii_lowercase(), "wal");
3499 source_conn
3500 .pragma_update(None, "wal_autocheckpoint", 0)
3501 .unwrap();
3502 source_conn
3503 .execute_batch(
3504 "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
3505 INSERT INTO snapshot_row(body) VALUES ('committed-only-in-wal');",
3506 )
3507 .unwrap();
3508 assert!(source_wal.exists() && source_shm.exists());
3509
3510 std::fs::copy(&source, &snapshot).unwrap();
3511 std::fs::copy(&source_wal, &snapshot_wal).unwrap();
3512 std::fs::copy(&source_shm, &snapshot_shm).unwrap();
3513
3514 let snapshot_paths = [&snapshot, &snapshot_wal, &snapshot_shm];
3515 let original_permissions =
3516 snapshot_paths.map(|path| std::fs::metadata(path).unwrap().permissions());
3517 for path in snapshot_paths {
3518 let mut permissions = std::fs::metadata(path).unwrap().permissions();
3519 permissions.set_readonly(true);
3520 std::fs::set_permissions(path, permissions).unwrap();
3521 }
3522
3523 let main_before = std::fs::read(&snapshot).unwrap();
3524 let wal_before = std::fs::read(&snapshot_wal).unwrap();
3525 let shm_before = std::fs::read(&snapshot_shm).unwrap();
3526 let entries_before = directory_entries(dir.path());
3527
3528 let pool = ConnectionPool::new(PoolConfig {
3529 path: Some(snapshot.clone()),
3530 read_only: true,
3531 write_queue_enabled: Some(false),
3532 ..PoolConfig::for_test()
3533 })
3534 .unwrap();
3535 let reader = pool.reader().unwrap();
3536 let body: String = reader
3537 .conn()
3538 .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
3539 .unwrap();
3540 assert_eq!(body, "committed-only-in-wal");
3541 drop(reader);
3542
3543 let standalone = pool
3544 .open_standalone_reader(StandaloneReaderPurpose::DiagnosticsIndependentSnapshot)
3545 .unwrap();
3546 let count: i64 = standalone
3547 .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
3548 .unwrap();
3549 assert_eq!(count, 1);
3550 drop(standalone);
3551 drop(pool);
3552
3553 assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
3554 assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
3555 assert_eq!(std::fs::read(&snapshot_shm).unwrap(), shm_before);
3556 assert_eq!(directory_entries(dir.path()), entries_before);
3557
3558 for (path, permissions) in snapshot_paths.into_iter().zip(original_permissions) {
3559 std::fs::set_permissions(path, permissions).unwrap();
3560 }
3561 drop(source_conn);
3562 }
3563
3564 #[cfg(unix)]
3569 #[test]
3570 fn read_only_frozen_wal_symlink_reads_target_frames_without_mutation() {
3571 use std::os::unix::fs::symlink;
3572
3573 let dir = tempfile::tempdir().unwrap();
3574 let source = dir.path().join("wal-source.db");
3575 let snapshot = dir.path().join("frozen-target.db");
3576 let alias = dir.path().join("frozen-alias.db");
3577 let source_wal = sqlite_sidecar(&source, "-wal");
3578 let source_shm = sqlite_sidecar(&source, "-shm");
3579 let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
3580 let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
3581 let alias_wal = sqlite_sidecar(&alias, "-wal");
3582 let alias_shm = sqlite_sidecar(&alias, "-shm");
3583
3584 let source_conn = Connection::open(&source).unwrap();
3585 let mode: String = source_conn
3586 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
3587 .unwrap();
3588 assert_eq!(mode.to_ascii_lowercase(), "wal");
3589 source_conn
3590 .pragma_update(None, "wal_autocheckpoint", 0)
3591 .unwrap();
3592 source_conn
3593 .execute_batch(
3594 "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
3595 INSERT INTO snapshot_row(body) VALUES ('visible-through-target-wal');",
3596 )
3597 .unwrap();
3598 assert!(source_wal.exists() && source_shm.exists());
3599
3600 std::fs::copy(&source, &snapshot).unwrap();
3601 std::fs::copy(&source_wal, &snapshot_wal).unwrap();
3602 std::fs::copy(&source_shm, &snapshot_shm).unwrap();
3603 symlink(&snapshot, &alias).unwrap();
3604 assert!(!alias_wal.exists() && !alias_shm.exists());
3605
3606 let snapshot_paths = [&snapshot, &snapshot_wal, &snapshot_shm];
3607 let original_permissions =
3608 snapshot_paths.map(|path| std::fs::metadata(path).unwrap().permissions());
3609 for path in snapshot_paths {
3610 let mut permissions = std::fs::metadata(path).unwrap().permissions();
3611 permissions.set_readonly(true);
3612 std::fs::set_permissions(path, permissions).unwrap();
3613 }
3614
3615 let main_before = std::fs::read(&snapshot).unwrap();
3616 let wal_before = std::fs::read(&snapshot_wal).unwrap();
3617 let shm_before = std::fs::read(&snapshot_shm).unwrap();
3618 let entries_before = directory_entries(dir.path());
3619
3620 let pool = ConnectionPool::new(PoolConfig {
3621 path: Some(alias.clone()),
3622 read_only: true,
3623 write_queue_enabled: Some(false),
3624 ..PoolConfig::for_test()
3625 })
3626 .unwrap();
3627 let reader = pool.reader().unwrap();
3628 let body: String = reader
3629 .conn()
3630 .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
3631 .unwrap();
3632 assert_eq!(body, "visible-through-target-wal");
3633 drop(reader);
3634 let standalone = pool
3635 .open_standalone_reader(StandaloneReaderPurpose::DiagnosticsIndependentSnapshot)
3636 .unwrap();
3637 let count: i64 = standalone
3638 .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
3639 .unwrap();
3640 assert_eq!(count, 1);
3641 drop(standalone);
3642 drop(pool);
3643
3644 assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
3645 assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
3646 assert_eq!(std::fs::read(&snapshot_shm).unwrap(), shm_before);
3647 assert_eq!(directory_entries(dir.path()), entries_before);
3648 assert!(!alias_wal.exists() && !alias_shm.exists());
3649
3650 for (path, permissions) in snapshot_paths.into_iter().zip(original_permissions) {
3651 std::fs::set_permissions(path, permissions).unwrap();
3652 }
3653 drop(source_conn);
3654 }
3655
3656 #[test]
3662 fn read_only_clean_wal_snapshot_is_sidecar_free() {
3663 let dir = tempfile::tempdir().unwrap();
3664 let path = dir.path().join("clean snapshot ?#%.db");
3665 {
3666 let conn = Connection::open(&path).unwrap();
3667 let mode: String = conn
3668 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
3669 .unwrap();
3670 assert_eq!(mode.to_ascii_lowercase(), "wal");
3671 conn.execute_batch(
3672 "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
3673 INSERT INTO snapshot_row(body) VALUES ('checkpointed');",
3674 )
3675 .unwrap();
3676 }
3677 let wal = sqlite_sidecar(&path, "-wal");
3678 let shm = sqlite_sidecar(&path, "-shm");
3679 assert!(!wal.exists() && !shm.exists());
3680 assert!(sqlite_header_uses_wal(&path).unwrap());
3681
3682 let original_permissions = std::fs::metadata(&path).unwrap().permissions();
3683 let mut read_only_permissions = original_permissions.clone();
3684 read_only_permissions.set_readonly(true);
3685 std::fs::set_permissions(&path, read_only_permissions).unwrap();
3686 let main_before = std::fs::read(&path).unwrap();
3687 let entries_before = directory_entries(dir.path());
3688
3689 let pool = ConnectionPool::new(PoolConfig {
3690 path: Some(path.clone()),
3691 read_only: true,
3692 write_queue_enabled: Some(false),
3693 ..PoolConfig::for_test()
3694 })
3695 .unwrap();
3696 let reader = pool.reader().unwrap();
3697 let body: String = reader
3698 .conn()
3699 .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
3700 .unwrap();
3701 assert_eq!(body, "checkpointed");
3702 drop(reader);
3703 let standalone = pool
3704 .open_standalone_reader(StandaloneReaderPurpose::DiagnosticsIndependentSnapshot)
3705 .unwrap();
3706 let count: i64 = standalone
3707 .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
3708 .unwrap();
3709 assert_eq!(count, 1);
3710 drop(standalone);
3711 drop(pool);
3712
3713 assert_eq!(std::fs::read(&path).unwrap(), main_before);
3714 assert_eq!(directory_entries(dir.path()), entries_before);
3715 assert!(!wal.exists() && !shm.exists());
3716 std::fs::set_permissions(&path, original_permissions).unwrap();
3717 }
3718
3719 #[test]
3724 fn read_only_live_rollback_journal_keeps_change_detection() {
3725 let dir = tempfile::tempdir().unwrap();
3726 let path = dir.path().join("live-delete-journal.db");
3727 let writer = Connection::open(&path).unwrap();
3728 writer
3729 .execute_batch("CREATE TABLE live_row(id INTEGER PRIMARY KEY);")
3730 .unwrap();
3731
3732 let pool = ConnectionPool::new(PoolConfig {
3733 path: Some(path),
3734 read_only: true,
3735 write_queue_enabled: Some(false),
3736 ..PoolConfig::for_test()
3737 })
3738 .unwrap();
3739 {
3740 let reader = pool.reader().unwrap();
3741 let count: i64 = reader
3742 .conn()
3743 .query_row("SELECT COUNT(*) FROM live_row", [], |row| row.get(0))
3744 .unwrap();
3745 assert_eq!(count, 0);
3746 }
3747
3748 writer
3749 .execute("INSERT INTO live_row DEFAULT VALUES", [])
3750 .unwrap();
3751 let reader = pool.reader().unwrap();
3752 let count: i64 = reader
3753 .conn()
3754 .query_row("SELECT COUNT(*) FROM live_row", [], |row| row.get(0))
3755 .unwrap();
3756 assert_eq!(
3757 count, 1,
3758 "rollback-journal read-only connections must retain live change detection"
3759 );
3760 }
3761
3762 #[test]
3771 fn pooled_reader_return_clears_temp_schema_and_attached_databases() {
3772 let dir = tempfile::tempdir().unwrap();
3773 let path = dir.path().join("pooled-reader-reset.db");
3774 {
3775 let seed = Connection::open(&path).unwrap();
3776 seed.execute_batch("CREATE TABLE main_row(id INTEGER PRIMARY KEY);")
3777 .unwrap();
3778 }
3779
3780 let secret_path = dir.path().join("secret.db");
3781 {
3782 let secret = Connection::open(&secret_path).unwrap();
3783 secret
3784 .execute_batch(
3785 "CREATE TABLE secret_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
3786 INSERT INTO secret_row(body) VALUES ('leaked-across-checkouts');",
3787 )
3788 .unwrap();
3789 }
3790
3791 let pool = ConnectionPool::new(PoolConfig {
3792 path: Some(path),
3793 max_readers: 1,
3794 write_queue_enabled: Some(false),
3795 ..PoolConfig::default()
3796 })
3797 .unwrap();
3798
3799 {
3800 let reader = pool.reader().unwrap();
3801 reader
3802 .conn()
3803 .execute_batch("CREATE TEMP TABLE leaked_temp(id INTEGER PRIMARY KEY);")
3804 .unwrap();
3805 reader
3806 .conn()
3807 .execute_batch(&format!(
3808 "ATTACH DATABASE '{}' AS secret;",
3809 secret_path.display()
3810 ))
3811 .unwrap();
3812 let leaked_count: i64 = reader
3813 .conn()
3814 .query_row("SELECT COUNT(*) FROM secret.secret_row", [], |row| {
3815 row.get(0)
3816 })
3817 .unwrap();
3818 assert_eq!(leaked_count, 1);
3819 reader.mark_dirty();
3825 }
3826
3827 let reader = pool.reader().unwrap();
3828 let temp_table_survived: i64 = reader
3829 .conn()
3830 .query_row(
3831 "SELECT COUNT(*) FROM sqlite_temp_master WHERE name = 'leaked_temp'",
3832 [],
3833 |row| row.get(0),
3834 )
3835 .unwrap();
3836 assert_eq!(
3837 temp_table_survived, 0,
3838 "a TEMP table from an earlier checkout must not survive pooled reader reuse"
3839 );
3840
3841 let attachment_survived: i64 = reader
3842 .conn()
3843 .query_row(
3844 "SELECT COUNT(*) FROM pragma_database_list WHERE name = 'secret'",
3845 [],
3846 |row| row.get(0),
3847 )
3848 .unwrap();
3849 assert_eq!(
3850 attachment_survived, 0,
3851 "an ATTACHed database from an earlier checkout must not survive pooled reader reuse"
3852 );
3853 }
3854
3855 #[test]
3863 fn writable_schema_evasion_of_the_temp_catalog_scan_still_disqualifies_reuse() {
3864 let dir = tempfile::tempdir().unwrap();
3865 let path = dir.path().join("writable-schema-evasion.db");
3866 {
3867 let seed = Connection::open(&path).unwrap();
3868 seed.execute_batch("CREATE TABLE main_row(id INTEGER PRIMARY KEY);")
3869 .unwrap();
3870 }
3871
3872 let pool = ConnectionPool::new(PoolConfig {
3873 path: Some(path),
3874 max_readers: 1,
3875 write_queue_enabled: Some(false),
3876 ..PoolConfig::default()
3877 })
3878 .unwrap();
3879
3880 {
3881 let reader = pool.reader().unwrap();
3882 reader
3883 .conn()
3884 .execute_batch("CREATE TEMP TABLE leaked_temp(id INTEGER PRIMARY KEY);")
3885 .unwrap();
3886 reader
3887 .conn()
3888 .execute_batch(
3889 "PRAGMA writable_schema = ON; \
3890 DELETE FROM sqlite_temp_master WHERE name = 'leaked_temp';",
3891 )
3892 .unwrap();
3893 let visible: i64 = reader
3894 .conn()
3895 .query_row(
3896 "SELECT COUNT(*) FROM sqlite_temp_master WHERE name = 'leaked_temp'",
3897 [],
3898 |row| row.get(0),
3899 )
3900 .unwrap();
3901 assert_eq!(
3902 visible, 0,
3903 "the evasion must actually hide the row from the catalog scan"
3904 );
3905 reader.mark_dirty();
3906 }
3907
3908 let reader = pool.reader().unwrap();
3909 let leaked_still_queryable = reader
3914 .conn()
3915 .query_row("SELECT COUNT(*) FROM leaked_temp", [], |row| {
3916 row.get::<_, i64>(0)
3917 })
3918 .is_ok();
3919 assert!(
3920 !leaked_still_queryable,
3921 "a writable_schema evasion of the catalog scan must still disqualify the \
3922 connection via the settings check"
3923 );
3924 }
3925
3926 #[test]
3934 fn a_checkout_that_never_marks_dirty_skips_the_catalog_scan_entirely() {
3935 let dir = tempfile::tempdir().unwrap();
3936 let path = dir.path().join("typed-read-skips-scan.db");
3937 {
3938 let seed = Connection::open(&path).unwrap();
3939 seed.execute_batch("CREATE TABLE main_row(id INTEGER PRIMARY KEY);")
3940 .unwrap();
3941 }
3942
3943 let pool = ConnectionPool::new(PoolConfig {
3944 path: Some(path),
3945 max_readers: 1,
3946 write_queue_enabled: Some(false),
3947 ..PoolConfig::default()
3948 })
3949 .unwrap();
3950
3951 {
3952 let reader = pool.reader().unwrap();
3953 reader
3956 .conn()
3957 .execute_batch("CREATE TEMP TABLE survivor(id INTEGER PRIMARY KEY);")
3958 .unwrap();
3959 }
3960
3961 let reader = pool.reader().unwrap();
3962 let survived = reader
3963 .conn()
3964 .query_row("SELECT COUNT(*) FROM survivor", [], |row| {
3965 row.get::<_, i64>(0)
3966 })
3967 .is_ok();
3968 assert!(
3969 survived,
3970 "a non-dirty checkout must return without running the catalog scan at all, \
3971 so a TEMP table it left behind is still visible on the next checkout"
3972 );
3973 }
3974
3975 #[test]
3982 fn busy_timeout_and_cache_size_changes_do_not_survive_a_dirty_pooled_checkout() {
3983 let dir = tempfile::tempdir().unwrap();
3984 let path = dir.path().join("settings-evasion.db");
3985 {
3986 let seed = Connection::open(&path).unwrap();
3987 seed.execute_batch("CREATE TABLE main_row(id INTEGER PRIMARY KEY);")
3988 .unwrap();
3989 }
3990
3991 let pool = ConnectionPool::new(PoolConfig {
3992 path: Some(path),
3993 max_readers: 1,
3994 write_queue_enabled: Some(false),
3995 ..PoolConfig::default()
3996 })
3997 .unwrap();
3998 let default_busy_timeout_ms = i64::try_from(pool.config().busy_timeout.as_millis())
3999 .expect("configured busy_timeout fits i64 millis");
4000
4001 {
4002 let reader = pool.reader().unwrap();
4003 reader
4004 .conn()
4005 .pragma_update(None, "busy_timeout", 1i64)
4006 .unwrap();
4007 reader
4008 .conn()
4009 .pragma_update(None, "cache_size", -64i64)
4010 .unwrap();
4011 reader
4012 .conn()
4013 .query_row("SELECT 1", [], |row| row.get::<_, i64>(0))
4014 .unwrap();
4015 reader.mark_dirty();
4016 }
4017
4018 let reader = pool.reader().unwrap();
4019 let busy_timeout: i64 = reader
4020 .conn()
4021 .query_row("PRAGMA busy_timeout", [], |row| row.get(0))
4022 .unwrap();
4023 let cache_size: i64 = reader
4024 .conn()
4025 .query_row("PRAGMA cache_size", [], |row| row.get(0))
4026 .unwrap();
4027 assert_eq!(
4028 busy_timeout, default_busy_timeout_ms,
4029 "busy_timeout must be restored to the pool's configured baseline"
4030 );
4031 assert_eq!(
4032 cache_size, -65536,
4033 "cache_size must be restored to the pool's configured baseline"
4034 );
4035 }
4036
4037 #[test]
4043 fn degraded_shared_reader_lease_clears_temp_schema_and_attached_databases_on_dirty_return() {
4044 let dir = tempfile::tempdir().unwrap();
4045 let path = dir.path().join("degraded-shared-reader-reset.db");
4046 {
4047 let seed = Connection::open(&path).unwrap();
4048 seed.execute_batch("CREATE TABLE main_row(id INTEGER PRIMARY KEY);")
4049 .unwrap();
4050 }
4051 let secret_path = dir.path().join("secret.db");
4052 {
4053 let secret = Connection::open(&secret_path).unwrap();
4054 secret
4055 .execute_batch(
4056 "CREATE TABLE secret_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
4057 INSERT INTO secret_row(body) VALUES ('leaked-across-checkouts');",
4058 )
4059 .unwrap();
4060 }
4061
4062 let pool = ConnectionPool::new(PoolConfig {
4063 path: Some(path),
4064 max_readers: 0,
4065 write_queue_enabled: Some(false),
4066 ..PoolConfig::default()
4067 })
4068 .unwrap();
4069
4070 {
4071 let reader = pool.reader().unwrap();
4072 reader
4073 .conn()
4074 .execute_batch("CREATE TEMP TABLE leaked_temp(id INTEGER PRIMARY KEY);")
4075 .unwrap();
4076 reader
4077 .conn()
4078 .execute_batch(&format!(
4079 "ATTACH DATABASE '{}' AS secret;",
4080 secret_path.display()
4081 ))
4082 .unwrap();
4083 reader.mark_dirty();
4084 }
4085
4086 let reader = pool.reader().unwrap();
4087 let temp_table_survived: i64 = reader
4088 .conn()
4089 .query_row(
4090 "SELECT COUNT(*) FROM sqlite_temp_master WHERE name = 'leaked_temp'",
4091 [],
4092 |row| row.get(0),
4093 )
4094 .unwrap();
4095 assert_eq!(
4096 temp_table_survived, 0,
4097 "a TEMP table must not survive a dirty degraded shared-reader-lease return"
4098 );
4099 let attachment_survived: i64 = reader
4100 .conn()
4101 .query_row(
4102 "SELECT COUNT(*) FROM pragma_database_list WHERE name = 'secret'",
4103 [],
4104 |row| row.get(0),
4105 )
4106 .unwrap();
4107 assert_eq!(
4108 attachment_survived, 0,
4109 "an ATTACHed database must not survive a dirty degraded shared-reader-lease return"
4110 );
4111 }
4112
4113 #[test]
4124 fn query_row_refuses_write_and_transaction_control_statements() {
4125 let pool = ConnectionPool::new(PoolConfig {
4126 path: None,
4127 ..PoolConfig::default()
4128 })
4129 .unwrap();
4130 pool.writer()
4131 .unwrap()
4132 .conn()
4133 .execute_batch("CREATE TABLE query_row_admission_probe(id INTEGER PRIMARY KEY);")
4134 .unwrap();
4135
4136 let reader = pool.reader().unwrap();
4137 for (label, sql) in [
4138 ("BEGIN", "BEGIN"),
4139 (
4140 "INSERT",
4141 "INSERT INTO query_row_admission_probe(id) VALUES (1)",
4142 ),
4143 (
4144 "CREATE TEMP TABLE",
4145 "CREATE TEMP TABLE query_row_admission_probe_temp(id INTEGER PRIMARY KEY)",
4146 ),
4147 ("setting PRAGMA", "PRAGMA journal_mode = OFF"),
4148 ] {
4149 let result = reader.query_row(sql, [], |row| row.get::<_, i64>(0));
4150 assert!(
4151 result.is_err(),
4152 "query_row must refuse {label} ({sql:?}); got {result:?}"
4153 );
4154 }
4155
4156 let row_count: i64 = reader
4157 .query_row(
4158 "SELECT COUNT(*) FROM query_row_admission_probe",
4159 [],
4160 |row| row.get(0),
4161 )
4162 .expect("an admitted SELECT must still succeed");
4163 assert_eq!(
4164 row_count, 0,
4165 "a refused INSERT must never have reached SQLite"
4166 );
4167 let temp_table_survived: i64 = reader
4168 .query_row(
4169 "SELECT COUNT(*) FROM sqlite_temp_master \
4170 WHERE name = 'query_row_admission_probe_temp'",
4171 [],
4172 |row| row.get(0),
4173 )
4174 .expect("an admitted SELECT must still succeed");
4175 assert_eq!(
4176 temp_table_survived, 0,
4177 "a refused CREATE TEMP TABLE must never have reached SQLite"
4178 );
4179 let journal_mode: String = reader
4180 .query_row("PRAGMA journal_mode", [], |row| row.get(0))
4181 .expect("the read-only journal_mode PRAGMA form must still be admitted");
4182 assert_ne!(
4183 journal_mode.to_ascii_lowercase(),
4184 "off",
4185 "a refused setting PRAGMA must never have reached SQLite"
4186 );
4187 }
4188
4189 #[test]
4197 fn discarded_reader_replacement_open_failure_is_recorded() {
4198 let dir = tempfile::tempdir().unwrap();
4199 let path = dir.path().join("discard-replacement-failure.db");
4200 {
4201 let seed = Connection::open(&path).unwrap();
4202 seed.execute_batch("CREATE TABLE t(id INTEGER PRIMARY KEY);")
4203 .unwrap();
4204 }
4205 let pool = ConnectionPool::new(PoolConfig {
4206 path: Some(path.clone()),
4207 max_readers: 1,
4208 write_queue_enabled: Some(false),
4209 ..PoolConfig::default()
4210 })
4211 .unwrap();
4212
4213 let before = pool.reader_acquisition_snapshot();
4214
4215 let reader = pool.reader().unwrap();
4216 reader.discard();
4217 std::fs::remove_file(&path).unwrap();
4220 for suffix in ["-wal", "-shm"] {
4221 let _ = std::fs::remove_file(sqlite_sidecar(&path, suffix));
4222 }
4223 drop(reader);
4224
4225 let after = pool.reader_acquisition_snapshot();
4226 assert_eq!(
4227 after.reader_replacement_open_failures - before.reader_replacement_open_failures,
4228 1,
4229 "a non-reusable checkout's failed replacement open must be recorded, not silently \
4230 swallowed"
4231 );
4232 }
4233
4234 #[test]
4239 fn read_only_live_wal_with_writable_shm_is_refused_without_mutation() {
4240 let dir = tempfile::tempdir().unwrap();
4241 let path = dir.path().join("live-wal.db");
4242 let wal = sqlite_sidecar(&path, "-wal");
4243 let shm = sqlite_sidecar(&path, "-shm");
4244 let writer = Connection::open(&path).unwrap();
4245 writer.pragma_update(None, "journal_mode", "WAL").unwrap();
4246 writer
4247 .execute_batch(
4248 "CREATE TABLE live_row(id INTEGER PRIMARY KEY);\
4249 INSERT INTO live_row DEFAULT VALUES;",
4250 )
4251 .unwrap();
4252 assert!(wal.exists() && shm.exists());
4253
4254 let main_before = std::fs::read(&path).unwrap();
4255 let wal_before = std::fs::read(&wal).unwrap();
4256 let shm_before = std::fs::read(&shm).unwrap();
4257 let error = match ConnectionPool::new(PoolConfig {
4258 path: Some(path.clone()),
4259 read_only: true,
4260 write_queue_enabled: Some(false),
4261 ..PoolConfig::for_test()
4262 }) {
4263 Ok(_) => panic!("a live WAL database with writable -shm must fail closed"),
4264 Err(error) => error,
4265 };
4266 assert!(
4267 error
4268 .to_string()
4269 .contains("writable WAL shared-memory sidecar"),
4270 "diagnostic must explain how to freeze the snapshot: {error}"
4271 );
4272 assert_eq!(std::fs::read(&path).unwrap(), main_before);
4273 assert_eq!(std::fs::read(&wal).unwrap(), wal_before);
4274 assert_eq!(std::fs::read(&shm).unwrap(), shm_before);
4275
4276 drop(writer);
4277 }
4278
4279 #[test]
4290 fn pool_drop_never_leaves_a_reader_as_the_last_connection_closed() {
4291 let dir = tempfile::tempdir().unwrap();
4292 let path = dir.path().join("close-order.db");
4293 let wal = sqlite_sidecar(&path, "-wal");
4294 let shm = sqlite_sidecar(&path, "-shm");
4295
4296 let pool = ConnectionPool::new(PoolConfig {
4297 path: Some(path.clone()),
4298 max_readers: 1,
4299 write_queue_enabled: Some(false),
4300 ..PoolConfig::default()
4301 })
4302 .unwrap();
4303
4304 pool.writer()
4305 .unwrap()
4306 .execute_batch(
4307 "CREATE TABLE close_order_row(id INTEGER PRIMARY KEY);\
4308 INSERT INTO close_order_row DEFAULT VALUES;",
4309 )
4310 .unwrap();
4311 assert!(
4312 wal.exists(),
4313 "a WAL-mode write must leave a -wal sidecar before the pool drops"
4314 );
4315
4316 drop(pool);
4317
4318 assert!(
4319 !wal.exists(),
4320 "the pool's last connection to close must be writable enough to checkpoint -wal away"
4321 );
4322 assert!(
4323 !shm.exists(),
4324 "the pool's last connection to close must be writable enough to checkpoint -shm away"
4325 );
4326 }
4327
4328 #[cfg(unix)]
4333 #[test]
4334 fn read_only_live_wal_symlink_rejects_target_writable_shm_without_mutation() {
4335 use std::os::unix::fs::symlink;
4336
4337 let dir = tempfile::tempdir().unwrap();
4338 let target = dir.path().join("live-target.db");
4339 let alias = dir.path().join("live-alias.db");
4340 let target_wal = sqlite_sidecar(&target, "-wal");
4341 let target_shm = sqlite_sidecar(&target, "-shm");
4342 let alias_wal = sqlite_sidecar(&alias, "-wal");
4343 let alias_shm = sqlite_sidecar(&alias, "-shm");
4344
4345 let writer = Connection::open(&target).unwrap();
4346 writer.pragma_update(None, "journal_mode", "WAL").unwrap();
4347 writer
4348 .execute_batch(
4349 "CREATE TABLE live_row(id INTEGER PRIMARY KEY);\
4350 INSERT INTO live_row DEFAULT VALUES;",
4351 )
4352 .unwrap();
4353 assert!(target_wal.exists() && target_shm.exists());
4354 symlink(&target, &alias).unwrap();
4355 assert!(!alias_wal.exists() && !alias_shm.exists());
4356
4357 let main_before = std::fs::read(&target).unwrap();
4358 let wal_before = std::fs::read(&target_wal).unwrap();
4359 let shm_before = std::fs::read(&target_shm).unwrap();
4360 let entries_before = directory_entries(dir.path());
4361
4362 let error = match ConnectionPool::new(PoolConfig {
4363 path: Some(alias),
4364 read_only: true,
4365 write_queue_enabled: Some(false),
4366 ..PoolConfig::for_test()
4367 }) {
4368 Ok(_) => panic!("a symlink must not hide the target's writable -shm"),
4369 Err(error) => error,
4370 };
4371 assert!(
4372 error
4373 .to_string()
4374 .contains("writable WAL shared-memory sidecar"),
4375 "diagnostic must identify the canonical target's live sidecar: {error}"
4376 );
4377 assert_eq!(std::fs::read(&target).unwrap(), main_before);
4378 assert_eq!(std::fs::read(&target_wal).unwrap(), wal_before);
4379 assert_eq!(std::fs::read(&target_shm).unwrap(), shm_before);
4380 assert_eq!(directory_entries(dir.path()), entries_before);
4381 assert!(!alias_wal.exists() && !alias_shm.exists());
4382
4383 drop(writer);
4384 }
4385
4386 #[test]
4387 #[serial]
4388 fn pool_config_default_values_match_constants() {
4389 let _pool_env = clear_pool_env();
4394 let cfg = PoolConfig::default();
4395 assert_eq!(
4396 cfg.journal_size_limit_bytes,
4397 DEFAULT_JOURNAL_SIZE_LIMIT_BYTES
4398 );
4399 assert_eq!(cfg.busy_timeout, Duration::from_secs(30));
4400 assert_eq!(cfg.checkout_timeout, Duration::from_secs(5));
4401 }
4402
4403 #[test]
4404 #[serial]
4405 fn legacy_env_cannot_change_wal_autocheckpoint() {
4406 let _pool_env = clear_pool_env();
4407 std::env::set_var("KHIVE_WAL_AUTOCHECKPOINT_PAGES", "8000");
4408 let dir = tempfile::tempdir().unwrap();
4409 let path = dir.path().join("legacy_autocheckpoint_env.db");
4410 let pool = ConnectionPool::new(PoolConfig {
4411 path: Some(path),
4412 ..PoolConfig::for_test()
4413 })
4414 .expect("pool open");
4415 {
4416 let writer = pool.writer().expect("writer");
4417 assert_eq!(
4418 wal_autocheckpoint_pages(writer.conn()),
4419 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
4420 "the removed env override must not change the unclaimed fallback"
4421 );
4422 }
4423 pool.claim_checkpoint_ownership().expect("claim ownership");
4424 let writer = pool.writer().expect("writer after claim");
4425 assert_eq!(
4426 wal_autocheckpoint_pages(writer.conn()),
4427 0,
4428 "the removed env override must not change the claimed-owner setting"
4429 );
4430 std::env::remove_var("KHIVE_WAL_AUTOCHECKPOINT_PAGES");
4431 }
4432
4433 #[test]
4434 #[serial]
4435 fn pool_config_env_override_journal_size_limit() {
4436 std::env::set_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES", "134217728");
4437 let cfg = PoolConfig::default();
4438 std::env::remove_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES");
4439 assert_eq!(cfg.journal_size_limit_bytes, 134_217_728);
4440 }
4441
4442 #[test]
4443 #[serial]
4444 fn pool_config_env_override_busy_timeout() {
4445 std::env::set_var("KHIVE_BUSY_TIMEOUT_SECS", "60");
4446 let cfg = PoolConfig::default();
4447 std::env::remove_var("KHIVE_BUSY_TIMEOUT_SECS");
4448 assert_eq!(cfg.busy_timeout, Duration::from_secs(60));
4449 }
4450
4451 #[test]
4452 #[serial]
4453 fn pool_config_env_override_checkout_timeout() {
4454 std::env::set_var("KHIVE_CHECKOUT_TIMEOUT_SECS", "10");
4455 let cfg = PoolConfig::default();
4456 std::env::remove_var("KHIVE_CHECKOUT_TIMEOUT_SECS");
4457 assert_eq!(cfg.checkout_timeout, Duration::from_secs(10));
4458 }
4459
4460 #[test]
4461 #[serial]
4462 fn pool_config_write_queue_defaults_unset() {
4463 let _pool_env = clear_pool_env();
4464 let cfg = PoolConfig::default();
4465 assert_eq!(cfg.write_queue_enabled, None);
4466 assert_eq!(cfg.write_queue_capacity, DEFAULT_WRITE_QUEUE_CAPACITY);
4467 }
4468
4469 #[test]
4470 #[serial]
4471 fn clear_pool_env_restores_overrides_on_drop() {
4472 let _ambient_env = PoolEnvGuard::capture();
4473 std::env::set_var("KHIVE_BUSY_TIMEOUT_SECS", "73");
4474
4475 {
4476 let _pool_env = clear_pool_env();
4477 assert_eq!(std::env::var_os("KHIVE_BUSY_TIMEOUT_SECS"), None);
4478 }
4479
4480 assert_eq!(
4481 std::env::var_os("KHIVE_BUSY_TIMEOUT_SECS"),
4482 Some(std::ffi::OsString::from("73"))
4483 );
4484 }
4485
4486 #[test]
4487 #[serial]
4488 fn pool_config_env_override_write_queue_enabled() {
4489 std::env::set_var("KHIVE_WRITE_QUEUE", "1");
4490 let cfg = PoolConfig::default();
4491 std::env::remove_var("KHIVE_WRITE_QUEUE");
4492 assert_eq!(cfg.write_queue_enabled, Some(true));
4493 }
4494
4495 #[test]
4496 #[serial]
4497 fn pool_config_env_override_write_queue_enabled_accepts_true_case_insensitive() {
4498 std::env::set_var("KHIVE_WRITE_QUEUE", "True");
4499 let cfg = PoolConfig::default();
4500 std::env::remove_var("KHIVE_WRITE_QUEUE");
4501 assert_eq!(cfg.write_queue_enabled, Some(true));
4502 }
4503
4504 #[test]
4505 #[serial]
4506 fn pool_config_env_override_write_queue_enabled_accepts_zero_as_explicit_off() {
4507 std::env::set_var("KHIVE_WRITE_QUEUE", "0");
4508 let cfg = PoolConfig::default();
4509 std::env::remove_var("KHIVE_WRITE_QUEUE");
4510 assert_eq!(cfg.write_queue_enabled, Some(false));
4511 }
4512
4513 #[cfg(unix)]
4518 #[test]
4519 #[serial]
4520 fn pool_config_env_override_write_queue_non_unicode_value_is_explicit_off() {
4521 use std::os::unix::ffi::OsStrExt;
4522 let _pool_env = clear_pool_env();
4523 std::env::set_var(
4524 "KHIVE_WRITE_QUEUE",
4525 std::ffi::OsStr::from_bytes(b"\xff\xfe"),
4526 );
4527 let cfg = PoolConfig::default();
4528 assert_eq!(cfg.write_queue_enabled, Some(false));
4529 }
4530
4531 #[test]
4532 #[serial]
4533 fn pool_config_env_override_write_queue_invalid_value_is_explicit_off() {
4534 std::env::set_var("KHIVE_WRITE_QUEUE", "banana");
4538 let cfg = PoolConfig::default();
4539 std::env::remove_var("KHIVE_WRITE_QUEUE");
4540 assert_eq!(cfg.write_queue_enabled, Some(false));
4541 }
4542
4543 #[test]
4544 #[serial]
4545 fn pool_config_write_routing_strict_defaults_off() {
4546 let _pool_env = clear_pool_env();
4547 let cfg = PoolConfig::default();
4548 assert!(!cfg.write_routing_strict);
4549 }
4550
4551 #[test]
4552 #[serial]
4553 fn pool_config_env_override_write_routing_strict() {
4554 std::env::set_var("KHIVE_WRITE_ROUTING", "strict");
4555 let cfg = PoolConfig::default();
4556 std::env::remove_var("KHIVE_WRITE_ROUTING");
4557 assert!(cfg.write_routing_strict);
4558 }
4559
4560 #[test]
4561 #[serial]
4562 fn pool_config_env_override_write_routing_strict_case_insensitive() {
4563 std::env::set_var("KHIVE_WRITE_ROUTING", "STRICT");
4564 let cfg = PoolConfig::default();
4565 std::env::remove_var("KHIVE_WRITE_ROUTING");
4566 assert!(cfg.write_routing_strict);
4567 }
4568
4569 #[test]
4570 #[serial]
4571 fn pool_config_env_write_routing_ignores_unrecognized_value() {
4572 std::env::set_var("KHIVE_WRITE_ROUTING", "eventual");
4573 let cfg = PoolConfig::default();
4574 std::env::remove_var("KHIVE_WRITE_ROUTING");
4575 assert!(!cfg.write_routing_strict);
4576 }
4577
4578 #[test]
4579 #[serial]
4580 fn pool_config_env_override_write_queue_capacity() {
4581 std::env::set_var("KHIVE_WRITE_QUEUE_CAPACITY", "64");
4582 let cfg = PoolConfig::default();
4583 std::env::remove_var("KHIVE_WRITE_QUEUE_CAPACITY");
4584 assert_eq!(cfg.write_queue_capacity, 64);
4585 }
4586
4587 #[test]
4588 #[serial]
4589 fn pool_config_env_invalid_write_queue_capacity_falls_back_to_default() {
4590 std::env::set_var("KHIVE_WRITE_QUEUE_CAPACITY", "0");
4591 let cfg = PoolConfig::default();
4592 std::env::remove_var("KHIVE_WRITE_QUEUE_CAPACITY");
4593 assert_eq!(cfg.write_queue_capacity, DEFAULT_WRITE_QUEUE_CAPACITY);
4594 }
4595
4596 #[test]
4597 #[serial]
4598 fn pool_config_invalid_journal_size_limit_falls_back_to_default() {
4599 std::env::set_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES", "");
4600 let cfg = PoolConfig::default();
4601 std::env::remove_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES");
4602 assert_eq!(
4603 cfg.journal_size_limit_bytes,
4604 DEFAULT_JOURNAL_SIZE_LIMIT_BYTES
4605 );
4606 }
4607
4608 #[test]
4609 fn file_backed_pool_opens_successfully() {
4610 let dir = tempfile::tempdir().unwrap();
4611 let path = dir.path().join("test_pool.db");
4612 let cfg = PoolConfig {
4613 path: Some(path.clone()),
4614 ..PoolConfig::default()
4615 };
4616 let pool = ConnectionPool::new(cfg).expect("file-backed pool should open");
4617 assert!(path.exists());
4618 assert!(pool.max_readers() > 0);
4619 }
4620
4621 #[test]
4622 fn standalone_wal_writer_uses_configured_journal_size_limit() {
4623 let dir = tempfile::tempdir().unwrap();
4624 let path = dir.path().join("standalone_wal_journal_limit.db");
4625 let configured_limit = 12_345_678;
4626 let pool = ConnectionPool::new(PoolConfig {
4627 path: Some(path),
4628 journal_size_limit_bytes: configured_limit,
4629 write_queue_enabled: Some(false),
4630 ..PoolConfig::for_test()
4631 })
4632 .expect("WAL pool open");
4633
4634 let standalone = pool
4635 .open_standalone_writer_untracked()
4636 .expect("standalone WAL writer open");
4637 assert_eq!(current_journal_mode(&standalone).unwrap(), "wal");
4638 assert_eq!(journal_size_limit_bytes(&standalone), configured_limit);
4639 }
4640
4641 #[test]
4642 fn standalone_rollback_writer_keeps_sqlite_journal_size_limit() {
4643 let dir = tempfile::tempdir().unwrap();
4644 let path = dir.path().join("standalone_rollback_journal_limit.db");
4645 let sqlite_default = {
4646 let conn = Connection::open(&path).expect("seed rollback-journal database");
4647 assert_eq!(current_journal_mode(&conn).unwrap(), "delete");
4648 journal_size_limit_bytes(&conn)
4649 };
4650 let configured_limit = if sqlite_default == 12_345_678 {
4651 23_456_789
4652 } else {
4653 12_345_678
4654 };
4655 let pool = ConnectionPool::new(PoolConfig {
4656 path: Some(path),
4657 wal_mode: false,
4658 journal_size_limit_bytes: configured_limit,
4659 write_queue_enabled: Some(false),
4660 ..PoolConfig::for_test()
4661 })
4662 .expect("rollback-journal pool open");
4663
4664 let standalone = pool
4665 .open_standalone_writer_untracked()
4666 .expect("standalone rollback-journal writer open");
4667 assert_eq!(current_journal_mode(&standalone).unwrap(), "delete");
4668 assert_eq!(journal_size_limit_bytes(&standalone), sqlite_default);
4669 }
4670
4671 #[test]
4672 fn writer_connections_follow_checkpoint_ownership_claim() {
4673 let dir = tempfile::tempdir().unwrap();
4674 let path = dir.path().join("writer_autocheckpoint.db");
4675 let pool = ConnectionPool::new(PoolConfig {
4676 path: Some(path),
4677 write_queue_enabled: Some(false),
4678 ..PoolConfig::for_test()
4679 })
4680 .expect("pool open");
4681
4682 {
4686 let writer = pool.writer().expect("pooled writer");
4687 assert_eq!(
4688 wal_autocheckpoint_pages(writer.conn()),
4689 FALLBACK_WAL_AUTOCHECKPOINT_PAGES
4690 );
4691 }
4692 let standalone = pool
4693 .open_standalone_writer()
4694 .expect("standalone writer opened before any claim");
4695 assert_eq!(
4696 wal_autocheckpoint_pages(&standalone),
4697 FALLBACK_WAL_AUTOCHECKPOINT_PAGES
4698 );
4699 drop(standalone);
4700
4701 pool.claim_checkpoint_ownership().expect("claim ownership");
4705 {
4706 let writer = pool.writer().expect("pooled writer after claim");
4707 assert_eq!(wal_autocheckpoint_pages(writer.conn()), 0);
4708 }
4709 let claimed_standalone = pool
4710 .open_standalone_writer()
4711 .expect("standalone writer opened after the claim");
4712 assert_eq!(wal_autocheckpoint_pages(&claimed_standalone), 0);
4713 drop(claimed_standalone);
4714
4715 let later_infrastructure = pool
4716 .open_standalone_writer_untracked()
4717 .expect("later infrastructure writer");
4718 assert_eq!(wal_autocheckpoint_pages(&later_infrastructure), 0);
4719
4720 let memory_pool = ConnectionPool::new(PoolConfig {
4721 write_queue_enabled: Some(false),
4722 ..PoolConfig::default()
4723 })
4724 .expect("in-memory pool open");
4725 let memory_writer = memory_pool.writer().expect("in-memory writer");
4726 assert_eq!(
4727 wal_autocheckpoint_pages(memory_writer.conn()),
4728 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
4729 "an unclaimed in-memory pool keeps the bounded fallback"
4730 );
4731 }
4732
4733 #[test]
4734 fn standalone_writer_waits_for_checkpoint_claim_resolution() {
4735 let dir = tempfile::tempdir().unwrap();
4736 let path = dir.path().join("checkpoint_claim_race.db");
4737 let pool = Arc::new(
4738 ConnectionPool::new(PoolConfig {
4739 path: Some(path),
4740 checkout_timeout: Duration::from_secs(5),
4741 write_queue_enabled: Some(false),
4742 ..PoolConfig::for_test()
4743 })
4744 .expect("pool open"),
4745 );
4746
4747 let legacy_conn = pool.legacy_conn();
4748 let held_writer = legacy_conn.lock();
4749 let claim_start = Arc::new(std::sync::Barrier::new(2));
4750 let claim_pool = Arc::clone(&pool);
4751 let claim_thread_start = Arc::clone(&claim_start);
4752 let claim_thread = thread::spawn(move || {
4753 claim_thread_start.wait();
4754 claim_pool.claim_checkpoint_ownership()
4755 });
4756 claim_start.wait();
4757
4758 {
4759 let mut state = pool.checkpoint_ownership.state.lock();
4760 while state.phase != CheckpointOwnership::Claiming {
4761 pool.checkpoint_ownership.changed.wait(&mut state);
4762 }
4763 }
4764
4765 let open_start = Arc::new(std::sync::Barrier::new(2));
4766 let open_pool = Arc::clone(&pool);
4767 let open_thread_start = Arc::clone(&open_start);
4768 let open_thread = thread::spawn(move || {
4769 open_thread_start.wait();
4770 let conn = open_pool
4771 .open_standalone_writer()
4772 .expect("standalone writer after claim resolution");
4773 wal_autocheckpoint_pages(&conn)
4774 });
4775 open_start.wait();
4776
4777 {
4778 let mut state = pool.checkpoint_ownership.state.lock();
4779 while state.connection_waiters == 0 {
4780 pool.checkpoint_ownership.changed.wait(&mut state);
4781 }
4782 assert_eq!(state.phase, CheckpointOwnership::Claiming);
4783 }
4784
4785 drop(held_writer);
4786 claim_thread
4787 .join()
4788 .expect("claim thread joins")
4789 .expect("claim succeeds");
4790 assert_eq!(
4791 open_thread.join().expect("standalone-open thread joins"),
4792 0,
4793 "a writer open concurrent with a successful claim must inherit claimed ownership"
4794 );
4795 }
4796
4797 #[test]
4798 fn standalone_fallback_application_linearizes_before_claim_publication() {
4799 let dir = tempfile::tempdir().unwrap();
4800 let path = dir.path().join("checkpoint_open_before_claim.db");
4801 let pool = Arc::new(
4802 ConnectionPool::new(PoolConfig {
4803 path: Some(path),
4804 checkout_timeout: Duration::from_secs(5),
4805 write_queue_enabled: Some(false),
4806 ..PoolConfig::for_test()
4807 })
4808 .expect("pool open"),
4809 );
4810 let pause = Arc::new(CheckpointConnectionConfigPause::new());
4811 *pool.checkpoint_ownership.connection_config_pause.lock() = Some(Arc::clone(&pause));
4812
4813 let open_pool = Arc::clone(&pool);
4814 let open_thread = thread::spawn(move || {
4815 let conn = open_pool
4816 .open_standalone_writer()
4817 .expect("standalone writer opens");
4818 wal_autocheckpoint_pages(&conn)
4819 });
4820 pause.selected.wait();
4821 assert!(
4822 pool.checkpoint_ownership.state.try_lock().is_none(),
4823 "standalone selection must retain the ownership gate until its PRAGMA is applied"
4824 );
4825
4826 let (claim_observed_tx, claim_observed_rx) = std::sync::mpsc::sync_channel(0);
4827 *pool.checkpoint_ownership.claim_lock_observed.lock() = Some(claim_observed_tx);
4828 let claim_pool = Arc::clone(&pool);
4829 let claim_thread = thread::spawn(move || claim_pool.claim_checkpoint_ownership());
4830 assert!(
4831 claim_observed_rx
4832 .recv()
4833 .expect("claim reports whether it observed gate contention"),
4834 "the claim must attempt the gate between fallback selection and PRAGMA application"
4835 );
4836 pause.resume.wait();
4837
4838 assert_eq!(
4839 open_thread.join().expect("standalone-open thread joins"),
4840 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
4841 "an open linearized before the claim keeps the fallback"
4842 );
4843 claim_thread
4844 .join()
4845 .expect("claim thread joins")
4846 .expect("claim succeeds after standalone configuration");
4847 assert_eq!(pool.effective_wal_autocheckpoint_pages(), 0);
4848 }
4849
4850 #[test]
4851 fn failed_checkpoint_ownership_claim_keeps_fallback_and_can_be_retried() {
4852 let dir = tempfile::tempdir().unwrap();
4853 let path = dir.path().join("checkpoint_claim_retry.db");
4854 let pool = ConnectionPool::new(PoolConfig {
4855 path: Some(path),
4856 checkout_timeout: Duration::from_millis(1),
4857 write_queue_enabled: Some(false),
4858 ..PoolConfig::for_test()
4859 })
4860 .expect("pool open");
4861
4862 let legacy_conn = pool.legacy_conn();
4863 let held_writer = legacy_conn.lock();
4864 let error = pool
4865 .claim_checkpoint_ownership()
4866 .expect_err("the held pooled writer must make the claim time out");
4867 assert!(matches!(
4868 error,
4869 SqliteError::WriterPoolCheckoutTimeout { .. }
4870 ));
4871 assert_eq!(
4872 pool.effective_wal_autocheckpoint_pages(),
4873 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
4874 "a failed claim must leave later writer connections fallback-safe"
4875 );
4876
4877 let fallback_writer = pool
4878 .open_standalone_writer()
4879 .expect("standalone writer after failed claim");
4880 assert_eq!(
4881 wal_autocheckpoint_pages(&fallback_writer),
4882 FALLBACK_WAL_AUTOCHECKPOINT_PAGES
4883 );
4884 drop(fallback_writer);
4885
4886 drop(held_writer);
4887 pool.claim_checkpoint_ownership()
4888 .expect("the ownership claim remains retryable");
4889 assert_eq!(pool.effective_wal_autocheckpoint_pages(), 0);
4890 let writer = pool.writer().expect("pooled writer after successful retry");
4891 assert_eq!(wal_autocheckpoint_pages(writer.conn()), 0);
4892 }
4893
4894 #[test]
4895 fn threshold_crossing_commits_do_not_run_an_implicit_checkpoint_once_claimed() {
4896 const FORMER_AUTOCHECKPOINT_THRESHOLD_PAGES: i64 = FALLBACK_WAL_AUTOCHECKPOINT_PAGES as i64;
4897
4898 let dir = tempfile::tempdir().unwrap();
4899 let path = dir.path().join("no_implicit_checkpoint.db");
4900 let pool = ConnectionPool::new(PoolConfig {
4901 path: Some(path),
4902 write_queue_enabled: Some(false),
4903 ..PoolConfig::for_test()
4904 })
4905 .expect("pool open");
4906 pool.claim_checkpoint_ownership()
4907 .expect("claim ownership for the dedicated-owner posture");
4908 let writer = pool.writer().expect("pooled writer");
4909 writer
4910 .execute_batch("CREATE TABLE blobs (value BLOB NOT NULL)")
4911 .expect("create fixture table");
4912
4913 let page_size: i64 = writer
4914 .pragma_query_value(None, "page_size", |row| row.get(0))
4915 .expect("read page size");
4916 let payload_bytes = page_size * 32;
4917 for _ in 0..160 {
4918 writer
4919 .execute(
4920 "INSERT INTO blobs (value) VALUES (zeroblob(?1))",
4921 [payload_bytes],
4922 )
4923 .expect("autocommit fixture row");
4924 }
4925
4926 let log_frames: i64 = writer
4927 .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| row.get(1))
4928 .expect("observe WAL frame count");
4929 assert!(
4930 log_frames > FORMER_AUTOCHECKPOINT_THRESHOLD_PAGES,
4931 "the commit sequence must retain more than the former automatic threshold; \
4932 observed {log_frames} frames"
4933 );
4934 }
4935
4936 #[test]
4944 fn unclaimed_pool_retains_bounded_autocheckpoint_reclamation() {
4945 let dir = tempfile::tempdir().unwrap();
4946 let path = dir.path().join("bounded_fallback_reclamation.db");
4947 let pool = ConnectionPool::new(PoolConfig {
4948 path: Some(path),
4949 write_queue_enabled: Some(false),
4950 ..PoolConfig::for_test()
4951 })
4952 .expect("pool open");
4953 let writer = pool.writer().expect("pooled writer");
4954 writer
4955 .execute_batch("CREATE TABLE blobs (value BLOB NOT NULL)")
4956 .expect("create fixture table");
4957
4958 let page_size: i64 = writer
4959 .pragma_query_value(None, "page_size", |row| row.get(0))
4960 .expect("read page size");
4961 let payload_bytes = page_size * 32;
4962 for _ in 0..160 {
4963 writer
4964 .execute(
4965 "INSERT INTO blobs (value) VALUES (zeroblob(?1))",
4966 [payload_bytes],
4967 )
4968 .expect("autocommit fixture row");
4969 }
4970
4971 let log_frames: i64 = writer
4977 .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| row.get(1))
4978 .expect("observe WAL frame count");
4979 assert!(
4980 log_frames < FALLBACK_WAL_AUTOCHECKPOINT_PAGES as i64,
4981 "an unclaimed pool must reclaim WAL frames via the bounded autocheckpoint; \
4982 observed {log_frames} retained frames"
4983 );
4984 }
4985
4986 #[tokio::test]
4987 #[serial]
4988 async fn unset_write_queue_resolves_on_for_file_backed_pool() {
4989 let _pool_env = clear_pool_env();
4990 let dir = tempfile::tempdir().unwrap();
4991 let path = dir.path().join("unset_file_backed.db");
4992 let pool = ConnectionPool::new(PoolConfig {
4993 path: Some(path),
4994 write_queue_enabled: None,
4995 ..PoolConfig::for_test()
4996 })
4997 .expect("file-backed pool should open");
4998 assert_eq!(pool.config().write_queue_enabled, Some(true));
4999 assert!(
5002 pool.writer_task_handle()
5003 .expect("spawn inside a runtime context must not error")
5004 .is_some(),
5005 "resolved-on file-backed pool must actually spawn the writer task"
5006 );
5007 }
5008
5009 #[tokio::test]
5010 #[serial]
5011 async fn unset_write_queue_resolves_off_for_memory_backed_pool() {
5012 let _pool_env = clear_pool_env();
5013 let pool = ConnectionPool::new(PoolConfig {
5014 path: None,
5015 write_queue_enabled: None,
5016 ..PoolConfig::default()
5017 })
5018 .expect("in-memory pool should open");
5019 assert_eq!(pool.config().write_queue_enabled, Some(false));
5020 assert!(
5023 pool.writer_task_handle()
5024 .expect("disabled queue must resolve without error")
5025 .is_none(),
5026 "resolved-off in-memory pool must not spawn a writer task"
5027 );
5028 }
5029
5030 #[test]
5031 #[serial]
5032 fn explicit_false_stays_off_for_file_backed_pool() {
5033 let _pool_env = clear_pool_env();
5034 let dir = tempfile::tempdir().unwrap();
5035 let path = dir.path().join("explicit_false_file_backed.db");
5036 let pool = ConnectionPool::new(PoolConfig {
5037 path: Some(path),
5038 write_queue_enabled: Some(false),
5039 ..PoolConfig::for_test()
5040 })
5041 .expect("file-backed pool should open");
5042 assert_eq!(pool.config().write_queue_enabled, Some(false));
5043 }
5044
5045 #[tokio::test]
5046 #[serial]
5047 async fn explicit_true_stays_on_for_memory_backed_pool() {
5048 let _pool_env = clear_pool_env();
5049 let pool = ConnectionPool::new(PoolConfig {
5050 path: None,
5051 write_queue_enabled: Some(true),
5052 ..PoolConfig::default()
5053 })
5054 .expect("in-memory pool should open");
5055 assert_eq!(pool.config().write_queue_enabled, Some(true));
5056 assert!(
5062 pool.writer_task_handle()
5063 .expect("spawn degrade must resolve without error")
5064 .is_none(),
5065 "explicit-on in-memory pool must degrade to no writer task"
5066 );
5067 assert_eq!(
5068 pool.writer_task_spawn_count(),
5069 1,
5070 "the spawn attempt must happen exactly once and degrade, not retry"
5071 );
5072 assert!(
5073 pool.take_writer_task_join().is_none(),
5074 "a degraded spawn stores no JoinHandle to drain"
5075 );
5076 }
5077
5078 #[test]
5079 #[serial]
5080 fn explicit_true_on_memory_pool_warns_but_false_and_none_do_not() {
5081 let _pool_env = clear_pool_env();
5082 let messages = Arc::new(std::sync::Mutex::new(Vec::new()));
5083 let subscriber = WarningCapture {
5084 messages: Arc::clone(&messages),
5085 };
5086
5087 tracing::subscriber::with_default(subscriber, || {
5088 let _explicit_true = ConnectionPool::new(PoolConfig {
5089 path: None,
5090 write_queue_enabled: Some(true),
5091 ..PoolConfig::default()
5092 })
5093 .expect("in-memory pool should open");
5094 let _explicit_false = ConnectionPool::new(PoolConfig {
5095 path: None,
5096 write_queue_enabled: Some(false),
5097 ..PoolConfig::default()
5098 })
5099 .expect("in-memory pool should open");
5100 let _unset = ConnectionPool::new(PoolConfig {
5101 path: None,
5102 write_queue_enabled: None,
5103 ..PoolConfig::default()
5104 })
5105 .expect("in-memory pool should open");
5106 });
5107
5108 let messages = messages.lock().unwrap();
5109 assert_eq!(
5110 messages
5111 .iter()
5112 .filter(|message| message.contains("write queue explicitly requested"))
5113 .count(),
5114 1,
5115 "only an explicit in-memory queue request should warn: {messages:?}"
5116 );
5117 let warning = messages
5118 .iter()
5119 .find(|message| message.contains("write queue explicitly requested"))
5120 .expect("explicit in-memory queue warning should be captured");
5121 assert!(
5122 warning.contains("in-memory pools cannot host a writer task"),
5123 "warning must explain why the request is inert: {messages:?}"
5124 );
5125 }
5126
5127 #[test]
5128 fn standalone_writer_open_counts_its_connection_class_once() {
5129 let dir = tempfile::tempdir().unwrap();
5130 let path = dir.path().join("standalone_writer_counter.db");
5131 let pool = ConnectionPool::new(PoolConfig {
5132 path: Some(path),
5133 ..PoolConfig::for_test()
5134 })
5135 .expect("file-backed pool");
5136
5137 let _standalone = pool
5138 .open_standalone_writer()
5139 .expect("standalone writer opens");
5140
5141 assert_eq!(
5142 pool.writer_acquisition_snapshot(),
5143 WriterAcquisitionSnapshot {
5144 acquisitions: 1,
5145 pooled_acquisitions: 0,
5146 standalone_acquisitions: 1,
5147 writer_task_acquisitions: 0,
5148 timeouts: 0,
5149 writer_task_begin_busy: 0,
5150 writer_task_begin_busy_absorbed: 0,
5151 writer_task_begin_errors: 0,
5152 writer_task_request_failures: 0,
5153 writer_task_side_effects_unknown: 0,
5154 },
5155 "the public standalone boundary must contribute to the aggregate exactly once"
5156 );
5157 }
5158
5159 #[test]
5160 fn reader_snapshot_tracks_pool_saturation_hold_lifecycle_and_exception_classes() {
5161 let dir = tempfile::tempdir().unwrap();
5162 let path = dir.path().join("reader_acquisition_counters.db");
5163 let pool = ConnectionPool::new(PoolConfig {
5164 path: Some(path),
5165 max_readers: 1,
5166 checkout_timeout: Duration::from_millis(2),
5167 ..PoolConfig::default()
5168 })
5169 .expect("file-backed pool");
5170
5171 assert_eq!(
5172 pool.reader_acquisition_snapshot(),
5173 ReaderAcquisitionSnapshot {
5174 reader_admission_capacity: 1,
5175 available_reader_admission_slots: 1,
5176 ..ReaderAcquisitionSnapshot::default()
5177 }
5178 );
5179
5180 let held = pool.reader().expect("first pooled checkout succeeds");
5181 assert_eq!(
5182 pool.reader_acquisition_snapshot(),
5183 ReaderAcquisitionSnapshot {
5184 reader_admission_capacity: 1,
5185 available_reader_admission_slots: 0,
5186 acquisitions: 1,
5187 pooled_checkouts: 1,
5188 active_pooled_checkouts: 1,
5189 peak_active_pooled_checkouts: 1,
5190 ..ReaderAcquisitionSnapshot::default()
5191 }
5192 );
5193
5194 let timeout = match pool.reader() {
5195 Ok(_) => panic!("the sole live checkout must exhaust bounded admission"),
5196 Err(error) => error,
5197 };
5198 assert!(
5199 matches!(
5200 &timeout,
5201 SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
5202 if code.code == rusqlite::ErrorCode::DatabaseBusy
5203 ),
5204 "reader saturation must keep the pool-exhausted classification: {timeout}"
5205 );
5206 drop(held);
5207
5208 let explicit = pool
5209 .open_standalone_reader(StandaloneReaderPurpose::ExplicitSqlReadTransaction)
5210 .expect("explicit read-transaction exception opens");
5211 drop(explicit);
5212 let infrastructure = pool
5213 .open_standalone_reader(StandaloneReaderPurpose::DiagnosticsIndependentSnapshot)
5214 .expect("infrastructure exception opens");
5215 drop(infrastructure);
5216
5217 let snapshot = pool.reader_acquisition_snapshot();
5218 assert_eq!(snapshot.acquisitions, 2);
5219 assert_eq!(snapshot.pooled_checkouts, 1);
5220 assert_eq!(snapshot.standalone_opens, 1);
5221 assert_eq!(snapshot.infrastructure_standalone_opens, 1);
5222 assert_eq!(snapshot.checkout_timeouts, 1);
5223 assert_eq!(snapshot.active_pooled_checkouts, 0);
5224 assert_eq!(snapshot.peak_active_pooled_checkouts, 1);
5225 assert_eq!(snapshot.completed_pooled_checkouts, 1);
5226 assert!(
5227 snapshot.max_completed_hold_micros > 0,
5228 "the held checkout's completed lifecycle must expose nonzero hold evidence"
5229 );
5230 }
5231
5232 #[test]
5233 fn in_memory_pool_degrades_to_single_connection() {
5234 let cfg = PoolConfig {
5235 path: None,
5236 ..PoolConfig::default()
5237 };
5238 let pool = ConnectionPool::new(cfg).expect("in-memory pool should open");
5239 assert_eq!(pool.max_readers(), 0);
5240 }
5241
5242 #[test]
5243 fn writer_checkout_and_release_works() {
5244 let cfg = PoolConfig {
5245 path: None,
5246 ..PoolConfig::default()
5247 };
5248 let pool = ConnectionPool::new(cfg).unwrap();
5249 {
5250 let _writer = pool.writer().expect("writer checkout should succeed");
5251 }
5252 let _writer2 = pool
5254 .writer()
5255 .expect("second writer checkout should succeed");
5256 }
5257
5258 #[test]
5259 fn writer_checkout_snapshot_counts_successes_and_timeouts_at_the_pool_boundary() {
5260 let cfg = PoolConfig {
5261 path: None,
5262 checkout_timeout: Duration::from_millis(1),
5263 ..PoolConfig::default()
5264 };
5265 let pool = ConnectionPool::new(cfg).unwrap();
5266
5267 assert_eq!(
5268 pool.writer_acquisition_snapshot(),
5269 WriterAcquisitionSnapshot::default()
5270 );
5271
5272 let held = pool.writer().expect("first checkout succeeds");
5273 let error = match pool.writer() {
5274 Ok(_) => panic!("the held pool mutex must force a finite-wait timeout"),
5275 Err(error) => error,
5276 };
5277 assert!(
5278 matches!(
5279 &error,
5280 SqliteError::WriterPoolCheckoutTimeout { timeout }
5281 if *timeout == Duration::from_millis(1)
5282 ),
5283 "timeout must have a stable, structurally matchable stage: {error}"
5284 );
5285 assert_eq!(
5286 pool.writer_acquisition_snapshot(),
5287 WriterAcquisitionSnapshot {
5288 acquisitions: 1,
5289 pooled_acquisitions: 1,
5290 standalone_acquisitions: 0,
5291 writer_task_acquisitions: 0,
5292 timeouts: 1,
5293 writer_task_begin_busy: 0,
5297 writer_task_begin_busy_absorbed: 0,
5298 writer_task_begin_errors: 0,
5299 writer_task_request_failures: 0,
5300 writer_task_side_effects_unknown: 0,
5301 }
5302 );
5303
5304 drop(held);
5305 let _reacquired = pool.writer().expect("checkout succeeds after release");
5306 assert_eq!(
5307 pool.writer_acquisition_snapshot(),
5308 WriterAcquisitionSnapshot {
5309 acquisitions: 2,
5310 pooled_acquisitions: 2,
5311 standalone_acquisitions: 0,
5312 writer_task_acquisitions: 0,
5313 timeouts: 1,
5314 writer_task_begin_busy: 0,
5315 writer_task_begin_busy_absorbed: 0,
5316 writer_task_begin_errors: 0,
5317 writer_task_request_failures: 0,
5318 writer_task_side_effects_unknown: 0,
5319 }
5320 );
5321 }
5322
5323 #[test]
5324 fn zero_wait_maintenance_skip_is_not_reported_as_a_checkout_timeout() {
5325 let pool = ConnectionPool::new(PoolConfig::default()).unwrap();
5326 let held = pool.writer().expect("finite-wait checkout succeeds");
5327 let before = pool.writer_acquisition_snapshot();
5328
5329 assert!(
5330 pool.try_writer_nowait().is_err(),
5331 "zero-wait maintenance checkout must skip while held"
5332 );
5333
5334 assert_eq!(
5335 pool.writer_acquisition_snapshot(),
5336 before,
5337 "a checkpoint-style zero-wait skip is not a finite-wait checkout timeout"
5338 );
5339 drop(held);
5340 }
5341
5342 #[test]
5346 #[serial(tx_registry)]
5347 fn writer_guard_transaction_registers_during_closure_only() {
5348 let cfg = PoolConfig {
5349 path: None,
5350 ..PoolConfig::default()
5351 };
5352 let pool = ConnectionPool::new(cfg).unwrap();
5353 let guard = pool.writer().unwrap();
5354
5355 let mut seen_during_closure = false;
5356 let result: Result<(), SqliteError> = guard.transaction(|_conn| {
5357 seen_during_closure = khive_storage::tx_registry::snapshot()
5358 .iter()
5359 .any(|(_, label)| label.as_deref() == Some("writer_guard_tx"));
5360 Ok(())
5361 });
5362 result.expect("transaction should commit");
5363
5364 assert!(
5365 seen_during_closure,
5366 "expected a writer_guard_tx entry visible inside the closure"
5367 );
5368 assert!(
5369 !khive_storage::tx_registry::snapshot()
5370 .iter()
5371 .any(|(_, label)| label.as_deref() == Some("writer_guard_tx")),
5372 "expected the entry to be gone after the transaction completes"
5373 );
5374 }
5375
5376 #[test]
5380 fn writer_task_handle_fails_loud_without_tokio_runtime() {
5381 let dir = tempfile::tempdir().unwrap();
5382 let path = dir.path().join("writer_task_no_runtime.db");
5383 let cfg = PoolConfig {
5384 path: Some(path),
5385 write_queue_enabled: Some(true),
5386 ..PoolConfig::for_test()
5387 };
5388 let pool = ConnectionPool::new(cfg).expect("file-backed pool should open");
5389
5390 let result = pool.writer_task_handle();
5391
5392 assert!(
5393 matches!(result, Err(StorageError::WriterTaskNoRuntime)),
5394 "expected Err(StorageError::WriterTaskNoRuntime) outside a Tokio \
5395 runtime, got {result:?}"
5396 );
5397 assert_eq!(
5398 pool.writer_task_spawn_count(),
5399 0,
5400 "the guard must reject before ever attempting tokio::spawn"
5401 );
5402 }
5403
5404 #[test]
5407 fn strict_writer_task_for_write_preserves_missing_runtime_error() {
5408 let dir = tempfile::tempdir().unwrap();
5409 let pool = ConnectionPool::new(PoolConfig {
5410 path: Some(dir.path().join("strict_writer_task_no_runtime.db")),
5411 write_queue_enabled: Some(true),
5412 write_routing_strict: true,
5413 ..PoolConfig::for_test()
5414 })
5415 .expect("file-backed pool should open");
5416
5417 let result = pool.writer_task_for_write(None, "strict_test_write");
5418
5419 assert!(
5420 matches!(result, Err(StorageError::WriterTaskNoRuntime)),
5421 "strict routing must preserve WriterTaskNoRuntime, got {result:?}"
5422 );
5423 assert_eq!(pool.writer_task_spawn_count(), 0);
5424 }
5425
5426 #[tokio::test]
5431 async fn take_writer_task_join_returns_some_once_then_none() {
5432 let dir = tempfile::tempdir().unwrap();
5433 let path = dir.path().join("join_lifecycle.db");
5434 let pool = ConnectionPool::new(PoolConfig {
5435 path: Some(path),
5436 write_queue_enabled: Some(true),
5437 ..PoolConfig::for_test()
5438 })
5439 .expect("file-backed pool should open");
5440
5441 assert!(
5444 pool.take_writer_task_join().is_none(),
5445 "before spawn there is no JoinHandle to take"
5446 );
5447 assert!(!pool.writer_task_join_was_stored());
5448 pool.writer_task_handle()
5449 .expect("runtime is present")
5450 .expect("write queue enabled must spawn a writer task");
5451 assert!(pool.writer_task_join_was_stored());
5452
5453 let join = pool
5454 .take_writer_task_join()
5455 .expect("the first take must return the spawned task's JoinHandle");
5456 assert!(
5457 pool.take_writer_task_join().is_none(),
5458 "the second take must return None — the handle is one-shot"
5459 );
5460
5461 drop(pool);
5467 tokio::time::timeout(Duration::from_secs(5), join)
5468 .await
5469 .expect("the writer task must exit once every handle clone is dropped")
5470 .expect("the writer task must not panic");
5471 }
5472
5473 #[cfg(debug_assertions)]
5477 #[tokio::test]
5478 #[should_panic(expected = "writer task JoinHandle stored twice")]
5479 async fn set_writer_task_join_second_store_trips_debug_assert() {
5480 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
5481 pool.set_writer_task_join(tokio::spawn(async {}));
5482 pool.set_writer_task_join(tokio::spawn(async {}));
5483 }
5484
5485 #[cfg(debug_assertions)]
5490 #[tokio::test]
5491 #[should_panic(expected = "writer task JoinHandle stored twice")]
5492 async fn set_writer_task_join_second_store_after_take_trips_debug_assert() {
5493 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
5494 pool.set_writer_task_join(tokio::spawn(async {}));
5495 assert!(pool.take_writer_task_join().is_some());
5496 pool.set_writer_task_join(tokio::spawn(async {}));
5497 }
5498
5499 #[cfg(not(debug_assertions))]
5505 #[tokio::test]
5506 async fn set_writer_task_join_first_wins_keeps_existing_handle() {
5507 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
5508
5509 let (first_done_tx, first_done_rx) = tokio::sync::oneshot::channel::<()>();
5512 let first = tokio::spawn(async move {
5513 let _ = first_done_tx.send(());
5514 });
5515 let (_never_sent, never_rx) = tokio::sync::oneshot::channel::<()>();
5519 let second = tokio::spawn(async move {
5520 let _ = never_rx.await;
5521 });
5522
5523 pool.set_writer_task_join(first);
5524 pool.set_writer_task_join(second);
5525
5526 let taken = pool
5527 .take_writer_task_join()
5528 .expect("the first handle must still be stored");
5529 tokio::time::timeout(Duration::from_secs(5), taken)
5530 .await
5531 .expect("stored handle must be the first task's; the second never completes")
5532 .expect("the first task must not panic");
5533 assert!(
5534 first_done_rx.await.is_ok(),
5535 "completing the taken handle must mean the FIRST task ran to completion"
5536 );
5537 }
5538
5539 #[test]
5544 #[serial(pool_cwd)]
5545 fn mint_db_identity_alias_convergence() {
5546 let dir = tempfile::tempdir().unwrap();
5547 let real_dir = dir.path().join("real");
5548 fs::create_dir(&real_dir).unwrap();
5549 let db_path = real_dir.join("khive.db");
5550 fs::write(&db_path, b"").unwrap();
5551
5552 #[cfg(unix)]
5553 let dir_symlink = dir.path().join("dir_link");
5554 #[cfg(unix)]
5555 let file_symlink = dir.path().join("file_link.db");
5556 #[cfg(unix)]
5557 {
5558 std::os::unix::fs::symlink(&real_dir, &dir_symlink).unwrap();
5559 std::os::unix::fs::symlink(&db_path, &file_symlink).unwrap();
5560 }
5561
5562 let (via_real, canonical_real) = mint_db_identity(&db_path).unwrap();
5563
5564 let relative_result = {
5566 let _cwd = CwdGuard::enter(&real_dir);
5567 mint_db_identity(&PathBuf::from("khive.db"))
5568 };
5569 let (via_relative, canonical_relative) = relative_result.unwrap();
5570 assert_eq!(canonical_real, canonical_relative);
5571 assert_eq!(via_real, via_relative);
5572
5573 #[cfg(unix)]
5574 {
5575 let (via_dir_symlink, canonical_dir_symlink) =
5576 mint_db_identity(&dir_symlink.join("khive.db")).unwrap();
5577 assert_eq!(canonical_real, canonical_dir_symlink);
5578 assert_eq!(via_real, via_dir_symlink);
5579
5580 let (via_file_symlink, canonical_file_symlink) =
5581 mint_db_identity(&file_symlink).unwrap();
5582 assert_eq!(canonical_real, canonical_file_symlink);
5583 assert_eq!(via_real, via_file_symlink);
5584 }
5585
5586 let bare_name_result = {
5588 let _cwd = CwdGuard::enter(&real_dir);
5589 mint_db_identity(&PathBuf::from("khive.db"))
5590 };
5591 let (via_bare_name, canonical_bare_name) = bare_name_result.unwrap();
5592 assert_eq!(canonical_real, canonical_bare_name);
5593 assert_eq!(via_real, via_bare_name);
5594 }
5595
5596 #[test]
5608 #[serial(pool_cwd)]
5609 fn sidecar_dir_for_alias_convergence() {
5610 let dir = tempfile::tempdir().unwrap();
5611 let real_dir = dir.path().join("real");
5612 fs::create_dir(&real_dir).unwrap();
5613 let db_path = real_dir.join("khive.db");
5614 fs::write(&db_path, b"").unwrap();
5615
5616 #[cfg(unix)]
5617 let dir_symlink = dir.path().join("dir_link");
5618 #[cfg(unix)]
5619 let file_symlink = dir.path().join("file_link.db");
5620 #[cfg(unix)]
5621 {
5622 std::os::unix::fs::symlink(&real_dir, &dir_symlink).unwrap();
5623 std::os::unix::fs::symlink(&db_path, &file_symlink).unwrap();
5624 }
5625
5626 let pool_for = |path: &Path| -> Arc<ConnectionPool> {
5627 let cfg = PoolConfig {
5628 path: Some(path.to_path_buf()),
5629 ..PoolConfig::for_test()
5630 };
5631 Arc::new(ConnectionPool::new(cfg).expect("file-backed pool should open"))
5632 };
5633 let sidecar_of = |pool: &ConnectionPool| -> PathBuf {
5634 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"))
5635 };
5636
5637 let via_real = pool_for(&db_path);
5638 let sidecar_real = sidecar_of(&via_real);
5639
5640 let via_relative = {
5641 let _cwd = CwdGuard::enter(&real_dir);
5642 pool_for(Path::new("khive.db"))
5643 };
5644 assert_eq!(
5645 sidecar_real,
5646 sidecar_of(&via_relative),
5647 "a relative spelling of the same database must derive the same sidecar directory"
5648 );
5649
5650 #[cfg(unix)]
5651 {
5652 let via_dir_symlink = pool_for(&dir_symlink.join("khive.db"));
5653 assert_eq!(
5654 sidecar_real,
5655 sidecar_of(&via_dir_symlink),
5656 "opening through a directory symlink must derive the same sidecar directory"
5657 );
5658
5659 let via_file_symlink = pool_for(&file_symlink);
5660 assert_eq!(
5661 sidecar_real,
5662 sidecar_of(&via_file_symlink),
5663 "opening through a file-level symlink must derive the same sidecar directory"
5664 );
5665 }
5666
5667 let via_bare_name = {
5668 let _cwd = CwdGuard::enter(&real_dir);
5669 pool_for(Path::new("khive.db"))
5670 };
5671 assert_eq!(
5672 sidecar_real,
5673 sidecar_of(&via_bare_name),
5674 "a bare file name resolved against the current directory must derive the same \
5675 sidecar directory"
5676 );
5677 }
5678
5679 #[cfg(unix)]
5685 #[test]
5686 fn mint_db_identity_dangling_symlink_first_open_convergence() {
5687 let dir = tempfile::tempdir().unwrap();
5688 let target = dir.path().join("target.db");
5689 let link = dir.path().join("link.db");
5690 std::os::unix::fs::symlink(&target, &link).unwrap();
5691 assert!(!target.exists(), "target must not exist yet (dangling)");
5692
5693 let (via_dangling_link, canonical_via_link) = mint_db_identity(&link).unwrap();
5694
5695 fs::write(&target, b"").unwrap();
5698 let (via_target, canonical_via_target) = mint_db_identity(&target).unwrap();
5699
5700 assert_eq!(canonical_via_link, canonical_via_target);
5701 assert_eq!(via_dangling_link, via_target);
5702 }
5703
5704 #[test]
5707 fn mint_db_identity_missing_parent_fails() {
5708 let dir = tempfile::tempdir().unwrap();
5709 let missing = dir.path().join("nonexistent_subdir").join("khive.db");
5710 let result = mint_db_identity(&missing);
5711 assert!(
5712 result.is_err(),
5713 "minting must fail when the parent directory does not exist"
5714 );
5715 }
5716
5717 #[cfg(unix)]
5720 #[test]
5721 fn mint_db_identity_non_utf8_path_round_trips() {
5722 use std::ffi::OsStr;
5723 use std::os::unix::ffi::OsStrExt;
5724
5725 let dir = tempfile::tempdir().unwrap();
5726 let raw_name = OsStr::from_bytes(b"khive-\xffdb.sqlite");
5728 let db_path = dir.path().join(raw_name);
5729 if let Err(e) = fs::write(&db_path, b"") {
5734 eprintln!(
5735 "skipping mint_db_identity_non_utf8_path_round_trips: filesystem rejected a \
5736 non-UTF-8 file name ({e}); this platform's filesystem does not support the \
5737 case under test"
5738 );
5739 return;
5740 }
5741
5742 let (identity, canonical) = mint_db_identity(&db_path).unwrap();
5743 assert_eq!(canonical.file_name().unwrap(), raw_name);
5744
5745 let (identity_again, canonical_again) = mint_db_identity(&db_path).unwrap();
5746 assert_eq!(identity, identity_again);
5747 assert_eq!(canonical, canonical_again);
5748 }
5749
5750 fn admission_identity_after_refusal(pool: &ConnectionPool) -> String {
5751 let held = pool.reader().expect("hold the sole reader");
5752 let Err(error) = pool.resolve_reader_checkout(
5753 StorageCapability::Sql,
5754 "identity_read",
5755 pool.reader_until(|| false),
5756 ) else {
5757 panic!("held reader must exhaust this pool's admission budget");
5758 };
5759 assert!(
5760 error.is_retryable(),
5761 "admission refusal must remain retryable"
5762 );
5763 let display = error.to_string();
5764 let StorageError::AdmissionTimeout {
5765 operation,
5766 timeout_ms,
5767 pool_identity,
5768 } = error
5769 else {
5770 panic!("pool refusal must retain its typed admission classification");
5771 };
5772 assert_eq!(
5773 operation, "identity_read",
5774 "pool identity must not alter operation"
5775 );
5776 assert_eq!(timeout_ms, 20);
5777 let identity = pool_identity.expect("typed admission error must name the pool");
5778 assert!(
5779 !identity.contains('/') && !identity.contains('\\'),
5780 "pool identity must never contain a directory or separator: {identity}"
5781 );
5782 assert_eq!(
5783 display,
5784 format!("admission timeout during identity_read after 20ms (pool: {identity})"),
5785 "admission error text must name the refusing pool"
5786 );
5787 drop(held);
5788 identity
5789 }
5790
5791 fn identity_test_pool(path: Option<PathBuf>, read_only: bool) -> ConnectionPool {
5792 ConnectionPool::new(PoolConfig {
5793 path,
5794 read_only,
5795 max_readers: 1,
5796 checkout_timeout: Duration::from_millis(20),
5797 ..PoolConfig::default()
5798 })
5799 .unwrap()
5800 }
5801
5802 #[test]
5803 fn reader_admission_timeout_identifies_the_refusing_pool() {
5804 let dir = tempfile::tempdir().unwrap();
5805 for read_only in [false, true] {
5806 let name = format!("identity-{}.db", uuid::Uuid::new_v4());
5807 let path = dir.path().join(&name);
5808 {
5809 let seed = Connection::open(&path).unwrap();
5810 seed.execute_batch("CREATE TABLE seed (id INTEGER)")
5811 .unwrap();
5812 }
5813 let canonical = fs::canonicalize(&path).unwrap();
5814 let configured = dir.path().join(".").join(&name);
5815 assert_ne!(configured.as_os_str(), canonical.as_os_str());
5816 let pool = identity_test_pool(Some(configured), read_only);
5817 assert_eq!(
5818 admission_identity_after_refusal(&pool),
5819 name,
5820 "typed admission field must contain only the canonical file name"
5821 );
5822 #[cfg(unix)]
5823 {
5824 let alias = dir
5825 .path()
5826 .join(format!("alias-{}.db", uuid::Uuid::new_v4()));
5827 std::os::unix::fs::symlink(&canonical, &alias).unwrap();
5828 let alias_pool = identity_test_pool(Some(alias), read_only);
5829 assert_eq!(
5830 admission_identity_after_refusal(&alias_pool),
5831 name,
5832 "symlink spelling must not change the canonical database file name"
5833 );
5834 }
5835 }
5836 let memory = identity_test_pool(None, false);
5837 assert_eq!(admission_identity_after_refusal(&memory), ":memory:");
5838 }
5839
5840 #[test]
5841 fn reader_admission_identity_hash_is_build_stable() {
5842 #[cfg(unix)]
5845 assert_eq!(
5846 pool_identity_suffix(Path::new("/khive/pool/khive.db")),
5847 "8fa8797b",
5848 "suffix must match the published Unix SHA-256 vector"
5849 );
5850 #[cfg(windows)]
5851 assert_eq!(
5852 pool_identity_suffix(Path::new("/khive/pool/khive.db")),
5853 "1186b990",
5854 "suffix must match the published Windows SHA-256 vector"
5855 );
5856 }
5857
5858 fn assert_disambiguated_identity(identity: &str, basename: &str) {
5859 let suffix = identity
5860 .strip_prefix(&format!("{basename}#"))
5861 .expect("different open files with the same basename need a hash suffix");
5862 assert_eq!(
5863 suffix.len(),
5864 8,
5865 "disambiguation needs exactly eight hex digits"
5866 );
5867 assert!(
5868 suffix.bytes().all(|b| b.is_ascii_hexdigit()),
5869 "disambiguation must contain only a hash, never directory text"
5870 );
5871 }
5872
5873 #[test]
5874 fn reader_admission_identity_disambiguates_open_files() {
5875 let first_dir = tempfile::tempdir().unwrap();
5876 let second_dir = tempfile::tempdir().unwrap();
5877 let basename = format!("collision-{}.db", uuid::Uuid::new_v4());
5878 let first = identity_test_pool(Some(first_dir.path().join(&basename)), false);
5879 assert_eq!(admission_identity_after_refusal(&first), basename);
5880 let second_path = second_dir.path().join(&basename);
5881 let second = identity_test_pool(Some(second_path.clone()), false);
5882 let first_identity = admission_identity_after_refusal(&first);
5883 let second_identity = admission_identity_after_refusal(&second);
5884 assert_disambiguated_identity(&first_identity, &basename);
5885 assert_disambiguated_identity(&second_identity, &basename);
5886 assert_ne!(
5887 first_identity, second_identity,
5888 "distinct files need distinct identities"
5889 );
5890 assert_eq!(admission_identity_after_refusal(&first), first_identity);
5891 drop(second);
5892 assert_eq!(
5893 admission_identity_after_refusal(&first),
5894 basename,
5895 "closing the colliding store must remove its registry entry"
5896 );
5897 let reopened = identity_test_pool(Some(second_path), false);
5898 assert_eq!(admission_identity_after_refusal(&first), first_identity);
5899 assert_eq!(admission_identity_after_refusal(&reopened), second_identity);
5900 }
5901
5902 #[test]
5903 fn reader_admission_identity_same_path_pools_share_label() {
5904 let first_dir = tempfile::tempdir().unwrap();
5905 let second_dir = tempfile::tempdir().unwrap();
5906 let basename = format!("same-path-{}.db", uuid::Uuid::new_v4());
5907 let first = identity_test_pool(Some(first_dir.path().join(&basename)), false);
5908 let duplicate = identity_test_pool(Some(first_dir.path().join(".").join(&basename)), false);
5909 assert_eq!(
5910 admission_identity_after_refusal(&first),
5911 basename,
5912 "two pools on the same canonical path must not get a suffix"
5913 );
5914 assert_eq!(admission_identity_after_refusal(&duplicate), basename);
5915 let other = identity_test_pool(Some(second_dir.path().join(&basename)), false);
5916 let first_identity = admission_identity_after_refusal(&first);
5917 let other_identity = admission_identity_after_refusal(&other);
5918 assert_disambiguated_identity(&first_identity, &basename);
5919 assert_disambiguated_identity(&other_identity, &basename);
5920 assert_eq!(admission_identity_after_refusal(&duplicate), first_identity);
5921 drop(first);
5922 assert_eq!(
5923 admission_identity_after_refusal(&other),
5924 other_identity,
5925 "dropping one pool must retain the other pool's path registration"
5926 );
5927 drop(duplicate);
5928 assert_eq!(
5929 admission_identity_after_refusal(&other),
5930 basename,
5931 "dropping the final pool must remove the path registration"
5932 );
5933 }
5934
5935 #[test]
5941 fn resolve_reader_checkout_maps_each_arm_distinctly() {
5942 let pool = ConnectionPool::new(PoolConfig {
5943 path: None,
5944 ..PoolConfig::default()
5945 })
5946 .unwrap();
5947
5948 let guard = pool
5949 .resolve_reader_checkout(
5950 StorageCapability::Sql,
5951 "arm_checked_out",
5952 pool.reader_until(|| false),
5953 )
5954 .expect("an uncontended checkout must pass the guard through");
5955 drop(guard);
5956
5957 let Err(cancelled) =
5958 pool.resolve_reader_checkout(StorageCapability::Sql, "arm_cancelled", Ok(None))
5959 else {
5960 panic!("a cancelled checkout must be refused");
5961 };
5962 assert!(
5963 matches!(cancelled, StorageError::Timeout { .. }),
5964 "cancellation/deadline before checkout must be the non-retryable \
5965 Timeout, got {cancelled:?}"
5966 );
5967
5968 let Err(exhausted) = pool.resolve_reader_checkout(
5969 StorageCapability::Sql,
5970 "arm_exhausted",
5971 Err(pool_exhausted_error(Duration::from_millis(5), 1)),
5972 ) else {
5973 panic!("an exhausted checkout must be refused");
5974 };
5975 assert!(
5976 matches!(exhausted, StorageError::AdmissionTimeout { .. }),
5977 "the pool's own SQLITE_BUSY (checkout_timeout exhausted) must be \
5978 the retryable AdmissionTimeout, got {exhausted:?}"
5979 );
5980
5981 let Err(opaque) = pool.resolve_reader_checkout(
5982 StorageCapability::Entities,
5983 "arm_driver",
5984 Err(SqliteError::InvalidData("retired pooled writer".into())),
5985 ) else {
5986 panic!("an opaque checkout error must be refused");
5987 };
5988 assert!(
5989 matches!(
5990 &opaque,
5991 StorageError::Driver { capability, .. }
5992 if *capability == StorageCapability::Entities
5993 ),
5994 "any other checkout error must stay a non-retryable Driver failure \
5995 under the caller's capability, got {opaque:?}"
5996 );
5997 }
5998
5999 #[test]
6004 fn the_longest_completed_hold_names_the_operation_that_held_it() {
6005 let pool = ConnectionPool::new(PoolConfig {
6006 path: None,
6007 ..PoolConfig::default()
6008 })
6009 .unwrap();
6010
6011 for _ in 0..3 {
6012 let guard = pool
6013 .resolve_reader_checkout(
6014 StorageCapability::Sql,
6015 "fast_read",
6016 pool.reader_until(|| false),
6017 )
6018 .expect("a fast checkout resolves");
6019 drop(guard);
6020 }
6021
6022 let slow = pool
6023 .resolve_reader_checkout(
6024 StorageCapability::Sql,
6025 "slow_read",
6026 pool.reader_until(|| false),
6027 )
6028 .expect("the slow checkout resolves");
6029 thread::sleep(Duration::from_millis(20));
6032 drop(slow);
6033
6034 let snapshot = pool.reader_acquisition_snapshot();
6035 assert_eq!(
6036 snapshot.completed_pooled_checkouts, 4,
6037 "all four checkouts must complete through the pooled route, or the \
6038 attribution below is reading a population of one"
6039 );
6040 assert_eq!(
6041 snapshot.max_completed_hold_operation,
6042 Some("slow_read"),
6043 "the longest hold must name the operation that held it; got {:?} at \
6044 {} micros",
6045 snapshot.max_completed_hold_operation,
6046 snapshot.max_completed_hold_micros
6047 );
6048 }
6049
6050 #[test]
6055 fn a_checkout_taken_outside_the_resolve_route_reports_no_operation() {
6056 let pool = ConnectionPool::new(PoolConfig {
6057 path: None,
6058 ..PoolConfig::default()
6059 })
6060 .unwrap();
6061
6062 let guard = pool
6063 .reader_until(|| false)
6064 .expect("the checkout succeeds")
6065 .expect("the checkout is not cancelled");
6066 drop(guard);
6067
6068 let snapshot = pool.reader_acquisition_snapshot();
6069 assert_eq!(snapshot.completed_pooled_checkouts, 1);
6070 assert_eq!(
6071 snapshot.max_completed_hold_operation, None,
6072 "an unlabelled route must report no operation rather than borrow one"
6073 );
6074 }
6075}