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

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

commit 06e11a116147e102119d3a86cce220bc51cc73ac
Author: Krishna Vishal <[email protected]>
AuthorDate: Sat Aug 15 00:10:03 2026 +0530

    feat(simulator): assert cross-replica committed log equality
    
    The harness compared each replica against itself (offsets and views never
    regress) and against the workload's shadow, but never against the OTHER
    replicas, which is the actual consensus property: two replicas that both
    committed op N must have committed the same op N. A divergence was
    invisible as long as every replica stayed internally consistent.
    
    `state_checker` follows TigerBeetle's `testing/cluster/state_checker.zig`
    rather than inventing a shape: one canonical commit chain, every replica
    asserted to agree with it wherever they overlap, and the chain asserted to
    be hash-linked so a replica cannot match op-by-op while having been built
    on a different prefix. It runs per tick alongside the other invariants, so
    a divergence is reported at the tick it appears rather than at the next
    quiesce.
    
    Equality is over the committed PREFIX, not equal heads: a replica that
    missed the last commit broadcast or rejoined recently may legitimately
    trail, and requiring equal heads would fail on ordinary lag while saying
    nothing more about safety.
    
    Two oracle checks were unsound once primaries can be crashed, and both
    failed on correct clusters rather than catching anything:
    
    The leader-relative partition check measured `commit_offset`, which is the
    highest durably PERSISTED message offset, not a committed watermark. A
    backup that persisted an op the view then settled below is ordinary VSR.
    It now measures the group's consensus `commit_min`. The check also needs a
    settled view to resolve a leader at all, so `settle_to_stable_view` waits
    for one agreed view on both planes before asserting, instead of guessing
    at a leader mid-view-change.
    
    The entity oracle filtered committed topics by their STREAM's name, so the
    harness's own `sim-topic-*` fillers counted as workload state once they
    landed inside a workload-created stream, and the shadow was blamed for not
    predicting them. Every level now filters on its own name.
    
    Non-vacuity is asserted, not assumed: the test requires ops witnessed on
    more than one replica, since an equality oracle that never finds a shared
    op passes in silence. A 4000-tick metadata run with crashes and restarts
    compares all 134 ops of its chain.
    
    Crash/restart sweeps: 24 of 24 without primary crashes, 17 of 20 with
    `--crash-primary`. The three remaining failures are production assertions
    (`advance_commit_min` sequentiality, and a `transmute_header` validity
    check), not oracle disagreements.
---
 core/simulator/src/bin/workload-fuzz.rs      |   9 ++
 core/simulator/src/lib.rs                    | 131 ++++++++++++++++
 core/simulator/src/workload/invariants.rs    |  19 ++-
 core/simulator/src/workload/mod.rs           |   1 +
 core/simulator/src/workload/oracle.rs        | 222 ++++++++++++++++++++++-----
 core/simulator/src/workload/state_checker.rs | 219 ++++++++++++++++++++++++++
 6 files changed, 559 insertions(+), 42 deletions(-)

diff --git a/core/simulator/src/bin/workload-fuzz.rs 
b/core/simulator/src/bin/workload-fuzz.rs
index c7abfecc5..c41d5be9f 100644
--- a/core/simulator/src/bin/workload-fuzz.rs
+++ b/core/simulator/src/bin/workload-fuzz.rs
@@ -437,6 +437,15 @@ fn main() {
             "{}",
             oracle::quiesce_failure_report(&sim, &workload),
         );
+        // Then wait for one agreed view before asserting. `assert_converged`
+        // resolves the leader as whichever live replica claims to be primary, 
so
+        // asserting mid-view-change either finds none or finds a deposed one 
--
+        // false failures rather than divergences.
+        assert!(
+            oracle::settle_to_stable_view(&mut sim, &mut workload, 50_000),
+            "metadata views never converged after the drain\n{}",
+            oracle::quiesce_failure_report(&sim, &workload),
+        );
         oracle::assert_converged(&sim, &workload);
         println!("quiesced and converged (leader-relative + entity oracle)");
         // Again after the drain: the drain both answers outstanding requests 
and
diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs
index 5d8a42f19..6223a4f42 100644
--- a/core/simulator/src/lib.rs
+++ b/core/simulator/src/lib.rs
@@ -141,6 +141,19 @@ impl SimReplica {
     }
 }
 
+/// One replica's view of a partition group's consensus. Read by the quiesce
+/// oracle; see [`Simulator::partition_consensus_state`].
+#[derive(Debug, Clone, Copy)]
+pub(crate) struct PartitionConsensusState {
+    pub status: consensus::Status,
+    pub view: u32,
+    pub is_primary: bool,
+    /// Ops this replica has committed in the group. The committed watermark, 
as
+    /// opposed to `PartitionOffsets::commit_offset`, which is the highest
+    /// durably PERSISTED message offset and so counts an uncommitted suffix 
too.
+    pub commit_min: u64,
+}
+
 pub struct Simulator {
     /// All replicas, indexed by replica id. Always fully populated — crashed
     /// replicas are kept alive but skipped during dispatch.
@@ -997,6 +1010,29 @@ impl Simulator {
         Some(u64::from(partition.consensus().view()))
     }
 
+    /// One replica's view of a partition group's consensus, or `None` when 
that
+    /// replica does not host the namespace.
+    ///
+    /// Read by the quiesce oracle to decide whether a group has settled into 
one
+    /// view, which its leader-relative checks depend on once partition 
primaries
+    /// can be crashed.
+    #[must_use]
+    pub(crate) fn partition_consensus_state(
+        &self,
+        replica_idx: usize,
+        namespace: IggyNamespace,
+    ) -> Option<PartitionConsensusState> {
+        let shard = self.replicas[replica_idx].partition_shard(namespace);
+        let partition = shard.plane.partitions().get_by_ns(&namespace)?;
+        let consensus = partition.consensus();
+        Some(PartitionConsensusState {
+            status: consensus.status(),
+            view: consensus.view(),
+            is_primary: consensus.is_primary(),
+            commit_min: consensus.commit_min(),
+        })
+    }
+
     /// Index of the current primary for `namespace`, as seen by the first live
     /// replica hosting it, or `None` if no live replica hosts it.
     ///
@@ -3208,6 +3244,101 @@ mod tests {
     /// pipeline to `PIPELINE_PREPARE_QUEUE_MAX`; a request on `ns_b`
     /// still commits while `ns_a` is wedged (no quorum without backup
     /// acks); lifting the block drains `ns_a` completely.
+    /// The cross-replica equality check actually compares replicas against 
each
+    /// other, and holds over a metadata workload with crashes and restarts.
+    ///
+    /// Non-vacuity is the point of the assertions on the chain. An equality
+    /// oracle that never finds two replicas at the same op passes in silence, 
so
+    /// a green run would say nothing: `ops_compared` counts only ops 
witnessed on
+    /// more than one replica, which is the subset that exercised the property.
+    #[test]
+    fn committed_metadata_agrees_across_replicas() {
+        use crate::workload::{
+            self, FaultInjector, Workload,
+            invariants::Invariants,
+            options::{ActionWeights, WorkloadOptions},
+            oracle,
+        };
+        
server_common::MemoryPool::init_pool(&server_common::MemoryPoolConfigOther {
+            enabled: false,
+            size: iggy_common::IggyByteSize::from(0u64),
+            bucket_capacity: 1,
+        });
+
+        let replica_count: u8 = 5;
+        let client_id: u128 = 1;
+        let seed = 0x57A7_E000;
+        let network_opts = packet::PacketSimulatorOptions {
+            node_count: replica_count,
+            client_count: 1,
+            seed,
+            ..packet::PacketSimulatorOptions::default()
+        };
+        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);
+
+        // Metadata ops, since the committed chain this checks is the metadata 
WAL.
+        // Crash and restart so replicas rejoin and repair, which is when a
+        // divergence would be introduced if one could be.
+        let mut options = WorkloadOptions::new(seed, replica_count, vec![ns]);
+        options.weights = ActionWeights::metadata_only();
+        options.crash_per_tick_ratio = 0.01;
+        options.restart_per_tick_ratio = 0.02;
+        let mut wl = Workload::new(options);
+
+        let clients = [client];
+        let mut injector = FaultInjector::new(seed, replica_count);
+        let mut invariants = Invariants::new();
+        // Driven here rather than through `workload::run` so the accumulated
+        // chain is readable afterwards; `run` builds its own `Invariants`.
+        for _ in 0..4_000u32 {
+            wl.tick();
+            injector.step(&mut sim, &wl);
+            workload::resubmit_due(&mut sim, &mut wl);
+            if let Some((target, msg)) = wl.build_request(&clients[0]) {
+                sim.submit_request(clients[0].client_id(), target, 
msg.into_generic());
+            }
+            for reply in sim.step() {
+                let cmds = wl.on_reply(&reply);
+                workload::apply_sim_commands(&mut sim, &cmds);
+            }
+            invariants.check(&sim, &wl);
+        }
+
+        assert!(
+            injector.restarts() > 0,
+            "no replica restarted, so rejoin and repair never ran"
+        );
+        let chain = invariants.state_checker();
+        assert!(
+            chain.chain_len() > 0,
+            "the canonical commit chain is empty: nothing was ever recorded"
+        );
+        assert!(
+            chain.ops_compared() > 0,
+            "no committed op was witnessed on two replicas, so the equality 
check \
+             never actually compared anything and would pass on a diverged 
cluster"
+        );
+
+        assert!(
+            oracle::drive_to_quiesce(&mut sim, &mut wl, 50_000),
+            "{}",
+            oracle::quiesce_failure_report(&sim, &wl),
+        );
+        assert!(
+            oracle::settle_to_stable_view(&mut sim, &mut wl, 50_000),
+            "metadata views never converged after the drain"
+        );
+        oracle::assert_converged(&sim, &wl);
+    }
+
     /// A restarted replica comes back with the partition data it had, so its
     /// `commit_offset` does not regress.
     ///
diff --git a/core/simulator/src/workload/invariants.rs 
b/core/simulator/src/workload/invariants.rs
index d7f1b8a8d..c4f57095d 100644
--- a/core/simulator/src/workload/invariants.rs
+++ b/core/simulator/src/workload/invariants.rs
@@ -24,6 +24,7 @@
 //! unchanged.
 
 use crate::Simulator;
+use crate::workload::state_checker::StateChecker;
 use crate::workload::{CLIENT_REQUEST_QUEUE_MAX, Workload};
 use server_common::sharding::IggyNamespace;
 use std::collections::HashMap;
@@ -34,6 +35,10 @@ use std::collections::HashMap;
 pub struct Invariants {
     commit_offset: HashMap<(u8, IggyNamespace), u64>,
     view: HashMap<(u8, IggyNamespace), u64>,
+    /// Cross-replica committed-log agreement. Lives here so it runs every tick
+    /// like the rest: a divergence is reported at the tick it appears rather 
than
+    /// at the next quiesce, by which point the run has moved on.
+    state_checker: StateChecker,
 }
 
 impl Invariants {
@@ -49,7 +54,9 @@ impl Invariants {
     /// - consensus `view` never regresses (a view change only advances it).
     ///
     /// Globally:
-    /// - total in-flight requests stay within the per-client queue ceiling.
+    /// - total in-flight requests stay within the per-client queue ceiling,
+    /// - live replicas agree on every committed metadata op they share, and 
the
+    ///   committed chain stays hash-linked (see [`StateChecker`]).
     ///
     /// Crashed replicas are skipped and their last-seen marks retained, which
     /// holds across a restart for both quantities: the superblock carries 
`view`,
@@ -94,6 +101,8 @@ impl Invariants {
              (client_count={}, queue_max={CLIENT_REQUEST_QUEUE_MAX}) 
(seed={seed:#x})",
             workload.options.client_count,
         );
+
+        self.state_checker.check(sim, seed);
     }
 
     /// Number of `(replica, namespace)` pairs observed so far. Used by tests 
to
@@ -103,6 +112,14 @@ impl Invariants {
     pub(crate) fn tracked_pairs(&self) -> usize {
         self.commit_offset.len()
     }
+
+    /// The canonical committed chain built so far. Tests read it to prove the
+    /// equality check compared replicas against each other rather than passing
+    /// over an empty chain.
+    #[must_use]
+    pub const fn state_checker(&self) -> &StateChecker {
+        &self.state_checker
+    }
 }
 
 /// Panic if `cur < prev`. Pure so the catch logic is unit-testable without a
diff --git a/core/simulator/src/workload/mod.rs 
b/core/simulator/src/workload/mod.rs
index 322d66c22..164367084 100644
--- a/core/simulator/src/workload/mod.rs
+++ b/core/simulator/src/workload/mod.rs
@@ -33,6 +33,7 @@ pub mod ops;
 pub mod options;
 pub mod oracle;
 pub mod shadow;
+pub mod state_checker;
 
 use crate::Simulator;
 use crate::client::SimClient;
diff --git a/core/simulator/src/workload/oracle.rs 
b/core/simulator/src/workload/oracle.rs
index 567172421..b1eba8e5d 100644
--- a/core/simulator/src/workload/oracle.rs
+++ b/core/simulator/src/workload/oracle.rs
@@ -23,22 +23,22 @@
 //!
 //! - no live replica is ahead of the leader on any namespace (a backup ahead 
of
 //!   the leader is a split-brain / divergence bug),
+//! - every live replica agrees with every other on each committed metadata op
+//!   they both hold, the real consensus property (see
+//!   [`super::state_checker`]),
 //! - on a serial run, the workload's predicted [`Shadow`] equals the metadata
 //!   committed on the leader, the payoff of the name-keyed shadow.
 //!
-//! Full cross-replica EQUALITY (every live replica holding the same committed
-//! log) is the real consensus property, but it is not asserted yet. Backups
-//! apply prepares in strict order and drop any gap (`op != current_op + 1` in
-//! `metadata::on_replicate` / `iggy_partition`), relying on the primary's
-//! retransmit and the repair sessions (`MetadataRepairSession` / partition
-//! `RepairSession`) to refill. Message repair has landed on both planes, so
-//! the equality assert is unblocked but not yet re-enabled: the sim must
-//! first drive quiesce long enough for repair rounds to converge.
+//! Equality is asserted over the committed PREFIX, not over equal heads: a
+//! replica that missed the last commit broadcast, or that rejoined recently, 
may
+//! legitimately trail. What it may not do is hold different history at an op 
it
+//! did commit. Requiring equal heads instead would fail on ordinary lag and 
say
+//! nothing extra about safety.
 
 use crate::Simulator;
 use crate::replica::Replica;
 use crate::workload::shadow::Shadow;
-use crate::workload::{Workload, apply_sim_commands, resubmit_due};
+use crate::workload::{Workload, apply_sim_commands, resubmit_due, 
state_checker};
 use consensus::{MetadataHandle, Status};
 use metadata::impls::metadata::StreamsFrontend;
 use std::collections::BTreeSet;
@@ -72,14 +72,25 @@ struct CommittedMetadata {
 impl CommittedMetadata {
     /// Restrict to workload-generated entities (see [`WORKLOAD_PREFIX`]), so 
the
     /// entity oracle compares like with like against the shadow.
+    ///
+    /// Every level is filtered on its OWN name, not on its stream's. The 
harness
+    /// seeds filler topics and partitions of its own (`sim-topic-*`, see
+    /// `Streams::seed_namespace`) to keep slab ids dense, and once the 
workload
+    /// has created enough streams those fillers land inside a stream named
+    /// `wl-...`. Filtering topics by their stream alone then admits harness 
state
+    /// into the comparison and the shadow is blamed for not predicting it.
     fn workload_owned(mut self) -> Self {
         self.streams
             .retain(|name| name.starts_with(WORKLOAD_PREFIX));
-        self.topics
-            .retain(|(stream, _)| stream.starts_with(WORKLOAD_PREFIX));
+        self.topics.retain(|(stream, topic)| {
+            stream.starts_with(WORKLOAD_PREFIX) && 
topic.starts_with(WORKLOAD_PREFIX)
+        });
         self.users.retain(|name| name.starts_with(WORKLOAD_PREFIX));
-        self.consumer_groups
-            .retain(|(stream, _, _)| stream.starts_with(WORKLOAD_PREFIX));
+        self.consumer_groups.retain(|(stream, topic, group)| {
+            stream.starts_with(WORKLOAD_PREFIX)
+                && topic.starts_with(WORKLOAD_PREFIX)
+                && group.starts_with(WORKLOAD_PREFIX)
+        });
         self
     }
 }
@@ -208,28 +219,141 @@ pub fn quiesce_failure_report(sim: &Simulator, workload: 
&Workload) -> String {
     report
 }
 
-/// Post-drain consensus checks that hold today.
+/// Step until every live replica's metadata plane is `Normal` in one shared
+/// view, or `max_ticks` elapses.
 ///
-/// Asserts no live replica is ahead of the leader, and (on a serial run) that
-/// the shadow equals the metadata committed on the leader. See the module docs
-/// for why full cross-replica equality is deferred.
+/// Needed before [`assert_converged`], which resolves the leader as "the live
+/// replica whose metadata consensus says it is primary". With primaries spared
+/// from crashes that was always the same replica in view 0. Once a primary can
+/// be crashed, live replicas transiently hold different views and there may 
be no
+/// `Normal` primary at all, so the leader lookup fails or names a deposed one 
--
+/// a false failure, not a divergence. Waiting for one view removes the 
ambiguity
+/// rather than guessing at a leader.
 ///
-/// Assumes one stable primary that every live replica agrees on: the leader is
-/// `Simulator::primary_index` (a single replica's view), and both checks treat
-/// it as the authoritative, most-advanced log. Sound today because the driver
-/// spares primaries from crashes, so no view change runs mid-test. Once
-/// primary-crash injection lands, live replicas can hold different views and
-/// this breaks: it may pick a stale or crashed leader (a correctly-ahead new
-/// primary then trips "exceeds leader"), or find no `Normal` primary mid-view
-/// change (`metadata_leader` returns `None`). Both are false failures. Fix
-/// then: resolve the leader by highest `(view, commit_offset)`, or quiesce
-/// until live replicas reconverge to one view before asserting. Crash 
injection
-/// already runs but spares primaries (`maybe_inject_crash`); this is deferred
-/// until primary-crash injection lands, itself gated on a request-resend path.
+/// Returns `false` if the views never converge, which is a real liveness 
failure
+/// and the caller should report it rather than assert against an unsettled
+/// cluster.
+#[must_use]
+pub fn settle_to_stable_view(sim: &mut Simulator, workload: &mut Workload, 
max_ticks: u64) -> bool {
+    for _ in 0..max_ticks {
+        if views_are_settled(sim, workload) {
+            return true;
+        }
+        workload.tick();
+        resubmit_due(sim, workload);
+        for reply in sim.step() {
+            let cmds = workload.on_reply(&reply);
+            apply_sim_commands(sim, &cmds);
+        }
+    }
+    views_are_settled(sim, workload)
+}
+
+/// Both planes settled: the metadata group and every tracked partition group.
+///
+/// The partition half matters for the leader-relative offset check in
+/// [`assert_converged`], which resolves its leader from 
`Simulator::primary_index`
+/// -- a single replica's view of that group's primary. Each partition group 
runs
+/// its own view change, so settling only the metadata plane leaves that check
+/// asserting against a deposed partition leader, and a correctly-ahead new one
+/// then trips "exceeds leader".
+fn views_are_settled(sim: &Simulator, workload: &Workload) -> bool {
+    if !metadata_view_is_settled(sim) {
+        return false;
+    }
+    workload
+        .options
+        .namespaces
+        .iter()
+        .all(|&ns| partition_view_is_settled(sim, ns))
+}
+
+/// True when every live replica's metadata consensus is `Normal` in the same
+/// view and exactly one of them claims to be primary.
+fn metadata_view_is_settled(sim: &Simulator) -> bool {
+    let mut view = None;
+    let mut primaries = 0usize;
+    let mut live = 0usize;
+    for replica_idx in 0..sim.replica_count {
+        if sim.is_crashed(replica_idx) {
+            continue;
+        }
+        let Some(consensus) = sim.replicas[usize::from(replica_idx)].shards[0]
+            .plane
+            .metadata()
+            .consensus
+            .as_ref()
+        else {
+            continue;
+        };
+        live += 1;
+        if consensus.status() != Status::Normal {
+            return false;
+        }
+        match view {
+            Some(agreed) if agreed != consensus.view() => return false,
+            Some(_) => {}
+            None => view = Some(consensus.view()),
+        }
+        if consensus.is_primary() {
+            primaries += 1;
+        }
+    }
+    live > 0 && primaries == 1
+}
+
+/// True when every live replica hosting `ns` has that group `Normal` in one
+/// shared view with exactly one primary.
+///
+/// A replica that does not host the namespace is skipped rather than treated 
as
+/// disagreement: a group only materialises on its hash-owning shard, and 
after a
+/// restart a group whose stream the workload deleted is not re-materialised at
+/// all.
+fn partition_view_is_settled(sim: &Simulator, ns: 
server_common::sharding::IggyNamespace) -> bool {
+    let mut view = None;
+    let mut primaries = 0usize;
+    let mut hosts = 0usize;
+    for replica_idx in 0..sim.replica_count {
+        if sim.is_crashed(replica_idx) {
+            continue;
+        }
+        let Some(state) = 
sim.partition_consensus_state(usize::from(replica_idx), ns) else {
+            continue;
+        };
+        hosts += 1;
+        if state.status != Status::Normal {
+            return false;
+        }
+        match view {
+            Some(agreed) if agreed != state.view => return false,
+            Some(_) => {}
+            None => view = Some(state.view),
+        }
+        if state.is_primary {
+            primaries += 1;
+        }
+    }
+    // No live host is settled by default: there is no leader for the offset 
check
+    // to resolve either, and it skips the namespace for the same reason.
+    hosts == 0 || primaries == 1
+}
+
+/// Post-drain consensus checks.
+///
+/// Asserts no live replica is ahead of the leader, that every live replica 
agrees
+/// with every other on each committed metadata op they share, and (on a serial
+/// run) that the shadow equals the metadata committed on the leader.
+///
+/// Assumes one stable primary that every live replica agrees on, which
+/// [`settle_to_stable_view`] establishes and which callers should run first 
once
+/// primaries can be crashed. Without it this may pick a deposed leader (a
+/// correctly-ahead new primary then trips "exceeds leader") or find no 
`Normal`
+/// primary at all mid-view-change. Both are false failures.
 ///
 /// # Panics
-/// If a replica is ahead of the leader or the shadow mismatches the leader. 
The
-/// workload seed is in the message so the failing run replays 
deterministically.
+/// If a replica is ahead of the leader, two replicas disagree on a committed 
op,
+/// or the shadow mismatches the leader. The workload seed is in the message 
so the
+/// failing run replays deterministically.
 pub fn assert_converged(sim: &Simulator, workload: &Workload) {
     let seed = workload.options.seed;
     let live: Vec<usize> = (0..sim.replica_count)
@@ -241,25 +365,41 @@ pub fn assert_converged(sim: &Simulator, workload: 
&Workload) {
         "no live replicas at quiesce (seed={seed:#x})"
     );
 
-    // Safety direction: no live replica may be ahead of the leader on any
-    // namespace. A backup may trail (no idle catch-up yet, see module docs),
-    // but a backup whose commit_offset exceeds the leader's is a divergence.
+    // The consensus property proper: replicas compared against each other, not
+    // against the workload's expectations. Runs before the leader-relative 
checks
+    // because a genuine divergence explains any leader confusion below it.
+    state_checker::assert_committed_prefixes_agree(sim, seed);
+
+    // Safety direction: no live replica may have COMMITTED more of a group 
than
+    // its leader has. A backup may trail (no idle catch-up yet, see module 
docs),
+    // but a backup committed past the leader is a divergence.
+    //
+    // Measured on the group's consensus `commit_min`, not on
+    // `PartitionOffsets::commit_offset`. The latter is the highest durably
+    // persisted message offset, which counts an uncommitted suffix: a backup 
that
+    // persisted op N while the view that elected the new leader settled on 
N-1 is
+    // ordinary VSR, not divergence. That only stayed invisible while primaries
+    // were spared, since a never-crashed primary is always the furthest ahead;
+    // with primary crashes it fires on a correct cluster.
     for &ns in &workload.options.namespaces {
         let Some(leader) = sim.primary_index(ns) else {
             continue;
         };
-        let Some(leader_offset) = sim
-            .offsets(usize::from(leader), ns)
-            .map(|o| o.commit_offset)
+        let Some(leader_committed) = sim
+            .partition_consensus_state(usize::from(leader), ns)
+            .map(|state| state.commit_min)
         else {
             continue;
         };
         for &replica_idx in &live {
-            if let Some(offset) = sim.offsets(replica_idx, ns).map(|o| 
o.commit_offset) {
+            if let Some(committed) = sim
+                .partition_consensus_state(replica_idx, ns)
+                .map(|state| state.commit_min)
+            {
                 assert!(
-                    offset <= leader_offset,
-                    "replica {replica_idx} commit_offset {offset} exceeds 
leader {leader} \
-                     ({leader_offset}) on ns {ns:?} at quiesce 
(seed={seed:#x})",
+                    committed <= leader_committed,
+                    "replica {replica_idx} committed {committed} ops exceeds 
leader {leader} \
+                     ({leader_committed}) on ns {ns:?} at quiesce 
(seed={seed:#x})",
                 );
             }
         }
diff --git a/core/simulator/src/workload/state_checker.rs 
b/core/simulator/src/workload/state_checker.rs
new file mode 100644
index 000000000..a3c5135f8
--- /dev/null
+++ b/core/simulator/src/workload/state_checker.rs
@@ -0,0 +1,219 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Cross-replica committed-log equality, after TigerBeetle's
+//! `testing/cluster/state_checker.zig`.
+//!
+//! The per-tick checks in [`super::invariants`] catch a single replica
+//! contradicting itself, and [`super::oracle`] compares committed metadata
+//! against the workload's shadow. Neither compares replicas to EACH OTHER, 
which
+//! is the actual consensus property: two replicas that both committed op N 
must
+//! have committed the same op N.
+//!
+//! Modelled on TigerBeetle's checker rather than invented: it keeps one 
canonical
+//! commit chain, asserts every replica agrees with it wherever they overlap
+//! (`(commit_a == commit_b) == (checksum_a == checksum_b)`), and asserts the 
chain
+//! is hash-linked (`header_b.parent == checksum_a`). Recording which replicas
+//! reached each op also makes the check provably non-vacuous, which matters:
+//! a chain nothing was ever compared against passes silently.
+
+use crate::Simulator;
+use consensus::MetadataHandle;
+use iggy_binary_protocol::PrepareHeader;
+use journal::Journal;
+use std::collections::{BTreeMap, BTreeSet};
+
+/// One op of the canonical committed chain.
+#[derive(Debug)]
+struct CanonicalCommit {
+    /// Identity of the prepare committed at this op. Two replicas disagreeing
+    /// here is a divergence: the same log position holds different history.
+    ///
+    /// The chain's hash link is checked against this rather than against a 
stored
+    /// `parent`: an arriving header's `parent` must equal the canonical 
previous
+    /// op's `checksum`, so keeping each entry's own parent as well would 
record a
+    /// value nothing ever reads.
+    checksum: u128,
+    /// Replicas observed committing this op, so the check can prove it 
compared
+    /// something rather than passing over an empty chain.
+    replicas: BTreeSet<u8>,
+}
+
+/// Canonical committed metadata chain, accumulated across ticks.
+#[derive(Debug, Default)]
+pub struct StateChecker {
+    commits: BTreeMap<u64, CanonicalCommit>,
+    /// Highest op already verified per replica, so a tick only walks what is
+    /// new. A high-water mark, never lowered: a restart recovers its commit
+    /// point from a lower bound (`SimJournal::recovery_commit_watermark`) and 
so
+    /// may report a smaller `commit_min` than it did before, and re-verifying
+    /// that prefix every tick would make this quadratic for no gain.
+    verified_upto: BTreeMap<u8, u64>,
+}
+
+impl StateChecker {
+    #[must_use]
+    pub fn new() -> Self {
+        Self::default()
+    }
+
+    /// Fold every live replica's newly committed metadata ops into the 
canonical
+    /// chain, asserting agreement.
+    ///
+    /// Reads committed state only (ops at or below `commit_min`), so a prepare
+    /// still in flight is never compared: replicas are allowed to disagree 
about
+    /// uncommitted tails, and that is what a view change resolves.
+    ///
+    /// Crashed replicas are skipped rather than dropped: their high-water 
mark is
+    /// kept, so a restart re-verifies only what it commits anew.
+    ///
+    /// # Panics
+    /// On any disagreement about a committed op, or a broken hash chain. The
+    /// message names both replicas and the op, and the seed replays the run.
+    pub fn check(&mut self, sim: &Simulator, seed: u64) {
+        for replica_idx in 0..sim.replica_count {
+            if sim.is_crashed(replica_idx) {
+                continue;
+            }
+            let replica = &sim.replicas[usize::from(replica_idx)];
+            let Some(consensus) = 
replica.shards[0].plane.metadata().consensus.as_ref() else {
+                continue;
+            };
+            let committed = consensus.commit_min();
+            let verified = 
self.verified_upto.get(&replica_idx).copied().unwrap_or(0);
+            for op in (verified + 1)..=committed {
+                // Absent header at a committed op: the sim never checkpoints, 
so
+                // nothing drains the prefix and this would be a real hole. 
Left to
+                // the journal's own invariants rather than asserted here, 
since a
+                // recovered replica legitimately reports a commit point one 
above
+                // its head (the watermark is a lower bound).
+                let Some(header) = journaled_header(replica, op) else {
+                    continue;
+                };
+                self.record(replica_idx, op, &header, seed);
+            }
+            self.verified_upto
+                .insert(replica_idx, verified.max(committed));
+        }
+    }
+
+    /// Number of ops in the canonical chain. Tests assert this is non-zero, 
so a
+    /// green run cannot mean "never compared anything".
+    #[must_use]
+    pub fn chain_len(&self) -> usize {
+        self.commits.len()
+    }
+
+    /// Ops witnessed by more than one replica. The only ops that actually
+    /// exercised the equality property: an op only ever seen on one replica 
was
+    /// recorded, never compared.
+    #[must_use]
+    pub fn ops_compared(&self) -> usize {
+        self.commits
+            .values()
+            .filter(|commit| commit.replicas.len() > 1)
+            .count()
+    }
+
+    fn record(&mut self, replica_idx: u8, op: u64, header: &PrepareHeader, 
seed: u64) {
+        // Hash-chain link, checked before the identity comparison so a 
diverged
+        // prefix is reported at the op where the chains part rather than at 
the
+        // first op whose contents happen to differ.
+        if let Some(previous) = self.commits.get(&(op - 1))
+            && header.parent != previous.checksum
+        {
+            panic!(
+                "replica {replica_idx} committed op {op} whose parent {:#x} is 
not the \
+                 canonical op {} checksum {:#x}: its committed history forked 
below this \
+                 op (seed={seed:#x})",
+                header.parent,
+                op - 1,
+                previous.checksum,
+            );
+        }
+        match self.commits.get_mut(&op) {
+            Some(canonical) => {
+                assert_eq!(
+                    canonical.checksum, header.checksum,
+                    "replicas disagree on committed op {op}: canonical 
checksum {:#x} \
+                     (committed by {:?}) vs replica {replica_idx}'s {:#x}. Two 
replicas \
+                     committed different history at the same log position 
(seed={seed:#x})",
+                    canonical.checksum, canonical.replicas, header.checksum,
+                );
+                canonical.replicas.insert(replica_idx);
+            }
+            None => {
+                self.commits.insert(
+                    op,
+                    CanonicalCommit {
+                        checksum: header.checksum,
+                        replicas: BTreeSet::from([replica_idx]),
+                    },
+                );
+            }
+        }
+    }
+}
+
+/// Assert every live replica's committed metadata prefix agrees, op for op.
+///
+/// The quiesce-time counterpart to [`StateChecker::check`]: that one folds 
ops in
+/// as they commit and so compares whatever happened to overlap, while this 
walks
+/// the full committed prefix of every live replica at rest and requires the
+/// shorter to be a genuine PREFIX of the longer. A replica may still trail (it
+/// may have missed the last commit broadcast), but where it has committed
+/// anything it must match.
+///
+/// # Panics
+/// If two live replicas disagree on any committed op.
+pub fn assert_committed_prefixes_agree(sim: &Simulator, seed: u64) {
+    let mut canonical: BTreeMap<u64, (u128, u8)> = BTreeMap::new();
+    for replica_idx in 0..sim.replica_count {
+        if sim.is_crashed(replica_idx) {
+            continue;
+        }
+        let replica = &sim.replicas[usize::from(replica_idx)];
+        let Some(consensus) = 
replica.shards[0].plane.metadata().consensus.as_ref() else {
+            continue;
+        };
+        for op in 1..=consensus.commit_min() {
+            let Some(header) = journaled_header(replica, op) else {
+                continue;
+            };
+            match canonical.get(&op) {
+                Some(&(checksum, owner)) => assert_eq!(
+                    checksum, header.checksum,
+                    "at quiesce replica {replica_idx} and replica {owner} 
disagree on \
+                     committed metadata op {op}: {:#x} vs {checksum:#x} 
(seed={seed:#x})",
+                    header.checksum,
+                ),
+                None => {
+                    canonical.insert(op, (header.checksum, replica_idx));
+                }
+            }
+        }
+    }
+}
+
+/// The header a replica has journaled at `op`, if any.
+///
+/// Reads shard 0's retained metadata WAL, which is where the committed 
metadata
+/// log lives; the journal is harness-owned so this also works across a 
restart.
+fn journaled_header(replica: &crate::SimReplica, op: u64) -> 
Option<PrepareHeader> {
+    let slot = usize::try_from(op).ok()?;
+    replica.metadata_journal.header(slot).copied()
+}

Reply via email to