From 27b09d3b4912cfd9fe7ec3a6b8090f297c65490f Mon Sep 17 00:00:00 2001 From: Christophe Bornet Date: Fri, 18 Nov 2022 15:31:03 +0100 Subject: [PATCH] Add test for reactive consumer properties --- .../PulsarReactivePropertiesTests.java | 120 ++++++++++++++++++ 1 file changed, 120 insertions(+) diff --git a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactivePropertiesTests.java b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactivePropertiesTests.java index 53501689..0613f76d 100644 --- a/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactivePropertiesTests.java +++ b/spring-pulsar-spring-boot-autoconfigure/src/test/java/org/springframework/pulsar/autoconfigure/PulsarReactivePropertiesTests.java @@ -19,22 +19,34 @@ package org.springframework.pulsar.autoconfigure; import static org.assertj.core.api.Assertions.assertThat; import java.time.Duration; +import java.util.Collections; import java.util.HashMap; import java.util.Map; import org.apache.pulsar.client.api.CompressionType; +import org.apache.pulsar.client.api.ConsumerCryptoFailureAction; import org.apache.pulsar.client.api.HashingScheme; +import org.apache.pulsar.client.api.KeySharedMode; import org.apache.pulsar.client.api.MessageRoutingMode; import org.apache.pulsar.client.api.ProducerAccessMode; import org.apache.pulsar.client.api.ProducerCryptoFailureAction; +import org.apache.pulsar.client.api.RegexSubscriptionMode; +import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.reactive.client.api.ReactiveMessageConsumerSpec; import org.apache.pulsar.reactive.client.api.ReactiveMessageSenderSpec; import org.junit.jupiter.api.Nested; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.EnumSource; +import org.junit.jupiter.params.provider.EnumSource.Mode; import org.springframework.boot.context.properties.bind.Bindable; import org.springframework.boot.context.properties.bind.Binder; import org.springframework.boot.context.properties.source.ConfigurationPropertySource; import org.springframework.boot.context.properties.source.MapConfigurationPropertySource; +import org.springframework.pulsar.autoconfigure.PulsarReactiveProperties.SchedulerType; + +import reactor.core.scheduler.Schedulers; /** * Unit tests for {@link PulsarReactiveProperties}. @@ -45,6 +57,10 @@ public class PulsarReactivePropertiesTests { private final PulsarReactiveProperties properties = new PulsarReactiveProperties(); + private void bind(String name, String value) { + bind(Collections.singletonMap(name, value)); + } + private void bind(Map map) { ConfigurationPropertySource source = new MapConfigurationPropertySource(map); new Binder(source).bind("spring.pulsar.reactive", Bindable.ofInstance(this.properties)); @@ -92,4 +108,108 @@ public class PulsarReactivePropertiesTests { } + @Nested + class ConsumerPropertiesTests { + + @Test + void consumerPropsToConsumerSpec() { + Map props = new HashMap<>(); + props.put("spring.pulsar.reactive.consumer.topics[0]", "my-topic"); + props.put("spring.pulsar.reactive.consumer.topics-pattern", "my-pattern"); + props.put("spring.pulsar.reactive.consumer.subscription-name", "my-subscription"); + props.put("spring.pulsar.reactive.consumer.subscription-type", "Shared"); + props.put("spring.pulsar.reactive.consumer.subscription-properties[my-sub-prop]", "my-sub-prop-value"); + props.put("spring.pulsar.reactive.consumer.receiver-queue-size", "1"); + props.put("spring.pulsar.reactive.consumer.acknowledgements-group-time", "2s"); + props.put("spring.pulsar.reactive.consumer.acknowledge-asynchronously", "false"); + props.put("spring.pulsar.reactive.consumer.negative-ack-redelivery-delay", "3s"); + props.put("spring.pulsar.reactive.consumer.dead-letter-policy.max-redeliver-count", "4"); + props.put("spring.pulsar.reactive.consumer.dead-letter-policy.retry-letter-topic", "my-retry-topic"); + props.put("spring.pulsar.reactive.consumer.dead-letter-policy.dead-letter-topic", "my-dlt-topic"); + props.put("spring.pulsar.reactive.consumer.dead-letter-policy.initial-subscription-name", + "my-initial-subscription"); + props.put("spring.pulsar.reactive.consumer.max-total-receiver-queue-size-across-partitions", "5"); + props.put("spring.pulsar.reactive.consumer.consumer-name", "my-consumer"); + props.put("spring.pulsar.reactive.consumer.ack-timeout", "6s"); + props.put("spring.pulsar.reactive.consumer.ack-timeout-tick-time", "7s"); + props.put("spring.pulsar.reactive.consumer.priority-level", "8"); + props.put("spring.pulsar.reactive.consumer.crypto-failure-action", "DISCARD"); + props.put("spring.pulsar.reactive.consumer.properties[my-prop]", "my-prop-value"); + props.put("spring.pulsar.reactive.consumer.read-compacted", "true"); + props.put("spring.pulsar.reactive.consumer.batch-index-ack-enabled", "true"); + props.put("spring.pulsar.reactive.consumer.topics-pattern-auto-discovery-period", "9s"); + props.put("spring.pulsar.reactive.consumer.topics-pattern-subscription-mode", "AllTopics"); + props.put("spring.pulsar.reactive.consumer.auto-update-partitions", "false"); + props.put("spring.pulsar.reactive.consumer.auto-update-partitions-interval", "10s"); + props.put("spring.pulsar.reactive.consumer.replicate-subscription-state", "true"); + props.put("spring.pulsar.reactive.consumer.auto-ack-oldest-chunked-message-on-queue-full", "false"); + props.put("spring.pulsar.reactive.consumer.max-pending-chunked-message", "11"); + props.put("spring.pulsar.reactive.consumer.expire-time-of-incomplete-chunked-message", "12s"); + + bind(props); + ReactiveMessageConsumerSpec consumerSpec = properties.buildReactiveMessageConsumerSpec(); + + assertThat(consumerSpec.getTopicNames()).containsExactly("my-topic"); + assertThat(consumerSpec.getTopicsPattern().toString()).isEqualTo("my-pattern"); + assertThat(consumerSpec.getSubscriptionName()).isEqualTo("my-subscription"); + assertThat(consumerSpec.getSubscriptionType()).isEqualTo(SubscriptionType.Shared); + assertThat(consumerSpec.getSubscriptionProperties()).hasSize(1).containsEntry("my-sub-prop", + "my-sub-prop-value"); + assertThat(consumerSpec.getReceiverQueueSize()).isEqualTo(1); + assertThat(consumerSpec.getAcknowledgementsGroupTime()).isEqualTo(Duration.ofSeconds(2)); + assertThat(consumerSpec.getAcknowledgeAsynchronously()).isFalse(); + assertThat(consumerSpec.getNegativeAckRedeliveryDelay()).isEqualTo(Duration.ofSeconds(3)); + assertThat(consumerSpec.getDeadLetterPolicy().getMaxRedeliverCount()).isEqualTo(4); + assertThat(consumerSpec.getDeadLetterPolicy().getRetryLetterTopic()).isEqualTo("my-retry-topic"); + assertThat(consumerSpec.getDeadLetterPolicy().getDeadLetterTopic()).isEqualTo("my-dlt-topic"); + assertThat(consumerSpec.getDeadLetterPolicy().getInitialSubscriptionName()) + .isEqualTo("my-initial-subscription"); + assertThat(consumerSpec.getMaxTotalReceiverQueueSizeAcrossPartitions()).isEqualTo(5); + assertThat(consumerSpec.getConsumerName()).isEqualTo("my-consumer"); + assertThat(consumerSpec.getAckTimeout()).isEqualTo(Duration.ofSeconds(6)); + assertThat(consumerSpec.getAckTimeoutTickTime()).isEqualTo(Duration.ofSeconds(7)); + assertThat(consumerSpec.getPriorityLevel()).isEqualTo(8); + assertThat(consumerSpec.getCryptoFailureAction()).isEqualTo(ConsumerCryptoFailureAction.DISCARD); + assertThat(consumerSpec.getProperties()).hasSize(1).containsEntry("my-prop", "my-prop-value"); + assertThat(consumerSpec.getReadCompacted()).isTrue(); + assertThat(consumerSpec.getBatchIndexAckEnabled()).isTrue(); + assertThat(consumerSpec.getTopicsPatternAutoDiscoveryPeriod()).isEqualTo(Duration.ofSeconds(9)); + assertThat(consumerSpec.getTopicsPatternSubscriptionMode()).isEqualTo(RegexSubscriptionMode.AllTopics); + assertThat(consumerSpec.getAutoUpdatePartitions()).isFalse(); + assertThat(consumerSpec.getAutoUpdatePartitionsInterval()).isEqualTo(Duration.ofSeconds(10)); + assertThat(consumerSpec.getReplicateSubscriptionState()).isTrue(); + assertThat(consumerSpec.getAutoAckOldestChunkedMessageOnQueueFull()).isFalse(); + assertThat(consumerSpec.getMaxPendingChunkedMessage()).isEqualTo(11); + assertThat(consumerSpec.getExpireTimeOfIncompleteChunkedMessage()).isEqualTo(Duration.ofSeconds(12)); + } + + @ParameterizedTest + @EnumSource(KeySharedMode.class) + void keySharedModeProperty(KeySharedMode keySharedMode) { + bind("spring.pulsar.reactive.consumer.key-shared-mode", keySharedMode.name()); + ReactiveMessageConsumerSpec consumerSpec = properties.buildReactiveMessageConsumerSpec(); + + assertThat(consumerSpec.getKeySharedPolicy().getKeySharedMode()).isEqualTo(keySharedMode); + } + + @ParameterizedTest + @EnumSource(value = SchedulerType.class, names = "immediate", mode = Mode.EXCLUDE) + void acknowledgeScheduler(SchedulerType acknowledgeSchedulerType) { + bind("spring.pulsar.reactive.consumer.acknowledge-scheduler-type", acknowledgeSchedulerType.name()); + ReactiveMessageConsumerSpec consumerSpec = properties.buildReactiveMessageConsumerSpec(); + + assertThat(consumerSpec.getAcknowledgeScheduler().toString()) + .isEqualTo("Schedulers.%s()".formatted(acknowledgeSchedulerType)); + } + + @Test + void acknowledgeSchedulerImmediate() { + bind("spring.pulsar.reactive.consumer.acknowledge-scheduler-type", "immediate"); + ReactiveMessageConsumerSpec consumerSpec = properties.buildReactiveMessageConsumerSpec(); + + assertThat(consumerSpec.getAcknowledgeScheduler()).isSameAs(Schedulers.immediate()); + } + + } + }