From 6bd69778b12ae382f67bf96ef2f8d49b45eadeed Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Wed, 05 Aug 2026 17:01:27 +0000
Subject: [PATCH] [#818] Drop the empty domain map a replica DB creation leaves behind when it bails out (#830)
---
opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/ECLMultiDomainDBCursorTest.java | 20 +
opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/FileChangelogDB.java | 90 ++++++-
opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/FileChangelogDBTest.java | 523 +++++++++++++++++++++++++++++++++++++++----
opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/SequentialDBCursor.java | 10
opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/MultiDomainDBCursor.java | 33 ++
5 files changed, 606 insertions(+), 70 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 6561be7..8a91472 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
@@ -81,6 +81,12 @@
* </ol>
* When creating a replicaDB, synchronize on the domainMap to avoid
* concurrent shutdown.
+ * <p>
+ * A creation which bails out while holding the domainMap monitor removes the still empty
+ * domainMap it inserted, so that no phantom domain is left behind. It may only do so after
+ * checking, under that same monitor, that the domainMap is still the one mapped to its baseDN:
+ * the removal is equality based and two empty maps are equal, so without the identity check it
+ * could unmap the fresh domainMap of another creation.
*/
private final ConcurrentMap<DN, ConcurrentMap<Integer, FileReplicaDB>> domainToReplicaDBs =
new ConcurrentHashMap<>();
@@ -278,25 +284,49 @@
// 1) a shutdown was initiated or 2) an initialize was called.
// Return will allow the code to:
// 1) shutdown properly or 2) lazily recreate the replicaDB
+ // There is nothing to clean up here: this domainMap is already unmapped, and whatever map
+ // is now associated to baseDN belongs to another creation.
return null;
}
- if (shutdown.get())
+ try
{
- // 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;
- }
+ 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);
+ final FileReplicaDB newDB = newReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv);
+ domainMap.put(serverId, newDB);
+ return Pair.of(newDB, true);
+ }
+ finally
+ {
+ // Leaving without having created the replica DB must not leave behind the empty domainMap
+ // inserted by getExistingOrNewDomainMap(): nothing would ever remove it, and every multi
+ // domain cursor created afterwards would walk a domain holding no replica DB at all.
+ // Only an empty map may be dropped: a populated one must stay mapped for the drain of
+ // shutdownDB() to find, even when the creation of this serverId failed.
+ // The identity check above passed under this monitor, and every removal site takes the
+ // monitor of the map it unmaps before removing it: the mapping cannot have changed since,
+ // so this equality based remove provably drops this domainMap and no other.
+ // Unlike removeDomain(), this path clears no ChangeNumberIndexer state, and a later
+ // creation broadcasts addDomain() anew to multi domain cursors which already incorporated
+ // the domain: MultiDomainDBCursor discards such an announcement when it incorporates new
+ // cursors, so no second cursor is opened over the domain.
+ if (domainMap.isEmpty())
+ {
+ domainToReplicaDBs.remove(baseDN, domainMap);
+ }
+ }
}
}
@@ -326,6 +356,20 @@
return new FileReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv);
}
+ /**
+ * Returns the map of replica DBs per domain.
+ * <p>
+ * Package private, for tests only: they both observe which domain maps this changelog holds and
+ * mutate the map to drive race interleavings, so this getter intentionally returns the live
+ * internal map, not a copy or an unmodifiable view.
+ *
+ * @return the live map of replica DBs per domain
+ */
+ ConcurrentMap<DN, ConcurrentMap<Integer, FileReplicaDB>> getDomainToReplicaDBs()
+ {
+ return domainToReplicaDBs;
+ }
+
@Override
public void initializeDB() throws ChangelogException
{
@@ -403,13 +447,23 @@
firstException = e;
}
- for (Iterator<ConcurrentMap<Integer, FileReplicaDB>> it =
- this.domainToReplicaDBs.values().iterator(); it.hasNext();)
+ for (Iterator<Map.Entry<DN, ConcurrentMap<Integer, FileReplicaDB>>> it =
+ this.domainToReplicaDBs.entrySet().iterator(); it.hasNext();)
{
- final ConcurrentMap<Integer, FileReplicaDB> domainMap = it.next();
+ final Map.Entry<DN, ConcurrentMap<Integer, FileReplicaDB>> entry = it.next();
+ final ConcurrentMap<Integer, FileReplicaDB> domainMap = entry.getValue();
synchronized (domainMap)
{
- it.remove();
+ // Follow the removal protocol documented on domainToReplicaDBs: unmap only the domainMap
+ // instance the monitor was taken on. The check is identity based, like removeDomain()'s:
+ // an equality based remove could still drop a map whose monitor is not held, since two
+ // empty maps are equal and removeDomain() plus a fresh creation may have swapped one for
+ // another since this iterator read its entry. Replica DBs another remover already visited
+ // are shut down again below, which is harmless: shutdown is a no-op the second time.
+ if (domainToReplicaDBs.get(entry.getKey()) == domainMap)
+ {
+ domainToReplicaDBs.remove(entry.getKey());
+ }
for (FileReplicaDB replicaDB : domainMap.values())
{
replicaDB.shutdown();
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/MultiDomainDBCursor.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/MultiDomainDBCursor.java
index 8a9ec2d..4019bd7 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/MultiDomainDBCursor.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/changelog/file/MultiDomainDBCursor.java
@@ -12,12 +12,14 @@
* information: "Portions Copyright [year] [name of copyright owner]".
*
* Copyright 2014-2016 ForgeRock AS.
+ * Portions Copyright 2026 3A Systems, LLC.
*/
package org.opends.server.replication.server.changelog.file;
import java.util.Iterator;
import java.util.Map.Entry;
import java.util.concurrent.ConcurrentSkipListMap;
+import java.util.concurrent.ConcurrentSkipListSet;
import net.jcip.annotations.NotThreadSafe;
@@ -34,6 +36,18 @@
{
private final ReplicationDomainDB domainDB;
private final ConcurrentSkipListMap<DN, ServerState> newDomains = new ConcurrentSkipListMap<>();
+ /**
+ * The domains this cursor already iterates over, so that a second announcement of a domain is
+ * discarded at incorporation: a second cursor over the same domain would either leak unclosed
+ * - the cursor tree of {@link CompositeDBCursor} collapses cursors comparing equal - or
+ * deliver every change twice. A domain may be announced again after the empty domainMap of a
+ * failed replica DB creation was dropped, while this cursor still iterates over it.
+ * <p>
+ * Added to and removed from on the thread iterating this cursor - {@link #incorporateNewCursors()}
+ * and {@link #removeDomain(DN)} run on it - read by the threads announcing domains, and
+ * cleared by {@link #close()}, which an ending ECL session may call from another thread.
+ */
+ private final ConcurrentSkipListSet<DN> incorporatedDomains = new ConcurrentSkipListSet<>();
private final CursorOptions options;
/**
@@ -52,6 +66,10 @@
/**
* Adds a replication domain for this cursor to iterate over. Added cursors
* will be created and iterated over on the next call to {@link #next()}.
+ * <p>
+ * Announcing a domain this cursor already iterates over has no effect: the announcement is
+ * discarded when cursors are incorporated, and new replica DBs of such a domain reach it
+ * through {@link DomainDBCursor#addReplicaDB(int, org.opends.server.replication.common.CSN)}.
*
* @param baseDN
* the replication domain's baseDN
@@ -60,6 +78,9 @@
*/
public void addDomain(DN baseDN, ServerState startAfterState)
{
+ // incorporateNewCursors() discards announcements of domains this cursor already iterates
+ // over: checking incorporatedDomains here would be a check-then-act against removeDomain()
+ // on the cursor's thread, able to drop an announcement the removal no longer covers
newDomains.put(baseDN, startAfterState != null ? startAfterState : new ServerState());
}
@@ -73,8 +94,14 @@
final Entry<DN, ServerState> entry = iter.next();
final DN baseDN = entry.getKey();
final ServerState serverState = entry.getValue();
- final DBCursor<UpdateMsg> domainDBCursor = domainDB.getCursorFrom(baseDN, serverState, options);
- addCursor(domainDBCursor, baseDN);
+ // discard the announcement of a domain this cursor already iterates over: this is the only
+ // thread adding to and removing from incorporatedDomains, so the check cannot race them
+ if (!incorporatedDomains.contains(baseDN))
+ {
+ final DBCursor<UpdateMsg> domainDBCursor = domainDB.getCursorFrom(baseDN, serverState, options);
+ addCursor(domainDBCursor, baseDN);
+ incorporatedDomains.add(baseDN);
+ }
iter.remove();
}
}
@@ -90,6 +117,7 @@
public void removeDomain(DN baseDN)
{
removeCursor(baseDN);
+ incorporatedDomains.remove(baseDN);
}
/** {@inheritDoc} */
@@ -99,6 +127,7 @@
super.close();
domainDB.unregisterCursor(this);
newDomains.clear();
+ incorporatedDomains.clear();
}
}
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/ECLMultiDomainDBCursorTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/ECLMultiDomainDBCursorTest.java
index 6ef8d2f..e6e3ffd 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/ECLMultiDomainDBCursorTest.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/ECLMultiDomainDBCursorTest.java
@@ -12,10 +12,13 @@
* information: "Portions Copyright [year] [name of copyright owner]".
*
* Copyright 2014-2016 ForgeRock AS.
+ * Portions Copyright 2026 3A Systems, LLC.
*/
package org.opends.server.replication.server.changelog.file;
+import java.util.HashMap;
import java.util.HashSet;
+import java.util.Map;
import java.util.Set;
import org.forgerock.opendj.ldap.DN;
@@ -45,6 +48,8 @@
private MultiDomainDBCursor multiDomainCursor;
private ECLMultiDomainDBCursor eclCursor;
private final Set<DN> eclEnabledDomains = new HashSet<>();
+ /** The long-lived cursor of each domain announced to {@link #multiDomainCursor}. */
+ private final Map<DN, SequentialDBCursor> domainCursors = new HashMap<>();
private ECLEnabledDomainPredicate predicate = new ECLEnabledDomainPredicate()
{
@Override
@@ -61,6 +66,9 @@
options = new CursorOptions(GREATER_THAN_OR_EQUAL_TO_KEY, ON_MATCHING_KEY);
multiDomainCursor = new MultiDomainDBCursor(domainDB, options);
eclCursor = new ECLMultiDomainDBCursor(predicate, multiDomainCursor);
+ // one test class instance runs all the methods: the domains announced to the previous
+ // method's multiDomainCursor must not be mistaken for domains this one iterates over
+ domainCursors.clear();
}
@AfterMethod
@@ -185,6 +193,18 @@
private void addDomainCursorToCursor(DN baseDN, SequentialDBCursor cursor) throws ChangelogException
{
+ final SequentialDBCursor existing = domainCursors.get(baseDN);
+ if (existing != null)
+ {
+ // already known to the cursor: its long-lived per-domain cursor receives the new changes,
+ // exactly as DomainDBCursor.addReplicaDB() does in production
+ for (UpdateMsg msg : cursor.drain())
+ {
+ existing.add(msg);
+ }
+ return;
+ }
+ domainCursors.put(baseDN, cursor);
final ServerState state = new ServerState();
when(domainDB.getCursorFrom(baseDN, state, options)).thenReturn(cursor);
multiDomainCursor.addDomain(baseDN, state);
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 595338a..31bfa61 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
@@ -20,43 +20,58 @@
import java.lang.management.ManagementFactory;
import java.lang.management.ThreadInfo;
import java.lang.management.ThreadMXBean;
-import java.lang.reflect.Field;
+import java.util.List;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.CopyOnWriteArrayList;
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.i18n.LocalizableMessage;
import org.forgerock.opendj.config.server.ConfigException;
import org.forgerock.opendj.ldap.DN;
import org.forgerock.opendj.server.config.server.MonitorProviderCfg;
+import org.forgerock.util.Pair;
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.common.CSN;
+import org.opends.server.replication.common.MultiDomainServerState;
+import org.opends.server.replication.common.ServerState;
+import org.opends.server.replication.protocol.DeleteMsg;
+import org.opends.server.replication.protocol.UpdateMsg;
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.opends.server.replication.server.changelog.api.DBCursor;
+import org.opends.server.replication.server.changelog.api.DBCursor.CursorOptions;
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.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.StaticUtils.toLowerCase;
import static org.testng.Assert.*;
/**
* Test the FileChangelogDB class: the races between a replica DB creation and
- * {@link FileChangelogDB#shutdownDB()}, and the window between
+ * {@link FileChangelogDB#shutdownDB()}; 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
- * unmapped the domain.
+ * unmapped the domain; and the cleanup of a creation which bails out without having created a
+ * replica DB - the empty domainMap it inserted must be dropped, without unmapping the fresh
+ * domainMap of a concurrent creation and without announcing the domain a second time to the
+ * multi domain cursors that were live at the time.
*/
@SuppressWarnings("javadoc")
public class FileChangelogDBTest extends ReplicationTestCase
@@ -67,8 +82,19 @@
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;
+ /** Server id of the replica DB whose creation bails out on the identity check. */
+ private static final int STALE_SERVER_ID = 815;
+ /** Server id of the replica DB created in the fresh domain map the bail-out must not unmap. */
+ private static final int FRESH_SERVER_ID = 816;
+ /** Server id of a replica DB created before the failing one, in the same domain. */
+ private static final int EXISTING_SERVER_ID = 817;
+ /** Server id of the replica DB whose creation is made to fail. */
+ private static final int FAILING_SERVER_ID = 818;
private static final long TIMEOUT_MS = 30000;
+ private static final LocalizableMessage CREATION_FAILURE =
+ LocalizableMessage.raw("FileChangelogDBTest replica DB creation failure");
+
private DN TEST_ROOT_DN;
@BeforeClass
@@ -191,7 +217,7 @@
}
join(creator);
join(shutdowner);
- deregisterLeakedReplicaDBMonitors(replicationServer);
+ deregisterLeakedReplicaDBMonitors(replicationServer, RACING_SERVER_ID, DRAINED_SERVER_ID);
remove(replicationServer);
TestCaseUtils.deleteDirectory(testRoot);
}
@@ -282,7 +308,7 @@
}
};
shutdowner.start();
- awaitBlockedOnAMonitor(shutdowner);
+ waitUntilBlockedOn(shutdowner, changelogDB.getDomainToReplicaDBs().get(TEST_ROOT_DN));
changelogDB.releaseCreatedReplicaDB();
creator.join(TIMEOUT_MS);
@@ -305,13 +331,358 @@
}
join(creator);
join(shutdowner);
- deregisterLeakedReplicaDBMonitors(replicationServer);
+ deregisterLeakedReplicaDBMonitors(replicationServer, RACING_SERVER_ID);
remove(replicationServer);
TestCaseUtils.deleteDirectory(testRoot);
}
}
/**
+ * A replica DB creation which fails must not leave behind the empty domain map it inserted:
+ * nothing would ever remove it from {@code domainToReplicaDBs}, and every multi domain cursor
+ * created afterwards would walk a domain holding no replica DB at all - the symptom this test
+ * also asserts, through the domains a new multi domain cursor asks the changelog to open.
+ */
+ @Test
+ public void failedReplicaDBCreationDropsTheDomainMapItInserted() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ ReplicationServer replicationServer = null;
+ RaceableChangelogDB changelogDB = null;
+ File testRoot = null;
+ try
+ {
+ replicationServer = configureReplicationServer(100, 5000);
+ testRoot = createCleanDir("FileChangelogDB");
+ changelogDB = new RaceableChangelogDB(replicationServer, testRoot.getPath(), createCryptoSuite(false));
+ changelogDB.initializeDB();
+
+ changelogDB.failNextReplicaDBCreation();
+ try
+ {
+ changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, FAILING_SERVER_ID, replicationServer);
+ failBecauseExceptionWasNotThrown(ChangelogException.class);
+ }
+ catch (ChangelogException expected)
+ {
+ assertThat(expected).hasMessage(CREATION_FAILURE.toString());
+ }
+ assertThat(changelogDB.getDomainToReplicaDBs())
+ .as("the empty domain map inserted for the creation which failed")
+ .doesNotContainKey(TEST_ROOT_DN);
+
+ // the symptom of the leftover map: a multi domain cursor created after the failure must not
+ // walk the phantom domain
+ changelogDB.walkedDomains.clear();
+ final MultiDomainDBCursor cursor = changelogDB.getCursorFrom(
+ new MultiDomainServerState(), new CursorOptions(GREATER_THAN_OR_EQUAL_TO_KEY, ON_MATCHING_KEY));
+ try
+ {
+ cursor.next();
+ assertThat(changelogDB.walkedDomains)
+ .as("the domains walked by a multi domain cursor created after the failed creation")
+ .doesNotContain(TEST_ROOT_DN);
+ }
+ finally
+ {
+ cursor.close();
+ }
+
+ // the next creation starts from scratch and repopulates the domain
+ final Pair<FileReplicaDB, Boolean> result =
+ changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, FAILING_SERVER_ID, replicationServer);
+ assertThat(result.getSecond()).as("the replica DB was created anew").isTrue();
+ assertThat(changelogDB.getDomainToReplicaDBs().get(TEST_ROOT_DN)).containsOnlyKeys(FAILING_SERVER_ID);
+ }
+ finally
+ {
+ try
+ {
+ if (changelogDB != null)
+ {
+ changelogDB.shutdownDB();
+ }
+ }
+ finally
+ {
+ remove(replicationServer);
+ TestCaseUtils.deleteDirectory(testRoot);
+ }
+ }
+ }
+
+ /**
+ * A replica DB creation which fails must only drop an <b>empty</b> domain map: a populated one
+ * must stay mapped, so that the drain of {@code shutdownDB()} finds the replica DBs it holds and
+ * shuts them down.
+ */
+ @Test
+ public void failedReplicaDBCreationKeepsAPopulatedDomainMap() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ ReplicationServer replicationServer = null;
+ RaceableChangelogDB changelogDB = null;
+ File testRoot = null;
+ try
+ {
+ replicationServer = configureReplicationServer(100, 5000);
+ testRoot = createCleanDir("FileChangelogDB");
+ changelogDB = new RaceableChangelogDB(replicationServer, testRoot.getPath(), createCryptoSuite(false));
+ changelogDB.initializeDB();
+
+ changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, EXISTING_SERVER_ID, replicationServer);
+ final ConcurrentMap<Integer, FileReplicaDB> domainMap =
+ changelogDB.getDomainToReplicaDBs().get(TEST_ROOT_DN);
+
+ changelogDB.failNextReplicaDBCreation();
+ try
+ {
+ changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, FAILING_SERVER_ID, replicationServer);
+ failBecauseExceptionWasNotThrown(ChangelogException.class);
+ }
+ catch (ChangelogException expected)
+ {
+ assertThat(expected).hasMessage(CREATION_FAILURE.toString());
+ }
+ assertThat(changelogDB.getDomainToReplicaDBs().get(TEST_ROOT_DN))
+ .as("the domain map holding the replica DB created before the failure")
+ .isSameAs(domainMap)
+ .containsOnlyKeys(EXISTING_SERVER_ID);
+ }
+ finally
+ {
+ try
+ {
+ if (changelogDB != null)
+ {
+ changelogDB.shutdownDB();
+ }
+ }
+ finally
+ {
+ remove(replicationServer);
+ TestCaseUtils.deleteDirectory(testRoot);
+ }
+ }
+ }
+
+ /**
+ * The cleanup of a failed creation leaves the domain announced to every multi domain cursor
+ * which was live at the time, and the next successful creation of the domain announces it to
+ * them again: the second announcement must not open a second cursor over the same domain. Such
+ * a cursor would either leak unclosed - the cursor tree of {@code CompositeDBCursor} collapses
+ * cursors comparing equal - or deliver every change twice, which kills the
+ * {@code ChangeNumberIndexer} thread with the {@code IllegalStateException} its cookie update
+ * throws on a replayed change.
+ * <p>
+ * This covers the drop path only: the {@code removeDomain()} path never re-announces a domain
+ * to a cursor which still holds it - the cursor drops the domain, through
+ * {@code indexer.clear()}, before the domain is unmapped.
+ */
+ @Test
+ public void announcingADomainTwiceToALiveCursorMustNotOpenASecondDomainCursor() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ ReplicationServer replicationServer = null;
+ RaceableChangelogDB changelogDB = null;
+ File testRoot = null;
+ try
+ {
+ replicationServer = configureReplicationServer(100, 5000);
+ testRoot = createCleanDir("FileChangelogDB");
+ changelogDB = new RaceableChangelogDB(replicationServer, testRoot.getPath(), createCryptoSuite(false));
+ changelogDB.initializeDB();
+
+ // the live cursor both announcements reach
+ final MultiDomainDBCursor cursor = changelogDB.getCursorFrom(
+ new MultiDomainServerState(), new CursorOptions(GREATER_THAN_OR_EQUAL_TO_KEY, ON_MATCHING_KEY));
+ try
+ {
+ // first announcement: the failed creation announces the domain before dropping its map
+ changelogDB.failNextReplicaDBCreation();
+ try
+ {
+ changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, FAILING_SERVER_ID, replicationServer);
+ failBecauseExceptionWasNotThrown(ChangelogException.class);
+ }
+ catch (ChangelogException expected)
+ {
+ assertThat(expected).hasMessage(CREATION_FAILURE.toString());
+ }
+ cursor.next();
+ assertThat(changelogDB.walkedDomains)
+ .as("the domains the live cursor iterates over after the first announcement")
+ .containsExactly(TEST_ROOT_DN);
+
+ // second announcement: the next creation of the domain announces it to the cursor again
+ final FileReplicaDB replicaDB =
+ changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, FAILING_SERVER_ID, replicationServer).getFirst();
+ final CSN csn = new CSN(System.currentTimeMillis(), 1, FAILING_SERVER_ID);
+ replicaDB.add(new DeleteMsg(TEST_ROOT_DN, csn, "uid"));
+ waitChangesArePersisted(replicaDB, 1);
+
+ assertThat(cursor.next()).as("the change published after the second announcement").isTrue();
+ assertThat(cursor.getRecord().getCSN()).isEqualTo(csn);
+ assertThat(changelogDB.walkedDomains)
+ .as("announcing an already incorporated domain again must not open a second cursor over it")
+ .containsExactly(TEST_ROOT_DN);
+ assertThat(cursor.next()).as("the single published change is delivered more than once").isFalse();
+ }
+ finally
+ {
+ cursor.close();
+ }
+ }
+ finally
+ {
+ try
+ {
+ if (changelogDB != null)
+ {
+ changelogDB.shutdownDB();
+ }
+ }
+ finally
+ {
+ remove(replicationServer);
+ TestCaseUtils.deleteDirectory(testRoot);
+ }
+ }
+ }
+
+ /**
+ * A creation which bails out on the identity check must not drop the domain map another creation
+ * has freshly inserted: the drop is equality based and two empty maps are equal, so an identity
+ * unaware cleanup would unmap the fresh map, and the replica DB about to be published into it
+ * would no longer be reachable from {@code domainToReplicaDBs} - nothing would ever shut it
+ * down, which is the leak of #813 all over again.
+ * <p>
+ * The interleaving is driven step by step:
+ * <ol>
+ * <li>the stale creator obtains its domain map and is held before entering its monitor;</li>
+ * <li>{@code removeDomain()} unmaps that domain map;</li>
+ * <li>a fresh creator inserts a new, still empty domain map and is held inside
+ * {@code newReplicaDB()}, under the fresh map's monitor;</li>
+ * <li>the stale creator is released: its identity check fails and it must bail out without
+ * touching the fresh map, then retry and block on the fresh map's monitor;</li>
+ * <li>the fresh creator is released: both creations complete into that same map.</li>
+ * </ol>
+ */
+ @Test
+ public void bailOutMustNotUnmapAnotherThreadsFreshDomainMap() throws Exception
+ {
+ TestCaseUtils.startServer();
+
+ ReplicationServer replicationServer = null;
+ RaceableChangelogDB changelogDB = null;
+ File testRoot = null;
+ Thread staleCreator = null;
+ Thread freshCreator = null;
+ final AtomicReference<Throwable> staleCreationFailure = new AtomicReference<>();
+ final AtomicReference<Throwable> freshCreationFailure = new AtomicReference<>();
+ try
+ {
+ replicationServer = configureReplicationServer(100, 5000);
+ testRoot = createCleanDir("FileChangelogDB");
+ changelogDB = new RaceableChangelogDB(replicationServer, testRoot.getPath(), createCryptoSuite(false));
+ changelogDB.initializeDB();
+
+ final RaceableChangelogDB racedChangelogDB = changelogDB;
+ final ReplicationServer racedReplicationServer = replicationServer;
+
+ // 1- the stale creator obtains the domain map about to be unmapped, and is parked there
+ changelogDB.holdNextCreationAfterItsDomainMapIsObtained();
+ staleCreator = new Thread("FileChangelogDBTest stale replica DB creator")
+ {
+ @Override
+ public void run()
+ {
+ try
+ {
+ racedChangelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, STALE_SERVER_ID, racedReplicationServer);
+ }
+ catch (Throwable t)
+ {
+ staleCreationFailure.set(t);
+ }
+ }
+ };
+ staleCreator.start();
+ changelogDB.awaitCreatorHoldingItsDomainMap();
+
+ // 2- the domain map the stale creator holds is unmapped
+ changelogDB.removeDomain(TEST_ROOT_DN);
+
+ // 3- the fresh creator inserts a new, still empty domain map, and is parked inside
+ // newReplicaDB(), under the monitor of that fresh map
+ changelogDB.holdNextReplicaDBOnceCreated();
+ freshCreator = new Thread("FileChangelogDBTest fresh replica DB creator")
+ {
+ @Override
+ public void run()
+ {
+ try
+ {
+ racedChangelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, FRESH_SERVER_ID, racedReplicationServer);
+ }
+ catch (Throwable t)
+ {
+ freshCreationFailure.set(t);
+ }
+ }
+ };
+ freshCreator.start();
+ changelogDB.awaitCreatorHoldingItsCreatedReplicaDB();
+ final ConcurrentMap<Integer, FileReplicaDB> freshDomainMap =
+ changelogDB.getDomainToReplicaDBs().get(TEST_ROOT_DN);
+ assertThat(freshDomainMap).as("the fresh domain map, not published into yet").isNotNull().isEmpty();
+
+ // 4- the stale creator bails out on its identity check, retries, and blocks on the monitor
+ // of the fresh domain map - without the identity check its cleanup would have unmapped the
+ // fresh map, and it would have completed into a third map instead of blocking
+ changelogDB.releaseCreatorHoldingItsDomainMap();
+ waitUntilBlockedOnOrCompleted(staleCreator, freshDomainMap);
+ assertThat(changelogDB.getDomainToReplicaDBs().get(TEST_ROOT_DN))
+ .as("the domain map the fresh creator is about to publish its replica DB into")
+ .isSameAs(freshDomainMap);
+
+ // 5- both creations complete into that same map
+ changelogDB.releaseCreatedReplicaDB();
+ staleCreator.join(TIMEOUT_MS);
+ freshCreator.join(TIMEOUT_MS);
+ assertThat(staleCreator.isAlive()).as("the stale creator thread did not complete").isFalse();
+ assertThat(freshCreator.isAlive()).as("the fresh creator thread did not complete").isFalse();
+ assertThat(staleCreationFailure.get()).isNull();
+ assertThat(freshCreationFailure.get()).isNull();
+ assertThat(changelogDB.getDomainToReplicaDBs().get(TEST_ROOT_DN))
+ .isSameAs(freshDomainMap)
+ .containsOnlyKeys(STALE_SERVER_ID, FRESH_SERVER_ID);
+ }
+ finally
+ {
+ try
+ {
+ if (changelogDB != null)
+ {
+ changelogDB.releaseAllHeldThreads();
+ changelogDB.shutdownDB();
+ }
+ }
+ finally
+ {
+ join(staleCreator);
+ join(freshCreator);
+ deregisterLeakedReplicaDBMonitors(replicationServer, STALE_SERVER_ID, FRESH_SERVER_ID);
+ 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}.
@@ -329,7 +700,7 @@
changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, SERVER_ID, replicationServer).getFirst();
final ConcurrentMap<DN, ConcurrentMap<Integer, FileReplicaDB>> domainToReplicaDBs =
- getDomainToReplicaDBs(changelogDB);
+ changelogDB.getDomainToReplicaDBs();
final ConcurrentMap<Integer, FileReplicaDB> domainMap = domainToReplicaDBs.get(TEST_ROOT_DN);
assertThat(domainMap).isNotNull();
@@ -373,7 +744,7 @@
changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, SERVER_ID, replicationServer).getFirst();
final ConcurrentMap<DN, ConcurrentMap<Integer, FileReplicaDB>> domainToReplicaDBs =
- getDomainToReplicaDBs(changelogDB);
+ changelogDB.getDomainToReplicaDBs();
final ConcurrentMap<Integer, FileReplicaDB> domainMap = domainToReplicaDBs.get(TEST_ROOT_DN);
assertThat(domainMap).isNotNull();
@@ -420,27 +791,27 @@
}, "removeDomain() under test");
}
- @SuppressWarnings("unchecked")
- private ConcurrentMap<DN, ConcurrentMap<Integer, FileReplicaDB>> getDomainToReplicaDBs(
- FileChangelogDB changelogDB) throws Exception
+ /** Waits until the provided replica DB has persisted the provided number of records. */
+ private void waitChangesArePersisted(FileReplicaDB replicaDB, int recordCount) throws Exception
{
- final Field field = FileChangelogDB.class.getDeclaredField("domainToReplicaDBs");
- field.setAccessible(true);
- return (ConcurrentMap<DN, ConcurrentMap<Integer, FileReplicaDB>>) field.get(changelogDB);
+ final long deadline = System.currentTimeMillis() + TIMEOUT_MS;
+ while (replicaDB.getNumberRecords() < recordCount)
+ {
+ if (System.currentTimeMillis() > deadline)
+ {
+ throw new AssertionError("Timed out waiting for " + recordCount + " records to be persisted");
+ }
+ Thread.sleep(10);
+ }
}
/** 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() + TIMEOUT_MS;
while (System.currentTimeMillis() < deadline)
{
- final ThreadInfo threadInfo = threadMXBean.getThreadInfo(thread.getId());
- final LockInfo lockInfo = threadInfo != null ? threadInfo.getLockInfo() : null;
- if (lockInfo != null
- && threadInfo.getThreadState() == Thread.State.BLOCKED
- && lockInfo.getIdentityHashCode() == System.identityHashCode(monitor))
+ if (isBlockedOn(thread, monitor))
{
return;
}
@@ -450,6 +821,35 @@
"Timed out waiting for " + thread.getName() + " to block on the domainMap monitor");
}
+ /**
+ * Waits until the provided thread is blocked acquiring the monitor of the provided object, or
+ * has completed: completion is left for the caller's assertions to diagnose.
+ */
+ private void waitUntilBlockedOnOrCompleted(Thread thread, Object monitor) throws Exception
+ {
+ final long deadline = System.currentTimeMillis() + TIMEOUT_MS;
+ while (System.currentTimeMillis() < deadline)
+ {
+ if (!thread.isAlive() || isBlockedOn(thread, monitor))
+ {
+ return;
+ }
+ Thread.sleep(1);
+ }
+ throw new AssertionError("Timed out waiting for " + thread.getName()
+ + " to block on the domainMap monitor or complete");
+ }
+
+ private static boolean isBlockedOn(Thread thread, Object monitor)
+ {
+ final ThreadMXBean threadMXBean = ManagementFactory.getThreadMXBean();
+ final ThreadInfo threadInfo = threadMXBean.getThreadInfo(thread.getId());
+ final LockInfo lockInfo = threadInfo != null ? threadInfo.getLockInfo() : null;
+ return lockInfo != null
+ && threadInfo.getThreadState() == Thread.State.BLOCKED
+ && lockInfo.getIdentityHashCode() == System.identityHashCode(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
{
@@ -468,31 +868,6 @@
}
/**
- * 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.
@@ -505,13 +880,13 @@
}
/** Releases the monitor providers a regression leaks, so that they do not outlive this test. */
- private void deregisterLeakedReplicaDBMonitors(final ReplicationServer replicationServer)
+ private void deregisterLeakedReplicaDBMonitors(final ReplicationServer replicationServer, final int... serverIds)
{
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 })
+ for (final int serverId : serverIds)
{
// deregister the provider instead of removing the map entry, so that the JMX MBean
// registered alongside it is released as well
@@ -526,21 +901,30 @@
/**
* 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}.
+ * the shutdown flag, hold it after it has obtained its domain map but before it enters the
+ * monitor, hold it again once the replica DB is created but not yet published into the domain
+ * map, hold the shutdown inside the drain of {@code domainToReplicaDBs}, make the next replica
+ * DB creation fail, and record the domains cursors are opened for.
*/
private static final class RaceableChangelogDB extends FileChangelogDB
{
private final AtomicBoolean holdNextCreation = new AtomicBoolean();
+ private final AtomicBoolean holdNextDomainMapObtained = new AtomicBoolean();
private final AtomicBoolean holdNextReplicaDBShutdown = new AtomicBoolean();
private final AtomicBoolean holdNextCreatedReplicaDB = new AtomicBoolean();
+ private final AtomicBoolean failNextCreation = new AtomicBoolean();
private final CountDownLatch creatorIsInWindow = new CountDownLatch(1);
private final CountDownLatch creatorIsReleased = new CountDownLatch(1);
+ private final CountDownLatch creatorHoldsItsDomainMap = new CountDownLatch(1);
+ private final CountDownLatch domainMapIsReleased = 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);
+ /** The baseDNs of the domains any cursor was opened for, one element per opening. */
+ private final List<DN> walkedDomains = new CopyOnWriteArrayList<>();
+
RaceableChangelogDB(final ReplicationServer replicationServer, final String dbDirectoryPath,
final CryptoSuite cryptoSuite) throws ConfigException
{
@@ -555,13 +939,23 @@
creatorIsInWindow.countDown();
await(creatorIsReleased);
}
- return super.getExistingOrNewDomainMap(baseDN);
+ final ConcurrentMap<Integer, FileReplicaDB> domainMap = super.getExistingOrNewDomainMap(baseDN);
+ if (holdNextDomainMapObtained.compareAndSet(true, false))
+ {
+ creatorHoldsItsDomainMap.countDown();
+ await(domainMapIsReleased);
+ }
+ return domainMap;
}
@Override
FileReplicaDB newReplicaDB(final int serverId, final DN baseDN, final ReplicationServer server,
final CryptoSuite cryptoSuite, final ReplicationEnvironment replicationEnv) throws ChangelogException
{
+ if (failNextCreation.compareAndSet(true, false))
+ {
+ throw new ChangelogException(CREATION_FAILURE);
+ }
if (holdNextReplicaDBShutdown.compareAndSet(true, false))
{
return new HeldOnShutdownReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv);
@@ -577,11 +971,24 @@
return replicaDB;
}
+ @Override
+ public DBCursor<UpdateMsg> getCursorFrom(final DN baseDN, final ServerState startState,
+ final CursorOptions options) throws ChangelogException
+ {
+ walkedDomains.add(baseDN);
+ return super.getCursorFrom(baseDN, startState, options);
+ }
+
void holdNextReplicaDBCreationBeforeItsDomainMapIsInserted()
{
holdNextCreation.set(true);
}
+ void holdNextCreationAfterItsDomainMapIsObtained()
+ {
+ holdNextDomainMapObtained.set(true);
+ }
+
void holdNextReplicaDBInItsShutdown()
{
holdNextReplicaDBShutdown.set(true);
@@ -592,11 +999,21 @@
holdNextCreatedReplicaDB.set(true);
}
+ void failNextReplicaDBCreation()
+ {
+ failNextCreation.set(true);
+ }
+
void awaitCreatorInWindow()
{
await(creatorIsInWindow);
}
+ void awaitCreatorHoldingItsDomainMap()
+ {
+ await(creatorHoldsItsDomainMap);
+ }
+
void awaitCreatorHoldingItsCreatedReplicaDB()
{
await(creatorHoldsItsCreatedReplicaDB);
@@ -612,6 +1029,11 @@
creatorIsReleased.countDown();
}
+ void releaseCreatorHoldingItsDomainMap()
+ {
+ domainMapIsReleased.countDown();
+ }
+
void releaseCreatedReplicaDB()
{
createdReplicaDBIsReleased.countDown();
@@ -625,6 +1047,7 @@
void releaseAllHeldThreads()
{
releaseCreator();
+ releaseCreatorHoldingItsDomainMap();
releaseCreatedReplicaDB();
releaseDrain();
}
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/SequentialDBCursor.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/SequentialDBCursor.java
index 333e73b..401ae13 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/SequentialDBCursor.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/server/changelog/file/SequentialDBCursor.java
@@ -12,9 +12,11 @@
* information: "Portions Copyright [year] [name of copyright owner]".
*
* Copyright 2013-2016 ForgeRock AS.
+ * Portions Copyright 2026 3A Systems, LLC.
*/
package org.opends.server.replication.server.changelog.file;
+import java.util.ArrayList;
import java.util.List;
import org.opends.server.replication.protocol.UpdateMsg;
@@ -44,6 +46,14 @@
this.msgs.add(msg);
}
+ /** Returns the messages this cursor has not consumed yet, leaving it empty. */
+ public List<UpdateMsg> drain()
+ {
+ final List<UpdateMsg> drained = new ArrayList<>(msgs);
+ msgs.clear();
+ return drained;
+ }
+
@Override
public UpdateMsg getRecord()
{
--
Gitblit v1.10.0