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 46df538e1d4ae4eec4a2d903297e9dce6a4aaca2 Author: Krishna Vishal <[email protected]> AuthorDate: Fri Aug 14 14:28:55 2026 +0530 fix(simulator): draw the wire consumer kind, not a bare boolean The four consumer-offset ops sampled `consumer_kind` as `u8::from(bool)`, so half of every request they emitted carried kind 0. `WireConsumer` recognises only 1 (consumer) and 2 (consumer group), so those requests failed to decode; the partition plane answers an unparseable consumer-offset request by logging a WARN and dropping the frame without a reply. A client blocked on that request never hears back, and with one in-flight slot per client the first malformed draw wedged the workload for the rest of the run. `ActionWeights::default` carries `StoreConsumerOffset2` at weight 25, so every default-weight workload hit it within a handful of requests; a 3000-tick run drained 6 replies instead of ~270. It went unnoticed because the fuzzer hardcoded `SendMessages: 100` and the in-tree workload tests only assert that some reply arrived, never that the run drains. `KIND_CONSUMER` becomes public alongside the sibling it is always compared against, so a synthesized consumer names the discriminant rather than guessing a small integer. Still one bool draw per sample, so the PRNG trace shape is unchanged and only the encoded bytes move. The determinism baseline is re-locked: the old value was recorded over a trace that wedged a few replies in. --- core/binary_protocol/src/lib.rs | 2 +- core/binary_protocol/src/primitives/consumer.rs | 6 +++++- core/simulator/src/lib.rs | 9 ++++++++- .../src/workload/ops/delete_consumer_offset.rs | 3 ++- .../src/workload/ops/delete_consumer_offset_2.rs | 3 ++- core/simulator/src/workload/ops/mod.rs | 21 ++++++++++++++++++++- .../src/workload/ops/store_consumer_offset.rs | 3 ++- .../src/workload/ops/store_consumer_offset_2.rs | 3 ++- 8 files changed, 42 insertions(+), 8 deletions(-) diff --git a/core/binary_protocol/src/lib.rs b/core/binary_protocol/src/lib.rs index 0825d0f44..4ed40c33e 100644 --- a/core/binary_protocol/src/lib.rs +++ b/core/binary_protocol/src/lib.rs @@ -88,7 +88,7 @@ pub use message_view::{ WireMessageIterator, WireMessageIteratorMut, WireMessageView, WireMessageViewMut, }; pub use primitives::ack_level::AckLevel; -pub use primitives::consumer::{KIND_CONSUMER_GROUP, WireConsumer}; +pub use primitives::consumer::{KIND_CONSUMER, KIND_CONSUMER_GROUP, WireConsumer}; pub use primitives::identifier::{MAX_WIRE_NAME_LENGTH, WireIdentifier, WireName}; pub use primitives::partition_assignment::CreatedPartitionAssignment; pub use primitives::partitioning::{MAX_MESSAGES_KEY_LENGTH, WirePartitioning}; diff --git a/core/binary_protocol/src/primitives/consumer.rs b/core/binary_protocol/src/primitives/consumer.rs index 202de84f1..1e7739b37 100644 --- a/core/binary_protocol/src/primitives/consumer.rs +++ b/core/binary_protocol/src/primitives/consumer.rs @@ -20,7 +20,11 @@ use crate::WireIdentifier; use crate::codec::{WireDecode, WireEncode, read_u8}; use bytes::{BufMut, BytesMut}; -const KIND_CONSUMER: u8 = 1; +/// 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`. pub const KIND_CONSUMER_GROUP: u8 = 2; diff --git a/core/simulator/src/lib.rs b/core/simulator/src/lib.rs index 70f24b7d3..2aafb34ee 100644 --- a/core/simulator/src/lib.rs +++ b/core/simulator/src/lib.rs @@ -1675,8 +1675,15 @@ mod tests { // Re-locked again when replies stopped echoing a group id (the // client wire lost its namespace field): the reply-hash tuple // dropped that component. + // Re-locked again when the consumer-offset ops began drawing the WIRE + // consumer kind (1 / 2) instead of a bare boolean (0 / 1). Kind 0 is + // not a `WireConsumer` discriminant, so every such request used to be + // dropped unparsed with no reply, wedging the client's only in-flight + // slot; the previous baseline was recorded over a trace that stalled a + // few replies in. Same single bool draw, so only the encoded bytes and + // the replies they now earn moved. See `ops::sample_consumer_kind`. assert_eq!( - h1, 0xCF1F_BC79_B44A_65F7, + h1, 0x25C7_0F6A_5D0F_A8B2, "workload reply hash drifted from locked baseline" ); } diff --git a/core/simulator/src/workload/ops/delete_consumer_offset.rs b/core/simulator/src/workload/ops/delete_consumer_offset.rs index c05c7960a..aac37c71b 100644 --- a/core/simulator/src/workload/ops/delete_consumer_offset.rs +++ b/core/simulator/src/workload/ops/delete_consumer_offset.rs @@ -25,6 +25,7 @@ use server_common::sharding::IggyNamespace; use crate::client::SimClient; use crate::workload::effect::Effect; +use crate::workload::ops::sample_consumer_kind; use crate::workload::options::WorkloadOptions; use crate::workload::shadow::Shadow; @@ -51,7 +52,7 @@ pub fn sample( match outcome { Outcome::Success => { let ns = shadow.pick_namespace(prng)?; - let consumer_kind: u8 = u8::from(prng.random::<bool>()); + let consumer_kind = sample_consumer_kind(prng); let consumer_id: u32 = prng.random_range(0..options.consumer_pool_size.max(1)); Some(Input { ns, diff --git a/core/simulator/src/workload/ops/delete_consumer_offset_2.rs b/core/simulator/src/workload/ops/delete_consumer_offset_2.rs index ad2a6201a..8cd936145 100644 --- a/core/simulator/src/workload/ops/delete_consumer_offset_2.rs +++ b/core/simulator/src/workload/ops/delete_consumer_offset_2.rs @@ -25,6 +25,7 @@ use server_common::sharding::IggyNamespace; use crate::client::SimClient; use crate::workload::effect::Effect; +use crate::workload::ops::sample_consumer_kind; use crate::workload::options::WorkloadOptions; use crate::workload::shadow::Shadow; @@ -52,7 +53,7 @@ pub fn sample( match outcome { Outcome::Success => { let ns = shadow.pick_namespace(prng)?; - let consumer_kind: u8 = u8::from(prng.random::<bool>()); + let consumer_kind = sample_consumer_kind(prng); let consumer_id: u32 = prng.random_range(0..options.consumer_pool_size.max(1)); let f: f32 = prng.random(); let ack = if f < options.ack_quorum_ratio { diff --git a/core/simulator/src/workload/ops/mod.rs b/core/simulator/src/workload/ops/mod.rs index df6ffed48..d74d5b0fa 100644 --- a/core/simulator/src/workload/ops/mod.rs +++ b/core/simulator/src/workload/ops/mod.rs @@ -55,7 +55,8 @@ pub mod update_stream; pub mod update_topic; pub mod update_user; -use iggy_binary_protocol::RoutedRequestHeader; +use iggy_binary_protocol::{KIND_CONSUMER, KIND_CONSUMER_GROUP, RoutedRequestHeader}; +use rand::RngExt; use rand_xoshiro::Xoshiro256Plus; use server_common::Message; @@ -65,6 +66,24 @@ use crate::workload::effect::Effect; use crate::workload::options::WorkloadOptions; use crate::workload::shadow::Shadow; +/// Draw a consumer kind for the four consumer-offset ops, as the WIRE +/// discriminant rather than a bare boolean. +/// +/// `WireConsumer::decode` accepts only [`KIND_CONSUMER`] (1) and +/// [`KIND_CONSUMER_GROUP`] (2); anything else is an `UnknownDiscriminant`, +/// which `parse_consumer_offset_request` maps to `IggyError::InvalidCommand` +/// and the partition plane answers by logging a WARN and dropping the frame +/// with NO reply. A client blocked on that request never gets one back, so a +/// single malformed draw wedges its in-flight slot for the rest of the run. +/// One bool draw either way, so the PRNG trace shape is unchanged. +pub(crate) fn sample_consumer_kind(prng: &mut Xoshiro256Plus) -> u8 { + if prng.random::<bool>() { + KIND_CONSUMER_GROUP + } else { + KIND_CONSUMER + } +} + /// Generates per-op enums (`InFlightInput`, `InFlightOutcome`) plus four /// dispatch fns over a fixed `(Action, module)` table. Missing variants /// are a compile error via the exhaustive `match` arms. diff --git a/core/simulator/src/workload/ops/store_consumer_offset.rs b/core/simulator/src/workload/ops/store_consumer_offset.rs index edefbfc72..2a9795b7d 100644 --- a/core/simulator/src/workload/ops/store_consumer_offset.rs +++ b/core/simulator/src/workload/ops/store_consumer_offset.rs @@ -26,6 +26,7 @@ use server_common::sharding::IggyNamespace; use crate::client::SimClient; use crate::workload::effect::Effect; +use crate::workload::ops::sample_consumer_kind; use crate::workload::options::WorkloadOptions; use crate::workload::shadow::Shadow; @@ -53,7 +54,7 @@ pub fn sample( match outcome { Outcome::Success => { let ns = shadow.pick_namespace(prng)?; - let consumer_kind: u8 = u8::from(prng.random::<bool>()); + let consumer_kind = sample_consumer_kind(prng); let consumer_id: u32 = prng.random_range(0..options.consumer_pool_size.max(1)); // Draw against the configured ceiling, then clamp to committed // reality so the offset is reachable. Clamping post-draw keeps diff --git a/core/simulator/src/workload/ops/store_consumer_offset_2.rs b/core/simulator/src/workload/ops/store_consumer_offset_2.rs index baac692d0..e745d8ced 100644 --- a/core/simulator/src/workload/ops/store_consumer_offset_2.rs +++ b/core/simulator/src/workload/ops/store_consumer_offset_2.rs @@ -31,6 +31,7 @@ use server_common::sharding::IggyNamespace; use crate::client::SimClient; use crate::workload::effect::Effect; +use crate::workload::ops::sample_consumer_kind; use crate::workload::options::WorkloadOptions; use crate::workload::shadow::Shadow; @@ -59,7 +60,7 @@ pub fn sample( match outcome { Outcome::Success => { let ns = shadow.pick_namespace(prng)?; - let consumer_kind: u8 = u8::from(prng.random::<bool>()); + let consumer_kind = sample_consumer_kind(prng); let consumer_id: u32 = prng.random_range(0..options.consumer_pool_size.max(1)); // Draw against the configured ceiling, then clamp to committed // reality so the offset is reachable. Clamping post-draw keeps
