1use crate::{Clock, Timer};
2use futures::channel::oneshot;
3use futures::future::BoxFuture;
4use jiff::Timestamp;
5use std::sync::{Arc, Mutex, PoisonError};
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(PoisonError::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(PoisonError::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(PoisonError::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: PendingSleeps = {
88 let mut pending = self.pending.lock().unwrap_or_else(PoisonError::into_inner);
89 let (fired, remaining): (PendingSleeps, PendingSleeps) = std::mem::take(&mut *pending)
90 .into_iter()
91 .partition(|(wake_at, _)| *wake_at <= now);
92 *pending = remaining;
93 fired
94 };
95 for (_, sender) in to_fire {
96 let _ = sender.send(());
98 }
99 }
100
101 pub(crate) fn register_sleep(&self, wake_at: Timestamp) -> oneshot::Receiver<()> {
102 let (tx, rx) = oneshot::channel();
103 let now = self.current();
104 if wake_at <= now {
105 let _ = tx.send(());
107 } else {
108 let mut pending = self.pending.lock().unwrap_or_else(PoisonError::into_inner);
109 pending.push((wake_at, tx));
110 }
111 rx
112 }
113}
114
115impl Clock for MockClock {
116 fn now(&self) -> Timestamp {
117 self.current()
118 }
119}
120
121#[derive(Debug)]
126pub struct AdvanceableTimer {
127 clock: Arc<MockClock>,
128}
129
130impl AdvanceableTimer {
131 pub fn new(clock: Arc<MockClock>) -> Self {
132 Self { clock }
133 }
134}
135
136impl Timer for AdvanceableTimer {
137 fn sleep(&self, dur: Duration) -> BoxFuture<'static, ()> {
138 let now = self.clock.current();
139 let wake_at = now.checked_add(dur).unwrap_or(now);
140 let rx = self.clock.register_sleep(wake_at);
141 Box::pin(async move {
142 let _ = rx.await;
144 })
145 }
146
147 fn next_tick(&self) -> BoxFuture<'static, ()> {
148 Box::pin(futures::future::ready(()))
149 }
150}
151
152pub fn frozen_at(rfc3339: &str) -> Arc<dyn Clock> {
156 Arc::new(FrozenClock::at(rfc3339))
157}
158
159pub fn mock_pair(start: Timestamp) -> (Arc<MockClock>, Arc<AdvanceableTimer>) {
161 let clock = Arc::new(MockClock::new(start));
162 let timer = Arc::new(AdvanceableTimer::new(Arc::clone(&clock)));
163 (clock, timer)
164}
165
166#[cfg(test)]
167mod tests {
168 #![allow(clippy::unwrap_used, clippy::expect_used, clippy::panic)]
169 use super::*;
170
171 fn t0() -> Timestamp {
172 "2026-05-10T00:00:00Z".parse().expect("valid RFC 3339")
173 }
174
175 #[test]
178 fn frozen_at_parses_rfc3339() {
179 let clock = FrozenClock::at("2026-05-10T00:00:00Z");
180 assert_eq!(clock.now(), t0());
181 }
182
183 #[test]
184 fn frozen_clock_at_instant() {
185 let clock = FrozenClock::at_instant(t0());
186 assert_eq!(clock.now(), t0());
187 }
188
189 #[test]
190 fn frozen_clock_never_advances() {
191 let clock = FrozenClock::at_instant(t0());
192 let a = clock.now();
193 let b = clock.now();
194 assert_eq!(a, b);
195 }
196
197 #[test]
198 fn frozen_clock_unix_millis() {
199 let clock = FrozenClock::at("1970-01-01T00:00:01Z");
200 assert_eq!(clock.now_unix_millis(), 1_000);
201 }
202
203 #[test]
204 fn frozen_at_helper_trait_dispatch() {
205 let c: Arc<dyn Clock> = frozen_at("2026-05-10T00:00:00Z");
206 assert_eq!(c.now(), t0());
207 }
208
209 #[test]
212 fn mock_clock_starts_at_given_instant() {
213 let clock = MockClock::new(t0());
214 assert_eq!(clock.now(), t0());
215 }
216
217 #[test]
218 fn mock_clock_advance_by_increases_time() {
219 let clock = MockClock::new(t0());
220 clock.advance_by(Duration::from_secs(1));
221 assert_eq!(clock.now(), t0() + jiff::SignedDuration::from_secs(1));
222 }
223
224 #[test]
225 fn mock_clock_set_now_updates_instant() {
226 let clock = MockClock::new(t0());
227 let new_now = t0() + jiff::SignedDuration::from_hours(5);
228 clock.set_now(new_now);
229 assert_eq!(clock.now(), new_now);
230 }
231
232 #[test]
233 fn mock_clock_trait_dispatch() {
234 let c: Arc<dyn Clock> = Arc::new(MockClock::new(t0()));
235 assert_eq!(c.now(), t0());
236 }
237
238 #[tokio::test]
241 async fn advanceable_timer_sleep_resolves_after_advance() {
242 let (clock, timer) = mock_pair(t0());
243 let sleep_fut = timer.sleep(Duration::from_secs(60));
244 let clock2 = Arc::clone(&clock);
246 tokio::spawn(async move {
247 clock2.advance_and_fire(Duration::from_secs(60));
248 })
249 .await
250 .ok();
251 sleep_fut.await; }
253
254 #[tokio::test]
255 async fn advanceable_timer_partial_advance_fires_only_due_sleeps() {
256 let (clock, timer) = mock_pair(t0());
257 let sleep_30 = timer.sleep(Duration::from_secs(30));
258 let sleep_60 = timer.sleep(Duration::from_secs(60));
259
260 clock.advance_and_fire(Duration::from_secs(45));
261
262 sleep_30.await;
264
265 clock.advance_and_fire(Duration::from_secs(15));
267 sleep_60.await;
268 }
269
270 #[tokio::test]
271 async fn advanceable_timer_next_tick_resolves_immediately() {
272 let (_, timer) = mock_pair(t0());
273 timer.next_tick().await;
274 }
275
276 #[test]
277 fn mock_pair_helper_returns_linked_pair() {
278 let (clock, timer) = mock_pair(t0());
279 let _: Arc<MockClock> = clock;
280 let _: Arc<AdvanceableTimer> = timer;
281 }
282
283 #[test]
284 fn frozen_clock_send_sync() {
285 fn assert_send_sync<T: Send + Sync>() {}
286 assert_send_sync::<FrozenClock>();
287 }
288
289 #[test]
290 fn mock_clock_send_sync() {
291 fn assert_send_sync<T: Send + Sync>() {}
292 assert_send_sync::<MockClock>();
293 assert_send_sync::<AdvanceableTimer>();
294 }
295}