import org.junit.Before;
import org.junit.Test;
import org.mockito.Mockito;
+import org.mockito.invocation.InvocationOnMock;
+import org.mockito.stubbing.Answer;
+import org.opendaylight.controller.cluster.datastore.DatastoreContext.Builder;
import org.opendaylight.controller.cluster.datastore.exceptions.NoShardLeaderException;
import org.opendaylight.controller.cluster.datastore.exceptions.ShardLeaderNotRespondingException;
import org.opendaylight.controller.cluster.datastore.messages.CommitTransactionReply;
import org.opendaylight.yangtools.yang.data.api.schema.NormalizedNode;
import org.opendaylight.yangtools.yang.data.api.schema.tree.DataTreeModification;
import org.opendaylight.yangtools.yang.data.api.schema.tree.TipProducingDataTree;
+import org.opendaylight.yangtools.yang.data.api.schema.tree.TreeType;
import org.opendaylight.yangtools.yang.data.impl.schema.ImmutableNodes;
import org.opendaylight.yangtools.yang.data.impl.schema.builder.api.CollectionNodeBuilder;
import org.opendaylight.yangtools.yang.data.impl.schema.builder.impl.ImmutableContainerNodeBuilder;
leaderTestKit.waitUntilLeader(leaderDistributedDataStore.getActorContext(), SHARD_NAMES);
}
- private void verifyCars(DOMStoreReadTransaction readTx, MapEntryNode... entries) throws Exception {
+ private static void verifyCars(DOMStoreReadTransaction readTx, MapEntryNode... entries) throws Exception {
Optional<NormalizedNode<?, ?>> optional = readTx.read(CarsModel.CAR_LIST_PATH).get(5, TimeUnit.SECONDS);
assertEquals("isPresent", true, optional.isPresent());
assertEquals("Car list node", listBuilder.build(), optional.get());
}
- private void verifyNode(DOMStoreReadTransaction readTx, YangInstanceIdentifier path, NormalizedNode<?, ?> expNode)
+ private static void verifyNode(DOMStoreReadTransaction readTx, YangInstanceIdentifier path, NormalizedNode<?, ?> expNode)
throws Exception {
Optional<NormalizedNode<?, ?>> optional = readTx.read(path).get(5, TimeUnit.SECONDS);
assertEquals("isPresent", true, optional.isPresent());
assertEquals("Data node", expNode, optional.get());
}
- private void verifyExists(DOMStoreReadTransaction readTx, YangInstanceIdentifier path) throws Exception {
+ private static void verifyExists(DOMStoreReadTransaction readTx, YangInstanceIdentifier path) throws Exception {
Boolean exists = readTx.exists(path).get(5, TimeUnit.SECONDS);
assertEquals("exists", true, exists);
}
// Switch the leader to the follower
followerDatastoreContextBuilder.shardElectionTimeoutFactor(1);
- followerDistributedDataStore.onDatastoreContextUpdated(followerDatastoreContextBuilder.build());
+ sendDatastoreContextUpdate(followerDistributedDataStore, followerDatastoreContextBuilder);
JavaTestKit.shutdownActorSystem(leaderSystem, null, true);
@Test
public void testReadyLocalTransactionForwardedToLeader() throws Exception {
initDatastores("testReadyLocalTransactionForwardedToLeader");
+ followerTestKit.waitUntilLeader(followerDistributedDataStore.getActorContext(), "cars");
Optional<ActorRef> carsFollowerShard = followerDistributedDataStore.getActorContext().findLocalShard("cars");
assertEquals("Cars follower shard found", true, carsFollowerShard.isPresent());
- TipProducingDataTree dataTree = InMemoryDataTreeFactory.getInstance().create();
+ TipProducingDataTree dataTree = InMemoryDataTreeFactory.getInstance().create(TreeType.OPERATIONAL);
dataTree.setSchemaContext(SchemaContextHelper.full());
DataTreeModification modification = dataTree.takeSnapshot().newModification();
ReadyLocalTransaction readyLocal = new ReadyLocalTransaction(transactionID , modification, true);
carsFollowerShard.get().tell(readyLocal, followerTestKit.getRef());
- followerTestKit.expectMsgClass(CommitTransactionReply.SERIALIZABLE_CLASS);
+ Object resp = followerTestKit.expectMsgClass(Object.class);
+ if(resp instanceof akka.actor.Status.Failure) {
+ throw new AssertionError("Unexpected failure response", ((akka.actor.Status.Failure)resp).cause());
+ }
+
+ assertTrue("Expected response of type " + CommitTransactionReply.SERIALIZABLE_CLASS,
+ CommitTransactionReply.SERIALIZABLE_CLASS.equals(resp.getClass()));
verifyCars(leaderDistributedDataStore.newReadOnlyTransaction(), car);
}
JavaTestKit.shutdownActorSystem(leaderSystem, null, true);
followerDatastoreContextBuilder.operationTimeoutInMillis(50).shardElectionTimeoutFactor(1);
- followerDistributedDataStore.onDatastoreContextUpdated(followerDatastoreContextBuilder.build());
+ sendDatastoreContextUpdate(followerDistributedDataStore, followerDatastoreContextBuilder);
DOMStoreReadWriteTransaction rwTx = followerDistributedDataStore.newReadWriteTransaction();
Uninterruptibles.sleepUninterruptibly(100, TimeUnit.MILLISECONDS);
followerDatastoreContextBuilder.operationTimeoutInMillis(10).shardElectionTimeoutFactor(1);
- followerDistributedDataStore.onDatastoreContextUpdated(followerDatastoreContextBuilder.build());
+ sendDatastoreContextUpdate(followerDistributedDataStore, followerDatastoreContextBuilder);
DOMStoreReadWriteTransaction rwTx = followerDistributedDataStore.newReadWriteTransaction();
JavaTestKit.shutdownActorSystem(leaderSystem, null, true);
followerDatastoreContextBuilder.operationTimeoutInMillis(500);
- followerDistributedDataStore.onDatastoreContextUpdated(followerDatastoreContextBuilder.build());
+ sendDatastoreContextUpdate(followerDistributedDataStore, followerDatastoreContextBuilder);
DOMStoreReadWriteTransaction rwTx = followerDistributedDataStore.newReadWriteTransaction();
followerTestKit.doCommit(rwTx.ready());
}
+
+ private static void sendDatastoreContextUpdate(DistributedDataStore dataStore, final Builder builder) {
+ DatastoreContextFactory mockContextFactory = Mockito.mock(DatastoreContextFactory.class);
+ Answer<DatastoreContext> answer = new Answer<DatastoreContext>() {
+ @Override
+ public DatastoreContext answer(InvocationOnMock invocation) {
+ return builder.build();
+ }
+ };
+ Mockito.doAnswer(answer).when(mockContextFactory).getBaseDatastoreContext();
+ Mockito.doAnswer(answer).when(mockContextFactory).getShardDatastoreContext(Mockito.anyString());
+ dataStore.onDatastoreContextUpdated(mockContextFactory);
+ }
}