From 1667ed007318c505817942af31b1af8e9942fb58 Mon Sep 17 00:00:00 2001 From: Thibault Wittemberg Date: Sat, 3 Oct 2026 13:06:34 +0200 Subject: [PATCH] Fix retention of consumed Just values in nested subscriptions --- CHANGELOG.md | 2 + README.md | 6 ++ Sources/Creators/AsyncJustSequence.swift | 22 +++---- Tests/Creators/AsyncJustSequenceTests.swift | 13 +++++ .../AsyncFlatMapLatestLifetimeTests.swift | 57 +++++++++++++++++++ 5 files changed, 86 insertions(+), 14 deletions(-) create mode 100644 Tests/Operators/AsyncFlatMapLatestLifetimeTests.swift diff --git a/CHANGELOG.md b/CHANGELOG.md index 230972e..1947404 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +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). + - 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. - SwiftPM: require Swift 5.8 or later for the Swift Async Algorithms test dependency. diff --git a/README.md b/README.md index 8fa7130..2067602 100644 --- a/README.md +++ b/README.md @@ -123,3 +123,9 @@ or replay consumer receives the latest stored state, followed by subsequent send * [`switchToLatest()`](./Sources/Operators/AsyncSwitchToLatestSequence.swift): Republishes elements sent by the most recently received `AsyncSequence` when self is an `AsyncSequence` of `AsyncSequence` More operators and extensions are to come. Pull requests are of course welcome. + +## Subscription lifetime + +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. diff --git a/Sources/Creators/AsyncJustSequence.swift b/Sources/Creators/AsyncJustSequence.swift index 22d65f0..392e07f 100644 --- a/Sources/Creators/AsyncJustSequence.swift +++ b/Sources/Creators/AsyncJustSequence.swift @@ -29,23 +29,17 @@ public struct AsyncJustSequence: AsyncSequence { } public struct Iterator: AsyncIteratorProtocol { - let element: Element? - let isConsumed = ManagedCriticalState(false) + let pendingElement: ManagedCriticalState + + init(element: Element?) { + self.pendingElement = ManagedCriticalState(element) + } public mutating func next() async -> Element? { - guard !Task.isCancelled else { return nil } - - let shouldEarlyReturn = self.isConsumed.withCriticalRegion { isConsumed -> Bool in - if !isConsumed { - isConsumed = true - return false - } - return true + self.pendingElement.withCriticalRegion { pendingElement in + defer { pendingElement = nil } + return Task.isCancelled ? nil : pendingElement } - - if shouldEarlyReturn { return nil } - - return self.element } } } diff --git a/Tests/Creators/AsyncJustSequenceTests.swift b/Tests/Creators/AsyncJustSequenceTests.swift index 731c634..51d55e7 100644 --- a/Tests/Creators/AsyncJustSequenceTests.swift +++ b/Tests/Creators/AsyncJustSequenceTests.swift @@ -9,6 +9,19 @@ import XCTest final class AsyncJustSequenceTests: XCTestCase { + func test_iterator_copies_share_consumption_but_new_iterators_are_independent() async { + let sequence = AsyncJustSequence(42) + var first = sequence.makeAsyncIterator() + var copy = first + var independent = sequence.makeAsyncIterator() + let firstValue = await first.next() + let copiedValue = await copy.next() + let independentValue = await independent.next() + XCTAssertEqual(firstValue, 42) + XCTAssertNil(copiedValue) + XCTAssertEqual(independentValue, 42) + } + func test_AsyncJustSequence_outputs_expected_element_and_finishes() async { var receivedResult = [Int]() diff --git a/Tests/Operators/AsyncFlatMapLatestLifetimeTests.swift b/Tests/Operators/AsyncFlatMapLatestLifetimeTests.swift new file mode 100644 index 0000000..58bf80b --- /dev/null +++ b/Tests/Operators/AsyncFlatMapLatestLifetimeTests.swift @@ -0,0 +1,57 @@ +@testable import AsyncExtensions +import XCTest + +private final class LifetimeFeature: Sendable { + let values: AnyAsyncSequence + let onDeinit: @Sendable () -> Void + + init(subject: AsyncCurrentValueSubject, onDeinit: @Sendable @escaping () -> Void) { + self.values = subject.eraseToAnyAsyncSequence() + self.onDeinit = onDeinit + } + + deinit { onDeinit() } +} + +private final class LifetimeDevice: @unchecked Sendable { + // The test removes the feature only after the consumer has built its inner iterator. + var features: [AnyAsyncSequence] + + init(feature: LifetimeFeature) { + features = [AsyncJustSequence(feature).eraseToAnyAsyncSequence()] + } +} + +final class AsyncFlatMapLatestLifetimeTests: XCTestCase { + func test_removing_feature_releases_it_while_its_inner_subject_remains_subscribed() async throws { + let subject = AsyncCurrentValueSubject(50) + let released = expectation(description: "Consumed feature released") + let device = LifetimeDevice(feature: LifetimeFeature(subject: subject) { released.fulfill() }) + let result: AnyAsyncSequence = AsyncJustSequence(device) + .eraseToAnyAsyncSequence() + .flatMapLatest { device -> AnyAsyncSequence in + device.features.first!.flatMap { $0.values }.eraseToAnyAsyncSequence() + } + .eraseToAnyAsyncSequence() + let receivedInitial = expectation(description: "Initial value received") + let receivedLater = expectation(description: "Inner subscription still active") + let task = Task { + var received = [Int]() + for try await value in result { + received.append(value) + if value == 50 { receivedInitial.fulfill() } + if value == 75 { receivedLater.fulfill() } + } + return received + } + + await fulfillment(of: [receivedInitial], timeout: 2) + device.features.removeFirst() + await fulfillment(of: [released], timeout: 2) + subject.send(75) + await fulfillment(of: [receivedLater], timeout: 2) + task.cancel() + let received = try await task.value + XCTAssertEqual(received, [50, 75]) + } +}