From dd6afa81d2f7dd9ae8f0a08c6ee01457fb3de62c Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 18 Sep 2018 16:05:01 -0400 Subject: [PATCH] Use RetrySynchronizationManager Use the `RetrySynchronizationManager` instead of a `RetryListener`. Fix `setAttributesIfNecessary` for gateway conversion errors (this, at least, should be cherry-picked). --- .../inbound/AmqpInboundChannelAdapter.java | 33 +++---------------- .../amqp/inbound/AmqpInboundGateway.java | 1 + .../amqp/inbound/InboundEndpointTests.java | 16 ++++++--- 3 files changed, 17 insertions(+), 33 deletions(-) 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 da5b5a011c..56170b3694 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 @@ -39,9 +39,7 @@ import org.springframework.integration.endpoint.MessageProducerSupport; import org.springframework.integration.support.ErrorMessageStrategy; import org.springframework.integration.support.ErrorMessageUtils; import org.springframework.retry.RecoveryCallback; -import org.springframework.retry.RetryCallback; -import org.springframework.retry.RetryContext; -import org.springframework.retry.RetryListener; +import org.springframework.retry.support.RetrySynchronizationManager; import org.springframework.retry.support.RetryTemplate; import org.springframework.util.Assert; @@ -132,9 +130,6 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements + "send an error message when retries are exhausted"); } Listener messageListener = new Listener(); - if (this.retryTemplate != null) { - this.retryTemplate.registerListener(messageListener); - } this.messageListenerContainer.setMessageListener(messageListener); this.messageListenerContainer.afterPropertiesSet(); super.onInit(); @@ -177,7 +172,9 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements attributesHolder.set(ErrorMessageUtils.getAttributeAccessor(null, null)); } if (needAttributes) { - AttributeAccessor attributes = attributesHolder.get(); + AttributeAccessor attributes = this.retryTemplate != null + ? RetrySynchronizationManager.getContext() + : attributesHolder.get(); if (attributes != null) { attributes.setAttribute(ErrorMessageUtils.INPUT_MESSAGE_CONTEXT_KEY, message); attributes.setAttribute(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE, amqpMessage); @@ -196,7 +193,7 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements } } - protected class Listener implements ChannelAwareMessageListener, RetryListener { + protected class Listener implements ChannelAwareMessageListener { @SuppressWarnings("unchecked") @Override @@ -259,26 +256,6 @@ public class AmqpInboundChannelAdapter extends MessageProducerSupport implements return messagingMessage; } - @Override - public boolean open(RetryContext context, RetryCallback callback) { - if (AmqpInboundChannelAdapter.this.recoveryCallback != null) { - attributesHolder.set(context); - } - return true; - } - - @Override - public void close(RetryContext context, RetryCallback callback, - Throwable throwable) { - attributesHolder.remove(); - } - - @Override - public void onError(RetryContext context, RetryCallback callback, - Throwable throwable) { - // Empty - } - } } diff --git a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java index b7e33b1804..f9919f3ffd 100644 --- a/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java +++ b/spring-integration-amqp/src/main/java/org/springframework/integration/amqp/inbound/AmqpInboundGateway.java @@ -296,6 +296,7 @@ public class AmqpInboundGateway extends MessagingGatewaySupport { } catch (RuntimeException e) { if (getErrorChannel() != null) { + setAttributesIfNecessary(message, null); AmqpInboundGateway.this.messagingTemplate.send(getErrorChannel(), buildErrorMessage(null, new ListenerExecutionFailedException("Message conversion failed", e, message))); } 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 13e24f1ac3..145b587c99 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 @@ -248,9 +248,12 @@ public class InboundEndpointTests { }); adapter.afterPropertiesSet(); - ((ChannelAwareMessageListener) container.getMessageListener()).onMessage(null, null); + ((ChannelAwareMessageListener) container.getMessageListener()) + .onMessage(mock(org.springframework.amqp.core.Message.class), null); assertNull(outputChannel.receive(0)); - assertNotNull(errorChannel.receive(0)); + Message received = errorChannel.receive(0); + assertNotNull(received); + assertNotNull(received.getHeaders().get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE)); } @Test @@ -276,14 +279,17 @@ public class InboundEndpointTests { @Override public Object fromMessage(org.springframework.amqp.core.Message message) throws MessageConversionException { - return null; + throw new MessageConversionException("intended"); } }); adapter.afterPropertiesSet(); - ((ChannelAwareMessageListener) container.getMessageListener()).onMessage(null, null); + ((ChannelAwareMessageListener) container.getMessageListener()) + .onMessage(mock(org.springframework.amqp.core.Message.class), null); assertNull(outputChannel.receive(0)); - assertNotNull(errorChannel.receive(0)); + Message received = errorChannel.receive(0); + assertNotNull(received); + assertNotNull(received.getHeaders().get(AmqpMessageHeaderErrorMessageStrategy.AMQP_RAW_MESSAGE)); } @Test