1use crossbeam_queue::ArrayQueue;
3use parking_lot::{Condvar, Mutex};
4use rusqlite::hooks::{AuthContext, Authorization};
5use rusqlite::{Connection, OpenFlags};
6use std::fs;
7use std::io::Read as _;
8use std::ops::{Deref, DerefMut};
9use std::path::{Path, PathBuf};
10use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
11use std::sync::{Arc, OnceLock};
12use std::thread;
13use std::time::{Duration, Instant};
14use tokio::sync::Semaphore;
15
16use crate::error::SqliteError;
17use crate::writer_task::WriterTaskHandle;
18use khive_storage::error::StorageError;
19use khive_storage::tx_registry::{DbIdentity, TxOrigin};
20use khive_storage::StorageCapability;
21
22const CACHE_SIZE_KIB: &str = "-65536";
23const MMAP_SIZE_BYTES: &str = "1073741824";
24const DEFAULT_READER_CAP: usize = 8;
25
26const DEFAULT_JOURNAL_SIZE_LIMIT_BYTES: i64 = 67_108_864; const DEFAULT_WRITE_QUEUE_CAPACITY: usize = 256;
28
29pub(crate) const FALLBACK_WAL_AUTOCHECKPOINT_PAGES: u32 = 4_000;
37
38#[derive(Clone, Copy, Debug, PartialEq, Eq)]
39enum CheckpointOwnership {
40 Unclaimed,
41 Claiming,
42 Claimed,
43}
44
45struct CheckpointOwnershipState {
46 phase: CheckpointOwnership,
47 #[cfg(test)]
48 connection_waiters: usize,
49}
50
51#[cfg(test)]
52struct CheckpointConnectionConfigPause {
53 selected: std::sync::Barrier,
54 resume: std::sync::Barrier,
55}
56
57#[cfg(test)]
58impl CheckpointConnectionConfigPause {
59 fn new() -> Self {
60 Self {
61 selected: std::sync::Barrier::new(2),
62 resume: std::sync::Barrier::new(2),
63 }
64 }
65}
66
67struct CheckpointOwnershipGate {
68 state: Mutex<CheckpointOwnershipState>,
69 changed: Condvar,
70 #[cfg(test)]
71 connection_config_pause: Mutex<Option<Arc<CheckpointConnectionConfigPause>>>,
72 #[cfg(test)]
73 claim_lock_observed: Mutex<Option<std::sync::mpsc::SyncSender<bool>>>,
74}
75
76impl CheckpointOwnershipGate {
77 fn new() -> Self {
78 Self {
79 state: Mutex::new(CheckpointOwnershipState {
80 phase: CheckpointOwnership::Unclaimed,
81 #[cfg(test)]
82 connection_waiters: 0,
83 }),
84 changed: Condvar::new(),
85 #[cfg(test)]
86 connection_config_pause: Mutex::new(None),
87 #[cfg(test)]
88 claim_lock_observed: Mutex::new(None),
89 }
90 }
91
92 fn begin_claim(&self) -> bool {
95 #[cfg(test)]
96 let claim_lock_observed = self.claim_lock_observed.lock().take();
97 #[cfg(test)]
98 let mut state = if let Some(observed) = claim_lock_observed {
99 match self.state.try_lock() {
100 Some(state) => {
101 let _ = observed.send(false);
102 state
103 }
104 None => {
105 let _ = observed.send(true);
106 self.state.lock()
107 }
108 }
109 } else {
110 self.state.lock()
111 };
112 #[cfg(not(test))]
113 let mut state = self.state.lock();
114 loop {
115 match state.phase {
116 CheckpointOwnership::Unclaimed => {
117 state.phase = CheckpointOwnership::Claiming;
118 self.changed.notify_all();
119 return true;
120 }
121 CheckpointOwnership::Claiming => self.changed.wait(&mut state),
122 CheckpointOwnership::Claimed => return false,
123 }
124 }
125 }
126
127 fn finish_claim(&self, succeeded: bool) {
128 let mut state = self.state.lock();
129 debug_assert_eq!(state.phase, CheckpointOwnership::Claiming);
130 state.phase = if succeeded {
131 CheckpointOwnership::Claimed
132 } else {
133 CheckpointOwnership::Unclaimed
134 };
135 self.changed.notify_all();
136 }
137
138 fn settled_state(&self) -> parking_lot::MutexGuard<'_, CheckpointOwnershipState> {
139 let mut state = self.state.lock();
140 while state.phase == CheckpointOwnership::Claiming {
141 #[cfg(test)]
142 {
143 state.connection_waiters += 1;
144 self.changed.notify_all();
145 }
146 self.changed.wait(&mut state);
147 #[cfg(test)]
148 {
149 state.connection_waiters -= 1;
150 self.changed.notify_all();
151 }
152 }
153 state
154 }
155
156 #[cfg(test)]
157 fn wal_autocheckpoint_pages(&self) -> u32 {
158 let state = self.settled_state();
159 match state.phase {
160 CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
161 CheckpointOwnership::Claimed => 0,
162 CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
163 }
164 }
165
166 fn configure_wal_autocheckpoint(&self, conn: &Connection) -> Result<(), SqliteError> {
171 let state = self.settled_state();
172 let pages = match state.phase {
173 CheckpointOwnership::Unclaimed => FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
174 CheckpointOwnership::Claimed => 0,
175 CheckpointOwnership::Claiming => unreachable!("claim wait must settle the state"),
176 };
177 #[cfg(test)]
178 if let Some(pause) = self.connection_config_pause.lock().take() {
179 pause.selected.wait();
180 pause.resume.wait();
181 }
182 conn.pragma_update(None, "wal_autocheckpoint", pages)?;
183 drop(state);
184 Ok(())
185 }
186}
187
188fn deny_retired_writer(_context: AuthContext<'_>) -> Authorization {
189 Authorization::Deny
190}
191
192pub(crate) const TEST_HARNESS_ENV: &str = "KHIVE_TEST_HARNESS";
193
194#[derive(Clone, Debug)]
196pub struct PoolConfig {
197 pub path: Option<PathBuf>,
199 pub max_readers: usize,
201 pub wal_mode: bool,
203 pub busy_timeout: Duration,
207 pub checkout_timeout: Duration,
211 pub journal_size_limit_bytes: i64,
217 pub read_only: bool,
225 pub write_queue_enabled: Option<bool>,
254 pub write_queue_capacity: usize,
259 pub write_routing_strict: bool,
270 pub write_admission_deadline_ms: u64,
284 pub read_tx_max_age: Duration,
294}
295
296const WRITE_ADMISSION_DEADLINE_MS_RANGE: std::ops::RangeInclusive<u64> = 100..=10_000;
298const DEFAULT_WRITE_ADMISSION_DEADLINE_MS: u64 = 2000;
299
300impl Default for PoolConfig {
301 fn default() -> Self {
302 Self {
303 path: None,
304 max_readers: std::thread::available_parallelism()
305 .map(|n| n.get())
306 .unwrap_or(1)
307 .clamp(1, DEFAULT_READER_CAP),
308 wal_mode: true,
309 busy_timeout: Duration::from_secs(
310 std::env::var("KHIVE_BUSY_TIMEOUT_SECS")
311 .ok()
312 .and_then(|v| v.parse::<u64>().ok())
313 .unwrap_or(30),
314 ),
315 checkout_timeout: Duration::from_secs(
316 std::env::var("KHIVE_CHECKOUT_TIMEOUT_SECS")
317 .ok()
318 .and_then(|v| v.parse::<u64>().ok())
319 .unwrap_or(5),
320 ),
321 journal_size_limit_bytes: std::env::var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES")
322 .ok()
323 .and_then(|v| v.parse::<i64>().ok())
324 .unwrap_or(DEFAULT_JOURNAL_SIZE_LIMIT_BYTES),
325 read_only: false,
326 write_queue_enabled: std::env::var_os("KHIVE_WRITE_QUEUE").map(|v| {
331 v.to_str()
332 .is_some_and(|v| v == "1" || v.eq_ignore_ascii_case("true"))
333 }),
334 write_queue_capacity: std::env::var("KHIVE_WRITE_QUEUE_CAPACITY")
335 .ok()
336 .and_then(|v| v.parse::<usize>().ok())
337 .filter(|&n| n > 0)
338 .unwrap_or(DEFAULT_WRITE_QUEUE_CAPACITY),
339 write_routing_strict: std::env::var("KHIVE_WRITE_ROUTING")
340 .map(|v| v.eq_ignore_ascii_case("strict"))
341 .unwrap_or(false),
342 write_admission_deadline_ms: std::env::var("KHIVE_WRITE_ADMISSION_DEADLINE_MS")
343 .ok()
344 .and_then(|v| v.parse::<u64>().ok())
345 .unwrap_or(DEFAULT_WRITE_ADMISSION_DEADLINE_MS),
346 read_tx_max_age: crate::checkpoint::tx_age_thresholds_from_env(
347 Duration::from_secs(30),
348 Duration::from_secs(120),
349 )
350 .1,
351 }
352 }
353}
354
355fn refuse_home_data_store_in_tests(config: &PoolConfig) -> Result<(), SqliteError> {
372 if std::env::var(TEST_HARNESS_ENV).as_deref() != Ok("1") {
373 return Ok(());
374 }
375
376 let Some(path) = config.path.as_deref() else {
377 return Ok(());
378 };
379 if path
380 .as_os_str()
381 .as_encoded_bytes()
382 .get(..5)
383 .is_some_and(|prefix| prefix.eq_ignore_ascii_case(b"file:"))
384 {
385 return Err(SqliteError::InvalidData(format!(
386 "test harness refused SQLite URI database path {}; use a filesystem path outside \
387 HOME/.khive (deliberate sessions against a real store run the built binary \
388 directly, outside the Cargo test environment)",
389 path.display()
390 )));
391 }
392
393 let Some(home) = std::env::var_os("HOME") else {
394 return Ok(());
395 };
396 let canonical_path = canonicalize_deepest_existing(path)?;
397 let canonical_home_data_dir =
398 canonicalize_deepest_existing(&PathBuf::from(home).join(".khive"))?;
399 if canonical_path.starts_with(&canonical_home_data_dir) {
400 return Err(SqliteError::InvalidData(format!(
401 "test harness refused to open SQLite database under HOME/.khive: {} \
402 (deliberate sessions against a real store run the built binary directly, \
403 outside the Cargo test environment)",
404 canonical_path.display()
405 )));
406 }
407 Ok(())
408}
409
410fn canonicalize_deepest_existing(path: &Path) -> Result<PathBuf, SqliteError> {
411 let absolute = if path.is_absolute() {
412 path.to_path_buf()
413 } else {
414 std::env::current_dir().map_err(SqliteError::Io)?.join(path)
415 };
416
417 for ancestor in absolute.ancestors() {
418 match fs::canonicalize(ancestor) {
419 Ok(mut canonical) => {
420 let missing = absolute.strip_prefix(ancestor).map_err(|error| {
421 SqliteError::InvalidData(format!(
422 "failed to preserve missing path components for {}: {error}",
423 absolute.display()
424 ))
425 })?;
426 canonical.push(missing);
427 return Ok(canonical);
428 }
429 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
430 Err(error) => {
431 return Err(SqliteError::InvalidData(format!(
432 "failed to canonicalize database path ancestor {}: {error}",
433 ancestor.display()
434 )));
435 }
436 }
437 }
438
439 Err(SqliteError::InvalidData(format!(
440 "database path has no canonicalizable ancestor: {}",
441 absolute.display()
442 )))
443}
444
445fn validate_write_admission_deadline(deadline_ms: u64) -> Result<(), SqliteError> {
450 if WRITE_ADMISSION_DEADLINE_MS_RANGE.contains(&deadline_ms) {
451 return Ok(());
452 }
453 Err(SqliteError::InvalidConfig(format!(
454 "write_admission_deadline_ms must be in [{}, {}] ms, got {deadline_ms}",
455 WRITE_ADMISSION_DEADLINE_MS_RANGE.start(),
456 WRITE_ADMISSION_DEADLINE_MS_RANGE.end()
457 )))
458}
459
460pub struct ConnectionPool {
474 writer: Arc<Mutex<Connection>>,
475 checkpoint_ownership: CheckpointOwnershipGate,
484 pooled_writer_retired: AtomicBool,
489 writer_acquisition_counters: Arc<WriterAcquisitionCounters>,
494 readers: ArrayQueue<Connection>,
495 max_readers: usize,
496 config: PoolConfig,
497 read_only_open_target: Option<PathBuf>,
505 sql_bridge_reader_slots: Arc<Semaphore>,
506 sql_bridge_writer_slots: Arc<Semaphore>,
507 writer_task: OnceLock<Option<WriterTaskHandle>>,
512 writer_task_join: Mutex<Option<tokio::task::JoinHandle<()>>>,
520 writer_task_join_stored: AtomicBool,
525 origin: TxOrigin,
531 identity_path: Option<PathBuf>,
537 #[cfg(test)]
543 writer_task_spawn_count: std::sync::atomic::AtomicUsize,
544}
545
546enum ReaderLease<'pool> {
547 Pooled(Connection),
548 Shared(parking_lot::MutexGuard<'pool, Connection>),
549}
550
551pub struct ReaderGuard<'pool> {
554 lease: Option<ReaderLease<'pool>>,
555 pool: &'pool ConnectionPool,
556 reusable: bool,
557}
558
559impl<'pool> ReaderGuard<'pool> {
560 pub fn conn(&self) -> &Connection {
562 match self
563 .lease
564 .as_ref()
565 .expect("reader guard missing connection")
566 {
567 ReaderLease::Pooled(conn) => conn,
568 ReaderLease::Shared(guard) => guard,
569 }
570 }
571
572 pub(crate) fn discard(&mut self) {
576 self.reusable = false;
577 }
578}
579
580impl<'pool> Deref for ReaderGuard<'pool> {
581 type Target = Connection;
582
583 fn deref(&self) -> &Self::Target {
584 self.conn()
585 }
586}
587
588impl<'pool> Drop for ReaderGuard<'pool> {
589 fn drop(&mut self) {
590 let Some(lease) = self.lease.take() else {
591 return;
592 };
593
594 match lease {
595 ReaderLease::Pooled(conn) if self.reusable => self.pool.return_reader(conn),
596 ReaderLease::Pooled(conn) => {
597 close_connection_quietly(conn);
598 if let Ok(conn) = self.pool.open_reader_connection() {
599 if let Err(conn) = self.pool.readers.push(conn) {
600 close_connection_quietly(conn);
601 }
602 }
603 }
604 ReaderLease::Shared(guard) if !self.reusable => {
605 self.pool.retire_pooled_writer(&guard);
606 }
607 ReaderLease::Shared(_guard) => {}
608 }
609 }
610}
611
612pub struct WriterGuard<'pool> {
615 guard: parking_lot::MutexGuard<'pool, Connection>,
616 origin: TxOrigin,
620}
621
622#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
631pub struct WriterAcquisitionSnapshot {
632 pub acquisitions: u64,
635 pub pooled_acquisitions: u64,
637 pub standalone_acquisitions: u64,
639 pub writer_task_acquisitions: u64,
642 pub timeouts: u64,
644 pub writer_task_begin_busy: u64,
648 pub writer_task_begin_errors: u64,
651 pub writer_task_request_failures: u64,
655 pub writer_task_side_effects_unknown: u64,
660}
661
662#[derive(Debug, Default)]
666pub(crate) struct WriterAcquisitionCounters {
667 pooled_acquisitions: AtomicU64,
668 standalone_acquisitions: AtomicU64,
669 writer_task_acquisitions: AtomicU64,
670 pooled_timeouts: AtomicU64,
671 writer_task_begin_busy: AtomicU64,
672 writer_task_begin_errors: AtomicU64,
673 writer_task_request_failures: AtomicU64,
674 writer_task_side_effects_unknown: AtomicU64,
675}
676
677impl WriterAcquisitionCounters {
678 pub(crate) fn record_writer_task_acquisition(&self) {
679 self.writer_task_acquisitions
680 .fetch_add(1, Ordering::Relaxed);
681 }
682
683 pub(crate) fn record_writer_task_begin_busy(&self) {
690 self.writer_task_begin_busy.fetch_add(1, Ordering::Relaxed);
691 }
692
693 pub(crate) fn record_writer_task_begin_error(&self) {
697 self.writer_task_begin_errors
698 .fetch_add(1, Ordering::Relaxed);
699 }
700
701 pub(crate) fn record_writer_task_request_failure(&self) {
705 self.writer_task_request_failures
706 .fetch_add(1, Ordering::Relaxed);
707 }
708
709 pub(crate) fn record_writer_task_side_effects_unknown(&self) {
714 self.writer_task_side_effects_unknown
715 .fetch_add(1, Ordering::Relaxed);
716 }
717
718 fn snapshot(&self) -> WriterAcquisitionSnapshot {
719 let pooled_acquisitions = self.pooled_acquisitions.load(Ordering::Relaxed);
720 let standalone_acquisitions = self.standalone_acquisitions.load(Ordering::Relaxed);
721 let writer_task_acquisitions = self.writer_task_acquisitions.load(Ordering::Relaxed);
722 WriterAcquisitionSnapshot {
723 acquisitions: pooled_acquisitions
724 .saturating_add(standalone_acquisitions)
725 .saturating_add(writer_task_acquisitions),
726 pooled_acquisitions,
727 standalone_acquisitions,
728 writer_task_acquisitions,
729 timeouts: self.pooled_timeouts.load(Ordering::Relaxed),
730 writer_task_begin_busy: self.writer_task_begin_busy.load(Ordering::Relaxed),
731 writer_task_begin_errors: self.writer_task_begin_errors.load(Ordering::Relaxed),
732 writer_task_request_failures: self.writer_task_request_failures.load(Ordering::Relaxed),
733 writer_task_side_effects_unknown: self
734 .writer_task_side_effects_unknown
735 .load(Ordering::Relaxed),
736 }
737 }
738}
739
740impl<'pool> WriterGuard<'pool> {
741 pub fn conn(&self) -> &Connection {
743 &self.guard
744 }
745
746 pub fn conn_mut(&mut self) -> &mut Connection {
748 &mut self.guard
749 }
750
751 pub fn transaction<F, R>(&self, f: F) -> Result<R, SqliteError>
754 where
755 F: FnOnce(&Connection) -> Result<R, SqliteError>,
756 {
757 self.guard.execute_batch("BEGIN IMMEDIATE")?;
758 let _tx_handle = khive_storage::tx_registry::register_scoped(
759 Some("writer_guard_tx".to_string()),
760 self.origin.clone(),
761 );
762
763 match f(&self.guard) {
764 Ok(result) => {
765 if let Err(err) = self.guard.execute_batch("COMMIT") {
766 let _ = self.guard.execute_batch("ROLLBACK");
767 return Err(err.into());
768 }
769 Ok(result)
770 }
771 Err(err) => {
772 let _ = self.guard.execute_batch("ROLLBACK");
773 Err(err)
774 }
775 }
776 }
777}
778
779impl<'pool> Deref for WriterGuard<'pool> {
780 type Target = Connection;
781
782 fn deref(&self) -> &Self::Target {
783 self.conn()
784 }
785}
786
787impl<'pool> DerefMut for WriterGuard<'pool> {
788 fn deref_mut(&mut self) -> &mut Self::Target {
789 self.conn_mut()
790 }
791}
792
793impl ConnectionPool {
794 pub fn new(config: PoolConfig) -> Result<Self, SqliteError> {
802 refuse_home_data_store_in_tests(&config)?;
803 validate_write_admission_deadline(config.write_admission_deadline_ms)?;
804
805 let mut config = config;
809 let inert_memory_queue_request =
810 config.path.is_none() && config.write_queue_enabled == Some(true);
811 config.write_queue_enabled =
812 Some(config.write_queue_enabled.unwrap_or(config.path.is_some()));
813 if inert_memory_queue_request {
814 tracing::warn!(
815 "write queue explicitly requested for an in-memory pool; it is inert because \
816 in-memory pools cannot host a writer task"
817 );
818 }
819
820 let (origin, identity_path) = match config.path.as_ref() {
825 Some(path) => {
826 let (identity, canonical) = mint_db_identity(path)?;
827 (TxOrigin::Database(identity), Some(canonical))
828 }
829 None => (TxOrigin::Memory, None),
830 };
831 let read_only_open_target = read_only_open_target(&config, identity_path.as_deref())?;
832 let writer = open_writer_connection(&config, read_only_open_target.as_deref())?;
833 let wal_enabled = configure_writer_connection(&writer, &config)?;
834 let max_readers = effective_reader_count(&config, wal_enabled);
835
836 let readers = ArrayQueue::new(max_readers.max(1));
837
838 let pool = Self {
839 writer: Arc::new(Mutex::new(writer)),
840 checkpoint_ownership: CheckpointOwnershipGate::new(),
841 pooled_writer_retired: AtomicBool::new(false),
842 writer_acquisition_counters: Arc::new(WriterAcquisitionCounters::default()),
843 readers,
844 max_readers,
845 config,
846 read_only_open_target,
847 sql_bridge_reader_slots: Arc::new(Semaphore::new(max_readers.max(1))),
848 sql_bridge_writer_slots: Arc::new(Semaphore::new(1)),
849 writer_task: OnceLock::new(),
850 writer_task_join: Mutex::new(None),
851 writer_task_join_stored: AtomicBool::new(false),
852 origin,
853 identity_path,
854 #[cfg(test)]
855 writer_task_spawn_count: std::sync::atomic::AtomicUsize::new(0),
856 };
857
858 for _ in 0..pool.max_readers {
859 let conn = pool.open_reader_connection()?;
860 pool.readers
861 .push(conn)
862 .expect("reader queue must have capacity during pool initialization");
863 }
864
865 if !pool.config.read_only {
870 crate::timeout_sink::init(
871 pool.canonical_path().and_then(Path::parent),
872 &crate::timeout_sink::db_label(&pool),
873 );
874 }
875
876 Ok(pool)
877 }
878
879 pub fn reader(&self) -> Result<ReaderGuard<'_>, SqliteError> {
889 self.reader_until(|| false)?.ok_or_else(|| {
890 SqliteError::InvalidData("uncancelled reader checkout stopped unexpectedly".into())
891 })
892 }
893
894 pub(crate) fn reader_until<C>(
900 &self,
901 should_stop: C,
902 ) -> Result<Option<ReaderGuard<'_>>, SqliteError>
903 where
904 C: Fn() -> bool,
905 {
906 if self.max_readers == 0 {
907 self.ensure_pooled_writer_active()?;
908 let started = Instant::now();
909 loop {
910 if should_stop() {
911 return Ok(None);
912 }
913 let remaining = self
914 .config
915 .checkout_timeout
916 .saturating_sub(started.elapsed());
917 if remaining.is_zero() {
918 return Err(pool_exhausted_error(
919 self.config.checkout_timeout,
920 self.max_readers,
921 ));
922 }
923 if let Some(guard) = self
924 .writer
925 .try_lock_for(remaining.min(Duration::from_millis(2)))
926 {
927 self.ensure_pooled_writer_active()?;
928 return Ok(Some(ReaderGuard {
929 lease: Some(ReaderLease::Shared(guard)),
930 pool: self,
931 reusable: true,
932 }));
933 }
934 }
935 }
936
937 let started = Instant::now();
938 let mut attempt = 0u32;
939
940 loop {
941 if should_stop() {
942 return Ok(None);
943 }
944 if let Some(conn) = self.readers.pop() {
945 return Ok(Some(ReaderGuard {
946 lease: Some(ReaderLease::Pooled(conn)),
947 pool: self,
948 reusable: true,
949 }));
950 }
951
952 if started.elapsed() >= self.config.checkout_timeout {
953 return Err(pool_exhausted_error(
954 self.config.checkout_timeout,
955 self.max_readers,
956 ));
957 }
958
959 match attempt {
960 0..=7 => {
961 let spins = 1usize << attempt;
962 for _ in 0..spins {
963 std::hint::spin_loop();
964 }
965 }
966 8..=15 => thread::yield_now(),
967 _ => {
968 let remaining = self
969 .config
970 .checkout_timeout
971 .saturating_sub(started.elapsed());
972 let sleep = Duration::from_micros(50 * (1u64 << (attempt - 16).min(6)));
973 thread::sleep(sleep.min(remaining).min(Duration::from_millis(2)));
974 }
975 }
976
977 attempt = attempt.saturating_add(1);
978 }
979 }
980
981 pub fn writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
987 self.ensure_pooled_writer_active()?;
988 let Some(guard) = self.writer.try_lock_for(self.config.checkout_timeout) else {
989 self.writer_acquisition_counters
990 .pooled_timeouts
991 .fetch_add(1, Ordering::Relaxed);
992 let message = format!(
993 "timed out after {:?} waiting for sqlite writer connection",
994 self.config.checkout_timeout
995 );
996 crate::timeout_sink::emit_timeout(
997 &crate::timeout_sink::db_label(self),
998 crate::timeout_sink::Site::PoolAdmission,
999 &message,
1000 Some(
1001 self.config
1002 .checkout_timeout
1003 .as_millis()
1004 .min(u128::from(u64::MAX)) as u64,
1005 ),
1006 );
1007 return Err(SqliteError::WriterPoolCheckoutTimeout {
1008 timeout: self.config.checkout_timeout,
1009 });
1010 };
1011 self.ensure_pooled_writer_active()?;
1012 self.writer_acquisition_counters
1013 .pooled_acquisitions
1014 .fetch_add(1, Ordering::Relaxed);
1015 Ok(WriterGuard {
1016 guard,
1017 origin: self.origin(),
1018 })
1019 }
1020
1021 pub fn try_writer(&self) -> Result<WriterGuard<'_>, SqliteError> {
1026 self.writer()
1027 }
1028
1029 pub fn try_writer_nowait(&self) -> Result<WriterGuard<'_>, SqliteError> {
1038 self.ensure_pooled_writer_active()?;
1039 let guard = self.writer.try_lock().ok_or_else(|| {
1040 SqliteError::InvalidData(
1041 "writer connection busy (checkpoint skipped this tick)".to_string(),
1042 )
1043 })?;
1044 self.ensure_pooled_writer_active()?;
1045 Ok(WriterGuard {
1046 guard,
1047 origin: self.origin(),
1048 })
1049 }
1050
1051 pub(crate) fn retire_pooled_writer(&self, conn: &Connection) {
1052 self.pooled_writer_retired.store(true, Ordering::Release);
1053 if let Err(error) = conn.authorizer(Some(deny_retired_writer)) {
1054 tracing::error!(
1055 %error,
1056 "failed to install the retired pooled-writer quarantine authorizer"
1057 );
1058 }
1059 }
1060
1061 fn ensure_pooled_writer_active(&self) -> Result<(), SqliteError> {
1062 if self.pooled_writer_retired.load(Ordering::Acquire) {
1063 return Err(SqliteError::InvalidData(
1064 "pooled writer connection retired after a terminal transaction fault".to_string(),
1065 ));
1066 }
1067 Ok(())
1068 }
1069
1070 pub fn writer_acquisition_snapshot(&self) -> WriterAcquisitionSnapshot {
1073 self.writer_acquisition_counters.snapshot()
1074 }
1075
1076 pub(crate) fn writer_acquisition_counters(&self) -> Arc<WriterAcquisitionCounters> {
1078 Arc::clone(&self.writer_acquisition_counters)
1079 }
1080
1081 pub fn available_readers(&self) -> usize {
1083 self.readers.len()
1084 }
1085
1086 pub fn max_readers(&self) -> usize {
1088 self.max_readers
1089 }
1090
1091 pub fn config(&self) -> &PoolConfig {
1093 &self.config
1094 }
1095
1096 pub(crate) fn reader_admission_timeout(&self, operation: &'static str) -> StorageError {
1100 StorageError::AdmissionTimeout {
1101 operation: operation.into(),
1102 timeout_ms: u64::try_from(self.config.checkout_timeout.as_millis()).unwrap_or(u64::MAX),
1103 }
1104 }
1105
1106 pub(crate) fn resolve_reader_checkout<'p>(
1125 &self,
1126 capability: StorageCapability,
1127 operation: &'static str,
1128 outcome: Result<Option<ReaderGuard<'p>>, SqliteError>,
1129 ) -> Result<ReaderGuard<'p>, StorageError> {
1130 match outcome {
1131 Ok(Some(guard)) => Ok(guard),
1132 Ok(None) => Err(StorageError::Timeout {
1133 operation: operation.into(),
1134 }),
1135 Err(error) => {
1136 let is_pool_exhausted = matches!(
1137 &error,
1138 SqliteError::Rusqlite(rusqlite::Error::SqliteFailure(code, _))
1139 if code.code == rusqlite::ErrorCode::DatabaseBusy
1140 );
1141 if is_pool_exhausted {
1142 Err(self.reader_admission_timeout(operation))
1143 } else {
1144 Err(StorageError::driver(capability, operation, error))
1145 }
1146 }
1147 }
1148 }
1149
1150 pub(crate) fn sql_bridge_reader_slots(&self) -> Arc<Semaphore> {
1152 Arc::clone(&self.sql_bridge_reader_slots)
1153 }
1154
1155 pub(crate) fn sql_bridge_writer_slots(&self) -> Arc<Semaphore> {
1157 Arc::clone(&self.sql_bridge_writer_slots)
1158 }
1159
1160 pub fn origin(&self) -> TxOrigin {
1166 self.origin.clone()
1167 }
1168
1169 pub fn canonical_path(&self) -> Option<&Path> {
1175 self.identity_path.as_deref()
1176 }
1177
1178 pub fn write_queue_active(&self) -> bool {
1190 debug_assert!(
1191 self.config.write_queue_enabled.is_some(),
1192 "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
1193 before any write_queue_active read"
1194 );
1195 self.config.write_queue_enabled.unwrap_or(false) && self.config.path.is_some()
1196 }
1197
1198 pub fn writer_task_join_was_stored(&self) -> bool {
1204 self.writer_task_join_stored.load(Ordering::SeqCst)
1205 }
1206
1207 pub fn writer_task_handle(&self) -> Result<Option<WriterTaskHandle>, StorageError> {
1231 debug_assert!(
1239 self.config.write_queue_enabled.is_some(),
1240 "write_queue_enabled must be resolved to Some(..) by ConnectionPool::new \
1241 before any writer_task_handle read"
1242 );
1243 if !self.config.write_queue_enabled.unwrap_or(false) {
1244 return Ok(None);
1245 }
1246 if let Some(existing) = self.writer_task.get() {
1249 return Ok(existing.clone());
1250 }
1251 if tokio::runtime::Handle::try_current().is_err() {
1255 return Err(StorageError::WriterTaskNoRuntime);
1256 }
1257 Ok(self
1258 .writer_task
1259 .get_or_init(|| {
1260 #[cfg(test)]
1261 self.writer_task_spawn_count
1262 .fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1263
1264 match crate::writer_task::spawn(self, self.config.write_queue_capacity) {
1265 Ok(handle) => Some(handle),
1266 Err(e) => {
1267 tracing::warn!(
1268 error = %e,
1269 "KHIVE_WRITE_QUEUE=1 but the writer task failed to spawn; \
1270 writes fall back to the pool-mutex path"
1271 );
1272 None
1273 }
1274 }
1275 })
1276 .clone())
1277 }
1278
1279 pub(crate) fn writer_task_for_write(
1289 &self,
1290 cached: Option<&WriterTaskHandle>,
1291 operation: &'static str,
1292 ) -> Result<Option<WriterTaskHandle>, StorageError> {
1293 let handle = match cached {
1294 Some(handle) => Some(handle.clone()),
1295 None => match self.writer_task_handle() {
1296 Ok(handle) => handle,
1297 Err(error) if self.config.write_routing_strict => return Err(error),
1298 Err(_) => None,
1299 },
1300 };
1301
1302 if handle.is_none() && self.config.write_routing_strict {
1303 return Err(StorageError::Pool {
1304 operation: operation.into(),
1305 message: "strict write routing requires a writer-task handle; no handle is \
1306 available, so the direct writer fallback was refused"
1307 .into(),
1308 });
1309 }
1310 Ok(handle)
1311 }
1312
1313 pub(crate) fn record_direct_route(&self, site: crate::timeout_sink::Site) {
1317 if self.write_queue_active() {
1318 crate::timeout_sink::emit_direct_route_violation(
1319 &crate::timeout_sink::db_label(self),
1320 site,
1321 );
1322 }
1323 }
1324
1325 #[cfg(test)]
1330 pub(crate) fn writer_task_spawn_count(&self) -> usize {
1331 self.writer_task_spawn_count
1332 .load(std::sync::atomic::Ordering::SeqCst)
1333 }
1334
1335 pub(crate) fn set_writer_task_join(&self, join: tokio::task::JoinHandle<()>) {
1348 let first_store = !self.writer_task_join_stored.swap(true, Ordering::SeqCst);
1353 debug_assert!(
1354 first_store,
1355 "writer task JoinHandle stored twice (even counting a taken one); \
1356 the writer_task OnceLock is supposed to make spawn at-most-once per pool"
1357 );
1358 if first_store {
1359 *self.writer_task_join.lock() = Some(join);
1360 }
1361 }
1362
1363 pub fn take_writer_task_join(&self) -> Option<tokio::task::JoinHandle<()>> {
1381 self.writer_task_join.lock().take()
1382 }
1383
1384 pub fn legacy_conn(&self) -> Arc<Mutex<Connection>> {
1389 Arc::clone(&self.writer)
1390 }
1391
1392 fn open_reader_connection(&self) -> Result<Connection, SqliteError> {
1393 let path = self.read_connection_path()?;
1394 open_reader_connection(path, &self.config)
1395 }
1396
1397 fn read_connection_path(&self) -> Result<&Path, SqliteError> {
1398 self.read_only_open_target
1399 .as_deref()
1400 .or(self.config.path.as_deref())
1401 .ok_or_else(|| {
1402 SqliteError::InvalidData(
1403 "in-memory databases do not support standalone connections".to_string(),
1404 )
1405 })
1406 }
1407
1408 pub fn open_standalone_writer(&self) -> Result<Connection, SqliteError> {
1419 let conn = self.open_standalone_writer_untracked()?;
1420 self.writer_acquisition_counters
1421 .standalone_acquisitions
1422 .fetch_add(1, Ordering::Relaxed);
1423 Ok(conn)
1424 }
1425
1426 pub(crate) fn open_standalone_writer_untracked(&self) -> Result<Connection, SqliteError> {
1436 let path = self.config.path.as_ref().ok_or_else(|| {
1437 SqliteError::InvalidData(
1438 "in-memory databases do not support standalone connections".to_string(),
1439 )
1440 })?;
1441
1442 if self.config.read_only {
1443 return Err(SqliteError::InvalidData(
1444 "database is read-only: standalone write connections are not permitted".to_string(),
1445 ));
1446 }
1447
1448 let conn = Connection::open_with_flags(
1449 path,
1450 OpenFlags::SQLITE_OPEN_READ_WRITE
1451 | OpenFlags::SQLITE_OPEN_NO_MUTEX
1452 | OpenFlags::SQLITE_OPEN_URI,
1453 )?;
1454 conn.busy_timeout(self.config.busy_timeout)?;
1455 self.checkpoint_ownership
1456 .configure_wal_autocheckpoint(&conn)?;
1457 conn.pragma_update(None, "foreign_keys", "ON")?;
1458 conn.pragma_update(None, "synchronous", "NORMAL")?;
1459
1460 let wal_enabled =
1461 self.config.wal_mode && current_journal_mode(&conn)?.eq_ignore_ascii_case("wal");
1462 if wal_enabled {
1463 conn.pragma_update(
1464 None,
1465 "journal_size_limit",
1466 self.config.journal_size_limit_bytes,
1467 )?;
1468 }
1469
1470 Ok(conn)
1471 }
1472
1473 #[cfg(test)]
1477 pub(crate) fn effective_wal_autocheckpoint_pages(&self) -> u32 {
1478 self.checkpoint_ownership.wal_autocheckpoint_pages()
1479 }
1480
1481 pub fn claim_checkpoint_ownership(&self) -> Result<(), SqliteError> {
1503 if !self.checkpoint_ownership.begin_claim() {
1504 return Ok(());
1505 }
1506 let result = (|| {
1507 if !self.config.read_only {
1508 let writer = self.writer()?;
1509 writer.conn().pragma_update(None, "wal_autocheckpoint", 0)?;
1510 }
1511 Ok(())
1512 })();
1513 self.checkpoint_ownership.finish_claim(result.is_ok());
1514 result
1515 }
1516
1517 pub async fn propagate_checkpoint_claim_to_writer_task(&self) -> Result<(), StorageError> {
1526 let Some(handle) = self.writer_task_handle()? else {
1527 return Ok(());
1528 };
1529 handle
1530 .send_top_level(|conn| {
1531 conn.pragma_update(None, "wal_autocheckpoint", 0)
1532 .map_err(|e| StorageError::Pool {
1533 operation: "claim_checkpoint_ownership".into(),
1534 message: e.to_string(),
1535 })
1536 })
1537 .await
1538 }
1539
1540 pub fn open_standalone_reader(&self) -> Result<Connection, SqliteError> {
1545 let path = self.read_connection_path()?;
1546
1547 let conn = Connection::open_with_flags(
1548 path,
1549 OpenFlags::SQLITE_OPEN_READ_ONLY
1550 | OpenFlags::SQLITE_OPEN_NO_MUTEX
1551 | OpenFlags::SQLITE_OPEN_URI,
1552 )?;
1553 configure_reader_connection(&conn, &self.config)?;
1554 conn.pragma_update(None, "synchronous", "NORMAL")?;
1555 Ok(conn)
1556 }
1557
1558 fn return_reader(&self, conn: Connection) {
1559 if self.max_readers == 0 {
1560 return;
1561 }
1562
1563 let conn = if reset_reader_connection(&conn) && reader_connection_is_healthy(&conn) {
1564 Some(conn)
1565 } else {
1566 close_connection_quietly(conn);
1567 self.open_reader_connection().ok()
1568 };
1569
1570 if let Some(conn) = conn {
1571 if let Err(conn) = self.readers.push(conn) {
1572 eprintln!(
1573 "[sqlite-pool] reader pool queue full, discarding replacement connection"
1574 );
1575 close_connection_quietly(conn);
1576 }
1577 }
1578 }
1579}
1580
1581const MAX_SYMLINK_DEPTH: u32 = 40;
1586
1587fn mint_db_identity(configured_path: &Path) -> Result<(DbIdentity, PathBuf), SqliteError> {
1620 let absolute = if configured_path.is_absolute() {
1621 configured_path.to_path_buf()
1622 } else {
1623 let cwd = std::env::current_dir().map_err(|e| {
1624 SqliteError::InvalidData(format!(
1625 "cannot mint database identity for {configured_path:?}: failed to resolve the \
1626 process current directory: {e}"
1627 ))
1628 })?;
1629 cwd.join(configured_path)
1630 };
1631
1632 if absolute.exists() {
1633 let canonical = absolute.canonicalize().map_err(|e| {
1634 SqliteError::InvalidData(format!(
1635 "cannot mint database identity: failed to canonicalize existing path \
1636 {absolute:?}: {e}"
1637 ))
1638 })?;
1639 return Ok((
1640 DbIdentity::new(canonical.clone().into_os_string()),
1641 canonical,
1642 ));
1643 }
1644
1645 let resolved_target = resolve_symlink_chain(&absolute)?;
1646 let parent = resolved_target.parent().ok_or_else(|| {
1647 SqliteError::InvalidData(format!(
1648 "cannot mint database identity for {resolved_target:?}: path has no parent \
1649 directory"
1650 ))
1651 })?;
1652 let file_name = resolved_target.file_name().ok_or_else(|| {
1653 SqliteError::InvalidData(format!(
1654 "cannot mint database identity for {resolved_target:?}: path has no file name"
1655 ))
1656 })?;
1657 let canonical_parent = parent.canonicalize().map_err(|e| {
1658 SqliteError::InvalidData(format!(
1659 "cannot mint database identity: parent directory {parent:?} of first-open path \
1660 {resolved_target:?} does not exist or is inaccessible: {e}"
1661 ))
1662 })?;
1663 let mut identity_path = canonical_parent;
1664 identity_path.push(file_name);
1665 Ok((
1666 DbIdentity::new(identity_path.clone().into_os_string()),
1667 identity_path,
1668 ))
1669}
1670
1671fn resolve_symlink_chain(path: &Path) -> Result<PathBuf, SqliteError> {
1677 let mut current = path.to_path_buf();
1678 for _ in 0..MAX_SYMLINK_DEPTH {
1679 match fs::symlink_metadata(¤t) {
1680 Ok(meta) if meta.file_type().is_symlink() => {
1681 let target = fs::read_link(¤t).map_err(|e| {
1682 SqliteError::InvalidData(format!(
1683 "cannot mint database identity: failed to read symlink {current:?}: {e}"
1684 ))
1685 })?;
1686 current = if target.is_absolute() {
1687 target
1688 } else {
1689 match current.parent() {
1690 Some(parent) => parent.join(&target),
1691 None => target,
1692 }
1693 };
1694 }
1695 _ => return Ok(current),
1696 }
1697 }
1698 Err(SqliteError::InvalidData(format!(
1699 "cannot mint database identity for {path:?}: symlink chain exceeds \
1700 {MAX_SYMLINK_DEPTH} levels"
1701 )))
1702}
1703
1704fn effective_reader_count(config: &PoolConfig, wal_enabled: bool) -> usize {
1705 if config.path.is_some() && config.read_only {
1706 config.max_readers.max(1)
1707 } else if config.path.is_some() && config.wal_mode && wal_enabled {
1708 config.max_readers
1709 } else {
1710 0
1711 }
1712}
1713
1714fn open_writer_connection(
1715 config: &PoolConfig,
1716 read_only_open_target: Option<&Path>,
1717) -> Result<Connection, SqliteError> {
1718 match config.path.as_ref() {
1719 Some(path) => {
1720 let flags = if config.read_only {
1721 writer_read_only_open_flags()
1722 } else {
1723 writer_open_flags()
1724 };
1725 let target = if config.read_only {
1726 read_only_open_target.ok_or_else(|| {
1727 SqliteError::InvalidData(
1728 "file-backed read-only pool has no canonical open target".to_string(),
1729 )
1730 })?
1731 } else {
1732 path
1733 };
1734 Connection::open_with_flags(target, flags).map_err(Into::into)
1735 }
1736 None => Connection::open_in_memory().map_err(Into::into),
1737 }
1738}
1739
1740fn read_only_open_target(
1763 config: &PoolConfig,
1764 physical_path: Option<&Path>,
1765) -> Result<Option<PathBuf>, SqliteError> {
1766 if !config.read_only {
1767 return Ok(None);
1768 }
1769 let Some(path) = physical_path else {
1770 return Ok(None);
1771 };
1772 read_only_wal_open_target_for_path(path).map(Some)
1773}
1774
1775fn read_only_wal_open_target_for_path(path: &Path) -> Result<PathBuf, SqliteError> {
1776 if !sqlite_header_uses_wal(path)? {
1777 return Ok(path.to_path_buf());
1778 }
1779
1780 let shm = sqlite_sidecar_path(path, "-shm");
1781 match fs::metadata(&shm) {
1782 Ok(metadata) if metadata.permissions().readonly() => {
1783 let wal = sqlite_sidecar_path(path, "-wal");
1784 match fs::metadata(&wal) {
1785 Ok(_) => Ok(path.to_path_buf()),
1786 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1787 Err(SqliteError::InvalidData(format!(
1788 "read-only WAL snapshot {} has a shared-memory sidecar {} but no WAL \
1789 sidecar {}; refusing the inconsistent sidecar set before SQLite open",
1790 path.display(),
1791 shm.display(),
1792 wal.display(),
1793 )))
1794 }
1795 Err(error) => Err(SqliteError::Io(error)),
1796 }
1797 }
1798 Ok(_) => Err(SqliteError::InvalidData(format!(
1799 "read-only WAL snapshot {} has a writable WAL shared-memory sidecar {}; close every \
1800 live writer and remove the transient -shm file (or make a genuinely frozen snapshot) \
1801 before inspection",
1802 path.display(),
1803 shm.display(),
1804 ))),
1805 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1806 let wal = sqlite_sidecar_path(path, "-wal");
1807 match fs::metadata(&wal) {
1808 Ok(metadata) if metadata.len() > 0 => Err(SqliteError::InvalidData(format!(
1809 "read-only WAL snapshot {} has a non-empty WAL sidecar {} but no read-only \
1810 shared-memory sidecar {}; refusing before SQLite open because immutable \
1811 mode would omit committed WAL frames and ordinary read-only mode would \
1812 create or mutate -shm; include the frozen read-only -shm beside this \
1813 snapshot, or checkpoint a writable copy before inspection",
1814 path.display(),
1815 wal.display(),
1816 shm.display(),
1817 ))),
1818 Ok(_) => sqlite_immutable_uri(path),
1819 Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
1820 sqlite_immutable_uri(path)
1821 }
1822 Err(error) => Err(SqliteError::Io(error)),
1823 }
1824 }
1825 Err(error) => Err(SqliteError::Io(error)),
1826 }
1827}
1828
1829pub(crate) fn open_read_only_snapshot_connection(path: &Path) -> Result<Connection, SqliteError> {
1830 let (_, physical_path) = mint_db_identity(path)?;
1831 let target = read_only_wal_open_target_for_path(&physical_path)?;
1832 Connection::open_with_flags(&target, reader_open_flags()).map_err(Into::into)
1833}
1834
1835fn sqlite_header_uses_wal(path: &Path) -> Result<bool, SqliteError> {
1836 let mut file = fs::File::open(path)?;
1837 let mut header = [0_u8; 20];
1838 if let Err(error) = file.read_exact(&mut header) {
1839 if error.kind() == std::io::ErrorKind::UnexpectedEof {
1840 return Ok(false);
1841 }
1842 return Err(SqliteError::Io(error));
1843 }
1844 Ok(&header[..16] == b"SQLite format 3\0" && header[18] == 2 && header[19] == 2)
1845}
1846
1847fn sqlite_sidecar_path(path: &Path, suffix: &str) -> PathBuf {
1848 let mut sidecar = path.as_os_str().to_os_string();
1849 sidecar.push(suffix);
1850 PathBuf::from(sidecar)
1851}
1852
1853fn sqlite_immutable_uri(path: &Path) -> Result<PathBuf, SqliteError> {
1854 let absolute = if path.is_absolute() {
1855 path.to_path_buf()
1856 } else {
1857 std::env::current_dir()?.join(path)
1858 };
1859 let mut uri = String::from("file:");
1860
1861 #[cfg(unix)]
1862 {
1863 use std::os::unix::ffi::OsStrExt as _;
1864 push_sqlite_uri_path(&mut uri, absolute.as_os_str().as_bytes());
1865 }
1866
1867 #[cfg(not(unix))]
1868 {
1869 let path = absolute.to_str().ok_or_else(|| {
1870 SqliteError::InvalidData(format!(
1871 "read-only WAL snapshot path is not representable as a SQLite URI: {}",
1872 absolute.display()
1873 ))
1874 })?;
1875 let normalized = path.replace('\\', "/");
1876 if cfg!(windows) && !normalized.starts_with('/') {
1877 uri.push('/');
1878 }
1879 push_sqlite_uri_path(&mut uri, normalized.as_bytes());
1880 }
1881
1882 uri.push_str("?mode=ro&immutable=1");
1883 Ok(PathBuf::from(uri))
1884}
1885
1886fn push_sqlite_uri_path(uri: &mut String, bytes: &[u8]) {
1887 const HEX: &[u8; 16] = b"0123456789ABCDEF";
1888 for &byte in bytes {
1889 if byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.' | b'~' | b'/') {
1890 uri.push(byte as char);
1891 } else {
1892 uri.push('%');
1893 uri.push(HEX[(byte >> 4) as usize] as char);
1894 uri.push(HEX[(byte & 0x0f) as usize] as char);
1895 }
1896 }
1897}
1898
1899fn open_reader_connection(path: &Path, config: &PoolConfig) -> Result<Connection, SqliteError> {
1900 let conn = Connection::open_with_flags(path, reader_open_flags())?;
1901 configure_reader_connection(&conn, config)?;
1902 Ok(conn)
1903}
1904
1905fn writer_open_flags() -> OpenFlags {
1906 OpenFlags::SQLITE_OPEN_READ_WRITE
1907 | OpenFlags::SQLITE_OPEN_CREATE
1908 | OpenFlags::SQLITE_OPEN_URI
1909 | OpenFlags::SQLITE_OPEN_NO_MUTEX
1910}
1911
1912fn writer_read_only_open_flags() -> OpenFlags {
1915 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
1916}
1917
1918fn reader_open_flags() -> OpenFlags {
1919 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_URI | OpenFlags::SQLITE_OPEN_NO_MUTEX
1920}
1921
1922fn configure_writer_connection(
1923 conn: &Connection,
1924 config: &PoolConfig,
1925) -> Result<bool, SqliteError> {
1926 if config.read_only {
1927 conn.pragma_update(None, "foreign_keys", "ON")?;
1931 conn.busy_timeout(config.busy_timeout)?;
1932 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
1933 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
1934 conn.pragma_update(None, "temp_store", "MEMORY")?;
1935 conn.pragma_update(None, "query_only", "ON")?;
1936
1937 let wal_enabled =
1938 config.wal_mode && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
1939 return Ok(wal_enabled);
1940 }
1941
1942 let wants_wal = config.path.is_some() && config.wal_mode;
1943
1944 if wants_wal {
1945 conn.pragma_update(None, "journal_mode", "WAL")?;
1946 }
1947
1948 conn.pragma_update(None, "synchronous", "NORMAL")?;
1949 conn.pragma_update(None, "foreign_keys", "ON")?;
1950 conn.busy_timeout(config.busy_timeout)?;
1951 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
1952 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
1953 conn.pragma_update(None, "temp_store", "MEMORY")?;
1954 conn.pragma_update(
1959 None,
1960 "wal_autocheckpoint",
1961 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
1962 )?;
1963
1964 let wal_enabled = wants_wal && current_journal_mode(conn)?.eq_ignore_ascii_case("wal");
1965
1966 if wal_enabled {
1967 conn.pragma_update(None, "journal_size_limit", config.journal_size_limit_bytes)?;
1968 }
1969
1970 Ok(wal_enabled)
1971}
1972
1973fn configure_reader_connection(conn: &Connection, config: &PoolConfig) -> Result<(), SqliteError> {
1974 conn.pragma_update(None, "foreign_keys", "ON")?;
1975 conn.busy_timeout(config.busy_timeout)?;
1976 conn.pragma_update(None, "cache_size", CACHE_SIZE_KIB)?;
1977 conn.pragma_update(None, "mmap_size", MMAP_SIZE_BYTES)?;
1978 conn.pragma_update(None, "temp_store", "MEMORY")?;
1979 Ok(())
1980}
1981
1982fn current_journal_mode(conn: &Connection) -> Result<String, SqliteError> {
1983 conn.pragma_query_value(None, "journal_mode", |row| row.get::<_, String>(0))
1984 .map(|mode| mode.to_ascii_lowercase())
1985 .map_err(Into::into)
1986}
1987
1988fn reset_reader_connection(conn: &Connection) -> bool {
1989 if conn.is_autocommit() {
1990 return true;
1991 }
1992
1993 match conn.execute_batch("ROLLBACK") {
1994 Ok(()) => conn.is_autocommit(),
1995 Err(rusqlite::Error::SqliteFailure(err, _)) => {
1996 if matches!(
1997 err.code,
1998 rusqlite::ErrorCode::CannotOpen
1999 | rusqlite::ErrorCode::DatabaseCorrupt
2000 | rusqlite::ErrorCode::NotADatabase
2001 | rusqlite::ErrorCode::DiskFull
2002 ) {
2003 return false;
2004 }
2005 conn.is_autocommit()
2006 }
2007 Err(_) => false,
2008 }
2009}
2010
2011fn reader_connection_is_healthy(conn: &Connection) -> bool {
2012 match conn.query_row("SELECT 1", [], |row| row.get::<_, i64>(0)) {
2013 Ok(_) => true,
2014 Err(rusqlite::Error::SqliteFailure(err, _)) => !matches!(
2015 err.code,
2016 rusqlite::ErrorCode::CannotOpen
2017 | rusqlite::ErrorCode::NotADatabase
2018 | rusqlite::ErrorCode::DatabaseCorrupt
2019 | rusqlite::ErrorCode::PermissionDenied
2020 | rusqlite::ErrorCode::SystemIoFailure
2021 ),
2022 Err(_) => true,
2023 }
2024}
2025
2026fn close_connection_quietly(conn: Connection) {
2027 match conn.close() {
2028 Ok(()) => {}
2029 Err((conn, _)) => drop(conn),
2030 }
2031}
2032
2033fn pool_exhausted_error(timeout: Duration, max_readers: usize) -> SqliteError {
2034 rusqlite::Error::SqliteFailure(
2035 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_BUSY),
2036 Some(format!(
2037 "Pool exhausted: no reader available after {timeout:?} (max_readers={max_readers})"
2038 )),
2039 )
2040 .into()
2041}
2042
2043#[cfg(test)]
2044mod tests {
2045 use super::*;
2046 use serial_test::serial;
2047
2048 struct WarningCapture {
2049 messages: Arc<std::sync::Mutex<Vec<String>>>,
2050 }
2051
2052 impl tracing::Subscriber for WarningCapture {
2053 fn enabled(&self, _: &tracing::Metadata<'_>) -> bool {
2054 true
2055 }
2056
2057 fn new_span(&self, _: &tracing::span::Attributes<'_>) -> tracing::span::Id {
2058 tracing::span::Id::from_u64(1)
2059 }
2060
2061 fn record(&self, _: &tracing::span::Id, _: &tracing::span::Record<'_>) {}
2062
2063 fn record_follows_from(&self, _: &tracing::span::Id, _: &tracing::span::Id) {}
2064
2065 fn event(&self, event: &tracing::Event<'_>) {
2066 struct Visitor(Option<String>);
2067
2068 impl tracing::field::Visit for Visitor {
2069 fn record_debug(
2070 &mut self,
2071 field: &tracing::field::Field,
2072 value: &dyn std::fmt::Debug,
2073 ) {
2074 if field.name() == "message" {
2075 self.0 = Some(format!("{value:?}"));
2076 }
2077 }
2078 }
2079
2080 let mut visitor = Visitor(None);
2081 event.record(&mut visitor);
2082 if let Some(message) = visitor.0 {
2083 self.messages.lock().unwrap().push(message);
2084 }
2085 }
2086
2087 fn enter(&self, _: &tracing::span::Id) {}
2088
2089 fn exit(&self, _: &tracing::span::Id) {}
2090 }
2091
2092 struct CwdGuard {
2097 original: PathBuf,
2098 }
2099
2100 impl CwdGuard {
2101 fn enter(dir: &Path) -> Self {
2102 let original = std::env::current_dir().unwrap();
2103 std::env::set_current_dir(dir).unwrap();
2104 Self { original }
2105 }
2106 }
2107
2108 impl Drop for CwdGuard {
2109 fn drop(&mut self) {
2110 let _ = std::env::set_current_dir(&self.original);
2111 }
2112 }
2113
2114 const POOL_ENV_VARS: [&str; 7] = [
2115 "KHIVE_BUSY_TIMEOUT_SECS",
2116 "KHIVE_CHECKOUT_TIMEOUT_SECS",
2117 "KHIVE_WAL_AUTOCHECKPOINT_PAGES",
2118 "KHIVE_JOURNAL_SIZE_LIMIT_BYTES",
2119 "KHIVE_WRITE_QUEUE",
2120 "KHIVE_WRITE_QUEUE_CAPACITY",
2121 "KHIVE_WRITE_ROUTING",
2122 ];
2123
2124 struct PoolEnvGuard {
2125 saved: Vec<(&'static str, Option<std::ffi::OsString>)>,
2126 }
2127
2128 impl PoolEnvGuard {
2129 fn capture() -> Self {
2130 Self {
2131 saved: POOL_ENV_VARS
2132 .into_iter()
2133 .map(|key| (key, std::env::var_os(key)))
2134 .collect(),
2135 }
2136 }
2137 }
2138
2139 impl Drop for PoolEnvGuard {
2140 fn drop(&mut self) {
2141 for (key, value) in &self.saved {
2142 match value {
2143 Some(value) => std::env::set_var(key, value),
2144 None => std::env::remove_var(key),
2145 }
2146 }
2147 }
2148 }
2149
2150 fn clear_pool_env() -> PoolEnvGuard {
2151 let guard = PoolEnvGuard::capture();
2152 for var in POOL_ENV_VARS {
2153 std::env::remove_var(var);
2154 }
2155 guard
2156 }
2157
2158 fn wal_autocheckpoint_pages(conn: &Connection) -> u32 {
2159 conn.pragma_query_value(None, "wal_autocheckpoint", |row| row.get(0))
2160 .expect("read PRAGMA wal_autocheckpoint")
2161 }
2162
2163 fn journal_size_limit_bytes(conn: &Connection) -> i64 {
2164 conn.pragma_query_value(None, "journal_size_limit", |row| row.get(0))
2165 .expect("read PRAGMA journal_size_limit")
2166 }
2167
2168 #[test]
2169 fn read_only_rollback_journal_pool_keeps_a_dedicated_reader() {
2170 let dir = tempfile::tempdir().unwrap();
2171 let path = dir.path().join("read_only_delete_journal.db");
2172 {
2173 let conn = Connection::open(&path).unwrap();
2174 conn.execute_batch("CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY);")
2175 .unwrap();
2176 let mode: String = conn
2177 .pragma_query_value(None, "journal_mode", |row| row.get(0))
2178 .unwrap();
2179 assert_eq!(mode.to_ascii_lowercase(), "delete");
2180 }
2181
2182 let pool = ConnectionPool::new(PoolConfig {
2183 path: Some(path),
2184 read_only: true,
2185 write_queue_enabled: Some(false),
2186 ..PoolConfig::default()
2187 })
2188 .unwrap();
2189
2190 assert!(
2191 pool.max_readers() > 0,
2192 "a read-only rollback-journal snapshot must use a genuine read-only reader, not \
2193 alias reader() onto the query-only writer slot"
2194 );
2195 let reader = pool.reader().expect("dedicated read-only reader checkout");
2196 let count: i64 = reader
2197 .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
2198 .unwrap();
2199 assert_eq!(count, 0);
2200 drop(reader);
2201 assert_eq!(
2202 pool.writer_acquisition_snapshot(),
2203 WriterAcquisitionSnapshot::default(),
2204 "constructing and reading a rollback-journal snapshot must never acquire the writer"
2205 );
2206 }
2207
2208 fn sqlite_sidecar(path: &Path, suffix: &str) -> PathBuf {
2209 let mut sidecar = path.as_os_str().to_os_string();
2210 sidecar.push(suffix);
2211 PathBuf::from(sidecar)
2212 }
2213
2214 fn directory_entries(path: &Path) -> Vec<std::ffi::OsString> {
2215 let mut entries = std::fs::read_dir(path)
2216 .unwrap()
2217 .map(|entry| entry.unwrap().file_name())
2218 .collect::<Vec<_>>();
2219 entries.sort();
2220 entries
2221 }
2222
2223 #[test]
2228 fn read_only_persistent_wal_without_shm_is_refused_without_mutation() {
2229 let dir = tempfile::tempdir().unwrap();
2230 let source = dir.path().join("wal-source.db");
2231 let snapshot = dir.path().join("snapshot ?#%.db");
2232 let source_wal = sqlite_sidecar(&source, "-wal");
2233 let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
2234 let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
2235
2236 let source_conn = Connection::open(&source).unwrap();
2237 let mode: String = source_conn
2238 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
2239 .unwrap();
2240 assert_eq!(mode.to_ascii_lowercase(), "wal");
2241 source_conn
2242 .pragma_update(None, "wal_autocheckpoint", 0)
2243 .unwrap();
2244 source_conn
2245 .execute_batch(
2246 "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
2247 INSERT INTO snapshot_row(body) VALUES ('committed-only-in-wal');",
2248 )
2249 .unwrap();
2250 assert!(source_wal.exists(), "fixture must retain a WAL sidecar");
2251
2252 std::fs::copy(&source, &snapshot).unwrap();
2253 std::fs::copy(&source_wal, &snapshot_wal).unwrap();
2254 assert!(
2255 !snapshot_shm.exists(),
2256 "fixture intentionally omits the transient shared-memory index"
2257 );
2258
2259 let main_before = std::fs::read(&snapshot).unwrap();
2260 let wal_before = std::fs::read(&snapshot_wal).unwrap();
2261 let entries_before = directory_entries(dir.path());
2262
2263 let error = match ConnectionPool::new(PoolConfig {
2264 path: Some(snapshot.clone()),
2265 read_only: true,
2266 write_queue_enabled: Some(false),
2267 ..PoolConfig::default()
2268 }) {
2269 Ok(_) => panic!("a non-empty WAL without its frozen -shm must fail closed"),
2270 Err(error) => error,
2271 };
2272 assert!(
2273 error
2274 .to_string()
2275 .contains("would omit committed WAL frames"),
2276 "diagnostic must explain why neither unsafe open mode is allowed: {error}"
2277 );
2278
2279 assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
2280 assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
2281 assert_eq!(directory_entries(dir.path()), entries_before);
2282 assert!(
2283 !snapshot_shm.exists(),
2284 "read-only admission and every reader must keep the source free of -shm"
2285 );
2286
2287 drop(source_conn);
2288 }
2289
2290 #[test]
2295 fn read_only_persistent_wal_with_read_only_shm_reads_without_mutation() {
2296 let dir = tempfile::tempdir().unwrap();
2297 let source = dir.path().join("wal-source.db");
2298 let snapshot = dir.path().join("frozen-wal-snapshot.db");
2299 let source_wal = sqlite_sidecar(&source, "-wal");
2300 let source_shm = sqlite_sidecar(&source, "-shm");
2301 let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
2302 let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
2303
2304 let source_conn = Connection::open(&source).unwrap();
2305 let mode: String = source_conn
2306 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
2307 .unwrap();
2308 assert_eq!(mode.to_ascii_lowercase(), "wal");
2309 source_conn
2310 .pragma_update(None, "wal_autocheckpoint", 0)
2311 .unwrap();
2312 source_conn
2313 .execute_batch(
2314 "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
2315 INSERT INTO snapshot_row(body) VALUES ('committed-only-in-wal');",
2316 )
2317 .unwrap();
2318 assert!(source_wal.exists() && source_shm.exists());
2319
2320 std::fs::copy(&source, &snapshot).unwrap();
2321 std::fs::copy(&source_wal, &snapshot_wal).unwrap();
2322 std::fs::copy(&source_shm, &snapshot_shm).unwrap();
2323
2324 let snapshot_paths = [&snapshot, &snapshot_wal, &snapshot_shm];
2325 let original_permissions =
2326 snapshot_paths.map(|path| std::fs::metadata(path).unwrap().permissions());
2327 for path in snapshot_paths {
2328 let mut permissions = std::fs::metadata(path).unwrap().permissions();
2329 permissions.set_readonly(true);
2330 std::fs::set_permissions(path, permissions).unwrap();
2331 }
2332
2333 let main_before = std::fs::read(&snapshot).unwrap();
2334 let wal_before = std::fs::read(&snapshot_wal).unwrap();
2335 let shm_before = std::fs::read(&snapshot_shm).unwrap();
2336 let entries_before = directory_entries(dir.path());
2337
2338 let pool = ConnectionPool::new(PoolConfig {
2339 path: Some(snapshot.clone()),
2340 read_only: true,
2341 write_queue_enabled: Some(false),
2342 ..PoolConfig::default()
2343 })
2344 .unwrap();
2345 let reader = pool.reader().unwrap();
2346 let body: String = reader
2347 .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
2348 .unwrap();
2349 assert_eq!(body, "committed-only-in-wal");
2350 drop(reader);
2351
2352 let standalone = pool.open_standalone_reader().unwrap();
2353 let count: i64 = standalone
2354 .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
2355 .unwrap();
2356 assert_eq!(count, 1);
2357 drop(standalone);
2358 drop(pool);
2359
2360 assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
2361 assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
2362 assert_eq!(std::fs::read(&snapshot_shm).unwrap(), shm_before);
2363 assert_eq!(directory_entries(dir.path()), entries_before);
2364
2365 for (path, permissions) in snapshot_paths.into_iter().zip(original_permissions) {
2366 std::fs::set_permissions(path, permissions).unwrap();
2367 }
2368 drop(source_conn);
2369 }
2370
2371 #[cfg(unix)]
2376 #[test]
2377 fn read_only_frozen_wal_symlink_reads_target_frames_without_mutation() {
2378 use std::os::unix::fs::symlink;
2379
2380 let dir = tempfile::tempdir().unwrap();
2381 let source = dir.path().join("wal-source.db");
2382 let snapshot = dir.path().join("frozen-target.db");
2383 let alias = dir.path().join("frozen-alias.db");
2384 let source_wal = sqlite_sidecar(&source, "-wal");
2385 let source_shm = sqlite_sidecar(&source, "-shm");
2386 let snapshot_wal = sqlite_sidecar(&snapshot, "-wal");
2387 let snapshot_shm = sqlite_sidecar(&snapshot, "-shm");
2388 let alias_wal = sqlite_sidecar(&alias, "-wal");
2389 let alias_shm = sqlite_sidecar(&alias, "-shm");
2390
2391 let source_conn = Connection::open(&source).unwrap();
2392 let mode: String = source_conn
2393 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
2394 .unwrap();
2395 assert_eq!(mode.to_ascii_lowercase(), "wal");
2396 source_conn
2397 .pragma_update(None, "wal_autocheckpoint", 0)
2398 .unwrap();
2399 source_conn
2400 .execute_batch(
2401 "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
2402 INSERT INTO snapshot_row(body) VALUES ('visible-through-target-wal');",
2403 )
2404 .unwrap();
2405 assert!(source_wal.exists() && source_shm.exists());
2406
2407 std::fs::copy(&source, &snapshot).unwrap();
2408 std::fs::copy(&source_wal, &snapshot_wal).unwrap();
2409 std::fs::copy(&source_shm, &snapshot_shm).unwrap();
2410 symlink(&snapshot, &alias).unwrap();
2411 assert!(!alias_wal.exists() && !alias_shm.exists());
2412
2413 let snapshot_paths = [&snapshot, &snapshot_wal, &snapshot_shm];
2414 let original_permissions =
2415 snapshot_paths.map(|path| std::fs::metadata(path).unwrap().permissions());
2416 for path in snapshot_paths {
2417 let mut permissions = std::fs::metadata(path).unwrap().permissions();
2418 permissions.set_readonly(true);
2419 std::fs::set_permissions(path, permissions).unwrap();
2420 }
2421
2422 let main_before = std::fs::read(&snapshot).unwrap();
2423 let wal_before = std::fs::read(&snapshot_wal).unwrap();
2424 let shm_before = std::fs::read(&snapshot_shm).unwrap();
2425 let entries_before = directory_entries(dir.path());
2426
2427 let pool = ConnectionPool::new(PoolConfig {
2428 path: Some(alias.clone()),
2429 read_only: true,
2430 write_queue_enabled: Some(false),
2431 ..PoolConfig::default()
2432 })
2433 .unwrap();
2434 let reader = pool.reader().unwrap();
2435 let body: String = reader
2436 .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
2437 .unwrap();
2438 assert_eq!(body, "visible-through-target-wal");
2439 drop(reader);
2440 let standalone = pool.open_standalone_reader().unwrap();
2441 let count: i64 = standalone
2442 .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
2443 .unwrap();
2444 assert_eq!(count, 1);
2445 drop(standalone);
2446 drop(pool);
2447
2448 assert_eq!(std::fs::read(&snapshot).unwrap(), main_before);
2449 assert_eq!(std::fs::read(&snapshot_wal).unwrap(), wal_before);
2450 assert_eq!(std::fs::read(&snapshot_shm).unwrap(), shm_before);
2451 assert_eq!(directory_entries(dir.path()), entries_before);
2452 assert!(!alias_wal.exists() && !alias_shm.exists());
2453
2454 for (path, permissions) in snapshot_paths.into_iter().zip(original_permissions) {
2455 std::fs::set_permissions(path, permissions).unwrap();
2456 }
2457 drop(source_conn);
2458 }
2459
2460 #[test]
2466 fn read_only_clean_wal_snapshot_is_sidecar_free() {
2467 let dir = tempfile::tempdir().unwrap();
2468 let path = dir.path().join("clean snapshot ?#%.db");
2469 {
2470 let conn = Connection::open(&path).unwrap();
2471 let mode: String = conn
2472 .pragma_update_and_check(None, "journal_mode", "WAL", |row| row.get(0))
2473 .unwrap();
2474 assert_eq!(mode.to_ascii_lowercase(), "wal");
2475 conn.execute_batch(
2476 "CREATE TABLE snapshot_row(id INTEGER PRIMARY KEY, body TEXT NOT NULL);\
2477 INSERT INTO snapshot_row(body) VALUES ('checkpointed');",
2478 )
2479 .unwrap();
2480 }
2481 let wal = sqlite_sidecar(&path, "-wal");
2482 let shm = sqlite_sidecar(&path, "-shm");
2483 assert!(!wal.exists() && !shm.exists());
2484 assert!(sqlite_header_uses_wal(&path).unwrap());
2485
2486 let original_permissions = std::fs::metadata(&path).unwrap().permissions();
2487 let mut read_only_permissions = original_permissions.clone();
2488 read_only_permissions.set_readonly(true);
2489 std::fs::set_permissions(&path, read_only_permissions).unwrap();
2490 let main_before = std::fs::read(&path).unwrap();
2491 let entries_before = directory_entries(dir.path());
2492
2493 let pool = ConnectionPool::new(PoolConfig {
2494 path: Some(path.clone()),
2495 read_only: true,
2496 write_queue_enabled: Some(false),
2497 ..PoolConfig::default()
2498 })
2499 .unwrap();
2500 let reader = pool.reader().unwrap();
2501 let body: String = reader
2502 .query_row("SELECT body FROM snapshot_row", [], |row| row.get(0))
2503 .unwrap();
2504 assert_eq!(body, "checkpointed");
2505 drop(reader);
2506 let standalone = pool.open_standalone_reader().unwrap();
2507 let count: i64 = standalone
2508 .query_row("SELECT COUNT(*) FROM snapshot_row", [], |row| row.get(0))
2509 .unwrap();
2510 assert_eq!(count, 1);
2511 drop(standalone);
2512 drop(pool);
2513
2514 assert_eq!(std::fs::read(&path).unwrap(), main_before);
2515 assert_eq!(directory_entries(dir.path()), entries_before);
2516 assert!(!wal.exists() && !shm.exists());
2517 std::fs::set_permissions(&path, original_permissions).unwrap();
2518 }
2519
2520 #[test]
2525 fn read_only_live_rollback_journal_keeps_change_detection() {
2526 let dir = tempfile::tempdir().unwrap();
2527 let path = dir.path().join("live-delete-journal.db");
2528 let writer = Connection::open(&path).unwrap();
2529 writer
2530 .execute_batch("CREATE TABLE live_row(id INTEGER PRIMARY KEY);")
2531 .unwrap();
2532
2533 let pool = ConnectionPool::new(PoolConfig {
2534 path: Some(path),
2535 read_only: true,
2536 write_queue_enabled: Some(false),
2537 ..PoolConfig::default()
2538 })
2539 .unwrap();
2540 {
2541 let reader = pool.reader().unwrap();
2542 let count: i64 = reader
2543 .query_row("SELECT COUNT(*) FROM live_row", [], |row| row.get(0))
2544 .unwrap();
2545 assert_eq!(count, 0);
2546 }
2547
2548 writer
2549 .execute("INSERT INTO live_row DEFAULT VALUES", [])
2550 .unwrap();
2551 let reader = pool.reader().unwrap();
2552 let count: i64 = reader
2553 .query_row("SELECT COUNT(*) FROM live_row", [], |row| row.get(0))
2554 .unwrap();
2555 assert_eq!(
2556 count, 1,
2557 "rollback-journal read-only connections must retain live change detection"
2558 );
2559 }
2560
2561 #[test]
2566 fn read_only_live_wal_with_writable_shm_is_refused_without_mutation() {
2567 let dir = tempfile::tempdir().unwrap();
2568 let path = dir.path().join("live-wal.db");
2569 let wal = sqlite_sidecar(&path, "-wal");
2570 let shm = sqlite_sidecar(&path, "-shm");
2571 let writer = Connection::open(&path).unwrap();
2572 writer.pragma_update(None, "journal_mode", "WAL").unwrap();
2573 writer
2574 .execute_batch(
2575 "CREATE TABLE live_row(id INTEGER PRIMARY KEY);\
2576 INSERT INTO live_row DEFAULT VALUES;",
2577 )
2578 .unwrap();
2579 assert!(wal.exists() && shm.exists());
2580
2581 let main_before = std::fs::read(&path).unwrap();
2582 let wal_before = std::fs::read(&wal).unwrap();
2583 let shm_before = std::fs::read(&shm).unwrap();
2584 let error = match ConnectionPool::new(PoolConfig {
2585 path: Some(path.clone()),
2586 read_only: true,
2587 write_queue_enabled: Some(false),
2588 ..PoolConfig::default()
2589 }) {
2590 Ok(_) => panic!("a live WAL database with writable -shm must fail closed"),
2591 Err(error) => error,
2592 };
2593 assert!(
2594 error
2595 .to_string()
2596 .contains("writable WAL shared-memory sidecar"),
2597 "diagnostic must explain how to freeze the snapshot: {error}"
2598 );
2599 assert_eq!(std::fs::read(&path).unwrap(), main_before);
2600 assert_eq!(std::fs::read(&wal).unwrap(), wal_before);
2601 assert_eq!(std::fs::read(&shm).unwrap(), shm_before);
2602
2603 drop(writer);
2604 }
2605
2606 #[cfg(unix)]
2611 #[test]
2612 fn read_only_live_wal_symlink_rejects_target_writable_shm_without_mutation() {
2613 use std::os::unix::fs::symlink;
2614
2615 let dir = tempfile::tempdir().unwrap();
2616 let target = dir.path().join("live-target.db");
2617 let alias = dir.path().join("live-alias.db");
2618 let target_wal = sqlite_sidecar(&target, "-wal");
2619 let target_shm = sqlite_sidecar(&target, "-shm");
2620 let alias_wal = sqlite_sidecar(&alias, "-wal");
2621 let alias_shm = sqlite_sidecar(&alias, "-shm");
2622
2623 let writer = Connection::open(&target).unwrap();
2624 writer.pragma_update(None, "journal_mode", "WAL").unwrap();
2625 writer
2626 .execute_batch(
2627 "CREATE TABLE live_row(id INTEGER PRIMARY KEY);\
2628 INSERT INTO live_row DEFAULT VALUES;",
2629 )
2630 .unwrap();
2631 assert!(target_wal.exists() && target_shm.exists());
2632 symlink(&target, &alias).unwrap();
2633 assert!(!alias_wal.exists() && !alias_shm.exists());
2634
2635 let main_before = std::fs::read(&target).unwrap();
2636 let wal_before = std::fs::read(&target_wal).unwrap();
2637 let shm_before = std::fs::read(&target_shm).unwrap();
2638 let entries_before = directory_entries(dir.path());
2639
2640 let error = match ConnectionPool::new(PoolConfig {
2641 path: Some(alias),
2642 read_only: true,
2643 write_queue_enabled: Some(false),
2644 ..PoolConfig::default()
2645 }) {
2646 Ok(_) => panic!("a symlink must not hide the target's writable -shm"),
2647 Err(error) => error,
2648 };
2649 assert!(
2650 error
2651 .to_string()
2652 .contains("writable WAL shared-memory sidecar"),
2653 "diagnostic must identify the canonical target's live sidecar: {error}"
2654 );
2655 assert_eq!(std::fs::read(&target).unwrap(), main_before);
2656 assert_eq!(std::fs::read(&target_wal).unwrap(), wal_before);
2657 assert_eq!(std::fs::read(&target_shm).unwrap(), shm_before);
2658 assert_eq!(directory_entries(dir.path()), entries_before);
2659 assert!(!alias_wal.exists() && !alias_shm.exists());
2660
2661 drop(writer);
2662 }
2663
2664 #[test]
2665 #[serial]
2666 fn pool_config_default_values_match_constants() {
2667 let _pool_env = clear_pool_env();
2671 let cfg = PoolConfig::default();
2672 assert_eq!(
2673 cfg.journal_size_limit_bytes,
2674 DEFAULT_JOURNAL_SIZE_LIMIT_BYTES
2675 );
2676 assert_eq!(cfg.busy_timeout, Duration::from_secs(30));
2677 assert_eq!(cfg.checkout_timeout, Duration::from_secs(5));
2678 }
2679
2680 #[test]
2681 #[serial]
2682 fn legacy_env_cannot_change_wal_autocheckpoint() {
2683 let _pool_env = clear_pool_env();
2684 std::env::set_var("KHIVE_WAL_AUTOCHECKPOINT_PAGES", "8000");
2685 let dir = tempfile::tempdir().unwrap();
2686 let path = dir.path().join("legacy_autocheckpoint_env.db");
2687 let pool = ConnectionPool::new(PoolConfig {
2688 path: Some(path),
2689 ..PoolConfig::default()
2690 })
2691 .expect("pool open");
2692 {
2693 let writer = pool.writer().expect("writer");
2694 assert_eq!(
2695 wal_autocheckpoint_pages(writer.conn()),
2696 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
2697 "the removed env override must not change the unclaimed fallback"
2698 );
2699 }
2700 pool.claim_checkpoint_ownership().expect("claim ownership");
2701 let writer = pool.writer().expect("writer after claim");
2702 assert_eq!(
2703 wal_autocheckpoint_pages(writer.conn()),
2704 0,
2705 "the removed env override must not change the claimed-owner setting"
2706 );
2707 std::env::remove_var("KHIVE_WAL_AUTOCHECKPOINT_PAGES");
2708 }
2709
2710 #[test]
2711 #[serial]
2712 fn pool_config_env_override_journal_size_limit() {
2713 std::env::set_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES", "134217728");
2714 let cfg = PoolConfig::default();
2715 std::env::remove_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES");
2716 assert_eq!(cfg.journal_size_limit_bytes, 134_217_728);
2717 }
2718
2719 #[test]
2720 #[serial]
2721 fn pool_config_env_override_busy_timeout() {
2722 std::env::set_var("KHIVE_BUSY_TIMEOUT_SECS", "60");
2723 let cfg = PoolConfig::default();
2724 std::env::remove_var("KHIVE_BUSY_TIMEOUT_SECS");
2725 assert_eq!(cfg.busy_timeout, Duration::from_secs(60));
2726 }
2727
2728 #[test]
2729 #[serial]
2730 fn pool_config_env_override_checkout_timeout() {
2731 std::env::set_var("KHIVE_CHECKOUT_TIMEOUT_SECS", "10");
2732 let cfg = PoolConfig::default();
2733 std::env::remove_var("KHIVE_CHECKOUT_TIMEOUT_SECS");
2734 assert_eq!(cfg.checkout_timeout, Duration::from_secs(10));
2735 }
2736
2737 #[test]
2738 #[serial]
2739 fn pool_config_write_queue_defaults_unset() {
2740 let _pool_env = clear_pool_env();
2741 let cfg = PoolConfig::default();
2742 assert_eq!(cfg.write_queue_enabled, None);
2743 assert_eq!(cfg.write_queue_capacity, DEFAULT_WRITE_QUEUE_CAPACITY);
2744 }
2745
2746 #[test]
2747 #[serial]
2748 fn clear_pool_env_restores_overrides_on_drop() {
2749 let _ambient_env = PoolEnvGuard::capture();
2750 std::env::set_var("KHIVE_BUSY_TIMEOUT_SECS", "73");
2751
2752 {
2753 let _pool_env = clear_pool_env();
2754 assert_eq!(std::env::var_os("KHIVE_BUSY_TIMEOUT_SECS"), None);
2755 }
2756
2757 assert_eq!(
2758 std::env::var_os("KHIVE_BUSY_TIMEOUT_SECS"),
2759 Some(std::ffi::OsString::from("73"))
2760 );
2761 }
2762
2763 #[test]
2764 #[serial]
2765 fn pool_config_env_override_write_queue_enabled() {
2766 std::env::set_var("KHIVE_WRITE_QUEUE", "1");
2767 let cfg = PoolConfig::default();
2768 std::env::remove_var("KHIVE_WRITE_QUEUE");
2769 assert_eq!(cfg.write_queue_enabled, Some(true));
2770 }
2771
2772 #[test]
2773 #[serial]
2774 fn pool_config_env_override_write_queue_enabled_accepts_true_case_insensitive() {
2775 std::env::set_var("KHIVE_WRITE_QUEUE", "True");
2776 let cfg = PoolConfig::default();
2777 std::env::remove_var("KHIVE_WRITE_QUEUE");
2778 assert_eq!(cfg.write_queue_enabled, Some(true));
2779 }
2780
2781 #[test]
2782 #[serial]
2783 fn pool_config_env_override_write_queue_enabled_accepts_zero_as_explicit_off() {
2784 std::env::set_var("KHIVE_WRITE_QUEUE", "0");
2785 let cfg = PoolConfig::default();
2786 std::env::remove_var("KHIVE_WRITE_QUEUE");
2787 assert_eq!(cfg.write_queue_enabled, Some(false));
2788 }
2789
2790 #[cfg(unix)]
2795 #[test]
2796 #[serial]
2797 fn pool_config_env_override_write_queue_non_unicode_value_is_explicit_off() {
2798 use std::os::unix::ffi::OsStrExt;
2799 let _pool_env = clear_pool_env();
2800 std::env::set_var(
2801 "KHIVE_WRITE_QUEUE",
2802 std::ffi::OsStr::from_bytes(b"\xff\xfe"),
2803 );
2804 let cfg = PoolConfig::default();
2805 assert_eq!(cfg.write_queue_enabled, Some(false));
2806 }
2807
2808 #[test]
2809 #[serial]
2810 fn pool_config_env_override_write_queue_invalid_value_is_explicit_off() {
2811 std::env::set_var("KHIVE_WRITE_QUEUE", "banana");
2815 let cfg = PoolConfig::default();
2816 std::env::remove_var("KHIVE_WRITE_QUEUE");
2817 assert_eq!(cfg.write_queue_enabled, Some(false));
2818 }
2819
2820 #[test]
2821 #[serial]
2822 fn pool_config_write_routing_strict_defaults_off() {
2823 let _pool_env = clear_pool_env();
2824 let cfg = PoolConfig::default();
2825 assert!(!cfg.write_routing_strict);
2826 }
2827
2828 #[test]
2829 #[serial]
2830 fn pool_config_env_override_write_routing_strict() {
2831 std::env::set_var("KHIVE_WRITE_ROUTING", "strict");
2832 let cfg = PoolConfig::default();
2833 std::env::remove_var("KHIVE_WRITE_ROUTING");
2834 assert!(cfg.write_routing_strict);
2835 }
2836
2837 #[test]
2838 #[serial]
2839 fn pool_config_env_override_write_routing_strict_case_insensitive() {
2840 std::env::set_var("KHIVE_WRITE_ROUTING", "STRICT");
2841 let cfg = PoolConfig::default();
2842 std::env::remove_var("KHIVE_WRITE_ROUTING");
2843 assert!(cfg.write_routing_strict);
2844 }
2845
2846 #[test]
2847 #[serial]
2848 fn pool_config_env_write_routing_ignores_unrecognized_value() {
2849 std::env::set_var("KHIVE_WRITE_ROUTING", "eventual");
2850 let cfg = PoolConfig::default();
2851 std::env::remove_var("KHIVE_WRITE_ROUTING");
2852 assert!(!cfg.write_routing_strict);
2853 }
2854
2855 #[test]
2856 #[serial]
2857 fn pool_config_env_override_write_queue_capacity() {
2858 std::env::set_var("KHIVE_WRITE_QUEUE_CAPACITY", "64");
2859 let cfg = PoolConfig::default();
2860 std::env::remove_var("KHIVE_WRITE_QUEUE_CAPACITY");
2861 assert_eq!(cfg.write_queue_capacity, 64);
2862 }
2863
2864 #[test]
2865 #[serial]
2866 fn pool_config_env_invalid_write_queue_capacity_falls_back_to_default() {
2867 std::env::set_var("KHIVE_WRITE_QUEUE_CAPACITY", "0");
2868 let cfg = PoolConfig::default();
2869 std::env::remove_var("KHIVE_WRITE_QUEUE_CAPACITY");
2870 assert_eq!(cfg.write_queue_capacity, DEFAULT_WRITE_QUEUE_CAPACITY);
2871 }
2872
2873 #[test]
2874 #[serial]
2875 fn pool_config_invalid_journal_size_limit_falls_back_to_default() {
2876 std::env::set_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES", "");
2877 let cfg = PoolConfig::default();
2878 std::env::remove_var("KHIVE_JOURNAL_SIZE_LIMIT_BYTES");
2879 assert_eq!(
2880 cfg.journal_size_limit_bytes,
2881 DEFAULT_JOURNAL_SIZE_LIMIT_BYTES
2882 );
2883 }
2884
2885 #[test]
2886 fn file_backed_pool_opens_successfully() {
2887 let dir = tempfile::tempdir().unwrap();
2888 let path = dir.path().join("test_pool.db");
2889 let cfg = PoolConfig {
2890 path: Some(path.clone()),
2891 ..PoolConfig::default()
2892 };
2893 let pool = ConnectionPool::new(cfg).expect("file-backed pool should open");
2894 assert!(path.exists());
2895 assert!(pool.max_readers() > 0);
2896 }
2897
2898 #[test]
2899 fn standalone_wal_writer_uses_configured_journal_size_limit() {
2900 let dir = tempfile::tempdir().unwrap();
2901 let path = dir.path().join("standalone_wal_journal_limit.db");
2902 let configured_limit = 12_345_678;
2903 let pool = ConnectionPool::new(PoolConfig {
2904 path: Some(path),
2905 journal_size_limit_bytes: configured_limit,
2906 write_queue_enabled: Some(false),
2907 ..PoolConfig::default()
2908 })
2909 .expect("WAL pool open");
2910
2911 let standalone = pool
2912 .open_standalone_writer_untracked()
2913 .expect("standalone WAL writer open");
2914 assert_eq!(current_journal_mode(&standalone).unwrap(), "wal");
2915 assert_eq!(journal_size_limit_bytes(&standalone), configured_limit);
2916 }
2917
2918 #[test]
2919 fn standalone_rollback_writer_keeps_sqlite_journal_size_limit() {
2920 let dir = tempfile::tempdir().unwrap();
2921 let path = dir.path().join("standalone_rollback_journal_limit.db");
2922 let sqlite_default = {
2923 let conn = Connection::open(&path).expect("seed rollback-journal database");
2924 assert_eq!(current_journal_mode(&conn).unwrap(), "delete");
2925 journal_size_limit_bytes(&conn)
2926 };
2927 let configured_limit = if sqlite_default == 12_345_678 {
2928 23_456_789
2929 } else {
2930 12_345_678
2931 };
2932 let pool = ConnectionPool::new(PoolConfig {
2933 path: Some(path),
2934 wal_mode: false,
2935 journal_size_limit_bytes: configured_limit,
2936 write_queue_enabled: Some(false),
2937 ..PoolConfig::default()
2938 })
2939 .expect("rollback-journal pool open");
2940
2941 let standalone = pool
2942 .open_standalone_writer_untracked()
2943 .expect("standalone rollback-journal writer open");
2944 assert_eq!(current_journal_mode(&standalone).unwrap(), "delete");
2945 assert_eq!(journal_size_limit_bytes(&standalone), sqlite_default);
2946 }
2947
2948 #[test]
2949 fn writer_connections_follow_checkpoint_ownership_claim() {
2950 let dir = tempfile::tempdir().unwrap();
2951 let path = dir.path().join("writer_autocheckpoint.db");
2952 let pool = ConnectionPool::new(PoolConfig {
2953 path: Some(path),
2954 write_queue_enabled: Some(false),
2955 ..PoolConfig::default()
2956 })
2957 .expect("pool open");
2958
2959 {
2963 let writer = pool.writer().expect("pooled writer");
2964 assert_eq!(
2965 wal_autocheckpoint_pages(writer.conn()),
2966 FALLBACK_WAL_AUTOCHECKPOINT_PAGES
2967 );
2968 }
2969 let standalone = pool
2970 .open_standalone_writer()
2971 .expect("standalone writer opened before any claim");
2972 assert_eq!(
2973 wal_autocheckpoint_pages(&standalone),
2974 FALLBACK_WAL_AUTOCHECKPOINT_PAGES
2975 );
2976 drop(standalone);
2977
2978 pool.claim_checkpoint_ownership().expect("claim ownership");
2982 {
2983 let writer = pool.writer().expect("pooled writer after claim");
2984 assert_eq!(wal_autocheckpoint_pages(writer.conn()), 0);
2985 }
2986 let claimed_standalone = pool
2987 .open_standalone_writer()
2988 .expect("standalone writer opened after the claim");
2989 assert_eq!(wal_autocheckpoint_pages(&claimed_standalone), 0);
2990 drop(claimed_standalone);
2991
2992 let later_infrastructure = pool
2993 .open_standalone_writer_untracked()
2994 .expect("later infrastructure writer");
2995 assert_eq!(wal_autocheckpoint_pages(&later_infrastructure), 0);
2996
2997 let memory_pool = ConnectionPool::new(PoolConfig {
2998 write_queue_enabled: Some(false),
2999 ..PoolConfig::default()
3000 })
3001 .expect("in-memory pool open");
3002 let memory_writer = memory_pool.writer().expect("in-memory writer");
3003 assert_eq!(
3004 wal_autocheckpoint_pages(memory_writer.conn()),
3005 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3006 "an unclaimed in-memory pool keeps the bounded fallback"
3007 );
3008 }
3009
3010 #[test]
3011 fn standalone_writer_waits_for_checkpoint_claim_resolution() {
3012 let dir = tempfile::tempdir().unwrap();
3013 let path = dir.path().join("checkpoint_claim_race.db");
3014 let pool = Arc::new(
3015 ConnectionPool::new(PoolConfig {
3016 path: Some(path),
3017 checkout_timeout: Duration::from_secs(5),
3018 write_queue_enabled: Some(false),
3019 ..PoolConfig::default()
3020 })
3021 .expect("pool open"),
3022 );
3023
3024 let legacy_conn = pool.legacy_conn();
3025 let held_writer = legacy_conn.lock();
3026 let claim_start = Arc::new(std::sync::Barrier::new(2));
3027 let claim_pool = Arc::clone(&pool);
3028 let claim_thread_start = Arc::clone(&claim_start);
3029 let claim_thread = thread::spawn(move || {
3030 claim_thread_start.wait();
3031 claim_pool.claim_checkpoint_ownership()
3032 });
3033 claim_start.wait();
3034
3035 {
3036 let mut state = pool.checkpoint_ownership.state.lock();
3037 while state.phase != CheckpointOwnership::Claiming {
3038 pool.checkpoint_ownership.changed.wait(&mut state);
3039 }
3040 }
3041
3042 let open_start = Arc::new(std::sync::Barrier::new(2));
3043 let open_pool = Arc::clone(&pool);
3044 let open_thread_start = Arc::clone(&open_start);
3045 let open_thread = thread::spawn(move || {
3046 open_thread_start.wait();
3047 let conn = open_pool
3048 .open_standalone_writer()
3049 .expect("standalone writer after claim resolution");
3050 wal_autocheckpoint_pages(&conn)
3051 });
3052 open_start.wait();
3053
3054 {
3055 let mut state = pool.checkpoint_ownership.state.lock();
3056 while state.connection_waiters == 0 {
3057 pool.checkpoint_ownership.changed.wait(&mut state);
3058 }
3059 assert_eq!(state.phase, CheckpointOwnership::Claiming);
3060 }
3061
3062 drop(held_writer);
3063 claim_thread
3064 .join()
3065 .expect("claim thread joins")
3066 .expect("claim succeeds");
3067 assert_eq!(
3068 open_thread.join().expect("standalone-open thread joins"),
3069 0,
3070 "a writer open concurrent with a successful claim must inherit claimed ownership"
3071 );
3072 }
3073
3074 #[test]
3075 fn standalone_fallback_application_linearizes_before_claim_publication() {
3076 let dir = tempfile::tempdir().unwrap();
3077 let path = dir.path().join("checkpoint_open_before_claim.db");
3078 let pool = Arc::new(
3079 ConnectionPool::new(PoolConfig {
3080 path: Some(path),
3081 checkout_timeout: Duration::from_secs(5),
3082 write_queue_enabled: Some(false),
3083 ..PoolConfig::default()
3084 })
3085 .expect("pool open"),
3086 );
3087 let pause = Arc::new(CheckpointConnectionConfigPause::new());
3088 *pool.checkpoint_ownership.connection_config_pause.lock() = Some(Arc::clone(&pause));
3089
3090 let open_pool = Arc::clone(&pool);
3091 let open_thread = thread::spawn(move || {
3092 let conn = open_pool
3093 .open_standalone_writer()
3094 .expect("standalone writer opens");
3095 wal_autocheckpoint_pages(&conn)
3096 });
3097 pause.selected.wait();
3098 assert!(
3099 pool.checkpoint_ownership.state.try_lock().is_none(),
3100 "standalone selection must retain the ownership gate until its PRAGMA is applied"
3101 );
3102
3103 let (claim_observed_tx, claim_observed_rx) = std::sync::mpsc::sync_channel(0);
3104 *pool.checkpoint_ownership.claim_lock_observed.lock() = Some(claim_observed_tx);
3105 let claim_pool = Arc::clone(&pool);
3106 let claim_thread = thread::spawn(move || claim_pool.claim_checkpoint_ownership());
3107 assert!(
3108 claim_observed_rx
3109 .recv()
3110 .expect("claim reports whether it observed gate contention"),
3111 "the claim must attempt the gate between fallback selection and PRAGMA application"
3112 );
3113 pause.resume.wait();
3114
3115 assert_eq!(
3116 open_thread.join().expect("standalone-open thread joins"),
3117 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3118 "an open linearized before the claim keeps the fallback"
3119 );
3120 claim_thread
3121 .join()
3122 .expect("claim thread joins")
3123 .expect("claim succeeds after standalone configuration");
3124 assert_eq!(pool.effective_wal_autocheckpoint_pages(), 0);
3125 }
3126
3127 #[test]
3128 fn failed_checkpoint_ownership_claim_keeps_fallback_and_can_be_retried() {
3129 let dir = tempfile::tempdir().unwrap();
3130 let path = dir.path().join("checkpoint_claim_retry.db");
3131 let pool = ConnectionPool::new(PoolConfig {
3132 path: Some(path),
3133 checkout_timeout: Duration::from_millis(1),
3134 write_queue_enabled: Some(false),
3135 ..PoolConfig::default()
3136 })
3137 .expect("pool open");
3138
3139 let legacy_conn = pool.legacy_conn();
3140 let held_writer = legacy_conn.lock();
3141 let error = pool
3142 .claim_checkpoint_ownership()
3143 .expect_err("the held pooled writer must make the claim time out");
3144 assert!(matches!(
3145 error,
3146 SqliteError::WriterPoolCheckoutTimeout { .. }
3147 ));
3148 assert_eq!(
3149 pool.effective_wal_autocheckpoint_pages(),
3150 FALLBACK_WAL_AUTOCHECKPOINT_PAGES,
3151 "a failed claim must leave later writer connections fallback-safe"
3152 );
3153
3154 let fallback_writer = pool
3155 .open_standalone_writer()
3156 .expect("standalone writer after failed claim");
3157 assert_eq!(
3158 wal_autocheckpoint_pages(&fallback_writer),
3159 FALLBACK_WAL_AUTOCHECKPOINT_PAGES
3160 );
3161 drop(fallback_writer);
3162
3163 drop(held_writer);
3164 pool.claim_checkpoint_ownership()
3165 .expect("the ownership claim remains retryable");
3166 assert_eq!(pool.effective_wal_autocheckpoint_pages(), 0);
3167 let writer = pool.writer().expect("pooled writer after successful retry");
3168 assert_eq!(wal_autocheckpoint_pages(writer.conn()), 0);
3169 }
3170
3171 #[test]
3172 fn threshold_crossing_commits_do_not_run_an_implicit_checkpoint_once_claimed() {
3173 const FORMER_AUTOCHECKPOINT_THRESHOLD_PAGES: i64 = FALLBACK_WAL_AUTOCHECKPOINT_PAGES as i64;
3174
3175 let dir = tempfile::tempdir().unwrap();
3176 let path = dir.path().join("no_implicit_checkpoint.db");
3177 let pool = ConnectionPool::new(PoolConfig {
3178 path: Some(path),
3179 write_queue_enabled: Some(false),
3180 ..PoolConfig::default()
3181 })
3182 .expect("pool open");
3183 pool.claim_checkpoint_ownership()
3184 .expect("claim ownership for the dedicated-owner posture");
3185 let writer = pool.writer().expect("pooled writer");
3186 writer
3187 .execute_batch("CREATE TABLE blobs (value BLOB NOT NULL)")
3188 .expect("create fixture table");
3189
3190 let page_size: i64 = writer
3191 .pragma_query_value(None, "page_size", |row| row.get(0))
3192 .expect("read page size");
3193 let payload_bytes = page_size * 32;
3194 for _ in 0..160 {
3195 writer
3196 .execute(
3197 "INSERT INTO blobs (value) VALUES (zeroblob(?1))",
3198 [payload_bytes],
3199 )
3200 .expect("autocommit fixture row");
3201 }
3202
3203 let log_frames: i64 = writer
3204 .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| row.get(1))
3205 .expect("observe WAL frame count");
3206 assert!(
3207 log_frames > FORMER_AUTOCHECKPOINT_THRESHOLD_PAGES,
3208 "the commit sequence must retain more than the former automatic threshold; \
3209 observed {log_frames} frames"
3210 );
3211 }
3212
3213 #[test]
3221 fn unclaimed_pool_retains_bounded_autocheckpoint_reclamation() {
3222 let dir = tempfile::tempdir().unwrap();
3223 let path = dir.path().join("bounded_fallback_reclamation.db");
3224 let pool = ConnectionPool::new(PoolConfig {
3225 path: Some(path),
3226 write_queue_enabled: Some(false),
3227 ..PoolConfig::default()
3228 })
3229 .expect("pool open");
3230 let writer = pool.writer().expect("pooled writer");
3231 writer
3232 .execute_batch("CREATE TABLE blobs (value BLOB NOT NULL)")
3233 .expect("create fixture table");
3234
3235 let page_size: i64 = writer
3236 .pragma_query_value(None, "page_size", |row| row.get(0))
3237 .expect("read page size");
3238 let payload_bytes = page_size * 32;
3239 for _ in 0..160 {
3240 writer
3241 .execute(
3242 "INSERT INTO blobs (value) VALUES (zeroblob(?1))",
3243 [payload_bytes],
3244 )
3245 .expect("autocommit fixture row");
3246 }
3247
3248 let log_frames: i64 = writer
3254 .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| row.get(1))
3255 .expect("observe WAL frame count");
3256 assert!(
3257 log_frames < FALLBACK_WAL_AUTOCHECKPOINT_PAGES as i64,
3258 "an unclaimed pool must reclaim WAL frames via the bounded autocheckpoint; \
3259 observed {log_frames} retained frames"
3260 );
3261 }
3262
3263 #[tokio::test]
3264 #[serial]
3265 async fn unset_write_queue_resolves_on_for_file_backed_pool() {
3266 let _pool_env = clear_pool_env();
3267 let dir = tempfile::tempdir().unwrap();
3268 let path = dir.path().join("unset_file_backed.db");
3269 let pool = ConnectionPool::new(PoolConfig {
3270 path: Some(path),
3271 write_queue_enabled: None,
3272 ..PoolConfig::default()
3273 })
3274 .expect("file-backed pool should open");
3275 assert_eq!(pool.config().write_queue_enabled, Some(true));
3276 assert!(
3279 pool.writer_task_handle()
3280 .expect("spawn inside a runtime context must not error")
3281 .is_some(),
3282 "resolved-on file-backed pool must actually spawn the writer task"
3283 );
3284 }
3285
3286 #[tokio::test]
3287 #[serial]
3288 async fn unset_write_queue_resolves_off_for_memory_backed_pool() {
3289 let _pool_env = clear_pool_env();
3290 let pool = ConnectionPool::new(PoolConfig {
3291 path: None,
3292 write_queue_enabled: None,
3293 ..PoolConfig::default()
3294 })
3295 .expect("in-memory pool should open");
3296 assert_eq!(pool.config().write_queue_enabled, Some(false));
3297 assert!(
3300 pool.writer_task_handle()
3301 .expect("disabled queue must resolve without error")
3302 .is_none(),
3303 "resolved-off in-memory pool must not spawn a writer task"
3304 );
3305 }
3306
3307 #[test]
3308 #[serial]
3309 fn explicit_false_stays_off_for_file_backed_pool() {
3310 let _pool_env = clear_pool_env();
3311 let dir = tempfile::tempdir().unwrap();
3312 let path = dir.path().join("explicit_false_file_backed.db");
3313 let pool = ConnectionPool::new(PoolConfig {
3314 path: Some(path),
3315 write_queue_enabled: Some(false),
3316 ..PoolConfig::default()
3317 })
3318 .expect("file-backed pool should open");
3319 assert_eq!(pool.config().write_queue_enabled, Some(false));
3320 }
3321
3322 #[tokio::test]
3323 #[serial]
3324 async fn explicit_true_stays_on_for_memory_backed_pool() {
3325 let _pool_env = clear_pool_env();
3326 let pool = ConnectionPool::new(PoolConfig {
3327 path: None,
3328 write_queue_enabled: Some(true),
3329 ..PoolConfig::default()
3330 })
3331 .expect("in-memory pool should open");
3332 assert_eq!(pool.config().write_queue_enabled, Some(true));
3333 assert!(
3339 pool.writer_task_handle()
3340 .expect("spawn degrade must resolve without error")
3341 .is_none(),
3342 "explicit-on in-memory pool must degrade to no writer task"
3343 );
3344 assert_eq!(
3345 pool.writer_task_spawn_count(),
3346 1,
3347 "the spawn attempt must happen exactly once and degrade, not retry"
3348 );
3349 assert!(
3350 pool.take_writer_task_join().is_none(),
3351 "a degraded spawn stores no JoinHandle to drain"
3352 );
3353 }
3354
3355 #[test]
3356 #[serial]
3357 fn explicit_true_on_memory_pool_warns_but_false_and_none_do_not() {
3358 let _pool_env = clear_pool_env();
3359 let messages = Arc::new(std::sync::Mutex::new(Vec::new()));
3360 let subscriber = WarningCapture {
3361 messages: Arc::clone(&messages),
3362 };
3363
3364 tracing::subscriber::with_default(subscriber, || {
3365 let _explicit_true = ConnectionPool::new(PoolConfig {
3366 path: None,
3367 write_queue_enabled: Some(true),
3368 ..PoolConfig::default()
3369 })
3370 .expect("in-memory pool should open");
3371 let _explicit_false = ConnectionPool::new(PoolConfig {
3372 path: None,
3373 write_queue_enabled: Some(false),
3374 ..PoolConfig::default()
3375 })
3376 .expect("in-memory pool should open");
3377 let _unset = ConnectionPool::new(PoolConfig {
3378 path: None,
3379 write_queue_enabled: None,
3380 ..PoolConfig::default()
3381 })
3382 .expect("in-memory pool should open");
3383 });
3384
3385 let messages = messages.lock().unwrap();
3386 assert_eq!(
3387 messages
3388 .iter()
3389 .filter(|message| message.contains("write queue explicitly requested"))
3390 .count(),
3391 1,
3392 "only an explicit in-memory queue request should warn: {messages:?}"
3393 );
3394 let warning = messages
3395 .iter()
3396 .find(|message| message.contains("write queue explicitly requested"))
3397 .expect("explicit in-memory queue warning should be captured");
3398 assert!(
3399 warning.contains("in-memory pools cannot host a writer task"),
3400 "warning must explain why the request is inert: {messages:?}"
3401 );
3402 }
3403
3404 #[test]
3405 fn standalone_writer_open_counts_its_connection_class_once() {
3406 let dir = tempfile::tempdir().unwrap();
3407 let path = dir.path().join("standalone_writer_counter.db");
3408 let pool = ConnectionPool::new(PoolConfig {
3409 path: Some(path),
3410 ..PoolConfig::default()
3411 })
3412 .expect("file-backed pool");
3413
3414 let _standalone = pool
3415 .open_standalone_writer()
3416 .expect("standalone writer opens");
3417
3418 assert_eq!(
3419 pool.writer_acquisition_snapshot(),
3420 WriterAcquisitionSnapshot {
3421 acquisitions: 1,
3422 pooled_acquisitions: 0,
3423 standalone_acquisitions: 1,
3424 writer_task_acquisitions: 0,
3425 timeouts: 0,
3426 writer_task_begin_busy: 0,
3427 writer_task_begin_errors: 0,
3428 writer_task_request_failures: 0,
3429 writer_task_side_effects_unknown: 0,
3430 },
3431 "the public standalone boundary must contribute to the aggregate exactly once"
3432 );
3433 }
3434
3435 #[test]
3436 fn in_memory_pool_degrades_to_single_connection() {
3437 let cfg = PoolConfig {
3438 path: None,
3439 ..PoolConfig::default()
3440 };
3441 let pool = ConnectionPool::new(cfg).expect("in-memory pool should open");
3442 assert_eq!(pool.max_readers(), 0);
3443 }
3444
3445 #[test]
3446 fn writer_checkout_and_release_works() {
3447 let cfg = PoolConfig {
3448 path: None,
3449 ..PoolConfig::default()
3450 };
3451 let pool = ConnectionPool::new(cfg).unwrap();
3452 {
3453 let _writer = pool.writer().expect("writer checkout should succeed");
3454 }
3455 let _writer2 = pool
3457 .writer()
3458 .expect("second writer checkout should succeed");
3459 }
3460
3461 #[test]
3462 fn writer_checkout_snapshot_counts_successes_and_timeouts_at_the_pool_boundary() {
3463 let cfg = PoolConfig {
3464 path: None,
3465 checkout_timeout: Duration::from_millis(1),
3466 ..PoolConfig::default()
3467 };
3468 let pool = ConnectionPool::new(cfg).unwrap();
3469
3470 assert_eq!(
3471 pool.writer_acquisition_snapshot(),
3472 WriterAcquisitionSnapshot::default()
3473 );
3474
3475 let held = pool.writer().expect("first checkout succeeds");
3476 let error = match pool.writer() {
3477 Ok(_) => panic!("the held pool mutex must force a finite-wait timeout"),
3478 Err(error) => error,
3479 };
3480 assert!(
3481 matches!(
3482 &error,
3483 SqliteError::WriterPoolCheckoutTimeout { timeout }
3484 if *timeout == Duration::from_millis(1)
3485 ),
3486 "timeout must have a stable, structurally matchable stage: {error}"
3487 );
3488 assert_eq!(
3489 pool.writer_acquisition_snapshot(),
3490 WriterAcquisitionSnapshot {
3491 acquisitions: 1,
3492 pooled_acquisitions: 1,
3493 standalone_acquisitions: 0,
3494 writer_task_acquisitions: 0,
3495 timeouts: 1,
3496 writer_task_begin_busy: 0,
3500 writer_task_begin_errors: 0,
3501 writer_task_request_failures: 0,
3502 writer_task_side_effects_unknown: 0,
3503 }
3504 );
3505
3506 drop(held);
3507 let _reacquired = pool.writer().expect("checkout succeeds after release");
3508 assert_eq!(
3509 pool.writer_acquisition_snapshot(),
3510 WriterAcquisitionSnapshot {
3511 acquisitions: 2,
3512 pooled_acquisitions: 2,
3513 standalone_acquisitions: 0,
3514 writer_task_acquisitions: 0,
3515 timeouts: 1,
3516 writer_task_begin_busy: 0,
3517 writer_task_begin_errors: 0,
3518 writer_task_request_failures: 0,
3519 writer_task_side_effects_unknown: 0,
3520 }
3521 );
3522 }
3523
3524 #[test]
3525 fn zero_wait_maintenance_skip_is_not_reported_as_a_checkout_timeout() {
3526 let pool = ConnectionPool::new(PoolConfig::default()).unwrap();
3527 let held = pool.writer().expect("finite-wait checkout succeeds");
3528 let before = pool.writer_acquisition_snapshot();
3529
3530 assert!(
3531 pool.try_writer_nowait().is_err(),
3532 "zero-wait maintenance checkout must skip while held"
3533 );
3534
3535 assert_eq!(
3536 pool.writer_acquisition_snapshot(),
3537 before,
3538 "a checkpoint-style zero-wait skip is not a finite-wait checkout timeout"
3539 );
3540 drop(held);
3541 }
3542
3543 #[test]
3547 #[serial(tx_registry)]
3548 fn writer_guard_transaction_registers_during_closure_only() {
3549 let cfg = PoolConfig {
3550 path: None,
3551 ..PoolConfig::default()
3552 };
3553 let pool = ConnectionPool::new(cfg).unwrap();
3554 let guard = pool.writer().unwrap();
3555
3556 let mut seen_during_closure = false;
3557 let result: Result<(), SqliteError> = guard.transaction(|_conn| {
3558 seen_during_closure = khive_storage::tx_registry::snapshot()
3559 .iter()
3560 .any(|(_, label)| label.as_deref() == Some("writer_guard_tx"));
3561 Ok(())
3562 });
3563 result.expect("transaction should commit");
3564
3565 assert!(
3566 seen_during_closure,
3567 "expected a writer_guard_tx entry visible inside the closure"
3568 );
3569 assert!(
3570 !khive_storage::tx_registry::snapshot()
3571 .iter()
3572 .any(|(_, label)| label.as_deref() == Some("writer_guard_tx")),
3573 "expected the entry to be gone after the transaction completes"
3574 );
3575 }
3576
3577 #[test]
3581 fn writer_task_handle_fails_loud_without_tokio_runtime() {
3582 let dir = tempfile::tempdir().unwrap();
3583 let path = dir.path().join("writer_task_no_runtime.db");
3584 let cfg = PoolConfig {
3585 path: Some(path),
3586 write_queue_enabled: Some(true),
3587 ..PoolConfig::default()
3588 };
3589 let pool = ConnectionPool::new(cfg).expect("file-backed pool should open");
3590
3591 let result = pool.writer_task_handle();
3592
3593 assert!(
3594 matches!(result, Err(StorageError::WriterTaskNoRuntime)),
3595 "expected Err(StorageError::WriterTaskNoRuntime) outside a Tokio \
3596 runtime, got {result:?}"
3597 );
3598 assert_eq!(
3599 pool.writer_task_spawn_count(),
3600 0,
3601 "the guard must reject before ever attempting tokio::spawn"
3602 );
3603 }
3604
3605 #[test]
3608 fn strict_writer_task_for_write_preserves_missing_runtime_error() {
3609 let dir = tempfile::tempdir().unwrap();
3610 let pool = ConnectionPool::new(PoolConfig {
3611 path: Some(dir.path().join("strict_writer_task_no_runtime.db")),
3612 write_queue_enabled: Some(true),
3613 write_routing_strict: true,
3614 ..PoolConfig::default()
3615 })
3616 .expect("file-backed pool should open");
3617
3618 let result = pool.writer_task_for_write(None, "strict_test_write");
3619
3620 assert!(
3621 matches!(result, Err(StorageError::WriterTaskNoRuntime)),
3622 "strict routing must preserve WriterTaskNoRuntime, got {result:?}"
3623 );
3624 assert_eq!(pool.writer_task_spawn_count(), 0);
3625 }
3626
3627 #[tokio::test]
3632 async fn take_writer_task_join_returns_some_once_then_none() {
3633 let dir = tempfile::tempdir().unwrap();
3634 let path = dir.path().join("join_lifecycle.db");
3635 let pool = ConnectionPool::new(PoolConfig {
3636 path: Some(path),
3637 write_queue_enabled: Some(true),
3638 ..PoolConfig::default()
3639 })
3640 .expect("file-backed pool should open");
3641
3642 assert!(
3645 pool.take_writer_task_join().is_none(),
3646 "before spawn there is no JoinHandle to take"
3647 );
3648 assert!(!pool.writer_task_join_was_stored());
3649 pool.writer_task_handle()
3650 .expect("runtime is present")
3651 .expect("write queue enabled must spawn a writer task");
3652 assert!(pool.writer_task_join_was_stored());
3653
3654 let join = pool
3655 .take_writer_task_join()
3656 .expect("the first take must return the spawned task's JoinHandle");
3657 assert!(
3658 pool.take_writer_task_join().is_none(),
3659 "the second take must return None — the handle is one-shot"
3660 );
3661
3662 drop(pool);
3668 tokio::time::timeout(Duration::from_secs(5), join)
3669 .await
3670 .expect("the writer task must exit once every handle clone is dropped")
3671 .expect("the writer task must not panic");
3672 }
3673
3674 #[cfg(debug_assertions)]
3678 #[tokio::test]
3679 #[should_panic(expected = "writer task JoinHandle stored twice")]
3680 async fn set_writer_task_join_second_store_trips_debug_assert() {
3681 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
3682 pool.set_writer_task_join(tokio::spawn(async {}));
3683 pool.set_writer_task_join(tokio::spawn(async {}));
3684 }
3685
3686 #[cfg(debug_assertions)]
3691 #[tokio::test]
3692 #[should_panic(expected = "writer task JoinHandle stored twice")]
3693 async fn set_writer_task_join_second_store_after_take_trips_debug_assert() {
3694 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
3695 pool.set_writer_task_join(tokio::spawn(async {}));
3696 assert!(pool.take_writer_task_join().is_some());
3697 pool.set_writer_task_join(tokio::spawn(async {}));
3698 }
3699
3700 #[cfg(not(debug_assertions))]
3706 #[tokio::test]
3707 async fn set_writer_task_join_first_wins_keeps_existing_handle() {
3708 let pool = ConnectionPool::new(PoolConfig::default()).expect("in-memory pool should open");
3709
3710 let (first_done_tx, first_done_rx) = tokio::sync::oneshot::channel::<()>();
3713 let first = tokio::spawn(async move {
3714 let _ = first_done_tx.send(());
3715 });
3716 let (_never_sent, never_rx) = tokio::sync::oneshot::channel::<()>();
3720 let second = tokio::spawn(async move {
3721 let _ = never_rx.await;
3722 });
3723
3724 pool.set_writer_task_join(first);
3725 pool.set_writer_task_join(second);
3726
3727 let taken = pool
3728 .take_writer_task_join()
3729 .expect("the first handle must still be stored");
3730 tokio::time::timeout(Duration::from_secs(5), taken)
3731 .await
3732 .expect("stored handle must be the first task's; the second never completes")
3733 .expect("the first task must not panic");
3734 assert!(
3735 first_done_rx.await.is_ok(),
3736 "completing the taken handle must mean the FIRST task ran to completion"
3737 );
3738 }
3739
3740 #[test]
3745 #[serial(pool_cwd)]
3746 fn mint_db_identity_alias_convergence() {
3747 let dir = tempfile::tempdir().unwrap();
3748 let real_dir = dir.path().join("real");
3749 fs::create_dir(&real_dir).unwrap();
3750 let db_path = real_dir.join("khive.db");
3751 fs::write(&db_path, b"").unwrap();
3752
3753 #[cfg(unix)]
3754 let dir_symlink = dir.path().join("dir_link");
3755 #[cfg(unix)]
3756 let file_symlink = dir.path().join("file_link.db");
3757 #[cfg(unix)]
3758 {
3759 std::os::unix::fs::symlink(&real_dir, &dir_symlink).unwrap();
3760 std::os::unix::fs::symlink(&db_path, &file_symlink).unwrap();
3761 }
3762
3763 let (via_real, canonical_real) = mint_db_identity(&db_path).unwrap();
3764
3765 let relative_result = {
3767 let _cwd = CwdGuard::enter(&real_dir);
3768 mint_db_identity(&PathBuf::from("khive.db"))
3769 };
3770 let (via_relative, canonical_relative) = relative_result.unwrap();
3771 assert_eq!(canonical_real, canonical_relative);
3772 assert_eq!(via_real, via_relative);
3773
3774 #[cfg(unix)]
3775 {
3776 let (via_dir_symlink, canonical_dir_symlink) =
3777 mint_db_identity(&dir_symlink.join("khive.db")).unwrap();
3778 assert_eq!(canonical_real, canonical_dir_symlink);
3779 assert_eq!(via_real, via_dir_symlink);
3780
3781 let (via_file_symlink, canonical_file_symlink) =
3782 mint_db_identity(&file_symlink).unwrap();
3783 assert_eq!(canonical_real, canonical_file_symlink);
3784 assert_eq!(via_real, via_file_symlink);
3785 }
3786
3787 let bare_name_result = {
3789 let _cwd = CwdGuard::enter(&real_dir);
3790 mint_db_identity(&PathBuf::from("khive.db"))
3791 };
3792 let (via_bare_name, canonical_bare_name) = bare_name_result.unwrap();
3793 assert_eq!(canonical_real, canonical_bare_name);
3794 assert_eq!(via_real, via_bare_name);
3795 }
3796
3797 #[test]
3809 #[serial(pool_cwd)]
3810 fn sidecar_dir_for_alias_convergence() {
3811 let dir = tempfile::tempdir().unwrap();
3812 let real_dir = dir.path().join("real");
3813 fs::create_dir(&real_dir).unwrap();
3814 let db_path = real_dir.join("khive.db");
3815 fs::write(&db_path, b"").unwrap();
3816
3817 #[cfg(unix)]
3818 let dir_symlink = dir.path().join("dir_link");
3819 #[cfg(unix)]
3820 let file_symlink = dir.path().join("file_link.db");
3821 #[cfg(unix)]
3822 {
3823 std::os::unix::fs::symlink(&real_dir, &dir_symlink).unwrap();
3824 std::os::unix::fs::symlink(&db_path, &file_symlink).unwrap();
3825 }
3826
3827 let pool_for = |path: &Path| -> Arc<ConnectionPool> {
3828 let cfg = PoolConfig {
3829 path: Some(path.to_path_buf()),
3830 ..PoolConfig::default()
3831 };
3832 Arc::new(ConnectionPool::new(cfg).expect("file-backed pool should open"))
3833 };
3834 let sidecar_of = |pool: &ConnectionPool| -> PathBuf {
3835 crate::walpin::sidecar_dir_for(pool.canonical_path().expect("file-backed pool"))
3836 };
3837
3838 let via_real = pool_for(&db_path);
3839 let sidecar_real = sidecar_of(&via_real);
3840
3841 let via_relative = {
3842 let _cwd = CwdGuard::enter(&real_dir);
3843 pool_for(Path::new("khive.db"))
3844 };
3845 assert_eq!(
3846 sidecar_real,
3847 sidecar_of(&via_relative),
3848 "a relative spelling of the same database must derive the same sidecar directory"
3849 );
3850
3851 #[cfg(unix)]
3852 {
3853 let via_dir_symlink = pool_for(&dir_symlink.join("khive.db"));
3854 assert_eq!(
3855 sidecar_real,
3856 sidecar_of(&via_dir_symlink),
3857 "opening through a directory symlink must derive the same sidecar directory"
3858 );
3859
3860 let via_file_symlink = pool_for(&file_symlink);
3861 assert_eq!(
3862 sidecar_real,
3863 sidecar_of(&via_file_symlink),
3864 "opening through a file-level symlink must derive the same sidecar directory"
3865 );
3866 }
3867
3868 let via_bare_name = {
3869 let _cwd = CwdGuard::enter(&real_dir);
3870 pool_for(Path::new("khive.db"))
3871 };
3872 assert_eq!(
3873 sidecar_real,
3874 sidecar_of(&via_bare_name),
3875 "a bare file name resolved against the current directory must derive the same \
3876 sidecar directory"
3877 );
3878 }
3879
3880 #[cfg(unix)]
3886 #[test]
3887 fn mint_db_identity_dangling_symlink_first_open_convergence() {
3888 let dir = tempfile::tempdir().unwrap();
3889 let target = dir.path().join("target.db");
3890 let link = dir.path().join("link.db");
3891 std::os::unix::fs::symlink(&target, &link).unwrap();
3892 assert!(!target.exists(), "target must not exist yet (dangling)");
3893
3894 let (via_dangling_link, canonical_via_link) = mint_db_identity(&link).unwrap();
3895
3896 fs::write(&target, b"").unwrap();
3899 let (via_target, canonical_via_target) = mint_db_identity(&target).unwrap();
3900
3901 assert_eq!(canonical_via_link, canonical_via_target);
3902 assert_eq!(via_dangling_link, via_target);
3903 }
3904
3905 #[test]
3908 fn mint_db_identity_missing_parent_fails() {
3909 let dir = tempfile::tempdir().unwrap();
3910 let missing = dir.path().join("nonexistent_subdir").join("khive.db");
3911 let result = mint_db_identity(&missing);
3912 assert!(
3913 result.is_err(),
3914 "minting must fail when the parent directory does not exist"
3915 );
3916 }
3917
3918 #[cfg(unix)]
3921 #[test]
3922 fn mint_db_identity_non_utf8_path_round_trips() {
3923 use std::ffi::OsStr;
3924 use std::os::unix::ffi::OsStrExt;
3925
3926 let dir = tempfile::tempdir().unwrap();
3927 let raw_name = OsStr::from_bytes(b"khive-\xffdb.sqlite");
3929 let db_path = dir.path().join(raw_name);
3930 if let Err(e) = fs::write(&db_path, b"") {
3935 eprintln!(
3936 "skipping mint_db_identity_non_utf8_path_round_trips: filesystem rejected a \
3937 non-UTF-8 file name ({e}); this platform's filesystem does not support the \
3938 case under test"
3939 );
3940 return;
3941 }
3942
3943 let (identity, canonical) = mint_db_identity(&db_path).unwrap();
3944 assert_eq!(canonical.file_name().unwrap(), raw_name);
3945
3946 let (identity_again, canonical_again) = mint_db_identity(&db_path).unwrap();
3947 assert_eq!(identity, identity_again);
3948 assert_eq!(canonical, canonical_again);
3949 }
3950
3951 #[test]
3957 fn resolve_reader_checkout_maps_each_arm_distinctly() {
3958 let pool = ConnectionPool::new(PoolConfig {
3959 path: None,
3960 ..PoolConfig::default()
3961 })
3962 .unwrap();
3963
3964 let guard = pool
3965 .resolve_reader_checkout(
3966 StorageCapability::Sql,
3967 "arm_checked_out",
3968 pool.reader_until(|| false),
3969 )
3970 .expect("an uncontended checkout must pass the guard through");
3971 drop(guard);
3972
3973 let Err(cancelled) =
3974 pool.resolve_reader_checkout(StorageCapability::Sql, "arm_cancelled", Ok(None))
3975 else {
3976 panic!("a cancelled checkout must be refused");
3977 };
3978 assert!(
3979 matches!(cancelled, StorageError::Timeout { .. }),
3980 "cancellation/deadline before checkout must be the non-retryable \
3981 Timeout, got {cancelled:?}"
3982 );
3983
3984 let Err(exhausted) = pool.resolve_reader_checkout(
3985 StorageCapability::Sql,
3986 "arm_exhausted",
3987 Err(pool_exhausted_error(Duration::from_millis(5), 1)),
3988 ) else {
3989 panic!("an exhausted checkout must be refused");
3990 };
3991 assert!(
3992 matches!(exhausted, StorageError::AdmissionTimeout { .. }),
3993 "the pool's own SQLITE_BUSY (checkout_timeout exhausted) must be \
3994 the retryable AdmissionTimeout, got {exhausted:?}"
3995 );
3996
3997 let Err(opaque) = pool.resolve_reader_checkout(
3998 StorageCapability::Entities,
3999 "arm_driver",
4000 Err(SqliteError::InvalidData("retired pooled writer".into())),
4001 ) else {
4002 panic!("an opaque checkout error must be refused");
4003 };
4004 assert!(
4005 matches!(
4006 &opaque,
4007 StorageError::Driver { capability, .. }
4008 if *capability == StorageCapability::Entities
4009 ),
4010 "any other checkout error must stay a non-retryable Driver failure \
4011 under the caller's capability, got {opaque:?}"
4012 );
4013 }
4014}