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() +}
