package org.opendaylight.controller.cluster.datastore;
import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertSame;
import static org.junit.Assert.assertTrue;
import com.google.common.base.Optional;
import com.google.common.collect.ImmutableMap;
import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Sets;
import com.google.common.util.concurrent.Uninterruptibles;
import com.typesafe.config.ConfigFactory;
+import java.net.URI;
import java.util.Arrays;
import java.util.List;
import java.util.Map;
import org.junit.Test;
import org.mockito.Mock;
import org.mockito.MockitoAnnotations;
+import org.opendaylight.controller.cluster.datastore.config.Configuration;
+import org.opendaylight.controller.cluster.datastore.config.ConfigurationImpl;
+import org.opendaylight.controller.cluster.datastore.config.EmptyModuleShardConfigProvider;
+import org.opendaylight.controller.cluster.datastore.config.ModuleShardConfiguration;
import org.opendaylight.controller.cluster.datastore.exceptions.NoShardLeaderException;
import org.opendaylight.controller.cluster.datastore.exceptions.NotInitializedException;
import org.opendaylight.controller.cluster.datastore.exceptions.PrimaryNotFoundException;
import org.opendaylight.controller.cluster.datastore.identifiers.ShardIdentifier;
import org.opendaylight.controller.cluster.datastore.identifiers.ShardManagerIdentifier;
import org.opendaylight.controller.cluster.datastore.messages.ActorInitialized;
+import org.opendaylight.controller.cluster.datastore.messages.CreateShard;
+import org.opendaylight.controller.cluster.datastore.messages.CreateShardReply;
import org.opendaylight.controller.cluster.datastore.messages.FindLocalShard;
import org.opendaylight.controller.cluster.datastore.messages.FindPrimary;
import org.opendaylight.controller.cluster.datastore.messages.LocalPrimaryShardFound;
import org.opendaylight.controller.cluster.datastore.messages.LocalShardFound;
import org.opendaylight.controller.cluster.datastore.messages.LocalShardNotFound;
+import org.opendaylight.controller.cluster.datastore.messages.PeerDown;
+import org.opendaylight.controller.cluster.datastore.messages.PeerUp;
import org.opendaylight.controller.cluster.datastore.messages.PrimaryShardInfo;
import org.opendaylight.controller.cluster.datastore.messages.RemotePrimaryShardFound;
import org.opendaylight.controller.cluster.datastore.messages.ShardLeaderStateChanged;
+import org.opendaylight.controller.cluster.datastore.messages.SwitchShardBehavior;
import org.opendaylight.controller.cluster.datastore.messages.UpdateSchemaContext;
-import org.opendaylight.controller.cluster.datastore.utils.MessageCollectorActor;
import org.opendaylight.controller.cluster.datastore.utils.MockClusterWrapper;
import org.opendaylight.controller.cluster.datastore.utils.MockConfiguration;
import org.opendaylight.controller.cluster.datastore.utils.PrimaryShardInfoFutureCache;
import org.opendaylight.controller.cluster.notifications.RoleChangeNotification;
import org.opendaylight.controller.cluster.raft.RaftState;
import org.opendaylight.controller.cluster.raft.base.messages.FollowerInitialSyncUpStatus;
+import org.opendaylight.controller.cluster.raft.base.messages.SwitchBehavior;
import org.opendaylight.controller.cluster.raft.utils.InMemoryJournal;
+import org.opendaylight.controller.cluster.raft.utils.MessageCollectorActor;
import org.opendaylight.controller.md.cluster.datastore.model.TestModel;
import org.opendaylight.yangtools.yang.data.api.schema.tree.DataTree;
import org.opendaylight.yangtools.yang.model.api.SchemaContext;
private static TestActorRef<MessageCollectorActor> mockShardActor;
+ private static String mockShardName;
+
private final DatastoreContext.Builder datastoreContextBuilder = DatastoreContext.newBuilder().
dataStoreType(shardMrgIDSuffix).shardInitializationTimeout(600, TimeUnit.MILLISECONDS)
.shardHeartbeatIntervalInMillis(100).shardElectionTimeoutFactor(6);
InMemoryJournal.clear();
if(mockShardActor == null) {
- String name = new ShardIdentifier(Shard.DEFAULT_NAME, "member-1", "config").toString();
- mockShardActor = TestActorRef.create(getSystem(), Props.create(MessageCollectorActor.class), name);
+ mockShardName = new ShardIdentifier(Shard.DEFAULT_NAME, "member-1", "config").toString();
+ mockShardActor = TestActorRef.create(getSystem(), Props.create(MessageCollectorActor.class), mockShardName);
}
mockShardActor.underlyingActor().clear();
InMemoryJournal.clear();
}
- private Props newShardMgrProps(boolean persistent) {
- return ShardManager.props(new MockClusterWrapper(), new MockConfiguration(),
- datastoreContextBuilder.persistent(persistent).build(), ready, primaryShardInfoCache);
+ private Props newShardMgrProps() {
+ return newShardMgrProps(new MockConfiguration());
+ }
+
+ private Props newShardMgrProps(Configuration config) {
+ return ShardManager.props(new MockClusterWrapper(), config, datastoreContextBuilder.build(), ready,
+ primaryShardInfoCache);
}
private Props newPropsShardMgrWithMockShardActor() {
shardManager1.underlyingActor().waitForUnreachableMember();
+ PeerDown peerDown = MessageCollectorActor.expectFirstMatching(mockShardActor1, PeerDown.class);
+ assertEquals("getMemberName", "member-2", peerDown.getMemberName());
+ MessageCollectorActor.clearMessages(mockShardActor1);
+
+ shardManager1.underlyingActor().onReceiveCommand(MockClusterWrapper.
+ createMemberRemoved("member-2", "akka.tcp://cluster-test@127.0.0.1:2558"));
+
+ MessageCollectorActor.expectFirstMatching(mockShardActor1, PeerDown.class);
+
shardManager1.tell(new FindPrimary("default", true), getRef());
expectMsgClass(duration("5 seconds"), NoShardLeaderException.class);
shardManager1.underlyingActor().waitForReachableMember();
+ PeerUp peerUp = MessageCollectorActor.expectFirstMatching(mockShardActor1, PeerUp.class);
+ assertEquals("getMemberName", "member-2", peerUp.getMemberName());
+ MessageCollectorActor.clearMessages(mockShardActor1);
+
shardManager1.tell(new FindPrimary("default", true), getRef());
RemotePrimaryShardFound found1 = expectMsgClass(duration("5 seconds"), RemotePrimaryShardFound.class);
String path1 = found1.getPrimaryPath();
assertTrue("Unexpected primary path " + path1, path1.contains("member-2-shard-default-config"));
+ shardManager1.underlyingActor().onReceiveCommand(MockClusterWrapper.
+ createMemberUp("member-2", "akka.tcp://cluster-test@127.0.0.1:2558"));
+
+ MessageCollectorActor.expectFirstMatching(mockShardActor1, PeerUp.class);
+
}};
JavaTestKit.shutdownActorSystem(system1);
public void testRoleChangeNotificationAndShardLeaderStateChangedReleaseReady() throws Exception {
new JavaTestKit(getSystem()) {
{
- TestActorRef<ShardManager> shardManager = TestActorRef.create(getSystem(), newShardMgrProps(true));
+ TestActorRef<ShardManager> shardManager = TestActorRef.create(getSystem(), newShardMgrProps());
String memberId = "member-1-shard-default-" + shardMrgIDSuffix;
shardManager.underlyingActor().onReceiveCommand(new RoleChangeNotification(
public void testRoleChangeNotificationToFollowerWithShardLeaderStateChangedReleaseReady() throws Exception {
new JavaTestKit(getSystem()) {
{
- TestActorRef<ShardManager> shardManager = TestActorRef.create(getSystem(), newShardMgrProps(true));
+ TestActorRef<ShardManager> shardManager = TestActorRef.create(getSystem(), newShardMgrProps());
String memberId = "member-1-shard-default-" + shardMrgIDSuffix;
shardManager.underlyingActor().onReceiveCommand(new RoleChangeNotification(
public void testReadyCountDownForMemberUpAfterLeaderStateChanged() throws Exception {
new JavaTestKit(getSystem()) {
{
- TestActorRef<ShardManager> shardManager = TestActorRef.create(getSystem(), newShardMgrProps(true));
+ TestActorRef<ShardManager> shardManager = TestActorRef.create(getSystem(), newShardMgrProps());
String memberId = "member-1-shard-default-" + shardMrgIDSuffix;
shardManager.underlyingActor().onReceiveCommand(new RoleChangeNotification(
public void testRoleChangeNotificationDoNothingForUnknownShard() throws Exception {
new JavaTestKit(getSystem()) {
{
- TestActorRef<ShardManager> shardManager = TestActorRef.create(getSystem(), newShardMgrProps(true));
+ TestActorRef<ShardManager> shardManager = TestActorRef.create(getSystem(), newShardMgrProps());
shardManager.underlyingActor().onReceiveCommand(new RoleChangeNotification(
"unknown", RaftState.Candidate.name(), RaftState.Leader.name()));
@Test
public void testByDefaultSyncStatusIsFalse() throws Exception{
- final Props persistentProps = newShardMgrProps(true);
+ final Props persistentProps = newShardMgrProps();
final TestActorRef<ShardManager> shardManager =
TestActorRef.create(getSystem(), persistentProps);
@Test
public void testWhenShardIsCandidateSyncStatusIsFalse() throws Exception{
- final Props persistentProps = newShardMgrProps(true);
+ final Props persistentProps = newShardMgrProps();
final TestActorRef<ShardManager> shardManager =
TestActorRef.create(getSystem(), persistentProps);
}
+ @Test
+ public void testOnReceiveSwitchShardBehavior() throws Exception {
+ new JavaTestKit(getSystem()) {{
+ final ActorRef shardManager = getSystem().actorOf(newPropsShardMgrWithMockShardActor());
+
+ shardManager.tell(new UpdateSchemaContext(TestModel.createTestContext()), getRef());
+ shardManager.tell(new ActorInitialized(), mockShardActor);
+
+ shardManager.tell(new SwitchShardBehavior(mockShardName, "Leader", 1000), getRef());
+
+ SwitchBehavior switchBehavior = MessageCollectorActor.expectFirstMatching(mockShardActor, SwitchBehavior.class);
+
+ assertEquals(RaftState.Leader, switchBehavior.getNewState());
+ assertEquals(1000, switchBehavior.getNewTerm());
+ }};
+ }
+
+ public void testOnReceiveCreateShard() {
+ new JavaTestKit(getSystem()) {{
+ datastoreContextBuilder.shardInitializationTimeout(1, TimeUnit.MINUTES).persistent(true);
+
+ ActorRef shardManager = getSystem().actorOf(newShardMgrProps(
+ new ConfigurationImpl(new EmptyModuleShardConfigProvider())));
+
+ SchemaContext schemaContext = TestModel.createTestContext();
+ shardManager.tell(new UpdateSchemaContext(schemaContext), ActorRef.noSender());
+
+ DatastoreContext datastoreContext = DatastoreContext.newBuilder().shardElectionTimeoutFactor(100).
+ persistent(false).build();
+ TestShardPropsCreator shardPropsCreator = new TestShardPropsCreator();
+
+ ModuleShardConfiguration config = new ModuleShardConfiguration(URI.create("foo-ns"), "foo-module",
+ "foo", null, Arrays.asList("member-1", "member-5", "member-6"));
+ shardManager.tell(new CreateShard(config, shardPropsCreator, datastoreContext), getRef());
+
+ expectMsgClass(duration("5 seconds"), CreateShardReply.class);
+
+ shardManager.tell(new FindLocalShard("foo", true), getRef());
+
+ expectMsgClass(duration("5 seconds"), LocalShardFound.class);
+
+ assertEquals("isRecoveryApplicable", false, shardPropsCreator.datastoreContext.isPersistent());
+ assertEquals("peerMembers", Sets.newHashSet(new ShardIdentifier("foo", "member-5", shardMrgIDSuffix).toString(),
+ new ShardIdentifier("foo", "member-6", shardMrgIDSuffix).toString()),
+ shardPropsCreator.peerAddresses.keySet());
+ assertEquals("ShardIdentifier", new ShardIdentifier("foo", "member-1", shardMrgIDSuffix),
+ shardPropsCreator.shardId);
+ assertSame("schemaContext", schemaContext, shardPropsCreator.schemaContext);
+
+ // Send CreateShard with same name - should fail.
+
+ shardManager.tell(new CreateShard(config, shardPropsCreator, null), getRef());
+
+ expectMsgClass(duration("5 seconds"), akka.actor.Status.Failure.class);
+ }};
+ }
+
+ @Test
+ public void testOnReceiveCreateShardWithNoInitialSchemaContext() {
+ new JavaTestKit(getSystem()) {{
+ ActorRef shardManager = getSystem().actorOf(newShardMgrProps(
+ new ConfigurationImpl(new EmptyModuleShardConfigProvider())));
+
+ TestShardPropsCreator shardPropsCreator = new TestShardPropsCreator();
+
+ ModuleShardConfiguration config = new ModuleShardConfiguration(URI.create("foo-ns"), "foo-module",
+ "foo", null, Arrays.asList("member-1"));
+ shardManager.tell(new CreateShard(config, shardPropsCreator, null), getRef());
+
+ expectMsgClass(duration("5 seconds"), CreateShardReply.class);
+
+ SchemaContext schemaContext = TestModel.createTestContext();
+ shardManager.tell(new UpdateSchemaContext(schemaContext), ActorRef.noSender());
+
+ shardManager.tell(new FindLocalShard("foo", true), getRef());
+
+ expectMsgClass(duration("5 seconds"), LocalShardFound.class);
+
+ assertSame("schemaContext", schemaContext, shardPropsCreator.schemaContext);
+ assertNotNull("schemaContext is null", shardPropsCreator.datastoreContext);
+ }};
+ }
+
+ private static class TestShardPropsCreator implements ShardPropsCreator {
+ ShardIdentifier shardId;
+ Map<String, String> peerAddresses;
+ SchemaContext schemaContext;
+ DatastoreContext datastoreContext;
+
+ @Override
+ public Props newProps(ShardIdentifier shardId, Map<String, String> peerAddresses,
+ DatastoreContext datastoreContext, SchemaContext schemaContext) {
+ this.shardId = shardId;
+ this.peerAddresses = peerAddresses;
+ this.schemaContext = schemaContext;
+ this.datastoreContext = datastoreContext;
+ return Shard.props(shardId, peerAddresses, datastoreContext, schemaContext);
+ }
+
+ }
+
private static class TestShardManager extends ShardManager {
private final CountDownLatch recoveryComplete = new CountDownLatch(1);