This is an automated email from the ASF dual-hosted git repository.

numinnex pushed a commit to branch fix_integration_flakes
in repository https://gitbox.apache.org/repos/asf/iggy.git

commit faf8d7a30e8b834e1201110eb325e9df4480d4d0
Author: Grzegorz Koszyk <[email protected]>
AuthorDate: Fri Aug 7 12:00:28 2026 +0200

    test(integration): fix vsr flakes from racing recovery routes
---
 .../tests/cluster/metadata_state_transfer.rs       |  28 ++-
 .../tests/cluster/partition_state_transfer.rs      |  67 ++++--
 .../integration/tests/connectors/http/http_sink.rs |  65 ++---
 core/integration/tests/data_integrity/mod.rs       |  11 +-
 .../verify_consumer_group_partition_assignment.rs  |  18 +-
 .../server/scenarios/message_cleanup_scenario.rs   | 129 ++++++----
 core/integration/tests/server/scenarios/mod.rs     |   2 -
 .../scenarios/stress_produce_consume_scenario.rs   | 264 ---------------------
 core/partitions/src/iggy_partition.rs              | 126 ++++++++++
 9 files changed, 339 insertions(+), 371 deletions(-)

diff --git a/core/integration/tests/cluster/metadata_state_transfer.rs 
b/core/integration/tests/cluster/metadata_state_transfer.rs
index b66bfacf7..275946412 100644
--- a/core/integration/tests/cluster/metadata_state_transfer.rs
+++ b/core/integration/tests/cluster/metadata_state_transfer.rs
@@ -301,21 +301,33 @@ async fn 
given_election_past_a_node_when_it_rejoins_stale_should_probe_then_stat
     // while the cluster sits at a higher view -- the stale-primary case.
     harness.restart_node_from_clean_slate(0).unwrap();
 
-    // It must first PROBE (its heartbeat-send timer converting, since it has 
no
-    // heartbeat-receive timer as a "primary"), then complete the transfer. The
-    // probe marker distinguishes this path from the same-view backstop.
+    // It must leave view 0 before it can transfer, and there are two correct
+    // routes off it. Either its heartbeat-SEND timer converts into a probe (it
+    // has no heartbeat-RECEIVE timer while it believes itself primary), or an
+    // unsolicited `StartView` from the live primary reaches it first and it
+    // adopts that. Which one wins is scheduler luck: the probe timer races the
+    // survivors' next StartView broadcast, and on an unloaded box the 
StartView
+    // lands ~450ms into boot and takes it. Pinning the probe marker alone made
+    // this spec fail whenever the node caught up the faster way. What must 
hold
+    // either way is that it left the stale view by a legitimate route and then
+    // converged rather than wedging.
     let deadline = Instant::now() + TRANSFER_BUDGET;
-    let mut probed = false;
+    let mut left_stale_view = false;
     loop {
-        probed = probed || harness.node(0).stdout_contains("probing to catch 
up");
+        left_stale_view = left_stale_view
+            || harness.node(0).stdout_contains("probing to catch up")
+            || harness
+                .node(0)
+                .stdout_contains("adopting view from StartView");
         if harness
             .node(0)
             .stdout_contains("metadata state transfer installed")
         {
             assert!(
-                probed,
-                "the rejoined node transferred without first probing; the \
-                 stale-view path must reach the transfer through a view probe"
+                left_stale_view,
+                "the rejoined node transferred while still believing view 0; 
it \
+                 must first leave the stale view, by its own probe or by 
adopting \
+                 a StartView"
             );
             break;
         }
diff --git a/core/integration/tests/cluster/partition_state_transfer.rs 
b/core/integration/tests/cluster/partition_state_transfer.rs
index 0d3647917..13413d886 100644
--- a/core/integration/tests/cluster/partition_state_transfer.rs
+++ b/core/integration/tests/cluster/partition_state_transfer.rs
@@ -68,6 +68,15 @@ const INSTALL_MARKER: &str = "partition state transfer 
installed";
 const FULLY_SERVED_MARKER: &str = "partition state transfer fully served";
 const ABANDON_MARKER: &str =
     "partition state transfer stalled past its retry budget; abandoning with a 
backed-off re-arm";
+/// The second route off a dead transfer peer. A replica parked in
+/// `AwaitingTarget` re-arms the moment it adopts a StartView, and with the
+/// stock config that beats the stall budget every time: the budget needs six
+/// rounds of `repair_retry_interval` (~6s) while `heartbeat_timeout` elects
+/// the new primary in 5s. Both routes prove the same thing - the replica did
+/// not retry into the corpse - so a spec that pins one is asserting on
+/// scheduler luck.
+const REARM_ON_VIEW_MARKER: &str =
+    "adopted a live view while awaiting transfer; requesting partition state 
transfer";
 
 /// Transfer end-to-end: adoption, repair round-trip, conversion, chunk pull,
 /// install, tail repair. CI runners are slow; bound without hanging the suite.
@@ -243,7 +252,7 @@ async fn 
given_evicted_ring_when_node_restarts_with_data_should_state_transfer_p
 // pull also runs the staged-segment reuse scan, but how much it can adopt
 // depends on how many artifacts completed before the kill, so nothing here
 // asserts on reuse.
-async fn 
given_transfer_peer_dies_when_stalled_should_abandon_and_recover_partition(
+async fn 
given_transfer_peer_dies_when_stalled_should_leave_dead_peer_and_recover_partition(
     harness: &mut TestHarness,
 ) {
     let client = harness
@@ -254,26 +263,26 @@ async fn 
given_transfer_peer_dies_when_stalled_should_abandon_and_recover_partit
     // Bulky payloads so the pull spans many 256 KiB chunks: the kill below
     // must land while the transfer is provably in flight, and a small
     // partition finishes inside the marker-poll latency, leaving the
-    // abandon path untested.
+    // dead-peer path untested.
     produce_bulky(&client, BULKY_MESSAGES_COUNT, BULKY_PAYLOAD_LEN).await;
     sleep(Duration::from_secs(1)).await;
     let _seed_client = client;
 
-    // Wipe node 2, wait until its rejoin CONVERTED to a transfer and the
-    // serving peer (the view-0 primary, node 0) proved it started serving
-    // the pull, then kill that peer mid-pull. Node 2 must not retry into
-    // the corpse forever: the stall budget abandons with a backed-off
-    // re-arm against the next peer, the survivors elect past node 0, and
-    // the transfer re-runs against the new primary (node 1).
+    // Wipe node 2, wait until its rejoin CONVERTED to a transfer, then kill
+    // the serving peer (the view-0 primary, node 0) mid-pull. Node 2 must not
+    // retry into the corpse forever: it re-arms against the next peer -- via
+    // the stall budget, or via the StartView it adopts once the survivors
+    // elect past node 0 -- and the transfer re-runs against the new primary
+    // (node 1).
     harness
         .restart_node_from_clean_slate(2)
         .expect("clean-slate restart of node 2");
     // Kill node 0 the moment node 2 CONVERTED: the transfer is then armed
     // at node 0 but the 64 MiB pull cannot possibly finish inside the kill
-    // latency, so node 2 deterministically ends up stalling against a dead
-    // peer -- whether the descriptor made it out or not, both funnels land
-    // in the stall budget. (Gating on the serving marker instead raced the
-    // pull itself: a release-build pull finishes in a few hundred ms.)
+    // latency, so node 2 deterministically ends up parked against a dead
+    // peer, whether the descriptor made it out or not. (Gating on the serving
+    // marker instead raced the pull itself: a release-build pull finishes in
+    // a few hundred ms.)
     let deadline = Instant::now() + TRANSFER_BUDGET;
     while !harness.node(2).stdout_contains(CONVERSION_MARKER) {
         assert!(
@@ -289,13 +298,22 @@ async fn 
given_transfer_peer_dies_when_stalled_should_abandon_and_recover_partit
     let installs_before_kill = 
harness.node(2).stdout_occurrences(INSTALL_MARKER);
     let served_before_kill = 
harness.node(1).stdout_occurrences(FULLY_SERVED_MARKER);
     let abandons_before_kill = 
harness.node(2).stdout_occurrences(ABANDON_MARKER);
+    let rearms_before_kill = 
harness.node(2).stdout_occurrences(REARM_ON_VIEW_MARKER);
     harness
         .stop_node(0)
         .expect("stop the serving peer (node 0)");
 
-    // The abandon is now deterministic: the pull was in flight against a
-    // peer that is gone, so the stall budget must exhaust.
-    await_new_marker(harness, 2, ABANDON_MARKER, abandons_before_kill).await;
+    // The pull was in flight against a peer that is gone, so node 2 must leave
+    // the dead session by one of the two routes off it, whichever fires first.
+    await_new_marker_any(
+        harness,
+        2,
+        &[
+            (ABANDON_MARKER, abandons_before_kill),
+            (REARM_ON_VIEW_MARKER, rearms_before_kill),
+        ],
+    )
+    .await;
 
     // Recovery: the scheduled re-arm targets the surviving primary. No
     // follow-up commit is asserted -- the cluster is quorum-marginal with
@@ -497,6 +515,25 @@ async fn await_new_marker(harness: &TestHarness, node: 
usize, marker: &str, base
     }
 }
 
+/// [`await_new_marker`] over alternative markers on one node, each carried
+/// with its own pre-fault baseline: satisfied by the first to advance.
+async fn await_new_marker_any(harness: &TestHarness, node: usize, markers: 
&[(&str, usize)]) {
+    let deadline = Instant::now() + TRANSFER_BUDGET;
+    loop {
+        if markers
+            .iter()
+            .any(|(marker, baseline)| 
harness.node(node).stdout_occurrences(marker) > *baseline)
+        {
+            return;
+        }
+        assert!(
+            Instant::now() < deadline,
+            "node {node} logged none of {markers:?} again after the fault 
within {TRANSFER_BUDGET:?}"
+        );
+        sleep(MARKER_POLL).await;
+    }
+}
+
 /// [`await_marker`] over a set of nodes: satisfied by the first one to log it.
 async fn await_marker_any(harness: &TestHarness, nodes: &[usize], marker: 
&str) {
     let deadline = Instant::now() + TRANSFER_BUDGET;
diff --git a/core/integration/tests/connectors/http/http_sink.rs 
b/core/integration/tests/connectors/http/http_sink.rs
index 7e20c25f8..0466e71f4 100644
--- a/core/integration/tests/connectors/http/http_sink.rs
+++ b/core/integration/tests/connectors/http/http_sink.rs
@@ -301,41 +301,36 @@ async fn 
individual_json_messages_delivered_as_separate_posts(
 // Test 2: NDJSON Batch Mode
 // ============================================================================
 
-/// Validates `batch_mode=ndjson`: all messages in one request as 
newline-delimited JSON.
-/// Checks single request, line count = message count, per-line envelope, 
`application/x-ndjson`.
-#[iggy_harness(
-    server(connectors_runtime(config_path = 
"tests/connectors/http/sink.toml")),
-    seed = seeds::connector_stream
-)]
-async fn ndjson_messages_delivered_as_single_request(
-    harness: &TestHarness,
-    fixture: HttpSinkNdjsonFixture,
-) {
-    let client = harness.root_client().await.unwrap();
-    let stream_id: Identifier = seeds::names::STREAM.try_into().unwrap();
-    let topic_id: Identifier = seeds::names::TOPIC.try_into().unwrap();
+/// Creates the connector stream/topic and pre-publishes the NDJSON test 
messages.
+///
+/// Same reason as [`connector_stream_with_json_array_messages`]: publishing 
from the
+/// test body races the sink's poll loop against the server's commit frontier, 
so a
+/// poll can observe a partial batch and flush it as its own request. Messages 
at rest
+/// before the first poll always arrive as one batch.
+async fn connector_stream_with_ndjson_messages(
+    client: &IggyClient,
+) -> Result<(), seeds::SeedError> {
+    seeds::connector_stream(client).await?;
+
+    let stream_id: Identifier = seeds::names::STREAM.try_into()?;
+    let topic_id: Identifier = seeds::names::TOPIC.try_into()?;
 
-    // Step 1: Build 3 JSON event messages
     let json_payloads: Vec<serde_json::Value> = vec![
         serde_json::json!({"event": "login", "user": 1}),
         serde_json::json!({"event": "click", "user": 2}),
         serde_json::json!({"event": "logout", "user": 3}),
     ];
 
-    let mut messages: Vec<IggyMessage> = json_payloads
-        .iter()
-        .enumerate()
-        .map(|(i, payload)| {
-            let bytes = serde_json::to_vec(payload).expect("Failed to 
serialize");
+    let mut messages: Vec<IggyMessage> = 
Vec::with_capacity(json_payloads.len());
+    for (i, payload) in json_payloads.iter().enumerate() {
+        messages.push(
             IggyMessage::builder()
                 .id((i + 1) as u128)
-                .payload(Bytes::from(bytes))
-                .build()
-                .expect("Failed to build message")
-        })
-        .collect();
+                .payload(Bytes::from(serde_json::to_vec(payload)?))
+                .build()?,
+        );
+    }
 
-    // Step 2: Publish messages to Iggy
     client
         .send_messages(
             &stream_id,
@@ -343,10 +338,20 @@ async fn ndjson_messages_delivered_as_single_request(
             &Partitioning::partition_id(0),
             &mut messages,
         )
-        .await
-        .expect("Failed to send messages");
+        .await?;
+
+    Ok(())
+}
 
-    // Step 3: Wait for single NDJSON request (all messages batched into one)
+/// Validates `batch_mode=ndjson`: all messages in one request as 
newline-delimited JSON.
+/// Checks single request, line count = message count, per-line envelope, 
`application/x-ndjson`.
+#[iggy_harness(
+    server(connectors_runtime(config_path = 
"tests/connectors/http/sink.toml")),
+    seed = connector_stream_with_ndjson_messages
+)]
+async fn ndjson_messages_delivered_as_single_request(fixture: 
HttpSinkNdjsonFixture) {
+    // Step 1: Wait for the single NDJSON request. The seed pre-published all
+    // messages, so the sink delivers them in one batch (see the seed docs 
above).
     let requests = fixture
         .container()
         .wait_for_requests(1)
@@ -357,7 +362,7 @@ async fn ndjson_messages_delivered_as_single_request(
     assert_eq!(req.method, "POST", "Expected POST method");
     assert_eq!(req.url, "/ingest", "Expected /ingest URL");
 
-    // Step 4: Parse NDJSON body — each line is a separate JSON envelope
+    // Step 2: Parse NDJSON body — each line is a separate JSON envelope
     let lines: Vec<&str> = req.body.trim().lines().collect();
     assert_eq!(
         lines.len(),
@@ -379,7 +384,7 @@ async fn ndjson_messages_delivered_as_single_request(
         );
     }
 
-    // Step 5: Verify NDJSON content type
+    // Step 3: Verify NDJSON content type
     let ct = req
         .header("Content-Type")
         .expect("Content-Type header must be present");
diff --git a/core/integration/tests/data_integrity/mod.rs 
b/core/integration/tests/data_integrity/mod.rs
index 620f8e872..5a3815d2d 100644
--- a/core/integration/tests/data_integrity/mod.rs
+++ b/core/integration/tests/data_integrity/mod.rs
@@ -15,10 +15,13 @@
 // specific language governing permissions and limitations
 // under the License.
 
-// Partially vsr-gated inside the module: the bench-fill test requires
-// PARTITION-plane state transfer (sub-floor stats/offset seeding), which is
-// not implemented yet; the metadata-only deletion/restart test runs under
-// vsr (metadata journal repair covers its rejoin window).
+// Partially vsr-gated inside the module: the remaining gates cover
+// `flush_unsaved_buffer`, which server-ng answers `FeatureUnavailable` and
+// which the eager-flush server envs replace under vsr. The bench-fill test
+// itself runs under vsr since PARTITION-plane state transfer landed, but the
+// harness spawns `iggy-bench` off disk with no cargo build-graph edge, so the
+// binary must have been built `--features vsr` or its login hangs on the
+// framing mismatch.
 mod verify_after_server_restart;
 mod verify_user_login_after_restart;
 
diff --git 
a/core/integration/tests/data_integrity/verify_consumer_group_partition_assignment.rs
 
b/core/integration/tests/data_integrity/verify_consumer_group_partition_assignment.rs
index f5f378837..b92c5bcf4 100644
--- 
a/core/integration/tests/data_integrity/verify_consumer_group_partition_assignment.rs
+++ 
b/core/integration/tests/data_integrity/verify_consumer_group_partition_assignment.rs
@@ -28,6 +28,11 @@ const STREAM_NAME: &str = "cg-partition-test-stream";
 const TOPIC_NAME: &str = "cg-partition-test-topic";
 const CONSUMER_GROUP_NAME: &str = "cg-partition-test-group";
 const PARTITIONS_COUNT: u32 = 3;
+/// Slices the slab-reuse wait so the surviving consumer can prove liveness
+/// inside the server's staleness window. Product of the two is the 3s the
+/// spec waits for the freed slab to become reusable.
+const CONSUMER1_KEEPALIVE_PINGS: u32 = 3;
+const CONSUMER1_KEEPALIVE_INTERVAL: Duration = Duration::from_secs(1);
 
 async fn create_stale_tcp_client(server_addr: &str) -> IggyClient {
     let config = TcpClientConfig {
@@ -3144,7 +3149,18 @@ async fn 
should_not_assign_partition_to_wrong_member_after_slab_reuse(harness: &
 
     // 3. Consumer2 (revocation target) disconnects — its slab is freed
     drop(client2);
-    sleep(Duration::from_secs(3)).await;
+    // Consumer1 must stay alive across the wait. This spec asks the server for
+    // `heartbeat.interval = 2s`, so its verifier evicts any consumer-group
+    // member idle past 1.2 intervals, and harness clients never ping on their
+    // own (the SDK pinger is spawned by `IggyClient::connect`, which the
+    // builder does not call). Silence here evicted consumer1 mid-wait and its
+    // offset store below came back `StaleClient`. Any request refreshes
+    // liveness; a ping is the cheapest. Consumer2 stays silent by construction
+    // - it is already dropped, and the socket close frees its slab.
+    for _ in 0..CONSUMER1_KEEPALIVE_PINGS {
+        sleep(CONSUMER1_KEEPALIVE_INTERVAL).await;
+        client1.ping().await.unwrap();
+    }
 
     // 4. Consumer3 joins — may reuse consumer2's old slab
     let client3 = harness.new_client().await.unwrap();
diff --git 
a/core/integration/tests/server/scenarios/message_cleanup_scenario.rs 
b/core/integration/tests/server/scenarios/message_cleanup_scenario.rs
index af2519729..51e67d6a4 100644
--- a/core/integration/tests/server/scenarios/message_cleanup_scenario.rs
+++ b/core/integration/tests/server/scenarios/message_cleanup_scenario.rs
@@ -17,9 +17,9 @@
 
 //! Tests for message retention policies (time-based and size-based).
 //!
-//! Configuration: 100KB segment size, 100ms cleaner interval, instant flush.
+//! Configuration: 10KiB segment size, 100ms cleaner interval, instant flush.
 //! Message size: 64B header + 936B payload = 1KB per message.
-//! Therefore: 100 messages = 1 segment, 101+ messages = 2+ segments.
+//! Therefore a segment holds 9 messages and every ~10 messages rotates one.
 
 use bytes::Bytes;
 use iggy::prelude::*;
@@ -48,7 +48,13 @@ pub async fn run_expiry_after_rotation(client: &IggyClient, 
data_path: &Path) {
     let stream = client.create_stream(STREAM_NAME).await.unwrap();
     let stream_id = stream.id;
 
-    let expiry = Duration::from_secs(2);
+    // The whole send burst has to land inside this window: expiry runs off 
each
+    // message's own timestamp, so a produce run that outlives it gets its
+    // oldest segments reclaimed before the pre-expiry poll ever runs. A 3-node
+    // vsr cluster in a debug build pays a consensus round-trip plus an fsync
+    // per request, which is what pushed the old one-message-per-request loop
+    // past the old 2s window.
+    let expiry = Duration::from_secs(4);
     let topic = client
         .create_topic(
             &Identifier::named(STREAM_NAME).unwrap(),
@@ -70,18 +76,24 @@ pub async fn run_expiry_after_rotation(client: &IggyClient, 
data_path: &Path) {
         .display()
         .to_string();
 
-    // Send 110 messages (1KB each) to create 2 segments (100KB segment size)
+    // Send 110 messages (1KB each) in batches, spanning several 10KiB 
segments.
+    // Batched rather than one request per message: the burst must fit inside
+    // `expiry` with room to spare, and each request costs a round-trip.
     let payload = make_payload('A');
-    let total_messages = 110;
-
-    for i in 0..total_messages {
-        let message = IggyMessage::builder()
-            .id(i as u128)
-            .payload(payload.clone())
-            .build()
-            .unwrap();
-
-        let mut messages = vec![message];
+    let total_messages: usize = 110;
+    let batch_size = 10;
+
+    for chunk_start in (0..total_messages).step_by(batch_size) {
+        let mut messages: Vec<IggyMessage> = (chunk_start
+            ..total_messages.min(chunk_start + batch_size))
+            .map(|i| {
+                IggyMessage::builder()
+                    .id(i as u128)
+                    .payload(payload.clone())
+                    .build()
+                    .unwrap()
+            })
+            .collect();
         client
             .send_messages(
                 &Identifier::named(STREAM_NAME).unwrap(),
@@ -117,7 +129,7 @@ pub async fn run_expiry_after_rotation(client: &IggyClient, 
data_path: &Path) {
 
     assert_eq!(
         polled_before.messages.len(),
-        total_messages as usize,
+        total_messages,
         "Should poll all messages before expiry"
     );
 
@@ -345,16 +357,24 @@ pub async fn run_combined_retention(client: &IggyClient, 
data_path: &Path) {
         .display()
         .to_string();
 
-    // Send 110 messages to create 2 segments (under size threshold, but will 
expire)
+    // Send 110 messages to create 2 segments (under size threshold, but will
+    // expire). Batched so the burst finishes well inside `expiry`: a 
per-message
+    // request pays a consensus round-trip plus an fsync, and a loop that
+    // outlives the window has its head reclaimed before the count below.
     let payload = make_payload('C');
-    for i in 0..110 {
-        let message = IggyMessage::builder()
-            .id(i as u128)
-            .payload(payload.clone())
-            .build()
-            .unwrap();
-
-        let mut messages = vec![message];
+    let total_messages: usize = 110;
+    let batch_size = 10;
+    for chunk_start in (0..total_messages).step_by(batch_size) {
+        let mut messages: Vec<IggyMessage> = (chunk_start
+            ..total_messages.min(chunk_start + batch_size))
+            .map(|i| {
+                IggyMessage::builder()
+                    .id(i as u128)
+                    .payload(payload.clone())
+                    .build()
+                    .unwrap()
+            })
+            .collect();
         client
             .send_messages(
                 &Identifier::named(STREAM_NAME).unwrap(),
@@ -412,19 +432,25 @@ pub async fn run_expiry_with_multiple_partitions(client: 
&IggyClient, data_path:
     let topic_id = topic.id;
 
     let payload = make_payload('D');
-    let messages_per_partition = 110;
+    let messages_per_partition: usize = 110;
+    let batch_size = 10;
 
-    // Send messages to all partitions
+    // Send messages to all partitions. Batched: a per-message request costs a
+    // consensus round-trip plus an fsync, and `PARTITIONS_COUNT` × 110 of 
those
+    // outlive `expiry`, so the cleaner would reclaim the first partition's
+    // sealed segments before the last one had even been written.
     for partition_id in 0..PARTITIONS_COUNT {
-        for i in 0..messages_per_partition {
-            let msg_id = partition_id as u128 * 1000 + i as u128;
-            let message = IggyMessage::builder()
-                .id(msg_id)
-                .payload(payload.clone())
-                .build()
-                .unwrap();
-
-            let mut messages = vec![message];
+        for chunk_start in (0..messages_per_partition).step_by(batch_size) {
+            let mut messages: Vec<IggyMessage> = (chunk_start
+                ..messages_per_partition.min(chunk_start + batch_size))
+                .map(|i| {
+                    IggyMessage::builder()
+                        .id(partition_id as u128 * 1000 + i as u128)
+                        .payload(payload.clone())
+                        .build()
+                        .unwrap()
+                })
+                .collect();
             client
                 .send_messages(
                     &Identifier::named(STREAM_NAME).unwrap(),
@@ -594,10 +620,11 @@ pub async fn run_expiry_respects_consumer_offset(client: 
&IggyClient, data_path:
     let stream = client.create_stream(TEST_STREAM).await.unwrap();
     let stream_id = stream.id;
 
-    // Expiry must outlast the send + first-poll phase: 300 serial sends with
-    // per-message fsync (and VSR quorum in cluster mode) take ~3s under load.
-    // If segments expire before the consumer commits its first offset, there 
is
-    // no barrier yet and the cleaner legally deletes them, breaking the 
premise.
+    // Expiry must outlast the send + first-poll phase. If segments expire
+    // before the consumer commits its first offset, there is no barrier yet
+    // and the cleaner legally deletes them, breaking the premise: the poll
+    // below then starts at the earliest surviving offset instead of 0.
+    // The sends are batched for the same reason (see below).
     let expiry = Duration::from_secs(4);
     let topic = client
         .create_topic(
@@ -620,21 +647,29 @@ pub async fn run_expiry_respects_consumer_offset(client: 
&IggyClient, data_path:
         .display()
         .to_string();
 
-    // Send 300 messages (1KB each) -> 3 sealed segments + active
+    // Send 300 messages (1KB each) -> 3 sealed segments + active. Batched:
+    // one request per message costs a consensus round-trip plus an fsync each,
+    // which on a 3-node debug cluster runs the burst well past `expiry`.
     let payload = make_payload('B');
     let total_messages = 300u32;
-    for i in 0..total_messages {
-        let message = IggyMessage::builder()
-            .id(i as u128)
-            .payload(payload.clone())
-            .build()
-            .unwrap();
+    let batch_size = 10u32;
+    for chunk_start in (0..total_messages).step_by(batch_size as usize) {
+        let mut messages: Vec<IggyMessage> = (chunk_start
+            ..total_messages.min(chunk_start + batch_size))
+            .map(|i| {
+                IggyMessage::builder()
+                    .id(i as u128)
+                    .payload(payload.clone())
+                    .build()
+                    .unwrap()
+            })
+            .collect();
         client
             .send_messages(
                 &Identifier::named(TEST_STREAM).unwrap(),
                 &Identifier::named(TEST_TOPIC).unwrap(),
                 &Partitioning::partition_id(PARTITION_ID),
-                &mut [message],
+                &mut messages,
             )
             .await
             .unwrap();
diff --git a/core/integration/tests/server/scenarios/mod.rs 
b/core/integration/tests/server/scenarios/mod.rs
index 99d20bcf8..5287c3053 100644
--- a/core/integration/tests/server/scenarios/mod.rs
+++ b/core/integration/tests/server/scenarios/mod.rs
@@ -50,8 +50,6 @@ pub mod single_message_per_batch_scenario;
 pub mod snapshot_scenario;
 pub mod stale_client_consumer_group_scenario;
 pub mod stream_size_validation_scenario;
-#[cfg(feature = "vsr")]
-pub mod stress_produce_consume_scenario;
 pub mod system_scenario;
 pub mod tcp_tls_scenario;
 pub mod timestamp_scenario;
diff --git 
a/core/integration/tests/server/scenarios/stress_produce_consume_scenario.rs 
b/core/integration/tests/server/scenarios/stress_produce_consume_scenario.rs
deleted file mode 100644
index 507df4300..000000000
--- a/core/integration/tests/server/scenarios/stress_produce_consume_scenario.rs
+++ /dev/null
@@ -1,264 +0,0 @@
-// Licensed to the Apache Software Foundation (ASF) under one
-// or more contributor license agreements.  See the NOTICE file
-// distributed with this work for additional information
-// regarding copyright ownership.  The ASF licenses this file
-// to you under the Apache License, Version 2.0 (the
-// "License"); you may not use this file except in compliance
-// with the License.  You may obtain a copy of the License at
-//
-//   http://www.apache.org/licenses/LICENSE-2.0
-//
-// Unless required by applicable law or agreed to in writing,
-// software distributed under the License is distributed on an
-// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
-// KIND, either express or implied.  See the License for the
-// specific language governing permissions and limitations
-// under the License.
-
-//! Data-plane produce+consume stress across server-ng topologies: `PRODUCERS`
-//! producers and `CONSUMERS` consumers all hammer a SINGLE partition, 
asserting
-//! no message loss and a strictly contiguous offset log.
-//!
-//! Targets the partition-ref-across-await UB fix: the consume poll path and 
the
-//! produce/commit pump run as sibling tasks over the same partition, so
-//! concentrating every producer and consumer on one partition maximizes the
-//! `&`/`&mut` aliasing window on that partition's pump that the fix closes. 
All
-//! producers run concurrently with all consumers for `HAMMER_DURATION`; the
-//! consumers then drain. A single partition lives on a single shard, so the
-//! multi-shard variants still spin up N shards but concentrate the load on the
-//! one owning shard.
-//!
-//! Oracle: producers interleave on the partition's shared offset sequence, so
-//! per-producer contiguity does not hold. Instead every consumer reads the
-//! partition in full and must observe a contiguous `0..total` (no gap = no 
loss,
-//! no dup) with a count equal to the sum of all producers' sends.
-//!
-//! Strictly data plane: polls by explicit offset with `auto_commit = false` 
and
-//! performs no mid-run topic/partition mutation, so it never drives the 
metadata
-//! consensus plane concurrently. That avoids a separate, still-open `on_ack`
-//! journal-durability race that panics the primary under concurrent metadata 
ops
-//! (see the gated `concurrent_produce_consume_scenario` in 
`scenarios/mod.rs`).
-
-use bytes::Bytes;
-use iggy::prelude::*;
-use integration::harness::TestHarness;
-use integration::iggy_harness;
-use std::sync::Arc;
-use std::sync::atomic::{AtomicBool, Ordering};
-use std::time::{Duration, Instant};
-
-const STREAM_NAME: &str = "stress-pc-stream";
-const TOPIC_NAME: &str = "stress-pc-topic";
-// All traffic targets this one partition to concentrate produce+consume
-// contention on a single partition pump.
-const PARTITION_ID: u32 = 0;
-const PRODUCERS: u32 = 4;
-const CONSUMERS: u32 = 4;
-const PRODUCER_BATCH: u32 = 16;
-const CONSUMER_BATCH: u32 = 64;
-const HAMMER_DURATION: Duration = Duration::from_secs(20);
-// Safety net so a wedged consumer fails loudly instead of hanging the suite.
-const MAX_TEST_DURATION: Duration = Duration::from_secs(120);
-// Whole-test wall-clock guard. A server that dies at boot leaves the harness
-// client retrying connect with no cap, and a parked poll never re-checks
-// MAX_TEST_DURATION, so without this the suite hangs indefinitely instead of
-// failing. Set above MAX_TEST_DURATION so a slow-but-progressing consumer
-// trips its own informative deadline first.
-const WALL_CLOCK_TIMEOUT: Duration = Duration::from_secs(150);
-// Empty polls observed after producers stop before the partition is declared 
drained.
-const DRAIN_EMPTY_POLLS: u32 = 20;
-
-/// Single-node, single shard (`"1"`) and multi shard (`"2"`).
-#[iggy_harness(
-    cluster_nodes = 1,
-    server(system.sharding.cpu_allocation = ["1", "2"])
-)]
-async fn 
given_single_node_when_produce_consume_hammered_should_not_lose_messages(
-    harness: &TestHarness,
-) {
-    run_hammer(harness).await;
-}
-
-/// Three-node cluster, single shard (`"1"`) and multi shard (`"2"`) per node.
-/// Heavy (3 servers * N shards); run on demand with `--ignored`.
-#[iggy_harness(
-    cluster_nodes = 3,
-    server(system.sharding.cpu_allocation = ["1", "2"])
-)]
-#[ignore = "3-node cluster: heavy, run on demand with --ignored"]
-async fn given_cluster_when_produce_consume_hammered_should_not_lose_messages(
-    harness: &TestHarness,
-) {
-    run_hammer(harness).await;
-}
-
-async fn run_hammer(harness: &TestHarness) {
-    tokio::time::timeout(WALL_CLOCK_TIMEOUT, run_hammer_inner(harness))
-        .await
-        .expect("stress test exceeded WALL_CLOCK_TIMEOUT; server likely 
crashed at boot or a poll wedged");
-}
-
-async fn run_hammer_inner(harness: &TestHarness) {
-    let stream_id = Identifier::named(STREAM_NAME).unwrap();
-
-    let setup = harness.tcp_root_client().await.unwrap();
-    setup.create_stream(STREAM_NAME).await.unwrap();
-    setup
-        .create_topic(
-            &stream_id,
-            TOPIC_NAME,
-            1,
-            CompressionAlgorithm::None,
-            None,
-            IggyExpiry::NeverExpire,
-            MaxTopicSize::ServerDefault,
-        )
-        .await
-        .unwrap();
-    drop(setup);
-
-    let producer_done = Arc::new(AtomicBool::new(false));
-
-    // Spawn consumers first so they poll concurrently with the producers from
-    // the very first send. Each reads the whole partition independently by
-    // explicit offset and asserts strict contiguity (no gap = no loss, no 
dup).
-    let mut consumers = Vec::with_capacity(CONSUMERS as usize);
-    for consumer_id in 0..CONSUMERS {
-        let client = harness.tcp_root_client().await.unwrap();
-        let done = producer_done.clone();
-        consumers.push(tokio::spawn(consume_partition(client, consumer_id, 
done)));
-    }
-
-    // All producers hammer the single partition for HAMMER_DURATION.
-    let mut producers = Vec::with_capacity(PRODUCERS as usize);
-    for producer_id in 0..PRODUCERS {
-        let client = harness.tcp_root_client().await.unwrap();
-        producers.push(tokio::spawn(produce_partition(client, producer_id)));
-    }
-
-    // Producers stop at the hammer deadline; sum their sends, then signal
-    // consumers to drain.
-    let mut total_sent = 0u64;
-    for handle in producers {
-        total_sent += handle.await.unwrap();
-    }
-    producer_done.store(true, Ordering::Relaxed);
-
-    assert!(
-        total_sent > 0,
-        "hammer produced no messages; workload wiring is broken"
-    );
-
-    // Every consumer independently read the full partition; each must have 
seen
-    // exactly the committed total, contiguously (asserted inside the task).
-    for (consumer_id, handle) in consumers.into_iter().enumerate() {
-        let received = handle.await.unwrap();
-        assert_eq!(
-            received, total_sent,
-            "consumer {consumer_id}: consumed {received} != produced 
{total_sent} (message loss)",
-        );
-    }
-
-    let cleanup = harness.tcp_root_client().await.unwrap();
-    cleanup.delete_stream(&stream_id).await.unwrap();
-}
-
-/// Send `PRODUCER_BATCH`-sized batches to the shared partition until the 
hammer
-/// deadline. Returns the number of messages sent (each send awaits commit).
-async fn produce_partition(client: IggyClient, producer_id: u32) -> u64 {
-    let stream = Identifier::named(STREAM_NAME).unwrap();
-    let topic = Identifier::named(TOPIC_NAME).unwrap();
-    let partitioning = Partitioning::partition_id(PARTITION_ID);
-    let deadline = Instant::now() + HAMMER_DURATION;
-    let mut sent = 0u64;
-
-    while Instant::now() < deadline {
-        let mut messages: Vec<IggyMessage> = (0..PRODUCER_BATCH)
-            .map(|i| {
-                IggyMessage::builder()
-                    .payload(Bytes::from(format!(
-                        "prod{producer_id}-{}",
-                        sent + u64::from(i)
-                    )))
-                    .build()
-                    .unwrap()
-            })
-            .collect();
-        client
-            .send_messages(&stream, &topic, &partitioning, &mut messages)
-            .await
-            .unwrap_or_else(|e| panic!("producer {producer_id} send failed at 
sent={sent}: {e}"));
-        sent += u64::from(PRODUCER_BATCH);
-    }
-    sent
-}
-
-/// Read the shared partition in full by explicit offset with `auto_commit =
-/// false`, asserting each message arrives at the next contiguous offset. 
Drains
-/// until producers are done and `DRAIN_EMPTY_POLLS` consecutive empty polls
-/// confirm the tail. Returns the number of messages received.
-async fn consume_partition(
-    client: IggyClient,
-    consumer_id: u32,
-    producer_done: Arc<AtomicBool>,
-) -> u64 {
-    let stream = Identifier::named(STREAM_NAME).unwrap();
-    let topic = Identifier::named(TOPIC_NAME).unwrap();
-    let consumer = Consumer::default();
-    let mut next_offset = 0u64;
-    let mut received = 0u64;
-    let mut consecutive_empty = 0u32;
-    let deadline = Instant::now() + MAX_TEST_DURATION;
-
-    loop {
-        assert!(
-            Instant::now() < deadline,
-            "consumer {consumer_id} timed out: received {received}, 
next_offset {next_offset}"
-        );
-
-        let polled = match client
-            .poll_messages(
-                &stream,
-                &topic,
-                Some(PARTITION_ID),
-                &consumer,
-                &PollingStrategy::offset(next_offset),
-                CONSUMER_BATCH,
-                false,
-            )
-            .await
-        {
-            Ok(polled) => polled,
-            Err(e) => {
-                // Transient under load; back off and retry.
-                eprintln!("consumer {consumer_id} poll error: {e:?}");
-                tokio::time::sleep(Duration::from_millis(10)).await;
-                continue;
-            }
-        };
-
-        if polled.messages.is_empty() {
-            if producer_done.load(Ordering::Relaxed) {
-                consecutive_empty += 1;
-                if consecutive_empty >= DRAIN_EMPTY_POLLS {
-                    break;
-                }
-            }
-            tokio::time::sleep(Duration::from_millis(5)).await;
-            continue;
-        }
-
-        consecutive_empty = 0;
-        for msg in &polled.messages {
-            assert_eq!(
-                msg.header.offset, next_offset,
-                "consumer {consumer_id} offset gap/dup: expected 
{next_offset}, got {}",
-                msg.header.offset
-            );
-            next_offset += 1;
-            received += 1;
-        }
-    }
-
-    received
-}
diff --git a/core/partitions/src/iggy_partition.rs 
b/core/partitions/src/iggy_partition.rs
index dad3f70e3..54a31ee67 100644
--- a/core/partitions/src/iggy_partition.rs
+++ b/core/partitions/src/iggy_partition.rs
@@ -158,6 +158,24 @@ where
     /// generation against this and resets only when it advances, so a 
redundant
     /// reconcile pass never re-wipes a partition already at this generation.
     pub(crate) applied_purge_generation: u64,
+    /// Highest partition op the last local purge superseded. A consumer-offset
+    /// op journaled before the purge but committed after it must NOT be
+    /// applied: `PurgeTopic` rides the metadata plane while
+    /// `StoreConsumerOffset*` rides this one, so the two logs carry no mutual
+    /// ordering, and a server-side auto-commit is dispatched with no client
+    /// waiting on it (`server-ng::dispatch::submit_auto_commit`) -- it can
+    /// still be in flight when the purge lands. Without this fence the commit
+    /// walk rewrites the offset map AND the offset file after the purge
+    /// cleared both, and nothing ever clears them again: this replica's
+    /// `applied_purge_generation` already matches the committed one, so the
+    /// reconciler sees no outstanding purge. The purge already fences the
+    /// other asynchronous writer of this state (the state-transfer install);
+    /// this is the symmetric fence for the offset-commit path.
+    ///
+    /// Memory-only on purpose, exactly like `applied_purge_generation`: a
+    /// restarted replica rehydrates its offsets from an already-empty
+    /// directory, so there is nothing for a persisted fence to protect.
+    pub(crate) consumer_offset_purge_fence: u64,
     /// Durable superblock for this partition's consensus group, recording
     /// `(view, log_view)` across a crash so this replica can never
     /// re-participate in a view older than one it advertised. `None` for
@@ -268,6 +286,10 @@ where
             .field("recovered_durable_offset", &self.recovered_durable_offset)
             .field("observed_view", &self.observed_view)
             .field("applied_purge_generation", &self.applied_purge_generation)
+            .field(
+                "consumer_offset_purge_fence",
+                &self.consumer_offset_purge_fence,
+            )
             .finish_non_exhaustive()
     }
 }
@@ -436,6 +458,7 @@ where
             persisted_offsets: RefCell::new(HashMap::new()),
             observed_view,
             applied_purge_generation: 0,
+            consumer_offset_purge_fence: 0,
             superblock: None,
             superblock_lock: LocalGate::new(),
             superblock_write_failures: Cell::new(0),
@@ -964,6 +987,17 @@ where
         &mut self,
         op: u64,
     ) -> Result<(), IggyError> {
+        // Superseded by a purge: this op was journaled before the partition 
was
+        // reset, so its offset describes messages that no longer exist. Drop 
it
+        // rather than let the commit walk write it back over the cleared state
+        // (see `consumer_offset_purge_fence`). A store admitted after the 
purge
+        // carries a higher op -- the primary assigns them monotonically at
+        // admission -- so it lands above the fence and applies normally, and a
+        // suppressed in-flight auto-commit is re-derived from the next poll.
+        if op <= self.consumer_offset_purge_fence {
+            self.pending_consumer_offset_commits.remove(&op);
+            return Ok(());
+        }
         // Peek (copy) instead of remove: if `persist_consumer_offset_commit`
         // fails (e.g. disk full, fd exhausted) the pending entry must remain
         // stageable for retry on the next apply. Removing first would strand
@@ -3823,6 +3857,16 @@ where
         self.installed_frontier = None;
         self.segment_checksum_cache.borrow_mut().clear();
 
+        // Fence the offset-commit path before clearing anything: every op the
+        // journal already holds predates this purge, so committing one later
+        // must not resurrect the offset the clear below is about to drop. The
+        // sequencer's current value is the right line -- it covers ops that 
are
+        // journaled but not yet committed, which is exactly where an in-flight
+        // auto-commit sits. Staged entries for those ops go too, so the commit
+        // walk cannot apply one straight out of the staging table.
+        self.consumer_offset_purge_fence = 
self.consensus().sequencer().current_sequence();
+        self.pending_consumer_offset_commits.clear();
+
         // Clear consumer + consumer-group offsets (memory + disk). Collect the
         // file paths before deleting so the map guard is not held across an
         // await.
@@ -5714,6 +5758,88 @@ mod tests {
         let _ = std::fs::remove_dir_all(&dir);
     }
 
+    /// A consumer-offset op journaled before a purge must not resurrect the
+    /// offset when its commit lands afterwards. `PurgeTopic` rides the 
metadata
+    /// plane while `StoreConsumerOffset*` rides this one, so the two logs 
carry
+    /// no mutual ordering; a server auto-commit in particular is dispatched
+    /// with nobody waiting on it and can still be in flight when the purge
+    /// runs. Applying it afterwards used to rewrite both the offset map and 
the
+    /// offset file over the cleared state, permanently: this replica's
+    /// `applied_purge_generation` already matches the committed one, so the
+    /// reconciler sees no outstanding purge to re-run.
+    #[compio::test]
+    async fn purge_fences_a_consumer_offset_op_journaled_before_it() {
+        let dir = std::env::temp_dir().join(format!(
+            "iggy-purge-offset-fence-{}-{}",
+            std::process::id(),
+            std::time::SystemTime::now()
+                .duration_since(std::time::UNIX_EPOCH)
+                .expect("system clock after epoch")
+                .as_nanos(),
+        ));
+        compio::fs::create_dir_all(&dir)
+            .await
+            .expect("create temp partition dir");
+
+        let mut partition = test_partition();
+        partition.set_partition_dir(dir.to_string_lossy().into_owned());
+        let consumer_id = 7u32;
+
+        // Journalled (staged) before the purge, not yet committed.
+        partition.consensus().sequencer().set_sequence(9);
+        partition.stage_consumer_offset_upsert(
+            9,
+            ConsumerKind::ConsumerGroup,
+            consumer_id,
+            9,
+            true,
+        );
+
+        partition
+            .purge(&repair_config(), 1)
+            .await
+            .expect("purge partition");
+
+        assert_eq!(
+            partition.consumer_offset_purge_fence, 9,
+            "the purge must fence every op the journal already held",
+        );
+        assert!(
+            partition.pending_consumer_offset_commits.is_empty(),
+            "the purge must drop staged pre-purge offset commits",
+        );
+
+        // The commit walk reaching that op afterwards is a no-op, not an 
error:
+        // refusing it would fail the commit and wedge the partition.
+        partition
+            .apply_staged_consumer_offset_commit(9)
+            .await
+            .expect("a fenced commit must not fail the commit walk");
+        assert!(
+            partition.consumer_group_offsets.pin().is_empty(),
+            "a pre-purge offset must not be resurrected after the purge",
+        );
+
+        // A store admitted after the purge carries a higher op -- the primary
+        // assigns them monotonically at admission -- so it still applies.
+        partition.stage_consumer_offset_upsert(
+            10,
+            ConsumerKind::ConsumerGroup,
+            consumer_id,
+            4,
+            false,
+        );
+        partition
+            .apply_staged_consumer_offset_commit(10)
+            .await
+            .expect("a post-purge commit applies");
+        assert_eq!(
+            partition.consumer_group_offsets.pin().len(),
+            1,
+            "a store admitted after the purge must survive the fence",
+        );
+    }
+
     /// Purge unlinks every segment and recreates the same paths, so a poll
     /// suspended across the purge must not keep serving the old inodes through
     /// its cached read state. The wipe reaches the in-flight clone through the

Reply via email to