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 use super::*;
178
179 fn t0() -> Timestamp {
180 "2026-05-10T00:00:00Z".parse().expect("valid RFC 3339")
181 }
182
183 #[test]
186 fn frozen_at_parses_rfc3339() {
187 let clock = FrozenClock::at("2026-05-10T00:00:00Z");
188 assert_eq!(clock.now(), t0());
189 }
190
191 #[test]
192 fn frozen_clock_at_instant() {
193 let clock = FrozenClock::at_instant(t0());
194 assert_eq!(clock.now(), t0());
195 }
196
197 #[test]
198 fn frozen_clock_never_advances() {
199 let clock = FrozenClock::at_instant(t0());
200 let a = clock.now();
201 let b = clock.now();
202 assert_eq!(a, b);
203 }
204
205 #[test]
206 fn frozen_clock_unix_millis() {
207 let clock = FrozenClock::at("1970-01-01T00:00:01Z");
208 assert_eq!(clock.now_unix_millis(), 1_000);
209 }
210
211 #[test]
212 fn frozen_at_helper_trait_dispatch() {
213 let c: StdArc<dyn Clock> = frozen_at("2026-05-10T00:00:00Z");
214 assert_eq!(c.now(), t0());
215 }
216
217 #[test]
220 fn mock_clock_starts_at_given_instant() {
221 let clock = MockClock::new(t0());
222 assert_eq!(clock.now(), t0());
223 }
224
225 #[test]
226 fn mock_clock_advance_by_increases_time() {
227 let clock = MockClock::new(t0());
228 clock.advance_by(Duration::from_secs(1));
229 assert_eq!(clock.now(), t0() + jiff::SignedDuration::from_secs(1));
230 }
231
232 #[test]
233 fn mock_clock_set_now_updates_instant() {
234 let clock = MockClock::new(t0());
235 let new_now = t0() + jiff::SignedDuration::from_hours(5);
236 clock.set_now(new_now);
237 assert_eq!(clock.now(), new_now);
238 }
239
240 #[test]
241 fn mock_clock_trait_dispatch() {
242 let c: StdArc<dyn Clock> = StdArc::new(MockClock::new(t0()));
243 assert_eq!(c.now(), t0());
244 }
245
246 #[tokio::test]
249 async fn advanceable_timer_sleep_resolves_after_advance() {
250 let (clock, timer) = mock_pair(t0());
251 let sleep_fut = timer.sleep(Duration::from_secs(60));
252 let clock2 = StdArc::clone(&clock);
254 tokio::spawn(async move {
255 clock2.advance_and_fire(Duration::from_secs(60));
256 })
257 .await
258 .ok();
259 sleep_fut.await; }
261
262 #[tokio::test]
263 async fn advanceable_timer_partial_advance_fires_only_due_sleeps() {
264 let (clock, timer) = mock_pair(t0());
265 let sleep_30 = timer.sleep(Duration::from_secs(30));
266 let sleep_60 = timer.sleep(Duration::from_secs(60));
267
268 clock.advance_and_fire(Duration::from_secs(45));
269
270 sleep_30.await;
272
273 clock.advance_and_fire(Duration::from_secs(15));
275 sleep_60.await;
276 }
277
278 #[tokio::test]
279 async fn advanceable_timer_next_tick_resolves_immediately() {
280 let (_, timer) = mock_pair(t0());
281 timer.next_tick().await;
282 }
283
284 #[test]
285 fn mock_pair_helper_returns_linked_pair() {
286 let (clock, timer) = mock_pair(t0());
287 let _: StdArc<MockClock> = clock;
288 let _: StdArc<AdvanceableTimer> = timer;
289 }
290
291 #[test]
292 fn frozen_clock_send_sync() {
293 fn assert_send_sync<T: Send + Sync>() {}
294 assert_send_sync::<FrozenClock>();
295 }
296
297 #[test]
298 fn mock_clock_send_sync() {
299 fn assert_send_sync<T: Send + Sync>() {}
300 assert_send_sync::<MockClock>();
301 assert_send_sync::<AdvanceableTimer>();
302 }
303}