Repository: cassandra
Updated Branches:
  refs/heads/trunk 59de35332 -> 167ebbcf4


Fail incremental repair prepare phase if it encounters sstables from 
un-finalized sessions

Patch by Blake Eggleston; Reviewed by Marcus Eriksson for CASSANDRA-14763


Project: http://git-wip-us.apache.org/repos/asf/cassandra/repo
Commit: http://git-wip-us.apache.org/repos/asf/cassandra/commit/167ebbcf
Tree: http://git-wip-us.apache.org/repos/asf/cassandra/tree/167ebbcf
Diff: http://git-wip-us.apache.org/repos/asf/cassandra/diff/167ebbcf

Branch: refs/heads/trunk
Commit: 167ebbcf4304512fa538e9cfc18da4295511d16c
Parents: 59de353
Author: Blake Eggleston <[email protected]>
Authored: Tue Sep 18 14:45:06 2018 -0700
Committer: Blake Eggleston <[email protected]>
Committed: Fri Sep 21 10:02:08 2018 -0700

----------------------------------------------------------------------
 CHANGES.txt                                     |  1 +
 .../db/repair/PendingAntiCompaction.java        | 79 ++++++++++++++++++--
 .../apache/cassandra/tools/RepairRunner.java    |  4 +-
 .../db/repair/PendingAntiCompactionTest.java    | 54 ++++++++++---
 .../cassandra/repair/AbstractRepairTest.java    |  4 +-
 5 files changed, 124 insertions(+), 18 deletions(-)
----------------------------------------------------------------------


http://git-wip-us.apache.org/repos/asf/cassandra/blob/167ebbcf/CHANGES.txt
----------------------------------------------------------------------
diff --git a/CHANGES.txt b/CHANGES.txt
index f9d2f3c..a053130 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,4 +1,5 @@
 4.0
+ * Fail incremental repair prepare phase if it encounters sstables from 
un-finalized sessions (CASSANDRA-14763)
  * Add a check for receiving digest response from transient node 
(CASSANDRA-14750)
  * Fail query on transient replica if coordinator only expects full data 
(CASSANDRA-14704)
  * Remove mentions of transient replication from repair path (CASSANDRA-14698)

http://git-wip-us.apache.org/repos/asf/cassandra/blob/167ebbcf/src/java/org/apache/cassandra/db/repair/PendingAntiCompaction.java
----------------------------------------------------------------------
diff --git a/src/java/org/apache/cassandra/db/repair/PendingAntiCompaction.java 
b/src/java/org/apache/cassandra/db/repair/PendingAntiCompaction.java
index a205c3c..2829f0f 100644
--- a/src/java/org/apache/cassandra/db/repair/PendingAntiCompaction.java
+++ b/src/java/org/apache/cassandra/db/repair/PendingAntiCompaction.java
@@ -20,10 +20,13 @@ package org.apache.cassandra.db.repair;
 
 import java.util.ArrayList;
 import java.util.Collection;
+import java.util.HashSet;
 import java.util.List;
+import java.util.Set;
 import java.util.UUID;
 import java.util.concurrent.Callable;
 import java.util.concurrent.ExecutorService;
+import java.util.stream.Collectors;
 
 import com.google.common.collect.Iterables;
 import com.google.common.collect.Lists;
@@ -43,9 +46,14 @@ import 
org.apache.cassandra.db.lifecycle.LifecycleTransaction;
 import org.apache.cassandra.dht.Range;
 import org.apache.cassandra.dht.Token;
 import org.apache.cassandra.io.sstable.format.SSTableReader;
+import org.apache.cassandra.io.sstable.metadata.StatsMetadata;
 import org.apache.cassandra.locator.RangesAtEndpoint;
+import org.apache.cassandra.service.ActiveRepairService;
 import org.apache.cassandra.utils.concurrent.Refs;
 
+import static 
org.apache.cassandra.service.ActiveRepairService.NO_PENDING_REPAIR;
+import static 
org.apache.cassandra.service.ActiveRepairService.UNREPAIRED_SSTABLE;
+
 /**
  * Performs an anti compaction on a set of tables and token ranges, isolating 
the unrepaired sstables
  * for a give token range into a pending repair group so they can't be 
compacted with other sstables
@@ -94,7 +102,47 @@ public class PendingAntiCompaction
 
         private Iterable<SSTableReader> getSSTables()
         {
-            return Iterables.filter(cfs.getLiveSSTables(), s -> 
!s.isRepaired() && !s.isPendingRepair() && s.intersects(ranges));
+            Set<UUID> conflictingSessions = new HashSet<>();
+
+            Iterable<SSTableReader> sstables = 
cfs.getLiveSSTables().stream().filter(sstable -> {
+                if (!sstable.intersects(ranges))
+                    return false;
+
+                StatsMetadata metadata = sstable.getSSTableMetadata();
+
+                // exclude repaired sstables
+                if (metadata.repairedAt != UNREPAIRED_SSTABLE)
+                    return false;
+
+                // exclude sstables pending repair, but record session ids for
+                // non-finalized sessions for a later error message
+                if (metadata.pendingRepair != NO_PENDING_REPAIR)
+                {
+                    if 
(!ActiveRepairService.instance.consistent.local.isSessionFinalized(metadata.pendingRepair))
+                    {
+                        conflictingSessions.add(metadata.pendingRepair);
+                    }
+                    return false;
+                }
+
+                return true;
+            }).collect(Collectors.toList());
+
+            // If there are sstables we'd like to acquire that are currently 
held by other sessions, we need to bail out. If we
+            // didn't bail out here and the other repair sessions we're seeing 
were to fail, incremental repair behavior would be
+            // confusing. You generally expect all data received before a 
repair session to be repaired when the session completes,
+            // and that wouldn't be the case if the other session failed and 
moved it's data back to unrepaired.
+            if (!conflictingSessions.isEmpty())
+            {
+                logger.warn("Prepare phase for incremental repair session {} 
has failed because it encountered " +
+                            "intersecting sstables belonging to another 
incremental repair session(s) ({}). This is " +
+                            "caused by starting an incremental repair session 
before a previous one has completed. " +
+                            "Check nodetool repair_admin for hung sessions and 
fix them.",
+                            sessionID, conflictingSessions);
+                throw new SSTableAcquisitionException();
+            }
+
+            return sstables;
         }
 
         @SuppressWarnings("resource")
@@ -114,9 +162,16 @@ public class PendingAntiCompaction
         public AcquireResult call() throws Exception
         {
             logger.debug("acquiring sstables for pending anti compaction on 
session {}", sessionID);
-            AcquireResult refTxn = acquireTuple();
-            if (refTxn != null)
-                return refTxn;
+            try
+            {
+                AcquireResult refTxn = acquireTuple();
+                if (refTxn != null)
+                    return refTxn;
+            }
+            catch (SSTableAcquisitionException e)
+            {
+                return null;
+            }
 
             // try to modify after cancelling running compactions. This will 
attempt to cancel in flight compactions for
             // up to a minute, after which point, null will be returned
@@ -140,9 +195,23 @@ public class PendingAntiCompaction
             return 
CompactionManager.instance.submitPendingAntiCompaction(result.cfs, tokenRanges, 
result.refs, result.txn, parentRepairSession);
         }
 
+        private static boolean shouldAbort(AcquireResult result)
+        {
+            if (result == null)
+                return true;
+
+            // sstables in the acquire result are now marked compacting and 
are locked to this anti compaction. If any
+            // of them are marked repaired or pending repair, acquisition 
raced with another pending anti-compaction, or
+            // possibly even a repair session, and we need to abort to prevent 
sstables from moving between sessions.
+            return result.refs != null && Iterables.any(result.refs, sstable 
-> {
+                StatsMetadata metadata = sstable.getSSTableMetadata();
+                return metadata.pendingRepair != NO_PENDING_REPAIR || 
metadata.repairedAt != UNREPAIRED_SSTABLE;
+            });
+        }
+
         public ListenableFuture apply(List<AcquireResult> results) throws 
Exception
         {
-            if (Iterables.any(results, t -> t == null))
+            if (Iterables.any(results, AcquisitionCallback::shouldAbort))
             {
                 // Release all sstables, and report failure back to coordinator
                 for (AcquireResult result : results)

http://git-wip-us.apache.org/repos/asf/cassandra/blob/167ebbcf/src/java/org/apache/cassandra/tools/RepairRunner.java
----------------------------------------------------------------------
diff --git a/src/java/org/apache/cassandra/tools/RepairRunner.java 
b/src/java/org/apache/cassandra/tools/RepairRunner.java
index ca894f0..d1b3409 100644
--- a/src/java/org/apache/cassandra/tools/RepairRunner.java
+++ b/src/java/org/apache/cassandra/tools/RepairRunner.java
@@ -119,7 +119,9 @@ public class RepairRunner extends 
JMXNotificationProgressListener
         out.println(message);
         if (type == ProgressEventType.ERROR)
         {
-            error = new RuntimeException("Repair job has failed with the error 
message: " + message);
+            error = new RuntimeException(String.format("Repair job has failed 
with the error message: %s. " +
+                                                       "Check the logs on the 
repair participants for further details",
+                                                       message));
         }
         if (type == ProgressEventType.COMPLETE)
         {

http://git-wip-us.apache.org/repos/asf/cassandra/blob/167ebbcf/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java
----------------------------------------------------------------------
diff --git 
a/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java 
b/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java
index 374a760..ab32e5b 100644
--- a/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java
+++ b/test/unit/org/apache/cassandra/db/repair/PendingAntiCompactionTest.java
@@ -18,7 +18,6 @@
 
 package org.apache.cassandra.db.repair;
 
-import java.net.UnknownHostException;
 import java.util.ArrayList;
 import java.util.Collection;
 import java.util.Collections;
@@ -30,33 +29,35 @@ import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 
 import com.google.common.collect.Lists;
+import com.google.common.collect.Sets;
 import com.google.common.util.concurrent.ListenableFuture;
 import com.google.common.util.concurrent.ListenableFutureTask;
 import org.junit.Assert;
 import org.junit.Before;
 import org.junit.BeforeClass;
 import org.junit.Test;
-
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import org.apache.cassandra.SchemaLoader;
 import org.apache.cassandra.config.DatabaseDescriptor;
-import org.apache.cassandra.db.compaction.CompactionManager;
-import org.apache.cassandra.locator.RangesAtEndpoint;
-import org.apache.cassandra.locator.Replica;
-import org.apache.cassandra.schema.TableId;
-import org.apache.cassandra.schema.TableMetadata;
-import org.apache.cassandra.schema.Schema;
 import org.apache.cassandra.cql3.QueryProcessor;
 import org.apache.cassandra.cql3.statements.schema.CreateTableStatement;
 import org.apache.cassandra.db.ColumnFamilyStore;
+import org.apache.cassandra.db.compaction.CompactionManager;
 import org.apache.cassandra.dht.ByteOrderedPartitioner;
 import org.apache.cassandra.dht.Range;
 import org.apache.cassandra.dht.Token;
 import org.apache.cassandra.io.sstable.format.SSTableReader;
 import org.apache.cassandra.locator.InetAddressAndPort;
+import org.apache.cassandra.locator.RangesAtEndpoint;
+import org.apache.cassandra.locator.Replica;
+import org.apache.cassandra.repair.AbstractRepairTest;
+import org.apache.cassandra.repair.consistent.LocalSessionAccessor;
 import org.apache.cassandra.schema.KeyspaceParams;
+import org.apache.cassandra.schema.Schema;
+import org.apache.cassandra.schema.TableId;
+import org.apache.cassandra.schema.TableMetadata;
 import org.apache.cassandra.service.ActiveRepairService;
 import org.apache.cassandra.streaming.PreviewKind;
 import org.apache.cassandra.utils.ByteBufferUtil;
@@ -87,6 +88,7 @@ public class PendingAntiCompactionTest
     {
         SchemaLoader.prepareServer();
         local = InetAddressAndPort.getByName("127.0.0.1");
+        ActiveRepairService.instance.consistent.local.start();
     }
 
     @Before
@@ -128,6 +130,13 @@ public class PendingAntiCompactionTest
         }
     }
 
+    private UUID prepareSession()
+    {
+        UUID sessionID = AbstractRepairTest.registerSession(cfs, true, true);
+        LocalSessionAccessor.prepareUnsafe(sessionID, 
AbstractRepairTest.COORDINATOR, 
Sets.newHashSet(AbstractRepairTest.COORDINATOR));
+        return sessionID;
+    }
+
     /**
      * verify the pending anti compaction happy path
      */
@@ -239,7 +248,7 @@ public class PendingAntiCompactionTest
     }
 
     @Test
-    public void pendingRepairSSTablesAreNotAcquired() throws Exception
+    public void finalizedPendingRepairSSTablesAreNotAcquired() throws Exception
     {
         cfs.disableAutoCompaction();
         makeSSTables(2);
@@ -251,7 +260,9 @@ public class PendingAntiCompactionTest
         Assert.assertTrue(repaired.intersects(FULL_RANGE));
         Assert.assertTrue(unrepaired.intersects(FULL_RANGE));
 
-        
repaired.descriptor.getMetadataSerializer().mutateRepairMetadata(repaired.descriptor,
 0, UUIDGen.getTimeUUID(), false);
+        UUID sessionId = prepareSession();
+        LocalSessionAccessor.finalizeUnsafe(sessionId);
+        
repaired.descriptor.getMetadataSerializer().mutateRepairMetadata(repaired.descriptor,
 0, sessionId, false);
         repaired.reloadSSTableMetadata();
         Assert.assertTrue(repaired.isPendingRepair());
 
@@ -266,6 +277,29 @@ public class PendingAntiCompactionTest
     }
 
     @Test
+    public void conflictingSessionAcquisitionFailure() throws Exception
+    {
+        cfs.disableAutoCompaction();
+        makeSSTables(2);
+
+        List<SSTableReader> sstables = new ArrayList<>(cfs.getLiveSSTables());
+        Assert.assertEquals(2, sstables.size());
+        SSTableReader repaired = sstables.get(0);
+        SSTableReader unrepaired = sstables.get(1);
+        Assert.assertTrue(repaired.intersects(FULL_RANGE));
+        Assert.assertTrue(unrepaired.intersects(FULL_RANGE));
+
+        UUID sessionId = prepareSession();
+        
repaired.descriptor.getMetadataSerializer().mutateRepairMetadata(repaired.descriptor,
 0, sessionId, false);
+        repaired.reloadSSTableMetadata();
+        Assert.assertTrue(repaired.isPendingRepair());
+
+        PendingAntiCompaction.AcquisitionCallable acquisitionCallable = new 
PendingAntiCompaction.AcquisitionCallable(cfs, FULL_RANGE, 
UUIDGen.getTimeUUID());
+        PendingAntiCompaction.AcquireResult result = 
acquisitionCallable.call();
+        Assert.assertNull(result);
+    }
+
+    @Test
     public void pendingRepairNoSSTablesExist() throws Exception
     {
         cfs.disableAutoCompaction();

http://git-wip-us.apache.org/repos/asf/cassandra/blob/167ebbcf/test/unit/org/apache/cassandra/repair/AbstractRepairTest.java
----------------------------------------------------------------------
diff --git a/test/unit/org/apache/cassandra/repair/AbstractRepairTest.java 
b/test/unit/org/apache/cassandra/repair/AbstractRepairTest.java
index 970e648..2c47137 100644
--- a/test/unit/org/apache/cassandra/repair/AbstractRepairTest.java
+++ b/test/unit/org/apache/cassandra/repair/AbstractRepairTest.java
@@ -41,7 +41,7 @@ import org.apache.cassandra.utils.UUIDGen;
 @Ignore
 public abstract class AbstractRepairTest
 {
-    protected static final InetAddressAndPort COORDINATOR;
+    public static final InetAddressAndPort COORDINATOR;
     protected static final InetAddressAndPort PARTICIPANT1;
     protected static final InetAddressAndPort PARTICIPANT2;
     protected static final InetAddressAndPort PARTICIPANT3;
@@ -77,7 +77,7 @@ public abstract class AbstractRepairTest
 
     protected static final Set<Range<Token>> ALL_RANGES = 
ImmutableSet.of(RANGE1, RANGE2, RANGE3);
 
-    protected static UUID registerSession(ColumnFamilyStore cfs, boolean 
isIncremental, boolean isGlobal)
+    public static UUID registerSession(ColumnFamilyStore cfs, boolean 
isIncremental, boolean isGlobal)
     {
         UUID sessionId = UUIDGen.getTimeUUID();
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to