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

Reply via email to