This is an automated email from the ASF dual-hosted git repository. spetz pushed a commit to branch rewind_wal_parent_fix in repository https://gitbox.apache.org/repos/asf/iggy.git
commit 35f6a3c0ea70524f4f3a4de82751abd47e8e4e6d Author: spetz <[email protected]> AuthorDate: Thu Aug 27 15:17:58 2026 +0200 fix(cluster): prevent repaired prepares from rewinding WAL parent --- core/consensus/src/plane_helpers.rs | 56 +++++++++++++++++++++++++++++++++++ core/partitions/src/iggy_partition.rs | 25 +++++++--------- core/shard/src/lib.rs | 18 ++++++----- 3 files changed, 77 insertions(+), 22 deletions(-) diff --git a/core/consensus/src/plane_helpers.rs b/core/consensus/src/plane_helpers.rs index 4e4bcd7ef..9a6f1667a 100644 --- a/core/consensus/src/plane_helpers.rs +++ b/core/consensus/src/plane_helpers.rs @@ -771,6 +771,26 @@ pub fn panic_if_hash_chain_would_break_in_same_view( } } +/// Resolve a repair's new contiguous head and its hash-chain checksum. +/// +/// A DVC can put the sequencer above holes that repair later backfills. Such a +/// lower entry changes durable coverage, but not the head the next prepare must +/// parent. Updating the checksum from the incoming entry would rewind the chain +/// while leaving the sequencer ahead of it. +pub fn repaired_frontier_update( + previous_frontier: u64, + mut header_at: impl FnMut(u64) -> Option<PrepareHeader>, +) -> Option<(u64, u128)> { + let mut frontier = previous_frontier; + while header_at(frontier + 1).is_some() { + frontier += 1; + } + if frontier == previous_frontier { + return None; + } + header_at(frontier).map(|header| (frontier, header.checksum)) +} + /// Ack a prepare back to its primary once the owning plane vouches for it. /// /// `is_persisted` is the caller's journal-containment verdict for `header`: @@ -856,6 +876,7 @@ mod tests { use iggy_common::calculate_checksum; use message_bus::SendError; use server_common::{MESSAGE_ALIGN, iobuf::Frozen}; + use std::collections::BTreeMap; /// `PrepareHeader`'s alignment, which every suffix body has to satisfy. const BODY_ALIGN: usize = align_of::<PrepareHeader>(); @@ -2022,4 +2043,39 @@ mod tests { ); assert!(header.validate().is_ok()); } + + fn header_with_checksum(op: u64, checksum: u128) -> PrepareHeader { + PrepareHeader { + command: Command::Prepare, + op, + checksum, + ..Default::default() + } + } + + #[test] + fn given_lower_backfill_when_resolving_repaired_frontier_should_not_rewind() { + let headers = BTreeMap::from([ + (169, header_with_checksum(169, 1690)), + (170, header_with_checksum(170, 1700)), + (171, header_with_checksum(171, 1710)), + ]); + + let update = repaired_frontier_update(171, |op| headers.get(&op).copied()); + + assert_eq!(update, None); + } + + #[test] + fn given_gap_closure_when_resolving_repaired_frontier_should_use_new_head_checksum() { + let headers = BTreeMap::from([ + (170, header_with_checksum(170, 1700)), + (171, header_with_checksum(171, 1710)), + (172, header_with_checksum(172, 1720)), + ]); + + let update = repaired_frontier_update(169, |op| headers.get(&op).copied()); + + assert_eq!(update, Some((172, 1720))); + } } diff --git a/core/partitions/src/iggy_partition.rs b/core/partitions/src/iggy_partition.rs index 17bf8aefa..213c8d72f 100644 --- a/core/partitions/src/iggy_partition.rs +++ b/core/partitions/src/iggy_partition.rs @@ -40,7 +40,7 @@ use consensus::{ ReplicaLogContext, RequestLogEvent, Sequencer, SimEventKind, VsrConsensus, ack_preflight, ack_quorum_reached, build_deny_reply_from_request, build_reply_from_request, build_reply_message, drain_committable_prefix, emit_namespace_progress_event, - emit_partition_diag, emit_sim_event, fence_old_prepare_by_commit, + emit_partition_diag, emit_sim_event, fence_old_prepare_by_commit, repaired_frontier_update, replicate_frozen_to_next_in_chain, replicate_preflight, restamp_prepare_view, send_prepare_ok as send_prepare_ok_common, verify_prepare_integrity, }; @@ -4447,21 +4447,18 @@ where // cannot walk. A dropped frame stalls the frontier here; the stall // retry refills the hole and the next apply resumes the advance // (walking over ops that were journaled out of order meanwhile). - let mut frontier = self.consensus().sequencer().current_sequence(); - while self - .log - .journal() - .inner - .header_by_op(frontier + 1) - .is_some() - { - frontier += 1; - } - let consensus = self.consensus(); - if frontier > consensus.sequencer().current_sequence() { + // The checksum moves only with the frontier and is read from the + // journal header at the new head: a lower backfill must not rewind + // the parent the next prepare chains onto. + let previous_frontier = self.consensus().sequencer().current_sequence(); + let update = repaired_frontier_update(previous_frontier, |op| { + self.log.journal().inner.header_by_op(op) + }); + if let Some((frontier, frontier_checksum)) = update { + let consensus = self.consensus(); consensus.sequencer().set_sequence(frontier); + consensus.set_last_prepare_checksum(frontier_checksum); } - consensus.set_last_prepare_checksum(header.checksum); } /// Conclude a repair stream: settle the commit floor at the serving diff --git a/core/shard/src/lib.rs b/core/shard/src/lib.rs index 1efaa598a..a6a93f75a 100644 --- a/core/shard/src/lib.rs +++ b/core/shard/src/lib.rs @@ -32,7 +32,8 @@ use consensus::{ DvcSuffix, FatalReason, MergedLog, MetadataHandle, MuxPlane, PartitionsHandle, Pipeline, Plane, PlaneKind, STATE_TRANSFER_MAX_DECODE_RETRIES, STATE_TRANSFER_MAX_STALL_RETRIES, Sequencer, Status, VsrAction, VsrConsensus, build_deny_reply_from_request_header, dvc_blank, - dvc_header_kind, encode_prepare_headers, fatal, restamp_prepare_view, verify_prepare_integrity, + dvc_header_kind, encode_prepare_headers, fatal, repaired_frontier_update, restamp_prepare_view, + verify_prepare_integrity, }; #[cfg(any(test, feature = "simulator"))] use crossfire::AsyncRxTrait; @@ -4473,15 +4474,15 @@ where // `apply_repaired_prepare`: DVC advertises the sequencer, so a // hole below a repaired op must stall the advance rather than // mint an election candidate with an unwalkable log. - let mut frontier = consensus.sequencer().current_sequence(); + let previous_frontier = consensus.sequencer().current_sequence(); #[allow(clippy::cast_possible_truncation)] - while journal.header((frontier + 1) as usize).is_some() { - frontier += 1; - } - if frontier > consensus.sequencer().current_sequence() { + let update = repaired_frontier_update(previous_frontier, |op| { + journal.header(op as usize).map(|header| *header) + }); + if let Some((frontier, frontier_checksum)) = update { consensus.sequencer().set_sequence(frontier); + consensus.set_last_prepare_checksum(frontier_checksum); } - consensus.set_last_prepare_checksum(header.checksum); return; } // A metadata-plane op that did not match above (no metadata consensus on @@ -9417,9 +9418,10 @@ mod persist_gate_tests { mod repair_scope_tests { //! Who parked the log decides what it means. - use super::{MergedLog, repair_op_in_scope, repair_serve_ceiling}; use iggy_binary_protocol::{Command, PrepareHeader}; + use super::{MergedLog, repair_op_in_scope, repair_serve_ceiling}; + fn header(op: u64) -> PrepareHeader { PrepareHeader { command: Command::Prepare,
