From 0d786c98f3fd76f136d62e1b2a5006bd2321b809 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Mon, 5 Oct 2026 15:00:44 +0200 Subject: [PATCH 1/2] Split `BackendConnector` trait --- CHANGELOG.md | 6 + examples/egui_todolist/src/database.rs | 49 +++- powersync/README.md | 10 +- powersync/src/lib.rs | 2 +- powersync/src/sync/checkpoint.rs | 20 +- powersync/src/sync/connector.rs | 60 +++-- powersync/src/sync/coordinator.rs | 7 +- powersync/src/sync/download/http.rs | 55 ++-- powersync/src/sync/download/mod.rs | 6 +- powersync/src/sync/download/sync_iteration.rs | 19 +- powersync/src/sync/options.rs | 69 ++++- powersync/src/sync/upload.rs | 153 +++++++---- powersync/tests/sync_test.rs | 249 +++++++++++++----- powersync_test_utils/src/mock_sync_service.rs | 30 ++- 14 files changed, 528 insertions(+), 207 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 37eaa86..9c2af87 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,3 +1,9 @@ +## 0.2.0 (unreleased) + +- __Breaking__: Split `BackendConnector` trait into `Authenticator` and `MutationUploader`. + - `Authenticator` is allowed to cache tokens now, and is informed when a token has expired. + - Add `SyncOptions::download_only` and `SyncOptions::upload_only` for unidirectional sync. + ## 0.1.0 - __Breaking__: `PowerSyncEnvironment::custom` now needs a way to spawn async diff --git a/examples/egui_todolist/src/database.rs b/examples/egui_todolist/src/database.rs index 5c0b2d6..d117d88 100644 --- a/examples/egui_todolist/src/database.rs +++ b/examples/egui_todolist/src/database.rs @@ -1,8 +1,10 @@ +use std::sync::{Arc, Mutex}; + use async_trait::async_trait; use futures_lite::StreamExt; use log::warn; use powersync::{ - BackendConnector, CheckpointMode, ConnectionPool, PowerSyncCredentials, PowerSyncDatabase, + Authenticator, CheckpointMode, ConnectionPool, MutationUploader, PowerSyncDatabase, SyncOptions, UpdateType, env::PowerSyncEnvironment, error::PowerSyncError, @@ -79,6 +81,7 @@ impl TodoList { #[derive(Clone)] pub struct TodoDatabase { pub db: PowerSyncDatabase, + cached_credentials: Arc>>>, } impl TodoDatabase { @@ -94,11 +97,14 @@ impl TodoDatabase { schema.tables.push(TodoEntry::schema()); let db = PowerSyncDatabase::new(env, schema); - Self { db } + Self { + db, + cached_credentials: Arc::default(), + } } pub async fn connect(&self) { - let mut options = SyncOptions::new(self.clone()); + let mut options = SyncOptions::new("http://localhost:8080", self.clone(), self.clone()); options.with_checkpoint_mode(CheckpointMode::Requests(Default::default())); self.db.connect(options).await } @@ -107,29 +113,46 @@ impl TodoDatabase { self.db.disconnect().await; } - async fn fetch_credentials_self_hosted(&self) -> Result { + async fn fetch_credentials_self_hosted(&self) -> Result, PowerSyncError> { let response = reqwest::get("http://localhost:6060/api/auth/token").await?; #[derive(Deserialize)] struct TokenResponse { - token: String, + token: Arc, } let token: TokenResponse = response.json().await?; - Ok(PowerSyncCredentials { - endpoint: "http://localhost:8080".to_string(), - token: token.token, - }) + Ok(token.token) } } #[async_trait] -impl BackendConnector for TodoDatabase { - async fn fetch_credentials(&self) -> Result { - self.fetch_credentials_self_hosted().await +impl Authenticator for TodoDatabase { + async fn resolve_credentials(&self) -> Result, PowerSyncError> { + { + let guard = self.cached_credentials.lock().unwrap(); + if let Some(ref cached) = *guard { + // Fast path, we have cached credentials. Concurrent fetch_credentials_self_hosted + // are okay too. + return Ok(cached.clone()); + } + } + + let credentials = self.fetch_credentials_self_hosted().await?; + let mut guard = self.cached_credentials.lock().unwrap(); + *guard = Some(credentials.clone()); + Ok(credentials) } - async fn upload_data(&self) -> Result<(), PowerSyncError> { + fn invalidate_credentials(&self) { + let mut guard = self.cached_credentials.lock().unwrap(); + *guard = None; + } +} + +#[async_trait] +impl MutationUploader for TodoDatabase { + async fn upload(&self) -> Result<(), PowerSyncError> { let mut transactions = self.db.crud_transactions(); let mut last_tx = None; diff --git a/powersync/README.md b/powersync/README.md index 9a8cf7b..9200fef 100644 --- a/powersync/README.md +++ b/powersync/README.md @@ -61,7 +61,15 @@ stream transformer. ### Connecting To automatically keep the local SQLite database in-sync with a backend process, call `PowerSyncDatabase::connect`. -This requires passing your own `BackendConnector` implementation, see the examples for a possible implementation. +This requires passing `SyncOptions`, which can be created for: + +- downloading only (via `SyncOptions::download_only`), which requires the endpoint of your PowerSync instance and an + `powersync::Authenticator` implementation returning JWT tokens for the client to use. +- uploading only (via `SyncOptions::upload_only`), which requires a `powersync::MutationUploader` responsible for + uploading local mutations to your backend. In this mode, no connection to a PowerSync service is made. +- bidirectional sync (recommended, via `SyncOptions::new`). This requires both an authenticator and a mutation uploader. + +Also see the examples for a possible implementation. ### Sync status diff --git a/powersync/src/lib.rs b/powersync/src/lib.rs index fc16497..92efc95 100644 --- a/powersync/src/lib.rs +++ b/powersync/src/lib.rs @@ -12,7 +12,7 @@ pub use db::streams::StreamSubscriptionOptions; pub use db::streams::SyncStream; pub use db::{DisconnectAndClearFlags, PowerSyncDatabase}; pub use sync::checkpoint::{CheckpointError, CheckpointRequest}; -pub use sync::connector::{BackendConnector, PowerSyncCredentials}; +pub use sync::connector::{Authenticator, MutationUploader}; pub use sync::options::{CheckpointMode, RequestsCheckpointMode, SyncOptions}; pub use sync::status::SyncStatusData; pub use sync::stream_priority::StreamPriority; diff --git a/powersync/src/sync/checkpoint.rs b/powersync/src/sync/checkpoint.rs index cc48e94..b94ed01 100644 --- a/powersync/src/sync/checkpoint.rs +++ b/powersync/src/sync/checkpoint.rs @@ -5,13 +5,14 @@ use log::{debug, warn}; use thiserror::Error; use crate::{ - BackendConnector, CheckpointMode, RequestsCheckpointMode, SyncOptions, SyncStatusData, + CheckpointMode, RequestsCheckpointMode, SyncOptions, SyncStatusData, db::internal::InnerPowerSyncState, error::{PowerSyncError, RawPowerSyncError}, sync::{ coordinator::{SyncChannels, SyncCoordinator}, download::http::checkpoint_request, instruction::CheckpointRequestPayload, + options::EndpointAndAuthenticator, upload::get_client_id, }, }; @@ -113,16 +114,15 @@ pub async fn repost_unacknowledged_checkpoints( let CheckpointMode::Requests(requests) = options.checkpoints else { return; }; + let Some(ref authenticator) = options.endpoint else { + return; + }; loop { // Make sure the system is seeded and ready. - let result = repost_unacknowledged_checkpoint_iteration( - &db, - &channels, - options.connector.as_ref(), - requests, - ) - .await; + let result = + repost_unacknowledged_checkpoint_iteration(&db, &channels, authenticator, requests) + .await; if let Err(err) = result { if let RawPowerSyncError::Checkpoint { @@ -141,7 +141,7 @@ pub async fn repost_unacknowledged_checkpoints( async fn repost_unacknowledged_checkpoint_iteration( db: &InnerPowerSyncState, channels: &SyncChannels, - connector: &dyn BackendConnector, + authenticator: &EndpointAndAuthenticator, mode: RequestsCheckpointMode, ) -> Result<(), PowerSyncError> { // Make sure the system is seeded and ready @@ -187,7 +187,7 @@ async fn repost_unacknowledged_checkpoint_iteration( debug!("Retrying checkpoint request id {request_id}"); checkpoint_request( &db, - connector, + authenticator, &CheckpointRequestPayload { client_id: get_client_id(db).await?, checkpoint_request_id: request_id, diff --git a/powersync/src/sync/connector.rs b/powersync/src/sync/connector.rs index f8429a7..04f9597 100644 --- a/powersync/src/sync/connector.rs +++ b/powersync/src/sync/connector.rs @@ -1,19 +1,26 @@ -use std::pin::Pin; +use std::{pin::Pin, sync::Arc}; use async_trait::async_trait; use url::Url; use crate::error::{PowerSyncError, RawPowerSyncError}; -/// A backend connector is responsible for upload local writes as well as resolving JWTs used to -/// connect to the PowerSync service. +/// Authenticates the PowerSync SDK against a PowerSync service, allowing it to sync changes from a +/// backend source database. #[async_trait] -pub trait BackendConnector: Send + Sync { - /// Fetches a fresh JWT from the backend to be used against the PowerSync service. - async fn fetch_credentials(&self) -> Result; +pub trait Authenticator: Send + Sync { + /// Resolves a JWT to use when connecting to a PowerSync service. + /// + /// The SDK does not cache this value and will call this method multiple times while a + /// database is connected. Implementations should consider caching credentials. + async fn resolve_credentials(&self) -> Result, PowerSyncError>; - /// Inspects completed CRUD transactions on a database and uploads them. - async fn upload_data(&self) -> Result<(), PowerSyncError>; + /// Invoked by the SDK when the PowerSync service has rejected credentials previously + /// returned by [Self::resolve_credentials]. + /// + /// If the ocnnector returns cached tokens, it can use this as a hint to refresh its local + /// state. + fn invalidate_credentials(&self) {} /// This is optional, and should only return a future for connectors capable of requesting /// checkpoints. @@ -36,15 +43,37 @@ pub trait BackendConnector: Send + Sync { } } +/// Uploads local mutations (from `INSERT`, `UPDATE` and `DELETE` statements against the +/// local database) to your source database. +/// +/// While simple cases might use a protocol like PostgREST to write directly into the database, +/// a custom backend is commonly used to validate uploaded mutations. +#[async_trait] +pub trait MutationUploader: Send + Sync { + async fn upload(&self) -> Result<(), PowerSyncError>; +} + +/// Allows using async functions and closures (e.g. `|| async { ... }`) as a [MutationUploader]. +#[async_trait] +impl MutationUploader for F +where + F: Fn() -> Fut + Send + Sync, + Fut: Future> + Send, +{ + async fn upload(&self) -> Result<(), PowerSyncError> { + self().await + } +} + /// Credentials used to connect to a PowerSync service instance. -pub struct PowerSyncCredentials { +pub struct PowerSyncCredentials<'a> { /// PowerSync endpoint, e.g. `https://myinstance.powersync.co`. - pub endpoint: String, + pub endpoint: &'a str, /// The token used to authenticate against the PowerSync service. - pub token: String, + pub token: Arc, } -impl PowerSyncCredentials { +impl<'a> PowerSyncCredentials<'a> { /// Parses the [Self::endpoint] into a URI. pub(crate) fn parsed_endpoint(&self, endpoint: &str) -> Result { let url = Url::parse(&self.endpoint) @@ -61,12 +90,13 @@ impl PowerSyncCredentials { #[cfg(test)] mod test { - use crate::PowerSyncCredentials; + use super::PowerSyncCredentials; + use std::sync::Arc; fn is_endpoint_valid(endpoint: &str) -> bool { PowerSyncCredentials { - token: "".to_string(), - endpoint: endpoint.to_string(), + token: Arc::new("test".to_string()), + endpoint, } .parsed_endpoint("") .is_ok() diff --git a/powersync/src/sync/coordinator.rs b/powersync/src/sync/coordinator.rs index b0e5a77..4fea136 100644 --- a/powersync/src/sync/coordinator.rs +++ b/powersync/src/sync/coordinator.rs @@ -166,9 +166,12 @@ impl SyncCoordinator { ) -> Result { let guard = self.task.lock().await; let tasks = Self::extract_connected_with_requests(guard.as_ref())?; + let Some(ref authenticator) = tasks.options.endpoint else { + return Err(CheckpointError::Disconnected); + }; + let authenticator = authenticator.clone(); let channels = tasks.channels.clone(); - let connector = tasks.options.connector.clone(); // Avoid holding the lock across a suspension point. It's fine if there's a concurrent // reconnect, post_checkpoint_request will return an error in that case. drop(guard); @@ -177,7 +180,7 @@ impl SyncCoordinator { .await .map_err(CheckpointError::as_request_error)?; let checkpoint_request_id = - post_checkpoint_request(client_id, connector.as_ref(), &channels, &db) + post_checkpoint_request(client_id, &authenticator, &channels, &db) .await .map_err(CheckpointError::as_request_error)?; diff --git a/powersync/src/sync/download/http.rs b/powersync/src/sync/download/http.rs index 5c35b6c..323a448 100644 --- a/powersync/src/sync/download/http.rs +++ b/powersync/src/sync/download/http.rs @@ -1,15 +1,16 @@ use std::borrow::Cow; use std::sync::Arc; -use crate::BackendConnector; +use crate::Authenticator; use crate::http::{Request, Response}; use crate::sync::checkpoint::CheckpointError; use crate::sync::instruction::CheckpointRequestPayload; +use crate::sync::options::EndpointAndAuthenticator; use crate::util::LineSplitter; use crate::{ db::internal::InnerPowerSyncState, error::{PowerSyncError, RawPowerSyncError}, - sync::{connector::PowerSyncCredentials, download::sync_iteration::DownloadEvent}, + sync::download::sync_iteration::DownloadEvent, util::BsonObjects, }; use futures_lite::{Stream, StreamExt, stream}; @@ -20,10 +21,12 @@ use serde_with::{DisplayFromStr, serde_as}; /// connection to the PowerSync service. pub fn sync_stream( db: Arc, - auth: PowerSyncCredentials, + authenticator: Arc, request_body: String, ) -> impl Stream> { let response = async move { + let auth = authenticator.fetch_credentials().await?; + let request = Request { method: "POST", url: auth.parsed_endpoint("sync/stream")?, @@ -42,7 +45,7 @@ pub fn sync_stream( }; let response = db.env.client.send(request).await?; - check_ok(response.status)?; + check_ok(response.status, authenticator.authenticator.as_ref())?; Ok::(response) }; @@ -61,8 +64,9 @@ pub fn sync_stream( pub async fn write_checkpoint( db: &InnerPowerSyncState, client_id: &str, - auth: PowerSyncCredentials, + authenticator: &EndpointAndAuthenticator, ) -> Result { + let auth = authenticator.fetch_credentials().await?; let mut url = auth.parsed_endpoint("write-checkpoint2.json")?; url.set_query(Some(&format!("client_id={}", client_id))); @@ -81,7 +85,7 @@ pub async fn write_checkpoint( }; let response = db.env.client.send(request).await?; - check_ok(response.status)?; + check_ok(response.status, authenticator.authenticator.as_ref())?; #[derive(Deserialize)] struct WriteCheckpointResponse { @@ -104,16 +108,17 @@ pub async fn write_checkpoint( /// Posts a checkpoint request to the PowerSync sync service. pub async fn checkpoint_request( db: &InnerPowerSyncState, - connector: &dyn BackendConnector, + authenticator: &EndpointAndAuthenticator, body: &CheckpointRequestPayload, ) -> Result { - if let Some(future) = - connector.post_checkpoint_request(&body.client_id, body.checkpoint_request_id) + if let Some(future) = authenticator + .authenticator + .post_checkpoint_request(&body.client_id, body.checkpoint_request_id) { return future.await; } - let auth = connector.fetch_credentials().await?; + let auth = authenticator.fetch_credentials().await?; let url = auth.parsed_endpoint("sync/checkpoint-request")?; let body = serde_json::to_vec(body)?; @@ -139,7 +144,7 @@ pub async fn checkpoint_request( .into()); } - check_ok(response.status)?; + check_ok(response.status, authenticator.authenticator.as_ref())?; #[derive(Deserialize)] struct CheckpointRequestResponse { @@ -159,10 +164,13 @@ pub async fn checkpoint_request( Ok(response.data.checkpoint_request_id) } -fn check_ok(code: u16) -> Result<(), PowerSyncError> { +fn check_ok(code: u16, authenticator: &dyn Authenticator) -> Result<(), PowerSyncError> { match code { 200 => Ok(()), - 401 => Err(RawPowerSyncError::InvalidCredentials.into()), + 401 => { + authenticator.invalidate_credentials(); + Err(RawPowerSyncError::InvalidCredentials.into()) + } _ => Err(RawPowerSyncError::UnexpectedStatusCode { code }.into()), } } @@ -257,11 +265,22 @@ mod tests { environment, Schema::default().into(), )); - let credentials = PowerSyncCredentials { - endpoint: "https://rust.unit.test.powersync.com/".to_string(), - token: "token".to_string(), - }; - let mut events = Box::pin(sync_stream(db, credentials, "{}".to_string())); + let endpoint = "https://rust.unit.test.powersync.com/"; + + struct StaticAuthenticator; + + #[async_trait] + impl Authenticator for StaticAuthenticator { + async fn resolve_credentials(&self) -> Result, PowerSyncError> { + Ok(Arc::new("token".to_string())) + } + } + + let authenticator = Arc::new(EndpointAndAuthenticator { + endpoint: endpoint.to_string(), + authenticator: Box::new(StaticAuthenticator), + }); + let mut events = Box::pin(sync_stream(db, authenticator, "{}".to_string())); future::block_on(events.as_mut().try_next()) } diff --git a/powersync/src/sync/download/mod.rs b/powersync/src/sync/download/mod.rs index 727198c..2f39661 100644 --- a/powersync/src/sync/download/mod.rs +++ b/powersync/src/sync/download/mod.rs @@ -19,8 +19,12 @@ pub async fn download_loop( channels.checkpoints.disconnected(); }; + let Some(ref endpoint) = options.endpoint else { + return; + }; + loop { - let download = DownloadClient::new(db.clone(), &channels, &events, &options); + let download = DownloadClient::new(db.clone(), &channels, &events, &options, endpoint); let delay_retry = match download.run().await { Ok(end) => !end.hide_disconnect, Err(e) => { diff --git a/powersync/src/sync/download/sync_iteration.rs b/powersync/src/sync/download/sync_iteration.rs index 502031b..bc125f0 100644 --- a/powersync/src/sync/download/sync_iteration.rs +++ b/powersync/src/sync/download/sync_iteration.rs @@ -9,13 +9,12 @@ use serde::Serialize; use serde_json::Map; use serde_json::value::RawValue; -use crate::BackendConnector; use crate::db::connection::{SqliteConnection, TransactionGuard}; use crate::schema::SchemaOrCustom; use crate::sync::coordinator::SyncChannels; use crate::sync::download::http::checkpoint_request; use crate::sync::instruction::CheckpointRequestPayload; -use crate::sync::options::CheckpointMode; +use crate::sync::options::{CheckpointMode, EndpointAndAuthenticator}; use crate::{ SyncOptions, db::internal::InnerPowerSyncState, @@ -32,6 +31,7 @@ pub struct DownloadClient<'a> { channels: &'a SyncChannels, receive_commands: &'a async_channel::Receiver, options: &'a SyncOptions, + endpoint: &'a Arc, stream: Option>>, checkpoint_seed: Option>>, @@ -43,12 +43,14 @@ impl<'a> DownloadClient<'a> { channels: &'a SyncChannels, events: &'a async_channel::Receiver, options: &'a SyncOptions, + endpoint: &'a Arc, ) -> Self { Self { db, channels, receive_commands: events, options, + endpoint, stream: None, checkpoint_seed: None, } @@ -117,7 +119,7 @@ impl<'a> DownloadClient<'a> { if let Some(seed_request) = checkpoint_request { let state = self.channels.checkpoints.clone(); - let connector = self.options.connector.clone(); + let connector = self.endpoint.clone(); let db = Arc::clone(&self.db); self.checkpoint_seed = Some( @@ -136,7 +138,7 @@ impl<'a> DownloadClient<'a> { Arc::clone(&self.db), &mut self.stream, request, - self.options, + self.endpoint.clone(), ) .await?; @@ -168,21 +170,20 @@ impl<'a> DownloadClient<'a> { db: Arc, stream: &mut Option>>, request: Box, - options: &SyncOptions, + endpoint: Arc, ) -> Result<(), PowerSyncError> { - let credentials = options.connector.fetch_credentials().await?; let request = request.get().to_string(); - *stream = Some(sync_stream(db, credentials, request).boxed()); + *stream = Some(sync_stream(db, endpoint, request).boxed()); Ok(()) } async fn seed_checkpoint_state( db: Arc, - connector: Arc, + connector: Arc, request: CheckpointRequestPayload, ) -> Result<(), PowerSyncError> { - let response = checkpoint_request(&db, connector.as_ref(), &request).await?; + let response = checkpoint_request(&db, &connector, &request).await?; db.seed_checkpoint_request_id(response).await?; Ok(()) diff --git a/powersync/src/sync/options.rs b/powersync/src/sync/options.rs index 3cb1c9e..7b1ed20 100644 --- a/powersync/src/sync/options.rs +++ b/powersync/src/sync/options.rs @@ -2,13 +2,19 @@ use std::{sync::Arc, time::Duration}; use futures_lite::future::yield_now; -use crate::{env::PowerSyncEnvironment, error::PowerSyncError, sync::connector::BackendConnector}; +use crate::{ + MutationUploader, + env::PowerSyncEnvironment, + error::PowerSyncError, + sync::connector::{Authenticator, PowerSyncCredentials}, +}; /// Options controlling how PowerSync connects to a sync service. #[derive(Clone)] pub struct SyncOptions { - /// The connector to fetch credentials from. - pub(crate) connector: Arc, + /// The authenticator to fetch credentials from, if downloading is enabled. + pub(crate) endpoint: Option>, + pub(crate) uploader: Option>, /// Whether to sync `auto_subscribe: true` streams automatically. pub(crate) include_default_streams: bool, /// The retry delay between sync iterations on errors. @@ -19,16 +25,50 @@ pub struct SyncOptions { } impl SyncOptions { - /// Creates new [SyncOptions] with default options given the [BackendConnector]. - pub fn new(connector: impl BackendConnector + 'static) -> Self { + fn empty() -> Self { Self { - connector: Arc::new(connector), + endpoint: None, + uploader: None, include_default_streams: true, retry_delay: Duration::from_secs(5), checkpoints: CheckpointMode::default(), } } + /// Creates new [SyncOptions] with default options given the [Authenticator] and + /// [MutationUploader]. + pub fn new( + endpoint: &str, + authenticator: impl Authenticator + 'static, + uploader: impl MutationUploader + 'static, + ) -> Self { + let mut downloads = Self::download_only(endpoint, authenticator); + downloads.uploader = Some(Arc::new(uploader)); + downloads + } + + /// Creates new [SyncOptions] for downloading only. + pub fn download_only(endpoint: &str, authenticator: impl Authenticator + 'static) -> Self { + Self { + endpoint: Some(Arc::new(EndpointAndAuthenticator { + endpoint: endpoint.to_owned(), + authenticator: Box::new(authenticator), + })), + ..Self::empty() + } + } + + /// Creates new sync options for uploading only. + /// + /// When connecting with these options, the sync client won't attempt to connect to a PowerSync + /// service. + pub fn upload_only(uploader: impl MutationUploader + 'static) -> Self { + Self { + uploader: Some(Arc::new(uploader)), + ..Self::empty() + } + } + /// Whether to sync streams that have `auto_subscribe: true`. /// /// This is enabled by default. @@ -70,6 +110,23 @@ impl SyncOptions { } } +pub(crate) struct EndpointAndAuthenticator { + pub endpoint: String, + pub authenticator: Box, +} + +impl EndpointAndAuthenticator { + pub async fn fetch_credentials<'a>( + &'a self, + ) -> Result, PowerSyncError> { + let jwt = self.authenticator.resolve_credentials().await?; + Ok(PowerSyncCredentials { + endpoint: &self.endpoint, + token: jwt, + }) + } +} + #[derive(Clone, Debug, Default, PartialEq, Eq)] pub enum CheckpointMode { /// Uses a legacy endpoint to request checkpoints after uploading data. diff --git a/powersync/src/sync/upload.rs b/powersync/src/sync/upload.rs index 47b6519..67c58fc 100644 --- a/powersync/src/sync/upload.rs +++ b/powersync/src/sync/upload.rs @@ -9,12 +9,13 @@ use log::{debug, info, warn}; use powersync_sqlite_nostd::{Destructor, ResultCode}; use crate::{ - BackendConnector, SyncOptions, + MutationUploader, SyncOptions, db::connection::{SqliteConnection, TransactionGuard}, error::RawPowerSyncError, sync::{ - download::http::checkpoint_request, instruction::CheckpointRequestPayload, - options::CheckpointMode, + download::http::checkpoint_request, + instruction::CheckpointRequestPayload, + options::{CheckpointMode, EndpointAndAuthenticator}, }, }; use crate::{ @@ -30,6 +31,24 @@ pub async fn crud_upload_loop( channels: SyncChannels, trigger_uploads: Receiver<()>, ) { + let common: CommonCrudUpload<'_> = CommonCrudUpload { + options: &options, + channels: &channels, + db: db.as_ref(), + }; + + let uploader = match (&options.endpoint, &options.uploader) { + (_, Some(uploader)) => uploader, + (Some(endpoint), None) => { + // We can't upload data, but we're connected for downloads. We might have been connected + // for uploads before, and thus need to request a checkpoint before syncing completed + // uploads. + common.request_target_checkpoint_once(endpoint).await; + return; + } + _ => return, + }; + let mut tables = HashSet::new(); tables.insert("ps_crud".to_string()); @@ -52,9 +71,8 @@ pub async fn crud_upload_loop( ); let mut upload = CrudUpload { - options: &options, - db: &db, - channels: &channels, + common, + uploader: uploader.as_ref(), }; upload.run().await; next_trigger.await; @@ -63,7 +81,7 @@ pub async fn crud_upload_loop( pub async fn post_checkpoint_request( client_id: String, - connector: &dyn BackendConnector, + authenticator: &EndpointAndAuthenticator, channels: &SyncChannels, db: &InnerPowerSyncState, ) -> Result { @@ -76,7 +94,7 @@ pub async fn post_checkpoint_request( let checkpoint_request_id = db.next_checkpoint_request_id().await?; checkpoint_request( db, - connector, + authenticator, &CheckpointRequestPayload { client_id, checkpoint_request_id, @@ -85,17 +103,23 @@ pub async fn post_checkpoint_request( .await } -struct CrudUpload<'a> { +#[derive(Clone, Copy)] +struct CommonCrudUpload<'a> { options: &'a SyncOptions, channels: &'a SyncChannels, db: &'a InnerPowerSyncState, } +struct CrudUpload<'a> { + common: CommonCrudUpload<'a>, + uploader: &'a dyn MutationUploader, +} + impl<'a> CrudUpload<'a> { pub async fn run(&mut self) { let mut last_item_id = None::; scopeguard::defer! { - self.db.status.update(|s| s.set_upload_state(UploadStatus::Idle)); + self.common.db.status.update(|s| s.set_upload_state(UploadStatus::Idle)); } // Invoke upload method on connector until there are no remaining CRUD items to upload. @@ -107,10 +131,10 @@ impl<'a> CrudUpload<'a> { last_item_id = None; info!("CRUD uploads failed, will retry, {e}"); - self.db - .status + let db = self.common.db; + db.status .update(|data| data.set_upload_state(UploadStatus::Error(e))); - self.options.retry_delay(&self.db.env).await; + self.common.options.retry_delay(&db.env).await; } } } @@ -122,19 +146,15 @@ impl<'a> CrudUpload<'a> { ) -> Result, PowerSyncError> { let Some(item) = self.oldest_crud_item_id().await? else { // Uploading is completed, advance write checkpoint. - if let Some(advance_target) = self.sequence_for_checkpoint().await? { - let write_checkpoint = self.get_write_checkpoint().await?; - advance_target.complete(write_checkpoint, &self.db).await?; + if let Some(ref endpoint) = self.common.options.endpoint { + self.common.request_checkpoint_if_needed(endpoint).await?; } - // It's possible that pending CRUD uploads were preventing data from syncing. So now - // that that's completed, notify the download client in case it needs to retry. - self.channels.mark_crud_uploads_completed().await; - return Ok(ControlFlow::Break(())); }; - self.db + self.common + .db .status .update(|data| data.set_upload_state(UploadStatus::Uploading)); if matches!(*last_item_id, Some(x) if x == item) { @@ -145,36 +165,16 @@ impl<'a> CrudUpload<'a> { } *last_item_id = Some(item); - self.options.connector.upload_data().await?; + self.uploader.upload().await?; Ok(ControlFlow::Continue(())) } async fn oldest_crud_item_id(&self) -> Result, PowerSyncError> { - let reader = self.db.reader().await?; + let reader = self.common.db.reader().await?; Self::read_oldest_crud_item_id(reader.sqlite_connection()) } - async fn get_write_checkpoint(&self) -> Result { - let client_id = get_client_id(&self.db).await?; - - match self.options.checkpoints { - CheckpointMode::Legacy => { - let credentials = self.options.connector.fetch_credentials().await?; - write_checkpoint(&self.db, &client_id, credentials).await - } - CheckpointMode::Requests(_) => { - post_checkpoint_request( - client_id, - self.options.connector.as_ref(), - &self.channels, - &self.db, - ) - .await - } - } - } - fn read_oldest_crud_item_id(conn: &SqliteConnection) -> Result, PowerSyncError> { let stmt = conn.prepare("SELECT id FROM ps_crud ORDER BY id LIMIT 1")?; @@ -184,6 +184,45 @@ impl<'a> CrudUpload<'a> { }) } + const DUPLICATE_ITEM_WARNING: &'static str = " +Potentially previously uploaded CRUD entries are still present in the upload queue. +Make sure to handle uploads and complete CRUD transactions or batches by calling and awaiting their +`complete()` method. +The next upload iteration will be delayed."; +} + +impl<'a> CommonCrudUpload<'a> { + async fn request_checkpoint_if_needed( + &self, + endpoint: &EndpointAndAuthenticator, + ) -> Result<(), PowerSyncError> { + let did_request_checkpoint = + if let Some(advance_target) = self.sequence_for_checkpoint().await? { + let write_checkpoint = self.get_write_checkpoint(endpoint).await?; + advance_target.complete(write_checkpoint, &self.db).await?; + }; + + // It's possible that pending CRUD uploads were preventing data from syncing. So now + // that that's completed, notify the download client in case it needs to retry. + self.channels.mark_crud_uploads_completed().await; + + Ok(did_request_checkpoint) + } + + async fn get_write_checkpoint( + &self, + endpoint: &EndpointAndAuthenticator, + ) -> Result { + let client_id = get_client_id(&self.db).await?; + + match self.options.checkpoints { + CheckpointMode::Legacy => write_checkpoint(&self.db, &client_id, endpoint).await, + CheckpointMode::Requests(_) => { + post_checkpoint_request(client_id, endpoint, &self.channels, &self.db).await + } + } + } + fn ps_crud_sequence(tx: &TransactionGuard) -> Result, PowerSyncError> { let seq_before = tx .inner @@ -216,11 +255,25 @@ impl<'a> CrudUpload<'a> { })) } - const DUPLICATE_ITEM_WARNING: &'static str = " -Potentially previously uploaded CRUD entries are still present in the upload queue. -Make sure to handle uploads and complete CRUD transactions or batches by calling and awaiting their -`complete()` method. -The next upload iteration will be delayed."; + /// Tries requesting a checkpoint until that is successful. + async fn request_target_checkpoint_once(&self, endpoint: &EndpointAndAuthenticator) { + scopeguard::defer! { + self.db.status.update(|s| s.set_upload_state(UploadStatus::Idle)); + } + + loop { + match self.request_checkpoint_if_needed(endpoint).await { + Ok(()) => break, + Err(e) => { + self.db + .status + .update(|s| s.set_upload_state(UploadStatus::Error(e))); + + self.options.retry_delay(&self.db.env).await; + } + } + } + } } struct PendingCheckpointRequest { @@ -243,8 +296,8 @@ impl PendingCheckpointRequest { return Ok(()); } - let seq_after = - CrudUpload::ps_crud_sequence(&writer)?.expect("sqlite sequence should not be empty"); + let seq_after = CommonCrudUpload::ps_crud_sequence(&writer)? + .expect("sqlite sequence should not be empty"); if seq_after != self.crud_sequence { debug!( diff --git a/powersync/tests/sync_test.rs b/powersync/tests/sync_test.rs index bee0f06..b8187f5 100644 --- a/powersync/tests/sync_test.rs +++ b/powersync/tests/sync_test.rs @@ -13,14 +13,13 @@ use futures_lite::{ future::{self, yield_now}, }; use powersync::{ - BackendConnector, CheckpointError, CheckpointMode, DisconnectAndClearFlags, - PowerSyncCredentials, PowerSyncDatabase, RequestsCheckpointMode, StreamPriority, - StreamSubscription, StreamSubscriptionOptions, SyncOptions, SyncStatusData, - error::PowerSyncError, + Authenticator, CheckpointError, CheckpointMode, DisconnectAndClearFlags, MutationUploader, + PowerSyncDatabase, RequestsCheckpointMode, StreamPriority, StreamSubscription, + StreamSubscriptionOptions, SyncOptions, SyncStatusData, error::PowerSyncError, }; use powersync_test_utils::{ DatabaseTest, execute, - mock_sync_service::TestConnector, + mock_sync_service::{TestConnector, WriteCheckpointResponse}, query_all, sync_line::{Checkpoint, SyncLine}, }; @@ -53,7 +52,8 @@ impl SyncStreamTest { } fn connect_options(&self, configure: impl FnOnce(&mut SyncOptions)) { - let mut options = SyncOptions::new(TestConnector::default()); + let connector = TestConnector::default(); + let mut options = SyncOptions::new(Self::TEST_ENDPOINT, connector.clone(), connector); configure(&mut options); self.run(self.db.connect(options)) @@ -88,6 +88,8 @@ impl SyncStreamTest { }) .await } + + const TEST_ENDPOINT: &str = "https://rust-unit-test.powersync.com"; } #[test] @@ -357,15 +359,8 @@ fn upload_retry() { struct FirstUploadFailure; #[async_trait] - impl BackendConnector for FailOnFirstUpload { - async fn fetch_credentials(&self) -> Result { - Ok(PowerSyncCredentials { - endpoint: "https://rust.unit.test.powersync.com/".to_string(), - token: "token".to_string(), - }) - } - - async fn upload_data(&self) -> Result<(), PowerSyncError> { + impl MutationUploader for FailOnFirstUpload { + async fn upload(&self) -> Result<(), PowerSyncError> { let Some(tx) = self.db.next_crud_transaction().await? else { return Ok(()); }; @@ -384,11 +379,15 @@ fn upload_retry() { let sync = SyncStreamTest::new(); let upload_counter = Arc::new(AtomicUsize::default()); let event = Arc::new(Event::new()); - let mut options = SyncOptions::new(FailOnFirstUpload { - db: sync.db.clone(), - counter: upload_counter.clone(), - completed_second: event.clone(), - }); + let mut options = SyncOptions::new( + SyncStreamTest::TEST_ENDPOINT, + TestConnector::default(), + FailOnFirstUpload { + db: sync.db.clone(), + counter: upload_counter.clone(), + completed_second: event.clone(), + }, + ); options.with_retry_delay(Duration::ZERO); // We can't use timers in tests sync.run(sync.db.connect(options)); @@ -431,28 +430,24 @@ fn fetching_credentials_does_not_hold_the_download_writer_lease() { } #[async_trait] - impl BackendConnector for WriterUsingConnector { - async fn fetch_credentials(&self) -> Result { + impl Authenticator for WriterUsingConnector { + async fn resolve_credentials(&self) -> Result, PowerSyncError> { self.entered.send(()).await.unwrap(); self.release.recv().await.unwrap(); - Ok(PowerSyncCredentials { - endpoint: "https://rust.unit.test.powersync.com/".to_string(), - token: "token".to_string(), - }) - } - - async fn upload_data(&self) -> Result<(), PowerSyncError> { - Ok(()) + Ok(Arc::new("token".to_string())) } } let sync = SyncStreamTest::new(); let (entered_tx, entered_rx) = async_channel::bounded(1); let (release_tx, release_rx) = async_channel::bounded(1); - sync.run(sync.db.connect(SyncOptions::new(WriterUsingConnector { - entered: entered_tx, - release: release_rx, - }))); + sync.run(sync.db.connect(SyncOptions::download_only( + SyncStreamTest::TEST_ENDPOINT, + WriterUsingConnector { + entered: entered_tx, + release: release_rx, + }, + ))); sync.run(async { entered_rx.recv().await.unwrap(); @@ -530,20 +525,13 @@ fn reconnects_on_failure() { #[test] fn requests_checkpoints_for_updates() { - struct TestConnector { + struct TestUploader { db: PowerSyncDatabase, } #[async_trait] - impl BackendConnector for TestConnector { - async fn fetch_credentials(&self) -> Result { - Ok(PowerSyncCredentials { - endpoint: "https://rust.unit.test.powersync.com/".to_string(), - token: "token".to_string(), - }) - } - - async fn upload_data(&self) -> Result<(), PowerSyncError> { + impl MutationUploader for TestUploader { + async fn upload(&self) -> Result<(), PowerSyncError> { let Some(tx) = self.db.next_crud_transaction().await? else { return Ok(()); }; @@ -554,9 +542,13 @@ fn requests_checkpoints_for_updates() { } let sync = SyncStreamTest::new(); - let mut options = SyncOptions::new(TestConnector { - db: sync.db.clone(), - }); + let mut options = SyncOptions::new( + SyncStreamTest::TEST_ENDPOINT, + TestConnector::default(), + TestUploader { + db: sync.db.clone(), + }, + ); options.with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); options.with_retry_delay(Duration::ZERO); sync.run(sync.db.connect(options)); @@ -657,15 +649,8 @@ fn download_is_retried_on_checkpoint_request() { } #[async_trait] - impl BackendConnector for Connector { - async fn fetch_credentials(&self) -> Result { - Ok(PowerSyncCredentials { - endpoint: "https://rust.unit.test.powersync.com/".to_string(), - token: "token".to_string(), - }) - } - - async fn upload_data(&self) -> Result<(), PowerSyncError> { + impl MutationUploader for Connector { + async fn upload(&self) -> Result<(), PowerSyncError> { let tx = self.db.next_crud_transaction().await?; if let Some(tx) = tx { tx.complete().await?; @@ -676,9 +661,13 @@ fn download_is_retried_on_checkpoint_request() { } let sync = SyncStreamTest::new(); - let mut options = SyncOptions::new(Connector { - db: sync.db.clone(), - }); + let mut options = SyncOptions::new( + SyncStreamTest::TEST_ENDPOINT, + TestConnector::default(), + Connector { + db: sync.db.clone(), + }, + ); options.with_retry_delay(Duration::from_hours(1)); options.with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); @@ -712,19 +701,23 @@ fn download_is_retried_on_checkpoint_request() { } #[test] -fn can_use_checkpoint_method_from_connector() { +fn can_use_checkpoint_method_from_authenticator() { let sync = SyncStreamTest::new(); let did_request_checkpoint = Event::new(); let listener = did_request_checkpoint.listen(); - let mut options = SyncOptions::new(TestConnector { - post_checkpoint_request: Box::new(move |request_id| { - assert_eq!(request_id, 1); - - did_request_checkpoint.notify(1); - return Some(async move { Ok(request_id) }.boxed()); - }), - }); + let mut options = SyncOptions::download_only( + SyncStreamTest::TEST_ENDPOINT, + TestConnector { + post_checkpoint_request: Arc::new(move |request_id| { + assert_eq!(request_id, 1); + + did_request_checkpoint.notify(1); + return Some(async move { Ok(request_id) }.boxed()); + }), + token: Arc::new("token".to_string()), + }, + ); options.with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); options.with_retry_delay(Duration::ZERO); sync.run(sync.db.connect(options)); @@ -1047,3 +1040,123 @@ fn disconnect_and_clear_soft() { ); }); } + +/// Creates a [MutationUploader] that completes the next CRUD transaction and then notifies +/// `did_upload`. +fn completing_uploader( + db: &PowerSyncDatabase, + did_upload: &Arc, +) -> impl MutationUploader + 'static { + let db = db.clone(); + let did_upload = did_upload.clone(); + + move || { + let db = db.clone(); + let did_upload = did_upload.clone(); + + async move { + if let Some(tx) = db.next_crud_transaction().await? { + tx.complete().await?; + did_upload.notify(usize::MAX); + } + + Ok(()) + } + } +} + +#[test] +fn connect_download_only() { + let sync = SyncStreamTest::new(); + let did_upload = Arc::new(Event::new()); + let upload_listener = did_upload.listen(); + + sync.run(async { + sync.db + .connect(SyncOptions::download_only( + SyncStreamTest::TEST_ENDPOINT, + TestConnector::default(), + )) + .await; + let _request = sync.test.http.receive_requests.recv().await.unwrap(); + sync.wait_for_status(|s| s.is_connected()).await; + + // Create a local mutation, which is never uploaded. + execute( + &sync.db, + "INSERT INTO users (id, name) VALUES (uuid(), ?)", + params!["local mutation"], + ) + .await; + sync.test.advance_time(Duration::from_hours(1)); + assert_eq!( + query_all(&sync.db, "SELECT * FROM ps_crud", params![]) + .await + .as_array() + .unwrap() + .len(), + 1 + ); + + // Upgrade to a full connection, which should upload. + sync.db + .connect(SyncOptions::new( + SyncStreamTest::TEST_ENDPOINT, + TestConnector::default(), + completing_uploader(&sync.db, &did_upload), + )) + .await; + upload_listener.await; + }); +} + +#[test] +fn connect_upload_only() { + let sync = SyncStreamTest::new(); + let did_upload = Arc::new(Event::new()); + let upload_listener = did_upload.listen(); + + let write_checkpoint_requests = Arc::new(AtomicUsize::new(0)); + *sync.test.http.write_checkpoints.lock().unwrap() = Box::new({ + let requests = write_checkpoint_requests.clone(); + move || { + requests.fetch_add(1, Ordering::SeqCst); + WriteCheckpointResponse::new("1".to_string()) + } + }); + + sync.run(async { + sync.db + .connect(SyncOptions::upload_only(completing_uploader( + &sync.db, + &did_upload, + ))) + .await; + + execute( + &sync.db, + "INSERT INTO users (id, name) VALUES (uuid(), ?)", + params!["local mutation"], + ) + .await; + upload_listener.await; + + // Connecting in upload-only mode should make no SDK-initiated HTTP requests. + sync.test.advance_time(Duration::from_hours(1)); + assert!(sync.test.http.receive_requests.is_empty()); + assert_eq!(write_checkpoint_requests.load(Ordering::SeqCst), 0); + + // Reconnect in download-only mode. This should request a write checkpoint because the + // upload-only mode can't. + sync.db + .connect(SyncOptions::download_only( + SyncStreamTest::TEST_ENDPOINT, + TestConnector::default(), + )) + .await; + let _request = sync.test.http.receive_requests.recv().await.unwrap(); + sync.wait_for_status(|s| s.is_connected()).await; + sync.test.advance_time(Duration::from_secs(1)); + assert_eq!(write_checkpoint_requests.load(Ordering::SeqCst), 1); + }); +} diff --git a/powersync_test_utils/src/mock_sync_service.rs b/powersync_test_utils/src/mock_sync_service.rs index 30d6ceb..0e7ec67 100644 --- a/powersync_test_utils/src/mock_sync_service.rs +++ b/powersync_test_utils/src/mock_sync_service.rs @@ -5,7 +5,8 @@ use futures_lite::future::Boxed; use futures_lite::{FutureExt, Stream, StreamExt, ready, stream}; use pin_project_lite::pin_project; use powersync::http::{HttpClient, Request, Response, ResponseBody}; -use powersync::{BackendConnector, PowerSyncCredentials, StreamPriority, error::PowerSyncError}; +use powersync::{Authenticator, MutationUploader}; +use powersync::{StreamPriority, error::PowerSyncError}; use serde::{Deserialize, Serialize}; use serde_json::json; use serde_with::{DisplayFromStr, serde_as}; @@ -290,33 +291,29 @@ pub struct WriteCheckpointResponseData { pub write_checkpoint: String, } +#[derive(Clone)] pub struct TestConnector { - pub post_checkpoint_request: Box< + pub post_checkpoint_request: Arc< dyn Fn(i64) -> Option> + Send>>> + Send + Sync, >, + pub token: Arc, } impl Default for TestConnector { fn default() -> Self { Self { - post_checkpoint_request: Box::new(|_| None), + post_checkpoint_request: Arc::new(|_| None), + token: Arc::new("token".to_string()), } } } #[async_trait] -impl BackendConnector for TestConnector { - async fn fetch_credentials(&self) -> Result { - Ok(PowerSyncCredentials { - endpoint: "https://rust.unit.test.powersync.com/".to_string(), - token: "token".to_string(), - }) - } - - async fn upload_data(&self) -> Result<(), PowerSyncError> { - Ok(()) +impl Authenticator for TestConnector { + async fn resolve_credentials(&self) -> Result, PowerSyncError> { + Ok(self.token.clone()) } fn post_checkpoint_request<'a>( @@ -327,3 +324,10 @@ impl BackendConnector for TestConnector { (self.post_checkpoint_request)(request_id) } } + +#[async_trait] +impl MutationUploader for TestConnector { + async fn upload(&self) -> Result<(), PowerSyncError> { + Ok(()) + } +} From 0e08ae9227fae668a94109d53b1ab169281cad71 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Mon, 5 Oct 2026 15:22:33 +0200 Subject: [PATCH 2/2] Implement prefetching --- powersync/src/sync/download/sync_iteration.rs | 5 +- powersync/tests/sync_test.rs | 53 ++++++++++++++++++- 2 files changed, 54 insertions(+), 4 deletions(-) diff --git a/powersync/src/sync/download/sync_iteration.rs b/powersync/src/sync/download/sync_iteration.rs index bc125f0..382f65c 100644 --- a/powersync/src/sync/download/sync_iteration.rs +++ b/powersync/src/sync/download/sync_iteration.rs @@ -146,9 +146,8 @@ impl<'a> DownloadClient<'a> { self.channels.trigger_crud_upload(); } Instruction::FetchCredentials { .. } => { - // TODO: Pre-fetching credentials - // If did_expire is true, the core extension will also emit a stop - // instruction. So we don't have to handle that separately. + // Credentials either have expired or are about to expire, invalidate. + self.endpoint.authenticator.invalidate_credentials(); } Instruction::CloseSyncStream(close) => { return Ok(Some(close)); diff --git a/powersync/tests/sync_test.rs b/powersync/tests/sync_test.rs index b8187f5..8f93c79 100644 --- a/powersync/tests/sync_test.rs +++ b/powersync/tests/sync_test.rs @@ -1,6 +1,6 @@ use std::{ sync::{ - Arc, + Arc, Mutex, atomic::{AtomicUsize, Ordering}, }, time::{Duration, SystemTime}, @@ -461,6 +461,57 @@ fn fetching_credentials_does_not_hold_the_download_writer_lease() { }); } +#[test] +fn invalidates_cached_credentials_on_token_expiry() { + #[derive(Default)] + struct CachingState { + cached: Mutex>>, + fetches: AtomicUsize, + invalidations: AtomicUsize, + } + + struct CachingAuthenticator(Arc); + + #[async_trait] + impl Authenticator for CachingAuthenticator { + async fn resolve_credentials(&self) -> Result, PowerSyncError> { + let mut cached = self.0.cached.lock().unwrap(); + let token = cached.get_or_insert_with(|| { + let id = self.0.fetches.fetch_add(1, Ordering::SeqCst); + Arc::new(format!("token-{id}")) + }); + + Ok(token.clone()) + } + + fn invalidate_credentials(&self) { + self.0.invalidations.fetch_add(1, Ordering::SeqCst); + *self.0.cached.lock().unwrap() = None; + } + } + + let sync = SyncStreamTest::new(); + let authenticator = Arc::new(CachingState::default()); + let options = SyncOptions::download_only( + SyncStreamTest::TEST_ENDPOINT, + CachingAuthenticator(authenticator.clone()), + ); + sync.run(sync.db.connect(options)); + + sync.run(async { + let request = sync.test.http.receive_requests.recv().await.unwrap(); + assert_eq!(authenticator.fetches.load(Ordering::SeqCst), 1); + assert_eq!(authenticator.invalidations.load(Ordering::SeqCst), 0); + + // A token about to expire makes the core extension request new credentials. + request.send_keepalive(5).await; + + while authenticator.invalidations.load(Ordering::SeqCst) == 0 { + yield_now().await; + } + }); +} + #[test] fn reports_correct_times() { let sync = SyncStreamTest::new();