Skip to main content

koan_core/db/
pool.rs

1//! Database connections, opened once and kept.
2//!
3//! Opening one costs a permissions syscall, the whole schema DDL and a WAL
4//! checkpoint before a single row comes back, and while downloads are writing
5//! the checkpoint contends with them.
6//!
7//! A pool rather than one shared connection, because rusqlite's `Connection` is
8//! `Send` but not `Sync`: sharing one means a mutex, and a mutex means every
9//! read waits for every other. SQLite in WAL mode reads concurrently across
10//! connections, so the way to keep that is to have several and hand them out.
11//!
12//! A connection is opened only when every existing one is busy, so the pool
13//! grows to whatever concurrency actually happens. What it *keeps* is capped:
14//! queueing a thousand tracks runs as many transfers at once as
15//! `download_workers` allows and no more, but a burst wider than the cap should
16//! not leave a thousand connections parked for the rest of the session, each
17//! holding its own page cache. Past the cap a returned connection is closed
18//! rather than kept.
19
20use std::ops::Deref;
21use std::path::{Path, PathBuf};
22use std::sync::atomic::{AtomicBool, Ordering};
23
24use super::connection::{Database, DbError};
25
26pub struct Pool {
27    path: PathBuf,
28    idle: parking_lot::Mutex<Vec<Database>>,
29    keep: usize,
30    /// Whether the schema has been applied in this process.
31    schema: AtomicBool,
32    /// Held while applying it, so two threads arriving at once do it once.
33    applying: parking_lot::Mutex<()>,
34}
35
36/// How many idle connections to hold on to.
37///
38/// Comfortably more than the concurrency anything here actually reaches — a
39/// handful of front-end reads alongside the configured download workers — so
40/// the steady state never opens one, while a burst still cannot park hundreds.
41const KEEP_IDLE: usize = 32;
42
43/// The process's pool for the default library.
44///
45/// One database, opened once, shared by everything that reads it: the front
46/// ends, the downloader, the background tasks. Callers with their own path —
47/// tests, mostly — build their own.
48pub fn shared() -> &'static Pool {
49    static POOL: std::sync::OnceLock<Pool> = std::sync::OnceLock::new();
50    POOL.get_or_init(|| Pool::new(crate::config::db_path()))
51}
52
53impl Pool {
54    /// The schema must already be applied — see [`Pool::get`].
55    pub fn new(path: PathBuf) -> Self {
56        Self {
57            path,
58            idle: parking_lot::Mutex::new(Vec::new()),
59            keep: KEEP_IDLE,
60            schema: AtomicBool::new(false),
61            applying: parking_lot::Mutex::new(()),
62        }
63    }
64
65    /// Apply the schema, once, before handing out the first connection.
66    ///
67    /// So the pool is safe as the only way anything reaches the database.
68    /// Pooled connections open with `open_existing`, which does no DDL — fine
69    /// for a running app that opened the library at startup, and wrong for a
70    /// one-shot command on a machine that has never run koan. Doing it here
71    /// rather than relying on a caller having done it first removes the
72    /// invariant instead of documenting it.
73    fn ensure_schema(&self) -> Result<(), DbError> {
74        if self.schema.load(Ordering::Acquire) {
75            return Ok(());
76        }
77        let _applying = self.applying.lock();
78        if self.schema.load(Ordering::Acquire) {
79            return Ok(());
80        }
81        Database::open(&self.path)?;
82        self.schema.store(true, Ordering::Release);
83        Ok(())
84    }
85
86    /// Borrow a connection, opening one only if none are free.
87    ///
88    /// `open_existing`, so this never re-runs the DDL or checkpoints: the
89    /// schema is applied once at startup, before any pool exists.
90    pub fn get(&self) -> Result<Handle<'_>, DbError> {
91        self.ensure_schema()?;
92        let pooled = self.idle.lock().pop();
93        let db = match pooled {
94            Some(db) => db,
95            None => Database::open_existing(&self.path)?,
96        };
97        Ok(Handle {
98            db: Some(db),
99            pool: self,
100        })
101    }
102
103    pub fn path(&self) -> &Path {
104        &self.path
105    }
106
107    fn put_back(&self, db: Database) {
108        let mut idle = self.idle.lock();
109        if idle.len() < self.keep {
110            idle.push(db);
111        }
112        // Otherwise it closes here, which is the point of the cap.
113    }
114}
115
116/// A borrowed connection, returned to the pool when it goes out of scope.
117///
118/// Derefs to `Database`, so callers reach `.conn` as on an owned connection.
119pub struct Handle<'a> {
120    db: Option<Database>,
121    pool: &'a Pool,
122}
123
124impl Deref for Handle<'_> {
125    type Target = Database;
126
127    fn deref(&self) -> &Database {
128        self.db.as_ref().expect("a handle holds its connection")
129    }
130}
131
132impl Drop for Handle<'_> {
133    fn drop(&mut self) {
134        if let Some(db) = self.db.take() {
135            self.pool.put_back(db);
136        }
137    }
138}
139
140#[cfg(test)]
141mod tests {
142    use super::*;
143
144    /// A pool over an empty directory — nothing has opened this database, which
145    /// is the state a one-shot command starts from.
146    fn pool() -> (tempfile::TempDir, Pool) {
147        let dir = tempfile::tempdir().unwrap();
148        let path = dir.path().join("koan.db");
149        (dir, Pool::new(path))
150    }
151
152    #[test]
153    fn the_first_connection_applies_the_schema() {
154        // A one-shot command on a machine that has never run koan reaches the
155        // database through here and nowhere else.
156        let (_dir, pool) = pool();
157        let db = pool.get().unwrap();
158        let tracks: i64 = db
159            .conn
160            .query_row("SELECT count(*) FROM tracks", [], |r| r.get(0))
161            .expect("the schema should be there");
162        assert_eq!(tracks, 0);
163    }
164
165    #[test]
166    fn a_connection_comes_back_and_is_reused() {
167        let (_dir, pool) = pool();
168        {
169            let db = pool.get().unwrap();
170            db.conn.execute_batch("SELECT 1").unwrap();
171        }
172        assert_eq!(pool.idle.lock().len(), 1, "returned on drop");
173        {
174            let _db = pool.get().unwrap();
175            assert_eq!(pool.idle.lock().len(), 0, "handed back out");
176        }
177        assert_eq!(pool.idle.lock().len(), 1);
178    }
179
180    #[test]
181    fn concurrent_borrowers_get_their_own() {
182        let (_dir, pool) = pool();
183        let first = pool.get().unwrap();
184        let second = pool.get().unwrap();
185        first.conn.execute_batch("SELECT 1").unwrap();
186        second.conn.execute_batch("SELECT 1").unwrap();
187        drop(first);
188        drop(second);
189        assert_eq!(pool.idle.lock().len(), 2, "both kept for next time");
190    }
191
192    #[test]
193    fn idle_connections_are_capped() {
194        // A burst wider than the cap must not park connections for the rest of
195        // the session.
196        let (_dir, mut pool) = pool();
197        pool.keep = 2;
198        let handles: Vec<_> = (0..6).map(|_| pool.get().unwrap()).collect();
199        assert_eq!(pool.idle.lock().len(), 0, "all of them are out");
200        drop(handles);
201        assert_eq!(pool.idle.lock().len(), 2, "the rest closed on return");
202    }
203
204    #[test]
205    fn a_pooled_connection_sees_what_another_wrote() {
206        // WAL readers are per-connection snapshots; a stale one would serve a
207        // library that is missing whatever just landed.
208        let (_dir, pool) = pool();
209        {
210            let db = pool.get().unwrap();
211            db.conn
212                .execute_batch("CREATE TABLE probe (v INTEGER); INSERT INTO probe VALUES (7)")
213                .unwrap();
214        }
215        let db = pool.get().unwrap();
216        let v: i64 = db
217            .conn
218            .query_row("SELECT v FROM probe", [], |r| r.get(0))
219            .unwrap();
220        assert_eq!(v, 7);
221    }
222}