diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/RoundRobinDispatcherConcurrentTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/RoundRobinDispatcherConcurrentTests.java index 886715fb98..84ff8ce0cc 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/RoundRobinDispatcherConcurrentTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dispatcher/RoundRobinDispatcherConcurrentTests.java @@ -17,6 +17,7 @@ package org.springframework.integration.dispatcher; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; import static org.mockito.Mockito.doThrow; import static org.mockito.Mockito.times; @@ -42,6 +43,7 @@ import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor; /** * @author Iwein Fuld + * @author Artem Bilan */ @RunWith(MockitoJUnitRunner.class) public class RoundRobinDispatcherConcurrentTests { @@ -68,7 +70,7 @@ public class RoundRobinDispatcherConcurrentTests { private Message message; @Before - public void initialize() throws Exception { + public void initialize() { dispatcher.setLoadBalancingStrategy(new RoundRobinLoadBalancingStrategy()); executor.setCorePoolSize(10); executor.setMaxPoolSize(10); @@ -80,7 +82,7 @@ public class RoundRobinDispatcherConcurrentTests { this.executor.shutdown(); } - @Test(timeout = 1000) + @Test public void noHandlerExhaustion() throws Exception { dispatcher.addHandler(handler1); dispatcher.addHandler(handler2); @@ -106,7 +108,7 @@ public class RoundRobinDispatcherConcurrentTests { executor.execute(messageSenderTask); } start.countDown(); - allDone.await(); + assertTrue(allDone.await(10, TimeUnit.SECONDS)); assertFalse("not all messages were accepted", failed.get()); verify(handler1, times(TOTAL_EXECUTIONS / 4)).handleMessage(message); verify(handler2, times(TOTAL_EXECUTIONS / 4)).handleMessage(message); @@ -114,7 +116,7 @@ public class RoundRobinDispatcherConcurrentTests { verify(handler4, times(TOTAL_EXECUTIONS / 4)).handleMessage(message); } - @Test(timeout = 2000) + @Test public void unlockOnFailure() throws Exception { // dispatcher has no subscribers (shouldn't lead to deadlock) final CountDownLatch start = new CountDownLatch(1); @@ -140,7 +142,7 @@ public class RoundRobinDispatcherConcurrentTests { executor.execute(messageSenderTask); } start.countDown(); - allDone.await(); + assertTrue(allDone.await(10, TimeUnit.SECONDS)); } @Test @@ -152,30 +154,28 @@ public class RoundRobinDispatcherConcurrentTests { final CountDownLatch allDone = new CountDownLatch(TOTAL_EXECUTIONS); final Message message = this.message; final AtomicBoolean failed = new AtomicBoolean(false); - Runnable messageSenderTask = new Runnable() { - @Override - public void run() { - try { - start.await(); - } - catch (InterruptedException e) { - Thread.currentThread().interrupt(); - } - if (!dispatcher.dispatch(message)) { - failed.set(true); - } - else { - allDone.countDown(); - } + Runnable messageSenderTask = () -> { + try { + start.await(); + } + catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + if (!dispatcher.dispatch(message)) { + failed.set(true); + } + else { + allDone.countDown(); } }; for (int i = 0; i < TOTAL_EXECUTIONS; i++) { executor.execute(messageSenderTask); } start.countDown(); - allDone.await(5000, TimeUnit.MILLISECONDS); + assertTrue(allDone.await(10, TimeUnit.SECONDS)); assertFalse("not all messages were accepted", failed.get()); verify(handler1, times(TOTAL_EXECUTIONS / 2)).handleMessage(message); verify(handler2, times(TOTAL_EXECUTIONS)).handleMessage(message); } + } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/dsl/manualflow/ManualFlowTests.java b/spring-integration-core/src/test/java/org/springframework/integration/dsl/manualflow/ManualFlowTests.java index ebaf260f9d..eafd8fe395 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/dsl/manualflow/ManualFlowTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/dsl/manualflow/ManualFlowTests.java @@ -18,6 +18,7 @@ package org.springframework.integration.dsl.manualflow; import static org.hamcrest.Matchers.containsString; import static org.hamcrest.Matchers.instanceOf; +import static org.hamcrest.Matchers.lessThan; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; @@ -191,8 +192,12 @@ public class ManualFlowTests { assertFalse(this.beanFactory.containsBean(flowRegistration.getId() + BeanFactoryHandler.class.getName() + "#0")); ThreadPoolTaskScheduler taskScheduler = this.beanFactory.getBean(ThreadPoolTaskScheduler.class); - Thread.sleep(100); - assertEquals(0, taskScheduler.getActiveCount()); + + int n = 0; + while (taskScheduler.getActiveCount() > 0 && n++ < 100) { + Thread.sleep(100); + } + assertThat(n, lessThan(100)); assertTrue(additionalBean.destroyed); } diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java index 24dda6f495..c9a26367d1 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/PollerAdviceTests.java @@ -25,13 +25,11 @@ import static org.junit.Assert.assertThat; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.atLeast; -import static org.mockito.Mockito.mock; import static org.mockito.Mockito.spy; import static org.mockito.Mockito.verify; import java.util.ArrayList; import java.util.Collections; -import java.util.Date; import java.util.LinkedList; import java.util.List; import java.util.concurrent.CountDownLatch; @@ -42,7 +40,6 @@ import java.util.concurrent.atomic.AtomicReference; import org.aopalliance.aop.Advice; import org.aopalliance.intercept.Joinpoint; import org.aopalliance.intercept.MethodInterceptor; -import org.junit.Rule; import org.junit.Test; import org.junit.runner.RunWith; @@ -66,7 +63,6 @@ import org.springframework.integration.config.ExpressionControlBusFactoryBean; import org.springframework.integration.core.MessageSource; import org.springframework.integration.scheduling.PollSkipAdvice; import org.springframework.integration.scheduling.SimplePollSkipStrategy; -import org.springframework.integration.test.rule.Log4j2LevelAdjuster; import org.springframework.integration.test.util.OnlyOnceTrigger; import org.springframework.integration.test.util.TestUtils; import org.springframework.integration.util.CompoundTrigger; @@ -94,15 +90,18 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; @DirtiesContext public class PollerAdviceTests { - @Rule - public Log4j2LevelAdjuster adjuster = Log4j2LevelAdjuster.trace(); - @Autowired private MessageChannel control; @Autowired private SimplePollSkipStrategy skipper; + @Autowired + private ThreadPoolTaskScheduler threadPoolTaskScheduler; + + @Autowired + private BeanFactory beanFactory; + @Test public void testDefaultDontSkip() throws Exception { SourcePollingChannelAdapter adapter = new SourcePollingChannelAdapter(); @@ -111,19 +110,9 @@ public class PollerAdviceTests { latch.countDown(); return null; }); - adapter.setTrigger(new Trigger() { - - private boolean done; - - @Override - public Date nextExecutionTime(TriggerContext triggerContext) { - Date date = done ? null : new Date(System.currentTimeMillis() + 10); - done = true; - return date; - } - }); + adapter.setTrigger(new OnlyOnceTrigger()); configure(adapter); - List adviceChain = new ArrayList(); + List adviceChain = new ArrayList<>(); PollSkipAdvice advice = new PollSkipAdvice(); adviceChain.add(advice); adapter.setAdviceChain(adviceChain); @@ -153,18 +142,7 @@ public class PollerAdviceTests { } CountDownLatch latch = new CountDownLatch(1); adapter.setSource(new LocalSource(latch)); - class OneAndDone10msTrigger implements Trigger { - - private boolean done; - - @Override - public Date nextExecutionTime(TriggerContext triggerContext) { - Date date = done ? null : new Date(System.currentTimeMillis() + 10); - done = true; - return date; - } - } - adapter.setTrigger(new OneAndDone10msTrigger()); + adapter.setTrigger(new OnlyOnceTrigger()); configure(adapter); List adviceChain = new ArrayList<>(); SimplePollSkipStrategy skipper = new SimplePollSkipStrategy(); @@ -179,7 +157,7 @@ public class PollerAdviceTests { skipper.reset(); latch = new CountDownLatch(1); adapter.setSource(new LocalSource(latch)); - adapter.setTrigger(new OneAndDone10msTrigger()); + adapter.setTrigger(new OnlyOnceTrigger()); adapter.start(); assertTrue(latch.await(10, TimeUnit.SECONDS)); adapter.stop(); @@ -276,7 +254,7 @@ public class PollerAdviceTests { public void testActiveIdleAdvice() throws Exception { SourcePollingChannelAdapter adapter = new SourcePollingChannelAdapter(); final CountDownLatch latch = new CountDownLatch(5); - final LinkedList triggerPeriods = new LinkedList(); + final LinkedList triggerPeriods = new LinkedList<>(); final DynamicPeriodicTrigger trigger = new DynamicPeriodicTrigger(10); adapter.setSource(() -> { triggerPeriods.add(trigger.getPeriod()); @@ -307,7 +285,7 @@ public class PollerAdviceTests { public void testCompoundTriggerAdvice() throws Exception { SourcePollingChannelAdapter adapter = new SourcePollingChannelAdapter(); final CountDownLatch latch = new CountDownLatch(5); - final LinkedList overridePresent = new LinkedList(); + final LinkedList overridePresent = new LinkedList<>(); final CompoundTrigger compoundTrigger = new CompoundTrigger(new PeriodicTrigger(10)); Trigger override = spy(new PeriodicTrigger(5)); final CompoundTriggerAdvice advice = new CompoundTriggerAdvice(compoundTrigger, override); @@ -336,10 +314,8 @@ public class PollerAdviceTests { private void configure(SourcePollingChannelAdapter adapter) { adapter.setOutputChannel(new NullChannel()); - adapter.setBeanFactory(mock(BeanFactory.class)); - ThreadPoolTaskScheduler scheduler = new ThreadPoolTaskScheduler(); - scheduler.afterPropertiesSet(); - adapter.setTaskScheduler(scheduler); + adapter.setBeanFactory(this.beanFactory); + adapter.setTaskScheduler(this.threadPoolTaskScheduler); } @Test diff --git a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/compound-trigger-context.xml b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/compound-trigger-context.xml index 428830e06b..4f87c01c0c 100644 --- a/spring-integration-core/src/test/java/org/springframework/integration/endpoint/compound-trigger-context.xml +++ b/spring-integration-core/src/test/java/org/springframework/integration/endpoint/compound-trigger-context.xml @@ -22,8 +22,8 @@ - - + +