From 0a1306cd1b11bfbdbd35768815595885b0e1e6ff Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 15 Aug 2017 12:31:52 -0400 Subject: [PATCH] INT-4328: AMQP: Returns/Nacks: Create ErrorMessage JIRA: https://jira.spring.io/browse/INT-4328 Add support for sending `ErrorMessage`s to the return and nack channels. **cherry-pick to 4.3.x, but change default EMS to null (will require minor adjustment to test - set the EMS in `adapterWithReturnsAndErrorMessageStrategy`)** --- .../AmqpOutboundChannelAdapterParser.java | 1 + .../config/AmqpOutboundGatewayParser.java | 1 + .../AbstractAmqpOutboundEndpoint.java | 167 ++++++++++++------ ...AmqpMessageHeaderErrorMessageStrategy.java | 7 +- .../support/NackedAmqpMessageException.java | 58 ++++++ .../support/ReturnedAmqpMessageException.java | 80 +++++++++ .../config/spring-integration-amqp-5.0.xsd | 13 ++ ...boundChannelAdapterParserTests-context.xml | 5 +- ...AmqpOutboundChannelAdapterParserTests.java | 1 + ...AmqpOutboundGatewayParserTests-context.xml | 5 +- .../AmqpOutboundGatewayParserTests.java | 1 + .../outbound/AmqpOutboundEndpointTests.java | 22 +++ .../amqp/outbound/AsyncAmqpGatewayTests.java | 15 +- .../support/DefaultErrorMessageStrategy.java | 3 +- src/reference/asciidoc/amqp.adoc | 42 +++-- 15 files changed, 339 insertions(+), 82 deletions(-) create mode 100644 spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/NackedAmqpMessageException.java create mode 100644 spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/ReturnedAmqpMessageException.java diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java index 684a0b2ad4..30aee77cb4 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParser.java @@ -81,6 +81,7 @@ public class AmqpOutboundChannelAdapterParser extends AbstractOutboundChannelAda IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "confirm-ack-channel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "confirm-nack-channel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "return-channel"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "error-message-strategy"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "delay-expression", "delayExpressionString"); IntegrationNamespaceUtils.setValueIfAttributeDefined(builder, element, "headers-last", "headersMappedLast"); diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java index 14bc54688e..6ae5b81531 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParser.java @@ -94,6 +94,7 @@ public class AmqpOutboundGatewayParser extends AbstractConsumerEndpointParser { IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "confirm-ack-channel"); IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "confirm-nack-channel"); + IntegrationNamespaceUtils.setReferenceIfAttributeDefined(builder, element, "error-message-strategy"); BeanDefinitionBuilder mapperBuilder = BeanDefinitionBuilder .genericBeanDefinition(DefaultAmqpHeaderMapper.class); diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AbstractAmqpOutboundEndpoint.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AbstractAmqpOutboundEndpoint.java index 0e33d54ebb..d398a11caf 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AbstractAmqpOutboundEndpoint.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/outbound/AbstractAmqpOutboundEndpoint.java @@ -30,11 +30,15 @@ import org.springframework.context.Lifecycle; import org.springframework.expression.Expression; import org.springframework.integration.amqp.support.AmqpHeaderMapper; import org.springframework.integration.amqp.support.DefaultAmqpHeaderMapper; +import org.springframework.integration.amqp.support.NackedAmqpMessageException; +import org.springframework.integration.amqp.support.ReturnedAmqpMessageException; import org.springframework.integration.channel.NullChannel; import org.springframework.integration.expression.ValueExpression; import org.springframework.integration.handler.AbstractReplyProducingMessageHandler; import org.springframework.integration.handler.ExpressionEvaluatingMessageProcessor; import org.springframework.integration.support.AbstractIntegrationMessageBuilder; +import org.springframework.integration.support.DefaultErrorMessageStrategy; +import org.springframework.integration.support.ErrorMessageStrategy; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.util.Assert; @@ -50,42 +54,48 @@ import org.springframework.util.StringUtils; public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducingMessageHandler implements Lifecycle { - private volatile String exchangeName; + private String exchangeName; - private volatile String routingKey; + private String routingKey; - private volatile Expression exchangeNameExpression; + private Expression exchangeNameExpression; - private volatile Expression routingKeyExpression; + private Expression routingKeyExpression; - private volatile ExpressionEvaluatingMessageProcessor routingKeyGenerator; + private ExpressionEvaluatingMessageProcessor routingKeyGenerator; - private volatile ExpressionEvaluatingMessageProcessor exchangeNameGenerator; + private ExpressionEvaluatingMessageProcessor exchangeNameGenerator; - private volatile AmqpHeaderMapper headerMapper = DefaultAmqpHeaderMapper.outboundMapper(); + private AmqpHeaderMapper headerMapper = DefaultAmqpHeaderMapper.outboundMapper(); - private volatile Expression confirmCorrelationExpression; + private Expression confirmCorrelationExpression; - private volatile ExpressionEvaluatingMessageProcessor correlationDataGenerator; + private ExpressionEvaluatingMessageProcessor correlationDataGenerator; - private volatile MessageChannel confirmAckChannel; + private MessageChannel confirmAckChannel; - private volatile MessageChannel confirmNackChannel; + private String confirmAckChannelName; - private volatile MessageChannel returnChannel; + private MessageChannel confirmNackChannel; - private volatile MessageDeliveryMode defaultDeliveryMode; + private String confirmNackChannelName; - private volatile boolean lazyConnect = true; + private MessageChannel returnChannel; - private volatile ConnectionFactory connectionFactory; + private MessageDeliveryMode defaultDeliveryMode; - private volatile Expression delayExpression; + private boolean lazyConnect = true; - private volatile ExpressionEvaluatingMessageProcessor delayGenerator; + private ConnectionFactory connectionFactory; + + private Expression delayExpression; + + private ExpressionEvaluatingMessageProcessor delayGenerator; private boolean headersMappedLast; + private ErrorMessageStrategy errorMessageStrategy = new DefaultErrorMessageStrategy(); + private volatile boolean running; public void setHeaderMapper(AmqpHeaderMapper headerMapper) { @@ -180,6 +190,15 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin this.confirmAckChannel = ackChannel; } + /** + * Set the channel name to which acks are send (publisher confirms). + * @param ackChannelName the channel name. + * @since 4.3.12 + */ + public void setConfirmAckChannelName(String ackChannelName) { + this.confirmAckChannelName = ackChannelName; + } + /** * Set the channel to which nacks are send (publisher confirms). * @param nackChannel the channel. @@ -188,6 +207,15 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin this.confirmNackChannel = nackChannel; } + /** + * Set the channel name to which nacks are send (publisher confirms). + * @param nackChannelName the channel name. + * @since 4.3.12 + */ + public void setConfirmNackChannelName(String nackChannelName) { + this.confirmNackChannelName = nackChannelName; + } + /** * Set the channel to which returned messages are sent. * @param returnChannel the channel. @@ -253,6 +281,16 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin } } + /** + * Set the error message strategy to use for returned (or negatively confirmed) + * messages. + * @param errorMessageStrategy the strategy. + * @since 4.3.12 + */ + public void setErrorMessageStrategy(ErrorMessageStrategy errorMessageStrategy) { + this.errorMessageStrategy = errorMessageStrategy; + } + protected final void setConnectionFactory(ConnectionFactory connectionFactory) { this.connectionFactory = connectionFactory; } @@ -294,10 +332,16 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin } protected MessageChannel getConfirmAckChannel() { + if (this.confirmAckChannel == null && this.confirmAckChannelName != null) { + this.confirmAckChannel = getChannelResolver().resolveDestination(confirmAckChannelName); + } return this.confirmAckChannel; } protected MessageChannel getConfirmNackChannel() { + if (this.confirmNackChannel == null && this.confirmNackChannelName != null) { + this.confirmNackChannel = getChannelResolver().resolveDestination(confirmNackChannelName); + } return this.confirmNackChannel; } @@ -347,10 +391,11 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin } else { NullChannel nullChannel = extractTypeIfPossible(this.confirmAckChannel, NullChannel.class); - Assert.state(this.confirmAckChannel == null || nullChannel != null, + Assert.state((this.confirmAckChannel == null || nullChannel != null) && this.confirmAckChannelName == null, "A 'confirmCorrelationExpression' is required when specifying a 'confirmAckChannel'"); nullChannel = extractTypeIfPossible(this.confirmNackChannel, NullChannel.class); - Assert.state(this.confirmNackChannel == null || nullChannel != null, + Assert.state( + (this.confirmNackChannel == null || nullChannel != null) && this.confirmNackChannelName == null, "A 'confirmCorrelationExpression' is required when specifying a 'confirmNackChannel'"); } if (this.delayExpression != null) { @@ -411,16 +456,8 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin protected CorrelationData generateCorrelationData(Message requestMessage) { CorrelationData correlationData = null; if (this.correlationDataGenerator != null) { - Object userCorrelationData = this.correlationDataGenerator.processMessage(requestMessage); - if (userCorrelationData != null) { - if (userCorrelationData instanceof CorrelationData) { - correlationData = (CorrelationData) userCorrelationData; - } - else { - correlationData = new CorrelationDataWrapper(requestMessage - .getHeaders().getId().toString(), userCorrelationData); - } - } + correlationData = new CorrelationDataWrapper(requestMessage.getHeaders().getId().toString(), + this.correlationDataGenerator.processMessage(requestMessage), requestMessage); } return correlationData; } @@ -465,44 +502,55 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin ? this.getMessageBuilderFactory().fromMessage((Message) returnedObject) : this.getMessageBuilderFactory().withPayload(returnedObject); Map headers = getHeaderMapper().toHeadersFromReply(message.getMessageProperties()); - builder.copyHeadersIfAbsent(headers) - .setHeader(AmqpHeaders.RETURN_REPLY_CODE, replyCode) - .setHeader(AmqpHeaders.RETURN_REPLY_TEXT, replyText) - .setHeader(AmqpHeaders.RETURN_EXCHANGE, exchange) - .setHeader(AmqpHeaders.RETURN_ROUTING_KEY, routingKey); - return builder.build(); + if (this.errorMessageStrategy == null) { + builder.copyHeadersIfAbsent(headers) + .setHeader(AmqpHeaders.RETURN_REPLY_CODE, replyCode) + .setHeader(AmqpHeaders.RETURN_REPLY_TEXT, replyText) + .setHeader(AmqpHeaders.RETURN_EXCHANGE, exchange) + .setHeader(AmqpHeaders.RETURN_ROUTING_KEY, routingKey); + } + Message returnedMessage = builder.build(); + if (this.errorMessageStrategy != null) { + returnedMessage = this.errorMessageStrategy.buildErrorMessage(new ReturnedAmqpMessageException( + returnedMessage, message, replyCode, replyText, exchange, routingKey), null); + } + return returnedMessage; } protected void handleConfirm(CorrelationData correlationData, boolean ack, String cause) { - Object userCorrelationData = correlationData; + CorrelationDataWrapper wrapper = (CorrelationDataWrapper) correlationData; if (correlationData == null) { if (logger.isDebugEnabled()) { logger.debug("No correlation data provided for ack: " + ack + " cause:" + cause); } return; } - if (correlationData instanceof CorrelationDataWrapper) { - userCorrelationData = ((CorrelationDataWrapper) correlationData).getUserData(); - } + Object userCorrelationData = wrapper.getUserData(); + Message confirmMessage; + if (this.errorMessageStrategy == null || ack) { + Map headers = new HashMap(); + headers.put(AmqpHeaders.PUBLISH_CONFIRM, ack); + if (!ack && StringUtils.hasText(cause)) { + headers.put(AmqpHeaders.PUBLISH_CONFIRM_NACK_CAUSE, cause); + } - Map headers = new HashMap(); - headers.put(AmqpHeaders.PUBLISH_CONFIRM, ack); - if (!ack && StringUtils.hasText(cause)) { - headers.put(AmqpHeaders.PUBLISH_CONFIRM_NACK_CAUSE, cause); - } + AbstractIntegrationMessageBuilder builder = userCorrelationData instanceof Message + ? this.getMessageBuilderFactory().fromMessage((Message) userCorrelationData) + : this.getMessageBuilderFactory().withPayload(userCorrelationData); - AbstractIntegrationMessageBuilder builder = userCorrelationData instanceof Message - ? this.getMessageBuilderFactory().fromMessage((Message) userCorrelationData) - : this.getMessageBuilderFactory().withPayload(userCorrelationData); - - Message confirmMessage = builder - .copyHeaders(headers) - .build(); - if (ack && this.confirmAckChannel != null) { - sendOutput(confirmMessage, this.confirmAckChannel, true); + confirmMessage = builder + .copyHeaders(headers) + .build(); } - else if (!ack && this.confirmNackChannel != null) { - sendOutput(confirmMessage, this.confirmNackChannel, true); + else { + confirmMessage = this.errorMessageStrategy.buildErrorMessage( + new NackedAmqpMessageException(wrapper.getMessage(), wrapper.getUserData(), cause), null); + } + if (ack && getConfirmAckChannel() != null) { + sendOutput(confirmMessage, getConfirmAckChannel(), true); + } + else if (!ack && getConfirmNackChannel() != null) { + sendOutput(confirmMessage, getConfirmNackChannel(), true); } else { if (logger.isInfoEnabled()) { @@ -517,15 +565,22 @@ public abstract class AbstractAmqpOutboundEndpoint extends AbstractReplyProducin private final Object userData; - private CorrelationDataWrapper(String id, Object userData) { + private final Message message; + + CorrelationDataWrapper(String id, Object userData, Message message) { super(id); this.userData = userData; + this.message = message; } public Object getUserData() { return this.userData; } + public Message getMessage() { + return this.message; + } + } } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/AmqpMessageHeaderErrorMessageStrategy.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/AmqpMessageHeaderErrorMessageStrategy.java index 1fcf01913b..e3fb99938e 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/AmqpMessageHeaderErrorMessageStrategy.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/AmqpMessageHeaderErrorMessageStrategy.java @@ -17,6 +17,7 @@ package org.springframework.integration.amqp.support; import java.util.Collections; +import java.util.HashMap; import java.util.Map; import org.springframework.amqp.support.AmqpHeaders; @@ -43,11 +44,11 @@ public class AmqpMessageHeaderErrorMessageStrategy implements ErrorMessageStrate */ public static final String AMQP_RAW_MESSAGE = AmqpHeaders.PREFIX + "raw_message"; - @SuppressWarnings("deprecation") @Override public ErrorMessage buildErrorMessage(Throwable throwable, AttributeAccessor context) { - Object inputMessage = context.getAttribute(ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KEY); - Map headers = + Object inputMessage = context == null ? null + : context.getAttribute(ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KEY); + Map headers = context == null ? new HashMap() : Collections.singletonMap(AMQP_RAW_MESSAGE, context.getAttribute(AMQP_RAW_MESSAGE)); return new ErrorMessage(throwable, headers, inputMessage instanceof Message ? (Message) inputMessage : null); } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/NackedAmqpMessageException.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/NackedAmqpMessageException.java new file mode 100644 index 0000000000..7b29c31e00 --- /dev/null +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/NackedAmqpMessageException.java @@ -0,0 +1,58 @@ +/* + * 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.integration.amqp.support; + +import org.springframework.messaging.Message; +import org.springframework.messaging.MessagingException; + +/** + * An exception representing a negatively acknowledged message from a + * publisher confirm. + * + * @author Gary Russell + * @since 4.3.12 + * + */ +public class NackedAmqpMessageException extends MessagingException { + + private static final long serialVersionUID = 1L; + + private final Object correlationData; + + private final String nackReason; + + public NackedAmqpMessageException(Message message, Object correlationData, String nackReason) { + super(message); + this.correlationData = correlationData; + this.nackReason = nackReason; + } + + public Object getCorrelationData() { + return this.correlationData; + } + + public String getNackReason() { + return this.nackReason; + } + + @Override + public String toString() { + return super.toString() + " [correlationData=" + this.correlationData + ", nackReason=" + this.nackReason + + "]"; + } + +} diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/ReturnedAmqpMessageException.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/ReturnedAmqpMessageException.java new file mode 100644 index 0000000000..a6a9ffaacb --- /dev/null +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/support/ReturnedAmqpMessageException.java @@ -0,0 +1,80 @@ +/* + * 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.integration.amqp.support; + +import org.springframework.amqp.core.Message; +import org.springframework.messaging.MessagingException; + +/** + * A MessagingException for a returned message. + * + * @author Gary Russell + * @since 4.3.12 + * + */ +public class ReturnedAmqpMessageException extends MessagingException { + + private static final long serialVersionUID = 1L; + + private final Message amqpMessage; + + private final int replyCode; + + private final String replyText; + + private final String exchange; + + private final String routingKey; + + public ReturnedAmqpMessageException(org.springframework.messaging.Message message, Message amqpMessage, + int replyCode, String replyText, String exchange, String routingKey) { + super(message); + this.amqpMessage = amqpMessage; + this.replyCode = replyCode; + this.replyText = replyText; + this.exchange = exchange; + this.routingKey = routingKey; + } + + public Message getAmqpMessage() { + return this.amqpMessage; + } + + public int getReplyCode() { + return this.replyCode; + } + + public String getReplyText() { + return this.replyText; + } + + public String getExchange() { + return this.exchange; + } + + public String getRoutingKey() { + return this.routingKey; + } + + @Override + public String toString() { + return super.toString() + " [amqpMessage=" + this.amqpMessage + ", replyCode=" + this.replyCode + + ", replyText=" + this.replyText + ", exchange=" + this.exchange + ", routingKey=" + this.routingKey + + "]"; + } + +} diff --git a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-5.0.xsd b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-5.0.xsd index 33d98eab40..3423339107 100644 --- a/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-5.0.xsd +++ b/spring-integration-amqp/src/main/resources/org/springframework/integration/amqp/config/spring-integration-amqp-5.0.xsd @@ -544,6 +544,19 @@ property set to TRUE. + + + + + + + + + + + confirm-ack-channel="ackChannel" + error-message-strategy="ems"/> + + diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java index 5f5eb611e1..7eb9af3657 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundChannelAdapterParserTests.java @@ -179,6 +179,7 @@ public class AmqpOutboundChannelAdapterParserTests { MessageChannel ackChannel = context.getBean("ackChannel", MessageChannel.class); assertSame(ackChannel, TestUtils.getPropertyValue(endpoint, "confirmAckChannel")); assertSame(nullChannel, TestUtils.getPropertyValue(endpoint, "confirmNackChannel")); + assertSame(context.getBean("ems"), TestUtils.getPropertyValue(endpoint, "errorMessageStrategy")); } @SuppressWarnings("rawtypes") diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml index 64903d9aed..6885607938 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests-context.xml @@ -18,10 +18,13 @@ delay-expression="42" auto-startup="false" order="5" - return-channel="returnChannel"> + return-channel="returnChannel" + error-message-strategy="ems"> + + diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java index a2942da215..516882be63 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/config/AmqpOutboundGatewayParserTests.java @@ -77,6 +77,7 @@ public class AmqpOutboundGatewayParserTests { assertEquals("amqp:outbound-async-gateway", async.getComponentType()); checkGWProps(context, async); assertSame(context.getBean("asyncTemplate"), TestUtils.getPropertyValue(async, "template")); + assertSame(context.getBean("ems"), TestUtils.getPropertyValue(gateway, "errorMessageStrategy")); context.close(); } diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpointTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpointTests.java index e2a997c149..68449635ce 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpointTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AmqpOutboundEndpointTests.java @@ -16,9 +16,11 @@ package org.springframework.integration.amqp.outbound; +import static org.hamcrest.Matchers.instanceOf; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertThat; import org.junit.Rule; import org.junit.Test; @@ -30,11 +32,14 @@ import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.rabbit.junit.BrokerRunning; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.beans.factory.annotation.Autowired; +import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.integration.amqp.support.ReturnedAmqpMessageException; import org.springframework.integration.mapping.support.JsonHeaders; import org.springframework.integration.support.MessageBuilder; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; import org.springframework.messaging.PollableChannel; +import org.springframework.messaging.support.ErrorMessage; import org.springframework.test.annotation.DirtiesContext; import org.springframework.test.annotation.DirtiesContext.ClassMode; import org.springframework.test.context.ContextConfiguration; @@ -87,6 +92,10 @@ public class AmqpOutboundEndpointTests { @Autowired private ConnectionFactory connectionFactory; + @Autowired + @Qualifier("withReturns.handler") + private AmqpOutboundEndpoint withReturns; + @Test public void testGatewayPublisherConfirms() throws Exception { while (this.amqpTemplateConfirms.receive(this.queue.getName()) != null) { @@ -137,6 +146,7 @@ public class AmqpOutboundEndpointTests { @Test public void adapterWithReturns() throws Exception { + this.withReturns.setErrorMessageStrategy(null); Message message = MessageBuilder.withPayload("hello").build(); this.returnRequestChannel.send(message); Message returned = returnChannel.receive(10000); @@ -144,6 +154,18 @@ public class AmqpOutboundEndpointTests { assertEquals(message.getPayload(), returned.getPayload()); } + @Test + public void adapterWithReturnsAndErrorMessageStrategy() throws Exception { + Message message = MessageBuilder.withPayload("hello").build(); + this.returnRequestChannel.send(message); + Message returned = returnChannel.receive(10000); + assertNotNull(returned); + assertThat(returned, instanceOf(ErrorMessage.class)); + assertThat(returned.getPayload(), instanceOf(ReturnedAmqpMessageException.class)); + ReturnedAmqpMessageException payload = (ReturnedAmqpMessageException) returned.getPayload(); + assertEquals(message.getPayload(), payload.getFailedMessage().getPayload()); + } + @Test public void adapterWithContentType() throws Exception { RabbitTemplate template = new RabbitTemplate(this.connectionFactory); diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AsyncAmqpGatewayTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AsyncAmqpGatewayTests.java index cbed32f6fa..99c03cbf97 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AsyncAmqpGatewayTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/outbound/AsyncAmqpGatewayTests.java @@ -54,6 +54,8 @@ import org.springframework.amqp.support.AmqpHeaders; import org.springframework.amqp.utils.test.TestUtils; import org.springframework.beans.DirectFieldAccessor; import org.springframework.beans.factory.BeanFactory; +import org.springframework.integration.amqp.support.NackedAmqpMessageException; +import org.springframework.integration.amqp.support.ReturnedAmqpMessageException; import org.springframework.integration.channel.DirectChannel; import org.springframework.integration.channel.QueueChannel; import org.springframework.integration.support.MessageBuilder; @@ -202,7 +204,10 @@ public class AsyncAmqpGatewayTests { gateway.handleMessage(message); Message returned = returnChannel.receive(10000); assertNotNull(returned); - assertEquals("fiz", returned.getPayload()); + assertThat(returned, instanceOf(ErrorMessage.class)); + assertThat(returned.getPayload(), instanceOf(ReturnedAmqpMessageException.class)); + ReturnedAmqpMessageException payload = (ReturnedAmqpMessageException) returned.getPayload(); + assertEquals("fiz", payload.getFailedMessage().getPayload()); ackChannel.receive(10000); ackChannel.purge(null); @@ -222,9 +227,11 @@ public class AsyncAmqpGatewayTests { ack = ackChannel.receive(10000); assertNotNull(ack); - assertEquals("buz", ack.getPayload()); - assertEquals("nacknack", ack.getHeaders().get(AmqpHeaders.PUBLISH_CONFIRM_NACK_CAUSE)); - assertEquals(false, ack.getHeaders().get(AmqpHeaders.PUBLISH_CONFIRM)); + assertThat(returned, instanceOf(ErrorMessage.class)); + assertThat(returned.getPayload(), instanceOf(ReturnedAmqpMessageException.class)); + NackedAmqpMessageException nack = (NackedAmqpMessageException) ack.getPayload(); + assertEquals("buz", nack.getFailedMessage().getPayload()); + assertEquals("nacknack", nack.getNackReason()); asyncTemplate.stop(); receiver.stop(); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/DefaultErrorMessageStrategy.java b/spring-integration-core/src/main/java/org/springframework/integration/support/DefaultErrorMessageStrategy.java index a286b77924..a9f6aa1313 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/support/DefaultErrorMessageStrategy.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/DefaultErrorMessageStrategy.java @@ -37,7 +37,8 @@ public class DefaultErrorMessageStrategy implements ErrorMessageStrategy { @Override public ErrorMessage buildErrorMessage(Throwable throwable, AttributeAccessor attributes) { - Object inputMessage = attributes.getAttribute(ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KEY); + Object inputMessage = attributes == null ? null + : attributes.getAttribute(ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KEY); return new ErrorMessage(throwable, inputMessage instanceof Message ? (Message) inputMessage : null); } diff --git a/src/reference/asciidoc/amqp.adoc b/src/reference/asciidoc/amqp.adoc index 16dd4604f4..66f8e10571 100644 --- a/src/reference/asciidoc/amqp.adoc +++ b/src/reference/asciidoc/amqp.adoc @@ -505,9 +505,10 @@ A configuration sample for an AMQP Outbound Channel Adapter is shown below. confirm-ack-channel="" <11> confirm-nack-channel="" <12> return-channel="" <13> - header-mapper="" <14> - mapped-request-headers="" <15> - lazy-connect="true" /> <16> + error-message-strategy="" <14> + header-mapper="" <15> + mapped-request-headers="" <16> + lazy-connect="true" /> <17> ---- @@ -574,21 +575,25 @@ Previously, a new message was created with the correlation data as its payload, _Optional_. <11> The channel to which positive (ack) publisher confirms are sent; payload is the correlation data defined by the _confirm-correlation-expression_. +If the expression is `#root` or `#this`, the message is built from the original message, with the `amqp_publishConfirm` header set to `true`. _Optional, default=nullChannel_. -<12> The channel to which negative (nack) publisher confirms are sent; payload is the correlation data defined by the _confirm-correlation-expression_. +<12> The channel to which negative (nack) publisher confirms are sent; payload is the correlation data defined by the _confirm-correlation-expression_ (if there is no `ErrorMessageStrategy` configured). +If the expression is `#root` or `#this`, the message is built from the original message, with the `amqp_publishConfirm` header set to `false`. +When there is an `ErrorMessageStrategy`, the message will be an `ErrorMessage` with a `NackedAmqpMessageException` payload. _Optional, default=nullChannel_. <13> The channel to which returned messages are sent. When provided, the underlying amqp template is configured to return undeliverable messages to the adapter. -The message will be constructed from the data received from amqp, with the following additional headers: _amqp_returnReplyCode, - amqp_returnReplyText, amqp_returnExchange, amqp_returnRoutingKey_. +When there is no `ErrorMessageStrategy` configured, the message will be constructed from the data received from amqp, with the following additional headers: _amqp_returnReplyCode, amqp_returnReplyText, amqp_returnExchange, amqp_returnRoutingKey_. +When there is an `ErrorMessageStrategy`, the message will be an `ErrorMessage` with a `ReturnedAmqpMessageException` payload. _Optional_. +<14> A reference to an `ErrorMessageStrategy` implementation used to build `ErrorMessage` s when sending returned or negatively acknowedged messages. -<14> A reference to an `AmqpHeaderMapper` to use when sending AMQP Messages. +<15> A reference to an `AmqpHeaderMapper` to use when sending AMQP Messages. By default only standard AMQP properties (e.g. `contentType`) will be copied to the Spring Integration `MessageHeaders`. Any user-defined headers will NOT be copied to the Message by the default`DefaultAmqpHeaderMapper`. @@ -596,12 +601,12 @@ Not allowed if 'request-header-names' is provided. _Optional_. -<15> Comma-separated list of names of AMQP Headers to be mapped from the `MessageHeaders` to the AMQP Message. +<16> Comma-separated list of names of AMQP Headers to be mapped from the `MessageHeaders` to the AMQP Message. Not allowed if the 'header-mapper' reference is provided. The values in this list can also be simple patterns to be matched against the header names (e.g. `"\*"` or `"foo*, bar"` or `"*foo"`). -<16> When set to `false`, the endpoint will attempt to connect to the broker during application context initialization. +<17> When set to `false`, the endpoint will attempt to connect to the broker during application context initialization. This allows "fail fast" detection of bad configuration, but will also cause initialization to fail if the broker is down. When true (default), the connection is established (if it doesn't already exist because some other component established it) when the first message is sent. @@ -718,7 +723,8 @@ Configuration for an AMQP Outbound Gateway is shown below. confirm-ack-channel="" <14> confirm-nack-channel="" <15> return-channel="" <16> - lazy-connect="true" /> <17> + error-message-strategy="" <17> + lazy-connect="true" /> <18> ---- @@ -795,20 +801,24 @@ emitted on the ack/nack channel is based on that message, with the additional he Previously, a new message was created with the correlation data as its payload, regardless of type. _Optional_. -<14> Since _version 4.2_. The channel to which positive (ack) publisher confirms are sent; payload is the correlation data defined by the _confirm-correlation-expression_. +<14> The channel to which positive (ack) publisher confirms are sent; payload is the correlation data defined by the _confirm-correlation-expression_. +If the expression is `#root` or `#this`, the message is built from the original message, with the `amqp_publishConfirm` header set to `true`. _Optional, default=nullChannel_. -<15> Since _version 4.2_. The channel to which negative (nack) publisher confirms are sent; payload is the correlation data defined by the _confirm-correlation-expression_. +<15> The channel to which negative (nack) publisher confirms are sent; payload is the correlation data defined by the _confirm-correlation-expression_ (if there is no `ErrorMessageStrategy` configured). +If the expression is `#root` or `#this`, the message is built from the original message, with the `amqp_publishConfirm` header set to `false`. +When there is an `ErrorMessageStrategy`, the message will be an `ErrorMessage` with a `NackedAmqpMessageException` payload. _Optional, default=nullChannel_. <16> The channel to which returned messages are sent. -When provided, the underlying amqp template is configured to return undeliverable messages to the gateway. -The message will be constructed from the data received from amqp, with the following additional headers: _amqp_returnReplyCode, - amqp_returnReplyText, amqp_returnExchange, amqp_returnRoutingKey_. +When provided, the underlying amqp template is configured to return undeliverable messages to the adapter. +When there is no `ErrorMessageStrategy` configured, the message will be constructed from the data received from amqp, with the following additional headers: _amqp_returnReplyCode, amqp_returnReplyText, amqp_returnExchange, amqp_returnRoutingKey_. +When there is an `ErrorMessageStrategy`, the message will be an `ErrorMessage` with a `ReturnedAmqpMessageException` payload. _Optional_. +<17> A reference to an `ErrorMessageStrategy` implementation used to build `ErrorMessage` s when sending returned or negatively acknowedged messages. -<17> When set to `false`, the endpoint will attempt to connect to the broker during application context initialization. +<18> When set to `false`, the endpoint will attempt to connect to the broker during application context initialization. This allows "fail fast" detection of bad configuration, by logging an error message if the broker is down. When true (default), the connection is established (if it doesn't already exist because some other component established it) when the first message is sent.