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