Skip to main content

laminar_core/cluster/control/
lease_deadline.rs

1use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
2use std::time::{Duration, Instant};
3
4/// Process-local monotonic deadline for a renewable durable lease.
5pub 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    /// Create an inactive, renewable deadline for a lease manager before acquisition.
15    #[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    /// Create an irreversibly fenced deadline.
27    #[must_use]
28    pub fn fenced() -> Self {
29        let deadline = Self::uninitialized();
30        deadline.fence();
31        deadline
32    }
33
34    /// Create a deadline live for `remaining`.
35    #[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    /// Irreversibly revoke the lease deadline and wake all waiters.
78    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    /// Withdraw a renewable grant without terminalizing its manager.
86    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    /// Whether the holder remains inside its last successful renewal deadline.
96    #[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    /// Wait until the deadline expires naturally or is terminally fenced.
106    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}