Merge branch 'master' into dotnet-protocol
diff --git a/core/bench/src/actors/consumer/client/low_level.rs b/core/bench/src/actors/consumer/client/low_level.rs
index 7f47748..964231f 100644
--- a/core/bench/src/actors/consumer/client/low_level.rs
+++ b/core/bench/src/actors/consumer/client/low_level.rs
@@ -22,6 +22,7 @@
 use crate::utils::ClientFactory;
 use crate::utils::{batch_total_size_bytes, batch_user_size_bytes};
 use iggy::prelude::*;
+use std::collections::HashMap;
 use std::sync::Arc;
 use std::time::Duration;
 use tokio::time::Instant;
@@ -36,7 +37,9 @@
     partition_id: Option<u32>,
     polling_strategy: PollingStrategy,
     auto_commit: bool,
-    offset: u64,
+    /// Where offset polling continues in each partition. A group member polls its partitions
+    /// round-robin, so one shared cursor would skip the start of every partition but the first.
+    next_offsets: HashMap<u32, u64>,
 }
 
 impl LowLevelConsumerClient {
@@ -51,7 +54,7 @@
             partition_id: None,
             polling_strategy: PollingStrategy::next(),
             auto_commit: true,
-            offset: 0,
+            next_offsets: HashMap::new(),
         }
     }
 }
@@ -62,14 +65,21 @@
         let consumer = self.consumer.as_ref().expect("consumer not initialized");
         let messages_to_receive = self.config.messages_per_batch.get();
 
+        let polling_strategy = self.polling_strategy;
+        let next_offsets = &self.next_offsets;
+        let strategy_for = |partition_id: u32| {
+            next_offsets
+                .get(&partition_id)
+                .map_or(polling_strategy, |offset| PollingStrategy::offset(*offset))
+        };
         let before_poll = Instant::now();
         let polled = client
-            .poll_messages(
+            .poll_messages_with_strategy_for(
                 &self.stream_id,
                 &self.topic_id,
                 self.partition_id,
                 consumer,
-                &self.polling_strategy,
+                &strategy_for,
                 messages_to_receive,
                 self.auto_commit,
             )
@@ -99,9 +109,11 @@
         let user_bytes = batch_user_size_bytes(&polled);
         let total_bytes = batch_total_size_bytes(&polled);
 
-        self.offset += messages_count;
-        if self.polling_strategy.kind == PollingKind::Offset {
-            self.polling_strategy.value += messages_count;
+        if self.polling_strategy.kind == PollingKind::Offset
+            && let Some(last) = polled.messages.last()
+        {
+            self.next_offsets
+                .insert(polled.partition_id, last.header.offset + 1);
         }
 
         Ok(Some(BatchMetrics {
@@ -137,7 +149,7 @@
         )
         .await;
         let (polling_strategy, auto_commit) = match self.config.polling_kind {
-            PollingKind::Offset => (PollingStrategy::offset(self.offset), false),
+            PollingKind::Offset => (PollingStrategy::offset(0), false),
             PollingKind::Next => (PollingStrategy::next(), true),
             _ => panic!("Unsupported polling kind: {:?}", self.config.polling_kind),
         };
diff --git a/core/common/src/consumer_group_client_state.rs b/core/common/src/consumer_group_client_state.rs
index 1e694cc..d4331cb 100644
--- a/core/common/src/consumer_group_client_state.rs
+++ b/core/common/src/consumer_group_client_state.rs
@@ -185,6 +185,24 @@
             .cloned()
             .collect()
     }
+
+    /// Drop what the consensus session owned. Membership is per connection
+    /// and the coordinator fences assignments by a generation it tracks per
+    /// session, so nothing synced under the old session holds once it is
+    /// reset. The balanced cursors and the partition counts stay: they belong
+    /// to a topic, not to a session, and clearing them would restart the
+    /// produce round-robin at partition 0 and cost a metadata round trip per
+    /// topic on every reconnect.
+    pub fn clear_session_scoped(&self) {
+        self.assignments
+            .lock()
+            .unwrap_or_else(std::sync::PoisonError::into_inner)
+            .clear();
+        self.joined_groups
+            .lock()
+            .unwrap_or_else(std::sync::PoisonError::into_inner)
+            .clear();
+    }
 }
 
 #[cfg(test)]
@@ -242,4 +260,24 @@
         state.deregister_group("s|t|g");
         assert!(!state.is_registered("s|t|g"));
     }
+
+    #[test]
+    fn session_reset_drops_membership_and_keeps_topic_state() {
+        let state = ConsumerGroupClientState::new();
+        let id = Identifier::named("g").unwrap();
+        state.register_group("s|t|g".to_owned(), id.clone(), id.clone(), id);
+        state.set_assignment("s|t|g".to_owned(), 1, vec![0, 1]);
+        assert_eq!(state.next_balanced_partition("s|t", 3), 0);
+        state.set_partition_count("s|t".to_owned(), 3);
+
+        state.clear_session_scoped();
+
+        assert!(state.registered_groups().is_empty());
+        assert!(!state.is_registered("s|t|g"));
+        assert!(!state.has_assignment("s|t|g"));
+        // Topic state outlives the session: the produce cursor carries on and
+        // the partition count is still cached.
+        assert_eq!(state.next_balanced_partition("s|t", 3), 1);
+        assert_eq!(state.partition_count("s|t"), Some(3));
+    }
 }
diff --git a/core/common/src/lib.rs b/core/common/src/lib.rs
index 89a9fc7..e74eae2 100644
--- a/core/common/src/lib.rs
+++ b/core/common/src/lib.rs
@@ -41,6 +41,13 @@
 /// a genuine end-of-partition empty poll, which echoes the real partition id.
 pub const RESYNC_REQUIRED_PARTITION_SENTINEL: u32 = u32::MAX;
 
+/// Client-side `partition_id` of an empty poll reply for a consumer-group
+/// member that currently holds no partitions. The server never sends it: the
+/// transport fills it in when the synced assignment is empty, so a caller can
+/// tell "nothing assigned" from a genuine empty poll, which echoes the real
+/// partition id. Same value as the Go and Node SDKs use for that case.
+pub const NO_ASSIGNED_PARTITION: u32 = u32::MAX - 1;
+
 /// Frozen ceiling on the widest batch record any admission path can persist.
 /// Two knobs bound admission and both are validated against this at boot:
 /// `message_bus.max_message_size` caps every bus-framed wire message, and
diff --git a/core/common/src/traits/binary_impls/messages.rs b/core/common/src/traits/binary_impls/messages.rs
index bf720a1..d8b1241 100644
--- a/core/common/src/traits/binary_impls/messages.rs
+++ b/core/common/src/traits/binary_impls/messages.rs
@@ -166,14 +166,15 @@
 }
 
 /// Poll a consumer group: select one of the member's assigned partitions
-/// (round-robin) and send an explicit-partition poll. A coordinator fence
-/// rejection (stale assignment after a rebalance) triggers one re-sync + retry.
+/// (round-robin), ask `strategy_for` where to read it from and send an
+/// explicit-partition poll. A coordinator fence rejection (stale assignment
+/// after a rebalance) triggers one re-sync + retry.
 async fn poll_group_messages<B: BinaryClient>(
     client: &B,
     stream_id: &Identifier,
     topic_id: &Identifier,
     consumer: &Consumer,
-    strategy: &PollingStrategy,
+    strategy_for: &(dyn Fn(u32) -> PollingStrategy + Send + Sync),
     count: u32,
     auto_commit: bool,
 ) -> Result<PolledMessages, IggyError> {
@@ -195,14 +196,19 @@
                     topic_id.clone(),
                 ));
             }
-            return Ok(PolledMessages::empty());
+            return Ok(PolledMessages {
+                partition_id: crate::NO_ASSIGNED_PARTITION,
+                ..PolledMessages::empty()
+            });
         };
+        // Resolved per attempt: a fence retry can land on another partition.
+        let strategy = strategy_for(partition_id);
         let request = PollMessagesRequest {
             consumer: consumer_to_wire(consumer)?,
             stream_id: identifier_to_wire(stream_id)?,
             topic_id: identifier_to_wire(topic_id)?,
             partition_id: Some(partition_id),
-            strategy: polling_strategy_to_wire(strategy),
+            strategy: polling_strategy_to_wire(&strategy),
             count,
             auto_commit,
         };
@@ -284,6 +290,28 @@
         count: u32,
         auto_commit: bool,
     ) -> Result<PolledMessages, IggyError> {
+        self.poll_messages_with_strategy_for(
+            stream_id,
+            topic_id,
+            partition_id,
+            consumer,
+            &|_: u32| *strategy,
+            count,
+            auto_commit,
+        )
+        .await
+    }
+
+    async fn poll_messages_with_strategy_for(
+        &self,
+        stream_id: &Identifier,
+        topic_id: &Identifier,
+        partition_id: Option<u32>,
+        consumer: &Consumer,
+        strategy_for: &(dyn Fn(u32) -> PollingStrategy + Send + Sync),
+        count: u32,
+        auto_commit: bool,
+    ) -> Result<PolledMessages, IggyError> {
         fail_if_not_authenticated(self).await?;
         // VSR: a consumer-group poll without an explicit partition is resolved
         // client-side from the member's cached assignment (the broker routes
@@ -294,18 +322,19 @@
                 stream_id,
                 topic_id,
                 consumer,
-                strategy,
+                strategy_for,
                 count,
                 auto_commit,
             )
             .await;
         }
+        let strategy = strategy_for(partition_id.unwrap_or(0));
         let req = PollMessagesRequest {
             consumer: consumer_to_wire(consumer)?,
             stream_id: identifier_to_wire(stream_id)?,
             topic_id: identifier_to_wire(topic_id)?,
             partition_id,
-            strategy: polling_strategy_to_wire(strategy),
+            strategy: polling_strategy_to_wire(&strategy),
             count,
             auto_commit,
         };
diff --git a/core/common/src/traits/message_client.rs b/core/common/src/traits/message_client.rs
index 0c5fa11..18c6d22 100644
--- a/core/common/src/traits/message_client.rs
+++ b/core/common/src/traits/message_client.rs
@@ -16,8 +16,8 @@
 // under the License.
 
 use crate::{
-    Consumer, Identifier, IggyError, IggyMessage, Partitioning, PolledMessages, PollingStrategy,
-    SendMessagesResponse,
+    Consumer, ConsumerKind, Identifier, IggyError, IggyMessage, Partitioning, PolledMessages,
+    PollingStrategy, SendMessagesResponse,
 };
 use async_trait::async_trait;
 
@@ -29,6 +29,7 @@
     /// Authentication is required, and the permission to poll the messages.
     ///
     /// Polling a consumer group the client is not (or no longer) a member of fails with `ConsumerGroupMemberNotFound` rather than returning an empty batch, so the caller can rejoin.
+    /// A member that holds no partitions gets an empty batch whose `partition_id` is [`NO_ASSIGNED_PARTITION`](crate::NO_ASSIGNED_PARTITION).
     #[allow(clippy::too_many_arguments)]
     async fn poll_messages(
         &self,
@@ -41,6 +42,42 @@
         auto_commit: bool,
     ) -> Result<PolledMessages, IggyError>;
 
+    /// [`poll_messages`](Self::poll_messages) whose strategy is chosen once the partition is
+    /// known. A consumer-group poll without a partition picks one of the member's assigned
+    /// partitions first and then asks `strategy_for` for it, so a caller can continue every
+    /// partition from its own position. Any other poll has its partition up front and asks
+    /// `strategy_for` for that one, or for `0` when none was given, which is the partition the
+    /// server reads then.
+    ///
+    /// The default implementation is for transports that cannot pick a partition client-side:
+    /// a consumer-group poll without a partition fails with `FeatureUnavailable`.
+    #[allow(clippy::too_many_arguments)]
+    async fn poll_messages_with_strategy_for(
+        &self,
+        stream_id: &Identifier,
+        topic_id: &Identifier,
+        partition_id: Option<u32>,
+        consumer: &Consumer,
+        strategy_for: &(dyn Fn(u32) -> PollingStrategy + Send + Sync),
+        count: u32,
+        auto_commit: bool,
+    ) -> Result<PolledMessages, IggyError> {
+        if consumer.kind == ConsumerKind::ConsumerGroup && partition_id.is_none() {
+            return Err(IggyError::FeatureUnavailable);
+        }
+        let strategy = strategy_for(partition_id.unwrap_or(0));
+        self.poll_messages(
+            stream_id,
+            topic_id,
+            partition_id,
+            consumer,
+            &strategy,
+            count,
+            auto_commit,
+        )
+        .await
+    }
+
     /// Send messages using specified partitioning strategy to the given stream and topic by unique IDs or names.
     ///
     /// Authentication is required, and the permission to send the messages.
diff --git a/core/common/src/types/message/polled_messages.rs b/core/common/src/types/message/polled_messages.rs
index af23d4e..d974fa4 100644
--- a/core/common/src/types/message/polled_messages.rs
+++ b/core/common/src/types/message/polled_messages.rs
@@ -29,7 +29,11 @@
 /// - `messages`: the collection of messages.
 #[derive(Debug, Serialize, Deserialize)]
 pub struct PolledMessages {
-    /// The identifier of the partition. If it's '0', then there's no partition assigned to the consumer group member.
+    /// The identifier of the partition. An empty reply can carry a sentinel instead of a real
+    /// id: [`NO_ASSIGNED_PARTITION`](crate::NO_ASSIGNED_PARTITION) for a consumer-group member
+    /// that holds no partitions, or
+    /// [`RESYNC_REQUIRED_PARTITION_SENTINEL`](crate::RESYNC_REQUIRED_PARTITION_SENTINEL) when
+    /// the server fenced a stale group assignment.
     pub partition_id: u32,
     /// The current offset of the partition.
     pub current_offset: u64,
diff --git a/core/integration/tests/sdk/consumer_group.rs b/core/integration/tests/sdk/consumer_group.rs
index b670de3..f233560 100644
--- a/core/integration/tests/sdk/consumer_group.rs
+++ b/core/integration/tests/sdk/consumer_group.rs
@@ -15,6 +15,7 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use std::collections::HashMap;
 use std::str::FromStr;
 use std::time::Duration;
 
@@ -29,6 +30,9 @@
 const CONSUMER_USERNAME: &str = "consumer-group-rejoin-user";
 const CONSUMER_PASSWORD: &str = "password123";
 const CONSUMER_REJOIN_TIMEOUT: Duration = Duration::from_secs(10);
+const PARTITIONS_COUNT: u32 = 2;
+const MESSAGES_PER_PARTITION: u32 = 5;
+const READ_TIMEOUT: Duration = Duration::from_secs(10);
 
 // Pins a 60s server heartbeat because harness clients never ping on their own:
 // the SDK pinger is spawned by `IggyClient::connect`, which the harness builder
@@ -169,3 +173,82 @@
     assert_eq!(group.members_count, 1);
     assert_eq!(group.members.len(), 1);
 }
+
+// A group member polls its partitions round-robin. Under a strategy other than `next()` the
+// continuation must be kept per partition: one shared cursor would ask the second partition for
+// the offset reached in the first one and skip its beginning.
+#[iggy_harness(
+    test_client_transport = [Tcp, WebSocket, Quic],
+    server(heartbeat.enabled = true, heartbeat.interval = "60s")
+)]
+async fn given_offset_strategy_when_member_polls_two_partitions_should_read_each_from_its_start(
+    harness: &TestHarness,
+) {
+    let root_client = harness
+        .root_client()
+        .await
+        .expect("Failed to get root client");
+    let stream_id = Identifier::named(STREAM_NAME).unwrap();
+    let topic_id = Identifier::named(TOPIC_NAME).unwrap();
+
+    root_client.create_stream(STREAM_NAME).await.unwrap();
+    root_client
+        .create_topic(
+            &stream_id,
+            TOPIC_NAME,
+            &TopicCreateOptions {
+                partitions_count: Some(PARTITIONS_COUNT),
+                message_expiry: Some(IggyExpiry::NeverExpire),
+                ..TopicCreateOptions::default()
+            },
+        )
+        .await
+        .unwrap();
+    for partition_id in 0..PARTITIONS_COUNT {
+        let mut messages: Vec<IggyMessage> = (0..MESSAGES_PER_PARTITION)
+            .map(|index| IggyMessage::from_str(&format!("{partition_id}-{index}")).unwrap())
+            .collect();
+        root_client
+            .send_messages(
+                &stream_id,
+                &topic_id,
+                &Partitioning::partition_id(partition_id),
+                &mut messages,
+            )
+            .await
+            .unwrap();
+    }
+
+    let mut consumer = root_client
+        .consumer_group(CONSUMER_GROUP_NAME, STREAM_NAME, TOPIC_NAME)
+        .unwrap()
+        .polling_strategy(PollingStrategy::offset(0))
+        .batch_length(MESSAGES_PER_PARTITION)
+        .auto_commit(AutoCommit::Disabled)
+        .build();
+    consumer.init().await.unwrap();
+
+    let expected_total = (PARTITIONS_COUNT * MESSAGES_PER_PARTITION) as usize;
+    let mut offsets_by_partition: HashMap<u32, Vec<u64>> = HashMap::new();
+    for _ in 0..expected_total {
+        let received = timeout(READ_TIMEOUT, consumer.next())
+            .await
+            .expect("every partition must be read from its start before the timeout")
+            .expect("consumer stream should remain open")
+            .expect("polling must not fail");
+        offsets_by_partition
+            .entry(received.partition_id)
+            .or_default()
+            .push(received.message.header.offset);
+    }
+    consumer.shutdown().await.unwrap();
+
+    let expected_offsets: Vec<u64> = (0..u64::from(MESSAGES_PER_PARTITION)).collect();
+    for partition_id in 0..PARTITIONS_COUNT {
+        assert_eq!(
+            offsets_by_partition.get(&partition_id),
+            Some(&expected_offsets),
+            "partition {partition_id} must be read from offset 0 without gaps"
+        );
+    }
+}
diff --git a/core/integration/tests/sdk/consumer_group_membership.rs b/core/integration/tests/sdk/consumer_group_membership.rs
index 429aea6..a8e7f59 100644
--- a/core/integration/tests/sdk/consumer_group_membership.rs
+++ b/core/integration/tests/sdk/consumer_group_membership.rs
@@ -15,10 +15,14 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use std::str::FromStr;
+use std::sync::Arc;
 use std::time::Duration;
 
 use futures::StreamExt;
+use iggy::prelude::locking::IggyRwLockFn;
 use iggy::prelude::*;
+use iggy_common::{BinaryTransport, ConsumerGroupClientState};
 use integration::iggy_harness;
 use tokio::time::timeout;
 
@@ -156,6 +160,161 @@
     }
 }
 
+// A member holding zero partitions gets an empty poll reply whose partition id
+// is the `NO_ASSIGNED_PARTITION` sentinel, so a caller can tell it from an
+// empty partition, which echoes its real id.
+#[iggy_harness(
+    test_client_transport = [Tcp, WebSocket, Quic],
+    server(heartbeat.enabled = true, heartbeat.interval = "60s")
+)]
+async fn given_member_holds_no_partitions_when_polled_should_report_no_assigned_partition(
+    harness: &TestHarness,
+) {
+    let stream_id = Identifier::named(STREAM_NAME).unwrap();
+    let topic_id = Identifier::named(TOPIC_NAME).unwrap();
+    let group_id = Identifier::named(CONSUMER_GROUP_NAME).unwrap();
+
+    let mut clients = harness
+        .root_clients(2)
+        .await
+        .expect("Failed to create root clients");
+    let first_member = clients.remove(0);
+    let second_member = clients.remove(0);
+
+    first_member.create_stream(STREAM_NAME).await.unwrap();
+    first_member
+        .create_topic(
+            &stream_id,
+            TOPIC_NAME,
+            &TopicCreateOptions {
+                partitions_count: Some(1),
+                message_expiry: Some(IggyExpiry::NeverExpire),
+                ..TopicCreateOptions::default()
+            },
+        )
+        .await
+        .unwrap();
+    first_member
+        .create_consumer_group(&stream_id, &topic_id, CONSUMER_GROUP_NAME)
+        .await
+        .unwrap();
+
+    // The first member to join keeps the topic's only partition.
+    first_member
+        .join_consumer_group(&stream_id, &topic_id, &group_id)
+        .await
+        .unwrap();
+    second_member
+        .join_consumer_group(&stream_id, &topic_id, &group_id)
+        .await
+        .unwrap();
+
+    let consumer = Consumer::group(group_id.clone());
+    let owner_poll = first_member
+        .poll_messages(
+            &stream_id,
+            &topic_id,
+            None,
+            &consumer,
+            &PollingStrategy::next(),
+            1,
+            false,
+        )
+        .await
+        .unwrap();
+    assert!(owner_poll.messages.is_empty());
+    assert_eq!(
+        owner_poll.partition_id, 0,
+        "an empty poll of an owned partition must echo its id"
+    );
+
+    let surplus_poll = second_member
+        .poll_messages(
+            &stream_id,
+            &topic_id,
+            None,
+            &consumer,
+            &PollingStrategy::next(),
+            1,
+            false,
+        )
+        .await
+        .unwrap();
+    assert!(surplus_poll.messages.is_empty());
+    assert_eq!(
+        surplus_poll.partition_id, NO_ASSIGNED_PARTITION,
+        "a member without partitions must report the sentinel"
+    );
+}
+
+// A member built with `do_not_auto_join_consumer_group()` relies on the caller for the join, so
+// polling must not wait for a join the consumer itself never performs.
+#[iggy_harness(
+    test_client_transport = [Tcp, WebSocket, Quic],
+    server(heartbeat.enabled = true, heartbeat.interval = "60s")
+)]
+async fn given_member_that_does_not_auto_join_when_joined_by_the_caller_should_receive_messages(
+    harness: &TestHarness,
+) {
+    let stream_id = Identifier::named(STREAM_NAME).unwrap();
+    let topic_id = Identifier::named(TOPIC_NAME).unwrap();
+    let group_id = Identifier::named(CONSUMER_GROUP_NAME).unwrap();
+
+    let client = harness.new_client().await.expect("Failed to create client");
+    client
+        .login_user(DEFAULT_ROOT_USERNAME, DEFAULT_ROOT_PASSWORD)
+        .await
+        .unwrap();
+    client.create_stream(STREAM_NAME).await.unwrap();
+    client
+        .create_topic(
+            &stream_id,
+            TOPIC_NAME,
+            &TopicCreateOptions {
+                partitions_count: Some(1),
+                message_expiry: Some(IggyExpiry::NeverExpire),
+                ..TopicCreateOptions::default()
+            },
+        )
+        .await
+        .unwrap();
+    client
+        .create_consumer_group(&stream_id, &topic_id, CONSUMER_GROUP_NAME)
+        .await
+        .unwrap();
+    let mut messages = vec![IggyMessage::from_str("message").unwrap()];
+    client
+        .send_messages(
+            &stream_id,
+            &topic_id,
+            &Partitioning::partition_id(0),
+            &mut messages,
+        )
+        .await
+        .unwrap();
+
+    // Membership is per connection, so the join goes through the client the consumer uses.
+    client
+        .join_consumer_group(&stream_id, &topic_id, &group_id)
+        .await
+        .unwrap();
+
+    let mut consumer = client
+        .consumer_group(CONSUMER_GROUP_NAME, STREAM_NAME, TOPIC_NAME)
+        .unwrap()
+        .batch_length(1)
+        .do_not_auto_join_consumer_group()
+        .build();
+    consumer.init().await.unwrap();
+
+    let received = timeout(RESOLVE_TIMEOUT, consumer.next())
+        .await
+        .expect("a member joined by the caller must poll instead of waiting for a join")
+        .expect("consumer stream should remain open")
+        .expect("the member must receive the message from its assigned partition");
+    assert_eq!(received.message.payload, "message");
+}
+
 // End-to-end wire pin for the consumer-group join/leave error ladder. The
 // metadata STM unit tests pin the committed result codes; this pins that
 // the server actually emits them over the wire, so a client observes the same
@@ -252,3 +411,120 @@
         error.as_code(),
     );
 }
+
+// Membership is per connection: a reconnect registers a new client identity
+// that the coordinator knows as a member of nothing, so the transport must not
+// carry the old session's membership and assignment into the new one. Carried
+// over, the first poll would run against a stale assignment instead of
+// reporting the missing membership straight away.
+#[iggy_harness(
+    test_client_transport = [Tcp, WebSocket, Quic],
+    server(heartbeat.enabled = true, heartbeat.interval = "60s")
+)]
+async fn given_group_member_when_session_reset_should_forget_membership(harness: &TestHarness) {
+    let stream_id = Identifier::named(STREAM_NAME).unwrap();
+    let topic_id = Identifier::named(TOPIC_NAME).unwrap();
+    let group_id = Identifier::named(CONSUMER_GROUP_NAME).unwrap();
+
+    let client = harness.new_client().await.expect("Failed to create client");
+    client
+        .login_user(DEFAULT_ROOT_USERNAME, DEFAULT_ROOT_PASSWORD)
+        .await
+        .unwrap();
+    client.create_stream(STREAM_NAME).await.unwrap();
+    client
+        .create_topic(
+            &stream_id,
+            TOPIC_NAME,
+            &TopicCreateOptions {
+                partitions_count: Some(1),
+                message_expiry: Some(IggyExpiry::NeverExpire),
+                ..TopicCreateOptions::default()
+            },
+        )
+        .await
+        .unwrap();
+    client
+        .create_consumer_group(&stream_id, &topic_id, CONSUMER_GROUP_NAME)
+        .await
+        .unwrap();
+    client
+        .join_consumer_group(&stream_id, &topic_id, &group_id)
+        .await
+        .unwrap();
+
+    let consumer = Consumer::group(group_id.clone());
+    // The first group poll syncs the assignment, which is what registers the
+    // membership in the transport cache.
+    client
+        .poll_messages(
+            &stream_id,
+            &topic_id,
+            None,
+            &consumer,
+            &PollingStrategy::next(),
+            1,
+            false,
+        )
+        .await
+        .unwrap();
+
+    let state = consumer_group_state(&client).await;
+    // The key the join and the group poll build from the same identifiers.
+    let key = format!("{stream_id}|{topic_id}|{group_id}");
+    assert!(
+        state.is_registered(&key),
+        "the sync must register the membership"
+    );
+    assert!(
+        state.has_assignment(&key),
+        "the sole member must hold the topic's only partition"
+    );
+
+    client.disconnect().await.unwrap();
+    // Reconnects through the transport rather than `IggyClient::connect`: the
+    // heartbeat the latter spawns re-syncs every registered group on its own,
+    // and this asserts on the reset alone.
+    client.client().read().await.connect().await.unwrap();
+    client
+        .login_user(DEFAULT_ROOT_USERNAME, DEFAULT_ROOT_PASSWORD)
+        .await
+        .unwrap();
+
+    assert!(
+        state.registered_groups().is_empty(),
+        "a reconnected client must not carry the old session's membership"
+    );
+    assert!(!state.is_registered(&key));
+    assert!(!state.has_assignment(&key));
+
+    // The new identity never joined, and the poll must say so instead of
+    // running against the old assignment.
+    let poll = client
+        .poll_messages(
+            &stream_id,
+            &topic_id,
+            None,
+            &consumer,
+            &PollingStrategy::next(),
+            1,
+            false,
+        )
+        .await;
+    assert!(
+        matches!(poll, Err(IggyError::ConsumerGroupMemberNotFound(..))),
+        "expected ConsumerGroupMemberNotFound for a poll without a rejoin, got {poll:?}"
+    );
+}
+
+/// The consumer-group cache lives on the transport `IggyClient` wraps.
+async fn consumer_group_state(client: &IggyClient) -> Arc<ConsumerGroupClientState> {
+    match &*client.client().read().await {
+        ClientWrapper::Tcp(client) => client.consumer_group_state(),
+        ClientWrapper::Quic(client) => client.consumer_group_state(),
+        ClientWrapper::WebSocket(client) => client.consumer_group_state(),
+        ClientWrapper::Http(_) | ClientWrapper::Iggy(_) => {
+            panic!("the consumer-group cache is a binary-transport concern")
+        }
+    }
+}
diff --git a/core/integration/tests/sdk/consumer_shutdown.rs b/core/integration/tests/sdk/consumer_shutdown.rs
new file mode 100644
index 0000000..172a569
--- /dev/null
+++ b/core/integration/tests/sdk/consumer_shutdown.rs
@@ -0,0 +1,112 @@
+// 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.
+
+use std::str::FromStr;
+
+use futures::StreamExt;
+use iggy::prelude::*;
+use integration::iggy_harness;
+use tokio::time::{Duration, Instant, timeout};
+
+const STREAM_NAME: &str = "consumer-shutdown-stream";
+const TOPIC_NAME: &str = "consumer-shutdown-topic";
+const CONSUMER_NAME: &str = "consumer-shutdown-consumer";
+const PARTITION_ID: u32 = 0;
+const POLL_TIMEOUT: Duration = Duration::from_secs(10);
+const OFFSET_DRAIN_TIMEOUT: Duration = Duration::from_secs(5);
+
+// `AutoCommit::Disabled` leaves every commit to the caller, so the final flush of `shutdown()`
+// must not run either: it would commit a message whose handler failed.
+#[iggy_harness]
+async fn given_disabled_auto_commit_when_shutdown_should_not_store_the_reading_position(
+    harness: &TestHarness,
+) {
+    let client = harness.root_client().await.unwrap();
+    let stream_id = Identifier::named(STREAM_NAME).unwrap();
+    let topic_id = Identifier::named(TOPIC_NAME).unwrap();
+
+    client.create_stream(STREAM_NAME).await.unwrap();
+    client
+        .create_topic(
+            &stream_id,
+            TOPIC_NAME,
+            &TopicCreateOptions {
+                partitions_count: Some(1),
+                message_expiry: Some(IggyExpiry::NeverExpire),
+                ..TopicCreateOptions::default()
+            },
+        )
+        .await
+        .unwrap();
+    let mut messages = vec![
+        IggyMessage::from_str("message_1").unwrap(),
+        IggyMessage::from_str("message_2").unwrap(),
+    ];
+    client
+        .send_messages(
+            &stream_id,
+            &topic_id,
+            &Partitioning::partition_id(PARTITION_ID),
+            &mut messages,
+        )
+        .await
+        .unwrap();
+
+    // Both messages must come in one batch: with nothing committed, `next()` serves the same
+    // batch again and the consumer's own filter drops it, so a second poll would stall.
+    let mut consumer = client
+        .consumer(CONSUMER_NAME, STREAM_NAME, TOPIC_NAME, PARTITION_ID)
+        .unwrap()
+        .auto_commit(AutoCommit::Disabled)
+        .batch_length(2)
+        .offset_drain_timeout(IggyDuration::from(OFFSET_DRAIN_TIMEOUT))
+        .build();
+    consumer.init().await.unwrap();
+
+    for expected_offset in [0, 1] {
+        let received = timeout(POLL_TIMEOUT, consumer.next())
+            .await
+            .expect("Consumer should receive a message before timeout")
+            .expect("Consumer stream should remain open")
+            .expect("Consumer should poll a message");
+        assert_eq!(received.message.header.offset, expected_offset);
+    }
+    assert_eq!(consumer.get_last_consumed_offset(PARTITION_ID), Some(1));
+
+    // The store task has to exit on the shutdown flag. One that never exits only costs the drain
+    // timeout and a warning, so the stored offset alone would not catch it.
+    let shutdown_started = Instant::now();
+    consumer.shutdown().await.unwrap();
+    assert!(
+        shutdown_started.elapsed() < OFFSET_DRAIN_TIMEOUT,
+        "shutdown() must not wait for the offset drain timeout"
+    );
+
+    let stored_offset = client
+        .get_consumer_offset(
+            &Consumer::new(Identifier::named(CONSUMER_NAME).unwrap()),
+            &stream_id,
+            &topic_id,
+            Some(PARTITION_ID),
+        )
+        .await
+        .unwrap();
+    assert!(
+        stored_offset.is_none(),
+        "nothing must be stored under AutoCommit::Disabled, got {stored_offset:?}"
+    );
+}
diff --git a/core/integration/tests/sdk/mod.rs b/core/integration/tests/sdk/mod.rs
index 065b257..0cec5f8 100644
--- a/core/integration/tests/sdk/mod.rs
+++ b/core/integration/tests/sdk/mod.rs
@@ -18,6 +18,7 @@
 mod consumer_group;
 mod consumer_group_membership;
 mod consumer_offset;
+mod consumer_shutdown;
 mod disconnect_relogin;
 mod hello_world;
 mod http_refresh;
diff --git a/core/sdk/src/client_wrappers/binary_message_client.rs b/core/sdk/src/client_wrappers/binary_message_client.rs
index eb3eed9..3fb7256 100644
--- a/core/sdk/src/client_wrappers/binary_message_client.rs
+++ b/core/sdk/src/client_wrappers/binary_message_client.rs
@@ -35,15 +35,37 @@
         count: u32,
         auto_commit: bool,
     ) -> Result<PolledMessages, IggyError> {
+        self.poll_messages_with_strategy_for(
+            stream_id,
+            topic_id,
+            partition_id,
+            consumer,
+            &|_: u32| *strategy,
+            count,
+            auto_commit,
+        )
+        .await
+    }
+
+    async fn poll_messages_with_strategy_for(
+        &self,
+        stream_id: &Identifier,
+        topic_id: &Identifier,
+        partition_id: Option<u32>,
+        consumer: &Consumer,
+        strategy_for: &(dyn Fn(u32) -> PollingStrategy + Send + Sync),
+        count: u32,
+        auto_commit: bool,
+    ) -> Result<PolledMessages, IggyError> {
         match self {
             ClientWrapper::Iggy(client) => {
                 client
-                    .poll_messages(
+                    .poll_messages_with_strategy_for(
                         stream_id,
                         topic_id,
                         partition_id,
                         consumer,
-                        strategy,
+                        strategy_for,
                         count,
                         auto_commit,
                     )
@@ -51,12 +73,12 @@
             }
             ClientWrapper::Http(client) => {
                 client
-                    .poll_messages(
+                    .poll_messages_with_strategy_for(
                         stream_id,
                         topic_id,
                         partition_id,
                         consumer,
-                        strategy,
+                        strategy_for,
                         count,
                         auto_commit,
                     )
@@ -64,12 +86,12 @@
             }
             ClientWrapper::Tcp(client) => {
                 client
-                    .poll_messages(
+                    .poll_messages_with_strategy_for(
                         stream_id,
                         topic_id,
                         partition_id,
                         consumer,
-                        strategy,
+                        strategy_for,
                         count,
                         auto_commit,
                     )
@@ -77,12 +99,12 @@
             }
             ClientWrapper::Quic(client) => {
                 client
-                    .poll_messages(
+                    .poll_messages_with_strategy_for(
                         stream_id,
                         topic_id,
                         partition_id,
                         consumer,
-                        strategy,
+                        strategy_for,
                         count,
                         auto_commit,
                     )
@@ -90,12 +112,12 @@
             }
             ClientWrapper::WebSocket(client) => {
                 client
-                    .poll_messages(
+                    .poll_messages_with_strategy_for(
                         stream_id,
                         topic_id,
                         partition_id,
                         consumer,
-                        strategy,
+                        strategy_for,
                         count,
                         auto_commit,
                     )
diff --git a/core/sdk/src/clients/binary_message.rs b/core/sdk/src/clients/binary_message.rs
index db722a0..e83321b 100644
--- a/core/sdk/src/clients/binary_message.rs
+++ b/core/sdk/src/clients/binary_message.rs
@@ -37,6 +37,28 @@
         count: u32,
         auto_commit: bool,
     ) -> Result<PolledMessages, IggyError> {
+        self.poll_messages_with_strategy_for(
+            stream_id,
+            topic_id,
+            partition_id,
+            consumer,
+            &|_: u32| *strategy,
+            count,
+            auto_commit,
+        )
+        .await
+    }
+
+    async fn poll_messages_with_strategy_for(
+        &self,
+        stream_id: &Identifier,
+        topic_id: &Identifier,
+        partition_id: Option<u32>,
+        consumer: &Consumer,
+        strategy_for: &(dyn Fn(u32) -> PollingStrategy + Send + Sync),
+        count: u32,
+        auto_commit: bool,
+    ) -> Result<PolledMessages, IggyError> {
         if count == 0 {
             return Err(IggyError::InvalidMessagesCount);
         }
@@ -45,12 +67,12 @@
             .client
             .read()
             .await
-            .poll_messages(
+            .poll_messages_with_strategy_for(
                 stream_id,
                 topic_id,
                 partition_id,
                 consumer,
-                strategy,
+                strategy_for,
                 count,
                 auto_commit,
             )
diff --git a/core/sdk/src/clients/consumer.rs b/core/sdk/src/clients/consumer.rs
index a7bf549..1f38f2c 100644
--- a/core/sdk/src/clients/consumer.rs
+++ b/core/sdk/src/clients/consumer.rs
@@ -26,8 +26,8 @@
 };
 use iggy_common::{
     Consumer, ConsumerKind, DiagnosticEvent, EncryptorKind, IdKind, Identifier, IggyDuration,
-    IggyError, IggyMessage, IggyTimestamp, NonZeroIggyDuration, PolledMessages, PollingKind,
-    PollingStrategy,
+    IggyError, IggyMessage, IggyTimestamp, NO_ASSIGNED_PARTITION, NonZeroIggyDuration,
+    PolledMessages, PollingKind, PollingStrategy,
 };
 use std::collections::VecDeque;
 use std::fmt::{self, Debug, Formatter};
@@ -400,7 +400,8 @@
 ///         .store_offset(received.message.header.offset, Some(received.partition_id))
 ///         .await?;
 /// }
-/// // No shutdown() here: it would commit the reading position, failed message included.
+///
+/// consumer.shutdown().await?;
 /// # Ok(())
 /// # }
 /// ```
@@ -420,13 +421,14 @@
 ///   creating the group first if [`create_consumer_group_if_not_exists()`] is set (the default).
 ///   It rejoins on its own after a reconnect and whenever the server reports that its membership
 ///   is gone.
-/// - Until the join has succeeded the consumer does not poll. It re-checks every
-///   [`polling_retry_interval()`] and polls once joined.
+/// - Such a member does not poll until it is in the group. A join that fails is yielded as
+///   `Some(Err(..))` after [`polling_retry_interval()`], and the next call tries again.
 /// - Partitions are redistributed whenever members join or leave, so a member reads different
 ///   partitions over time and messages from several partitions interleave in its stream.
 /// - More members than partitions leaves the surplus members without partitions. Such a member
-///   still polls, one group sync round trip per attempt, so give it a [`poll_interval()`]. The
-///   partition count of the topic is the ceiling on how far one group can be scaled out.
+///   keeps asking the server for an assignment, parking for [`polling_retry_interval()`] between
+///   attempts. The partition count of the topic is the ceiling on how far one group can be
+///   scaled out.
 /// - The group shares one set of stored offsets, kept under the group name. Thus,
 ///   a partition taken over by another member continues where the previous one
 ///   committed.
@@ -452,13 +454,12 @@
 /// | [`PollingStrategy::timestamp()`] | the first message at or after a given point in time |
 ///
 /// Only [`PollingStrategy::next()`] consults the offset stored on the server.
-/// Use this if you want to resume where a previous run stopped. The other four are starting points
-/// for the first request only. From the second request onwards, the consumer asks for whatever
-/// follows the last message it handed over, and it keeps one such continuation point for all
-/// partitions. That makes them fit for standalone consumers only: a group member polls a different
-/// partition on every request, so a continuation point taken from one partition is applied to the
-/// next, where it skips or repeats messages, and under the default [`auto_commit()`] the skipped
-/// range is committed as read.
+/// Use this if you want to resume where a previous run stopped. The other four are the starting
+/// point for the first request to each partition. From then on the consumer asks that partition
+/// for whatever follows the last message it handed over from it, and it keeps that position per
+/// partition. A partition that moves to another member and back therefore continues from this
+/// consumer's own position, not from where the other member got to, so a rebalance can repeat
+/// messages under these strategies, which ignore the group's stored offsets by definition.
 ///
 /// [`StreamExt::next`] yields `None` once [`shutdown()`](Self::shutdown) has been called, and never
 /// otherwise: not when the topic is empty and not while the client is disconnected. A request that
@@ -495,10 +496,10 @@
 ///
 /// | Setting | Commits |
 /// | --- | --- |
-/// | [`AutoCommit::Disabled`] | never on its own, decide manually with [`store_offset()`](Self::store_offset). [`shutdown()`](Self::shutdown) still commits the reading position |
+/// | [`AutoCommit::Disabled`] | never, not even on [`shutdown()`](Self::shutdown). Commit with [`store_offset()`](Self::store_offset) |
 /// | [`AutoCommit::Interval`] | on every tick, the reading position of every partition read so far |
 /// | [`AutoCommitWhen::PollingMessages`] | sends the commit with the poll request itself, before your code sees the batch |
-/// | [`AutoCommitWhen::ConsumingEachMessage`] | queued just before every message is handed over to the calling code, at one round trip per message and without backpressure |
+/// | [`AutoCommitWhen::ConsumingEachMessage`] | queued just before every message is handed over to the calling code. Commits queued faster than they are sent collapse into the latest one per partition |
 /// | [`AutoCommitWhen::ConsumingEveryNthMessage`] | queued just before a message whose offset divides by `n` is handed over |
 /// | [`AutoCommitWhen::ConsumingAllMessages`] | queued when the buffer of the current batch runs empty |
 /// | [`AutoCommitAfter`] variants | once the handler returned, `Ok` or `Err`, and only under [`IggyConsumerMessageExt::consume_messages`], see below |
@@ -529,9 +530,8 @@
 ///   and store the offset using [`Self::store_offset()`] after handling a message. Every other
 ///   setting except the plain [`AutoCommit::After`] variants can commit a message before your
 ///   handler is done with it, so a crash in the handler loses it. [`AutoCommit::IntervalOrAfter`]
-///   still commits on its interval tick. [`shutdown()`](Self::shutdown) commits the reading
-///   position under every setting, [`AutoCommit::Disabled`] included, so it also commits a message
-///   whose handler failed.
+///   still commits on its interval tick, and [`shutdown()`](Self::shutdown) commits the reading
+///   position under every setting but [`AutoCommit::Disabled`], a failed message included.
 ///
 /// # Options and defaults
 ///
@@ -540,18 +540,18 @@
 ///
 /// | Option | Default | Controls |
 /// | --- | --- | --- |
-/// | [`stream()`], [`topic()`], [`partition()`] | the values passed to the entry point | what is read. [`partition()`] is for standalone consumers, on a group member it pins every poll to that partition instead of the server's assignment |
+/// | [`stream()`], [`topic()`], [`partition()`] | the values passed to the entry point | what is read. [`partition()`] is for standalone consumers, a group member ignores it with a warning and reads the server's assignment |
 /// | [`batch_length()`] | 1000 | messages fetched per request |
 /// | [`poll_interval()`] | none | smallest gap between two requests |
-/// | [`polling_strategy()`] | [`PollingStrategy::next()`] | where reading starts. Anything but [`PollingStrategy::next()`] is for standalone consumers only |
-/// | [`auto_commit()`] | [`AutoCommit::IntervalOrWhen`], one second, [`AutoCommitWhen::PollingMessages`] | when offsets are committed. [`commit_failed_messages()`] is a synonym for [`AutoCommit::Disabled`] |
+/// | [`polling_strategy()`] | [`PollingStrategy::next()`] | where reading each partition starts |
+/// | [`auto_commit()`] | [`AutoCommit::IntervalOrWhen`], one second, [`AutoCommitWhen::PollingMessages`] | when offsets are committed |
 /// | [`allow_replay()`] | off | whether a message can be handed over again |
-/// | [`auto_join_consumer_group()`] | on | joining the group during [`init()`](Self::init) and after a reconnect. A group member built with [`do_not_auto_join_consumer_group()`] never polls, since polling waits for the join |
+/// | [`auto_join_consumer_group()`] | on | joining the group during [`init()`](Self::init) and again whenever the membership is lost. With [`do_not_auto_join_consumer_group()`] joining is up to the caller, and a poll without a membership fails with [`IggyError::ConsumerGroupMemberNotFound`] |
 /// | [`create_consumer_group_if_not_exists()`] | on | creating the group when it is missing |
-/// | [`polling_retry_interval()`] | one second | wait between attempts while polling is blocked |
+/// | [`polling_retry_interval()`] | one second | wait between attempts while polling is blocked or the member holds no partitions |
 /// | [`init_retries()`] | none, one second apart | retries when the stream or topic is missing at [`init()`](Self::init) |
 /// | [`offset_drain_timeout()`] | five seconds | how long [`shutdown()`](Self::shutdown) waits for pending commits |
-/// | [`encryptor()`] | inherited from the client | decrypting payloads and user headers |
+/// | [`encryptor()`] | inherited from the client | decrypting payloads and user headers, see [Encryption](#encryption) |
 ///
 /// The switches have inverse setters as well, such as [`without_poll_interval()`],
 /// [`without_encryptor()`], [`do_not_auto_join_consumer_group()`] and
@@ -565,12 +565,12 @@
 /// client share unless one of them overrides it on its builder. Without an encryptor the consumer
 /// yields payloads as stored, encrypted or not.
 ///
-/// A message that cannot be decrypted is yielded as an `Err` and the whole batch is dropped. What
-/// happens next depends on [`auto_commit()`]. Under [`AutoCommitWhen::PollingMessages`] (the
-/// default) the server committed the batch with the poll, so it is skipped for good. Under every
-/// other setting the next request fetches the same batch and fails the same way until
-/// [`store_offset()`](Self::store_offset) moves the offset past it. Pick a setting other than
-/// [`AutoCommitWhen::PollingMessages`] if a batch that fails to decrypt must not be lost silently.
+/// A message that cannot be decrypted is yielded as an `Err` and the whole batch is dropped. The
+/// next request fetches the same batch and fails the same way until
+/// [`store_offset()`](Self::store_offset) moves the offset past it. Under
+/// [`AutoCommitWhen::PollingMessages`] the server would have committed the batch with the poll and
+/// skipped it for good, so [`init()`](Self::init) rejects that setting, the default included, with
+/// [`IggyError::InvalidConfiguration`] when the consumer has an encryptor.
 ///
 /// # Concurrency
 ///
@@ -587,11 +587,11 @@
 /// # Shutting down
 ///
 /// Call [`shutdown()`](Self::shutdown) once done consuming. It drains the commit tasks, commits
-/// the reading position of every partition, [`AutoCommit::Disabled`] included, and leaves the
-/// consumer group. Dropping an `IggyConsumer` instead skips that final commit and the group
-/// leave, so the server reassigns the member's partitions only once the connection is gone.
-/// Commits already queued are still sent. Neither stops the lifecycle task, which runs until the
-/// client shuts down.
+/// the reading position of every partition unless [`auto_commit()`] is [`AutoCommit::Disabled`],
+/// leaves the consumer group and stops the background tasks. Dropping an `IggyConsumer` instead
+/// skips the final commit and the group leave, so the server reassigns the member's partitions
+/// only once the connection is gone. Commits already queued are still sent and the background
+/// tasks still stop.
 ///
 /// [`IggyClient`]: crate::prelude::IggyClient
 /// [`IggyClient::consumer()`]: crate::prelude::IggyClient::consumer
@@ -604,7 +604,6 @@
 /// [`auto_join_consumer_group()`]: crate::prelude::IggyConsumerBuilder::auto_join_consumer_group
 /// [`batch_length()`]: crate::prelude::IggyConsumerBuilder::batch_length
 /// [`build()`]: crate::prelude::IggyConsumerBuilder::build
-/// [`commit_failed_messages()`]: crate::prelude::IggyConsumerBuilder::commit_failed_messages
 /// [`create_consumer_group_if_not_exists()`]: crate::prelude::IggyConsumerBuilder::create_consumer_group_if_not_exists
 /// [`do_not_auto_join_consumer_group()`]: crate::prelude::IggyConsumerBuilder::do_not_auto_join_consumer_group
 /// [`do_not_create_consumer_group_if_not_exists()`]: crate::prelude::IggyConsumerBuilder::do_not_create_consumer_group_if_not_exists
@@ -632,6 +631,9 @@
     topic_id: Arc<Identifier>,
     partition_id: Option<u32>,
     polling_strategy: PollingStrategy,
+    /// The next offset to ask each partition for. Empty under [`PollingStrategy::next()`], which
+    /// leaves the continuation to the offset stored on the server.
+    next_offsets: Arc<DashMap<u32, u64>>,
     poll_interval_micros: u64,
     batch_length: u32,
     auto_commit: AutoCommit,
@@ -643,10 +645,14 @@
     poll_future: Option<PollMessagesFuture>,
     buffered_messages: VecDeque<IggyMessage>,
     encryptor: Option<Arc<EncryptorKind>>,
-    store_offset_sender: flume::Sender<(u32, u64)>,
+    /// The latest offset each message trigger asked to commit, per partition. The store task
+    /// drains it, so a burst of triggers costs one round trip per partition instead of one each.
+    pending_commits: Arc<DashMap<u32, u64>>,
+    store_offset_notify: Arc<Notify>,
     store_offset_task: Option<JoinHandle<()>>,
     background_commit_task: Option<JoinHandle<()>>,
     background_commit_notify: Arc<Notify>,
+    events_task: Option<JoinHandle<()>>,
     store_offset_after_each_message: bool,
     store_offset_after_all_messages: bool,
     store_after_every_nth_message: u64,
@@ -680,8 +686,15 @@
         allow_replay: bool,
         offset_drain_timeout: IggyDuration,
     ) -> Self {
-        let (store_offset_sender, _) = flume::unbounded();
         let is_consumer_group = consumer.kind == ConsumerKind::ConsumerGroup;
+        let partition_id = if is_consumer_group && partition_id.is_some() {
+            warn!(
+                "Consumer group member: {consumer_name} ignores the partition set on the builder and reads the server's assignment"
+            );
+            None
+        } else {
+            partition_id
+        };
         let consumer = Arc::new(consumer);
         let stream_id = Arc::new(stream_id);
         let topic_id = Arc::new(topic_id);
@@ -706,6 +719,7 @@
             topic_id,
             partition_id,
             polling_strategy,
+            next_offsets: Arc::new(DashMap::new()),
             poll_interval_micros: polling_interval.map_or(0, |interval| interval.as_micros()),
             state,
             current_offsets: Arc::new(DashMap::new()),
@@ -721,10 +735,12 @@
             create_consumer_group_if_not_exists,
             buffered_messages: VecDeque::new(),
             encryptor,
-            store_offset_sender,
+            pending_commits: Arc::new(DashMap::new()),
+            store_offset_notify: Arc::new(Notify::new()),
             store_offset_task: None,
             background_commit_task: None,
             background_commit_notify: Arc::new(Notify::new()),
+            events_task: None,
             store_offset_after_each_message: matches!(
                 auto_commit,
                 AutoCommit::When(AutoCommitWhen::ConsumingEachMessage)
@@ -772,11 +788,11 @@
         &self.stream_id
     }
 
-    /// Returns the partition the most recent poll response came from.
+    /// Returns the partition the most recent poll response with messages came from, or `0` before
+    /// the first one.
     ///
-    /// This is `0` before the first response, and an empty response can report `0` as well. For a
-    /// consumer group the value changes over time, as the server hands different partitions to
-    /// this member. To commit for the partition a message came from, pass
+    /// For a consumer group the value changes over time, as the server hands different partitions
+    /// to this member. To commit for the partition a message came from, pass
     /// [`ReceivedMessage::partition_id`] to [`store_offset()`](Self::store_offset) instead.
     pub fn partition_id(&self) -> u32 {
         self.state.partition_id()
@@ -874,14 +890,15 @@
     /// # Lifecycle events
     ///
     /// Calling init spawns a background task that listens for lifecycle changes ([`DiagnosticEvent`]s) of the
-    /// client connection. It runs until the client shuts down and is not stopped by
-    /// [`shutdown()`](Self::shutdown).
+    /// client connection. It runs until [`shutdown()`](Self::shutdown) or until the client shuts
+    /// down.
     /// - [`DiagnosticEvent::Connected`]: a fresh connection has not joined anything yet.
     ///   Polling resumes immediately only for a consumer that is not a group member.
-    /// - [`DiagnosticEvent::SignedIn`]: re-enables polling. A group member signing in after a
-    ///   reconnect rejoins its group first and only polls once that succeeded. A failed rejoin is
-    ///   logged and leaves polling disabled until the next reconnect or an explicit login.
-    /// - [`DiagnosticEvent::Disconnected`] and [`DiagnosticEvent::SignedOut`] disable polling.
+    /// - [`DiagnosticEvent::SignedIn`]: re-enables polling. A group member whose membership is
+    ///   gone rejoins its group on the next poll, before the request goes out. A failed rejoin is
+    ///   yielded as a poll error and tried again on the poll after.
+    /// - [`DiagnosticEvent::Disconnected`] and [`DiagnosticEvent::SignedOut`] disable polling and
+    ///   forget the group membership.
     /// - [`DiagnosticEvent::Shutdown`] disables polling and terminates the background task listening
     ///   for lifecycle changes. It does not flush in-flight commits; that only happens when
     ///   [`shutdown()`](Self::shutdown) itself is called.
@@ -896,7 +913,8 @@
     ///   [`AutoCommit::IntervalOrWhen`], [`AutoCommit::IntervalOrAfter`]). Every tick it stores the
     ///   reading position of every partition read so far.
     /// - An offset store task, always. It sends the commits queued by the [`AutoCommitWhen`] and
-    ///   [`AutoCommitAfter`] triggers one at a time and stays idle under [`AutoCommit::Disabled`].
+    ///   [`AutoCommitAfter`] triggers, keeping only the latest queued offset per partition, and
+    ///   stays idle under [`AutoCommit::Disabled`].
     ///
     /// Both skip an offset that is not ahead of this consumer's own record of what it stored
     /// ([`get_last_stored_offset()`](Self::get_last_stored_offset)). Only offset `0` is always sent.
@@ -906,6 +924,10 @@
     ///
     /// # Errors
     ///
+    /// - [`IggyError::InvalidConfiguration`] when the consumer has an encryptor and
+    ///   [`auto_commit()`](crate::prelude::IggyConsumerBuilder::auto_commit) is
+    ///   [`AutoCommitWhen::PollingMessages`], checked before anything is sent. See
+    ///   [Encryption](IggyConsumer#encryption).
     /// - [`IggyError::StreamNameNotFound`] or [`IggyError::TopicNameNotFound`] when the
     ///   stream or the topic still does not exist once the retries are exhausted.
     /// - [`IggyError::ConsumerGroupNameNotFound`] when the consumer group does not exist
@@ -922,6 +944,13 @@
         let topic_id = self.topic_id.clone();
         let consumer_name = &self.consumer_name;
 
+        if self.encryptor.is_some() && self.auto_commit_after_polling {
+            error!(
+                "Consumer: {consumer_name} has an encryptor and auto-commit on polling. That commits a batch before it is decrypted, so a batch that fails to decrypt would be lost. Pick another auto-commit setting."
+            );
+            return Err(IggyError::InvalidConfiguration);
+        }
+
         info!(
             "Initializing consumer: {consumer_name} for stream: {stream_id}, topic: {topic_id}..."
         );
@@ -993,8 +1022,10 @@
             }
         }
 
-        self.subscribe_events().await;
-        // No-op if either is_consumer_group or auto_join_consumer_group is false
+        // A retried init() after a failed join must not leave the earlier task behind.
+        if let Some(previous) = self.events_task.replace(self.subscribe_events().await) {
+            previous.abort();
+        }
         self.init_consumer_group().await?;
 
         match self.auto_commit {
@@ -1006,24 +1037,7 @@
             _ => {}
         }
 
-        let state = self.state.clone();
-        let (store_offset_sender, store_offset_receiver) = flume::unbounded();
-        self.store_offset_sender = store_offset_sender;
-
-        // Message-triggered commits from `poll_next` and `consume_messages` queue here and go out
-        // one at a time. The interval task above and the poll request's own `auto_commit` flag
-        // are the other commit paths.
-        self.store_offset_task = Some(tokio::spawn(async move {
-            while let Ok((partition_id, offset)) = store_offset_receiver.recv_async().await {
-                trace!(
-                    "Received offset to store: {offset}, partition ID: {partition_id}, stream: {}, topic: {}",
-                    state.stream_id, state.topic_id
-                );
-                _ = state
-                    .store_consumer_offset(partition_id, offset, false)
-                    .await
-            }
-        }));
+        self.store_offset_task = Some(self.store_pending_commits_in_background());
 
         self.initialized = true;
         info!(
@@ -1062,12 +1076,47 @@
         })
     }
 
+    /// Sends the commits queued by the message triggers of `poll_next` and `consume_messages`.
+    /// The interval task and the poll request's own `auto_commit` flag are the other commit paths.
+    fn store_pending_commits_in_background(&self) -> JoinHandle<()> {
+        let state = self.state.clone();
+        let pending_commits = self.pending_commits.clone();
+        let shutdown = self.shutdown.clone();
+        let notify = self.store_offset_notify.clone();
+        tokio::spawn(async move {
+            loop {
+                notify.notified().await;
+                // Keys first, so no map guard is held across a round trip. An offset queued
+                // meanwhile stays in the map and the permit its trigger leaves wakes the next turn.
+                let partitions: Vec<u32> =
+                    pending_commits.iter().map(|entry| *entry.key()).collect();
+                for partition_id in partitions {
+                    let Some((_, offset)) = pending_commits.remove(&partition_id) else {
+                        continue;
+                    };
+                    _ = state
+                        .store_consumer_offset(partition_id, offset, false)
+                        .await;
+                }
+                if shutdown.load(ORDERING) && pending_commits.is_empty() {
+                    break;
+                }
+            }
+        })
+    }
+
+    /// Queues a commit for the store task. A later offset for the same partition replaces a
+    /// queued one that has not been sent yet.
     pub(crate) fn send_store_offset(&self, partition_id: u32, offset: u64) {
-        if let Err(error) = self.store_offset_sender.send((partition_id, offset)) {
+        if !self.initialized || self.shutdown.load(ORDERING) {
             error!(
-                "Failed to send offset to store: {error}, please verify if `init()` on IggyConsumer object has been called."
+                "Offset: {offset} for partition ID: {partition_id} was not queued for storing, consumer: {} is not initialized or has been shut down.",
+                self.consumer_name
             );
+            return;
         }
+        self.pending_commits.insert(partition_id, offset);
+        self.store_offset_notify.notify_one();
     }
 
     async fn init_consumer_group(&self) -> Result<(), IggyError> {
@@ -1098,7 +1147,9 @@
         .await
     }
 
-    async fn subscribe_events(&self) {
+    /// Keeps the polling flags in step with the connection. Joining the group again after a
+    /// reconnect is left to the poll path, which retries it and reports a failure as a poll error.
+    async fn subscribe_events(&self) -> JoinHandle<()> {
         trace!("Subscribing to diagnostic events");
         let mut receiver;
         {
@@ -1107,17 +1158,8 @@
         }
 
         let is_consumer_group = self.is_consumer_group;
-        let can_join_consumer_group = is_consumer_group && self.auto_join_consumer_group;
-        let client = self.client.clone();
-        let create_consumer_group_if_not_exists = self.create_consumer_group_if_not_exists;
-        let stream_id = self.stream_id.clone();
-        let topic_id = self.topic_id.clone();
-        let consumer = self.consumer.clone();
-        let consumer_name = self.consumer_name.clone();
         let can_poll = self.can_poll.clone();
         let joined_consumer_group = self.joined_consumer_group.clone();
-        let mut reconnected = false;
-        let mut disconnected = false;
 
         tokio::spawn(async move {
             while let Some(event) = receiver.next().await {
@@ -1129,69 +1171,19 @@
                         can_poll.store(false, ORDERING);
                         break;
                     }
-
                     DiagnosticEvent::Connected => {
                         trace!("Connected to the server");
                         joined_consumer_group.store(false, ORDERING);
                         if !is_consumer_group {
                             can_poll.store(true, ORDERING);
                         }
-                        if disconnected {
-                            reconnected = true;
-                            disconnected = false;
-                        }
                     }
                     DiagnosticEvent::Disconnected => {
-                        disconnected = true;
-                        reconnected = false;
                         joined_consumer_group.store(false, ORDERING);
                         can_poll.store(false, ORDERING);
                         warn!("Disconnected from the server");
                     }
                     DiagnosticEvent::SignedIn => {
-                        if !is_consumer_group {
-                            can_poll.store(true, ORDERING);
-                            continue;
-                        }
-
-                        if !can_join_consumer_group {
-                            can_poll.store(true, ORDERING);
-                            trace!("Auto join consumer group is disabled");
-                            continue;
-                        }
-
-                        if !reconnected {
-                            can_poll.store(true, ORDERING);
-                            continue;
-                        }
-
-                        if joined_consumer_group.load(ORDERING) {
-                            can_poll.store(true, ORDERING);
-                            continue;
-                        }
-
-                        info!(
-                            "Rejoining consumer group: {consumer_name} for stream: {stream_id}, topic: {topic_id}..."
-                        );
-                        if let Err(error) = Self::initialize_consumer_group(
-                            client.clone(),
-                            create_consumer_group_if_not_exists,
-                            stream_id.clone(),
-                            topic_id.clone(),
-                            consumer.clone(),
-                            &consumer_name,
-                            joined_consumer_group.clone(),
-                        )
-                        .await
-                        {
-                            error!(
-                                "Failed to join consumer group: {consumer_name} for stream: {stream_id}, topic: {topic_id}. {error}"
-                            );
-                            continue;
-                        }
-                        info!(
-                            "Rejoined consumer group: {consumer_name} for stream: {stream_id}, topic: {topic_id}"
-                        );
                         can_poll.store(true, ORDERING);
                     }
                     DiagnosticEvent::SignedOut => {
@@ -1200,7 +1192,7 @@
                     }
                 }
             }
-        });
+        })
     }
 
     fn create_poll_messages_future(
@@ -1211,10 +1203,10 @@
         let partition_id = self.partition_id;
         let consumer = self.consumer.clone();
         let polling_strategy = self.polling_strategy;
+        let next_offsets = self.next_offsets.clone();
         let client = self.client.clone();
         let count = self.batch_length;
         let auto_commit_after_polling = self.auto_commit_after_polling;
-        let auto_commit_enabled = self.auto_commit != AutoCommit::Disabled;
         let interval = self.poll_interval_micros;
         let last_polled_at = self.last_polled_at.clone();
         let can_poll = self.can_poll.clone();
@@ -1226,39 +1218,74 @@
         let auto_join_consumer_group = self.auto_join_consumer_group;
         let create_consumer_group_if_not_exists = self.create_consumer_group_if_not_exists;
         let joined_consumer_group = self.joined_consumer_group.clone();
+        let consumer_name = self.consumer_name.clone();
 
         async move {
             if interval > 0 {
                 Self::wait_before_polling(interval, last_polled_at.load(ORDERING)).await;
             }
 
-            while !can_poll.load(ORDERING)
-                || (is_consumer_group && !joined_consumer_group.load(ORDERING))
+            while !can_poll.load(ORDERING) {
+                trace!("Cannot poll yet, waiting {retry_interval}...");
+                sleep(retry_interval.get_duration()).await;
+            }
+
+            // A member that joins on its own is in the group before it polls. One built with
+            // `do_not_auto_join_consumer_group()` polls right away and gets a missing
+            // membership reported as a poll error.
+            if is_consumer_group
+                && auto_join_consumer_group
+                && !joined_consumer_group.load(ORDERING)
+                && let Err(error) = Self::initialize_consumer_group(
+                    client.clone(),
+                    create_consumer_group_if_not_exists,
+                    stream_id.clone(),
+                    topic_id.clone(),
+                    consumer.clone(),
+                    &consumer_name,
+                    joined_consumer_group.clone(),
+                )
+                .await
             {
-                trace!(
-                    "Cannot poll yet (can_poll={}, joined_cg={}), waiting {retry_interval}...",
-                    can_poll.load(ORDERING),
-                    joined_consumer_group.load(ORDERING)
+                error!(
+                    "Failed to join consumer group: {consumer_name} for stream: {stream_id}, topic: {topic_id}. {error}"
                 );
                 sleep(retry_interval.get_duration()).await;
+                return Err(error);
             }
 
             trace!("Sending poll messages request");
             last_polled_at.store(IggyTimestamp::now().into(), ORDERING);
+            // The map guard is dropped inside `map_or`, and the only writer is `poll_next`,
+            // which runs after this future has returned, so the lookup cannot block.
+            let strategy_for = |partition: u32| {
+                next_offsets
+                    .get(&partition)
+                    .map_or(polling_strategy, |offset| PollingStrategy::offset(*offset))
+            };
             let polled_messages = client
                 .read()
                 .await
-                .poll_messages(
+                .poll_messages_with_strategy_for(
                     &stream_id,
                     &topic_id,
                     partition_id,
                     &consumer,
-                    &polling_strategy,
+                    &strategy_for,
                     count,
                     auto_commit_after_polling,
                 )
                 .await;
 
+            if let Ok(polled) = &polled_messages
+                && polled.partition_id == NO_ASSIGNED_PARTITION
+            {
+                trace!(
+                    "No partition assigned to consumer: {consumer_name}, waiting {retry_interval}..."
+                );
+                sleep(retry_interval.get_duration()).await;
+            }
+
             if let Ok(mut polled_messages) = polled_messages {
                 if polled_messages.messages.is_empty() {
                     return Ok(polled_messages);
@@ -1280,8 +1307,9 @@
                     polled_messages
                         .messages
                         .retain(|message| message.header.offset > consumed_offset);
+                    polled_messages.count = polled_messages.messages.len() as u32;
                     if polled_messages.messages.is_empty() {
-                        return Ok(PolledMessages::empty());
+                        return Ok(polled_messages);
                     }
                 }
 
@@ -1306,44 +1334,6 @@
                     "Last consumed offset: {consumed_offset}, current offset: {}, stored offset: {stored_offset}, in partition ID: {partition_id}, topic: {topic_id}, stream: {stream_id}, consumer: {consumer}",
                     polled_messages.current_offset
                 );
-
-                if !allow_replay
-                    && (has_consumed_offset && polled_messages.current_offset == consumed_offset)
-                {
-                    trace!(
-                        "No new messages to consume in partition ID: {partition_id}, topic: {topic_id}, stream: {stream_id}, consumer: {consumer}"
-                    );
-                    if auto_commit_enabled && stored_offset < consumed_offset {
-                        trace!(
-                            "Auto-committing the offset: {consumed_offset} in partition ID: {partition_id}, topic: {topic_id}, stream: {stream_id}, consumer: {consumer}"
-                        );
-                        client
-                            .read()
-                            .await
-                            .store_consumer_offset(
-                                &consumer,
-                                &stream_id,
-                                &topic_id,
-                                Some(partition_id),
-                                consumed_offset,
-                            )
-                            .await?;
-                        if let Some(stored_offset_entry) = last_stored_offset.get(&partition_id) {
-                            stored_offset_entry.store(consumed_offset, ORDERING);
-                        } else {
-                            last_stored_offset
-                                .insert(partition_id, AtomicU64::new(consumed_offset));
-                        }
-                    }
-
-                    return Ok(PolledMessages {
-                        messages: vec![],
-                        current_offset: polled_messages.current_offset,
-                        partition_id,
-                        count: 0,
-                    });
-                }
-
                 return Ok(polled_messages);
             }
 
@@ -1354,26 +1344,10 @@
                 && auto_join_consumer_group
                 && matches!(&error, IggyError::ConsumerGroupMemberNotFound(..))
             {
-                joined_consumer_group.store(false, ORDERING);
-                let consumer_name = consumer.id.as_string();
                 info!(
-                    "Consumer group membership was revoked for consumer: {consumer_name}, stream: {stream_id}, topic: {topic_id}. Rejoining..."
+                    "Consumer group membership was revoked for consumer: {consumer_name}, stream: {stream_id}, topic: {topic_id}. Rejoining on the next poll..."
                 );
-                if let Err(error) = Self::initialize_consumer_group(
-                    client,
-                    create_consumer_group_if_not_exists,
-                    stream_id,
-                    topic_id,
-                    consumer,
-                    &consumer_name,
-                    joined_consumer_group.clone(),
-                )
-                .await
-                {
-                    // Allow the next poll to retry rejoining
-                    joined_consumer_group.store(true, ORDERING);
-                    return Err(error);
-                }
+                joined_consumer_group.store(false, ORDERING);
                 return Ok(PolledMessages::empty());
             }
 
@@ -1565,14 +1539,13 @@
                 }
             }
 
-            // Popping above may have left the buffer empty.
-            // The next turn will therefore poll messages from the server.
-            // With `PollingStrategy` the user defines the starting point where to poll from.
-            // After that, each poll must read the next sequential offset. Hence, strategy is
-            // set to `PollingKind::Offset` and the next offset to read from is the last consumed message + 1.
+            // Popping above may have left the buffer empty, so the next turn polls the server.
+            // `polling_strategy` is only where reading a partition starts; from then on each
+            // poll continues after the last message handed over from that partition.
             if self.buffered_messages.is_empty() {
                 if self.polling_strategy.kind != PollingKind::Next {
-                    self.polling_strategy = PollingStrategy::offset(message.header.offset + 1);
+                    self.next_offsets
+                        .insert(partition_id, message.header.offset + 1);
                 }
 
                 if self.store_offset_after_all_messages {
@@ -1604,95 +1577,98 @@
 
         while let Some(future) = self.poll_future.as_mut() {
             match future.poll_unpin(cx) {
-                Poll::Ready(Ok(mut polled_messages)) => {
-                    let partition_id = polled_messages.partition_id;
+                Poll::Ready(Ok(polled_messages)) => {
+                    let PolledMessages {
+                        partition_id,
+                        current_offset,
+                        messages,
+                        ..
+                    } = polled_messages;
+                    let mut messages = VecDeque::from(messages);
+                    let Some(mut first) = messages.pop_front() else {
+                        self.poll_future = Some(Box::pin(self.create_poll_messages_future()));
+                        continue;
+                    };
+
+                    // Only a response that carries messages names a partition; an empty one can
+                    // carry a sentinel instead of a real id.
                     self.state
                         .current_partition_id
                         .store(partition_id, ORDERING);
-                    if polled_messages.messages.is_empty() {
-                        self.poll_future = Some(Box::pin(self.create_poll_messages_future()));
-                    } else {
-                        if let Some(ref encryptor) = self.encryptor {
-                            for message in &mut polled_messages.messages {
-                                let offset = message.header.offset;
-                                let payload = encryptor.decrypt(&message.payload);
-                                if let Err(error) = payload {
+
+                    if let Some(ref encryptor) = self.encryptor {
+                        for message in std::iter::once(&mut first).chain(messages.iter_mut()) {
+                            let offset = message.header.offset;
+                            let payload = encryptor.decrypt(&message.payload);
+                            if let Err(error) = payload {
+                                self.poll_future = None;
+                                error!(
+                                    "Failed to decrypt the message payload at offset: {offset}, partition ID: {partition_id}",
+                                );
+                                return Poll::Ready(Some(Err(error)));
+                            }
+
+                            let payload = payload.unwrap();
+                            message.payload = Bytes::from(payload);
+                            message.header.payload_length = message.payload.len() as u32;
+
+                            if let Some(ref user_headers) = message.user_headers {
+                                let decrypted_headers = encryptor.decrypt(user_headers);
+                                if let Err(error) = decrypted_headers {
                                     self.poll_future = None;
                                     error!(
-                                        "Failed to decrypt the message payload at offset: {offset}, partition ID: {partition_id}",
+                                        "Failed to decrypt the message user headers at offset: {offset}, partition ID: {partition_id}",
                                     );
                                     return Poll::Ready(Some(Err(error)));
                                 }
-
-                                let payload = payload.unwrap();
-                                message.payload = Bytes::from(payload);
-                                message.header.payload_length = message.payload.len() as u32;
-
-                                if let Some(ref user_headers) = message.user_headers {
-                                    let decrypted_headers = encryptor.decrypt(user_headers);
-                                    if let Err(error) = decrypted_headers {
-                                        self.poll_future = None;
-                                        error!(
-                                            "Failed to decrypt the message user headers at offset: {offset}, partition ID: {partition_id}",
-                                        );
-                                        return Poll::Ready(Some(Err(error)));
-                                    }
-                                    let decrypted_headers = decrypted_headers.unwrap();
-                                    message.header.user_headers_length =
-                                        decrypted_headers.len() as u32;
-                                    message.user_headers = Some(Bytes::from(decrypted_headers));
-                                }
+                                let decrypted_headers = decrypted_headers.unwrap();
+                                message.header.user_headers_length = decrypted_headers.len() as u32;
+                                message.user_headers = Some(Bytes::from(decrypted_headers));
                             }
                         }
-
-                        if let Some(current_offset_entry) = self.current_offsets.get(&partition_id)
-                        {
-                            current_offset_entry.store(polled_messages.current_offset, ORDERING);
-                        } else {
-                            self.current_offsets.insert(
-                                partition_id,
-                                AtomicU64::new(polled_messages.current_offset),
-                            );
-                        }
-
-                        let message = polled_messages.messages.remove(0);
-                        self.buffered_messages.extend(polled_messages.messages);
-
-                        if self.polling_strategy.kind != PollingKind::Next {
-                            self.polling_strategy =
-                                PollingStrategy::offset(message.header.offset + 1);
-                        }
-
-                        if let Some(last_consumed_offset_entry) =
-                            self.state.last_consumed_offsets.get(&partition_id)
-                        {
-                            last_consumed_offset_entry.store(message.header.offset, ORDERING);
-                        } else {
-                            self.state
-                                .last_consumed_offsets
-                                .insert(partition_id, AtomicU64::new(message.header.offset));
-                        }
-
-                        if (self.store_after_every_nth_message > 0
-                            && message.header.offset % self.store_after_every_nth_message == 0)
-                            || self.store_offset_after_each_message
-                            || (self.store_offset_after_all_messages
-                                && self.buffered_messages.is_empty())
-                        {
-                            self.send_store_offset(
-                                polled_messages.partition_id,
-                                message.header.offset,
-                            );
-                        }
-
-                        // Drop future since it is [invalid after being ready](https://doc.rust-lang.org/std/future/trait.Future.html#panics)
-                        self.poll_future = None;
-                        return Poll::Ready(Some(Ok(ReceivedMessage::new(
-                            message,
-                            polled_messages.current_offset,
-                            polled_messages.partition_id,
-                        ))));
                     }
+
+                    if let Some(current_offset_entry) = self.current_offsets.get(&partition_id) {
+                        current_offset_entry.store(current_offset, ORDERING);
+                    } else {
+                        self.current_offsets
+                            .insert(partition_id, AtomicU64::new(current_offset));
+                    }
+
+                    // A poll is only sent once the buffer has run empty, so nothing is overwritten.
+                    self.buffered_messages = messages;
+
+                    if self.polling_strategy.kind != PollingKind::Next {
+                        self.next_offsets
+                            .insert(partition_id, first.header.offset + 1);
+                    }
+
+                    if let Some(last_consumed_offset_entry) =
+                        self.state.last_consumed_offsets.get(&partition_id)
+                    {
+                        last_consumed_offset_entry.store(first.header.offset, ORDERING);
+                    } else {
+                        self.state
+                            .last_consumed_offsets
+                            .insert(partition_id, AtomicU64::new(first.header.offset));
+                    }
+
+                    if (self.store_after_every_nth_message > 0
+                        && first.header.offset % self.store_after_every_nth_message == 0)
+                        || self.store_offset_after_each_message
+                        || (self.store_offset_after_all_messages
+                            && self.buffered_messages.is_empty())
+                    {
+                        self.send_store_offset(partition_id, first.header.offset);
+                    }
+
+                    // Drop future since it is [invalid after being ready](https://doc.rust-lang.org/std/future/trait.Future.html#panics)
+                    self.poll_future = None;
+                    return Poll::Ready(Some(Ok(ReceivedMessage::new(
+                        first,
+                        current_offset,
+                        partition_id,
+                    ))));
                 }
                 Poll::Ready(Err(err)) => {
                     self.poll_future = None;
@@ -1715,14 +1691,15 @@
     ///   commits in flight. The consumer waits for `offset_drain_timeout` on each in turn before
     ///   forcing it to abort.
     /// - commit the reading position of every partition where it is ahead of this consumer's own
-    ///   record of what it stored, under every [`AutoCommit`] setting, [`AutoCommit::Disabled`]
-    ///   included. Under auto-commit-on-poll (the default) the poll already committed the whole
-    ///   batch, so this store moves the server offset back to the last message handed over, and
-    ///   the next run resumes right after it instead of after the last batch fetched.
+    ///   record of what it stored, unless [`auto_commit()`] is [`AutoCommit::Disabled`]. Under
+    ///   auto-commit-on-poll (the default) the poll already committed the whole batch, so this
+    ///   store moves the server offset back to the last message handed over, and the next run
+    ///   resumes right after it instead of after the last batch fetched.
     /// - leave the consumer group, if this consumer is a group member. This lets the server give its partitions to
     ///   the remaining members immediately instead of waiting for the connection to time out.
+    /// - stop the task watching the connection lifecycle.
     ///
-    /// The lifecycle event task is not stopped. It runs until the client shuts down.
+    /// [`auto_commit()`]: crate::prelude::IggyConsumerBuilder::auto_commit
     ///
     /// # Errors
     ///
@@ -1757,17 +1734,8 @@
             );
         }
 
-        // Drop the sending end of the store offset task to end the `recv_async()` loop in `init()`.
-        // Offsets in queue will still be committed. This prevents loading additional offsets into a channel
-        // that is not read anymore.
-        // Replace with a new (hanging) channel, since `store_offset_sender` is not optional.
-        let (closed_sender, _) = flume::bounded(0);
-        drop(std::mem::replace(
-            &mut self.store_offset_sender,
-            closed_sender,
-        ));
-
-        // This task never sleeps, so no need to notify.
+        // Wakes the store task, which sends what is still queued and exits on the shutdown flag.
+        self.store_offset_notify.notify_one();
         if let Some(mut task) = self.store_offset_task.take()
             && time::timeout(self.offset_drain_timeout.get_duration(), &mut task)
                 .await
@@ -1780,18 +1748,19 @@
             );
         }
 
-        for (partition_id, consumed_offset) in self.state.last_consumed_offsets() {
-            let stored_offset = self.state.get_last_stored_offset(partition_id).unwrap_or(0);
-
-            if consumed_offset > stored_offset {
-                trace!(
-                    "Flushing final offset: {consumed_offset} for partition: {partition_id}, stream: {}, topic: {}",
-                    self.stream_id, self.topic_id
-                );
-                let _ = self
-                    .state
-                    .store_consumer_offset(partition_id, consumed_offset, self.allow_replay)
-                    .await;
+        if self.auto_commit != AutoCommit::Disabled {
+            for (partition_id, consumed_offset) in self.state.last_consumed_offsets() {
+                let stored_offset = self.state.get_last_stored_offset(partition_id).unwrap_or(0);
+                if consumed_offset > stored_offset {
+                    trace!(
+                        "Flushing final offset: {consumed_offset} for partition: {partition_id}, stream: {}, topic: {}",
+                        self.stream_id, self.topic_id
+                    );
+                    let _ = self
+                        .state
+                        .store_consumer_offset(partition_id, consumed_offset, self.allow_replay)
+                        .await;
+                }
             }
         }
 
@@ -1820,19 +1789,26 @@
             }
         }
 
+        if let Some(task) = self.events_task.take() {
+            task.abort();
+        }
+
         info!("Consumer: {} has been shut down.", self.consumer_name);
         Ok(())
     }
 }
 
-/// Wakes the interval commit task so it exits. Commits already queued still go out.
-///
-/// Nothing is flushed and the consumer group is not left. Await [`IggyConsumer::shutdown`] first,
-/// see [Shutting down](IggyConsumer#shutting-down).
+/// Stops the background tasks. Commits already queued still go out, nothing else is flushed and
+/// the consumer group is not left. Await [`IggyConsumer::shutdown`] first, see
+/// [Shutting down](IggyConsumer#shutting-down).
 impl Drop for IggyConsumer {
     fn drop(&mut self) {
         self.shutdown.store(true, ORDERING);
         self.background_commit_notify.notify_one();
+        self.store_offset_notify.notify_one();
+        if let Some(task) = self.events_task.take() {
+            task.abort();
+        }
         trace!(
             "Consumer {} has been dropped, shutdown signal sent",
             self.consumer_name
@@ -1846,9 +1822,14 @@
     use crate::client_wrappers::client_wrapper::ClientWrapper;
     use crate::clients::consumer_builder::IggyConsumerBuilder;
     use crate::tcp::tcp_client::TcpClient;
+    use iggy_common::Aes256GcmEncryptor;
     use iggy_common::locking::IggyRwLockFn;
     use std::str::FromStr;
     use std::task::Waker;
+    use tokio::time::timeout;
+
+    const POLL_RETRY_INTERVAL: Duration = Duration::from_millis(10);
+    const POLL_TIMEOUT: Duration = Duration::from_secs(2);
 
     fn builder_for(consumer: Consumer) -> IggyConsumerBuilder {
         IggyConsumerBuilder::new(
@@ -1921,6 +1902,165 @@
         assert!(consumer.poll_future.is_none());
     }
 
+    #[test]
+    fn group_member_should_ignore_the_partition_set_on_the_builder() {
+        let consumer = builder_for(Consumer::group(Identifier::numeric(1).unwrap()))
+            .partition(Some(1))
+            .build();
+
+        assert_eq!(consumer.partition_id, None);
+    }
+
+    #[test]
+    fn standalone_consumer_should_keep_the_partition_set_on_the_builder() {
+        let consumer = builder().partition(Some(1)).build();
+
+        assert_eq!(consumer.partition_id, Some(1));
+    }
+
+    fn message_at(offset: u64) -> IggyMessage {
+        let mut message = IggyMessage::from_str("payload").unwrap();
+        message.header.offset = offset;
+        message
+    }
+
+    /// Hands over `messages` as one buffered batch read from `partition_id`.
+    fn hand_over_batch(consumer: &mut IggyConsumer, partition_id: u32, messages: Vec<IggyMessage>) {
+        consumer
+            .state
+            .current_partition_id
+            .store(partition_id, ORDERING);
+        consumer.buffered_messages = VecDeque::from(messages);
+        let mut context = Context::from_waker(Waker::noop());
+        while !consumer.buffered_messages.is_empty() {
+            assert!(matches!(
+                Pin::new(&mut *consumer).poll_next(&mut context),
+                Poll::Ready(Some(Ok(_)))
+            ));
+        }
+    }
+
+    fn next_offset(consumer: &IggyConsumer, partition_id: u32) -> Option<u64> {
+        consumer
+            .next_offsets
+            .get(&partition_id)
+            .map(|offset| *offset)
+    }
+
+    #[test]
+    fn group_member_should_continue_each_partition_after_its_last_message() {
+        let mut consumer = builder_for(Consumer::group(Identifier::numeric(1).unwrap()))
+            .polling_strategy(PollingStrategy::first())
+            .auto_commit(AutoCommit::Disabled)
+            .build();
+
+        hand_over_batch(&mut consumer, 3, vec![message_at(10), message_at(11)]);
+        assert_eq!(next_offset(&consumer, 3), Some(12));
+        assert_eq!(consumer.polling_strategy, PollingStrategy::first());
+
+        hand_over_batch(&mut consumer, 4, vec![message_at(7)]);
+        assert_eq!(next_offset(&consumer, 4), Some(8));
+        assert_eq!(next_offset(&consumer, 3), Some(12));
+    }
+
+    #[test]
+    fn next_strategy_should_leave_the_continuation_to_the_server() {
+        let mut consumer = builder_for(Consumer::group(Identifier::numeric(1).unwrap()))
+            .auto_commit(AutoCommit::Disabled)
+            .build();
+
+        hand_over_batch(&mut consumer, 3, vec![message_at(10), message_at(11)]);
+
+        assert!(consumer.next_offsets.is_empty());
+    }
+
+    /// Polls once as a group member on a client that is not connected. The outcome must be an
+    /// error, never an endless wait for a join.
+    async fn poll_once_as_group_member(
+        builder: IggyConsumerBuilder,
+    ) -> Option<Result<ReceivedMessage, IggyError>> {
+        let mut consumer = builder
+            .polling_retry_interval(NonZeroIggyDuration::new(POLL_RETRY_INTERVAL).unwrap())
+            .build();
+        timeout(POLL_TIMEOUT, consumer.next())
+            .await
+            .expect("a group member must poll or report an error instead of waiting for a join")
+    }
+
+    #[tokio::test]
+    async fn group_member_without_auto_join_should_poll_instead_of_waiting_for_the_join() {
+        let builder = builder_for(Consumer::group(Identifier::numeric(1).unwrap()))
+            .do_not_auto_join_consumer_group();
+
+        assert!(matches!(
+            poll_once_as_group_member(builder).await,
+            Some(Err(_))
+        ));
+    }
+
+    #[tokio::test]
+    async fn group_member_should_report_a_failed_join_as_a_poll_error() {
+        let builder = builder_for(Consumer::group(Identifier::numeric(1).unwrap()))
+            .auto_join_consumer_group();
+
+        assert!(matches!(
+            poll_once_as_group_member(builder).await,
+            Some(Err(_))
+        ));
+    }
+
+    #[tokio::test]
+    async fn init_should_reject_an_encryptor_with_auto_commit_on_polling() {
+        let encryptor = Arc::new(EncryptorKind::Aes256Gcm(
+            Aes256GcmEncryptor::new(&[1; 32]).unwrap(),
+        ));
+        for auto_commit in [
+            AutoCommit::When(AutoCommitWhen::PollingMessages),
+            AutoCommit::IntervalOrWhen(
+                NonZeroIggyDuration::ONE_SECOND,
+                AutoCommitWhen::PollingMessages,
+            ),
+        ] {
+            let mut consumer = builder()
+                .encryptor(encryptor.clone())
+                .auto_commit(auto_commit)
+                .build();
+
+            assert!(
+                matches!(consumer.init().await, Err(IggyError::InvalidConfiguration)),
+                "{auto_commit:?} must be rejected with an encryptor"
+            );
+        }
+
+        let mut consumer = builder()
+            .encryptor(encryptor)
+            .auto_commit(AutoCommit::When(AutoCommitWhen::ConsumingEachMessage))
+            .build();
+
+        assert!(!matches!(
+            consumer.init().await,
+            Err(IggyError::InvalidConfiguration)
+        ));
+    }
+
+    #[test]
+    fn send_store_offset_should_keep_the_latest_offset_per_partition() {
+        let mut consumer = builder().build();
+        consumer.initialized = true;
+
+        consumer.send_store_offset(1, 5);
+        consumer.send_store_offset(1, 7);
+        consumer.send_store_offset(2, 3);
+
+        let mut queued: Vec<(u32, u64)> = consumer
+            .pending_commits
+            .iter()
+            .map(|entry| (*entry.key(), *entry.value()))
+            .collect();
+        queued.sort_unstable();
+        assert_eq!(queued, vec![(1, 7), (2, 3)]);
+    }
+
     #[tokio::test]
     async fn should_accept_every_auto_commit_mode() {
         for auto_commit in [
diff --git a/core/sdk/src/clients/consumer_builder.rs b/core/sdk/src/clients/consumer_builder.rs
index 87b427d..91cb9ac 100644
--- a/core/sdk/src/clients/consumer_builder.rs
+++ b/core/sdk/src/clients/consumer_builder.rs
@@ -93,8 +93,8 @@
     }
 
     /// Sets the partition to read. `None` lets a consumer group read its assigned partitions and
-    /// makes the server read partition `0` for a standalone consumer. `Some(n)` on a group member
-    /// pins every poll to that partition instead of the assignment.
+    /// makes the server read partition `0` for a standalone consumer. `Some(n)` is for standalone
+    /// consumers. A group member ignores it with a warning and reads its assignment.
     pub fn partition(self, partition: Option<u32>) -> Self {
         Self { partition, ..self }
     }
@@ -123,15 +123,8 @@
         }
     }
 
-    /// Same as [`auto_commit`](Self::auto_commit) with [`AutoCommit::Disabled`].
-    pub fn commit_failed_messages(self) -> Self {
-        Self {
-            auto_commit: AutoCommit::Disabled,
-            ..self
-        }
-    }
-
-    /// Automatically joins the consumer group if the consumer is a part of a consumer group.
+    /// Joins the consumer group during `init()` and again after the membership was lost, for
+    /// example after a reconnect. On by default.
     pub fn auto_join_consumer_group(self) -> Self {
         Self {
             auto_join_consumer_group: true,
@@ -139,7 +132,9 @@
         }
     }
 
-    /// Does not automatically join the consumer group if the consumer is a part of a consumer group.
+    /// Leaves joining the consumer group to the caller. The member polls as soon as `init()`
+    /// returns, and a poll without a membership fails with
+    /// [`IggyError::ConsumerGroupMemberNotFound`](iggy_common::IggyError::ConsumerGroupMemberNotFound).
     pub fn do_not_auto_join_consumer_group(self) -> Self {
         Self {
             auto_join_consumer_group: false,
@@ -195,7 +190,9 @@
         }
     }
 
-    /// Sets the polling retry interval in case of server disconnection.
+    /// Sets how long a poll waits before the next attempt while it is blocked: after a
+    /// disconnect, after a failed group join, or while the group member holds no partitions.
+    /// One second by default.
     pub fn polling_retry_interval(self, interval: NonZeroIggyDuration) -> Self {
         Self {
             polling_retry_interval: interval,
@@ -233,7 +230,7 @@
 
     /// Builds the consumer.
     ///
-    /// Note: After building the consumer, `init()` must be invoked before producing messages.
+    /// Note: After building the consumer, `init()` must be invoked before consuming messages.
     pub fn build(self) -> IggyConsumer {
         IggyConsumer::new(
             self.client,
diff --git a/core/sdk/src/prelude.rs b/core/sdk/src/prelude.rs
index 81f7d8c..72e2503 100644
--- a/core/sdk/src/prelude.rs
+++ b/core/sdk/src/prelude.rs
@@ -78,6 +78,7 @@
     IGGY_MESSAGE_HEADERS_LENGTH_OFFSET_RANGE, IGGY_MESSAGE_ID_OFFSET_RANGE,
     IGGY_MESSAGE_OFFSET_OFFSET_RANGE, IGGY_MESSAGE_ORIGIN_TIMESTAMP_OFFSET_RANGE,
     IGGY_MESSAGE_PAYLOAD_LENGTH_OFFSET_RANGE, IGGY_MESSAGE_TIMESTAMP_OFFSET_RANGE, INDEX_SIZE,
-    MAX_PAYLOAD_SIZE, MAX_USER_HEADERS_SIZE, SEC_IN_MICRO,
+    MAX_PAYLOAD_SIZE, MAX_USER_HEADERS_SIZE, NO_ASSIGNED_PARTITION,
+    RESYNC_REQUIRED_PARTITION_SENTINEL, SEC_IN_MICRO,
     defaults::{DEFAULT_ROOT_PASSWORD, DEFAULT_ROOT_USER_ID, DEFAULT_ROOT_USERNAME},
 };
diff --git a/core/sdk/src/quic/quic_client.rs b/core/sdk/src/quic/quic_client.rs
index e69e94a..5dcadc1 100644
--- a/core/sdk/src/quic/quic_client.rs
+++ b/core/sdk/src/quic/quic_client.rs
@@ -400,6 +400,12 @@
         }
 
         consensus_session.bind(session);
+        drop(consensus_session);
+        // Every fresh client identity passes through here, including one that
+        // replaces a session the transport never reset: a connection lost
+        // mid-request can leave the old session in place until this sign-in
+        // re-mints it.
+        self.consumer_group_state.clear_session_scoped();
         Ok(())
     }
 
@@ -408,6 +414,7 @@
             .consensus_session
             .lock()
             .expect("consensus session mutex poisoned") = ConsensusSession::new();
+        self.consumer_group_state.clear_session_scoped();
         Ok(())
     }
 
diff --git a/core/sdk/src/tcp/tcp_client.rs b/core/sdk/src/tcp/tcp_client.rs
index 23f70f7..1546708 100644
--- a/core/sdk/src/tcp/tcp_client.rs
+++ b/core/sdk/src/tcp/tcp_client.rs
@@ -344,6 +344,12 @@
         }
 
         consensus_session.bind(session);
+        drop(consensus_session);
+        // Every fresh client identity passes through here, including one that
+        // replaces a session the transport never reset: a connection lost
+        // mid-request can leave the old session in place until this sign-in
+        // re-mints it.
+        self.consumer_group_state.clear_session_scoped();
         Ok(())
     }
 
@@ -352,6 +358,7 @@
             .consensus_session
             .lock()
             .expect("consensus session mutex poisoned") = ConsensusSession::new();
+        self.consumer_group_state.clear_session_scoped();
         Ok(())
     }
 
diff --git a/core/sdk/src/websocket/websocket_client.rs b/core/sdk/src/websocket/websocket_client.rs
index f40bc91..2a22683 100644
--- a/core/sdk/src/websocket/websocket_client.rs
+++ b/core/sdk/src/websocket/websocket_client.rs
@@ -395,6 +395,12 @@
         }
 
         consensus_session.bind(session);
+        drop(consensus_session);
+        // Every fresh client identity passes through here, including one that
+        // replaces a session the transport never reset: a connection lost
+        // mid-request can leave the old session in place until this sign-in
+        // re-mints it.
+        self.consumer_group_state.clear_session_scoped();
         Ok(())
     }
 
@@ -403,6 +409,7 @@
             .consensus_session
             .lock()
             .expect("consensus session mutex poisoned") = ConsensusSession::new();
+        self.consumer_group_state.clear_session_scoped();
         Ok(())
     }
 
diff --git a/foreign/php/README.md b/foreign/php/README.md
index 795b3be..15007bb 100644
--- a/foreign/php/README.md
+++ b/foreign/php/README.md
@@ -117,7 +117,9 @@
 }
 ```
 
-Consumer group callbacks require a finite message limit:
+Consumer group callbacks require a finite message limit. The partition id
+argument is ignored for a consumer group, since the member reads the partitions
+the server assigns to it:
 
 ```php
 <?php
@@ -126,7 +128,7 @@
     'php-consumer',
     $stream,
     $topic,
-    $partitionId,
+    null,
     \Iggy\PollingStrategy::next(),
     10,
     \Iggy\AutoCommit::disabled(),
diff --git a/foreign/php/iggy-php.stubs.php b/foreign/php/iggy-php.stubs.php
index 211e6f1..be8437c 100644
--- a/foreign/php/iggy-php.stubs.php
+++ b/foreign/php/iggy-php.stubs.php
@@ -94,6 +94,9 @@
         /**
          * Creates and initializes a consumer group consumer.
          *
+         * `$partition_id` is ignored for a consumer group: the member reads the partitions
+         * the server assigns to it.
+         *
          * @param string $name
          * @param string $stream
          * @param string $topic
diff --git a/foreign/php/src/client.rs b/foreign/php/src/client.rs
index fee9752..23e5a61 100644
--- a/foreign/php/src/client.rs
+++ b/foreign/php/src/client.rs
@@ -293,6 +293,9 @@
     }
 
     /// Creates and initializes a consumer group consumer.
+    ///
+    /// `$partition_id` is ignored for a consumer group: the member reads the partitions
+    /// the server assigns to it.
     #[allow(clippy::too_many_arguments)]
     #[php(defaults(
         create_consumer_group_if_not_exists = true,
diff --git a/foreign/python/apache_iggy.pyi b/foreign/python/apache_iggy.pyi
index 6c3715b..c0f603e 100644
--- a/foreign/python/apache_iggy.pyi
+++ b/foreign/python/apache_iggy.pyi
@@ -1377,6 +1377,8 @@
     ) -> collections.abc.Awaitable[IggyConsumer]:
         r"""
         Creates a new consumer group consumer.
+        `partition_id` is ignored for a consumer group: the member reads the partitions
+        the server assigns to it.
         Returns the consumer or a RuntimeError on failure. Raises `ValueError` if
         `poll_interval`, `polling_retry_interval`, `init_retry_interval` or an
         `AutoCommit` interval is negative, or if any of those except `poll_interval`
diff --git a/foreign/python/src/client.rs b/foreign/python/src/client.rs
index f669a87..10c5c77 100644
--- a/foreign/python/src/client.rs
+++ b/foreign/python/src/client.rs
@@ -1087,6 +1087,8 @@
     }
 
     /// Creates a new consumer group consumer.
+    /// `partition_id` is ignored for a consumer group: the member reads the partitions
+    /// the server assigns to it.
     /// Returns the consumer or a RuntimeError on failure. Raises `ValueError` if
     /// `poll_interval`, `polling_retry_interval`, `init_retry_interval` or an
     /// `AutoCommit` interval is negative, or if any of those except `poll_interval`