Skip to main content

cliban_core/
store.rs

1//! The store: a [`Store`] handle (Clone + Send + Sync) over a single writer
2//! thread that owns the one rusqlite [`Connection`].
3//!
4//! The handle holds an mpsc sender; a background worker owns the resource and
5//! serves jobs one at a time. SQLite is a single writer regardless, so
6//! funnelling every read and write through one connection is both correct and
7//! the simplest thing that preserves the Elixir `Repo.transaction` semantics —
8//! each job runs to completion before the next starts, so a context function
9//! that opens a transaction has the connection entirely to itself.
10//!
11//! rusqlite is blocking and `Connection` is `!Sync`, so the worker is a
12//! dedicated OS thread (not a tokio task). Async callers submit a closure +
13//! await a `oneshot`; the closure runs on the worker thread with `&Connection`
14//! and returns a value back through the channel. WAL is enabled on open.
15
16use std::path::Path;
17use std::sync::mpsc as std_mpsc;
18use std::thread;
19
20use rusqlite::Connection;
21use tokio::sync::oneshot;
22
23use crate::error::{Error, Result};
24
25/// A unit of work for the writer thread: a boxed closure given the live
26/// connection. We erase the return type into the closure itself (it owns its
27/// own oneshot sender), so the worker loop stays monomorphic.
28type Job = Box<dyn FnOnce(&Connection) + Send + 'static>;
29
30/// Clone-able handle to the store. Cheap to clone (just an `mpsc::Sender`);
31/// safe to share across tasks and threads. Dropping all clones shuts the
32/// worker down.
33#[derive(Clone)]
34pub struct Store {
35    tx: std_mpsc::Sender<Job>,
36}
37
38impl Store {
39    /// Open (or create) the DB at `path`, run migrations, enable WAL, and spawn
40    /// the writer thread. Returns once the worker is ready (the open +
41    /// migration happen synchronously on the worker and the result is awaited).
42    pub fn open(path: impl AsRef<Path>) -> Result<Store> {
43        let path = path.as_ref().to_path_buf();
44        if let Some(parent) = path.parent() {
45            // Make the data dir if the DB lives somewhere that doesn't exist
46            // yet. Ignore failures here; the open below will surface a real
47            // error.
48            let _ = std::fs::create_dir_all(parent);
49        }
50        Self::spawn(move || Connection::open(&path))
51    }
52
53    /// Open the store at the default [`crate::paths::db_path`] location. The CLI
54    /// and TUI entry points will call this; tests use
55    /// [`Store::open_in_memory`] instead.
56    pub fn open_default() -> Result<Store> {
57        Self::open(crate::paths::db_path())
58    }
59
60    /// Open an in-memory store. Test convenience; the contents vanish on drop.
61    pub fn open_in_memory() -> Result<Store> {
62        Self::spawn(Connection::open_in_memory)
63    }
64
65    fn spawn<F>(open: F) -> Result<Store>
66    where
67        F: FnOnce() -> rusqlite::Result<Connection> + Send + 'static,
68    {
69        let (tx, rx) = std_mpsc::channel::<Job>();
70        let (ready_tx, ready_rx) = std_mpsc::channel::<Result<()>>();
71
72        thread::Builder::new()
73            .name("cliban-store".into())
74            .spawn(move || {
75                let conn = match open().map_err(Error::from).and_then(init_connection) {
76                    Ok(c) => {
77                        let _ = ready_tx.send(Ok(()));
78                        c
79                    }
80                    Err(e) => {
81                        let _ = ready_tx.send(Err(e));
82                        return;
83                    }
84                };
85
86                // Serve jobs until every handle is dropped.
87                while let Ok(job) = rx.recv() {
88                    job(&conn);
89                }
90            })
91            .expect("spawn cliban-store thread");
92
93        match ready_rx.recv() {
94            Ok(Ok(())) => Ok(Store { tx }),
95            Ok(Err(e)) => Err(e),
96            Err(_) => Err(Error::WriterGone),
97        }
98    }
99
100    /// Submit `f` to run on the writer thread with exclusive access to the
101    /// connection, and await its result. This is the single primitive every
102    /// context method is built on; a context function that needs a transaction
103    /// simply opens one inside `f` (it has the connection to itself for the
104    /// duration of the call).
105    pub async fn call<F, T>(&self, f: F) -> Result<T>
106    where
107        F: FnOnce(&Connection) -> Result<T> + Send + 'static,
108        T: Send + 'static,
109    {
110        let (reply_tx, reply_rx) = oneshot::channel::<Result<T>>();
111        let job: Job = Box::new(move |conn| {
112            let out = f(conn);
113            let _ = reply_tx.send(out);
114        });
115        self.tx.send(job).map_err(|_| Error::WriterGone)?;
116        reply_rx.await.map_err(|_| Error::WriterGone)?
117    }
118}
119
120/// Per-connection setup, run once on open. WAL for concurrent readers, foreign
121/// keys ON (Ecto runs with `PRAGMA foreign_key = ON` per-connection;
122/// ecto_sqlite3 enables it by default), busy_timeout so a momentarily-locked
123/// DB retries instead of erroring, and the migration baseline.
124fn init_connection(conn: Connection) -> Result<Connection> {
125    conn.pragma_update(None, "journal_mode", "WAL")?;
126    conn.pragma_update(None, "foreign_keys", "ON")?;
127    conn.pragma_update(None, "busy_timeout", 5000)?;
128    crate::migrations::run(&conn)?;
129    Ok(conn)
130}