diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java index 1238dbd63f..678fc25c68 100644 --- a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java +++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java @@ -2786,7 +2786,7 @@ void replay(LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown) * another road asked for with the backoff, or one given back with it, keeps its * wait whichever thread runs it. */ - final boolean parkedGivenBack = giveBackParkedChanges( + final List parkedGivenBack = giveBackParkedChanges( replayThreadShutdown.get() || t instanceof OutOfMemoryError ? SessionRestart.NOW : SessionRestart.AFTER_BACKOFF); if (owned != null) @@ -2822,8 +2822,17 @@ void replay(LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown) */ recoverFromReplayFailure(owned, replayThreadShutdown, t instanceof OutOfMemoryError); } + /* + * The change is handed back and counted as the road it took counts it: the last + * resort below speaks for a give-back which did not run, and a throw out of what + * follows - the restart the parked changes are run with, or the line which says it + * is held - is not one. Left set, that throw would have the last resort report the + * change as "released without its failure being counted", which is the one line an + * operator acts on, and it would be false. + */ + owned = null; } - if (parkedGivenBack && !replayThreadShutdown.get() && !sessionHasAnOwner()) + if (!parkedGivenBack.isEmpty() && !replayThreadShutdown.get() && !sessionHasAnOwner()) { /* * The road the change this thread was replaying took may have run the restart the @@ -2838,11 +2847,31 @@ void replay(LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown) * the state checkpointer runs one restart for every change the threads of the * pool hand back on their way out, rather than each of them running one while * the configuration change which is stopping them waits. A domain whose session - * has an owner is left alone the way the give-back left it: nothing was asked - * for on that road, and a request another thread left standing is not this - * one's to spend on a restart which is refused where it runs. + * has an owner is left alone the way the give-back left it - it asked for nothing + * there, which is why it handed back no change to run a restart for - and a + * request another thread left standing is not this one's to spend on a restart + * which is refused where it runs. */ - runRequestedSessionRestarts(); + if (!runRequestedSessionRestarts()) + { + /* + * The changes above were reported as given back to a replication server which + * "still owns it and sends it again", and it does not send them yet: a total + * update is being processed over the session, the restart which brings them back + * waits for it, and the state checkpointer runs it once it is over - the same + * hold, and the same line, a change whose replay failed is reported with. One + * line per change, the way the give-back reports them: these are the changes + * this one is about, and nothing else says they wait. + * + * Built on the road out of a JVM which has run out of memory too, where the line + * of a failed replay is not: the give-back has already built one line per change + * on that road, so what this asks the JVM for is not memory it was spared. + */ + for (CSN csn : parkedGivenBack) + { + logger.info(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE, csn, getBaseDN()); + } + } } } catch (Throwable recoveryFailure) @@ -2856,7 +2885,8 @@ void replay(LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown) * whatever happens here, and the thread this runs on may well be ending on it. A * restart which can not run here leaves its request standing, and the state * checkpointer of this domain runs it: one which threw asks for itself again on - * its way out, and one a domain whose session has an owner would refuse is not + * its way out, one a total update is being processed over the session is not run + * while it lasts, and one a domain whose session has an owner would refuse is not * run at all rather than spent on the refusal. * * The restart is asked for once the change is released and not before, the way @@ -2888,12 +2918,14 @@ void replay(LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown) try { /* - * Outside the guard above: two roads reach here with a request standing and no - * change of this thread's to hand back, and both are the parked changes' - the + * Outside the guard above: three roads reach here with a request standing and no + * change of this thread's to hand back, and all are the parked changes' - the * give-back which released them asks for the restart before it reports them, and * a throw out of the report - the JVM which unwound this replay is out of memory - * - leaves the request standing; and a restart the parked road ran and which - * threw has asked for one again on its way out. The changes it released are + * - leaves the request standing; a restart the parked road ran and which threw + * has asked for one again on its way out; and a throw out of the line which says + * a total update holds that restart leaves the request standing too, where it is + * refused until the total update is over. The changes it released are * listed, uncommitted and unowned, so the request is what brings them back, and * this thread is the one there to run it (issue #954). A give-back which threw * before it released anything left the parked changes as they were, owned by this @@ -3540,8 +3572,10 @@ else if (op instanceof ModifyDNOperation) if (replayFailed && recoverFromReplayFailure(msg.getCSN(), replayThreadShutdown)) { // The ack has been published and the change is given back: the replication server - // delivers it again, now or - while a total update owns the session - after the - // import restarts it. There is nothing left to replay here. + // delivers it again, now or - while a total update is being processed over the + // session - once it is over, when the restart which was held runs or, on the import + // direction, when the session is started from the reloaded state. There is nothing + // left to replay here. return; } @@ -3822,8 +3856,11 @@ public void resetUnreplayedChangeAlertThrottle() * * @param csn the CSN of the change which could not be replayed * @param failure how long, and over how many deliveries, its replay has been failing + * @return whether the warning was written, so that a line which qualifies it - the one + * which says the restart it announced is held - is written with it and folded + * with it */ - private void logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure failure) + private boolean logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure failure) { final long now = monotonicNowInMs(); final long lastLogged = lastReplayRetryWarningTime.get(); @@ -3832,6 +3869,7 @@ private void logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure f { logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts(), failure.getFailingForMs(), foldedReplayRetryWarnings.getAndSet(0)); + return true; } else { @@ -3845,6 +3883,7 @@ private void logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure f foldedReplayRetryWarnings.incrementAndGet(); logger.trace("Could not replay change %s in domain %s: delivery %d, failing for %d ms", csn, getBaseDN(), failure.getAttempts(), failure.getFailingForMs()); + return false; } } @@ -4001,6 +4040,7 @@ private boolean recoverFromReplayFailure( return true; } + boolean warned = false; if (!outOfMemory) { /* @@ -4012,7 +4052,7 @@ private boolean recoverFromReplayFailure( * unlogged: the error ends the replay thread, and the uncaught exception handler of * DirectoryThread writes the line and raises the alert for it. */ - logReplayRetryWarning(csn, failure); + warned = logReplayRetryWarning(csn, failure); } /* * This change is not owned by anyone anymore, so the session has to be restarted for @@ -4031,7 +4071,32 @@ private boolean recoverFromReplayFailure( */ sessionRestarts.request(replayThreadShutdown.get() || outOfMemory ? SessionRestart.NOW : SessionRestart.AFTER_BACKOFF); - runRequestedSessionRestarts(); + if (!runRequestedSessionRestarts() && warned) + { + /* + * The warning above said the session is being restarted for the change, and it is not + * yet: a total update is being processed over that session - almost always an export + * from this replica, since a total update into it owns the session and is refused + * above, except by an import which claims its context between that read and this one + * (issue #1041) - and the restart waits for it, for as long as the total update takes. + * Said on its own, so that a change which is not delivered again for minutes is not a + * change nobody asked for. + * + * Written where that warning was written and nowhere else: it qualifies that line, so + * a backend which fails every delivery of a long export would otherwise be one of + * these per delivery while the warnings they qualify are folded into a count - the + * very repetition the throttle is there to fold. It is not built at all on the road + * out of a JVM which has run out of memory, for the reason the warning is not built + * there: warned is false on it. + * + * On the import road of the window above the request is cleared by importBackend() + * rather than run, so the restart this line announces never runs there. The change is + * delivered again all the same, which is what an operator reads this for: the import + * loads the ServerState of the exporter, and the session started at its end asks for + * everything that state does not cover. + */ + logger.info(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE, csn, getBaseDN()); + } return true; } @@ -4051,15 +4116,48 @@ private boolean recoverFromReplayFailure( */ void giveBackChangesParkedByStoppingThread() { + // What it handed back is not read here: this road runs no restart for them, so it has + // nothing to say about one which waits. giveBackParkedChanges(SessionRestart.NOW); } /** * Restarts the session as long as changes which could not be replayed are waiting to be - * delivered again. + * delivered again - unless a total update is being processed, in which case the requests + * are left standing for the state checkpointer to run once it is over. + *

+ * A restart stops the session the total update runs over, in either direction. An import + * into this replica reads its entries from that session and would end on the ones which + * had arrived - and {@code disabled} does not say an import is running, since + * {@code preBackendImport()} keeps the backend events this domain is the cause of from + * disabling it. An export from this replica publishes its entries over it, and + * {@code exportLDIFEntry()} gives the export up as + * {@code ERR_INIT_RS_DISCONNECTION_DURING_EXPORT} once the broker has been stopped under + * it, which leaves the replica it was initializing to be initialized again: minutes on a + * large backend, spent for a change which would have waited. So the change waits: the + * request stays standing, the state checkpointer comes for it once a second and runs it + * as soon as the total update is over ({@link #runPendingSessionRestart()}), and the + * replication server delivers the change again then. The ServerState waits with it, and + * the replay of this domain keeps running in the meantime. + *

+ * A total update which begins between this read and any of the stops this call makes is + * not seen here, and is cut by it: the read and the claim of the import/export context + * share no lock, which is issue #1041 on the import side. The window is the call rather + * than a few statements of it: {@code restartSession()} stops the session before it waits + * its backoff out, so a request taken later in the loop stops the session a backoff wait + * after the read - up to {@link #MAX_REPLAY_RETRY_DELAY_IN_MS}, and the restart of a + * request another thread has just been told is held is one of those. Before this the whole + * of the total update was that window. + * + * @return {@code false} when a total update is being processed and the requests were left + * standing for the state checkpointer, {@code true} otherwise */ - private void runRequestedSessionRestarts() + private boolean runRequestedSessionRestarts() { + if (ieRunning()) + { + return false; + } /* * The outer loop is what makes a request which was made while this thread was giving * up the recovery its own: the thread which made it found the recovery taken and left @@ -4099,6 +4197,7 @@ private void runRequestedSessionRestarts() replayFailureRecovery.set(false); } } + return true; } /** @@ -4114,22 +4213,15 @@ private void runRequestedSessionRestarts() * topology, with the changes it did not replay owned by the replication server and its * ServerState stopped behind them. *

- * Not run while a total update is being processed, in either direction: a restart stops - * the session the total update runs over. An import into this replica reads its entries - * from that session and would end on the ones which had arrived - and {@code disabled} - * does not say an import is running, since {@code preBackendImport()} keeps the backend - * events this domain is the cause of from disabling it. An export from this replica - * publishes its entries over it, and {@code exportLDIFEntry()} gives the export up as - * {@code ERR_INIT_RS_DISCONNECTION_DURING_EXPORT} once the broker has been stopped - * under it, which leaves the replica it was initializing to be initialized again. This - * thread is the one which can afford to wait: the request stays standing, and it comes - * back here once a second, so the restart is run as soon as the total update is over. - * The change the restart was asked for waits for as long as the total update takes, and - * the ServerState with it; the replay of this domain keeps running in the meantime. + * While a total update is being processed, in either direction, the restart is not run - + * no restart asked for by a released change is, see {@link #runRequestedSessionRestarts()} + * - and this thread is the one which can afford to wait for it: the request stays + * standing, and it comes back here once a second, so the restart is run as soon as the + * total update is over. */ private void runPendingSessionRestart() { - if (shutdown.get() || disabled || ieRunning() || !sessionRestarts.isPending()) + if (shutdown.get() || disabled || !sessionRestarts.isPending()) { return; } @@ -4210,22 +4302,27 @@ private void runPendingSessionRestart() * @param restart what the session restart is asked for as: with the backoff a failing * backend is owed, or without it on a thread which is stopping or which an * OutOfMemoryError is ending - * @return whether any change was handed back: a change which nobody owns is one only a - * new delivery brings back, so the caller runs the restart asked for them - on - * a thread which is not stopping, and on a domain whose session has no owner + * @return the changes it handed back and asked the restart for, oldest first: a change + * which nobody owns is one only a new delivery brings back, so the caller runs the + * restart asked for them - on a thread which is not stopping, and on a domain whose + * session has no owner - and reports what that restart did for them. Empty when + * this thread had parked none, and empty on a domain whose session has an owner, + * where nothing is asked for. The list is the one the release allocated: nothing is + * allocated for the answer on the road out of a JVM which has run out of memory */ - private boolean giveBackParkedChanges(SessionRestart restart) + private List giveBackParkedChanges(SessionRestart restart) { final List parked = remotePendingChanges.releaseParkedChangesOwnedByCurrentThread(); if (parked.isEmpty()) { - return false; + return parked; } if (sessionHasAnOwner()) { // The domain owns its session, or a total update does: both forget the pending - // changes, and neither leaves a session for this thread to restart. - return true; + // changes, and neither leaves a session for this thread to restart. Nothing was asked + // for here, so the caller has nothing of this road's to run or to report. + return Collections.emptyList(); } /* * Asked for before the changes are reported: a throw out of the report - the JVM which @@ -4238,7 +4335,7 @@ private boolean giveBackParkedChanges(SessionRestart restart) incProcessedUpdates(); logger.info(NOTE_REPLAY_PARKED_CHANGE_GIVEN_BACK, csn, getBaseDN()); } - return true; + return parked; } /** diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/service/ReplicationDomain.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/service/ReplicationDomain.java index 80c6bbf5bc..3dbca34167 100644 --- a/opendj-server-legacy/src/main/java/org/opends/server/replication/service/ReplicationDomain.java +++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/service/ReplicationDomain.java @@ -3521,8 +3521,12 @@ public final void disableService() * {@code SESSION_BEING_STOPPED}, held for the length of the stop and released once the * listener thread is gone - it is the one thread which claims a total update this replica * did not ask for, and {@link #disableService()} waits for it. An export in the context is - * not an owner: the session is stopped from under it and the exporter reports the cut, as - * it does for every other stop. A total update which lands between the end of that export + * not an owner here: the session is stopped from under it and the exporter reports the + * cut, as it does for every other stop. The session restart a replay asks for does not + * get this far while a total update runs - {@code runRequestedSessionRestarts()} leaves + * the request standing until it is over (issue #1048) - so an export reaches this arm only + * when it begins between that read and the claim below. A total update which lands + * between the end of that export * and the stop is refused by the listener when it reads the broker as stopping after its * claim; a stop which lands after that read still has the import run over a session which * is going down, and end as a failed import over the suffix it has replaced (issue diff --git a/opendj-server-legacy/src/messages/org/opends/messages/replication.properties b/opendj-server-legacy/src/messages/org/opends/messages/replication.properties index 0276c4e193..28420418a7 100644 --- a/opendj-server-legacy/src/messages/org/opends/messages/replication.properties +++ b/opendj-server-legacy/src/messages/org/opends/messages/replication.properties @@ -715,3 +715,7 @@ WARN_IGNORING_UPDATE_UNSUPPORTED_BY_PEER_328=Replication server RS(%d) not sendi ERR_INIT_REJECTED_SESSION_STOPPING_330=The total update of domain "%s" was refused by directory \ server %d: its session to the replication server is being stopped, and the entries would have \ been streamed over that session. Ask for the total update again once the session is back +NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE_331=The session restart asked for by change %s \ + in domain "%s" is held: a total update is being processed over that session, and stopping it \ + would end the total update. The restart runs once the total update is over, and the change is \ + sent again then; until then it is not recorded as replayed diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ModifyMsgWhoseOperationRefusesAControl.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ModifyMsgWhoseOperationRefusesAControl.java index 9db7874b03..26cfe101e0 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ModifyMsgWhoseOperationRefusesAControl.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ModifyMsgWhoseOperationRefusesAControl.java @@ -46,7 +46,7 @@ * so this one is handed to the domain rather than published. The twin of the fixture * {@code UpdateOperationTest} holds its barrier with. */ -final class ModifyMsgWhoseOperationRefusesAControl extends ModifyMsg +class ModifyMsgWhoseOperationRefusesAControl extends ModifyMsg { ModifyMsgWhoseOperationRefusesAControl( CSN csn, DN dn, List mods, String entryUUID) diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ParkedChangeGiveBackTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ParkedChangeGiveBackTest.java index 614b021bb5..096f01081d 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ParkedChangeGiveBackTest.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ParkedChangeGiveBackTest.java @@ -25,6 +25,7 @@ import java.util.SortedSet; import java.util.TreeSet; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; import org.forgerock.opendj.ldap.DN; import org.forgerock.opendj.server.config.meta.ReplicationDomainCfgDefn.IsolationPolicy; @@ -38,6 +39,7 @@ import org.opends.server.replication.server.ReplServerFakeConfiguration; import org.opends.server.replication.server.ReplicationServer; import org.opends.server.types.Entry; +import org.opends.server.types.Modification; import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -90,6 +92,12 @@ public class ParkedChangeGiveBackTest extends ReplicationTestCase private LDAPReplicationDomain domain; private TestSynchronousReplayQueue queue; private CSNGenerator gen; + /** + * The budget the replay of a change is retried for, read by the domain at every decision: + * the default until a case lowers it, so that a case can give up on one change while the + * barrier it set up before stays asked for. + */ + private final AtomicLong replayGiveUpDelayInMs = new AtomicLong(); @BeforeMethod public void setUpLocal() throws Exception @@ -103,7 +111,16 @@ public void setUpLocal() throws Exception final SortedSet replServers = new TreeSet<>(); replServers.add("localhost:" + rsPort); - final DomainFakeCfg conf = new DomainFakeCfg(baseDN, DS_ID, replServers); + final DomainFakeCfg conf = new DomainFakeCfg(baseDN, DS_ID, replServers) + { + @Override + public long getReplayGiveUpDelay() + { + return replayGiveUpDelayInMs.get(); + } + }; + // What the fake configuration spells out as the default of the property. + replayGiveUpDelayInMs.set(new DomainFakeCfg(baseDN, DS_ID, replServers).getReplayGiveUpDelay()); conf.setIsolationPolicy(IsolationPolicy.ACCEPT_ALL_UPDATES); queue = new TestSynchronousReplayQueue(); domain = MultimasterReplication.createNewDomain(conf, queue); @@ -133,11 +150,15 @@ public void tearDown() throws Exception * brings the session back before the replay returns. The restart it asked for is run at * once: the one run again after the failure is the one which waits its backoff out, and * that is the one move of the count of the restarts in a row. + *

+ * The restart ran, so nothing of it is reported as held: no total update is being processed + * over the session here, and the line which says a restart waits for one belongs to the + * domain which is exporting or importing (issue #1048). */ @Test(timeOut = 120_000) public void theThreadWhichGaveBackAParkedChangeRunsTheRestartItAskedFor() throws Exception { - parkAChangeBehindABarrier(addEntry("waitedOn")); + final CSN parked = parkAChangeBehindABarrier(addEntry("waitedOn")); final Entry other = addEntry("other"); final int restartsBefore = domain.getConsecutiveSessionRestarts(); domain.failNextSessionRestarts(1); @@ -157,6 +178,11 @@ public void theThreadWhichGaveBackAParkedChangeRunsTheRestartItAskedFor() throws assertEquals(domain.getConsecutiveSessionRestarts(), restartsBefore + 1, "the restart the give-back asked for must be run at once, without the backoff: the one" + " run again after the failure is the one which waits it out"); + assertThat(errorLogRecordsOf( + NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE.ordinal(), parked)) + .as("the restart the parked change was given back with was reported as held: no total" + + " update is being processed over the session, and it ran") + .isEmpty(); } /** @@ -231,6 +257,60 @@ public void aStoppingThreadHandsBackTheChangeItDidNotApply() throws Exception "the change parked by the thread which is stopping must be given back"); } + /** + * A change which the road of a failed replay has already taken is not reported as one whose + * give-back failed when what follows that road throws: the last resort of {@code replay()} + * speaks for a give-back which did not run, and this one ran. + *

+ * The change fails and the ack of its delivery runs out of memory, so the replay is unwound + * with the change still owned by this thread. Its give-up budget is spent at once, so the + * road it takes gives it up and runs no restart; the restart the parked change was given + * back with is run after it, by the same thread, and is asked to fail - the throw out of + * what follows that road. + */ + @Test(timeOut = 120_000) + public void aThrowAfterTheChangeWasTakenCareOfIsNotReportedAsAFailedGiveBack() + throws Exception + { + parkAChangeBehindABarrier(addEntry("waitedOn")); + final Entry other = addEntry("other"); + final CSN givenUpOn = gen.newCSN(); + replayGiveUpDelayInMs.set(0); + domain.failNextSessionRestarts(1); + + OutOfMemoryError unwinding = null; + try + { + replayMsg(new ModifyMsgWhoseAckRunsOutOfMemory(givenUpOn, other.getName(), + generatemods("description", "the replay of this change fails and its ack runs out" + + " of memory"), getEntryUUID(other.getName())), RUNNING); + } + catch (OutOfMemoryError e) + { + // The error is the fixture's own, and this is the thread it would have ended. + unwinding = e; + } + assertNotNull(unwinding, "the replay was not unwound: the ack of the delivery must run out" + + " of memory"); + + assertThat(errorLogRecordsOf(ERR_REPLAY_SKIPPING_CHANGE.ordinal(), givenUpOn)) + .as("the change was not given up on: its budget was spent, so the road it took must" + + " have skipped it rather than run a restart for it") + .isNotEmpty(); + assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 0, + "the change parked by the replay which was unwound must be given back"); + assertThat(injectedRestartFailuresAmong(unwinding.getSuppressed())) + .as("the restart the parked change was given back with must have been run after the" + + " road of the failed change, and have met the failure it was asked to meet") + .hasSize(1); + assertThat(errorLogRecordsOf(ERR_REPLAY_GIVE_BACK_FAILED.ordinal(), givenUpOn)) + .as("the change was reported as released without its failure being counted: it was" + + " counted and given up on, and the throw came after that") + .isEmpty(); + awaitConnected(RESTART_BOUND_IN_MS, "the restart run again after the one which failed did" + + " not bring the session back"); + } + private void awaitConnected(long boundInMs, String message) throws Exception { final long deadline = System.currentTimeMillis() + boundInMs; @@ -256,9 +336,12 @@ private Entry addEntry(String cn) throws Exception * operation is built and then refused, so it is asked for again rather than stepped over - * and then a change on the same entry, which is parked as waiting for it and owned by this * thread from then on. The restart the failed change asks for is run on this thread as - * well, so the session is back once this returns. + * well, so the session is back once this returns - and nothing of it is reported as held, + * which is the negative arm of that line (issue #1048). + * + * @return the CSN of the change which is left parked */ - private void parkAChangeBehindABarrier(Entry entry) throws Exception + private CSN parkAChangeBehindABarrier(Entry entry) throws Exception { final String entryUUID = getEntryUUID(entry.getName()); final CSN failing = gen.newCSN(); @@ -268,12 +351,19 @@ private void parkAChangeBehindABarrier(Entry entry) throws Exception "the change whose replay fails must stay listed as one which is not in the data"); awaitConnected(RESTART_BOUND_IN_MS, "the session was not brought back for the change whose replay failed"); + assertThat(errorLogRecordsOf( + NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE.ordinal(), failing)) + .as("the restart the change whose replay failed was asked for again with was reported" + + " as held: no total update is being processed over the session, and it ran") + .isEmpty(); - replayMsg(new ModifyMsg(gen.newCSN(), entry.getName(), + final CSN parked = gen.newCSN(); + replayMsg(new ModifyMsg(parked, entry.getName(), generatemods("description", "the change which was parked as a dependency"), entryUUID), RUNNING); assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 1, "a change which waits for one that is not in the data must be parked"); + return parked; } /** When the thread a replay runs on is stopped, if it is. */ @@ -370,4 +460,26 @@ private void replayMsg(UpdateMsg updateMsg, AtomicBoolean stopping) throws Inter domain.markInProgress(ldapUpdate); domain.replay(ldapUpdate, stopping); } + + /** + * A ModifyMsg whose replay fails and whose ack runs out of memory: the ack is published once + * the failure is decided and before the change is given back, so the replay is unwound with + * the change still owned by the thread which was replaying it. + */ + private static final class ModifyMsgWhoseAckRunsOutOfMemory + extends ModifyMsgWhoseOperationRefusesAControl + { + private ModifyMsgWhoseAckRunsOutOfMemory( + CSN csn, DN dn, List 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"); + } + } } diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ReplayDuringExportTest.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ReplayDuringExportTest.java new file mode 100644 index 0000000000..713a942847 --- /dev/null +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ReplayDuringExportTest.java @@ -0,0 +1,717 @@ +/* + * The contents of this file are subject to the terms of the Common Development and + * Distribution License (the License). You may not use this file except in compliance with the + * License. + * + * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the + * specific language governing permission and limitations under the License. + * + * When distributing Covered Software, include this CDDL Header Notice in each file and include + * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL + * Header, with the fields enclosed by brackets [] replaced by your own identifying + * information: "Portions copyright [year] [name of copyright owner]". + * + * Copyright 2026 3A Systems, LLC. + */ +package org.opends.server.replication.plugin; + +import static java.nio.charset.StandardCharsets.*; +import static org.assertj.core.api.Assertions.*; +import static org.opends.messages.ReplicationMessages.*; +import static org.opends.server.TestCaseUtils.*; +import static org.opends.server.core.DirectoryServer.*; +import static org.testng.Assert.*; + +import java.lang.reflect.Field; +import java.lang.reflect.Method; +import java.net.SocketTimeoutException; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; +import java.util.SortedSet; +import java.util.TreeSet; +import java.util.concurrent.atomic.AtomicBoolean; + +import org.forgerock.opendj.ldap.DN; +import org.forgerock.opendj.ldap.ResultCode; +import org.forgerock.opendj.server.config.meta.ReplicationDomainCfgDefn.IsolationPolicy; +import org.opends.server.TestCaseUtils; +import org.opends.server.core.DirectoryServer; +import org.opends.server.plugins.ShortCircuitPlugin; +import org.opends.server.replication.ReplicationTestCase; +import org.opends.server.replication.common.CSN; +import org.opends.server.replication.common.CSNGenerator; +import org.opends.server.replication.common.ServerStatus; +import org.opends.server.replication.protocol.DoneMsg; +import org.opends.server.replication.protocol.EntryMsg; +import org.opends.server.replication.protocol.ErrorMsg; +import org.opends.server.replication.protocol.InitializeRcvAckMsg; +import org.opends.server.replication.protocol.InitializeRequestMsg; +import org.opends.server.replication.protocol.InitializeTargetMsg; +import org.opends.server.replication.protocol.LDAPUpdateMsg; +import org.opends.server.replication.protocol.ModifyMsg; +import org.opends.server.replication.protocol.ReplicationMsg; +import org.opends.server.replication.protocol.UpdateMsg; +import org.opends.server.replication.server.ReplServerFakeConfiguration; +import org.opends.server.replication.server.ReplicationServer; +import org.opends.server.replication.service.ReplicationBroker; +import org.opends.server.replication.service.ReplicationDomain; +import org.opends.server.types.Entry; +import org.opends.server.types.OperationType; +import org.testng.Assert; +import org.testng.annotations.AfterMethod; +import org.testng.annotations.BeforeMethod; +import org.testng.annotations.Test; + +/** + * Tests the replay of a change while this replica is the source of a total update. + *

+ * The export of a total update publishes its entries over the session of the domain, from a + * thread of the export pool, while the replay of the domain keeps running. A change which can + * not be replayed meanwhile is given back for the replication server to send again, and that + * takes a session restart. Run by the thread which released the change, the restart stops the + * broker the export publishes over, and {@code exportLDIFEntry()} gives the export up on it: the + * replica being initialized is left to be initialized again, for a change which would have + * waited (issue #1048). The restart has to wait for the export instead, and the state + * checkpointer runs it once the export is over. + *

+ * The importer is a broker of this test, so that the test says when the export moves: the + * exporter publishes no more than the initialization window ahead of the importer's + * acknowledgements, and the change is replayed while the export waits for one. + */ +@SuppressWarnings("javadoc") +public class ReplayDuringExportTest extends ReplicationTestCase +{ + /** + * A total update needs a backend which keeps its data across the export, and one which + * the exporter can lock: the {@code userRoot} backend, as for the import direction. + */ + private static final String EXAMPLE_DN = "dc=example,dc=com"; + private static final int RS_ID = 612; + private static final int DS_ID = 1; + private static final int IMPORTER_ID = 2; + /** How many entry messages the exporter publishes ahead of the importer's acknowledgements. */ + private static final int INIT_WINDOW = 2; + /** + * An entry message carries a buffer of the export stream rather than one entry, so the data + * has to outgrow the window by that much before the exporter waits for an acknowledgement. + */ + private static final int ENTRY_MSG_BYTES = 8192; + private static final int BULK_ENTRY_BYTES = 4096; + private static final int BULK_ENTRIES = 2 * (INIT_WINDOW + 2); + private static final AtomicBoolean SHUTDOWN = new AtomicBoolean(false); + /** + * How long the export is given to release its context once the importer has left the full + * update status - the exporter waits for that status to go, and releases it then. + *

+ * Short enough that the receive of the stream and this wait fit inside the timeout of the + * case with room to spare: a case which ends on the timeout prints none of the messages + * which say what went wrong. + */ + private static final long EXPORT_END_BOUND_IN_MS = 30_000; + /** + * How long the session is given to come back once a restart has been run for a change which + * was given back: the wait that restart is owed, and the start of the session. + */ + private static final long SESSION_BACK_BOUND_IN_MS = 30_000; + + private DN baseDN; + private ReplicationServer replicationServer; + private LDAPReplicationDomain domain; + private TestSynchronousReplayQueue queue; + private ReplicationBroker importer; + private CSNGenerator gen; + + @BeforeMethod + public void setUpLocal() throws Exception + { + baseDN = DN.valueOf(EXAMPLE_DN); + TestCaseUtils.clearBackend("userRoot", EXAMPLE_DN); + + final int rsPort = TestCaseUtils.findFreePort(); + replicationServer = new ReplicationServer(new ReplServerFakeConfiguration( + rsPort, "replayDuringExportTestDb", 0, RS_ID, 0, 100, new TreeSet())); + + final SortedSet replServers = new TreeSet<>(); + replServers.add("localhost:" + rsPort); + final DomainFakeCfg conf = new DomainFakeCfg(baseDN, DS_ID, replServers); + conf.setIsolationPolicy(IsolationPolicy.ACCEPT_ALL_UPDATES); + queue = new TestSynchronousReplayQueue(); + domain = MultimasterReplication.createNewDomain(conf, queue); + domain.start(); + assertTrue(domain.isConnected(), "the domain did not connect to the replication server"); + + // A short socket timeout: the test bounds its own waits, and receive() returns to it on it. + importer = openReplicationSession(baseDN, IMPORTER_ID, 100, rsPort, 2000); + gen = new CSNGenerator(IMPORTER_ID, 0); + } + + @AfterMethod + public void tearDown() throws Exception + { + try + { + stop(importer); + MultimasterReplication.deleteDomain(baseDN); + } + finally + { + remove(replicationServer); + } + } + + /** + * A change which can not be replayed while the export streams must leave the session to + * the export, and be delivered again once the export is over. + *

+ * The attempts in place are spent - the backend is live, an export takes nothing away - and + * the change is given back and asked for again, as it is when nothing else is going on: what + * waits is the session restart that takes. The restart stands as a request for as long as + * the export runs, and the state checkpointer, which holds its own restarts back for the + * same reason, runs it when the export is over. Without the hold the replay thread stops the + * broker the exporter publishes over: the export ends on the entries which had been + * published, with {@code ERR_INIT_RS_DISCONNECTION_DURING_EXPORT}, the rest never reaches + * the importer, and the importer has to be initialized again. + */ + @Test(timeOut = 120_000) + public void aReplayWhichFailsDuringTheExportLeavesTheSessionToTheExport() 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 foldedInto = addPersonEntry("folded"); + addEntriesWorthMoreThanTheWindow(); + final long exportedEntries = countEntriesOfTheDomain(); + + /* + * The change goes through the replication server, which is what has it to deliver again + * once the session has been restarted for it; the replay queue of the domain is the + * test's, so the change is replayed when the test says, which is during the export. + */ + final CSN csn = gen.newCSN(); + importer.publish(new ModifyMsg(csn, DN.valueOf("cn=movedAway," + EXAMPLE_DN), + generatemods("description", "replayed during the export"), entryUUID)); + final LDAPUpdateMsg delivered = awaitDelivery(csn, 30_000, "the change was not delivered"); + + startExport(); + final List held = receiveEntryMsgsWithoutAcknowledging(INIT_WINDOW); + assertTrue(domain.ieRunning(), "the export is not being processed"); + + // Replayed while the exporter waits for an acknowledgement: every attempt in place ends on + // an entryUUID search which does not run, and the change is given back. + ShortCircuitPlugin.registerShortCircuit( + OperationType.SEARCH, "PreParse", ResultCode.UNAVAILABLE.intValue()); + try + { + replay(delivered); + assertTrue(ShortCircuitPlugin.getShortCircuitCount(OperationType.SEARCH, "PreParse") + >= LDAPReplicationDomain.IN_PLACE_REPLAY_ATTEMPTS, + "every attempt in place must have made its search: the backend is live while the" + + " export runs, so nothing holds the replay off"); + } + finally + { + ShortCircuitPlugin.deregisterShortCircuit(OperationType.SEARCH, "PreParse"); + } + + assertThat(errorLogRecordsOf(WARN_REPLAY_RETRYING_CHANGE.ordinal(), csn)) + .as("the change was not asked for again: an export is not a total update into this" + + " replica, whose state would cover the change once it is loaded") + .isNotEmpty(); + + /* + * A second change failing during the same export is warned about by the count the next + * warning carries and not by a line of its own, and the line which says its restart is + * held is folded with the warning it qualifies: a backend which fails every delivery of + * a long export would otherwise write one of them per delivery, which is the repetition + * the throttle is there to fold. + */ + final CSN folded = failAReplayOf(foldedInto); + assertThat(errorLogRecordsOf(WARN_REPLAY_RETRYING_CHANGE.ordinal(), folded)) + .as("the warning of the second change was written: the throttle must fold it into the" + + " count the next warning carries") + .isEmpty(); + assertThat(errorLogRecordsOf(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE.ordinal(), + folded)) + .as("the hold was reported for a change whose warning was folded: the line qualifies" + + " that warning, so it is folded with it") + .isEmpty(); + + /* + * The export is held across a tick of the state checkpointer, which comes for every + * restart left standing once a second: the request is standing now, and whichever thread + * comes for it while the export runs has to leave it standing. + */ + Thread.sleep(1500); + + finishExport(held, exportedEntries); + + assertThat(errorLogRecordsOf(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE.ordinal(), csn)) + .as("the restart the change was asked for again with was not reported as held") + .isNotEmpty(); + + // The change is delivered again once the export is over, and applied. + final LDAPUpdateMsg again = awaitDelivery(csn, 30_000, "the change was not delivered again" + + " once the export was over: the session restart it was asked for again with was" + + " not run"); + replay(again); + assertThat(DirectoryServer.getEntry(entry.getName()).getAllAttributes("description")) + .as("the change delivered again after the export was not applied").isNotEmpty(); + assertTrue(domain.getServerState().cover(csn), + "the change delivered again after the export was applied and not recorded"); + /* + * Read only now: until the first change was committed the ServerState could cover + * nothing newer from this server, whatever became of the second one. It was handed to + * the domain rather than published, so nothing delivers it again, and it holds the state + * back where it stands. + */ + assertFalse(domain.getServerState().cover(folded), + "the second change was recorded as replayed: its replay failed, so it must stay listed" + + " as one which is not in the data"); + } + + /** + * The changes a replay which was unwound had parked are given back while the export streams, + * and the restart they are handed back with waits for the export the way the restart a failed + * replay asks for does: the line which says so names them, and the export streams to its end. + *

+ * The replay is unwound by the ack of a change it had applied, which is the road the give-back + * of the parked changes is reached from (issue #954): what that road hands back is reported as + * given back to a replication server "which still owns it and sends it again", and during an + * export it does not send it yet. These changes are handed to the domain rather than published, + * so the redelivery the restart brings is the case above's to assert; what is asserted here is + * that the export is not cut for them and that their wait is reported. + *

+ * The same road is walked once before the export, which is the negative arm of the line: the + * restart of that give-back runs, so nothing of it waits and nothing says it does. + *

+ * The fixture is the shape {@code ParkedChangeGiveBackTest} gives that road - a change parked + * behind one whose operation is refused, and a replay unwound by an ack - over the backend + * this case exports, and with a total update running over the session. + */ + @Test(timeOut = 120_000) + public void theParkedChangesGivenBackDuringTheExportAreReportedAsHeld() throws Exception + { + final Entry waitedOn = addPersonEntry("waitedOn"); + final Entry unwoundBeforeTheExport = addPersonEntry("unwoundBefore"); + final Entry unwoundDuringTheExport = addPersonEntry("unwoundDuring"); + final CSN parkedBeforeTheExport = parkAChangeBehindABarrier(waitedOn); + + final long generationBefore = sessionGeneration(); + unwindTheReplayOf(unwoundBeforeTheExport); + /* + * Read the moment the replay returns: the restart of this give-back is run by this thread + * before the replay returns, and one left standing would be run by the state checkpointer + * within its tick, where awaitConnected() below could not tell the two apart - it would + * find the session up either way, restarted or never stopped. + */ + assertThat(sessionGeneration()) + .as("the thread which gave the parked change back did not run the restart it asked" + + " for: no total update is being processed, so nothing holds it") + .isGreaterThan(generationBefore); + + assertThat(errorLogRecordsOf( + NOTE_REPLAY_PARKED_CHANGE_GIVEN_BACK.ordinal(), parkedBeforeTheExport)) + .as("the change parked by the replay which was unwound was not given back") + .isNotEmpty(); + assertThat(errorLogRecordsOf( + NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE.ordinal(), parkedBeforeTheExport)) + .as("the restart of the give-back was reported as held while no total update was being" + + " processed: it ran") + .isEmpty(); + awaitConnected("the session was not brought back by the restart the give-back ran"); + + final CSN parked = parkAChangeBehind(waitedOn); + addEntriesWorthMoreThanTheWindow(); + final long exportedEntries = countEntriesOfTheDomain(); + + startExport(); + final List held = receiveEntryMsgsWithoutAcknowledging(INIT_WINDOW); + assertTrue(domain.ieRunning(), "the export is not being processed"); + + unwindTheReplayOf(unwoundDuringTheExport); + + assertThat(errorLogRecordsOf(NOTE_REPLAY_PARKED_CHANGE_GIVEN_BACK.ordinal(), parked)) + .as("the change parked by the replay which was unwound was not given back") + .isNotEmpty(); + assertThat(errorLogRecordsOf(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE.ordinal(), parked)) + .as("the restart the parked change was given back with was not reported as held: the" + + " change waits for the export with nothing said about it") + .isNotEmpty(); + + // The export was not cut by the give-back, and streams to its end. + finishExport(held, exportedEntries); + } + + /** + * Replays a change handed to the domain rather than published, whose every attempt in place + * ends on an entryUUID search which does not run: it is given back and asked for again, the + * way the published change of the case above is. + *

+ * The DN it carries is not in the data and its entryUUID is the provided entry's, which is + * what has the replay look the entry up by that UUID - the search which is short-circuited. + * + * @return the CSN of the change whose replay failed + */ + private CSN failAReplayOf(Entry entry) throws Exception + { + final CSN csn = gen.newCSN(); + // The registration counts the searches it refuses from zero, so the count below is this + // replay's own. + ShortCircuitPlugin.registerShortCircuit( + OperationType.SEARCH, "PreParse", ResultCode.UNAVAILABLE.intValue()); + try + { + replayHandedOver(new ModifyMsg(csn, DN.valueOf("cn=alsoMovedAway," + EXAMPLE_DN), + generatemods("description", "replayed during the export as well"), + getEntryUUID(entry.getName()))); + // Read before the short circuit is deregistered, which forgets the count with it. + assertThat(ShortCircuitPlugin.getShortCircuitCount(OperationType.SEARCH, "PreParse")) + .as("the replay of this change must have failed: every attempt in place makes the" + + " entryUUID search which does not run") + .isGreaterThanOrEqualTo(LDAPReplicationDomain.IN_PLACE_REPLAY_ATTEMPTS); + } + finally + { + ShortCircuitPlugin.deregisterShortCircuit(OperationType.SEARCH, "PreParse"); + } + return csn; + } + + + + private static Entry addPersonEntry(String cn) throws Exception + { + return TestCaseUtils.addEntry( + "dn: cn=" + cn + "," + EXAMPLE_DN, + "objectClass: top", + "objectClass: person", + "cn: " + cn, + "sn: " + cn); + } + + /** + * Replays, on the thread of this test, a change whose operation is refused - it stays listed + * as one which is not in the data, and every change on that entry waits for it - and then a + * change which is parked behind it. + * + * @return the CSN of the change which is left parked + */ + private CSN parkAChangeBehindABarrier(Entry entry) throws Exception + { + final CSN failing = gen.newCSN(); + replayHandedOver(new ModifyMsgWhoseOperationRefusesAControl(failing, entry.getName(), + generatemods("description", "the replay of this change fails"), + getEntryUUID(entry.getName()))); + assertFalse(domain.getServerState().cover(failing), + "the change whose replay fails must stay listed as one which is not in the data"); + awaitConnected("the session was not brought back for the change whose replay failed"); + return parkAChangeBehind(entry); + } + + /** + * Replays, on the thread of this test, a change on an entry whose barrier is still missing + * from the data: it is parked as waiting for that one and owned by this thread from then on, + * since nothing hands a parked change out again while what it waits for is missing + * (issue #954). + * + * @return the CSN of the change which is left parked + */ + private CSN parkAChangeBehind(Entry entry) throws Exception + { + final CSN parked = gen.newCSN(); + replayHandedOver(new ModifyMsg(parked, entry.getName(), + generatemods("description", "the change which waits for the one which failed"), + getEntryUUID(entry.getName()))); + assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 1, + "a change which waits for one that is not in the data must be parked"); + return parked; + } + + /** + * Reads the generation of the session of the domain, which every stop and start of it bumps, + * under {@code serviceStateLock}, as {@code getSessionGeneration()} asks. + */ + private long sessionGeneration() throws Exception + { + final Field lockField = ReplicationDomain.class.getDeclaredField("serviceStateLock"); + lockField.setAccessible(true); + final Method getSessionGeneration = + ReplicationDomain.class.getDeclaredMethod("getSessionGeneration"); + getSessionGeneration.setAccessible(true); + synchronized (lockField.get(domain)) + { + return (Long) getSessionGeneration.invoke(domain); + } + } + + /** Waits for the session of the domain to be up, which a restart leaves it. */ + private void awaitConnected(String orElse) throws Exception + { + final long deadline = System.currentTimeMillis() + SESSION_BACK_BOUND_IN_MS; + while (!domain.isConnected()) + { + assertTrue(System.currentTimeMillis() < deadline, orElse); + Thread.sleep(50); + } + } + + /** + * Replays, on the thread of this test, a change whose ack runs out of memory once it is + * applied: the replay is unwound with that change in the data and owned by nobody, so what + * the give-back on the way out of {@code replay()} has to hand back is what this thread + * parked. + */ + private void unwindTheReplayOf(Entry entry) throws Exception + { + try + { + replayHandedOver(new ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied(gen.newCSN(), + entry.getName(), generatemods("description", "the replay of this change is unwound"), + getEntryUUID(entry.getName()))); + } + catch (OutOfMemoryError unwound) + { + // The error is the fixture's own, and this is the thread it would have ended. + return; + } + throw new AssertionError( + "the replay was not unwound: the ack of the delivery must run out of memory"); + } + + /** + * Hands a change to the domain rather than publishing it, and replays it on the thread of + * this test: what the replication server has to deliver again is the published change of the + * case above, and these are the changes whose give-back this one is about. + */ + private void replayHandedOver(UpdateMsg msg) throws Exception + { + domain.processUpdate(msg); + replay(queue.take().getUpdateMessage()); + } + + /** Adds entries whose export outgrows the initialization window, so that the exporter waits. */ + private void addEntriesWorthMoreThanTheWindow() throws Exception + { + assertThat(BULK_ENTRIES * BULK_ENTRY_BYTES) + .as("the data must outgrow the window for the exporter to wait for an acknowledgement") + .isGreaterThan((INIT_WINDOW + 1) * ENTRY_MSG_BYTES); + final char[] padding = new char[BULK_ENTRY_BYTES]; + Arrays.fill(padding, 'x'); + for (int i = 0; i < BULK_ENTRIES; i++) + { + TestCaseUtils.addEntry( + "dn: cn=bulk" + i + "," + EXAMPLE_DN, + "objectClass: top", + "objectClass: person", + "cn: bulk" + i, + "sn: bulk" + i, + "description: " + new String(padding)); + } + } + + private long countEntriesOfTheDomain() throws Exception + { + return getServerContext().getBackendConfigManager().findLocalBackendForEntry(baseDN) + .getNumberOfEntriesInBaseDN(baseDN); + } + + /** + * Has the importer ask this replica for a total update, and returns once the export has + * begun: the {@code InitializeTargetMsg} which starts it has arrived. + */ + private void startExport() throws Exception + { + // The export is refused while this replica does not see the importer in its topology. + final long deadline = System.currentTimeMillis() + 30_000; + while (!domain.getReplicaInfos().containsKey(IMPORTER_ID)) + { + assertTrue(System.currentTimeMillis() < deadline, + "the domain did not see the importer in its topology"); + Thread.sleep(20); + } + importer.publish(new InitializeRequestMsg(baseDN, IMPORTER_ID, DS_ID, INIT_WINDOW)); + // The exporter waits for the importer to be in the full update status before it streams. + importer.signalStatusChange(ServerStatus.FULL_UPDATE_STATUS); + final ReplicationMsg msg = receiveTotalUpdateMsg(30_000); + assertThat(msg).as("the total update did not begin").isInstanceOf(InitializeTargetMsg.class); + } + + /** + * Receives entry messages up to the window and acknowledges none of them: the exporter + * publishes no more than the window ahead of the last acknowledgement, so its next entry + * message waits for one from now on. + */ + private List receiveEntryMsgsWithoutAcknowledging(int window) throws Exception + { + final List received = new ArrayList<>(); + while (received.size() < window) + { + final ReplicationMsg msg = receiveTotalUpdateMsg(30_000); + assertThat(msg).as("the export did not stream up to the window").isInstanceOf(EntryMsg.class); + received.add((EntryMsg) msg); + } + return received; + } + + /** + * Acknowledges what arrived while the export was held and everything after it as it + * arrives, up to the {@code DoneMsg}, and checks that every entry of the domain arrived. A + * total update which was cut streams no further: the rest of its entries never arrives, or + * an {@code ErrorMsg} arrives in their place, and either fails here. The importer then + * leaves the full update status, which the exporter waits for before it releases its + * context - and it leaves it whatever happened, or the export never ends. + *

+ * An export which does not end all the same is printed rather than asserted: it is what + * the assertions which follow this call wait for - the change the restart brings back once + * the export is over - and a throw out of this {@code finally} would replace the failure + * of the stream above it, which is the one worth reading. + */ + private void finishExport(List held, long exportedEntries) throws Exception + { + try + { + final StringBuilder ldif = new StringBuilder(); + int lastMsgId = 0; + for (EntryMsg entryMsg : held) + { + ldif.append(new String(entryMsg.getEntryBytes(), UTF_8)); + lastMsgId = entryMsg.getMsgId(); + } + importer.publish(new InitializeRcvAckMsg(IMPORTER_ID, DS_ID, lastMsgId)); + final int heldAt = lastMsgId; + while (true) + { + final ReplicationMsg msg = receiveTotalUpdateMsg(60_000); + if (msg instanceof DoneMsg) + { + break; + } + assertThat(msg).as("the export was cut instead of streaming to its end") + .isInstanceOf(EntryMsg.class); + final EntryMsg entryMsg = (EntryMsg) msg; + ldif.append(new String(entryMsg.getEntryBytes(), UTF_8)); + lastMsgId = entryMsg.getMsgId(); + importer.publish(new InitializeRcvAckMsg(IMPORTER_ID, DS_ID, lastMsgId)); + } + assertThat(lastMsgId).as("the export did not stream past the window it was held at") + .isGreaterThan(heldAt); + assertThat(countEntries(ldif)).as("the export did not stream every entry of the domain") + .isEqualTo(exportedEntries); + } + finally + { + leaveTheFullUpdateStatus(); + final long deadline = System.currentTimeMillis() + EXPORT_END_BOUND_IN_MS; + while (domain.ieRunning() && System.currentTimeMillis() < deadline) + { + Thread.sleep(50); + } + if (domain.ieRunning()) + { + System.err.println("the export of " + baseDN + " did not end within " + + EXPORT_END_BOUND_IN_MS + " ms"); + } + } + } + + /** Counts the entries of an LDIF stream by the blank line which separates them. */ + private static long countEntries(CharSequence ldif) + { + long count = 0; + for (int i = ldif.length() - 1; i > 0; i--) + { + if (ldif.charAt(i) == '\n' && ldif.charAt(i - 1) == '\n') + { + count++; + } + } + return count; + } + + /** + * The importer reconnects once its import is over - it comes back with the generation ID of + * the data it loaded, which is the one it was opened with here - and the exporter waits for + * the importer to leave the full update status before it releases its context. + */ + private void leaveTheFullUpdateStatus() + { + importer.reStart(true); + } + + /** + * Receives the next message of the total update on the importer: the updates of this + * replica's own and the topology are not it. + */ + private ReplicationMsg receiveTotalUpdateMsg(long timeoutMs) throws Exception + { + final long deadline = System.currentTimeMillis() + timeoutMs; + final List others = new ArrayList<>(); + while (System.currentTimeMillis() < deadline) + { + final ReplicationMsg msg; + try + { + msg = importer.receive(); + } + catch (SocketTimeoutException e) + { + continue; + } + if (msg instanceof InitializeTargetMsg || msg instanceof EntryMsg + || msg instanceof DoneMsg) + { + return msg; + } + if (msg instanceof ErrorMsg) + { + Assert.fail("the total update was given up: " + ((ErrorMsg) msg).getDetails()); + } + others.add(msg); + } + Assert.fail("nothing of the total update arrived within " + timeoutMs + " ms; received " + + others); + return null; + } + + /** + * 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 published was delivered"); + return msg; + } + + /** The records of the error log which carry the provided message id and the provided CSN. */ + private static List errorLogRecordsOf(int msgId, CSN csn) + { + final List records = new ArrayList<>(); + for (String record : TestCaseUtils.ERROR_TEXT_WRITER.getMessages()) + { + if (record.contains("msgID=" + msgId) && record.contains(csn.toString())) + { + records.add(record); + } + } + return records; + } + + private void replay(LDAPUpdateMsg ldapUpdate) + { + domain.markInProgress(ldapUpdate); + domain.replay(ldapUpdate, SHUTDOWN); + } +} 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 538c7e1060..e921beb02b 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 @@ -830,111 +830,6 @@ public void aDomainDisabledBeforeTheImportItAskedForStartsRefusesTheImport() thr .hasSize(listenerDeathsBefore); } - /** - * A session restart decided while a total update out of this replica is running stops the - * session that export streams over (issue #1041). - *

- * What the restart must leave alone is a total update into this replica: the data it is - * about to replace is read over the session, and the import is the thread the stop waits - * for. An export is not that: it streams out of a backend nothing is taking away, on a - * thread of its own, and a session stopped under it is the cut it reports to whoever asked - * for the total update - the same cut every other stop of the session is. The claim the - * restart makes for the import is not made here, and the session is stopped as it was - * before the claim. - *

- * The export holds the context by standing where it waits for its target to report the - * start of the total update: the target is a broker of this test, and reports nothing. - */ - @Test(timeOut = 120_000) - public void aRestartDecidedWhileAnExportRunsStopsTheSessionItStreamsOver() 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()); - waitUntil(() -> domain.getReplicaInfos().containsKey(EXPORTER_ID), - "the exporter is not in the replicas of the domain: nothing to export into"); - - final AtomicReference exportFailure = new AtomicReference<>(); - final Thread export = new Thread(() -> { - try - { - domain.initializeRemote(EXPORTER_ID, null); - } - catch (Throwable t) - { - exportFailure.set(t); - } - }, "export of " + EXAMPLE_DN); - - // The restart is held after its decision, before the stop: what the case reads is the - // decision the export was found by, not the session which is down a moment later. - final CountDownLatch stopHeld = new CountDownLatch(1); - final CountDownLatch releaseStop = new CountDownLatch(1); - domain.setServiceStopHook(() -> { - stopHeld.countDown(); - awaitUninterruptibly(releaseStop); - }); - final CSN csn = gen.newCSN(); - final AtomicReference replayFailure = new AtomicReference<>(); - final Thread replay = new Thread(() -> { - try - { - replayMsg(new ModifyMsg(csn, DN.valueOf("cn=movedAway," + EXAMPLE_DN), - generatemods("description", "replayed while the export was running"), entryUUID)); - } - catch (Throwable t) - { - replayFailure.set(t); - } - }, "replay of " + csn); - try - { - export.start(); - waitUntil(() -> domain.ieRunning() || exportFailure.get() != null, - "the export did not claim the import context"); - assertNull(exportFailure.get(), - "the export failed before it claimed the context: " + exportFailure.get()); - - // A change whose entryUUID search never runs spends its attempts in place and asks - // for the session to be restarted, the way it does in the case above. - ShortCircuitPlugin.registerShortCircuit( - OperationType.SEARCH, "PreParse", ResultCode.UNAVAILABLE.intValue()); - try - { - replay.start(); - assertTrue(stopHeld.await(30, TimeUnit.SECONDS), - "the restart left the session to the export: an export is not the owner a total" - + " update into this replica is"); - } - finally - { - ShortCircuitPlugin.deregisterShortCircuit(OperationType.SEARCH, "PreParse"); - } - assertTrue(domain.ieRunning(), "the export ended before the restart was decided"); - } - finally - { - releaseStop.countDown(); - domain.setServiceStopHook(null); - } - - replay.join(60_000); - assertFalse(replay.isAlive(), "the restart did not end"); - assertNull(replayFailure.get(), "the replay failed: " + replayFailure.get()); - export.join(60_000); - assertFalse(export.isAlive(), "the export did not end once the session it streams over" - + " was stopped"); - assertThat(exportFailure.get()) - .as("the export was not told that the session it streams over was cut") - .isInstanceOf(DirectoryException.class); - waitUntil(domain::isConnected, "the session was not started back after the restart"); - assertFalse(domain.ieRunning(), "the export which was cut left its context claimed"); - } - /** * Has the exporter start a total update into this replica, and returns once the backend * of the domain is deregistered for it: from then on the import is reading the session, diff --git a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/TestSynchronousReplayQueue.java b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/TestSynchronousReplayQueue.java index dd908c0623..a921f38b0d 100644 --- a/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/TestSynchronousReplayQueue.java +++ b/opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/TestSynchronousReplayQueue.java @@ -13,13 +13,15 @@ * * Copyright 2009 Sun Microsystems, Inc. * Portions copyright 2013-2015 ForgeRock AS. + * Portions Copyright 2026 3A Systems, LLC. */ package org.opends.server.replication.plugin; import java.util.Collection; +import java.util.Deque; import java.util.Iterator; -import java.util.LinkedList; import java.util.concurrent.BlockingQueue; +import java.util.concurrent.ConcurrentLinkedDeque; import java.util.concurrent.TimeUnit; /** @@ -30,7 +32,13 @@ */ public class TestSynchronousReplayQueue implements BlockingQueue { - private LinkedList list = new LinkedList<>(); + /** + * Written by the listener thread of the domain - a change the replication server delivers + * is offered here - and read by the thread of the test, which replays it: the two share no + * lock, so the deque has to be safe for that hand-off. {@code take()} is still synchronous - + * it throws when nothing was offered - which is what makes the queue a test one. + */ + private final Deque list = new ConcurrentLinkedDeque<>(); @Override public boolean add(UpdateToReplay e)