diff --git a/CHANGELOG.md b/CHANGELOG.md index be0cce2..12f8705 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -12,6 +12,8 @@ - Subjects: fix a deadlock when sending values or termination concurrently with consumer cancellation (https://github.com/sideeffect-io/AsyncExtensions/issues/52). - Subjects: preserve a shared order for concurrent sends, keep current-value/replay delivery consistent with stored state, and queue termination after accepted values without locking during delivery (https://github.com/sideeffect-io/AsyncExtensions/issues/61). Concurrent sends may return while another sender drains their queued delivery. - Subjects: unregister consumers when their last iterator copy is released, and finish and await test consumers to remove the reproduced channel/task leaks (https://github.com/sideeffect-io/AsyncExtensions/issues/43). +- Subjects: preserve the first termination and ignore later values or termination, keeping current values frozen and replay history cleared. +- Subjects: make zero-capacity replay retain no history while continuing to deliver live values to existing consumers. - Multicast: preserve upstream element order and prevent termination from overtaking values when consumers advance concurrently. - Multicast: let cancelled subscribers finish while upstream is suspended, and document reusing one multicast/share instance to avoid overlapping upstream iterators (https://github.com/sideeffect-io/AsyncExtensions/issues/31). diff --git a/README.md b/README.md index 68009ca..9a31f06 100644 --- a/README.md +++ b/README.md @@ -114,15 +114,39 @@ AsyncStream) Subjects serialize concurrent sends in the order their state lock is acquired. Which producer goes first is unspecified, but consumers registered for the same sends receive the same order, and termination follows previously accepted values. Current-value and replay state use that -same order, so consumers that remain subscribed catch up to the stored state. +same order, so consumers that remain subscribed catch up to the stored state while the subject is active. + +The first finish or failure is permanent. Later values, assignments to a current-value subject's +`value`, and further termination signals are ignored. A current-value subject keeps its last +accepted value; a replay subject clears its history on termination. Iterators created after +termination receive only that original finish or failure, without any current value or replayed history. + +Registration happens synchronously in `makeAsyncIterator()`. A passthrough iterator receives only +values accepted after registration. A current-value or replay iterator receives the stored state +at registration followed by subsequent sends. Replay capacity zero retains no history and still +delivers live values to registered consumers. Replay capacity limits history for new subscribers; +existing subscribers have unbounded buffers if they consume more slowly than values are produced. A subscription unregisters when its consuming task is cancelled or its last iterator copy is released. This also removes abandoned consumer buffers after a loop exits early. An active loop still needs an owner: retain and cancel its `Task`, or call `subject.send(.finished)` when the producer ends. Dropping a task handle or capturing an owner weakly does not stop an active loop. -State updates and subscriber registration are synchronous. Delivery runs outside the state lock +Accepted state updates and subscriber registration are synchronous. Delivery runs outside the state lock to allow cancellation handlers to send back into the subject. If another sender is already delivering, `send` queues its delivery and returns; it does not wait for consumers to receive -the value. The active sender drains pending deliveries before returning. A new current-value -or replay consumer receives the latest stored state, followed by subsequent sends. +the value. The active sender drains pending deliveries before returning. Stored state can therefore +lead delivery temporarily. Delivery buffers a value or resumes a waiting iterator; it does not +wait for the consumer's application code to process the value. + +Every change to subjects must pass the complete subject regression set and the full test suite: + +| Guarantee | Regression coverage | +| --- | --- | +| Atomic registration with values and termination | [Subject registration suites](./Tests/AsyncSubjets/) (`test_subscription_racing_*`) | +| Cancellation completes, including handlers that send back into the subject | [AsyncSubjectCancellationTests](./Tests/AsyncSubjets/AsyncSubjectCancellationTests.swift) | +| Shared concurrent order and current-value/replay consistency | [AsyncSubjectConcurrentSendOrderingTests](./Tests/AsyncSubjets/AsyncSubjectConcurrentSendOrderingTests.swift) | +| Queued sends return, accepted values precede termination, and later sends cannot replace it | [AsyncSubjectQueuedDeliveryTests](./Tests/AsyncSubjets/AsyncSubjectQueuedDeliveryTests.swift) | +| First termination is permanent, including races, and current value freezes | [AsyncSubjectTerminationTests](./Tests/AsyncSubjets/AsyncSubjectTerminationTests.swift) | +| Zero capacity delivers live values without historical replay | [AsyncReplaySubjectTests](./Tests/AsyncSubjets/AsyncReplaySubjectTests.swift) and [AsyncThrowingReplaySubjectTests](./Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift) | +| Abandoned subscriptions and ignored values release their storage | [AsyncSubjectLifetimeTests](./Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift) | ### Combiners * [`zip(_:)`](./Sources/Combiners/Zip/AsyncZipSequence.swift): Zips any number of async sequences into arrays of elements diff --git a/Sources/AsyncSubjects/AsyncCurrentValueSubject.swift b/Sources/AsyncSubjects/AsyncCurrentValueSubject.swift index b2135bb..2b47511 100644 --- a/Sources/AsyncSubjects/AsyncCurrentValueSubject.swift +++ b/Sources/AsyncSubjects/AsyncCurrentValueSubject.swift @@ -5,11 +5,13 @@ // Created by Thibault Wittemberg on 07/01/2022. // -/// A n`AsyncCurrentValueSubject` is an async sequence in which one can send values over time. +/// An `AsyncCurrentValueSubject` is an async sequence in which one can send values over time. /// The current value is always accessible as an instance variable. -/// The current value is replayed in any new async for in loops. +/// The current value is replayed to new consumers while the subject is active. /// When the `AsyncCurrentValueSubject` is terminated, new consumers will /// immediately resume with this termination. +/// The first termination is permanent; subsequent values and termination are ignored. +/// The current value remains the last value accepted before termination. /// /// ``` /// let currentValue = AsyncCurrentValueSubject(1) @@ -28,9 +30,9 @@ /// /// .. later in the application flow /// -/// await currentValue.send(2) +/// currentValue.send(2) /// -/// print(currentValue.element) // will print 2 +/// print(currentValue.value) // will print 2 /// ``` public final class AsyncCurrentValueSubject: AsyncSubject where Element: Sendable { public typealias Element = Element @@ -69,6 +71,7 @@ public final class AsyncCurrentValueSubject: AsyncSubject where Element /// - Parameter element: the value to send public func send(_ element: Element) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } state.current = element let channels = Array(state.channels.values) return state.deliveries.enqueue { @@ -84,6 +87,7 @@ public final class AsyncCurrentValueSubject: AsyncSubject where Element /// - Parameter termination: The termination to finish the subject. public func send(_ termination: Termination) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } state.terminalState = termination let channels = Array(state.channels.values) state.channels.removeAll() diff --git a/Sources/AsyncSubjects/AsyncPassthroughSubject.swift b/Sources/AsyncSubjects/AsyncPassthroughSubject.swift index 8bd7dce..c3ee8ac 100644 --- a/Sources/AsyncSubjects/AsyncPassthroughSubject.swift +++ b/Sources/AsyncSubjects/AsyncPassthroughSubject.swift @@ -8,6 +8,7 @@ /// An `AsyncPassthroughSubject` is an async sequence in which one can send values over time. /// When the `AsyncPassthroughSubject` is terminated, new consumers will /// immediately resume with this termination. +/// The first termination is permanent; subsequent values and termination are ignored. /// /// ``` /// let passthrough = AsyncPassthroughSubject() @@ -54,6 +55,7 @@ public final class AsyncPassthroughSubject: AsyncSubject { /// - Parameter element: the value to send public func send(_ element: Element) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } let channels = Array(state.channels.values) return state.deliveries.enqueue { for channel in channels { @@ -68,6 +70,7 @@ public final class AsyncPassthroughSubject: AsyncSubject { /// - Parameter termination: The termination to finish the subject public func send(_ termination: Termination) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } state.terminalState = termination let channels = Array(state.channels.values) state.channels.removeAll() diff --git a/Sources/AsyncSubjects/AsyncReplaySubject.swift b/Sources/AsyncSubjects/AsyncReplaySubject.swift index 2f7f7f0..c3d92a8 100644 --- a/Sources/AsyncSubjects/AsyncReplaySubject.swift +++ b/Sources/AsyncSubjects/AsyncReplaySubject.swift @@ -8,8 +8,10 @@ /// An `AsyncReplaySubject` is an async sequence in which one can send values over time. /// Values are buffered in a FIFO fashion so they can be replayed by new consumers. /// When the `bufferSize` is outreached the oldest value is dropped. +/// A buffer size of zero retains no history and still delivers values to existing consumers. /// When the `AsyncReplaySubject` is terminated, new consumers will -/// immediately resume with this termination, whether it is a finish or a failure. +/// immediately finish without replaying buffered values. +/// The first termination is permanent; subsequent values and termination are ignored. /// /// ``` /// let replay = AsyncReplaySubject(bufferSize: 3) @@ -48,10 +50,13 @@ public final class AsyncReplaySubject: AsyncSubject where Element: Send /// - Parameter element: the value to send public func send(_ element: Element) { let shouldDrain = self.state.withCriticalRegion { state in - if state.buffer.count >= state.bufferSize && !state.buffer.isEmpty { - state.buffer.removeFirst() + guard state.terminalState == nil else { return false } + if state.bufferSize > 0 { + if state.buffer.count >= state.bufferSize { + state.buffer.removeFirst() + } + state.buffer.append(element) } - state.buffer.append(element) let channels = Array(state.channels.values) return state.deliveries.enqueue { for channel in channels { @@ -66,6 +71,7 @@ public final class AsyncReplaySubject: AsyncSubject where Element: Send /// - Parameter termination: The termination to finish the subject. public func send(_ termination: Termination) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } state.terminalState = termination let channels = Array(state.channels.values) state.channels.removeAll() diff --git a/Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift b/Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift index 483e006..21db021 100644 --- a/Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift +++ b/Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift @@ -5,11 +5,13 @@ // Created by Thibault Wittemberg on 07/01/2022. // -/// An`AsyncThrowingCurrentValueSubject` is an async sequence in which one can send values over time. +/// An `AsyncThrowingCurrentValueSubject` is an async sequence in which one can send values over time. /// The current value is always accessible as an instance variable. -/// The current value is replayed in any new async for in loops. +/// The current value is replayed to new consumers while the subject is active. /// When the `AsyncThrowingCurrentValueSubject` is terminated, new consumers will /// immediately resume with this termination, whether it is a finish or a failure. +/// The first termination is permanent; subsequent values and termination are ignored. +/// The current value remains the last value accepted before termination. /// /// ``` /// let currentValue = AsyncThrowingCurrentValueSubject(1) @@ -28,10 +30,10 @@ /// /// .. later in the application flow /// -/// await currentValue.send(2) +/// currentValue.send(2) /// -/// print(currentValue.element) // will print 2 -/// await currentValue.send(.failure(error)) +/// print(currentValue.value) // will print 2 +/// currentValue.send(.failure(error)) /// /// ``` public final class AsyncThrowingCurrentValueSubject: AsyncSubject where Element: Sendable { @@ -69,6 +71,7 @@ public final class AsyncThrowingCurrentValueSubject: As /// - Parameter element: the value to send public func send(_ element: Element) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } state.current = element let channels = Array(state.channels.values) return state.deliveries.enqueue { @@ -84,6 +87,7 @@ public final class AsyncThrowingCurrentValueSubject: As /// - Parameter termination: The termination to finish the subject. public func send(_ termination: Termination) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } state.terminalState = termination let channels = Array(state.channels.values) state.channels.removeAll() diff --git a/Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift b/Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift index 3cc093d..6b21ebf 100644 --- a/Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift +++ b/Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift @@ -8,6 +8,7 @@ /// An `AsyncThrowingPassthroughSubject` is an async sequence in which one can send values over time. /// When the `AsyncThrowingPassthroughSubject` is terminated, new consumers will /// immediately resume with this termination, whether it is a finish or a failure. +/// The first termination is permanent; subsequent values and termination are ignored. /// /// ``` /// let passthrough = AsyncThrowingPassthroughSubject() @@ -55,6 +56,7 @@ public final class AsyncThrowingPassthroughSubject: Asy /// - Parameter element: the value to send public func send(_ element: Element) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } let channels = Array(state.channels.values) return state.deliveries.enqueue { for channel in channels { @@ -69,6 +71,7 @@ public final class AsyncThrowingPassthroughSubject: Asy /// - Parameter termination: The termination to finish the subject public func send(_ termination: Termination) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } state.terminalState = termination let channels = Array(state.channels.values) state.channels.removeAll() diff --git a/Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift b/Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift index b654875..44f9030 100644 --- a/Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift +++ b/Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift @@ -8,17 +8,21 @@ /// An `AsyncThrowingReplaySubject` is an async sequence in which one can send values over time. /// Values are buffered in a FIFO fashion so they can be replayed by new consumers. /// When the `bufferSize` is outreached the oldest value is dropped. +/// A buffer size of zero retains no history and still delivers values to existing consumers. /// When the `AsyncThrowingReplaySubject` is terminated, new consumers will /// immediately resume with this termination, whether it is a finish or a failure. +/// Buffered values are not replayed after termination. +/// The first termination is permanent; subsequent values and termination are ignored. /// /// ``` /// let replay = AsyncThrowingReplaySubject(bufferSize: 3) /// /// for i in (1...5) { replay.send(i) } -/// replay.senf(.failure(error)) +/// replay.send(.failure(error)) /// +/// // Iteration throws immediately; the buffered values are not replayed. /// for try await element in replay { -/// print(element) // will print 3, 4, 5 and throw +/// print(element) /// } /// ``` public final class AsyncThrowingReplaySubject: AsyncSubject where Element: Sendable { @@ -47,10 +51,13 @@ public final class AsyncThrowingReplaySubject: AsyncSub /// - Parameter element: the value to send public func send(_ element: Element) { let shouldDrain = self.state.withCriticalRegion { state in - if state.buffer.count >= state.bufferSize && !state.buffer.isEmpty { - state.buffer.removeFirst() + guard state.terminalState == nil else { return false } + if state.bufferSize > 0 { + if state.buffer.count >= state.bufferSize { + state.buffer.removeFirst() + } + state.buffer.append(element) } - state.buffer.append(element) let channels = Array(state.channels.values) return state.deliveries.enqueue { for channel in channels { @@ -65,6 +72,7 @@ public final class AsyncThrowingReplaySubject: AsyncSub /// - Parameter termination: The termination to finish the subject public func send(_ termination: Termination) { let shouldDrain = self.state.withCriticalRegion { state in + guard state.terminalState == nil else { return false } state.terminalState = termination let channels = Array(state.channels.values) state.channels.removeAll() diff --git a/Tests/AsyncSubjets/AsyncReplaySubjectTests.swift b/Tests/AsyncSubjets/AsyncReplaySubjectTests.swift index 5452bb8..6da40af 100644 --- a/Tests/AsyncSubjets/AsyncReplaySubjectTests.swift +++ b/Tests/AsyncSubjets/AsyncReplaySubjectTests.swift @@ -9,6 +9,23 @@ import XCTest final class AsyncReplaySubjectTests: XCTestCase { + func test_zero_capacity_does_not_replay_history_but_delivers_live_values() async { + let subject = AsyncReplaySubject(bufferSize: 0) + let existing = subject.makeAsyncIterator() + subject.send(1) + subject.send(2) + let late = subject.makeAsyncIterator() + subject.send(3) + subject.send(.finished) + + let existingResult = await drainBufferedElements(of: existing) + let lateResult = await drainBufferedElements(of: late) + XCTAssertEqual(existingResult.elements, [1, 2, 3]) + XCTAssertEqual(lateResult.elements, [3]) + XCTAssertTrue(existingResult.isTerminated) + XCTAssertTrue(lateResult.isTerminated) + } + func test_send_replays_buffered_elements() async { let exp = expectation(description: "Send has stacked elements in the replay the buffer") exp.expectedFulfillmentCount = 2 diff --git a/Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift b/Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift index b84b6b0..01ee3c0 100644 --- a/Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift +++ b/Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift @@ -8,6 +8,46 @@ private final class LifetimePayload: Sendable { } final class AsyncSubjectLifetimeTests: XCTestCase { + func test_terminated_subjects_do_not_retain_new_values() { + assertTerminatedValueReleased(AsyncPassthroughSubject()) + assertTerminatedValueReleased(AsyncCurrentValueSubject(LifetimePayload {})) + assertTerminatedValueReleased(AsyncReplaySubject(bufferSize: 2)) + for termination: Termination in [.finished, .failure(MockError(code: 1))] { + assertTerminatedValueReleased(AsyncThrowingPassthroughSubject(), termination: termination) + assertTerminatedValueReleased(AsyncThrowingCurrentValueSubject(LifetimePayload {}), termination: termination) + assertTerminatedValueReleased(AsyncThrowingReplaySubject(bufferSize: 2), termination: termination) + } + } + + private func assertTerminatedValueReleased( + _ subject: S, + termination: Termination = .finished, + file: StaticString = #filePath, + line: UInt = #line + ) where S.Element == LifetimePayload { + subject.send(termination) + assertSentValueReleased(subject, file: file, line: line) + } + + func test_zero_capacity_replay_subjects_do_not_retain_values_without_subscribers() { + assertSentValueReleased(AsyncReplaySubject(bufferSize: 0)) + assertSentValueReleased(AsyncThrowingReplaySubject(bufferSize: 0)) + } + + private func assertSentValueReleased( + _ subject: S, + file: StaticString = #filePath, + line: UInt = #line + ) where S.Element == LifetimePayload { + let released = ManagedCriticalState(false) + var payload: LifetimePayload? = LifetimePayload { released.apply(criticalState: true) } + subject.send(payload!) + payload = nil + withExtendedLifetime(subject) { + XCTAssertTrue(released.criticalState, "The subject retained an ignored value", file: file, line: line) + } + } + func test_abandoned_passthrough_iterators_release_their_buffered_values() { assertBufferedValueReleased(AsyncPassthroughSubject()) assertBufferedValueReleased(AsyncThrowingPassthroughSubject()) diff --git a/Tests/AsyncSubjets/AsyncSubjectQueuedDeliveryTests.swift b/Tests/AsyncSubjets/AsyncSubjectQueuedDeliveryTests.swift index 769379e..fe3d8aa 100644 --- a/Tests/AsyncSubjets/AsyncSubjectQueuedDeliveryTests.swift +++ b/Tests/AsyncSubjets/AsyncSubjectQueuedDeliveryTests.swift @@ -110,6 +110,9 @@ final class AsyncSubjectQueuedDeliveryTests: XCTestCase { subject.send(2) let late = subject.makeAsyncIterator() subject.send(termination) + // Rejected sends must neither replace termination nor discard earlier queued values. + subject.send(3) + subject.send(.finished) let afterTermination = subject.makeAsyncIterator() XCTAssertFalse(first.hasBufferedElements, file: file, line: line) XCTAssertFalse(second.hasBufferedElements, file: file, line: line) diff --git a/Tests/AsyncSubjets/AsyncSubjectTerminationTests.swift b/Tests/AsyncSubjets/AsyncSubjectTerminationTests.swift new file mode 100644 index 0000000..2ae0dd6 --- /dev/null +++ b/Tests/AsyncSubjets/AsyncSubjectTerminationTests.swift @@ -0,0 +1,128 @@ +@testable import AsyncExtensions +import XCTest + +final class AsyncSubjectTerminationTests: XCTestCase { + func test_throwing_passthrough_preserves_the_first_termination() async { + await assertFirstTermination(makeSubject: { AsyncThrowingPassthroughSubject() }) + } + + func test_throwing_current_value_preserves_the_first_termination() async { + await assertFirstTermination(makeSubject: { AsyncThrowingCurrentValueSubject(0) }, initial: [0]) + } + + func test_throwing_replay_preserves_the_first_termination() async { + await assertFirstTermination(makeSubject: { AsyncThrowingReplaySubject(bufferSize: 2) }) + } + + func test_throwing_passthrough_consumers_agree_on_concurrent_termination() async { + await assertConcurrentTermination(makeSubject: { AsyncThrowingPassthroughSubject() }) + } + + func test_throwing_current_value_consumers_agree_on_concurrent_termination() async { + await assertConcurrentTermination(makeSubject: { AsyncThrowingCurrentValueSubject(0) }, initial: [0]) + } + + func test_throwing_replay_consumers_agree_on_concurrent_termination() async { + await assertConcurrentTermination(makeSubject: { AsyncThrowingReplaySubject(bufferSize: 2) }) + } + + func test_current_value_ignores_sends_and_assignment_after_finish() { + let subject = AsyncCurrentValueSubject(0) + subject.send(1) + subject.send(.finished) + subject.send(2) + subject.value = 3 + XCTAssertEqual(subject.value, 1) + } + + func test_throwing_current_value_ignores_sends_and_assignment_after_termination() { + for termination: Termination in [.finished, .failure(MockError(code: 1))] { + let subject = AsyncThrowingCurrentValueSubject(0) + subject.send(1) + subject.send(termination) + subject.send(2) + subject.value = 3 + XCTAssertEqual(subject.value, 1) + } + } + + private func assertFirstTermination( + makeSubject: () -> Subject, + initial: [Int] = [], + file: StaticString = #filePath, + line: UInt = #line + ) async where Subject.Element == Int, Subject.Failure == Error { + let firstError = MockError(code: 1) + let secondError = MockError(code: 2) + let cases: [(Termination, Termination, MockError?)] = [ + (.finished, .finished, nil), + (.finished, .failure(secondError), nil), + (.failure(firstError), .finished, firstError), + (.failure(firstError), .failure(secondError), firstError) + ] + for (firstTermination, secondTermination, expectedError) in cases { + let subject = makeSubject() + let existing = subject.makeAsyncIterator() + subject.send(1) + subject.send(firstTermination) + subject.send(2) + subject.send(secondTermination) + let late = subject.makeAsyncIterator() + let existingResult = await receive(existing, file: file, line: line) + let lateResult = await receive(late, file: file, line: line) + XCTAssertEqual(existingResult.values, initial + [1], file: file, line: line) + XCTAssertEqual(existingResult.error, expectedError, file: file, line: line) + XCTAssertEqual(lateResult.values, [], file: file, line: line) + XCTAssertEqual(lateResult.error, expectedError, file: file, line: line) + } + } + + private func assertConcurrentTermination( + makeSubject: () -> Subject, + initial: [Int] = [], + file: StaticString = #filePath, + line: UInt = #line + ) async where Subject.Element == Int, Subject.Failure == Error { + for competingTermination: Termination in [.finished, .failure(MockError(code: 2))] { + for _ in 0..<200 { + let subject = makeSubject() + let existing = subject.makeAsyncIterator() + subject.send(1) + race({ subject.send(.failure(MockError(code: 1))) }, { subject.send(competingTermination) }) + let late = subject.makeAsyncIterator() + let existingResult = await receive(existing, file: file, line: line) + let lateResult = await receive(late, file: file, line: line) + XCTAssertEqual(existingResult.values, initial + [1], file: file, line: line) + XCTAssertEqual(lateResult.values, [], file: file, line: line) + switch competingTermination { + case .finished: + XCTAssertTrue(existingResult.error == nil || existingResult.error == MockError(code: 1), file: file, line: line) + case .failure: + XCTAssertTrue(existingResult.error == MockError(code: 1) || existingResult.error == MockError(code: 2), file: file, line: line) + } + guard existingResult.error == lateResult.error else { + return XCTFail("Existing and late consumers received different terminal outcomes", file: file, line: line) + } + } + } + } + + private func receive( + _ input: Iterator, + file: StaticString, + line: UInt + ) async -> (values: [Int], error: MockError?) where Iterator.Element == Int { + var iterator = input + var values: [Int] = [] + do { + while let value = try await iterator.next() { values.append(value) } + return (values, nil) + } catch { + guard let error = error as? MockError else { + XCTFail("Unexpected error: \(error)", file: file, line: line) + return (values, nil) + } + return (values, error) + } + } +} diff --git a/Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift b/Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift index 19db9af..b6bd60d 100644 --- a/Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift +++ b/Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift @@ -9,6 +9,23 @@ import XCTest final class AsyncThrowingReplaySubjectTests: XCTestCase { + func test_zero_capacity_does_not_replay_history_but_delivers_live_values() async { + let subject = AsyncThrowingReplaySubject(bufferSize: 0) + let existing = subject.makeAsyncIterator() + subject.send(1) + subject.send(2) + let late = subject.makeAsyncIterator() + subject.send(3) + subject.send(.finished) + + let existingResult = await drainBufferedElements(of: existing) + let lateResult = await drainBufferedElements(of: late) + XCTAssertEqual(existingResult.elements, [1, 2, 3]) + XCTAssertEqual(lateResult.elements, [3]) + XCTAssertTrue(existingResult.isTerminated) + XCTAssertTrue(lateResult.isTerminated) + } + func test_send_replays_buffered_elements() async throws { let exp = expectation(description: "Send has stacked elements in the replay the buffer") exp.expectedFulfillmentCount = 2