This is an automated email from the ASF dual-hosted git repository.
krishvishal pushed a commit to branch consensus-prefix-contiguity
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/consensus-prefix-contiguity by
this push:
new c4ce76b54 fix: address review comments round 3
c4ce76b54 is described below
commit c4ce76b54081f264db693c1c8cea4662a82f36a4
Author: Krishna Vishal <[email protected]>
AuthorDate: Tue Sep 8 18:00:25 2026 +0530
fix: address review comments round 3
---
core/consensus/src/impls.rs | 103 +++++++++++++++-----------------
core/shard/src/lib.rs | 78 +++++++++++++++++++-----
core/simulator/src/bin/workload-fuzz.rs | 7 ++-
3 files changed, 115 insertions(+), 73 deletions(-)
diff --git a/core/consensus/src/impls.rs b/core/consensus/src/impls.rs
index 789329ab7..0d1ebce94 100644
--- a/core/consensus/src/impls.rs
+++ b/core/consensus/src/impls.rs
@@ -4484,14 +4484,47 @@ mod pipeline_entry_tests {
}
}
-/// A [`MessageBus`] that accepts everything and remembers nothing. Shared by
the
-/// test modules, none of which care what happens to a frame.
+/// Fixtures every consensus test module needs.
#[cfg(test)]
pub mod test_bus {
+ use super::{Command, METADATA_GROUP, Message, StartViewHeader};
use message_bus::{BusMessage, MessageBus};
use server_common::MESSAGE_ALIGN;
use server_common::iobuf::Frozen;
+ /// A `StartView` at `view` announcing head `op`. `commit == op` is the
steady
+ /// case (no suffix); a lower `commit` keeps an uncommitted suffix.
+ ///
+ /// # Panics
+ /// Never: a zeroed buffer of the right size is a valid `StartViewHeader`.
+ #[must_use]
+ #[allow(clippy::cast_possible_truncation)]
+ pub fn make_start_view(
+ view: u32,
+ op: u64,
+ commit: u64,
+ replica: u8,
+ incarnation: u128,
+ ) -> Message<StartViewHeader> {
+ let size = std::mem::size_of::<StartViewHeader>();
+ let mut msg = Message::<StartViewHeader>::new(size);
+ let header = bytemuck::checked::try_from_bytes_mut::<StartViewHeader>(
+ &mut msg.as_mut_slice()[..size],
+ )
+ .expect("zeroed bytes are a valid StartViewHeader");
+ header.command = Command::StartView;
+ header.cluster = 1;
+ header.view = view;
+ header.op = op;
+ header.commit = commit;
+ header.replica = replica;
+ header.incarnation = incarnation;
+ header.group = METADATA_GROUP;
+ header.size = size as u32;
+ msg
+ }
+
+ /// A [`MessageBus`] that accepts everything and remembers nothing.
pub struct NoopBus;
impl MessageBus for NoopBus {
@@ -4539,7 +4572,7 @@ mod timestamp_clamp_tests {
}
}
- use crate::test_bus::NoopBus;
+ use crate::test_bus::{NoopBus, make_start_view};
#[test]
fn observed_log_timestamp_floors_new_primary_stamps() {
@@ -4593,34 +4626,6 @@ mod timestamp_clamp_tests {
);
}
- #[allow(clippy::cast_possible_truncation)]
- /// `commit == op` is the steady case (no suffix); a lower `commit` keeps
an
- /// uncommitted suffix.
- fn make_start_view(
- view: u32,
- op: u64,
- commit: u64,
- replica: u8,
- incarnation: u128,
- ) -> Message<StartViewHeader> {
- let size = std::mem::size_of::<StartViewHeader>();
- let mut msg = Message::<StartViewHeader>::new(size);
- let header = bytemuck::checked::try_from_bytes_mut::<StartViewHeader>(
- &mut msg.as_mut_slice()[..size],
- )
- .expect("zeroed bytes are a valid StartViewHeader");
- header.command = Command::StartView;
- header.cluster = 1;
- header.view = view;
- header.op = op;
- header.commit = commit;
- header.replica = replica;
- header.incarnation = incarnation;
- header.group = METADATA_GROUP;
- header.size = size as u32;
- msg
- }
-
#[test]
fn
given_recovering_replica_when_start_view_incarnation_foreign_should_ignore() {
// A StartView addressed to a PREVIOUS incarnation, still in flight
when the
@@ -5568,29 +5573,7 @@ mod recovery_barrier_tests {
use super::*;
use crate::LocalPipeline;
- use crate::test_bus::NoopBus;
-
- /// A `StartView` at `view` announcing head `op` with commit point
`commit`.
- fn start_view(view: u32, op: u64, commit: u64) -> Message<StartViewHeader>
{
- let size = std::mem::size_of::<StartViewHeader>();
- let mut msg = Message::<StartViewHeader>::new(size);
- let header = bytemuck::checked::try_from_bytes_mut::<StartViewHeader>(
- &mut msg.as_mut_slice()[..size],
- )
- .expect("zeroed bytes are a valid StartViewHeader");
- header.command = Command::StartView;
- header.cluster = 1;
- header.view = view;
- header.op = op;
- header.commit = commit;
- header.replica = 1;
- header.group = METADATA_GROUP;
- #[allow(clippy::cast_possible_truncation)]
- {
- header.size = size as u32;
- }
- msg
- }
+ use crate::test_bus::{NoopBus, make_start_view};
/// Recovered at head 120, proven committed only through 100, in view 7.
fn recovered_with_gated_suffix() -> VsrConsensus<NoopBus, LocalPipeline> {
@@ -5624,7 +5607,11 @@ mod recovery_barrier_tests {
// them and nothing re-prepares them.
assert!(
!consensus
- .handle_start_view(PlaneKind::Metadata, start_view(7, 105,
105).header(), &[])
+ .handle_start_view(
+ PlaneKind::Metadata,
+ make_start_view(7, 105, 105, 1, 0).header(),
+ &[]
+ )
.is_empty(),
"the StartView at the commit floor must be adopted"
);
@@ -5654,7 +5641,11 @@ mod recovery_barrier_tests {
// Head 120, commit still 100: 101..=120 re-replicate under the new
view.
assert!(
!consensus
- .handle_start_view(PlaneKind::Metadata, start_view(7, 120,
100).header(), &[])
+ .handle_start_view(
+ PlaneKind::Metadata,
+ make_start_view(7, 120, 100, 1, 0).header(),
+ &[]
+ )
.is_empty(),
"the StartView carrying the surviving suffix must be adopted"
);
diff --git a/core/shard/src/lib.rs b/core/shard/src/lib.rs
index 69cf2db25..e408541c8 100644
--- a/core/shard/src/lib.rs
+++ b/core/shard/src/lib.rs
@@ -5331,13 +5331,23 @@ where
"metadata repair peer retained nothing in the
requested range; \
re-requesting rather than converting to state
transfer"
);
+ // Charge the round, and ACT on exhaustion.
`gap_repair_peer`
+ // re-picks the primary deterministically, so dropping
the
+ // session on its own re-arms the same peer at the
debounce
+ // interval forever: the stall path never runs, so
rotation
+ // is never reached and the state-transfer escalation
this
+ // guard replaced stays out of reach.
+ if self.burn_metadata_repair_attempt() {
+ self.rotate_stalled_metadata_repair(
+ consensus,
+ header.replica,
+ session.from_op,
+ session.to_op,
+ )
+ .await;
+ return;
+ }
*self.metadata_repair.borrow_mut() = None;
- // Charged, not cleared. `gap_repair_peer` re-picks the
- // primary deterministically, so clearing makes this
- // request / `RangeEvicted` / re-arm cycle unbounded
at the
- // debounce interval and puts the state-transfer
escalation
- // this guard replaced out of reach.
- self.burn_metadata_repair_attempt();
return;
}
if consensus.state_transfer_stage() ==
consensus::StateTransferStage::Idle {
@@ -5796,13 +5806,8 @@ where
})
};
if let Some((peer, nonce, session_from_op, to_op)) = stalled {
- // What the walk has not reached, floored at the op this session
was
- // armed for. Recomputing would ask for a different window than
the one
- // reported missing: the merged-log scan opens above the snapshot
floor
- // and its `committed_elsewhere` fallback reports ops below even
that,
- // both already folded into `session.from_op`. It also carries the
- // initial arm's snapshot clamp, so a retry cannot ask for
compacted ops.
- let from_op = session_from_op.max(consensus.commit_min() + 1);
+ let from_op =
+ stalled_repair_from_op(session_from_op,
consensus.commit_min(), repairing_view);
if from_op > to_op {
// `from_op` past `to_op` without `commit_min` reaching it: the
// primary-elect window above starts at the merged log's commit
@@ -10403,6 +10408,26 @@ where
}
}
+/// Where a stalled repair session reopens its window.
+///
+/// A merged-log session reopens exactly where it was armed. Its floor is
COVERAGE,
+/// not the walk: `first_op_not_covered` reports ops whose journal entry is
absent
+/// or diverging, and an op can be applied (`commit_min` past it) while its
entry is
+/// gone. Raising the floor to `commit_min + 1` there skips the very op the
scan
+/// reported, and the view change parks on it forever.
+///
+/// A tail-repair session is the other way round. Its window IS the commit
gap, so
+/// ops the walk has since consumed must not be asked for again.
`session.from_op`
+/// still floors it, carrying the initial arm's snapshot clamp so no retry
asks for
+/// compacted ops.
+fn stalled_repair_from_op(session_from_op: u64, commit_min: u64,
repairing_view: bool) -> u64 {
+ if repairing_view {
+ session_from_op
+ } else {
+ session_from_op.max(commit_min + 1)
+ }
+}
+
/// Walk a merged-log source list one step past `avoid`, wrapping.
///
/// A ring, not a filter. The list is `log_view`-ordered and identical on every
@@ -12917,7 +12942,7 @@ mod metadata_repair_session_tests {
use super::{
MetadataRepairSession, gap_repair_peer, metadata_repair_superseded,
next_transfer_peer,
- repair_chunk_walked,
+ repair_chunk_walked, stalled_repair_from_op,
};
/// Armed at view 3, against the window `11..=20`.
@@ -12932,6 +12957,31 @@ mod metadata_repair_session_tests {
}
}
+ /// A merged-log session must re-ask for the op the coverage scan
reported, even
+ /// once the walk has passed it. Coverage is about the journal ENTRY; an
op can
+ /// be applied and still have no entry to serve, which is exactly what
+ /// `committed_elsewhere` reports.
+ #[test]
+ fn
given_a_dropped_response_below_commit_min_when_retrying_should_still_ask_for_it()
{
+ assert_eq!(
+ stalled_repair_from_op(5, 6, true),
+ 5,
+ "clamping to commit_min + 1 would retry from 7 and skip the
reported hole"
+ );
+ }
+
+ /// The tail-repair session is the other way round: its window is the
commit gap,
+ /// so ops the walk consumed must not be re-requested.
+ #[test]
+ fn
given_a_walked_window_when_retrying_a_tail_session_should_open_above_it() {
+ assert_eq!(stalled_repair_from_op(5, 6, false), 7);
+ assert_eq!(
+ stalled_repair_from_op(11, 3, false),
+ 11,
+ "the arm floor still holds, so no retry asks for compacted ops"
+ );
+ }
+
#[test]
fn given_a_gap_stopped_backup_when_picking_a_peer_should_ask_the_primary()
{
assert_eq!(gap_repair_peer(2, 3, 0), Some(0));
diff --git a/core/simulator/src/bin/workload-fuzz.rs
b/core/simulator/src/bin/workload-fuzz.rs
index e5379577e..fb6ab966a 100644
--- a/core/simulator/src/bin/workload-fuzz.rs
+++ b/core/simulator/src/bin/workload-fuzz.rs
@@ -142,9 +142,10 @@ struct Args {
///
/// Separate because a partition-only run satisfies the combined floor
while the
/// metadata oracle compares an empty chain against an empty chain and
agrees.
- /// Default `0`: a partition-focused run legitimately commits no metadata,
so a
- /// campaign that wants the metadata property tested asks for it.
- #[arg(long, default_value_t = 0)]
+ /// Defaults to `1`, which is the floor `--min-ops-compared` carried
before it
+ /// counted both planes; a campaign that wants no metadata coverage opts
out
+ /// with `0`.
+ #[arg(long, default_value_t = 1)]
min_metadata_ops_compared: usize,
/// Fail the run if crash or restart injection was requested but never
happened.
/// Off by default, since a short run at low probability may legitimately
draw