This is an automated email from the ASF dual-hosted git repository. asf-gitbox-commits pushed a commit to branch cassandra-5.0 in repository https://gitbox.apache.org/repos/asf/cassandra.git
commit 3944e9748293087e514c154784499823370c6d42 Merge: 5d84080193 be37969e89 Author: Caleb Rackliffe <[email protected]> AuthorDate: Mon Aug 17 14:37:25 2026 -0500 Merge branch 'cassandra-4.1' into cassandra-5.0 * cassandra-4.1: Make runWithCompactionsDisabled return non-null on success CHANGES.txt | 1 + .../org/apache/cassandra/db/ColumnFamilyStore.java | 88 +++++--- .../cassandra/db/compaction/CompactionManager.java | 5 +- .../org/apache/cassandra/db/view/TableViews.java | 3 + .../cassandra/exceptions/RequestFailureReason.java | 8 +- src/java/org/apache/cassandra/net/InboundSink.java | 5 +- .../apache/cassandra/service/StorageService.java | 2 + .../apache/cassandra/db/TruncateBlockingTest.java | 226 +++++++++++++++++++++ .../db/compaction/CancelCompactionsTest.java | 112 ++++++++++ 9 files changed, 419 insertions(+), 31 deletions(-) diff --cc CHANGES.txt index 2376d79d1b,75813be2a5..6e9059397f --- a/CHANGES.txt +++ b/CHANGES.txt @@@ -1,9 -1,6 +1,10 @@@ -4.1.13 +5.0.10 + * Avoid rebuilding per-SSTable SAI components unless missing or corrupted (CASSANDRA-21515) + * Propagate trickle_fsync settings to compressed SSTable writers (CASSANDRA-21487) + * Allow DatabaseDescriptor.setCompressedReadAheadBufferSizeInKb(0) to disable read-ahead buffer (CASSANDRA-21522) + * Return CorruptSSTableException if chunk metadata and file size are out of sync (CASSANDRA-21519) Merged from 4.0: + * Make runWithCompactionsDisabled return non-null on success (CASSANDRA-21527) * Validate authz before performing role check in LIST ROLES/PERMISSIONS (CASSANDRA-21560) * Ensure transferred_ranges reset on decommision re-attempt when pending ranges cannot be proven continous (CASSANDRA-16290) * Add validation to uncompressed length during decompression (CASSANDRA-21567) diff --cc src/java/org/apache/cassandra/db/ColumnFamilyStore.java index fd47b8e5bd,6192781ed7..0ab15db8b3 --- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java +++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java @@@ -112,8 -112,8 +112,9 @@@ import org.apache.cassandra.dht.Range import org.apache.cassandra.dht.Splitter; import org.apache.cassandra.dht.Token; import org.apache.cassandra.exceptions.ConfigurationException; +import org.apache.cassandra.exceptions.InvalidRequestException; import org.apache.cassandra.exceptions.StartupException; + import org.apache.cassandra.exceptions.TruncateException; import org.apache.cassandra.index.SecondaryIndexManager; import org.apache.cassandra.index.internal.CassandraIndex; import org.apache.cassandra.index.transactions.UpdateTransaction; @@@ -1904,8 -1813,15 +1913,15 @@@ public class ColumnFamilyStore implemen TimeUUID session = sst.getPendingRepair(); return session != null && sessions.contains(session); }; - return runWithCompactionsDisabled(() -> compactionStrategyManager.releaseRepairData(sessions), - predicate, OperationType.STREAM, false, true, true); + CleanupSummary summary = runWithCompactionsDisabled(() -> compactionStrategyManager.releaseRepairData(sessions), - predicate, false, true, true); ++ predicate, OperationType.STREAM, false, true, true); + if (summary == null) + { + logger.warn("Unable to cancel in-progress compactions for {}.{}, could not force release repair data for sessions {}", - keyspace.getName(), name, sessions); ++ getKeyspaceName(), name, sessions); + return new CleanupSummary(this, Collections.emptySet(), new HashSet<>(sessions)); + } + return summary; } else { @@@ -2794,39 -2681,58 +2810,57 @@@ for (SSTableReader sstable : cfs.getLiveSSTables()) now = Math.max(now, sstable.maxDataAge); truncatedAt = now; - - Runnable truncateRunnable = new Runnable() + Throwable failure = null; + try { - public void run() - { - logger.info("Truncating {}.{} with truncatedAt={}", getKeyspaceName(), getTableName(), truncatedAt); - // since truncation can happen at different times on different nodes, we need to make sure - // that any repairs are aborted, otherwise we might clear the data on one node and then - // stream in data that is actually supposed to have been deleted - ActiveRepairService.instance().abort((prs) -> prs.getTableIds().contains(metadata.id), - "Stopping parent sessions {} due to truncation of tableId="+metadata.id); - data.notifyTruncated(truncatedAt); - Boolean succeeded = runWithCompactionsDisabled(() -> runTruncate(truncatedAt, noSnapshot, replayAfter), - true, true); ++ Boolean succeeded = runWithCompactionsDisabled(() -> runTruncate(truncatedAt, noSnapshot, replayAfter), OperationType.P0, true, true); + // null means compactions couldn't be disabled, not failure of the truncate work itself + if (succeeded == null) + failure = new TruncateException("Unable to stop in-progress compactions. Please retry when current compaction tasks have completed or overall system load has decreased."); + else if (!succeeded) + failure = new TruncateException("Truncate failed."); + } + catch (Throwable t) + { + failure = t; + } - if (!noSnapshot && isAutoSnapshotEnabled()) - snapshot(Keyspace.getTimestampedSnapshotNameWithPrefix(name, SNAPSHOT_TRUNCATE_PREFIX), DatabaseDescriptor.getAutoSnapshotTtl()); + try + { + viewManager.build(); + } + catch (Throwable t) + { + failure = merge(failure, t); + } - discardSSTables(truncatedAt); + maybeFail(failure); - indexManager.truncateAllIndexesBlocking(truncatedAt); - viewManager.truncateBlocking(replayAfter, truncatedAt); - logger.info("Truncate of {}.{} is complete", keyspace.getName(), name); ++ logger.info("Truncate of {}.{} is complete", getKeyspaceName(), name); + } - SystemKeyspace.saveTruncationRecord(ColumnFamilyStore.this, truncatedAt, replayAfter); - logger.trace("cleaning out row cache"); - invalidateCaches(); + private boolean runTruncate(long truncatedAt, boolean noSnapshot, CommitLogPosition replayAfter) + { - logger.info("Truncating {}.{} with truncatedAt={}", keyspace.getName(), getTableName(), truncatedAt); ++ logger.info("Truncating {}.{} with truncatedAt={}", getKeyspaceName(), getTableName(), truncatedAt); + // since truncation can happen at different times on different nodes, we need to make sure + // that any repairs are aborted, otherwise we might clear the data on one node and then + // stream in data that is actually supposed to have been deleted - ActiveRepairService.instance.abort((prs) -> prs.getTableIds().contains(metadata.id), - "Stopping parent sessions {} due to truncation of tableId="+metadata.id); ++ ActiveRepairService.instance().abort((prs) -> prs.getTableIds().contains(metadata.id), ++ "Stopping parent sessions {} due to truncation of tableId=" + metadata.id); + data.notifyTruncated(truncatedAt); - } - }; - if (!noSnapshot && DatabaseDescriptor.isAutoSnapshot()) ++ if (!noSnapshot && isAutoSnapshotEnabled()) + snapshot(Keyspace.getTimestampedSnapshotNameWithPrefix(name, SNAPSHOT_TRUNCATE_PREFIX), DatabaseDescriptor.getAutoSnapshotTtl()); - runWithCompactionsDisabled(FutureTask.callable(truncateRunnable), OperationType.P0, true, true); + discardSSTables(truncatedAt); - viewManager.build(); + indexManager.truncateAllIndexesBlocking(truncatedAt); + viewManager.truncateBlocking(replayAfter, truncatedAt); - logger.info("Truncate of {}.{} is complete", getKeyspaceName(), name); + SystemKeyspace.saveTruncationRecord(ColumnFamilyStore.this, truncatedAt, replayAfter); + logger.trace("cleaning out row cache"); + invalidateCaches(); + return true; } /** diff --cc src/java/org/apache/cassandra/db/compaction/CompactionManager.java index 7c93558004,d6422bdd67..a20dc32a6e --- a/src/java/org/apache/cassandra/db/compaction/CompactionManager.java +++ b/src/java/org/apache/cassandra/db/compaction/CompactionManager.java @@@ -998,9 -981,9 +998,9 @@@ public class CompactionManager implemen // here we compute the task off the compaction executor, so having that present doesn't // confuse runWithCompactionsDisabled -- i.e., we don't want to deadlock ourselves, waiting // for ourselves to finish/acknowledge cancellation before continuing. - CompactionTasks tasks = cfStore.getCompactionStrategyManager().getMaximalTasks(gcBefore, splitOutput); + CompactionTasks tasks = cfStore.getCompactionStrategyManager().getMaximalTasks(gcBefore, splitOutput, operationType); - if (tasks.isEmpty()) + if (tasks == null || tasks.isEmpty()) return Collections.emptyList(); List<Future<?>> futures = new ArrayList<>(); @@@ -1048,6 -1030,9 +1048,9 @@@ false, false)) { + if (tasks == null) - throw new RuntimeException("Unable to cancel in-progress compactions for " + cfStore.keyspace.getName() + '.' + cfStore.getTableName() + ". Usually retrying will work."); ++ throw new RuntimeException("Unable to cancel in-progress compactions for " + cfStore.getKeyspaceName() + '.' + cfStore.getTableName() + ". Usually retrying will work."); + if (tasks.isEmpty()) return; diff --cc src/java/org/apache/cassandra/exceptions/RequestFailureReason.java index 5e9900280c,3b911b68eb..99cc07be59 --- a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java +++ b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java @@@ -36,8 -38,7 +36,9 @@@ public enum RequestFailureReaso INCOMPATIBLE_SCHEMA (3), READ_SIZE (4), NODE_DOWN (5), + INDEX_NOT_AVAILABLE (6), - READ_TOO_MANY_INDEXES (7); ++ READ_TOO_MANY_INDEXES (7), + TRUNCATE_FAILED (12); public static final Serializer serializer = new Serializer(); diff --cc src/java/org/apache/cassandra/service/StorageService.java index a358841e2d,26e52ecc5f..f428501217 --- a/src/java/org/apache/cassandra/service/StorageService.java +++ b/src/java/org/apache/cassandra/service/StorageService.java @@@ -7756,40 -7190,4 +7756,42 @@@ public class StorageService extends Not return DatabaseDescriptor.getPaxosRepairRaceWait(); } + @Override + public List<String> getTablesForKeyspace(String keyspace) + { + return Keyspace.open(keyspace).getColumnFamilyStores().stream().map(cfs -> cfs.name).collect(Collectors.toList()); + } + + @Override + public List<String> mutateSSTableRepairedState(boolean repaired, boolean preview, String keyspace, List<String> tableNames) + { + Map<String, ColumnFamilyStore> tables = Keyspace.open(keyspace).getColumnFamilyStores() + .stream().collect(Collectors.toMap(c -> c.name, c -> c)); + for (String tableName : tableNames) + { + if (!tables.containsKey(tableName)) + throw new RuntimeException("Table " + tableName + " does not exist in keyspace " + keyspace); + } + + // only select SSTables that are unrepaired when repaired is true and vice versa + Predicate<SSTableReader> predicate = sst -> repaired != sst.isRepaired(); + + // mutate SSTables + long repairedAt = !repaired ? 0 : currentTimeMillis(); + List<String> sstablesTouched = new ArrayList<>(); + for (String tableName : tableNames) + { + ColumnFamilyStore table = tables.get(tableName); + Set<SSTableReader> result = table.runWithCompactionsDisabled(() -> { + Set<SSTableReader> sstables = table.getLiveSSTables().stream().filter(predicate).collect(Collectors.toSet()); + if (!preview) + table.getCompactionStrategyManager().mutateRepaired(sstables, repairedAt, null, false); + return sstables; + }, predicate, OperationType.ANTICOMPACTION, true, false, true); ++ if (result == null) ++ throw new RuntimeException("Unable to cancel in-progress compactions for " + keyspace + '.' + tableName + ". Usually retrying will work."); + sstablesTouched.addAll(result.stream().map(sst -> sst.descriptor.baseFile().name()).collect(Collectors.toList())); + } + return sstablesTouched; + } } diff --cc test/unit/org/apache/cassandra/db/TruncateBlockingTest.java index 0000000000,24433a6306..63ac47145e mode 000000,100644..100644 --- a/test/unit/org/apache/cassandra/db/TruncateBlockingTest.java +++ b/test/unit/org/apache/cassandra/db/TruncateBlockingTest.java @@@ -1,0 -1,99 +1,226 @@@ + /* + * 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.db; + ++import java.util.Collections; ++ + import org.assertj.core.api.Assertions; + import org.jboss.byteman.contrib.bmunit.BMRule; + import org.jboss.byteman.contrib.bmunit.BMUnitRunner; + import org.junit.Test; + import org.junit.runner.RunWith; + + import org.apache.cassandra.cql3.CQLTester; ++import org.apache.cassandra.db.compaction.CompactionInfo; ++import org.apache.cassandra.db.compaction.CompactionManager; + import org.apache.cassandra.db.compaction.OperationType; + import org.apache.cassandra.db.lifecycle.LifecycleTransaction; + import org.apache.cassandra.exceptions.TruncateException; + import org.apache.cassandra.io.sstable.format.SSTableReader; ++import org.apache.cassandra.service.StorageService; + ++import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID; + import static org.junit.Assert.assertFalse; + import static org.junit.Assert.assertNotNull; + + @RunWith(BMUnitRunner.class) + public class TruncateBlockingTest extends CQLTester + { ++ @Test ++ public void testTruncateFailsWhenCompactionsCannotBeDisabled() ++ { ++ createTable("CREATE TABLE %s (id int PRIMARY KEY, v text)"); ++ ++ execute("INSERT INTO %s (id, v) VALUES (1, 'a')"); ++ execute("INSERT INTO %s (id, v) VALUES (2, 'b')"); ++ execute("INSERT INTO %s (id, v) VALUES (3, 'c')"); ++ flush(); ++ ++ ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); ++ ++ // Register a P0-priority compaction holder for this table to force ++ // runWithCompactionsDisabled to return null immediately. ++ CompactionInfo.Holder holder = new CompactionInfo.Holder() ++ { ++ public CompactionInfo getCompactionInfo() ++ { ++ return new CompactionInfo(cfs.metadata(), ++ OperationType.P0, ++ 0, ++ 100, ++ 100, ++ nextTimeUUID(), ++ Collections.emptySet()); ++ } ++ ++ public boolean isGlobal() ++ { ++ return false; ++ } ++ }; ++ ++ CompactionManager.instance.active.beginCompaction(holder); ++ try ++ { ++ Assertions.assertThatThrownBy(cfs::truncateBlocking) ++ .as("Unable to stop compaction. Usually retrying truncate will work") ++ .isInstanceOf(TruncateException.class); ++ ++ assertRows(execute("SELECT * FROM %s WHERE id = 1"), row(1, "a")); ++ assertRows(execute("SELECT * FROM %s WHERE id = 2"), row(2, "b")); ++ assertRows(execute("SELECT * FROM %s WHERE id = 3"), row(3, "c")); ++ assertFalse("SSTables should still be present after truncation failure", ++ cfs.getLiveSSTables().isEmpty()); ++ ++ } ++ finally ++ { ++ CompactionManager.instance.active.finishCompaction(holder); ++ } ++ } ++ + @Test + @BMRule(name = "no-op waitForCessation", + targetClass = "org.apache.cassandra.db.compaction.CompactionManager", + targetMethod = "waitForCessation", + action = "return;") - public void testTruncateFailsWhenCompactionsDoNotStopInTime() throws Throwable ++ public void testTruncateFailsWhenCompactionsDoNotStopInTime() + { + createTable("CREATE TABLE %s (id int PRIMARY KEY, v text)"); + + execute("INSERT INTO %s (id, v) VALUES (1, 'a')"); + execute("INSERT INTO %s (id, v) VALUES (2, 'b')"); + execute("INSERT INTO %s (id, v) VALUES (3, 'c')"); + flush(); + + ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); + SSTableReader sstable = cfs.getLiveSSTables().iterator().next(); + + // Mark the sstable as compacting directly in the tracker, without registering a + // CompactionInfo.Holder. There is nothing for interruptCompactionForCFs to stop, so + // runWithCompactionsDisabled falls through to waitForCessation, then finds the sstable + // still in the compacting set and returns null. + try (LifecycleTransaction txn = cfs.getTracker().tryModify(sstable, OperationType.ANTICOMPACTION)) + { + assertNotNull("Unable to mark sstable compacting", txn); + + Assertions.assertThatThrownBy(cfs::truncateBlocking) + .as("Unable to stop compaction. Usually retrying truncate will work") + .isInstanceOf(TruncateException.class); + + assertRows(execute("SELECT * FROM %s WHERE id = 1"), row(1, "a")); + assertRows(execute("SELECT * FROM %s WHERE id = 2"), row(2, "b")); + assertRows(execute("SELECT * FROM %s WHERE id = 3"), row(3, "c")); + assertFalse("SSTables should still be present after truncation failure", + cfs.getLiveSSTables().isEmpty()); + } + } + + @Test - public void testRebuildOnFailedScrubReturnsFalseWhenTruncateFails() throws Throwable ++ public void testRebuildOnFailedScrubReturnsFalseWhenTruncateFails() + { + createTable("CREATE TABLE %s (id int PRIMARY KEY, v text)"); + // rebuildOnFailedScrub only applies to indexes with their own backing table - createIndex("CREATE INDEX ON %s (v)"); ++ createIndex("CREATE INDEX ON %s (v) USING 'legacy_local_table'"); + + execute("INSERT INTO %s (id, v) VALUES (1, 'a')"); + flush(); + + ColumnFamilyStore baseCfs = getCurrentColumnFamilyStore(); + ColumnFamilyStore indexCfs = baseCfs.indexManager.getAllIndexColumnFamilyStores().iterator().next(); - SSTableReader sstable = indexCfs.getLiveSSTables().iterator().next(); + - try (LifecycleTransaction txn = indexCfs.getTracker().tryModify(sstable, OperationType.ANTICOMPACTION)) ++ // Register a P0-priority compaction holder for the index cfs to force truncateBlocking to fail. ++ CompactionInfo.Holder holder = new CompactionInfo.Holder() + { - assertNotNull("Unable to mark sstable compacting", txn); ++ public CompactionInfo getCompactionInfo() ++ { ++ return new CompactionInfo(indexCfs.metadata(), ++ OperationType.P0, ++ 0, ++ 100, ++ 100, ++ nextTimeUUID(), ++ Collections.emptySet()); ++ } ++ ++ public boolean isGlobal() ++ { ++ return false; ++ } ++ }; + ++ CompactionManager.instance.active.beginCompaction(holder); ++ try ++ { + RuntimeException scrubFailure = new RuntimeException("original scrub failure"); + // rebuildOnFailedScrub should report the rebuild as unsuccessful + assertFalse("rebuildOnFailedScrub should return false when it can't truncate the index", + indexCfs.rebuildOnFailedScrub(scrubFailure)); + } ++ finally ++ { ++ CompactionManager.instance.active.finishCompaction(holder); ++ } ++ } ++ ++ @Test ++ public void testMutateSSTableRepairedStateThrowsWhenCompactionsCannotBeDisabled() ++ { ++ createTable("CREATE TABLE %s (id int PRIMARY KEY, v text)"); ++ ++ execute("INSERT INTO %s (id, v) VALUES (1, 'a')"); ++ flush(); ++ ++ ColumnFamilyStore cfs = getCurrentColumnFamilyStore(); ++ ++ // Register a P0-priority compaction holder for this table to force runWithCompactionsDisabled ++ // to return null ++ CompactionInfo.Holder holder = new CompactionInfo.Holder() ++ { ++ public CompactionInfo getCompactionInfo() ++ { ++ return new CompactionInfo(cfs.metadata(), ++ OperationType.P0, ++ 0, ++ 100, ++ 100, ++ nextTimeUUID(), ++ Collections.emptySet()); ++ } ++ ++ public boolean isGlobal() ++ { ++ return false; ++ } ++ }; ++ ++ CompactionManager.instance.active.beginCompaction(holder); ++ try ++ { ++ // mutateSSTableRepairedState should report the null runWithCompactionsDisabled result as a ++ // failure to the caller ++ Assertions.assertThatThrownBy(() -> StorageService.instance.mutateSSTableRepairedState(true, false, keyspace(), Collections.singletonList(currentTable()))) ++ .as("Unable to cancel in-progress compactions. Usually retrying will work") ++ .isInstanceOf(RuntimeException.class); ++ } ++ finally ++ { ++ CompactionManager.instance.active.finishCompaction(holder); ++ } + } -} ++} diff --cc test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java index b3a182f186,12af15c05e..88a3987b33 --- a/test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java @@@ -32,10 -32,14 +32,11 @@@ import java.util.stream.Collectors import com.google.common.collect.ImmutableSet; import com.google.common.util.concurrent.Uninterruptibles; -import net.bytebuddy.ByteBuddy; -import net.bytebuddy.agent.ByteBuddyAgent; -import net.bytebuddy.dynamic.loading.ClassReloadingStrategy; -import net.bytebuddy.implementation.StubMethod; - + import org.assertj.core.api.Assertions; +import org.junit.Assume; import org.junit.Test; +import org.apache.cassandra.config.DatabaseDescriptor; import org.apache.cassandra.cql3.CQLTester; import org.apache.cassandra.db.ColumnFamilyStore; import org.apache.cassandra.db.lifecycle.LifecycleTransaction; @@@ -372,6 -372,118 +375,115 @@@ public class CancelCompactionsTest exte return new Murmur3Partitioner.LongToken(t); } - private LifecycleTransaction blockCompactions(ColumnFamilyStore cfs, List<SSTableReader> sstables) ++ private CompactionInfo.Holder p0Holder(ColumnFamilyStore cfs) + { - LifecycleTransaction txn = cfs.getTracker().tryModify(sstables, OperationType.COMPACTION); - assertNotNull(txn); - return txn; - } ++ return new CompactionInfo.Holder() ++ { ++ public CompactionInfo getCompactionInfo() ++ { ++ return new CompactionInfo(cfs.metadata(), OperationType.P0, 0, 100, 100, nextTimeUUID(), Collections.emptySet()); ++ } + - private ClassReloadingStrategy stubWaitForCessation() - { - ByteBuddyAgent.install(); - ClassReloadingStrategy strategy = ClassReloadingStrategy.fromInstalledAgent(); - new ByteBuddy().redefine(CompactionManager.class) - .method(named("waitForCessation")) - .intercept(StubMethod.INSTANCE) - .make() - .load(CompactionManager.class.getClassLoader(), strategy); - return strategy; ++ public boolean isGlobal() ++ { ++ return false; ++ } ++ }; + } + + @Test - public void testForceCompactionThrowsWhenCompactionsCannotBeDisabled() throws Exception ++ public void testForceCompactionThrowsWhenCompactionsCannotBeDisabled() + { + ColumnFamilyStore cfs = MockSchema.newCFS(); - List<SSTableReader> sstables = createSSTables(cfs, 3, 0); ++ createSSTables(cfs, 3, 0); + - LifecycleTransaction txn = blockCompactions(cfs, sstables); - ClassReloadingStrategy strategy = stubWaitForCessation(); ++ // Register a P0-priority compaction holder for this table to force runWithCompactionsDisabled ++ // to return null ++ CompactionInfo.Holder holder = p0Holder(cfs); ++ CompactionManager.instance.active.beginCompaction(holder); + try + { + Range<Token> allData = new Range<>(cfs.getPartitioner().getMinimumToken(), cfs.getPartitioner().getMaximumToken()); + + // forceCompaction should fail at the null runWithCompactionsDisabled result + Assertions.assertThatThrownBy(() -> cfs.forceCompactionForTokenRange(Collections.singleton(allData))) + .as("Unable to cancel in-progress compactions. Usually retrying will work") + .isInstanceOf(RuntimeException.class); + } + finally + { - strategy.reset(CompactionManager.class); - txn.abort(); ++ CompactionManager.instance.active.finishCompaction(holder); + } + } + + @Test + public void testGarbageCollectReturnsUnableToCancelWhenCompactionsCannotBeDisabled() throws Throwable + { + ColumnFamilyStore cfs = MockSchema.newCFS(); - List<SSTableReader> sstables = createSSTables(cfs, 3, 0); ++ createSSTables(cfs, 3, 0); + - LifecycleTransaction txn = blockCompactions(cfs, sstables); - ClassReloadingStrategy strategy = stubWaitForCessation(); ++ // Register a P0-priority compaction holder for this table to force runWithCompactionsDisabled ++ // to return null immediately. ++ CompactionInfo.Holder holder = p0Holder(cfs); ++ CompactionManager.instance.active.beginCompaction(holder); + try + { - // garbageCollect goes through parallelAllSSTableOperation, which must handle a null - // LifecycleTransaction from markAllCompacting without NPEing. ++ // garbageCollect goes through withAllSSTables, which must be able to pass a null ++ // LifecycleTransaction to its caller-supplied op without NPEing. + assertEquals(CompactionManager.AllSSTableOpStatus.UNABLE_TO_CANCEL, cfs.garbageCollect(TombstoneOption.ROW, 0)); + } + finally + { - strategy.reset(CompactionManager.class); - txn.abort(); ++ CompactionManager.instance.active.finishCompaction(holder); + } + } + + @Test - public void testReleaseRepairDataReturnsUnsuccessfulWhenCompactionsCannotBeDisabled() throws Exception ++ public void testReleaseRepairDataReturnsUnsuccessfulWhenCompactionsCannotBeDisabled() + { + ColumnFamilyStore cfs = MockSchema.newCFS(); - List<SSTableReader> sstables = createSSTables(cfs, 3, 0); ++ createSSTables(cfs, 3, 0); + - // releaseRepairData only tries to cancel compactions on sstables pending one of these sessions, - // so tag one of our blocked sstables with a pending session to make it match - TimeUUID pendingSession = nextTimeUUID(); - Set<TimeUUID> sessions = ImmutableSet.of(pendingSession, nextTimeUUID()); - AbstractPendingRepairTest.mutateRepaired(sstables.get(0), pendingSession, false); ++ Set<TimeUUID> sessions = ImmutableSet.of(nextTimeUUID(), nextTimeUUID()); + - LifecycleTransaction txn = blockCompactions(cfs, sstables); - ClassReloadingStrategy strategy = stubWaitForCessation(); ++ // Register a P0-priority compaction holder for this table to force runWithCompactionsDisabled ++ // to return null immediately. ++ CompactionInfo.Holder holder = p0Holder(cfs); ++ CompactionManager.instance.active.beginCompaction(holder); + try + { + // force release should report the sessions as unsuccessful rather than throwing + CleanupSummary summary = cfs.releaseRepairData(sessions, true); + assertTrue(summary.successful.isEmpty()); + assertEquals(sessions, summary.unsuccessful); + } + finally + { - strategy.reset(CompactionManager.class); - txn.abort(); ++ CompactionManager.instance.active.finishCompaction(holder); + } + } + + @Test - public void testSubmitMaximalNoOpsWhenCompactionsCannotBeDisabled() throws Exception ++ public void testSubmitMaximalNoOpsWhenCompactionsCannotBeDisabled() + { + ColumnFamilyStore cfs = MockSchema.newCFS(); - List<SSTableReader> sstables = createSSTables(cfs, 3, 0); ++ createSSTables(cfs, 3, 0); + - LifecycleTransaction txn = blockCompactions(cfs, sstables); - ClassReloadingStrategy strategy = stubWaitForCessation(); ++ // Register a P0-priority compaction holder for this table to force runWithCompactionsDisabled ++ // to return null immediately. ++ CompactionInfo.Holder holder = p0Holder(cfs); ++ CompactionManager.instance.active.beginCompaction(holder); + try + { + // submitMaximal should no-op on null getMaximalTasks result + assertTrue(CompactionManager.instance.submitMaximal(cfs, -1, false).isEmpty()); + } + finally + { - strategy.reset(CompactionManager.class); - txn.abort(); ++ CompactionManager.instance.active.finishCompaction(holder); + } + } + private List<SSTableReader> createSSTables(ColumnFamilyStore cfs, int count, int startGeneration) { List<SSTableReader> sstables = new ArrayList<>(); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
