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}