From 07b9adb0dc53e63ba6bdfd29af3ab81ec13db5b9 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Wed, 12 Dec 2018 16:44:41 +0100 Subject: [PATCH] GH-1554 Introduced flag to disable 1.3 content type propagation Introduced `spring.cloud.stream.propagateOriginalContentType` boolean property on BindingServiceProperties Resolves #1554 --- .../config/TextPlainConversionTest.java | 4 +- .../config/contentType/ContentTypeTests.java | 2 - .../binding/MessageConverterConfigurer.java | 51 ++++++++++++++----- .../config/BindingServiceProperties.java | 22 ++++++++ .../binder/tck/ContentTypeTckTests.java | 15 +++--- 5 files changed, 68 insertions(+), 26 deletions(-) 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 e9438646d..6f607c06d 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"); + assertThat(received.getPayload()).isEqualTo("Bar".getBytes()); } @Test diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/contentType/ContentTypeTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/contentType/ContentTypeTests.java index 874c83101..37725a42a 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/contentType/ContentTypeTests.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/contentType/ContentTypeTests.java @@ -148,8 +148,6 @@ public class ContentTypeTests { .build()); Message message = (Message) collector .forChannel(source.output()).poll(1, TimeUnit.SECONDS); - assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE, MimeType.class) - .includes(MimeTypeUtils.IMAGE_JPEG)); assertThat(message.getPayload()).isEqualTo(data); } } 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 41bcab1c7..63745d084 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,6 +45,7 @@ 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; @@ -307,35 +308,57 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig this.messageConverter = messageConverter; } + @Override public Message doPreSend(Message message, MessageChannel channel) { + @SuppressWarnings("deprecation") + boolean propagateOriginalContentType = + MessageConverterConfigurer.this.bindingServiceProperties.isPropagateOriginalContentType(); + + boolean contentTypeHeaderSet = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE); + // ===== 1.3 backward compatibility code part-1 === - String oct = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE) ? message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString() : null; - String ct = oct; - if (message.getPayload() instanceof String) { - ct = JavaClassMimeTypeUtils.mimeTypeFromObject(message.getPayload(), ObjectUtils.nullSafeToString(oct)).toString(); + String ct = null; + String oct = null; + if (propagateOriginalContentType) { + oct = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE) ? message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString() : null; + ct = message.getPayload() instanceof String + ? ct = JavaClassMimeTypeUtils.mimeTypeFromObject(message.getPayload(), ObjectUtils.nullSafeToString(oct)).toString() + : oct; } // ===== END 1.3 backward compatibility code part-1 === - 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); + + MutableMessageHeaders headers = new MutableMessageHeaders(message.getHeaders()); + if (!headers.containsKey(MessageHeaders.CONTENT_TYPE)) { + headers.put(MessageHeaders.CONTENT_TYPE, this.mimeType); } @SuppressWarnings("unchecked") Message outboundMessage = message.getPayload() instanceof byte[] - ? (Message)message : (Message) this.messageConverter.toMessage(message.getPayload(), message.getHeaders()); + ? (Message)message : (Message) this.messageConverter.toMessage(message.getPayload(), headers); if (outboundMessage == null) { throw new IllegalStateException("Failed to convert message: '" + message + "' to outbound message."); } /// ===== 1.3 backward compatibility code part-2 === - if (ct != null && !ct.equals(oct) && oct != null) { - @SuppressWarnings("unchecked") - Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders()); - headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf(ct)); - headersMap.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, MimeType.valueOf(oct)); + if (propagateOriginalContentType) { + if (ct != null && !ct.equals(oct) && oct != null) { + @SuppressWarnings("unchecked") + Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders()); + headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf(ct)); + headersMap.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, MimeType.valueOf(oct)); + } + } + else { + if (!contentTypeHeaderSet) { + @SuppressWarnings("unchecked") + Map headersMap = (Map) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.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 9ec1dbd90..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 @@ -56,6 +56,18 @@ public class BindingServiceProperties implements ApplicationContextAware, Initia private static final int DEFAULT_BINDING_RETRY_INTERVAL = 30; + /** + * Setting it to true ensures that the original content-type of the message is propagated + * to the outgoing message as `originalContentType` header. + * + * This deprecated feature primarily exists for backward compatibility + * and will not be supported in future versions. + * + * Default: true + */ + @Deprecated + private boolean propagateOriginalContentType = true; + /** * The instance id of the application: a number from 0 to instanceCount-1. * Used for partitioning and with Kafka. @@ -284,4 +296,14 @@ public class BindingServiceProperties implements ApplicationContextAware, Initia this.bindings.put(binding, bindingPropertiesTarget); } + @Deprecated + public boolean isPropagateOriginalContentType() { + return propagateOriginalContentType; + } + + @Deprecated + public void setPropagateOriginalContentType(boolean propagateOriginalContentType) { + this.propagateOriginalContentType = propagateOriginalContentType; + } + } 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 447da480c..22e79ad9e 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,6 +61,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.assertNotNull; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; /** @@ -188,7 +189,7 @@ public class ContentTypeTckTests { String jsonPayload = "{\"name\":\"oleg\"}"; source.send(new GenericMessage<>(jsonPayload.getBytes())); Message outputMessage = target.receive(); - assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); + assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -202,7 +203,7 @@ public class ContentTypeTckTests { String jsonPayload = "{\"name\":\"oleg\"}"; source.send(new GenericMessage<>(jsonPayload.getBytes())); Message outputMessage = target.receive(); - assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); + assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -256,11 +257,10 @@ public class ContentTypeTckTests { InputDestination source = context.getBean(InputDestination.class); OutputDestination target = context.getBean(OutputDestination.class); String jsonPayload = "{\"name\":\"oleg\"}"; - //source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, MimeType.valueOf("text/*")).build()); source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader("contentType", new MimeType("text", "plain")).build()); Message outputMessage = target.receive(); - //assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); + assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -272,11 +272,10 @@ public class ContentTypeTckTests { InputDestination source = context.getBean(InputDestination.class); OutputDestination target = context.getBean(OutputDestination.class); String jsonPayload = "{\"name\":\"oleg\"}"; - //source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, MimeType.valueOf("text/*")).build()); source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader("contentType", new MimeType("text")).build()); Message outputMessage = target.receive(); - //assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); + assertEquals("text/*", outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()); assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -334,7 +333,7 @@ public class ContentTypeTckTests { String jsonPayload = "{\"name\":\"oleg\"}"; source.send(new GenericMessage<>(jsonPayload.getBytes())); Message outputMessage = target.receive(); - assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); + assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); } @@ -348,7 +347,7 @@ public class ContentTypeTckTests { String jsonPayload = "{\"name\":\"oleg\"}"; source.send(new GenericMessage<>(jsonPayload.getBytes())); Message outputMessage = target.receive(); - assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); + assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE)); assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8)); }