diff --git a/lib/ProducerImpl.cc b/lib/ProducerImpl.cc index f1da6f59..083230f3 100644 --- a/lib/ProducerImpl.cc +++ b/lib/ProducerImpl.cc @@ -609,6 +609,10 @@ void ProducerImpl::sendAsyncWithStatsUpdate(const Message& msg, SendCallback&& c failures.complete(); } } else { + if (batchMessageContainer_) { + batchMessageAndSend().complete(); + } + const bool sendChunks = (totalChunks > 1); ChunkMessageIdListPtr chunkMessageIdList; if (sendChunks) { diff --git a/tests/ClientDeduplicationTest.cc b/tests/ClientDeduplicationTest.cc index 90991b99..0b652716 100644 --- a/tests/ClientDeduplicationTest.cc +++ b/tests/ClientDeduplicationTest.cc @@ -19,8 +19,12 @@ #include #include +#include +#include +#include #include #include +#include #include "HttpHelper.h" @@ -149,3 +153,80 @@ TEST(ClientDeduplicationTest, testProducerDeduplication) { client.close(); } + +class ClientDeduplicationBatchingTest : public ::testing::TestWithParam { +}; + +TEST_P(ClientDeduplicationBatchingTest, testBatchedMessagesBeforeDelayedMessage) { + Client client(serviceUrl); + + std::string topicName = "persistent://public/dedup-3/testBatchedMessagesBeforeDelayedMessage-" + + std::to_string(static_cast(GetParam())) + "-" + std::to_string(time(NULL)); + + std::string url = adminUrl + "admin/v2/namespaces/public/dedup-3"; + int res = makePutRequest(url, R"({"replication_clusters": ["standalone"]})"); + ASSERT_TRUE(res == 204 || res == 409); + + url = adminUrl + "admin/v2/namespaces/public/dedup-3/permissions/anonymous"; + res = makePostRequest(url, R"(["produce","consume"])"); + ASSERT_TRUE(res == 204 || res == 409); + + url = adminUrl + "admin/v2/namespaces/public/dedup-3/deduplication"; + res = makePostRequest(url, "true"); + ASSERT_TRUE(res == 204 || res == 409); + + // Ensure dedup status was refreshed + std::this_thread::sleep_for(std::chrono::seconds(1)); + + Producer producer; + ProducerConfiguration producerConf; + producerConf.setBatchingEnabled(true); + producerConf.setBatchingType(GetParam()); + // Long enough for the batch to be still pending when the delayed message is sent + producerConf.setBatchingMaxPublishDelayMs(3000); + ASSERT_EQ(client.createProducer(topicName, producerConf, producer), ResultOk); + + constexpr int numBatchedMessages = 10; + std::vector>> promises; + for (int i = 0; i < numBatchedMessages; i++) { + auto promise = std::make_shared>(); + promises.emplace_back(promise); + producer.sendAsync(MessageBuilder() + .setContent("batched-" + std::to_string(i)) + .setPartitionKey("key-" + std::to_string(i % 2)) + .build(), + [promise](Result result, const MessageId&) { promise->set_value(result); }); + } + + // A delayed message is never batched, and it takes the next sequence id + ASSERT_EQ( + producer.send( + MessageBuilder().setContent("delayed").setDeliverAfter(std::chrono::milliseconds(1)).build()), + ResultOk); + + for (auto& promise : promises) { + ASSERT_EQ(promise->get_future().get(), ResultOk); + } + + // Every message must have been stored, none dropped as a duplicate + Reader reader; + ASSERT_EQ(client.createReader(topicName, MessageId::earliest(), {}, reader), ResultOk); + + std::set received; + Message msg; + while (reader.readNext(msg, 3000) == ResultOk) { + received.emplace(msg.getDataAsString()); + } + + ASSERT_EQ(received.size(), numBatchedMessages + 1); + for (int i = 0; i < numBatchedMessages; i++) { + ASSERT_EQ(received.count("batched-" + std::to_string(i)), 1u); + } + ASSERT_EQ(received.count("delayed"), 1u); + + client.close(); +} + +INSTANTIATE_TEST_SUITE_P(Pulsar, ClientDeduplicationBatchingTest, + ::testing::Values(ProducerConfiguration::DefaultBatching, + ProducerConfiguration::KeyBasedBatching));