added another lookup method accounting for MimeTypes
polishing
This commit is contained in:
@@ -29,8 +29,9 @@ public interface FunctionCatalog {
|
|||||||
|
|
||||||
|
|
||||||
|
|
||||||
default <T> T lookupRaw(String name, MimeType... acceptedOutputTypes) {
|
default <T> T lookup(String name, MimeType... acceptedOutputTypes) {
|
||||||
return null;
|
throw new UnsupportedOperationException(
|
||||||
|
"This instance of FunctionCatalog does not support this operation");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package org.springframework.cloud.function.context.catalog;
|
|||||||
import java.lang.reflect.ParameterizedType;
|
import java.lang.reflect.ParameterizedType;
|
||||||
import java.lang.reflect.Type;
|
import java.lang.reflect.Type;
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
|
import java.util.Collections;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
|
|
||||||
import org.apache.commons.logging.Log;
|
import org.apache.commons.logging.Log;
|
||||||
@@ -11,10 +12,12 @@ import org.reactivestreams.Publisher;
|
|||||||
import org.springframework.cloud.function.context.FunctionRegistration;
|
import org.springframework.cloud.function.context.FunctionRegistration;
|
||||||
import org.springframework.core.convert.ConversionService;
|
import org.springframework.core.convert.ConversionService;
|
||||||
import org.springframework.messaging.Message;
|
import org.springframework.messaging.Message;
|
||||||
|
import org.springframework.messaging.MessageHeaders;
|
||||||
import org.springframework.messaging.converter.MessageConverter;
|
import org.springframework.messaging.converter.MessageConverter;
|
||||||
import org.springframework.messaging.support.MessageBuilder;
|
import org.springframework.messaging.support.MessageBuilder;
|
||||||
import org.springframework.util.Assert;
|
import org.springframework.util.Assert;
|
||||||
import org.springframework.util.CollectionUtils;
|
import org.springframework.util.CollectionUtils;
|
||||||
|
import org.springframework.util.MimeType;
|
||||||
|
|
||||||
import reactor.core.publisher.Flux;
|
import reactor.core.publisher.Flux;
|
||||||
import reactor.core.publisher.Mono;
|
import reactor.core.publisher.Mono;
|
||||||
@@ -90,34 +93,35 @@ class FunctionTypeConversionHelper {
|
|||||||
}
|
}
|
||||||
|
|
||||||
@SuppressWarnings("rawtypes")
|
@SuppressWarnings("rawtypes")
|
||||||
Object convertOutputIfNecessary(Object output) {
|
Object convertOutputIfNecessary(Object output, MimeType... acceptedOutputTypes) {
|
||||||
List<Object> convertedResults = new ArrayList<Object>();
|
List<Object> convertedResults = new ArrayList<Object>();
|
||||||
if (output instanceof Tuple2) {
|
if (output instanceof Tuple2) {
|
||||||
convertedResults.add(this.doConvert(((Tuple2)output).getT1(), byte[].class, true));
|
convertedResults.add(this.doConvert(((Tuple2)output).getT1(), acceptedOutputTypes[0]));
|
||||||
convertedResults.add(this.doConvert(((Tuple2)output).getT2(), byte[].class, true));
|
convertedResults.add(this.doConvert(((Tuple2)output).getT2(), acceptedOutputTypes[1]));
|
||||||
}
|
}
|
||||||
if (output instanceof Tuple3) {
|
if (output instanceof Tuple3) {
|
||||||
convertedResults.add(this.doConvert(((Tuple3)output).getT3(), byte[].class, true));
|
convertedResults.add(this.doConvert(((Tuple3)output).getT3(), acceptedOutputTypes[2]));
|
||||||
}
|
}
|
||||||
if (output instanceof Tuple4) {
|
if (output instanceof Tuple4) {
|
||||||
convertedResults.add(this.doConvert(((Tuple4)output).getT4(), byte[].class, true));
|
convertedResults.add(this.doConvert(((Tuple4)output).getT4(), acceptedOutputTypes[3]));
|
||||||
}
|
}
|
||||||
if (output instanceof Tuple5) {
|
if (output instanceof Tuple5) {
|
||||||
convertedResults.add(this.doConvert(((Tuple5)output).getT5(), byte[].class, true));
|
convertedResults.add(this.doConvert(((Tuple5)output).getT5(), acceptedOutputTypes[4]));
|
||||||
}
|
}
|
||||||
if (output instanceof Tuple6) {
|
if (output instanceof Tuple6) {
|
||||||
convertedResults.add(this.doConvert(((Tuple6)output).getT6(), byte[].class, true));
|
convertedResults.add(this.doConvert(((Tuple6)output).getT6(), acceptedOutputTypes[5]));
|
||||||
}
|
}
|
||||||
if (output instanceof Tuple7) {
|
if (output instanceof Tuple7) {
|
||||||
convertedResults.add(this.doConvert(((Tuple7)output).getT7(), byte[].class, true));
|
convertedResults.add(this.doConvert(((Tuple7)output).getT7(), acceptedOutputTypes[6]));
|
||||||
}
|
}
|
||||||
if (output instanceof Tuple8) {
|
if (output instanceof Tuple8) {
|
||||||
convertedResults.add(this.doConvert(((Tuple8)output).getT8(), byte[].class, true));
|
convertedResults.add(this.doConvert(((Tuple8)output).getT8(), acceptedOutputTypes[7]));
|
||||||
|
}
|
||||||
|
|
||||||
|
if (!CollectionUtils.isEmpty(convertedResults)) {
|
||||||
|
output = Tuples.fromArray(convertedResults.toArray());
|
||||||
}
|
}
|
||||||
|
|
||||||
output = CollectionUtils.isEmpty(convertedResults)
|
|
||||||
? this.doConvert(output, byte[].class, true)
|
|
||||||
: Tuples.fromArray(convertedResults.toArray());
|
|
||||||
return output;
|
return output;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -167,30 +171,43 @@ class FunctionTypeConversionHelper {
|
|||||||
return (Class<?>) targetType;
|
return (Class<?>) targetType;
|
||||||
}
|
}
|
||||||
|
|
||||||
private Object doConvert(Object incoming, Type targetType) {
|
|
||||||
return this.doConvert(incoming, targetType, false);
|
|
||||||
}
|
|
||||||
|
|
||||||
@SuppressWarnings({ "unchecked", "rawtypes" })
|
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||||
private Object doConvert(Object incoming, Type targetType, boolean toMessage) {
|
private Object doConvert(Object incoming, Type targetType) {
|
||||||
Class<?> actualType = this.getRawType(targetType);
|
Class<?> actualType = this.getRawType(targetType);
|
||||||
if (incoming instanceof Publisher) {
|
if (incoming instanceof Publisher) {
|
||||||
if (!actualType.isAssignableFrom(Void.class)) {
|
if (!actualType.isAssignableFrom(Void.class)) {
|
||||||
incoming = incoming instanceof Mono
|
incoming = incoming instanceof Mono
|
||||||
? Mono.from((Publisher) incoming).map(value -> this.doConvertArgument(value, targetType, actualType, toMessage))
|
? Mono.from((Publisher) incoming).map(value -> this.doConvertArgument(value, targetType, actualType))
|
||||||
: Flux.from((Publisher) incoming).map(value -> this.doConvertArgument(value, targetType, actualType, toMessage));
|
: Flux.from((Publisher) incoming).map(value -> this.doConvertArgument(value, targetType, actualType));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
else {
|
else {
|
||||||
Assert.isTrue(!Publisher.class.isAssignableFrom(this.functionRegistration.getType().getInputWrapper()),
|
Assert.isTrue(!Publisher.class.isAssignableFrom(this.functionRegistration.getType().getInputWrapper()),
|
||||||
"Invoking reactive function as imperative is not allowed. Function name(s): "
|
"Invoking reactive function as imperative is not allowed. Function name(s): "
|
||||||
+ this.functionRegistration.getNames());
|
+ this.functionRegistration.getNames());
|
||||||
incoming = this.doConvertArgument(incoming, targetType, actualType, toMessage);
|
incoming = this.doConvertArgument(incoming, targetType, actualType);
|
||||||
}
|
}
|
||||||
return incoming;
|
return incoming;
|
||||||
}
|
}
|
||||||
|
|
||||||
private Object doConvertArgument(Object incomingValue, Type targetType, Class<?> actualInputType, boolean toMessage) {
|
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||||
|
private Object doConvert(Object incoming, MimeType mimeType) {
|
||||||
|
MessageHeaders headers = new MessageHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, mimeType));
|
||||||
|
if (incoming instanceof Publisher) {
|
||||||
|
incoming = incoming instanceof Mono
|
||||||
|
? Mono.from((Publisher) incoming).map(value -> this.messageConverter.toMessage(value, headers))
|
||||||
|
: Flux.from((Publisher) incoming).map(value -> this.messageConverter.toMessage(value, headers));
|
||||||
|
}
|
||||||
|
else {
|
||||||
|
Assert.isTrue(!Publisher.class.isAssignableFrom(this.functionRegistration.getType().getInputWrapper()),
|
||||||
|
"Invoking reactive function as imperative is not allowed. Function name(s): "
|
||||||
|
+ this.functionRegistration.getNames());
|
||||||
|
incoming = this.messageConverter.toMessage(incoming, headers);
|
||||||
|
}
|
||||||
|
return incoming;
|
||||||
|
}
|
||||||
|
|
||||||
|
private Object doConvertArgument(Object incomingValue, Type targetType, Class<?> actualInputType) {
|
||||||
if (!Void.class.isAssignableFrom(actualInputType)) {
|
if (!Void.class.isAssignableFrom(actualInputType)) {
|
||||||
if (incomingValue instanceof Message<?>) {
|
if (incomingValue instanceof Message<?>) {
|
||||||
incomingValue = this.isMessage(targetType)
|
incomingValue = this.isMessage(targetType)
|
||||||
@@ -204,9 +221,6 @@ class FunctionTypeConversionHelper {
|
|||||||
incomingValue = this.conversionService.convert(incomingValue, actualInputType);
|
incomingValue = this.conversionService.convert(incomingValue, actualInputType);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (toMessage) {
|
|
||||||
incomingValue = MessageBuilder.withPayload(incomingValue).build();
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
else {
|
else {
|
||||||
incomingValue = null;
|
incomingValue = null;
|
||||||
|
|||||||
@@ -42,6 +42,8 @@ import org.springframework.context.ConfigurableApplicationContext;
|
|||||||
import org.springframework.core.convert.ConversionService;
|
import org.springframework.core.convert.ConversionService;
|
||||||
import org.springframework.lang.Nullable;
|
import org.springframework.lang.Nullable;
|
||||||
import org.springframework.messaging.converter.CompositeMessageConverter;
|
import org.springframework.messaging.converter.CompositeMessageConverter;
|
||||||
|
import org.springframework.util.MimeType;
|
||||||
|
import org.springframework.util.ObjectUtils;
|
||||||
import org.springframework.util.StringUtils;
|
import org.springframework.util.StringUtils;
|
||||||
|
|
||||||
import reactor.core.publisher.Flux;
|
import reactor.core.publisher.Flux;
|
||||||
@@ -76,7 +78,12 @@ public class LazyFunctionRegistry implements FunctionRegistry, FunctionInspector
|
|||||||
@SuppressWarnings("unchecked")
|
@SuppressWarnings("unchecked")
|
||||||
@Override
|
@Override
|
||||||
public <T> T lookup(Class<?> type, String definition) {
|
public <T> T lookup(Class<?> type, String definition) {
|
||||||
return (T) this.compose(type, definition, false);
|
return (T) this.compose(type, definition);
|
||||||
|
}
|
||||||
|
|
||||||
|
@SuppressWarnings("unchecked")
|
||||||
|
public <T> T lookup(String definition, MimeType... acceptedOutputTypes) {
|
||||||
|
return (T) this.compose(null, definition, acceptedOutputTypes);
|
||||||
}
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
@@ -163,10 +170,10 @@ public class LazyFunctionRegistry implements FunctionRegistry, FunctionInspector
|
|||||||
// }
|
// }
|
||||||
|
|
||||||
@SuppressWarnings({ "unchecked", "rawtypes" })
|
@SuppressWarnings({ "unchecked", "rawtypes" })
|
||||||
private Function<?,?> compose(Class<?> type, String definition, boolean raw) {
|
private Function<?,?> compose(Class<?> type, String definition, MimeType... acceptedOutputTypes) {
|
||||||
Function<?,?> resultFunction = null;
|
Function<?,?> resultFunction = null;
|
||||||
if (this.registrationsByName.containsKey(definition)) {
|
if (this.registrationsByName.containsKey(definition)) {
|
||||||
resultFunction = new FunctionInvocationWrapper(this.registrationsByName.get(definition), false);
|
resultFunction = new FunctionInvocationWrapper(this.registrationsByName.get(definition), false, acceptedOutputTypes);
|
||||||
}
|
}
|
||||||
else {
|
else {
|
||||||
String[] names = StringUtils.delimitedListToStringArray(definition.replaceAll(",", "|").trim(), "|");
|
String[] names = StringUtils.delimitedListToStringArray(definition.replaceAll(",", "|").trim(), "|");
|
||||||
@@ -195,7 +202,7 @@ public class LazyFunctionRegistry implements FunctionRegistry, FunctionInspector
|
|||||||
FunctionRegistration<Object> registration = new FunctionRegistration<>(function, name).type(funcType);
|
FunctionRegistration<Object> registration = new FunctionRegistration<>(function, name).type(funcType);
|
||||||
registrationsByFunction.putIfAbsent(function, registration);
|
registrationsByFunction.putIfAbsent(function, registration);
|
||||||
registrationsByName.putIfAbsent(name, registration);
|
registrationsByName.putIfAbsent(name, registration);
|
||||||
function = new FunctionInvocationWrapper(registration, false);
|
function = new FunctionInvocationWrapper(registration, false, acceptedOutputTypes);
|
||||||
if (resultFunction == null) {
|
if (resultFunction == null) {
|
||||||
resultFunction = (Function<?,?>) function;
|
resultFunction = (Function<?,?>) function;
|
||||||
}
|
}
|
||||||
@@ -216,7 +223,7 @@ public class LazyFunctionRegistry implements FunctionRegistry, FunctionInspector
|
|||||||
registration = new FunctionRegistration<Object>(resultFunction, composedNameBuilder.toString()).type(funcType);
|
registration = new FunctionRegistration<Object>(resultFunction, composedNameBuilder.toString()).type(funcType);
|
||||||
registrationsByFunction.putIfAbsent(resultFunction, registration);
|
registrationsByFunction.putIfAbsent(resultFunction, registration);
|
||||||
registrationsByName.putIfAbsent(composedNameBuilder.toString(), registration);
|
registrationsByName.putIfAbsent(composedNameBuilder.toString(), registration);
|
||||||
resultFunction = new FunctionInvocationWrapper(registration, true);
|
resultFunction = new FunctionInvocationWrapper(registration, true, acceptedOutputTypes);
|
||||||
}
|
}
|
||||||
previousFunctionType = funcType;
|
previousFunctionType = funcType;
|
||||||
prefix = "|";
|
prefix = "|";
|
||||||
@@ -247,10 +254,13 @@ public class LazyFunctionRegistry implements FunctionRegistry, FunctionInspector
|
|||||||
|
|
||||||
private final FunctionTypeConversionHelper functionTypeConversionHelper;
|
private final FunctionTypeConversionHelper functionTypeConversionHelper;
|
||||||
|
|
||||||
FunctionInvocationWrapper(FunctionRegistration<?> functionRegistration, boolean composed) {
|
private final MimeType[] acceptedOutputTypes;
|
||||||
|
|
||||||
|
FunctionInvocationWrapper(FunctionRegistration<?> functionRegistration, boolean composed, MimeType... acceptedOutputTypes) {
|
||||||
this.target = functionRegistration.getTarget();
|
this.target = functionRegistration.getTarget();
|
||||||
this.functionRegistration = functionRegistration;
|
this.functionRegistration = functionRegistration;
|
||||||
this.composed = composed;
|
this.composed = composed;
|
||||||
|
this.acceptedOutputTypes = acceptedOutputTypes;
|
||||||
this.functionTypeConversionHelper = new FunctionTypeConversionHelper(this.functionRegistration,
|
this.functionTypeConversionHelper = new FunctionTypeConversionHelper(this.functionRegistration,
|
||||||
conversionService, messageConverter);
|
conversionService, messageConverter);
|
||||||
}
|
}
|
||||||
@@ -306,8 +316,9 @@ public class LazyFunctionRegistry implements FunctionRegistry, FunctionInspector
|
|||||||
}
|
}
|
||||||
|
|
||||||
// ====
|
// ====
|
||||||
//result = this.functionTypeConversionHelper.convertOutputIfNecessary(result);
|
if (!ObjectUtils.isEmpty(this.acceptedOutputTypes)) {
|
||||||
//
|
result = this.functionTypeConversionHelper.convertOutputIfNecessary(result, this.acceptedOutputTypes);
|
||||||
|
}
|
||||||
|
|
||||||
return this.wrapOutputToReactiveIfNecessary(result);
|
return this.wrapOutputToReactiveIfNecessary(result);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -231,7 +231,7 @@ public class LazyFunctionRegistryMultiInOutTests {
|
|||||||
public void testMultiToMultiWithMessageByteArrayPayload() {
|
public void testMultiToMultiWithMessageByteArrayPayload() {
|
||||||
FunctionCatalog catalog = this.configureCatalog();
|
FunctionCatalog catalog = this.configureCatalog();
|
||||||
Function<Tuple3<Flux<Message<byte[]>>, Flux<Message<byte[]>>, Flux<Message<byte[]>>>, Tuple2<Flux<Message<byte[]>>, Mono<Message<byte[]>>>> multiTuMulti =
|
Function<Tuple3<Flux<Message<byte[]>>, Flux<Message<byte[]>>, Flux<Message<byte[]>>>, Tuple2<Flux<Message<byte[]>>, Mono<Message<byte[]>>>> multiTuMulti =
|
||||||
catalog.lookupRaw("multiTuMulti", MimeTypeUtils.parseMimeType("foo/bar"), MimeTypeUtils.parseMimeType("bar/*"));
|
catalog.lookup("multiTuMulti", MimeTypeUtils.parseMimeType("application/json"), MimeTypeUtils.parseMimeType("application/json"));
|
||||||
|
|
||||||
Flux<Message<byte[]>> firstFlux = Flux.just(
|
Flux<Message<byte[]>> firstFlux = Flux.just(
|
||||||
MessageBuilder.withPayload("Unlce".getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, "text/plain").build(),
|
MessageBuilder.withPayload("Unlce".getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, "text/plain").build(),
|
||||||
@@ -249,7 +249,6 @@ public class LazyFunctionRegistryMultiInOutTests {
|
|||||||
MessageBuilder.withPayload(one.array()).setHeader(MessageHeaders.CONTENT_TYPE, "octet-stream/integer").build(),
|
MessageBuilder.withPayload(one.array()).setHeader(MessageHeaders.CONTENT_TYPE, "octet-stream/integer").build(),
|
||||||
MessageBuilder.withPayload(two.array()).setHeader(MessageHeaders.CONTENT_TYPE, "octet-stream/integer").build());
|
MessageBuilder.withPayload(two.array()).setHeader(MessageHeaders.CONTENT_TYPE, "octet-stream/integer").build());
|
||||||
|
|
||||||
|
|
||||||
Tuple2<Flux<Message<byte[]>>, Mono<Message<byte[]>>> result = multiTuMulti.apply(Tuples.of(firstFlux, secondFlux, thirdFlux));
|
Tuple2<Flux<Message<byte[]>>, Mono<Message<byte[]>>> result = multiTuMulti.apply(Tuples.of(firstFlux, secondFlux, thirdFlux));
|
||||||
result.getT1().subscribe(v -> System.out.println("=> 1: " + v));
|
result.getT1().subscribe(v -> System.out.println("=> 1: " + v));
|
||||||
result.getT2().subscribe(v -> System.out.println("=> 2: " + v));
|
result.getT2().subscribe(v -> System.out.println("=> 2: " + v));
|
||||||
|
|||||||
Reference in New Issue
Block a user