From 83f791a07ac287fc425ce7e0e92971ef760dc95d Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 27 Aug 2018 18:06:26 +0200 Subject: [PATCH] Added support for Function binding Added support for binding Suppliers, Function and Consumers as message producers and handlers for apps annotated with EnableBinding Changed '.name' property to '.definition' Resolves #1445 --- .../config/BindingServiceConfiguration.java | 5 +- .../function/FunctionConfiguration.java | 34 ++--- .../IntegrationFlowFunctionSupport.java | 20 +-- ...ies.java => StreamFunctionProperties.java} | 15 +- .../stream/function/FunctionInvokerTests.java | 2 +- .../GreenfieldFunctionEnableBindingTests.java | 127 +++++++++++++++++ .../function/NewSourceAsSupplierTests.java | 132 ------------------ .../ProcessorToFunctionsSupportTests.java | 6 +- .../SourceToFunctionsSupportTests.java | 20 +-- 9 files changed, 170 insertions(+), 191 deletions(-) rename spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/{FunctionProperties.java => StreamFunctionProperties.java} (72%) create mode 100644 spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java delete 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 12124b684..cbe058ae9 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 @@ -44,8 +44,7 @@ import org.springframework.cloud.stream.binding.InputBindingLifecycle; import org.springframework.cloud.stream.binding.MessageChannelStreamListenerResultAdapter; import org.springframework.cloud.stream.binding.OutputBindingLifecycle; import org.springframework.cloud.stream.binding.StreamListenerAnnotationBeanPostProcessor; -import org.springframework.cloud.stream.function.FunctionConfiguration; -import org.springframework.cloud.stream.function.FunctionProperties; +import org.springframework.cloud.stream.function.StreamFunctionProperties; import org.springframework.cloud.stream.micrometer.DestinationPublishingMetricsAutoConfiguration; import org.springframework.context.ApplicationListener; import org.springframework.context.annotation.Bean; @@ -78,7 +77,7 @@ import org.springframework.util.Assert; * @author Soby Chacko */ @Configuration -@EnableConfigurationProperties({ BindingServiceProperties.class, SpringIntegrationProperties.class, FunctionProperties.class }) +@EnableConfigurationProperties({ BindingServiceProperties.class, SpringIntegrationProperties.class, StreamFunctionProperties.class }) @Import({ DestinationPublishingMetricsAutoConfiguration.class, SpelExpressionConverterConfiguration.class }) @Role(BeanDefinition.ROLE_INFRASTRUCTURE) @ConditionalOnBean(value = BinderTypeRegistry.class, search = SearchStrategy.CURRENT) 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 37b4a6551..c42654ef4 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 @@ -17,14 +17,10 @@ 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; @@ -32,9 +28,6 @@ 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; /** * @@ -43,7 +36,7 @@ import org.springframework.util.ClassUtils; * @since 2.1 */ @Configuration -@ConditionalOnProperty("spring.cloud.stream.function.name") +@ConditionalOnProperty("spring.cloud.stream.function.definition") public class FunctionConfiguration { @Autowired(required=false) @@ -55,13 +48,10 @@ public class FunctionConfiguration { @Autowired(required=false) private Sink sink; - @Autowired - private ConfigurableListableBeanFactory registry; - @Bean public IntegrationFlowFunctionSupport functionSupport(FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory, - FunctionProperties functionProperties) { + StreamFunctionProperties functionProperties) { return new IntegrationFlowFunctionSupport(functionCatalog, functionInspector, messageConverterFactory, functionProperties); @@ -72,11 +62,13 @@ public class FunctionConfiguration { return new FunctionCatalogWrapper(catalog); } - - @ConditionalOnProperty("spring.cloud.stream.function.name") - @ConditionalOnMissingBean + /** + * This configuration creates an instance of {@link IntegrationFlow} appropriate for binding declared using EnableBinding. + * At the moment only Source, Processor and Sink are supported. + */ + @ConditionalOnMissingBean // starter apps typically already provide and instance of IntegrationFlow, so we don't need this one. @Bean - public IntegrationFlow foo(IntegrationFlowFunctionSupport functionSupport) { + public IntegrationFlow integrationFlowCreator(IntegrationFlowFunctionSupport functionSupport) { if (processor != null) { return functionSupport.integrationFlowForFunction(processor.input(), processor.output()).get(); } @@ -86,14 +78,6 @@ public class FunctionConfiguration { 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"); + throw new UnsupportedOperationException("Bindings other then Source, Processor and Sink are not currently supported"); } } 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 2bd645dc3..27e977f17 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 @@ -53,7 +53,7 @@ public class IntegrationFlowFunctionSupport { private final CompositeMessageConverterFactory messageConverterFactory; - private final FunctionProperties functionProperties; + private final StreamFunctionProperties functionProperties; @Autowired private MessageChannel errorChannel; @@ -65,7 +65,7 @@ public class IntegrationFlowFunctionSupport { * @param functionProperties */ public IntegrationFlowFunctionSupport(FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector, - CompositeMessageConverterFactory messageConverterFactory, FunctionProperties functionProperties) { + CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties) { Assert.notNull(functionCatalog, "'functionCatalog' must not be null"); Assert.notNull(functionInspector, "'functionInspector' must not be null"); @@ -78,18 +78,18 @@ public class IntegrationFlowFunctionSupport { } public FunctionType getCurrentFunctionType() { - return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getName())).getType(); + return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getDefinition())).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. + * The name of the bean must be provided via `spring.cloud.stream.function.definition` property. * @return instance of {@link IntegrationFlowBuilder} * @throws IllegalStateException if the named Supplier can not be located. */ public IntegrationFlowBuilder integrationFlowFromNamedSupplier() { - if (StringUtils.hasText(this.functionProperties.getName())) { - Supplier supplier = functionCatalog.lookup(Supplier.class, this.functionProperties.getName()); + if (StringUtils.hasText(this.functionProperties.getDefinition())) { + Supplier supplier = functionCatalog.lookup(Supplier.class, this.functionProperties.getDefinition()); if (supplier instanceof FluxSupplier) { supplier = ((FluxSupplier)supplier).getTarget(); } @@ -98,7 +98,7 @@ public class IntegrationFlowFunctionSupport { } throw new IllegalStateException( - "A Supplier is not specified in the 'spring.cloud.stream.function.name' property."); + "A Supplier is not specified in the 'spring.cloud.stream.function.definition' property."); } /** @@ -135,7 +135,7 @@ public class IntegrationFlowFunctionSupport { /** * 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. + * The name of the bean must be provided via `spring.cloud.stream.function.definition` property. *

* NOTE: If this method returns true, the integration flow is now represented * as a Reactive Streams {@link Publisher} bean. @@ -146,9 +146,9 @@ public class IntegrationFlowFunctionSupport { * @return true if {@link Function} was located and added and false if it wasn't. */ public boolean andThenFunction(IntegrationFlowBuilder flowBuilder, MessageChannel outputChannel) { - if (StringUtils.hasText(this.functionProperties.getName())) { + if (StringUtils.hasText(this.functionProperties.getDefinition())) { FunctionInvoker functionInvoker = - new FunctionInvoker<>(this.functionProperties.getName(), this.functionCatalog, + new FunctionInvoker<>(this.functionProperties.getDefinition(), this.functionCatalog, this.functionInspector, this.messageConverterFactory, this.errorChannel); if (outputChannel != null) { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionProperties.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamFunctionProperties.java similarity index 72% rename from spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionProperties.java rename to spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamFunctionProperties.java index 39488dae1..74677c86e 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionProperties.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/StreamFunctionProperties.java @@ -25,20 +25,21 @@ import org.springframework.boot.context.properties.ConfigurationProperties; * @since 2.1 */ @ConfigurationProperties("spring.cloud.stream.function") -public class FunctionProperties { +public class StreamFunctionProperties { /** - * Name of functions to bind. If several functions need to be composed into one, use pipes (e.g., 'fooFunc|barFunc') + * Definition of functions to bind. If several functions need to be composed + * into one, use pipes (e.g., 'fooFunc|barFunc') */ - private String name; + private String definition; - public String getName() { - return this.name; + public String getDefinition() { + return this.definition; } - public void setName(String name) { - this.name = name; + public void setDefinition(String definition) { + this.definition = definition; } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java index 6ec512a3d..e53e346ee 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/FunctionInvokerTests.java @@ -84,7 +84,7 @@ public class FunctionInvokerTests { } @Bean - public Function messageToMessageNoType() { + public Function, Message> messageToMessageNoType() { return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build(); } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java new file mode 100644 index 000000000..458e9a637 --- /dev/null +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/GreenfieldFunctionEnableBindingTests.java @@ -0,0 +1,127 @@ +/* + * Copyright 2018 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 + * + * http://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 java.nio.charset.StandardCharsets; +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.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.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.integration.channel.QueueChannel; +import org.springframework.messaging.Message; +import org.springframework.messaging.PollableChannel; +import org.springframework.messaging.support.GenericMessage; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * This test validates proper function binding for applications where EnableBinding is declared. + * + * @author Oleg Zhurakousky + */ +public class GreenfieldFunctionEnableBindingTests { + + @Test + public void testSourceFromSupplier() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SourceFromSupplier.class)).web( + WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=date", "--spring.jmx.enabled=false")) { + + OutputDestination target = context.getBean(OutputDestination.class); + Message sourceMessage = target.receive(10000); + Date date = (Date) new CompositeMessageConverterFactory().getMessageConverterForAllRegistered().fromMessage(sourceMessage, Date.class); + assertThat(date).isEqualTo(new Date(12345L)); + } + } + + @Test + public void testProcessorFromFunction() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(ProcessorFromFunction.class)).web( + WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false")) { + + InputDestination source = context.getBean(InputDestination.class); + source.send(new GenericMessage("John Doe".getBytes())); + OutputDestination target = context.getBean(OutputDestination.class); + assertThat(target.receive(10000).getPayload()).isEqualTo("JOHN DOE".getBytes(StandardCharsets.UTF_8)); + } + } + + @Test + public void testSinkFromConsumer() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(SinkFromConsumer.class)).web( + WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=sink", "--spring.jmx.enabled=false")) { + + InputDestination source = context.getBean(InputDestination.class); + PollableChannel result = context.getBean("result", PollableChannel.class); + source.send(new GenericMessage("John Doe".getBytes())); + assertThat(result.receive(10000).getPayload()).isEqualTo("John Doe"); + } + } + + + @EnableAutoConfiguration + @EnableBinding(Source.class) + public static class SourceFromSupplier { + @Bean + public Supplier date() { + return () -> new Date(12345L); + } + } + + @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 PollableChannel result() { + return new QueueChannel(); + } + @Bean + public Consumer sink(PollableChannel result) { + return s -> { + result.send(new GenericMessage(s)); + System.out.println(s); + }; + } + } +} 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 deleted file mode 100644 index 5cd6b3274..000000000 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/NewSourceAsSupplierTests.java +++ /dev/null @@ -1,132 +0,0 @@ -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); - } - } -} diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ProcessorToFunctionsSupportTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ProcessorToFunctionsSupportTests.java index 0e738d8d2..d74602d66 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ProcessorToFunctionsSupportTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ProcessorToFunctionsSupportTests.java @@ -74,7 +74,7 @@ public class ProcessorToFunctionsSupportTests { new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class)) .web(WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false"); + .run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false"); InputDestination source = context.getBean(InputDestination.class); OutputDestination target = context.getBean(OutputDestination.class); @@ -88,7 +88,7 @@ public class ProcessorToFunctionsSupportTests { new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class)) .web(WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false"); + .run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false"); InputDestination source = context.getBean(InputDestination.class); OutputDestination target = context.getBean(OutputDestination.class); @@ -102,7 +102,7 @@ public class ProcessorToFunctionsSupportTests { new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(ConsumerConfiguration.class)) .web(WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=log", "--spring.jmx.enabled=false"); + .run("--spring.cloud.stream.function.definition=log", "--spring.jmx.enabled=false"); InputDestination source = context.getBean(InputDestination.class); OutputDestination target = context.getBean(OutputDestination.class); 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 47b40bbb7..ae4a70c4b 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 @@ -72,7 +72,7 @@ public class SourceToFunctionsSupportTests { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class)).web( WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=toUpperCase", "--spring.jmx.enabled=false")) { + .run("--spring.cloud.stream.function.definition=toUpperCase", "--spring.jmx.enabled=false")) { OutputDestination target = context.getBean(OutputDestination.class); assertThat(target.receive(1000).getPayload()).isEqualTo("HELLO FUNCTION".getBytes(StandardCharsets.UTF_8)); @@ -84,7 +84,7 @@ public class SourceToFunctionsSupportTests { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfiguration.class)).web( WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false")) { + .run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false")) { OutputDestination target = context.getBean(OutputDestination.class); assertThat(target.receive(1000).getPayload()).isEqualTo( "HELLO FUNCTION:HELLO FUNCTION".getBytes(StandardCharsets.UTF_8)); @@ -97,7 +97,7 @@ public class SourceToFunctionsSupportTests { new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfigurationNoConversionPossible.class)) .web(WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf", + .run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false")) { PollableChannel errorChannel = context.getBean("errorChannel", PollableChannel.class); OutputDestination target = context.getBean(OutputDestination.class); @@ -112,7 +112,7 @@ public class SourceToFunctionsSupportTests { new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(FunctionsConfigurationNoConversionPossible.class)) .web(WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=toUpperCase|concatWithSelf", + .run("--spring.cloud.stream.function.definition=toUpperCase|concatWithSelf", "--spring.jmx.enabled=false")) { PollableChannel errorChannel = context.getBean("errorChannel", PollableChannel.class); OutputDestination target = context.getBean(OutputDestination.class); @@ -125,7 +125,7 @@ public class SourceToFunctionsSupportTests { public void testMessageSourceIsCreatedFromProvidedSupplier() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web( - WebApplicationType.NONE).run("--spring.cloud.stream.function.name=number", "--spring.jmx.enabled=false")) { + WebApplicationType.NONE).run("--spring.cloud.stream.function.definition=number", "--spring.jmx.enabled=false")) { OutputDestination target = context.getBean(OutputDestination.class); assertThat(target.receive(10000).getPayload()).isEqualTo("1".getBytes(StandardCharsets.UTF_8)); @@ -140,7 +140,7 @@ public class SourceToFunctionsSupportTests { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web( WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=number|concatWithSelf", "--spring.jmx.enabled=false")) { + .run("--spring.cloud.stream.function.definition=number|concatWithSelf", "--spring.jmx.enabled=false")) { OutputDestination target = context.getBean(OutputDestination.class); assertThat(target.receive(10000).getPayload()).isEqualTo("11".getBytes(StandardCharsets.UTF_8)); @@ -155,7 +155,7 @@ public class SourceToFunctionsSupportTests { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web( WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=number|concatWithSelf|multiplyByTwo", + .run("--spring.cloud.stream.function.definition=number|concatWithSelf|multiplyByTwo", "--spring.jmx.enabled=false")) { OutputDestination target = context.getBean(OutputDestination.class); @@ -177,7 +177,7 @@ public class SourceToFunctionsSupportTests { new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)).web( WebApplicationType.NONE) - .run("--spring.cloud.stream.function.name=doesNotExist", "--spring.jmx.enabled=false"); + .run("--spring.cloud.stream.function.definition=doesNotExist", "--spring.jmx.enabled=false"); } @EnableAutoConfiguration @@ -297,11 +297,11 @@ public class SourceToFunctionsSupportTests { private Source source; @Autowired - private FunctionProperties functionProperties; + private StreamFunctionProperties functionProperties; @Bean public IntegrationFlow messageSourceFlow(IntegrationFlowFunctionSupport functionSupport) { - Assert.hasText(this.functionProperties.getName(), "Supplier name must be provided"); + Assert.hasText(this.functionProperties.getDefinition(), "Supplier name must be provided"); return functionSupport.integrationFlowFromNamedSupplier().channel(this.source.output()).get(); }