sijie closed pull request #1861: Cpp client: add getLastMessageId and
hasMessageAvailable in consmer and reader
URL: https://github.com/apache/incubator-pulsar/pull/1861
This is a PR merged from a forked repository.
As GitHub hides the original diff on merge, it is displayed below for
the sake of provenance:
As this is a foreign pull request (from a fork), the diff is supplied
below (as it won't show otherwise due to GitHub magic):
diff --git a/pulsar-client-cpp/include/pulsar/MessageId.h
b/pulsar-client-cpp/include/pulsar/MessageId.h
index 149d177eae..dfe3a51fca 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 2985f5066f..9ce9f44b0b 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 f7b6b5093e..8d8243cd6a 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 2176e27d19..860ca6a58b 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 fcf81175cc..13bf99a870 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 501811688c..53fb1bb247 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 a50394f6c6..55ded8f2e3 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 d16c8b350d..fcdaed122c 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 533319bb0a..cd86d626e8 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 45c509b265..9507a0b25a 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 61216b84d7..d346316575 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 be9678daa6..1f293c9364 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
----------------------------------------------------------------
This is an automated message from the Apache Git Service.
To respond to the message, please log on GitHub and use the
URL above to go to the specific comment.
For queries about this service, please contact Infrastructure at:
[email protected]
With regards,
Apache Git Services