This is an automated email from the ASF dual-hosted git repository.
petrov-mg pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ignite.git
The following commit(s) were added to refs/heads/master by this push:
new 97e2d106efb IGNITE-29001 Removed unused Message parameter for
TcpDiscoverySpi#writeToSocket methods (#13505)
97e2d106efb is described below
commit 97e2d106efb1b8fa77b19897aff24334d5dcf35b
Author: Mikhail Petrov <[email protected]>
AuthorDate: Wed Aug 26 16:04:02 2026 +0300
IGNITE-29001 Removed unused Message parameter for
TcpDiscoverySpi#writeToSocket methods (#13505)
---
.../ignite/spi/discovery/tcp/ServerImpl.java | 26 ++++++-------
.../ignite/spi/discovery/tcp/TcpDiscoverySpi.java | 8 +---
.../ignite/internal/IgniteClientRejoinTest.java | 6 +--
.../IgniteDiscoveryMassiveNodeFailTest.java | 28 ++++++--------
.../dht/IgniteCacheTopologySplitAbstractTest.java | 6 +--
.../processors/rest/RestProcessorHangTest.java | 4 +-
.../IgniteTcpCommunicationConnectOnInitTest.java | 6 +--
.../spi/discovery/tcp/BlockTcpDiscoverySpi.java | 11 ++++--
.../discovery/tcp/MultiDataCenterSplitTest.java | 5 +--
.../spi/discovery/tcp/ReceivedMessagesTracker.java | 45 ++++++++++++++++++++++
...cpClientDiscoverySpiFailureTimeoutSelfTest.java | 11 ++----
.../tcp/TcpClientDiscoverySpiSelfTest.java | 29 ++++++++++----
.../tcp/TcpDiscoveryCoordinatorFailureTest.java | 32 ---------------
.../discovery/tcp/TcpDiscoveryFailedJoinTest.java | 6 +--
.../tcp/TcpDiscoveryNetworkIssuesTest.java | 6 +--
.../TcpDiscoveryPendingMessageDeliveryTest.java | 6 +--
.../tcp/TcpDiscoverySpiFailureTimeoutSelfTest.java | 12 +++---
.../tcp/TcpDiscoverySpiReconnectDelayTest.java | 18 +++++++--
.../spi/discovery/tcp/TestTcpDiscoverySpi.java | 30 +++++++++++++++
19 files changed, 169 insertions(+), 126 deletions(-)
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
index c97661427c1..fe86a1c789b 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ServerImpl.java
@@ -6861,7 +6861,7 @@ class ServerImpl extends TcpDiscoveryImpl {
if (msg instanceof TcpDiscoveryConnectionCheckMessage)
{
ringMessageReceived();
- spi.writeToSocket(msg, sock, RES_OK, sockTimeout);
+ spi.writeToSocket(sock, RES_OK, sockTimeout);
continue;
}
@@ -6887,7 +6887,7 @@ class ServerImpl extends TcpDiscoveryImpl {
TcpDiscoverySpiState state = spiStateCopy();
if (state == CONNECTED) {
- spi.writeToSocket(msg, sock, RES_OK,
sockTimeout);
+ spi.writeToSocket(sock, RES_OK, sockTimeout);
if (clientMsgWrk != null &&
clientMsgWrk.runner() == null && !clientMsgWrk.isDone())
new
MessageWorkerThreadWithCleanup<>(clientMsgWrk, log).start();
@@ -6901,21 +6901,21 @@ class ServerImpl extends TcpDiscoveryImpl {
// If message is received from previous node
and node is connecting forward to next node.
if
(!getLocalNodeId().equals(msg0.routerNodeId()) && state == CONNECTING) {
- spi.writeToSocket(msg, sock, RES_OK,
sockTimeout);
+ spi.writeToSocket(sock, RES_OK,
sockTimeout);
msgWorker.addMessage(msg);
continue;
}
- spi.writeToSocket(msg, sock,
RES_CONTINUE_JOIN, sockTimeout);
+ spi.writeToSocket(sock, RES_CONTINUE_JOIN,
sockTimeout);
break;
}
}
else if (msg instanceof
TcpDiscoveryDuplicateIdMessage) {
// Send receipt back.
- spi.writeToSocket(msg, sock, RES_OK, sockTimeout);
+ spi.writeToSocket(sock, RES_OK, sockTimeout);
boolean ignored = false;
@@ -6944,7 +6944,7 @@ class ServerImpl extends TcpDiscoveryImpl {
}
else if (msg instanceof TcpDiscoveryAuthFailedMessage)
{
// Send receipt back.
- spi.writeToSocket(msg, sock, RES_OK, sockTimeout);
+ spi.writeToSocket(sock, RES_OK, sockTimeout);
synchronized (mux) {
if (spiState == CONNECTING) {
@@ -6972,7 +6972,7 @@ class ServerImpl extends TcpDiscoveryImpl {
}
else if (msg instanceof
TcpDiscoveryCheckFailedMessage) {
// Send receipt back.
- spi.writeToSocket(msg, sock, RES_OK, sockTimeout);
+ spi.writeToSocket(sock, RES_OK, sockTimeout);
boolean ignored = false;
@@ -7015,7 +7015,7 @@ class ServerImpl extends TcpDiscoveryImpl {
}
else if (msg instanceof
TcpDiscoveryLoopbackProblemMessage) {
// Send receipt back.
- spi.writeToSocket(msg, sock, RES_OK, sockTimeout);
+ spi.writeToSocket(sock, RES_OK, sockTimeout);
boolean ignored = false;
@@ -7080,7 +7080,7 @@ class ServerImpl extends TcpDiscoveryImpl {
clientMsgWrk.addMessage(ack);
}
else
- spi.writeToSocket(msg, sock, RES_OK, sockTimeout);
+ spi.writeToSocket(sock, RES_OK, sockTimeout);
if (metricsUpdateMsg != null)
processClientMetricsUpdateMessage(metricsUpdateMsg);
@@ -7398,13 +7398,13 @@ class ServerImpl extends TcpDiscoveryImpl {
// Check that joining node can accept incoming connections.
if (node.clientRouterNodeId() == null) {
if (!pingJoiningNode(node)) {
- spi.writeToSocket(msg, sock, RES_JOIN_IMPOSSIBLE,
sockTimeout);
+ spi.writeToSocket(sock, RES_JOIN_IMPOSSIBLE,
sockTimeout);
return false;
}
}
- spi.writeToSocket(msg, sock, RES_OK, sockTimeout);
+ spi.writeToSocket(sock, RES_OK, sockTimeout);
if (log.isDebugEnabled())
log.debug("Responded to join request message [msg=" + msg
+ ", res=" + RES_OK + ']');
@@ -7441,7 +7441,7 @@ class ServerImpl extends TcpDiscoveryImpl {
// Local node is stopping. Remote node should try next one.
res = RES_CONTINUE_JOIN;
- spi.writeToSocket(msg, sock, res, sockTimeout);
+ spi.writeToSocket(sock, res, sockTimeout);
if (log.isDebugEnabled())
log.debug("Responded to join request message [msg=" + msg
+ ", res=" + res + ']');
@@ -7719,7 +7719,7 @@ class ServerImpl extends TcpDiscoveryImpl {
throws IgniteCheckedException, IOException {
byte[] msgBytes = msgT.get2() == null ?
clientMsgSer.serializeMessage(msgT.get1()) : msgT.get2();
- spi.writeToSocket(sock, msgT.get1(), msgBytes, timeout);
+ spi.writeToSocket(sock, msgBytes, timeout);
}
/**
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
index 84da71e3cec..a20ac2d7f01 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpi.java
@@ -1617,7 +1617,7 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
sock.connect(resolved,
(int)timeoutHelper.nextTimeoutChunk(sockTimeout));
- writeToSocket(sock, null, U.IGNITE_HEADER,
timeoutHelper.nextTimeoutChunk(sockTimeout));
+ writeToSocket(sock, U.IGNITE_HEADER,
timeoutHelper.nextTimeoutChunk(sockTimeout));
return sock;
}
@@ -1722,10 +1722,9 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
}
/**
- * Writes message to the socket.
+ * Writes raw data to the socket.
*
* @param sock Socket.
- * @param msg Message.
* @param data Raw data to write.
* @param timeout Socket write timeout.
* @throws IOException If IO failed or write timed out.
@@ -1733,7 +1732,6 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
*/
protected void writeToSocket(
Socket sock,
- @Nullable TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
@@ -1808,7 +1806,6 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
/**
* Writes response to the socket.
*
- * @param msg Received message.
* @param sock Socket.
* @param res Integer response.
* @param timeout Socket timeout.
@@ -1816,7 +1813,6 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
* @throws IgniteCheckedException If node is not yet initialized or is
stopping.
*/
protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/IgniteClientRejoinTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/IgniteClientRejoinTest.java
index 4a5ced6eb64..23be4f0f8b0 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/IgniteClientRejoinTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/IgniteClientRejoinTest.java
@@ -355,14 +355,13 @@ public class IgniteClientRejoinTest extends
GridCommonAbstractTest {
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
if (blockAll || block && sock.getPort() == 47500)
throw new SocketException("Test discovery exception");
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/** {@inheritDoc} */
@@ -379,7 +378,6 @@ public class IgniteClientRejoinTest extends
GridCommonAbstractTest {
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
@@ -387,7 +385,7 @@ public class IgniteClientRejoinTest extends
GridCommonAbstractTest {
if (blockAll || block && sock.getPort() == 47500)
throw new SocketException("Test discovery exception");
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
/** {@inheritDoc} */
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/IgniteDiscoveryMassiveNodeFailTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/IgniteDiscoveryMassiveNodeFailTest.java
index f69d068efe1..2ddece5dc92 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/IgniteDiscoveryMassiveNodeFailTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/IgniteDiscoveryMassiveNodeFailTest.java
@@ -304,16 +304,15 @@ public class IgniteDiscoveryMassiveNodeFailTest extends
GridCommonAbstractTest {
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
assertNotFailedNode(sock);
- if (isDrop(msg))
+ if (isDrop())
return;
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/** {@inheritDoc} */
@@ -321,37 +320,32 @@ public class IgniteDiscoveryMassiveNodeFailTest extends
GridCommonAbstractTest {
long timeout) throws IOException, IgniteCheckedException {
assertNotFailedNode(ses.socket());
- if (isDrop(msg))
+ if (isDrop()) {
+ ignite.log().info(">> Drop message " + msg);
+
return;
+ }
super.writeMessage(ses, msg, timeout);
}
- /**
- *
- */
- private boolean isDrop(TcpDiscoveryAbstractMessage msg) {
- boolean drop = failNodes && forceFailConnectivity &&
failedNodes.contains(ignite.cluster().localNode());
-
- if (drop)
- ignite.log().info(">> Drop message " + msg);
-
- return drop;
+ /** */
+ private boolean isDrop() {
+ return failNodes && forceFailConnectivity &&
failedNodes.contains(ignite.cluster().localNode());
}
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
) throws IOException, IgniteCheckedException {
assertNotFailedNode(sock);
- if (isDrop(msg))
+ if (isDrop())
return;
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
/**
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/dht/IgniteCacheTopologySplitAbstractTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/dht/IgniteCacheTopologySplitAbstractTest.java
index 14cb10543ca..f85acbd4a42 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/dht/IgniteCacheTopologySplitAbstractTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/distributed/dht/IgniteCacheTopologySplitAbstractTest.java
@@ -220,13 +220,12 @@ public abstract class
IgniteCacheTopologySplitAbstractTest extends GridCommonAbs
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
checkSegmented((InetSocketAddress)sock.getRemoteSocketAddress(),
timeout);
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/** {@inheritDoc} */
@@ -240,14 +239,13 @@ public abstract class
IgniteCacheTopologySplitAbstractTest extends GridCommonAbs
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
) throws IOException, IgniteCheckedException {
checkSegmented((InetSocketAddress)sock.getRemoteSocketAddress(),
timeout);
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
}
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/processors/rest/RestProcessorHangTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/processors/rest/RestProcessorHangTest.java
index 1fb395f619f..b473603babc 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/processors/rest/RestProcessorHangTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/processors/rest/RestProcessorHangTest.java
@@ -31,7 +31,6 @@ import org.apache.ignite.internal.IgniteKernal;
import org.apache.ignite.internal.IgnitionEx;
import org.apache.ignite.internal.processors.rest.request.GridRestCacheRequest;
import org.apache.ignite.spi.discovery.tcp.TestTcpDiscoverySpi;
-import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryAbstractMessage;
import org.apache.ignite.testframework.GridTestUtils;
import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
import org.junit.Test;
@@ -73,7 +72,6 @@ public class RestProcessorHangTest extends
GridCommonAbstractTest {
// Discovery spi that never allows connecting.
TestTcpDiscoverySpi discoSpi = new TestTcpDiscoverySpi() {
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
@@ -86,7 +84,7 @@ public class RestProcessorHangTest extends
GridCommonAbstractTest {
// No-op.
}
- super.writeToSocket(msg, sock, 255, timeout);
+ super.writeToSocket(sock, 255, timeout);
}
};
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/communication/tcp/IgniteTcpCommunicationConnectOnInitTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/communication/tcp/IgniteTcpCommunicationConnectOnInitTest.java
index 5955faaf9a9..ba7507bfb64 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/communication/tcp/IgniteTcpCommunicationConnectOnInitTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/communication/tcp/IgniteTcpCommunicationConnectOnInitTest.java
@@ -202,13 +202,12 @@ public class IgniteTcpCommunicationConnectOnInitTest
extends GridCommonAbstractT
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
awaitLatch();
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/** {@inheritDoc} */
@@ -224,14 +223,13 @@ public class IgniteTcpCommunicationConnectOnInitTest
extends GridCommonAbstractT
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
) throws IOException, IgniteCheckedException {
awaitLatch();
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
/**
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/BlockTcpDiscoverySpi.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/BlockTcpDiscoverySpi.java
index a74f2f2e7fd..5c6060fa60f 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/BlockTcpDiscoverySpi.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/BlockTcpDiscoverySpi.java
@@ -67,14 +67,17 @@ public class BlockTcpDiscoverySpi extends
TestTcpDiscoverySpi {
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
- if (spiCtx != null)
- apply(spiCtx.localNode(), msg);
+ if (spiCtx != null) {
+ TcpDiscoveryAbstractMessage msg = decodeMessage(this, data);
+
+ if (msg != null)
+ apply(spiCtx.localNode(), msg);
+ }
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/** {@inheritDoc} */
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/MultiDataCenterSplitTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/MultiDataCenterSplitTest.java
index 73acc8b9a37..96bff0e22b0 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/MultiDataCenterSplitTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/MultiDataCenterSplitTest.java
@@ -351,11 +351,10 @@ public class MultiDataCenterSplitTest extends
GridCommonAbstractTest {
}
/** {@inheritDoc} */
- @Override protected void writeToSocket(Socket sock, @Nullable
TcpDiscoveryAbstractMessage msg, byte[] data,
- long timeout) throws IOException, IgniteCheckedException {
+ @Override protected void writeToSocket(Socket sock, byte[] data, long
timeout) throws IOException, IgniteCheckedException {
tryToBlock(sock, data, timeout);
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/** */
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/ReceivedMessagesTracker.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/ReceivedMessagesTracker.java
new file mode 100644
index 00000000000..f7b93eba19f
--- /dev/null
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/ReceivedMessagesTracker.java
@@ -0,0 +1,45 @@
+/*
+ * 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.spi.discovery.tcp;
+
+import java.net.Socket;
+import java.util.Collections;
+import java.util.Map;
+import java.util.WeakHashMap;
+import org.apache.ignite.plugin.extensions.communication.Message;
+import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryAbstractMessage;
+import org.jetbrains.annotations.Nullable;
+
+/** */
+public class ReceivedMessagesTracker {
+ /** */
+ private final Map<Socket, TcpDiscoveryAbstractMessage> msgs =
Collections.synchronizedMap(new WeakHashMap<>());
+
+ /** */
+ public <T extends Message> T track(TcpDiscoveryIoSession ses, T msg) {
+ if (msg instanceof TcpDiscoveryAbstractMessage)
+ msgs.put(ses.socket(), (TcpDiscoveryAbstractMessage)msg);
+
+ return msg;
+ }
+
+ /** */
+ public @Nullable TcpDiscoveryAbstractMessage lastFor(Socket sock) {
+ return msgs.get(sock);
+ }
+}
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpClientDiscoverySpiFailureTimeoutSelfTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpClientDiscoverySpiFailureTimeoutSelfTest.java
index 49ba4f26982..3fd9b589bc9 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpClientDiscoverySpiFailureTimeoutSelfTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpClientDiscoverySpiFailureTimeoutSelfTest.java
@@ -447,14 +447,12 @@ public class TcpClientDiscoverySpiFailureTimeoutSelfTest
extends TcpClientDiscov
/** */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
if (writeToSocketDelay > 0) {
try {
- U.dumpStack(log, "Before sleep [msg=" + msg +
- ", arrLen=" + (data != null ? data.length : "n/a") +
']');
+ U.dumpStack(log, "Before sleep [arrLen=" + (data != null ?
data.length : "n/a") + ']');
Thread.sleep(writeToSocketDelay);
}
@@ -464,7 +462,7 @@ public class TcpClientDiscoverySpiFailureTimeoutSelfTest
extends TcpClientDiscov
}
if (sock.getSoTimeout() >= writeToSocketDelay)
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
else
throw new SocketTimeoutException("Write to socket delay
timeout exception.");
}
@@ -494,14 +492,13 @@ public class TcpClientDiscoverySpiFailureTimeoutSelfTest
extends TcpClientDiscov
/** */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
) throws IOException, IgniteCheckedException {
if (writeToSocketDelay > 0) {
try {
- U.dumpStack(log, "Before sleep [msg=" + msg + ']');
+ U.dumpStack(log, "Before sleep [res=" + res + ']');
Thread.sleep(writeToSocketDelay);
}
@@ -511,7 +508,7 @@ public class TcpClientDiscoverySpiFailureTimeoutSelfTest
extends TcpClientDiscov
}
if (sock.getSoTimeout() >= writeToSocketDelay)
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
else
throw new SocketTimeoutException("Write to socket delay
timeout exception.");
}
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpClientDiscoverySpiSelfTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpClientDiscoverySpiSelfTest.java
index 8bc4664977a..f8a0410bcfa 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpClientDiscoverySpiSelfTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpClientDiscoverySpiSelfTest.java
@@ -58,6 +58,7 @@ import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.lang.IgniteBiPredicate;
import org.apache.ignite.lang.IgniteInClosure;
import org.apache.ignite.lang.IgnitePredicate;
+import org.apache.ignite.plugin.extensions.communication.Message;
import org.apache.ignite.plugin.segmentation.SegmentationPolicy;
import org.apache.ignite.resources.IgniteInstanceResource;
import org.apache.ignite.spi.IgniteSpiException;
@@ -84,6 +85,7 @@ import static
org.apache.ignite.events.EventType.EVT_NODE_FAILED;
import static org.apache.ignite.events.EventType.EVT_NODE_JOINED;
import static org.apache.ignite.events.EventType.EVT_NODE_LEFT;
import static org.apache.ignite.events.EventType.EVT_NODE_SEGMENTED;
+import static
org.apache.ignite.spi.discovery.tcp.TestTcpDiscoverySpi.decodeMessage;
import static org.apache.ignite.testframework.GridTestUtils.noop;
/**
@@ -2505,6 +2507,9 @@ public class TcpClientDiscoverySpiSelfTest extends
GridCommonAbstractTest {
/** */
private volatile boolean skipNodeAdded;
+ /** */
+ private final ReceivedMessagesTracker msgTracker = new
ReceivedMessagesTracker();
+
/**
* @param lock Lock.
*/
@@ -2588,18 +2593,19 @@ public class TcpClientDiscoverySpiSelfTest extends
GridCommonAbstractTest {
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
- byte[] msgBytes,
+ byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
waitFor(writeLock);
- if (!onMessage(sock, msg))
+ TcpDiscoveryAbstractMessage msg = decodeMessage(this, data);
+
+ if (msg != null && !onMessage(sock, msg))
return;
- super.writeToSocket(sock, msg, msgBytes, timeout);
+ super.writeToSocket(sock, data, timeout);
- if (afterWrite != null)
+ if (msg != null && afterWrite != null)
afterWrite.apply(msg, sock);
}
@@ -2674,13 +2680,22 @@ public class TcpClientDiscoverySpiSelfTest extends
GridCommonAbstractTest {
t.resume();
}
+ /** {@inheritDoc} */
+ @Override protected <T extends Message> T readMessage(
+ TcpDiscoveryIoSession ses,
+ long timeout
+ ) throws IOException, IgniteCheckedException {
+ return msgTracker.track(ses, super.readMessage(ses, timeout));
+ }
+
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
) throws IOException, IgniteCheckedException {
+ TcpDiscoveryAbstractMessage msg = msgTracker.lastFor(sock);
+
if (delayJoinAckFor != null && msg instanceof
TcpDiscoveryJoinRequestMessage) {
TcpDiscoveryJoinRequestMessage msg0 =
(TcpDiscoveryJoinRequestMessage)msg;
@@ -2698,7 +2713,7 @@ public class TcpClientDiscoverySpiSelfTest extends
GridCommonAbstractTest {
}
}
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
/** {@inheritDoc} */
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryCoordinatorFailureTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryCoordinatorFailureTest.java
index 2421a74c181..43633ef5474 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryCoordinatorFailureTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryCoordinatorFailureTest.java
@@ -285,22 +285,6 @@ public class TcpDiscoveryCoordinatorFailureTest extends
GridCommonAbstractTest {
throw new IgniteCheckedException("Failed to wait for
NodeAddFinishedMessage");
}
- /** {@inheritDoc} */
- @Override protected void writeToSocket(
- Socket sock,
- TcpDiscoveryAbstractMessage msg,
- byte[] data,
- long timeout
- ) throws IOException, IgniteCheckedException {
- if (isDrop(msg)) {
- // Replace logic routine message with a stub to update
last-sent-time to avoid segmentation on
- // connRecoveryTimeout.
- msg = new TcpDiscoveryConnectionCheckMessage(locNode);
- }
-
- super.writeToSocket(sock, msg, data, timeout);
- }
-
/** {@inheritDoc} */
@Override protected void writeMessage(
TcpDiscoveryIoSession ses,
@@ -316,22 +300,6 @@ public class TcpDiscoveryCoordinatorFailureTest extends
GridCommonAbstractTest {
super.writeMessage(ses, msg, timeout);
}
- /** {@inheritDoc} */
- @Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
- Socket sock,
- int res,
- long timeout
- ) throws IOException, IgniteCheckedException {
- if (isDrop(msg)) {
- // Replace logic routine message with a stub to update
last-sent-time to avoid segmentation on
- // connRecoveryTimeout.
- msg = new TcpDiscoveryConnectionCheckMessage(locNode);
- }
-
- super.writeToSocket(msg, sock, res, timeout);
- }
-
/**
*
*/
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryFailedJoinTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryFailedJoinTest.java
index 9b6b8e4aa2f..0b5eb6891e7 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryFailedJoinTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryFailedJoinTest.java
@@ -197,12 +197,11 @@ public class TcpDiscoveryFailedJoinTest extends
GridCommonAbstractTest {
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
if (sock.getPort() != FAIL_PORT)
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/** {@inheritDoc} */
@@ -214,13 +213,12 @@ public class TcpDiscoveryFailedJoinTest extends
GridCommonAbstractTest {
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
) throws IOException, IgniteCheckedException {
if (sock.getPort() != FAIL_PORT)
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
}
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryNetworkIssuesTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryNetworkIssuesTest.java
index f416144d49b..37496f5d3b2 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryNetworkIssuesTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryNetworkIssuesTest.java
@@ -616,7 +616,6 @@ public class TcpDiscoveryNetworkIssuesTest extends
GridCommonAbstractTest {
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
@@ -624,7 +623,7 @@ public class TcpDiscoveryNetworkIssuesTest extends
GridCommonAbstractTest {
if (dropMsg(sock))
return;
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
/** {@inheritDoc} */
@@ -648,14 +647,13 @@ public class TcpDiscoveryNetworkIssuesTest extends
GridCommonAbstractTest {
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
if (dropMsg(sock))
return;
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/**
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryPendingMessageDeliveryTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryPendingMessageDeliveryTest.java
index 503d1198e5f..ec2f0c30ab2 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryPendingMessageDeliveryTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryPendingMessageDeliveryTest.java
@@ -270,12 +270,11 @@ public class TcpDiscoveryPendingMessageDeliveryTest
extends GridCommonAbstractTe
/** {@inheritDoc} */
@Override protected void writeToSocket(
Socket sock,
- TcpDiscoveryAbstractMessage msg,
byte[] data,
long timeout
) throws IOException, IgniteCheckedException {
if (!blockMsgs)
- super.writeToSocket(sock, msg, data, timeout);
+ super.writeToSocket(sock, data, timeout);
}
/** {@inheritDoc} */
@@ -287,13 +286,12 @@ public class TcpDiscoveryPendingMessageDeliveryTest
extends GridCommonAbstractTe
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
) throws IOException, IgniteCheckedException {
if (!blockMsgs)
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
}
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpiFailureTimeoutSelfTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpiFailureTimeoutSelfTest.java
index b7d34e25441..c71c26b0540 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpiFailureTimeoutSelfTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpiFailureTimeoutSelfTest.java
@@ -23,6 +23,7 @@ import java.net.Socket;
import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.configuration.IgniteConfiguration;
+import org.apache.ignite.plugin.extensions.communication.Message;
import org.apache.ignite.spi.IgniteSpiOperationTimeoutException;
import org.apache.ignite.spi.IgniteSpiOperationTimeoutHelper;
import org.apache.ignite.spi.discovery.AbstractDiscoverySelfTest;
@@ -261,7 +262,6 @@ public class TcpDiscoverySpiFailureTimeoutSelfTest extends
AbstractDiscoverySelf
/** */
private volatile IgniteSpiOperationTimeoutException err;
-
/** {@inheritDoc} */
@Override protected Socket openSocket(
Socket sock,
@@ -334,16 +334,16 @@ public class TcpDiscoverySpiFailureTimeoutSelfTest
extends AbstractDiscoverySelf
}
/** {@inheritDoc} */
- @Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
- Socket sock,
- int res,
+ @Override protected <T extends Message> T readMessage(
+ TcpDiscoveryIoSession ses,
long timeout
) throws IOException, IgniteCheckedException {
+ T msg = super.readMessage(ses, timeout);
+
if (cntConnCheckMsg && msg instanceof
TcpDiscoveryConnectionCheckMessage)
connCheckStatusMsgCntReceived++;
- super.writeToSocket(msg, sock, res, timeout);
+ return msg;
}
/**
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpiReconnectDelayTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpiReconnectDelayTest.java
index 4191afddc4e..21c7207c6b0 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpiReconnectDelayTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoverySpiReconnectDelayTest.java
@@ -31,6 +31,7 @@ import org.apache.ignite.events.Event;
import org.apache.ignite.internal.util.typedef.G;
import org.apache.ignite.lang.IgniteBiPredicate;
import org.apache.ignite.lang.IgnitePredicate;
+import org.apache.ignite.plugin.extensions.communication.Message;
import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryAbstractMessage;
import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryClientReconnectMessage;
import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryJoinRequestMessage;
@@ -398,6 +399,9 @@ public class TcpDiscoverySpiReconnectDelayTest extends
GridCommonAbstractTest {
/** */
private final AtomicInteger failReconReq = new AtomicInteger();
+ /** */
+ private final ReceivedMessagesTracker msgTracker = new
ReceivedMessagesTracker();
+
/** {@inheritDoc} */
@Override protected void writeMessage(TcpDiscoveryIoSession ses,
TcpDiscoveryAbstractMessage msg,
long timeout) throws IOException, IgniteCheckedException {
@@ -408,18 +412,24 @@ public class TcpDiscoverySpiReconnectDelayTest extends
GridCommonAbstractTest {
super.writeMessage(ses, msg, timeout);
}
+ /** {@inheritDoc} */
+ @Override protected <T extends Message> T readMessage(
+ TcpDiscoveryIoSession ses,
+ long timeout
+ ) throws IOException, IgniteCheckedException {
+ return msgTracker.track(ses, super.readMessage(ses, timeout));
+ }
+
/** {@inheritDoc} */
@Override protected void writeToSocket(
- TcpDiscoveryAbstractMessage msg,
Socket sock,
int res,
long timeout
) throws IOException, IgniteCheckedException {
-
- if (msg instanceof TcpDiscoveryJoinRequestMessage &&
failJoinReqRes.getAndDecrement() > 0)
+ if (msgTracker.lastFor(sock) instanceof
TcpDiscoveryJoinRequestMessage && failJoinReqRes.getAndDecrement() > 0)
res = RES_WAIT;
- super.writeToSocket(msg, sock, res, timeout);
+ super.writeToSocket(sock, res, timeout);
}
/**
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TestTcpDiscoverySpi.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TestTcpDiscoverySpi.java
index e1133dfd06c..866bbc84217 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TestTcpDiscoverySpi.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/TestTcpDiscoverySpi.java
@@ -17,12 +17,19 @@
package org.apache.ignite.spi.discovery.tcp;
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
import java.io.IOException;
+import java.io.InputStream;
+import java.io.OutputStream;
+import java.net.Socket;
+import java.util.Arrays;
import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.IgniteException;
import org.apache.ignite.internal.CoreMessagesProvider;
import
org.apache.ignite.internal.managers.communication.IgniteMessageFactoryImpl;
import
org.apache.ignite.internal.managers.discovery.IgniteDiscoverySpiInternalListener;
+import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.plugin.extensions.communication.MessageFactory;
import
org.apache.ignite.plugin.extensions.communication.MessageFactoryProvider;
import org.apache.ignite.spi.discovery.DiscoverySpiCustomMessage;
@@ -145,4 +152,27 @@ public class TestTcpDiscoverySpi extends TcpDiscoverySpi
implements IgniteDiscov
@Override public MessageFactory messageFactory() {
return msgFactory != null ? msgFactory : super.messageFactory();
}
+
+ /** */
+ public static @Nullable TcpDiscoveryAbstractMessage
decodeMessage(TcpDiscoverySpi spi, byte[] data) {
+ if (Arrays.equals(U.IGNITE_HEADER, data))
+ return null;
+
+ Socket dataSock = new Socket() {
+ @Override public InputStream getInputStream() {
+ return new ByteArrayInputStream(data);
+ }
+
+ @Override public OutputStream getOutputStream() {
+ return new ByteArrayOutputStream();
+ }
+ };
+
+ try (dataSock) {
+ return new TcpDiscoveryIoSession(dataSock, spi).readMessage();
+ }
+ catch (Exception e) {
+ throw new IgniteException("Failed to decode a message", e);
+ }
+ }
}