|
18 | 18 | */ |
19 | 19 | #include <gtest/gtest.h> |
20 | 20 | #include <pulsar/Client.h> |
| 21 | +#include <pulsar/ClientConfiguration.h> |
21 | 22 | #include <pulsar/Reader.h> |
22 | 23 | #include <time.h> |
23 | 24 |
|
@@ -982,5 +983,49 @@ TEST_F(ReaderSeekTest, testSeekInclusiveChunkMessage) { |
982 | 983 | assertStartMessageId(false, secondMsgId); |
983 | 984 | } |
984 | 985 |
|
| 986 | +// Regression test for segfault when Reader is used with messageListenerThreads=0. |
| 987 | +// Verifies ExecutorServiceProvider(0) does not cause undefined behavior and |
| 988 | +// ConsumerImpl::messageReceived does not dereference null listenerExecutor_. |
| 989 | +TEST(ReaderTest, testReaderWithZeroMessageListenerThreads) { |
| 990 | + ClientConfiguration clientConf; |
| 991 | + clientConf.setMessageListenerThreads(0); |
| 992 | + Client client(serviceUrl, clientConf); |
| 993 | + |
| 994 | + const std::string topicName = "testReaderWithZeroMessageListenerThreads-" + std::to_string(time(nullptr)); |
| 995 | + |
| 996 | + Producer producer; |
| 997 | + ASSERT_EQ(ResultOk, client.createProducer(topicName, producer)); |
| 998 | + |
| 999 | + ReaderConfiguration readerConf; |
| 1000 | + Reader reader; |
| 1001 | + ASSERT_EQ(ResultOk, client.createReader(topicName, MessageId::earliest(), readerConf, reader)); |
| 1002 | + |
| 1003 | + constexpr int numMessages = 5; |
| 1004 | + for (int i = 0; i < numMessages; i++) { |
| 1005 | + Message msg = MessageBuilder().setContent("msg-" + std::to_string(i)).build(); |
| 1006 | + ASSERT_EQ(ResultOk, producer.send(msg)); |
| 1007 | + } |
| 1008 | + |
| 1009 | + int received = 0; |
| 1010 | + for (int i = 0; i < numMessages + 2; i++) { |
| 1011 | + bool hasMessageAvailable = false; |
| 1012 | + ASSERT_EQ(ResultOk, reader.hasMessageAvailable(hasMessageAvailable)); |
| 1013 | + if (!hasMessageAvailable) { |
| 1014 | + break; |
| 1015 | + } |
| 1016 | + Message msg; |
| 1017 | + Result res = reader.readNext(msg, 3000); |
| 1018 | + ASSERT_EQ(ResultOk, res) << "readNext failed at iteration " << i; |
| 1019 | + std::string content = msg.getDataAsString(); |
| 1020 | + EXPECT_EQ("msg-" + std::to_string(received), content); |
| 1021 | + ++received; |
| 1022 | + } |
| 1023 | + EXPECT_EQ(received, numMessages); |
| 1024 | + |
| 1025 | + producer.close(); |
| 1026 | + reader.close(); |
| 1027 | + client.close(); |
| 1028 | +} |
| 1029 | + |
985 | 1030 | INSTANTIATE_TEST_SUITE_P(Pulsar, ReaderTest, ::testing::Values(true, false)); |
986 | 1031 | INSTANTIATE_TEST_SUITE_P(Pulsar, ReaderSeekTest, ::testing::Values(true, false)); |
0 commit comments