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}