#![warn(missing_docs)]
mod delivery_state;
mod enqueue;
mod error;
mod finalize;
mod hydrate;
mod import;
mod lifecycle;
mod lifecycle_mutation;
mod migration;
mod page;
mod schema;
mod scope;
pub use error::{
ClaimError, EnqueueError, FinalizeError, ImportError, MutationError, PageError, SchemaError,
TransientKind,
};
pub use migration::{
CrateVersion, LEGACY_MIGRATION, MIGRATIONS, Migration, MigrationCompatibility,
MigrationCompatibilityError, SCHEMA_VERSION, V1_TENANT_ACTIVATE_SQL, V1_TENANT_PREPARE_SQL,
};
pub use page::SnapshotPager;
pub use schema::check_schema;
pub use scope::{AdminDovecote, TenantDovecote};
use sqlx::{AssertSqlSafe, SqlSafeStr, Sqlite, SqlitePool, Transaction, query, query_scalar};
use std::time::Duration;
pub const DEFAULT_BUSY_RETRIES: u32 = 3;
pub const DEFAULT_BUSY_TIMEOUT: Duration = Duration::from_secs(5);
pub async fn begin_write(pool: &SqlitePool) -> Result<Transaction<'static, Sqlite>, EnqueueError> {
begin_write_with_config(pool, BusyConfig::default()).await
}
pub async fn begin_enqueue(
pool: &SqlitePool,
) -> Result<Transaction<'static, Sqlite>, EnqueueError> {
begin_write(pool).await
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct BusyConfig {
timeout: Duration,
retries: u32,
}
impl BusyConfig {
pub const fn new(timeout: Duration, retries: u32) -> Self {
Self { timeout, retries }
}
pub const fn with_retries(retries: u32) -> Self {
Self::new(DEFAULT_BUSY_TIMEOUT, retries)
}
pub const fn timeout(self) -> Duration {
self.timeout
}
pub const fn retries(self) -> u32 {
self.retries
}
}
impl Default for BusyConfig {
fn default() -> Self {
Self::new(DEFAULT_BUSY_TIMEOUT, DEFAULT_BUSY_RETRIES)
}
}
#[derive(Clone)]
pub struct SqliteDovecote {
pool: SqlitePool,
busy: BusyConfig,
}
impl SqliteDovecote {
pub fn new(pool: SqlitePool) -> Self {
Self {
pool,
busy: BusyConfig::default(),
}
}
pub const fn with_busy_config(pool: SqlitePool, busy: BusyConfig) -> Self {
Self { pool, busy }
}
pub fn pool(&self) -> &SqlitePool {
&self.pool
}
pub const fn busy_config(&self) -> BusyConfig {
self.busy
}
pub async fn begin_write(&self) -> Result<Transaction<'static, Sqlite>, EnqueueError> {
begin_write_with_config(&self.pool, self.busy).await
}
pub async fn begin_enqueue(&self) -> Result<Transaction<'static, Sqlite>, EnqueueError> {
self.begin_write().await
}
pub fn for_tenant(&self, tenant_id: dovecote::TenantId) -> TenantDovecote {
TenantDovecote::new(self.pool.clone(), tenant_id, self.busy)
}
pub fn admin(&self) -> AdminDovecote {
AdminDovecote::new(self.pool.clone(), self.busy)
}
pub async fn check_schema(&self) -> Result<(), SchemaError> {
check_schema(&self.pool).await
}
}
pub(crate) fn checked_milliseconds(value: Duration) -> Result<i64, String> {
if !value.is_zero() && !value.subsec_nanos().is_multiple_of(1_000_000) {
return Err("duration must be an exact whole number of milliseconds".to_owned());
}
i64::try_from(value.as_millis()).map_err(|_| "duration exceeds SQLite integer range".to_owned())
}
pub(crate) fn checked_busy_timeout(value: Duration) -> Result<i64, String> {
let milliseconds = checked_milliseconds(value)?;
if milliseconds > i64::from(i32::MAX) {
return Err("busy timeout exceeds SQLite's signed 32-bit millisecond range".to_owned());
}
Ok(milliseconds)
}
pub(crate) async fn begin_immediate(
pool: &SqlitePool,
busy: BusyConfig,
_operation: &'static str,
) -> Result<Transaction<'static, Sqlite>, sqlx::Error> {
let mut connection = pool.acquire().await?;
install_foreign_keys(&mut connection).await?;
install_busy_timeout(&mut connection, busy).await?;
Transaction::begin(
sqlx::pool::MaybePoolConnection::PoolConnection(connection),
Some(AssertSqlSafe("BEGIN IMMEDIATE").into_sql_str()),
)
.await
}
async fn begin_write_with_config(
pool: &SqlitePool,
busy: BusyConfig,
) -> Result<Transaction<'static, Sqlite>, EnqueueError> {
validate_busy_config(busy).map_err(|detail| EnqueueError::Configuration { detail })?;
let mut tries = 0;
loop {
match begin_immediate(pool, busy, "begin write transaction").await {
Ok(transaction) => return Ok(transaction),
Err(source) if error::is_busy(&source) && tries < busy.retries() => {
tries += 1;
}
Err(source) => return Err(EnqueueError::sql("begin write transaction", source)),
}
}
}
pub(crate) fn validate_busy_config(busy: BusyConfig) -> Result<(), String> {
checked_busy_timeout(busy.timeout()).map(|_| ())
}
pub(crate) async fn begin_read(
pool: &SqlitePool,
) -> Result<Transaction<'static, Sqlite>, sqlx::Error> {
let mut connection = pool.acquire().await?;
install_foreign_keys(&mut connection).await?;
Transaction::begin(
sqlx::pool::MaybePoolConnection::PoolConnection(connection),
Some(AssertSqlSafe("BEGIN").into_sql_str()),
)
.await
}
pub(crate) async fn commit_transaction(
mut transaction: Transaction<'static, Sqlite>,
) -> Result<(), sqlx::Error> {
use sqlx_core::transaction::TransactionManager;
let result =
<sqlx::sqlite::SqliteTransactionManager as TransactionManager>::commit(&mut *transaction)
.await;
if result.is_err() {
let _ = <sqlx::sqlite::SqliteTransactionManager as TransactionManager>::rollback(
&mut *transaction,
)
.await;
}
result
}
pub(crate) async fn install_foreign_keys(
connection: &mut sqlx::pool::PoolConnection<Sqlite>,
) -> Result<(), sqlx::Error> {
query("PRAGMA foreign_keys = ON")
.execute(&mut **connection)
.await?;
let enabled: i64 = query_scalar("PRAGMA foreign_keys")
.fetch_one(&mut **connection)
.await?;
if enabled != 1 {
return Err(sqlx::Error::Protocol(
"SQLite foreign-key enforcement could not be enabled".to_owned(),
));
}
Ok(())
}
pub(crate) async fn install_busy_timeout(
connection: &mut sqlx::pool::PoolConnection<Sqlite>,
busy: BusyConfig,
) -> Result<(), sqlx::Error> {
let milliseconds = checked_busy_timeout(busy.timeout())
.map_err(|detail| sqlx::Error::Protocol(format!("invalid busy configuration: {detail}")))?;
let statement = AssertSqlSafe(format!("PRAGMA busy_timeout = {milliseconds}"));
query(statement).execute(&mut **connection).await?;
let installed: i64 = query_scalar("PRAGMA busy_timeout")
.fetch_one(&mut **connection)
.await?;
if installed != milliseconds {
return Err(sqlx::Error::Protocol(format!(
"SQLite busy timeout installation mismatch: requested {milliseconds}, installed {installed}"
)));
}
Ok(())
}
pub(crate) async fn transaction_is_write(
transaction: &mut Transaction<'_, Sqlite>,
) -> Result<bool, sqlx::Error> {
let mut handle = transaction.lock_handle().await?;
let state = unsafe {
libsqlite3_sys::sqlite3_txn_state(handle.as_raw_handle().as_ptr(), std::ptr::null())
};
Ok(state == libsqlite3_sys::SQLITE_TXN_WRITE)
}