diff --git a/binders/pom.xml b/binders/pom.xml index bbf7f97af..5178cd931 100644 --- a/binders/pom.xml +++ b/binders/pom.xml @@ -16,6 +16,7 @@ kafka-binder rabbit-binder + pulsar-binder diff --git a/binders/pulsar-binder/pom.xml b/binders/pulsar-binder/pom.xml new file mode 100644 index 000000000..61002d155 --- /dev/null +++ b/binders/pulsar-binder/pom.xml @@ -0,0 +1,170 @@ + + + 4.0.0 + spring-cloud-stream-binder-pulsar-parent + 4.1.0-SNAPSHOT + pom + + org.springframework.cloud + spring-cloud-stream-binders + 4.1.0-SNAPSHOT + + + + 17 + true + true + true + + + spring-cloud-stream-binder-pulsar + + + + + org.junit.vintage + junit-vintage-engine + test + + + + + + + + org.apache.maven.plugins + maven-antrun-plugin + + + + org.apache.maven.plugins + maven-javadoc-plugin + + true + + + + org.apache.maven.plugins + maven-surefire-plugin + 3.0.0-M7 + + true + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + ${maven-compiler-plugin.version} + + ${java.version} + ${java.version} + -parameters + + + + org.apache.maven.plugins + maven-checkstyle-plugin + + + + + + + spring + + + + + coverage + + + env.TRAVIS + true + + + + + + org.jacoco + jacoco-maven-plugin + 0.7.9 + + + agent + + prepare-agent + + + + report + test + + report + + + + + + + + + + + spring-snapshots + Spring Snapshots + https://repo.spring.io/libs-snapshot-local + + + spring-milestones + Spring milestones + https://repo.spring.io/libs-milestone-local + + + rsocket-snapshots + RSocket Snapshots + https://oss.jfrog.org/oss-snapshot-local + + true + + + + spring-releases + Spring Releases + https://repo.spring.io/release + + + + + spring-snapshots + Spring Snapshots + https://repo.spring.io/snapshot + + true + + + + spring-milestones + Spring Milestones + https://repo.spring.io/milestone + + false + + + + spring-releases + Spring Releases + https://repo.spring.io/release + + + + + + org.apache.maven.plugins + maven-checkstyle-plugin + + + + diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/pom.xml b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/pom.xml new file mode 100644 index 000000000..5f473cca7 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/pom.xml @@ -0,0 +1,50 @@ + + + 4.0.0 + + spring-cloud-stream-binder-pulsar + jar + spring-cloud-stream-binder-pulsar + PUlsar binder implementation + + + org.springframework.cloud + spring-cloud-stream-binder-pulsar-parent + 4.1.0-SNAPSHOT + + + + + org.springframework.pulsar + spring-pulsar-spring-boot-starter + 0.2.1-SNAPSHOT + + + org.springframework.cloud + spring-cloud-stream + + + org.springframework.cloud + spring-cloud-stream-test-support + test + + + org.testcontainers + pulsar + 1.17.6 + test + + + org.testcontainers + junit-jupiter + 1.17.6 + test + + + org.awaitility + awaitility + test + + + + diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderHeaderMapper.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderHeaderMapper.java new file mode 100644 index 000000000..32b1bc02e --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderHeaderMapper.java @@ -0,0 +1,70 @@ +/* + * Copyright 2022-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.util.Map; + +import org.apache.pulsar.client.api.Message; + +import org.springframework.cloud.stream.binder.BinderHeaders; +import org.springframework.integration.IntegrationMessageHeaderAccessor; +import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.support.MessageHeaderAccessor; +import org.springframework.pulsar.support.header.PulsarHeaderMapper; + + +/** + * A delegating {@code PulsarHeaderMapper} that ensures the delegate mapper never includes + * internal binder specific headers during outbound mapping. + * + * @author Chris Bono + */ +class PulsarBinderHeaderMapper implements PulsarHeaderMapper { + + private final PulsarHeaderMapper delegate; + + /** + * Construct a mapper with the specified delegate. + * @param delegate the delegate mapper + */ + PulsarBinderHeaderMapper(PulsarHeaderMapper delegate) { + this.delegate = delegate; + } + + @Override + public Map toPulsarHeaders(MessageHeaders springHeaders) { + Map pulsarHeaders = this.delegate.toPulsarHeaders(springHeaders); + pulsarHeaders.remove(MessageHeaders.ID); + pulsarHeaders.remove(MessageHeaders.TIMESTAMP); + pulsarHeaders.remove(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT); + pulsarHeaders.remove(BinderHeaders.NATIVE_HEADERS_PRESENT); + return pulsarHeaders; + } + + @Override + public MessageHeaders toSpringHeaders(Message pulsarMessage) { + var springHeaders = this.delegate.toSpringHeaders(pulsarMessage); + if (!springHeaders.isEmpty()) { + MessageHeaderAccessor mutableHeaders = new MessageHeaderAccessor(); + mutableHeaders.copyHeaders(springHeaders); + mutableHeaders.setHeader(BinderHeaders.NATIVE_HEADERS_PRESENT, Boolean.TRUE); + springHeaders = mutableHeaders.getMessageHeaders(); + } + return springHeaders; + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderUtils.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderUtils.java new file mode 100644 index 000000000..229372d7c --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderUtils.java @@ -0,0 +1,207 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.time.Duration; +import java.util.HashMap; +import java.util.Map; +import java.util.Objects; +import java.util.UUID; + +import org.springframework.boot.context.properties.PropertyMapper; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarConsumerProperties; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.core.log.LogAccessor; +import org.springframework.pulsar.autoconfigure.ConsumerConfigProperties; +import org.springframework.pulsar.autoconfigure.ProducerConfigProperties; +import org.springframework.util.StringUtils; +import org.springframework.util.unit.DataSize; + +/** + * Binder utility methods. + * + * @author Soby Chacko + * @author Chris Bono + */ +final class PulsarBinderUtils { + + private static final LogAccessor LOGGER = new LogAccessor(PulsarBinderUtils.class); + + private static final String SUBSCRIPTION_NAME_FORMAT_STR = "%s-anon-subscription-%s"; + + private PulsarBinderUtils() { + } + + /** + * Gets the subscription name to use for the binder. + * @param consumerProps the pulsar consumer props + * @param consumerDestination the destination being subscribed to + * @return the subscription name from the consumer properties or a generated name in + * the format {@link #SUBSCRIPTION_NAME_FORMAT_STR} when the name is not set on the + * consumer properties + */ + static String subscriptionName(PulsarConsumerProperties consumerProps, ConsumerDestination consumerDestination) { + if (StringUtils.hasText(consumerProps.getSubscriptionName())) { + return consumerProps.getSubscriptionName(); + } + return SUBSCRIPTION_NAME_FORMAT_STR.formatted(consumerDestination.getName(), UUID.randomUUID()); + } + + /** + * Merges properties defined at the binder and binding level (binding properties + * override binder properties). + *

+ * NOTE: Properties whose value is not different from the default value in the + * {@code baseProps} are not included in the merged result. + * @param baseProps the map of base level properties (eg. 'spring.pulsar.consumer.*') + * @param binderProps the map of binder level properties (eg. + * 'spring.cloud.stream.pulsar.binder.consumer.*') + * @param bindingProps the map of binding level properties (eg. + * 'spring.cloud.stream.pulsar.bindings.myBinding-in-0.consumer.*') + * @return map of merged binder and binding properties including only properties whose + * value has changed from the same property in the base properties + */ + static Map mergePropertiesWithPrecedence(Map baseProps, + Map binderProps, Map bindingProps) { + Objects.requireNonNull(baseProps, "baseProps must be specified"); + Objects.requireNonNull(binderProps, "binderProps must be specified"); + Objects.requireNonNull(bindingProps, "bindingProps must be specified"); + + Map newOrModifiedBinderProps = extractNewOrModifiedProperties(binderProps, baseProps); + LOGGER.trace(() -> "New or modified binder props: %s".formatted(newOrModifiedBinderProps)); + + Map newOrModifiedBindingProps = extractNewOrModifiedProperties(bindingProps, baseProps); + LOGGER.trace(() -> "New or modified binding props: %s".formatted(newOrModifiedBindingProps)); + + Map mergedProps = new HashMap<>(newOrModifiedBinderProps); + mergedProps.putAll(newOrModifiedBindingProps); + LOGGER.trace(() -> "Final merged props: %s".formatted(mergedProps)); + + return mergedProps; + } + + private static Map extractNewOrModifiedProperties(Map candidateProps, + Map baseProps) { + Map newOrModifiedProps = new HashMap<>(); + candidateProps.forEach((propName, propValue) -> { + if (!baseProps.containsKey(propName) || (!Objects.equals(propValue, baseProps.get(propName)))) { + newOrModifiedProps.put(propName, propValue); + } + }); + return newOrModifiedProps; + } + + /** + * Gets a map representation of a {@link ProducerConfigProperties}. + * @param producerProps the producer props + * @return map representation of producer props where each entry is a field and its + * associated value + */ + static Map convertProducerPropertiesToMap(ProducerConfigProperties producerProps) { + var properties = new PulsarBinderUtils.Properties(); + var map = PropertyMapper.get().alwaysApplyingWhenNonNull(); + map.from(producerProps::getTopicName).to(properties.in("topicName")); + map.from(producerProps::getProducerName).to(properties.in("producerName")); + map.from(producerProps::getSendTimeout).asInt(Duration::toMillis).to(properties.in("sendTimeoutMs")); + map.from(producerProps::getBlockIfQueueFull).to(properties.in("blockIfQueueFull")); + map.from(producerProps::getMaxPendingMessages).to(properties.in("maxPendingMessages")); + map.from(producerProps::getMaxPendingMessagesAcrossPartitions) + .to(properties.in("maxPendingMessagesAcrossPartitions")); + map.from(producerProps::getMessageRoutingMode).to(properties.in("messageRoutingMode")); + map.from(producerProps::getHashingScheme).to(properties.in("hashingScheme")); + map.from(producerProps::getCryptoFailureAction).to(properties.in("cryptoFailureAction")); + map.from(producerProps::getBatchingMaxPublishDelay).as(it -> it.toNanos() / 1000) + .to(properties.in("batchingMaxPublishDelayMicros")); + map.from(producerProps::getBatchingPartitionSwitchFrequencyByPublishDelay) + .to(properties.in("batchingPartitionSwitchFrequencyByPublishDelay")); + map.from(producerProps::getBatchingMaxMessages).to(properties.in("batchingMaxMessages")); + map.from(producerProps::getBatchingMaxBytes).asInt(DataSize::toBytes).to(properties.in("batchingMaxBytes")); + map.from(producerProps::getBatchingEnabled).to(properties.in("batchingEnabled")); + map.from(producerProps::getChunkingEnabled).to(properties.in("chunkingEnabled")); + map.from(producerProps::getEncryptionKeys).to(properties.in("encryptionKeys")); + map.from(producerProps::getCompressionType).to(properties.in("compressionType")); + map.from(producerProps::getInitialSequenceId).to(properties.in("initialSequenceId")); + map.from(producerProps::getAutoUpdatePartitions).to(properties.in("autoUpdatePartitions")); + map.from(producerProps::getAutoUpdatePartitionsInterval).as(Duration::toSeconds) + .to(properties.in("autoUpdatePartitionsIntervalSeconds")); + map.from(producerProps::getMultiSchema).to(properties.in("multiSchema")); + map.from(producerProps::getProducerAccessMode).to(properties.in("accessMode")); + map.from(producerProps::getLazyStartPartitionedProducers).to(properties.in("lazyStartPartitionedProducers")); + map.from(producerProps::getProperties).to(properties.in("properties")); + return properties; + } + + /** + * Gets a map representation of a {@link ConsumerConfigProperties}. + * @param consumerProps the consumer props + * @return map representation of consumer props where each entry is a field and its + * associated value + */ + static Map convertConsumerPropertiesToMap(ConsumerConfigProperties consumerProps) { + var properties = new PulsarBinderUtils.Properties(); + var map = PropertyMapper.get().alwaysApplyingWhenNonNull(); + map.from(consumerProps::getTopics).to(properties.in("topicNames")); + map.from(consumerProps::getTopicsPattern).to(properties.in("topicsPattern")); + map.from(consumerProps::getSubscriptionName).to(properties.in("subscriptionName")); + map.from(consumerProps::getSubscriptionType).to(properties.in("subscriptionType")); + map.from(consumerProps::getSubscriptionProperties).to(properties.in("subscriptionProperties")); + map.from(consumerProps::getSubscriptionMode).to(properties.in("subscriptionMode")); + map.from(consumerProps::getReceiverQueueSize).to(properties.in("receiverQueueSize")); + map.from(consumerProps::getAcknowledgementsGroupTime).as(it -> it.toNanos() / 1000) + .to(properties.in("acknowledgementsGroupTimeMicros")); + map.from(consumerProps::getNegativeAckRedeliveryDelay).as(it -> it.toNanos() / 1000) + .to(properties.in("negativeAckRedeliveryDelayMicros")); + map.from(consumerProps::getMaxTotalReceiverQueueSizeAcrossPartitions) + .to(properties.in("maxTotalReceiverQueueSizeAcrossPartitions")); + map.from(consumerProps::getConsumerName).to(properties.in("consumerName")); + map.from(consumerProps::getAckTimeout).as(Duration::toMillis).to(properties.in("ackTimeoutMillis")); + map.from(consumerProps::getTickDuration).as(Duration::toMillis).to(properties.in("tickDurationMillis")); + map.from(consumerProps::getPriorityLevel).to(properties.in("priorityLevel")); + map.from(consumerProps::getCryptoFailureAction).to(properties.in("cryptoFailureAction")); + map.from(consumerProps::getProperties).to(properties.in("properties")); + map.from(consumerProps::getReadCompacted).to(properties.in("readCompacted")); + map.from(consumerProps::getSubscriptionInitialPosition).to(properties.in("subscriptionInitialPosition")); + map.from(consumerProps::getPatternAutoDiscoveryPeriod).to(properties.in("patternAutoDiscoveryPeriod")); + map.from(consumerProps::getRegexSubscriptionMode).to(properties.in("regexSubscriptionMode")); + map.from(consumerProps::getDeadLetterPolicy).to(properties.in("deadLetterPolicy")); + map.from(consumerProps::getRetryEnable).to(properties.in("retryEnable")); + map.from(consumerProps::getAutoUpdatePartitions).to(properties.in("autoUpdatePartitions")); + map.from(consumerProps::getAutoUpdatePartitionsInterval).as(Duration::toSeconds) + .to(properties.in("autoUpdatePartitionsIntervalSeconds")); + map.from(consumerProps::getReplicateSubscriptionState).to(properties.in("replicateSubscriptionState")); + map.from(consumerProps::getResetIncludeHead).to(properties.in("resetIncludeHead")); + map.from(consumerProps::getBatchIndexAckEnabled).to(properties.in("batchIndexAckEnabled")); + map.from(consumerProps::getAckReceiptEnabled).to(properties.in("ackReceiptEnabled")); + map.from(consumerProps::getPoolMessages).to(properties.in("poolMessages")); + map.from(consumerProps::getStartPaused).to(properties.in("startPaused")); + map.from(consumerProps::getAutoAckOldestChunkedMessageOnQueueFull) + .to(properties.in("autoAckOldestChunkedMessageOnQueueFull")); + map.from(consumerProps::getMaxPendingChunkedMessage).to(properties.in("maxPendingChunkedMessage")); + map.from(consumerProps::getExpireTimeOfIncompleteChunkedMessage).as(Duration::toMillis) + .to(properties.in("expireTimeOfIncompleteChunkedMessageMillis")); + return properties; + } + + static class Properties extends HashMap { + + java.util.function.Consumer in(String key) { + return (value) -> put(key, value); + } + + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarMessageChannelBinder.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarMessageChannelBinder.java new file mode 100644 index 000000000..d6eafe2eb --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/PulsarMessageChannelBinder.java @@ -0,0 +1,320 @@ +/* + * Copyright 2022-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.util.Optional; +import java.util.Set; + +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.common.schema.SchemaType; + +import org.springframework.cloud.stream.binder.AbstractMessageChannelBinder; +import org.springframework.cloud.stream.binder.Binder; +import org.springframework.cloud.stream.binder.BinderSpecificPropertiesProvider; +import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; +import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.ExtendedPropertiesBinder; +import org.springframework.cloud.stream.binder.HeaderMode; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarConsumerProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarExtendedBindingProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarProducerProperties; +import org.springframework.cloud.stream.binder.pulsar.provisioning.PulsarTopicProvisioner; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.cloud.stream.provisioning.ProducerDestination; +import org.springframework.integration.core.MessageProducer; +import org.springframework.integration.endpoint.MessageProducerSupport; +import org.springframework.integration.handler.AbstractMessageProducingHandler; +import org.springframework.integration.support.management.ManageableLifecycle; +import org.springframework.lang.Nullable; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHandler; +import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.support.MessageBuilder; +import org.springframework.pulsar.autoconfigure.ConsumerConfigProperties; +import org.springframework.pulsar.autoconfigure.ProducerConfigProperties; +import org.springframework.pulsar.core.ProducerBuilderConfigurationUtil; +import org.springframework.pulsar.core.ProducerBuilderCustomizer; +import org.springframework.pulsar.core.PulsarConsumerFactory; +import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.core.SchemaResolver; +import org.springframework.pulsar.core.TypedMessageBuilderCustomizer; +import org.springframework.pulsar.listener.AbstractPulsarMessageListenerContainer; +import org.springframework.pulsar.listener.DefaultPulsarMessageListenerContainer; +import org.springframework.pulsar.listener.PulsarContainerProperties; +import org.springframework.pulsar.listener.PulsarRecordMessageListener; +import org.springframework.pulsar.support.header.PulsarHeaderMapper; + + +/** + * {@link Binder} implementation for Apache Pulsar. + * + * @author Soby Chacko + * @author Chris Bono + */ +public class PulsarMessageChannelBinder extends + AbstractMessageChannelBinder, ExtendedProducerProperties, PulsarTopicProvisioner> + implements ExtendedPropertiesBinder { + + private final PulsarTemplate pulsarTemplate; + + private final PulsarConsumerFactory pulsarConsumerFactory; + + private final PulsarBinderConfigurationProperties binderConfigProps; + + private final SchemaResolver schemaResolver; + + private final PulsarHeaderMapper headerMapper; + + private PulsarExtendedBindingProperties extendedBindingProperties = new PulsarExtendedBindingProperties(); + + public PulsarMessageChannelBinder(PulsarTopicProvisioner provisioningProvider, + PulsarTemplate pulsarTemplate, PulsarConsumerFactory pulsarConsumerFactory, + PulsarBinderConfigurationProperties binderConfigProps, SchemaResolver schemaResolver, + PulsarHeaderMapper headerMapper) { + super(null, provisioningProvider); + this.pulsarTemplate = pulsarTemplate; + this.pulsarConsumerFactory = pulsarConsumerFactory; + this.binderConfigProps = binderConfigProps; + this.schemaResolver = schemaResolver; + this.headerMapper = headerMapper; + } + + @Override + protected MessageHandler createProducerMessageHandler(ProducerDestination destination, + ExtendedProducerProperties producerProperties, MessageChannel errorChannel) { + final Schema schema; + if (producerProperties.isUseNativeEncoding()) { + var schemaType = Optional.ofNullable(producerProperties.getExtension().getSchemaType()) + .orElse(SchemaType.NONE); + schema = this.schemaResolver + .resolveSchema(schemaType, producerProperties.getExtension().getMessageType(), + producerProperties.getExtension().getMessageKeyType()) + .orElseThrow(() -> "Could not determine producer schema for " + destination.getName()); + } + else { + schema = null; + } + + var baseProducerProps = PulsarBinderUtils.convertProducerPropertiesToMap(new ProducerConfigProperties()); + var binderProducerProps = PulsarBinderUtils + .convertProducerPropertiesToMap(this.binderConfigProps.getProducer()); + var bindingProducerProps = PulsarBinderUtils.convertProducerPropertiesToMap(producerProperties.getExtension()); + var mergedProducerProps = PulsarBinderUtils.mergePropertiesWithPrecedence(baseProducerProps, + binderProducerProps, bindingProducerProps); + + var handler = new PulsarProducerConfigurationMessageHandler(this.pulsarTemplate, schema, destination.getName(), + (builder) -> ProducerBuilderConfigurationUtil.loadConf(builder, mergedProducerProps), + determineOutboundHeaderMapper(producerProperties)); + handler.setApplicationContext(getApplicationContext()); + handler.setBeanFactory(getBeanFactory()); + + return handler; + } + + @Nullable + private PulsarBinderHeaderMapper determineOutboundHeaderMapper( + ExtendedProducerProperties extProducerProps) { + if (HeaderMode.none.equals(extProducerProps.getHeaderMode())) { + return null; + } + return new PulsarBinderHeaderMapper(this.headerMapper); + } + + @Override + protected MessageProducer createConsumerEndpoint(ConsumerDestination destination, String group, + ExtendedConsumerProperties properties) { + var containerProperties = new PulsarContainerProperties(); + containerProperties.setTopics(Set.of(destination.getName())); + + var inboundHeaderMapper = determineInboundHeaderMapper(properties); + + var messageDrivenChannelAdapter = new PulsarMessageDrivenChannelAdapter(); + containerProperties.setMessageListener((PulsarRecordMessageListener) (consumer, pulsarMsg) -> { + var springMessage = (inboundHeaderMapper != null) + ? MessageBuilder.createMessage(pulsarMsg.getValue(), inboundHeaderMapper.toSpringHeaders(pulsarMsg)) + : MessageBuilder.withPayload(pulsarMsg.getValue()).build(); + messageDrivenChannelAdapter.send(springMessage); + }); + + if (properties.isUseNativeDecoding()) { + var schemaType = Optional.ofNullable(properties.getExtension().getSchemaType()).orElse(SchemaType.NONE); + var schema = this.schemaResolver + .resolveSchema(schemaType, properties.getExtension().getMessageType(), + properties.getExtension().getMessageKeyType()) + .orElseThrow(() -> "Could not determine consumer schema for " + destination.getName()); + containerProperties.setSchema(schema); + } + else { + containerProperties.setSchema(Schema.BYTES); + } + var subscriptionName = PulsarBinderUtils.subscriptionName(properties.getExtension(), destination); + containerProperties.setSubscriptionName(subscriptionName); + + var baseConsumerProps = PulsarBinderUtils.convertConsumerPropertiesToMap(new ConsumerConfigProperties()); + var binderConsumerProps = PulsarBinderUtils + .convertConsumerPropertiesToMap(this.binderConfigProps.getConsumer()); + var bindingConsumerProps = PulsarBinderUtils.convertConsumerPropertiesToMap(properties.getExtension()); + var mergedConsumerProps = PulsarBinderUtils.mergePropertiesWithPrecedence(baseConsumerProps, + binderConsumerProps, bindingConsumerProps); + containerProperties.getPulsarConsumerProperties().putAll(mergedConsumerProps); + containerProperties.updateContainerProperties(); + + var container = new DefaultPulsarMessageListenerContainer<>(this.pulsarConsumerFactory, containerProperties); + messageDrivenChannelAdapter.setMessageListenerContainer(container); + + return messageDrivenChannelAdapter; + } + + @Nullable + private PulsarBinderHeaderMapper determineInboundHeaderMapper( + ExtendedConsumerProperties extConsumerProps) { + if (HeaderMode.none.equals(extConsumerProps.getHeaderMode())) { + return null; + } + return new PulsarBinderHeaderMapper(this.headerMapper); + } + + @Override + public PulsarConsumerProperties getExtendedConsumerProperties(String channelName) { + return this.extendedBindingProperties.getExtendedConsumerProperties(channelName); + } + + @Override + public PulsarProducerProperties getExtendedProducerProperties(String channelName) { + return this.extendedBindingProperties.getExtendedProducerProperties(channelName); + } + + @Override + public String getDefaultsPrefix() { + return null; + } + + @Override + public Class getExtendedPropertiesEntryClass() { + return null; + } + + public PulsarExtendedBindingProperties getExtendedBindingProperties() { + return this.extendedBindingProperties; + } + + public void setExtendedBindingProperties(PulsarExtendedBindingProperties extendedBindingProperties) { + this.extendedBindingProperties = extendedBindingProperties; + } + + static class PulsarMessageDrivenChannelAdapter extends MessageProducerSupport { + + AbstractPulsarMessageListenerContainer messageListenerContainer; + + public void send(Message message) { + sendMessage(message); + } + + @Override + protected void doStart() { + this.messageListenerContainer.start(); + } + + @Override + protected void doStop() { + this.messageListenerContainer.stop(); + } + + public void setMessageListenerContainer(AbstractPulsarMessageListenerContainer messageListenerContainer) { + this.messageListenerContainer = messageListenerContainer; + } + + } + + static class PulsarProducerConfigurationMessageHandler extends AbstractMessageProducingHandler + implements ManageableLifecycle { + + private final PulsarTemplate pulsarTemplate; + + private final Schema schema; + + private final String destination; + + private final ProducerBuilderCustomizer layeredProducerPropsCustomizer; + + private final PulsarHeaderMapper headerMapper; + + private boolean running = true; + + PulsarProducerConfigurationMessageHandler(PulsarTemplate pulsarTemplate, Schema schema, + String destination, ProducerBuilderCustomizer layeredProducerPropsCustomizer, + PulsarHeaderMapper headerMapper) { + this.pulsarTemplate = pulsarTemplate; + this.schema = schema; + this.destination = destination; + this.layeredProducerPropsCustomizer = layeredProducerPropsCustomizer; + this.headerMapper = headerMapper; + } + + @Override + public void start() { + try { + super.onInit(); + } + catch (Exception ex) { + this.logger.error(ex, "Initialization errors: "); + throw new RuntimeException(ex); + } + } + + @Override + public void stop() { + // TODO - should we close the underlyiung producer? + this.running = false; + } + + @Override + public boolean isRunning() { + return this.running; + } + + @Override + protected void handleMessageInternal(Message message) { + try { + // @formatter:off + this.pulsarTemplate.newMessage(message.getPayload()) + .withTopic(this.destination) + .withSchema(this.schema) + .withProducerCustomizer(this.layeredProducerPropsCustomizer) + .withMessageCustomizer(this.applySpringHeadersAsPulsarProperties(message.getHeaders())) + .sendAsync(); + // @formatter:on + } + catch (PulsarClientException ex) { + logger.trace(ex, "Failed to send message to destination: " + this.destination); + } + } + + private TypedMessageBuilderCustomizer applySpringHeadersAsPulsarProperties(MessageHeaders headers) { + return (mb) -> { + if (this.headerMapper != null) { + this.headerMapper.toPulsarHeaders(headers).forEach(mb::property); + } + }; + } + + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/config/PulsarBinderConfiguration.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/config/PulsarBinderConfiguration.java new file mode 100644 index 000000000..6e77e06a6 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/config/PulsarBinderConfiguration.java @@ -0,0 +1,75 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar.config; + +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; +import org.springframework.boot.context.properties.EnableConfigurationProperties; +import org.springframework.cloud.stream.binder.Binder; +import org.springframework.cloud.stream.binder.pulsar.PulsarMessageChannelBinder; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarExtendedBindingProperties; +import org.springframework.cloud.stream.binder.pulsar.provisioning.PulsarTopicProvisioner; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Configuration; +import org.springframework.pulsar.autoconfigure.PulsarProperties; +import org.springframework.pulsar.core.PulsarAdministration; +import org.springframework.pulsar.core.PulsarConsumerFactory; +import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.core.SchemaResolver; +import org.springframework.pulsar.support.header.JacksonUtils; +import org.springframework.pulsar.support.header.JsonPulsarHeaderMapper; +import org.springframework.pulsar.support.header.PulsarHeaderMapper; +import org.springframework.pulsar.support.header.ToStringPulsarHeaderMapper; + +/** + * Pulsar binder {@link Configuration}. + * + * @author Soby Chacko + */ +@Configuration(proxyBeanMethods = false) +@ConditionalOnMissingBean(Binder.class) +@EnableConfigurationProperties({ PulsarProperties.class, PulsarExtendedBindingProperties.class, + PulsarBinderConfigurationProperties.class }) +public class PulsarBinderConfiguration { + + @Bean + public PulsarTopicProvisioner pulsarTopicProvisioner(PulsarAdministration pulsarAdministration, + PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties) { + return new PulsarTopicProvisioner(pulsarAdministration, pulsarBinderConfigurationProperties); + } + + @Bean + @ConditionalOnMissingBean + public PulsarHeaderMapper pulsarHeaderMapper() { + if (JacksonUtils.isJacksonPresent()) { + return JsonPulsarHeaderMapper.builder().build(); + } + return new ToStringPulsarHeaderMapper(); + } + + @Bean + public PulsarMessageChannelBinder pulsarMessageChannelBinder(PulsarTopicProvisioner pulsarTopicProvisioner, + PulsarTemplate pulsarTemplate, PulsarConsumerFactory pulsarConsumerFactory, + PulsarBinderConfigurationProperties binderConfigProps, PulsarExtendedBindingProperties bindingConfigProps, + SchemaResolver schemaResolver, PulsarHeaderMapper headerMapper) { + PulsarMessageChannelBinder pulsarMessageChannelBinder = new PulsarMessageChannelBinder(pulsarTopicProvisioner, + pulsarTemplate, pulsarConsumerFactory, binderConfigProps, schemaResolver, headerMapper); + pulsarMessageChannelBinder.setExtendedBindingProperties(bindingConfigProps); + return pulsarMessageChannelBinder; + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/package-info.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/package-info.java new file mode 100644 index 000000000..95966720b --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/package-info.java @@ -0,0 +1,25 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +/** + * Package containing Spring Cloud Stream binder classes for Apache Pulsar. + */ +@NonNullApi +@NonNullFields +package org.springframework.cloud.stream.binder.pulsar; + +import org.springframework.lang.NonNullApi; +import org.springframework.lang.NonNullFields; diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarBinderConfigurationProperties.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarBinderConfigurationProperties.java new file mode 100644 index 000000000..7722b7c6e --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarBinderConfigurationProperties.java @@ -0,0 +1,71 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar.properties; + +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.boot.context.properties.NestedConfigurationProperty; +import org.springframework.lang.Nullable; +import org.springframework.pulsar.autoconfigure.ConsumerConfigProperties; +import org.springframework.pulsar.autoconfigure.ProducerConfigProperties; + +/** + * {@link ConfigurationProperties @ConfigurationProperties} for the Pulsar binder. + *

+ * These properties are applied at the binder level (to all bindings). + * + * @author Soby Chacko + * @author Chris Bono + */ +@ConfigurationProperties(prefix = "spring.cloud.stream.pulsar.binder") +public class PulsarBinderConfigurationProperties { + + /** + * Pulsar consumer specific binder-level properties (applied to all bindings). + */ + @NestedConfigurationProperty + private final ConsumerConfigProperties consumer = new ConsumerConfigProperties(); + + /** + * Pulsar producer specific binder-level properties (applied to all bindings). + */ + @NestedConfigurationProperty + private final ProducerConfigProperties producer = new ProducerConfigProperties(); + + /** + * Number of topic partitions. + */ + @Nullable + private Integer partitionCount; + + public ConsumerConfigProperties getConsumer() { + return this.consumer; + } + + public ProducerConfigProperties getProducer() { + return this.producer; + } + + @Nullable + public Integer getPartitionCount() { + return this.partitionCount; + } + + public void setPartitionCount(Integer partitionCount) { + this.partitionCount = partitionCount; + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarBindingProperties.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarBindingProperties.java new file mode 100644 index 000000000..8abe593ca --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarBindingProperties.java @@ -0,0 +1,71 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar.properties; + +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.boot.context.properties.NestedConfigurationProperty; +import org.springframework.cloud.stream.binder.BinderSpecificPropertiesProvider; + +/** + * Container for Pulsar specific extended producer and consumer binding properties. + *

+ * These properties are applied to individual bindings and will override any binder-level + * setting. + * + *

+ * NOTE: This class is only referenced as a value in the + * {@link PulsarExtendedBindingProperties#getBindings() bindings map} and therefore, by + * default is not included in the generated configuration metadata. To get around this + * limitation it is annotated with {@code @ConfigurationProperties}. However, that is the + * only reason it is annotated and is not intended to be used directly. + * + * @author Soby Chacko + * @author Chris Bono + */ +@SuppressWarnings("ConfigurationProperties") +@ConfigurationProperties("spring.cloud.stream.pulsar.bindings.for-docs-only") +public class PulsarBindingProperties implements BinderSpecificPropertiesProvider { + + /** + * Pulsar consumer specific binding properties. + */ + @NestedConfigurationProperty + private PulsarConsumerProperties consumer = new PulsarConsumerProperties(); + + /** + * Pulsar producer specific binding properties. + */ + @NestedConfigurationProperty + private PulsarProducerProperties producer = new PulsarProducerProperties(); + + public PulsarConsumerProperties getConsumer() { + return this.consumer; + } + + public void setConsumer(PulsarConsumerProperties consumer) { + this.consumer = consumer; + } + + public PulsarProducerProperties getProducer() { + return this.producer; + } + + public void setProducer(PulsarProducerProperties producer) { + this.producer = producer; + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarConsumerProperties.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarConsumerProperties.java new file mode 100644 index 000000000..ea4684aac --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarConsumerProperties.java @@ -0,0 +1,93 @@ +/* + * Copyright 2022-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar.properties; + +import org.apache.pulsar.common.schema.SchemaType; + +import org.springframework.lang.Nullable; +import org.springframework.pulsar.autoconfigure.ConsumerConfigProperties; + +/** + * Pulsar consumer properties used by the binder. + * + * @author Soby Chacko + * @author Chris Bono + */ +public class PulsarConsumerProperties extends ConsumerConfigProperties { + + /** + * Pulsar {@link SchemaType} for this binding. + */ + @Nullable + private SchemaType schemaType; + + /** + * Pulsar message type for this binding. + */ + @Nullable + private Class messageType; + + /** + * Pulsar message key type for this binding (only used when schema type is + * {@code }KEY_VALUE}). + */ + @Nullable + private Class messageKeyType; + + /** + * Number of topic partitions. + */ + @Nullable + private Integer partitionCount; + + @Nullable + public SchemaType getSchemaType() { + return this.schemaType; + } + + public void setSchemaType(@Nullable SchemaType schemaType) { + this.schemaType = schemaType; + } + + @Nullable + public Class getMessageType() { + return this.messageType; + } + + public void setMessageType(@Nullable Class messageType) { + this.messageType = messageType; + } + + @Nullable + public Class getMessageKeyType() { + return this.messageKeyType; + } + + public void setMessageKeyType(@Nullable Class messageKeyType) { + this.messageKeyType = messageKeyType; + } + + @Nullable + public Integer getPartitionCount() { + return this.partitionCount; + } + + public void setPartitionCount(Integer partitionCount) { + this.partitionCount = partitionCount; + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarExtendedBindingProperties.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarExtendedBindingProperties.java new file mode 100644 index 000000000..4e487a3c0 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarExtendedBindingProperties.java @@ -0,0 +1,60 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar.properties; + +import java.util.Map; + +import org.springframework.boot.context.properties.ConfigurationProperties; +import org.springframework.cloud.stream.binder.AbstractExtendedBindingProperties; +import org.springframework.cloud.stream.binder.BinderSpecificPropertiesProvider; + + +/** + * {@link ConfigurationProperties @ConfigurationProperties} for Pulsar binder specific + * extensions to the common binding properties. + *

+ * These properties are applied to individual bindings and will override any binder-level + * settings. + * + * @author Soby Chacko + * @author Chris Bono + */ +@ConfigurationProperties("spring.cloud.stream.pulsar") +public class PulsarExtendedBindingProperties extends + AbstractExtendedBindingProperties { + + private static final String DEFAULTS_PREFIX = "spring.cloud.stream.pulsar.default"; + + @Override + public String getDefaultsPrefix() { + return DEFAULTS_PREFIX; + } + + /** + * Properties per individual binding name (e.g. 'mySink-in-0'). + */ + @Override + public Map getBindings() { + return this.doGetBindings(); + } + + @Override + public Class getExtendedPropertiesEntryClass() { + return PulsarBindingProperties.class; + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarProducerProperties.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarProducerProperties.java new file mode 100644 index 000000000..db73378e2 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/PulsarProducerProperties.java @@ -0,0 +1,93 @@ +/* + * Copyright 2022-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar.properties; + +import org.apache.pulsar.common.schema.SchemaType; + +import org.springframework.lang.Nullable; +import org.springframework.pulsar.autoconfigure.ProducerConfigProperties; + +/** + * Pulsar producer properties used by the binder. + * + * @author Soby Chacko + * @author Chris Bono + */ +public class PulsarProducerProperties extends ProducerConfigProperties { + + /** + * Pulsar {@link SchemaType} for this binding. + */ + @Nullable + private SchemaType schemaType; + + /** + * Pulsar message type for this binding. + */ + @Nullable + private Class messageType; + + /** + * Pulsar message key type for this binding (only used when schema type is + * {@code }KEY_VALUE}). + */ + @Nullable + private Class messageKeyType; + + /** + * Number of topic partitions. + */ + @Nullable + private Integer partitionCount; + + @Nullable + public SchemaType getSchemaType() { + return this.schemaType; + } + + public void setSchemaType(@Nullable SchemaType schemaType) { + this.schemaType = schemaType; + } + + @Nullable + public Class getMessageType() { + return this.messageType; + } + + public void setMessageType(@Nullable Class messageType) { + this.messageType = messageType; + } + + @Nullable + public Class getMessageKeyType() { + return this.messageKeyType; + } + + public void setMessageKeyType(@Nullable Class messageKeyType) { + this.messageKeyType = messageKeyType; + } + + @Nullable + public Integer getPartitionCount() { + return this.partitionCount; + } + + public void setPartitionCount(Integer partitionCount) { + this.partitionCount = partitionCount; + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/package-info.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/package-info.java new file mode 100644 index 000000000..688bc2473 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/properties/package-info.java @@ -0,0 +1,25 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +/** + * Package containing Spring Cloud Stream binder properties classes for Apache Pulsar. + */ +@NonNullApi +@NonNullFields +package org.springframework.cloud.stream.binder.pulsar.properties; + +import org.springframework.lang.NonNullApi; +import org.springframework.lang.NonNullFields; diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/provisioning/PulsarTopicProvisioner.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/provisioning/PulsarTopicProvisioner.java new file mode 100644 index 000000000..46a2f674e --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/provisioning/PulsarTopicProvisioner.java @@ -0,0 +1,94 @@ +/* + * Copyright 2022-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar.provisioning; + +import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; +import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarConsumerProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarProducerProperties; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.cloud.stream.provisioning.ProducerDestination; +import org.springframework.cloud.stream.provisioning.ProvisioningException; +import org.springframework.cloud.stream.provisioning.ProvisioningProvider; +import org.springframework.lang.Nullable; +import org.springframework.pulsar.core.PulsarAdministration; +import org.springframework.pulsar.core.PulsarTopic; + +/** + * Pulsar topic provisioner. + * + * @author Soby Chacko + */ +public class PulsarTopicProvisioner implements + ProvisioningProvider, ExtendedProducerProperties> { + + private final PulsarAdministration pulsarAdministration; + + private final PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties; + + public PulsarTopicProvisioner(PulsarAdministration pulsarAdministration, + PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties) { + this.pulsarAdministration = pulsarAdministration; + this.pulsarBinderConfigurationProperties = pulsarBinderConfigurationProperties; + } + + @Override + public ProducerDestination provisionProducerDestination(String name, + ExtendedProducerProperties pulsarProducerProperties) + throws ProvisioningException { + Integer partitionCountFromBinding = pulsarProducerProperties.getExtension().getPartitionCount(); + var partitionCount = getPartitionCount(partitionCountFromBinding); + var pulsarTopic = PulsarTopic.builder(name).numberOfPartitions(partitionCount).build(); + this.pulsarAdministration.createOrModifyTopics(pulsarTopic); + return new PulsarDestination(pulsarTopic.topicName(), pulsarTopic.numberOfPartitions()); + } + + private int getPartitionCount(@Nullable Integer partitionCountConfig) { + var partitionCount = this.pulsarBinderConfigurationProperties.getPartitionCount(); + if (partitionCountConfig != null && partitionCountConfig > 0) { + partitionCount = partitionCountConfig; + } + return partitionCount == null ? 0 : partitionCount; + } + + @Override + public ConsumerDestination provisionConsumerDestination(String name, String group, + ExtendedConsumerProperties pulsarConsumerProperties) + throws ProvisioningException { + var partitionCountFromBinding = pulsarConsumerProperties.getExtension().getPartitionCount(); + var partitionCount = getPartitionCount(partitionCountFromBinding); + var pulsarTopic = PulsarTopic.builder(name).numberOfPartitions(partitionCount).build(); + this.pulsarAdministration.createOrModifyTopics(pulsarTopic); + return new PulsarDestination(pulsarTopic.topicName(), pulsarTopic.numberOfPartitions()); + } + + private record PulsarDestination(String destinationName, + Integer partitions) implements ProducerDestination, ConsumerDestination { + + @Override + public String getName() { + return this.destinationName; + } + + @Override + public String getNameForPartition(int partition) { + return this.destinationName; + } + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/provisioning/package-info.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/provisioning/package-info.java new file mode 100644 index 000000000..c6a016779 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/java/org/springframework/cloud/stream/binder/pulsar/provisioning/package-info.java @@ -0,0 +1,25 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +/** + * Package containing Spring Cloud Stream binder provisioning classes for Apache Pulsar. + */ +@NonNullApi +@NonNullFields +package org.springframework.cloud.stream.binder.pulsar.provisioning; + +import org.springframework.lang.NonNullApi; +import org.springframework.lang.NonNullFields; diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/resources/META-INF/spring.binders b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/resources/META-INF/spring.binders new file mode 100644 index 000000000..778e113a6 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/main/resources/META-INF/spring.binders @@ -0,0 +1,2 @@ +pulsar:\ +org.springframework.cloud.stream.binder.pulsar.config.PulsarBinderConfiguration diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/AbstractPulsarTestBinder.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/AbstractPulsarTestBinder.java new file mode 100644 index 000000000..a6191c2f7 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/AbstractPulsarTestBinder.java @@ -0,0 +1,48 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import org.springframework.cloud.stream.binder.AbstractPollableConsumerTestBinder; +import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; +import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarConsumerProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarProducerProperties; +import org.springframework.context.ApplicationContext; + +/** + * Base class for {@link PulsarTestBinder}. + * + * @author Soby Chacko + */ +public abstract class AbstractPulsarTestBinder extends + AbstractPollableConsumerTestBinder, ExtendedProducerProperties> { + + private ApplicationContext applicationContext; + + @Override + public void cleanup() { + } + + protected final void setApplicationContext(ApplicationContext context) { + this.applicationContext = context; + } + + public ApplicationContext getApplicationContext() { + return this.applicationContext; + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderConfigurationPropertiesTests.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderConfigurationPropertiesTests.java new file mode 100644 index 000000000..d21b40335 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderConfigurationPropertiesTests.java @@ -0,0 +1,121 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.util.HashMap; +import java.util.Map; + +import org.apache.pulsar.client.api.ProducerAccessMode; +import org.apache.pulsar.client.api.SubscriptionMode; +import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils; +import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; +import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; +import org.assertj.core.api.InstanceOfAssertFactories; +import org.junit.jupiter.api.Test; + +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.cloud.stream.binder.pulsar.properties.PulsarBinderConfigurationProperties; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatNoException; + +/** + * Tests for {@link PulsarBinderConfigurationProperties}. + * + * @author Chris Bono + */ +public class PulsarBinderConfigurationPropertiesTests { + + private final PulsarBinderConfigurationProperties properties = new PulsarBinderConfigurationProperties(); + + private void bind(Map map) { + ConfigurationPropertySource source = new MapConfigurationPropertySource(map); + new Binder(source).bind("spring.cloud.stream.pulsar.binder", Bindable.ofInstance(this.properties)); + } + + @Test + void partitionCountProperty() { + assertThat(properties.getPartitionCount()).isNull(); + bind(Map.of("spring.cloud.stream.pulsar.binder.partition-count", "5150")); + assertThat(properties.getPartitionCount()).isEqualTo(5150); + } + + @Test + void producerProperties() { + // Only spot check a few values (PulsarPropertiesTests does the heavy lifting) + Map props = new HashMap<>(); + props.put("spring.cloud.stream.pulsar.binder.producer.topic-name", "my-topic"); + props.put("spring.cloud.stream.pulsar.binder.producer.send-timeout", "2s"); + props.put("spring.cloud.stream.pulsar.binder.producer.max-pending-messages", "3"); + props.put("spring.cloud.stream.pulsar.binder.producer.producer-access-mode", "exclusive"); + props.put("spring.cloud.stream.pulsar.binder.producer.properties[my-prop]", "my-prop-value"); + + bind(props); + Map producerProps = PulsarBinderUtils.convertProducerPropertiesToMap(properties.getProducer()); + + // Verify that the props can be loaded in a ProducerBuilder + assertThatNoException().isThrownBy(() -> ConfigurationDataUtils.loadData(producerProps, + new ProducerConfigurationData(), ProducerConfigurationData.class)); + + // @formatter:off + assertThat(producerProps) + .containsEntry("topicName", "my-topic") + .containsEntry("sendTimeoutMs", 2_000) + .containsEntry("maxPendingMessages", 3) + .containsEntry("accessMode", ProducerAccessMode.Exclusive) + .hasEntrySatisfying("properties", properties -> + assertThat(properties) + .asInstanceOf(InstanceOfAssertFactories.map(String.class, String.class)) + .containsEntry("my-prop", "my-prop-value")); + // @formatter:on + } + + @Test + void consumerProperties() { + // Only spot check a few values (PulsarPropertiesTests does the heavy lifting) + Map props = new HashMap<>(); + props.put("spring.cloud.stream.pulsar.binder.consumer.topics[0]", "my-topic"); + props.put("spring.cloud.stream.pulsar.binder.consumer.subscription-properties[my-sub-prop]", + "my-sub-prop-value"); + props.put("spring.cloud.stream.pulsar.binder.consumer.subscription-mode", "nondurable"); + props.put("spring.cloud.stream.pulsar.binder.consumer.receiver-queue-size", "1"); + + bind(props); + Map consumerProps = PulsarBinderUtils.convertConsumerPropertiesToMap(properties.getConsumer()); + + // Verify that the props can be loaded in a ConsumerBuilder + assertThatNoException().isThrownBy(() -> ConfigurationDataUtils.loadData(consumerProps, + new ConsumerConfigurationData<>(), ConsumerConfigurationData.class)); + + // @formatter:off + assertThat(consumerProps) + .hasEntrySatisfying("topicNames", + topics -> assertThat(topics).asInstanceOf(InstanceOfAssertFactories.collection(String.class)) + .containsExactly("my-topic")) + .hasEntrySatisfying("subscriptionProperties", + properties -> assertThat(properties) + .asInstanceOf(InstanceOfAssertFactories.map(String.class, String.class)) + .containsEntry("my-sub-prop", "my-sub-prop-value")) + .containsEntry("subscriptionMode", SubscriptionMode.NonDurable) + .containsEntry("receiverQueueSize", 1); + // @formatter:on + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderHeaderMapperTests.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderHeaderMapperTests.java new file mode 100644 index 000000000..6a2b5f1cf --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderHeaderMapperTests.java @@ -0,0 +1,115 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.util.HashMap; +import java.util.Map; + +import org.apache.pulsar.client.api.Message; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.mockito.InjectMocks; +import org.mockito.Mock; +import org.mockito.junit.jupiter.MockitoExtension; + +import org.springframework.cloud.stream.binder.BinderHeaders; +import org.springframework.integration.IntegrationMessageHeaderAccessor; +import org.springframework.messaging.MessageHeaders; +import org.springframework.pulsar.support.header.PulsarHeaderMapper; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.AssertionsForClassTypes.entry; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; +import static org.mockito.Mockito.when; + +/** + * Tests for {@link PulsarBinderHeaderMapper}. + * + * @author Chris Bono + */ +@ExtendWith(MockitoExtension.class) +class PulsarBinderHeaderMapperTests { + + @Mock + private PulsarHeaderMapper delegateMapper; + + @InjectMocks + private PulsarBinderHeaderMapper binderHeaderMapper; + + @Nested + class ToPulsarHeadersOutboundTests { + + @Test + void delegateReturnsEmptyHeaders() { + var delegatePulsarHeaders = new HashMap(); + when(delegateMapper.toPulsarHeaders(any(MessageHeaders.class))).thenReturn(delegatePulsarHeaders); + var springHeaders = mock(MessageHeaders.class); + var pulsarHeaders = binderHeaderMapper.toPulsarHeaders(springHeaders); + verify(delegateMapper).toPulsarHeaders(springHeaders); + assertThat(pulsarHeaders).isEmpty(); + } + + @Test + void neverHeadersRemovedFromDelegateHeaders() { + var delegatePulsarHeaders = new HashMap(); + delegatePulsarHeaders.put(MessageHeaders.ID, "5150"); + delegatePulsarHeaders.put(MessageHeaders.TIMESTAMP, "12345"); + delegatePulsarHeaders.put(IntegrationMessageHeaderAccessor.DELIVERY_ATTEMPT, "5"); + delegatePulsarHeaders.put(BinderHeaders.NATIVE_HEADERS_PRESENT, "true"); + delegatePulsarHeaders.put("foo", "bar"); + when(delegateMapper.toPulsarHeaders(any(MessageHeaders.class))).thenReturn(delegatePulsarHeaders); + var springHeaders = mock(MessageHeaders.class); + var pulsarHeaders = binderHeaderMapper.toPulsarHeaders(springHeaders); + verify(delegateMapper).toPulsarHeaders(springHeaders); + assertThat(pulsarHeaders).containsOnly(entry("foo", "bar")); + } + + } + + @Nested + class ToSpringHeadersInboundTests { + + @Test + void delegateReturnsEmptyHeaders() { + var emptyDelegateHeaders = mock(MessageHeaders.class); + when(emptyDelegateHeaders.isEmpty()).thenReturn(true); + when(delegateMapper.toSpringHeaders(any(Message.class))).thenReturn(emptyDelegateHeaders); + var springHeaders = binderHeaderMapper.toSpringHeaders(mock(Message.class)); + assertThat(springHeaders).isSameAs(emptyDelegateHeaders); + verify(springHeaders).isEmpty(); + verifyNoMoreInteractions(springHeaders); + } + + @Test + void nativeHeadersIndicatorAddedToDelegateHeaders() { + var delegateSpringHeaders = new MessageHeaders(Map.of("foo", "bar")); + when(delegateMapper.toSpringHeaders(any(Message.class))).thenReturn(delegateSpringHeaders); + var pulsarMessage = mock(Message.class); + var springHeaders = binderHeaderMapper.toSpringHeaders(pulsarMessage); + verify(delegateMapper).toSpringHeaders(pulsarMessage); + assertThat(springHeaders).containsEntry("foo", "bar").containsEntry(BinderHeaders.NATIVE_HEADERS_PRESENT, + Boolean.TRUE); + + } + + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderIntegrationTests.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderIntegrationTests.java new file mode 100644 index 000000000..08ec2731b --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderIntegrationTests.java @@ -0,0 +1,892 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.time.Duration; +import java.util.ArrayList; +import java.util.Collection; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.function.Consumer; +import java.util.function.Supplier; + +import org.apache.pulsar.client.api.Producer; +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.impl.schema.JSONSchema; +import org.apache.pulsar.common.schema.KeyValue; +import org.awaitility.Awaitility; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.extension.ExtendWith; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +import org.springframework.boot.SpringApplication; +import org.springframework.boot.SpringBootConfiguration; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.test.system.CapturedOutput; +import org.springframework.boot.test.system.OutputCaptureExtension; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.context.annotation.Import; +import org.springframework.lang.Nullable; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.support.MessageBuilder; +import org.springframework.pulsar.autoconfigure.PulsarProperties; +import org.springframework.pulsar.core.ConsumerBuilderCustomizer; +import org.springframework.pulsar.core.DefaultPulsarConsumerFactory; +import org.springframework.pulsar.core.DefaultPulsarProducerFactory; +import org.springframework.pulsar.core.DefaultSchemaResolver; +import org.springframework.pulsar.core.ProducerBuilderCustomizer; +import org.springframework.pulsar.core.PulsarConsumerFactory; +import org.springframework.pulsar.core.PulsarProducerFactory; +import org.springframework.pulsar.core.SchemaResolver.SchemaResolverCustomizer; +import org.springframework.pulsar.core.TopicResolver; +import org.springframework.pulsar.support.header.PulsarHeaderMapper; +import org.springframework.pulsar.support.header.ToStringPulsarHeaderMapper; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Integration tests for {@link PulsarBinderIntegrationTests}. + * + * @author Soby Chacko + * @author Chris Bono + */ +@ExtendWith(OutputCaptureExtension.class) +@SuppressWarnings("JUnitMalformedDeclaration") +class PulsarBinderIntegrationTests implements PulsarTestContainerSupport { + + private static final int AWAIT_DURATION = 10; + + @Test + void binderAndBindingPropsAreAppliedAndRespected(CapturedOutput output) { + SpringApplication app = new SpringApplication(BinderAndBindingPropsTestConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext context = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=textSupplier;textLogger", + "--spring.cloud.stream.bindings.textLogger-in-0.destination=textSupplier-out-0", + "--spring.pulsar.producer.producer-name=textSupplierProducer-fromBase", + "--spring.cloud.stream.pulsar.binder.producer.producer-name=textSupplierProducer-fromBinder", + "--spring.cloud.stream.pulsar.bindings.textSupplier-out-0.producer.producer-name=textSupplierProducer-fromBinding", + "--spring.cloud.stream.pulsar.binder.producer.max-pending-messages=1100", + "--spring.pulsar.producer.block-if-queue-full=true", + "--spring.cloud.stream.pulsar.binder.consumer.subscription-name=textLoggerSub-fromBinder", + "--spring.cloud.stream.pulsar.binder.consumer.consumer-name=textLogger-fromBinder", + "--spring.cloud.stream.pulsar.bindings.textLogger-in-0.consumer.consumer-name=textLogger-fromBinding")) { + + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: test-basic-scenario")); + + // now verify the properties were set onto producer and consumer as expected + TrackingProducerFactory producerFactory = context.getBean(TrackingProducerFactory.class); + assertThat(producerFactory.producersCreated).isNotEmpty().element(0) + .hasFieldOrPropertyWithValue("producerName", "textSupplierProducer-fromBinding") + .hasFieldOrPropertyWithValue("conf.maxPendingMessages", 1100) + .hasFieldOrPropertyWithValue("conf.blockIfQueueFull", true); + + TrackingConsumerFactory consumerFactory = context.getBean(TrackingConsumerFactory.class); + assertThat(consumerFactory.consumersCreated).isNotEmpty().element(0) + .hasFieldOrPropertyWithValue("consumerName", "textLogger-fromBinding") + .hasFieldOrPropertyWithValue("conf.subscriptionName", "textLoggerSub-fromBinder"); + } + } + + @Nested + class DefaultEncoding { + + @Test + void primitiveTypeString(CapturedOutput output) { + SpringApplication app = new SpringApplication(PrimitiveTextConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=textSupplier;textLogger", + "--spring.cloud.stream.bindings.textLogger-in-0.destination=textSupplier-out-0", + "--spring.cloud.stream.pulsar.bindings.textLogger-in-0.consumer.subscription-name=pbit-text-sub1")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: test-basic-scenario")); + } + } + + @Test + void primitiveTypeFloat(CapturedOutput output) { + SpringApplication app = new SpringApplication(PrimitiveFloatConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=piSupplier;piLogger", + "--spring.cloud.stream.bindings.piSupplier-out-0.destination=pi-stream", + "--spring.cloud.stream.bindings.piLogger-in-0.destination=pi-stream", + "--spring.cloud.stream.pulsar.bindings.piLogger-in-0.consumer.subscription-name=pbit-float-sub1")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: 3.14")); + } + } + + } + + @Nested + class NativeEncoding { + + @Test + void primitiveTypeFloat(CapturedOutput output) { + SpringApplication app = new SpringApplication(PrimitiveFloatConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=piSupplier;piLogger", + "--spring.cloud.stream.bindings.piLogger-in-0.destination=piSupplier-out-0", + "--spring.cloud.stream.bindings.piSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.piSupplier-out-0.producer.schema-type=FLOAT", + "--spring.cloud.stream.bindings.piLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.piLogger-in-0.consumer.schema-type=FLOAT", + "--spring.cloud.stream.pulsar.bindings.piLogger-in-0.consumer.subscription-name=pbit-float-sub2")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: 3.14")); + } + } + + @Test + void jsonTypeFooWithSchemaType(CapturedOutput output) { + SpringApplication app = new SpringApplication(JsonFooConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=fooSupplier;fooLogger", + "--spring.cloud.stream.bindings.fooSupplier-out-0.destination=foo-stream-1", + "--spring.cloud.stream.bindings.fooLogger-in-0.destination=foo-stream-1", + "--spring.cloud.stream.bindings.fooSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.fooSupplier-out-0.producer.schema-type=JSON", + "--spring.cloud.stream.pulsar.bindings.fooSupplier-out-0.producer.message-type=" + + Foo.class.getName(), + "--spring.cloud.stream.bindings.fooLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.fooLogger-in-0.consumer.schema-type=JSON", + "--spring.cloud.stream.pulsar.bindings.fooLogger-in-0.consumer.message-type=" + Foo.class.getName(), + "--spring.cloud.stream.pulsar.bindings.fooLogger-in-0.consumer.subscription-name=pbit-foo-sub1")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: Foo[value=5150]")); + } + } + + @Test + void jsonTypeFooWithoutSchemaTypeDefaultsToJsonSchema(CapturedOutput output) { + SpringApplication app = new SpringApplication(JsonFooConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=fooSupplier;fooLogger", + "--spring.cloud.stream.bindings.fooSupplier-out-0.destination=foo-stream-2", + "--spring.cloud.stream.bindings.fooLogger-in-0.destination=foo-stream-2", + "--spring.cloud.stream.bindings.fooSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.fooSupplier-out-0.producer.message-type=" + + Foo.class.getName(), + "--spring.cloud.stream.bindings.fooLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.fooLogger-in-0.consumer.message-type=" + Foo.class.getName(), + "--spring.cloud.stream.pulsar.bindings.fooLogger-in-0.consumer.subscription-name=pbit-foo-sub2")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: Foo[value=5150]")); + } + } + + @Test + void avroTypeUserWithSchemaType(CapturedOutput output) { + SpringApplication app = new SpringApplication(AvroUserConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=userSupplier;userLogger", + "--spring.cloud.stream.bindings.userSupplier-out-0.destination=user-stream-1", + "--spring.cloud.stream.bindings.userLogger-in-0.destination=user-stream-1", + "--spring.cloud.stream.bindings.userSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.schema-type=AVRO", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.bindings.userLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.schema-type=AVRO", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.subscription-name=pbit-user-sub1")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: User{name='user21', age=21}")); + } + } + + @Test + void avroTypeUserWithoutSchemaTypeWithCustomMappingsViaProps(CapturedOutput output) { + SpringApplication app = new SpringApplication(AvroUserConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=userSupplier;userLogger", + "--spring.cloud.stream.bindings.userSupplier-out-0.destination=user-stream-2", + "--spring.cloud.stream.bindings.userLogger-in-0.destination=user-stream-2", + "--spring.cloud.stream.bindings.userSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.bindings.userLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.subscription-name=pbit-user-sub2", + "--spring.pulsar.defaults.type-mappings[0].message-type=%s".formatted(User.class.getName()), + "--spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=AVRO")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: User{name='user21', age=21}")); + } + } + + @Test + void avroTypeUserWithoutSchemaTypeWithCustomMappingsViaCustomizer(CapturedOutput output) { + SpringApplication app = new SpringApplication(AvroUserConfigCustomMappings.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=userSupplier;userLogger", + "--spring.cloud.stream.bindings.userSupplier-out-0.destination=user-stream-3", + "--spring.cloud.stream.bindings.userLogger-in-0.destination=user-stream-3", + "--spring.cloud.stream.bindings.userSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.bindings.userLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.subscription-name=pbit-user-sub3")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: User{name='user21', age=21}")); + } + } + + @Test + void keyValueAvroTypeWithSchemaTypeAndCustomTypeMappingsViaProps(CapturedOutput output) { + SpringApplication app = new SpringApplication(KeyValueAvroUserConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=userSupplier;userLogger", + "--spring.cloud.stream.bindings.userSupplier-out-0.destination=kv-stream-1", + "--spring.cloud.stream.bindings.userLogger-in-0.destination=kv-stream-1", + "--spring.cloud.stream.bindings.userSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.schema-type=KEY_VALUE", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-key-type=" + + String.class.getName(), + "--spring.cloud.stream.bindings.userLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.schema-type=KEY_VALUE", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-key-type=" + + String.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.subscription-name=pbit-kv-sub1", + "--spring.pulsar.defaults.type-mappings[0].message-type=%s".formatted(User.class.getName()), + "--spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=AVRO")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: 21->User{name='user21', age=21}")); + } + } + + @Test + void keyValueAvroTypeWithoutSchemaTypeAndCustomTypeMappingsViaProps(CapturedOutput output) { + SpringApplication app = new SpringApplication(KeyValueAvroUserConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=userSupplier;userLogger", + "--spring.cloud.stream.bindings.userSupplier-out-0.destination=kv-stream-2", + "--spring.cloud.stream.bindings.userLogger-in-0.destination=kv-stream-2", + "--spring.cloud.stream.bindings.userSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-key-type=" + + String.class.getName(), + "--spring.cloud.stream.bindings.userLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-key-type=" + + String.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.subscription-name=pbit-kv-sub2", + "--spring.pulsar.defaults.type-mappings[0].message-type=%s".formatted(User.class.getName()), + "--spring.pulsar.defaults.type-mappings[0].schema-info.schema-type=AVRO")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: 21->User{name='user21', age=21}")); + } + } + + @Test + void keyValueAvroTypeWithSchemaTypeAndCustomTypeMappingsViaCustomizer(CapturedOutput output) { + SpringApplication app = new SpringApplication(KeyValueAvroUserConfigCustomMappings.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=userSupplier;userLogger", + "--spring.cloud.stream.bindings.userSupplier-out-0.destination=kv-stream-3", + "--spring.cloud.stream.bindings.userLogger-in-0.destination=kv-stream-3", + "--spring.cloud.stream.bindings.userSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.schema-type=KEY_VALUE", + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userSupplier-out-0.producer.message-key-type=" + + String.class.getName(), + "--spring.cloud.stream.bindings.userLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.schema-type=KEY_VALUE", + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-type=" + + User.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.message-key-type=" + + String.class.getName(), + "--spring.cloud.stream.pulsar.bindings.userLogger-in-0.consumer.subscription-name=pbit-kv-sub3")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: 21->User{name='user21', age=21}")); + } + } + + @Test + void keyValueJsonTypeWithoutSchemaTypeAndWithoutCustomTypeMappings(CapturedOutput output) { + SpringApplication app = new SpringApplication(KeyValueJsonFooConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=fooSupplier;fooLogger", + "--spring.cloud.stream.bindings.fooSupplier-out-0.destination=kv-stream-4", + "--spring.cloud.stream.bindings.fooLogger-in-0.destination=kv-stream-4", + "--spring.cloud.stream.bindings.fooSupplier-out-0.producer.use-native-encoding=true", + "--spring.cloud.stream.pulsar.bindings.fooSupplier-out-0.producer.message-type=" + + Foo.class.getName(), + "--spring.cloud.stream.pulsar.bindings.fooSupplier-out-0.producer.message-key-type=" + + String.class.getName(), + "--spring.cloud.stream.bindings.fooLogger-in-0.consumer.use-native-decoding=true", + "--spring.cloud.stream.pulsar.bindings.fooLogger-in-0.consumer.message-type=" + Foo.class.getName(), + "--spring.cloud.stream.pulsar.bindings.fooLogger-in-0.consumer.message-key-type=" + + String.class.getName(), + "--spring.cloud.stream.pulsar.bindings.fooLogger-in-0.consumer.subscription-name=pbit-kv-sub4")) { + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: 5150->Foo[value=5150]")); + } + } + + } + + @Nested + class CustomMessageHeaders { + + @Test + void headersPropagatedSendAndReceive(CapturedOutput output) { + SpringApplication app = new SpringApplication(CustomSimpleHeadersConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=springMessageSupplier;springMessageLogger", + "--spring.cloud.stream.bindings.springMessageSupplier-out-0.destination=cmh-1", + "--spring.cloud.stream.bindings.springMessageLogger-in-0.destination=cmh-1", + "--spring.cloud.stream.pulsar.bindings.springMessageLogger-in-0.consumer.subscription-name=pbit-cmh1-sub1")) { + // Wait for a few of the messages to flow through (check for index = 5) + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)).until( + () -> output.toString().contains("Hello binder: test-headers-msg-5 w/ custom-id: 5150-5")); + } + } + + @Test + void complexHeadersAreEncodedAndPropagated(CapturedOutput output) { + SpringApplication app = new SpringApplication(CustomComplexHeadersConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=springMessageSupplier;springMessageLogger", + "--spring.cloud.stream.bindings.springMessageSupplier-out-0.destination=cmh-2", + "--spring.cloud.stream.bindings.springMessageLogger-in-0.destination=cmh-2", + "--spring.cloud.stream.pulsar.bindings.springMessageLogger-in-0.consumer.subscription-name=pbit-cmh2-sub1")) { + // Wait for a few of the messages to flow through (check for index = 5) + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)).until(() -> output.toString() + .contains("Hello binder: test-headers-msg-5 w/ custom-id: FooHeader[value=5150-5]")); + } + } + + @Test + void producerHeaderModeNone(CapturedOutput output) { + SpringApplication app = new SpringApplication(CustomComplexHeadersConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=springMessageSupplier;springMessageLogger", + "--spring.cloud.stream.bindings.springMessageSupplier-out-0.destination=cmh-3", + "--spring.cloud.stream.bindings.springMessageSupplier-out-0.producer.header-mode=none", + "--spring.cloud.stream.bindings.springMessageLogger-in-0.destination=cmh-3", + "--spring.cloud.stream.pulsar.bindings.springMessageLogger-in-0.consumer.subscription-name=pbit-cmh3-sub1")) { + // Wait for a few of the messages to flow through (check for index = 5) + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: test-headers-msg-5 w/ custom-id: null")); + } + } + + @Test + void consumerHeaderModeNone(CapturedOutput output) { + SpringApplication app = new SpringApplication(CustomComplexHeadersConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=springMessageSupplier;springMessageLogger", + "--spring.cloud.stream.bindings.springMessageSupplier-out-0.destination=cmh-4", + "--spring.cloud.stream.bindings.springMessageLogger-in-0.destination=cmh-4", + "--spring.cloud.stream.bindings.springMessageLogger-in-0.consumer.header-mode=none", + "--spring.cloud.stream.pulsar.bindings.springMessageLogger-in-0.consumer.subscription-name=pbit-cmh4-sub1")) { + // Wait for a few of the messages to flow through (check for index = 5) + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)) + .until(() -> output.toString().contains("Hello binder: test-headers-msg-5 w/ custom-id: null")); + } + } + + @Test + void customHeaderMapperRespected(CapturedOutput output) { + SpringApplication app = new SpringApplication(CustomHeaderMapperConfig.class); + app.setWebApplicationType(WebApplicationType.NONE); + try (ConfigurableApplicationContext ignored = app.run( + "--spring.pulsar.client.service-url=" + PulsarTestContainerSupport.getPulsarBrokerUrl(), + "--spring.pulsar.administration.service-url=" + PulsarTestContainerSupport.getHttpServiceUrl(), + "--spring.cloud.function.definition=springMessageSupplier;springMessageLogger", + "--spring.cloud.stream.bindings.springMessageSupplier-out-0.destination=cmh-5", + "--spring.cloud.stream.bindings.springMessageLogger-in-0.destination=cmh-5", + "--spring.cloud.stream.pulsar.bindings.springMessageLogger-in-0.consumer.subscription-name=pbit-cmh5-sub1")) { + // Wait for a few of the messages to flow through (check for index = 5) + Awaitility.await().atMost(Duration.ofSeconds(AWAIT_DURATION)).until(() -> output.toString() + .contains("Hello binder: test-headers-msg-5 w/ custom-id: tsh->tph->FooHeader[value=5150-5]")); + } + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class CustomSimpleHeadersConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + private int msgCount = 0; + + @Bean + public Supplier> springMessageSupplier() { + return () -> { + msgCount++; + return MessageBuilder.withPayload("test-headers-msg-" + msgCount) + .setHeader("custom-id", "5150-" + msgCount).build(); + }; + } + + @Bean + public Consumer> springMessageLogger() { + return s -> this.logger.info("Hello binder: {} w/ custom-id: {}", s.getPayload(), + s.getHeaders().get("custom-id")); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class CustomComplexHeadersConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + private int msgCount = 0; + + @Bean + public Supplier> springMessageSupplier() { + return () -> { + msgCount++; + return MessageBuilder.withPayload("test-headers-msg-" + msgCount) + .setHeader("custom-id", new FooHeader("5150-" + msgCount)).build(); + }; + } + + @Bean + public Consumer> springMessageLogger() { + return s -> { + var header = s.getHeaders().get("custom-id"); + if (header != null) { + assertThat(header).isInstanceOf(FooHeader.class); + } + this.logger.info("Hello binder: {} w/ custom-id: {}", s.getPayload(), header); + }; + } + + record FooHeader(String value) { + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class CustomHeaderMapperConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + private int msgCount = 0; + + @Bean + public PulsarHeaderMapper extendedToStringHeaderMapper() { + return new ToStringPulsarHeaderMapper(List.of("custom-id"), List.of("foo", "custom-id")) { + @Override + public Map toPulsarHeaders(MessageHeaders springHeaders) { + Map pulsarHeaders = super.toPulsarHeaders(springHeaders); + // foo and custom-id are allowed and expected + assertThat(pulsarHeaders).containsKeys("foo", "custom-id"); + return pulsarHeaders; + } + + @Override + public MessageHeaders toSpringHeaders(org.apache.pulsar.client.api.Message pulsarMessage) { + MessageHeaders springHeaders = super.toSpringHeaders(pulsarMessage); + // foo not allowed, custom-id allowed + assertThat(springHeaders).doesNotContainKey("foo").containsKey("custom-id"); + return springHeaders; + } + + @Override + protected String toPulsarHeaderValue(String name, Object value, Object context) { + return "tph->" + super.toPulsarHeaderValue(name, value, context); + } + + @Override + protected Object toSpringHeaderValue(String headerName, String rawHeader, Object context) { + return "tsh->" + super.toSpringHeaderValue(headerName, rawHeader, context); + } + }; + } + + @Bean + public Supplier> springMessageSupplier() { + return () -> { + msgCount++; + return MessageBuilder.withPayload("test-headers-msg-" + msgCount) + .setHeader("foo", "bar-" + msgCount) + .setHeader("custom-id", new FooHeader("5150-" + msgCount)).build(); + }; + } + + @Bean + public Consumer> springMessageLogger() { + return s -> { + var header = s.getHeaders().get("custom-id"); + if (header != null) { + assertThat(header).isInstanceOf(String.class); + } + var fooHeader = s.getHeaders().get("foo"); + assertThat(fooHeader).isNull(); + this.logger.info("Hello binder: {} w/ custom-id: {}", s.getPayload(), header); + }; + } + + record FooHeader(String value) { + } + + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class PrimitiveTextConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + @Bean + public Supplier textSupplier() { + return () -> "test-basic-scenario"; + } + + @Bean + public Consumer textLogger() { + return s -> this.logger.info("Hello binder: " + s); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + @Import(PrimitiveTextConfig.class) + static class BinderAndBindingPropsTestConfig { + + @SuppressWarnings("unchecked") + @Bean + public PulsarProducerFactory pulsarProducerFactory(PulsarClient pulsarClient, + PulsarProperties pulsarProperties, TopicResolver topicResolver) { + var customizer = (ProducerBuilderCustomizer) pulsarProperties.getProducer() + .toProducerBuilderCustomizer(); + return new TrackingProducerFactory(pulsarClient, pulsarProperties.getProducer().getTopicName(), customizer, + topicResolver); + } + + @SuppressWarnings("unchecked") + @Bean + public PulsarConsumerFactory pulsarConsumerFactory(PulsarClient pulsarClient, + PulsarProperties pulsarProperties) { + var customizer = (ConsumerBuilderCustomizer) pulsarProperties.getConsumer() + .toConsumerBuilderCustomizer(); + return new TrackingConsumerFactory(pulsarClient, customizer); + } + + } + + static class TrackingProducerFactory extends DefaultPulsarProducerFactory { + + List> producersCreated = new ArrayList<>(); + + TrackingProducerFactory(PulsarClient pulsarClient, @Nullable String defaultTopic, + ProducerBuilderCustomizer defaultConfigCustomizer, TopicResolver topicResolver) { + super(pulsarClient, defaultTopic, defaultConfigCustomizer, topicResolver); + } + + @Override + protected Producer doCreateProducer(Schema schema, @Nullable String topic, + @Nullable Collection encryptionKeys, + @Nullable List> producerBuilderCustomizers) + throws PulsarClientException { + Producer producer = super.doCreateProducer(schema, topic, encryptionKeys, + producerBuilderCustomizers); + producersCreated.add(producer); + return producer; + } + + } + + static class TrackingConsumerFactory extends DefaultPulsarConsumerFactory { + + List> consumersCreated = new ArrayList<>(); + + TrackingConsumerFactory(PulsarClient pulsarClient, ConsumerBuilderCustomizer defaultConsumerConfig) { + super(pulsarClient, defaultConsumerConfig); + } + + @Override + public org.apache.pulsar.client.api.Consumer createConsumer(Schema schema, + @Nullable Collection topics, @Nullable String subscriptionName, + @Nullable Map metadataProperties, + @Nullable List> consumerBuilderCustomizers) + throws PulsarClientException { + org.apache.pulsar.client.api.Consumer consumer = super.createConsumer(schema, topics, + subscriptionName, metadataProperties, consumerBuilderCustomizers); + consumersCreated.add(consumer); + return consumer; + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class PrimitiveFloatConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + @Bean + public Supplier piSupplier() { + return () -> 3.14f; + } + + @Bean + public Consumer piLogger() { + return f -> this.logger.info("Hello binder: " + f); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class JsonFooConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + @Bean + public Supplier fooSupplier() { + return () -> new Foo("5150"); + } + + @Bean + public Consumer fooLogger() { + return f -> this.logger.info("Hello binder: " + f); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + @Import(JsonFooConfig.class) + static class JsonFooWithCustomMappingConfig { + + @Bean + public SchemaResolverCustomizer customMappings() { + return (resolver) -> resolver.addCustomSchemaMapping(Foo.class, JSONSchema.of(Foo.class)); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class AvroUserConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + @Bean + public Supplier userSupplier() { + return () -> new User("user21", 21); + } + + @Bean + public Consumer userLogger() { + return f -> this.logger.info("Hello binder: " + f); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + @Import(AvroUserConfig.class) + static class AvroUserConfigCustomMappings { + + @Bean + public SchemaResolverCustomizer customMappings() { + return (resolver) -> resolver.addCustomSchemaMapping(User.class, Schema.AVRO(User.class)); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class KeyValueAvroUserConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + @Bean + public Supplier> userSupplier() { + return () -> new KeyValue<>("21", new User("user21", 21)); + } + + @Bean + public Consumer> userLogger() { + return f -> this.logger.info("Hello binder: " + f.getKey() + "->" + f.getValue()); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + @Import(KeyValueAvroUserConfig.class) + static class KeyValueAvroUserConfigCustomMappings { + + @Bean + public SchemaResolverCustomizer customMappings() { + return (resolver) -> resolver.addCustomSchemaMapping(User.class, Schema.AVRO(User.class)); + } + + } + + @EnableAutoConfiguration + @SpringBootConfiguration + static class KeyValueJsonFooConfig { + + private final Logger logger = LoggerFactory.getLogger(getClass()); + + @Bean + public Supplier> fooSupplier() { + return () -> new KeyValue<>("5150", new Foo("5150")); + } + + @Bean + public Consumer> fooLogger() { + return f -> this.logger.info("Hello binder: " + f.getKey() + "->" + f.getValue()); + } + + } + + record Foo(String value) { + } + + /** + * Do not convert this to a Record as Avro does not seem to work well w/ records. + */ + static class User { + + private String name; + + private int age; + + User() { + } + + User(String name, int age) { + this.name = name; + this.age = age; + } + + public String getName() { + return name; + } + + public void setName(String name) { + this.name = name; + } + + public int getAge() { + return age; + } + + public void setAge(int age) { + this.age = age; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (o == null || getClass() != o.getClass()) { + return false; + } + User user = (User) o; + return age == user.age && Objects.equals(name, user.name); + } + + @Override + public int hashCode() { + return Objects.hash(name, age); + } + + @Override + public String toString() { + return "User{" + "name='" + name + '\'' + ", age=" + age + '}'; + } + + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderTests.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderTests.java new file mode 100644 index 000000000..a412d3890 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderTests.java @@ -0,0 +1,252 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.nio.charset.StandardCharsets; +import java.util.UUID; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.pulsar.client.api.PulsarClient; +import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.SubscriptionInitialPosition; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Disabled; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInfo; + +import org.springframework.cloud.stream.binder.Binder; +import org.springframework.cloud.stream.binder.Binding; +import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; +import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.PartitionCapableBinderTests; +import org.springframework.cloud.stream.binder.Spy; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarConsumerProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarProducerProperties; +import org.springframework.cloud.stream.binder.pulsar.provisioning.PulsarTopicProvisioner; +import org.springframework.cloud.stream.config.BindingProperties; +import org.springframework.integration.channel.DirectChannel; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.support.MessageBuilder; +import org.springframework.lang.Nullable; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.MessageHeaders; +import org.springframework.pulsar.core.DefaultPulsarConsumerFactory; +import org.springframework.pulsar.core.DefaultPulsarProducerFactory; +import org.springframework.pulsar.core.DefaultSchemaResolver; +import org.springframework.pulsar.core.PulsarAdministration; +import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.support.header.JsonPulsarHeaderMapper; +import org.springframework.util.Assert; +import org.springframework.util.MimeTypeUtils; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests for {@link PulsarMessageChannelBinder}. + * + * @author Soby Chacko + */ +public class PulsarBinderTests extends + PartitionCapableBinderTests, ExtendedProducerProperties> + implements PulsarTestContainerSupport { + + private PulsarTestBinder binder; + + @Nullable + protected PulsarClient pulsarClient; + + @BeforeEach + void createPulsarClient() throws PulsarClientException { + pulsarClient = PulsarClient.builder().serviceUrl(PulsarTestContainerSupport.getPulsarBrokerUrl()).build(); + } + + @AfterEach + void closePulsarClient() throws PulsarClientException { + if (pulsarClient != null && !pulsarClient.isClosed()) { + pulsarClient.close(); + } + } + + @Override + protected boolean usesExplicitRouting() { + return false; + } + + @Override + protected String getClassUnderTestName() { + return PulsarMessageChannelBinder.class.getSimpleName(); + } + + @Override + protected PulsarTestBinder getBinder() { + var pulsarAdministration = new PulsarAdministration(PulsarTestContainerSupport.getHttpServiceUrl()); + var configProps = new PulsarBinderConfigurationProperties(); + var provisioner = new PulsarTopicProvisioner(pulsarAdministration, configProps); + var producerFactory = new DefaultPulsarProducerFactory<>(pulsarClient); + var pulsarTemplate = new PulsarTemplate<>(producerFactory); + var consumerFactory = new DefaultPulsarConsumerFactory<>(pulsarClient, + (consumerBuilder -> consumerBuilder.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest))); + if (this.binder == null) { + this.binder = new PulsarTestBinder(provisioner, pulsarTemplate, consumerFactory, configProps, + new DefaultSchemaResolver(), JsonPulsarHeaderMapper.builder().build()); + } + return this.binder; + } + + @Override + protected ExtendedConsumerProperties createConsumerProperties() { + final ExtendedConsumerProperties pulsarConsumerProperties = new ExtendedConsumerProperties<>( + new PulsarConsumerProperties()); + return pulsarConsumerProperties; + } + + @Override + public Spy spyOn(String name) { + return null; + } + + private ExtendedProducerProperties createProducerProperties() { + return this.createProducerProperties(null); + } + + @Override + protected ExtendedProducerProperties createProducerProperties(TestInfo testInto) { + return new ExtendedProducerProperties<>(new PulsarProducerProperties()); + } + + @Override + protected void binderBindUnbindLatency() throws InterruptedException { + Thread.sleep(500); + } + + @Test + @Override + @SuppressWarnings({ "unchecked", "rawtypes" }) + public void testSendAndReceive(TestInfo testInfo) throws Exception { + Binder binder = getBinder(); + BindingProperties outputBindingProperties = createProducerBindingProperties(createProducerProperties()); + + DirectChannel moduleOutputChannel = createBindableChannel("output", outputBindingProperties); + ExtendedConsumerProperties consumerProperties = createConsumerProperties(); + DirectChannel moduleInputChannel = createBindableChannel("input", + createConsumerBindingProperties(consumerProperties)); + + Binding producerBinding = binder.bindProducer("foo.bar", moduleOutputChannel, + outputBindingProperties.getProducer()); + Binding consumerBinding = binder.bindConsumer("foo.bar", null, moduleInputChannel, + consumerProperties); + + Message message = MessageBuilder + .withPayload("foo".getBytes(StandardCharsets.UTF_8)).build(); + + // Let the consumer actually bind to the producer before sending a msg + binderBindUnbindLatency(); + moduleOutputChannel.send(message); + CountDownLatch latch = new CountDownLatch(1); + AtomicReference> inboundMessageRef = new AtomicReference<>(); + moduleInputChannel.subscribe(message1 -> { + try { + inboundMessageRef.set((Message) message1); + } + finally { + latch.countDown(); + } + }); + Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message"); + + assertThat(inboundMessageRef.get()).isNotNull(); + assertThat(new String(inboundMessageRef.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo("foo"); + + producerBinding.unbind(); + consumerBinding.unbind(); + } + + @Test + @Override + @SuppressWarnings({ "unchecked", "rawtypes" }) + public void testAnonymousGroup(TestInfo testInfo) throws Exception { + Binder binder = getBinder(); + BindingProperties producerBindingProperties = createProducerBindingProperties( + createProducerProperties(testInfo)); + DirectChannel output = createBindableChannel("output", producerBindingProperties); + Binding producerBinding = binder.bindProducer( + String.format("defaultGroup%s0", getDestinationNameDelimiter()), output, + producerBindingProperties.getProducer()); + + QueueChannel input1 = new QueueChannel(); + Binding binding1 = binder.bindConsumer( + String.format("defaultGroup%s0", getDestinationNameDelimiter()), null, input1, + createConsumerProperties()); + + QueueChannel input2 = new QueueChannel(); + Binding binding2 = binder.bindConsumer( + String.format("defaultGroup%s0", getDestinationNameDelimiter()), null, input2, + createConsumerProperties()); + + String testPayload1 = "foo-" + UUID.randomUUID(); + output.send(MessageBuilder.withPayload(testPayload1) + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build()); + + Message receivedMessage1 = (Message) receive(input1); + assertThat(receivedMessage1).isNotNull(); + assertThat(new String(receivedMessage1.getPayload())).isEqualTo(testPayload1); + + Message receivedMessage2 = (Message) receive(input2); + assertThat(receivedMessage2).isNotNull(); + assertThat(new String(receivedMessage2.getPayload())).isEqualTo(testPayload1); + + binding2.unbind(); + + String testPayload2 = "foo-" + UUID.randomUUID(); + output.send(MessageBuilder.withPayload(testPayload2) + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build()); + + binding2 = binder.bindConsumer(String.format("defaultGroup%s0", getDestinationNameDelimiter()), null, input2, + createConsumerProperties()); + String testPayload3 = "foo-" + UUID.randomUUID(); + output.send(MessageBuilder.withPayload(testPayload3) + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build()); + + receivedMessage1 = (Message) receive(input1); + assertThat(receivedMessage1).isNotNull(); + assertThat(new String(receivedMessage1.getPayload())).isEqualTo(testPayload2); + receivedMessage1 = (Message) receive(input1); + assertThat(receivedMessage1).isNotNull(); + assertThat(new String(receivedMessage1.getPayload())).isEqualTo(testPayload3); + + receivedMessage2 = (Message) receive(input2); + assertThat(receivedMessage2).isNotNull(); + assertThat(new String(receivedMessage2.getPayload())).isEqualTo(testPayload1); + + producerBinding.unbind(); + binding1.unbind(); + binding2.unbind(); + } + + @Test + @Override + @Disabled + public void testPartitionedModuleSpEL(TestInfo testInfo) { + // This use-case needs to be further evaluated for Pulsar binder. + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderUtilsTests.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderUtilsTests.java new file mode 100644 index 000000000..17805d225 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarBinderUtilsTests.java @@ -0,0 +1,324 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.stream.Stream; + +import org.apache.pulsar.client.api.CompressionType; +import org.apache.pulsar.client.api.ConsumerCryptoFailureAction; +import org.apache.pulsar.client.api.DeadLetterPolicy; +import org.apache.pulsar.client.api.HashingScheme; +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.SubscriptionInitialPosition; +import org.apache.pulsar.client.api.SubscriptionMode; +import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils; +import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; +import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; +import org.assertj.core.api.InstanceOfAssertFactories; +import org.junit.jupiter.api.Nested; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.Arguments; +import org.junit.jupiter.params.provider.MethodSource; + +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.cloud.stream.binder.pulsar.properties.PulsarConsumerProperties; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.pulsar.autoconfigure.ConsumerConfigProperties; +import org.springframework.pulsar.autoconfigure.ProducerConfigProperties; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatNoException; +import static org.junit.jupiter.params.provider.Arguments.arguments; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +/** + * Unit tests for {@link PulsarBinderUtils}. + * + * @author Soby Chacko + * @author Chris Bono + */ +public class PulsarBinderUtilsTests { + + @Nested + class SubscriptionNameTests { + + @Test + void respectsValueWhenSetAsProperty() { + var consumerDestination = mock(ConsumerDestination.class); + var pulsarConsumerProperties = mock(PulsarConsumerProperties.class); + when(pulsarConsumerProperties.getSubscriptionName()).thenReturn("my-sub"); + assertThat(PulsarBinderUtils.subscriptionName(pulsarConsumerProperties, consumerDestination)) + .isEqualTo("my-sub"); + } + + @Test + void generatesValueWhenNotSetAsProperty() { + var consumerDestination = mock(ConsumerDestination.class); + var pulsarConsumerProperties = mock(PulsarConsumerProperties.class); + when(pulsarConsumerProperties.getSubscriptionName()).thenReturn(null); + when(consumerDestination.getName()).thenReturn("my-topic"); + assertThat(PulsarBinderUtils.subscriptionName(pulsarConsumerProperties, consumerDestination)) + .startsWith("my-topic-anon-subscription-"); + } + + } + + @Nested + class MergedPropertiesTests { + + @ParameterizedTest(name = "{0}") + @MethodSource("mergePropertiesTestProvider") + void mergePropertiesTest(String testName, Map baseProps, Map binderProps, + Map bindingProps, Map expectedMergedProps) { + assertThat(PulsarBinderUtils.mergePropertiesWithPrecedence(baseProps, binderProps, bindingProps)) + .containsExactlyInAnyOrderEntriesOf(expectedMergedProps); + } + + // @formatter:off + static Stream mergePropertiesTestProvider() { + return Stream.of( + arguments("binderLevelContainsSamePropAsBaseWithDiffValue", + Map.of("foo", "foo-base"), + Map.of("foo", "foo-binder"), + Collections.emptyMap(), + Map.of("foo", "foo-binder")), + arguments("binderLevelContainsNewPropNotInBase", + Collections.emptyMap(), + Map.of("foo", "foo-binder"), + Collections.emptyMap(), + Map.of("foo", "foo-binder")), + arguments("binderLevelContainsSamePropAsBaseWithSameValue", + Map.of("foo", "foo-base"), + Map.of("foo", "foo-base"), + Collections.emptyMap(), + Collections.emptyMap()), + arguments("bindingLevelContainsSamePropAsBaseWithDiffValue", + Map.of("foo", "foo-base"), + Collections.emptyMap(), + Map.of("foo", "foo-binding"), + Map.of("foo", "foo-binding")), + arguments("bindingLevelContainsNewPropNotInBase", + Collections.emptyMap(), + Map.of("foo", "foo-binding"), + Collections.emptyMap(), + Map.of("foo", "foo-binding")), + arguments("bindingLevelContainsSamePropAsBaseWithSameValue", + Map.of("foo", "foo-base"), + Collections.emptyMap(), + Map.of("foo", "foo-base"), + Collections.emptyMap()), + arguments("bindingOverridesBinder", + Map.of("bar", "bar-base"), + Map.of("foo", "foo-binder"), + Map.of("foo", "foo-binding"), + Map.of("foo", "foo-binding")), + arguments("binderOverridesBaseAndBindingOverridesBinder", + Map.of("foo", "foo-base"), + Map.of("foo", "foo-binder"), + Map.of("foo", "foo-binding"), + Map.of("foo", "foo-binding")), + arguments("onlyBaseProps", + Map.of("foo", "foo-base"), + Collections.emptyMap(), + Collections.emptyMap(), + Collections.emptyMap())); + } + // @formatter:on + + } + + @Nested + class ConvertedPropertiesTests { + + private final ProducerConfigProperties properties = new ProducerConfigProperties(); + + private void bind(Map map) { + ConfigurationPropertySource source = new MapConfigurationPropertySource(map); + new Binder(source).bind("spring.pulsar.producer", Bindable.ofInstance(this.properties)); + } + + @Test + void producerPropertiesToMap() { + Map props = new HashMap<>(); + props.put("spring.pulsar.producer.topic-name", "my-topic"); + props.put("spring.pulsar.producer.producer-name", "my-producer"); + props.put("spring.pulsar.producer.send-timeout", "2s"); + props.put("spring.pulsar.producer.block-if-queue-full", "true"); + props.put("spring.pulsar.producer.max-pending-messages", "3"); + props.put("spring.pulsar.producer.max-pending-messages-across-partitions", "4"); + props.put("spring.pulsar.producer.message-routing-mode", "custompartition"); + props.put("spring.pulsar.producer.hashing-scheme", "murmur3_32hash"); + props.put("spring.pulsar.producer.crypto-failure-action", "send"); + props.put("spring.pulsar.producer.batching-max-publish-delay", "5s"); + props.put("spring.pulsar.producer.batching-partition-switch-frequency-by-publish-delay", "6"); + props.put("spring.pulsar.producer.batching-max-messages", "7"); + props.put("spring.pulsar.producer.batching-max-bytes", "8"); + props.put("spring.pulsar.producer.batching-enabled", "false"); + props.put("spring.pulsar.producer.chunking-enabled", "true"); + props.put("spring.pulsar.producer.encryption-keys[0]", "my-key"); + props.put("spring.pulsar.producer.compression-type", "lz4"); + props.put("spring.pulsar.producer.initial-sequence-id", "9"); + props.put("spring.pulsar.producer.producer-access-mode", "exclusive"); + props.put("spring.pulsar.producer.lazy-start=partitioned-producers", "true"); + props.put("spring.pulsar.producer.properties[my-prop]", "my-prop-value"); + + bind(props); + Map producerProps = PulsarBinderUtils.convertProducerPropertiesToMap(properties); + + // Verify that the props can be loaded in a ProducerBuilder + assertThatNoException().isThrownBy(() -> ConfigurationDataUtils.loadData(producerProps, + new ProducerConfigurationData(), ProducerConfigurationData.class)); + + assertThat(producerProps).containsEntry("topicName", "my-topic") + .containsEntry("producerName", "my-producer").containsEntry("sendTimeoutMs", 2_000) + .containsEntry("blockIfQueueFull", true).containsEntry("maxPendingMessages", 3) + .containsEntry("maxPendingMessagesAcrossPartitions", 4) + .containsEntry("messageRoutingMode", MessageRoutingMode.CustomPartition) + .containsEntry("hashingScheme", HashingScheme.Murmur3_32Hash) + .containsEntry("cryptoFailureAction", ProducerCryptoFailureAction.SEND) + .containsEntry("batchingMaxPublishDelayMicros", 5_000_000L) + .containsEntry("batchingPartitionSwitchFrequencyByPublishDelay", 6) + .containsEntry("batchingMaxMessages", 7).containsEntry("batchingMaxBytes", 8) + .containsEntry("batchingEnabled", false).containsEntry("chunkingEnabled", true) + .hasEntrySatisfying("encryptionKeys", + keys -> assertThat(keys).asInstanceOf(InstanceOfAssertFactories.collection(String.class)) + .containsExactly("my-key")) + .containsEntry("compressionType", CompressionType.LZ4).containsEntry("initialSequenceId", 9L) + .containsEntry("accessMode", ProducerAccessMode.Exclusive) + .containsEntry("lazyStartPartitionedProducers", true).hasEntrySatisfying("properties", + properties -> assertThat(properties) + .asInstanceOf(InstanceOfAssertFactories.map(String.class, String.class)) + .containsEntry("my-prop", "my-prop-value")); + } + + } + + @Nested + class ConvertedConsumerPropertiesTests { + + private final ConsumerConfigProperties properties = new ConsumerConfigProperties(); + + private void bind(Map map) { + ConfigurationPropertySource source = new MapConfigurationPropertySource(map); + new Binder(source).bind("spring.pulsar.consumer", Bindable.ofInstance(this.properties)); + } + + @Test + void consumerPropertiesToMap() { + Map props = new HashMap<>(); + props.put("spring.pulsar.consumer.topics[0]", "my-topic"); + props.put("spring.pulsar.consumer.topics-pattern", "my-pattern"); + props.put("spring.pulsar.consumer.subscription-name", "my-subscription"); + props.put("spring.pulsar.consumer.subscription-type", "shared"); + props.put("spring.pulsar.consumer.subscription-properties[my-sub-prop]", "my-sub-prop-value"); + props.put("spring.pulsar.consumer.subscription-mode", "nondurable"); + props.put("spring.pulsar.consumer.receiver-queue-size", "1"); + props.put("spring.pulsar.consumer.acknowledgements-group-time", "2s"); + props.put("spring.pulsar.consumer.negative-ack-redelivery-delay", "3s"); + props.put("spring.pulsar.consumer.max-total-receiver-queue-size-across-partitions", "5"); + props.put("spring.pulsar.consumer.consumer-name", "my-consumer"); + props.put("spring.pulsar.consumer.ack-timeout", "6s"); + props.put("spring.pulsar.consumer.tick-duration", "7s"); + props.put("spring.pulsar.consumer.priority-level", "8"); + props.put("spring.pulsar.consumer.crypto-failure-action", "discard"); + props.put("spring.pulsar.consumer.properties[my-prop]", "my-prop-value"); + props.put("spring.pulsar.consumer.read-compacted", "true"); + props.put("spring.pulsar.consumer.subscription-initial-position", "earliest"); + props.put("spring.pulsar.consumer.pattern-auto-discovery-period", "9"); + props.put("spring.pulsar.consumer.regex-subscription-mode", "all-topics"); + props.put("spring.pulsar.consumer.dead-letter-policy.max-redeliver-count", "4"); + props.put("spring.pulsar.consumer.dead-letter-policy.retry-letter-topic", "my-retry-topic"); + props.put("spring.pulsar.consumer.dead-letter-policy.dead-letter-topic", "my-dlt-topic"); + props.put("spring.pulsar.consumer.dead-letter-policy.initial-subscription-name", "my-initial-subscription"); + props.put("spring.pulsar.consumer.retry-enable", "true"); + props.put("spring.pulsar.consumer.auto-update-partitions", "false"); + props.put("spring.pulsar.consumer.auto-update-partitions-interval", "10s"); + props.put("spring.pulsar.consumer.replicate-subscription-state", "true"); + props.put("spring.pulsar.consumer.reset-include-head", "true"); + props.put("spring.pulsar.consumer.batch-index-ack-enabled", "true"); + props.put("spring.pulsar.consumer.ack-receipt-enabled", "true"); + props.put("spring.pulsar.consumer.pool-messages", "true"); + props.put("spring.pulsar.consumer.start-paused", "true"); + props.put("spring.pulsar.consumer.auto-ack-oldest-chunked-message-on-queue-full", "false"); + props.put("spring.pulsar.consumer.max-pending-chunked-message", "11"); + props.put("spring.pulsar.consumer.expire-time-of-incomplete-chunked-message", "12s"); + + bind(props); + Map consumerProps = PulsarBinderUtils.convertConsumerPropertiesToMap(properties); + + // Verify that the props can be loaded in a ConsumerBuilder + assertThatNoException().isThrownBy(() -> ConfigurationDataUtils.loadData(consumerProps, + new ConsumerConfigurationData<>(), ConsumerConfigurationData.class)); + + assertThat(consumerProps) + .hasEntrySatisfying("topicNames", + topics -> assertThat(topics) + .asInstanceOf(InstanceOfAssertFactories.collection(String.class)) + .containsExactly("my-topic")) + .hasEntrySatisfying("topicsPattern", p -> assertThat(p.toString()).isEqualTo("my-pattern")) + .containsEntry("subscriptionName", "my-subscription") + .containsEntry("subscriptionType", SubscriptionType.Shared) + .hasEntrySatisfying("subscriptionProperties", + properties -> assertThat(properties) + .asInstanceOf(InstanceOfAssertFactories.map(String.class, String.class)) + .containsEntry("my-sub-prop", "my-sub-prop-value")) + .containsEntry("subscriptionMode", SubscriptionMode.NonDurable) + .containsEntry("receiverQueueSize", 1).containsEntry("acknowledgementsGroupTimeMicros", 2_000_000L) + .containsEntry("negativeAckRedeliveryDelayMicros", 3_000_000L) + .containsEntry("maxTotalReceiverQueueSizeAcrossPartitions", 5) + .containsEntry("consumerName", "my-consumer").containsEntry("ackTimeoutMillis", 6_000L) + .containsEntry("tickDurationMillis", 7_000L).containsEntry("priorityLevel", 8) + .containsEntry("cryptoFailureAction", ConsumerCryptoFailureAction.DISCARD) + .hasEntrySatisfying("properties", + properties -> assertThat(properties) + .asInstanceOf(InstanceOfAssertFactories.map(String.class, String.class)) + .containsEntry("my-prop", "my-prop-value")) + .containsEntry("readCompacted", true) + .containsEntry("subscriptionInitialPosition", SubscriptionInitialPosition.Earliest) + .containsEntry("patternAutoDiscoveryPeriod", 9) + .containsEntry("regexSubscriptionMode", RegexSubscriptionMode.AllTopics) + .hasEntrySatisfying("deadLetterPolicy", dlp -> { + DeadLetterPolicy deadLetterPolicy = (DeadLetterPolicy) dlp; + assertThat(deadLetterPolicy.getMaxRedeliverCount()).isEqualTo(4); + assertThat(deadLetterPolicy.getRetryLetterTopic()).isEqualTo("my-retry-topic"); + assertThat(deadLetterPolicy.getDeadLetterTopic()).isEqualTo("my-dlt-topic"); + assertThat(deadLetterPolicy.getInitialSubscriptionName()).isEqualTo("my-initial-subscription"); + }).containsEntry("retryEnable", true).containsEntry("autoUpdatePartitions", false) + .containsEntry("autoUpdatePartitionsIntervalSeconds", 10L) + .containsEntry("replicateSubscriptionState", true).containsEntry("resetIncludeHead", true) + .containsEntry("batchIndexAckEnabled", true).containsEntry("ackReceiptEnabled", true) + .containsEntry("poolMessages", true).containsEntry("startPaused", true) + .containsEntry("autoAckOldestChunkedMessageOnQueueFull", false) + .containsEntry("maxPendingChunkedMessage", 11) + .containsEntry("expireTimeOfIncompleteChunkedMessageMillis", 12_000L); + } + + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarExtendedBindingPropertiesTests.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarExtendedBindingPropertiesTests.java new file mode 100644 index 000000000..c272e064b --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarExtendedBindingPropertiesTests.java @@ -0,0 +1,140 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.util.HashMap; +import java.util.Map; + +import org.apache.pulsar.client.api.ProducerAccessMode; +import org.apache.pulsar.client.api.SubscriptionMode; +import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.client.impl.conf.ConfigurationDataUtils; +import org.apache.pulsar.client.impl.conf.ConsumerConfigurationData; +import org.apache.pulsar.client.impl.conf.ProducerConfigurationData; +import org.assertj.core.api.InstanceOfAssertFactories; +import org.junit.jupiter.api.Test; + +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.cloud.stream.binder.pulsar.properties.PulsarExtendedBindingProperties; +import org.springframework.pulsar.listener.PulsarContainerProperties; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatNoException; + +/** + * Tests for {@link PulsarExtendedBindingProperties}. + * + * @author Chris Bono + */ +public class PulsarExtendedBindingPropertiesTests { + + private final PulsarExtendedBindingProperties properties = new PulsarExtendedBindingProperties(); + + private void bind(Map map) { + ConfigurationPropertySource source = new MapConfigurationPropertySource(map); + new Binder(source).bind("spring.cloud.stream.pulsar", Bindable.ofInstance(this.properties)); + } + + @Test + void producerProperties() { + // Only spot check a few values (PulsarPropertiesTests does the heavy lifting) + Map props = new HashMap<>(); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.producer.topic-name", "my-topic"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.producer.send-timeout", "2s"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.producer.max-pending-messages", "3"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.producer.producer-access-mode", "exclusive"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.producer.properties[my-prop]", "my-prop-value"); + + bind(props); + + assertThat(properties.getBindings()).containsOnlyKeys("my-foo"); + Map producerProps = PulsarBinderUtils + .convertProducerPropertiesToMap(properties.getExtendedProducerProperties("my-foo")); + // Verify that the props can be loaded in a ProducerBuilder + assertThatNoException().isThrownBy(() -> ConfigurationDataUtils.loadData(producerProps, + new ProducerConfigurationData(), ProducerConfigurationData.class)); + // @formatter:off + assertThat(producerProps) + .containsEntry("topicName", "my-topic") + .containsEntry("sendTimeoutMs", 2_000) + .containsEntry("maxPendingMessages", 3) + .containsEntry("accessMode", ProducerAccessMode.Exclusive) + .hasEntrySatisfying("properties", properties -> + assertThat(properties) + .asInstanceOf(InstanceOfAssertFactories.map(String.class, String.class)) + .containsEntry("my-prop", "my-prop-value")); + // @formatter:on + } + + @Test + void consumerProperties() { + // Only spot check a few values (PulsarPropertiesTests does the heavy lifting) + Map props = new HashMap<>(); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.consumer.topics[0]", "my-topic"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.consumer.subscription-properties[my-sub-prop]", + "my-sub-prop-value"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.consumer.subscription-mode", "nondurable"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.consumer.receiver-queue-size", "1"); + + bind(props); + + assertThat(properties.getBindings()).containsOnlyKeys("my-foo"); + Map consumerProps = PulsarBinderUtils + .convertConsumerPropertiesToMap(properties.getExtendedConsumerProperties("my-foo")); + // Verify that the props can be loaded in a ConsumerBuilder + assertThatNoException().isThrownBy(() -> ConfigurationDataUtils.loadData(consumerProps, + new ConsumerConfigurationData<>(), ConsumerConfigurationData.class)); + // @formatter:off + assertThat(consumerProps) + .hasEntrySatisfying("topicNames", + topics -> assertThat(topics).asInstanceOf(InstanceOfAssertFactories.collection(String.class)) + .containsExactly("my-topic")) + .hasEntrySatisfying("subscriptionProperties", + properties -> assertThat(properties) + .asInstanceOf(InstanceOfAssertFactories.map(String.class, String.class)) + .containsEntry("my-sub-prop", "my-sub-prop-value")) + .containsEntry("subscriptionMode", SubscriptionMode.NonDurable) + .containsEntry("receiverQueueSize", 1); + // @formatter:on + } + + @Test + void extendedBindingsArePropagatedToContainerProperties() { + Map props = new HashMap<>(); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.consumer.subscription-name", "my-foo-sbscription"); + props.put("spring.cloud.stream.pulsar.bindings.my-foo.consumer.subscription-type", "Shared"); + + bind(props); + + var bindingConsumerProps = PulsarBinderUtils + .convertConsumerPropertiesToMap(properties.getExtendedConsumerProperties("my-foo")); + PulsarContainerProperties pulsarContainerProperties = new PulsarContainerProperties(); + pulsarContainerProperties.getPulsarConsumerProperties().putAll(bindingConsumerProps); + + assertThat(pulsarContainerProperties.getSubscriptionName()).isNull(); + assertThat(pulsarContainerProperties.getSubscriptionType()).isEqualTo(SubscriptionType.Exclusive); + + pulsarContainerProperties.updateContainerProperties(); + + assertThat(pulsarContainerProperties.getSubscriptionName()).isEqualTo("my-foo-sbscription"); + assertThat(pulsarContainerProperties.getSubscriptionType()).isEqualTo(SubscriptionType.Shared); + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTestBinder.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTestBinder.java new file mode 100644 index 000000000..6040c54cc --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTestBinder.java @@ -0,0 +1,63 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.pulsar.provisioning.PulsarTopicProvisioner; +import org.springframework.context.annotation.AnnotationConfigApplicationContext; +import org.springframework.context.annotation.Configuration; +import org.springframework.integration.config.EnableIntegration; +import org.springframework.pulsar.core.PulsarConsumerFactory; +import org.springframework.pulsar.core.PulsarTemplate; +import org.springframework.pulsar.core.SchemaResolver; +import org.springframework.pulsar.support.header.PulsarHeaderMapper; + +/** + * Test binder to exercise producer/consumer bindings in + * {@link PulsarMessageChannelBinder}. + * + * @author Soby Chacko + * @author Chris Bono + */ +public class PulsarTestBinder extends AbstractPulsarTestBinder { + + @SuppressWarnings({ "unchecked" }) + PulsarTestBinder(PulsarTopicProvisioner pulsarTopicProvisioner, PulsarTemplate pulsarTemplate, + PulsarConsumerFactory pulsarConsumerFactory, PulsarBinderConfigurationProperties binderConfigProps, + SchemaResolver schemaResolver, PulsarHeaderMapper headerMapper) { + + try { + var binder = new PulsarMessageChannelBinder(pulsarTopicProvisioner, (PulsarTemplate) pulsarTemplate, + pulsarConsumerFactory, binderConfigProps, schemaResolver, headerMapper); + var context = new AnnotationConfigApplicationContext(Config.class); + setApplicationContext(context); + binder.setApplicationContext(context); + binder.afterPropertiesSet(); + this.setPollableConsumerBinder(binder); + } + catch (Exception e) { + throw new RuntimeException(e); + } + } + + @Configuration + @EnableIntegration + static class Config { + + } + +} diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTestContainerSupport.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTestContainerSupport.java new file mode 100644 index 000000000..b8bddb838 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTestContainerSupport.java @@ -0,0 +1,68 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import java.util.Locale; + +import org.junit.jupiter.api.BeforeAll; +import org.testcontainers.containers.PulsarContainer; +import org.testcontainers.junit.jupiter.Testcontainers; +import org.testcontainers.utility.DockerImageName; + +/** + * Provides a static {@link PulsarContainer} that can be shared across test classes. + * + * @author Chris Bono + */ +@Testcontainers(disabledWithoutDocker = true) +public interface PulsarTestContainerSupport { + + PulsarContainer PULSAR_CONTAINER = new PulsarContainer(getPulsarImage()); + + @BeforeAll + static void startContainer() { + PULSAR_CONTAINER.start(); + } + + static String getPulsarBrokerUrl() { + return PULSAR_CONTAINER.getPulsarBrokerUrl(); + } + + static DockerImageName getPulsarImage() { + return isRunningOnMacM1() ? getMacM1PulsarImage() : getStandardPulsarImage(); + } + + static String getHttpServiceUrl() { + return PULSAR_CONTAINER.getHttpServiceUrl(); + } + + private static boolean isRunningOnMacM1() { + String osName = System.getProperty("os.name").toLowerCase(Locale.ENGLISH); + String osArchitecture = System.getProperty("os.arch").toLowerCase(Locale.ENGLISH); + return osName.contains("mac") && osArchitecture.equals("aarch64"); + } + + private static DockerImageName getStandardPulsarImage() { + return DockerImageName.parse("apachepulsar/pulsar:2.11.0"); + } + + private static DockerImageName getMacM1PulsarImage() { + return DockerImageName.parse("kezhenxu94/pulsar").asCompatibleSubstituteFor("apachepulsar/pulsar"); + } + +} + diff --git a/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTopicProvisionerTests.java b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTopicProvisionerTests.java new file mode 100644 index 000000000..fb60f0e71 --- /dev/null +++ b/binders/pulsar-binder/spring-cloud-stream-binder-pulsar/src/test/java/org/springframework/cloud/stream/binder/pulsar/PulsarTopicProvisionerTests.java @@ -0,0 +1,136 @@ +/* + * Copyright 2023-2023 the original author or authors. + * + * Licensed 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 + * + * https://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. + */ + +package org.springframework.cloud.stream.binder.pulsar; + +import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; + +import org.springframework.cloud.stream.binder.ExtendedConsumerProperties; +import org.springframework.cloud.stream.binder.ExtendedProducerProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarBinderConfigurationProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarConsumerProperties; +import org.springframework.cloud.stream.binder.pulsar.properties.PulsarProducerProperties; +import org.springframework.cloud.stream.binder.pulsar.provisioning.PulsarTopicProvisioner; +import org.springframework.cloud.stream.provisioning.ConsumerDestination; +import org.springframework.cloud.stream.provisioning.ProducerDestination; +import org.springframework.pulsar.core.PulsarAdministration; +import org.springframework.pulsar.core.PulsarTopic; + +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.times; +import static org.mockito.Mockito.verify; + +/** + * @author Soby Chacko + */ +public class PulsarTopicProvisionerTests { + + @Test + void provisionThroughProducerBindingWithDefaultPartitioning() { + PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties = new PulsarBinderConfigurationProperties(); + PulsarAdministration pulsarAdministration = mock(PulsarAdministration.class); + PulsarTopicProvisioner pulsarTopicProvisioner = new PulsarTopicProvisioner(pulsarAdministration, + pulsarBinderConfigurationProperties); + ExtendedProducerProperties properties = new ExtendedProducerProperties<>( + new PulsarProducerProperties()); + ProducerDestination producerDestination = pulsarTopicProvisioner.provisionProducerDestination("foo", + properties); + verifyAndAssert(pulsarAdministration, producerDestination.getName(), "foo", 0); + } + + private static void verifyAndAssert(PulsarAdministration pulsarAdministration, String actualProducerDestination, + String expectedProducerDestination, int expectedPartitionCount) { + ArgumentCaptor pulsarTopicArgumentCaptor = ArgumentCaptor.forClass(PulsarTopic.class); + verify(pulsarAdministration, times(1)).createOrModifyTopics(pulsarTopicArgumentCaptor.capture()); + assertThat(actualProducerDestination).isEqualTo(expectedProducerDestination); + PulsarTopic pulsarTopic = pulsarTopicArgumentCaptor.getValue(); + assertThat(pulsarTopic.topicName()).isEqualTo(expectedProducerDestination); + assertThat(pulsarTopic.numberOfPartitions()).isEqualTo(expectedPartitionCount); + } + + @Test + void provisionThroughConsumerBindingWithDefaultPartitioning() { + PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties = new PulsarBinderConfigurationProperties(); + PulsarAdministration pulsarAdministration = mock(PulsarAdministration.class); + PulsarTopicProvisioner pulsarTopicProvisioner = new PulsarTopicProvisioner(pulsarAdministration, + pulsarBinderConfigurationProperties); + ExtendedConsumerProperties properties = new ExtendedConsumerProperties<>( + new PulsarConsumerProperties()); + ConsumerDestination consumerDestination = pulsarTopicProvisioner.provisionConsumerDestination("bar", "", + properties); + verifyAndAssert(pulsarAdministration, consumerDestination.getName(), "bar", 0); + } + + @Test + void provisioningOnProducerBindingWithPartitionsSetAtTheBinderProperties() { + PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties = new PulsarBinderConfigurationProperties(); + pulsarBinderConfigurationProperties.setPartitionCount(4); + PulsarAdministration pulsarAdministration = mock(PulsarAdministration.class); + PulsarTopicProvisioner pulsarTopicProvisioner = new PulsarTopicProvisioner(pulsarAdministration, + pulsarBinderConfigurationProperties); + ExtendedProducerProperties properties = new ExtendedProducerProperties<>( + new PulsarProducerProperties()); + ProducerDestination producerDestination = pulsarTopicProvisioner.provisionProducerDestination("foo", + properties); + verifyAndAssert(pulsarAdministration, producerDestination.getName(), "foo", 4); + } + + @Test + void provisioningOnProducerBindingWithPartitionsSetAtTheBindingProperties() { + PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties = new PulsarBinderConfigurationProperties(); + PulsarAdministration pulsarAdministration = mock(PulsarAdministration.class); + PulsarTopicProvisioner pulsarTopicProvisioner = new PulsarTopicProvisioner(pulsarAdministration, + pulsarBinderConfigurationProperties); + ExtendedProducerProperties properties = new ExtendedProducerProperties<>( + new PulsarProducerProperties()); + properties.getExtension().setPartitionCount(4); + ProducerDestination producerDestination = pulsarTopicProvisioner.provisionProducerDestination("foo", + properties); + verifyAndAssert(pulsarAdministration, producerDestination.getName(), "foo", 4); + } + + @Test + void provisionThroughConsumerBindingWithPartitionsSetAtTheBinderProperties() { + PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties = new PulsarBinderConfigurationProperties(); + pulsarBinderConfigurationProperties.setPartitionCount(4); + PulsarAdministration pulsarAdministration = mock(PulsarAdministration.class); + PulsarTopicProvisioner pulsarTopicProvisioner = new PulsarTopicProvisioner(pulsarAdministration, + pulsarBinderConfigurationProperties); + ExtendedConsumerProperties properties = new ExtendedConsumerProperties<>( + new PulsarConsumerProperties()); + ConsumerDestination consumerDestination = pulsarTopicProvisioner.provisionConsumerDestination("bar", "", + properties); + verifyAndAssert(pulsarAdministration, consumerDestination.getName(), "bar", 4); + } + + @Test + void provisionThroughConsumerBindingWithPartitionsSetAtTheBindingProperties() { + PulsarBinderConfigurationProperties pulsarBinderConfigurationProperties = new PulsarBinderConfigurationProperties(); + PulsarAdministration pulsarAdministration = mock(PulsarAdministration.class); + PulsarTopicProvisioner pulsarTopicProvisioner = new PulsarTopicProvisioner(pulsarAdministration, + pulsarBinderConfigurationProperties); + PulsarConsumerProperties pulsarConsumerProperties = new PulsarConsumerProperties(); + pulsarConsumerProperties.setPartitionCount(4); + ExtendedConsumerProperties properties = new ExtendedConsumerProperties<>( + pulsarConsumerProperties); + ConsumerDestination consumerDestination = pulsarTopicProvisioner.provisionConsumerDestination("bar", "", + properties); + verifyAndAssert(pulsarAdministration, consumerDestination.getName(), "bar", 4); + } + +}