1use std::collections::BTreeSet;
9use std::time::Duration;
10
11use crate::replication::{VoteDecision, VoteRequest};
12
13#[derive(Debug, Clone, PartialEq, Eq)]
14pub enum ReplicationControlMessage {
15 ElectionVoteRequest(VoteRequest),
16 ElectionVoteDecision(VoteDecision),
17 LogicalCommit {
18 term: u64,
19 lsn: u64,
20 payload_hash: String,
21 },
22 LeaseProbe {
23 holder_id: String,
24 term: u64,
25 },
26}
27
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub struct NetworkFaults {
30 pub loss_per_million: u32,
31 pub max_delay_ms: u64,
32 pub reorder: bool,
33}
34
35impl NetworkFaults {
36 pub fn reliable() -> Self {
37 Self {
38 loss_per_million: 0,
39 max_delay_ms: 0,
40 reorder: false,
41 }
42 }
43
44 pub fn lossy(loss_per_million: u32) -> Self {
45 Self {
46 loss_per_million,
47 max_delay_ms: 0,
48 reorder: false,
49 }
50 }
51}
52
53#[derive(Debug, Clone, Copy, PartialEq, Eq)]
54pub struct SimulationClock {
55 now_ms: u64,
56}
57
58impl SimulationClock {
59 pub fn new() -> Self {
60 Self { now_ms: 0 }
61 }
62
63 pub fn now_ms(&self) -> u64 {
64 self.now_ms
65 }
66
67 pub fn elapsed(&self) -> Duration {
68 Duration::from_millis(self.now_ms)
69 }
70
71 pub fn advance(&mut self, by: Duration) {
72 let millis = u64::try_from(by.as_millis()).unwrap_or(u64::MAX);
73 self.now_ms = self.now_ms.saturating_add(millis);
74 }
75
76 fn advance_to(&mut self, now_ms: u64) {
77 self.now_ms = self.now_ms.max(now_ms);
78 }
79}
80
81impl Default for SimulationClock {
82 fn default() -> Self {
83 Self::new()
84 }
85}
86
87#[derive(Debug, Clone, PartialEq, Eq)]
88pub struct Delivered<M> {
89 pub from: String,
90 pub to: String,
91 pub message: M,
92 pub delivered_at_ms: u64,
93}
94
95#[derive(Debug, Clone, PartialEq, Eq)]
96pub enum SendOutcome {
97 Accepted { deliver_at_ms: u64 },
98 Dropped(DropReason),
99}
100
101#[derive(Debug, Clone, Copy, PartialEq, Eq)]
102pub enum DropReason {
103 Partition,
104 Loss,
105}
106
107#[derive(Debug, Clone)]
108struct Pending<M> {
109 from: String,
110 to: String,
111 message: M,
112 deliver_at_ms: u64,
113 order: u64,
114}
115
116#[derive(Debug, Clone)]
117pub struct InProcessReplicationNetwork<M> {
118 clock: SimulationClock,
119 rng: SeededRng,
120 faults: NetworkFaults,
121 partitions: BTreeSet<(String, String)>,
122 pending: Vec<Pending<M>>,
123 sequence: u64,
124}
125
126impl<M> InProcessReplicationNetwork<M> {
127 pub fn new(seed: u64, faults: NetworkFaults) -> Self {
128 Self {
129 clock: SimulationClock::new(),
130 rng: SeededRng::new(seed),
131 faults,
132 partitions: BTreeSet::new(),
133 pending: Vec::new(),
134 sequence: 0,
135 }
136 }
137
138 pub fn clock(&self) -> SimulationClock {
139 self.clock
140 }
141
142 pub fn advance(&mut self, by: Duration) {
143 self.clock.advance(by);
144 }
145
146 pub fn partition(&mut self, a: impl Into<String>, b: impl Into<String>) {
147 self.partitions.insert(partition_key(a.into(), b.into()));
148 }
149
150 pub fn heal(&mut self, a: &str, b: &str) {
151 self.partitions
152 .remove(&partition_key(a.to_string(), b.to_string()));
153 }
154
155 pub fn send(
156 &mut self,
157 from: impl Into<String>,
158 to: impl Into<String>,
159 message: M,
160 ) -> SendOutcome {
161 let from = from.into();
162 let to = to.into();
163 if self
164 .partitions
165 .contains(&partition_key(from.clone(), to.clone()))
166 {
167 return SendOutcome::Dropped(DropReason::Partition);
168 }
169 if self.faults.loss_per_million > 0
170 && self.rng.next_bounded(1_000_000) < u64::from(self.faults.loss_per_million)
171 {
172 return SendOutcome::Dropped(DropReason::Loss);
173 }
174
175 self.sequence = self.sequence.saturating_add(1);
176 let delay = if self.faults.max_delay_ms == 0 {
177 0
178 } else {
179 self.rng
180 .next_bounded(self.faults.max_delay_ms.saturating_add(1))
181 };
182 let deliver_at_ms = self.clock.now_ms().saturating_add(delay);
183 let order = if self.faults.reorder {
184 self.rng.next_u64()
185 } else {
186 self.sequence
187 };
188 self.pending.push(Pending {
189 from,
190 to,
191 message,
192 deliver_at_ms,
193 order,
194 });
195 SendOutcome::Accepted { deliver_at_ms }
196 }
197
198 pub fn advance_to_next_delivery(&mut self) -> bool {
199 let Some(next) = self.pending.iter().map(|p| p.deliver_at_ms).min() else {
200 return false;
201 };
202 self.clock.advance_to(next);
203 true
204 }
205
206 pub fn drain_ready_for(&mut self, recipient: &str) -> Vec<Delivered<M>> {
207 let now = self.clock.now_ms();
208 let mut ready = Vec::new();
209 let mut pending = Vec::with_capacity(self.pending.len());
210 for msg in self.pending.drain(..) {
211 if msg.to == recipient && msg.deliver_at_ms <= now {
212 ready.push(msg);
213 } else {
214 pending.push(msg);
215 }
216 }
217 self.pending = pending;
218 ready.sort_by_key(|msg| (msg.deliver_at_ms, msg.order));
219 ready
220 .into_iter()
221 .map(|msg| Delivered {
222 from: msg.from,
223 to: msg.to,
224 message: msg.message,
225 delivered_at_ms: msg.deliver_at_ms,
226 })
227 .collect()
228 }
229}
230
231#[derive(Debug, Clone)]
232struct SeededRng {
233 state: u64,
234}
235
236impl SeededRng {
237 fn new(seed: u64) -> Self {
238 let state = if seed == 0 {
239 0x9E37_79B9_7F4A_7C15
240 } else {
241 seed
242 };
243 Self { state }
244 }
245
246 fn next_u64(&mut self) -> u64 {
247 let mut x = self.state;
248 x ^= x << 13;
249 x ^= x >> 7;
250 x ^= x << 17;
251 self.state = x;
252 x
253 }
254
255 fn next_bounded(&mut self, upper_exclusive: u64) -> u64 {
256 if upper_exclusive == 0 {
257 0
258 } else {
259 self.next_u64() % upper_exclusive
260 }
261 }
262}
263
264fn partition_key(a: String, b: String) -> (String, String) {
265 if a <= b {
266 (a, b)
267 } else {
268 (b, a)
269 }
270}
271
272#[cfg(test)]
273mod tests {
274 use std::collections::{BTreeMap, BTreeSet};
275 use std::rc::Rc;
276 use std::sync::Arc;
277
278 use super::*;
279 use crate::replication::{
280 ElectionCoordinator, ElectionOutcome, ElectionRequest, ElectionTransport, LastVote,
281 LastVoteError, LastVoteStore, LeaseError, LeaseStore, Member, MemoryLastVoteStore,
282 RefusalReason, Voter, WriterLease,
283 };
284
285 #[test]
286 fn fault_injection_is_seed_reproducible() {
287 let trace_a = delivery_trace(0xD57, 0);
288 let trace_b = delivery_trace(0xD57, 0);
289 let trace_c = delivery_trace(0xD58, 0);
290
291 assert_eq!(trace_a, trace_b, "same seed must reproduce the trace");
292 assert_ne!(
293 trace_a, trace_c,
294 "different seed should explore a different trace"
295 );
296
297 let mut partitioned = InProcessReplicationNetwork::new(1, NetworkFaults::reliable());
298 partitioned.partition("a", "b");
299 assert_eq!(
300 partitioned.send("a", "b", 1u64),
301 SendOutcome::Dropped(DropReason::Partition)
302 );
303
304 let mut lossy = InProcessReplicationNetwork::new(1, NetworkFaults::lossy(1_000_000));
305 assert_eq!(
306 lossy.send("a", "b", 1u64),
307 SendOutcome::Dropped(DropReason::Loss)
308 );
309 }
310
311 fn delivery_trace(seed: u64, loss_per_million: u32) -> Vec<(u64, u64)> {
312 let faults = NetworkFaults {
313 loss_per_million,
314 max_delay_ms: 25,
315 reorder: true,
316 };
317 let mut network = InProcessReplicationNetwork::new(seed, faults);
318 for value in 0..12u64 {
319 let _ = network.send("a", "b", value);
320 }
321 network.advance(Duration::from_millis(25));
322 network
323 .drain_ready_for("b")
324 .into_iter()
325 .map(|msg| (msg.delivered_at_ms, msg.message))
326 .collect()
327 }
328
329 #[test]
330 fn election_safety_under_partition_has_at_most_one_leader_per_term() {
331 let members = five_voters();
332 let stores = shared_vote_stores(&members);
333 let mut network = partitioned_network(0x1358);
334 for peer in ["d", "e"] {
335 network.partition("a", peer);
336 network.partition("b", peer);
337 network.partition("c", peer);
338 }
339
340 let mut leaders = BTreeMap::new();
341 for candidate in [
342 candidate_request("a", 4, 120, 100),
343 candidate_request("d", 4, 120, 100),
344 ] {
345 let mut tx = NetworkElectionTransport::new(
346 &mut network,
347 members.clone(),
348 stores.clone(),
349 candidate.candidate.id.clone(),
350 100,
351 );
352 let outcome = ElectionCoordinator::run(&candidate, &mut tx, Duration::from_secs(60));
353 if let ElectionOutcome::Elected { term, .. } = outcome {
354 let previous = leaders.insert(term, candidate.candidate.id.clone());
355 assert_eq!(
356 previous, None,
357 "two leaders elected in term {term}: {previous:?} and {:?}",
358 candidate.candidate.id
359 );
360 }
361 }
362
363 assert_eq!(leaders.get(&5), Some(&"a".to_string()));
364 }
365
366 #[test]
367 fn partitioned_elections_do_not_split_brain_or_lose_committed_writes() {
368 let committed_watermark = 100;
369 let committed_writes: BTreeSet<u64> = (1..=committed_watermark).collect();
370
371 for seed in 1..=24 {
372 let members = five_voters();
373 let stores = shared_vote_stores(&members);
374 let mut network = partitioned_network(seed);
375 let mut elected = Vec::new();
376 let candidates = if seed % 2 == 0 {
377 [("a", 120), ("d", 80)]
378 } else {
379 [("d", 80), ("a", 120)]
380 };
381
382 for (id, lsn) in candidates {
383 let req = candidate_request(id, 4, lsn, committed_watermark);
384 let mut tx = NetworkElectionTransport::new(
385 &mut network,
386 members.clone(),
387 stores.clone(),
388 id.to_string(),
389 committed_watermark,
390 );
391 let outcome = ElectionCoordinator::run(&req, &mut tx, Duration::from_secs(60));
392 if let ElectionOutcome::Elected { term, .. } = outcome {
393 elected.push((term, id.to_string(), lsn));
394 }
395 }
396
397 let mut leaders_by_term = BTreeSet::new();
398 for (term, id, lsn) in elected {
399 assert!(
400 leaders_by_term.insert(term),
401 "split-brain in term {term} under seed {seed}"
402 );
403 assert!(
404 lsn >= committed_watermark,
405 "leader {id} lost committed writes under seed {seed}"
406 );
407 assert!(
408 committed_writes.iter().all(|committed| *committed <= lsn),
409 "leader {id} does not cover all committed writes under seed {seed}"
410 );
411 }
412 }
413 }
414
415 #[test]
416 fn lease_fencing_holds_when_a_partitioned_primary_returns_stale() {
417 let members = five_voters();
418 let stores = shared_vote_stores(&members);
419 let mut network = partitioned_network(0x715);
420 for peer in ["a", "b"] {
421 network.partition("old-primary", peer);
422 }
423 let promoted = candidate_request("a", 4, 150, 100);
424 let mut tx =
425 NetworkElectionTransport::new(&mut network, members, stores, "a".to_string(), 100);
426
427 let outcome = ElectionCoordinator::run(&promoted, &mut tx, Duration::from_secs(60));
428 let ElectionOutcome::Elected { term: new_term, .. } = outcome else {
429 panic!("expected a replacement primary, got {outcome:?}");
430 };
431
432 let store = lease_store("dst-fence");
433 let lease = store
434 .try_acquire_for_term("main", "new-primary", 60_000, new_term)
435 .expect("new primary lease");
436 assert_eq!(lease.term, new_term);
437
438 let err = store
439 .try_acquire_for_term("main", "old-primary", 60_000, new_term - 1)
440 .expect_err("stale partitioned primary must be fenced");
441 assert!(
442 matches!(
443 err,
444 LeaseError::Fenced {
445 current_term,
446 ..
447 } if current_term == new_term
448 ),
449 "got {err:?}"
450 );
451
452 let stale_lease = WriterLease {
453 database_key: "main".to_string(),
454 holder_id: "old-primary".to_string(),
455 term: new_term - 1,
456 generation: 1,
457 acquired_at_ms: 0,
458 expires_at_ms: u64::MAX,
459 };
460 assert!(stale_lease.fenced_by_term(new_term));
461 }
462
463 #[test]
464 fn partition_clock_skew_lease_expiry_preserves_ownership_and_sync_acks() {
465 let report = run_partition_clock_skew_lease_expiry_scenario(0x1846);
466
467 assert_eq!(report.promoted_owner, "replica-a");
468 assert_eq!(report.self_fenced_owner, "old-primary");
469 assert!(
470 report.old_owner_late_write_refused,
471 "the deposed owner must be refused by the admission gate"
472 );
473 assert!(
474 report.local_policy_losses > 0,
475 "the schedule should exercise documented local-policy loss"
476 );
477 assert_eq!(
478 report.sync_ack_loss_count, 0,
479 "synchronous acknowledgements must survive recovery"
480 );
481 }
482
483 #[test]
484 #[ignore = "heavy seed sweep runs nightly in CI"]
485 fn dst_seed_sweep_election_safety_no_split_brain_no_lost_committed_writes() {
486 for seed in 1..=256 {
487 let members = five_voters();
488 let stores = shared_vote_stores(&members);
489 let mut network = partitioned_network(seed);
490 let mut leaders = BTreeMap::new();
491 for (id, lsn) in [("a", 125), ("b", 130), ("d", 90), ("e", 95)] {
492 let req = candidate_request(id, 7, lsn, 100);
493 let mut tx = NetworkElectionTransport::new(
494 &mut network,
495 members.clone(),
496 stores.clone(),
497 id.to_string(),
498 100,
499 );
500 if let ElectionOutcome::Elected { term, .. } =
501 ElectionCoordinator::run(&req, &mut tx, Duration::from_secs(60))
502 {
503 assert!(lsn >= 100, "seed {seed}: elected {id} below watermark");
504 assert_eq!(
505 leaders.insert(term, id.to_string()),
506 None,
507 "seed {seed}: more than one leader in term {term}"
508 );
509 }
510 }
511 }
512 }
513
514 #[test]
515 #[ignore = "heavy seed sweep runs nightly in CI"]
516 fn dst_seed_sweep_partition_clock_skew_lease_expiry_oracles() {
517 for seed in 1..=128 {
518 let report = run_partition_clock_skew_lease_expiry_scenario(seed);
519 assert_eq!(
520 report.sync_ack_loss_count, 0,
521 "seed {seed}: synchronous acked writes were lost"
522 );
523 assert!(
524 report.old_owner_late_write_refused,
525 "seed {seed}: stale owner was not refused"
526 );
527 }
528 }
529
530 fn run_partition_clock_skew_lease_expiry_scenario(seed: u64) -> DstLeaseExpiryReport {
531 let mut scenario = LeaseExpiryScenario::new(seed);
532 scenario.bootstrap_sync_write();
533 scenario.partition_old_owner();
534 scenario.accept_old_owner_local_write_before_expiry();
535 scenario.advance_until_old_owner_self_fences();
536 scenario.promote_covered_replica();
537 scenario.accept_promoted_owner_sync_write();
538 scenario.refuse_old_owner_late_write();
539 scenario.assert_no_double_owner_window();
540 scenario.assert_acked_write_loss_oracle()
541 }
542
543 struct LeaseExpiryScenario {
544 network: InProcessReplicationNetwork<ReplicationControlMessage>,
545 old_owner: SimOwner,
546 current_epoch: u64,
547 committed_watermark: u64,
548 next_write_id: u64,
549 replica_logs: BTreeMap<String, BTreeSet<u64>>,
550 accepted: Vec<SimAcceptedWrite>,
551 self_fenced_at_ms: Option<u64>,
552 promoted_at_ms: Option<u64>,
553 promoted_owner: Option<String>,
554 old_owner_late_write_refused: bool,
555 }
556
557 impl LeaseExpiryScenario {
558 fn new(seed: u64) -> Self {
559 let mut replica_logs = BTreeMap::new();
560 for member in [
561 "old-primary",
562 "replica-a",
563 "replica-b",
564 "replica-c",
565 "replica-d",
566 ] {
567 replica_logs.insert(member.to_string(), BTreeSet::new());
568 }
569
570 Self {
571 network: InProcessReplicationNetwork::new(
572 seed,
573 NetworkFaults {
574 loss_per_million: 0,
575 max_delay_ms: 25,
576 reorder: true,
577 },
578 ),
579 old_owner: SimOwner {
580 id: "old-primary".to_string(),
581 term: 7,
582 epoch: 1,
583 lease_expires_local_ms: 100,
584 clock_skew_ms: -25,
585 self_fenced: false,
586 },
587 current_epoch: 1,
588 committed_watermark: 0,
589 next_write_id: 1,
590 replica_logs,
591 accepted: Vec::new(),
592 self_fenced_at_ms: None,
593 promoted_at_ms: None,
594 promoted_owner: None,
595 old_owner_late_write_refused: false,
596 }
597 }
598
599 fn bootstrap_sync_write(&mut self) {
600 let write_id = self.next_write_id();
601 self.record_owner_durable_write(
602 self.old_owner.id.clone(),
603 self.old_owner.term,
604 self.old_owner.epoch,
605 write_id,
606 AckPolicy::Synchronous,
607 );
608 let old_owner_id = self.old_owner.id.clone();
609 let replicated = self.replicate_commit(
610 &old_owner_id,
611 self.old_owner.term,
612 write_id,
613 &["replica-a", "replica-b"],
614 );
615 assert_eq!(
616 replicated.len(),
617 2,
618 "bootstrap synchronous write must reach a covered quorum"
619 );
620 self.committed_watermark = write_id;
621 }
622
623 fn partition_old_owner(&mut self) {
624 for peer in ["replica-a", "replica-b", "replica-c", "replica-d"] {
625 self.network.partition(&self.old_owner.id, peer);
626 }
627 self.network.advance(Duration::from_millis(40));
628 }
629
630 fn accept_old_owner_local_write_before_expiry(&mut self) {
631 assert!(
632 self.old_owner.local_now_ms(self.network.clock().now_ms())
633 < self.old_owner.lease_expires_local_ms,
634 "old owner should still believe its lease is alive under skew"
635 );
636 let write_id = self.next_write_id();
637 self.record_owner_durable_write(
638 self.old_owner.id.clone(),
639 self.old_owner.term,
640 self.old_owner.epoch,
641 write_id,
642 AckPolicy::Local,
643 );
644 }
645
646 fn advance_until_old_owner_self_fences(&mut self) {
647 while self.old_owner.local_now_ms(self.network.clock().now_ms())
648 < self.old_owner.lease_expires_local_ms
649 {
650 self.network.advance(Duration::from_millis(10));
651 }
652 self.old_owner.self_fenced = true;
653 self.self_fenced_at_ms = Some(self.network.clock().now_ms());
654 }
655
656 fn promote_covered_replica(&mut self) {
657 self.network.advance(Duration::from_millis(10));
658 let members = replica_members_with_old_owner();
659 let stores = shared_vote_stores(&members);
660 let request = candidate_request(
661 "replica-a",
662 self.old_owner.term,
663 self.committed_watermark,
664 self.committed_watermark,
665 );
666 let mut tx = NetworkElectionTransport::new(
667 &mut self.network,
668 members,
669 stores,
670 "replica-a".to_string(),
671 self.committed_watermark,
672 );
673 let outcome = ElectionCoordinator::run(&request, &mut tx, Duration::from_secs(60));
674 let ElectionOutcome::Elected { term, .. } = outcome else {
675 panic!("covered replica must be promoted, got {outcome:?}");
676 };
677
678 self.current_epoch += 1;
679 self.promoted_at_ms = Some(self.network.clock().now_ms());
680 self.promoted_owner = Some("replica-a".to_string());
681 assert_eq!(
682 term,
683 self.old_owner.term + 1,
684 "supervisor promotion should advance the term"
685 );
686 }
687
688 fn accept_promoted_owner_sync_write(&mut self) {
689 let write_id = self.next_write_id();
690 let term = self.old_owner.term + 1;
691 self.record_owner_durable_write(
692 "replica-a".to_string(),
693 term,
694 self.current_epoch,
695 write_id,
696 AckPolicy::Synchronous,
697 );
698 let replicated =
699 self.replicate_commit("replica-a", term, write_id, &["replica-b", "replica-c"]);
700 assert!(
701 replicated.len() >= 2,
702 "promoted owner must synchronously replicate write {write_id}"
703 );
704 }
705
706 fn refuse_old_owner_late_write(&mut self) {
707 let stale_lease = WriterLease {
708 database_key: "main".to_string(),
709 holder_id: self.old_owner.id.clone(),
710 term: self.old_owner.term,
711 generation: 1,
712 acquired_at_ms: 0,
713 expires_at_ms: u64::MAX,
714 };
715 self.old_owner_late_write_refused =
716 self.old_owner.self_fenced || stale_lease.fenced_by_term(self.old_owner.term + 1);
717 assert!(
718 self.old_owner_late_write_refused,
719 "old owner must not admit writes after promotion"
720 );
721 }
722
723 fn assert_no_double_owner_window(&self) {
724 let self_fenced_at_ms = self.self_fenced_at_ms.expect("self fence happened");
725 let promoted_at_ms = self.promoted_at_ms.expect("promotion happened");
726 assert!(
727 self_fenced_at_ms <= promoted_at_ms,
728 "promotion at {promoted_at_ms} overlapped old owner until {self_fenced_at_ms}"
729 );
730
731 let mut writers_by_instant_and_epoch = BTreeMap::new();
732 for write in &self.accepted {
733 assert!(write.term > 0, "accepted writes must be term-stamped");
734 let previous = writers_by_instant_and_epoch
735 .insert((write.accepted_at_ms, write.epoch), write.owner_id.clone());
736 assert!(
737 previous
738 .as_ref()
739 .is_none_or(|owner| owner == &write.owner_id),
740 "two owners accepted durable writes at t={} epoch={}: {:?} and {}",
741 write.accepted_at_ms,
742 write.epoch,
743 previous,
744 write.owner_id
745 );
746 }
747 }
748
749 fn assert_acked_write_loss_oracle(&self) -> DstLeaseExpiryReport {
750 let surviving: BTreeSet<u64> = self
751 .replica_logs
752 .get(self.promoted_owner.as_deref().expect("promoted owner"))
753 .expect("promoted owner log")
754 .clone();
755 let mut sync_ack_loss_count = 0;
756 let mut local_policy_losses = 0;
757
758 for write in &self.accepted {
759 let survived = surviving.contains(&write.write_id);
760 match write.policy {
761 AckPolicy::Synchronous => {
762 if !survived {
763 sync_ack_loss_count += 1;
764 }
765 }
766 AckPolicy::Local => {
767 if !survived {
768 local_policy_losses += 1;
769 }
770 }
771 }
772 }
773
774 DstLeaseExpiryReport {
775 promoted_owner: self.promoted_owner.clone().expect("promoted owner"),
776 self_fenced_owner: self.old_owner.id.clone(),
777 old_owner_late_write_refused: self.old_owner_late_write_refused,
778 sync_ack_loss_count,
779 local_policy_losses,
780 }
781 }
782
783 fn replicate_commit(
784 &mut self,
785 owner_id: &str,
786 term: u64,
787 write_id: u64,
788 peers: &[&str],
789 ) -> BTreeSet<String> {
790 for peer in peers {
791 let _ = self.network.send(
792 owner_id.to_string(),
793 (*peer).to_string(),
794 ReplicationControlMessage::LogicalCommit {
795 term,
796 lsn: write_id,
797 payload_hash: format!("write-{write_id}"),
798 },
799 );
800 }
801 self.network.advance(Duration::from_millis(25));
802
803 let mut replicated = BTreeSet::new();
804 for peer in peers {
805 for delivery in self.network.drain_ready_for(peer) {
806 if let ReplicationControlMessage::LogicalCommit { lsn, .. } = delivery.message {
807 self.replica_logs
808 .get_mut(*peer)
809 .expect("known peer")
810 .insert(lsn);
811 replicated.insert((*peer).to_string());
812 }
813 }
814 }
815 replicated
816 }
817
818 fn record_owner_durable_write(
819 &mut self,
820 owner_id: String,
821 term: u64,
822 epoch: u64,
823 write_id: u64,
824 policy: AckPolicy,
825 ) {
826 self.replica_logs
827 .get_mut(&owner_id)
828 .expect("known owner")
829 .insert(write_id);
830 self.accepted.push(SimAcceptedWrite {
831 owner_id,
832 term,
833 epoch,
834 write_id,
835 policy,
836 accepted_at_ms: self.network.clock().now_ms(),
837 });
838 }
839
840 fn next_write_id(&mut self) -> u64 {
841 let write_id = self.next_write_id;
842 self.next_write_id += 1;
843 write_id
844 }
845 }
846
847 #[derive(Debug)]
848 struct DstLeaseExpiryReport {
849 promoted_owner: String,
850 self_fenced_owner: String,
851 old_owner_late_write_refused: bool,
852 sync_ack_loss_count: usize,
853 local_policy_losses: usize,
854 }
855
856 #[derive(Clone, Copy, Debug, Eq, PartialEq)]
857 enum AckPolicy {
858 Synchronous,
859 Local,
860 }
861
862 #[derive(Debug)]
863 struct SimOwner {
864 id: String,
865 term: u64,
866 epoch: u64,
867 lease_expires_local_ms: i64,
868 clock_skew_ms: i64,
869 self_fenced: bool,
870 }
871
872 impl SimOwner {
873 fn local_now_ms(&self, simulated_now_ms: u64) -> i64 {
874 simulated_now_ms as i64 + self.clock_skew_ms
875 }
876 }
877
878 #[derive(Debug)]
879 struct SimAcceptedWrite {
880 owner_id: String,
881 term: u64,
882 epoch: u64,
883 write_id: u64,
884 policy: AckPolicy,
885 accepted_at_ms: u64,
886 }
887
888 struct NetworkElectionTransport<'a> {
889 network: &'a mut InProcessReplicationNetwork<ReplicationControlMessage>,
890 members: Vec<Member>,
891 stores: BTreeMap<String, Rc<MemoryLastVoteStore>>,
892 candidate_id: String,
893 watermark: u64,
894 bumped_term: Option<u64>,
895 promoted_term: Option<u64>,
896 }
897
898 impl<'a> NetworkElectionTransport<'a> {
899 fn new(
900 network: &'a mut InProcessReplicationNetwork<ReplicationControlMessage>,
901 members: Vec<Member>,
902 stores: BTreeMap<String, Rc<MemoryLastVoteStore>>,
903 candidate_id: String,
904 watermark: u64,
905 ) -> Self {
906 Self {
907 network,
908 members,
909 stores,
910 candidate_id,
911 watermark,
912 bumped_term: None,
913 promoted_term: None,
914 }
915 }
916 }
917
918 impl ElectionTransport for NetworkElectionTransport<'_> {
919 fn members(&self) -> Vec<Member> {
920 self.members.clone()
921 }
922
923 fn request_vote(&mut self, peer_id: &str, req: &VoteRequest) -> VoteDecision {
924 let outcome = self.network.send(
925 self.candidate_id.clone(),
926 peer_id.to_string(),
927 ReplicationControlMessage::ElectionVoteRequest(req.clone()),
928 );
929 if !matches!(outcome, SendOutcome::Accepted { .. }) {
930 return unreachable_refusal(req);
931 }
932 if !self.network.advance_to_next_delivery() {
933 return unreachable_refusal(req);
934 }
935 let requests = self.network.drain_ready_for(peer_id);
936 let Some(request) = requests
937 .into_iter()
938 .find_map(|delivery| match delivery.message {
939 ReplicationControlMessage::ElectionVoteRequest(request) => Some(request),
940 _ => None,
941 })
942 else {
943 return unreachable_refusal(req);
944 };
945
946 let store = self.stores.get(peer_id).expect("known voter").clone();
947 let voter = Voter::new(peer_id, RcStore(store));
948 let decision = voter
949 .consider(&request, self.watermark)
950 .expect("memory vote store");
951 let outcome = self.network.send(
952 peer_id.to_string(),
953 self.candidate_id.clone(),
954 ReplicationControlMessage::ElectionVoteDecision(decision.clone()),
955 );
956 if !matches!(outcome, SendOutcome::Accepted { .. }) {
957 return unreachable_refusal(req);
958 }
959 if !self.network.advance_to_next_delivery() {
960 return unreachable_refusal(req);
961 }
962 self.network
963 .drain_ready_for(&self.candidate_id)
964 .into_iter()
965 .find_map(|delivery| match delivery.message {
966 ReplicationControlMessage::ElectionVoteDecision(decision) => Some(decision),
967 _ => None,
968 })
969 .unwrap_or_else(|| unreachable_refusal(req))
970 }
971
972 fn elapsed(&self) -> Duration {
973 self.network.clock().elapsed()
974 }
975
976 fn bump_term(&mut self, new_term: u64) {
977 self.bumped_term = Some(new_term);
978 }
979
980 fn promote(&mut self, new_term: u64) {
981 self.promoted_term = Some(new_term);
982 }
983 }
984
985 #[derive(Clone)]
986 struct RcStore(Rc<MemoryLastVoteStore>);
987
988 impl LastVoteStore for RcStore {
989 fn load(&self) -> Result<LastVote, LastVoteError> {
990 self.0.load()
991 }
992
993 fn persist(&self, vote: &LastVote) -> Result<(), LastVoteError> {
994 self.0.persist(vote)
995 }
996 }
997
998 fn unreachable_refusal(req: &VoteRequest) -> VoteDecision {
999 VoteDecision::Refused(RefusalReason::StaleTerm {
1000 candidate_term: req.term,
1001 voter_term: u64::MAX,
1002 })
1003 }
1004
1005 fn five_voters() -> Vec<Member> {
1006 vec![
1007 Member::data_voting("a"),
1008 Member::data_voting("b"),
1009 Member::data_voting("c"),
1010 Member::data_voting("d"),
1011 Member::data_voting("e"),
1012 ]
1013 }
1014
1015 fn replica_members_with_old_owner() -> Vec<Member> {
1016 vec![
1017 Member::data_voting("old-primary"),
1018 Member::data_voting("replica-a"),
1019 Member::data_voting("replica-b"),
1020 Member::data_voting("replica-c"),
1021 Member::data_voting("replica-d"),
1022 ]
1023 }
1024
1025 fn shared_vote_stores(members: &[Member]) -> BTreeMap<String, Rc<MemoryLastVoteStore>> {
1026 members
1027 .iter()
1028 .map(|member| (member.id.clone(), Rc::new(MemoryLastVoteStore::new())))
1029 .collect()
1030 }
1031
1032 fn partitioned_network(seed: u64) -> InProcessReplicationNetwork<ReplicationControlMessage> {
1033 let mut network = InProcessReplicationNetwork::new(
1034 seed,
1035 NetworkFaults {
1036 loss_per_million: 0,
1037 max_delay_ms: 20,
1038 reorder: true,
1039 },
1040 );
1041 for left in ["a", "b", "c"] {
1042 for right in ["d", "e"] {
1043 network.partition(left, right);
1044 }
1045 }
1046 network
1047 }
1048
1049 fn candidate_request(id: &str, current_term: u64, lsn: u64, watermark: u64) -> ElectionRequest {
1050 ElectionRequest {
1051 candidate: Member::data_voting(id),
1052 current_term,
1053 last_log_lsn: lsn,
1054 commit_watermark: watermark,
1055 }
1056 }
1057
1058 fn lease_store(tag: &str) -> LeaseStore {
1059 use crate::storage::backend::LocalBackend;
1060
1061 LeaseStore::new(Arc::new(LocalBackend)).with_prefix(format!(
1062 "{}/reddb-{tag}-{}",
1063 std::env::temp_dir().to_string_lossy(),
1064 crate::utils::now_unix_nanos(),
1065 ))
1066 }
1067}