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 aa8797ea1..c5d137622 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 @@ -127,12 +127,15 @@ public class FunctionConfiguration { BindingServiceProperties serviceProperties, ConfigurableApplicationContext applicationContext, FunctionBindingRegistrar bindingHolder, BinderAwareChannelResolver dynamicDestinationResolver) { - boolean shouldCreateInitializer = bindableProxyFactories != null - && (applicationContext.containsBean("output") // need this to compose to existing legacy message source - || ObjectUtils.isEmpty(applicationContext.getBeanNamesForAnnotation(EnableBinding.class))); +// boolean shouldCreateInitializer = bindableProxyFactories != null +// && (applicationContext.containsBean("output") // need this to compose to existing legacy message source +// || ObjectUtils.isEmpty(applicationContext.getBeanNamesForAnnotation(EnableBinding.class))); + + boolean shouldCreateInitializer = applicationContext.containsBean("output") + || ObjectUtils.isEmpty(applicationContext.getBeanNamesForAnnotation(EnableBinding.class)); return shouldCreateInitializer - ? new FunctionToDestinationBinder(functionCatalog, functionProperties, bindableProxyFactories, + ? new FunctionToDestinationBinder(functionCatalog, functionProperties, serviceProperties, dynamicDestinationResolver) : null; @@ -282,7 +285,7 @@ public class FunctionConfiguration { private GenericApplicationContext applicationContext; - private final BindableProxyFactory[] bindableProxyFactories; + private BindableProxyFactory[] bindableProxyFactories; private final FunctionCatalog functionCatalog; @@ -293,9 +296,7 @@ public class FunctionConfiguration { private final BinderAwareChannelResolver dynamicDestinationResolver; FunctionToDestinationBinder(FunctionCatalog functionCatalog, StreamFunctionProperties functionProperties, - BindableProxyFactory[] bindableProxyFactories, BindingServiceProperties serviceProperties, - BinderAwareChannelResolver dynamicDestinationResolver) { - this.bindableProxyFactories = bindableProxyFactories; + BindingServiceProperties serviceProperties, BinderAwareChannelResolver dynamicDestinationResolver) { this.functionCatalog = functionCatalog; this.functionProperties = functionProperties; this.serviceProperties = serviceProperties; @@ -309,6 +310,8 @@ public class FunctionConfiguration { @Override public void afterPropertiesSet() throws Exception { + Map beansOfType = applicationContext.getBeansOfType(BindableProxyFactory.class); + this.bindableProxyFactories = beansOfType.values().toArray(new BindableProxyFactory[0]); for (BindableProxyFactory bindableProxyFactory : this.bindableProxyFactories) { String functionDefinition = bindableProxyFactory instanceof BindableFunctionProxyFactory ? ((BindableFunctionProxyFactory) bindableProxyFactory).getFunctionDefinition() diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/FunctionBindingTestUtils.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/FunctionBindingTestUtils.java new file mode 100644 index 000000000..da0fee63d --- /dev/null +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/test/FunctionBindingTestUtils.java @@ -0,0 +1,67 @@ +package org.springframework.cloud.stream.binder.test; + +import java.util.Map; +import java.util.function.Consumer; +import java.util.function.Function; +import java.util.function.Supplier; + +import org.springframework.beans.factory.InitializingBean; +import org.springframework.cloud.stream.binder.Binding; +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.config.BindingProperties; +import org.springframework.cloud.stream.config.BindingServiceProperties; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.messaging.MessageChannel; + +public class FunctionBindingTestUtils { + + public static void bind(ConfigurableApplicationContext applicationContext, Object function) { + try { + String functionName = function instanceof Function ? "function" : (function instanceof Consumer ? "consumer" : "supplier"); + System.setProperty("spring.cloud.function.definition", functionName); + applicationContext.getBeanFactory().registerSingleton(functionName, function); + + + InitializingBean functionBindingRegistrar = applicationContext.getBean("functionBindingRegistrar", InitializingBean.class); + functionBindingRegistrar.afterPropertiesSet(); + + BindableProxyFactory bindingProxy = applicationContext.getBean("&" + functionName + "_binding", BindableProxyFactory.class); + bindingProxy.afterPropertiesSet(); + + InitializingBean functionBinder = applicationContext.getBean("functionInitializer", InitializingBean.class); + functionBinder.afterPropertiesSet(); + + BindingServiceProperties bindingProperties = applicationContext.getBean(BindingServiceProperties.class); + String inputBindingName = functionName + "-in-0"; + String outputBindingName = functionName + "-out-0"; + Map bindings = bindingProperties.getBindings(); + BindingProperties inputProperties = bindings.get(inputBindingName); + BindingProperties outputProperties = bindings.get(outputBindingName); + ConsumerProperties consumerProperties = inputProperties.getConsumer(); + ProducerProperties producerProperties = outputProperties.getProducer(); + + + TestChannelBinder binder = applicationContext.getBean(TestChannelBinder.class); + if (function instanceof Supplier || function instanceof Function) { + Binding bindProducer = binder.bindProducer(outputProperties.getDestination(), + applicationContext.getBean(outputBindingName, MessageChannel.class), + producerProperties == null ? new ProducerProperties() : producerProperties); + bindProducer.start(); + } + if (function instanceof Consumer || function instanceof Function) { + Binding bindConsumer = binder.bindConsumer(inputProperties.getDestination(), null, + applicationContext.getBean(inputBindingName, MessageChannel.class), + consumerProperties == null ? new ConsumerProperties() : consumerProperties); + bindConsumer.start(); + } + } + catch (Exception e) { + throw new IllegalStateException("Failed to bind function", e); + } + finally { + System.clearProperty("spring.cloud.function.definition"); + } + } +} diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index 5a4499098..36604af71 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -30,15 +30,22 @@ import org.junit.After; import org.junit.Test; import reactor.core.publisher.Flux; +import org.springframework.beans.factory.InitializingBean; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; import org.springframework.cloud.function.context.config.ContextFunctionCatalogAutoConfiguration; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.annotation.StreamListener; +import org.springframework.cloud.stream.binder.Binding; +import org.springframework.cloud.stream.binder.ConsumerProperties; +import org.springframework.cloud.stream.binder.ProducerProperties; +import org.springframework.cloud.stream.binder.test.FunctionBindingTestUtils; import org.springframework.cloud.stream.binder.test.InputDestination; import org.springframework.cloud.stream.binder.test.OutputDestination; +import org.springframework.cloud.stream.binder.test.TestChannelBinder; import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +import org.springframework.cloud.stream.binding.BindableProxyFactory; import org.springframework.cloud.stream.messaging.Sink; import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.annotation.Bean; @@ -48,6 +55,7 @@ import org.springframework.integration.handler.LoggingHandler; import org.springframework.integration.scheduling.PollerMetadata; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageChannel; import org.springframework.messaging.support.GenericMessage; import org.springframework.scheduling.support.PeriodicTrigger; @@ -67,6 +75,30 @@ public class ImplicitFunctionBindingTests { System.clearProperty("spring.cloud.function.definition"); } + @Test + public void marcinsTest() { + ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(EmptyConfiguration.class)) + .web(WebApplicationType.NONE).run("--spring.jmx.enabled=false"); + + InputDestination input = context.getBean(InputDestination.class); + try { + input.send(new GenericMessage("hello".getBytes())); + fail(); // it should since there are no functions and no bindings + } + catch (Exception e) { + // good, we expected it + } + + Function function = v -> v.toUpperCase(); + FunctionBindingTestUtils.bind(context, function); + + input.send(new GenericMessage("hello".getBytes())); + + OutputDestination output = context.getBean(OutputDestination.class); + assertThat(new String(output.receive().getPayload())).isEqualTo("HELLO"); + } + @Test public void testEmptyConfiguration() {