From 9b10928c159317160fb3049727679a48232b6041 Mon Sep 17 00:00:00 2001 From: Sylvain Lebresne Date: Mon, 25 May 2015 18:26:56 +0200 Subject: [PATCH 1/2] Always mark sstable suspected on corruption patch by slebresne; reviewed by benedict for CASSANDRA-9478 --- CHANGES.txt | 1 + .../io/sstable/SSTableIdentityIterator.java | 47 +++++++++++++++---- .../BlacklistingCompactionsTest.java | 16 +++++-- 3 files changed, 49 insertions(+), 15 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index d23661dc00..1aad965e5f 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.0.16: + * Always mark sstable suspect when corrupted (CASSANDRA-9478) * Add database users and permissions to CQL3 documentation (CASSANDRA-7558) * Allow JVM_OPTS to be passed to standalone tools (CASSANDRA-5969) * Fix bad condition in RangeTombstoneList (CASSANDRA-9485) diff --git a/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java b/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java index 52da9bb620..8b45005f3b 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java @@ -50,6 +50,9 @@ public class SSTableIdentityIterator implements Comparable iter = cf.metadata().getOnDiskIterator(in, columnCount, flag, expireBefore, dataVersion); - while (iter.hasNext()) - cf.addAtom(iter.next()); - - if (validateColumns) + try { - try + Iterator iter = cf.metadata().getOnDiskIterator(in, columnCount, flag, expireBefore, dataVersion); + while (iter.hasNext()) + cf.addAtom(iter.next()); + + if (validateColumns) { - cf.metadata().validateColumns(cf); + try + { + cf.metadata().validateColumns(cf); + } + catch (MarshalException e) + { + throw new RuntimeException("Error validating row " + key, e); + } } - catch (MarshalException e) + return cf; + } + catch (IOError e) + { + // catch here b/c atomIterator is an AbstractIterator; hasNext reads the value + if (e.getCause() instanceof IOException) { - throw new RuntimeException("Error validating row " + key, e); + if (sstable != null) + sstable.markSuspect(); + throw new CorruptSSTableException((IOException)e.getCause(), filename); + } + else + { + throw e; } } - return cf; } public int compareTo(SSTableIdentityIterator o) diff --git a/test/unit/org/apache/cassandra/db/compaction/BlacklistingCompactionsTest.java b/test/unit/org/apache/cassandra/db/compaction/BlacklistingCompactionsTest.java index e392a4bd58..08d1d66d3c 100644 --- a/test/unit/org/apache/cassandra/db/compaction/BlacklistingCompactionsTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/BlacklistingCompactionsTest.java @@ -22,9 +22,7 @@ package org.apache.cassandra.db.compaction; import java.io.RandomAccessFile; -import java.util.Collection; -import java.util.HashSet; -import java.util.Set; +import java.util.*; import org.junit.BeforeClass; import org.junit.Test; @@ -41,6 +39,7 @@ import org.apache.cassandra.utils.ByteBufferUtil; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertTrue; public class BlacklistingCompactionsTest extends SchemaLoader { @@ -121,7 +120,14 @@ public class BlacklistingCompactionsTest extends SchemaLoader { raf = new RandomAccessFile(sstable.getFilename(), "rw"); assertNotNull(raf); - raf.write(0xFFFFFF); + assertTrue(raf.length() > 20); + raf.seek(new Random().nextInt((int)(raf.length() - 20))); + // We want to write something large enough that the corruption cannot get undetected + // (even without compression) + byte[] corruption = new byte[20]; + Arrays.fill(corruption, (byte)0xFF); + raf.write(corruption); + } finally { @@ -155,6 +161,6 @@ public class BlacklistingCompactionsTest extends SchemaLoader cfs.truncateBlocking(); - assertEquals(failures, sstablesToCorrupt); + assertEquals(sstablesToCorrupt, failures); } } From 6f93bd1f65888104e33da2f9f01056b6115952e5 Mon Sep 17 00:00:00 2001 From: Benjamin Lerer Date: Tue, 2 Jun 2015 15:34:31 +0200 Subject: [PATCH 2/2] Fix custom payload encoding decoding to match protocol spec patch by blerer; reviewed by slebresne for CASSANDRA-9515 --- CHANGES.txt | 1 + .../CustomPayloadMirroringQueryHandler.java | 18 ++++- .../apache/cassandra/cql3/QueryHandler.java | 20 ++++- .../apache/cassandra/cql3/QueryProcessor.java | 22 +++++- .../apache/cassandra/transport/CBUtil.java | 18 ++--- .../apache/cassandra/transport/Message.java | 13 ++-- .../transport/MessagePayloadTest.java | 73 ++++++++++++------- 7 files changed, 110 insertions(+), 55 deletions(-) diff --git a/CHANGES.txt b/CHANGES.txt index 7f0ef51a51..db94c76468 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.2 + * Fix custom payload coding/decoding to match the spec (CASSANDRA-9515) * ant test-all results incomplete when parsed (CASSANDRA-9463) * Disallow frozen<> types in function arguments and return types for clarity (CASSANDRA-9411) diff --git a/src/java/org/apache/cassandra/cql3/CustomPayloadMirroringQueryHandler.java b/src/java/org/apache/cassandra/cql3/CustomPayloadMirroringQueryHandler.java index 3930e9c0f6..02a6df91d1 100644 --- a/src/java/org/apache/cassandra/cql3/CustomPayloadMirroringQueryHandler.java +++ b/src/java/org/apache/cassandra/cql3/CustomPayloadMirroringQueryHandler.java @@ -17,6 +17,7 @@ */ package org.apache.cassandra.cql3; +import java.nio.ByteBuffer; import java.util.Map; import org.apache.cassandra.cql3.statements.BatchStatement; @@ -34,14 +35,17 @@ public class CustomPayloadMirroringQueryHandler implements QueryHandler { static QueryProcessor queryProcessor = QueryProcessor.instance; - public ResultMessage process(String query, QueryState state, QueryOptions options, Map customPayload) + public ResultMessage process(String query, + QueryState state, + QueryOptions options, + Map customPayload) { ResultMessage result = queryProcessor.process(query, state, options, customPayload); result.setCustomPayload(customPayload); return result; } - public ResultMessage.Prepared prepare(String query, QueryState state, Map customPayload) + public ResultMessage.Prepared prepare(String query, QueryState state, Map customPayload) { ResultMessage.Prepared prepared = queryProcessor.prepare(query, state, customPayload); prepared.setCustomPayload(customPayload); @@ -58,14 +62,20 @@ public class CustomPayloadMirroringQueryHandler implements QueryHandler return queryProcessor.getPreparedForThrift(id); } - public ResultMessage processPrepared(CQLStatement statement, QueryState state, QueryOptions options, Map customPayload) + public ResultMessage processPrepared(CQLStatement statement, + QueryState state, + QueryOptions options, + Map customPayload) { ResultMessage result = queryProcessor.processPrepared(statement, state, options, customPayload); result.setCustomPayload(customPayload); return result; } - public ResultMessage processBatch(BatchStatement statement, QueryState state, BatchQueryOptions options, Map customPayload) + public ResultMessage processBatch(BatchStatement statement, + QueryState state, + BatchQueryOptions options, + Map customPayload) { ResultMessage result = queryProcessor.processBatch(statement, state, options, customPayload); result.setCustomPayload(customPayload); diff --git a/src/java/org/apache/cassandra/cql3/QueryHandler.java b/src/java/org/apache/cassandra/cql3/QueryHandler.java index 8b579d7eaa..3c11c0e74b 100644 --- a/src/java/org/apache/cassandra/cql3/QueryHandler.java +++ b/src/java/org/apache/cassandra/cql3/QueryHandler.java @@ -17,6 +17,7 @@ */ package org.apache.cassandra.cql3; +import java.nio.ByteBuffer; import java.util.Map; import org.apache.cassandra.cql3.statements.BatchStatement; @@ -29,15 +30,26 @@ import org.apache.cassandra.utils.MD5Digest; public interface QueryHandler { - ResultMessage process(String query, QueryState state, QueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException; + ResultMessage process(String query, + QueryState state, + QueryOptions options, + Map customPayload) throws RequestExecutionException, RequestValidationException; - ResultMessage.Prepared prepare(String query, QueryState state, Map customPayload) throws RequestValidationException; + ResultMessage.Prepared prepare(String query, + QueryState state, + Map customPayload) throws RequestValidationException; ParsedStatement.Prepared getPrepared(MD5Digest id); ParsedStatement.Prepared getPreparedForThrift(Integer id); - ResultMessage processPrepared(CQLStatement statement, QueryState state, QueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException; + ResultMessage processPrepared(CQLStatement statement, + QueryState state, + QueryOptions options, + Map customPayload) throws RequestExecutionException, RequestValidationException; - ResultMessage processBatch(BatchStatement statement, QueryState state, BatchQueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException; + ResultMessage processBatch(BatchStatement statement, + QueryState state, + BatchQueryOptions options, + Map customPayload) throws RequestExecutionException, RequestValidationException; } diff --git a/src/java/org/apache/cassandra/cql3/QueryProcessor.java b/src/java/org/apache/cassandra/cql3/QueryProcessor.java index 7b9261c2f5..3170932df0 100644 --- a/src/java/org/apache/cassandra/cql3/QueryProcessor.java +++ b/src/java/org/apache/cassandra/cql3/QueryProcessor.java @@ -232,7 +232,11 @@ public class QueryProcessor implements QueryHandler return instance.process(queryString, queryState, QueryOptions.forInternalCalls(cl, Collections.emptyList())); } - public ResultMessage process(String query, QueryState state, QueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException + public ResultMessage process(String query, + QueryState state, + QueryOptions options, + Map customPayload) + throws RequestExecutionException, RequestValidationException { return process(query, state, options); } @@ -342,7 +346,9 @@ public class QueryProcessor implements QueryHandler return UntypedResultSet.create(cqlRows); } - public ResultMessage.Prepared prepare(String query, QueryState state, Map customPayload) throws RequestValidationException + public ResultMessage.Prepared prepare(String query, + QueryState state, + Map customPayload) throws RequestValidationException { return prepare(query, state); } @@ -422,7 +428,11 @@ public class QueryProcessor implements QueryHandler } } - public ResultMessage processPrepared(CQLStatement statement, QueryState state, QueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException + public ResultMessage processPrepared(CQLStatement statement, + QueryState state, + QueryOptions options, + Map customPayload) + throws RequestExecutionException, RequestValidationException { return processPrepared(statement, state, options); } @@ -450,7 +460,11 @@ public class QueryProcessor implements QueryHandler return processStatement(statement, queryState, options); } - public ResultMessage processBatch(BatchStatement statement, QueryState state, BatchQueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException + public ResultMessage processBatch(BatchStatement statement, + QueryState state, + BatchQueryOptions options, + Map customPayload) + throws RequestExecutionException, RequestValidationException { return processBatch(statement, state, options); } diff --git a/src/java/org/apache/cassandra/transport/CBUtil.java b/src/java/org/apache/cassandra/transport/CBUtil.java index 48beea020e..92e289145a 100644 --- a/src/java/org/apache/cassandra/transport/CBUtil.java +++ b/src/java/org/apache/cassandra/transport/CBUtil.java @@ -183,36 +183,36 @@ public abstract class CBUtil return 2 + bytes.length; } - public static Map readBytesMap(ByteBuf cb) + public static Map readBytesMap(ByteBuf cb) { int length = cb.readUnsignedShort(); - Map m = new HashMap<>(length); + Map m = new HashMap<>(length); for (int i = 0; i < length; i++) { String k = readString(cb); - byte[] v = readBytes(cb); + ByteBuffer v = readValue(cb); m.put(k, v); } return m; } - public static void writeBytesMap(Map m, ByteBuf cb) + public static void writeBytesMap(Map m, ByteBuf cb) { cb.writeShort(m.size()); - for (Map.Entry entry : m.entrySet()) + for (Map.Entry entry : m.entrySet()) { writeString(entry.getKey(), cb); - writeBytes(entry.getValue(), cb); + writeValue(entry.getValue(), cb); } } - public static int sizeOfBytesMap(Map m) + public static int sizeOfBytesMap(Map m) { int size = 2; - for (Map.Entry entry : m.entrySet()) + for (Map.Entry entry : m.entrySet()) { size += sizeOfString(entry.getKey()); - size += sizeOfBytes(entry.getValue()); + size += sizeOfValue(entry.getValue()); } return size; } diff --git a/src/java/org/apache/cassandra/transport/Message.java b/src/java/org/apache/cassandra/transport/Message.java index b6d5a956f4..440d48165f 100644 --- a/src/java/org/apache/cassandra/transport/Message.java +++ b/src/java/org/apache/cassandra/transport/Message.java @@ -19,6 +19,7 @@ package org.apache.cassandra.transport; import java.util.ArrayList; import java.io.IOException; +import java.nio.ByteBuffer; import java.util.EnumSet; import java.util.HashSet; import java.util.List; @@ -148,7 +149,7 @@ public abstract class Message protected Connection connection; private int streamId; private Frame sourceFrame; - private Map customPayload; + private Map customPayload; protected Message(Type type) { @@ -186,12 +187,12 @@ public abstract class Message return sourceFrame; } - public Map getCustomPayload() + public Map getCustomPayload() { return customPayload; } - public void setCustomPayload(Map customPayload) + public void setCustomPayload(Map customPayload) { this.customPayload = customPayload; } @@ -269,7 +270,7 @@ public abstract class Message UUID tracingId = isRequest || !isTracing ? null : CBUtil.readUUID(frame.body); List warnings = isRequest || !hasWarning ? null : CBUtil.readStringList(frame.body); - Map customPayload = !isCustomPayload ? null : CBUtil.readBytesMap(frame.body); + Map customPayload = !isCustomPayload ? null : CBUtil.readBytesMap(frame.body); try { @@ -329,7 +330,7 @@ public abstract class Message if (message instanceof Response) { UUID tracingId = ((Response)message).getTracingId(); - Map customPayload = message.getCustomPayload(); + Map customPayload = message.getCustomPayload(); if (tracingId != null) messageSize += CBUtil.sizeOfUUID(tracingId); List warnings = ((Response)message).getWarnings(); @@ -367,7 +368,7 @@ public abstract class Message assert message instanceof Request; if (((Request)message).isTracingRequested()) flags.add(Frame.Header.Flag.TRACING); - Map payload = message.getCustomPayload(); + Map payload = message.getCustomPayload(); if (payload != null) messageSize += CBUtil.sizeOfBytesMap(payload); body = CBUtil.allocator.buffer(messageSize); diff --git a/test/unit/org/apache/cassandra/transport/MessagePayloadTest.java b/test/unit/org/apache/cassandra/transport/MessagePayloadTest.java index 1049d63a95..73daa48e7c 100644 --- a/test/unit/org/apache/cassandra/transport/MessagePayloadTest.java +++ b/test/unit/org/apache/cassandra/transport/MessagePayloadTest.java @@ -48,10 +48,12 @@ import org.apache.cassandra.transport.messages.QueryMessage; import org.apache.cassandra.transport.messages.ResultMessage; import org.apache.cassandra.utils.MD5Digest; +import static org.apache.cassandra.utils.ByteBufferUtil.bytes; + public class MessagePayloadTest extends CQLTester { - public static Map requestPayload; - public static Map responsePayload; + public static Map requestPayload; + public static Map responsePayload; private static Field cqlQueryHandlerField; private static boolean modifiersAccessible; @@ -125,8 +127,8 @@ public class MessagePayloadTest extends CQLTester { client.connect(false); - Map reqMap; - Map respMap; + Map reqMap; + Map respMap; QueryMessage queryMessage = new QueryMessage( "CREATE TABLE " + KEYSPACE + ".atable (pk int PRIMARY KEY, v text)", @@ -134,23 +136,23 @@ public class MessagePayloadTest extends CQLTester ); PrepareMessage prepareMessage = new PrepareMessage("SELECT * FROM " + KEYSPACE + ".atable"); - reqMap = Collections.singletonMap("foo", "42".getBytes()); - responsePayload = respMap = Collections.singletonMap("bar", "42".getBytes()); + reqMap = Collections.singletonMap("foo", bytes(42)); + responsePayload = respMap = Collections.singletonMap("bar", bytes(42)); queryMessage.setCustomPayload(reqMap); Message.Response queryResponse = client.execute(queryMessage); payloadEquals(reqMap, requestPayload); payloadEquals(respMap, queryResponse.getCustomPayload()); - reqMap = Collections.singletonMap("foo", "43".getBytes()); - responsePayload = respMap = Collections.singletonMap("bar", "43".getBytes()); + reqMap = Collections.singletonMap("foo", bytes(43)); + responsePayload = respMap = Collections.singletonMap("bar", bytes(43)); prepareMessage.setCustomPayload(reqMap); ResultMessage.Prepared prepareResponse = (ResultMessage.Prepared) client.execute(prepareMessage); payloadEquals(reqMap, requestPayload); payloadEquals(respMap, prepareResponse.getCustomPayload()); ExecuteMessage executeMessage = new ExecuteMessage(prepareResponse.statementId, QueryOptions.DEFAULT); - reqMap = Collections.singletonMap("foo", "44".getBytes()); - responsePayload = respMap = Collections.singletonMap("bar", "44".getBytes()); + reqMap = Collections.singletonMap("foo", bytes(44)); + responsePayload = respMap = Collections.singletonMap("bar", bytes(44)); executeMessage.setCustomPayload(reqMap); Message.Response executeResponse = client.execute(executeMessage); payloadEquals(reqMap, requestPayload); @@ -160,8 +162,8 @@ public class MessagePayloadTest extends CQLTester Collections.singletonList("INSERT INTO " + KEYSPACE + ".atable (pk,v) VALUES (1, 'foo')"), Collections.singletonList(Collections.emptyList()), QueryOptions.DEFAULT); - reqMap = Collections.singletonMap("foo", "45".getBytes()); - responsePayload = respMap = Collections.singletonMap("bar", "45".getBytes()); + reqMap = Collections.singletonMap("foo", bytes(45)); + responsePayload = respMap = Collections.singletonMap("bar", bytes(45)); batchMessage.setCustomPayload(reqMap); Message.Response batchResponse = client.execute(batchMessage); payloadEquals(reqMap, requestPayload); @@ -194,7 +196,7 @@ public class MessagePayloadTest extends CQLTester { client.connect(false); - Map reqMap; + Map reqMap; QueryMessage queryMessage = new QueryMessage( "CREATE TABLE " + KEYSPACE + ".atable (pk int PRIMARY KEY, v text)", @@ -202,8 +204,8 @@ public class MessagePayloadTest extends CQLTester ); PrepareMessage prepareMessage = new PrepareMessage("SELECT * FROM " + KEYSPACE + ".atable"); - reqMap = Collections.singletonMap("foo", "42".getBytes()); - responsePayload = Collections.singletonMap("bar", "42".getBytes()); + reqMap = Collections.singletonMap("foo", bytes(42)); + responsePayload = Collections.singletonMap("bar", bytes(42)); queryMessage.setCustomPayload(reqMap); try { @@ -217,8 +219,8 @@ public class MessagePayloadTest extends CQLTester queryMessage.setCustomPayload(null); client.execute(queryMessage); - reqMap = Collections.singletonMap("foo", "43".getBytes()); - responsePayload = Collections.singletonMap("bar", "43".getBytes()); + reqMap = Collections.singletonMap("foo", bytes(43)); + responsePayload = Collections.singletonMap("bar", bytes(43)); prepareMessage.setCustomPayload(reqMap); try { @@ -233,8 +235,8 @@ public class MessagePayloadTest extends CQLTester ResultMessage.Prepared prepareResponse = (ResultMessage.Prepared) client.execute(prepareMessage); ExecuteMessage executeMessage = new ExecuteMessage(prepareResponse.statementId, QueryOptions.DEFAULT); - reqMap = Collections.singletonMap("foo", "44".getBytes()); - responsePayload = Collections.singletonMap("bar", "44".getBytes()); + reqMap = Collections.singletonMap("foo", bytes(44)); + responsePayload = Collections.singletonMap("bar", bytes(44)); executeMessage.setCustomPayload(reqMap); try { @@ -250,8 +252,8 @@ public class MessagePayloadTest extends CQLTester Collections.singletonList("INSERT INTO " + KEYSPACE + ".atable (pk,v) VALUES (1, 'foo')"), Collections.singletonList(Collections.emptyList()), QueryOptions.DEFAULT); - reqMap = Collections.singletonMap("foo", "45".getBytes()); - responsePayload = Collections.singletonMap("bar", "45".getBytes()); + reqMap = Collections.singletonMap("foo", bytes(45)); + responsePayload = Collections.singletonMap("bar", bytes(45)); batchMessage.setCustomPayload(reqMap); try { @@ -274,13 +276,13 @@ public class MessagePayloadTest extends CQLTester } } - private static void payloadEquals(Map map1, Map map2) + private static void payloadEquals(Map map1, Map map2) { Assert.assertNotNull(map1); Assert.assertNotNull(map2); Assert.assertEquals(map1.keySet(), map2.keySet()); - for (Map.Entry e : map1.entrySet()) - Assert.assertArrayEquals(e.getValue(), map2.get(e.getKey())); + for (Map.Entry e : map1.entrySet()) + Assert.assertEquals(e.getValue(), map2.get(e.getKey())); } public static class TestQueryHandler implements QueryHandler @@ -295,7 +297,10 @@ public class MessagePayloadTest extends CQLTester return QueryProcessor.instance.getPreparedForThrift(id); } - public ResultMessage.Prepared prepare(String query, QueryState state, Map customPayload) throws RequestValidationException + public ResultMessage.Prepared prepare(String query, + QueryState state, + Map customPayload) + throws RequestValidationException { if (customPayload != null) requestPayload = customPayload; @@ -308,7 +313,11 @@ public class MessagePayloadTest extends CQLTester return result; } - public ResultMessage process(String query, QueryState state, QueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException + public ResultMessage process(String query, + QueryState state, + QueryOptions options, + Map customPayload) + throws RequestExecutionException, RequestValidationException { if (customPayload != null) requestPayload = customPayload; @@ -321,7 +330,11 @@ public class MessagePayloadTest extends CQLTester return result; } - public ResultMessage processBatch(BatchStatement statement, QueryState state, BatchQueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException + public ResultMessage processBatch(BatchStatement statement, + QueryState state, + BatchQueryOptions options, + Map customPayload) + throws RequestExecutionException, RequestValidationException { if (customPayload != null) requestPayload = customPayload; @@ -334,7 +347,11 @@ public class MessagePayloadTest extends CQLTester return result; } - public ResultMessage processPrepared(CQLStatement statement, QueryState state, QueryOptions options, Map customPayload) throws RequestExecutionException, RequestValidationException + public ResultMessage processPrepared(CQLStatement statement, + QueryState state, + QueryOptions options, + Map customPayload) + throws RequestExecutionException, RequestValidationException { if (customPayload != null) requestPayload = customPayload;