Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
55 changes: 55 additions & 0 deletions nodedb-cluster/src/raft_loop/tick/apply_committed.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,10 +12,31 @@ use crate::forward::PlanExecutor;

use super::super::loop_core::{CommitApplier, RaftLoop};

/// The first committed index that is not greater than its predecessor, when
/// the batch carries one.
///
/// Committed entries are contiguous and strictly ascending by construction, so
/// a non-increasing pair is a producer regression. The applier's delivery guard
/// stays the boundary that absorbs a repeat; this makes the regression
/// observable instead of silent.
fn first_non_increasing_committed_index(entries: &[nodedb_raft::LogEntry]) -> Option<u64> {
entries
.windows(2)
.find(|pair| pair[1].index <= pair[0].index)
.map(|pair| pair[1].index)
}

impl<A: CommitApplier, P: PlanExecutor> RaftLoop<A, P> {
/// Apply one group's committed entries from this tick's `Ready` output.
/// Called only when `!group_ready.committed_entries.is_empty()`.
pub(super) fn apply_group_commits(&self, group_id: u64, group_ready: &nodedb_raft::Ready) {
if let Some(index) = first_non_increasing_committed_index(&group_ready.committed_entries) {
warn!(
group_id,
index,
"committed batch is not strictly increasing; the applier guard absorbs a repeat"
);
}
for entry in &group_ready.committed_entries {
if let Some(cc) = ConfChange::from_entry_data(&entry.data) {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
Expand Down Expand Up @@ -111,3 +132,37 @@ impl<A: CommitApplier, P: PlanExecutor> RaftLoop<A, P> {
}
}
}

#[cfg(test)]
mod tests {
use super::*;
use nodedb_raft::LogEntry;

fn entry(index: u64) -> LogEntry {
LogEntry {
term: 1,
index,
data: Vec::new(),
}
}

#[test]
fn detects_a_non_increasing_committed_index() {
assert_eq!(first_non_increasing_committed_index(&[]), None);
assert_eq!(first_non_increasing_committed_index(&[entry(1)]), None);
assert_eq!(
first_non_increasing_committed_index(&[entry(1), entry(2)]),
None
);
// A repeat inside the batch: the first index not greater than its
// predecessor is reported.
assert_eq!(
first_non_increasing_committed_index(&[entry(2), entry(2), entry(3)]),
Some(2)
);
assert_eq!(
first_non_increasing_committed_index(&[entry(1), entry(2), entry(1)]),
Some(1)
);
}
}
50 changes: 49 additions & 1 deletion nodedb-raft/src/node/internal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -412,7 +412,19 @@ impl<S: LogStorage> RaftNode<S> {
}

pub(super) fn collect_committed_entries(&mut self) {
let from = self.volatile.last_applied + 1;
// Resume from the furthest index already queued into `Ready`, not only
// from `last_applied`. This runs on every commit-index advance — a
// follower's AppendEntries, a single voter's propose — while the loop
// drains `Ready` and advances `last_applied` later. Two advances inside
// one window would queue the same committed range twice, and the
// applier then receives the same committed index twice in one batch.
let queued_through = self
.ready
.committed_entries
.last()
.map(|entry| entry.index)
.unwrap_or(0);
let from = self.volatile.last_applied.max(queued_through) + 1;
let to = self.volatile.commit_index;
if from > to {
return;
Expand All @@ -434,3 +446,39 @@ impl<S: LogStorage> RaftNode<S> {
self.election_deadline = Instant::now() + timeout;
}
}

#[cfg(test)]
mod tests {
use crate::node::core::RaftNode;
use crate::storage::MemStorage;
use crate::test_support::{force_election, test_config};

/// Two commit advances inside one `Ready` window queue each index once:
/// the second collect resumes past what the first already queued.
#[test]
fn repeated_collection_queues_each_index_once() {
let mut node = RaftNode::new(test_config(1, vec![]), MemStorage::new());
force_election(&mut node);
let election = node.take_ready();
if let Some(last) = election.committed_entries.last() {
node.advance_applied(last.index);
}

node.propose(b"one".to_vec())
.expect("single voter commits immediately");
node.propose(b"two".to_vec())
.expect("single voter commits immediately");

let ready = node.take_ready();
let indices: Vec<u64> = ready
.committed_entries
.iter()
.map(|entry| entry.index)
.collect();
assert_eq!(indices.len(), 2, "one entry per propose: {indices:?}");
assert!(
indices.windows(2).all(|pair| pair[0] < pair[1]),
"queued indices must be strictly increasing: {indices:?}"
);
}
}
Loading