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 a33f9b61a..10f58e796 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 @@ -43,6 +43,7 @@ import org.springframework.boot.context.properties.EnableConfigurationProperties import org.springframework.cloud.function.context.FunctionCatalog; import org.springframework.cloud.function.context.FunctionRegistry; import org.springframework.cloud.function.context.catalog.FunctionInspector; +import org.springframework.cloud.function.context.config.RoutingFunction; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BinderType; import org.springframework.cloud.stream.binder.BinderTypeRegistry; @@ -278,14 +279,17 @@ public class BinderFactoryAutoConfiguration { if (!StringUtils.hasText(name) && catalog.size() == 0) { ((SmartInitializingSingleton) catalog).afterSingletonsInstantiated(); } + if (!StringUtils.hasText(name) && Boolean.parseBoolean( + environment.getProperty("spring.cloud.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); - } } - + if (StringUtils.hasText(name)) { + ((StandardEnvironment) environment).getSystemProperties() + .putIfAbsent("spring.cloud.stream.function.definition", name); + } return name; } 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 5a69dc55c..0a78376c5 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 @@ -60,6 +60,15 @@ public class FunctionConfiguration { StreamFunctionProperties functionProperties, BindingServiceProperties bindingServiceProperties) { ((SmartInitializingSingleton) functionCatalog).afterSingletonsInstantiated(); + +// if (functionCatalog.size() > 0) { +// String name = StringUtils.hasText(functionProperties.getDefinition()) ? functionProperties.getDefinition() : ""; +// Assert.notNull(functionCatalog.lookup(name), +// "Failed to locate function `" + functionProperties.getDefinition() +// + "' in function catalog. Available functions are " +// + functionCatalog.getNames(Function.class)); +// } + return new IntegrationFlowFunctionSupport(functionCatalog, functionInspector, messageConverterFactory, functionProperties, bindingServiceProperties); } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/RoutingFunctionEnvironmentPostProcessor.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/RoutingFunctionEnvironmentPostProcessor.java new file mode 100644 index 000000000..f0216bfdc --- /dev/null +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/RoutingFunctionEnvironmentPostProcessor.java @@ -0,0 +1,46 @@ +/* + * 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.boot.SpringApplication; +import org.springframework.boot.env.EnvironmentPostProcessor; +import org.springframework.cloud.function.context.config.RoutingFunction; +import org.springframework.core.env.ConfigurableEnvironment; +import org.springframework.core.env.StandardEnvironment; +import org.springframework.util.StringUtils; +/** + * + * @author Oleg Zhurakousky + * @since 2.2.1 + */ +class RoutingFunctionEnvironmentPostProcessor implements EnvironmentPostProcessor { + + @Override + public void postProcessEnvironment(ConfigurableEnvironment environment, SpringApplication application) { + String name = environment.getProperty("spring.cloud.stream.function.definition"); + if (StringUtils.hasText(name) && ( + name.equals(RoutingFunction.FUNCTION_NAME) || + name.contains(RoutingFunction.FUNCTION_NAME + "|") || + name.contains("|" + RoutingFunction.FUNCTION_NAME) + )) { + ((StandardEnvironment) environment).getSystemProperties() + .putIfAbsent("spring.cloud.function.routing.enabled", "true"); + } + } + +} 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 f7a9928b7..03891cfa7 100644 --- a/spring-cloud-stream/src/main/resources/META-INF/spring.factories +++ b/spring-cloud-stream/src/main/resources/META-INF/spring.factories @@ -6,5 +6,8 @@ org.springframework.cloud.stream.config.BindingsEndpointAutoConfiguration,\ org.springframework.cloud.stream.config.BindingServiceConfiguration,\ org.springframework.cloud.stream.function.FunctionConfiguration +org.springframework.boot.env.EnvironmentPostProcessor:\ +org.springframework.cloud.stream.function.RoutingFunctionEnvironmentPostProcessor + diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java new file mode 100644 index 000000000..c47018ed8 --- /dev/null +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/RoutingFunctionTests.java @@ -0,0 +1,268 @@ +/* + * 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 java.util.function.Function; + +import org.junit.After; +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.binder.test.InputDestination; +import org.springframework.cloud.stream.binder.test.OutputDestination; +import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration; +import org.springframework.context.ConfigurableApplicationContext; +import org.springframework.context.annotation.Bean; +import org.springframework.integration.support.MessageBuilder; +import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; +import org.springframework.util.MimeTypeUtils; + +import reactor.core.publisher.Flux; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * + * @author Oleg Zhurakousky + * @since 2.2.1 + */ +public class RoutingFunctionTests { + + @After + public void after() { + System.getProperties().remove("spring.cloud.function.routing.enabled"); + System.getProperties().remove("spring.cloud.stream.function.definition"); + } + + @Test + public void testDefaultRoutingFunctionBinding() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + RoutingFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.function.routing.enabled=true")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder + .withPayload("Hello".getBytes()) + .setHeader("function.name", "echo") + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN) + .build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes()); + + } + } + + @Test + public void testDefaultRoutingFunctionBindingFlux() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + RoutingFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.function.routing.enabled=true")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder + .withPayload("Hello".getBytes()) + .setHeader("function.name", "echoFlux") + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN) + .build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes()); + + } + } + + @Test // see RoutingFunctionEnvironmentPostProcessor + public void testExplicitRoutingFunctionBinding() { + + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + RoutingFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.stream.function.definition=router")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder + .withPayload("Hello".getBytes()) + .setHeader("function.name", "echo") + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN) + .build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("Hello".getBytes()); + + } + } + + @Test + public void testExplicitRoutingFunctionBindingWithCompositionAndRoutingEnabledExplicitly() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + RoutingFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.stream.function.definition=enrich|router", + "--spring.cloud.function.routing.enabled=true")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder + .withPayload("Hello".getBytes()).build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("HELLO".getBytes()); + } + } + + @Test + public void testExplicitRoutingFunctionBindingWithCompositionAndRoutingEnabledImplicitly() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + RoutingFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.stream.function.definition=enrich|router")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder + .withPayload("Hello".getBytes()).build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("HELLO".getBytes()); + } + } + + @Test + public void testExplicitRoutingFunctionBindingWithCompositionAndRoutingEnabledExplicitlyAndMoreComposition() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + RoutingFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.stream.function.definition=enrich|router|reverse", + "--spring.cloud.function.routing.enabled=true")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder + .withPayload("Hello".getBytes()).build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("OLLEH".getBytes()); + } + } + + @Test + public void testExplicitRoutingFunctionBindingWithCompositionAndRoutingEnabledImplicitlyAndMoreComposition() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + RoutingFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.stream.function.definition=enrich|router|reverse")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder + .withPayload("Hello".getBytes()).build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("OLLEH".getBytes()); + } + } + + + + @EnableAutoConfiguration + public static class RoutingFunctionConfiguration { + + @Bean + public Function echo() { + return x -> { + System.out.println("===> echo"); + return x; + }; + } + + @Bean + public Function, Flux> echoFlux() { + return flux -> flux.map(x -> { + System.out.println("===> echoFlux"); + return x; + }); + } + + @Bean + public Function, Message> enrich() { + return x -> { + System.out.println("===> enrich"); + return MessageBuilder.withPayload(x.getPayload()).setHeader("function.name", "uppercase").build(); + }; + } + + @Bean + public Function uppercase() { + return x -> { + System.out.println("===> uppercase"); + return x.toUpperCase(); + }; + } + + @Bean + public Function reverse() { + return x -> { + System.out.println("===> reverse"); + return new StringBuilder(x).reverse().toString(); + }; + } + } +}