field_ios

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

TeraSubmissionStore.swift (12679B)


      1 import Foundation
      2 
      3 /// One explicit native waiter. The runtime retains admission for a late FFI
      4 /// callback after a waiter deadline; Rust owns durable receipts and effects.
      5 /// Editing changes never invalidate this worker or replace its captured request.
      6 @MainActor
      7 final class TeraSubmissionStore: ObservableObject {
      8   @Published private(set) var request: TeraSubmissionRequest?
      9   @Published private(set) var status: TeraSubmissionStatus?
     10   @Published private(set) var isWorking = false
     11   @Published private(set) var message: String?
     12   @Published private(set) var failureCode: String?
     13   let inventory: TeraSubmissionInventory
     14   var changed: () -> Void = {}
     15   private let client: TeraRuntimeClient
     16   private let composer: TeraComposerAutosave
     17   private let media: (any TeraAddMediaHandling)?
     18   private var capture: TeraComposerCapture?
     19   private var scope: TeraComposerScope?
     20   private var context: TeraLocalNetwork?
     21   private var generation = TeraSessionGeneration.initial
     22   private var worker: Task<Void, Never>?
     23   private var paused = false
     24   private let stopControl = TeraSubmissionStopControl()
     25   private(set) var isNewComposer = false
     26 
     27   init(client: TeraRuntimeClient, composer: TeraComposerAutosave, media: (any TeraAddMediaHandling)?) {
     28     self.client = client
     29     self.composer = composer
     30     self.media = media
     31     inventory = TeraSubmissionInventory(client: client)
     32   }
     33 
     34   deinit { worker?.cancel() }
     35 
     36   var hasAction: Bool {
     37     request != nil || capture != nil
     38   }
     39 
     40   var usesCurrentAction: Bool {
     41     hasAction && !isNewComposer
     42   }
     43 
     44   var canReplaceEditing: Bool {
     45     worker == nil || (capture == nil && status != nil)
     46   }
     47 
     48   func configure(scope: TeraComposerScope, context: TeraLocalNetwork) {
     49     self.context = context
     50     guard self.scope != scope else { return }
     51     stop()
     52     self.scope = scope
     53     capture = nil
     54     request = nil
     55     status = nil
     56     stopControl.reset()
     57     message = nil
     58     failureCode = nil
     59     isNewComposer = false
     60     inventory.configure(scope: scope)
     61     changed()
     62   }
     63 
     64   func start() {
     65     paused = false
     66     inventory.start()
     67   }
     68 
     69   func refreshSelected() async {
     70     if request != nil, worker == nil {
     71       await run(advancing: false)
     72     }
     73   }
     74 
     75   /// Explicit continuation uses only the retained original request/capture.
     76   /// The runtime remains authoritative for every effect and retry admission.
     77   func continueSelected() async {
     78     guard hasAction, status?.canOfferContinuation != false else { return }
     79     await run(advancing: true)
     80   }
     81 
     82   func stop() {
     83     generation = generation.invalidated()
     84     paused = true
     85     worker?.cancel()
     86     inventory.stop()
     87     // Keep this waiter until it returns. Runtime admission separately remains
     88     // occupied until an abandoned FFI callback has actually completed.
     89   }
     90 
     91   func stopWaiting() {
     92     worker?.cancel()
     93     message = "Waiting stopped. The original request is retained; remote effects may still complete."
     94     changed()
     95   }
     96 
     97   func requestStop() async {
     98     guard !paused else { return }
     99     stopControl.request()
    100     let requested = generation
    101     message = "Stop is pending confirmation. The original request and its effects are retained."
    102     changed()
    103     guard let request else { return }
    104     do {
    105       let stopped = try await client.requestSubmissionStop(request: request)
    106       try accept(stopped, generation: requested)
    107       message = status?.summary
    108     } catch {
    109       guard generation == requested, !paused, self.request == request else { return }
    110       message = "Stop is not yet confirmed. Reconcile the original request to retain the stop."
    111     }
    112     changed()
    113   }
    114 
    115   func submit(form: TeraAddForm) async {
    116     guard !paused, !Task.isCancelled else { return }
    117     if let worker {
    118       await worker.value; return
    119     }
    120     if isNewComposer {
    121       newAction()
    122     }
    123     do {
    124       if !hasAction {
    125         capture = try composer.beginSubmissionCapture(TeraComposerForm(editing: form))
    126       }
    127     } catch {
    128       message = TeraAddPresentation.message(for: error)
    129       changed()
    130       return
    131     }
    132     await run(advancing: true)
    133   }
    134 
    135   /// Viewing saved work does not change the editing buffer or begin an effect.
    136   func select(_ summary: TeraSubmissionSummary) async {
    137     guard worker == nil, !paused, summary.request.scope == scope else { return }
    138     guard capture == nil else {
    139       message = "Reconcile the original submission or choose New before selecting another operation."
    140       changed()
    141       return
    142     }
    143     request = summary.request
    144     status = nil
    145     stopControl.reset()
    146     await run(advancing: false)
    147   }
    148 
    149   /// New is explicit. Preserve the old captured source in its own composer row
    150   /// before the editing interlock saves newer input under a fresh composer ID.
    151   func preserveForReplacement(currentForm: () -> TeraAddForm) async throws {
    152     guard canReplaceEditing else {
    153       throw TeraRuntimeFailure.local(operation: "submission.new", code: "operation_in_progress",
    154                                      safeMessage: "The original submission is still returning. Keep editing and retry New when it finishes.")
    155     }
    156     if capture != nil {
    157       let requested = generation
    158       if let request {
    159         _ = try await client.recoverSubmission(request: request)
    160       }
    161       try ensureCurrent(requested)
    162       composer.reset(scope: scope)
    163       composer.change(TeraComposerForm(editing: currentForm()))
    164       capture = nil
    165     }
    166   }
    167 
    168   func newAction() {
    169     guard canReplaceEditing else { return }
    170     if worker != nil {
    171       isNewComposer = true
    172       changed()
    173       return
    174     }
    175     isNewComposer = false
    176     capture = nil
    177     request = nil
    178     status = nil
    179     stopControl.reset()
    180     message = nil
    181     failureCode = nil
    182     changed()
    183     inventory.start()
    184   }
    185 
    186   private func run(advancing: Bool) async {
    187     guard worker == nil, !paused else { return }
    188     let requested = generation
    189     isWorking = true
    190     message = nil
    191     failureCode = nil
    192     changed()
    193     let task = Task { @MainActor [weak self] in
    194       guard let self else { return }
    195       await execute(advancing: advancing, generation: requested)
    196       worker = nil
    197       isWorking = false
    198       changed()
    199     }
    200     worker = task
    201     // Cancelling the view's waiter does not cancel durable work. Stop waiting
    202     // is an explicit host action; both paths retain the same request.
    203     await task.value
    204   }
    205 
    206   private func execute(advancing: Bool, generation requested: TeraSessionGeneration) async {
    207     do {
    208       try ensureCurrent(requested)
    209       let current = try await resolve(advancing: advancing, generation: requested)
    210       guard var current else {
    211         message = "The original submission is reserved. Retry it to finish local preparation."
    212         changed()
    213         return
    214       }
    215       try accept(current, generation: requested)
    216       if stopControl.requested, !current.delivery.isStopped {
    217         let stopped = try await client.requestSubmissionStop(request: current.request)
    218         try accept(stopped, generation: requested)
    219       }
    220       current = status ?? current
    221       try await reconcileLocal(current, generation: requested)
    222       if advancing {
    223         let effects = TeraSubmissionEffects(client: client, media: media,
    224                                             ensure: { try self.ensureCurrent(requested) },
    225                                             accept: { try self.accept($0, generation: requested) },
    226                                             mayStart: { !self.stopControl.requested }, stopControl: stopControl)
    227         try await effects.advance(current)
    228       } else {
    229         try await reconcileBackground(current, generation: requested)
    230       }
    231       try ensureCurrent(requested)
    232       if let status {
    233         try await reconcileLocal(status, generation: requested)
    234       }
    235       message = status?.summary
    236     } catch {
    237       guard generation == requested, !paused else { return }
    238       if let request, !Task.isCancelled {
    239         // Read after every uncertain result. Never infer absence from a read
    240         // error, open newer form media, or create another operation on retry.
    241         if let recovered = try? await client.recoverSubmission(request: request) {
    242           try? accept(recovered, generation: requested)
    243           try? await reconcileLocal(recovered, generation: requested)
    244         }
    245       }
    246       guard generation == requested, !paused else { return }
    247       failureCode = TeraAddPresentation.failure(for: error)?.code
    248       message = Task.isCancelled
    249         ? "Waiting stopped. Retry the original request to reconcile its outcome."
    250         : "Original request retained. \(TeraAddPresentation.message(for: error))"
    251     }
    252     guard generation == requested, !paused else { return }
    253     changed()
    254     inventory.start()
    255   }
    256 
    257   private func reconcileLocal(_ current: TeraSubmissionStatus, generation requested: TeraSessionGeneration) async throws {
    258     try ensureCurrent(requested)
    259     guard current.settlement.signed > 0 else { return }
    260     guard let context, context.id == current.request.scope.localNetworkID else { throw TeraComposerAcknowledgment.unconfirmed }
    261     let reconciled = try await client.reconcileSubmissionLocal(request: current.request, context: context)
    262     try accept(reconciled, generation: requested)
    263   }
    264 
    265   private func accept(_ value: TeraSubmissionStatus, generation requested: TeraSessionGeneration) throws {
    266     try ensureCurrent(requested)
    267     guard value.request == request, value.request.scope == scope,
    268           status == nil || (status?.operationID == value.operationID && status?.intentID == value.intentID
    269             && status?.captured == value.captured)
    270     else {
    271       throw TeraComposerAcknowledgment.unconfirmed
    272     }
    273     if let status, value.revision < status.revision || !value.delivery.follows(status.delivery)
    274       || !value.targetDetails.follows(status.targetDetails)
    275       || value.settlement.signed < status.settlement.signed || value.settlement.admitted < status.settlement.admitted {
    276         return
    277       }
    278     status = value
    279     if let capture {
    280       composer.releaseSubmissionCapture(capture)
    281       self.capture = nil
    282     }
    283     changed()
    284   }
    285 
    286   private func ensureCurrent(_ requested: TeraSessionGeneration) throws {
    287     guard requested == generation, generation.isActive, !paused, !Task.isCancelled else { throw CancellationError() }
    288   }
    289 }
    290 
    291 private extension TeraSubmissionStore {
    292   private func resolve(advancing: Bool, generation requested: TeraSessionGeneration) async throws -> TeraSubmissionStatus? {
    293       if request == nil {
    294         guard let capture else { throw TeraComposerAcknowledgment.unconfirmed }
    295         let saved = try await composer.saveSubmissionCapture(capture)
    296         try ensureCurrent(requested)
    297         let commandID = try await client.reserveSubmissionID()
    298         try ensureCurrent(requested)
    299         request = TeraSubmissionRequest(commandID: commandID, scope: saved.scope,
    300                                         composerID: saved.id, expectedRevision: saved.revision)
    301         changed()
    302       }
    303       guard let request else { throw TeraComposerAcknowledgment.unconfirmed }
    304       var current = try await client.recoverSubmission(request: request)
    305       try ensureCurrent(requested)
    306       if current == nil, advancing {
    307         let reservation = try await client.reserveSubmission(request: request)
    308         try ensureCurrent(requested)
    309         if let capture {
    310           guard reservation.commandID == request.commandID,
    311                 reservation.captured.scope == request.scope,
    312                 reservation.captured.id == request.composerID,
    313                 reservation.captured.revision == request.expectedRevision
    314           else {
    315             throw TeraComposerAcknowledgment.unconfirmed
    316           }
    317           try composer.continueEditing(after: capture, reserved: reservation.captured)
    318           self.capture = nil
    319         }
    320         let opened = try await TeraOpenedMedia.open(reservation.captured.form.editingValue.media, using: media)
    321         defer { opened.close() }
    322         try ensureCurrent(requested)
    323         current = try await client.prepareSubmission(request: request, media: opened.handles)
    324       }
    325       try ensureCurrent(requested)
    326       return current
    327   }
    328 
    329   private func reconcileBackground(_ current: TeraSubmissionStatus, generation requested: TeraSessionGeneration) async throws {
    330     if current.delivery.isStopped {
    331       try await TeraSubmissionEffects(client: client, media: media,
    332                                       ensure: { try self.ensureCurrent(requested) },
    333                                       accept: { try self.accept($0, generation: requested) }).reconcileStopped(current)
    334     } else {
    335       try await media?.reconcileBackgroundSubmissions([current], client: client)
    336     }
    337   }
    338 }