/* * 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 static org.opends.server.TestCaseUtils.TEST_ROOT_DN_STRING; import java.io.Closeable; import java.lang.reflect.Field; import java.net.ServerSocket; import java.net.Socket; import java.util.LinkedHashSet; import java.util.List; import java.util.Queue; import java.util.Set; import java.util.TreeSet; import java.util.concurrent.Callable; import java.util.concurrent.CopyOnWriteArrayList; 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.AtomicReference; import org.forgerock.opendj.ldap.DN; 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.replication.server.ReplServerFakeConfiguration; import org.opends.server.replication.server.ReplicationServer; import org.opends.server.replication.service.ReplicationBroker; import org.opends.server.util.StaticUtils; import org.testng.annotations.Test; /** * Which end of a replication session has a publisher thread, and what a close of the session * does to the messages that thread has not sent yet. *
* {@link Session#publish(ReplicationMsg)} has two branches, and which one a message takes used * to decide whether a close could lose it. With a publisher thread running the call is an enqueue * onto {@code sendQueue}; without one it is a synchronous write of the socket. {@link * Session#close()} used to drain neither: it set the flag the publisher loops on, interrupted it * and joined, so anything still queued was dropped while the {@code StopMsg} published afterwards * 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 * 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 * it into a queue nothing sends. *
* The first two cases pin which end could ever pay it. Only {@code ServerHandler} starts a session's
* publisher, so it is the replication-server end of a session which has one; the broker of a
* directory server never starts its own. A change a directory server publishes is therefore on
* the wire by the time {@code publish()} returns, and no close of that session could drop it -
* which is what rules the send queue out as the explanation of #963.
*
* @see issue #963
*/
@SuppressWarnings("javadoc")
public class SessionPublisherDrainTest extends ReplicationTestCase
{
private static final int DS_ID = 123;
private static final int RS_ID = 104;
private static final int SOCKET_TIMEOUT_MS = 5000;
/**
* The number of messages the queued-message test publishes. It has to be enough for
* {@code publish()} to outrun the publisher thread and leave a backlog behind it - the socket
* buffers of a loopback pair swallow all of it, so it is not a full buffer which leaves one -
* and stay under the 4000 the send queue holds, past which {@code publish()} would block
* instead of queueing.
*/
private static final int MESSAGES_PUBLISHED = 3000;
/**
* The number of messages the case of a queue which cannot be sent leaves on the sender. It is
* small and exact: what that case is about is the count the close reports, not a backlog.
*/
private static final int MESSAGES_LEFT_UNSENT = 7;
/** What the line a close writes about a queue it could not send is recognised by. */
private static final String NOT_SENT_REPORT = "was closed with";
/**
* The session a directory server publishes its changes on has no publisher thread: nothing
* calls {@link Session#start()} on it, so the thread is still {@code NEW} once the broker is
* connected and has completed its handshake.
*/
@Test
public void theSessionOfADirectoryServerBrokerHasNoPublisherThread() throws Exception
{
final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
ReplicationServer replicationServer = null;
ReplicationBroker broker = null;
try
{
final int replicationPort = TestCaseUtils.findFreePort();
replicationServer = new ReplicationServer(new ReplServerFakeConfiguration(
replicationPort, "sessionPublisherDrainDb", 0, RS_ID, 0, 100, new TreeSet
* The peer reads the changes and then the {@code StopMsg}, which is what {@link #drain(Session)}
* stops on: a queue sent after that message would not be counted, so the size below pins the
* order as well as the delivery.
*/
@Test
public void aSessionWithAPublisherThreadSendsWhatIsStillQueuedWhenItIsClosed() 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();
/*
* A reader on the *sending* end, which is what the server keeps running across a close:
* ServerHandler.shutdown() closes the session (ServerHandler.java:946) and only joins its
* ServerReader afterwards (:966). It is not decoration here.
*
* Without it this end reaches close() with inbound bytes nobody ever read, and a close in
* that state ends the connection with a reset instead of a FIN - which discards what the
* peer has not read yet, the whole of what the drain just wrote included. Measured with a
* bare socket pair, 8 MiB written to a peer reading behind the writer, the only difference
* between the runs being whether the closing side drained its own inbound:
*
* linux 6.12 / jdk 11 no reader -> 53.4% arrived, "Connection reset"
* reader -> 100% arrived, end of stream
* macos 15.7 / jdk 26 no reader -> 98.3% arrived, "Connection reset"
* reader -> 100% arrived, end of stream
*
* Which is why a case without it measures the teardown rather than the drain, and measures
* it differently per platform: three ubuntu legs of CI lost between 866 and 1694 of these
* messages while every macos and windows leg passed. This case says so itself - with the
* start() below commented out and the class run on linux/jdk11, it fails with "the peer
* received 2873 of the 3000 messages published; the read ended by java.net.SocketException:
* Connection reset", and passes with it in.
*
* Its soTimeout has to go, too: receive() hands a read timeout to setSessionError(), and a
* session carrying an error skips the drain exactly as it skips the StopMsg.
*/
sender.setSoTimeout(0);
final Thread senderReader = newInboundReader(sender);
senderReader.start();
/*
* Nothing reads the peer end while these are published, but that is not what leaves the
* backlog: 3000 frames of ~150 B are ~440 KB, which the socket buffers of a loopback pair
* swallow, so the publisher thread is not blocked inside a write. What leaves the backlog
* is publish() - encode and offer - outrunning the publisher, which writes and flushes one
* frame at a time through the TLS layer. That is a race rather than a state, so the case
* asserts below that it was still won when the close ran: with an empty queue there is
* nothing for the close to drain and this pins nothing.
*/
for (int i = 0; i < MESSAGES_PUBLISHED; i++)
{
sender.publish(new DeleteMsg(DN.valueOf("uid=queued" + i + "," + TEST_ROOT_DN_STRING),
csns.newCSN(), "00000000-0000-0000-0000-000000000000"));
}
final int queuedAtClose = sendQueueOf(sender).size();
assertThat(queuedAtClose)
.as("the publisher had sent everything before the close ran, so this case drained "
+ "nothing: what it measures is the socket rather than close()")
.isGreaterThan(0);
final Future> closed = executor.submit(new Callable
* The road is a peer which is gone by the time the close reaches the queue - a directory server
* whose {@code StopMsg} brought the {@code ServerReader} of its handler to the {@code close()}
* of its finally, with the publisher of that session still holding a backlog. Here the sockets
* of the sender are closed under it instead, which is the same failed write with an exact
* count: the first {@code send()} of the drain throws, so the report has to name every message
* the queue held. A peer closed from the outside gives up somewhere inside the TCP buffers
* instead, and would pin no number at all.
*
* The queue is filled through the field rather than by publishing on a started session for the
* same reason: what a publisher thread has left behind is a race, and this case is the count.
*/
@Test
public void aCloseWhichCannotSendTheQueueReportsEveryMessageTheQueueHeld() 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 Queue
* The error is set through the field, with the sockets left intact: a close which wrote the
* queue regardless would then get it through, and the peer reading a change is what says so.
* With the sockets closed under it, as in the case above, such a close would fail at its first
* write and look the same from the peer.
*/
@Test
public void aCloseOfAFailedSessionWritesNothingAndReportsTheQueue() 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 Queue
* The count is exact whatever the thread got through before the close: every message published
* is either still in the queue or was taken and failed. The case waits for the queue to empty so
* that it is the publisher's own count the report stands on - the queue alone would report
* nothing.
*/
@Test
public void aCloseOfAStartedSessionWhoseWritesFailedReportsWhatThePublisherTookAsWell()
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
{
sender.start();
sender.waitForStartup();
closeTheSocketsUnder(sender);
for (int i = 0; i < MESSAGES_LEFT_UNSENT; i++)
{
sender.publish(new DeleteMsg(DN.valueOf("uid=takenandlost" + i + ","
+ TEST_ROOT_DN_STRING), csns.newCSN(), "00000000-0000-0000-0000-000000000000"));
}
final Queue
* It is what it reads at the socket rather than what it returns that matters: the peer of these
* cases publishes nothing, so this returns no message at all, while the read it sits in is what
* keeps the receive queue of this end empty - which is the condition a close needs to end the
* connection in an orderly way.
*/
private Thread newInboundReader(final Session session)
{
final Thread reader = new Thread(new Runnable()
{
@Override
public void run()
{
try
{
while (true)
{
session.receive();
}
}
catch (final Exception ignored)
{
// The close of the session ends the read, which is the end of this thread.
}
}
}, "inbound reader of " + session.getName());
reader.setDaemon(true);
return reader;
}
/**
* Reads until the session gives nothing back, and answers the CSNs of the changes it read
* together with what ended the read.
*
* Why the reason is carried rather than swallowed: a short read is exactly the failure this
* suite is about, and "the peer received fewer than were sent" does not say whether the stream
* ended orderly at a {@code StopMsg} or was cut off - which are different defects.
*/
private static final class Drained
{
private final List