RadrootsAppleBackgroundTransfer.swift (14832B)
1 import Foundation 2 3 public actor RadrootsAppleBackgroundTransfer: RadrootsBackgroundTransfer { 4 private let store: any RadrootsBackgroundTransferStore 5 private let adapters: RadrootsAppleBackgroundTransferAdapters 6 private var admissions: Set<RadrootsBackgroundTransferIdentifier> = [] 7 private var stoppedAdmissions: [RadrootsBackgroundTransferIdentifier: RadrootsBackgroundTransferState] = [:] 8 private struct InactiveExecutionBodyError: Error { 9 let underlying: any Error 10 } 11 12 public init(store: any RadrootsBackgroundTransferStore, adapters: RadrootsAppleBackgroundTransferAdapters) { 13 self.store = store 14 self.adapters = adapters 15 } 16 17 public init(roots: RadrootsAppleFileRoots, sessionIdentifier: String) throws { 18 let store = RadrootsAppleBackgroundTransferStore(roots: roots) 19 let resolver = RadrootsAppleBackgroundTransferFileResolver(roots: roots) 20 let name = try RadrootsBackgroundTransferValidation.normalizedIdentifier(sessionIdentifier) 21 let downloadRoot = try roots.resolvedURL(for: RadrootsFileReference( 22 scope: .temporary, relativePath: "background_transfers/\(name)/downloads" 23 ), allowRootDirectory: true) 24 self.store = store 25 adapters = try .live(sessionIdentifier: name, store: store, fileResolver: resolver, 26 downloadStagingRoot: downloadRoot) 27 } 28 29 public func enqueue(_ request: RadrootsBackgroundTransferRequest) async throws -> RadrootsBackgroundTransferHandle { 30 try reserve(request.identifier) 31 defer { release(request.identifier) } 32 return try await admission(request, retry: false) 33 } 34 35 public func retry(_ request: RadrootsBackgroundTransferRequest) async throws -> RadrootsBackgroundTransferHandle { 36 try reserve(request.identifier) 37 defer { release(request.identifier) } 38 return try await admission(request, retry: true) 39 } 40 41 public func withInactiveExecution<Result: Sendable>( 42 for identifier: RadrootsBackgroundTransferIdentifier, 43 operation: @escaping @Sendable (RadrootsBackgroundTransferSnapshot?) async throws -> Result 44 ) async throws -> Result { 45 try Task.checkCancellation() 46 try reserve(identifier) 47 defer { release(identifier) } 48 do { 49 return try await store.withAdmission(for: identifier) { 50 let snapshot = try await self.inactiveSnapshot(identifier) 51 try Task.checkCancellation() 52 do { return try await operation(snapshot) } 53 catch { throw InactiveExecutionBodyError(underlying: error) } 54 } 55 } catch let error as InactiveExecutionBodyError { 56 throw error.underlying 57 } catch is CancellationError { 58 throw CancellationError() 59 } catch let error as RadrootsBackgroundTransferError { 60 throw error 61 } catch { 62 throw RadrootsBackgroundTransferError.persistence(error) 63 } 64 } 65 66 private func inactiveSnapshot( 67 _ identifier: RadrootsBackgroundTransferIdentifier 68 ) async throws -> RadrootsBackgroundTransferSnapshot? { 69 guard try await !activeIdentifiers().contains(identifier), stoppedAdmissions[identifier] == nil else { 70 throw RadrootsBackgroundTransferError.transferFailure 71 } 72 let snapshot = try await load(identifier) 73 if let snapshot { 74 guard [.failed, .interrupted, .cancelled, .expired].contains(snapshot.state), 75 snapshot.response == nil, snapshot.downloadedArtifact == nil 76 else { throw RadrootsBackgroundTransferError.transferFailure } 77 } 78 // Do not rewrite state, drop a receipt or cancel an OS task. Completion 79 // callbacks may still persist late evidence while admission is held. 80 return snapshot 81 } 82 83 private func admission( 84 _ request: RadrootsBackgroundTransferRequest, retry: Bool 85 ) async throws -> RadrootsBackgroundTransferHandle { 86 do { 87 return try await store.withAdmission(for: request.identifier) { 88 try await self.admitReserved(request, retry: retry) 89 } 90 } catch let error as RadrootsBackgroundTransferError { 91 throw error 92 } catch { 93 throw RadrootsBackgroundTransferError.persistence(error) 94 } 95 } 96 97 private func admitReserved( 98 _ request: RadrootsBackgroundTransferRequest, retry: Bool 99 ) async throws -> RadrootsBackgroundTransferHandle { 100 let existing = try await load(request.identifier) 101 if retry { 102 guard let existing, [.failed, .interrupted, .cancelled, .expired].contains(existing.state), 103 try existing.request.redactedForPersistence() == request.redactedForPersistence() 104 else { throw RadrootsBackgroundTransferError.invalidRequest } 105 } else if existing != nil { 106 throw RadrootsBackgroundTransferError.invalidRequest 107 } 108 guard try await !activeIdentifiers().contains(request.identifier) else { 109 throw RadrootsBackgroundTransferError.invalidRequest 110 } 111 return try await schedule(request, replacing: existing) 112 } 113 114 private func reserve(_ identifier: RadrootsBackgroundTransferIdentifier) throws { 115 guard admissions.count < 256, admissions.insert(identifier).inserted else { 116 throw RadrootsBackgroundTransferError.invalidRequest 117 } 118 } 119 120 private func release(_ identifier: RadrootsBackgroundTransferIdentifier) { 121 admissions.remove(identifier) 122 stoppedAdmissions.removeValue(forKey: identifier) 123 } 124 125 private func schedule( 126 _ request: RadrootsBackgroundTransferRequest, replacing existing: RadrootsBackgroundTransferSnapshot? 127 ) async throws -> RadrootsBackgroundTransferHandle { 128 guard !Task.isCancelled, stoppedAdmissions[request.identifier] == nil else { 129 throw RadrootsBackgroundTransferError.transferFailure 130 } 131 let queued = try RadrootsBackgroundTransferSnapshot( 132 request: request.redactedForPersistence(), updatedAt: adapters.now(), executionID: UUID() 133 ) 134 guard try await exchange(existing, queued) else { throw RadrootsBackgroundTransferError.invalidRequest } 135 if Task.isCancelled { 136 stoppedAdmissions[request.identifier] = .cancelled 137 } 138 if let state = stoppedAdmissions[request.identifier] { 139 _ = try await exchange(queued, queued.transitioned(to: state, at: adapters.now(), 140 failure: state == .expired ? .expired : nil, 141 possibleRemoteOrphan: request.isUpload)) 142 throw RadrootsBackgroundTransferError.transferFailure 143 } 144 do { 145 guard let executionID = queued.executionID else { throw RadrootsBackgroundTransferError.invalidRequest } 146 try await adapters.enqueue(request, executionID) 147 } catch { 148 let persistenceFailure = RadrootsBackgroundTransferError.persistence(error) 149 if stoppedAdmissions[request.identifier] != nil { 150 try? await adapters.cancel(request.identifier) 151 } 152 if let current = try await load(request.identifier), current.executionID == queued.executionID, 153 current.state == .queued || current.state == .running 154 { 155 _ = try await exchange(current, current.transitioned(to: .failed, at: adapters.now(), 156 failure: .enqueueFailed, 157 possibleRemoteOrphan: request.isUpload)) 158 } 159 if persistenceFailure == .spaceInsufficient || persistenceFailure == .receiptCapacityExceeded { 160 throw persistenceFailure 161 } 162 throw RadrootsBackgroundTransferError.transferFailure 163 } 164 if stoppedAdmissions[request.identifier] != nil || Task.isCancelled { 165 try await stop(request.identifier, state: stoppedAdmissions[request.identifier] ?? .cancelled) 166 } else if let current = try await load(request.identifier), current.executionID == queued.executionID, 167 current.state == .queued 168 { 169 _ = try await exchange(current, current.transitioned(to: .running, at: adapters.now())) 170 } 171 return RadrootsBackgroundTransferHandle(request: request) 172 } 173 174 public func cancel(_ identifier: RadrootsBackgroundTransferIdentifier) async throws { 175 try await stop(identifier, state: .cancelled) 176 } 177 178 public func expire(_ identifier: RadrootsBackgroundTransferIdentifier) async throws { 179 try await stop(identifier, state: .expired) 180 } 181 182 private func stop( 183 _ identifier: RadrootsBackgroundTransferIdentifier, 184 state: RadrootsBackgroundTransferState 185 ) async throws { 186 if admissions.contains(identifier) { 187 stoppedAdmissions[identifier] = state 188 } 189 if let existing = try await load(identifier) { 190 guard [.queued, .running, .interrupted].contains(existing.state) else { 191 if admissions.contains(identifier), [.cancelled, .expired].contains(existing.state) { 192 do { 193 try await adapters.cancel(identifier) 194 } catch { 195 throw RadrootsBackgroundTransferError.transferFailure 196 } 197 } 198 return 199 } 200 let desired = try existing.transitioned(to: state, at: adapters.now(), 201 failure: state == .expired ? .expired : nil, 202 possibleRemoteOrphan: existing.request.isUpload) 203 guard try await exchange(existing, desired) else { 204 throw RadrootsBackgroundTransferError.transferFailure 205 } 206 } 207 do { 208 try await adapters.cancel(identifier) 209 } catch { 210 throw RadrootsBackgroundTransferError.transferFailure 211 } 212 } 213 214 public func settle( 215 _ identifier: RadrootsBackgroundTransferIdentifier, verification: RadrootsBackgroundTransferVerification 216 ) async throws { 217 guard let existing = try await load(identifier), existing.state == .awaitingVerification else { 218 throw RadrootsBackgroundTransferError.invalidRequest 219 } 220 let desired: RadrootsBackgroundTransferSnapshot = switch verification { 221 case .accepted: 222 try existing.transitioned(to: .completed, at: adapters.now()) 223 case let .rejected(failure): 224 try existing.transitioned(to: .failed, at: adapters.now(), failure: failure, 225 possibleRemoteOrphan: existing.request.isUpload) 226 } 227 guard try await exchange(existing, desired) else { throw RadrootsBackgroundTransferError.transferFailure } 228 } 229 230 public func snapshot(for identifier: RadrootsBackgroundTransferIdentifier) async throws 231 -> RadrootsBackgroundTransferSnapshot? 232 { 233 try await snapshots().first { $0.identifier == identifier } 234 } 235 236 public func snapshots() async throws -> [RadrootsBackgroundTransferSnapshot] { 237 let active = try await activeIdentifiers() 238 let stored = try await loadAll() 239 let identifiers = Set(stored.map(\.identifier)) 240 for orphan in active.subtracting(identifiers) where !admissions.contains(orphan) { 241 guard try await !admissionIsActive(orphan) else { continue } 242 do { 243 try await adapters.cancel(orphan) 244 } catch { 245 throw RadrootsBackgroundTransferError.transferFailure 246 } 247 } 248 for existing in stored where !admissions.contains(existing.identifier) { 249 guard try await !admissionIsActive(existing.identifier) else { continue } 250 let desired: RadrootsBackgroundTransferSnapshot 251 if active.contains(existing.identifier), [.queued, .interrupted].contains(existing.state) { 252 desired = try existing.transitioned(to: .running, at: adapters.now(), failure: existing.failure) 253 } else if !active.contains(existing.identifier), [.queued, .running].contains(existing.state) { 254 desired = try existing.transitioned(to: .interrupted, at: adapters.now(), failure: .interrupted, 255 possibleRemoteOrphan: existing.request.isUpload) 256 } else { 257 continue 258 } 259 _ = try await exchange(existing, desired) 260 } 261 return try await loadAll() 262 } 263 264 public func handleEventsForBackgroundURLSession( 265 identifier: String, completionHandler: @escaping @Sendable () -> Void 266 ) async { 267 await adapters.handleBackgroundEvents(identifier, completionHandler) 268 } 269 270 private func load(_ identifier: RadrootsBackgroundTransferIdentifier) async throws 271 -> RadrootsBackgroundTransferSnapshot? 272 { 273 try await loadAll().first { $0.identifier == identifier } 274 } 275 276 private func loadAll() async throws -> [RadrootsBackgroundTransferSnapshot] { 277 do { 278 return try await store.loadSnapshots() 279 } catch let error as RadrootsBackgroundTransferError { 280 throw error 281 } catch { 282 throw RadrootsBackgroundTransferError.persistence(error) 283 } 284 } 285 286 private func exchange( 287 _ expected: RadrootsBackgroundTransferSnapshot?, _ desired: RadrootsBackgroundTransferSnapshot 288 ) async throws -> Bool { 289 do { 290 return try await store.compareExchangeSnapshot(expected: expected, desired: desired) 291 } catch let error as RadrootsBackgroundTransferError { 292 throw error 293 } catch { 294 throw RadrootsBackgroundTransferError.persistence(error) 295 } 296 } 297 298 private func activeIdentifiers() async throws -> Set<RadrootsBackgroundTransferIdentifier> { 299 do { 300 return try await adapters.activeTransferIdentifiers() 301 } catch let error as RadrootsBackgroundTransferError { 302 throw error 303 } catch { 304 throw RadrootsBackgroundTransferError.transferFailure 305 } 306 } 307 308 private func admissionIsActive(_ identifier: RadrootsBackgroundTransferIdentifier) async throws -> Bool { 309 do { 310 return try await store.admissionIsActive(for: identifier) 311 } catch let error as RadrootsBackgroundTransferError { 312 throw error 313 } catch { 314 throw RadrootsBackgroundTransferError.persistence(error) 315 } 316 } 317 }