use std::collections::{BTreeMap, BTreeSet, VecDeque};
use super::bootstrap::{
inspect, BootDecision, BootstrapEffect, Integrity, QuarantineReason, QuarantinedNode,
RecoveredState, RejoinMessage,
};
use super::replica::{Effect, LogEntry, Message, ReplicaCore, Role, NOOP_PAYLOAD};
use super::types::{ClusterId, HardState, LogPosition, NodeId, Term};
const CLUSTER: ClusterId = ClusterId(7);
#[derive(Debug, Clone, PartialEq, Eq)]
enum Wire {
R(Message),
B(RejoinMessage),
}
struct Rng(u64);
impl Rng {
fn new(seed: u64) -> Self {
Rng(seed.max(1))
}
fn next(&mut self) -> u64 {
let mut x = self.0;
x ^= x >> 12;
x ^= x << 25;
x ^= x >> 27;
self.0 = x;
x.wrapping_mul(0x2545F4914F6CDD1D)
}
fn chance(&mut self, pct: u64) -> bool {
self.next() % 100 < pct
}
fn pick(&mut self, n: usize) -> usize {
(self.next() % n as u64) as usize
}
}
struct InFlight {
from: NodeId,
to: NodeId,
msg: Wire,
}
struct SimNode {
core: Option<ReplicaCore>,
quarantined: Option<QuarantinedNode>,
disk_hard: HardState,
disk_log: Vec<LogEntry>,
torn_hard: bool,
corrupt_log: bool,
applied: u64,
}
struct Sim {
nodes: BTreeMap<NodeId, SimNode>,
voters: BTreeSet<NodeId>,
net: VecDeque<InFlight>,
rng: Rng,
cut: BTreeSet<(NodeId, NodeId)>,
grants_sent: BTreeMap<(NodeId, Term), BTreeSet<NodeId>>,
leaders: BTreeMap<Term, BTreeSet<NodeId>>,
committed_at: BTreeMap<u64, LogEntry>,
advertised: BTreeMap<(NodeId, Term), LogPosition>,
freshness_at_grant: Vec<(LogPosition, LogPosition)>,
preserves: Vec<NodeId>,
alarms: Vec<(NodeId, Vec<QuarantineReason>)>,
}
impl Sim {
fn new(
node_logs: &[(u64, Vec<LogEntry>)],
start_term: u64,
seed: u64,
use_pre_vote: bool,
) -> Self {
let voters: BTreeSet<NodeId> = node_logs.iter().map(|(id, _)| NodeId(*id)).collect();
let hard = HardState {
current_term: Term(start_term),
voted_for: None,
};
let mut nodes = BTreeMap::new();
for (id, log) in node_logs {
let nid = NodeId(*id);
let core = ReplicaCore::new(nid, voters.clone(), hard, log.clone(), use_pre_vote);
nodes.insert(
nid,
SimNode {
core: Some(core),
quarantined: None,
disk_hard: hard,
disk_log: log.clone(),
torn_hard: false,
corrupt_log: false,
applied: 0,
},
);
}
Sim {
nodes,
voters,
net: VecDeque::new(),
rng: Rng::new(seed),
cut: BTreeSet::new(),
grants_sent: BTreeMap::new(),
leaders: BTreeMap::new(),
committed_at: BTreeMap::new(),
advertised: BTreeMap::new(),
freshness_at_grant: Vec::new(),
preserves: Vec::new(),
alarms: Vec::new(),
}
}
fn run_effects(&mut self, id: NodeId, effects: Vec<Effect>, crash_at: Option<u8>) {
let mut queue: VecDeque<Effect> = effects.into();
while let Some(eff) = queue.pop_front() {
match eff {
Effect::Persist { hard, log } => {
if crash_at == Some(0) {
self.crash(id);
return;
}
{
let node = self.nodes.get_mut(&id).unwrap();
node.disk_hard = hard;
node.disk_log = log;
}
if crash_at == Some(1) {
self.crash(id);
return;
}
let flushed = self
.nodes
.get_mut(&id)
.unwrap()
.core
.as_mut()
.unwrap()
.state_persisted();
for f in flushed {
queue.push_back(f);
}
}
Effect::Send { to, msg } => self.record_and_route(id, to, msg),
Effect::Broadcast { msg } => {
let peers: Vec<NodeId> =
self.voters.iter().copied().filter(|p| *p != id).collect();
for to in peers {
self.record_and_route(id, to, msg.clone());
}
}
Effect::BecameLeader { term } => {
self.leaders.entry(term).or_default().insert(id);
}
Effect::SteppedDown { .. } => {}
Effect::CommitAdvanced { to } => self.ledger_commit(id, to),
}
}
}
fn ledger_commit(&mut self, id: NodeId, to: u64) {
let from = self.nodes[&id].applied + 1;
for i in from..=to {
let e = *self.nodes[&id]
.core
.as_ref()
.expect("commit on live node")
.entry(i)
.expect("committed index must be in log");
match self.committed_at.get(&i) {
None => {
self.committed_at.insert(i, e);
}
Some(prev) => assert_eq!(
*prev, e,
"AUTHORITY SAFETY VIOLATED: index {i} applied with two \
different entries ({prev:?} vs {e:?}, second by {id:?})"
),
}
}
let node = self.nodes.get_mut(&id).unwrap();
node.applied = node.applied.max(to);
}
fn record_and_route(&mut self, from: NodeId, to: NodeId, msg: Message) {
match &msg {
Message::VoteRequest {
term,
candidate,
last_log,
} => {
self.advertised.insert((*candidate, *term), *last_log);
}
Message::VoteResponse {
term,
granted: true,
} => {
self.grants_sent
.entry((from, *term))
.or_default()
.insert(to);
let voter_log = self.nodes[&from]
.core
.as_ref()
.map(|c| c.last_log())
.unwrap_or(LogPosition::ZERO);
if let Some(cand) = self.advertised.get(&(to, *term)) {
self.freshness_at_grant.push((voter_log, *cand));
}
}
_ => {}
}
self.net.push_back(InFlight {
from,
to,
msg: Wire::R(msg),
});
}
fn crash(&mut self, id: NodeId) {
self.nodes.get_mut(&id).unwrap().core = None;
}
fn crash_torn(&mut self, id: NodeId) {
let node = self.nodes.get_mut(&id).unwrap();
node.core = None;
node.quarantined = None;
node.torn_hard = true;
}
fn crash_corrupt(&mut self, id: NodeId) {
let node = self.nodes.get_mut(&id).unwrap();
node.core = None;
node.quarantined = None;
node.corrupt_log = true;
}
fn restart_via_bootstrap(&mut self, id: NodeId) {
let voters = self.voters.clone();
let node = self.nodes.get_mut(&id).unwrap();
let recovered = RecoveredState {
cluster_id: Some(CLUSTER),
hard: Some(node.disk_hard),
log: Some(node.disk_log.clone()),
commit_marker: 0, integrity: Integrity {
hard_state_verified: !node.torn_hard,
log_verified: !node.corrupt_log,
},
};
match inspect(CLUSTER, &recovered) {
BootDecision::Healthy { hard, log } => {
node.core = Some(ReplicaCore::new(id, voters, hard, log, false));
node.quarantined = None;
node.applied = node.applied.min(node.disk_log.len() as u64);
}
BootDecision::Quarantine { reasons, term_hint } => {
node.core = None;
node.quarantined = Some(QuarantinedNode::new(id, CLUSTER, reasons, term_hint));
}
}
}
fn tick_rejoin(&mut self, id: NodeId, leader_hint: NodeId) {
let Some(q) = self.nodes.get_mut(&id).unwrap().quarantined.as_mut() else {
return;
};
let effects = q.tick_rejoin(leader_hint);
self.run_bootstrap_effects(id, effects);
}
fn run_bootstrap_effects(&mut self, id: NodeId, effects: Vec<BootstrapEffect>) {
for eff in effects {
match eff {
BootstrapEffect::PreserveOldState => {
self.preserves.push(id);
}
BootstrapEffect::Alarm { reasons } => {
self.alarms.push((id, reasons));
}
BootstrapEffect::Send { to, msg } => {
self.net.push_back(InFlight {
from: id,
to,
msg: Wire::B(msg),
});
}
BootstrapEffect::AdoptSnapshot {
cluster_id: _,
hard,
log,
} => {
assert!(
self.preserves.contains(&id),
"adopt without preserving old state first"
);
let voters = self.voters.clone();
let node = self.nodes.get_mut(&id).unwrap();
node.disk_hard = hard;
node.disk_log = log.clone();
node.torn_hard = false;
node.corrupt_log = false;
node.quarantined = None;
node.core = Some(ReplicaCore::new(id, voters, hard, log, false));
node.applied = node.applied.min(node.disk_log.len() as u64);
}
}
}
}
fn restart(&mut self, id: NodeId, use_pre_vote: bool) {
let voters = self.voters.clone();
let node = self.nodes.get_mut(&id).unwrap();
node.core = Some(ReplicaCore::new(
id,
voters,
node.disk_hard,
node.disk_log.clone(),
use_pre_vote,
));
node.quarantined = None;
node.applied = node.applied.min(node.disk_log.len() as u64);
}
fn timeout(&mut self, id: NodeId, crash_at: Option<u8>) {
if let Some(core) = self.nodes.get_mut(&id).unwrap().core.as_mut() {
let effects = core.on_election_timeout();
self.run_effects(id, effects, crash_at);
}
}
fn propose(&mut self, id: NodeId, payload: u64) -> bool {
let Some(core) = self.nodes.get_mut(&id).unwrap().core.as_mut() else {
return false;
};
match core.propose(payload) {
Some(effects) => {
self.run_effects(id, effects, None);
true
}
None => false,
}
}
fn heartbeat(&mut self, id: NodeId) {
if let Some(core) = self.nodes.get_mut(&id).unwrap().core.as_mut() {
let effects = core.tick_heartbeat();
self.run_effects(id, effects, None);
}
}
fn deliver_one(&mut self, crash_at: Option<u8>) -> bool {
let Some(m) = self.net.pop_front() else {
return false;
};
let key = (m.from.min(m.to), m.from.max(m.to));
if self.cut.contains(&key) {
return true;
}
match m.msg {
Wire::R(msg) => {
let Some(node) = self.nodes.get_mut(&m.to) else {
return true;
};
let Some(core) = node.core.as_mut() else {
return true;
};
let effects = core.on_message(m.from, msg, false);
self.run_effects(m.to, effects, crash_at);
}
Wire::B(RejoinMessage::Request { node: asker }) => {
let grant = self
.nodes
.get(&m.to)
.and_then(|n| n.core.as_ref())
.and_then(|c| c.rejoin_grant());
if let Some((term, log, commit)) = grant {
self.net.push_back(InFlight {
from: m.to,
to: asker,
msg: Wire::B(RejoinMessage::Grant {
cluster_id: CLUSTER,
term,
log,
commit,
verified: true,
}),
});
}
}
Wire::B(grant @ RejoinMessage::Grant { .. }) => {
let effects = {
let Some(node) = self.nodes.get_mut(&m.to) else {
return true;
};
let Some(q) = node.quarantined.as_mut() else {
return true; };
q.on_grant(m.from, grant)
};
self.run_bootstrap_effects(m.to, effects);
}
}
true
}
fn drain(&mut self) {
while self.deliver_one(None) {}
}
fn current_leader(&self) -> Option<NodeId> {
self.nodes
.iter()
.find(|(_, n)| n.core.as_ref().is_some_and(|c| c.role() == Role::Leader))
.map(|(id, _)| *id)
}
fn check_vote_safety(&self) {
for ((voter, term), cands) in &self.grants_sent {
assert!(
cands.len() <= 1,
"VOTE SAFETY VIOLATED: voter {voter:?} granted {cands:?} in {term:?}"
);
}
}
fn check_suffix_protection(&self) {
for (voter_log, cand_log) in &self.freshness_at_grant {
assert!(
cand_log.is_at_least_as_up_to_date_as(voter_log),
"SUFFIX PROTECTION VIOLATED: granted candidate at {cand_log:?} \
while voter was at {voter_log:?}"
);
}
}
fn check_single_leader_per_term(&self) {
for (term, ls) in &self.leaders {
assert!(ls.len() <= 1, "TWO LEADERS IN {term:?}: {ls:?}");
}
}
fn check_committed_prefix_integrity(&self) {
for (id, node) in &self.nodes {
let Some(core) = node.core.as_ref() else {
continue;
};
for i in 1..=node.applied {
if let Some(expected) = self.committed_at.get(&i) {
let actual = core.entry(i);
assert_eq!(
actual,
Some(expected),
"COMMITTED PREFIX DAMAGED on {id:?} at index {i}"
);
}
}
}
}
fn check_all(&self) {
self.check_vote_safety();
self.check_suffix_protection();
self.check_single_leader_per_term();
self.check_committed_prefix_integrity();
}
}
fn entries(terms: &[u64]) -> Vec<LogEntry> {
terms
.iter()
.enumerate()
.map(|(i, t)| LogEntry {
term: Term(*t),
payload: 1000 + i as u64, })
.collect()
}
fn empty() -> Vec<LogEntry> {
Vec::new()
}
#[test]
fn three_nodes_elect_exactly_one_leader() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 42, true);
sim.timeout(NodeId(1), None);
sim.drain();
sim.check_all();
assert_eq!(
sim.leaders.values().map(|s| s.len()).sum::<usize>(),
1,
"exactly one leadership event expected, got {:?}",
sim.leaders
);
assert_eq!(
sim.committed_at.get(&1).map(|e| e.payload),
Some(NOOP_PAYLOAD)
);
}
#[test]
fn r2_crash_before_persist_never_leaks_the_grant() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 7, false);
sim.timeout(NodeId(1), None);
let mut injected = false;
while !sim.net.is_empty() {
let to = sim.net.front().unwrap().to;
let inject = if to == NodeId(3) && !injected {
injected = true;
Some(0u8)
} else {
None
};
sim.deliver_one(inject);
}
sim.restart(NodeId(3), false);
assert_eq!(sim.nodes[&NodeId(3)].disk_hard, HardState::default());
sim.timeout(NodeId(2), None);
sim.drain();
sim.check_all();
}
#[test]
fn r2_crash_after_persist_binds_the_restarted_node() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 9, false);
sim.timeout(NodeId(1), None);
let mut injected = false;
while !sim.net.is_empty() {
let to = sim.net.front().unwrap().to;
let inject = if to == NodeId(3) && !injected {
injected = true;
Some(1u8)
} else {
None
};
sim.deliver_one(inject);
}
assert_eq!(
sim.nodes[&NodeId(3)].disk_hard,
HardState {
current_term: Term(1),
voted_for: Some(NodeId(1)),
}
);
sim.restart(NodeId(3), false);
let effects = sim
.nodes
.get_mut(&NodeId(3))
.unwrap()
.core
.as_mut()
.unwrap()
.on_message(
NodeId(2),
Message::VoteRequest {
term: Term(1),
candidate: NodeId(2),
last_log: LogPosition::ZERO,
},
false,
);
sim.run_effects(NodeId(3), effects, None);
sim.drain();
let grants = sim
.grants_sent
.get(&(NodeId(3), Term(1)))
.cloned()
.unwrap_or_default();
assert!(
!grants.contains(&NodeId(2)),
"restarted node granted a second candidate in the same term: {grants:?}"
);
sim.check_all();
}
#[test]
fn r3_stale_log_candidate_is_refused_by_fresher_voters() {
let stale = entries(&[1, 1, 1, 1, 1, 1, 1, 1, 1]); let fresh = entries(&[1, 1, 1, 1, 2]); let mut sim = Sim::new(&[(1, stale), (2, fresh.clone()), (3, fresh)], 2, 21, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.check_all();
assert!(
sim.leaders.values().all(|s| !s.contains(&NodeId(1))),
"stale-log candidate won an election: {:?}",
sim.leaders
);
sim.timeout(NodeId(2), None);
sim.drain();
sim.check_all();
assert!(
sim.leaders.values().any(|s| s.contains(&NodeId(2))),
"fresh candidate failed to win: {:?}",
sim.leaders
);
}
#[test]
fn partition_minority_cannot_elect_and_heal_converges() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 63, false);
sim.cut.insert((NodeId(1), NodeId(2)));
sim.cut.insert((NodeId(1), NodeId(3)));
sim.timeout(NodeId(1), None);
sim.timeout(NodeId(2), None);
sim.drain();
sim.check_all();
assert!(
sim.leaders.values().all(|s| !s.contains(&NodeId(1))),
"partitioned minority elected itself: {:?}",
sim.leaders
);
sim.cut.clear();
sim.drain();
sim.check_all();
}
#[test]
fn pre_vote_probe_never_burns_terms() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 5, true);
sim.cut.insert((NodeId(1), NodeId(2)));
sim.cut.insert((NodeId(1), NodeId(3)));
for _ in 0..25 {
sim.timeout(NodeId(1), None);
sim.drain();
}
assert_eq!(sim.nodes[&NodeId(1)].disk_hard.current_term, Term(0));
sim.check_all();
}
#[test]
fn replication_commits_across_cluster() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 11, false);
sim.timeout(NodeId(1), None);
sim.drain();
let leader = sim.current_leader().expect("leader elected");
for p in [101, 102, 103] {
assert!(sim.propose(leader, p), "propose on leader");
}
sim.drain();
sim.heartbeat(leader); sim.drain();
sim.check_all();
assert_eq!(sim.committed_at.get(&2).map(|e| e.payload), Some(101));
assert_eq!(sim.committed_at.get(&3).map(|e| e.payload), Some(102));
assert_eq!(sim.committed_at.get(&4).map(|e| e.payload), Some(103));
for (id, node) in &sim.nodes {
assert!(node.applied >= 4, "{id:?} applied only to {}", node.applied);
}
}
#[test]
fn r1_stale_leader_cannot_commit_and_gets_fenced() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 17, false);
sim.timeout(NodeId(1), None);
sim.drain();
assert_eq!(sim.current_leader(), Some(NodeId(1)));
let applied_before = sim.nodes[&NodeId(1)].applied;
sim.cut.insert((NodeId(1), NodeId(2)));
sim.cut.insert((NodeId(1), NodeId(3)));
assert!(sim.propose(NodeId(1), 201));
assert!(sim.propose(NodeId(1), 202));
sim.drain();
assert_eq!(
sim.nodes[&NodeId(1)].applied,
applied_before,
"stale leader advanced commit without a quorum"
);
sim.timeout(NodeId(2), None);
sim.drain();
assert!(sim.propose(NodeId(2), 301));
sim.drain();
sim.cut.clear();
sim.heartbeat(NodeId(2));
sim.drain();
sim.heartbeat(NodeId(2));
sim.drain();
sim.check_all();
assert!(
sim.committed_at
.values()
.all(|e| e.payload != 201 && e.payload != 202),
"stale leader's tentative writes leaked into committed history"
);
let committed_301 = sim
.committed_at
.iter()
.find(|(_, e)| e.payload == 301)
.map(|(i, _)| *i)
.expect("301 committed");
let n1 = sim.nodes[&NodeId(1)].core.as_ref().unwrap();
assert_eq!(n1.entry(committed_301).map(|e| e.payload), Some(301));
assert_eq!(sim.current_leader(), Some(NodeId(2)));
}
#[test]
fn follower_catchup_after_crash() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 29, false);
sim.timeout(NodeId(1), None);
sim.drain();
assert!(sim.propose(NodeId(1), 401));
let mut injected = false;
while !sim.net.is_empty() {
let to = sim.net.front().unwrap().to;
let inject = if to == NodeId(3) && !injected {
injected = true;
Some(0u8)
} else {
None
};
sim.deliver_one(inject);
}
assert!(sim.committed_at.values().any(|e| e.payload == 401));
sim.restart(NodeId(3), false);
sim.heartbeat(NodeId(1));
sim.drain();
sim.heartbeat(NodeId(1));
sim.drain();
sim.check_all();
let n3 = &sim.nodes[&NodeId(3)];
assert!(
n3.applied >= 2,
"restarted follower failed to catch up (applied {})",
n3.applied
);
}
#[test]
fn seeded_soak_invariants_hold() {
for seed in 1..30u64 {
let mut sim = Sim::new(
&[(1, empty()), (2, empty()), (3, empty())],
0,
seed,
seed % 2 == 0,
);
let mut payload = 100;
for _step in 0..300 {
let ids = [NodeId(1), NodeId(2), NodeId(3)];
match sim.rng.next() % 12 {
0 => {
let id = ids[sim.rng.pick(3)];
if sim.nodes[&id].core.is_some() {
let inject = sim.rng.chance(20).then(|| (sim.rng.next() % 2) as u8);
sim.timeout(id, inject);
}
}
1 => {
let id = ids[sim.rng.pick(3)];
if sim.nodes[&id].core.is_some() && sim.rng.chance(15) {
sim.crash(id);
}
}
2 => {
let id = ids[sim.rng.pick(3)];
if sim.nodes[&id].core.is_none() {
let pv = sim.rng.chance(50);
sim.restart(id, pv);
}
}
3 => {
let id = ids[sim.rng.pick(3)];
payload += 1;
let _ = sim.propose(id, payload);
}
4 => {
let id = ids[sim.rng.pick(3)];
sim.heartbeat(id);
}
5 => {
let a = ids[sim.rng.pick(3)];
let b = ids[sim.rng.pick(3)];
if a != b {
let key = (a.min(b), a.max(b));
if !sim.cut.remove(&key) {
sim.cut.insert(key);
}
}
}
_ => {
let inject = sim.rng.chance(10).then(|| (sim.rng.next() % 2) as u8);
sim.deliver_one(inject);
}
}
}
sim.cut.clear();
sim.drain();
sim.check_all();
}
}
#[test]
fn ct141_torn_node_quarantines_then_rejoins_via_leader() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 31, false);
sim.timeout(NodeId(1), None);
sim.drain();
assert!(sim.propose(NodeId(1), 501));
sim.drain();
sim.crash_torn(NodeId(3));
sim.restart_via_bootstrap(NodeId(3));
{
let n3 = &sim.nodes[&NodeId(3)];
let q = n3.quarantined.as_ref().expect("torn node must quarantine");
assert!(q.reasons().contains(&QuarantineReason::TornHardState));
assert!(q.stale_reads_allowed());
assert!(n3.core.is_none(), "quarantined node must not run a core");
}
sim.tick_rejoin(NodeId(3), NodeId(1));
sim.drain();
{
let n3 = &sim.nodes[&NodeId(3)];
assert!(n3.quarantined.is_none(), "rejoin did not complete");
assert!(n3.core.is_some());
assert!(!n3.torn_hard);
}
assert!(
sim.preserves.contains(&NodeId(3)),
"old state must be preserved before resync"
);
sim.heartbeat(NodeId(1));
sim.drain();
sim.check_all();
assert!(
sim.nodes[&NodeId(3)].applied >= 2,
"rejoined node failed to catch up (applied {})",
sim.nodes[&NodeId(3)].applied
);
}
#[test]
fn quarantined_node_cannot_vote() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 37, false);
sim.crash_torn(NodeId(3));
sim.restart_via_bootstrap(NodeId(3));
sim.cut.insert((NodeId(1), NodeId(2)));
sim.timeout(NodeId(1), None);
sim.timeout(NodeId(2), None);
sim.drain();
sim.check_all();
assert!(
sim.leaders.is_empty(),
"an election succeeded without a legal quorum: {:?}",
sim.leaders
);
assert!(sim.grants_sent.keys().all(|(voter, _)| *voter != NodeId(3)));
}
#[test]
fn corrupted_log_alarms_and_rejoins_from_verified_source() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 41, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.crash_corrupt(NodeId(3));
sim.restart_via_bootstrap(NodeId(3));
{
let q = sim.nodes[&NodeId(3)].quarantined.as_ref().unwrap();
assert!(q.reasons().contains(&QuarantineReason::LogCorruption));
assert!(!q.stale_reads_allowed(), "corrupt data must not be served");
}
sim.tick_rejoin(NodeId(3), NodeId(1));
sim.drain();
assert!(
sim.alarms.iter().any(|(id, _)| *id == NodeId(3)),
"corruption evidence must alarm"
);
assert!(sim.nodes[&NodeId(3)].quarantined.is_none());
sim.heartbeat(NodeId(1));
sim.drain();
sim.check_all();
}
#[test]
fn rejoin_requires_leadership_certificate() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 43, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.crash_torn(NodeId(3));
sim.restart_via_bootstrap(NodeId(3));
sim.tick_rejoin(NodeId(3), NodeId(2));
sim.drain();
assert!(sim.nodes[&NodeId(3)].quarantined.is_some());
sim.tick_rejoin(NodeId(3), NodeId(1));
sim.drain();
assert!(sim.nodes[&NodeId(3)].quarantined.is_none());
sim.check_all();
}
#[test]
fn seeded_soak_with_quarantine_invariants_hold() {
for seed in 1..20u64 {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, seed, false);
let mut payload = 5000;
for _step in 0..300 {
let ids = [NodeId(1), NodeId(2), NodeId(3)];
match sim.rng.next() % 14 {
0 => {
let id = ids[sim.rng.pick(3)];
if sim.nodes[&id].core.is_some() {
sim.timeout(id, None);
}
}
1 => {
let id = ids[sim.rng.pick(3)];
if sim.nodes[&id].core.is_some() && sim.rng.chance(10) {
if sim.rng.chance(50) {
sim.crash_torn(id);
} else {
sim.crash(id);
}
}
}
2 => {
let id = ids[sim.rng.pick(3)];
if sim.nodes[&id].core.is_none() && sim.nodes[&id].quarantined.is_none() {
sim.restart_via_bootstrap(id);
}
}
3 => {
let id = ids[sim.rng.pick(3)];
payload += 1;
let _ = sim.propose(id, payload);
}
4 => {
let id = ids[sim.rng.pick(3)];
sim.heartbeat(id);
}
5 => {
let id = ids[sim.rng.pick(3)];
let hint = ids[sim.rng.pick(3)];
if id != hint {
sim.tick_rejoin(id, hint);
}
}
6 => {
let a = ids[sim.rng.pick(3)];
let b = ids[sim.rng.pick(3)];
if a != b {
let key = (a.min(b), a.max(b));
if !sim.cut.remove(&key) {
sim.cut.insert(key);
}
}
}
_ => {
sim.deliver_one(None);
}
}
}
sim.cut.clear();
sim.drain();
sim.check_all();
}
}
#[test]
fn codex_stale_grant_below_term_hint_is_refused() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 47, false);
sim.timeout(NodeId(1), None);
sim.drain();
assert_eq!(sim.current_leader(), Some(NodeId(1)));
sim.cut.insert((NodeId(1), NodeId(2)));
sim.cut.insert((NodeId(1), NodeId(3)));
sim.timeout(NodeId(2), None);
sim.drain();
sim.crash_torn(NodeId(3));
sim.restart_via_bootstrap(NodeId(3));
assert!(sim.nodes[&NodeId(3)].quarantined.is_some());
sim.cut.clear();
sim.tick_rejoin(NodeId(3), NodeId(1));
sim.drain();
assert!(
sim.nodes[&NodeId(3)].quarantined.is_some(),
"below-hint grant was adopted — term regression"
);
sim.heartbeat(NodeId(2));
sim.drain(); sim.tick_rejoin(NodeId(3), NodeId(2));
sim.drain();
assert!(sim.nodes[&NodeId(3)].quarantined.is_none());
let grants = sim
.grants_sent
.get(&(NodeId(3), Term(2)))
.cloned()
.unwrap_or_default();
assert!(grants.len() <= 1 && !grants.contains(&NodeId(1)));
sim.check_all();
}
#[test]
fn codex_duplicate_response_does_not_regress_next_index() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 53, false);
sim.timeout(NodeId(1), None);
sim.drain();
for p in [601, 602, 603] {
assert!(sim.propose(NodeId(1), p));
}
sim.drain();
let next_before = sim.nodes[&NodeId(1)]
.core
.as_ref()
.unwrap()
.next_index_of(NodeId(2))
.unwrap();
assert_eq!(next_before, 5);
let effects = sim
.nodes
.get_mut(&NodeId(1))
.unwrap()
.core
.as_mut()
.unwrap()
.on_message(
NodeId(2),
Message::AppendResponse {
term: Term(1),
success: true,
last_index: 1,
},
false,
);
sim.run_effects(NodeId(1), effects, None);
let next_after = sim.nodes[&NodeId(1)]
.core
.as_ref()
.unwrap()
.next_index_of(NodeId(2))
.unwrap();
assert_eq!(
next_after, 5,
"duplicate old success regressed next_index to {next_after}"
);
let effects = sim
.nodes
.get_mut(&NodeId(1))
.unwrap()
.core
.as_mut()
.unwrap()
.on_message(
NodeId(2),
Message::AppendResponse {
term: Term(1),
success: false,
last_index: 0,
},
false,
);
sim.run_effects(NodeId(1), effects, None);
let next_final = sim.nodes[&NodeId(1)]
.core
.as_ref()
.unwrap()
.next_index_of(NodeId(2))
.unwrap();
assert!(
next_final >= 5,
"stale failure regressed next_index to {next_final}"
);
sim.drain();
sim.check_all();
}