GH-2014 Add support for NewDestinationBindingCallback back to StreamBridge

Resolves #2014
This commit is contained in:
Oleg Zhurakousky
2020-09-22 13:58:52 +02:00
parent 6fed9dc91e
commit 29fb69a2cd
3 changed files with 64 additions and 3 deletions

View File

@@ -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

View File

@@ -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<String, SubscribableChannel>() {
@Override
protected boolean removeEldestEntry(Map.Entry<String, SubscribableChannel> 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);
}

View File

@@ -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<String, String> echo() {
return v -> v;
}
@Bean
public Function<String, String> 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 {