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));