diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/TextPlainConversionTest.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/TextPlainConversionTest.java index 6f607c06d..e9438646d 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/TextPlainConversionTest.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/TextPlainConversionTest.java @@ -68,10 +68,10 @@ public class TextPlainConversionTest { public void testByteArrayConversionOnOutput() throws Exception { testProcessor.output().send(MessageBuilder.withPayload("Bar".getBytes()).build()); @SuppressWarnings("unchecked") - Message received = (Message)((TestSupportBinder) binderFactory.getBinder(null, MessageChannel.class)) + Message received = (Message)((TestSupportBinder) binderFactory.getBinder(null, MessageChannel.class)) .messageCollector().forChannel(testProcessor.output()).poll(1, TimeUnit.SECONDS); assertThat(received).isNotNull(); - assertThat(received.getPayload()).isEqualTo("Bar".getBytes()); + assertThat(received.getPayload()).isEqualTo("Bar"); } @Test diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java index 63745d084..dd48f3030 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageConverterConfigurer.java @@ -45,7 +45,6 @@ import org.springframework.integration.channel.AbstractMessageChannel; import org.springframework.integration.expression.ExpressionUtils; import org.springframework.integration.support.MessageBuilderFactory; import org.springframework.integration.support.MutableMessageBuilderFactory; -import org.springframework.integration.support.MutableMessageHeaders; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHeaders; @@ -329,14 +328,15 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig // ===== END 1.3 backward compatibility code part-1 === - MutableMessageHeaders headers = new MutableMessageHeaders(message.getHeaders()); - if (!headers.containsKey(MessageHeaders.CONTENT_TYPE)) { - headers.put(MessageHeaders.CONTENT_TYPE, this.mimeType); + if (!message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE)) { + @SuppressWarnings("unchecked") + Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders()); + headersMap.put(MessageHeaders.CONTENT_TYPE, this.mimeType); } @SuppressWarnings("unchecked") Message outboundMessage = message.getPayload() instanceof byte[] - ? (Message)message : (Message) this.messageConverter.toMessage(message.getPayload(), headers); + ? (Message)message : (Message) this.messageConverter.toMessage(message.getPayload(), message.getHeaders()); if (outboundMessage == null) { throw new IllegalStateException("Failed to convert message: '" + message + "' to outbound message."); } @@ -345,7 +345,8 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig if (propagateOriginalContentType) { if (ct != null && !ct.equals(oct) && oct != null) { @SuppressWarnings("unchecked") - Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders()); + Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, + outboundMessage.getHeaders()); headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf(ct)); headersMap.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, MimeType.valueOf(oct)); } @@ -353,12 +354,10 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig else { if (!contentTypeHeaderSet) { @SuppressWarnings("unchecked") - Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders()); + Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, + outboundMessage.getHeaders()); headersMap.remove(MessageHeaders.CONTENT_TYPE); } - else { - System.out.println(); - } } // ===== END 1.3 backward compatibility code part-2 === return outboundMessage; diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java index 1273bc596..5a72880df 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingServiceProperties.java @@ -66,7 +66,7 @@ public class BindingServiceProperties implements ApplicationContextAware, Initia * Default: true */ @Deprecated - private boolean propagateOriginalContentType; + private boolean propagateOriginalContentType = true; /** * The instance id of the application: a number from 0 to instanceCount-1. 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 22e79ad9e..dbea12ed0 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 @@ -61,7 +61,6 @@ import org.springframework.util.MimeTypeUtils; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; -import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; /** @@ -189,7 +188,6 @@ public class ContentTypeTckTests { String jsonPayload = "{\"name\":\"oleg\"}"; source.send(new GenericMessage<>(jsonPayload.getBytes())); Message outputMessage = target.receive(); - assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -203,7 +201,6 @@ public class ContentTypeTckTests { String jsonPayload = "{\"name\":\"oleg\"}"; source.send(new GenericMessage<>(jsonPayload.getBytes())); Message outputMessage = target.receive(); - assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -275,7 +272,7 @@ public class ContentTypeTckTests { source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader("contentType", new MimeType("text")).build()); Message outputMessage = target.receive(); - assertEquals("text/*", outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()); + assertEquals("text/plain", outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()); assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -333,7 +330,6 @@ public class ContentTypeTckTests { String jsonPayload = "{\"name\":\"oleg\"}"; source.send(new GenericMessage<>(jsonPayload.getBytes())); Message outputMessage = target.receive(); - assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -341,13 +337,13 @@ public class ContentTypeTckTests { public void byteArrayToByteArrayInboundOutboundContentTypeBinding() { ApplicationContext context = new SpringApplicationBuilder(ByteArrayToByteArrayStreamListener.class) .web(WebApplicationType.NONE) - .run("--spring.cloud.stream.bindings.input.contentType=text/plain", "--spring.cloud.stream.bindings.output.contentType=text/plain", "--spring.jmx.enabled=false"); + .run("--spring.cloud.stream.bindings.input.contentType=text/plain", "--spring.cloud.stream.bindings.output.contentType=text/plain", + "--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(); - assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); }