use std::path::Path;
use std::sync::mpsc as std_mpsc;
use std::thread;
use rusqlite::Connection;
use tokio::sync::oneshot;
use crate::error::{Error, Result};
type Job = Box<dyn FnOnce(&Connection) + Send + 'static>;
#[derive(Clone)]
pub struct Store {
tx: std_mpsc::Sender<Job>,
}
impl Store {
pub fn open(path: impl AsRef<Path>) -> Result<Store> {
let path = path.as_ref().to_path_buf();
if let Some(parent) = path.parent() {
let _ = std::fs::create_dir_all(parent);
}
Self::spawn(move || Connection::open(&path))
}
pub fn open_default() -> Result<Store> {
Self::open(crate::paths::db_path())
}
pub fn open_in_memory() -> Result<Store> {
Self::spawn(Connection::open_in_memory)
}
fn spawn<F>(open: F) -> Result<Store>
where
F: FnOnce() -> rusqlite::Result<Connection> + Send + 'static,
{
let (tx, rx) = std_mpsc::channel::<Job>();
let (ready_tx, ready_rx) = std_mpsc::channel::<Result<()>>();
thread::Builder::new()
.name("cliban-store".into())
.spawn(move || {
let conn = match open().map_err(Error::from).and_then(init_connection) {
Ok(c) => {
let _ = ready_tx.send(Ok(()));
c
}
Err(e) => {
let _ = ready_tx.send(Err(e));
return;
}
};
while let Ok(job) = rx.recv() {
job(&conn);
}
})
.expect("spawn cliban-store thread");
match ready_rx.recv() {
Ok(Ok(())) => Ok(Store { tx }),
Ok(Err(e)) => Err(e),
Err(_) => Err(Error::WriterGone),
}
}
pub async fn call<F, T>(&self, f: F) -> Result<T>
where
F: FnOnce(&Connection) -> Result<T> + Send + 'static,
T: Send + 'static,
{
let (reply_tx, reply_rx) = oneshot::channel::<Result<T>>();
let job: Job = Box::new(move |conn| {
let out = f(conn);
let _ = reply_tx.send(out);
});
self.tx.send(job).map_err(|_| Error::WriterGone)?;
reply_rx.await.map_err(|_| Error::WriterGone)?
}
}
fn init_connection(conn: Connection) -> Result<Connection> {
conn.pragma_update(None, "journal_mode", "WAL")?;
conn.pragma_update(None, "foreign_keys", "ON")?;
conn.pragma_update(None, "busy_timeout", 5000)?;
crate::migrations::run(&conn)?;
Ok(conn)
}