From 76bbbdd771fe396f4707f049fb67f8c6d965553b Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 5 Jun 2018 16:34:58 -0400 Subject: [PATCH] INT-4482: AMQP: Fix Double ErrorMessage JIRA: https://jira.spring.io/browse/INT-4482 The outer try/catch sends an `ErrorMessage` for all exceptions; it should only do so for `MessageConversionException`. Integration flow exceptions will have been already handled by `MessageProducerSupport`. Also, populate the raw message header consistently - previously it only was populated for flow exceptions. Although the LEFE contains the raw message, it should be in the `ErrorMessage` header for consistency. **cherry-pick to 5.0.x, 4.3.x** * Polishing - PR Comments --- .../inbound/AmqpInboundChannelAdapter.java | 19 +-- .../integration/amqp/dsl/AmqpTests.java | 127 +++++++++++++++++- .../amqp/inbound/InboundEndpointTests.java | 10 +- .../support/StringObjectMapBuilder.java | 28 ++++ 4 files changed, 164 insertions(+), 20 deletions(-) create mode 100644 spring-integration-core/src/main/java/org/springframework/integration/support/StringObjectMapBuilder.java diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java index c50c672a2d..60abef7ee4 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundChannelAdapter.java @@ -25,6 +25,7 @@ import org.springframework.amqp.rabbit.core.ChannelAwareMessageListener; import org.springframework.amqp.rabbit.listener.AbstractMessageListenerContainer; import org.springframework.amqp.rabbit.listener.exception.ListenerExecutionFailedException; import org.springframework.amqp.support.AmqpHeaders; +import org.springframework.amqp.support.converter.MessageConversionException; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.amqp.support.converter.SimpleMessageConverter; import org.springframework.core.AttributeAccessor; @@ -200,14 +201,10 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements @SuppressWarnings("unchecked") @Override public void onMessage(final Message message, final Channel channel) throws Exception { + boolean retryDisabled = AmqpInboundChannelAdapter.this.retryTemplate == null; try { - if (AmqpInboundChannelAdapter.this.retryTemplate == null) { - try { - createAndSend(message, channel); - } - finally { - attributesHolder.remove(); - } + if (retryDisabled) { + createAndSend(message, channel); } else { final org.springframework.messaging.Message toSend = createMessage(message, channel); @@ -220,8 +217,9 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements (RecoveryCallback) AmqpInboundChannelAdapter.this.recoveryCallback); } } - catch (RuntimeException e) { + catch (MessageConversionException e) { if (getErrorChannel() != null) { + setAttributesIfNecessary(message, null); getMessagingTemplate().send(getErrorChannel(), buildErrorMessage(null, new ListenerExecutionFailedException("Message conversion failed", e, message))); } @@ -229,6 +227,11 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements throw e; } } + finally { + if (retryDisabled) { + attributesHolder.remove(); + } + } } private void createAndSend(Message message, Channel channel) { diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/AmqpTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/AmqpTests.java index 00187f1ae5..ca3dcbfd08 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/AmqpTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/dsl/AmqpTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2017 the original author or authors. + * Copyright 2014-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,11 +16,16 @@ package org.springframework.integration.amqp.dsl; +import static org.hamcrest.Matchers.instanceOf; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertSame; +import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; +import java.util.concurrent.atomic.AtomicReference; + import org.junit.jupiter.api.AfterAll; import org.junit.jupiter.api.Test; @@ -37,14 +42,19 @@ import org.springframework.amqp.rabbit.junit.RabbitAvailable; 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.MessageConversionException; +import org.springframework.amqp.support.converter.SimpleMessageConverter; import org.springframework.beans.factory.annotation.Autowired; import org.springframework.beans.factory.annotation.Qualifier; +import org.springframework.context.ConfigurableApplicationContext; import org.springframework.context.Lifecycle; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import org.springframework.integration.amqp.channel.AbstractAmqpChannel; import org.springframework.integration.amqp.inbound.AmqpInboundGateway; 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.channel.QueueChannel; import org.springframework.integration.config.EnableIntegration; @@ -52,6 +62,7 @@ import org.springframework.integration.dsl.IntegrationFlow; import org.springframework.integration.dsl.IntegrationFlowBuilder; import org.springframework.integration.dsl.IntegrationFlows; import org.springframework.integration.support.MessageBuilder; +import org.springframework.integration.support.StringObjectMapBuilder; import org.springframework.integration.test.util.TestUtils; import org.springframework.messaging.Message; import org.springframework.messaging.MessageChannel; @@ -67,7 +78,8 @@ import org.springframework.test.context.junit.jupiter.SpringJUnitConfig; */ @SpringJUnitConfig @RabbitAvailable(queues = { "amqpOutboundInput", "amqpReplyChannel", "asyncReplies", - "defaultReplyTo", "si.dsl.test", "testTemplateChannelTransacted" }) + "defaultReplyTo", "si.dsl.test", "si.dsl.exception.test.dlq", + "si.dsl.conv.exception.test.dlq", "testTemplateChannelTransacted" }) @DirtiesContext public class AmqpTests { @@ -92,8 +104,10 @@ public class AmqpTests { private Lifecycle asyncOutboundGateway; @AfterAll - public static void tearDown() { - RabbitAvailableCondition.getBrokerRunning().removeTestQueues(); + public static void tearDown(ConfigurableApplicationContext context) { + context.stop(); // prevent queues from being redeclared after deletion + RabbitAvailableCondition.getBrokerRunning().removeTestQueues("si.dsl.exception.test", + "si.dsl.conv.exception.test"); } @Test @@ -173,6 +187,33 @@ public class AmqpTests { this.asyncOutboundGateway.stop(); } + @Autowired + private AtomicReference lefe; + + @Autowired + private AtomicReference raw; + + + @Test + public void testInboundMessagingExceptionFlow() { + this.amqpTemplate.convertAndSend("si.dsl.exception.test", "foo"); + assertNotNull(this.amqpTemplate.receive("si.dsl.exception.test.dlq", 30_000)); + assertNull(this.lefe.get()); + assertNotNull(this.raw.get()); + this.raw.set(null); + } + + @Test + public void testInboundConversionExceptionFlow() { + this.amqpTemplate.convertAndSend("si.dsl.conv.exception.test", "foo"); + assertNotNull(this.amqpTemplate.receive("si.dsl.conv.exception.test.dlq", 30_000)); + assertNotNull(this.lefe.get()); + assertThat(this.lefe.get().getCause(), instanceOf(MessageConversionException.class)); + assertNotNull(this.raw.get()); + this.raw.set(null); + this.lefe.set(null); + } + @Autowired private AbstractAmqpChannel unitChannel; @@ -277,6 +318,84 @@ public class AmqpTests { .get(); } + @Bean + public AtomicReference lefe() { + return new AtomicReference<>(); + } + + @Bean + public AtomicReference raw() { + return new AtomicReference<>(); + } + + @Bean + 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()); + } + + @Bean + public Queue exDLQ() { + return new Queue("si.dsl.exception.test.dlq"); + } + + @Bean + public IntegrationFlow inboundWithExceptionFlow(ConnectionFactory cf) { + return IntegrationFlows.from(Amqp.inboundAdapter(cf, exQueue()) + .configureContainer(c -> c.defaultRequeueRejected(false)) + .errorChannel("errors.input")) + .handle(m -> { + throw new RuntimeException("fail"); + }) + .get(); + } + + @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(); + }); + } + + @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()); + } + + @Bean + public Queue exConvDLQ() { + return new Queue("si.dsl.conv.exception.test.dlq"); + } + + @Bean + public IntegrationFlow inboundWithConvExceptionFlow(ConnectionFactory cf) { + return IntegrationFlows.from(Amqp.inboundAdapter(cf, exConvQueue()) + .configureContainer(c -> c.defaultRequeueRejected(false)) + .messageConverter(new SimpleMessageConverter() { + + @Override + public Object fromMessage(org.springframework.amqp.core.Message message) + throws MessageConversionException { + throw new MessageConversionException("fail"); + } + + }) + .errorChannel("errors.input")) + .get(); + } + @Bean public Queue asyncReplies() { return new Queue("asyncReplies"); diff --git a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/InboundEndpointTests.java b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/InboundEndpointTests.java index 639fe22a4f..9cc1fc2408 100644 --- a/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/InboundEndpointTests.java +++ b/spring-integration-amqp/src/test/java/org/springframework/integration/amqp/inbound/InboundEndpointTests.java @@ -238,17 +238,11 @@ public class InboundEndpointTests { adapter.setOutputChannel(outputChannel); QueueChannel errorChannel = new QueueChannel(); adapter.setErrorChannel(errorChannel); - adapter.setMessageConverter(new MessageConverter() { - - @Override - public org.springframework.amqp.core.Message toMessage(Object object, MessageProperties messageProperties) - throws MessageConversionException { - throw new MessageConversionException("intended"); - } + adapter.setMessageConverter(new SimpleMessageConverter() { @Override public Object fromMessage(org.springframework.amqp.core.Message message) throws MessageConversionException { - return null; + throw new MessageConversionException("intended"); } }); diff --git a/spring-integration-core/src/main/java/org/springframework/integration/support/StringObjectMapBuilder.java b/spring-integration-core/src/main/java/org/springframework/integration/support/StringObjectMapBuilder.java new file mode 100644 index 0000000000..974ae0b6be --- /dev/null +++ b/spring-integration-core/src/main/java/org/springframework/integration/support/StringObjectMapBuilder.java @@ -0,0 +1,28 @@ +/* + * 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.integration.support; + +/** + * A map builder creating a map with String keys and values. + * + * @author Gary Russell + * + * @since 5.0.6 + */ +public class StringObjectMapBuilder extends MapBuilder { + +}