TeraRuntimeBackpressureTests.swift (6214B)
1 @testable import TeraApp 2 import TeraKitBindings 3 import XCTest 4 5 final class TeraRuntimeBackpressureTests: XCTestCase { 6 func testIndependentSubscriptionsUseBoundedNewestBuffers() async throws { 7 let harness = RuntimeHarness() 8 let client = TeraRuntimeClient(factory: harness.start) 9 _ = try await client.start(configuration: TeraRuntimeClientTests().makeConfiguration(generation: "02")) 10 let first = try await client.changes(bufferCapacity: 2) 11 let second = try await client.changes(bufferCapacity: 4) 12 13 for generation in 1 ... 10 { 14 await harness.emitRevision(UInt64(generation)) 15 } 16 17 _ = try await client.stop() 18 var firstValues: [TeraRuntimeChange] = [] 19 var secondValues: [TeraRuntimeChange] = [] 20 for await value in first { 21 firstValues.append(value) 22 } 23 for await value in second { 24 secondValues.append(value) 25 } 26 XCTAssertEqual(firstValues.count, 2) 27 XCTAssertEqual(secondValues.count, 4) 28 for values in [firstValues, secondValues] { 29 XCTAssertEqual(values.last?.delivery, .resnapshotRequired) 30 XCTAssertEqual(values.last?.scope.context, nil) 31 XCTAssertEqual(values.last(where: { $0.delivery == .change })?.revision.rawValue, 10) 32 } 33 let cancelCount = await harness.cancelCount() 34 XCTAssertEqual(cancelCount, 2) 35 } 36 37 func testGeneratedCallbackRetainsFinalGapAtEverySupportedBufferBoundary() async { 38 for capacity in [1, 2, 16, 64] { 39 for overflow in [false, true] { 40 let pair = AsyncStream.makeStream(of: TeraRuntimeChange.self, bufferingPolicy: .bufferingNewest(capacity)) 41 let observer = TeraGeneratedRuntimeObserver(continuation: pair.continuation) 42 for revision in 1 ... (capacity + (overflow ? 1 : 0)) { 43 observer.onChange(change: generated(revision: UInt64(revision))) 44 } 45 observer.finish() 46 var values: [TeraRuntimeChange] = [] 47 for await value in pair.stream { 48 values.append(value) 49 } 50 XCTAssertEqual(values.count, capacity) 51 if overflow { 52 XCTAssertEqual(values.last, generated(revision: 1).appValue.requiringResnapshot()) 53 } else { 54 XCTAssertTrue(values.allSatisfy { $0.delivery == .change }) 55 XCTAssertEqual(values.map(\.revision.rawValue), (1 ... capacity).map { UInt64($0) }) 56 } 57 } 58 } 59 } 60 61 func testRepeatedGapDoesNotAdvanceDomainWatermarksOrAcceptAnotherEpoch() async throws { 62 let harness = RuntimeHarness() 63 let client = TeraRuntimeClient(factory: harness.start) 64 _ = try await client.start(configuration: TeraRuntimeClientTests().makeConfiguration(generation: "66")) 65 let stream = try await client.changes() 66 await harness.emit(generated(revision: 5).appValue) 67 let gap = generated(revision: 99).appValue.requiringResnapshot() 68 await harness.emit(gap) 69 await harness.emit(gap) 70 await harness.emit(generated(revision: 6).appValue) 71 var otherEpoch = generated(revision: 100) 72 otherEpoch.epoch = String(repeating: "2", count: 32) 73 await harness.emit(otherEpoch.appValue.requiringResnapshot()) 74 await harness.emit(generated(revision: 7).appValue) 75 _ = try await client.stop() 76 var values: [TeraRuntimeChange] = [] 77 for await value in stream { 78 values.append(value) 79 } 80 XCTAssertEqual(values.map(\.delivery), [.change, .resnapshotRequired, .resnapshotRequired, .change, .change]) 81 XCTAssertEqual(values.filter { $0.delivery == .change }.map(\.revision.rawValue), [5, 6, 7]) 82 } 83 84 func testOldScopeGapCannotCrossRuntimeReplacement() async throws { 85 let harness = RuntimeHarness() 86 let client = TeraRuntimeClient(factory: harness.start) 87 _ = try await client.start(configuration: TeraRuntimeClientTests().makeConfiguration(generation: "66")) 88 let first = try await client.changes() 89 await harness.emit(generated(revision: 5).appValue) 90 _ = try await client.start(configuration: TeraRuntimeClientTests().makeConfiguration(generation: "77")) 91 let second = try await client.changes() 92 await harness.emit(generated(revision: 99).appValue.requiringResnapshot()) 93 let current = generated(revision: 1, account: "77").appValue 94 await harness.emit(current) 95 _ = try await client.stop() 96 var firstValues: [TeraRuntimeChange] = [] 97 var secondValues: [TeraRuntimeChange] = [] 98 for await value in first { 99 firstValues.append(value) 100 } 101 for await value in second { 102 secondValues.append(value) 103 } 104 XCTAssertEqual(firstValues, [generated(revision: 5).appValue]) 105 XCTAssertEqual(secondValues, [current]) 106 } 107 108 func testGeneratedDeliveryVariantsRoundTripWithoutChangingScopeOrEpoch() throws { 109 for delivery in [FfiRuntimeChangeDelivery.change, .resnapshotRequired] { 110 var original = generated(revision: .max) 111 original.delivery = delivery 112 let decoded = try FfiConverterTypeFfiRuntimeChangeRecord_lift(FfiConverterTypeFfiRuntimeChangeRecord_lower(original)) 113 XCTAssertEqual(decoded, original) 114 XCTAssertEqual(decoded.appValue.delivery, delivery == .change ? .change : .resnapshotRequired) 115 let gap = decoded.appValue.requiringResnapshot() 116 XCTAssertEqual(gap.scope.publicKey, original.scope.publicKey) 117 XCTAssertEqual(gap.scope.sourceGeneration, original.scope.sourceGeneration) 118 XCTAssertEqual(gap.epoch, original.epoch) 119 XCTAssertNil(gap.scope.context) 120 XCTAssertNil(gap.entityID) 121 XCTAssertEqual(gap.kind, .initial) 122 } 123 } 124 125 private func generated(revision: UInt64, account: String = "66") -> FfiRuntimeChangeRecord { 126 FfiRuntimeChangeRecord( 127 schemaVersion: 3, 128 scope: FfiRuntimeChangeScope( 129 publicKey: String(repeating: account, count: 32), sourceGeneration: String(repeating: account, count: 32), 130 context: FfiLocalNetworkRecord(schemaVersion: 1, id: "local", label: "Local network", relayUrls: ["wss://relay.example"], locality: nil, followedAuthors: [], generation: 1) 131 ), 132 epoch: String(repeating: "1", count: 32), revision: .current(value: revision), delivery: .change, kind: .today, entityId: "card" 133 ) 134 } 135 }