diff --git a/CHANGELOG.md b/CHANGELOG.md index 1947404..3462919 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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:** diff --git a/README.md b/README.md index 2067602..e019c5d 100644 --- a/README.md +++ b/README.md @@ -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.Continuation + private let shared: AsyncShareSequence> + + init() { + let (continuation, stream) = AsyncThrowingStream.pipe() + self.continuation = continuation + self.shared = stream.share() + } + + func events() -> AsyncShareSequence> { + 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 diff --git a/Sources/Operators/AsyncMulticastSequence.swift b/Sources/Operators/AsyncMulticastSequence.swift index 697912a..19cf68c 100644 --- a/Sources/Operators/AsyncMulticastSequence.swift +++ b/Sources/Operators/AsyncMulticastSequence.swift @@ -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` @@ -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 do { let element = try await iterator.next() @@ -137,7 +141,7 @@ where Base.Element == Subject.Element, Subject.Failure == Error, Base.AsyncItera } state = .available(iterator) } - }.value + } } public func makeAsyncIterator() -> AsyncIterator { @@ -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() diff --git a/Sources/Operators/AsyncSequence+Share.swift b/Sources/Operators/AsyncSequence+Share.swift index d946c21..e2044ed 100644 --- a/Sources/Operators/AsyncSequence+Share.swift +++ b/Sources/Operators/AsyncSequence+Share.swift @@ -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()``. diff --git a/Tests/Operators/AsyncMulticastCancellationTests.swift b/Tests/Operators/AsyncMulticastCancellationTests.swift new file mode 100644 index 0000000..c505f6c --- /dev/null +++ b/Tests/Operators/AsyncMulticastCancellationTests.swift @@ -0,0 +1,79 @@ +@testable import AsyncExtensions +import XCTest + +private actor SharedConnection { + let continuation: AsyncThrowingStream.Continuation + let shared: AsyncShareSequence> + + init() { + let (continuation, stream) = AsyncThrowingStream.pipe() + self.continuation = continuation + self.shared = stream.share() + } + + func events() -> AsyncShareSequence> { 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.pipe() + let shared = upstream.multicast(AsyncThrowingPassthroughSubject()).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) + } +}