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 628e1e060..f18f28fe8 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 @@ -169,8 +169,7 @@ class FunctionInvoker implements Function>, Flux) (value instanceof Message ? value - : this.messageConverter.toMessage(value, originalMessage.getHeaders(), - this.outputClass)); + : this.messageConverter.toMessage(value, originalMessage.getHeaders())); if (returnMessage == null && value.getClass().isAssignableFrom(this.outputClass)) { returnMessage = wrapOutputToMessage(value, originalMessage); 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 87cf1fd3a..99be5c010 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 @@ -49,6 +49,7 @@ import org.springframework.util.ReflectionUtils; import static org.assertj.core.api.Assertions.assertThat; + /** * @author Oleg Zhurakousky * @author Tolga Kavukcu @@ -58,6 +59,30 @@ public class FunctionInvokerTests { private static String testWithFluxedConsumerValue; + @Test + public void testSimpleEchoConfiguration() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + SimpleEchoConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.stream.function.definition=func")) { + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context + .getBean(OutputDestination.class); + + Message inputMessage = MessageBuilder + .withPayload("{\"name\":\"bob\"}".getBytes()).build(); + inputDestination.send(inputMessage); + + Message outputMessage = outputDestination.receive(); + System.out.println("Received: " + new String(outputMessage.getPayload())); + assertThat(outputMessage.getPayload()).isEqualTo("{\"name\":\"bob\"}".getBytes()); + + } + } + @Test public void testFunctionHonorsOutboundBindingContentType() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( @@ -304,6 +329,28 @@ public class FunctionInvokerTests { } } + @EnableAutoConfiguration + @EnableBinding(Processor.class) + public static class SimpleEchoConfiguration { + + @Bean + public Function func() { + return x -> x; + } + + public static class Person { + private String name; + + public String getName() { + return name; + } + + public void setName(String name) { + this.name = name; + } + } + } + @EnableAutoConfiguration @EnableBinding(Processor.class) public static class ConverterDoesNotProduceCTConfiguration {