This is an automated email from the ASF dual-hosted git repository.
hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new ca0b578a2 fix(sdk): correct consumer group polling and commits in
IggyConsumer (#4034)
ca0b578a2 is described below
commit ca0b578a225bf4b9169e4ae78844f3f0423a4e99
Author: Hubert Gruszecki <[email protected]>
AuthorDate: Wed Sep 2 17:37:18 2026 +0200
fix(sdk): correct consumer group polling and commits in IggyConsumer (#4034)
---
core/bench/src/actors/consumer/client/low_level.rs | 28 +-
core/common/src/consumer_group_client_state.rs | 38 +
core/common/src/lib.rs | 7 +
core/common/src/traits/binary_impls/messages.rs | 43 +-
core/common/src/traits/message_client.rs | 41 +-
core/common/src/types/message/polled_messages.rs | 6 +-
core/integration/tests/sdk/consumer_group.rs | 83 +++
.../tests/sdk/consumer_group_membership.rs | 276 ++++++++
core/integration/tests/sdk/consumer_shutdown.rs | 112 +++
core/integration/tests/sdk/mod.rs | 1 +
.../src/client_wrappers/binary_message_client.rs | 42 +-
core/sdk/src/clients/binary_message.rs | 26 +-
core/sdk/src/clients/consumer.rs | 762 ++++++++++++---------
core/sdk/src/clients/consumer_builder.rs | 25 +-
core/sdk/src/prelude.rs | 3 +-
core/sdk/src/quic/quic_client.rs | 7 +
core/sdk/src/tcp/tcp_client.rs | 7 +
core/sdk/src/websocket/websocket_client.rs | 7 +
foreign/php/README.md | 6 +-
foreign/php/iggy-php.stubs.php | 3 +
foreign/php/src/client.rs | 3 +
foreign/python/apache_iggy.pyi | 2 +
foreign/python/src/client.rs | 2 +
23 files changed, 1172 insertions(+), 358 deletions(-)
diff --git a/core/bench/src/actors/consumer/client/low_level.rs
b/core/bench/src/actors/consumer/client/low_level.rs
index 7f4774895..964231f62 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::benchmarks::common::create_consumer;
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 @@ pub struct LowLevelConsumerClient {
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 @@ impl LowLevelConsumerClient {
partition_id: None,
polling_strategy: PollingStrategy::next(),
auto_commit: true,
- offset: 0,
+ next_offsets: HashMap::new(),
}
}
}
@@ -62,14 +65,21 @@ impl ConsumerClient for LowLevelConsumerClient {
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 @@ impl ConsumerClient for LowLevelConsumerClient {
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 @@ impl BenchmarkInit for LowLevelConsumerClient {
)
.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 1e694cc70..d4331cb12 100644
--- a/core/common/src/consumer_group_client_state.rs
+++ b/core/common/src/consumer_group_client_state.rs
@@ -185,6 +185,24 @@ impl ConsumerGroupClientState {
.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 @@ mod tests {
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 89a9fc74b..e74eae292 100644
--- a/core/common/src/lib.rs
+++ b/core/common/src/lib.rs
@@ -41,6 +41,13 @@ pub use
consumer_group_client_state::ConsumerGroupClientState;
/// 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 bf720a1fe..d8b12410d 100644
--- a/core/common/src/traits/binary_impls/messages.rs
+++ b/core/common/src/traits/binary_impls/messages.rs
@@ -166,14 +166,15 @@ async fn resolve_partitioning<B: BinaryClient>(
}
/// 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 @@ async fn poll_group_messages<B: BinaryClient>(
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,
};
@@ -283,6 +289,28 @@ impl<B: BinaryClient> MessageClient for B {
strategy: &PollingStrategy,
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
@@ -294,18 +322,19 @@ impl<B: BinaryClient> MessageClient for B {
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 0c5fa11e4..18c6d2294 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 @@ pub trait MessageClient {
/// 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 @@ pub trait MessageClient {
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 af23d4ec7..d974fa479 100644
--- a/core/common/src/types/message/polled_messages.rs
+++ b/core/common/src/types/message/polled_messages.rs
@@ -29,7 +29,11 @@ use tracing::error;
/// - `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 b670de3ee..f233560da 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_GROUP_NAME: &str =
"consumer-group-rejoin-group";
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 @@ async fn
consumer_group_retries_rejoin_after_failure(harness: &TestHarness) {
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 429aea649..a8e7f59e7 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 @@ async fn
given_group_member_holds_no_partitions_when_group_deleted_should_surfac
}
}
+// 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 @@ fn assert_rejected(result: Result<(), IggyError>,
expected_code: u32, context: &
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 000000000..172a56985
--- /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 065b257a7..0cec5f837 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 eb3eed94f..3fb72563a 100644
--- a/core/sdk/src/client_wrappers/binary_message_client.rs
+++ b/core/sdk/src/client_wrappers/binary_message_client.rs
@@ -34,16 +34,38 @@ impl MessageClient for ClientWrapper {
strategy: &PollingStrategy,
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 @@ impl MessageClient for ClientWrapper {
}
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 @@ impl MessageClient for ClientWrapper {
}
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 @@ impl MessageClient for ClientWrapper {
}
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 @@ impl MessageClient for ClientWrapper {
}
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 db722a006..e83321bbd 100644
--- a/core/sdk/src/clients/binary_message.rs
+++ b/core/sdk/src/clients/binary_message.rs
@@ -36,6 +36,28 @@ impl MessageClient for IggyClient {
strategy: &PollingStrategy,
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 @@ impl MessageClient for IggyClient {
.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 a7bf54983..1f38f2c76 100644
--- a/core/sdk/src/clients/consumer.rs
+++ b/core/sdk/src/clients/consumer.rs
@@ -26,8 +26,8 @@ use iggy_common::{
};
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 @@ unsafe impl Sync for IggyConsumer {}
/// .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 @@ unsafe impl Sync for IggyConsumer {}
/// 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 @@ unsafe impl Sync for IggyConsumer {}
/// | [`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 @@ unsafe impl Sync for IggyConsumer {}
///
/// | 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 @@ unsafe impl Sync for IggyConsumer {}
/// 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 @@ unsafe impl Sync for IggyConsumer {}
///
/// | 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 @@ unsafe impl Sync for IggyConsumer {}
/// 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 @@ unsafe impl Sync for IggyConsumer {}
/// # 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 @@ unsafe impl Sync for IggyConsumer {}
/// [`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 @@ pub struct IggyConsumer {
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 @@ pub struct IggyConsumer {
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 @@ impl IggyConsumer {
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 @@ impl IggyConsumer {
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 @@ impl IggyConsumer {
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 @@ impl IggyConsumer {
&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 @@ impl IggyConsumer {
/// # 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 @@ impl IggyConsumer {
/// [`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 @@ impl IggyConsumer {
///
/// # 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 @@ impl IggyConsumer {
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 @@ impl IggyConsumer {
}
}
- 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 @@ impl IggyConsumer {
_ => {}
}
- 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 @@ impl IggyConsumer {
})
}
+ /// 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 @@ impl IggyConsumer {
.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 @@ impl IggyConsumer {
}
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 @@ impl IggyConsumer {
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 @@ impl IggyConsumer {
}
}
}
- });
+ })
}
fn create_poll_messages_future(
@@ -1211,10 +1203,10 @@ impl IggyConsumer {
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 @@ impl IggyConsumer {
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 @@ impl IggyConsumer {
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 @@ impl IggyConsumer {
"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 @@ impl IggyConsumer {
&& 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 @@ impl Stream for IggyConsumer {
}
}
- // 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 @@ impl Stream for IggyConsumer {
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 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));
+ }
- if self.polling_strategy.kind != PollingKind::Next {
- self.polling_strategy =
- PollingStrategy::offset(message.header.offset
+ 1);
- }
+ // A poll is only sent once the buffer has run empty, so
nothing is overwritten.
+ self.buffered_messages = messages;
- 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.polling_strategy.kind != PollingKind::Next {
+ self.next_offsets
+ .insert(partition_id, first.header.offset + 1);
+ }
- 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,
- );
- }
+ 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));
+ }
- // 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 (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 @@ impl IggyConsumer {
/// 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 @@ impl IggyConsumer {
);
}
- // 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 @@ impl IggyConsumer {
);
}
- 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 @@ impl IggyConsumer {
}
}
+ 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 @@ mod tests {
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 @@ mod tests {
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 87b427d83..91cb9ac83 100644
--- a/core/sdk/src/clients/consumer_builder.rs
+++ b/core/sdk/src/clients/consumer_builder.rs
@@ -93,8 +93,8 @@ impl IggyConsumerBuilder {
}
/// 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 @@ impl IggyConsumerBuilder {
}
}
- /// 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 @@ impl IggyConsumerBuilder {
}
}
- /// 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 @@ impl IggyConsumerBuilder {
}
}
- /// 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 @@ impl IggyConsumerBuilder {
/// 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 81f7d8c16..72e2503d1 100644
--- a/core/sdk/src/prelude.rs
+++ b/core/sdk/src/prelude.rs
@@ -78,6 +78,7 @@ pub use iggy_common::{
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 e69e94af6..5dcadc1d6 100644
--- a/core/sdk/src/quic/quic_client.rs
+++ b/core/sdk/src/quic/quic_client.rs
@@ -400,6 +400,12 @@ impl iggy_common::VsrSessionControl for QuicClient {
}
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 @@ impl iggy_common::VsrSessionControl for QuicClient {
.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 23f70f791..154670833 100644
--- a/core/sdk/src/tcp/tcp_client.rs
+++ b/core/sdk/src/tcp/tcp_client.rs
@@ -344,6 +344,12 @@ impl iggy_common::VsrSessionControl for TcpClient {
}
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 @@ impl iggy_common::VsrSessionControl for TcpClient {
.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 f40bc9178..2a22683e0 100644
--- a/core/sdk/src/websocket/websocket_client.rs
+++ b/core/sdk/src/websocket/websocket_client.rs
@@ -395,6 +395,12 @@ impl iggy_common::VsrSessionControl for WebSocketClient {
}
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 @@ impl iggy_common::VsrSessionControl for WebSocketClient {
.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 795b3bebb..15007bb90 100644
--- a/foreign/php/README.md
+++ b/foreign/php/README.md
@@ -117,7 +117,9 @@ foreach ($messages as $message) {
}
```
-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 @@ $consumer = $client->consumerGroup(
'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 211e6f120..be8437cb0 100644
--- a/foreign/php/iggy-php.stubs.php
+++ b/foreign/php/iggy-php.stubs.php
@@ -94,6 +94,9 @@ namespace Iggy {
/**
* 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 fee9752ed..23e5a6124 100644
--- a/foreign/php/src/client.rs
+++ b/foreign/php/src/client.rs
@@ -293,6 +293,9 @@ impl IggyClient {
}
/// 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 6c3715b54..c0f603e58 100644
--- a/foreign/python/apache_iggy.pyi
+++ b/foreign/python/apache_iggy.pyi
@@ -1377,6 +1377,8 @@ class IggyClient:
) -> 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 f669a8748..10c5c776b 100644
--- a/foreign/python/src/client.rs
+++ b/foreign/python/src/client.rs
@@ -1087,6 +1087,8 @@ impl IggyClient {
}
/// 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`