From 22cbe9e0dd81c00ec9891d817ec34e15f3afd878 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 19 Dec 2011 08:58:59 -0500 Subject: [PATCH 1/2] INT-2291 polished broken test --- .../channel/MixedDispatcherConfigurationScenarioTests.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests.java index a0a57432fc..ef796377a5 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests.java @@ -24,7 +24,6 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; import org.junit.Before; -import org.junit.Ignore; import org.junit.Test; import org.junit.runner.RunWith; import org.mockito.InOrder; @@ -429,8 +428,7 @@ public class MixedDispatcherConfigurationScenarioTests { verify(exceptionRegistry, never()).add((Exception) anyObject()); } - - @Test(timeout = 5000) @Ignore + @Test(timeout = 5000) public void failoverNoLoadBalancingWithExecutorConcurrent() throws Exception { final ExecutorChannel channel = (ExecutorChannel) ac.getBean("noLoadBalancerFailoverExecutor"); final UnicastingDispatcher dispatcher = channel.getDispatcher(); @@ -475,6 +473,8 @@ public class MixedDispatcherConfigurationScenarioTests { start.countDown(); allDone.await(); + Thread.sleep(500); + scheduler.shutdown(); scheduler.awaitTermination(5, TimeUnit.SECONDS); From 8869da7a958d34b56a554a52f735dc771c182bd1 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 19 Dec 2011 09:22:47 -0500 Subject: [PATCH 2/2] INT-2291 polished MixedDispatcherConfigurationScenarioTests.failoverNoLoadBalancingWithExecutorConcurrent() test even more, introduced 'taskSubmissionMonitor' as CountDownLatch to ensure that shutdown command is not called untill all tasks are submitted. Also made sure that ExecutorChannel is using the same tasks executor so shutdown command actually shuts down the same executor. --- ...cherConfigurationScenarioTests-context.xml | 4 +- ...dDispatcherConfigurationScenarioTests.java | 52 ++++++++++--------- 2 files changed, 29 insertions(+), 27 deletions(-) diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests-context.xml index ac8ca06ada..0916a38ba4 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests-context.xml @@ -28,9 +28,7 @@ - - - + diff --git a/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests.java b/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests.java index ef796377a5..bd4910ab75 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/channel/MixedDispatcherConfigurationScenarioTests.java @@ -19,7 +19,6 @@ package org.springframework.integration.channel; import java.util.List; import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -61,7 +60,7 @@ public class MixedDispatcherConfigurationScenarioTests { private static final int TOTAL_EXECUTIONS = 40; - private ExecutorService scheduler; + private ExecutorService executor; private CountDownLatch allDone; private CountDownLatch start; @@ -91,9 +90,10 @@ public class MixedDispatcherConfigurationScenarioTests { Mockito.reset(handlerA); Mockito.reset(handlerB); Mockito.reset(handlerC); - scheduler = Executors.newCachedThreadPool(); + ac = new ClassPathXmlApplicationContext("MixedDispatcherConfigurationScenarioTests-context.xml", MixedDispatcherConfigurationScenarioTests.class); + executor = ac.getBean("taskExecutor", ExecutorService.class); allDone = new CountDownLatch(TOTAL_EXECUTIONS); start = new CountDownLatch(1); failed = new AtomicBoolean(false); @@ -146,13 +146,13 @@ public class MixedDispatcherConfigurationScenarioTests { } }; for (int i = 0; i < TOTAL_EXECUTIONS; i++) { - scheduler.execute(messageSenderTask); + executor.execute(messageSenderTask); } start.countDown(); allDone.await(); - scheduler.shutdown(); - scheduler.awaitTermination(5, TimeUnit.SECONDS); + executor.shutdown(); + executor.awaitTermination(5, TimeUnit.SECONDS); assertTrue("not all messages were accepted", failed.get()); verify(handlerA, times(TOTAL_EXECUTIONS)).handleMessage(message); @@ -196,13 +196,13 @@ public class MixedDispatcherConfigurationScenarioTests { } }; for (int i = 0; i < TOTAL_EXECUTIONS; i++) { - scheduler.execute(messageSenderTask); + executor.execute(messageSenderTask); } start.countDown(); allDone.await(); - scheduler.shutdown(); - scheduler.awaitTermination(5, TimeUnit.SECONDS); + executor.shutdown(); + executor.awaitTermination(5, TimeUnit.SECONDS); assertTrue("not all messages were accepted", failed.get()); verify(handlerA, times(TOTAL_EXECUTIONS)).handleMessage(message); @@ -274,13 +274,13 @@ public class MixedDispatcherConfigurationScenarioTests { } }; for (int i = 0; i < TOTAL_EXECUTIONS; i++) { - scheduler.execute(messageSenderTask); + executor.execute(messageSenderTask); } start.countDown(); allDone.await(); - scheduler.shutdown(); - scheduler.awaitTermination(5, TimeUnit.SECONDS); + executor.shutdown(); + executor.awaitTermination(5, TimeUnit.SECONDS); assertTrue("not all messages were accepted", failed.get()); verify(handlerA, times(14)).handleMessage(message); @@ -334,13 +334,13 @@ public class MixedDispatcherConfigurationScenarioTests { } }; for (int i = 0; i < TOTAL_EXECUTIONS; i++) { - scheduler.execute(messageSenderTask); + executor.execute(messageSenderTask); } start.countDown(); allDone.await(); - scheduler.shutdown(); - scheduler.awaitTermination(5, TimeUnit.SECONDS); + executor.shutdown(); + executor.awaitTermination(5, TimeUnit.SECONDS); assertTrue("not all messages were accepted", failed.get()); verify(handlerA, times(14)).handleMessage(message); @@ -413,13 +413,13 @@ public class MixedDispatcherConfigurationScenarioTests { } }; for (int i = 0; i < TOTAL_EXECUTIONS; i++) { - scheduler.execute(messageSenderTask); + executor.execute(messageSenderTask); } start.countDown(); allDone.await(); - scheduler.shutdown(); - scheduler.awaitTermination(5, TimeUnit.SECONDS); + executor.shutdown(); + executor.awaitTermination(5, TimeUnit.SECONDS); assertFalse("not all messages were accepted", failed.get()); verify(handlerA, times(TOTAL_EXECUTIONS)).handleMessage(message); @@ -434,7 +434,9 @@ public class MixedDispatcherConfigurationScenarioTests { final UnicastingDispatcher dispatcher = channel.getDispatcher(); dispatcher.addHandler(handlerA); dispatcher.addHandler(handlerB); - dispatcher.addHandler(handlerC); + dispatcher.addHandler(handlerC); + + final CountDownLatch taskSubmissionMonitor = new CountDownLatch(TOTAL_EXECUTIONS); doAnswer(new Answer() { public Object answer(InvocationOnMock invocation) { @@ -460,23 +462,25 @@ public class MixedDispatcherConfigurationScenarioTests { Runnable messageSenderTask = new Runnable() { public void run() { try { - start.await(); + start.await(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } channel.send(message); + taskSubmissionMonitor.countDown(); } }; for (int i = 0; i < TOTAL_EXECUTIONS; i++) { - scheduler.execute(messageSenderTask); + executor.execute(messageSenderTask); + } start.countDown(); allDone.await(); - Thread.sleep(500); + taskSubmissionMonitor.await(); - scheduler.shutdown(); - scheduler.awaitTermination(5, TimeUnit.SECONDS); + executor.shutdown(); + executor.awaitTermination(5, TimeUnit.SECONDS); verify(handlerA, times(TOTAL_EXECUTIONS)).handleMessage(message); verify(handlerB, times(TOTAL_EXECUTIONS)).handleMessage(message);