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 a62bfd0fe..8172b0193 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 @@ -239,7 +239,7 @@ public abstract class AbstractBinderTests customMessageConverters; + private List customMessageConverters; @Test public void testCustomMessageConverter() throws Exception { - assertThat(customMessageConverters).hasSize(2); - assertThat(customMessageConverters).extracting("class").contains(FooToBarConverter.class, - BarToFooConverter.class); + assertThat(customMessageConverters).hasSize(3); + assertThat(customMessageConverters).extracting("class").contains(FooConverter.class, + BarConverter.class, DefaultDatatypeChannelMessageConverter.class); testSource.output().send(MessageBuilder.withPayload(new Foo("hi")).build()); @SuppressWarnings("unchecked") Message received = (Message) ((TestSupportBinder) binderFactory.getBinder(null)) .messageCollector().forChannel(testSource.output()).poll(1, TimeUnit.SECONDS); Assert.assertThat(received, notNullValue()); - assertThat(received.getHeaders().get(MessageHeaders.CONTENT_TYPE)).isEqualTo("test/bar"); + assertThat(received.getHeaders().get(MessageHeaders.CONTENT_TYPE)).isEqualTo(MimeType.valueOf("test/foo")); } @EnableBinding(Source.class) @@ -81,39 +83,34 @@ public class CustomMessageConverterTests { public static class TestSource { @Bean - public AbstractFromMessageConverter fooConverter() { - return new FooToBarConverter(); + public MessageConverter fooConverter() { + return new FooConverter(); } @Bean - public AbstractFromMessageConverter barConverter() { - return new BarToFooConverter(); + public MessageConverter barConverter() { + return new BarConverter(); } } - public static class FooToBarConverter extends AbstractFromMessageConverter { + public static class FooConverter extends AbstractMessageConverter { - public FooToBarConverter() { - super(MimeType.valueOf("test/bar")); + public FooConverter() { + super(MimeType.valueOf("test/foo")); } @Override - protected Class[] supportedTargetTypes() { - return new Class[] {Bar.class}; + protected boolean supports(Class clazz) { + return clazz.equals(Foo.class); } @Override - protected Class[] supportedPayloadTypes() { - return new Class[] {Foo.class}; - } - - @Override - public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { + protected Object convertToInternal(Object payload, MessageHeaders headers, Object conversionHint) { Object result = null; try { - if (message.getPayload() instanceof Foo) { - Foo fooPayload = (Foo) message.getPayload(); - result = new Bar(fooPayload.test); + if (payload instanceof Foo) { + Foo fooPayload = (Foo) payload; + result = fooPayload.test.getBytes(); } } catch (Exception e) { @@ -124,29 +121,25 @@ public class CustomMessageConverterTests { } } - public static class BarToFooConverter extends AbstractFromMessageConverter { + public static class BarConverter extends AbstractMessageConverter { - public BarToFooConverter() { - super(MimeType.valueOf("test/foo")); + public BarConverter() { + super(MimeType.valueOf("test/bar")); + } + + + @Override + protected boolean supports(Class clazz) { + return clazz.equals(Bar.class); } @Override - protected Class[] supportedTargetTypes() { - return new Class[] {Foo.class}; - } - - @Override - protected Class[] supportedPayloadTypes() { - return new Class[] {Bar.class}; - } - - @Override - public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { + protected Object convertToInternal(Object payload, MessageHeaders headers, Object conversionHint) { Object result = null; try { - if (message.getPayload() instanceof Bar) { - Bar barPayload = (Bar) message.getPayload(); - result = new Foo(barPayload.testing); + if (payload instanceof Bar) { + Bar barPayload = (Bar) payload; + result = barPayload.testing.getBytes(); } } catch (Exception e) { diff --git a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/StreamListenerTests.java b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/StreamListenerTests.java index ccd350cd2..696f3d8ab 100644 --- a/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/StreamListenerTests.java +++ b/spring-cloud-stream-integration-tests/src/test/java/org/springframework/cloud/stream/config/StreamListenerTests.java @@ -42,6 +42,8 @@ import org.springframework.messaging.handler.annotation.Header; import org.springframework.messaging.handler.annotation.Headers; import org.springframework.messaging.handler.annotation.Payload; import org.springframework.messaging.handler.annotation.SendTo; +import org.springframework.util.MimeType; +import org.springframework.util.MimeTypeUtils; import static org.assertj.core.api.Assertions.assertThat; import static org.junit.Assert.fail; @@ -129,7 +131,7 @@ public class StreamListenerTests { Message message = (Message) collector.forChannel(processor.output()).poll(1, TimeUnit.SECONDS); assertThat(message).isNotNull(); assertThat(message.getPayload()).isEqualTo("{\"qux\":\"barbar" + id + "\"}"); - assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE)).isEqualTo("application/json"); + assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE, MimeType.class).includes(MimeTypeUtils.APPLICATION_JSON)); context.close(); } @@ -224,7 +226,7 @@ public class StreamListenerTests { Message message = (Message) collector.forChannel(processor.output()).poll(1, TimeUnit.SECONDS); assertThat(message).isNotNull(); assertThat(message.getPayload()).isEqualTo("{\"qux\":\"barbar" + id + "\"}"); - assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE)).isEqualTo("application/json"); + assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE, MimeType.class).includes(MimeTypeUtils.APPLICATION_JSON)); context.close(); } diff --git a/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/custom/source-channel-configurers.properties b/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/custom/source-channel-configurers.properties index 052071a95..34a8481dc 100644 --- a/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/custom/source-channel-configurers.properties +++ b/spring-cloud-stream-integration-tests/src/test/resources/org/springframework/cloud/stream/config/custom/source-channel-configurers.properties @@ -1,2 +1,2 @@ spring.cloud.stream.bindings.output.destination=configure1 -spring.cloud.stream.bindings.output.contentType=test/bar +spring.cloud.stream.bindings.output.contentType=test/foo 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 eb6a437f9..096616eab 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 @@ -31,6 +31,7 @@ import org.springframework.messaging.MessageHandler; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.SubscribableChannel; import org.springframework.util.Assert; +import org.springframework.util.MimeType; /** * {@link AbstractBinder} that serves as base class for {@link MessageChannel} @@ -119,7 +120,7 @@ public abstract class AbstractMessageChannelBinder - *
  • {@link BinderHeaders#PARTITION_HEADER} - indicates the target - * partition where the message must be sent
  • + *
  • {@link BinderHeaders#PARTITION_HEADER} - indicates the target + * partition where the message must be sent
  • * *

    * @param destination the name of the target destination @@ -210,8 +211,8 @@ public abstract class AbstractMessageChannelBinder binding = this.channelBindingService.bindProducer(channel, channelName); diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/CompositeMessageChannelConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/CompositeMessageChannelConfigurer.java index 720580427..fd235af99 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/CompositeMessageChannelConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/CompositeMessageChannelConfigurer.java @@ -33,9 +33,17 @@ public class CompositeMessageChannelConfigurer implements MessageChannelConfigur } @Override - public void configureMessageChannel(MessageChannel messageChannel, String channelName) { - for (MessageChannelConfigurer messageChannelConfigurer : messageChannelConfigurers) { - messageChannelConfigurer.configureMessageChannel(messageChannel, channelName); + public void configureInputChannel(MessageChannel messageChannel, String channelName) { + for (MessageChannelConfigurer messageChannelConfigurer : this.messageChannelConfigurers) { + messageChannelConfigurer.configureInputChannel(messageChannel, channelName); } } + + @Override + public void configureOutputChannel(MessageChannel messageChannel, String channelName) { + for (MessageChannelConfigurer messageChannelConfigurer : this.messageChannelConfigurers) { + messageChannelConfigurer.configureOutputChannel(messageChannel, channelName); + } + } + } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/DefaultBindableChannelFactory.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/DefaultBindableChannelFactory.java index 5b71eb678..8af6fd79c 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/DefaultBindableChannelFactory.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/DefaultBindableChannelFactory.java @@ -35,9 +35,17 @@ public class DefaultBindableChannelFactory implements BindableChannelFactory { } @Override - public SubscribableChannel createSubscribableChannel(String name) { + public SubscribableChannel createInputChannel(String name) { SubscribableChannel subscribableChannel = new DirectChannel(); - messageChannelConfigurer.configureMessageChannel(subscribableChannel, name); + this.messageChannelConfigurer.configureInputChannel(subscribableChannel, name); return subscribableChannel; } + + @Override + public SubscribableChannel createOutputChannel(String name) { + SubscribableChannel subscribableChannel = new DirectChannel(); + this.messageChannelConfigurer.configureOutputChannel(subscribableChannel, name); + return subscribableChannel; + } + } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageChannelConfigurer.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageChannelConfigurer.java index c696b4fef..936c6c106 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageChannelConfigurer.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/binding/MessageChannelConfigurer.java @@ -25,9 +25,16 @@ import org.springframework.messaging.MessageChannel; public interface MessageChannelConfigurer { /** - * Configure the given message channel. + * Configure the given input message channel. * @param messageChannel the message channel * @param channelName name of the message channel */ - void configureMessageChannel(MessageChannel messageChannel, String channelName); + void configureInputChannel(MessageChannel messageChannel, String channelName); + + /** + * Configure the given output message channel. + * @param messageChannel the message channel + * @param channelName name of the message channel + */ + void configureOutputChannel(MessageChannel messageChannel, String channelName); } 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 addf953bb..1bc954e05 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 @@ -17,8 +17,6 @@ package org.springframework.cloud.stream.binding; -import java.util.Collections; - import org.springframework.beans.BeansException; import org.springframework.beans.factory.BeanFactory; import org.springframework.beans.factory.BeanFactoryAware; @@ -39,8 +37,9 @@ import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.MessageHeaders; import org.springframework.messaging.converter.MessageConversionException; -import org.springframework.messaging.converter.SmartMessageConverter; +import org.springframework.messaging.converter.MessageConverter; import org.springframework.messaging.support.ChannelInterceptorAdapter; +import org.springframework.tuple.Tuple; import org.springframework.util.Assert; import org.springframework.util.MimeType; import org.springframework.util.StringUtils; @@ -81,92 +80,35 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea Assert.notNull(this.beanFactory, "Bean factory cannot be empty"); } + @Override + public void configureInputChannel(MessageChannel messageChannel, String channelName) { + configureMessageChannel(messageChannel, channelName, true); + } + + @Override + public void configureOutputChannel(MessageChannel messageChannel, String channelName) { + configureMessageChannel(messageChannel, channelName, false); + } + /** * Setup data-type and message converters for the given message channel. * @param channel message channel to set the data-type and message converters * @param channelName the channel name */ - @Override - public void configureMessageChannel(MessageChannel channel, String channelName) { + private void configureMessageChannel(MessageChannel channel, String channelName, boolean input) { Assert.isAssignable(AbstractMessageChannel.class, channel.getClass()); AbstractMessageChannel messageChannel = (AbstractMessageChannel) channel; final BindingProperties bindingProperties = this.channelBindingServiceProperties.getBindingProperties( channelName); final String contentType = bindingProperties.getContentType(); - if (bindingProperties.getProducer() != null && bindingProperties.getProducer().isPartitioned()) { + if (!input && bindingProperties.getProducer() != null && bindingProperties.getProducer().isPartitioned()) { messageChannel.addInterceptor(new PartitioningInterceptor(bindingProperties)); } if (StringUtils.hasText(contentType)) { - messageChannel.addInterceptor(new ContentTypeConvertingInterceptor(contentType)); + messageChannel.addInterceptor(new ContentTypeConvertingInterceptor(contentType, input)); } } - /** - * A {@link SmartMessageConverter} that delegates to another {@link SmartMessageConverter} for conversion. - * - * Will wrap the returning result of the conversion into a {@link Message} if it is not a {@link Message} - * instance already. - */ - private final class MessageWrappingMessageConverter implements SmartMessageConverter { - - private final MimeType contentType; - - private final SmartMessageConverter delegate; - - private MessageWrappingMessageConverter(SmartMessageConverter delegate, MimeType contentType) { - Assert.notNull(delegate, "Delegate converter cannot be null"); - Assert.notNull(contentType, "Content type cannot be null"); - this.delegate = delegate; - this.contentType = contentType; - } - - @Override - public Object fromMessage(Message message, Class targetClass) { - Object converted = this.delegate.fromMessage(message, targetClass); - if (converted instanceof Message) { - return converted; - } - else { - return build(converted, message.getHeaders()); - } - } - - @Override - public Object fromMessage(Message message, Class targetClass, Object conversionHint) { - Object converted = this.delegate.fromMessage(message, targetClass, conversionHint); - if (converted == null || converted instanceof Message) { - return converted; - } - else { - return build(converted, message.getHeaders()); - } - } - - @Override - public Message toMessage(Object payload, MessageHeaders headers) { - return this.delegate.toMessage(payload, headers); - } - - @Override - public Message toMessage(Object payload, MessageHeaders headers, Object conversionHint) { - return this.delegate.toMessage(payload, headers, conversionHint); - } - - /** - * Convenience method to construct a converted message - * @param payload the converted payload - * @param headers the existing message headers - * @return the converted message - */ - protected Object build(Object payload, MessageHeaders headers) { - MimeType messageContentType = MessageConverterUtils.X_JAVA_OBJECT.equals(this.contentType) ? - MessageConverterUtils.javaObjectMimeType(payload.getClass()) : this.contentType; - return MessageConverterConfigurer.this.messageBuilderFactory.withPayload(payload).copyHeaders(headers) - .copyHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, - messageContentType.toString())) - .build(); - } - } private final class ContentTypeConvertingInterceptor extends ChannelInterceptorAdapter { @@ -174,41 +116,68 @@ public class MessageConverterConfigurer implements MessageChannelConfigurer, Bea private final MimeType mimeType; - private ContentTypeConvertingInterceptor(String contentType) { + private final boolean input; + + private final Class klazz; + + private final MessageConverter messageConverter; + + private ContentTypeConvertingInterceptor(String contentType, boolean input) { this.contentType = contentType; this.mimeType = MessageConverterUtils.getMimeType(contentType); + this.input = input; + if (MessageConverterUtils.X_JAVA_OBJECT.equals(this.mimeType)) { + this.klazz = + MessageConverterUtils + .getJavaTypeForJavaObjectContentType(this.mimeType); + } + else if (this.mimeType.equals(MessageConverterUtils.X_SPRING_TUPLE)) { + this.klazz = Tuple.class; + } + else if (this.mimeType.getType().equals("text") || this.mimeType.getSubtype().equals( + "json") || this.mimeType.getSubtype().equals("xml")) { + this.klazz = String.class; + } + else { + this.klazz = byte[].class; + } + this.messageConverter = + MessageConverterConfigurer.this.compositeMessageConverterFactory + .getMessageConverterForType(this.mimeType); } @Override public Message preSend(Message message, MessageChannel channel) { - Class[] classes = - MessageConverterConfigurer.this.compositeMessageConverterFactory.supportedDataTypes( - this.mimeType); - MessageWrappingMessageConverter messageConverter = - new MessageWrappingMessageConverter( - MessageConverterConfigurer.this.compositeMessageConverterFactory - .getMessageConverterForType(this.mimeType), this.mimeType); - for (Class aClass : classes) { - if (aClass.isAssignableFrom(message.getPayload().getClass())) { - Object contentTypeFromMessage = message.getHeaders().get(MessageHeaders.CONTENT_TYPE); - if (contentTypeFromMessage == null) { - return MessageConverterConfigurer.this.messageBuilderFactory - .fromMessage(message) - .setHeader(MessageHeaders.CONTENT_TYPE, this.contentType) - .build(); - } - else { - return message; - } + Message sentMessage = null; + if (this.klazz.isAssignableFrom(message.getPayload().getClass())) { + Object contentTypeFromMessage = message.getHeaders().get(MessageHeaders.CONTENT_TYPE); + if (contentTypeFromMessage == null) { + sentMessage = MessageConverterConfigurer.this.messageBuilderFactory + .fromMessage(message) + .setHeader(MessageHeaders.CONTENT_TYPE, this.contentType) + .build(); } else { - Object converted = messageConverter.fromMessage(message, aClass); - if (converted != null) { - return (Message) converted; - } + sentMessage = message; } } - throw new MessageConversionException("Cannot convert " + message + " to " + this.contentType); + else { + Object converted = this.input ? this.messageConverter.fromMessage(message, this.klazz) + : this.messageConverter.toMessage(message.getPayload(), message.getHeaders()); + if (converted instanceof Message) { + sentMessage = (Message) converted; + } + else { + sentMessage = MessageConverterConfigurer.this.messageBuilderFactory.withPayload(converted) + .copyHeaders(message.getHeaders()).setHeaderIfAbsent(MessageHeaders.CONTENT_TYPE, + + this.mimeType).build(); + } + } + if (sentMessage == null) { + throw new MessageConversionException("Cannot convert " + message + " to " + this.contentType); + } + return sentMessage; } } diff --git a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceConfiguration.java b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceConfiguration.java index 123bf3cb4..b3e3ecb63 100644 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceConfiguration.java +++ b/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/config/ChannelBindingServiceConfiguration.java @@ -48,7 +48,6 @@ import org.springframework.cloud.stream.binding.MessageConverterConfigurer; import org.springframework.cloud.stream.binding.OutputBindingLifecycle; import org.springframework.cloud.stream.binding.SingleChannelBindable; import org.springframework.cloud.stream.binding.StreamListenerAnnotationBeanPostProcessor; -import org.springframework.cloud.stream.converter.AbstractFromMessageConverter; import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; @@ -60,6 +59,7 @@ import org.springframework.integration.config.IntegrationEvaluationContextFactor import org.springframework.integration.context.IntegrationContextUtils; import org.springframework.integration.json.JsonPropertyAccessor; import org.springframework.messaging.MessageChannel; +import org.springframework.messaging.converter.MessageConverter; import org.springframework.messaging.core.DestinationResolutionException; import org.springframework.messaging.core.DestinationResolver; import org.springframework.messaging.handler.annotation.support.DefaultMessageHandlerMethodFactory; @@ -88,7 +88,7 @@ public class ChannelBindingServiceConfiguration { * User defined custom message converters */ @Autowired(required = false) - private List customMessageConverters; + private List customMessageConverters; @Bean // This conditional is intentionally not in an autoconfig (usually a bad idea) because @@ -160,7 +160,7 @@ public class ChannelBindingServiceConfiguration { @Bean public CompositeMessageConverterFactory compositeMessageConverterFactory() { - List messageConverters = new ArrayList<>(); + List messageConverters = new ArrayList<>(); if (!CollectionUtils.isEmpty(this.customMessageConverters)) { messageConverters.addAll(Collections.unmodifiableCollection(this.customMessageConverters)); } 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 deleted file mode 100644 index 29c6aab52..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/AbstractFromMessageConverter.java +++ /dev/null @@ -1,171 +0,0 @@ -/* - * Copyright 2015-2016 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.ArrayList; -import java.util.Collection; -import java.util.Collections; -import java.util.List; - -import org.apache.commons.logging.Log; -import org.apache.commons.logging.LogFactory; - -import org.springframework.messaging.Message; -import org.springframework.messaging.MessageHeaders; -import org.springframework.messaging.converter.AbstractMessageConverter; -import org.springframework.util.Assert; -import org.springframework.util.ClassUtils; -import org.springframework.util.MimeType; - - -/** - * Base class for converters applied via Spring Integration 4.x data type channels. - * - * Extend this class to implement {@link org.springframework.messaging.converter.MessageConverter MessageConverters} - * used with custom Message conversion. Only {@link #fromMessage} is supported. - * @author David Turanski - * @author Ilayaperumal Gopinathan - * @author Marius Bogoevici - */ -public abstract class AbstractFromMessageConverter extends AbstractMessageConverter { - - protected final Log logger = LogFactory.getLog(getClass()); - - protected final List targetMimeTypes; - - /** - * Creates a converter that ignores content-type message headers - * @param targetMimeType the required target type - */ - protected AbstractFromMessageConverter(MimeType targetMimeType) { - this(new ArrayList(), targetMimeType); - } - - protected AbstractFromMessageConverter(Collection targetMimeTypes) { - this(new ArrayList(), targetMimeTypes); - } - - /** - * Creates a converter that handles one or more content-type message headers - * @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) { - super(supportedSourceMimeTypes); - Assert.notNull(targetMimeType, "'targetMimeType' cannot be null"); - this.targetMimeTypes = Collections.singletonList(targetMimeType); - } - - /** - * Creates a converter that handles one or more content-type message headers and one - * or more target MIME types - * @param supportedSourceMimeTypes a list of supported content types - * @param targetMimeTypes a list of supported target types - */ - protected AbstractFromMessageConverter(Collection supportedSourceMimeTypes, - Collection targetMimeTypes) { - super(supportedSourceMimeTypes); - Assert.notNull(targetMimeTypes, "'targetMimeTypes' cannot be null"); - this.targetMimeTypes = new ArrayList<>(targetMimeTypes); - } - - /** - * Creates a converter that requires a specific content-type message header - * @param supportedSourceMimeType {@link MimeType} that must be present in - * content-type header - * @param targetMimeType the required target type - */ - protected AbstractFromMessageConverter(MimeType supportedSourceMimeType, MimeType targetMimeType) { - this(Collections.singletonList(supportedSourceMimeType), targetMimeType); - } - - /** - * Creates a converter that requires a specific content-type message header and - * supports multiple target MIME types. - * @param supportedSourceMimeType {@link MimeType} that must be present in - * content-type header - * @param targetMimeTypes a list of supported target types - */ - protected AbstractFromMessageConverter(MimeType supportedSourceMimeType, Collection targetMimeTypes) { - this(Collections.singletonList(supportedSourceMimeType), targetMimeTypes); - } - - /** - * Subclasses implement this to specify supported target types - * @return an array of supported classes or null if any target type is supported - */ - protected abstract Class[] supportedTargetTypes(); - - /** - * Subclasses implement this to specify supported payload types - * @return an array of supported classes or null if any target type is supported - */ - protected abstract Class[] supportedPayloadTypes(); - - protected boolean supportsPayloadType(Class clazz) { - return supportsType(clazz, supportedPayloadTypes()); - } - - @Override - protected boolean supports(Class clazz) { - return supportsType(clazz, supportedTargetTypes()); - } - - private boolean supportsType(Class clazz, Class[] supportedTypes) { - if (supportedTypes != null) { - for (Class targetType : supportedTypes) { - if (ClassUtils.isAssignable(targetType, clazz)) { - return true; - } - } - return false; - } - return true; - } - - @Override - protected boolean canConvertFrom(Message message, Class targetClass) { - return super.canConvertFrom(message, targetClass) && supportsPayloadType(message.getPayload().getClass()); - } - - public boolean supportsTargetMimeType(MimeType mimeType) { - for (MimeType targetMimeType : targetMimeTypes) { - if (targetMimeType.includes(mimeType)) { - return true; - } - } - return false; - } - - /** - * Not supported by default - */ - @Override - protected boolean canConvertTo(Object payload, MessageHeaders headers) { - return false; - } - - /** - * Not supported by default - */ - @Override - public Object convertToInternal(Object payload, MessageHeaders headers, Object conversionHint) { - throw new UnsupportedOperationException("'convertTo' not supported"); - } - -} 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 deleted file mode 100644 index de2d35e70..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/ByteArrayToStringMessageConverter.java +++ /dev/null @@ -1,85 +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.io.UnsupportedEncodingException; -import java.util.ArrayList; -import java.util.Arrays; -import java.util.List; - -import org.springframework.messaging.Message; -import org.springframework.util.MimeType; -import org.springframework.util.MimeTypeUtils; - - -/** - * A {@link org.springframework.messaging.converter.MessageConverter} - * to convert from byte[] to String applying the Charset provided in - * the content-type header if any. - * - * @author David Turanski - * @author Ilayaperumal Gopinathan - */ -public class ByteArrayToStringMessageConverter extends AbstractFromMessageConverter { - - private final static List targetMimeTypes = new ArrayList(); - static { - 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); - } - - @Override - protected Class[] supportedTargetTypes() { - return new Class[] { String.class }; - } - - @Override - protected Class[] supportedPayloadTypes() { - return new Class[] { byte[].class }; - } - - /** - * Don't need to manipulate message headers. Just return payload - */ - @Override - public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { - MimeType mimeType = getContentTypeResolver().resolve(message.getHeaders()); - - String converted = null; - - if (mimeType == null || mimeType.getParameter("Charset") == null) { - converted = new String((byte[]) message.getPayload()); - } - else { - String encoding = mimeType.getParameter("Charset"); - if (encoding != null) { - try { - converted = new String((byte[]) message.getPayload(), encoding); - } - catch (UnsupportedEncodingException e) { - logger.error(e.getMessage(), e); - } - } - } - return converted; - } -} 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 09c3296cc..fad90e011 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 @@ -17,19 +17,21 @@ package org.springframework.cloud.stream.converter; import java.util.ArrayList; -import java.util.Arrays; import java.util.Collections; -import java.util.HashSet; import java.util.List; -import java.util.Set; import com.fasterxml.jackson.databind.ObjectMapper; +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import org.springframework.messaging.converter.AbstractMessageConverter; +import org.springframework.messaging.converter.ByteArrayMessageConverter; import org.springframework.messaging.converter.CompositeMessageConverter; +import org.springframework.messaging.converter.MappingJackson2MessageConverter; import org.springframework.messaging.converter.MessageConverter; +import org.springframework.messaging.converter.StringMessageConverter; import org.springframework.util.CollectionUtils; import org.springframework.util.MimeType; -import org.springframework.util.ObjectUtils; /** @@ -40,18 +42,20 @@ import org.springframework.util.ObjectUtils; */ public class CompositeMessageConverterFactory { + private final Log log = LogFactory.getLog(CompositeMessageConverterFactory.class); + private final ObjectMapper objectMapper; - private final List converters; + private final List converters; public CompositeMessageConverterFactory() { - this(Collections.emptyList(), new ObjectMapper()); + this(Collections.emptyList(), new ObjectMapper()); } /** - * @param customConverters a list of {@link AbstractFromMessageConverter} + * @param customConverters a list of {@link AbstractMessageConverter} */ - public CompositeMessageConverterFactory(List customConverters, + public CompositeMessageConverterFactory(List customConverters, ObjectMapper objectMapper) { this.objectMapper = objectMapper; if (!CollectionUtils.isEmpty(customConverters)) { @@ -65,56 +69,52 @@ public class CompositeMessageConverterFactory { private void initDefaultConverters() { - this.converters.add(new JsonToTupleMessageConverter()); - this.converters.add(new TupleToJsonMessageConverter(objectMapper)); - this.converters.add(new JsonToPojoMessageConverter(objectMapper)); - this.converters.add(new PojoToJsonMessageConverter(objectMapper)); - this.converters.add(new ByteArrayToStringMessageConverter()); - this.converters.add(new StringToByteArrayMessageConverter()); - this.converters.add(new PojoToStringMessageConverter()); - this.converters.add(new JavaToSerializedMessageConverter()); - this.converters.add(new SerializedToJavaMessageConverter()); + this.converters.add(new TupleJsonMessageConverter(this.objectMapper)); + + MappingJackson2MessageConverter jsonMessageConverter = new MappingJackson2MessageConverter(); + jsonMessageConverter.setSerializedPayloadClass(String.class); + if (this.objectMapper != null) { + jsonMessageConverter.setObjectMapper(this.objectMapper); + } + + this.converters.add(jsonMessageConverter); + this.converters.add(new ByteArrayMessageConverter()); + + this.converters.add(new StringMessageConverter()); + this.converters.add(new JavaSerializationMessageConverter()); } /** * Creation method. - * @param targetMimeType the target MIME type + * @param mimeType the target MIME type * @return a converter for the target MIME type */ - public CompositeMessageConverter getMessageConverterForType(MimeType targetMimeType) { - List targetMimeTypeConverters = new ArrayList(); - for (AbstractFromMessageConverter converter : converters) { - if (converter.supportsTargetMimeType(targetMimeType)) { - targetMimeTypeConverters.add(converter); - } - } - if (CollectionUtils.isEmpty(targetMimeTypeConverters)) { - throw new ConversionException("No message converter is registered for " - + targetMimeType.toString()); - } - return new CompositeMessageConverter(targetMimeTypeConverters); - } - - public CompositeMessageConverter getMessageConverterForAllRegistered() { - return new CompositeMessageConverter(new ArrayList(converters)); - } - - public Class[] supportedDataTypes(MimeType targetMimeType) { - Set> supportedDataTypes = new HashSet<>(); - // Make sure to check if the target type is of explicit java object type. - if (MessageConverterUtils.X_JAVA_OBJECT.includes(targetMimeType)) { - supportedDataTypes.add(MessageConverterUtils.getJavaTypeForJavaObjectContentType(targetMimeType)); - } - else { - for (AbstractFromMessageConverter converter : converters) { - if (converter.supportsTargetMimeType(targetMimeType)) { - Class[] targetTypes = converter.supportedTargetTypes(); - if (!ObjectUtils.isEmpty(targetTypes)) { - supportedDataTypes.addAll(Arrays.asList(targetTypes)); + public CompositeMessageConverter getMessageConverterForType(MimeType mimeType) { + List converters = new ArrayList<>(); + for (MessageConverter converter : this.converters) { + if (converter instanceof AbstractMessageConverter) { + for (MimeType type : ((AbstractMessageConverter) converter).getSupportedMimeTypes()) { + if (type.includes(mimeType)) { + converters.add(converter); } } } + else { + if (this.log.isDebugEnabled()) { + this.log.debug("Ommitted " + converter + " of type " + converter.getClass().toString() + + " for '" + mimeType.toString() + "' as it is not an AbstractMessageConverter"); + } + } } - return supportedDataTypes.toArray(new Class[supportedDataTypes.size()]); + if (CollectionUtils.isEmpty(converters)) { + throw new ConversionException("No message converter is registered for " + + mimeType.toString()); + } + return new CompositeMessageConverter(converters); } + + public CompositeMessageConverter getMessageConverterForAllRegistered() { + return new CompositeMessageConverter(new ArrayList<>(this.converters)); + } + } 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/JavaSerializationMessageConverter.java similarity index 50% rename from spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JavaToSerializedMessageConverter.java rename to spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JavaSerializationMessageConverter.java index a22afe575..61e078422 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/JavaSerializationMessageConverter.java @@ -1,5 +1,5 @@ /* - * Copyright 2015 the original author or authors. + * Copyright 2016 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. @@ -16,44 +16,55 @@ package org.springframework.cloud.stream.converter; +import java.io.ByteArrayInputStream; import java.io.ByteArrayOutputStream; import java.io.IOException; +import java.io.ObjectInputStream; import java.io.ObjectOutputStream; import java.io.Serializable; +import java.util.Arrays; import org.springframework.messaging.Message; - +import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.converter.AbstractMessageConverter; /** - * A {@link org.springframework.messaging.converter.MessageConverter} - * to convert from a POJO to byte[] with Java.io serialization if any. - * - * @author David Turanski + * @author Marius Bogoevici */ -public class JavaToSerializedMessageConverter extends AbstractFromMessageConverter { +public class JavaSerializationMessageConverter extends AbstractMessageConverter { - public JavaToSerializedMessageConverter() { - super(MessageConverterUtils.X_JAVA_OBJECT, MessageConverterUtils.X_JAVA_SERIALIZED_OBJECT); + public JavaSerializationMessageConverter() { + super(Arrays.asList(MessageConverterUtils.X_JAVA_SERIALIZED_OBJECT)); } @Override - protected Class[] supportedTargetTypes() { - return new Class[] { byte[].class }; - } - - @Override - protected Class[] supportedPayloadTypes() { - return new Class[] { Serializable.class }; + protected boolean supports(Class clazz) { + return Serializable.class.isAssignableFrom(clazz); } @Override public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { + if (!(message.getPayload() instanceof byte[])) { + return null; + } + try { + ByteArrayInputStream bis = new ByteArrayInputStream((byte[]) (message.getPayload())); + return new ObjectInputStream(bis).readObject(); + } + catch (Exception e) { + this.logger.error(e.getMessage(), e); + } + return null; + } + + @Override + protected Object convertToInternal(Object payload, MessageHeaders headers, Object conversionHint) { ByteArrayOutputStream bos = new ByteArrayOutputStream(); try { - new ObjectOutputStream(bos).writeObject(message.getPayload()); + new ObjectOutputStream(bos).writeObject(payload); } catch (IOException e) { - logger.error(e.getMessage(), e); + this.logger.error(e.getMessage(), e); return null; } return bos.toByteArray(); 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 deleted file mode 100644 index a3d137ce7..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToPojoMessageConverter.java +++ /dev/null @@ -1,67 +0,0 @@ -/* - * Copyright 2015-2016 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 com.fasterxml.jackson.databind.ObjectMapper; - -import org.springframework.messaging.Message; -import org.springframework.util.MimeTypeUtils; - -/** - * - * @author David Turanski - * @author Ilayaperumal Gopinathan - */ -public class JsonToPojoMessageConverter extends AbstractFromMessageConverter { - - private final ObjectMapper objectMapper; - - public JsonToPojoMessageConverter(ObjectMapper objectMapper) { - super(MimeTypeUtils.APPLICATION_JSON, MessageConverterUtils.X_JAVA_OBJECT); - this.objectMapper = (objectMapper != null) ? objectMapper : new ObjectMapper(); - } - - @Override - protected Class[] supportedPayloadTypes() { - return new Class[] {String.class, byte[].class}; - } - - @Override - protected Class[] supportedTargetTypes() { - return null; // any type - } - - @Override - public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { - Object result = null; - try { - Object payload = message.getPayload(); - - if (payload instanceof byte[]) { - result = objectMapper.readValue((byte[]) payload, targetClass); - } - else if (payload instanceof String) { - result = objectMapper.readValue((String) payload, targetClass); - } - } - catch (Exception e) { - logger.error(e.getMessage(), e); - return null; - } - return result; - } -} 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 deleted file mode 100644 index e656e0f1b..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/JsonToTupleMessageConverter.java +++ /dev/null @@ -1,69 +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.ArrayList; -import java.util.List; - -import org.springframework.messaging.Message; -import org.springframework.tuple.Tuple; -import org.springframework.tuple.TupleBuilder; -import org.springframework.util.MimeType; -import org.springframework.util.MimeTypeUtils; - - -/** - * A {@link org.springframework.messaging.converter.MessageConverter} - * to convert from a JSON (byte[] or String) to a {@link Tuple}. - * - * @author David Turanski - */ -public class JsonToTupleMessageConverter extends AbstractFromMessageConverter { - - private final static List targetMimeTypes = new ArrayList(); - - static { - targetMimeTypes.add(MessageConverterUtils.X_SPRING_TUPLE); - targetMimeTypes.add(MessageConverterUtils.X_JAVA_OBJECT); - } - - public JsonToTupleMessageConverter() { - super(MimeTypeUtils.APPLICATION_JSON, targetMimeTypes); - } - - @Override - protected Class[] supportedTargetTypes() { - return new Class[] { Tuple.class }; - } - - @Override - protected Class[] supportedPayloadTypes() { - return new Class[] { byte[].class, String.class }; - } - - @Override - public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { - String source = null; - if (message.getPayload() instanceof byte[]) { - source = new String((byte[]) message.getPayload()); - } - else { - source = (String) message.getPayload(); - } - return TupleBuilder.fromString(source); - } -} 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 deleted file mode 100644 index a8f380231..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToJsonMessageConverter.java +++ /dev/null @@ -1,79 +0,0 @@ -/* - * Copyright 2015-2016 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 com.fasterxml.jackson.core.JsonProcessingException; -import com.fasterxml.jackson.databind.ObjectMapper; - -import org.springframework.beans.factory.annotation.Value; -import org.springframework.messaging.Message; -import org.springframework.util.MimeTypeUtils; - - -/** - * A {@link org.springframework.messaging.converter.MessageConverter} - * to convert a Java object to a JSON String - * - * @author David Turanski - * @author David Liu - * @author Ilayaperumal Gopinathan - */ -public class PojoToJsonMessageConverter extends AbstractFromMessageConverter { - - private final ObjectMapper objectMapper; - - @Value("${typeconversion.json.prettyPrint:false}") - private volatile boolean prettyPrint; - - public PojoToJsonMessageConverter(ObjectMapper objectMapper) { - super(MimeTypeUtils.APPLICATION_JSON); - this.objectMapper = (objectMapper != null) ? objectMapper : new ObjectMapper(); - } - - @Override - protected Class[] supportedTargetTypes() { - return new Class[] {String.class}; - } - - @Override - protected Class[] supportedPayloadTypes() { - return null; - } - - - public void setPrettyPrint(boolean prettyPrint) { - this.prettyPrint = prettyPrint; - } - - @Override - public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { - Object result; - try { - if (prettyPrint) { - result = objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(message.getPayload()); - } - else { - result = objectMapper.writeValueAsString(message.getPayload()); - } - } - catch (JsonProcessingException e) { - logger.error(e.getMessage(), e); - return null; - } - return result; - } -} 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 deleted file mode 100644 index 6f06509d1..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/PojoToStringMessageConverter.java +++ /dev/null @@ -1,62 +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.messaging.Message; -import org.springframework.tuple.Tuple; -import org.springframework.tuple.TupleBuilder; -import org.springframework.util.MimeTypeUtils; - - -/** - * A {@link org.springframework.messaging.converter.MessageConverter} - * to convert a Java object to a String using toString() - * - * @author David Turanski - * @author Ilayaperumal Gopinathan - */ -public class PojoToStringMessageConverter extends AbstractFromMessageConverter { - - public PojoToStringMessageConverter() { - super(MimeTypeUtils.TEXT_PLAIN); - } - - @Override - protected Class[] supportedTargetTypes() { - return new Class[] { String.class }; - } - - @Override - protected Class[] supportedPayloadTypes() { - return null; - } - - @Override - 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 payloadString; - } - -} 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 deleted file mode 100644 index f92d6a062..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/SerializedToJavaMessageConverter.java +++ /dev/null @@ -1,62 +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.io.ByteArrayInputStream; -import java.io.ObjectInputStream; -import java.io.Serializable; - -import org.springframework.messaging.Message; - - -/** - * A {@link org.springframework.messaging.converter.MessageConverter} - * to deserialize {@link Serializable} Java objects. - * - * @author David Turanski - */ -public class SerializedToJavaMessageConverter extends AbstractFromMessageConverter { - - public SerializedToJavaMessageConverter() { - super(MessageConverterUtils.X_JAVA_SERIALIZED_OBJECT, MessageConverterUtils.X_JAVA_OBJECT); - } - - @Override - protected Class[] supportedTargetTypes() { - return new Class[] { Serializable.class }; - } - - @Override - protected Class[] supportedPayloadTypes() { - return new Class[] { byte[].class }; - } - - @Override - public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { - ByteArrayInputStream bis = new ByteArrayInputStream((byte[]) (message.getPayload())); - Object result = null; - try { - result = new ObjectInputStream(bis).readObject(); - } - catch (Exception e) { - logger.error(e.getMessage(), e); - return null; - } - - return result; - } -} 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 deleted file mode 100644 index cf9bd9a59..000000000 --- a/spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/StringToByteArrayMessageConverter.java +++ /dev/null @@ -1,81 +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.io.UnsupportedEncodingException; -import java.util.ArrayList; -import java.util.List; - -import org.springframework.messaging.Message; -import org.springframework.util.MimeType; -import org.springframework.util.MimeTypeUtils; - - -/** - * A {@link org.springframework.messaging.converter.MessageConverter} - * to convert a String to a byte[], applying the provided Charset in - * the content-type header if any. - * - * @author David Turanski - * @author Ilayaperumal Gopinathan - */ -public class StringToByteArrayMessageConverter extends AbstractFromMessageConverter { - - private final static List targetMimeTypes = new ArrayList(); - static { - targetMimeTypes.add(MimeTypeUtils.APPLICATION_OCTET_STREAM); - } - - public StringToByteArrayMessageConverter() { - super(targetMimeTypes); - } - - @Override - protected Class[] supportedTargetTypes() { - return new Class[] { byte[].class }; - } - - @Override - protected Class[] supportedPayloadTypes() { - return new Class[] { String.class }; - } - - /** - * Don't need to manipulate message headers. Just return the payload - */ - @Override - 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(); - } - else { - String encoding = mimeType.getParameter("Charset"); - if (encoding != null) { - try { - converted = ((String) message.getPayload()).getBytes(encoding); - } - catch (UnsupportedEncodingException e) { - logger.error(e.getMessage(), e); - } - } - } - return converted; - - } -} 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/TupleJsonMessageConverter.java similarity index 60% rename from spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/TupleToJsonMessageConverter.java rename to spring-cloud-stream/src/main/java/org/springframework/cloud/stream/converter/TupleJsonMessageConverter.java index 67b8f88b3..c2a12d23b 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/TupleJsonMessageConverter.java @@ -17,23 +17,24 @@ package org.springframework.cloud.stream.converter; import java.io.IOException; +import java.nio.charset.Charset; import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.beans.factory.annotation.Value; import org.springframework.messaging.Message; +import org.springframework.messaging.MessageHeaders; +import org.springframework.messaging.converter.AbstractMessageConverter; import org.springframework.tuple.Tuple; -import org.springframework.util.MimeTypeUtils; - +import org.springframework.tuple.TupleBuilder; /** * A {@link org.springframework.messaging.converter.MessageConverter} * to convert a {@link Tuple} to a JSON String - * * @author David Turanski * @author Ilayaperumal Gopinathan */ -public class TupleToJsonMessageConverter extends AbstractFromMessageConverter { +public class TupleJsonMessageConverter extends AbstractMessageConverter { @Value("${typeconversion.json.prettyPrint:false}") private volatile boolean prettyPrint; @@ -44,32 +45,27 @@ public class TupleToJsonMessageConverter extends AbstractFromMessageConverter { this.prettyPrint = prettyPrint; } - public TupleToJsonMessageConverter(ObjectMapper objectMapper) { - super(MimeTypeUtils.APPLICATION_JSON); + public TupleJsonMessageConverter(ObjectMapper objectMapper) { + super(MessageConverterUtils.X_SPRING_TUPLE); this.objectMapper = (objectMapper != null) ? objectMapper : new ObjectMapper(); } @Override - protected Class[] supportedTargetTypes() { - return new Class[] { String.class }; + protected boolean supports(Class clazz) { + return Tuple.class.isAssignableFrom(clazz); } @Override - protected Class[] supportedPayloadTypes() { - return new Class[] { Tuple.class }; - } - - @Override - public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { - Tuple t = (Tuple) message.getPayload(); + protected Object convertToInternal(Object payload, MessageHeaders headers, Object conversionHint) { + Tuple t = (Tuple) payload; String json; - if (prettyPrint) { + if (this.prettyPrint) { try { - Object tmp = objectMapper.readValue(t.toString(), Object.class); - json = objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(tmp); + Object tmp = this.objectMapper.readValue(t.toString(), Object.class); + json = this.objectMapper.writerWithDefaultPrettyPrinter().writeValueAsString(tmp); } catch (IOException e) { - logger.error(e.getMessage(), e); + this.logger.error(e.getMessage(), e); return null; } } @@ -79,4 +75,16 @@ public class TupleToJsonMessageConverter extends AbstractFromMessageConverter { return json; } + @Override + public Object convertFromInternal(Message message, Class targetClass, Object conversionHint) { + String source; + if (message.getPayload() instanceof byte[]) { + source = new String((byte[]) message.getPayload(), Charset.forName("UTF-8")); + } + else { + source = message.getPayload().toString(); + } + return TupleBuilder.fromString(source); + } + }