Skip to main content

ConnectionPool

Struct ConnectionPool 

Source
pub struct ConnectionPool { /* private fields */ }
Expand description

A read-write connection pool for SQLite.

Architecture:

  • 1 writer connection protected by a Mutex (exclusive access)
  • N reader connections in a lock-free queue (concurrent access)
  • All connections share the same database file in WAL mode

Writable in-memory databases, or writable file databases when WAL mode is disabled/unavailable, degrade to single-connection mode and route all operations through the writer connection. A file-backed read-only pool always retains at least one dedicated read-only connection: rollback-journal snapshots do not need WAL to support concurrent readers, and inspection must never alias a read onto the query-only writer slot.

Implementations§

Source§

impl ConnectionPool

Source

pub fn new(config: PoolConfig) -> Result<ConnectionPool, SqliteError>

Create a new connection pool.

Opens 1 writer + N reader connections to the same database when pooling is enabled. All connections are configured consistently (busy timeout, foreign keys, cache, mmap, temp store). Writable in-memory databases and writable non-WAL files fall back to single-connection mode. Read-only files retain a dedicated reader regardless of journal mode.

Source

pub fn reader(&self) -> Result<ReaderGuard<'_>, SqliteError>

Check out a reader connection.

Tries to pop from the lock-free queue. If empty, spins briefly then waits with exponential backoff up to checkout_timeout.

In degraded mode (WAL unavailable, max_readers == 0), this method checks the shared writer mutex in bounded slices and returns pool exhaustion after checkout_timeout; it never blocks indefinitely on the non-reentrant mutex.

Source

pub fn writer(&self) -> Result<WriterGuard<'_>, SqliteError>

Check out the writer connection.

Waits up to checkout_timeout for the writer Mutex and returns Err(SqliteError::WriterPoolCheckoutTimeout) if the timeout is exceeded.

Source

pub fn try_writer(&self) -> Result<WriterGuard<'_>, SqliteError>

Non-panicking writer checkout.

Returns Err on timeout instead of panicking. Use this in request handlers where a 500 is preferable to crashing the process.

Source

pub fn try_writer_nowait(&self) -> Result<WriterGuard<'_>, SqliteError>

Zero-wait writer checkout for background tasks.

Uses try_lock() (no timeout, no spin) — returns Err immediately when any other caller holds the writer Mutex. Background tasks (e.g. the WAL checkpoint task) MUST use this instead of try_writer so that a busy writer causes the background task to skip its current tick rather than stalling for up to checkout_timeout (default 5s) while write traffic is in progress.

Source

pub fn writer_acquisition_snapshot(&self) -> WriterAcquisitionSnapshot

Snapshot all instrumented writer acquisition outcomes since this pool was constructed.

Source

pub fn reader_acquisition_snapshot(&self) -> ReaderAcquisitionSnapshot

Snapshot reader acquisition, saturation, and hold lifecycle outcomes since this pool was constructed. Counters reset only with pool reconstruction.

Source

pub fn available_readers(&self) -> usize

Get the current number of available reader connections.

Source

pub fn max_readers(&self) -> usize

Get the total number of reader connections in the pool.

Source

pub fn config(&self) -> &PoolConfig

Return the pool configuration.

Source

pub fn main_pool_generation(&self) -> u64

Identify this pool’s counter window when it is designated as main. Repeated runtime handles and diagnostics reads reuse the same generation; constructing secondary pools does not consume main-pool generations.

Source

pub fn origin(&self) -> TxOrigin

This pool’s ADR-091 backend-scoped attribution origin (ADR-091, backend-scoped WAL-pin attribution design note): Database(_) for a file-backed pool, Memory for an in-memory pool. Every tx_registry::register_scoped call site threaded in this crate passes this value as the span’s origin.

Source

pub fn canonical_path(&self) -> Option<&Path>

The canonical path this pool’s origin() identity was minted from, None for an in-memory pool. DbIdentity has no path accessor by design; sidecar derivation and other filesystem consumers use this — the same canonical value the identity was minted from — instead of re-deriving a path from the raw configured one.

Source

pub fn write_queue_active(&self) -> bool

Whether the write queue is effectively enabled for this pool: the resolved write_queue_enabled flag AND file-backed.

ConnectionPool::new resolves the “no preference” (None) preference to a concrete Some(..) once path is known, so every reader of config.write_queue_enabled sees a resolved value; the debug_assert pins that invariant and a None that slipped past would read as disabled. Bypassing ConnectionPool::new to construct a pool is a construction-path bug. Use this instead of repeating config().write_queue_enabled.unwrap_or(false) && config().path.is_some() at every routing/violation site.

Source

pub fn writer_task_join_was_stored(&self) -> bool

Whether a writer-task JoinHandle has been stored at least once.

Unlike Self::take_writer_task_join, this remains true after the one-shot handle slot is emptied, distinguishing a task that never spawned from a handle another caller already consumed.

Source

pub fn writer_task_handle( &self, ) -> Result<Option<WriterTaskHandle>, StorageError>

Return the pool-wide ADR-067 Component A writer task, spawning it lazily on first access if PoolConfig::write_queue_enabled is set. Exactly one writer task exists per ConnectionPool (per DB file); see crates/khive-db/docs/api/pool.md#connectionpoolwriter_task_handle–single-writer-task-rationale for why a per-store writer task would defeat the single-writer guarantee.

Returns Ok(None) if the flag is off, or if the writer task failed to spawn for a reason other than a missing runtime (for example, an in-memory pool has no standalone-connection support) — callers fall back to the legacy pool-mutex write path in either case. A spawn failure is logged once here (at first access), not once per store.

Returns Err(StorageError::WriterTaskNoRuntime) instead of panicking when write_queue_enabled is set but this is the first access and no Tokio runtime is available on the calling thread (checked via tokio::runtime::Handle::try_current) — spawning the writer task requires tokio::spawn, which panics outside a runtime. Callers that already treat a missing writer task as best-effort (construction-time degrade to the legacy path, matching slice 1’s documented policy) can collapse this into None with .ok().flatten(); callers that need to fail loud on a genuine misconfiguration (write queue requested but no runtime to run it on) can propagate the Err directly.

Source

pub fn writer_task_for_runtime_write( &self, operation: RuntimeWriteOperation, ) -> Result<Option<WriterTaskHandle>, StorageError>

Resolve a runtime-owned transaction through the same strict/compatibility policy as store writes. None permits the caller’s direct transaction and records its compatibility fallback when the file-backed queue is enabled. A strict refusal returns before any direct-writer acquisition or telemetry.

Source

pub fn take_writer_task_join(&self) -> Option<JoinHandle<()>>

Take the writer task’s JoinHandle, if a writer task was spawned and the handle has not already been taken.

Intended for short-lived batch callers that drop every WriterTaskHandle clone (closing the queue) and then need to await the task’s exit before treating the database file as settled: the task’s connection close fires SQLite’s close-time WAL checkpoint, so until the task exits the file bytes can still move after the caller’s last write returned.

One-shot: None means either the write queue never spawned (disabled, or spawn degraded) or another caller already took the handle — in both cases there is nothing further to await here. Exactly one subsystem may own the drain: the single caller that receives Some(_) is the sole owner of the task-exit await (and of the close-time WAL checkpoint that settles the database file); every later caller receives None and must not arrange its own await.

Source

pub fn legacy_conn(&self) -> Arc<Mutex<RawMutex, Connection>> ⓘ

Compatibility method: returns the writer connection wrapped in Arc<Mutex>.

WARNING: This exists only for backward compatibility with code that calls store.conn(). New code should use reader() and writer().

Source

pub fn open_standalone_writer(&self) -> Result<Connection, SqliteError>

Open a standalone read-write connection to the same file-backed database.

Stores whose trait methods take Send + 'static closures (executed via spawn_blocking) cannot hold the pooled WriterGuard’s MutexGuard across the call — it opens an independent connection instead. This must still honor PoolConfig::read_only: opening SQLITE_OPEN_READ_WRITE unconditionally here would let a read-only backend’s graph/event/text stores bypass the flag that the pooled writer enforces via query_only. A fully configured successful open increments the standalone acquisition class exactly once.

Source

pub fn claim_checkpoint_ownership(&self) -> Result<(), SqliteError>

Claim routine WAL-checkpoint ownership for this pool.

Called by the scheduled checkpoint task at startup — the one caller that actually replaces SQLite’s per-commit autocheckpoint with dedicated PASSIVE checkpointing (ADR-091 Amendment 10). The claim makes every subsequently opened writer-capable connection set PRAGMA wal_autocheckpoint = 0, and re-applies that pragma on the already-open pooled writer under the writer mutex. A writer task spawned before the claim keeps its own long-lived connection; Self::propagate_checkpoint_claim_to_writer_task reaches that one.

Without a claim, writer-capable connections keep the bounded FALLBACK_WAL_AUTOCHECKPOINT_PAGES threshold, so a writable pool in a process that never runs the checkpoint task (embedded runtimes, one-shot CLI executions) retains SQLite’s own WAL reclamation instead of growing its WAL without bound.

Read-only pools record the claim but have no writer-capable connections to reconfigure. Writable pools publish the claim only after the pooled writer is configured successfully; a failed attempt keeps the bounded fallback active and remains retryable.

Source

pub async fn propagate_checkpoint_claim_to_writer_task( &self, ) -> Result<(), StorageError>

Flip an already-running writer task’s long-lived connection to the claimed-owner setting.

Connections opened after Self::claim_checkpoint_ownership inherit wal_autocheckpoint = 0 at open; only a writer task spawned before the claim still holds a connection on the bounded fallback. Returns Ok(()) without side effects when the pool’s write queue is disabled.

Trait Implementations§

Source§

impl Drop for ConnectionPool

Source§

fn drop(&mut self)

Close every read-only reader before the fields below it drop in declaration order (writer first, readers well before writer_task). A read-only connection cannot take the EXCLUSIVE lock SQLite needs to checkpoint on close, so if a reader were left to close last, WAL mode would leave -wal/-shm behind. Draining readers here, before that field-order drop runs, makes the writable writer connection close after every reader instead of before it.

Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more