From 19976a091a49c4ca4b105f4d8166454224f07558 Mon Sep 17 00:00:00 2001 From: "Asier G. Morato" Date: Wed, 9 Sep 2026 09:15:21 +0200 Subject: [PATCH 1/3] Apply subscription changes made while the sync stream is being established A sync stream subscription added after `connect()` returned but before the first `/sync/stream` response arrived only took effect on the next keep-alive line from the service, typically 20 seconds later. On a fresh database that window covers the common "connect, then subscribe" pattern, so a client's first stream ran with no subscriptions at all and delivered nothing until then. `SyncIteration.run()` started `watchSyncStreams` and `watchCompletedCrudUploads` before the request went out, but only subscribed to `localEvents` after `fetchSyncLines` returned. `BroadcastStream.dispatch` drops events with no listener, so an `updateSubscriptions` raised in that gap never reached the core extension, which fell back to its periodic check on keep-alive. Subscribe to `localEvents` before the watchers start. The stream buffers everything dispatched until the control loop consumes it, so the change is applied as soon as the response arrives and the iteration reconnects with the new subscription set right away. The same fix covers a CRUD upload completing during connection establishment. Adds a regression test that subscribes while the first request is held open and expects a second request carrying the subscription without any keep-alive. --- CHANGELOG.md | 8 +++ .../sync/StreamingSyncClient.swift | 11 +++- Tests/PowerSyncTests/SyncTests.swift | 51 +++++++++++++++++++ 3 files changed, 69 insertions(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 4c2fe61f..15ab93c0 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,13 @@ # Changelog +## Unreleased + +* Fix sync stream subscriptions changed while a connection was being established (after + `connect()` returned and before the first `/sync/stream` response arrived) taking effect only + on the next keep-alive line from the service, typically 20 seconds later. The change + notification was dispatched before the sync iteration started listening for local events, so + it was dropped and the core extension never learned about it until its periodic check. + ## 1.16.1 * Fix a re-query storm on databases opened at an absolute path (App Group container). The diff --git a/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift b/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift index 69bf910f..9df1d1b1 100644 --- a/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift +++ b/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift @@ -514,6 +514,15 @@ private struct ActiveSyncIteration: Sendable { signals.markPendingCheckpointRequestsRequiringAffirmation() } + // Listen for local events BEFORE the watchers below start dispatching them. The stream + // buffers whatever is dispatched until the control loop consumes it, so a subscription + // change or a completed upload that lands while the request is still being established + // takes effect as soon as the response arrives. Subscribing only after `fetchSyncLines` + // returned left those events without a listener (`BroadcastStream.dispatch` drops them), + // and the core extension then only noticed the change on the next keep-alive line, + // typically 20 seconds later. + let pendingLocalEvents = localEvents.subscribe() + // Notify the core extension for changed Sync Stream subscriptions, as we might have to reconnect. let (currentStreams, streamChanges) = syncClient.db.group.syncCoordinator.streams.observeActiveStreams() async let _ = watchSyncStreams(changes: streamChanges) @@ -557,7 +566,7 @@ private struct ActiveSyncIteration: Sendable { controlArgs = AsyncAlgorithms.merge( serviceEvents, checkpointRequestStateValidationEvents(task: checkpointRequestStateSeed), - localEvents.subscribe() + pendingLocalEvents ) } catch { let streamError = error diff --git a/Tests/PowerSyncTests/SyncTests.swift b/Tests/PowerSyncTests/SyncTests.swift index ec27e3c3..17a2c532 100644 --- a/Tests/PowerSyncTests/SyncTests.swift +++ b/Tests/PowerSyncTests/SyncTests.swift @@ -2192,6 +2192,57 @@ class InMemorySyncIntegrationTests { try await db.close() } + @Test func subscriptionAddedWhileConnectingReconnects() async throws { + // A subscription added after `connect()` returned but before the first `/sync/stream` + // response arrived used to be applied only on the next keep-alive line from the service, + // typically 20 seconds later: the change was dispatched before the sync iteration was + // listening for local events. Nothing here sends a keep-alive, so the reconnect has to + // come from the subscription change itself. + let firstRequestStarted = Signal() + let releaseFirstResponse = Signal() + let requests = AsyncMutex<[JsonParam]>([]) + let db = openDatabase(MockHttpSession { request in + let body = try StreamingSyncClient.jsonDecoder.decode(JsonParam.self, from: try #require(request.httpBody)) + let count = await requests.withMutex { requests in + requests.append(body) + return requests.count + } + if count == 1 { + await firstRequestStarted.complete() + await releaseFirstResponse.await() + } + return AsyncThrowingChannel() + }) + + try await db.connect(connector: TestConnector(), options: ConnectOptions()) + await firstRequestStarted.await() + + // The first request is in flight and its response has not arrived yet. + let subscription = try await db.syncStream(name: "a", params: nil).subscribe() + await releaseFirstResponse.complete() + + try await waitUntilAsync { await requests.inner.count == 2 } + let allRequests = await requests.inner + if case let .object(streams) = allRequests[0]["streams"] { + try #require(streams["subscriptions"] == .array([])) + } else { + Issue.record("Should have streams key in first body") + } + if case let .object(streams) = allRequests[1]["streams"] { + try #require(streams["subscriptions"] == .array([ + .object([ + "stream": .string("a"), + "parameters": .null, + "override_priority": .null, + ]) + ])) + } else { + Issue.record("Should have streams key in second body") + } + let _ = consume subscription + try await db.close() + } + @Test func subscriptionsUpdateWhileOffline() async throws { let db = openDatabase(MockHttpSession { request in throw PowerSyncError.operationFailed(message: "Unexpected connection", underlyingError: nil) From 89ff70646d04778e02075f74aad3c2903625c26c Mon Sep 17 00:00:00 2001 From: "Asier G. Morato" Date: Wed, 9 Sep 2026 09:57:08 +0200 Subject: [PATCH 2/3] Tighten the wording and decouple the regression test from the retry delay State the invariant in the comment instead of the history, and drop the claim about completed uploads: that watcher still subscribes inside its own task, and an upload finishing before any checkpoint exists is a no-op in the core anyway. Give the test a short retry delay so its 5 s polling budget no longer coincides with the default retry interval. --- CHANGELOG.md | 7 ++----- .../Implementation/sync/StreamingSyncClient.swift | 11 ++++------- Tests/PowerSyncTests/SyncTests.swift | 12 ++++++------ 3 files changed, 12 insertions(+), 18 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 15ab93c0..304f30d8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -2,11 +2,8 @@ ## Unreleased -* Fix sync stream subscriptions changed while a connection was being established (after - `connect()` returned and before the first `/sync/stream` response arrived) taking effect only - on the next keep-alive line from the service, typically 20 seconds later. The change - notification was dispatched before the sync iteration started listening for local events, so - it was dropped and the core extension never learned about it until its periodic check. +* Fix sync stream subscriptions changed while the connection was still being established only + taking effect on the next keep-alive line from the service, typically 20 seconds later. ## 1.16.1 diff --git a/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift b/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift index 9df1d1b1..e93eac33 100644 --- a/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift +++ b/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift @@ -514,13 +514,10 @@ private struct ActiveSyncIteration: Sendable { signals.markPendingCheckpointRequestsRequiringAffirmation() } - // Listen for local events BEFORE the watchers below start dispatching them. The stream - // buffers whatever is dispatched until the control loop consumes it, so a subscription - // change or a completed upload that lands while the request is still being established - // takes effect as soon as the response arrives. Subscribing only after `fetchSyncLines` - // returned left those events without a listener (`BroadcastStream.dispatch` drops them), - // and the core extension then only noticed the change on the next keep-alive line, - // typically 20 seconds later. + // Subscribe to local events before the watchers below start dispatching them: `BroadcastStream` + // only delivers to listeners that already exist, and the subscription buffers everything until + // the control loop drains it once the sync stream is established. Otherwise a subscription + // change made while the request is in flight is only picked up on the next keep-alive line. let pendingLocalEvents = localEvents.subscribe() // Notify the core extension for changed Sync Stream subscriptions, as we might have to reconnect. diff --git a/Tests/PowerSyncTests/SyncTests.swift b/Tests/PowerSyncTests/SyncTests.swift index 17a2c532..cb07900c 100644 --- a/Tests/PowerSyncTests/SyncTests.swift +++ b/Tests/PowerSyncTests/SyncTests.swift @@ -2193,11 +2193,9 @@ class InMemorySyncIntegrationTests { } @Test func subscriptionAddedWhileConnectingReconnects() async throws { - // A subscription added after `connect()` returned but before the first `/sync/stream` - // response arrived used to be applied only on the next keep-alive line from the service, - // typically 20 seconds later: the change was dispatched before the sync iteration was - // listening for local events. Nothing here sends a keep-alive, so the reconnect has to - // come from the subscription change itself. + // Regression test: a subscription added after `connect()` returned but before the first + // `/sync/stream` response arrived must trigger a reconnect on its own. Nothing here sends a + // keep-alive, so the second request can only come from the subscription change. let firstRequestStarted = Signal() let releaseFirstResponse = Signal() let requests = AsyncMutex<[JsonParam]>([]) @@ -2214,7 +2212,9 @@ class InMemorySyncIntegrationTests { return AsyncThrowingChannel() }) - try await db.connect(connector: TestConnector(), options: ConnectOptions()) + // A short retry delay keeps the 5 s polling budget below independent of how the reconnect + // is scheduled. + try await db.connect(connector: TestConnector(), options: ConnectOptions(retryDelay: 0.1)) await firstRequestStarted.await() // The first request is in flight and its response has not arrived yet. From ff7e1347916bd9fbb2956c4942fc1f1e4ff9bcdc Mon Sep 17 00:00:00 2001 From: "Asier G. Morato" Date: Wed, 9 Sep 2026 12:56:00 +0200 Subject: [PATCH 3/3] Address review: drop the comment and the retry delay, prepare 1.16.2 The comment reads as obvious once the ordering is known, and the test's short retry delay was defensive rather than needed: the close a subscription change produces carries hide_disconnect, so the reconnect never goes through the retry path and the test runs in about 60 ms either way. Sets the CHANGELOG heading and libraryVersion to 1.16.2 for the release. --- CHANGELOG.md | 2 +- Sources/PowerSync/CurrentVersion.swift | 2 +- .../PowerSync/Implementation/sync/StreamingSyncClient.swift | 4 ---- Tests/PowerSyncTests/SyncTests.swift | 4 +--- 4 files changed, 3 insertions(+), 9 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 304f30d8..3a2dfefe 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,6 @@ # Changelog -## Unreleased +## 1.16.2 * Fix sync stream subscriptions changed while the connection was still being established only taking effect on the next keep-alive line from the service, typically 20 seconds later. diff --git a/Sources/PowerSync/CurrentVersion.swift b/Sources/PowerSync/CurrentVersion.swift index 350972f5..ecbceee7 100644 --- a/Sources/PowerSync/CurrentVersion.swift +++ b/Sources/PowerSync/CurrentVersion.swift @@ -1,2 +1,2 @@ // The current version of the PowerSync Swift SDK. This should be updated to the latest version in `CHANGELOG.md` when a new version is released. -let libraryVersion = "1.16.1" +let libraryVersion = "1.16.2" diff --git a/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift b/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift index e93eac33..18bf0d81 100644 --- a/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift +++ b/Sources/PowerSync/Implementation/sync/StreamingSyncClient.swift @@ -514,10 +514,6 @@ private struct ActiveSyncIteration: Sendable { signals.markPendingCheckpointRequestsRequiringAffirmation() } - // Subscribe to local events before the watchers below start dispatching them: `BroadcastStream` - // only delivers to listeners that already exist, and the subscription buffers everything until - // the control loop drains it once the sync stream is established. Otherwise a subscription - // change made while the request is in flight is only picked up on the next keep-alive line. let pendingLocalEvents = localEvents.subscribe() // Notify the core extension for changed Sync Stream subscriptions, as we might have to reconnect. diff --git a/Tests/PowerSyncTests/SyncTests.swift b/Tests/PowerSyncTests/SyncTests.swift index cb07900c..a0677149 100644 --- a/Tests/PowerSyncTests/SyncTests.swift +++ b/Tests/PowerSyncTests/SyncTests.swift @@ -2212,9 +2212,7 @@ class InMemorySyncIntegrationTests { return AsyncThrowingChannel() }) - // A short retry delay keeps the 5 s polling budget below independent of how the reconnect - // is scheduled. - try await db.connect(connector: TestConnector(), options: ConnectOptions(retryDelay: 0.1)) + try await db.connect(connector: TestConnector(), options: ConnectOptions()) await firstRequestStarted.await() // The first request is in flight and its response has not arrived yet.