Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -99,7 +99,6 @@
import org.apache.ignite.internal.util.typedef.C1;
import org.apache.ignite.internal.util.typedef.F;
import org.apache.ignite.internal.util.typedef.P1;
import org.apache.ignite.internal.util.typedef.T2;
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.S;
Expand All @@ -122,8 +121,10 @@
import org.apache.ignite.spi.discovery.DiscoverySpiCustomMessage;
import org.apache.ignite.spi.discovery.DiscoverySpiListener;
import org.apache.ignite.spi.discovery.IgniteDiscoveryThread;
import org.apache.ignite.spi.discovery.tcp.internal.ClientMessageHolder;
import org.apache.ignite.spi.discovery.tcp.internal.DiscoveryDataPacket;
import org.apache.ignite.spi.discovery.tcp.internal.FutureTask;
import org.apache.ignite.spi.discovery.tcp.internal.TcpDiscoveryMessageSerializer;
import org.apache.ignite.spi.discovery.tcp.internal.TcpDiscoveryNode;
import org.apache.ignite.spi.discovery.tcp.internal.TcpDiscoveryNodesRing;
import org.apache.ignite.spi.discovery.tcp.internal.TcpDiscoverySpiState;
Expand Down Expand Up @@ -2859,7 +2860,7 @@ protected class RingMessageWorker extends MessageWorker<TcpDiscoveryAbstractMess
// 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(ctx);
private final TcpDiscoveryMessageSerializer cliMsgSer = new TcpDiscoveryMessageSerializer(ctx);

/** IO session. */
private TcpDiscoveryIoSession ses;
Expand Down Expand Up @@ -3229,50 +3230,52 @@ private void processAuthFailedMessage(TcpDiscoveryAuthFailedMessage authFailedMs
* @param msg Message.
*/
private void sendMessageToClients(TcpDiscoveryAbstractMessage msg) {
if (redirectToClients(msg)) {
if (spi.ensured(msg))
msgHist.add(msg);
if (!redirectToClients(msg))
return;

if (clientMsgWorkers.isEmpty())
return;
if (spi.ensured(msg))
msgHist.add(msg);

if (clientMsgWorkers.isEmpty())
return;

byte[] msgBytes;
ClientMessageHolder sharedMsgHolder = new ClientMessageHolder(msg);

for (ClientMessageWorker worker : clientMsgWorkers.values()) {
TcpDiscoveryAbstractMessage rebuiltMsg = rebuildForClient(msg, worker.clientNodeId);

ClientMessageHolder msgToSend = rebuiltMsg == msg
? sharedMsgHolder
: new ClientMessageHolder(rebuiltMsg);

try {
msgBytes = clientMsgSer.serializeMessage(msg);
msgToSend.serialize(cliMsgSer);
}
catch (IgniteCheckedException | IOException e) {
U.error(log, "Failed to serialize message: " + msg, e);
catch (IgniteCheckedException e) {
U.error(log, "Failed to serialize message: " + msgToSend, e);

return;
}

for (ClientMessageWorker clientMsgWorker : clientMsgWorkers.values()) {
TcpDiscoveryAbstractMessage msg0 = msg;
byte[] msgBytes0 = msgBytes;
worker.addMessage(msgToSend);
}
}

if (msg instanceof TcpDiscoveryNodeAddedMessage) {
TcpDiscoveryNodeAddedMessage nodeAddedMsg = (TcpDiscoveryNodeAddedMessage)msg;
/** */
private TcpDiscoveryAbstractMessage rebuildForClient(TcpDiscoveryAbstractMessage msg, UUID clientNodeId) {
if (!(msg instanceof TcpDiscoveryNodeAddedMessage))
return msg;

if (clientMsgWorker.clientNodeId.equals(nodeAddedMsg.node().id())) {
msg0 = new TcpDiscoveryNodeAddedMessage(nodeAddedMsg);
TcpDiscoveryNodeAddedMessage nodeAddedMsg = (TcpDiscoveryNodeAddedMessage)msg;

prepareNodeAddedMessage(msg0, clientMsgWorker.clientNodeId, null);
if (!clientNodeId.equals(nodeAddedMsg.node().id()))
return msg;

try {
msgBytes0 = clientMsgSer.serializeMessage(msg0);
}
catch (IgniteCheckedException | IOException e) {
U.error(log, "Failed to serialize message: " + msg0, e);
TcpDiscoveryNodeAddedMessage res = new TcpDiscoveryNodeAddedMessage(nodeAddedMsg);

return;
}
}
}
prepareNodeAddedMessage(res, clientNodeId, null);

clientMsgWorker.addMessage(msg0, msgBytes0);
}
}
return res;
}

/**
Expand Down Expand Up @@ -7469,21 +7472,10 @@ private class StatisticsPrinter extends IgniteSpiThread {
}

/** */
private class ClientMessageWorker extends MessageWorker<T2<TcpDiscoveryAbstractMessage, byte[]>> {
private class ClientMessageWorker extends MessageWorker<ClientMessageHolder> {
/** Node ID. */
private final UUID clientNodeId;

// The code responsible for sending and receiving messages to and from client nodes represents a special case in ServerImpl,
// as it is split into two separate components.
// One part, ClientMessageWorker, handles only message sending to clients and does not process responses.
// The other part, which reads messages from clients, is implemented in SocketReader.
// Due to this separation, we don't require a full TcpDiscoveryIoSession here
// and can instead extract just the message-writing functionality.
// At the same time, we aim to keep both reading and writing logic encapsulated within TcpDiscoveryIoSession.
// As a result, we need to copy some code from TcpDiscoveryIoSession into the new class, TcpDiscoveryMessageSerializer.
/** */
private final TcpDiscoveryMessageSerializer clientMsgSer;

/** Session shared with the socket reader serving the same client connection. */
private final TcpDiscoveryIoSession ses;

Expand Down Expand Up @@ -7518,8 +7510,6 @@ private ClientMessageWorker(TcpDiscoveryIoSession ses, UUID clientNodeId, Ignite
this.ses = ses;
this.clientNodeId = clientNodeId;

clientMsgSer = new TcpDiscoveryMessageSerializer(ctx);

lastMetricsUpdateMsgTimeNanos = System.nanoTime();
}

Expand All @@ -7546,24 +7536,19 @@ void metrics(ClusterMetrics metrics) {
this.metrics = metrics;
}

/**
* @param msg Message.
*/
/** @param msg Discovery Message. */
void addMessage(TcpDiscoveryAbstractMessage msg) {
addMessage(msg, null);
addMessage(new ClientMessageHolder(msg));
}

/**
* @param msg Message.
* @param msgBytes Optional message bytes.
*/
void addMessage(TcpDiscoveryAbstractMessage msg, @Nullable byte[] msgBytes) {
T2<TcpDiscoveryAbstractMessage, byte[]> t = new T2<>(msg, msgBytes);
/** @param msgHolder Holder of a Discovery Message to send to the client. */
void addMessage(ClientMessageHolder msgHolder) {
TcpDiscoveryAbstractMessage msg = msgHolder.message();

if (msg.highPriority())
queue.addFirst(t);
queue.addFirst(msgHolder);
else
queue.add(t);
queue.add(msgHolder);

DebugLogger log = messageLogger(msg);

Expand All @@ -7572,10 +7557,10 @@ void addMessage(TcpDiscoveryAbstractMessage msg, @Nullable byte[] msgBytes) {
}

/** {@inheritDoc} */
@Override protected void processMessage(T2<TcpDiscoveryAbstractMessage, byte[]> msgT) {
@Override protected void processMessage(ClientMessageHolder msgHolder) {
boolean success = false;

TcpDiscoveryAbstractMessage msg = msgT.get1();
TcpDiscoveryAbstractMessage msg = msgHolder.message();

try {
assert msg.verified() : msg;
Expand All @@ -7601,8 +7586,9 @@ else if (msgLog.isDebugEnabled()) {
+ getLocalNodeId() + ", rmtNodeId=" + clientNodeId + ", msg=" + msg + ']');
}

writeToSocket(msgT, spi.failureDetectionTimeoutEnabled() ? spi.clientFailureDetectionTimeout() :
spi.getSocketTimeout());
long timeout = spi.failureDetectionTimeoutEnabled() ? spi.clientFailureDetectionTimeout() : spi.getSocketTimeout();

writeMessage(msgHolder, timeout);
}
}
else {
Expand All @@ -7613,7 +7599,7 @@ else if (msgLog.isDebugEnabled()) {

assert topologyInitialized(msg) : msg;

writeToSocket(msgT, spi.getEffectiveSocketTimeout(false));
writeMessage(msgHolder, spi.getEffectiveSocketTimeout(false));
}

boolean clientFailed = msg instanceof TcpDiscoveryNodeFailedMessage &&
Expand Down Expand Up @@ -7643,14 +7629,16 @@ else if (msgLog.isDebugEnabled()) {
}

/**
* @param msgT Message tuple.
* @param msgHolder Message holder.
* @param timeout Timeout.
*/
private void writeToSocket(T2<TcpDiscoveryAbstractMessage, byte[]> msgT, long timeout)
throws IgniteCheckedException, IOException {
byte[] msgBytes = msgT.get2() == null ? clientMsgSer.serializeMessage(msgT.get1()) : msgT.get2();
private void writeMessage(ClientMessageHolder msgHolder, long timeout) throws IgniteCheckedException, IOException {
byte[] msgBytes = msgHolder.messageBytes();

spi.write(ses, msgBytes, timeout);
if (msgBytes != null)
spi.write(ses, msgBytes, timeout);
else
spi.writeMessage(ses, msgHolder.message(), timeout);
}

/**
Expand Down
Loading
Loading