From fd6753e0a12890a9455488b0ddd6f4b780b56b2b Mon Sep 17 00:00:00 2001 From: Thibault Wittemberg Date: Sat, 3 Oct 2026 13:11:53 +0200 Subject: [PATCH 1/2] Fix multicast subscriber cancellation while upstream is suspended --- CHANGELOG.md | 1 + README.md | 28 +++++++ .../Operators/AsyncMulticastSequence.swift | 31 +++++--- Sources/Operators/AsyncSequence+Share.swift | 2 + .../AsyncMulticastCancellationTests.swift | 79 +++++++++++++++++++ 5 files changed, 128 insertions(+), 13 deletions(-) create mode 100644 Tests/Operators/AsyncMulticastCancellationTests.swift diff --git a/CHANGELOG.md b/CHANGELOG.md index 1947404..dbfdc57 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). +- 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). - 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..b690a78 100644 --- a/README.md +++ b/README.md @@ -129,3 +129,31 @@ More operators and extensions are to come. Pull requests are of course welcome. An `AsyncJustSequence` iterator releases its stored value after emitting it. An inner sequence can continue producing values after the object that provided it is released. Removing an object or sequence from an array does not cancel an existing subscription. Retain the consuming `Task` and cancel it when its owner stops observing; a weak capture of the owner alone does not stop the task. + +## 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. 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) + } +} From 3cde9070aed8ba98d8c29ee0e215140de0b42bdd Mon Sep 17 00:00:00 2001 From: Thibault Wittemberg Date: Sat, 3 Oct 2026 13:17:17 +0200 Subject: [PATCH 2/2] Keep multicast guidance separate from subscription lifetime changes --- CHANGELOG.md | 2 +- README.md | 56 ++++++++++++++++++++++++++-------------------------- 2 files changed, 29 insertions(+), 29 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index dbfdc57..3462919 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). -- 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). - 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. @@ -11,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 b690a78..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 @@ -129,31 +157,3 @@ More operators and extensions are to come. Pull requests are of course welcome. An `AsyncJustSequence` iterator releases its stored value after emitting it. An inner sequence can continue producing values after the object that provided it is released. Removing an object or sequence from an array does not cancel an existing subscription. Retain the consuming `Task` and cancel it when its owner stops observing; a weak capture of the owner alone does not stop the task. - -## 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.