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 63457695229 IGNITE-29037 Refactored Message Factory usage in TCP
Discovery implementations (#13551)
63457695229 is described below
commit 63457695229c9c4a967ef9eb62cf925e9c491140
Author: Mikhail Petrov <[email protected]>
AuthorDate: Fri Sep 4 16:53:40 2026 +0300
IGNITE-29037 Refactored Message Factory usage in TCP Discovery
implementations (#13551)
---
.../ignite/spi/discovery/tcp/ClientImpl.java | 32 +++------
.../ignite/spi/discovery/tcp/ServerImpl.java | 80 +++++++++-------------
.../ignite/spi/discovery/tcp/TcpDiscoveryImpl.java | 19 ++---
.../spi/discovery/tcp/TcpDiscoveryIoSession.java | 40 ++++++-----
.../tcp/TcpDiscoveryMessageSerializer.java | 9 +--
.../ignite/spi/discovery/tcp/TcpDiscoverySpi.java | 21 +-----
.../cache/CacheMetricsCacheSizeTest.java | 3 +-
.../apache/ignite/spi/MessagesPluginProvider.java | 13 ----
.../spi/discovery/tcp/BlockTcpDiscoverySpi.java | 2 +-
.../tcp/DiscoveryUnmarshalVulnerabilityTest.java | 2 +-
.../tcp/TcpClientDiscoverySpiSelfTest.java | 2 +-
.../spi/discovery/tcp/TestTcpDiscoverySpi.java | 50 +-------------
12 files changed, 85 insertions(+), 188 deletions(-)
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
index 621a234ee54..dd1aed7a2a2 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/ClientImpl.java
@@ -47,7 +47,6 @@ import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.atomic.AtomicReference;
import javax.net.ssl.SSLException;
-import org.apache.ignite.Ignite;
import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.IgniteClientDisconnectedException;
import org.apache.ignite.IgniteException;
@@ -60,7 +59,6 @@ import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.configuration.IgniteConfiguration;
import org.apache.ignite.failure.FailureContext;
import org.apache.ignite.internal.IgniteClientDisconnectedCheckedException;
-import org.apache.ignite.internal.IgniteEx;
import org.apache.ignite.internal.IgniteInterruptedCheckedException;
import org.apache.ignite.internal.IgniteNodeAttributes;
import
org.apache.ignite.internal.managers.discovery.DiscoveryServerOnlyCustomMessage;
@@ -73,7 +71,6 @@ import org.apache.ignite.internal.util.typedef.X;
import org.apache.ignite.internal.util.typedef.internal.LT;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.internal.util.worker.GridWorker;
-import org.apache.ignite.internal.worker.WorkersRegistry;
import org.apache.ignite.lang.IgniteInClosure;
import org.apache.ignite.lang.IgniteUuid;
import org.apache.ignite.spi.IgniteSpiAdapter;
@@ -208,10 +205,9 @@ class ClientImpl extends TcpDiscoveryImpl {
ClientImpl(TcpDiscoverySpi adapter) {
super(adapter);
- String instanceName = adapter.ignite() == null ||
adapter.ignite().name() == null
- ? "client-node" : adapter.ignite().name();
+ String instanceName = ctx.igniteInstanceName();
- executorSrvc = newSingleThreadScheduledExecutor("tcp-discovery-exec",
instanceName);
+ executorSrvc = newSingleThreadScheduledExecutor("tcp-discovery-exec",
instanceName == null ? "client-node" : instanceName);
}
/** {@inheritDoc} */
@@ -1052,13 +1048,6 @@ class ClientImpl extends TcpDiscoveryImpl {
joinLatch.countDown();
}
- /** */
- private WorkersRegistry getWorkersRegistry() {
- Ignite ignite = spi.ignite();
-
- return ignite instanceof IgniteEx ?
((IgniteEx)ignite).context().workersRegistry() : null;
- }
-
/** */
private Collection<ClusterNode> remoteVisibleNodes() {
return U.arrayList(rmtNodes.values(),
TcpDiscoveryNodesRing.VISIBLE_NODES);
@@ -1101,7 +1090,7 @@ class ClientImpl extends TcpDiscoveryImpl {
/**
*/
SocketReader() {
- super(spi.ignite().name(), "tcp-client-disco-sock-reader-[]", log);
+ super(ctx.igniteInstanceName(), "tcp-client-disco-sock-reader-[]",
log);
}
/**
@@ -1274,7 +1263,7 @@ class ClientImpl extends TcpDiscoveryImpl {
*
*/
SocketWriter() {
- super(spi.ignite().name(), "tcp-client-disco-sock-writer", log);
+ super(ctx.igniteInstanceName(), "tcp-client-disco-sock-writer",
log);
sockTimeout = spi.failureDetectionTimeoutEnabled() ?
spi.failureDetectionTimeout() :
spi.getSocketTimeout();
@@ -1524,7 +1513,7 @@ class ClientImpl extends TcpDiscoveryImpl {
* @param prevAddr Address of the node, that this client was
previously connected to.
*/
protected Reconnector(boolean join, InetSocketAddress prevAddr) {
- super(spi.ignite().name(), "tcp-client-disco-reconnector", log);
+ super(ctx.igniteInstanceName(), "tcp-client-disco-reconnector",
log);
this.join = join;
this.prevAddr = prevAddr;
@@ -1689,7 +1678,7 @@ class ClientImpl extends TcpDiscoveryImpl {
* @param log Logger.
*/
private MessageWorker(IgniteLogger log) {
- super(spi.ignite().name(), "tcp-client-disco-msg-worker", log,
getWorkersRegistry());
+ super(ctx.igniteInstanceName(), "tcp-client-disco-msg-worker",
log, ctx.workersRegistry());
}
/** {@inheritDoc} */
@@ -1957,8 +1946,7 @@ class ClientImpl extends TcpDiscoveryImpl {
Thread.currentThread().interrupt();
}
catch (Throwable t) {
- if (spi.ignite() instanceof IgniteEx)
- ((IgniteEx)spi.ignite()).context().failure().process(new
FailureContext(CRITICAL_ERROR, t));
+ ctx.failure().process(new FailureContext(CRITICAL_ERROR, t));
}
finally {
TcpDiscoveryIoSession ses = this.currSes;
@@ -2210,7 +2198,7 @@ class ClientImpl extends TcpDiscoveryImpl {
if (joining())
delayDiscoData.add(dataPacket);
else
- spi.onExchange(dataPacket,
U.resolveClassLoader(spi.ignite().configuration()));
+ spi.onExchange(dataPacket,
U.resolveClassLoader(ctx.config()));
}
}
}
@@ -2244,11 +2232,11 @@ class ClientImpl extends TcpDiscoveryImpl {
DiscoveryDataPacket dataContainer = msg.clientDiscoData();
if (dataContainer != null)
- spi.onExchange(dataContainer,
U.resolveClassLoader(spi.ignite().configuration()));
+ spi.onExchange(dataContainer,
U.resolveClassLoader(ctx.config()));
if (!delayDiscoData.isEmpty()) {
for (DiscoveryDataPacket data : delayDiscoData)
- spi.onExchange(data,
U.resolveClassLoader(spi.ignite().configuration()));
+ spi.onExchange(data,
U.resolveClassLoader(ctx.config()));
delayDiscoData.clear();
}
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 e8a9e5c2394..5974e617807 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
@@ -76,7 +76,6 @@ import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.events.NodeValidationFailedEvent;
import org.apache.ignite.failure.FailureContext;
import org.apache.ignite.internal.ClusterMetricsSnapshot;
-import org.apache.ignite.internal.IgniteEx;
import org.apache.ignite.internal.IgniteFutureTimeoutCheckedException;
import org.apache.ignite.internal.IgniteInterruptedCheckedException;
import org.apache.ignite.internal.IgniteNodeAttributes;
@@ -107,7 +106,6 @@ import org.apache.ignite.internal.util.typedef.internal.S;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.internal.util.worker.GridWorker;
import org.apache.ignite.internal.util.worker.GridWorkerListener;
-import org.apache.ignite.internal.worker.WorkersRegistry;
import org.apache.ignite.lang.IgniteBiTuple;
import org.apache.ignite.lang.IgniteFuture;
import org.apache.ignite.lang.IgniteInClosure;
@@ -362,14 +360,14 @@ class ServerImpl extends TcpDiscoveryImpl {
super(adapter);
utilityPool = new IgniteThreadPoolExecutor("disco-pool",
- spi.ignite().name(),
+ ctx.igniteInstanceName(),
0,
utilityPoolSize,
2000,
new LinkedBlockingQueue<>());
List<DistributedBooleanProperty> props = newConnectionEnabledProperty(
- ((IgniteEx)spi.ignite()).context().internalSubscriptionProcessor(),
+ ctx.internalSubscriptionProcessor(),
log,
"ClientNode",
"ServerNode"
@@ -2249,7 +2247,7 @@ class ServerImpl extends TcpDiscoveryImpl {
* Constructor.
*/
private IpFinderCleaner() {
- super(spi.ignite().name(), "tcp-disco-ip-finder-cleaner", log);
+ super(ctx.igniteInstanceName(), "tcp-disco-ip-finder-cleaner",
log);
setPriority(spi.threadPri);
}
@@ -2458,11 +2456,6 @@ class ServerImpl extends TcpDiscoveryImpl {
node.setAttributes(attrs);
}
- /** */
- private static WorkersRegistry getWorkerRegistry(TcpDiscoverySpi spi) {
- return spi.ignite() instanceof IgniteEx ?
((IgniteEx)spi.ignite()).context().workersRegistry() : null;
- }
-
/**
* Upcasts collection type.
*
@@ -2866,7 +2859,7 @@ class ServerImpl extends TcpDiscoveryImpl {
// To address this, we use TcpDiscoveryMessageSerializer, which
includes some code copied from TcpDiscoveryIoSession
// and can be instantiated independently of any active session.
/** */
- private final TcpDiscoveryMessageSerializer clientMsgSer = new
TcpDiscoveryMessageSerializer(spi);
+ private final TcpDiscoveryMessageSerializer clientMsgSer = new
TcpDiscoveryMessageSerializer(ctx);
/** IO session. */
private TcpDiscoveryIoSession ses;
@@ -2902,7 +2895,7 @@ class ServerImpl extends TcpDiscoveryImpl {
/** */
protected RingMessageWorker(IgniteLogger log,
BlockingDeque<TcpDiscoveryAbstractMessage> queue) {
- super("tcp-disco-msg-worker-[]", log, 10, getWorkerRegistry(spi),
queue);
+ super("tcp-disco-msg-worker-[]", log, 10, ctx.workersRegistry(),
queue);
setBeforeEachPollAction(() -> {
updateHeartbeat();
@@ -3049,17 +3042,15 @@ class ServerImpl extends TcpDiscoveryImpl {
throw e;
}
finally {
- if (spi.ignite() instanceof IgniteEx) {
- if (err == null && !spi.isNodeStopping0() &&
spiStateCopy() != DISCONNECTING)
- err = new IllegalStateException("Worker " + name() + "
is terminated unexpectedly.");
+ if (err == null && !spi.isNodeStopping0() && spiStateCopy() !=
DISCONNECTING)
+ err = new IllegalStateException("Worker " + name() + " is
terminated unexpectedly.");
- FailureProcessor failure =
((IgniteEx)spi.ignite()).context().failure();
+ FailureProcessor failure = ctx.failure();
- if (err instanceof OutOfMemoryError)
- failure.process(new FailureContext(CRITICAL_ERROR,
err));
- else if (err != null)
- failure.process(new
FailureContext(SYSTEM_WORKER_TERMINATION, err));
- }
+ if (err instanceof OutOfMemoryError)
+ failure.process(new FailureContext(CRITICAL_ERROR, err));
+ else if (err != null)
+ failure.process(new
FailureContext(SYSTEM_WORKER_TERMINATION, err));
}
}
@@ -4591,8 +4582,7 @@ class ServerImpl extends TcpDiscoveryImpl {
DiscoveryDataPacket packet = req.gridDiscoveryData();
try {
- DiscoveryDataBag dataBag =
packet.bagWithJoiningNodeData(spi.ignite().log(),
- spi.ignite().configuration().isClientMode());
+ DiscoveryDataBag dataBag = packet.bagWithJoiningNodeData(log,
ctx.config().isClientMode());
return spi.getSpiContext().validateNode(req.node(), dataBag);
}
@@ -4936,7 +4926,7 @@ class ServerImpl extends TcpDiscoveryImpl {
if (dataPacket.hasJoiningNodeData()) {
if (spiState == CONNECTED) {
// Node already connected to the cluster can apply
joining nodes' disco data immediately
- spi.onExchange(dataPacket,
U.resolveClassLoader(spi.ignite().configuration()));
+ spi.onExchange(dataPacket,
U.resolveClassLoader(ctx.config()));
spi.collectExchangeData(dataPacket);
}
@@ -5154,11 +5144,11 @@ class ServerImpl extends TcpDiscoveryImpl {
}
if (gridDiscoveryData != null)
- spi.onExchange(gridDiscoveryData,
U.resolveClassLoader(spi.ignite().configuration()));
+ spi.onExchange(gridDiscoveryData,
U.resolveClassLoader(ctx.config()));
if (joiningNodesDiscoDataList != null) {
for (DiscoveryDataPacket dataPacket :
joiningNodesDiscoDataList)
- spi.onExchange(dataPacket,
U.resolveClassLoader(spi.ignite().configuration()));
+ spi.onExchange(dataPacket,
U.resolveClassLoader(ctx.config()));
}
nullifyDiscoData();
@@ -6310,7 +6300,7 @@ class ServerImpl extends TcpDiscoveryImpl {
* @throws IgniteSpiException In case of error.
*/
TcpServer(IgniteLogger log) throws IgniteSpiException {
- super(spi.ignite().name(), "tcp-disco-srvr-[]", log,
getWorkerRegistry(spi));
+ super(ctx.igniteInstanceName(), "tcp-disco-srvr-[]", log,
ctx.workersRegistry());
int lastPort = spi.locPortRange == 0 ? spi.locPort : spi.locPort +
spi.locPortRange - 1;
@@ -6330,7 +6320,7 @@ class ServerImpl extends TcpDiscoveryImpl {
if (log.isInfoEnabled()) {
log.info("Successfully bound to TCP port [port=" +
port +
", localHost=" + spi.locHost +
- ", locNodeId=" +
spi.ignite().configuration().getNodeId() +
+ ", locNodeId=" + spi.cfgNodeId +
']');
}
@@ -6413,17 +6403,15 @@ class ServerImpl extends TcpDiscoveryImpl {
throw t;
}
finally {
- if (spi.ignite() instanceof IgniteEx) {
- if (err == null && !spi.isNodeStopping0() &&
spiStateCopy() != DISCONNECTING)
- err = new IllegalStateException("Worker " + name() + "
is terminated unexpectedly.");
+ if (err == null && !spi.isNodeStopping0() && spiStateCopy() !=
DISCONNECTING)
+ err = new IllegalStateException("Worker " + name() + " is
terminated unexpectedly.");
- FailureProcessor failure =
((IgniteEx)spi.ignite()).context().failure();
+ FailureProcessor failure = ctx.failure();
- if (err instanceof OutOfMemoryError)
- failure.process(new FailureContext(CRITICAL_ERROR,
err));
- else if (err != null)
- failure.process(new
FailureContext(SYSTEM_WORKER_TERMINATION, err));
- }
+ if (err instanceof OutOfMemoryError)
+ failure.process(new FailureContext(CRITICAL_ERROR, err));
+ else if (err != null)
+ failure.process(new
FailureContext(SYSTEM_WORKER_TERMINATION, err));
U.closeQuiet(srvrSock);
}
@@ -6458,11 +6446,11 @@ class ServerImpl extends TcpDiscoveryImpl {
* @param sock Socket to read data from.
*/
SocketReader(Socket sock) {
- super(spi.ignite().name(), "tcp-disco-sock-reader-[]", log);
+ super(ctx.igniteInstanceName(), "tcp-disco-sock-reader-[]", log);
this.sock = sock;
- ses = createSession(sock);
+ ses = new TcpDiscoveryIoSession(ctx, sock);
setPriority(spi.threadPri);
}
@@ -7085,11 +7073,7 @@ class ServerImpl extends TcpDiscoveryImpl {
}
}
catch (UnknownMessageException e) {
- if (spi.ignite() instanceof IgniteEx) {
- FailureProcessor failure =
((IgniteEx)spi.ignite()).context().failure();
-
- failure.process(new
FailureContext(SYSTEM_WORKER_TERMINATION, e));
- }
+ ctx.failure().process(new
FailureContext(SYSTEM_WORKER_TERMINATION, e));
}
finally {
if (clientMsgWrk != null) {
@@ -7461,7 +7445,7 @@ class ServerImpl extends TcpDiscoveryImpl {
* Constructor.
*/
StatisticsPrinter() {
- super(spi.ignite().name(), "tcp-disco-stats-printer", log);
+ super(ctx.igniteInstanceName(), "tcp-disco-stats-printer", log);
assert spi.statsPrintFreq > 0;
@@ -7534,7 +7518,7 @@ class ServerImpl extends TcpDiscoveryImpl {
this.ses = ses;
this.clientNodeId = clientNodeId;
- clientMsgSer = new TcpDiscoveryMessageSerializer(spi);
+ clientMsgSer = new TcpDiscoveryMessageSerializer(ctx);
lastMetricsUpdateMsgTimeNanos = System.nanoTime();
}
@@ -7878,7 +7862,7 @@ class ServerImpl extends TcpDiscoveryImpl {
@Nullable GridWorkerListener lsnr,
BlockingDeque<T> queue
) {
- super(spi.ignite().name(), name, log, lsnr);
+ super(ctx.igniteInstanceName(), name, log, lsnr);
this.queue = queue;
this.pollingTimeout = pollingTimeout;
@@ -8101,7 +8085,7 @@ class ServerImpl extends TcpDiscoveryImpl {
rmtDcPingPool = new IgniteThreadPoolExecutor(
"disco-remote-dc-ping-worker",
- spi.ignite().name(),
+ ctx.igniteInstanceName(),
pingRmtDcPoolSz,
pingRmtDcPoolSz,
0,
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryImpl.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryImpl.java
index b83e2130ae2..0f2324361d5 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryImpl.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryImpl.java
@@ -18,7 +18,6 @@
package org.apache.ignite.spi.discovery.tcp;
import java.net.InetSocketAddress;
-import java.net.Socket;
import java.time.Instant;
import java.time.ZoneId;
import java.time.format.DateTimeFormatter;
@@ -36,6 +35,7 @@ import org.apache.ignite.cache.CacheMetrics;
import org.apache.ignite.cluster.ClusterMetrics;
import org.apache.ignite.cluster.ClusterNode;
import org.apache.ignite.internal.ClusterMetricsSnapshot;
+import org.apache.ignite.internal.GridKernalContext;
import org.apache.ignite.internal.IgniteEx;
import org.apache.ignite.internal.processors.cache.CacheMetricsSnapshot;
import org.apache.ignite.internal.processors.cluster.CacheMetricsMessage;
@@ -87,6 +87,9 @@ abstract class TcpDiscoveryImpl {
/** */
protected final TcpDiscoverySpi spi;
+ /** */
+ protected final GridKernalContext ctx;
+
/** */
protected final IgniteLogger log;
@@ -145,7 +148,9 @@ abstract class TcpDiscoveryImpl {
log = spi.log;
- operationCtxDispatcher =
((IgniteEx)spi.ignite()).context().operationContextDispatcher();
+ ctx = ((IgniteEx)spi.ignite()).context();
+
+ operationCtxDispatcher = ctx.operationContextDispatcher();
}
/**
@@ -459,16 +464,6 @@ abstract class TcpDiscoveryImpl {
return res;
}
- /**
- * Instantiates IO session for exchanging discovery messages with remote
node.
- *
- * @param sock Socket to remote node.
- * @return IO session for writing and reading {@link
TcpDiscoveryAbstractMessage}.
- */
- TcpDiscoveryIoSession createSession(Socket sock) {
- return new TcpDiscoveryIoSession(sock, spi);
- }
-
/**
* @param msg Message.
* @return Message logger.
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryIoSession.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryIoSession.java
index 71a4995eeab..0251712fc95 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryIoSession.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryIoSession.java
@@ -35,7 +35,6 @@ import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.IgniteException;
import org.apache.ignite.IgniteLogger;
import org.apache.ignite.internal.GridKernalContext;
-import org.apache.ignite.internal.IgniteEx;
import org.apache.ignite.internal.direct.DirectMessageReader;
import org.apache.ignite.internal.direct.DirectMessageWriter;
import org.apache.ignite.internal.managers.communication.DiscoveryMarshalling;
@@ -46,6 +45,7 @@ import org.apache.ignite.internal.util.typedef.X;
import org.apache.ignite.internal.util.typedef.internal.U;
import org.apache.ignite.marshaller.jdk.JdkMarshaller;
import org.apache.ignite.plugin.extensions.communication.Message;
+import org.apache.ignite.plugin.extensions.communication.MessageFactory;
import org.apache.ignite.plugin.extensions.communication.MessageSerializer;
import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryAbstractMessage;
import org.jetbrains.annotations.NotNull;
@@ -70,7 +70,13 @@ public class TcpDiscoveryIoSession implements AutoCloseable {
private static final int MSG_BUFFER_SIZE = 100;
/** */
- private final TcpDiscoverySpi spi;
+ private final GridKernalContext ctx;
+
+ /** */
+ private final MessageFactory<?> msgFactory;
+
+ /** */
+ private final IgniteLogger log;
/** */
private final Socket sock;
@@ -96,19 +102,21 @@ public class TcpDiscoveryIoSession implements
AutoCloseable {
/**
* Creates a new discovery I/O session bound to the given socket.
*
+ * @param ctx Kernal context.
* @param sock Socket connected to a remote discovery node.
- * @param spi Discovery SPI instance owning this session.
* @throws IgniteException If an I/O error occurs while initializing
buffers.
*/
- TcpDiscoveryIoSession(Socket sock, TcpDiscoverySpi spi) {
+ TcpDiscoveryIoSession(GridKernalContext ctx, Socket sock) {
this.sock = sock;
- this.spi = spi;
+ this.ctx = ctx;
+ this.msgFactory = ctx.messageFactory();
+ this.log = ctx.log(getClass());
readBuf = ByteBuffer.allocate(MSG_BUFFER_SIZE);
writeBuf = ByteBuffer.allocate(MSG_BUFFER_SIZE);
- msgWriter = new DirectMessageWriter(spi.messageFactory());
- msgReader = new DirectMessageReader(spi.messageFactory(), null);
+ msgWriter = new DirectMessageWriter(msgFactory);
+ msgReader = new DirectMessageReader(msgFactory, null);
try {
int sendBufSize = sock.getSendBufferSize() > 0 ?
sock.getSendBufferSize() : DFLT_SOCK_BUFFER_SIZE;
@@ -178,7 +186,7 @@ public class TcpDiscoveryIoSession implements AutoCloseable
{
Message msg;
try {
- msg = spi.messageFactory().create(msgType);
+ msg = msgFactory.create(msgType);
}
catch (IgniteException e) {
detectSslAlert(b0, b1);
@@ -202,7 +210,7 @@ public class TcpDiscoveryIoSession implements AutoCloseable
{
readBuf.limit(read);
- finished = MessageSerialization.readFrom(spi.messageFactory(),
msg, msgReader);
+ finished = MessageSerialization.readFrom(msgFactory, msg,
msgReader);
// Server Discovery only sends next message to next Server
upon receiving a receipt for the previous one.
// This behaviour guarantees that we never read a next message
from the buffer right after the end of
@@ -218,9 +226,7 @@ public class TcpDiscoveryIoSession implements AutoCloseable
{
}
while (!finished);
- GridKernalContext kctx = ((IgniteEx)spi.ignite()).context();
-
- DiscoveryMarshalling.unmarshal(msg, kctx);
+ DiscoveryMarshalling.unmarshal(msg, ctx);
return (T)msg;
}
@@ -238,14 +244,14 @@ public class TcpDiscoveryIoSession implements
AutoCloseable {
/** @return SSL certificate this session is established with. {@code null}
if SSL is disabled or certificate validation failed. */
@Nullable Certificate[] extractCertificates() {
- if (!spi.isSslEnabled())
+ if (!(sock instanceof SSLSocket))
return null;
try {
return ((SSLSocket)sock).getSession().getPeerCertificates();
}
catch (SSLPeerUnverifiedException e) {
- U.error(spi.log, "Failed to extract discovery IO session
certificates", e);
+ U.error(log, "Failed to extract discovery IO session
certificates", e);
return null;
}
@@ -264,9 +270,7 @@ public class TcpDiscoveryIoSession implements AutoCloseable
{
* @throws IOException If serialization fails.
*/
void serializeMessage(Message m, OutputStream out) throws IOException,
IgniteCheckedException {
- GridKernalContext kctx = ((IgniteEx)spi.ignite()).context();
-
- DiscoveryMarshalling.marshal(m, kctx, null);
+ DiscoveryMarshalling.marshal(m, ctx, null);
msgWriter.reset();
msgWriter.setBuffer(writeBuf);
@@ -277,7 +281,7 @@ public class TcpDiscoveryIoSession implements AutoCloseable
{
// Should be cleared before first operation.
writeBuf.clear();
- finished = MessageSerialization.writeTo(spi.messageFactory(), m,
msgWriter);
+ finished = MessageSerialization.writeTo(msgFactory, m, msgWriter);
out.write(writeBuf.array(), 0, writeBuf.position());
}
diff --git
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryMessageSerializer.java
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryMessageSerializer.java
index 8c871f9e2eb..ec7cdc569f0 100644
---
a/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryMessageSerializer.java
+++
b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryMessageSerializer.java
@@ -22,6 +22,7 @@ import java.io.InputStream;
import java.io.OutputStream;
import java.net.Socket;
import org.apache.ignite.IgniteCheckedException;
+import org.apache.ignite.internal.GridKernalContext;
import org.apache.ignite.plugin.extensions.communication.Message;
import org.apache.ignite.plugin.extensions.communication.MessageSerializer;
import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryAbstractMessage;
@@ -35,10 +36,10 @@ import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryAbstractMessage;
*/
class TcpDiscoveryMessageSerializer extends TcpDiscoveryIoSession {
/**
- * @param spi Discovery SPI instance.
+ * @param ctx Kernal context.
*/
- public TcpDiscoveryMessageSerializer(TcpDiscoverySpi spi) {
- super(new Socket() {
+ public TcpDiscoveryMessageSerializer(GridKernalContext ctx) {
+ super(ctx, new Socket() {
@Override public OutputStream getOutputStream() throws IOException
{
return null;
}
@@ -46,7 +47,7 @@ class TcpDiscoveryMessageSerializer extends
TcpDiscoveryIoSession {
@Override public InputStream getInputStream() throws IOException {
return null;
}
- }, spi);
+ });
}
/**
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 c4c74289af3..7b8490ded83 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
@@ -57,7 +57,6 @@ import org.apache.ignite.internal.IgniteEx;
import org.apache.ignite.internal.IgniteInterruptedCheckedException;
import
org.apache.ignite.internal.managers.communication.UnknownMessageException;
import org.apache.ignite.internal.managers.discovery.IgniteDiscoverySpi;
-import org.apache.ignite.internal.processors.failure.FailureProcessor;
import org.apache.ignite.internal.processors.metric.MetricRegistryImpl;
import
org.apache.ignite.internal.processors.rollingupgrade.feature.IgniteComponentFeatureSet;
import
org.apache.ignite.internal.processors.rollingupgrade.feature.IgniteNodeFeatureSet;
@@ -74,7 +73,6 @@ import org.apache.ignite.lang.IgniteProductVersion;
import org.apache.ignite.lang.IgniteUuid;
import org.apache.ignite.marshaller.Marshaller;
import org.apache.ignite.plugin.extensions.communication.Message;
-import org.apache.ignite.plugin.extensions.communication.MessageFactory;
import org.apache.ignite.resources.IgniteInstanceResource;
import org.apache.ignite.resources.LoggerResource;
import org.apache.ignite.spi.IgniteSpiAdapter;
@@ -463,10 +461,6 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
@GridToStringExclude
protected IgniteSpiContext spiCtx;
- /** Discovery messages factory. */
- @GridToStringExclude
- private MessageFactory msgFactory;
-
/** For test purposes. */
private boolean skipAddrsRandomization = false;
@@ -600,8 +594,6 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
setAddressResolver(ignite.configuration().getAddressResolver());
marsh =
((IgniteEx)ignite).context().marshallerContext().jdkMarshaller();
-
- msgFactory = ((IgniteEx)ignite).context().messageFactory();
}
}
@@ -1122,11 +1114,6 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
locNodeVer = ver;
}
- /** @return Discovery messages factory. */
- public MessageFactory messageFactory() {
- return msgFactory;
- }
-
/**
* Gets ID of the local node.
*
@@ -1624,7 +1611,7 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
sock.connect(resolved,
(int)timeoutHelper.nextTimeoutChunk(sockTimeout));
- TcpDiscoveryIoSession ses = new TcpDiscoveryIoSession(sock, this);
+ TcpDiscoveryIoSession ses = new
TcpDiscoveryIoSession(ignite.context(), sock);
write(ses, U.IGNITE_HEADER,
timeoutHelper.nextTimeoutChunk(sockTimeout));
@@ -2131,11 +2118,7 @@ public class TcpDiscoverySpi extends IgniteSpiAdapter
implements IgniteDiscovery
dataBag = dataPacket.bagWithJoiningNodeData(ignite.log(),
ignite.configuration().isClientMode());
}
catch (IgniteCheckedException e) {
- if (ignite() instanceof IgniteEx) {
- FailureProcessor failure =
((IgniteEx)ignite()).context().failure();
-
- failure.process(new FailureContext(CRITICAL_ERROR, e));
- }
+ ignite.context().failure().process(new
FailureContext(CRITICAL_ERROR, e));
throw new IgniteException(e);
}
diff --git
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheMetricsCacheSizeTest.java
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheMetricsCacheSizeTest.java
index 9aed75febee..a4a0db43330 100644
---
a/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheMetricsCacheSizeTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/internal/processors/cache/CacheMetricsCacheSizeTest.java
@@ -35,7 +35,6 @@ import org.apache.ignite.internal.direct.DirectMessageWriter;
import org.apache.ignite.internal.processors.cluster.CacheMetricsMessage;
import org.apache.ignite.internal.util.nio.MessageSerialization;
import org.apache.ignite.plugin.extensions.communication.MessageFactory;
-import org.apache.ignite.spi.discovery.tcp.TcpDiscoverySpi;
import
org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryMetricsUpdateMessage;
import org.apache.ignite.testframework.junits.common.GridCommonAbstractTest;
import org.junit.Test;
@@ -105,7 +104,7 @@ public class CacheMetricsCacheSizeTest extends
GridCommonAbstractTest {
msg.addServerMetrics(srvrId, new ClusterMetricsSnapshot());
msg.addServerCacheMetrics(srvrId, cacheMetrics);
- MessageFactory msgFactory =
((TcpDiscoverySpi)grid(0).context().discovery().getInjectedDiscoverySpi()).messageFactory();
+ MessageFactory msgFactory = grid(0).context().messageFactory();
// First time we write initial message type which is not read by the
reader because the message type is known.
// We have to skip this header at the further message reading.
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/MessagesPluginProvider.java
b/modules/core/src/test/java/org/apache/ignite/spi/MessagesPluginProvider.java
index 1198584d229..0090d11cc3c 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/MessagesPluginProvider.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/MessagesPluginProvider.java
@@ -17,7 +17,6 @@
package org.apache.ignite.spi;
-import org.apache.ignite.IgniteCheckedException;
import org.apache.ignite.internal.CoreMessagesProvider;
import org.apache.ignite.plugin.AbstractTestPluginProvider;
import org.apache.ignite.plugin.ExtensionRegistry;
@@ -25,8 +24,6 @@ import org.apache.ignite.plugin.PluginContext;
import org.apache.ignite.plugin.extensions.communication.Message;
import
org.apache.ignite.plugin.extensions.communication.MessageFactoryProvider;
import org.apache.ignite.plugin.extensions.communication.MessageMarshaller;
-import org.apache.ignite.spi.discovery.DiscoverySpi;
-import org.apache.ignite.spi.discovery.tcp.TestTcpDiscoverySpi;
import org.jetbrains.annotations.Nullable;
import static org.apache.ignite.testframework.GridTestUtils.loadMarshaller;
@@ -78,14 +75,4 @@ public class MessagesPluginProvider extends
AbstractTestPluginProvider {
// Register messages into the communication protocol.
registry.registerExtension(MessageFactoryProvider.class,
msgFactoryProvider);
}
-
- /** {@inheritDoc} */
- @Override public void start(PluginContext ctx) throws
IgniteCheckedException {
- DiscoverySpi discoSpi = ctx.igniteConfiguration().getDiscoverySpi();
-
- if (discoSpi instanceof TestTcpDiscoverySpi testDiscoSpi) {
- // Register messages into the discovery protocol.
- testDiscoSpi.messageFactory(msgFactoryProvider);
- }
- }
}
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 52e81b05df9..4904df46a40 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
@@ -70,7 +70,7 @@ public class BlockTcpDiscoverySpi extends TestTcpDiscoverySpi
{
long timeout
) throws IOException, IgniteCheckedException {
if (spiCtx != null) {
- TcpDiscoveryAbstractMessage msg = decodeMessage(this, data);
+ TcpDiscoveryAbstractMessage msg = decodeMessage(ignite.context(),
data);
if (msg != null)
apply(spiCtx.localNode(), msg);
diff --git
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/DiscoveryUnmarshalVulnerabilityTest.java
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/DiscoveryUnmarshalVulnerabilityTest.java
index 5a4e18fe8d4..8e8029c295e 100644
---
a/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/DiscoveryUnmarshalVulnerabilityTest.java
+++
b/modules/core/src/test/java/org/apache/ignite/spi/discovery/tcp/DiscoveryUnmarshalVulnerabilityTest.java
@@ -245,7 +245,7 @@ public abstract class DiscoveryUnmarshalVulnerabilityTest
extends GridCommonAbst
private byte[] serializedMessage() throws IgniteCheckedException {
ByteBuffer buf = ByteBuffer.allocate(4096);
- MessageFactory msgFactory =
((TcpDiscoverySpi)grid(0).configuration().getDiscoverySpi()).messageFactory();
+ MessageFactory msgFactory = grid(0).context().messageFactory();
DirectMessageWriter writer = new DirectMessageWriter(msgFactory);
writer.setBuffer(buf);
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 ca96f33f159..0d5222fca2d 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
@@ -2600,7 +2600,7 @@ public class TcpClientDiscoverySpiSelfTest extends
GridCommonAbstractTest {
waitFor(writeLock);
- TcpDiscoveryAbstractMessage msg = decodeMessage(this, data);
+ TcpDiscoveryAbstractMessage msg = decodeMessage(ignite.context(),
data);
if (msg != null && !onMessage(sock, msg))
return;
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 866bbc84217..feedf6a4c6e 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
@@ -26,19 +26,15 @@ 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.GridKernalContext;
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;
import org.apache.ignite.spi.discovery.DiscoverySpiListener;
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;
import org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryPingResponse;
-import org.apache.ignite.testframework.GridTestUtils;
import org.apache.ignite.testframework.GridTestUtils.DiscoveryHook;
import org.jetbrains.annotations.Nullable;
@@ -57,12 +53,6 @@ public class TestTcpDiscoverySpi extends TcpDiscoverySpi
implements IgniteDiscov
/** */
private IgniteDiscoverySpiInternalListener internalLsnr;
- /** */
- private MessageFactory msgFactory;
-
- /** */
- private MessageFactoryProvider provider;
-
/** {@inheritDoc} */
@Override protected void writeMessage(TcpDiscoveryIoSession ses,
TcpDiscoveryAbstractMessage msg, long timeout) throws IOException,
IgniteCheckedException {
@@ -119,42 +109,8 @@ public class TestTcpDiscoverySpi extends TcpDiscoverySpi
implements IgniteDiscov
this.discoHook = discoHook;
}
- /**
- * Sets test discovery messages factory provider. Note that {@link
MessageFactoryProvider} must be set before SPI start.
- * Otherwise, this method call will take no effect.
- *
- * @param msgFactoryProvider Discovery messages factory provider.
- */
- public void messageFactory(MessageFactoryProvider msgFactoryProvider) {
- provider = msgFactoryProvider;
- assert !started();
-
- msgFactory = new IgniteMessageFactoryImpl(new MessageFactoryProvider[]
{
- new CoreMessagesProvider(),
- msgFactoryProvider
- });
- }
-
- /** {@inheritDoc} */
- @Override public MessageFactoryProvider messageFactoryProvider() {
- return provider;
- }
-
- /** {@inheritDoc} */
- @Override protected void initLocalNode(int srvPort, boolean
addExtAddrAttr) {
- if (msgFactory != null)
- GridTestUtils.setFieldValue(this, TcpDiscoverySpi.class,
"msgFactory", msgFactory);
-
- super.initLocalNode(srvPort, addExtAddrAttr);
- }
-
- /** {@inheritDoc} */
- @Override public MessageFactory messageFactory() {
- return msgFactory != null ? msgFactory : super.messageFactory();
- }
-
/** */
- public static @Nullable TcpDiscoveryAbstractMessage
decodeMessage(TcpDiscoverySpi spi, byte[] data) {
+ public static @Nullable TcpDiscoveryAbstractMessage
decodeMessage(GridKernalContext ctx, byte[] data) {
if (Arrays.equals(U.IGNITE_HEADER, data))
return null;
@@ -169,7 +125,7 @@ public class TestTcpDiscoverySpi extends TcpDiscoverySpi
implements IgniteDiscov
};
try (dataSock) {
- return new TcpDiscoveryIoSession(dataSock, spi).readMessage();
+ return new TcpDiscoveryIoSession(ctx, dataSock).readMessage();
}
catch (Exception e) {
throw new IgniteException("Failed to decode a message", e);