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 b081d74e28..94e882a17f 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 @@ -123,6 +123,7 @@ public class AdvisedMessageHandlerTests { public void successFailureAdvice() { final AtomicBoolean doFail = new AtomicBoolean(); AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + @Override protected Object handleRequestMessage(Message requestMessage) { if (doFail.get()) { @@ -209,6 +210,7 @@ public class AdvisedMessageHandlerTests { public void propagateOnSuccessExpressionFailures() { final AtomicBoolean doFail = new AtomicBoolean(); AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + @Override protected Object handleRequestMessage(Message requestMessage) { if (doFail.get()) { @@ -272,6 +274,7 @@ public class AdvisedMessageHandlerTests { public void propagateOnFailureExpressionFailures() { final AtomicBoolean doFail = new AtomicBoolean(true); AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + @Override protected Object handleRequestMessage(Message requestMessage) { if (doFail.get()) { @@ -310,7 +313,7 @@ public class AdvisedMessageHandlerTests { Message reply = replies.receive(1); assertNull(reply); - Message failure = failureChannel.receive(1000); + Message failure = failureChannel.receive(10000); assertNotNull(failure); assertEquals("Hello, world!", ((MessagingException) failure.getPayload()).getFailedMessage().getPayload()); assertEquals(MessageHandlingExpressionEvaluatingAdviceException.class, failure.getPayload().getClass()); @@ -328,7 +331,7 @@ public class AdvisedMessageHandlerTests { reply = replies.receive(1); assertNull(reply); - failure = failureChannel.receive(1000); + failure = failureChannel.receive(10000); assertNotNull(failure); assertEquals("Hello, world!", ((MessagingException) failure.getPayload()).getFailedMessage().getPayload()); assertEquals(MessageHandlingExpressionEvaluatingAdviceException.class, failure.getPayload().getClass()); @@ -337,6 +340,7 @@ public class AdvisedMessageHandlerTests { } @Test + @SuppressWarnings("rawtypes") public void circuitBreakerTests() throws Exception { final AtomicBoolean doFail = new AtomicBoolean(); AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { @@ -359,7 +363,7 @@ public class AdvisedMessageHandlerTests { * we reset the failure counter. */ advice.setThreshold(2); - advice.setHalfOpenAfter(100); + advice.setHalfOpenAfter(1000); List adviceChain = new ArrayList(); adviceChain.add(advice); @@ -390,7 +394,15 @@ public class AdvisedMessageHandlerTests { catch (Exception e) { assertEquals("Circuit Breaker is Open for baz", e.getCause().getMessage()); } - Thread.sleep(100); + + Map metadataMap = TestUtils.getPropertyValue(advice, "metadataMap", Map.class); + Object metadata = metadataMap.values().iterator().next(); + + DirectFieldAccessor metadataDfa = new DirectFieldAccessor(metadata); + + // Simulate some timeout in between requests + metadataDfa.setPropertyValue("lastFailure", System.currentTimeMillis() - 10000); + try { handler.handleMessage(message); fail("Expected failure"); @@ -405,7 +417,10 @@ public class AdvisedMessageHandlerTests { catch (Exception e) { assertEquals("Circuit Breaker is Open for baz", e.getCause().getMessage()); } - Thread.sleep(100); + + // Simulate some timeout in between requests + metadataDfa.setPropertyValue("lastFailure", System.currentTimeMillis() - 10000); + doFail.set(false); handler.handleMessage(message); doFail.set(true); @@ -458,7 +473,7 @@ public class AdvisedMessageHandlerTests { Message message = new GenericMessage("Hello, world!"); handler.handleMessage(message); assertTrue(counter.get() == -1); - Message reply = replies.receive(1000); + Message reply = replies.receive(10000); assertNotNull(reply); assertEquals("bar", reply.getPayload()); @@ -482,6 +497,7 @@ public class AdvisedMessageHandlerTests { RequestHandlerRetryAdvice advice = new RequestHandlerRetryAdvice(); advice.setRetryStateGenerator(new RetryStateGenerator() { + @Override public RetryState determineRetryState(Message message) { return new DefaultRetryState(message.getHeaders().getId()); @@ -504,7 +520,7 @@ public class AdvisedMessageHandlerTests { } } assertTrue(counter.get() == -1); - Message reply = replies.receive(1000); + Message reply = replies.receive(10000); assertNotNull(reply); assertEquals("bar", reply.getPayload()); @@ -528,6 +544,7 @@ public class AdvisedMessageHandlerTests { RequestHandlerRetryAdvice advice = new RequestHandlerRetryAdvice(); advice.setRetryStateGenerator(new RetryStateGenerator() { + @Override public RetryState determineRetryState(Message message) { return new DefaultRetryState(message.getHeaders().getId()); @@ -562,13 +579,14 @@ public class AdvisedMessageHandlerTests { } private void defaultStatefulRetryRecoverAfterThirdTryGuts(final AtomicInteger counter, - AbstractReplyProducingMessageHandler handler, QueueChannel replies, RequestHandlerRetryAdvice advice) { + AbstractReplyProducingMessageHandler handler, QueueChannel replies, RequestHandlerRetryAdvice advice) { advice.setRecoveryCallback(new RecoveryCallback() { @Override public Object recover(RetryContext context) throws Exception { return "baz"; } + }); List adviceChain = new ArrayList(); @@ -586,7 +604,7 @@ public class AdvisedMessageHandlerTests { } } assertTrue(counter.get() == 0); - Message reply = replies.receive(1000); + Message reply = replies.receive(10000); assertNotNull(reply); assertEquals("baz", reply.getPayload()); } @@ -613,7 +631,7 @@ public class AdvisedMessageHandlerTests { Message message = new GenericMessage("Hello, world!"); handler.handleMessage(message); - Message error = errors.receive(1000); + Message error = errors.receive(10000); assertNotNull(error); assertEquals("fooException", ((Exception) error.getPayload()).getCause().getMessage()); @@ -652,7 +670,7 @@ public class AdvisedMessageHandlerTests { Message message = new GenericMessage("Hello, world!"); handler.handleMessage(message); - Message error = errors.receive(1000); + Message error = errors.receive(10000); assertNotNull(error); assertTrue(error.getPayload() instanceof ErrorMessageSendingRecoverer.RetryExceptionNotAvailableException); assertNotNull(((MessagingException) error.getPayload()).getFailedMessage()); @@ -664,6 +682,7 @@ public class AdvisedMessageHandlerTests { final AtomicInteger counter = new AtomicInteger(3); AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + @Override protected Object handleRequestMessage(Message requestMessage) { return "foo"; @@ -674,6 +693,7 @@ public class AdvisedMessageHandlerTests { adviceChain.add(new RequestHandlerRetryAdvice()); adviceChain.add(new MethodInterceptor() { + @Override public Object invoke(MethodInvocation invocation) throws Throwable { counter.getAndDecrement(); @@ -702,6 +722,7 @@ public class AdvisedMessageHandlerTests { final AtomicInteger counter = new AtomicInteger(0); AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + @Override protected Object handleRequestMessage(Message requestMessage) { return "foo"; @@ -730,6 +751,7 @@ public class AdvisedMessageHandlerTests { adviceChain.add(expressionAdvice); adviceChain.add(new RequestHandlerRetryAdvice()); adviceChain.add(new MethodInterceptor() { + @Override public Object invoke(MethodInvocation invocation) throws Throwable { throw new RuntimeException("intentional: " + counter.incrementAndGet()); @@ -741,7 +763,7 @@ public class AdvisedMessageHandlerTests { handler.afterPropertiesSet(); handler.handleMessage(new GenericMessage("test")); - Message receive = replies.receive(1000); + Message receive = replies.receive(10000); assertNotNull(receive); assertEquals("intentional: 3", receive.getPayload()); assertEquals(1, outerCounter.get()); @@ -753,6 +775,7 @@ public class AdvisedMessageHandlerTests { final AtomicInteger counter = new AtomicInteger(0); AbstractReplyProducingMessageHandler handler = new AbstractReplyProducingMessageHandler() { + @Override protected Object handleRequestMessage(Message requestMessage) { return "foo"; @@ -771,6 +794,7 @@ public class AdvisedMessageHandlerTests { adviceChain.add(new RequestHandlerRetryAdvice()); adviceChain.add(expressionAdvice); adviceChain.add(new MethodInterceptor() { + @Override public Object invoke(MethodInvocation invocation) throws Throwable { throw new RuntimeException("intentional: " + counter.incrementAndGet()); @@ -789,9 +813,10 @@ public class AdvisedMessageHandlerTests { } for (int i = 1; i <= 3; i++) { - Message receive = errors.receive(1000); + Message receive = errors.receive(10000); assertNotNull(receive); - assertEquals("intentional: " + i, ((MessageHandlingExpressionEvaluatingAdviceException) receive.getPayload()).getEvaluationResult()); + assertEquals("intentional: " + i, + ((MessageHandlingExpressionEvaluatingAdviceException) receive.getPayload()).getEvaluationResult()); } assertNull(errors.receive(1)); @@ -821,11 +846,13 @@ public class AdvisedMessageHandlerTests { final Throwable theThrowable = new Throwable("foo"); MethodInvocation methodInvocation = mock(MethodInvocation.class); - Method method = AbstractReplyProducingMessageHandler.class.getDeclaredMethod("handleRequestMessage", Message.class); + Method method = AbstractReplyProducingMessageHandler.class.getDeclaredMethod("handleRequestMessage", + Message.class); when(methodInvocation.getMethod()).thenReturn(method); - when(methodInvocation.getArguments()).thenReturn(new Object[] {new GenericMessage("foo")}); + when(methodInvocation.getArguments()).thenReturn(new Object[] { new GenericMessage("foo") }); try { doAnswer(new Answer() { + @Override public Object answer(InvocationOnMock invocation) throws Throwable { throw theThrowable; @@ -864,7 +891,7 @@ public class AdvisedMessageHandlerTests { } catch (Throwable t) { assertSame(theThrowable, t); - ErrorMessage error = (ErrorMessage) errors.receive(1000); + ErrorMessage error = (ErrorMessage) errors.receive(10000); assertNotNull(error); assertSame(theThrowable, error.getPayload().getCause()); } @@ -874,6 +901,7 @@ public class AdvisedMessageHandlerTests { public void testInappropriateAdvice() throws Exception { final AtomicBoolean called = new AtomicBoolean(false); Advice advice = new AbstractRequestHandlerAdvice() { + @Override protected Object doInvoke(ExecutionCallback callback, Object target, Message message) throws Exception { called.set(true); @@ -882,6 +910,7 @@ public class AdvisedMessageHandlerTests { }; PollableChannel inputChannel = new QueueChannel(); PollingConsumer consumer = new PollingConsumer(inputChannel, new MessageHandler() { + @Override public void handleMessage(Message message) throws MessagingException { } @@ -890,6 +919,7 @@ public class AdvisedMessageHandlerTests { consumer.setTaskExecutor(new ErrorHandlingTaskExecutor( Executors.newSingleThreadExecutor(), new ErrorHandler() { + @Override public void handleError(Throwable t) { } @@ -906,6 +936,7 @@ public class AdvisedMessageHandlerTests { when(logger.isWarnEnabled()).thenReturn(Boolean.TRUE); final AtomicReference logMessage = new AtomicReference(); doAnswer(new Answer() { + @Override public Object answer(InvocationOnMock invocation) throws Throwable { logMessage.set((String) invocation.getArguments()[0]); @@ -926,6 +957,7 @@ public class AdvisedMessageHandlerTests { public void filterDiscardNoAdvice() { MessageFilter filter = new MessageFilter(new MessageSelector() { + @Override public boolean accept(Message message) { return false; @@ -940,6 +972,7 @@ public class AdvisedMessageHandlerTests { @Test public void filterDiscardWithinAdvice() { MessageFilter filter = new MessageFilter(new MessageSelector() { + @Override public boolean accept(Message message) { return false; @@ -950,6 +983,7 @@ public class AdvisedMessageHandlerTests { List adviceChain = new ArrayList(); final AtomicReference> discardedWithinAdvice = new AtomicReference>(); adviceChain.add(new AbstractRequestHandlerAdvice() { + @Override protected Object doInvoke(ExecutionCallback callback, Object target, Message message) throws Exception { Object result = callback.execute(); @@ -968,6 +1002,7 @@ public class AdvisedMessageHandlerTests { @Test public void filterDiscardOutsideAdvice() { MessageFilter filter = new MessageFilter(new MessageSelector() { + @Override public boolean accept(Message message) { return false; @@ -979,6 +1014,7 @@ public class AdvisedMessageHandlerTests { final AtomicReference> discardedWithinAdvice = new AtomicReference>(); final AtomicBoolean adviceCalled = new AtomicBoolean(); adviceChain.add(new AbstractRequestHandlerAdvice() { + @Override protected Object doInvoke(ExecutionCallback callback, Object target, Message message) throws Exception { Object result = callback.execute(); @@ -1059,7 +1095,9 @@ public class AdvisedMessageHandlerTests { } private interface Bar { + Object handleRequestMessage(Message message) throws Throwable; + } private class Foo implements Bar { @@ -1074,5 +1112,7 @@ public class AdvisedMessageHandlerTests { public Object handleRequestMessage(Message message) throws Throwable { throw this.throwable; } + } + }