Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion rust/crates/truapi/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,7 @@ uniffi = { workspace = true, optional = true }
subxt = { workspace = true, features = ["native"], optional = true }
subxt-rpcs = { workspace = true, features = ["jsonrpsee", "native"], optional = true }
base64 = { workspace = true, optional = true }
rusqlite = { workspace = true, features = ["bundled"], optional = true }
rusqlite = { workspace = true, features = ["bundled", "hooks"], optional = true }
async-sqlite = { workspace = true, features = ["bundled"], optional = true }
rusqlite_migration = { workspace = true, optional = true }

Expand Down
57 changes: 48 additions & 9 deletions rust/crates/truapi/src/store.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,21 @@
//! and async access that works on any executor.
//!
//! Each connection runs on its own thread (via `async-sqlite`), so SQLite work
//! never blocks the runtime's executor.
//! never blocks the runtime's executor. [`Db::observe`] turns a query into a
//! stream that follows every commit.

mod observe;

use std::path::{Path, PathBuf};
use std::sync::Arc;

use async_sqlite::{JournalMode, Pool, PoolBuilder};
use rusqlite::{OpenFlags, TransactionBehavior};
use rusqlite_migration::Migrations;

pub use observe::ObservedStatement;
use observe::{Invalidation, authorize_for_change_tracking, track_changes};

/// Where a database lives.
#[derive(Debug, Clone)]
pub enum DbLocation {
Expand Down Expand Up @@ -73,6 +80,9 @@ pub enum DbError {
/// The host configured no database location.
#[error("no database configured")]
NotConfigured,
/// An observed query reads no table, so no commit could ever change it.
#[error("observed query reads no table: {0}")]
Unobservable(&'static str),
}

impl From<async_sqlite::Error> for DbError {
Expand Down Expand Up @@ -101,6 +111,7 @@ pub struct DbStatus {
pub struct Db {
writer: Pool,
readers: Pool,
invalidation: Arc<Invalidation>,
}

impl Db {
Expand All @@ -117,6 +128,8 @@ impl Db {
.map_err(|error| DbError::Open(error.to_string()))?;

let migrations = config.migrations;
let invalidation = Arc::new(Invalidation::default());
let hook = invalidation.clone();
writer
.conn_mut_and_then(move |conn| {
conn.pragma_update(None, "synchronous", "FULL")?;
Expand All @@ -129,7 +142,8 @@ impl Db {
rusqlite_migration::MigrationDefinitionError::NoMigrationsDefined,
)) => Ok(()),
Err(error) => Err(DbError::Migration(error.to_string())),
}
}?;
track_changes(conn, hook).map_err(DbError::from)
})
.await?;

Expand All @@ -143,7 +157,10 @@ impl Db {
.await
.map_err(|error| DbError::Open(error.to_string()))?;
readers
.conn_for_each(|conn| conn.busy_timeout(BUSY_TIMEOUT))
.conn_for_each(|conn| {
conn.busy_timeout(BUSY_TIMEOUT)?;
conn.authorizer(Some(authorize_for_change_tracking))
})
.await
.into_iter()
.collect::<Result<Vec<()>, _>>()?;
Expand All @@ -152,22 +169,32 @@ impl Db {
DbLocation::Memory => writer.clone(),
};

Ok(Self { writer, readers })
Ok(Self {
writer,
readers,
invalidation,
})
}

/// Runs `f` in one `BEGIN IMMEDIATE` transaction on the writer. Commits
/// when `f` returns `Ok` and rolls back when it returns `Err`.
/// when `f` returns `Ok` and rolls back when it returns `Err`. A commit
/// wakes the observers of every table it changed.
pub async fn write<T, F>(&self, f: F) -> Result<T, DbError>
where
F: FnOnce(&rusqlite::Transaction<'_>) -> Result<T, DbError> + Send + 'static,
T: Send + 'static,
{
let invalidation = self.invalidation.clone();
self.writer
.conn_mut_and_then(move |conn| {
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let value = f(&tx)?;
tx.commit()?;
Ok(value)
let committed = in_transaction(conn, f);
// Publish only once readers can see the rows: `commit_hook`
// runs before the commit is visible.
match committed {
Ok(_) => invalidation.publish(),
Err(_) => invalidation.discard(),
}
committed
})
.await
}
Expand Down Expand Up @@ -203,10 +230,22 @@ impl Db {
.await?;
self.readers.close().await?;
self.writer.close().await?;
self.invalidation.close();
Ok(())
}
}

/// Runs `f` in one `BEGIN IMMEDIATE` transaction, committing on `Ok`.
fn in_transaction<T>(
conn: &mut rusqlite::Connection,
f: impl FnOnce(&rusqlite::Transaction<'_>) -> Result<T, DbError>,
) -> Result<T, DbError> {
let tx = conn.transaction_with_behavior(TransactionBehavior::Immediate)?;
let value = f(&tx)?;
tx.commit()?;
Ok(value)
}

const BUSY_TIMEOUT: core::time::Duration = core::time::Duration::from_secs(5);

#[cfg(test)]
Expand Down
Loading
Loading