TeraMutationAdmissionTests.swift (8895B)
1 import RadrootsKit 2 @testable import TeraApp 3 import XCTest 4 5 @MainActor 6 final class TeraMutationAdmissionTests: XCTestCase { 7 func testUploadReservesBeforeDiscoveryAndKeepsOtherDraftsIndependent() async throws { 8 let fixture = try BackgroundUploadFixture() 9 let other = try BackgroundUploadFixture(draftID: String(repeating: "2", count: 32)) 10 defer { fixture.remove(); other.remove() } 11 let transfer = BackgroundTransferHarness() 12 let coordinator = fixture.coordinator(transfer: transfer) 13 let pause = ResourceTestPause() 14 await transfer.pauseDiscovery(pause) 15 let job = fixture.job(revision: 2, operation: String(repeating: "a", count: 32)) 16 let owner = Task { try await coordinator.uploadInBackground(job: job, media: fixture.media) } 17 await entered(pause) 18 await assertBusy(coordinator, job: job, media: fixture.media) 19 await assertBusy( 20 coordinator, 21 job: fixture.job(revision: 3, operation: String(repeating: "b", count: 32)), 22 media: fixture.media 23 ) 24 let independent = try await coordinator.uploadInBackground( 25 job: other.job(revision: 2, operation: String(repeating: "c", count: 32)), 26 media: other.media 27 ) 28 XCTAssertEqual(independent.draftID, other.draftID) 29 let discoveries = await transfer.discoveryCount 30 XCTAssertEqual(discoveries, 2) 31 await pause.resume.open() 32 let receipt = try await owner.value 33 XCTAssertEqual(receipt.draftID, fixture.draftID) 34 let counts = await transfer.counts 35 XCTAssertEqual(counts.enqueue, 2) 36 XCTAssertEqual(counts.retry, 0) 37 } 38 39 func testRetryRetainsAdmissionAcrossTheNativeCallback() async throws { 40 let fixture = try BackgroundUploadFixture() 41 defer { fixture.remove() } 42 let transfer = BackgroundTransferHarness() 43 let coordinator = fixture.coordinator(transfer: transfer) 44 let job = fixture.job(revision: 2, operation: String(repeating: "a", count: 32)) 45 try await transfer.seed(request: fixture.request(job: job), state: .interrupted) 46 let pause = ResourceTestPause() 47 await transfer.pauseRetry(pause) 48 let owner = Task { try await coordinator.uploadInBackground(job: job, media: fixture.media) } 49 await entered(pause) 50 await assertBusy(coordinator, job: job, media: fixture.media) 51 await pause.resume.open() 52 _ = try await owner.value 53 let counts = await transfer.counts 54 XCTAssertEqual(counts.retry, 1) 55 XCTAssertEqual(counts.enqueue, 0) 56 } 57 58 func testCancellationBeforeAdmissionDoesNotEnterNativeWorkOrRetainTheScope() async throws { 59 let fixture = try BackgroundUploadFixture() 60 defer { fixture.remove() } 61 let transfer = BackgroundTransferHarness() 62 let coordinator = fixture.coordinator(transfer: transfer) 63 let job = fixture.job(revision: 2, operation: String(repeating: "a", count: 32)) 64 let queued = ResourceTestPause() 65 let waiter = Task { 66 await queued.wait() 67 return try await coordinator.uploadInBackground(job: job, media: fixture.media) 68 } 69 await entered(queued) 70 waiter.cancel() 71 await queued.resume.open() 72 do { 73 _ = try await waiter.value 74 XCTFail("A cancelled queued caller must not enter native work") 75 } catch is CancellationError {} 76 let discoveries = await transfer.discoveryCount 77 XCTAssertEqual(discoveries, 0) 78 _ = try await coordinator.uploadInBackground(job: job, media: fixture.media) 79 let counts = await transfer.counts 80 XCTAssertEqual(counts.enqueue, 1) 81 } 82 83 func testCancelledOwnerKeepsAdmissionUntilItsNativeWaitReturns() async throws { 84 let fixture = try BackgroundUploadFixture() 85 defer { fixture.remove() } 86 let transfer = BackgroundTransferHarness() 87 let coordinator = fixture.coordinator(transfer: transfer) 88 let job = fixture.job(revision: 2, operation: String(repeating: "a", count: 32)) 89 let pause = ResourceTestPause() 90 await transfer.pauseDiscovery(pause) 91 let owner = Task { try await coordinator.uploadInBackground(job: job, media: fixture.media) } 92 await entered(pause) 93 owner.cancel() 94 await assertBusy(coordinator, job: job, media: fixture.media) 95 await pause.resume.open() 96 do { 97 _ = try await owner.value 98 XCTFail("Expected cancellation before enqueue") 99 } catch is CancellationError {} 100 var counts = await transfer.counts 101 XCTAssertEqual(counts.enqueue, 0) 102 _ = try await coordinator.uploadInBackground(job: job, media: fixture.media) 103 counts = await transfer.counts 104 XCTAssertEqual(counts.enqueue, 1) 105 } 106 107 func testReentrantNativeCallbackCannotReadmitTheSameDraft() async throws { 108 let fixture = try BackgroundUploadFixture() 109 defer { fixture.remove() } 110 let transfer = BackgroundTransferHarness() 111 let coordinator = fixture.coordinator(transfer: transfer) 112 let job = fixture.job(revision: 2, operation: String(repeating: "a", count: 32)) 113 let callback = expectation(description: "Reentrant callback returned") 114 await transfer.onDiscovery { 115 do { 116 _ = try await coordinator.uploadInBackground(job: job, media: fixture.media) 117 XCTFail("A callback must not acquire its caller's scope") 118 } catch let failure as TeraRuntimeFailure { 119 XCTAssertEqual(failure.code, "operation_in_progress") 120 } catch { 121 XCTFail("Unexpected reentrant failure") 122 } 123 callback.fulfill() 124 } 125 _ = try await coordinator.uploadInBackground(job: job, media: fixture.media) 126 await fulfillment(of: [callback], timeout: 2) 127 let counts = await transfer.counts 128 XCTAssertEqual(counts.enqueue, 1) 129 let discoveries = await transfer.discoveryCount 130 XCTAssertEqual(discoveries, 1) 131 } 132 133 private func assertBusy( 134 _ coordinator: TeraAddMediaCoordinator, 135 job: TeraNativeUploadJob, 136 media: TeraPreparedMedia 137 ) async { 138 do { 139 _ = try await coordinator.uploadInBackground(job: job, media: media) 140 XCTFail("Overlapping calls must preserve the admitted operation") 141 } catch let failure as TeraRuntimeFailure { 142 XCTAssertEqual(failure.code, "operation_in_progress") 143 XCTAssertEqual(failure.category, "operation") 144 XCTAssertTrue(failure.retryable) 145 XCTAssertEqual(failure.recoveryActions, ["retry_operation_with_same_idempotency_key"]) 146 } catch { 147 XCTFail("Unexpected admission failure") 148 } 149 } 150 151 private func entered(_ pause: ResourceTestPause) async { 152 let arrived = expectation(description: "Entered the explicit native pause") 153 let observer = Task { await pause.entered.wait(); arrived.fulfill() } 154 await fulfillment(of: [arrived], timeout: 2) 155 await pause.entered.open() 156 await observer.value 157 } 158 } 159 160 extension TeraMutationAdmissionTests { 161 @MainActor 162 func testProfileAdmissionSurvivesStopUntilTheAdmittedCallReturns() async throws { 163 let configuration = TeraRuntimeClientTests().makeConfiguration(generation: "41") 164 let backend = ResourceTestBackend(publicKeyHex: configuration.publicKeyHex) 165 let client = TeraRuntimeClient(factory: { _ in await backend.start() }) 166 _ = try await client.start(configuration: configuration) 167 let store = TeraSettingsStore(runtimeClient: client) 168 store.profileName = "Farm profile" 169 let pause = ResourceTestPause() 170 await backend.pauseProfile(pause) 171 let owner = Task { await store.saveProfile() } 172 await entered(pause) 173 store.profileName = "A conflicting profile update" 174 await store.saveProfile() 175 store.stop() 176 await store.saveProfile() 177 var count = await backend.profileMutations 178 XCTAssertEqual(count, 1) 179 await pause.resume.open() 180 await owner.value 181 XCTAssertNil(store.profileStatus) 182 await store.saveProfile() 183 count = await backend.profileMutations 184 XCTAssertEqual(count, 2) 185 XCTAssertNotNil(store.profileStatus) 186 _ = try await client.stop() 187 } 188 189 @MainActor 190 func testProfileAdvanceRejectsConflictsAndCancelledQueuedCallers() async throws { 191 let configuration = TeraRuntimeClientTests().makeConfiguration(generation: "41") 192 let backend = ResourceTestBackend(publicKeyHex: configuration.publicKeyHex) 193 let client = TeraRuntimeClient(factory: { _ in await backend.start() }) 194 _ = try await client.start(configuration: configuration) 195 let store = TeraSettingsStore(runtimeClient: client) 196 await store.saveProfile() 197 let pause = ResourceTestPause() 198 await backend.pauseProfile(pause) 199 let owner = Task { await store.advanceProfile() } 200 await entered(pause) 201 await store.cancelProfile() 202 await store.saveProfile() 203 await pause.resume.open() 204 await owner.value 205 let queued = ResourceTestPause() 206 let waiter = Task { await queued.wait(); await store.saveProfile() } 207 await entered(queued) 208 waiter.cancel() 209 await queued.resume.open() 210 await waiter.value 211 let count = await backend.profileMutations 212 XCTAssertEqual(count, 2) 213 XCTAssertFalse(store.isWorking) 214 _ = try await client.stop() 215 } 216 }