maksaska commented on code in PR #13625:
URL: https://github.com/apache/ignite/pull/13625#discussion_r4166803632


##########
modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/dht/MdcTopologySplitAbstractTest.java:
##########
@@ -0,0 +1,607 @@
+/*
+ * 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.ignite.internal.processors.cache.distributed.dht;
+
+import java.net.InetSocketAddress;
+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.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.UUID;
+import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
+import org.apache.ignite.Ignite;
+import org.apache.ignite.IgniteCache;
+import org.apache.ignite.cache.CacheAtomicityMode;
+import org.apache.ignite.cache.CachePeekMode;
+import org.apache.ignite.cache.affinity.rendezvous.MdcAffinityBackupFilter;
+import org.apache.ignite.cache.affinity.rendezvous.RendezvousAffinityFunction;
+import org.apache.ignite.cluster.ClusterNode;
+import org.apache.ignite.configuration.CacheConfiguration;
+import org.apache.ignite.configuration.IgniteConfiguration;
+import org.apache.ignite.internal.IgniteEx;
+import org.apache.ignite.internal.TestRecordingCommunicationSpi;
+import org.apache.ignite.internal.processors.cache.CacheInvalidStateException;
+import 
org.apache.ignite.internal.processors.cache.distributed.dht.preloader.GridDhtPartitionsExchangeFuture;
+import org.apache.ignite.internal.util.typedef.F;
+import org.apache.ignite.internal.util.typedef.G;
+import org.apache.ignite.internal.util.typedef.X;
+import org.apache.ignite.lang.IgniteBiPredicate;
+import org.apache.ignite.plugin.extensions.communication.Message;
+import org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi;
+import org.apache.ignite.spi.discovery.tcp.ipfinder.vm.TcpDiscoveryVmIpFinder;
+import org.apache.ignite.testframework.GridTestUtils;
+import org.apache.ignite.topology.MdcTopologyValidator;
+import org.jetbrains.annotations.Nullable;
+
+import static org.apache.ignite.IgniteSystemProperties.IGNITE_DATA_CENTER_ID;
+import static org.apache.ignite.cache.CacheWriteSynchronizationMode.FULL_SYNC;
+
+/**
+ * Base for tests that cut the network between data centers (DCs) of one 
cluster.
+ *
+ * <p>The cluster has {@link #serversPerDc()} servers and one client in each 
DC of {@link #dataCenters()}.
+ * Servers {@code 0 .. dcs * serversPerDc - 1} go DC by DC; the client of DC 
{@code i} has index
+ * {@code dcs * serversPerDc + i}. A client knows only the servers of its own 
DC, so it stays on its DC's side of
+ * a split.</p>
+ *
+ * <p>{@link #split(String...)} cuts the given DCs off from the rest, and 
{@link #splitInto(List)} splits the DCs into
+ * any number of sides: discovery connections between sides fail, and 
communication messages between them are held.
+ * Either then waits, with a timeout, until every node sees exactly its own 
side. {@link #heal(String...)} drops the
+ * held messages (delivering them would replay one side's view on another) and 
restarts the given DCs, which then
+ * join the cluster again. Subclasses can hold more messages while the cluster 
is split through
+ * {@link #blockMessage(ClusterNode, ClusterNode, Message)}.</p>
+ *
+ * <p>Unlike {@link IgniteCacheTopologySplitAbstractTest#splitAndWait()}, 
nothing here waits for one exact topology
+ * version without a timeout, so a split that ends in an unexpected topology 
fails the test instead of hanging it.</p>
+ */
+public abstract class MdcTopologySplitAbstractTest extends 
IgniteCacheTopologySplitAbstractTest {
+    /** */
+    protected static final String DC1 = "DC1";
+
+    /** */
+    protected static final String DC2 = "DC2";
+
+    /** */
+    protected static final String DC3 = "DC3";
+
+    /** Time for the sides of a split to see only themselves, and for a healed 
cluster to become whole. */
+    protected static final long TOPOLOGY_TIMEOUT = 30_000;
+
+    /** */
+    private static final String LOCAL_IP = "127.0.0.1";
+
+    /** Message of the exception a write gets when the topology validator 
rejects it. */
+    private static final String VALIDATOR_REJECTION = "cache topology is not 
valid";
+
+    /** Side of each DC in the current split, numbered from 0; empty when the 
cluster is whole. */
+    private volatile Map<String, Integer> sideOfDc = Collections.emptyMap();
+
+    /** @return DCs of the cluster, in the order their servers are numbered. */
+    protected abstract List<String> dataCenters();
+
+    /** @return Number of servers in each DC. */
+    protected int serversPerDc() {
+        return 2;
+    }
+
+    /** {@inheritDoc} */
+    @Override protected IgniteConfiguration getConfiguration(String 
igniteInstanceName) throws Exception {
+        IgniteConfiguration cfg = super.getConfiguration(igniteInstanceName);
+
+        int idx = getTestIgniteInstanceIndex(igniteInstanceName);
+
+        String dc = dataCenterOf(idx);
+
+        cfg.setUserAttributes(F.asMap(IGNITE_DATA_CENTER_ID, dc));
+
+        TcpDiscoverySpi disco = new MdcSplitDiscoverySpi(dc);
+
+        
disco.setReconnectCount(((TcpDiscoverySpi)cfg.getDiscoverySpi()).getReconnectCount());
+
+        if (isServer(idx)) {
+            disco.setLocalPort(discoveryPort(idx));
+            disco.setLocalPortRange(0);
+            disco.setIpFinder(new 
TcpDiscoveryVmIpFinder().setAddresses(serverAddresses(dataCenters())));
+        }
+        else
+            disco.setIpFinder(new 
TcpDiscoveryVmIpFinder().setAddresses(serverAddresses(Collections.singleton(dc))));
+
+        cfg.setDiscoverySpi(disco);
+
+        return cfg;
+    }
+
+    /** {@inheritDoc} */
+    @Override protected void afterTest() throws Exception {
+        sideOfDc = Collections.emptyMap();
+
+        stopAllGrids();
+
+        super.afterTest();
+    }
+
+    /**
+     * Starts every server, then a client in each DC.
+     *
+     * @throws Exception If failed.
+     */
+    protected void startCluster() throws Exception {
+        startGridsMultiThreaded(serverCount());
+
+        for (String dc : dataCenters())
+            startClientGrid(clientIndex(dc));
+
+        awaitPartitionMapExchange();
+    }
+
+    /**
+     * Cuts the given DCs off from the rest of the cluster and waits until 
each side sees only itself.
+     *
+     * @param dcs DCs to cut off.
+     */
+    protected void split(String... dcs) throws Exception {
+        List<String> cut = Arrays.asList(dcs);
+
+        List<String> rest = new ArrayList<>(dataCenters());
+
+        rest.removeAll(cut);
+
+        splitInto(Arrays.asList(rest, cut));
+    }
+
+    /**
+     * Splits the DCs into the given sides and waits until each side sees only 
itself.
+     *
+     * @param sides Sides of the split: two at least, together holding every 
DC exactly once.
+     */
+    protected void splitInto(List<List<String>> sides) throws Exception {
+        assertTrue("Cluster is split already: " + sideOfDc, 
sideOfDc.isEmpty());
+        assertTrue("A split needs two sides at least: " + sides, sides.size() 
>= 2);
+
+        Map<String, Integer> split = new HashMap<>();
+
+        for (int side = 0; side < sides.size(); side++) {
+            assertFalse("Every side must have a DC: " + sides, 
sides.get(side).isEmpty());
+
+            for (String dc : sides.get(side)) {
+                assertTrue("Unknown DC: " + dc, dataCenters().contains(dc));
+                assertNull("DC is listed twice: " + dc, split.put(dc, side));
+            }
+        }
+
+        assertEquals("Every DC must be on a side: " + sides, 
dataCenters().size(), split.size());
+
+        log.info(">>> Splitting DCs into " + sides);
+
+        Map<String, Integer> sideMap = Collections.unmodifiableMap(split);
+
+        // Hold communication first, then cut discovery: no message crosses 
the split once any node sees it.
+        for (Ignite ignite : G.allGrids())
+            communication(ignite).blockMessages(new 
MdcSplitBlocker(ignite.cluster().localNode(), sideMap));
+
+        sideOfDc = sideMap;
+
+        long start = System.nanoTime();
+
+        awaitSidesSeeThemselves();
+
+        log.info(">>> Split done in " + 
TimeUnit.NANOSECONDS.toMillis(System.nanoTime() - start) + " ms");
+    }
+
+    /**
+     * Restarts every node of the given DCs and waits until the cluster is 
whole again. Every held message is
+     * dropped and every message filter is removed, including those added by 
{@link #blockMessage}. The given DCs
+     * must be whole sides of the split, every side but one, normally the 
sides that lost writes: they rejoin the
+     * remaining side and rebalance from it. The sides form separate rings, 
which never merge on their own.
+     *
+     * @param restartDcs DCs to restart.
+     */
+    protected void heal(String... restartDcs) throws Exception {
+        Map<String, Integer> split = sideOfDc;
+
+        assertFalse("Cluster is not split", split.isEmpty());
+
+        Set<String> restartSet = new HashSet<>(Arrays.asList(restartDcs));
+
+        Set<Integer> kept = new HashSet<>();
+
+        for (Map.Entry<String, Integer> e : split.entrySet()) {
+            if (!restartSet.contains(e.getKey()))
+                kept.add(e.getValue());
+        }
+
+        boolean wholeSides = split.keySet().containsAll(restartSet)
+            && restartSet.stream().noneMatch(dc -> 
kept.contains(split.get(dc)));
+
+        assertTrue("DCs to restart must be whole sides of the split " + split 
+ ", every side but one: " +
+            restartSet, wholeSides && kept.size() == 1);
+
+        log.info(">>> Healing the split, restarting DCs " + 
Arrays.toString(restartDcs));
+
+        List<Integer> restart = new ArrayList<>();
+
+        for (String dc : dataCenters()) {
+            if (!restartSet.contains(dc))
+                continue;
+
+            restart.addAll(serverIndexes(dc));
+            restart.add(clientIndex(dc));
+        }
+
+        for (int idx : restart)
+            stopGrid(idx, true);
+
+        sideOfDc = Collections.emptyMap();
+
+        for (Ignite ignite : G.allGrids())
+            communication(ignite).stopBlock(false);
+
+        for (int idx : restart) {
+            if (isServer(idx))
+                startGrid(idx);
+        }
+
+        for (int idx : restart) {
+            if (!isServer(idx))
+                startClientGrid(idx);
+        }
+
+        int total = serverCount() + dataCenters().size();
+
+        assertTrue("Cluster did not become whole: " + topologyViews(),
+            GridTestUtils.waitForCondition(() -> G.allGrids().stream()
+                .allMatch(ignite -> ignite.cluster().nodes().size() == total), 
TOPOLOGY_TIMEOUT));
+
+        awaitPartitionMapExchange();
+
+        log.info(">>> Heal done");
+    }
+
+    /**
+     * @param name Cache name.
+     * @param atomicityMode Atomicity mode.
+     * @param validator Topology validator, or {@code null} for none.
+     * @return Cache configuration that keeps one copy of each partition in 
every DC.
+     */
+    protected CacheConfiguration<Integer, Integer> cacheConfiguration(
+        String name,
+        CacheAtomicityMode atomicityMode,
+        @Nullable MdcTopologyValidator validator
+    ) {
+        int dcs = dataCenters().size();
+
+        int backups = dcs - 1;
+
+        return new CacheConfiguration<Integer, Integer>(name)
+            .setAtomicityMode(atomicityMode)
+            .setWriteSynchronizationMode(FULL_SYNC)
+            .setBackups(backups)
+            .setAffinity(new RendezvousAffinityFunction(false, 32)
+                .setAffinityBackupFilter(new MdcAffinityBackupFilter(dcs, 
backups)))
+            .setTopologyValidator(validator);
+    }
+
+    /**
+     * @param dcs Every DC of the cluster.
+     * @return Validator that lets a side write while it sees a majority of 
the DCs.
+     */
+    protected static MdcTopologyValidator majorityValidator(String... dcs) {

Review Comment:
   majorityValidator now fails at once on an even number of DCs and points to 
mainDcValidator(). 



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to