import static org.junit.Assert.assertNotNull;
import static org.junit.Assert.fail;
import static org.mockito.Matchers.any;
+import static org.mockito.Matchers.anyBoolean;
import static org.mockito.Matchers.eq;
import static org.mockito.Mockito.atLeastOnce;
import static org.mockito.Mockito.doNothing;
import static org.opendaylight.controller.cluster.datastore.ShardDataTreeMocking.coordinatedCanCommit;
import static org.opendaylight.controller.cluster.datastore.ShardDataTreeMocking.coordinatedCommit;
import static org.opendaylight.controller.cluster.datastore.ShardDataTreeMocking.coordinatedPreCommit;
+import static org.opendaylight.controller.cluster.datastore.ShardDataTreeMocking.immediate3PhaseCommit;
import static org.opendaylight.controller.cluster.datastore.ShardDataTreeMocking.immediateCanCommit;
import static org.opendaylight.controller.cluster.datastore.ShardDataTreeMocking.immediateCommit;
+import static org.opendaylight.controller.cluster.datastore.ShardDataTreeMocking.immediatePayloadReplication;
import static org.opendaylight.controller.cluster.datastore.ShardDataTreeMocking.immediatePreCommit;
import com.google.common.base.Optional;
@Before
public void setUp() {
- doReturn(true).when(mockShard).canSkipPayload();
doReturn(Ticker.systemTicker()).when(mockShard).ticker();
doReturn(Mockito.mock(ShardStats.class)).when(mockShard).getShardMBean();
private void modify(final boolean merge, final boolean expectedCarsPresent, final boolean expectedPeoplePresent)
throws ExecutionException, InterruptedException {
+ immediatePayloadReplication(shardDataTree, mockShard);
assertEquals(fullSchema, shardDataTree.getSchemaContext());
@Test
public void bug4359AddRemoveCarOnce() throws ExecutionException, InterruptedException {
+ immediatePayloadReplication(shardDataTree, mockShard);
+
final List<DataTreeCandidate> candidates = new ArrayList<>();
candidates.add(addCar(shardDataTree));
candidates.add(removeCar(shardDataTree));
@Test
public void bug4359AddRemoveCarTwice() throws ExecutionException, InterruptedException {
+ immediatePayloadReplication(shardDataTree, mockShard);
+
final List<DataTreeCandidate> candidates = new ArrayList<>();
candidates.add(addCar(shardDataTree));
candidates.add(removeCar(shardDataTree));
@Test
public void testListenerNotifiedOnApplySnapshot() throws Exception {
+ immediatePayloadReplication(shardDataTree, mockShard);
+
DOMDataTreeChangeListener listener = mock(DOMDataTreeChangeListener.class);
- shardDataTree.registerTreeChangeListener(CarsModel.CAR_LIST_PATH.node(CarsModel.CAR_QNAME), listener);
+ shardDataTree.registerTreeChangeListener(CarsModel.CAR_LIST_PATH.node(CarsModel.CAR_QNAME), listener,
+ Optional.absent(), noop -> { });
addCar(shardDataTree, "optima");
});
ShardDataTree newDataTree = new ShardDataTree(mockShard, fullSchema, TreeType.OPERATIONAL);
+ immediatePayloadReplication(newDataTree, mockShard);
addCar(newDataTree, "optima");
addCar(newDataTree, "murano");
}
@Test
- public void testPipelinedTransactions() throws Exception {
- doReturn(false).when(mockShard).canSkipPayload();
-
+ public void testPipelinedTransactionsWithCoordinatedCommits() throws Exception {
final ShardDataTreeCohort cohort1 = newShardDataTreeCohort(snapshot ->
snapshot.write(CarsModel.BASE_PATH, CarsModel.emptyContainer()));
verify(preCommitCallback4).onSuccess(cohort4.getCandidate());
final FutureCallback<UnsignedLong> commitCallback2 = coordinatedCommit(cohort2);
- verify(mockShard, never()).persistPayload(eq(cohort1.getIdentifier()), any(CommitTransactionPayload.class));
+ verify(mockShard, never()).persistPayload(eq(cohort1.getIdentifier()), any(CommitTransactionPayload.class),
+ anyBoolean());
verifyNoMoreInteractions(commitCallback2);
final FutureCallback<UnsignedLong> commitCallback4 = coordinatedCommit(cohort4);
- verify(mockShard, never()).persistPayload(eq(cohort4.getIdentifier()), any(CommitTransactionPayload.class));
+ verify(mockShard, never()).persistPayload(eq(cohort4.getIdentifier()), any(CommitTransactionPayload.class),
+ anyBoolean());
verifyNoMoreInteractions(commitCallback4);
final FutureCallback<UnsignedLong> commitCallback1 = coordinatedCommit(cohort1);
InOrder inOrder = inOrder(mockShard);
- inOrder.verify(mockShard).persistPayload(eq(cohort1.getIdentifier()), any(CommitTransactionPayload.class));
- inOrder.verify(mockShard).persistPayload(eq(cohort2.getIdentifier()), any(CommitTransactionPayload.class));
- inOrder.verifyNoMoreInteractions();
+ inOrder.verify(mockShard).persistPayload(eq(cohort1.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(true));
+ inOrder.verify(mockShard).persistPayload(eq(cohort2.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(false));
verifyNoMoreInteractions(commitCallback1);
verifyNoMoreInteractions(commitCallback2);
final FutureCallback<UnsignedLong> commitCallback3 = coordinatedCommit(cohort3);
inOrder = inOrder(mockShard);
- inOrder.verify(mockShard).persistPayload(eq(cohort3.getIdentifier()), any(CommitTransactionPayload.class));
- inOrder.verify(mockShard).persistPayload(eq(cohort4.getIdentifier()), any(CommitTransactionPayload.class));
- inOrder.verifyNoMoreInteractions();
+ inOrder.verify(mockShard).persistPayload(eq(cohort3.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(true));
+ inOrder.verify(mockShard).persistPayload(eq(cohort4.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(false));
verifyNoMoreInteractions(commitCallback3);
verifyNoMoreInteractions(commitCallback4);
assertEquals("People node", peopleNode, optional.get());
}
+ @Test
+ public void testPipelinedTransactionsWithImmediateCommits() throws Exception {
+ final ShardDataTreeCohort cohort1 = newShardDataTreeCohort(snapshot ->
+ snapshot.write(CarsModel.BASE_PATH, CarsModel.emptyContainer()));
+
+ final ShardDataTreeCohort cohort2 = newShardDataTreeCohort(snapshot ->
+ snapshot.write(CarsModel.CAR_LIST_PATH, CarsModel.newCarMapNode()));
+
+ YangInstanceIdentifier carPath = CarsModel.newCarPath("optima");
+ MapEntryNode carNode = CarsModel.newCarEntry("optima", new BigInteger("100"));
+ final ShardDataTreeCohort cohort3 = newShardDataTreeCohort(snapshot -> snapshot.write(carPath, carNode));
+
+ final FutureCallback<UnsignedLong> commitCallback2 = immediate3PhaseCommit(cohort2);
+ final FutureCallback<UnsignedLong> commitCallback3 = immediate3PhaseCommit(cohort3);
+ final FutureCallback<UnsignedLong> commitCallback1 = immediate3PhaseCommit(cohort1);
+
+ InOrder inOrder = inOrder(mockShard);
+ inOrder.verify(mockShard).persistPayload(eq(cohort1.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(true));
+ inOrder.verify(mockShard).persistPayload(eq(cohort2.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(true));
+ inOrder.verify(mockShard).persistPayload(eq(cohort3.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(false));
+
+ // The payload instance doesn't matter - it just needs to be of type CommitTransactionPayload.
+ CommitTransactionPayload mockPayload = CommitTransactionPayload.create(nextTransactionId(),
+ cohort1.getCandidate());
+ shardDataTree.applyReplicatedPayload(cohort1.getIdentifier(), mockPayload);
+ shardDataTree.applyReplicatedPayload(cohort2.getIdentifier(), mockPayload);
+ shardDataTree.applyReplicatedPayload(cohort3.getIdentifier(), mockPayload);
+
+ inOrder = inOrder(commitCallback1, commitCallback2, commitCallback3);
+ inOrder.verify(commitCallback1).onSuccess(any(UnsignedLong.class));
+ inOrder.verify(commitCallback2).onSuccess(any(UnsignedLong.class));
+ inOrder.verify(commitCallback3).onSuccess(any(UnsignedLong.class));
+
+ final DataTreeSnapshot snapshot =
+ shardDataTree.newReadOnlyTransaction(nextTransactionId()).getSnapshot();
+ Optional<NormalizedNode<?, ?>> optional = snapshot.readNode(carPath);
+ assertEquals("Car node present", true, optional.isPresent());
+ assertEquals("Car node", carNode, optional.get());
+ }
+
+ @Test
+ public void testPipelinedTransactionsWithImmediateReplication() throws Exception {
+ immediatePayloadReplication(shardDataTree, mockShard);
+
+ final ShardDataTreeCohort cohort1 = newShardDataTreeCohort(snapshot ->
+ snapshot.write(CarsModel.BASE_PATH, CarsModel.emptyContainer()));
+
+ final ShardDataTreeCohort cohort2 = newShardDataTreeCohort(snapshot ->
+ snapshot.write(CarsModel.CAR_LIST_PATH, CarsModel.newCarMapNode()));
+
+ YangInstanceIdentifier carPath = CarsModel.newCarPath("optima");
+ MapEntryNode carNode = CarsModel.newCarEntry("optima", new BigInteger("100"));
+ final ShardDataTreeCohort cohort3 = newShardDataTreeCohort(snapshot -> snapshot.write(carPath, carNode));
+
+ final FutureCallback<UnsignedLong> commitCallback1 = immediate3PhaseCommit(cohort1);
+ final FutureCallback<UnsignedLong> commitCallback2 = immediate3PhaseCommit(cohort2);
+ final FutureCallback<UnsignedLong> commitCallback3 = immediate3PhaseCommit(cohort3);
+
+ InOrder inOrder = inOrder(commitCallback1, commitCallback2, commitCallback3);
+ inOrder.verify(commitCallback1).onSuccess(any(UnsignedLong.class));
+ inOrder.verify(commitCallback2).onSuccess(any(UnsignedLong.class));
+ inOrder.verify(commitCallback3).onSuccess(any(UnsignedLong.class));
+
+ final DataTreeSnapshot snapshot =
+ shardDataTree.newReadOnlyTransaction(nextTransactionId()).getSnapshot();
+ Optional<NormalizedNode<?, ?>> optional = snapshot.readNode(CarsModel.BASE_PATH);
+ assertEquals("Car node present", true, optional.isPresent());
+ }
+
@SuppressWarnings("unchecked")
@Test
public void testAbortWithPendingCommits() throws Exception {
- doReturn(false).when(mockShard).canSkipPayload();
-
final ShardDataTreeCohort cohort1 = newShardDataTreeCohort(snapshot ->
snapshot.write(CarsModel.BASE_PATH, CarsModel.emptyContainer()));
coordinatedCommit(cohort4);
InOrder inOrder = inOrder(mockShard);
- inOrder.verify(mockShard).persistPayload(eq(cohort1.getIdentifier()), any(CommitTransactionPayload.class));
- inOrder.verify(mockShard).persistPayload(eq(cohort3.getIdentifier()), any(CommitTransactionPayload.class));
- inOrder.verify(mockShard).persistPayload(eq(cohort4.getIdentifier()), any(CommitTransactionPayload.class));
- inOrder.verifyNoMoreInteractions();
+ inOrder.verify(mockShard).persistPayload(eq(cohort1.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(false));
+ inOrder.verify(mockShard).persistPayload(eq(cohort3.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(false));
+ inOrder.verify(mockShard).persistPayload(eq(cohort4.getIdentifier()), any(CommitTransactionPayload.class),
+ eq(false));
// The payload instance doesn't matter - it just needs to be of type CommitTransactionPayload.
CommitTransactionPayload mockPayload = CommitTransactionPayload.create(nextTransactionId(),
@SuppressWarnings("unchecked")
@Test
public void testAbortWithFailedRebase() throws Exception {
+ immediatePayloadReplication(shardDataTree, mockShard);
+
final ShardDataTreeCohort cohort1 = newShardDataTreeCohort(snapshot ->
snapshot.write(CarsModel.BASE_PATH, CarsModel.emptyContainer()));