From 720dce0b181a0203b625c570ce3c80d1bcc0d68d Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Thu, 10 Dec 2020 15:56:35 +0100 Subject: [PATCH] GH-2065 Add test to validate deserialization issue The actual fix was on s-c-function side https://github.com/spring-cloud/spring-cloud-function/commit/7403a51464d6f7318b3b41d50c2fbebdb05dbcc0 Resolves #2065 --- pom.xml | 2 +- .../ImplicitFunctionBindingTests.java | 5 ++- .../MultipleInputOutputFunctionTests.java | 37 +++++++++++++++++++ 3 files changed, 42 insertions(+), 2 deletions(-) diff --git a/pom.xml b/pom.xml index 89efbd080..25ad5c116 100644 --- a/pom.xml +++ b/pom.xml @@ -25,7 +25,7 @@ 1.8 2020.0.0-RC2 2.1 - 3.1.0-M5 + 3.1.0-SNAPSHOT true true true diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java index de2e624e1..e1b3ccedb 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ImplicitFunctionBindingTests.java @@ -1361,7 +1361,10 @@ public class ImplicitFunctionBindingTests { @Bean public Function>>, Message>>> funcB() { - return v -> v; + return v -> { + System.out.println(v); + return v; + }; } } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java index 4281cdf95..80bff1bfb 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/MultipleInputOutputFunctionTests.java @@ -151,6 +151,31 @@ public class MultipleInputOutputFunctionTests { } } + @Test + public void testMultiInputMessageSingleOutput() { + try (ConfigurableApplicationContext context = new SpringApplicationBuilder( + TestChannelBinderConfiguration.getCompleteConfiguration( + ReactiveFunctionConfiguration.class)) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false", + "--spring.cloud.function.definition=multiInputSingleOutputMessage")) { + context.getBean(InputDestination.class); + + InputDestination inputDestination = context.getBean(InputDestination.class); + OutputDestination outputDestination = context.getBean(OutputDestination.class); + + Message stringInputMessage = MessageBuilder.withPayload("one".getBytes()).build(); + Message integerInputMessage = MessageBuilder.withPayload("1".getBytes()).build(); + inputDestination.send(stringInputMessage, "multiInputSingleOutputMessage-in-0"); + inputDestination.send(integerInputMessage, "multiInputSingleOutputMessage-in-1"); + + Message outputMessage = outputDestination.receive(); + assertThat(outputMessage.getPayload()).isEqualTo("one".getBytes()); + outputMessage = outputDestination.receive(0, "multiInputSingleOutputMessage-out-0"); + assertThat(outputMessage.getPayload()).isEqualTo("1".getBytes()); + } + } + @Test public void testSingleInputMultiOutput() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( @@ -336,6 +361,18 @@ public class MultipleInputOutputFunctionTests { }; } + @Bean + public Function>, Flux>>, Flux> multiInputSingleOutputMessage() { + return tuple -> { + Flux stringStream = tuple.getT1().map(m -> m.getPayload()); + Flux intStream = tuple.getT2().map(i -> { + int v = i.getPayload(); + return String.valueOf(v); + }); + return Flux.merge(stringStream, intStream); + }; + } + @Bean @SuppressWarnings({ "unchecked", "rawtypes" }) public static Function, Tuple2, Flux>> singleInputMultipleOutputs() {