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    pub fn create(path: impl AsRef<Path>) -> Result<Self, LeaderError> {
91        let file = OpenOptions::new()
92            .read(true).write(true).create(true).truncate(true)
93            .open(path.as_ref())?;
94        file.set_len(LEADER_FILE_SIZE as u64)?;
95        let mut mmap = unsafe { MmapOptions::new().len(LEADER_FILE_SIZE).map_mut(&file)? };
96        let hdr = mmap.as_mut_ptr() as *mut LeaderHeader;
97        unsafe {
98            std::ptr::write(hdr, LeaderHeader {
99                magic: LEADER_MAGIC,
100                current_leader_pid: AtomicU32::new(NO_LEADER),
101                election_term: AtomicU32::new(0),
102                leader_heartbeat: AtomicU64::new(0),
103                global_epoch: AtomicU64::new(0),
104                _pad: [0; 32],
105            });
106        }
107        Ok(Self {
108            _file: file, mmap,
109            header_sidecar: subetha_core::HandshakeHeader::new(),
110            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
111        })
112    }
113
114    pub fn open(path: impl AsRef<Path>) -> Result<Self, LeaderError> {
115        let file = OpenOptions::new().read(true).write(true).open(path.as_ref())?;
116        if file.metadata()?.len() < LEADER_FILE_SIZE as u64 {
117            return Err(LeaderError::LayoutMismatch);
118        }
119        let mmap = unsafe { MmapOptions::new().len(LEADER_FILE_SIZE).map_mut(&file)? };
120        let hdr = unsafe { &*(mmap.as_ptr() as *const LeaderHeader) };
121        if hdr.magic != LEADER_MAGIC {
122            return Err(LeaderError::LayoutMismatch);
123        }
124        Ok(Self {
125            _file: file, mmap,
126            header_sidecar: subetha_core::HandshakeHeader::new(),
127            ring_sidecar: Box::new(subetha_core::ObservationRing::new()),
128        })
129    }
130
131    pub fn header(&self) -> &LeaderHeader {
132        unsafe { &*(self.mmap.as_ptr() as *const LeaderHeader) }
133    }
134
135    /// Attempt to claim leadership for `my_pid`. Returns `true` if
136    /// successful (now leader) or already leader; `false` if the
137    /// current leader is alive AND has a lower or equal PID.
138    ///
139    /// `grace_epochs` is the staleness window: when the current
140    /// leader's heartbeat is more than this many epochs behind the
141    /// global epoch, the leader is presumed dead and any process
142    /// can claim.
143    pub fn try_claim_leadership(&self, my_pid: u32, grace_epochs: u64) -> bool {
144        assert!(my_pid != NO_LEADER, "PID 0 is reserved for NO_LEADER sentinel");
145        let header = self.header();
146        loop {
147            let cur_pid = header.current_leader_pid.load(Ordering::Acquire);
148            let can_claim = if cur_pid == NO_LEADER {
149                true
150            } else if cur_pid == my_pid {
151                self.ring_sidecar
152                    .push_op(crate::sidecar_ops::ownership::OP_CLAIM, 0);
153                return true;  // already leader
154            } else if my_pid < cur_pid {
155                true  // lower PID preempts
156            } else {
157                // Higher PID: only claim if leader is stale.
158                let beat = header.leader_heartbeat.load(Ordering::Acquire);
159                let global = header.global_epoch.load(Ordering::Acquire);
160                global.saturating_sub(beat) > grace_epochs
161            };
162            if !can_claim {
163                self.ring_sidecar
164                    .push_op(crate::sidecar_ops::ownership::OP_CLAIM, 1);
165                return false;
166            }
167            if header.current_leader_pid.compare_exchange(
168                cur_pid, my_pid, Ordering::AcqRel, Ordering::Acquire,
169            ).is_ok() {
170                header.election_term.fetch_add(1, Ordering::AcqRel);
171                let global = header.global_epoch.load(Ordering::Acquire);
172                header.leader_heartbeat.store(global, Ordering::Release);
173                self.ring_sidecar
174                    .push_op(crate::sidecar_ops::ownership::OP_CLAIM, 0);
175                return true;
176            }
177            std::hint::spin_loop();
178        }
179    }
180
181    /// Heartbeat as the current leader. Updates the heartbeat to
182    /// the current global epoch. Returns `true` if the caller is
183    /// still leader (so the heartbeat counts), `false` if another
184    /// process has taken over (caller is no longer leader).
185    pub fn beat_as_leader(&self, my_pid: u32) -> bool {
186        let header = self.header();
187        if header.current_leader_pid.load(Ordering::Acquire) != my_pid {
188            self.ring_sidecar
189                .push_op(crate::sidecar_ops::ownership::OP_BEAT, 1);
190            return false;
191        }
192        let global = header.global_epoch.load(Ordering::Acquire);
193        header.leader_heartbeat.store(global, Ordering::Release);
194        self.ring_sidecar
195            .push_op(crate::sidecar_ops::ownership::OP_BEAT, 0);
196        true
197    }
198
199    /// Advance the global epoch by 1 and return the new value.
200    /// Typically called by the leader once per scan tick.
201    pub fn tick_epoch(&self) -> u64 {
202        self.header().global_epoch.fetch_add(1, Ordering::AcqRel) + 1
203    }
204
205    /// Current global epoch value.
206    pub fn global_epoch(&self) -> u64 {
207        self.header().global_epoch.load(Ordering::Acquire)
208    }
209
210    /// Current leader PID, or `None` when there is no leader.
211    pub fn current_leader(&self) -> Option<u32> {
212        let pid = self.header().current_leader_pid.load(Ordering::Acquire);
213        if pid == NO_LEADER { None } else { Some(pid) }
214    }
215
216    /// Convenience: is `my_pid` the current leader?
217    pub fn am_i_leader(&self, my_pid: u32) -> bool {
218        self.header().current_leader_pid.load(Ordering::Acquire) == my_pid
219    }
220
221    /// Current election term. Increments on each leadership change.
222    /// Subscribers can poll this to detect handovers.
223    pub fn election_term(&self) -> u32 {
224        self.header().election_term.load(Ordering::Acquire)
225    }
226
227    /// Voluntarily release leadership. Returns `true` if the caller
228    /// was the leader at the moment of release.
229    pub fn step_down(&self, my_pid: u32) -> bool {
230        let ok = self.header().current_leader_pid
231            .compare_exchange(my_pid, NO_LEADER, Ordering::AcqRel, Ordering::Acquire)
232            .is_ok();
233        self.ring_sidecar.push_op(
234            crate::sidecar_ops::ownership::OP_RELEASE,
235            if ok { 0 } else { 1 },
236        );
237        ok
238    }
239
240    pub fn flush(&self) -> Result<(), LeaderError> {
241        self.mmap.flush()?;
242        Ok(())
243    }
244
245    /// Non-blocking flush: schedules a writeback via the OS.
246    /// Note: Windows is only partially async (sync to page cache,
247    /// not to disk).
248    pub fn flush_async(&self) -> Result<(), LeaderError> {
249        self.mmap.flush_async()?;
250        Ok(())
251    }
252}
253
254#[cfg(test)]
255mod tests {
256    use super::*;
257
258    fn tmp(name: &str) -> std::path::PathBuf {
259        let mut p = std::env::temp_dir();
260        let pid = std::process::id();
261        p.push(format!("subetha-leader-{name}-{pid}.bin"));
262        p
263    }
264
265    #[test]
266    fn empty_election_first_claimer_wins() {
267        let p = tmp("empty");
268        let e = SharedLeaderElection::create(&p).unwrap();
269        assert_eq!(e.current_leader(), None);
270        assert!(e.try_claim_leadership(42, 3));
271        assert_eq!(e.current_leader(), Some(42));
272        assert!(e.am_i_leader(42));
273        assert!(!e.am_i_leader(99));
274        std::fs::remove_file(&p).ok();
275    }
276
277    #[test]
278    fn lower_pid_preempts_higher() {
279        let p = tmp("preempt");
280        let e = SharedLeaderElection::create(&p).unwrap();
281        assert!(e.try_claim_leadership(500, 3));
282        assert!(e.am_i_leader(500));
283        // Lower PID claims; preempts.
284        assert!(e.try_claim_leadership(100, 3));
285        assert_eq!(e.current_leader(), Some(100));
286        // Higher PID cannot reclaim while 100 is alive (in the same epoch).
287        assert!(!e.try_claim_leadership(500, 3));
288        assert_eq!(e.current_leader(), Some(100));
289        std::fs::remove_file(&p).ok();
290    }
291
292    #[test]
293    fn equal_pid_returns_true_idempotent() {
294        let p = tmp("equal-pid");
295        let e = SharedLeaderElection::create(&p).unwrap();
296        assert!(e.try_claim_leadership(42, 3));
297        // Same PID claiming again is a no-op success.
298        assert!(e.try_claim_leadership(42, 3));
299        std::fs::remove_file(&p).ok();
300    }
301
302    #[test]
303    fn stale_leader_replaced_after_grace_window() {
304        let p = tmp("stale");
305        let e = SharedLeaderElection::create(&p).unwrap();
306        assert!(e.try_claim_leadership(100, 1));
307        // Tick global epoch beyond grace without beating.
308        e.tick_epoch();
309        e.tick_epoch();
310        // Higher PID can now claim because heartbeat is stale.
311        assert!(e.try_claim_leadership(500, 1));
312        assert_eq!(e.current_leader(), Some(500));
313        std::fs::remove_file(&p).ok();
314    }
315
316    #[test]
317    fn beat_keeps_leader_alive() {
318        let p = tmp("beat");
319        let e = SharedLeaderElection::create(&p).unwrap();
320        assert!(e.try_claim_leadership(100, 1));
321        e.tick_epoch();
322        assert!(e.beat_as_leader(100));
323        e.tick_epoch();
324        assert!(e.beat_as_leader(100));
325        // Higher PID still can't preempt (heartbeat fresh).
326        assert!(!e.try_claim_leadership(500, 1));
327        std::fs::remove_file(&p).ok();
328    }
329
330    #[test]
331    fn election_term_increments_on_each_handover() {
332        let p = tmp("term");
333        let e = SharedLeaderElection::create(&p).unwrap();
334        let t0 = e.election_term();
335        assert!(e.try_claim_leadership(500, 3));
336        let t1 = e.election_term();
337        assert_eq!(t1, t0 + 1);
338        // Preemption by lower PID is a handover; term advances.
339        assert!(e.try_claim_leadership(100, 3));
340        let t2 = e.election_term();
341        assert_eq!(t2, t1 + 1);
342        // Same-PID re-claim is a no-op; term stays.
343        assert!(e.try_claim_leadership(100, 3));
344        assert_eq!(e.election_term(), t2);
345        std::fs::remove_file(&p).ok();
346    }
347
348    #[test]
349    fn step_down_clears_leadership() {
350        let p = tmp("step-down");
351        let e = SharedLeaderElection::create(&p).unwrap();
352        e.try_claim_leadership(42, 3);
353        assert!(e.step_down(42));
354        assert_eq!(e.current_leader(), None);
355        // Non-leader step_down is a no-op.
356        assert!(!e.step_down(99));
357        std::fs::remove_file(&p).ok();
358    }
359
360    #[test]
361    fn cross_handle_visibility() {
362        let p = tmp("cross-handle");
363        let e_a = SharedLeaderElection::create(&p).unwrap();
364        let e_b = SharedLeaderElection::open(&p).unwrap();
365        assert!(e_a.try_claim_leadership(100, 3));
366        assert_eq!(e_b.current_leader(), Some(100));
367        assert!(e_b.am_i_leader(100));
368        assert!(!e_b.am_i_leader(200));
369        std::fs::remove_file(&p).ok();
370    }
371
372    #[test]
373    fn beat_returns_false_after_preemption() {
374        let p = tmp("beat-after-preempt");
375        let e = SharedLeaderElection::create(&p).unwrap();
376        e.try_claim_leadership(500, 3);
377        // Lower PID preempts.
378        e.try_claim_leadership(100, 3);
379        // 500's beat now returns false.
380        assert!(!e.beat_as_leader(500));
381        assert!(e.beat_as_leader(100));
382        std::fs::remove_file(&p).ok();
383    }
384
385    #[test]
386    fn disk_persistence_survives_reopen() {
387        let p = tmp("disk-persist");
388        {
389            let e = SharedLeaderElection::create(&p).unwrap();
390            e.try_claim_leadership(42, 3);
391            e.flush().unwrap();
392        }
393        let e2 = SharedLeaderElection::open(&p).unwrap();
394        assert_eq!(e2.current_leader(), Some(42));
395        std::fs::remove_file(&p).ok();
396    }
397}