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