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 6561be7146..8a914728b7 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 @@ public class FileChangelogDB implements ChangelogDB, ReplicationDomainDB * * When creating a replicaDB, synchronize on the domainMap to avoid * concurrent shutdown. + *

+ * 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> domainToReplicaDBs = new ConcurrentHashMap<>(); @@ -278,25 +284,49 @@ private Pair getExistingOrNewReplicaDB(final ConcurrentM // 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 @@ FileReplicaDB newReplicaDB(final int serverId, final DN baseDN, final Replicatio return new FileReplicaDB(serverId, baseDN, server, cryptoSuite, replicationEnv); } + /** + * Returns the map of replica DBs per domain. + *

+ * 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> getDomainToReplicaDBs() + { + return domainToReplicaDBs; + } + @Override public void initializeDB() throws ChangelogException { @@ -403,13 +447,23 @@ public void shutdownDB() throws ChangelogException firstException = e; } - for (Iterator> it = - this.domainToReplicaDBs.values().iterator(); it.hasNext();) + for (Iterator>> it = + this.domainToReplicaDBs.entrySet().iterator(); it.hasNext();) { - final ConcurrentMap domainMap = it.next(); + final Map.Entry> entry = it.next(); + final ConcurrentMap 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 8a9ec2d221..4019bd762c 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 @@ public class MultiDomainDBCursor extends CompositeDBCursor { private final ReplicationDomainDB domainDB; private final ConcurrentSkipListMap 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. + *

+ * 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 incorporatedDomains = new ConcurrentSkipListSet<>(); private final CursorOptions options; /** @@ -52,6 +66,10 @@ public MultiDomainDBCursor(final ReplicationDomainDB domainDB, CursorOptions opt /** * 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()}. + *

+ * 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 MultiDomainDBCursor(final ReplicationDomainDB domainDB, CursorOptions opt */ 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 @@ protected void incorporateNewCursors() throws ChangelogException final Entry entry = iter.next(); final DN baseDN = entry.getKey(); final ServerState serverState = entry.getValue(); - final DBCursor 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 domainDBCursor = domainDB.getCursorFrom(baseDN, serverState, options); + addCursor(domainDBCursor, baseDN); + incorporatedDomains.add(baseDN); + } iter.remove(); } } @@ -90,6 +117,7 @@ protected void incorporateNewCursors() throws ChangelogException public void removeDomain(DN baseDN) { removeCursor(baseDN); + incorporatedDomains.remove(baseDN); } /** {@inheritDoc} */ @@ -99,6 +127,7 @@ public void close() 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 6ef8d2f3ab..e6e3ffd4fd 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 @@ public class ECLMultiDomainDBCursorTest extends DirectoryServerTestCase private MultiDomainDBCursor multiDomainCursor; private ECLMultiDomainDBCursor eclCursor; private final Set eclEnabledDomains = new HashSet<>(); + /** The long-lived cursor of each domain announced to {@link #multiDomainCursor}. */ + private final Map domainCursors = new HashMap<>(); private ECLEnabledDomainPredicate predicate = new ECLEnabledDomainPredicate() { @Override @@ -61,6 +66,9 @@ public void setup() throws Exception 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 assertMessagesInOrder(DN baseDN, UpdateMsg msg1, UpdateMsg msg2) th 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 595338a0d2..31bfa61765 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 @@ public class FileChangelogDBTest extends ReplicationTestCase 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 @@ public void run() } join(creator); join(shutdowner); - deregisterLeakedReplicaDBMonitors(replicationServer); + deregisterLeakedReplicaDBMonitors(replicationServer, RACING_SERVER_ID, DRAINED_SERVER_ID); remove(replicationServer); TestCaseUtils.deleteDirectory(testRoot); } @@ -282,7 +308,7 @@ public void run() } }; shutdowner.start(); - awaitBlockedOnAMonitor(shutdowner); + waitUntilBlockedOn(shutdowner, changelogDB.getDomainToReplicaDBs().get(TEST_ROOT_DN)); changelogDB.releaseCreatedReplicaDB(); creator.join(TIMEOUT_MS); @@ -305,12 +331,357 @@ public void run() } 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 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 empty 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 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. + *

+ * 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. + *

+ * The interleaving is driven step by step: + *

    + *
  1. the stale creator obtains its domain map and is held before entering its monitor;
  2. + *
  3. {@code removeDomain()} unmaps that domain map;
  4. + *
  5. a fresh creator inserts a new, still empty domain map and is held inside + * {@code newReplicaDB()}, under the fresh map's monitor;
  6. + *
  7. 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;
  8. + *
  9. the fresh creator is released: both creations complete into that same map.
  10. + *
+ */ + @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 staleCreationFailure = new AtomicReference<>(); + final AtomicReference 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 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 @@ -329,7 +700,7 @@ public void removeDomainRacingConcurrentRemovalMustNotThrowNPE() throws Exceptio changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, SERVER_ID, replicationServer).getFirst(); final ConcurrentMap> domainToReplicaDBs = - getDomainToReplicaDBs(changelogDB); + changelogDB.getDomainToReplicaDBs(); final ConcurrentMap domainMap = domainToReplicaDBs.get(TEST_ROOT_DN); assertThat(domainMap).isNotNull(); @@ -373,7 +744,7 @@ public void removeDomainMustNotUnmapConcurrentlyRecreatedDomain() throws Excepti changelogDB.getOrCreateReplicaDB(TEST_ROOT_DN, SERVER_ID, replicationServer).getFirst(); final ConcurrentMap> domainToReplicaDBs = - getDomainToReplicaDBs(changelogDB); + changelogDB.getDomainToReplicaDBs(); final ConcurrentMap domainMap = domainToReplicaDBs.get(TEST_ROOT_DN); assertThat(domainMap).isNotNull(); @@ -420,27 +791,27 @@ public void run() }, "removeDomain() under test"); } - @SuppressWarnings("unchecked") - private ConcurrentMap> 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>) 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 @@ private void waitUntilBlockedOn(Thread thread, Object monitor) throws Exception "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 { @@ -467,31 +867,6 @@ private void join(final Thread thread) throws InterruptedException } } - /** - * 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 @@ -505,13 +880,13 @@ private String replicaDBMonitorName(final ReplicationServer replicationServer, f } /** 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 @@ private void deregisterLeakedReplicaDBMonitors(final ReplicationServer replicati /** * 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 walkedDomains = new CopyOnWriteArrayList<>(); + RaceableChangelogDB(final ReplicationServer replicationServer, final String dbDirectoryPath, final CryptoSuite cryptoSuite) throws ConfigException { @@ -555,13 +939,23 @@ ConcurrentMap getExistingOrNewDomainMap(final DN baseDN) creatorIsInWindow.countDown(); await(creatorIsReleased); } - return super.getExistingOrNewDomainMap(baseDN); + final ConcurrentMap 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 @@ FileReplicaDB newReplicaDB(final int serverId, final DN baseDN, final Replicatio return replicaDB; } + @Override + public DBCursor 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 @@ void holdNextReplicaDBOnceCreated() holdNextCreatedReplicaDB.set(true); } + void failNextReplicaDBCreation() + { + failNextCreation.set(true); + } + void awaitCreatorInWindow() { await(creatorIsInWindow); } + void awaitCreatorHoldingItsDomainMap() + { + await(creatorHoldsItsDomainMap); + } + void awaitCreatorHoldingItsCreatedReplicaDB() { await(creatorHoldsItsCreatedReplicaDB); @@ -612,6 +1029,11 @@ void releaseCreator() creatorIsReleased.countDown(); } + void releaseCreatorHoldingItsDomainMap() + { + domainMapIsReleased.countDown(); + } + void releaseCreatedReplicaDB() { createdReplicaDBIsReleased.countDown(); @@ -625,6 +1047,7 @@ void releaseDrain() 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 333e73bd4a..401ae1313f 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 @@ public void add(UpdateMsg msg) this.msgs.add(msg); } + /** Returns the messages this cursor has not consumed yet, leaving it empty. */ + public List drain() + { + final List drained = new ArrayList<>(msgs); + msgs.clear(); + return drained; + } + @Override public UpdateMsg getRecord() {