import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertThat;
import static org.junit.Assert.assertTrue;
+import static org.mockito.AdditionalMatchers.or;
import static org.mockito.Mockito.doNothing;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.reset;
import akka.actor.Status.Failure;
import akka.actor.Status.Success;
import akka.cluster.Cluster;
-import akka.testkit.JavaTestKit;
-import com.google.common.base.Optional;
+import akka.pattern.Patterns;
+import akka.util.Timeout;
import com.google.common.base.Stopwatch;
+import com.google.common.collect.ImmutableMap;
import com.google.common.collect.Iterables;
import com.google.common.collect.Lists;
import com.google.common.util.concurrent.Uninterruptibles;
import java.util.Arrays;
import java.util.Collection;
import java.util.List;
+import java.util.Optional;
import java.util.concurrent.TimeUnit;
import org.junit.After;
import org.junit.Before;
import org.opendaylight.controller.cluster.datastore.MemberNode;
import org.opendaylight.controller.cluster.datastore.entityownership.selectionstrategy.EntityOwnerSelectionStrategyConfig;
import org.opendaylight.controller.cluster.datastore.messages.AddShardReplica;
-import org.opendaylight.controller.cluster.raft.RaftState;
+import org.opendaylight.controller.cluster.datastore.messages.ChangeShardMembersVotingStatus;
import org.opendaylight.controller.cluster.raft.policy.DisableElectionsRaftPolicy;
import org.opendaylight.controller.cluster.raft.utils.InMemoryJournal;
import org.opendaylight.controller.cluster.raft.utils.InMemorySnapshotStore;
import org.opendaylight.controller.md.cluster.datastore.model.SchemaContextHelper;
-import org.opendaylight.mdsal.eos.common.api.CandidateAlreadyRegisteredException;
import org.opendaylight.mdsal.eos.common.api.EntityOwnershipState;
import org.opendaylight.mdsal.eos.dom.api.DOMEntity;
import org.opendaylight.mdsal.eos.dom.api.DOMEntityOwnershipCandidateRegistration;
import org.opendaylight.yangtools.yang.data.api.schema.MapNode;
import org.opendaylight.yangtools.yang.data.api.schema.NormalizedNode;
import org.opendaylight.yangtools.yang.model.api.SchemaContext;
+import scala.concurrent.Await;
+import scala.concurrent.Future;
+import scala.concurrent.duration.FiniteDuration;
/**
* End-to-end integration tests for the entity ownership functionality.
@Test
public void testFunctionalityWithThreeNodes() throws Exception {
- String name = "test";
+ String name = "testFunctionalityWithThreeNodes";
MemberNode leaderNode = MemberNode.builder(memberNodes).akkaConfig("Member1").testName(name)
.moduleShardsConfig(MODULE_SHARDS_CONFIG).schemaContext(SCHEMA_CONTEXT).createOperDatastore(false)
.datastoreContextBuilder(leaderDatastoreContextBuilder).build();
follower1EntityOwnershipService.registerCandidate(ENTITY1);
verifyCandidates(leaderDistributedDataStore, ENTITY1, "member-1", "member-2");
verifyOwner(leaderDistributedDataStore, ENTITY1, "member-1");
+ verifyOwner(follower2Node.configDataStore(), ENTITY1, "member-1");
Uninterruptibles.sleepUninterruptibly(300, TimeUnit.MILLISECONDS);
verify(leaderMockListener, never()).ownershipChanged(ownershipChange(ENTITY1));
verify(follower1MockListener, never()).ownershipChanged(ownershipChange(ENTITY1));
follower1EntityOwnershipService.registerCandidate(ENTITY2);
verify(follower1MockListener, timeout(5000)).ownershipChanged(ownershipChange(ENTITY2, false, true, true));
verify(leaderMockListener, timeout(5000)).ownershipChanged(ownershipChange(ENTITY2, false, false, true));
+ verifyOwner(follower2Node.configDataStore(), ENTITY2, "member-2");
reset(leaderMockListener, follower1MockListener);
// Register follower2 candidate for entity2 and verify it gets added but doesn't become owner
follower2EntityOwnershipService.registerListener(ENTITY_TYPE1, follower2MockListener);
- verify(follower2MockListener, timeout(5000)).ownershipChanged(ownershipChange(ENTITY2, false, false, true));
- verify(follower2MockListener, timeout(5000)).ownershipChanged(ownershipChange(ENTITY1, false, false, true));
+ verify(follower2MockListener, timeout(5000).times(2)).ownershipChanged(or(
+ ownershipChange(ENTITY1, false, false, true), ownershipChange(ENTITY2, false, false, true)));
follower2EntityOwnershipService.registerCandidate(ENTITY2);
verifyCandidates(leaderDistributedDataStore, ENTITY2, "member-2", "member-3");
followerDatastoreContextBuilder.shardElectionTimeoutFactor(5)
.customRaftPolicyImplementation(DisableElectionsRaftPolicy.class.getName());
- String name = "test";
+ String name = "testLeaderEntityOwnersReassignedAfterShutdown";
MemberNode leaderNode = MemberNode.builder(memberNodes).akkaConfig("Member1").testName(name)
.moduleShardsConfig(MODULE_SHARDS_CONFIG).schemaContext(SCHEMA_CONTEXT).createOperDatastore(false)
.datastoreContextBuilder(leaderDatastoreContextBuilder).build();
follower1Node.configDataStore().waitTillReady();
follower2Node.configDataStore().waitTillReady();
+ follower1Node.waitForMembersUp("member-1", "member-3");
+
final DOMEntityOwnershipService leaderEntityOwnershipService = newOwnershipService(leaderDistributedDataStore);
final DOMEntityOwnershipService follower1EntityOwnershipService =
newOwnershipService(follower1Node.configDataStore());
verifyCandidates(leaderDistributedDataStore, ENTITY2, "member-1", "member-3");
verifyOwner(leaderDistributedDataStore, ENTITY2, "member-1");
- // Shutdown the leader and verify its removed from the candidate list
-
- leaderNode.cleanup();
- follower1Node.waitForMemberDown("member-1");
-
- // Re-enable elections on follower1 so it becomes the leader
+ // Re-enable elections on all remaining followers so one becomes the new leader
ActorRef follower1Shard = IntegrationTestKit.findLocalShard(follower1Node.configDataStore().getActorContext(),
ENTITY_OWNERSHIP_SHARD_NAME);
follower1Shard.tell(DatastoreContext.newBuilderFrom(followerDatastoreContextBuilder.build())
.customRaftPolicyImplementation(null).build(), ActorRef.noSender());
- MemberNode.verifyRaftState(follower1Node.configDataStore(), ENTITY_OWNERSHIP_SHARD_NAME,
- raftState -> assertEquals("Raft state", RaftState.Leader.toString(), raftState.getRaftState()));
+ ActorRef follower2Shard = IntegrationTestKit.findLocalShard(follower2Node.configDataStore().getActorContext(),
+ ENTITY_OWNERSHIP_SHARD_NAME);
+ follower2Shard.tell(DatastoreContext.newBuilderFrom(followerDatastoreContextBuilder.build())
+ .customRaftPolicyImplementation(null).build(), ActorRef.noSender());
+
+ // Shutdown the leader and verify its removed from the candidate list
+
+ leaderNode.cleanup();
+ follower1Node.waitForMemberDown("member-1");
+ follower2Node.waitForMemberDown("member-1");
// Verify the prior leader's entity owners are re-assigned.
followerDatastoreContextBuilder.shardElectionTimeoutFactor(5)
.customRaftPolicyImplementation(DisableElectionsRaftPolicy.class.getName());
- String name = "test";
+ String name = "testLeaderAndFollowerEntityOwnersReassignedAfterShutdown";
final MemberNode leaderNode = MemberNode.builder(memberNodes).akkaConfig("Member1")
.useAkkaArtery(false).testName(name)
.moduleShardsConfig(MODULE_SHARDS_5_NODE_CONFIG).schemaContext(SCHEMA_CONTEXT)
leaderDistributedDataStore.waitTillReady();
follower1Node.configDataStore().waitTillReady();
follower2Node.configDataStore().waitTillReady();
+ follower3Node.configDataStore().waitTillReady();
+ follower4Node.configDataStore().waitTillReady();
+
+ leaderNode.waitForMembersUp("member-2", "member-3", "member-4", "member-5");
+ follower1Node.waitForMembersUp("member-1", "member-3", "member-4", "member-5");
final DOMEntityOwnershipService leaderEntityOwnershipService = newOwnershipService(leaderDistributedDataStore);
final DOMEntityOwnershipService follower1EntityOwnershipService =
verifyCandidates(leaderDistributedDataStore, ENTITY2, "member-1", "member-3", "member-4");
verifyOwner(leaderDistributedDataStore, ENTITY2, "member-1");
+ // Re-enable elections on all remaining followers so one becomes the new leader
+
+ ActorRef follower1Shard = IntegrationTestKit.findLocalShard(follower1Node.configDataStore().getActorContext(),
+ ENTITY_OWNERSHIP_SHARD_NAME);
+ follower1Shard.tell(DatastoreContext.newBuilderFrom(followerDatastoreContextBuilder.build())
+ .customRaftPolicyImplementation(null).build(), ActorRef.noSender());
+
+ ActorRef follower2Shard = IntegrationTestKit.findLocalShard(follower2Node.configDataStore().getActorContext(),
+ ENTITY_OWNERSHIP_SHARD_NAME);
+ follower2Shard.tell(DatastoreContext.newBuilderFrom(followerDatastoreContextBuilder.build())
+ .customRaftPolicyImplementation(null).build(), ActorRef.noSender());
+
+ ActorRef follower4Shard = IntegrationTestKit.findLocalShard(follower4Node.configDataStore().getActorContext(),
+ ENTITY_OWNERSHIP_SHARD_NAME);
+ follower4Shard.tell(DatastoreContext.newBuilderFrom(followerDatastoreContextBuilder.build())
+ .customRaftPolicyImplementation(null).build(), ActorRef.noSender());
+
// Shutdown the leader and follower3
leaderNode.cleanup();
follower1Node.waitForMemberDown("member-1");
follower1Node.waitForMemberDown("member-4");
-
- // Re-enable elections on follower1 so it becomes the leader
-
- ActorRef follower1Shard = IntegrationTestKit.findLocalShard(follower1Node.configDataStore().getActorContext(),
- ENTITY_OWNERSHIP_SHARD_NAME);
- follower1Shard.tell(DatastoreContext.newBuilderFrom(followerDatastoreContextBuilder.build())
- .customRaftPolicyImplementation(null).build(), ActorRef.noSender());
-
- MemberNode.verifyRaftState(follower1Node.configDataStore(), ENTITY_OWNERSHIP_SHARD_NAME,
- raftState -> assertEquals("Raft state", RaftState.Leader.toString(), raftState.getRaftState()));
+ follower2Node.waitForMemberDown("member-1");
+ follower2Node.waitForMemberDown("member-4");
+ follower4Node.waitForMemberDown("member-1");
+ follower4Node.waitForMemberDown("member-4");
// Verify the prior leader's and follower3 entity owners are re-assigned.
* Reproduces bug <a href="https://bugs.opendaylight.org/show_bug.cgi?id=4554">4554</a>.
*/
@Test
- public void testCloseCandidateRegistrationInQuickSuccession() throws CandidateAlreadyRegisteredException {
+ public void testCloseCandidateRegistrationInQuickSuccession() throws Exception {
String name = "testCloseCandidateRegistrationInQuickSuccession";
MemberNode leaderNode = MemberNode.builder(memberNodes).akkaConfig("Member1").testName(name)
.moduleShardsConfig(MODULE_SHARDS_CONFIG).schemaContext(SCHEMA_CONTEXT).createOperDatastore(false)
private static Optional<DOMEntityOwnershipChange> getValueSafely(ArgumentCaptor<DOMEntityOwnershipChange> captor) {
try {
- return Optional.fromNullable(captor.getValue());
+ return Optional.ofNullable(captor.getValue());
} catch (MockitoException e) {
// No value was captured
- return Optional.absent();
+ return Optional.empty();
}
}
AddShardReplica addReplica = new AddShardReplica(ENTITY_OWNERSHIP_SHARD_NAME);
follower1DistributedDataStore.getActorContext().getShardManager().tell(addReplica,
follower1Node.kit().getRef());
- Object reply = follower1Node.kit().expectMsgAnyClassOf(JavaTestKit.duration("5 sec"),
+ Object reply = follower1Node.kit().expectMsgAnyClassOf(follower1Node.kit().duration("5 sec"),
Success.class, Failure.class);
if (reply instanceof Failure) {
throw new AssertionError("AddShardReplica failed", ((Failure)reply).cause());
@Test
public void testOwnerSelectedOnRapidUnregisteringAndRegisteringOfCandidates() throws Exception {
- String name = "test";
+ String name = "testOwnerSelectedOnRapidUnregisteringAndRegisteringOfCandidates";
MemberNode leaderNode = MemberNode.builder(memberNodes).akkaConfig("Member1").testName(name)
.moduleShardsConfig(MODULE_SHARDS_CONFIG).schemaContext(SCHEMA_CONTEXT).createOperDatastore(false)
.datastoreContextBuilder(leaderDatastoreContextBuilder).build();
@Test
public void testOwnerSelectedOnRapidRegisteringAndUnregisteringOfCandidates() throws Exception {
- String name = "test";
+ String name = "testOwnerSelectedOnRapidRegisteringAndUnregisteringOfCandidates";
MemberNode leaderNode = MemberNode.builder(memberNodes).akkaConfig("Member1").testName(name)
.moduleShardsConfig(MODULE_SHARDS_CONFIG).schemaContext(SCHEMA_CONTEXT).createOperDatastore(false)
.datastoreContextBuilder(leaderDatastoreContextBuilder).build();
verifyOwner(leaderDistributedDataStore, ENTITY1, "member-2");
}
+ @Test
+ public void testEntityOwnershipWithNonVotingMembers() throws Exception {
+ followerDatastoreContextBuilder.shardElectionTimeoutFactor(5)
+ .customRaftPolicyImplementation(DisableElectionsRaftPolicy.class.getName());
+
+ String name = "testEntityOwnershipWithNonVotingMembers";
+ final MemberNode member1LeaderNode = MemberNode.builder(memberNodes).akkaConfig("Member1")
+ .useAkkaArtery(false).testName(name)
+ .moduleShardsConfig(MODULE_SHARDS_5_NODE_CONFIG).schemaContext(SCHEMA_CONTEXT)
+ .createOperDatastore(false).datastoreContextBuilder(leaderDatastoreContextBuilder).build();
+
+ final MemberNode member2FollowerNode = MemberNode.builder(memberNodes).akkaConfig("Member2")
+ .useAkkaArtery(false).testName(name)
+ .moduleShardsConfig(MODULE_SHARDS_5_NODE_CONFIG).schemaContext(SCHEMA_CONTEXT)
+ .createOperDatastore(false).datastoreContextBuilder(followerDatastoreContextBuilder).build();
+
+ final MemberNode member3FollowerNode = MemberNode.builder(memberNodes).akkaConfig("Member3")
+ .useAkkaArtery(false).testName(name)
+ .moduleShardsConfig(MODULE_SHARDS_5_NODE_CONFIG).schemaContext(SCHEMA_CONTEXT)
+ .createOperDatastore(false).datastoreContextBuilder(followerDatastoreContextBuilder).build();
+
+ final MemberNode member4FollowerNode = MemberNode.builder(memberNodes).akkaConfig("Member4")
+ .useAkkaArtery(false).testName(name)
+ .moduleShardsConfig(MODULE_SHARDS_5_NODE_CONFIG).schemaContext(SCHEMA_CONTEXT)
+ .createOperDatastore(false).datastoreContextBuilder(followerDatastoreContextBuilder).build();
+
+ final MemberNode member5FollowerNode = MemberNode.builder(memberNodes).akkaConfig("Member5")
+ .useAkkaArtery(false).testName(name)
+ .moduleShardsConfig(MODULE_SHARDS_5_NODE_CONFIG).schemaContext(SCHEMA_CONTEXT)
+ .createOperDatastore(false).datastoreContextBuilder(followerDatastoreContextBuilder).build();
+
+ AbstractDataStore leaderDistributedDataStore = member1LeaderNode.configDataStore();
+
+ leaderDistributedDataStore.waitTillReady();
+ member2FollowerNode.configDataStore().waitTillReady();
+ member3FollowerNode.configDataStore().waitTillReady();
+ member4FollowerNode.configDataStore().waitTillReady();
+ member5FollowerNode.configDataStore().waitTillReady();
+
+ member1LeaderNode.waitForMembersUp("member-2", "member-3", "member-4", "member-5");
+
+ final DOMEntityOwnershipService member3EntityOwnershipService =
+ newOwnershipService(member3FollowerNode.configDataStore());
+ final DOMEntityOwnershipService member4EntityOwnershipService =
+ newOwnershipService(member4FollowerNode.configDataStore());
+ final DOMEntityOwnershipService member5EntityOwnershipService =
+ newOwnershipService(member5FollowerNode.configDataStore());
+
+ newOwnershipService(member1LeaderNode.configDataStore());
+ member1LeaderNode.kit().waitUntilLeader(member1LeaderNode.configDataStore().getActorContext(),
+ ENTITY_OWNERSHIP_SHARD_NAME);
+
+ // Make member4 and member5 non-voting
+
+ Future<Object> future = Patterns.ask(leaderDistributedDataStore.getActorContext().getShardManager(),
+ new ChangeShardMembersVotingStatus(ENTITY_OWNERSHIP_SHARD_NAME,
+ ImmutableMap.of("member-4", Boolean.FALSE, "member-5", Boolean.FALSE)),
+ new Timeout(10, TimeUnit.SECONDS));
+ Object response = Await.result(future, FiniteDuration.apply(10, TimeUnit.SECONDS));
+ if (response instanceof Throwable) {
+ throw new AssertionError("ChangeShardMembersVotingStatus failed", (Throwable)response);
+ }
+
+ assertNull("Expected null Success response. Actual " + response, response);
+
+ // Register member4 candidate for entity1 - it should not become owner since it's non-voting
+
+ member4EntityOwnershipService.registerCandidate(ENTITY1);
+ verifyCandidates(leaderDistributedDataStore, ENTITY1, "member-4");
+
+ // Register member5 candidate for entity2 - it should not become owner since it's non-voting
+
+ member5EntityOwnershipService.registerCandidate(ENTITY2);
+ verifyCandidates(leaderDistributedDataStore, ENTITY2, "member-5");
+
+ Uninterruptibles.sleepUninterruptibly(500, TimeUnit.MILLISECONDS);
+ verifyOwner(leaderDistributedDataStore, ENTITY1, "");
+ verifyOwner(leaderDistributedDataStore, ENTITY2, "");
+
+ // Register member3 candidate for entity1 - it should become owner since it's voting
+
+ member3EntityOwnershipService.registerCandidate(ENTITY1);
+ verifyCandidates(leaderDistributedDataStore, ENTITY1, "member-4", "member-3");
+ verifyOwner(leaderDistributedDataStore, ENTITY1, "member-3");
+
+ // Switch member4 and member5 back to voting and member3 non-voting. This should result in member4 and member5
+ // to become entity owners.
+
+ future = Patterns.ask(leaderDistributedDataStore.getActorContext().getShardManager(),
+ new ChangeShardMembersVotingStatus(ENTITY_OWNERSHIP_SHARD_NAME,
+ ImmutableMap.of("member-3", Boolean.FALSE, "member-4", Boolean.TRUE, "member-5", Boolean.TRUE)),
+ new Timeout(10, TimeUnit.SECONDS));
+ response = Await.result(future, FiniteDuration.apply(10, TimeUnit.SECONDS));
+ if (response instanceof Throwable) {
+ throw new AssertionError("ChangeShardMembersVotingStatus failed", (Throwable)response);
+ }
+
+ assertNull("Expected null Success response. Actual " + response, response);
+
+ verifyOwner(leaderDistributedDataStore, ENTITY1, "member-4");
+ verifyOwner(leaderDistributedDataStore, ENTITY2, "member-5");
+ }
+
private static void verifyGetOwnershipState(final DOMEntityOwnershipService service, final DOMEntity entity,
final EntityOwnershipState expState) {
Optional<EntityOwnershipState> state = service.getOwnershipState(entity);
- assertEquals("getOwnershipState present", true, state.isPresent());
+ assertTrue("getOwnershipState present", state.isPresent());
assertEquals("EntityOwnershipState", expState, state.get());
}
.read(entityPath(entity.getType(), entity.getIdentifier()).node(Candidate.QNAME))
.get(5, TimeUnit.SECONDS);
try {
- assertEquals("Candidates not found for " + entity, true, possible.isPresent());
+ assertTrue("Candidates not found for " + entity, possible.isPresent());
Collection<String> actual = new ArrayList<>();
for (MapEntryNode candidate: ((MapNode)possible.get()).getValue()) {
actual.add(candidate.getChild(CANDIDATE_NAME_NODE_ID).get().getValue().toString());