diff --git a/src/java/org/apache/cassandra/repair/consistent/LocalSessions.java b/src/java/org/apache/cassandra/repair/consistent/LocalSessions.java index 5fcfe986c5..32ce747310 100644 --- a/src/java/org/apache/cassandra/repair/consistent/LocalSessions.java +++ b/src/java/org/apache/cassandra/repair/consistent/LocalSessions.java @@ -936,7 +936,7 @@ public class LocalSessions } @VisibleForTesting - protected void sessionCompleted(LocalSession session) + public void sessionCompleted(LocalSession session) { for (TableId tid: session.tableIds) { diff --git a/test/unit/org/apache/cassandra/db/compaction/AbstractPendingRepairTest.java b/test/unit/org/apache/cassandra/db/compaction/AbstractPendingRepairTest.java index b0bb23cc64..834ea3e264 100644 --- a/test/unit/org/apache/cassandra/db/compaction/AbstractPendingRepairTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/AbstractPendingRepairTest.java @@ -20,6 +20,7 @@ package org.apache.cassandra.db.compaction; import java.io.IOException; import java.util.HashSet; +import java.util.List; import java.util.Set; import org.junit.Before; @@ -122,4 +123,10 @@ public class AbstractPendingRepairTest extends AbstractRepairTest { mutateRepaired(sstable, ActiveRepairService.UNREPAIRED_SSTABLE, pendingRepair, isTransient); } + + public static void mutateRepaired(List sstables, TimeUUID pendingRepair, boolean isTransient) + { + for (SSTableReader sstable : sstables) + mutateRepaired(sstable, ActiveRepairService.UNREPAIRED_SSTABLE, pendingRepair, isTransient); + } } diff --git a/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerPendingRepairTest.java b/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerPendingRepairTest.java index 76bbd88e3f..45011f1e47 100644 --- a/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerPendingRepairTest.java +++ b/test/unit/org/apache/cassandra/db/compaction/CompactionStrategyManagerPendingRepairTest.java @@ -18,6 +18,7 @@ package org.apache.cassandra.db.compaction; +import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -25,12 +26,14 @@ import com.google.common.collect.Iterables; import org.junit.Assert; import org.junit.Test; +import org.apache.cassandra.io.sstable.Component; import org.apache.cassandra.io.sstable.format.SSTableReader; import org.apache.cassandra.notifications.SSTableAddedNotification; import org.apache.cassandra.notifications.SSTableDeletingNotification; import org.apache.cassandra.notifications.SSTableListChangedNotification; import org.apache.cassandra.notifications.SSTableRepairStatusChanged; import org.apache.cassandra.repair.NoSuchRepairSessionException; +import org.apache.cassandra.repair.consistent.LocalSession; import org.apache.cassandra.repair.consistent.LocalSessionAccessor; import org.apache.cassandra.service.ActiveRepairService; import org.apache.cassandra.utils.FBUtilities; @@ -307,6 +310,96 @@ public class CompactionStrategyManagerPendingRepairTest extends AbstractPendingR Assert.assertEquals(ActiveRepairService.UNREPAIRED_SSTABLE, sstable.getSSTableMetadata().repairedAt); } + /** + * Tests that finalized repairs racing with compactions on the same set of sstables don't leave unrepaired sstables behind + * + * This test checks that when a repair has been finalized but there are still pending sstables a finalize repair + * compaction task is issued for that repair session. + */ + @Test + public void testFinalizedAndCompactionRace() throws NoSuchRepairSessionException + { + TimeUUID repairID = registerSession(cfs, true, true); + LocalSessionAccessor.prepareUnsafe(repairID, COORDINATOR, PARTICIPANTS); + + int numberOfSStables = 4; + List sstables = new ArrayList<>(numberOfSStables); + for (int i = 0; i < numberOfSStables; i++) + { + SSTableReader sstable = makeSSTable(true); + sstables.add(sstable); + Assert.assertFalse(sstable.isRepaired()); + Assert.assertFalse(sstable.isPendingRepair()); + } + + // change to pending repair + mutateRepaired(sstables, repairID, false); + csm.handleNotification(new SSTableAddedNotification(sstables, null), cfs.getTracker()); + for (SSTableReader sstable : sstables) + { + Assert.assertFalse(sstable.isRepaired()); + Assert.assertTrue(sstable.isPendingRepair()); + Assert.assertEquals(repairID, sstable.getPendingRepair()); + } + + // Get a compaction taks based on the sstables marked as pending repair + cfs.getCompactionStrategyManager().enable(); + for (SSTableReader sstable : sstables) + pendingContains(sstable); + AbstractCompactionTask compactionTask = csm.getNextBackgroundTask(FBUtilities.nowInSeconds()); + Assert.assertNotNull(compactionTask); + + // Finalize the repair session + LocalSessionAccessor.finalizeUnsafe(repairID); + LocalSession session = ARS.consistent.local.getSession(repairID); + ARS.consistent.local.sessionCompleted(session); + Assert.assertTrue(hasPendingStrategiesFor(repairID)); + + // run the compaction + compactionTask.execute(ActiveCompactionsTracker.NOOP); + + // The repair session is finalized but there is an sstable left behind pending repair! + SSTableReader compactedSSTable = cfs.getLiveSSTables().iterator().next(); + Assert.assertEquals(repairID, compactedSSTable.getPendingRepair()); + Assert.assertEquals(1, cfs.getLiveSSTables().size()); + + System.out.println("*********************************************************************************************"); + System.out.println(compactedSSTable); + System.out.println("Pending repair UUID: " + compactedSSTable.getPendingRepair()); + System.out.println("Repaired at: " + compactedSSTable.getRepairedAt()); + System.out.println("Creation time: " + compactedSSTable.getCreationTimeFor(Component.DATA)); + System.out.println("Live sstables: " + cfs.getLiveSSTables().size()); + System.out.println("*********************************************************************************************"); + + // Run compaction again. It should pick up the pending repair sstable + compactionTask = csm.getNextBackgroundTask(FBUtilities.nowInSeconds()); + if (compactionTask != null) + { + Assert.assertSame(PendingRepairManager.RepairFinishedCompactionTask.class, compactionTask.getClass()); + compactionTask.execute(ActiveCompactionsTracker.NOOP); + } + + System.out.println("*********************************************************************************************"); + System.out.println(compactedSSTable); + System.out.println("Pending repair UUID: " + compactedSSTable.getPendingRepair()); + System.out.println("Repaired at: " + compactedSSTable.getRepairedAt()); + System.out.println("Creation time: " + compactedSSTable.getCreationTimeFor(Component.DATA)); + System.out.println("Live sstables: " + cfs.getLiveSSTables().size()); + System.out.println("*********************************************************************************************"); + + Assert.assertEquals(1, cfs.getLiveSSTables().size()); + Assert.assertFalse(hasPendingStrategiesFor(repairID)); + Assert.assertFalse(hasTransientStrategiesFor(repairID)); + Assert.assertTrue(repairedContains(compactedSSTable)); + Assert.assertFalse(unrepairedContains(compactedSSTable)); + Assert.assertFalse(pendingContains(compactedSSTable)); + // sstable should have pendingRepair cleared, and repairedAt set correctly + long expectedRepairedAt = ActiveRepairService.instance.getParentRepairSession(repairID).repairedAt; + Assert.assertFalse(compactedSSTable.isPendingRepair()); + Assert.assertTrue(compactedSSTable.isRepaired()); + Assert.assertEquals(expectedRepairedAt, compactedSSTable.getSSTableMetadata().repairedAt); + } + @Test public void finalizedSessionTransientCleanup() { diff --git a/test/unit/org/apache/cassandra/repair/consistent/LocalSessionTest.java b/test/unit/org/apache/cassandra/repair/consistent/LocalSessionTest.java index c83335bfbd..745edfec52 100644 --- a/test/unit/org/apache/cassandra/repair/consistent/LocalSessionTest.java +++ b/test/unit/org/apache/cassandra/repair/consistent/LocalSessionTest.java @@ -196,7 +196,7 @@ public class LocalSessionTest extends AbstractRepairTest public Map completedSessions = new HashMap<>(); - protected void sessionCompleted(LocalSession session) + public void sessionCompleted(LocalSession session) { TimeUUID sessionID = session.sessionID; int calls = completedSessions.getOrDefault(sessionID, 0);