From 81755865c8b5c172104a61261d6d940d25b7eecd Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Tue, 7 Apr 2020 08:04:45 +0200 Subject: [PATCH] GH-1941 Add ability to read content type as byte[] Modified implementation of DefaultContentTypeResolver to ensure it can recognize content type header as byte[] Resolves #1941 --- ...cationJsonMessageMarshallingConverter.java | 17 +------------ .../CompositeMessageConverterFactory.java | 24 +++++++++++++++++-- 2 files changed, 23 insertions(+), 18 deletions(-) diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java index e5f50595f..b5d3ee1d4 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ApplicationJsonMessageMarshallingConverter.java @@ -31,13 +31,11 @@ import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.cloud.function.context.catalog.FunctionTypeUtils; import org.springframework.core.MethodParameter; import org.springframework.core.ParameterizedTypeReference; -import org.springframework.integration.support.MutableMessageHeaders; import org.springframework.lang.Nullable; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.MappingJackson2MessageConverter; import org.springframework.messaging.converter.MessageConversionException; -import org.springframework.util.MimeType; /** * Variation of {@link MappingJackson2MessageConverter} to support marshalling and @@ -138,7 +136,7 @@ class ApplicationJsonMessageMarshallingConverter extends MappingJackson2MessageC else { final JavaType typeToUse = type; if (payload instanceof Collection) { - Collection collection = (Collection) ((Collection) payload).stream() + Collection collection = ((Collection) payload).stream() .map(value -> { try { if (value instanceof byte[]) { @@ -168,17 +166,4 @@ class ApplicationJsonMessageMarshallingConverter extends MappingJackson2MessageC throw new MessageConversionException("Cannot parse payload ", e); } } - - @Override - @Nullable - protected MimeType getMimeType(@Nullable MessageHeaders headers) { - Object contentType = headers.get(MessageHeaders.CONTENT_TYPE); - if (contentType instanceof byte[]) { - contentType = new String((byte[]) contentType, StandardCharsets.UTF_8); - contentType = ((String) contentType).replace("\"", ""); - headers = new MutableMessageHeaders(headers); - headers.put(MessageHeaders.CONTENT_TYPE, contentType); - } - return super.getMimeType(headers); - } } 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 7bc77af97..9a132f582 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 @@ -16,15 +16,20 @@ package org.springframework.cloud.stream.converter; +import java.lang.reflect.Field; +import java.nio.charset.StandardCharsets; import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Map; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.springframework.cloud.stream.config.BindingProperties; +import org.springframework.lang.Nullable; +import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.AbstractMessageConverter; import org.springframework.messaging.converter.ByteArrayMessageConverter; import org.springframework.messaging.converter.CompositeMessageConverter; @@ -32,6 +37,7 @@ import org.springframework.messaging.converter.DefaultContentTypeResolver; import org.springframework.messaging.converter.MessageConverter; import org.springframework.util.CollectionUtils; import org.springframework.util.MimeType; +import org.springframework.util.ReflectionUtils; /** * A factory for creating an instance of {@link CompositeMessageConverter} for a given @@ -71,7 +77,22 @@ public class CompositeMessageConverterFactory { } initDefaultConverters(); - DefaultContentTypeResolver resolver = new DefaultContentTypeResolver(); + Field headersField = ReflectionUtils.findField(MessageHeaders.class, "headers"); + headersField.setAccessible(true); + DefaultContentTypeResolver resolver = new DefaultContentTypeResolver() { + @Override + @SuppressWarnings("unchecked") + public MimeType resolve(@Nullable MessageHeaders headers) { + Object contentType = headers.get(MessageHeaders.CONTENT_TYPE); + if (contentType instanceof byte[]) { + contentType = new String((byte[]) contentType, StandardCharsets.UTF_8); + contentType = ((String) contentType).replace("\"", ""); + Map headersMap = (Map) ReflectionUtils.getField(headersField, headers); + headersMap.put(MessageHeaders.CONTENT_TYPE, contentType); + } + return super.resolve(headers); + } + }; resolver.setDefaultMimeType(BindingProperties.DEFAULT_CONTENT_TYPE); this.converters.stream().filter(mc -> mc instanceof AbstractMessageConverter) .forEach(mc -> ((AbstractMessageConverter) mc) @@ -135,5 +156,4 @@ public class CompositeMessageConverterFactory { public CompositeMessageConverter getMessageConverterForAllRegistered() { return new CompositeMessageConverter(new ArrayList<>(this.converters)); } - }