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
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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).

Expand Down
32 changes: 28 additions & 4 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
12 changes: 8 additions & 4 deletions Sources/AsyncSubjects/AsyncCurrentValueSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Int>(1)
Expand All @@ -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<Element>: AsyncSubject where Element: Sendable {
public typealias Element = Element
Expand Down Expand Up @@ -69,6 +71,7 @@ public final class AsyncCurrentValueSubject<Element>: 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 {
Expand All @@ -84,6 +87,7 @@ public final class AsyncCurrentValueSubject<Element>: AsyncSubject where Element
/// - Parameter termination: The termination to finish the subject.
public func send(_ termination: Termination<Failure>) {
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()
Expand Down
3 changes: 3 additions & 0 deletions Sources/AsyncSubjects/AsyncPassthroughSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Int>()
Expand Down Expand Up @@ -54,6 +55,7 @@ public final class AsyncPassthroughSubject<Element: Sendable>: 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 {
Expand All @@ -68,6 +70,7 @@ public final class AsyncPassthroughSubject<Element: Sendable>: AsyncSubject {
/// - Parameter termination: The termination to finish the subject
public func send(_ termination: Termination<Failure>) {
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()
Expand Down
14 changes: 10 additions & 4 deletions Sources/AsyncSubjects/AsyncReplaySubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Int>(bufferSize: 3)
Expand Down Expand Up @@ -48,10 +50,13 @@ public final class AsyncReplaySubject<Element>: 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 {
Expand All @@ -66,6 +71,7 @@ public final class AsyncReplaySubject<Element>: AsyncSubject where Element: Send
/// - Parameter termination: The termination to finish the subject.
public func send(_ termination: Termination<Failure>) {
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()
Expand Down
14 changes: 9 additions & 5 deletions Sources/AsyncSubjects/AsyncThrowingCurrentValueSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Int, Error>(1)
Expand All @@ -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<Element, Failure: Error>: AsyncSubject where Element: Sendable {
Expand Down Expand Up @@ -69,6 +71,7 @@ public final class AsyncThrowingCurrentValueSubject<Element, Failure: Error>: 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 {
Expand All @@ -84,6 +87,7 @@ public final class AsyncThrowingCurrentValueSubject<Element, Failure: Error>: As
/// - Parameter termination: The termination to finish the subject.
public func send(_ termination: Termination<Failure>) {
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()
Expand Down
3 changes: 3 additions & 0 deletions Sources/AsyncSubjects/AsyncThrowingPassthroughSubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Int, Error>()
Expand Down Expand Up @@ -55,6 +56,7 @@ public final class AsyncThrowingPassthroughSubject<Element, Failure: Error>: 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 {
Expand All @@ -69,6 +71,7 @@ public final class AsyncThrowingPassthroughSubject<Element, Failure: Error>: Asy
/// - Parameter termination: The termination to finish the subject
public func send(_ termination: Termination<Failure>) {
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()
Expand Down
18 changes: 13 additions & 5 deletions Sources/AsyncSubjects/AsyncThrowingReplaySubject.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Int, Error>(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<Element, Failure: Error>: AsyncSubject where Element: Sendable {
Expand Down Expand Up @@ -47,10 +51,13 @@ public final class AsyncThrowingReplaySubject<Element, Failure: Error>: 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 {
Expand All @@ -65,6 +72,7 @@ public final class AsyncThrowingReplaySubject<Element, Failure: Error>: AsyncSub
/// - Parameter termination: The termination to finish the subject
public func send(_ termination: Termination<Failure>) {
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()
Expand Down
17 changes: 17 additions & 0 deletions Tests/AsyncSubjets/AsyncReplaySubjectTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -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<Int>(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
Expand Down
40 changes: 40 additions & 0 deletions Tests/AsyncSubjets/AsyncSubjectLifetimeTests.swift
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,46 @@ private final class LifetimePayload: Sendable {
}

final class AsyncSubjectLifetimeTests: XCTestCase {
func test_terminated_subjects_do_not_retain_new_values() {
assertTerminatedValueReleased(AsyncPassthroughSubject<LifetimePayload>())
assertTerminatedValueReleased(AsyncCurrentValueSubject(LifetimePayload {}))
assertTerminatedValueReleased(AsyncReplaySubject<LifetimePayload>(bufferSize: 2))
for termination: Termination<Error> in [.finished, .failure(MockError(code: 1))] {
assertTerminatedValueReleased(AsyncThrowingPassthroughSubject<LifetimePayload, Error>(), termination: termination)
assertTerminatedValueReleased(AsyncThrowingCurrentValueSubject<LifetimePayload, Error>(LifetimePayload {}), termination: termination)
assertTerminatedValueReleased(AsyncThrowingReplaySubject<LifetimePayload, Error>(bufferSize: 2), termination: termination)
}
}

private func assertTerminatedValueReleased<S: AsyncSubject>(
_ subject: S,
termination: Termination<S.Failure> = .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<LifetimePayload>(bufferSize: 0))
assertSentValueReleased(AsyncThrowingReplaySubject<LifetimePayload, Error>(bufferSize: 0))
}

private func assertSentValueReleased<S: AsyncSubject>(
_ 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<LifetimePayload>())
assertBufferedValueReleased(AsyncThrowingPassthroughSubject<LifetimePayload, Error>())
Expand Down
Loading
Loading