From 1a1ebe7b7830663d6c45023bf39559e0cab0091d Mon Sep 17 00:00:00 2001 From: Ilayaperumal Gopinathan Date: Mon, 28 Sep 2015 18:00:01 -0700 Subject: [PATCH] type conversion: address review comments Fix supports check on FromMessageConverter: make sure the targetType is on the LHS Fix PojoToStringMessageConverter: handle message payload tuple separately so that the tupleToString converter is set --- .../integration/MapToTupleTransformer.java | 13 ++--- .../TupleKryoRegistrar.java | 12 ++-- .../tuple/integration/package-info.java | 2 +- .../tuple/kryo/DefaultTupleSerializer.java | 9 ++- .../stream/binding/ChannelBindingService.java | 26 +++++++-- .../AbstractFromMessageConverter.java | 31 +++------- .../ByteArrayToStringMessageConverter.java | 11 ++-- .../JavaToSerializedMessageConverter.java | 2 +- .../converter/JsonToPojoMessageConverter.java | 2 +- .../JsonToTupleMessageConverter.java | 2 +- .../converter/MessageConverterUtils.java | 10 +--- .../converter/PojoToJsonMessageConverter.java | 2 +- .../PojoToStringMessageConverter.java | 16 +++++- .../SerializedToJavaMessageConverter.java | 2 +- .../converter/StrictContentTypeResolver.java | 37 ------------ .../StringConvertingContentTypeResolver.java | 57 ------------------- .../StringToByteArrayMessageConverter.java | 8 +-- .../TupleToJsonMessageConverter.java | 2 +- 18 files changed, 74 insertions(+), 170 deletions(-) rename spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/{kryo => integration}/TupleKryoRegistrar.java (83%) delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StrictContentTypeResolver.java delete mode 100644 spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringConvertingContentTypeResolver.java diff --git a/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/MapToTupleTransformer.java b/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/MapToTupleTransformer.java index daec0dfc3..ba3d014b8 100644 --- a/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/MapToTupleTransformer.java +++ b/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/MapToTupleTransformer.java @@ -1,5 +1,5 @@ /* - * Copyright 2013 the original author or authors. + * Copyright 2013-2015 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -20,14 +20,15 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; -import org.springframework.integration.transformer.AbstractPayloadTransformer; import org.springframework.cloud.stream.tuple.Tuple; import org.springframework.cloud.stream.tuple.TupleBuilder; +import org.springframework.integration.transformer.AbstractPayloadTransformer; /** * Converts from a Map to the Tuple data structure. * * @author Mark Pollack + * @author Ilayaperumal Gopinathan */ public class MapToTupleTransformer extends AbstractPayloadTransformer, Tuple> { @@ -36,11 +37,9 @@ public class MapToTupleTransformer extends AbstractPayloadTransformer newNames = new ArrayList(); List newValues = new ArrayList(); - for (Object name : map.keySet()) { - newNames.add(name.toString()); - } - for (Object value : map.values()) { - newValues.add(value); + for (Map.Entry entry: map.entrySet()) { + newNames.add(entry.getKey().toString()); + newValues.add(entry.getValue()); } return TupleBuilder.tuple().ofNamesAndValues(newNames, newValues); diff --git a/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/kryo/TupleKryoRegistrar.java b/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/TupleKryoRegistrar.java similarity index 83% rename from spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/kryo/TupleKryoRegistrar.java rename to spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/TupleKryoRegistrar.java index 43c2469e8..f9268c598 100644 --- a/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/kryo/TupleKryoRegistrar.java +++ b/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/TupleKryoRegistrar.java @@ -13,11 +13,12 @@ * limitations under the License. */ -package org.springframework.cloud.stream.tuple.kryo; +package org.springframework.cloud.stream.tuple.integration; import java.util.ArrayList; import java.util.List; +import org.springframework.cloud.stream.tuple.kryo.DefaultTupleSerializer; import org.springframework.integration.codec.kryo.AbstractKryoRegistrar; import org.springframework.integration.codec.kryo.KryoRegistrar; import org.springframework.cloud.stream.tuple.DefaultTuple; @@ -26,16 +27,15 @@ import com.esotericsoftware.kryo.Registration; import com.esotericsoftware.kryo.serializers.CollectionSerializer; /** - * A {@link KryoRegistrar} - * used to register a Tuple serializer. + * A {@link KryoRegistrar} used to register a Tuple serializer. + * * @author David Turanski - * @since 1.2 */ public class TupleKryoRegistrar extends AbstractKryoRegistrar { - private final static int TUPLE_REGISTRATION_ID = 41; + private final static int TUPLE_REGISTRATION_ID = 43; - private final static int ARRAY_LIST_REGISTRATION_ID = 42; + private final static int ARRAY_LIST_REGISTRATION_ID = 44; private final DefaultTupleSerializer defaultTupleSerializer = new DefaultTupleSerializer(); diff --git a/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/package-info.java b/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/package-info.java index 97950fb5b..f2c407d50 100644 --- a/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/package-info.java +++ b/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/integration/package-info.java @@ -1,5 +1,5 @@ /** - * Contains classes that supports tuple integration such as tuple transformers etc., + * Contains classes that support tuple integration such as tuple transformers. */ package org.springframework.cloud.stream.tuple.integration; diff --git a/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/kryo/DefaultTupleSerializer.java b/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/kryo/DefaultTupleSerializer.java index dfdcbb6ba..65dca4bed 100644 --- a/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/kryo/DefaultTupleSerializer.java +++ b/spring-cloud-stream-tuple/src/main/java/org/springframework/cloud/stream/tuple/kryo/DefaultTupleSerializer.java @@ -18,17 +18,16 @@ package org.springframework.cloud.stream.tuple.kryo; import java.util.ArrayList; import java.util.List; +import org.springframework.cloud.stream.tuple.Tuple; +import org.springframework.cloud.stream.tuple.TupleBuilder; + import com.esotericsoftware.kryo.Kryo; import com.esotericsoftware.kryo.Serializer; import com.esotericsoftware.kryo.io.Input; import com.esotericsoftware.kryo.io.Output; -import org.springframework.cloud.stream.tuple.Tuple; -import org.springframework.cloud.stream.tuple.TupleBuilder; - /** - * Deserializes Tuples by writing the field names and then the values as class/object pairs - * followed by the tuple Id and timestamp. + * Serializes Tuples by writing the field names and then the values as class/object pairs. * * @author David Turanski */ diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/ChannelBindingService.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/ChannelBindingService.java index 35fc8416e..897789941 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/ChannelBindingService.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/ChannelBindingService.java @@ -19,6 +19,8 @@ package org.springframework.cloud.stream.binding; import java.util.HashSet; import java.util.Set; +import org.springframework.aop.framework.Advised; +import org.springframework.aop.support.AopUtils; import org.springframework.beans.factory.InitializingBean; import org.springframework.cloud.stream.binder.Binder; import org.springframework.cloud.stream.config.BindingProperties; @@ -124,11 +126,17 @@ public class ChannelBindingService implements InitializingBean { /** * Setup data-type and message converters for the given message channel. * - * @param messageChannel message channel to set the data-type and message converters + * @param channel message channel to set the data-type and message converters * @param channelName the channel name */ - public void configureMessageConverters(MessageChannel messageChannel, String channelName) { - Assert.isAssignable(AbstractMessageChannel.class, messageChannel.getClass()); + public void configureMessageConverters(Object channel, String channelName) { + AbstractMessageChannel messageChannel = null; + try { + messageChannel = getMessageChannel(channel); + } + catch (Exception e) { + throw new IllegalStateException("Could not get the message channel to configure message converters" + e); + } BindingProperties bindingProperties = channelBindingServiceProperties.getBindings().get(channelName); if (bindingProperties != null) { String contentType = bindingProperties.getContentType(); @@ -137,9 +145,17 @@ public class ChannelBindingService implements InitializingBean { MessageConverter messageConverter = messageConverterFactory.newInstance(mimeType); Class dataType = MessageConverterUtils.getJavaTypeForContentType(mimeType, Thread.currentThread().getContextClassLoader()); - ((AbstractMessageChannel)messageChannel).setDatatypes(dataType); - ((AbstractMessageChannel)messageChannel).setMessageConverter(messageConverter); + messageChannel.setDatatypes(dataType); + messageChannel.setMessageConverter(messageConverter); } } } + + private AbstractMessageChannel getMessageChannel(Object channel) throws Exception { + if (AopUtils.isJdkDynamicProxy(channel)) { + return (AbstractMessageChannel) (((Advised) channel).getTargetSource().getTarget()); + } + Assert.isAssignable(AbstractMessageChannel.class, channel.getClass()); + return (AbstractMessageChannel) channel; + } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/AbstractFromMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/AbstractFromMessageConverter.java index ea8c7e768..0c9cbdcb3 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/AbstractFromMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/AbstractFromMessageConverter.java @@ -28,7 +28,6 @@ import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.AbstractMessageConverter; -import org.springframework.messaging.converter.ContentTypeResolver; import org.springframework.util.Assert; import org.springframework.util.ClassUtils; import org.springframework.util.MimeType; @@ -41,6 +40,7 @@ import org.springframework.util.MimeType; * used with custom Message conversion. Only {@link #fromMessage} is supported. * * @author David Turanski + * @author Ilayaperumal Gopinathan */ public abstract class AbstractFromMessageConverter extends AbstractMessageConverter { @@ -54,11 +54,11 @@ public abstract class AbstractFromMessageConverter extends AbstractMessageConver * @param targetMimeType the required target type */ protected AbstractFromMessageConverter(MimeType targetMimeType) { - this(new ArrayList(), targetMimeType, new StrictContentTypeResolver(targetMimeType)); + this(new ArrayList(), targetMimeType); } protected AbstractFromMessageConverter(Collection targetMimeTypes) { - this(new ArrayList(), targetMimeTypes, new StringConvertingContentTypeResolver()); + this(new ArrayList(), targetMimeTypes); } /** @@ -67,11 +67,9 @@ public abstract class AbstractFromMessageConverter extends AbstractMessageConver * @param supportedSourceMimeTypes list of {@link MimeType} that may present in content-type header * @param targetMimeType the required target type */ - protected AbstractFromMessageConverter(Collection supportedSourceMimeTypes, MimeType targetMimeType, - ContentTypeResolver contentTypeResolver) { + protected AbstractFromMessageConverter(Collection supportedSourceMimeTypes, MimeType targetMimeType) { super(supportedSourceMimeTypes); Assert.notNull(targetMimeType, "'targetMimeType' cannot be null"); - setContentTypeResolver(contentTypeResolver); this.targetMimeTypes = Collections.singletonList(targetMimeType); } @@ -80,14 +78,11 @@ public abstract class AbstractFromMessageConverter extends AbstractMessageConver * * @param supportedSourceMimeTypes a list of supported content types * @param targetMimeTypes a list of supported target types - * @param contentTypeResolver the {@link ContentTypeResolver} to use */ protected AbstractFromMessageConverter(Collection supportedSourceMimeTypes, - Collection targetMimeTypes, - ContentTypeResolver contentTypeResolver) { + Collection targetMimeTypes) { super(supportedSourceMimeTypes); Assert.notNull(targetMimeTypes, "'targetMimeTypes' cannot be null"); - setContentTypeResolver(contentTypeResolver); this.targetMimeTypes = new ArrayList(targetMimeTypes); } @@ -98,8 +93,7 @@ public abstract class AbstractFromMessageConverter extends AbstractMessageConver * @param targetMimeType the required target type */ protected AbstractFromMessageConverter(MimeType supportedSourceMimeType, MimeType targetMimeType) { - this(Collections.singletonList(supportedSourceMimeType), targetMimeType, new StrictContentTypeResolver( - supportedSourceMimeType)); + this(Collections.singletonList(supportedSourceMimeType), targetMimeType); } /** @@ -109,8 +103,7 @@ public abstract class AbstractFromMessageConverter extends AbstractMessageConver * @param targetMimeTypes a list of supported target types */ protected AbstractFromMessageConverter(MimeType supportedSourceMimeType, Collection targetMimeTypes) { - this(Collections.singletonList(supportedSourceMimeType), targetMimeTypes, new StrictContentTypeResolver( - supportedSourceMimeType)); + this(Collections.singletonList(supportedSourceMimeType), targetMimeTypes); } /** @@ -139,7 +132,7 @@ public abstract class AbstractFromMessageConverter extends AbstractMessageConver private boolean supportsType(Class clazz, Class[] supportedTypes) { if (supportedTypes != null) { for (Class targetType : supportedTypes) { - if (ClassUtils.isAssignable(clazz, targetType)) { + if (ClassUtils.isAssignable(targetType, clazz)) { return true; } } @@ -163,14 +156,6 @@ public abstract class AbstractFromMessageConverter extends AbstractMessageConver return false; } - @Override - // TODO: This will likely be fixed in core Spring - public void setContentTypeResolver(ContentTypeResolver resolver) { - if (getContentTypeResolver() == null) { - super.setContentTypeResolver(resolver); - } - } - /** * Not supported by default */ diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ByteArrayToStringMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ByteArrayToStringMessageConverter.java index e1346e2a1..a152492a7 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ByteArrayToStringMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ByteArrayToStringMessageConverter.java @@ -22,7 +22,6 @@ import java.util.Arrays; import java.util.List; import org.springframework.messaging.Message; -import org.springframework.messaging.converter.ContentTypeResolver; import org.springframework.util.MimeType; import org.springframework.util.MimeTypeUtils; @@ -33,21 +32,19 @@ import org.springframework.util.MimeTypeUtils; * the content-type header if any. * * @author David Turanski + * @author Ilayaperumal Gopinathan */ public class ByteArrayToStringMessageConverter extends AbstractFromMessageConverter { - private final static ContentTypeResolver contentTypeResolver = new StringConvertingContentTypeResolver(); - private final static List targetMimeTypes = new ArrayList(); static { - targetMimeTypes.add(MessageConverterUtils.X_SPRING_STRING); targetMimeTypes.add(MessageConverterUtils.X_JAVA_OBJECT); targetMimeTypes.add(MimeTypeUtils.TEXT_PLAIN); } public ByteArrayToStringMessageConverter() { super(Arrays.asList(new MimeType[] { MimeTypeUtils.APPLICATION_OCTET_STREAM, MimeTypeUtils.TEXT_PLAIN }), - targetMimeTypes, contentTypeResolver); + targetMimeTypes); } @Override @@ -64,8 +61,8 @@ public class ByteArrayToStringMessageConverter extends AbstractFromMessageConver * Don't need to manipulate message headers. Just return payload */ @Override - public Object convertFromInternal(Message message, Class targetClass) { - MimeType mimeType = contentTypeResolver.resolve(message.getHeaders()); + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { + MimeType mimeType = getContentTypeResolver().resolve(message.getHeaders()); String converted = null; diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JavaToSerializedMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JavaToSerializedMessageConverter.java index 427da06d9..d2be552c9 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JavaToSerializedMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JavaToSerializedMessageConverter.java @@ -46,7 +46,7 @@ public class JavaToSerializedMessageConverter extends AbstractFromMessageConvert } @Override - public Object convertFromInternal(Message message, Class targetClass) { + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { ByteArrayOutputStream bos = new ByteArrayOutputStream(); try { new ObjectOutputStream(bos).writeObject(message.getPayload()); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToPojoMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToPojoMessageConverter.java index 652611dcd..129336c1d 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToPojoMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToPojoMessageConverter.java @@ -47,7 +47,7 @@ public class JsonToPojoMessageConverter extends AbstractFromMessageConverter { } @Override - public Object convertFromInternal(Message message, Class targetClass) { + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { Object result = null; try { Object payload = message.getPayload(); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToTupleMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToTupleMessageConverter.java index a9aff1ee0..c6650cba3 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToTupleMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToTupleMessageConverter.java @@ -56,7 +56,7 @@ public class JsonToTupleMessageConverter extends AbstractFromMessageConverter { } @Override - public Object convertFromInternal(Message message, Class targetClass) { + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { String source = null; if (message.getPayload() instanceof byte[]) { source = new String((byte[]) message.getPayload()); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/MessageConverterUtils.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/MessageConverterUtils.java index f2ee9588e..cfb854811 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/MessageConverterUtils.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/MessageConverterUtils.java @@ -20,10 +20,10 @@ import static org.springframework.util.MimeType.valueOf; import static org.springframework.util.MimeTypeUtils.APPLICATION_JSON; import static org.springframework.util.MimeTypeUtils.APPLICATION_OCTET_STREAM; -import org.springframework.cloud.stream.tuple.Tuple; import org.springframework.util.ClassUtils; import org.springframework.util.MimeType; import org.springframework.cloud.stream.tuple.DefaultTuple; +import org.springframework.cloud.stream.tuple.Tuple; import org.springframework.util.StringUtils; @@ -40,11 +40,6 @@ public class MessageConverterUtils { */ public static final MimeType X_SPRING_TUPLE = MimeType.valueOf("application/x-spring-tuple"); - /** - * An MimeType for specifying a String. - */ - public static final MimeType X_SPRING_STRING = MimeType.valueOf("application/x-spring-string"); - /** * A general MimeType for Java Types. */ @@ -91,9 +86,6 @@ public class MessageConverterUtils { else if (X_JAVA_SERIALIZED_OBJECT.includes(contentType)) { return byte[].class; } - else if (X_SPRING_STRING.includes(contentType)) { - return String.class; - } return null; } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToJsonMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToJsonMessageConverter.java index 89b118ce0..a92b931ea 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToJsonMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToJsonMessageConverter.java @@ -60,7 +60,7 @@ public class PojoToJsonMessageConverter extends AbstractFromMessageConverter { } @Override - public Object convertFromInternal(Message message, Class targetClass) { + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { Object result; try { if (prettyPrint) { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToStringMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToStringMessageConverter.java index 292c9000d..6b420e773 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToStringMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToStringMessageConverter.java @@ -15,6 +15,8 @@ */ package org.springframework.cloud.stream.converter; +import org.springframework.cloud.stream.tuple.Tuple; +import org.springframework.cloud.stream.tuple.TupleBuilder; import org.springframework.messaging.Message; import org.springframework.util.MimeTypeUtils; @@ -24,6 +26,7 @@ import org.springframework.util.MimeTypeUtils; * to convert a Java object to a String using toString() * * @author David Turanski + * @author Ilayaperumal Gopinathan */ public class PojoToStringMessageConverter extends AbstractFromMessageConverter { @@ -42,8 +45,17 @@ public class PojoToStringMessageConverter extends AbstractFromMessageConverter { } @Override - public Object convertFromInternal(Message message, Class targetClass) { - return buildConvertedMessage(message.getPayload().toString(), message.getHeaders(), MimeTypeUtils.TEXT_PLAIN); + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { + String payloadString = null; + if (message.getPayload() instanceof Tuple) { + TupleBuilder builder = TupleBuilder.tuple(); + builder.putAll((Tuple)message.getPayload()); + payloadString = builder.build().toString(); + } + else { + payloadString = message.getPayload().toString(); + } + return buildConvertedMessage(payloadString, message.getHeaders(), MimeTypeUtils.TEXT_PLAIN); } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/SerializedToJavaMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/SerializedToJavaMessageConverter.java index 61b94e7bf..5bb3d557d 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/SerializedToJavaMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/SerializedToJavaMessageConverter.java @@ -46,7 +46,7 @@ public class SerializedToJavaMessageConverter extends AbstractFromMessageConvert } @Override - public Object convertFromInternal(Message message, Class targetClass) { + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { ByteArrayInputStream bis = new ByteArrayInputStream((byte[]) (message.getPayload())); Object result = null; try { diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StrictContentTypeResolver.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StrictContentTypeResolver.java deleted file mode 100644 index 0155008f6..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StrictContentTypeResolver.java +++ /dev/null @@ -1,37 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.stream.converter; - -import org.springframework.util.MimeType; - - -/** - * A {@link StringConvertingContentTypeResolver} that requires a the content-type to be present. - * - * @author David Turanski - */ -// TODO: This will likely be pushed to core Spring -public class StrictContentTypeResolver extends StringConvertingContentTypeResolver { - - /** - * @param defaultMimeType the required {@link MimeType} - */ - public StrictContentTypeResolver(MimeType defaultMimeType) { - super(); - setDefaultMimeType(defaultMimeType); - } -} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringConvertingContentTypeResolver.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringConvertingContentTypeResolver.java deleted file mode 100644 index 949df0a45..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringConvertingContentTypeResolver.java +++ /dev/null @@ -1,57 +0,0 @@ -/* - * Copyright 2015 the original author or authors. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package org.springframework.cloud.stream.converter; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; - -import org.springframework.messaging.MessageHeaders; -import org.springframework.messaging.converter.DefaultContentTypeResolver; -import org.springframework.util.MimeType; - -/** - * A {@link DefaultContentTypeResolver} that can parse String values. - * - * @author David Turanski - */ -public class StringConvertingContentTypeResolver extends DefaultContentTypeResolver { - - private ConcurrentMap mimeTypeCache = new ConcurrentHashMap<>(); - - @Override - public MimeType resolve(MessageHeaders headers) { - return resolve((Map) headers); - } - - public MimeType resolve(Map headers) { - Object value = headers.get(MessageHeaders.CONTENT_TYPE); - if (value instanceof MimeType) { - return (MimeType) value; - } - else if (value instanceof String) { - MimeType mimeType = mimeTypeCache.get(value); - if (mimeType == null) { - String valueAsString = (String) value; - mimeType = MimeType.valueOf(valueAsString); - mimeTypeCache.put(valueAsString,mimeType); - } - return mimeType; - } - return getDefaultMimeType(); - } -} diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringToByteArrayMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringToByteArrayMessageConverter.java index 17440a9c9..cf9bd9a59 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringToByteArrayMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringToByteArrayMessageConverter.java @@ -21,7 +21,6 @@ import java.util.ArrayList; import java.util.List; import org.springframework.messaging.Message; -import org.springframework.messaging.converter.ContentTypeResolver; import org.springframework.util.MimeType; import org.springframework.util.MimeTypeUtils; @@ -32,11 +31,10 @@ import org.springframework.util.MimeTypeUtils; * the content-type header if any. * * @author David Turanski + * @author Ilayaperumal Gopinathan */ public class StringToByteArrayMessageConverter extends AbstractFromMessageConverter { - private final static ContentTypeResolver contentTypeResolver = new StringConvertingContentTypeResolver(); - private final static List targetMimeTypes = new ArrayList(); static { targetMimeTypes.add(MimeTypeUtils.APPLICATION_OCTET_STREAM); @@ -60,8 +58,8 @@ public class StringToByteArrayMessageConverter extends AbstractFromMessageConver * Don't need to manipulate message headers. Just return the payload */ @Override - public Object convertFromInternal(Message message, Class targetClass) { - MimeType mimeType = contentTypeResolver.resolve(message.getHeaders()); + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { + MimeType mimeType = getContentTypeResolver().resolve(message.getHeaders()); byte[] converted = null; if (mimeType == null || mimeType.getParameter("Charset") == null) { converted = ((String) message.getPayload()).getBytes(); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/TupleToJsonMessageConverter.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/TupleToJsonMessageConverter.java index 9fce50f66..a6b9c1835 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/TupleToJsonMessageConverter.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/TupleToJsonMessageConverter.java @@ -57,7 +57,7 @@ public class TupleToJsonMessageConverter extends AbstractFromMessageConverter { } @Override - public Object convertFromInternal(Message message, Class targetClass) { + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { Tuple t = (Tuple) message.getPayload(); String json; if (prettyPrint) {