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]
