Repository: cassandra-dtest Updated Branches: refs/heads/master 6e80b1846 -> ae79495a0
add tests for CASSANDRA-10726 Project: http://git-wip-us.apache.org/repos/asf/cassandra-dtest/repo Commit: http://git-wip-us.apache.org/repos/asf/cassandra-dtest/commit/ae79495a Tree: http://git-wip-us.apache.org/repos/asf/cassandra-dtest/tree/ae79495a Diff: http://git-wip-us.apache.org/repos/asf/cassandra-dtest/diff/ae79495a Branch: refs/heads/master Commit: ae79495a0a0cb0249e6b7375f68e5d2691c3e47d Parents: 6e80b18 Author: Blake Eggleston <[email protected]> Authored: Fri Mar 30 12:29:27 2018 -0700 Committer: Blake Eggleston <[email protected]> Committed: Fri Aug 24 09:49:34 2018 -0700 ---------------------------------------------------------------------- byteman/read_repair/sorted_live_endpoints.btm | 15 + byteman/read_repair/stop_data_reads.btm | 10 + byteman/read_repair/stop_digest_reads.btm | 10 + byteman/read_repair/stop_rr_writes.btm | 8 + byteman/read_repair/stop_writes.btm | 8 + read_repair_test.py | 305 ++++++++++++++++++++- 6 files changed, 355 insertions(+), 1 deletion(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/cassandra-dtest/blob/ae79495a/byteman/read_repair/sorted_live_endpoints.btm ---------------------------------------------------------------------- diff --git a/byteman/read_repair/sorted_live_endpoints.btm b/byteman/read_repair/sorted_live_endpoints.btm new file mode 100644 index 0000000..221e958 --- /dev/null +++ b/byteman/read_repair/sorted_live_endpoints.btm @@ -0,0 +1,15 @@ +RULE sorted live endpoints +CLASS org.apache.cassandra.service.StorageProxy +METHOD getLiveSortedEndpoints +AT ENTRY +BIND ep1 = org.apache.cassandra.locator.InetAddressAndPort.getByName("127.0.0.1"); + ep2 = org.apache.cassandra.locator.InetAddressAndPort.getByName("127.0.0.2"); + ep3 = org.apache.cassandra.locator.InetAddressAndPort.getByName("127.0.0.3"); + eps = new java.util.ArrayList(); +IF true +DO + eps.add(ep1); + eps.add(ep2); + eps.add(ep3); + return eps; +ENDRULE \ No newline at end of file http://git-wip-us.apache.org/repos/asf/cassandra-dtest/blob/ae79495a/byteman/read_repair/stop_data_reads.btm ---------------------------------------------------------------------- diff --git a/byteman/read_repair/stop_data_reads.btm b/byteman/read_repair/stop_data_reads.btm new file mode 100644 index 0000000..9506aba --- /dev/null +++ b/byteman/read_repair/stop_data_reads.btm @@ -0,0 +1,10 @@ +# block data (but not digest) reads +RULE disable data reads +CLASS org.apache.cassandra.db.ReadCommandVerbHandler +METHOD doVerb +# wait until command is declared locally. because generics +AFTER WRITE $command +# bail out if it's not a digest request +IF NOT $command.isDigestQuery() +DO return; +ENDRULE http://git-wip-us.apache.org/repos/asf/cassandra-dtest/blob/ae79495a/byteman/read_repair/stop_digest_reads.btm ---------------------------------------------------------------------- diff --git a/byteman/read_repair/stop_digest_reads.btm b/byteman/read_repair/stop_digest_reads.btm new file mode 100644 index 0000000..141d3ff --- /dev/null +++ b/byteman/read_repair/stop_digest_reads.btm @@ -0,0 +1,10 @@ +# block data (but not digest) reads +RULE disable data reads +CLASS org.apache.cassandra.db.ReadCommandVerbHandler +METHOD doVerb +# wait until command is declared locally. because generics +AFTER WRITE $command +# bail out if it's not a digest request +IF $command.isDigestQuery() +DO return; +ENDRULE http://git-wip-us.apache.org/repos/asf/cassandra-dtest/blob/ae79495a/byteman/read_repair/stop_rr_writes.btm ---------------------------------------------------------------------- diff --git a/byteman/read_repair/stop_rr_writes.btm b/byteman/read_repair/stop_rr_writes.btm new file mode 100644 index 0000000..bb8aab7 --- /dev/null +++ b/byteman/read_repair/stop_rr_writes.btm @@ -0,0 +1,8 @@ +# block remote read repair mutation messages +RULE disable read repair mutations +CLASS org.apache.cassandra.db.ReadRepairVerbHandler +METHOD doVerb +AT ENTRY +IF true +DO return; +ENDRULE http://git-wip-us.apache.org/repos/asf/cassandra-dtest/blob/ae79495a/byteman/read_repair/stop_writes.btm ---------------------------------------------------------------------- diff --git a/byteman/read_repair/stop_writes.btm b/byteman/read_repair/stop_writes.btm new file mode 100644 index 0000000..39d1765 --- /dev/null +++ b/byteman/read_repair/stop_writes.btm @@ -0,0 +1,8 @@ +# block remote read repair writes +RULE disable mutations +CLASS org.apache.cassandra.db.MutationVerbHandler +METHOD doVerb +AT ENTRY +IF true +DO return; +ENDRULE \ No newline at end of file http://git-wip-us.apache.org/repos/asf/cassandra-dtest/blob/ae79495a/read_repair_test.py ---------------------------------------------------------------------- diff --git a/read_repair_test.py b/read_repair_test.py index 175e19e..a7dbf14 100644 --- a/read_repair_test.py +++ b/read_repair_test.py @@ -3,12 +3,15 @@ import time import pytest import logging -from cassandra import ConsistencyLevel +from cassandra import ConsistencyLevel, WriteTimeout, ReadTimeout from cassandra.query import SimpleStatement +from ccmlib.node import Node +from pytest import raises from dtest import Tester, create_ks from tools.assertions import assert_one from tools.data import rows_to_list +from tools.jmxutils import JolokiaAgent, make_mbean from tools.misc import retry_till_success since = pytest.mark.since @@ -314,6 +317,306 @@ class TestReadRepair(Tester): print(("-" * 40)) +def quorum(query_string): + return SimpleStatement(query_string=query_string, consistency_level=ConsistencyLevel.QUORUM) + + +kcv = lambda k, c, v: [k, c, v] + + +listify = lambda results: [list(r) for r in results] + + +class StorageProxy(object): + + def __init__(self, node): + assert isinstance(node, Node) + self.node = node + self.jmx = JolokiaAgent(node) + + def start(self): + self.jmx.start() + + def stop(self): + self.jmx.stop() + + def _get_metric(self, metric): + mbean = make_mbean("metrics", type="ReadRepair", name=metric) + return self.jmx.read_attribute(mbean, "Count") + + @property + def blocking_read_repair(self): + return self._get_metric("RepairedBlocking") + + @property + def speculated_rr_read(self): + return self._get_metric("SpeculatedRead") + + @property + def speculated_rr_write(self): + return self._get_metric("SpeculatedWrite") + + def get_table_metric(self, keyspace, table, metric, attr="Count"): + mbean = make_mbean("metrics", keyspace=keyspace, scope=table, type="Table", name=metric) + return self.jmx.read_attribute(mbean, attr) + + def __enter__(self): + """ For contextmanager-style usage. """ + self.start() + return self + + def __exit__(self, exc_type, value, traceback): + """ For contextmanager-style usage. """ + self.stop() + + +class TestSpeculativeReadRepair(Tester): + + @pytest.fixture(scope='function', autouse=True) + def fixture_set_cluster_settings(self, fixture_dtest_setup): + cluster = fixture_dtest_setup.cluster + cluster.set_configuration_options(values={'hinted_handoff_enabled': False, + 'dynamic_snitch': False, + 'write_request_timeout_in_ms': 500, + 'read_request_timeout_in_ms': 500}) + cluster.populate(3, install_byteman=True, debug=True).start(wait_for_binary_proto=True, + jvm_args=['-XX:-PerfDisableSharedMem']) + session = fixture_dtest_setup.patient_exclusive_cql_connection(cluster.nodelist()[0], timeout=2) + + session.execute("CREATE KEYSPACE ks WITH replication = {'class': 'SimpleStrategy', 'replication_factor': 3}") + session.execute("CREATE TABLE ks.tbl (k int, c int, v int, primary key (k, c)) WITH speculative_retry = '250ms';") + + def get_cql_connection(self, node, **kwargs): + return self.patient_exclusive_cql_connection(node, retry_policy=None, **kwargs) + + @since('4.0') + def test_failed_read_repair(self): + """ + If none of the disagreeing nodes ack the repair mutation, the read should fail + """ + node1, node2, node3 = self.cluster.nodelist() + assert isinstance(node1, Node) + assert isinstance(node2, Node) + assert isinstance(node3, Node) + + session = self.get_cql_connection(node1, timeout=2) + session.execute(quorum("INSERT INTO ks.tbl (k, c, v) VALUES (1, 0, 1)")) + + node2.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + node2.byteman_submit(['./byteman/read_repair/stop_rr_writes.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_rr_writes.btm']) + + with raises(WriteTimeout): + session.execute(quorum("INSERT INTO ks.tbl (k, c, v) VALUES (1, 1, 2)")) + + node2.byteman_submit(['./byteman/read_repair/sorted_live_endpoints.btm']) + session = self.get_cql_connection(node2) + with StorageProxy(node2) as storage_proxy: + assert storage_proxy.blocking_read_repair == 0 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write == 0 + + with raises(ReadTimeout): + session.execute(quorum("SELECT * FROM ks.tbl WHERE k=1")) + + assert storage_proxy.blocking_read_repair > 0 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write > 0 + + @since('4.0') + def test_normal_read_repair(self): + """ test the normal case """ + node1, node2, node3 = self.cluster.nodelist() + assert isinstance(node1, Node) + assert isinstance(node2, Node) + assert isinstance(node3, Node) + session = self.get_cql_connection(node1, timeout=2) + + session.execute(quorum("INSERT INTO ks.tbl (k, c, v) VALUES (1, 0, 1)")) + + node2.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + + session.execute("INSERT INTO ks.tbl (k, c, v) VALUES (1, 1, 2)") + + # re-enable writes + node2.byteman_submit(['-u', './byteman/read_repair/stop_writes.btm']) + + node2.byteman_submit(['./byteman/read_repair/sorted_live_endpoints.btm']) + with StorageProxy(node2) as storage_proxy: + assert storage_proxy.blocking_read_repair == 0 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write == 0 + + session = self.get_cql_connection(node2) + expected = [kcv(1, 0, 1), kcv(1, 1, 2)] + results = session.execute(quorum("SELECT * FROM ks.tbl WHERE k=1")) + assert listify(results) == expected + + assert storage_proxy.blocking_read_repair == 1 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write == 0 + + @since('4.0') + def test_speculative_data_request(self): + """ If one node doesn't respond to a full data request, it should query the other """ + node1, node2, node3 = self.cluster.nodelist() + assert isinstance(node1, Node) + assert isinstance(node2, Node) + assert isinstance(node3, Node) + session = self.get_cql_connection(node1, timeout=2) + + session.execute(quorum("INSERT INTO ks.tbl (k, c, v) VALUES (1, 0, 1)")) + + node2.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + + session.execute("INSERT INTO ks.tbl (k, c, v) VALUES (1, 1, 2)") + + # re-enable writes + node2.byteman_submit(['-u', './byteman/read_repair/stop_writes.btm']) + + node1.byteman_submit(['./byteman/read_repair/sorted_live_endpoints.btm']) + with StorageProxy(node1) as storage_proxy: + assert storage_proxy.blocking_read_repair == 0 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write == 0 + + session = self.get_cql_connection(node1) + node2.byteman_submit(['./byteman/read_repair/stop_data_reads.btm']) + results = session.execute(quorum("SELECT * FROM ks.tbl WHERE k=1")) + assert listify(results) == [kcv(1, 0, 1), kcv(1, 1, 2)] + + assert storage_proxy.blocking_read_repair == 1 + assert storage_proxy.speculated_rr_read == 1 + assert storage_proxy.speculated_rr_write == 0 + + @since('4.0') + def test_speculative_write(self): + """ if one node doesn't respond to a read repair mutation, it should be sent to the remaining node """ + node1, node2, node3 = self.cluster.nodelist() + assert isinstance(node1, Node) + assert isinstance(node2, Node) + assert isinstance(node3, Node) + session = self.get_cql_connection(node1, timeout=2) + + session.execute(quorum("INSERT INTO ks.tbl (k, c, v) VALUES (1, 0, 1)")) + + node2.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + + session.execute("INSERT INTO ks.tbl (k, c, v) VALUES (1, 1, 2)") + + # re-enable writes on node 3, leave them off on node2 + node2.byteman_submit(['./byteman/read_repair/stop_rr_writes.btm']) + + node1.byteman_submit(['./byteman/read_repair/sorted_live_endpoints.btm']) + with StorageProxy(node1) as storage_proxy: + assert storage_proxy.blocking_read_repair == 0 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write == 0 + + session = self.get_cql_connection(node1) + expected = [kcv(1, 0, 1), kcv(1, 1, 2)] + results = session.execute(quorum("SELECT * FROM ks.tbl WHERE k=1")) + assert listify(results) == expected + + assert storage_proxy.blocking_read_repair == 1 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write == 1 + + @since('4.0') + def test_quorum_requirement(self): + """ + Even if we speculate on every stage, we should still only require a quorum of responses for success + """ + node1, node2, node3 = self.cluster.nodelist() + assert isinstance(node1, Node) + assert isinstance(node2, Node) + assert isinstance(node3, Node) + session = self.get_cql_connection(node1, timeout=2) + + session.execute(quorum("INSERT INTO ks.tbl (k, c, v) VALUES (1, 0, 1)")) + + node2.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + + session.execute("INSERT INTO ks.tbl (k, c, v) VALUES (1, 1, 2)") + + # re-enable writes + node2.byteman_submit(['-u', './byteman/read_repair/stop_writes.btm']) + node3.byteman_submit(['-u', './byteman/read_repair/stop_writes.btm']) + + # force endpoint order + node1.byteman_submit(['./byteman/read_repair/sorted_live_endpoints.btm']) + + # node2.byteman_submit(['./byteman/read_repair/stop_digest_reads.btm']) + node2.byteman_submit(['./byteman/read_repair/stop_data_reads.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_rr_writes.btm']) + + with StorageProxy(node1) as storage_proxy: + assert storage_proxy.get_table_metric("ks", "tbl", "SpeculativeRetries") == 0 + assert storage_proxy.blocking_read_repair == 0 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write == 0 + + session = self.get_cql_connection(node1) + expected = [kcv(1, 0, 1), kcv(1, 1, 2)] + results = session.execute(quorum("SELECT * FROM ks.tbl WHERE k=1")) + assert listify(results) == expected + + assert storage_proxy.get_table_metric("ks", "tbl", "SpeculativeRetries") == 0 + assert storage_proxy.blocking_read_repair == 1 + assert storage_proxy.speculated_rr_read == 1 + assert storage_proxy.speculated_rr_write == 1 + + @since('4.0') + def test_quorum_requirement_on_speculated_read(self): + """ + Even if we speculate on every stage, we should still only require a quorum of responses for success + """ + node1, node2, node3 = self.cluster.nodelist() + assert isinstance(node1, Node) + assert isinstance(node2, Node) + assert isinstance(node3, Node) + session = self.get_cql_connection(node1, timeout=2) + + session.execute(quorum("INSERT INTO ks.tbl (k, c, v) VALUES (1, 0, 1)")) + + node2.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_writes.btm']) + + session.execute("INSERT INTO ks.tbl (k, c, v) VALUES (1, 1, 2)") + + # re-enable writes + node2.byteman_submit(['-u', './byteman/read_repair/stop_writes.btm']) + node3.byteman_submit(['-u', './byteman/read_repair/stop_writes.btm']) + + # force endpoint order + node1.byteman_submit(['./byteman/read_repair/sorted_live_endpoints.btm']) + + node2.byteman_submit(['./byteman/read_repair/stop_digest_reads.btm']) + node3.byteman_submit(['./byteman/read_repair/stop_data_reads.btm']) + node2.byteman_submit(['./byteman/read_repair/stop_rr_writes.btm']) + + with StorageProxy(node1) as storage_proxy: + assert storage_proxy.get_table_metric("ks", "tbl", "SpeculativeRetries") == 0 + assert storage_proxy.blocking_read_repair == 0 + assert storage_proxy.speculated_rr_read == 0 + assert storage_proxy.speculated_rr_write == 0 + + session = self.get_cql_connection(node1) + expected = [kcv(1, 0, 1), kcv(1, 1, 2)] + results = session.execute(quorum("SELECT * FROM ks.tbl WHERE k=1")) + assert listify(results) == expected + + assert storage_proxy.get_table_metric("ks", "tbl", "SpeculativeRetries") == 1 + assert storage_proxy.blocking_read_repair == 1 + assert storage_proxy.speculated_rr_read == 0 # there shouldn't be any replicas to speculate on + assert storage_proxy.speculated_rr_write == 1 + + class NotRepairedException(Exception): """ Thrown to indicate that the data on a replica hasn't been doesn't match what we'd expect if a --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
