Skip to main content

subetha_cxc/
owner_lease.rs

1//! `OwnerLease<T>` - cross-process Mutex with auto-failover.
2//!
3//! Composite primitive: combines a SeqLock-protected payload cell
4//! with the lowest-live-PID + heartbeat-based ownership tracking
5//! pattern. Provides Mutex-like exclusive access to a shared value
6//! T, with the additional guarantee that if the current owner dies,
7//! another process can claim ownership within `grace_epochs` and
8//! continue.
9//!
10//! # API shape
11//!
12//! - `try_acquire(my_pid, grace_epochs)` -> bool: CAS-claim ownership
13//! - `release(my_pid)` -> bool: voluntarily release
14//! - `with_lease(my_pid, grace_epochs, |&mut T| ...)` -> closure-based RAII pattern
15//! - `read_as_owner(my_pid)` -> `Option<T>`: read the value, owner only
16//! - `write_as_owner(my_pid, T)` -> bool: write the value, owner only
17//! - `beat(my_pid)` -> bool: refresh the heartbeat, returns false if no longer owner
18//! - `tick_epoch()`: advance the global epoch (caller responsible for periodic ticks)
19//!
20//! # Layout
21//!
22//! ```text
23//! +-----------------------------+
24//! | LeaseHeader (64B)           |
25//! |   - magic                   |
26//! |   - payload_size            |
27//! |   - seq_version: AtomicU32  | (for payload SeqLock)
28//! |   - owner_pid: AtomicU32    |
29//! |   - lease_term: AtomicU32   |
30//! |   - heartbeat: AtomicU64    |
31//! |   - global_epoch: AtomicU64 |
32//! +-----------------------------+
33//! | payload [u8; PAYLOAD_BYTES] | (one cache line, 48 bytes usable)
34//! +-----------------------------+
35//! ```
36
37use std::fs::{File, OpenOptions};
38use std::marker::PhantomData;
39use std::mem::{align_of, size_of};
40use std::path::Path;
41use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
42
43use memmap2::{MmapMut, MmapOptions};
44
45pub const LEASE_MAGIC: u64 = 0x4150_4D46_4F57_4C53;
46
47pub const PAYLOAD_BYTES: usize = 48;
48
49pub const NO_OWNER: u32 = 0;
50
51#[repr(C, align(64))]
52pub struct LeaseHeader {
53    pub magic: u64,
54    pub payload_size: u32,
55    pub seq_version: AtomicU32,
56    pub owner_pid: AtomicU32,
57    pub lease_term: AtomicU32,
58    pub heartbeat_epoch: AtomicU64,
59    pub global_epoch: AtomicU64,
60    _pad: [u8; 24],
61}
62
63#[repr(C, align(64))]
64pub struct LeasePayload {
65    pub bytes: [u8; PAYLOAD_BYTES],
66    _pad: [u8; 16],
67}
68
69pub const LEASE_FILE_SIZE: usize = size_of::<LeaseHeader>() + size_of::<LeasePayload>();
70
71#[derive(Debug, Clone, Copy, PartialEq, Eq)]
72pub enum LeaseError {
73    LayoutMismatch,
74    PayloadTooLarge,
75    NotOwner,
76    Contention,
77    IoError(std::io::ErrorKind),
78}
79
80impl From<std::io::Error> for LeaseError {
81    fn from(e: std::io::Error) -> Self { Self::IoError(e.kind()) }
82}
83
84pub struct OwnerLease<T: Copy + 'static> {
85    _file: File,
86    mmap: MmapMut,
87    _phantom: PhantomData<T>,
88    header_sidecar: subetha_core::HandshakeHeader,
89    ring_sidecar: Box<subetha_core::ObservationRing>,
90}
91
92unsafe impl<T: Copy + Send + 'static> Send for OwnerLease<T> {}
93unsafe impl<T: Copy + Sync + 'static> Sync for OwnerLease<T> {}
94
95impl<T: Copy + Send + Sync + 'static> subetha_sidecar::AdaptiveInstance for OwnerLease<T> {
96    fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
97    fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
98    fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
99        Box::new(subetha_sidecar::NoMigrationPolicy)
100    }
101}
102
103impl<T: Copy + 'static> OwnerLease<T> {
104    /// Obtain the lease at `path`, initializing it with `initial` only when
105    /// the path does not yet exist. Attaching leaves the current owner and
106    /// term in place; `initial` is then unused. A region built for a
107    /// different payload type is a `LayoutMismatch`. Use
108    /// [`reset`](Self::reset) to deliberately strip a lease.
109    pub fn create(path: impl AsRef<Path>, initial: T) -> Result<Self, LeaseError> {
110        Self::check_layout()?;
111        let (file, mmap) = crate::mmf_attach::create_or_attach(
112            path.as_ref(),
113            LEASE_FILE_SIZE,
114            |ptr| unsafe { Self::init_region(ptr, initial) },
115            |ptr| unsafe { (*(ptr as *const LeaseHeader)).magic == LEASE_MAGIC },
116        )?;
117        let hdr = unsafe { &*(mmap.as_ptr() as *const LeaseHeader) };
118        if hdr.payload_size as usize != size_of::<T>() {
119            return Err(LeaseError::LayoutMismatch);
120        }
121        Ok(Self {
122            _file: file, mmap, _phantom: PhantomData,
123            header_sidecar: subetha_core::HandshakeHeader::new(),
124            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
125        })
126    }
127
128    /// Reinitialise the lease at `path`, stripping whatever owner and term a
129    /// live holder has. For a caller that knows it owns the path.
130    pub fn reset(path: impl AsRef<Path>, initial: T) -> Result<Self, LeaseError> {
131        Self::check_layout()?;
132        let (file, mmap) =
133            crate::mmf_attach::reset(path.as_ref(), LEASE_FILE_SIZE, |ptr| unsafe {
134                Self::init_region(ptr, initial)
135            })?;
136        Ok(Self {
137            _file: file, mmap, _phantom: PhantomData,
138            header_sidecar: subetha_core::HandshakeHeader::new(),
139            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
140        })
141    }
142
143    /// Lay out a fresh lease region: payload first, magic last, because
144    /// attachers spin on the magic and must not observe it before the payload
145    /// is in place.
146    ///
147    /// # Safety
148    /// `ptr` addresses at least [`LEASE_FILE_SIZE`] writable zeroed bytes.
149    unsafe fn init_region(ptr: *mut u8, initial: T) {
150        let hdr = ptr as *mut LeaseHeader;
151        unsafe {
152            std::ptr::write(hdr, LeaseHeader {
153                magic: 0,
154                payload_size: size_of::<T>() as u32,
155                seq_version: AtomicU32::new(0),
156                owner_pid: AtomicU32::new(NO_OWNER),
157                lease_term: AtomicU32::new(0),
158                heartbeat_epoch: AtomicU64::new(0),
159                global_epoch: AtomicU64::new(0),
160                _pad: [0; 24],
161            });
162            let payload_ptr = ptr.add(size_of::<LeaseHeader>()) as *mut LeasePayload;
163            std::ptr::write(payload_ptr, LeasePayload { bytes: [0; PAYLOAD_BYTES], _pad: [0; 16] });
164            let dst = (*payload_ptr).bytes.as_mut_ptr() as *mut T;
165            std::ptr::write_unaligned(dst, initial);
166            std::ptr::write_volatile(std::ptr::addr_of_mut!((*hdr).magic), LEASE_MAGIC);
167        }
168    }
169
170    pub fn open(path: impl AsRef<Path>) -> Result<Self, LeaseError> {
171        Self::check_layout()?;
172        let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
173        if file.metadata()?.len() < LEASE_FILE_SIZE as u64 {
174            return Err(LeaseError::LayoutMismatch);
175        }
176        let mmap = unsafe { MmapOptions::new().len(LEASE_FILE_SIZE).map_mut(&file)? };
177        let hdr = unsafe { &*(mmap.as_ptr() as *const LeaseHeader) };
178        if hdr.magic != LEASE_MAGIC || hdr.payload_size as usize != size_of::<T>() {
179            return Err(LeaseError::LayoutMismatch);
180        }
181        Ok(Self {
182            _file: file, mmap, _phantom: PhantomData,
183            header_sidecar: subetha_core::HandshakeHeader::new(),
184            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
185        })
186    }
187
188    fn check_layout() -> Result<(), LeaseError> {
189        if size_of::<T>() > PAYLOAD_BYTES {
190            return Err(LeaseError::PayloadTooLarge);
191        }
192        if align_of::<T>() > 8 {
193            return Err(LeaseError::PayloadTooLarge);
194        }
195        Ok(())
196    }
197
198    fn header(&self) -> &LeaseHeader {
199        unsafe { &*(self.mmap.as_ptr() as *const LeaseHeader) }
200    }
201
202    fn payload_ptr(&self) -> *mut u8 {
203        unsafe { self.mmap.as_ptr().add(size_of::<LeaseHeader>()) as *mut u8 }
204    }
205
206    /// Try to claim ownership. Succeeds when (a) no current owner,
207    /// (b) my_pid < current owner's PID (preemption), or
208    /// (c) current owner's heartbeat is more than grace_epochs stale.
209    pub fn try_acquire(&self, my_pid: u32, grace_epochs: u64) -> bool {
210        assert!(my_pid != NO_OWNER, "PID 0 reserved for NO_OWNER");
211        let header = self.header();
212        loop {
213            let cur = header.owner_pid.load(Ordering::Acquire);
214            let can_claim = if cur == NO_OWNER {
215                true
216            } else if cur == my_pid {
217                self.ring_sidecar
218                    .push_op(crate::sidecar_ops::ownership::OP_ACQUIRE, 0);
219                return true;
220            } else if my_pid < cur {
221                true
222            } else {
223                let beat = header.heartbeat_epoch.load(Ordering::Acquire);
224                let global = header.global_epoch.load(Ordering::Acquire);
225                global.saturating_sub(beat) > grace_epochs
226            };
227            if !can_claim {
228                self.ring_sidecar
229                    .push_op(crate::sidecar_ops::ownership::OP_ACQUIRE, 1);
230                return false;
231            }
232            if header.owner_pid.compare_exchange(
233                cur, my_pid, Ordering::AcqRel, Ordering::Acquire,
234            ).is_ok() {
235                header.lease_term.fetch_add(1, Ordering::AcqRel);
236                let global = header.global_epoch.load(Ordering::Acquire);
237                header.heartbeat_epoch.store(global, Ordering::Release);
238                self.ring_sidecar
239                    .push_op(crate::sidecar_ops::ownership::OP_ACQUIRE, 0);
240                return true;
241            }
242            std::hint::spin_loop();
243        }
244    }
245
246    pub fn release(&self, my_pid: u32) -> bool {
247        let ok = self.header().owner_pid
248            .compare_exchange(my_pid, NO_OWNER, Ordering::AcqRel, Ordering::Acquire)
249            .is_ok();
250        self.ring_sidecar.push_op(
251            crate::sidecar_ops::ownership::OP_RELEASE,
252            if ok { 0 } else { 1 },
253        );
254        ok
255    }
256
257    /// Closure-based lease scope. Tries to acquire; if successful,
258    /// runs `f(&mut T)` with the payload, then releases. Returns
259    /// `None` if acquisition failed.
260    pub fn with_lease<R, F: FnOnce(&mut T) -> R>(
261        &self,
262        my_pid: u32,
263        grace_epochs: u64,
264        f: F,
265    ) -> Option<R> {
266        if !self.try_acquire(my_pid, grace_epochs) { return None; }
267        let result = {
268            // Read current payload, hand to closure, then write back.
269            let header = self.header();
270            let mut value: T = unsafe {
271                let src = self.payload_ptr() as *const T;
272                std::ptr::read_unaligned(src)
273            };
274            let r = f(&mut value);
275            // Publish updated payload via SeqLock so external readers
276            // (when added) see consistent state.
277            header.seq_version.fetch_add(1, Ordering::AcqRel);
278            unsafe {
279                let dst = self.payload_ptr() as *mut T;
280                std::ptr::write_unaligned(dst, value);
281            }
282            header.seq_version.fetch_add(1, Ordering::Release);
283            r
284        };
285        self.release(my_pid);
286        Some(result)
287    }
288
289    /// Read the payload only when the caller holds the lease.
290    pub fn read_as_owner(&self, my_pid: u32) -> Option<T> {
291        if !self.am_i_owner(my_pid) {
292            self.ring_sidecar
293                .push_op(crate::sidecar_ops::ownership::OP_GET, 1);
294            return None;
295        }
296        let value: T = unsafe {
297            let src = self.payload_ptr() as *const T;
298            std::ptr::read_unaligned(src)
299        };
300        self.ring_sidecar
301            .push_op(crate::sidecar_ops::ownership::OP_GET, 0);
302        Some(value)
303    }
304
305    /// Write the payload only when the caller holds the lease.
306    pub fn write_as_owner(&self, my_pid: u32, value: T) -> bool {
307        if !self.am_i_owner(my_pid) {
308            self.ring_sidecar
309                .push_op(crate::sidecar_ops::ownership::OP_GET, 1);
310            return false;
311        }
312        let header = self.header();
313        header.seq_version.fetch_add(1, Ordering::AcqRel);
314        unsafe {
315            let dst = self.payload_ptr() as *mut T;
316            std::ptr::write_unaligned(dst, value);
317        }
318        header.seq_version.fetch_add(1, Ordering::Release);
319        self.ring_sidecar
320            .push_op(crate::sidecar_ops::ownership::OP_GET, 0);
321        true
322    }
323
324    /// Refresh the heartbeat. Returns `false` if the caller has
325    /// been preempted (no longer owner).
326    pub fn beat(&self, my_pid: u32) -> bool {
327        let header = self.header();
328        if header.owner_pid.load(Ordering::Acquire) != my_pid {
329            self.ring_sidecar
330                .push_op(crate::sidecar_ops::ownership::OP_BEAT, 1);
331            return false;
332        }
333        let global = header.global_epoch.load(Ordering::Acquire);
334        header.heartbeat_epoch.store(global, Ordering::Release);
335        self.ring_sidecar
336            .push_op(crate::sidecar_ops::ownership::OP_BEAT, 0);
337        true
338    }
339
340    pub fn tick_epoch(&self) -> u64 {
341        self.header().global_epoch.fetch_add(1, Ordering::AcqRel) + 1
342    }
343
344    pub fn current_owner(&self) -> Option<u32> {
345        let pid = self.header().owner_pid.load(Ordering::Acquire);
346        if pid == NO_OWNER { None } else { Some(pid) }
347    }
348
349    pub fn am_i_owner(&self, my_pid: u32) -> bool {
350        self.header().owner_pid.load(Ordering::Acquire) == my_pid
351    }
352
353    pub fn lease_term(&self) -> u32 {
354        self.header().lease_term.load(Ordering::Acquire)
355    }
356
357    pub fn flush(&self) -> Result<(), LeaseError> {
358        self.mmap.flush()?;
359        Ok(())
360    }
361
362    /// Non-blocking flush: schedules a writeback via the OS.
363    /// Note: Windows is only partially async (sync to page cache,
364    /// not to disk).
365    pub fn flush_async(&self) -> Result<(), LeaseError> {
366        self.mmap.flush_async()?;
367        Ok(())
368    }
369}
370
371#[cfg(test)]
372mod tests {
373    use super::*;
374
375    fn tmp(name: &str) -> std::path::PathBuf {
376        let mut p = std::env::temp_dir();
377        let pid = std::process::id();
378        p.push(format!("subetha-lease-{name}-{pid}.bin"));
379        p
380    }
381
382    /// A second `create` on the same path attaches: the holder keeps its
383    /// lease and the payload the owner wrote survives. `reset` is the call
384    /// that strips both.
385    #[test]
386    fn second_create_attaches_and_keeps_the_holder() {
387        let p = tmp("attach");
388        std::fs::remove_file(&p).ok();
389
390        let holder: OwnerLease<u64> = OwnerLease::create(&p, 42).unwrap();
391        assert!(holder.try_acquire(100, 3));
392        assert!(holder.write_as_owner(100, 777));
393
394        let late: OwnerLease<u64> = OwnerLease::create(&p, 42).unwrap();
395        assert_eq!(late.current_owner(), Some(100), "attach stripped the lease");
396        assert_eq!(late.read_as_owner(100), Some(777), "attach cleared the payload");
397
398        // Windows refuses to truncate a file with live mappings, so reset only
399        // works once every handle is gone - which is the ownership reset
400        // demands anyway.
401        drop(late);
402        drop(holder);
403        let fresh: OwnerLease<u64> = OwnerLease::reset(&p, 42).unwrap();
404        assert_eq!(fresh.current_owner(), None);
405        std::fs::remove_file(&p).ok();
406    }
407
408    #[test]
409    fn acquire_release_round_trip() {
410        let p = tmp("rt");
411        let l: OwnerLease<u64> = OwnerLease::create(&p, 42).unwrap();
412        assert_eq!(l.current_owner(), None);
413        assert!(l.try_acquire(100, 3));
414        assert_eq!(l.current_owner(), Some(100));
415        assert!(l.am_i_owner(100));
416        assert!(l.release(100));
417        assert_eq!(l.current_owner(), None);
418        std::fs::remove_file(&p).ok();
419    }
420
421    #[test]
422    fn lower_pid_preempts_higher() {
423        let p = tmp("preempt");
424        let l: OwnerLease<u64> = OwnerLease::create(&p, 0).unwrap();
425        assert!(l.try_acquire(500, 3));
426        assert!(l.try_acquire(100, 3));
427        assert_eq!(l.current_owner(), Some(100));
428        assert!(!l.try_acquire(500, 3));
429        std::fs::remove_file(&p).ok();
430    }
431
432    #[test]
433    fn read_and_write_require_ownership() {
434        let p = tmp("read-write");
435        let l: OwnerLease<u64> = OwnerLease::create(&p, 100).unwrap();
436        // No owner; reads/writes fail.
437        assert_eq!(l.read_as_owner(50), None);
438        assert!(!l.write_as_owner(50, 999));
439        // Acquire.
440        assert!(l.try_acquire(50, 3));
441        assert_eq!(l.read_as_owner(50), Some(100));
442        assert!(l.write_as_owner(50, 999));
443        assert_eq!(l.read_as_owner(50), Some(999));
444        // Non-owner reads fail even after another's acquire.
445        assert_eq!(l.read_as_owner(999), None);
446        std::fs::remove_file(&p).ok();
447    }
448
449    #[test]
450    fn with_lease_acquires_runs_releases() {
451        let p = tmp("with-lease");
452        let l: OwnerLease<u64> = OwnerLease::create(&p, 10).unwrap();
453        let r = l.with_lease(100, 3, |v| {
454            *v += 5;
455            *v
456        });
457        assert_eq!(r, Some(15));
458        assert_eq!(l.current_owner(), None, "with_lease releases on return");
459        // Verify persistence.
460        assert!(l.try_acquire(100, 3));
461        assert_eq!(l.read_as_owner(100), Some(15));
462        std::fs::remove_file(&p).ok();
463    }
464
465    #[test]
466    fn stale_owner_failover_within_grace() {
467        let p = tmp("failover");
468        let l: OwnerLease<u64> = OwnerLease::create(&p, 0).unwrap();
469        assert!(l.try_acquire(100, 1));
470        // Tick beyond grace without 100 beating.
471        l.tick_epoch();
472        l.tick_epoch();
473        // Higher PID 500 can now preempt because 100's heartbeat is stale.
474        assert!(l.try_acquire(500, 1));
475        assert_eq!(l.current_owner(), Some(500));
476        std::fs::remove_file(&p).ok();
477    }
478
479    #[test]
480    fn beat_keeps_owner_alive() {
481        let p = tmp("beat");
482        let l: OwnerLease<u64> = OwnerLease::create(&p, 0).unwrap();
483        assert!(l.try_acquire(100, 1));
484        l.tick_epoch();
485        assert!(l.beat(100));
486        // 500 cannot preempt while 100 beats.
487        assert!(!l.try_acquire(500, 1));
488        std::fs::remove_file(&p).ok();
489    }
490
491    #[test]
492    fn lease_term_advances_on_handover() {
493        let p = tmp("term");
494        let l: OwnerLease<u64> = OwnerLease::create(&p, 0).unwrap();
495        let t0 = l.lease_term();
496        l.try_acquire(500, 3);
497        let t1 = l.lease_term();
498        assert_eq!(t1, t0 + 1);
499        l.try_acquire(100, 3);  // preempt
500        let t2 = l.lease_term();
501        assert_eq!(t2, t1 + 1);
502        l.try_acquire(100, 3);  // idempotent
503        assert_eq!(l.lease_term(), t2);
504        std::fs::remove_file(&p).ok();
505    }
506
507    #[test]
508    fn cross_handle_lease_visibility() {
509        let p = tmp("cross-handle");
510        let a: OwnerLease<u64> = OwnerLease::create(&p, 7).unwrap();
511        let b: OwnerLease<u64> = OwnerLease::open(&p).unwrap();
512        assert!(a.try_acquire(100, 3));
513        assert_eq!(b.current_owner(), Some(100));
514        assert!(b.am_i_owner(100));
515        assert!(!b.am_i_owner(200));
516        // Process B (who isn't owner) can't write.
517        assert!(!b.write_as_owner(200, 99));
518        // Process A writes; B sees via own acquire later.
519        a.write_as_owner(100, 88);
520        a.release(100);
521        assert!(b.try_acquire(200, 3));
522        assert_eq!(b.read_as_owner(200), Some(88));
523        std::fs::remove_file(&p).ok();
524    }
525
526    #[test]
527    fn disk_persistence_survives_reopen() {
528        let p = tmp("disk-persist");
529        {
530            let l: OwnerLease<u64> = OwnerLease::create(&p, 1234).unwrap();
531            l.try_acquire(42, 3);
532            l.write_as_owner(42, 5678);
533            l.flush().unwrap();
534        }
535        let l2: OwnerLease<u64> = OwnerLease::open(&p).unwrap();
536        assert_eq!(l2.current_owner(), Some(42));
537        assert_eq!(l2.read_as_owner(42), Some(5678));
538        std::fs::remove_file(&p).ok();
539    }
540
541    #[test]
542    fn beat_after_preemption_returns_false() {
543        let p = tmp("beat-preempt");
544        let l: OwnerLease<u64> = OwnerLease::create(&p, 0).unwrap();
545        l.try_acquire(500, 3);
546        l.try_acquire(100, 3);  // 100 preempts
547        assert!(!l.beat(500));
548        assert!(l.beat(100));
549        std::fs::remove_file(&p).ok();
550    }
551
552    #[test]
553    fn struct_payload_round_trip() {
554        #[derive(Clone, Copy, Debug, PartialEq)]
555        #[repr(C)]
556        struct State { active: u32, count: u32, score: f64 }
557        let p = tmp("struct");
558        let l: OwnerLease<State> = OwnerLease::create(&p,
559            State { active: 0, count: 0, score: 0.0 }).unwrap();
560        l.with_lease(7, 3, |s| {
561            s.active = 1;
562            s.count = 42;
563            s.score = 2.5;
564        });
565        // Re-acquire to read.
566        l.try_acquire(7, 3);
567        let s = l.read_as_owner(7).unwrap();
568        assert_eq!(s, State { active: 1, count: 42, score: 2.5 });
569        std::fs::remove_file(&p).ok();
570    }
571}