From cd7ef9eefb3e3227914695ae9532564c09b82bb0 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Fri, 4 Jan 2019 07:52:40 +0100 Subject: [PATCH] GH-1575 Fixed MessageHeaders regression Added additional tests for both MessageHeaders and Map/Collections Fixed ByteArrayMessageConverter to ensure it handles Object if CT is octet-stream and added tests for it Resolves #1575 --- .../config/SmartPayloadArgumentResolver.java | 2 + .../CompositeMessageConverterFactory.java | 10 +- .../binder/tck/ContentTypeTckTests.java | 126 ++++++++++++++++++ 3 files changed, 137 insertions(+), 1 deletion(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java index 35a56f787..47739eba7 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/SmartPayloadArgumentResolver.java @@ -19,6 +19,7 @@ package org.springframework.cloud.stream.config; import org.springframework.core.MethodParameter; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.MessageConversionException; import org.springframework.messaging.converter.MessageConverter; import org.springframework.messaging.converter.SmartMessageConverter; @@ -64,6 +65,7 @@ class SmartPayloadArgumentResolver extends PayloadArgumentResolver { @Override public boolean supportsParameter(MethodParameter parameter) { return (!Message.class.isAssignableFrom(parameter.getParameterType()) + && !MessageHeaders.class.isAssignableFrom(parameter.getParameterType()) && !parameter.hasParameterAnnotation(Header.class) && !parameter.hasParameterAnnotation(Headers.class)); } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/CompositeMessageConverterFactory.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/CompositeMessageConverterFactory.java index 53fd3c29e..738943a1e 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/CompositeMessageConverterFactory.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/CompositeMessageConverterFactory.java @@ -82,7 +82,15 @@ public class CompositeMessageConverterFactory { applicationJsonConverter.setStrictContentTypeMatch(true); this.converters.add(applicationJsonConverter); this.converters.add(new TupleJsonMessageConverter(this.objectMapper)); - this.converters.add(new ByteArrayMessageConverter()); + this.converters.add(new ByteArrayMessageConverter() { + @Override + protected boolean supports(Class clazz) { + if (!super.supports(clazz)) { + return (Object.class == clazz); + } + return true; + } + }); this.converters.add(new ObjectStringMessageConverter()); // Deprecated converters 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 7e066fbed..c3e0e264c 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 @@ -60,6 +60,7 @@ import org.springframework.util.MimeTypeUtils; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertTrue; @@ -925,4 +926,129 @@ public class ContentTypeTckTests { } } + //====== + @Test + public void testWithMapInputParameter() { + ApplicationContext context = new SpringApplicationBuilder(MapInputConfiguration.class) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false"); + InputDestination source = context.getBean(InputDestination.class); + OutputDestination target = context.getBean(OutputDestination.class); + String jsonPayload = "{\"name\":\"oleg\"}"; + source.send(new GenericMessage<>(jsonPayload.getBytes())); + Message outputMessage = target.receive(); + assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); + } + + @EnableBinding(Processor.class) + @Import(TestChannelBinderConfiguration.class) + @EnableAutoConfiguration + public static class MapInputConfiguration { + @StreamListener(Processor.INPUT) + @SendTo(Processor.OUTPUT) + public Map echo(Map value) throws Exception { + return value; + } + } + + @Test + public void testWithMapPayloadParameter() { + ApplicationContext context = new SpringApplicationBuilder(MapInputConfiguration.class) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false"); + InputDestination source = context.getBean(InputDestination.class); + OutputDestination target = context.getBean(OutputDestination.class); + String jsonPayload = "{\"name\":\"oleg\"}"; + source.send(new GenericMessage<>(jsonPayload.getBytes())); + Message outputMessage = target.receive(); + assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); + } + + @EnableBinding(Processor.class) + @Import(TestChannelBinderConfiguration.class) + @EnableAutoConfiguration + public static class MapPayloadConfiguration { + @StreamListener(Processor.INPUT) + @SendTo(Processor.OUTPUT) + public Map echo(Message> value) throws Exception { + return value.getPayload(); + } + } + + @Test + public void testWithListInputParameter() { + ApplicationContext context = new SpringApplicationBuilder(ListInputConfiguration.class) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false"); + InputDestination source = context.getBean(InputDestination.class); + OutputDestination target = context.getBean(OutputDestination.class); + String jsonPayload = "[\"foo\",\"bar\"]"; + source.send(new GenericMessage<>(jsonPayload.getBytes())); + Message outputMessage = target.receive(); + assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); + } + + @EnableBinding(Processor.class) + @Import(TestChannelBinderConfiguration.class) + @EnableAutoConfiguration + public static class ListInputConfiguration { + @StreamListener(Processor.INPUT) + @SendTo(Processor.OUTPUT) + public List echo(List value) throws Exception { + return value; + } + } + + + @Test + public void testWithMessageHeadersInputParameter() { + ApplicationContext context = new SpringApplicationBuilder(MessageHeadersInputConfiguration.class) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false"); + InputDestination source = context.getBean(InputDestination.class); + OutputDestination target = context.getBean(OutputDestination.class); + String jsonPayload = "{\"name\":\"oleg\"}"; + source.send(new GenericMessage<>(jsonPayload.getBytes())); + Message outputMessage = target.receive(); + assertNotEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); + assertTrue(outputMessage.getHeaders().containsKey(MessageHeaders.ID)); + assertTrue(outputMessage.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE)); + } + + @EnableBinding(Processor.class) + @Import(TestChannelBinderConfiguration.class) + @EnableAutoConfiguration + public static class MessageHeadersInputConfiguration { + @StreamListener(Processor.INPUT) + @SendTo(Processor.OUTPUT) + public Map echo(MessageHeaders value) throws Exception { + return value; + } + } + + @Test + public void testWithTypelessInputParameterAndOctetStream() { + ApplicationContext context = new SpringApplicationBuilder(TypelessPayloadConfiguration.class) + .web(WebApplicationType.NONE) + .run("--spring.jmx.enabled=false"); + InputDestination source = context.getBean(InputDestination.class); + OutputDestination target = context.getBean(OutputDestination.class); + String jsonPayload = "[\"foo\",\"bar\"]"; + source.send(MessageBuilder.withPayload(jsonPayload.getBytes()) + .setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_OCTET_STREAM).build()); + Message outputMessage = target.receive(); + assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); + } + + @EnableBinding(Processor.class) + @Import(TestChannelBinderConfiguration.class) + @EnableAutoConfiguration + public static class TypelessPayloadConfiguration { + @StreamListener(Processor.INPUT) + @SendTo(Processor.OUTPUT) + public Object echo(Object value) throws Exception { + System.out.println(value); + return value; + } + } }