field_ios

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

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 }