From f1c4d914fc196b0dec84b9b512ef52e40b4c861e Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 25 Aug 2026 14:43:10 +0200 Subject: [PATCH 1/6] Support explicit checkpoint requests --- packages/powersync/CHANGELOG.md | 2 + packages/powersync/lib/powersync.dart | 1 + .../native/native_powersync_database.dart | 32 ++++++- .../native/sync_isolate_protocol.dart | 9 ++ .../lib/src/database/powersync_database.dart | 18 +++- .../database/web/web_powersync_database.dart | 3 +- .../powersync/lib/src/devtools/protocol.dart | 1 + .../lib/src/platform_specific/web.dart | 31 ++++++- .../lib/src/sync/checkpoint_request.dart | 87 +++++++++++++++++- .../lib/src/sync/connection_manager.dart | 76 +++++++++++---- .../lib/src/sync/mutable_sync_status.dart | 1 + .../lib/src/sync/streaming_sync.dart | 19 +++- .../powersync/lib/src/sync/sync_status.dart | 22 ++++- .../lib/src/web/sync_controller.dart | 8 ++ .../lib/src/web/sync_worker_protocol.dart | 23 +++++ .../test/sync/in_memory_sync_test.dart | 92 +++++++++++++++++++ .../test/utils/abstract_test_utils.dart | 3 +- 17 files changed, 399 insertions(+), 29 deletions(-) diff --git a/packages/powersync/CHANGELOG.md b/packages/powersync/CHANGELOG.md index 3d3ed2eb..4854dd19 100644 --- a/packages/powersync/CHANGELOG.md +++ b/packages/powersync/CHANGELOG.md @@ -2,6 +2,8 @@ - Add `SyncOptions.checkpointMode`, which can be used to replace a legacy implementation for write checkpoints with a newer and more efficient endpoint. +- Add `PowerSyncDatabase.requestCheckpoint`, which can be used to "sync now": The checkpoint completes once all + source data created before the checkpoint has synced. - Reduce log noise from attachments service ([#408](https://github.com/powersync-ja/powersync.dart/issues/408)). - Fix errors that can't be sent across isolates disrupting the sync isolate. - Update PowerSync SQLite core extension to 0.5.3, improving error messages for failed internal SQL statements. diff --git a/packages/powersync/lib/powersync.dart b/packages/powersync/lib/powersync.dart index 545b06d3..b32fca6b 100644 --- a/packages/powersync/lib/powersync.dart +++ b/packages/powersync/lib/powersync.dart @@ -11,6 +11,7 @@ export 'src/database/encryption_options.dart'; export 'src/exceptions.dart'; export 'src/log.dart'; export 'src/schema.dart'; +export 'src/sync/checkpoint_request.dart' show CheckpointRequestException; export 'src/sync/options.dart' hide ResolvedSyncOptions, LegacyCheckpointMode, RequestsCheckpointMode; export 'src/sync/stream.dart' hide CoreActiveStreamSubscription; diff --git a/packages/powersync/lib/src/database/native/native_powersync_database.dart b/packages/powersync/lib/src/database/native/native_powersync_database.dart index 7bd57c1b..1c6450b4 100644 --- a/packages/powersync/lib/src/database/native/native_powersync_database.dart +++ b/packages/powersync/lib/src/database/native/native_powersync_database.dart @@ -5,6 +5,7 @@ import 'package:meta/meta.dart'; import 'package:logging/logging.dart'; import 'package:powersync/src/abort_controller.dart'; +import 'package:powersync/src/platform_specific/int64.dart'; import 'package:powersync/src/sync/bucket_storage.dart'; import 'package:powersync/src/connector.dart'; import 'package:powersync/src/database/powersync_database.dart'; @@ -16,6 +17,7 @@ import 'package:powersync/src/sync/streaming_sync.dart'; import 'package:powersync/src/sync/sync_status.dart'; import 'package:sqlite_async/sqlite_async.dart'; +import '../../sync/checkpoint_request.dart'; import 'sync_isolate_protocol.dart'; /// A PowerSync managed database. @@ -38,7 +40,7 @@ final class NativePowerSyncDatabase extends BasePowerSyncDatabase { @override @internal - Future connectInternal({ + Future connectInternal({ required PowerSyncBackendConnector connector, required ResolvedSyncOptions options, required List initiallyActiveStreams, @@ -50,6 +52,7 @@ final class NativePowerSyncDatabase extends BasePowerSyncDatabase { StreamSubscription? activeStreamsSubscription; final receiveMessages = ReceivePort(); final receiveUnhandledErrors = ReceivePort(); + final results = IsolateResultCollection(); final receiveExit = ReceivePort(); final mutexServer = MutexServer({ 'sync': group.syncMutex, @@ -68,6 +71,7 @@ final class NativePowerSyncDatabase extends BasePowerSyncDatabase { } // Cleanup + results.close(); activeStreamsSubscription?.cancel(); receiveMessages.close(); receiveUnhandledErrors.close(); @@ -186,7 +190,8 @@ final class NativePowerSyncDatabase extends BasePowerSyncDatabase { // Don't spawn isolate if this operation was cancelled already. if (abort.aborted) { - return waitForShutdown(); + await waitForShutdown(); + throw AbortException('connectInternal'); } receiveExit.listen((message) { @@ -218,6 +223,8 @@ final class NativePowerSyncDatabase extends BasePowerSyncDatabase { receivedIsolateExit.future, ]).whenComplete(close), ); + + return _IsolateSyncHandle(initPort!, results); } } @@ -277,6 +284,13 @@ Future _syncIsolate(_PowerSyncDatabaseIsolateArgs args) async { ); case ClientToSyncIsolateMessageType.mutexGranted: mutexes.markGranted(payload as int); + case ClientToSyncIsolateMessageType.requestCheckpoint: + (payload as PortCompleter).handle(() async { + final sync = openedStreamingSync; + if (sync == null) throw disconnected; + + return await sync.requestCheckpoint(); + }); } }); sPort.sendInit(rPort.sendPort); @@ -359,3 +373,17 @@ Future _syncIsolate(_PowerSyncDatabaseIsolateArgs args) async { }, ); } + +final class _IsolateSyncHandle implements RemoteStreamingSyncHandle { + final SyncIsolatePort _syncIsolate; + final IsolateResultCollection _results; + + _IsolateSyncHandle(this._syncIsolate, this._results); + + @override + Future requestCheckpoint() { + final pending = _results.createPending(); + _syncIsolate.requestCheckpoint(pending.completer); + return pending.future; + } +} diff --git a/packages/powersync/lib/src/database/native/sync_isolate_protocol.dart b/packages/powersync/lib/src/database/native/sync_isolate_protocol.dart index 41d19e65..d864da94 100644 --- a/packages/powersync/lib/src/database/native/sync_isolate_protocol.dart +++ b/packages/powersync/lib/src/database/native/sync_isolate_protocol.dart @@ -6,6 +6,7 @@ import 'package:sqlite_async/sqlite_async.dart'; import '../../connector.dart'; import '../../isolate_completer.dart'; +import '../../platform_specific/int64.dart'; import '../../sync/streaming_sync.dart' show SubscribedStream; import '../../sync/sync_status.dart'; @@ -58,6 +59,10 @@ enum ClientToSyncIsolateMessageType { /// A completed [SyncIsolateToClientMessageType.mutexAcquire] call, payload is /// the request id being completed. mutexGranted, + + /// Client requests a checkpoint from the sync client, payload is a + /// [PortCompleter]. + requestCheckpoint, } /// Typed client-side view over a [SendPort] used to send messages to a sync @@ -78,6 +83,10 @@ extension type SyncIsolatePort(SendPort port) { void sendMutexGranted(int requestId) { send(ClientToSyncIsolateMessageType.mutexGranted, requestId); } + + void requestCheckpoint(PortCompleter completer) { + send(ClientToSyncIsolateMessageType.requestCheckpoint, completer); + } } /// Sync isolate view over a [SendPort] used to send messages to a client diff --git a/packages/powersync/lib/src/database/powersync_database.dart b/packages/powersync/lib/src/database/powersync_database.dart index c9cb51a2..f7407e28 100644 --- a/packages/powersync/lib/src/database/powersync_database.dart +++ b/packages/powersync/lib/src/database/powersync_database.dart @@ -21,6 +21,7 @@ import '../powersync_update_notification.dart'; import '../schema.dart'; import '../schema_logic.dart' as schema_logic; import '../sync/bucket_storage.dart'; +import '../sync/checkpoint_request.dart'; import '../sync/connection_manager.dart'; import '../sync/options.dart'; import '../sync/stream.dart'; @@ -337,7 +338,7 @@ abstract base class PowerSyncDatabase extends SqliteConnection { /// active, so the method should not call [disconnect] itself. @protected @internal - Future connectInternal({ + Future connectInternal({ required PowerSyncBackendConnector connector, required ResolvedSyncOptions options, required List initiallyActiveStreams, @@ -630,6 +631,21 @@ SELECT * FROM crud_entries; SyncStream syncStream(String name, [Map? parameters]) { return _connections.syncStream(name, parameters); } + + /// Requests a checkpoint from the PowerSync service. + /// + /// The returned request can be awaited (using [CheckpointRequest.waitForSync]) + /// to confirm that the local database has applieed server-side changes up to + /// the checkpoint. This method requires an active or connecting sync client + /// connected with [CheckpointMode.requests] and PowerSync service version + /// 1.24.0 or later. + /// + /// It throws a [CheckpointRequestException] if requesting the checkpoint has + /// failed, e.g. because the database is not connected. + @experimental + Future requestCheckpoint() { + return _connections.requestCheckpoint(); + } } @internal diff --git a/packages/powersync/lib/src/database/web/web_powersync_database.dart b/packages/powersync/lib/src/database/web/web_powersync_database.dart index 975b6388..c039e3af 100644 --- a/packages/powersync/lib/src/database/web/web_powersync_database.dart +++ b/packages/powersync/lib/src/database/web/web_powersync_database.dart @@ -32,7 +32,7 @@ final class WebPowerSyncDatabase extends BasePowerSyncDatabase { @override @internal - Future connectInternal({ + Future connectInternal({ required PowerSyncBackendConnector connector, required AbortController abort, required List initiallyActiveStreams, @@ -98,6 +98,7 @@ final class WebPowerSyncDatabase extends BasePowerSyncDatabase { await sync.abort(); abort.completeAbort(); }).ignore(); + return sync; } /// Takes a read lock, without starting a transaction. diff --git a/packages/powersync/lib/src/devtools/protocol.dart b/packages/powersync/lib/src/devtools/protocol.dart index be27866a..61ba7ec3 100644 --- a/packages/powersync/lib/src/devtools/protocol.dart +++ b/packages/powersync/lib/src/devtools/protocol.dart @@ -94,5 +94,6 @@ SyncStatus deserializeSyncStatus(Map serialized) { .toList(), _ => null, }, + lastAppliedCheckpoint: null, ); } diff --git a/packages/powersync/lib/src/platform_specific/web.dart b/packages/powersync/lib/src/platform_specific/web.dart index fb5d3fc1..767bed9d 100644 --- a/packages/powersync/lib/src/platform_specific/web.dart +++ b/packages/powersync/lib/src/platform_specific/web.dart @@ -44,11 +44,28 @@ BasePowerSyncDatabase openPowerSyncDatabase( ); } +const _has64BitInts = !identical(0, 0.0); + +final _intConversionBuffer = JSArrayBuffer(8); +final _intConversionView = JSDataView(_intConversionBuffer); + Int64 parseInt64(String s) { if (Int64.has64BitIntegers) { return NativeInt64.parse(s); } else { - return JsBigInt64(_stringToBigInt(s.toJS)); + return JsBigInt64.parse(s); + } +} + +Int64 int64FromBigInt(JSBigInt value) { + if (_has64BitInts) { + _intConversionView.setBigInt64(0, value); + final high = _intConversionView.getUint32(0).toDartInt; + final low = _intConversionView.getUint32(4).toDartInt; + + return NativeInt64((high << 32) | low); + } else { + return JsBigInt64(value); } } @@ -57,6 +74,10 @@ final class JsBigInt64 implements Int64 { JsBigInt64(this.value) : assert(!Int64.has64BitIntegers); + static JsBigInt64 parse(String s) { + return JsBigInt64(_stringToBigInt(s.toJS)); + } + @override bool operator >=(Int64 other) { return other is JsBigInt64 && @@ -85,3 +106,11 @@ external JSNumber _bigIntToDouble(JSBigInt value); @JS('BigInt') external JSBigInt _stringToBigInt(JSString value); + +extension on JSDataView { + @JS() + external JSNumber getUint32(int byteOffset); + + @JS() + external void setBigInt64(int byteoffset, JSBigInt value); +} diff --git a/packages/powersync/lib/src/sync/checkpoint_request.dart b/packages/powersync/lib/src/sync/checkpoint_request.dart index 80c5abab..655f49d3 100644 --- a/packages/powersync/lib/src/sync/checkpoint_request.dart +++ b/packages/powersync/lib/src/sync/checkpoint_request.dart @@ -1,13 +1,96 @@ +library; + import 'package:meta/meta.dart'; +import '../database/powersync_database.dart'; +import '../platform_specific/int64.dart'; +import 'connection_manager.dart'; +import 'sync_status.dart'; + +/// A checkpoint request created by [PowerSyncDatabase.requestCheckpoint]. +/// +/// Use this to wait until the local database has applied server-side changes up +/// to the requested checkpoint. This is useful for explicit refresh flows where +/// the caller wants confirmation that the local view has caught up to the +/// service. +/// +/// Checkpoint requests are backed by request ids tracked in the local database, +/// so they are reusable across disconnect and reconnect cycles. A [waitForSync] +/// interrupted by a disconnect throws an error, but the same request can be +/// awaited again once a new connection is established. +/// +/// Requests do not survive [PowerSyncDatabase.disconnectAndClear], instances +/// created before a clear should be discarded and requested again. +abstract final class CheckpointRequest { + CheckpointRequest._(); + + /// Whether this checkpoint request has synced before. + @experimental + bool get hasSynced; + + /// Waits until this checkpoint has been synced locally. + /// + /// This method faisl on sync errors: If a download or upload error occurs + /// before this checkpoint request has synced, that error is rethrown here. + /// This makes it easier to observe sync errors when relying on checkpoints. + /// Once sync has recovered, it is valid to call this method again to await + /// the checkpoint. + @experimental + Future waitForSync({Future? abortTrigger}); +} + +@internal +final class CheckpointRequestImpl extends CheckpointRequest { + final Int64 _requestId; + final ConnectionManager _connections; + + CheckpointRequestImpl(this._requestId, this._connections) : super._(); + + @override + bool get hasSynced { + return _connections.currentStatus.hasApplied(_requestId); + } + + @override + Future waitForSync({Future? abortTrigger}) async { + if (hasSynced) return; + + _connections.checkConnectedWithRequestsMode(); + + await _connections.firstStatusMatching(abort: abortTrigger, (status) { + if (status.hasApplied(_requestId)) return true; + + if (status.anyError case final anyError?) { + throw CheckpointRequestException._( + 'Sync error while waiting for checkpoint request', + anyError, + ); + } + + if (!status.connected && !status.connecting) { + throw disconnected; + } + + return false; + }); + } +} + /// An exception related to requested checkpoints for PowerSync. final class CheckpointRequestException implements Exception { final String _message; + final Object? cause; - const CheckpointRequestException._(this._message); + const CheckpointRequestException._(this._message, [this.cause]); @override - String toString() => _message; + String toString() { + if (cause case final cause?) { + return '$_message ($cause)'; + } + + return _message; + } } @internal diff --git a/packages/powersync/lib/src/sync/connection_manager.dart b/packages/powersync/lib/src/sync/connection_manager.dart index c51fade0..fbf8f588 100644 --- a/packages/powersync/lib/src/sync/connection_manager.dart +++ b/packages/powersync/lib/src/sync/connection_manager.dart @@ -8,8 +8,10 @@ import 'package:powersync/src/database/active_instances.dart'; import 'package:powersync/src/sync/options.dart'; import 'package:powersync/src/sync/stream.dart'; import 'package:powersync/src/sync/sync_status.dart'; +import 'package:sqlite_async/sqlite_async.dart'; import '../database/powersync_database.dart'; +import 'checkpoint_request.dart'; import 'instruction.dart'; import 'mutable_sync_status.dart'; import 'streaming_sync.dart'; @@ -23,6 +25,8 @@ final class ConnectionManager { final PowerSyncDatabase db; final ActiveDatabaseGroup _activeGroup; + ResolvedSyncOptions _currentOptions = ResolvedSyncOptions.resolve(null); + /// All streams (with parameters) for which a subscription has been requested /// explicitly. final Map<_RawStreamKey, _ActiveSubscription> _locallyActiveSubscriptions = @@ -46,7 +50,7 @@ final class ConnectionManager { /// /// The controller must only be accessed from within a critical section of the /// sync mutex. - AbortController? _abortActiveSync; + (AbortController, RemoteStreamingSyncHandle)? _abortActiveSync; ConnectionManager(this.db) : _activeGroup = db.group; @@ -56,8 +60,19 @@ final class ConnectionManager { } } + void checkConnectedWithRequestsMode() { + final client = _abortActiveSync?.$2; + if (client == null) { + throw disconnected; + } + + if (_currentOptions.checkpointMode is! RequestsCheckpointMode) { + throw disabled; + } + } + Future _abortCurrentSync() async { - if (_abortActiveSync case final disconnector?) { + if (_abortActiveSync case (final disconnector, _)?) { /// Checking `disconnecter.aborted` prevents race conditions /// where multiple calls to `disconnect` can attempt to abort /// the controller more than once before it has finished aborting. @@ -83,15 +98,36 @@ final class ConnectionManager { }); } - Future firstStatusMatching(bool Function(SyncStatus) predicate) async { - if (predicate(currentStatus)) { - return; - } - await for (final result in statusStream) { - if (predicate(result)) { - break; + Future firstStatusMatching( + bool Function(SyncStatus) predicate, { + Future? abort, + }) { + final completer = Completer(); + final subscription = statusStream.listen(null); + + void checkPredicate(SyncStatus status) { + try { + if (predicate(status)) { + completer.complete(); + subscription.cancel(); + } + } catch (e, s) { + completer.completeError(e, s); + subscription.cancel(); } } + + subscription.onData(checkPredicate); + checkPredicate(currentStatus); + + abort?.whenComplete(() { + if (!completer.isCompleted) { + completer.completeError(AbortException('firstStatusMatching')); + subscription.cancel(); + } + }); + + return completer.future; } List get _subscribedStreams => [ @@ -116,9 +152,6 @@ final class ConnectionManager { late void Function() retryHandler; Future connectWithSyncLock() async { - // Ensure there has not been a subsequent connect() call installing a new - // sync client. - assert(identical(_abortActiveSync, thisConnectAborter)); assert(!thisConnectAborter.aborted); // This needs to be a single-subscription controller to ensure we won't @@ -127,8 +160,9 @@ final class ConnectionManager { // need a new controller per internal connect attempt. final subscriptionsChanged = _subscriptionsChanged = StreamController(); + _currentOptions = options; // ignore: invalid_use_of_protected_member - await db.connectInternal( + final handle = await db.connectInternal( connector: connector, options: options, abort: thisConnectAborter, @@ -140,6 +174,7 @@ final class ConnectionManager { // while we hold the lock (and async tasks won't hold the sync lock). asyncWorkZone: zone, ); + _abortActiveSync = (thisConnectAborter, handle); thisConnectAborter.onCompletion.whenComplete(retryHandler); } @@ -158,7 +193,7 @@ final class ConnectionManager { assert(identical(_abortActiveSync, thisConnectAborter)); // We need a new abort controller for this attempt - _abortActiveSync = thisConnectAborter = AbortController(); + thisConnectAborter = AbortController(); db.logger.warning('Sync client failed, retrying...'); await connectWithSyncLock(); @@ -171,9 +206,6 @@ final class ConnectionManager { await _abortCurrentSync(); assert(_abortActiveSync == null); - // Install the abort controller for this particular connect call, allowing - // it to be disconnected. - _abortActiveSync = thisConnectAborter; await connectWithSyncLock(); }); } @@ -278,6 +310,16 @@ final class ConnectionManager { return _SyncStreamImplementation(this, name, parameters); } + Future requestCheckpoint() async { + return _activeGroup.syncConnectMutex.lock(() async { + final client = _abortActiveSync?.$2; + checkConnectedWithRequestsMode(); + + final checkpoint = await client!.requestCheckpoint(); + return CheckpointRequestImpl(checkpoint, this); + }); + } + void close() { _statusController.close(); } diff --git a/packages/powersync/lib/src/sync/mutable_sync_status.dart b/packages/powersync/lib/src/sync/mutable_sync_status.dart index b10bd70c..5a113f4c 100644 --- a/packages/powersync/lib/src/sync/mutable_sync_status.dart +++ b/packages/powersync/lib/src/sync/mutable_sync_status.dart @@ -64,6 +64,7 @@ final class MutableSyncStatus { uploadError: uploadError, downloadError: downloadError, streamSubscriptions: streams, + lastAppliedCheckpoint: lastAppliedCheckpointRequestId, ); } } diff --git a/packages/powersync/lib/src/sync/streaming_sync.dart b/packages/powersync/lib/src/sync/streaming_sync.dart index 96478cd6..55b6e17f 100644 --- a/packages/powersync/lib/src/sync/streaming_sync.dart +++ b/packages/powersync/lib/src/sync/streaming_sync.dart @@ -27,7 +27,13 @@ import 'sync_status.dart'; typedef SubscribedStream = ({String name, String parameters}); -abstract interface class StreamingSync { +/// A superinterface for [StreamingSync] providing methods that can be called +/// in remote contexts. +abstract interface class RemoteStreamingSyncHandle { + Future requestCheckpoint(); +} + +abstract interface class StreamingSync implements RemoteStreamingSyncHandle { Stream get statusStream; Future streamingSync(); @@ -157,6 +163,17 @@ class StreamingSyncImplementation implements StreamingSync { } } + @override + Future requestCheckpoint() async { + if (options.checkpointMode is! RequestsCheckpointMode) { + throw checkpoint.disabled; + } + + final neverAbort = AbortController(); + final checkpointId = await _requestNextCheckpointFromService(neverAbort); + return Int64.parse(checkpointId); + } + Future _resolveClientId() async { return (_clientId ??= await adapter.getClientId()); } diff --git a/packages/powersync/lib/src/sync/sync_status.dart b/packages/powersync/lib/src/sync/sync_status.dart index 99d7a457..521aee06 100644 --- a/packages/powersync/lib/src/sync/sync_status.dart +++ b/packages/powersync/lib/src/sync/sync_status.dart @@ -4,6 +4,7 @@ library; import 'package:collection/collection.dart'; import 'package:meta/meta.dart'; +import '../platform_specific/int64.dart'; import 'stream.dart'; final class SyncStatus { @@ -56,6 +57,8 @@ final class SyncStatus { final List? _internalSubscriptions; + final Int64? _lastAppliedCheckpoint; + @internal const SyncStatus({ required this.connected, @@ -68,8 +71,10 @@ final class SyncStatus { required this.uploadError, required this.priorityStatusEntries, required List? streamSubscriptions, + required Int64? lastAppliedCheckpoint, }) : hasSynced = lastSyncedAt != null, - _internalSubscriptions = streamSubscriptions; + _internalSubscriptions = streamSubscriptions, + _lastAppliedCheckpoint = lastAppliedCheckpoint; @internal const SyncStatus.uninitialized() @@ -83,7 +88,8 @@ final class SyncStatus { downloadError = null, priorityStatusEntries = const [], _internalSubscriptions = null, - hasSynced = null; + hasSynced = null, + _lastAppliedCheckpoint = null; @override bool operator ==(Object other) { @@ -104,7 +110,8 @@ final class SyncStatus { other._internalSubscriptions, _internalSubscriptions, ) && - other.downloadProgress == downloadProgress); + other.downloadProgress == downloadProgress && + other._lastAppliedCheckpoint == _lastAppliedCheckpoint); } // Deprecated because it can't set fields back to null @@ -131,6 +138,7 @@ final class SyncStatus { priorityStatusEntries ?? this.priorityStatusEntries, downloadProgress: downloadProgress, streamSubscriptions: _internalSubscriptions, + lastAppliedCheckpoint: _lastAppliedCheckpoint, ); } @@ -267,6 +275,8 @@ extension InternalSyncStatusAccess on SyncStatus { List? get internalSubscriptions => _internalSubscriptions; + Int64? get lastAppliedCheckpoint => _lastAppliedCheckpoint; + SyncStatus changeErrors({ required Object? downloadError, required Object? uploadError, @@ -282,8 +292,14 @@ extension InternalSyncStatusAccess on SyncStatus { uploadError: uploadError, priorityStatusEntries: priorityStatusEntries, streamSubscriptions: _internalSubscriptions, + lastAppliedCheckpoint: _lastAppliedCheckpoint, ); } + + bool hasApplied(Int64 checkpointRequest) { + final lastApplied = _lastAppliedCheckpoint; + return lastApplied != null && lastApplied >= checkpointRequest; + } } /// Current information about a [SyncStream] that the sync client is subscribed diff --git a/packages/powersync/lib/src/web/sync_controller.dart b/packages/powersync/lib/src/web/sync_controller.dart index 5bfa75a3..74bcff37 100644 --- a/packages/powersync/lib/src/web/sync_controller.dart +++ b/packages/powersync/lib/src/web/sync_controller.dart @@ -8,6 +8,8 @@ import 'package:web/web.dart' hide Client; import '../connector.dart'; import '../database/powersync_database.dart'; +import '../platform_specific/int64.dart'; +import '../platform_specific/web.dart'; import '../sync/options.dart'; import '../sync/streaming_sync.dart'; import '../sync/sync_status.dart'; @@ -146,6 +148,12 @@ class SyncWorkerHandle implements StreamingSync { await _channel.abortSynchronization(); } + @override + Future requestCheckpoint() async { + final checkpoint = await _channel.requestCheckpoint(); + return int64FromBigInt(checkpoint); + } + @override Stream get statusStream => _status.stream; diff --git a/packages/powersync/lib/src/web/sync_worker_protocol.dart b/packages/powersync/lib/src/web/sync_worker_protocol.dart index 8c51b17b..955c0ccf 100644 --- a/packages/powersync/lib/src/web/sync_worker_protocol.dart +++ b/packages/powersync/lib/src/web/sync_worker_protocol.dart @@ -56,6 +56,10 @@ enum SyncWorkerMessageType { /// For requests, the payload is a [CustomCheckpointRequest]. The response is /// a nullable string. customCheckpointRequest, + + /// Request a checkpoint from the sync client. + requestCheckpoint, + invalidCredentialsCallback, credentialsCallback, @@ -294,6 +298,7 @@ extension type SerializedSyncStatus._(JSObject _) implements JSObject { required JSArray? priorityStatusEntries, required JSArray? syncProgress, required JSString streamSubscriptions, + required JSBigInt? lastAppliedCheckpoint, }); factory SerializedSyncStatus.from(SyncStatus status) { @@ -321,6 +326,11 @@ extension type SerializedSyncStatus._(JSObject _) implements JSObject { ), }, streamSubscriptions: json.encode(status.internalSubscriptions).toJS, + lastAppliedCheckpoint: switch (status.lastAppliedCheckpoint) { + null => null, + JsBigInt64(:final value) => value, + final other => JsBigInt64.parse(other.toString()).value, + }, ); } @@ -336,6 +346,7 @@ extension type SerializedSyncStatus._(JSObject _) implements JSObject { external JSArray? priorityStatusEntries; external JSArray? syncProgress; external JSString? streamSubscriptions; + external JSBigInt? lastAppliedCheckpoint; SyncStatus asSyncStatus() { final streamSubscriptions = this.streamSubscriptions?.toDart; @@ -383,6 +394,10 @@ extension type SerializedSyncStatus._(JSObject _) implements JSObject { ) .toList(), }, + lastAppliedCheckpoint: switch (lastAppliedCheckpoint) { + null => null, + final id => JsBigInt64(id), + }, ); } } @@ -468,6 +483,7 @@ final class WorkerCommunicationChannel { case SyncWorkerMessageType.credentialsCallback: case SyncWorkerMessageType.invalidCredentialsCallback: case SyncWorkerMessageType.uploadCrud: + case SyncWorkerMessageType.requestCheckpoint: requestId = (message.payload as JSNumber).toDartInt; case SyncWorkerMessageType.customCheckpointRequest: requestId = @@ -719,6 +735,13 @@ final class WorkerCommunicationChannel { return (await completion) as JSArrayBuffer?; } + Future requestCheckpoint() async { + final response = await _numericRequest( + SyncWorkerMessageType.requestCheckpoint, + ); + return response as JSBigInt; + } + Future close() async { if (!_closed.isCompleted) { _incomingMessages?.cancel(); diff --git a/packages/powersync/test/sync/in_memory_sync_test.dart b/packages/powersync/test/sync/in_memory_sync_test.dart index f2208500..46a0c5ed 100644 --- a/packages/powersync/test/sync/in_memory_sync_test.dart +++ b/packages/powersync/test/sync/in_memory_sync_test.dart @@ -4,9 +4,11 @@ import 'dart:convert'; import 'package:async/async.dart'; import 'package:logging/logging.dart'; import 'package:powersync/powersync.dart'; +import 'package:powersync/src/sync/checkpoint_request.dart'; import 'package:powersync/src/sync/options.dart'; import 'package:shelf/shelf.dart'; import 'package:shelf_router/shelf_router.dart'; +import 'package:sqlite_async/sqlite_async.dart'; import 'package:test/test.dart'; import '../server/sync_server/in_memory_sync_server.dart'; @@ -1216,6 +1218,96 @@ void _declareTests(String name, bool bson) { completeInitialRequest.complete(); }); + + group('requestCheckpoint', () { + test('fails when disconnected', () async { + expect(database.requestCheckpoint(), throwsA(disconnected)); + }); + + test('fails when connected with legacy mode', () async { + await waitForConnection(); + expect(database.requestCheckpoint(), throwsA(disabled)); + }); + + test('waits until data is applied', () async { + await waitForConnection(options: options); + + final requested = await database.requestCheckpoint(); + syncService.addLine(checkpoint(lastOpId: 0, writeCheckpoint: '2')); + expect(requested.hasSynced, false); + syncService.addLine(checkpointComplete(lastOpId: '0')); + + await requested.waitForSync(); + expect(requested.hasSynced, isTrue); + }); + + test('throws on disconnect but can request again', () async { + await waitForConnection(options: options); + final request = await database.requestCheckpoint(); + + final failureExpectation = expectLater( + request.waitForSync(), + throwsA(disconnected), + ); + await database.disconnect(); + syncService.endCurrentListener(); + await failureExpectation; + + await waitForConnection(options: options); + syncService.addLine(checkpoint(lastOpId: 0, writeCheckpoint: '2')); + syncService.addLine(checkpointComplete(lastOpId: '0')); + + await request.waitForSync(); + }); + + test('fails when reconnecting with legacy mode', () async { + await waitForConnection(options: options); + final request = await database.requestCheckpoint(); + + final failureExpectation = expectLater( + request.waitForSync(), + throwsA(disconnected), + ); + await database.disconnect(); + syncService.endCurrentListener(); + await failureExpectation; + + await waitForConnection(); + await expectLater(request.waitForSync(), throwsA(disabled)); + }); + + test('fails on sync errors', () async { + await waitForConnection(options: options, expectNoWarnings: false); + final checkpoint = await database.requestCheckpoint(); + + final failureExpectation = expectLater( + checkpoint.waitForSync(), + throwsA( + isA().having( + (e) => e.toString(), + 'toString()', + contains('Sync error while waiting for checkpoint request'), + ), + ), + ); + syncService.addLine({'checkpoint': 'invalid sync line'}); + await failureExpectation; + }); + + test('can abort waiting for requests', () async { + await waitForConnection(options: options, expectNoWarnings: false); + final checkpoint = await database.requestCheckpoint(); + + final abort = Completer(); + final wasAborted = expectLater( + checkpoint.waitForSync(abortTrigger: abort.future), + throwsA(isA()), + ); + + abort.complete(); + await wasAborted; + }); + }); }); test('checkpoint request can be aborted', () async { diff --git a/packages/powersync/test/utils/abstract_test_utils.dart b/packages/powersync/test/utils/abstract_test_utils.dart index 99fade0d..131aa0a3 100644 --- a/packages/powersync/test/utils/abstract_test_utils.dart +++ b/packages/powersync/test/utils/abstract_test_utils.dart @@ -184,7 +184,7 @@ final class TestDatabase extends BasePowerSyncDatabase { }); @override - Future connectInternal({ + Future connectInternal({ required PowerSyncBackendConnector connector, required ResolvedSyncOptions options, required List initiallyActiveStreams, @@ -217,6 +217,7 @@ final class TestDatabase extends BasePowerSyncDatabase { await impl.abort(); abort.completeAbort(); }).ignore(); + return impl; } @override From f2587f05315755f08e0f87499b62de30d6570971 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 25 Aug 2026 14:59:12 +0200 Subject: [PATCH 2/6] Add tests --- packages/powersync/lib/src/web/sync_worker.dart | 6 ++++++ .../powersync/test/sync/streaming_sync_test.dart | 16 ++++++++++++++++ packages/powersync/test/test_server.dart | 11 +++++++++++ .../powersync/test/web/sync_worker_test.dart | 5 ++++- 4 files changed, 37 insertions(+), 1 deletion(-) diff --git a/packages/powersync/lib/src/web/sync_worker.dart b/packages/powersync/lib/src/web/sync_worker.dart index 564dccfa..0df7d207 100644 --- a/packages/powersync/lib/src/web/sync_worker.dart +++ b/packages/powersync/lib/src/web/sync_worker.dart @@ -21,6 +21,7 @@ import 'package:sqlite_async/web.dart'; import 'package:web/web.dart' hide RequestMode; import '../database/powersync_database.dart'; +import '../platform_specific/web.dart'; import '../sync/bucket_storage.dart'; import 'http/client.dart'; import 'sync_worker_protocol.dart'; @@ -117,6 +118,11 @@ class ConnectedClient { (payload as UpdateSubscriptions).toDart, ); return (JSObject(), null); + case SyncWorkerMessageType.requestCheckpoint: + final sync = _runner!.sync!; + final checkpoint = await sync.requestCheckpoint(); + + return ((checkpoint as JsBigInt64).value, null); default: throw StateError('Unexpected message type $type'); } diff --git a/packages/powersync/test/sync/streaming_sync_test.dart b/packages/powersync/test/sync/streaming_sync_test.dart index d56df1f5..766bbf8e 100644 --- a/packages/powersync/test/sync/streaming_sync_test.dart +++ b/packages/powersync/test/sync/streaming_sync_test.dart @@ -301,6 +301,22 @@ void main() { ); await hadRequest.future; }); + + test('can request checkpoints', () async { + final pdb = await testUtils.setupPowerSync(path: path); + final server = await createServer(); + final connector = TestConnector( + () async => + PowerSyncCredentials(endpoint: server.endpoint, token: 'token'), + ); + + await pdb.connect( + connector: connector, + options: SyncOptions(checkpointMode: CheckpointMode.requests()), + ); + + await pdb.requestCheckpoint(); + }); }); } diff --git a/packages/powersync/test/test_server.dart b/packages/powersync/test/test_server.dart index 3b8eb0d2..287e0fe7 100644 --- a/packages/powersync/test/test_server.dart +++ b/packages/powersync/test/test_server.dart @@ -19,7 +19,18 @@ final class TestServer { TestServer({this.tokenExpiresIn = 65}); Future init({MockSyncService? mockSyncService}) async { + // TODO: Unify this with MockSyncService app.post('/sync/stream', handleSyncStream); + app.post('/sync/checkpoint-request', (Request request) async { + final body = convert.jsonDecode(await request.readAsString()); + final checkpoint = int.parse(body['checkpoint_request_id'] as String); + + return Response.ok( + convert.jsonEncode({ + 'data': {'checkpoint_request_id': '$checkpoint'}, + }), + ); + }); // Open on an arbitrary open port server = await shelf_io.serve( mockSyncService?.router.call ?? app.call, diff --git a/packages/powersync/test/web/sync_worker_test.dart b/packages/powersync/test/web/sync_worker_test.dart index f181994c..b9c94ca5 100644 --- a/packages/powersync/test/web/sync_worker_test.dart +++ b/packages/powersync/test/web/sync_worker_test.dart @@ -156,9 +156,12 @@ void main() { ); await handle.streamingSync(); - service.waitForCheckpointRequest( + await service.waitForCheckpointRequest( () => service.amountOfCheckpointRequests > 0, ); + + final checkpoint = await handle.requestCheckpoint(); + expect(checkpoint.toString(), '2'); await handle.abort(); }); From d3b8578376ea6afd9c5558c584aacbda6d17b7c5 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 25 Aug 2026 16:05:23 +0200 Subject: [PATCH 3/6] Add checkpoint requests to demo app --- .../lib/app_config_template.dart | 4 ++ .../lib/components/app_bar.dart | 48 ++++++++++++- .../lib/components/page_layout.dart | 70 ++++++++++++++++++- .../lib/navigation.gr.dart | 2 +- .../lib/powersync/powersync.dart | 12 +++- .../lib/stores/refresh.dart | 30 ++++++++ .../macos/Runner.xcodeproj/project.pbxproj | 6 +- .../xcshareddata/swiftpm/Package.resolved | 13 ---- .../xcshareddata/swiftpm/Package.resolved | 13 ---- 9 files changed, 163 insertions(+), 35 deletions(-) create mode 100644 demos/supabase-todolist-drift/lib/stores/refresh.dart delete mode 100644 demos/supabase-todolist-drift/macos/Runner.xcodeproj/project.xcworkspace/xcshareddata/swiftpm/Package.resolved delete mode 100644 demos/supabase-todolist-drift/macos/Runner.xcworkspace/xcshareddata/swiftpm/Package.resolved diff --git a/demos/supabase-todolist-drift/lib/app_config_template.dart b/demos/supabase-todolist-drift/lib/app_config_template.dart index 4f35559e..a1384542 100644 --- a/demos/supabase-todolist-drift/lib/app_config_template.dart +++ b/demos/supabase-todolist-drift/lib/app_config_template.dart @@ -6,4 +6,8 @@ class AppConfig { static const String powersyncUrl = 'https://foo.powersync.journeyapps.com'; static const String supabaseStorageBucket = ''; // Optional. Only required when syncing attachments and using Supabase Storage. See packages/powersync_attachments_helper. + + // Whether to enable swipe-to-refresh through through checkpoint requests + // (https://docs.powersync.com/client-sdks/advanced/checkpoint-requests). + static const bool enableCheckpointRequests = false; } diff --git a/demos/supabase-todolist-drift/lib/components/app_bar.dart b/demos/supabase-todolist-drift/lib/components/app_bar.dart index 1c454323..8df86b7c 100644 --- a/demos/supabase-todolist-drift/lib/components/app_bar.dart +++ b/demos/supabase-todolist-drift/lib/components/app_bar.dart @@ -1,11 +1,14 @@ import 'package:auto_route/auto_route.dart'; import 'package:flutter/foundation.dart'; +import 'package:flutter_hooks/flutter_hooks.dart'; +import 'package:hooks_riverpod/hooks_riverpod.dart'; import 'package:material_ui/material_ui.dart'; -import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'package:powersync/powersync.dart'; +import '../app_config.dart'; import '../powersync/powersync.dart'; import '../screens/search.dart'; +import '../stores/refresh.dart'; final appBar = AppBar( title: const Text('PowerSync Flutter Demo'), @@ -28,6 +31,7 @@ final class StatusAppBar extends ConsumerWidget implements PreferredSizeWidget { leading: const AutoLeadingButton(), title: title, actions: [ + const _RefreshButton(), IconButton( onPressed: () { showSearch(context: context, delegate: FtsSearchDelegate()); @@ -142,3 +146,45 @@ class _AutoLeadingButtonState extends State { setState(() {}); } } + +/// A refresh button implemented with [Sync Catch-Up](https://docs.powersync.com/client-sdks/advanced/checkpoint-requests). +class _RefreshButton extends HookConsumerWidget { + const _RefreshButton(); + + @override + Widget build(BuildContext context, WidgetRef ref) { + if (!AppConfig.enableCheckpointRequests) { + return const SizedBox.shrink(); + } + + final canRefreshNow = ref.watch(canRefresh); + final isRefreshing = ref.watch(refresh).isPending; + final controller = useAnimationController( + duration: const Duration(seconds: 1), + ); + + useEffect(() { + if (isRefreshing) { + controller.repeat(); + } else { + controller.reset(); + } + return null; + }, [isRefreshing]); + + void triggerRefresh() { + syncNow(ref); + } + + return Tooltip( + message: 'Sync now', + child: IconButton( + onPressed: canRefreshNow ? triggerRefresh : null, + icon: RotationTransition( + turns: controller, + child: const Icon(Icons.refresh), + ), + ), + ); + } +} diff --git a/demos/supabase-todolist-drift/lib/components/page_layout.dart b/demos/supabase-todolist-drift/lib/components/page_layout.dart index fe9017a4..1ba0bf14 100644 --- a/demos/supabase-todolist-drift/lib/components/page_layout.dart +++ b/demos/supabase-todolist-drift/lib/components/page_layout.dart @@ -1,12 +1,19 @@ +import 'dart:async'; + import 'package:auto_route/auto_route.dart'; +import 'package:flutter/scheduler.dart'; +import 'package:flutter_hooks/flutter_hooks.dart'; import 'package:material_ui/material_ui.dart'; import 'package:hooks_riverpod/hooks_riverpod.dart'; +import 'package:riverpod/experimental/mutation.dart'; +import '../app_config.dart'; import '../navigation.gr.dart'; +import '../stores/refresh.dart'; import '../supabase.dart'; import 'app_bar.dart'; -final class PageLayout extends ConsumerWidget { +final class PageLayout extends HookConsumerWidget { final Widget content; final Widget? title; final Widget? floatingActionButton; @@ -22,11 +29,70 @@ final class PageLayout extends ConsumerWidget { @override Widget build(BuildContext context, WidgetRef ref) { + Widget body = Center(child: content); + + if (AppConfig.enableCheckpointRequests) { + final refreshIndicatorKey = + useMemoized(() => GlobalKey()); + final isRefreshing = ref.watch(refresh).isPending; + final refreshingBefore = usePrevious(isRefreshing) ?? false; + final scaffoldMessenger = ScaffoldMessenger.of(context); + + useEffect(() { + // Show the refresh indicator when a refresh is triggered. + if (!refreshingBefore && isRefreshing) { + refreshIndicatorKey.currentState?.show(); + } + + // Show a status snackbar on completed sync. + if (refreshingBefore && !isRefreshing) { + SchedulerBinding.instance.addPostFrameCallback((_) { + switch (ref.read(refresh)) { + case MutationError(): + scaffoldMessenger + .showSnackBar(SnackBar(content: Text('Sync failed.'))); + case MutationSuccess(): + scaffoldMessenger.showSnackBar( + SnackBar(content: Text('Sync refresh complete!'))); + default: + break; + } + }); + } + return null; + }, [refreshIndicatorKey, refreshingBefore, isRefreshing]); + + body = RefreshIndicator( + key: refreshIndicatorKey, + onRefresh: () { + // This is called both when a child scrollable is swiped to refresh, + // and by the useEffect hook above. Depending on what triggered it, we + // either request a checkpoint or wait for the previous one to + // complete. + final isPending = ref.read(refresh).isPending; + if (isPending) { + final completer = Completer(); + late final ProviderSubscription subscription; + subscription = ref.listenManual(refresh, (previous, next) { + if (!next.isPending) { + subscription.close(); + completer.complete(); + } + }); + return completer.future; + } else { + return syncNow(ref); + } + }, + child: body, + ); + } + return Scaffold( appBar: StatusAppBar( title: title ?? const Text('PowerSync Demo'), ), - body: Center(child: content), + body: body, floatingActionButton: floatingActionButton, drawer: showDrawer ? Drawer( diff --git a/demos/supabase-todolist-drift/lib/navigation.gr.dart b/demos/supabase-todolist-drift/lib/navigation.gr.dart index 63d0f812..971eeb47 100644 --- a/demos/supabase-todolist-drift/lib/navigation.gr.dart +++ b/demos/supabase-todolist-drift/lib/navigation.gr.dart @@ -11,7 +11,7 @@ // ignore_for_file: no_leading_underscores_for_library_prefixes import 'package:auto_route/auto_route.dart' as _i10; import 'package:camera/camera.dart' as _i12; -import 'package:flutter/material.dart' as _i11; +import 'package:material_ui/material_ui.dart' as _i11; import 'package:supabase_todolist_drift/navigation.dart' as _i5; import 'package:supabase_todolist_drift/screens/add_item_dialog.dart' as _i1; import 'package:supabase_todolist_drift/screens/add_list_dialog.dart' as _i2; diff --git a/demos/supabase-todolist-drift/lib/powersync/powersync.dart b/demos/supabase-todolist-drift/lib/powersync/powersync.dart index ff2f97a3..6efc4c7e 100644 --- a/demos/supabase-todolist-drift/lib/powersync/powersync.dart +++ b/demos/supabase-todolist-drift/lib/powersync/powersync.dart @@ -7,6 +7,7 @@ import 'package:riverpod_annotation/riverpod_annotation.dart'; import 'package:stream_transform/stream_transform.dart'; import 'package:supabase_flutter/supabase_flutter.dart'; +import '../app_config.dart'; import '../supabase.dart'; import 'connector.dart'; import 'schema.dart'; @@ -22,10 +23,17 @@ Future powerSyncInstance(Ref ref) async { ); await db.initialize(); + final syncOptions = SyncOptions( + checkpointMode: AppConfig.enableCheckpointRequests + // ignore: experimental_member_use + ? CheckpointMode.requests() + : const CheckpointMode.legacy(), + ); + SupabaseConnector? currentConnector; if (ref.read(sessionProvider).value != null) { currentConnector = SupabaseConnector(); - db.connect(connector: currentConnector); + db.connect(connector: currentConnector, options: syncOptions); } final instance = Supabase.instance.client.auth; @@ -33,7 +41,7 @@ Future powerSyncInstance(Ref ref) async { final event = data.event; if (event == AuthChangeEvent.signedIn) { currentConnector = SupabaseConnector(); - db.connect(connector: currentConnector!); + db.connect(connector: currentConnector!, options: syncOptions); } else if (event == AuthChangeEvent.signedOut) { currentConnector = null; await db.disconnect(); diff --git a/demos/supabase-todolist-drift/lib/stores/refresh.dart b/demos/supabase-todolist-drift/lib/stores/refresh.dart new file mode 100644 index 00000000..35fe278c --- /dev/null +++ b/demos/supabase-todolist-drift/lib/stores/refresh.dart @@ -0,0 +1,30 @@ +// ignore_for_file: experimental_member_use + +import 'package:riverpod/experimental/mutation.dart'; +import 'package:riverpod/riverpod.dart'; + +import '../app_config.dart'; +import '../powersync/powersync.dart'; + +final refresh = Mutation(label: 'explicit sync'); + +/// Whether a checkpoint request can be sent from the current database state. +/// +/// This is the case if checkpoint requests are enabled and the database is +/// connected. For this demo, we also avoid concurrent checkpoint refresh runs. +final canRefresh = Provider((ref) { + if (!AppConfig.enableCheckpointRequests) return false; + + if (ref.watch(refresh).isPending) return false; + + final status = ref.watch(syncStatus); + return status.connected || status.connecting; +}); + +Future syncNow(MutationTarget target) { + return refresh.run(target, (tx) async { + final powersync = await tx.get(powerSyncInstanceProvider.future); + final checkpoint = await powersync.requestCheckpoint(); + await checkpoint.waitForSync(); + }); +} diff --git a/demos/supabase-todolist-drift/macos/Runner.xcodeproj/project.pbxproj b/demos/supabase-todolist-drift/macos/Runner.xcodeproj/project.pbxproj index 5c87fe22..e27fb300 100644 --- a/demos/supabase-todolist-drift/macos/Runner.xcodeproj/project.pbxproj +++ b/demos/supabase-todolist-drift/macos/Runner.xcodeproj/project.pbxproj @@ -463,7 +463,7 @@ GCC_WARN_UNINITIALIZED_AUTOS = YES_AGGRESSIVE; GCC_WARN_UNUSED_FUNCTION = YES; GCC_WARN_UNUSED_VARIABLE = YES; - MACOSX_DEPLOYMENT_TARGET = 10.15; + MACOSX_DEPLOYMENT_TARGET = 12.0; MTL_ENABLE_DEBUG_INFO = NO; SDKROOT = macosx; SWIFT_COMPILATION_MODE = wholemodule; @@ -545,7 +545,7 @@ GCC_WARN_UNINITIALIZED_AUTOS = YES_AGGRESSIVE; GCC_WARN_UNUSED_FUNCTION = YES; GCC_WARN_UNUSED_VARIABLE = YES; - MACOSX_DEPLOYMENT_TARGET = 10.15; + MACOSX_DEPLOYMENT_TARGET = 12.0; MTL_ENABLE_DEBUG_INFO = YES; ONLY_ACTIVE_ARCH = YES; SDKROOT = macosx; @@ -595,7 +595,7 @@ GCC_WARN_UNINITIALIZED_AUTOS = YES_AGGRESSIVE; GCC_WARN_UNUSED_FUNCTION = YES; GCC_WARN_UNUSED_VARIABLE = YES; - MACOSX_DEPLOYMENT_TARGET = 10.15; + MACOSX_DEPLOYMENT_TARGET = 12.0; MTL_ENABLE_DEBUG_INFO = NO; SDKROOT = macosx; SWIFT_COMPILATION_MODE = wholemodule; diff --git a/demos/supabase-todolist-drift/macos/Runner.xcodeproj/project.xcworkspace/xcshareddata/swiftpm/Package.resolved b/demos/supabase-todolist-drift/macos/Runner.xcodeproj/project.xcworkspace/xcshareddata/swiftpm/Package.resolved deleted file mode 100644 index adfc2b19..00000000 --- a/demos/supabase-todolist-drift/macos/Runner.xcodeproj/project.xcworkspace/xcshareddata/swiftpm/Package.resolved +++ /dev/null @@ -1,13 +0,0 @@ -{ - "pins" : [ - { - "identity" : "csqlite", - "kind" : "remoteSourceControl", - "location" : "https://github.com/simolus3/CSQLite.git", - "state" : { - "revision" : "a268235ae86718e66d6a29feef3bd22c772eb82b" - } - } - ], - "version" : 2 -} diff --git a/demos/supabase-todolist-drift/macos/Runner.xcworkspace/xcshareddata/swiftpm/Package.resolved b/demos/supabase-todolist-drift/macos/Runner.xcworkspace/xcshareddata/swiftpm/Package.resolved deleted file mode 100644 index adfc2b19..00000000 --- a/demos/supabase-todolist-drift/macos/Runner.xcworkspace/xcshareddata/swiftpm/Package.resolved +++ /dev/null @@ -1,13 +0,0 @@ -{ - "pins" : [ - { - "identity" : "csqlite", - "kind" : "remoteSourceControl", - "location" : "https://github.com/simolus3/CSQLite.git", - "state" : { - "revision" : "a268235ae86718e66d6a29feef3bd22c772eb82b" - } - } - ], - "version" : 2 -} From 373e3070bf12e260cffa8808df50f8dbd9820125 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 25 Aug 2026 16:52:09 +0200 Subject: [PATCH 4/6] Slightly simplify --- packages/powersync/lib/src/sync/checkpoint_request.dart | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/packages/powersync/lib/src/sync/checkpoint_request.dart b/packages/powersync/lib/src/sync/checkpoint_request.dart index 655f49d3..71a0df61 100644 --- a/packages/powersync/lib/src/sync/checkpoint_request.dart +++ b/packages/powersync/lib/src/sync/checkpoint_request.dart @@ -22,8 +22,6 @@ import 'sync_status.dart'; /// Requests do not survive [PowerSyncDatabase.disconnectAndClear], instances /// created before a clear should be discarded and requested again. abstract final class CheckpointRequest { - CheckpointRequest._(); - /// Whether this checkpoint request has synced before. @experimental bool get hasSynced; @@ -44,7 +42,7 @@ final class CheckpointRequestImpl extends CheckpointRequest { final Int64 _requestId; final ConnectionManager _connections; - CheckpointRequestImpl(this._requestId, this._connections) : super._(); + CheckpointRequestImpl(this._requestId, this._connections); @override bool get hasSynced { From bd7c569db82928c24b48522c9dac03d95cbfdfef Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Tue, 1 Sep 2026 13:16:19 +0200 Subject: [PATCH 5/6] AI feedback --- .../lib/src/platform_specific/web.dart | 4 +- .../lib/src/sync/checkpoint_state.dart | 54 ++++++++----------- .../lib/src/sync/connection_manager.dart | 36 +++---------- .../powersync/lib/src/sync/stream_utils.dart | 38 +++++++++++++ .../powersync/lib/src/web/sync_worker.dart | 2 + .../lib/src/web/sync_worker_protocol.dart | 2 +- 6 files changed, 72 insertions(+), 64 deletions(-) diff --git a/packages/powersync/lib/src/platform_specific/web.dart b/packages/powersync/lib/src/platform_specific/web.dart index 767bed9d..a5258e60 100644 --- a/packages/powersync/lib/src/platform_specific/web.dart +++ b/packages/powersync/lib/src/platform_specific/web.dart @@ -44,8 +44,6 @@ BasePowerSyncDatabase openPowerSyncDatabase( ); } -const _has64BitInts = !identical(0, 0.0); - final _intConversionBuffer = JSArrayBuffer(8); final _intConversionView = JSDataView(_intConversionBuffer); @@ -58,7 +56,7 @@ Int64 parseInt64(String s) { } Int64 int64FromBigInt(JSBigInt value) { - if (_has64BitInts) { + if (Int64.has64BitIntegers) { _intConversionView.setBigInt64(0, value); final high = _intConversionView.getUint32(0).toDartInt; final low = _intConversionView.getUint32(4).toDartInt; diff --git a/packages/powersync/lib/src/sync/checkpoint_state.dart b/packages/powersync/lib/src/sync/checkpoint_state.dart index 774f2ec8..fa87e8a1 100644 --- a/packages/powersync/lib/src/sync/checkpoint_state.dart +++ b/packages/powersync/lib/src/sync/checkpoint_state.dart @@ -1,9 +1,9 @@ import 'dart:async'; import 'package:async/async.dart'; -import 'package:sqlite_async/sqlite_async.dart' show AbortException; import 'checkpoint_request.dart'; +import 'stream_utils.dart'; final class CheckpointStateSignals { _CheckpointState _currentState = const _Pending(); @@ -57,37 +57,27 @@ final class CheckpointStateSignals { required Future abort, bool wakeDownloadLoop = true, }) { - final completer = Completer(); - final subscription = _stateController.stream.listen(null); - - void handleState(_CheckpointState state) { - switch (state) { - case _Disconnected(): - completer.completeError(disconnected); - subscription.cancel(); - case _CheckpointSeeded(:final result): - completer.complete(result.asFuture); - subscription.cancel(); - case _Pending(): - if (wakeDownloadLoop && !_waitingForWaiter.isCompleted) { - _waitingForWaiter.complete(); - } - } - } - - subscription.onData(handleState); - handleState(_currentState); - - abort.whenComplete(() { - if (!completer.isCompleted) { - completer.completeError( - AbortException('waitForCheckpointRequestsReady'), - ); - subscription.cancel(); - } - }); - - return completer.future; + return _stateController.stream.waitForFirstMatching( + predicate: (state) { + switch (state) { + case _Disconnected(): + throw disconnected; + case _CheckpointSeeded(:final result): + if (result.asError case final error?) { + Error.throwWithStackTrace(error.error, error.stackTrace); + } + return true; + case _Pending(): + if (wakeDownloadLoop && !_waitingForWaiter.isCompleted) { + _waitingForWaiter.complete(); + } + return false; + } + }, + debugName: 'waitForCheckpointRequestsReady', + abort: abort, + current: _currentState, + ); } } diff --git a/packages/powersync/lib/src/sync/connection_manager.dart b/packages/powersync/lib/src/sync/connection_manager.dart index fbf8f588..9fa22ab3 100644 --- a/packages/powersync/lib/src/sync/connection_manager.dart +++ b/packages/powersync/lib/src/sync/connection_manager.dart @@ -8,12 +8,12 @@ import 'package:powersync/src/database/active_instances.dart'; import 'package:powersync/src/sync/options.dart'; import 'package:powersync/src/sync/stream.dart'; import 'package:powersync/src/sync/sync_status.dart'; -import 'package:sqlite_async/sqlite_async.dart'; import '../database/powersync_database.dart'; import 'checkpoint_request.dart'; import 'instruction.dart'; import 'mutable_sync_status.dart'; +import 'stream_utils.dart'; import 'streaming_sync.dart'; /// A (stream name, JSON parameters) pair that uniquely identifies a stream @@ -102,32 +102,12 @@ final class ConnectionManager { bool Function(SyncStatus) predicate, { Future? abort, }) { - final completer = Completer(); - final subscription = statusStream.listen(null); - - void checkPredicate(SyncStatus status) { - try { - if (predicate(status)) { - completer.complete(); - subscription.cancel(); - } - } catch (e, s) { - completer.completeError(e, s); - subscription.cancel(); - } - } - - subscription.onData(checkPredicate); - checkPredicate(currentStatus); - - abort?.whenComplete(() { - if (!completer.isCompleted) { - completer.completeError(AbortException('firstStatusMatching')); - subscription.cancel(); - } - }); - - return completer.future; + return statusStream.waitForFirstMatching( + predicate: predicate, + debugName: 'firstStatusMatching', + abort: abort, + current: currentStatus, + ); } List get _subscribedStreams => [ @@ -190,7 +170,7 @@ final class ConnectionManager { if (!thisConnectAborter.aborted) { // We only change _abortActiveSync after disconnecting, which resets // the abort controller. - assert(identical(_abortActiveSync, thisConnectAborter)); + assert(identical(_abortActiveSync?.$1, thisConnectAborter)); // We need a new abort controller for this attempt thisConnectAborter = AbortController(); diff --git a/packages/powersync/lib/src/sync/stream_utils.dart b/packages/powersync/lib/src/sync/stream_utils.dart index feba0a37..2a85ad84 100644 --- a/packages/powersync/lib/src/sync/stream_utils.dart +++ b/packages/powersync/lib/src/sync/stream_utils.dart @@ -4,6 +4,7 @@ import 'dart:convert' as convert; import 'dart:math'; import 'dart:typed_data'; +import 'package:sqlite_async/sqlite_async.dart'; import 'package:typed_data/typed_buffers.dart'; import '../exceptions.dart'; @@ -166,6 +167,43 @@ Stream streamFromFutureAwaitInCancellation(Future future) { return controller.stream; } +extension WaitForFirst on Stream { + /// An abortable variant of [Stream.first]. + Future waitForFirstMatching({ + required bool Function(T) predicate, + required String debugName, + Future? abort, + T? current, + }) { + final completer = Completer(); + final subscription = listen(null); + + void check(T current) { + try { + if (predicate(current)) { + completer.complete(); + subscription.cancel(); + } + } catch (e, s) { + completer.completeError(e, s); + subscription.cancel(); + } + } + + subscription.onData(check); + if (current != null) check(current); + + abort?.whenComplete(() { + if (!completer.isCompleted) { + completer.completeError(AbortException(debugName)); + subscription.cancel(); + } + }); + + return completer.future; + } +} + /// An [EventSink] that takes raw bytes as inputs, buffers them internally by /// reading a 4-byte length prefix for each message and then emits them as /// chunks. diff --git a/packages/powersync/lib/src/web/sync_worker.dart b/packages/powersync/lib/src/web/sync_worker.dart index 0df7d207..eb67d612 100644 --- a/packages/powersync/lib/src/web/sync_worker.dart +++ b/packages/powersync/lib/src/web/sync_worker.dart @@ -122,6 +122,8 @@ class ConnectedClient { final sync = _runner!.sync!; final checkpoint = await sync.requestCheckpoint(); + // Workers are always compiled to JS, so we can safely assume the + // value here is a JS bigint. return ((checkpoint as JsBigInt64).value, null); default: throw StateError('Unexpected message type $type'); diff --git a/packages/powersync/lib/src/web/sync_worker_protocol.dart b/packages/powersync/lib/src/web/sync_worker_protocol.dart index 955c0ccf..0197f80e 100644 --- a/packages/powersync/lib/src/web/sync_worker_protocol.dart +++ b/packages/powersync/lib/src/web/sync_worker_protocol.dart @@ -396,7 +396,7 @@ extension type SerializedSyncStatus._(JSObject _) implements JSObject { }, lastAppliedCheckpoint: switch (lastAppliedCheckpoint) { null => null, - final id => JsBigInt64(id), + final id => int64FromBigInt(id), }, ); } From 1959878fdb36dd069f3d34232dda46a73dbab0c3 Mon Sep 17 00:00:00 2001 From: Simon Binder Date: Wed, 2 Sep 2026 16:23:28 +0200 Subject: [PATCH 6/6] Add todo comment to merge client interfaces --- packages/powersync/lib/src/sync/streaming_sync.dart | 2 ++ 1 file changed, 2 insertions(+) diff --git a/packages/powersync/lib/src/sync/streaming_sync.dart b/packages/powersync/lib/src/sync/streaming_sync.dart index 55b6e17f..3a27b65e 100644 --- a/packages/powersync/lib/src/sync/streaming_sync.dart +++ b/packages/powersync/lib/src/sync/streaming_sync.dart @@ -29,6 +29,8 @@ typedef SubscribedStream = ({String name, String parameters}); /// A superinterface for [StreamingSync] providing methods that can be called /// in remote contexts. +// TODO: Merge this into StreamingSync to reduce platform-specific code when +// managing sync clients. abstract interface class RemoteStreamingSyncHandle { Future requestCheckpoint(); }