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]

Reply via email to