Added support for type-less Message to FunctionInvoker
This commit is contained in:
@@ -87,6 +87,7 @@ class FunctionInvoker<I, O> implements Function<Flux<Message<I>>, Flux<Message<O
|
||||
.map(this::resolveArgument) // resolves argument type before invocation of user function
|
||||
.onErrorContinue((x, y) -> onError(x, (Message<I>) y))
|
||||
.transform(this.userFunction::apply) // invoke user function
|
||||
.onErrorContinue((x, y) -> onError(x, (Message<I>) y))
|
||||
.map(resultMessage -> toMessage(resultMessage, originalMessageRef.get())); // create output message
|
||||
}
|
||||
|
||||
@@ -136,7 +137,8 @@ class FunctionInvoker<I, O> implements Function<Flux<Message<I>>, Flux<Message<O
|
||||
}
|
||||
|
||||
private boolean shouldConvertFromMessage(Message<?> message) {
|
||||
return !message.getPayload().getClass().isAssignableFrom(this.inputClass) &&
|
||||
return !this.inputClass.isAssignableFrom(Message.class) &&
|
||||
!message.getPayload().getClass().isAssignableFrom(this.inputClass) &&
|
||||
!this.inputClass.isAssignableFrom(Object.class);
|
||||
}
|
||||
|
||||
|
||||
@@ -62,6 +62,11 @@ public class FunctionInvokerTests {
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class));
|
||||
outputMessage = pojoToPojoSameType.apply(Flux.just(inputMessage)).blockFirst();
|
||||
assertThat(inputMessage.getPayload()).isEqualTo(outputMessage.getPayload());
|
||||
|
||||
FunctionInvoker<Foo, Foo> messageToMessageNoType = new FunctionInvoker<>("messageToMessageNoType",
|
||||
new FunctionCatalogWrapper(context.getBean(FunctionCatalog.class)), context.getBean(FunctionInspector.class), context.getBean(CompositeMessageConverterFactory.class));
|
||||
outputMessage = messageToMessageNoType.apply(Flux.just(inputMessage)).blockFirst();
|
||||
assertThat(outputMessage).isInstanceOf(Message.class);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -72,6 +77,17 @@ public class FunctionInvokerTests {
|
||||
public Function<Message<Foo>, Message<Bar>> messageToMessageDifferentType() {
|
||||
return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<Message<?>, Message<?>> messageToMessageAnyType() {
|
||||
return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<Message, Message> messageToMessageNoType() {
|
||||
return x -> MessageBuilder.withPayload(new Bar()).copyHeaders(x.getHeaders()).build();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Function<Message<Foo>, Message<Foo>> messageToMessageSameType() {
|
||||
return x -> x;
|
||||
@@ -81,6 +97,7 @@ public class FunctionInvokerTests {
|
||||
public Function<Foo, Foo> pojoToPojoSameType() {
|
||||
return x -> x;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
private static class Foo {
|
||||
|
||||
Reference in New Issue
Block a user