mirror of https://github.com/apache/cassandra
Merge branch 'cassandra-4.1' into trunk
This commit is contained in:
commit
73c0f7f2af
|
|
@ -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)
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
Loading…
Reference in New Issue