Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions include/pulsar/Reader.h
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,7 @@ class PulsarFriend;
class ReaderImpl;

typedef std::function<void(Result result, bool hasMessageAvailable)> HasMessageAvailableCallback;
typedef std::function<void(Result result, const Message& message)> ReadNextCallback;

/**
* A Reader can be used to scan through all the messages currently available in a topic.
Expand Down Expand Up @@ -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
*
Expand Down
1 change: 1 addition & 0 deletions include/pulsar/ReaderConfiguration.h
Original file line number Diff line number Diff line change
Expand Up @@ -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);

Expand Down
2 changes: 1 addition & 1 deletion lib/ConsumerImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion lib/ConsumerImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
2 changes: 1 addition & 1 deletion lib/ConsumerImplBase.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
2 changes: 1 addition & 1 deletion lib/MultiTopicsConsumerImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
2 changes: 1 addition & 1 deletion lib/MultiTopicsConsumerImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
8 changes: 8 additions & 0 deletions lib/Reader.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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<bool, Result> promise;
closeAsync(WaitForCallback(promise));
Expand Down
8 changes: 8 additions & 0 deletions lib/ReaderImpl.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand Down
1 change: 1 addition & 0 deletions lib/ReaderImpl.h
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ class PULSAR_PUBLIC ReaderImpl : public std::enable_shared_from_this<ReaderImpl>

Result readNext(Message& msg);
Result readNext(Message& msg, int timeoutMs);
void readNextAsync(ReceiveCallback callback);

void closeAsync(ResultCallback callback);

Expand Down
45 changes: 45 additions & 0 deletions tests/ReaderTest.cc
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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);

Expand Down