diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java index d94f169e48..4bea16769d 100644 --- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java +++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java @@ -106,7 +106,8 @@ public class ReplicationServer private static final int LISTEN_PORT_PROBE_TIMEOUT_MS = 200; private volatile ServerSocket listenSocket; - private Thread listenThread; + /** Volatile like its socket above: a port change reads it from the configuration thread. */ + private volatile Thread listenThread; private Thread connectThread; /** The current configuration of this replication server. */ @@ -130,8 +131,6 @@ public class ReplicationServer private boolean externalChangelogRegistered; private final AtomicBoolean shutdown = new AtomicBoolean(); - /** Written by the thread applying a configuration change, read by the listen thread. */ - private volatile boolean stopListen; private final ReplSessionSecurity replSessionSecurity; private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass(); @@ -155,6 +154,15 @@ public class ReplicationServer */ static final AtomicInteger listenPortBindFailures = new AtomicInteger(); + /** + * Number of listen port changes whose wait for the previous listen thread was interrupted. + *
+ * This is required for unit testing: a wait which is interrupted and a wait which is over
+ * before it starts leave this replication server in the same state, so nothing else tells
+ * a test that it exercised the interruption instead of passing over it.
+ */
+ static final AtomicInteger interruptedListenThreadStops = new AtomicInteger();
+
/** Monitors for synchronizing domain creation with the connect thread. */
private final Object domainTicketLock = new Object();
private final Object connectThreadLock = new Object();
@@ -260,15 +268,22 @@ public static List
+ * The socket is the one this thread was created for, not the one this replication server
+ * currently listens on: closing it is what stops this thread, and it is then the only
+ * thread it stops, even when another one is already listening on another port.
+ *
+ * @param socket
+ * the bound socket this thread accepts connections on
*/
- void runListen()
+ void runListen(ServerSocket socket)
{
logger.info(NOTE_REPLICATION_SERVER_LISTENING,
getServerId(),
- listenSocket.getInetAddress().getHostAddress(),
- listenSocket.getLocalPort());
+ socket.getInetAddress().getHostAddress(),
+ socket.getLocalPort());
- while (!shutdown.get() && !stopListen)
+ while (!shutdown.get() && !socket.isClosed())
{
// Wait on the replicationServer port.
// Read incoming messages and create LDAP or ReplicationServer listener
@@ -279,7 +294,7 @@ void runListen()
Socket newSocket = null;
try
{
- newSocket = listenSocket.accept();
+ newSocket = socket.accept();
newSocket.setTcpNoDelay(true);
newSocket.setKeepAlive(true);
int timeoutMS = MultimasterReplication.getConnectionTimeoutMS();
@@ -486,8 +501,8 @@ private boolean connect(HostPort remoteServerAddress, DN baseDN)
* Initialization function for the replicationServer.
*
* @throws ConfigException
- * when the replication server cannot be started, in particular when its listen
- * port cannot be bound.
+ * when the replication server cannot be started, in particular when its changelog
+ * cannot be read or when its listen port cannot be bound.
*/
private void initialize() throws ConfigException
{
@@ -495,9 +510,13 @@ private void initialize() throws ConfigException
try
{
+ // Assigned before the changelog is opened: the monitor instance name of the domains it
+ // restores, and of their changelogs, embeds it, and a provider registered under a name
+ // which later changes can never be deregistered again.
+ setServerURL();
+
this.changelogDB.initializeDB();
- setServerURL();
// Assigned before the threads are created, so that a failure below still releases it.
listenSocket = bindListenPort(getReplicationPort());
@@ -520,6 +539,12 @@ private void initialize() throws ConfigException
{
logger.trace("RS " + getMonitorInstanceName() + " successfully initialized");
}
+ } catch (ChangelogException e)
+ {
+ // A replication server which cannot read its changelog is as dead as one which cannot
+ // bind its listen port (issue #802). The message already names the changelog directory.
+ logger.traceException(e);
+ throw new ConfigException(e.getMessageObject(), e);
} catch (UnknownHostException e)
{
// Not logged here: the caller reports the ConfigException, logging it once.
@@ -527,8 +552,8 @@ private void initialize() throws ConfigException
throw new ConfigException(ERR_UNKNOWN_HOSTNAME.get(), e);
} catch (IOException e)
{
- // A replication server whose listen port is not bound is dead: every consumer would
- // otherwise only learn about it as a "connection refused" somewhere else.
+ // Every consumer would otherwise only learn about it as a "connection refused"
+ // somewhere else (issue #792).
logger.traceException(e);
throw new ConfigException(bindFailureMessage(getReplicationPort(), e), e);
}
@@ -582,28 +607,29 @@ private ServerSocket bindListenPort(int port) throws IOException
}
/**
- * Stops the listen thread and releases the listen port.
+ * Stops the listen thread of the provided listen socket and releases that port.
+ *
+ * Both are the ones the thread was started on rather than the current ones, so this stops
+ * that thread only, even when another one is already listening on another port.
*
+ * @param socket
+ * the listen socket to close, which is what stops its thread
+ * @param thread
+ * the listen thread of that socket, {@code null} when it was never started
* @throws InterruptedException
* if this thread is interrupted while waiting for the listen thread to stop
*/
- private void stopListenThread() throws InterruptedException
+ private void stopListenThread(ServerSocket socket, Thread thread) throws InterruptedException
{
- stopListen = true;
- close(listenSocket);
- if (listenThread != null)
+ close(socket);
+ if (thread != null)
{
- listenThread.join();
- listenThread = null;
+ thread.join();
}
}
/**
* Starts a listen thread on the provided listen socket.
- *
- * {@code stopListen} is only cleared here, i.e. once the listen port is bound, so a
- * failure to bind leaves this replication server consistently stopped rather than with a
- * listen thread which would spin on a closed socket.
*
* @param boundListenSocket
* the bound socket the listen thread will accept connections on
@@ -611,30 +637,40 @@ private void stopListenThread() throws InterruptedException
private void startListenThread(ServerSocket boundListenSocket)
{
listenSocket = boundListenSocket;
- stopListen = false;
- listenThread = new ReplicationServerListenThread(this);
+ listenThread = new ReplicationServerListenThread(this, boundListenSocket);
listenThread.start();
}
/**
* Switches the listen port to the one of the provided configuration.
*
- * The new port is bound while the current one is still open and serving, so a failure
- * leaves this replication server listening on its current port, with its current
- * configuration: there is nothing to roll back, and no window during which this
- * replication server advertises a port that nothing listens to.
+ * The new port is bound, and its listen thread started, while the current one is still
+ * open and serving. A failure therefore leaves this replication server listening on its
+ * current port, with its current configuration: there is nothing to roll back, and there
+ * is no window during which this replication server listens on no port at all.
+ *
+ * The trade is a window during which both ports accept, so a peer which connects to the
+ * previous port just before it is released gets a session which outlives the change. That
+ * is the deliberate inverse of a window during which nothing listens at all.
*
* @param newConfig
* the configuration being applied, whose listen port differs from the current one
* @param ccr
* the result of the configuration change, to which a failure is added
+ * @param listenThreadStopInterrupted
+ * set when the wait for the previous listen thread was interrupted, instead of
+ * restoring the interrupt status here: the caller restores it once the rest of
+ * the change, some of which is interruptible, has run
* @return {@code true} when this replication server listens on the new port, in which
* case {@code newConfig} has become its configuration
*/
- private boolean switchListenPort(ReplicationServerCfg newConfig, ConfigChangeResult ccr)
+ private boolean switchListenPort(ReplicationServerCfg newConfig, ConfigChangeResult ccr,
+ AtomicBoolean listenThreadStopInterrupted)
{
final ReplicationServerCfg previousConfig = this.config;
final String previousServerURL = serverURL;
+ final ServerSocket previousListenSocket = listenSocket;
+ final Thread previousListenThread = listenThread;
final int newPort = newConfig.getReplicationPort();
ServerSocket newListenSocket = null;
try
@@ -645,12 +681,28 @@ private boolean switchListenPort(ReplicationServerCfg newConfig, ConfigChangeRes
this.config = newConfig;
setServerURL();
- stopListenThread();
+ // The new port is served before the current one is released, so that this replication
+ // server is never left with no listener at all, whatever happens next.
startListenThread(newListenSocket);
newListenSocket = null;
+ // In step with getReplicationPort(), which answers the new port from here on: the
+ // wait below blocks for as long as the previous thread takes to serve its current
+ // connection, and localPorts must not trail it for that whole window.
localPorts.remove(previousConfig.getReplicationPort());
localPorts.add(newPort);
+ try
+ {
+ stopListenThread(previousListenSocket, previousListenThread);
+ }
+ catch (InterruptedException e)
+ {
+ // The previous port is already released and its thread stops on its own as soon as
+ // it wakes up on its closed socket: only the wait for it was cut short.
+ interruptedListenThreadStops.incrementAndGet();
+ listenThreadStopInterrupted.set(true);
+ logger.traceException(e);
+ }
return true;
}
catch (UnknownHostException e)
@@ -666,15 +718,6 @@ private boolean switchListenPort(ReplicationServerCfg newConfig, ConfigChangeRes
ccr.setResultCode(ResultCode.OPERATIONS_ERROR);
ccr.addMessage(bindFailureMessage(newPort, e));
}
- catch (InterruptedException e)
- {
- // The previous listen thread may still be running, so do not hand it a new socket:
- // stopListen is still set, which makes that thread stop as soon as it wakes up.
- Thread.currentThread().interrupt();
- logger.traceException(e);
- ccr.setResultCode(ResultCode.OPERATIONS_ERROR);
- ccr.addMessage(ERR_COULD_NOT_STOP_LISTEN_THREAD.get(getExceptionMessage(e)));
- }
// The failure is reported through the ConfigChangeResult, which the configuration
// handler logs: nothing of the new configuration was applied.
this.config = previousConfig;
@@ -754,6 +797,22 @@ private void abortInitialization()
{
listenThread.interrupt();
}
+ // Opening the changelog restores one domain per domain it holds, and each of them starts
+ // its threads and registers its monitor provider: a failure after that point, such as a
+ // listen port which cannot be bound, would otherwise leave them behind. Shut them down
+ // before the changelog they write to, and one unchecked exception at a time: the changelog
+ // this one is built on is known to be broken, and what follows still has to run.
+ for (ReplicationServerDomain domain : getReplicationServerDomains())
+ {
+ try
+ {
+ domain.shutdown();
+ }
+ catch (RuntimeException ignored)
+ {
+ logger.traceException(ignored);
+ }
+ }
shutdownExternalChangelog();
if (this.changelogDB != null)
{
@@ -1195,8 +1254,9 @@ public ConfigChangeResult applyConfigurationChange(
// done first, and the new port is bound before the current one is released, so that a
// change which cannot be applied leaves this replication server as it was, instead of
// half configured and, worse, without any listener.
+ final AtomicBoolean listenThreadStopInterrupted = new AtomicBoolean();
if (configuration.getReplicationPort() != oldConfig.getReplicationPort()
- && !switchListenPort(configuration, ccr))
+ && !switchListenPort(configuration, ccr, listenThreadStopInterrupted))
{
return ccr;
}
@@ -1262,6 +1322,15 @@ public ConfigChangeResult applyConfigurationChange(
{
ccr.setAdminActionRequired(true);
}
+
+ // The interrupt which cut short the wait for the previous listen thread, deferred by
+ // switchListenPort(): restored only now, because the steps above include interruptible
+ // ones — stopping the handlers of removed replication servers locks interruptibly —
+ // which an interrupt status left set would have failed while the change reports SUCCESS.
+ if (listenThreadStopInterrupted.get())
+ {
+ Thread.currentThread().interrupt();
+ }
return ccr;
}
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerListenThread.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerListenThread.java
index f2a554be5e..0f91647dc7 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerListenThread.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServerListenThread.java
@@ -13,9 +13,12 @@
*
* Copyright 2008 Sun Microsystems, Inc.
* Portions Copyright 2011-2015 ForgeRock AS.
+ * Portions Copyright 2026 3A Systems, LLC.
*/
package org.opends.server.replication.server;
+import java.net.ServerSocket;
+
import org.opends.server.api.DirectoryThread;
/**
@@ -30,25 +33,30 @@ public class ReplicationServerListenThread extends DirectoryThread
*/
private final ReplicationServer server;
+ /** The socket this thread accepts connections on, and whose closing stops it. */
+ private final ServerSocket listenSocket;
+
/**
* Creates a new instance of this directory thread with the
* specified name.
*
* @param server The ReplicationServer that will be called to
* handle the connections.
+ * @param listenSocket The bound socket this thread will accept connections on.
*/
- public ReplicationServerListenThread(ReplicationServer server)
+ public ReplicationServerListenThread(ReplicationServer server, ServerSocket listenSocket)
{
super("Replication server RS(" + server.getServerId()
+ ") connection listener on port "
- + server.getReplicationPort());
+ + listenSocket.getLocalPort());
this.server = server;
+ this.listenSocket = listenSocket;
}
/** {@inheritDoc} */
@Override
public void run()
{
- server.runListen();
+ server.runListen(listenSocket);
}
}
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/api/ChangelogDB.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/api/ChangelogDB.java
index 9b9e3efe39..a75f7dad18 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/api/ChangelogDB.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/api/ChangelogDB.java
@@ -12,6 +12,7 @@
* information: "Portions Copyright [year] [name of copyright owner]".
*
* Copyright 2013 ForgeRock AS.
+ * Portions Copyright 2026 3A Systems, LLC.
*/
package org.opends.server.replication.server.changelog.api;
@@ -29,8 +30,13 @@ public interface ChangelogDB
* Initializes the replication database by reading its previous state and
* building the relevant ReplicaDBs according to the previous state. This
* method must be called once before using the ChangelogDB.
+ *
+ * @throws ChangelogException
+ * If the previous state could not be read. The database is then
+ * unusable, possibly half open, and the caller must release it by
+ * calling {@link #shutdownDB()}.
*/
- void initializeDB();
+ void initializeDB() throws ChangelogException;
/**
* Sets the purge delay for the replication database. Can be called while the
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/FileChangelogDB.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/FileChangelogDB.java
index 2164ae965f..768d127e19 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/FileChangelogDB.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/FileChangelogDB.java
@@ -279,7 +279,7 @@ private Pair
+ * It is also the shape which restores domains before it fails, so it is the one where the
+ * aborted initialization has something to release.
+ */
+ @Test
+ public void replServerFailsWhenAReplicaChangelogCannotBeRead() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ final int rsServerId = 8021;
+ final String dbDirName = "replServerFailsWhenAReplicaChangelogCannotBeReadDb";
+ final File dbDirectory = createPopulatedChangelog(dbDirName);
+ try
+ {
+ // The head log file of the replica is replaced by a directory: the changelog state is
+ // then still readable, and the changes of the domain it names are not.
+ final File headLogFile = findReplicaLogFile(dbDirectory);
+ assertNotNull(headLogFile, "no replica changelog was written under " + dbDirectory);
+ assertTrue(headLogFile.delete() && headLogFile.mkdir(), "could not replace " + headLogFile);
+
+ final int[] ports = TestCaseUtils.findFreePorts(1);
+ final int instancesBefore = ReplicationServer.getAllInstances().size();
+ try
+ {
+ final ReplicationServer replicationServer = new ReplicationServer(
+ new ReplServerFakeConfiguration(ports[0], dbDirName, 0, rsServerId, 0, 0, null));
+ remove(replicationServer);
+ fail("Creating a replication server over an unreadable replica changelog should have failed");
+ }
+ catch (ConfigException expected)
+ {
+ assertTrue(expected.getCause() instanceof ChangelogException,
+ "the failure should be the one of the changelog, but was: " + expected.getCause());
+ assertTrue(expected.getMessage().contains(dbDirectory.getAbsolutePath()),
+ "the failure should name the changelog directory, but was: " + expected.getMessage());
+ // The log of the replica, named after its directory, and not the one of the change
+ // number index: otherwise this test exercises another failure shape than its name.
+ assertTrue(expected.getMessage().contains(headLogFile.getParentFile().getPath()),
+ "the failure should be the one of the replica changelog which was corrupted,"
+ + " but was: " + expected.getMessage());
+ assertEquals(ReplicationServer.getAllInstances().size(), instancesBefore,
+ "the failed replication server must not be left registered");
+ assertNothingLeftBehind(rsServerId);
+ }
+ }
+ finally
+ {
+ recursiveDelete(dbDirectory);
+ }
+ }
+
+ /**
+ * Tests that a replication server which cannot bind its listen port over a changelog it
+ * did read releases the domains that reading restored: each of them holds a timer thread
+ * and registers monitor providers, and the aborted instance is never handed to anything
+ * which could shut them down later.
+ */
+ @Test
+ public void abortedStartReleasesTheRestoredDomains() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ final int rsServerId = 8022;
+ final String dbDirName = "abortedStartReleasesTheRestoredDomainsDb";
+ final File dbDirectory = createPopulatedChangelog(dbDirName);
+ try (ServerSocket portHolder = TestCaseUtils.bindFreePort())
+ {
+ try
+ {
+ final ReplicationServer replicationServer = new ReplicationServer(
+ new ReplServerFakeConfiguration(portHolder.getLocalPort(), dbDirName, 0, rsServerId, 0, 0, null));
+ remove(replicationServer);
+ fail("Creating a replication server on a port already in use should have failed");
+ }
+ catch (ConfigException expected)
+ {
+ assertNothingLeftBehind(rsServerId);
+ }
+ }
+ finally
+ {
+ recursiveDelete(dbDirectory);
+ }
+ }
+
+ /**
+ * Tests that a replication server which restarted over an existing changelog releases the
+ * domains that reading it restored when it stops: their monitor instance name embeds the
+ * URL of their replication server, which used to be assigned only after the changelog had
+ * been read, so they were registered under a name holding a null URL and the name looked
+ * up to deregister them, built from the assigned URL, could never match it again.
+ */
+ @Test
+ public void restartedReplServerReleasesTheRestoredDomains() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ final int rsServerId = 8023;
+ final String dbDirName = "restartedReplServerReleasesTheRestoredDomainsDb";
+ final File dbDirectory = createPopulatedChangelog(dbDirName);
+ ReplicationServer replicationServer = null;
+ try
+ {
+ final int[] ports = TestCaseUtils.findFreePorts(1);
+ replicationServer = new ReplicationServer(
+ new ReplServerFakeConfiguration(ports[0], dbDirName, 0, rsServerId, 0, 0, null));
+ assertTrue(replicationServer.isListening());
+ assertFalse(domainRegistrationsOf(rsServerId).isEmpty(),
+ "the restored domain should hold a timer thread and monitor providers, otherwise"
+ + " this test does not test that they are released");
+ }
+ finally
+ {
+ remove(replicationServer);
+ recursiveDelete(dbDirectory);
+ }
+ assertNothingLeftBehind(rsServerId);
+ }
+
+ /**
+ * Tests that a port change whose wait for the previous listen thread is interrupted still
+ * leaves this replication server listening: the new port is served before the previous one
+ * is released, so an interrupt can only cut short the wait for a thread which is already
+ * stopping, never leave the replication server with no listener at all.
+ *
+ * A connection which is accepted and then says nothing keeps the previous listen thread in
+ * its handshake instead of at {@code accept()}, so closing its socket does not stop it
+ * before the wait even begins: {@code Thread.join()} only throws while the thread it waits
+ * for is alive, so without that connection this test would pass over the interruption
+ * instead of exercising it. {@code interruptedListenThreadStops} tells the two apart.
+ *
+ * That the connection is accepted is waited for on both sides of it, see
+ * {@link #waitForListenThread(int, int, boolean)}: connecting only proves the port is bound.
+ */
+ @Test
+ public void replServerKeepsListeningWhenAPortChangeIsInterrupted() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ ReplicationServer replicationServer = null;
+ try
+ {
+ final int[] ports = TestCaseUtils.findFreePorts(2);
+ final String dbDirName = "replServerKeepsListeningWhenAPortChangeIsInterruptedDb";
+ replicationServer = new ReplicationServer(
+ new ReplServerFakeConfiguration(ports[0], dbDirName, 0, 1, 0, 0, null));
+ assertTrue(replicationServer.isListening());
+
+ final int interruptsBefore = ReplicationServer.interruptedListenThreadStops.get();
+ final ConfigChangeResult ccr;
+ try (Socket silent = new Socket())
+ {
+ // The listen thread has to be inside accept() before the connection is made, and to
+ // have left it afterwards: a connection completes against the listen backlog of the
+ // kernel, so it does not prove that the thread which serves it ever ran.
+ final int serverId = replicationServer.getServerId();
+ waitForListenThread(serverId, ports[0], true);
+ silent.connect(new InetSocketAddress("127.0.0.1", ports[0]), 5000);
+ waitForListenThread(serverId, ports[0], false);
+
+ // Thread.join() throws InterruptedException at once when the interrupt status is
+ // already set, i.e. this interrupts the port change in its wait for the listen thread.
+ // What runs until that wait — binding the new port, resolving the server URL, starting
+ // the new listen thread — has to fit in the handshake timeout of the silent connection,
+ // MultimasterReplication.getConnectionTimeoutMS(), 5s by default: that handshake is
+ // what keeps the previous listen thread alive, hence what makes the wait for it block.
+ Thread.currentThread().interrupt();
+ ccr = replicationServer.applyConfigurationChange(
+ new ReplServerFakeConfiguration(ports[1], dbDirName, 0, 1, 0, 0, null));
+ }
+ // Cleared for the rest of this test, and for whatever runs next in this thread.
+ final boolean interrupted = Thread.interrupted();
+
+ assertEquals(ReplicationServer.interruptedListenThreadStops.get(), interruptsBefore + 1,
+ "the port change should have been interrupted in its wait for the previous listen"
+ + " thread, otherwise this test does not test that interruption");
+ assertTrue(interrupted, "the interrupted port change should have restored the interrupt status");
+ assertEquals(ccr.getResultCode(), ResultCode.SUCCESS);
+ assertEquals(replicationServer.getReplicationPort(), ports[1],
+ "the replication server should have switched to the new listen port");
+ assertTrue(replicationServer.isListening(),
+ "an interrupted port change must not leave the replication server without a listener");
+
+ // and it must be usable on the new port.
+ ReplicationBroker broker = openReplicationSession(
+ DN.valueOf(TEST_ROOT_DN_STRING), 1, 10, ports[1], 1000);
+ assertTrue(broker.getCurrentSendWindow() != 0);
+ }
+ finally
+ {
+ remove(replicationServer);
+ }
+ }
+
+ /**
+ * Waits for the listen thread of the provided replication server and port to be inside
+ * {@code accept()}, or to have left it.
+ *
+ * A connection completes against the listen backlog of the kernel, and
+ * {@code isListening()} only tells that the socket is bound — the listen thread is
+ * started after it — so neither of them proves that thread ever ran. Waiting for it to be
+ * inside {@code accept()} before connecting is what makes the second wait conclusive: the
+ * thread can then only have left {@code accept()} for the connection this test made. That
+ * rests on it being the only connection the port ever gets — a stale broker of an earlier
+ * test reconnecting to a recycled port would satisfy the second wait spuriously.
+ *
+ * Having left {@code accept()} does not mean the connection is served yet — the thread is
+ * typically still warming up towards its handshake. What the second wait establishes is
+ * that the thread is off {@code accept()} and cannot terminate until its socket is
+ * closed, which is what stopping it does.
+ *
+ * @param serverId
+ * the server id of the replication server whose listen thread is waited for
+ * @param port
+ * the port that listen thread listens on
+ * @param accepting
+ * {@code true} to wait for that thread to be inside {@code accept()},
+ * {@code false} to wait for it to have left it
+ */
+ private void waitForListenThread(int serverId, int port, boolean accepting) throws Exception
+ {
+ final String listenThread =
+ "replication server rs(" + serverId + ") connection listener on port " + port;
+ final long deadline = System.currentTimeMillis() + 60000;
+ while (System.currentTimeMillis() < deadline)
+ {
+ for (Map.Entry
+ * It is not the one of the replication servers under test, so that what it leaves behind —
+ * it is the only one to which a replica ever connects, and the monitor providers of a
+ * connection are deregistered when its handler notices it is gone — cannot be mistaken for
+ * what they leave behind.
+ */
+ private static final int CHANGELOG_WRITER_RS_ID = 8020;
+
+ /** The suffix of the directory a changelog holds per domain, see {@code ReplicationEnvironment}. */
+ private static final String DOMAIN_DIRECTORY_SUFFIX = ".dom";
+
+ /**
+ * Creates the changelog of a replication server which ran and served one replica, and
+ * returns its directory: a changelog whose reading restores a domain, i.e. one over which
+ * an initialization has something to release when it fails.
+ */
+ private File createPopulatedChangelog(String dbDirName) throws Exception
+ {
+ final File dbDirectory = getFileForPath(dbDirName);
+ recursiveDelete(dbDirectory);
+
+ final int[] ports = TestCaseUtils.findFreePorts(1);
+ final ReplicationServer replicationServer = new ReplicationServer(
+ new ReplServerFakeConfiguration(ports[0], dbDirName, 0, CHANGELOG_WRITER_RS_ID, 0, 0, null));
+ ReplicationBroker broker = null;
+ try
+ {
+ final DN baseDN = DN.valueOf(TEST_ROOT_DN_STRING);
+ broker = openReplicationSession(baseDN, 42, 100, ports[0], 1000);
+ broker.publish(new DeleteMsg(baseDN, new CSNGenerator(42, 0).newCSN(), "uid"));
+
+ // The changelog is written by the replication server, so the calls above are only over
+ // once the log of the replica exists. It is the last of the four files which
+ // ReplicationEnvironment.getOrCreateReplicaDB() writes, after domains.state, the server
+ // id directory and the generation id, so waiting for it waits for all of them.
+ waitForReplicaLogFile(dbDirectory);
+ }
+ finally
+ {
+ stop(broker);
+ replicationServer.shutdown();
+ }
+ return dbDirectory;
+ }
+
+ /**
+ * Asserts that the replication server with the provided server id left neither a thread of
+ * a domain nor a monitor provider of a domain or of its changelog behind.
+ */
+ private void assertNothingLeftBehind(int rsServerId) throws Exception
+ {
+ // Stopping a thread only asks it to stop, so give the ones being stopped the time to
+ // actually stop before reporting them as left behind.
+ final long deadline = System.currentTimeMillis() + 10000;
+ List
+ * The wait is long because it is only ever reached on a machine which is slow enough for
+ * the replication server to still be writing that file: a test which passes never waits.
+ */
+ private void waitForReplicaLogFile(File dbDirectory) throws Exception
+ {
+ final long deadline = System.currentTimeMillis() + 60000;
+ while (findReplicaLogFile(dbDirectory) == null && System.currentTimeMillis() < deadline)
+ {
+ Thread.sleep(10);
+ }
+ assertNotNull(findReplicaLogFile(dbDirectory),
+ "no replica changelog, i.e. no head log file under a '" + DOMAIN_DIRECTORY_SUFFIX
+ + "' directory, was written under " + dbDirectory
+ + ": check that suffix against ReplicationEnvironment, which owns it");
+ }
+
+ /**
+ * Returns the head log file of a replica changelog, i.e. the one under a domain directory.
+ *
+ * The changelog of the change number index holds a head log file of its own, and it is
+ * created when the replication server starts: a lookup which is not scoped to a domain
+ * directory can match it instead, and then waits for nothing and corrupts the wrong log.
+ */
+ private File findReplicaLogFile(File dbDirectory)
+ {
+ final File[] entries = dbDirectory.listFiles();
+ if (entries == null)
+ {
+ return null;
+ }
+ for (File entry : entries)
+ {
+ if (entry.isDirectory() && entry.getName().endsWith(DOMAIN_DIRECTORY_SUFFIX))
+ {
+ final File replicaLogFile = findFile(entry, "head", ".log");
+ if (replicaLogFile != null)
+ {
+ return replicaLogFile;
+ }
+ }
+ }
+ return null;
+ }
+
+ /** Returns the first file whose name matches, at any depth of the provided directory. */
+ private File findFile(File directory, String prefix, String suffix)
+ {
+ final File[] files = directory.listFiles();
+ if (files == null)
+ {
+ return null;
+ }
+ for (File file : files)
+ {
+ final String name = file.getName();
+ if (name.startsWith(prefix) && name.endsWith(suffix))
+ {
+ return file;
+ }
+ final File found = file.isDirectory() ? findFile(file, prefix, suffix) : null;
+ if (found != null)
+ {
+ return found;
+ }
+ }
+ return null;
+ }
+
/** Returns the names of the virtual attributes provided by the external changelog. */
private List