diff --git a/Sources/CodexBar/Sync/CloudSyncEngine.swift b/Sources/CodexBar/Sync/CloudSyncEngine.swift index 48d6ddf845..d9861e1d25 100644 --- a/Sources/CodexBar/Sync/CloudSyncEngine.swift +++ b/Sources/CodexBar/Sync/CloudSyncEngine.swift @@ -150,6 +150,225 @@ enum CloudSyncDirtyState { } } +enum CloudSyncSnapshotMigration { + static func obsoleteRecordNames( + liveSnapshots: [AccountSnapshotSyncPayload], + hashes: [String: String], + envelope: CloudSyncPersistence.Envelope) -> Set + { + AccountSnapshotSyncPayload.obsoleteEmailKeyedRecordNames( + liveSnapshots: liveSnapshots, + knownRecordNames: Set(hashes.keys).union(envelope.fleetSnapshots.keys)) + } + + static func drop( + _ names: Set, + hashes: inout [String: String], + envelope: inout CloudSyncPersistence.Envelope, + desiredRecords: inout [CKRecord.ID: CKRecord], + zoneID: CKRecordZone.ID) -> [CKRecord.ID] + { + names.map { name in + let recordID = CKRecord.ID(recordName: name, zoneID: zoneID) + desiredRecords.removeValue(forKey: recordID) + hashes.removeValue(forKey: name) + envelope.fleetSnapshots.removeValue(forKey: name) + envelope.encodedSystemFields.removeValue(forKey: name) + envelope.recordMetadata.removeValue(forKey: name) + return recordID + } + } + + static func predecessorNames( + for snapshot: AccountSnapshotSyncPayload, + obsoleteNames: Set) -> Set + { + guard let predecessor = snapshot.emailKeyedPredecessorRecordName(), + obsoleteNames.contains(predecessor) + else { + return [] + } + return [predecessor] + } + + static func takeDeletes( + forSavedRecordNames savedNames: [String], + pending: inout [String: Set], + afterLiveSnapshotReconciliation hasReconciledLiveSnapshots: Bool) -> Set + { + guard hasReconciledLiveSnapshots else { return [] } + return self.takeDeletes(forSavedRecordNames: savedNames, pending: &pending) + } + + static func takeDeletes( + forSavedRecordNames savedNames: [String], + pending: inout [String: Set]) -> Set + { + var toDrop: Set = [] + for name in savedNames { + if let obsolete = pending.removeValue(forKey: name) { + toDrop.formUnion(obsolete) + } + } + let stillReferenced = Set(pending.values.joined()) + return toDrop.subtracting(stillReferenced) + } + + static func retainingObsoletePredecessors( + in pending: inout [String: Set], + obsoleteNames: Set) + { + pending = pending.compactMapValues { predecessors in + let live = predecessors.intersection(obsoleteNames) + return live.isEmpty ? nil : live + } + } + + static func assigningPredecessors( + _ predecessors: Set, + to replacement: String, + pending: inout [String: Set]) + { + if predecessors.isEmpty { + pending.removeValue(forKey: replacement) + } else { + pending[replacement] = predecessors + } + } + + static func cancelledPersistedDeletes( + pendingDeletes: Set, + liveNames: Set) -> Set + { + pendingDeletes.intersection(liveNames) + } + + static func pendingDeletesToRequeue( + pendingDeletes: Set, + liveNames: Set) -> Set + { + pendingDeletes.subtracting(liveNames) + } + + static func liveSnapshotRecordNames( + pendingRecordNames: some Sequence, + storedRecordNames: some Sequence) -> Set + { + Set(pendingRecordNames).union(storedRecordNames) + } + + static func retryableFailedDeletes( + _ failures: [CKRecord.ID: CKError], + liveNames: Set = []) -> [CKRecord.ID] + { + failures.compactMap { recordID, error in + guard self.retryDelay(for: error) != nil else { return nil } + guard !liveNames.contains(recordID.recordName) else { return nil } + return recordID + } + } + + static func reportableFailedDeletes(_ failures: [CKRecord.ID: CKError]) -> [CKError] { + failures.values.filter { error in + error.code != .unknownItem && self.retryDelay(for: error) == nil + } + } + + /// Delayed retries only for recoverable CloudKit failures. Terminal per-record errors such as + /// `permissionFailure`, `notAuthenticated`, and `invalidArguments` are reported once. + static func retryDelay(for error: CKError) -> TimeInterval? { + switch error.code { + case .unknownItem, .permissionFailure, .notAuthenticated, .invalidArguments: + return nil + case .networkUnavailable, .networkFailure, .serviceUnavailable, .requestRateLimited, .zoneBusy, .quotaExceeded, + .serverResponseLost, .accountTemporarilyUnavailable: + return max(error.retryAfterSeconds ?? 1, 1) + default: + guard let retryAfter = error.retryAfterSeconds else { return nil } + return max(retryAfter, 1) + } + } + + static func finishedFailedDeleteNames(_ failures: [CKRecord.ID: CKError]) -> Set { + Set(failures.compactMap { recordID, error in + self.retryDelay(for: error) == nil ? recordID.recordName : nil + }) + } + + static func abandonedReplacementNames( + failures: [String: CKError], + pendingReplacements: Set) -> Set + { + Set(failures.compactMap { name, error in + pendingReplacements.contains(name) && self.retryDelay(for: error) == nil ? name : nil + }) + } + + static func applyConfirmedSaveHashes( + savedRecordNames: [String], + pendingSaveHashes: inout [String: String], + lastSnapshotHashes: inout [String: String]) + { + for name in savedRecordNames { + guard let hash = pendingSaveHashes.removeValue(forKey: name) else { continue } + lastSnapshotHashes[name] = hash + } + } + + static func applyTerminalSaveSkip( + recordName: String, + error: CKError, + pendingSaveHashes: inout [String: String], + skippedTerminalReplacementHashes: inout [String: String]) + { + guard self.retryDelay(for: error) == nil else { return } + guard let hash = pendingSaveHashes.removeValue(forKey: recordName) else { return } + skippedTerminalReplacementHashes[recordName] = hash + } + + static func hasInFlightSave(recordName: String, pendingSaveHashes: [String: String]) -> Bool { + pendingSaveHashes[recordName] != nil + } + + static func mergingPendingSnapshots( + _ pending: [AccountSnapshotSyncPayload], + with extras: [AccountSnapshotSyncPayload]) -> [AccountSnapshotSyncPayload] + { + var byName: [String: Int] = [:] + var result = pending + for (index, payload) in pending.enumerated() { + byName[payload.recordName] = index + } + for payload in extras where byName[payload.recordName] == nil { + byName[payload.recordName] = result.count + result.append(payload) + } + return result + } + + static func unpublishedFleetSnapshots( + savedRecordNames: [String], + fleetSnapshots: [String: AccountSnapshotSyncPayload], + lastSnapshotHashes: [String: String]) -> [AccountSnapshotSyncPayload] + { + savedRecordNames.compactMap { name in + guard let payload = fleetSnapshots[name], + let hash = try? CanonicalSyncJSON.hash(payload), + lastSnapshotHashes[name] != hash + else { return nil } + return payload + } + } + + static func shouldResumeDelayedRetry( + originatingEngine: ObjectIdentifier?, + currentEngine: ObjectIdentifier?) -> Bool + { + guard let originatingEngine, let currentEngine else { return false } + return originatingEngine == currentEngine + } +} + enum CloudSyncEntitlementGate { static let entitlement = "com.apple.developer.icloud-services" @@ -184,11 +403,15 @@ actor CloudSyncEngine: CKSyncEngineDelegate { private var lastSnapshotPushAt: Date? private var pendingSnapshots: [AccountSnapshotSyncPayload] = [] private var lastSnapshotHashes: [String: String] = [:] + private var skippedTerminalReplacementHashes: [String: String] = [:] + private var pendingSaveHashes: [String: String] = [:] private var lastKnownProviderConfigs: [ProviderInstanceID: ProviderConfig] = [:] private var lastKnownPreferences: SyncedPreferences? private var lastKnownIncludeSecrets: Bool? private var quotaRetryState = CloudSyncQuotaRetryState() private var didRehydrateFleetState = false + /// Restored predecessor mappings must not delete until local live snapshots have been applied. + private var hasReconciledLiveSnapshots = false private let logger = CodexBarLog.logger(LogCategories.settings) init( @@ -524,35 +747,6 @@ actor CloudSyncEngine: CKSyncEngineDelegate { self.persistEnvelope() } - private func pushPendingSnapshots() async { - guard let engine = self.engine, !self.pendingSnapshots.isEmpty else { return } - guard await MainActor.run(body: { !self.state.status.needsAppUpdate }) else { return } - do { - for payload in self.pendingSnapshots { - let hash = try CanonicalSyncJSON.hash(payload) - guard self.lastSnapshotHashes[payload.recordName] != hash else { continue } - let recordID = self.recordID(named: payload.recordName) - let record = self.record(type: .accountSnapshot, id: recordID) - record["schemaVersion"] = payload.schemaVersion as CKRecordValue - record["provider"] = payload.provider.rawValue as CKRecordValue - record["deviceID"] = payload.deviceID as CKRecordValue - record["accountKey"] = payload.accountKey as CKRecordValue - record["fetchedAt"] = payload.fetchedAt as CKRecordValue - record.encryptedValues["displayLabel"] = payload.displayLabel as CKRecordValue - record.encryptedValues["usagePayload"] = try CanonicalSyncJSON.string(payload.usage) as CKRecordValue - self.desiredRecords[recordID] = record - self.lastSnapshotHashes[payload.recordName] = hash - self.persistenceEnvelope.fleetSnapshots[payload.recordName] = payload - engine.state.add(pendingRecordZoneChanges: [.saveRecord(recordID)]) - } - self.pendingSnapshots = [] - self.lastSnapshotPushAt = Date() - self.persistEnvelope() - } catch { - await self.record(error: error) - } - } - // MARK: CKSyncEngineDelegate nonisolated func handleEvent(_ event: CKSyncEngine.Event, syncEngine: CKSyncEngine) async { @@ -603,13 +797,26 @@ actor CloudSyncEngine: CKSyncEngineDelegate { CloudSyncDirtyState.clearSavedRecords( changes.savedRecords.map(\.recordID.recordName), envelope: &self.persistenceEnvelope) + // Evaluate confirmed saves before abandoning terminal failures so a mixed batch + // still counts a failed sibling's shared predecessor as referenced. + await self.finishConfirmedSnapshotMigrations( + savedRecordNames: changes.savedRecords.map(\.recordID.recordName), + syncEngine: syncEngine) for failure in changes.failedRecordSaves { await self.handleSaveFailure(failure, syncEngine: syncEngine) } + await self.handleSentRecordDeletes( + deletedIDs: changes.deletedRecordIDs, + failures: changes.failedRecordDeletes) self.persistEnvelope() if !changes.savedRecords.isEmpty { await MainActor.run { self.state.status.lastSuccessfulPushAt = Date() } } + if !self.pendingSnapshots.isEmpty { + Task { [weak self] in + await self?.pushPendingSnapshots() + } + } case let .fetchedDatabaseChanges(changes): if changes.deletions.contains(where: { $0.zoneID == Self.zoneID }) { syncEngine.state.add(pendingDatabaseChanges: [.saveZone(CKRecordZone(zoneID: Self.zoneID))]) @@ -761,6 +968,10 @@ actor CloudSyncEngine: CKSyncEngineDelegate { } private func applyAccountSnapshot(_ record: CKRecord) async throws { + guard !CloudSyncSnapshotMigration.hasInFlightSave( + recordName: record.recordID.recordName, + pendingSaveHashes: self.pendingSaveHashes) + else { return } guard let providerRaw = record["provider"] as? String, let provider = ProviderInstanceID(rawValue: providerRaw), let deviceID = record["deviceID"] as? String, @@ -790,9 +1001,13 @@ actor CloudSyncEngine: CKSyncEngineDelegate { case .quotaExceeded: let retry = self.quotaRetryState.nextDelay(serverRetryAfter: failure.error.retryAfterSeconds) self.scheduleRetry(recordID: failure.record.recordID, after: retry) + case .accountTemporarilyUnavailable: + let retry = CloudSyncSnapshotMigration.retryDelay(for: failure.error) ?? 1 + self.scheduleRetry(recordID: failure.record.recordID, after: retry) case .serverRecordChanged: guard let server = failure.error.serverRecord else { await self.record(error: failure.error) + self.pendingSaveHashes.removeValue(forKey: failure.record.recordID.recordName) return } await self.resolveConflict(with: server, syncEngine: syncEngine) @@ -805,6 +1020,7 @@ actor CloudSyncEngine: CKSyncEngineDelegate { self.recreateZoneAndRequeue(failure.record, syncEngine: syncEngine) } else { await self.record(error: failure.error) + self.abandonTerminalReplacementSave(failure) } } } @@ -837,6 +1053,7 @@ actor CloudSyncEngine: CKSyncEngineDelegate { guard self.schemaVersion(server) <= CodexBarSyncSchema.currentVersion else { syncEngine.state.remove(pendingRecordZoneChanges: [.saveRecord(server.recordID)]) self.desiredRecords.removeValue(forKey: server.recordID) + self.pendingSaveHashes.removeValue(forKey: server.recordID.recordName) self.cacheSystemFields(server) self.persistEnvelope() await MainActor.run { self.state.status.needsAppUpdate = true } @@ -858,6 +1075,7 @@ actor CloudSyncEngine: CKSyncEngineDelegate { } else { syncEngine.state.remove(pendingRecordZoneChanges: [.saveRecord(server.recordID)]) self.desiredRecords.removeValue(forKey: server.recordID) + self.pendingSaveHashes.removeValue(forKey: server.recordID.recordName) await self.applyFetchedRecords([server]) } } @@ -908,8 +1126,13 @@ actor CloudSyncEngine: CKSyncEngineDelegate { self.engine = nil self.desiredRecords = [:] self.quotaRetryState.reset() + self.pendingSaveHashes = [:] + self.pendingSnapshots = [] + self.hasReconciledLiveSnapshots = false if clearPersistence { self.persistenceEnvelope = .init(stateSerialization: nil, encodedSystemFields: [:]) + self.lastSnapshotHashes = [:] + self.skippedTerminalReplacementHashes = [:] self.didRehydrateFleetState = false try? self.persistence.delete() await MainActor.run { @@ -1024,3 +1247,215 @@ actor CloudSyncEngine: CKSyncEngineDelegate { return String(bytes: bytes, encoding: .utf8) ?? "unknown" } } + +extension CloudSyncEngine { + private func finishConfirmedSnapshotMigrations( + savedRecordNames: [String], + syncEngine: CKSyncEngine) async + { + let toDrop = CloudSyncSnapshotMigration.takeDeletes( + forSavedRecordNames: savedRecordNames, + pending: &self.persistenceEnvelope.pendingPredecessorDeletes, + afterLiveSnapshotReconciliation: self.hasReconciledLiveSnapshots) + CloudSyncSnapshotMigration.applyConfirmedSaveHashes( + savedRecordNames: savedRecordNames, + pendingSaveHashes: &self.pendingSaveHashes, + lastSnapshotHashes: &self.lastSnapshotHashes) + self.pendingSnapshots = CloudSyncSnapshotMigration.mergingPendingSnapshots( + self.pendingSnapshots, + with: CloudSyncSnapshotMigration.unpublishedFleetSnapshots( + savedRecordNames: savedRecordNames, + fleetSnapshots: self.persistenceEnvelope.fleetSnapshots, + lastSnapshotHashes: self.lastSnapshotHashes)) + guard !toDrop.isEmpty else { return } + let recordIDs = CloudSyncSnapshotMigration.drop( + toDrop, + hashes: &self.lastSnapshotHashes, + envelope: &self.persistenceEnvelope, + desiredRecords: &self.desiredRecords, + zoneID: Self.zoneID) + for recordID in recordIDs { + syncEngine.state.add(pendingRecordZoneChanges: [.deleteRecord(recordID)]) + } + self.rememberPendingSnapshotDeletes(toDrop) + await MainActor.run { + toDrop.forEach { self.state.fleetSnapshots.removeValue(forKey: $0) } + } + } + + private func handleSentRecordDeletes(deletedIDs: [CKRecord.ID], failures: [CKRecord.ID: CKError]) async { + var finished = Set(deletedIDs.map(\.recordName)) + finished.formUnion(CloudSyncSnapshotMigration.finishedFailedDeleteNames(failures)) + self.forgetPendingSnapshotDeletes(finished) + for error in CloudSyncSnapshotMigration.reportableFailedDeletes(failures) { + await self.record(error: error) + } + let liveNames = CloudSyncSnapshotMigration.liveSnapshotRecordNames( + pendingRecordNames: self.pendingSnapshots.map(\.recordName) + + self.desiredRecords.keys.map(\.recordName), + storedRecordNames: self.lastSnapshotHashes.keys) + for recordID in CloudSyncSnapshotMigration.retryableFailedDeletes(failures, liveNames: liveNames) { + self.rememberPendingSnapshotDeletes([recordID.recordName]) + let delay = failures[recordID].flatMap(CloudSyncSnapshotMigration.retryDelay(for:)) ?? 1 + self.scheduleDeleteRetry(recordID: recordID, after: delay) + } + } + + private func rememberPendingSnapshotDeletes(_ names: Set) { + guard !names.isEmpty else { return } + self.persistenceEnvelope.pendingSnapshotDeletes.formUnion(names) + self.persistEnvelope() + } + + private func forgetPendingSnapshotDeletes(_ names: Set) { + let remaining = self.persistenceEnvelope.pendingSnapshotDeletes.subtracting(names) + guard remaining != self.persistenceEnvelope.pendingSnapshotDeletes else { return } + self.persistenceEnvelope.pendingSnapshotDeletes = remaining + self.persistEnvelope() + } + + private func cancelPendingSnapshotDeletes(_ names: Set) { + guard !names.isEmpty else { return } + self.forgetPendingSnapshotDeletes(names) + guard let engine = self.engine else { return } + engine.state.remove(pendingRecordZoneChanges: names.map { name in + .deleteRecord(self.recordID(named: name)) + }) + } + + private func requeuePendingSnapshotDeletes() { + guard let engine = self.engine else { return } + let liveNames = CloudSyncSnapshotMigration.liveSnapshotRecordNames( + pendingRecordNames: self.pendingSnapshots.map(\.recordName) + + self.desiredRecords.keys.map(\.recordName), + storedRecordNames: self.lastSnapshotHashes.keys) + let names = CloudSyncSnapshotMigration.pendingDeletesToRequeue( + pendingDeletes: self.persistenceEnvelope.pendingSnapshotDeletes, + liveNames: liveNames) + for name in names { + engine.state.add(pendingRecordZoneChanges: [.deleteRecord(self.recordID(named: name))]) + } + } + + private func abandonTerminalReplacementSave( + _ failure: CKSyncEngine.Event.SentRecordZoneChanges.FailedRecordSave) + { + let name = failure.record.recordID.recordName + let abandoned = CloudSyncSnapshotMigration.abandonedReplacementNames( + failures: [name: failure.error], + pendingReplacements: Set(self.persistenceEnvelope.pendingPredecessorDeletes.keys)) + if abandoned.contains(name) { + self.persistenceEnvelope.pendingPredecessorDeletes.removeValue(forKey: name) + } + CloudSyncSnapshotMigration.applyTerminalSaveSkip( + recordName: name, + error: failure.error, + pendingSaveHashes: &self.pendingSaveHashes, + skippedTerminalReplacementHashes: &self.skippedTerminalReplacementHashes) + } + + private func scheduleDeleteRetry(recordID: CKRecord.ID, after delay: TimeInterval) { + Task { [weak self] in + let originatingEngine = await self?.engine.map { ObjectIdentifier($0) } + do { + if delay > 0 { + try await Task.sleep(for: .seconds(delay)) + } + await Task.yield() + guard let self, await self.enabled else { return } + guard await self.persistenceEnvelope.pendingSnapshotDeletes.contains(recordID.recordName) else { + return + } + guard let engine = await self.engine, + CloudSyncSnapshotMigration.shouldResumeDelayedRetry( + originatingEngine: originatingEngine, + currentEngine: ObjectIdentifier(engine)) + else { return } + engine.state.add(pendingRecordZoneChanges: [.deleteRecord(recordID)]) + try await engine.sendChanges(.init(scope: .recordIDs([recordID]))) + } catch is CancellationError { + return + } catch { + await self?.record(error: error) + } + } + } + + private func pushPendingSnapshots() async { + guard let engine = self.engine else { return } + guard !self.pendingSnapshots.isEmpty || !self.persistenceEnvelope.pendingSnapshotDeletes.isEmpty else { + return + } + guard await MainActor.run(body: { !self.state.status.needsAppUpdate }) else { return } + do { + let obsoleteNames = CloudSyncSnapshotMigration.obsoleteRecordNames( + liveSnapshots: self.pendingSnapshots, + hashes: self.lastSnapshotHashes, + envelope: self.persistenceEnvelope) + if !self.pendingSnapshots.isEmpty { + CloudSyncSnapshotMigration.retainingObsoletePredecessors( + in: &self.persistenceEnvelope.pendingPredecessorDeletes, + obsoleteNames: obsoleteNames) + self.cancelPendingSnapshotDeletes( + CloudSyncSnapshotMigration.cancelledPersistedDeletes( + pendingDeletes: self.persistenceEnvelope.pendingSnapshotDeletes, + liveNames: Set(self.pendingSnapshots.map(\.recordName)))) + self.hasReconciledLiveSnapshots = true + } + self.requeuePendingSnapshotDeletes() + var stillPending: [AccountSnapshotSyncPayload] = [] + for payload in self.pendingSnapshots { + let hash = try CanonicalSyncJSON.hash(payload) + if self.skippedTerminalReplacementHashes[payload.recordName] == hash { + continue + } + let predecessors = CloudSyncSnapshotMigration.predecessorNames( + for: payload, + obsoleteNames: obsoleteNames) + // Replace, don't union: a later live email-keyed snapshot must not stay queued + // for delete after the slot-keyed save is confirmed. + CloudSyncSnapshotMigration.assigningPredecessors( + predecessors, + to: payload.recordName, + pending: &self.persistenceEnvelope.pendingPredecessorDeletes) + // Hashes are recorded only after CloudKit confirms a save. + let alreadyPublished = self.lastSnapshotHashes[payload.recordName] == hash + if alreadyPublished { + if !predecessors.isEmpty { + await self.finishConfirmedSnapshotMigrations( + savedRecordNames: [payload.recordName], + syncEngine: engine) + } + continue + } + if CloudSyncSnapshotMigration.hasInFlightSave( + recordName: payload.recordName, + pendingSaveHashes: self.pendingSaveHashes) + { + self.persistenceEnvelope.fleetSnapshots[payload.recordName] = payload + stillPending.append(payload) + continue + } + let recordID = self.recordID(named: payload.recordName) + let record = self.record(type: .accountSnapshot, id: recordID) + record["schemaVersion"] = payload.schemaVersion as CKRecordValue + record["provider"] = payload.provider.rawValue as CKRecordValue + record["deviceID"] = payload.deviceID as CKRecordValue + record["accountKey"] = payload.accountKey as CKRecordValue + record["fetchedAt"] = payload.fetchedAt as CKRecordValue + record.encryptedValues["displayLabel"] = payload.displayLabel as CKRecordValue + record.encryptedValues["usagePayload"] = try CanonicalSyncJSON.string(payload.usage) as CKRecordValue + self.desiredRecords[recordID] = record + self.persistenceEnvelope.fleetSnapshots[payload.recordName] = payload + self.skippedTerminalReplacementHashes.removeValue(forKey: payload.recordName) + self.pendingSaveHashes[payload.recordName] = hash + engine.state.add(pendingRecordZoneChanges: [.saveRecord(recordID)]) + } + self.pendingSnapshots = stillPending + self.lastSnapshotPushAt = Date() + self.persistEnvelope() + } catch { + await self.record(error: error) + } + } +} diff --git a/Sources/CodexBar/Sync/CloudSyncPersistence.swift b/Sources/CodexBar/Sync/CloudSyncPersistence.swift index 4338429b71..bc0324bdaa 100644 --- a/Sources/CodexBar/Sync/CloudSyncPersistence.swift +++ b/Sources/CodexBar/Sync/CloudSyncPersistence.swift @@ -19,6 +19,8 @@ struct CloudSyncPersistence: Sendable { var preferencesDirty: Bool var fleetDevices: [String: DeviceSyncPayload] var fleetSnapshots: [String: AccountSnapshotSyncPayload] + var pendingSnapshotDeletes: Set + var pendingPredecessorDeletes: [String: Set] init( stateSerialization: CKSyncEngine.State.Serialization?, @@ -28,7 +30,9 @@ struct CloudSyncPersistence: Sendable { dirtyProviders: Set = [], preferencesDirty: Bool = false, fleetDevices: [String: DeviceSyncPayload] = [:], - fleetSnapshots: [String: AccountSnapshotSyncPayload] = [:]) + fleetSnapshots: [String: AccountSnapshotSyncPayload] = [:], + pendingSnapshotDeletes: Set = [], + pendingPredecessorDeletes: [String: Set] = [:]) { self.stateSerialization = stateSerialization self.encodedSystemFields = encodedSystemFields @@ -38,6 +42,8 @@ struct CloudSyncPersistence: Sendable { self.preferencesDirty = preferencesDirty self.fleetDevices = fleetDevices self.fleetSnapshots = fleetSnapshots + self.pendingSnapshotDeletes = pendingSnapshotDeletes + self.pendingPredecessorDeletes = pendingPredecessorDeletes } private enum CodingKeys: String, CodingKey { @@ -49,6 +55,8 @@ struct CloudSyncPersistence: Sendable { case preferencesDirty case fleetDevices case fleetSnapshots + case pendingSnapshotDeletes + case pendingPredecessorDeletes } init(from decoder: any Decoder) throws { @@ -77,6 +85,12 @@ struct CloudSyncPersistence: Sendable { self.fleetSnapshots = try container.decodeIfPresent( [String: AccountSnapshotSyncPayload].self, forKey: .fleetSnapshots) ?? [:] + self.pendingSnapshotDeletes = try container.decodeIfPresent( + Set.self, + forKey: .pendingSnapshotDeletes) ?? [] + self.pendingPredecessorDeletes = try container.decodeIfPresent( + [String: Set].self, + forKey: .pendingPredecessorDeletes) ?? [:] } } diff --git a/Sources/CodexBarCore/Sync/SyncModels.swift b/Sources/CodexBarCore/Sync/SyncModels.swift index b814d435cb..5025950f27 100644 --- a/Sources/CodexBarCore/Sync/SyncModels.swift +++ b/Sources/CodexBarCore/Sync/SyncModels.swift @@ -225,6 +225,43 @@ public struct AccountSnapshotSyncPayload: Codable, Sendable { } return CanonicalSyncJSON.hash(data: Data(identity.lowercased().utf8)) } + + /// CloudKit record IDs cannot be renamed. Claude Swap snapshots keyed by + /// `claude-swap:` replace a same-device email-keyed record; that + /// predecessor is deleted only after the slot-keyed replacement is saved. + /// Other providers must not classify an email-to-durable-ID change as obsolete. + public func emailKeyedPredecessorRecordName() -> String? { + // Provider-specific by design: only Claude Swap slot keys retire leftover email-keyed CloudKit snapshots. + guard self.provider == .claude, + self.usage.identity?.loginMethod == ClaudeSwapAccountProjection.sourceLabel + else { return nil } + let accountID = self.usage.identity?.accountID?.trimmingCharacters(in: .whitespacesAndNewlines) ?? "" + guard accountID.hasPrefix("\(ClaudeSwapAccountProjection.sourceName):") else { return nil } + let emailKey = Self.accountKey(for: self.usage.identity?.accountEmail) + guard emailKey != "default", + Self.accountKey(for: accountID) == self.accountKey, + emailKey != self.accountKey + else { + return nil + } + return "snap-\(self.provider.rawValue)-\(emailKey)-\(self.deviceID)" + } + + public static func obsoleteEmailKeyedRecordNames( + liveSnapshots: [AccountSnapshotSyncPayload], + knownRecordNames: Set) -> Set + { + let liveNames = Set(liveSnapshots.map(\.recordName)) + var obsolete: Set = [] + for snapshot in liveSnapshots { + guard let predecessor = snapshot.emailKeyedPredecessorRecordName(), + !liveNames.contains(predecessor), + knownRecordNames.contains(predecessor) + else { continue } + obsolete.insert(predecessor) + } + return obsolete + } } public struct SyncedPreferences: Codable, Sendable { diff --git a/Tests/CodexBarTests/CloudSyncSettingsTests.swift b/Tests/CodexBarTests/CloudSyncSettingsTests.swift index 6b090461fa..9f76b604c3 100644 --- a/Tests/CodexBarTests/CloudSyncSettingsTests.swift +++ b/Tests/CodexBarTests/CloudSyncSettingsTests.swift @@ -162,6 +162,38 @@ struct CloudSyncSettingsTests { #expect(envelope.dirtyProviders.isEmpty) #expect(!envelope.preferencesDirty) + #expect(envelope.pendingSnapshotDeletes.isEmpty) + } + + @Test + func `pending snapshot deletes survive persistence round trip`() throws { + let directory = FileManager.default.temporaryDirectory + .appendingPathComponent("CloudSyncPendingDeletesTests-\(UUID().uuidString)", isDirectory: true) + defer { try? FileManager.default.removeItem(at: directory) } + let fileURL = directory.appendingPathComponent("engine-state.json") + let persistence = CloudSyncPersistence(fileURL: fileURL) + var envelope = CloudSyncPersistence.Envelope(stateSerialization: nil, encodedSystemFields: [:]) + envelope.pendingSnapshotDeletes = ["snap-claude-old-device-id"] + try persistence.save(envelope) + + #expect(persistence.load().pendingSnapshotDeletes == ["snap-claude-old-device-id"]) + } + + @Test + func `pending predecessor deletes survive persistence round trip`() throws { + let directory = FileManager.default.temporaryDirectory + .appendingPathComponent("CloudSyncPendingPredecessorsTests-\(UUID().uuidString)", isDirectory: true) + defer { try? FileManager.default.removeItem(at: directory) } + let fileURL = directory.appendingPathComponent("engine-state.json") + let persistence = CloudSyncPersistence(fileURL: fileURL) + var envelope = CloudSyncPersistence.Envelope(stateSerialization: nil, encodedSystemFields: [:]) + envelope.pendingPredecessorDeletes = ["snap-claude-slot-device-id": ["snap-claude-old-device-id"]] + try persistence.save(envelope) + + #expect( + persistence.load().pendingPredecessorDeletes["snap-claude-slot-device-id"] == [ + "snap-claude-old-device-id", + ]) } @Test diff --git a/Tests/CodexBarTests/SyncModelTests.swift b/Tests/CodexBarTests/SyncModelTests.swift index 1ceaeb96e9..83a6a17422 100644 --- a/Tests/CodexBarTests/SyncModelTests.swift +++ b/Tests/CodexBarTests/SyncModelTests.swift @@ -176,6 +176,83 @@ struct SyncModelTests { #expect(decoded.fetchedAt == usage.updatedAt) } + @Test + func `slot keyed snapshot names the leftover email keyed CloudKit record`() { + let payload = Self.claudeSnapshot(accountID: "claude-swap:2", email: "Owner@Example.com") + let emailKey = AccountSnapshotSyncPayload.accountKey(for: "owner@example.com") + + #expect(payload.emailKeyedPredecessorRecordName() == "snap-claude-\(emailKey)-device-id") + #expect(payload.recordName != payload.emailKeyedPredecessorRecordName()) + } + + @Test + func `email keyed snapshot has no CloudKit predecessor`() { + let payload = Self.claudeSnapshot(accountID: "owner@example.com", email: "owner@example.com") + + #expect(payload.emailKeyedPredecessorRecordName() == nil) + } + + @Test + func `obsolete email keyed names skip live records and unknown CloudKit keys`() throws { + let slot = Self.claudeSnapshot(accountID: "claude-swap:2", email: "owner@example.com") + let oauth = Self.claudeSnapshot(accountID: "owner@example.com", email: "owner@example.com") + let predecessor = try #require(slot.emailKeyedPredecessorRecordName()) + + #expect( + AccountSnapshotSyncPayload.obsoleteEmailKeyedRecordNames( + liveSnapshots: [slot], + knownRecordNames: [predecessor]) == [predecessor]) + #expect( + AccountSnapshotSyncPayload.obsoleteEmailKeyedRecordNames( + liveSnapshots: [slot, oauth], + knownRecordNames: [predecessor]).isEmpty) + #expect( + AccountSnapshotSyncPayload.obsoleteEmailKeyedRecordNames( + liveSnapshots: [slot], + knownRecordNames: []).isEmpty) + } + + @Test + func `duplicate swap slots sharing a mailbox retire one email keyed record`() throws { + let first = Self.claudeSnapshot(accountID: "claude-swap:1", email: "shared@example.com") + let second = Self.claudeSnapshot(accountID: "claude-swap:2", email: "shared@example.com") + let predecessor = try #require(first.emailKeyedPredecessorRecordName()) + + #expect(second.emailKeyedPredecessorRecordName() == predecessor) + #expect( + AccountSnapshotSyncPayload.obsoleteEmailKeyedRecordNames( + liveSnapshots: [first, second], + knownRecordNames: [predecessor]) == [predecessor]) + } + + @Test + func `non Claude snapshot does not name an email keyed CloudKit predecessor`() { + let payload = Self.snapshot( + provider: .codex, + loginMethod: "pro", + accountID: "user-workspace-1", + email: "owner@example.com") + let emailKey = AccountSnapshotSyncPayload.accountKey(for: "owner@example.com") + let leftover = "snap-codex-\(emailKey)-device-id" + + #expect(payload.emailKeyedPredecessorRecordName() == nil) + #expect( + AccountSnapshotSyncPayload.obsoleteEmailKeyedRecordNames( + liveSnapshots: [payload], + knownRecordNames: [leftover]).isEmpty) + } + + @Test + func `claude subscription snapshot does not name an email keyed CloudKit predecessor`() { + let payload = Self.snapshot( + provider: .claude, + loginMethod: "Claude.ai", + accountID: "user_abc", + email: "owner@example.com") + + #expect(payload.emailKeyedPredecessorRecordName() == nil) + } + @Test func `account snapshot ignores retired provider payload keys`() throws { let legacy = #""" @@ -219,6 +296,38 @@ struct SyncModelTests { #expect(decoded.usage.details.isEmpty) #expect(decoded.usage.updatedAt == decoded.fetchedAt) } + + private static func claudeSnapshot(accountID: String, email: String) -> AccountSnapshotSyncPayload { + self.snapshot( + provider: .claude, + loginMethod: ClaudeSwapAccountProjection.sourceLabel, + accountID: accountID, + email: email) + } + + private static func snapshot( + provider: ProviderInstanceID, + loginMethod: String, + accountID: String, + email: String) -> AccountSnapshotSyncPayload + { + let usage = UsageSnapshot( + primary: nil, + secondary: nil, + updatedAt: Date(timeIntervalSince1970: 100), + identity: ProviderIdentitySnapshot( + providerID: provider, + accountEmail: email, + accountOrganization: nil, + loginMethod: loginMethod, + accountID: accountID)) + return AccountSnapshotSyncPayload( + provider: provider, + deviceID: "device-id", + accountIdentity: accountID, + displayLabel: email, + usage: usage) + } } import CloudKit @@ -249,3 +358,365 @@ struct CloudSyncRecordRebaseTests { #expect(rebased.encryptedValues["cookieHeader"] as? String == "cookie=1") } } + +struct CloudSyncSnapshotMigrationDeleteRetryTests { + @Test + func `terminal CloudKit delete errors are reported once and not retried`() { + let zoneID = CloudSyncEngine.zoneID + func recordID(_ name: String) -> CKRecord.ID { + CKRecord.ID(recordName: name, zoneID: zoneID) + } + + let failures: [CKRecord.ID: CKError] = [ + recordID("unknown"): Self.cloudKitError(.unknownItem), + recordID("denied"): Self.cloudKitError(.permissionFailure), + recordID("unauth"): Self.cloudKitError(.notAuthenticated), + recordID("invalid"): Self.cloudKitError(.invalidArguments), + recordID("network"): Self.cloudKitError(.networkFailure), + recordID("quota"): Self.cloudKitError(.quotaExceeded, retryAfter: 30), + recordID("lost"): Self.cloudKitError(.serverResponseLost), + recordID("temp"): Self.cloudKitError(.accountTemporarilyUnavailable), + ] + + let retryable = Set(CloudSyncSnapshotMigration.retryableFailedDeletes(failures).map(\.recordName)) + let reported = Set(CloudSyncSnapshotMigration.reportableFailedDeletes(failures).map(\.code)) + + #expect(retryable == ["network", "quota", "lost", "temp"]) + #expect(reported == [.permissionFailure, .notAuthenticated, .invalidArguments]) + #expect(CloudSyncSnapshotMigration.retryDelay(for: Self.cloudKitError(.unknownItem)) == nil) + #expect(CloudSyncSnapshotMigration.retryDelay(for: Self.cloudKitError(.networkFailure)) == 1) + #expect(CloudSyncSnapshotMigration.retryDelay(for: Self.cloudKitError(.quotaExceeded, retryAfter: 30)) == 30) + #expect(CloudSyncSnapshotMigration.retryDelay(for: Self.cloudKitError(.serverResponseLost)) == 1) + #expect(CloudSyncSnapshotMigration.retryDelay(for: Self.cloudKitError(.accountTemporarilyUnavailable)) == 1) + #expect(CloudSyncSnapshotMigration.retryDelay(for: Self.cloudKitError(.permissionFailure)) == nil) + #expect( + CloudSyncSnapshotMigration.finishedFailedDeleteNames(failures) == [ + "unknown", "denied", "unauth", "invalid", + ]) + } + + @Test + func `transient delete failures do not retry live snapshots`() { + let zoneID = CloudSyncEngine.zoneID + func recordID(_ name: String) -> CKRecord.ID { + CKRecord.ID(recordName: name, zoneID: zoneID) + } + + let failures: [CKRecord.ID: CKError] = [ + recordID("snap-live"): Self.cloudKitError(.networkFailure), + recordID("snap-obsolete"): Self.cloudKitError(.networkFailure), + recordID("snap-denied"): Self.cloudKitError(.permissionFailure), + ] + + let liveNames = CloudSyncSnapshotMigration.liveSnapshotRecordNames( + pendingRecordNames: ["snap-pending"], + storedRecordNames: ["snap-confirmed"]) + #expect(liveNames == ["snap-pending", "snap-confirmed"]) + #expect(!liveNames.contains("snap-remote-cache")) + + let retryable = Set( + CloudSyncSnapshotMigration.retryableFailedDeletes( + failures, + liveNames: ["snap-live"]).map(\.recordName)) + #expect(retryable == ["snap-obsolete"]) + } + + private static func cloudKitError(_ code: CKError.Code, retryAfter: TimeInterval? = nil) -> CKError { + var userInfo: [String: Any] = [:] + if let retryAfter { + userInfo[CKErrorRetryAfterKey] = NSNumber(value: retryAfter) + } + let nsError = NSError(domain: CKErrorDomain, code: code.rawValue, userInfo: userInfo) + return CKError(_nsError: nsError) + } +} + +struct CloudSyncSnapshotMigrationSaveThenDeleteTests { + @Test + func `predecessor deletes wait until the replacement record is saved`() throws { + let slot = Self.claudeSnapshot(accountID: "claude-swap:2", email: "owner@example.com") + let predecessor = try #require(slot.emailKeyedPredecessorRecordName()) + let obsolete: Set = [predecessor] + + #expect(CloudSyncSnapshotMigration.predecessorNames(for: slot, obsoleteNames: obsolete) == [predecessor]) + #expect(CloudSyncSnapshotMigration.predecessorNames(for: slot, obsoleteNames: []).isEmpty) + + var pending = [slot.recordName: Set([predecessor])] + #expect(CloudSyncSnapshotMigration.takeDeletes(forSavedRecordNames: [], pending: &pending).isEmpty) + #expect(pending[slot.recordName] == [predecessor]) + #expect( + CloudSyncSnapshotMigration.takeDeletes( + forSavedRecordNames: [slot.recordName], + pending: &pending) == [predecessor]) + #expect(pending.isEmpty) + } + + @Test + func `shared predecessor waits for every slot replacement to save`() throws { + let first = Self.claudeSnapshot(accountID: "claude-swap:1", email: "shared@example.com") + let second = Self.claudeSnapshot(accountID: "claude-swap:2", email: "shared@example.com") + let predecessor = try #require(first.emailKeyedPredecessorRecordName()) + var pending = [ + first.recordName: Set([predecessor]), + second.recordName: Set([predecessor]), + ] + + #expect( + CloudSyncSnapshotMigration.takeDeletes( + forSavedRecordNames: [first.recordName], + pending: &pending).isEmpty) + #expect(pending[second.recordName] == [predecessor]) + #expect( + CloudSyncSnapshotMigration.takeDeletes( + forSavedRecordNames: [second.recordName], + pending: &pending) == [predecessor]) + #expect(pending.isEmpty) + } + + @Test + func `failed sibling replacement keeps a shared predecessor after a confirmed save`() throws { + let first = Self.claudeSnapshot(accountID: "claude-swap:1", email: "shared@example.com") + let second = Self.claudeSnapshot(accountID: "claude-swap:2", email: "shared@example.com") + let predecessor = try #require(first.emailKeyedPredecessorRecordName()) + var pending = [ + first.recordName: Set([predecessor]), + second.recordName: Set([predecessor]), + ] + + #expect( + CloudSyncSnapshotMigration.takeDeletes( + forSavedRecordNames: [first.recordName], + pending: &pending).isEmpty) + #expect(pending[second.recordName] == [predecessor]) + + let abandoned = CloudSyncSnapshotMigration.abandonedReplacementNames( + failures: [second.recordName: Self.cloudKitError(.permissionFailure)], + pendingReplacements: Set(pending.keys)) + #expect(abandoned == [second.recordName]) + for name in abandoned { + pending.removeValue(forKey: name) + } + #expect(pending.isEmpty) + } + + @Test + func `already published replacements still delete newly obsolete predecessors`() throws { + let slot = Self.claudeSnapshot(accountID: "claude-swap:2", email: "owner@example.com") + let predecessor = try #require(slot.emailKeyedPredecessorRecordName()) + var pending: [String: Set] = [:] + + CloudSyncSnapshotMigration.assigningPredecessors( + [predecessor], + to: slot.recordName, + pending: &pending) + #expect( + CloudSyncSnapshotMigration.takeDeletes( + forSavedRecordNames: [slot.recordName], + pending: &pending) == [predecessor]) + #expect(pending.isEmpty) + } + + @Test + func `pending predecessors drop names that are live again`() throws { + let slot = Self.claudeSnapshot(accountID: "claude-swap:2", email: "owner@example.com") + let predecessor = try #require(slot.emailKeyedPredecessorRecordName()) + var pending = [slot.recordName: Set([predecessor])] + + CloudSyncSnapshotMigration.retainingObsoletePredecessors( + in: &pending, + obsoleteNames: []) + #expect(pending.isEmpty) + + pending = [slot.recordName: [predecessor, "snap-stale"]] + CloudSyncSnapshotMigration.assigningPredecessors( + [predecessor], + to: slot.recordName, + pending: &pending) + #expect(pending[slot.recordName] == [predecessor]) + + CloudSyncSnapshotMigration.assigningPredecessors([], to: slot.recordName, pending: &pending) + #expect(pending.isEmpty) + } + + @Test + func `restored predecessor deletes wait until live snapshots reconcile`() throws { + let slot = Self.claudeSnapshot(accountID: "claude-swap:2", email: "owner@example.com") + let predecessor = try #require(slot.emailKeyedPredecessorRecordName()) + var pending = [slot.recordName: Set([predecessor])] + + #expect( + CloudSyncSnapshotMigration.takeDeletes( + forSavedRecordNames: [slot.recordName], + pending: &pending, + afterLiveSnapshotReconciliation: false).isEmpty) + #expect(pending[slot.recordName] == [predecessor]) + #expect( + CloudSyncSnapshotMigration.takeDeletes( + forSavedRecordNames: [slot.recordName], + pending: &pending, + afterLiveSnapshotReconciliation: true) == [predecessor]) + #expect(pending.isEmpty) + } + + @Test + func `persisted deletes are cancelled when the predecessor is live again`() throws { + let slot = Self.claudeSnapshot(accountID: "claude-swap:2", email: "owner@example.com") + let predecessor = try #require(slot.emailKeyedPredecessorRecordName()) + + #expect( + CloudSyncSnapshotMigration.cancelledPersistedDeletes( + pendingDeletes: [predecessor, "snap-claude-stale-device-id"], + liveNames: [predecessor]) == [predecessor]) + #expect( + CloudSyncSnapshotMigration.cancelledPersistedDeletes( + pendingDeletes: ["snap-claude-stale-device-id"], + liveNames: [predecessor]).isEmpty) + #expect( + CloudSyncSnapshotMigration.pendingDeletesToRequeue( + pendingDeletes: [predecessor, "snap-claude-stale-device-id"], + liveNames: [predecessor]) == ["snap-claude-stale-device-id"]) + } + + @Test + func `terminal replacement save failures stop retrying the same payload`() { + let slot = "snap-claude-slot-device-id" + let failures = [ + slot: Self.cloudKitError(.permissionFailure), + "snap-other": Self.cloudKitError(.networkFailure), + "snap-quota": Self.cloudKitError(.quotaExceeded, retryAfter: 12), + ] + + #expect( + CloudSyncSnapshotMigration.abandonedReplacementNames( + failures: failures, + pendingReplacements: [slot, "snap-quota"]) == [slot]) + #expect( + CloudSyncSnapshotMigration.abandonedReplacementNames( + failures: [slot: Self.cloudKitError(.networkFailure)], + pendingReplacements: [slot]).isEmpty) + } + + @Test + func `confirmed save hashes replace the previously stored version`() { + var pending = ["snap-a": "hash-new"] + var last = ["snap-a": "hash-old", "snap-b": "hash-other"] + + CloudSyncSnapshotMigration.applyConfirmedSaveHashes( + savedRecordNames: ["snap-a"], + pendingSaveHashes: &pending, + lastSnapshotHashes: &last) + + #expect(last["snap-a"] == "hash-new") + #expect(last["snap-b"] == "hash-other") + #expect(pending.isEmpty) + } + + @Test + func `terminal save skips apply without a predecessor mapping`() { + var pending = ["snap-new": "hash-sent"] + var skipped: [String: String] = [:] + + CloudSyncSnapshotMigration.applyTerminalSaveSkip( + recordName: "snap-new", + error: Self.cloudKitError(.permissionFailure), + pendingSaveHashes: &pending, + skippedTerminalReplacementHashes: &skipped) + #expect(skipped["snap-new"] == "hash-sent") + #expect(pending.isEmpty) + + pending = ["snap-new": "hash-sent"] + skipped = [:] + CloudSyncSnapshotMigration.applyTerminalSaveSkip( + recordName: "snap-new", + error: Self.cloudKitError(.networkFailure), + pendingSaveHashes: &pending, + skippedTerminalReplacementHashes: &skipped) + #expect(skipped.isEmpty) + #expect(pending["snap-new"] == "hash-sent") + } + + @Test + func `newer in-flight snapshot payloads stay pending until confirmed`() throws { + let older = Self.claudeSnapshot(accountID: "claude-swap:2", email: "owner@example.com") + let newer = AccountSnapshotSyncPayload( + provider: older.provider, + deviceID: older.deviceID, + accountKey: older.accountKey, + fetchedAt: Date(timeIntervalSince1970: 200), + displayLabel: older.displayLabel, + usage: UsageSnapshot( + primary: nil, + secondary: nil, + updatedAt: Date(timeIntervalSince1970: 200), + identity: older.usage.identity), + schemaVersion: older.schemaVersion) + let olderHash = try CanonicalSyncJSON.hash(older) + let newerHash = try CanonicalSyncJSON.hash(newer) + #expect(olderHash != newerHash) + + let merged = CloudSyncSnapshotMigration.mergingPendingSnapshots([], with: [newer]) + #expect(merged.map(\.recordName) == [newer.recordName]) + #expect( + CloudSyncSnapshotMigration.mergingPendingSnapshots([newer], with: [older]).first?.fetchedAt + == newer.fetchedAt) + + var last = [newer.recordName: olderHash] + let unpublished = CloudSyncSnapshotMigration.unpublishedFleetSnapshots( + savedRecordNames: [newer.recordName], + fleetSnapshots: [newer.recordName: newer], + lastSnapshotHashes: last) + #expect(unpublished.map(\.recordName) == [newer.recordName]) + + last[newer.recordName] = newerHash + #expect( + CloudSyncSnapshotMigration.unpublishedFleetSnapshots( + savedRecordNames: [newer.recordName], + fleetSnapshots: [newer.recordName: newer], + lastSnapshotHashes: last).isEmpty) + } + + @Test + func `delayed delete retries do not resume on a replacement sync engine`() { + let original = NSObject() + #expect( + CloudSyncSnapshotMigration.shouldResumeDelayedRetry( + originatingEngine: ObjectIdentifier(original), + currentEngine: ObjectIdentifier(original))) + #expect( + !CloudSyncSnapshotMigration.shouldResumeDelayedRetry( + originatingEngine: ObjectIdentifier(original), + currentEngine: ObjectIdentifier(NSObject()))) + #expect( + !CloudSyncSnapshotMigration.shouldResumeDelayedRetry( + originatingEngine: ObjectIdentifier(original), + currentEngine: nil)) + } + + private static func cloudKitError(_ code: CKError.Code, retryAfter: TimeInterval? = nil) -> CKError { + var userInfo: [String: Any] = [:] + if let retryAfter { + userInfo[CKErrorRetryAfterKey] = NSNumber(value: retryAfter) + } + let nsError = NSError(domain: CKErrorDomain, code: code.rawValue, userInfo: userInfo) + return CKError(_nsError: nsError) + } + + private static func claudeSnapshot(accountID: String, email: String) -> AccountSnapshotSyncPayload { + let usage = UsageSnapshot( + primary: nil, + secondary: nil, + updatedAt: Date(timeIntervalSince1970: 100), + identity: ProviderIdentitySnapshot( + providerID: .claude, + accountEmail: email, + accountOrganization: nil, + loginMethod: "claude-swap", + accountID: accountID)) + return AccountSnapshotSyncPayload( + provider: .claude, + deviceID: "device-id", + accountIdentity: accountID, + displayLabel: email, + usage: usage) + } +}