This is an automated email from the ASF dual-hosted git repository.

asf-gitbox-commits pushed a commit to branch cassandra-4.1
in repository https://gitbox.apache.org/repos/asf/cassandra.git

commit be37969e89a5741cc41687a24eeee77d461917b4
Merge: 4827aafd02 7b4170491d
Author: Caleb Rackliffe <[email protected]>
AuthorDate: Mon Aug 17 14:15:19 2026 -0500

    Merge branch 'cassandra-4.0' into cassandra-4.1
    
    * cassandra-4.0:
      Make runWithCompactionsDisabled return non-null on success

 CHANGES.txt                                        |   1 +
 .../org/apache/cassandra/db/ColumnFamilyStore.java |  89 ++++++++++-----
 .../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/db/TruncateBlockingTest.java  |  99 +++++++++++++++++
 .../db/compaction/CancelCompactionsTest.java       | 121 +++++++++++++++++++++
 8 files changed, 300 insertions(+), 31 deletions(-)

diff --cc CHANGES.txt
index b76d2a3754,50ad389065..75813be2a5
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@@ -1,5 -1,5 +1,6 @@@
 -4.0.22
 +4.1.13
 +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 2d23522280,2ac4783d89..6192781ed7
--- a/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
+++ b/src/java/org/apache/cassandra/db/ColumnFamilyStore.java
@@@ -73,46 -41,25 +73,47 @@@ import com.google.common.util.concurren
  import org.slf4j.Logger;
  import org.slf4j.LoggerFactory;
  
 -import org.apache.cassandra.cache.*;
 -import org.apache.cassandra.concurrent.*;
 -import org.apache.cassandra.config.*;
 +import org.apache.cassandra.cache.CounterCacheKey;
 +import org.apache.cassandra.cache.IRowCacheEntry;
 +import org.apache.cassandra.cache.RowCacheKey;
 +import org.apache.cassandra.cache.RowCacheSentinel;
 +import org.apache.cassandra.concurrent.ExecutorPlus;
 +import org.apache.cassandra.concurrent.FutureTask;
 +import org.apache.cassandra.config.CassandraRelevantProperties;
 +import org.apache.cassandra.config.DatabaseDescriptor;
 +import org.apache.cassandra.config.DurationSpec;
  import org.apache.cassandra.db.commitlog.CommitLog;
  import org.apache.cassandra.db.commitlog.CommitLogPosition;
 -import org.apache.cassandra.db.compaction.*;
 +import org.apache.cassandra.db.compaction.AbstractCompactionStrategy;
 +import org.apache.cassandra.db.compaction.CompactionManager;
 +import org.apache.cassandra.db.compaction.CompactionStrategyManager;
 +import org.apache.cassandra.db.compaction.OperationType;
 +import org.apache.cassandra.db.compaction.Verifier;
  import org.apache.cassandra.db.filter.ClusteringIndexFilter;
  import org.apache.cassandra.db.filter.DataLimits;
 -import org.apache.cassandra.db.streaming.CassandraStreamManager;
 -import org.apache.cassandra.db.repair.CassandraTableRepairManager;
 -import org.apache.cassandra.db.view.TableViews;
 -import org.apache.cassandra.db.lifecycle.*;
 +import org.apache.cassandra.db.memtable.Flushing;
 +import org.apache.cassandra.db.memtable.Memtable;
 +import org.apache.cassandra.db.memtable.ShardBoundaries;
 +import org.apache.cassandra.db.lifecycle.LifecycleNewTracker;
 +import org.apache.cassandra.db.lifecycle.LifecycleTransaction;
 +import org.apache.cassandra.db.lifecycle.SSTableSet;
 +import org.apache.cassandra.db.lifecycle.Tracker;
 +import org.apache.cassandra.db.lifecycle.View;
  import org.apache.cassandra.db.partitions.CachedPartition;
  import org.apache.cassandra.db.partitions.PartitionUpdate;
 +import org.apache.cassandra.db.repair.CassandraTableRepairManager;
  import org.apache.cassandra.db.rows.CellPath;
 -import org.apache.cassandra.dht.*;
 +import org.apache.cassandra.db.streaming.CassandraStreamManager;
 +import org.apache.cassandra.db.view.TableViews;
 +import org.apache.cassandra.dht.AbstractBounds;
 +import org.apache.cassandra.dht.Bounds;
 +import org.apache.cassandra.dht.IPartitioner;
  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.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;
@@@ -1801,11 -1670,18 +1810,18 @@@ public class ColumnFamilyStore implemen
          if (force)
          {
              Predicate<SSTableReader> predicate = sst -> {
 -                UUID session = sst.getPendingRepair();
 +                TimeUUID session = sst.getPendingRepair();
                  return session != null && sessions.contains(session);
              };
-             return runWithCompactionsDisabled(() -> 
compactionStrategyManager.releaseRepairData(sessions),
-                                               predicate, false, true, true);
+             CleanupSummary summary = runWithCompactionsDisabled(() -> 
compactionStrategyManager.releaseRepairData(sessions),
+                                                                 predicate, 
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);
+                 return new CleanupSummary(this, Collections.emptySet(), new 
HashSet<>(sessions));
+             }
+             return summary;
          }
          else
          {
@@@ -2665,39 -2344,58 +2681,58 @@@
              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={}", 
keyspace.getName(), 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);
+             // 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 && DatabaseDescriptor.isAutoSnapshot())
-                 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);
+     }
  
-                 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);
+         // 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);
  
-             }
-         };
+         if (!noSnapshot && DatabaseDescriptor.isAutoSnapshot())
 -            snapshot(Keyspace.getTimestampedSnapshotNameWithPrefix(name, 
SNAPSHOT_TRUNCATE_PREFIX));
++            snapshot(Keyspace.getTimestampedSnapshotNameWithPrefix(name, 
SNAPSHOT_TRUNCATE_PREFIX), DatabaseDescriptor.getAutoSnapshotTtl());
  
-         runWithCompactionsDisabled(FutureTask.callable(truncateRunnable), 
true, true);
+         discardSSTables(truncatedAt);
  
-         viewManager.build();
+         indexManager.truncateAllIndexesBlocking(truncatedAt);
+         viewManager.truncateBlocking(replayAfter, truncatedAt);
  
-         logger.info("Truncate of {}.{} is complete", keyspace.getName(), 
name);
+         SystemKeyspace.saveTruncationRecord(ColumnFamilyStore.this, 
truncatedAt, replayAfter);
+         logger.trace("cleaning out row cache");
+         invalidateCaches();
+         return true;
      }
  
      /**
diff --cc src/java/org/apache/cassandra/exceptions/RequestFailureReason.java
index 3d3476a139,6664070e97..3b911b68eb
--- a/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java
+++ b/src/java/org/apache/cassandra/exceptions/RequestFailureReason.java
@@@ -36,8 -36,7 +36,9 @@@ public enum RequestFailureReaso
      READ_TOO_MANY_TOMBSTONES (1),
      TIMEOUT                  (2),
      INCOMPATIBLE_SCHEMA      (3),
 +    READ_SIZE                (4),
-     NODE_DOWN                (5);
++    NODE_DOWN                (5),
+     TRUNCATE_FAILED          (12);
  
      public static final Serializer serializer = new Serializer();
  
diff --cc 
test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java
index 67421ba6d2,7ea27a1367..12af15c05e
--- a/test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java
+++ b/test/unit/org/apache/cassandra/db/compaction/CancelCompactionsTest.java
@@@ -32,8 -33,14 +32,14 @@@ 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.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;
@@@ -51,14 -58,12 +57,17 @@@ import org.apache.cassandra.io.sstable.
  import org.apache.cassandra.locator.InetAddressAndPort;
  import org.apache.cassandra.locator.RangesAtEndpoint;
  import org.apache.cassandra.locator.Replica;
++import org.apache.cassandra.repair.consistent.admin.CleanupSummary;
++import org.apache.cassandra.schema.CompactionParams.TombstoneOption;
  import org.apache.cassandra.schema.MockSchema;
  import org.apache.cassandra.schema.TableMetadata;
  import org.apache.cassandra.service.ActiveRepairService;
  import org.apache.cassandra.streaming.PreviewKind;
  import org.apache.cassandra.utils.FBUtilities;
 +import org.apache.cassandra.utils.TimeUUID;
  
++import static net.bytebuddy.matcher.ElementMatchers.named;
 +import static org.apache.cassandra.utils.TimeUUID.Generator.nextTimeUUID;
  import static org.junit.Assert.assertEquals;
  import static org.junit.Assert.assertFalse;
  import static org.junit.Assert.assertNotNull;
@@@ -363,6 -372,118 +372,118 @@@ public class CancelCompactionsTest exte
          return new Murmur3Partitioner.LongToken(t);
      }
  
+     private LifecycleTransaction blockCompactions(ColumnFamilyStore cfs, 
List<SSTableReader> sstables)
+     {
+         LifecycleTransaction txn = cfs.getTracker().tryModify(sstables, 
OperationType.COMPACTION);
+         assertNotNull(txn);
+         return txn;
+     }
+ 
+     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;
+     }
+ 
+     @Test
+     public void testForceCompactionThrowsWhenCompactionsCannotBeDisabled() 
throws Exception
+     {
+         ColumnFamilyStore cfs = MockSchema.newCFS();
+         List<SSTableReader> sstables = createSSTables(cfs, 3, 0);
+ 
+         LifecycleTransaction txn = blockCompactions(cfs, sstables);
+         ClassReloadingStrategy strategy = stubWaitForCessation();
+         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();
+         }
+     }
+ 
+     @Test
+     public void 
testGarbageCollectReturnsUnableToCancelWhenCompactionsCannotBeDisabled() throws 
Throwable
+     {
+         ColumnFamilyStore cfs = MockSchema.newCFS();
+         List<SSTableReader> sstables = createSSTables(cfs, 3, 0);
+ 
+         LifecycleTransaction txn = blockCompactions(cfs, sstables);
+         ClassReloadingStrategy strategy = stubWaitForCessation();
+         try
+         {
+             // garbageCollect goes through parallelAllSSTableOperation, which 
must handle a null
+             // LifecycleTransaction from markAllCompacting without NPEing.
+             
assertEquals(CompactionManager.AllSSTableOpStatus.UNABLE_TO_CANCEL, 
cfs.garbageCollect(TombstoneOption.ROW, 0));
+         }
+         finally
+         {
+             strategy.reset(CompactionManager.class);
+             txn.abort();
+         }
+     }
+ 
+     @Test
+     public void 
testReleaseRepairDataReturnsUnsuccessfulWhenCompactionsCannotBeDisabled() 
throws Exception
+     {
+         ColumnFamilyStore cfs = MockSchema.newCFS();
+         List<SSTableReader> sstables = 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
 -        UUID pendingSession = UUID.randomUUID();
 -        Set<UUID> sessions = ImmutableSet.of(pendingSession, 
UUID.randomUUID());
++        TimeUUID pendingSession = nextTimeUUID();
++        Set<TimeUUID> sessions = ImmutableSet.of(pendingSession, 
nextTimeUUID());
+         AbstractPendingRepairTest.mutateRepaired(sstables.get(0), 
pendingSession, false);
+ 
+         LifecycleTransaction txn = blockCompactions(cfs, sstables);
+         ClassReloadingStrategy strategy = stubWaitForCessation();
+         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();
+         }
+     }
+ 
+     @Test
+     public void testSubmitMaximalNoOpsWhenCompactionsCannotBeDisabled() 
throws Exception
+     {
+         ColumnFamilyStore cfs = MockSchema.newCFS();
+         List<SSTableReader> sstables = createSSTables(cfs, 3, 0);
+ 
+         LifecycleTransaction txn = blockCompactions(cfs, sstables);
+         ClassReloadingStrategy strategy = stubWaitForCessation();
+         try
+         {
+             // submitMaximal should no-op on null getMaximalTasks result
+             assertTrue(CompactionManager.instance.submitMaximal(cfs, -1, 
false).isEmpty());
+         }
+         finally
+         {
+             strategy.reset(CompactionManager.class);
+             txn.abort();
+         }
+     }
+ 
      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