From 4f9a40dde2b38772f9f05054db8965b517455350 Mon Sep 17 00:00:00 2001 From: komalatammal Date: Mon, 27 Jun 2022 21:51:56 -0400 Subject: [PATCH 1/7] python cc binding for getLastMessageId --- pulsar-client-cpp/include/pulsar/c/consumer.h | 2 ++ pulsar-client-cpp/lib/c/c_Consumer.cc | 4 ++++ pulsar-client-cpp/python/src/consumer.cc | 13 ++++++++++++- 3 files changed, 18 insertions(+), 1 deletion(-) diff --git a/pulsar-client-cpp/include/pulsar/c/consumer.h b/pulsar-client-cpp/include/pulsar/c/consumer.h index 03f80f3239479..a239376092373 100644 --- a/pulsar-client-cpp/include/pulsar/c/consumer.h +++ b/pulsar-client-cpp/include/pulsar/c/consumer.h @@ -236,6 +236,8 @@ PULSAR_PUBLIC pulsar_result pulsar_consumer_seek(pulsar_consumer_t *consumer, pu PULSAR_PUBLIC int pulsar_consumer_is_connected(pulsar_consumer_t *consumer); +PULSAR_PUBLIC pulsar_result pulsar_consumer_get_last_message_id(pulsar_consumer_t *consumer, pulsar_message_id_t *messageId); + #ifdef __cplusplus } #endif diff --git a/pulsar-client-cpp/lib/c/c_Consumer.cc b/pulsar-client-cpp/lib/c/c_Consumer.cc index 9917e8cfad6e3..f6fb953d97496 100644 --- a/pulsar-client-cpp/lib/c/c_Consumer.cc +++ b/pulsar-client-cpp/lib/c/c_Consumer.cc @@ -143,3 +143,7 @@ pulsar_result pulsar_consumer_seek(pulsar_consumer_t *consumer, pulsar_message_i } int pulsar_consumer_is_connected(pulsar_consumer_t *consumer) { return consumer->consumer.isConnected(); } + +pulsar_result pulsar_consumer_get_last_message_id(pulsar_consumer_t *consumer, pulsar_message_id_t *messageId) { + return (pulsar_result)consumer->consumer.getLastMessageId(messageId->messageId); +} diff --git a/pulsar-client-cpp/python/src/consumer.cc b/pulsar-client-cpp/python/src/consumer.cc index 10ffd07496f78..811ceb3ddf553 100644 --- a/pulsar-client-cpp/python/src/consumer.cc +++ b/pulsar-client-cpp/python/src/consumer.cc @@ -83,6 +83,16 @@ void Consumer_seek_timestamp(Consumer& consumer, uint64_t timestamp) { bool Consumer_is_connected(Consumer& consumer) { return consumer.isConnected(); } +MessageId Consumer_get_last_message_id(Consumer& consumer) { + MessageId msgId; + Result res; + Py_BEGIN_ALLOW_THREADS res = consumer.getLastMessageId(msgId); + Py_END_ALLOW_THREADS + + CHECK_RESULT(res); + return msgId; +} + void export_consumer() { using namespace boost::python; @@ -105,5 +115,6 @@ void export_consumer() { .def("redeliver_unacknowledged_messages", &Consumer::redeliverUnacknowledgedMessages) .def("seek", &Consumer_seek) .def("seek", &Consumer_seek_timestamp) - .def("is_connected", &Consumer_is_connected); + .def("is_connected", &Consumer_is_connected) + .def("get_last_message_id", &Consumer_get_last_message_id); } From 6abf17ad03b205716aecb897d0f41d8cf8237059 Mon Sep 17 00:00:00 2001 From: komalatammal Date: Mon, 27 Jun 2022 22:21:20 -0400 Subject: [PATCH 2/7] add python Consumer class method and doc --- pulsar-client-cpp/python/pulsar/__init__.py | 7 ++++++- 1 file changed, 6 insertions(+), 1 deletion(-) diff --git a/pulsar-client-cpp/python/pulsar/__init__.py b/pulsar-client-cpp/python/pulsar/__init__.py index e79955b57dd79..3832c3e69e23f 100644 --- a/pulsar-client-cpp/python/pulsar/__init__.py +++ b/pulsar-client-cpp/python/pulsar/__init__.py @@ -1253,7 +1253,12 @@ def is_connected(self): Check if the consumer is connected or not. """ return self._consumer.is_connected() - + + def get_last_message_id(self): + """ + Get the last message id. + """ + return self._consumer.get_last_message_id() class Reader: From 22ee6d65aef1d8b3de9e0549f1d5ddccfddb7c88 Mon Sep 17 00:00:00 2001 From: komalatammal Date: Mon, 27 Jun 2022 22:55:26 -0400 Subject: [PATCH 3/7] fix linter issues based on clang-format --- pulsar-client-cpp/README.md | 2 +- pulsar-client-cpp/include/pulsar/c/consumer.h | 3 ++- pulsar-client-cpp/lib/c/c_Consumer.cc | 3 ++- pulsar-client-cpp/lib/checksum/crc32c_arm.h | 10 +++++----- 4 files changed, 10 insertions(+), 8 deletions(-) diff --git a/pulsar-client-cpp/README.md b/pulsar-client-cpp/README.md index f5265577e458e..d83e6cbe0ad03 100644 --- a/pulsar-client-cpp/README.md +++ b/pulsar-client-cpp/README.md @@ -280,7 +280,7 @@ ${PULSAR_PATH}/pulsar-test-service-stop.sh ## Requirements for Contributors -It's recommended to install [LLVM](https://llvm.org/builds/) for `clang-tidy` and `clang-format`. Pulsar C++ client use `clang-format` 6.0+ to format files. +It's required to install [LLVM](https://llvm.org/builds/) for `clang-tidy` and `clang-format`. Pulsar C++ client use `clang-format` 6.0+ to format files. `make format` will automatically format the files. Use `pulsar-client-cpp/docker-format.sh` to ensure the C++ sources are correctly formatted. diff --git a/pulsar-client-cpp/include/pulsar/c/consumer.h b/pulsar-client-cpp/include/pulsar/c/consumer.h index a239376092373..37fd2acf5b8a6 100644 --- a/pulsar-client-cpp/include/pulsar/c/consumer.h +++ b/pulsar-client-cpp/include/pulsar/c/consumer.h @@ -236,7 +236,8 @@ PULSAR_PUBLIC pulsar_result pulsar_consumer_seek(pulsar_consumer_t *consumer, pu PULSAR_PUBLIC int pulsar_consumer_is_connected(pulsar_consumer_t *consumer); -PULSAR_PUBLIC pulsar_result pulsar_consumer_get_last_message_id(pulsar_consumer_t *consumer, pulsar_message_id_t *messageId); +PULSAR_PUBLIC pulsar_result pulsar_consumer_get_last_message_id(pulsar_consumer_t *consumer, + pulsar_message_id_t *messageId); #ifdef __cplusplus } diff --git a/pulsar-client-cpp/lib/c/c_Consumer.cc b/pulsar-client-cpp/lib/c/c_Consumer.cc index f6fb953d97496..00d8311f13278 100644 --- a/pulsar-client-cpp/lib/c/c_Consumer.cc +++ b/pulsar-client-cpp/lib/c/c_Consumer.cc @@ -144,6 +144,7 @@ pulsar_result pulsar_consumer_seek(pulsar_consumer_t *consumer, pulsar_message_i int pulsar_consumer_is_connected(pulsar_consumer_t *consumer) { return consumer->consumer.isConnected(); } -pulsar_result pulsar_consumer_get_last_message_id(pulsar_consumer_t *consumer, pulsar_message_id_t *messageId) { +pulsar_result pulsar_consumer_get_last_message_id(pulsar_consumer_t *consumer, + pulsar_message_id_t *messageId) { return (pulsar_result)consumer->consumer.getLastMessageId(messageId->messageId); } diff --git a/pulsar-client-cpp/lib/checksum/crc32c_arm.h b/pulsar-client-cpp/lib/checksum/crc32c_arm.h index 862215288394c..4848fc04c18fe 100644 --- a/pulsar-client-cpp/lib/checksum/crc32c_arm.h +++ b/pulsar-client-cpp/lib/checksum/crc32c_arm.h @@ -37,11 +37,11 @@ #define crc32c_u16(crc, v) __crc32ch(crc, v) #define crc32c_u32(crc, v) __crc32cw(crc, v) #define crc32c_u64(crc, v) __crc32cd(crc, v) -#define PREF4X64L1(buffer, PREF_OFFSET, ITR) \ - __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [ c ] "I"((PREF_OFFSET) + ((ITR) + 0) * 64)); \ - __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [ c ] "I"((PREF_OFFSET) + ((ITR) + 1) * 64)); \ - __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [ c ] "I"((PREF_OFFSET) + ((ITR) + 2) * 64)); \ - __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [ c ] "I"((PREF_OFFSET) + ((ITR) + 3) * 64)); +#define PREF4X64L1(buffer, PREF_OFFSET, ITR) \ + __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [c] "I"((PREF_OFFSET) + ((ITR) + 0) * 64)); \ + __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [c] "I"((PREF_OFFSET) + ((ITR) + 1) * 64)); \ + __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [c] "I"((PREF_OFFSET) + ((ITR) + 2) * 64)); \ + __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [c] "I"((PREF_OFFSET) + ((ITR) + 3) * 64)); #define PREF1KL1(buffer, PREF_OFFSET) \ PREF4X64L1(buffer, (PREF_OFFSET), 0) \ From 76bd90fe0f63ccda78475a6e56a4fe52c62d0cf1 Mon Sep 17 00:00:00 2001 From: komalatammal Date: Mon, 27 Jun 2022 23:20:07 -0400 Subject: [PATCH 4/7] ubuntu linter fix --- pulsar-client-cpp/lib/checksum/crc32c_arm.h | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/pulsar-client-cpp/lib/checksum/crc32c_arm.h b/pulsar-client-cpp/lib/checksum/crc32c_arm.h index 4848fc04c18fe..862215288394c 100644 --- a/pulsar-client-cpp/lib/checksum/crc32c_arm.h +++ b/pulsar-client-cpp/lib/checksum/crc32c_arm.h @@ -37,11 +37,11 @@ #define crc32c_u16(crc, v) __crc32ch(crc, v) #define crc32c_u32(crc, v) __crc32cw(crc, v) #define crc32c_u64(crc, v) __crc32cd(crc, v) -#define PREF4X64L1(buffer, PREF_OFFSET, ITR) \ - __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [c] "I"((PREF_OFFSET) + ((ITR) + 0) * 64)); \ - __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [c] "I"((PREF_OFFSET) + ((ITR) + 1) * 64)); \ - __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [c] "I"((PREF_OFFSET) + ((ITR) + 2) * 64)); \ - __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [c] "I"((PREF_OFFSET) + ((ITR) + 3) * 64)); +#define PREF4X64L1(buffer, PREF_OFFSET, ITR) \ + __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [ c ] "I"((PREF_OFFSET) + ((ITR) + 0) * 64)); \ + __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [ c ] "I"((PREF_OFFSET) + ((ITR) + 1) * 64)); \ + __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [ c ] "I"((PREF_OFFSET) + ((ITR) + 2) * 64)); \ + __asm__("PRFM PLDL1KEEP, [%x[v],%[c]]" ::[v] "r"(buffer), [ c ] "I"((PREF_OFFSET) + ((ITR) + 3) * 64)); #define PREF1KL1(buffer, PREF_OFFSET) \ PREF4X64L1(buffer, (PREF_OFFSET), 0) \ From 32f1ba6647a5c4d618fb6b467d922c7c3aafb532 Mon Sep 17 00:00:00 2001 From: komalatammal Date: Tue, 28 Jun 2022 17:50:22 -0400 Subject: [PATCH 5/7] try run unit test in ci --- pulsar-client-cpp/python/pulsar_test.py | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/pulsar-client-cpp/python/pulsar_test.py b/pulsar-client-cpp/python/pulsar_test.py index dbdd6be59c7a6..127ecc4247cca 100755 --- a/pulsar-client-cpp/python/pulsar_test.py +++ b/pulsar-client-cpp/python/pulsar_test.py @@ -753,6 +753,18 @@ def test_reader_argument_errors(self): self._check_value_error(lambda: client.create_reader(topic, MessageId.earliest, reader_name=5)) client.close() + def test_get_last_message_id(self): + client = Client(self.serviceUrl) + consumer = client.subscribe( + "persistent://public/default/topic_name_test", "topic_name_test_sub", consumer_type=ConsumerType.Shared + ) + producer = client.create_producer("persistent://public/default/topic_name_test") + msg_id = producer.send(b"hello") + + msg = consumer.receive(TM) + self.assertEqual(msg.message_id(), msg_id) + client.close() + def test_publish_compact_and_consume(self): client = Client(self.serviceUrl) topic = "compaction_%s" % (uuid.uuid4()) From 366a16542a7121f331e612ca2ed16464eeeaf453 Mon Sep 17 00:00:00 2001 From: komalatammal Date: Tue, 28 Jun 2022 21:15:52 -0400 Subject: [PATCH 6/7] fix doc comment --- pulsar-client-cpp/README.md | 2 +- pulsar-client-cpp/python/pulsar_test.py | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-client-cpp/README.md b/pulsar-client-cpp/README.md index d83e6cbe0ad03..2af50a0ffafbe 100644 --- a/pulsar-client-cpp/README.md +++ b/pulsar-client-cpp/README.md @@ -280,7 +280,7 @@ ${PULSAR_PATH}/pulsar-test-service-stop.sh ## Requirements for Contributors -It's required to install [LLVM](https://llvm.org/builds/) for `clang-tidy` and `clang-format`. Pulsar C++ client use `clang-format` 6.0+ to format files. `make format` will automatically format the files. +It's required to install [LLVM](https://llvm.org/builds/) for `clang-tidy` and `clang-format`. Pulsar C++ client use `clang-format` 6.0+ to format files. `make format` automatically formats the files. Use `pulsar-client-cpp/docker-format.sh` to ensure the C++ sources are correctly formatted. diff --git a/pulsar-client-cpp/python/pulsar_test.py b/pulsar-client-cpp/python/pulsar_test.py index 127ecc4247cca..007c859a0463f 100755 --- a/pulsar-client-cpp/python/pulsar_test.py +++ b/pulsar-client-cpp/python/pulsar_test.py @@ -762,6 +762,7 @@ def test_get_last_message_id(self): msg_id = producer.send(b"hello") msg = consumer.receive(TM) + msg_id = consumer.get_last_message_id() self.assertEqual(msg.message_id(), msg_id) client.close() From 12ab0cbe48ddbcea82942619aab4368468a8d9f1 Mon Sep 17 00:00:00 2001 From: komalatammal Date: Wed, 29 Jun 2022 13:15:51 -0400 Subject: [PATCH 7/7] test the test case can be ran --- pulsar-client-cpp/python/pulsar_test.py | 1 - 1 file changed, 1 deletion(-) diff --git a/pulsar-client-cpp/python/pulsar_test.py b/pulsar-client-cpp/python/pulsar_test.py index 007c859a0463f..127ecc4247cca 100755 --- a/pulsar-client-cpp/python/pulsar_test.py +++ b/pulsar-client-cpp/python/pulsar_test.py @@ -762,7 +762,6 @@ def test_get_last_message_id(self): msg_id = producer.send(b"hello") msg = consumer.receive(TM) - msg_id = consumer.get_last_message_id() self.assertEqual(msg.message_id(), msg_id) client.close()