From f5f97d81539a931653e8cbb8ba6d6be8421f6c29 Mon Sep 17 00:00:00 2001 From: Joshua McKenzie Date: Tue, 24 Mar 2015 13:07:58 -0500 Subject: [PATCH] Fix potential data loss in CompressedSequentialWriter Patch by jmckenzie; reviewed by tjake for CASSANDRA-8949 --- CHANGES.txt | 1 + .../compress/CompressedSequentialWriter.java | 3 + .../cassandra/io/compress/LZ4Compressor.java | 6 +- .../CompressedSequentialWriterTest.java | 126 ++++++++++++++++++ 4 files changed, 135 insertions(+), 1 deletion(-) create mode 100644 test/unit/org/apache/cassandra/io/compress/CompressedSequentialWriterTest.java diff --git a/CHANGES.txt b/CHANGES.txt index 2073b00257..f0555dd38c 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 2.0.14: + * Fix potential data loss in CompressedSequentialWriter (CASSANDRA-8949) * (cqlsh) Allow increasing CSV field size limit through cqlshrc config option (CASSANDRA-8934) * Stop logging range tombstones when exceeding the threshold diff --git a/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java b/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java index 909d82276e..97200e0a4c 100644 --- a/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java +++ b/src/java/org/apache/cassandra/io/compress/CompressedSequentialWriter.java @@ -229,6 +229,9 @@ public class CompressedSequentialWriter extends SequentialWriter bufferOffset = current - validBufferBytes; chunkCount = realMark.nextChunkIndex - 1; + // mark as dirty so we don't lose the bytes on subsequent reBuffer calls + isDirty = true; + // truncate data and index file truncate(chunkOffset); metadataWriter.resetAndTruncate(realMark.nextChunkIndex - 1); diff --git a/src/java/org/apache/cassandra/io/compress/LZ4Compressor.java b/src/java/org/apache/cassandra/io/compress/LZ4Compressor.java index 0cf36c11ba..6f56f302d9 100644 --- a/src/java/org/apache/cassandra/io/compress/LZ4Compressor.java +++ b/src/java/org/apache/cassandra/io/compress/LZ4Compressor.java @@ -23,6 +23,8 @@ import java.util.HashSet; import java.util.Map; import java.util.Set; +import com.google.common.annotations.VisibleForTesting; + import net.jpountz.lz4.LZ4Exception; import net.jpountz.lz4.LZ4Factory; @@ -30,7 +32,9 @@ public class LZ4Compressor implements ICompressor { private static final int INTEGER_BYTES = 4; - private static final LZ4Compressor instance = new LZ4Compressor(); + + @VisibleForTesting + public static final LZ4Compressor instance = new LZ4Compressor(); public static LZ4Compressor create(Map args) { diff --git a/test/unit/org/apache/cassandra/io/compress/CompressedSequentialWriterTest.java b/test/unit/org/apache/cassandra/io/compress/CompressedSequentialWriterTest.java new file mode 100644 index 0000000000..acc4faa0eb --- /dev/null +++ b/test/unit/org/apache/cassandra/io/compress/CompressedSequentialWriterTest.java @@ -0,0 +1,126 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.cassandra.io.compress; + +import java.io.File; +import java.io.IOException; +import java.util.Arrays; +import java.util.Collections; +import java.util.Random; + +import static org.junit.Assert.assertEquals; + +import org.apache.cassandra.exceptions.ConfigurationException; +import org.apache.cassandra.io.sstable.SSTableMetadata; +import org.junit.Test; + +import org.apache.cassandra.db.marshal.BytesType; +import org.apache.cassandra.io.util.FileMark; + +public class CompressedSequentialWriterTest +{ + private ICompressor compressor; + + private void runTests(String testName) throws IOException, ConfigurationException + { + // Test small < 1 chunk data set + testWrite(File.createTempFile(testName + "_small", "1"), 25); + + // Test to confirm pipeline w/chunk-aligned data writes works + testWrite(File.createTempFile(testName + "_chunkAligned", "1"), CompressionParameters.DEFAULT_CHUNK_LENGTH); + + // Test to confirm pipeline on non-chunk boundaries works + testWrite(File.createTempFile(testName + "_large", "1"), CompressionParameters.DEFAULT_CHUNK_LENGTH * 3 + 100); + } + + @Test + public void testLZ4Writer() throws IOException, ConfigurationException + { + compressor = LZ4Compressor.instance; + runTests("LZ4"); + } + + @Test + public void testDeflateWriter() throws IOException, ConfigurationException + { + compressor = DeflateCompressor.instance; + runTests("Deflate"); + } + + @Test + public void testSnappyWriter() throws IOException, ConfigurationException + { + compressor = SnappyCompressor.instance; + runTests("Snappy"); + } + + private void testWrite(File f, int bytesToTest) throws IOException, ConfigurationException + { + try + { + final String filename = f.getAbsolutePath(); + + SSTableMetadata.Collector sstableMetadataCollector = SSTableMetadata.createCollector(BytesType.instance).replayPosition(null); + CompressedSequentialWriter writer = new CompressedSequentialWriter( f, filename + ".metadata", false, new CompressionParameters( compressor, 32, Collections.emptyMap() ), sstableMetadataCollector); + + byte[] dataPre = new byte[bytesToTest]; + byte[] rawPost = new byte[bytesToTest]; + Random r = new Random(); + + // Test both write with byte[] and ByteBuffer + r.nextBytes(dataPre); + r.nextBytes(rawPost); + + writer.write(dataPre); + FileMark mark = writer.mark(); + + // Write enough garbage to transition chunk + for (int i = 0; i < CompressionParameters.DEFAULT_CHUNK_LENGTH; i++) + { + writer.write((byte)i); + } + writer.resetAndTruncate(mark); + writer.write(rawPost); + writer.close(); + + assert f.exists(); + CompressedRandomAccessReader reader = CompressedRandomAccessReader.open(filename, new CompressionMetadata(filename + ".metadata", f.length(), true)); + assertEquals(dataPre.length + rawPost.length, reader.length()); + byte[] result = new byte[(int)reader.length()]; + + reader.readFully(result); + + assert(reader.isEOF()); + reader.close(); + + byte[] fullInput = new byte[bytesToTest * 2]; + System.arraycopy(dataPre, 0, fullInput, 0, dataPre.length); + System.arraycopy(rawPost, 0, fullInput, bytesToTest, rawPost.length); + assert Arrays.equals(result, fullInput); + } + finally + { + // cleanup + if (f.exists()) + f.delete(); + File metadata = new File(f + ".metadata"); + if (metadata.exists()) + metadata.delete(); + } + } +}