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 c0d9a2104ac18c8eb02b7252cad5561fb5cf3ecb
Author: Krishna Vishal <[email protected]>
AuthorDate: Fri Aug 14 14:29:42 2026 +0530

    feat(simulator): select the fuzzer op mix and surface server diagnostics
    
    The fuzzer hardcoded `SendMessages: 100` and its module docs claimed
    metadata and mixed workloads were gated on the metadata request-gap.
    That gap was a `SimClient` numbering bug, fixed when the client began
    numbering requests per plane, so the gate is gone and the doc was stale.
    
    `--plane` now picks between partition-only, metadata-only, the mixed
    default and uniform-across-every-action, as named `ActionWeights`
    presets rather than a table spelled out at each call site. `uniform`
    spreads the remainder rather than requiring the action count to divide
    100, so appending an action never breaks it. `--ack-quorum-ratio`
    exposes the knob that decides whether an offset store takes the
    replicated path or the leader-local one.
    
    The binary also installs a tracing subscriber, `RUST_LOG`-gated and off
    by default. Server-side diagnostics are the only record of a request
    dropped after logging, which is precisely the shape that wedges a
    client's in-flight slot, and without a subscriber such a run reads as an
    unexplained stall.
---
 Cargo.lock                              |  1 +
 core/simulator/Cargo.toml               |  1 +
 core/simulator/src/bin/workload-fuzz.rs | 66 ++++++++++++++++++++++++------
 core/simulator/src/workload/options.rs  | 72 +++++++++++++++++++++++++++++++++
 4 files changed, 127 insertions(+), 13 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock
index 9a69b24d1..14607f416 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -12191,6 +12191,7 @@ dependencies = [
  "shard",
  "strum 0.28.0",
  "tracing",
+ "tracing-subscriber",
 ]
 
 [[package]]
diff --git a/core/simulator/Cargo.toml b/core/simulator/Cargo.toml
index 37df9a52e..adc3582a0 100644
--- a/core/simulator/Cargo.toml
+++ b/core/simulator/Cargo.toml
@@ -48,6 +48,7 @@ server_common = { path = "../server_common", features = 
["simulator"] }
 shard = { path = "../shard", features = ["simulator"] }
 strum = { workspace = true }
 tracing = { workspace = true }
+tracing-subscriber = { workspace = true }
 
 [lints.clippy]
 enum_glob_use = "deny"
diff --git a/core/simulator/src/bin/workload-fuzz.rs 
b/core/simulator/src/bin/workload-fuzz.rs
index fc495f8af..e3eebfb50 100644
--- a/core/simulator/src/bin/workload-fuzz.rs
+++ b/core/simulator/src/bin/workload-fuzz.rs
@@ -24,18 +24,15 @@
 //!
 //! ```text
 //! workload-fuzz [--seed N] [--ticks N] [--clients N] [--replicas N]
+//!               [--plane partition|metadata|mixed|uniform]
 //!               [--crash-prob F] [--no-quiesce]
 //! ```
 //!
-//! The default workload is partition-plane (`SendMessages`): it drains and
-//! converges. Metadata and mixed-plane workloads are gated on the metadata
-//! request-gap: a client's replicated-metadata request ids must arrive
-//! contiguously (`committed + 1`), so a dropped or reordered metadata request
-//! opens a permanent `RequestGap` that wedges that client's metadata plane.
-//! Broader op coverage lands once the workload generator models that
-//! constraint.
-
-use clap::Parser;
+//! `--plane` selects the op mix (see [`ActionWeights`]). Partition-plane runs
+//! drain and converge most readily; `uniform` is the widest per-tick op
+//! coverage.
+
+use clap::{Parser, ValueEnum};
 use iggy_common::IggyByteSize;
 use server_common::sharding::IggyNamespace;
 use server_common::{MemoryPool, MemoryPoolConfigOther};
@@ -59,12 +56,43 @@ struct Args {
     clients: u8,
     #[arg(long, default_value_t = 3, value_parser = 
clap::value_parser!(u8).range(1..))]
     replicas: u8,
+    /// Op mix to draw from.
+    #[arg(long, value_enum, default_value_t = Plane::Partition)]
+    plane: Plane,
+    /// Probability a consumer-offset store asks for `Quorum` rather than
+    /// `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,
     #[arg(long, default_value_t = 0.0, value_parser = parse_unit_interval)]
     crash_prob: f32,
     #[arg(long)]
     no_quiesce: bool,
 }
 
+/// Which plane the sampled ops target. Maps onto an [`ActionWeights`] preset.
+#[derive(Clone, Copy, Debug, PartialEq, Eq, ValueEnum)]
+enum Plane {
+    /// Writes and consumer offsets only.
+    Partition,
+    /// Replicated metadata mutations only.
+    Metadata,
+    /// Stream creates over a write-heavy base.
+    Mixed,
+    /// Every action equally likely.
+    Uniform,
+}
+
+impl Plane {
+    fn weights(self) -> ActionWeights {
+        match self {
+            Self::Partition => ActionWeights::partition_only(),
+            Self::Metadata => ActionWeights::metadata_only(),
+            Self::Mixed => ActionWeights::default(),
+            Self::Uniform => ActionWeights::uniform(),
+        }
+    }
+}
+
 /// Clap value parser: accept a probability in `[0.0, 1.0]`.
 fn parse_unit_interval(raw: &str) -> Result<f32, String> {
     let value: f32 = raw
@@ -80,12 +108,23 @@ fn parse_unit_interval(raw: &str) -> Result<f32, String> {
 fn main() {
     let args = Args::parse();
 
+    // Server-side diagnostics (`emit_partition_diag` and friends) are the only
+    // record of a request the server dropped after logging, which is exactly 
the
+    // shape that wedges a client's in-flight slot. Without a subscriber they 
go
+    // nowhere and the run looks like an unexplained stall, so install one and 
let
+    // `RUST_LOG` select. Off by default: a WARN per dropped frame drowns the 
run.
+    tracing_subscriber::fmt()
+        .with_env_filter(tracing_subscriber::EnvFilter::from_default_env())
+        .with_writer(std::io::stderr)
+        .init();
+
     // A provided seed reproduces a prior run exactly; otherwise draw one and
     // log it. Both the network and workload PRNGs derive from it.
     let seed = args.seed.unwrap_or_else(rand::random);
     let ticks = args.ticks;
     let clients = args.clients;
     let replicas = args.replicas;
+    let plane = args.plane;
     let crash_prob = args.crash_prob;
     let quiesce = !args.no_quiesce;
 
@@ -97,7 +136,7 @@ fn main() {
 
     println!(
         "workload-fuzz: seed={seed} ticks={ticks} clients={clients} 
replicas={replicas} \
-         crash_prob={crash_prob} quiesce={quiesce}"
+         plane={plane:?} crash_prob={crash_prob} quiesce={quiesce}"
     );
 
     // poll_messages / reply paths panic without an initialized pool; disabled
@@ -131,7 +170,8 @@ fn main() {
     let mut options = WorkloadOptions::new(seed, replicas, vec![ns]);
     options.client_count = clients;
     options.crash_per_tick_ratio = crash_prob;
-    options.weights = ActionWeights::new(&[(Action::SendMessages, 100)]);
+    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);
@@ -146,8 +186,8 @@ fn main() {
             println!("quiesced and converged (leader-relative + entity 
oracle)");
         } else {
             println!(
-                "WARN: did not quiesce within budget — expected when crashing 
to bare quorum \
-                 or under the metadata request-gap limitation; per-tick 
invariants still held"
+                "WARN: did not quiesce within budget — expected when crashing 
to bare quorum; \
+                 per-tick invariants still held"
             );
         }
     }
diff --git a/core/simulator/src/workload/options.rs 
b/core/simulator/src/workload/options.rs
index a6c0bbece..a4a29447a 100644
--- a/core/simulator/src/workload/options.rs
+++ b/core/simulator/src/workload/options.rs
@@ -47,6 +47,78 @@ impl ActionWeights {
         Self { weights }
     }
 
+    /// Partition plane only: writes plus consumer-offset traffic, no metadata
+    /// mutation. The regime that drains and converges most readily, so it is
+    /// what a run reaches for when the question is about replication rather
+    /// than about the state machine.
+    #[must_use]
+    pub fn partition_only() -> Self {
+        Self::new(&[
+            (Action::SendMessages, 60),
+            (Action::StoreConsumerOffset2, 25),
+            (Action::DeleteConsumerOffset2, 5),
+            (Action::StoreConsumerOffset, 7),
+            (Action::DeleteConsumerOffset, 3),
+        ])
+    }
+
+    /// Metadata plane only: every replicated metadata mutation, weighted so
+    /// creates outrun deletes and the shadow keeps a live population to sample
+    /// duplicate- and missing-target outcomes against. `DeleteSegments` is
+    /// excluded: it resolves against partition state the partition-plane
+    /// presets build, so it belongs to a mixed run.
+    #[must_use]
+    pub fn metadata_only() -> Self {
+        Self::new(&[
+            (Action::CreateStream, 12),
+            (Action::UpdateStream, 6),
+            (Action::DeleteStream, 6),
+            (Action::PurgeStream, 4),
+            (Action::CreateTopic, 12),
+            (Action::UpdateTopic, 6),
+            (Action::DeleteTopic, 6),
+            (Action::PurgeTopic, 4),
+            (Action::CreatePartitions, 6),
+            (Action::DeletePartitions, 4),
+            (Action::CreateConsumerGroup, 6),
+            (Action::DeleteConsumerGroup, 4),
+            (Action::CreateUser, 6),
+            (Action::UpdateUser, 3),
+            (Action::DeleteUser, 3),
+            (Action::ChangePassword, 3),
+            (Action::UpdatePermissions, 3),
+            (Action::CreatePersonalAccessToken, 3),
+            (Action::DeletePersonalAccessToken, 3),
+        ])
+    }
+
+    /// Every action equally likely. Widest op coverage per tick, at the cost 
of
+    /// a shallow population per entity kind.
+    ///
+    /// `Action::COUNT` does not divide 100, so the first `100 % COUNT` actions
+    /// carry one extra point. Spread that way rather than asserting the count
+    /// divides evenly, so appending an `Action` never breaks this preset.
+    #[must_use]
+    pub fn uniform() -> Self {
+        use strum::IntoEnumIterator;
+
+        let count = u32::try_from(Action::COUNT).expect("Action::COUNT fits 
u32");
+        let base = 100 / count;
+        let remainder = 100 % count;
+        let entries: Vec<(Action, u8)> = Action::iter()
+            .enumerate()
+            .map(|(idx, action)| {
+                let extra = u32::try_from(idx).expect("action index fits u32") 
< remainder;
+                let weight = base + u32::from(extra);
+                (
+                    action,
+                    u8::try_from(weight).expect("per-action weight is at most 
100"),
+                )
+            })
+            .collect();
+        Self::new(&entries)
+    }
+
     #[must_use]
     pub const fn weight(&self, action: Action) -> u8 {
         self.weights[action as usize]

Reply via email to