From 6c259be62bb6d198612fd1b52fff75e3cb0afa62 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Wed, 25 Oct 2017 18:10:11 -0400 Subject: [PATCH] Legacy cotent-type changes for the consumer - Remove binding config property `legacyContentTypeHeaderEnabled` that was introduced to enable legacy content type handling - Enable LegacyContentTypeInterceptor in 2.0 always, but bypass any legacy content type handling if the message received is from a 2.0 producer by checking on the version header - Introdce a new BinderHeader property for version - Fix the LegacyContentType related tests - Remove the check for originalContentType in ReceivingHandler when the payload received is a byte[] - Remove unnecessary deserializePayloadIfNecessary calls in ReceivingHandler - Remove deprecated deserializePayload methods in AbstractBinder and MessageSerializationUtils Partly fixes #1106 Fixes #1110 --- .../stream/config/LegacyContentTypeTests.java | 7 +++-- ...legacy-sink-channel-configurers.properties | 4 --- .../cloud/stream/binder/AbstractBinder.java | 10 ------- .../binder/AbstractMessageChannelBinder.java | 21 ++++---------- .../cloud/stream/binder/BinderHeaders.java | 8 ++++++ .../binder/MessageSerializationUtils.java | 28 ------------------- .../binding/MessageConverterConfigurer.java | 11 ++++---- .../stream/config/BindingProperties.java | 10 ------- .../MessageConverterConfigurerTests.java | 9 ++++-- 9 files changed, 30 insertions(+), 78 deletions(-) delete mode 100644 spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/channel/legacy-sink-channel-configurers.properties diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/LegacyContentTypeTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/LegacyContentTypeTests.java index 0209ec9ae..37e87a52d 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/LegacyContentTypeTests.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/LegacyContentTypeTests.java @@ -28,7 +28,6 @@ import org.springframework.boot.test.context.SpringBootTest; import org.springframework.cloud.stream.annotation.EnableBinding; import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.messaging.Sink; -import org.springframework.context.annotation.PropertySource; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHandler; @@ -61,14 +60,16 @@ public class LegacyContentTypeTests { } }; testSink.input().subscribe(messageHandler); - testSink.input().send(MessageBuilder.withPayload("{\"message\":\"Hi\"}".getBytes()).setHeader(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, "application/json").build()); + testSink.input().send(MessageBuilder.withPayload("{\"message\":\"Hi\"}".getBytes()) + .setHeader(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, "application/json") + .setHeader(BinderHeaders.SCST_VERSION, "1.x") + .build()); assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); testSink.input().unsubscribe(messageHandler); } @EnableBinding(Sink.class) @EnableAutoConfiguration - @PropertySource("classpath:/org/springframework/cloud/stream/config/channel/legacy-sink-channel-configurers.properties") public static class LegacyTestSink { } diff --git a/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/channel/legacy-sink-channel-configurers.properties b/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/channel/legacy-sink-channel-configurers.properties deleted file mode 100644 index 78c379dd0..000000000 --- a/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/channel/legacy-sink-channel-configurers.properties +++ /dev/null @@ -1,4 +0,0 @@ -spring.cloud.stream.bindings.input.destination=configure1 -spring.cloud.stream.bindings.input.legacyContentTypeHeaderEnabled=true -spring.cloud.stream.bindings.input.contentType=application/x-spring-tuple - diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java index c5540c355..7df341865 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractBinder.java @@ -151,16 +151,6 @@ public abstract class AbstractBinder message) { - return MessageSerializationUtils.deserializePayload(new MessageValues(message), this.contentTypeResolver); - } - - @Deprecated - protected final MessageValues deserializePayloadIfNecessary(MessageValues messageValues) { - return MessageSerializationUtils.deserializePayload(messageValues, this.contentTypeResolver); - } - @Deprecated protected String buildPartitionRoutingExpression(String expressionRoot) { return "'" + expressionRoot + "-' + headers['" + BinderHeaders.PARTITION_HEADER + "']"; diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java index d40e18d60..ced59706e 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binder/AbstractMessageChannelBinder.java @@ -45,13 +45,12 @@ import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.SubscribableChannel; import org.springframework.util.Assert; import org.springframework.util.MimeType; -import org.springframework.util.MimeTypeUtils; /** * {@link AbstractBinder} that serves as base class for {@link MessageChannel} binders. * Implementors must implement the following methods: * @@ -94,7 +93,7 @@ public abstract class AbstractMessageChannelBinder requestMessage) { - if (!(requestMessage.getPayload() instanceof byte[]) - && !requestMessage.getHeaders().containsKey(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)) { + if (!(requestMessage.getPayload() instanceof byte[])) { return requestMessage; } - MessageValues messageValues; if (this.extractEmbeddedHeaders && !requestMessage.getHeaders().containsKey(BinderHeaders.NATIVE_HEADERS_PRESENT) && EmbeddedHeaderUtils.mayHaveEmbeddedHeaders((byte[]) requestMessage.getPayload())) { + MessageValues messageValues; try { messageValues = EmbeddedHeaderUtils.extractHeaders((Message) requestMessage, true); @@ -554,18 +552,11 @@ public abstract class AbstractMessageChannelBinder> payloadTypeCache = new ConcurrentHashMap<>(); - /** * Serialize the message payload unless it is a byte array. * @@ -49,25 +42,4 @@ public abstract class MessageSerializationUtils { return messageValues; } - - - /** - * De-serialize the message payload if necessary. - * - * @param messageValues message with the payload to deserialize - * @param contentTypeResolver used for resolving the mime type. - * @return Deserialized Message. - */ - public static MessageValues deserializePayload(MessageValues messageValues, ContentTypeResolver contentTypeResolver) { - Object payload = messageValues.getPayload(); - MimeType contentType = contentTypeResolver.resolve(new MessageHeaders(messageValues.getHeaders())); - if (payload != null) { - messageValues.setPayload(payload); - messageValues.put(MessageHeaders.CONTENT_TYPE, contentType); - } - return messageValues; - } - - - } 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 418a4951c..0772e9380 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 @@ -121,7 +121,7 @@ public class MessageConverterConfigurer getPartitionKeyExtractorStrategy(producerProperties), getPartitionSelectorStrategy(producerProperties))); } - if (input && bindingProperties.isLegacyContentTypeHeaderEnabled()) { + if (input) { messageChannel.addInterceptor(new LegacyContentTypeHeaderInterceptor()); } // TODO: Set all interceptors in the correct order for input/output channels @@ -306,11 +306,12 @@ public class MessageConverterConfigurer @Override public Message preSend(Message message, MessageChannel channel) { - Object originalContentType = message.getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE); - if (originalContentType != null) { - return MessageConverterConfigurer.this.messageBuilderFactory + if (!message.getHeaders().containsKey(BinderHeaders.SCST_VERSION) || + !message.getHeaders().get(BinderHeaders.SCST_VERSION).equals("2.x")) { + Object originalContentType = message.getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE); + return originalContentType != null ? MessageConverterConfigurer.this.messageBuilderFactory .fromMessage(message) - .setHeader(MessageHeaders.CONTENT_TYPE, originalContentType).build(); + .setHeader(MessageHeaders.CONTENT_TYPE, originalContentType).build() : message; } return message; } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingProperties.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingProperties.java index 676de5d54..c0f880313 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingProperties.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/BindingProperties.java @@ -58,8 +58,6 @@ public class BindingProperties { private String contentType = MimeTypeUtils.APPLICATION_JSON_VALUE; - private boolean legacyContentTypeHeaderEnabled = false; - private String binder; private ConsumerProperties consumer; @@ -90,14 +88,6 @@ public class BindingProperties { this.contentType = contentType; } - public boolean isLegacyContentTypeHeaderEnabled() { - return legacyContentTypeHeaderEnabled; - } - - public void setLegacyContentTypeHeaderEnabled(boolean legacyContentTypeHeaderEnabled) { - this.legacyContentTypeHeaderEnabled = legacyContentTypeHeaderEnabled; - } - public String getBinder() { return binder; } diff --git a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/MessageConverterConfigurerTests.java b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/MessageConverterConfigurerTests.java index 386df1582..b479d02dd 100644 --- a/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/MessageConverterConfigurerTests.java +++ b/spring-cloud-stream/src/test/java/org/springframework/cloud/stream/binding/MessageConverterConfigurerTests.java @@ -31,6 +31,7 @@ import org.springframework.messaging.converter.AbstractMessageConverter; import org.springframework.messaging.converter.MessageConversionException; import org.springframework.messaging.converter.MessageConverter; import org.springframework.messaging.support.GenericMessage; +import org.springframework.messaging.support.MessageBuilder; import org.springframework.util.MimeType; import static org.assertj.core.api.Assertions.assertThat; @@ -103,7 +104,6 @@ public class MessageConverterConfigurerTests { BindingServiceProperties props = new BindingServiceProperties(); BindingProperties bindingProps = new BindingProperties(); bindingProps.setContentType("foo/bar"); - bindingProps.setLegacyContentTypeHeaderEnabled(true); props.setBindings(Collections.singletonMap("foo", bindingProps)); CompositeMessageConverterFactory converterFactory = new CompositeMessageConverterFactory( Collections.emptyList(), null); @@ -111,8 +111,11 @@ public class MessageConverterConfigurerTests { QueueChannel in = new QueueChannel(); configurer.configureInputChannel(in, "foo"); Foo foo = new Foo(); - in.send(new GenericMessage<>(foo, - Collections.singletonMap(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, "application/json"))); + in.send( + MessageBuilder.withPayload(foo) + .setHeader(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, "application/json") + .setHeader(BinderHeaders.SCST_VERSION, "1.x") + .build()); Message received = in.receive(0); assertThat(received).isNotNull(); assertThat(received.getPayload()).isEqualTo(foo);