From b7a7470e924c12d2dca1762fd2b23384b02f2802 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 29 Aug 2018 05:31:11 -0500 Subject: [PATCH] GH-1458 added support for MessageingGateway to Source Added support for wiring MessagingGateway as Source to provide the same function support as for the Supplier Resolves #1458 --- spring-cloud-stream/pom.xml | 12 +++++ .../function/FunctionCatalogWrapper.java | 4 ++ .../function/FunctionConfiguration.java | 7 ++- .../IntegrationFlowFunctionSupport.java | 45 +++++++++++++++--- .../GreenfieldFunctionEnableBindingTests.java | 46 +++++++++++++++++++ 5 files changed, 106 insertions(+), 8 deletions(-) diff --git a/spring-cloud-stream/pom.xml b/spring-cloud-stream/pom.xml index 80e66b46b..b42cbd0db 100644 --- a/spring-cloud-stream/pom.xml +++ b/spring-cloud-stream/pom.xml @@ -72,6 +72,18 @@ spring-boot-autoconfigure-processor true + + + org.springframework.integration + spring-integration-http + 5.1.0.M2 + test + + + org.springframework.boot + spring-boot-starter-web + test + diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionCatalogWrapper.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionCatalogWrapper.java index f6f7b6701..6d5c4af62 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionCatalogWrapper.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionCatalogWrapper.java @@ -21,6 +21,7 @@ import org.springframework.util.Assert; /** * @author David Turanski + * @author Oleg Zhurakousky * * @since 2.1 **/ @@ -44,4 +45,7 @@ class FunctionCatalogWrapper { return lookup(null, name); } + boolean contains(Class functionType, String name) { + return catalog.lookup(functionType, name) != null; + } } 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 c42654ef4..8572b92ef 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,6 +16,8 @@ package org.springframework.cloud.stream.function; +import java.util.function.Supplier; + import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty; @@ -52,7 +54,6 @@ public class FunctionConfiguration { public IntegrationFlowFunctionSupport functionSupport(FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties) { - return new IntegrationFlowFunctionSupport(functionCatalog, functionInspector, messageConverterFactory, functionProperties); } @@ -76,7 +77,9 @@ public class FunctionConfiguration { return functionSupport.integrationFlowForFunction(sink.input(), null).get(); } else if (source != null) { - return functionSupport.integrationFlowFromNamedSupplier().channel(this.source.output()).get(); + return functionSupport.containsFunction(Supplier.class) + ? functionSupport.integrationFlowFromNamedSupplier().channel(this.source.output()).get() + : null; } 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 27e977f17..8f974ad8e 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,10 +26,12 @@ import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import org.springframework.beans.factory.annotation.Autowired; +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.function.core.FluxSupplier; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; +import org.springframework.integration.channel.FluxMessageChannel; import org.springframework.integration.dsl.IntegrationFlowBuilder; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.messaging.Message; @@ -45,7 +47,7 @@ import org.springframework.util.StringUtils; * * @since 2.1 */ -public class IntegrationFlowFunctionSupport { +class IntegrationFlowFunctionSupport { private final FunctionCatalogWrapper functionCatalog; @@ -64,7 +66,7 @@ public class IntegrationFlowFunctionSupport { * @param messageConverterFactory * @param functionProperties */ - public IntegrationFlowFunctionSupport(FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector, + IntegrationFlowFunctionSupport(FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties) { Assert.notNull(functionCatalog, "'functionCatalog' must not be null"); @@ -77,8 +79,21 @@ public class IntegrationFlowFunctionSupport { this.functionProperties = functionProperties; } + /** + * Determines if function specified via 'spring.cloud.stream.function.definition' + * property can be located in {@link FunctionCatalog} + * + * @param typeOfFunction must be Supplier, Function or Consumer + * @return + */ + public boolean containsFunction(Class typeOfFunction) { + return StringUtils.hasText(this.functionProperties.getDefinition()) + && this.functionCatalog.contains(typeOfFunction, this.functionProperties.getDefinition()); + } + public FunctionType getCurrentFunctionType() { - return functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getDefinition())).getType(); + FunctionType functionType = functionInspector.getRegistration(functionCatalog.lookup(this.functionProperties.getDefinition())).getType(); + return functionType; } /** @@ -162,6 +177,23 @@ public class IntegrationFlowFunctionSupport { return false; } + public boolean andThenFunction(FluxMessageChannel fluxChannel, MessageChannel outputChannel) { + if (StringUtils.hasText(this.functionProperties.getDefinition())) { + FunctionInvoker functionInvoker = + new FunctionInvoker<>(this.functionProperties.getDefinition(), this.functionCatalog, + this.functionInspector, this.messageConverterFactory, this.errorChannel); + + if (outputChannel != null) { + subscribeToInput(functionInvoker, fluxChannel, outputChannel::send); + } + else { + subscribeToInput(functionInvoker, fluxChannel, null); + } + return true; + } + return false; + } + private Mono subscribeToOutput(Consumer> outputProcessor, Publisher> outputPublisher) { @@ -171,11 +203,12 @@ public class IntegrationFlowFunctionSupport { return output.then(); } - private void subscribeToInput(FunctionInvoker functionInvoker, Publisher> publisher, + @SuppressWarnings("unchecked") + private void subscribeToInput(FunctionInvoker functionInvoker, Publisher publisher, Consumer> outputProcessor) { - Flux> inputPublisher = Flux.from(publisher); - subscribeToOutput(outputProcessor, functionInvoker.apply(inputPublisher)).subscribe(); + Flux inputPublisher = Flux.from(publisher); + subscribeToOutput(outputProcessor, functionInvoker.apply((Flux>) inputPublisher)).subscribe(); } } 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 index 5eac13e0e..4fc42eaa9 100644 --- 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 @@ -24,9 +24,11 @@ import java.util.function.Supplier; import org.junit.Test; +import org.springframework.beans.factory.annotation.Autowired; import org.springframework.boot.WebApplicationType; import org.springframework.boot.autoconfigure.EnableAutoConfiguration; import org.springframework.boot.builder.SpringApplicationBuilder; +import org.springframework.boot.test.web.client.TestRestTemplate; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.test.InputDestination; import org.springframework.cloud.stream.binder.test.OutputDestination; @@ -37,7 +39,12 @@ 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.http.HttpMethod; +import org.springframework.integration.channel.FluxMessageChannel; import org.springframework.integration.channel.QueueChannel; +import org.springframework.integration.http.dsl.Http; +import org.springframework.integration.http.dsl.HttpRequestHandlerEndpointSpec; +import org.springframework.integration.http.inbound.HttpRequestHandlingEndpointSupport; import org.springframework.messaging.Message; import org.springframework.messaging.PollableChannel; import org.springframework.messaging.support.GenericMessage; @@ -94,6 +101,19 @@ public class GreenfieldFunctionEnableBindingTests { } } + @Test + public void testHttpEndpoint() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration(HttpInboundEndpoint.class)).web( + WebApplicationType.SERVLET).run("--spring.cloud.stream.function.definition=upperCase", "--spring.jmx.enabled=false")) { + TestRestTemplate restTemplate = new TestRestTemplate(); + restTemplate.postForLocation("http://localhost:8080", "hello"); + + OutputDestination target = context.getBean(OutputDestination.class); + assertThat(target.receive(10000).getPayload()).isEqualTo("HELLO".getBytes(StandardCharsets.UTF_8)); + } + } + @EnableAutoConfiguration @EnableBinding(Source.class) @@ -128,4 +148,30 @@ public class GreenfieldFunctionEnableBindingTests { }; } } + + @EnableAutoConfiguration + @EnableBinding(Source.class) + public static class HttpInboundEndpoint { + + @Autowired + private Source source; + + @Bean + public Function upperCase() { + return s -> s.toUpperCase(); + } + + @Bean + public HttpRequestHandlingEndpointSupport doFoo(IntegrationFlowFunctionSupport functionSupport) { + FluxMessageChannel fluxChannel = new FluxMessageChannel(); + HttpRequestHandlerEndpointSpec httpRequestHandler = Http + .inboundChannelAdapter("/*") + .requestMapping(requestMapping -> requestMapping.methods(HttpMethod.POST) + .consumes("*/*")) + .requestChannel(fluxChannel); + + functionSupport.andThenFunction(fluxChannel, source.output()); + return httpRequestHandler.get(); + } + } }