GH-3089: Add AmqpInGateway.replyHeadersMappedLast (#3091)

* GH-3089: Add AmqpInGateway.replyHeadersMappedLast

Fixes https://github.com/spring-projects/spring-integration/issues/3089

In some use-case we would like to control when headers from SI message
should be populated into an AMQP message.
One of the use-case is like a `SimpleMessageConverter` and its `plain/text`
for the String reply, meanwhile we know that this content is an
`application/json`.
So, with a new `replyHeadersMappedLast` we can override the mentioned
`content-type` header, populated by the `MessageConverter` with an
actual value from the message headers populated in the flow upstream

* Introduce an `AmqpInboundGateway.replyHeadersMappedLast`; expose it
on the DSL and XML level
* Use newly introduced `MappingUtils.mapReplyMessage()`
* Optimize `DefaultAmqpHeaderMapper` to not parse JSON headers at all
when `JsonHeaders.TYPE_ID` is already present (e.g. `MessageConverter`
result)
* Also skip `JsonHeaders` when we `populateUserDefinedHeader()`

**Cherry-pick to 5.1.x**

* * Fix language and package typos
* Add missed `@param` in JavaDoc of the `AmqpBaseInboundGatewaySpec.batchingStrategy()`
* Extract a `RabbitTemplate` `MessageConverter` to use for reply messages
conversion - pursue a backward compatibility
This commit is contained in:
Artem Bilan
2019-10-31 16:25:40 -04:00
committed by Gary Russell
parent 315fafdaf2
commit 54de7a2209
11 changed files with 248 additions and 100 deletions

View File

@@ -30,6 +30,7 @@ import org.springframework.util.StringUtils;
* @author Mark Fisher
* @author Gary Russell
* @author Artem Bilan
*
* @since 2.1
*/
public class AmqpInboundGatewayParser extends AbstractAmqpInboundAdapterParser {
@@ -48,6 +49,8 @@ public class AmqpInboundGatewayParser extends AbstractAmqpInboundAdapterParser {
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "request-channel");
IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "reply-channel");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "default-reply-to");
IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "reply-headers-last",
"replyHeadersMappedLast");
}
}

View File

@@ -138,4 +138,29 @@ public class AmqpBaseInboundGatewaySpec<S extends AmqpBaseInboundGatewaySpec<S>>
return _this();
}
/**
* Set to true to bind the source message in the headers.
* @param bindSourceMessage true to bind.
* @return the spec.
* @since 5.1.9
* @see AmqpInboundGateway#setBindSourceMessage(boolean)
*/
public S bindSourceMessage(boolean bindSourceMessage) {
this.target.setBindSourceMessage(bindSourceMessage);
return _this();
}
/**
* When mapping headers for the outbound (reply) message, determine whether the headers are
* mapped before the message is converted, or afterwards.
* @param replyHeadersMappedLast true if reply headers are mapped after conversion.
* @return the spec.
* @since 5.1.9
* @see AmqpInboundGateway#setReplyHeadersMappedLast(boolean)
*/
public S replyHeadersMappedLast(boolean replyHeadersMappedLast) {
this.target.setReplyHeadersMappedLast(replyHeadersMappedLast);
return _this();
}
}

View File

@@ -23,8 +23,6 @@ import org.springframework.amqp.core.AcknowledgeMode;
import org.springframework.amqp.core.Address;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.core.MessagePostProcessor;
import org.springframework.amqp.core.MessageProperties;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer;
import org.springframework.amqp.rabbit.listener.api.ChannelAwareMessageListener;
@@ -38,6 +36,7 @@ import org.springframework.integration.amqp.support.AmqpHeaderMapper;
import org.springframework.integration.amqp.support.AmqpMessageHeaderErrorMessageStrategy;
import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper;
import org.springframework.integration.amqp.support.EndpointUtils;
import org.springframework.integration.amqp.support.MappingUtils;
import org.springframework.integration.gateway.MessagingGatewaySupport;
import org.springframework.integration.support.ErrorMessageUtils;
import org.springframework.messaging.MessageChannel;
@@ -45,7 +44,6 @@ import org.springframework.retry.RecoveryCallback;
import org.springframework.retry.support.RetrySynchronizationManager;
import org.springframework.retry.support.RetryTemplate;
import org.springframework.util.Assert;
import org.springframework.util.StringUtils;
import com.rabbitmq.client.Channel;
@@ -71,9 +69,11 @@ public class AmqpInboundGateway extends MessagingGatewaySupport {
private final boolean amqpTemplateExplicitlySet;
private volatile MessageConverter amqpMessageConverter = new SimpleMessageConverter();
private MessageConverter amqpMessageConverter = new SimpleMessageConverter();
private volatile AmqpHeaderMapper headerMapper = DefaultAmqpHeaderMapper.inboundMapper();
private MessageConverter templateMessageConverter = this.amqpMessageConverter;
private AmqpHeaderMapper headerMapper = DefaultAmqpHeaderMapper.inboundMapper();
private Address defaultReplyTo;
@@ -83,6 +83,8 @@ public class AmqpInboundGateway extends MessagingGatewaySupport {
private boolean bindSourceMessage;
private boolean replyHeadersMappedLast;
public AmqpInboundGateway(AbstractMessageListenerContainer listenerContainer) {
this(listenerContainer, new RabbitTemplate(listenerContainer.getConnectionFactory()), false);
}
@@ -110,6 +112,9 @@ public class AmqpInboundGateway extends MessagingGatewaySupport {
this.messageListenerContainer.setAutoStartup(false);
this.amqpTemplate = amqpTemplate;
this.amqpTemplateExplicitlySet = amqpTemplateExplicitlySet;
if (this.amqpTemplateExplicitlySet && this.amqpTemplate instanceof RabbitTemplate) {
this.templateMessageConverter = ((RabbitTemplate) this.amqpTemplate).getMessageConverter();
}
setErrorMessageStrategy(new AmqpMessageHeaderErrorMessageStrategy());
}
@@ -125,6 +130,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport {
this.amqpMessageConverter = messageConverter;
if (!this.amqpTemplateExplicitlySet) {
((RabbitTemplate) this.amqpTemplate).setMessageConverter(messageConverter);
this.templateMessageConverter = messageConverter;
}
}
@@ -187,6 +193,24 @@ public class AmqpInboundGateway extends MessagingGatewaySupport {
this.bindSourceMessage = bindSourceMessage;
}
/**
* When mapping headers for the outbound (reply) message, determine whether the headers are
* mapped before the message is converted, or afterwards. This only affects headers
* that might be added by the message converter. When false, the converter's headers
* win; when true, any headers added by the converter will be overridden (if the
* source message has a header that maps to those headers). You might wish to set this
* to true, for example, when using a
* {@link org.springframework.amqp.support.converter.SimpleMessageConverter} with a
* String payload that contains json; the converter will set the content type to
* {@code text/plain} which can be overridden to {@code application/json} by setting
* the {@link AmqpHeaders#CONTENT_TYPE} message header. Default: false.
* @param replyHeadersMappedLast true if reply headers are mapped after conversion.
* @since 5.1.9
*/
public void setReplyHeadersMappedLast(boolean replyHeadersMappedLast) {
this.replyHeadersMappedLast = replyHeadersMappedLast;
}
@Override
public String getComponentType() {
return "amqp:inbound-gateway";
@@ -331,7 +355,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport {
private void process(Message message, org.springframework.messaging.Message<Object> messagingMessage) {
setAttributesIfNecessary(message, messagingMessage);
final org.springframework.messaging.Message<?> reply = sendAndReceiveMessage(messagingMessage);
org.springframework.messaging.Message<?> reply = sendAndReceiveMessage(messagingMessage);
if (reply != null) {
Address replyTo;
String replyToProperty = message.getMessageProperties().getReplyTo();
@@ -342,30 +366,15 @@ public class AmqpInboundGateway extends MessagingGatewaySupport {
replyTo = AmqpInboundGateway.this.defaultReplyTo;
}
MessagePostProcessor messagePostProcessor =
message1 -> {
MessageProperties messageProperties = message1.getMessageProperties();
String contentEncoding = messageProperties.getContentEncoding();
long contentLength = messageProperties.getContentLength();
String contentType = messageProperties.getContentType();
AmqpInboundGateway.this.headerMapper.fromHeadersToReply(reply.getHeaders(),
messageProperties);
// clear the replyTo from the original message since we are using it now
messageProperties.setReplyTo(null);
// reset the content-* properties as determined by the MessageConverter
if (StringUtils.hasText(contentEncoding)) {
messageProperties.setContentEncoding(contentEncoding);
}
messageProperties.setContentLength(contentLength);
if (contentType != null) {
messageProperties.setContentType(contentType);
}
return message1;
};
org.springframework.amqp.core.Message amqpMessage =
MappingUtils.mapReplyMessage(reply, AmqpInboundGateway.this.templateMessageConverter,
AmqpInboundGateway.this.headerMapper,
message.getMessageProperties().getReceivedDeliveryMode(),
AmqpInboundGateway.this.replyHeadersMappedLast);
if (replyTo != null) {
AmqpInboundGateway.this.amqpTemplate.convertAndSend(replyTo.getExchangeName(),
replyTo.getRoutingKey(), reply.getPayload(), messagePostProcessor);
AmqpInboundGateway.this.amqpTemplate.send(replyTo.getExchangeName(), replyTo.getRoutingKey(),
amqpMessage);
}
else {
if (!AmqpInboundGateway.this.amqpTemplateExplicitlySet) {
@@ -373,8 +382,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport {
"and the `defaultReplyTo` hasn't been configured.");
}
else {
AmqpInboundGateway.this.amqpTemplate.convertAndSend(reply.getPayload(),
messagePostProcessor);
AmqpInboundGateway.this.amqpTemplate.send(amqpMessage);
}
}
}

View File

@@ -105,7 +105,7 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
*/
@Override
protected Map<String, Object> extractStandardHeaders(MessageProperties amqpMessageProperties) {
Map<String, Object> headers = new HashMap<String, Object>();
Map<String, Object> headers = new HashMap<>();
try {
String appId = amqpMessageProperties.getAppId();
if (StringUtils.hasText(appId)) {
@@ -325,24 +325,23 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
amqpMessageProperties.setUserId(userId);
}
Map<String, String> jsonHeaders = new HashMap<String, String>();
for (String jsonHeader : JsonHeaders.HEADERS) {
Object value = getHeaderIfAvailable(headers, jsonHeader, Object.class);
if (value != null) {
headers.remove(jsonHeader);
if (value instanceof Class<?>) {
value = ((Class<?>) value).getName();
}
jsonHeaders.put(jsonHeader.replaceFirst(JsonHeaders.PREFIX, ""), value.toString());
}
}
/*
* If the MessageProperties already contains JsonHeaders, don't overwrite them here because they were
* set up by a message converter.
*/
if (!amqpMessageProperties.getHeaders().containsKey(JsonHeaders.TYPE_ID.replaceFirst(JsonHeaders.PREFIX, ""))) {
Map<String, String> jsonHeaders = new HashMap<>();
for (String jsonHeader : JsonHeaders.HEADERS) {
Object value = getHeaderIfAvailable(headers, jsonHeader, Object.class);
if (value != null) {
headers.remove(jsonHeader);
if (value instanceof Class<?>) {
value = ((Class<?>) value).getName();
}
jsonHeaders.put(jsonHeader.replaceFirst(JsonHeaders.PREFIX, ""), value.toString());
}
}
amqpMessageProperties.getHeaders().putAll(jsonHeaders);
}
@@ -361,8 +360,9 @@ public class DefaultAmqpHeaderMapper extends AbstractHeaderMapper<MessagePropert
MessageProperties amqpMessageProperties) {
// do not overwrite an existing header with the same key
// TODO: do we need to expose a boolean 'overwrite' flag?
if (!amqpMessageProperties.getHeaders().containsKey(headerName)
&& !AmqpHeaders.CONTENT_TYPE.equals(headerName)) {
if (!amqpMessageProperties.getHeaders().containsKey(headerName) &&
!AmqpHeaders.CONTENT_TYPE.equals(headerName) &&
!headerName.startsWith(JsonHeaders.PREFIX)) {
amqpMessageProperties.setHeader(headerName, headerValue);
}
}

View File

@@ -30,6 +30,8 @@ import org.springframework.util.MimeType;
* Utility methods used during message mapping.
*
* @author Gary Russell
* @author Artem Bilan
*
* @since 4.3
*
*/
@@ -40,7 +42,7 @@ public final class MappingUtils {
}
/**
* Map an o.s.Message to an o.s.a.core.Message. When using a
* Map an o.s.m.Message to an o.s.a.core.Message. When using a
* {@link ContentTypeDelegatingMessageConverter}, {@link AmqpHeaders#CONTENT_TYPE} and
* {@link MessageHeaders#CONTENT_TYPE} will be used for the selection, with the AMQP
* header taking precedence.
@@ -54,25 +56,64 @@ public final class MappingUtils {
public static org.springframework.amqp.core.Message mapMessage(Message<?> requestMessage,
MessageConverter converter, AmqpHeaderMapper headerMapper, MessageDeliveryMode defaultDeliveryMode,
boolean headersMappedLast) {
return doMapMessage(requestMessage, converter, headerMapper, defaultDeliveryMode, headersMappedLast, false);
}
/**
* Map a reply o.s.m.Message to an o.s.a.core.Message. When using a
* {@link ContentTypeDelegatingMessageConverter}, {@link AmqpHeaders#CONTENT_TYPE} and
* {@link MessageHeaders#CONTENT_TYPE} will be used for the selection, with the AMQP
* header taking precedence.
* @param replyMessage the reply message.
* @param converter the message converter to use.
* @param headerMapper the header mapper to use.
* @param defaultDeliveryMode the default delivery mode.
* @param headersMappedLast true if headers are mapped after conversion.
* @return the mapped Message.
* @since 5.1.9
*/
public static org.springframework.amqp.core.Message mapReplyMessage(Message<?> replyMessage,
MessageConverter converter, AmqpHeaderMapper headerMapper, MessageDeliveryMode defaultDeliveryMode,
boolean headersMappedLast) {
return doMapMessage(replyMessage, converter, headerMapper, defaultDeliveryMode, headersMappedLast, true);
}
private static org.springframework.amqp.core.Message doMapMessage(Message<?> message,
MessageConverter converter, AmqpHeaderMapper headerMapper, MessageDeliveryMode defaultDeliveryMode,
boolean headersMappedLast, boolean reply) {
MessageProperties amqpMessageProperties = new MessageProperties();
org.springframework.amqp.core.Message amqpMessage;
if (!headersMappedLast) {
headerMapper.fromHeadersToRequest(requestMessage.getHeaders(), amqpMessageProperties);
mapHeaders(message.getHeaders(), amqpMessageProperties, headerMapper, reply);
}
if (converter instanceof ContentTypeDelegatingMessageConverter && headersMappedLast) {
String contentType = contentTypeAsString(requestMessage.getHeaders());
String contentType = contentTypeAsString(message.getHeaders());
if (contentType != null) {
amqpMessageProperties.setContentType(contentType);
}
}
amqpMessage = converter.toMessage(requestMessage.getPayload(), amqpMessageProperties);
amqpMessage = converter.toMessage(message.getPayload(), amqpMessageProperties);
if (headersMappedLast) {
headerMapper.fromHeadersToRequest(requestMessage.getHeaders(), amqpMessageProperties);
mapHeaders(message.getHeaders(), amqpMessageProperties, headerMapper, reply);
}
checkDeliveryMode(requestMessage, amqpMessageProperties, defaultDeliveryMode);
checkDeliveryMode(message, amqpMessageProperties, defaultDeliveryMode);
return amqpMessage;
}
private static void mapHeaders(MessageHeaders messageHeaders, MessageProperties amqpMessageProperties,
AmqpHeaderMapper headerMapper, boolean reply) {
if (reply) {
headerMapper.fromHeadersToReply(messageHeaders, amqpMessageProperties);
}
else {
headerMapper.fromHeadersToRequest(messageHeaders, amqpMessageProperties);
}
}
private static String contentTypeAsString(MessageHeaders headers) {
Object contentType = headers.get(AmqpHeaders.CONTENT_TYPE);
if (contentType instanceof MimeType) {

View File

@@ -239,6 +239,18 @@
</xsd:documentation>
</xsd:annotation>
</xsd:attribute>
<xsd:attribute name="reply-headers-last">
<xsd:annotation>
<xsd:documentation>
Whether reply headers are mapped before or after conversion from a messaging Message to
a spring amqp Message. Set to true, for example, if you wish to override the
contentType header set by the converter.
</xsd:documentation>
</xsd:annotation>
<xsd:simpleType>
<xsd:union memberTypes="xsd:boolean xsd:string" />
</xsd:simpleType>
</xsd:attribute>
</xsd:extension>
</xsd:complexContent>
</xsd:complexType>

View File

@@ -1,38 +1,41 @@
<?xml version="1.0" encoding="UTF-8"?>
<beans xmlns="http://www.springframework.org/schema/beans"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:si-amqp="http://www.springframework.org/schema/integration/amqp"
xmlns:int="http://www.springframework.org/schema/integration"
xsi:schemaLocation="http://www.springframework.org/schema/integration/amqp https://www.springframework.org/schema/integration/amqp/spring-integration-amqp.xsd
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xmlns:si-amqp="http://www.springframework.org/schema/integration/amqp"
xmlns:int="http://www.springframework.org/schema/integration"
xsi:schemaLocation="http://www.springframework.org/schema/integration/amqp https://www.springframework.org/schema/integration/amqp/spring-integration-amqp.xsd
http://www.springframework.org/schema/integration https://www.springframework.org/schema/integration/spring-integration.xsd
http://www.springframework.org/schema/beans https://www.springframework.org/schema/beans/spring-beans.xsd">
<int:channel id="requests"/>
<si-amqp:inbound-gateway id="gateway" request-channel="requests" queue-names="test" reply-timeout="1234"
connection-factory="rabbitConnectionFactory" message-converter="testConverter"/>
connection-factory="rabbitConnectionFactory" message-converter="testConverter"/>
<bean id="rabbitConnectionFactory" class="org.mockito.Mockito" factory-method="mock">
<constructor-arg value="org.springframework.amqp.rabbit.connection.ConnectionFactory"/>
</bean>
<bean id="testConverter" class="org.springframework.integration.amqp.config.AmqpInboundGatewayParserTests$TestConverter"/>
<bean id="testConverter"
class="org.springframework.integration.amqp.config.AmqpInboundGatewayParserTests$TestConverter"/>
<bean id="amqpTemplate" class="org.mockito.Mockito" factory-method="mock">
<constructor-arg value="org.springframework.amqp.core.AmqpTemplate"/>
</bean>
<si-amqp:inbound-gateway id="autoStartFalseGateway" request-channel="requests" queue-names="test"
connection-factory="rabbitConnectionFactory" message-converter="testConverter"
missing-queues-fatal="false"
amqp-template="amqpTemplate"
default-reply-to="fooExchange/barRoutingKey"
auto-startup="false" phase="123"/>
connection-factory="rabbitConnectionFactory" message-converter="testConverter"
missing-queues-fatal="false"
amqp-template="amqpTemplate"
default-reply-to="fooExchange/barRoutingKey"
auto-startup="false" phase="123"/>
<si-amqp:inbound-gateway id="withHeaderMapper" request-channel="requestChannel" queue-names="inboundchanneladapter.test.2"
auto-startup="false" phase="123"
mapped-request-headers="foo*, STANDARD_REQUEST_HEADERS"
mapped-reply-headers="bar*"/>
<si-amqp:inbound-gateway id="withHeaderMapper" request-channel="requestChannel"
queue-names="inboundchanneladapter.test.2"
auto-startup="false" phase="123"
mapped-request-headers="foo*, STANDARD_REQUEST_HEADERS"
mapped-reply-headers="bar*"
reply-headers-last="true"/>
<int:channel id="requestChannel"/>

View File

@@ -16,6 +16,7 @@
package org.springframework.integration.amqp.config;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertSame;
@@ -107,7 +108,7 @@ public class AmqpInboundGatewayParserTests {
});
final AmqpInboundGateway gateway = context.getBean("withHeaderMapper", AmqpInboundGateway.class);
assertThat(TestUtils.getPropertyValue(gateway, "replyHeadersMappedLast", Boolean.class)).isTrue();
Field amqpTemplateField = ReflectionUtils.findField(AmqpInboundGateway.class, "amqpTemplate");
amqpTemplateField.setAccessible(true);
RabbitTemplate amqpTemplate = TestUtils.getPropertyValue(gateway, "amqpTemplate", RabbitTemplate.class);

View File

@@ -16,6 +16,7 @@
package org.springframework.integration.amqp.dsl;
import static org.assertj.core.api.Assertions.assertThat;
import static org.hamcrest.Matchers.instanceOf;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertNotNull;
@@ -24,6 +25,9 @@ import static org.junit.Assert.assertSame;
import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
import java.util.Collections;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.atomic.AtomicReference;
import org.junit.jupiter.api.AfterAll;
@@ -43,6 +47,7 @@ import org.springframework.amqp.rabbit.junit.RabbitAvailableCondition;
import org.springframework.amqp.rabbit.listener.DirectMessageListenerContainer;
import org.springframework.amqp.rabbit.listener.SimpleMessageListenerContainer;
import org.springframework.amqp.rabbit.listener.exception.ListenerExecutionFailedException;
import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter;
import org.springframework.amqp.support.converter.MessageConversionException;
import org.springframework.amqp.support.converter.SimpleMessageConverter;
import org.springframework.beans.factory.annotation.Autowired;
@@ -61,6 +66,8 @@ import org.springframework.integration.config.EnableIntegration;
import org.springframework.integration.dsl.IntegrationFlow;
import org.springframework.integration.dsl.IntegrationFlowBuilder;
import org.springframework.integration.dsl.IntegrationFlows;
import org.springframework.integration.dsl.Transformers;
import org.springframework.integration.dsl.context.IntegrationFlowContext;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.support.StringObjectMapBuilder;
import org.springframework.integration.test.util.TestUtils;
@@ -78,14 +85,17 @@ import org.springframework.test.context.junit.jupiter.SpringJUnitConfig;
*/
@SpringJUnitConfig
@RabbitAvailable(queues = { "amqpOutboundInput", "amqpReplyChannel", "asyncReplies",
"defaultReplyTo", "si.dsl.test", "si.dsl.exception.test.dlq",
"si.dsl.conv.exception.test.dlq", "testTemplateChannelTransacted" })
"defaultReplyTo", "si.dsl.test", "si.dsl.exception.test.dlq",
"si.dsl.conv.exception.test.dlq", "testTemplateChannelTransacted" })
@DirtiesContext
public class AmqpTests {
@Autowired
private ConnectionFactory rabbitConnectionFactory;
@Autowired
private IntegrationFlowContext integrationFlowContext;
@Autowired
private AmqpTemplate amqpTemplate;
@@ -93,6 +103,10 @@ public class AmqpTests {
@Qualifier("queue")
private Queue amqpQueue;
@Autowired
@Qualifier("queue2")
private Queue amqpQueue2;
@Autowired
private AmqpInboundGateway amqpInboundGateway;
@@ -232,6 +246,38 @@ public class AmqpTests {
assertTrue(TestUtils.getPropertyValue(this.unitChannel, "extractPayload", Boolean.class));
}
@Test
void testContentTypeOverrideWithReplyHeadersMappedLast() {
IntegrationFlow testFlow =
IntegrationFlows
.from(Amqp.inboundGateway(this.rabbitConnectionFactory, this.amqpQueue2)
.replyHeadersMappedLast(true))
.transform(Transformers.fromJson())
.enrich((enricher) -> enricher.property("REPLY_KEY", "REPLY_VALUE"))
.transform(Transformers.toJson())
.get();
IntegrationFlowContext.IntegrationFlowRegistration registration =
this.integrationFlowContext.registration(testFlow).register();
RabbitTemplate rabbitTemplate = new RabbitTemplate(this.rabbitConnectionFactory);
rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter());
Object result = rabbitTemplate.convertSendAndReceive(this.amqpQueue2.getName(),
new HashMap<>(Collections.singletonMap("TEST_KEY", "TEST_VALUE")));
assertThat(result).isInstanceOf(Map.class);
@SuppressWarnings("unchecked")
Map<String, String> resultMap = (Map<String, String>) result;
assertThat(resultMap)
.containsEntry("TEST_KEY", "TEST_VALUE")
.containsEntry("REPLY_KEY", "REPLY_VALUE");
registration.destroy();
}
@Configuration
@EnableIntegration
public static class ContextConfiguration {
@@ -256,6 +302,11 @@ public class AmqpTests {
return new AnonymousQueue();
}
@Bean
public Queue queue2() {
return new AnonymousQueue();
}
@Bean
public Queue defaultReplyTo() {
return new Queue("defaultReplyTo");
@@ -267,9 +318,9 @@ public class AmqpTests {
.from(Amqp.inboundGateway(rabbitConnectionFactory, amqpTemplate, queue())
.id("amqpInboundGateway")
.configureContainer(c -> c
.id("amqpInboundGatewayContainer")
.recoveryInterval(5000)
.concurrentConsumers(1))
.id("amqpInboundGatewayContainer")
.recoveryInterval(5000)
.concurrentConsumers(1))
.defaultReplyTo(defaultReplyTo().getName()))
.transform("hello "::concat)
.transform(String.class, String::toUpperCase)
@@ -282,8 +333,8 @@ public class AmqpTests {
.from(Amqp.inboundGateway(new DirectMessageListenerContainer())
.id("amqpInboundGateway")
.configureContainer(c -> c
.recoveryInterval(5000)
.consumersPerQueue(1))
.recoveryInterval(5000)
.consumersPerQueue(1))
.defaultReplyTo(defaultReplyTo().getName()))
.transform("hello "::concat)
.transform(String.class, String::toUpperCase)
@@ -332,9 +383,9 @@ public class AmqpTests {
public Queue exQueue() {
return new Queue("si.dsl.exception.test", true, false, false,
new StringObjectMapBuilder()
.put("x-dead-letter-exchange", "")
.put("x-dead-letter-routing-key", exDLQ().getName())
.get());
.put("x-dead-letter-exchange", "")
.put("x-dead-letter-routing-key", exDLQ().getName())
.get());
}
@Bean
@@ -345,8 +396,8 @@ public class AmqpTests {
@Bean
public IntegrationFlow inboundWithExceptionFlow(ConnectionFactory cf) {
return IntegrationFlows.from(Amqp.inboundAdapter(cf, exQueue())
.configureContainer(c -> c.defaultRequeueRejected(false))
.errorChannel("errors.input"))
.configureContainer(c -> c.defaultRequeueRejected(false))
.errorChannel("errors.input"))
.handle(m -> {
throw new RuntimeException("fail");
})
@@ -356,22 +407,22 @@ public class AmqpTests {
@Bean
public IntegrationFlow errors() {
return f -> f.handle(m -> {
raw().set(m.getHeaders().get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE,
org.springframework.amqp.core.Message.class));
if (m.getPayload() instanceof ListenerExecutionFailedException) {
lefe().set((ListenerExecutionFailedException) m.getPayload());
}
throw (RuntimeException) m.getPayload();
});
raw().set(m.getHeaders().get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE,
org.springframework.amqp.core.Message.class));
if (m.getPayload() instanceof ListenerExecutionFailedException) {
lefe().set((ListenerExecutionFailedException) m.getPayload());
}
throw (RuntimeException) m.getPayload();
});
}
@Bean
public Queue exConvQueue() {
return new Queue("si.dsl.conv.exception.test", true, false, false,
new StringObjectMapBuilder()
.put("x-dead-letter-exchange", "")
.put("x-dead-letter-routing-key", exConvDLQ().getName())
.get());
.put("x-dead-letter-exchange", "")
.put("x-dead-letter-routing-key", exConvDLQ().getName())
.get());
}
@Bean
@@ -382,17 +433,17 @@ public class AmqpTests {
@Bean
public IntegrationFlow inboundWithConvExceptionFlow(ConnectionFactory cf) {
return IntegrationFlows.from(Amqp.inboundAdapter(cf, exConvQueue())
.configureContainer(c -> c.defaultRequeueRejected(false))
.messageConverter(new SimpleMessageConverter() {
.configureContainer(c -> c.defaultRequeueRejected(false))
.messageConverter(new SimpleMessageConverter() {
@Override
public Object fromMessage(org.springframework.amqp.core.Message message)
throws MessageConversionException {
throw new MessageConversionException("fail");
}
@Override
public Object fromMessage(org.springframework.amqp.core.Message message)
throws MessageConversionException {
throw new MessageConversionException("fail");
}
})
.errorChannel("errors.input"))
})
.errorChannel("errors.input"))
.get();
}
@@ -410,7 +461,7 @@ public class AmqpTests {
public IntegrationFlow amqpAsyncOutboundFlow(AsyncRabbitTemplate asyncRabbitTemplate) {
return f -> f
.handle(Amqp.asyncOutboundGateway(asyncRabbitTemplate)
.routingKeyFunction(m -> queue().getName()),
.routingKeyFunction(m -> queue().getName()),
e -> e.id("asyncOutboundGateway"));
}

View File

@@ -219,6 +219,7 @@ public class InboundEndpointTests {
gateway.setRequestChannel(channel);
gateway.setBeanFactory(mock(BeanFactory.class));
gateway.setDefaultReplyTo("foo");
gateway.setReplyHeadersMappedLast(true);
gateway.afterPropertiesSet();

View File

@@ -1144,6 +1144,9 @@ The `ObjectToJsonTransformer` does exactly that (by default).
There is now a property called `headersMappedLast` on the outbound channel adapter and gateway (as well as on AMQP-backed channels).
Setting this to `true` restores the behavior of overwriting the property added by the converter.
Starting with version 5.1.9, a similar `replyHeadersMappedLast` is provided for the `AmqpInboundGateway` when we produce a reply and would like to override headers populated by the converter.
See its JavaDocs for more information.
====
[[amqp-user-id]]