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]
