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
This commit is contained in:
+14
-14
@@ -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)"
|
||||
|
||||
@@ -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<Output, Failure: Error>: Subject {
|
||||
|
||||
@@ -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<Output, Failure>: Publisher where Failure: Error {
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<Downstream: Subject>
|
||||
: Subscriber,
|
||||
|
||||
@@ -67,19 +67,126 @@ public final class ObservableObjectPublisher: Publisher {
|
||||
|
||||
public typealias Failure = Never
|
||||
|
||||
private let subject: PassthroughSubject<Void, Never>
|
||||
private let lock = UnfairLock.allocate()
|
||||
|
||||
public init() {
|
||||
subject = .init()
|
||||
private var connections = Set<Conduit>()
|
||||
|
||||
// TODO: Combine needs this for some reason
|
||||
private var identifier: ObjectIdentifier?
|
||||
|
||||
public init() {}
|
||||
|
||||
deinit {
|
||||
lock.deallocate()
|
||||
}
|
||||
|
||||
public func receive<Downstream: Subscriber>(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<Downstream: Subscriber>
|
||||
: 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<Mirror.Child>(("downstream", downstream))
|
||||
return Mirror(self, children: children)
|
||||
}
|
||||
|
||||
var playgroundDescription: Any {
|
||||
return description
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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.
|
||||
///
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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()
|
||||
|
||||
@@ -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.
|
||||
///
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<Output, Failure: Error>: Publisher {
|
||||
|
||||
@@ -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()
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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")])
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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<TestingError>?
|
||||
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)
|
||||
}
|
||||
|
||||
|
||||
@@ -390,11 +390,15 @@ final class FlatMapTests: XCTestCase {
|
||||
}
|
||||
|
||||
func testChildValueReceivedWhileSendingValue() throws {
|
||||
let upstreamPublisher = PassthroughSubject<AnyPublisher<Int, TestingError>,
|
||||
TestingError>()
|
||||
let upstreamSubscription = CustomSubscription()
|
||||
let upstreamPublisher = CustomPublisherBase<CustomPublisher, TestingError>(
|
||||
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() {
|
||||
|
||||
@@ -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()
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user