TeraSubmissionTestBackend.swift (16195B)
1 import Foundation 2 @testable import TeraApp 3 4 /// Controllable native-boundary test double. Production policy is exercised by 5 /// the separate generated FFI/SQLite tests, never inferred from this scheduler. 6 actor SubmissionTestBackend { 7 let recoveryScheduleStorage = NativeRecoveryScheduleTestStorage() 8 let composer: ComposerTestStorage 9 let writable: Bool 10 let offline: Bool 11 let delayedPhase: AddDelayPhase? 12 let delayAfterCompletion: Bool 13 private(set) var idCount = 0 14 private(set) var prepareCount = 0 15 private(set) var advanceCount = 0 16 private(set) var uploadCount = 0 17 private(set) var foregroundCount = 0 18 private var foregroundPause: ResourceTestPause? 19 private(set) var completionPersisted = false 20 private var reservations: [String: TeraSubmissionReservation] = [:] 21 private var operations: [String: TeraSubmissionStatus] = [:] 22 var operationCount: Int { 23 operations.count 24 } 25 26 private var preparePause: ResourceTestPause? 27 private var advancePause: ResourceTestPause? 28 private var failAfterCommit = false 29 private var failReads = false 30 private var delayed = false 31 32 init(composer: ComposerTestStorage, writable: Bool, offline: Bool, delayedPhase: AddDelayPhase?, delayAfterCompletion: Bool) { 33 self.composer = composer 34 self.writable = writable 35 self.offline = offline 36 self.delayedPhase = delayedPhase 37 self.delayAfterCompletion = delayAfterCompletion 38 } 39 40 func pausePrepare(_ pause: ResourceTestPause, loseReceipt: Bool = false) { 41 preparePause = pause 42 failAfterCommit = loseReceipt 43 } 44 45 func pauseAdvance(_ pause: ResourceTestPause) { 46 advancePause = pause 47 } 48 49 func pauseForeground(_ pause: ResourceTestPause) { 50 foregroundPause = pause 51 } 52 53 func foreground(_ input: TeraSubmissionMediaRequest) async throws -> TeraSubmissionStatus { 54 let value = try status(input.request) 55 guard !value.delivery.isStopped, value.revision == input.expectedRevision else { throw failure("submission_stopped") } 56 foregroundCount += 1 57 await foregroundPause?.wait() 58 let current = try replacing(status(input.request), state: .readyToSign, media: value.media.map { 59 $0.opaqueReference == input.media.media.opaqueReference 60 ? TeraSubmissionMedia(opaqueReference: $0.opaqueReference, progress: progress(input.media.media, stage: .verified)) : $0 61 }) 62 operations[input.request.commandID] = current 63 completionPersisted = true 64 return current 65 } 66 67 func unreadable(_ value: Bool) { 68 failReads = value 69 } 70 71 /// Injects saved facts for presentation tests; it does not simulate policy. 72 func installPresentationFixture(_ value: TeraSubmissionStatus) { 73 precondition(operations[value.request.commandID]?.request == value.request) 74 operations[value.request.commandID] = value 75 } 76 77 func reserveID() -> String { 78 idCount += 1 79 return String(format: "%032x", 1000 + idCount) 80 } 81 82 func reserve(_ request: TeraSubmissionRequest) async throws -> TeraSubmissionReservation { 83 if let existing = reservations[request.commandID] { 84 guard existing.captured.scope == request.scope, existing.captured.id == request.composerID, 85 existing.captured.revision == request.expectedRevision else { throw failure("idempotency_conflict") } 86 return existing 87 } 88 let source = try await composer.load(request.scope, id: request.composerID) 89 guard source.revision == request.expectedRevision else { throw failure("composer_revision_conflict") } 90 let value = TeraSubmissionReservation(commandID: request.commandID, reservationID: request.commandID, 91 captured: source, reservedAtUnixMilliseconds: 1_800_000_000_000, replayed: false) 92 reservations[request.commandID] = value 93 return value 94 } 95 96 func recover(_ request: TeraSubmissionRequest) throws -> TeraSubmissionStatus? { 97 if failReads { 98 throw failure("storage_unavailable") 99 } 100 guard let value = operations[request.commandID] else { return nil } 101 guard value.request == request else { throw failure("idempotency_conflict") } 102 return value 103 } 104 105 func status(_ request: TeraSubmissionRequest) throws -> TeraSubmissionStatus { 106 guard let value = try recover(request) else { throw failure("submission_not_found") } 107 return value 108 } 109 110 func prepare(_ request: TeraSubmissionRequest, media: [TeraPreparedMediaHandle]) async throws -> TeraSubmissionStatus { 111 if let value = try recover(request) { 112 return value 113 } 114 let reservation = try await reserve(request) 115 guard writable else { throw failure("submission_policy_unavailable") } 116 guard media.map(\.media.opaqueReference) == reservation.captured.form.media.map(\.opaqueReference) else { 117 throw failure("submission_media_invalid") 118 } 119 prepareCount += 1 120 let initial = TeraSubmissionStatus( 121 request: request, intentID: String(format: "%032x", 2000 + prepareCount), 122 operationID: String(format: "%032x", 3000 + prepareCount), revision: 1, captured: reservation.captured, 123 state: media.isEmpty ? .readyToSign : .mediaPreparing, 124 committedAtUnixMilliseconds: 1_800_000_000_000, updatedAtUnixMilliseconds: 1_800_000_000_000, 125 media: media.map { TeraSubmissionMedia(opaqueReference: $0.media.opaqueReference, 126 progress: progress($0.media, stage: .pending)) }, settlement: settlement(complete: false), 127 delivery: TeraPublicationEvidence(state: .notIssued, stopRequestedAtUnixMilliseconds: nil, 128 schedulingRevision: 1, retainedFacts: 0, recordedAttempts: 0, unresolvedClaims: false), 129 targetDetails: .fixture(), retry: .ready 130 ) 131 if delayedPhase == .queue, !delayed { 132 delayed = true 133 try await Task.sleep(for: .milliseconds(50)) 134 } 135 operations[request.commandID] = initial 136 if let pause = preparePause { 137 preparePause = nil; await pause.wait() 138 } 139 if failAfterCommit { 140 failAfterCommit = false; throw failure("storage_unavailable") 141 } 142 return initial 143 } 144 145 func advance(_ request: TeraSubmissionRequest, revision: UInt64) async throws -> TeraSubmissionStatus { 146 let value = try status(request) 147 guard !value.delivery.isStopped else { throw failure("submission_stopped") } 148 guard value.revision == revision else { throw failure("draft_revision_conflict") } 149 advanceCount += 1 150 let queued = replacing(value, state: .queued) 151 operations[request.commandID] = queued 152 if let pause = advancePause { 153 advancePause = nil; await pause.wait() 154 } 155 if delayedPhase == .advance, !delayed { 156 delayed = true 157 try await Task.sleep(for: .milliseconds(50)) 158 } 159 if offline { 160 throw failure("relay_offline") 161 } 162 let complete = try replacing(status(request), state: .complete) 163 operations[request.commandID] = complete 164 return complete 165 } 166 167 func upload(_ input: TeraSubmissionMediaRequest) throws -> TeraSubmissionUploadJob { 168 let value = try status(input.request) 169 guard !value.delivery.isStopped else { throw failure("submission_stopped") } 170 guard value.revision == input.expectedRevision else { throw failure("draft_revision_conflict") } 171 uploadCount += 1 172 let current = replacing(value, state: .mediaUploading, media: value.media.map { 173 $0.opaqueReference == input.media.media.opaqueReference 174 ? TeraSubmissionMedia(opaqueReference: $0.opaqueReference, progress: progress(input.media.media, stage: .uploading)) : $0 175 }) 176 operations[input.request.commandID] = current 177 return TeraSubmissionUploadJob(submission: current, transfer: TeraNativeTransferJob( 178 ownerID: current.intentID, expectedRevision: current.revision, 179 operationID: String(format: "%032x", 4000 + uploadCount), 180 remoteURL: "http://127.0.0.1:3000/\(input.media.media.sha256).png", uploadURL: "http://127.0.0.1:3000/upload", 181 authorizationHeader: "Nostr test-token", 182 expectedSHA256: input.media.media.sha256, mediaType: input.media.media.mediaType, byteSize: input.media.media.byteSize 183 )) 184 } 185 186 func complete(_ input: TeraSubmissionMediaRequest, response: TeraAddBackgroundUploadReceipt) async throws -> TeraSubmissionStatus { 187 let value = try status(input.request) 188 guard value.revision == input.expectedRevision, response.expectedRevision == value.revision, 189 response.draftID == value.intentID else { throw failure("submission_media_invalid") } 190 let current = replacing(value, state: .readyToSign, media: value.media.map { 191 $0.opaqueReference == input.media.media.opaqueReference 192 ? TeraSubmissionMedia(opaqueReference: $0.opaqueReference, progress: progress(input.media.media, stage: .verified)) : $0 193 }) 194 operations[input.request.commandID] = current 195 completionPersisted = true 196 if delayAfterCompletion { 197 try await Task.sleep(for: .milliseconds(50)) 198 } 199 return current 200 } 201 202 func page(scope: TeraComposerScope, limit: UInt16, cursor: String?) throws -> TeraSubmissionPage { 203 if failReads { 204 throw failure("storage_unavailable") 205 } 206 let matches = reservations.values.filter { $0.captured.scope == scope && $0.commandID > (cursor ?? "") } 207 .sorted { $0.commandID < $1.commandID } 208 let entries = matches.prefix(Int(limit)).map { value in 209 let request = TeraSubmissionRequest(commandID: value.commandID, scope: scope, 210 composerID: value.captured.id, expectedRevision: value.captured.revision) 211 let state = operations[value.commandID].map { 212 TeraSubmissionSummaryState.operation(intentID: $0.intentID, operationID: $0.operationID, revision: $0.revision, state: $0.state) 213 } ?? .reserved 214 return TeraSubmissionEntry.submission(TeraSubmissionSummary(request: request, reservationID: value.reservationID, 215 reservedAtUnixMilliseconds: value.reservedAtUnixMilliseconds, state: state)) 216 } 217 return TeraSubmissionPage(scope: scope, entries: entries, 218 nextCursor: matches.count > Int(limit) ? entries.last?.id : nil) 219 } 220 221 private func replacing(_ value: TeraSubmissionStatus, state: TeraOutboxState, media: [TeraSubmissionMedia]? = nil) -> TeraSubmissionStatus { 222 TeraSubmissionStatus(request: value.request, intentID: value.intentID, operationID: value.operationID, 223 revision: value.revision + 1, captured: value.captured, state: state, 224 committedAtUnixMilliseconds: value.committedAtUnixMilliseconds, 225 updatedAtUnixMilliseconds: value.updatedAtUnixMilliseconds + 1, media: media ?? value.media, 226 settlement: settlement(complete: state == .complete), 227 delivery: TeraPublicationEvidence(state: state == .complete ? .accepted : value.delivery.state, 228 stopRequestedAtUnixMilliseconds: value.delivery.stopRequestedAtUnixMilliseconds, 229 schedulingRevision: value.delivery.schedulingRevision + 1, 230 retainedFacts: state == .complete ? 1 : value.delivery.retainedFacts, 231 recordedAttempts: state == .complete ? 1 : value.delivery.recordedAttempts, unresolvedClaims: false), 232 targetDetails: state == .complete ? .fixture(accepted: true) : value.targetDetails, 233 retry: state == .complete ? .complete : value.retry) 234 } 235 236 func requestStop(_ request: TeraSubmissionRequest) throws -> TeraSubmissionStatus { 237 let value = try status(request) 238 if value.delivery.isStopped { 239 return value 240 } 241 let stopped = TeraSubmissionStatus(request: value.request, intentID: value.intentID, operationID: value.operationID, 242 revision: value.revision, captured: value.captured, state: value.state == .complete ? .complete : .cancelled, 243 committedAtUnixMilliseconds: value.committedAtUnixMilliseconds, updatedAtUnixMilliseconds: value.updatedAtUnixMilliseconds, 244 media: value.media, settlement: value.settlement, 245 delivery: TeraPublicationEvidence(state: value.delivery.state, stopRequestedAtUnixMilliseconds: 1_800_000_000_100, 246 schedulingRevision: value.delivery.schedulingRevision + 1, retainedFacts: value.delivery.retainedFacts, 247 recordedAttempts: value.delivery.recordedAttempts, unresolvedClaims: value.delivery.unresolvedClaims), 248 targetDetails: value.targetDetails, retry: .stopped) 249 operations[request.commandID] = stopped 250 return stopped 251 } 252 253 private func progress(_ media: TeraPreparedMedia, stage: TeraDraftMediaStage) -> TeraDraftMediaStatus { 254 TeraDraftMediaStatus(url: "http://127.0.0.1:3000/\(media.sha256).png", stage: stage, uploadAttempts: stage == .pending ? 0 : 1, 255 verifiedAtUnixMilliseconds: stage == .verified ? 1_800_000_000_000 : nil, 256 possibleOrphan: false, orphanReasonCode: nil, orphanRecordedAtUnixMilliseconds: nil) 257 } 258 259 private func settlement(complete: Bool) -> TeraOperationSettlement { 260 TeraOperationSettlement(artifacts: 1, signed: complete ? 1 : 0, admitted: complete ? 1 : 0, 261 pending: complete ? 0 : 1, retryable: 0, indeterminate: 0, failedTerminal: 0, cancelled: 0, 262 deliveryPlans: 1, deliverySatisfied: complete ? 1 : 0, deliveryPending: complete ? 0 : 1, 263 deliveryRetryable: 0, deliveryExhausted: 0, deliveryFailedTerminal: 0, deliveryCancelled: 0) 264 } 265 266 private func failure(_ code: String) -> TeraRuntimeFailure { 267 .local(operation: "test.submission", code: code, safeMessage: "Submission needs attention.") 268 } 269 } 270 271 extension AddBackend { 272 func requestSubmissionStop(request: TeraSubmissionRequest) async throws -> TeraSubmissionStatus { 273 try await submissionBackend.requestStop(request) 274 } 275 276 func reconcileSubmissionLocal(request: TeraSubmissionRequest, context _: TeraLocalNetwork) async throws -> TeraSubmissionStatus { 277 try await submissionBackend.status(request) 278 } 279 280 func reserveSubmissionID() async -> String { 281 await submissionBackend.reserveID() 282 } 283 284 func reserveSubmission(request: TeraSubmissionRequest) async throws -> TeraSubmissionReservation { 285 try await submissionBackend.reserve(request) 286 } 287 288 func prepareSubmission(request: TeraSubmissionRequest, media: [TeraPreparedMediaHandle]) async throws -> TeraSubmissionStatus { 289 try await submissionBackend.prepare(request, media: media) 290 } 291 292 func recoverSubmission(request: TeraSubmissionRequest) async throws -> TeraSubmissionStatus? { 293 try await submissionBackend.recover(request) 294 } 295 296 func submissionStatus(request: TeraSubmissionRequest) async throws -> TeraSubmissionStatus { 297 try await submissionBackend.status(request) 298 } 299 300 func advanceSubmission(request: TeraSubmissionRequest, expectedRevision: UInt64) async throws -> TeraSubmissionStatus { 301 try await submissionBackend.advance(request, revision: expectedRevision) 302 } 303 304 func listSubmissions(scope: TeraComposerScope, limit: UInt16, cursor: String?) async throws -> TeraSubmissionPage { 305 try await submissionBackend.page(scope: scope, limit: limit, cursor: cursor) 306 } 307 308 func prepareSubmissionUpload(input: TeraSubmissionMediaRequest) async throws -> TeraSubmissionUploadJob { 309 try await submissionBackend.upload(input) 310 } 311 312 func uploadSubmissionMedia(input: TeraSubmissionMediaRequest) async throws -> TeraSubmissionStatus { 313 try await submissionBackend.foreground(input) 314 } 315 316 func completeSubmissionUpload(input: TeraSubmissionMediaRequest, response: TeraAddBackgroundUploadReceipt) async throws -> TeraSubmissionStatus { 317 try await submissionBackend.complete(input, response: response) 318 } 319 }