diff --git a/docs/src/main/asciidoc/preface.adoc b/docs/src/main/asciidoc/preface.adoc index 08edc8d80..d026cd6be 100644 --- a/docs/src/main/asciidoc/preface.adoc +++ b/docs/src/main/asciidoc/preface.adoc @@ -174,3 +174,4 @@ TBD - Reactive module in favor of native support via spring-cloud-function. [Details to follow] - Test support module with MessageCollector [Details to follow] - @StreamMessageConverter [Details to follow] +- Original content type - removed diff --git a/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java b/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java index 81259dfb4..9f31ec9f5 100644 --- a/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java +++ b/spring-cloud-stream-binder-test/src/main/java/org/springframework/cloud/stream/binder/AbstractBinderTests.java @@ -208,8 +208,6 @@ public abstract class AbstractBinderTests message) throws MessagingException { - assertThat(message.getPayload()).isInstanceOf(byte[].class); - assertThat(((byte[]) message.getPayload())).isEqualTo( - "{\"message\":\"Hi\"}".getBytes(StandardCharsets.UTF_8)); - assertThat( - message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()) - .isEqualTo("application/json"); - latch.countDown(); - } - }; - this.testSink.input().subscribe(messageHandler); - this.testSink.input().send(MessageBuilder - .withPayload("{\"message\":\"Hi\"}".getBytes()) - .setHeader(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, "application/json") - .build()); - assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); - this.testSink.input().unsubscribe(messageHandler); - } - - @EnableBinding(Sink.class) - @EnableAutoConfiguration - public static class LegacyTestSink { - - } - -} diff --git a/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/binder/TestSupportBinder.java b/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/binder/TestSupportBinder.java index f0479c384..a251a9d40 100644 --- a/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/binder/TestSupportBinder.java +++ b/spring-cloud-stream-test-support/src/main/java/org/springframework/cloud/stream/test/binder/TestSupportBinder.java @@ -25,7 +25,6 @@ import java.util.concurrent.ConcurrentMap; import java.util.concurrent.LinkedBlockingDeque; import org.springframework.cloud.stream.binder.Binder; -import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.binder.Binding; import org.springframework.cloud.stream.binder.ConsumerProperties; import org.springframework.cloud.stream.binder.ProducerProperties; @@ -193,12 +192,7 @@ public class TestSupportBinder public Message preSend(Message message, MessageChannel channel) { Class targetClass = null; MessageConverter converter = null; - MimeType contentType = message.getHeaders() - .containsKey(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE) - ? MimeType.valueOf(message.getHeaders() - .get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE) - .toString()) - : MimeType.valueOf(this.contentTypeResolver + MimeType contentType = MimeType.valueOf(this.contentTypeResolver .resolve(message.getHeaders()).toString()); if (contentType != null) { @@ -225,8 +219,7 @@ public class TestSupportBinder catch (Exception e) { throw new IllegalStateException( "Failed to determine class name for contentType: " - + message.getHeaders().get( - BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE), + + message.getHeaders(), e); } } @@ -256,7 +249,7 @@ public class TestSupportBinder message = MessageBuilder.withPayload(payload) .copyHeaders(message.getHeaders()) .setHeader(MessageHeaders.CONTENT_TYPE, contentType) - .removeHeader(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE).build(); + .build(); return message; } 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 7d6e19afa..1e0f2203b 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 @@ -1099,12 +1099,6 @@ public abstract class AbstractMessageChannelBinder doPreSend(Message message, MessageChannel channel) { - // If handler is a function, FunctionInvoker will already perform message - // conversion. - // In fact in the future we should consider propagating knowledge of the - // default content type - // to MessageConverters instead of interceptors if (message.getPayload() instanceof byte[] && message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE)) { return message; } - // ===== 1.3 backward compatibility code part-1 === String oct = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE) ? message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString() : null; @@ -329,7 +323,6 @@ public class MessageConverterConfigurer ? 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") @@ -348,17 +341,13 @@ public class MessageConverterConfigurer + "' 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, outboundMessage.getHeaders()); headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf(ct)); - headersMap.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, - MimeType.valueOf(oct)); } - // ===== END 1.3 backward compatibility code part-2 === return outboundMessage; } 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 883d34d78..ede576b16 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 @@ -94,26 +94,6 @@ public class ContentTypeTckTests { assertThat(new String(outputMessage.getPayload())).isEqualTo("oleg"); } - @Test - // emulates 1.3 behavior - public void stringToMapMessageStreamListenerOriginalContentType() { - ApplicationContext context = new SpringApplicationBuilder( - StringToMapMessageStreamListener.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\"}"; - - Message message = MessageBuilder.withPayload(jsonPayload.getBytes()) - .setHeader(MessageHeaders.CONTENT_TYPE, "text/plain") - .setHeader("originalContentType", "application/json;charset=UTF-8") - .build(); - - source.send(message); - Message outputMessage = target.receive(); - assertThat(new String(outputMessage.getPayload())).isEqualTo("oleg"); - } - @Test public void withInternalPipeline() { ApplicationContext context = new SpringApplicationBuilder(InternalPipeLine.class) 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 7c0693500..561976a41 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 @@ -21,7 +21,6 @@ import java.util.Collections; import org.junit.Ignore; import org.junit.Test; -import org.springframework.cloud.stream.binder.BinderHeaders; import org.springframework.cloud.stream.config.BindingProperties; import org.springframework.cloud.stream.config.BindingServiceProperties; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; @@ -32,7 +31,6 @@ 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; @@ -105,29 +103,6 @@ public class MessageConverterConfigurerTests { } } - @Test - public void testConfigureInputChannelWithLegacyContentType() { - BindingServiceProperties props = new BindingServiceProperties(); - BindingProperties bindingProps = new BindingProperties(); - bindingProps.setContentType("foo/bar"); - props.setBindings(Collections.singletonMap("foo", bindingProps)); - CompositeMessageConverterFactory converterFactory = new CompositeMessageConverterFactory( - Collections.emptyList(), null); - MessageConverterConfigurer configurer = new MessageConverterConfigurer(props, - converterFactory.getMessageConverterForAllRegistered()); - QueueChannel in = new QueueChannel(); - configurer.configureInputChannel(in, "foo"); - Foo foo = new Foo(); - 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); - assertThat(received.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()) - .isEqualTo("application/json"); - } - public static class Foo { private String bar = "bar";