Fix some sporadic tests failures
* Increase timeouts in the `RoundRobinDispatcherConcurrentTests` and `ManualFlowTests` * Fix `PollerAdviceTests` to re-use `TaskScheduler` from the ctx instead of local, not closed instance * Use `OnlyOnceTrigger` instead of local implementations * Change the `primary` `Trigger` bean to the `PeriodicTrigger` as well. The minimum interval for the `CronTrigger` is 1 seconds - it doesn't matter for this test-case
This commit is contained in:
@@ -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);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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<Advice> adviceChain = new ArrayList<Advice>();
|
||||
List<Advice> 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<Advice> 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<Long> triggerPeriods = new LinkedList<Long>();
|
||||
final LinkedList<Long> 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<Object> overridePresent = new LinkedList<Object>();
|
||||
final LinkedList<Object> 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
|
||||
|
||||
@@ -22,8 +22,8 @@
|
||||
<constructor-arg ref="primary" />
|
||||
</bean>
|
||||
|
||||
<bean id="primary" class="org.springframework.scheduling.support.CronTrigger">
|
||||
<constructor-arg value="*/1 * * * * *" />
|
||||
<bean id="primary" class="org.springframework.scheduling.support.PeriodicTrigger">
|
||||
<constructor-arg value="10" />
|
||||
</bean>
|
||||
|
||||
<bean id="secondary" class="org.springframework.scheduling.support.PeriodicTrigger">
|
||||
|
||||
Reference in New Issue
Block a user