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]

Reply via email to