diff --git a/CHANGELOG.md b/CHANGELOG.md index 4656d68..945ad69 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,10 @@ +## 0.0.8 (unreleased) + +- __Breaking__: Timers on `PowerSyncEnvironment::custom` are now passed by value. +- Don't mark sync status as connected when connection fails. +- Fix sync client blocking writer for longer than necessary. +- Set cache size and busy timeout on all connections instead of just the writer. + ## 0.0.7 - Update PowerSync core extension to version 0.5.2. diff --git a/powersync/src/db/internal.rs b/powersync/src/db/internal.rs index ad0fb33..9d688c5 100644 --- a/powersync/src/db/internal.rs +++ b/powersync/src/db/internal.rs @@ -10,11 +10,9 @@ use crate::{ util::SharedFuture, }; use event_listener::EventListener; -use futures_lite::future::yield_now; use futures_lite::{FutureExt, Stream, StreamExt, ready}; use powersync_sqlite_nostd::{ColumnType, Destructor, ResultCode}; -use std::sync::{Mutex, Weak}; -use std::time::Duration; +use std::sync::Weak; use std::{ pin::Pin, sync::Arc, @@ -38,7 +36,6 @@ pub struct InnerPowerSyncState { /// reference to [InnerPowerSyncState], we only keep a weak reference here to ensure we can drop /// actors through the channels owned by [SyncCoordinator]. pub(crate) sync: Weak, - pub(crate) retry_delay: Mutex>, } impl InnerPowerSyncState { @@ -53,7 +50,6 @@ impl InnerPowerSyncState { schema: Arc::new(schema), status: SyncStatus::new(), current_streams: SyncStreamTracker::default(), - retry_delay: Default::default(), sync: Arc::downgrade(sync), } } @@ -156,21 +152,6 @@ impl InnerPowerSyncState { Ok(self.env.pool.writer().await) } - pub async fn sync_iteration_delay(&self) { - let delay = { - let guard = self.retry_delay.lock().unwrap(); - *guard - }; - - if let Some(delay) = delay - && delay > Duration::ZERO - { - self.env.timer.delay_once(delay).await - } else { - yield_now().await - } - } - pub fn watch_status<'a>(&'a self) -> impl Stream> + 'a { struct StreamImpl<'a> { db: &'a InnerPowerSyncState, diff --git a/powersync/src/db/mod.rs b/powersync/src/db/mod.rs index bae02a9..e423dd7 100644 --- a/powersync/src/db/mod.rs +++ b/powersync/src/db/mod.rs @@ -84,7 +84,7 @@ impl PowerSyncDatabase { /// Requests the download actor, started with [Self::download_actor], to start establishing a /// connection to the PowerSync service. pub async fn connect(&self, options: SyncOptions) { - self.sync.connect(options, &self.inner).await + self.sync.connect(options).await } /// If the sync client is currently connected, requests it to disconnect. diff --git a/powersync/src/db/pool.rs b/powersync/src/db/pool.rs index 763c51f..0a2af80 100644 --- a/powersync/src/db/pool.rs +++ b/powersync/src/db/pool.rs @@ -35,10 +35,15 @@ impl ConnectionPool { SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE, )?); + fn configure_common(connection: &SqliteConnection) -> Result<(), PowerSyncError> { + connection.exec(c"PRAGMA busy_timeout = 30000")?; + connection.exec(c"PRAGMA cache_size = -51200")?; // -(50 * 1024) + Ok(()) + } + writer.exec(c"PRAGMA journal_mode = WAL")?; writer.exec(c"PRAGMA journal_size_limit = 6291456")?; // 6 * 1024 * 1024 - writer.exec(c"PRAGMA busy_timeout = 30000")?; - writer.exec(c"PRAGMA cache_size = -51200")?; // -(50 * 1024) + configure_common(&writer)?; let mut readers = vec![]; for _ in 0..5 { @@ -46,7 +51,7 @@ impl ConnectionPool { &path, SQLITE_OPEN_READONLY, )?); - reader.exec(c"PRAGMA query_only = 1")?; + configure_common(&reader)?; readers.push(reader); } diff --git a/powersync/src/env.rs b/powersync/src/env.rs index edc5d60..06b2b37 100644 --- a/powersync/src/env.rs +++ b/powersync/src/env.rs @@ -4,6 +4,7 @@ use crate::http::HttpClient; use num_traits::FromPrimitive; use powersync_core::powersync_init_static; use powersync_sqlite_nostd::ResultCode; +use std::sync::Arc; use std::{pin::Pin, time::Duration}; /// All external dependencies required for the PowerSync SDK. @@ -17,19 +18,15 @@ pub struct PowerSyncEnvironment { /// The [ConnectionPool] used to obtain connections for queries asynchronously. pub(crate) pool: ConnectionPool, /// The [Timer] implementation used to delay sync iterations after errors. - pub(crate) timer: &'static (dyn Timer + Send + Sync), + pub(crate) timer: Arc, } impl PowerSyncEnvironment { - pub fn custom( - client: C, - pool: ConnectionPool, - timer: &'static (dyn Timer + Send + Sync), - ) -> Self { + pub fn custom(client: C, pool: ConnectionPool, timer: T) -> Self { Self { client: Box::new(client), pool, - timer, + timer: Arc::new(timer), } } @@ -50,7 +47,7 @@ impl PowerSyncEnvironment { /// A [Timer] implementation based on [async_io::Timer]. #[cfg(feature = "smol")] - pub fn async_io_timer() -> &'static (dyn Timer + Send + Sync) { + pub fn async_io_timer() -> impl Timer { use async_io::Timer as PlatformTimer; struct AsyncIoTimer; @@ -64,12 +61,12 @@ impl PowerSyncEnvironment { .boxed() } } - &AsyncIoTimer + AsyncIoTimer } /// A [Timer] implementation based on [tokio::time::sleep]. #[cfg(feature = "tokio")] - pub fn tokio_timer() -> &'static (dyn Timer + Send + Sync) { + pub fn tokio_timer() -> impl Timer { use tokio::time::sleep; struct TokioTimer; @@ -80,7 +77,7 @@ impl PowerSyncEnvironment { sleep(duration).boxed() } } - &TokioTimer + TokioTimer } } @@ -90,7 +87,7 @@ impl PowerSyncEnvironment { /// Because the native PowerSync SDK is executor-agnostic, it can't use a builtin function to retry /// sync after a delay to recover from errors. This trait, as part of the [PowerSyncEnvironment], /// is thus used to schedule the delay. -pub trait Timer { +pub trait Timer: Send + Sync + 'static { /// Returns a future that returns [Poll::Pending] when being polled the first time and schedules /// the context's waker to be woken after the specified `duration`. fn delay_once(&self, duration: Duration) -> Pin + Send>>; diff --git a/powersync/src/sync/coordinator.rs b/powersync/src/sync/coordinator.rs index f5b5f12..441802d 100644 --- a/powersync/src/sync/coordinator.rs +++ b/powersync/src/sync/coordinator.rs @@ -5,7 +5,6 @@ use async_oneshot::oneshot; use crate::{ SyncOptions, - db::internal::InnerPowerSyncState, sync::{ download::DownloadActorCommand, streams::ChangedSyncSubscriptions, upload::UploadActorCommand, @@ -42,16 +41,10 @@ pub struct SyncCoordinator { } impl SyncCoordinator { - pub async fn connect(&self, options: SyncOptions, db: &InnerPowerSyncState) { - { - let mut lock = db.retry_delay.lock().unwrap(); - *lock = Some(options.retry_delay); - } - - let connector = options.connector.clone(); - self.download_actor_request(DownloadActorCommand::Connect(options)) + pub async fn connect(&self, options: SyncOptions) { + self.download_actor_request(DownloadActorCommand::Connect(options.clone())) .await; - self.upload_actor_request(UploadActorCommand::Connect(connector)) + self.upload_actor_request(UploadActorCommand::Connect(options)) .await; } diff --git a/powersync/src/sync/download/actor.rs b/powersync/src/sync/download/actor.rs index ba31b09..9e5cc22 100644 --- a/powersync/src/sync/download/actor.rs +++ b/powersync/src/sync/download/actor.rs @@ -75,6 +75,14 @@ impl DownloadActor { }; } + fn retry_delay(&self) -> Boxed<()> { + if let Some(ref options) = self.options { + options.retry_delay(&self.db.env).boxed() + } else { + async {}.boxed() + } + } + async fn handle_event(&mut self) { match &mut self.state { DownloadActorState::Idle => { @@ -176,18 +184,15 @@ impl DownloadActor { let timeout = if close.hide_disconnect { async {}.boxed() } else { - let db = self.db.clone(); - - async move { db.sync_iteration_delay().await }.boxed() + self.retry_delay() }; self.state = DownloadActorState::WaitingForReconnect { timeout } } Event::SyncIterationError(e) => { self.db.status.update(|status| status.set_download_error(e)); - let db = self.db.clone(); self.state = DownloadActorState::WaitingForReconnect { - timeout: async move { db.sync_iteration_delay().await }.boxed(), + timeout: self.retry_delay(), } } } diff --git a/powersync/src/sync/download/http.rs b/powersync/src/sync/download/http.rs index 9e9f0c9..3a307f7 100644 --- a/powersync/src/sync/download/http.rs +++ b/powersync/src/sync/download/http.rs @@ -188,7 +188,7 @@ mod tests { fn first_event(client: impl HttpClient) -> Result, PowerSyncError> { PowerSyncEnvironment::powersync_auto_extension().unwrap(); let pool = ConnectionPool::single_connection(Connection::open_in_memory().unwrap()); - let environment = PowerSyncEnvironment::custom(client, pool, &UnusedTimer); + let environment = PowerSyncEnvironment::custom(client, pool, UnusedTimer); let coordinator = Arc::new(SyncCoordinator::default()); let db = Arc::new(InnerPowerSyncState::new( environment, diff --git a/powersync/src/sync/options.rs b/powersync/src/sync/options.rs index bc8a104..65c0788 100644 --- a/powersync/src/sync/options.rs +++ b/powersync/src/sync/options.rs @@ -1,6 +1,8 @@ use std::{sync::Arc, time::Duration}; -use crate::sync::connector::BackendConnector; +use futures_lite::future::yield_now; + +use crate::{env::PowerSyncEnvironment, sync::connector::BackendConnector}; /// Options controlling how PowerSync connects to a sync service. #[derive(Clone)] @@ -34,4 +36,20 @@ impl SyncOptions { pub fn with_retry_delay(&mut self, delay: Duration) { self.retry_delay = delay; } + + pub(crate) fn retry_delay( + &self, + env: &PowerSyncEnvironment, + ) -> impl Future + 'static { + let delay = self.retry_delay; + let timer = env.timer.clone(); + + async move { + if delay > Duration::ZERO { + timer.delay_once(delay).await + } else { + yield_now().await + } + } + } } diff --git a/powersync/src/sync/upload.rs b/powersync/src/sync/upload.rs index 2b0000e..a398446 100644 --- a/powersync/src/sync/upload.rs +++ b/powersync/src/sync/upload.rs @@ -7,11 +7,13 @@ use futures_lite::{ use log::{debug, info, warn}; use powersync_sqlite_nostd::{Destructor, ResultCode}; -use crate::db::connection::{SqliteConnection, TransactionGuard}; use crate::db::watch::ListenerConfiguration; use crate::sync::coordinator::SyncCoordinator; use crate::{ - BackendConnector, + SyncOptions, + db::connection::{SqliteConnection, TransactionGuard}, +}; +use crate::{ db::internal::InnerPowerSyncState, error::PowerSyncError, sync::{ @@ -21,7 +23,7 @@ use crate::{ }; pub enum UploadActorCommand { - Connect(Arc), + Connect(SyncOptions), TriggerCrudUpload, Disconnect, } @@ -51,7 +53,7 @@ impl UploadActor { fn connected_state( db: &Arc, - connector: Arc, + options: SyncOptions, ) -> ConnectedUploadActor { let mut tables = HashSet::new(); tables.insert("ps_crud".to_string()); @@ -62,7 +64,7 @@ impl UploadActor { .update_notifiers() .listen(ListenerConfiguration::if_matches(tables, false)); ConnectedUploadActor { - connector, + options, crud_stream: stream.map(|_| ()).boxed(), } } @@ -77,10 +79,10 @@ impl UploadActor { // Already in progress, don't start another. None } - UploadActorCommand::Connect(connector) => { - // TODO: Only abort if the connector has changed? + UploadActorCommand::Connect(options) => { + // TODO: Only abort if options have changed? Some(UploadActorState::Connected(Self::connected_state( - db, connector, + db, options, ))) } UploadActorCommand::Disconnect => Some(UploadActorState::Idle), @@ -104,12 +106,12 @@ impl UploadActor { }; match command.command { - UploadActorCommand::Connect(connector) => { + UploadActorCommand::Connect(options) => { let _ = command.response.send(()); - UploadActorState::Connected(Self::connected_state(&self.db, connector)) + UploadActorState::Connected(Self::connected_state(&self.db, options)) } UploadActorCommand::TriggerCrudUpload => { - // We can't upload because we're not connector + // We can't upload because we're not connected old_state } UploadActorCommand::Disconnect => { @@ -138,8 +140,8 @@ impl UploadActor { let _ = command.response.send(()); match command.command { - UploadActorCommand::Connect(connector) => Transition::Abort( - UploadActorState::Connected(Self::connected_state(&self.db, connector)), + UploadActorCommand::Connect(options) => Transition::Abort( + UploadActorState::Connected(Self::connected_state(&self.db, options)), ), UploadActorCommand::TriggerCrudUpload => Transition::StartUpload, UploadActorCommand::Disconnect => Transition::Abort(UploadActorState::Idle), @@ -181,7 +183,7 @@ impl UploadActor { UploadActorState::RunningUpload { result: async move { let mut upload = CrudUpload { - connector: state.connector.as_ref(), + options: &state.options, db, }; upload.run().await; @@ -207,14 +209,13 @@ impl UploadActorState { } struct ConnectedUploadActor { - /// The connector to use when uploading changes. - connector: Arc, + options: SyncOptions, /// A stream emitting changes when the `ps_crud` table is updated locally. crud_stream: futures_lite::stream::Boxed<()>, } struct CrudUpload<'a> { - connector: &'a dyn BackendConnector, + options: &'a SyncOptions, db: Arc, } @@ -234,7 +235,7 @@ impl<'a> CrudUpload<'a> { self.db .status .update(|data| data.set_upload_state(UploadStatus::Error(e))); - self.db.sync_iteration_delay().await; + self.options.retry_delay(&self.db.env).await; } } } @@ -271,7 +272,7 @@ impl<'a> CrudUpload<'a> { } *last_item_id = Some(item); - self.connector.upload_data().await?; + self.options.connector.upload_data().await?; Ok(ControlFlow::Continue(())) } @@ -295,7 +296,7 @@ impl<'a> CrudUpload<'a> { stmt.column_text(0)?.to_string() }; - let credentials = self.connector.fetch_credentials().await?; + let credentials = self.options.connector.fetch_credentials().await?; write_checkpoint(&self.db, &client_id, credentials).await } diff --git a/powersync/tests/sync_test.rs b/powersync/tests/sync_test.rs index 8a7b613..e0895ef 100644 --- a/powersync/tests/sync_test.rs +++ b/powersync/tests/sync_test.rs @@ -496,3 +496,36 @@ fn reports_correct_times() { assert!(delta < Duration::from_secs(5)); }); } + +#[test] +fn reconnects_on_failure() { + let sync = SyncStreamTest::new(); + sync.connect_options(|options| { + options.with_retry_delay(Duration::from_hours(1)); + }); + + sync.run(async { + let request = sync.test.http.receive_requests.recv().await.unwrap(); + sync.wait_for_status(|s| s.is_connected()).await; + + // Send a line causing an error + request + .channel + .send(SyncLine::Custom(json!("invalid sync line"))) + .await + .unwrap(); + + sync.wait_for_status(|s| s.download_error().is_some()).await; + }); + + let task = sync.test.ex.spawn({ + let http = sync.test.http.clone(); + async move { http.receive_requests.recv().await } + }); + + // Should reconnect after the configured delay. + sync.test.advance_time(Duration::from_mins(30)); + assert!(!task.is_finished()); + sync.test.advance_time(Duration::from_mins(30)); + assert!(task.is_finished()); +} diff --git a/powersync_test_utils/src/lib.rs b/powersync_test_utils/src/lib.rs index 20c2ccc..3fa31ca 100644 --- a/powersync_test_utils/src/lib.rs +++ b/powersync_test_utils/src/lib.rs @@ -1,6 +1,13 @@ -use std::{sync::Arc, vec}; +use std::{ + collections::VecDeque, + sync::{Arc, Mutex}, + task::{Context, Poll, Waker}, + time::Duration, + vec, +}; use async_executor::Executor; +use futures_lite::FutureExt; use log::LevelFilter; use powersync::{ env::{PowerSyncEnvironment, Timer}, @@ -19,6 +26,7 @@ pub mod sync_line; pub struct DatabaseTest { pub dir: TempDir, pub http: Arc, + timer: Arc>, pub ex: Executor<'static>, } @@ -32,6 +40,7 @@ impl Default for DatabaseTest { Self { dir: TempDir::new("powersync_rust").expect("should create test directory"), http: Arc::new(MockSyncService::new()), + timer: Default::default(), ex: Executor::new(), } } @@ -65,21 +74,88 @@ impl DatabaseTest { PowerSyncDatabase::new(self.in_memory(), Self::default_schema()) } + pub fn advance_time(&self, duration: Duration) { + let mut state = self.timer.lock().unwrap(); + let end_timestamp = state.time_passed + duration; + state.time_passed = end_timestamp; + + while let Some((_, waker)) = state + .scheduled + .pop_front_if(|&mut (time, _)| time <= end_timestamp) + { + waker.wake(); + } + + drop(state); + + // Drive async tasks until idle + while self.ex.try_tick() {} + } + fn env(&self, pool: ConnectionPool) -> PowerSyncEnvironment { PowerSyncEnvironment::powersync_auto_extension().expect("should load core extension"); - struct DisabledTimer; + let timer = self.timer.clone(); + + struct TestTimer { + state: Arc>, + } + + struct TestDelay { + state: Arc>, + end_timestamp: Duration, + did_register: bool, + } - impl Timer for DisabledTimer { + impl Timer for TestTimer { fn delay_once( &self, - _duration: std::time::Duration, + duration: Duration, ) -> std::pin::Pin + Send>> { - panic!("Tests should not run into a delay") + let state = self.state.lock().unwrap(); + let end_timestamp = state.time_passed + duration; + + TestDelay { + state: self.state.clone(), + end_timestamp, + did_register: false, + } + .boxed() } } - PowerSyncEnvironment::custom(self.http.clone().client(), pool, &DisabledTimer) + impl Future for TestDelay { + type Output = (); + + fn poll(mut self: std::pin::Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<()> { + let state = self.state.clone(); + let mut state = state.lock().unwrap(); + if state.time_passed >= self.end_timestamp { + return Poll::Ready(()); + } + + if !self.did_register { + self.did_register = true; + + let search = state + .scheduled + .binary_search_by_key(&self.end_timestamp, |&(ts, _)| ts); + + let schedule = (self.end_timestamp, cx.waker().clone()); + state.scheduled.insert( + match search { + Ok(existing) => existing + 1, + Err(expected_index) => expected_index, + }, + schedule, + ); + } + + Poll::Pending + } + } + + PowerSyncEnvironment::custom(self.http.clone().client(), pool, TestTimer { state: timer }) } pub fn default_schema() -> Schema { @@ -90,6 +166,12 @@ impl DatabaseTest { } } +#[derive(Default)] +struct MockTimer { + time_passed: Duration, + scheduled: VecDeque<(Duration, Waker)>, +} + /// Runs a query and returns rows as a `serde_json` array. pub async fn query_all(db: &PowerSyncDatabase, sql: &str, params: impl Params) -> Value { let reader = db.reader().await.unwrap();