Add support for MonoSupplier
This commit is contained in:
@@ -37,6 +37,7 @@ import org.springframework.cloud.function.core.FluxSupplier;
|
|||||||
import org.springframework.cloud.function.core.FluxToMonoFunction;
|
import org.springframework.cloud.function.core.FluxToMonoFunction;
|
||||||
import org.springframework.cloud.function.core.FluxedConsumer;
|
import org.springframework.cloud.function.core.FluxedConsumer;
|
||||||
import org.springframework.cloud.function.core.FluxedFunction;
|
import org.springframework.cloud.function.core.FluxedFunction;
|
||||||
|
import org.springframework.cloud.function.core.MonoSupplier;
|
||||||
import org.springframework.cloud.function.core.MonoToFluxFunction;
|
import org.springframework.cloud.function.core.MonoToFluxFunction;
|
||||||
import org.springframework.util.Assert;
|
import org.springframework.util.Assert;
|
||||||
import org.springframework.util.CollectionUtils;
|
import org.springframework.util.CollectionUtils;
|
||||||
@@ -155,17 +156,19 @@ public class FunctionRegistration<T> implements BeanNameAware {
|
|||||||
result = new FunctionRegistration<S>(target);
|
result = new FunctionRegistration<S>(target);
|
||||||
result.type(this.type.getType());
|
result.type(this.type.getType());
|
||||||
|
|
||||||
if (!type.isWrapper()) {
|
if (!this.type.isWrapper()) {
|
||||||
target = target instanceof Supplier
|
target = target instanceof Supplier
|
||||||
? (S) new FluxSupplier((Supplier<?>) target)
|
? (S) new FluxSupplier((Supplier<?>) target)
|
||||||
: target instanceof Function
|
: target instanceof Function
|
||||||
? (S) new FluxFunction((Function<?, ?>) target)
|
? (S) new FluxFunction((Function<?, ?>) target)
|
||||||
: (S) new FluxConsumer((Consumer<?>) target);
|
: (S) new FluxConsumer((Consumer<?>) target);
|
||||||
}
|
}
|
||||||
else if (Mono.class.isAssignableFrom(type.getOutputWrapper())) {
|
else if (Mono.class.isAssignableFrom(this.type.getOutputWrapper())) {
|
||||||
target = (S) new FluxToMonoFunction((Function) target);
|
target = target instanceof Supplier
|
||||||
|
? (S) new MonoSupplier((Supplier<?>) target)
|
||||||
|
: (S) new FluxToMonoFunction((Function<?, ?>) target);
|
||||||
}
|
}
|
||||||
else if (Mono.class.isAssignableFrom(type.getInputWrapper())) {
|
else if (Mono.class.isAssignableFrom(this.type.getInputWrapper())) {
|
||||||
target = (S) new MonoToFluxFunction((Function) target);
|
target = (S) new MonoToFluxFunction((Function) target);
|
||||||
}
|
}
|
||||||
else if (target instanceof Consumer) {
|
else if (target instanceof Consumer) {
|
||||||
|
|||||||
@@ -0,0 +1,50 @@
|
|||||||
|
/*
|
||||||
|
* Copyright 2012-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.function.core;
|
||||||
|
|
||||||
|
import java.util.function.Supplier;
|
||||||
|
|
||||||
|
import reactor.core.publisher.Mono;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* {@link Supplier} implementation that wraps a target Supplier so that the target's
|
||||||
|
* simple output type will be wrapped in a {@link Mono} instance.
|
||||||
|
*
|
||||||
|
* @param <T> output type of target supplier
|
||||||
|
* @author Mark Fisher
|
||||||
|
*/
|
||||||
|
public class MonoSupplier<T> implements Supplier<Mono<T>>, FluxWrapper<Supplier<T>> {
|
||||||
|
|
||||||
|
private final Supplier<T> supplier;
|
||||||
|
|
||||||
|
public MonoSupplier(Supplier<T> supplier) {
|
||||||
|
this.supplier = supplier;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public Supplier<T> getTarget() {
|
||||||
|
return this.supplier;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||||
|
public Mono<T> get() {
|
||||||
|
Object result = this.supplier.get();
|
||||||
|
return Mono.just((T) result);
|
||||||
|
}
|
||||||
|
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user