From dfb1034eed283c42296cd008ad1330def714206c Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Thu, 10 Sep 2026 09:12:24 +0200 Subject: [PATCH 1/4] Start adopting checkpoint request protocol --- Cargo.lock | 1 + powersync/Cargo.toml | 1 + powersync/src/db/internal.rs | 47 ++++- powersync/src/error.rs | 4 + powersync/src/sync/checkpoint.rs | 17 ++ powersync/src/sync/coordinator.rs | 3 + powersync/src/sync/download/http.rs | 45 +++++ powersync/src/sync/download/mod.rs | 8 +- powersync/src/sync/download/sync_iteration.rs | 78 +++++++- powersync/src/sync/instruction.rs | 15 ++ powersync/src/sync/mod.rs | 2 + powersync/src/sync/options.rs | 45 ++++- powersync/src/sync/state.rs | 188 ++++++++++++++++++ powersync/src/sync/upload.rs | 28 ++- 14 files changed, 476 insertions(+), 6 deletions(-) create mode 100644 powersync/src/sync/checkpoint.rs create mode 100644 powersync/src/sync/state.rs diff --git a/Cargo.lock b/Cargo.lock index f39632f..2193341 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3086,6 +3086,7 @@ checksum = "439ee305def115ba05938db6eb1644ff94165c5ab5e9420d1c1bcedbba909391" name = "powersync" version = "0.0.7" dependencies = [ + "async-broadcast", "async-channel", "async-executor", "async-io", diff --git a/powersync/Cargo.toml b/powersync/Cargo.toml index 060180b..e97ce4b 100644 --- a/powersync/Cargo.toml +++ b/powersync/Cargo.toml @@ -24,6 +24,7 @@ rusqlite = ["dep:rusqlite"] ffi = [] [dependencies] +async-broadcast = "0.7" async-channel = "2.5.0" async-lock = "3.4.1" async-executor = { version = "1.14.0", optional = true } diff --git a/powersync/src/db/internal.rs b/powersync/src/db/internal.rs index 4bf5369..6c6d4e8 100644 --- a/powersync/src/db/internal.rs +++ b/powersync/src/db/internal.rs @@ -111,12 +111,38 @@ impl InnerPowerSyncState { writer.commit() } + pub async fn seed_checkpoint_request_id(&self, id: i64) -> Result<(), PowerSyncError> { + let mut writer = self.writer().await?; + let writer = TransactionGuard::new(writer.sqlite_connection_mut())?; + + Self::checkpoint_request_control(&writer, CheckpointCounter::Seed, Some(id))?; + writer.commit() + } + + pub async fn next_checkpoint_request_id(&self) -> Result { + let mut writer = self.writer().await?; + let writer = TransactionGuard::new(writer.sqlite_connection_mut())?; + + let id = Self::checkpoint_request_control(&writer, CheckpointCounter::Next, None)?; + writer.commit()?; + Ok(id.expect("Core extension should return next checkpoint request id")) + } + pub fn target_checkpoint_request_id( writer: &TransactionGuard, update: Option, + ) -> Result, PowerSyncError> { + Self::checkpoint_request_control(writer, CheckpointCounter::Target, update) + } + + fn checkpoint_request_control( + writer: &TransactionGuard, + counter: CheckpointCounter, + update: Option, ) -> Result, PowerSyncError> { let stmt = writer.inner.prepare("SELECT powersync_control(?, ?);")?; - stmt.bind_text(1, "target_checkpoint_request_id", Destructor::STATIC)?; + + stmt.bind_text(1, counter.control_op(), Destructor::STATIC)?; if let Some(update) = update { stmt.bind_int64(2, update)?; } else { @@ -205,3 +231,22 @@ impl InnerPowerSyncState { } } } + +#[derive(Clone, Copy)] +pub enum CheckpointCounter { + Target, + Seed, + Current, + Next, +} + +impl CheckpointCounter { + fn control_op(&self) -> &'static str { + match self { + CheckpointCounter::Target => "target_checkpoint_request_id", + CheckpointCounter::Seed => "seed_checkpoint_request_id", + CheckpointCounter::Current => "current_checkpoint_request_id", + CheckpointCounter::Next => "next_checkpoint_request_id", + } + } +} diff --git a/powersync/src/error.rs b/powersync/src/error.rs index f15f817..7ec468b 100644 --- a/powersync/src/error.rs +++ b/powersync/src/error.rs @@ -5,6 +5,8 @@ use std::sync::Arc; use std::{borrow::Cow, fmt::Display}; use thiserror::Error; +use crate::sync::checkpoint::CheckpointError; + pub type Result = std::result::Result; /// A [RawPowerSyncError], but boxed. @@ -122,6 +124,8 @@ pub(crate) enum RawPowerSyncError { #[source] source: Box, }, + #[error("Checkpoint error: {error}")] + Checkpoint { error: CheckpointError }, } impl From for PowerSyncError { diff --git a/powersync/src/sync/checkpoint.rs b/powersync/src/sync/checkpoint.rs new file mode 100644 index 0000000..a7a10d1 --- /dev/null +++ b/powersync/src/sync/checkpoint.rs @@ -0,0 +1,17 @@ +use thiserror::Error; + +use crate::error::PowerSyncError; + +#[derive(Error, Debug)] +pub enum CheckpointError { + #[error( + "The PowerSync service does not support checkpoint requests. Update to PowerSync service version 1.24.0 or later to use this API." + )] + InstanceNotSupported, + #[error("Cannot request checkpoints, sync client is disconnected")] + Disconnected, + #[error("Connected with legacy checkpoint mode, cannot request checkpoints")] + Disabled, + #[error("Error on sync status before checkpoint was applied: {cause}")] + StatusError { cause: PowerSyncError }, +} diff --git a/powersync/src/sync/coordinator.rs b/powersync/src/sync/coordinator.rs index cd460b2..5a8de28 100644 --- a/powersync/src/sync/coordinator.rs +++ b/powersync/src/sync/coordinator.rs @@ -10,6 +10,7 @@ use crate::{ error::PowerSyncError, sync::{ download::{DownloadEvent, download_loop}, + state::CheckpointStateSignals, streams::ChangedSyncSubscriptions, upload::crud_upload_loop, }, @@ -133,6 +134,7 @@ impl Drop for SyncTasks { #[derive(Clone)] pub struct SyncChannels { + pub checkpoints: Arc, local_download_events: Sender, trigger_upload: Sender<()>, } @@ -144,6 +146,7 @@ impl SyncChannels { ( Self { + checkpoints: Default::default(), local_download_events: download_send, trigger_upload: uploads_send, }, diff --git a/powersync/src/sync/download/http.rs b/powersync/src/sync/download/http.rs index c3de724..819d6c2 100644 --- a/powersync/src/sync/download/http.rs +++ b/powersync/src/sync/download/http.rs @@ -2,6 +2,7 @@ use std::borrow::Cow; use std::sync::Arc; use crate::http::{Request, Response}; +use crate::sync::instruction::CheckpointRequestPayload; use crate::util::LineSplitter; use crate::{ db::internal::InnerPowerSyncState, @@ -98,6 +99,50 @@ pub async fn write_checkpoint( Ok(response.data.write_checkpoint) } +/// Posts a checkpoint request to the PowerSync sync service. +pub async fn checkpoint_request( + db: &InnerPowerSyncState, + body: &CheckpointRequestPayload, + auth: PowerSyncCredentials, +) -> Result { + let url = auth.parsed_endpoint("sync/checkpoint-request")?; + + let body = serde_json::to_vec(body)?; + let request = Request { + method: "POST", + url, + headers: { + let mut headers: Vec<(&str, Cow<'_, str>)> = vec![]; + headers.push(("Content-Type", "application/json".into())); + headers.push(("Authorization", format!("Token {}", auth.token).into())); + headers.push(("Accept", "application/json".into())); + + headers + }, + body: Some(body), + }; + + let response = db.env.client.send(request).await?; + check_ok(response.status)?; + + #[derive(Deserialize)] + struct CheckpointRequestResponse { + data: CheckpointData, + } + + #[serde_as] + #[derive(Deserialize)] + struct CheckpointData { + #[serde_as(as = "DisplayFromStr")] + checkpoint_request_id: i64, + } + + let body_bytes = response.body.read_fully().await?; + + let response: CheckpointRequestResponse = serde_json::from_slice(&body_bytes)?; + Ok(response.data.checkpoint_request_id) +} + fn check_ok(code: u16) -> Result<(), PowerSyncError> { match code { 200 => Ok(()), diff --git a/powersync/src/sync/download/mod.rs b/powersync/src/sync/download/mod.rs index 97983f9..727198c 100644 --- a/powersync/src/sync/download/mod.rs +++ b/powersync/src/sync/download/mod.rs @@ -3,6 +3,7 @@ mod sync_iteration; use std::sync::Arc; +use futures_lite::future; use log::warn; pub use sync_iteration::{DownloadClient, DownloadEvent}; @@ -14,6 +15,10 @@ pub async fn download_loop( options: SyncOptions, events: async_channel::Receiver, ) { + scopeguard::defer! { + channels.checkpoints.disconnected(); + }; + loop { let download = DownloadClient::new(db.clone(), &channels, &events, &options); let delay_retry = match download.run().await { @@ -26,8 +31,9 @@ pub async fn download_loop( } }; + let iteration_ended = channels.checkpoints.download_iteration_ended(); if delay_retry { - options.retry_delay(&db.env).await + future::or(options.retry_delay(&db.env), iteration_ended).await } } } diff --git a/powersync/src/sync/download/sync_iteration.rs b/powersync/src/sync/download/sync_iteration.rs index 688e423..9b2a6f6 100644 --- a/powersync/src/sync/download/sync_iteration.rs +++ b/powersync/src/sync/download/sync_iteration.rs @@ -1,5 +1,7 @@ use std::sync::Arc; +use futures_lite::FutureExt; +use futures_lite::future::Boxed; use futures_lite::{StreamExt, future, stream::Boxed as BoxedStream}; use log::{debug, info, trace, warn}; use powersync_sqlite_nostd::{Destructor, ManagedStmt, ResultCode}; @@ -7,9 +9,13 @@ 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::{ SyncOptions, db::internal::InnerPowerSyncState, @@ -26,7 +32,9 @@ pub struct DownloadClient<'a> { channels: &'a SyncChannels, receive_commands: &'a async_channel::Receiver, options: &'a SyncOptions, + stream: Option>>, + checkpoint_seed: Option>>, } impl<'a> DownloadClient<'a> { @@ -42,6 +50,7 @@ impl<'a> DownloadClient<'a> { receive_commands: events, options, stream: None, + checkpoint_seed: None, } } @@ -51,6 +60,7 @@ impl<'a> DownloadClient<'a> { schema: self.db.schema.clone(), include_defaults: self.options.include_default_streams, active_streams: self.db.current_streams.collect_active_streams(), + checkpoint_mode: CoreCheckpointMode::from(&self.options.checkpoints), }; if let Some(end) = self.handle_event(DownloadEvent::Start(start)).await? { return Ok(end); @@ -61,7 +71,10 @@ impl<'a> DownloadClient<'a> { Some(stream) => { future::or( Self::receive_command(&self.receive_commands), - Self::receive_on_stream(stream), + future::or( + Self::receive_on_stream(stream), + Self::wait_for_checkpoint_seed(&mut self.checkpoint_seed), + ), ) .await } @@ -96,8 +109,29 @@ impl<'a> DownloadClient<'a> { Instruction::UpdateSyncStatus { status } => { self.db.status.update(|s| s.update_from_core(status)) } - Instruction::EstablishSyncStream { request } => { + Instruction::EstablishSyncStream { + request, + checkpoint_request, + } => { trace!("Establishing sync stream with {request}"); + + if let Some(seed_request) = checkpoint_request { + let state = self.channels.checkpoints.clone(); + let connector = self.options.connector.clone(); + let db = Arc::clone(&self.db); + + self.checkpoint_seed = Some( + async move { + let result = + Self::seed_checkpoint_state(db, connector, seed_request).await; + + state.mark_checkpoints_ready(result.clone()); + result + } + .boxed(), + ); + } + Self::establish_sync_stream( Arc::clone(&self.db), &mut self.stream, @@ -143,6 +177,18 @@ impl<'a> DownloadClient<'a> { Ok(()) } + async fn seed_checkpoint_state( + db: Arc, + connector: Arc, + request: CheckpointRequestPayload, + ) -> Result<(), PowerSyncError> { + let credentials = connector.fetch_credentials().await?; + let response = checkpoint_request(&db, &request, credentials).await?; + db.seed_checkpoint_request_id(response).await?; + + Ok(()) + } + async fn receive_command( channel: &async_channel::Receiver, ) -> Result { @@ -157,6 +203,17 @@ impl<'a> DownloadClient<'a> { .await? .unwrap_or(DownloadEvent::ResponseStreamEnd)) } + + async fn wait_for_checkpoint_seed> + Unpin>( + future: &mut Option, + ) -> Result { + if let Some(future) = future { + future.await?; + } + + // This doesn't generate events, we should just poll the future. + future::pending().await + } } /// An event that triggers the downloading client to advance. @@ -269,4 +326,21 @@ pub struct StartDownloadIteration { pub schema: Arc, pub include_defaults: bool, pub active_streams: Vec, + pub checkpoint_mode: CoreCheckpointMode, +} + +#[derive(Serialize, Debug, PartialEq, Eq)] +#[serde(rename_all = "snake_case")] +pub enum CoreCheckpointMode { + Legacy, + Requests, +} + +impl From<&CheckpointMode> for CoreCheckpointMode { + fn from(value: &CheckpointMode) -> Self { + match value { + CheckpointMode::Legacy => Self::Legacy, + CheckpointMode::Requests(_) => Self::Requests, + } + } } diff --git a/powersync/src/sync/instruction.rs b/powersync/src/sync/instruction.rs index 6ae570b..42ca2c1 100644 --- a/powersync/src/sync/instruction.rs +++ b/powersync/src/sync/instruction.rs @@ -2,6 +2,7 @@ use std::time::{Duration, SystemTime}; use serde::{Deserialize, Serialize, de::IgnoredAny}; use serde_json::value::RawValue; +use serde_with::{DisplayFromStr, serde_as}; use crate::{sync::progress::ProgressCounters, util::SerializedJsonObject}; @@ -20,6 +21,7 @@ pub enum Instruction { /// and then forward received lines via [SyncEvent::TextLine] and [SyncEvent::BinaryLine]. EstablishSyncStream { request: Box, + checkpoint_request: Option, }, FetchCredentials { /// Whether the credentials currently used have expired. @@ -56,6 +58,7 @@ pub enum LogSeverity { } /// Information about a progressing download. +#[serde_as] #[derive(Deserialize, Default, Debug)] pub struct DownloadSyncStatus { /// Whether the socket to the sync service is currently open and connected. @@ -68,6 +71,10 @@ pub struct DownloadSyncStatus { pub connecting: bool, pub streams: Vec, pub downloading: Option, + + #[serde(rename = "internal_last_applied_checkpoint_request_id")] + #[serde_as(as = "Option")] + pub last_applied_checkpoint_request: Option, } #[derive(Deserialize, Serialize, Debug)] @@ -93,3 +100,11 @@ impl From for SystemTime { SystemTime::UNIX_EPOCH + since_epoch } } + +#[serde_as] +#[derive(Debug, Deserialize, Serialize)] +pub struct CheckpointRequestPayload { + pub client_id: String, + #[serde_as(as = "DisplayFromStr")] + pub checkpoint_request_id: i64, +} diff --git a/powersync/src/sync/mod.rs b/powersync/src/sync/mod.rs index d304a11..c288783 100644 --- a/powersync/src/sync/mod.rs +++ b/powersync/src/sync/mod.rs @@ -1,9 +1,11 @@ +pub mod checkpoint; pub mod connector; pub mod coordinator; pub mod download; mod instruction; pub mod options; pub mod progress; +mod state; pub mod status; pub mod stream_priority; pub mod streams; diff --git a/powersync/src/sync/options.rs b/powersync/src/sync/options.rs index 5620156..2d61696 100644 --- a/powersync/src/sync/options.rs +++ b/powersync/src/sync/options.rs @@ -2,7 +2,7 @@ use std::{sync::Arc, time::Duration}; use futures_lite::future::yield_now; -use crate::{env::PowerSyncEnvironment, sync::connector::BackendConnector}; +use crate::{env::PowerSyncEnvironment, error::PowerSyncError, sync::connector::BackendConnector}; /// Options controlling how PowerSync connects to a sync service. #[derive(Clone)] @@ -13,6 +13,9 @@ pub struct SyncOptions { pub(crate) include_default_streams: bool, /// The retry delay between sync iterations on errors. pub(crate) retry_delay: Duration, + + /// How to request checkpoints after completing uploads. + pub(crate) checkpoints: CheckpointMode, } impl SyncOptions { @@ -22,6 +25,7 @@ impl SyncOptions { connector: Arc::new(connector), include_default_streams: true, retry_delay: Duration::from_secs(5), + checkpoints: CheckpointMode::default(), } } @@ -57,3 +61,42 @@ impl SyncOptions { } } } + +#[derive(Clone, Debug, Default, PartialEq, Eq)] +pub enum CheckpointMode { + #[default] + Legacy, + Requests(RequestsCheckpointMode), +} + +#[derive(Clone, Debug, PartialEq, Eq)] +pub struct RequestsCheckpointMode { + pub(crate) retry_delay: Duration, +} + +impl RequestsCheckpointMode { + const DEFAULT_RETRY: Duration = Duration::from_secs(10); + const MIN_RETRY: Duration = Self::DEFAULT_RETRY; +} + +impl Default for RequestsCheckpointMode { + fn default() -> Self { + Self { + retry_delay: Self::DEFAULT_RETRY, + } + } +} + +impl TryFrom for RequestsCheckpointMode { + type Error = PowerSyncError; + + fn try_from(value: Duration) -> Result { + if value < Self::MIN_RETRY { + return Err(PowerSyncError::argument_error(format!( + "Minimum retry delay is 10s, got {value:?}" + ))); + } + + Ok(Self { retry_delay: value }) + } +} diff --git a/powersync/src/sync/state.rs b/powersync/src/sync/state.rs new file mode 100644 index 0000000..de87f6f --- /dev/null +++ b/powersync/src/sync/state.rs @@ -0,0 +1,188 @@ +use std::sync::Mutex; + +use crate::{error::PowerSyncError, sync::checkpoint::CheckpointError}; +use async_broadcast::{InactiveReceiver, Sender, broadcast}; +use event_listener::{Event, EventListener}; + +#[derive(Default, Debug, Clone)] +enum CheckpointState { + Disconnected, + #[default] + Pending, + DidSeed(Result<(), PowerSyncError>), +} + +pub struct CheckpointStateSignals { + state_and_waiter: Mutex<(CheckpointState, Option)>, + state_channel: (Sender, InactiveReceiver), +} + +impl Default for CheckpointStateSignals { + fn default() -> Self { + let (mut sender, receiver) = broadcast(1); + sender.set_overflow(true); + + Self { + state_and_waiter: Default::default(), + state_channel: (sender, receiver.deactivate()), + } + } +} + +impl CheckpointStateSignals { + fn notify_state_channel(&self, state: CheckpointState) { + let _ = self.state_channel.0.try_broadcast(state); + } + + fn update_status(&self, state: CheckpointState) { + let mut waiters = self.state_and_waiter.lock().unwrap(); + waiters.0 = state.clone(); + self.notify_state_channel(state); + } + + /// Marks the current download iteration as ended, blocking new checkpoint requests until the + /// seed was performed in the next iteration. + /// + /// Returns a receiver that will receive an event when another actor waits for checkpoints, + /// which allows resuming immediately. + pub fn download_iteration_ended(&self) -> EventListener { + // Checkpoint waiters called after this should be able to resume the download iteration. + let event = Event::default(); + let listener = event.listen(); + + let mut waiters = self.state_and_waiter.lock().unwrap(); + *waiters = (CheckpointState::Pending, Some(event)); + self.notify_state_channel(CheckpointState::Pending); + + listener + } + + /// Marks the sync client as disconnected, failing all outstanding checkpoint requests and + /// preventing new ones. + pub fn disconnected(&self) { + self.update_status(CheckpointState::Disconnected); + } + + pub fn mark_checkpoints_ready(&self, result: Result<(), PowerSyncError>) { + self.update_status(CheckpointState::DidSeed(result)); + } + + /// Waits until a download iteration is active and has seeded the checkpoint state, meaning that + /// checkpoint ids can safely be allocated. + pub async fn wait_for_checkpoint_requests_ready( + &self, + wake_download_loop: bool, + ) -> Result<(), CheckpointError> { + let mut state = { + let waiters = self.state_and_waiter.lock().unwrap(); + waiters.0.clone() + }; + let mut receiver = self.state_channel.1.activate_cloned(); + + loop { + match state { + CheckpointState::DidSeed(result) => { + break result.map_err(|e| CheckpointError::StatusError { cause: e }); + } + CheckpointState::Disconnected => break Err(CheckpointError::Disconnected), + CheckpointState::Pending => { + if wake_download_loop { + let mut waiters = self.state_and_waiter.lock().unwrap(); + if let Some(sender) = waiters.1.take() { + sender.notify(1); + } + } + } + } + + state = receiver + .recv_direct() + .await + .map_err(|_| CheckpointError::Disconnected)?; + } + } +} + +#[cfg(test)] +mod test { + use std::task::Poll; + + use futures_lite::{FutureExt, future::block_on}; + use futures_test::task::noop_context; + + use crate::{error::PowerSyncError, sync::checkpoint::CheckpointError}; + + use super::CheckpointStateSignals; + + #[test] + fn mark_ready_success() { + let signals = CheckpointStateSignals::default(); + signals.mark_checkpoints_ready(Ok(())); + block_on(signals.wait_for_checkpoint_requests_ready(false)).unwrap(); + } + + #[test] + fn mark_ready_failure() { + let signals = CheckpointStateSignals::default(); + signals.mark_checkpoints_ready(Err(PowerSyncError::argument_error("test"))); + block_on(signals.wait_for_checkpoint_requests_ready(false)).unwrap_err(); + } + + #[test] + fn mark_disconnected() { + let signals = CheckpointStateSignals::default(); + + signals.disconnected(); + assert!(matches!( + block_on(signals.wait_for_checkpoint_requests_ready(false)).unwrap_err(), + CheckpointError::Disconnected + )) + } + + #[test] + fn supports_concurrent_waiters() { + let signals = CheckpointStateSignals::default(); + let mut a = signals + .wait_for_checkpoint_requests_ready(true) + .boxed_local(); + let mut b = signals + .wait_for_checkpoint_requests_ready(true) + .boxed_local(); + + assert!(matches!(a.poll(&mut noop_context()), Poll::Pending)); + assert!(matches!(b.poll(&mut noop_context()), Poll::Pending)); + + signals.mark_checkpoints_ready(Ok(())); + block_on(a).unwrap(); + block_on(b).unwrap(); + } + + #[test] + fn waits_for_checkpoint_waiter_is_notified() { + let signals = CheckpointStateSignals::default(); + let notified = signals.download_iteration_ended(); + + let mut future = signals + .wait_for_checkpoint_requests_ready(true) + .boxed_local(); + assert!(matches!(future.poll(&mut noop_context()), Poll::Pending)); + + block_on(notified); + signals.mark_checkpoints_ready(Ok(())); + block_on(future).unwrap(); + } + + #[test] + fn does_not_notify_waiter_when_wake_download_loop_is_false() { + let signals = CheckpointStateSignals::default(); + let mut listener = signals.download_iteration_ended().boxed_local(); + + assert!(matches!(listener.poll(&mut noop_context()), Poll::Pending)); + + let mut future = signals + .wait_for_checkpoint_requests_ready(false) + .boxed_local(); + assert!(matches!(future.poll(&mut noop_context()), Poll::Pending)); + assert!(matches!(listener.poll(&mut noop_context()), Poll::Pending)); + } +} diff --git a/powersync/src/sync/upload.rs b/powersync/src/sync/upload.rs index 5294fee..d1209e4 100644 --- a/powersync/src/sync/upload.rs +++ b/powersync/src/sync/upload.rs @@ -11,6 +11,11 @@ use powersync_sqlite_nostd::{Destructor, ResultCode}; use crate::{ SyncOptions, db::connection::{SqliteConnection, TransactionGuard}, + error::RawPowerSyncError, + sync::{ + download::http::checkpoint_request, instruction::CheckpointRequestPayload, + options::CheckpointMode, + }, }; use crate::{ db::internal::InnerPowerSyncState, @@ -141,7 +146,28 @@ impl<'a> CrudUpload<'a> { }; let credentials = self.options.connector.fetch_credentials().await?; - write_checkpoint(&self.db, &client_id, credentials).await + + match self.options.checkpoints { + CheckpointMode::Legacy => write_checkpoint(&self.db, &client_id, credentials).await, + CheckpointMode::Requests(_) => { + self.channels + .checkpoints + .wait_for_checkpoint_requests_ready(true) + .await + .map_err(|e| RawPowerSyncError::Checkpoint { error: e })?; + + let checkpoint_request_id = self.db.next_checkpoint_request_id().await?; + checkpoint_request( + &self.db, + &CheckpointRequestPayload { + client_id, + checkpoint_request_id, + }, + credentials, + ) + .await + } + } } fn read_oldest_crud_item_id(conn: &SqliteConnection) -> Result, PowerSyncError> { From b0f68db2e9d25cc58ea63940359fd93f9dfbedf1 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 15 Sep 2026 11:11:43 +0200 Subject: [PATCH 2/4] Port a few tests --- powersync/src/db/internal.rs | 2 - powersync/src/error.rs | 6 + powersync/src/lib.rs | 2 +- powersync/src/sync/connector.rs | 22 ++ powersync/src/sync/download/http.rs | 8 + powersync/src/sync/download/sync_iteration.rs | 3 +- powersync/src/sync/options.rs | 5 + powersync/src/sync/upload.rs | 1 + powersync/tests/sync_test.rs | 264 +++++++++++++++++- powersync_test_utils/src/mock_sync_service.rs | 118 +++++++- powersync_test_utils/src/sync_line.rs | 6 + 11 files changed, 420 insertions(+), 17 deletions(-) diff --git a/powersync/src/db/internal.rs b/powersync/src/db/internal.rs index 6c6d4e8..e38de66 100644 --- a/powersync/src/db/internal.rs +++ b/powersync/src/db/internal.rs @@ -236,7 +236,6 @@ impl InnerPowerSyncState { pub enum CheckpointCounter { Target, Seed, - Current, Next, } @@ -245,7 +244,6 @@ impl CheckpointCounter { match self { CheckpointCounter::Target => "target_checkpoint_request_id", CheckpointCounter::Seed => "seed_checkpoint_request_id", - CheckpointCounter::Current => "current_checkpoint_request_id", CheckpointCounter::Next => "next_checkpoint_request_id", } } diff --git a/powersync/src/error.rs b/powersync/src/error.rs index 7ec468b..20c34ea 100644 --- a/powersync/src/error.rs +++ b/powersync/src/error.rs @@ -53,6 +53,12 @@ impl From for PowerSyncError { } } +impl From for PowerSyncError { + fn from(value: CheckpointError) -> Self { + RawPowerSyncError::Checkpoint { error: value }.into() + } +} + impl From for PowerSyncError { fn from(value: RawPowerSyncError) -> Self { PowerSyncError { diff --git a/powersync/src/lib.rs b/powersync/src/lib.rs index 2a68d1b..0279d98 100644 --- a/powersync/src/lib.rs +++ b/powersync/src/lib.rs @@ -12,7 +12,7 @@ pub use db::streams::StreamSubscription; pub use db::streams::StreamSubscriptionOptions; pub use db::streams::SyncStream; pub use sync::connector::{BackendConnector, PowerSyncCredentials}; -pub use sync::options::SyncOptions; +pub use sync::options::{CheckpointMode, RequestsCheckpointMode, SyncOptions}; pub use sync::status::SyncStatusData; pub use sync::stream_priority::StreamPriority; pub mod error; diff --git a/powersync/src/sync/connector.rs b/powersync/src/sync/connector.rs index b5dc376..3ce5e4a 100644 --- a/powersync/src/sync/connector.rs +++ b/powersync/src/sync/connector.rs @@ -1,3 +1,5 @@ +use std::pin::Pin; + use async_trait::async_trait; use url::Url; @@ -12,6 +14,26 @@ pub trait BackendConnector: Send + Sync { /// Inspects completed CRUD transactions on a database and uploads them. async fn upload_data(&self) -> Result<(), PowerSyncError>; + + /// This is optional, and should only return a future for connectors capable of requesting + /// checkpoints. + /// + /// For upploads that are processed asynchronously by a backend (for example through a message + /// queue): The sync client as part of the PowerSync Rust SDK generates a checkpoint request id + /// and hands it to your backend via this function, which is responsible for creaeting a + /// matching checkpoint once the uploads preceeding the request have been processed. + /// + /// For more details, see [asynchronous backend uploads](https://docs.powersync.com/client-sdks/advanced/checkpoint-requests#asynchronous-upload-backends). + /// + /// To use this connector, using [crate::sync::options::CheckpointMode::Requests] is required. + /// Note that this requires PowerSync service version 1.24.0 or later. + fn post_checkpoint_request<'a>( + &'a self, + _client_id: &'a str, + _request_id: i64, + ) -> Option> + Send + 'a>>> { + None + } } /// Credentials used to connect to a PowerSync service instance. diff --git a/powersync/src/sync/download/http.rs b/powersync/src/sync/download/http.rs index 819d6c2..9886740 100644 --- a/powersync/src/sync/download/http.rs +++ b/powersync/src/sync/download/http.rs @@ -1,6 +1,7 @@ use std::borrow::Cow; use std::sync::Arc; +use crate::BackendConnector; use crate::http::{Request, Response}; use crate::sync::instruction::CheckpointRequestPayload; use crate::util::LineSplitter; @@ -102,9 +103,16 @@ pub async fn write_checkpoint( /// Posts a checkpoint request to the PowerSync sync service. pub async fn checkpoint_request( db: &InnerPowerSyncState, + connector: &dyn BackendConnector, body: &CheckpointRequestPayload, auth: PowerSyncCredentials, ) -> Result { + if let Some(future) = + connector.post_checkpoint_request(&body.client_id, body.checkpoint_request_id) + { + return future.await; + } + let url = auth.parsed_endpoint("sync/checkpoint-request")?; let body = serde_json::to_vec(body)?; diff --git a/powersync/src/sync/download/sync_iteration.rs b/powersync/src/sync/download/sync_iteration.rs index 9b2a6f6..0be46f2 100644 --- a/powersync/src/sync/download/sync_iteration.rs +++ b/powersync/src/sync/download/sync_iteration.rs @@ -183,7 +183,7 @@ impl<'a> DownloadClient<'a> { request: CheckpointRequestPayload, ) -> Result<(), PowerSyncError> { let credentials = connector.fetch_credentials().await?; - let response = checkpoint_request(&db, &request, credentials).await?; + let response = checkpoint_request(&db, connector.as_ref(), &request, credentials).await?; db.seed_checkpoint_request_id(response).await?; Ok(()) @@ -211,6 +211,7 @@ impl<'a> DownloadClient<'a> { future.await?; } + *future = None; // This doesn't generate events, we should just poll the future. future::pending().await } diff --git a/powersync/src/sync/options.rs b/powersync/src/sync/options.rs index 2d61696..c43e8d1 100644 --- a/powersync/src/sync/options.rs +++ b/powersync/src/sync/options.rs @@ -60,6 +60,11 @@ impl SyncOptions { } } } + + /// Configures the [CheckpointMode] used to request checkpoints after completed uploads. + pub fn with_checkpoint_mode(&mut self, mode: CheckpointMode) { + self.checkpoints = mode; + } } #[derive(Clone, Debug, Default, PartialEq, Eq)] diff --git a/powersync/src/sync/upload.rs b/powersync/src/sync/upload.rs index d1209e4..ea9ab8d 100644 --- a/powersync/src/sync/upload.rs +++ b/powersync/src/sync/upload.rs @@ -159,6 +159,7 @@ impl<'a> CrudUpload<'a> { let checkpoint_request_id = self.db.next_checkpoint_request_id().await?; checkpoint_request( &self.db, + self.options.connector.as_ref(), &CheckpointRequestPayload { client_id, checkpoint_request_id, diff --git a/powersync/tests/sync_test.rs b/powersync/tests/sync_test.rs index 4ee0fb9..79fa3fa 100644 --- a/powersync/tests/sync_test.rs +++ b/powersync/tests/sync_test.rs @@ -8,10 +8,14 @@ use std::{ use async_trait::async_trait; use event_listener::Event; -use futures_lite::{StreamExt, future}; +use futures_lite::{ + FutureExt, StreamExt, + future::{self, yield_now}, +}; use powersync::{ - BackendConnector, PowerSyncCredentials, PowerSyncDatabase, StreamPriority, StreamSubscription, - StreamSubscriptionOptions, SyncOptions, SyncStatusData, error::PowerSyncError, + BackendConnector, CheckpointMode, PowerSyncCredentials, PowerSyncDatabase, + RequestsCheckpointMode, StreamPriority, StreamSubscription, StreamSubscriptionOptions, + SyncOptions, SyncStatusData, error::PowerSyncError, }; use powersync_test_utils::{ DatabaseTest, @@ -39,8 +43,15 @@ impl SyncStreamTest { self.connect_options(|_| {}); } + fn connect_with_checkpoints(&self) { + self.connect_options(|options| { + options + .with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); + }); + } + fn connect_options(&self, configure: impl FnOnce(&mut SyncOptions)) { - let mut options = SyncOptions::new(TestConnector); + let mut options = SyncOptions::new(TestConnector::default()); configure(&mut options); self.run(self.db.connect(options)) @@ -514,3 +525,248 @@ fn reconnects_on_failure() { sync.test.advance_time(Duration::from_mins(30)); assert!(task.is_finished()); } + +#[test] +fn requests_checkpoints_for_updates() { + struct TestConnector { + 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> { + let Some(tx) = self.db.next_crud_transaction().await? else { + return Ok(()); + }; + + tx.complete().await?; + Ok(()) + } + } + + let sync = SyncStreamTest::new(); + let mut options = SyncOptions::new(TestConnector { + db: sync.db.clone(), + }); + options.with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); + options.with_retry_delay(Duration::ZERO); + sync.run(sync.db.connect(options)); + + { + let writer = sync.run(sync.db.writer()).unwrap(); + writer + .execute( + "INSERT INTO users (id, name) VALUES (uuid(), ?)", + params!["local user"], + ) + .unwrap(); + } + + // The local write should eventually be uploaded. + sync.run(async { + while sync + .test + .http + .last_checkpoint_request + .load(Ordering::SeqCst) + < 2 + { + yield_now().await; + } + }); +} + +#[test] +fn reports_download_error_when_seeding_checkpoint_fails() { + let sync = SyncStreamTest::new(); + sync.test + .http + .checkpoint_requests_supported + .store(false, Ordering::SeqCst); + + sync.connect_options(|options| { + options.with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); + options.with_retry_delay(Duration::ZERO); + }); + sync.run(sync.wait_for_status(|s| s.download_error().is_some())); +} + +/* +test('reposts current checkpoint until applied', () async { + await waitForConnection( + options: SyncOptions( + retryDelay: retryDelay, + checkpointMode: RequestsCheckpointMode.unverifiedDuration( + Duration(milliseconds: 50), + ), + ), + ); + + // Because we didn't include the checkpoint in a sync response, it + // should keep getting requested. + await syncService.waitForCheckpointRequest( + () => syncService.amountOfCheckpointRequests == 10, + ); + + // Finally, include the checkpoint. + syncService.addLine(checkpoint(lastOpId: 0, writeCheckpoint: '1')); + syncService.addLine(checkpointComplete(lastOpId: '0')); + await database.waitForFirstSync(); + + final requestsBefore = syncService.amountOfCheckpointRequests; + await Future.delayed(Duration(seconds: 1)); + expect( + syncService.amountOfCheckpointRequests, + requestsBefore, + reason: 'Should not keep posting checkpoint requests', + ); +}); + +test('download is retried on checkpoint request', () async { + final status = await waitForConnection( + options: SyncOptions( + retryDelay: Duration(seconds: 10), + checkpointMode: CheckpointMode.requests(), + ), + expectNoWarnings: false, + ); + uploadData = (db) async { + if (await db.getCrudBatch() case final batch?) { + await batch.complete(); + } + }; + + // Destroy the initial connection by sending a bogus line. + final sw = Stopwatch()..start(); + syncService.addLine( + checkpoint(lastOpId: 1, writeCheckpoint: 'invalid line'), + ); + await expectLater( + status, + emitsThrough(isSyncStatus(downloadError: isNotNull)), + ); + syncService.endCurrentListener(); + + // Trigger an upload here. Because the upload needs a seeded sync + // iteration, we should reconnect immediately instead of after the + // configured 10s delay. + await database.execute( + 'INSERT INTO customers (id, email, name) VALUES (uuid(), uuid(), uuid())', + ); + await expectLater(status, emitsThrough(isSyncStatus(connected: true))); + final elapsed = sw.elapsed; + + expect(elapsed, lessThan(Duration(seconds: 5))); +}); +*/ + +#[test] +fn can_use_checkpoint_method_from_connector() { + 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()); + }), + }); + options.with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); + options.with_retry_delay(Duration::ZERO); + sync.run(sync.db.connect(options)); + + sync.run(listener); +} + +#[test] +fn reconciles_checkpoint_state_on_token_expiry() { + let sync = SyncStreamTest::new(); + sync.test + .http + .last_checkpoint_request + .store(100, Ordering::SeqCst); + + sync.connect_options(|options| { + options.with_retry_delay(Duration::ZERO); + options.with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); + }); + sync.run(async { + while sync + .test + .http + .last_checkpoint_request + .load(Ordering::SeqCst) + == 0 + { + yield_now().await; + } + + let request = sync.test.http.receive_requests.recv().await.unwrap(); + + // Simulate what would happen if we suddenly switched users after the old token expired. + // The client expects a checkpoint of 100, for another user the service wouldn't have that + // counter yet. The client must request a checkpoint with the existing id, allowing the + // service to recognize that this device + user combo needs higher checkpoint ids. + sync.test + .http + .last_checkpoint_request + .store(0, Ordering::SeqCst); + + request.send_keepalive(0).await; + request.channel.close(); + + while sync + .test + .http + .amount_of_checkpoint_requests + .load(Ordering::SeqCst) + < 2 + { + yield_now().await; + } + assert_eq!( + sync.test + .http + .last_checkpoint_request + .load(Ordering::SeqCst), + 100 + ); + }); +} + +#[test] +fn reads_sync_lines_before_checkpoint_requests_are_ready() { + let sync = SyncStreamTest::new(); + { + let mut guard = sync.test.http.before_checkpoint_response.lock().unwrap(); + // Make /sync/checkpoint-request never return + *guard = Box::new(|| future::pending().boxed()); + } + + sync.connect_with_checkpoints(); + + sync.run(async { + let request = sync.test.http.receive_requests.recv().await.unwrap(); + sync.wait_for_status(|s| s.is_connected()).await; + + request + .send_checkpoint(Checkpoint { + last_op_id: 0, + write_checkpoint: None, + buckets: vec![], + streams: vec![], + }) + .await; + sync.wait_for_status(|s| s.is_downloading()).await; + }); +} diff --git a/powersync_test_utils/src/mock_sync_service.rs b/powersync_test_utils/src/mock_sync_service.rs index 7af82cf..30d6ceb 100644 --- a/powersync_test_utils/src/mock_sync_service.rs +++ b/powersync_test_utils/src/mock_sync_service.rs @@ -1,13 +1,16 @@ use crate::sync_line::{Checkpoint, DataLine, OplogEntry, SyncLine}; use async_trait::async_trait; use bytes::Bytes; -use futures_lite::{Stream, StreamExt, ready, stream}; +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 serde::Serialize; +use serde::{Deserialize, Serialize}; use serde_json::json; +use serde_with::{DisplayFromStr, serde_as}; use std::pin::Pin; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::task::Context; use std::{ sync::{Arc, Mutex}, @@ -18,6 +21,10 @@ pub struct MockSyncService { pub receive_requests: async_channel::Receiver, send_requests: async_channel::Sender, pub write_checkpoints: Mutex WriteCheckpointResponse + Send>>, + pub before_checkpoint_response: Mutex Boxed<()> + Send>>, + pub checkpoint_requests_supported: AtomicBool, + pub last_checkpoint_request: AtomicUsize, + pub amount_of_checkpoint_requests: AtomicUsize, } impl Default for MockSyncService { @@ -30,6 +37,10 @@ impl Default for MockSyncService { write_checkpoints: Mutex::new(Box::new(|| { WriteCheckpointResponse::new("10".to_string()) })), + before_checkpoint_response: Mutex::new(Box::new(|| async {}.boxed())), + checkpoint_requests_supported: AtomicBool::new(true), + last_checkpoint_request: Default::default(), + amount_of_checkpoint_requests: Default::default(), } } } @@ -53,13 +64,22 @@ impl MockSyncService { #[async_trait] impl HttpClient for MockClient { async fn send(&self, req: Request) -> Result { - match req.url.path() { - "/sync/stream" => Ok(self.service.sync_stream(req).await), - "/write-checkpoint2.json" => { - Ok(self.service.generate_write_checkpoint_response()) + Ok(match req.url.path() { + "/sync/stream" => self.service.sync_stream(req).await, + "/sync/checkpoint-request" => { + if !self + .service + .checkpoint_requests_supported + .load(Ordering::SeqCst) + { + return Ok(MockSyncService::generate_not_found()); + } + + self.service.generate_checkpoint_request_response(req) } - _ => Ok(MockSyncService::generate_bad_request()), - } + "/write-checkpoint2.json" => self.service.generate_write_checkpoint_response(), + _ => MockSyncService::generate_bad_request(), + }) } } @@ -105,6 +125,46 @@ impl MockSyncService { } } + fn generate_checkpoint_request_response(&self, req: Request) -> Response { + #[serde_as] + #[derive(Deserialize, Serialize)] + struct RequestData { + #[serde_as(as = "DisplayFromStr")] + checkpoint_request_id: usize, + } + + #[derive(Serialize)] + struct ResponseBodyData { + data: RequestData, + } + + let body: RequestData = serde_json::from_slice(&req.body.unwrap_or_default()).unwrap(); + let before = self + .last_checkpoint_request + .fetch_max(body.checkpoint_request_id, Ordering::SeqCst); + let request_id = before.max(body.checkpoint_request_id); + + self.amount_of_checkpoint_requests + .fetch_add(1, Ordering::SeqCst); + + let data = Bytes::from( + serde_json::to_vec(&ResponseBodyData { + data: RequestData { + checkpoint_request_id: request_id, + }, + }) + .unwrap(), + ); + Response { + status: 200, + content_type: Some("application/json".to_string()), + body: ResponseBody { + length: Some(data.len() as u64), + reader: stream::once(Ok(data)).boxed(), + }, + } + } + fn generate_bad_request() -> Response { Response { status: 400, @@ -115,6 +175,17 @@ impl MockSyncService { }, } } + + fn generate_not_found() -> Response { + Response { + status: 404, + content_type: None, + body: ResponseBody { + reader: stream::empty().boxed(), + length: Some(0), + }, + } + } } pub struct PendingSyncResponse { @@ -144,6 +215,13 @@ impl PendingSyncResponse { self.channel.send(msg).await.unwrap() } + pub async fn send_keepalive(&self, token_expires_in: u32) { + self.channel + .send(SyncLine::TokenExpiresIn(token_expires_in)) + .await + .unwrap(); + } + pub async fn bogus_data_line(&self, last_id: &mut i64, bucket: &'static str, amount: usize) { let mut oplog = vec![]; for _ in 0..amount { @@ -212,7 +290,21 @@ pub struct WriteCheckpointResponseData { pub write_checkpoint: String, } -pub struct TestConnector; +pub struct TestConnector { + pub post_checkpoint_request: Box< + dyn Fn(i64) -> Option> + Send>>> + + Send + + Sync, + >, +} + +impl Default for TestConnector { + fn default() -> Self { + Self { + post_checkpoint_request: Box::new(|_| None), + } + } +} #[async_trait] impl BackendConnector for TestConnector { @@ -226,4 +318,12 @@ impl BackendConnector for TestConnector { async fn upload_data(&self) -> Result<(), PowerSyncError> { Ok(()) } + + fn post_checkpoint_request<'a>( + &'a self, + _client_id: &'a str, + request_id: i64, + ) -> Option> + Send + 'a>>> { + (self.post_checkpoint_request)(request_id) + } } diff --git a/powersync_test_utils/src/sync_line.rs b/powersync_test_utils/src/sync_line.rs index 77144e8..37b792e 100644 --- a/powersync_test_utils/src/sync_line.rs +++ b/powersync_test_utils/src/sync_line.rs @@ -4,6 +4,7 @@ use serde_with::{DisplayFromStr, serde_as}; pub enum SyncLine<'a> { Checkpoint(Checkpoint<'a>), + TokenExpiresIn(u32), Data(DataLine<'a>), Custom(serde_json::Value), } @@ -19,6 +20,11 @@ impl<'a> Serialize for SyncLine<'a> { map.serialize_entry("checkpoint", cp)?; map.end() } + SyncLine::TokenExpiresIn(remaining) => { + let mut map = serializer.serialize_map(Some(1))?; + map.serialize_entry("token_expires_in", remaining)?; + map.end() + } SyncLine::Data(data) => { let mut map = serializer.serialize_map(Some(1))?; map.serialize_entry("data", data)?; From 148c6eb15ec5f4aa1d68b58ecf760c3fdf6b998e Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Fri, 18 Sep 2026 17:45:07 +0200 Subject: [PATCH 3/4] Support checkpoint request protocol --- powersync/src/db/internal.rs | 27 ++- powersync/src/error.rs | 2 +- powersync/src/sync/checkpoint.rs | 106 ++++++++++- powersync/src/sync/coordinator.rs | 36 ++-- powersync/src/sync/download/http.rs | 2 +- powersync/src/sync/download/sync_iteration.rs | 3 +- powersync/src/sync/options.rs | 2 +- powersync/src/sync/status.rs | 7 + powersync/src/sync/upload.rs | 34 ++-- powersync/tests/sync_test.rs | 176 +++++++++++------- 10 files changed, 290 insertions(+), 105 deletions(-) diff --git a/powersync/src/db/internal.rs b/powersync/src/db/internal.rs index e38de66..180a734 100644 --- a/powersync/src/db/internal.rs +++ b/powersync/src/db/internal.rs @@ -120,14 +120,17 @@ impl InnerPowerSyncState { } pub async fn next_checkpoint_request_id(&self) -> Result { - let mut writer = self.writer().await?; - let writer = TransactionGuard::new(writer.sqlite_connection_mut())?; - - let id = Self::checkpoint_request_control(&writer, CheckpointCounter::Next, None)?; - writer.commit()?; + let id = self + .read_checkpoint_request_id(CheckpointCounter::Next) + .await?; Ok(id.expect("Core extension should return next checkpoint request id")) } + pub async fn current_checkpoint_request_id(&self) -> Result, PowerSyncError> { + self.read_checkpoint_request_id(CheckpointCounter::Current) + .await + } + pub fn target_checkpoint_request_id( writer: &TransactionGuard, update: Option, @@ -135,6 +138,18 @@ impl InnerPowerSyncState { Self::checkpoint_request_control(writer, CheckpointCounter::Target, update) } + async fn read_checkpoint_request_id( + &self, + counter: CheckpointCounter, + ) -> Result, PowerSyncError> { + let mut writer = self.writer().await?; + let writer = TransactionGuard::new(writer.sqlite_connection_mut())?; + + let id = Self::checkpoint_request_control(&writer, counter, None)?; + writer.commit()?; + Ok(id) + } + fn checkpoint_request_control( writer: &TransactionGuard, counter: CheckpointCounter, @@ -236,6 +251,7 @@ impl InnerPowerSyncState { pub enum CheckpointCounter { Target, Seed, + Current, Next, } @@ -244,6 +260,7 @@ impl CheckpointCounter { match self { CheckpointCounter::Target => "target_checkpoint_request_id", CheckpointCounter::Seed => "seed_checkpoint_request_id", + CheckpointCounter::Current => "current_checkpoint_request_id", CheckpointCounter::Next => "next_checkpoint_request_id", } } diff --git a/powersync/src/error.rs b/powersync/src/error.rs index 20c34ea..c5d24e3 100644 --- a/powersync/src/error.rs +++ b/powersync/src/error.rs @@ -15,7 +15,7 @@ pub type Result = std::result::Result; /// [RawPowerSyncError] enum type). #[derive(Debug, Clone)] pub struct PowerSyncError { - inner: Arc, + pub(crate) inner: Arc, } impl PowerSyncError { diff --git a/powersync/src/sync/checkpoint.rs b/powersync/src/sync/checkpoint.rs index a7a10d1..90a7837 100644 --- a/powersync/src/sync/checkpoint.rs +++ b/powersync/src/sync/checkpoint.rs @@ -1,6 +1,17 @@ +use std::sync::Arc; + +use log::{debug, warn}; use thiserror::Error; -use crate::error::PowerSyncError; +use crate::{ + BackendConnector, CheckpointMode, RequestsCheckpointMode, SyncOptions, + db::internal::InnerPowerSyncState, + error::{PowerSyncError, RawPowerSyncError}, + sync::{ + coordinator::SyncChannels, download::http::checkpoint_request, + instruction::CheckpointRequestPayload, upload::get_client_id, + }, +}; #[derive(Error, Debug)] pub enum CheckpointError { @@ -15,3 +26,96 @@ pub enum CheckpointError { #[error("Error on sync status before checkpoint was applied: {cause}")] StatusError { cause: PowerSyncError }, } + +pub async fn repost_unacknowledged_checkpoints( + db: Arc, + channels: SyncChannels, + options: SyncOptions, +) { + let CheckpointMode::Requests(requests) = options.checkpoints 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; + + if let Err(err) = result { + if let RawPowerSyncError::Checkpoint { + error: CheckpointError::InstanceNotSupported, + } = err.inner.as_ref() + { + return; + } + + warn!("Error retrying checkpoint request: {err}"); + db.env.runtime.delay_once(requests.retry_delay).await; + } + } +} + +async fn repost_unacknowledged_checkpoint_iteration( + db: &InnerPowerSyncState, + channels: &SyncChannels, + connector: &dyn BackendConnector, + mode: RequestsCheckpointMode, +) -> Result<(), PowerSyncError> { + // Make sure the system is seeded and ready + channels + .checkpoints + .wait_for_checkpoint_requests_ready(false) + .await?; + + // Get the current checkpoint_request_id + let Some(request_id) = db + .current_checkpoint_request_id() + .await? + .take_if(|id| *id > 0) + else { + // This should not be reached. For completeness sake - wait a bit. + db.env.runtime.delay_once(mode.retry_delay).await; + return Ok(()); + }; + + // Give the request some time to sync + db.env.runtime.delay_once(mode.retry_delay).await; + + if db.current_checkpoint_request_id().await? != Some(request_id) { + return Ok(()); + } + + // If the request was applied, we don't need to retry + if db + .status + .current_snapshot() + .is_checkpoint_request_applied(request_id) + { + return Ok(()); + } + + // Make sure we are online and ready before making the request + channels + .checkpoints + .wait_for_checkpoint_requests_ready(false) + .await?; + + // It's safe if this request races with a new one. The service will reject it. + debug!("Retrying checkpoint request id {request_id}"); + checkpoint_request( + &db, + connector, + &CheckpointRequestPayload { + client_id: get_client_id(db).await?, + checkpoint_request_id: request_id, + }, + ) + .await?; + + Ok(()) +} diff --git a/powersync/src/sync/coordinator.rs b/powersync/src/sync/coordinator.rs index 5a8de28..c262098 100644 --- a/powersync/src/sync/coordinator.rs +++ b/powersync/src/sync/coordinator.rs @@ -9,6 +9,7 @@ use crate::{ env::PowerSyncTask, error::PowerSyncError, sync::{ + checkpoint::repost_unacknowledged_checkpoints, download::{DownloadEvent, download_loop}, state::CheckpointStateSignals, streams::ChangedSyncSubscriptions, @@ -40,15 +41,21 @@ impl SyncCoordinator { )); let uploads = db.env.spawn(crud_upload_loop( db.clone(), - options, + options.clone(), channels.clone(), uploads_receive, )); + let checkpoints = db.env.spawn(repost_unacknowledged_checkpoints( + db.clone(), + channels.clone(), + options, + )); *guard = Some(SyncTasks { channels, uploads: Some(uploads), downloads: Some(downloads), + retried_checkpoints: Some(checkpoints), }); } @@ -108,26 +115,33 @@ struct SyncTasks { channels: SyncChannels, downloads: Option, uploads: Option, + retried_checkpoints: Option, } impl SyncTasks { + fn tasks(&mut self) -> [&mut Option; 3] { + [ + &mut self.downloads, + &mut self.uploads, + &mut self.retried_checkpoints, + ] + } + pub async fn cancel(mut self) { - if let Some(task) = self.downloads.take() { - task.cancel_and_join().await; - } - if let Some(task) = self.uploads.take() { - task.cancel_and_join().await; + for maybe_task in self.tasks() { + if let Some(task) = maybe_task.take() { + task.cancel_and_join().await; + } } } } impl Drop for SyncTasks { fn drop(&mut self) { - if let Some(task) = self.downloads.take() { - task.cancel(); - } - if let Some(task) = self.uploads.take() { - task.cancel(); + for maybe_task in self.tasks() { + if let Some(task) = maybe_task.take() { + task.cancel(); + } } } } diff --git a/powersync/src/sync/download/http.rs b/powersync/src/sync/download/http.rs index 9886740..bc5a9bf 100644 --- a/powersync/src/sync/download/http.rs +++ b/powersync/src/sync/download/http.rs @@ -105,7 +105,6 @@ pub async fn checkpoint_request( db: &InnerPowerSyncState, connector: &dyn BackendConnector, body: &CheckpointRequestPayload, - auth: PowerSyncCredentials, ) -> Result { if let Some(future) = connector.post_checkpoint_request(&body.client_id, body.checkpoint_request_id) @@ -113,6 +112,7 @@ pub async fn checkpoint_request( return future.await; } + let auth = connector.fetch_credentials().await?; let url = auth.parsed_endpoint("sync/checkpoint-request")?; let body = serde_json::to_vec(body)?; diff --git a/powersync/src/sync/download/sync_iteration.rs b/powersync/src/sync/download/sync_iteration.rs index 0be46f2..ffe1ba3 100644 --- a/powersync/src/sync/download/sync_iteration.rs +++ b/powersync/src/sync/download/sync_iteration.rs @@ -182,8 +182,7 @@ impl<'a> DownloadClient<'a> { connector: Arc, request: CheckpointRequestPayload, ) -> Result<(), PowerSyncError> { - let credentials = connector.fetch_credentials().await?; - let response = checkpoint_request(&db, connector.as_ref(), &request, credentials).await?; + let response = checkpoint_request(&db, connector.as_ref(), &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 c43e8d1..a881626 100644 --- a/powersync/src/sync/options.rs +++ b/powersync/src/sync/options.rs @@ -74,7 +74,7 @@ pub enum CheckpointMode { Requests(RequestsCheckpointMode), } -#[derive(Clone, Debug, PartialEq, Eq)] +#[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct RequestsCheckpointMode { pub(crate) retry_delay: Duration, } diff --git a/powersync/src/sync/status.rs b/powersync/src/sync/status.rs index 20eecb0..141a369 100644 --- a/powersync/src/sync/status.rs +++ b/powersync/src/sync/status.rs @@ -154,6 +154,13 @@ impl SyncStatusData { } } + pub(crate) fn is_checkpoint_request_applied(&self, id: i64) -> bool { + let Some(last_applied) = self.downloading.last_applied_checkpoint_request else { + return false; + }; + return last_applied >= id; + } + /// Returns an [EventListener] that completes once this data is stale, or returns [None] /// immediately if this is already stale. pub(crate) fn listen_for_changes(&self) -> Option { diff --git a/powersync/src/sync/upload.rs b/powersync/src/sync/upload.rs index ea9ab8d..69b17d7 100644 --- a/powersync/src/sync/upload.rs +++ b/powersync/src/sync/upload.rs @@ -132,23 +132,13 @@ impl<'a> CrudUpload<'a> { } async fn get_write_checkpoint(&self) -> Result { - let client_id = { - let reader = self.db.reader().await?; - - let stmt = reader - .sqlite_connection() - .prepare("SELECT powersync_client_id()")?; - let ResultCode::ROW = stmt.step()? else { - panic!("Expected row"); // Can't happen, scalar select - }; - - stmt.column_text(0)?.to_string() - }; - - let credentials = self.options.connector.fetch_credentials().await?; + let client_id = get_client_id(&self.db).await?; match self.options.checkpoints { - CheckpointMode::Legacy => write_checkpoint(&self.db, &client_id, credentials).await, + CheckpointMode::Legacy => { + let credentials = self.options.connector.fetch_credentials().await?; + write_checkpoint(&self.db, &client_id, credentials).await + } CheckpointMode::Requests(_) => { self.channels .checkpoints @@ -164,7 +154,6 @@ impl<'a> CrudUpload<'a> { client_id, checkpoint_request_id, }, - credentials, ) .await } @@ -256,3 +245,16 @@ impl PendingCheckpointRequest { Ok(()) } } + +pub async fn get_client_id(db: &InnerPowerSyncState) -> Result { + let reader = db.reader().await?; + + let stmt = reader + .sqlite_connection() + .prepare("SELECT powersync_client_id()")?; + let ResultCode::ROW = stmt.step()? else { + panic!("Expected row"); // Can't happen, scalar select + }; + + Ok(stmt.column_text(0)?.to_string()) +} diff --git a/powersync/tests/sync_test.rs b/powersync/tests/sync_test.rs index 79fa3fa..98fd1d2 100644 --- a/powersync/tests/sync_test.rs +++ b/powersync/tests/sync_test.rs @@ -598,74 +598,116 @@ fn reports_download_error_when_seeding_checkpoint_fails() { sync.run(sync.wait_for_status(|s| s.download_error().is_some())); } -/* -test('reposts current checkpoint until applied', () async { - await waitForConnection( - options: SyncOptions( - retryDelay: retryDelay, - checkpointMode: RequestsCheckpointMode.unverifiedDuration( - Duration(milliseconds: 50), - ), - ), - ); - - // Because we didn't include the checkpoint in a sync response, it - // should keep getting requested. - await syncService.waitForCheckpointRequest( - () => syncService.amountOfCheckpointRequests == 10, - ); - - // Finally, include the checkpoint. - syncService.addLine(checkpoint(lastOpId: 0, writeCheckpoint: '1')); - syncService.addLine(checkpointComplete(lastOpId: '0')); - await database.waitForFirstSync(); - - final requestsBefore = syncService.amountOfCheckpointRequests; - await Future.delayed(Duration(seconds: 1)); - expect( - syncService.amountOfCheckpointRequests, - requestsBefore, - reason: 'Should not keep posting checkpoint requests', - ); -}); - -test('download is retried on checkpoint request', () async { - final status = await waitForConnection( - options: SyncOptions( - retryDelay: Duration(seconds: 10), - checkpointMode: CheckpointMode.requests(), - ), - expectNoWarnings: false, - ); - uploadData = (db) async { - if (await db.getCrudBatch() case final batch?) { - await batch.complete(); +#[test] +fn reposts_current_checkpoint_until_applied() { + let sync = SyncStreamTest::new(); + + sync.connect_options(|options| { + options.with_checkpoint_mode(CheckpointMode::Requests( + Duration::from_hours(1).try_into().unwrap(), + )); + }); + + sync.run(async { + let request = sync.test.http.receive_requests.recv().await.unwrap(); + + // Because we didn't include the checkpoint in a sync response, it should keep getting + // requested. + for i in 2..=10 { + sync.test.advance_time(Duration::from_hours(1)); + assert_eq!( + sync.test + .http + .amount_of_checkpoint_requests + .load(Ordering::SeqCst), + i + ); + } + + // Finally, include the checkpoint + request + .send_checkpoint(Checkpoint { + last_op_id: 0, + write_checkpoint: Some(1), + buckets: vec![], + streams: vec![], + }) + .await; + request.send_checkpoint_complete(0, None).await; + sync.wait_for_status(|s| !s.is_downloading()).await; + + // After which no further checkpoints should be requested + sync.test.advance_time(Duration::from_hours(10)); + assert_eq!( + sync.test + .http + .amount_of_checkpoint_requests + .load(Ordering::SeqCst), + 10 + ); + }); +} + +#[test] +fn download_is_retried_on_checkpoint_request() { + struct Connector { + db: PowerSyncDatabase, } - }; - - // Destroy the initial connection by sending a bogus line. - final sw = Stopwatch()..start(); - syncService.addLine( - checkpoint(lastOpId: 1, writeCheckpoint: 'invalid line'), - ); - await expectLater( - status, - emitsThrough(isSyncStatus(downloadError: isNotNull)), - ); - syncService.endCurrentListener(); - - // Trigger an upload here. Because the upload needs a seeded sync - // iteration, we should reconnect immediately instead of after the - // configured 10s delay. - await database.execute( - 'INSERT INTO customers (id, email, name) VALUES (uuid(), uuid(), uuid())', - ); - await expectLater(status, emitsThrough(isSyncStatus(connected: true))); - final elapsed = sw.elapsed; - - expect(elapsed, lessThan(Duration(seconds: 5))); -}); -*/ + + #[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> { + let tx = self.db.next_crud_transaction().await?; + if let Some(tx) = tx { + tx.complete().await?; + } + + Ok(()) + } + } + + let sync = SyncStreamTest::new(); + let mut options = SyncOptions::new(Connector { + db: sync.db.clone(), + }); + options.with_retry_delay(Duration::from_hours(1)); + options.with_checkpoint_mode(CheckpointMode::Requests(RequestsCheckpointMode::default())); + + // Destroy the initial connection by sending a bogus line. + sync.run(async { + sync.db.connect(options).await; + + let request = sync.test.http.receive_requests.recv().await.unwrap(); + request + .channel + .send(SyncLine::Custom(json!("invalid sync line"))) + .await + .unwrap(); + + sync.wait_for_status(|s| s.download_error().is_some()).await; + + // Trigger an upload here. Because the upload needs a seeded sync iteration, we should + // reconnect immediately instead of after the configured delay. + { + let writer = sync.db.writer().await.unwrap(); + writer + .execute( + "INSERT INTO users (id, name) VALUES (uuid(), 'local user')", + params![], + ) + .unwrap(); + } + + sync.test.http.receive_requests.recv().await.unwrap(); + }); +} #[test] fn can_use_checkpoint_method_from_connector() { From 9cbbf458c1b447f679444920d657d99ca9bee89d Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Fri, 18 Sep 2026 18:13:17 +0200 Subject: [PATCH 4/4] AI feedback --- powersync/src/sync/download/http.rs | 8 ++++++++ powersync/src/sync/state.rs | 2 +- 2 files changed, 9 insertions(+), 1 deletion(-) diff --git a/powersync/src/sync/download/http.rs b/powersync/src/sync/download/http.rs index bc5a9bf..5c35b6c 100644 --- a/powersync/src/sync/download/http.rs +++ b/powersync/src/sync/download/http.rs @@ -3,6 +3,7 @@ use std::sync::Arc; use crate::BackendConnector; use crate::http::{Request, Response}; +use crate::sync::checkpoint::CheckpointError; use crate::sync::instruction::CheckpointRequestPayload; use crate::util::LineSplitter; use crate::{ @@ -131,6 +132,13 @@ pub async fn checkpoint_request( }; let response = db.env.client.send(request).await?; + if response.status == 404 { + return Err(RawPowerSyncError::Checkpoint { + error: CheckpointError::InstanceNotSupported, + } + .into()); + } + check_ok(response.status)?; #[derive(Deserialize)] diff --git a/powersync/src/sync/state.rs b/powersync/src/sync/state.rs index de87f6f..76aec42 100644 --- a/powersync/src/sync/state.rs +++ b/powersync/src/sync/state.rs @@ -73,11 +73,11 @@ impl CheckpointStateSignals { &self, wake_download_loop: bool, ) -> Result<(), CheckpointError> { + let mut receiver = self.state_channel.1.activate_cloned(); let mut state = { let waiters = self.state_and_waiter.lock().unwrap(); waiters.0.clone() }; - let mut receiver = self.state_channel.1.activate_cloned(); loop { match state {