diff --git a/nodedb-cluster/src/raft_loop/tick/apply_committed.rs b/nodedb-cluster/src/raft_loop/tick/apply_committed.rs index f7540364f..61028aeb9 100644 --- a/nodedb-cluster/src/raft_loop/tick/apply_committed.rs +++ b/nodedb-cluster/src/raft_loop/tick/apply_committed.rs @@ -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 { + entries + .windows(2) + .find(|pair| pair[1].index <= pair[0].index) + .map(|pair| pair[1].index) +} + impl RaftLoop { /// 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()); @@ -111,3 +132,37 @@ impl RaftLoop { } } } + +#[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) + ); + } +} diff --git a/nodedb-raft/src/node/internal.rs b/nodedb-raft/src/node/internal.rs index 16f823202..6af9a3551 100644 --- a/nodedb-raft/src/node/internal.rs +++ b/nodedb-raft/src/node/internal.rs @@ -412,7 +412,19 @@ impl RaftNode { } 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; @@ -434,3 +446,39 @@ impl RaftNode { 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 = 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:?}" + ); + } +}