diff --git a/include/pulsar/Reader.h b/include/pulsar/Reader.h index 233da4ff..4c124c9d 100644 --- a/include/pulsar/Reader.h +++ b/include/pulsar/Reader.h @@ -29,6 +29,7 @@ class PulsarFriend; class ReaderImpl; typedef std::function HasMessageAvailableCallback; +typedef std::function ReadNextCallback; /** * A Reader can be used to scan through all the messages currently available in a topic. @@ -68,6 +69,13 @@ class PULSAR_PUBLIC Reader { */ Result readNext(Message& msg, int timeoutMs); + /** + * Read asynchronously the next message in the topic. + * + * @param callback + */ + void readNextAsync(ReadNextCallback callback); + /** * Close the reader and stop the broker to push more messages * diff --git a/include/pulsar/ReaderConfiguration.h b/include/pulsar/ReaderConfiguration.h index 9ae8e1c4..4f6464f7 100644 --- a/include/pulsar/ReaderConfiguration.h +++ b/include/pulsar/ReaderConfiguration.h @@ -162,6 +162,7 @@ class PULSAR_PUBLIC ReaderConfiguration { * Set the internal subscription name. * * @param internal subscriptionName + * Default value is reader-{random string}. */ void setInternalSubscriptionName(std::string internalSubscriptionName); diff --git a/lib/ConsumerImpl.cc b/lib/ConsumerImpl.cc index ccd082b6..e7f6cf47 100644 --- a/lib/ConsumerImpl.cc +++ b/lib/ConsumerImpl.cc @@ -841,7 +841,7 @@ Result ConsumerImpl::receive(Message& msg) { return res; } -void ConsumerImpl::receiveAsync(ReceiveCallback& callback) { +void ConsumerImpl::receiveAsync(ReceiveCallback callback) { Message msg; // fail the callback if consumer is closing or closed diff --git a/lib/ConsumerImpl.h b/lib/ConsumerImpl.h index 4832f4e6..6e1b8ff3 100644 --- a/lib/ConsumerImpl.h +++ b/lib/ConsumerImpl.h @@ -94,7 +94,7 @@ class ConsumerImpl : public ConsumerImplBase { const std::string& getTopic() const override; Result receive(Message& msg) override; Result receive(Message& msg, int timeout) override; - void receiveAsync(ReceiveCallback& callback) override; + void receiveAsync(ReceiveCallback callback) override; void unsubscribeAsync(ResultCallback callback) override; void acknowledgeAsync(const MessageId& msgId, ResultCallback callback) override; void acknowledgeAsync(const MessageIdList& messageIdList, ResultCallback callback) override; diff --git a/lib/ConsumerImplBase.h b/lib/ConsumerImplBase.h index 5bc7e1b8..9cf63a38 100644 --- a/lib/ConsumerImplBase.h +++ b/lib/ConsumerImplBase.h @@ -51,7 +51,7 @@ class ConsumerImplBase : public HandlerBase, public std::enable_shared_from_this virtual const std::string& getSubscriptionName() const = 0; virtual Result receive(Message& msg) = 0; virtual Result receive(Message& msg, int timeout) = 0; - virtual void receiveAsync(ReceiveCallback& callback) = 0; + virtual void receiveAsync(ReceiveCallback callback) = 0; void batchReceiveAsync(BatchReceiveCallback callback); virtual void unsubscribeAsync(ResultCallback callback) = 0; virtual void acknowledgeAsync(const MessageId& msgId, ResultCallback callback) = 0; diff --git a/lib/MultiTopicsConsumerImpl.cc b/lib/MultiTopicsConsumerImpl.cc index a0135667..c7a656c0 100644 --- a/lib/MultiTopicsConsumerImpl.cc +++ b/lib/MultiTopicsConsumerImpl.cc @@ -583,7 +583,7 @@ Result MultiTopicsConsumerImpl::receive(Message& msg, int timeout) { } } -void MultiTopicsConsumerImpl::receiveAsync(ReceiveCallback& callback) { +void MultiTopicsConsumerImpl::receiveAsync(ReceiveCallback callback) { Message msg; // fail the callback if consumer is closing or closed diff --git a/lib/MultiTopicsConsumerImpl.h b/lib/MultiTopicsConsumerImpl.h index da42b748..50cdecf3 100644 --- a/lib/MultiTopicsConsumerImpl.h +++ b/lib/MultiTopicsConsumerImpl.h @@ -65,7 +65,7 @@ class MultiTopicsConsumerImpl : public ConsumerImplBase { const std::string& getTopic() const override; Result receive(Message& msg) override; Result receive(Message& msg, int timeout) override; - void receiveAsync(ReceiveCallback& callback) override; + void receiveAsync(ReceiveCallback callback) override; void unsubscribeAsync(ResultCallback callback) override; void acknowledgeAsync(const MessageId& msgId, ResultCallback callback) override; void acknowledgeAsync(const MessageIdList& messageIdList, ResultCallback callback) override; diff --git a/lib/Reader.cc b/lib/Reader.cc index 261c0fac..c02fb2e5 100644 --- a/lib/Reader.cc +++ b/lib/Reader.cc @@ -49,6 +49,14 @@ Result Reader::readNext(Message& msg, int timeoutMs) { return impl_->readNext(msg, timeoutMs); } +void Reader::readNextAsync(ReadNextCallback callback) { + if (!impl_) { + return callback(ResultConsumerNotInitialized, {}); + } + + impl_->readNextAsync(callback); +} + Result Reader::close() { Promise promise; closeAsync(WaitForCallback(promise)); diff --git a/lib/ReaderImpl.cc b/lib/ReaderImpl.cc index a80e2e50..da1d95ec 100644 --- a/lib/ReaderImpl.cc +++ b/lib/ReaderImpl.cc @@ -111,6 +111,14 @@ Result ReaderImpl::readNext(Message& msg, int timeoutMs) { return res; } +void ReaderImpl::readNextAsync(ReceiveCallback callback) { + auto self = shared_from_this(); + consumer_->receiveAsync([self, callback](Result result, const Message& message) { + self->acknowledgeIfNecessary(result, message); + callback(result, message); + }); +} + void ReaderImpl::messageListener(Consumer consumer, const Message& msg) { readerListener_(Reader(shared_from_this()), msg); acknowledgeIfNecessary(ResultOk, msg); diff --git a/lib/ReaderImpl.h b/lib/ReaderImpl.h index ed16c7d3..e216241d 100644 --- a/lib/ReaderImpl.h +++ b/lib/ReaderImpl.h @@ -67,6 +67,7 @@ class PULSAR_PUBLIC ReaderImpl : public std::enable_shared_from_this Result readNext(Message& msg); Result readNext(Message& msg, int timeoutMs); + void readNextAsync(ReceiveCallback callback); void closeAsync(ResultCallback callback); diff --git a/tests/ReaderTest.cc b/tests/ReaderTest.cc index b88da998..eefe1bcf 100644 --- a/tests/ReaderTest.cc +++ b/tests/ReaderTest.cc @@ -25,6 +25,7 @@ #include "HttpHelper.h" #include "PulsarFriend.h" +#include "WaitUtils.h" #include "lib/ClientConnection.h" #include "lib/Latch.h" #include "lib/LogUtils.h" @@ -68,6 +69,50 @@ TEST(ReaderTest, testSimpleReader) { client.close(); } +TEST(ReaderTest, testAsyncRead) { + Client client(serviceUrl); + + std::string topicName = "persistent://public/default/test-simple-reader" + std::to_string(time(nullptr)); + + ReaderConfiguration readerConf; + Reader reader; + ASSERT_EQ(ResultOk, client.createReader(topicName, MessageId::earliest(), readerConf, reader)); + + Producer producer; + ASSERT_EQ(ResultOk, client.createProducer(topicName, producer)); + + for (int i = 0; i < 10; i++) { + std::string content = "my-message-" + std::to_string(i); + Message msg = MessageBuilder().setContent(content).build(); + ASSERT_EQ(ResultOk, producer.send(msg)); + } + + for (int i = 0; i < 10; i++) { + reader.readNextAsync([i](Result result, const Message& msg) { + ASSERT_EQ(ResultOk, result); + std::string content = msg.getDataAsString(); + std::string expected = "my-message-" + std::to_string(i); + ASSERT_EQ(expected, content); + }); + } + + waitUntil( + std::chrono::seconds(5), + [&]() { + bool hasMsg; + reader.hasMessageAvailable(hasMsg); + return !hasMsg; + }, + 1000); + bool hasMsg; + reader.hasMessageAvailable(hasMsg); + ASSERT_FALSE(hasMsg); + + producer.close(); + reader.close(); + client.close(); +} + TEST(ReaderTest, testReaderAfterMessagesWerePublished) { Client client(serviceUrl);