From ff281c56c412dc41bd7b81e930558ee907080ac5 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 18 Jan 2022 10:32:12 +0100 Subject: [PATCH] GH-2266 Fix channel interceptor application for dynamic destinations Resolves #2266 Resolves #2269 --- .../cloud/stream/binding/MessageConverterConfigurer.java | 4 +++- .../cloud/stream/config/BindingServiceProperties.java | 5 +++++ .../cloud/stream/function/FunctionConfiguration.java | 2 -- .../springframework/cloud/stream/function/StreamBridge.java | 1 - .../cloud/stream/function/StreamBridgeTests.java | 4 +++- 5 files changed, 11 insertions(+), 5 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java index d9355de8f..15c3f481f 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java @@ -147,7 +147,9 @@ public class MessageConverterConfigurer ProducerProperties producerProperties = bindingProperties.getProducer(); boolean partitioned = !inbound && producerProperties != null && producerProperties.isPartitioned(); boolean functional = streamFunctionProperties != null - && (StringUtils.hasText(streamFunctionProperties.getDefinition()) || StringUtils.hasText(bindingServiceProperties.getSource())); + && (StringUtils.hasText(streamFunctionProperties.getDefinition()) + || StringUtils.hasText(bindingServiceProperties.getInputBindings()) + || StringUtils.hasText(bindingServiceProperties.getOutputBindings())); if (partitioned) { if (inbound || !functional) { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java index 015d5dea1..a4698fabc 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java @@ -44,6 +44,7 @@ import org.springframework.core.convert.converter.Converter; import org.springframework.core.convert.support.GenericConversionService; import org.springframework.integration.support.utils.IntegrationUtils; import org.springframework.util.Assert; +import org.springframework.util.StringUtils; /** * @author Dave Syer @@ -330,6 +331,7 @@ public class BindingServiceProperties @Deprecated public void setSource(String source) { this.source = source; + this.outputBindings = source; } public void updateProducerProperties(String bindingName, @@ -356,10 +358,13 @@ public class BindingServiceProperties } public String getOutputBindings() { + return outputBindings; } public void setOutputBindings(String outputBindings) { + Assert.state(!StringUtils.hasText(this.source), "Setting 'source' and 'output-binding' is not allowed " + + "because 'source' is deprecated in favor of 'output-binding'."); this.outputBindings = outputBindings; } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java index f9bda9147..71a55ed32 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionConfiguration.java @@ -129,8 +129,6 @@ public class FunctionConfiguration { private final static String SOURCE_PROPERY = "spring.cloud.stream.source"; -// private final static String OUT_BINDINGS = "spring.cloud.stream.output-bindings"; - @Bean public StreamBridge streamBridgeUtils(FunctionCatalog functionCatalog, FunctionRegistry functionRegistry, BindingServiceProperties bindingServiceProperties, ConfigurableApplicationContext applicationContext, diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java index 6552e7dde..087ceb6ca 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamBridge.java @@ -256,7 +256,6 @@ public final class StreamBridge implements SmartInitializingSingleton { SubscribableChannel messageChannel = this.channelCache.get(destinationName); if (messageChannel == null && this.applicationContext.containsBean(destinationName)) { messageChannel = this.applicationContext.getBean(destinationName, SubscribableChannel.class); - this.addInterceptors((AbstractMessageChannel) messageChannel, destinationName); } if (messageChannel == null) { messageChannel = new DirectWithAttributesChannel(); diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java index 8b5c51b68..74126e089 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/StreamBridgeTests.java @@ -193,11 +193,13 @@ public class StreamBridgeTests { .web(WebApplicationType.NONE).run( "--spring.jmx.enabled=false", "--spring.cloud.stream.dynamic-destination-cache-size=1", - "--spring.cloud.stream.source=outputA;outputB", + "--spring.cloud.stream.output-bindings=outputA;outputB", "--spring.cloud.stream.bindings.outputA-out-0.destination=outputA", "--spring.cloud.stream.bindings.outputB-out-0.destination=outputB")) { StreamBridge bridge = context.getBean(StreamBridge.class); + bridge.send("outputA-out-0", "hello foo"); + bridge.send("outputA-out-0", "hello foo"); bridge.send("outputA-out-0", "hello foo"); bridge.send("outputA-out-0", "hello foo");