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

Reply via email to