GameArchiver.swift (38425B)
1 import CloudKit 2 import CoreData 3 import Foundation 4 5 enum CompletedMetadataPageWalker { 6 struct Result<Item, Cursor> { 7 let selected: [Item] 8 let buffered: [Item] 9 let cursor: Cursor? 10 } 11 12 /// Walks ordered metadata pages until the first item outside the initial 13 /// date window. Passing the returned cursor into every subsequent fetch is 14 /// the invariant that prevents page one from being re-read indefinitely. 15 @MainActor 16 static func recent<Item, Cursor>( 17 cutoff: Date, 18 completedAt: (Item) -> Date, 19 fetch: @MainActor (Cursor?) async throws -> ( 20 records: [Item], 21 cursor: Cursor? 22 ) 23 ) async rethrows -> Result<Item, Cursor> { 24 var cursor: Cursor? 25 var selected: [Item] = [] 26 var buffered: [Item] = [] 27 repeat { 28 let page = try await fetch(cursor) 29 cursor = page.cursor 30 if let firstOlder = page.records.firstIndex(where: { 31 completedAt($0) < cutoff 32 }) { 33 selected.append(contentsOf: page.records[..<firstOlder]) 34 buffered.append(contentsOf: page.records[firstOlder...]) 35 break 36 } 37 selected.append(contentsOf: page.records) 38 } while cursor != nil 39 return Result(selected: selected, buffered: buffered, cursor: cursor) 40 } 41 } 42 43 /// Compacts completed games into the account's private archive zone and retires 44 /// their multi-record live zones once every participant has archived, or once 45 /// the hard retention deadline expires. 46 @MainActor 47 final class GameArchiver { 48 nonisolated static let archiveRetryWindow: TimeInterval = 14 * 24 * 60 * 60 49 nonisolated static let completedPageSize = 7 50 51 struct CompletedPage: Sendable { 52 let oldestCompletedAt: Date? 53 let hasMore: Bool 54 } 55 56 private enum CompletedSource { 57 case chronicle(CKRecord.ID) 58 case game( 59 gameID: UUID, 60 zoneID: CKRecordZone.ID, 61 scope: DatabaseScope 62 ) 63 } 64 65 private struct CompletedMetadata { 66 let originalGameID: UUID 67 let completedAt: Date 68 let source: CompletedSource 69 70 var isLiveGame: Bool { 71 if case .game = source { return true } 72 return false 73 } 74 } 75 76 private struct LocalGame { 77 let snapshot: Archive.Snapshot 78 let databaseScope: DatabaseScope 79 let liveZoneID: CKRecordZone.ID 80 let isShared: Bool 81 /// nil means the owner has not received an authoritative CKShare roster 82 /// yet. An empty set is a known solo/no-participant game. 83 let acceptedParticipants: Set<String>? 84 let archiveAcknowledgedAt: Date? 85 } 86 87 private struct StoredArchive { 88 let payload: Archive.Payload 89 let isLegacy: Bool 90 } 91 92 private let container: CKContainer 93 private let persistence: PersistenceController 94 private let syncEngine: SyncEngine 95 private let syncMonitor: SyncMonitor? 96 private let eventLog: EventLog? 97 private let localIdentity: () -> (authorID: String, playerName: String)? 98 private let localDefaults: UserDefaults 99 private let ubiquitousStore: NSUbiquitousKeyValueStore? 100 private var ensuredArchiveZone = false 101 private var chronicleCursor: CKQueryOperation.Cursor? 102 private var bufferedCompleted: [CompletedMetadata] = [] 103 104 init( 105 container: CKContainer, 106 persistence: PersistenceController, 107 syncEngine: SyncEngine, 108 syncMonitor: SyncMonitor? = nil, 109 eventLog: EventLog? = nil, 110 localIdentity: @escaping () -> (authorID: String, playerName: String)? = { nil }, 111 localDefaults: UserDefaults = .standard, 112 ubiquitousStore: NSUbiquitousKeyValueStore? = .default 113 ) { 114 self.container = container 115 self.persistence = persistence 116 self.syncEngine = syncEngine 117 self.syncMonitor = syncMonitor 118 self.eventLog = eventLog 119 self.localIdentity = localIdentity 120 self.localDefaults = localDefaults 121 self.ubiquitousStore = ubiquitousStore 122 } 123 124 // MARK: - Reconciliation 125 126 /// Immediate completion path. It writes/refreshes the archive and emits a 127 /// participant acknowledgement, but leaves zone retirement to the cold- 128 /// launch reconciliation path so an open Success Panel is never replaced 129 /// underneath the user. 130 func archiveIfNeeded(gameID: UUID) async { 131 guard let graceStart = await accountGraceStart() else { return } 132 _ = await reconcileArchive(gameID: gameID, graceStart: graceStart) 133 } 134 135 /// Cold-launch backstop for new completions, migrations, acknowledgements, 136 /// and owner-side zone retirement. 137 func reconcileUnarchived() async { 138 guard let graceStart = await accountGraceStart() else { return } 139 await migrateMaterializedLegacyArchives() 140 141 let ctx = persistence.container.newBackgroundContext() 142 let ids: [UUID] = await ctx.perform { 143 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 144 req.predicate = NSPredicate( 145 format: "completedAt != nil AND isAccessRevoked == NO AND ckRecordName BEGINSWITH %@", 146 "game-" 147 ) 148 return ((try? ctx.fetch(req)) ?? []).compactMap(\.id) 149 } 150 for id in ids { 151 guard let result = await reconcileArchive(gameID: id, graceStart: graceStart), 152 result.local.databaseScope == .private 153 else { continue } 154 await retireOwnedGameIfEligible( 155 result.local, 156 snapshot: result.snapshot, 157 archiveComplete: result.complete, 158 graceStart: graceStart 159 ) 160 } 161 } 162 163 /// Whether the Chronicle in CloudKit still says what this pass concluded. 164 /// 165 /// The journal comparison is deliberately limited to `.available`: only a 166 /// complete Chronicle embeds journals, so a provisional one is always 167 /// stored with an empty journal. Comparing that against the merged local 168 /// history would report "missing" on every pass and re-upload the whole 169 /// payload — a fetch, a replay query, and a compressed asset — for the 170 /// entire multi-day waiting window, without ever changing what's stored. 171 nonisolated static func chronicleNeedsWrite( 172 stored: ( 173 isLegacy: Bool, 174 formatVersion: Int, 175 replayState: Archive.ReplayState, 176 journalKeys: Set<JournalDeviceKey> 177 )?, 178 replayState: Archive.ReplayState, 179 presentJournalKeys: Set<JournalDeviceKey> 180 ) -> Bool { 181 guard let stored else { return true } 182 if stored.isLegacy { return true } 183 if stored.formatVersion != Archive.currentPayloadFormatVersion { return true } 184 if stored.replayState != replayState { return true } 185 return replayState == .available 186 && !presentJournalKeys.isSubset(of: stored.journalKeys) 187 } 188 189 nonisolated static func hasArchiveRetryExpired( 190 completedAt: Date, 191 graceStart: Date = .distantPast, 192 now: Date = Date() 193 ) -> Bool { 194 now >= max(completedAt, graceStart).addingTimeInterval(archiveRetryWindow) 195 } 196 197 private func localGame(gameID: UUID) async -> LocalGame? { 198 let ctx = persistence.container.newBackgroundContext() 199 return await ctx.perform { 200 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 201 req.predicate = NSPredicate(format: "id == %@", gameID as CVarArg) 202 req.fetchLimit = 1 203 guard let entity = try? ctx.fetch(req).first, 204 entity.completedAt != nil, 205 !entity.isAccessRevoked, 206 entity.ckRecordName?.hasPrefix("game-") == true, 207 let snapshot = Archive.snapshot(forGameID: gameID, in: ctx) 208 else { return nil } 209 210 let scope = DatabaseScope(entityValue: entity.databaseScope) 211 let isShared = entity.ckShareRecordName != nil || scope == .shared 212 let participants: Set<String>? 213 if !isShared { 214 participants = [] 215 if entity.archiveParticipants == nil { entity.archiveParticipants = "" } 216 } else if let encoded = entity.archiveParticipants { 217 participants = Set(encoded.split(separator: ",").map(String.init)) 218 } else if let encoded = entity.shareParticipants { 219 entity.archiveParticipants = encoded 220 participants = Set(encoded.split(separator: ",").map(String.init)) 221 } else { 222 participants = nil 223 } 224 if ctx.hasChanges { try? ctx.save() } 225 return LocalGame( 226 snapshot: snapshot, 227 databaseScope: scope, 228 liveZoneID: CKRecordZone.ID( 229 zoneName: entity.ckZoneName ?? "game-\(gameID.uuidString)", 230 ownerName: entity.ckZoneOwnerName ?? CKCurrentUserDefaultName 231 ), 232 isShared: isShared, 233 acceptedParticipants: participants, 234 archiveAcknowledgedAt: entity.archiveAcknowledgedAt 235 ) 236 } 237 } 238 239 private func reconcileArchive( 240 gameID: UUID, 241 graceStart: Date 242 ) async -> (local: LocalGame, snapshot: Archive.Snapshot, complete: Bool)? { 243 guard let local = await localGame(gameID: gameID) else { return nil } 244 245 let stored = await fetchArchive(originalGameID: gameID) 246 // A complete compact Chronicle is authoritative: it was written only 247 // after every expected device journal was present. Re-reading the live 248 // replay cannot improve it and makes overlapping migration backstops 249 // download the same Moves and Journal assets repeatedly. 250 let storedComplete = stored.map { 251 !$0.isLegacy && $0.payload.replayAvailable 252 } ?? false 253 let canSkipReplayFetch = storedComplete 254 && stored?.payload.formatVersion == Archive.currentPayloadFormatVersion 255 let fetch = canSkipReplayFetch 256 ? nil 257 : try? await syncEngine.fetchReplay(forGameID: gameID) 258 var snapshot = local.snapshot 259 if let stored { snapshot = Archive.merging(snapshot, peerJournals: stored.payload.journal) } 260 if let fetch { snapshot = Archive.merging(snapshot, peerJournals: fetch.journals) } 261 262 let present = Set(snapshot.journal.map(\.key)) 263 let fetchedMissing = fetch.map { 264 $0.expectedDevices.subtracting(present).count 265 } 266 let fetchedComplete = fetchedMissing == 0 267 // Only a compact payload's explicit flag is authoritative. Legacy 268 // archives used archivedAt for both completeness and a timed best-effort 269 // fallback, so they must be checked against the still-live zone again. 270 let complete = storedComplete || fetchedComplete 271 let replayState: Archive.ReplayState 272 if complete { 273 replayState = .available 274 } else if case .waiting(let storedMissing) = stored?.payload.replayState { 275 replayState = .waiting(missing: fetchedMissing ?? storedMissing) 276 } else { 277 // A failed first fetch cannot determine the exact count yet, but it 278 // is still retryable rather than a retention fallback. 279 replayState = .waiting(missing: fetchedMissing ?? 1) 280 } 281 282 let needsWrite = Self.chronicleNeedsWrite( 283 stored: stored.map { 284 ( 285 isLegacy: $0.isLegacy, 286 formatVersion: $0.payload.formatVersion, 287 replayState: $0.payload.replayState, 288 journalKeys: Set($0.payload.journal.map(\.key)) 289 ) 290 }, 291 replayState: replayState, 292 presentJournalKeys: present 293 ) 294 if needsWrite { 295 guard await write(snapshot, replayState: replayState) else { return nil } 296 } 297 298 if stored?.isLegacy == true { 299 await syncEngine.enqueueDeleteLegacyArchiveZone( 300 Archive.legacyZoneID(forOriginalGameID: gameID) 301 ) 302 } 303 304 if complete { 305 await markArchived(originalGameID: gameID) 306 if local.databaseScope == .shared, local.archiveAcknowledgedAt == nil { 307 await acknowledgeChronicle(gameID: gameID) 308 } 309 } 310 return (local, snapshot, complete) 311 } 312 313 // MARK: - Retirement 314 315 private func retireOwnedGameIfEligible( 316 _ local: LocalGame, 317 snapshot: Archive.Snapshot, 318 archiveComplete: Bool, 319 graceStart: Date 320 ) async { 321 let expired = Self.hasArchiveRetryExpired( 322 completedAt: snapshot.completedAt, 323 graceStart: graceStart 324 ) 325 326 var quorum = false 327 if archiveComplete, let expected = local.acceptedParticipants { 328 if expected.isEmpty { 329 quorum = true 330 } else if let acknowledged = try? await syncEngine.fetchChronicleAcknowledgements( 331 forGameID: snapshot.originalGameID 332 ) { 333 quorum = expected.isSubset(of: acknowledged) 334 } 335 } 336 guard quorum || expired else { return } 337 338 let keepReplay = quorum && archiveComplete 339 if !keepReplay { 340 syncMonitor?.note( 341 "archive \(snapshot.originalGameID.uuidString.prefix(8)): " + 342 "retention deadline reached; retiring without replay" 343 ) 344 guard await write(snapshot, replayState: .unavailable) else { return } 345 } 346 guard await promoteOwnedBeforeRetirement( 347 snapshot, 348 replayState: keepReplay ? .available : .unavailable 349 ) else { 350 return 351 } 352 await syncEngine.enqueueRetireOwnedGameZone(local.liveZoneID) 353 } 354 355 private func promoteOwnedBeforeRetirement( 356 _ snapshot: Archive.Snapshot, 357 replayState: Archive.ReplayState 358 ) async -> Bool { 359 let ctx = persistence.container.newBackgroundContext() 360 return await ctx.perform { 361 let payload = Archive.payload(from: snapshot, replayState: replayState) 362 guard Archive.materialize(payload, in: ctx) != nil else { return false } 363 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 364 req.predicate = NSPredicate( 365 format: "id == %@", snapshot.originalGameID as CVarArg 366 ) 367 req.fetchLimit = 1 368 if let original = try? ctx.fetch(req).first { ctx.delete(original) } 369 do { 370 if ctx.hasChanges { try ctx.save() } 371 return true 372 } catch { 373 return false 374 } 375 } 376 } 377 378 // MARK: - Acknowledgement 379 380 private func acknowledgeChronicle(gameID: UUID) async { 381 guard let identity = localIdentity() else { return } 382 let enqueued = await syncEngine.enqueuePing( 383 kind: .chronicled, 384 gameID: gameID, 385 authorID: identity.authorID, 386 playerName: identity.playerName 387 ) 388 guard enqueued else { return } 389 let ctx = persistence.container.newBackgroundContext() 390 await ctx.perform { 391 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 392 req.predicate = NSPredicate(format: "id == %@", gameID as CVarArg) 393 req.fetchLimit = 1 394 guard let entity = try? ctx.fetch(req).first else { return } 395 entity.archiveAcknowledgedAt = Date() 396 try? ctx.save() 397 } 398 } 399 400 private func markArchived(originalGameID: UUID) async { 401 let ctx = persistence.container.newBackgroundContext() 402 await ctx.perform { 403 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 404 req.predicate = NSPredicate(format: "id == %@", originalGameID as CVarArg) 405 req.fetchLimit = 1 406 guard let entity = try? ctx.fetch(req).first else { return } 407 if entity.archivedAt == nil { entity.archivedAt = Date() } 408 entity.archiveGameID = Archive.archiveGameID(for: originalGameID) 409 try? ctx.save() 410 } 411 } 412 413 // MARK: - Restore / legacy migration 414 415 /// Rebuilds a completed game after another owner device retired its live 416 /// zone and the sync applier removed the local live row first. 417 func restoreRetired(gameID: UUID) async { 418 guard let stored = await fetchArchive(originalGameID: gameID) else { return } 419 let ctx = persistence.container.newBackgroundContext() 420 await ctx.perform { 421 _ = Archive.materialize(stored.payload, in: ctx) 422 if ctx.hasChanges { try? ctx.save() } 423 } 424 } 425 426 /// Promotes a participant's private archive when the owner retires the live 427 /// shared zone. Falls back to the local snapshot only when CloudKit has not 428 /// delivered the compact record yet. 429 func promoteRevoked(gameID: UUID) async { 430 let ctx = persistence.container.newBackgroundContext() 431 let local: Archive.Snapshot? = await ctx.perform { 432 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 433 req.predicate = NSPredicate(format: "id == %@", gameID as CVarArg) 434 req.fetchLimit = 1 435 guard (try? ctx.fetch(req).first)?.completedAt != nil else { return nil } 436 return Archive.snapshot(forGameID: gameID, in: ctx) 437 } 438 guard let local else { return } 439 let payload = await fetchArchive(originalGameID: gameID)?.payload 440 ?? Archive.payload(from: local, replayState: .unavailable) 441 let promoteCtx = persistence.container.newBackgroundContext() 442 await promoteCtx.perform { 443 guard Archive.materialize(payload, in: promoteCtx) != nil else { return } 444 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 445 req.predicate = NSPredicate(format: "id == %@", gameID as CVarArg) 446 req.fetchLimit = 1 447 if let original = try? promoteCtx.fetch(req).first { promoteCtx.delete(original) } 448 if promoteCtx.hasChanges { try? promoteCtx.save() } 449 } 450 } 451 452 private func migrateMaterializedLegacyArchives() async { 453 let ctx = persistence.container.newBackgroundContext() 454 let candidates: [(localID: UUID, originalID: UUID)] = await ctx.perform { 455 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 456 req.predicate = NSPredicate(format: "ckRecordName BEGINSWITH %@", "archive-") 457 return ((try? ctx.fetch(req)) ?? []).compactMap { entity in 458 guard let localID = entity.id, 459 entity.ckZoneName?.hasPrefix("archive-") == true, 460 let name = entity.ckRecordName, 461 let originalID = Archive.originalGameID(fromName: name) 462 else { return nil } 463 return (localID, originalID) 464 } 465 } 466 for candidate in candidates { 467 let snapshot: Archive.Snapshot? = await ctx.perform { 468 Archive.snapshot( 469 forGameID: candidate.localID, 470 originalGameID: candidate.originalID, 471 in: ctx 472 ) 473 } 474 guard let snapshot, 475 await write(snapshot, replayState: .available) 476 else { continue } 477 await ctx.perform { 478 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 479 req.predicate = NSPredicate(format: "id == %@", candidate.localID as CVarArg) 480 req.fetchLimit = 1 481 if let entity = try? ctx.fetch(req).first { 482 entity.ckRecordName = Archive.recordName( 483 forOriginalGameID: candidate.originalID 484 ) 485 entity.ckZoneName = Archive.zoneName 486 try? ctx.save() 487 } 488 } 489 await syncEngine.enqueueDeleteLegacyArchiveZone( 490 Archive.legacyZoneID(forOriginalGameID: candidate.originalID) 491 ) 492 } 493 } 494 495 // MARK: - CloudKit archive I/O 496 497 /// Starts a fresh descending scan across Chronicle and completed Game 498 /// records. Only scalar metadata is read while locating the initial 499 /// seven-day window; full zones/assets are fetched for visible records. 500 func loadRecentCompleted(since cutoff: Date) async -> CompletedPage { 501 chronicleCursor = nil 502 bufferedCompleted = [] 503 504 do { 505 let walk = try await CompletedMetadataPageWalker.recent( 506 cutoff: cutoff, 507 completedAt: \.completedAt, 508 fetch: fetchChronicleMetadataPage(continuing:) 509 ) 510 chronicleCursor = walk.cursor 511 bufferedCompleted.append(contentsOf: walk.buffered) 512 513 let gameMetadata = try await fetchCompletedGameMetadata() 514 let recentGames = gameMetadata.filter { $0.completedAt >= cutoff } 515 bufferedCompleted.append( 516 contentsOf: gameMetadata.filter { $0.completedAt < cutoff } 517 ) 518 bufferedCompleted = mergedMetadata(bufferedCompleted) 519 520 let selected = mergedMetadata(walk.selected + recentGames) 521 await hydrateCompleted(selected) 522 await trimMaterializedChronicles(before: cutoff) 523 return CompletedPage( 524 oldestCompletedAt: selected.last?.completedAt, 525 hasMore: !bufferedCompleted.isEmpty || chronicleCursor != nil 526 ) 527 } catch { 528 syncMonitor?.recordError("load recent completed games", error) 529 return CompletedPage( 530 oldestCompletedAt: nil, 531 hasMore: !bufferedCompleted.isEmpty || chronicleCursor != nil 532 ) 533 } 534 } 535 536 /// Continues the metadata scan and hydrates the next available records, 537 /// crossing arbitrarily large date gaps in one request. 538 func loadMoreCompleted() async -> CompletedPage { 539 do { 540 while uniqueGameCount(in: bufferedCompleted) < Self.completedPageSize, 541 let cursor = chronicleCursor { 542 let page = try await fetchChronicleMetadataPage( 543 continuing: cursor 544 ) 545 chronicleCursor = page.cursor 546 bufferedCompleted.append(contentsOf: page.records) 547 bufferedCompleted = mergedMetadata(bufferedCompleted) 548 } 549 550 let selected = Array(bufferedCompleted.prefix(Self.completedPageSize)) 551 let selectedIDs = Set(selected.map(\.originalGameID)) 552 bufferedCompleted.removeAll { 553 selectedIDs.contains($0.originalGameID) 554 } 555 await hydrateCompleted(selected) 556 return CompletedPage( 557 oldestCompletedAt: selected.last?.completedAt, 558 hasMore: !bufferedCompleted.isEmpty || chronicleCursor != nil 559 ) 560 } catch { 561 syncMonitor?.recordError("load more completed games", error) 562 return CompletedPage( 563 oldestCompletedAt: nil, 564 hasMore: !bufferedCompleted.isEmpty || chronicleCursor != nil 565 ) 566 } 567 } 568 569 /// One-launch migration backstop for completed live zones that this device 570 /// has never materialised. It begins only after the bounded Game List load 571 /// has returned, so old payloads cannot delay or expand the initial list. 572 /// Each fetched game is immediately compacted; later launches skip it 573 /// because the retained live row is already local while retirement settles. 574 func migrateMissingCompletedGames() async { 575 do { 576 let metadata = try await fetchCompletedGameMetadata() 577 let ctx = persistence.container.newBackgroundContext() 578 var localLiveIDs: Set<UUID> = await ctx.perform { 579 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 580 req.predicate = NSPredicate( 581 format: "completedAt != nil AND isAccessRevoked == NO " + 582 "AND ckRecordName BEGINSWITH %@", 583 "game-" 584 ) 585 return Set(((try? ctx.fetch(req)) ?? []).compactMap(\.id)) 586 } 587 588 for item in mergedMetadata(metadata) { 589 guard !localLiveIDs.contains(item.originalGameID), 590 case .game(let gameID, let zoneID, let scope) = item.source 591 else { continue } 592 do { 593 guard try await syncEngine.fetchCompletedGameDirect( 594 gameID: gameID, 595 zoneID: zoneID, 596 scope: scope 597 ) else { continue } 598 localLiveIDs.insert(gameID) 599 await archiveIfNeeded(gameID: gameID) 600 } catch { 601 syncMonitor?.recordError("migrate completed game", error) 602 } 603 } 604 await reconcileUnarchived() 605 } catch { 606 syncMonitor?.recordError("scan completed-game migration", error) 607 } 608 } 609 610 private func fetchChronicleMetadataPage( 611 continuing cursor: CKQueryOperation.Cursor? = nil 612 ) async throws -> (records: [CompletedMetadata], cursor: CKQueryOperation.Cursor?) { 613 let database = container.privateCloudDatabase 614 let result: (matchResults: [(CKRecord.ID, Result<CKRecord, any Error>)], queryCursor: CKQueryOperation.Cursor?) 615 if let cursor { 616 result = try await database.records( 617 continuingMatchFrom: cursor, 618 desiredKeys: ["completedAt"], 619 resultsLimit: 50 620 ) 621 } else { 622 let query = CKQuery(recordType: Archive.recordType, predicate: NSPredicate(value: true)) 623 query.sortDescriptors = [NSSortDescriptor(key: "completedAt", ascending: false)] 624 result = try await database.records( 625 matching: query, 626 inZoneWith: Archive.zoneID, 627 desiredKeys: ["completedAt"], 628 resultsLimit: 50 629 ) 630 } 631 let records = result.matchResults.compactMap { _, item -> CompletedMetadata? in 632 guard let record = try? item.get(), 633 let originalGameID = Archive.originalGameID( 634 fromName: record.recordID.recordName 635 ), 636 let completedAt = record["completedAt"] as? Date 637 else { return nil } 638 return CompletedMetadata( 639 originalGameID: originalGameID, 640 completedAt: completedAt, 641 source: .chronicle(record.recordID) 642 ) 643 } 644 return (records, result.queryCursor) 645 } 646 647 private func fetchCompletedGameMetadata() async throws -> [CompletedMetadata] { 648 async let privateMetadata = fetchCompletedGameMetadata( 649 database: container.privateCloudDatabase, 650 scope: .private 651 ) 652 async let sharedMetadata = fetchCompletedGameMetadata( 653 database: container.sharedCloudDatabase, 654 scope: .shared 655 ) 656 return try await privateMetadata + sharedMetadata 657 } 658 659 private func fetchCompletedGameMetadata( 660 database: CKDatabase, 661 scope: DatabaseScope 662 ) async throws -> [CompletedMetadata] { 663 let zoneIDs = try await database.allRecordZones() 664 .map(\.zoneID) 665 .filter { RecordSerializer.gameID(fromGameRecordName: $0.zoneName) != nil } 666 let recordIDs = zoneIDs.map { 667 CKRecord.ID(recordName: $0.zoneName, zoneID: $0) 668 } 669 670 var metadata: [CompletedMetadata] = [] 671 for start in stride(from: 0, to: recordIDs.count, by: 200) { 672 let end = min(start + 200, recordIDs.count) 673 let batch = Array(recordIDs[start..<end]) 674 let results = try await database.records( 675 for: batch, 676 desiredKeys: ["completedAt"] 677 ) 678 for (recordID, result) in results { 679 guard let record = try? result.get(), 680 let gameID = RecordSerializer.gameID( 681 fromGameRecordName: recordID.recordName 682 ), 683 let completedAt = record["completedAt"] as? Date 684 else { continue } 685 metadata.append(CompletedMetadata( 686 originalGameID: gameID, 687 completedAt: completedAt, 688 source: .game( 689 gameID: gameID, 690 zoneID: recordID.zoneID, 691 scope: scope 692 ) 693 )) 694 } 695 } 696 return metadata 697 } 698 699 /// Collapses a live Game and Chronicle for the same original puzzle into 700 /// one candidate. Chronicle is the canonical completed representation; 701 /// the live row remains only for acknowledgement and zone retirement. 702 private func mergedMetadata( 703 _ metadata: [CompletedMetadata] 704 ) -> [CompletedMetadata] { 705 var byGameID: [UUID: CompletedMetadata] = [:] 706 for item in metadata { 707 if let existing = byGameID[item.originalGameID] { 708 if !item.isLiveGame && existing.isLiveGame { 709 byGameID[item.originalGameID] = item 710 } 711 } else { 712 byGameID[item.originalGameID] = item 713 } 714 } 715 return byGameID.values.sorted { lhs, rhs in 716 if lhs.completedAt != rhs.completedAt { 717 return lhs.completedAt > rhs.completedAt 718 } 719 return lhs.originalGameID.uuidString < rhs.originalGameID.uuidString 720 } 721 } 722 723 private func uniqueGameCount(in metadata: [CompletedMetadata]) -> Int { 724 Set(metadata.map(\.originalGameID)).count 725 } 726 727 private func hydrateCompleted(_ metadata: [CompletedMetadata]) async { 728 let chronicles = metadata.compactMap { item -> CKRecord.ID? in 729 if case .chronicle(let recordID) = item.source { return recordID } 730 return nil 731 } 732 await materializeChronicles(chronicles) 733 734 for item in metadata { 735 guard case .game(let gameID, let zoneID, let scope) = item.source 736 else { continue } 737 do { 738 _ = try await syncEngine.fetchCompletedGameDirect( 739 gameID: gameID, 740 zoneID: zoneID, 741 scope: scope 742 ) 743 } catch { 744 syncMonitor?.recordError("load completed game", error) 745 } 746 } 747 } 748 749 private func materializeChronicles(_ recordIDs: [CKRecord.ID]) async { 750 guard !recordIDs.isEmpty else { return } 751 let database = container.privateCloudDatabase 752 let result = try? await database.records( 753 for: recordIDs, 754 desiredKeys: ["completedAt", Archive.payloadKey] 755 ) 756 guard let result else { return } 757 let records = result.compactMap { _, item in try? item.get() } 758 let ctx = persistence.container.newBackgroundContext() 759 await ctx.perform { 760 for record in records { 761 _ = self.syncEngine.applyPreferredArchiveRecord(record, in: ctx) 762 } 763 if ctx.hasChanges { try? ctx.save() } 764 } 765 } 766 767 /// Evicts only rows materialized from Chronicle. Live completed games stay 768 /// local until their archive/acknowledgement/retirement work has converged. 769 private func trimMaterializedChronicles(before cutoff: Date) async { 770 let ctx = persistence.container.newBackgroundContext() 771 await ctx.perform { 772 let req = NSFetchRequest<GameEntity>(entityName: "GameEntity") 773 req.predicate = NSPredicate( 774 format: "completedAt < %@ AND ckRecordName BEGINSWITH %@", 775 cutoff as NSDate, 776 "chronicle-" 777 ) 778 for game in (try? ctx.fetch(req)) ?? [] { 779 if let name = game.ckRecordName, 780 let originalID = Archive.originalGameID(fromName: name) { 781 let liveReq = NSFetchRequest<GameEntity>(entityName: "GameEntity") 782 liveReq.predicate = NSPredicate( 783 format: "id == %@", 784 originalID as CVarArg 785 ) 786 liveReq.fetchLimit = 1 787 (try? ctx.fetch(liveReq).first)? 788 .isSupersededByChronicle = false 789 } 790 ctx.delete(game) 791 } 792 if ctx.hasChanges { try? ctx.save() } 793 } 794 } 795 796 private func fetchArchive(originalGameID: UUID) async -> StoredArchive? { 797 let name = Archive.recordName(forOriginalGameID: originalGameID) 798 let commonID = CKRecord.ID(recordName: name, zoneID: Archive.zoneID) 799 if let record = try? await container.privateCloudDatabase.record(for: commonID), 800 let payload = Archive.payload(from: record) { 801 return StoredArchive(payload: payload, isLegacy: false) 802 } 803 let legacyID = CKRecord.ID( 804 recordName: Archive.legacyRecordName(forOriginalGameID: originalGameID), 805 zoneID: Archive.legacyZoneID(forOriginalGameID: originalGameID) 806 ) 807 if let record = try? await container.privateCloudDatabase.record(for: legacyID), 808 let payload = Archive.payload(from: record) { 809 return StoredArchive(payload: payload, isLegacy: true) 810 } 811 return nil 812 } 813 814 @discardableResult 815 private func write( 816 _ snapshot: Archive.Snapshot, 817 replayState: Archive.ReplayState 818 ) async -> Bool { 819 do { 820 try await ensureArchiveZone() 821 let package = try Archive.recordPackage( 822 from: snapshot, 823 replayState: replayState 824 ) 825 defer { removeTemporaryArchiveFiles(package.temporaryAssetFileURLs) } 826 try await save(package.record) 827 return true 828 } catch { 829 syncMonitor?.recordError("archive game", error) 830 eventLog?.note( 831 "GameArchiver: write deferred for \(snapshot.originalGameID.uuidString); " + 832 "will retry on cold launch — \(error)", 833 level: "error" 834 ) 835 return false 836 } 837 } 838 839 private func ensureArchiveZone() async throws { 840 guard !ensuredArchiveZone else { return } 841 do { 842 try await createArchiveZone(Archive.zoneID) 843 } catch { 844 guard Self.isZoneAlreadyExists(error) else { throw error } 845 } 846 ensuredArchiveZone = true 847 } 848 849 private func createArchiveZone(_ zoneID: CKRecordZone.ID) async throws { 850 try await withCheckedThrowingContinuation { (continuation: CheckedContinuation<Void, Error>) in 851 let operation = CKModifyRecordZonesOperation( 852 recordZonesToSave: [CKRecordZone(zoneID: zoneID)], 853 recordZoneIDsToDelete: nil 854 ) 855 operation.qualityOfService = .utility 856 operation.modifyRecordZonesResultBlock = { continuation.resume(with: $0) } 857 container.privateCloudDatabase.add(operation) 858 } 859 } 860 861 private func save(_ record: CKRecord) async throws { 862 try await withCheckedThrowingContinuation { (continuation: CheckedContinuation<Void, Error>) in 863 let operation = CKModifyRecordsOperation(recordsToSave: [record], recordIDsToDelete: nil) 864 operation.savePolicy = .allKeys 865 operation.qualityOfService = .utility 866 operation.modifyRecordsResultBlock = { continuation.resume(with: $0) } 867 container.privateCloudDatabase.add(operation) 868 } 869 } 870 871 private nonisolated static func isZoneAlreadyExists(_ error: Error) -> Bool { 872 guard let ckError = error as? CKError else { return false } 873 if ckError.code == .serverRejectedRequest { 874 return ckError.localizedDescription.lowercased().contains("already exist") 875 } 876 if ckError.code == .partialFailure { 877 return ckError.partialErrorsByItemID?.values.contains { 878 isZoneAlreadyExists($0) 879 } ?? false 880 } 881 return false 882 } 883 884 private func removeTemporaryArchiveFiles(_ urls: [URL]) { 885 for url in urls { 886 do { 887 try FileManager.default.removeItem(at: url) 888 } catch { 889 eventLog?.note( 890 "GameArchiver: failed to remove temporary archive asset " + 891 "\(url.lastPathComponent) — \(error)", 892 level: "error" 893 ) 894 } 895 } 896 } 897 898 // MARK: - Account-wide migration grace 899 900 private static let graceKey = "archiveGraceStart.v1.1.0" 901 902 /// Uses ubiquitous key-value storage as an account-wide shared default, 903 /// with an author-keyed local fallback for immediate reads and account 904 /// switches. Devices continually publish the earliest value they have seen, 905 /// so a delayed KVS update can postpone retirement but can never make the 906 /// grace period shorter than the first v1.1.0 launch observed by a device. 907 private func accountGraceStart() async -> Date? { 908 guard let authorID = localIdentity()?.authorID, !authorID.isEmpty else { return nil } 909 let localKey = "\(Self.graceKey).\(authorID)" 910 ubiquitousStore?.synchronize() 911 let localValue = localDefaults.object(forKey: localKey) as? Double 912 let cloudValue = ubiquitousStore?.object(forKey: Self.graceKey) as? Double 913 let earliest = [localValue, cloudValue].compactMap { $0 }.min() 914 ?? Date().timeIntervalSince1970 915 localDefaults.set(earliest, forKey: localKey) 916 if cloudValue == nil || earliest < cloudValue! { 917 ubiquitousStore?.set(earliest, forKey: Self.graceKey) 918 ubiquitousStore?.synchronize() 919 } 920 return Date(timeIntervalSince1970: earliest) 921 } 922 }