From 62645e51813461cf53c35d26ee20e93ae99fedc6 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Mon, 3 Aug 2026 14:07:03 +0300 Subject: [PATCH] Update Pulsar client from `4.1.2` to `4.2.4` Pulsar 4.2 added the `DecryptFailListener` feature and changed how `cryptoFailureAction` is defaulted. `ConsumerConfigurationData` no longer initializes the field to `ConsumerCryptoFailureAction.FAIL`; it now defaults to `null`, and `ConsumerBuilderImpl.subscribeAsync` applies `FAIL` at subscribe time only when neither `cryptoFailureAction` nor `decryptFailListener` is set. `AdaptedReactiveMessageConsumerTests` stubs `subscribeAsync` with an expected `ConsumerConfigurationData` built by hand, so the expected conf held `null` while the conf passed by the builder held `FAIL`. The stub no longer matched and the tests failed with "didn't match expected consumer conf". Set the expected default explicitly in `keySharedPolicy` and `topicsPattern`. No production change is needed: the adapter correctly leaves `cryptoFailureAction` unset when the spec does not specify it, which lets Pulsar apply its own default and keeps `decryptFailListener` usable. Signed-off-by: Lari Hotari --- gradle/libs.versions.toml | 2 +- .../adapter/AdaptedReactiveMessageConsumerTests.java | 6 ++++++ 2 files changed, 7 insertions(+), 1 deletion(-) diff --git a/gradle/libs.versions.toml b/gradle/libs.versions.toml index 8e44fe94..ecad3a88 100644 --- a/gradle/libs.versions.toml +++ b/gradle/libs.versions.toml @@ -29,7 +29,7 @@ junit-jupiter = "5.11.4" licenser = "0.6.1" log4j = "2.24.3" mockito = "5.18.0" -pulsar = "4.1.2" +pulsar = "4.2.4" rat-gradle = "0.8.0" reactor = "3.7.14" slf4j = "2.0.17" diff --git a/pulsar-client-reactive-adapter/src/test/java/org/apache/pulsar/reactive/client/internal/adapter/AdaptedReactiveMessageConsumerTests.java b/pulsar-client-reactive-adapter/src/test/java/org/apache/pulsar/reactive/client/internal/adapter/AdaptedReactiveMessageConsumerTests.java index b6c1afec..65fe5a53 100644 --- a/pulsar-client-reactive-adapter/src/test/java/org/apache/pulsar/reactive/client/internal/adapter/AdaptedReactiveMessageConsumerTests.java +++ b/pulsar-client-reactive-adapter/src/test/java/org/apache/pulsar/reactive/client/internal/adapter/AdaptedReactiveMessageConsumerTests.java @@ -193,6 +193,9 @@ void keySharedPolicy() throws Exception { expectedConsumerConf.setSubscriptionName("my-sub"); expectedConsumerConf.setSubscriptionType(SubscriptionType.Key_Shared); expectedConsumerConf.setKeySharedPolicy(keySharedPolicy); + // ConsumerBuilderImpl applies this default at subscribe time when neither + // cryptoFailureAction nor decryptFailListener is configured + expectedConsumerConf.setCryptoFailureAction(ConsumerCryptoFailureAction.FAIL); CompletableFuture failedConsumer = new CompletableFuture<>(); failedConsumer.completeExceptionally(new RuntimeException("didn't match expected consumer conf")); @@ -229,6 +232,9 @@ void topicsPattern() throws Exception { expectedConsumerConf.setTopicsPattern(topicsPattern); expectedConsumerConf.setRegexSubscriptionMode(RegexSubscriptionMode.AllTopics); expectedConsumerConf.setPatternAutoDiscoveryPeriod(1); + // ConsumerBuilderImpl applies this default at subscribe time when neither + // cryptoFailureAction nor decryptFailListener is configured + expectedConsumerConf.setCryptoFailureAction(ConsumerCryptoFailureAction.FAIL); CompletableFuture failedConsumer = new CompletableFuture<>(); failedConsumer.completeExceptionally(new RuntimeException("didn't match expected consumer conf"));