From bcba9a19d4aa164858d1ba76721c3fad9174998c Mon Sep 17 00:00:00 2001 From: Sergej Jaskiewicz Date: Sat, 14 Dec 2019 23:11:47 +0300 Subject: [PATCH] Update for Xcode 11.3 (#123) - Send subscription synchronously in ReceiveOn and Delay operators - Some locks made recursive, as they should be - ObservableObjectPublisher doesn't use PassthroughSubject under the hood anymore --- .circleci/config.yml | 28 ++--- Sources/OpenCombine/CurrentValueSubject.swift | 4 - Sources/OpenCombine/Future.swift | 4 - .../OpenCombine/Helpers/FilterProducer.swift | 4 - .../OpenCombine/Helpers/ReduceProducer.swift | 4 - .../Helpers/SubjectSubscriber.swift | 4 - Sources/OpenCombine/ObservableObject.swift | 117 +++++++++++++++++- Sources/OpenCombine/PassthroughSubject.swift | 4 - .../Publishers/Publishers.Autoconnect.swift | 4 - .../Publishers/Publishers.Concatenate.swift | 4 +- .../Publishers/Publishers.Delay.swift | 34 +---- .../Publishers/Publishers.Drop.swift | 4 - .../Publishers/Publishers.DropWhile.swift | 4 - .../Publishers/Publishers.FlatMap.swift | 7 +- .../Publishers/Publishers.HandleEvents.swift | 4 - .../Publishers/Publishers.IgnoreOutput.swift | 4 - .../Publishers/Publishers.Map.swift | 4 - .../Publishers.MeasureInterval.swift | 4 - .../Publishers/Publishers.Multicast.swift | 4 - .../Publishers/Publishers.Output.swift | 4 - .../Publishers/Publishers.Print.swift | 4 - .../Publishers/Publishers.ReceiveOn.swift | 18 +-- .../Publishers/Publishers.ReplaceError.swift | 4 - .../Publishers/Publishers.Scan.swift | 4 - .../Publishers/Publishers.Sequence.swift | 4 - .../Publishers/Publishers.SubscribeOn.swift | 4 - Sources/OpenCombine/Publishers/Record.swift | 4 - .../ObservableObjectPublisherTests.swift | 45 +++++-- Tests/OpenCombineTests/PublishedTests.swift | 6 +- .../PublisherTests/ConcatenateTests.swift | 18 ++- .../PublisherTests/DelayTests.swift | 102 ++++++++------- .../PublisherTests/FlatMapTests.swift | 41 +++--- .../PublisherTests/ReceiveOnTests.swift | 77 +++++------- .../PublisherTests/SubscribeOnTests.swift | 3 + 34 files changed, 299 insertions(+), 285 deletions(-) diff --git a/.circleci/config.yml b/.circleci/config.yml index 6fce731..96889c0 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -1,10 +1,10 @@ version: 2 jobs: - "Execute tests on macOS 10.15.0 (Xcode 11.2.0, Swift 5.1.2)": + "Execute tests on macOS 10.15.0 (Xcode 11.3.0, Swift 5.1.3)": macos: - xcode: "11.2.0" + xcode: "11.3.0" environment: - SWIFT_VERSION: "5.1.2" + SWIFT_VERSION: "5.1.3" steps: - checkout - run: @@ -60,35 +60,35 @@ jobs: command: | bash <(curl -s https://codecov.io/bash) -D DerivedData - "Execute compatibility tests on iOS 13.2.2 (Xcode 11.2.0, Swift 5.1.2)": + "Execute compatibility tests on iOS 13.3 (Xcode 11.3.0, Swift 5.1.3)": macos: - xcode: "11.2.0" + xcode: "11.3.0" environment: - SWIFT_VERSION: "5.1.2" + SWIFT_VERSION: "5.1.3" steps: - checkout - run: name: Generating Xcode project command: make generate-compatibility-xcodeproj - run: - name: Building for testing on iOS 13.2.2 with xcodebuild + name: Building for testing on iOS 13.3 with xcodebuild command: | set -o pipefail \ && xcodebuild build-for-testing \ -scheme OpenCombine-Package \ - -destination "platform=iOS Simulator,name=iPhone 11,OS=13.2.2" \ + -destination "platform=iOS Simulator,name=iPhone 11,OS=13.3" \ -derivedDataPath DerivedData \ | tee xcodebuild_build-for-testing.log \ | xcpretty - store_artifacts: path: xcodebuild_build-for-testing.log - run: - name: Testing against Combine on iOS 13.2.2 with xcodebuild + name: Testing against Combine on iOS 13.3 with xcodebuild command: | set -o pipefail \ && xcodebuild test-without-building \ -scheme OpenCombine-Package \ - -destination "platform=iOS Simulator,name=iPhone 11,OS=13.2.2" \ + -destination "platform=iOS Simulator,name=iPhone 11,OS=13.3" \ -derivedDataPath DerivedData \ | tee xcodebuild_test-without-building.log \ | xcpretty --report junit -o build/reports/results.xml @@ -217,7 +217,7 @@ jobs: "Run SwiftLint and Danger": macos: - xcode: "11.2.0" + xcode: "11.3.0" environment: HOMEBREW_NO_AUTO_UPDATE: "1" steps: @@ -236,7 +236,7 @@ jobs: "Run Pod spec lint": macos: - xcode: "11.2.0" + xcode: "11.3.0" environment: HOMEBREW_NO_AUTO_UPDATE: "1" steps: @@ -250,10 +250,10 @@ workflows: version: 2 "OpenCombine: execute tests on macOS": jobs: - - "Execute tests on macOS 10.15.0 (Xcode 11.2.0, Swift 5.1.2)" + - "Execute tests on macOS 10.15.0 (Xcode 11.3.0, Swift 5.1.3)" "OpenCombine: execute compatibility tests": jobs: - - "Execute compatibility tests on iOS 13.2.2 (Xcode 11.2.0, Swift 5.1.2)" + - "Execute compatibility tests on iOS 13.3 (Xcode 11.3.0, Swift 5.1.3)" "OpenCombine: execute tests on iOS": jobs: - "Execute tests on iOS 9.3 (Xcode 10.2.1, Swift 5.0.1)" diff --git a/Sources/OpenCombine/CurrentValueSubject.swift b/Sources/OpenCombine/CurrentValueSubject.swift index 87dca6e..58692e6 100644 --- a/Sources/OpenCombine/CurrentValueSubject.swift +++ b/Sources/OpenCombine/CurrentValueSubject.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 11.06.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - /// A subject that wraps a single value and publishes a new element whenever the value /// changes. public final class CurrentValueSubject: Subject { diff --git a/Sources/OpenCombine/Future.swift b/Sources/OpenCombine/Future.swift index 18ee64e..0f54a26 100644 --- a/Sources/OpenCombine/Future.swift +++ b/Sources/OpenCombine/Future.swift @@ -5,10 +5,6 @@ // Created by Max Desiatov on 24/11/2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - /// A publisher that eventually produces one value and then finishes or fails. public final class Future: Publisher where Failure: Error { diff --git a/Sources/OpenCombine/Helpers/FilterProducer.swift b/Sources/OpenCombine/Helpers/FilterProducer.swift index 2e4d20d..226938b 100644 --- a/Sources/OpenCombine/Helpers/FilterProducer.swift +++ b/Sources/OpenCombine/Helpers/FilterProducer.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 23.10.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - /// A helper class that acts like both subscriber and subscription. /// /// Filter-like operators send an instance of their `Inner` class that is subclass diff --git a/Sources/OpenCombine/Helpers/ReduceProducer.swift b/Sources/OpenCombine/Helpers/ReduceProducer.swift index f86e728..1a6a528 100644 --- a/Sources/OpenCombine/Helpers/ReduceProducer.swift +++ b/Sources/OpenCombine/Helpers/ReduceProducer.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 22.09.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - /// A helper class that acts like both subscriber and subscription. /// /// Reduce-like operators send an instance of their `Inner` class that is subclass diff --git a/Sources/OpenCombine/Helpers/SubjectSubscriber.swift b/Sources/OpenCombine/Helpers/SubjectSubscriber.swift index 61b2054..e95fae5 100644 --- a/Sources/OpenCombine/Helpers/SubjectSubscriber.swift +++ b/Sources/OpenCombine/Helpers/SubjectSubscriber.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 16/09/2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - // NOTE: This class has been audited for thread safety. internal final class SubjectSubscriber : Subscriber, diff --git a/Sources/OpenCombine/ObservableObject.swift b/Sources/OpenCombine/ObservableObject.swift index d8133a4..d449f0a 100644 --- a/Sources/OpenCombine/ObservableObject.swift +++ b/Sources/OpenCombine/ObservableObject.swift @@ -67,19 +67,126 @@ public final class ObservableObjectPublisher: Publisher { public typealias Failure = Never - private let subject: PassthroughSubject + private let lock = UnfairLock.allocate() - public init() { - subject = .init() + private var connections = Set() + + // TODO: Combine needs this for some reason + private var identifier: ObjectIdentifier? + + public init() {} + + deinit { + lock.deallocate() } public func receive(subscriber: Downstream) where Downstream.Input == Void, Downstream.Failure == Never { - subject.subscribe(subscriber) + let inner = Inner(downstream: subscriber, parent: self) + lock.lock() + connections.insert(inner) + lock.unlock() + subscriber.receive(subscription: inner) } public func send() { - subject.send() + lock.lock() + let connections = self.connections + lock.unlock() + for connection in connections { + connection.send() + } + } + + private func remove(_ conduit: Conduit) { + lock.lock() + connections.remove(conduit) + lock.unlock() + } +} + +extension ObservableObjectPublisher { + private class Conduit: Hashable { + + fileprivate func send() { + abstractMethod() + } + + fileprivate static func == (lhs: Conduit, rhs: Conduit) -> Bool { + return lhs === rhs + } + + fileprivate func hash(into hasher: inout Hasher) { + hasher.combine(ObjectIdentifier(self)) + } + } + + private final class Inner + : Conduit, + Subscription, + CustomStringConvertible, + CustomReflectable, + CustomPlaygroundDisplayConvertible + where Downstream.Input == Void, Downstream.Failure == Never + { + private enum State { + case initialized + case active + case terminal + } + + private weak var parent: ObservableObjectPublisher? + private let downstream: Downstream + private let downstreamLock = UnfairRecursiveLock.allocate() + private let lock = UnfairLock.allocate() + private var state = State.initialized + + init(downstream: Downstream, parent: ObservableObjectPublisher) { + self.parent = parent + self.downstream = downstream + } + + deinit { + downstreamLock.deallocate() + lock.deallocate() + } + + override func send() { + lock.lock() + let state = self.state + lock.unlock() + if state == .active { + downstreamLock.lock() + _ = downstream.receive() + downstreamLock.unlock() + } + } + + func request(_ demand: Subscribers.Demand) { + lock.lock() + if state == .initialized { + state = .active + } + lock.unlock() + } + + func cancel() { + lock.lock() + state = .terminal + lock.unlock() + parent?.remove(self) + } + + var description: String { return "ObservableObjectPublisher" } + + var customMirror: Mirror { + let children = CollectionOfOne(("downstream", downstream)) + return Mirror(self, children: children) + } + + var playgroundDescription: Any { + return description + } } } diff --git a/Sources/OpenCombine/PassthroughSubject.swift b/Sources/OpenCombine/PassthroughSubject.swift index 49892e8..63ee24f 100644 --- a/Sources/OpenCombine/PassthroughSubject.swift +++ b/Sources/OpenCombine/PassthroughSubject.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 11.06.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - /// A subject that passes along values and completion. /// /// Use a `PassthroughSubject` in unit tests when you want a publisher than can publish diff --git a/Sources/OpenCombine/Publishers/Publishers.Autoconnect.swift b/Sources/OpenCombine/Publishers/Publishers.Autoconnect.swift index f60741b..72f48ec 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Autoconnect.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Autoconnect.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 18/09/2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension ConnectablePublisher { /// Automates the process of connecting or disconnecting from this connectable diff --git a/Sources/OpenCombine/Publishers/Publishers.Concatenate.swift b/Sources/OpenCombine/Publishers/Publishers.Concatenate.swift index c0d2c7e..4b85d2e 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Concatenate.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Concatenate.swift @@ -145,9 +145,7 @@ extension Publishers.Concatenate { private let lock = UnfairLock.allocate() - // ??? This lock is non-recursive in Combine, but it should be! - // (FB7404824 if Apple folks are watching) - private let downstreamLock = UnfairLock.allocate() + private let downstreamLock = UnfairRecursiveLock.allocate() fileprivate init(downstream: Downstream, suffix: Suffix) { self.downstream = downstream diff --git a/Sources/OpenCombine/Publishers/Publishers.Delay.swift b/Sources/OpenCombine/Publishers/Publishers.Delay.swift index 5a6f100..5386e75 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Delay.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Delay.swift @@ -5,10 +5,6 @@ // Created by Евгений Богомолов on 07/09/2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Delays delivery of all output to the downstream receiver by a specified amount @@ -104,7 +100,7 @@ extension Publishers.Delay { private let lock = UnfairLock.allocate() private var state: State - private let downstreamLock = UnfairLock.allocate() + private let downstreamLock = UnfairRecursiveLock.allocate() fileprivate init(_ publisher: Delay, downstream: Downstream) { state = .ready(publisher, downstream) @@ -115,13 +111,7 @@ extension Publishers.Delay { downstreamLock.deallocate() } - private func schedule(_ delay: Delay, - immediate: Bool, - work: @escaping () -> Void) { - if immediate { - delay.scheduler.schedule(options: delay.options, work) - return - } + private func schedule(_ delay: Delay, work: @escaping () -> Void) { delay .scheduler .schedule(after: delay.scheduler.now.advanced(by: delay.interval), @@ -139,18 +129,6 @@ extension Publishers.Delay { } state = .subscribed(delay, downstream, subscription) lock.unlock() - schedule(delay, immediate: true) { [weak self] in - self?.scheduledReceive(subscription: subscription) - } - } - - private func scheduledReceive(subscription: Subscription) { - lock.lock() - guard case let .subscribed(_, downstream, _) = state else { - lock.unlock() - return - } - lock.unlock() downstreamLock.lock() downstream.receive(subscription: self) downstreamLock.unlock() @@ -163,8 +141,8 @@ extension Publishers.Delay { return .none } lock.unlock() - schedule(delay, immediate: false) { [weak self] in - self?.scheduledReceive(input, downstream: downstream) + schedule(delay) { + self.scheduledReceive(input, downstream: downstream) } return .none } @@ -193,8 +171,8 @@ extension Publishers.Delay { } state = .terminal lock.unlock() - schedule(delay, immediate: false) { [weak self] in - self?.scheduledReceive(completion: completion, downstream: downstream) + schedule(delay) { + self.scheduledReceive(completion: completion, downstream: downstream) } } diff --git a/Sources/OpenCombine/Publishers/Publishers.Drop.swift b/Sources/OpenCombine/Publishers/Publishers.Drop.swift index f0db237..fb168ed 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Drop.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Drop.swift @@ -5,10 +5,6 @@ // Created by Sven Weidauer on 03.10.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Omits the specified number of elements before republishing subsequent elements. /// diff --git a/Sources/OpenCombine/Publishers/Publishers.DropWhile.swift b/Sources/OpenCombine/Publishers/Publishers.DropWhile.swift index 362b8c7..d215094 100644 --- a/Sources/OpenCombine/Publishers/Publishers.DropWhile.swift +++ b/Sources/OpenCombine/Publishers/Publishers.DropWhile.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 16.06.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Omits elements from the upstream publisher until a given closure returns false, diff --git a/Sources/OpenCombine/Publishers/Publishers.FlatMap.swift b/Sources/OpenCombine/Publishers/Publishers.FlatMap.swift index 2288614..cbbe78f 100644 --- a/Sources/OpenCombine/Publishers/Publishers.FlatMap.swift +++ b/Sources/OpenCombine/Publishers/Publishers.FlatMap.swift @@ -4,10 +4,6 @@ // Created by Eric Patey on 16.08.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Transforms all elements from an upstream publisher into a new or existing /// publisher. @@ -92,10 +88,9 @@ extension Publishers.FlatMap { /// by the `downstreamLock`. private let lock = UnfairLock.allocate() - // Must be recursive lock. Probably a bug in Combine. /// All the calls to the downstream subscriber should be made with this lock /// acquired. - private let downstreamLock = UnfairLock.allocate() + private let downstreamLock = UnfairRecursiveLock.allocate() private let downstream: Downstream diff --git a/Sources/OpenCombine/Publishers/Publishers.HandleEvents.swift b/Sources/OpenCombine/Publishers/Publishers.HandleEvents.swift index 29b62e9..58bb2bc 100644 --- a/Sources/OpenCombine/Publishers/Publishers.HandleEvents.swift +++ b/Sources/OpenCombine/Publishers/Publishers.HandleEvents.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 03.12.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Performs the specified closures when publisher events occur. diff --git a/Sources/OpenCombine/Publishers/Publishers.IgnoreOutput.swift b/Sources/OpenCombine/Publishers/Publishers.IgnoreOutput.swift index 8fce777..1f5f2b3 100644 --- a/Sources/OpenCombine/Publishers/Publishers.IgnoreOutput.swift +++ b/Sources/OpenCombine/Publishers/Publishers.IgnoreOutput.swift @@ -4,10 +4,6 @@ // Created by Eric Patey on 16.08.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Ingores all upstream elements, but passes along a completion diff --git a/Sources/OpenCombine/Publishers/Publishers.Map.swift b/Sources/OpenCombine/Publishers/Publishers.Map.swift index e1bca80..a3ce569 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Map.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Map.swift @@ -5,10 +5,6 @@ // Created by Anton Nazarov on 25.06.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Transforms all elements from the upstream publisher with a provided closure. diff --git a/Sources/OpenCombine/Publishers/Publishers.MeasureInterval.swift b/Sources/OpenCombine/Publishers/Publishers.MeasureInterval.swift index 85a958b..54d32dd 100644 --- a/Sources/OpenCombine/Publishers/Publishers.MeasureInterval.swift +++ b/Sources/OpenCombine/Publishers/Publishers.MeasureInterval.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 03.12.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Measures and emits the time interval between events received from an upstream diff --git a/Sources/OpenCombine/Publishers/Publishers.Multicast.swift b/Sources/OpenCombine/Publishers/Publishers.Multicast.swift index c6a3197..e54b33e 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Multicast.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Multicast.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 14.06.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Applies a closure to create a subject that delivers elements to subscribers. diff --git a/Sources/OpenCombine/Publishers/Publishers.Output.swift b/Sources/OpenCombine/Publishers/Publishers.Output.swift index ba1ba94..e2debfa 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Output.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Output.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 24.10.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Republishes elements up to the specified maximum count. diff --git a/Sources/OpenCombine/Publishers/Publishers.Print.swift b/Sources/OpenCombine/Publishers/Publishers.Print.swift index 9347948..77ac222 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Print.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Print.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 16.06.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Prints log messages for all publishing events. diff --git a/Sources/OpenCombine/Publishers/Publishers.ReceiveOn.swift b/Sources/OpenCombine/Publishers/Publishers.ReceiveOn.swift index 8cdca56..87f1d17 100644 --- a/Sources/OpenCombine/Publishers/Publishers.ReceiveOn.swift +++ b/Sources/OpenCombine/Publishers/Publishers.ReceiveOn.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 02.12.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Specifies the scheduler on which to receive elements from the publisher. /// @@ -101,7 +97,7 @@ extension Publishers.ReceiveOn { private let lock = UnfairLock.allocate() private var state: State - private let downstreamLock = UnfairLock.allocate() + private let downstreamLock = UnfairRecursiveLock.allocate() init(_ receiveOn: ReceiveOn, downstream: Downstream) { state = .ready(receiveOn, downstream) @@ -121,18 +117,6 @@ extension Publishers.ReceiveOn { } state = .subscribed(receiveOn, downstream, subscription) lock.unlock() - receiveOn.scheduler.schedule(options: receiveOn.options) { [weak self] in - self?.scheduledReceive(subscription: subscription) - } - } - - private func scheduledReceive(subscription: Subscription) { - lock.lock() - guard case let .subscribed(_, downstream, _) = state else { - lock.unlock() - return - } - lock.unlock() downstreamLock.lock() downstream.receive(subscription: self) downstreamLock.unlock() diff --git a/Sources/OpenCombine/Publishers/Publishers.ReplaceError.swift b/Sources/OpenCombine/Publishers/Publishers.ReplaceError.swift index 3de0b75..6dcead8 100644 --- a/Sources/OpenCombine/Publishers/Publishers.ReplaceError.swift +++ b/Sources/OpenCombine/Publishers/Publishers.ReplaceError.swift @@ -5,10 +5,6 @@ // Created by Bogdan Vlad on 8/29/19. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Replaces any errors in the stream with the provided element. /// diff --git a/Sources/OpenCombine/Publishers/Publishers.Scan.swift b/Sources/OpenCombine/Publishers/Publishers.Scan.swift index 3a96a8c..4d24b15 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Scan.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Scan.swift @@ -4,10 +4,6 @@ // Created by Eric Patey on 26.08.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Transforms elements from the upstream publisher by providing the current element diff --git a/Sources/OpenCombine/Publishers/Publishers.Sequence.swift b/Sources/OpenCombine/Publishers/Publishers.Sequence.swift index 0036d4f..0a21198 100644 --- a/Sources/OpenCombine/Publishers/Publishers.Sequence.swift +++ b/Sources/OpenCombine/Publishers/Publishers.Sequence.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 19.06.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publishers { /// A publisher that publishes a given sequence of elements. diff --git a/Sources/OpenCombine/Publishers/Publishers.SubscribeOn.swift b/Sources/OpenCombine/Publishers/Publishers.SubscribeOn.swift index 7b4f02f..5bbfeea 100644 --- a/Sources/OpenCombine/Publishers/Publishers.SubscribeOn.swift +++ b/Sources/OpenCombine/Publishers/Publishers.SubscribeOn.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 02.12.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - extension Publisher { /// Specifies the scheduler on which to perform subscribe, cancel, and request diff --git a/Sources/OpenCombine/Publishers/Record.swift b/Sources/OpenCombine/Publishers/Record.swift index daa824f..1ad5181 100644 --- a/Sources/OpenCombine/Publishers/Record.swift +++ b/Sources/OpenCombine/Publishers/Record.swift @@ -5,10 +5,6 @@ // Created by Sergej Jaskiewicz on 12.11.2019. // -#if canImport(COpenCombineHelpers) -import COpenCombineHelpers -#endif - /// A publisher that allows for recording a series of inputs and a completion for later /// playback to each subscriber. public struct Record: Publisher { diff --git a/Tests/OpenCombineTests/ObservableObjectPublisherTests.swift b/Tests/OpenCombineTests/ObservableObjectPublisherTests.swift index 34ae2d3..27ac9fa 100644 --- a/Tests/OpenCombineTests/ObservableObjectPublisherTests.swift +++ b/Tests/OpenCombineTests/ObservableObjectPublisherTests.swift @@ -23,22 +23,27 @@ final class ObservableObjectPublisherTests: XCTestCase { receiveSubscription: { downstreamSubscription1 = $0 } ) publisher.subscribe(tracking1) - tracking1.assertHistoryEqual([.subscription("PassthroughSubject")]) + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher")]) downstreamSubscription1?.request(.max(1)) - tracking1.assertHistoryEqual([.subscription("PassthroughSubject")]) + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher")]) publisher.send() - tracking1.assertHistoryEqual([.subscription("PassthroughSubject"), + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher"), .signal]) publisher.send() publisher.send() downstreamSubscription1?.request(.max(3)) - tracking1.assertHistoryEqual([.subscription("PassthroughSubject"), + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher"), + .signal, + .signal, .signal]) publisher.send() publisher.send() publisher.send() publisher.send() - tracking1.assertHistoryEqual([.subscription("PassthroughSubject"), + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher"), + .signal, + .signal, + .signal, .signal, .signal, .signal, @@ -49,28 +54,48 @@ final class ObservableObjectPublisherTests: XCTestCase { receiveSubscription: { $0.request(.unlimited) } ) publisher.subscribe(tracking2) - tracking2.assertHistoryEqual([.subscription("PassthroughSubject")]) + tracking2.assertHistoryEqual([.subscription("ObservableObjectPublisher")]) publisher.send() - tracking1.assertHistoryEqual([.subscription("PassthroughSubject"), + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher"), + .signal, + .signal, + .signal, .signal, .signal, .signal, .signal, .signal]) - tracking2.assertHistoryEqual([.subscription("PassthroughSubject"), + tracking2.assertHistoryEqual([.subscription("ObservableObjectPublisher"), .signal]) downstreamSubscription1?.cancel() publisher.send() - tracking1.assertHistoryEqual([.subscription("PassthroughSubject"), + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher"), + .signal, + .signal, + .signal, .signal, .signal, .signal, .signal, .signal]) - tracking2.assertHistoryEqual([.subscription("PassthroughSubject"), + tracking2.assertHistoryEqual([.subscription("ObservableObjectPublisher"), .signal, .signal]) + + tracking1.cancel() + tracking2.cancel() + } + + func testObservableObjectPublisherReflection() throws { + try testSubscriptionReflection( + description: "ObservableObjectPublisher", + customMirror: expectedChildren( + ("downstream", .contains("TrackingSubscriberBase")) + ), + playgroundDescription: "ObservableObjectPublisher", + sut: ObservableObjectPublisher() + ) } } diff --git a/Tests/OpenCombineTests/PublishedTests.swift b/Tests/OpenCombineTests/PublishedTests.swift index e512428..ce5ce4f 100644 --- a/Tests/OpenCombineTests/PublishedTests.swift +++ b/Tests/OpenCombineTests/PublishedTests.swift @@ -88,11 +88,11 @@ final class PublishedTests: XCTestCase { receiveSubscription: { downstreamSubscription = $0 } ) testObject.objectWillChange.subscribe(tracking1) - tracking1.assertHistoryEqual([.subscription("PassthroughSubject")]) + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher")]) downstreamSubscription?.request(.max(2)) - tracking1.assertHistoryEqual([.subscription("PassthroughSubject")]) + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher")]) testObject.state = 100 - tracking1.assertHistoryEqual([.subscription("PassthroughSubject")]) + tracking1.assertHistoryEqual([.subscription("ObservableObjectPublisher")]) } } diff --git a/Tests/OpenCombineTests/PublisherTests/ConcatenateTests.swift b/Tests/OpenCombineTests/PublisherTests/ConcatenateTests.swift index b976888..5ea41ee 100644 --- a/Tests/OpenCombineTests/PublisherTests/ConcatenateTests.swift +++ b/Tests/OpenCombineTests/PublisherTests/ConcatenateTests.swift @@ -359,9 +359,21 @@ final class ConcatenateTests: XCTestCase { helper.publisher.send(completion: .failure(.oops)) } - assertCrashes { - helper.publisher.send(completion: .failure(.oops)) - } + helper.publisher.send(completion: .failure(.oops)) + + XCTAssertEqual(helper.tracking.history, [.subscription("Concatenate"), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops)), + .completion(.failure(.oops))]) + XCTAssertEqual(helper.subscription.history, []) } func testHelperMethods() { diff --git a/Tests/OpenCombineTests/PublisherTests/DelayTests.swift b/Tests/OpenCombineTests/PublisherTests/DelayTests.swift index 71097b7..d5398ac 100644 --- a/Tests/OpenCombineTests/PublisherTests/DelayTests.swift +++ b/Tests/OpenCombineTests/PublisherTests/DelayTests.swift @@ -18,7 +18,7 @@ final class DelayTests: XCTestCase { // Delay's Inner doesn't conform to CustomStringConvertible, so we can't compare // subscriptions using their descriptions - let delaySubscription: StringSubscription = { + private let delaySubscription: StringSubscription = { let tracking = TrackingSubscriber() let scheduler = VirtualTimeScheduler() CustomPublisher(subscription: CustomSubscription()) @@ -29,6 +29,15 @@ final class DelayTests: XCTestCase { ?? "Delay" }() + private let delaySubscriptionImmediateScheduler: StringSubscription = { + let tracking = TrackingSubscriber() + CustomPublisher(subscription: CustomSubscription()) + .delay(for: 0, scheduler: ImmediateScheduler.shared) + .subscribe(tracking) + return tracking.subscriptions.first.map(StringSubscription.subscription) + ?? "Delay" + }() + func testBasicBehavior() { let scheduler = VirtualTimeScheduler() let helper = OperatorTestHelper(publisherType: CustomPublisher.self, @@ -41,15 +50,15 @@ final class DelayTests: XCTestCase { } XCTAssertNotNil(helper.publisher.subscriber) - XCTAssertEqual(helper.tracking.history, []) - XCTAssertEqual(helper.subscription.history, []) - XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions)]) + XCTAssertEqual(helper.tracking.history, [.subscription(delaySubscription)]) + XCTAssertEqual(helper.subscription.history, [.requested(.max(100))]) + XCTAssertEqual(scheduler.history, []) scheduler.executeScheduledActions() XCTAssertEqual(helper.tracking.history, [.subscription(delaySubscription)]) XCTAssertEqual(helper.subscription.history, [.requested(.max(100))]) - XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions)]) + XCTAssertEqual(scheduler.history, []) XCTAssertEqual(helper.publisher.send(1), .none) XCTAssertEqual(helper.publisher.send(2), .none) @@ -62,8 +71,7 @@ final class DelayTests: XCTestCase { .nanoseconds(200)]) XCTAssertEqual(scheduler.history, - [.schedule(options: .nontrivialOptions), - .now, + [.now, .scheduleAfterDate(.nanoseconds(200), tolerance: .nanoseconds(5), options: .nontrivialOptions), @@ -89,8 +97,7 @@ final class DelayTests: XCTestCase { .requested(.max(12))]) XCTAssertEqual(scheduler.history, - [.schedule(options: .nontrivialOptions), - .now, + [.now, .scheduleAfterDate(.nanoseconds(200), tolerance: .nanoseconds(5), options: .nontrivialOptions), @@ -117,8 +124,7 @@ final class DelayTests: XCTestCase { .requested(.max(12))]) XCTAssertEqual(scheduler.scheduledDates, [.nanoseconds(400)]) XCTAssertEqual(scheduler.history, - [.schedule(options: .nontrivialOptions), - .now, + [.now, .scheduleAfterDate(.nanoseconds(200), tolerance: .nanoseconds(5), options: .nontrivialOptions), @@ -146,8 +152,7 @@ final class DelayTests: XCTestCase { .requested(.max(12)), .requested(.max(12))]) XCTAssertEqual(scheduler.history, - [.schedule(options: .nontrivialOptions), - .now, + [.now, .scheduleAfterDate(.nanoseconds(200), tolerance: .nanoseconds(5), options: .nontrivialOptions), @@ -223,14 +228,14 @@ final class DelayTests: XCTestCase { XCTAssertEqual(helper.subscription.history, [.requested(.unlimited), .cancelled]) XCTAssertEqual(helper.tracking.history, [.subscription(delaySubscription)]) - XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions)]) + XCTAssertEqual(scheduler.history, []) XCTAssertEqual(helper.publisher.send(0), .none) helper.publisher.send(completion: .finished) XCTAssertEqual(helper.subscription.history, [.requested(.unlimited), .cancelled]) XCTAssertEqual(helper.tracking.history, [.subscription(delaySubscription)]) - XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions)]) + XCTAssertEqual(scheduler.history, []) } func testReceiveCompletionImmediatelyAfterSubscription() { @@ -246,19 +251,19 @@ final class DelayTests: XCTestCase { helper.publisher.send(completion: .failure(.oops)) - XCTAssertEqual(helper.tracking.history, []) - XCTAssertEqual(helper.subscription.history, []) + XCTAssertEqual(helper.tracking.history, [.subscription(delaySubscription)]) + XCTAssertEqual(helper.subscription.history, [.requested(.unlimited)]) XCTAssertEqual(scheduler.history, - [.schedule(options: .nontrivialOptions), - .now, + [.now, .scheduleAfterDate(.nanoseconds(123), tolerance: .nanoseconds(5), options: .nontrivialOptions)]) scheduler.executeScheduledActions() - XCTAssertEqual(helper.tracking.history, [.completion(.failure(.oops))]) - XCTAssertEqual(helper.subscription.history, []) + XCTAssertEqual(helper.tracking.history, [.subscription(delaySubscription), + .completion(.failure(.oops))]) + XCTAssertEqual(helper.subscription.history, [.requested(.unlimited)]) } func testReceiveCompletionImmediatelyAfterValue() { @@ -282,8 +287,7 @@ final class DelayTests: XCTestCase { XCTAssertEqual(helper.subscription.history, [.requested(.unlimited), .requested(.max(418))]) XCTAssertEqual(scheduler.history, - [.schedule(options: .nontrivialOptions), - .now, + [.now, .scheduleAfterDate(.nanoseconds(123), tolerance: .nanoseconds(5), options: .nontrivialOptions), @@ -313,13 +317,29 @@ final class DelayTests: XCTestCase { $0.delay(for: .nanoseconds(123), scheduler: ImmediateScheduler.shared) } + var recursionCounter = 5 helper.tracking.onValue = { _ in + if recursionCounter == 0 { return } + recursionCounter -= 1 _ = helper.publisher.send(-1) } - assertCrashes { - _ = helper.publisher.send(0) - } + XCTAssertEqual(helper.publisher.send(0), .none) + XCTAssertEqual(helper.tracking.history, + [.subscription(delaySubscriptionImmediateScheduler), + .value(0), + .value(-1), + .value(-1), + .value(-1), + .value(-1), + .value(-1)]) + XCTAssertEqual(helper.subscription.history, [.requested(.unlimited), + .requested(.max(418)), + .requested(.max(418)), + .requested(.max(418)), + .requested(.max(418)), + .requested(.max(418)), + .requested(.max(418))]) } func testReceiveCompletionRecursively() { @@ -334,27 +354,7 @@ final class DelayTests: XCTestCase { helper.publisher.send(completion: .finished) } - func testWeakCaptureWhenSchedulingSubscription() { - let scheduler = VirtualTimeScheduler() - var subscription: Subscription? - var subscriberReleased = false - do { - let publisher = CustomPublisher(subscription: CustomSubscription()) - let delay = publisher.delay(for: 0.35, scheduler: scheduler) - let tracking = TrackingSubscriber(receiveSubscription: { subscription = $0 }, - onDeinit: { subscriberReleased = true }) - delay.subscribe(tracking) - XCTAssertEqual(tracking.history, []) - XCTAssertEqual(scheduler.history, [.minimumTolerance, - .schedule(options: nil)]) - publisher.cancel() - } - XCTAssertTrue(subscriberReleased) - scheduler.executeScheduledActions() - XCTAssertNil(subscription) - } - - func testWeakCaptureWhenSchedulingValue() { + func testStrongCaptureWhenSchedulingValue() { let scheduler = VirtualTimeScheduler() var value: Int? var subscriberReleased = false @@ -370,7 +370,6 @@ final class DelayTests: XCTestCase { XCTAssertEqual(tracking.history, [.subscription(delaySubscription)]) XCTAssertEqual(scheduler.history, [.minimumTolerance, - .schedule(options: nil), .now, .scheduleAfterDate(.seconds(0.35), tolerance: 0, @@ -380,11 +379,11 @@ final class DelayTests: XCTestCase { } XCTAssertFalse(subscriberReleased) scheduler.executeScheduledActions() - XCTAssertNil(value) + XCTAssertEqual(value, 42) XCTAssertTrue(subscriberReleased) } - func testWeakCaptureWhenSchedulingCompletion() { + func testStrongCaptureWhenSchedulingCompletion() { let scheduler = VirtualTimeScheduler() var completion: Subscribers.Completion? var subscriberReleased = false @@ -400,7 +399,6 @@ final class DelayTests: XCTestCase { XCTAssertEqual(tracking.history, [.subscription(delaySubscription)]) XCTAssertEqual(scheduler.history, [.minimumTolerance, - .schedule(options: nil), .now, .scheduleAfterDate(.seconds(0.35), tolerance: 0, @@ -410,7 +408,7 @@ final class DelayTests: XCTestCase { } XCTAssertFalse(subscriberReleased) scheduler.executeScheduledActions() - XCTAssertNil(completion) + XCTAssertEqual(completion, .finished) XCTAssertTrue(subscriberReleased) } diff --git a/Tests/OpenCombineTests/PublisherTests/FlatMapTests.swift b/Tests/OpenCombineTests/PublisherTests/FlatMapTests.swift index 102e439..b649a3b 100644 --- a/Tests/OpenCombineTests/PublisherTests/FlatMapTests.swift +++ b/Tests/OpenCombineTests/PublisherTests/FlatMapTests.swift @@ -390,11 +390,15 @@ final class FlatMapTests: XCTestCase { } func testChildValueReceivedWhileSendingValue() throws { - let upstreamPublisher = PassthroughSubject, - TestingError>() + let upstreamSubscription = CustomSubscription() + let upstreamPublisher = CustomPublisherBase( + subscription: upstreamSubscription + ) - let child1Publisher = CustomPublisher(subscription: CustomSubscription()) - let child2Publisher = CustomPublisher(subscription: CustomSubscription()) + let childSubscription1 = CustomSubscription() + let childSubscription2 = CustomSubscription() + let child1Publisher = CustomPublisher(subscription: childSubscription1) + let child2Publisher = CustomPublisher(subscription: childSubscription2) let flatMap = upstreamPublisher.flatMap { $0 } @@ -408,12 +412,17 @@ final class FlatMapTests: XCTestCase { flatMap.subscribe(downstreamSubscriber) - upstreamPublisher.send(AnyPublisher(child1Publisher)) - upstreamPublisher.send(AnyPublisher(child2Publisher)) + XCTAssertEqual(upstreamPublisher.send(child1Publisher), .none) + XCTAssertEqual(upstreamPublisher.send(child2Publisher), .none) - assertCrashes { - XCTAssertEqual(child1Publisher.send(666), .max(1)) - } + XCTAssertEqual(child1Publisher.send(666), .max(1)) + + XCTAssertEqual(upstreamSubscription.history, [.requested(.unlimited)]) + XCTAssertEqual(downstreamSubscriber.history, [.subscription("FlatMap"), + .value(666), + .value(777)]) + XCTAssertEqual(childSubscription1.history, [.requested(.max(1))]) + XCTAssertEqual(childSubscription2.history, [.requested(.max(1))]) } func testOuterLockReentrance() { @@ -462,9 +471,6 @@ final class FlatMapTests: XCTestCase { // Create some downstream demand try XCTUnwrap(helper.downstreamSubscription).request(.max(5)) - // If Apple changes the implementation to use recursive lock, - // we must make sure no stack overflow occurs here, - // which will also be detected as a crash, which is not what we want. var recursionDepth = 10 helper.tracking.onFailure = { _ in if recursionDepth <= 0 { @@ -474,10 +480,13 @@ final class FlatMapTests: XCTestCase { _ = child.send(1) } - // Expected deadlock - assertCrashes { - child.send(completion: .failure(.oops)) - } + child.send(completion: .failure(.oops)) + + XCTAssertEqual(helper.tracking.history, [.subscription("FlatMap"), + .completion(.failure(.oops)), + .value(1)]) + XCTAssertEqual(helper.subscription.history, [.requested(.max(1))]) + XCTAssertEqual(childSubscription.history, [.requested(.max(1))]) } func testCompletesProperlyWhenUpstreamOutlivesChildren() { diff --git a/Tests/OpenCombineTests/PublisherTests/ReceiveOnTests.swift b/Tests/OpenCombineTests/PublisherTests/ReceiveOnTests.swift index 49fa41b..56825d0 100644 --- a/Tests/OpenCombineTests/PublisherTests/ReceiveOnTests.swift +++ b/Tests/OpenCombineTests/PublisherTests/ReceiveOnTests.swift @@ -26,15 +26,15 @@ final class ReceiveOnTests: XCTestCase { XCTAssertNotNil(helper.publisher.subscriber, "Subscription must be performed synchronously") - XCTAssertEqual(helper.tracking.history, []) - XCTAssertEqual(helper.subscription.history, []) - XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions)]) + XCTAssertEqual(helper.tracking.history, [.subscription("ReceiveOn")]) + XCTAssertEqual(helper.subscription.history, [.requested(.max(100))]) + XCTAssertEqual(scheduler.history, []) scheduler.executeScheduledActions() XCTAssertEqual(helper.tracking.history, [.subscription("ReceiveOn")]) XCTAssertEqual(helper.subscription.history, [.requested(.max(100))]) - XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions)]) + XCTAssertEqual(scheduler.history, []) XCTAssertEqual(helper.publisher.send(1), .none) XCTAssertEqual(helper.publisher.send(2), .none) @@ -47,7 +47,6 @@ final class ReceiveOnTests: XCTestCase { .nanoseconds(0)]) XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions), - .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions)]) @@ -64,7 +63,6 @@ final class ReceiveOnTests: XCTestCase { .requested(.max(12))]) XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions), - .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions)]) @@ -82,7 +80,6 @@ final class ReceiveOnTests: XCTestCase { .requested(.max(12))]) XCTAssertEqual(scheduler.scheduledDates, [.nanoseconds(0)]) XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions), - .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions)]) @@ -97,7 +94,6 @@ final class ReceiveOnTests: XCTestCase { .requested(.max(12)), .requested(.max(12))]) XCTAssertEqual(scheduler.history, [.schedule(options: .nontrivialOptions), - .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions), .schedule(options: .nontrivialOptions)]) @@ -155,14 +151,14 @@ final class ReceiveOnTests: XCTestCase { XCTAssertEqual(helper.subscription.history, [.requested(.unlimited), .cancelled]) XCTAssertEqual(helper.tracking.history, [.subscription("ReceiveOn")]) - XCTAssertEqual(scheduler.history, [.schedule(options: nil)]) + XCTAssertEqual(scheduler.history, []) XCTAssertEqual(helper.publisher.send(0), .none) helper.publisher.send(completion: .finished) XCTAssertEqual(helper.subscription.history, [.requested(.unlimited), .cancelled]) XCTAssertEqual(helper.tracking.history, [.subscription("ReceiveOn")]) - XCTAssertEqual(scheduler.history, [.schedule(options: nil)]) + XCTAssertEqual(scheduler.history, []) } func testReceiveCompletionImmediatelyAfterSubscription() { @@ -175,15 +171,15 @@ final class ReceiveOnTests: XCTestCase { helper.publisher.send(completion: .failure(.oops)) - XCTAssertEqual(helper.tracking.history, []) - XCTAssertEqual(helper.subscription.history, []) - XCTAssertEqual(scheduler.history, [.schedule(options: nil), - .schedule(options: nil)]) + XCTAssertEqual(helper.tracking.history, [.subscription("ReceiveOn")]) + XCTAssertEqual(helper.subscription.history, [.requested(.unlimited)]) + XCTAssertEqual(scheduler.history, [.schedule(options: nil)]) scheduler.executeScheduledActions() - XCTAssertEqual(helper.tracking.history, [.completion(.failure(.oops))]) - XCTAssertEqual(helper.subscription.history, []) + XCTAssertEqual(helper.tracking.history, [.subscription("ReceiveOn"), + .completion(.failure(.oops))]) + XCTAssertEqual(helper.subscription.history, [.requested(.unlimited)]) } func testReceiveCompletionImmediatelyAfterValue() { @@ -204,7 +200,6 @@ final class ReceiveOnTests: XCTestCase { XCTAssertEqual(helper.subscription.history, [.requested(.unlimited), .requested(.max(418))]) XCTAssertEqual(scheduler.history, [.schedule(options: nil), - .schedule(options: nil), .schedule(options: nil), .schedule(options: nil)]) @@ -218,20 +213,35 @@ final class ReceiveOnTests: XCTestCase { .requested(.max(418))]) } - func testCrashesWhenReceivingInputRecursively() { + func testReceiveInputRecursively() { let helper = OperatorTestHelper(publisherType: CustomPublisher.self, initialDemand: .unlimited, receiveValueDemand: .max(418)) { $0.receive(on: ImmediateScheduler.shared) } + var recursionCounter = 5 helper.tracking.onValue = { _ in + if recursionCounter == 0 { return } + recursionCounter -= 1 _ = helper.publisher.send(-1) } - assertCrashes { - _ = helper.publisher.send(0) - } + XCTAssertEqual(helper.publisher.send(0), .none) + XCTAssertEqual(helper.tracking.history, [.subscription("ReceiveOn"), + .value(0), + .value(-1), + .value(-1), + .value(-1), + .value(-1), + .value(-1)]) + XCTAssertEqual(helper.subscription.history, [.requested(.unlimited), + .requested(.max(418)), + .requested(.max(418)), + .requested(.max(418)), + .requested(.max(418)), + .requested(.max(418)), + .requested(.max(418))]) } func testReceiveCompletionRecursively() { @@ -246,25 +256,6 @@ final class ReceiveOnTests: XCTestCase { helper.publisher.send(completion: .finished) } - func testWeakCaptureWhenSchedulingSubscription() { - let scheduler = VirtualTimeScheduler() - var subscription: Subscription? - var subscriberReleased = false - do { - let publisher = CustomPublisher(subscription: CustomSubscription()) - let receiveOn = publisher.receive(on: scheduler) - let tracking = TrackingSubscriber(receiveSubscription: { subscription = $0 }, - onDeinit: { subscriberReleased = true }) - receiveOn.subscribe(tracking) - XCTAssertEqual(tracking.history, []) - XCTAssertEqual(scheduler.history, [.schedule(options: nil)]) - publisher.cancel() - } - XCTAssertTrue(subscriberReleased) - scheduler.executeScheduledActions() - XCTAssertNil(subscription) - } - func testWeakCaptureWhenSchedulingValue() { let scheduler = VirtualTimeScheduler() var value: Int? @@ -279,8 +270,7 @@ final class ReceiveOnTests: XCTestCase { XCTAssertEqual(tracking.history, [.subscription("ReceiveOn")]) XCTAssertEqual(publisher.send(42), .none) XCTAssertEqual(tracking.history, [.subscription("ReceiveOn")]) - XCTAssertEqual(scheduler.history, [.schedule(options: nil), - .schedule(options: nil)]) + XCTAssertEqual(scheduler.history, [.schedule(options: nil)]) tracking.cancel() publisher.cancel() } @@ -304,8 +294,7 @@ final class ReceiveOnTests: XCTestCase { XCTAssertEqual(tracking.history, [.subscription("ReceiveOn")]) publisher.send(completion: .finished) XCTAssertEqual(tracking.history, [.subscription("ReceiveOn")]) - XCTAssertEqual(scheduler.history, [.schedule(options: nil), - .schedule(options: nil)]) + XCTAssertEqual(scheduler.history, [.schedule(options: nil)]) tracking.cancel() publisher.cancel() } diff --git a/Tests/OpenCombineTests/PublisherTests/SubscribeOnTests.swift b/Tests/OpenCombineTests/PublisherTests/SubscribeOnTests.swift index c1332e6..e6f5745 100644 --- a/Tests/OpenCombineTests/PublisherTests/SubscribeOnTests.swift +++ b/Tests/OpenCombineTests/PublisherTests/SubscribeOnTests.swift @@ -206,7 +206,10 @@ final class SubscribeOnTests: XCTestCase { $0.subscribe(on: ImmediateScheduler.shared) } + var recursionCounter = 5 helper.subscription.onRequest = { _ in + if recursionCounter == 0 { return } + recursionCounter -= 1 helper.downstreamSubscription?.request(.unlimited) }