From 29fb69a2cda05ec387a52a87002072bf17aa3b74 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 22 Sep 2020 13:58:52 +0200 Subject: [PATCH] GH-2014 Add support for NewDestinationBindingCallback back to StreamBridge Resolves #2014 --- .../function/FunctionConfiguration.java | 6 ++- .../cloud/stream/function/StreamBridge.java | 17 ++++++- .../stream/function/StreamBridgeTests.java | 44 +++++++++++++++++++ 3 files changed, 64 insertions(+), 3 deletions(-) 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 739d62531..fe70fe6c9 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 @@ -64,6 +64,7 @@ import org.springframework.cloud.stream.binder.BindingCreatedEvent; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.ProducerProperties; import org.springframework.cloud.stream.binding.BindableProxyFactory; +import org.springframework.cloud.stream.binding.BinderAwareChannelResolver.NewDestinationBindingCallback; import org.springframework.cloud.stream.config.BinderFactoryAutoConfiguration; import org.springframework.cloud.stream.config.BindingBeansRegistrar; import org.springframework.cloud.stream.config.BindingProperties; @@ -120,8 +121,9 @@ public class FunctionConfiguration { @Bean public StreamBridge streamBridgeUtils(FunctionCatalog functionCatalog, FunctionRegistry functionRegistry, - BindingServiceProperties bindingServiceProperties, ConfigurableApplicationContext applicationContext) { - return new StreamBridge(functionCatalog, functionRegistry, bindingServiceProperties, applicationContext); + BindingServiceProperties bindingServiceProperties, ConfigurableApplicationContext applicationContext, + @Nullable NewDestinationBindingCallback callback) { + return new StreamBridge(functionCatalog, functionRegistry, bindingServiceProperties, applicationContext, callback); } @Bean 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 d96989415..e77adb2a8 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 @@ -32,11 +32,13 @@ import org.springframework.cloud.function.context.FunctionRegistry; import org.springframework.cloud.function.context.FunctionType; import org.springframework.cloud.function.context.catalog.SimpleFunctionRegistry.FunctionInvocationWrapper; import org.springframework.cloud.stream.binder.ProducerProperties; +import org.springframework.cloud.stream.binding.BinderAwareChannelResolver.NewDestinationBindingCallback; import org.springframework.cloud.stream.binding.BindingService; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.messaging.DirectWithAttributesChannel; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.integration.support.MessageBuilder; +import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.SubscribableChannel; import org.springframework.util.MimeType; @@ -58,6 +60,7 @@ import org.springframework.util.MimeTypeUtils; * @since 3.0.3 * */ +@SuppressWarnings("deprecation") public final class StreamBridge implements SmartInitializingSingleton { private static String STREAM_BRIDGE_FUNC_NAME = "streamBridge"; @@ -70,6 +73,8 @@ public final class StreamBridge implements SmartInitializingSingleton { private final FunctionRegistry functionRegistry; + private final NewDestinationBindingCallback destinationBindingCallback; + private BindingServiceProperties bindingServiceProperties; private ConfigurableApplicationContext applicationContext; @@ -89,11 +94,13 @@ public final class StreamBridge implements SmartInitializingSingleton { */ @SuppressWarnings("serial") StreamBridge(FunctionCatalog functionCatalog, FunctionRegistry functionRegistry, - BindingServiceProperties bindingServiceProperties, ConfigurableApplicationContext applicationContext) { + BindingServiceProperties bindingServiceProperties, ConfigurableApplicationContext applicationContext, + @Nullable NewDestinationBindingCallback destinationBindingCallback) { this.functionCatalog = functionCatalog; this.functionRegistry = functionRegistry; this.applicationContext = applicationContext; this.bindingServiceProperties = bindingServiceProperties; + this.destinationBindingCallback = destinationBindingCallback; this.channelCache = new LinkedHashMap() { @Override protected boolean removeEldestEntry(Map.Entry eldest) { @@ -165,6 +172,7 @@ public final class StreamBridge implements SmartInitializingSingleton { this.initialized = true; } + @SuppressWarnings({ "unchecked", "deprecation" }) SubscribableChannel resolveDestination(String destinationName, ProducerProperties producerProperties) { SubscribableChannel messageChannel = this.channelCache.get(destinationName); if (messageChannel == null && this.applicationContext.containsBean(destinationName)) { @@ -172,6 +180,13 @@ public final class StreamBridge implements SmartInitializingSingleton { } if (messageChannel == null) { messageChannel = new DirectWithAttributesChannel(); + if (this.destinationBindingCallback != null) { + Object extendedProducerProperties = this.bindingService + .getExtendedProducerProperties(messageChannel, destinationName); + this.destinationBindingCallback.configure(destinationName, messageChannel, + producerProperties, extendedProducerProperties); + } + this.bindingService.bindProducer(messageChannel, destinationName, false); this.channelCache.put(destinationName, messageChannel); } 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 e1b2ddba4..b72d86140 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 @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.function; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.function.Function; import java.util.function.Supplier; @@ -28,6 +29,7 @@ import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.cloud.stream.binder.test.OutputDestination; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +import org.springframework.cloud.stream.binding.BinderAwareChannelResolver.NewDestinationBindingCallback; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; import org.springframework.integration.dsl.IntegrationFlow; @@ -43,6 +45,7 @@ import static org.junit.Assert.fail; * @author Oleg Zhurakousky * */ +@SuppressWarnings("deprecation") public class StreamBridgeTests { @Before @@ -206,6 +209,19 @@ public class StreamBridgeTests { } } + @Test + public void testNewBindingCallback() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration + .getCompleteConfiguration(BindingCallbackConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.cloud.stream.source=uppercase", + "--spring.jmx.enabled=false")) { + + StreamBridge bridge = context.getBean(StreamBridge.class); + bridge.send("uppercase-in-0", "hello"); + assertThat(context.getBean("callbackVerifier", AtomicBoolean.class)).isTrue(); + } + } + @EnableAutoConfiguration public static class EmptyConfiguration { @@ -234,6 +250,34 @@ public class StreamBridgeTests { } } + + @EnableAutoConfiguration + public static class BindingCallbackConfiguration { + + @Bean + public Function echo() { + return v -> v; + } + + @Bean + public Function uppercase() { + return v -> v.toUpperCase(); + } + + @Bean + public AtomicBoolean callbackVerifier() { + return new AtomicBoolean(); + } + + @Bean + public NewDestinationBindingCallback callback(AtomicBoolean callbackVerifier) { + + return (name, channel, props, extended) -> { + callbackVerifier.set(true); + }; + } + } + @EnableAutoConfiguration public static class IntegrationFlowConfiguration {