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 b1b911ffb fix(sdk): stream should terminate after consumer shutdown
(#3964)
b1b911ffb is described below
commit b1b911ffb894fb9d29381fd7fb8b2ef5933302bb
Author: haubur <[email protected]>
AuthorDate: Tue Sep 1 08:25:19 2026 +0200
fix(sdk): stream should terminate after consumer shutdown (#3964)
---
core/sdk/src/clients/consumer.rs | 68 ++++++++++++++++++++++++++++++++++++++--
1 file changed, 66 insertions(+), 2 deletions(-)
diff --git a/core/sdk/src/clients/consumer.rs b/core/sdk/src/clients/consumer.rs
index b9d23a684..a4477e6c8 100644
--- a/core/sdk/src/clients/consumer.rs
+++ b/core/sdk/src/clients/consumer.rs
@@ -1055,6 +1055,10 @@ impl Stream for IggyConsumer {
type Item = Result<ReceivedMessage, IggyError>;
fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) ->
Poll<Option<Self::Item>> {
+ if self.shutdown.load(ORDERING) {
+ return Poll::Ready(None);
+ }
+
let partition_id = self.state.partition_id();
if let Some(message) = self.buffered_messages.pop_front() {
{
@@ -1314,12 +1318,14 @@ mod tests {
use crate::clients::consumer_builder::IggyConsumerBuilder;
use crate::tcp::tcp_client::TcpClient;
use iggy_common::locking::IggyRwLockFn;
+ use std::str::FromStr;
+ use std::task::Waker;
- fn builder() -> IggyConsumerBuilder {
+ fn builder_for(consumer: Consumer) -> IggyConsumerBuilder {
IggyConsumerBuilder::new(
IggyRwLock::new(ClientWrapper::Tcp(TcpClient::default())),
"consumer".to_owned(),
- Consumer::new(Identifier::numeric(1).unwrap()),
+ consumer,
Identifier::numeric(1).unwrap(),
Identifier::numeric(1).unwrap(),
None,
@@ -1328,6 +1334,64 @@ mod tests {
)
}
+ fn builder() -> IggyConsumerBuilder {
+ builder_for(Consumer::new(Identifier::numeric(1).unwrap()))
+ }
+
+ async fn assert_stream_terminates_after_shutdown(consumer: Consumer) {
+ let mut consumer = builder_for(consumer)
+ .partition(Some(1))
+ .batch_length(1)
+ .auto_commit(AutoCommit::Disabled)
+ .build();
+ consumer.buffered_messages.extend([
+ IggyMessage::from_str("a").unwrap(),
+ IggyMessage::from_str("b").unwrap(),
+ ]);
+ let mut context = Context::from_waker(Waker::noop());
+
+ assert!(matches!(
+ Pin::new(&mut consumer).poll_next(&mut context),
+ Poll::Ready(Some(Ok(_)))
+ ));
+
+ consumer.shutdown().await.unwrap();
+
+ assert_eq!(consumer.buffered_messages.len(), 1);
+ assert!(matches!(
+ Pin::new(&mut consumer).poll_next(&mut context),
+ Poll::Ready(None)
+ ));
+ }
+
+ #[tokio::test]
+ async fn
standalone_consumer_should_stop_yielding_messages_after_shutdown() {
+
assert_stream_terminates_after_shutdown(Consumer::new(Identifier::numeric(1).unwrap()))
+ .await;
+ }
+
+ #[tokio::test]
+ async fn consumer_group_should_stop_yielding_messages_after_shutdown() {
+
assert_stream_terminates_after_shutdown(Consumer::group(Identifier::numeric(1).unwrap()))
+ .await;
+ }
+
+ #[tokio::test]
+ async fn consumer_group_should_not_create_poll_future_after_shutdown() {
+ let mut consumer =
builder_for(Consumer::group(Identifier::numeric(1).unwrap()))
+ .auto_commit(AutoCommit::Disabled)
+ .build();
+ let mut context = Context::from_waker(Waker::noop());
+
+ consumer.shutdown().await.unwrap();
+
+ assert!(matches!(
+ Pin::new(&mut consumer).poll_next(&mut context),
+ Poll::Ready(None)
+ ));
+ assert!(consumer.poll_future.is_none());
+ }
+
#[tokio::test]
async fn should_accept_every_auto_commit_mode() {
for auto_commit in [