Skip to main content

miden_node_db/sqlite/
pool.rs

1//! Async connection pool over raw `rusqlite`.
2//!
3//! SQLite permits only a single writer at a time, so the pool is split into a **single** writer
4//! connection and a pool of read-only connections. Writes (`write`/`begin_write`) serialize on the
5//! one writer; reads (`read`/`begin_read`) run concurrently on the reader pool. This makes the
6//! single-writer model structural (rather than relying on lock contention) and lets a held write
7//! transaction stay open without starving readers.
8
9use std::num::NonZeroUsize;
10use std::path::{Path, PathBuf};
11
12use deadpool::Runtime;
13use deadpool::managed::{Manager, Metrics, Object, Pool, RecycleError, RecycleResult};
14use deadpool_sync::SyncWrapper;
15use miden_node_tracing::Instrument;
16use rusqlite::{Connection, OpenFlags, TransactionBehavior};
17
18use crate::sqlite::tx::{ReadTx, WriteTx};
19use crate::{DatabaseError, default_connection_pool_size};
20
21/// Per-connection prepared-statement cache capacity. Raised well above rusqlite's default of 16
22/// because we keep a large set of distinct statements; the bounded connection pools cap total
23/// cached-statement memory.
24const STATEMENT_CACHE_CAPACITY: usize = 512;
25
26// CONNECTION MANAGER
27// =================================================================================================
28
29/// Errors raised while creating or recycling a pooled connection.
30///
31/// Internal to the pool: callers only ever observe a [`DatabaseError`] (pool failures are boxed into
32/// [`DatabaseError::ConnectionPoolObtainError`]), so this type is not part of the public API.
33#[derive(Debug, thiserror::Error)]
34pub(crate) enum SqliteManagerError {
35    /// Opening the database file failed.
36    #[error("failed to open the sqlite database")]
37    Open(#[source] rusqlite::Error),
38    /// Applying the per-connection PRAGMAs failed.
39    #[error("failed to configure the sqlite connection")]
40    Configure(#[source] rusqlite::Error),
41    /// The pooled connection's mutex was poisoned by a panic during a previous interaction.
42    #[error("the pooled sqlite connection is poisoned")]
43    Poisoned,
44}
45
46struct SqliteManager {
47    path: PathBuf,
48    /// When set, connections are configured `PRAGMA query_only = ON` and skip the writer-only
49    /// `journal_mode` setup — used for the reader pool.
50    read_only: bool,
51}
52
53impl Manager for SqliteManager {
54    type Type = SyncWrapper<Connection>;
55    type Error = SqliteManagerError;
56
57    async fn create(&self) -> Result<Self::Type, Self::Error> {
58        let path = self.path.clone();
59        let read_only = self.read_only;
60        SyncWrapper::new(Runtime::Tokio1, move || {
61            let conn = Connection::open_with_flags(&path, OpenFlags::SQLITE_OPEN_READ_WRITE)
62                .map_err(SqliteManagerError::Open)?;
63            configure_connection(&conn, read_only).map_err(SqliteManagerError::Configure)?;
64            Ok(conn)
65        })
66        .await
67    }
68
69    async fn recycle(
70        &self,
71        conn: &mut Self::Type,
72        _metrics: &Metrics,
73    ) -> RecycleResult<Self::Error> {
74        if conn.is_mutex_poisoned() {
75            return Err(RecycleError::Backend(SqliteManagerError::Poisoned));
76        }
77        // Safety net for a held transaction handle dropped without `commit`/`rollback`: roll back
78        // any still-open transaction so the next user gets a clean connection.
79        conn.interact(|conn| {
80            if !conn.is_autocommit() {
81                let _ = conn.execute_batch("ROLLBACK");
82            }
83        })
84        .await
85        .map_err(|_| RecycleError::Backend(SqliteManagerError::Poisoned))?;
86        Ok(())
87    }
88}
89
90/// Applies the per-connection PRAGMAs and statement-cache sizing.
91///
92/// Both pools open the file `READ_WRITE`; reader connections are made read-only at runtime with
93/// `PRAGMA query_only = ON` (which, unlike opening `READ_ONLY`, still lets them create the WAL
94/// `-shm` file and read a WAL database).
95fn configure_connection(conn: &Connection, read_only: bool) -> rusqlite::Result<()> {
96    // busy_timeout makes concurrent writers wait instead of failing immediately; foreign keys
97    // enforce referential integrity.
98    if read_only {
99        // A query_only connection cannot set `journal_mode` (it is a write); WAL is already
100        // persisted in the file header by the writer / migration path.
101        conn.execute_batch(
102            "PRAGMA busy_timeout = 5000;
103             PRAGMA foreign_keys = ON;
104             PRAGMA query_only = ON;",
105        )?;
106    } else {
107        // WAL allows concurrent readers while the writer holds the lock.
108        conn.execute_batch(
109            "PRAGMA busy_timeout = 5000;
110             PRAGMA journal_mode = WAL;
111             PRAGMA foreign_keys = ON;",
112        )?;
113    }
114    conn.set_prepared_statement_cache_capacity(STATEMENT_CACHE_CAPACITY);
115    // Register the `array` extension so the cacheable IN-list helpers can bind lists via
116    // `rarray(?)` (see `crate::sqlite::in_list`).
117    rusqlite::vtab::array::load_module(conn)?;
118    Ok(())
119}
120
121// DATABASE HANDLES
122// =================================================================================================
123
124/// Opens a database over `database_filepath` with the default reader-pool size.
125///
126/// Returns a `(DbWriter, DbReader)` pair (see [`open_with_pool_size`]).
127pub fn open(database_filepath: &Path) -> Result<(DbWriter, DbReader), DatabaseError> {
128    open_with_pool_size(database_filepath, default_connection_pool_size())
129}
130
131/// Opens a database over `database_filepath` with the given reader-pool size. The writer is always a
132/// single connection.
133///
134/// Returns the write handle and the reader handle as **separate** values rather than one shared
135/// object, so read-only and write access are distinct, non-interchangeable capabilities: a component
136/// handed only a [`DbReader`] has no way to write. The [`DbWriter`] is not `Clone`, so there is at
137/// most one writer owner (mirroring the `max_size(1)` writer pool).
138pub fn open_with_pool_size(
139    database_filepath: &Path,
140    connection_pool_size: NonZeroUsize,
141) -> Result<(DbWriter, DbReader), DatabaseError> {
142    let writer = Pool::builder(SqliteManager {
143        path: database_filepath.to_path_buf(),
144        read_only: false,
145    })
146    .max_size(1)
147    .build()?;
148    let readers = Pool::builder(SqliteManager {
149        path: database_filepath.to_path_buf(),
150        read_only: true,
151    })
152    .max_size(connection_pool_size.get())
153    .build()?;
154    Ok((DbWriter { writer }, DbReader { readers }))
155}
156
157/// Read-only handle over the reader connection pool. Cloning shares the underlying pool.
158///
159/// Exposes only read access ([`read`](Self::read)/[`begin_read`](Self::begin_read)); it has no way to
160/// mutate the database. Hand this to components that must never write.
161#[derive(Clone)]
162pub struct DbReader {
163    readers: Pool<SqliteManager>,
164}
165
166impl DbReader {
167    /// Checks a reader connection out of the pool.
168    async fn checkout_reader(&self) -> Result<Object<SqliteManager>, DatabaseError> {
169        self.readers
170            .get()
171            .in_current_span()
172            .await
173            .map_err(|err| DatabaseError::ConnectionPoolObtainError(Box::new(err)))
174    }
175
176    /// Runs `query` inside a read-only (`DEFERRED`, never committed) transaction on a reader
177    /// connection.
178    pub async fn read<R, E, F>(&self, msg: impl ToString + Send, query: F) -> Result<R, E>
179    where
180        F: FnOnce(&ReadTx<'_>) -> Result<R, E> + Send + 'static,
181        R: Send + 'static,
182        E: From<DatabaseError> + Send + 'static,
183    {
184        let conn = self.checkout_reader().await.map_err(E::from)?;
185        let msg = msg.to_string();
186        let span = miden_node_tracing::Span::current();
187        conn.interact(move |conn| {
188            let _guard = span.enter();
189            let tx = conn
190                .transaction_with_behavior(TransactionBehavior::Deferred)
191                .map_err(|err| E::from(DatabaseError::from(err)))?;
192            query(&ReadTx::new(&tx))
193            // `tx` is dropped here without a commit, rolling back any writes.
194        })
195        .await
196        .map_err(|err| E::from(DatabaseError::interact(&msg, &err)))?
197    }
198
199    /// Begins a read-only (`DEFERRED`) transaction on a reader connection and returns a handle held
200    /// across `.await` points. See [`ReadTransaction`].
201    pub async fn begin_read(&self) -> Result<ReadTransaction, DatabaseError> {
202        let conn = self.checkout_reader().await?;
203        run_tx_stmt(&conn, "BEGIN DEFERRED").await?;
204        Ok(ReadTransaction { conn })
205    }
206}
207
208/// Write handle over the single writer connection.
209///
210/// **Not `Clone`**: SQLite permits only one writer, so at most one owner holds write access at a time
211/// (the writer pool is `max_size(1)`). Exposes only write access
212/// ([`write`](Self::write)/[`begin_write`](Self::begin_write)); pair it with a [`DbReader`] when the
213/// owner also needs to read.
214pub struct DbWriter {
215    writer: Pool<SqliteManager>,
216}
217
218impl DbWriter {
219    /// Checks the single writer connection out of the pool.
220    async fn checkout_writer(&self) -> Result<Object<SqliteManager>, DatabaseError> {
221        self.writer
222            .get()
223            .in_current_span()
224            .await
225            .map_err(|err| DatabaseError::ConnectionPoolObtainError(Box::new(err)))
226    }
227
228    /// Runs `query` inside a read-write (`IMMEDIATE`) transaction on the single writer connection,
229    /// committing on `Ok`.
230    pub async fn write<R, E, F>(&self, msg: impl ToString + Send, query: F) -> Result<R, E>
231    where
232        F: FnOnce(&WriteTx<'_>) -> Result<R, E> + Send + 'static,
233        R: Send + 'static,
234        E: From<DatabaseError> + Send + 'static,
235    {
236        let conn = self.checkout_writer().await.map_err(E::from)?;
237        let msg = msg.to_string();
238        let span = miden_node_tracing::Span::current();
239        conn.interact(move |conn| {
240            let _guard = span.enter();
241            let tx = conn
242                .transaction_with_behavior(TransactionBehavior::Immediate)
243                .map_err(|err| E::from(DatabaseError::from(err)))?;
244            let result = query(&WriteTx::new(&tx))?;
245            tx.commit().map_err(|err| E::from(DatabaseError::from(err)))?;
246            Ok(result)
247        })
248        .await
249        .map_err(|err| E::from(DatabaseError::interact(&msg, &err)))?
250    }
251
252    /// Begins a read-write (`IMMEDIATE`) transaction on the single writer connection and returns a
253    /// handle held across `.await` points. The handle must be committed (or it rolls back). See
254    /// [`WriteTransaction`].
255    pub async fn begin_write(&self) -> Result<WriteTransaction, DatabaseError> {
256        let conn = self.checkout_writer().await?;
257        run_tx_stmt(&conn, "BEGIN IMMEDIATE").await?;
258        Ok(WriteTransaction { conn })
259    }
260}
261
262// HELD TRANSACTIONS
263// =================================================================================================
264
265/// Runs a transaction-control statement (`BEGIN`/`COMMIT`/`ROLLBACK`) on a checked-out connection.
266async fn run_tx_stmt(
267    conn: &Object<SqliteManager>,
268    stmt: &'static str,
269) -> Result<(), DatabaseError> {
270    conn.interact(move |conn| conn.execute_batch(stmt))
271        .await
272        .map_err(|err| DatabaseError::interact(stmt, &err))?
273        .map_err(DatabaseError::from)
274}
275
276/// A read transaction (`DEFERRED`) held across `.await` points, on a reader connection.
277///
278/// Run batches of synchronous queries with [`run`](Self::run); the transaction stays open between
279/// calls, so a request handler can interleave queries with async work on a single consistent
280/// snapshot. The transaction is read-only and ends (rolls back) when the handle is dropped, or
281/// explicitly via [`close`](Self::close).
282pub struct ReadTransaction {
283    conn: Object<SqliteManager>,
284}
285
286impl ReadTransaction {
287    /// Runs a batch of read queries against the open transaction.
288    pub async fn run<R, E, F>(&self, msg: impl ToString + Send, query: F) -> Result<R, E>
289    where
290        F: FnOnce(&ReadTx<'_>) -> Result<R, E> + Send + 'static,
291        R: Send + 'static,
292        E: From<DatabaseError> + Send + 'static,
293    {
294        let msg = msg.to_string();
295        let span = miden_node_tracing::Span::current();
296        self.conn
297            .interact(move |conn| {
298                let _guard = span.enter();
299                query(&ReadTx::new(conn))
300            })
301            .await
302            .map_err(|err| E::from(DatabaseError::interact(&msg, &err)))?
303    }
304
305    /// Ends the transaction explicitly (rolls back; a read transaction has nothing to commit).
306    pub async fn close(self) -> Result<(), DatabaseError> {
307        run_tx_stmt(&self.conn, "ROLLBACK").await
308    }
309}
310
311/// A read-write transaction (`IMMEDIATE`) held across `.await` points, on the single writer
312/// connection.
313///
314/// Run batches of synchronous queries with [`run`](Self::run); the transaction stays open between
315/// calls, so a request handler can interleave reads and writes with async work atomically. Finish
316/// with [`commit`](Self::commit) to persist, or [`rollback`](Self::rollback) to discard; if the
317/// handle is dropped without either, the pool rolls the transaction back when the connection is
318/// recycled.
319///
320/// The handle holds the sole writer connection for its whole lifetime.
321pub struct WriteTransaction {
322    conn: Object<SqliteManager>,
323}
324
325impl WriteTransaction {
326    /// Runs a batch of read/write queries against the open transaction.
327    pub async fn run<R, E, F>(&self, msg: impl ToString + Send, query: F) -> Result<R, E>
328    where
329        F: FnOnce(&WriteTx<'_>) -> Result<R, E> + Send + 'static,
330        R: Send + 'static,
331        E: From<DatabaseError> + Send + 'static,
332    {
333        let msg = msg.to_string();
334        let span = miden_node_tracing::Span::current();
335        self.conn
336            .interact(move |conn| {
337                let _guard = span.enter();
338                query(&WriteTx::new(conn))
339            })
340            .await
341            .map_err(|err| E::from(DatabaseError::interact(&msg, &err)))?
342    }
343
344    /// Commits the transaction, persisting all writes.
345    pub async fn commit(self) -> Result<(), DatabaseError> {
346        run_tx_stmt(&self.conn, "COMMIT").await
347    }
348
349    /// Rolls back the transaction, discarding all writes.
350    pub async fn rollback(self) -> Result<(), DatabaseError> {
351        run_tx_stmt(&self.conn, "ROLLBACK").await
352    }
353}
354
355#[cfg(test)]
356mod tests {
357    use std::num::NonZeroUsize;
358    use std::path::{Path, PathBuf};
359
360    use rusqlite::Connection;
361
362    use super::{DbReader, DbWriter, open_with_pool_size};
363    use crate::DatabaseError;
364
365    /// A throwaway file-backed database; the pools open existing files `READ_WRITE` only, so the
366    /// file and schema are created up front.
367    struct TempDb {
368        path: PathBuf,
369    }
370
371    impl TempDb {
372        fn new(name: &str) -> Self {
373            let path = std::env::temp_dir()
374                .join(format!("miden-node-db-pool-{name}-{}.sqlite3", std::process::id()));
375            let db = Self { path };
376            db.remove_files();
377            let conn = Connection::open(&db.path).expect("create db file");
378            conn.execute_batch("CREATE TABLE items (id INTEGER PRIMARY KEY);")
379                .expect("create table");
380            db
381        }
382
383        fn path(&self) -> &Path {
384            &self.path
385        }
386
387        fn remove_files(&self) {
388            let _ = fs_err::remove_file(&self.path);
389            let _ = fs_err::remove_file(self.path.with_extension("sqlite3-wal"));
390            let _ = fs_err::remove_file(self.path.with_extension("sqlite3-shm"));
391        }
392    }
393
394    impl Drop for TempDb {
395        fn drop(&mut self) {
396            self.remove_files();
397        }
398    }
399
400    fn open_db(temp: &TempDb) -> (DbWriter, DbReader) {
401        open_with_pool_size(temp.path(), NonZeroUsize::new(4).unwrap()).unwrap()
402    }
403
404    async fn count_items(reader: &DbReader) -> i64 {
405        reader
406            .read::<_, DatabaseError, _>("count", |r| {
407                Ok(r.query("SELECT COUNT(*) FROM items", &[], |row| row.get::<i64>(0))?
408                    .into_iter()
409                    .next()
410                    .unwrap_or(0))
411            })
412            .await
413            .unwrap()
414    }
415
416    async fn insert_committed(writer: &DbWriter, id: i64) {
417        let tx = writer.begin_write().await.unwrap();
418        tx.run::<_, DatabaseError, _>("insert", move |w| {
419            w.execute("INSERT INTO items (id) VALUES (?1)", &[&id])?;
420            Ok(())
421        })
422        .await
423        .unwrap();
424        tx.commit().await.unwrap();
425    }
426
427    #[tokio::test]
428    async fn held_write_transaction_commits_across_awaits() {
429        let temp = TempDb::new("commit");
430        let (writer, reader) = open_db(&temp);
431
432        let tx = writer.begin_write().await.unwrap();
433        tx.run::<_, DatabaseError, _>("insert-1", |w| {
434            w.execute("INSERT INTO items (id) VALUES (?1)", &[&1i64])?;
435            Ok(())
436        })
437        .await
438        .unwrap();
439
440        // Interleave async work between statements on the same still-open transaction.
441        tokio::task::yield_now().await;
442
443        tx.run::<_, DatabaseError, _>("insert-2", |w| {
444            w.execute("INSERT INTO items (id) VALUES (?1)", &[&2i64])?;
445            Ok(())
446        })
447        .await
448        .unwrap();
449
450        tx.commit().await.unwrap();
451
452        assert_eq!(count_items(&reader).await, 2);
453    }
454
455    #[tokio::test]
456    async fn dropped_write_transaction_rolls_back() {
457        let temp = TempDb::new("rollback");
458        let (writer, reader) = open_db(&temp);
459
460        {
461            let tx = writer.begin_write().await.unwrap();
462            tx.run::<_, DatabaseError, _>("insert", |w| {
463                w.execute("INSERT INTO items (id) VALUES (?1)", &[&1i64])?;
464                Ok(())
465            })
466            .await
467            .unwrap();
468            // `tx` is dropped here without a commit.
469        }
470
471        // The sole writer connection is reused; `recycle` must have rolled back the orphaned
472        // transaction, otherwise this `BEGIN IMMEDIATE` would fail with "cannot start a transaction
473        // within a transaction". The first insert must not have persisted.
474        insert_committed(&writer, 2).await;
475        assert_eq!(count_items(&reader).await, 1);
476    }
477
478    #[tokio::test]
479    async fn reads_proceed_while_write_transaction_is_held() {
480        let temp = TempDb::new("concurrent");
481        let (writer, reader) = open_db(&temp);
482        insert_committed(&writer, 1).await;
483
484        // Hold an open write transaction with an uncommitted insert.
485        let tx = writer.begin_write().await.unwrap();
486        tx.run::<_, DatabaseError, _>("insert-uncommitted", |w| {
487            w.execute("INSERT INTO items (id) VALUES (?1)", &[&2i64])?;
488            Ok(())
489        })
490        .await
491        .unwrap();
492
493        // A read on the reader pool proceeds (does not block on the writer) and does not see the
494        // uncommitted row.
495        assert_eq!(count_items(&reader).await, 1);
496
497        tx.commit().await.unwrap();
498        assert_eq!(count_items(&reader).await, 2);
499    }
500
501    #[tokio::test]
502    async fn reader_connections_are_query_only() {
503        let temp = TempDb::new("query_only");
504        let (_writer, reader) = open_db(&temp);
505
506        let query_only = reader
507            .read::<_, DatabaseError, _>("pragma", |r| {
508                Ok(r.query("PRAGMA query_only", &[], |row| row.get::<i64>(0))?
509                    .into_iter()
510                    .next()
511                    .unwrap_or(0))
512            })
513            .await
514            .unwrap();
515        assert_eq!(query_only, 1, "reader connections must be query_only");
516
517        // A write attempted on a reader connection is rejected.
518        let result = reader
519            .read::<(), DatabaseError, _>("rejected-write", |r| {
520                r.query("INSERT INTO items (id) VALUES (99)", &[], |_| Ok(()))?;
521                Ok(())
522            })
523            .await;
524        assert!(result.is_err(), "writes on a reader connection must fail");
525    }
526}