From 6ae52a56b6447872587737c6df5fbaa33f72f023 Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Wed, 16 Sep 2026 13:36:19 +0300 Subject: [PATCH 1/2] [#1055] Report a ReplicaOfflineMsg forwarded once it is written to the peer, not once it is queued ServerWriter reported the forward as soon as Session.publish() had handed the message to the send queue of the session. A session thread busy writing an earlier buffer when the message was queued let the shutdown, released by that report, end its wait for a peer which had not been told and close the session: Session.close() sends what is still queued only once the write it joins has returned, and only within a budget of its own, not within the grace period the shutdown is meant to spend on that message. Session.publish(msg, whenWritten) runs the callback once, on the thread which wrote the message - the thread of the session, or the close for a message it sends out of the queue - after the write returned, and never for a message which was not written; it returns false for a message it neither wrote nor queued, and for one a close took back after draining the queue. ServerWriter reports the forward from that callback, and gives the peer up when the session refuses the message. --- .../server/replication/protocol/Session.java | 106 +++- .../replication/server/ServerWriter.java | 54 +- .../protocol/SessionPublisherDrainTest.java | 106 +++- .../replication/protocol/SessionTest.java | 245 +++++++++ .../ReplicationServerShutdownSyncTest.java | 483 ++++++++++++++++-- 5 files changed, 903 insertions(+), 91 deletions(-) create mode 100644 opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java index bd867c0220..cc381edfec 100644 --- a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java +++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java @@ -112,7 +112,20 @@ public final class Session extends DirectoryThread implements Closeable */ private BufferedOutputStream output; - private final LinkedBlockingQueue sendQueue = new LinkedBlockingQueue<>(4000); + /** A message queued for the thread of this session, and what to run once it is written. */ + private static final class Outgoing + { + private final byte[] buffer; + private final Runnable whenWritten; + + private Outgoing(byte[] buffer, Runnable whenWritten) + { + this.buffer = buffer; + this.whenWritten = whenWritten; + } + } + + private final LinkedBlockingQueue sendQueue = new LinkedBlockingQueue<>(4000); private AtomicBoolean isRunning = new AtomicBoolean(false); private final CountDownLatch latch = new CountDownLatch(1); @@ -166,8 +179,10 @@ public Session(final Socket socket, * This object won't be used anymore after this method is called. *

* 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()}. + * is sent here rather than dropped, within the budget of {@link #DRAIN_BUDGET_MS}, and its + * callback runs here once it is written. See {@link #sendWhatThePublisherLeftQueued()}. What the + * close gives up on is not written, and the callbacks of those messages never run - see + * {@link #publish(ReplicationMsg, Runnable)}. */ @Override public void close() @@ -341,8 +356,8 @@ public void close() private void sendWhatThePublisherLeftQueued() { final long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(DRAIN_BUDGET_MS); - byte[] buffer; - while ((buffer = sendQueue.poll()) != null) + Outgoing outgoing; + while ((outgoing = sendQueue.poll()) != null) { if (System.nanoTime() - deadline >= 0) { @@ -352,7 +367,7 @@ private void sendWhatThePublisherLeftQueued() } try { - send(buffer); + send(outgoing.buffer); } catch (final IOException e) { @@ -367,6 +382,7 @@ private void sendWhatThePublisherLeftQueued() "the write failed with " + e.getClass().getName() + ": " + e.getMessage()); return; } + written(outgoing.whenWritten); } } @@ -495,27 +511,53 @@ public boolean isEncrypted() * If an IO error occurred. */ public void publish(final ReplicationMsg msg) throws IOException + { + publish(msg, null); + } + + /** + * Sends a replication message to the remote peer, and runs the provided callback once the + * message has been written to the socket. + *

+ * While the thread of this session runs, a message published is queued for it and written + * later, so the return of this method says only that the message is queued. The callback is + * the only word that the message has left this server: it runs once, on the thread which wrote + * the message, after the write returned - the thread of the session, or the one closing it for + * a message the close sends out of the queue - and never for a message which was not written, + * which is what becomes of a message the write of which fails, and of what a close gives up on + * (see {@link #close()}). It must be short and must not block: the session writes nothing else + * until it returns. + * + * @param msg + * The message to be sent. + * @param whenWritten + * What to run once the message has been written, or null. + * @return whether the message was written or queued to be written; false when it was neither, + * because it has no encoding for the protocol version of the peer or because the + * session is being closed - the callback then never runs. A message queued after a + * close drained the queue is taken back, and counts as neither. + * @throws IOException + * If an IO error occurred. + */ + public boolean publish(final ReplicationMsg msg, final Runnable whenWritten) throws IOException { final byte[] buffer = msg.getBytes(protocolVersion); if (buffer == null) { // skip anything that cannot be encoded for this peer. - return; + return false; } if (isRunning.get()) { + final Outgoing outgoing = new Outgoing(buffer, whenWritten); while (!closeInitiated) { try { // Avoid blocking forever so that we can check for session closure. - if (sendQueue.offer(buffer, 100, TimeUnit.MILLISECONDS)) + if (sendQueue.offer(outgoing, 100, TimeUnit.MILLISECONDS)) { - if (!isRunning.get()) - { - takeBackWhatWasQueuedTooLate(buffer); - } - return; + return isRunning.get() || !takeBackWhatWasQueuedTooLate(outgoing); } } catch (final InterruptedException e) @@ -524,10 +566,27 @@ public void publish(final ReplicationMsg msg) throws IOException throw new IOException(e.getMessage()); } } + return false; } - else + send(buffer); + written(whenWritten); + return true; + } + + /** Runs what was to run once a message is written; a callback which fails takes nothing down. */ + private void written(final Runnable whenWritten) + { + if (whenWritten != null) { - send(buffer); + try + { + whenWritten.run(); + } + catch (final RuntimeException e) + { + logger.error(LocalizableMessage.raw("The callback of a message written to %s failed: %s", + readableRemoteAddress, stackTraceToSingleLineString(e))); + } } } @@ -542,16 +601,21 @@ public void publish(final ReplicationMsg msg) throws IOException * 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. + * + * @return whether the message was taken back - it is then never written, and its callback never + * runs */ - private void takeBackWhatWasQueuedTooLate(final byte[] buffer) + private boolean takeBackWhatWasQueuedTooLate(final Outgoing outgoing) { publishLock.lock(); try { - if (sendQueue.remove(buffer)) + if (sendQueue.remove(outgoing)) { reportQueueNotSent(1, "it was queued after the publisher of the session had stopped"); + return true; } + return false; } finally { @@ -769,10 +833,10 @@ public void run() boolean needClosing = false; while (!closeInitiated) { - byte[] buffer; + Outgoing outgoing; try { - buffer = sendQueue.take(); + outgoing = sendQueue.take(); } catch (InterruptedException ie) { @@ -780,14 +844,16 @@ public void run() } try { - send(buffer); + send(outgoing.buffer); } catch (IOException e) { setSessionError(e); publisherFailedWrites.incrementAndGet(); needClosing = true; + continue; } + written(outgoing.whenWritten); } /* * A close clears the flag itself, under publishLock, once it has joined this thread - see diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java index d1edac1b8e..9f68e4bbc1 100644 --- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java +++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ServerWriter.java @@ -17,9 +17,11 @@ */ package org.opends.server.replication.server; +import java.io.IOException; import java.net.SocketException; import org.forgerock.i18n.LocalizableMessage; +import org.forgerock.opendj.ldap.DN; import org.opends.server.api.DirectoryThread; import org.forgerock.i18n.slf4j.LocalizedLogger; import org.opends.server.replication.common.ServerStatus; @@ -120,24 +122,14 @@ public void run() replicationServerDomain.getBaseDN(), handler.getServerId()); } } + else if (updateMsg instanceof ReplicaOfflineMsg && !handler.isDataServer()) + { + forwardReplicaOfflineMsg((ReplicaOfflineMsg) updateMsg); + } else { // Publish the update to the remote server using a protocol version it supports session.publish(updateMsg); - /* - * Only the forward to a peer RS ends the wait of the shutdown: what the grace period - * buys is the rest of the topology learning that the replica went offline. A directory - * server is never handed this message - ReplicationServerDomain.put() does not queue - * it for one, and DataServerHandler.updateServerState() drops the one the changelog - * cursor of a directory server which is catching up synthesizes from the offline CSN - * of the replica (issue #1029) - so the guard says whose forward counts rather than - * telling two deliveries apart. - */ - if (updateMsg instanceof ReplicaOfflineMsg && !handler.isDataServer()) - { - dsrsShutdownSync.replicaOfflineMsgForwarded( - replicationServerDomain.getBaseDN(), updateMsg.getCSN(), handler.getServerId()); - } } } } @@ -170,6 +162,40 @@ public void run() } } + /** + * Publishes a ReplicaOfflineMsg to the peer replication server, and reports the forward to the + * shutdown which may be waiting for it. + *

+ * Only the forward to a peer RS ends the wait of the shutdown: what the grace period buys is + * the rest of the topology learning that the replica went offline. A directory server is never + * handed this message - ReplicationServerDomain.put() does not queue it for one, and + * DataServerHandler.updateServerState() drops the one the changelog cursor of a directory + * server which is catching up synthesizes from the offline CSN of the replica (issue #1029) - + * so the guard of the caller says whose forward counts rather than telling two deliveries apart. + *

+ * The forward is reported once the message has been written to the peer, not once it is queued + * for the thread of the session: the shutdown closes the session as soon as its wait ends, and + * Session.close() sends what is still queued only once the write it joins has returned, and + * only within a budget of its own, so a message reported forwarded while it was queued behind + * one the peer had not read yet would end the wait for a peer which had not been told, and + * leave its delivery to that budget rather than to the grace period. A message the session + * refuses - one published while the session is being closed - will never be written, and the + * shutdown must not wait for it. One the protocol version of the peer cannot carry is refused + * by the session as well, but does not get this far: isUpdateMsgFiltered() drops it and says so + * first (issue #1014). + */ + private void forwardReplicaOfflineMsg(final ReplicaOfflineMsg msg) throws IOException + { + final DN baseDN = replicationServerDomain.getBaseDN(); + final int serverId = handler.getServerId(); + final boolean accepted = session.publish(msg, + () -> dsrsShutdownSync.replicaOfflineMsgForwarded(baseDN, msg.getCSN(), serverId)); + if (!accepted) + { + dsrsShutdownSync.replicaOfflineMsgNotForwarded(baseDN, serverId); + } + } + private boolean isUpdateMsgFiltered(UpdateMsg updateMsg) { if (!updateMsg.isEncodableFor(handler.getProtocolVersion())) diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionPublisherDrainTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionPublisherDrainTest.java index 915a312ee4..d67b207b87 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionPublisherDrainTest.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionPublisherDrainTest.java @@ -19,9 +19,11 @@ import static org.opends.server.TestCaseUtils.TEST_ROOT_DN_STRING; import java.io.Closeable; +import java.lang.reflect.Constructor; import java.lang.reflect.Field; import java.net.ServerSocket; import java.net.Socket; +import java.util.ArrayList; import java.util.LinkedHashSet; import java.util.List; import java.util.Queue; @@ -33,6 +35,7 @@ import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import org.forgerock.opendj.ldap.DN; @@ -58,7 +61,9 @@ * still went out - leaving the peer with an orderly close and no sign that something was lost. * That is the limitation PR #919 recorded, and * {@code aSessionWithAPublisherThreadSendsWhatIsStillQueuedWhenItIsClosed} is what holds the - * close to sending that queue instead. The cases after it pin what a close gives up on: a + * close to sending that queue instead, and + * {@code aCloseRunsTheCallbackOfEachMessageItSendsOutOfTheQueue} to telling the publisher of each + * of those messages that it was written. The cases after it pin what a close gives up on: a * queue whose write fails is reported with every message it held; a session which had already * failed is written nothing more and still has its queue reported, together with what its * publisher took and failed to write, and nothing when there is nothing; and a closed session - @@ -317,12 +322,14 @@ public void aCloseWhichCannotSendTheQueueReportsEveryMessageTheQueueHeld() throw final Session receiver = pair[1]; try { - final Queue sendQueue = sendQueueOf(sender); + final AtomicInteger callbacks = new AtomicInteger(); + final Queue sendQueue = sendQueueOf(sender); for (int i = 0; i < MESSAGES_LEFT_UNSENT; i++) { - sendQueue.add(new DeleteMsg(DN.valueOf("uid=unsent" + i + "," + TEST_ROOT_DN_STRING), + sendQueue.add(queued(new DeleteMsg( + DN.valueOf("uid=unsent" + i + "," + TEST_ROOT_DN_STRING), csns.newCSN(), "00000000-0000-0000-0000-000000000000") - .getBytes(sender.getProtocolVersion())); + .getBytes(sender.getProtocolVersion()), callbacks::incrementAndGet)); } closeTheSocketsUnder(sender); @@ -347,6 +354,75 @@ public Void call() + "as well as for the ones left in it") .contains(MESSAGES_LEFT_UNSENT + " message(s)") .contains("the write failed with"); + assertThat(callbacks.get()) + .as("the callback of a message the close could not write ran, which tells its " + + "publisher that the peer was sent a message it never was") + .isZero(); + } + finally + { + StaticUtils.close(sender, receiver); + } + } + } + + /** + * A message the close sends out of the queue has its callback run, once, by the close, after it + * is written - which is how the writer of a replication server learns that a ReplicaOfflineMsg + * it had queued reached the peer only on the way out (issue #1055). Without that the message + * goes out and its publisher is never told, and a shutdown still waiting on it spends the rest + * of its grace period for a peer which was told. + *

+ * The queue is filled through the field of a session which was never started, so the close is + * the only thread which writes it and the callbacks run on the thread of this test. + */ + @Test + public void aCloseRunsTheCallbackOfEachMessageItSendsOutOfTheQueue() throws Exception + { + final CSNGenerator csns = new CSNGenerator(RS_ID, 0); + try (ServerSocket listen = new ServerSocket(0)) + { + final Session[] pair = connectSessionPair(listen); + final Session sender = pair[0]; + final Session receiver = pair[1]; + try + { + final List sent = new ArrayList<>(); + final List reported = new CopyOnWriteArrayList<>(); + final List reportedBy = new CopyOnWriteArrayList<>(); + final Queue sendQueue = sendQueueOf(sender); + for (int i = 0; i < MESSAGES_LEFT_UNSENT; i++) + { + final CSN csn = csns.newCSN(); + sent.add(csn); + sendQueue.add(queued(new DeleteMsg( + DN.valueOf("uid=drained" + i + "," + TEST_ROOT_DN_STRING), + csn, "00000000-0000-0000-0000-000000000000") + .getBytes(sender.getProtocolVersion()), new Runnable() + { + @Override + public void run() + { + reported.add(csn); + reportedBy.add(Thread.currentThread()); + } + })); + } + + sender.close(); + + assertThat(reported) + .as("the close sent the queue without running the callback of each message it wrote, " + + "once and in the order it wrote them") + .containsExactlyElementsOf(sent); + assertThat(reportedBy) + .as("the callbacks ran on a thread other than the one which closed the session") + .containsOnly(Thread.currentThread()); + final Drained drained = drain(receiver); + assertThat(drained.received) + .as("the peer did not receive what the callbacks report as written; the read ended " + + "by %s", drained.endedBy) + .containsExactlyElementsOf(sent); } finally { @@ -375,12 +451,13 @@ public void aCloseOfAFailedSessionWritesNothingAndReportsTheQueue() throws Excep final Session receiver = pair[1]; try { - final Queue sendQueue = sendQueueOf(sender); + final Queue sendQueue = sendQueueOf(sender); for (int i = 0; i < MESSAGES_LEFT_UNSENT; i++) { - sendQueue.add(new DeleteMsg(DN.valueOf("uid=failed" + i + "," + TEST_ROOT_DN_STRING), + sendQueue.add(queued(new DeleteMsg( + DN.valueOf("uid=failed" + i + "," + TEST_ROOT_DN_STRING), csns.newCSN(), "00000000-0000-0000-0000-000000000000") - .getBytes(sender.getProtocolVersion())); + .getBytes(sender.getProtocolVersion()), null)); } final Field sessionError = Session.class.getDeclaredField("sessionError"); sessionError.setAccessible(true); @@ -497,7 +574,7 @@ public void aCloseOfAStartedSessionWhoseWritesFailedReportsWhatThePublisherTookA sender.publish(new DeleteMsg(DN.valueOf("uid=takenandlost" + i + "," + TEST_ROOT_DN_STRING), csns.newCSN(), "00000000-0000-0000-0000-000000000000")); } - final Queue sendQueue = sendQueueOf(sender); + final Queue sendQueue = sendQueueOf(sender); final long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(SOCKET_TIMEOUT_MS); while (!sendQueue.isEmpty() && System.nanoTime() - deadline < 0) { @@ -693,11 +770,20 @@ else if (msg instanceof StopMsg) /** The queue a started session's publisher thread takes its buffers from. */ @SuppressWarnings("unchecked") - private static Queue sendQueueOf(final Session session) throws Exception + private static Queue sendQueueOf(final Session session) throws Exception { final Field sendQueue = Session.class.getDeclaredField("sendQueue"); sendQueue.setAccessible(true); - return (Queue) sendQueue.get(session); + return (Queue) sendQueue.get(session); + } + + /** An already encoded buffer as {@code publish()} queues it, with what to run once written. */ + private static Object queued(final byte[] buffer, final Runnable whenWritten) throws Exception + { + final Constructor outgoing = Class.forName(Session.class.getName() + "$Outgoing") + .getDeclaredConstructor(byte[].class, Runnable.class); + outgoing.setAccessible(true); + return outgoing.newInstance(buffer, whenWritten); } /** diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java new file mode 100644 index 0000000000..7525fdd7a3 --- /dev/null +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java @@ -0,0 +1,245 @@ +/* + * The contents of this file are subject to the terms of the Common Development and + * Distribution License (the License). You may not use this file except in compliance with the + * License. + * + * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the + * specific language governing permission and limitations under the License. + * + * When distributing Covered Software, include this CDDL Header Notice in each file and include + * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL + * Header, with the fields enclosed by brackets [] replaced by your own identifying + * information: "Portions copyright [year] [name of copyright owner]". + * + * Copyright 2026 3A Systems, LLC. + */ +package org.opends.server.replication.protocol; + +import static org.assertj.core.api.Assertions.assertThat; + +import java.io.Closeable; +import java.net.InetSocketAddress; +import java.net.ServerSocket; +import java.net.Socket; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; + +import org.opends.server.TestCaseUtils; +import org.opends.server.replication.ReplicationTestCase; +import org.opends.server.replication.common.CSN; +import org.opends.server.replication.common.CSNGenerator; +import org.opends.server.util.StaticUtils; +import org.testng.annotations.Test; + +/** + * A message published to a running session is handed to the session thread, which writes it to + * the socket later. A caller which needs to know when that happened - the replication server + * forwarding a ReplicaOfflineMsg, whose shutdown must not close the session before the message + * is on the wire - attaches a callback to the message, and the session runs it once, from the + * thread which wrote it, only after the write returned. + *

+ * The peer of each test reads nothing until the test lets it, and both ends of the connection + * have socket buffers far smaller than the first message published, so that the session thread + * of the end under test is held inside the write of that message for as long as the test wants + * - the state in which a message published behind it is queued and not written. + */ +@SuppressWarnings("javadoc") +public class SessionTest extends ReplicationTestCase +{ + private static final int SOCKET_TIMEOUT_MS = 30000; + /** + * Socket buffers small enough that {@link #BLOCKING_MESSAGE_SIZE} bytes cannot be written + * through them: the write blocks until the peer reads. Set explicitly on both ends, since the + * buffers the kernel picks on its own grow well beyond it on a loopback link - and the message + * is larger by far than what they hold, since a kernel which does not honour the size asked + * for on the receiving side, as macOS does not, must still be unable to take the whole of it. + */ + private static final int SOCKET_BUFFER_SIZE = 8 * 1024; + private static final int BLOCKING_MESSAGE_SIZE = 4 * 1024 * 1024; + /** Time given to a callback which must not run, to see that it does not. */ + private static final long SETTLE_MS = 500; + private static final int SENDER_ID = 1; + private static final int PEER_ID = 2; + + @Test + public void theCallbackRunsOnceTheMessageIsWrittenAndNotWhenItIsQueued() throws Exception + { + try (SessionPair pair = connectSessionPair()) + { + pair.publisher.start(); + pair.publisher.waitForStartup(); + + // The session thread is inside the write of this message until the peer reads it. + pair.publisher.publish(newBlockingMsg()); + pair.awaitBytesReachedThePeer(); + + final CountDownLatch written = new CountDownLatch(1); + final boolean accepted = pair.publisher.publish(new HeartbeatMsg(), written::countDown); + + assertThat(accepted).as("the message was refused by a running session").isTrue(); + assertThat(written.await(SETTLE_MS, TimeUnit.MILLISECONDS)) + .as("the callback ran while the message was still queued behind a message the peer " + + "had not read") + .isFalse(); + + assertThat(pair.peer.receive()).isInstanceOf(EntryMsg.class); + assertThat(written.await(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS)) + .as("the callback did not run once the message had been written") + .isTrue(); + assertThat(pair.peer.receive()).isInstanceOf(HeartbeatMsg.class); + } + } + + @Test + public void aMessageWhichCannotBeEncodedForThePeerIsRefusedWithoutRunningTheCallback() + throws Exception + { + try (SessionPair pair = connectSessionPair()) + { + pair.publisher.start(); + pair.publisher.waitForStartup(); + // A ReplicaOfflineMsg has no encoding before protocol version 8. + pair.publisher.setProtocolVersion(ProtocolVersion.REPLICATION_PROTOCOL_V7); + + final CountDownLatch written = new CountDownLatch(1); + final boolean accepted = + pair.publisher.publish(new ReplicaOfflineMsg(newCSN()), written::countDown); + + assertThat(accepted) + .as("a message the peer cannot decode was reported as accepted") + .isFalse(); + assertThat(written.await(SETTLE_MS, TimeUnit.MILLISECONDS)) + .as("the callback ran for a message which was never written") + .isFalse(); + } + } + + /** + * Before its thread is started, and once that thread is gone, a session writes on the + * publishing thread itself; the callback then runs on that same thread, after the write. + */ + @Test + public void theCallbackRunsAfterAMessageWrittenOnThePublishingThread() throws Exception + { + try (SessionPair pair = connectSessionPair()) + { + final CountDownLatch written = new CountDownLatch(1); + final boolean accepted = pair.publisher.publish(new HeartbeatMsg(), written::countDown); + + assertThat(accepted).as("the message was refused by a session with no thread").isTrue(); + assertThat(written.getCount()) + .as("the callback had not run when the publish which wrote the message returned") + .isZero(); + assertThat(pair.peer.receive()).isInstanceOf(HeartbeatMsg.class); + } + } + + private static EntryMsg newBlockingMsg() + { + return new EntryMsg(SENDER_ID, PEER_ID, new byte[BLOCKING_MESSAGE_SIZE], 1); + } + + private static CSN newCSN() + { + return new CSNGenerator(SENDER_ID, 0).newCSN(); + } + + /** + * Connects the end under test, in the server role of the replication protocol, with a peer + * which reads only when a test does. Both ends exchange one message under TLS and then drop the + * security layer, as the replication handshake does when encryption is not required, so that + * the socket buffers alone decide when a write blocks: a message read under TLS by each end is + * what consumes the records TLS itself sends after its negotiation, which would otherwise be + * read as the start of a replication message once the layer is gone. + */ + private static SessionPair connectSessionPair() throws Exception + { + final ReplSessionSecurity security = getReplSessionSecurity(); + final ExecutorService executor = Executors.newSingleThreadExecutor(); + final Socket peerSocket = new Socket(); + Socket publisherSocket = null; + Session publisher = null; + boolean connected = false; + try (ServerSocket listen = TestCaseUtils.bindFreePort()) + { + listen.setSoTimeout(SOCKET_TIMEOUT_MS); + peerSocket.setReceiveBufferSize(SOCKET_BUFFER_SIZE); + peerSocket.setTcpNoDelay(true); + peerSocket.connect(new InetSocketAddress("127.0.0.1", listen.getLocalPort()), SOCKET_TIMEOUT_MS); + // The TLS negotiation needs both ends handshaking at the same time. + final Future peerEnd = + executor.submit(() -> security.createClientSession(peerSocket, SOCKET_TIMEOUT_MS)); + + publisherSocket = listen.accept(); + publisherSocket.setSendBufferSize(SOCKET_BUFFER_SIZE); + publisherSocket.setTcpNoDelay(true); + publisher = security.createServerSession(publisherSocket, SOCKET_TIMEOUT_MS); + assertThat(publisher).as("could not create the session under test").isNotNull(); + final Session peer = peerEnd.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS); + + publisher.publish(new HeartbeatMsg()); + assertThat(peer.receive()).isInstanceOf(HeartbeatMsg.class); + peer.publish(new HeartbeatMsg()); + assertThat(publisher.receive()).isInstanceOf(HeartbeatMsg.class); + publisher.stopEncryption(); + peer.stopEncryption(); + connected = true; + return new SessionPair(publisher, peer, peerSocket); + } + finally + { + executor.shutdownNow(); + if (!connected) + { + if (publisher != null) + { + publisher.close(); + } + StaticUtils.close(publisherSocket, peerSocket); + } + } + } + + private static final class SessionPair implements Closeable + { + private final Session publisher; + private final Session peer; + private final Socket peerSocket; + + private SessionPair(Session publisher, Session peer, Socket peerSocket) + { + this.publisher = publisher; + this.peer = peer; + this.peerSocket = peerSocket; + } + + /** + * Waits for the first bytes of a message to reach the peer: the session thread of the end + * under test is then inside the write of that message, and stays there until the peer + * reads, since the message is larger than the buffers on both sides of the connection. + */ + void awaitBytesReachedThePeer() throws Exception + { + final long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(SOCKET_TIMEOUT_MS); + while (peerSocket.getInputStream().available() == 0) + { + assertThat(System.nanoTime() < deadline) + .as("nothing was written to the peer") + .isTrue(); + Thread.sleep(10); + } + } + + @Override + public void close() + { + // The peer first: a session thread held inside a write is released by the peer going away, + // and close() joins that thread. + peer.close(); + publisher.close(); + } + } +} diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java index 31eabf6874..b532ed27de 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/ReplicationServerShutdownSyncTest.java @@ -25,6 +25,7 @@ import java.net.ServerSocket; import java.net.Socket; import java.util.ArrayList; +import java.util.Arrays; import java.util.Collection; import java.util.List; import java.util.TreeSet; @@ -40,6 +41,7 @@ import java.util.concurrent.atomic.AtomicReference; import org.forgerock.opendj.ldap.DN; +import org.forgerock.opendj.ldap.ModificationType; import org.opends.server.TestCaseUtils; import org.opends.server.core.DirectoryServer; import org.opends.server.replication.ReplicationTestCase; @@ -48,6 +50,7 @@ import org.opends.server.replication.common.RSInfo; import org.opends.server.replication.common.ServerState; import org.opends.server.replication.protocol.DeleteMsg; +import org.opends.server.replication.protocol.ModifyMsg; import org.opends.server.replication.protocol.ProtocolVersion; import org.opends.server.replication.protocol.ReplServerStartMsg; import org.opends.server.replication.protocol.ReplSessionSecurity; @@ -59,6 +62,8 @@ import org.opends.server.replication.protocol.WindowMsg; import org.opends.server.replication.service.DSRSShutdownSync; import org.opends.server.replication.service.ReplicationBroker; +import org.opends.server.types.Attributes; +import org.opends.server.types.Modification; import org.opends.server.util.StaticUtils; import org.opends.server.util.TestTimer; import org.testng.annotations.DataProvider; @@ -109,8 +114,38 @@ public class ReplicationServerShutdownSyncTest extends ReplicationTestCase * share a server id. */ private static final int PEER_VERSION_RS_ID = 8240; + /** The peer replication server whose session thread is busy writing an earlier change. */ + private static final int BUSY_RS_ID = 100; + /** The peer replication server whose protocol version predates the ReplicaOfflineMsg. */ + private static final int LEGACY_RS_ID = 101; /** Send window a peer advertises when nothing has to hold its writer back. */ private static final int PEER_WINDOW = 100; + /** + * Socket buffers of the connection to the peer which does not read: small enough that + * {@link #SOCKET_FILLING_CHANGE_SIZE} bytes cannot be written through them, so that the session + * thread writing that change is held inside the write until the peer reads. Set on both ends, + * since the buffers the kernel picks on its own grow well beyond it on a loopback link. + */ + private static final int SMALL_SOCKET_BUFFER_SIZE = 8 * 1024; + /** + * Larger by far than what the socket buffers hold: a kernel which does not honour the size + * asked for on the receiving side - macOS keeps a few hundred kilobytes there - must still + * be unable to take the whole change. + */ + private static final int SOCKET_FILLING_CHANGE_SIZE = 4 * 1024 * 1024; + /** + * What the peer which does not read must have been sent, and not read, for the session thread + * serving it to be inside a write: well below its receive buffer, since the kernel advertises + * less than the whole of it, and well above any of the small messages a replication server + * sends a peer on its own. + */ + private static final int SOCKET_BUFFER_FILL_MARK = SMALL_SOCKET_BUFFER_SIZE / 4; + /** + * Time a shutdown released by the ReplicaOfflineMsg being queued, rather than written, is given + * to close the session of the peer. A shutdown which waits for the write cannot close the + * session before the peer reads, so it spends this time and no more. + */ + private static final long EARLY_CLOSE_TIMEOUT_MS = 1000; /** * Send window of the peer which is held back: one change fills it, and the message which * follows stays with its writer until the peer gives it credit again. @@ -578,6 +613,98 @@ public void call() throws Exception } } + /** + * The forward the shutdown waits for must mean that the message has been written to the peer, + * not that it has been queued for the thread of its session: a forward reported on the queue + * ends the wait of the shutdown before the peer has been told, and leaves the message to what + * Session.close() sends of the queue once the write it joins has returned, within a budget of + * its own rather than the grace period. The session thread is busy with an earlier message when + * the ReplicaOfflineMsg is queued behind it whenever the peer reads slower than the + * replication server writes. Here the peer does not read at all, and the + * socket buffers on both sides of its connection are far smaller than the change which fills + * them, so the session thread is held inside the write of that change until the test lets the + * peer read. + */ + @Test + public void thePeerStillReadingAnEarlierChangeIsToldTheReplicaWentOfflineBeforeItIsStopped() + throws Exception + { + final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING); + final RecordingShutdownSync shutdownSync = new RecordingShutdownSync(); + final ExecutorService executor = Executors.newFixedThreadPool(2); + ReplicationServer replicationServer = null; + ReplicationBroker broker = null; + FakePeerReplicationServer peer = null; + Future shutdown = null; + try (ServerSocket listen = TestCaseUtils.bindFreePort()) + { + listen.setSoTimeout(SOCKET_TIMEOUT_MS); + final int replicationPort = TestCaseUtils.findFreePort(); + replicationServer = + newReplicationServer(shutdownSync, "shutdownSyncBusySessionDb", 8238, replicationPort); + broker = + openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID); + + final ConnectedSessions connection = + connectSessionPair(listen, getReplSessionSecurity(), SMALL_SOCKET_BUFFER_SIZE); + final Future served = + serveAsTheListenThreadWould(replicationServer, connection.localEnd, executor); + peer = FakePeerReplicationServer.connected(connection.remoteEnd, connection.remoteSocket, + BUSY_RS_ID, baseDN, EMPTY_DN_GENID, PEER_WINDOW); + served.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS); + final ReplicationServerDomain domain = + replicationServer.getReplicationServerDomain(baseDN, true); + waitForConnectedReplicationServer(domain, BUSY_RS_ID); + + /* + * One generator for the change and the announcement: a CSN which does not follow the one + * of the change would be dropped by the handler of the peer as already seen. + */ + final CSNGenerator csns = new CSNGenerator(LOCAL_DS_ID, 0); + broker.publish(newChangeLargerThanTheSocketBuffers(csns.newCSN())); + peer.awaitReceiveBufferFilled(); + + final CSN offlineCSN = csns.newCSN(); + shutdownSync.replicaOfflineMsgSent(baseDN, offlineCSN); + broker.publish(new ReplicaOfflineMsg(offlineCSN)); + shutdownSync.awaitDispatch(); + + shutdown = executor.submit(newShutdown(replicationServer)); + awaitCloseInitiated(connection.localEnd, EARLY_CLOSE_TIMEOUT_MS); + final List forwardedBeforeThePeerRead = new ArrayList<>(shutdownSync.forwardedBy()); + + final Future received = peer.receive(ReplicaOfflineMsg.class); + final ReplicaOfflineMsg forwarded = received.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS); + final long elapsed = shutdown.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS); + + assertThat(forwarded) + .as("the peer was never told that the replica went offline: its session was closed " + + "with the message still queued behind the change it was reading, and its read " + + "ended with: %s (forward reported by %s, the shutdown took %d ms)", + peer.failure(), shutdownSync.forwardedBy(), elapsed) + .isNotNull(); + assertThat(forwardedBeforeThePeerRead) + .as("the writer reported the message forwarded while it was still queued behind a " + + "change the peer had not read") + .doesNotContain(BUSY_RS_ID); + assertThat(shutdownSync.forwardedBy()) + .as("the message was written to the peer and nothing reported the forward") + .contains(BUSY_RS_ID); + assertThat(elapsed) + .as("the shutdown waited out the grace period after the message had been written") + .isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD); + } + finally + { + // Closing the peer releases a session thread held inside a write, and the shutdown with it. + closeQuietly(peer); + awaitQuietly(shutdown); + stop(broker); + removeQuietly(replicationServer); + executor.shutdownNow(); + } + } + /** * Only a peer replication server learning about the offline replica ends the wait, and a * directory server is never told. ReplicationServerDomain.put() never queues a @@ -725,14 +852,11 @@ public void theShutdownWaitsForEveryPeerToBeToldTheReplicaWentOffline() throws E + "with: %s", peer.failure()) .isNotNull(); /* - * The forward asserted above proves the message reached the Session, not the wire: close() - * now sends what its publisher left queued, but only within its own budget, and it writes - * that queue under publishLock so that a message published meanwhile lands after it rather - * than between two of its own. If this is the only assertion which fails, the close is - * where to look before the granularity of the barrier: the warning close() writes for a - * queue it could not hand over says why it gave that queue up, and its absence does not - * prove the message left this end - a publish() concurrent with the close is still dropped - * at the door without one. + * The forward asserted above proves the message was written to the socket of the peer: the + * writer reports it from the callback of Session.publish(), which runs only once the write + * has returned, whether the thread of the session wrote it or the close sent it out of the + * queue. If this is the only assertion which fails, the read of the peer is where to look + * before the granularity of the barrier. */ assertThat(receivedWhenHeldBack.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS)) .as("the peer which was held back never learned that the replica went offline, " @@ -827,6 +951,65 @@ public void theShutdownStopsWaitingForAPeerWhoseMessageTheWriterDropped() throws } } + /** + * A peer whose protocol version has no encoding for the ReplicaOfflineMsg cannot be told that + * the replica went offline - it predates the message, and has nothing to do with it - and + * nothing will ever report a forward to it: the writer drops the message before the session is + * given it (issue #1014), and must strike the peer off rather than let the shutdown wait out the + * grace period for it. + */ + @Test + public void theShutdownStopsWaitingForAPeerWhoseProtocolCannotCarryTheMessage() throws Exception + { + final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING); + final RecordingShutdownSync shutdownSync = new RecordingShutdownSync(); + ReplicationServer replicationServer = null; + ReplicationBroker broker = null; + FakePeerReplicationServer peer = null; + try + { + final int replicationPort = TestCaseUtils.findFreePort(); + replicationServer = newReplicationServer( + shutdownSync, "shutdownSyncLegacyProtocolDb", 8239, replicationPort); + broker = + openReplicationSession(baseDN, LOCAL_DS_ID, 100, replicationPort, 5000, EMPTY_DN_GENID); + peer = FakePeerReplicationServer.connected(replicationPort, LEGACY_RS_ID, baseDN, EMPTY_DN_GENID, + PEER_WINDOW, ProtocolVersion.REPLICATION_PROTOCOL_V7); + + final ReplicationServerDomain domain = + replicationServer.getReplicationServerDomain(baseDN, true); + waitForConnectedReplicationServer(domain, LEGACY_RS_ID); + + final CSN offlineCSN = newOfflineCSN(); + shutdownSync.replicaOfflineMsgSent(baseDN, offlineCSN); + broker.publish(new ReplicaOfflineMsg(offlineCSN)); + shutdownSync.awaitDispatch(); + awaitGiveUpOn(shutdownSync, LEGACY_RS_ID, + "the writer let the shutdown wait for a peer whose protocol cannot carry the message"); + + final long startTime = System.nanoTime(); + replicationServer.shutdown(); + final long elapsed = elapsedMillis(startTime); + + assertThat(shutdownSync.dispatchedTo()) + .as("the message was not queued for the peer, so this test never reproduced the " + + "drop it is about") + .contains(LEGACY_RS_ID); + assertThat(shutdownSync.forwardedBy()) + .as("a message the peer cannot decode was reported forwarded to it") + .doesNotContain(LEGACY_RS_ID); + assertThat(elapsed) + .as("the shutdown waited for a forward to a peer whose protocol cannot carry the message") + .isLessThan(DSRSShutdownSync.REPLICA_OFFLINE_GRACE_PERIOD); + } + finally + { + closeQuietly(peer); + stop(broker); + removeQuietly(replicationServer); + } + } + /** * A peer whose handshake is aborted while the message is being pushed must not be waited for. *

@@ -1583,6 +1766,73 @@ private static CSN newOfflineCSN(int serverId) return new CSNGenerator(serverId, 0).newCSN(); } + /** + * A change of the collocated replica larger than the socket buffers of the connection to the + * peer which does not read, so that the session thread writing it to that peer is held inside + * the write. + */ + private static ModifyMsg newChangeLargerThanTheSocketBuffers(CSN csn) + { + final char[] value = new char[SOCKET_FILLING_CHANGE_SIZE]; + Arrays.fill(value, 'x'); + final List mods = newArrayList( + new Modification(ModificationType.REPLACE, Attributes.create("description", new String(value)))); + return new ModifyMsg(csn, DN.valueOf("uid=busy," + TEST_ROOT_DN_STRING), mods, "busy-entry-uuid"); + } + + /** The shutdown of the replication server, reporting how long it took. */ + private static Callable newShutdown(final ReplicationServer replicationServer) + { + return new Callable() + { + @Override + public Long call() + { + final long startTime = System.nanoTime(); + replicationServer.shutdown(); + return elapsedMillis(startTime); + } + }; + } + + /** + * Serves a connection to the replication server as its listen thread does - see + * ReplicationServer.runListen() - over a session the test established itself, so that the + * sockets underneath are its own to configure: the start message of the peer is read, and the + * handler is created and started from it. The start blocks until the handshake is over, so it + * runs on a thread of its own, as it does on the listen thread. + */ + private static Future serveAsTheListenThreadWould( + final ReplicationServer replicationServer, final Session session, ExecutorService executor) + { + return executor.submit(new Callable() + { + @Override + public ReplicationServerHandler call() throws Exception + { + final ReplServerStartMsg startMsg = (ReplServerStartMsg) session.receive(); + final ReplicationServerHandler rsHandler = + new ReplicationServerHandler(session, 100, replicationServer, 100); + rsHandler.startFromRemoteRS(startMsg); + return rsHandler; + } + }); + } + + /** + * Waits for the close of the session to have been initiated, and gives up quietly once the + * timeout is over: the caller says what a close within the timeout, or none, means. + */ + private static void awaitCloseInitiated(Session session, long timeoutMillis) + throws InterruptedException + { + final long deadline = System.nanoTime() + TimeUnit.MILLISECONDS.toNanos(timeoutMillis); + while (!session.closeInitiated() && System.nanoTime() < deadline) + { + Thread.sleep(10); + } + } + private static boolean sleepQuietly(long millis) { try @@ -1632,6 +1882,24 @@ private void closeQuietly(FakePeerReplicationServer peer) } } + private static void awaitQuietly(Future future) + { + if (future != null) + { + try + { + future.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS); + } + catch (InterruptedException e) + { + Thread.currentThread().interrupt(); + } + catch (Exception ignored) + { + } + } + } + /** * Establishes a connected session pair over the given listen socket, as a remote server * connecting to the RS would. The TLS negotiation performed by the session factories needs both @@ -1643,8 +1911,20 @@ private void closeQuietly(FakePeerReplicationServer peer) private Session[] connectSessionPair(ServerSocket listenSocket, final ReplSessionSecurity security) throws Exception { - final Socket clientSocket = new Socket("127.0.0.1", listenSocket.getLocalPort()); - clientSocket.setTcpNoDelay(true); + final ConnectedSessions connection = connectSessionPair(listenSocket, security, 0); + return new Session[] { connection.remoteEnd, connection.localEnd }; + } + + /** + * Establishes a connected session pair over the given listen socket, with the send buffer of + * the local end and the receive buffer of the remote end bounded by the given size: a message + * larger than both then holds the thread writing it until the remote end reads. A size of 0 + * leaves the buffers to the kernel. + */ + private ConnectedSessions connectSessionPair(ServerSocket listenSocket, + final ReplSessionSecurity security, int socketBufferSize) throws Exception + { + final Socket clientSocket = new Socket(); final ExecutorService executor = Executors.newSingleThreadExecutor(); Future clientEnd = null; Socket serverSocket = null; @@ -1652,6 +1932,14 @@ private Session[] connectSessionPair(ServerSocket listenSocket, final ReplSessio boolean connected = false; try { + if (socketBufferSize > 0) + { + // Before the connection is made: the window the local end is told is sized from it. + clientSocket.setReceiveBufferSize(socketBufferSize); + } + clientSocket.setTcpNoDelay(true); + clientSocket.connect( + new InetSocketAddress("127.0.0.1", listenSocket.getLocalPort()), SOCKET_TIMEOUT_MS); clientEnd = executor.submit(new Callable() { @Override @@ -1662,14 +1950,18 @@ public Session call() throws Exception }); serverSocket = listenSocket.accept(); + if (socketBufferSize > 0) + { + serverSocket.setSendBufferSize(socketBufferSize); + } serverSocket.setTcpNoDelay(true); serverEnd = security.createServerSession(serverSocket, SOCKET_TIMEOUT_MS); assertThat(serverEnd).as("could not create a session for the handler under test").isNotNull(); - final Session[] sessionPair = - new Session[] { clientEnd.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS), serverEnd }; + final ConnectedSessions connection = new ConnectedSessions( + clientEnd.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS), clientSocket, serverEnd); connected = true; - return sessionPair; + return connection; } finally { @@ -1683,6 +1975,21 @@ public Session call() throws Exception } } + /** The two ends of a connection to the replication server, and the socket of the remote one. */ + private static final class ConnectedSessions + { + private final Session remoteEnd; + private final Socket remoteSocket; + private final Session localEnd; + + private ConnectedSessions(Session remoteEnd, Socket remoteSocket, Session localEnd) + { + this.remoteEnd = remoteEnd; + this.remoteSocket = remoteSocket; + this.localEnd = localEnd; + } + } + private void closeServerEndQuietly(Session serverEnd, Socket serverSocket) { if (serverEnd != null) @@ -1815,6 +2122,9 @@ private static final class FakePeerReplicationServer private final long generationId; private final String serverURL; private final Session session; + /** The socket under the session: what it has received and not read is what the replication + * server has written to this peer. */ + private final Socket socket; private final ExecutorService reader = Executors.newSingleThreadExecutor(); /** * What ended the exchange with the replication server, so that a message which never arrived @@ -1849,8 +2159,32 @@ static FakePeerReplicationServer connected(int replicationPort, int serverId, DN static FakePeerReplicationServer connected(int replicationPort, int serverId, DN baseDN, long generationId, int windowSize, short protocolVersion) throws Exception { - final FakePeerReplicationServer peer = new FakePeerReplicationServer( - replicationPort, serverId, baseDN, generationId, windowSize, protocolVersion); + return completed(new FakePeerReplicationServer( + replicationPort, serverId, baseDN, generationId, windowSize, protocolVersion)); + } + + /** + * A connected peer over a session the test established itself - one whose sockets it + * configured - which the replication server serves as its listen thread would. + */ + static FakePeerReplicationServer connected(Session newSession, Socket newSocket, int serverId, + DN baseDN, long generationId, int windowSize) throws Exception + { + return completed(new FakePeerReplicationServer(newSession, newSocket, serverId, baseDN, + generationId, windowSize, ProtocolVersion.getCurrentVersion())); + } + + /** A peer whose handshake stops after its first phase, before it sends its TopologyMsg. */ + static FakePeerReplicationServer handshaking( + int replicationPort, int serverId, DN baseDN, long generationId) throws Exception + { + return new FakePeerReplicationServer(replicationPort, serverId, baseDN, generationId, + PEER_WINDOW, ProtocolVersion.getCurrentVersion()); + } + + private static FakePeerReplicationServer completed(FakePeerReplicationServer peer) + throws Exception + { boolean handshaken = false; try { @@ -1868,60 +2202,115 @@ static FakePeerReplicationServer connected(int replicationPort, int serverId, DN return peer; } - /** A peer whose handshake stops after its first phase, before it sends its TopologyMsg. */ - static FakePeerReplicationServer handshaking( - int replicationPort, int serverId, DN baseDN, long generationId) throws Exception - { - return new FakePeerReplicationServer(replicationPort, serverId, baseDN, generationId, - PEER_WINDOW, ProtocolVersion.getCurrentVersion()); - } - private FakePeerReplicationServer(int replicationPort, int serverId, DN baseDN, long generationId, int windowSize, short protocolVersion) throws Exception { this.serverId = serverId; this.generationId = generationId; - final Socket socket = new Socket(); + final Socket newSocket = new Socket(); Session newSession = null; String newServerURL = null; boolean started = false; try { - socket.setTcpNoDelay(true); - socket.connect(new InetSocketAddress("127.0.0.1", replicationPort), SOCKET_TIMEOUT_MS); - newSession = getReplSessionSecurity().createClientSession(socket, SOCKET_TIMEOUT_MS); - // the version this peer speaks: the replication server negotiates the oldest of the two - newSession.setProtocolVersion(protocolVersion); - - newServerURL = "127.0.0.1:" + socket.getLocalPort(); - newSession.publish(new ReplServerStartMsg(serverId, newServerURL, baseDN, windowSize, - new ServerState(), generationId, false, GROUP_ID, 5000)); - final ReplServerStartMsg inStartMsg = - waitForSpecificMsg(newSession, ReplServerStartMsg.class); - if (!inStartMsg.getSSLEncryption()) + newSocket.setTcpNoDelay(true); + newSocket.connect(new InetSocketAddress("127.0.0.1", replicationPort), SOCKET_TIMEOUT_MS); + newSession = getReplSessionSecurity().createClientSession(newSocket, SOCKET_TIMEOUT_MS); + newServerURL = start(newSession, newSocket, serverId, baseDN, generationId, windowSize, + protocolVersion); + started = true; + } + finally + { + if (!started) { - newSession.stopEncryption(); + abandon(newSession, newSocket); } + } + serverURL = newServerURL; + session = newSession; + socket = newSocket; + } + + private FakePeerReplicationServer(Session newSession, Socket newSocket, int serverId, + DN baseDN, long generationId, int windowSize, short protocolVersion) throws Exception + { + this.serverId = serverId; + this.generationId = generationId; + String newServerURL = null; + boolean started = false; + try + { + newServerURL = start(newSession, newSocket, serverId, baseDN, generationId, windowSize, + protocolVersion); started = true; } finally { if (!started) { - // The caller has no handle on this peer yet, so nothing else would close it. - reader.shutdownNow(); - if (newSession != null) - { - newSession.close(); - } - else - { - StaticUtils.close(socket); - } + abandon(newSession, newSocket); } } serverURL = newServerURL; session = newSession; + socket = newSocket; + } + + /** + * Runs the first phase of the handshake, the exchange of the start messages, and returns the + * URL this peer announced itself under. + */ + private static String start(Session newSession, Socket newSocket, int serverId, DN baseDN, + long generationId, int windowSize, short protocolVersion) throws Exception + { + // The replication server speaks the older of the two versions from the start message on. + newSession.setProtocolVersion(protocolVersion); + final String newServerURL = "127.0.0.1:" + newSocket.getLocalPort(); + newSession.publish(new ReplServerStartMsg(serverId, newServerURL, baseDN, windowSize, + new ServerState(), generationId, false, GROUP_ID, 5000)); + final ReplServerStartMsg inStartMsg = + waitForSpecificMsg(newSession, ReplServerStartMsg.class); + if (!inStartMsg.getSSLEncryption()) + { + newSession.stopEncryption(); + } + return newServerURL; + } + + /** The caller has no handle on this peer yet, so nothing else would close it. */ + private void abandon(Session newSession, Socket newSocket) + { + reader.shutdownNow(); + if (newSession != null) + { + newSession.close(); + } + else + { + StaticUtils.close(newSocket); + } + } + + /** + * Waits for the replication server to have filled the receive buffer of this peer, which + * reads nothing meanwhile, up to {@link #SOCKET_BUFFER_FILL_MARK}: the session thread + * serving this peer is then inside the write of a message larger than the buffers on both + * sides of the connection, and stays there until this peer reads. + */ + void awaitReceiveBufferFilled() throws Exception + { + newConnectionTimer().repeatUntilSuccess(new Callable() + { + @Override + public Void call() throws Exception + { + assertThat(socket.getInputStream().available()) + .as("the replication server never filled the receive buffer of the peer") + .isGreaterThanOrEqualTo(SOCKET_BUFFER_FILL_MARK); + return null; + } + }); } /** From 7c7f0558905d94f3b8e4273ea7440fe6001c6b7d Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Thu, 24 Sep 2026 17:50:46 +0300 Subject: [PATCH 2/2] [#1055] Pin the take-back of a message queued after the close drained the queue A publish() which read the close as not yet begun and was descheduled before its offer can queue a message after the drain; publish() then takes it back, reports it and answers false, which is what the writer of a replication server reads as a refusal. No case reached that arm, so answering true there survived. Session gets a hook, run between the check of the close and the offer and set only by the tests, which holds the publishing thread at that spot while the close runs to its end. SessionTest's class javadoc no longer says every case holds the session thread inside a write: only the first one publishes a message of that size. --- .../server/replication/protocol/Session.java | 30 +++++ .../protocol/SessionPublisherDrainTest.java | 117 +++++++++++++++++- .../replication/protocol/SessionTest.java | 7 +- 3 files changed, 150 insertions(+), 4 deletions(-) diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java index cc381edfec..f92154f569 100644 --- a/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java +++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/protocol/Session.java @@ -40,6 +40,7 @@ import org.forgerock.i18n.LocalizableMessage; import org.forgerock.i18n.slf4j.LocalizedLogger; +import org.forgerock.util.annotations.VisibleForTesting; import org.opends.server.api.DirectoryThread; import org.opends.server.types.HostPort; import org.opends.server.util.StaticUtils; @@ -127,6 +128,12 @@ private Outgoing(byte[] buffer, Runnable whenWritten) private final LinkedBlockingQueue sendQueue = new LinkedBlockingQueue<>(4000); private AtomicBoolean isRunning = new AtomicBoolean(false); + /** + * What {@link #publish(ReplicationMsg, Runnable)} runs between its check that no close has + * begun and the offer of the message to {@code sendQueue}, or null. Only the tests set it - see + * {@link #beforeQueueing(Runnable)}. + */ + private volatile Runnable beforeQueueing; private final CountDownLatch latch = new CountDownLatch(1); /** @@ -552,6 +559,11 @@ public boolean publish(final ReplicationMsg msg, final Runnable whenWritten) thr final Outgoing outgoing = new Outgoing(buffer, whenWritten); while (!closeInitiated) { + final Runnable hook = beforeQueueing; + if (hook != null) + { + hook.run(); + } try { // Avoid blocking forever so that we can check for session closure. @@ -623,6 +635,24 @@ private boolean takeBackWhatWasQueuedTooLate(final Outgoing outgoing) } } + /** + * Sets what {@link #publish(ReplicationMsg, Runnable)} runs between its check that no close has + * begun and the offer of the message to the queue. + *

+ * Only there for the tests of {@link #takeBackWhatWasQueuedTooLate(Outgoing)}: a + * {@code publish()} descheduled at that spot is the only one which can queue a message after a + * close has drained the queue, and nothing else holds a thread there on cue while the close + * runs to its end. + * + * @param hook + * What to run there, on the publishing thread, or null for nothing. + */ + @VisibleForTesting + void beforeQueueing(final Runnable hook) + { + beforeQueueing = hook; + } + /** Sends a replication message already encoded to the socket. * * @param buffer diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionPublisherDrainTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionPublisherDrainTest.java index d67b207b87..4111a8f9a2 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionPublisherDrainTest.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionPublisherDrainTest.java @@ -31,6 +31,7 @@ import java.util.TreeSet; import java.util.concurrent.Callable; import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; @@ -64,7 +65,8 @@ * close to sending that queue instead, and * {@code aCloseRunsTheCallbackOfEachMessageItSendsOutOfTheQueue} to telling the publisher of each * of those messages that it was written. The cases after it pin what a close gives up on: a - * queue whose write fails is reported with every message it held; a session which had already + * queue whose write fails is reported with every message it held; a message queued after the + * close drained the queue is taken back and reported, and refused; a session which had already * failed is written nothing more and still has its queue reported, together with what its * publisher took and failed to write, and nothing when there is nothing; and a closed session - * one closed before it was started included - fails a later {@code publish()} rather than take @@ -431,6 +433,119 @@ public void run() } } + /** + * A message queued after the close drained the queue is taken back and reported, rather than + * left in a queue nothing sends: {@code publish()} answers that it neither wrote nor queued it, + * and its callback never runs. Answered as queued, the message would leave the writer of a + * replication server waiting on a callback which never comes, for the rest of the grace period + * of the shutdown. + *

+ * Only a {@code publish()} which read the close as not yet begun and was descheduled before its + * offer gets there. The session holds the publishing thread at that spot through + * {@link Session#beforeQueueing(Runnable)} while the close runs to its end, and then lets the + * offer go. + */ + @Test + public void aMessageQueuedAfterTheCloseDrainedTheQueueIsTakenBackAndReported() throws Exception + { + final CSNGenerator csns = new CSNGenerator(RS_ID, 0); + final ExecutorService executor = Executors.newSingleThreadExecutor(); + try (ServerSocket listen = new ServerSocket(0)) + { + final Session[] pair = connectSessionPair(listen); + final Session sender = pair[0]; + final Session receiver = pair[1]; + try + { + sender.start(); + sender.waitForStartup(); + // So that the close ends the connection with a FIN rather than a reset, which would cut the + // StopMsg off - see aSessionWithAPublisherThreadSendsWhatIsStillQueuedWhenItIsClosed. + sender.setSoTimeout(0); + newInboundReader(sender).start(); + + final CountDownLatch atTheOffer = new CountDownLatch(1); + final CountDownLatch closed = new CountDownLatch(1); + sender.beforeQueueing(new Runnable() + { + @Override + public void run() + { + atTheOffer.countDown(); + try + { + closed.await(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS); + } + catch (final InterruptedException e) + { + Thread.currentThread().interrupt(); + } + } + }); + final AtomicInteger callbacks = new AtomicInteger(); + final Future published = executor.submit(new Callable() + { + @Override + public Boolean call() throws Exception + { + return sender.publish(new DeleteMsg( + DN.valueOf("uid=queuedtoolate," + TEST_ROOT_DN_STRING), + csns.newCSN(), "00000000-0000-0000-0000-000000000000"), + callbacks::incrementAndGet); + } + }); + assertThat(atTheOffer.await(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS)) + .as("the publish did not reach the offer of its message") + .isTrue(); + + final AtomicReference accepted = new AtomicReference<>(); + final List records = errorLogRecordsOf(new Callable() + { + @Override + public Void call() throws Exception + { + sender.close(); + closed.countDown(); + accepted.set(published.get(SOCKET_TIMEOUT_MS, TimeUnit.MILLISECONDS)); + return null; + } + }); + + assertThat(accepted.get()) + .as("a message queued after the close drained the queue was answered as queued, " + + "though nothing is left to send it") + .isFalse(); + assertThat(callbacks.get()) + .as("the callback of a message which was never written ran") + .isZero(); + assertThat(sendQueueOf(sender)) + .as("the message queued after the close was left in the queue") + .isEmpty(); + final Set reported = reportsIn(records); + assertThat(reported) + .as("a message taken back is reported once, and here the reports were: " + reported) + .hasSize(1); + assertThat(reported.iterator().next()) + .contains("1 message(s)") + .contains("it was queued after the publisher of the session had stopped"); + final Drained drained = drain(receiver); + assertThat(drained.received) + .as("the peer received the message taken back; the read ended by %s", drained.endedBy) + .isEmpty(); + assertThat(drained.endedBy).isEqualTo("a StopMsg"); + } + finally + { + sender.beforeQueueing(null); + StaticUtils.close(sender, receiver); + } + } + finally + { + executor.shutdownNow(); + } + } + /** * A close of a session which has already failed writes nothing more to it - neither the queue * nor the {@code StopMsg} - and still reports the queue it gives up on, once. diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java index 7525fdd7a3..f736e7b761 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java @@ -42,9 +42,10 @@ * thread which wrote it, only after the write returned. *

* The peer of each test reads nothing until the test lets it, and both ends of the connection - * have socket buffers far smaller than the first message published, so that the session thread - * of the end under test is held inside the write of that message for as long as the test wants - * - the state in which a message published behind it is queued and not written. + * have socket buffers far smaller than {@link #BLOCKING_MESSAGE_SIZE}, so that a test which + * publishes a message of that size holds the session thread of the end under test inside its + * write for as long as it wants - the state in which a message published behind it is queued + * and not written. */ @SuppressWarnings("javadoc") public class SessionTest extends ReplicationTestCase