1use crate::{Clock, Timer};
2use futures::channel::oneshot;
3use futures::future::BoxFuture;
4use jiff::Timestamp;
5use std::sync::{Arc, Mutex};
6use std::time::Duration;
7
8#[derive(Debug)]
12pub struct FrozenClock {
13 instant: Timestamp,
14}
15
16impl FrozenClock {
17 pub fn at(rfc3339: &str) -> Self {
19 let instant = rfc3339.parse().unwrap_or(Timestamp::UNIX_EPOCH);
20 Self { instant }
21 }
22
23 pub fn at_instant(instant: Timestamp) -> Self {
24 Self { instant }
25 }
26}
27
28impl Clock for FrozenClock {
29 fn now(&self) -> Timestamp {
30 self.instant
31 }
32}
33
34type PendingSleeps = Vec<(Timestamp, oneshot::Sender<()>)>;
37
38#[derive(Debug)]
40pub struct MockClock {
41 inner: Arc<Mutex<Timestamp>>,
42 pending: Arc<Mutex<PendingSleeps>>,
43}
44
45impl MockClock {
46 pub fn new(start: Timestamp) -> Self {
47 Self {
48 inner: Arc::new(Mutex::new(start)),
49 pending: Arc::new(Mutex::new(Vec::new())),
50 }
51 }
52
53 pub fn set_now(&self, instant: Timestamp) {
54 let mut guard = self.inner.lock().unwrap_or_else(|e| e.into_inner());
55 *guard = instant;
56 drop(guard);
57 self.fire_pending();
58 }
59
60 pub fn advance_by(&self, dur: Duration) {
61 let new_now = {
62 let mut guard = self.inner.lock().unwrap_or_else(|e| e.into_inner());
63 *guard = guard.checked_add(dur).unwrap_or(*guard);
64 *guard
65 };
66 self.fire_ready(new_now);
67 }
68
69 pub fn advance_and_fire(&self, dur: Duration) {
72 self.advance_by(dur);
73 }
74
75 fn current(&self) -> Timestamp {
76 *self.inner.lock().unwrap_or_else(|e| e.into_inner())
77 }
78
79 fn fire_pending(&self) {
80 let now = self.current();
81 self.fire_ready(now);
82 }
83
84 fn fire_ready(&self, now: Timestamp) {
85 let to_fire: Vec<oneshot::Sender<()>> = {
87 let mut pending = self.pending.lock().unwrap_or_else(|e| e.into_inner());
88 let mut fired = Vec::new();
89 let mut remaining: PendingSleeps = Vec::new();
90 let mut original: PendingSleeps = Vec::new();
91 std::mem::swap(&mut *pending, &mut original);
92 for (wake_at, sender) in original {
93 if wake_at <= now {
94 fired.push(sender);
95 } else {
96 remaining.push((wake_at, sender));
97 }
98 }
99 *pending = remaining;
100 fired
101 };
102 for sender in to_fire {
103 let _ = sender.send(());
105 }
106 }
107
108 pub(crate) fn register_sleep(&self, wake_at: Timestamp) -> oneshot::Receiver<()> {
109 let (tx, rx) = oneshot::channel();
110 let now = self.current();
111 if wake_at <= now {
112 let _ = tx.send(());
114 } else {
115 let mut pending = self.pending.lock().unwrap_or_else(|e| e.into_inner());
116 pending.push((wake_at, tx));
117 }
118 rx
119 }
120}
121
122impl Clock for MockClock {
123 fn now(&self) -> Timestamp {
124 self.current()
125 }
126}
127
128#[derive(Debug)]
133pub struct AdvanceableTimer {
134 clock: Arc<MockClock>,
135}
136
137impl AdvanceableTimer {
138 pub fn new(clock: Arc<MockClock>) -> Self {
139 Self { clock }
140 }
141}
142
143impl Timer for AdvanceableTimer {
144 fn sleep(&self, dur: Duration) -> BoxFuture<'static, ()> {
145 let now = self.clock.current();
146 let wake_at = now.checked_add(dur).unwrap_or(now);
147 let rx = self.clock.register_sleep(wake_at);
148 Box::pin(async move {
149 let _ = rx.await;
151 })
152 }
153
154 fn next_tick(&self) -> BoxFuture<'static, ()> {
155 Box::pin(futures::future::ready(()))
156 }
157}
158
159use std::sync::Arc as StdArc;
162
163pub fn frozen_at(rfc3339: &str) -> StdArc<dyn Clock> {
165 StdArc::new(FrozenClock::at(rfc3339))
166}
167
168pub fn mock_pair(start: Timestamp) -> (StdArc<MockClock>, StdArc<AdvanceableTimer>) {
170 let clock = StdArc::new(MockClock::new(start));
171 let timer = StdArc::new(AdvanceableTimer::new(StdArc::clone(&clock)));
172 (clock, timer)
173}
174
175#[cfg(test)]
176mod tests {
177 #![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
178 use super::*;
179
180 fn t0() -> Timestamp {
181 "2026-05-10T00:00:00Z".parse().expect("valid RFC 3339")
182 }
183
184 #[test]
187 fn frozen_at_parses_rfc3339() {
188 let clock = FrozenClock::at("2026-05-10T00:00:00Z");
189 assert_eq!(clock.now(), t0());
190 }
191
192 #[test]
193 fn frozen_clock_at_instant() {
194 let clock = FrozenClock::at_instant(t0());
195 assert_eq!(clock.now(), t0());
196 }
197
198 #[test]
199 fn frozen_clock_never_advances() {
200 let clock = FrozenClock::at_instant(t0());
201 let a = clock.now();
202 let b = clock.now();
203 assert_eq!(a, b);
204 }
205
206 #[test]
207 fn frozen_clock_unix_millis() {
208 let clock = FrozenClock::at("1970-01-01T00:00:01Z");
209 assert_eq!(clock.now_unix_millis(), 1_000);
210 }
211
212 #[test]
213 fn frozen_at_helper_trait_dispatch() {
214 let c: StdArc<dyn Clock> = frozen_at("2026-05-10T00:00:00Z");
215 assert_eq!(c.now(), t0());
216 }
217
218 #[test]
221 fn mock_clock_starts_at_given_instant() {
222 let clock = MockClock::new(t0());
223 assert_eq!(clock.now(), t0());
224 }
225
226 #[test]
227 fn mock_clock_advance_by_increases_time() {
228 let clock = MockClock::new(t0());
229 clock.advance_by(Duration::from_secs(1));
230 assert_eq!(clock.now(), t0() + jiff::SignedDuration::from_secs(1));
231 }
232
233 #[test]
234 fn mock_clock_set_now_updates_instant() {
235 let clock = MockClock::new(t0());
236 let new_now = t0() + jiff::SignedDuration::from_hours(5);
237 clock.set_now(new_now);
238 assert_eq!(clock.now(), new_now);
239 }
240
241 #[test]
242 fn mock_clock_trait_dispatch() {
243 let c: StdArc<dyn Clock> = StdArc::new(MockClock::new(t0()));
244 assert_eq!(c.now(), t0());
245 }
246
247 #[tokio::test]
250 async fn advanceable_timer_sleep_resolves_after_advance() {
251 let (clock, timer) = mock_pair(t0());
252 let sleep_fut = timer.sleep(Duration::from_secs(60));
253 let clock2 = StdArc::clone(&clock);
255 tokio::spawn(async move {
256 clock2.advance_and_fire(Duration::from_secs(60));
257 })
258 .await
259 .ok();
260 sleep_fut.await; }
262
263 #[tokio::test]
264 async fn advanceable_timer_partial_advance_fires_only_due_sleeps() {
265 let (clock, timer) = mock_pair(t0());
266 let sleep_30 = timer.sleep(Duration::from_secs(30));
267 let sleep_60 = timer.sleep(Duration::from_secs(60));
268
269 clock.advance_and_fire(Duration::from_secs(45));
270
271 sleep_30.await;
273
274 clock.advance_and_fire(Duration::from_secs(15));
276 sleep_60.await;
277 }
278
279 #[tokio::test]
280 async fn advanceable_timer_next_tick_resolves_immediately() {
281 let (_, timer) = mock_pair(t0());
282 timer.next_tick().await;
283 }
284
285 #[test]
286 fn mock_pair_helper_returns_linked_pair() {
287 let (clock, timer) = mock_pair(t0());
288 let _: StdArc<MockClock> = clock;
289 let _: StdArc<AdvanceableTimer> = timer;
290 }
291
292 #[test]
293 fn frozen_clock_send_sync() {
294 fn assert_send_sync<T: Send + Sync>() {}
295 assert_send_sync::<FrozenClock>();
296 }
297
298 #[test]
299 fn mock_clock_send_sync() {
300 fn assert_send_sync<T: Send + Sync>() {}
301 assert_send_sync::<MockClock>();
302 assert_send_sync::<AdvanceableTimer>();
303 }
304}