TeraObservationCoalescingTests.swift (6918B)
1 @testable import TeraApp 2 import XCTest 3 4 @MainActor 5 final class TeraObservationCoalescingTests: XCTestCase { 6 func testRecoveryResnapshotsWithoutAnInitialHintOrAnotherEvent() async throws { 7 let backend = try TeraScopeBackend() 8 let client = try await TeraScopeFixtures.client(backend) 9 let observation = TeraStoreObservation() 10 let retry = ResourceTestPause() 11 var batches: [TeraObservationBatch] = [] 12 observation.start( 13 client: client, buffer: (capacity: 8, delay: { _ in await retry.wait() }), state: { _ in }, 14 accepts: { _ in true }, refresh: { batches.append($0) } 15 ) 16 await TeraScopeFixtures.eventually { batches.count == 1 } 17 await client.suspend() 18 await retry.entered.wait() 19 XCTAssertEqual(batches.count, 1) 20 await retry.resume.open() 21 await TeraScopeFixtures.eventually { batches.count == 2 } 22 XCTAssertTrue(batches.allSatisfy(\.resnapshot)) 23 let subscriptions = await backend.counts[.subscribe] 24 XCTAssertEqual(subscriptions, 2) 25 observation.stop() 26 _ = try await client.stop() 27 } 28 29 func testStormCoalescesFixedDomainsAndDisplaysOneFinalRevision() async throws { 30 let backend = try TeraScopeBackend() 31 let client = try await TeraScopeFixtures.client(backend) 32 let observation = TeraStoreObservation() 33 let firstRead = ResourceTestPause() 34 let kinds: [TeraRuntimeChangeKind] = [.today, .drafts, .media, .settings, .identity, .profile, .relay, .lifecycle] 35 var batches: [TeraObservationBatch] = [] 36 let durable = ObservationTestState() 37 var visibleRevisions: [Int] = [] 38 var seen: UInt64 = 0 39 observation.start( 40 client: client, buffer: (capacity: 8, delay: { _ in throw CancellationError() }), state: { _ in }, 41 accepts: { seen = $0.revision.rawValue ?? seen; return true }, 42 refresh: { batch in 43 batches.append(batch) 44 let captured = durable.revision 45 if batches.count == 1 { 46 await firstRead.wait() 47 } 48 visibleRevisions.append(captured) 49 } 50 ) 51 await firstRead.entered.wait() 52 for index in 1 ... 256 { 53 durable.revision = index 54 await backend.emit(kinds[(index - 1) % kinds.count]) 55 await TeraScopeFixtures.eventually { seen == UInt64(index) } 56 } 57 XCTAssertEqual(batches.count, 1) 58 await firstRead.resume.open() 59 await TeraScopeFixtures.eventually { visibleRevisions.count == 2 } 60 XCTAssertEqual(visibleRevisions, [0, 256]) 61 XCTAssertEqual(try XCTUnwrap(batches.last).domains, Set(kinds)) 62 XCTAssertFalse(try XCTUnwrap(batches.last).resnapshot) 63 observation.stop() 64 _ = try await client.stop() 65 } 66 67 func testGapOrExhaustedRevisionForcesFinalSnapshotAfterSilence() async throws { 68 for exhausted in [false, true] { 69 let backend = try TeraScopeBackend() 70 let client = try await TeraScopeFixtures.client(backend) 71 let observation = TeraStoreObservation() 72 let first = ResourceTestPause() 73 var seen = false 74 var batches: [TeraObservationBatch] = [] 75 observation.start( 76 client: client, buffer: (capacity: 1, delay: { _ in throw CancellationError() }), state: { _ in }, 77 accepts: { _ in seen = true; return true }, refresh: { 78 batches.append($0) 79 if batches.count == 1 { 80 await first.wait() 81 } 82 } 83 ) 84 await first.entered.wait() 85 await backend.emit(.initial, delivery: exhausted ? .change : .resnapshotRequired, exhausted: exhausted) 86 await TeraScopeFixtures.eventually { seen } 87 XCTAssertEqual(batches.count, 1) 88 await first.resume.open() 89 await TeraScopeFixtures.eventually { batches.count == 2 } 90 XCTAssertTrue(try XCTUnwrap(batches.last).resnapshot) 91 XCTAssertTrue(try XCTUnwrap(batches.last).domains.isEmpty) 92 observation.stop() 93 _ = try await client.stop() 94 } 95 } 96 97 func testOldRefreshCompletionCannotClearReplacementOrItsPendingWork() async throws { 98 let backend = try TeraScopeBackend() 99 let client = try await TeraScopeFixtures.client(backend) 100 let observation = TeraStoreObservation() 101 let oldRead = ResourceTestPause() 102 let newRead = ResourceTestPause() 103 let oldFinished = ResourceTestGate() 104 observation.start( 105 client: client, buffer: (capacity: 8, delay: { _ in throw CancellationError() }), state: { _ in }, 106 accepts: { _ in true }, refresh: { _ in await oldRead.wait(); await oldFinished.open() } 107 ) 108 await oldRead.entered.wait() 109 observation.stop() 110 var seen = false 111 var newBatches: [TeraObservationBatch] = [] 112 observation.start( 113 client: client, buffer: (capacity: 8, delay: { _ in throw CancellationError() }), state: { _ in }, 114 accepts: { _ in seen = true; return true }, refresh: { 115 newBatches.append($0) 116 if newBatches.count == 1 { 117 await newRead.wait() 118 } 119 } 120 ) 121 await newRead.entered.wait() 122 await oldRead.resume.open() 123 await oldFinished.wait() 124 await backend.emit(.drafts) 125 await TeraScopeFixtures.eventually { seen } 126 XCTAssertEqual(newBatches.count, 1) 127 await newRead.resume.open() 128 await TeraScopeFixtures.eventually { newBatches.count == 2 } 129 XCTAssertEqual(try XCTUnwrap(newBatches.last).domains, [.drafts]) 130 observation.stop() 131 _ = try await client.stop() 132 } 133 134 func testOldCompletionCannotOrphanReplacementRefreshFromStop() async throws { 135 let backend = try TeraScopeBackend() 136 let client = try await TeraScopeFixtures.client(backend) 137 let observation = TeraStoreObservation() 138 let oldRead = ResourceTestPause() 139 let newRead = ResourceTestPause() 140 let newFinished = ResourceTestGate() 141 weak var oldLifetime: ObservationTaskLifetime? 142 do { 143 let lifetime = ObservationTaskLifetime() 144 oldLifetime = lifetime 145 observation.start( 146 client: client, buffer: (capacity: 8, delay: { _ in throw CancellationError() }), state: { _ in }, 147 accepts: { _ in true }, refresh: { [lifetime] _ in 148 await oldRead.wait() 149 withExtendedLifetime(lifetime) {} 150 } 151 ) 152 } 153 await oldRead.entered.wait() 154 observation.stop() 155 var cancelled: Bool? 156 observation.start( 157 client: client, buffer: (capacity: 8, delay: { _ in throw CancellationError() }), state: { _ in }, 158 accepts: { _ in true }, refresh: { _ in 159 await newRead.wait() 160 cancelled = Task.isCancelled 161 await newFinished.open() 162 } 163 ) 164 await newRead.entered.wait() 165 await oldRead.resume.open() 166 await TeraScopeFixtures.eventually { oldLifetime == nil } 167 observation.stop() 168 await newRead.resume.open() 169 await newFinished.wait() 170 XCTAssertEqual(cancelled, true, "Stop must still own and cancel the replacement refresh") 171 _ = try await client.stop() 172 } 173 } 174 175 @MainActor 176 private final class ObservationTaskLifetime {} 177 178 @MainActor 179 private final class ObservationTestState { 180 var revision = 0 181 }