diff --git a/CHANGES.txt b/CHANGES.txt index 01e0260626..512459f5e8 100644 --- a/CHANGES.txt +++ b/CHANGES.txt @@ -1,4 +1,5 @@ 5.0.5 + * Fix MAX_SEGMENT_SIZE < chunkSize in MmappedRegions::updateState (CASSANDRA-20636) * Full Java 17 support (CASSANDRA-20681) * Ensure replica filtering protection does not trigger unnecessary short read protection reads (CASSANDRA-20639) * Unified Compaction does not properly validate min and target sizes (CASSANDRA-20398) diff --git a/src/java/org/apache/cassandra/io/util/MmappedRegions.java b/src/java/org/apache/cassandra/io/util/MmappedRegions.java index 578217e279..5cbc65163e 100644 --- a/src/java/org/apache/cassandra/io/util/MmappedRegions.java +++ b/src/java/org/apache/cassandra/io/util/MmappedRegions.java @@ -183,7 +183,7 @@ public class MmappedRegions extends SharedCloseableImpl private void updateState(long length, int chunkSize) { // make sure the regions span whole chunks - long maxSize = (long) (MAX_SEGMENT_SIZE / chunkSize) * chunkSize; + long maxSize = Math.max(chunkSize, (long) (MAX_SEGMENT_SIZE / chunkSize) * chunkSize); state.length = length; long pos = state.getPosition(); while (pos < length) diff --git a/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java b/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java index f17301e064..a4b81e1018 100644 --- a/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java +++ b/test/unit/org/apache/cassandra/io/sstable/SSTableReaderTest.java @@ -30,12 +30,14 @@ import java.util.concurrent.Future; import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; +import java.util.function.Supplier; import java.util.stream.Collectors; import java.util.stream.IntStream; import java.util.stream.Stream; import com.google.common.collect.ImmutableList; import com.google.common.collect.Sets; +import org.junit.After; import org.junit.Assume; import org.junit.BeforeClass; import org.junit.Rule; @@ -88,6 +90,9 @@ import org.apache.cassandra.schema.CompressionParams; import org.apache.cassandra.schema.KeyspaceParams; import org.apache.cassandra.service.CacheService; import org.apache.cassandra.utils.ByteBufferUtil; +import org.apache.cassandra.utils.Throwables; +import org.apache.cassandra.utils.concurrent.Ref; +import org.apache.cassandra.utils.concurrent.SelfRefCounted; import org.mockito.Mockito; import static java.lang.String.format; @@ -97,10 +102,12 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; import static org.junit.Assert.assertNotNull; import static org.junit.Assert.assertTrue; +import static org.junit.Assert.fail; import static org.junit.Assume.assumeTrue; public class SSTableReaderTest { + public static final String KEYSPACE1 = "SSTableReaderTest"; public static final String CF_STANDARD = "Standard1"; public static final String CF_STANDARD2 = "Standard2"; @@ -111,6 +118,7 @@ public class SSTableReaderTest public static final String CF_STANDARD_LOW_INDEX_INTERVAL = "StandardLowIndexInterval"; public static final String CF_STANDARD_SMALL_BLOOM_FILTER = "StandardSmallBloomFilter"; + private final List> refsToRelease = new ArrayList<>(); private IPartitioner partitioner; Token t(int i) @@ -143,6 +151,28 @@ public class SSTableReaderTest CompactionManager.instance.disableAutoCompaction(); } + @After + public void teardown() + { + Throwable exceptions = null; + for (Ref ref : refsToRelease) + { + try + { + ref.release(); + } + catch (Throwable exc) + { + exceptions = Throwables.merge(exceptions, exc); + } + } + + if (exceptions != null) + fail("Unable to release all tracked references " + exceptions); + + refsToRelease.clear(); + } + @Test public void testGetPositionsForRanges() { @@ -310,7 +340,6 @@ public class SSTableReaderTest } } - @Test public void testOnDiskSizeCompressedBoundaries() { @@ -344,7 +373,6 @@ public class SSTableReaderTest onDiskSizeForRanges(sstable, Collections.singleton(new Range<>(partitioner.getMinimumToken(), t0(k - 1))))); // inclusive end } - long onDiskSizeForRanges(SSTableReader sstable, Collection> ranges) { return sstable.onDiskSizeForPartitionPositions(sstable.getPositionsForRanges(ranges)); @@ -375,12 +403,22 @@ public class SSTableReaderTest return s.substring(0, s.length() - n); } - @Test public void testSpannedIndexPositions() throws IOException + { + doTestSpannedIndexPositions(PageAware.PAGE_SIZE); + } + + @Test + public void testSpannedIndexPositionsWithMaxSegmentSizeSmallerThanPageSize() throws IOException + { + doTestSpannedIndexPositions(PageAware.PAGE_SIZE - 1); + } + + public void doTestSpannedIndexPositions(int maxSegmentSize) throws IOException { int originalMaxSegmentSize = MmappedRegions.MAX_SEGMENT_SIZE; - MmappedRegions.MAX_SEGMENT_SIZE = PageAware.PAGE_SIZE; + MmappedRegions.MAX_SEGMENT_SIZE = maxSegmentSize; try { @@ -655,7 +693,6 @@ public class SSTableReaderTest assertEquals(fpCount + 2, sstable.getFilterTracker().getFalsePositiveCount()); } - @Test public void testGetPositionsListenerCalls() { @@ -775,12 +812,12 @@ public class SSTableReaderTest Util.flush(store); CompactionManager.instance.performMaximal(store, false); - SSTableReaderWithFilter sstable = (SSTableReaderWithFilter) store.getLiveSSTables().iterator().next(); - sstable = (SSTableReaderWithFilter) sstable.cloneWithNewStart(dk(3)); - return sstable; + return (SSTableReaderWithFilter) trackReleaseableRef(() -> { + SSTableReader reader = store.getLiveSSTables().iterator().next(); + return reader.cloneWithNewStart(dk(3)); + }); } - @Test public void testOpeningSSTable() throws Exception { @@ -895,7 +932,8 @@ public class SSTableReaderTest components.remove(Components.SUMMARY); target = SSTableReader.open(store, desc, components, store.metadata); - try { + try + { assertTrue("Bloomfilter was not recreated", bloomModified < bloomFile.lastModified()); assertTrue("Summary was not recreated", summaryModified < summaryFile.lastModified()); } @@ -1054,9 +1092,8 @@ public class SSTableReaderTest if (sstable instanceof IndexSummarySupport) { new IndexSummaryComponent(((IndexSummarySupport) sstable).getIndexSummary(), sstable.getFirst(), sstable.getLast()).save(sstable.descriptor.fileFor(Components.SUMMARY), true); - SSTableReader reopened = SSTableReader.open(store, sstable.descriptor); + SSTableReader reopened = trackReleaseableRef(() -> SSTableReader.open(store, sstable.descriptor)); assert reopened.getFirst().getToken() instanceof LocalToken; - reopened.selfRef().release(); } } } @@ -1065,7 +1102,7 @@ public class SSTableReaderTest * see CASSANDRA-5407 */ @Test - public void testGetScannerForNoIntersectingRanges() throws Exception + public void testGetScannerForNoIntersectingRanges() { ColumnFamilyStore store = discardSSTables(KEYSPACE1, CF_STANDARD); partitioner = store.getPartitioner(); @@ -1240,7 +1277,7 @@ public class SSTableReaderTest txn.update(replacement, true); txn.finish(); } - R reopen = (R) SSTableReader.open(store, sstable.descriptor); + R reopen = (R) trackReleaseableRef(() -> SSTableReader.open(store, sstable.descriptor)); assert reopen.getIndexSummary().getSamplingLevel() == sstable.getIndexSummary().getSamplingLevel() + 1; } @@ -1311,9 +1348,9 @@ public class SSTableReaderTest { Keyspace keyspace = Keyspace.open(KEYSPACE1); ColumnFamilyStore cfs = keyspace.getColumnFamilyStore(CF_MOVE_AND_OPEN); - SSTableReader sstable = getNewSSTable(cfs); + SSTableReader sstable = trackReleaseableRef(() -> getNewSSTable(cfs)); + cfs.clearUnsafe(); - sstable.selfRef().release(); File tmpdir = new File(Files.createTempDirectory("testMoveAndOpen")); tmpdir.deleteOnExit(); SSTableId id = SSTableIdFactory.instance.defaultBuilder().generator(Stream.empty()).get(); @@ -1325,7 +1362,8 @@ public class SSTableReaderTest assertFalse(f.exists()); assertTrue(sstable.descriptor.fileFor(c).exists()); } - SSTableReader.moveAndOpenSSTable(cfs, sstable.descriptor, notLiveDesc, sstable.components, false); + trackReleaseableRef(() -> SSTableReader.moveAndOpenSSTable(cfs, sstable.descriptor, notLiveDesc, sstable.components, false)); + // make sure the files were moved: for (Component c : sstable.components) { @@ -1435,4 +1473,11 @@ public class SSTableReaderTest cfs.discardSSTables(System.currentTimeMillis()); return cfs; } + + private > T trackReleaseableRef(Supplier refSupplier) + { + T ref = refSupplier.get(); + refsToRelease.add(ref.selfRef()); + return ref; + } }