| | |
| | | * 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 |
| | |
| | | |
| | | /** |
| | | * 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, |
| | |
| | | * 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, |
| | |
| | | "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"); |
| | | |
| | |
| | | } |
| | | |
| | | /** |
| | | * 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 |
| | |
| | | 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"); |
| | | } |
| | | } |