Skip to content

Commit a806ba9

Browse files
committed
Add Publishers.AgainAt.Timer.time(at:)
1 parent e8e118a commit a806ba9

2 files changed

Lines changed: 107 additions & 0 deletions

File tree

Sources/CombineExtensions/AgainAt.swift

Lines changed: 22 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -50,19 +50,27 @@ public extension Publishers.AgainAt {
5050
final class Timer: @unchecked Sendable {
5151
public let now: Context.SchedulerTimeType
5252

53+
private let scheduler: Context
5354
private let onRepublishAt: @Sendable (Context.SchedulerTimeType) -> Void
5455

5556
fileprivate init(
5657
now: Context.SchedulerTimeType,
58+
scheduler: Context,
5759
onRepublishAt: @escaping @Sendable (Context.SchedulerTimeType) -> Void
5860
) {
5961
self.now = now
62+
self.scheduler = scheduler
6063
self.onRepublishAt = onRepublishAt
6164
}
6265

6366
public func republish(at time: Context.SchedulerTimeType) {
6467
onRepublishAt(time)
6568
}
69+
70+
public func time(at date: Date) -> Context.SchedulerTimeType {
71+
let nanoseconds = date.timeIntervalSinceNow.nanoseconds
72+
return scheduler.now.advanced(by: .nanoseconds(nanoseconds))
73+
}
6674
}
6775
}
6876

@@ -300,6 +308,7 @@ private final class AgainAtSubscription<Upstream, Context, Downstream>:
300308
case let .value(downstream, output):
301309
let timer = Publishers.AgainAt<Upstream, Context>.Timer(
302310
now: scheduler.now,
311+
scheduler: scheduler,
303312
onRepublishAt: { [weak self] time in
304313
self?.republish(at: time)
305314
}
@@ -342,3 +351,16 @@ private final class AgainAtSubscription<Upstream, Context, Downstream>:
342351
}
343352
}
344353
}
354+
355+
private extension TimeInterval {
356+
var nanoseconds: Int {
357+
guard self > 0 else { return 0 }
358+
let nanoseconds = self * 1_000_000_000
359+
guard nanoseconds < Double(maxNanoseconds) else {
360+
return maxNanoseconds
361+
}
362+
return Int(nanoseconds)
363+
}
364+
}
365+
366+
private let maxNanoseconds = Int.max - 1024

Tests/CombineExtensionsTests/AgainAtTests.swift

Lines changed: 85 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -115,6 +115,90 @@ final class AgainAtTests: XCTestCase {
115115
XCTAssertEqual(values, [1, 1])
116116
}
117117

118+
func testRepublishAtFutureDatePublishesAfterConvertedTime() {
119+
let scheduler = DispatchQueue.test
120+
let subject = PassthroughSubject<Int, Never>()
121+
var values = [Int]()
122+
var republishers = [DateRepublisher]()
123+
124+
let subscription = subject
125+
.againAt(scheduler: scheduler)
126+
.sink(
127+
receiveCompletion: { _ in },
128+
receiveValue: { value, timer in
129+
values.append(value)
130+
republishers.append { timer.republish(at: timer.time(at: $0)) }
131+
}
132+
)
133+
defer { subscription.cancel() }
134+
135+
subject.send(1)
136+
scheduler.advance()
137+
138+
republishers[0](Date(timeIntervalSinceNow: 1))
139+
140+
scheduler.advance(by: .milliseconds(500))
141+
XCTAssertEqual(values, [1])
142+
143+
scheduler.advance(by: .seconds(1))
144+
XCTAssertEqual(values, [1, 1])
145+
}
146+
147+
func testRepublishAtPastDateUsesCurrentSchedulerTime() {
148+
let scheduler = DispatchQueue.test
149+
let subject = PassthroughSubject<Int, Never>()
150+
var values = [Int]()
151+
var republishers = [DateRepublisher]()
152+
var timerNows = [DispatchQueue.SchedulerTimeType]()
153+
154+
let subscription = subject
155+
.againAt(scheduler: scheduler)
156+
.sink(
157+
receiveCompletion: { _ in },
158+
receiveValue: { value, timer in
159+
values.append(value)
160+
timerNows.append(timer.now)
161+
republishers.append { timer.republish(at: timer.time(at: $0)) }
162+
}
163+
)
164+
defer { subscription.cancel() }
165+
166+
subject.send(1)
167+
scheduler.advance()
168+
169+
let firstNow = scheduler.now
170+
scheduler.advance(by: .seconds(5))
171+
let republishNow = scheduler.now
172+
173+
republishers[0](.distantPast)
174+
scheduler.advance()
175+
176+
XCTAssertEqual(values, [1, 1])
177+
XCTAssertEqual(timerNows, [firstNow, republishNow])
178+
}
179+
180+
func testTimeAtDistantFutureDoesNotOverflow() {
181+
let scheduler = DispatchQueue.test
182+
let subject = PassthroughSubject<Int, Never>()
183+
var convertedTimes = [DispatchQueue.SchedulerTimeType]()
184+
185+
let subscription = subject
186+
.againAt(scheduler: scheduler)
187+
.sink(
188+
receiveCompletion: { _ in },
189+
receiveValue: { _, timer in
190+
convertedTimes.append(timer.time(at: .distantFuture))
191+
}
192+
)
193+
defer { subscription.cancel() }
194+
195+
subject.send(1)
196+
scheduler.advance()
197+
198+
XCTAssertEqual(convertedTimes.count, 1)
199+
XCTAssertTrue(convertedTimes[0] > scheduler.now)
200+
}
201+
118202
func testRepublishPublishesNewestUpstreamOutputAtFireTime() {
119203
let scheduler = DispatchQueue.test
120204
let subject = PassthroughSubject<Int, Never>()
@@ -642,6 +726,7 @@ final class AgainAtTests: XCTestCase {
642726
}
643727

644728
private typealias Republisher = (DispatchQueue.SchedulerTimeType) -> Void
729+
private typealias DateRepublisher = (Date) -> Void
645730

646731
private enum TestError: Error, Equatable {
647732
case failed

0 commit comments

Comments
 (0)