use std::collections::{BTreeMap, BTreeSet, VecDeque};
use super::bootstrap::{
inspect, BootDecision, BootstrapEffect, Integrity, QuarantineReason, QuarantinedNode,
RecoveredState, RejoinMessage,
};
use super::replica::{
Effect, KeyedProposal, LogEntry, Message, Payload, ReplicaCore, Role, Snapshot,
};
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_base: LogPosition,
disk_log: Vec<LogEntry>,
disk_claims: BTreeMap<u64, u64>,
disk_active: u32,
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>)>,
keyed_committed: BTreeMap<u64, (u64, Payload)>,
witnesses: BTreeSet<NodeId>,
incompat_alarms: Vec<(NodeId, NodeId)>,
}
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_base: LogPosition::ZERO,
disk_log: log.clone(),
disk_claims: log
.iter()
.enumerate()
.filter_map(|(i, e)| e.key.map(|k| (k, i as u64 + 1)))
.collect(),
disk_active: 0,
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(),
keyed_committed: BTreeMap::new(),
witnesses: BTreeSet::new(),
incompat_alarms: Vec::new(),
}
}
fn new_with_witnesses(
node_logs: &[(u64, Vec<LogEntry>)],
start_term: u64,
seed: u64,
use_pre_vote: bool,
witness_ids: &[u64],
) -> Self {
let mut sim = Sim::new(node_logs, start_term, seed, use_pre_vote);
sim.witnesses = witness_ids.iter().map(|w| NodeId(*w)).collect();
let w = sim.witnesses.clone();
for node in sim.nodes.values_mut() {
if let Some(core) = node.core.as_mut() {
core.set_witnesses(w.clone());
}
}
sim
}
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,
base,
log,
claims,
active,
} => {
if crash_at == Some(0) {
self.crash(id);
return;
}
{
let node = self.nodes.get_mut(&id).unwrap();
node.disk_hard = hard;
node.disk_base = base;
node.disk_log = log;
node.disk_claims = claims;
node.disk_active = active;
}
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),
Effect::InstallState { last_index } => {
let node = self.nodes.get_mut(&id).unwrap();
node.applied = node.applied.max(last_index);
}
Effect::PeerIncompatible { peer } => {
self.incompat_alarms.push((id, peer));
}
}
}
}
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")
.clone();
if let Some(k) = e.key {
match self.keyed_committed.get(&k) {
None => {
self.keyed_committed.insert(k, (i, e.payload.clone()));
}
Some((pi, pp)) => assert_eq!(
(*pi, pp),
(i, &e.payload),
"CLAIM DOUBLE-WRITE: key {k} committed at two places"
),
}
}
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()),
active: node.disk_active,
commit_marker: 0, integrity: Integrity {
hard_state_verified: !node.torn_hard,
log_verified: !node.corrupt_log,
},
};
match inspect(CLUSTER, u32::MAX, &recovered) {
BootDecision::Healthy { hard, log } => {
let mut core = ReplicaCore::new(id, voters, hard, log, false);
core.set_witnesses(self.witnesses.clone());
node.core = Some(core);
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,
base,
log,
claims,
active,
} => {
assert!(
self.preserves.contains(&id),
"adopt without preserving old state first"
);
let voters = self.voters.clone();
let w = self.witnesses.clone();
let node = self.nodes.get_mut(&id).unwrap();
node.disk_hard = hard;
node.disk_base = base;
node.disk_log = log.clone();
node.disk_claims = claims.clone();
node.disk_active = active;
node.torn_hard = false;
node.corrupt_log = false;
node.quarantined = None;
let mut core = ReplicaCore::new_from_durable(
id, voters, hard, base, log, claims, active, false,
);
core.set_witnesses(w);
node.core = Some(core);
node.applied = node
.applied
.min(node.disk_base.index + node.disk_log.len() as u64)
.max(node.disk_base.index);
}
}
}
}
fn restart(&mut self, id: NodeId, use_pre_vote: bool) {
let voters = self.voters.clone();
let node = self.nodes.get_mut(&id).unwrap();
let mut core = ReplicaCore::new_from_durable(
id,
voters,
node.disk_hard,
node.disk_base,
node.disk_log.clone(),
node.disk_claims.clone(),
node.disk_active,
use_pre_vote,
);
core.set_witnesses(self.witnesses.clone());
node.core = Some(core);
node.quarantined = None;
node.applied = node
.applied
.min(node.disk_base.index + node.disk_log.len() as u64)
.max(node.disk_base.index);
}
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::Test(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, base, log, claims, active, commit)) = grant {
self.net.push_back(InFlight {
from: m.to,
to: asker,
msg: Wire::B(RejoinMessage::Grant {
cluster_id: CLUSTER,
term,
base,
log,
claims,
active,
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 (core.base().index + 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::unkeyed(Term(*t), Payload::Test(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.clone()),
Some(Payload::Noop)
);
}
#[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,
supported: u32::MAX,
},
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!(sim.committed_at.get(&2).is_some_and(|e| e.payload == 101));
assert!(sim.committed_at.get(&3).is_some_and(|e| e.payload == 102));
assert!(sim.committed_at.get(&4).is_some_and(|e| e.payload == 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!(n1.entry(committed_301).is_some_and(|e| e.payload == 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,
unsupported: false,
},
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,
unsupported: false,
},
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();
}
impl Sim {
fn propose_keyed(&mut self, id: NodeId, key: u64, payload: u64) -> Option<KeyedProposal> {
let core = self.nodes.get_mut(&id).unwrap().core.as_mut()?;
let outcome = core.propose_keyed(key, Payload::Test(payload))?;
if let KeyedProposal::Appended { effects, index } = outcome {
self.run_effects(id, effects, None);
return Some(KeyedProposal::Appended {
index,
effects: Vec::new(), });
}
Some(outcome)
}
}
#[test]
fn keyed_commit_survives_failover_and_dedupes_retry() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 71, false);
sim.timeout(NodeId(1), None);
sim.drain();
let out = sim.propose_keyed(NodeId(1), 77, 707).expect("leader");
let orig_index = match out {
KeyedProposal::Appended { index, .. } => index,
other => panic!("expected fresh append, got {other:?}"),
};
sim.drain();
assert!(
sim.keyed_committed.contains_key(&77),
"keyed entry committed"
);
sim.crash(NodeId(1));
sim.timeout(NodeId(2), None);
sim.drain();
let retry = sim.propose_keyed(NodeId(2), 77, 707).expect("new leader");
match retry {
KeyedProposal::DuplicateCommitted { index } => {
assert_eq!(index, orig_index, "dedupe must return the ORIGINAL entry");
}
other => panic!("retry after committed failover must dedupe, got {other:?}"),
}
sim.check_all();
assert_eq!(
sim.keyed_committed.get(&77).map(|(i, p)| (*i, p.clone())),
Some((orig_index, Payload::Test(707))),
"exactly one committed effect for the key"
);
}
#[test]
fn keyed_tentative_loss_reexecutes_cleanly() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 73, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.cut.insert((NodeId(1), NodeId(2)));
sim.cut.insert((NodeId(1), NodeId(3)));
let out = sim.propose_keyed(NodeId(1), 88, 808).expect("stale leader");
assert!(matches!(out, KeyedProposal::Appended { .. }));
sim.drain();
assert!(
!sim.keyed_committed.contains_key(&88),
"tentative keyed write must not commit without a quorum"
);
sim.timeout(NodeId(2), None);
sim.drain();
let retry = sim.propose_keyed(NodeId(2), 88, 808).expect("new leader");
let new_index = match retry {
KeyedProposal::Appended { index, .. } => index,
other => panic!("retry after tentative loss must re-execute, got {other:?}"),
};
sim.drain();
sim.cut.clear();
sim.heartbeat(NodeId(2));
sim.drain();
sim.heartbeat(NodeId(2));
sim.drain();
sim.check_all();
assert_eq!(
sim.keyed_committed.get(&88).map(|(i, p)| (*i, p.clone())),
Some((new_index, Payload::Test(808))),
"exactly one committed effect, at the canonical index"
);
let n1 = sim.nodes[&NodeId(1)].core.as_ref().unwrap();
assert_eq!(
n1.entry(new_index).and_then(|e| e.key),
Some(88),
"healed node holds the canonical keyed entry"
);
}
#[test]
fn keyed_retry_pending_parks_then_dedupes() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 79, false);
sim.timeout(NodeId(1), None);
sim.drain();
let core = sim
.nodes
.get_mut(&NodeId(1))
.unwrap()
.core
.as_mut()
.unwrap();
let out1 = core.propose_keyed(99, Payload::Test(909)).unwrap();
let index = match &out1 {
KeyedProposal::Appended { index, .. } => *index,
other => panic!("fresh append expected, got {other:?}"),
};
let out2 = core.propose_keyed(99, Payload::Test(909)).unwrap();
assert_eq!(out2, KeyedProposal::DuplicatePending { index });
let log_len_before = core.log_len();
if let KeyedProposal::Appended { effects, .. } = out1 {
sim.run_effects(NodeId(1), effects, None);
}
sim.drain();
let core = sim
.nodes
.get_mut(&NodeId(1))
.unwrap()
.core
.as_mut()
.unwrap();
assert_eq!(
core.log_len(),
log_len_before,
"no second append for the key"
);
let out3 = core.propose_keyed(99, Payload::Test(909)).unwrap();
assert_eq!(out3, KeyedProposal::DuplicateCommitted { index });
sim.check_all();
}
#[test]
fn seeded_soak_keyed_claims_hold_under_chaos() {
for seed in 1..15u64 {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, seed, false);
let keys = [11u64, 22, 33, 44];
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(40) {
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 | 4 => {
let id = ids[sim.rng.pick(3)];
let k = keys[sim.rng.pick(keys.len())];
if let Some(KeyedProposal::DuplicateCommitted { .. }) =
sim.propose_keyed(id, k, k * 10)
{
assert!(
sim.keyed_committed.contains_key(&k),
"DuplicateCommitted for key {k} without a \
committed effect (ghost success)"
);
}
}
5 => {
let id = ids[sim.rng.pick(3)];
sim.heartbeat(id);
}
6 => {
let id = ids[sim.rng.pick(3)];
let hint = ids[sim.rng.pick(3)];
if id != hint {
sim.tick_rejoin(id, hint);
}
}
7 => {
let a = ids[sim.rng.pick(3)];
let b = ids[sim.rng.pick(3)];
if a != b {
let kk = (a.min(b), a.max(b));
if !sim.cut.remove(&kk) {
sim.cut.insert(kk);
}
}
}
_ => {
sim.deliver_one(None);
}
}
}
sim.cut.clear();
sim.drain();
sim.check_all(); }
}
#[test]
fn witness_never_campaigns() {
let mut sim = Sim::new_with_witnesses(
&[(1, empty()), (2, empty()), (3, empty())],
0,
83,
false,
&[3],
);
for _ in 0..5 {
sim.timeout(NodeId(3), None);
sim.drain();
}
assert!(
sim.leaders.is_empty(),
"witness campaigned: {:?}",
sim.leaders
);
assert_eq!(sim.nodes[&NodeId(3)].disk_hard.current_term, Term(0));
}
#[test]
fn witness_votes_but_never_counts_for_commit() {
let mut sim = Sim::new_with_witnesses(
&[(1, empty()), (2, empty()), (3, empty())],
0,
89,
false,
&[3],
);
sim.cut.insert((NodeId(1), NodeId(2))); sim.timeout(NodeId(1), None);
sim.drain();
assert!(
sim.leaders.values().any(|s| s.contains(&NodeId(1))),
"witness vote must elect: {:?}",
sim.leaders
);
assert!(sim.propose(NodeId(1), 901));
sim.drain();
sim.heartbeat(NodeId(1));
sim.drain();
assert_eq!(
sim.nodes[&NodeId(1)].applied,
0,
"commit advanced on witness acks alone"
);
assert!(sim.committed_at.is_empty());
sim.cut.clear();
sim.heartbeat(NodeId(1));
sim.drain();
sim.check_all();
assert!(
sim.committed_at.values().any(|e| e.payload == 901),
"entry must commit once the data quorum is reachable"
);
}
#[test]
fn witness_tiebreak_preserves_all_committed_entries() {
let mut sim = Sim::new_with_witnesses(
&[(1, empty()), (2, empty()), (3, empty())],
0,
97,
false,
&[3],
);
sim.timeout(NodeId(1), None);
sim.drain();
for p in [911, 912, 913] {
assert!(sim.propose(NodeId(1), p));
}
sim.drain();
let committed_before: Vec<Payload> = sim
.committed_at
.values()
.map(|e| e.payload.clone())
.collect();
assert!(
committed_before.iter().any(|p| *p == 913),
"writes committed"
);
sim.crash(NodeId(1));
sim.timeout(NodeId(2), None);
sim.drain();
assert!(
sim.leaders.values().any(|s| s.contains(&NodeId(2))),
"surviving data node must win with the witness vote"
);
sim.propose(NodeId(2), 914);
sim.drain();
assert!(
!sim.committed_at.values().any(|e| e.payload == 914),
"write committed without a data quorum"
);
sim.check_all(); for p in [911, 912, 913] {
assert!(
sim.committed_at.values().any(|e| e.payload == p),
"committed entry {p} lost across witness-assisted failover"
);
}
sim.restart(NodeId(1), false);
sim.heartbeat(NodeId(2));
sim.drain();
sim.heartbeat(NodeId(2));
sim.drain();
sim.check_all();
assert!(
sim.committed_at.values().any(|e| e.payload == 914),
"stalled write must commit once the data quorum returns"
);
}
#[test]
fn seeded_soak_witness_topology_keyed_claims_hold() {
for seed in 1..10u64 {
let mut sim = Sim::new_with_witnesses(
&[(1, empty()), (2, empty()), (3, empty())],
0,
seed,
false,
&[3],
);
let keys = [55u64, 66];
for _step in 0..250 {
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() {
sim.timeout(id, None);
}
}
1 => {
let id = ids[sim.rng.pick(3)];
if sim.nodes[&id].core.is_some() && sim.rng.chance(10) {
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(id, false);
}
}
3 | 4 => {
let id = ids[sim.rng.pick(3)];
let k = keys[sim.rng.pick(keys.len())];
if let Some(KeyedProposal::DuplicateCommitted { .. }) =
sim.propose_keyed(id, k, k * 10)
{
assert!(sim.keyed_committed.contains_key(&k));
}
}
5 => {
let id = ids[sim.rng.pick(3)];
sim.heartbeat(id);
}
6 => {
let a = ids[sim.rng.pick(3)];
let b = ids[sim.rng.pick(3)];
if a != b {
let kk = (a.min(b), a.max(b));
if !sim.cut.remove(&kk) {
sim.cut.insert(kk);
}
}
}
_ => {
sim.deliver_one(None);
}
}
}
sim.cut.clear();
sim.drain();
sim.check_all();
}
}
#[test]
fn compaction_preserves_claims_no_replay_window() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 101, false);
sim.timeout(NodeId(1), None);
sim.drain();
let out = sim.propose_keyed(NodeId(1), 121, 1210).expect("leader");
let orig_index = match out {
KeyedProposal::Appended { index, .. } => index,
other => panic!("fresh append expected, got {other:?}"),
};
sim.drain();
{
let core = sim
.nodes
.get_mut(&NodeId(1))
.unwrap()
.core
.as_mut()
.unwrap();
let commit = core.commit_index();
let (snap, effects) = core.compact(commit).expect("compaction through commit");
assert_eq!(snap.last.index, commit);
assert!(
snap.claims.contains_key(&121),
"claim must ride the snapshot"
);
assert!(core.entry(orig_index).is_none(), "entry compacted away");
sim.run_effects(NodeId(1), effects, None);
}
let retry = sim.propose_keyed(NodeId(1), 121, 1210).expect("leader");
assert_eq!(
retry,
KeyedProposal::DuplicateCommitted { index: orig_index },
"compaction opened a claim replay window"
);
sim.check_all();
}
#[test]
fn compaction_refuses_uncommitted() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 103, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.cut.insert((NodeId(1), NodeId(2)));
sim.cut.insert((NodeId(1), NodeId(3)));
assert!(sim.propose(NodeId(1), 131));
sim.drain();
let core = sim
.nodes
.get_mut(&NodeId(1))
.unwrap()
.core
.as_mut()
.unwrap();
let commit = core.commit_index();
let tentative_tip = core.last_index();
assert!(tentative_tip > commit);
assert!(
core.compact(tentative_tip).is_none(),
"compacted an uncommitted entry"
);
assert!(
core.compact(commit).is_some(),
"committed compaction refused"
);
}
#[test]
fn straggler_beyond_gc_recovers_via_snapshot_install() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 107, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.crash(NodeId(3));
let out = sim.propose_keyed(NodeId(1), 141, 1410).expect("leader");
let keyed_index = match out {
KeyedProposal::Appended { index, .. } => index,
other => panic!("{other:?}"),
};
for p in [142, 143] {
assert!(sim.propose(NodeId(1), p));
}
sim.drain();
{
let core = sim
.nodes
.get_mut(&NodeId(1))
.unwrap()
.core
.as_mut()
.unwrap();
let commit = core.commit_index();
let (_snap, effects) = core.compact(commit).expect("compact");
sim.run_effects(NodeId(1), effects, None);
}
assert!(sim.propose(NodeId(1), 144));
sim.drain();
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)].core.as_ref().unwrap();
assert!(
n3.base().index >= keyed_index,
"straggler did not adopt the snapshot (base {:?})",
n3.base()
);
assert!(
sim.nodes[&NodeId(3)].applied >= keyed_index,
"adopted state not applied"
);
assert!(
sim.committed_at.values().any(|e| e.payload == 144),
"post-snapshot suffix replicated"
);
}
#[test]
fn stale_snapshot_never_regresses() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 109, false);
sim.timeout(NodeId(1), None);
sim.drain();
for p in [151, 152] {
assert!(sim.propose(NodeId(1), p));
}
sim.drain();
let n2_commit_before = sim.nodes[&NodeId(2)].core.as_ref().unwrap().commit_index();
let n2_last_before = sim.nodes[&NodeId(2)].core.as_ref().unwrap().last_index();
let effects = sim
.nodes
.get_mut(&NodeId(2))
.unwrap()
.core
.as_mut()
.unwrap()
.on_message(
NodeId(1),
Message::InstallSnapshot {
term: Term(1),
leader: NodeId(1),
snapshot: Snapshot {
last: LogPosition { term: 1, index: 1 },
claims: BTreeMap::new(),
active: 0,
},
},
false,
);
sim.run_effects(NodeId(2), effects, None);
sim.drain();
let n2 = sim.nodes[&NodeId(2)].core.as_ref().unwrap();
assert_eq!(n2.commit_index(), n2_commit_before, "commit regressed");
assert_eq!(n2.last_index(), n2_last_before, "log regressed");
sim.check_all();
}
#[test]
fn seeded_soak_with_compaction_invariants_hold() {
for seed in 1..12u64 {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, seed, false);
let keys = [61u64, 62, 63];
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(12) {
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(id, false);
}
}
3 | 4 => {
let id = ids[sim.rng.pick(3)];
let k = keys[sim.rng.pick(keys.len())];
if let Some(KeyedProposal::DuplicateCommitted { .. }) =
sim.propose_keyed(id, k, k * 10)
{
assert!(sim.keyed_committed.contains_key(&k));
}
}
5 => {
let id = ids[sim.rng.pick(3)];
let effects = sim
.nodes
.get_mut(&id)
.unwrap()
.core
.as_mut()
.and_then(|core| {
let c = core.commit_index();
core.compact(c).map(|(_s, e)| e)
});
if let Some(effects) = effects {
sim.run_effects(id, effects, None);
}
sim.heartbeat(id);
}
6 => {
let id = ids[sim.rng.pick(3)];
sim.heartbeat(id);
}
7 => {
let a = ids[sim.rng.pick(3)];
let b = ids[sim.rng.pick(3)];
if a != b {
let kk = (a.min(b), a.max(b));
if !sim.cut.remove(&kk) {
sim.cut.insert(kk);
}
}
}
_ => {
sim.deliver_one(None);
}
}
}
sim.cut.clear();
sim.drain();
sim.check_all();
}
}
impl Sim {
fn set_node_caps(&mut self, id: NodeId, supported: u32) {
if let Some(core) = self.nodes.get_mut(&id).unwrap().core.as_mut() {
core.set_supported(supported);
}
}
fn feed_peer_caps(&mut self, id: NodeId, caps: &[(u64, u32)]) {
if let Some(core) = self.nodes.get_mut(&id).unwrap().core.as_mut() {
core.set_peer_caps(caps.iter().map(|(n, c)| (NodeId(*n), *c)).collect());
}
}
fn propose_activation(&mut self, id: NodeId, bits: u32) -> bool {
let Some(core) = self.nodes.get_mut(&id).unwrap().core.as_mut() else {
return false;
};
match core.propose_activation(bits) {
Some(effects) => {
self.run_effects(id, effects, None);
true
}
None => false,
}
}
fn restart_via_bootstrap_with(&mut self, id: NodeId, supported: u32) {
let voters = self.voters.clone();
let w = self.witnesses.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()),
active: node.disk_active,
commit_marker: 0,
integrity: Integrity {
hard_state_verified: !node.torn_hard,
log_verified: !node.corrupt_log,
},
};
match inspect(CLUSTER, supported, &recovered) {
BootDecision::Healthy { hard, log } => {
let mut core = ReplicaCore::new(id, voters, hard, log, false);
core.set_witnesses(w);
core.set_supported(supported);
node.core = Some(core);
node.quarantined = None;
}
BootDecision::Quarantine { reasons, term_hint } => {
node.core = None;
node.quarantined = Some(QuarantinedNode::new(id, CLUSTER, reasons, term_hint));
}
}
}
}
const CAP_X: u32 = 0b0001;
#[test]
fn activation_requires_unanimous_support() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 113, false);
sim.timeout(NodeId(1), None);
sim.drain();
assert!(!sim.propose_activation(NodeId(1), CAP_X));
sim.feed_peer_caps(NodeId(1), &[(2, CAP_X), (3, 0)]);
assert!(!sim.propose_activation(NodeId(1), CAP_X));
sim.feed_peer_caps(NodeId(1), &[(2, CAP_X), (3, CAP_X)]);
assert!(sim.propose_activation(NodeId(1), CAP_X));
sim.drain();
sim.check_all();
for id in [1u64, 2, 3] {
assert_eq!(
sim.nodes[&NodeId(id)].core.as_ref().unwrap().active_caps(),
CAP_X,
"node {id} did not activate on commit"
);
}
}
#[test]
fn stale_advertisement_stalls_incompatible_follower_visibly() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 127, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.set_node_caps(NodeId(3), 0);
sim.feed_peer_caps(NodeId(1), &[(2, CAP_X), (3, CAP_X)]);
assert!(sim.propose_activation(NodeId(1), CAP_X));
sim.drain();
assert_eq!(
sim.nodes[&NodeId(1)].core.as_ref().unwrap().active_caps(),
CAP_X
);
let n3 = sim.nodes[&NodeId(3)].core.as_ref().unwrap();
assert_eq!(n3.active_caps(), 0);
assert_eq!(
sim.incompat_alarms,
vec![(NodeId(1), NodeId(3))],
"expected exactly one alarm"
);
sim.heartbeat(NodeId(1));
sim.drain();
assert_eq!(sim.incompat_alarms.len(), 1);
sim.crash(NodeId(1));
sim.timeout(NodeId(3), None);
sim.drain();
assert!(
sim.leaders.values().all(|s| !s.contains(&NodeId(3))),
"incompatible node won an election"
);
sim.check_all();
}
#[test]
fn codex_f2_downgraded_node_with_retained_activation_quarantines() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 131, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.feed_peer_caps(NodeId(1), &[(2, CAP_X), (3, CAP_X)]);
assert!(sim.propose_activation(NodeId(1), CAP_X));
sim.drain();
sim.crash(NodeId(2));
sim.restart_via_bootstrap_with(NodeId(2), 0);
let n2 = &sim.nodes[&NodeId(2)];
let q = n2
.quarantined
.as_ref()
.expect("downgraded node must quarantine");
assert!(q
.reasons()
.contains(&QuarantineReason::UnsupportedCapability));
sim.restart_via_bootstrap_with(NodeId(2), CAP_X);
assert!(sim.nodes[&NodeId(2)].quarantined.is_none());
sim.heartbeat(NodeId(1));
sim.drain();
sim.check_all();
}
#[test]
fn unsupported_snapshot_install_refused() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 137, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.set_node_caps(NodeId(2), 0); let before_base = sim.nodes[&NodeId(2)].core.as_ref().unwrap().base();
let effects = sim
.nodes
.get_mut(&NodeId(2))
.unwrap()
.core
.as_mut()
.unwrap()
.on_message(
NodeId(1),
Message::InstallSnapshot {
term: Term(1),
leader: NodeId(1),
snapshot: Snapshot {
last: LogPosition { term: 1, index: 9 },
claims: BTreeMap::new(),
active: CAP_X,
},
},
false,
);
sim.run_effects(NodeId(2), effects, None);
let n2 = sim.nodes[&NodeId(2)].core.as_ref().unwrap();
assert_eq!(n2.base(), before_base, "unsupported snapshot was adopted");
assert_eq!(n2.active_caps(), 0);
sim.check_all();
}
#[test]
fn activation_survives_compaction_and_restart() {
let mut sim = Sim::new(&[(1, empty()), (2, empty()), (3, empty())], 0, 139, false);
sim.timeout(NodeId(1), None);
sim.drain();
sim.feed_peer_caps(NodeId(1), &[(2, CAP_X), (3, CAP_X)]);
assert!(sim.propose_activation(NodeId(1), CAP_X));
sim.drain();
{
let core = sim
.nodes
.get_mut(&NodeId(1))
.unwrap()
.core
.as_mut()
.unwrap();
let c = core.commit_index();
let (_snap, effects) = core.compact(c).expect("compact");
sim.run_effects(NodeId(1), effects, None);
}
sim.heartbeat(NodeId(1));
sim.drain();
sim.crash(NodeId(1));
sim.restart(NodeId(1), false);
assert_eq!(
sim.nodes[&NodeId(1)].core.as_ref().unwrap().active_caps(),
CAP_X,
"active lost across compaction + restart"
);
sim.check_all();
}