From 716f90c1ab26fcefdf16c1da297e3caacf6b8399 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 24 Aug 2018 09:21:59 +0200 Subject: [PATCH] interim --- .../config/BindingServiceConfiguration.java | 2 +- .../function/FunctionConfiguration.java | 49 +++++++ .../IntegrationFlowFunctionSupport.java | 27 +++- .../main/resources/META-INF/spring.factories | 4 +- .../function/NewSourceAsSupplierTests.java | 132 ++++++++++++++++++ 5 files changed, 211 insertions(+), 3 deletions(-) create mode 100644 spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/NewSourceAsSupplierTests.java diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceConfiguration.java index 643d3763a..12124b684 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceConfiguration.java @@ -79,7 +79,7 @@ import org.springframework.util.Assert; */ @Configuration @EnableConfigurationProperties({ BindingServiceProperties.class, SpringIntegrationProperties.class, FunctionProperties.class }) -@Import({ DestinationPublishingMetricsAutoConfiguration.class, SpelExpressionConverterConfiguration.class, FunctionConfiguration.class }) +@Import({ DestinationPublishingMetricsAutoConfiguration.class, SpelExpressionConverterConfiguration.class }) @Role(BeanDefinition.ROLE_INFRASTRUCTURE) @ConditionalOnBean(value = BinderTypeRegistry.class, search = SearchStrategy.CURRENT) public class BindingServiceConfiguration { 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 8d5ef48e5..37b4a6551 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 @@ -16,12 +16,25 @@ package org.springframework.cloud.stream.function; +import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.config.ConfigurableListableBeanFactory; +import org.springframework.beans.factory.support.BeanDefinitionRegistry; +import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; import org.springframework.cloud.function.context.FunctionCatalog; +import org.springframework.cloud.function.context.FunctionType; import org.springframework.cloud.function.context.catalog.FunctionInspector; +import org.springframework.cloud.stream.binding.BindingBeanDefinitionRegistryUtils; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; +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.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.integration.dsl.IntegrationFlow; +import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.SubscribableChannel; +import org.springframework.util.ClassUtils; /** * @@ -33,6 +46,18 @@ import org.springframework.context.annotation.Configuration; @ConditionalOnProperty("spring.cloud.stream.function.name") public class FunctionConfiguration { + @Autowired(required=false) + private Source source; + + @Autowired(required=false) + private Processor processor; + + @Autowired(required=false) + private Sink sink; + + @Autowired + private ConfigurableListableBeanFactory registry; + @Bean public IntegrationFlowFunctionSupport functionSupport(FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory, @@ -47,4 +72,28 @@ public class FunctionConfiguration { return new FunctionCatalogWrapper(catalog); } + + @ConditionalOnProperty("spring.cloud.stream.function.name") + @ConditionalOnMissingBean + @Bean + public IntegrationFlow foo(IntegrationFlowFunctionSupport functionSupport) { + if (processor != null) { + return functionSupport.integrationFlowForFunction(processor.input(), processor.output()).get(); + } + else if (sink != null) { + return functionSupport.integrationFlowForFunction(sink.input(), null).get(); + } + else if (source != null) { + return functionSupport.integrationFlowFromNamedSupplier().channel(this.source.output()).get(); + } + + FunctionType ft = functionSupport.getCurrentFunctionType(); + BindingBeanDefinitionRegistryUtils.registerBindingTargetBeanDefinitions(Sink.class, + Sink.class.getName(), (BeanDefinitionRegistry) registry); + BindingBeanDefinitionRegistryUtils.registerBindingTargetsQualifiedBeanDefinitions( + ClassUtils.resolveClassName(this.getClass().getName(), null), Sink.class, + (BeanDefinitionRegistry) registry); + return functionSupport.integrationFlowForFunction(registry.getBean("input", SubscribableChannel.class), null).get(); + //throw new UnsupportedOperationException("Not yet supotrted"); + } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java index f9d6f4b6b..2bd645dc3 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/IntegrationFlowFunctionSupport.java @@ -26,6 +26,7 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.cloud.function.context.FunctionType; import org.springframework.cloud.function.context.catalog.FunctionInspector; import org.springframework.cloud.function.core.FluxSupplier; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; @@ -76,6 +77,10 @@ public class IntegrationFlowFunctionSupport { this.functionProperties = functionProperties; } + public FunctionType getCurrentFunctionType() { + return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getName())).getType(); + } + /** * Create an instance of the {@link IntegrationFlowBuilder} from a {@link Supplier} bean available in the context. * The name of the bean must be provided via `spring.cloud.stream.function.name` property. @@ -113,6 +118,21 @@ public class IntegrationFlowFunctionSupport { return flowBuilder; } + /** + * + * @param inputChannel + * @param outputChannel + * @return + */ + public IntegrationFlowBuilder integrationFlowForFunction(SubscribableChannel inputChannel, MessageChannel outputChannel) { + IntegrationFlowBuilder flowBuilder = IntegrationFlows.from(inputChannel).bridge(); + + if (!this.andThenFunction(flowBuilder, outputChannel)) { + flowBuilder = flowBuilder.channel(outputChannel); + } + return flowBuilder; + } + /** * Add a {@link Function} bean to the end of an integration flow. * The name of the bean must be provided via `spring.cloud.stream.function.name` property. @@ -131,7 +151,12 @@ public class IntegrationFlowFunctionSupport { new FunctionInvoker<>(this.functionProperties.getName(), this.functionCatalog, this.functionInspector, this.messageConverterFactory, this.errorChannel); - subscribeToInput(functionInvoker, flowBuilder.toReactivePublisher(), outputChannel::send); + if (outputChannel != null) { + subscribeToInput(functionInvoker, flowBuilder.toReactivePublisher(), outputChannel::send); + } + else { + subscribeToInput(functionInvoker, flowBuilder.toReactivePublisher(), null); + } return true; } return false; diff --git a/spring-cloud-stream/src/main/resources/META-INF/spring.factories b/spring-cloud-stream/src/main/resources/META-INF/spring.factories index 4fe70199f..f7a9928b7 100644 --- a/spring-cloud-stream/src/main/resources/META-INF/spring.factories +++ b/spring-cloud-stream/src/main/resources/META-INF/spring.factories @@ -3,6 +3,8 @@ org.springframework.cloud.stream.config.ChannelBindingAutoConfiguration,\ org.springframework.cloud.stream.config.BindersHealthIndicatorAutoConfiguration,\ org.springframework.cloud.stream.config.ChannelsEndpointAutoConfiguration,\ org.springframework.cloud.stream.config.BindingsEndpointAutoConfiguration,\ -org.springframework.cloud.stream.config.BindingServiceConfiguration +org.springframework.cloud.stream.config.BindingServiceConfiguration,\ +org.springframework.cloud.stream.function.FunctionConfiguration + diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/NewSourceAsSupplierTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/NewSourceAsSupplierTests.java new file mode 100644 index 000000000..5cd6b3274 --- /dev/null +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/NewSourceAsSupplierTests.java @@ -0,0 +1,132 @@ +package org.springframework.cloud.stream.function; + +import java.util.Date; +import java.util.function.Consumer; +import java.util.function.Function; +import java.util.function.Supplier; + +import org.junit.Test; +import org.springframework.boot.WebApplicationType; +import org.springframework.boot.autoconfigure.EnableAutoConfiguration; +import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.cloud.stream.annotation.EnableBinding; +import org.springframework.cloud.stream.binder.test.InputDestination; +import org.springframework.cloud.stream.binder.test.OutputDestination; +import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +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.annotation.Bean; +import org.springframework.messaging.Message; +import org.springframework.messaging.support.GenericMessage; + +public class NewSourceAsSupplierTests { + + @Test + public void testSourceFromSupplier() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SourceFromSupplier.class)).web( + WebApplicationType.NONE).run("--spring.cloud.stream.function.name=date", "--spring.jmx.enabled=false")) { + + OutputDestination target = context.getBean(OutputDestination.class); + Message sourceMessage = target.receive(10000); + System.out.println(sourceMessage); +// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8)); +// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8)); +// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8)); + //etc + } + } + + @Test + public void testProcessorFromFunction() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(ProcessorFromFunction.class)).web( + WebApplicationType.NONE).run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false")) { + + InputDestination source = context.getBean(InputDestination.class); + source.send(new GenericMessage("fopo".getBytes())); + OutputDestination target = context.getBean(OutputDestination.class); + Message targetMessage = target.receive(10000); + System.out.println(new String(targetMessage.getPayload())); +// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8)); +// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8)); +// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8)); + //etc + } + } + + @Test + public void testSinkFromConsumer() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumer.class)).web( + WebApplicationType.NONE).run("--spring.cloud.stream.function.name=sink", "--spring.jmx.enabled=false")) { + + InputDestination source = context.getBean(InputDestination.class); + source.send(new GenericMessage("fopo".getBytes())); +// OutputDestination target = context.getBean(OutputDestination.class); +// Message targetMessage = target.receive(10000); +// System.out.println(new String(targetMessage.getPayload())); +// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8)); +// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8)); +// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8)); + //etc + } + } + + @Test + public void testSinkFromConsumerNoEnableBinding() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumerNoEnableBinding.class)).web( + WebApplicationType.NONE).run("--spring.cloud.stream.function.name=sink", "--spring.jmx.enabled=false")) { + + InputDestination source = context.getBean(InputDestination.class); + source.send(new GenericMessage("Hello No Binding".getBytes())); +// OutputDestination target = context.getBean(OutputDestination.class); +// Message targetMessage = target.receive(10000); +// System.out.println(new String(targetMessage.getPayload())); +// assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8)); +// assertThat(target.receive(10000).getPayload()).isEqualTo("2".getBytes(StandardCharsets.UTF_8)); +// assertThat(target.receive(10000).getPayload()).isEqualTo("3".getBytes(StandardCharsets.UTF_8)); + //etc + } + } + + + @EnableAutoConfiguration + @EnableBinding(Source.class) + public static class SourceFromSupplier { + @Bean + public Supplier date() { + return () -> new Date(); + } + } + + @EnableAutoConfiguration + @EnableBinding(Processor.class) + public static class ProcessorFromFunction { + @Bean + public Function toUpperCase() { + return s -> s.toUpperCase(); + } + } + + @EnableAutoConfiguration + @EnableBinding(Sink.class) + public static class SinkFromConsumer { + @Bean + public Consumer sink() { + return s -> System.out.println(s); + } + } + + @EnableAutoConfiguration +// @EnableBinding(Sink.class) + public static class SinkFromConsumerNoEnableBinding { + @Bean + public Consumer sink() { + return s -> System.out.println("==> " + s); + } + } +}