1use super::*;
2
3#[cfg(test)]
4pub(super) type SpaceProbe = dyn Fn(&Path) -> std::io::Result<u64> + Send + Sync;
5
6#[cfg(test)]
7thread_local! {
8 pub(super) static STARTUP_SPACE_PROBE: std::cell::RefCell<Option<(u64, Arc<SpaceProbe>)>> =
9 const { std::cell::RefCell::new(None) };
10}
11
12pub(crate) struct WriteAdmission {
17 database_path: Option<PathBuf>,
18 volume: Option<PathBuf>,
19 expected_volume: Option<VolumeIdentity>,
20 volume_lock_dir: Option<PathBuf>,
21 floor_bytes: u64,
22 guard_deadline: Duration,
23 settlement_unknown: AtomicBool,
26 lease_timeouts: AtomicU64,
29 #[cfg(test)]
30 space_probe: Mutex<Option<Arc<SpaceProbe>>>,
31 #[cfg(test)]
32 current_volume_override: Mutex<Option<VolumeIdentity>>,
33}
34
35impl WriteAdmission {
36 pub(crate) fn new(
37 database_path: Option<PathBuf>,
38 floor_bytes: u64,
39 guard_deadline_ms: u64,
40 volume_lock_dir: Option<PathBuf>,
41 ) -> Result<Self, SqliteError> {
42 let volume = database_path.clone();
43 let expected_volume = volume.as_deref().map(VolumeIdentity::resolve).transpose()?;
44 #[cfg(test)]
45 let (floor_bytes, space_probe) = STARTUP_SPACE_PROBE.with(|probe| {
46 probe
47 .borrow()
48 .as_ref()
49 .map(|(floor, probe)| (*floor, Some(Arc::clone(probe))))
50 .unwrap_or((floor_bytes, None))
51 });
52 Ok(Self {
53 database_path,
54 volume,
55 expected_volume,
56 volume_lock_dir,
57 floor_bytes,
58 guard_deadline: Duration::from_millis(guard_deadline_ms),
59 settlement_unknown: AtomicBool::new(false),
60 lease_timeouts: AtomicU64::new(0),
61 #[cfg(test)]
62 space_probe: Mutex::new(space_probe),
63 #[cfg(test)]
64 current_volume_override: Mutex::new(None),
65 })
66 }
67
68 pub(crate) fn for_canonical_path(database_path: Option<PathBuf>) -> Result<Self, SqliteError> {
71 if database_path.is_none() {
72 return Self::new(None, 0, DEFAULT_DISK_GUARD_DEADLINE_MS, None);
73 }
74 let policy = crate::migrations::MigrationWritePolicy::from_environment()?;
75 Self::for_migration_policy(database_path, &policy)
76 }
77
78 pub(crate) fn for_migration_policy(
79 database_path: Option<PathBuf>,
80 policy: &crate::migrations::MigrationWritePolicy,
81 ) -> Result<Self, SqliteError> {
82 let effective = policy.disk_guard_config();
83 if effective.legacy_environment_present {
84 tracing::warn!(
85 "legacy SQLite reserve setting is deprecated; use KHIVE_SQLITE_DISK_RESERVE_BYTES"
86 );
87 }
88 if effective.reserve_bytes == 0 {
89 tracing::warn!("SQLite disk reserve is zero; logical writes have no floor refusal");
90 }
91 Self::new(
92 database_path,
93 effective.reserve_bytes,
94 effective.guard_deadline_ms,
95 Some(policy.volume_lock_dir().to_path_buf()),
96 )
97 }
98
99 pub(super) fn verify_current_volume(&self) -> Result<(), SqliteError> {
100 let (Some(path), Some(expected)) = (self.volume.as_deref(), self.expected_volume.as_ref())
101 else {
102 return Ok(());
103 };
104 #[cfg(test)]
105 let current = self.current_volume_override.lock().clone();
106 #[cfg(not(test))]
107 let current: Option<VolumeIdentity> = None;
108 let current = match current {
109 Some(identity) => identity,
110 None => VolumeIdentity::resolve(path)?,
111 };
112 if ¤t != expected {
113 return Err(SqliteError::CapacityUnavailable {
114 phase: CapacityUnavailablePhase::Identity,
115 message: "database volume changed after admission identity was captured"
116 .to_string(),
117 });
118 }
119 Ok(())
120 }
121
122 pub(crate) fn close_retired_connection(&self, conn: Connection) -> Result<(), SqliteError> {
128 let database_path = self
129 .database_path
130 .as_ref()
131 .map(|path| path.display().to_string())
132 .unwrap_or_else(|| ":memory:".to_string());
133 let volume_key = self
134 .expected_volume
135 .as_ref()
136 .map(VolumeIdentity::diagnostic_key)
137 .unwrap_or_else(|| "in-memory".to_string());
138 let settled = crate::connection_settlement::close_retired_connection(
139 conn,
140 &database_path,
141 &volume_key,
142 );
143 if settled.is_err() {
144 self.settlement_unknown.store(true, Ordering::Release);
145 }
146 settled
147 }
148
149 pub(crate) fn ensure_settled(&self) -> Result<(), SqliteError> {
150 if self.settlement_unknown.load(Ordering::Acquire) {
151 return Err(SqliteError::WriterPoisoned);
152 }
153 Ok(())
154 }
155
156 pub(crate) fn volume_identity(&self) -> Result<Option<VolumeIdentity>, SqliteError> {
157 self.verify_current_volume()?;
158 Ok(self.expected_volume.clone())
159 }
160
161 fn available_space(&self, volume: &Path) -> std::io::Result<u64> {
162 #[cfg(test)]
163 if let Some(probe) = self.space_probe.lock().as_ref() {
164 return probe(volume);
165 }
166 fs4::available_space(volume)
167 }
168
169 pub(crate) fn check(&self) -> Result<(), SqliteError> {
170 self.check_with_headroom(0)
171 }
172
173 pub(crate) fn check_with_headroom(
174 &self,
175 required_headroom_bytes: u64,
176 ) -> Result<(), SqliteError> {
177 let Some(volume) = self.volume.as_deref() else {
178 return Ok(());
179 };
180 self.verify_current_volume()?;
181 let identity =
182 self.expected_volume
183 .as_ref()
184 .ok_or_else(|| SqliteError::CapacityUnavailable {
185 phase: CapacityUnavailablePhase::Identity,
186 message: "file-backed admission has no captured volume identity".to_string(),
187 })?;
188 let available = self
189 .available_space(identity.probe_path())
190 .map_err(|error| SqliteError::CapacityUnavailable {
191 phase: CapacityUnavailablePhase::Probe,
192 message: format!(
193 "cannot sample available space on {}: {error}",
194 volume.display()
195 ),
196 })?;
197 self.verify_current_volume()?;
198 let threshold = self.floor_bytes.checked_add(required_headroom_bytes);
204 if threshold.is_none()
205 || (threshold != Some(0) && available <= threshold.unwrap_or(u64::MAX))
206 {
207 return Err(SqliteError::CapacityFloor {
208 volume: volume.display().to_string(),
209 available_bytes: available,
210 floor_bytes: self.floor_bytes,
211 required_headroom_bytes,
212 });
213 }
214 Ok(())
215 }
216
217 #[track_caller]
222 pub(crate) fn acquire(&self) -> Result<Option<VolumeLease>, SqliteError> {
223 self.ensure_settled()?;
224 let Some(identity) = self.expected_volume.as_ref() else {
225 return Ok(None);
226 };
227 self.verify_current_volume()?;
228 #[cfg(test)]
229 if crate::disk_guard::SKIP_VOLUME_LEASE_FOR_MEASUREMENT.load(Ordering::Relaxed) {
230 self.verify_current_volume()?;
231 return Ok(Some(VolumeLease::unheld_for_measurement()));
232 }
233 let lease = match identity.acquire(self.guard_deadline, self.volume_lock_dir.as_deref()) {
234 Ok(lease) => lease,
235 Err(LeaseRefusal::TimedOut(error)) => {
236 self.record_lease_timeout(&error);
237 return Err(error);
238 }
239 Err(refusal) => return Err(refusal.into()),
240 };
241 self.verify_current_volume()?;
242 Ok(Some(lease))
243 }
244
245 fn record_lease_timeout(&self, error: &SqliteError) {
250 self.lease_timeouts.fetch_add(1, Ordering::Relaxed);
251 let db = self
252 .database_path
253 .as_ref()
254 .map(|path| path.display().to_string())
255 .unwrap_or_else(|| "memory".to_string());
256 crate::timeout_sink::emit_lease_timeout(
257 &db,
258 &error.to_string(),
259 self.guard_deadline.as_millis().min(u128::from(u64::MAX)) as u64,
260 );
261 }
262
263 pub(crate) fn lease_timeouts(&self) -> u64 {
264 self.lease_timeouts.load(Ordering::Relaxed)
265 }
266
267 pub(crate) fn check_for_vacuum(&self) -> Result<(), SqliteError> {
271 let headroom = self
272 .database_path
273 .as_deref()
274 .map(crate::vacuum_capacity::estimate_vacuum_headroom)
275 .transpose()?
276 .unwrap_or(0);
277 self.check_with_headroom(headroom)
278 }
279
280 #[cfg(test)]
281 pub(super) fn set_test_space_probe(
282 &self,
283 probe: impl Fn(&Path) -> std::io::Result<u64> + Send + Sync + 'static,
284 ) {
285 *self.space_probe.lock() = Some(Arc::new(probe));
286 }
287
288 #[cfg(test)]
289 pub(crate) fn set_test_current_volume(&self, identity: Option<VolumeIdentity>) {
290 *self.current_volume_override.lock() = identity;
291 }
292
293 #[cfg(test)]
294 pub(crate) fn captured_volume_for_test(&self) -> Option<VolumeIdentity> {
295 self.expected_volume.clone()
296 }
297}
298
299pub struct WriterGuard<'pool> {
311 pub(super) guard: parking_lot::MutexGuard<'pool, Connection>,
312 pub(super) origin: TxOrigin,
316 pub(super) pool: &'pool ConnectionPool,
317 pub(super) admission: &'pool WriteAdmission,
318 pub(super) _volume_lease: Option<VolumeLease>,
321}
322
323pub(crate) struct PooledAutocommitWriteUnit<'pool> {
327 writer: WriterGuard<'pool>,
328}
329
330impl PooledAutocommitWriteUnit<'_> {
331 pub(crate) fn conn(&self) -> &Connection {
332 self.writer.conn()
333 }
334}
335
336pub(crate) struct PooledTransactionWriteUnit<'pool> {
339 writer: WriterGuard<'pool>,
340}
341
342impl PooledTransactionWriteUnit<'_> {
343 pub(crate) fn conn(&self) -> &Connection {
344 self.writer.conn()
345 }
346}
347
348pub(crate) struct StandaloneTransactionWriteUnit {
352 conn: Option<Connection>,
353 admission: Arc<WriteAdmission>,
354 _volume_lease: Option<VolumeLease>,
355}
356
357impl StandaloneTransactionWriteUnit {
358 pub(crate) fn conn(&self) -> &Connection {
359 self.conn
360 .as_ref()
361 .expect("standalone write unit owns its connection")
362 }
363}
364
365impl Drop for StandaloneTransactionWriteUnit {
366 fn drop(&mut self) {
367 let conn = self
368 .conn
369 .take()
370 .expect("standalone write unit retirement owner");
371 if conn.is_autocommit() {
372 drop(conn);
373 } else {
374 let _ = self.admission.close_retired_connection(conn);
377 }
378 }
379}
380
381pub struct CheckpointGuard<'pool> {
392 pub(super) guard: parking_lot::MutexGuard<'pool, Connection>,
393}
394
395#[derive(Debug, Clone, Copy, PartialEq, Eq)]
397pub struct CheckpointResult {
398 pub busy: i64,
400 pub log_frames: i64,
402 pub checkpointed_frames: i64,
404}
405
406impl CheckpointGuard<'_> {
407 pub fn passive(&self) -> Result<CheckpointResult, SqliteError> {
409 self.guard
410 .query_row("PRAGMA wal_checkpoint(PASSIVE)", [], |row| {
411 Ok(CheckpointResult {
412 busy: row.get(0)?,
413 log_frames: row.get(1)?,
414 checkpointed_frames: row.get(2)?,
415 })
416 })
417 .map_err(Into::into)
418 }
419
420 pub fn truncate(&self) -> Result<CheckpointResult, SqliteError> {
422 self.guard
423 .query_row("PRAGMA wal_checkpoint(TRUNCATE)", [], |row| {
424 Ok(CheckpointResult {
425 busy: row.get(0)?,
426 log_frames: row.get(1)?,
427 checkpointed_frames: row.get(2)?,
428 })
429 })
430 .map_err(Into::into)
431 }
432}
433
434#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
443pub struct WriterAcquisitionSnapshot {
444 pub acquisitions: u64,
447 pub pooled_acquisitions: u64,
449 pub standalone_acquisitions: u64,
451 pub writer_task_acquisitions: u64,
454 pub timeouts: u64,
459 pub lease_timeouts: u64,
465 pub direct_busy_refusals: u64,
470 pub writer_task_begin_busy: u64,
475 pub writer_task_begin_busy_absorbed: u64,
480 pub writer_task_begin_errors: u64,
483 pub writer_task_request_failures: u64,
487 pub writer_task_side_effects_unknown: u64,
492 pub writer_guard_drop_rollbacks: u64,
499}
500
501#[derive(Debug, Default)]
505pub(crate) struct WriterAcquisitionCounters {
506 pub(super) pooled_acquisitions: AtomicU64,
507 pub(super) standalone_acquisitions: AtomicU64,
508 pub(super) writer_task_acquisitions: AtomicU64,
509 pub(super) pooled_timeouts: AtomicU64,
510 pub(super) direct_busy_refusals: AtomicU64,
511 pub(super) writer_task_begin_busy: AtomicU64,
512 pub(super) writer_task_begin_busy_absorbed: AtomicU64,
513 pub(super) writer_task_begin_errors: AtomicU64,
514 pub(super) writer_task_request_failures: AtomicU64,
515 pub(super) writer_task_side_effects_unknown: AtomicU64,
516 pub(super) writer_guard_drop_rollbacks: AtomicU64,
517}
518
519impl<'pool> WriterGuard<'pool> {
520 pub(crate) fn admit_autocommit(self) -> Result<PooledAutocommitWriteUnit<'pool>, SqliteError> {
524 if !self.guard.is_autocommit() {
525 self.pool.retire_pooled_writer(&self.guard);
526 return Err(SqliteError::InvalidData(
527 "pooled autocommit write began on a connection in a transaction".to_string(),
528 ));
529 }
530 if self._volume_lease.is_none() && self.pool.canonical_path().is_some() {
531 return Err(SqliteError::InvalidData(
532 "checkpoint writer checkout cannot admit a logical write".to_string(),
533 ));
534 }
535 self.admission.check()?;
536 Ok(PooledAutocommitWriteUnit { writer: self })
537 }
538
539 pub fn conn(&self) -> &Connection {
541 &self.guard
542 }
543
544 pub fn conn_mut(&mut self) -> &mut Connection {
546 &mut self.guard
547 }
548
549 pub fn transaction<F, R>(&self, f: F) -> Result<R, SqliteError>
552 where
553 F: FnOnce(&Connection) -> Result<R, SqliteError>,
554 {
555 if self._volume_lease.is_none() && self.pool.canonical_path().is_some() {
556 return Err(SqliteError::InvalidData(
557 "checkpoint writer checkout cannot start a logical transaction".to_string(),
558 ));
559 }
560 self.guard.execute_batch("BEGIN IMMEDIATE")?;
561 if let Err(error) = self.admission.check() {
562 self.rollback_or_retire("capacity admission")?;
563 return Err(error);
564 }
565 let _tx_handle = khive_storage::tx_registry::register_scoped(
566 Some("writer_guard_tx".to_string()),
567 self.origin.clone(),
568 );
569
570 match f(&self.guard) {
571 Ok(result) => {
572 if let Err(err) = self.guard.execute_batch("COMMIT") {
573 self.rollback_or_retire("commit failure")?;
574 return Err(err.into());
575 }
576 Ok(result)
577 }
578 Err(err) => {
579 self.rollback_or_retire("transaction body failure")?;
580 Err(err)
581 }
582 }
583 }
584
585 pub(crate) fn rollback_or_retire(&self, context: &str) -> Result<(), SqliteError> {
586 let rollback = self.guard.execute_batch("ROLLBACK");
587 if rollback.is_err() || !self.guard.is_autocommit() {
588 self.pool.retire_pooled_writer(&self.guard);
589 tracing::error!(context, "pooled writer rollback did not prove autocommit");
590 return Err(SqliteError::WriterSettlementUnknown);
591 }
592 Ok(())
593 }
594}
595
596pub(crate) struct PooledMigrationTransactions<'pool> {
598 pool: &'pool ConnectionPool,
599}
600
601impl crate::migrations::MigrationTransactions for PooledMigrationTransactions<'_> {
602 fn admitted<T>(
603 &mut self,
604 operation: impl FnOnce(&mut Connection) -> Result<T, SqliteError>,
605 ) -> Result<T, SqliteError> {
606 let mut writer = self.pool.writer_for_admitted_operation()?;
607 let result = operation(writer.conn_mut());
608 if writer.guard.is_autocommit() {
609 return result;
610 }
611 writer.rollback_or_retire("schema migration left its transaction open")?;
612 result.and(Err(SqliteError::WriterSettlementUnknown))
614 }
615}
616
617impl Drop for WriterGuard<'_> {
618 fn drop(&mut self) {
619 if !self.guard.is_autocommit() {
620 let settlement = self.rollback_or_retire("guard dropped with open transaction");
624 self.pool.record_writer_guard_drop(&settlement);
625 }
626 if self.pool.pooled_writer_retired.load(Ordering::Acquire) {
627 if let Some(replacement) = self.pool.retirement_connection.lock().take() {
628 let original = std::mem::replace(&mut *self.guard, replacement);
630 let _ = self.admission.close_retired_connection(original);
631 }
632 }
633 }
634}
635
636impl<'pool> Deref for WriterGuard<'pool> {
637 type Target = Connection;
638
639 fn deref(&self) -> &Self::Target {
640 self.conn()
641 }
642}
643
644impl<'pool> DerefMut for WriterGuard<'pool> {
645 fn deref_mut(&mut self) -> &mut Self::Target {
646 self.conn_mut()
647 }
648}
649
650impl ConnectionPool {
651 #[track_caller]
652 pub(crate) fn autocommit_write_unit(
653 &self,
654 ) -> Result<PooledAutocommitWriteUnit<'_>, SqliteError> {
655 self.writer_for_admitted_operation()?.admit_autocommit()
656 }
657
658 #[track_caller]
659 pub(crate) fn transaction_write_unit(
660 &self,
661 ) -> Result<super::PooledTransactionWriteUnit<'_>, SqliteError> {
662 let writer = self.writer_for_admitted_operation()?;
663 if !writer.guard.is_autocommit() {
664 writer.pool.retire_pooled_writer(&writer.guard);
665 return Err(SqliteError::InheritedWriterTransaction);
666 }
667 if let Err(error) = writer.guard.execute_batch("BEGIN IMMEDIATE") {
668 if !writer.guard.is_autocommit() {
669 writer.pool.retire_pooled_writer(&writer.guard);
670 return Err(SqliteError::WriterSettlementUnknown);
671 }
672 return Err(error.into());
673 }
674 if let Err(error) = writer.admission.check() {
675 writer.rollback_or_retire("capacity admission")?;
676 return Err(error);
677 }
678 Ok(PooledTransactionWriteUnit { writer })
679 }
680
681 #[track_caller]
686 pub(crate) fn execute_direct_transaction<R, F>(
687 &self,
688 capability: StorageCapability,
689 operation: &'static str,
690 f: F,
691 ) -> Result<R, StorageError>
692 where
693 F: FnOnce(&Connection) -> Result<R, StorageError>,
694 {
695 let _tx_handle =
696 khive_storage::tx_registry::register_scoped(Some(operation.to_string()), self.origin());
697 let db_label = crate::timeout_sink::db_label(self);
698 let map_admission_error = |error: SqliteError| {
699 crate::timeout_sink::maybe_emit_sqlite_full(&db_label, &error);
700 let error = error.into_storage_error(capability, operation);
701 self.record_direct_writer_error(&error);
702 error
703 };
704 let result = if self.canonical_path().is_some() {
705 let unit = self
706 .standalone_transaction_write_unit()
707 .map_err(map_admission_error)?;
708 let (result, _) = execute_wrapped_transaction(unit.conn(), operation, f);
709 result
710 } else {
711 let unit = self.transaction_write_unit().map_err(map_admission_error)?;
712 let conn = unit.conn();
713 let (result, terminal_state) = execute_wrapped_transaction(conn, operation, f);
714 if terminal_state.is_some() {
715 self.retire_pooled_writer(conn);
716 }
717 result
718 };
719 if let Err(error) = &result {
720 crate::timeout_sink::maybe_emit_sqlite_full(&db_label, error);
721 self.record_direct_writer_error(error);
722 }
723 result
724 }
725
726 pub fn run_migrations(&self) -> Result<u32, SqliteError> {
729 if self.config.read_only {
730 return Err(SqliteError::InvalidData(
731 "cannot run migrations on a read-only pool".to_string(),
732 ));
733 }
734 let owner = crate::stores::blob::acquire_database_gc_owner_for_path_blocking(
735 self.canonical_path().map(Path::to_path_buf),
736 )
737 .map_err(|error| {
738 SqliteError::InvalidData(format!(
739 "failed to acquire database GC owner before schema migration: {error}"
740 ))
741 })?;
742 crate::migrations::run_migrations_with_database_gc_owner(
743 &mut self.migration_transactions(),
744 &owner,
745 &self.write_admission,
746 )
747 }
748
749 pub(crate) fn migration_transactions(&self) -> PooledMigrationTransactions<'_> {
753 PooledMigrationTransactions { pool: self }
754 }
755
756 #[track_caller]
760 pub(crate) fn standalone_transaction_write_unit(
761 &self,
762 ) -> Result<super::StandaloneTransactionWriteUnit, SqliteError> {
763 if self.config.read_only {
764 return Err(SqliteError::InvalidData(
765 "database is read-only: standalone write transactions are not permitted".into(),
766 ));
767 }
768 let volume_lease = self.write_admission.acquire()?;
769 let conn = self.open_standalone_writer_for_admitted_operation()?;
770 let unit = StandaloneTransactionWriteUnit {
771 conn: Some(conn),
772 admission: Arc::clone(&self.write_admission),
773 _volume_lease: volume_lease,
774 };
775 if !unit.conn().is_autocommit() {
776 return Err(SqliteError::WriterSettlementUnknown);
777 }
778 if let Err(error) = unit.conn().execute_batch("BEGIN IMMEDIATE") {
779 if !unit.conn().is_autocommit() {
780 return Err(SqliteError::WriterSettlementUnknown);
781 }
782 return Err(error.into());
783 }
784 if let Err(error) = self.write_admission.check() {
785 if unit.conn().execute_batch("ROLLBACK").is_err() || !unit.conn().is_autocommit() {
786 return Err(SqliteError::WriterSettlementUnknown);
787 }
788 return Err(error);
789 }
790 Ok(unit)
791 }
792}