1use 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
46pub 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 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; } else if my_pid < cur_pid {
155 true } else {
157 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 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 pub fn tick_epoch(&self) -> u64 {
202 self.header().global_epoch.fetch_add(1, Ordering::AcqRel) + 1
203 }
204
205 pub fn global_epoch(&self) -> u64 {
207 self.header().global_epoch.load(Ordering::Acquire)
208 }
209
210 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 pub fn am_i_leader(&self, my_pid: u32) -> bool {
218 self.header().current_leader_pid.load(Ordering::Acquire) == my_pid
219 }
220
221 pub fn election_term(&self) -> u32 {
224 self.header().election_term.load(Ordering::Acquire)
225 }
226
227 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 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 assert!(e.try_claim_leadership(100, 3));
285 assert_eq!(e.current_leader(), Some(100));
286 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 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 e.tick_epoch();
309 e.tick_epoch();
310 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 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 assert!(e.try_claim_leadership(100, 3));
340 let t2 = e.election_term();
341 assert_eq!(t2, t1 + 1);
342 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 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 e.try_claim_leadership(100, 3);
379 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}