use std::collections::{BTreeMap, BTreeSet};
use crafty_proto::{
AppendEntries, AppendEntriesReply, CatalogCommand, EntryPayload, InstallSnapshot,
InstallSnapshotReply, LogEntry, LogId, LogIndex, Membership, NodeId,
QueueAutoscalePolicyCommand, RaftRpc, RaftRpcReply, RequestVote, RequestVoteReply, Round,
SagaJournalCommand, Term, TwoPhaseAbortCommand, TwoPhaseJournalCommand, TwoPhasePrepareCommand,
};
use crate::config::Configuration;
use crate::failure_detector::{
AckWindowLiveness, FailureDetectorKind, PhiAccrualLiveness, ReachabilityConfig,
};
use crate::log::Log;
use crate::rng::Rng;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Role {
Follower,
PreCandidate,
Candidate,
Leader,
}
#[derive(Debug, Clone)]
pub struct Config {
pub election_timeout_min: u64,
pub election_timeout_max: u64,
pub heartbeat_interval: u64,
pub seed: u64,
pub reachability: ReachabilityConfig,
}
impl Default for Config {
fn default() -> Self {
Self {
election_timeout_min: 10,
election_timeout_max: 20,
heartbeat_interval: 3,
seed: 0,
reachability: ReachabilityConfig::default(),
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Committed {
pub index: LogIndex,
pub command: Vec<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)]
pub struct ReadId(pub u64);
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Output {
Send(NodeId, RaftRpc),
Reply(NodeId, RaftRpcReply),
Apply(Committed),
RoleChanged(Role),
ReadReady {
id: ReadId,
index: LogIndex,
},
ReadFailed {
id: ReadId,
},
LoadSnapshot {
index: LogIndex,
data: Vec<u8>,
},
CatalogApplied {
index: LogIndex,
command: CatalogCommand,
},
SagaJournalApplied {
index: LogIndex,
command: SagaJournalCommand,
},
TwoPhasePrepareApplied {
index: LogIndex,
command: TwoPhasePrepareCommand,
},
TwoPhaseAbortApplied {
index: LogIndex,
command: TwoPhaseAbortCommand,
},
TwoPhaseJournalApplied {
index: LogIndex,
command: TwoPhaseJournalCommand,
},
QueueAutoscalePolicyApplied {
index: LogIndex,
command: QueueAutoscalePolicyCommand,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct NotLeader {
pub leader: Option<NodeId>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MembershipError {
NotLeader {
leader: Option<NodeId>,
},
InProgress,
EmptyVoters,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CatalogProposeError {
NotLeader {
leader: Option<NodeId>,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Persist {
pub term: Term,
pub voted_for: Option<NodeId>,
pub hard_state_dirty: bool,
pub truncate_from: Option<LogIndex>,
pub entries: Vec<LogEntry>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SnapshotState {
pub last_included: LogId,
pub membership: Membership,
pub data: Vec<u8>,
}
#[derive(Debug, Clone)]
struct PendingRead {
id: ReadId,
index: LogIndex,
round: Round,
acks: BTreeSet<NodeId>,
}
#[derive(Debug, Clone)]
struct StoredSnapshot {
last_index: LogIndex,
last_term: Term,
membership: Membership,
data: Vec<u8>,
}
#[derive(Debug, Clone)]
pub struct RaftNode {
id: NodeId,
initial: Membership,
config: Config,
current_term: Term,
voted_for: Option<NodeId>,
log: Log,
persisted_term: Term,
persisted_vote: Option<NodeId>,
log_dirty_from: Option<LogIndex>,
role: Role,
leader_id: Option<NodeId>,
commit_index: LogIndex,
last_applied: LogIndex,
votes: BTreeSet<NodeId>,
next_index: BTreeMap<NodeId, LogIndex>,
match_index: BTreeMap<NodeId, LogIndex>,
sent_upper: BTreeMap<NodeId, LogIndex>,
heartbeat_round: Round,
pending_reads: Vec<PendingRead>,
snapshot: Option<StoredSnapshot>,
last_ack_clock: BTreeMap<NodeId, u64>,
ack_liveness: AckWindowLiveness,
phi_liveness: PhiAccrualLiveness,
lease_round: Round,
lease_round_clock: u64,
lease_acks: BTreeSet<NodeId>,
lease_expiry: u64,
elapsed: u64,
heartbeat_elapsed: u64,
election_timeout: u64,
logical_clock: u64,
rng: Rng,
outbox: Vec<Output>,
}
impl RaftNode {
#[must_use]
pub fn new(id: NodeId, members: impl IntoIterator<Item = NodeId>, config: Config) -> Self {
let mut voters: Vec<NodeId> = members.into_iter().collect();
voters.sort();
voters.dedup();
let membership = Membership {
voters,
voters_outgoing: Vec::new(),
learners: Vec::new(),
};
Self::with_membership(id, membership, config)
}
#[must_use]
pub fn with_membership(id: NodeId, membership: Membership, config: Config) -> Self {
let mut rng = Rng::new(config.seed ^ id.0 ^ 0x9E37_79B9_7F4A_7C15);
let election_timeout = rng.range(config.election_timeout_min, config.election_timeout_max);
let phi_threshold = config.reachability.phi_threshold;
Self {
id,
initial: membership,
config,
current_term: Term::ZERO,
voted_for: None,
log: Log::default(),
persisted_term: Term::ZERO,
persisted_vote: None,
log_dirty_from: None,
role: Role::Follower,
leader_id: None,
commit_index: LogIndex::ZERO,
last_applied: LogIndex::ZERO,
votes: BTreeSet::new(),
next_index: BTreeMap::new(),
match_index: BTreeMap::new(),
sent_upper: BTreeMap::new(),
heartbeat_round: Round::ZERO,
pending_reads: Vec::new(),
snapshot: None,
last_ack_clock: BTreeMap::new(),
ack_liveness: AckWindowLiveness::default(),
phi_liveness: PhiAccrualLiveness::new(phi_threshold),
lease_round: Round::ZERO,
lease_round_clock: 0,
lease_acks: BTreeSet::new(),
lease_expiry: 0,
elapsed: 0,
heartbeat_elapsed: 0,
election_timeout,
logical_clock: 0,
rng,
outbox: Vec::new(),
}
}
#[must_use]
pub fn restore(
id: NodeId,
members: impl IntoIterator<Item = NodeId>,
config: Config,
term: Term,
voted_for: Option<NodeId>,
entries: impl IntoIterator<Item = LogEntry>,
) -> Self {
let mut node = Self::new(id, members, config);
node.current_term = term;
node.voted_for = voted_for;
for entry in entries {
node.log.push_entry(entry);
}
node.persisted_term = term;
node.persisted_vote = voted_for;
node.log_dirty_from = None;
node
}
#[must_use]
pub fn restore_with_snapshot(
id: NodeId,
members: impl IntoIterator<Item = NodeId>,
config: Config,
term: Term,
voted_for: Option<NodeId>,
snapshot: SnapshotState,
entries: impl IntoIterator<Item = LogEntry>,
) -> Self {
let mut node = Self::new(id, members, config);
node.current_term = term;
node.voted_for = voted_for;
let last = snapshot.last_included;
node.log.install_snapshot(last.index, last.term);
node.snapshot = Some(StoredSnapshot {
last_index: last.index,
last_term: last.term,
membership: snapshot.membership,
data: snapshot.data,
});
for entry in entries {
node.log.push_entry(entry);
}
node.commit_index = last.index;
node.last_applied = last.index;
node.persisted_term = term;
node.persisted_vote = voted_for;
node.log_dirty_from = None;
node
}
#[must_use]
pub fn id(&self) -> NodeId {
self.id
}
#[must_use]
pub fn role(&self) -> Role {
self.role
}
#[must_use]
pub fn is_leader(&self) -> bool {
self.role == Role::Leader
}
#[must_use]
pub fn current_term(&self) -> Term {
self.current_term
}
#[must_use]
pub fn leader_id(&self) -> Option<NodeId> {
self.leader_id
}
#[must_use]
pub fn commit_index(&self) -> LogIndex {
self.commit_index
}
#[must_use]
pub fn last_applied(&self) -> LogIndex {
self.last_applied
}
#[must_use]
pub fn last_log_index(&self) -> LogIndex {
self.log.last_index()
}
#[must_use]
pub fn voted_for(&self) -> Option<NodeId> {
self.voted_for
}
#[must_use]
pub fn term_at(&self, idx: LogIndex) -> Option<Term> {
self.log.term_at(idx)
}
#[must_use]
pub(crate) fn configuration(&self) -> Configuration {
let membership = self
.log
.last_membership()
.map(|(_, m)| m)
.or_else(|| self.snapshot.as_ref().map(|s| &s.membership))
.unwrap_or(&self.initial);
Configuration::from_membership(membership)
}
#[must_use]
pub fn committed_membership(&self) -> crafty_proto::Membership {
self.configuration().to_membership()
}
#[must_use]
pub fn snapshot_index(&self) -> LogIndex {
self.log.snapshot_index()
}
#[must_use]
pub fn compactable_entries(&self) -> u64 {
self.last_applied
.0
.saturating_sub(self.log.snapshot_index().0)
}
#[must_use]
pub fn compactable_log_bytes(&self) -> u64 {
self.log.bytes_up_to(self.last_applied)
}
#[must_use]
pub fn stored_snapshot(&self) -> Option<SnapshotState> {
self.snapshot.as_ref().map(|s| SnapshotState {
last_included: LogId::new(s.last_term, s.last_index),
membership: s.membership.clone(),
data: s.data.clone(),
})
}
#[must_use]
pub fn log_entries_from(&self, from: LogIndex) -> Vec<LogEntry> {
self.log.entries_from(from).to_vec()
}
#[must_use]
pub fn voters(&self) -> Vec<NodeId> {
self.configuration().voters()
}
#[must_use]
pub fn is_joint(&self) -> bool {
self.configuration().is_joint()
}
#[must_use]
pub fn take_outputs(&mut self) -> Vec<Output> {
std::mem::take(&mut self.outbox)
}
#[must_use]
pub fn take_persist(&mut self) -> Option<Persist> {
let hard_state_dirty =
self.current_term != self.persisted_term || self.voted_for != self.persisted_vote;
let log_from = self.log_dirty_from.take();
if !hard_state_dirty && log_from.is_none() {
return None;
}
self.persisted_term = self.current_term;
self.persisted_vote = self.voted_for;
let (truncate_from, entries) = match log_from {
Some(from) => {
let from = LogIndex(from.0.max(self.log.snapshot_index().0 + 1));
(Some(from), self.log.entries_from(from).to_vec())
}
None => (None, Vec::new()),
};
Some(Persist {
term: self.current_term,
voted_for: self.voted_for,
hard_state_dirty,
truncate_from,
entries,
})
}
fn mark_log_dirty(&mut self, from: LogIndex) {
self.log_dirty_from = Some(match self.log_dirty_from {
Some(cur) if cur.0 <= from.0 => cur,
_ => from,
});
}
fn log_append(&mut self, term: Term, payload: EntryPayload) -> LogIndex {
let idx = self.log.append(term, payload);
self.mark_log_dirty(idx);
idx
}
fn log_push(&mut self, entry: LogEntry) {
let idx = entry.index;
self.log.push_entry(entry);
self.mark_log_dirty(idx);
}
fn log_truncate_from(&mut self, idx: LogIndex) {
self.log.truncate_from(idx);
self.mark_log_dirty(idx);
}
fn config_index(&self) -> LogIndex {
self.log
.last_membership()
.map_or(LogIndex::ZERO, |(idx, _)| idx)
}
fn is_voter(&self, id: NodeId) -> bool {
self.configuration().is_voter(id)
}
fn peers(&self) -> Vec<NodeId> {
self.configuration().peers(self.id)
}
fn quorum_ok(&self, acked: &BTreeSet<NodeId>) -> bool {
self.configuration().has_quorum(acked)
}
fn quorum_of_votes(&self) -> bool {
let votes = self.votes.clone();
self.quorum_ok(&votes)
}
pub fn tick(&mut self) {
self.logical_clock += 1;
if self.role == Role::Leader {
self.update_liveness();
self.heartbeat_elapsed += 1;
if self.heartbeat_elapsed >= self.config.heartbeat_interval {
self.heartbeat_elapsed = 0;
self.broadcast_append();
}
} else {
self.elapsed += 1;
if self.elapsed >= self.election_timeout {
self.start_pre_election();
}
}
}
pub fn campaign(&mut self) {
self.start_real_election();
}
pub fn receive(&mut self, from: NodeId, rpc: RaftRpc) {
match rpc {
RaftRpc::RequestVote(rv) => self.handle_request_vote(from, &rv),
RaftRpc::AppendEntries(ae) => self.handle_append_entries(from, &ae),
RaftRpc::InstallSnapshot(is) => self.handle_install_snapshot(from, is),
}
}
pub fn receive_reply(&mut self, from: NodeId, reply: RaftRpcReply) {
let term = match &reply {
RaftRpcReply::RequestVote(r) => r.term,
RaftRpcReply::AppendEntries(r) => r.term,
RaftRpcReply::InstallSnapshot(r) => r.term,
};
if term > self.current_term {
self.become_follower(term);
return;
}
match reply {
RaftRpcReply::RequestVote(r) => self.handle_vote_reply(from, &r),
RaftRpcReply::AppendEntries(r) => self.handle_append_reply(from, &r),
RaftRpcReply::InstallSnapshot(r) => self.handle_snapshot_reply(from, &r),
}
}
pub fn propose(&mut self, command: Vec<u8>) -> Result<LogIndex, NotLeader> {
if self.role != Role::Leader {
return Err(NotLeader {
leader: self.leader_id,
});
}
let idx = self.log_append(self.current_term, EntryPayload::Command(command));
self.broadcast_append();
self.maybe_advance_commit();
Ok(idx)
}
pub fn read_index(&mut self, id: ReadId) -> Result<(), NotLeader> {
if self.role != Role::Leader {
return Err(NotLeader {
leader: self.leader_id,
});
}
self.broadcast_append();
let round = self.heartbeat_round;
let mut acks = BTreeSet::new();
acks.insert(self.id);
self.pending_reads.push(PendingRead {
id,
index: self.commit_index,
round,
acks,
});
self.try_complete_reads();
Ok(())
}
#[must_use]
pub fn compact(&mut self, up_to: LogIndex, data: Vec<u8>) -> bool {
if up_to <= self.log.snapshot_index() || up_to > self.last_applied {
return false;
}
let Some(term) = self.log.term_at(up_to) else {
return false;
};
let membership = self.membership_at(up_to);
self.log.compact(up_to, term);
self.snapshot = Some(StoredSnapshot {
last_index: up_to,
last_term: term,
membership,
data,
});
true
}
fn membership_at(&self, idx: LogIndex) -> Membership {
for i in (self.log.snapshot_index().0 + 1..=idx.0).rev() {
if let Some(EntryPayload::Membership(m)) = self.log.get(LogIndex(i)).map(|e| &e.payload)
{
return m.clone();
}
}
self.snapshot
.as_ref()
.map_or_else(|| self.initial.clone(), |s| s.membership.clone())
}
pub fn propose_membership(
&mut self,
new_voters: impl IntoIterator<Item = NodeId>,
learners: impl IntoIterator<Item = NodeId>,
) -> Result<LogIndex, MembershipError> {
if self.role != Role::Leader {
return Err(MembershipError::NotLeader {
leader: self.leader_id,
});
}
let current = self.configuration();
if current.is_joint() || self.config_index() > self.commit_index {
return Err(MembershipError::InProgress);
}
let mut voters: Vec<NodeId> = new_voters.into_iter().collect();
voters.sort();
voters.dedup();
if voters.is_empty() {
return Err(MembershipError::EmptyVoters);
}
let mut learners: Vec<NodeId> = learners.into_iter().collect();
learners.sort();
learners.dedup();
learners.retain(|l| !voters.contains(l));
let joint = Membership {
voters,
voters_outgoing: current.voters(),
learners,
};
let idx = self.log_append(self.current_term, EntryPayload::Membership(joint));
self.broadcast_append();
self.maybe_advance_commit();
Ok(idx)
}
pub fn propose_catalog(
&mut self,
command: CatalogCommand,
) -> Result<LogIndex, CatalogProposeError> {
if self.role != Role::Leader {
return Err(CatalogProposeError::NotLeader {
leader: self.leader_id,
});
}
let idx = self.log_append(self.current_term, EntryPayload::Catalog(command));
self.broadcast_append();
self.maybe_advance_commit();
Ok(idx)
}
pub fn propose_saga_journal(
&mut self,
command: SagaJournalCommand,
) -> Result<LogIndex, CatalogProposeError> {
if self.role != Role::Leader {
return Err(CatalogProposeError::NotLeader {
leader: self.leader_id,
});
}
let idx = self.log_append(self.current_term, EntryPayload::SagaJournal(command));
self.broadcast_append();
self.maybe_advance_commit();
Ok(idx)
}
pub fn propose_two_phase_prepare(
&mut self,
command: TwoPhasePrepareCommand,
) -> Result<LogIndex, CatalogProposeError> {
if self.role != Role::Leader {
return Err(CatalogProposeError::NotLeader {
leader: self.leader_id,
});
}
let idx = self.log_append(self.current_term, EntryPayload::TwoPhasePrepare(command));
self.broadcast_append();
self.maybe_advance_commit();
Ok(idx)
}
pub fn propose_two_phase_abort(
&mut self,
command: TwoPhaseAbortCommand,
) -> Result<LogIndex, CatalogProposeError> {
if self.role != Role::Leader {
return Err(CatalogProposeError::NotLeader {
leader: self.leader_id,
});
}
let idx = self.log_append(self.current_term, EntryPayload::TwoPhaseAbort(command));
self.broadcast_append();
self.maybe_advance_commit();
Ok(idx)
}
pub fn propose_two_phase_journal(
&mut self,
command: TwoPhaseJournalCommand,
) -> Result<LogIndex, CatalogProposeError> {
if self.role != Role::Leader {
return Err(CatalogProposeError::NotLeader {
leader: self.leader_id,
});
}
let idx = self.log_append(self.current_term, EntryPayload::TwoPhaseJournal(command));
self.broadcast_append();
self.maybe_advance_commit();
Ok(idx)
}
pub fn propose_queue_autoscale_policy(
&mut self,
command: QueueAutoscalePolicyCommand,
) -> Result<LogIndex, CatalogProposeError> {
if self.role != Role::Leader {
return Err(CatalogProposeError::NotLeader {
leader: self.leader_id,
});
}
let idx = self.log_append(
self.current_term,
EntryPayload::QueueAutoscalePolicy(command),
);
self.broadcast_append();
self.maybe_advance_commit();
Ok(idx)
}
fn set_role(&mut self, role: Role) {
if self.role != role {
tracing::debug!(
target: "crafty::raft",
node = self.id.0,
term = self.current_term.0,
?role,
"raft role changed"
);
self.role = role;
self.outbox.push(Output::RoleChanged(role));
}
}
fn become_follower(&mut self, term: Term) {
if term > self.current_term {
self.current_term = term;
self.voted_for = None;
}
self.votes.clear();
self.fail_pending_reads();
self.lease_expiry = 0;
self.lease_acks.clear();
self.set_role(Role::Follower);
}
fn start_pre_election(&mut self) {
if !self.is_voter(self.id) {
self.reset_election_timer();
return;
}
self.set_role(Role::PreCandidate);
self.votes.clear();
self.votes.insert(self.id);
self.reset_election_timer();
if self.quorum_of_votes() {
self.start_real_election();
return;
}
let rv = RequestVote {
term: self.current_term.next(),
candidate_id: self.id,
last_log: self.log.last_id(),
pre_vote: true,
};
self.send_vote_requests(&rv);
}
fn start_real_election(&mut self) {
if !self.is_voter(self.id) {
self.reset_election_timer();
return;
}
self.current_term = self.current_term.next();
self.set_role(Role::Candidate);
self.voted_for = Some(self.id);
self.votes.clear();
self.votes.insert(self.id);
self.leader_id = None;
self.reset_election_timer();
if self.quorum_of_votes() {
self.become_leader();
return;
}
let rv = RequestVote {
term: self.current_term,
candidate_id: self.id,
last_log: self.log.last_id(),
pre_vote: false,
};
self.send_vote_requests(&rv);
}
fn send_vote_requests(&mut self, rv: &RequestVote) {
for p in self.configuration().voter_peers(self.id) {
self.outbox
.push(Output::Send(p, RaftRpc::RequestVote(rv.clone())));
}
}
fn become_leader(&mut self) {
self.set_role(Role::Leader);
self.leader_id = Some(self.id);
self.lease_expiry = 0;
let next = self.log.last_index().next();
self.next_index.clear();
self.match_index.clear();
self.sent_upper.clear();
self.last_ack_clock.clear();
for p in self.peers() {
self.next_index.insert(p, next);
self.match_index.insert(p, LogIndex::ZERO);
}
self.log_append(self.current_term, EntryPayload::Noop);
self.heartbeat_elapsed = 0;
self.broadcast_append();
self.maybe_advance_commit();
}
fn reset_election_timer(&mut self) {
self.elapsed = 0;
self.election_timeout = self.rng.range(
self.config.election_timeout_min,
self.config.election_timeout_max,
);
}
fn handle_request_vote(&mut self, from: NodeId, rv: &RequestVote) {
if !self.is_voter(self.id) {
self.reply_vote(from, false, rv.pre_vote);
return;
}
let up_to_date = rv.last_log >= self.log.last_id();
if rv.pre_vote {
let leader_recent =
self.leader_id.is_some() && self.elapsed < self.config.election_timeout_min;
let granted = rv.term >= self.current_term && up_to_date && !leader_recent;
self.reply_vote(from, granted, true);
return;
}
if rv.term > self.current_term {
self.become_follower(rv.term);
}
let mut granted = false;
if rv.term >= self.current_term {
let can_vote = self.voted_for.is_none() || self.voted_for == Some(rv.candidate_id);
if can_vote && up_to_date {
granted = true;
self.voted_for = Some(rv.candidate_id);
self.reset_election_timer();
}
}
self.reply_vote(from, granted, false);
}
fn reply_vote(&mut self, to: NodeId, vote_granted: bool, pre_vote: bool) {
let reply = RequestVoteReply {
term: self.current_term,
vote_granted,
pre_vote,
};
self.outbox
.push(Output::Reply(to, RaftRpcReply::RequestVote(reply)));
}
fn handle_vote_reply(&mut self, from: NodeId, reply: &RequestVoteReply) {
if reply.pre_vote {
if self.role == Role::PreCandidate && reply.vote_granted {
self.votes.insert(from);
if self.quorum_of_votes() {
self.start_real_election();
}
}
return;
}
if self.role != Role::Candidate || reply.term != self.current_term {
return;
}
if reply.vote_granted {
self.votes.insert(from);
if self.quorum_of_votes() {
self.become_leader();
}
}
}
fn handle_append_entries(&mut self, from: NodeId, ae: &AppendEntries) {
if ae.term < self.current_term {
self.reply_append(from, false, None, None, ae.round);
return;
}
if ae.term > self.current_term {
self.become_follower(ae.term);
} else if self.role != Role::Follower {
self.set_role(Role::Follower);
}
self.leader_id = Some(ae.leader_id);
self.reset_election_timer();
if ae.prev_log.index.0 > 0 {
match self.log.term_at(ae.prev_log.index) {
None => {
let hint = self.log.last_index().next();
self.reply_append(from, false, Some(hint), None, ae.round);
return;
}
Some(t) if t != ae.prev_log.term => {
let first = self.log.first_index_of_term(t).unwrap_or(ae.prev_log.index);
self.reply_append(from, false, Some(first), Some(t), ae.round);
return;
}
_ => {}
}
}
let mut idx = ae.prev_log.index;
for entry in &ae.entries {
idx = idx.next();
match self.log.term_at(idx) {
Some(t) if t == entry.term => {}
Some(_) => {
self.log_truncate_from(idx);
self.log_push(LogEntry {
term: entry.term,
index: idx,
payload: entry.payload.clone(),
});
}
None => {
self.log_push(LogEntry {
term: entry.term,
index: idx,
payload: entry.payload.clone(),
});
}
}
}
if ae.leader_commit > self.commit_index {
self.commit_index = ae.leader_commit.min(idx);
self.apply_committed();
}
self.reply_append(from, true, None, None, ae.round);
}
fn reply_append(
&mut self,
to: NodeId,
success: bool,
conflict_index: Option<LogIndex>,
conflict_term: Option<Term>,
round: Round,
) {
let reply = AppendEntriesReply {
term: self.current_term,
success,
conflict_index,
conflict_term,
round,
};
self.outbox
.push(Output::Reply(to, RaftRpcReply::AppendEntries(reply)));
}
fn handle_append_reply(&mut self, from: NodeId, reply: &AppendEntriesReply) {
if self.role != Role::Leader || reply.term != self.current_term {
return;
}
if reply.success {
let upper = self
.sent_upper
.get(&from)
.copied()
.unwrap_or(LogIndex::ZERO);
let current = self
.match_index
.get(&from)
.copied()
.unwrap_or(LogIndex::ZERO);
if upper > current {
self.match_index.insert(from, upper);
}
self.next_index.insert(from, upper.next());
self.last_ack_clock.insert(from, self.logical_clock);
if self.config.reachability.detector == FailureDetectorKind::PhiAccrual {
self.phi_liveness.record_heartbeat(from, self.logical_clock);
}
self.confirm_reads(from, reply.round);
if reply.round >= self.lease_round {
self.lease_acks.insert(from);
self.maybe_extend_lease();
}
self.maybe_advance_commit();
self.try_complete_reads();
} else {
let ni = if let Some(ci) = reply.conflict_index {
LogIndex(ci.0.max(1))
} else {
let cur = self.next_index.get(&from).copied().unwrap_or(LogIndex(1)).0;
LogIndex(cur.saturating_sub(1).max(1))
};
self.next_index.insert(from, ni);
self.send_append(from);
}
}
fn broadcast_append(&mut self) {
self.heartbeat_round = self.heartbeat_round.next();
self.lease_round = self.heartbeat_round;
self.lease_round_clock = self.logical_clock;
self.lease_acks.clear();
self.lease_acks.insert(self.id);
self.maybe_extend_lease();
for p in self.peers() {
self.send_append(p);
}
}
fn send_append(&mut self, peer: NodeId) {
let ni = self
.next_index
.get(&peer)
.copied()
.unwrap_or_else(|| self.log.last_index().next());
if ni.0 <= self.log.snapshot_index().0 && self.snapshot.is_some() {
self.send_snapshot(peer);
return;
}
let prev_index = LogIndex(ni.0.saturating_sub(1));
let prev_term = self.log.term_at(prev_index).unwrap_or(Term::ZERO);
let entries = self.log.entries_from(ni).to_vec();
let upper = LogIndex(prev_index.0 + entries.len() as u64);
self.sent_upper.insert(peer, upper);
let ae = AppendEntries {
term: self.current_term,
leader_id: self.id,
prev_log: LogId::new(prev_term, prev_index),
entries,
leader_commit: self.commit_index,
round: self.heartbeat_round,
};
self.outbox
.push(Output::Send(peer, RaftRpc::AppendEntries(ae)));
}
fn send_snapshot(&mut self, peer: NodeId) {
let Some(snap) = self.snapshot.as_ref() else {
return;
};
let is = InstallSnapshot {
term: self.current_term,
leader_id: self.id,
last_included: LogId::new(snap.last_term, snap.last_index),
last_config: snap.membership.clone(),
offset: 0,
data: snap.data.clone(),
done: true,
};
self.sent_upper.insert(peer, snap.last_index);
self.outbox
.push(Output::Send(peer, RaftRpc::InstallSnapshot(is)));
}
fn maybe_advance_commit(&mut self) {
if self.role != Role::Leader {
return;
}
let last = self.log.last_index().0;
let mut new_commit = self.commit_index;
for n in (self.commit_index.0 + 1)..=last {
let idx = LogIndex(n);
if self.log.term_at(idx) != Some(self.current_term) {
continue;
}
let mut acked: BTreeSet<NodeId> = BTreeSet::new();
acked.insert(self.id);
for (peer, m) in &self.match_index {
if m.0 >= n {
acked.insert(*peer);
}
}
if self.quorum_ok(&acked) {
new_commit = idx;
}
}
if new_commit > self.commit_index {
self.commit_index = new_commit;
self.apply_committed();
self.maybe_finalize_membership();
self.maybe_step_down_if_removed();
self.try_complete_reads();
}
}
fn apply_committed(&mut self) {
while self.last_applied < self.commit_index {
let next = self.last_applied.next();
match self.log.get(next).map(|e| &e.payload) {
Some(EntryPayload::Command(c)) => {
self.outbox.push(Output::Apply(Committed {
index: next,
command: c.clone(),
}));
}
Some(EntryPayload::Catalog(command)) => {
self.outbox.push(Output::CatalogApplied {
index: next,
command: command.clone(),
});
}
Some(EntryPayload::SagaJournal(command)) => {
self.outbox.push(Output::SagaJournalApplied {
index: next,
command: command.clone(),
});
}
Some(EntryPayload::TwoPhasePrepare(command)) => {
self.outbox.push(Output::TwoPhasePrepareApplied {
index: next,
command: command.clone(),
});
}
Some(EntryPayload::TwoPhaseAbort(command)) => {
self.outbox.push(Output::TwoPhaseAbortApplied {
index: next,
command: command.clone(),
});
}
Some(EntryPayload::TwoPhaseJournal(command)) => {
self.outbox.push(Output::TwoPhaseJournalApplied {
index: next,
command: command.clone(),
});
}
Some(EntryPayload::QueueAutoscalePolicy(command)) => {
self.outbox.push(Output::QueueAutoscalePolicyApplied {
index: next,
command: command.clone(),
});
}
_ => {}
}
self.last_applied = next;
}
}
fn confirm_reads(&mut self, from: NodeId, round: Round) {
for r in &mut self.pending_reads {
if round >= r.round {
r.acks.insert(from);
}
}
}
fn try_complete_reads(&mut self) {
if self.role != Role::Leader || self.pending_reads.is_empty() {
return;
}
if self.log.term_at(self.commit_index) != Some(self.current_term) {
return;
}
let conf = self.configuration();
let applied = self.last_applied;
let mut ready = Vec::new();
self.pending_reads.retain(|r| {
if conf.has_quorum(&r.acks) && applied >= r.index {
ready.push((r.id, r.index));
false
} else {
true
}
});
for (id, index) in ready {
self.outbox.push(Output::ReadReady { id, index });
}
}
fn fail_pending_reads(&mut self) {
for r in std::mem::take(&mut self.pending_reads) {
self.outbox.push(Output::ReadFailed { id: r.id });
}
}
fn lease_ticks(&self) -> u64 {
self.config.election_timeout_min / 2
}
fn maybe_extend_lease(&mut self) {
if self.role != Role::Leader {
return;
}
if self.configuration().has_quorum(&self.lease_acks) {
let candidate = self.lease_round_clock.saturating_add(self.lease_ticks());
if candidate > self.lease_expiry {
self.lease_expiry = candidate;
}
}
}
#[must_use]
pub fn lease_valid(&self) -> bool {
self.role == Role::Leader && self.logical_clock < self.lease_expiry
}
pub fn lease_read(&self) -> Result<Option<LogIndex>, NotLeader> {
if self.role != Role::Leader {
return Err(NotLeader {
leader: self.leader_id,
});
}
let authoritative = self.log.term_at(self.commit_index) == Some(self.current_term);
if self.logical_clock < self.lease_expiry && authoritative {
Ok(Some(self.commit_index))
} else {
Ok(None)
}
}
#[must_use]
pub fn reachable(&self, window: u64) -> Vec<NodeId> {
let voters = self.configuration().voters();
if self.role != Role::Leader {
return voters;
}
let now = self.logical_clock;
voters
.into_iter()
.filter(|&v| {
v == self.id
|| self
.last_ack_clock
.get(&v)
.is_some_and(|&acked| now.saturating_sub(acked) <= window)
})
.collect()
}
#[must_use]
pub fn reachable_now(&self) -> Vec<NodeId> {
let voters = self.configuration().voters();
if self.role != Role::Leader {
return voters;
}
let now = self.logical_clock;
voters
.into_iter()
.filter(|&v| match self.config.reachability.detector {
FailureDetectorKind::AckWindow => v == self.id || self.ack_liveness.is_reachable(v),
FailureDetectorKind::PhiAccrual => {
v == self.id || self.phi_liveness.is_reachable(v, now)
}
})
.collect()
}
fn update_liveness(&mut self) {
if self.role != Role::Leader {
return;
}
let voters = self.configuration().voters();
let now = self.logical_clock;
match self.config.reachability.detector {
FailureDetectorKind::AckWindow => {
let window = self
.config
.reachability
.window(self.config.election_timeout_max);
let hysteresis = self
.config
.reachability
.hysteresis(self.config.election_timeout_min);
self.ack_liveness.update(
now,
self.id,
&voters,
&self.last_ack_clock,
window,
hysteresis,
);
}
FailureDetectorKind::PhiAccrual => {}
}
}
fn maybe_finalize_membership(&mut self) {
if self.role != Role::Leader {
return;
}
let conf = self.configuration();
let cfg_idx = self.config_index();
if conf.is_joint() && cfg_idx.0 != 0 && cfg_idx <= self.commit_index {
let final_config = Membership {
voters: conf.voters(),
voters_outgoing: Vec::new(),
learners: conf.to_membership().learners,
};
self.log_append(self.current_term, EntryPayload::Membership(final_config));
self.broadcast_append();
}
}
fn maybe_step_down_if_removed(&mut self) {
if self.role != Role::Leader {
return;
}
let conf = self.configuration();
if !conf.is_joint() && self.config_index() <= self.commit_index && !conf.is_voter(self.id) {
self.become_follower(self.current_term);
self.leader_id = None;
}
}
fn handle_install_snapshot(&mut self, from: NodeId, is: InstallSnapshot) {
if is.term < self.current_term {
self.reply_snapshot(from);
return;
}
if is.term > self.current_term {
self.become_follower(is.term);
} else if self.role != Role::Follower {
self.set_role(Role::Follower);
}
self.leader_id = Some(is.leader_id);
self.reset_election_timer();
let last = is.last_included;
if last.index.0 <= self.log.snapshot_index().0 || last.index <= self.last_applied {
self.reply_snapshot(from);
return;
}
self.log.install_snapshot(last.index, last.term);
self.mark_log_dirty(LogIndex(last.index.0 + 1));
self.snapshot = Some(StoredSnapshot {
last_index: last.index,
last_term: last.term,
membership: is.last_config.clone(),
data: is.data.clone(),
});
if self.commit_index < last.index {
self.commit_index = last.index;
}
self.last_applied = last.index;
self.outbox.push(Output::LoadSnapshot {
index: last.index,
data: is.data,
});
self.reply_snapshot(from);
}
fn reply_snapshot(&mut self, to: NodeId) {
let reply = InstallSnapshotReply {
term: self.current_term,
};
self.outbox
.push(Output::Reply(to, RaftRpcReply::InstallSnapshot(reply)));
}
fn handle_snapshot_reply(&mut self, from: NodeId, reply: &InstallSnapshotReply) {
if self.role != Role::Leader || reply.term != self.current_term {
return;
}
let upper = self
.sent_upper
.get(&from)
.copied()
.unwrap_or(LogIndex::ZERO);
let current = self
.match_index
.get(&from)
.copied()
.unwrap_or(LogIndex::ZERO);
if upper > current {
self.match_index.insert(from, upper);
}
self.next_index.insert(from, upper.next());
self.maybe_advance_commit();
}
}