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

Reply via email to