diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java index 18e560afd..2f0c82600 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java @@ -105,7 +105,7 @@ class ApplicationJsonMessageMarshallingConverter extends MappingJackson2MessageC } if (result == null) { if (message.getPayload() instanceof byte[] - && targetClass.isAssignableFrom(String.class)) { + && String.class.isAssignableFrom(targetClass)) { result = new String((byte[]) message.getPayload(), StandardCharsets.UTF_8); } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/tck/ContentTypeTckTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/tck/ContentTypeTckTests.java index ede576b16..d6644ea35 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/tck/ContentTypeTckTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binder/tck/ContentTypeTckTests.java @@ -745,7 +745,9 @@ public class ContentTypeTckTests { public Person echo(Object value) throws Exception { ObjectMapper mapper = new ObjectMapper(); // assume it is string because CT is text/plain - return mapper.readValue((String) value, Person.class); + return value instanceof byte[] + ? mapper.readValue((byte[]) value, Person.class) + : mapper.readValue((String) value, Person.class); } } @@ -760,7 +762,9 @@ public class ContentTypeTckTests { public Person echo(Message message) throws Exception { ObjectMapper mapper = new ObjectMapper(); // assume it is string because CT is text/plain - return mapper.readValue((String) message.getPayload(), Person.class); + return message.getPayload() instanceof byte[] + ? mapper.readValue((byte[]) message.getPayload(), Person.class) + : mapper.readValue((String) message.getPayload(), Person.class); } } 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 bba972ea7..086853df1 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 @@ -186,7 +186,7 @@ public class ImplicitFunctionBindingTests { // good, we expected it } - Function function = v -> v.toUpperCase(); + Function function = v -> new String(v).toUpperCase(); FunctionBindingTestUtils.bind(context, function); input.send(new GenericMessage("hello".getBytes())); @@ -1015,7 +1015,7 @@ public class ImplicitFunctionBindingTests { Message result = outputDestination.receive(2000); assertThat(result.getPayload()).isInstanceOf(byte[].class); // check output type - assertThat(new String((byte[]) result.getPayload())).isEqualTo("String"); // check input type + assertThat(new String((byte[]) result.getPayload())).isEqualTo("byte[]"); // check input type } try (ConfigurableApplicationContext context = new SpringApplicationBuilder( @@ -1053,7 +1053,7 @@ public class ImplicitFunctionBindingTests { Message result = outputDestination.receive(2000); assertThat(result.getPayload()).isInstanceOf(byte[].class); // check output type - assertThat(new String((byte[]) result.getPayload())).isEqualTo("String"); // check input type + assertThat(new String((byte[]) result.getPayload())).isEqualTo("byte[]"); // check input type } try (ConfigurableApplicationContext context = new SpringApplicationBuilder( TestChannelBinderConfiguration.getCompleteConfiguration(SingleFunctionConfiguration2.class)) diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java index 31956e13e..9f36ef470 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/function/ScenarioTests.java @@ -46,7 +46,7 @@ public class ScenarioTests { @Test public void test2106() { try (ConfigurableApplicationContext context = new SpringApplicationBuilder(TestChannelBinderConfiguration - .getCompleteConfiguration(ConsumerConfiguration.class, ConsumerConfiguration.class)) + .getCompleteConfiguration(ConsumerConfiguration.class)) .web(WebApplicationType.NONE).run( "--spring.cloud.function.definition=consume;echo", "--spring.cloud.stream.bindings.consume-in-0.destination=input", @@ -82,17 +82,17 @@ public class ScenarioTests { } @Test - public void testComposingSupplierWuthTypelessMessageFunction() { + public void testComposingSupplierWuthTypelessMessageFunction() throws Exception { try (ConfigurableApplicationContext context = new SpringApplicationBuilder( - TestChannelBinderConfiguration.getCompleteConfiguration(TestConfiguration.class)) + TestChannelBinderConfiguration.getCompleteConfiguration(SupplierConfiguration.class)) .web(WebApplicationType.NONE) .run("--spring.jmx.enabled=false", "--spring.cloud.function.definition=messageSupplier|messageFunction")) { OutputDestination output = context.getBean(OutputDestination.class); - assertThat(output.receive(1000)).isNotNull(); - assertThat(output.receive(1100)).isNotNull(); - assertThat(output.receive(1200)).isNotNull(); + assertThat(output.receive(1000, "messageSuppliermessageFunction-out-0")).isNotNull(); + assertThat(output.receive(1200, "messageSuppliermessageFunction-out-0")).isNotNull(); + assertThat(output.receive(1300, "messageSuppliermessageFunction-out-0")).isNotNull(); } } @@ -134,23 +134,29 @@ public class ScenarioTests { @EnableAutoConfiguration @Configuration public static class TestConfiguration { + @SuppressWarnings("unchecked") + @Bean + public Function genericTypeFunction() { + return v -> { + return (O) ("hello_" + new String((byte[]) v)); + }; + } + } + + @EnableAutoConfiguration + @Configuration + public static class SupplierConfiguration { @Bean public Supplier> messageSupplier() { return () -> new GenericMessage<>("10/27/20 07:20:01"); } @Bean public Function, Message> messageFunction() { - return message -> message; - } - - @SuppressWarnings("unchecked") - @Bean - public Function genericTypeFunction() { - return v -> { - System.out.println(v); - return (O) ("hello_" + v); + return message -> { + return message; }; } + } @EnableAutoConfiguration