This is an automated email from the ASF dual-hosted git repository.

rpuch pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/ignite-3.git


The following commit(s) were added to refs/heads/main by this push:
     new b7782257f22 IGNITE-25678 Fix 
laggingSchemasPreventPartitionDataReplication (#6045)
b7782257f22 is described below

commit b7782257f2299884b45fdea6f2b29edb3b49b2be
Author: Roman Puchkovskiy <[email protected]>
AuthorDate: Mon Jun 16 11:38:58 2025 +0400

    IGNITE-25678 Fix laggingSchemasPreventPartitionDataReplication (#6045)
---
 .../raftsnapshot/ItTableRaftSnapshotsTest.java     | 28 +---------------
 .../schemasync/ItSchemaSyncAndReplicationTest.java |  6 +++-
 .../java/org/apache/ignite/internal/Cluster.java   | 39 ++++++++++++++++++++--
 3 files changed, 43 insertions(+), 30 deletions(-)

diff --git 
a/modules/raft/src/integrationTest/java/org/apache/ignite/internal/raftsnapshot/ItTableRaftSnapshotsTest.java
 
b/modules/raft/src/integrationTest/java/org/apache/ignite/internal/raftsnapshot/ItTableRaftSnapshotsTest.java
index ba2eeb16bce..b6702745ca2 100644
--- 
a/modules/raft/src/integrationTest/java/org/apache/ignite/internal/raftsnapshot/ItTableRaftSnapshotsTest.java
+++ 
b/modules/raft/src/integrationTest/java/org/apache/ignite/internal/raftsnapshot/ItTableRaftSnapshotsTest.java
@@ -59,7 +59,6 @@ import java.util.stream.Stream;
 import org.apache.ignite.Ignite;
 import org.apache.ignite.internal.Cluster;
 import org.apache.ignite.internal.ClusterPerTestIntegrationTest;
-import org.apache.ignite.internal.TestWrappers;
 import org.apache.ignite.internal.app.IgniteImpl;
 import org.apache.ignite.internal.cluster.management.CmgGroupId;
 import org.apache.ignite.internal.logger.IgniteLogger;
@@ -80,7 +79,6 @@ import 
org.apache.ignite.internal.storage.pagememory.PersistentPageMemoryStorage
 import 
org.apache.ignite.internal.storage.pagememory.VolatilePageMemoryStorageEngine;
 import org.apache.ignite.internal.storage.rocksdb.RocksDbStorageEngine;
 import org.apache.ignite.internal.table.InternalTable;
-import org.apache.ignite.internal.table.NodeUtils;
 import 
org.apache.ignite.internal.table.distributed.schema.PartitionCommandsMarshallerImpl;
 import org.apache.ignite.internal.testframework.IgniteTestUtils;
 import org.apache.ignite.internal.testframework.log4j2.LogInspector;
@@ -325,19 +323,6 @@ class ItTableRaftSnapshotsTest extends 
ClusterPerTestIntegrationTest {
         LOG.info("Lease is accepted by [nodeConsistentId={}].", 
primary.join().getLeaseholder());
     }
 
-    private @Nullable String getPrimaryReplicaName() {
-        IgniteImpl node = unwrapIgniteImpl(cluster.node(0));
-
-        CompletableFuture<ReplicaMeta> primary = 
node.placementDriver().getPrimaryReplica(
-                cluster.solePartitionId(TEST_ZONE_NAME, TEST_TABLE_NAME),
-                node.clockService().now()
-        );
-
-        assertThat(primary, willCompleteSuccessfully());
-
-        return primary.join().getLeaseholder();
-    }
-
     private void putToNode(int nodeIndex, int key, String value) {
         putToNode(nodeIndex, key, value, null);
     }
@@ -456,18 +441,7 @@ class ItTableRaftSnapshotsTest extends 
ClusterPerTestIntegrationTest {
     }
 
     private void transferPrimaryOnSolePartitionTo(int nodeIndex) throws 
InterruptedException {
-        String proposedPrimaryName = cluster.node(nodeIndex).name();
-
-        if (!proposedPrimaryName.equals(getPrimaryReplicaName())) {
-
-            String newPrimaryName = NodeUtils.transferPrimary(
-                    
cluster.runningNodes().map(TestWrappers::unwrapIgniteImpl).collect(toList()),
-                    cluster.solePartitionId(TEST_ZONE_NAME, TEST_TABLE_NAME),
-                    proposedPrimaryName
-            );
-
-            assertEquals(proposedPrimaryName, newPrimaryName);
-        }
+        cluster.transferPrimaryTo(nodeIndex, 
cluster.solePartitionId(TEST_ZONE_NAME, TEST_TABLE_NAME));
     }
 
     /**
diff --git 
a/modules/runner/src/integrationTest/java/org/apache/ignite/internal/schemasync/ItSchemaSyncAndReplicationTest.java
 
b/modules/runner/src/integrationTest/java/org/apache/ignite/internal/schemasync/ItSchemaSyncAndReplicationTest.java
index 449ec1a32f1..c216102f31a 100644
--- 
a/modules/runner/src/integrationTest/java/org/apache/ignite/internal/schemasync/ItSchemaSyncAndReplicationTest.java
+++ 
b/modules/runner/src/integrationTest/java/org/apache/ignite/internal/schemasync/ItSchemaSyncAndReplicationTest.java
@@ -35,6 +35,7 @@ import 
org.apache.ignite.internal.ClusterPerTestIntegrationTest;
 import org.apache.ignite.internal.app.IgniteImpl;
 import org.apache.ignite.internal.metastorage.server.WatchListenerInhibitor;
 import org.apache.ignite.internal.metastorage.server.raft.MetastorageGroupId;
+import org.apache.ignite.internal.replicator.ReplicationGroupId;
 import org.apache.ignite.internal.storage.MvPartitionStorage;
 import org.apache.ignite.internal.storage.RowId;
 import org.apache.ignite.internal.table.TableViewInternal;
@@ -118,7 +119,10 @@ class ItSchemaSyncAndReplicationTest extends 
ClusterPerTestIntegrationTest {
 
     private void transferLeadershipsTo(int nodeIndex) throws 
InterruptedException {
         cluster.transferLeadershipTo(nodeIndex, MetastorageGroupId.INSTANCE);
-        cluster.transferLeadershipTo(nodeIndex, 
cluster.solePartitionId(ZONE_NAME, TABLE_NAME));
+
+        ReplicationGroupId solePartitionId = 
cluster.solePartitionId(ZONE_NAME, TABLE_NAME);
+        cluster.transferLeadershipTo(nodeIndex, solePartitionId);
+        cluster.transferPrimaryTo(nodeIndex, solePartitionId);
     }
 
     private CompletableFuture<?> rejectionDueToMetadataLagTriggered() {
diff --git 
a/modules/runner/src/testFixtures/java/org/apache/ignite/internal/Cluster.java 
b/modules/runner/src/testFixtures/java/org/apache/ignite/internal/Cluster.java
index 52b0082569b..6350451d9a6 100644
--- 
a/modules/runner/src/testFixtures/java/org/apache/ignite/internal/Cluster.java
+++ 
b/modules/runner/src/testFixtures/java/org/apache/ignite/internal/Cluster.java
@@ -18,6 +18,7 @@
 package org.apache.ignite.internal;
 
 import static java.util.Collections.nCopies;
+import static java.util.concurrent.TimeUnit.SECONDS;
 import static java.util.stream.Collectors.joining;
 import static java.util.stream.Collectors.toList;
 import static org.apache.ignite.internal.ClusterConfiguration.configOverrides;
@@ -47,7 +48,6 @@ import java.util.Set;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.CopyOnWriteArrayList;
-import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.BiPredicate;
 import java.util.function.BooleanSupplier;
@@ -68,10 +68,12 @@ import 
org.apache.ignite.internal.lang.IgniteStringFormatter;
 import org.apache.ignite.internal.logger.IgniteLogger;
 import org.apache.ignite.internal.logger.Loggers;
 import org.apache.ignite.internal.network.NetworkMessage;
+import org.apache.ignite.internal.placementdriver.ReplicaMeta;
 import org.apache.ignite.internal.raft.RaftNodeId;
 import org.apache.ignite.internal.raft.server.impl.JraftServerImpl;
 import org.apache.ignite.internal.replicator.ReplicationGroupId;
 import org.apache.ignite.internal.sql.SqlCommon;
+import org.apache.ignite.internal.table.NodeUtils;
 import org.apache.ignite.internal.testframework.TestIgnitionManager;
 import org.apache.ignite.raft.jraft.RaftGroupService;
 import org.apache.ignite.raft.jraft.Status;
@@ -453,7 +455,7 @@ public class Cluster {
      */
     public Ignite startNode(int index, String nodeBootstrapConfigTemplate) {
         ServerRegistration registration = startEmbeddedNode(index, 
nodeBootstrapConfigTemplate);
-        assertThat("nodeIndex=" + index, registration.registrationFuture(), 
willSucceedIn(20, TimeUnit.SECONDS));
+        assertThat("nodeIndex=" + index, registration.registrationFuture(), 
willSucceedIn(20, SECONDS));
         Ignite newIgniteNode = registration.server().api();
 
         assertEquals(newIgniteNode, nodes.get(index));
@@ -755,6 +757,39 @@ public class Cluster {
         }
     }
 
+    /**
+     * Transfers primary replica of given replication group to the node with 
given index.
+     *
+     * @param nodeIndex Destination node index.
+     * @param groupId ID of the replication group.
+     */
+    public void transferPrimaryTo(int nodeIndex, ReplicationGroupId groupId) 
throws InterruptedException {
+        String proposedPrimaryName = node(nodeIndex).name();
+
+        if (!proposedPrimaryName.equals(getPrimaryReplicaName(groupId))) {
+
+            String newPrimaryName = NodeUtils.transferPrimary(
+                    
runningNodes().map(TestWrappers::unwrapIgniteImpl).collect(toList()),
+                    groupId,
+                    proposedPrimaryName
+            );
+
+            assertEquals(proposedPrimaryName, newPrimaryName);
+        }
+    }
+
+    private @Nullable String getPrimaryReplicaName(ReplicationGroupId groupId) 
{
+        IgniteImpl node = unwrapIgniteImpl(aliveNode());
+
+        CompletableFuture<ReplicaMeta> primary = node.placementDriver()
+                .awaitPrimaryReplica(groupId, node.clockService().now(), 30, 
SECONDS);
+
+        assertThat(primary, willCompleteSuccessfully());
+
+        @Nullable ReplicaMeta replicaMeta = primary.join();
+        return replicaMeta != null ? replicaMeta.getLeaseholder() : null;
+    }
+
     /**
      * Returns the ID of the sole partition that exists in the cluster or 
throws if there are less than one
      * or more than one partitions.

Reply via email to