From 8733350c43f0360194a9d4fe73da4f1e539e669f Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 9 Oct 2012 13:48:40 -0400 Subject: [PATCH] INT-2781 Retry Advice; Fix Recovery For Zero Tries Retry Advice (ErrorMessageSendingRecoverer) tried to send a message with a null payload when the RetryPolicy allowed zero attempts. Create a MessageHandlingException with the failedMessage, when no attempts were made to call the handler. --- .../advice/ErrorMessageSendingRecoverer.java | 16 +++++ .../advice/RequestHandlerRetryAdvice.java | 60 ++++++++++++---- .../advice/AdvisedMessageHandlerTests.java | 70 +++++++++++++++++++ 3 files changed, 131 insertions(+), 15 deletions(-) diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/advice/ErrorMessageSendingRecoverer.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/advice/ErrorMessageSendingRecoverer.java index 96621c5e58..59c296bceb 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/advice/ErrorMessageSendingRecoverer.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/advice/ErrorMessageSendingRecoverer.java @@ -17,6 +17,7 @@ package org.springframework.integration.handler.advice; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; +import org.springframework.integration.Message; import org.springframework.integration.MessageChannel; import org.springframework.integration.MessagingException; import org.springframework.integration.core.MessagingTemplate; @@ -49,6 +50,13 @@ public class ErrorMessageSendingRecoverer implements RecoveryCallback { public Object recover(RetryContext context) throws Exception { Throwable lastThrowable = context.getLastThrowable(); + if (lastThrowable == null) { + lastThrowable = new RetryExceptionNotAvailableException( + (Message) context.getAttribute("message"), + "No retry exception available; " + + "this can occur, for example, if the RetryPolicy allowed zero attempts to execute the handler; " + + "RetryContext: " + context.toString()); + } if (logger.isDebugEnabled()) { String supplement = ""; if (lastThrowable instanceof MessagingException) { @@ -60,4 +68,12 @@ public class ErrorMessageSendingRecoverer implements RecoveryCallback { return null; } + public static class RetryExceptionNotAvailableException extends MessagingException { + + private static final long serialVersionUID = 1L; + + public RetryExceptionNotAvailableException(Message message, String description) { + super(message, description); + } + } } diff --git a/spring-integration-core/src/main/java/org/springframework/integration/handler/advice/RequestHandlerRetryAdvice.java b/spring-integration-core/src/main/java/org/springframework/integration/handler/advice/RequestHandlerRetryAdvice.java index 837f78d5bf..a89660b799 100644 --- a/spring-integration-core/src/main/java/org/springframework/integration/handler/advice/RequestHandlerRetryAdvice.java +++ b/spring-integration-core/src/main/java/org/springframework/integration/handler/advice/RequestHandlerRetryAdvice.java @@ -20,8 +20,10 @@ import org.springframework.integration.MessagingException; import org.springframework.retry.RecoveryCallback; import org.springframework.retry.RetryCallback; import org.springframework.retry.RetryContext; +import org.springframework.retry.RetryListener; import org.springframework.retry.RetryState; import org.springframework.retry.support.RetryTemplate; +import org.springframework.util.Assert; /** * Uses spring-retry to perform stateless or stateful retry. @@ -34,12 +36,15 @@ import org.springframework.retry.support.RetryTemplate; * @since 2.2 * */ -public class RequestHandlerRetryAdvice extends AbstractRequestHandlerAdvice { +public class RequestHandlerRetryAdvice extends AbstractRequestHandlerAdvice + implements RetryListener { private volatile RetryTemplate retryTemplate = new RetryTemplate(); private volatile RecoveryCallback recoveryCallback; + private static final ThreadLocal> messageHolder = new ThreadLocal>(); + // Stateless unless a state generator is provided private volatile RetryStateGenerator retryStateGenerator = new RetryStateGenerator() { @@ -49,6 +54,7 @@ public class RequestHandlerRetryAdvice extends AbstractRequestHandlerAdvice { }; public void setRetryTemplate(RetryTemplate retryTemplate) { + Assert.notNull(retryTemplate, "'retryTemplate' cannot be null"); this.retryTemplate = retryTemplate; } @@ -57,30 +63,54 @@ public class RequestHandlerRetryAdvice extends AbstractRequestHandlerAdvice { } public void setRetryStateGenerator(RetryStateGenerator retryStateGenerator) { + Assert.notNull(retryStateGenerator, "'retryStateGenerator' cannot be null"); this.retryStateGenerator = retryStateGenerator; } + @Override + protected void onInit() throws Exception { + super.onInit(); + this.retryTemplate.registerListener(this); + } + @Override protected Object doInvoke(final ExecutionCallback callback, Object target, final Message message) throws Exception { RetryState retryState = null; retryState = this.retryStateGenerator.determineRetryState(message); + messageHolder.set(message); - return retryTemplate.execute(new RetryCallback(){ - public Object doWithRetry(RetryContext context) throws Exception { - try { - return callback.execute(); - } - catch (MessagingException e) { - if (e.getFailedMessage() == null) { - e.setFailedMessage(message); + try { + return retryTemplate.execute(new RetryCallback(){ + public Object doWithRetry(RetryContext context) throws Exception { + try { + return callback.execute(); + } + catch (MessagingException e) { + if (e.getFailedMessage() == null) { + e.setFailedMessage(message); + } + throw e; + } + catch (Exception e) { + throw new MessagingException(message, "Failed to invoke handler", e); } - throw e; } - catch (Exception e) { - throw new MessagingException(message, "Failed to invoke handler", e); - } - } - }, this.recoveryCallback, retryState); + }, this.recoveryCallback, retryState); + } + finally { + messageHolder.remove(); + } + } + + public boolean open(RetryContext context, RetryCallback callback) { + context.setAttribute("message", messageHolder.get()); + return true; + } + + public void close(RetryContext context, RetryCallback callback, Throwable throwable) { + } + + public void onError(RetryContext context, RetryCallback callback, Throwable throwable) { } } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/handler/advice/AdvisedMessageHandlerTests.java b/spring-integration-core/src/test/java/org/springframework/integration/handler/advice/AdvisedMessageHandlerTests.java index 26b1f64864..5a48f708c3 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/handler/advice/AdvisedMessageHandlerTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/handler/advice/AdvisedMessageHandlerTests.java @@ -18,6 +18,7 @@ package org.springframework.integration.handler.advice; 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.assertTrue; import static org.junit.Assert.fail; @@ -40,7 +41,9 @@ import org.springframework.integration.message.GenericMessage; import org.springframework.retry.RecoveryCallback; import org.springframework.retry.RetryContext; import org.springframework.retry.RetryState; +import org.springframework.retry.policy.SimpleRetryPolicy; import org.springframework.retry.support.DefaultRetryState; +import org.springframework.retry.support.RetryTemplate; /** * @author Gary Russell @@ -507,4 +510,71 @@ public class AdvisedMessageHandlerTests { assertNotNull(reply); assertEquals("baz", reply.getPayload()); } + + @Test + public void errorMessageSendingRecovererTests() { + AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + + @Override + protected Object handleRequestMessage(Message requestMessage) { + throw new RuntimeException("fooException"); + } + }; + QueueChannel errors = new QueueChannel(); + RequestHandlerRetryAdvice advice = new RequestHandlerRetryAdvice(); + ErrorMessageSendingRecoverer recoverer = new ErrorMessageSendingRecoverer(errors); + advice.setRecoveryCallback(recoverer); + + List adviceChain = new ArrayList(); + adviceChain.add(advice); + handler.setAdviceChain(adviceChain); + handler.afterPropertiesSet(); + + Message message = new GenericMessage("Hello, world!"); + handler.handleMessage(message); + Message error = errors.receive(1000); + assertNotNull(error); + assertEquals("fooException", ((Exception) error.getPayload()).getCause().getCause().getMessage()); + + } + + @Test + public void errorMessageSendingRecovererTestsNoThrowable() { + AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + + @Override + protected Object handleRequestMessage(Message requestMessage) { + throw new RuntimeException("fooException"); + } + }; + QueueChannel errors = new QueueChannel(); + RequestHandlerRetryAdvice advice = new RequestHandlerRetryAdvice(); + ErrorMessageSendingRecoverer recoverer = new ErrorMessageSendingRecoverer(errors); + advice.setRecoveryCallback(recoverer); + RetryTemplate retryTemplate = new RetryTemplate(); + retryTemplate.setRetryPolicy(new SimpleRetryPolicy() { + + @Override + public boolean canRetry(RetryContext context) { + return false; + } + }); + advice.setRetryTemplate(retryTemplate); + advice.afterPropertiesSet(); + + List adviceChain = new ArrayList(); + adviceChain.add(advice); + handler.setAdviceChain(adviceChain); + handler.afterPropertiesSet(); + + Message message = new GenericMessage("Hello, world!"); + handler.handleMessage(message); + Message error = errors.receive(1000); + assertNotNull(error); + assertTrue(error.getPayload() instanceof ErrorMessageSendingRecoverer.RetryExceptionNotAvailableException); + assertNotNull(((MessagingException) error.getPayload()).getFailedMessage()); + assertSame(message, ((MessagingException) error.getPayload()).getFailedMessage()); + } + + }