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

Reply via email to