use std::collections::{BTreeMap, BTreeSet};
use super::types::{quorum, HardState, LogPosition, NodeId, Term};
pub const NOOP_PAYLOAD: u64 = 0;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LogEntry {
pub term: Term,
pub payload: u64,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Message {
PreVoteRequest {
term: Term,
candidate: NodeId,
last_log: LogPosition,
},
PreVoteResponse {
term: Term,
granted: bool,
},
VoteRequest {
term: Term,
candidate: NodeId,
last_log: LogPosition,
},
VoteResponse {
term: Term,
granted: bool,
},
AppendEntries {
term: Term,
leader: NodeId,
prev: LogPosition,
entries: Vec<LogEntry>,
commit: u64,
},
AppendResponse {
term: Term,
success: bool,
last_index: u64,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Effect {
Persist {
hard: HardState,
log: Vec<LogEntry>,
},
Send {
to: NodeId,
msg: Message,
},
Broadcast {
msg: Message,
},
BecameLeader {
term: Term,
},
SteppedDown {
term: Term,
},
CommitAdvanced {
to: u64,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Role {
Follower,
PreCandidate,
Candidate,
Leader,
}
#[derive(Debug, Clone)]
struct Pending {
hard: HardState,
held: Vec<Effect>,
}
pub struct ReplicaCore {
id: NodeId,
voters: BTreeSet<NodeId>,
persisted_hard: HardState,
durable_len: u64,
log: Vec<LogEntry>,
pending: Option<Pending>,
role: Role,
commit: u64,
match_index: BTreeMap<NodeId, u64>,
next_index: BTreeMap<NodeId, u64>,
votes: BTreeSet<NodeId>,
pre_votes: BTreeSet<NodeId>,
use_pre_vote: bool,
}
impl ReplicaCore {
pub fn new(
id: NodeId,
voters: BTreeSet<NodeId>,
restored_hard: HardState,
restored_log: Vec<LogEntry>,
use_pre_vote: bool,
) -> Self {
debug_assert!(voters.contains(&id), "a core's own id must be a voter");
let durable_len = restored_log.len() as u64;
Self {
id,
voters,
persisted_hard: restored_hard,
durable_len,
log: restored_log,
pending: None,
role: Role::Follower,
commit: 0,
match_index: BTreeMap::new(),
next_index: BTreeMap::new(),
votes: BTreeSet::new(),
pre_votes: BTreeSet::new(),
use_pre_vote,
}
}
pub fn role(&self) -> Role {
self.role
}
pub fn current_term(&self) -> Term {
self.effective().current_term
}
pub fn persisted_hard_state(&self) -> HardState {
self.persisted_hard
}
pub fn commit_index(&self) -> u64 {
self.commit
}
pub fn next_index_of(&self, peer: NodeId) -> Option<u64> {
self.next_index.get(&peer).copied()
}
pub fn log_len(&self) -> u64 {
self.log.len() as u64
}
pub fn entry(&self, index: u64) -> Option<&LogEntry> {
if index == 0 {
return None;
}
self.log.get(index as usize - 1)
}
pub fn last_log(&self) -> LogPosition {
match self.log.last() {
Some(e) => LogPosition {
term: e.term.0,
index: self.log.len() as u64,
},
None => LogPosition::ZERO,
}
}
fn effective(&self) -> HardState {
self.pending
.as_ref()
.map(|p| p.hard)
.unwrap_or(self.persisted_hard)
}
fn stage(&mut self, hard: HardState) -> Effect {
match &mut self.pending {
Some(p) => p.hard = hard,
None => {
self.pending = Some(Pending {
hard,
held: Vec::new(),
})
}
}
Effect::Persist {
hard,
log: self.log.clone(),
}
}
fn stage_log(&mut self) -> Effect {
let hard = self.effective();
self.stage(hard)
}
fn hold(&mut self, eff: Effect) {
debug_assert!(self.pending.is_some(), "hold() requires a staged persist");
if let Some(p) = &mut self.pending {
p.held.push(eff);
}
}
fn hold_or(&mut self, eff: Effect, out: &mut Vec<Effect>) {
if self.pending.is_some() {
self.hold(eff);
} else {
out.push(eff);
}
}
pub fn state_persisted(&mut self) -> Vec<Effect> {
let mut out = match self.pending.take() {
Some(p) => {
self.persisted_hard = p.hard;
self.durable_len = self.log.len() as u64;
p.held
}
None => return Vec::new(),
};
if self.role == Role::Leader {
self.match_index.insert(self.id, self.durable_len);
self.try_advance_commit(&mut out);
}
out
}
pub fn propose(&mut self, payload: u64) -> Option<Vec<Effect>> {
if self.role != Role::Leader {
return None;
}
let term = self.current_term();
self.log.push(LogEntry { term, payload });
let mut out = vec![self.stage_log()];
self.append_effects_for_followers(&mut out);
Some(out)
}
pub fn tick_heartbeat(&mut self) -> Vec<Effect> {
if self.role != Role::Leader {
return Vec::new();
}
let mut out = Vec::new();
self.append_effects_for_followers(&mut out);
out
}
fn append_effects_for_followers(&mut self, out: &mut Vec<Effect>) {
let peers: Vec<NodeId> = self
.voters
.iter()
.copied()
.filter(|p| *p != self.id)
.collect();
for peer in peers {
self.send_append_to(peer, out);
}
}
fn send_append_to(&mut self, peer: NodeId, out: &mut Vec<Effect>) {
let term = self.current_term();
let next = *self.next_index.get(&peer).unwrap_or(&1);
let prev_index = next.saturating_sub(1);
let prev = if prev_index == 0 {
LogPosition::ZERO
} else {
match self.entry(prev_index) {
Some(e) => LogPosition {
term: e.term.0,
index: prev_index,
},
None => LogPosition::ZERO,
}
};
let entries: Vec<LogEntry> = self.log.iter().skip(prev_index as usize).copied().collect();
let msg = Message::AppendEntries {
term,
leader: self.id,
prev,
entries,
commit: self.commit,
};
self.hold_or(Effect::Send { to: peer, msg }, out);
}
pub fn on_election_timeout(&mut self) -> Vec<Effect> {
match self.role {
Role::Leader => Vec::new(),
_ if self.use_pre_vote => self.start_pre_vote(),
_ => self.start_campaign(),
}
}
pub fn on_message(
&mut self,
from: NodeId,
msg: Message,
leader_recently_seen: bool,
) -> Vec<Effect> {
match msg {
Message::PreVoteRequest {
term,
candidate,
last_log,
} => self.on_pre_vote_request(from, term, candidate, last_log, leader_recently_seen),
Message::PreVoteResponse { term, granted } => {
self.on_pre_vote_response(from, term, granted)
}
Message::VoteRequest {
term,
candidate,
last_log,
} => self.on_vote_request(from, term, candidate, last_log),
Message::VoteResponse { term, granted } => self.on_vote_response(from, term, granted),
Message::AppendEntries {
term,
leader,
prev,
entries,
commit,
} => self.on_append_entries(from, term, leader, prev, entries, commit),
Message::AppendResponse {
term,
success,
last_index,
} => self.on_append_response(from, term, success, last_index),
}
}
fn start_pre_vote(&mut self) -> Vec<Effect> {
self.role = Role::PreCandidate;
self.pre_votes.clear();
self.pre_votes.insert(self.id);
if self.pre_quorum_reached() {
return self.start_campaign();
}
vec![Effect::Broadcast {
msg: Message::PreVoteRequest {
term: self.current_term().next(),
candidate: self.id,
last_log: self.last_log(),
},
}]
}
fn start_campaign(&mut self) -> Vec<Effect> {
self.role = Role::Candidate;
self.votes.clear();
self.votes.insert(self.id);
let term = self.current_term().next();
let persist = self.stage(HardState {
current_term: term,
voted_for: Some(self.id),
});
if self.vote_quorum_reached() {
self.become_leader_held(term);
return vec![persist];
}
self.hold(Effect::Broadcast {
msg: Message::VoteRequest {
term,
candidate: self.id,
last_log: self.last_log(),
},
});
vec![persist]
}
fn on_pre_vote_request(
&mut self,
from: NodeId,
term: Term,
_candidate: NodeId,
last_log: LogPosition,
leader_recently_seen: bool,
) -> Vec<Effect> {
let grant = !leader_recently_seen
&& term > self.current_term()
&& last_log.is_at_least_as_up_to_date_as(&self.last_log());
vec![Effect::Send {
to: from,
msg: Message::PreVoteResponse {
term,
granted: grant,
},
}]
}
fn on_pre_vote_response(&mut self, from: NodeId, term: Term, granted: bool) -> Vec<Effect> {
if self.role != Role::PreCandidate || term != self.current_term().next() || !granted {
return Vec::new();
}
if self.voters.contains(&from) {
self.pre_votes.insert(from);
}
if self.pre_quorum_reached() {
return self.start_campaign();
}
Vec::new()
}
fn on_vote_request(
&mut self,
from: NodeId,
term: Term,
candidate: NodeId,
last_log: LogPosition,
) -> Vec<Effect> {
let mut effects = Vec::new();
let cur = self.current_term();
if term < cur {
return vec![Effect::Send {
to: from,
msg: Message::VoteResponse {
term: cur,
granted: false,
},
}];
}
let newer = term > cur;
if newer {
if matches!(self.role, Role::Leader | Role::Candidate) {
effects.push(Effect::SteppedDown { term: cur });
}
self.role = Role::Follower;
}
let prior_vote = if newer {
None
} else {
self.effective().voted_for
};
let may_vote = prior_vote.is_none() || prior_vote == Some(candidate);
let fresh = last_log.is_at_least_as_up_to_date_as(&self.last_log());
if may_vote && fresh {
effects.push(self.stage(HardState {
current_term: term,
voted_for: Some(candidate),
}));
self.hold(Effect::Send {
to: from,
msg: Message::VoteResponse {
term,
granted: true,
},
});
} else {
let refusal = Effect::Send {
to: from,
msg: Message::VoteResponse {
term,
granted: false,
},
};
if newer {
effects.push(self.stage(HardState {
current_term: term,
voted_for: None,
}));
self.hold(refusal);
} else {
self.hold_or(refusal, &mut effects);
}
}
effects
}
fn on_vote_response(&mut self, from: NodeId, term: Term, granted: bool) -> Vec<Effect> {
let cur = self.current_term();
if term > cur {
return self.step_down_to(term, cur);
}
if self.role != Role::Candidate || term != cur || !granted {
return Vec::new();
}
if self.voters.contains(&from) {
self.votes.insert(from);
}
if self.vote_quorum_reached() {
let mut out = Vec::new();
self.become_leader(cur, &mut out);
return out;
}
Vec::new()
}
fn on_append_entries(
&mut self,
from: NodeId,
term: Term,
_leader: NodeId,
prev: LogPosition,
entries: Vec<LogEntry>,
leader_commit: u64,
) -> Vec<Effect> {
let mut effects = Vec::new();
let cur = self.current_term();
if term < cur {
return vec![Effect::Send {
to: from,
msg: Message::AppendResponse {
term: cur,
success: false,
last_index: self.durable_len,
},
}];
}
if matches!(self.role, Role::Leader | Role::Candidate) {
effects.push(Effect::SteppedDown { term: cur });
}
self.role = Role::Follower;
let newer = term > cur;
let adopted = if newer {
HardState {
current_term: term,
voted_for: None,
}
} else {
self.effective()
};
let prev_ok = prev.index == 0
|| self
.entry(prev.index)
.is_some_and(|e| e.term.0 == prev.term);
if !prev_ok {
let refusal = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: false,
last_index: self.durable_len.min(prev.index.saturating_sub(1)),
},
};
if newer {
effects.push(self.stage(adopted));
self.hold(refusal);
} else {
self.hold_or(refusal, &mut effects);
}
return effects;
}
let mut idx = prev.index;
let mut changed = false;
for e in &entries {
idx += 1;
match self.entry(idx) {
Some(existing) if *existing == *e => {}
Some(_) => {
if idx <= self.commit {
debug_assert!(false, "conflict at/below commit index — corruption");
let refusal = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: false,
last_index: self.durable_len,
},
};
if newer {
effects.push(self.stage(adopted));
self.hold(refusal);
} else {
self.hold_or(refusal, &mut effects);
}
return effects;
}
self.log.truncate(idx as usize - 1);
self.log.push(*e);
changed = true;
}
None => {
self.log.push(*e);
changed = true;
}
}
}
let covered = idx.max(prev.index);
if changed || newer {
effects.push(self.stage(adopted));
self.hold(Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: true,
last_index: if changed {
covered
} else {
self.durable_len.min(covered)
},
},
});
} else {
let ack = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: true,
last_index: self.durable_len.min(covered),
},
};
self.hold_or(ack, &mut effects);
}
let new_commit = leader_commit.min(covered);
if new_commit > self.commit {
self.commit = new_commit;
let advanced = Effect::CommitAdvanced { to: new_commit };
self.hold_or(advanced, &mut effects);
}
effects
}
fn on_append_response(
&mut self,
from: NodeId,
term: Term,
success: bool,
last_index: u64,
) -> Vec<Effect> {
let cur = self.current_term();
if term > cur {
return self.step_down_to(term, cur);
}
if self.role != Role::Leader || term != cur {
return Vec::new();
}
let mut out = Vec::new();
if success {
let m = self.match_index.entry(from).or_insert(0);
*m = (*m).max(last_index);
let floor = *m + 1;
self.next_index.insert(from, floor);
self.try_advance_commit(&mut out);
} else {
let floor = self.match_index.get(&from).copied().unwrap_or(0) + 1;
let next = self
.next_index
.get(&from)
.copied()
.unwrap_or(self.log_len() + 1);
let backed = next.saturating_sub(1).clamp(1, last_index + 1).max(floor);
self.next_index.insert(from, backed);
self.send_append_to(from, &mut out);
}
out
}
fn try_advance_commit(&mut self, out: &mut Vec<Effect>) {
let cur = self.current_term();
let mut n = self.log_len();
while n > self.commit {
if self.entry(n).map(|e| e.term) == Some(cur) {
let acks = self
.voters
.iter()
.filter(|v| self.match_index.get(v).copied().unwrap_or(0) >= n)
.count();
if acks >= quorum(self.voters.len()) {
self.commit = n;
let advanced = Effect::CommitAdvanced { to: n };
self.hold_or(advanced, out);
self.append_effects_for_followers(out);
return;
}
}
n -= 1;
}
}
fn become_leader(&mut self, term: Term, out: &mut Vec<Effect>) {
self.init_leader_state(term);
out.push(Effect::BecameLeader { term });
self.log.push(LogEntry {
term,
payload: NOOP_PAYLOAD,
});
out.push(self.stage_log());
self.append_effects_for_followers(out);
}
fn become_leader_held(&mut self, term: Term) {
debug_assert!(self.pending.is_some());
self.init_leader_state(term);
self.hold(Effect::BecameLeader { term });
self.log.push(LogEntry {
term,
payload: NOOP_PAYLOAD,
});
let _ = self.stage_log(); }
fn init_leader_state(&mut self, _term: Term) {
self.role = Role::Leader;
let last = self.log_len();
self.next_index.clear();
self.match_index.clear();
for v in self.voters.iter().copied() {
self.next_index.insert(v, last + 1);
self.match_index.insert(v, 0);
}
self.match_index.insert(self.id, self.durable_len);
}
fn step_down_to(&mut self, newer: Term, cur: Term) -> Vec<Effect> {
let mut effects = Vec::new();
if matches!(self.role, Role::Leader | Role::Candidate) {
effects.push(Effect::SteppedDown { term: cur });
}
self.role = Role::Follower;
effects.push(self.stage(HardState {
current_term: newer,
voted_for: None,
}));
effects
}
pub fn rejoin_grant(&self) -> Option<(Term, Vec<LogEntry>, u64)> {
if self.role != Role::Leader {
return None;
}
let cur = self.current_term();
let committed_in_current_term =
self.commit > 0 && self.entry(self.commit).map(|e| e.term) == Some(cur);
if !committed_in_current_term {
return None;
}
let durable: Vec<LogEntry> = self
.log
.iter()
.take(self.durable_len as usize)
.copied()
.collect();
let commit = self.commit.min(self.durable_len);
Some((cur, durable, commit))
}
fn pre_quorum_reached(&self) -> bool {
self.pre_votes.len() >= quorum(self.voters.len())
}
fn vote_quorum_reached(&self) -> bool {
self.votes.len() >= quorum(self.voters.len())
}
}