Handle possible deadlock
- Remove synchronization from scheduleEventQueueProcessing method in executor. Looks like this sync is not really needed and indeed may cause jvm level deadlocks if threads are used for execution. - Fixes #360
This commit is contained in:
@@ -0,0 +1,97 @@
|
||||
/*
|
||||
* Copyright 2017 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
* You may obtain a copy of the License at
|
||||
*
|
||||
* http://www.apache.org/licenses/LICENSE-2.0
|
||||
*
|
||||
* Unless required by applicable law or agreed to in writing, software
|
||||
* distributed under the License is distributed on an "AS IS" BASIS,
|
||||
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
* See the License for the specific language governing permissions and
|
||||
* limitations under the License.
|
||||
*/
|
||||
package org.springframework.statemachine.buildtests;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskExecutor;
|
||||
import org.springframework.statemachine.StateMachine;
|
||||
import org.springframework.statemachine.config.StateMachineBuilder;
|
||||
import org.springframework.statemachine.test.StateMachineTestPlan;
|
||||
import org.springframework.statemachine.test.StateMachineTestPlanBuilder;
|
||||
|
||||
public class TimerSmokeTests {
|
||||
|
||||
private static ThreadPoolTaskExecutor taskExecutor = new ThreadPoolTaskExecutor();
|
||||
{
|
||||
taskExecutor.initialize();
|
||||
}
|
||||
|
||||
private StateMachine<String, String> buildMachine() throws Exception {
|
||||
|
||||
StateMachineBuilder.Builder<String, String> builder = StateMachineBuilder.builder();
|
||||
|
||||
builder.configureConfiguration()
|
||||
.withConfiguration()
|
||||
.taskExecutor(taskExecutor);
|
||||
|
||||
builder.configureStates()
|
||||
.withStates()
|
||||
.initial("initial")
|
||||
.end("end");
|
||||
|
||||
builder.configureTransitions()
|
||||
.withExternal()
|
||||
.source("initial")
|
||||
.target("end")
|
||||
.timerOnce(30)
|
||||
.and()
|
||||
.withLocal()
|
||||
.source("initial")
|
||||
.event("repeate");
|
||||
|
||||
return builder.build();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testNPE() throws Exception {
|
||||
StateMachine<String, String> stateMachine;
|
||||
for (int i = 0; i < 20; i++) {
|
||||
stateMachine = buildMachine();
|
||||
stateMachine.start();
|
||||
while (!stateMachine.isComplete()) {
|
||||
stateMachine.sendEvent("repeate");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testDeadlock() throws Exception {
|
||||
StateMachineTestPlan<String, String> plan;
|
||||
for (int i = 0; i < 20; i++) {
|
||||
plan = StateMachineTestPlanBuilder.<String, String> builder()
|
||||
.defaultAwaitTime(1)
|
||||
.stateMachine(buildMachine())
|
||||
.step()
|
||||
.expectStateMachineStarted(1)
|
||||
.expectStateEntered(1)
|
||||
.expectStateEntered("initial")
|
||||
.and()
|
||||
.step()
|
||||
.sendEvent("repeate")
|
||||
.expectStates("initial")
|
||||
.and()
|
||||
.step()
|
||||
.expectStateEntered(1)
|
||||
.expectStateEntered("end")
|
||||
.and()
|
||||
.step()
|
||||
.expectStateMachineStopped(1)
|
||||
.and()
|
||||
.build();
|
||||
plan.test();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -263,7 +263,7 @@ public class DefaultStateMachineExecutor<S, E> extends LifecycleObjectSupport im
|
||||
stateMachineExecutorTransit.transit(tran, stateContext, queuedMessage);
|
||||
}
|
||||
|
||||
private synchronized void scheduleEventQueueProcessing() {
|
||||
private void scheduleEventQueueProcessing() {
|
||||
TaskExecutor executor = getTaskExecutor();
|
||||
if (executor == null) {
|
||||
return;
|
||||
|
||||
@@ -27,6 +27,8 @@ import org.junit.Test;
|
||||
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.scheduling.TaskScheduler;
|
||||
import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
|
||||
import org.springframework.statemachine.AbstractStateMachineTests;
|
||||
import org.springframework.statemachine.StateContext;
|
||||
import org.springframework.statemachine.StateMachine;
|
||||
@@ -53,15 +55,38 @@ public class ActionAndTimerTests extends AbstractStateMachineTests {
|
||||
assertThat(machine.getState().getIds(), containsInAnyOrder(TestStates.S1));
|
||||
machine.sendEvent(TestEvents.E1);
|
||||
assertThat(machine.getState().getIds(), containsInAnyOrder(TestStates.S2));
|
||||
// sleep so that action with timerOnce(1000) is fired before event is send
|
||||
// event sending is happening on a main thread, but actions are executed on a same
|
||||
// pool than the DefaultStateMachineExecutor is using. existing action execution is interrupted.
|
||||
Thread.sleep(2000);
|
||||
|
||||
assertThat(testTimerAction.latch.await(4, TimeUnit.SECONDS), is(true));
|
||||
assertThat(testTimerAction.e, nullValue());
|
||||
|
||||
machine.sendEvent(TestEvents.E2);
|
||||
assertThat(testListener.s3EnteredLatch.await(2, TimeUnit.SECONDS), is(true));
|
||||
assertThat(machine.getState().getIds(), containsInAnyOrder(TestStates.S3));
|
||||
assertThat(testTimerAction.latch.await(2, TimeUnit.SECONDS), is(true));
|
||||
assertThat(testExitAction.latch.await(2, TimeUnit.SECONDS), is(true));
|
||||
assertThat(testExitAction.e, nullValue());
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@Test
|
||||
public void testExitActionWithTimerOnceThreadPoolTaskScheduler() throws Exception {
|
||||
context.register(Config2.class);
|
||||
context.refresh();
|
||||
StateMachine<TestStates, TestEvents> machine = context.getBean(StateMachine.class);
|
||||
TestTimerAction testTimerAction = context.getBean(TestTimerAction.class);
|
||||
TestExitAction testExitAction = context.getBean(TestExitAction.class);
|
||||
TestListener testListener = new TestListener();
|
||||
machine.addStateListener(testListener);
|
||||
machine.start();
|
||||
assertThat(machine.getState().getIds(), containsInAnyOrder(TestStates.S1));
|
||||
machine.sendEvent(TestEvents.E1);
|
||||
assertThat(machine.getState().getIds(), containsInAnyOrder(TestStates.S2));
|
||||
|
||||
assertThat(testTimerAction.latch.await(4, TimeUnit.SECONDS), is(true));
|
||||
assertThat(testTimerAction.e, nullValue());
|
||||
|
||||
machine.sendEvent(TestEvents.E2);
|
||||
assertThat(testListener.s3EnteredLatch.await(2, TimeUnit.SECONDS), is(true));
|
||||
assertThat(machine.getState().getIds(), containsInAnyOrder(TestStates.S3));
|
||||
assertThat(testExitAction.latch.await(2, TimeUnit.SECONDS), is(true));
|
||||
assertThat(testExitAction.e, nullValue());
|
||||
}
|
||||
@@ -109,6 +134,55 @@ public class ActionAndTimerTests extends AbstractStateMachineTests {
|
||||
}
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableStateMachine
|
||||
static class Config2 extends EnumStateMachineConfigurerAdapter<TestStates, TestEvents> {
|
||||
|
||||
@Override
|
||||
public void configure(StateMachineStateConfigurer<TestStates, TestEvents> states) throws Exception {
|
||||
states
|
||||
.withStates()
|
||||
.initial(TestStates.S1)
|
||||
.state(TestStates.S2, null, testExitAction())
|
||||
.state(TestStates.S3);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void configure(StateMachineTransitionConfigurer<TestStates, TestEvents> transitions) throws Exception {
|
||||
transitions
|
||||
.withExternal()
|
||||
.source(TestStates.S1)
|
||||
.target(TestStates.S2)
|
||||
.event(TestEvents.E1)
|
||||
.and()
|
||||
.withExternal()
|
||||
.source(TestStates.S2)
|
||||
.target(TestStates.S3)
|
||||
.event(TestEvents.E2)
|
||||
.and()
|
||||
.withInternal()
|
||||
.source(TestStates.S2)
|
||||
.action(testTimerAction())
|
||||
.timerOnce(1000);
|
||||
}
|
||||
|
||||
@Bean
|
||||
public TaskScheduler taskScheduler() {
|
||||
ThreadPoolTaskScheduler taskScheduler = new ThreadPoolTaskScheduler();
|
||||
return taskScheduler;
|
||||
}
|
||||
|
||||
@Bean
|
||||
public TestExitAction testExitAction() {
|
||||
return new TestExitAction();
|
||||
}
|
||||
|
||||
@Bean
|
||||
public TestTimerAction testTimerAction() {
|
||||
return new TestTimerAction();
|
||||
}
|
||||
}
|
||||
|
||||
private static class TestListener extends StateMachineListenerAdapter<TestStates, TestEvents> {
|
||||
|
||||
volatile CountDownLatch s3EnteredLatch = new CountDownLatch(1);
|
||||
|
||||
Reference in New Issue
Block a user