diff --git a/crates/cloudsync/src/api.rs b/crates/cloudsync/src/api.rs index 1e6d46cbe4c..857c2f048ba 100644 --- a/crates/cloudsync/src/api.rs +++ b/crates/cloudsync/src/api.rs @@ -8,7 +8,10 @@ use sqlx::{Executor, Sqlite, SqliteConnection}; use crate::error::Error; -const EXTENSION_LOAD_BUSY_RETRIES: usize = 4; +// PRAGMA busy_timeout = 50 keeps SQLite's C busy-handler short so after_connect +// does not park a Tokio worker. Retry SQLITE_BUSY asynchronously for the same +// 5s budget used by the rest of the app instead of failing the pool connect. +const EXTENSION_LOAD_BUSY_TIMEOUT: Duration = Duration::from_secs(5); type ExtensionLoadLock = tokio::sync::Mutex<()>; static EXTENSION_LOAD_LOCKS: OnceLock>>> = @@ -107,56 +110,40 @@ async fn load_and_validate_extension( connection: &mut SqliteConnection, extension_path: &Path, ) -> Result<(), Error> { + // Loading the extension is per-connection and does not need a reserved + // lock. BEGIN IMMEDIATE here fights in-flight writers and makes sqlx log + // after_connect SQLITE_BUSY while the pool grows. sqlx::query("PRAGMA busy_timeout = 50") .execute(&mut *connection) .await?; - let begin_result = begin_extension_load(connection).await; - if let Err(error) = begin_result { - restore_busy_timeout(connection).await?; - return Err(error); - } - - let load_result = async { - load_extension(connection, extension_path).await?; - validate_extension(connection).await - } - .await; - let finish_result = if load_result.is_ok() { - sqlx::query("COMMIT").execute(&mut *connection).await - } else { - sqlx::query("ROLLBACK").execute(&mut *connection).await - }; - let restore_result = restore_busy_timeout(connection).await; - - load_result?; - finish_result?; - restore_result?; - Ok(()) -} - -async fn begin_extension_load(connection: &mut SqliteConnection) -> Result<(), Error> { - let mut last_error = None; - for attempt in 0..EXTENSION_LOAD_BUSY_RETRIES { - match sqlx::query("BEGIN IMMEDIATE") - .execute(&mut *connection) - .await - { - Ok(_) => return Ok(()), - Err(error) if is_sqlite_busy(&error) => { - last_error = Some(error); - if attempt + 1 < EXTENSION_LOAD_BUSY_RETRIES { - let delay_ms = 10_u64 << attempt; - tokio::time::sleep(Duration::from_millis(delay_ms)).await; + let deadline = tokio::time::Instant::now() + EXTENSION_LOAD_BUSY_TIMEOUT; + let mut attempt = 0u32; + let result = loop { + match load_extension_and_validate(connection, extension_path).await { + Ok(()) => break Ok(()), + Err(error) if is_busy_error(&error) => { + let now = tokio::time::Instant::now(); + if now >= deadline { + break Err(error); } + let delay = Duration::from_millis(10_u64 << attempt.min(5)); + tokio::time::sleep(delay.min(deadline - now)).await; + attempt = attempt.saturating_add(1); } - Err(error) => return Err(error.into()), + Err(error) => break Err(error), } - } + }; + restore_busy_timeout(connection).await?; + result +} - Err(last_error - .expect("extension load retry loop must capture the final SQLITE_BUSY error") - .into()) +async fn load_extension_and_validate( + connection: &mut SqliteConnection, + extension_path: &Path, +) -> Result<(), Error> { + load_extension(connection, extension_path).await?; + validate_extension(connection).await } fn is_sqlite_busy(error: &sqlx::Error) -> bool { @@ -168,6 +155,19 @@ fn is_sqlite_busy(error: &sqlx::Error) -> bool { }) } +fn is_busy_error(error: &Error) -> bool { + match error { + Error::Sqlx(error) => is_sqlite_busy(error), + error => { + let message = error.to_string().to_ascii_lowercase(); + message.contains("database is locked") + || message.contains("database table is locked") + || message.contains("sqlite error 5") + || message.contains("sqlite error 6") + } + } +} + async fn restore_busy_timeout(connection: &mut SqliteConnection) -> Result<(), Error> { sqlx::query("PRAGMA busy_timeout = 5000") .execute(connection) @@ -497,3 +497,18 @@ where Ok(()) } + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn sqlite_busy_messages_are_retried() { + assert!(is_busy_error(&Error::ExtensionInitialization( + "failed to load sqlite-sync (sqlite error 5): database is locked".to_string(), + ))); + assert!(!is_busy_error(&Error::ExtensionInitialization( + "expected sqlite-sync 1.1.2, loaded 1.0.0".to_string(), + ))); + } +} diff --git a/crates/db-core/src/tests.rs b/crates/db-core/src/tests.rs index ee1c9201d1c..b259042d803 100644 --- a/crates/db-core/src/tests.rs +++ b/crates/db-core/src/tests.rs @@ -242,17 +242,13 @@ async fn cloudsync_pool_serializes_and_validates_every_extension_load() { .await .unwrap(); - let started_at = std::time::Instant::now(); let first_pool = db.pool().clone(); let second_pool = db.pool().clone(); let first = tokio::spawn(async move { first_pool.acquire().await }); let second = tokio::spawn(async move { second_pool.acquire().await }); - tokio::time::sleep(Duration::from_millis(90)).await; - sqlx::query("COMMIT").execute(&mut *writer).await.unwrap(); - let mut first = first.await.unwrap().unwrap(); let mut second = second.await.unwrap().unwrap(); - assert!(started_at.elapsed() >= Duration::from_millis(90)); + sqlx::query("COMMIT").execute(&mut *writer).await.unwrap(); assert_cloudsync_connection_ready( &mut writer,