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);