/* * 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.util.ArrayList; 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.protocol.DeleteMsg; import org.opends.server.replication.protocol.DoneMsg; import org.opends.server.replication.protocol.EntryMsg; 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.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.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 target of a total update. *
* The import of a total update streams over the session of the domain, on its listener * thread, and the backend it replaces is deregistered for the length of it. A change which * was queued for replay before the {@code InitializeTargetMsg} arrived is replayed into no * backend: whatever such a replay decides is about to be overwritten by the import, and the * 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). *
* 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. *
* The {@code timeOut} each case declares is what it is expected to take at the most; it is
* not what bounds it. {@code TestListener} sets the timeout of every test method from the
* {@code org.opends.test.timeout} property, ten minutes under Maven and none outside it.
*/
@SuppressWarnings("javadoc")
public class ReplayDuringImportTest extends ReplicationTestCase
{
/**
* The memory backend of {@code o=test} loses its data when it is disabled and enabled
* back, which is what an import does to the backend it replaces: a total update needs a
* backend which keeps what was imported into it.
*/
private static final String EXAMPLE_DN = "dc=example,dc=com";
private static final int RS_ID = 611;
private static final int DS_ID = 1;
private static final int EXPORTER_ID = 2;
private static final int INIT_WINDOW = 100;
private static final AtomicBoolean SHUTDOWN = new AtomicBoolean(false);
/** An entry of the exporter's data, and its entryUUID. */
private static final String IMPORTED_ENTRY_DN = "cn=imported,ou=People," + EXAMPLE_DN;
private static final String IMPORTED_ENTRY_UUID = "21111111-1111-1111-1111-111111111113";
private DN baseDN;
private ReplicationServer replicationServer;
private LDAPReplicationDomain domain;
private TestSynchronousReplayQueue queue;
private ReplicationBroker exporter;
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, "replayDuringImportTestDb", 0, RS_ID, 0, 100, new TreeSet
* The change is given back at the top of its first attempt: the data it would be applied
* to is being replaced, so nothing is attempted into the backend the import took away,
* nothing is reported, and the session is left to the import - which streams every entry
* to its end. Without the hold-off the operation is refused with NO_SUCH_OBJECT - nothing
* serves the base DN - and the entryUUID search conflict resolution reads the data with
* can not run either: the attempts in place are spent into no backend and the exit
* reports the change; without the owner the total update is, the session is then
* restarted for the change to be delivered again, which stops the broker the import is
* reading, and the import ends on the entries which had arrived with nothing to say it.
*/
@Test(timeOut = 120_000)
public void aReplayDuringTheImportLeavesTheSessionToTheImport() 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 String[] exported = exportedEntries();
startImportInto(exported.length);
// Queued before the InitializeTargetMsg arrived, replayed into no backend.
final CSN csn = gen.newCSN();
replayMsg(new ModifyMsg(csn, DN.valueOf("cn=movedAway," + EXAMPLE_DN),
generatemods("description", "replayed during the import"), entryUUID));
finishImport(exported);
for (String ldif : exported)
{
final DN dn = dnOf(ldif);
assertTrue(entryExists(dn), "the import ended before " + dn
+ " arrived: the session it streams over was stopped from under it");
}
/*
* The two roads which leave the session to the import are told apart here: the
* hold-off gives the change back before an attempt is made, the guard on the restart
* after the attempts are spent. The exhaustion exit is the one thing the first road
* leaves no record of.
*/
assertThat(errorLogRecordsOf(ERR_ERROR_REPLAYING_OPERATION.ordinal(), csn))
.as("the change was attempted into no backend instead of being given back at once")
.isEmpty();
assertThat(errorLogRecordsOf(WARN_REPLAY_RETRYING_CHANGE.ordinal(), csn))
.as("the change was asked for again, which restarts the session the import streams over")
.isEmpty();
}
/**
* A change given back while the import ran must not hold the ServerState back once the
* import has replaced the data.
*
* A change which is given back stays listed as pending and uncommitted - that is what
* has the replication server send it again - and a commit advances the ServerState no
* further than the oldest uncommitted change. The state the import loads is the
* exporter's, which covers the change already, so nothing sends it again: left listed,
* it would stop the ServerState of this replica for good.
*/
@Test(timeOut = 120_000)
public void aChangeGivenBackDuringTheImportDoesNotHoldTheServerStateBack() 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 String[] exported = exportedEntries();
startImportInto(exported.length);
replayMsg(new ModifyMsg(gen.newCSN(), DN.valueOf("cn=movedAway," + EXAMPLE_DN),
generatemods("description", "replayed during the import"), entryUUID));
finishImport(exported);
// A change on an entry the import brought, replayed once the import is over.
final DN importedDN = DN.valueOf(IMPORTED_ENTRY_DN);
final CSN csn = gen.newCSN();
replayMsg(new ModifyMsg(csn, importedDN,
generatemods("description", "replayed after the import"), IMPORTED_ENTRY_UUID));
assertThat(DirectoryServer.getEntry(importedDN).getAllAttributes("description"))
.as("a change replayed after the import was not applied").isNotEmpty();
assertTrue(domain.getServerState().cover(csn), "a change applied after the import was not"
+ " recorded: the change given back during the import is still listed and holds the"
+ " ServerState back");
}
/**
* A total update this replica asked for owns the session from the request, not from the
* first entry: the {@code InitializeTargetMsg} which answers the request arrives over
* that session, and a restart made while the answer is on its way loses it.
*
* The backend is live for the length of the request - nothing has been taken away yet -
* so the change is attempted, every attempt ends on an entryUUID search which does not
* run, and the exhaustion exit reports it: what is refused is the restart which would
* have followed, and the retry warning which goes with it. The exporter then answers the
* request, and its entries stream to their end over the session which was left alone.
*/
@Test(timeOut = 120_000)
public void aRequestOnItsWayOwnsTheSessionTheAnswerArrivesOver() 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 String[] exported = exportedEntries();
// The request is out, and the exporter holds it until the change below has been replayed.
domain.initializeFromRemote(EXPORTER_ID, null);
assertNotNull(waitForSpecificMsg(exporter, InitializeRequestMsg.class));
final CSN csn = gen.newCSN();
ShortCircuitPlugin.registerShortCircuit(
OperationType.SEARCH, "PreParse", ResultCode.UNAVAILABLE.intValue());
try
{
replayMsg(new ModifyMsg(csn, DN.valueOf("cn=movedAway," + EXAMPLE_DN),
generatemods("description", "replayed while the request was on its way"), entryUUID));
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"
+ " request is on its way, so nothing holds the replay off");
}
finally
{
ShortCircuitPlugin.deregisterShortCircuit(OperationType.SEARCH, "PreParse");
}
assertThat(errorLogRecordsOf(ERR_ERROR_REPLAYING_OPERATION.ordinal(), csn))
.as("the attempts in place were spent, which the exhaustion exit reports").isNotEmpty();
assertThat(errorLogRecordsOf(WARN_REPLAY_RETRYING_CHANGE.ordinal(), csn))
.as("the change was asked for again, which restarts the session the answer to the"
+ " request arrives over")
.isEmpty();
assertTrue(domain.isConnected(), "the session the request was made over was stopped");
answerImportRequest(exported.length);
finishImport(exported);
for (String ldif : exported)
{
final DN dn = dnOf(ldif);
assertTrue(entryExists(dn), "the import ended before " + dn
+ " arrived: the answer to the request was lost with the session it was made over");
}
}
/**
* 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.
*
* 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,
* and the backend is live meanwhile, so the change which is parked and the replay which
* is unwound run as they would on any domain. The domain going away, or being disabled,
* forgets its pending changes a moment after it takes the session and clears every
* request and every count on its way, so a give-back on that road is a race with the
* forgetting and leaves nothing to read.
*
* The replay is unwound on the thread of this test - it applied its change, and the ack
* of its delivery runs out of memory - so the parked change is this thread's to give
* back, and the error which ends a replay thread is caught here instead.
*/
@Test(timeOut = 120_000)
public void aParkedChangeGivenBackWhileTheRequestIsOnItsWayIsNotAskedForAgain() 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 String[] exported = exportedEntries();
// The request is out, and the exporter holds it until the give-back below has run.
domain.initializeFromRemote(EXPORTER_ID, null);
assertNotNull(waitForSpecificMsg(exporter, InitializeRequestMsg.class));
/*
* The barrier: a change whose replay fails stays listed and uncommitted - the attempts
* in place end on an entryUUID search which does not run, the way they do in the case
* above - and stays among the changes the newer ones are checked against, so a change
* which follows it on the same entry has to wait for it. The restart which would have
* followed is refused, the total update owning the session, and the search is let
* through again before anything below reads a monitor.
*/
final DN movedAway = DN.valueOf("cn=movedAway," + EXAMPLE_DN);
final CSN failing = gen.newCSN();
ShortCircuitPlugin.registerShortCircuit(
OperationType.SEARCH, "PreParse", ResultCode.UNAVAILABLE.intValue());
try
{
replayMsg(new ModifyMsg(failing, movedAway,
generatemods("description", "the replay of this change fails"), entryUUID));
}
finally
{
ShortCircuitPlugin.deregisterShortCircuit(OperationType.SEARCH, "PreParse");
}
assertFalse(domain.getServerState().cover(failing),
"the change whose replay fails must stay listed as one which is not in the data");
// Parked as waiting for it by this thread, which owns it from here on.
final CSN parked = gen.newCSN();
replayMsg(new ModifyMsg(parked, movedAway,
generatemods("description", "the change which was parked as a dependency"), entryUUID));
assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 1,
"a change which waits for one that is not in the data must be parked");
/*
* The replay which is unwound while this thread still holds the parked change: its own
* change is applied and committed, so the give-back on the way out finds the parked
* change alone. The count is read once the parked change is listed, since a parked
* change publishes no ack and is not counted until the delivery which replays it is.
*/
final long processed = getMonitorAttrValue(baseDN, "replayed-updates");
final CSN unwound = gen.newCSN();
try
{
replayMsg(new ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied(unwound, entry.getName(),
generatemods("description", "the replay of this change is unwound once it is applied"),
entryUUID));
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(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");
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")
.isEmpty();
assertTrue(domain.isConnected(), "the session the answer to the request arrives over was stopped");
answerImportRequest(exported.length);
finishImport(exported);
for (String ldif : exported)
{
final DN dn = dnOf(ldif);
assertTrue(entryExists(dn), "the import ended before " + dn
+ " arrived: the answer to the request was lost with the session it was made over");
}
}
/**
* 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.
*
* 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
* spent: a restart which ran would have spent it, and would have left the session it
* stopped down.
*/
@Test(timeOut = 120_000)
public void aRequestWhichStoodWhileTheImportRanIsNotRunOnceItIsOver() throws Exception
{
final String[] exported = exportedEntries();
startImportInto(exported.length);
domain.requestSessionRestart();
domain.failNextSessionRestarts(1);
try
{
finishImport(exported);
// Two ticks of the checkpointer: a request standing when the import ends is run on the first.
Thread.sleep(2000);
assertEquals(domain.getSessionRestartFailuresLeft(), 1, "the request which stood while"
+ " the import ran was run against the session started back at its end");
assertTrue(domain.isConnected(), "the session started back at the end of the import"
+ " was stopped for a request made before it");
}
finally
{
domain.failNextSessionRestarts(0);
}
}
/**
* A total update forgets the deliveries which were folded into no warning, along with
* the changes they were deliveries of (issue #942).
*
* The changes listed as pending do not outlive the ServerState the import replaces, and
* the recovery from a failed replay goes with them - the session restart backoff, and the
* count the next warning about a change being asked for again says it stands for. The
* first warning over the imported data must not count the deliveries of a change which
* is not listed anymore.
*
* Nothing sends a change of this test again - the exporter never had it - so the count
* is fed by two changes failing within one interval rather than by one change delivered
* twice: the first is warned about, the second is folded into no warning. The changes
* are deletes: a short circuit on the modifies would be tripped by the ServerState being
* saved to the base entry and by the import disabling the backend it replaces.
*/
@Test(timeOut = 120_000)
public void aWarningAfterTheImportDoesNotCountTheDeliveriesBefore() throws Exception
{
final Entry warnedAbout = TestCaseUtils.addEntry(
"dn: cn=warnedAbout," + EXAMPLE_DN,
"objectClass: top",
"objectClass: person",
"cn: warnedAbout",
"sn: warnedAbout");
final Entry folded = TestCaseUtils.addEntry(
"dn: cn=folded," + EXAMPLE_DN,
"objectClass: top",
"objectClass: person",
"cn: folded",
"sn: folded");
final String warnedAboutUUID = getEntryUUID(warnedAbout.getName());
final String foldedUUID = getEntryUUID(folded.getName());
final String[] exported = exportedEntries();
ShortCircuitPlugin.registerShortCircuit(
OperationType.DELETE, "PreParse", ResultCode.UNAVAILABLE.intValue());
try
{
replayMsg(new DeleteMsg(warnedAbout.getName(), gen.newCSN(), warnedAboutUUID));
replayMsg(new DeleteMsg(folded.getName(), gen.newCSN(), foldedUUID));
startImportInto(exported.length);
finishImport(exported);
/*
* Only the timestamp of the throttle is put back, so that the failure over the
* imported data is warned about straight away: the count is the domain's to keep or
* to forget.
*/
domain.resetReplayRetryWarningThrottle();
final CSN csn = gen.newCSN();
replayMsg(new DeleteMsg(DN.valueOf(IMPORTED_ENTRY_DN), csn, IMPORTED_ENTRY_UUID));
final List