From 44ebda86d53d2c1da92ae198741c0732a72169f2 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Wed, 05 Aug 2026 10:10:43 +0000
Subject: [PATCH] [#813] Refuse to create a replica DB once the changelog shutdown has started (#820)
---
opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/FileChangelogDB.java | 55 ++++
opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogDBTest.java | 490 +++++++++++++++++++++++++++++++++++++++++++-
opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogTestFixtures.java | 68 ++++++
opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileReplicaDBTest.java | 37 ---
4 files changed, 601 insertions(+), 49 deletions(-)
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 98f5e21..6561be7 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
@@ -212,7 +212,19 @@
throw new ChangelogException(ERR_CANNOT_CREATE_REPLICA_DB_BECAUSE_CHANGELOG_DB_SHUTDOWN.get());
}
- private ConcurrentMap<Integer, FileReplicaDB> getExistingOrNewDomainMap(final DN baseDN)
+ /**
+ * Returns the map holding the replica DBs of the provided domain, inserting a new one if it does
+ * not exist yet.
+ * <p>
+ * Package private and overridable so that tests can stop a thread right after it has read the
+ * shutdown flag in {@link #getOrCreateReplicaDB(DN, int, ReplicationServer)}, i.e. inside the
+ * window {@link #shutdownDB()} races with.
+ *
+ * @param baseDN
+ * the baseDN whose map of replica DBs must be returned
+ * @return the map of replica DBs of the provided domain
+ */
+ ConcurrentMap<Integer, FileReplicaDB> getExistingOrNewDomainMap(final DN baseDN)
{
// happy path: the domainMap already exists
final ConcurrentMap<Integer, FileReplicaDB> currentValue = domainToReplicaDBs.get(baseDN);
@@ -269,12 +281,51 @@
return null;
}
- final FileReplicaDB newDB = new FileReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv);
+ if (shutdown.get())
+ {
+ // A shutdown was initiated after the shutdown flag was read by getOrCreateReplicaDB():
+ // it may already have drained domainToReplicaDBs before this domainMap was inserted into
+ // it, in which case nothing would ever shutdown a replicaDB created here.
+ // Reading false instead means shutdownDB() has not flipped the flag yet, hence has not
+ // created its iterator yet either: since ConcurrentHashMap iterators traverse the
+ // elements as they existed upon construction of the iterator, it will see this domainMap,
+ // which was inserted before this monitor was acquired, and will have to block on this
+ // same monitor to drain it.
+ return null;
+ }
+
+ final FileReplicaDB newDB = newReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv);
domainMap.put(serverId, newDB);
return Pair.of(newDB, true);
}
}
+ /**
+ * Creates a new replica DB.
+ * <p>
+ * Package private and overridable so that tests can control the creation and the shutdown of the
+ * replica DBs this changelog holds.
+ *
+ * @param serverId
+ * the serverId for which to create a replica DB
+ * @param baseDN
+ * the baseDN for which to create a replica DB
+ * @param server
+ * the ReplicationServer
+ * @param cryptoSuite
+ * the cryptosuite to use for encryption
+ * @param replicationEnv
+ * the replication environment holding the log of the replica DB
+ * @return the newly created replica DB
+ * @throws ChangelogException
+ * if a problem occurred with the database
+ */
+ FileReplicaDB newReplicaDB(final int serverId, final DN baseDN, final ReplicationServer server,
+ final CryptoSuite cryptoSuite, final ReplicationEnvironment replicationEnv) throws ChangelogException
+ {
+ return new FileReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv);
+ }
+
@Override
public void initializeDB() throws ChangelogException
{
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogDBTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogDBTest.java
index 6e86009..595338a 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogDBTest.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogDBTest.java
@@ -15,6 +15,7 @@
*/
package org.opends.server.replication.server.changelog.file;
+import java.io.File;
import java.lang.management.LockInfo;
import java.lang.management.ManagementFactory;
import java.lang.management.ThreadInfo;
@@ -22,23 +23,36 @@
import java.lang.reflect.Field;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicReference;
+import org.assertj.core.api.SoftAssertions;
+import org.forgerock.opendj.config.server.ConfigException;
import org.forgerock.opendj.ldap.DN;
+import org.forgerock.opendj.server.config.server.MonitorProviderCfg;
import org.opends.server.TestCaseUtils;
+import org.opends.server.api.MonitorProvider;
+import org.opends.server.core.DirectoryServer;
+import org.opends.server.crypto.CryptoSuite;
import org.opends.server.replication.ReplicationTestCase;
-import org.opends.server.replication.server.ReplServerFakeConfiguration;
import org.opends.server.replication.server.ReplicationServer;
+import org.opends.server.replication.server.ReplicationServerDomain;
+import org.opends.server.replication.server.changelog.api.ChangelogException;
import org.testng.annotations.BeforeClass;
import org.testng.annotations.Test;
import static org.assertj.core.api.Assertions.*;
+import static org.opends.messages.ReplicationMessages.*;
import static org.opends.server.TestCaseUtils.*;
+import static org.opends.server.replication.server.changelog.file.FileChangelogTestFixtures.*;
+import static org.opends.server.util.StaticUtils.toLowerCase;
import static org.testng.Assert.*;
/**
- * Tests the {@link FileChangelogDB} class, and especially the window between
+ * Test the FileChangelogDB class: the races between a replica DB creation and
+ * {@link FileChangelogDB#shutdownDB()}, and the window between
* {@link FileChangelogDB#removeDomain(DN)}'s unlocked read of the domainMap and its
* acquisition of the domainMap monitor, during which a concurrent remover
* ({@code shutdownDB()}, {@code clearDB()} or another {@code removeDomain()}) may have
@@ -47,7 +61,13 @@
@SuppressWarnings("javadoc")
public class FileChangelogDBTest extends ReplicationTestCase
{
+ /** Server id of the replica DB which is shut down by the drain of the changelog. */
+ private static final int DRAINED_SERVER_ID = 814;
+ /** Server id of the replica DB whose creation races that drain. */
+ private static final int RACING_SERVER_ID = 813;
+ /** Server id of the replica DB whose domain removal races a concurrent remover. */
private static final int SERVER_ID = 1;
+ private static final long TIMEOUT_MS = 30000;
private DN TEST_ROOT_DN;
@@ -58,6 +78,240 @@
}
/**
+ * A replica DB whose creation loses the race against {@code shutdownDB()} must not be created at
+ * all: it would be held by a domain map the shutdown has already drained, so nothing would ever
+ * shut it down, and its monitor provider would stay registered for the lifetime of the process.
+ * <p>
+ * The interleaving is driven step by step:
+ * <ol>
+ * <li>the creator thread reads the shutdown flag, sees {@code false}, and is held there, before
+ * it inserts the domain map it needs;</li>
+ * <li>the shutdown flips the flag and drains {@code domainToReplicaDBs}, and is held inside the
+ * shutdown of the replica DB it found, i.e. once that domain map has been removed and while the
+ * replication environment is still open;</li>
+ * <li>the creator is released into that window.</li>
+ * </ol>
+ */
+ @Test
+ public void replicaDBLosingTheRaceAgainstShutdownIsNotCreated() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ ReplicationServer replicationServer = null;
+ RaceableChangelogDB changelogDB = null;
+ File testRoot = null;
+ Thread creator = null;
+ Thread shutdowner = null;
+ final AtomicReference<Throwable> creationFailure = new AtomicReference<>();
+ final AtomicReference<Throwable> shutdownFailure = new AtomicReference<>();
+ try
+ {
+ replicationServer = configureReplicationServer(100, 5000);
+ testRoot = createCleanDir("FileChangelogDB");
+ changelogDB = new RaceableChangelogDB(replicationServer, testRoot.getPath(), createCryptoSuite(false));
+ changelogDB.initializeDB();
+
+ // the replica DB the drain will be held in, and which is the only one registered so far
+ changelogDB.holdNextReplicaDBInItsShutdown();
+ changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, DRAINED_SERVER_ID, replicationServer);
+ // asserted, so that the test cannot pass by looking for a registration it cannot see
+ assertThat(DirectoryServer.getMonitorProviders().keySet())
+ .as("the replica DB held by the drain is not registered")
+ .contains(replicaDBMonitorName(replicationServer, DRAINED_SERVER_ID));
+ assertThat(DirectoryServer.getMonitorProviders().keySet())
+ .doesNotContain(replicaDBMonitorName(replicationServer, RACING_SERVER_ID));
+
+ final FileChangelogDB racedChangelogDB = changelogDB;
+ final ReplicationServer racedReplicationServer = replicationServer;
+ changelogDB.holdNextReplicaDBCreationBeforeItsDomainMapIsInserted();
+ creator = new Thread("FileChangelogDBTest replica DB creator")
+ {
+ @Override
+ public void run()
+ {
+ try
+ {
+ racedChangelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, RACING_SERVER_ID, racedReplicationServer);
+ }
+ catch (Throwable t)
+ {
+ creationFailure.set(t);
+ }
+ }
+ };
+ creator.start();
+ changelogDB.awaitCreatorInWindow();
+
+ shutdowner = new Thread("FileChangelogDBTest changelog shutdown")
+ {
+ @Override
+ public void run()
+ {
+ try
+ {
+ racedChangelogDB.shutdownDB();
+ }
+ catch (Throwable t)
+ {
+ shutdownFailure.set(t);
+ }
+ }
+ };
+ shutdowner.start();
+ changelogDB.awaitDrainInReplicaDBShutdown();
+
+ changelogDB.releaseCreator();
+ creator.join(TIMEOUT_MS);
+
+ assertThat(creator.isAlive()).as("the creator thread did not complete").isFalse();
+ final SoftAssertions softly = new SoftAssertions();
+ softly.assertThat(creationFailure.get())
+ .as("a replica DB created while the changelog is being drained is released by nobody")
+ .isInstanceOf(ChangelogException.class)
+ .hasMessage(ERR_CANNOT_CREATE_REPLICA_DB_BECAUSE_CHANGELOG_DB_SHUTDOWN.get().toString());
+ softly.assertThat(DirectoryServer.getMonitorProviders().keySet())
+ .as("monitor providers of the replica DBs created during the shutdown")
+ .doesNotContain(replicaDBMonitorName(replicationServer, RACING_SERVER_ID));
+ softly.assertAll();
+
+ changelogDB.releaseDrain();
+ shutdowner.join(TIMEOUT_MS);
+ assertThat(shutdowner.isAlive()).as("the shutdown thread did not complete").isFalse();
+ assertThat(shutdownFailure.get()).isNull();
+ assertThat(DirectoryServer.getMonitorProviders().keySet())
+ .as("the drained replica DB is still registered")
+ .doesNotContain(replicaDBMonitorName(replicationServer, DRAINED_SERVER_ID));
+ }
+ finally
+ {
+ if (changelogDB != null)
+ {
+ changelogDB.releaseAllHeldThreads();
+ changelogDB.shutdownDB();
+ }
+ join(creator);
+ join(shutdowner);
+ deregisterLeakedReplicaDBMonitors(replicationServer);
+ remove(replicationServer);
+ TestCaseUtils.deleteDirectory(testRoot);
+ }
+ }
+
+ /**
+ * A replica DB whose creation wins the race against {@code shutdownDB()} - i.e. reads the
+ * shutdown flag as {@code false} under the domain map monitor - must be shut down by the drain:
+ * the drain builds its iterator over {@code domainToReplicaDBs} after the flag is flipped, so it
+ * sees the domain map inserted before that monitor was taken, and blocks on the monitor until
+ * the creation has published the new replica DB.
+ * <p>
+ * This is the branch the fix in {@code getExistingOrNewReplicaDB()} relies on: a drain rewritten
+ * to no longer traverse the map as it existed when the flag was flipped - snapshotting the keys
+ * beforehand, shutting the replication environment down first - would silently reintroduce the
+ * leak this test guards against.
+ * <p>
+ * The interleaving is driven step by step:
+ * <ol>
+ * <li>the creator thread creates its replica DB - the monitor provider is now registered - and
+ * is held before the DB is published into the domain map, still under the domain map
+ * monitor;</li>
+ * <li>the shutdown starts, flips the flag, and blocks on the domain map monitor the creator
+ * holds;</li>
+ * <li>the creator is released: it publishes the replica DB and exits the monitor, and the drain
+ * must then shut that replica DB down.</li>
+ * </ol>
+ */
+ @Test
+ public void replicaDBWinningTheRaceAgainstShutdownIsShutDownByTheDrain() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ ReplicationServer replicationServer = null;
+ RaceableChangelogDB changelogDB = null;
+ File testRoot = null;
+ Thread creator = null;
+ Thread shutdowner = null;
+ final AtomicReference<Throwable> creationFailure = new AtomicReference<>();
+ final AtomicReference<Throwable> shutdownFailure = new AtomicReference<>();
+ final AtomicReference<FileReplicaDB> createdReplicaDB = new AtomicReference<>();
+ try
+ {
+ replicationServer = configureReplicationServer(100, 5000);
+ testRoot = createCleanDir("FileChangelogDB");
+ changelogDB = new RaceableChangelogDB(replicationServer, testRoot.getPath(), createCryptoSuite(false));
+ changelogDB.initializeDB();
+
+ final FileChangelogDB racedChangelogDB = changelogDB;
+ final ReplicationServer racedReplicationServer = replicationServer;
+ changelogDB.holdNextReplicaDBOnceCreated();
+ creator = new Thread("FileChangelogDBTest replica DB creator")
+ {
+ @Override
+ public void run()
+ {
+ try
+ {
+ createdReplicaDB.set(racedChangelogDB
+ .getOrCreateReplicaDB(TEST_ROOT_DN, RACING_SERVER_ID, racedReplicationServer).getFirst());
+ }
+ catch (Throwable t)
+ {
+ creationFailure.set(t);
+ }
+ }
+ };
+ creator.start();
+ changelogDB.awaitCreatorHoldingItsCreatedReplicaDB();
+ // asserted, so that the deregistration below cannot pass by never having seen a registration
+ assertThat(DirectoryServer.getMonitorProviders().keySet())
+ .as("the racing replica DB is not registered")
+ .contains(replicaDBMonitorName(replicationServer, RACING_SERVER_ID));
+
+ shutdowner = new Thread("FileChangelogDBTest changelog shutdown")
+ {
+ @Override
+ public void run()
+ {
+ try
+ {
+ racedChangelogDB.shutdownDB();
+ }
+ catch (Throwable t)
+ {
+ shutdownFailure.set(t);
+ }
+ }
+ };
+ shutdowner.start();
+ awaitBlockedOnAMonitor(shutdowner);
+
+ changelogDB.releaseCreatedReplicaDB();
+ creator.join(TIMEOUT_MS);
+ shutdowner.join(TIMEOUT_MS);
+ assertThat(creator.isAlive()).as("the creator thread did not complete").isFalse();
+ assertThat(shutdowner.isAlive()).as("the shutdown thread did not complete").isFalse();
+ assertThat(creationFailure.get()).as("a creation which won the race must succeed").isNull();
+ assertThat(createdReplicaDB.get()).as("the replica DB which won the race").isNotNull();
+ assertThat(shutdownFailure.get()).isNull();
+ assertThat(DirectoryServer.getMonitorProviders().keySet())
+ .as("the replica DB which won the race is not shut down by the drain")
+ .doesNotContain(replicaDBMonitorName(replicationServer, RACING_SERVER_ID));
+ }
+ finally
+ {
+ if (changelogDB != null)
+ {
+ changelogDB.releaseAllHeldThreads();
+ changelogDB.shutdownDB();
+ }
+ join(creator);
+ join(shutdowner);
+ deregisterLeakedReplicaDBMonitors(replicationServer);
+ remove(replicationServer);
+ TestCaseUtils.deleteDirectory(testRoot);
+ }
+ }
+
+ /**
* The concurrent remover unmapped the domain and shut its replica DBs down, exactly like
* the {@code shutdownDB()} drain does: {@code removeDomain()} must complete without
* throwing a {@link NullPointerException}.
@@ -69,7 +323,7 @@
try
{
TestCaseUtils.startServer();
- replicationServer = newReplicationServer();
+ replicationServer = configureReplicationServer(100, 100);
final FileChangelogDB changelogDB = (FileChangelogDB) replicationServer.getChangelogDB();
final FileReplicaDB replicaDB =
changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, SERVER_ID, replicationServer).getFirst();
@@ -90,7 +344,7 @@
domainToReplicaDBs.remove(TEST_ROOT_DN);
replicaDB.shutdown();
}
- remover.join(TimeUnit.SECONDS.toMillis(30));
+ remover.join(TIMEOUT_MS);
assertFalse(remover.isAlive(), "removeDomain() did not complete");
assertThat(thrown.get()).isNull();
@@ -113,7 +367,7 @@
try
{
TestCaseUtils.startServer();
- replicationServer = newReplicationServer();
+ replicationServer = configureReplicationServer(100, 100);
final FileChangelogDB changelogDB = (FileChangelogDB) replicationServer.getChangelogDB();
final FileReplicaDB replicaDB =
changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, SERVER_ID, replicationServer).getFirst();
@@ -135,7 +389,7 @@
replicaDB.shutdown();
domainToReplicaDBs.put(TEST_ROOT_DN, recreatedDomainMap);
}
- remover.join(TimeUnit.SECONDS.toMillis(30));
+ remover.join(TIMEOUT_MS);
assertFalse(remover.isAlive(), "removeDomain() did not complete");
assertThat(thrown.get()).isNull();
@@ -147,13 +401,6 @@
}
}
- private ReplicationServer newReplicationServer() throws Exception
- {
- final int changelogPort = findFreePort();
- return new ReplicationServer(
- new ReplServerFakeConfiguration(changelogPort, null, 0, 2, 100, 100, null));
- }
-
private Thread newRemoverThread(final FileChangelogDB changelogDB, final AtomicReference<Throwable> thrown)
{
return new Thread(new Runnable()
@@ -182,10 +429,11 @@
return (ConcurrentMap<DN, ConcurrentMap<Integer, FileReplicaDB>>) field.get(changelogDB);
}
+ /** Waits until the provided thread is blocked acquiring the monitor of the provided object. */
private void waitUntilBlockedOn(Thread thread, Object monitor) throws Exception
{
final ThreadMXBean threadMXBean = ManagementFactory.getThreadMXBean();
- final long deadline = System.currentTimeMillis() + TimeUnit.SECONDS.toMillis(30);
+ final long deadline = System.currentTimeMillis() + TIMEOUT_MS;
while (System.currentTimeMillis() < deadline)
{
final ThreadInfo threadInfo = threadMXBean.getThreadInfo(thread.getId());
@@ -201,4 +449,218 @@
throw new AssertionError(
"Timed out waiting for " + thread.getName() + " to block on the domainMap monitor");
}
+
+ /** Joins the provided thread, leaving a signal behind when it did not die within the timeout. */
+ private void join(final Thread thread) throws InterruptedException
+ {
+ if (thread != null)
+ {
+ thread.join(TIMEOUT_MS);
+ if (thread.isAlive())
+ {
+ final IllegalStateException hung = new IllegalStateException("Test thread " + thread.getName()
+ + " is still alive after " + TIMEOUT_MS + " ms: it may leak a live changelog into later tests");
+ hung.setStackTrace(thread.getStackTrace());
+ hung.printStackTrace();
+ thread.interrupt();
+ }
+ }
+ }
+
+ /**
+ * Waits until the provided thread is blocked acquiring a monitor: the domain map monitor held by
+ * the creator is the only one it can stay blocked on - the other locks on its way to the drain
+ * are only transiently contended, hence the two consecutive observations.
+ */
+ private static void awaitBlockedOnAMonitor(final Thread thread) throws InterruptedException
+ {
+ final long deadline = System.currentTimeMillis() + TIMEOUT_MS;
+ int blockedObservations = 0;
+ while (blockedObservations < 2)
+ {
+ if (!thread.isAlive())
+ {
+ throw new IllegalStateException(thread.getName() + " completed without blocking on the domain map monitor");
+ }
+ if (System.currentTimeMillis() > deadline)
+ {
+ throw new IllegalStateException(
+ "timed out waiting for " + thread.getName() + " to block on the domain map monitor");
+ }
+ blockedObservations = thread.getState() == Thread.State.BLOCKED ? blockedObservations + 1 : 0;
+ Thread.sleep(1);
+ }
+ }
+
+ /**
+ * Returns the name the monitor provider of the provided replica DB is registered under, i.e. the
+ * name built by {@code FileReplicaDB.DbMonitorProvider.getMonitorInstanceName()}, lower-cased
+ * the way {@code DirectoryServer.registerMonitorProvider()} stores it.
+ */
+ private String replicaDBMonitorName(final ReplicationServer replicationServer, final int serverId)
+ {
+ final ReplicationServerDomain domain = replicationServer.getReplicationServerDomain(TEST_ROOT_DN);
+ assertThat(domain).as("the domain scoping the monitor name of DS(" + serverId + ")").isNotNull();
+ return toLowerCase("Changelog for DS(" + serverId + "),cn=" + domain.getMonitorInstanceName());
+ }
+
+ /** Releases the monitor providers a regression leaks, so that they do not outlive this test. */
+ private void deregisterLeakedReplicaDBMonitors(final ReplicationServer replicationServer)
+ {
+ if (replicationServer == null || replicationServer.getReplicationServerDomain(TEST_ROOT_DN) == null)
+ {
+ return; // no replica DB was ever created, hence no monitor provider was ever registered
+ }
+ for (final int serverId : new int[] { RACING_SERVER_ID, DRAINED_SERVER_ID })
+ {
+ // deregister the provider instead of removing the map entry, so that the JMX MBean
+ // registered alongside it is released as well
+ final MonitorProvider<? extends MonitorProviderCfg> provider =
+ DirectoryServer.getMonitorProviders().get(replicaDBMonitorName(replicationServer, serverId));
+ if (provider != null)
+ {
+ DirectoryServer.deregisterMonitorProvider(provider);
+ }
+ }
+ }
+
+ /**
+ * A changelog DB which lets a test hold a thread creating a replica DB right after it has read
+ * the shutdown flag, hold it again once the replica DB is created but not yet published into the
+ * domain map, and hold the shutdown inside the drain of {@code domainToReplicaDBs}.
+ */
+ private static final class RaceableChangelogDB extends FileChangelogDB
+ {
+ private final AtomicBoolean holdNextCreation = new AtomicBoolean();
+ private final AtomicBoolean holdNextReplicaDBShutdown = new AtomicBoolean();
+ private final AtomicBoolean holdNextCreatedReplicaDB = new AtomicBoolean();
+ private final CountDownLatch creatorIsInWindow = new CountDownLatch(1);
+ private final CountDownLatch creatorIsReleased = new CountDownLatch(1);
+ private final CountDownLatch creatorHoldsItsCreatedReplicaDB = new CountDownLatch(1);
+ private final CountDownLatch createdReplicaDBIsReleased = new CountDownLatch(1);
+ private final CountDownLatch drainIsInReplicaDBShutdown = new CountDownLatch(1);
+ private final CountDownLatch drainIsReleased = new CountDownLatch(1);
+
+ RaceableChangelogDB(final ReplicationServer replicationServer, final String dbDirectoryPath,
+ final CryptoSuite cryptoSuite) throws ConfigException
+ {
+ super(replicationServer, dbDirectoryPath, cryptoSuite);
+ }
+
+ @Override
+ ConcurrentMap<Integer, FileReplicaDB> getExistingOrNewDomainMap(final DN baseDN)
+ {
+ if (holdNextCreation.compareAndSet(true, false))
+ {
+ creatorIsInWindow.countDown();
+ await(creatorIsReleased);
+ }
+ return super.getExistingOrNewDomainMap(baseDN);
+ }
+
+ @Override
+ FileReplicaDB newReplicaDB(final int serverId, final DN baseDN, final ReplicationServer server,
+ final CryptoSuite cryptoSuite, final ReplicationEnvironment replicationEnv) throws ChangelogException
+ {
+ if (holdNextReplicaDBShutdown.compareAndSet(true, false))
+ {
+ return new HeldOnShutdownReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv);
+ }
+ final FileReplicaDB replicaDB = super.newReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv);
+ if (holdNextCreatedReplicaDB.compareAndSet(true, false))
+ {
+ // the replica DB exists and its monitor provider is registered, but it is not published
+ // into the domain map yet: hold the creator there, under the domain map monitor
+ creatorHoldsItsCreatedReplicaDB.countDown();
+ await(createdReplicaDBIsReleased);
+ }
+ return replicaDB;
+ }
+
+ void holdNextReplicaDBCreationBeforeItsDomainMapIsInserted()
+ {
+ holdNextCreation.set(true);
+ }
+
+ void holdNextReplicaDBInItsShutdown()
+ {
+ holdNextReplicaDBShutdown.set(true);
+ }
+
+ void holdNextReplicaDBOnceCreated()
+ {
+ holdNextCreatedReplicaDB.set(true);
+ }
+
+ void awaitCreatorInWindow()
+ {
+ await(creatorIsInWindow);
+ }
+
+ void awaitCreatorHoldingItsCreatedReplicaDB()
+ {
+ await(creatorHoldsItsCreatedReplicaDB);
+ }
+
+ void awaitDrainInReplicaDBShutdown()
+ {
+ await(drainIsInReplicaDBShutdown);
+ }
+
+ void releaseCreator()
+ {
+ creatorIsReleased.countDown();
+ }
+
+ void releaseCreatedReplicaDB()
+ {
+ createdReplicaDBIsReleased.countDown();
+ }
+
+ void releaseDrain()
+ {
+ drainIsReleased.countDown();
+ }
+
+ void releaseAllHeldThreads()
+ {
+ releaseCreator();
+ releaseCreatedReplicaDB();
+ releaseDrain();
+ }
+
+ /** A replica DB which holds the thread shutting it down until the test releases it. */
+ private final class HeldOnShutdownReplicaDB extends FileReplicaDB
+ {
+ HeldOnShutdownReplicaDB(final int serverId, final DN baseDN, final ReplicationServer server,
+ final CryptoSuite cryptoSuite, final ReplicationEnvironment replicationEnv) throws ChangelogException
+ {
+ super(serverId, baseDN, server, cryptoSuite, replicationEnv);
+ }
+
+ @Override
+ void shutdown()
+ {
+ drainIsInReplicaDBShutdown.countDown();
+ await(drainIsReleased);
+ super.shutdown();
+ }
+ }
+
+ private static void await(final CountDownLatch latch)
+ {
+ try
+ {
+ if (!latch.await(TIMEOUT_MS, TimeUnit.MILLISECONDS))
+ {
+ throw new IllegalStateException("timed out waiting for the replica DB creation race");
+ }
+ }
+ catch (InterruptedException e)
+ {
+ Thread.currentThread().interrupt();
+ throw new IllegalStateException(e);
+ }
+ }
+ }
}
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogTestFixtures.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogTestFixtures.java
new file mode 100644
index 0000000..1accbc8
--- /dev/null
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogTestFixtures.java
@@ -0,0 +1,68 @@
+/*
+ * 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.server.changelog.file;
+
+import java.io.File;
+import java.io.IOException;
+
+import org.forgerock.opendj.config.server.ConfigException;
+import org.opends.server.TestCaseUtils;
+import org.opends.server.crypto.CryptoSuite;
+import org.opends.server.replication.server.ReplServerFakeConfiguration;
+import org.opends.server.replication.server.ReplicationServer;
+
+import static org.opends.server.TestCaseUtils.*;
+
+/** Fixtures shared by the tests of the file based changelog. */
+final class FileChangelogTestFixtures
+{
+ static final String CIPHER_TRANSFORMATION = "AES/CBC/PKCS5Padding";
+ static final int KEY_LENGTH = 128;
+
+ private FileChangelogTestFixtures()
+ {
+ // static helpers only
+ }
+
+ /** Returns a replication server listening on a free port, with no connected replica. */
+ static ReplicationServer configureReplicationServer(int windowSize, int queueSize)
+ throws IOException, ConfigException
+ {
+ final int changelogPort = findFreePort();
+ ReplServerFakeConfiguration replServerFakeCfg =
+ new ReplServerFakeConfiguration(changelogPort, null, 0, 2, queueSize, windowSize, null);
+ return new ReplicationServer(replServerFakeCfg);
+ }
+
+ /** Returns a crypto suite the changelog can encrypt its records with. */
+ static CryptoSuite createCryptoSuite(boolean confidential)
+ {
+ return getServerContext().getCryptoManager().newCryptoSuite(CIPHER_TRANSFORMATION, KEY_LENGTH, confidential);
+ }
+
+ /** Returns an empty directory of the provided name under the unit test build directory. */
+ static File createCleanDir(String directoryName) throws IOException
+ {
+ String buildRoot = System.getProperty(TestCaseUtils.PROPERTY_BUILD_ROOT);
+ String path = System.getProperty(TestCaseUtils.PROPERTY_BUILD_DIR, buildRoot
+ + File.separator + "build");
+ path = path + File.separator + "unit-tests" + File.separator + directoryName;
+ final File testRoot = new File(path);
+ TestCaseUtils.deleteDirectory(testRoot);
+ testRoot.mkdirs();
+ return testRoot;
+ }
+}
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileReplicaDBTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileReplicaDBTest.java
index 5e3ce5e..30eea61 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileReplicaDBTest.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileReplicaDBTest.java
@@ -18,12 +18,10 @@
import java.io.File;
import java.io.FileOutputStream;
-import java.io.IOException;
import java.util.ArrayList;
import org.assertj.core.api.SoftAssertions;
import org.forgerock.i18n.slf4j.LocalizedLogger;
-import org.forgerock.opendj.config.server.ConfigException;
import org.forgerock.opendj.ldap.ByteString;
import org.forgerock.opendj.ldap.ByteStringBuilder;
import org.forgerock.opendj.ldap.DN;
@@ -35,7 +33,6 @@
import org.opends.server.replication.common.CSNGenerator;
import org.opends.server.replication.protocol.DeleteMsg;
import org.opends.server.replication.protocol.UpdateMsg;
-import org.opends.server.replication.server.ReplServerFakeConfiguration;
import org.opends.server.replication.server.ReplicationServer;
import org.opends.server.replication.server.changelog.api.ChangelogException;
import org.opends.server.replication.server.changelog.api.DBCursor;
@@ -48,6 +45,7 @@
import static org.opends.server.TestCaseUtils.*;
import static org.opends.server.replication.server.changelog.api.DBCursor.KeyMatchingStrategy.*;
import static org.opends.server.replication.server.changelog.api.DBCursor.PositionStrategy.*;
+import static org.opends.server.replication.server.changelog.file.FileChangelogTestFixtures.*;
import static org.opends.server.util.CollectionUtils.*;
import static org.testng.Assert.*;
@@ -58,8 +56,6 @@
public class FileReplicaDBTest extends ReplicationTestCase
{
private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();
- private final String cipherTransformation = "AES/CBC/PKCS5Padding";
- private final int keyLength = 128;
private DN TEST_ROOT_DN;
/**
@@ -108,16 +104,12 @@
RecordParser<CSN, UpdateMsg> parser = FileReplicaDB.newReplicaDBParser(cryptoSuite);
ByteString data1 = parser.encodeRecord(Record.from(msg.getCSN(), msg));
- cryptoSuite.newParameters(cipherTransformation, keyLength, !confidential);
+ cryptoSuite.newParameters(CIPHER_TRANSFORMATION, KEY_LENGTH, !confidential);
ByteString data2 = parser.encodeRecord(Record.from(msg.getCSN(), msg));
assertFalse(data1.equals(data2));
}
- private CryptoSuite createCryptoSuite(boolean confidential)
- {
- return getServerContext().getCryptoManager().newCryptoSuite(cipherTransformation, keyLength, confidential);
- }
@Test
public void testDomainDNWithForwardSlashes() throws Exception
{
@@ -378,7 +370,7 @@
{
TestCaseUtils.startServer();
replicationServer = configureReplicationServer(100000, 10);
- testRoot = createCleanDir();
+ testRoot = createCleanDir("FileReplicaDB");
dbEnv = new ReplicationEnvironment(testRoot.getPath(), replicationServer, TimeService.SYSTEM);
replicaDB = new FileReplicaDB(1, TEST_ROOT_DN, replicationServer, createCryptoSuite(false), dbEnv);
@@ -536,7 +528,7 @@
TestCaseUtils.startServer();
replicationServer = configureReplicationServer(100000, 10);
- testRoot = createCleanDir();
+ testRoot = createCleanDir("FileReplicaDB");
dbEnv = new ReplicationEnvironment(testRoot.getPath(), replicationServer, TimeService.SYSTEM);
replicaDB = new FileReplicaDB(1, TEST_ROOT_DN, replicationServer, createCryptoSuite(false), dbEnv);
@@ -648,33 +640,12 @@
assertEquals(replicaDB.getNumberRecords(), expectedNbRecords);
}
- private ReplicationServer configureReplicationServer(int windowSize, int queueSize)
- throws IOException, ConfigException
- {
- final int changelogPort = findFreePort();
- ReplServerFakeConfiguration replServerFakeCfg =
- new ReplServerFakeConfiguration(changelogPort, null, 0, 2, queueSize, windowSize, null);
- return new ReplicationServer(replServerFakeCfg);
- }
-
private FileReplicaDB newReplicaDB(ReplicationServer rs) throws Exception
{
final FileChangelogDB changelogDB = (FileChangelogDB) rs.getChangelogDB();
return changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, 1, rs).getFirst();
}
- private File createCleanDir() throws IOException
- {
- String buildRoot = System.getProperty(TestCaseUtils.PROPERTY_BUILD_ROOT);
- String path = System.getProperty(TestCaseUtils.PROPERTY_BUILD_DIR, buildRoot
- + File.separator + "build");
- path = path + File.separator + "unit-tests" + File.separator + "FileReplicaDB";
- final File testRoot = new File(path);
- TestCaseUtils.deleteDirectory(testRoot);
- testRoot.mkdirs();
- return testRoot;
- }
-
private void assertFoundInOrder(FileReplicaDB replicaDB, CSN... csns) throws Exception
{
if (csns.length == 0)
--
Gitblit v1.10.0