laminar_core/cluster/control/
lease_deadline.rs1use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
2use std::time::{Duration, Instant};
3
4pub struct LeaseDeadline {
6 origin: Instant,
7 valid_until_ns: AtomicU64,
8 terminal: AtomicBool,
9 transition: parking_lot::Mutex<()>,
10 changed: tokio::sync::Notify,
11}
12
13impl LeaseDeadline {
14 #[must_use]
16 pub(crate) fn uninitialized() -> Self {
17 Self {
18 origin: Instant::now(),
19 valid_until_ns: AtomicU64::new(0),
20 terminal: AtomicBool::new(false),
21 transition: parking_lot::Mutex::new(()),
22 changed: tokio::sync::Notify::new(),
23 }
24 }
25
26 #[must_use]
28 pub fn fenced() -> Self {
29 let deadline = Self::uninitialized();
30 deadline.fence();
31 deadline
32 }
33
34 #[must_use]
36 pub fn live_for(remaining: Duration) -> Self {
37 let deadline = Self::uninitialized();
38 deadline.extend(remaining);
39 deadline
40 }
41
42 pub(crate) fn extend(&self, remaining: Duration) {
43 let until = self.origin.elapsed().saturating_add(remaining).as_nanos();
44 let until = u64::try_from(until).unwrap_or(u64::MAX).max(1);
45 self.update_valid_until(until);
46 }
47
48 pub(crate) fn extend_until(&self, valid_until: Instant) {
49 let Some(remaining_from_origin) = valid_until.checked_duration_since(self.origin) else {
50 self.fence();
51 return;
52 };
53 if remaining_from_origin.is_zero() {
54 self.fence();
55 return;
56 }
57 let until = u64::try_from(remaining_from_origin.as_nanos()).unwrap_or(u64::MAX);
58 self.update_valid_until(until);
59 }
60
61 fn update_valid_until(&self, valid_until_ns: u64) {
62 let _transition = self.transition.lock();
63 if self.terminal.load(Ordering::Acquire) {
64 return;
65 }
66 let current = self.valid_until_ns.load(Ordering::Acquire);
67 if current != 0 && self.origin.elapsed().as_nanos() >= u128::from(current) {
68 self.terminal.store(true, Ordering::Release);
69 self.valid_until_ns.store(0, Ordering::Release);
70 self.changed.notify_waiters();
71 return;
72 }
73 self.valid_until_ns.store(valid_until_ns, Ordering::Release);
74 self.changed.notify_waiters();
75 }
76
77 pub fn fence(&self) {
79 let _transition = self.transition.lock();
80 self.terminal.store(true, Ordering::Release);
81 self.valid_until_ns.store(0, Ordering::Release);
82 self.changed.notify_waiters();
83 }
84
85 pub(crate) fn withdraw(&self) {
87 let _transition = self.transition.lock();
88 if self.terminal.load(Ordering::Acquire) {
89 return;
90 }
91 self.valid_until_ns.store(0, Ordering::Release);
92 self.changed.notify_waiters();
93 }
94
95 #[must_use]
97 pub fn is_live(&self) -> bool {
98 if self.terminal.load(Ordering::Acquire) {
99 return false;
100 }
101 let deadline = self.valid_until_ns.load(Ordering::Acquire);
102 deadline != 0 && self.origin.elapsed().as_nanos() < u128::from(deadline)
103 }
104
105 pub async fn wait_until_expired(&self) {
107 loop {
108 let changed = self.changed.notified();
109 tokio::pin!(changed);
110 changed.as_mut().enable();
111
112 if self.terminal.load(Ordering::Acquire) {
113 return;
114 }
115 let valid_until_ns = self.valid_until_ns.load(Ordering::Acquire);
116 let elapsed_ns = self.origin.elapsed().as_nanos();
117 if valid_until_ns == 0 {
118 return;
119 }
120 if elapsed_ns >= u128::from(valid_until_ns) {
121 let _transition = self.transition.lock();
122 if self.terminal.load(Ordering::Acquire) {
123 return;
124 }
125 let current = self.valid_until_ns.load(Ordering::Acquire);
126 if current != 0 && self.origin.elapsed().as_nanos() < u128::from(current) {
127 continue;
128 }
129 self.terminal.store(true, Ordering::Release);
130 self.valid_until_ns.store(0, Ordering::Release);
131 self.changed.notify_waiters();
132 return;
133 }
134 let remaining_ns = u128::from(valid_until_ns).saturating_sub(elapsed_ns);
135 let remaining = Duration::from_nanos(u64::try_from(remaining_ns).unwrap_or(u64::MAX));
136 tokio::select! {
137 () = tokio::time::sleep(remaining) => {}
138 () = &mut changed => {}
139 }
140 }
141 }
142}
143
144impl std::fmt::Debug for LeaseDeadline {
145 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
146 formatter
147 .debug_struct("LeaseDeadline")
148 .field("terminal", &self.terminal.load(Ordering::Acquire))
149 .field("live", &self.is_live())
150 .finish_non_exhaustive()
151 }
152}
153
154#[cfg(test)]
155mod tests {
156 use super::*;
157
158 #[test]
159 fn absolute_extension_preserves_the_original_deadline() {
160 let deadline = LeaseDeadline::uninitialized();
161 let valid_for = Duration::from_millis(37);
162 let valid_until = deadline.origin.checked_add(valid_for).unwrap();
163
164 deadline.extend_until(valid_until);
165
166 assert_eq!(
167 deadline.valid_until_ns.load(Ordering::Acquire),
168 u64::try_from(valid_for.as_nanos()).unwrap()
169 );
170 }
171
172 #[test]
173 fn naturally_expired_deadline_cannot_be_resurrected() {
174 let deadline = LeaseDeadline::live_for(Duration::from_nanos(1));
175 while deadline.is_live() {
176 std::hint::spin_loop();
177 }
178
179 assert!(!deadline.is_live());
180 deadline.extend(Duration::from_secs(60));
181
182 assert!(!deadline.is_live());
183 assert!(deadline.terminal.load(Ordering::Acquire));
184 assert_eq!(deadline.valid_until_ns.load(Ordering::Acquire), 0);
185 }
186
187 #[tokio::test]
188 async fn renewal_before_expiry_extends_an_existing_waiter() {
189 let deadline = std::sync::Arc::new(LeaseDeadline::live_for(Duration::from_millis(100)));
190 let mut waiting = {
191 let deadline = std::sync::Arc::clone(&deadline);
192 tokio::spawn(async move { deadline.wait_until_expired().await })
193 };
194 tokio::time::sleep(Duration::from_millis(20)).await;
195
196 deadline.extend(Duration::from_secs(60));
197
198 assert!(
199 tokio::time::timeout(Duration::from_millis(120), &mut waiting)
200 .await
201 .is_err(),
202 "the waiter used the superseded pre-renewal deadline"
203 );
204 assert!(deadline.is_live());
205 deadline.fence();
206 tokio::time::timeout(Duration::from_secs(1), waiting)
207 .await
208 .expect("fencing did not stop the renewed waiter")
209 .unwrap();
210 }
211
212 #[test]
213 fn terminal_fence_rejects_every_later_extension() {
214 let deadline = LeaseDeadline::live_for(Duration::from_secs(60));
215 deadline.fence();
216
217 deadline.extend(Duration::from_secs(60));
218 deadline.extend_until(Instant::now() + Duration::from_secs(60));
219
220 assert!(deadline.terminal.load(Ordering::Acquire));
221 assert!(!deadline.is_live());
222 assert_eq!(deadline.valid_until_ns.load(Ordering::Acquire), 0);
223 }
224
225 #[test]
226 fn withdrawn_renewable_grant_can_be_reacquired() {
227 let deadline = LeaseDeadline::live_for(Duration::from_secs(60));
228
229 deadline.withdraw();
230 assert!(!deadline.is_live());
231 assert!(!deadline.terminal.load(Ordering::Acquire));
232
233 deadline.extend(Duration::from_secs(60));
234 assert!(deadline.is_live());
235 }
236
237 #[test]
238 fn terminal_fence_wins_concurrent_extension() {
239 for _ in 0..64 {
240 let deadline = std::sync::Arc::new(LeaseDeadline::live_for(Duration::from_secs(60)));
241 let barrier = std::sync::Arc::new(std::sync::Barrier::new(3));
242 std::thread::scope(|scope| {
243 let extending = std::sync::Arc::clone(&deadline);
244 let extending_barrier = std::sync::Arc::clone(&barrier);
245 scope.spawn(move || {
246 extending_barrier.wait();
247 extending.extend(Duration::from_secs(120));
248 extending.extend_until(Instant::now() + Duration::from_secs(120));
249 });
250
251 let fencing = std::sync::Arc::clone(&deadline);
252 let fencing_barrier = std::sync::Arc::clone(&barrier);
253 scope.spawn(move || {
254 fencing_barrier.wait();
255 fencing.fence();
256 });
257
258 barrier.wait();
259 });
260
261 assert!(deadline.terminal.load(Ordering::Acquire));
262 assert!(!deadline.is_live());
263 assert_eq!(deadline.valid_until_ns.load(Ordering::Acquire), 0);
264 }
265 }
266
267 #[tokio::test]
268 async fn terminal_fence_wakes_expiry_waiter() {
269 let deadline = std::sync::Arc::new(LeaseDeadline::live_for(Duration::from_secs(60)));
270 let waiting = {
271 let deadline = std::sync::Arc::clone(&deadline);
272 tokio::spawn(async move { deadline.wait_until_expired().await })
273 };
274 tokio::task::yield_now().await;
275
276 deadline.fence();
277
278 tokio::time::timeout(Duration::from_secs(1), waiting)
279 .await
280 .expect("terminal fence did not wake the expiry waiter")
281 .unwrap();
282 }
283}