Repository: cassandra
Updated Branches:
  refs/heads/trunk a6f39834b -> cb56d9fc3


http://git-wip-us.apache.org/repos/asf/cassandra/blob/cb56d9fc/test/unit/org/apache/cassandra/repair/RepairSessionTest.java
----------------------------------------------------------------------
diff --git a/test/unit/org/apache/cassandra/repair/RepairSessionTest.java 
b/test/unit/org/apache/cassandra/repair/RepairSessionTest.java
index efae538..984218d 100644
--- a/test/unit/org/apache/cassandra/repair/RepairSessionTest.java
+++ b/test/unit/org/apache/cassandra/repair/RepairSessionTest.java
@@ -66,7 +66,7 @@ public class RepairSessionTest
         RepairSession session = new RepairSession(parentSessionId, sessionId, 
Arrays.asList(repairRange),
                                                   "Keyspace1", 
RepairParallelism.SEQUENTIAL,
                                                   endpoints, false, false, 
false,
-                                                  PreviewKind.NONE, 
"Standard1");
+                                                  PreviewKind.NONE, false, 
"Standard1");
 
         // perform convict
         session.convict(remote, Double.MAX_VALUE);

http://git-wip-us.apache.org/repos/asf/cassandra/blob/cb56d9fc/test/unit/org/apache/cassandra/repair/StreamingRepairTaskTest.java
----------------------------------------------------------------------
diff --git a/test/unit/org/apache/cassandra/repair/StreamingRepairTaskTest.java 
b/test/unit/org/apache/cassandra/repair/StreamingRepairTaskTest.java
index f433f2e..f0ff0e0 100644
--- a/test/unit/org/apache/cassandra/repair/StreamingRepairTaskTest.java
+++ b/test/unit/org/apache/cassandra/repair/StreamingRepairTaskTest.java
@@ -65,8 +65,9 @@ public class StreamingRepairTaskTest extends 
AbstractRepairTest
         UUID sessionID = registerSession(cfs, true, true);
         ActiveRepairService.ParentRepairSession prs = 
ActiveRepairService.instance.getParentRepairSession(sessionID);
         RepairJobDesc desc = new RepairJobDesc(sessionID, 
UUIDGen.getTimeUUID(), ks, tbl, prs.getRanges());
+
         SyncRequest request = new SyncRequest(desc, PARTICIPANT1, 
PARTICIPANT2, PARTICIPANT3, prs.getRanges(), PreviewKind.NONE);
-        StreamingRepairTask task = new StreamingRepairTask(desc, request, 
desc.sessionId, PreviewKind.NONE);
+        StreamingRepairTask task = new StreamingRepairTask(desc, 
request.initiator, request.src, request.dst, request.ranges, desc.sessionId, 
PreviewKind.NONE, false);
 
         StreamPlan plan = task.createStreamPlan(request.src, request.dst);
         Assert.assertFalse(plan.getFlushBeforeTransfer());
@@ -79,7 +80,7 @@ public class StreamingRepairTaskTest extends 
AbstractRepairTest
         ActiveRepairService.ParentRepairSession prs = 
ActiveRepairService.instance.getParentRepairSession(sessionID);
         RepairJobDesc desc = new RepairJobDesc(sessionID, 
UUIDGen.getTimeUUID(), ks, tbl, prs.getRanges());
         SyncRequest request = new SyncRequest(desc, PARTICIPANT1, 
PARTICIPANT2, PARTICIPANT3, prs.getRanges(), PreviewKind.NONE);
-        StreamingRepairTask task = new StreamingRepairTask(desc, request, 
null, PreviewKind.NONE);
+        StreamingRepairTask task = new StreamingRepairTask(desc, 
request.initiator, request.src, request.dst, request.ranges, null, 
PreviewKind.NONE, false);
 
         StreamPlan plan = task.createStreamPlan(request.src, request.dst);
         Assert.assertTrue(plan.getFlushBeforeTransfer());

http://git-wip-us.apache.org/repos/asf/cassandra/blob/cb56d9fc/test/unit/org/apache/cassandra/repair/asymmetric/DifferenceHolderTest.java
----------------------------------------------------------------------
diff --git 
a/test/unit/org/apache/cassandra/repair/asymmetric/DifferenceHolderTest.java 
b/test/unit/org/apache/cassandra/repair/asymmetric/DifferenceHolderTest.java
new file mode 100644
index 0000000..52a43e6
--- /dev/null
+++ b/test/unit/org/apache/cassandra/repair/asymmetric/DifferenceHolderTest.java
@@ -0,0 +1,106 @@
+/*
+ * 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.repair.asymmetric;
+
+import java.net.InetAddress;
+import java.net.UnknownHostException;
+import java.util.Iterator;
+
+import com.google.common.collect.Lists;
+import com.google.common.collect.Sets;
+import org.junit.Test;
+
+import org.apache.cassandra.dht.IPartitioner;
+import org.apache.cassandra.dht.Murmur3Partitioner;
+import org.apache.cassandra.dht.Range;
+import org.apache.cassandra.dht.Token;
+import org.apache.cassandra.repair.TreeResponse;
+import org.apache.cassandra.utils.MerkleTree;
+import org.apache.cassandra.utils.MerkleTrees;
+import org.apache.cassandra.utils.MerkleTreesTest;
+
+import static junit.framework.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+public class DifferenceHolderTest
+{
+    @Test
+    public void testFromEmptyMerkleTrees() throws UnknownHostException
+    {
+        InetAddress a1 = InetAddress.getByName("127.0.0.1");
+        InetAddress a2 = InetAddress.getByName("127.0.0.2");
+
+        MerkleTrees mt1 = new MerkleTrees(Murmur3Partitioner.instance);
+        MerkleTrees mt2 = new MerkleTrees(Murmur3Partitioner.instance);
+        mt1.init();
+        mt2.init();
+
+        TreeResponse tr1 = new TreeResponse(a1, mt1);
+        TreeResponse tr2 = new TreeResponse(a2, mt2);
+
+        DifferenceHolder dh = new DifferenceHolder(Lists.newArrayList(tr1, 
tr2));
+        assertTrue(dh.get(a1).get(a2).isEmpty());
+    }
+
+    @Test
+    public void testFromMismatchedMerkleTrees() throws UnknownHostException
+    {
+        IPartitioner partitioner = Murmur3Partitioner.instance;
+        Range<Token> fullRange = new Range<>(partitioner.getMinimumToken(), 
partitioner.getMinimumToken());
+        int maxsize = 16;
+        InetAddress a1 = InetAddress.getByName("127.0.0.1");
+        InetAddress a2 = InetAddress.getByName("127.0.0.2");
+        // merkle tree building stolen from MerkleTreesTest:
+        MerkleTrees mt1 = new MerkleTrees(partitioner);
+        MerkleTrees mt2 = new MerkleTrees(partitioner);
+        mt1.addMerkleTree(32, fullRange);
+        mt2.addMerkleTree(32, fullRange);
+        mt1.init();
+        mt2.init();
+        // add dummy hashes to both trees
+        for (MerkleTree.TreeRange range : mt1.invalids())
+            range.addAll(new MerkleTreesTest.HIterator(range.right));
+        for (MerkleTree.TreeRange range : mt2.invalids())
+            range.addAll(new MerkleTreesTest.HIterator(range.right));
+
+        MerkleTree.TreeRange leftmost = null;
+        MerkleTree.TreeRange middle = null;
+
+        mt1.maxsize(fullRange, maxsize + 2); // give some room for splitting
+
+        // split the leftmost
+        Iterator<MerkleTree.TreeRange> ranges = mt1.invalids();
+        leftmost = ranges.next();
+        mt1.split(leftmost.right);
+
+        // set the hashes for the leaf of the created split
+        middle = mt1.get(leftmost.right);
+        middle.hash("arbitrary!".getBytes());
+        mt1.get(partitioner.midpoint(leftmost.left, 
leftmost.right)).hash("even more arbitrary!".getBytes());
+
+        TreeResponse tr1 = new TreeResponse(a1, mt1);
+        TreeResponse tr2 = new TreeResponse(a2, mt2);
+
+        DifferenceHolder dh = new DifferenceHolder(Lists.newArrayList(tr1, 
tr2));
+        assertTrue(dh.get(a1).get(a2).size() == 1);
+        assertTrue(dh.hasDifferenceBetween(a1, a2, fullRange));
+        // only a1 is added as a key - see comment in dh.keyHosts()
+        assertEquals(Sets.newHashSet(a1), dh.keyHosts());
+    }
+}

http://git-wip-us.apache.org/repos/asf/cassandra/blob/cb56d9fc/test/unit/org/apache/cassandra/repair/asymmetric/RangeDenormalizerTest.java
----------------------------------------------------------------------
diff --git 
a/test/unit/org/apache/cassandra/repair/asymmetric/RangeDenormalizerTest.java 
b/test/unit/org/apache/cassandra/repair/asymmetric/RangeDenormalizerTest.java
new file mode 100644
index 0000000..a128f2b
--- /dev/null
+++ 
b/test/unit/org/apache/cassandra/repair/asymmetric/RangeDenormalizerTest.java
@@ -0,0 +1,86 @@
+/*
+ * 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.repair.asymmetric;
+
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Set;
+
+import com.google.common.collect.Sets;
+import org.junit.Test;
+
+import org.apache.cassandra.dht.Range;
+import org.apache.cassandra.dht.Token;
+
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+import static org.apache.cassandra.repair.asymmetric.ReduceHelperTest.range;
+
+public class RangeDenormalizerTest
+{
+    @Test
+    public void testDenormalize()
+    {
+        // test when the new incoming range is fully contained within an 
existing incoming range
+        StreamFromOptions dummy = new StreamFromOptions(null, range(0, 100));
+        Map<Range<Token>, StreamFromOptions> incoming = new HashMap<>();
+        incoming.put(range(0, 100), dummy);
+        Set<Range<Token>> newInput = RangeDenormalizer.denormalize(range(30, 
40), incoming);
+        assertEquals(3, incoming.size());
+        assertTrue(incoming.containsKey(range(0, 30)));
+        assertTrue(incoming.containsKey(range(30, 40)));
+        assertTrue(incoming.containsKey(range(40, 100)));
+        assertEquals(1, newInput.size());
+        assertTrue(newInput.contains(range(30, 40)));
+    }
+
+    @Test
+    public void testDenormalize2()
+    {
+        // test when the new incoming range fully contains an existing 
incoming range
+        StreamFromOptions dummy = new StreamFromOptions(null, range(40, 50));
+        Map<Range<Token>, StreamFromOptions> incoming = new HashMap<>();
+        incoming.put(range(40, 50), dummy);
+        Set<Range<Token>> newInput = RangeDenormalizer.denormalize(range(0, 
100), incoming);
+        assertEquals(1, incoming.size());
+        assertTrue(incoming.containsKey(range(40, 50)));
+        assertEquals(3, newInput.size());
+        assertTrue(newInput.contains(range(0, 40)));
+        assertTrue(newInput.contains(range(40, 50)));
+        assertTrue(newInput.contains(range(50, 100)));
+    }
+
+    @Test
+    public void testDenormalize3()
+    {
+        // test when there are multiple existing incoming ranges and the new 
incoming overlaps some and contains some
+        StreamFromOptions dummy = new StreamFromOptions(null, range(0, 100));
+        StreamFromOptions dummy2 = new StreamFromOptions(null, range(200, 
300));
+        StreamFromOptions dummy3 = new StreamFromOptions(null, range(500, 
600));
+        Map<Range<Token>, StreamFromOptions> incoming = new HashMap<>();
+        incoming.put(range(0, 100), dummy);
+        incoming.put(range(200, 300), dummy2);
+        incoming.put(range(500, 600), dummy3);
+        Set<Range<Token>> expectedNewInput = Sets.newHashSet(range(50, 100), 
range(100, 200), range(200, 300), range(300, 350));
+        Set<Range<Token>> expectedIncomingKeys = Sets.newHashSet(range(0, 50), 
range(50, 100), range(200, 300), range(500, 600));
+        Set<Range<Token>> newInput = RangeDenormalizer.denormalize(range(50, 
350), incoming);
+        assertEquals(expectedNewInput, newInput);
+        assertEquals(expectedIncomingKeys, incoming.keySet());
+    }
+}

http://git-wip-us.apache.org/repos/asf/cassandra/blob/cb56d9fc/test/unit/org/apache/cassandra/repair/asymmetric/ReduceHelperTest.java
----------------------------------------------------------------------
diff --git 
a/test/unit/org/apache/cassandra/repair/asymmetric/ReduceHelperTest.java 
b/test/unit/org/apache/cassandra/repair/asymmetric/ReduceHelperTest.java
new file mode 100644
index 0000000..19c42fb
--- /dev/null
+++ b/test/unit/org/apache/cassandra/repair/asymmetric/ReduceHelperTest.java
@@ -0,0 +1,425 @@
+/*
+ * 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.repair.asymmetric;
+
+
+import java.net.InetAddress;
+import java.net.UnknownHostException;
+import java.util.ArrayList;
+import java.util.Arrays;
+import java.util.Collection;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.HashSet;
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+import com.google.common.collect.ImmutableMap;
+import com.google.common.collect.ImmutableSet;
+import com.google.common.collect.Sets;
+import org.junit.Test;
+
+import org.apache.cassandra.dht.Murmur3Partitioner;
+import org.apache.cassandra.dht.Range;
+import org.apache.cassandra.dht.Token;
+
+import static junit.framework.TestCase.fail;
+import static org.junit.Assert.assertEquals;
+import static org.junit.Assert.assertTrue;
+
+public class ReduceHelperTest
+{
+    private static final InetAddress[] addresses;
+    private static final InetAddress A;
+    private static final InetAddress B;
+    private static final InetAddress C;
+    private static final InetAddress D;
+    private static final InetAddress E;
+
+    static
+    {
+        try
+        {
+            A = InetAddress.getByName("127.0.0.0");
+            B = InetAddress.getByName("127.0.0.1");
+            C = InetAddress.getByName("127.0.0.2");
+            D = InetAddress.getByName("127.0.0.3");
+            E = InetAddress.getByName("127.0.0.4");
+            // for diff creation in loops:
+            addresses = new InetAddress[]{ A, B, C, D, E };
+        }
+        catch (UnknownHostException e)
+        {
+            throw new RuntimeException(e);
+        }
+    }
+
+    @Test
+    public void testSimpleReducing()
+    {
+        /*
+        A == B and D == E =>
+        A streams from C, {D, E} since D==E
+        B streams from C, {D, E} since D==E
+        C streams from {A, B}, {D, E} since A==B and D==E
+        D streams from {A, B}, C since A==B
+        E streams from {A, B}, C since A==B
+
+          A   B   C   D   E
+        A     =   x   x   x
+        B         x   x   x
+        C             x   x
+        D                 =
+         */
+        Map<InetAddress, HostDifferences> differences = new HashMap<>();
+        for (int i = 0; i < 4; i++)
+        {
+            HostDifferences hostDiffs = new HostDifferences();
+            for (int j = i + 1; j < 5; j++)
+            {
+                // no diffs between A, B and D, E:
+                if (addresses[i] == A && addresses[j] == B || addresses[i] == 
D && addresses[j] == E)
+                    continue;
+                List<Range<Token>> diff = list(new Range<>(new 
Murmur3Partitioner.LongToken(0), new Murmur3Partitioner.LongToken(10)));
+                hostDiffs.add(addresses[j], diff);
+            }
+            differences.put(addresses[i], hostDiffs);
+
+        }
+        DifferenceHolder differenceHolder = new DifferenceHolder(differences);
+        Map<InetAddress, IncomingRepairStreamTracker> tracker = 
ReduceHelper.createIncomingRepairStreamTrackers(differenceHolder);
+
+        assertEquals(set(set(C), set(E,D)), streams(tracker.get(A)));
+        assertEquals(set(set(C), set(E,D)), streams(tracker.get(B)));
+        assertEquals(set(set(A,B), set(E,D)), streams(tracker.get(C)));
+        assertEquals(set(set(A,B), set(C)), streams(tracker.get(D)));
+        assertEquals(set(set(A,B), set(C)), streams(tracker.get(E)));
+
+        ImmutableMap<InetAddress, HostDifferences> reduced = 
ReduceHelper.reduce(differenceHolder, (x,y) -> y);
+
+        HostDifferences n0 = reduced.get(A);
+        assertEquals(0, n0.get(A).size());
+        assertEquals(0, n0.get(B).size());
+        assertTrue(n0.get(C).size() > 0);
+        assertStreamFromEither(n0.get(D), n0.get(E));
+
+        HostDifferences n1 = reduced.get(B);
+        assertEquals(0, n1.get(A).size());
+        assertEquals(0, n1.get(B).size());
+        assertTrue(n1.get(C).size() > 0);
+        assertStreamFromEither(n1.get(D), n1.get(E));
+
+        HostDifferences n2 = reduced.get(C);
+        // we are either streaming from node 0 or node 1, not both:
+        assertStreamFromEither(n2.get(A), n2.get(B));
+        assertEquals(0, n2.get(C).size());
+        assertStreamFromEither(n2.get(D), n2.get(E));
+
+        HostDifferences n3 = reduced.get(D);
+        assertStreamFromEither(n3.get(A), n3.get(B));
+        assertTrue(n3.get(C).size() > 0);
+        assertEquals(0, n3.get(D).size());
+        assertEquals(0, n3.get(E).size());
+
+        HostDifferences n4 = reduced.get(E);
+        assertStreamFromEither(n4.get(A), n4.get(B));
+        assertTrue(n4.get(C).size() > 0);
+        assertEquals(0, n4.get(D).size());
+        assertEquals(0, n4.get(E).size());
+    }
+
+    @Test
+    public void testSimpleReducingWithPreferedNodes()
+    {
+        /*
+        A == B and D == E =>
+        A streams from C, {D, E} since D==E
+        B streams from C, {D, E} since D==E
+        C streams from {A, B}, {D, E} since A==B and D==E
+        D streams from {A, B}, C since A==B
+        E streams from {A, B}, C since A==B
+
+          A   B   C   D   E
+        A     =   x   x   x
+        B         x   x   x
+        C             x   x
+        D                 =
+         */
+        Map<InetAddress, HostDifferences> differences = new HashMap<>();
+        for (int i = 0; i < 4; i++)
+        {
+            HostDifferences hostDifferences = new HostDifferences();
+            for (int j = i + 1; j < 5; j++)
+            {
+                // no diffs between A, B and D, E:
+                if (addresses[i] == A && addresses[j] == B || addresses[i] == 
D && addresses[j] == E)
+                    continue;
+                List<Range<Token>> diff = list(new Range<>(new 
Murmur3Partitioner.LongToken(0), new Murmur3Partitioner.LongToken(10)));
+                hostDifferences.add(addresses[j], diff);
+            }
+            differences.put(addresses[i], hostDifferences);
+        }
+
+        DifferenceHolder differenceHolder = new DifferenceHolder(differences);
+        Map<InetAddress, IncomingRepairStreamTracker> tracker = 
ReduceHelper.createIncomingRepairStreamTrackers(differenceHolder);
+        assertEquals(set(set(C), set(E, D)), streams(tracker.get(A)));
+        assertEquals(set(set(C), set(E, D)), streams(tracker.get(B)));
+        assertEquals(set(set(A, B), set(E, D)), streams(tracker.get(C)));
+        assertEquals(set(set(A, B), set(C)), streams(tracker.get(D)));
+        assertEquals(set(set(A, B), set(C)), streams(tracker.get(E)));
+
+        // if there is an option, never stream from node 1:
+        ImmutableMap<InetAddress, HostDifferences> reduced = 
ReduceHelper.reduce(differenceHolder, (x,y) -> Sets.difference(y, set(B)));
+
+        HostDifferences n0 = reduced.get(A);
+        assertEquals(0, n0.get(A).size());
+        assertEquals(0, n0.get(B).size());
+        assertTrue(n0.get(C).size() > 0);
+        assertStreamFromEither(n0.get(D), n0.get(E));
+
+        HostDifferences n1 = reduced.get(B);
+        assertEquals(0, n1.get(A).size());
+        assertEquals(0, n1.get(B).size());
+        assertTrue(n1.get(C).size() > 0);
+        assertStreamFromEither(n1.get(D), n1.get(E));
+
+
+        HostDifferences n2 = reduced.get(C);
+        assertTrue(n2.get(A).size() > 0);
+        assertEquals(0, n2.get(B).size());
+        assertEquals(0, n2.get(C).size());
+        assertStreamFromEither(n2.get(D), n2.get(E));
+
+        HostDifferences n3 = reduced.get(D);
+        assertTrue(n3.get(A).size() > 0);
+        assertEquals(0, n3.get(B).size());
+        assertTrue(n3.get(C).size() > 0);
+        assertEquals(0, n3.get(D).size());
+        assertEquals(0, n3.get(E).size());
+
+        HostDifferences n4 = reduced.get(E);
+        assertTrue(n4.get(A).size() > 0);
+        assertEquals(0, n4.get(B).size());
+        assertTrue(n4.get(C).size() > 0);
+        assertEquals(0, n4.get(D).size());
+        assertEquals(0, n4.get(E).size());
+    }
+
+    private Iterable<Set<InetAddress>> streams(IncomingRepairStreamTracker 
incomingRepairStreamTracker)
+    {
+        return 
incomingRepairStreamTracker.getIncoming().values().iterator().next().allStreams();
+    }
+
+    @Test
+    public void testOverlapDifference()
+    {
+        /*
+            |A     |B     |C
+         ---+------+------+--------
+         A  |=     |50,100|0,50
+         B  |      |=     |0,100
+         C  |      |      |=
+
+         A needs to stream (50, 100] from B, (0, 50] from C
+         B needs to stream (50, 100] from A, (0, 100] from C
+         C needs to stream (0, 50] from A, (0, 100] from B
+         A == B on (0, 50]   => C can stream (0, 50] from either A or B
+         A == C on (50, 100] => B can stream (50, 100] from either A or C
+         =>
+         A streams (50, 100] from {B}, (0, 50] from C
+         B streams (0, 50] from {C}, (50, 100] from {A, C}
+         C streams (0, 50] from {A, B}, (50, 100] from B
+         */
+        Map<InetAddress, HostDifferences> differences = new HashMap<>();
+        addDifference(A, differences, B, list(range(50, 100)));
+        addDifference(A, differences, C, list(range(0, 50)));
+        addDifference(B, differences, C, list(range(0, 100)));
+        DifferenceHolder differenceHolder = new DifferenceHolder(differences);
+        Map<InetAddress, IncomingRepairStreamTracker> tracker = 
ReduceHelper.createIncomingRepairStreamTrackers(differenceHolder);
+        assertEquals(set(set(C)), tracker.get(A).getIncoming().get(range(0, 
50)).allStreams());
+        assertEquals(set(set(B)), tracker.get(A).getIncoming().get(range(50, 
100)).allStreams());
+        assertEquals(set(set(C)), tracker.get(B).getIncoming().get(range(0, 
50)).allStreams());
+        assertEquals(set(set(A,C)), tracker.get(B).getIncoming().get(range(50, 
100)).allStreams());
+        assertEquals(set(set(A,B)), tracker.get(C).getIncoming().get(range(0, 
50)).allStreams());
+        assertEquals(set(set(B)), tracker.get(C).getIncoming().get(range(50, 
100)).allStreams());
+
+        ImmutableMap<InetAddress, HostDifferences> reduced = 
ReduceHelper.reduce(differenceHolder, (x, y) -> y);
+
+        HostDifferences n0 = reduced.get(A);
+
+        assertTrue(n0.get(B).equals(list(range(50, 100))));
+        assertTrue(n0.get(C).equals(list(range(0, 50))));
+
+        HostDifferences n1 = reduced.get(B);
+        assertEquals(0, n1.get(B).size());
+        if (n1.get(A) != null)
+        {
+            assertTrue(n1.get(C).equals(list(range(0, 50))));
+            assertTrue(n1.get(A).equals(list(range(50, 100))));
+        }
+        else
+        {
+            assertTrue(n1.get(C).equals(list(range(0, 50), range(50, 100))));
+        }
+        HostDifferences n2 = reduced.get(C);
+        assertEquals(0, n2.get(C).size());
+        if (n2.get(A) != null)
+        {
+            assertTrue(n2.get(A).equals(list(range(0,50))));
+            assertTrue(n2.get(B).equals(list(range(50, 100))));
+        }
+        else
+        {
+            assertTrue(n2.get(A).equals(list(range(0, 50), range(50, 100))));
+        }
+
+
+    }
+
+    @Test
+    public void testOverlapDifference2()
+    {
+        /*
+            |A               |B               |C
+         ---+----------------+----------------+------------------
+         A  |=               |5,45            |0,10 40,50
+         B  |                |=               |0,5 10,40 45,50
+         C  |                |                |=
+
+         A needs to stream (5, 45] from B, (0, 10], (40, 50) from C
+         B needs to stream (5, 45] from A, (0, 5], (10, 40], (45, 50] from C
+         C needs to stream (0, 10], (40,50] from A, (0,5], (10,40], (45,50] 
from B
+         A == B on (0, 5], (45, 50]
+         A == C on (10, 40]
+         B == C on (5, 10], (40, 45]
+         */
+
+        Map<InetAddress, HostDifferences> differences = new HashMap<>();
+        addDifference(A, differences, B, list(range(5, 45)));
+        addDifference(A, differences, C, list(range(0, 10), range(40,50)));
+        addDifference(B, differences, C, list(range(0, 5), range(10,40), 
range(45,50)));
+
+        DifferenceHolder differenceHolder = new DifferenceHolder(differences);
+        Map<InetAddress, IncomingRepairStreamTracker> tracker = 
ReduceHelper.createIncomingRepairStreamTrackers(differenceHolder);
+
+        Map<Range<Token>, StreamFromOptions> ranges = 
tracker.get(A).getIncoming();
+        assertEquals(5, ranges.size());
+
+        assertEquals(set(set(C)), ranges.get(range(0, 5)).allStreams());
+        assertEquals(set(set(B, C)), ranges.get(range(5, 10)).allStreams());
+        assertEquals(set(set(B)), ranges.get(range(10, 40)).allStreams());
+        assertEquals(set(set(B, C)), ranges.get(range(40, 45)).allStreams());
+        assertEquals(set(set(C)), ranges.get(range(45, 50)).allStreams());
+
+        ranges = tracker.get(B).getIncoming();
+        assertEquals(5, ranges.size());
+        assertEquals(set(set(C)), ranges.get(range(0, 5)).allStreams());
+        assertEquals(set(set(A)), ranges.get(range(5, 10)).allStreams());
+        assertEquals(set(set(A, C)), ranges.get(range(10, 40)).allStreams());
+        assertEquals(set(set(A)), ranges.get(range(40, 45)).allStreams());
+        assertEquals(set(set(C)), ranges.get(range(45, 50)).allStreams());
+
+        ranges = tracker.get(C).getIncoming();
+        assertEquals(5, ranges.size());
+        assertEquals(set(set(A, B)), ranges.get(range(0, 5)).allStreams());
+        assertEquals(set(set(A)), ranges.get(range(5, 10)).allStreams());
+        assertEquals(set(set(B)), ranges.get(range(10, 40)).allStreams());
+        assertEquals(set(set(A)), ranges.get(range(40, 45)).allStreams());
+        assertEquals(set(set(A,B)), ranges.get(range(45, 50)).allStreams());
+        ImmutableMap<InetAddress, HostDifferences> reduced = 
ReduceHelper.reduce(differenceHolder, (x, y) -> y);
+
+        assertNoOverlap(A, reduced.get(A), list(range(0, 50)));
+        assertNoOverlap(B, reduced.get(B), list(range(0, 50)));
+        assertNoOverlap(C, reduced.get(C), list(range(0, 50)));
+    }
+
+    private void assertNoOverlap(InetAddress incomingNode, HostDifferences 
node, List<Range<Token>> expectedAfterNormalize)
+    {
+        Set<Range<Token>> allRanges = new HashSet<>();
+        Set<InetAddress> remoteNodes = Sets.newHashSet(A,B,C);
+        remoteNodes.remove(incomingNode);
+        Iterator<InetAddress> iter = remoteNodes.iterator();
+        allRanges.addAll(node.get(iter.next()));
+        InetAddress i = iter.next();
+        for (Range<Token> r : node.get(i))
+        {
+            for (Range<Token> existing : allRanges)
+                if (r.intersects(existing))
+                    fail();
+        }
+        allRanges.addAll(node.get(i));
+        List<Range<Token>> normalized = Range.normalize(allRanges);
+        assertEquals(expectedAfterNormalize, normalized);
+    }
+
+    @SafeVarargs
+    private static List<Range<Token>> list(Range<Token> r, Range<Token> ... rs)
+    {
+        List<Range<Token>> ranges = new ArrayList<>();
+        ranges.add(r);
+        Collections.addAll(ranges, rs);
+        return ranges;
+    }
+
+    private static Set<InetAddress> set(InetAddress ... elem)
+    {
+        return Sets.newHashSet(elem);
+    }
+    @SafeVarargs
+    private static Set<Set<InetAddress>> set(Set<InetAddress> ... elem)
+    {
+        Set<Set<InetAddress>> ret = Sets.newHashSet();
+        ret.addAll(Arrays.asList(elem));
+        return ret;
+    }
+
+    static Murmur3Partitioner.LongToken longtok(long l)
+    {
+        return new Murmur3Partitioner.LongToken(l);
+    }
+
+    static Range<Token> range(long t, long t2)
+    {
+        return new Range<>(longtok(t), longtok(t2));
+    }
+
+    @Test
+    public void testSubtractAllRanges()
+    {
+        Set<Range<Token>> ranges = new HashSet<>();
+        ranges.add(range(10, 20)); ranges.add(range(40, 60));
+        assertEquals(0, RangeDenormalizer.subtractFromAllRanges(ranges, 
range(0, 100)).size());
+        ranges.add(range(90, 110));
+        assertEquals(Sets.newHashSet(range(100, 110)), 
RangeDenormalizer.subtractFromAllRanges(ranges, range(0, 100)));
+        ranges.add(range(-10, 10));
+        assertEquals(Sets.newHashSet(range(-10, 0), range(100, 110)), 
RangeDenormalizer.subtractFromAllRanges(ranges, range(0, 100)));
+    }
+
+    private void assertStreamFromEither(List<Range<Token>> r1, 
List<Range<Token>> r2)
+    {
+        assertTrue(r1.size() > 0 ^ r2.size() > 0);
+    }
+
+    private void addDifference(InetAddress host1, Map<InetAddress, 
HostDifferences> differences, InetAddress host2, List<Range<Token>> ranges)
+    {
+        differences.computeIfAbsent(host1, (x) -> new 
HostDifferences()).add(host2, ranges);
+    }
+}

http://git-wip-us.apache.org/repos/asf/cassandra/blob/cb56d9fc/test/unit/org/apache/cassandra/repair/asymmetric/StreamFromOptionsTest.java
----------------------------------------------------------------------
diff --git 
a/test/unit/org/apache/cassandra/repair/asymmetric/StreamFromOptionsTest.java 
b/test/unit/org/apache/cassandra/repair/asymmetric/StreamFromOptionsTest.java
new file mode 100644
index 0000000..3ba3cfe
--- /dev/null
+++ 
b/test/unit/org/apache/cassandra/repair/asymmetric/StreamFromOptionsTest.java
@@ -0,0 +1,124 @@
+/*
+ * 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.repair.asymmetric;
+
+import java.net.InetAddress;
+import java.net.UnknownHostException;
+import java.util.Collections;
+import java.util.HashSet;
+import java.util.Set;
+
+import com.google.common.collect.Iterables;
+import org.junit.Test;
+
+import org.apache.cassandra.dht.Murmur3Partitioner;
+import org.apache.cassandra.dht.Range;
+import org.apache.cassandra.dht.Token;
+
+import static junit.framework.Assert.assertEquals;
+import static org.junit.Assert.assertFalse;
+import static org.junit.Assert.assertTrue;
+
+public class StreamFromOptionsTest
+{
+    @Test
+    public void addAllDiffingTest() throws UnknownHostException
+    {
+        StreamFromOptions sfo = new StreamFromOptions(new MockDiffs(true), 
range(0, 10));
+        Set<InetAddress> toAdd = new HashSet<>();
+        toAdd.add(InetAddress.getByName("127.0.0.1"));
+        toAdd.add(InetAddress.getByName("127.0.0.2"));
+        toAdd.add(InetAddress.getByName("127.0.0.3"));
+        toAdd.forEach(sfo::add);
+
+        // if all added have differences, each set will contain a single host
+        assertEquals(3, Iterables.size(sfo.allStreams()));
+        Set<InetAddress> allStreams = new HashSet<>();
+        for (Set<InetAddress> streams : sfo.allStreams())
+        {
+            assertEquals(1, streams.size());
+            allStreams.addAll(streams);
+        }
+        assertEquals(toAdd, allStreams);
+    }
+
+    @Test
+    public void addAllMatchingTest() throws UnknownHostException
+    {
+        StreamFromOptions sfo = new StreamFromOptions(new MockDiffs(false), 
range(0, 10));
+        Set<InetAddress> toAdd = new HashSet<>();
+        toAdd.add(InetAddress.getByName("127.0.0.1"));
+        toAdd.add(InetAddress.getByName("127.0.0.2"));
+        toAdd.add(InetAddress.getByName("127.0.0.3"));
+        toAdd.forEach(sfo::add);
+
+        // if all added match, the set will contain all hosts
+        assertEquals(1, Iterables.size(sfo.allStreams()));
+        assertEquals(toAdd, sfo.allStreams().iterator().next());
+    }
+
+    @Test
+    public void splitTest() throws UnknownHostException
+    {
+        splitTestHelper(true);
+        splitTestHelper(false);
+    }
+
+    private void splitTestHelper(boolean diffing) throws UnknownHostException
+    {
+        StreamFromOptions sfo = new StreamFromOptions(new MockDiffs(diffing), 
range(0, 10));
+        Set<InetAddress> toAdd = new HashSet<>();
+        toAdd.add(InetAddress.getByName("127.0.0.1"));
+        toAdd.add(InetAddress.getByName("127.0.0.2"));
+        toAdd.add(InetAddress.getByName("127.0.0.3"));
+        toAdd.forEach(sfo::add);
+        StreamFromOptions sfo1 = sfo.copy(range(0, 5));
+        StreamFromOptions sfo2 = sfo.copy(range(5, 10));
+        assertEquals(range(0, 10), sfo.range);
+        assertEquals(range(0, 5), sfo1.range);
+        assertEquals(range(5, 10), sfo2.range);
+        assertTrue(Iterables.elementsEqual(sfo1.allStreams(), 
sfo2.allStreams()));
+        // verify the backing set is not shared between the copies:
+        sfo1.add(InetAddress.getByName("127.0.0.4"));
+        sfo2.add(InetAddress.getByName("127.0.0.5"));
+        assertFalse(Iterables.elementsEqual(sfo1.allStreams(), 
sfo2.allStreams()));
+    }
+
+    private Range<Token> range(long left, long right)
+    {
+        return new Range<>(new Murmur3Partitioner.LongToken(left), new 
Murmur3Partitioner.LongToken(right));
+    }
+
+    private static class MockDiffs extends DifferenceHolder
+    {
+        private final boolean hasDifference;
+
+        public MockDiffs(boolean hasDifference)
+        {
+            super(Collections.emptyMap());
+            this.hasDifference = hasDifference;
+        }
+
+        @Override
+        public boolean hasDifferenceBetween(InetAddress node1, InetAddress 
node2, Range<Token> range)
+        {
+            return hasDifference;
+        }
+    }
+}

http://git-wip-us.apache.org/repos/asf/cassandra/blob/cb56d9fc/test/unit/org/apache/cassandra/utils/MerkleTreesTest.java
----------------------------------------------------------------------
diff --git a/test/unit/org/apache/cassandra/utils/MerkleTreesTest.java 
b/test/unit/org/apache/cassandra/utils/MerkleTreesTest.java
index 8d7284c..b40f6c4 100644
--- a/test/unit/org/apache/cassandra/utils/MerkleTreesTest.java
+++ b/test/unit/org/apache/cassandra/utils/MerkleTreesTest.java
@@ -514,7 +514,7 @@ public class MerkleTreesTest
         return hstack.pop();
     }
 
-    static class HIterator extends AbstractIterator<RowHash>
+    public static class HIterator extends AbstractIterator<RowHash>
     {
         private Iterator<Token> tokens;
 


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

Reply via email to