From 281029604d4409571bf306060da69781c4d9187e Mon Sep 17 00:00:00 2001 From: Mikhail Petrov Date: Thu, 3 Sep 2026 12:29:48 +0300 Subject: [PATCH 1/2] IGNITE-29037 --- .../ignite/spi/discovery/tcp/ClientImpl.java | 18 +----- .../ignite/spi/discovery/tcp/ServerImpl.java | 57 +++++++------------ .../spi/discovery/tcp/TcpDiscoveryImpl.java | 19 +++---- .../discovery/tcp/TcpDiscoveryIoSession.java | 38 +++++++------ .../tcp/TcpDiscoveryMessageSerializer.java | 9 +-- .../spi/discovery/tcp/TcpDiscoverySpi.java | 14 +---- .../cache/CacheMetricsCacheSizeTest.java | 3 +- .../ignite/spi/MessagesPluginProvider.java | 13 ----- .../discovery/tcp/BlockTcpDiscoverySpi.java | 2 +- .../DiscoveryUnmarshalVulnerabilityTest.java | 2 +- .../tcp/TcpClientDiscoverySpiSelfTest.java | 2 +- .../discovery/tcp/TestTcpDiscoverySpi.java | 50 +--------------- 12 files changed, 64 insertions(+), 163 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 621a234ee5428..147dc198f6c33 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.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.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.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,8 +205,7 @@ 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 = adapter.ignite().name() == null ? "client-node" : adapter.ignite().name(); executorSrvc = newSingleThreadScheduledExecutor("tcp-discovery-exec", instanceName); } @@ -1052,13 +1048,6 @@ private void joinError(IgniteSpiException err) { joinLatch.countDown(); } - /** */ - private WorkersRegistry getWorkersRegistry() { - Ignite ignite = spi.ignite(); - - return ignite instanceof IgniteEx ? ((IgniteEx)ignite).context().workersRegistry() : null; - } - /** */ private Collection remoteVisibleNodes() { return U.arrayList(rmtNodes.values(), TcpDiscoveryNodesRing.VISIBLE_NODES); @@ -1689,7 +1678,7 @@ protected class MessageWorker extends GridWorker { * @param log Logger. */ private MessageWorker(IgniteLogger log) { - super(spi.ignite().name(), "tcp-client-disco-msg-worker", log, getWorkersRegistry()); + super(spi.ignite().name(), "tcp-client-disco-msg-worker", log, ctx.workersRegistry()); } /** {@inheritDoc} */ @@ -1957,8 +1946,7 @@ else if (discoMsg instanceof TcpDiscoveryCheckFailedMessage) 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; 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 e8a9e5c239400..bce7eabd7c392 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.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.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; @@ -369,7 +367,7 @@ class ServerImpl extends TcpDiscoveryImpl { new LinkedBlockingQueue<>()); List props = newConnectionEnabledProperty( - ((IgniteEx)spi.ignite()).context().internalSubscriptionProcessor(), + ctx.internalSubscriptionProcessor(), log, "ClientNode", "ServerNode" @@ -2458,11 +2456,6 @@ private static void enrichNodeWithAttribute(TcpDiscoveryNode node, String attrNa 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 @@ protected class RingMessageWorker extends MessageWorker 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 @@ private void addToQueue(TcpDiscoveryAbstractMessage msg, boolean addFirst) { 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)); } } @@ -6310,7 +6301,7 @@ private class TcpServer extends GridWorker { * @throws IgniteSpiException In case of error. */ TcpServer(IgniteLogger log) throws IgniteSpiException { - super(spi.ignite().name(), "tcp-disco-srvr-[]", log, getWorkerRegistry(spi)); + super(spi.ignite().name(), "tcp-disco-srvr-[]", log, ctx.workersRegistry()); int lastPort = spi.locPortRange == 0 ? spi.locPort : spi.locPort + spi.locPortRange - 1; @@ -6413,17 +6404,15 @@ private class TcpServer extends GridWorker { 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); } @@ -6462,7 +6451,7 @@ private class SocketReader extends IgniteSpiThread { this.sock = sock; - ses = createSession(sock); + ses = new TcpDiscoveryIoSession(ctx, sock); setPriority(spi.threadPri); } @@ -7085,11 +7074,7 @@ else if (msg instanceof TcpDiscoveryRingLatencyCheckMessage) { } } 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) { @@ -7534,7 +7519,7 @@ private ClientMessageWorker(TcpDiscoveryIoSession ses, UUID clientNodeId, Ignite this.ses = ses; this.clientNodeId = clientNodeId; - clientMsgSer = new TcpDiscoveryMessageSerializer(spi); + clientMsgSer = new TcpDiscoveryMessageSerializer(ctx); lastMetricsUpdateMsgTimeNanos = System.nanoTime(); } 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 b83e2130ae2e0..0f2324361d5dc 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.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 @@ protected static List toOrderedList(Collection addrs) 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 71a4995eeabcb..7689b921a9e80 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.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.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,10 @@ 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 Socket sock; @@ -96,19 +99,20 @@ public class TcpDiscoveryIoSession implements AutoCloseable { /** * Creates a new discovery I/O session bound to the given socket. * + * @param ctx Kernel 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(); 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 +182,7 @@ T readMessage() throws IgniteCheckedException, IOException { Message msg; try { - msg = spi.messageFactory().create(msgType); + msg = msgFactory.create(msgType); } catch (IgniteException e) { detectSslAlert(b0, b1); @@ -202,7 +206,7 @@ T readMessage() throws IgniteCheckedException, IOException { 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 +222,7 @@ T readMessage() throws IgniteCheckedException, IOException { } while (!finished); - GridKernalContext kctx = ((IgniteEx)spi.ignite()).context(); - - DiscoveryMarshalling.unmarshal(msg, kctx); + DiscoveryMarshalling.unmarshal(msg, ctx); return (T)msg; } @@ -238,14 +240,16 @@ T readMessage() throws IgniteCheckedException, IOException { /** @return SSL certificate this session is established with. {@code null} if SSL is disabled or certificate validation failed. */ @Nullable Certificate[] extractCertificates() { - if (!spi.isSslEnabled()) + boolean isSslEnabled = ctx.config().getSslContextFactory() != null; + + if (!isSslEnabled) 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(ctx.log(getClass()), "Failed to extract discovery IO session certificates", e); return null; } @@ -264,9 +268,7 @@ public Socket socket() { * @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 +279,7 @@ void serializeMessage(Message m, OutputStream out) throws IOException, IgniteChe // 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 8c871f9e2eb46..ca06882921485 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.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 @@ */ class TcpDiscoveryMessageSerializer extends TcpDiscoveryIoSession { /** - * @param spi Discovery SPI instance. + * @param ctx Kernel 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 @@ public TcpDiscoveryMessageSerializer(TcpDiscoverySpi spi) { @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 c4c74289af30b..0f9c627b6bcad 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 @@ -74,7 +74,6 @@ 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 +462,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 +595,6 @@ public void setClientReconnectDisabled(boolean clientReconnectDisabled) { setAddressResolver(ignite.configuration().getAddressResolver()); marsh = ((IgniteEx)ignite).context().marshallerContext().jdkMarshaller(); - - msgFactory = ((IgniteEx)ignite).context().messageFactory(); } } @@ -1122,11 +1115,6 @@ public void setConnectionRecoveryTimeout(long connRecoveryTimeout) { locNodeVer = ver; } - /** @return Discovery messages factory. */ - public MessageFactory messageFactory() { - return msgFactory; - } - /** * Gets ID of the local node. * @@ -1624,7 +1612,7 @@ protected TcpDiscoveryIoSession openSession( 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)); 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 9aed75febee59..a4a0db433309d 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.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 void testCacheSize() throws Exception { 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 1198584d22972..0090d11cc3c56 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.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 MessagesPluginProvider(Class... msgs) { // 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 52e81b05df960..4904df46a40d3 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 @@ private synchronized void apply(ClusterNode addr, TcpDiscoveryAbstractMessage ms 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 5a4e18fe8d421..8e8029c295ee5 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 @@ void attack(byte[] data) throws IOException { 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 ca96f33f159cf..0d5222fca2d2a 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 @@ private void pauseResumeOperation(boolean isPause, AtomicBoolean... locks) { 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 866bbc842178b..feedf6a4c6ed5 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.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 void discoveryHook(DiscoveryHook discoHook) { 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 void messageFactory(MessageFactoryProvider msgFactoryProvider) { }; 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); From 957c5d028809b5008000c8a9452896bd97469cad Mon Sep 17 00:00:00 2001 From: Mikhail Petrov Date: Thu, 3 Sep 2026 19:58:23 +0300 Subject: [PATCH 2/2] IGNITE-29037 --- .../ignite/spi/discovery/tcp/ClientImpl.java | 18 ++++++------- .../ignite/spi/discovery/tcp/ServerImpl.java | 25 +++++++++---------- .../discovery/tcp/TcpDiscoveryIoSession.java | 12 +++++---- .../tcp/TcpDiscoveryMessageSerializer.java | 2 +- .../spi/discovery/tcp/TcpDiscoverySpi.java | 7 +----- 5 files changed, 30 insertions(+), 34 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 147dc198f6c33..dd1aed7a2a24f 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 @@ -205,9 +205,9 @@ class ClientImpl extends TcpDiscoveryImpl { ClientImpl(TcpDiscoverySpi adapter) { super(adapter); - String instanceName = 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} */ @@ -1090,7 +1090,7 @@ private class SocketReader extends IgniteSpiThread { /** */ SocketReader() { - super(spi.ignite().name(), "tcp-client-disco-sock-reader-[]", log); + super(ctx.igniteInstanceName(), "tcp-client-disco-sock-reader-[]", log); } /** @@ -1263,7 +1263,7 @@ private class SocketWriter extends IgniteSpiThread { * */ 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(); @@ -1513,7 +1513,7 @@ private class Reconnector extends IgniteSpiThread { * @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; @@ -1678,7 +1678,7 @@ protected class MessageWorker extends GridWorker { * @param log Logger. */ private MessageWorker(IgniteLogger log) { - super(spi.ignite().name(), "tcp-client-disco-msg-worker", log, ctx.workersRegistry()); + super(ctx.igniteInstanceName(), "tcp-client-disco-msg-worker", log, ctx.workersRegistry()); } /** {@inheritDoc} */ @@ -2198,7 +2198,7 @@ else if (log.isDebugEnabled()) if (joining()) delayDiscoData.add(dataPacket); else - spi.onExchange(dataPacket, U.resolveClassLoader(spi.ignite().configuration())); + spi.onExchange(dataPacket, U.resolveClassLoader(ctx.config())); } } } @@ -2232,11 +2232,11 @@ private void processNodeAddFinishedMessage(TcpDiscoveryNodeAddFinishedMessage ms 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 bce7eabd7c392..5974e61780743 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 @@ -360,7 +360,7 @@ class ServerImpl extends TcpDiscoveryImpl { super(adapter); utilityPool = new IgniteThreadPoolExecutor("disco-pool", - spi.ignite().name(), + ctx.igniteInstanceName(), 0, utilityPoolSize, 2000, @@ -2247,7 +2247,7 @@ private class IpFinderCleaner extends IgniteSpiThread { * 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); } @@ -4582,8 +4582,7 @@ private IgniteNodeValidationResult validateByIgniteComponentsWithJoiningNodeData 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); } @@ -4927,7 +4926,7 @@ else if (!locNodeId.equals(node.id()) && ring.node(node.id()) != null) { 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); } @@ -5145,11 +5144,11 @@ private void processNodeAddFinishedMessage(TcpDiscoveryNodeAddFinishedMessage ms } 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(); @@ -6301,7 +6300,7 @@ private class TcpServer extends GridWorker { * @throws IgniteSpiException In case of error. */ TcpServer(IgniteLogger log) throws IgniteSpiException { - super(spi.ignite().name(), "tcp-disco-srvr-[]", log, ctx.workersRegistry()); + super(ctx.igniteInstanceName(), "tcp-disco-srvr-[]", log, ctx.workersRegistry()); int lastPort = spi.locPortRange == 0 ? spi.locPort : spi.locPort + spi.locPortRange - 1; @@ -6321,7 +6320,7 @@ private class TcpServer extends GridWorker { if (log.isInfoEnabled()) { log.info("Successfully bound to TCP port [port=" + port + ", localHost=" + spi.locHost + - ", locNodeId=" + spi.ignite().configuration().getNodeId() + + ", locNodeId=" + spi.cfgNodeId + ']'); } @@ -6447,7 +6446,7 @@ private class SocketReader extends IgniteSpiThread { * @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; @@ -7446,7 +7445,7 @@ private class StatisticsPrinter extends IgniteSpiThread { * Constructor. */ StatisticsPrinter() { - super(spi.ignite().name(), "tcp-disco-stats-printer", log); + super(ctx.igniteInstanceName(), "tcp-disco-stats-printer", log); assert spi.statsPrintFreq > 0; @@ -7863,7 +7862,7 @@ protected MessageWorker( @Nullable GridWorkerListener lsnr, BlockingDeque queue ) { - super(spi.ignite().name(), name, log, lsnr); + super(ctx.igniteInstanceName(), name, log, lsnr); this.queue = queue; this.pollingTimeout = pollingTimeout; @@ -8086,7 +8085,7 @@ void pingRemoteDCs(List nodesToPing) { 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/TcpDiscoveryIoSession.java b/modules/core/src/main/java/org/apache/ignite/spi/discovery/tcp/TcpDiscoveryIoSession.java index 7689b921a9e80..0251712fc95b7 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 @@ -75,6 +75,9 @@ public class TcpDiscoveryIoSession implements AutoCloseable { /** */ private final MessageFactory msgFactory; + /** */ + private final IgniteLogger log; + /** */ private final Socket sock; @@ -99,7 +102,7 @@ public class TcpDiscoveryIoSession implements AutoCloseable { /** * Creates a new discovery I/O session bound to the given socket. * - * @param ctx Kernel context. + * @param ctx Kernal context. * @param sock Socket connected to a remote discovery node. * @throws IgniteException If an I/O error occurs while initializing buffers. */ @@ -107,6 +110,7 @@ public class TcpDiscoveryIoSession implements AutoCloseable { this.sock = sock; this.ctx = ctx; this.msgFactory = ctx.messageFactory(); + this.log = ctx.log(getClass()); readBuf = ByteBuffer.allocate(MSG_BUFFER_SIZE); writeBuf = ByteBuffer.allocate(MSG_BUFFER_SIZE); @@ -240,16 +244,14 @@ T readMessage() throws IgniteCheckedException, IOException { /** @return SSL certificate this session is established with. {@code null} if SSL is disabled or certificate validation failed. */ @Nullable Certificate[] extractCertificates() { - boolean isSslEnabled = ctx.config().getSslContextFactory() != null; - - if (!isSslEnabled) + if (!(sock instanceof SSLSocket)) return null; try { return ((SSLSocket)sock).getSession().getPeerCertificates(); } catch (SSLPeerUnverifiedException e) { - U.error(ctx.log(getClass()), "Failed to extract discovery IO session certificates", e); + U.error(log, "Failed to extract discovery IO session certificates", e); return null; } 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 ca06882921485..ec7cdc569f0c7 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 @@ -36,7 +36,7 @@ */ class TcpDiscoveryMessageSerializer extends TcpDiscoveryIoSession { /** - * @param ctx Kernel context. + * @param ctx Kernal context. */ public TcpDiscoveryMessageSerializer(GridKernalContext ctx) { super(ctx, new Socket() { 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 0f9c627b6bcad..7b8490ded8385 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.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; @@ -2119,11 +2118,7 @@ protected void onExchange(DiscoveryDataPacket dataPacket, ClassLoader clsLdr) { 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); }