Merge pull request #643 from garyrussell/INT-2781
INT-2781 Retry Advice; Fix Recovery For Zero Tries
This commit is contained in:
@@ -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<Object> {
|
||||
|
||||
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<Object> {
|
||||
return null;
|
||||
}
|
||||
|
||||
public static class RetryExceptionNotAvailableException extends MessagingException {
|
||||
|
||||
private static final long serialVersionUID = 1L;
|
||||
|
||||
public RetryExceptionNotAvailableException(Message<?> message, String description) {
|
||||
super(message, description);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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<Object> recoveryCallback;
|
||||
|
||||
private static final ThreadLocal<Message<?>> messageHolder = new ThreadLocal<Message<?>>();
|
||||
|
||||
// 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<Object>(){
|
||||
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<Object>(){
|
||||
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 <T> boolean open(RetryContext context, RetryCallback<T> callback) {
|
||||
context.setAttribute("message", messageHolder.get());
|
||||
return true;
|
||||
}
|
||||
|
||||
public <T> void close(RetryContext context, RetryCallback<T> callback, Throwable throwable) {
|
||||
}
|
||||
|
||||
public <T> void onError(RetryContext context, RetryCallback<T> callback, Throwable throwable) {
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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<Advice> adviceChain = new ArrayList<Advice>();
|
||||
adviceChain.add(advice);
|
||||
handler.setAdviceChain(adviceChain);
|
||||
handler.afterPropertiesSet();
|
||||
|
||||
Message<String> message = new GenericMessage<String>("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<Advice> adviceChain = new ArrayList<Advice>();
|
||||
adviceChain.add(advice);
|
||||
handler.setAdviceChain(adviceChain);
|
||||
handler.afterPropertiesSet();
|
||||
|
||||
Message<String> message = new GenericMessage<String>("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());
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user