GH-1554 Introduced flag to disable 1.3 content type propagation
Introduced `spring.cloud.stream.propagateOriginalContentType` boolean property on BindingServiceProperties Resolves #1554
This commit is contained in:
@@ -68,10 +68,10 @@ public class TextPlainConversionTest {
|
||||
public void testByteArrayConversionOnOutput() throws Exception {
|
||||
testProcessor.output().send(MessageBuilder.withPayload("Bar".getBytes()).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(received.getPayload()).isEqualTo("Bar");
|
||||
assertThat(received.getPayload()).isEqualTo("Bar".getBytes());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
||||
@@ -148,8 +148,6 @@ public class ContentTypeTests {
|
||||
.build());
|
||||
Message<byte[]> message = (Message<byte[]>) collector
|
||||
.forChannel(source.output()).poll(1, TimeUnit.SECONDS);
|
||||
assertThat(message.getHeaders().get(MessageHeaders.CONTENT_TYPE, MimeType.class)
|
||||
.includes(MimeTypeUtils.IMAGE_JPEG));
|
||||
assertThat(message.getPayload()).isEqualTo(data);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -45,6 +45,7 @@ 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;
|
||||
@@ -307,35 +308,57 @@ public class MessageConverterConfigurer implements MessageChannelAndSourceConfig
|
||||
this.messageConverter = messageConverter;
|
||||
}
|
||||
|
||||
|
||||
@Override
|
||||
public Message<?> doPreSend(Message<?> message, MessageChannel channel) {
|
||||
@SuppressWarnings("deprecation")
|
||||
boolean propagateOriginalContentType =
|
||||
MessageConverterConfigurer.this.bindingServiceProperties.isPropagateOriginalContentType();
|
||||
|
||||
boolean contentTypeHeaderSet = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE);
|
||||
|
||||
// ===== 1.3 backward compatibility code part-1 ===
|
||||
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();
|
||||
String ct = null;
|
||||
String oct = null;
|
||||
if (propagateOriginalContentType) {
|
||||
oct = message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE) ? message.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString() : null;
|
||||
ct = message.getPayload() instanceof String
|
||||
? ct = JavaClassMimeTypeUtils.mimeTypeFromObject(message.getPayload(), ObjectUtils.nullSafeToString(oct)).toString()
|
||||
: oct;
|
||||
}
|
||||
// ===== END 1.3 backward compatibility code part-1 ===
|
||||
|
||||
if (!message.getHeaders().containsKey(MessageHeaders.CONTENT_TYPE)) {
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Object> headersMap = (Map<String, Object>) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders());
|
||||
headersMap.put(MessageHeaders.CONTENT_TYPE, this.mimeType);
|
||||
|
||||
MutableMessageHeaders headers = new MutableMessageHeaders(message.getHeaders());
|
||||
if (!headers.containsKey(MessageHeaders.CONTENT_TYPE)) {
|
||||
headers.put(MessageHeaders.CONTENT_TYPE, this.mimeType);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
Message<byte[]> outboundMessage = message.getPayload() instanceof byte[]
|
||||
? (Message<byte[]>)message : (Message<byte[]>) this.messageConverter.toMessage(message.getPayload(), message.getHeaders());
|
||||
? (Message<byte[]>)message : (Message<byte[]>) this.messageConverter.toMessage(message.getPayload(), headers);
|
||||
if (outboundMessage == null) {
|
||||
throw new IllegalStateException("Failed to convert message: '" + message + "' to outbound message.");
|
||||
}
|
||||
|
||||
/// ===== 1.3 backward compatibility code part-2 ===
|
||||
if (ct != null && !ct.equals(oct) && oct != null) {
|
||||
@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));
|
||||
if (propagateOriginalContentType) {
|
||||
if (ct != null && !ct.equals(oct) && oct != null) {
|
||||
@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));
|
||||
}
|
||||
}
|
||||
else {
|
||||
if (!contentTypeHeaderSet) {
|
||||
@SuppressWarnings("unchecked")
|
||||
Map<String, Object> headersMap = (Map<String, Object>) ReflectionUtils.getField(MessageConverterConfigurer.this.headersField, message.getHeaders());
|
||||
headersMap.remove(MessageHeaders.CONTENT_TYPE);
|
||||
}
|
||||
else {
|
||||
System.out.println();
|
||||
}
|
||||
}
|
||||
// ===== END 1.3 backward compatibility code part-2 ===
|
||||
return outboundMessage;
|
||||
|
||||
@@ -56,6 +56,18 @@ public class BindingServiceProperties implements ApplicationContextAware, Initia
|
||||
|
||||
private static final int DEFAULT_BINDING_RETRY_INTERVAL = 30;
|
||||
|
||||
/**
|
||||
* Setting it to true ensures that the original content-type of the message is propagated
|
||||
* to the outgoing message as `originalContentType` header.
|
||||
*
|
||||
* This deprecated feature primarily exists for backward compatibility
|
||||
* and will not be supported in future versions.
|
||||
*
|
||||
* Default: true
|
||||
*/
|
||||
@Deprecated
|
||||
private boolean propagateOriginalContentType = true;
|
||||
|
||||
/**
|
||||
* The instance id of the application: a number from 0 to instanceCount-1.
|
||||
* Used for partitioning and with Kafka.
|
||||
@@ -284,4 +296,14 @@ public class BindingServiceProperties implements ApplicationContextAware, Initia
|
||||
this.bindings.put(binding, bindingPropertiesTarget);
|
||||
}
|
||||
|
||||
@Deprecated
|
||||
public boolean isPropagateOriginalContentType() {
|
||||
return propagateOriginalContentType;
|
||||
}
|
||||
|
||||
@Deprecated
|
||||
public void setPropagateOriginalContentType(boolean propagateOriginalContentType) {
|
||||
this.propagateOriginalContentType = propagateOriginalContentType;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -61,6 +61,7 @@ import org.springframework.util.MimeTypeUtils;
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.junit.Assert.assertEquals;
|
||||
import static org.junit.Assert.assertNotNull;
|
||||
import static org.junit.Assert.assertNull;
|
||||
import static org.junit.Assert.assertTrue;
|
||||
|
||||
/**
|
||||
@@ -188,7 +189,7 @@ public class ContentTypeTckTests {
|
||||
String jsonPayload = "{\"name\":\"oleg\"}";
|
||||
source.send(new GenericMessage<>(jsonPayload.getBytes()));
|
||||
Message<byte[]> outputMessage = target.receive();
|
||||
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
|
||||
}
|
||||
|
||||
@@ -202,7 +203,7 @@ public class ContentTypeTckTests {
|
||||
String jsonPayload = "{\"name\":\"oleg\"}";
|
||||
source.send(new GenericMessage<>(jsonPayload.getBytes()));
|
||||
Message<byte[]> outputMessage = target.receive();
|
||||
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertEquals("oleg", new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
|
||||
}
|
||||
|
||||
@@ -256,11 +257,10 @@ public class ContentTypeTckTests {
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
String jsonPayload = "{\"name\":\"oleg\"}";
|
||||
//source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, MimeType.valueOf("text/*")).build());
|
||||
source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader("contentType", new MimeType("text", "plain")).build());
|
||||
|
||||
Message<byte[]> outputMessage = target.receive();
|
||||
//assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
|
||||
}
|
||||
|
||||
@@ -272,11 +272,10 @@ public class ContentTypeTckTests {
|
||||
InputDestination source = context.getBean(InputDestination.class);
|
||||
OutputDestination target = context.getBean(OutputDestination.class);
|
||||
String jsonPayload = "{\"name\":\"oleg\"}";
|
||||
//source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader(MessageHeaders.CONTENT_TYPE, MimeType.valueOf("text/*")).build());
|
||||
source.send(MessageBuilder.withPayload(jsonPayload.getBytes()).setHeader("contentType", new MimeType("text")).build());
|
||||
|
||||
Message<byte[]> outputMessage = target.receive();
|
||||
//assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertEquals("text/*", outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE).toString());
|
||||
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
|
||||
}
|
||||
|
||||
@@ -334,7 +333,7 @@ public class ContentTypeTckTests {
|
||||
String jsonPayload = "{\"name\":\"oleg\"}";
|
||||
source.send(new GenericMessage<>(jsonPayload.getBytes()));
|
||||
Message<byte[]> outputMessage = target.receive();
|
||||
assertEquals(MimeTypeUtils.APPLICATION_JSON, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
|
||||
}
|
||||
|
||||
@@ -348,7 +347,7 @@ public class ContentTypeTckTests {
|
||||
String jsonPayload = "{\"name\":\"oleg\"}";
|
||||
source.send(new GenericMessage<>(jsonPayload.getBytes()));
|
||||
Message<byte[]> outputMessage = target.receive();
|
||||
assertEquals(MimeTypeUtils.TEXT_PLAIN, outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertNull( outputMessage.getHeaders().get(MessageHeaders.CONTENT_TYPE));
|
||||
assertEquals(jsonPayload, new String(outputMessage.getPayload(), StandardCharsets.UTF_8));
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user