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
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,6 @@
import org.apache.ignite.IgniteLogger;
import org.apache.ignite.internal.GridKernalContext;
import org.apache.ignite.internal.direct.DirectMessageReader;
import org.apache.ignite.internal.direct.DirectMessageWriter;
import org.apache.ignite.internal.managers.communication.DiscoveryMarshalling;
import org.apache.ignite.internal.managers.communication.UnknownMessageException;
import org.apache.ignite.internal.util.CommonUtils;
Expand All @@ -47,6 +46,7 @@
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.internal.TcpDiscoveryMessageSerializer;
import org.apache.ignite.spi.discovery.tcp.messages.TcpDiscoveryAbstractMessage;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
Expand All @@ -66,8 +66,8 @@ public class TcpDiscoveryIoSession implements AutoCloseable {
/** Default size of buffer used for buffering socket in/out. */
private static final int DFLT_SOCK_BUFFER_SIZE = 8192;

/** Size for an intermediate buffer for serializing discovery messages. */
private static final int MSG_BUFFER_SIZE = 100;
/** Size of the intermediate buffer a message is deserialized through. */
private static final int READ_BUFFER_SIZE = 100;

/** */
private final GridKernalContext ctx;
Expand All @@ -82,23 +82,20 @@ public class TcpDiscoveryIoSession implements AutoCloseable {
private final Socket sock;

/** */
private final DirectMessageWriter msgWriter;
private final TcpDiscoveryMessageSerializer msgSer;

/** */
private final DirectMessageReader msgReader;

/** */
private final ByteBuffer readBuf;

/** Buffered socket output stream. */
private final OutputStream out;

/** Buffered socket input stream. */
private final CompositeInputStream in;

/** */
private final ByteBuffer readBuf;

/** */
private final ByteBuffer writeBuf;

/**
* Creates a new discovery I/O session bound to the given socket.
*
Expand All @@ -112,12 +109,11 @@ public class TcpDiscoveryIoSession implements AutoCloseable {
this.msgFactory = ctx.messageFactory();
this.log = ctx.log(getClass());

readBuf = ByteBuffer.allocate(MSG_BUFFER_SIZE);
writeBuf = ByteBuffer.allocate(MSG_BUFFER_SIZE);

msgWriter = new DirectMessageWriter(msgFactory);
readBuf = ByteBuffer.allocate(READ_BUFFER_SIZE);
msgReader = new DirectMessageReader(msgFactory, null);

msgSer = new TcpDiscoveryMessageSerializer(ctx);

try {
int sendBufSize = sock.getSendBufferSize() > 0 ? sock.getSendBufferSize() : DFLT_SOCK_BUFFER_SIZE;
int rcvBufSize = sock.getReceiveBufferSize() > 0 ? sock.getReceiveBufferSize() : DFLT_SOCK_BUFFER_SIZE;
Expand All @@ -136,9 +132,9 @@ public class TcpDiscoveryIoSession implements AutoCloseable {
* @param msg Message to send to the remote node.
* @throws IgniteCheckedException If serialization fails.
*/
void writeMessage(TcpDiscoveryAbstractMessage msg) throws IgniteCheckedException, IOException {
synchronized void writeMessage(TcpDiscoveryAbstractMessage msg) throws IgniteCheckedException, IOException {
try {
serializeMessage((Message)msg, out);
msgSer.writeTo(msg, out);

out.flush();
}
Expand Down Expand Up @@ -262,39 +258,13 @@ public Socket socket() {
return sock;
}

/**
* Serializes a discovery message into given output stream.
*
* @param m Discovery message to serialize.
* @param out Output stream to write serialized message.
* @throws IOException If serialization fails.
*/
void serializeMessage(Message m, OutputStream out) throws IOException, IgniteCheckedException {
DiscoveryMarshalling.marshal(m, ctx, null);

msgWriter.reset();
msgWriter.setBuffer(writeBuf);

boolean finished;

do {
// Should be cleared before first operation.
writeBuf.clear();

finished = MessageSerialization.writeTo(msgFactory, m, msgWriter);

out.write(writeBuf.array(), 0, writeBuf.position());
}
while (!finished);
}

/**
* Writes raw data to the underlying socket output stream.
*
* @param data Raw data to write.
* @throws IOException If failed.
*/
void write(byte[] data) throws IOException {
synchronized void write(byte[] data) throws IOException {
out.write(data);

out.flush();
Expand All @@ -306,7 +276,7 @@ void write(byte[] data) throws IOException {
* @param b Integer response.
* @throws IOException If failed.
*/
void write(int b) throws IOException {
synchronized void write(int b) throws IOException {
out.write(b);

out.flush();
Expand Down
Loading
Loading