This is an automated email from the ASF dual-hosted git repository.

krishvishal pushed a commit to branch simulator-liveness
in repository https://gitbox.apache.org/repos/asf/iggy.git

commit 48df814973ca75bd71c2131cd25d4ef9ccec8408
Author: Krishna Vishal <[email protected]>
AuthorDate: Thu Sep 10 12:16:31 2026 +0530

    fix(consensus): let a probe answer expose a divergent head
    
    A replica that adopts a view through a RequestStartView probe learns only
    the primary's numbers. It can hold a different entry at an op at or below
    the announced commit point, left from a view whose entry the change
    truncated, and with no canonical headers to compare against it adopts the
    commit point and applies its own stale entry. Two replicas then hold
    different history at the same log position.
    
    Seed 144 of the uniform swarm lane reaches this at op 7 with no crash
    injected: replica 0 prepares op 7 in view 0 and cannot get it acked,
    view 1 truncates it and prepares a different op 7, and replica 0 adopts
    view 1 by probe while keeping its own.
    
    The probe answer now carries the same suffix a DoViewChange does, which
    is what arms adopt_start_view_suffix and gives the reconcile sweep
    something to walk.
    
    Two adjacent bounds had to move with it. The announced commit is now
    clamped to the head, as DoViewChange already clamped it: commit_max
    legitimately runs ahead of the head, StartViewHeader::validate refuses
    commit > op, and the dispatcher turned that refusal into a panic. And the
    reconcile floor on both planes was announced_commit.max(commit_min),
    against the reasoning in the comment directly above it: the announced
    commit is what the cluster committed, not what this replica applied, so
    the divergent op sat at the floor and was skipped as unreconcilable.
    
    Over 7000 runs of 14 lanes at seeds 1..=500, failures fall from 62 to 55:
    the commit > op panic goes to zero and metadata divergence drops from 12
    to 8. The partition-plane floor is corrected for consistency and changes
    no outcome, so the partition divergences have another cause.
---
 core/consensus/src/impls.rs | 113 ++++++++++++++++++++++++++++++++++++++++++--
 core/shard/src/lib.rs       |   4 +-
 core/simulator/src/lib.rs   |  85 +++++++++++++++++++++++++++++++++
 3 files changed, 196 insertions(+), 6 deletions(-)

diff --git a/core/consensus/src/impls.rs b/core/consensus/src/impls.rs
index 64e28825a..d9d366779 100644
--- a/core/consensus/src/impls.rs
+++ b/core/consensus/src/impls.rs
@@ -3235,12 +3235,14 @@ impl<B: MessageBus, P: Pipeline<Entry = PipelineEntry>> 
VsrConsensus<B, P> {
         vec![VsrAction::SendStartView {
             view: self.view.get(),
             op: self.sequencer.current_sequence(),
-            commit: self.commit_max.get(),
+            commit: self.dvc_commit(),
             incarnation: header.incarnation,
             target: Some(header.replica),
-            // A probe answer reports this primary's settled frontier, not a
-            // freshly merged log, so there is no canonical suffix to publish.
-            suffix: Vec::new(),
+            // A prober can hold a different entry at an op at or below this
+            // commit point, left from a view whose entry the change truncated.
+            // Without the canonical headers it cannot know, so it adopts the
+            // commit point and applies its own stale entry.
+            suffix: self.local_dvc_suffix().headers().to_vec(),
             group: self.group,
         }]
     }
@@ -5895,3 +5897,106 @@ mod recovery_barrier_tests {
         );
     }
 }
+
+#[cfg(test)]
+mod probe_answer_tests {
+    //! What a primary publishes when it answers a `RequestStartView` probe.
+    //!
+    //! The answer is not a bare frontier report. A prober can hold a different
+    //! entry at an op at or below the announced commit point, left from a view
+    //! whose entry the change truncated, and the canonical headers are its 
only
+    //! way to find that out before it adopts the commit point.
+
+    use super::*;
+    use crate::LocalPipeline;
+    use crate::test_bus::NoopBus;
+    use crate::view_change_quorum::dvc_blank;
+
+    const PROBER: u8 = 1;
+
+    /// Replica 0 of 3, primary in view 0 with head `head` and commit point
+    /// `commit`, carrying a fresh suffix snapshot over `[commit + 1, head]`.
+    fn primary_with_suffix(head: u64, commit: u64) -> VsrConsensus<NoopBus, 
LocalPipeline> {
+        let consensus = VsrConsensus::new(1, 0, 3, METADATA_GROUP, NoopBus, 
LocalPipeline::new());
+        consensus.init();
+        consensus.sequencer().set_sequence(head);
+        consensus.advance_commit_max(commit);
+        // Tagged on `(head, commit)` as they stand now, so `local_dvc_suffix`
+        // returns it rather than falling back to empty.
+        let headers: Vec<PrepareHeader> = (commit + 
1..=head).rev().map(dvc_blank).collect();
+        let present = (1u128 << headers.len()) - 1;
+        consensus.set_local_dvc_suffix(DvcSuffix::new(headers, 0, present));
+        consensus
+    }
+
+    #[allow(clippy::cast_possible_truncation)]
+    fn probe(view: u32) -> RequestStartViewHeader {
+        RequestStartViewHeader {
+            checksum: 0,
+            checksum_body: 0,
+            cluster: 1,
+            size: size_of::<RequestStartViewHeader>() as u32,
+            view,
+            release: 0,
+            command: Command::RequestStartView,
+            replica: PROBER,
+            reserved_frame: [0; 66],
+            group: METADATA_GROUP,
+            reserved: [0; 104],
+            incarnation: 0,
+        }
+    }
+
+    /// The suffix is the whole point: without it the prober cannot tell its 
own
+    /// op from the view's op at the same number, adopts the commit point, and
+    /// commits the stale entry.
+    #[test]
+    fn 
given_a_cached_suffix_when_answering_a_probe_should_publish_its_headers() {
+        let consensus = primary_with_suffix(7, 5);
+
+        let actions = consensus.handle_request_start_view(PlaneKind::Metadata, 
&probe(0));
+
+        let [
+            VsrAction::SendStartView {
+                op,
+                commit,
+                suffix,
+                target,
+                ..
+            },
+        ] = &actions[..]
+        else {
+            panic!("a probe from a backup must be answered with one StartView: 
{actions:?}");
+        };
+        assert_eq!(*op, 7);
+        assert_eq!(*commit, 5);
+        assert_eq!(*target, Some(PROBER));
+        assert_eq!(
+            suffix.iter().map(|header| header.op).collect::<Vec<_>>(),
+            vec![7, 6],
+            "the answer must carry the view's canonical headers, high op first"
+        );
+    }
+
+    /// `commit_max` legitimately runs ahead of the head, since a replica 
learns
+    /// the commit point before it holds the prepares. 
`StartViewHeader::validate`
+    /// refuses `commit > op`, and the dispatcher turns that refusal into a 
panic,
+    /// so the announcement clamps the way `DoViewChange` already does.
+    #[test]
+    fn 
given_a_commit_point_above_the_head_when_answering_a_probe_should_clamp_it() {
+        let consensus = primary_with_suffix(7, 5);
+        consensus.advance_commit_max(9);
+        assert!(consensus.commit_max() > 
consensus.sequencer().current_sequence());
+
+        let actions = consensus.handle_request_start_view(PlaneKind::Metadata, 
&probe(0));
+
+        let [VsrAction::SendStartView { op, commit, .. }] = &actions[..] else {
+            panic!("expected one StartView: {actions:?}");
+        };
+        assert!(
+            commit <= op,
+            "announced commit {commit} exceeds head {op}, which 
StartViewHeader::validate \
+             rejects and the dispatcher panics on"
+        );
+    }
+}
diff --git a/core/shard/src/lib.rs b/core/shard/src/lib.rs
index 90121ace8..27d15c7a6 100644
--- a/core/shard/src/lib.rs
+++ b/core/shard/src/lib.rs
@@ -6162,7 +6162,7 @@ where
         // number and a backup can sit above it. Splitting on the view's number
         // would drop already-executed ops with no rollback, and silently.
         let announced_commit = pending.as_ref().map_or(0, |pending| 
pending.commit_max);
-        let applied_floor = announced_commit.max(consensus.commit_min());
+        let applied_floor = consensus.commit_min();
 
         let mut repairable_from: Option<u64> = None;
         for canonical in pending.as_ref().map_or(&[][..], |pending| 
&pending.headers) {
@@ -11311,7 +11311,7 @@ async fn reconcile_partition_view_divergence<B, SB>(
     // Truncation is safe only above what this replica has *applied*, which is 
not
     // the view's commit point: a backup can sit above it.
     let announced_commit = pending.map_or(0, |pending| pending.commit_max);
-    let applied_floor = 
announced_commit.max(partition.consensus().commit_min());
+    let applied_floor = partition.consensus().commit_min();
 
     let mut repairable_from: Option<u64> = None;
     for canonical in pending.map_or(&[][..], |pending| &pending.headers) {
diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs
index 08719b44b..8077668ed 100644
--- a/core/simulator/src/lib.rs
+++ b/core/simulator/src/lib.rs
@@ -8042,3 +8042,88 @@ mod review_4092_dst_tests {
         );
     }
 }
+
+#[cfg(test)]
+mod probe_answer_divergence_tests {
+    //! End-to-end seeds for the probe-answer `StartView`. The mechanism 
itself is
+    //! pinned by `consensus::impls::probe_answer_tests`; these replay the 
runs that
+    //! found it.
+
+    use super::*;
+
+    /// Seed 144 of the uniform swarm lane: replica 0 prepared op 7 in view 0
+    /// without acks, view 1 truncated it and prepared a different op 7, and
+    /// replica 0 then adopted view 1 through a probe answer that carried no
+    /// canonical headers, so it kept its own op 7 and committed that instead.
+    ///
+    /// Network faults only, no crash needed, which is why it lands at op 7 and
+    /// replays fast. The mechanism is pinned separately by
+    /// `consensus::impls::probe_answer_tests`; this is the end-to-end seed.
+    #[test]
+    fn 
given_a_probe_adopted_view_when_the_head_diverges_should_not_commit_the_stale_entry()
 {
+        use crate::workload::{
+            self, FaultInjector, Workload,
+            options::{ActionWeights, WorkloadOptions},
+            oracle,
+        };
+        
server_common::MemoryPool::init_pool(&server_common::MemoryPoolSettings {
+            enabled: false,
+            size: iggy_common::IggyByteSize::from(0u64),
+            bucket_capacity: 1,
+        });
+
+        let replica_count: u8 = 3;
+        let client_id: u128 = 1;
+        let seed = 144;
+        // Swarm, not a fixed profile: the asymmetric partitions and clogs this
+        // seed draws are what let a primary prepare an op it cannot get acked.
+        let mut network_opts = packet::PacketSimulatorOptions::swarm(seed);
+        network_opts.node_count = replica_count;
+        network_opts.client_count = 1;
+        let mut sim = Simulator::new(
+            usize::from(replica_count),
+            std::iter::once(client_id),
+            network_opts,
+        );
+        let client = SimClient::new(client_id);
+        let ns = IggyNamespace::new(1, 1, 0);
+        sim.init_partition(ns);
+        sim.register_client_with_primary(&client);
+
+        let mut options = WorkloadOptions::new(seed, replica_count, vec![ns]);
+        options.weights = ActionWeights::uniform();
+        options.crash_per_tick_ratio = 0.01;
+        options.restart_per_tick_ratio = 0.05;
+        let mut wl = Workload::new(options);
+
+        let mut injector = FaultInjector::new(seed, replica_count);
+        let mut invariants = crate::workload::invariants::Invariants::new();
+        // `run_with_faults` runs the per-tick invariants and the live state
+        // checker, which is where the divergence fired.
+        workload::run_with_faults(
+            &mut sim,
+            &mut wl,
+            &[client],
+            6_000,
+            u64::MAX,
+            &mut injector,
+            &mut invariants,
+        );
+
+        assert!(
+            oracle::drive_to_quiesce(&mut sim, &mut wl, 50_000, &mut 
invariants),
+            "{}",
+            oracle::quiesce_failure_report(&sim, &wl),
+        );
+        assert!(
+            oracle::settle_to_stable_view(&mut sim, &mut wl, 50_000, &mut 
invariants),
+            "metadata views never converged after the drain"
+        );
+        let report = oracle::assert_converged(&sim, &mut wl);
+        assert!(
+            report.ops_compared > 0,
+            "no committed metadata op was witnessed on two replicas, so this 
seed \
+             would pass on a diverged cluster"
+        );
+    }
+}

Reply via email to