This is an automated email from the ASF dual-hosted git repository.

mmerli 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 8ab6c34  Cpp client: add readCompacted in consumer config (#1858)
8ab6c34 is described below

commit 8ab6c34541eac5ee14e38764c2d4e3db13000da5
Author: Jia Zhai <[email protected]>
AuthorDate: Tue Jun 12 01:18:23 2018 +0800

    Cpp client: add readCompacted in consumer config (#1858)
    
    * add readCompacted in consumer config
    
    * change following comments
    
    * change following comments
    
    * add config in c wrapper
    
    * add e2e test in pulsar_test
    
    * add readerconfig in c wrapper
    
    * change following @ivan's comments
    
    * fix rebase error
---
 .../include/pulsar/ConsumerConfiguration.h         |   3 +
 .../include/pulsar/ReaderConfiguration.h           |   3 +
 .../include/pulsar/c/consumer_configuration.h      |   5 +
 .../include/pulsar/c/reader_configuration.h        |   5 +
 pulsar-client-cpp/lib/ClientImpl.cc                |   6 ++
 pulsar-client-cpp/lib/Commands.cc                  |   3 +-
 pulsar-client-cpp/lib/Commands.h                   |   3 +-
 pulsar-client-cpp/lib/ConsumerConfiguration.cc     |   4 +
 pulsar-client-cpp/lib/ConsumerConfigurationImpl.h  |   4 +-
 pulsar-client-cpp/lib/ConsumerImpl.cc              |  10 +-
 pulsar-client-cpp/lib/ConsumerImpl.h               |   2 +
 pulsar-client-cpp/lib/ReaderConfiguration.cc       |   4 +
 pulsar-client-cpp/lib/ReaderConfigurationImpl.h    |   7 +-
 pulsar-client-cpp/lib/ReaderImpl.cc                |   1 +
 pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc |   9 ++
 pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc   |   9 ++
 pulsar-client-cpp/python/pulsar/__init__.py        |   5 +-
 pulsar-client-cpp/python/pulsar_test.py            |  55 +++++++++++
 pulsar-client-cpp/python/src/config.cc             |   4 +
 .../tests/ConsumerConfigurationTest.cc             | 105 +++++++++++++++++++++
 20 files changed, 239 insertions(+), 8 deletions(-)

diff --git a/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h 
b/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
index 274abd5..c9584e3 100644
--- a/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
+++ b/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
@@ -149,6 +149,9 @@ class ConsumerConfiguration {
     ConsumerCryptoFailureAction getCryptoFailureAction() const;
     ConsumerConfiguration& setCryptoFailureAction(ConsumerCryptoFailureAction 
action);
 
+    bool isReadCompacted() const;
+    void setReadCompacted(bool compacted);
+
     friend class PulsarWrapper;
 
    private:
diff --git a/pulsar-client-cpp/include/pulsar/ReaderConfiguration.h 
b/pulsar-client-cpp/include/pulsar/ReaderConfiguration.h
index 3102cf5..8d365ab 100644
--- a/pulsar-client-cpp/include/pulsar/ReaderConfiguration.h
+++ b/pulsar-client-cpp/include/pulsar/ReaderConfiguration.h
@@ -86,6 +86,9 @@ class ReaderConfiguration {
     void setSubscriptionRolePrefix(const std::string& subscriptionRolePrefix);
     const std::string& getSubscriptionRolePrefix() const;
 
+    void setReadCompacted(bool compacted);
+    bool isReadCompacted() const;
+
    private:
     boost::shared_ptr<ReaderConfigurationImpl> impl_;
 };
diff --git a/pulsar-client-cpp/include/pulsar/c/consumer_configuration.h 
b/pulsar-client-cpp/include/pulsar/c/consumer_configuration.h
index fcefc97..7299867 100644
--- a/pulsar-client-cpp/include/pulsar/c/consumer_configuration.h
+++ b/pulsar-client-cpp/include/pulsar/c/consumer_configuration.h
@@ -147,6 +147,11 @@ long 
pulsar_consumer_get_unacked_messages_timeout_ms(pulsar_consumer_configurati
 
 int pulsar_consumer_is_encryption_enabled(pulsar_consumer_configuration_t 
*consumer_configuration);
 
+int pulsar_consumer_is_read_compacted(pulsar_consumer_configuration_t 
*consumer_configuration);
+
+void pulsar_consumer_set_read_compacted(pulsar_consumer_configuration_t 
*consumer_configuration,
+                                        int compacted);
+
 // const CryptoKeyReaderPtr getCryptoKeyReader()
 //
 // const;
diff --git a/pulsar-client-cpp/include/pulsar/c/reader_configuration.h 
b/pulsar-client-cpp/include/pulsar/c/reader_configuration.h
index faac37c..327fef1 100644
--- a/pulsar-client-cpp/include/pulsar/c/reader_configuration.h
+++ b/pulsar-client-cpp/include/pulsar/c/reader_configuration.h
@@ -79,6 +79,11 @@ void 
pulsar_reader_configuration_set_subscription_role_prefix(pulsar_reader_conf
 const char *pulsar_reader_configuration_get_subscription_role_prefix(
     pulsar_reader_configuration_t *configuration);
 
+void 
pulsar_reader_configuration_set_read_compacted(pulsar_reader_configuration_t 
*configuration,
+                                                    int readCompacted);
+
+int 
pulsar_reader_configuration_is_read_compacted(pulsar_reader_configuration_t 
*configuration);
+
 #pragma GCC visibility pop
 
 #ifdef __cplusplus
diff --git a/pulsar-client-cpp/lib/ClientImpl.cc 
b/pulsar-client-cpp/lib/ClientImpl.cc
index 80f947b..d8d07e7 100644
--- a/pulsar-client-cpp/lib/ClientImpl.cc
+++ b/pulsar-client-cpp/lib/ClientImpl.cc
@@ -222,6 +222,12 @@ void ClientImpl::subscribeAsync(const std::string& topic, 
const std::string& con
             lock.unlock();
             callback(ResultInvalidTopicName, Consumer());
             return;
+        } else if (conf.isReadCompacted() && 
(topicName->getDomain().compare("persistent") != 0 ||
+                                              (conf.getConsumerType() != 
ConsumerExclusive &&
+                                               conf.getConsumerType() != 
ConsumerFailover))) {
+            lock.unlock();
+            callback(ResultInvalidConfiguration, Consumer());
+            return;
         }
     }
 
diff --git a/pulsar-client-cpp/lib/Commands.cc 
b/pulsar-client-cpp/lib/Commands.cc
index 677fce0..fcf8117 100644
--- a/pulsar-client-cpp/lib/Commands.cc
+++ b/pulsar-client-cpp/lib/Commands.cc
@@ -185,7 +185,7 @@ SharedBuffer Commands::newConnect(const AuthenticationPtr& 
authentication, const
 SharedBuffer Commands::newSubscribe(const std::string& topic, const 
std::string& subscription,
                                     uint64_t consumerId, uint64_t requestId, 
CommandSubscribe_SubType subType,
                                     const std::string& consumerName, 
SubscriptionMode subscriptionMode,
-                                    Optional<MessageId> startMessageId) {
+                                    Optional<MessageId> startMessageId, bool 
readCompacted) {
     BaseCommand cmd;
     cmd.set_type(BaseCommand::SUBSCRIBE);
     CommandSubscribe* subscribe = cmd.mutable_subscribe();
@@ -196,6 +196,7 @@ SharedBuffer Commands::newSubscribe(const std::string& 
topic, const std::string&
     subscribe->set_request_id(requestId);
     subscribe->set_consumer_name(consumerName);
     subscribe->set_durable(subscriptionMode == SubscriptionModeDurable);
+    subscribe->set_read_compacted(readCompacted);
     if (startMessageId.is_present()) {
         MessageIdData& messageIdData = *subscribe->mutable_start_message_id();
         messageIdData.set_ledgerid(startMessageId.value().ledgerId());
diff --git a/pulsar-client-cpp/lib/Commands.h b/pulsar-client-cpp/lib/Commands.h
index 38c498a..5018116 100644
--- a/pulsar-client-cpp/lib/Commands.h
+++ b/pulsar-client-cpp/lib/Commands.h
@@ -77,7 +77,8 @@ class Commands {
     static SharedBuffer newSubscribe(const std::string& topic, const 
std::string& subscription,
                                      uint64_t consumerId, uint64_t requestId,
                                      proto::CommandSubscribe_SubType subType, 
const std::string& consumerName,
-                                     SubscriptionMode subscriptionMode, 
Optional<MessageId> startMessageId);
+                                     SubscriptionMode subscriptionMode, 
Optional<MessageId> startMessageId,
+                                     bool readCompacted);
 
     static SharedBuffer newUnsubscribe(uint64_t consumerId, uint64_t 
requestId);
 
diff --git a/pulsar-client-cpp/lib/ConsumerConfiguration.cc 
b/pulsar-client-cpp/lib/ConsumerConfiguration.cc
index 1b99f58..0c145c1 100644
--- a/pulsar-client-cpp/lib/ConsumerConfiguration.cc
+++ b/pulsar-client-cpp/lib/ConsumerConfiguration.cc
@@ -101,4 +101,8 @@ ConsumerConfiguration& 
ConsumerConfiguration::setCryptoFailureAction(ConsumerCry
     return *this;
 }
 
+bool ConsumerConfiguration::isReadCompacted() const { return 
impl_->readCompacted; }
+
+void ConsumerConfiguration::setReadCompacted(bool compacted) { 
impl_->readCompacted = compacted; }
+
 }  // namespace pulsar
diff --git a/pulsar-client-cpp/lib/ConsumerConfigurationImpl.h 
b/pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
index 91434a2..eb0c374 100644
--- a/pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
+++ b/pulsar-client-cpp/lib/ConsumerConfigurationImpl.h
@@ -34,6 +34,7 @@ struct ConsumerConfigurationImpl {
     long brokerConsumerStatsCacheTimeInMs;
     CryptoKeyReaderPtr cryptoKeyReader;
     ConsumerCryptoFailureAction cryptoFailureAction;
+    bool readCompacted;
     ConsumerConfigurationImpl()
         : unAckedMessagesTimeoutMs(0),
           consumerType(ConsumerExclusive),
@@ -43,7 +44,8 @@ struct ConsumerConfigurationImpl {
           receiverQueueSize(1000),
           maxTotalReceiverQueueSizeAcrossPartitions(50000),
           cryptoKeyReader(),
-          cryptoFailureAction(ConsumerCryptoFailureAction::FAIL) {}
+          cryptoFailureAction(ConsumerCryptoFailureAction::FAIL),
+          readCompacted(false) {}
 };
 }  // namespace pulsar
 #endif /* LIB_CONSUMERCONFIGURATIONIMPL_H_ */
diff --git a/pulsar-client-cpp/lib/ConsumerImpl.cc 
b/pulsar-client-cpp/lib/ConsumerImpl.cc
index 0425238..a50394f 100644
--- a/pulsar-client-cpp/lib/ConsumerImpl.cc
+++ b/pulsar-client-cpp/lib/ConsumerImpl.cc
@@ -57,7 +57,8 @@ ConsumerImpl::ConsumerImpl(const ClientImplPtr client, const 
std::string& topic,
       batchAcknowledgementTracker_(topic_, subscription, (long)consumerId_),
       brokerConsumerStats_(),
       consumerStatsBasePtr_(),
-      msgCrypto_() {
+      msgCrypto_(),
+      readCompacted_(conf.isReadCompacted()) {
     std::stringstream consumerStrStream;
     consumerStrStream << "[" << topic_ << ", " << subscription_ << ", " << 
consumerId_ << "] ";
     consumerStr_ = consumerStrStream.str();
@@ -135,8 +136,9 @@ void ConsumerImpl::connectionOpened(const 
ClientConnectionPtr& cnx) {
 
     ClientImplPtr client = client_.lock();
     uint64_t requestId = client->newRequestId();
-    SharedBuffer cmd = Commands::newSubscribe(topic_, subscription_, 
consumerId_, requestId, getSubType(),
-                                              consumerName_, 
subscriptionMode_, startMessageId_);
+    SharedBuffer cmd =
+        Commands::newSubscribe(topic_, subscription_, consumerId_, requestId, 
getSubType(), consumerName_,
+                               subscriptionMode_, startMessageId_, 
readCompacted_);
     cnx->sendRequestWithId(cmd, requestId)
         .addListener(boost::bind(&ConsumerImpl::handleCreateConsumer, 
shared_from_this(), cnx, _1));
 }
@@ -919,4 +921,6 @@ void ConsumerImpl::seekAsync(const MessageId& msgId, 
ResultCallback callback) {
     callback(ResultNotConnected);
 }
 
+bool ConsumerImpl::isReadCompacted() { return readCompacted_; }
+
 } /* namespace pulsar */
diff --git a/pulsar-client-cpp/lib/ConsumerImpl.h 
b/pulsar-client-cpp/lib/ConsumerImpl.h
index e35f444..d16c8b35 100644
--- a/pulsar-client-cpp/lib/ConsumerImpl.h
+++ b/pulsar-client-cpp/lib/ConsumerImpl.h
@@ -104,6 +104,7 @@ class ConsumerImpl : public ConsumerImplBase,
     virtual void getBrokerConsumerStatsAsync(BrokerConsumerStatsCallback 
callback);
     void handleSeek(Result result, ResultCallback callback);
     virtual void seekAsync(const MessageId& msgId, ResultCallback callback);
+    virtual bool isReadCompacted();
 
    protected:
     void connectionOpened(const ClientConnectionPtr& cnx);
@@ -165,6 +166,7 @@ class ConsumerImpl : public ConsumerImplBase,
     BrokerConsumerStatsImpl brokerConsumerStats_;
 
     MessageCryptoPtr msgCrypto_;
+    const bool readCompacted_;
 
     friend class PulsarFriend;
 };
diff --git a/pulsar-client-cpp/lib/ReaderConfiguration.cc 
b/pulsar-client-cpp/lib/ReaderConfiguration.cc
index e44714f..a649b33 100644
--- a/pulsar-client-cpp/lib/ReaderConfiguration.cc
+++ b/pulsar-client-cpp/lib/ReaderConfiguration.cc
@@ -56,4 +56,8 @@ const std::string& 
ReaderConfiguration::getSubscriptionRolePrefix() const {
 void ReaderConfiguration::setSubscriptionRolePrefix(const std::string& 
subscriptionRolePrefix) {
     impl_->subscriptionRolePrefix = subscriptionRolePrefix;
 }
+
+bool ReaderConfiguration::isReadCompacted() const { return 
impl_->readCompacted; }
+
+void ReaderConfiguration::setReadCompacted(bool compacted) { 
impl_->readCompacted = compacted; }
 }  // namespace pulsar
diff --git a/pulsar-client-cpp/lib/ReaderConfigurationImpl.h 
b/pulsar-client-cpp/lib/ReaderConfigurationImpl.h
index 67ddb3f..5dca8e3 100644
--- a/pulsar-client-cpp/lib/ReaderConfigurationImpl.h
+++ b/pulsar-client-cpp/lib/ReaderConfigurationImpl.h
@@ -29,8 +29,13 @@ struct ReaderConfigurationImpl {
     int receiverQueueSize;
     std::string readerName;
     std::string subscriptionRolePrefix;
+    bool readCompacted;
     ReaderConfigurationImpl()
-        : hasReaderListener(false), receiverQueueSize(1000), readerName(), 
subscriptionRolePrefix() {}
+        : hasReaderListener(false),
+          receiverQueueSize(1000),
+          readerName(),
+          subscriptionRolePrefix(),
+          readCompacted(false) {}
 };
 }  // namespace pulsar
 #endif /* LIB_READERCONFIGURATIONIMPL_H_ */
diff --git a/pulsar-client-cpp/lib/ReaderImpl.cc 
b/pulsar-client-cpp/lib/ReaderImpl.cc
index 5e69a72..45c509b 100644
--- a/pulsar-client-cpp/lib/ReaderImpl.cc
+++ b/pulsar-client-cpp/lib/ReaderImpl.cc
@@ -32,6 +32,7 @@ void ReaderImpl::start(const MessageId& startMessageId) {
     ConsumerConfiguration consumerConf;
     consumerConf.setConsumerType(ConsumerExclusive);
     consumerConf.setReceiverQueueSize(readerConf_.getReceiverQueueSize());
+    consumerConf.setReadCompacted(readerConf_.isReadCompacted());
 
     if (readerConf_.getReaderName().length() > 0) {
         consumerConf.setConsumerName(readerConf_.getReaderName());
diff --git a/pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc 
b/pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc
index 42c7c73..c8d5453 100644
--- a/pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc
+++ b/pulsar-client-cpp/lib/c/c_ConsumerConfiguration.cc
@@ -104,3 +104,12 @@ long pulsar_consumer_get_unacked_messages_timeout_ms(
 int pulsar_consumer_is_encryption_enabled(pulsar_consumer_configuration_t 
*consumer_configuration) {
     return consumer_configuration->consumerConfiguration.isEncryptionEnabled();
 }
+
+int pulsar_consumer_is_read_compacted(pulsar_consumer_configuration_t 
*consumer_configuration) {
+    return consumer_configuration->consumerConfiguration.isReadCompacted();
+}
+
+void pulsar_consumer_set_read_compacted(pulsar_consumer_configuration_t 
*consumer_configuration,
+                                        int compacted) {
+    consumer_configuration->consumerConfiguration.setReadCompacted(compacted);
+}
diff --git a/pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc 
b/pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc
index 55a6dc5..9b764f5 100644
--- a/pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc
+++ b/pulsar-client-cpp/lib/c/c_ReaderConfiguration.cc
@@ -75,4 +75,13 @@ void 
pulsar_reader_configuration_set_subscription_role_prefix(pulsar_reader_conf
 const char *pulsar_reader_configuration_get_subscription_role_prefix(
     pulsar_reader_configuration_t *configuration) {
     return configuration->conf.getSubscriptionRolePrefix().c_str();
+}
+
+void 
pulsar_reader_configuration_set_read_compacted(pulsar_reader_configuration_t 
*configuration,
+                                                    int readCompacted) {
+    configuration->conf.setReadCompacted(readCompacted);
+}
+
+int 
pulsar_reader_configuration_is_read_compacted(pulsar_reader_configuration_t 
*configuration) {
+    return configuration->conf.isReadCompacted();
 }
\ No newline at end of file
diff --git a/pulsar-client-cpp/python/pulsar/__init__.py 
b/pulsar-client-cpp/python/pulsar/__init__.py
index d8254d3..8c36c08 100644
--- a/pulsar-client-cpp/python/pulsar/__init__.py
+++ b/pulsar-client-cpp/python/pulsar/__init__.py
@@ -357,7 +357,8 @@ class Client:
                   receiver_queue_size=1000,
                   consumer_name=None,
                   unacked_messages_timeout_ms=None,
-                  broker_consumer_stats_cache_time_ms=30000
+                  broker_consumer_stats_cache_time_ms=30000,
+                  is_read_compacted=False
                   ):
         """
         Subscribe to the given topic and subscription combination.
@@ -415,9 +416,11 @@ class Client:
         _check_type_or_none(str, consumer_name, 'consumer_name')
         _check_type_or_none(int, unacked_messages_timeout_ms, 
'unacked_messages_timeout_ms')
         _check_type(int, broker_consumer_stats_cache_time_ms, 
'broker_consumer_stats_cache_time_ms')
+        _check_type(bool, is_read_compacted, 'is_read_compacted')
 
         conf = _pulsar.ConsumerConfiguration()
         conf.consumer_type(consumer_type)
+        conf.read_compacted(is_read_compacted)
         if message_listener:
             conf.message_listener(message_listener)
         conf.receiver_queue_size(receiver_queue_size)
diff --git a/pulsar-client-cpp/python/pulsar_test.py 
b/pulsar-client-cpp/python/pulsar_test.py
index e2f3e25..31163f8 100755
--- a/pulsar-client-cpp/python/pulsar_test.py
+++ b/pulsar-client-cpp/python/pulsar_test.py
@@ -39,6 +39,13 @@ def doHttpPost(url, data):
     req.add_header('Content-Type', 'application/json')
     urlopen(req)
 
+import urllib2
+def doHttpPut(url, data):
+    opener = urllib2.build_opener(urllib2.HTTPHandler)
+    request = urllib2.Request(url, data=data.encode())
+    request.add_header('Content-Type', 'application/json')
+    request.get_method = lambda: 'PUT'
+    opener.open(request)
 
 class PulsarTest(TestCase):
 
@@ -405,6 +412,54 @@ class PulsarTest(TestCase):
         self._check_value_error(lambda: client.create_reader(topic, 
MessageId.earliest, reader_name=5))
         client.close()
 
+    def test_publish_compact_and_consume(self):
+        client = Client(self.serviceUrl)
+        topic = 
'persistent://sample/standalone/ns1/my-python-test_publish_compact_and_consume'
+        producer = client.create_producer(topic, 
producer_name='my-producer-name', batching_enabled=False)
+        self.assertEqual(producer.last_sequence_id(), -1)
+        consumer = client.subscribe(topic, 'my-sub1', is_read_compacted=True)
+        consumer.close()
+        consumer2 = client.subscribe(topic, 'my-sub2', is_read_compacted=False)
+
+        # producer create 2 messages with same key.
+        producer.send('hello-0', partition_key='key0')
+        producer.send('hello-1', partition_key='key0')
+        producer.close()
+
+        # issue compact command, and wait success
+        url=self.adminUrl + 
'/admin/persistent/sample/standalone/ns1/my-python-test_publish_compact_and_consume/compaction'
+        doHttpPut(url, '')
+        while True:
+            req = urllib2.Request(url)
+            response = urllib2.urlopen(req)
+            s=response.read()
+            if 'RUNNING' in s:
+                print("Compact still running")
+                print(s)
+                time.sleep(0.2)
+            else:
+                self.assertTrue('SUCCESS' in s)
+                print("Compact Complete now")
+                print(s)
+                break
+
+        # after compact, consumer with `is_read_compacted=True`, expected read 
only the second message for same key.
+        consumer1 = client.subscribe(topic, 'my-sub1', is_read_compacted=True)
+        msg0 = consumer1.receive()
+        self.assertEqual(msg0.data(), b'hello-1')
+        consumer1.acknowledge(msg0)
+        consumer1.close()
+
+        # after compact, consumer with `is_read_compacted=False`, expected 
read 2 messages for same key.
+        msg0 = consumer2.receive()
+        self.assertEqual(msg0.data(), b'hello-0')
+        consumer2.acknowledge(msg0)
+        msg1 = consumer2.receive()
+        self.assertEqual(msg1.data(), b'hello-1')
+        consumer2.acknowledge(msg1)
+        consumer2.close()
+        client.close()
+
     def _check_value_error(self, fun):
         try:
             fun()
diff --git a/pulsar-client-cpp/python/src/config.cc 
b/pulsar-client-cpp/python/src/config.cc
index f03c883..9149add 100644
--- a/pulsar-client-cpp/python/src/config.cc
+++ b/pulsar-client-cpp/python/src/config.cc
@@ -134,6 +134,8 @@ void export_config() {
             .def("unacked_messages_timeout_ms", 
&ConsumerConfiguration::setUnAckedMessagesTimeoutMs)
             .def("broker_consumer_stats_cache_time_ms", 
&ConsumerConfiguration::getBrokerConsumerStatsCacheTimeInMs)
             .def("broker_consumer_stats_cache_time_ms", 
&ConsumerConfiguration::setBrokerConsumerStatsCacheTimeInMs)
+            .def("read_compacted", &ConsumerConfiguration::isReadCompacted)
+            .def("read_compacted", &ConsumerConfiguration::setReadCompacted)
             ;
 
     class_<ReaderConfiguration>("ReaderConfiguration")
@@ -144,5 +146,7 @@ void export_config() {
             .def("reader_name", &ReaderConfiguration::setReaderName)
             .def("subscription_role_prefix", 
&ReaderConfiguration::getSubscriptionRolePrefix, 
return_value_policy<copy_const_reference>())
             .def("subscription_role_prefix", 
&ReaderConfiguration::setSubscriptionRolePrefix)
+            .def("read_compacted", &ReaderConfiguration::isReadCompacted)
+            .def("read_compacted", &ReaderConfiguration::setReadCompacted)
             ;
 }
diff --git a/pulsar-client-cpp/tests/ConsumerConfigurationTest.cc 
b/pulsar-client-cpp/tests/ConsumerConfigurationTest.cc
new file mode 100644
index 0000000..cede638
--- /dev/null
+++ b/pulsar-client-cpp/tests/ConsumerConfigurationTest.cc
@@ -0,0 +1,105 @@
+/**
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+#include <pulsar/Client.h>
+#include <gtest/gtest.h>
+
+#include "../lib/Future.h"
+#include "../lib/Utils.h"
+
+using namespace pulsar;
+
+TEST(ConsumerConfigurationTest, testReadCompactPersistentExclusive) {
+    std::string lookupUrl = "pulsar://localhost:8885";
+    std::string topicName = "persistent://prop/unit/ns1/persist-topic";
+    std::string subName = "test-persist-exclusive";
+
+    Result result;
+
+    ConsumerConfiguration config;
+    config.setReadCompacted(true);
+    config.setConsumerType(ConsumerExclusive);
+
+    ClientConfiguration clientConfig;
+    Client client(lookupUrl, clientConfig);
+
+    Consumer consumer;
+    result = client.subscribe(topicName, subName, config, consumer);
+    ASSERT_EQ(ResultOk, result);
+    consumer.close();
+}
+
+TEST(ConsumerConfigurationTest, testReadCompactPersistentFailover) {
+    std::string lookupUrl = "pulsar://localhost:8885";
+    std::string topicName = "persistent://prop/unit/ns1/persist-topic";
+    std::string subName = "test-persist-fail-over";
+
+    Result result;
+
+    ConsumerConfiguration config;
+    config.setReadCompacted(true);
+    config.setConsumerType(ConsumerFailover);
+
+    ClientConfiguration clientConfig;
+    Client client(lookupUrl, clientConfig);
+
+    Consumer consumer;
+    result = client.subscribe(topicName, subName, config, consumer);
+    ASSERT_EQ(ResultOk, result);
+    consumer.close();
+}
+
+TEST(ConsumerConfigurationTest, testReadCompactPersistentShared) {
+    std::string lookupUrl = "pulsar://localhost:8885";
+    std::string topicName = "persistent://prop/unit/ns1/persist-topic";
+    std::string subName = "test-persist-shared";
+
+    Result result;
+
+    ConsumerConfiguration config;
+    config.setReadCompacted(true);
+    config.setConsumerType(ConsumerShared);
+
+    ClientConfiguration clientConfig;
+    Client client(lookupUrl, clientConfig);
+
+    Consumer consumer;
+    result = client.subscribe(topicName, subName, config, consumer);
+    ASSERT_EQ(ResultInvalidConfiguration, result);
+    consumer.close();
+}
+
+TEST(ConsumerConfigurationTest, testReadCompactNonPersistentExclusive) {
+    std::string lookupUrl = "pulsar://localhost:8885";
+    std::string topicName = 
"non-persistent://prop/unit/ns1/testNonPersistentTopic";
+    std::string subName = "test-non-persist-exclusive";
+
+    Result result;
+
+    ConsumerConfiguration config;
+    config.setReadCompacted(true);
+    config.setConsumerType(ConsumerExclusive);
+
+    ClientConfiguration clientConfig;
+    Client client(lookupUrl, clientConfig);
+
+    Consumer consumer;
+    result = client.subscribe(topicName, subName, config, consumer);
+    ASSERT_EQ(ResultInvalidConfiguration, result);
+    consumer.close();
+}
\ No newline at end of file

-- 
To stop receiving notification emails like this one, please contact
[email protected].

Reply via email to