| | |
| | | import org.opends.server.types.Operation; |
| | | import org.opends.server.types.OperationType; |
| | | import org.opends.server.types.RawModification; |
| | | import org.opends.server.util.StaticUtils; |
| | | import org.opends.server.util.TestTimer; |
| | | import org.opends.server.util.TestTimer.CallableVoid; |
| | | import org.opends.server.util.TimeThread; |
| | |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Test case for [Issue 954]: a change parked as waiting for another one is given back |
| | | * when the replay which parked it is unwound, and the session is restarted for it. |
| | | * <p> |
| | | * A change which waits for another one is parked and stays owned by the replay thread |
| | | * which parked it, while that thread goes on to the changes which follow: it is handed |
| | | * out again by {@code getNextUpdate()}, which every replay loop of this domain runs once |
| | | * it is done, so it is replayed by whichever thread clears the change it was waiting for. |
| | | * A replay which is unwound leaves the thread which parked it without that road - it |
| | | * takes the next delivery off the shared queue instead, and never comes back to the |
| | | * change it parked - and every redelivery of a change a replay thread owns is refused as |
| | | * a duplicate. On a domain which then goes quiet that change is where this replica's |
| | | * ServerState, and every change behind it from every master, stops. |
| | | * <p> |
| | | * The replay which is unwound here is one which applied its change: the change it was |
| | | * replaying is committed and owns nothing anymore by the time the give-back on the way |
| | | * out of {@code replay()} runs, so the restart that give-back asks for is the one thing |
| | | * which has the parked change delivered again - a replay which failed would have asked |
| | | * for the same restart on the road of its own change. The parked change travels the |
| | | * replication server, and nothing but a new session brings it back. |
| | | */ |
| | | @Test |
| | | public void aChangeParkedByAnUnwoundReplayIsDeliveredAgain() throws Exception |
| | | { |
| | | testSetUp("aChangeParkedByAnUnwoundReplayIsDeliveredAgain"); |
| | | logger.error(LocalizableMessage.raw( |
| | | "Starting replication test : aChangeParkedByAnUnwoundReplayIsDeliveredAgain")); |
| | | |
| | | final DN waitedOn = addEntryForChange("user.954.1"); |
| | | final String waitedOnUUID = getEntry(waitedOn, 1, true).parseAttribute("entryuuid").asString(); |
| | | final DN other = addEntryForChange("user.954.2"); |
| | | final String otherUUID = getEntry(other, 1, true).parseAttribute("entryuuid").asString(); |
| | | |
| | | final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null); |
| | | domain.resetUnreplayedChangeAlertThrottle(); |
| | | final long inProgress = getMonitorAttrValue(baseDN, "changes-in-progress-size"); |
| | | assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 0, |
| | | "no change of this domain is waiting for another one when this test starts"); |
| | | |
| | | // The replica the changes come from: one of its own, so that the CSNs of this test are |
| | | // never covered by the ServerState another one left behind. |
| | | final int serverId = 30; |
| | | final CSNGenerator gen = new CSNGenerator(serverId, TimeThread.getTime()); |
| | | final CSN failing = gen.newCSN(); |
| | | final CSN parked = gen.newCSN(); |
| | | final CSN unwound = gen.newCSN(); |
| | | final List<Modification> failingMods = generatemods("description", "the replay of this change fails"); |
| | | final String parkedDescription = "the change which was parked as a dependency"; |
| | | final List<Modification> parkedMods = generatemods("description", parkedDescription); |
| | | final List<Modification> unwoundMods = |
| | | generatemods("description", "the replay of this change is unwound once it is applied"); |
| | | |
| | | /* |
| | | * The change whose replay fails is the barrier the parked change waits behind, so its |
| | | * budget must not be spent while this test is setting up: it is shortened once the |
| | | * change which was parked is back, and put back in the finally below. |
| | | */ |
| | | setReplayGiveUpDelay("unlimited"); |
| | | try |
| | | { |
| | | /* |
| | | * One replay thread, so that the change which is parked and the replay which is |
| | | * unwound after it are the same thread's: a parked change is left owned by the thread |
| | | * which parked it, and this is about a thread which does not come back to it. |
| | | */ |
| | | setNumUpdateReplayThreads(1); |
| | | try |
| | | { |
| | | /* |
| | | * The change which is parked is published to the replication server rather than |
| | | * handed to the domain: a change which travelled a session is one the replication |
| | | * server sends again over the next session of this domain, and over nothing else - |
| | | * which is what the restart the give-back asks for has to be pinned against. |
| | | */ |
| | | final ReplicationBroker broker = |
| | | openReplicationSession(baseDN, serverId, 100, replServerPort, 1000); |
| | | try |
| | | { |
| | | /* |
| | | * A change whose replay failed stays listed and uncommitted - it is what holds this |
| | | * domain's ServerState back - 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. It has |
| | | * to be a change which is asked for again rather than stepped over: one whose |
| | | * operation is built and then refused, since #928 has a modify whose entry DN does |
| | | * not parse reported once and recorded as replayed. |
| | | */ |
| | | deliverUntilMonitorReaches(domain, "changes-in-progress-size", inProgress + 1, |
| | | () -> new ModifyMsgWhoseOperationRefusesAControl(failing, waitedOn, failingMods, waitedOnUUID), |
| | | "the change whose replay fails must be listed as one which is not in the data"); |
| | | |
| | | // The change which is parked as waiting for it by the replay thread it was given to. |
| | | broker.publish(new ModifyMsg(parked, waitedOn, parkedMods, waitedOnUUID)); |
| | | assertMonitorAttrValueEventually(baseDN, "dependent-changes-size", 1, |
| | | "a change which waits for one that is not in the data must be parked"); |
| | | |
| | | /* |
| | | * How many restarts in a row the session has been through, the barrier's among |
| | | * them: the restart the give-back asks for on the road out of a JVM which has run |
| | | * out of memory is owed no backoff, and a restart which waits the backoff out is |
| | | * the one road which moves this count. It can not move otherwise between here and |
| | | * the reading below: a replay which made it puts the count back to zero only once |
| | | * nothing is failing anymore, and the barrier keeps failing until it is given up. |
| | | */ |
| | | final int restarts = domain.getConsecutiveSessionRestarts(); |
| | | |
| | | /* |
| | | * The replay which is unwound while that same thread still holds the parked change. |
| | | * It applies its change, and the ack of its delivery is where the JVM runs out of |
| | | * memory: that is the one throw from the ack which is not caught, past every catch |
| | | * the replay itself has, and it leaves replay() with the change this thread was |
| | | * replaying committed and owned by nobody. What the give-back on the way out finds |
| | | * to give back is the parked change alone, and the restart it asks for is the one |
| | | * road which has that change delivered again: the road a failed replay takes to |
| | | * ask for its own change is not run. An Error met replaying a change does not get |
| | | * here, and neither does any other throw from the ack: both are reported and the |
| | | * replay carries on to the road the change itself decided (issue #922). |
| | | * |
| | | * It is made on another entry, so that it is replayed rather than parked in its |
| | | * turn, and it is delivered until the parked change has been given back: a delivery |
| | | * of a change which is being replayed, or was applied, is refused as the duplicate |
| | | * it is. |
| | | */ |
| | | deliverUntilMonitorReaches(domain, "dependent-changes-size", 0, |
| | | () -> new ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied(unwound, other, unwoundMods, otherUUID), |
| | | "a change parked by a replay which was unwound must be given back"); |
| | | |
| | | /* |
| | | * The thread which met the error is gone - an OutOfMemoryError ends the replay |
| | | * thread which met it (issue #923) - and it was the whole pool, so the delivery the |
| | | * restarted session brings waits in the replay queue: the pool is brought back for |
| | | * it, at any number now that the thread it had to be the same as is over. The change |
| | | * is parked again by the thread which takes it, the barrier being still there: that |
| | | * it is parked at all is what says a new session delivered it. |
| | | */ |
| | | setNumUpdateReplayThreads(2); |
| | | assertMonitorAttrValueEventually(baseDN, "dependent-changes-size", 1, |
| | | "the change which was given back must be delivered again by the session which" |
| | | + " was restarted for it"); |
| | | assertEquals(domain.getConsecutiveSessionRestarts(), restarts, |
| | | "a thread an OutOfMemoryError is ending must not wait the backoff out before" |
| | | + " it restarts the session for the changes it gave back"); |
| | | |
| | | /* |
| | | * The barrier is lifted, which lets the ServerState past it and hands the parked |
| | | * change to the thread which cleared it. |
| | | */ |
| | | giveUpOn(domain, failing, waitedOn, failingMods, waitedOnUUID); |
| | | checkEntryHasAttributeValue(waitedOn, "description", parkedDescription, 30, |
| | | "the change which was parked must be applied by the delivery which took it over"); |
| | | } |
| | | finally |
| | | { |
| | | broker.stop(); |
| | | } |
| | | assertEquals(resetNumUpdateReplayThreads(), 0, |
| | | "the number of replay threads could not be put back"); |
| | | } |
| | | finally |
| | | { |
| | | /* |
| | | * Best-effort on the way out of a red, with no assertion to replace the failure |
| | | * which is being reported: on the normal road the number is put back, and checked, |
| | | * above. |
| | | */ |
| | | resetNumUpdateReplayThreads(); |
| | | } |
| | | } |
| | | finally |
| | | { |
| | | try |
| | | { |
| | | /* |
| | | * A red before the barrier was given up would leave it listed for the rest of this |
| | | * class: it never travelled the replication server, so no session restart brings a |
| | | * delivery which could give it up, and a change which is failing keeps every later |
| | | * session restart at its backoff and the ServerState behind it. The failure which |
| | | * is being reported is the one that matters, so this does not replace it. |
| | | */ |
| | | if (!domain.getServerState().cover(failing)) |
| | | { |
| | | giveUpOn(domain, failing, waitedOn, failingMods, waitedOnUUID); |
| | | } |
| | | } |
| | | catch (Throwable cleanupFailure) |
| | | { |
| | | logger.error(LocalizableMessage.raw( |
| | | "the barrier of aChangeParkedByAnUnwoundReplayIsDeliveredAgain could not be given" |
| | | + " up on the way out: %s", StaticUtils.stackTraceToSingleLineString(cleanupFailure))); |
| | | } |
| | | finally |
| | | { |
| | | resetReplayGiveUpDelay(); |
| | | } |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Gives up on the change no delivery can replay: the budget is shortened to nothing, so |
| | | * that the change is given up on as soon as one more delivery of it fails. Nothing sends |
| | | * that change again - it never travelled a session - so its deliveries are made here. |
| | | */ |
| | | private void giveUpOn(final LDAPReplicationDomain domain, final CSN failing, final DN dn, |
| | | final List<Modification> mods, final String entryUUID) throws Exception |
| | | { |
| | | setReplayGiveUpDelay("0ms"); |
| | | TestTimer timer = new TestTimer.Builder() |
| | | .maxSleep(120, SECONDS) |
| | | .sleepTimes(500, MILLISECONDS) |
| | | .toTimer(); |
| | | timer.repeatUntilSuccess(new CallableVoid() |
| | | { |
| | | @Override |
| | | public void call() throws Exception |
| | | { |
| | | if (!domain.getServerState().cover(failing)) |
| | | { |
| | | domain.processUpdate(new ModifyMsgWhoseOperationRefusesAControl(failing, dn, mods, entryUUID)); |
| | | } |
| | | assertTrue(domain.getServerState().cover(failing), |
| | | "the change no delivery can replay must be given up on"); |
| | | } |
| | | }); |
| | | } |
| | | |
| | | /** |
| | | * Test case for [Issue 986]: a change parked as waiting for another one is given back |
| | | * when the replay thread which parked it is stopped with the pool. |
| | | * <p> |
| | | * A parked change stays owned by the replay thread which parked it while that thread |
| | | * goes back to the pool and takes the changes which follow: {@code getNextUpdate()} is |
| | | * what hands it out again, to whichever replay thread clears the change it was waiting |
| | | * for. Changing the number of replay threads stops the whole pool and creates another |
| | | * one, so a thread which parked a change and went back to the queue is joined while it |
| | | * is idle, and it would end still recorded as the owner of that change - a thread which |
| | | * does not exist anymore, while every redelivery of a change a replay thread owns is |
| | | * refused as a duplicate. On a domain which then goes quiet that change is where this |
| | | * replica's ServerState, and every change behind it from every master, stops. |
| | | */ |
| | | @Test |
| | | public void aChangeParkedByAThreadThePoolStoppedIsDeliveredAgain() throws Exception |
| | | { |
| | | testSetUp("aChangeParkedByAThreadThePoolStoppedIsDeliveredAgain"); |
| | | logger.error(LocalizableMessage.raw( |
| | | "Starting replication test : aChangeParkedByAThreadThePoolStoppedIsDeliveredAgain")); |
| | | |
| | | final DN waitedOn = addEntryForChange("user.986.1"); |
| | | final String waitedOnUUID = getEntry(waitedOn, 1, true).parseAttribute("entryuuid").asString(); |
| | | final DN other = addEntryForChange("user.986.2"); |
| | | final String otherUUID = getEntry(other, 1, true).parseAttribute("entryuuid").asString(); |
| | | |
| | | final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null); |
| | | domain.resetUnreplayedChangeAlertThrottle(); |
| | | final long inProgress = getMonitorAttrValue(baseDN, "changes-in-progress-size"); |
| | | assertEquals(getMonitorAttrValue(baseDN, "dependent-changes-size"), 0, |
| | | "no change of this domain is waiting for another one when this test starts"); |
| | | |
| | | final CSNGenerator gen = new CSNGenerator(26, TimeThread.getTime()); |
| | | final CSN failing = gen.newCSN(); |
| | | final CSN parked = gen.newCSN(); |
| | | final CSN applied = gen.newCSN(); |
| | | final List<Modification> failingMods = generatemods("description", "the replay of this change fails"); |
| | | final String parkedDescription = "the change which was parked by a thread the pool stopped"; |
| | | final List<Modification> parkedMods = generatemods("description", parkedDescription); |
| | | final String appliedDescription = "the change which is applied before the pool is stopped"; |
| | | final List<Modification> appliedMods = generatemods("description", appliedDescription); |
| | | |
| | | /* |
| | | * One replay thread, so that the change which is parked is parked by the thread the |
| | | * pool then stops: a parked change is left owned by the thread which parked it, and |
| | | * this is about a thread which is not there anymore to be given it back. |
| | | */ |
| | | setNumUpdateReplayThreads(1); |
| | | try |
| | | { |
| | | /* |
| | | * A change whose replay failed stays listed and uncommitted - it is what holds this |
| | | * domain's ServerState back - 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. |
| | | */ |
| | | deliverUntilMonitorReaches(domain, "changes-in-progress-size", inProgress + 1, |
| | | () -> new ModifyMsgWhoseOperationRefusesAControl(failing, waitedOn, failingMods, waitedOnUUID), |
| | | "the change whose replay fails must stay listed as one which is not in the data"); |
| | | |
| | | /* |
| | | * The change which is parked as waiting for it. The replay thread it was given to |
| | | * parks it and goes back to the pool: nothing is waiting for it there, so it is idle |
| | | * and it still owns the change it parked. |
| | | */ |
| | | deliverUntilMonitorReaches(domain, "dependent-changes-size", 1, |
| | | () -> new ModifyMsg(parked, waitedOn, parkedMods, waitedOnUUID), |
| | | "a change which waits for one that is not in the data must be parked"); |
| | | |
| | | /* |
| | | * A change which is applied, on another entry so that it is replayed rather than |
| | | * parked in its turn. It is there for the count of the deliveries this session took |
| | | * off it, replayed-updates: a session which is started counts from zero, so that |
| | | * count going back to zero below is what says the session was restarted - and it is |
| | | * not left to the give-back alone to put the count above zero before the pool is |
| | | * stopped. The restart the change whose replay failed asked for is over by now: the |
| | | * thread which ran it is the one which parked the change above. |
| | | */ |
| | | domain.processUpdate(new ModifyMsg(applied, other, appliedMods, otherUUID)); |
| | | checkEntryHasAttributeValue(other, "description", appliedDescription, 30, |
| | | "the change made on the other entry must be applied"); |
| | | // Counted once its ack is out, which is after the change is in the data. |
| | | final TestTimer counted = new TestTimer.Builder() |
| | | .maxSleep(30, SECONDS) |
| | | .sleepTimes(200, MILLISECONDS) |
| | | .toTimer(); |
| | | counted.repeatUntilSuccess(new CallableVoid() |
| | | { |
| | | @Override |
| | | public void call() throws Exception |
| | | { |
| | | assertTrue(getMonitorAttrValue(baseDN, "replayed-updates") > 0, |
| | | "the session must have counted the deliveries it took off before the pool is stopped"); |
| | | } |
| | | }); |
| | | |
| | | /* |
| | | * How many restarts in a row the session has been through, the barrier's among them: |
| | | * the restart the give-back of a stopped thread asks for is owed no backoff - what |
| | | * went away is a replay thread, not the backend - and a restart which waits the |
| | | * backoff out is the one road which moves this count. It can not move otherwise |
| | | * between here and the reading below: a replay which made it puts the count back to |
| | | * zero only once nothing is failing anymore, and the barrier keeps failing until it |
| | | * is given up, below. The restart the barrier asked for is over: the thread which |
| | | * ran it is the one which parked the change and applied the other, above. |
| | | */ |
| | | final int restarts = domain.getConsecutiveSessionRestarts(); |
| | | |
| | | /* |
| | | * The pool is stopped and created again, the way an administrator changing the |
| | | * number of replay threads has it: the thread which parked the change is joined |
| | | * where it waits for the next delivery, and the change it owns is nobody's. |
| | | */ |
| | | setNumUpdateReplayThreads(2); |
| | | |
| | | assertMonitorAttrValueEventually(baseDN, "dependent-changes-size", 0, |
| | | "a change parked by a replay thread the pool stopped must be given back"); |
| | | /* |
| | | * The give-back alone leaves the count where it was, plus the change it handed back: |
| | | * only a session which is started counts from zero, and the changes of this test |
| | | * never travelled a session, so nothing is delivered over the one which is started |
| | | * before the deliveries made below. |
| | | */ |
| | | assertMonitorAttrValueEventually(baseDN, "replayed-updates", 0, |
| | | "the session must be restarted for the changes which were given back"); |
| | | assertEquals(domain.getConsecutiveSessionRestarts(), restarts, |
| | | "the restart run for the changes a stopped thread gave back must not wait the" |
| | | + " backoff out: what went away is a replay thread, not the backend"); |
| | | |
| | | /* |
| | | * Nothing sends these changes again - they never travelled a session - so the |
| | | * deliveries which take over from the ones the stopped pool left behind are made |
| | | * here. The change no delivery can replay is given up on, which lets the ServerState |
| | | * past it, and the change which was parked behind it is applied. |
| | | */ |
| | | setReplayGiveUpDelay(TEST_GIVE_UP_DELAY); |
| | | TestTimer timer = new TestTimer.Builder() |
| | | .maxSleep(120, SECONDS) |
| | | .sleepTimes(500, MILLISECONDS) |
| | | .toTimer(); |
| | | timer.repeatUntilSuccess(new CallableVoid() |
| | | { |
| | | @Override |
| | | public void call() throws Exception |
| | | { |
| | | final ServerState state = domain.getServerState(); |
| | | if (!state.cover(failing)) |
| | | { |
| | | domain.processUpdate( |
| | | new ModifyMsgWhoseOperationRefusesAControl(failing, waitedOn, failingMods, waitedOnUUID)); |
| | | } |
| | | if (!state.cover(parked)) |
| | | { |
| | | domain.processUpdate(new ModifyMsg(parked, waitedOn, parkedMods, waitedOnUUID)); |
| | | } |
| | | assertTrue(state.cover(parked), |
| | | "the change which was given back must be replayed by the delivery which takes it over"); |
| | | } |
| | | }); |
| | | checkEntryHasAttributeValue(waitedOn, "description", parkedDescription, 30, |
| | | "the change which was parked must be applied by the delivery which took it over"); |
| | | assertEquals(resetNumUpdateReplayThreads(), 0, |
| | | "the number of replay threads could not be put back"); |
| | | } |
| | | finally |
| | | { |
| | | try |
| | | { |
| | | /* |
| | | * A red before the barrier was given up would leave it listed for the rest of this |
| | | * class: it never travelled the replication server, so no session restart brings a |
| | | * delivery which could give it up, and a change which is failing keeps every later |
| | | * session restart at its backoff and the ServerState behind it. The failure which |
| | | * is being reported is the one that matters, so this does not replace it. |
| | | */ |
| | | if (!domain.getServerState().cover(failing)) |
| | | { |
| | | giveUpOn(domain, failing, waitedOn, failingMods, waitedOnUUID); |
| | | } |
| | | } |
| | | catch (Throwable cleanupFailure) |
| | | { |
| | | logger.error(LocalizableMessage.raw( |
| | | "the barrier of aChangeParkedByAThreadThePoolStoppedIsDeliveredAgain could not be" |
| | | + " given up on the way out: %s", StaticUtils.stackTraceToSingleLineString(cleanupFailure))); |
| | | } |
| | | finally |
| | | { |
| | | try |
| | | { |
| | | resetReplayGiveUpDelay(); |
| | | } |
| | | finally |
| | | { |
| | | /* |
| | | * Best-effort on the way out of a red, with no assertion to replace the failure |
| | | * which is being reported: on the normal road the number is put back, and checked, |
| | | * above. |
| | | */ |
| | | resetNumUpdateReplayThreads(); |
| | | } |
| | | } |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Delivers a change until a monitor attribute of the domain reaches the expected value. |
| | | * <p> |
| | | * A delivery is dropped rather than queued while the listener thread is down, which it |
| | | * is for as long as a recovery is restarting the session, so a change which has to reach |
| | | * a replay thread is delivered until it does. A delivery of a change a replay thread |
| | | * owns is refused as the duplicate it is, so the deliveries which follow the one that |
| | | * was taken cost nothing. |
| | | */ |
| | | private void deliverUntilMonitorReaches(final LDAPReplicationDomain domain, |
| | | final String attributeName, final long expected, |
| | | final Supplier<? extends LDAPUpdateMsg> delivery, final String message) throws Exception |
| | | { |
| | | TestTimer timer = new TestTimer.Builder() |
| | | .maxSleep(20, SECONDS) |
| | | .sleepTimes(500, MILLISECONDS) |
| | | .toTimer(); |
| | | timer.repeatUntilSuccess(new CallableVoid() |
| | | { |
| | | @Override |
| | | public void call() throws Exception |
| | | { |
| | | if (getMonitorAttrValue(baseDN, attributeName) != expected) |
| | | { |
| | | domain.processUpdate(delivery.get()); |
| | | } |
| | | assertEquals(getMonitorAttrValue(baseDN, attributeName), expected, message); |
| | | } |
| | | }); |
| | | } |
| | | |
| | | /** |
| | | * Sets how many replay threads this server runs, the way an administrator would: the |
| | | * pool is stopped and created again with that number. |
| | | */ |
| | | private static void setNumUpdateReplayThreads(int threads) throws Exception |
| | | { |
| | | assertEquals(TestCaseUtils.applyModifications(true, |
| | | "dn: " + SYNCHRO_PLUGIN_DN, |
| | | "changetype: modify", |
| | | "replace: ds-cfg-num-update-replay-threads", |
| | | "ds-cfg-num-update-replay-threads: " + threads), 0, |
| | | "the number of replay threads could not be changed"); |
| | | } |
| | | |
| | | /** |
| | | * Puts the number of replay threads back to what this server computes for itself, which |
| | | * is what it runs with when the configuration carries no number of its own. |
| | | * |
| | | * @return the result code of the change, so that a finally can call this without an |
| | | * assertion which would replace the failure it is on the way out of |
| | | */ |
| | | private static int resetNumUpdateReplayThreads() throws Exception |
| | | { |
| | | return TestCaseUtils.applyModifications(true, |
| | | "dn: " + SYNCHRO_PLUGIN_DN, |
| | | "changetype: modify", |
| | | "delete: ds-cfg-num-update-replay-threads"); |
| | | } |
| | | |
| | | /** Adds the entry a change of these tests is made on. */ |
| | | private DN addEntryForChange(String uid) throws Exception |
| | | { |
| | | return TestCaseUtils.addEntry( |
| | | "dn: uid=" + uid + "," + baseDN, |
| | | "objectClass: top", |
| | | "objectClass: person", |
| | | "objectClass: organizationalPerson", |
| | | "objectClass: inetOrgPerson", |
| | | "uid: " + uid, |
| | | "cn: Aaccf Amar", |
| | | "sn: Amar").getName(); |
| | | } |
| | | |
| | | /** A delivery of a change whose replay does not run to its end. */ |
| | | private interface FailingDelivery |
| | | { |
| | |
| | | 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(); |
| | | final DN dn = addEntryForChange(uid); |
| | | final String uuid = getEntry(dn, 1, true).parseAttribute("entryuuid").asString(); |
| | | |
| | | final LDAPReplicationDomain domain = MultimasterReplication.findDomain(baseDN, null); |
| | |
| | | } |
| | | |
| | | /** |
| | | * A ModifyMsg whose ack runs out of memory on the way out of a replay which applied it. |
| | | * <p> |
| | | * The change is committed before the ack of its delivery is published, and the ack is |
| | | * where the JVM runs out of memory: the one throw from there which is not caught, so the |
| | | * replay is unwound with the change it was replaying in the data and owned by nobody - |
| | | * commit() cleared the owner - and the thread ends on the error. What the give-back on |
| | | * the way out of {@code replay()} has left to give back is the changes this thread parked |
| | | * as waiting for another one, and the restart it asks for them is the one which is run: |
| | | * the road a failed replay takes to ask for its own change again is not on the way. |
| | | * <p> |
| | | * Nothing on the way in reads what throws here: a message handed to the domain rather |
| | | * than published is not one this server acknowledges to anybody. |
| | | */ |
| | | private static final class ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied extends ModifyMsg |
| | | { |
| | | private ModifyMsgWhoseAckRunsOutOfMemoryOnceApplied( |
| | | 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 applied 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 |