From 41f5692c778b797fe09b5a658af37e4bc1da83ad Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Thu, 24 Sep 2026 11:56:08 +0000
Subject: [PATCH] [#1061] Ask for the session restart under an owner as well, and leave it to the state checkpointer (#1062)
---
opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ReplayDuringImportTest.java | 373 +++++++++++++++++++++++++++++++++++++++++++++++++---
1 files changed, 349 insertions(+), 24 deletions(-)
diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ReplayDuringImportTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ReplayDuringImportTest.java
index e921beb..0919df6 100644
--- a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ReplayDuringImportTest.java
+++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ReplayDuringImportTest.java
@@ -72,15 +72,20 @@
* one thing it must not do is stop the session the import is reading (issue #956). The same
* holds from the moment the total update is asked for: the answer to the request arrives
* over that session, so a replay which fails while it is on its way must not restart it.
- * A restart asked for before the total update took the session, and left standing for the
- * length of it, is not run once it is over either: the change it was asked for is gone with
- * the ServerState the import replaced. The changes a replay which is unwound had parked as
- * waiting for another one are released on the same terms, and nothing more is done for them
- * (issue #954).
+ * A restart asked for while the total update owns the session is left standing for the
+ * length of it, and is not run once it is over: the change it was asked for is gone with the
+ * ServerState the import replaced. The changes a replay which is unwound had parked as
+ * waiting for another one are released on the same terms (issue #954). A total update which
+ * is asked for and never begins - the request is refused, or gives up waiting for its answer
+ * - replaces nothing, and the request left standing under it is what has the changes
+ * released under it delivered again (issue #1061).
* <p>
* The exporter is a broker of this test, so that the test says when the entries arrive: the
* change is replayed while the import is waiting for them - or, for the request, while the
- * exporter is holding the answer.
+ * exporter is holding the answer, which it may never give. A change which has to be delivered
+ * again is published through the replication server, which is what has it to send again;
+ * the replay queue of the domain is the test's, so every delivery is replayed when, and on
+ * the thread, the test says.
* <p>
* The claim of a total update this replica did not ask for is made by the listener thread
* under no lock, so a restart of the session which reads no owner a moment before that claim
@@ -309,12 +314,13 @@
/**
* The changes a replay which is unwound had parked as waiting for another change are
- * released and nothing more while a total update owns the session (issue #954): no
- * session restart is asked for them - the one it would ask for is refused where it runs,
- * and the request would be spent on it - and they are neither reported as changes the
- * replication server sends again, which it does not before the import has replaced the
- * data, nor counted as processed. That is the road a change a stopping replay thread
- * abandons takes on this domain, and the give-back of the parked changes takes it too.
+ * released while a total update owns the session, and the session is left to the owner
+ * (issue #954): the restart asked for them is not run - it is refused where it runs, and
+ * the request would be spent on it - and they are neither reported as changes the
+ * replication server sends again, which it does not before the total update has let go
+ * of the session, nor counted as processed. That is the road a change a stopping replay
+ * thread abandons takes on this domain, and the give-back of the parked changes takes it
+ * too.
* <p>
* Pinned on the import road because it is the one road with an owner which a test holds
* open for as long as it needs: the request is on its way until the exporter answers it,
@@ -329,7 +335,7 @@
* back, and the error which ends a replay thread is caught here instead.
*/
@Test(timeOut = 120_000)
- public void aParkedChangeGivenBackWhileTheRequestIsOnItsWayIsNotAskedForAgain() throws Exception
+ public void aParkedChangeGivenBackWhileTheRequestIsOnItsWayLeavesTheSessionToTheOwner() throws Exception
{
final Entry entry = TestCaseUtils.addEntry(
"dn: cn=renamedSince," + EXAMPLE_DN,
@@ -399,10 +405,10 @@
"the change parked by the replay which was unwound must be given back");
assertEquals(getMonitorAttrValue(baseDN, "replayed-updates"), processed,
"a change released while a total update owns the session must not be counted as"
- + " processed: no session sends it again before the import has replaced the data");
+ + " processed: no session sends it again before the total update lets go of it");
assertThat(errorLogRecordsOf(NOTE_REPLAY_PARKED_CHANGE_GIVEN_BACK.ordinal(), parked))
.as("the change was reported as one the replication server sends again, which it does"
- + " not before the import has replaced the data")
+ + " not before the total update lets go of the session")
.isEmpty();
assertTrue(domain.isConnected(), "the session the answer to the request arrives over was stopped");
@@ -417,13 +423,287 @@
}
/**
+ * A change released while a total update which never begins owns the session is asked for
+ * again under the owner, and delivered again over the session the state checkpointer
+ * restarts once the owner is gone (issue #1061).
+ * <p>
+ * A total update this replica asked for owns the session from the request on, and a change
+ * whose replay fails meanwhile is released and left to the owner: the domain forgets its
+ * pending changes on its way down, and the import forgets them at its end - but a request
+ * which is refused, or which gives up waiting for its answer, replaces nothing and forgets
+ * nothing. Released and asked for by nobody, the change would stay listed and uncommitted
+ * until the next failed replay of this domain restarted the session, and the ServerState -
+ * which a commit moves no further than the oldest uncommitted change - would stop at it
+ * with everything behind it. So the restart is asked for under the owner as well, and the
+ * state checkpointer, which holds every request for as long as the total update owns the
+ * session, runs it within its tick of the owner letting go.
+ * <p>
+ * The road pinned here is the one a change whose attempts in place are spent takes: every
+ * attempt ends on an entryUUID search which does not run. The request gives up through
+ * the watchdog of the initialize task, which is the one road out of an unanswered request
+ * a test can take at a time of its choosing.
+ */
+ @Test(timeOut = 120_000)
+ public void aChangeReleasedUnderARequestWhichIsNeverAnsweredIsDeliveredAgain() throws Exception
+ {
+ final Entry entry = TestCaseUtils.addEntry(
+ "dn: cn=renamedSince," + EXAMPLE_DN,
+ "objectClass: top",
+ "objectClass: person",
+ "cn: renamedSince",
+ "sn: renamedSince");
+ final String entryUUID = getEntryUUID(entry.getName());
+
+ // The request is out, and the exporter never answers it.
+ domain.initializeFromRemote(EXPORTER_ID, null);
+ assertNotNull(waitForSpecificMsg(exporter, InitializeRequestMsg.class));
+
+ final CSN csn = gen.newCSN();
+ final LDAPUpdateMsg delivered = publishAndAwaitDelivery(new ModifyMsg(csn,
+ DN.valueOf("cn=movedAway," + EXAMPLE_DN),
+ generatemods("description", "released while the request was on its way"), entryUUID));
+ // A restart run on this thread would be back up before any read of the session: this counts it.
+ domain.failNextSessionRestarts(1);
+ ShortCircuitPlugin.registerShortCircuit(
+ OperationType.SEARCH, "PreParse", ResultCode.UNAVAILABLE.intValue());
+ try
+ {
+ replay(delivered);
+ }
+ finally
+ {
+ ShortCircuitPlugin.deregisterShortCircuit(OperationType.SEARCH, "PreParse");
+ }
+ assertFalse(domain.getServerState().cover(csn),
+ "the change whose replay fails must stay listed as one which is not in the data");
+ assertThat(errorLogRecordsOf(WARN_REPLAY_RETRYING_CHANGE.ordinal(), csn))
+ .as("the change was warned about as one the replication server sends again, which it"
+ + " does not while the total update owns the session")
+ .isEmpty();
+ assertEquals(domain.getSessionRestartFailuresLeft(), 1,
+ "a session restart was run under the owner by the replay whose attempts were spent");
+ domain.failNextSessionRestarts(0);
+
+ giveUpTheRequest();
+
+ final LDAPUpdateMsg again = awaitDelivery(csn, 30_000, "the change released under the"
+ + " request was not delivered again once the request gave up: nothing asked for the"
+ + " session restart which has the replication server send it again");
+ replay(again);
+ assertThat(DirectoryServer.getEntry(entry.getName()).getAllAttributes("description"))
+ .as("the change delivered again was not applied").isNotEmpty();
+ assertTrue(domain.getServerState().cover(csn),
+ "the change delivered again was applied and not recorded: it is still listed");
+ }
+
+ /**
+ * A change a stopping replay thread abandons while a total update which never begins owns
+ * the session takes the same road (issue #1061): abandoned at the top of its first attempt
+ * without being counted against its budget, released, asked for again under the owner,
+ * and delivered again once the owner is gone.
+ */
+ @Test(timeOut = 120_000)
+ public void aChangeAbandonedUnderARequestWhichIsNeverAnsweredIsDeliveredAgain() throws Exception
+ {
+ final Entry entry = TestCaseUtils.addEntry(
+ "dn: cn=renamedSince," + EXAMPLE_DN,
+ "objectClass: top",
+ "objectClass: person",
+ "cn: renamedSince",
+ "sn: renamedSince");
+ final String entryUUID = getEntryUUID(entry.getName());
+
+ domain.initializeFromRemote(EXPORTER_ID, null);
+ assertNotNull(waitForSpecificMsg(exporter, InitializeRequestMsg.class));
+
+ final CSN csn = gen.newCSN();
+ final LDAPUpdateMsg delivered = publishAndAwaitDelivery(new ModifyMsg(csn, entry.getName(),
+ generatemods("description", "abandoned while the request was on its way"), entryUUID));
+ // The thread of this test is one which is stopping: the change is abandoned unapplied.
+ assertTrue(domain.markInProgress(delivered), "the delivery was not handed to this thread");
+ domain.failNextSessionRestarts(1);
+ domain.replay(delivered, new AtomicBoolean(true));
+ assertThat(DirectoryServer.getEntry(entry.getName()).getAllAttributes("description"))
+ .as("a change abandoned by a stopping thread was applied").isEmpty();
+ assertThat(errorLogRecordsOf(NOTE_REPLAY_ABANDONED_CHANGE.ordinal(), csn))
+ .as("the change was reported as one the replication server sends again, which it"
+ + " does not while the total update owns the session")
+ .isEmpty();
+ assertEquals(getMonitorAttrValue(baseDN, "changes-with-failed-replay"), 0,
+ "the abandoned change was counted against its budget: it was never attempted");
+ assertEquals(domain.getSessionRestartFailuresLeft(), 1,
+ "a session restart was run under the owner by the thread which abandoned the change");
+ domain.failNextSessionRestarts(0);
+
+ giveUpTheRequest();
+
+ final LDAPUpdateMsg again = awaitDelivery(csn, 30_000, "the change abandoned under the"
+ + " request was not delivered again once the request gave up: nothing asked for the"
+ + " session restart which has the replication server send it again");
+ replay(again);
+ assertTrue(domain.getServerState().cover(csn),
+ "the change delivered again was applied and not recorded: it is still listed");
+ }
+
+ /**
+ * A change a stopping replay thread abandons once the import has forgotten it asks for no
+ * session restart (issue #1061).
+ * <p>
+ * The import forgets the pending changes at its end, and the session it starts asks for
+ * everything the ServerState it loaded does not cover. A thread which read the flag while
+ * the import ran and reaches the give-back only after that finds its change unlisted: a
+ * restart asked for then would outlive the clear, and the state checkpointer would stop
+ * and start the session the import has just started, for a change which is gone.
+ */
+ @Test(timeOut = 120_000)
+ public void aChangeAbandonedOnceTheImportForgotItAsksForNoRestart() throws Exception
+ {
+ final String[] exported = exportedEntries();
+
+ // Owned by this thread before the import begins, and still owned when the import ends.
+ final CSN csn = gen.newCSN();
+ domain.processUpdate(new ModifyMsg(csn, DN.valueOf(IMPORTED_ENTRY_DN),
+ generatemods("description", "abandoned once the import forgot it"), IMPORTED_ENTRY_UUID));
+ final LDAPUpdateMsg delivered = queue.take().getUpdateMessage();
+ assertTrue(domain.markInProgress(delivered), "the delivery was not handed to this thread");
+
+ startImportInto(exported.length);
+ finishImport(exported);
+
+ domain.failNextSessionRestarts(1);
+ try
+ {
+ domain.replay(delivered, new AtomicBoolean(true));
+
+ // Two ticks of the checkpointer: a request standing now is run on the first.
+ Thread.sleep(2000);
+ assertEquals(domain.getSessionRestartFailuresLeft(), 1, "a change the import forgot"
+ + " asked for a restart of the session the import started at its end");
+ assertThat(errorLogRecordsOf(NOTE_REPLAY_ABANDONED_CHANGE.ordinal(), csn))
+ .as("a change the import forgot was reported as one the replication server sends again")
+ .isEmpty();
+ }
+ finally
+ {
+ domain.failNextSessionRestarts(0);
+ }
+ }
+
+ /**
+ * A change the give-back released while a total update which never begins owns the session
+ * is asked for again under the owner by the give-back itself, and the request is left
+ * standing rather than spent (issue #1061).
+ * <p>
+ * The change it waited for is one another thread is replaying, held before its operation
+ * is built: nothing has failed, so no road but the give-back has asked for anything, and
+ * the request found standing once the owner is gone is the give-back's own. No restart is
+ * run on the thread of the give-back under the owner: one run there would be refused where
+ * it runs, and the case counts the restarts which get that far.
+ * <p>
+ * The replay which is unwound is this thread's, as in the case above: its change is
+ * applied, and the ack of its delivery runs out of memory.
+ */
+ @Test(timeOut = 120_000)
+ public void aParkedChangeGivenBackUnderARequestWhichIsNeverAnsweredIsDeliveredAgain() throws Exception
+ {
+ final Entry entry = TestCaseUtils.addEntry(
+ "dn: cn=renamedSince," + EXAMPLE_DN,
+ "objectClass: top",
+ "objectClass: person",
+ "cn: renamedSince",
+ "sn: renamedSince");
+ final String entryUUID = getEntryUUID(entry.getName());
+ final Entry other = TestCaseUtils.addEntry(
+ "dn: cn=unwound," + EXAMPLE_DN,
+ "objectClass: top",
+ "objectClass: person",
+ "cn: unwound",
+ "sn: unwound");
+ final String otherUUID = getEntryUUID(other.getName());
+
+ domain.initializeFromRemote(EXPORTER_ID, null);
+ assertNotNull(waitForSpecificMsg(exporter, InitializeRequestMsg.class));
+
+ /*
+ * The change the parked one waits for: replayed by a thread of the test which is held
+ * before the operation is built, so the change is being replayed - listed, uncommitted,
+ * owned - for as long as the latch holds, and nothing has failed.
+ */
+ final CountDownLatch letGo = new CountDownLatch(1);
+ final CSN held = gen.newCSN();
+ domain.processUpdate(new ModifyMsgWhoseOperationWaitsToBeBuilt(held, entry.getName(),
+ generatemods("description", "the change being replayed by another thread"), entryUUID,
+ letGo));
+ final LDAPUpdateMsg heldMsg = queue.take().getUpdateMessage();
+ final Thread otherThread = new Thread(() ->
+ {
+ assertTrue(domain.markInProgress(heldMsg), "the held change must be the one listed");
+ domain.replay(heldMsg, SHUTDOWN);
+ }, "ReplayDuringImportTest replay held before its operation is built");
+ otherThread.start();
+ try
+ {
+ // Parked as waiting for the held change by this thread, which owns it from here on.
+ final CSN parked = gen.newCSN();
+ final LDAPUpdateMsg parkedDelivery = publishAndAwaitDelivery(new ModifyMsg(parked,
+ entry.getName(),
+ generatemods("description", "the change which was parked as a dependency"), entryUUID));
+ replay(parkedDelivery);
+ assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 1,
+ "a change which waits for one being replayed by another thread must be parked");
+
+ // The replay which is unwound while this thread still holds the parked change.
+ final CSN unwound = gen.newCSN();
+ domain.failNextSessionRestarts(1);
+ try
+ {
+ replayMsg(new ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied(unwound, other.getName(),
+ generatemods("description", "the replay of this change is unwound once it is applied"),
+ otherUUID));
+ Assert.fail("the replay was not unwound: the ack of the delivery must run out of memory");
+ }
+ catch (OutOfMemoryError unwinding)
+ {
+ // The error is the fixture's own, and this is the thread it would have ended.
+ }
+ assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 0,
+ "the change parked by the replay which was unwound must be given back");
+ assertEquals(domain.getSessionRestartFailuresLeft(), 1,
+ "a session restart was run under the owner by the replay which was unwound");
+ domain.failNextSessionRestarts(0);
+
+ // The held change runs to its end: nothing fails, nothing asks for a restart.
+ letGo.countDown();
+ otherThread.join(30_000);
+ assertFalse(otherThread.isAlive(), "the held replay did not end once let go");
+ assertTrue(domain.getServerState().cover(held), "the held change was not recorded");
+ assertFalse(domain.getServerState().cover(parked),
+ "the parked change must stay listed as one which is not in the data");
+
+ giveUpTheRequest();
+
+ final LDAPUpdateMsg again = awaitDelivery(parked, 30_000, "the parked change given back"
+ + " under the request was not delivered again once the request gave up: the session"
+ + " restart the give-back asked for under the owner was not run, or was spent");
+ replay(again);
+ assertTrue(domain.getServerState().cover(parked),
+ "the parked change delivered again was applied and not recorded: it is still listed");
+ }
+ finally
+ {
+ letGo.countDown();
+ otherThread.join(30_000);
+ }
+ }
+
+ /**
* A session restart which stood while the import ran was asked for by a replay thread
- * for a change given back before the total update owned the session, and that change is
- * forgotten with the pending changes when the imported data replaces the ServerState:
- * the session started back at the end of the import asks for everything the imported
- * state does not cover. Run, the request would stop that session once for a delivery
- * which can not come. The request is made here by hand, in the place of one made
- * between a replay thread's read of the owner and the import claiming the session.
+ * for a change it gave back - before the total update owned the session, or under the
+ * owner - and that change is forgotten with the pending changes when the imported data
+ * replaces the ServerState: the session started back at the end of the import asks for
+ * everything the imported state does not cover. Run, the request would stop that session
+ * once for a delivery which can not come. The request is made here by hand, in the place
+ * of the one a failed replay makes.
* <p>
* The restart is the state checkpointer's to run, within its first tick after the total
* update has released the session, so the pin is that the failure it would meet is never
@@ -983,8 +1263,53 @@
private void replayMsg(UpdateMsg updateMsg) throws InterruptedException
{
domain.processUpdate(updateMsg);
- final LDAPUpdateMsg ldapUpdate = queue.take().getUpdateMessage();
- domain.markInProgress(ldapUpdate);
- domain.replay(ldapUpdate, SHUTDOWN);
+ replay(queue.take().getUpdateMessage());
+ }
+
+ /** Replays a delivery on the thread of this test, as a replay thread would. */
+ private void replay(LDAPUpdateMsg delivery)
+ {
+ assertTrue(domain.markInProgress(delivery), "the delivery is not the one listed: " + delivery);
+ domain.replay(delivery, SHUTDOWN);
+ }
+
+ /**
+ * Publishes a change through the replication server, which is what has it to deliver again
+ * once the session is restarted for it, and waits for the delivery to this replica.
+ */
+ private LDAPUpdateMsg publishAndAwaitDelivery(LDAPUpdateMsg msg) throws Exception
+ {
+ exporter.publish(msg);
+ return awaitDelivery(msg.getCSN(), 30_000, "the change published was not delivered");
+ }
+
+ /**
+ * Waits for the replication server to deliver the change to this replica: the listener
+ * thread of the domain puts it in the replay queue of the test, which takes it out.
+ */
+ private LDAPUpdateMsg awaitDelivery(CSN csn, long timeoutMs, String orElse) throws Exception
+ {
+ final long deadline = System.currentTimeMillis() + timeoutMs;
+ while (queue.peek() == null)
+ {
+ assertTrue(System.currentTimeMillis() < deadline, orElse + " within " + timeoutMs + " ms");
+ Thread.sleep(50);
+ }
+ final LDAPUpdateMsg msg = queue.take().getUpdateMessage();
+ assertEquals(msg.getCSN(), csn, "another change than the one awaited was delivered");
+ return msg;
+ }
+
+ /**
+ * Has the total update this replica asked for give up on its request, the way the
+ * watchdog of the initialize task does once the request has waited two minutes for an
+ * answer: the total update never begins, and its context is released with nothing
+ * replaced, so the session has no owner anymore.
+ */
+ private void giveUpTheRequest()
+ {
+ assertTrue(domain.abortStalledInitializeFromRemote(0),
+ "the request was not the one waiting for an answer");
+ assertFalse(domain.ieRunning(), "the total update was given up and is still being processed");
}
}
--
Gitblit v1.10.0