This is an automated email from the ASF dual-hosted git repository.
ibessonov 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 9fd225eb48e IGNITE-25692 Return old names in handshake messages (#6063)
9fd225eb48e is described below
commit 9fd225eb48e34aed0f63a52c8d60aabb86edcced
Author: Ivan Bessonov <[email protected]>
AuthorDate: Tue Jun 17 16:17:42 2025 +0300
IGNITE-25692 Return old names in handshake messages (#6063)
---
.../recovery/RecoveryAcceptorHandshakeManager.java | 12 ++++++------
.../RecoveryInitiatorHandshakeManager.java | 22 +++++++++++-----------
.../recovery/message/HandshakeStartMessage.java | 4 ++--
.../message/HandshakeStartResponseMessage.java | 2 +-
.../RecoveryAcceptorHandshakeManagerTest.java | 2 +-
.../RecoveryInitiatorHandshakeManagerTest.java | 4 ++--
6 files changed, 23 insertions(+), 23 deletions(-)
diff --git
a/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/RecoveryAcceptorHandshakeManager.java
b/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/RecoveryAcceptorHandshakeManager.java
index eaf839d9047..b66752fc783 100644
---
a/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/RecoveryAcceptorHandshakeManager.java
+++
b/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/RecoveryAcceptorHandshakeManager.java
@@ -168,8 +168,8 @@ public class RecoveryAcceptorHandshakeManager implements
HandshakeManager {
@Override
public void onConnectionOpen() {
HandshakeStartMessage handshakeStartMessage =
messageFactory.handshakeStartMessage()
-
.acceptorNode(HandshakeManagerUtils.clusterNodeToMessage(localNode))
- .acceptorClusterId(clusterIdSupplier.clusterId())
+
.serverNode(HandshakeManagerUtils.clusterNodeToMessage(localNode))
+ .serverClusterId(clusterIdSupplier.clusterId())
.productName(productVersionSource.productName())
.productVersion(productVersionSource.productVersion().toString())
.build();
@@ -224,7 +224,7 @@ public class RecoveryAcceptorHandshakeManager implements
HandshakeManager {
return;
}
- this.remoteNode = message.initiatorNode().asClusterNode();
+ this.remoteNode = message.clientNode().asClusterNode();
this.receivedCount = message.receivedCount();
this.remoteChannelId = message.connectionId();
@@ -233,7 +233,7 @@ public class RecoveryAcceptorHandshakeManager implements
HandshakeManager {
}
private boolean
possiblyRejectHandshakeStartResponse(HandshakeStartResponseMessage message) {
- if (staleIdDetector.isIdStale(message.initiatorNode().id())) {
+ if (staleIdDetector.isIdStale(message.clientNode().id())) {
handleStaleInitiatorId(message);
return true;
@@ -250,7 +250,7 @@ public class RecoveryAcceptorHandshakeManager implements
HandshakeManager {
private void handleStaleInitiatorId(HandshakeStartResponseMessage msg) {
String message = String.format("%s:%s is stale, it should be restarted
to be allowed to connect",
- msg.initiatorNode().name(), msg.initiatorNode().id()
+ msg.clientNode().name(), msg.clientNode().id()
);
sendRejectionMessageAndFailHandshake(message,
HandshakeRejectionReason.STALE_LAUNCH_ID, HandshakeException::new);
@@ -258,7 +258,7 @@ public class RecoveryAcceptorHandshakeManager implements
HandshakeManager {
private void
handleRefusalToEstablishConnectionDueToStopping(HandshakeStartResponseMessage
msg) {
String message = String.format("%s:%s tried to establish a connection
with %s, but it's stopping",
- msg.initiatorNode().name(), msg.initiatorNode().id(),
localNode.name()
+ msg.clientNode().name(), msg.clientNode().id(),
localNode.name()
);
sendRejectionMessageAndFailHandshake(message,
HandshakeRejectionReason.STOPPING, m -> new NodeStoppingException());
diff --git
a/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/RecoveryInitiatorHandshakeManager.java
b/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/RecoveryInitiatorHandshakeManager.java
index ac6785e5822..72e88cf9f0d 100644
---
a/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/RecoveryInitiatorHandshakeManager.java
+++
b/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/RecoveryInitiatorHandshakeManager.java
@@ -262,7 +262,7 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
return;
}
- this.remoteNode = handshakeStartMessage.acceptorNode().asClusterNode();
+ this.remoteNode = handshakeStartMessage.serverNode().asClusterNode();
ChannelKey channelKey = new ChannelKey(remoteNode.name(),
remoteNode.id(), connectionId);
switchEventLoopIfNeeded(channel, channelKey, channelEventLoopsSource,
() -> proceedAfterSavingIds(handshakeStartMessage));
@@ -305,19 +305,19 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
}
private boolean possiblyRejectHandshakeStart(HandshakeStartMessage
message) {
- if (message.acceptorNode().id().equals(localNode.id())) {
+ if (message.serverNode().id().equals(localNode.id())) {
handleLoopConnection(message);
return true;
}
- if (staleIdDetector.isIdStale(message.acceptorNode().id())) {
+ if (staleIdDetector.isIdStale(message.serverNode().id())) {
handleStaleAcceptorId(message);
return true;
}
- if (clusterIdMismatch(message.acceptorClusterId(),
clusterIdSupplier.clusterId())) {
+ if (clusterIdMismatch(message.serverClusterId(),
clusterIdSupplier.clusterId())) {
handleClusterIdMismatch(message);
return true;
@@ -348,7 +348,7 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
String message = String.format(
"Got handshake start from self, this should never happen; this
is a programming error [localNode=%s, acceptorNode=%s]",
localNode,
- msg.acceptorNode()
+ msg.serverNode()
);
sendRejectionMessageAndFailHandshake(message,
HandshakeRejectionReason.LOOP, CriticalHandshakeException::new);
@@ -367,7 +367,7 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
private void handleStaleAcceptorId(HandshakeStartMessage msg) {
String message = String.format("%s:%s is stale, node should be
restarted so that other nodes can connect",
- msg.acceptorNode().name(), msg.acceptorNode().id()
+ msg.serverNode().name(), msg.serverNode().id()
);
sendRejectionMessageAndFailHandshake(message,
HandshakeRejectionReason.STALE_LAUNCH_ID, HandshakeException::new);
@@ -376,7 +376,7 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
private void handleClusterIdMismatch(HandshakeStartMessage msg) {
String message = String.format(
"%s:%s belongs to cluster %s which is different from this one
%s, connection rejected; should CMG/MG repair be finished?",
- msg.acceptorNode().name(), msg.acceptorNode().id(),
msg.acceptorClusterId(), clusterIdSupplier.clusterId()
+ msg.serverNode().name(), msg.serverNode().id(),
msg.serverClusterId(), clusterIdSupplier.clusterId()
);
sendRejectionMessageAndFailHandshake(message,
HandshakeRejectionReason.CLUSTER_ID_MISMATCH, HandshakeException::new);
@@ -384,7 +384,7 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
private void handleProductNameMismatch(HandshakeStartMessage msg) {
String message = String.format("%s:%s runs product '%s' which is
different from this one '%s', connection rejected",
- msg.acceptorNode().name(), msg.acceptorNode().id(),
msg.productName(), productVersionSource.productName()
+ msg.serverNode().name(), msg.serverNode().id(),
msg.productName(), productVersionSource.productName()
);
sendRejectionMessageAndFailHandshake(message,
HandshakeRejectionReason.PRODUCT_MISMATCH, HandshakeException::new);
@@ -392,7 +392,7 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
private void handleProductVersionMismatch(HandshakeStartMessage msg) {
String message = String.format("%s:%s runs product version '%s' which
is different from this one '%s', connection rejected",
- msg.acceptorNode().name(), msg.acceptorNode().id(),
msg.productVersion(), productVersionSource.productVersion()
+ msg.serverNode().name(), msg.serverNode().id(),
msg.productVersion(), productVersionSource.productVersion()
);
sendRejectionMessageAndFailHandshake(message,
HandshakeRejectionReason.VERSION_MISMATCH, HandshakeException::new);
@@ -400,7 +400,7 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
private void
handleRefusalToEstablishConnectionDueToStopping(HandshakeStartMessage msg) {
String message = String.format("%s:%s tried to establish a connection
with %s, but it's stopping",
- msg.acceptorNode().name(), msg.acceptorNode().id(),
localNode.name()
+ msg.serverNode().name(), msg.serverNode().id(),
localNode.name()
);
sendRejectionMessageAndFailHandshake(message,
HandshakeRejectionReason.STOPPING, m -> new NodeStoppingException());
@@ -475,7 +475,7 @@ public class RecoveryInitiatorHandshakeManager implements
HandshakeManager {
PipelineUtils.afterHandshake(ctx.pipeline(), descriptor,
createMessageHandler(), MESSAGE_FACTORY);
HandshakeStartResponseMessage response =
MESSAGE_FACTORY.handshakeStartResponseMessage()
- .initiatorNode(clusterNodeToMessage(localNode))
+ .clientNode(clusterNodeToMessage(localNode))
.receivedCount(descriptor.receivedCount())
.connectionId(connectionId)
.build();
diff --git
a/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/message/HandshakeStartMessage.java
b/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/message/HandshakeStartMessage.java
index b1030e4fd18..2bd7f956bc3 100644
---
a/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/message/HandshakeStartMessage.java
+++
b/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/message/HandshakeStartMessage.java
@@ -30,11 +30,11 @@ import org.jetbrains.annotations.Nullable;
@Transferable(NetworkMessageTypes.HANDSHAKE_START)
public interface HandshakeStartMessage extends InternalMessage {
/** Returns the acceptor node that sends this. */
- ClusterNodeMessage acceptorNode();
+ ClusterNodeMessage serverNode();
/** ID of the cluster to which the acceptor node belongs ({@code null} if
it's not initialized yet. */
@Nullable
- UUID acceptorClusterId();
+ UUID serverClusterId();
/** Product name of the node that sends the message. */
String productName();
diff --git
a/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/message/HandshakeStartResponseMessage.java
b/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/message/HandshakeStartResponseMessage.java
index c6c2ecf95ad..7b6761f213a 100644
---
a/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/message/HandshakeStartResponseMessage.java
+++
b/modules/network/src/main/java/org/apache/ignite/internal/network/recovery/message/HandshakeStartResponseMessage.java
@@ -28,7 +28,7 @@ import
org.apache.ignite.internal.network.message.ClusterNodeMessage;
@Transferable(NetworkMessageTypes.HANDSHAKE_START_RESPONSE)
public interface HandshakeStartResponseMessage extends InternalMessage {
/** Returns the initiator node that sends this. */
- ClusterNodeMessage initiatorNode();
+ ClusterNodeMessage clientNode();
/**
* Returns connection id.
diff --git
a/modules/network/src/test/java/org/apache/ignite/internal/network/recovery/RecoveryAcceptorHandshakeManagerTest.java
b/modules/network/src/test/java/org/apache/ignite/internal/network/recovery/RecoveryAcceptorHandshakeManagerTest.java
index 187197de46f..9c867e83193 100644
---
a/modules/network/src/test/java/org/apache/ignite/internal/network/recovery/RecoveryAcceptorHandshakeManagerTest.java
+++
b/modules/network/src/test/java/org/apache/ignite/internal/network/recovery/RecoveryAcceptorHandshakeManagerTest.java
@@ -192,7 +192,7 @@ class RecoveryAcceptorHandshakeManagerTest extends
BaseIgniteAbstractTest {
private static HandshakeStartResponseMessage
handshakeStartResponseMessageFrom(UUID initiatorLaunchId) {
return MESSAGE_FACTORY.handshakeStartResponseMessage()
- .initiatorNode(
+ .clientNode(
MESSAGE_FACTORY.clusterNodeMessage()
.id(initiatorLaunchId)
.name(INITIATOR_CONSISTENT_ID)
diff --git
a/modules/network/src/test/java/org/apache/ignite/internal/network/recovery/RecoveryInitiatorHandshakeManagerTest.java
b/modules/network/src/test/java/org/apache/ignite/internal/network/recovery/RecoveryInitiatorHandshakeManagerTest.java
index 4f24df16432..02c4414c92b 100644
---
a/modules/network/src/test/java/org/apache/ignite/internal/network/recovery/RecoveryInitiatorHandshakeManagerTest.java
+++
b/modules/network/src/test/java/org/apache/ignite/internal/network/recovery/RecoveryInitiatorHandshakeManagerTest.java
@@ -207,7 +207,7 @@ class RecoveryInitiatorHandshakeManagerTest extends
BaseIgniteAbstractTest {
private static HandshakeStartMessage handshakeStartMessageFrom(UUID
acceptorLaunchId, UUID acceptorClusterId) {
return MESSAGE_FACTORY.handshakeStartMessage()
- .acceptorNode(
+ .serverNode(
MESSAGE_FACTORY.clusterNodeMessage()
.id(acceptorLaunchId)
.name(ACCEPTOR_CONSISTENT_ID)
@@ -215,7 +215,7 @@ class RecoveryInitiatorHandshakeManagerTest extends
BaseIgniteAbstractTest {
.port(PORT)
.build()
)
- .acceptorClusterId(acceptorClusterId)
+ .serverClusterId(acceptorClusterId)
.productName(IgniteProductVersion.CURRENT_PRODUCT)
.productVersion(IgniteProductVersion.CURRENT_VERSION.toString())
.build();