diff --git a/CHANGES.txt b/CHANGES.txt index 41f5cfbe00..42846d7391 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,5 +1,6 @@ 2.2 * Clean up gossiper logic for old versions (CASSANDRA-9370) + * 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) @@ -167,6 +168,7 @@ Merged from 2.0: * Fix streaming not holding ref when stream error (CASSANDRA-9295) * Fix canonical view returning early opened SSTables (CASSANDRA-9396) Merged from 2.0: + * 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/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/io/sstable/SSTableIdentityIterator.java b/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java index 6300749ec3..17f9a8d33e 100644 --- a/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java +++ b/src/java/org/apache/cassandra/io/sstable/SSTableIdentityIterator.java @@ -41,6 +41,9 @@ import org.apache.cassandra.serializers.MarshalException; private final boolean validateColumns; private final String filename; + // Not every SSTableIdentifyIterator is attached to a sstable, so this can be null. + private final SSTableReader sstable; + /** * Used to iterate through the columns of a row. * @param sstable SSTable we are reading ffrom. @@ -80,6 +83,7 @@ import org.apache.cassandra.serializers.MarshalException; this.key = key; this.flag = flag; this.validateColumns = checkData; + this.sstable = sstable; Version dataVersion = sstable == null ? DatabaseDescriptor.getSSTableFormat().info.getLatestVersion() : sstable.descriptor.version; int expireBefore = (int) (System.currentTimeMillis() / 1000); @@ -118,9 +122,15 @@ import org.apache.cassandra.serializers.MarshalException; { // catch here b/c atomIterator is an AbstractIterator; hasNext reads the value if (e.getCause() instanceof IOException) + { + if (sstable != null) + sstable.markSuspect(); throw new CorruptSSTableException((IOException)e.getCause(), filename); + } else + { throw e; + } } } diff --git a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java index ba634b64bd..f4a59a2ca2 100644 --- a/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java +++ b/src/java/org/apache/cassandra/io/sstable/format/SSTableReader.java @@ -940,14 +940,54 @@ public abstract class SSTableReader extends SSTable implements SelfRefCounted 0) @@ -966,9 +1006,7 @@ public abstract class SSTableReader extends SSTable implements SelfRefCounted 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/db/compaction/BlacklistingCompactionsTest.java b/test/unit/org/apache/cassandra/db/compaction/BlacklistingCompactionsTest.java index 5420b1b5e2..2b6a62ae3f 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.apache.cassandra.io.sstable.format.SSTableReader; import org.junit.After; @@ -43,6 +41,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; import static org.apache.cassandra.Util.cellname; public class BlacklistingCompactionsTest @@ -142,7 +141,14 @@ public class BlacklistingCompactionsTest { 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 { @@ -174,8 +180,7 @@ public class BlacklistingCompactionsTest break; } - cfs.truncateBlocking(); - assertEquals(failures, sstablesToCorrupt); + assertEquals(sstablesToCorrupt, failures); } } 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;