Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
52 changes: 51 additions & 1 deletion Sources/SQLiteData/CloudKit/SyncEngine.swift
Original file line number Diff line number Diff line change
Expand Up @@ -436,11 +436,20 @@
$0 = nil
}
#endif
observationRegistrar.withMutation(of: self, keyPath: \.isRunning) {
let released = observationRegistrar.withMutation(of: self, keyPath: \.isRunning) {
syncEngines.withValue {
let released = ($0.private, $0.shared)
$0 = SyncEngines()
return released
}
}
// Released engines never balance the activity counts, so they would stay stuck above zero.
fetchingChangesCount = 0
sendingChangesCount = 0
Task {
await released.0?.cancelOperations()
await released.1?.cancelOperations()
}
}

/// Determines if the sync engine is currently running or not.
Expand Down Expand Up @@ -575,6 +584,31 @@
_ = try await (`private`, shared)
}

/// Applies records fetched from CloudKit outside the sync engine, through the same path the
/// engine's own `fetchedRecordZoneChanges` events take.
///
/// A manual `CKSyncEngine.fetchChanges` skips the server unless a push, a scheduled sync or an
/// app activation already flagged changes. This lets an app read a zone's changes itself and
/// still keep the sync metadata consistent. Applying a record the engine later delivers again
/// is a no-op.
public func applyFetchedRecordZoneChanges(
modifications: [CKRecord],
deletions: [(recordID: CKRecord.ID, recordType: CKRecord.RecordType)] = [],
scope: CKDatabase.Scope
) async {
await startTask.withValue(\.self)?.value
let syncEngine: (any SyncEngineProtocol)? = syncEngines.withValue {
guard $0.isRunning else { return nil }
return scope == .shared ? $0.shared : $0.private
}
guard let syncEngine else { return }
await handleFetchedRecordZoneChanges(
modifications: modifications,
deletions: deletions,
syncEngine: syncEngine
)
}

/// Sends pending local changes to the server.
///
/// Use this method to ensure the sync engine sends all pending local changes to the server
Expand Down Expand Up @@ -997,6 +1031,13 @@
await handleEvent(event, syncEngine: syncEngine)
}

private func isCurrent(_ syncEngine: any SyncEngineProtocol) -> Bool {
syncEngines.withValue {
guard $0.isRunning else { return false }
return $0.private === syncEngine || $0.shared === syncEngine
}
}

package func handleEvent(_ event: Event, syncEngine: any SyncEngineProtocol) async {
#if DEBUG
logger.log(event, syncEngine: syncEngine)
Expand All @@ -1006,6 +1047,9 @@
case .accountChange(let changeType):
await handleAccountChange(changeType: changeType, syncEngine: syncEngine)
case .stateUpdate(let stateSerialization):
// An engine released by `stop()` can still report state; saving it would overwrite the
// running engine's state.
guard isCurrent(syncEngine) else { return }
await handleStateUpdate(stateSerialization: stateSerialization, syncEngine: syncEngine)
case .fetchedDatabaseChanges(let modifications, let deletions):
await handleFetchedDatabaseChanges(
Expand Down Expand Up @@ -1036,28 +1080,34 @@
)

case .willFetchRecordZoneChanges:
guard isCurrent(syncEngine) else { return }
await MainActor.run {
fetchingChangesCount += 1
}
case .didFetchRecordZoneChanges:
guard isCurrent(syncEngine) else { return }
await MainActor.run {
fetchingChangesCount -= 1
}

case .willFetchChanges:
guard isCurrent(syncEngine) else { return }
await MainActor.run {
fetchingChangesCount += 1
}
case .didFetchChanges:
guard isCurrent(syncEngine) else { return }
await MainActor.run {
fetchingChangesCount -= 1
}

case .willSendChanges:
guard isCurrent(syncEngine) else { return }
await MainActor.run {
sendingChangesCount += 1
}
case .didSendChanges:
guard isCurrent(syncEngine) else { return }
await MainActor.run {
sendingChangesCount -= 1
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,35 @@
@MainActor
@Suite
final class FetchRecordZoneChangeTests: BaseCloudKitTests, @unchecked Sendable {
@available(iOS 17, macOS 14, tvOS 17, watchOS 10, *)
@Test func applyRecordsFetchedOutsideTheEngine() async throws {
try await userDatabase.userWrite { db in
try db.seed {
RemindersList(id: 1, title: "Personal")
Reminder(id: 1, title: "Get milk", remindersListID: 1)
}
}
try await syncEngine.processPendingRecordZoneChanges(scope: .private)

let record = try syncEngine.private.database.record(for: Reminder.recordID(for: 1))
record.setValue("Buy milk", forKey: "title", at: 60)
let (saveResults, _) = try syncEngine.private.database.modifyRecords(
saving: [record],
deleting: [],
atomically: true
)
let serverRecord = try #require(try saveResults[record.recordID]?.get())

await syncEngine.applyFetchedRecordZoneChanges(modifications: [serverRecord], scope: .private)
await syncEngine.applyFetchedRecordZoneChanges(modifications: [serverRecord], scope: .private)

let title = try await userDatabase.read { db in
try Reminder.find(1).select(\.title).fetchOne(db)
}
#expect(title == "Buy milk")
#expect(syncEngine.private.state.pendingRecordZoneChanges.isEmpty)
}

@available(iOS 17, macOS 14, tvOS 17, watchOS 10, *)
@Test func saveExtraFieldsToSyncMetadata() async throws {
try await userDatabase.userWrite { db in
Expand Down