From f3a3a5c9fe3425fd1bc1f1f3e5494912fea156d5 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 12 Feb 2019 15:51:25 +0100 Subject: [PATCH] GH-1611 Removed FluxedConsumerWrapper and FunctionCatalogWrapper - Upgraded to spring-cloud-function 2.0.2.BUILD-SNAPSHOT - General cleanup and polishing Resolves #1611 --- pom.xml | 2 +- .../function/FluxedConsumerWrapper.java | 52 ------------------ .../function/FunctionCatalogWrapper.java | 54 ------------------- .../function/FunctionConfiguration.java | 11 ++-- .../stream/function/FunctionInvoker.java | 12 ++--- .../IntegrationFlowFunctionSupport.java | 12 +++-- .../stream/function/FunctionInvokerTests.java | 18 +++---- 7 files changed, 29 insertions(+), 132 deletions(-) delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FluxedConsumerWrapper.java delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionCatalogWrapper.java diff --git a/pom.xml b/pom.xml index c9f18427a..7854a496a 100644 --- a/pom.xml +++ b/pom.xml @@ -28,7 +28,7 @@ Californium-RELEASE 3.0.3 2.1 - 2.0.0.BUILD-SNAPSHOT + 2.0.2.BUILD-SNAPSHOT true diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FluxedConsumerWrapper.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FluxedConsumerWrapper.java deleted file mode 100644 index e9019aa65..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FluxedConsumerWrapper.java +++ /dev/null @@ -1,52 +0,0 @@ -/* - * Copyright 2018-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 - * - * 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.util.function.Consumer; -import java.util.function.Function; - -import reactor.core.publisher.Flux; -import reactor.core.publisher.Mono; - -import org.springframework.cloud.function.core.FluxWrapper; - -/** - * @author Oleg Zhurakousky - * @since 2.1 - * - * Will most likely be moved to SCF - */ -class FluxedConsumerWrapper - implements Function, Mono>, FluxWrapper>> { - - private final Consumer> consumer; - - FluxedConsumerWrapper(Consumer> consumer) { - this.consumer = consumer; - } - - @Override - public Consumer> getTarget() { - return this.consumer; - } - - @Override - public Mono apply(Flux t) { - return Mono.fromRunnable(() -> this.consumer.accept(t)); - } - -} 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 deleted file mode 100644 index c3c82ac70..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionCatalogWrapper.java +++ /dev/null @@ -1,54 +0,0 @@ -/* - * Copyright 2018-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 - * - * 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 org.springframework.cloud.function.context.FunctionCatalog; -import org.springframework.util.Assert; - -/** - * @author David Turanski - * @author Oleg Zhurakousky - * @since 2.1 - **/ -class FunctionCatalogWrapper { - - private final FunctionCatalog catalog; - - FunctionCatalogWrapper(FunctionCatalog catalog) { - this.catalog = catalog; - } - - T lookup(Class functionType, String name) { - T function = this.catalog.lookup(functionType, name); - Assert.notNull(function, - functionType == null - ? String.format("User provided Function '%s' cannot be located.", - name) - : String.format("User provided %s '%s' cannot be located.", - functionType.getSimpleName(), name)); - return function; - } - - T lookup(String name) { - return lookup(null, name); - } - - boolean contains(Class functionType, String name) { - return this.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 f1f8cfbef..1407d9a37 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 @@ -49,6 +49,7 @@ import org.springframework.messaging.SubscribableChannel; @EnableConfigurationProperties(StreamFunctionProperties.class) public class FunctionConfiguration { + @Autowired(required = false) private Source source; @@ -60,7 +61,7 @@ public class FunctionConfiguration { @Bean public IntegrationFlowFunctionSupport functionSupport( - FunctionCatalogWrapper functionCatalog, FunctionInspector functionInspector, + FunctionCatalog functionCatalog, FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties, BindingServiceProperties bindingServiceProperties) { @@ -68,10 +69,10 @@ public class FunctionConfiguration { messageConverterFactory, functionProperties, bindingServiceProperties); } - @Bean - public FunctionCatalogWrapper functionCatalogWrapper(FunctionCatalog catalog) { - return new FunctionCatalogWrapper(catalog); - } +// @Bean +// public FunctionCatalogWrapper functionCatalogWrapper(FunctionCatalog catalog) { +// return new FunctionCatalogWrapper(catalog); +// } /** * This configuration creates an instance of the {@link IntegrationFlow} from standard diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java index 337d63929..628e1e060 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/function/FunctionInvoker.java @@ -20,7 +20,6 @@ import java.lang.reflect.Field; import java.time.Duration; import java.util.Map; import java.util.concurrent.atomic.AtomicReference; -import java.util.function.Consumer; import java.util.function.Function; import org.apache.commons.logging.Log; @@ -28,6 +27,7 @@ import org.apache.commons.logging.LogFactory; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; +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.binder.ConsumerProperties; @@ -85,15 +85,15 @@ class FunctionInvoker implements Function>, Flux implements Function>, Flux((Consumer) originalUserFunction) - : (Function, Flux>) originalUserFunction; + this.userFunction = (Function, Flux>) originalUserFunction; Assert.isInstanceOf(Function.class, this.userFunction); this.messageConverter = compositeMessageConverterFactory 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 4de6aa1d7..2918c67d3 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 @@ -48,7 +48,7 @@ import org.springframework.util.StringUtils; */ public class IntegrationFlowFunctionSupport { - private final FunctionCatalogWrapper functionCatalog; + private final FunctionCatalog functionCatalog; private final FunctionInspector functionInspector; @@ -59,7 +59,7 @@ public class IntegrationFlowFunctionSupport { @Autowired private MessageChannel errorChannel; - IntegrationFlowFunctionSupport(FunctionCatalogWrapper functionCatalog, + IntegrationFlowFunctionSupport(FunctionCatalog functionCatalog, FunctionInspector functionInspector, CompositeMessageConverterFactory messageConverterFactory, StreamFunctionProperties functionProperties, @@ -85,7 +85,7 @@ public class IntegrationFlowFunctionSupport { */ public boolean containsFunction(Class typeOfFunction) { return StringUtils.hasText(this.functionProperties.getDefinition()) - && this.functionCatalog.contains(typeOfFunction, + && this.catalogContains(typeOfFunction, this.functionProperties.getDefinition()); } @@ -99,7 +99,7 @@ public class IntegrationFlowFunctionSupport { */ public boolean containsFunction(Class typeOfFunction, String functionName) { return StringUtils.hasText(functionName) - && this.functionCatalog.contains(typeOfFunction, functionName); + && this.catalogContains(typeOfFunction, functionName); } /** @@ -231,6 +231,10 @@ public class IntegrationFlowFunctionSupport { return true; } + private boolean catalogContains(Class functionType, String name) { + return this.functionCatalog.lookup(functionType, name) != null; + } + private Mono subscribeToOutput(Consumer> outputProcessor, Publisher> outputPublisher) { 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 fbe8eb629..87cf1fd3a 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 @@ -125,7 +125,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("messageToMessageSameType"); FunctionInvoker messageToMessageSameType = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); Message outputMessage = messageToMessageSameType @@ -135,7 +135,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("pojoToPojoSameType"); FunctionInvoker pojoToPojoSameType = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); outputMessage = pojoToPojoSameType.apply(Flux.just(inputMessage)) @@ -145,7 +145,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("messageToMessageNoType"); FunctionInvoker messageToMessageNoType = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); outputMessage = messageToMessageNoType.apply(Flux.just(inputMessage)) @@ -155,7 +155,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("withException"); FunctionInvoker withException = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); @@ -181,7 +181,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("pojoToPojoNonEmptyPojo"); FunctionInvoker pojoToPojoSameType = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); Message outputMessage = pojoToPojoSameType.apply(Flux.just(inputMessage)) @@ -193,7 +193,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("messageToMessageNoType"); FunctionInvoker messageToMessageNoType = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); outputMessage = messageToMessageNoType.apply(Flux.just(inputMessageWithBaz)) @@ -203,7 +203,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("withExceptionNativeEncodingEnabled"); FunctionInvoker withException = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); @@ -229,7 +229,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("pojoToPojoNonEmptyPojo"); FunctionInvoker pojoToPojoSameType = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class)); Message outputMessage = pojoToPojoSameType.apply(Flux.just(inputMessage)) @@ -256,7 +256,7 @@ public class FunctionInvokerTests { functionProperties.setDefinition("fluxConsumer"); FunctionInvoker fluxedConsumer = new FunctionInvoker<>( functionProperties, - new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), + context.getBean(FunctionCatalog.class), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class));