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