From 8869da7a958d34b56a554a52f735dc771c182bd1 Mon Sep 17 00:00:00 2001 From: Oleg Zhurakousky Date: Mon, 19 Dec 2011 09:22:47 -0500 Subject: [PATCH] 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);