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 9b0dd0586dfbddc6c98c56fe24715ec1f65a188b Author: Krishna Vishal <[email protected]> AuthorDate: Sat Aug 15 12:18:58 2026 +0530 style: satisfy the pedantic lints the new simulator code introduced Ten diagnostics across the batch's own code: seven denied under the crate's `pedantic` setting and three warnings that CI's `-D warnings` would have turned into failures too. Missing `#[must_use]` and `# Panics`, doc paragraphs and backticks, an argument taken by value and never consumed, and a return type wrapped in an `Option` that was never `None`. `await_setup_reply` returning `Option` was the one worth more than a lint fix: it panics when the handshake never completes, so the `None` its callers unwrapped through could not occur, and both of them dealt with it by defaulting a session id to zero. Returning the reply directly removes a branch that only ever obscured the real failure. --- core/binary_protocol/src/primitives/consumer.rs | 9 +-- core/metadata/src/impls/metadata.rs | 6 +- core/simulator/src/bin/workload-fuzz.rs | 88 ++++++++++++++----------- core/simulator/src/lib.rs | 57 ++++++++-------- core/simulator/src/workload/auditor.rs | 1 + core/simulator/src/workload/mod.rs | 9 +-- core/simulator/src/workload/options.rs | 12 +++- core/simulator/src/workload/state_checker.rs | 7 +- 8 files changed, 106 insertions(+), 83 deletions(-) diff --git a/core/binary_protocol/src/primitives/consumer.rs b/core/binary_protocol/src/primitives/consumer.rs index 1e7739b37..cc44f4479 100644 --- a/core/binary_protocol/src/primitives/consumer.rs +++ b/core/binary_protocol/src/primitives/consumer.rs @@ -20,10 +20,11 @@ use crate::WireIdentifier; use crate::codec::{WireDecode, WireEncode, read_u8}; use bytes::{BufMut, BytesMut}; -/// Wire discriminant for a single consumer (vs a `ConsumerGroup`). Public for -/// the same reason as its sibling below, plus one more: `decode` accepts only -/// these two values, so anything synthesizing a consumer needs to name them -/// rather than guess a small integer. +/// Wire discriminant for a single consumer (vs a `ConsumerGroup`). +/// +/// Public for the same reason as its sibling below, plus one more: `decode` +/// accepts only these two values, so anything synthesizing a consumer needs to +/// name them rather than guess a small integer. pub const KIND_CONSUMER: u8 = 1; /// Wire discriminant for a consumer-group consumer (vs a single `Consumer`). /// Public so the server dispatch can match on it by name instead of a raw `2`. diff --git a/core/metadata/src/impls/metadata.rs b/core/metadata/src/impls/metadata.rs index b3b682d02..700b8d87e 100644 --- a/core/metadata/src/impls/metadata.rs +++ b/core/metadata/src/impls/metadata.rs @@ -4787,9 +4787,9 @@ mod tests { /// metadata workload under crash/restart injection until this was gated on the /// journal instead. /// - /// TigerBeetle sidesteps the question by keeping a single frontier: `self.op` - /// IS the log head, and `on_prepare` routes anything at or below it to - /// `on_repair`, which re-acks a prepare already held. + /// `TigerBeetle` sidesteps the question by keeping a single frontier: + /// `self.op` IS the log head, and `on_prepare` routes anything at or below it + /// to `on_repair`, which re-acks a prepare already held. #[compio::test] async fn backup_admits_the_prepare_its_journal_needs_despite_a_leading_sequencer() { const CLIENT: u128 = 1; diff --git a/core/simulator/src/bin/workload-fuzz.rs b/core/simulator/src/bin/workload-fuzz.rs index 940120387..dd0720cf4 100644 --- a/core/simulator/src/bin/workload-fuzz.rs +++ b/core/simulator/src/bin/workload-fuzz.rs @@ -136,7 +136,7 @@ struct Args { clog_duration_mean: Option<u64>, } -/// Named network fault profile, in the spirit of TigerBeetle's VOPR modes: one +/// Named network fault profile, in the spirit of `TigerBeetle`'s VOPR modes: one /// flag for "how hostile is the network", rather than eleven. /// /// Progress falls off steeply with severity, because every lost frame costs a @@ -388,40 +388,7 @@ fn main() { bucket_capacity: 1, }); - let client_ids: Vec<u128> = (1..=u128::from(clients)).collect(); - let mut sim = if args.shell { - Simulator::with_shards_shell( - usize::from(replicas), - 1, - client_ids.iter().copied(), - network_opts, - ) - } else { - Simulator::new( - usize::from(replicas), - client_ids.iter().copied(), - network_opts, - ) - }; - let sim_clients: Vec<SimClient> = client_ids.iter().map(|&id| SimClient::new(id)).collect(); - - let ns = IggyNamespace::new(1, 1, 0); - sim.init_partition(ns); - if args.shell { - // The shell resolves a partition request's namespace against committed - // metadata, so the stream and topic behind it have to exist as well as the - // partition group. - sim.seed_stream_topic_partition(ns); - } - for client in &sim_clients { - if args.shell { - // Log in rather than bare-register: the dispatch path admits a request - // only from a bound session, and the login is what mints one. - sim.shell_login(client); - } else { - sim.register_client_with_primary(client); - } - } + let (mut sim, sim_clients, ns) = build_cluster(&args, replicas, clients, network_opts); let mut options = WorkloadOptions::new(seed, replicas, vec![ns]); options.client_count = clients; @@ -484,6 +451,51 @@ fn main() { println!("workload-fuzz: OK (seed={seed})"); } +/// Stand up the cluster, seed its namespace, and get every client a session. +/// +/// Returns the simulator, its clients, and the namespace the workload drives. +/// The shell path differs in two ways that have to agree: a partition request's +/// namespace is resolved against committed metadata, so the stream and topic +/// behind it must exist and not just the partition group; and dispatch admits a +/// request only from a bound session, which only a login mints. +fn build_cluster( + args: &Args, + replicas: u8, + clients: u8, + network_opts: PacketSimulatorOptions, +) -> (Simulator, Vec<SimClient>, IggyNamespace) { + let client_ids: Vec<u128> = (1..=u128::from(clients)).collect(); + let mut sim = if args.shell { + Simulator::with_shards_shell( + usize::from(replicas), + 1, + client_ids.iter().copied(), + network_opts, + ) + } else { + Simulator::new( + usize::from(replicas), + client_ids.iter().copied(), + network_opts, + ) + }; + let sim_clients: Vec<SimClient> = client_ids.iter().map(|&id| SimClient::new(id)).collect(); + + let ns = IggyNamespace::new(1, 1, 0); + sim.init_partition(ns); + if args.shell { + sim.seed_stream_topic_partition(ns); + } + for client in &sim_clients { + if args.shell { + sim.shell_login(client); + } else { + sim.register_client_with_primary(client); + } + } + (sim, sim_clients, ns) +} + /// Which protocol commands the run actually delivered, and which it never /// reached. /// @@ -524,9 +536,9 @@ fn print_coverage(workload: &Workload) { ); for action in Action::iter() { let commits = stats.commits(action); - let (denied, status) = stats.denials_per_action[action as usize]; - if commits > 0 || denied > 0 { - println!(" {action:?}: {commits} commits, {denied} denied (last status {status})"); + let (refused, code) = stats.denials_per_action[action as usize]; + if commits > 0 || refused > 0 { + println!(" {action:?}: {commits} commits, {refused} denied (last status {code})"); } } } diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs index a266dc65d..4a0b3a08d 100644 --- a/core/simulator/src/lib.rs +++ b/core/simulator/src/lib.rs @@ -281,7 +281,7 @@ impl Simulator { replica_count: usize, clients: impl Iterator<Item = u128>, network_options: PacketSimulatorOptions, - data_dir_root: std::path::PathBuf, + data_dir_root: &std::path::Path, ) -> Self { Self::build_inner( replica_count, @@ -326,7 +326,7 @@ impl Simulator { clients: impl Iterator<Item = u128>, network_options: PacketSimulatorOptions, shell: bool, - data_dir_root: Option<std::path::PathBuf>, + data_dir_root: Option<&std::path::Path>, ) -> Self { assert!( shards_per_replica >= 1, @@ -583,7 +583,8 @@ impl Simulator { // (`build_reply_with_body` maps the session field to `op`). let session = self .await_setup_reply(client.client_id(), target, &msg, "shell_login") - .map_or(0, |reply| reply.header().op); + .header() + .op; assert!(session > 0, "shell_login: login reply carried no session"); client.bind_session(session); } @@ -591,22 +592,26 @@ impl Simulator { /// Submit `message` to `target` and step until a client reply arrives, /// resubmitting every [`SETUP_RETRY_STEPS`] steps. /// - /// Returns the first reply, or `None` once [`SETUP_TOTAL_STEPS`] is spent. /// The same message is resubmitted verbatim, so the request id is stable and /// the metadata client table treats a retry as a duplicate. + /// + /// # Panics + /// If no reply arrives within [`SETUP_TOTAL_STEPS`]. A setup handshake that + /// never completes leaves the fixture unusable, so there is no useful + /// `None` for a caller to handle. fn await_setup_reply( &mut self, client_id: u128, target: u8, message: &Message<GenericHeader>, label: &str, - ) -> Option<Message<ReplyHeader>> { + ) -> Message<ReplyHeader> { for step in 0..SETUP_TOTAL_STEPS { if step % SETUP_RETRY_STEPS == 0 { self.submit_request(client_id, target, message.deep_copy()); } if let Some(reply) = self.step().into_iter().next() { - return Some(reply); + return reply; } } panic!( @@ -852,26 +857,22 @@ impl Simulator { #[allow(clippy::cast_possible_truncation)] pub fn register_client_with_primary(&mut self, client: &SimClient) { let msg = client.register().into_generic(); - let session = self - .await_setup_reply(client.client_id(), 0, &msg, "register_client_with_primary") - .map_or(0, |reply| { - let header = reply.header(); - debug_assert_eq!( - header.operation, - iggy_binary_protocol::Operation::Register, - "register_client_with_primary: first reply was not Register" - ); - assert_eq!( - header.client, - client.client_id(), - "register_client_with_primary: reply client_id mismatch \ - (expected {}, got {})", - client.client_id(), - header.client, - ); - header.commit - }); - client.bind_session(session); + let reply = + self.await_setup_reply(client.client_id(), 0, &msg, "register_client_with_primary"); + let header = reply.header(); + debug_assert_eq!( + header.operation, + iggy_binary_protocol::Operation::Register, + "register_client_with_primary: first reply was not Register" + ); + assert_eq!( + header.client, + client.client_id(), + "register_client_with_primary: reply client_id mismatch (expected {}, got {})", + client.client_id(), + header.client, + ); + client.bind_session(header.commit); // Partition has no `client_table`: at-least-once, no per-client // dedup. Consumers dedup via message id / content / producer-id+seq. @@ -3441,7 +3442,7 @@ mod tests { usize::from(replica_count), std::iter::once(client_id), network_opts, - root.path().to_path_buf(), + root.path(), ); // Small enough that the ops below cross the margin; the coordinator forces // a checkpoint once free slots fall to its margin (64 by default). @@ -3687,7 +3688,7 @@ mod tests { let network_opts = packet::PacketSimulatorOptions { node_count: replica_count, client_count: 1, - seed: 0xF0_2D_0001, + seed: 0xF02D_0001, ..packet::PacketSimulatorOptions::default() }; let mut sim = Simulator::with_shards_shell( diff --git a/core/simulator/src/workload/auditor.rs b/core/simulator/src/workload/auditor.rs index 903117935..96cd09258 100644 --- a/core/simulator/src/workload/auditor.rs +++ b/core/simulator/src/workload/auditor.rs @@ -217,6 +217,7 @@ impl ServerAuditor { self.in_flight.get(&key).map(|entry| entry.action) } + #[must_use] pub fn in_flight_count(&self) -> usize { self.in_flight.len() } diff --git a/core/simulator/src/workload/mod.rs b/core/simulator/src/workload/mod.rs index 564082bfd..9e25f9f82 100644 --- a/core/simulator/src/workload/mod.rs +++ b/core/simulator/src/workload/mod.rs @@ -141,7 +141,7 @@ impl Workload { /// Advance the resend clock by one tick. Called once per driver iteration; /// [`Self::due_resends`] measures against it. - pub fn tick(&mut self) { + pub const fn tick(&mut self) { self.now += 1; } @@ -550,9 +550,10 @@ pub fn run_with_faults( replies_seen } -/// 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. +/// Crash and restart injection with stability windows. +/// +/// Shaped after `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. /// /// Owns the fault PRNG so crash scheduling stays reproducible from the seed yet /// independent of the traffic draw order. Draws nothing while both probabilities diff --git a/core/simulator/src/workload/options.rs b/core/simulator/src/workload/options.rs index f5fd38acb..c21b12d50 100644 --- a/core/simulator/src/workload/options.rs +++ b/core/simulator/src/workload/options.rs @@ -32,9 +32,11 @@ pub const DEFAULT_REQUEST_TIMEOUT_TICKS: u64 = 200; /// 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. +/// 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 @@ -116,6 +118,10 @@ impl ActionWeights { /// `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. + /// + /// # Panics + /// Panics if the spread weights do not sum to 100, which would mean the + /// remainder arithmetic above is wrong rather than the caller. #[must_use] pub fn uniform() -> Self { use strum::IntoEnumIterator; diff --git a/core/simulator/src/workload/state_checker.rs b/core/simulator/src/workload/state_checker.rs index a3c5135f8..f2ffa3fef 100644 --- a/core/simulator/src/workload/state_checker.rs +++ b/core/simulator/src/workload/state_checker.rs @@ -15,8 +15,9 @@ // specific language governing permissions and limitations // under the License. -//! Cross-replica committed-log equality, after TigerBeetle's -//! `testing/cluster/state_checker.zig`. +//! 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 @@ -24,7 +25,7 @@ //! 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 +//! 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
