From 776339a8c63c0bd8c83cfe8a619f6146794358e8 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Sat, 12 Sep 2026 11:53:58 +0000
Subject: [PATCH] [#922] Give a change back when the replay which owns it is unwound (#958)

---
 opendj-server-legacy/src/test/java/org/opends/server/replication/UpdateOperationTest.java |  977 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
 1 files changed, 957 insertions(+), 20 deletions(-)

diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/UpdateOperationTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/UpdateOperationTest.java
index b2d03a7..b791c2a 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/UpdateOperationTest.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/UpdateOperationTest.java
@@ -33,9 +33,12 @@
 import java.net.SocketTimeoutException;
 import java.util.ArrayList;
 import java.util.Collections;
+import java.util.HashSet;
 import java.util.List;
+import java.util.Set;
 import java.util.concurrent.TimeoutException;
 import java.util.concurrent.atomic.AtomicReference;
+import java.util.function.Supplier;
 
 import org.assertj.core.api.Assertions;
 import org.forgerock.i18n.LocalizableMessage;
@@ -58,6 +61,7 @@
 import org.opends.server.plugins.PausePreParsePlugin;
 import org.opends.server.plugins.ShortCircuitPlugin;
 import org.opends.server.plugins.ShortCircuitPlugin.ParkedReplay;
+import org.opends.server.plugins.ShortCircuitPlugin.ThrownFromReplay;
 import org.opends.server.protocols.internal.InternalClientConnection;
 import org.opends.server.replication.common.AssuredMode;
 import org.opends.server.replication.common.CSN;
@@ -2584,13 +2588,738 @@
     logger.error(LocalizableMessage.raw(
         "Starting replication test : aChangeWhoseOperationWasBuiltIsNotGivenUpOnWhereItFailed"));
 
+    assertChangeIsDeliveredAgainAfter(ModifyMsgWhoseOperationRefusesAControl::new,
+        18, "user.889.7", "the replay must fail once the operation is built");
+  }
+
+  /**
+   * The errors a replay meets which say nothing about the changes which follow: a class
+   * which could not be linked, and a stack which ran out on the entry being replayed. The
+   * second is an error of the JVM, but the thread which met it is whole again once the
+   * stack has unwound and the entry is what raised it, so ending the thread would have one
+   * change this replica can not replay cost it one replay thread per delivery.
+   * <p>
+   * Each row builds its error where it is thrown rather than here, so that the stack trace
+   * it carries is the one of the replay it unwound.
+   */
+  @DataProvider(name = "recoverableReplayErrors")
+  public Object[][] recoverableReplayErrors()
+  {
+    return new Object[][] {
+      { (Supplier<Error>) () -> new LinkageError("the replay of this change meets an Error"),
+        19, "user.922.1", "a LinkageError" },
+      { (Supplier<Error>) StackOverflowError::new, 23, "user.922.5", "a StackOverflowError" },
+    };
+  }
+
+  /**
+   * Test case for [Issue 922] and [Issue 923]: a change whose replay threw an Error is
+   * given back, and the thread which was replaying it is still there to take the changes
+   * which follow.
+   * <p>
+   * A change is owned by the replay thread which took it, and that ownership is what
+   * keeps a change being replayed from being replayed a second time (OPENDJ-1115): every
+   * later delivery of it is refused as a duplicate. Ownership is given back on the roads
+   * which run to their end, so an Error - which unwinds the replay out of every one of
+   * them - would leave the change listed, uncommitted and owned by a thread which is not
+   * replaying it anymore: nobody could replay it, and this domain's ServerState would
+   * never move past it again.
+   * <p>
+   * The error is thrown where the operation is built, so what these rows pin is the arm
+   * which reports it and takes the road of a failed replay: it is caught inside the replay
+   * and never leaves it. The roads which do leave it - the give-back of a replay which was
+   * unwound, and the widened catch of the replay thread which is what keeps that thread
+   * alive - are pinned by
+   * {@link #aChangeWhoseReplayIsUnwoundAfterItsAckIsDeliveredAgain()}, where the throw is
+   * made past the point any catch of the replay runs on.
+   */
+  @Test(dataProvider = "recoverableReplayErrors")
+  public void aChangeWhoseReplayThrewAnErrorIsDeliveredAgain(
+      final Supplier<Error> error, int serverId, String uid, String what) throws Exception
+  {
+    testSetUp("aChangeWhoseReplayThrewAnErrorIsDeliveredAgain." + uid);
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : aChangeWhoseReplayThrewAnErrorIsDeliveredAgain " + uid));
+
+    assertChangeIsDeliveredAgainAfter(
+        (csn, dn, mods, entryUUID) ->
+            new ModifyMsgWhoseReplayThrows(csn, dn, mods, entryUUID, error),
+        serverId, uid, "the replay of this change throws " + what);
+  }
+
+  /**
+   * Test case for [Issue 922] and [Issue 923]: the change a replay thread was replaying
+   * when the JVM ran out of memory is given back before the error is left to end that
+   * thread.
+   * <p>
+   * An OutOfMemoryError is not turned into a failed replay and reported the way the other
+   * errors are - building the report asks for more of what the JVM has run out of - and
+   * the thread it unwinds is not replaced. The change it was replaying must not go with
+   * it: it is handed back so that the delivery which follows can replay it.
+   */
+  @Test
+  public void aChangeWhoseReplayRanOutOfMemoryIsDeliveredAgain() throws Exception
+  {
+    testSetUp("aChangeWhoseReplayRanOutOfMemoryIsDeliveredAgain");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : aChangeWhoseReplayRanOutOfMemoryIsDeliveredAgain"));
+
+    /*
+     * The threads are read by identity rather than counted: what this test is about is the
+     * thread which met the error being gone, which is what #923 sanctions. Whether the pool
+     * is refilled afterwards is the half of that issue which is left open, and a count
+     * would freeze it here as the behaviour which is wanted.
+     */
+    final Set<Long> replayThreadsBefore = replayThreadIds();
+    final int initialUncaughtAlerts = DummyAlertHandler.getAlertCount(ALERT_TYPE_UNCAUGHT_EXCEPTION);
+    assertChangeIsDeliveredAgainAfter(
+        (csn, dn, mods, entryUUID) -> new ModifyMsgWhoseReplayThrows(csn, dn, mods, entryUUID,
+            () -> new OutOfMemoryError("the JVM is out of memory")),
+        20, "user.922.2", "the replay of this change meets an OutOfMemoryError");
+    /*
+     * The change is given back, and so replayed by another thread, while the thread which
+     * met the error is still unwinding - through the session restart, and then through the
+     * uncaught exception handler which raises the alert - so its end is waited for rather
+     * than read off the pool the moment the change lands. The alert is raised by that
+     * handler once run() has returned, so it lands after the thread is gone, and it is
+     * waited for in the same breath.
+     */
+    TestTimer timer = new TestTimer.Builder()
+      .maxSleep(30, SECONDS)
+      .sleepTimes(200, MILLISECONDS)
+      .toTimer();
+    timer.repeatUntilSuccess(new CallableVoid()
+    {
+      @Override
+      public void call() throws Exception
+      {
+        assertFalse(replayThreadIds().containsAll(replayThreadsBefore),
+            "an OutOfMemoryError must end the replay thread which met it");
+        assertUncaughtExceptionAlertRaisedSince(initialUncaughtAlerts,
+            "the replay thread an OutOfMemoryError ended must have raised the alert #923 asks"
+                + " for on its way out");
+      }
+    });
+  }
+
+  /**
+   * Asserts that a {@code DirectoryThread} has ended on an uncaught throwable since the
+   * provided count of those alerts was read: the uncaught exception handler of its thread
+   * group is what raises it, and it is the one line a replay thread which an
+   * OutOfMemoryError ended leaves behind (issue #923).
+   * <p>
+   * At least one rather than exactly one: the alert is raised for every thread of the
+   * server which ends that way, and a thread of some other component ending during the
+   * test must not turn this into a failure of the wrong test.
+   */
+  private static void assertUncaughtExceptionAlertRaisedSince(int initialAlerts, String message)
+  {
+    Assertions.assertThat(DummyAlertHandler.getAlertCount(ALERT_TYPE_UNCAUGHT_EXCEPTION))
+        .as(message)
+        .isGreaterThanOrEqualTo(initialAlerts + 1);
+  }
+
+  /**
+   * Test case for [Issue 922]: a change whose ack could not be published still takes the
+   * road its own replay decided.
+   * <p>
+   * The ack of a delivery is published in a finally which every road out of the replay runs
+   * through, and it is published on the session that delivery came over - which is being
+   * torn down when the replay of the change failed. Whether the change was applied is not
+   * something that publish can tell, so a throw there is reported and the replay carries on
+   * to the road the change itself decided: this one failed, so it is kept out of the
+   * ServerState, counted and asked for again.
+   * <p>
+   * Left to unwind, that throw would step over the give-back which follows it and leave the
+   * change owned by a thread which is not replaying it anymore.
+   */
+  @Test
+  public void aChangeWhoseAckCouldNotBePublishedIsDeliveredAgain() throws Exception
+  {
+    testSetUp("aChangeWhoseAckCouldNotBePublishedIsDeliveredAgain");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : aChangeWhoseAckCouldNotBePublishedIsDeliveredAgain"));
+
+    assertChangeIsDeliveredAgainAfter(ModifyMsgWhoseAckThrows::new,
+        21, "user.922.3", "the ack of this change throws on the way out of the replay");
+  }
+
+  /**
+   * Test case for [Issue 922] and [Issue 923]: an OutOfMemoryError met where the ack of a
+   * delivery is published ends the replay thread, the way one met by the replay itself does.
+   * <p>
+   * A throw from the ack is caught so that it does not unwind the replay past the give-back
+   * of the change and past the hand-out of the changes which were waiting for it. A JVM
+   * which has run out of memory is the one exception to that: it is not something to carry
+   * on replaying from, so it is left to end this thread - the change is given back on the
+   * way out, and the uncaught exception handler of DirectoryThread writes the line and
+   * raises the alert #923 is about. Caught like every other throw from there, it would have
+   * this thread replay the changes which follow on an exhausted heap, with nothing said
+   * anywhere.
+   */
+  @Test
+  public void aChangeWhoseAckRanOutOfMemoryIsDeliveredAgain() throws Exception
+  {
+    testSetUp("aChangeWhoseAckRanOutOfMemoryIsDeliveredAgain");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : aChangeWhoseAckRanOutOfMemoryIsDeliveredAgain"));
+
+    final Set<Long> replayThreadsBefore = replayThreadIds();
+    final int initialUncaughtAlerts = DummyAlertHandler.getAlertCount(ALERT_TYPE_UNCAUGHT_EXCEPTION);
+    assertChangeIsDeliveredAgainAfter(ModifyMsgWhoseAckRunsOutOfMemory::new,
+        26, "user.922.10", "the ack of this change runs out of memory");
+    /*
+     * Waited for rather than read the moment the change lands: the change is given back -
+     * and so replayed by another thread - while the thread which met the error is still
+     * unwinding, through the session restart and then through the handler which raises the
+     * alert. The alert is what the rethrow is for, and the handler raises it once run() has
+     * returned, so it lands after the thread is gone and is waited for in the same breath.
+     */
+    TestTimer timer = new TestTimer.Builder()
+      .maxSleep(30, SECONDS)
+      .sleepTimes(200, MILLISECONDS)
+      .toTimer();
+    timer.repeatUntilSuccess(new CallableVoid()
+    {
+      @Override
+      public void call() throws Exception
+      {
+        assertFalse(replayThreadIds().containsAll(replayThreadsBefore),
+            "an OutOfMemoryError met where the ack is published must end the replay thread"
+                + " which met it");
+        assertUncaughtExceptionAlertRaisedSince(initialUncaughtAlerts,
+            "the replay thread an OutOfMemoryError met where the ack is published ended must"
+                + " have raised the alert #923 asks for on its way out");
+      }
+    });
+  }
+
+  /**
+   * Test case for [Issue 922] and [Issue 923]: a change whose replay is unwound once the
+   * ack of its delivery is out is given back, and the thread which was replaying it stays.
+   * <p>
+   * The catches of the replay itself span the roads which decide what became of the change,
+   * and the ack is published once that is decided. What the replay runs afterwards - the
+   * give-back of a change which failed, and the hand-out of the changes which were waiting
+   * for it - is past every one of them: a throw there unwinds {@code replay()} with the
+   * change still owned by this thread, and a change owned by a thread which is not replaying
+   * it anymore is refused as a duplicate on every later delivery. So it is given back on the
+   * way out, and the error is left to the replay thread, whose catch is what keeps it alive
+   * for the changes which follow (issue #923).
+   */
+  @Test
+  public void aChangeWhoseReplayIsUnwoundAfterItsAckIsDeliveredAgain() throws Exception
+  {
+    testSetUp("aChangeWhoseReplayIsUnwoundAfterItsAckIsDeliveredAgain");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : aChangeWhoseReplayIsUnwoundAfterItsAckIsDeliveredAgain"));
+
+    final Set<Long> replayThreadsBefore = replayThreadIds();
+    assertChangeIsDeliveredAgainAfter(ModifyMsgWhoseReplayIsUnwoundAfterItsAck::new,
+        24, "user.922.6", "the replay of this change is unwound once its ack is out");
+    Assertions.assertThat(replayThreadIds())
+        .as("an Error which unwinds a replay must not end the thread which met it (issue #923)")
+        .containsAll(replayThreadsBefore);
+  }
+
+  /**
+   * Test case for [Issue 922]: a change whose replay keeps being unwound after its ack is
+   * given up on rather than asked for forever.
+   * <p>
+   * The change is given back on the road a failed replay takes rather than handed back bare,
+   * so the failure counts against the give-up budget of the change: a change which keeps
+   * unwinding the replays it is given to is eventually recorded as one this replica could
+   * not apply, and the administrator is told that it now diverges. Handed back bare it would
+   * be asked for, and have this domain restart its session for it, for as long as the server
+   * is up - which is the wedge this issue is about wearing another face.
+   */
+  @Test
+  public void aChangeWhoseReplayKeepsBeingUnwoundAfterItsAckIsGivenUpOn() throws Exception
+  {
+    testSetUp("aChangeWhoseReplayKeepsBeingUnwoundAfterItsAckIsGivenUpOn");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : aChangeWhoseReplayKeepsBeingUnwoundAfterItsAckIsGivenUpOn"));
+
     Entry tmp = TestCaseUtils.addEntry(
-        "dn: uid=user.889.7," + baseDN,
+        "dn: uid=user.922.7," + baseDN,
         "objectClass: top",
         "objectClass: person",
         "objectClass: organizationalPerson",
         "objectClass: inetOrgPerson",
-        "uid: user.889.7",
+        "uid: user.922.7",
+        "cn: Aaccf Amar",
+        "sn: Amar");
+    final DN dn = tmp.getName();
+    final String uuid = getEntry(dn, 1, true).parseAttribute("entryuuid").asString();
+
+    final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null);
+    final long initialFailures = getMonitorAttrValue(baseDN, "replayed-updates-failed");
+    domain.resetUnreplayedChangeAlertThrottle();
+    final int initialAlerts = DummyAlertHandler.getAlertCount(ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE);
+    setReplayGiveUpDelay(TEST_GIVE_UP_DELAY);
+    try
+    {
+      final CSN csn = new CSNGenerator(25, TimeThread.getTime()).newCSN();
+      final List<Modification> mods =
+          generatemods("description", "the replay of this change is unwound after its ack");
+
+      /*
+       * Nothing sends this change again - it never travelled a session - so every delivery
+       * of it is made here, until this replica gives up on the change whose replay it can
+       * not run to its end.
+       */
+      TestTimer timer = new TestTimer.Builder()
+        .maxSleep(120, SECONDS)
+        .sleepTimes(200, MILLISECONDS)
+        .toTimer();
+      timer.repeatUntilSuccess(new CallableVoid()
+      {
+        @Override
+        public void call() throws Exception
+        {
+          if (!domain.getServerState().cover(csn))
+          {
+            domain.processUpdate(new ModifyMsgWhoseReplayIsUnwoundAfterItsAck(csn, dn, mods, uuid));
+          }
+          assertTrue(domain.getServerState().cover(csn),
+              "a change whose replay keeps being unwound must be given up on");
+        }
+      });
+      assertMonitorAttrValueEventually(baseDN, "replayed-updates-failed", initialFailures + 1,
+          "the change which was given up on must be counted as failed, once");
+      Assertions.assertThat(DummyAlertHandler.getAlertCount(ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE))
+          .as("the administrator must be told that this replica now diverges")
+          .isGreaterThan(initialAlerts);
+    }
+    finally
+    {
+      resetReplayGiveUpDelay();
+    }
+  }
+
+  /**
+   * Test case for [Issue 922]: the changes parked behind a change whose ack could not be
+   * published are replayed rather than left waiting for a thread which is gone.
+   * <p>
+   * A change which is waiting for another one is handed out by {@code getNextUpdate()},
+   * which the replay runs once it is done with the change it was given - after the ack of
+   * that delivery has been published. A throw from that ack used to unwind the replay past
+   * it, and the change had been committed by then, so nothing was owed back and nothing was
+   * handed out: the changes parked behind it stayed parked, and the ServerState of this
+   * domain stayed behind them until some other change was replayed here.
+   * <p>
+   * What tells the two apart is which thread replays the parked change. It is handed to
+   * whichever thread cleared the change it was waiting for, so it is replayed by the very
+   * thread which has just committed the change whose ack threw - a change nobody handed out
+   * is replayed by no one at all, and the wait below is what says so.
+   */
+  @Test
+  public void theChangesParkedBehindAChangeWhoseAckFailedAreReplayed() throws Exception
+  {
+    testSetUp("theChangesParkedBehindAChangeWhoseAckFailedAreReplayed");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : theChangesParkedBehindAChangeWhoseAckFailedAreReplayed"));
+
+    final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null);
+    final CSNGenerator gen = new CSNGenerator(26, TimeThread.getTime());
+    final String parentUUID = "26262626-2626-2626-2626-262626262626";
+    final String childUUID = "27272727-2727-2727-2727-272727272727";
+    final Entry parent = TestCaseUtils.makeEntry(
+        "dn: ou=parked.922," + baseDN,
+        "objectClass: top",
+        "objectClass: organizationalUnit",
+        "ou: parked.922",
+        "entryUUID: " + parentUUID);
+    final Entry child = TestCaseUtils.makeEntry(
+        "dn: uid=user.922.8,ou=parked.922," + baseDN,
+        "objectClass: top",
+        "objectClass: person",
+        "objectClass: organizationalPerson",
+        "objectClass: inetOrgPerson",
+        "uid: user.922.8",
+        "cn: Aaccf Amar",
+        "sn: Amar",
+        "entryUUID: " + childUUID);
+
+    /*
+     * Both adds are held at the pre-parse plugin point, one at a time. The first park is
+     * what keeps the parent listed as pending - and owned by the thread replaying it - while
+     * the child is checked for dependencies, so the child is parked behind a change which is
+     * in flight rather than behind one which has already been applied.
+     */
+    final CSN parentCsn = gen.newCSN();
+    final CSN childCsn = gen.newCSN();
+    final ParkedReplay parked = ShortCircuitPlugin.parkReplayedOperations(
+        OperationType.ADD, "PreParse",
+        op -> parentCsn.equals(OperationContext.getCSN(op))
+            || childCsn.equals(OperationContext.getCSN(op)));
+    try
+    {
+      domain.processUpdate(new AddMsgWhoseAckThrows(parentCsn, parent.getName(), parentUUID,
+          baseUUID, parent.getObjectClassAttribute(), parent.getAllAttributes()));
+      final Thread replayingParent = parked.awaitParked(60, SECONDS);
+
+      domain.processUpdate(new AddMsg(childCsn, child.getName(), childUUID, parentUUID,
+          child.getObjectClassAttribute(), child.getAllAttributes(), null));
+
+      /*
+       * The parent is applied and its ack throws where it is published. The replay carries
+       * on all the same, and the child is the change it hands itself next.
+       */
+      parked.release();
+      final Thread replayingChild = parked.awaitParked(60, SECONDS);
+      Assertions.assertThat(replayingChild)
+          .as("the change which was parked must be replayed by the thread which cleared what"
+              + " it was waiting for, rather than be left waiting")
+          .isSameAs(replayingParent);
+      parked.release();
+
+      assertNotNull(getEntry(child.getName(), 30000, true),
+          "the change which was parked behind the one whose ack threw must be applied");
+    }
+    finally
+    {
+      parked.deregister();
+    }
+  }
+
+  /**
+   * Test case for [Issue 922]: a change handed out as a dependency is given back when the
+   * replay it was handed to is unwound.
+   * <p>
+   * A change which was parked behind another one is handed out by {@code getNextUpdate()}
+   * to the thread which cleared what it was waiting for, and that thread owns it from then
+   * on. The give-back on the way out of an unwound replay asks which change this thread
+   * owns, so the hand-out has to be recorded where that question is answered, not only on
+   * the change: left out, the change would stay owned by a thread which is not replaying
+   * it anymore, and every later delivery of it would be refused as a duplicate - the wedge
+   * of this issue, on the dependency road.
+   * <p>
+   * The parent is held at the pre-parse plugin point while the child is delivered, so the
+   * child is parked behind a change in flight, and the thread which is thrown out of the
+   * child is read at the same plugin point: it must be the one which committed the parent,
+   * which is what says the child was handed out rather than taken off the queue. The child
+   * is thrown out of once, inside the replay, and then unwound past its ack, on the road
+   * every catch of the replay has already run on.
+   */
+  @Test
+  public void aChangeHandedOutAsADependencyIsGivenBackWhenItsReplayIsUnwound() throws Exception
+  {
+    testSetUp("aChangeHandedOutAsADependencyIsGivenBackWhenItsReplayIsUnwound");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : aChangeHandedOutAsADependencyIsGivenBackWhenItsReplayIsUnwound"));
+
+    final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null);
+    final CSNGenerator gen = new CSNGenerator(29, TimeThread.getTime());
+    final String parentUUID = "29292929-2929-2929-2929-292929292929";
+    final String childUUID = "30303030-3030-3030-3030-303030303030";
+    final Entry parent = TestCaseUtils.makeEntry(
+        "dn: ou=handed-out.922," + baseDN,
+        "objectClass: top",
+        "objectClass: organizationalUnit",
+        "ou: handed-out.922",
+        "entryUUID: " + parentUUID);
+    final Entry child = TestCaseUtils.makeEntry(
+        "dn: uid=user.922.11,ou=handed-out.922," + baseDN,
+        "objectClass: top",
+        "objectClass: person",
+        "objectClass: organizationalPerson",
+        "objectClass: inetOrgPerson",
+        "uid: user.922.11",
+        "cn: Aaccf Amar",
+        "sn: Amar",
+        "entryUUID: " + childUUID);
+
+    final CSN parentCsn = gen.newCSN();
+    final CSN childCsn = gen.newCSN();
+    final long initialFailures = getMonitorAttrValue(baseDN, "replayed-updates-failed");
+    final ParkedReplay parked = ShortCircuitPlugin.parkReplayedOperations(
+        OperationType.ADD, "PreParse", op -> parentCsn.equals(OperationContext.getCSN(op)));
+    /*
+     * The thread which reaches the plugin point with the child is the one replaying it, and
+     * it is read there rather than parked: a park would hold the child where the throw is
+     * made, and it is the throw which is wanted.
+     */
+    final AtomicReference<Thread> replayingChild = new AtomicReference<>();
+    final ThrownFromReplay thrown = ShortCircuitPlugin.throwFromReplayedOperations(
+        OperationType.ADD, "PreParse",
+        op ->
+        {
+          if (!childCsn.equals(OperationContext.getCSN(op)))
+          {
+            return false;
+          }
+          replayingChild.set(Thread.currentThread());
+          return true;
+        },
+        () -> new LinkageError("the replay of the change which was handed out meets an Error"),
+        1);
+    try
+    {
+      domain.processUpdate(new AddMsg(parentCsn, parent.getName(), parentUUID, baseUUID,
+          parent.getObjectClassAttribute(), parent.getAllAttributes(), null));
+      final Thread replayingParent = parked.awaitParked(60, SECONDS);
+
+      domain.processUpdate(new AddMsgWhoseReplayIsUnwoundAfterItsAck(childCsn, child.getName(),
+          childUUID, parentUUID, child.getObjectClassAttribute(), child.getAllAttributes()));
+
+      // The parent is applied, and the child is the change its thread hands itself next.
+      parked.release();
+
+      TestTimer timer = new TestTimer.Builder()
+        .maxSleep(60, SECONDS)
+        .sleepTimes(200, MILLISECONDS)
+        .toTimer();
+      timer.repeatUntilSuccess(new CallableVoid()
+      {
+        @Override
+        public void call() throws Exception
+        {
+          assertEquals(thrown.thrownCount(), 1,
+              "the change which was handed out must have been thrown out of");
+        }
+      });
+      Assertions.assertThat(replayingChild.get())
+          .as("the change which was parked must be replayed by the thread which cleared what"
+              + " it was waiting for: that is the hand-out this test is about")
+          .isSameAs(replayingParent);
+
+      /*
+       * The child was thrown out of and its replay was then unwound, so it is not in the
+       * data and must not be in the ServerState - and it must not be given up on either: it
+       * is asked for again. The delivery which asks for it is made here, the way
+       * assertChangeIsDeliveredAgainAfter() makes it: nothing sends the change again, since
+       * it never travelled a session, and a delivery is dropped rather than queued while the
+       * session is being restarted. A delivery of a change a thread still owns is refused as
+       * the duplicate it is, which is where a change left owned by a thread which is not
+       * replaying it never comes back.
+       */
+      assertFalse(domain.getServerState().cover(childCsn),
+          "a change whose replay was unwound must be asked for again, not recorded as replayed");
+      timer.repeatUntilSuccess(new CallableVoid()
+      {
+        @Override
+        public void call() throws Exception
+        {
+          if (!domain.getServerState().cover(childCsn))
+          {
+            domain.processUpdate(new AddMsg(childCsn, child.getName(), childUUID, parentUUID,
+                child.getObjectClassAttribute(), child.getAllAttributes(), null));
+          }
+          assertTrue(domain.getServerState().cover(childCsn),
+              "the change must be recorded as replayed once it has been delivered again");
+        }
+      });
+      assertNotNull(getEntry(child.getName(), 30000, true),
+          "the change which was handed out must be applied by the delivery which took over"
+              + " from the one which was unwound");
+      assertEquals(getMonitorAttrValue(baseDN, "replayed-updates-failed"), initialFailures,
+          "a change which was delivered again must not be counted as one this replica gave up on");
+      assertTrue(replayingParent.isAlive(),
+          "an Error which unwinds a replay must not end the thread which met it (issue #923)");
+    }
+    finally
+    {
+      thrown.deregister();
+      parked.deregister();
+    }
+  }
+
+  /**
+   * Test case for [Issue 922]: the ack of a delivery whose replay threw an Error says that
+   * the change was not applied.
+   * <p>
+   * An assured write in SAFE_READ mode is told that its change is durable here by the ack
+   * this replica publishes, and it is published in a finally which every road out of the
+   * replay runs through - the roads an Error unwinds among them. A replica which is asking
+   * for a change again must never have told a master that the change is in its data, so the
+   * ack of a delivery whose replay threw reports the error and names this replica.
+   * <p>
+   * The error is thrown from a plugin point which runs inside the operation, so the change
+   * travels a real session and the ack can be read off the broker which published it - the
+   * counters can not report it, since handing the change back restarts the session and that
+   * resets every one of them.
+   */
+  @Test
+  public void theAckOfADeliveryWhoseReplayThrewAnErrorReportsIt() throws Exception
+  {
+    testSetUp("theAckOfADeliveryWhoseReplayThrewAnErrorReportsIt");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : theAckOfADeliveryWhoseReplayThrewAnErrorReportsIt"));
+
+    final int serverId = 28;
+    /*
+     * In the group of the replication server, so that what is published below is one this
+     * domain has to acknowledge: an assured update from a broker of another group is
+     * acknowledged by the replication server itself, and says nothing about the replay.
+     */
+    ReplicationBroker broker =
+        openAssuredReplicationSession(baseDN, serverId, 100, replServerPort, 1000);
+    try
+    {
+      CSNGenerator gen = new CSNGenerator(serverId, 0);
+      Entry tmp = TestCaseUtils.addEntry(
+          "dn: uid=user.922.9," + baseDN,
+          "objectClass: top",
+          "objectClass: person",
+          "objectClass: organizationalPerson",
+          "objectClass: inetOrgPerson",
+          "uid: user.922.9",
+          "cn: Aaccf Amar",
+          "sn: Amar");
+      final String uuid = getEntry(tmp.getName(), 1, true).parseAttribute("entryuuid").asString();
+      final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null);
+
+      final CSN csn = gen.newCSN();
+      /*
+       * Thrown out of one replay and no more: the delivery which takes over from the one
+       * which was unwound is what applies the change, and a change this replica could never
+       * replay would be given up on rather than acknowledged twice.
+       */
+      final ThrownFromReplay thrown = ShortCircuitPlugin.throwFromReplayedOperations(
+          OperationType.DELETE, "PreParse",
+          op -> csn.equals(OperationContext.getCSN(op)),
+          () -> new LinkageError("the replay of this change meets an Error"), 1);
+      try
+      {
+        final DeleteMsg delete = new DeleteMsg(tmp.getName(), csn, uuid);
+        delete.setAssured(true);
+        delete.setAssuredMode(AssuredMode.SAFE_READ_MODE);
+        broker.publish(delete);
+
+        final AckMsg ack = awaitAck(broker, csn);
+        assertTrue(ack.hasReplayError(),
+            "the ack of a delivery whose replay threw an Error must report it rather than be"
+                + " the plain ack a master would take for a durable write");
+        assertFalse(ack.hasTimeout(),
+            "the ack must be the one the delivery published, not the one the replication"
+                + " server makes up when it gives up waiting for it");
+        Assertions.assertThat(ack.getFailedServers())
+            .as("the replica whose replay threw must be the one the ack names")
+            .containsExactly(domainSid);
+        assertEquals(thrown.thrownCount(), 1, "the replay of the change must have been unwound");
+      }
+      finally
+      {
+        thrown.deregister();
+      }
+
+      /*
+       * The change was given back rather than recorded as replayed, so the delivery which
+       * follows applies it - which is what the ack above said had not happened yet.
+       */
+      assertNull(getEntry(tmp.getName(), 30000, false),
+          "the change must be applied by the delivery which took over from the one whose"
+              + " replay threw");
+      assertTrue(domain.getServerState().cover(csn),
+          "the change must be recorded as replayed once it has been applied");
+    }
+    finally
+    {
+      broker.stop();
+    }
+  }
+
+  /**
+   * Test case for [Issue 922]: a change whose replay keeps throwing an Error is given up
+   * on rather than asked for forever.
+   * <p>
+   * An Error takes the road every other failed replay takes, so the failures it leaves
+   * behind count against the give-up budget of the change: a change this replica can never
+   * apply must not hold its ServerState - and every change which follows it, from every
+   * master - back for good, whether its replay reported the failure or threw it.
+   */
+  @Test
+  public void aChangeWhoseReplayKeepsThrowingAnErrorIsGivenUpOn() throws Exception
+  {
+    testSetUp("aChangeWhoseReplayKeepsThrowingAnErrorIsGivenUpOn");
+    logger.error(LocalizableMessage.raw(
+        "Starting replication test : aChangeWhoseReplayKeepsThrowingAnErrorIsGivenUpOn"));
+
+    Entry tmp = TestCaseUtils.addEntry(
+        "dn: uid=user.922.4," + baseDN,
+        "objectClass: top",
+        "objectClass: person",
+        "objectClass: organizationalPerson",
+        "objectClass: inetOrgPerson",
+        "uid: user.922.4",
+        "cn: Aaccf Amar",
+        "sn: Amar");
+    final DN dn = tmp.getName();
+    final String uuid = getEntry(dn, 1, true).parseAttribute("entryuuid").asString();
+
+    final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null);
+    final long initialFailures = getMonitorAttrValue(baseDN, "replayed-updates-failed");
+    domain.resetUnreplayedChangeAlertThrottle();
+    final int initialAlerts = DummyAlertHandler.getAlertCount(ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE);
+    setReplayGiveUpDelay(TEST_GIVE_UP_DELAY);
+    try
+    {
+      final CSN csn = new CSNGenerator(22, TimeThread.getTime()).newCSN();
+      final List<Modification> mods =
+          generatemods("description", "the replay of this change keeps throwing");
+
+      /*
+       * Nothing sends this change again - it never travelled a session - so every delivery
+       * of it is made here, until this replica gives up on the change it can not replay.
+       */
+      TestTimer timer = new TestTimer.Builder()
+        .maxSleep(120, SECONDS)
+        .sleepTimes(200, MILLISECONDS)
+        .toTimer();
+      timer.repeatUntilSuccess(new CallableVoid()
+      {
+        @Override
+        public void call() throws Exception
+        {
+          if (!domain.getServerState().cover(csn))
+          {
+            domain.processUpdate(new ModifyMsgWhoseReplayThrows(csn, dn, mods, uuid,
+                () -> new LinkageError("the replay of this change meets an Error")));
+          }
+          assertTrue(domain.getServerState().cover(csn),
+              "a change whose replay keeps throwing must be given up on");
+        }
+      });
+      assertMonitorAttrValueEventually(baseDN, "replayed-updates-failed", initialFailures + 1,
+          "the change which was given up on must be counted as failed, once");
+      Assertions.assertThat(DummyAlertHandler.getAlertCount(ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE))
+          .as("the administrator must be told that this replica now diverges")
+          .isGreaterThan(initialAlerts);
+    }
+    finally
+    {
+      resetReplayGiveUpDelay();
+    }
+  }
+
+  /** A delivery of a change whose replay does not run to its end. */
+  private interface FailingDelivery
+  {
+    ModifyMsg newDelivery(CSN csn, DN dn, List<Modification> mods, String entryUUID);
+  }
+
+  /**
+   * Delivers a change whose replay does not run to its end, then checks that the change is
+   * neither recorded as replayed nor lost: the delivery which follows must be able to
+   * replay it.
+   *
+   * @param delivery the delivery whose replay is to be unwound
+   * @param serverId the replica the change comes from, one per test so that the CSNs of
+   *                 one are never covered by the ServerState another left behind
+   * @param uid the entry the change is made on
+   * @param description the value the change writes, once it is replayed
+   */
+  private void assertChangeIsDeliveredAgainAfter(
+      FailingDelivery delivery, int serverId, String uid, final String description) throws Exception
+  {
+    Entry tmp = TestCaseUtils.addEntry(
+        "dn: uid=" + uid + "," + baseDN,
+        "objectClass: top",
+        "objectClass: person",
+        "objectClass: organizationalPerson",
+        "objectClass: inetOrgPerson",
+        "uid: " + uid,
         "cn: Aaccf Amar",
         "sn: Amar");
     final DN dn = tmp.getName();
@@ -2601,12 +3330,10 @@
     domain.resetUnreplayedChangeAlertThrottle();
     final int initialAlerts = DummyAlertHandler.getAlertCount(ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE);
 
-    final CSNGenerator gen = new CSNGenerator(18, TimeThread.getTime());
-    final CSN csn = gen.newCSN();
-    final String description = "the replay must fail once the operation is built";
+    final CSN csn = new CSNGenerator(serverId, TimeThread.getTime()).newCSN();
     final List<Modification> mods = generatemods("description", description);
 
-    domain.processUpdate(new ModifyMsgWhoseOperationRefusesAControl(csn, dn, mods, uuid));
+    domain.processUpdate(delivery.newDelivery(csn, dn, mods, uuid));
 
     /*
      * Long enough to outlast the session restart the failure asks for: a change which is
@@ -2615,24 +3342,22 @@
     for (int i = 0; i < MONITOR_ATTR_SAMPLES_ACROSS_A_REDELIVERY; i++)
     {
       assertFalse(domain.getServerState().cover(csn),
-          "a change whose operation was built must be asked for again, not recorded as replayed");
+          "a change whose replay was unwound must be asked for again, not recorded as replayed");
       Thread.sleep(200);
     }
     assertMonitorAttrValueStays(baseDN, "replayed-updates-failed", initialFailures,
         "a change which is still to be delivered again must not be counted as given up on");
-    assertEquals(DummyAlertHandler.getAlertCount(ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE), initialAlerts,
-        "a change which is still to be delivered again must not be alerted on as a divergence");
+    assertEquals(DummyAlertHandler.getAlertCount(ALERT_TYPE_REPLICATION_UNREPLAYED_CHANGE),
+        initialAlerts, "a change which is still to be delivered again must not be alerted on"
+            + " as a divergence");
 
     /*
-     * The failed change is the barrier which holds this domain's ServerState back until
-     * it is replayed, and the replication server sending it again is what replays it.
-     * Nothing sends this one - it never travelled a session - so the delivery which takes
-     * over from the one which failed is made here, and it is made until it is taken: a
-     * delivery is dropped rather than queued while the listener thread is down, which it
-     * is for as long as the recovery is restarting the session, and the monitor entry
-     * read above comes back with the broker rather than with the listener. A delivery of
-     * a change a replay thread owns is refused as the duplicate it is, and the ServerState
-     * keeps this from delivering a change which was replayed a second time.
+     * Nothing sends this change again - it never travelled a session - so the delivery
+     * which takes over from the one which was unwound is made here, and it is made until
+     * it is taken: a delivery is dropped rather than queued while the listener thread is
+     * down, which it is for as long as the recovery is restarting the session. A delivery
+     * of a change a replay thread owns is refused as the duplicate it is, so this is where
+     * a change left owned by a thread which is gone never comes back.
      */
     TestTimer timer = new TestTimer.Builder()
       .maxSleep(60, SECONDS)
@@ -2652,7 +3377,83 @@
       }
     });
     checkEntryHasAttributeValue(dn, "description", description, 30,
-        "the change must be applied by the delivery which took over from the failed one");
+        "the change must be applied by the delivery which took over from the one which failed");
+  }
+
+  /**
+   * Returns the identities of the replay threads which are running.
+   * <p>
+   * Read by identity rather than counted where a test is about one thread in particular
+   * having ended: the pool belongs to the server rather than to a test, so a count says
+   * whether it is the size it was, not whether the thread which met the error is the one
+   * which is gone.
+   */
+  private static Set<Long> replayThreadIds()
+  {
+    final Set<Long> running = new HashSet<>();
+    for (Thread thread : Thread.getAllStackTraces().keySet())
+    {
+      if (thread.isAlive() && thread.getName().startsWith("Replica replay thread "))
+      {
+        running.add(thread.getId());
+      }
+    }
+    return running;
+  }
+
+  /**
+   * A ModifyMsg whose replay throws an Error.
+   * <p>
+   * It is thrown where the operation is built, which is inside the replay and past the
+   * point where the change was marked as being replayed by the thread which took it: what
+   * this pins is the road out of a replay which no {@code catch} of the replay itself used
+   * to run on.
+   */
+  private static final class ModifyMsgWhoseReplayThrows extends ModifyMsg
+  {
+    private final Supplier<Error> error;
+
+    private ModifyMsgWhoseReplayThrows(
+        CSN csn, DN dn, List<Modification> mods, String entryUUID, Supplier<Error> error)
+    {
+      super(csn, dn, mods, entryUUID);
+      this.error = error;
+    }
+
+    @Override
+    public ModifyOperation createOperation(InternalClientConnection connection, DN newDN)
+    {
+      throw error.get();
+    }
+  }
+
+  /**
+   * A ModifyMsg whose replay fails and whose ack throws on the way out of it.
+   * <p>
+   * Its operation is built and can not be prepared for its replay, the way
+   * {@code ModifyMsgWhoseOperationRefusesAControl} has it, so the replay fails with the
+   * change owned by the replay thread. The ack of the delivery is then published in the
+   * finally every road out of the replay runs through, and this one throws there - which is
+   * what a session being torn down does.
+   */
+  private static final class ModifyMsgWhoseAckThrows
+      extends ModifyMsgWhoseOperationRefusesAControl
+  {
+    private ModifyMsgWhoseAckThrows(CSN csn, DN dn, List<Modification> mods, String entryUUID)
+    {
+      super(csn, dn, mods, entryUUID);
+    }
+
+    @Override
+    public boolean isAssured()
+    {
+      /*
+       * An Error rather than the exception a session being torn down raises: the two take
+       * the same road, and an Error is what a catch of Exception would let past - the guard
+       * around the ack has to hold whatever publishing it threw.
+       */
+      throw new LinkageError("the ack of this delivery can not be published");
+    }
   }
 
   /**
@@ -2725,6 +3526,96 @@
   }
 
   /**
+   * A ModifyMsg whose ack runs out of memory on the way out of a replay which failed.
+   * <p>
+   * The replay fails first - its operation can not be prepared for the replay, the way
+   * {@code ModifyMsgWhoseOperationRefusesAControl} has it fail - so the change is one this
+   * replica asks for again, and the ack which says so is where the JVM runs out of memory.
+   * That is the one throw from there which is not caught: it ends the replay thread, and
+   * the change is given back on the way out.
+   */
+  private static final class ModifyMsgWhoseAckRunsOutOfMemory
+      extends ModifyMsgWhoseOperationRefusesAControl
+  {
+    private ModifyMsgWhoseAckRunsOutOfMemory(
+        CSN csn, DN dn, List<Modification> mods, String entryUUID)
+    {
+      super(csn, dn, mods, entryUUID);
+    }
+
+    @Override
+    public boolean isAssured()
+    {
+      // Read first thing by processUpdateDone(), which is what publishes the ack.
+      throw new OutOfMemoryError("the ack of this delivery runs out of memory");
+    }
+  }
+
+  /**
+   * An AddMsg whose ack throws on the way out of a replay which applied it.
+   * <p>
+   * The change is committed before the ack of its delivery is published, so this is the
+   * road on which nothing is owed back to the replication server - and on which the changes
+   * parked behind this one are still waiting to be handed out.
+   */
+  private static final class AddMsgWhoseAckThrows extends AddMsg
+  {
+    private AddMsgWhoseAckThrows(CSN csn, DN dn, String entryUUID, String parentEntryUUID,
+        Attribute objectClasses, Iterable<Attribute> userAttributes)
+    {
+      super(csn, dn, entryUUID, parentEntryUUID, objectClasses, userAttributes, null);
+    }
+
+    @Override
+    public boolean isAssured()
+    {
+      // Read first thing by processUpdateDone(), which is what publishes the ack.
+      throw new LinkageError("the ack of this delivery can not be published");
+    }
+  }
+
+  /**
+   * An AddMsg whose replay is unwound once the ack of its delivery is out, the way
+   * {@link ModifyMsgWhoseReplayIsUnwoundAfterItsAck} is.
+   * <p>
+   * What has its replay fail is not the message but the test which delivers it, through a
+   * throw at the pre-parse plugin point: an add which is parked behind its parent has to be
+   * one whose operation builds and runs, or it would be given up on where no operation could
+   * be built from it. The throw is caught inside the replay, so it is what the replay runs
+   * once the ack is out - the give-back of the failed change - which reads the CSN off this
+   * message and is unwound by it.
+   */
+  private static final class AddMsgWhoseReplayIsUnwoundAfterItsAck extends AddMsg
+  {
+    private volatile boolean ackPublished;
+
+    private AddMsgWhoseReplayIsUnwoundAfterItsAck(CSN csn, DN dn, String entryUUID,
+        String parentEntryUUID, Attribute objectClasses, Iterable<Attribute> userAttributes)
+    {
+      super(csn, dn, entryUUID, parentEntryUUID, objectClasses, userAttributes, null);
+    }
+
+    @Override
+    public boolean isAssured()
+    {
+      // Read first thing by processUpdateDone(), and by nothing on the way in: a message
+      // handed to the domain rather than published is not one this server acknowledges.
+      ackPublished = true;
+      return super.isAssured();
+    }
+
+    @Override
+    public CSN getCSN()
+    {
+      if (ackPublished)
+      {
+        throw new LinkageError("the replay of this change is unwound once its ack is out");
+      }
+      return super.getCSN();
+    }
+  }
+
+  /**
    * Test case for [Issue 908]: a domain being disabled - for an LDIF import, a restore, or
    * a backend being taken offline - must not save its ServerState while a replay thread is
    * half way through applying one of its changes.
@@ -3057,6 +3948,52 @@
   }
 
   /**
+   * A ModifyMsg whose replay is unwound once the ack of its delivery is out.
+   * <p>
+   * Its operation is built and can not be prepared for its replay, the way
+   * {@code ModifyMsgWhoseOperationRefusesAControl} has it, so the replay fails with the
+   * change owned by the replay thread. The CSN of the change is then read again - by the
+   * road which gives it back and asks for it again - and it is that read which throws here:
+   * past the ack, past the finally it is published in, and past every catch the replay
+   * itself has. So the only thing left to give the change back is the road out of
+   * {@code replay()}.
+   */
+  private static final class ModifyMsgWhoseReplayIsUnwoundAfterItsAck
+      extends ModifyMsgWhoseOperationRefusesAControl
+  {
+    private volatile boolean ackPublished;
+
+    private ModifyMsgWhoseReplayIsUnwoundAfterItsAck(
+        CSN csn, DN dn, List<Modification> mods, String entryUUID)
+    {
+      super(csn, dn, mods, entryUUID);
+    }
+
+    @Override
+    public boolean isAssured()
+    {
+      /*
+       * processUpdateDone() reads this first and reads nothing else of a delivery which is
+       * not assured, so the ack of this one is out by the time it returns. Nothing on the
+       * way in reads it: a message handed to the domain rather than published is not one
+       * this server acknowledges to anybody.
+       */
+      ackPublished = true;
+      return super.isAssured();
+    }
+
+    @Override
+    public CSN getCSN()
+    {
+      if (ackPublished)
+      {
+        throw new LinkageError("the replay of this change is unwound once its ack is out");
+      }
+      return super.getCSN();
+    }
+  }
+
+  /**
    * A ModifyMsg whose operation can not be prepared for its replay.
    * <p>
    * The operation is built - so the replay is past the point where a message is given up
@@ -3066,7 +4003,7 @@
    * protocol: {@code ModifyMsg.createOperation()} builds an operation whose controls can
    * be added to, so this one is handed to the domain rather than published.
    */
-  private static final class ModifyMsgWhoseOperationRefusesAControl extends ModifyMsg
+  private static class ModifyMsgWhoseOperationRefusesAControl extends ModifyMsg
   {
     private ModifyMsgWhoseOperationRefusesAControl(
         CSN csn, DN dn, List<Modification> mods, String entryUUID)

--
Gitblit v1.10.0