From 5013103fcff9e9901621871cef50aa7319e2556a Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 13 Sep 2019 17:45:25 +0200 Subject: [PATCH] GH-1794 Additional fixes and refactoring --- .../BinderFactoryAutoConfiguration.java | 116 -------------- .../BindableFunctionProxyFactory.java | 100 ++++++++++++ .../function/FunctionConfiguration.java | 149 +++++++++++++++++- .../ImplicitFunctionBindingTests.java | 4 +- .../SourceToFunctionsSupportTests.java | 4 +- 5 files changed, 246 insertions(+), 127 deletions(-) create mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/BindableFunctionProxyFactory.java diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java index ea3e196ba..c9826d542 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BinderFactoryAutoConfiguration.java @@ -26,49 +26,29 @@ import java.util.LinkedList; import java.util.List; import java.util.Map; import java.util.Properties; -import java.util.stream.Stream; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; -import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanCreationException; -import org.springframework.beans.factory.BeanFactory; -import org.springframework.beans.factory.BeanFactoryAware; -import org.springframework.beans.factory.InitializingBean; import org.springframework.beans.factory.annotation.Qualifier; import org.springframework.beans.factory.annotation.Value; import org.springframework.beans.factory.config.BeanDefinition; import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; -import org.springframework.beans.factory.support.BeanDefinitionRegistry; -import org.springframework.beans.factory.support.RootBeanDefinition; import org.springframework.boot.context.properties.EnableConfigurationProperties; -import org.springframework.cloud.function.context.FunctionCatalog; -import org.springframework.cloud.function.context.catalog.FunctionInspector; -import org.springframework.cloud.function.context.config.RoutingFunction; -import org.springframework.cloud.stream.annotation.BindingProvider; -import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BinderType; import org.springframework.cloud.stream.binder.BinderTypeRegistry; import org.springframework.cloud.stream.binder.DefaultBinderTypeRegistry; -import org.springframework.cloud.stream.binding.BindableProxyFactory; import org.springframework.cloud.stream.binding.CompositeMessageChannelConfigurer; import org.springframework.cloud.stream.binding.MessageChannelConfigurer; import org.springframework.cloud.stream.binding.MessageConverterConfigurer; import org.springframework.cloud.stream.binding.MessageSourceBindingTargetFactory; import org.springframework.cloud.stream.binding.SubscribableChannelBindingTargetFactory; -import org.springframework.cloud.stream.messaging.Processor; -import org.springframework.cloud.stream.messaging.Sink; -import org.springframework.cloud.stream.messaging.Source; import org.springframework.context.ConfigurableApplicationContext; -import org.springframework.context.EnvironmentAware; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; import org.springframework.context.annotation.Role; -import org.springframework.core.annotation.AnnotationUtils; -import org.springframework.core.env.Environment; -import org.springframework.core.env.StandardEnvironment; import org.springframework.core.io.Resource; import org.springframework.core.io.UrlResource; import org.springframework.core.io.support.PropertiesLoaderUtils; @@ -85,7 +65,6 @@ import org.springframework.messaging.handler.annotation.support.HeadersMethodArg import org.springframework.messaging.handler.annotation.support.MessageHandlerMethodFactory; import org.springframework.messaging.handler.invocation.HandlerMethodArgumentResolver; import org.springframework.util.ClassUtils; -import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; import org.springframework.validation.Validator; @@ -248,99 +227,4 @@ public class BinderFactoryAutoConfiguration { return new CompositeMessageChannelConfigurer(configurerList); } - @Bean - public InitializingBean functionToChannelBindingInitializer(@Nullable FunctionCatalog functionCatalog, - @Nullable FunctionInspector functionInspector, BinderTypeRegistry bfac) { - return new ImplicitFunctionToChannelBindingInitializer(functionCatalog, functionInspector, bfac); - } - - - private static class ImplicitFunctionToChannelBindingInitializer implements InitializingBean, BeanFactoryAware, EnvironmentAware { - - private ConfigurableListableBeanFactory beanFactory; - - private Environment environment; - - private final FunctionCatalog functionCatalog; - - private final FunctionInspector functionInspector; - - private final BinderTypeRegistry bfac; - - ImplicitFunctionToChannelBindingInitializer(FunctionCatalog functionCatalog, - FunctionInspector functionInspector, BinderTypeRegistry bfac) { - this.functionCatalog = functionCatalog; - this.functionInspector = functionInspector; - this.bfac = bfac; - } - @Override - public void afterPropertiesSet() { - Class[] configurationClasses = bfac.getAll().values().iterator().next().getConfigurationClasses(); - boolean bindingProvider = Stream.of(configurationClasses) - .filter(clazz -> AnnotationUtils.findAnnotation(clazz, BindingProvider.class) != null) - .findFirst().isPresent(); - if (functionCatalog != null && ObjectUtils.isEmpty(beanFactory.getBeanNamesForAnnotation(EnableBinding.class))) { - BeanDefinitionRegistry registry = (BeanDefinitionRegistry) beanFactory; - String name = determineFunctionName(functionCatalog, environment); - if (StringUtils.hasText(name)) { - Object definedFunction = functionCatalog.lookup(name); - Class inputType = functionInspector.getInputType(definedFunction); - Class outputType = functionInspector.getOutputType(definedFunction); - - if (!bindingProvider) { - if (Void.class.isAssignableFrom(outputType)) { - bind(Sink.class, registry); - } - else if (Void.class.isAssignableFrom(inputType)) { - bind(Source.class, registry); - } - else { - bind(Processor.class, registry); - } - } - } - } - } - - @Override - public void setBeanFactory(BeanFactory beanFactory) throws BeansException { - this.beanFactory = (ConfigurableListableBeanFactory) beanFactory; - } - - @Override - public void setEnvironment(Environment environment) { - this.environment = environment; - } - - private String determineFunctionName(FunctionCatalog catalog, Environment environment) { - String name = environment.getProperty("spring.cloud.stream.function.definition"); - if (!StringUtils.hasText(name)) { - name = environment.getProperty("spring.cloud.function.definition"); - } - if (!StringUtils.hasText(name) && Boolean.parseBoolean( - environment.getProperty("spring.cloud.stream.function.routing.enabled", "false"))) { - name = RoutingFunction.FUNCTION_NAME; - } - if (!StringUtils.hasText(name) && catalog.size() >= 1 && catalog.size() <= 2) { - name = ((FunctionInspector) catalog).getName(catalog.lookup("")); - } - if (StringUtils.hasText(name)) { - ((StandardEnvironment) environment).getSystemProperties() - .putIfAbsent("spring.cloud.stream.function.definition", name); - } - return name; - } - - private void bind(Class type, BeanDefinitionRegistry registry) { - if (!registry.containsBeanDefinition(type.getName())) { - RootBeanDefinition rootBeanDefinition = new RootBeanDefinition( - BindableProxyFactory.class); - rootBeanDefinition.getConstructorArgumentValues() - .addGenericArgumentValue(type); - registry.registerBeanDefinition(type.getName(), rootBeanDefinition); - } - } - } - - } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/BindableFunctionProxyFactory.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/BindableFunctionProxyFactory.java new file mode 100644 index 000000000..000d68e26 --- /dev/null +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/BindableFunctionProxyFactory.java @@ -0,0 +1,100 @@ +/* + * Copyright 2019-2019 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.function; + +import org.springframework.beans.factory.FactoryBean; +import org.springframework.cloud.stream.binding.BindableProxyFactory; +import org.springframework.cloud.stream.binding.BoundTargetHolder; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.SubscribableChannel; +import org.springframework.util.Assert; + +/** + * {@link FactoryBean} for creating inputs/outputs destinations to be bound to + * function arguments. It is an extension to {@link BindableProxyFactory} which + * operates on Bindable interfaces (e.g., Source, Processor, Sink) which internally + * define inputs and output channels. Unlike BindableProxyFactory, this class simply + * operates based on the count of provided inputs and outputs and the names of inputs and outputs + * are based on convention. + * + * @author Oleg Zhurakousky + * + * + * + * @since 3.0 + */ +public class BindableFunctionProxyFactory extends BindableProxyFactory { + + private final int inputCount; + + private final int outputCount; + + public BindableFunctionProxyFactory(int inputCount, int outputCount) { + super(null); + this.inputCount = inputCount; + this.outputCount = outputCount; + } + + + @Override + public void afterPropertiesSet() { + Assert.notEmpty(BindableFunctionProxyFactory.this.bindingTargetFactories, + "'bindingTargetFactories' cannot be empty"); + + if (this.inputCount > 0) { + if (this.inputCount == 1) { + this.createInput("input"); + } + else { + throw new UnsupportedOperationException("Multiple inputs are not currently supported"); + } + } + + if (this.outputCount > 0) { + if (this.outputCount == 1) { + this.createOutput("output"); + } + else { + throw new UnsupportedOperationException("Multiple outputs are not currently supported"); + } + } + } + + private void createInput(String name) { + BindableFunctionProxyFactory.this.inputHolders.put(name, + new BoundTargetHolder(getBindingTargetFactory(SubscribableChannel.class) + .createInput(name), true)); + } + + private void createOutput(String name) { + BindableFunctionProxyFactory.this.outputHolders.put(name, + new BoundTargetHolder(getBindingTargetFactory(MessageChannel.class) + .createOutput(name), true)); + } + + + @Override + public Class getObjectType() { + return this.type; + } + + @Override + public boolean isSingleton() { + return true; + } + +} 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 230b7dc76..be3a6d8df 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 @@ -23,6 +23,7 @@ import java.util.concurrent.atomic.AtomicReference; import java.util.function.Consumer; import java.util.function.Function; import java.util.function.Supplier; +import java.util.stream.Stream; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; @@ -33,6 +34,7 @@ import reactor.core.publisher.MonoSink; import org.springframework.beans.BeansException; import org.springframework.beans.factory.InitializingBean; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; import org.springframework.beans.factory.support.RootBeanDefinition; import org.springframework.boot.autoconfigure.AutoConfigureBefore; import org.springframework.boot.context.properties.EnableConfigurationProperties; @@ -42,7 +44,10 @@ import org.springframework.cloud.function.context.catalog.BeanFactoryAwareFuncti import org.springframework.cloud.function.context.catalog.FunctionInspector; import org.springframework.cloud.function.context.catalog.FunctionTypeUtils; import org.springframework.cloud.function.context.config.FunctionContextUtils; +import org.springframework.cloud.function.context.config.RoutingFunction; +import org.springframework.cloud.stream.annotation.BindingProvider; import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.binder.BinderTypeRegistry; import org.springframework.cloud.stream.binder.BindingCreatedEvent; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binding.BindableProxyFactory; @@ -55,11 +60,14 @@ import org.springframework.cloud.stream.messaging.Sink; import org.springframework.cloud.stream.messaging.Source; import org.springframework.context.ApplicationContext; import org.springframework.context.ApplicationContextAware; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.EnvironmentAware; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.context.annotation.Import; import org.springframework.context.support.GenericApplicationContext; import org.springframework.core.annotation.AnnotationUtils; +import org.springframework.core.env.Environment; import org.springframework.core.type.MethodMetadata; import org.springframework.integration.channel.MessageChannelReactiveUtils; import org.springframework.integration.dsl.IntegrationFlow; @@ -77,6 +85,7 @@ import org.springframework.util.ClassUtils; import org.springframework.util.MimeTypeUtils; import org.springframework.util.ObjectUtils; import org.springframework.util.ReflectionUtils; +import org.springframework.util.StringUtils; /** * @author Oleg Zhurakousky @@ -88,18 +97,45 @@ import org.springframework.util.ReflectionUtils; @EnableConfigurationProperties(StreamFunctionProperties.class) @Import(BinderFactoryAutoConfiguration.class) @AutoConfigureBefore(BindingServiceConfiguration.class) -class FunctionConfiguration { +public class FunctionConfiguration { + /* + * Creates an effective representation of Bindable interfaces by maintaining the count of inputs and + * outputs based on the provided function, thus preserving the contract and the infrastructure code used + * by current EnableBinding/StreamListener combination. + * It is then used buy `functionInitializer` or 'supplierInitializer` where functions are actually bound to channels. + * + * Also, see the BindableFunctionProxyFactory + */ @Bean - public InitializingBean functionChannelBindingInitializer(FunctionCatalog functionCatalog, FunctionInspector functionInspector, - StreamFunctionProperties functionProperties, @Nullable BindableProxyFactory[] bindableProxyFactory, BindingServiceProperties serviceProperties) { - return new FunctionChannelBindingInitializer(functionCatalog, functionInspector, functionProperties, - ObjectUtils.isEmpty(bindableProxyFactory) ? null : bindableProxyFactory[0], serviceProperties); + public InitializingBean functionBindingHolder(Environment environment, FunctionCatalog functionCatalog, + StreamFunctionProperties streamFunctionProperties, BinderTypeRegistry binderTypeRegistry) { + return new FunctionBindingHolder(binderTypeRegistry, functionCatalog, streamFunctionProperties); } @Bean - public IntegrationFlow standAloneSupplierFlow(FunctionCatalog functionCatalog, FunctionInspector functionInspector, - StreamFunctionProperties functionProperties, GenericApplicationContext context) { + public InitializingBean functionInitializer(FunctionCatalog functionCatalog, FunctionInspector functionInspector, + StreamFunctionProperties functionProperties, @Nullable BindableProxyFactory[] bpfs, BindingServiceProperties serviceProperties, + ConfigurableApplicationContext applicationContext, FunctionBindingHolder bindingHolder) { + + if (bpfs == null || bpfs.length > 1) { + return null; // basically we're not dealing with multiple EnableBinding which is how multiple BindableProxyFactory are created + } + BindableProxyFactory bindableProxyFactory = bpfs[0]; + + return bindingHolder.getInputCount() > 0 // basically not a Supplier + || !ObjectUtils.isEmpty(applicationContext.getBeanNamesForAnnotation(EnableBinding.class)) // implies existing binding to which we are going to 'compose to' + ? new FunctionChannelBindingInitializer(functionCatalog, functionInspector, functionProperties, bindableProxyFactory, serviceProperties) + : null; + } + + @Bean + public IntegrationFlow supplierInitializer(FunctionCatalog functionCatalog, FunctionInspector functionInspector, + StreamFunctionProperties functionProperties, GenericApplicationContext context, FunctionBindingHolder bindingHolder) { + if (bindingHolder.getInputCount() > 0) { + return null; + } + FunctionInvocationWrapper functionWrapper = functionCatalog.lookup(functionProperties.getDefinition()); IntegrationFlow integrationFlow = null; if (ObjectUtils.isEmpty(context.getBeanNamesForAnnotation(EnableBinding.class)) && functionWrapper != null && functionWrapper.isSupplier()) { @@ -352,4 +388,103 @@ class FunctionConfiguration { return (Message) result; } } + + /* + * This class will effectively create a different representation of Bindable interfaces (e.g., Source, Processor...). + * It's main goal is to determine the count of inputs and outputs based on the provided function. + */ + private static class FunctionBindingHolder implements InitializingBean, ApplicationContextAware, EnvironmentAware { + + private final BinderTypeRegistry binderTypeRegistry; + + private final FunctionCatalog functionCatalog; + + private final StreamFunctionProperties streamFunctionProperties; + + private ConfigurableApplicationContext applicationContext; + + private Environment environment; + + private int inputCount; + + private int outputCount; + + FunctionBindingHolder(BinderTypeRegistry binderTypeRegistry, FunctionCatalog functionCatalog, StreamFunctionProperties streamFunctionProperties) { + this.binderTypeRegistry = binderTypeRegistry; + this.functionCatalog = functionCatalog; + this.streamFunctionProperties = streamFunctionProperties; + } + + @Override + public void afterPropertiesSet() throws Exception { + Class[] configurationClasses = binderTypeRegistry.getAll().values().iterator().next() + .getConfigurationClasses(); + boolean bindingProvider = Stream.of(configurationClasses) + .filter(clazz -> AnnotationUtils.findAnnotation(clazz, BindingProvider.class) != null) + .findFirst().isPresent(); + if (!bindingProvider + && ObjectUtils.isEmpty(applicationContext.getBeanNamesForAnnotation(EnableBinding.class))) { + this.determineFunctionName(functionCatalog, environment); + BeanDefinitionRegistry registry = (BeanDefinitionRegistry) applicationContext.getBeanFactory(); + RootBeanDefinition rootBeanDefinition = new RootBeanDefinition(BindableFunctionProxyFactory.class); + FunctionInvocationWrapper function = functionCatalog + .lookup(streamFunctionProperties.getDefinition()); + if (function != null) { + if (function.isSupplier()) { + this.inputCount = 0; + this.outputCount = 1; + } + else if (function.isConsumer()) { + this.inputCount = 1; + this.outputCount = 0; + } + else { + this.inputCount = 1; + this.outputCount = 1; + } + rootBeanDefinition.getConstructorArgumentValues().addGenericArgumentValue(this.inputCount); + rootBeanDefinition.getConstructorArgumentValues().addGenericArgumentValue(this.outputCount); + registry.registerBeanDefinition(streamFunctionProperties.getDefinition() + "_binding", + rootBeanDefinition); + } + + } + } + + int getInputCount() { + return this.inputCount; + } + + int getOutputCount() { + return this.outputCount; + } + + private void determineFunctionName(FunctionCatalog catalog, Environment environment) { + String definition = streamFunctionProperties.getDefinition(); + if (!StringUtils.hasText(definition)) { + definition = environment.getProperty("spring.cloud.function.definition"); + } + + if (StringUtils.hasText(definition)) { + streamFunctionProperties.setDefinition(definition); + } + else if (Boolean.parseBoolean(environment.getProperty("spring.cloud.stream.function.routing.enabled", "false"))) { + streamFunctionProperties.setDefinition(RoutingFunction.FUNCTION_NAME); + } + else { + streamFunctionProperties.setDefinition(((FunctionInspector) functionCatalog).getName(functionCatalog.lookup(""))); + } + } + + @Override + public void setApplicationContext(ApplicationContext applicationContext) throws BeansException { + this.applicationContext = (ConfigurableApplicationContext) applicationContext; + } + + @Override + public void setEnvironment(Environment environment) { + this.environment = environment; + } + + } } 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 481226a8c..ff0080156 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 @@ -147,7 +147,7 @@ public class ImplicitFunctionBindingTests { public void testConsumer() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration - .getCompleteConfiguration(SingleFunctionConfiguration.class)) + .getCompleteConfiguration(SingleConsumerConfiguration.class)) .web(WebApplicationType.NONE) .run("--spring.cloud.stream.function.definition=consumer", "--spring.jmx.enabled=false")) { @@ -193,7 +193,7 @@ public class ImplicitFunctionBindingTests { .web(WebApplicationType.NONE) .run("--spring.jmx.enabled=false")) { - assertThat(context.getBean("standAloneSupplierFlow")).isEqualTo(null); + assertThat(context.getBean("supplierInitializer")).isEqualTo(null); } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java index 1fffb2d71..811069243 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/SourceToFunctionsSupportTests.java @@ -233,8 +233,8 @@ public class SourceToFunctionsSupportTests { assertThat(new String(target.receive(2000).getPayload())).isEqualTo("5"); assertThat(new String(target.receive(2000).getPayload())).isEqualTo("6"); - assertThat(context.getBean("standAloneSupplierFlow")).isNotEqualTo(null); - assertThat(context.getBean("functionChannelBindingInitializer")).isNotEqualTo(null); + assertThat(context.getBean("supplierInitializer")).isNotEqualTo(null); + assertThat(context.getBean("functionInitializer")).isEqualTo(null); } }