From 6ecdd4b5d249567f24d18278f3080c13534a7ff7 Mon Sep 17 00:00:00 2001 From: Thibault Wittemberg Date: Sat, 3 Oct 2026 13:15:01 +0200 Subject: [PATCH 1/2] Release abandoned subject subscriptions and finish test consumers --- CHANGELOG.md | 1 + README.md | 2 + .../AsyncCurrentValueSubject.swift | 10 ++-- .../AsyncPassthroughSubject.swift | 10 ++-- .../AsyncSubjects/AsyncReplaySubject.swift | 10 ++-- .../AsyncThrowingCurrentValueSubject.swift | 10 ++-- .../AsyncThrowingPassthroughSubject.swift | 10 ++-- .../AsyncThrowingReplaySubject.swift | 10 ++-- Sources/Supporting/SubjectSubscription.swift | 11 ++++ .../AsyncCurrentValueSubjectTests.swift | 14 +++-- .../AsyncPassthroughSubjectTests.swift | 14 +++-- .../AsyncReplaySubjectTests.swift | 26 ++++++---- .../AsyncSubjectLifetimeTests.swift | 51 +++++++++++++++++++ ...syncThrowingCurrentValueSubjectTests.swift | 14 +++-- ...AsyncThrowingPassthroughSubjectTests.swift | 14 +++-- .../AsyncThrowingReplaySubjectTests.swift | 26 ++++++---- 16 files changed, 171 insertions(+), 62 deletions(-) create mode 100644 Sources/Supporting/SubjectSubscription.swift create mode 100644 Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift diff --git a/CHANGELOG.md b/CHANGELOG.md index 1947404..1500dbb 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,6 +1,7 @@ **Unreleased:** - Just: release consumed values from iterators so nested `flatMapLatest` subscriptions do not retain removed feature objects (https://github.com/sideeffect-io/AsyncExtensions/issues/35). +- 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). - Breaking: use Swift Async Algorithms for `Sequence.async` and two- or three-input `zip`/`merge`, removing import ambiguities. Add the `AsyncAlgorithms` product dependency and import when migrating these APIs. Preserve variadic `zip`/`merge` and the explicit `AsyncLazySequence` constructor. - Breaking: rename the buffered Date timer from `AsyncTimerSequence` to `AsyncBufferedTimerSequence` to avoid ambiguity with Apple's clock-based timer. diff --git a/README.md b/README.md index 2067602..245a8d6 100644 --- a/README.md +++ b/README.md @@ -88,6 +88,8 @@ goes first is unspecified, but consumers registered for the same sends receive t 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. +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 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 diff --git a/Sources/AsyncSubjects/AsyncCurrentValueSubject.swift b/Sources/AsyncSubjects/AsyncCurrentValueSubject.swift index 43a88ee..b2135bb 100644 --- a/Sources/AsyncSubjects/AsyncCurrentValueSubject.swift +++ b/Sources/AsyncSubjects/AsyncCurrentValueSubject.swift @@ -139,10 +139,12 @@ public final class AsyncCurrentValueSubject: AsyncSubject where Element public struct Iterator: AsyncSubjectIterator { var iterator: AsyncBufferedChannel.Iterator - let unregister: @Sendable () -> Void + let subscription: SubjectSubscription init(asyncSubject: AsyncCurrentValueSubject) { - (self.iterator, self.unregister) = asyncSubject.handleNewConsumer() + let consumer = asyncSubject.handleNewConsumer() + self.iterator = consumer.iterator + self.subscription = SubjectSubscription(unregister: consumer.unregister) } public var hasBufferedElements: Bool { @@ -152,8 +154,8 @@ public final class AsyncCurrentValueSubject: AsyncSubject where Element public mutating func next() async -> Element? { await withTaskCancellationHandler { await self.iterator.next() - } onCancel: { [unregister] in - unregister() + } onCancel: { [subscription] in + subscription.unregister() } } } diff --git a/Sources/AsyncSubjects/AsyncPassthroughSubject.swift b/Sources/AsyncSubjects/AsyncPassthroughSubject.swift index 49bfe0a..8bd7dce 100644 --- a/Sources/AsyncSubjects/AsyncPassthroughSubject.swift +++ b/Sources/AsyncSubjects/AsyncPassthroughSubject.swift @@ -121,10 +121,12 @@ public final class AsyncPassthroughSubject: AsyncSubject { public struct Iterator: AsyncSubjectIterator { var iterator: AsyncBufferedChannel.Iterator - let unregister: @Sendable () -> Void + let subscription: SubjectSubscription init(asyncSubject: AsyncPassthroughSubject) { - (self.iterator, self.unregister) = asyncSubject.handleNewConsumer() + let consumer = asyncSubject.handleNewConsumer() + self.iterator = consumer.iterator + self.subscription = SubjectSubscription(unregister: consumer.unregister) } public var hasBufferedElements: Bool { @@ -134,8 +136,8 @@ public final class AsyncPassthroughSubject: AsyncSubject { public mutating func next() async -> Element? { await withTaskCancellationHandler { await self.iterator.next() - } onCancel: { [unregister] in - unregister() + } onCancel: { [subscription] in + subscription.unregister() } } } diff --git a/Sources/AsyncSubjects/AsyncReplaySubject.swift b/Sources/AsyncSubjects/AsyncReplaySubject.swift index ea3bf46..2f7f7f0 100644 --- a/Sources/AsyncSubjects/AsyncReplaySubject.swift +++ b/Sources/AsyncSubjects/AsyncReplaySubject.swift @@ -125,10 +125,12 @@ public final class AsyncReplaySubject: AsyncSubject where Element: Send public struct Iterator: AsyncSubjectIterator { var iterator: AsyncBufferedChannel.Iterator - let unregister: @Sendable () -> Void + let subscription: SubjectSubscription init(asyncSubject: AsyncReplaySubject) { - (self.iterator, self.unregister) = asyncSubject.handleNewConsumer() + let consumer = asyncSubject.handleNewConsumer() + self.iterator = consumer.iterator + self.subscription = SubjectSubscription(unregister: consumer.unregister) } public var hasBufferedElements: Bool { @@ -138,8 +140,8 @@ public final class AsyncReplaySubject: AsyncSubject where Element: Send public mutating func next() async -> Element? { await withTaskCancellationHandler { await self.iterator.next() - } onCancel: { [unregister] in - unregister() + } onCancel: { [subscription] in + subscription.unregister() } } } diff --git a/Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift b/Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift index 18e5ee5..483e006 100644 --- a/Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift +++ b/Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift @@ -150,10 +150,12 @@ public final class AsyncThrowingCurrentValueSubject: As public struct Iterator: AsyncSubjectIterator { var iterator: AsyncThrowingBufferedChannel.Iterator - let unregister: @Sendable () -> Void + let subscription: SubjectSubscription init(asyncSubject: AsyncThrowingCurrentValueSubject) { - (self.iterator, self.unregister) = asyncSubject.handleNewConsumer() + let consumer = asyncSubject.handleNewConsumer() + self.iterator = consumer.iterator + self.subscription = SubjectSubscription(unregister: consumer.unregister) } public var hasBufferedElements: Bool { @@ -163,8 +165,8 @@ public final class AsyncThrowingCurrentValueSubject: As public mutating func next() async throws -> Element? { try await withTaskCancellationHandler { try await self.iterator.next() - } onCancel: { [unregister] in - unregister() + } onCancel: { [subscription] in + subscription.unregister() } } } diff --git a/Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift b/Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift index 8c2ff40..3cc093d 100644 --- a/Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift +++ b/Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift @@ -133,10 +133,12 @@ public final class AsyncThrowingPassthroughSubject: Asy public struct Iterator: AsyncSubjectIterator { var iterator: AsyncThrowingBufferedChannel.Iterator - let unregister: @Sendable () -> Void + let subscription: SubjectSubscription init(asyncSubject: AsyncThrowingPassthroughSubject) { - (self.iterator, self.unregister) = asyncSubject.handleNewConsumer() + let consumer = asyncSubject.handleNewConsumer() + self.iterator = consumer.iterator + self.subscription = SubjectSubscription(unregister: consumer.unregister) } public var hasBufferedElements: Bool { @@ -146,8 +148,8 @@ public final class AsyncThrowingPassthroughSubject: Asy public mutating func next() async throws -> Element? { try await withTaskCancellationHandler { try await self.iterator.next() - } onCancel: { [unregister] in - unregister() + } onCancel: { [subscription] in + subscription.unregister() } } } diff --git a/Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift b/Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift index f316151..b654875 100644 --- a/Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift +++ b/Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift @@ -135,10 +135,12 @@ public final class AsyncThrowingReplaySubject: AsyncSub public struct Iterator: AsyncSubjectIterator { var iterator: AsyncThrowingBufferedChannel.Iterator - let unregister: @Sendable () -> Void + let subscription: SubjectSubscription init(asyncSubject: AsyncThrowingReplaySubject) { - (self.iterator, self.unregister) = asyncSubject.handleNewConsumer() + let consumer = asyncSubject.handleNewConsumer() + self.iterator = consumer.iterator + self.subscription = SubjectSubscription(unregister: consumer.unregister) } public var hasBufferedElements: Bool { @@ -148,8 +150,8 @@ public final class AsyncThrowingReplaySubject: AsyncSub public mutating func next() async throws -> Element? { try await withTaskCancellationHandler { try await self.iterator.next() - } onCancel: { [unregister] in - unregister() + } onCancel: { [subscription] in + subscription.unregister() } } } diff --git a/Sources/Supporting/SubjectSubscription.swift b/Sources/Supporting/SubjectSubscription.swift new file mode 100644 index 0000000..ab547d1 --- /dev/null +++ b/Sources/Supporting/SubjectSubscription.swift @@ -0,0 +1,11 @@ +// Iterator copies share a registration. Remove it when their last copy is released, or when +// iteration is cancelled. Subject unregistration is idempotent and protected by its state lock. +final class SubjectSubscription: Sendable { + let unregister: @Sendable () -> Void + + init(unregister: @Sendable @escaping () -> Void) { + self.unregister = unregister + } + + deinit { unregister() } +} diff --git a/Tests/AsyncSubjets/AsyncCurrentValueSubjectTests.swift b/Tests/AsyncSubjets/AsyncCurrentValueSubjectTests.swift index a7ccadd..513d7e0 100644 --- a/Tests/AsyncSubjets/AsyncCurrentValueSubjectTests.swift +++ b/Tests/AsyncSubjets/AsyncCurrentValueSubjectTests.swift @@ -31,7 +31,7 @@ final class AsyncCurrentValueSubjectTests: XCTestCase { XCTAssertEqual(received2, 1) } - func test_send_pushes_values_in_the_subject() { + func test_send_pushes_values_in_the_subject() async { let hasReceivedOneElementExpectation = expectation(description: "One element has been iterated in the async sequence") hasReceivedOneElementExpectation.expectedFulfillmentCount = 2 @@ -42,7 +42,7 @@ final class AsyncCurrentValueSubjectTests: XCTestCase { let sut = AsyncCurrentValueSubject(1) - Task { + let firstConsumer = Task { var receivedElements = [Int]() for await element in sut { @@ -57,7 +57,7 @@ final class AsyncCurrentValueSubjectTests: XCTestCase { } } - Task { + let secondConsumer = Task { var receivedElements = [Int]() for await element in sut { @@ -72,12 +72,16 @@ final class AsyncCurrentValueSubjectTests: XCTestCase { } } - wait(for: [hasReceivedOneElementExpectation], timeout: 1) + await fulfillment(of: [hasReceivedOneElementExpectation], timeout: 1) sut.send(2) sut.value = 3 - wait(for: [hasReceivedSentElementsExpectation], timeout: 1) + await fulfillment(of: [hasReceivedSentElementsExpectation], timeout: 1) + sut.send(.finished) + await firstConsumer.value + await secondConsumer.value + } func test_sendFinished_ends_the_subject_and_immediately_resumes_futur_consumer() async { diff --git a/Tests/AsyncSubjets/AsyncPassthroughSubjectTests.swift b/Tests/AsyncSubjets/AsyncPassthroughSubjectTests.swift index 3f098a9..e4c711e 100644 --- a/Tests/AsyncSubjets/AsyncPassthroughSubjectTests.swift +++ b/Tests/AsyncSubjets/AsyncPassthroughSubjectTests.swift @@ -9,7 +9,7 @@ import XCTest final class AsyncPassthroughSubjectTests: XCTestCase { - func test_send_pushes_elements_in_the_subject() { + func test_send_pushes_elements_in_the_subject() async { let isReadyToBeIteratedExpectation = expectation(description: "Passthrough subject iterators are ready for iteration") isReadyToBeIteratedExpectation.expectedFulfillmentCount = 2 @@ -20,7 +20,7 @@ final class AsyncPassthroughSubjectTests: XCTestCase { let sut = AsyncPassthroughSubject() - Task { + let firstConsumer = Task { var receivedElements = [Int]() var it = sut.makeAsyncIterator() @@ -34,7 +34,7 @@ final class AsyncPassthroughSubjectTests: XCTestCase { } } - Task { + let secondConsumer = Task { var receivedElements = [Int]() var it = sut.makeAsyncIterator() @@ -48,13 +48,17 @@ final class AsyncPassthroughSubjectTests: XCTestCase { } } - wait(for: [isReadyToBeIteratedExpectation], timeout: 1) + await fulfillment(of: [isReadyToBeIteratedExpectation], timeout: 1) sut.send(1) sut.send(2) sut.send(3) - wait(for: [hasReceivedSentElementsExpectation], timeout: 1) + await fulfillment(of: [hasReceivedSentElementsExpectation], timeout: 1) + sut.send(.finished) + await firstConsumer.value + await secondConsumer.value + } func test_sendFinished_ends_the_subject_and_immediately_resumes_futur_consumer() async { diff --git a/Tests/AsyncSubjets/AsyncReplaySubjectTests.swift b/Tests/AsyncSubjets/AsyncReplaySubjectTests.swift index aba3a58..5452bb8 100644 --- a/Tests/AsyncSubjets/AsyncReplaySubjectTests.swift +++ b/Tests/AsyncSubjets/AsyncReplaySubjectTests.swift @@ -9,7 +9,7 @@ import XCTest final class AsyncReplaySubjectTests: XCTestCase { - func test_send_replays_buffered_elements() { + func test_send_replays_buffered_elements() async { let exp = expectation(description: "Send has stacked elements in the replay the buffer") exp.expectedFulfillmentCount = 2 @@ -23,7 +23,7 @@ final class AsyncReplaySubjectTests: XCTestCase { sut.send(5) sut.send(6) - Task { + let firstConsumer = Task { var receivedElements = [Int]() for await element in sut { @@ -35,7 +35,7 @@ final class AsyncReplaySubjectTests: XCTestCase { } } - Task { + let secondConsumer = Task { var receivedElements = [Int]() for await element in sut { @@ -47,10 +47,14 @@ final class AsyncReplaySubjectTests: XCTestCase { } } - waitForExpectations(timeout: 0.5) + await fulfillment(of: [exp], timeout: 1) + sut.send(.finished) + await firstConsumer.value + await secondConsumer.value + } - func test_send_pushes_elements_in_the_subject() { + func test_send_pushes_elements_in_the_subject() async { let hasReceivedOneElementExpectation = expectation(description: "One element has been iterated in the async sequence") hasReceivedOneElementExpectation.expectedFulfillmentCount = 2 @@ -63,7 +67,7 @@ final class AsyncReplaySubjectTests: XCTestCase { sut.send(1) - Task { + let firstConsumer = Task { var receivedElements = [Int]() for await element in sut { @@ -78,7 +82,7 @@ final class AsyncReplaySubjectTests: XCTestCase { } } - Task { + let secondConsumer = Task { var receivedElements = [Int]() for await element in sut { @@ -93,12 +97,16 @@ final class AsyncReplaySubjectTests: XCTestCase { } } - wait(for: [hasReceivedOneElementExpectation], timeout: 1) + await fulfillment(of: [hasReceivedOneElementExpectation], timeout: 1) sut.send(2) sut.send(3) - wait(for: [hasReceivedSentElementsExpectation], timeout: 1) + await fulfillment(of: [hasReceivedSentElementsExpectation], timeout: 1) + sut.send(.finished) + await firstConsumer.value + await secondConsumer.value + } func test_sendFinished_ends_the_subject_and_immediately_resumes_futur_consumer() async { diff --git a/Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift b/Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift new file mode 100644 index 0000000..b84b6b0 --- /dev/null +++ b/Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift @@ -0,0 +1,51 @@ +@testable import AsyncExtensions +import XCTest + +private final class LifetimePayload: Sendable { + let onDeinit: @Sendable () -> Void + init(onDeinit: @Sendable @escaping () -> Void) { self.onDeinit = onDeinit } + deinit { onDeinit() } +} + +final class AsyncSubjectLifetimeTests: XCTestCase { + func test_abandoned_passthrough_iterators_release_their_buffered_values() { + assertBufferedValueReleased(AsyncPassthroughSubject()) + assertBufferedValueReleased(AsyncThrowingPassthroughSubject()) + } + + private func assertBufferedValueReleased(_ subject: S) where S.Element == LifetimePayload { + let released = ManagedCriticalState(false) + var iterator: S.AsyncIterator? = subject.makeAsyncIterator() + var payload: LifetimePayload? = LifetimePayload { released.apply(criticalState: true) } + subject.send(payload!) + payload = nil + withExtendedLifetime(iterator) { XCTAssertFalse(released.criticalState) } + iterator = nil + XCTAssertTrue(released.criticalState) + } + + func test_last_iterator_copy_unregisters_each_subject_subscription() { + let passthrough = AsyncPassthroughSubject() + assertUnregister(passthrough) { passthrough.state.withCriticalRegion { $0.channels.count } } + let current = AsyncCurrentValueSubject(0) + assertUnregister(current) { current.state.withCriticalRegion { $0.channels.count } } + let replay = AsyncReplaySubject(bufferSize: 1) + assertUnregister(replay) { replay.state.withCriticalRegion { $0.channels.count } } + let throwingPassthrough = AsyncThrowingPassthroughSubject() + assertUnregister(throwingPassthrough) { throwingPassthrough.state.withCriticalRegion { $0.channels.count } } + let throwingCurrent = AsyncThrowingCurrentValueSubject(0) + assertUnregister(throwingCurrent) { throwingCurrent.state.withCriticalRegion { $0.channels.count } } + let throwingReplay = AsyncThrowingReplaySubject(bufferSize: 1) + assertUnregister(throwingReplay) { throwingReplay.state.withCriticalRegion { $0.channels.count } } + } + + private func assertUnregister(_ subject: S, subscriberCount: () -> Int) { + var iterator: S.AsyncIterator? = subject.makeAsyncIterator() + var copy = iterator + XCTAssertEqual(subscriberCount(), 1) + iterator = nil + withExtendedLifetime(copy) { XCTAssertEqual(subscriberCount(), 1) } + copy = nil + XCTAssertEqual(subscriberCount(), 0) + } +} diff --git a/Tests/AsyncSubjets/AsyncThrowingCurrentValueSubjectTests.swift b/Tests/AsyncSubjets/AsyncThrowingCurrentValueSubjectTests.swift index d933970..81baa48 100644 --- a/Tests/AsyncSubjets/AsyncThrowingCurrentValueSubjectTests.swift +++ b/Tests/AsyncSubjets/AsyncThrowingCurrentValueSubjectTests.swift @@ -31,7 +31,7 @@ final class AsyncThrowingCurrentValueSubjectTests: XCTestCase { XCTAssertEqual(received2, 1) } - func test_send_pushes_values_in_the_subject() { + func test_send_pushes_values_in_the_subject() async throws { let hasReceivedOneElementExpectation = expectation(description: "One element has been iterated in the async sequence") hasReceivedOneElementExpectation.expectedFulfillmentCount = 2 @@ -42,7 +42,7 @@ final class AsyncThrowingCurrentValueSubjectTests: XCTestCase { let sut = AsyncThrowingCurrentValueSubject(1) - Task { + let firstConsumer = Task { var receivedElements = [Int]() for try await element in sut { @@ -57,7 +57,7 @@ final class AsyncThrowingCurrentValueSubjectTests: XCTestCase { } } - Task { + let secondConsumer = Task { var receivedElements = [Int]() for try await element in sut { @@ -72,12 +72,16 @@ final class AsyncThrowingCurrentValueSubjectTests: XCTestCase { } } - wait(for: [hasReceivedOneElementExpectation], timeout: 1) + await fulfillment(of: [hasReceivedOneElementExpectation], timeout: 1) sut.send(2) sut.value = 3 - wait(for: [hasReceivedSentElementsExpectation], timeout: 1) + await fulfillment(of: [hasReceivedSentElementsExpectation], timeout: 1) + sut.send(.finished) + try await firstConsumer.value + try await secondConsumer.value + } func test_sendFinished_ends_the_subject_and_immediately_resumes_futur_consumer() async throws { diff --git a/Tests/AsyncSubjets/AsyncThrowingPassthroughSubjectTests.swift b/Tests/AsyncSubjets/AsyncThrowingPassthroughSubjectTests.swift index b50c8ac..95972d8 100644 --- a/Tests/AsyncSubjets/AsyncThrowingPassthroughSubjectTests.swift +++ b/Tests/AsyncSubjets/AsyncThrowingPassthroughSubjectTests.swift @@ -9,7 +9,7 @@ import XCTest final class AsyncThrowingPassthroughSubjectTests: XCTestCase { - func test_send_pushes_elements_in_the_subject() { + func test_send_pushes_elements_in_the_subject() async throws { let isReadyToBeIteratedExpectation = expectation(description: "Passthrough subject iterators are ready for iteration") isReadyToBeIteratedExpectation.expectedFulfillmentCount = 2 @@ -20,7 +20,7 @@ final class AsyncThrowingPassthroughSubjectTests: XCTestCase { let sut = AsyncThrowingPassthroughSubject() - Task { + let firstConsumer = Task { var receivedElements = [Int]() var it = sut.makeAsyncIterator() @@ -34,7 +34,7 @@ final class AsyncThrowingPassthroughSubjectTests: XCTestCase { } } - Task { + let secondConsumer = Task { var receivedElements = [Int]() var it = sut.makeAsyncIterator() @@ -48,13 +48,17 @@ final class AsyncThrowingPassthroughSubjectTests: XCTestCase { } } - wait(for: [isReadyToBeIteratedExpectation], timeout: 1) + await fulfillment(of: [isReadyToBeIteratedExpectation], timeout: 1) sut.send(1) sut.send(2) sut.send(3) - wait(for: [hasReceivedSentElementsExpectation], timeout: 1) + await fulfillment(of: [hasReceivedSentElementsExpectation], timeout: 1) + sut.send(.finished) + try await firstConsumer.value + try await secondConsumer.value + } func test_sendFinished_ends_the_subject_and_immediately_resumes_futur_consumer() async throws { diff --git a/Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift b/Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift index b1d2d85..19db9af 100644 --- a/Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift +++ b/Tests/AsyncSubjets/AsyncThrowingReplaySubjectTests.swift @@ -9,7 +9,7 @@ import XCTest final class AsyncThrowingReplaySubjectTests: XCTestCase { - func test_send_replays_buffered_elements() { + func test_send_replays_buffered_elements() async throws { let exp = expectation(description: "Send has stacked elements in the replay the buffer") exp.expectedFulfillmentCount = 2 @@ -23,7 +23,7 @@ final class AsyncThrowingReplaySubjectTests: XCTestCase { sut.send(5) sut.send(6) - Task { + let firstConsumer = Task { var receivedElements = [Int]() for try await element in sut { @@ -35,7 +35,7 @@ final class AsyncThrowingReplaySubjectTests: XCTestCase { } } - Task { + let secondConsumer = Task { var receivedElements = [Int]() for try await element in sut { @@ -47,10 +47,14 @@ final class AsyncThrowingReplaySubjectTests: XCTestCase { } } - waitForExpectations(timeout: 0.5) + await fulfillment(of: [exp], timeout: 1) + sut.send(.finished) + try await firstConsumer.value + try await secondConsumer.value + } - func test_send_pushes_elements_in_the_subject() { + func test_send_pushes_elements_in_the_subject() async throws { let hasReceivedOneElementExpectation = expectation(description: "One element has been iterated in the async sequence") hasReceivedOneElementExpectation.expectedFulfillmentCount = 2 @@ -63,7 +67,7 @@ final class AsyncThrowingReplaySubjectTests: XCTestCase { sut.send(1) - Task { + let firstConsumer = Task { var receivedElements = [Int]() for try await element in sut { @@ -78,7 +82,7 @@ final class AsyncThrowingReplaySubjectTests: XCTestCase { } } - Task { + let secondConsumer = Task { var receivedElements = [Int]() for try await element in sut { @@ -93,12 +97,16 @@ final class AsyncThrowingReplaySubjectTests: XCTestCase { } } - wait(for: [hasReceivedOneElementExpectation], timeout: 1) + await fulfillment(of: [hasReceivedOneElementExpectation], timeout: 1) sut.send(2) sut.send(3) - wait(for: [hasReceivedSentElementsExpectation], timeout: 1) + await fulfillment(of: [hasReceivedSentElementsExpectation], timeout: 1) + sut.send(.finished) + try await firstConsumer.value + try await secondConsumer.value + } func test_sendFinished_ends_the_subject_and_immediately_resumes_futur_consumer() async throws { From d50174756db27090d4d80dbd5edbdbf573a200f0 Mon Sep 17 00:00:00 2001 From: Thibault Wittemberg Date: Sat, 3 Oct 2026 13:17:18 +0200 Subject: [PATCH 2/2] Group the subject lifetime entry with existing subject fixes --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1500dbb..a4634de 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,7 +1,6 @@ **Unreleased:** - Just: release consumed values from iterators so nested `flatMapLatest` subscriptions do not retain removed feature objects (https://github.com/sideeffect-io/AsyncExtensions/issues/35). -- 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). - Breaking: use Swift Async Algorithms for `Sequence.async` and two- or three-input `zip`/`merge`, removing import ambiguities. Add the `AsyncAlgorithms` product dependency and import when migrating these APIs. Preserve variadic `zip`/`merge` and the explicit `AsyncLazySequence` constructor. - Breaking: rename the buffered Date timer from `AsyncTimerSequence` to `AsyncBufferedTimerSequence` to avoid ambiguity with Apple's clock-based timer. @@ -10,6 +9,7 @@ - SwitchToLatest: finish cancelled collection while the latest channel or outer sequence remains open, and discard late producer results (https://github.com/sideeffect-io/AsyncExtensions/issues/53). - 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). - Multicast: preserve upstream element order and prevent termination from overtaking values when consumers advance concurrently. **v0.5.2 - Oxygen:**