This is an automated email from the ASF dual-hosted git repository.
sijie pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/incubator-pulsar.git
The following commit(s) were added to refs/heads/master by this push:
new e7cc354 Cpp client: add getLastMessageId and hasMessageAvailable in
consmer and reader (#1861)
e7cc354 is described below
commit e7cc3544c2ea94dda29120e83a86ed650ce20e2f
Author: Jia Zhai <[email protected]>
AuthorDate: Thu Jun 21 06:56:21 2018 +0800
Cpp client: add getLastMessageId and hasMessageAvailable in consmer and
reader (#1861)
* add getlastmessageid and hasMessageAvailable in consmer and reader
* change following comments
* rebase master, change followinw @ivan's comments
---
pulsar-client-cpp/include/pulsar/MessageId.h | 2 +-
pulsar-client-cpp/include/pulsar/Reader.h | 6 ++
pulsar-client-cpp/lib/ClientConnection.cc | 59 ++++++++++-
pulsar-client-cpp/lib/ClientConnection.h | 5 +
pulsar-client-cpp/lib/Commands.cc | 18 ++++
pulsar-client-cpp/lib/Commands.h | 1 +
pulsar-client-cpp/lib/ConsumerImpl.cc | 74 +++++++++++++-
pulsar-client-cpp/lib/ConsumerImpl.h | 16 +++
pulsar-client-cpp/lib/Reader.cc | 21 ++++
pulsar-client-cpp/lib/ReaderImpl.cc | 5 +
pulsar-client-cpp/lib/ReaderImpl.h | 2 +
pulsar-client-cpp/tests/ReaderTest.cc | 144 +++++++++++++++++++++++++++
12 files changed, 350 insertions(+), 3 deletions(-)
diff --git a/pulsar-client-cpp/include/pulsar/MessageId.h
b/pulsar-client-cpp/include/pulsar/MessageId.h
index 149d177..dfe3a51 100644
--- a/pulsar-client-cpp/include/pulsar/MessageId.h
+++ b/pulsar-client-cpp/include/pulsar/MessageId.h
@@ -33,6 +33,7 @@ class MessageId {
public:
MessageId& operator=(const MessageId&);
MessageId();
+ explicit MessageId(int32_t partition, int64_t ledgerId, int64_t entryId,
int32_t batchIndex);
/**
* MessageId representing the "earliest" or "oldest available" message
stored in the topic
@@ -75,7 +76,6 @@ class MessageId {
friend class PulsarWrapper;
friend class PulsarFriend;
- explicit MessageId(int32_t partition, int64_t ledgerId, int64_t entryId,
int32_t batchIndex);
friend std::ostream& operator<<(std::ostream& s, const MessageId&
messageId);
int64_t ledgerId() const;
diff --git a/pulsar-client-cpp/include/pulsar/Reader.h
b/pulsar-client-cpp/include/pulsar/Reader.h
index 2985f50..9ce9f44 100644
--- a/pulsar-client-cpp/include/pulsar/Reader.h
+++ b/pulsar-client-cpp/include/pulsar/Reader.h
@@ -29,6 +29,8 @@ class PulsarWrapper;
class PulsarFriend;
class ReaderImpl;
+typedef boost::function<void(Result result, bool hasMessageAvailable)>
HasMessageAvailableCallback;
+
/**
* A Reader can be used to scan through all the messages currently available
in a topic.
*/
@@ -71,6 +73,10 @@ class Reader {
void closeAsync(ResultCallback callback);
+ void hasMessageAvailableAsync(HasMessageAvailableCallback callback);
+
+ Result hasMessageAvailable(bool& hasMessageAvailable);
+
private:
typedef boost::shared_ptr<ReaderImpl> ReaderImplPtr;
ReaderImplPtr impl_;
diff --git a/pulsar-client-cpp/lib/ClientConnection.cc
b/pulsar-client-cpp/lib/ClientConnection.cc
index f7b6b50..8d8243c 100644
--- a/pulsar-client-cpp/lib/ClientConnection.cc
+++ b/pulsar-client-cpp/lib/ClientConnection.cc
@@ -855,8 +855,19 @@ void ClientConnection::handleIncomingCommand() {
requestData.promise.setFailed(getResult(error.error()));
requestData.timer->cancel();
} else {
- lock.unlock();
+ PendingGetLastMessageIdRequestsMap::iterator it2 =
+
pendingGetLastMessageIdRequests_.find(error.request_id());
+ if (it2 != pendingGetLastMessageIdRequests_.end()) {
+ Promise<Result, MessageId> getLastMessageIdPromise
= it2->second;
+ pendingGetLastMessageIdRequests_.erase(it2);
+ lock.unlock();
+
+
getLastMessageIdPromise.setFailed(getResult(error.error()));
+ } else {
+ lock.unlock();
+ }
}
+
break;
}
@@ -928,6 +939,35 @@ void ClientConnection::handleIncomingCommand() {
break;
}
+ case BaseCommand::GET_LAST_MESSAGE_ID_RESPONSE: {
+ const CommandGetLastMessageIdResponse&
getLastMessageIdResponse =
+ incomingCmd_.getlastmessageidresponse();
+ LOG_DEBUG(cnxString_ << "Received getLastMessageIdResponse
from server. req_id: "
+ <<
getLastMessageIdResponse.request_id());
+
+ Lock lock(mutex_);
+ PendingGetLastMessageIdRequestsMap::iterator it =
+
pendingGetLastMessageIdRequests_.find(getLastMessageIdResponse.request_id());
+
+ if (it != pendingGetLastMessageIdRequests_.end()) {
+ Promise<Result, MessageId> getLastMessageIdPromise =
it->second;
+ pendingGetLastMessageIdRequests_.erase(it);
+ lock.unlock();
+
+ MessageIdData messageIdData =
getLastMessageIdResponse.last_message_id();
+ MessageId messageId =
MessageId(messageIdData.partition(), messageIdData.ledgerid(),
+
messageIdData.entryid(), messageIdData.batch_index());
+
+ getLastMessageIdPromise.setValue(messageId);
+ } else {
+ lock.unlock();
+ LOG_WARN(
+ "getLastMessageIdResponse command - Received
unknown request id from server: "
+ << getLastMessageIdResponse.request_id());
+ }
+ break;
+ }
+
default: {
LOG_WARN(cnxString_ << "Received invalid message from
server");
close();
@@ -1214,4 +1254,21 @@ int ClientConnection::getServerProtocolVersion() const {
return serverProtocolVe
Commands::ChecksumType ClientConnection::getChecksumType() const {
return getServerProtocolVersion() >= proto::v6 ? Commands::Crc32c :
Commands::None;
}
+
+Future<Result, MessageId> ClientConnection::newGetLastMessageId(uint64_t
consumerId, uint64_t requestId) {
+ Lock lock(mutex_);
+ Promise<Result, MessageId> promise;
+ if (isClosed()) {
+ lock.unlock();
+ LOG_ERROR(cnxString_ << " Client is not connected to the broker");
+ promise.setFailed(ResultNotConnected);
+ return promise.getFuture();
+ }
+
+ pendingGetLastMessageIdRequests_.insert(std::make_pair(requestId,
promise));
+ lock.unlock();
+ sendCommand(Commands::newGetLastMessageId(consumerId, requestId));
+ return promise.getFuture();
+}
+
} // namespace pulsar
diff --git a/pulsar-client-cpp/lib/ClientConnection.h
b/pulsar-client-cpp/lib/ClientConnection.h
index 2176e27..860ca6a 100644
--- a/pulsar-client-cpp/lib/ClientConnection.h
+++ b/pulsar-client-cpp/lib/ClientConnection.h
@@ -142,6 +142,8 @@ class ClientConnection : public
boost::enable_shared_from_this<ClientConnection>
Future<Result, BrokerConsumerStatsImpl> newConsumerStats(uint64_t
consumerId, uint64_t requestId);
+ Future<Result, MessageId> newGetLastMessageId(uint64_t consumerId,
uint64_t requestId);
+
private:
struct PendingRequestData {
Promise<Result, ResponseData> promise;
@@ -265,6 +267,9 @@ class ClientConnection : public
boost::enable_shared_from_this<ClientConnection>
typedef std::map<uint64_t, Promise<Result, BrokerConsumerStatsImpl> >
PendingConsumerStatsMap;
PendingConsumerStatsMap pendingConsumerStatsMap_;
+ typedef std::map<long, Promise<Result, MessageId> >
PendingGetLastMessageIdRequestsMap;
+ PendingGetLastMessageIdRequestsMap pendingGetLastMessageIdRequests_;
+
boost::mutex mutex_;
typedef boost::unique_lock<boost::mutex> Lock;
diff --git a/pulsar-client-cpp/lib/Commands.cc
b/pulsar-client-cpp/lib/Commands.cc
index fcf8117..13bf99a 100644
--- a/pulsar-client-cpp/lib/Commands.cc
+++ b/pulsar-client-cpp/lib/Commands.cc
@@ -312,6 +312,18 @@ SharedBuffer Commands::newSeek(uint64_t consumerId,
uint64_t requestId, const Me
return writeMessageWithSize(cmd);
}
+SharedBuffer Commands::newGetLastMessageId(uint64_t consumerId, uint64_t
requestId) {
+ BaseCommand cmd;
+ cmd.set_type(BaseCommand::GET_LAST_MESSAGE_ID);
+
+ CommandGetLastMessageId* getLastMessageId = cmd.mutable_getlastmessageid();
+ getLastMessageId->set_consumer_id(consumerId);
+ getLastMessageId->set_request_id(requestId);
+ const SharedBuffer buffer = writeMessageWithSize(cmd);
+ cmd.clear_getlastmessageid();
+ return buffer;
+}
+
std::string Commands::messageType(BaseCommand_Type type) {
switch (type) {
case BaseCommand::CONNECT:
@@ -398,6 +410,12 @@ std::string Commands::messageType(BaseCommand_Type type) {
case BaseCommand::ACTIVE_CONSUMER_CHANGE:
return "ACTIVE_CONSUMER_CHANGE";
break;
+ case BaseCommand::GET_LAST_MESSAGE_ID:
+ return "GET_LAST_MESSAGE_ID";
+ break;
+ case BaseCommand::GET_LAST_MESSAGE_ID_RESPONSE:
+ return "GET_LAST_MESSAGE_ID_RESPONSE";
+ break;
};
}
diff --git a/pulsar-client-cpp/lib/Commands.h b/pulsar-client-cpp/lib/Commands.h
index 5018116..53fb1bb 100644
--- a/pulsar-client-cpp/lib/Commands.h
+++ b/pulsar-client-cpp/lib/Commands.h
@@ -111,6 +111,7 @@ class Commands {
static SharedBuffer newConsumerStats(uint64_t consumerId, uint64_t
requestId);
static SharedBuffer newSeek(uint64_t consumerId, uint64_t requestId, const
MessageId& messageId);
+ static SharedBuffer newGetLastMessageId(uint64_t consumerId, uint64_t
requestId);
private:
Commands();
diff --git a/pulsar-client-cpp/lib/ConsumerImpl.cc
b/pulsar-client-cpp/lib/ConsumerImpl.cc
index a50394f..55ded8f 100644
--- a/pulsar-client-cpp/lib/ConsumerImpl.cc
+++ b/pulsar-client-cpp/lib/ConsumerImpl.cc
@@ -58,7 +58,8 @@ ConsumerImpl::ConsumerImpl(const ClientImplPtr client, const
std::string& topic,
brokerConsumerStats_(),
consumerStatsBasePtr_(),
msgCrypto_(),
- readCompacted_(conf.isReadCompacted()) {
+ readCompacted_(conf.isReadCompacted()),
+ lastMessageInBroker_(Optional<MessageId>::of(MessageId())) {
std::stringstream consumerStrStream;
consumerStrStream << "[" << topic_ << ", " << subscription_ << ", " <<
consumerId_ << "] ";
consumerStr_ = consumerStrStream.str();
@@ -923,4 +924,75 @@ void ConsumerImpl::seekAsync(const MessageId& msgId,
ResultCallback callback) {
bool ConsumerImpl::isReadCompacted() { return readCompacted_; }
+void ConsumerImpl::hasMessageAvailableAsync(HasMessageAvailableCallback
callback) {
+ MessageId lastDequed = this->lastMessageIdDequed();
+ MessageId lastInBroker = this->lastMessageIdInBroker();
+ if (lastInBroker > lastDequed && lastInBroker.entryId() != -1) {
+ callback(ResultOk, true);
+ return;
+ }
+
+ BrokerGetLastMessageIdCallback callback1 = [this, lastDequed,
callback](Result result,
+
MessageId messageId) {
+ if (result == ResultOk) {
+ if (messageId > lastDequed && messageId.entryId() != -1) {
+ callback(ResultOk, true);
+ } else {
+ callback(ResultOk, false);
+ }
+ } else {
+ callback(result, false);
+ }
+ };
+
+ getLastMessageIdAsync(callback1);
+}
+
+void ConsumerImpl::brokerGetLastMessageIdListener(Result res, MessageId
messageId,
+
BrokerGetLastMessageIdCallback callback) {
+ Lock lock(mutex_);
+ if (messageId > lastMessageIdInBroker()) {
+ lastMessageInBroker_ = Optional<MessageId>::of(messageId);
+ lock.unlock();
+ callback(res, messageId);
+ } else {
+ lock.unlock();
+ callback(res, lastMessageIdInBroker());
+ }
+}
+
+void ConsumerImpl::getLastMessageIdAsync(BrokerGetLastMessageIdCallback
callback) {
+ Lock lock(mutex_);
+ if (state_ == Closed || state_ == Closing) {
+ lock.unlock();
+ LOG_ERROR(getName() << "Client connection already closed.");
+ if (!callback.empty()) {
+ callback(ResultAlreadyClosed, MessageId());
+ }
+ return;
+ }
+ lock.unlock();
+
+ ClientConnectionPtr cnx = getCnx().lock();
+ if (cnx) {
+ if (cnx->getServerProtocolVersion() >= proto::v12) {
+ ClientImplPtr client = client_.lock();
+ uint64_t requestId = client->newRequestId();
+ LOG_DEBUG(getName() << " Sending getLastMessageId Command for
Consumer - " << getConsumerId()
+ << ", requestId - " << requestId);
+
+ cnx->newGetLastMessageId(consumerId_, requestId)
+
.addListener(boost::bind(&ConsumerImpl::brokerGetLastMessageIdListener,
shared_from_this(),
+ _1, _2, callback));
+ } else {
+ LOG_ERROR(getName() << " Operation not supported since server
protobuf version "
+ << cnx->getServerProtocolVersion() << " is
older than proto::v12");
+ callback(ResultUnsupportedVersionError, MessageId());
+ }
+ } else {
+ LOG_ERROR(getName() << " Client Connection not ready for Consumer");
+ callback(ResultNotConnected, MessageId());
+ }
+}
+
} /* namespace pulsar */
diff --git a/pulsar-client-cpp/lib/ConsumerImpl.h
b/pulsar-client-cpp/lib/ConsumerImpl.h
index d16c8b35..fcdaed1 100644
--- a/pulsar-client-cpp/lib/ConsumerImpl.h
+++ b/pulsar-client-cpp/lib/ConsumerImpl.h
@@ -53,6 +53,8 @@ class BatchAcknowledgementTracker;
typedef boost::shared_ptr<ConsumerImpl> ConsumerImplPtr;
typedef boost::weak_ptr<ConsumerImpl> ConsumerImplWeakPtr;
typedef boost::shared_ptr<MessageCrypto> MessageCryptoPtr;
+typedef boost::function<void(Result result, MessageId messageId)>
BrokerGetLastMessageIdCallback;
+typedef boost::function<void(Result result, bool hasMessageAvailable)>
HasMessageAvailableCallback;
enum ConsumerTopicType
{
@@ -105,6 +107,8 @@ class ConsumerImpl : public ConsumerImplBase,
void handleSeek(Result result, ResultCallback callback);
virtual void seekAsync(const MessageId& msgId, ResultCallback callback);
virtual bool isReadCompacted();
+ virtual void hasMessageAvailableAsync(HasMessageAvailableCallback
callback);
+ virtual void getLastMessageIdAsync(BrokerGetLastMessageIdCallback
callback);
protected:
void connectionOpened(const ClientConnectionPtr& cnx);
@@ -168,6 +172,18 @@ class ConsumerImpl : public ConsumerImplBase,
MessageCryptoPtr msgCrypto_;
const bool readCompacted_;
+ Optional<MessageId> lastMessageInBroker_;
+ void brokerGetLastMessageIdListener(Result res, MessageId messageId,
+ BrokerGetLastMessageIdCallback
callback);
+
+ MessageId lastMessageIdDequed() {
+ return lastDequedMessage_.is_present() ? lastDequedMessage_.value() :
MessageId();
+ }
+
+ MessageId lastMessageIdInBroker() {
+ return lastMessageInBroker_.is_present() ?
lastMessageInBroker_.value() : MessageId();
+ }
+
friend class PulsarFriend;
};
diff --git a/pulsar-client-cpp/lib/Reader.cc b/pulsar-client-cpp/lib/Reader.cc
index 533319b..cd86d62 100644
--- a/pulsar-client-cpp/lib/Reader.cc
+++ b/pulsar-client-cpp/lib/Reader.cc
@@ -66,4 +66,25 @@ void Reader::closeAsync(ResultCallback callback) {
impl_->closeAsync(callback);
}
+
+void Reader::hasMessageAvailableAsync(HasMessageAvailableCallback callback) {
+ if (!impl_) {
+ callback(ResultConsumerNotInitialized, false);
+ return;
+ }
+
+ impl_->hasMessageAvailableAsync(callback);
+}
+
+Result Reader::hasMessageAvailable(bool& hasMessageAvailable) {
+ if (!impl_) {
+ return ResultConsumerNotInitialized;
+ }
+
+ Promise<Result, bool> promise;
+
+ impl_->hasMessageAvailableAsync(WaitForCallbackValue<bool>(promise));
+ return promise.getFuture().get(hasMessageAvailable);
+}
+
} // namespace pulsar
diff --git a/pulsar-client-cpp/lib/ReaderImpl.cc
b/pulsar-client-cpp/lib/ReaderImpl.cc
index 45c509b..9507a0b 100644
--- a/pulsar-client-cpp/lib/ReaderImpl.cc
+++ b/pulsar-client-cpp/lib/ReaderImpl.cc
@@ -97,4 +97,9 @@ void ReaderImpl::acknowledgeIfNecessary(Result result, const
Message& msg) {
}
void ReaderImpl::closeAsync(ResultCallback callback) {
consumer_->closeAsync(callback); }
+
+void ReaderImpl::hasMessageAvailableAsync(HasMessageAvailableCallback
callback) {
+ consumer_->hasMessageAvailableAsync(callback);
+}
+
} // namespace pulsar
diff --git a/pulsar-client-cpp/lib/ReaderImpl.h
b/pulsar-client-cpp/lib/ReaderImpl.h
index 61216b8..d346316 100644
--- a/pulsar-client-cpp/lib/ReaderImpl.h
+++ b/pulsar-client-cpp/lib/ReaderImpl.h
@@ -49,6 +49,8 @@ class ReaderImpl : public
boost::enable_shared_from_this<ReaderImpl> {
ConsumerImplPtr getConsumer();
+ void hasMessageAvailableAsync(HasMessageAvailableCallback callback);
+
private:
void handleConsumerCreated(Result result, ConsumerImplBaseWeakPtr
consumer);
diff --git a/pulsar-client-cpp/tests/ReaderTest.cc
b/pulsar-client-cpp/tests/ReaderTest.cc
index be9678d..1f293c9 100644
--- a/pulsar-client-cpp/tests/ReaderTest.cc
+++ b/pulsar-client-cpp/tests/ReaderTest.cc
@@ -23,6 +23,9 @@
#include <string>
+#include <lib/LogUtils.h>
+DECLARE_LOG_OBJECT()
+
using namespace pulsar;
static std::string serviceUrl = "pulsar://localhost:8885";
@@ -54,6 +57,8 @@ TEST(ReaderTest, testSimpleReader) {
ASSERT_EQ(expected, content);
}
+ producer.close();
+ reader.close();
client.close();
}
@@ -84,6 +89,8 @@ TEST(ReaderTest, testReaderAfterMessagesWerePublished) {
ASSERT_EQ(expected, content);
}
+ producer.close();
+ reader.close();
client.close();
}
@@ -126,6 +133,9 @@ TEST(ReaderTest, testMultipleReaders) {
ASSERT_EQ(expected, content);
}
+ producer.close();
+ reader1.close();
+ reader2.close();
client.close();
}
@@ -162,6 +172,8 @@ TEST(ReaderTest, testReaderOnLastMessage) {
ASSERT_EQ(expected, content);
}
+ producer.close();
+ reader.close();
client.close();
}
@@ -208,6 +220,8 @@ TEST(ReaderTest, testReaderOnSpecificMessage) {
ASSERT_EQ(expected, content);
}
+ producer.close();
+ reader.close();
client.close();
}
@@ -268,7 +282,137 @@ TEST(ReaderTest, testReaderOnSpecificMessageWithBatches) {
ASSERT_EQ(expected, content);
}
+ producer.close();
reader.close();
reader2.close();
client.close();
}
+
+TEST(ReaderTest, testReaderReachEndOfTopic) {
+ Client client(serviceUrl);
+
+ std::string topicName =
"persistent://property/cluster/namespace/testReaderReachEndOfTopic";
+
+ // 1. create producer
+ Producer producer;
+ // Enable batching
+ ProducerConfiguration producerConf;
+ producerConf.setBatchingEnabled(true);
+ producerConf.setBatchingMaxPublishDelayMs(1000);
+ ASSERT_EQ(ResultOk, client.createProducer(topicName, producerConf,
producer));
+
+ // 2. create reader, and expect hasMessageAvailable return false since no
message produced.
+ ReaderConfiguration readerConf;
+ Reader reader;
+ ASSERT_EQ(ResultOk, client.createReader(topicName, MessageId::latest(),
readerConf, reader));
+
+ bool hasMessageAvailable;
+ ASSERT_EQ(ResultOk, reader.hasMessageAvailable(hasMessageAvailable));
+ ASSERT_FALSE(hasMessageAvailable);
+
+ // 3. produce 10 messages.
+ for (int i = 0; i < 10; i++) {
+ std::string content = "my-message-" +
boost::lexical_cast<std::string>(i);
+ Message msg = MessageBuilder().setContent(content).build();
+ ASSERT_EQ(ResultOk, producer.send(msg));
+ }
+
+ // 4. expect hasMessageAvailable return true, and after read 10 messages
out, it return false.
+ ASSERT_EQ(ResultOk, reader.hasMessageAvailable(hasMessageAvailable));
+ ASSERT_TRUE(hasMessageAvailable);
+
+ int readMessageCount = 0;
+ for (; hasMessageAvailable; readMessageCount++) {
+ Message msg;
+ ASSERT_EQ(ResultOk, reader.readNext(msg));
+
+ std::string content = msg.getDataAsString();
+ std::string expected = "my-message-" +
boost::lexical_cast<std::string>(readMessageCount);
+ ASSERT_EQ(expected, content);
+ reader.hasMessageAvailable(hasMessageAvailable);
+ }
+
+ ASSERT_EQ(readMessageCount, 10);
+ ASSERT_FALSE(hasMessageAvailable);
+
+ // 5. produce another 10 messages, expect hasMessageAvailable return true,
+ // and after read these 10 messages out, it return false.
+ for (int i = 10; i < 20; i++) {
+ std::string content = "my-message-" +
boost::lexical_cast<std::string>(i);
+ Message msg = MessageBuilder().setContent(content).build();
+ ASSERT_EQ(ResultOk, producer.send(msg));
+ }
+
+ ASSERT_EQ(ResultOk, reader.hasMessageAvailable(hasMessageAvailable));
+ ASSERT_TRUE(hasMessageAvailable);
+
+ for (; hasMessageAvailable; readMessageCount++) {
+ Message msg;
+ ASSERT_EQ(ResultOk, reader.readNext(msg));
+
+ std::string content = msg.getDataAsString();
+ std::string expected = "my-message-" +
boost::lexical_cast<std::string>(readMessageCount);
+ ASSERT_EQ(expected, content);
+ reader.hasMessageAvailable(hasMessageAvailable);
+ }
+ ASSERT_EQ(readMessageCount, 20);
+ ASSERT_FALSE(hasMessageAvailable);
+
+ producer.close();
+ reader.close();
+ client.close();
+}
+
+TEST(ReaderTest, testReaderReachEndOfTopicMessageWithBatches) {
+ Client client(serviceUrl);
+
+ std::string topicName =
+
"persistent://property/cluster/namespace/testReaderReachEndOfTopicMessageWithBatches";
+
+ // 1. create producer
+ Producer producer;
+ ASSERT_EQ(ResultOk, client.createProducer(topicName, producer));
+
+ // 2. create reader, and expect hasMessageAvailable return false since no
message produced.
+ ReaderConfiguration readerConf;
+ Reader reader;
+ ASSERT_EQ(ResultOk, client.createReader(topicName, MessageId::latest(),
readerConf, reader));
+
+ bool hasMessageAvailable;
+ ASSERT_EQ(ResultOk, reader.hasMessageAvailable(hasMessageAvailable));
+ ASSERT_FALSE(hasMessageAvailable);
+
+ // 3. produce 10 messages in batches way.
+ for (int i = 0; i < 10; i++) {
+ std::string content = "my-message-" +
boost::lexical_cast<std::string>(i);
+ Message msg = MessageBuilder().setContent(content).build();
+ producer.sendAsync(msg, NULL);
+ }
+ // Send one sync message, to wait for everything before to be persisted as
well
+ std::string content = "my-message-10";
+ Message msg = MessageBuilder().setContent(content).build();
+ ASSERT_EQ(ResultOk, producer.send(msg));
+
+ // 4. expect hasMessageAvailable return true, and after read 11 messages
out, it return false.
+ ASSERT_EQ(ResultOk, reader.hasMessageAvailable(hasMessageAvailable));
+ ASSERT_TRUE(hasMessageAvailable);
+
+ std::string lastMessageId;
+ int readMessageCount = 0;
+ for (; hasMessageAvailable; readMessageCount++) {
+ Message msg;
+ ASSERT_EQ(ResultOk, reader.readNext(msg));
+
+ std::string content = msg.getDataAsString();
+ std::string expected = "my-message-" +
boost::lexical_cast<std::string>(readMessageCount);
+ ASSERT_EQ(expected, content);
+ reader.hasMessageAvailable(hasMessageAvailable);
+ msg.getMessageId().serialize(lastMessageId);
+ }
+ ASSERT_FALSE(hasMessageAvailable);
+ ASSERT_EQ(readMessageCount, 11);
+
+ producer.close();
+ reader.close();
+ client.close();
+}
\ No newline at end of file