Skip to main content

reddb_server/replication/
dst.rs

1//! Deterministic simulation testing helpers for the replication control plane.
2//!
3//! The simulator is intentionally in-process and single-threaded. It gives
4//! control-plane tests a transport that can inject partition, delay, reorder,
5//! and message loss without opening sockets or depending on Tokio scheduler
6//! ordering.
7
8use 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}