field_ios

In-the-field app for Radroots on iOS
git clone https://radroots.dev/git/field_ios.git
Log | Files | Refs | README | LICENSE

TeraNativeRepairStore.swift (6304B)


      1 import Foundation
      2 
      3 /// Caller-owned bounded recovery; advisory status never grants effect authority.
      4 @MainActor
      5 final class TeraNativeRepairStore: ObservableObject {
      6   static let previewLimit = 64
      7   static let batchLimit = 4
      8   @Published private(set) var progress: TeraNativeRecoveryProgress?
      9   @Published private(set) var issues: [TeraNativeRecoveryIssue] = []
     10   @Published private(set) var message: String?
     11   @Published private(set) var isRunning = false
     12   private let client: TeraRuntimeClient
     13   private let media: (any TeraAddMediaHandling)?
     14   private var author: String?
     15   private var generation = TeraSessionGeneration.initial
     16   private var task: Task<Void, Never>?
     17   private enum Request { case sweep, transfer(String) }
     18   private var pending: Request?
     19 
     20   init(client: TeraRuntimeClient, media: (any TeraAddMediaHandling)?) {
     21     self.client = client
     22     self.media = media
     23   }
     24 
     25   deinit { task?.cancel() }
     26 
     27   func configure(author: String) {
     28     guard self.author != author else { return }
     29     stop()
     30     self.author = author
     31     progress = nil
     32     issues = []
     33     message = nil
     34   }
     35 
     36   func stop() {
     37     generation = generation.invalidated()
     38     pending = nil
     39     task?.cancel()
     40   }
     41 
     42   func retry() {
     43     enqueue(.sweep)
     44   }
     45 
     46   func check(_ issue: TeraNativeRecoveryIssue) {
     47     guard issues.contains(where: { $0.key == issue.key }) else { return }
     48     enqueue(.transfer(issue.key))
     49   }
     50 
     51   private func enqueue(_ request: Request) {
     52     guard !Task.isCancelled, generation.isActive else { return }
     53     // Coalesce live requests. A request after stop is retained until the old
     54     // worker actually returns, so cancellation cannot release its ownership.
     55     guard task == nil || task?.isCancelled == true else { return }
     56     pending = request
     57     startWorker()
     58   }
     59 
     60   private func startWorker() {
     61     guard task == nil, let request = pending, generation.isActive else { return }
     62     pending = nil
     63     let requested = generation
     64     isRunning = true
     65     task = Task { [weak self] in
     66       guard let self else { return }
     67       await run(requested, request: request)
     68       isRunning = false
     69       task = nil
     70       startWorker()
     71     }
     72   }
     73 
     74   @discardableResult
     75   func reconcile() async -> String? {
     76     guard !Task.isCancelled else { return message }
     77     retry()
     78     guard let task else { return message }
     79     await withTaskCancellationHandler { await task.value } onCancel: { task.cancel() }
     80     return message
     81   }
     82 
     83   /// Foreground resumes must recover even when the event observer is already
     84   /// live; neither a new draft nor a native transfer callback is required.
     85   func resumeIfObserving(_ observing: Bool) async -> Bool {
     86     guard observing else { return false }
     87     await reconcile()
     88     return true
     89   }
     90 
     91   private func run(_ requested: TeraSessionGeneration, request: Request) async {
     92     for _ in 0 ..< Self.batchLimit {
     93       guard requested == generation, !Task.isCancelled else { return }
     94       await batch(requested, request: request)
     95       if case .transfer = request {
     96         return
     97       }
     98       guard requested == generation, !Task.isCancelled,
     99             let progress, progress.pause == nil, progress.remaining > 0, progress.visited > 0
    100       else { return }
    101       await Task.yield()
    102     }
    103   }
    104 
    105   private func batch(_ requested: TeraSessionGeneration, request: Request) async {
    106     do {
    107       let result: TeraNativeRecoveryProgress? = switch request {
    108       case .sweep: try await media?.recoverNativeUploads(client: client)
    109       case let .transfer(key): try await media?.recoverNativeUpload(key: key, client: client)
    110       }
    111       guard requested == generation, !Task.isCancelled else { return }
    112       if let result {
    113         guard (0 ... TeraNativeRecoveryInventory.passLimit).contains(result.visited), result.remaining >= 0,
    114               result.issues.count <= Self.previewLimit else { throw TeraComposerAcknowledgment.unconfirmed }
    115       }
    116       let refreshed = await Self.refresh(issues, incoming: result?.issues ?? [], client: client)
    117       guard requested == generation, !Task.isCancelled else { return }
    118       issues = refreshed
    119       if case .transfer = request, let result {
    120         // A selected check does not advance the persistent sweep cursor.
    121         progress = .init(visited: result.visited, remaining: progress?.remaining ?? 0,
    122                          needsAttention: result.needsAttention, issues: result.issues, pause: result.pause)
    123       } else {
    124         progress = result
    125       }
    126       message = Self.message(progress, retained: !issues.isEmpty)
    127     } catch {
    128       guard requested == generation, !Task.isCancelled else { return }
    129       let pause = TeraNativeRecoveryClassification.pause(error)
    130       progress = .init(visited: 0, remaining: progress?.remaining ?? 0, needsAttention: pause == nil, pause: pause)
    131       message = Self.message(progress, retained: !issues.isEmpty)
    132     }
    133   }
    134 
    135   private static func refresh(_ previous: [TeraNativeRecoveryIssue], incoming: [TeraNativeRecoveryIssue],
    136                               client: TeraRuntimeClient) async -> [TeraNativeRecoveryIssue]
    137   {
    138     var values: [String: TeraNativeRecoveryIssue] = [:]
    139     let resolved = Set(incoming.filter { $0.reason == .resolved }.map(\.key))
    140     for issue in previous + incoming where issue.reason != .resolved && !resolved.contains(issue.key) {
    141       if values.count < previewLimit || values[issue.key] != nil {
    142         values[issue.key] = issue
    143       }
    144     }
    145     for key in values.keys.sorted() {
    146       if Task.isCancelled {
    147         break
    148       }
    149       guard let status = try? await client.nativeRecoveryStatus(key: key) else { continue }
    150       if status.reason == .resolved {
    151         values.removeValue(forKey: key)
    152       } else {
    153         values[key] = .init(key: key, reason: status.reason, status: status)
    154       }
    155     }
    156     return values.values.sorted { $0.key < $1.key }
    157   }
    158 
    159   private static func message(_ progress: TeraNativeRecoveryProgress?, retained: Bool) -> String? {
    160     guard let progress else { return nil }
    161     if let pause = progress.pause {
    162       return pause.message
    163     }
    164     if progress.needsAttention || retained {
    165       return "Photo recovery needs attention. Saved editing is still available."
    166     }
    167     if progress.remaining > 0 {
    168       return "More saved photo recovery remains. Saved editing is still available."
    169     }
    170     return nil
    171   }
    172 }