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..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; @@ -112,8 +113,27 @@ 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); + /** + * 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); /** @@ -166,8 +186,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 +363,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 +374,7 @@ private void sendWhatThePublisherLeftQueued() } try { - send(buffer); + send(outgoing.buffer); } catch (final IOException e) { @@ -367,6 +389,7 @@ private void sendWhatThePublisherLeftQueued() "the write failed with " + e.getClass().getName() + ": " + e.getMessage()); return; } + written(outgoing.whenWritten); } } @@ -495,27 +518,58 @@ 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) { + final Runnable hook = beforeQueueing; + if (hook != null) + { + hook.run(); + } 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 +578,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 +613,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 { @@ -559,6 +635,24 @@ private void takeBackWhatWasQueuedTooLate(final byte[] buffer) } } + /** + * 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 @@ -769,10 +863,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 +874,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..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 @@ -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; @@ -29,10 +31,12 @@ 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; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicInteger; import java.util.concurrent.atomic.AtomicReference; import org.forgerock.opendj.ldap.DN; @@ -58,8 +62,11 @@ * 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 - * queue whose write fails is reported with every message it held; a session which had already + * 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 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 @@ -317,12 +324,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 +356,10 @@ 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 { @@ -355,6 +368,184 @@ public Void call() } } + /** + * 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 + { + StaticUtils.close(sender, receiver); + } + } + } + + /** + * 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. @@ -375,12 +566,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 +689,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 +885,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..f736e7b761 --- /dev/null +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/protocol/SessionTest.java @@ -0,0 +1,246 @@ +/* + * 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 {@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 +{ + 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; + } + }); } /**