diff --git a/test/distributed/org/apache/cassandra/distributed/test/CASTest.java b/test/distributed/org/apache/cassandra/distributed/test/CASTest.java index c57b10d0f3..2007f10b2f 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/CASTest.java +++ b/test/distributed/org/apache/cassandra/distributed/test/CASTest.java @@ -31,11 +31,14 @@ import org.junit.Ignore; import org.junit.Test; import org.apache.cassandra.distributed.Cluster; +import org.apache.cassandra.distributed.api.ConsistencyLevel; import org.apache.cassandra.distributed.api.ICoordinator; import org.apache.cassandra.distributed.api.IInstanceConfig; import org.apache.cassandra.distributed.api.IMessageFilters; import org.apache.cassandra.exceptions.CasWriteTimeoutException; +import org.apache.cassandra.utils.FBUtilities; + import static org.apache.cassandra.distributed.api.ConsistencyLevel.ANY; import static org.apache.cassandra.distributed.api.ConsistencyLevel.ONE; import static org.apache.cassandra.distributed.api.ConsistencyLevel.QUORUM; @@ -342,6 +345,30 @@ public class CASTest extends CASCommonTestCases assertFalse("Expected CAS to not be applied, but was applied.", (Boolean)wasApplied); } + + private static Object[][] executeWithRetry(int attempts, ICoordinator coordinator, String query, ConsistencyLevel consistencyLevel, Object... boundValues) + { + while (--attempts > 0) + { + try + { + return coordinator.execute(query, consistencyLevel, boundValues); + } + catch (Throwable t) + { + if (!t.getClass().getName().equals(CasWriteTimeoutException.class.getName())) + throw t; + FBUtilities.sleepQuietly(100); + } + } + return coordinator.execute(query, consistencyLevel, boundValues); + } + + private static Object[][] executeWithRetry(ICoordinator coordinator, String query, ConsistencyLevel consistencyLevel, Object... boundValues) + { + return executeWithRetry(2, coordinator, query, consistencyLevel, boundValues); + } + // failed write (by node that did not yet witness a range movement via gossip) is witnessed later as successful // conflicting with another successful write performed by a node that did witness the range movement // A Promised, Accepted and Committed by {1, 2} @@ -360,15 +387,20 @@ public class CASTest extends CASCommonTestCases // {1} promises and accepts on !{3} => {1, 2}; commits on !{2,3} => {1} drop(FOUR_NODES, 1, to(3), to(3), to(2, 3)); - assertRows(FOUR_NODES.coordinator(1).execute("INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v1) VALUES (?, 1, 1) IF NOT EXISTS", ONE, pk), + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); + assertRows(executeWithRetry(FOUR_NODES.coordinator(1), "INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v1) VALUES (?, 1, 1) IF NOT EXISTS", ONE, pk), row(true)); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); for (int i = 1; i <= 3; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::addToRingNormal).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); + } // {4} reads from !{2} => {3, 4} drop(FOUR_NODES, 4, to(2), to(2), to()); - assertRows(FOUR_NODES.coordinator(4).execute("INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v2) VALUES (?, 1, 2) IF NOT EXISTS", ONE, pk), + assertRows(executeWithRetry(FOUR_NODES.coordinator(4), "INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v2) VALUES (?, 1, 2) IF NOT EXISTS", ONE, pk), row(false, pk, 1, 1, null)); } @@ -387,17 +419,21 @@ public class CASTest extends CASCommonTestCases // make it so {1} is unaware (yet) that {4} is an owner of the token FOUR_NODES.get(1).acceptOnInstance(CASTestBase::removeFromRing, FOUR_NODES.get(4)); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); // {4} promises, accepts and commits on !{2} => {3, 4} int pk = pk(FOUR_NODES, 1, 2); drop(FOUR_NODES, 4, to(2), to(2), to(2)); - assertRows(FOUR_NODES.coordinator(4).execute("INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v1) VALUES (?, 1, 1) IF NOT EXISTS", ONE, pk), + assertRows(executeWithRetry(FOUR_NODES.coordinator(4), "INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v1) VALUES (?, 1, 1) IF NOT EXISTS", ONE, pk), row(true)); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); // {1} promises, accepts and commmits on !{3} => {1, 2} drop(FOUR_NODES, 1, to(3), to(3), to(3)); - assertRows(FOUR_NODES.coordinator(1).execute("INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v2) VALUES (?, 1, 2) IF NOT EXISTS", ONE, pk), + assertRows(executeWithRetry(FOUR_NODES.coordinator(1), "INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v2) VALUES (?, 1, 2) IF NOT EXISTS", ONE, pk), row(false, pk, 1, 1, null)); + + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); } /** @@ -417,16 +453,24 @@ public class CASTest extends CASCommonTestCases // make it so {4} is bootstrapping, and this has propagated to only a quorum of other nodes for (int i = 1 ; i <= 4 ; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::removeFromRing).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); + } for (int i = 2 ; i <= 4 ; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::addToRingBootstrapping).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); + } int pk = pk(FOUR_NODES, 1, 2); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); // {1} promises and accepts on !{3} => {1, 2}; commmits on !{2, 3} => {1} drop(FOUR_NODES, 1, to(3), to(3), to(2, 3)); - assertRows(FOUR_NODES.coordinator(1).execute("INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v) VALUES (?, 1, 1) IF NOT EXISTS", ONE, pk), + assertRows(executeWithRetry(FOUR_NODES.coordinator(1), "INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v) VALUES (?, 1, 1) IF NOT EXISTS", ONE, pk), row(true)); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); // finish topology change for (int i = 1 ; i <= 4 ; ++i) @@ -453,16 +497,24 @@ public class CASTest extends CASCommonTestCases // make it so {4} is bootstrapping, and this has propagated to only a quorum of other nodes for (int i = 1 ; i <= 4 ; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::removeFromRing).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); + } for (int i = 2 ; i <= 4 ; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::addToRingBootstrapping).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); + } int pk = pk(FOUR_NODES, 1, 2); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); // {1} promises and accepts on !{3} => {1, 2}; commits on !{2, 3} => {1} drop(FOUR_NODES, 1, to(3), to(3), to(2, 3)); - assertRows(FOUR_NODES.coordinator(1).execute("INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1) VALUES (?, 1, 1) IF NOT EXISTS", ONE, pk), + assertRows(executeWithRetry(FOUR_NODES.coordinator(1), "INSERT INTO " + KEYSPACE + ".tbl (pk, ck, v1) VALUES (?, 1, 1) IF NOT EXISTS", ONE, pk), row(true)); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); // finish topology change for (int i = 1 ; i <= 4 ; ++i) @@ -498,9 +550,15 @@ public class CASTest extends CASCommonTestCases // make it so {4} is bootstrapping, and this has propagated to only a quorum of other nodes for (int i = 1 ; i <= 4 ; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::removeFromRing).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); + } for (int i = 2 ; i <= 4 ; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::addToRingBootstrapping).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); + } int pk = pk(FOUR_NODES, 1, 2); @@ -519,7 +577,8 @@ public class CASTest extends CASCommonTestCases // {1} promises and accepts on !{3} => {1, 2}; commits on !{2, 3} => {1} drop(FOUR_NODES, 1, to(3), to(3), to(2, 3)); // two options: either we can invalidate the previous operation and succeed, or we can complete the previous operation - Object[][] result = FOUR_NODES.coordinator(1).execute("INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v2) VALUES (?, 1, 2) IF NOT EXISTS", ONE, pk); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); + Object[][] result = executeWithRetry(FOUR_NODES.coordinator(1), "INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v2) VALUES (?, 1, 2) IF NOT EXISTS", ONE, pk); Object[] expectRow; if (result[0].length == 1) { @@ -531,6 +590,7 @@ public class CASTest extends CASCommonTestCases assertRows(result, row(false, pk, 1, 1, null)); expectRow = row(pk, 1, 1, null); } + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); // finish topology change for (int i = 1 ; i <= 4 ; ++i) @@ -564,9 +624,15 @@ public class CASTest extends CASCommonTestCases // make it so {4} is bootstrapping, and this has propagated to only a quorum of other nodes for (int i = 1; i <= 4; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::removeFromRing).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); + } for (int i = 2; i <= 4; ++i) + { FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::addToRingBootstrapping).accept(FOUR_NODES.get(4)); + FOUR_NODES.get(i).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); + } int pk = pk(FOUR_NODES, 1, 2); @@ -586,7 +652,8 @@ public class CASTest extends CASCommonTestCases // {1} promises and accepts on !{3} => {1, 2}; commits on !{2, 3} => {1} drop(FOUR_NODES, 1, to(3), to(3), to(2, 3)); // two options: either we can invalidate the previous operation and succeed, or we can complete the previous operation - Object[][] result = FOUR_NODES.coordinator(1).execute("INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v2) VALUES (?, 1, 2) IF NOT EXISTS", ONE, pk); + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertNotVisibleInRing).accept(FOUR_NODES.get(4)); + Object[][] result = executeWithRetry(FOUR_NODES.coordinator(1), "INSERT INTO " + KEYSPACE + "." + tableName + " (pk, ck, v2) VALUES (?, 1, 2) IF NOT EXISTS", ONE, pk); Object[] expectRow; if (result[0].length == 1) { @@ -598,6 +665,7 @@ public class CASTest extends CASCommonTestCases assertRows(result, row(false, pk, 1, 1, null)); expectRow = row(false, pk, 1, 1, null); } + FOUR_NODES.get(1).acceptsOnInstance(CASTestBase::assertVisibleInRing).accept(FOUR_NODES.get(4)); // finish topology change for (int i = 1; i <= 4; ++i) diff --git a/test/distributed/org/apache/cassandra/distributed/test/CASTestBase.java b/test/distributed/org/apache/cassandra/distributed/test/CASTestBase.java index 59bcb9e385..7ed59da0f1 100644 --- a/test/distributed/org/apache/cassandra/distributed/test/CASTestBase.java +++ b/test/distributed/org/apache/cassandra/distributed/test/CASTestBase.java @@ -22,6 +22,7 @@ import java.util.Collections; import java.util.UUID; import java.util.concurrent.atomic.AtomicInteger; +import org.junit.Assert; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -183,6 +184,12 @@ public abstract class CASTestBase extends TestBaseImpl } } + public static void assertVisibleInRing(IInstance peer) + { + InetAddressAndPort endpoint = InetAddressAndPort.getByAddress(peer.broadcastAddress()); + Assert.assertTrue(Gossiper.instance.isAlive(endpoint)); + } + // reset gossip state so we know of the node being alive only public static void removeFromRing(IInstance peer) { @@ -208,6 +215,12 @@ public abstract class CASTestBase extends TestBaseImpl } } + public static void assertNotVisibleInRing(IInstance peer) + { + InetAddressAndPort endpoint = InetAddressAndPort.getByAddress(peer.broadcastAddress()); + Assert.assertFalse(Gossiper.instance.isAlive(endpoint)); + } + public static void addToRingNormal(IInstance peer) { addToRing(false, peer);