diff --git a/lib/ConsumerImpl.cc b/lib/ConsumerImpl.cc index d932c165..f75d852c 100644 --- a/lib/ConsumerImpl.cc +++ b/lib/ConsumerImpl.cc @@ -1186,6 +1186,7 @@ void ConsumerImpl::closeAsync(ResultCallback originalCallback) { if (ackGroupingTrackerPtr_) { ackGroupingTrackerPtr_->close(); } + negativeAcksTracker_.close(); ClientConnectionPtr cnx = getCnx().lock(); if (!cnx) { @@ -1219,6 +1220,7 @@ void ConsumerImpl::shutdown() { if (client) { client->cleanupConsumer(this); } + negativeAcksTracker_.close(); cancelTimers(); consumerCreatedPromise_.setFailed(ResultAlreadyClosed); failPendingReceiveCallback(); diff --git a/lib/ConsumerImpl.h b/lib/ConsumerImpl.h index 6437472d..08a2b1fd 100644 --- a/lib/ConsumerImpl.h +++ b/lib/ConsumerImpl.h @@ -332,6 +332,7 @@ class ConsumerImpl : public ConsumerImplBase { FRIEND_TEST(ConsumerTest, testPartitionedConsumerUnAckedMessageRedelivery); FRIEND_TEST(ConsumerTest, testMultiTopicsConsumerUnAckedMessageRedelivery); FRIEND_TEST(ConsumerTest, testBatchUnAckedMessageTracker); + FRIEND_TEST(ConsumerTest, testNegativeAcksTrackerClose); FRIEND_TEST(DeadLetterQueueTest, testAutoSetDLQTopicName); }; diff --git a/lib/NegativeAcksTracker.cc b/lib/NegativeAcksTracker.cc index 9dcca20f..78078087 100644 --- a/lib/NegativeAcksTracker.cc +++ b/lib/NegativeAcksTracker.cc @@ -105,6 +105,8 @@ void NegativeAcksTracker::close() { boost::system::error_code ec; timer_->cancel(ec); } + timer_ = nullptr; + nackedMessages_.clear(); } void NegativeAcksTracker::setEnabledForTesting(bool enabled) { diff --git a/lib/NegativeAcksTracker.h b/lib/NegativeAcksTracker.h index c5a945b2..f8b334b3 100644 --- a/lib/NegativeAcksTracker.h +++ b/lib/NegativeAcksTracker.h @@ -28,6 +28,8 @@ #include #include +#include "TestUtil.h" + namespace pulsar { class ConsumerImpl; @@ -66,6 +68,8 @@ class NegativeAcksTracker { ExecutorServicePtr executor_; DeadlineTimerPtr timer_; bool enabledForTesting_; // to be able to test deterministically + + FRIEND_TEST(ConsumerTest, testNegativeAcksTrackerClose); }; } // namespace pulsar diff --git a/tests/ConsumerTest.cc b/tests/ConsumerTest.cc index f3c8abbe..ee9bb3d5 100644 --- a/tests/ConsumerTest.cc +++ b/tests/ConsumerTest.cc @@ -958,6 +958,36 @@ TEST_P(ConsumerSeekTest, testSeekForMessageId) { producer.close(); } +TEST(ConsumerTest, testNegativeAcksTrackerClose) { + Client client(lookupUrl); + auto topicName = "testNegativeAcksTrackerClose"; + + ConsumerConfiguration consumerConfig; + consumerConfig.setNegativeAckRedeliveryDelayMs(100); + Consumer consumer; + client.subscribe(topicName, "test-sub", consumerConfig, consumer); + + Producer producer; + client.createProducer(topicName, producer); + + for (int i = 0; i < 10; ++i) { + producer.send(MessageBuilder().setContent(std::to_string(i)).build()); + } + + Message msg; + PulsarFriend::setNegativeAckEnabled(consumer, false); + for (int i = 0; i < 10; ++i) { + consumer.receive(msg); + consumer.negativeAcknowledge(msg); + } + + consumer.close(); + auto consumerImplPtr = PulsarFriend::getConsumerImplPtr(consumer); + ASSERT_TRUE(consumerImplPtr->negativeAcksTracker_.nackedMessages_.empty()); + + client.close(); +} + INSTANTIATE_TEST_CASE_P(Pulsar, ConsumerSeekTest, ::testing::Values(true, false)); } // namespace pulsar