Add distibuted action
- New module spring-statemachine-cluster which is based on spring-cloud-cluster to provide leader election. - Ensemble now has a concept of a leader if implementation supports it. - New DistributedLeaderAction can use leader info to execute action only on a leader. - Tweak web sample with these new concepts. - Fixes #176
This commit is contained in:
@@ -0,0 +1,126 @@
|
||||
/*
|
||||
* Copyright 2016 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.cluster;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
import org.apache.curator.framework.CuratorFramework;
|
||||
import org.springframework.cloud.cluster.leader.Context;
|
||||
import org.springframework.cloud.cluster.leader.DefaultCandidate;
|
||||
import org.springframework.cloud.cluster.zk.leader.LeaderInitiator;
|
||||
import org.springframework.statemachine.StateMachine;
|
||||
import org.springframework.statemachine.ensemble.StateMachineEnsemble;
|
||||
import org.springframework.statemachine.zookeeper.ZookeeperStateMachineEnsemble;
|
||||
|
||||
/**
|
||||
* {@link StateMachineEnsemble} backed by a zookeeper and leader functionality
|
||||
* from a Spring Cloud Cluster.
|
||||
*
|
||||
* @author Janne Valkealahti
|
||||
*
|
||||
* @param <S> the type of state
|
||||
* @param <E> the type of event
|
||||
*/
|
||||
public class LeaderZookeeperStateMachineEnsemble<S, E> extends ZookeeperStateMachineEnsemble<S, E> {
|
||||
|
||||
private final Map<StateMachine<S, E>, InitiatorHolder> holders = new HashMap<>();
|
||||
private final CuratorFramework curatorClient;
|
||||
private final String basePath;
|
||||
private StateMachine<S, E> leader;
|
||||
|
||||
/**
|
||||
* Instantiates a new leader zookeeper state machine ensemble.
|
||||
*
|
||||
* @param curatorClient the curator client
|
||||
* @param basePath the base zookeeper path
|
||||
*/
|
||||
public LeaderZookeeperStateMachineEnsemble(CuratorFramework curatorClient, String basePath) {
|
||||
super(curatorClient, basePath);
|
||||
this.curatorClient = curatorClient;
|
||||
this.basePath = basePath;
|
||||
}
|
||||
|
||||
/**
|
||||
* Instantiates a new leader zookeeper state machine ensemble.
|
||||
*
|
||||
* @param curatorClient the curator client
|
||||
* @param basePath the base zookeeper path
|
||||
* @param cleanState if true clean existing state
|
||||
* @param logSize the log size
|
||||
*/
|
||||
public LeaderZookeeperStateMachineEnsemble(CuratorFramework curatorClient, String basePath, boolean cleanState, int logSize) {
|
||||
super(curatorClient, basePath, cleanState, logSize);
|
||||
this.curatorClient = curatorClient;
|
||||
this.basePath = basePath;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void join(StateMachine<S, E> stateMachine) {
|
||||
super.join(stateMachine);
|
||||
StateMachineCandidate candidate = new StateMachineCandidate(stateMachine);
|
||||
LeaderInitiator initiator = new LeaderInitiator(curatorClient, candidate, basePath + "/leader");
|
||||
initiator.start();
|
||||
holders.put(stateMachine, new InitiatorHolder(candidate, initiator));
|
||||
}
|
||||
|
||||
@Override
|
||||
public void leave(StateMachine<S, E> stateMachine) {
|
||||
super.leave(stateMachine);
|
||||
InitiatorHolder holder = holders.get(stateMachine);
|
||||
holder.candidate.yieldLeadership();
|
||||
holder.initiator.stop();
|
||||
holders.remove(stateMachine);
|
||||
}
|
||||
|
||||
@Override
|
||||
public StateMachine<S, E> getLeader() {
|
||||
return leader;
|
||||
}
|
||||
|
||||
private class InitiatorHolder {
|
||||
final StateMachineCandidate candidate;
|
||||
final LeaderInitiator initiator;
|
||||
|
||||
public InitiatorHolder(StateMachineCandidate candidate, LeaderInitiator initiator) {
|
||||
this.candidate = candidate;
|
||||
this.initiator = initiator;
|
||||
}
|
||||
}
|
||||
|
||||
private class StateMachineCandidate extends DefaultCandidate {
|
||||
final StateMachine<S, E> stateMachine;
|
||||
|
||||
public StateMachineCandidate(StateMachine<S, E> stateMachine) {
|
||||
super();
|
||||
this.stateMachine = stateMachine;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onGranted(Context ctx) {
|
||||
super.onGranted(ctx);
|
||||
leader = stateMachine;
|
||||
notifyGranted(stateMachine);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void onRevoked(Context ctx) {
|
||||
super.onRevoked(ctx);
|
||||
leader = null;
|
||||
notifyRevoked(stateMachine);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,104 @@
|
||||
/*
|
||||
* Copyright 2015 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.cluster;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
import org.apache.curator.framework.CuratorFramework;
|
||||
import org.apache.curator.framework.CuratorFrameworkFactory;
|
||||
import org.apache.curator.retry.ExponentialBackoffRetry;
|
||||
import org.apache.curator.test.TestingServer;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.springframework.beans.factory.DisposableBean;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
|
||||
public abstract class AbstractZookeeperTests {
|
||||
|
||||
protected AnnotationConfigApplicationContext context;
|
||||
|
||||
@Before
|
||||
public void setup() {
|
||||
context = buildContext();
|
||||
}
|
||||
|
||||
@After
|
||||
public void clean() {
|
||||
if (context != null) {
|
||||
context.close();
|
||||
}
|
||||
}
|
||||
|
||||
protected AnnotationConfigApplicationContext buildContext() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Configuration
|
||||
protected static class ZkServerConfig {
|
||||
|
||||
@Bean
|
||||
public TestingServerWrapper testingServerWrapper() throws Exception {
|
||||
return new TestingServerWrapper();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Configuration
|
||||
protected static class BaseConfig {
|
||||
|
||||
@Autowired
|
||||
TestingServerWrapper testingServerWrapper;
|
||||
|
||||
@Bean(destroyMethod = "close")
|
||||
public CuratorFramework curatorClient() throws Exception {
|
||||
CuratorFramework client = CuratorFrameworkFactory.builder().defaultData(new byte[0])
|
||||
.retryPolicy(new ExponentialBackoffRetry(1000, 3))
|
||||
.connectString("localhost:" + testingServerWrapper.getPort()).build();
|
||||
// for testing we start it here, thought initiator
|
||||
// is trying to start it if not already done
|
||||
client.start();
|
||||
return client;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
protected static class TestingServerWrapper implements DisposableBean {
|
||||
|
||||
TestingServer testingServer;
|
||||
|
||||
public TestingServerWrapper() throws Exception {
|
||||
this.testingServer = new TestingServer(true);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void destroy() throws Exception {
|
||||
try {
|
||||
testingServer.close();
|
||||
}
|
||||
catch (IOException e) {
|
||||
}
|
||||
}
|
||||
|
||||
public int getPort() {
|
||||
return testingServer.getPort();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,139 @@
|
||||
/*
|
||||
* Copyright 2016 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.cluster;
|
||||
|
||||
import static org.hamcrest.Matchers.contains;
|
||||
import static org.hamcrest.Matchers.is;
|
||||
import static org.junit.Assert.assertThat;
|
||||
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import org.apache.curator.framework.CuratorFramework;
|
||||
import org.junit.Test;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.context.annotation.AnnotationConfigApplicationContext;
|
||||
import org.springframework.context.annotation.Bean;
|
||||
import org.springframework.context.annotation.Configuration;
|
||||
import org.springframework.statemachine.StateMachine;
|
||||
import org.springframework.statemachine.config.EnableStateMachineFactory;
|
||||
import org.springframework.statemachine.config.StateMachineConfigurerAdapter;
|
||||
import org.springframework.statemachine.config.StateMachineFactory;
|
||||
import org.springframework.statemachine.config.builders.StateMachineConfigurationConfigurer;
|
||||
import org.springframework.statemachine.config.builders.StateMachineStateConfigurer;
|
||||
import org.springframework.statemachine.config.builders.StateMachineTransitionConfigurer;
|
||||
import org.springframework.statemachine.ensemble.EnsembleListenerAdapter;
|
||||
import org.springframework.statemachine.ensemble.StateMachineEnsemble;
|
||||
|
||||
public class LeaderZookeeperStateMachineEnsembleTests extends AbstractZookeeperTests {
|
||||
|
||||
@Override
|
||||
protected AnnotationConfigApplicationContext buildContext() {
|
||||
return new AnnotationConfigApplicationContext();
|
||||
}
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void testLeader() throws Exception {
|
||||
context.register(ZkServerConfig.class, BaseConfig.class, Config1.class);
|
||||
context.refresh();
|
||||
|
||||
StateMachineFactory<String, String> factory = context.getBean(StateMachineFactory.class);
|
||||
StateMachineEnsemble<String, String> stateMachineEnsemble = context.getBean(StateMachineEnsemble.class);
|
||||
TestEnsembleListener listener = context.getBean(TestEnsembleListener.class);
|
||||
|
||||
StateMachine<String, String> machine1 = factory.getStateMachine();
|
||||
assertThat(machine1.getState().getIds(), contains("S1"));
|
||||
assertThat(listener.latch.await(2, TimeUnit.SECONDS), is(true));
|
||||
assertThat(stateMachineEnsemble.getLeader(), is(machine1));
|
||||
|
||||
listener.reset(1);
|
||||
StateMachine<String, String> machine2 = factory.getStateMachine();
|
||||
stateMachineEnsemble.leave(machine1);
|
||||
assertThat(listener.latch.await(2, TimeUnit.SECONDS), is(true));
|
||||
assertThat(stateMachineEnsemble.getLeader(), is(machine2));
|
||||
}
|
||||
|
||||
@Configuration
|
||||
@EnableStateMachineFactory
|
||||
static class Config1 extends StateMachineConfigurerAdapter<String, String> {
|
||||
|
||||
@Autowired
|
||||
private CuratorFramework curatorClient;
|
||||
|
||||
@Override
|
||||
public void configure(StateMachineConfigurationConfigurer<String, String> config) throws Exception {
|
||||
config
|
||||
.withConfiguration()
|
||||
.autoStartup(true)
|
||||
.and()
|
||||
.withDistributed()
|
||||
.ensemble(stateMachineEnsemble());
|
||||
}
|
||||
|
||||
@Override
|
||||
public void configure(StateMachineStateConfigurer<String, String> states) throws Exception {
|
||||
states
|
||||
.withStates()
|
||||
.initial("S1")
|
||||
.state("S2")
|
||||
.state("S3");
|
||||
}
|
||||
|
||||
@Override
|
||||
public void configure(StateMachineTransitionConfigurer<String, String> transitions) throws Exception {
|
||||
transitions
|
||||
.withExternal()
|
||||
.source("S1").target("S2")
|
||||
.event("E1")
|
||||
.and()
|
||||
.withExternal()
|
||||
.source("S2").target("S3")
|
||||
.event("E2")
|
||||
.and()
|
||||
.withExternal()
|
||||
.source("S3").target("S1")
|
||||
.event("E3");
|
||||
}
|
||||
|
||||
@Bean
|
||||
public StateMachineEnsemble<String, String> stateMachineEnsemble() throws Exception {
|
||||
LeaderZookeeperStateMachineEnsemble<String,String> ensemble = new LeaderZookeeperStateMachineEnsemble<String, String>(curatorClient, "/foo");
|
||||
ensemble.addEnsembleListener(testEnsembleListener());
|
||||
return ensemble;
|
||||
}
|
||||
|
||||
@Bean
|
||||
public TestEnsembleListener testEnsembleListener() {
|
||||
return new TestEnsembleListener();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
static class TestEnsembleListener extends EnsembleListenerAdapter<String, String> {
|
||||
CountDownLatch latch = new CountDownLatch(1);
|
||||
|
||||
@Override
|
||||
public void ensembleLeaderGranted(StateMachine<String, String> stateMachine) {
|
||||
latch.countDown();
|
||||
}
|
||||
|
||||
void reset(int a1) {
|
||||
latch = new CountDownLatch(a1);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
Reference in New Issue
Block a user