Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,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:**
Expand Down
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
10 changes: 6 additions & 4 deletions Sources/AsyncSubjects/AsyncCurrentValueSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -139,10 +139,12 @@ public final class AsyncCurrentValueSubject<Element>: AsyncSubject where Element

public struct Iterator: AsyncSubjectIterator {
var iterator: AsyncBufferedChannel<Element>.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 {
Expand All @@ -152,8 +154,8 @@ public final class AsyncCurrentValueSubject<Element>: AsyncSubject where Element
public mutating func next() async -> Element? {
await withTaskCancellationHandler {
await self.iterator.next()
} onCancel: { [unregister] in
unregister()
} onCancel: { [subscription] in
subscription.unregister()
}
}
}
Expand Down
10 changes: 6 additions & 4 deletions Sources/AsyncSubjects/AsyncPassthroughSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -121,10 +121,12 @@ public final class AsyncPassthroughSubject<Element: Sendable>: AsyncSubject {

public struct Iterator: AsyncSubjectIterator {
var iterator: AsyncBufferedChannel<Element>.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 {
Expand All @@ -134,8 +136,8 @@ public final class AsyncPassthroughSubject<Element: Sendable>: AsyncSubject {
public mutating func next() async -> Element? {
await withTaskCancellationHandler {
await self.iterator.next()
} onCancel: { [unregister] in
unregister()
} onCancel: { [subscription] in
subscription.unregister()
}
}
}
Expand Down
10 changes: 6 additions & 4 deletions Sources/AsyncSubjects/AsyncReplaySubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -125,10 +125,12 @@ public final class AsyncReplaySubject<Element>: AsyncSubject where Element: Send

public struct Iterator: AsyncSubjectIterator {
var iterator: AsyncBufferedChannel<Element>.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 {
Expand All @@ -138,8 +140,8 @@ public final class AsyncReplaySubject<Element>: 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()
}
}
}
Expand Down
10 changes: 6 additions & 4 deletions Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -150,10 +150,12 @@ public final class AsyncThrowingCurrentValueSubject<Element, Failure: Error>: As

public struct Iterator: AsyncSubjectIterator {
var iterator: AsyncThrowingBufferedChannel<Element, Error>.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 {
Expand All @@ -163,8 +165,8 @@ public final class AsyncThrowingCurrentValueSubject<Element, Failure: Error>: 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()
}
}
}
Expand Down
10 changes: 6 additions & 4 deletions Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -133,10 +133,12 @@ public final class AsyncThrowingPassthroughSubject<Element, Failure: Error>: Asy

public struct Iterator: AsyncSubjectIterator {
var iterator: AsyncThrowingBufferedChannel<Element, Error>.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 {
Expand All @@ -146,8 +148,8 @@ public final class AsyncThrowingPassthroughSubject<Element, Failure: Error>: 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()
}
}
}
Expand Down
10 changes: 6 additions & 4 deletions Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -135,10 +135,12 @@ public final class AsyncThrowingReplaySubject<Element, Failure: Error>: AsyncSub

public struct Iterator: AsyncSubjectIterator {
var iterator: AsyncThrowingBufferedChannel<Element, Error>.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 {
Expand All @@ -148,8 +150,8 @@ public final class AsyncThrowingReplaySubject<Element, Failure: Error>: 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()
}
}
}
Expand Down
11 changes: 11 additions & 0 deletions Sources/Supporting/SubjectSubscription.swift
Original file line number Diff line number Diff line change
@@ -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() }
}
14 changes: 9 additions & 5 deletions Tests/AsyncSubjets/AsyncCurrentValueSubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -42,7 +42,7 @@ final class AsyncCurrentValueSubjectTests: XCTestCase {

let sut = AsyncCurrentValueSubject<Int>(1)

Task {
let firstConsumer = Task {
var receivedElements = [Int]()

for await element in sut {
Expand All @@ -57,7 +57,7 @@ final class AsyncCurrentValueSubjectTests: XCTestCase {
}
}

Task {
let secondConsumer = Task {
var receivedElements = [Int]()

for await element in sut {
Expand All @@ -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 {
Expand Down
14 changes: 9 additions & 5 deletions Tests/AsyncSubjets/AsyncPassthroughSubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -20,7 +20,7 @@ final class AsyncPassthroughSubjectTests: XCTestCase {

let sut = AsyncPassthroughSubject<Int>()

Task {
let firstConsumer = Task {
var receivedElements = [Int]()

var it = sut.makeAsyncIterator()
Expand All @@ -34,7 +34,7 @@ final class AsyncPassthroughSubjectTests: XCTestCase {
}
}

Task {
let secondConsumer = Task {
var receivedElements = [Int]()

var it = sut.makeAsyncIterator()
Expand All @@ -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 {
Expand Down
26 changes: 17 additions & 9 deletions Tests/AsyncSubjets/AsyncReplaySubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand All @@ -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 {
Expand All @@ -35,7 +35,7 @@ final class AsyncReplaySubjectTests: XCTestCase {
}
}

Task {
let secondConsumer = Task {
var receivedElements = [Int]()

for await element in sut {
Expand All @@ -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

Expand All @@ -63,7 +67,7 @@ final class AsyncReplaySubjectTests: XCTestCase {

sut.send(1)

Task {
let firstConsumer = Task {
var receivedElements = [Int]()

for await element in sut {
Expand All @@ -78,7 +82,7 @@ final class AsyncReplaySubjectTests: XCTestCase {
}
}

Task {
let secondConsumer = Task {
var receivedElements = [Int]()

for await element in sut {
Expand All @@ -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 {
Expand Down
Loading
Loading