Skip to main content

koan_core/db/
pool.rs

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