Skip to main content

khive_db/pool/
write_units.rs

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
12/// The SQLite write reserve is sampled at each operation admission. SQLite
13/// does not expose the size of an arbitrary upcoming transaction, so the
14/// reserve is a warning boundary, not a guarantee that a single very large
15/// transaction cannot consume more than the remaining headroom.
16pub(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    /// Set once a retired connection could neither roll back nor close: the
24    /// outcome of its transaction is unknown, so every later write is refused.
25    settlement_unknown: AtomicBool,
26    /// Lease waits that ran out `guard_deadline` because another writer held
27    /// the volume. Read through [`WriterAcquisitionSnapshot::lease_timeouts`].
28    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    /// Build admission for a raw connection with no backend-local override.
69    /// Pooled callers must use their pool's effective policy instead.
70    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 &current != 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    /// Settle a retired connection while the caller still holds its lease. A
123    /// terminal failure returns the typed outcome-unknown error for this write
124    /// and poisons this admission, so every later write on the pool, the
125    /// writer task and standalone writers is refused before it starts; the
126    /// caller may then release its lease.
127    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        // SQL does not tell admission how many bytes a generic transaction
199        // will append. VACUUM supplies its known copy-sized estimate here.
200        // A zero reserve removes the floor term only: an operation-specific
201        // headroom is still compared, and with neither term nothing is.
202        // Overflow is a refusal, never a wrapped low threshold.
203        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    /// Acquire the physical-volume lease before SQLite writer acquisition.
218    /// The caller must retain the returned guard through settlement. A request
219    /// from a thread that already holds this volume's lease fails at once as
220    /// re-entry, naming both call sites, instead of waiting out the deadline.
221    #[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    /// The lease is taken before the pool mutex, so ordinary same-volume
246    /// writer contention ends here rather than at `checkout_timeout`. Count it
247    /// and write a `timeout` row with `phase=lock`, as the pool mutex stage
248    /// does for its own timeouts, so the contention stays visible.
249    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    /// Admission for VACUUM: its copy-sized headroom is compared even under a
268    /// zero reserve, and a headroom that cannot be estimated refuses the
269    /// operation instead of admitting it with none.
270    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
299/// A writer connection checked out from the pool.
300/// The Mutex ensures only one writer at a time.
301///
302/// # Terminal settlement
303/// Dropping a retired writer settles its connection before the volume lease is
304/// released. If cleanup cannot restore autocommit and owned SQLite close also
305/// fails, the database path, volume, and errors are emitted to stderr and
306/// tracing, that write reports [`SqliteError::WriterSettlementUnknown`], the
307/// pool is poisoned so every later write is refused with
308/// [`SqliteError::WriterPoisoned`] before it starts, and the lease is then
309/// released.
310pub struct WriterGuard<'pool> {
311    pub(super) guard: parking_lot::MutexGuard<'pool, Connection>,
312    /// The origin (ADR-091 backend-scoped attribution) of the pool this
313    /// guard was checked out from, carried so `transaction` can register its
314    /// span with the correct origin without holding a `&ConnectionPool`.
315    pub(super) origin: TxOrigin,
316    pub(super) pool: &'pool ConnectionPool,
317    pub(super) admission: &'pool WriteAdmission,
318    /// Normal writer checkout acquires this before taking the writer mutex.
319    /// Maintenance-only nowait checkout leaves it absent to bypass the floor.
320    pub(super) _volume_lease: Option<VolumeLease>,
321}
322
323/// One synchronous pooled autocommit write or script. Construction samples
324/// capacity under the volume lease, and the owned writer guard retains that
325/// lease until this unit is dropped after SQLite returns to autocommit.
326pub(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
336/// One pooled transaction after lease, BEGIN, and capacity admission. The
337/// owned writer guard rolls back on drop if the caller fails to settle it.
338pub(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
348/// One standalone transaction on an owned connection. Settlement (autocommit
349/// or successful owned close) precedes release of the volume lease, including
350/// when the operation unwinds or its initial rollback fails.
351pub(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            // A terminal failure has poisoned the admission; Drop cannot
375            // return it, and the lease is released after this returns.
376            let _ = self.admission.close_retired_connection(conn);
377        }
378    }
379}
380
381/// A zero-wait checkout that can run only the fixed checkpoint recovery
382/// pragmas. The connection remains private: exposing it would let a caller
383/// execute logical writes without disk-reserve admission (ADR-154 §5).
384///
385/// Ordinary SQL is deliberately unavailable through this capability:
386/// ```compile_fail
387/// use khive_db::{ConnectionPool, PoolConfig};
388/// let pool = ConnectionPool::new(PoolConfig::default()).unwrap();
389/// pool.try_checkpoint_nowait().unwrap().execute_batch("CREATE TABLE bypass (id INTEGER)");
390/// ```
391pub struct CheckpointGuard<'pool> {
392    pub(super) guard: parking_lot::MutexGuard<'pool, Connection>,
393}
394
395/// SQLite's three-column result from a fixed WAL checkpoint pragma.
396#[derive(Debug, Clone, Copy, PartialEq, Eq)]
397pub struct CheckpointResult {
398    /// Whether SQLite reported a busy checkpoint.
399    pub busy: i64,
400    /// WAL frames observed by SQLite (`-1` when there is no WAL).
401    pub log_frames: i64,
402    /// WAL frames copied back into the database.
403    pub checkpointed_frames: i64,
404}
405
406impl CheckpointGuard<'_> {
407    /// Run a PASSIVE checkpoint without disk-reserve admission.
408    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    /// Run a TRUNCATE checkpoint without disk-reserve admission.
421    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/// Process-local monotonic counters for every instrumented writer acquisition
435/// boundary owned by one [`ConnectionPool`].
436///
437/// The aggregate `acquisitions` is the saturating sum of its three explicit
438/// connection classes. Infrastructure-only opens (the diagnostics PASSIVE
439/// probe, the writer task's one-time lifetime connection, and the checkpoint
440/// task's dedicated long-lived connection) are excluded; zero-wait
441/// maintenance probes also remain outside these request-traffic counters.
442#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
443pub struct WriterAcquisitionSnapshot {
444    /// Successful acquisitions across pooled, standalone, and writer-task
445    /// connection classes.
446    pub acquisitions: u64,
447    /// Successful finite-wait pool-mutex writer checkouts.
448    pub pooled_acquisitions: u64,
449    /// Successful per-operation standalone writer connection opens.
450    pub standalone_acquisitions: u64,
451    /// Successful writer-task ownership acquisitions (one per dequeued
452    /// top-level request or successful `BEGIN IMMEDIATE`).
453    pub writer_task_acquisitions: u64,
454    /// Finite-wait pool writer checkouts that exhausted `checkout_timeout`
455    /// waiting for the pool mutex. A file-backed writer reaches that wait only
456    /// after it holds the volume lease; a wait that ends at the lease is
457    /// counted in `lease_timeouts` instead.
458    pub timeouts: u64,
459    /// Writer acquisitions refused because another writer held the volume
460    /// lease past `disk_guard_deadline_ms`, returned as
461    /// `CapacityUnavailable { phase: Lock }`. Every write path takes the lease
462    /// before its connection, so same-volume writer contention is counted
463    /// here, across all of this pool's writer classes.
464    pub lease_timeouts: u64,
465    /// Instrumented direct executions whose final returned error retains SQLite's
466    /// primary DatabaseBusy code, once per operation after its busy handler.
467    /// Excludes LOCKED, checkout/open/admission failures, readers, writer tasks,
468    /// infrastructure probes and uninstrumented raw connection escapes.
469    pub direct_busy_refusals: u64,
470    /// Every writer-task `BEGIN IMMEDIATE` attempt refused busy or locked,
471    /// including refusals a subsequent bounded retry went on to absorb.
472    /// Counted separately from `timeouts` because that counter names the
473    /// pool-mutex checkout stage; folding the two would mislabel the stage.
474    pub writer_task_begin_busy: u64,
475    /// Subset of `writer_task_begin_busy` that a subsequent bounded retry
476    /// absorbed before the request closure ran, so the refusal never
477    /// reached the caller. `writer_task_begin_busy - writer_task_begin_busy_absorbed`
478    /// is the count of refusals a caller actually observed.
479    pub writer_task_begin_busy_absorbed: u64,
480    /// Writer-task `BEGIN IMMEDIATE` attempts that failed for a reason other
481    /// than busy or locked, and so surface as `StorageError::Pool`.
482    pub writer_task_begin_errors: u64,
483    /// Dequeued writer-task requests that reached the writer seam (executed
484    /// or attempted to execute their operation) and terminated in error,
485    /// counted once per request regardless of the specific terminal state.
486    pub writer_task_request_failures: u64,
487    /// Subset of `writer_task_request_failures` whose terminal state was
488    /// `WriterTaskRequestState::SideEffectsUnknown` — the commit or rollback
489    /// outcome could not be established, so the request's side effects on
490    /// the database are unknown.
491    pub writer_task_side_effects_unknown: u64,
492    /// Pooled writer guards released while their connection was still inside
493    /// a transaction, which the guard's drop then rolled back (or, when the
494    /// rollback could not prove autocommit, retired). Every writer path settles
495    /// its transaction before releasing the guard, so a non-zero count is a
496    /// caller that left one open; each drop also writes a `writer_guard_drop`
497    /// row to the timeout sink.
498    pub writer_guard_drop_rollbacks: u64,
499}
500
501/// Atomics backing [`WriterAcquisitionSnapshot`]. The writer task retains an
502/// `Arc` after spawn so its per-request acquisition site can update the same
503/// pool-scoped snapshot without retaining the whole pool.
504#[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    /// Bind the lease already owned by this guard to one autocommit write
521    /// unit. The caller must construct the unit immediately before its first
522    /// SQLite write, then keep it through the last call in that unit.
523    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    /// Returns a shared reference to the underlying connection.
540    pub fn conn(&self) -> &Connection {
541        &self.guard
542    }
543
544    /// Returns a mutable reference to the underlying connection.
545    pub fn conn_mut(&mut self) -> &mut Connection {
546        &mut self.guard
547    }
548
549    /// Execute a write transaction.
550    /// Wraps the closure in BEGIN IMMEDIATE ... COMMIT.
551    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
596/// See [`ConnectionPool::migration_transactions`].
597pub(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        // Rolled back: an operation that reported success did not commit.
613        result.and(Err(SqliteError::WriterSettlementUnknown))
614    }
615}
616
617impl Drop for WriterGuard<'_> {
618    fn drop(&mut self) {
619        if !self.guard.is_autocommit() {
620            // Last-resort settlement: the writer that opened this transaction
621            // returned without settling it. Its caller has already been
622            // answered, so the outcome is recorded rather than returned.
623            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                // The mutex excludes all aliases; settlement precedes lease release.
629                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    /// Execute DML on a typed direct transaction, choosing the standalone
682    /// connection required by file-backed async stores or the pooled writer
683    /// used by in-memory stores. The callback cannot begin or commit its own
684    /// outer transaction; this seam owns settlement and the volume lease.
685    #[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    /// Run core migrations while preserving owner-before-lease-before-writer
727    /// order for callers that hold a pool but not a backend wrapper.
728    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    /// Migration writes through this pool's writer, checked out in
750    /// lease-then-writer order for each admitted unit and returned once that
751    /// unit has settled.
752    pub(crate) fn migration_transactions(&self) -> PooledMigrationTransactions<'_> {
753        PooledMigrationTransactions { pool: self }
754    }
755
756    /// Own one standalone write transaction from lease acquisition through
757    /// settlement. The free-space sample occurs only after BEGIN IMMEDIATE
758    /// has acquired SQLite's writer slot.
759    #[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}