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 43d293c986dd308a63fdc2a1fe20d97dbfb4a5f8
Author: Krishna Vishal <[email protected]>
AuthorDate: Fri Aug 14 14:56:32 2026 +0530

    feat(simulator): inject replica restarts with stability windows
    
    Crash injection was permanent: a crashed replica never came back, so
    nothing in the harness ever exercised rejoin, the view probe, or log
    repair. The reason was the same missing resend path that forced primaries
    to be spared, and that landed in the previous commit.
    
    Crash and restart now run through one FaultInjector owning the fault
    PRNG, with TigerBeetle-style stability windows: a crash lasts long enough
    for the surviving primary to commit past the victim's log, so its rejoin
    has something to repair, and a rejoined replica runs long enough to catch
    up before becoming a candidate again. Restart is considered before crash
    so a tick never both revives and kills. `spare_primary` becomes a knob
    rather than an unconditional rule, since crashing the primary under live
    traffic is what puts a view change on the critical path.
    
    Restart injection is metadata-plane only for now. The harness retains the
    metadata WAL and the superblocks across a restart but not partition data,
    so a restarted replica reports commit_offset 0 and the monotonicity
    invariant trips on a harness limitation rather than a real regression.
    The invariant doc now says so instead of claiming both quantities survive
    a restart. Retaining the partition journal per namespace is the faithful
    fix, since a real server's segments are on disk.
    
    Two tests pin the rejoin paths that do work: a lost PrepareOk recovers by
    retransmit, and a prepare journaled but unacked commits after its backups
    restart. Both are worth stating because the backup's admission check
    drops a prepare it already holds, so neither outcome follows from reading
    that check alone.
    
    The fuzzer finds a metadata-plane wedge these do not explain
    (--seed 42 --replicas 5 --plane metadata --crash-prob 0.01
    --restart-prob 0.02): two live backups hold op 45, the primary sits at
    commit 44, and both log the gap drop for the whole drain while the
    client's request goes unanswered across 276 attempts. Neither reduction
    above reproduces it, so it stays an open finding rather than a test.
---
 core/simulator/src/bin/workload-fuzz.rs   |  66 ++++++---
 core/simulator/src/lib.rs                 | 221 ++++++++++++++++++++++++++++++
 core/simulator/src/workload/auditor.rs    |   8 ++
 core/simulator/src/workload/invariants.rs |  16 ++-
 core/simulator/src/workload/mod.rs        | 221 +++++++++++++++++++++++-------
 core/simulator/src/workload/options.rs    |  39 +++++-
 core/simulator/src/workload/oracle.rs     |  44 +++++-
 7 files changed, 542 insertions(+), 73 deletions(-)

diff --git a/core/simulator/src/bin/workload-fuzz.rs 
b/core/simulator/src/bin/workload-fuzz.rs
index 4ac434628..c7abfecc5 100644
--- a/core/simulator/src/bin/workload-fuzz.rs
+++ b/core/simulator/src/bin/workload-fuzz.rs
@@ -17,15 +17,17 @@
 
 //! Deterministic workload fuzzer for the Iggy simulator.
 //!
-//! Drives [`simulator::workload::run`] (per-tick invariants + optional crash
-//! injection) for a number of ticks, then optionally quiesces and asserts the
-//! Phase C consensus checks. Everything is a function of `--seed`, logged at
-//! start and on panic so any failure replays with `--seed <value>`.
+//! Drives [`simulator::workload::run_with_faults`] (per-tick invariants plus
+//! crash, restart and network fault injection) for a number of ticks, then
+//! optionally quiesces and asserts the Phase C consensus checks. Everything 
is a
+//! function of `--seed`, logged at start and on panic so any failure replays
+//! with `--seed <value>`.
 //!
 //! ```text
 //! workload-fuzz [--seed N] [--ticks N] [--clients N] [--replicas N]
 //!               [--plane partition|metadata|mixed|uniform]
-//!               [--faults none|light|heavy] [--crash-prob F] [--no-quiesce]
+//!               [--faults none|light|heavy] [--no-quiesce]
+//!               [--crash-prob F] [--restart-prob F] [--crash-primary]
 //!               [network overrides: --packet-loss, --replay, 
--partition-mode,
 //!                --partition-prob, --unpartition-prob, --clog-prob, ...]
 //! ```
@@ -48,7 +50,7 @@ use simulator::client::SimClient;
 use simulator::packet::{PacketSimulatorOptions, PartitionMode, 
PartitionSymmetry};
 use simulator::workload::actions::Action;
 use simulator::workload::options::{ActionWeights, WorkloadOptions};
-use simulator::workload::{Workload, oracle, run};
+use simulator::workload::{FaultInjector, Workload, oracle, run_with_faults};
 use strum::IntoEnumIterator;
 
 #[derive(Parser)]
@@ -70,8 +72,16 @@ struct Args {
     /// `NoAck`. `1.0` keeps every offset op on the replicated path.
     #[arg(long, default_value_t = 0.5, value_parser = parse_unit_interval)]
     ack_quorum_ratio: f32,
+    /// Per-tick chance one eligible replica is crashed.
     #[arg(long, default_value_t = 0.0, value_parser = parse_unit_interval)]
     crash_prob: f32,
+    /// Per-tick chance one crashed replica is restarted. Without this a crash
+    /// is permanent and nothing exercises rejoin or log repair.
+    #[arg(long, default_value_t = 0.0, value_parser = parse_unit_interval)]
+    restart_prob: f32,
+    /// Crash the primary too, putting a view change under live traffic.
+    #[arg(long)]
+    crash_primary: bool,
     #[arg(long)]
     no_quiesce: bool,
 
@@ -123,12 +133,11 @@ struct Args {
 /// Named network fault profile, in the spirit of TigerBeetle's VOPR modes: one
 /// flag for "how hostile is the network", rather than eleven.
 ///
-/// Anything but [`Faults::None`] currently stalls the run, and not because the
-/// cluster fails to make progress: `SimClient` has no request timeout, so a
-/// client holds its single in-flight slot forever once the request or its 
reply
-/// is dropped. Every profile below is therefore write-once-read-later until 
the
-/// client grows a resend path; they are wired now so the fault space is
-/// described in one place rather than rediscovered later.
+/// Progress falls off steeply with severity, because every lost frame costs a
+/// resend timeout: on one namespace with one client, a 3-replica cluster 
drains
+/// roughly 440 replies in 5000 ticks on a perfect network, 240 under `light` 
and
+/// 40 under `heavy`. All three still drain and converge; budget ticks
+/// accordingly rather than reading a low reply count as a stall.
 #[derive(Clone, Copy, Debug, PartialEq, Eq, ValueEnum)]
 enum Faults {
     /// Perfect network. Delays only, no loss and no partitions.
@@ -390,16 +399,33 @@ fn main() {
     let mut options = WorkloadOptions::new(seed, replicas, vec![ns]);
     options.client_count = clients;
     options.crash_per_tick_ratio = crash_prob;
+    options.restart_per_tick_ratio = args.restart_prob;
+    options.spare_primary = !args.crash_primary;
     options.ack_quorum_ratio = args.ack_quorum_ratio;
     options.weights = plane.weights();
     let mut workload = Workload::new(options);
 
-    let replies = run(&mut sim, &mut workload, &sim_clients, ticks, u64::MAX);
+    let mut injector = FaultInjector::new(seed, replicas);
+    let replies = run_with_faults(
+        &mut sim,
+        &mut workload,
+        &sim_clients,
+        ticks,
+        u64::MAX,
+        &mut injector,
+    );
     println!(
-        "ran {ticks} ticks; {replies} replies; crashed replicas: {}",
-        sim.crashed.len()
+        "ran {ticks} ticks; {replies} replies; crashes={} restarts={} still 
down: {}",
+        injector.crashes(),
+        injector.restarts(),
+        sim.crashed.len(),
     );
 
+    // Printed before the quiesce assert, so a failed drain still reports what
+    // the run managed to do. Reading it after the assert meant the failure 
that
+    // most needs the numbers is the one that never shows them.
+    print_coverage(&workload);
+
     if quiesce {
         // A failed drain is a hard failure, not a warning. It used to be one
         // because a lost request could not be retried, so a stall was expected
@@ -413,8 +439,16 @@ fn main() {
         );
         oracle::assert_converged(&sim, &workload);
         println!("quiesced and converged (leader-relative + entity oracle)");
+        // Again after the drain: the drain both answers outstanding requests 
and
+        // issues its own resends, so the pre-drain numbers are not the final 
ones.
+        print_coverage(&workload);
     }
 
+    println!("workload-fuzz: OK (seed={seed})");
+}
+
+/// Reply, rejection and resend counters plus per-action commits.
+fn print_coverage(workload: &Workload) {
     let stats = workload.auditor.stats();
     println!(
         "coverage: replies_seen={} replies_unknown={} committed_rejections={} \
@@ -431,6 +465,4 @@ fn main() {
             println!("  {action:?}: {commits} commits");
         }
     }
-
-    println!("workload-fuzz: OK (seed={seed})");
 }
diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs
index 3b8d621ae..7c38b1668 100644
--- a/core/simulator/src/lib.rs
+++ b/core/simulator/src/lib.rs
@@ -3135,6 +3135,227 @@ 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.
+    /// A lost `PrepareOk` does not wedge the metadata plane: once the acks 
flow
+    /// again the primary reaches its commit quorum without client involvement.
+    ///
+    /// Worth pinning because the obvious reading of the code says otherwise. 
The
+    /// backup's admission check in `metadata::on_replicate` is a single
+    /// `header.op != current_op + 1`, so a prepare the backup ALREADY HOLDS is
+    /// dropped exactly like a forward gap, logging "dropping out-of-order
+    /// prepare (gap)" — and the primary's retransmit of an unacked prepare is
+    /// precisely such a duplicate. Reading only that check, a lost ack should
+    /// livelock: the primary retransmits forever and every retransmit is
+    /// dropped.
+    ///
+    /// It recovers anyway, so recovery does not depend on that check accepting
+    /// the duplicate, and those gap warnings are benign rather than evidence 
of
+    /// a stall. This test exists to keep that distinction honest: it fails if
+    /// recovery ever does come to rest on the retransmit being re-acked.
+    ///
+    /// Drops acks rather than crashing anyone, so the property under test is
+    /// about lost acks in general, not about restart recovery.
+    #[test]
+    fn lost_prepare_ok_is_recovered_by_retransmit() {
+        use iggy_binary_protocol::Command2;
+
+        
server_common::MemoryPool::init_pool(&server_common::MemoryPoolConfigOther {
+            enabled: false,
+            size: iggy_common::IggyByteSize::from(0u64),
+            bucket_capacity: 1,
+        });
+
+        let replica_count: u8 = 3;
+        let client_id: u128 = 1;
+        let network_opts = packet::PacketSimulatorOptions {
+            node_count: replica_count,
+            client_count: 1,
+            seed: 0x5EED_0077,
+            ..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);
+        sim.register_client_with_primary(&client);
+
+        let committed_before = metadata_commit(&sim, 0);
+
+        // Drop only PrepareOk on both backup links. Everything else still 
flows,
+        // so the backups receive and journal the prepare; only the primary's
+        // evidence of that is lost.
+        for backup in 1..replica_count {
+            sim.network
+                .link_filter_mut(ProcessId::Replica(backup), 
ProcessId::Replica(0))
+                .remove(Command2::PrepareOk);
+        }
+
+        let msg = client.create_stream("wl-lost-ack");
+        sim.submit_request(client_id, 0, msg.into_generic());
+
+        // Long enough for the prepare to reach and be journaled by both 
backups
+        // while the primary sees no acks.
+        for _ in 0..200 {
+            sim.step();
+        }
+        assert_eq!(
+            metadata_commit(&sim, 0),
+            committed_before,
+            "the primary must not commit while every backup ack is dropped"
+        );
+        for backup in 1..replica_count {
+            assert!(
+                metadata_op(&sim, usize::from(backup)) > committed_before,
+                "backup {backup} must have journaled the prepare, else this 
test \
+                 proves nothing about a LOST ack"
+            );
+        }
+
+        // Restore the acks. From here the primary's retransmit is the only 
route
+        // to a commit, which is exactly the mechanism under test.
+        for backup in 1..replica_count {
+            sim.network
+                .link_filter_mut(ProcessId::Replica(backup), 
ProcessId::Replica(0))
+                .insert(Command2::PrepareOk);
+        }
+
+        for _ in 0..5_000 {
+            sim.step();
+            if metadata_commit(&sim, 0) > committed_before {
+                return;
+            }
+        }
+        panic!(
+            "metadata commit stuck at {} after 5000 ticks with healthy links: 
a lost \
+             PrepareOk is no longer recovered, so the backup's gap check is 
now \
+             swallowing the primary's retransmit",
+            metadata_commit(&sim, 0),
+        );
+    }
+
+    /// A prepare that every backup journaled but never acked still commits 
after
+    /// those backups restart.
+    ///
+    /// The rejoin path is what recovers it: a restarted replica comes back 
with
+    /// `current_op` at N from its own WAL, rejoins as a probing backup
+    /// (`Status::Recovering`, see `new_shard`), and its probe draws a targeted
+    /// `StartView` from the primary that returns it to `Normal` and gets the
+    /// tail acked. The primary's own retransmit cannot do it: the backup's
+    /// admission check in `metadata::on_replicate` is a single
+    /// `header.op != current_op + 1`, so a prepare it already holds is dropped
+    /// like a forward gap.
+    ///
+    /// Pinned because the fuzzer finds a run where this recovery does NOT 
happen
+    /// (see the module-level note on seed 42): there, two live backups hold 
op 45
+    /// with the primary at commit 44, and both log the gap drop for the whole
+    /// drain. This test covers the case that works, so a regression here 
narrows
+    /// where that one diverges rather than leaving both unexplained.
+    #[test]
+    fn unacked_prepare_commits_after_the_backup_restarts() {
+        use iggy_binary_protocol::Command2;
+
+        
server_common::MemoryPool::init_pool(&server_common::MemoryPoolConfigOther {
+            enabled: false,
+            size: iggy_common::IggyByteSize::from(0u64),
+            bucket_capacity: 1,
+        });
+
+        let replica_count: u8 = 3;
+        let client_id: u128 = 1;
+        let network_opts = packet::PacketSimulatorOptions {
+            node_count: replica_count,
+            client_count: 1,
+            seed: 0x5EED_0078,
+            ..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);
+        sim.register_client_with_primary(&client);
+        let committed_before = metadata_commit(&sim, 0);
+
+        // Lose every backup ack, so the prepare is journaled cluster-wide 
while
+        // the primary stays one short of its commit quorum.
+        for backup in 1..replica_count {
+            sim.network
+                .link_filter_mut(ProcessId::Replica(backup), 
ProcessId::Replica(0))
+                .remove(Command2::PrepareOk);
+        }
+
+        let msg = client.create_stream("wl-unacked");
+        sim.submit_request(client_id, 0, msg.into_generic());
+        for _ in 0..200 {
+            sim.step();
+        }
+        for backup in 1..replica_count {
+            assert!(
+                metadata_op(&sim, usize::from(backup)) > committed_before,
+                "backup {backup} must hold the prepare before it is restarted"
+            );
+        }
+        assert_eq!(
+            metadata_commit(&sim, 0),
+            committed_before,
+            "the primary must not have committed while its acks were dropped"
+        );
+
+        // Restart every backup. Each recovers the unacked op from its own WAL 
and
+        // rejoins as a probing backup, which is the state the primary's
+        // retransmit cannot get an ack out of.
+        for backup in 1..replica_count {
+            sim.replica_crash(backup);
+            for _ in 0..50 {
+                sim.tick();
+            }
+            sim.replica_restart(backup);
+        }
+
+        // Healthy links from here: nothing but the protocol stands between the
+        // primary and its quorum.
+        for backup in 1..replica_count {
+            sim.network
+                .link_filter_mut(ProcessId::Replica(backup), 
ProcessId::Replica(0))
+                .insert(Command2::PrepareOk);
+        }
+
+        for _ in 0..10_000 {
+            sim.step();
+            if metadata_commit(&sim, 0) > committed_before {
+                return;
+            }
+        }
+        panic!(
+            "metadata commit stuck at {} after 10000 ticks with healthy links 
and \
+             every replica holding op {}: the restarted backups never re-acked 
the \
+             prepare they recovered from their own WALs",
+            metadata_commit(&sim, 0),
+            metadata_op(&sim, 1),
+        );
+    }
+
+    /// Committed metadata op on a replica's shard 0.
+    fn metadata_commit(sim: &Simulator, replica_idx: usize) -> u64 {
+        sim.replicas[replica_idx].shards[0]
+            .plane
+            .metadata()
+            .consensus
+            .as_ref()
+            .expect("shard 0 owns metadata consensus")
+            .commit_min()
+    }
+
+    /// Highest metadata op a replica has journaled.
+    fn metadata_op(sim: &Simulator, replica_idx: usize) -> u64 {
+        sim.replicas[replica_idx]
+            .metadata_journal
+            .last_op()
+            .unwrap_or(0)
+    }
+
     #[test]
     fn per_partition_consensus_independence() {
         use consensus::PIPELINE_PREPARE_QUEUE_MAX;
diff --git a/core/simulator/src/workload/auditor.rs 
b/core/simulator/src/workload/auditor.rs
index d9616ec4e..ca51a6016 100644
--- a/core/simulator/src/workload/auditor.rs
+++ b/core/simulator/src/workload/auditor.rs
@@ -184,6 +184,14 @@ impl ServerAuditor {
     }
 
     #[must_use]
+    /// The action of an outstanding request, if one is recorded for `key`.
+    /// Diagnostic only: names what a stalled run is waiting on, which the bare
+    /// `(client, request)` pair cannot.
+    #[must_use]
+    pub fn in_flight_action(&self, key: (u128, u64)) -> Option<Action> {
+        self.in_flight.get(&key).map(|entry| entry.action)
+    }
+
     pub fn in_flight_count(&self) -> usize {
         self.in_flight.len()
     }
diff --git a/core/simulator/src/workload/invariants.rs 
b/core/simulator/src/workload/invariants.rs
index bf00d9286..a20b99281 100644
--- a/core/simulator/src/workload/invariants.rs
+++ b/core/simulator/src/workload/invariants.rs
@@ -51,8 +51,20 @@ impl Invariants {
     /// Globally:
     /// - total in-flight requests stay within the per-client queue ceiling.
     ///
-    /// Crashed replicas are skipped: their last-seen marks are retained, which
-    /// stays correct because both quantities are monotonic across a restart.
+    /// Crashed replicas are skipped and their last-seen marks retained, which 
is
+    /// correct for `view` (the superblock carries it across a restart) but NOT
+    /// for partition `commit_offset`, because the simulator does not retain
+    /// partition data across a restart the way it retains the metadata WAL and
+    /// the superblocks. `IggyPartitions` is dropped and rebuilt empty, so a
+    /// restarted replica reports `commit_offset` 0 and this trips, reporting a
+    /// harness limitation as a consensus regression.
+    ///
+    /// So restart injection (`WorkloadOptions::restart_per_tick_ratio`) is
+    /// metadata-plane only today. Driving it with partition traffic needs the
+    /// partition journal held by the harness, per-namespace, exactly as
+    /// `SimReplica::metadata_journal` already is; a real server's segments 
are on
+    /// disk and do survive, so retaining them is the faithful model, and the
+    /// alternative of relaxing this check would instead model total data loss.
     ///
     /// # Panics
     /// On any regression or in-flight overflow. The workload seed is in the
diff --git a/core/simulator/src/workload/mod.rs 
b/core/simulator/src/workload/mod.rs
index 4ac4735da..322d66c22 100644
--- a/core/simulator/src/workload/mod.rs
+++ b/core/simulator/src/workload/mod.rs
@@ -188,14 +188,23 @@ impl Workload {
         due
     }
 
-    /// Outstanding requests as `(client, request, target, attempts)`, in key
-    /// order. Diagnostic only: names what a run was still waiting on when it
-    /// failed to drain.
+    /// Outstanding requests as `(client, request, action, target, attempts)`, 
in
+    /// key order. Diagnostic only: names what a run was still waiting on when 
it
+    /// failed to drain, including which op, since some ops cannot be answered
+    /// twice and a resend of those is a permanent stall rather than a slow 
one.
     #[must_use]
-    pub(crate) fn outstanding_summary(&self) -> Vec<(u128, u64, u8, u32)> {
+    pub(crate) fn outstanding_summary(&self) -> Vec<(u128, u64, 
Option<Action>, u8, u32)> {
         self.outstanding
             .iter()
-            .map(|(&(client, request), entry)| (client, request, entry.target, 
entry.attempts))
+            .map(|(&key, entry)| {
+                (
+                    key.0,
+                    key.1,
+                    self.auditor.in_flight_action(key),
+                    entry.target,
+                    entry.attempts,
+                )
+            })
             .collect()
     }
 
@@ -463,23 +472,44 @@ const FAULT_SEED_SALT: u64 = 0x5A1A_F0E5_FACE_0001;
 ///
 /// The invariants are asserted after every tick, so a consensus or
 /// workload regression panics at the tick it occurs (the seed in the message
-/// replays it). When `crash_per_tick_ratio > 0` the driver also injects
-/// crash-only faults via [`maybe_inject_crash`].
+/// replays it). Crash and restart injection runs through [`FaultInjector`],
+/// idle unless one of the two probabilities is set.
+///
+/// Discards the injector; use [`run_with_faults`] to read the crash and 
restart
+/// counts back.
 pub fn run(
     sim: &mut Simulator,
     workload: &mut Workload,
     clients: &[SimClient],
     tick_budget: u64,
     replies_target: u64,
+) -> u64 {
+    let mut injector = FaultInjector::new(workload.options.seed, 
sim.replica_count);
+    run_with_faults(
+        sim,
+        workload,
+        clients,
+        tick_budget,
+        replies_target,
+        &mut injector,
+    )
+}
+
+/// [`run`] against a caller-owned [`FaultInjector`], so a test can assert what
+/// was actually injected instead of trusting the probabilities to have fired.
+pub fn run_with_faults(
+    sim: &mut Simulator,
+    workload: &mut Workload,
+    clients: &[SimClient],
+    tick_budget: u64,
+    replies_target: u64,
+    injector: &mut FaultInjector,
 ) -> u64 {
     let mut invariants = Invariants::new();
-    let mut fault_prng = Xoshiro256Plus::seed_from_u64(workload.options.seed ^ 
FAULT_SEED_SALT);
     let mut replies_seen = 0u64;
     for _ in 0..tick_budget {
         workload.tick();
-        if workload.options.crash_per_tick_ratio > 0.0 {
-            maybe_inject_crash(sim, workload, &mut fault_prng);
-        }
+        injector.step(sim, workload);
         // Resend before sampling: a timed-out request still holds the client's
         // slot, so `build_request` would decline it anyway.
         resubmit_due(sim, workload);
@@ -501,48 +531,141 @@ pub fn run(
     replies_seen
 }
 
-/// With probability `crash_per_tick_ratio`, crash one live non-primary 
replica,
-/// provided doing so leaves at least `min_survivors` live. Crash-only: a
-/// crashed replica is never restarted (that needs consensus durability).
-///
-/// "Non-primary" is partition-plane only: the exclusion set comes from
-/// `Simulator::primary_index`, which reads `partitions()`. The metadata-plane
-/// primary is not consulted; it is spared only by co-location, since every 
group
-/// starts at view 0 with `primary = view % replica_count` (so replica 0 leads
-/// both planes) and `min_survivors` keeps a commit quorum, so no view change
-/// moves it. Were the two planes' primaries to diverge, the metadata primary
-/// could be crashed.
+/// Crash and restart injection with stability windows, in the shape of
+/// TigerBeetle's VOPR: a crash must last a while before it may be repaired, 
and
+/// a repaired replica must run a while before it may fail again.
 ///
-/// Primaries are spared at all because the driver has no 
request-timeout/resend
-/// path: a request lost to a crashed primary would wedge the client's only
-/// in-flight slot. Forcing primary crashes (and the view change they trigger)
-/// while keeping traffic flowing is future work gated on that resend path.
-fn maybe_inject_crash(sim: &mut Simulator, workload: &Workload, prng: &mut 
Xoshiro256Plus) {
-    let live: Vec<u8> = (0..sim.replica_count)
-        .filter(|replica_idx| !sim.is_crashed(*replica_idx))
-        .collect();
-    if live.len() <= usize::from(workload.options.min_survivors) {
-        return;
+/// Owns the fault PRNG so crash scheduling stays reproducible from the seed 
yet
+/// independent of the traffic draw order. Draws nothing while both 
probabilities
+/// are zero, so a fault-free run replays bit-identically.
+pub struct FaultInjector {
+    prng: Xoshiro256Plus,
+    /// Tick of each replica's last crash or restart, indexed by replica id.
+    /// Compared against the stability windows to decide eligibility.
+    last_transition: Vec<u64>,
+    now: u64,
+    crashes: u64,
+    restarts: u64,
+}
+
+impl FaultInjector {
+    #[must_use]
+    pub fn new(seed: u64, replica_count: u8) -> Self {
+        Self {
+            prng: Xoshiro256Plus::seed_from_u64(seed ^ FAULT_SEED_SALT),
+            last_transition: vec![0; usize::from(replica_count)],
+            now: 0,
+            crashes: 0,
+            restarts: 0,
+        }
+    }
+
+    #[must_use]
+    pub const fn crashes(&self) -> u64 {
+        self.crashes
     }
-    let roll: f32 = prng.random();
-    if roll >= workload.options.crash_per_tick_ratio {
-        return;
+
+    #[must_use]
+    pub const fn restarts(&self) -> u64 {
+        self.restarts
+    }
+
+    /// Advance one tick and maybe crash or restart one replica.
+    ///
+    /// Restart is considered before crash so a single tick never both revives
+    /// and kills, which would make the stability windows meaningless.
+    pub fn step(&mut self, sim: &mut Simulator, workload: &Workload) {
+        self.now += 1;
+        self.maybe_restart(sim, workload);
+        self.maybe_crash(sim, workload);
+    }
+
+    /// With probability `restart_per_tick_ratio`, restart one replica that has
+    /// been down at least `crash_stability_ticks`.
+    ///
+    /// This is what exercises rejoin: the replica comes back with its durable
+    /// superblock and metadata WAL but no volatile consensus state, asks the
+    /// current view's primary for a `StartView`, and repairs the log it 
missed.
+    fn maybe_restart(&mut self, sim: &mut Simulator, workload: &Workload) {
+        if workload.options.restart_per_tick_ratio <= 0.0 {
+            return;
+        }
+        let eligible: Vec<u8> = (0..sim.replica_count)
+            .filter(|replica_idx| sim.is_crashed(*replica_idx))
+            .filter(|replica_idx| {
+                self.stable_for(*replica_idx) >= 
workload.options.crash_stability_ticks
+            })
+            .collect();
+        if eligible.is_empty() {
+            return;
+        }
+        let roll: f32 = self.prng.random();
+        if roll >= workload.options.restart_per_tick_ratio {
+            return;
+        }
+        let revived = eligible[self.prng.random_range(0..eligible.len())];
+        sim.replica_restart(revived);
+        self.last_transition[usize::from(revived)] = self.now;
+        self.restarts += 1;
     }
-    let primaries: HashSet<u8> = workload
-        .options
-        .namespaces
-        .iter()
-        .filter_map(|ns| sim.primary_index(*ns))
-        .collect();
-    let eligible: Vec<u8> = live
-        .into_iter()
-        .filter(|replica_idx| !primaries.contains(replica_idx))
-        .collect();
-    if eligible.is_empty() {
-        return;
+
+    /// With probability `crash_per_tick_ratio`, crash one live replica that 
has
+    /// been up at least `restart_stability_ticks`, provided doing so leaves at
+    /// least `min_survivors` live.
+    ///
+    /// Primaries are excluded unless `spare_primary` is off. The exclusion set
+    /// comes from `Simulator::primary_index`, which reads `partitions()`, so 
it
+    /// names partition-plane primaries; the metadata primary is spared only by
+    /// co-location, since every group starts at view 0 with
+    /// `primary = view % replica_count`. Once views diverge across planes the
+    /// metadata primary can be crashed even with this on.
+    fn maybe_crash(&mut self, sim: &mut Simulator, workload: &Workload) {
+        if workload.options.crash_per_tick_ratio <= 0.0 {
+            return;
+        }
+        let live: Vec<u8> = (0..sim.replica_count)
+            .filter(|replica_idx| !sim.is_crashed(*replica_idx))
+            .collect();
+        if live.len() <= usize::from(workload.options.min_survivors) {
+            return;
+        }
+        let roll: f32 = self.prng.random();
+        if roll >= workload.options.crash_per_tick_ratio {
+            return;
+        }
+        let primaries: HashSet<u8> = if workload.options.spare_primary {
+            workload
+                .options
+                .namespaces
+                .iter()
+                .filter_map(|ns| sim.primary_index(*ns))
+                .collect()
+        } else {
+            HashSet::new()
+        };
+        let eligible: Vec<u8> = live
+            .into_iter()
+            .filter(|replica_idx| !primaries.contains(replica_idx))
+            .filter(|replica_idx| {
+                self.stable_for(*replica_idx) >= 
workload.options.restart_stability_ticks
+            })
+            .collect();
+        if eligible.is_empty() {
+            return;
+        }
+        let victim = eligible[self.prng.random_range(0..eligible.len())];
+        sim.replica_crash(victim);
+        self.last_transition[usize::from(victim)] = self.now;
+        self.crashes += 1;
+    }
+
+    /// Ticks since this replica last changed state. A replica that never
+    /// transitioned counts from tick 0, so the first crash still has to wait 
out
+    /// `restart_stability_ticks`.
+    fn stable_for(&self, replica_idx: u8) -> u64 {
+        self.now
+            .saturating_sub(self.last_transition[usize::from(replica_idx)])
     }
-    let victim = eligible[prng.random_range(0..eligible.len())];
-    sim.replica_crash(victim);
 }
 
 /// Submit every request whose reply is overdue (see 
[`Workload::due_resends`]).
diff --git a/core/simulator/src/workload/options.rs 
b/core/simulator/src/workload/options.rs
index d707b2961..f5fd38acb 100644
--- a/core/simulator/src/workload/options.rs
+++ b/core/simulator/src/workload/options.rs
@@ -27,6 +27,16 @@ use strum::EnumCount;
 /// a crashed primary is retried long before the run's budget runs out.
 pub const DEFAULT_REQUEST_TIMEOUT_TICKS: u64 = 200;
 
+/// Default [`WorkloadOptions::crash_stability_ticks`]. Long enough that the
+/// surviving primary commits past the crashed replica's log, so its rejoin has
+/// something to repair.
+pub const DEFAULT_CRASH_STABILITY_TICKS: u64 = 300;
+
+/// Default [`WorkloadOptions::restart_stability_ticks`]. Long enough for a
+/// rejoined replica to finish catching up before it becomes a crash candidate
+/// again, so a run does not consist entirely of half-repaired replicas.
+pub const DEFAULT_RESTART_STABILITY_TICKS: u64 = 500;
+
 /// Per-action sampling weights as percentages. Unlisted variants default
 /// to 0 (never picked). Listed weights must sum to 100.
 #[derive(Debug, Clone, Copy)]
@@ -172,10 +182,29 @@ pub struct WorkloadOptions {
     pub consumer_pool_size: u32,
     /// Upper bound on offset carried by `StoreConsumerOffset2`.
     pub max_offset: u64,
-    /// Probability per tick that the driver crashes one live non-primary
-    /// replica (crash-only, no restart). `0.0` disables injection: the fault
-    /// PRNG draws nothing, so traffic stays bit-identical.
+    /// Probability per tick that the driver crashes one eligible replica.
+    /// `0.0` disables injection entirely: the fault PRNG draws nothing, so
+    /// traffic stays bit-identical.
     pub crash_per_tick_ratio: f32,
+    /// Probability per tick that the driver restarts one crashed replica.
+    /// Meaningless without `crash_per_tick_ratio`, since nothing is ever down.
+    pub restart_per_tick_ratio: f32,
+    /// Ticks a replica must stay down before it may be restarted. Keeps a 
crash
+    /// long enough to actually matter: a replica restarted the tick after it
+    /// crashed never falls behind, so nothing needs repairing.
+    pub crash_stability_ticks: u64,
+    /// Ticks a replica must stay up before it may be crashed again. Stops a
+    /// single unlucky replica from being crash-looped while its peers never
+    /// fail.
+    pub restart_stability_ticks: u64,
+    /// Leave the primary of every tracked namespace out of the crash pool.
+    ///
+    /// Defaults to `true`, which is what the driver did unconditionally before
+    /// clients could resend: a request lost to a crashed primary was never
+    /// retried, so it stranded the client's only in-flight slot. With 
resending
+    /// in place, setting this to `false` is what puts a view change under live
+    /// traffic, the scenario the harness exists for.
+    pub spare_primary: bool,
     /// Floor on live replicas the driver will not crash below, preserving a
     /// commit quorum. Defaults to `replica_count / 2 + 1`.
     pub min_survivors: u8,
@@ -206,6 +235,10 @@ impl WorkloadOptions {
             consumer_pool_size: 4,
             max_offset: 1_000_000,
             crash_per_tick_ratio: 0.0,
+            restart_per_tick_ratio: 0.0,
+            crash_stability_ticks: DEFAULT_CRASH_STABILITY_TICKS,
+            restart_stability_ticks: DEFAULT_RESTART_STABILITY_TICKS,
+            spare_primary: true,
             min_survivors: replica_count / 2 + 1,
             request_timeout_ticks: DEFAULT_REQUEST_TIMEOUT_TICKS,
         }
diff --git a/core/simulator/src/workload/oracle.rs 
b/core/simulator/src/workload/oracle.rs
index 4bdfc50f5..567172421 100644
--- a/core/simulator/src/workload/oracle.rs
+++ b/core/simulator/src/workload/oracle.rs
@@ -138,10 +138,10 @@ pub fn quiesce_failure_report(sim: &Simulator, workload: 
&Workload) -> String {
         workload.total_in_flight(),
         workload.options.seed,
     );
-    for (client, request, target, attempts) in workload.outstanding_summary() {
+    for (client, request, action, target, attempts) in 
workload.outstanding_summary() {
         let _ = writeln!(
             report,
-            "  outstanding client={client} request={request} \
+            "  outstanding client={client} request={request} action={action:?} 
\
              last_target=replica {target} attempts={attempts}",
         );
     }
@@ -152,6 +152,26 @@ pub fn quiesce_failure_report(sim: &Simulator, workload: 
&Workload) -> String {
             continue;
         }
         let _ = write!(report, "  replica {replica_idx}: live");
+        // Metadata plane first: a rejoining replica is quorum-invisible until 
it
+        // completes its view probe, so its status is the difference between 
"the
+        // cluster is slow" and "the cluster has no quorum despite enough live
+        // replicas".
+        if let Some(consensus) = 
sim.replicas[usize::from(replica_idx)].shards[0]
+            .plane
+            .metadata()
+            .consensus
+            .as_ref()
+        {
+            let _ = write!(
+                report,
+                " | metadata status={:?} view={} log_view={} commit={} 
primary={}",
+                consensus.status(),
+                consensus.view(),
+                consensus.log_view(),
+                consensus.commit_min(),
+                consensus.is_primary(),
+            );
+        }
         for &ns in &workload.options.namespaces {
             let view = sim.consensus_view(usize::from(replica_idx), ns);
             let commit = sim
@@ -164,6 +184,26 @@ pub fn quiesce_failure_report(sim: &Simulator, workload: 
&Workload) -> String {
             );
         }
         report.push('\n');
+        // Per-client table state. `check_request` admits a metadata request 
only
+        // when it is exactly `watermark + 1`, so the watermark says whether an
+        // outstanding request is still expected, already answered (and thus 
owed
+        // a cached-reply replay), or ahead of what this replica will accept.
+        let table = sim.replicas[usize::from(replica_idx)].shards[0]
+            .plane
+            .metadata()
+            .client_table
+            .borrow();
+        for client_id in table.client_ids() {
+            let _ = writeln!(
+                report,
+                "    client {client_id}: watermark={:?} epoch={:?} 
cached_reply_request={:?}",
+                table.get_watermark(client_id),
+                table.get_epoch(client_id),
+                table
+                    .get_reply(client_id)
+                    .map(|reply| reply.header().request),
+            );
+        }
     }
     report
 }

Reply via email to