diff --git a/Sources/CombineExtensions/AgainAt.swift b/Sources/CombineExtensions/AgainAt.swift index 8e56c74..ee92a78 100644 --- a/Sources/CombineExtensions/AgainAt.swift +++ b/Sources/CombineExtensions/AgainAt.swift @@ -50,19 +50,27 @@ public extension Publishers.AgainAt { final class Timer: @unchecked Sendable { public let now: Context.SchedulerTimeType + private let scheduler: Context private let onRepublishAt: @Sendable (Context.SchedulerTimeType) -> Void fileprivate init( now: Context.SchedulerTimeType, + scheduler: Context, onRepublishAt: @escaping @Sendable (Context.SchedulerTimeType) -> Void ) { self.now = now + self.scheduler = scheduler self.onRepublishAt = onRepublishAt } public func republish(at time: Context.SchedulerTimeType) { onRepublishAt(time) } + + public func time(at date: Date) -> Context.SchedulerTimeType { + let nanoseconds = date.timeIntervalSinceNow.nanoseconds + return scheduler.now.advanced(by: .nanoseconds(nanoseconds)) + } } } @@ -300,6 +308,7 @@ private final class AgainAtSubscription: case let .value(downstream, output): let timer = Publishers.AgainAt.Timer( now: scheduler.now, + scheduler: scheduler, onRepublishAt: { [weak self] time in self?.republish(at: time) } @@ -342,3 +351,16 @@ private final class AgainAtSubscription: } } } + +private extension TimeInterval { + var nanoseconds: Int { + guard self > 0 else { return 0 } + let nanoseconds = (self * 1_000_000_000).rounded(.up) + guard nanoseconds < Double(maxNanoseconds) else { + return maxNanoseconds + } + return Int(nanoseconds) + } +} + +private let maxNanoseconds = Int.max - 1024 diff --git a/Tests/CombineExtensionsTests/AgainAtTests.swift b/Tests/CombineExtensionsTests/AgainAtTests.swift index 71f0f89..33c7327 100644 --- a/Tests/CombineExtensionsTests/AgainAtTests.swift +++ b/Tests/CombineExtensionsTests/AgainAtTests.swift @@ -115,6 +115,90 @@ final class AgainAtTests: XCTestCase { XCTAssertEqual(values, [1, 1]) } + func testRepublishAtFutureDatePublishesAfterConvertedTime() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [DateRepublisher]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { value, timer in + values.append(value) + republishers.append { timer.republish(at: timer.time(at: $0)) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + republishers[0](Date(timeIntervalSinceNow: 2)) + + scheduler.advance(by: .seconds(1)) + XCTAssertEqual(values, [1]) + + scheduler.advance(by: .seconds(2)) + XCTAssertEqual(values, [1, 1]) + } + + func testRepublishAtPastDateUsesCurrentSchedulerTime() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var values = [Int]() + var republishers = [DateRepublisher]() + 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: timer.time(at: $0)) } + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + let firstNow = scheduler.now + scheduler.advance(by: .seconds(5)) + let republishNow = scheduler.now + + republishers[0](.distantPast) + scheduler.advance() + + XCTAssertEqual(values, [1, 1]) + XCTAssertEqual(timerNows, [firstNow, republishNow]) + } + + func testTimeAtDistantFutureDoesNotOverflow() { + let scheduler = DispatchQueue.test + let subject = PassthroughSubject() + var convertedTimes = [DispatchQueue.SchedulerTimeType]() + + let subscription = subject + .againAt(scheduler: scheduler) + .sink( + receiveCompletion: { _ in }, + receiveValue: { _, timer in + convertedTimes.append(timer.time(at: .distantFuture)) + } + ) + defer { subscription.cancel() } + + subject.send(1) + scheduler.advance() + + XCTAssertEqual(convertedTimes.count, 1) + XCTAssertTrue(convertedTimes[0] > scheduler.now) + } + func testRepublishPublishesNewestUpstreamOutputAtFireTime() { let scheduler = DispatchQueue.test let subject = PassthroughSubject() @@ -642,6 +726,7 @@ final class AgainAtTests: XCTestCase { } private typealias Republisher = (DispatchQueue.SchedulerTimeType) -> Void +private typealias DateRepublisher = (Date) -> Void private enum TestError: Error, Equatable { case failed