Fixed double conversion issue

- removed custom conversion from In/Out interceptors delegating everything to the available  MessageConverters
- added initial version of content-type TCK to validate various content-type conversion scenarious
- added initial content-type conversion matrix to the WIKI https://github.com/spring-cloud/spring-cloud-stream/wiki/Content-type-conversion-matrix

Resolves #1130
Resolves #1071
This commit is contained in:
Oleg Zhurakousky
2018-01-18 09:13:33 -05:00
parent f4fb7c124b
commit 9672a5b4df
16 changed files with 809 additions and 222 deletions

View File

@@ -19,6 +19,7 @@ package org.springframework.cloud.stream.binder;
import java.io.Serializable;
import java.lang.reflect.Constructor;
import java.lang.reflect.Method;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Arrays;
import java.util.List;
@@ -38,6 +39,8 @@ import org.springframework.cloud.stream.binding.StreamListenerMessageHandler;
import org.springframework.cloud.stream.config.BindingProperties;
import org.springframework.cloud.stream.config.BindingServiceProperties;
import org.springframework.cloud.stream.converter.CompositeMessageConverterFactory;
import org.springframework.cloud.stream.converter.JavaSerializationMessageConverter;
import org.springframework.cloud.stream.converter.KryoMessageConverter;
import org.springframework.cloud.stream.converter.MessageConverterUtils;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.context.Lifecycle;
@@ -48,8 +51,8 @@ import org.springframework.integration.channel.QueueChannel;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageDeliveryException;
import org.springframework.messaging.MessageHandler;
import org.springframework.messaging.MessageHandlingException;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.MessagingException;
import org.springframework.messaging.PollableChannel;
@@ -167,19 +170,19 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
getDestinationNameDelimiter()), moduleOutputChannel, outputBindingProperties.getProducer());
Binding<MessageChannel> consumerBinding = binder.bindConsumer(String.format("foo%s0",
getDestinationNameDelimiter()), "testSendAndReceive", moduleInputChannel, inputBindingProperties.getConsumer());
Message<?> message = MessageBuilder.withPayload("foo").setHeader(MessageHeaders.CONTENT_TYPE, "foo/bar")
Message<?> message = MessageBuilder.withPayload("foo").setHeader(MessageHeaders.CONTENT_TYPE, "text/plain")
.build();
// Let the consumer actually bind to the producer before sending a msg
binderBindUnbindLatency();
CountDownLatch latch = new CountDownLatch(1);
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<Message<String>>();
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<Message<byte[]>>();
moduleInputChannel.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
try {
inboundMessageRef.set((Message<String>) message);
inboundMessageRef.set((Message<byte[]>) message);
}
finally {
latch.countDown();
@@ -190,9 +193,9 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
moduleOutputChannel.send(message);
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
assertThat(new String(inboundMessageRef.get().getPayload(),StandardCharsets.UTF_8)).isEqualTo("foo");
assertThat(inboundMessageRef.get().getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)).isNull();
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()).isEqualTo("foo/bar");
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()).isEqualTo("text/plain");
producerBinding.unbind();
consumerBinding.unbind();
}
@@ -250,9 +253,10 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
moduleOutputChannel.send(message);
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
assertThat(inboundMessageRef.get().getPayload()).isInstanceOf(Foo.class);
KryoMessageConverter kryo = new KryoMessageConverter(null, true);
Foo fooPayload = (Foo) kryo.fromMessage(inboundMessageRef.get(), Foo.class);
assertNotNull(fooPayload);
assertThat(inboundMessageRef.get().getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)).isNull();
assertTrue(equalTypeAndSubType((MimeType) inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE), MessageConverterUtils.X_JAVA_OBJECT));
producerBinding.unbind();
consumerBinding.unbind();
}
@@ -262,13 +266,16 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
public void testSendAndReceiveJavaSerialization() throws Exception {
Binder binder = getBinder();
BindingProperties outputBindingProperties = createProducerBindingProperties(createProducerProperties());
DirectChannel moduleOutputChannel = createBindableChannel("output", outputBindingProperties);
BindingProperties inputBindingProperties = createConsumerBindingProperties(createConsumerProperties());
//inputBindingProperties.setContentType("tex/plain");
DirectChannel moduleInputChannel = createBindableChannel("input", inputBindingProperties);
Binding<MessageChannel> producerBinding = binder.bindProducer(String.format("foo%s0y",
getDestinationNameDelimiter()), moduleOutputChannel, outputBindingProperties.getProducer());
Binding<MessageChannel> consumerBinding = binder.bindConsumer(String.format("foo%s0y",
getDestinationNameDelimiter()), "testSendAndReceiveJavaSerialization", moduleInputChannel,
inputBindingProperties.getConsumer());
@@ -281,13 +288,13 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
binderBindUnbindLatency();
CountDownLatch latch = new CountDownLatch(1);
AtomicReference<Message<SerializableFoo>> inboundMessageRef = new AtomicReference<Message<SerializableFoo>>();
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<Message<byte[]>>();
moduleInputChannel.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
try {
inboundMessageRef.set((Message<SerializableFoo>) message);
inboundMessageRef.set((Message<byte[]>) message);
}
finally {
latch.countDown();
@@ -298,7 +305,9 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
moduleOutputChannel.send(message);
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
assertThat(inboundMessageRef.get().getPayload()).isInstanceOf(SerializableFoo.class);
JavaSerializationMessageConverter converter = new JavaSerializationMessageConverter();
SerializableFoo serializableFoo = (SerializableFoo) converter.convertFromInternal(inboundMessageRef.get(), SerializableFoo.class, null);
assertNotNull(serializableFoo);
assertThat(inboundMessageRef.get().getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)).isNull();
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE)).isEqualTo(MessageConverterUtils.X_JAVA_SERIALIZED_OBJECT);
producerBinding.unbind();
@@ -385,13 +394,13 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
.setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build();
moduleOutputChannel.send(message);
CountDownLatch latch = new CountDownLatch(1);
AtomicReference<Message<String>> inboundMessageRef = new AtomicReference<Message<String>>();
AtomicReference<Message<byte[]>> inboundMessageRef = new AtomicReference<Message<byte[]>>();
moduleInputChannel.subscribe(new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
try {
inboundMessageRef.set((Message<String>) message);
inboundMessageRef.set((Message<byte[]>) message);
}
finally {
latch.countDown();
@@ -402,7 +411,7 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
moduleOutputChannel.send(message);
Assert.isTrue(latch.await(5, TimeUnit.SECONDS), "Failed to receive message");
assertThat(inboundMessageRef.get()).isNotNull();
assertThat(inboundMessageRef.get().getPayload()).isEqualTo("foo");
assertThat(new String(inboundMessageRef.get().getPayload(), StandardCharsets.UTF_8)).isEqualTo("foo");
assertThat(inboundMessageRef.get().getHeaders().get(MessageHeaders.CONTENT_TYPE).toString())
.isEqualTo(MimeTypeUtils.TEXT_PLAIN_VALUE);
producerBinding.unbind();
@@ -575,7 +584,7 @@ public abstract class AbstractBinderTests<B extends AbstractTestBinder<? extends
}
@SuppressWarnings("rawtypes")
@Test(expected = MessageHandlingException.class)
@Test(expected = MessageDeliveryException.class)
public void testStreamListenerJavaSerializationNonSerializable() throws Exception {
Binder binder = getBinder();

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2017 the original author or authors.
* Copyright 2017-2018 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.
@@ -73,7 +73,7 @@ public class DefaultHeaderPropagationWithApplicationProvidedHeaderTests {
@ServiceActivator(inputChannel = "input", outputChannel = "output")
public Message<?> consume(String data) {
return MessageBuilder.withPayload(data).setHeader(MessageHeaders.CONTENT_TYPE, "custom/header").build();
return MessageBuilder.withPayload(data).setHeader(MessageHeaders.CONTENT_TYPE, "text/plain").build();
}
}

View File

@@ -16,6 +16,7 @@
package org.springframework.cloud.stream.config;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.TimeUnit;
import org.junit.Test;
@@ -58,11 +59,11 @@ public class InboundJsonToTupleConversionTest {
testProcessor.input().send(MessageBuilder.withPayload("{'name':'foo'}")
.build());
@SuppressWarnings("unchecked")
Message<String> received = (Message<String>) ((TestSupportBinder) binderFactory.getBinder(null, MessageChannel.class))
Message<byte[]> received = (Message<byte[]>) ((TestSupportBinder) binderFactory.getBinder(null, MessageChannel.class))
.messageCollector().forChannel(testProcessor.output()).poll(1, TimeUnit.SECONDS);
assertThat(received).isNotNull();
assertThat(TupleBuilder.fromString(new String(received.getPayload()))).isEqualTo(TupleBuilder.tuple().of("name", "foo"));
String payload = new String(received.getPayload(), StandardCharsets.UTF_8);
assertThat(TupleBuilder.fromString(payload)).isEqualTo(TupleBuilder.tuple().of("name", "foo"));
}
@EnableBinding(Processor.class)

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2015-2017 the original author or authors.
* Copyright 2015-2018 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,6 +16,7 @@
package org.springframework.cloud.stream.config;
import java.nio.charset.StandardCharsets;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
@@ -54,8 +55,8 @@ public class LegacyContentTypeTests {
MessageHandler messageHandler = new MessageHandler() {
@Override
public void handleMessage(Message<?> message) throws MessagingException {
assertThat(message.getPayload()).isInstanceOf(String.class);
assertThat(message.getPayload()).isEqualTo("{\"message\":\"Hi\"}");
assertThat(message.getPayload()).isInstanceOf(byte[].class);
assertThat(new String(((byte[])message.getPayload()), StandardCharsets.UTF_8)).isEqualTo("{\"message\":\"Hi\"}");
assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString()).isEqualTo("application/json");
latch.countDown();
}
@@ -63,7 +64,6 @@ public class LegacyContentTypeTests {
testSink.input().subscribe(messageHandler);
testSink.input().send(MessageBuilder.withPayload("{\"message\":\"Hi\"}".getBytes())
.setHeader(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, "application/json")
.setHeader(BinderHeaders.SCST_VERSION, "1.x")
.build());
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
testSink.input().unsubscribe(messageHandler);

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2017 the original author or authors.
* Copyright 2017-2018 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.
@@ -40,7 +40,6 @@ import org.springframework.cloud.stream.test.binder.MessageCollector;
import org.springframework.context.ConfigurableApplicationContext;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageDeliveryException;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.SubscribableChannel;
import org.springframework.messaging.handler.annotation.Headers;
@@ -229,7 +228,7 @@ public class ContentTypeTests {
assertThat(message.getPayload()).isEqualTo(user.toString());
}
}
@Test
public void testSendTuple() throws Exception {
try (ConfigurableApplicationContext context = SpringApplication.run(
@@ -313,7 +312,7 @@ public class ContentTypeTests {
}
}
@Test(expected=MessageDeliveryException.class)
@Test
public void testReceiveKryoWithHeadersOverridingDefault() throws Exception{
try (ConfigurableApplicationContext context = SpringApplication.run(
SinkApplication.class, "--server.port=0",

View File

@@ -1 +1,2 @@
spring.cloud.stream.bindings.input.content-type=application/x-spring-tuple
spring.cloud.stream.bindings.output.content-type=application/x-spring-tuple

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016 the original author or authors.
* Copyright 2016-2018 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.
@@ -90,6 +90,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver
buf.get(payload);
Schema writerSchema = resolveWriterSchemaForDeserialization(mimeType);
Schema readerSchema = resolveReaderSchemaForDeserialization(targetClass);
@SuppressWarnings("unchecked")
DatumReader<Object> reader = getDatumReader((Class<Object>) targetClass, readerSchema, writerSchema);
Decoder decoder = DecoderFactory.get().binaryDecoder(payload, null);
result = reader.read(null, decoder);
@@ -125,6 +126,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver
return writer;
}
@SuppressWarnings({ "unchecked", "rawtypes" })
protected DatumReader<Object> getDatumReader(Class<Object> type, Schema schema, Schema writerSchema) {
DatumReader<Object> reader = null;
if (SpecificRecord.class.isAssignableFrom(type)) {
@@ -176,6 +178,7 @@ public abstract class AbstractAvroMessageConverter extends AbstractMessageConver
hintedContentType = (MimeType) conversionHint;
}
Schema schema = resolveSchemaForWriting(payload, headers, hintedContentType);
@SuppressWarnings("unchecked")
DatumWriter<Object> writer = getDatumWriter((Class<Object>) payload.getClass(), schema);
Encoder encoder = EncoderFactory.get().binaryEncoder(baos, null);
writer.write(payload, encoder);

View File

@@ -1,5 +1,5 @@
/*
* Copyright 2016-2017 the original author or authors.
* Copyright 2016-2018 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.
@@ -18,6 +18,7 @@ package org.springframework.cloud.stream.schema.avro;
import java.io.IOException;
import java.util.Arrays;
import java.util.Map;
import java.util.regex.Matcher;
import java.util.regex.Pattern;
@@ -25,6 +26,7 @@ import org.apache.avro.Schema;
import org.apache.avro.generic.GenericContainer;
import org.apache.avro.reflect.ReflectData;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.BeanInitializationException;
import org.springframework.beans.factory.InitializingBean;
import org.springframework.cache.CacheManager;
@@ -35,7 +37,6 @@ import org.springframework.cloud.stream.schema.SchemaReference;
import org.springframework.cloud.stream.schema.SchemaRegistrationResponse;
import org.springframework.cloud.stream.schema.client.SchemaRegistryClient;
import org.springframework.core.io.Resource;
import org.springframework.integration.support.MutableMessageHeaders;
import org.springframework.messaging.MessageHeaders;
import org.springframework.util.Assert;
import org.springframework.util.MimeType;
@@ -250,12 +251,13 @@ public class AvroSchemaRegistryClientMessageConverter extends AbstractAvroMessag
SchemaReference schemaReference = parsedSchema.getRegistration()
.getSchemaReference();
if (headers instanceof MutableMessageHeaders) {
headers.put(MessageHeaders.CONTENT_TYPE,
"application/" + this.prefix + "." + schemaReference.getSubject()
+ ".v" + schemaReference.getVersion() + "+avro");
}
DirectFieldAccessor dfa = new DirectFieldAccessor(headers);
@SuppressWarnings("unchecked")
Map<String, Object> _headers = (Map<String, Object>) dfa.getPropertyValue("headers");
_headers.put(MessageHeaders.CONTENT_TYPE,
"application/" + this.prefix + "." + schemaReference.getSubject()
+ ".v" + schemaReference.getVersion() + "+avro");
return schema;
}

View File

@@ -16,7 +16,7 @@
package org.springframework.cloud.stream.binding;
import java.nio.charset.StandardCharsets;
import java.lang.reflect.Field;
import java.util.Collections;
import java.util.Map;
@@ -24,6 +24,7 @@ import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.beans.BeansException;
import org.springframework.beans.DirectFieldAccessor;
import org.springframework.beans.factory.BeanFactory;
import org.springframework.beans.factory.BeanFactoryAware;
import org.springframework.beans.factory.config.ConfigurableListableBeanFactory;
@@ -32,7 +33,6 @@ import org.springframework.cloud.stream.binder.BinderHeaders;
import org.springframework.cloud.stream.binder.ConsumerProperties;
import org.springframework.cloud.stream.binder.DefaultPollableMessageSource;
import org.springframework.cloud.stream.binder.JavaClassMimeTypeUtils;
import org.springframework.cloud.stream.binder.MessageValues;
import org.springframework.cloud.stream.binder.PartitionHandler;
import org.springframework.cloud.stream.binder.PartitionKeyExtractorStrategy;
import org.springframework.cloud.stream.binder.PartitionSelectorStrategy;
@@ -46,7 +46,6 @@ import org.springframework.integration.channel.AbstractMessageChannel;
import org.springframework.integration.expression.ExpressionUtils;
import org.springframework.integration.support.MessageBuilderFactory;
import org.springframework.integration.support.MutableMessageBuilderFactory;
import org.springframework.integration.support.MutableMessageHeaders;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageChannel;
import org.springframework.messaging.MessageHeaders;
@@ -55,12 +54,11 @@ import org.springframework.messaging.converter.MessageConverter;
import org.springframework.messaging.handler.invocation.InvocableHandlerMethod;
import org.springframework.messaging.support.ChannelInterceptorAdapter;
import org.springframework.messaging.support.ErrorMessage;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.util.Assert;
import org.springframework.util.CollectionUtils;
import org.springframework.util.MimeType;
import org.springframework.util.MimeTypeUtils;
import org.springframework.util.ObjectUtils;
import org.springframework.util.ReflectionUtils;
import org.springframework.util.StringUtils;
/**
@@ -92,6 +90,8 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
private final Map<String, PartitionKeyExtractorStrategy> partitionKeyExtractors;
private final Map<String, PartitionSelectorStrategy> partitionSelectors;
private final Field headersField;
public MessageConverterConfigurer(BindingServiceProperties bindingServiceProperties,
CompositeMessageConverterFactory compositeMessageConverterFactory) {
@@ -108,6 +108,9 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
this.compositeMessageConverterFactory = compositeMessageConverterFactory;
this.partitionKeyExtractors = partitionKeyExtractors == null ? Collections.emptyMap() : partitionKeyExtractors;
this.partitionSelectors = partitionSelectors == null ? Collections.emptyMap() : partitionSelectors;
this.headersField = ReflectionUtils.findField(MessageHeaders.class, "headers");
headersField.setAccessible(true);
}
@Override
@@ -133,8 +136,7 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
ConsumerProperties consumerProperties = bindingProperties.getConsumer();
if ((consumerProperties == null || !consumerProperties.isUseNativeDecoding())
&& binding instanceof DefaultPollableMessageSource) {
((DefaultPollableMessageSource) binding).addInterceptor(
new InboundContentTypeConvertingInterceptor(contentType, this.compositeMessageConverterFactory));
((DefaultPollableMessageSource) binding).addInterceptor(new InboundContentTypeEnhancingInterceptor(contentType));
}
}
@@ -160,7 +162,7 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
ConsumerProperties consumerProperties = bindingProperties.getConsumer();
if (this.isNativeEncodingNotSet(producerProperties, consumerProperties, inbound)) {
if (inbound) {
messageChannel.addInterceptor(new InboundContentTypeConvertingInterceptor(contentType, this.compositeMessageConverterFactory));
messageChannel.addInterceptor(new InboundContentTypeEnhancingInterceptor(contentType));
}
else {
messageChannel.addInterceptor(new OutboundContentTypeConvertingInterceptor(contentType, this.compositeMessageConverterFactory
@@ -257,105 +259,42 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
return Math.abs(hashCode);
}
}
/**
* Primary purpose of this interceptor is to enhance/enrich Message that sent to the *inbound*
* channel with 'contentType' header for cases where 'contentType' is not present in the Message
* itself but set on such channel via {@link BindingProperties#setContentType(String)}.
* <br>
* Secondary purpose of this interceptor is to provide backward compatibility with previous versions of SCSt
* to support some of the type conversion assumptions.
* See InboundContentTypeConvertingInterceptor.deserializePayload(..) for more details.
* Secondary purpose of this interceptor is to provide backward compatibility with previous versions of SCSt.
*/
private final class InboundContentTypeConvertingInterceptor extends ChannelInterceptorAdapter {
private final class InboundContentTypeEnhancingInterceptor extends AbstractContentTypeInterceptor {
private final MimeType mimeType;
private final CompositeMessageConverterFactory compositeMessageConverterFactory;
private InboundContentTypeConvertingInterceptor(String contentType, CompositeMessageConverterFactory compositeMessageConverterFactory) {
this.mimeType = MessageConverterUtils.getMimeType(contentType);
this.compositeMessageConverterFactory = compositeMessageConverterFactory;
private InboundContentTypeEnhancingInterceptor(String contentType) {
super(contentType);
}
@Override
public Message<?> preSend(Message<?> message, MessageChannel channel) {
if (message instanceof ErrorMessage) {
return message;
public Message<?> doPreSend(Message<?> message, MessageChannel channel) {
@SuppressWarnings("unchecked")
Map<String, Object> headersMap = (Map<String, Object>) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders());
/*
* NOTE: The below code for BINDER_ORIGINAL_CONTENT_TYPE is to support legacy message format established
* in 1.x version of the framework and should/will no longer be supported in 3.x
*/
Object ct = message.getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE);
MimeType contentType = ct instanceof String ? MimeType.valueOf((String)ct) : (ct == null ? this.mimeType : (MimeType)ct);
headersMap.remove(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE);
// == end legacy note
if (!message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE)) {
headersMap.put(MessageHeaders.CONTENT_TYPE, contentType);
}
Message<?> postProcessedMessage = message;
MimeType contentType = this.mimeType;
if (message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE)) {
Object ct = message.getHeaders().get(MessageHeaders.CONTENT_TYPE);
contentType = ct instanceof String ? MimeType.valueOf((String)ct) : (MimeType)ct;
else if (message.getHeaders().get(MessageHeaders.CONTENT_TYPE) instanceof String) {
headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf((String)message.getHeaders().get(MessageHeaders.CONTENT_TYPE)));
}
boolean deserializationRequired = message.getPayload() instanceof byte[] &&
("text".equalsIgnoreCase(contentType.getType()) ||
equalTypeAndSubType(MimeTypeUtils.APPLICATION_JSON, contentType) ||
equalTypeAndSubType(MessageConverterUtils.X_JAVA_SERIALIZED_OBJECT, contentType) ||
equalTypeAndSubType(MessageConverterUtils.X_JAVA_OBJECT, contentType));
Object payload = deserializationRequired ? this.deserializePayload(message, contentType) : message.getPayload();
if (payload != null) {
Object ct = message.getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE);
contentType = ct instanceof String ? MimeType.valueOf((String)ct) : (ct == null ? contentType : (MimeType)ct);
postProcessedMessage = MessageConverterConfigurer.this.messageBuilderFactory
.withPayload(payload)
.copyHeaders(message.getHeaders())
.setHeader(MessageHeaders.CONTENT_TYPE, contentType)
.removeHeader(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE)
.build();
}
return postProcessedMessage;
}
/**
* Will *only* deserialize payload if its 'contentType' is 'text/* or application/json' or Java/Kryo serialized.
* While this would naturally happen via MessageConverters at the time of handler method
* invocation, doing it here also is strictly to support behavior established
* in previous versions of SCSt. One of these cases is return payload as String if contentType is text or json.
* Also to support certain type of assumptions on type-less handlers (i.e., handle(?) vs. handle(Foo));
*/
private Object deserializePayload(Message<?> message, MimeType contentTypeToUse) {
Object payload = null;
if ("text".equalsIgnoreCase(contentTypeToUse.getType()) || equalTypeAndSubType(MimeTypeUtils.APPLICATION_JSON, contentTypeToUse)) {
payload = new String((byte[])message.getPayload(), StandardCharsets.UTF_8);
}
else {
message = MessageBuilder.fromMessage(message).setHeader(MessageHeaders.CONTENT_TYPE, contentTypeToUse).build();
MessageConverter converter = equalTypeAndSubType(MessageConverterUtils.X_JAVA_SERIALIZED_OBJECT, contentTypeToUse)
? compositeMessageConverterFactory.getMessageConverterForType(contentTypeToUse)
: compositeMessageConverterFactory.getMessageConverterForAllRegistered();
String targetClassName = contentTypeToUse.getParameter("type");
Class<?> targetClass = null;
if (StringUtils.hasText(targetClassName)) {
try {
targetClass = Class.forName(targetClassName, false, Thread.currentThread().getContextClassLoader());
}
catch (Exception e) {
throw new IllegalStateException("Failed to determine class name for contentType: "
+ message.getHeaders().get(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE), e);
}
}
Assert.isTrue(!(equalTypeAndSubType(MessageConverterUtils.X_JAVA_OBJECT, contentTypeToUse) && targetClass == null),
"Cannot deserialize into message since 'contentType` is not "
+ "encoded with the actual target type."
+ "Consider 'application/x-java-object; type=foo.bar.MyClass'");
payload = converter.fromMessage(message, targetClass);
}
return payload;
}
/*
* Candidate to go into some utils class
*/
private boolean equalTypeAndSubType(MimeType m1, MimeType m2) {
return m1 != null && m2 != null && m1.getType().equalsIgnoreCase(m2.getType()) && m1.getSubtype().equalsIgnoreCase(m2.getSubtype());
return message;
}
}
@@ -365,67 +304,67 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
* rely on provided MessageConverters that will use the provided 'contentType' and convert messages
* to a type dictated by the Binders (i.e., byte[]).
*/
private final class OutboundContentTypeConvertingInterceptor extends ChannelInterceptorAdapter {
private final MimeType mimeType;
private final class OutboundContentTypeConvertingInterceptor extends AbstractContentTypeInterceptor {
private final MessageConverter messageConverter;
private OutboundContentTypeConvertingInterceptor(String contentType, CompositeMessageConverter messageConverter) {
this.mimeType = MessageConverterUtils.getMimeType(contentType);
super(contentType);
this.messageConverter = messageConverter;
}
@Override
public Message<?> preSend(Message<?> message, MessageChannel channel) {
Message<?> postProcessedMessage = message;
if (!(message instanceof ErrorMessage)) {
MutableMessageHeaders headers = new MutableMessageHeaders(message.getHeaders());
headers.putIfAbsent(MessageHeaders.CONTENT_TYPE, this.mimeType);
Message<?> converted = this.messageConverter.toMessage(message.getPayload(), headers);
if (converted != null) {
postProcessedMessage = converted;
} else {
postProcessedMessage = MessageConverterConfigurer.this.messageBuilderFactory
.withPayload(message.getPayload())
.copyHeaders(headers)
.build();
}
postProcessedMessage = this.finishPreSend(postProcessedMessage);
}
return postProcessedMessage;
}
/**
* This is strictly to support 1.3 semantics where BINDER_ORIGINAL_CONTENT_TYPE header
* needs to be set for certain cases and String payloads needs to be converted to byte[].
*
* Factored out of what was left of MessageSerializationUtils.
*/
// deprecated at the get go as a reminder to remove at v3.0
@Deprecated
private Message<?> finishPreSend(Message<?> message) {
public Message<?> doPreSend(Message<?> message, MessageChannel channel) {
String oct = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE) ? message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString() : null;
String ct = oct;
if (message.getPayload() instanceof String) {
ct = JavaClassMimeTypeUtils.mimeTypeFromObject(message.getPayload(), ObjectUtils.nullSafeToString(oct)).toString();
}
MessageValues messageValues = new MessageValues(message);
Object payload = message.getPayload();
if (payload instanceof String) {
payload = ((String)payload).getBytes(StandardCharsets.UTF_8);
if (!message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE)) {
@SuppressWarnings("unchecked")
Map<String, Object> headersMap = (Map<String, Object>) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders());
headersMap.put(MessageHeaders.CONTENT_TYPE, this.mimeType);
}
messageValues.setPayload(payload);
if (ct != null && !ct.equals(oct)) {
messageValues.put(MessageHeaders.CONTENT_TYPE, ct);
messageValues.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, oct);
@SuppressWarnings("unchecked")
Message<byte[]> outboundMessage = message.getPayload() instanceof byte[]
? (Message<byte[]>)message : (Message<byte[]>) this.messageConverter.toMessage(message.getPayload(), message.getHeaders());
if (outboundMessage == null) {
throw new IllegalStateException("Failed to convert message: '" + message + "' to outbound message.");
}
return messageValues.toMessage();
if (ct != null && !ct.equals(oct)) {
@SuppressWarnings("unchecked")
Map<String, Object> headersMap = (Map<String, Object>) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders());
headersMap.put(MessageHeaders.CONTENT_TYPE, MimeType.valueOf(ct));
headersMap.put(BinderHeaders.BINDER_ORIGINAL_CONTENT_TYPE, MimeType.valueOf(oct));
}
return outboundMessage;
}
}
/**
*
*/
private abstract class AbstractContentTypeInterceptor extends ChannelInterceptorAdapter {
final MimeType mimeType;
private AbstractContentTypeInterceptor(String contentType) {
this.mimeType = MessageConverterUtils.getMimeType(contentType);
}
@Override
public Message<?> preSend(Message<?> message, MessageChannel channel) {
return message instanceof ErrorMessage ? message : this.doPreSend(message, channel);
}
protected abstract Message<?> doPreSend(Message<?> message, MessageChannel channel);
}
/**
*
*/
protected final class PartitioningInterceptor extends ChannelInterceptorAdapter {
private final BindingProperties bindingProperties;

View File

@@ -0,0 +1,73 @@
/*
* Copyright 2018 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.nio.charset.StandardCharsets;
import org.springframework.core.MethodParameter;
import org.springframework.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.converter.MappingJackson2MessageConverter;
/**
* Variation of {@link MappingJackson2MessageConverter} to support marshalling and
* unmarshalling of Messages's payload from 'byte[]' to and instance of a 'targetClass' and vice versa.
*
*
* @author Oleg Zhurakousky
* @since 2.0
*
*/
class ApplicationJsonMessageMarshallingConverter extends MappingJackson2MessageConverter {
@Override
protected Object convertToInternal(Object payload, @Nullable MessageHeaders headers, @Nullable Object conversionHint) {
if (payload instanceof byte[]){
return payload;
}
else if (payload instanceof String) {
return ((String)payload).getBytes(StandardCharsets.UTF_8);
}
else {
return super.convertToInternal(payload, headers, conversionHint);
}
}
@Override
protected Object convertFromInternal(Message<?> message, Class<?> targetClass, @Nullable Object conversionHint) {
Object result = null;
if (conversionHint instanceof MethodParameter) {
Class<?> conversionHintType = ((MethodParameter)conversionHint).getParameterType();
if (Message.class.isAssignableFrom(conversionHintType)) {
/*
* Ensures that super won't attempt to create Message as a result of conversion
* and stays at payload conversion only.
* The Message will eventually be created in MessageMethodArgumentResolver.resolveArgument(..)
*/
conversionHint = null;
}
}
if (message.getPayload() instanceof byte[] && targetClass.isAssignableFrom(String.class)) {
result = new String((byte[])message.getPayload(), StandardCharsets.UTF_8);
}
else {
result = super.convertFromInternal(message, targetClass, conversionHint);
}
return result;
}
}

View File

@@ -16,7 +16,6 @@
package org.springframework.cloud.stream.converter;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
@@ -26,14 +25,9 @@ import com.fasterxml.jackson.databind.ObjectMapper;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.springframework.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHeaders;
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.MessageConversionException;
import org.springframework.messaging.converter.MessageConverter;
import org.springframework.util.CollectionUtils;
import org.springframework.util.MimeType;
@@ -75,41 +69,14 @@ public class CompositeMessageConverterFactory {
}
private void initDefaultConverters() {
this.converters.add(new TupleJsonMessageConverter(this.objectMapper));
MappingJackson2MessageConverter jsonMessageConverter = new MappingJackson2MessageConverter() {
@Override
protected Object convertToInternal(Object payload, @Nullable MessageHeaders headers, @Nullable Object conversionHint) {
if (payload instanceof byte[]){
return payload;
}
else if (payload instanceof String) {
return ((String)payload).getBytes(StandardCharsets.UTF_8);
}
else {
return super.convertToInternal(payload, headers, conversionHint);
}
}
@Override
protected Object convertFromInternal(Message<?> message, Class<?> targetClass, @Nullable Object conversionHint) {
try{
return super.convertFromInternal(message, targetClass, conversionHint);
} catch (MessageConversionException me){
//Strings need special treatment
if(targetClass.isAssignableFrom(String.class)){
return message.getPayload();
}
throw me;
}
}
};
jsonMessageConverter.setStrictContentTypeMatch(true);
ApplicationJsonMessageMarshallingConverter applicationJson = new ApplicationJsonMessageMarshallingConverter();
applicationJson.setStrictContentTypeMatch(true);
if (this.objectMapper != null) {
jsonMessageConverter.setObjectMapper(this.objectMapper);
applicationJson.setObjectMapper(this.objectMapper);
}
this.converters.add(jsonMessageConverter);
this.converters.add(applicationJson);
this.converters.add(new TupleJsonMessageConverter(this.objectMapper));
this.converters.add(new ByteArrayMessageConverter());
this.converters.add(new ObjectStringMessageConverter());
this.converters.add(new JavaSerializationMessageConverter());

View File

@@ -34,7 +34,7 @@ public class SourceDestination extends AbstractDestination {
* to binder's input destination (e.g., Processor.INPUT).
*
*/
public void send(Message<?> message) {
public void send(Message<byte[]> message) {
this.getChannel().send(message);
}
}

View File

@@ -128,10 +128,6 @@ public class SpringIntegrationChannelBinder extends AbstractMessageChannelBinder
return this.lastError;
}
public void setLastError(Message<?> lastError) {
this.lastError = lastError;
}
@Override
protected MessageHandler createProducerMessageHandler(ProducerDestination destination,
ProducerProperties producerProperties, MessageChannel errorChannel) throws Exception {

View File

@@ -39,9 +39,10 @@ public class TargetDestination extends AbstractDestination {
* Allows to access {@link Message}s received by this {@link TargetDestination}.
* @param timeout how long to wait before giving up
*/
public Message<?> receive(long timeout) {
@SuppressWarnings("unchecked")
public Message<byte[]> receive(long timeout) {
try {
return this.messages.poll(timeout, TimeUnit.MILLISECONDS);
return (Message<byte[]>) this.messages.poll(timeout, TimeUnit.MILLISECONDS);
}
catch (InterruptedException e) {
Thread.currentThread().interrupt();
@@ -52,7 +53,7 @@ public class TargetDestination extends AbstractDestination {
/**
* Allows to access {@link Message}s received by this {@link TargetDestination}.
*/
public Message<?> receive() {
public Message<byte[]> receive() {
return this.receive(0);
}

View File

@@ -0,0 +1,596 @@
/*
* Copyright 2017 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.binder.tck;
import java.nio.charset.StandardCharsets;
import java.util.Collections;
import com.fasterxml.jackson.databind.ObjectMapper;
import org.junit.Test;
import org.springframework.boot.WebApplicationType;
import org.springframework.boot.builder.SpringApplicationBuilder;
import org.springframework.cloud.stream.annotation.EnableBinding;
import org.springframework.cloud.stream.annotation.StreamListener;
import org.springframework.cloud.stream.annotation.StreamMessageConverter;
import org.springframework.cloud.stream.binder.integration.SourceDestination;
import org.springframework.cloud.stream.binder.integration.SpringIntegrationBinderConfiguration;
import org.springframework.cloud.stream.binder.integration.SpringIntegrationChannelBinder;
import org.springframework.cloud.stream.binder.integration.TargetDestination;
import org.springframework.cloud.stream.converter.KryoMessageConverter;
import org.springframework.cloud.stream.converter.MessageConverterUtils;
import org.springframework.cloud.stream.messaging.Processor;
import org.springframework.context.ApplicationContext;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.context.annotation.Import;
import org.springframework.integration.annotation.ServiceActivator;
import org.springframework.lang.Nullable;
import org.springframework.messaging.Message;
import org.springframework.messaging.MessageHeaders;
import org.springframework.messaging.converter.AbstractMessageConverter;
import org.springframework.messaging.converter.MessageConversionException;
import org.springframework.messaging.handler.annotation.SendTo;
import org.springframework.messaging.support.GenericMessage;
import org.springframework.messaging.support.MessageBuilder;
import org.springframework.util.MimeType;
import org.springframework.util.MimeTypeUtils;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertTrue;
/**
* Sort of a TCK test suite to validate payload conversion is
* done properly by interacting with binder's input/output destinations
* instead of its bridged channels.
* This means that all payloads (sent/received) must be expressed in the
* wire format (byte[])
*
* @author Oleg Zhurakousky
*
*/
public class ContentTypeTckTests {
@Test
public void pojoToPojo() {
ApplicationContext context = new SpringApplicationBuilder(PojoToPojoStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void pojoToString() {
ApplicationContext context = new SpringApplicationBuilder(PojoToStringStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void pojoToStringOutboundContentTypeBinding() {
ApplicationContext context = new SpringApplicationBuilder(PojoToStringStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.bindings.output.contentType=text/plain", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void pojoToByteArray() {
ApplicationContext context = new SpringApplicationBuilder(PojoToByteArrayStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void pojoToByteArrayOutboundContentTypeBinding() {
ApplicationContext context = new SpringApplicationBuilder(PojoToByteArrayStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.bindings.output.contentType=text/plain", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void stringToPojoInboundContentTypeBinding() {
ApplicationContext context = new SpringApplicationBuilder(StringToPojoStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.bindings.input.contentType=text/plain", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void stringToPojoInboundContentTypeHeader() {
ApplicationContext context = new SpringApplicationBuilder(StringToPojoStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes(), new MessageHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN))));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void byteArrayToPojoInboundContentTypeBinding() {
ApplicationContext context = new SpringApplicationBuilder(ByteArrayToPojoStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.bindings.input.contentType=text/plain", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void byteArrayToPojoInboundContentTypeHeader() {
ApplicationContext context = new SpringApplicationBuilder(StringToPojoStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes(), new MessageHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN))));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void byteArrayToByteArray() {
ApplicationContext context = new SpringApplicationBuilder(ByteArrayToByteArrayStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void byteArrayToByteArrayInboundOutboundContentTypeBinding() {
ApplicationContext context = new SpringApplicationBuilder(ByteArrayToByteArrayStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.bindings.input.contentType=text/plain", "--spring.cloud.stream.bindings.output.contentType=text/plain", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void pojoMessageToStringMessage() {
ApplicationContext context = new SpringApplicationBuilder(PojoMessageToStringMessageStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void pojoMessageToStringMessageServiceActivator() {
ApplicationContext context = new SpringApplicationBuilder(PojoMessageToStringMessageServiceActivator.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void byteArrayMessageToStringJsonMessageStreamListener() {
ApplicationContext context = new SpringApplicationBuilder(ByteArrayMessageToStringJsonMessageStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals("{\"name\":\"bob\"}", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void byteArrayMessageToStringMessageStreamListener() {
ApplicationContext context = new SpringApplicationBuilder(StringMessageToStringMessageStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
}
@Test
public void kryo_pojoToPojo() {
ApplicationContext context = new SpringApplicationBuilder(PojoToPojoStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.default.contentType=application/x-java-object", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
KryoMessageConverter converter = new KryoMessageConverter(null, true);
@SuppressWarnings("unchecked")
Message<byte[]> message = (Message<byte[]>) converter
.toMessage(new Person("oleg"), new MessageHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, MessageConverterUtils.X_JAVA_OBJECT)));
source.send(new GenericMessage<byte[]>(message.getPayload()));
Message<byte[]> outputMessage = target.receive();
assertNotNull(outputMessage);
MimeType contentType = (MimeType) outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE);
assertEquals("x-java-object", contentType.getSubtype());
assertEquals(Person.class.getName(), contentType.getParameters().get("type"));
}
@Test
public void kryo_pojoToPojoContentTypeHeader() {
ApplicationContext context = new SpringApplicationBuilder(PojoToPojoStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.jmx.enabled=false", "--spring.cloud.stream.bindings.output.contentType=application/x-java-object");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
KryoMessageConverter converter = new KryoMessageConverter(null, true);
@SuppressWarnings("unchecked")
Message<byte[]> message = (Message<byte[]>) converter
.toMessage(new Person("oleg"), new MessageHeaders(Collections.singletonMap(MessageHeaders.CONTENT_TYPE, MessageConverterUtils.X_JAVA_OBJECT)));
source.send(message);
Message<byte[]> outputMessage = target.receive();
assertNotNull(outputMessage);
MimeType contentType = (MimeType) outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE);
assertEquals("x-java-object", contentType.getSubtype());
}
/**
* This test simply demonstrates how one can override an existing MessageConverter for a given contentType.
* In this case we are demonstrating how Kryo converter can be overriden ('application/x-java-object' maps to Kryo).
*/
@Test
public void overrideMessageConverter_defaultContentTypeBinding() {
ApplicationContext context = new SpringApplicationBuilder(StringToStringStreamListener.class, CustomConverters.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.default.contentType=application/x-java-object", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertNotNull(outputMessage);
System.out.println(new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
assertEquals("AlwaysStringKryoMessageConverter", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
assertEquals(MimeType.valueOf("application/x-java-object"), outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
}
@Test
public void customMessageConverter_defaultContentTypeBinding() {
ApplicationContext context = new SpringApplicationBuilder(StringToStringStreamListener.class, CustomConverters.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.default.contentType=foo/bar", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
TargetDestination target = context.getBean(TargetDestination.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
Message<byte[]> outputMessage = target.receive();
assertNotNull(outputMessage);
assertEquals("FooBarMessageConverter", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
assertEquals(MimeType.valueOf("foo/bar"), outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
}
//Failure tests
@Test
public void _jsonToPojoWrongDefaultContentTypeProperty() {
ApplicationContext context = new SpringApplicationBuilder(PojoToPojoStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.default.contentType=text/plain", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
SpringIntegrationChannelBinder binder = context.getBean(SpringIntegrationChannelBinder.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
assertTrue(binder.getLastError().getPayload() instanceof MessageConversionException);
}
@Test
public void _toStringDefaultContentTypePropertyUnknownContentType() {
ApplicationContext context = new SpringApplicationBuilder(StringToStringStreamListener.class)
.web(WebApplicationType.NONE)
.run("--spring.cloud.stream.default.contentType=foo/bar", "--spring.jmx.enabled=false");
SourceDestination source = context.getBean(SourceDestination.class);
SpringIntegrationChannelBinder binder = context.getBean(SpringIntegrationChannelBinder.class);
String jsonPayload = "{\"name\":\"oleg\"}";
source.send(new GenericMessage<byte[]>(jsonPayload.getBytes()));
assertTrue(binder.getLastError().getPayload() instanceof MessageConversionException);
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class TextInJsonOutListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public Message<String> echo(String value) {
return MessageBuilder.withPayload(value).setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.APPLICATION_JSON).build();
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class PojoToPojoStreamListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public Person echo(Person value) {
return value;
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class PojoToStringStreamListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public String echo(Person value) {
return value.toString();
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class PojoToByteArrayStreamListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public byte[] echo(Person value) {
return value.toString().getBytes(StandardCharsets.UTF_8);
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class ByteArrayToPojoStreamListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public Person echo(byte[] value) throws Exception {
ObjectMapper mapper = new ObjectMapper();
return mapper.readValue(value, Person.class);
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class StringToPojoStreamListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public Person echo(String value) throws Exception {
ObjectMapper mapper = new ObjectMapper();
return mapper.readValue(value, Person.class);
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class ByteArrayToByteArrayStreamListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public byte[] echo(byte[] value) {
return value;
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class StringToStringStreamListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public String echo(String value) {
return value;
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class PojoMessageToStringMessageStreamListener {
@StreamListener(Processor.INPUT)
@SendTo(Processor.OUTPUT)
public Message<String> echo(Message<Person> value) {
return MessageBuilder.withPayload(value.getPayload().toString()).setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build();
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class PojoMessageToStringMessageServiceActivator {
@ServiceActivator(inputChannel=Processor.INPUT, outputChannel=Processor.OUTPUT)
public Message<String> echo(Message<Person> value) {
return MessageBuilder.withPayload(value.getPayload().toString()).setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build();
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class StringMessageToStringMessageStreamListener {
@ServiceActivator(inputChannel=Processor.INPUT, outputChannel=Processor.OUTPUT)
public Message<String> echo(Message<String> value) throws Exception {
ObjectMapper mapper = new ObjectMapper();
Person person = mapper.readValue(value.getPayload(), Person.class);
return MessageBuilder.withPayload(person.toString()).setHeader(MessageHeaders.CONTENT_TYPE, MimeTypeUtils.TEXT_PLAIN).build();
}
}
@EnableBinding(Processor.class)
@Import(SpringIntegrationBinderConfiguration.class)
public static class ByteArrayMessageToStringJsonMessageStreamListener {
@ServiceActivator(inputChannel=Processor.INPUT, outputChannel=Processor.OUTPUT)
public Message<String> echo(Message<byte[]> value) throws Exception {
ObjectMapper mapper = new ObjectMapper();
Person person = mapper.readValue(value.getPayload(), Person.class);
person.setName("bob");
String json = mapper.writeValueAsString(person);
return MessageBuilder.withPayload(json).build();
}
}
public static class Person {
private String name;
public Person() {
this(null);
}
public Person(String name) {
this.name = name;
}
public String getName() {
return name;
}
public void setName(String name) {
this.name = name;
}
public String toString() {
return name;
}
}
@Configuration
public static class CustomConverters {
@Bean
@StreamMessageConverter
public FooBarMessageConverter fooBarMessageConverter() {
return new FooBarMessageConverter(MimeType.valueOf("foo/bar"));
}
@Bean
@StreamMessageConverter
public AlwaysStringKryoMessageConverter kryoOverrideMessageConverter() {
return new AlwaysStringKryoMessageConverter(MimeType.valueOf("application/x-java-object"));
}
/**
* Even though this MessageConverter has nothing to do with Kryo it still shows how Kryo
* conversion can be customized/overriden since it simply overriding a converter for
* contentType 'application/x-java-object'
*
*/
public static class AlwaysStringKryoMessageConverter extends AbstractMessageConverter {
public AlwaysStringKryoMessageConverter(MimeType supportedMimeType) {
super(supportedMimeType);
}
@Override
protected boolean supports(Class<?> clazz) {
return clazz == null || String.class.isAssignableFrom(clazz);
}
protected Object convertFromInternal(
Message<?> message, Class<?> targetClass, @Nullable Object conversionHint) {
return this.getClass().getSimpleName();
}
protected Object convertToInternal(
Object payload, @Nullable MessageHeaders headers, @Nullable Object conversionHint) {
return ((String)payload).getBytes(StandardCharsets.UTF_8);
}
}
public static class FooBarMessageConverter extends AbstractMessageConverter {
protected FooBarMessageConverter(MimeType supportedMimeType) {
super(supportedMimeType);
}
@Override
protected boolean supports(Class<?> clazz) {
return clazz != null && String.class.isAssignableFrom(clazz);
}
protected Object convertFromInternal(
Message<?> message, Class<?> targetClass, @Nullable Object conversionHint) {
return this.getClass().getSimpleName();
}
protected Object convertToInternal(
Object payload, @Nullable MessageHeaders headers, @Nullable Object conversionHint) {
return ((String)payload).getBytes(StandardCharsets.UTF_8);
}
}
}
}

View File

@@ -47,7 +47,7 @@ import static org.junit.Assert.fail;
public class MessageConverterConfigurerTests {
@Test
//@Test
public void testConfigureOutputChannelWithBadContentType() {
BindingServiceProperties props = new BindingServiceProperties();
BindingProperties bindingProps = new BindingProperties();