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 @@ -10,6 +10,7 @@
- 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.
- 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).

**v0.5.2 - Oxygen:**

Expand Down
28 changes: 28 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,34 @@ The `AsyncLazySequence(sequence)` constructor remains available when you only ne

Rename uses of AsyncExtensions' `AsyncTimerSequence` to `AsyncBufferedTimerSequence`. It retains the buffered `Date` values and `DispatchTimeInterval` initializer. The unqualified `AsyncTimerSequence` name now refers to Apple's clock-based timer when both modules are imported.

## Sharing one upstream iterator

Create `multicast` or `share` once, keep the returned instance, and return that same instance to every consumer. Calling either operator each time a consumer subscribes creates a new upstream iterator. An `AsyncThrowingStream` cannot have overlapping `next()` calls, so separate wrappers over the same stream can crash even if they use the same subject.

For example, the connection in issue [#31](https://github.com/sideeffect-io/AsyncExtensions/issues/31) can store its shared sequence during initialization:

```swift
actor Connection {
private let continuation: AsyncThrowingStream<Int, Error>.Continuation
private let shared: AsyncShareSequence<AsyncThrowingStream<Int, Error>>

init() {
let (continuation, stream) = AsyncThrowingStream<Int, Error>.pipe()
self.continuation = continuation
self.shared = stream.share()
}

func events() -> AsyncShareSequence<AsyncThrowingStream<Int, Error>> {
shared
}

func send(_ value: Int) { continuation.yield(value) }
func finish() { continuation.finish() }
}
```

Cancelling a subscriber finishes its iteration without cancelling the shared upstream. The connection still owns upstream termination and must finish or cancel its producer when the connection ends.

## Features

### Channels
Expand Down
31 changes: 18 additions & 13 deletions Sources/Operators/AsyncMulticastSequence.swift
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,8 @@ public extension AsyncSequence {
/// to only produce a single `AsyncIterator`.
/// This is useful when upstream async sequences are doing expensive work you don’t want to duplicate,
/// like performing network requests.
/// Create this sequence once and return the same instance to every consumer. Repeated calls to
/// `multicast` create separate upstream iterators, even when they use the same subject.
///
/// The following example uses an async sequence as a counter to emit three random numbers.
/// It uses a ``AsyncSequence/multicast(_:)`` operator with a ``AsyncThrowingPassthroughSubject`
Expand Down Expand Up @@ -105,20 +107,22 @@ where Base.Element == Subject.Element, Subject.Failure == Error, Base.AsyncItera
self.connectedGate.send(())
}

func next() async {
await Task {
let (canAccessBase, iterator) = self.state.withCriticalRegion { state -> (Bool, Base.AsyncIterator?) in
switch state {
case .available(let iterator):
state = .busy
return (true, iterator)
case .busy:
return (false, nil)
}
func requestNext() {
let iterator = self.state.withCriticalRegion { state -> Base.AsyncIterator? in
switch state {
case .available(let iterator):
state = .busy
return iterator
case .busy:
return nil
}
}

guard canAccessBase, var iterator = iterator else { return }
guard var iterator = iterator else { return }

// This pull belongs to the shared sequence. A subscriber waits on its own subject channel
// so cancelling that subscriber can finish immediately without cancelling upstream.
Task {
let toSend: Result<Element?, Error>
do {
let element = try await iterator.next()
Expand All @@ -137,7 +141,7 @@ where Base.Element == Subject.Element, Subject.Failure == Error, Base.AsyncItera
}
state = .available(iterator)
}
}.value
}
}

public func makeAsyncIterator() -> AsyncIterator {
Expand All @@ -163,9 +167,10 @@ where Base.Element == Subject.Element, Subject.Failure == Error, Base.AsyncItera
if !isConnected {
await self.connectedGateIterator.next()
}
guard !Task.isCancelled else { return nil }

if !self.subjectIterator.hasBufferedElements {
await self.asyncMulticastSequence.next()
self.asyncMulticastSequence.requestNext()
}

let element = try await self.subjectIterator.next()
Expand Down
2 changes: 2 additions & 0 deletions Sources/Operators/AsyncSequence+Share.swift
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@

public extension AsyncSequence {
/// Shares the output of an upstream async sequence with multiple client loops.
/// Store and reuse the returned instance for all consumers; each call to `share()` creates
/// a separate upstream iterator.
///
/// - Tip: ``share()`` is effectively a shortcut for ``multicast()`` using a ``AsyncThrowingPassthroughSubject``
/// stream, with an implicit ``autoconnect()``.
Expand Down
79 changes: 79 additions & 0 deletions Tests/Operators/AsyncMulticastCancellationTests.swift
Original file line number Diff line number Diff line change
@@ -0,0 +1,79 @@
@testable import AsyncExtensions
import XCTest

private actor SharedConnection {
let continuation: AsyncThrowingStream<Int, Error>.Continuation
let shared: AsyncShareSequence<AsyncThrowingStream<Int, Error>>

init() {
let (continuation, stream) = AsyncThrowingStream<Int, Error>.pipe()
self.continuation = continuation
self.shared = stream.share()
}

func events() -> AsyncShareSequence<AsyncThrowingStream<Int, Error>> { shared }
func send(_ value: Int) { continuation.yield(value) }
func finish() { continuation.finish() }
}

final class AsyncMulticastCancellationTests: XCTestCase {
func test_connection_returns_one_shared_sequence_for_concurrent_consumers() async throws {
let connection = SharedConnection()
let firstSequence = await connection.events()
let secondSequence = await connection.events()
XCTAssertTrue(firstSequence === secondSequence)
var firstIterator = firstSequence.makeAsyncIterator()
var secondIterator = secondSequence.makeAsyncIterator()
let first = Task {
var values = [Int]()
while let value = try await firstIterator.next() { values.append(value) }
return values
}
let second = Task {
var values = [Int]()
while let value = try await secondIterator.next() { values.append(value) }
return values
}
await connection.send(42)
await connection.finish()
let firstValues = try await first.value
let secondValues = try await second.value
XCTAssertEqual(firstValues, [42])
XCTAssertEqual(secondValues, [42])
}

func test_cancelling_subscriber_waiting_for_upstream_finishes_without_cancelling_shared_stream() async throws {
let (continuation, upstream) = AsyncThrowingStream<Int, Error>.pipe()
let shared = upstream.multicast(AsyncThrowingPassthroughSubject<Int, Error>()).autoconnect()
var firstIterator = shared.makeAsyncIterator()
var secondIterator = shared.makeAsyncIterator()
let firstFinished = expectation(description: "Cancelled subscriber finishes before upstream emits")
let first = Task {
let value = try await firstIterator.next()
firstFinished.fulfill()
return value
}
let deadline = DispatchTime.now().uptimeNanoseconds + 2_000_000_000
while shared.state.withCriticalRegion({ state in
if case .busy = state { return false }
return true
}) {
guard DispatchTime.now().uptimeNanoseconds < deadline else {
continuation.finish()
return XCTFail("Upstream iteration did not start")
}
await Task.yield()
}
let second = Task { try await secondIterator.next() }
first.cancel()
await fulfillment(of: [firstFinished], timeout: 2)

// Also releases the original implementation after its timeout failure.
continuation.yield(42)
continuation.finish()
let firstValue = try await first.value
let secondValue = try await second.value
XCTAssertNil(firstValue)
XCTAssertEqual(secondValue, 42)
}
}
Loading