diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9e1184b..e4a82f3 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -1,14 +1,22 @@ name: CI -on: push +on: + push: + pull_request: + +permissions: + contents: read jobs: test: - runs-on: macos-12 - + name: Test + runs-on: macos-15 + timeout-minutes: 15 + steps: - - uses: actions/checkout@v3 - - name: Select Xcode 14 - run: sudo xcode-select -s /Applications/Xcode_14.1.app + - name: Checkout + uses: actions/checkout@v6 + - name: Print Swift version + run: swift --version - name: Test run: swift test diff --git a/README.md b/README.md index dc85cbe..7a38196 100644 --- a/README.md +++ b/README.md @@ -6,14 +6,27 @@ A collection of useful extensions for Apple's [Combine framework](https://develo ### Publishers +- [AgainAt](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/AgainAt.swift) - [AnyConnectablePublisher](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/AnyConnectablePublisher.swift) - [BufferPassthroughSubject](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/BufferPassthroughSubject.swift) -- [EnumeratedPublisher](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/Enumerated.swift) +- [Distinct](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/Distinct.swift) +- [Enumerated](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/Enumerated.swift) - [InputStreamPublisher](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/InputStreamPublisher.swift) +- [MulticastLatest](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/MulticastLatest.swift) - [OutputStreamPublisher](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/OutputStreamPublisher.swift) -- [ReduceLatestPublisher](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/ReduceLatest.swift) -- [RetryIfPublisher](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/RetryIf.swift) -- [ThrottleWhilePublisher](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/ThrottleWhile.swift) +- [ReduceLatest](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/ReduceLatest.swift) +- [RetryIf](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/RetryIf.swift) +- [ThrottleWhile](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/ThrottleWhile.swift) + +### Extensions + +- [Cancellable.onCancel()](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/OnCancel.swift) +- [Publisher.sink()](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/Publisher+Sink.swift) + +### Schedulers + +- [TestScheduler](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/TestScheduler.swift) +- [UIScheduler](https://github.com/shareup/combine-extensions/blob/main/Sources/CombineExtensions/UIScheduler.swift) ### Thread-safe subscription management @@ -25,7 +38,7 @@ A collection of useful extensions for Apple's [Combine framework](https://develo Add CombineExtensions to the dependencies section of your package.swift file. ```swift -.package(url: "https://github.com/shareup/combine-extensions.git", from: "6.0.0") +.package(url: "https://github.com/shareup/combine-extensions.git", from: "6.1.0") ``` ## License diff --git a/Sources/CombineExtensions/AgainAt.swift b/Sources/CombineExtensions/AgainAt.swift new file mode 100644 index 0000000..8e56c74 --- /dev/null +++ b/Sources/CombineExtensions/AgainAt.swift @@ -0,0 +1,344 @@ +import Combine +import Foundation +import Synchronized + +public extension Publisher { + func againAt( + scheduler: Context, + options: Context.SchedulerOptions? = nil + ) -> Publishers.AgainAt { + Publishers.AgainAt(upstream: self, scheduler: scheduler, options: options) + } +} + +public extension Publishers { + struct AgainAt: Publisher { + public typealias Output = (Upstream.Output, Timer) + public typealias Failure = Upstream.Failure + + private let upstream: Upstream + private let scheduler: Context + private let options: Context.SchedulerOptions? + + public init( + upstream: Upstream, + scheduler: Context, + options: Context.SchedulerOptions? + ) { + self.upstream = upstream + self.scheduler = scheduler + self.options = options + } + + public func receive( + subscriber: S + ) where S.Input == Output, S.Failure == Failure { + let subscription = AgainAtSubscription( + scheduler: scheduler, + options: options, + subscriber: subscriber + ) + + upstream + .receive(on: scheduler, options: options) + .subscribe(subscription) + } + } +} + +public extension Publishers.AgainAt { + final class Timer: @unchecked Sendable { + public let now: Context.SchedulerTimeType + + private let onRepublishAt: @Sendable (Context.SchedulerTimeType) -> Void + + fileprivate init( + now: Context.SchedulerTimeType, + onRepublishAt: @escaping @Sendable (Context.SchedulerTimeType) -> Void + ) { + self.now = now + self.onRepublishAt = onRepublishAt + } + + public func republish(at time: Context.SchedulerTimeType) { + onRepublishAt(time) + } + } +} + +private final class AgainAtSubscription: + Subscription, + Subscriber, + @unchecked Sendable + where + Upstream: Publisher, + Context: Scheduler, + Downstream: Subscriber, + Downstream.Input == Publishers.AgainAt.Output, + Downstream.Failure == Upstream.Failure +{ + typealias Input = Upstream.Output + typealias Failure = Upstream.Failure + + private final class RepublishToken {} + + private struct State { + var downstream: Downstream? + var upstream: Subscription? + var demand: Subscribers.Demand = .none + + var latestOutput: Upstream.Output? + var pendingOutput: Upstream.Output? + + var activeRepublish: RepublishToken? + var pendingCompletion: Subscribers.Completion? + + var isDraining = false + var isDrainScheduled = false + } + + private enum DrainAction { + case completion(Downstream, Subscribers.Completion) + case stop + case value(Downstream, Upstream.Output) + } + + private let scheduler: Context + private let options: Context.SchedulerOptions? + private let state: Locked + + init( + scheduler: Context, + options: Context.SchedulerOptions?, + subscriber: Downstream + ) { + self.scheduler = scheduler + self.options = options + state = Locked(State(downstream: subscriber)) + } + + deinit { + cancel() + } + + func receive(subscription: Subscription) { + let downstream = state.access { state -> Downstream? in + guard state.downstream != nil, + state.upstream == nil + else { return nil } + state.upstream = subscription + return state.downstream + } + + guard let downstream else { + subscription.cancel() + return + } + + downstream.receive(subscription: self) + + let shouldRequest = state.access { state in + state.downstream != nil && state.upstream != nil + } + + if shouldRequest { + subscription.request(.unlimited) + } + } + + func request(_ demand: Subscribers.Demand) { + guard demand > .none else { return } + + state.access { state in + guard state.downstream != nil else { return } + state.demand += demand + } + + scheduleDrain() + } + + func cancel() { + let upstream = state.access { state -> Subscription? in + let upstream = state.upstream + + state.downstream = nil + state.upstream = nil + state.latestOutput = nil + state.pendingOutput = nil + state.activeRepublish = nil + state.pendingCompletion = nil + state.isDraining = false + state.isDrainScheduled = false + + return upstream + } + + upstream?.cancel() + } + + func receive(_ input: Upstream.Output) -> Subscribers.Demand { + state.access { state in + guard state.downstream != nil, state.pendingCompletion == nil else { + return + } + + state.latestOutput = input + state.pendingOutput = input + } + + drain() + return .none + } + + func receive(completion: Subscribers.Completion) { + state.access { state in + guard state.downstream != nil, state.pendingCompletion == nil else { + return + } + + state.pendingCompletion = completion + state.activeRepublish = nil + state.latestOutput = nil + } + + drain() + } + + private func republish(at time: Context.SchedulerTimeType) { + let token = RepublishToken() + + let shouldSchedule = state.access { state in + guard + state.downstream != nil, + state.pendingCompletion == nil, + state.latestOutput != nil + else { + return false + } + + state.activeRepublish = token + return true + } + + guard shouldSchedule else { return } + + scheduler.schedule( + after: time, + tolerance: scheduler.minimumTolerance, + options: options + ) { [weak self] in + self?.fire(token) + } + } + + private func fire(_ token: RepublishToken) { + state.access { state in + guard state.downstream != nil, + state.pendingCompletion == nil, + state.activeRepublish === token, + let latestOutput = state.latestOutput + else { return } + + state.activeRepublish = nil + state.pendingOutput = latestOutput + } + + drain() + } + + private func drain() { + let shouldDrain = state.access { state in + state.isDrainScheduled = false + + guard state.downstream != nil, + !state.isDraining, + hasWork(state) + else { return false } + + state.isDraining = true + return true + } + + guard shouldDrain else { return } + + while true { + let action = state.access { state -> DrainAction in + guard let downstream = state.downstream else { + state.isDraining = false + return .stop + } + + if let output = state.pendingOutput, state.demand > .none { + state.pendingOutput = nil + state.demand -= .max(1) + return .value(downstream, output) + } + + if let completion = state.pendingCompletion { + state.downstream = nil + state.upstream = nil + state.latestOutput = nil + state.pendingOutput = nil + state.activeRepublish = nil + state.isDraining = false + state.isDrainScheduled = false + return .completion(downstream, completion) + } + + state.isDraining = false + return .stop + } + + switch action { + case let .completion(downstream, completion): + downstream.receive(completion: completion) + return + + case .stop: + return + + case let .value(downstream, output): + let timer = Publishers.AgainAt.Timer( + now: scheduler.now, + onRepublishAt: { [weak self] time in + self?.republish(at: time) + } + ) + + let newDemand = downstream.receive((output, timer)) + + if newDemand > .none { + state.access { state in + if state.downstream != nil { + state.demand += newDemand + } + } + } + } + } + } + + private func hasWork(_ state: State) -> Bool { + (state.pendingOutput != nil && state.demand > .none) + || state.pendingCompletion != nil + } + + private func scheduleDrain() { + let shouldSchedule = state.access { state in + guard state.downstream != nil, + !state.isDraining, + !state.isDrainScheduled, + hasWork(state) + else { return false } + + state.isDrainScheduled = true + return true + } + + if shouldSchedule { + scheduler.schedule(options: options) { [weak self] in + self?.drain() + } + } + } +} diff --git a/Sources/CombineExtensions/MulticastLatest.swift b/Sources/CombineExtensions/MulticastLatest.swift index 407b5b0..4139c85 100644 --- a/Sources/CombineExtensions/MulticastLatest.swift +++ b/Sources/CombineExtensions/MulticastLatest.swift @@ -3,7 +3,7 @@ import Combine public extension Publisher { func multicastLatest() -> some Publisher { map(Optional.some) - .multicast({ CurrentValueSubject(nil) }) + .multicast { CurrentValueSubject(nil) } .autoconnect() .compactMap { $0 } } diff --git a/Sources/CombineTestExtensions/Publisher+Test.swift b/Sources/CombineTestExtensions/Publisher+Test.swift index 51386d1..fa5416b 100644 --- a/Sources/CombineTestExtensions/Publisher+Test.swift +++ b/Sources/CombineTestExtensions/Publisher+Test.swift @@ -516,7 +516,7 @@ private enum _Completion { } } -private class _Expectation: XCTestExpectation { +private class _Expectation: XCTestExpectation, @unchecked Sendable { var token: AnyCancellable? private var fulfillmentCount: Int = 0 diff --git a/Tests/CombineExtensionsTests/AgainAtTests.swift b/Tests/CombineExtensionsTests/AgainAtTests.swift new file mode 100644 index 0000000..71f0f89 --- /dev/null +++ b/Tests/CombineExtensionsTests/AgainAtTests.swift @@ -0,0 +1,648 @@ +import Combine +import CombineExtensions +import XCTest + +final class AgainAtTests: XCTestCase { + func testPublisherExtensionPublishesUpstreamOutputAndTimer() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var timerNows = [DispatchQueue.SchedulerTimeType]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + timerNows.append(timer.now) + } + ) + defer { subscription.cancel() } + + subject.send(1) + + XCTAssertTrue(values.isEmpty) + + scheduler.advance() + + XCTAssertEqual(values, [1]) + XCTAssertEqual(timerNows, [scheduler.now]) + } + + func testPublishersAgainAtInitializerPublishesUpstreamOutputAndTimer() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var timerNows = [DispatchQueue.SchedulerTimeType]() + + let subscription = Publishers.AgainAt( + upstream: subject, + scheduler: scheduler, + options: nil + ) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + timerNows.append(timer.now) + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + XCTAssertEqual(values, [1]) + XCTAssertEqual(timerNows, [scheduler.now]) + } + + func testRepublishPublishesAfterRequestedTime() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + + scheduler.advance(by: .milliseconds(999)) + XCTAssertEqual(values, [1]) + + scheduler.advance(by: .milliseconds(1)) + XCTAssertEqual(values, [1, 1]) + } + + func testRepublishAtCurrentSchedulerTimeFiresOnNextAdvance() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + republishers[0](scheduler.now) + + XCTAssertEqual(values, [1]) + + scheduler.advance() + + XCTAssertEqual(values, [1, 1]) + } + + func testRepublishPublishesNewestUpstreamOutputAtFireTime() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + + subject.send(2) + scheduler.advance() + + XCTAssertEqual(values, [1, 2]) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1, 2, 2]) + } + + func testOnlyMostRecentlySetRepublishTimeFires() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + republishers[0](scheduler.now.advanced(by: .seconds(2))) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1]) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1, 1]) + } + + func testMultipleRepublishersAtSameDateOnlyFireOnce() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + let fireDate = scheduler.now.advanced(by: .seconds(1)) + republishers[0](fireDate) + republishers[0](fireDate) + republishers[0](fireDate) + + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual(values, [1, 1]) + } + + func testNewerTimerCanReplaceOlderScheduledRepublish() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + republishers[0](scheduler.now.advanced(by: .seconds(1))) + + subject.send(2) + scheduler.advance() + republishers[1](scheduler.now.advanced(by: .seconds(2))) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1, 2]) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1, 2, 2]) + } + + func testOlderRetainedTimerCanReplaceNewerScheduledRepublish() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + subject.send(2) + scheduler.advance() + + republishers[1](scheduler.now.advanced(by: .seconds(2))) + republishers[0](scheduler.now.advanced(by: .seconds(1))) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1, 2, 2]) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1, 2, 2]) + } + + func testRepublishedOutputReceivesNewTimerThatCanRepublishAgain() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + var timerNows = [DispatchQueue.SchedulerTimeType]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + timerNows.append(timer.now) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + let firstNow = scheduler.now + republishers[0](firstNow.advanced(by: .seconds(1))) + + scheduler.advance(by: .seconds(1)) + let secondNow = scheduler.now + + republishers[1](secondNow.advanced(by: .seconds(1))) + + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual(values, [1, 1, 1]) + XCTAssertEqual(timerNows, [firstNow, secondNow, scheduler.now]) + } + + func testCompletionInvalidatesPendingRepublish() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + var didFinish = false + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in didFinish = true }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + subject.send(completion: .finished) + scheduler.advance() + + XCTAssertTrue(didFinish) + + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual(values, [1]) + } + + func testFailureInvalidatesPendingRepublish() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + var completion: Subscribers.Completion? + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { completion = $0 }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + subject.send(completion: .failure(.failed)) + scheduler.advance() + + XCTAssertEqual(completion, .failure(.failed)) + + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual(values, [1]) + } + + func testCancellationInvalidatesPendingRepublish() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + + subject.send(1) + scheduler.advance() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + subscription.cancel() + + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual(values, [1]) + } + + func testRetainedTimerDoesNotRepublishAfterCancellation() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + + subject.send(1) + scheduler.advance() + + subscription.cancel() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual(values, [1]) + } + + func testRetainedTimerDoesNotRepublishAfterCompletion() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + subject.send(completion: .finished) + scheduler.advance() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual(values, [1]) + } + + func testUpstreamOutputsDropsEarlierOutputWhileThereIsNoDemand() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var upstreamSubscription: Subscription? + + let subscriber = AnySubscriber< + Publishers.AgainAt, TestSchedulerOf> + .Output, + Never + >( + receiveSubscription: { subscription in + upstreamSubscription = subscription + subscription.request(.max(1)) + }, + receiveValue: { value, _ in + values.append(value) + return .none + }, + receiveCompletion: { _ in XCTFail("Should not complete") } + ) + + subject + .againAt(scheduler: scheduler) + .receive(subscriber: subscriber) + + subject.send(1) + subject.send(2) + subject.send(3) + + scheduler.advance() + + XCTAssertEqual(values, [1]) + + upstreamSubscription?.request(.max(1)) + scheduler.advance() + + XCTAssertEqual(values, [1, 3]) + } + + func testNewerUpstreamOutputReplacesPendingRepublishWithoutDemand() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [Republisher]() + var upstreamSubscription: Subscription? + + let subscriber = AnySubscriber< + Publishers.AgainAt, TestSchedulerOf> + .Output, + Never + >( + receiveSubscription: { subscription in + upstreamSubscription = subscription + subscription.request(.max(1)) + }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: $0) } + return .none + }, + receiveCompletion: { _ in XCTFail("Should not complete") } + ) + + subject + .againAt(scheduler: scheduler) + .receive(subscriber: subscriber) + + subject.send(1) + scheduler.advance() + + republishers[0](scheduler.now.advanced(by: .seconds(1))) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1]) + + subject.send(2) + scheduler.advance() + XCTAssertEqual(values, [1]) + + upstreamSubscription?.request(.max(1)) + scheduler.advance() + + XCTAssertEqual(values, [1, 2]) + } + + func testPendingValueIsDroppedWhenCompletionArrivesWithoutDemand() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var didFinish = false + + let subscriber = AnySubscriber< + Publishers.AgainAt, TestSchedulerOf> + .Output, + Never + >( + receiveSubscription: { _ in }, + receiveValue: { value, _ in + values.append(value) + return .none + }, + receiveCompletion: { _ in didFinish = true } + ) + + subject + .againAt(scheduler: scheduler) + .receive(subscriber: subscriber) + + subject.send(1) + subject.send(completion: .finished) + + scheduler.advance() + + XCTAssertTrue(didFinish) + XCTAssertEqual(values, []) + } + + func testPendingValueIsDroppedWhenFailureArrivesWithoutDemand() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var completion: Subscribers.Completion? + + let subscriber = AnySubscriber< + Publishers.AgainAt< + PassthroughSubject, + TestSchedulerOf + >.Output, + TestError + >( + receiveSubscription: { _ in }, + receiveValue: { value, _ in + values.append(value) + return .none + }, + receiveCompletion: { completion = $0 } + ) + + subject + .againAt(scheduler: scheduler) + .receive(subscriber: subscriber) + + subject.send(1) + subject.send(completion: .failure(.failed)) + + scheduler.advance() + + XCTAssertEqual(completion, .failure(.failed)) + XCTAssertEqual(values, []) + } + + func testPendingValueIsDeliveredBeforeCompletionWhenDemandExists() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var didFinish = false + + let subscriber = AnySubscriber< + Publishers.AgainAt, TestSchedulerOf> + .Output, + Never + >( + receiveSubscription: { subscription in + subscription.request(.max(1)) + }, + receiveValue: { value, _ in + values.append(value) + return .none + }, + receiveCompletion: { _ in didFinish = true } + ) + + subject + .againAt(scheduler: scheduler) + .receive(subscriber: subscriber) + + subject.send(1) + subject.send(completion: .finished) + + scheduler.advance() + + XCTAssertTrue(didFinish) + XCTAssertEqual(values, [1]) + } +} + +private typealias Republisher = (DispatchQueue.SchedulerTimeType) -> Void + +private enum TestError: Error, Equatable { + case failed +} diff --git a/Tests/CombineExtensionsTests/AnyConnectablePublisherTests.swift b/Tests/CombineExtensionsTests/AnyConnectablePublisherTests.swift index ec82c4d..8a7acbb 100644 --- a/Tests/CombineExtensionsTests/AnyConnectablePublisherTests.swift +++ b/Tests/CombineExtensionsTests/AnyConnectablePublisherTests.swift @@ -4,6 +4,44 @@ import CombineTestExtensions import XCTest final class AnyConnectablePublisherTests: XCTestCase { + func testErasedPublisherForwardsSubscribersAndConnects() throws { + let wrapped = TrackingConnectablePublisher() + let publisher = wrapped.eraseToAnyConnectablePublisher() + + var values = [Int]() + let subscription = publisher.sink { values.append($0) } + + XCTAssertEqual(1, wrapped.subscriberCount) + + wrapped.subject.send(1) + XCTAssertEqual([1], values) + + let connection = publisher.connect() + XCTAssertEqual(1, wrapped.connectCount) + + connection.cancel() + XCTAssertEqual(1, wrapped.connectionCancelCount) + + subscription.cancel() + } + + func testAutoconnectConnectsAndCancelsErasedPublisher() throws { + let wrapped = TrackingConnectablePublisher() + let publisher = wrapped.eraseToAnyConnectablePublisher() + + let subscription = publisher + .autoconnect() + .sink { (_: Int) in } + + XCTAssertEqual(1, wrapped.subscriberCount) + XCTAssertEqual(1, wrapped.connectCount) + XCTAssertEqual(0, wrapped.connectionCancelCount) + + subscription.cancel() + + XCTAssertEqual(1, wrapped.connectionCancelCount) + } + func testErasedTimerCanStillBeConnectedTo() throws { let pub = Timer.publish(every: 0.01, on: .main, in: .common) .eraseToAnyConnectablePublisher() @@ -44,3 +82,28 @@ final class AnyConnectablePublisherTests: XCTestCase { XCTAssertTrue(subscriptions.isEmpty) } } + +private final class TrackingConnectablePublisher: ConnectablePublisher { + typealias Output = Int + typealias Failure = Never + + let subject = PassthroughSubject() + private(set) var subscriberCount = 0 + private(set) var connectCount = 0 + private(set) var connectionCancelCount = 0 + + func receive( + subscriber: S + ) where S.Input == Int, S.Failure == Never { + subscriberCount += 1 + subject.receive(subscriber: subscriber) + } + + func connect() -> Cancellable { + connectCount += 1 + + return AnyCancellable { [weak self] in + self?.connectionCancelCount += 1 + } + } +} diff --git a/Tests/CombineExtensionsTests/BufferPassthroughSubjectTests.swift b/Tests/CombineExtensionsTests/BufferPassthroughSubjectTests.swift index ff0994c..70c76b1 100644 --- a/Tests/CombineExtensionsTests/BufferPassthroughSubjectTests.swift +++ b/Tests/CombineExtensionsTests/BufferPassthroughSubjectTests.swift @@ -25,6 +25,61 @@ class BufferPassthroughSubjectTests: XCTestCase { wait(for: [ex], timeout: 2) } + func testBufferedValuesRespectFirstSubscriberDemand() throws { + let subject = BufferPassthroughSubject() + subject.send(0) + subject.send(1) + subject.send(2) + + let subscriber = BufferManualDemandSubscriber(initialDemand: .max(1)) + subject.subscribe(subscriber) + + XCTAssertEqual([0], subscriber.values) + XCTAssertEqual([], subscriber.completions) + + subscriber.subscription?.request(.max(1)) + subject.send(3) + + XCTAssertEqual([0, 3], subscriber.values) + } + + func testBufferedCompletionWithoutValuesIsDeliveredToFirstSubscriber() throws { + let subject = BufferPassthroughSubject() + subject.send(completion: .finished) + + let ex = subject.expectToFinish(failsOnOutput: true) + + wait(for: [ex], timeout: 2) + } + + func testValuesSentAfterBufferedCompletionAreIgnored() throws { + let subject = BufferPassthroughSubject() + + subject.send(0) + subject.send(completion: .finished) + subject.send(1) + + let ex = subject.expectOutput([0], expectToFinish: true) + + wait(for: [ex], timeout: 2) + } + + func testRepeatedCompletionAfterPassingThroughIsOnlyDeliveredOnce() throws { + let subject = BufferPassthroughSubject() + + var completions = 0 + let subscription = subject.sink( + receiveCompletion: { _ in completions += 1 }, + receiveValue: { _ in } + ) + + subject.send(completion: .finished) + subject.send(completion: .finished) + + XCTAssertEqual(1, completions) + subscription.cancel() + } + func testPassesThroughValuesAfterReceivingSubscriber() throws { let subject = BufferPassthroughSubject() @@ -94,3 +149,34 @@ class BufferPassthroughSubjectTests: XCTestCase { private enum TestError: Error, Equatable { case error } + +private final class BufferManualDemandSubscriber: Subscriber { + var subscription: Subscription? + private(set) var values = [Input]() + private(set) var completions = [Subscribers.Completion]() + + private let initialDemand: Subscribers.Demand + private let demandOnValue: Subscribers.Demand + + init( + initialDemand: Subscribers.Demand, + demandOnValue: Subscribers.Demand = .none + ) { + self.initialDemand = initialDemand + self.demandOnValue = demandOnValue + } + + func receive(subscription: Subscription) { + self.subscription = subscription + subscription.request(initialDemand) + } + + func receive(_ input: Input) -> Subscribers.Demand { + values.append(input) + return demandOnValue + } + + func receive(completion: Subscribers.Completion) { + completions.append(completion) + } +} diff --git a/Tests/CombineExtensionsTests/DistinctTests.swift b/Tests/CombineExtensionsTests/DistinctTests.swift index 82ad3a3..18e6d51 100644 --- a/Tests/CombineExtensionsTests/DistinctTests.swift +++ b/Tests/CombineExtensionsTests/DistinctTests.swift @@ -93,6 +93,99 @@ final class DistinctTests: XCTestCase { wait(for: [ex], timeout: 2) } + func testDistinctDropsEmptyOutputsAndDuplicatesWithinSingleOutput() throws { + let pub = [ + [1, 1, 2, 2], + [1, 2], + [], + [3, 3, 4], + [4, 3], + [5], + ] + .publisher + .distinct() + + let expected = [ + [1, 2], + [3, 4], + [5], + ] + + let ex = pub.expectOutput(expected, expectToFinish: true) + + wait(for: [ex], timeout: 2) + } + + func testDistinctStateIsSharedByMultipleSubscribersToTheSamePublisherInstance() throws { + let subject = PassthroughSubject<[Int], Never>() + let publisher = subject.distinct() + + var firstValues = [[Int]]() + var secondValues = [[Int]]() + + let first = publisher.sink { firstValues.append($0) } + let second = publisher.sink { secondValues.append($0) } + + defer { + first.cancel() + second.cancel() + } + + subject.send([1, 2]) + subject.send([2, 3, 4]) + subject.send([5]) + + let subscriberOutputs = [firstValues, secondValues] + let expectedUniqueOutput = [[1, 2], [3, 4], [5]] + + XCTAssertEqual( + 1, + subscriberOutputs.filter { $0 == expectedUniqueOutput }.count + ) + XCTAssertEqual( + 1, + subscriberOutputs.filter(\.isEmpty).count + ) + } + + func testDistinctStateIsIndependentForSeparatePublisherInstances() throws { + let subject = PassthroughSubject<[Int], Never>() + + let firstPublisher = subject.distinct() + let secondPublisher = subject.distinct() + + var firstValues = [[Int]]() + var secondValues = [[Int]]() + + let first = firstPublisher.sink { firstValues.append($0) } + let second = secondPublisher.sink { secondValues.append($0) } + + defer { + first.cancel() + second.cancel() + } + + subject.send([1, 2]) + subject.send([2, 3]) + + XCTAssertEqual([[1, 2], [3]], firstValues) + XCTAssertEqual([[1, 2], [3]], secondValues) + } + + func testDistinctPropagatesFailureAfterDroppingDuplicateOutputs() throws { + let subject = PassthroughSubject<[Int], DistinctError>() + + let ex = subject + .distinct() + .expectOutput([[1]], completion: .failure(.failed)) + + subject.send([1]) + subject.send([1]) + subject.send(completion: .failure(.failed)) + + wait(for: [ex], timeout: 2) + } + func testDistinctWithDictArraysPublisher() throws { let dicts: [[[String: AnyHashable]]] = [ [["first": 1], ["second": "2"], ["third": 3.0]], @@ -161,7 +254,7 @@ final class DistinctTests: XCTestCase { } let lastIndex = input.count - 1 - input.enumerated().forEach { i, v in + for (i, v) in input.enumerated() { queue.async { subject.send(v) if i == lastIndex { subject.send(completion: .finished) } @@ -173,3 +266,7 @@ final class DistinctTests: XCTestCase { wait(for: [ex1, ex2, ex3], timeout: 2) } } + +private enum DistinctError: Error, Equatable { + case failed +} diff --git a/Tests/CombineExtensionsTests/EnumeratedTests.swift b/Tests/CombineExtensionsTests/EnumeratedTests.swift index cd43441..0f4e164 100644 --- a/Tests/CombineExtensionsTests/EnumeratedTests.swift +++ b/Tests/CombineExtensionsTests/EnumeratedTests.swift @@ -41,6 +41,73 @@ final class EnumeratedTests: XCTestCase { wait(for: [ex], timeout: 2) } + func testEnumeratedStartsAtNegativeIndex() throws { + let pub = ["a", "b", "c"] + .publisher + .enumerated(startIndex: -2) + + var expected = [(-2, "a"), (-1, "b"), (0, "c")] + let ex = pub.expectOutput( + { index, value in + let output = expected.removeFirst() + XCTAssertEqual(output.0, index) + XCTAssertEqual(output.1, value) + + return expected.isEmpty ? .finished : .moreExpected + }, + expectToFinish: true + ) + + wait(for: [ex], timeout: 2) + } + + func testEnumeratedStateIsIndependentForEachSubscriber() throws { + let subject = PassthroughSubject() + let publisher = subject.enumerated(startIndex: 5) + + var firstValues = [(Int, String)]() + var secondValues = [(Int, String)]() + + let first = publisher.sink { firstValues.append($0) } + let second = publisher.sink { secondValues.append($0) } + + defer { + first.cancel() + second.cancel() + } + + subject.send("a") + subject.send("b") + + XCTAssertEqual(firstValues.map(\.0), [5, 6]) + XCTAssertEqual(firstValues.map(\.1), ["a", "b"]) + XCTAssertEqual(secondValues.map(\.0), [5, 6]) + XCTAssertEqual(secondValues.map(\.1), ["a", "b"]) + } + + func testEnumeratedPropagatesFailure() throws { + let subject = PassthroughSubject() + + var expected = [(0, "a"), (1, "b")] + let ex = subject + .enumerated() + .expectOutputAndFailure { index, value in + let output = expected.removeFirst() + XCTAssertEqual(output.0, index) + XCTAssertEqual(output.1, value) + + return expected.isEmpty ? .finished : .moreExpected + } failureEvaluator: { error in + XCTAssertEqual(.failed, error) + } + + subject.send("a") + subject.send("b") + subject.send(completion: .failure(.failed)) + + wait(for: [ex], timeout: 2) + } + func testEnumeratedWithCustomStartIndexWithArrayPublisher() throws { let pub = [10, 11, 12, 13, 14, 15] .publisher @@ -106,3 +173,7 @@ final class EnumeratedTests: XCTestCase { wait(for: [ex], timeout: 2) } } + +private enum EnumeratedError: Error, Equatable { + case failed +} diff --git a/Tests/CombineExtensionsTests/InputStreamPublisherTests.swift b/Tests/CombineExtensionsTests/InputStreamPublisherTests.swift index 36f438d..0e459d3 100644 --- a/Tests/CombineExtensionsTests/InputStreamPublisherTests.swift +++ b/Tests/CombineExtensionsTests/InputStreamPublisherTests.swift @@ -20,6 +20,95 @@ class InputStreamPublisherTests: XCTestCase { wait(for: [ex], timeout: 2) } + func testInputStreamPublisherWithEmptyDataFinishesWithoutOutput() throws { + let data = Data() + + let ex = data + .publisher(maxChunkSize: 2) + .expectToFinish(failsOnOutput: true) + + wait(for: [ex], timeout: 2) + } + + func testInputStreamPublisherWithZeroChunkSizeFinishesWithoutOutput() throws { + let data = Data("Hello!".utf8) + + let ex = data + .publisher(maxChunkSize: 0) + .expectToFinish(failsOnOutput: true) + + wait(for: [ex], timeout: 2) + } + + func testInputStreamPublisherPublishesMoreChunksWhenAdditionalDemandIsRequested() throws { + let data = Data("Hello".utf8) + let subscriber = InputStreamManualSubscriber(initialDemand: .max(1)) + + data + .publisher(maxChunkSize: 2) + .receive(subscriber: subscriber) + + XCTAssertEqual( + [ + Array("He".utf8), + ], + subscriber.values + ) + XCTAssertTrue(subscriber.completions.isEmpty) + + subscriber.subscription?.request(.max(1)) + XCTAssertEqual( + [ + Array("He".utf8), + Array("ll".utf8), + ], + subscriber.values + ) + XCTAssertTrue(subscriber.completions.isEmpty) + + subscriber.subscription?.request(.max(1)) + XCTAssertEqual( + [ + Array("He".utf8), + Array("ll".utf8), + Array("o".utf8), + ], + subscriber.values + ) + XCTAssertTrue(subscriber.completions.isEmpty) + + subscriber.subscription?.request(.max(1)) + XCTAssertEqual(1, subscriber.completions.count) + + guard case .finished? = subscriber.completions.first + else { return XCTFail("Expected finished completion") } + } + + func testInputStreamPublisherUsesAdditionalDemandReturnedBySubscriber() throws { + let data = Data("Hello".utf8) + let subscriber = InputStreamManualSubscriber( + initialDemand: .max(1), + demandOnValue: .max(1) + ) + + data + .publisher(maxChunkSize: 2) + .receive(subscriber: subscriber) + + XCTAssertEqual( + [ + Array("He".utf8), + Array("ll".utf8), + Array("o".utf8), + ], + subscriber.values + ) + XCTAssertEqual(1, subscriber.completions.count) + + guard case .finished? = subscriber.completions.first + else { return XCTFail("Expected finished completion") } + } + func testInputStreamPublisherWithLessData() throws { let data = Data("Hello".utf8) @@ -61,6 +150,30 @@ class InputStreamPublisherTests: XCTestCase { wait(for: [ex], timeout: 2) } + func testInputStreamURLPublisherWithEmptyFileFinishesWithoutOutput() throws { + let tempDir = FileManager.default + .temporaryDirectory + .appendingPathComponent( + "testInputStreamURLPublisherWithEmptyFileFinishesWithoutOutput" + ) + + try? FileManager.default.removeItem(at: tempDir) + try FileManager.default.createDirectory(at: tempDir, withIntermediateDirectories: true) + + let url = tempDir.appendingPathComponent("empty.txt") + try Data().write(to: url, options: .atomic) + + defer { + try? FileManager.default.removeItem(at: tempDir) + } + + let ex = url + .publisher(maxChunkSize: 3) + .expectToFinish(failsOnOutput: true) + + wait(for: [ex], timeout: 2) + } + func testInputStreamPublisherFailsForInvalidURL() throws { let url = FileManager.default .temporaryDirectory @@ -141,3 +254,37 @@ class InputStreamPublisherTests: XCTestCase { wait(for: [completionEx], timeout: 0.1) } } + +private final class InputStreamManualSubscriber: Subscriber { + typealias Input = [UInt8] + typealias Failure = Error + + var subscription: Subscription? + private(set) var values = [[UInt8]]() + private(set) var completions = [Subscribers.Completion]() + + private let initialDemand: Subscribers.Demand + private let demandOnValue: Subscribers.Demand + + init( + initialDemand: Subscribers.Demand, + demandOnValue: Subscribers.Demand = .none + ) { + self.initialDemand = initialDemand + self.demandOnValue = demandOnValue + } + + func receive(subscription: Subscription) { + self.subscription = subscription + subscription.request(initialDemand) + } + + func receive(_ input: [UInt8]) -> Subscribers.Demand { + values.append(input) + return demandOnValue + } + + func receive(completion: Subscribers.Completion) { + completions.append(completion) + } +} diff --git a/Tests/CombineExtensionsTests/KeyedSubscriptionStoreTests.swift b/Tests/CombineExtensionsTests/KeyedSubscriptionStoreTests.swift index 1a8899a..6baa8fd 100644 --- a/Tests/CombineExtensionsTests/KeyedSubscriptionStoreTests.swift +++ b/Tests/CombineExtensionsTests/KeyedSubscriptionStoreTests.swift @@ -30,6 +30,70 @@ class KeyedSubscriptionStoreTests: XCTestCase { XCTAssertFalse(store.containsSubscription(forKey: "third")) } + func testKeyedSubscriptionStoreInitializerStoresSubscriptions() throws { + var cancelCount = 0 + + let store = KeyedSubscriptionStore( + subscriptions: ["initial": AnyCancellable { cancelCount += 1 }] + ) + + XCTAssertFalse(store.isEmpty) + XCTAssertTrue(store.containsSubscription(forKey: "initial")) + + store.removeAll() + + XCTAssertEqual(1, cancelCount) + } + + func testStoringSubscriptionForExistingKeyCancelsPreviousSubscription() throws { + let store = KeyedSubscriptionStore() + + var firstCancelCount = 0 + var secondCancelCount = 0 + + store.store( + subscription: AnyCancellable { firstCancelCount += 1 }, + forKey: "shared" + ) + + store.store( + subscription: AnyCancellable { secondCancelCount += 1 }, + forKey: "shared" + ) + + XCTAssertEqual(1, firstCancelCount) + XCTAssertEqual(0, secondCancelCount) + + store.removeAll() + + XCTAssertEqual(1, firstCancelCount) + XCTAssertEqual(1, secondCancelCount) + } + + func testRemovedSubscriptionStaysActiveWhileReturnedValueIsRetained() throws { + let store = KeyedSubscriptionStore() + let subject = PassthroughSubject() + + var receivedValues = [Int]() + subject + .sink { receivedValues.append($0) } + .store(in: store, key: "subject") + + var removedSubscription = store.removeSubscription(forKey: "subject") + + XCTAssertNotNil(removedSubscription) + XCTAssertTrue(store.isEmpty) + + subject.send(1) + + XCTAssertEqual([1], receivedValues) + + removedSubscription = nil + subject.send(2) + + XCTAssertEqual([1], receivedValues) + } + func testKeyedSubscriptionStoreEquatableAndHashable() throws { let one = KeyedSubscriptionStore() let two = KeyedSubscriptionStore() diff --git a/Tests/CombineExtensionsTests/MulticastLatestTests.swift b/Tests/CombineExtensionsTests/MulticastLatestTests.swift index 27c7f23..f748b17 100644 --- a/Tests/CombineExtensionsTests/MulticastLatestTests.swift +++ b/Tests/CombineExtensionsTests/MulticastLatestTests.swift @@ -4,11 +4,92 @@ import CombineTestExtensions import XCTest class MulticastLatestSubjectTests: XCTestCase { + func testMulticastLatestDoesNotEmitUntilUpstreamPublishes() { + let subject = PassthroughSubject() + + var values = [Int]() + let subscription = subject + .multicastLatest() + .sink { values.append($0) } + + XCTAssertEqual([], values) + + subject.send(1) + + XCTAssertEqual([1], values) + + subscription.cancel() + } + + func testLateSubscriberReceivesLatestValueImmediately() { + let subject = PassthroughSubject() + let publisher = subject.multicastLatest() + + var firstValues = [Int]() + let firstSubscription = publisher.sink { firstValues.append($0) } + + subject.send(1) + subject.send(2) + + var secondValues = [Int]() + let secondSubscription = publisher.sink { secondValues.append($0) } + + XCTAssertEqual([1, 2], firstValues) + XCTAssertEqual([2], secondValues) + + firstSubscription.cancel() + secondSubscription.cancel() + } + + func testMulticastLatestForwardsFinishedCompletion() { + let subject = PassthroughSubject() + + var values = [Int]() + var completion: Subscribers.Completion? + let subscription = subject + .multicastLatest() + .sink( + receiveCompletion: { completion = $0 }, + receiveValue: { values.append($0) } + ) + + subject.send(1) + subject.send(completion: .finished) + subject.send(2) + + XCTAssertEqual([1], values) + XCTAssertEqual(.finished, completion) + + subscription.cancel() + } + + func testMulticastLatestForwardsFailure() { + let subject = PassthroughSubject() + + var values = [Int]() + var completion: Subscribers.Completion? + let subscription = subject + .multicastLatest() + .sink( + receiveCompletion: { completion = $0 }, + receiveValue: { values.append($0) } + ) + + subject.send(1) + subject.send(completion: .failure(.failed)) + subject.send(2) + + XCTAssertEqual([1], values) + XCTAssertEqual(.failure(.failed), completion) + + subscription.cancel() + } + func testMulticastLatest() { let voidSubject = PassthroughSubject() let intSubject = PassthroughSubject() var mapCalledCount = 0 - + let multicastedPublisher = voidSubject .map { mapCalledCount += 1 @@ -16,28 +97,32 @@ class MulticastLatestSubjectTests: XCTestCase { } .switchToLatest() .multicastLatest() - + let outputExpectation = multicastedPublisher.expectOutput([1, 2]) voidSubject.send(()) intSubject.send(1) intSubject.send(2) - + let expectation1 = expectation(description: "first subscriber") let subscription1 = multicastedPublisher.sink { value in XCTAssertEqual(value, 2) expectation1.fulfill() } defer { subscription1.cancel() } - + let expectation2 = expectation(description: "second subscriber") let subscription2 = multicastedPublisher.sink { value in XCTAssertEqual(value, 2) expectation2.fulfill() } defer { subscription2.cancel() } - + wait(for: [outputExpectation, expectation1, expectation2], timeout: 0.5) XCTAssertEqual(mapCalledCount, 1) } } + +private enum MulticastError: Error, Equatable { + case failed +} diff --git a/Tests/CombineExtensionsTests/OnCancelTests.swift b/Tests/CombineExtensionsTests/OnCancelTests.swift index 1cea34b..7eb903c 100644 --- a/Tests/CombineExtensionsTests/OnCancelTests.swift +++ b/Tests/CombineExtensionsTests/OnCancelTests.swift @@ -3,6 +3,48 @@ import CombineExtensions import XCTest class OnCancelTests: XCTestCase { + func testOnCancelCancelsWrappedCancellableBeforeRunningBlock() throws { + var events = [String]() + let wrapped = SpyCancellable { events.append("cancel") } + + let cancellable = wrapped.onCancel { + XCTAssertTrue(wrapped.isCancelled) + events.append("block") + } + + cancellable.cancel() + + XCTAssertEqual(["cancel", "block"], events) + } + + func testOnCancelBlockRunsOnlyOnceWhenCancelledMultipleTimes() throws { + var cancelCount = 0 + var blockCount = 0 + + let wrapped = SpyCancellable { cancelCount += 1 } + let cancellable = wrapped.onCancel { blockCount += 1 } + + cancellable.cancel() + cancellable.cancel() + + XCTAssertEqual(1, cancelCount) + XCTAssertEqual(1, blockCount) + } + + func testOnCancelRunsWhenReturnedCancellableIsDeallocated() throws { + var cancelCount = 0 + var blockCount = 0 + + let wrapped = SpyCancellable { cancelCount += 1 } + var cancellable: AnyCancellable? = wrapped.onCancel { blockCount += 1 } + + XCTAssertNotNil(cancellable) + cancellable = nil + + XCTAssertEqual(1, cancelCount) + XCTAssertEqual(1, blockCount) + } + func testOnCancelIsCalledWhenCancelled() throws { let subject = PassthroughSubject() @@ -53,3 +95,17 @@ class OnCancelTests: XCTestCase { } private struct _Err: Error, Equatable {} + +private final class SpyCancellable: Cancellable { + private let onCancel: () -> Void + private(set) var isCancelled = false + + init(onCancel: @escaping () -> Void) { + self.onCancel = onCancel + } + + func cancel() { + isCancelled = true + onCancel() + } +} diff --git a/Tests/CombineExtensionsTests/OutputStreamPublisherTests.swift b/Tests/CombineExtensionsTests/OutputStreamPublisherTests.swift index 565ea36..5dfa11e 100644 --- a/Tests/CombineExtensionsTests/OutputStreamPublisherTests.swift +++ b/Tests/CombineExtensionsTests/OutputStreamPublisherTests.swift @@ -97,6 +97,93 @@ class OutputStreamPublisherTests: XCTestCase { XCTAssertEqual([72, 105, 33, 72, 105, 33, 0, 0, 0, 0] as [UInt8], bytes) } + func testLargeChunkFailsWithoutWritingOrPublishingByteCount() throws { + let subject = PassthroughSubject<[UInt8], Error>() + let input = (0 ..< 12).map(UInt8.init) + + let ex = subject + .stream(toBuffer: buffer, capacity: bufferCapacity) + .expectFailure( + { error in + let nserror = error as NSError + XCTAssertEqual(NSPOSIXErrorDomain, nserror.domain) + XCTAssertEqual(12, nserror.code) + }, + failsOnOutput: true + ) + + subject.send(input) + + wait(for: [ex], timeout: 2) + + XCTAssertEqual([0, 0, 0, 0, 0, 0, 0, 0, 0, 0] as [UInt8], bytes) + } + + func testDownstreamDemandIsForwardedToUpstream() throws { + let subject = PassthroughSubject<[UInt8], Error>() + let subscriber = OutputStreamManualSubscriber(initialDemand: .max(2)) + var requests = [Subscribers.Demand]() + + subject + .handleEvents(receiveRequest: { requests.append($0) }) + .stream(toBuffer: buffer, capacity: bufferCapacity) + .receive(subscriber: subscriber) + + XCTAssertEqual([.max(2)], requests) + + subject.send([1]) + subject.send([2]) + subject.send([3]) + + XCTAssertEqual([1, 1], subscriber.values) + XCTAssertEqual([1, 2, 0, 0, 0, 0, 0, 0, 0, 0] as [UInt8], bytes) + } + + func testNoUpstreamDemandIsRequestedUntilDownstreamRequestsDemand() throws { + let subject = PassthroughSubject<[UInt8], Error>() + let subscriber = OutputStreamManualSubscriber(initialDemand: .none) + var requests = [Subscribers.Demand]() + + subject + .handleEvents(receiveRequest: { requests.append($0) }) + .stream(toBuffer: buffer, capacity: bufferCapacity) + .receive(subscriber: subscriber) + + subject.send([1]) + + XCTAssertEqual([], requests) + XCTAssertEqual([], subscriber.values) + XCTAssertEqual([0, 0, 0, 0, 0, 0, 0, 0, 0, 0] as [UInt8], bytes) + + subscriber.subscription?.request(.max(1)) + subject.send([2]) + + XCTAssertEqual([.max(1)], requests) + XCTAssertEqual([1], subscriber.values) + XCTAssertEqual([2, 0, 0, 0, 0, 0, 0, 0, 0, 0] as [UInt8], bytes) + } + + func testUpstreamFailureIsForwardedAfterWritingAvailableOutput() throws { + let subject = PassthroughSubject<[UInt8], Error>() + + let ex = subject + .stream(toBuffer: buffer, capacity: bufferCapacity) + .expectOutputAndFailure { value in + XCTAssertEqual(1, value) + return .finished + } failureEvaluator: { error in + XCTAssertEqual(.failed, error as? OutputStreamTestError) + } + + subject.send([42]) + subject.send(completion: .failure(OutputStreamTestError.failed)) + subject.send([43]) + + wait(for: [ex], timeout: 2) + + XCTAssertEqual([42, 0, 0, 0, 0, 0, 0, 0, 0, 0] as [UInt8], bytes) + } + func testWithTooMuchData() throws { let subject = PassthroughSubject<[UInt8], Error>() let input: [UInt8] = [ @@ -113,17 +200,14 @@ class OutputStreamPublisherTests: XCTestCase { toBuffer: buffer, capacity: bufferCapacity ) - .expectOutputAndFailure( - { value in - expectedOutputCount -= value - return expectedOutputCount == 0 ? .finished : .moreExpected - }, - failureEvaluator: { error in - let nserror = error as NSError - XCTAssertEqual(NSPOSIXErrorDomain, nserror.domain) - XCTAssertEqual(12, nserror.code) - } - ) + .expectOutputAndFailure { value in + expectedOutputCount -= value + return expectedOutputCount == 0 ? .finished : .moreExpected + } failureEvaluator: { error in + let nserror = error as NSError + XCTAssertEqual(NSPOSIXErrorDomain, nserror.domain) + XCTAssertEqual(12, nserror.code) + } input.forEach { subject.send([$0]) } @@ -189,6 +273,35 @@ class OutputStreamPublisherTests: XCTestCase { XCTAssertEqual("Hello!", try String(contentsOf: url)) } + func testStreamToURLAppendAddsToExistingFile() throws { + let subject = PassthroughSubject<[UInt8], Error>() + + let tempDir = FileManager.default + .temporaryDirectory + .appendingPathComponent("testStreamToURLAppendAddsToExistingFile") + + try? FileManager.default.removeItem(at: tempDir) + try FileManager.default.createDirectory(at: tempDir, withIntermediateDirectories: true) + defer { try? FileManager.default.removeItem(at: tempDir) } + + let url = tempDir.appendingPathComponent("text.txt") + try Data("Hello".utf8).write(to: url, options: .atomic) + + let ex = subject + .stream( + toURL: url, + append: true + ) + .expectToFinish() + + subject.send(Array("!".utf8)) + subject.send(completion: .finished) + + wait(for: [ex], timeout: 2) + + XCTAssertEqual("Hello!", try String(contentsOf: url)) + } + func testFailsForInvalidURL() throws { let subject = PassthroughSubject<[UInt8], Error>() @@ -228,7 +341,7 @@ class OutputStreamPublisherTests: XCTestCase { completionEx.isInverted = true let demandEx = expectation(description: "Should have received demand") - let _ = subject + _ = subject .handleEvents( receiveSubscription: { _ in subscriptionEx.fulfill() }, receiveRequest: { demand in @@ -296,3 +409,43 @@ class OutputStreamPublisherTests: XCTestCase { XCTAssertEqual([42, 0, 0, 0, 0, 0, 0, 0, 0, 0] as [UInt8], bytes) } } + +private final class OutputStreamManualSubscriber: Subscriber { + typealias Input = Int + typealias Failure = Error + + var subscription: Subscription? + private(set) var values = [Int]() + private(set) var completions = [Subscribers.Completion]() + + private let initialDemand: Subscribers.Demand + private let demandOnValue: Subscribers.Demand + + init( + initialDemand: Subscribers.Demand, + demandOnValue: Subscribers.Demand = .none + ) { + self.initialDemand = initialDemand + self.demandOnValue = demandOnValue + } + + func receive(subscription: Subscription) { + self.subscription = subscription + if initialDemand > .none { + subscription.request(initialDemand) + } + } + + func receive(_ input: Int) -> Subscribers.Demand { + values.append(input) + return demandOnValue + } + + func receive(completion: Subscribers.Completion) { + completions.append(completion) + } +} + +private enum OutputStreamTestError: Error, Equatable { + case failed +} diff --git a/Tests/CombineExtensionsTests/PublisherSinkTests.swift b/Tests/CombineExtensionsTests/PublisherSinkTests.swift new file mode 100644 index 0000000..a5bcbe4 --- /dev/null +++ b/Tests/CombineExtensionsTests/PublisherSinkTests.swift @@ -0,0 +1,47 @@ +import Combine +import CombineExtensions +import XCTest + +final class PublisherSinkTests: XCTestCase { + func testSinkReceiveValueReceiveCompletionOverloadReceivesValuesThenCompletion() throws { + let subject = PassthroughSubject() + + var values = [Int]() + var completion: Subscribers.Completion? + + let subscription = subject.sink( + receiveValue: { values.append($0) }, + receiveCompletion: { completion = $0 } + ) + + subject.send(1) + subject.send(2) + subject.send(completion: .failure(.failed)) + subject.send(3) + + XCTAssertEqual([1, 2], values) + XCTAssertEqual(.failure(.failed), completion) + + subscription.cancel() + } + + func testVoidSinkCompletionOverloadIgnoresValuesAndReceivesCompletion() throws { + let subject = PassthroughSubject() + + var completion: Subscribers.Completion? + + let subscription = subject.sink { completion = $0 } + + subject.send(()) + subject.send(()) + subject.send(completion: .finished) + + XCTAssertEqual(.finished, completion) + + subscription.cancel() + } +} + +private enum SinkError: Error, Equatable { + case failed +} diff --git a/Tests/CombineExtensionsTests/PublisherTestExtensionsTests.swift b/Tests/CombineExtensionsTests/PublisherTestExtensionsTests.swift index ba9b257..46a87d0 100644 --- a/Tests/CombineExtensionsTests/PublisherTestExtensionsTests.swift +++ b/Tests/CombineExtensionsTests/PublisherTestExtensionsTests.swift @@ -3,6 +3,50 @@ import CombineTestExtensions import XCTest class PublisherTestExtensionsTests: XCTestCase { + func testExpectOutputSingleValueOverload() throws { + let ex = Just(1).expectOutput(1) + + wait(for: [ex], timeout: 0.5) + } + + func testExpectOutputCount() throws { + let ex = [1, 2, 3] + .publisher + .expectOutput(count: 3) + + wait(for: [ex], timeout: 0.5) + } + + func testExpectToFinish() throws { + let ex = Empty(completeImmediately: true) + .expectToFinish(failsOnOutput: true) + + wait(for: [ex], timeout: 0.5) + } + + func testExpectAnyFailure() throws { + let ex = Fail(error: .wrong) + .expectAnyFailure(failsOnOutput: true) + + wait(for: [ex], timeout: 0.5) + } + + func testExpectOutputWithCustomOutputAndFailureComparators() throws { + let subject = PassthroughSubject() + + let ex = subject.expectOutput( + ["hello"], + outputComparator: { $0.caseInsensitiveCompare($1) == .orderedSame }, + completion: .failure(.correct), + failureComparator: == + ) + + subject.send("HELLO") + subject.send(completion: .failure(.correct)) + + wait(for: [ex], timeout: 0.5) + } + func testExpectOutputWithEvaluator() throws { var ints = [1, 2, 3] let evaluator = { (input: Int) -> OutputExpectation in @@ -21,13 +65,12 @@ class PublisherTestExtensionsTests: XCTestCase { let ex = subject .receive(on: DispatchQueue(label: "test")) - .expectOutputAndFailure( - { output -> OutputExpectation in - XCTAssertEqual(ints.removeFirst(), output) - return ints.isEmpty ? .finished : .moreExpected - }, - failureEvaluator: { XCTAssertEqual(.correct, $0) } - ) + .expectOutputAndFailure { output -> OutputExpectation in + XCTAssertEqual(ints.removeFirst(), output) + return ints.isEmpty ? .finished : .moreExpected + } failureEvaluator: { + XCTAssertEqual(.correct, $0) + } subject.send(1) subject.send(2) diff --git a/Tests/CombineExtensionsTests/ReduceLatestTests.swift b/Tests/CombineExtensionsTests/ReduceLatestTests.swift index b28f0ab..c0d5c60 100644 --- a/Tests/CombineExtensionsTests/ReduceLatestTests.swift +++ b/Tests/CombineExtensionsTests/ReduceLatestTests.swift @@ -79,6 +79,120 @@ class ReduceLatestTests: XCTestCase { wait(for: [ex], timeout: 2) } + func testReduceLatestDoesNotPublishUntilDemandIsRequested() throws { + let pub1 = PassthroughSubject() + let pub2 = PassthroughSubject() + let subscriber = ReduceLatestManualSubscriber() + + pub1 + .reduceLatest(pub2, Result.init) + .receive(subscriber: subscriber) + + pub1.send(1) + pub1.send(2) + pub2.send("A") + + XCTAssertEqual([], subscriber.values) + + subscriber.subscription?.request(.max(1)) + + XCTAssertEqual([Result(2, "A")], subscriber.values) + } + + func testReduceLatestUsesAdditionalDemandReturnedBySubscriber() throws { + let pub1 = PassthroughSubject() + let pub2 = PassthroughSubject() + let subscriber = ReduceLatestManualSubscriber( + demandOnValue: .max(1) + ) + + pub1 + .reduceLatest(pub2, Result.init) + .receive(subscriber: subscriber) + + subscriber.subscription?.request(.max(1)) + pub1.send(1) + pub2.send("A") + + XCTAssertEqual( + [ + Result(), + Result(1), + Result(1, "A"), + ], + subscriber.values + ) + } + + func testReduceLatestCompletesOnlyAfterBothUpstreamsFinish() throws { + let pub1 = PassthroughSubject() + let pub2 = PassthroughSubject() + let subscriber = ReduceLatestManualSubscriber() + + pub1 + .reduceLatest(pub2, Result.init) + .receive(subscriber: subscriber) + + subscriber.subscription?.request(.unlimited) + + pub1.send(1) + pub1.send(completion: .finished) + + XCTAssertEqual([], subscriber.completions) + + pub2.send("A") + pub2.send(completion: .finished) + + XCTAssertEqual( + [ + Result(), + Result(1), + Result(1, "A"), + ], + subscriber.values + ) + XCTAssertEqual(1, subscriber.completions.count) + + guard case .finished? = subscriber.completions.first + else { return XCTFail("Expected finished completion") } + } + + func testReduceLatestForwardsFailureEvenWithoutDemand() throws { + let pub1 = PassthroughSubject() + let pub2 = PassthroughSubject() + let subscriber = ReduceLatestManualSubscriber() + + pub1 + .reduceLatest(pub2, Result.init) + .receive(subscriber: subscriber) + + pub1.send(completion: .failure(.anError)) + + XCTAssertEqual([], subscriber.values) + XCTAssertEqual([.failure(.anError)], subscriber.completions) + } + + func testReduceLatestCancelStopsOutputAndCompletion() throws { + let pub1 = PassthroughSubject() + let pub2 = PassthroughSubject() + let subscriber = ReduceLatestManualSubscriber() + + pub1 + .reduceLatest(pub2, Result.init) + .receive(subscriber: subscriber) + + subscriber.subscription?.request(.unlimited) + subscriber.subscription?.cancel() + + pub1.send(1) + pub2.send("A") + pub1.send(completion: .finished) + pub2.send(completion: .finished) + + XCTAssertEqual([Result()], subscriber.values) + XCTAssertEqual([], subscriber.completions) + } + func testReduceLatest3() throws { let pub1 = PassthroughSubject() let pub2 = PassthroughSubject() @@ -145,6 +259,31 @@ private enum ReducerError: Error, Equatable { case anError } +private final class ReduceLatestManualSubscriber: Subscriber { + var subscription: Subscription? + private(set) var values = [Input]() + private(set) var completions = [Subscribers.Completion]() + + private let demandOnValue: Subscribers.Demand + + init(demandOnValue: Subscribers.Demand = .none) { + self.demandOnValue = demandOnValue + } + + func receive(subscription: Subscription) { + self.subscription = subscription + } + + func receive(_ input: Input) -> Subscribers.Demand { + values.append(input) + return demandOnValue + } + + func receive(completion: Subscribers.Completion) { + completions.append(completion) + } +} + private struct Result: CustomStringConvertible, Equatable { let number: Int? let letter: Character? diff --git a/Tests/CombineExtensionsTests/RetryIfTests.swift b/Tests/CombineExtensionsTests/RetryIfTests.swift index 480d2d7..be73e9e 100644 --- a/Tests/CombineExtensionsTests/RetryIfTests.swift +++ b/Tests/CombineExtensionsTests/RetryIfTests.swift @@ -4,6 +4,138 @@ import CombineTestExtensions import XCTest final class RetryIfTests: XCTestCase { + func testRetryAfterConvenienceRetriesEveryFailure() throws { + let scheduler = DispatchQueue.test + + let ex = [1] + .publisher + .failOnOutputIndex([0], error: .two) + .retry( + after: { _ in .seconds(1) }, + scheduler: scheduler + ) + .expectOutput(1, expectToFinish: true) + + scheduler.advance(by: .seconds(1)) + + wait(for: [ex], timeout: 1) + } + + func testFinishedCompletionDoesNotRetry() throws { + let scheduler = DispatchQueue.test + var attempts = 0 + + let publisher = Deferred { + attempts += 1 + return Empty(completeImmediately: true) + } + + let ex = publisher + .retryIf( + { _ in true }, + after: .seconds(1), + scheduler: scheduler + ) + .expectToFinish(failsOnOutput: true) + + scheduler.run() + + wait(for: [ex], timeout: 1) + XCTAssertEqual(1, attempts) + } + + func testRejectedFailureAfterRetryCompletesWithRejectedFailure() throws { + let scheduler = DispatchQueue.test + var attempts = 0 + + let publisher = Deferred { + attempts += 1 + return Fail(error: attempts == 1 ? .one : .two) + } + + let ex = publisher + .retryIf( + { $0 == .one }, + after: .seconds(1), + scheduler: scheduler + ) + .expectFailure(.two, failsOnOutput: true) + + scheduler.advance(by: .seconds(1)) + + wait(for: [ex], timeout: 1) + XCTAssertEqual(2, attempts) + } + + func testPredicateAndIntervalReceiveFailuresAndRetryCounts() throws { + let scheduler = DispatchQueue.test + + var attempts = 0 + var predicateErrors = [Err]() + var intervalCounts = [Int]() + + let publisher = Deferred { + attempts += 1 + if attempts < 3 { + return Fail(error: .one).eraseToAnyPublisher() + } else { + return Just(42).setFailureType(to: Err.self).eraseToAnyPublisher() + } + } + + let ex = publisher + .retryIf( + { + predicateErrors.append($0) + return true + }, + after: { + intervalCounts.append($0) + return .seconds($0) + }, + scheduler: scheduler + ) + .expectOutput(42, expectToFinish: true) + + scheduler.advance(by: .seconds(1)) + scheduler.advance(by: .seconds(2)) + + wait(for: [ex], timeout: 1) + XCTAssertEqual([.one, .one], predicateErrors) + XCTAssertEqual([1, 2], intervalCounts) + } + + func testScheduledRetryAfterCancellationDoesNotPublishOutputOrCompletion() throws { + let scheduler = DispatchQueue.test + var attempts = 0 + var values = [Int]() + var completion: Subscribers.Completion? + + let publisher = Deferred { + attempts += 1 + return Fail(error: .one) + } + + let subscription = publisher + .retryIf( + { $0 == .one }, + after: .seconds(1), + scheduler: scheduler + ) + .sink( + receiveCompletion: { completion = $0 }, + receiveValue: { values.append($0) } + ) + + XCTAssertEqual(1, attempts) + + subscription.cancel() + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual([], values) + XCTAssertNil(completion) + } + func testRetryIfAfterFixedInterval() throws { let scheduler = DispatchQueue.test diff --git a/Tests/CombineExtensionsTests/SingleSubscriptionStoreTests.swift b/Tests/CombineExtensionsTests/SingleSubscriptionStoreTests.swift index d95ebc3..e8e5ba1 100644 --- a/Tests/CombineExtensionsTests/SingleSubscriptionStoreTests.swift +++ b/Tests/CombineExtensionsTests/SingleSubscriptionStoreTests.swift @@ -17,6 +17,60 @@ class SingleSubscriptionStoreTests: XCTestCase { XCTAssertTrue(store.isEmpty) } + func testSingleSubscriptionStoreInitializerStoresSubscription() throws { + var cancelCount = 0 + let store = SingleSubscriptionStore(AnyCancellable { cancelCount += 1 }) + + XCTAssertFalse(store.isEmpty) + + store.removeSubscription() + + XCTAssertEqual(1, cancelCount) + XCTAssertTrue(store.isEmpty) + } + + func testStoringNewSubscriptionCancelsPreviousSubscription() throws { + let store = SingleSubscriptionStore() + + var firstCancelCount = 0 + var secondCancelCount = 0 + + store.store(subscription: AnyCancellable { firstCancelCount += 1 }) + store.store(subscription: AnyCancellable { secondCancelCount += 1 }) + + XCTAssertEqual(1, firstCancelCount) + XCTAssertEqual(0, secondCancelCount) + + store.removeSubscription() + + XCTAssertEqual(1, firstCancelCount) + XCTAssertEqual(1, secondCancelCount) + } + + func testRemovedSubscriptionStaysActiveWhileReturnedValueIsRetained() throws { + let store = SingleSubscriptionStore() + let subject = PassthroughSubject() + + var receivedValues = [Int]() + subject + .sink { receivedValues.append($0) } + .store(in: store) + + var removedSubscription = store.removeSubscription() + + XCTAssertNotNil(removedSubscription) + XCTAssertTrue(store.isEmpty) + + subject.send(1) + + XCTAssertEqual([1], receivedValues) + + removedSubscription = nil + subject.send(2) + + XCTAssertEqual([1], receivedValues) + } + func testSingleSubscriptionStoreEquatableAndHashable() throws { let one = SingleSubscriptionStore() let two = SingleSubscriptionStore() diff --git a/Tests/CombineExtensionsTests/TestSchedulerTests.swift b/Tests/CombineExtensionsTests/TestSchedulerTests.swift index 458963b..4489369 100644 --- a/Tests/CombineExtensionsTests/TestSchedulerTests.swift +++ b/Tests/CombineExtensionsTests/TestSchedulerTests.swift @@ -48,6 +48,79 @@ final class TestSchedulerTests: XCTestCase { XCTAssertEqual(value, 1) } + func testRunWithNoScheduledActionsDoesNotChangeNow() { + let scheduler = DispatchQueue.test + let now = scheduler.now + + scheduler.run() + + XCTAssertEqual(now, scheduler.now) + } + + func testAdvanceWithNoScheduledActionsMovesNowToFinalDate() { + let scheduler = DispatchQueue.test + let now = scheduler.now + + scheduler.advance(by: .seconds(3)) + + XCTAssertEqual(now.advanced(by: .seconds(3)), scheduler.now) + } + + func testOneOffActionsAtSameDateRunInSchedulingOrder() { + let scheduler = DispatchQueue.test + var values = [Int]() + + scheduler.schedule(after: scheduler.now.advanced(by: .seconds(1))) { + values.append(1) + } + + scheduler.schedule(after: scheduler.now.advanced(by: .seconds(1))) { + values.append(2) + } + + scheduler.schedule(after: scheduler.now.advanced(by: .seconds(1))) { + values.append(3) + } + + scheduler.advance(by: .seconds(1)) + + XCTAssertEqual([1, 2, 3], values) + } + + func testActionsScheduledForNowDuringAdvanceRunDuringSameAdvance() { + let scheduler = DispatchQueue.test + var values = [Int]() + + scheduler.schedule { + values.append(1) + scheduler.schedule { values.append(3) } + } + + scheduler.schedule { values.append(2) } + + scheduler.advance() + + XCTAssertEqual([1, 2, 3], values) + } + + func testCancellingIntervalRemovesFutureScheduledActions() { + let scheduler = DispatchQueue.test + var values = [Int]() + + let cancellable = scheduler.schedule( + after: scheduler.now, + interval: .seconds(1) + ) { + values.append(1) + } + + scheduler.advance() + cancellable.cancel() + scheduler.advance(by: .seconds(5)) + + XCTAssertEqual([1], values) + } + func testDelay0Advance() { let scheduler = DispatchQueue.test diff --git a/Tests/CombineExtensionsTests/ThrottleWhileTests.swift b/Tests/CombineExtensionsTests/ThrottleWhileTests.swift index dfef4a6..7cc79c8 100644 --- a/Tests/CombineExtensionsTests/ThrottleWhileTests.swift +++ b/Tests/CombineExtensionsTests/ThrottleWhileTests.swift @@ -82,6 +82,120 @@ final class ThrottleWhileTests: XCTestCase { wait(for: [ex], timeout: 2) } + func testLatestDropsBufferedValueWhenUpstreamCompletesWhileThrottled() throws { + let subject = PassthroughSubject() + let regulator = CurrentValueSubject(true) + + var values = [Int]() + var completion: Subscribers.Completion? + let subscription = subject + .throttle(while: regulator, latest: true) + .sink( + receiveCompletion: { completion = $0 }, + receiveValue: { values.append($0) } + ) + + subject.send(1) + subject.send(2) + subject.send(completion: .finished) + regulator.send(false) + + XCTAssertEqual([], values) + XCTAssertEqual(.finished, completion) + + subscription.cancel() + } + + func testEarliestDropsBufferedValueWhenUpstreamCompletesWhileThrottled() throws { + let subject = PassthroughSubject() + let regulator = CurrentValueSubject(true) + + var values = [Int]() + var completion: Subscribers.Completion? + let subscription = subject + .throttle(while: regulator, latest: false) + .sink( + receiveCompletion: { completion = $0 }, + receiveValue: { values.append($0) } + ) + + subject.send(1) + subject.send(2) + subject.send(completion: .finished) + regulator.send(false) + + XCTAssertEqual([], values) + XCTAssertEqual(.finished, completion) + + subscription.cancel() + } + + func testUpstreamFailureWhileThrottledDropsBufferedValueAndFails() throws { + let subject = PassthroughSubject() + let regulator = CurrentValueSubject(true) + + var values = [Int]() + var completion: Subscribers.Completion? + let subscription = subject + .throttle(while: regulator, latest: true) + .sink( + receiveCompletion: { completion = $0 }, + receiveValue: { values.append($0) } + ) + + subject.send(1) + subject.send(completion: .failure(.failed)) + regulator.send(false) + + XCTAssertEqual([], values) + XCTAssertEqual(.failure(.failed), completion) + + subscription.cancel() + } + + func testRegulatorFailureWhileThrottledDropsBufferedValueAndFails() throws { + let subject = PassthroughSubject() + let regulator = CurrentValueSubject(true) + + var values = [Int]() + var completion: Subscribers.Completion? + let subscription = subject + .throttle(while: regulator, latest: true) + .sink( + receiveCompletion: { completion = $0 }, + receiveValue: { values.append($0) } + ) + + subject.send(1) + regulator.send(completion: .failure(.failed)) + subject.send(2) + + XCTAssertEqual([], values) + XCTAssertEqual(.failure(.failed), completion) + + subscription.cancel() + } + + func testRegulatorFalseBeforeUpstreamDoesNotPublishAnything() throws { + let subject = PassthroughSubject() + let regulator = PassthroughSubject() + + var values = [Int]() + let subscription = subject + .throttle(while: regulator, latest: true) + .sink { values.append($0) } + + regulator.send(false) + + XCTAssertEqual([], values) + + subject.send(1) + + XCTAssertEqual([1], values) + + subscription.cancel() + } + func testLatestHandlesEngagingAndReleasingThrottle() throws { let subject = PassthroughSubject() let regulator = PassthroughSubject() @@ -228,7 +342,7 @@ final class ThrottleWhileTests: XCTestCase { XCTAssertEqual(values, [1, 2]) } - + func testEarliestPublishesWhenRegulatorFiresFromSink() throws { let subject = PassthroughSubject() let regulator = PassthroughSubject() @@ -250,14 +364,14 @@ final class ThrottleWhileTests: XCTestCase { XCTAssertEqual(values, [1, 2]) } - + func testLatestFlippingRegulatorDoesNotResendSameEmission() throws { func pause() { RunLoop.main.run(until: Date(timeIntervalSinceNow: 0.01)) } - + let subject = PassthroughSubject() let regulator = PassthroughSubject() var values = [String]() - + let firstEx = expectation(description: "received first emission") let compEx = expectation(description: "stream is completed") let sub = subject @@ -271,41 +385,41 @@ final class ThrottleWhileTests: XCTestCase { receiveCompletion: { _ in compEx.fulfill() } ) defer { sub.cancel() } - + subject.send("one") wait(for: [firstEx], timeout: 2) XCTAssertEqual(["one"], values) - + regulator.send(true) regulator.send(false) regulator.send(true) regulator.send(false) - + pause() - + XCTAssertEqual(["one"], values) - + subject.send("one") - + regulator.send(true) subject.send("two") subject.send("three") subject.send("four") regulator.send(false) - + subject.send(completion: .finished) wait(for: [compEx], timeout: 2) - + XCTAssertEqual(["one", "one", "four"], values) } - + func testEarliestFlippingRegulatorDoesNotResendSameEmission() throws { func pause() { RunLoop.main.run(until: Date(timeIntervalSinceNow: 0.01)) } - + let subject = PassthroughSubject() let regulator = PassthroughSubject() var values = [String]() - + let firstEx = expectation(description: "received first emission") let compEx = expectation(description: "stream is completed") let sub = subject @@ -319,31 +433,35 @@ final class ThrottleWhileTests: XCTestCase { receiveCompletion: { _ in compEx.fulfill() } ) defer { sub.cancel() } - + subject.send("one") wait(for: [firstEx], timeout: 2) XCTAssertEqual(["one"], values) - + regulator.send(true) regulator.send(false) regulator.send(true) regulator.send(false) - + pause() - + XCTAssertEqual(["one"], values) - + subject.send("one") - + regulator.send(true) subject.send("two") subject.send("three") subject.send("four") regulator.send(false) - + subject.send(completion: .finished) wait(for: [compEx], timeout: 2) - + XCTAssertEqual(["one", "one", "two"], values) } } + +private enum ThrottleError: Error, Equatable { + case failed +} diff --git a/Tests/CombineExtensionsTests/UISchedulerTests.swift b/Tests/CombineExtensionsTests/UISchedulerTests.swift index 7dac9ec..a8d5703 100644 --- a/Tests/CombineExtensionsTests/UISchedulerTests.swift +++ b/Tests/CombineExtensionsTests/UISchedulerTests.swift @@ -4,6 +4,17 @@ import XCTest // Taken from https://github.com/pointfreeco/combine-schedulers/blob/main/Tests/CombineSchedulersTests/UISchedulerTests.swift final class UISchedulerTests: XCTestCase { + func testNowAndMinimumToleranceMirrorMainQueue() throws { + XCTAssertEqual(DispatchQueue.main.minimumTolerance, UIScheduler.shared.minimumTolerance) + + let before = DispatchQueue.main.now + let schedulerNow = UIScheduler.shared.now + let after = DispatchQueue.main.now + + XCTAssertGreaterThanOrEqual(schedulerNow, before) + XCTAssertLessThanOrEqual(schedulerNow, after) + } + func testPublishersOnTheMainThreadPublishImmediately() throws { var didWork = false UIScheduler.shared.schedule { didWork = true }