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 }