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> {
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 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 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 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; } else if my_pid < cur_pid {
178 true } else {
180 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 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 pub fn tick_epoch(&self) -> u64 {
225 self.header().global_epoch.fetch_add(1, Ordering::AcqRel) + 1
226 }
227
228 pub fn global_epoch(&self) -> u64 {
230 self.header().global_epoch.load(Ordering::Acquire)
231 }
232
233 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 pub fn am_i_leader(&self, my_pid: u32) -> bool {
241 self.header().current_leader_pid.load(Ordering::Acquire) == my_pid
242 }
243
244 pub fn election_term(&self) -> u32 {
247 self.header().election_term.load(Ordering::Acquire)
248 }
249
250 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 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 #[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 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 assert!(e.try_claim_leadership(100, 3));
330 assert_eq!(e.current_leader(), Some(100));
331 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 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 e.tick_epoch();
354 e.tick_epoch();
355 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 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 assert!(e.try_claim_leadership(100, 3));
385 let t2 = e.election_term();
386 assert_eq!(t2, t1 + 1);
387 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 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 e.try_claim_leadership(100, 3);
424 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}