merlimat closed pull request #1858: Cpp client: add readCompacted in consumer 
config
URL: https://github.com/apache/incubator-pulsar/pull/1858
 
 
   

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/ConsumerConfiguration.h 
b/pulsar-client-cpp/include/pulsar/ConsumerConfiguration.h
index 274abd52eb..c9584e38f2 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 3102cf57a1..8d365ab813 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 fcefc970c0..7299867f38 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 faac37c13e..327fef1124 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 80f947b607..d8d07e79e3 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 677fce0e80..fcf81175cc 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 38c498a463..501811688c 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 1b99f584f6..0c145c1cf7 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 91434a2fd8..eb0c3746bf 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 0425238567..a50394f6c6 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 e35f44464b..d16c8b350d 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 e44714fc06..a649b33e50 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 67ddb3fe54..5dca8e3cdd 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 5e69a72fcc..45c509b265 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 42c7c73096..c8d545383b 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 55a6dc5db3..9b764f5274 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 d8254d3bad..8c36c08a3d 100644
--- a/pulsar-client-cpp/python/pulsar/__init__.py
+++ b/pulsar-client-cpp/python/pulsar/__init__.py
@@ -357,7 +357,8 @@ def subscribe(self, topic, subscription_name,
                   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 @@ def my_listener(consumer, message):
         _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 e2f3e251bc..31163f845a 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 @@ def test_reader_argument_errors(self):
         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 f03c883730..9149addc6f 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 0000000000..cede638ec0
--- /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


 

----------------------------------------------------------------
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