Skip to content
Merged
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 @@ -13,6 +13,7 @@
*
* Copyright 2006-2009 Sun Microsystems, Inc.
* Portions Copyright 2011-2016 ForgeRock AS.
* Portions Copyright 2026 3A Systems, LLC.
*/
package org.opends.server.replication.protocol;

Expand All @@ -30,12 +31,14 @@
import java.util.concurrent.LinkedBlockingQueue;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.zip.DataFormatException;

import javax.net.ssl.SSLSocket;

import org.forgerock.i18n.LocalizableMessage;
import org.forgerock.i18n.slf4j.LocalizedLogger;
import org.opends.server.api.DirectoryThread;
import org.opends.server.types.HostPort;
Expand All @@ -48,6 +51,19 @@ public final class Session extends DirectoryThread implements Closeable
{
private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();

/**
* How long a close spends sending what the publisher thread left queued in {@code sendQueue},
* in milliseconds.
* <p>
* The same 5 s as {@code DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD}, which is how long a
* shutdown is already willing to wait for one of these messages - the announcement that a
* replica went offline - to be forwarded. A close has no reason to wait longer for it than the
* shutdown which is waiting on the close. It is a budget per close, though, not per shutdown:
* a shutdown closes its sessions one after another, so it can pay this once for each peer which
* is alive and not reading.
*/
private static final long DRAIN_BUDGET_MS = 5000;

private final Socket plainSocket;
private final SSLSocket secureSocket;
private final InputStream plainInput;
Expand Down Expand Up @@ -100,6 +116,14 @@ public final class Session extends DirectoryThread implements Closeable
private AtomicBoolean isRunning = new AtomicBoolean(false);
private final CountDownLatch latch = new CountDownLatch(1);

/**
* How many buffers the publisher thread took off {@code sendQueue} and failed to write. After a
* failed write that thread goes on taking the queue until the session is closed, and every
* write after the first fails as well, so these are messages the peer was not told about just
* as much as what is still queued - see {@link #close()}.
*/
private final AtomicInteger publisherFailedWrites = new AtomicInteger();

/**
* Creates a new Session.
*
Expand Down Expand Up @@ -140,6 +164,10 @@ public Session(final Socket socket,
/**
* This method is called when the session with the remote must be closed.
* This object won't be used anymore after this method is called.
* <p>
* A message which was published on this session but which its publisher thread had not sent yet
* is sent here rather than dropped, within the budget of {@link #DRAIN_BUDGET_MS}. See {@link
* #sendWhatThePublisherLeftQueued()}.
*/
@Override
public void close()
Expand All @@ -165,6 +193,18 @@ public void close()
Thread.currentThread().interrupt();
}

/*
* Re-read the error rather than answer with the snapshot taken before the join: a publisher
* whose send() failed while this thread was joining it recorded the error there, and what
* follows - the drain and the StopMsg - is what must not be written to a socket which has
* already failed. Reading it before the join left both writing to one, the drain naming its
* own failure rather than the one the publisher had recorded.
*/
synchronized (stateLock)
{
localSessionError = sessionError;
}

// Perform close outside of critical section.
if (logger.isTraceEnabled())
{
Expand All @@ -186,25 +226,174 @@ public void close()
}
}

// V4 protocol introduces a StopMsg to properly end communications.
if (localSessionError == null
&& protocolVersion >= ProtocolVersion.REPLICATION_PROTOCOL_V4)
if (localSessionError != null)
{
try
/*
* Nothing more is written to a session which already failed, neither what the publisher left
* queued nor the StopMsg: writing more to it cannot work. What it gives up on is still
* reported: what is left in the queue, and what the publisher took off it and failed to
* write - the write which failed first, and every one it went on failing until this close
* stopped it, which the join above makes complete. It is done without taking publishLock,
* so that a thread blocked in a write of this socket is released by the close of the
* sockets below rather than waited for.
*/
isRunning.set(false);
final int notSent = publisherFailedWrites.get() + takeWhatIsLeftQueued();
if (notSent > 0)
{
publish(new StopMsg());
reportQueueNotSent(notSent, "the session had already failed: " + localSessionError);
}
catch (final IOException ignored)
StaticUtils.close(plainSocket, secureSocket);
return;
}

/*
* The publisher thread has stopped and what it had not sent is still in the queue. Send it,
* rather than let the close drop it: nothing publishes these again, and the StopMsg which
* follows leaves the peer reading an orderly close with no sign that anything was missing.
*
* This thread is not the only one which can write the socket here - a ServerWriter or a
* HeartbeatThread outlives this close. It is this thread, not run(), which takes the session
* off the queueing branch of publish(), and it does so under publishLock and keeps the lock
* across the queue and the StopMsg: until then a publish() concurrent with the close takes
* the queueing branch, and from then on it takes the synchronous branch and waits for the
* lock. Either way nothing newer is written between two of the drained messages - see
* sendWhatThePublisherLeftQueued() below.
*/
publishLock.lock();
try
{
isRunning.set(false);
sendWhatThePublisherLeftQueued();

/*
* Re-read again: a write of the drain which failed recorded its error, and the StopMsg must
* not be written to a socket which has already failed either.
*/
synchronized (stateLock)
{
localSessionError = sessionError;
}

// V4 protocol introduces a StopMsg to properly end communications.
if (localSessionError == null
&& protocolVersion >= ProtocolVersion.REPLICATION_PROTOCOL_V4)
{
// Ignore errors on close.
try
{
publish(new StopMsg());
}
catch (final IOException ignored)
{
// Ignore errors on close.
}
}
}
finally
{
publishLock.unlock();
}

StaticUtils.close(plainSocket, secureSocket);
}



/**
* Sends the buffers the publisher thread had not sent when it stopped, so that a close of the
* session does not drop them.
* <p>
* Called from {@link #close()} once the publisher has been joined, with {@code publishLock}
* held. A queued message is already encoded for this peer's protocol version - {@link
* #publish(ReplicationMsg)} did that before queueing it - so there is nothing to decide here
* beyond how long to keep trying.
* <p>
* The whole queue goes out under {@code publishLock}, because the publisher is not the only
* thread which writes this socket: a {@code ServerWriter} is joined only after the close which
* gets here (ServerHandler.shutdown() closes the session before joining it, and ServerReader's
* finally closes it before stopping the handler), and a {@code HeartbeatThread} is shut down
* after it too. Their {@code publish()} takes the synchronous branch once the session is off
* the queueing one, so without the lock a newer message could be written between two of these
* older ones - which a peer replication server answers by dropping the older ones at debug
* level, its log file refusing a record which would break its key ordering. The close takes
* the session off the queueing branch under the same lock it holds here, so such a
* {@code publish()} either queues ahead of that and is drained here, or waits for the drain and
* lands after it, or returns at the door, having seen the close. One which read the close as
* not yet begun can still be descheduled before its buffer is queued and queue it after the
* last poll below; it checks for that once it has, and takes the buffer back and reports it -
* see {@link #publish(ReplicationMsg)}. The wait for the lock itself is not part of the budget
* below, no more than it is for the {@code StopMsg} which follows.
* <p>
* The budget bounds how many messages a close spends on a peer which is reading slowly. A peer
* which is reading pays none of it; a peer which answers with a reset pays one failed write. It
* is checked between messages, so a single write which blocks past the budget still runs to
* completion - for a peer which has stopped reading, or which vanished without a reset, that is
* for as long as TCP keeps the connection alive: bounding it needs a non-blocking socket, which
* this session is not. The budget is per close, and a shutdown closes its sessions one after
* another - ReplicationServerDomain.stopAllServers() stops each handler in turn on one thread -
* so a domain with several peers which are alive and not reading pays it once per such peer.
* On the road which does not shut the whole server down the wait is paid under the lock of the
* replication domain - ReplicationServerDomain.stopServer() holds it across the handler
* shutdown which closes this session - where a handshake meanwhile waiting on that lock times
* out and the broker retries. What is given up on is reported rather than dropped in silence,
* that being the part of this which cost the most to diagnose.
*/
private void sendWhatThePublisherLeftQueued()
{
final long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(DRAIN_BUDGET_MS);
byte[] buffer;
while ((buffer = sendQueue.poll()) != null)
{
if (System.nanoTime() - deadline >= 0)
{
reportQueueNotSent(1 + takeWhatIsLeftQueued(), "the peer did not read them within "
+ DRAIN_BUDGET_MS + " ms");
return;
}
try
{
send(buffer);
}
catch (final IOException e)
{
/*
* send() has recorded the error; the rest of the queue cannot go out either. The
* exception is named rather than traced: this is reached whenever a peer which has
* announced it is leaving closes before the drain reaches it, where what the write
* failed with is the whole of what a reader of the log needs - a directory server
* re-reads these from the changelog when it reconnects, a replication server does not.
*/
reportQueueNotSent(1 + takeWhatIsLeftQueued(),
"the write failed with " + e.getClass().getName() + ": " + e.getMessage());
return;
}
}
}

/**
* Empties the queue a close gives up on, and answers how many buffers it held. Taking them
* rather than counting them is what keeps a {@code publish()} which queued a buffer late from
* reporting it a second time: it reports only a buffer it still finds in the queue.
*/
private int takeWhatIsLeftQueued()
{
int taken = 0;
while (sendQueue.poll() != null)
{
taken++;
}
return taken;
}

/** Says which messages a close could not hand to the peer, and why. */
private void reportQueueNotSent(final int count, final String reason)
{
logger.warn(LocalizableMessage.raw(
"The replication session %s was closed with %d message(s) which had been published on it "
+ "but not yet sent, and the peer was not told about them: %s",
getName(), count, reason));
}

/**
* This methods allows to determine if the session close was initiated
* on this Session.
Expand Down Expand Up @@ -322,6 +511,10 @@ public void publish(final ReplicationMsg msg) throws IOException
// Avoid blocking forever so that we can check for session closure.
if (sendQueue.offer(buffer, 100, TimeUnit.MILLISECONDS))
{
if (!isRunning.get())
{
takeBackWhatWasQueuedTooLate(buffer);
}
return;
}
}
Expand All @@ -338,6 +531,34 @@ public void publish(final ReplicationMsg msg) throws IOException
}
}

/**
* Takes a buffer back out of the queue if nothing is going to send it, and reports it.
* <p>
* {@code publish()} reads the close as not yet begun and queues the buffer after that, with no
* lock in between: descheduled there, it can queue the buffer after a close has stopped the
* publisher thread and drained the queue, and nothing would then send it or say so. Once the
* session is off the queueing branch, the lock here is the one the close holds while it takes
* the session off that branch and drains the queue, so what is still in the queue after it is
* what nothing sends. A buffer queued before the session came off the queueing branch is left
* to the drain - it cannot have seen the flag cleared - and a buffer the drain or the close
* already took is not found here, so nothing is reported twice.
*/
private void takeBackWhatWasQueuedTooLate(final byte[] buffer)
{
publishLock.lock();
try
{
if (sendQueue.remove(buffer))
{
reportQueueNotSent(1, "it was queued after the publisher of the session had stopped");
}
}
finally
{
publishLock.unlock();
}
}

/** Sends a replication message already encoded to the socket.
*
* @param buffer
Expand Down Expand Up @@ -526,7 +747,20 @@ private void setSessionError(final Exception e)
@Override
public void run()
{
isRunning.set(true);
synchronized (stateLock)
{
/*
* A close which came before the start has already run, and it is not coming back to clear
* the flag - nor is the end of this method, which leaves that to the close. Set, the flag
* would keep every later publish() on the queueing branch, which returns at once on a
* closed session as if the message had been queued; unset, publish() stays on the
* synchronous branch, which fails on the closed socket.
*/
if (!closeInitiated)
{
isRunning.set(true);
}
}
latch.countDown();
if (logger.isTraceEnabled())
{
Expand All @@ -551,10 +785,21 @@ public void run()
catch (IOException e)
{
setSessionError(e);
publisherFailedWrites.incrementAndGet();
needClosing = true;
}
}
isRunning.set(false);
/*
* A close clears the flag itself, under publishLock, once it has joined this thread - see
* close(). Clearing it here would open a window between the end of this thread and the drain
* of the close, in which a publish() takes the synchronous branch and writes a newer message
* ahead of the queue. Only a loop which ended without a close - an interrupt from elsewhere -
* clears it here, so that publish() does not go on queueing onto a queue nobody sends.
*/
if (!closeInitiated)
{
isRunning.set(false);
}
if (needClosing)
{
close();
Expand Down
Loading
Loading