use std::collections::{BTreeMap, BTreeSet};
use serde::{Deserialize, Serialize};
use super::types::{quorum, HardState, LogPosition, NodeId, Term};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum Payload {
Noop,
Test(u64),
Op(Vec<u8>),
Control(Vec<u8>),
}
#[cfg(test)]
impl PartialEq<u64> for Payload {
fn eq(&self, other: &u64) -> bool {
matches!(self, Payload::Test(n) if n == other)
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LogEntry {
pub term: Term,
pub payload: Payload,
pub key: Option<u64>,
pub activate: Option<u32>,
}
impl LogEntry {
pub fn unkeyed(term: Term, payload: Payload) -> Self {
Self {
term,
payload,
key: None,
activate: None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum KeyedProposal {
Appended { index: u64, effects: Vec<Effect> },
DuplicateCommitted { index: u64 },
DuplicatePending { index: u64 },
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct Snapshot {
pub last: LogPosition,
pub claims: BTreeMap<u64, u64>,
pub active: u32,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum Message {
PreVoteRequest {
term: Term,
candidate: NodeId,
last_log: LogPosition,
},
PreVoteResponse {
term: Term,
granted: bool,
},
VoteRequest {
term: Term,
candidate: NodeId,
last_log: LogPosition,
supported: u32,
},
VoteResponse {
term: Term,
granted: bool,
},
AppendEntries {
term: Term,
leader: NodeId,
prev: LogPosition,
entries: Vec<LogEntry>,
commit: u64,
},
InstallSnapshot {
term: Term,
leader: NodeId,
snapshot: Snapshot,
},
AppendResponse {
term: Term,
success: bool,
last_index: u64,
unsupported: bool,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Effect {
Persist {
hard: HardState,
base: LogPosition,
log: Vec<LogEntry>,
claims: BTreeMap<u64, u64>,
active: u32,
},
Send {
to: NodeId,
msg: Message,
},
Broadcast {
msg: Message,
},
BecameLeader {
term: Term,
},
SteppedDown {
term: Term,
},
CommitAdvanced {
to: u64,
},
InstallState {
last_index: u64,
},
PeerIncompatible {
peer: NodeId,
},
}
#[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,
base: LogPosition,
log: Vec<LogEntry>,
pending: Option<Pending>,
role: Role,
commit: u64,
match_index: BTreeMap<NodeId, u64>,
next_index: BTreeMap<NodeId, u64>,
claims: BTreeMap<u64, u64>,
votes: BTreeSet<NodeId>,
pre_votes: BTreeSet<NodeId>,
use_pre_vote: bool,
witnesses: BTreeSet<NodeId>,
supported: u32,
active: u32,
peer_caps: BTreeMap<NodeId, u32>,
stalled: BTreeSet<NodeId>,
}
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;
let claims = restored_log
.iter()
.enumerate()
.filter_map(|(i, e)| e.key.map(|k| (k, i as u64 + 1)))
.collect();
Self {
id,
voters,
persisted_hard: restored_hard,
durable_len,
base: LogPosition::ZERO,
log: restored_log,
pending: None,
role: Role::Follower,
commit: 0,
match_index: BTreeMap::new(),
next_index: BTreeMap::new(),
claims,
votes: BTreeSet::new(),
pre_votes: BTreeSet::new(),
use_pre_vote,
witnesses: BTreeSet::new(),
supported: u32::MAX,
active: 0,
peer_caps: BTreeMap::new(),
stalled: BTreeSet::new(),
}
}
pub fn set_supported(&mut self, supported: u32) {
self.supported = supported;
}
pub fn set_peer_caps(&mut self, caps: BTreeMap<NodeId, u32>) {
self.peer_caps = caps;
self.stalled.clear();
}
pub fn active_caps(&self) -> u32 {
self.active
}
pub fn propose_activation(&mut self, bits: u32) -> Option<Vec<Effect>> {
if self.role != Role::Leader || bits == 0 {
return None;
}
if self.supported & bits != bits {
return None;
}
for v in self.voters.iter() {
if *v == self.id {
continue;
}
let caps = self.peer_caps.get(v).copied().unwrap_or(0);
if caps & bits != bits {
return None;
}
}
let term = self.current_term();
self.log.push(LogEntry {
term,
payload: Payload::Noop,
key: None,
activate: Some(bits),
});
let mut out = vec![self.stage_log()];
self.append_effects_for_followers(&mut out);
Some(out)
}
fn apply_committed_caps(&mut self, from: u64, to: u64) {
for i in from.max(self.base.index + 1)..=to {
if let Some(bits) = self.entry(i).and_then(|e| e.activate) {
self.active |= bits;
}
}
}
pub fn set_witnesses(&mut self, witnesses: BTreeSet<NodeId>) {
debug_assert!(witnesses.iter().all(|w| self.voters.contains(w)));
debug_assert!(
witnesses.len() <= 1,
"multi-witness topologies break election/data quorum intersection"
);
debug_assert!(
!witnesses.contains(&self.id) || self.role == Role::Follower,
"a witness cannot already hold a data role"
);
self.witnesses = witnesses;
}
#[allow(clippy::too_many_arguments)]
pub fn new_from_durable(
id: NodeId,
voters: BTreeSet<NodeId>,
restored_hard: HardState,
base: LogPosition,
restored_log: Vec<LogEntry>,
snapshot_claims: BTreeMap<u64, u64>,
snapshot_active: u32,
use_pre_vote: bool,
) -> Self {
let mut core = Self::new(id, voters, restored_hard, Vec::new(), use_pre_vote);
core.base = base;
core.claims = snapshot_claims;
core.active = snapshot_active;
core.commit = base.index;
for e in restored_log {
let key = e.key;
core.log.push(e);
if let Some(k) = key {
core.claims.insert(k, core.last_index());
}
}
core.durable_len = core.log.len() as u64;
core
}
pub fn last_index(&self) -> u64 {
self.base.index + self.log.len() as u64
}
pub fn base(&self) -> LogPosition {
self.base
}
pub fn compact(&mut self, up_to: u64) -> Option<(Snapshot, Vec<Effect>)> {
if up_to <= self.base.index || up_to > self.commit {
return None;
}
let last_term = self.entry(up_to)?.term;
let drop = (up_to - self.base.index) as usize;
self.log.drain(..drop);
self.base = LogPosition {
term: last_term.0,
index: up_to,
};
self.durable_len = self.durable_len.saturating_sub(drop as u64);
let persist = self.stage_log();
Some((
Snapshot {
last: self.base,
claims: self.claims.clone(),
active: self.active,
},
vec![persist],
))
}
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 <= self.base.index {
return None; }
self.log.get((index - self.base.index) as usize - 1)
}
pub fn last_log(&self) -> LogPosition {
match self.log.last() {
Some(e) => LogPosition {
term: e.term.0,
index: self.last_index(),
},
None => self.base,
}
}
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,
base: self.base,
log: self.log.clone(),
claims: self.claims.clone(),
active: self.active,
}
}
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.base.index + self.durable_len);
self.try_advance_commit(&mut out);
}
out
}
pub fn propose(&mut self, payload: Payload) -> Option<Vec<Effect>> {
if self.role != Role::Leader {
return None;
}
let term = self.current_term();
self.log.push(LogEntry::unkeyed(term, payload));
let mut out = vec![self.stage_log()];
self.append_effects_for_followers(&mut out);
Some(out)
}
pub fn propose_keyed(&mut self, key: u64, payload: Payload) -> Option<KeyedProposal> {
if self.role != Role::Leader {
return None;
}
if let Some(&index) = self.claims.get(&key) {
return Some(if index <= self.commit {
KeyedProposal::DuplicateCommitted { index }
} else {
KeyedProposal::DuplicatePending { index }
});
}
let term = self.current_term();
self.log.push(LogEntry {
term,
payload,
key: Some(key),
activate: None,
});
let index = self.last_index();
self.claims.insert(key, index);
let mut effects = vec![self.stage_log()];
self.append_effects_for_followers(&mut effects);
Some(KeyedProposal::Appended { index, effects })
}
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>) {
if self.stalled.contains(&peer) {
return; }
let term = self.current_term();
let next = *self.next_index.get(&peer).unwrap_or(&1);
if next <= self.base.index {
let msg = Message::InstallSnapshot {
term,
leader: self.id,
snapshot: Snapshot {
last: self.base,
claims: self.claims.clone(),
active: self.active,
},
};
self.hold_or(Effect::Send { to: peer, msg }, out);
return;
}
let prev_index = next.saturating_sub(1);
let prev = if prev_index == self.base.index {
self.base } else {
match self.entry(prev_index) {
Some(e) => LogPosition {
term: e.term.0,
index: prev_index,
},
None => self.base,
}
};
let entries: Vec<LogEntry> = self
.log
.iter()
.skip((prev.index - self.base.index) as usize)
.cloned()
.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> {
if self.witnesses.contains(&self.id) {
return Vec::new();
}
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,
supported,
} => self.on_vote_request(from, term, candidate, last_log, supported),
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,
unsupported,
} => self.on_append_response(from, term, success, last_index, unsupported),
Message::InstallSnapshot {
term,
leader: _,
snapshot,
} => self.on_install_snapshot(from, term, snapshot),
}
}
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(),
supported: self.supported,
},
});
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,
candidate_supported: u32,
) -> 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());
let capable = candidate_supported & self.active == self.active;
if may_vote && fresh && capable {
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.base.index + self.durable_len,
unsupported: false,
},
}];
}
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()
};
if prev.index < self.base.index {
let ack = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: true,
last_index: self.base.index,
unsupported: false,
},
};
if newer {
effects.push(self.stage(adopted));
self.hold(ack);
} else {
self.hold_or(ack, &mut effects);
}
return effects;
}
let prev_ok = (prev.index == self.base.index && prev.term == self.base.term)
|| 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.base.index + self.durable_len)
.min(prev.index.saturating_sub(1)),
unsupported: false,
},
};
if newer {
effects.push(self.stage(adopted));
self.hold(refusal);
} else {
self.hold_or(refusal, &mut effects);
}
return effects;
}
if entries
.iter()
.any(|e| e.activate.is_some_and(|b| self.supported & b != b))
{
let refusal = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: false,
last_index: self.base.index + self.durable_len,
unsupported: true,
},
};
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.base.index + self.durable_len,
unsupported: false,
},
};
if newer {
effects.push(self.stage(adopted));
self.hold(refusal);
} else {
self.hold_or(refusal, &mut effects);
}
return effects;
}
for removed in self.log.iter().skip((idx - self.base.index) as usize - 1) {
if let Some(k) = removed.key {
if self.claims.get(&k) >= Some(&idx) {
self.claims.remove(&k);
}
}
}
self.log.truncate((idx - self.base.index) as usize - 1);
self.log.push(e.clone());
if let Some(k) = e.key {
self.claims.insert(k, idx);
}
changed = true;
}
None => {
self.log.push(e.clone());
if let Some(k) = e.key {
self.claims.insert(k, idx);
}
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.base.index + self.durable_len).min(covered)
},
unsupported: false,
},
});
} else {
let ack = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: true,
last_index: (self.base.index + self.durable_len).min(covered),
unsupported: false,
},
};
self.hold_or(ack, &mut effects);
}
let new_commit = leader_commit.min(covered);
if new_commit > self.commit {
let from_idx = self.commit + 1;
self.commit = new_commit;
self.apply_committed_caps(from_idx, 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,
unsupported: bool,
) -> 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 && unsupported {
if self.stalled.insert(from) {
out.push(Effect::PeerIncompatible { peer: from });
}
return out;
}
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.last_index() + 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 on_install_snapshot(&mut self, from: NodeId, term: Term, snapshot: Snapshot) -> 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.base.index + self.durable_len,
unsupported: false,
},
}];
}
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()
};
if self.supported & snapshot.active != snapshot.active {
let refusal = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: false,
last_index: self.base.index + self.durable_len,
unsupported: true,
},
};
if newer {
effects.push(self.stage(adopted));
self.hold(refusal);
} else {
self.hold_or(refusal, &mut effects);
}
return effects;
}
if snapshot.active & self.active != self.active {
let refusal = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: false,
last_index: self.base.index + self.durable_len,
unsupported: false,
},
};
if newer {
effects.push(self.stage(adopted));
self.hold(refusal);
} else {
self.hold_or(refusal, &mut effects);
}
return effects;
}
if snapshot.last.index <= self.commit {
let ack = Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: true,
last_index: self.base.index + self.durable_len,
unsupported: false,
},
};
if newer {
effects.push(self.stage(adopted));
self.hold(ack);
} else {
self.hold_or(ack, &mut effects);
}
return effects;
}
let last_index = snapshot.last.index;
self.base = snapshot.last;
self.log.clear();
self.claims = snapshot.claims;
self.active |= snapshot.active;
self.commit = last_index;
effects.push(self.stage(adopted));
self.hold(Effect::InstallState { last_index });
self.hold(Effect::Send {
to: from,
msg: Message::AppendResponse {
term,
success: true,
last_index,
unsupported: false,
},
});
effects
}
fn try_advance_commit(&mut self, out: &mut Vec<Effect>) {
let cur = self.current_term();
let mut n = self.last_index();
while n > self.commit {
if self.entry(n).map(|e| e.term) == Some(cur) {
let data_voters = || self.voters.iter().filter(|v| !self.witnesses.contains(v));
let acks = data_voters()
.filter(|v| self.match_index.get(v).copied().unwrap_or(0) >= n)
.count();
if acks >= quorum(data_voters().count()) {
let from_idx = self.commit + 1;
self.commit = n;
self.apply_committed_caps(from_idx, 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::unkeyed(term, Payload::Noop));
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::unkeyed(term, Payload::Noop));
let _ = self.stage_log(); }
fn init_leader_state(&mut self, _term: Term) {
self.role = Role::Leader;
let last = self.last_index();
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.base.index + 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
}
#[allow(clippy::type_complexity)]
pub fn rejoin_grant(
&self,
) -> Option<(
Term,
LogPosition,
Vec<LogEntry>,
BTreeMap<u64, u64>,
u32,
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)
.cloned()
.collect();
let commit = self.commit.min(self.base.index + self.durable_len);
Some((
cur,
self.base,
durable,
self.claims.clone(),
self.active,
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())
}
}