JAkutenshi commented on code in PR #5366:
URL: https://github.com/apache/ignite-3/pull/5366#discussion_r1984987136
##########
modules/partition-replicator/src/main/java/org/apache/ignite/internal/partition/replicator/handlers/ReplicaSafeTimeSyncRequestHandler.java:
##########
@@ -54,9 +56,14 @@ public ReplicaSafeTimeSyncRequestHandler(
* Handles {@link ReplicaSafeTimeSyncRequest}.
*
* @param request Request to handle.
+ * @param isPrimary Whether current node is a primary replica.
* @return Future that will be completed when the request is handled.
*/
- public CompletableFuture<?> handle(ReplicaSafeTimeSyncRequest request) {
+ public CompletableFuture<?> handle(ReplicaSafeTimeSyncRequest request,
boolean isPrimary) {
+ if (!isPrimary) {
+ return nullCompletedFuture();
Review Comment:
Thank you
##########
modules/partition-replicator/src/integrationTest/java/org/apache/ignite/internal/partition/replicator/ItColocationTxRecoveryTest.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.ignite.internal.partition.replicator;
+
+import static
org.apache.ignite.internal.testframework.matchers.CompletableFutureMatcher.willCompleteSuccessfully;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+
+import java.util.concurrent.CompletableFuture;
+import org.apache.ignite.internal.partition.replicator.fixtures.Node;
+import org.apache.ignite.internal.placementdriver.ReplicaMeta;
+import org.apache.ignite.internal.replicator.ZonePartitionId;
+import org.apache.ignite.table.KeyValueView;
+import org.apache.ignite.tx.Transaction;
+import org.junit.jupiter.api.Test;
+
+class ItColocationTxRecoveryTest extends ItAbstractColocationTest {
+ private static final long KEY = 1;
+
+ /**
+ * Tests that tx recovery works. Scenario:
+ *
+ * <ol>
+ * <li>A transaction tx1 is started, it takes a shared lock on a key
and never gets finished</li>
+ * <li>Its coordinator (different from the node hosting the touched
partition primary) is stopped, so the transaction becomes
+ * abandoned</li>
+ * <li>Transaction tx2 tries to write to the same key, founds an
incompatible lock, realizes that it's held by an abandoned
+ * transaction, and does tx recovery to remove the lock on the
partition primary</li>
+ * <li>tx2 should succeed</li>
+ * </ol>
+ */
+ @Test
+ void abandonedTransactionGetsAbortedOnTouch() throws Exception {
+ assertThat(txConfiguration.abandonedCheckTs().update(600_000L),
willCompleteSuccessfully());
+
+ startCluster(3);
+
+ Node node0 = getNode(0);
+
+ // Create a zone with a single partition on every node.
+ int zoneId = createZone(node0, TEST_ZONE_NAME, 1, cluster.size());
+
+ createTable(node0, TEST_ZONE_NAME, TEST_TABLE_NAME1);
+
+ cluster.forEach(Node::waitForMetadataCompletenessAtNow);
+
+ putInitialValue(node0);
+
+ ReplicaMeta primaryReplica = getPrimaryReplica(zoneId);
+
+ Node coordinatorNodeToBeStopped = findAnyOtherNode(primaryReplica);
+ Transaction txToBeAbandoned =
coordinatorNodeToBeStopped.transactions().begin();
+ // Trigger a shared lock to be taken on the key.
+ coordinatorNodeToBeStopped.tableManager.table(TEST_TABLE_NAME1)
+ .keyValueView(Long.class, Integer.class)
+ .get(txToBeAbandoned, KEY);
+
+ coordinatorNodeToBeStopped.stop();
+ cluster.remove(coordinatorNodeToBeStopped);
+
+ Node runningNode = cluster.get(0);
+
+ KeyValueView<Long, Integer> kvView =
runningNode.tableManager.table(TEST_TABLE_NAME1).keyValueView(Long.class,
Integer.class);
+
+ Transaction conflictingTx = runningNode.transactions().begin();
+ assertDoesNotThrow(() -> kvView.put(conflictingTx, KEY, 111));
+ }
+
+ private static void putInitialValue(Node node0) {
+ node0.tableManager.table(TEST_TABLE_NAME1).keyValueView(Long.class,
Integer.class).put(null, KEY, 42);
+ }
+
+ private ReplicaMeta getPrimaryReplica(int zoneId) {
+ Node node = cluster.get(0);
+
+ CompletableFuture<ReplicaMeta> primaryReplicaFuture =
node.placementDriverManager.placementDriver().getPrimaryReplica(
+ new ZonePartitionId(zoneId, 0),
+ node.hybridClock.now()
+ );
+
+ assertThat(primaryReplicaFuture, willCompleteSuccessfully());
+ return primaryReplicaFuture.join();
+ }
+
+ private Node findAnyOtherNode(ReplicaMeta primaryReplica) {
Review Comment:
`findAnyNonPrimaryColocatedNode` or `findAnyNonPrimaryNode` is better, but
opinionated.
##########
modules/partition-replicator/src/main/java/org/apache/ignite/internal/partition/replicator/ZonePartitionReplicaListener.java:
##########
@@ -138,26 +148,23 @@ public ZonePartitionReplicaListener(
replicationGroupId
);
- minimumActiveTxTimeReplicaRequestHandler = new
MinimumActiveTxTimeReplicaRequestHandler(
- clockService,
- raftCommandApplicator
- );
-
- vacuumTxStateReplicaRequestHandler = new
VacuumTxStateReplicaRequestHandler(raftCommandApplicator);
-
txStateCommitPartitionReplicaRequestHandler = new
TxStateCommitPartitionReplicaRequestHandler(
txStatePartitionStorage,
txManager,
clusterNodeResolver,
localNode,
- new TxRecoveryEngine(
- txManager,
- clusterNodeResolver,
- replicationGroupId,
-
ZonePartitionReplicaListener::createAbandonedTxRecoveryEnlistment
- )
+ txRecoveryEngine
);
+ txRecoveryMessageHandler = new
TxRecoveryMessageHandler(txStatePartitionStorage, replicationGroupId,
txRecoveryEngine);
+
+ minimumActiveTxTimeReplicaRequestHandler = new
MinimumActiveTxTimeReplicaRequestHandler(
+ clockService,
+ raftCommandApplicator
+ );
+
+ vacuumTxStateReplicaRequestHandler = new
VacuumTxStateReplicaRequestHandler(raftCommandApplicator);
Review Comment:
Can I ask you to delete an excessive newline whitespace inside
`VacuumTxStateReplicaRequestHandler#handle`?
##########
modules/partition-replicator/src/integrationTest/java/org/apache/ignite/internal/partition/replicator/ItColocationTxRecoveryTest.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.ignite.internal.partition.replicator;
+
+import static
org.apache.ignite.internal.testframework.matchers.CompletableFutureMatcher.willCompleteSuccessfully;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+
+import java.util.concurrent.CompletableFuture;
+import org.apache.ignite.internal.partition.replicator.fixtures.Node;
+import org.apache.ignite.internal.placementdriver.ReplicaMeta;
+import org.apache.ignite.internal.replicator.ZonePartitionId;
+import org.apache.ignite.table.KeyValueView;
+import org.apache.ignite.tx.Transaction;
+import org.junit.jupiter.api.Test;
+
+class ItColocationTxRecoveryTest extends ItAbstractColocationTest {
+ private static final long KEY = 1;
+
+ /**
+ * Tests that tx recovery works. Scenario:
+ *
+ * <ol>
+ * <li>A transaction tx1 is started, it takes a shared lock on a key
and never gets finished</li>
+ * <li>Its coordinator (different from the node hosting the touched
partition primary) is stopped, so the transaction becomes
+ * abandoned</li>
+ * <li>Transaction tx2 tries to write to the same key, founds an
incompatible lock, realizes that it's held by an abandoned
+ * transaction, and does tx recovery to remove the lock on the
partition primary</li>
+ * <li>tx2 should succeed</li>
+ * </ol>
+ */
+ @Test
+ void abandonedTransactionGetsAbortedOnTouch() throws Exception {
+ assertThat(txConfiguration.abandonedCheckTs().update(600_000L),
willCompleteSuccessfully());
+
+ startCluster(3);
+
+ Node node0 = getNode(0);
+
+ // Create a zone with a single partition on every node.
+ int zoneId = createZone(node0, TEST_ZONE_NAME, 1, cluster.size());
+
+ createTable(node0, TEST_ZONE_NAME, TEST_TABLE_NAME1);
+
+ cluster.forEach(Node::waitForMetadataCompletenessAtNow);
+
+ putInitialValue(node0);
+
+ ReplicaMeta primaryReplica = getPrimaryReplica(zoneId);
+
+ Node coordinatorNodeToBeStopped = findAnyOtherNode(primaryReplica);
+ Transaction txToBeAbandoned =
coordinatorNodeToBeStopped.transactions().begin();
+ // Trigger a shared lock to be taken on the key.
+ coordinatorNodeToBeStopped.tableManager.table(TEST_TABLE_NAME1)
+ .keyValueView(Long.class, Integer.class)
+ .get(txToBeAbandoned, KEY);
+
+ coordinatorNodeToBeStopped.stop();
+ cluster.remove(coordinatorNodeToBeStopped);
+
+ Node runningNode = cluster.get(0);
+
+ KeyValueView<Long, Integer> kvView =
runningNode.tableManager.table(TEST_TABLE_NAME1).keyValueView(Long.class,
Integer.class);
+
+ Transaction conflictingTx = runningNode.transactions().begin();
+ assertDoesNotThrow(() -> kvView.put(conflictingTx, KEY, 111));
+ }
+
+ private static void putInitialValue(Node node0) {
Review Comment:
Why the separate method? And argument just a `node` is better there
##########
modules/partition-replicator/src/integrationTest/java/org/apache/ignite/internal/partition/replicator/ItColocationTxRecoveryTest.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.ignite.internal.partition.replicator;
+
+import static
org.apache.ignite.internal.testframework.matchers.CompletableFutureMatcher.willCompleteSuccessfully;
+import static org.hamcrest.MatcherAssert.assertThat;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+
+import java.util.concurrent.CompletableFuture;
+import org.apache.ignite.internal.partition.replicator.fixtures.Node;
+import org.apache.ignite.internal.placementdriver.ReplicaMeta;
+import org.apache.ignite.internal.replicator.ZonePartitionId;
+import org.apache.ignite.table.KeyValueView;
+import org.apache.ignite.tx.Transaction;
+import org.junit.jupiter.api.Test;
+
+class ItColocationTxRecoveryTest extends ItAbstractColocationTest {
+ private static final long KEY = 1;
+
+ /**
+ * Tests that tx recovery works. Scenario:
+ *
+ * <ol>
+ * <li>A transaction tx1 is started, it takes a shared lock on a key
and never gets finished</li>
+ * <li>Its coordinator (different from the node hosting the touched
partition primary) is stopped, so the transaction becomes
+ * abandoned</li>
+ * <li>Transaction tx2 tries to write to the same key, founds an
incompatible lock, realizes that it's held by an abandoned
+ * transaction, and does tx recovery to remove the lock on the
partition primary</li>
+ * <li>tx2 should succeed</li>
+ * </ol>
+ */
+ @Test
+ void abandonedTransactionGetsAbortedOnTouch() throws Exception {
+ assertThat(txConfiguration.abandonedCheckTs().update(600_000L),
willCompleteSuccessfully());
+
+ startCluster(3);
+
+ Node node0 = getNode(0);
+
+ // Create a zone with a single partition on every node.
+ int zoneId = createZone(node0, TEST_ZONE_NAME, 1, cluster.size());
+
+ createTable(node0, TEST_ZONE_NAME, TEST_TABLE_NAME1);
+
+ cluster.forEach(Node::waitForMetadataCompletenessAtNow);
+
+ putInitialValue(node0);
+
+ ReplicaMeta primaryReplica = getPrimaryReplica(zoneId);
+
+ Node coordinatorNodeToBeStopped = findAnyOtherNode(primaryReplica);
+ Transaction txToBeAbandoned =
coordinatorNodeToBeStopped.transactions().begin();
+ // Trigger a shared lock to be taken on the key.
+ coordinatorNodeToBeStopped.tableManager.table(TEST_TABLE_NAME1)
+ .keyValueView(Long.class, Integer.class)
+ .get(txToBeAbandoned, KEY);
+
+ coordinatorNodeToBeStopped.stop();
+ cluster.remove(coordinatorNodeToBeStopped);
+
+ Node runningNode = cluster.get(0);
+
+ KeyValueView<Long, Integer> kvView =
runningNode.tableManager.table(TEST_TABLE_NAME1).keyValueView(Long.class,
Integer.class);
+
+ Transaction conflictingTx = runningNode.transactions().begin();
+ assertDoesNotThrow(() -> kvView.put(conflictingTx, KEY, 111));
+ }
+
+ private static void putInitialValue(Node node0) {
+ node0.tableManager.table(TEST_TABLE_NAME1).keyValueView(Long.class,
Integer.class).put(null, KEY, 42);
+ }
+
+ private ReplicaMeta getPrimaryReplica(int zoneId) {
+ Node node = cluster.get(0);
+
+ CompletableFuture<ReplicaMeta> primaryReplicaFuture =
node.placementDriverManager.placementDriver().getPrimaryReplica(
+ new ZonePartitionId(zoneId, 0),
+ node.hybridClock.now()
+ );
+
+ assertThat(primaryReplicaFuture, willCompleteSuccessfully());
+ return primaryReplicaFuture.join();
Review Comment:
`ReplicaMeta` from `getPrimaryReplica` may be `null`. Should we consider it
and check?
--
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]