Skip to main content

subetha_cxc/
shared_leader_election.rs

1//! `SharedLeaderElection` - cross-process leader election with
2//! lowest-live-PID semantics and heartbeat-driven failover.
3//!
4//! Each process can call [`SharedLeaderElection::try_claim_leadership`]
5//! which atomically claims the leader role if:
6//! 1. There is no current leader (PID == 0), OR
7//! 2. The caller's PID is strictly lower than the current leader's
8//!    (lower PIDs preempt higher; the lowest live PID always wins), OR
9//! 3. The current leader's heartbeat has gone stale (last beat is
10//!    more than `grace_epochs` behind the global epoch).
11//!
12//! Election term increments on each handover so processes can
13//! observe leadership changes by polling the term.
14//!
15//! # Why lowest-live-PID
16//!
17//! It's the simplest deterministic election that always converges:
18//! given a set of live processes, exactly one (the lowest PID)
19//! deserves leadership. PIDs are unique within a host. No quorum
20//! needed; no two-round Paxos; no Raft term advance vote, just a
21//! CAS protocol that anyone reading the same MMF agrees on.
22//!
23//! # Layout
24//!
25//! ```text
26//! +-----------------------------+
27//! | LeaderHeader (64B)          |
28//! |   - magic                   |
29//! |   - current_leader_pid: u32 |
30//! |   - election_term: u32      |
31//! |   - leader_heartbeat: u64   |  (global epoch at last beat)
32//! |   - global_epoch: u64       |  (monotonic; bumped by leader's scan)
33//! +-----------------------------+
34//! ```
35
36use std::fs::{File, OpenOptions};
37use std::path::Path;
38use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
39
40use memmap2::{MmapMut, MmapOptions};
41
42pub const LEADER_MAGIC: u64 = 0x4150_4D46_4C44_5253;
43
44pub const DEFAULT_GRACE_EPOCHS: u64 = 3;
45
46/// Reserved PID value meaning "no leader claimed".
47pub const NO_LEADER: u32 = 0;
48
49#[repr(C, align(64))]
50pub struct LeaderHeader {
51    pub magic: u64,
52    pub current_leader_pid: AtomicU32,
53    pub election_term: AtomicU32,
54    pub leader_heartbeat: AtomicU64,
55    pub global_epoch: AtomicU64,
56    _pad: [u8; 32],
57}
58
59pub const LEADER_FILE_SIZE: usize = std::mem::size_of::<LeaderHeader>();
60
61#[derive(Debug, Clone, Copy, PartialEq, Eq)]
62pub enum LeaderError {
63    LayoutMismatch,
64    IoError(std::io::ErrorKind),
65}
66
67impl From<std::io::Error> for LeaderError {
68    fn from(e: std::io::Error) -> Self { Self::IoError(e.kind()) }
69}
70
71pub struct SharedLeaderElection {
72    _file: File,
73    mmap: MmapMut,
74    header_sidecar: subetha_core::HandshakeHeader,
75    ring_sidecar: Box<subetha_core::ObservationRing>,
76}
77
78unsafe impl Send for SharedLeaderElection {}
79unsafe impl Sync for SharedLeaderElection {}
80
81impl subetha_sidecar::AdaptiveInstance for SharedLeaderElection {
82    fn header(&self) -> &subetha_core::HandshakeHeader { &self.header_sidecar }
83    fn ring(&self) -> &subetha_core::ObservationRing { &self.ring_sidecar }
84    fn make_policy(&self) -> Box<dyn subetha_sidecar::Policy> {
85        Box::new(subetha_sidecar::NoMigrationPolicy)
86    }
87}
88
89impl SharedLeaderElection {
90    /// Obtain the election region at `path`, initializing it with no
91    /// leader if the path does not yet exist and attaching to it if it
92    /// does. Attaching leaves the current leader, term and epoch in
93    /// place. [`reset`](Self::reset) reinitializes.
94    pub fn create(path: impl AsRef<Path>) -> Result<Self, LeaderError> {
95        let (file, mmap) = crate::mmf_attach::create_or_attach(
96            path.as_ref(),
97            LEADER_FILE_SIZE,
98            |ptr| unsafe { Self::init_region(ptr) },
99            |ptr| unsafe { (*(ptr as *const LeaderHeader)).magic == LEADER_MAGIC },
100        )?;
101        Ok(Self {
102            _file: file, mmap,
103            header_sidecar: subetha_core::HandshakeHeader::new(),
104            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
105        })
106    }
107
108    /// Truncate the region at `path` and initialize a leaderless one,
109    /// deposing whatever leader a live peer holds. For a caller that
110    /// knows it owns the path.
111    pub fn reset(path: impl AsRef<Path>) -> Result<Self, LeaderError> {
112        let (file, mmap) =
113            crate::mmf_attach::reset(path.as_ref(), LEADER_FILE_SIZE, |ptr| unsafe {
114                Self::init_region(ptr)
115            })?;
116        Ok(Self {
117            _file: file, mmap,
118            header_sidecar: subetha_core::HandshakeHeader::new(),
119            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
120        })
121    }
122
123    /// Lay out a leaderless region: the zeroed region is already
124    /// `NO_LEADER`, term 0, heartbeat 0 and epoch 0, so only the magic
125    /// is written, last, because attachers spin on it.
126    ///
127    /// # Safety
128    /// `ptr` addresses at least [`LEADER_FILE_SIZE`] writable zeroed
129    /// bytes.
130    unsafe fn init_region(ptr: *mut u8) {
131        let hdr = ptr as *mut LeaderHeader;
132        unsafe {
133            std::ptr::write_volatile(&raw mut (*hdr).magic, LEADER_MAGIC);
134        }
135    }
136
137    pub fn open(path: impl AsRef<Path>) -> Result<Self, LeaderError> {
138        let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
139        if file.metadata()?.len() < LEADER_FILE_SIZE as u64 {
140            return Err(LeaderError::LayoutMismatch);
141        }
142        let mmap = unsafe { MmapOptions::new().len(LEADER_FILE_SIZE).map_mut(&file)? };
143        let hdr = unsafe { &*(mmap.as_ptr() as *const LeaderHeader) };
144        if hdr.magic != LEADER_MAGIC {
145            return Err(LeaderError::LayoutMismatch);
146        }
147        Ok(Self {
148            _file: file, mmap,
149            header_sidecar: subetha_core::HandshakeHeader::new(),
150            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
151        })
152    }
153
154    pub fn header(&self) -> &LeaderHeader {
155        unsafe { &*(self.mmap.as_ptr() as *const LeaderHeader) }
156    }
157
158    /// Attempt to claim leadership for `my_pid`. Returns `true` if
159    /// successful (now leader) or already leader; `false` if the
160    /// current leader is alive AND has a lower or equal PID.
161    ///
162    /// `grace_epochs` is the staleness window: when the current
163    /// leader's heartbeat is more than this many epochs behind the
164    /// global epoch, the leader is presumed dead and any process
165    /// can claim.
166    pub fn try_claim_leadership(&self, my_pid: u32, grace_epochs: u64) -> bool {
167        assert!(my_pid != NO_LEADER, "PID 0 is reserved for NO_LEADER sentinel");
168        let header = self.header();
169        loop {
170            let cur_pid = header.current_leader_pid.load(Ordering::Acquire);
171            let can_claim = if cur_pid == NO_LEADER {
172                true
173            } else if cur_pid == my_pid {
174                self.ring_sidecar
175                    .push_op(crate::sidecar_ops::ownership::OP_CLAIM, 0);
176                return true;  // already leader
177            } else if my_pid < cur_pid {
178                true  // lower PID preempts
179            } else {
180                // Higher PID: only claim if leader is stale.
181                let beat = header.leader_heartbeat.load(Ordering::Acquire);
182                let global = header.global_epoch.load(Ordering::Acquire);
183                global.saturating_sub(beat) > grace_epochs
184            };
185            if !can_claim {
186                self.ring_sidecar
187                    .push_op(crate::sidecar_ops::ownership::OP_CLAIM, 1);
188                return false;
189            }
190            if header.current_leader_pid.compare_exchange(
191                cur_pid, my_pid, Ordering::AcqRel, Ordering::Acquire,
192            ).is_ok() {
193                header.election_term.fetch_add(1, Ordering::AcqRel);
194                let global = header.global_epoch.load(Ordering::Acquire);
195                header.leader_heartbeat.store(global, Ordering::Release);
196                self.ring_sidecar
197                    .push_op(crate::sidecar_ops::ownership::OP_CLAIM, 0);
198                return true;
199            }
200            std::hint::spin_loop();
201        }
202    }
203
204    /// Heartbeat as the current leader. Updates the heartbeat to
205    /// the current global epoch. Returns `true` if the caller is
206    /// still leader (so the heartbeat counts), `false` if another
207    /// process has taken over (caller is no longer leader).
208    pub fn beat_as_leader(&self, my_pid: u32) -> bool {
209        let header = self.header();
210        if header.current_leader_pid.load(Ordering::Acquire) != my_pid {
211            self.ring_sidecar
212                .push_op(crate::sidecar_ops::ownership::OP_BEAT, 1);
213            return false;
214        }
215        let global = header.global_epoch.load(Ordering::Acquire);
216        header.leader_heartbeat.store(global, Ordering::Release);
217        self.ring_sidecar
218            .push_op(crate::sidecar_ops::ownership::OP_BEAT, 0);
219        true
220    }
221
222    /// Advance the global epoch by 1 and return the new value.
223    /// Typically called by the leader once per scan tick.
224    pub fn tick_epoch(&self) -> u64 {
225        self.header().global_epoch.fetch_add(1, Ordering::AcqRel) + 1
226    }
227
228    /// Current global epoch value.
229    pub fn global_epoch(&self) -> u64 {
230        self.header().global_epoch.load(Ordering::Acquire)
231    }
232
233    /// Current leader PID, or `None` when there is no leader.
234    pub fn current_leader(&self) -> Option<u32> {
235        let pid = self.header().current_leader_pid.load(Ordering::Acquire);
236        if pid == NO_LEADER { None } else { Some(pid) }
237    }
238
239    /// Convenience: is `my_pid` the current leader?
240    pub fn am_i_leader(&self, my_pid: u32) -> bool {
241        self.header().current_leader_pid.load(Ordering::Acquire) == my_pid
242    }
243
244    /// Current election term. Increments on each leadership change.
245    /// Subscribers can poll this to detect handovers.
246    pub fn election_term(&self) -> u32 {
247        self.header().election_term.load(Ordering::Acquire)
248    }
249
250    /// Voluntarily release leadership. Returns `true` if the caller
251    /// was the leader at the moment of release.
252    pub fn step_down(&self, my_pid: u32) -> bool {
253        let ok = self.header().current_leader_pid
254            .compare_exchange(my_pid, NO_LEADER, Ordering::AcqRel, Ordering::Acquire)
255            .is_ok();
256        self.ring_sidecar.push_op(
257            crate::sidecar_ops::ownership::OP_RELEASE,
258            if ok { 0 } else { 1 },
259        );
260        ok
261    }
262
263    pub fn flush(&self) -> Result<(), LeaderError> {
264        self.mmap.flush()?;
265        Ok(())
266    }
267
268    /// Non-blocking flush: schedules a writeback via the OS.
269    /// Note: Windows is only partially async (sync to page cache,
270    /// not to disk).
271    pub fn flush_async(&self) -> Result<(), LeaderError> {
272        self.mmap.flush_async()?;
273        Ok(())
274    }
275}
276
277#[cfg(test)]
278mod tests {
279    use super::*;
280
281    fn tmp(name: &str) -> std::path::PathBuf {
282        let mut p = std::env::temp_dir();
283        let pid = std::process::id();
284        p.push(format!("subetha-leader-{name}-{pid}.bin"));
285        p
286    }
287
288    #[test]
289    fn empty_election_first_claimer_wins() {
290        let p = tmp("empty");
291        let e = SharedLeaderElection::create(&p).unwrap();
292        assert_eq!(e.current_leader(), None);
293        assert!(e.try_claim_leadership(42, 3));
294        assert_eq!(e.current_leader(), Some(42));
295        assert!(e.am_i_leader(42));
296        assert!(!e.am_i_leader(99));
297        std::fs::remove_file(&p).ok();
298    }
299
300    /// A second create attaches with the sitting leader in place;
301    /// reset is what deposes.
302    #[test]
303    fn second_create_attaches_and_keeps_the_leader() {
304        let p = tmp("attach");
305        std::fs::remove_file(&p).ok();
306        let e = SharedLeaderElection::create(&p).unwrap();
307        assert!(e.try_claim_leadership(42, 3));
308
309        let e2 = SharedLeaderElection::create(&p).unwrap();
310        assert_eq!(e2.current_leader(), Some(42), "attach deposed the leader");
311
312        // Windows refuses to truncate a mapped file, so every handle goes
313        // before the reset.
314        drop(e);
315        drop(e2);
316        let fresh = SharedLeaderElection::reset(&p).unwrap();
317        assert_eq!(fresh.current_leader(), None, "reset left a leader seated");
318        drop(fresh);
319        std::fs::remove_file(&p).ok();
320    }
321
322    #[test]
323    fn lower_pid_preempts_higher() {
324        let p = tmp("preempt");
325        let e = SharedLeaderElection::create(&p).unwrap();
326        assert!(e.try_claim_leadership(500, 3));
327        assert!(e.am_i_leader(500));
328        // Lower PID claims; preempts.
329        assert!(e.try_claim_leadership(100, 3));
330        assert_eq!(e.current_leader(), Some(100));
331        // Higher PID cannot reclaim while 100 is alive (in the same epoch).
332        assert!(!e.try_claim_leadership(500, 3));
333        assert_eq!(e.current_leader(), Some(100));
334        std::fs::remove_file(&p).ok();
335    }
336
337    #[test]
338    fn equal_pid_returns_true_idempotent() {
339        let p = tmp("equal-pid");
340        let e = SharedLeaderElection::create(&p).unwrap();
341        assert!(e.try_claim_leadership(42, 3));
342        // Same PID claiming again is a no-op success.
343        assert!(e.try_claim_leadership(42, 3));
344        std::fs::remove_file(&p).ok();
345    }
346
347    #[test]
348    fn stale_leader_replaced_after_grace_window() {
349        let p = tmp("stale");
350        let e = SharedLeaderElection::create(&p).unwrap();
351        assert!(e.try_claim_leadership(100, 1));
352        // Tick global epoch beyond grace without beating.
353        e.tick_epoch();
354        e.tick_epoch();
355        // Higher PID can now claim because heartbeat is stale.
356        assert!(e.try_claim_leadership(500, 1));
357        assert_eq!(e.current_leader(), Some(500));
358        std::fs::remove_file(&p).ok();
359    }
360
361    #[test]
362    fn beat_keeps_leader_alive() {
363        let p = tmp("beat");
364        let e = SharedLeaderElection::create(&p).unwrap();
365        assert!(e.try_claim_leadership(100, 1));
366        e.tick_epoch();
367        assert!(e.beat_as_leader(100));
368        e.tick_epoch();
369        assert!(e.beat_as_leader(100));
370        // Higher PID still can't preempt (heartbeat fresh).
371        assert!(!e.try_claim_leadership(500, 1));
372        std::fs::remove_file(&p).ok();
373    }
374
375    #[test]
376    fn election_term_increments_on_each_handover() {
377        let p = tmp("term");
378        let e = SharedLeaderElection::create(&p).unwrap();
379        let t0 = e.election_term();
380        assert!(e.try_claim_leadership(500, 3));
381        let t1 = e.election_term();
382        assert_eq!(t1, t0 + 1);
383        // Preemption by lower PID is a handover; term advances.
384        assert!(e.try_claim_leadership(100, 3));
385        let t2 = e.election_term();
386        assert_eq!(t2, t1 + 1);
387        // Same-PID re-claim is a no-op; term stays.
388        assert!(e.try_claim_leadership(100, 3));
389        assert_eq!(e.election_term(), t2);
390        std::fs::remove_file(&p).ok();
391    }
392
393    #[test]
394    fn step_down_clears_leadership() {
395        let p = tmp("step-down");
396        let e = SharedLeaderElection::create(&p).unwrap();
397        e.try_claim_leadership(42, 3);
398        assert!(e.step_down(42));
399        assert_eq!(e.current_leader(), None);
400        // Non-leader step_down is a no-op.
401        assert!(!e.step_down(99));
402        std::fs::remove_file(&p).ok();
403    }
404
405    #[test]
406    fn cross_handle_visibility() {
407        let p = tmp("cross-handle");
408        let e_a = SharedLeaderElection::create(&p).unwrap();
409        let e_b = SharedLeaderElection::open(&p).unwrap();
410        assert!(e_a.try_claim_leadership(100, 3));
411        assert_eq!(e_b.current_leader(), Some(100));
412        assert!(e_b.am_i_leader(100));
413        assert!(!e_b.am_i_leader(200));
414        std::fs::remove_file(&p).ok();
415    }
416
417    #[test]
418    fn beat_returns_false_after_preemption() {
419        let p = tmp("beat-after-preempt");
420        let e = SharedLeaderElection::create(&p).unwrap();
421        e.try_claim_leadership(500, 3);
422        // Lower PID preempts.
423        e.try_claim_leadership(100, 3);
424        // 500's beat now returns false.
425        assert!(!e.beat_as_leader(500));
426        assert!(e.beat_as_leader(100));
427        std::fs::remove_file(&p).ok();
428    }
429
430    #[test]
431    fn disk_persistence_survives_reopen() {
432        let p = tmp("disk-persist");
433        {
434            let e = SharedLeaderElection::create(&p).unwrap();
435            e.try_claim_leadership(42, 3);
436            e.flush().unwrap();
437        }
438        let e2 = SharedLeaderElection::open(&p).unwrap();
439        assert_eq!(e2.current_leader(), Some(42));
440        std::fs::remove_file(&p).ok();
441    }
442}