| | |
| | | * in flight unreplayable, and one alert per change would be a storm. |
| | | */ |
| | | private static final long UNREPLAYED_CHANGE_ALERT_INTERVAL_IN_MS = 60000; |
| | | /** |
| | | * What the ack of a delivery whose replay ran out of memory says did not apply the |
| | | * change. |
| | | * <p> |
| | | * A constant rather than a message built from the change, because this is read on the |
| | | * road out of an {@code OutOfMemoryError}: formatting one asks the JVM for the memory |
| | | * it has just refused, and what {@code processUpdateDone()} needs of this string is |
| | | * that there is one - it sets {@code hasReplayError} on the ack and never sends the |
| | | * text, which is why no ordinal is spent on it either. |
| | | */ |
| | | private static final String REPLAY_RAN_OUT_OF_MEMORY = "the replay of this change ran out of memory"; |
| | | /** The number of updates this replica gave up replaying. */ |
| | | private final AtomicInteger numFailedReplayedUpdates = new AtomicInteger(); |
| | | /** Set while a replay thread is restarting the session after a failed replay. */ |
| | |
| | | |
| | | /** |
| | | * Create and replay a synchronized Operation from an UpdateMsg. |
| | | * <p> |
| | | * The change is given back on the way out of a replay which was unwound before it |
| | | * reached one of the roads which give it back: a change which stays owned by a thread |
| | | * which is not replaying it anymore is refused as a duplicate on every later delivery, |
| | | * so this domain's ServerState would never move past it (issue #922). |
| | | * |
| | | * @param msg |
| | | * The UpdateMsg to be replayed. |
| | |
| | | */ |
| | | void replay(LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown) |
| | | { |
| | | try |
| | | { |
| | | replayChangeAndTheChangesWaitingForIt(msg, replayThreadShutdown); |
| | | } |
| | | catch (Throwable t) |
| | | { |
| | | /* |
| | | * The roads which run to their end give the change back themselves, so what is left |
| | | * to give back here is a change whose replay was unwound over them: by an |
| | | * OutOfMemoryError, which is left to end this thread, or by a throw from what the |
| | | * replay runs once the ack of the delivery has been published - the give-back of the |
| | | * change and the hand-out of the changes which were waiting for it are on that road. |
| | | * |
| | | * It takes the road of a failed replay rather than being handed back on the spot, so |
| | | * that a change which keeps unwinding the replays it is given to is counted as |
| | | * failing and is eventually given up on: handing it back bare would have this domain |
| | | * ask for it, and restart its session for it, for as long as the server is up. |
| | | * |
| | | * The changes this thread parked as waiting for another change are left alone: they |
| | | * are handed to whichever thread clears the change they are waiting for, and that |
| | | * thread takes them over. |
| | | * |
| | | * Which change this thread owns is read before anything is done with it, and that |
| | | * read takes no lock and allocates nothing: everything below is gated on the answer, |
| | | * so a lookup which threw in its turn - on the road out of a JVM which has just |
| | | * refused an allocation - would leave the change listed, uncommitted and owned by a |
| | | * thread which is about to end, which is the state this whole issue is about. |
| | | */ |
| | | CSN owned = null; |
| | | try |
| | | { |
| | | owned = remotePendingChanges.getChangeOwnedByCurrentThread(); |
| | | if (owned != null) |
| | | { |
| | | if (replayThreadShutdown.get() || shutdown.get() || disabled) |
| | | { |
| | | /* |
| | | * The replay was not failing, it was being abandoned: this thread is stopping |
| | | * because their number is being changed, or the domain is going away or is |
| | | * being imported into. The change is handed back without being counted against |
| | | * its give-up budget, which is the road a replay abandoned that way takes. |
| | | */ |
| | | abandonReplay(owned); |
| | | } |
| | | else |
| | | { |
| | | /* |
| | | * A JVM which has run out of memory is told apart here rather than reported on |
| | | * its own road: the change is given back counted, the way every other unwound |
| | | * replay gives it back, but the line which says it is being asked for again is |
| | | * not built - that asks the JVM for the memory it has just refused - and the |
| | | * session is restarted without sitting through the backoff, since the thread |
| | | * which is doing it is on its way out. |
| | | * |
| | | * A change whose budget is spent is still reported and still raises its alert |
| | | * on this road: it is the one line which says this replica has diverged, and an |
| | | * operator who is not told would be left with a replica which is silently |
| | | * behind. That is the deliberate exception to the rule above. |
| | | */ |
| | | recoverFromReplayFailure(owned, replayThreadShutdown, t instanceof OutOfMemoryError); |
| | | } |
| | | } |
| | | } |
| | | catch (Throwable recoveryFailure) |
| | | { |
| | | /* |
| | | * The give-back is the road which hands the change over, so a throw out of it - an |
| | | * allocation which fails in its turn, where what unwound the replay was a JVM out |
| | | * of memory - would leave the change owned by this thread after all. Hand it back |
| | | * bare, without the failure count that road did not reach, and restart the session |
| | | * so that it is delivered again: this is the last resort, the throwable is rethrown |
| | | * whatever happens here, and the thread this runs on may well be ending on it - so |
| | | * a request left for the next recovery of this domain to pick up is a request which |
| | | * may never be run. |
| | | */ |
| | | if (owned != null) |
| | | { |
| | | remotePendingChanges.replayFailed(owned); |
| | | sessionRestartRequested.set(true); |
| | | /* |
| | | * Reported and restarted under guards of their own, and in that order: the report |
| | | * is the line an operator acts on, and the restart is what has the change |
| | | * delivered again - so a report which can not be formatted, on a road an |
| | | * OutOfMemoryError leads to, must not cost the restart. |
| | | */ |
| | | try |
| | | { |
| | | logger.error(ERR_REPLAY_GIVE_BACK_FAILED, owned, getBaseDN(), |
| | | stackTraceToSingleLineString(recoveryFailure)); |
| | | } |
| | | catch (Throwable reportFailure) |
| | | { |
| | | suppress(recoveryFailure, reportFailure); |
| | | } |
| | | try |
| | | { |
| | | runRequestedSessionRestarts(false); |
| | | } |
| | | catch (Throwable restartFailure) |
| | | { |
| | | /* |
| | | * Nothing is left to try: the change is listed, uncommitted and unowned, so any |
| | | * later session restart of this domain delivers it again. This goes with the |
| | | * throwable which is rethrown below rather than being reported on its own. |
| | | */ |
| | | suppress(recoveryFailure, restartFailure); |
| | | } |
| | | } |
| | | // The error which unwound the replay is the one reported, whatever the give-back |
| | | // ran into on top of it. |
| | | suppress(t, recoveryFailure); |
| | | } |
| | | throw t; |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Records the second throwable as one the first suppressed, unless the two are one and |
| | | * the same. |
| | | * <p> |
| | | * A JVM which has run out of memory hands out the error it prepared before it ran out as |
| | | * often as it is asked for one, so the roads out of a replay can carry the same instance |
| | | * twice - and a throwable can not suppress itself. |
| | | * |
| | | * @param thrown the throwable which is reported |
| | | * @param alsoThrown what was met on the way out of it |
| | | */ |
| | | private static void suppress(Throwable thrown, Throwable alsoThrown) |
| | | { |
| | | if (thrown != alsoThrown) |
| | | { |
| | | thrown.addSuppressed(alsoThrown); |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * Replays the change of the provided message, then the changes which were waiting for |
| | | * it, for as long as there are some. |
| | | * |
| | | * @param msg |
| | | * The UpdateMsg to be replayed. |
| | | * @param replayThreadShutdown |
| | | * whether the replay thread was asked to stop |
| | | */ |
| | | private void replayChangeAndTheChangesWaitingForIt( |
| | | LDAPUpdateMsg msg, AtomicBoolean replayThreadShutdown) |
| | | { |
| | | // Try replay the operation, then flush (replaying) any pending operation |
| | | // whose dependency has been replayed until no more left. |
| | | do |
| | |
| | | boolean replayAbandoned = false; |
| | | String replayErrorMsg = null; |
| | | CSN csn = null; |
| | | /* |
| | | * Read once, before anything of this delivery has run, so that the report of an ack |
| | | * which could not be published names the change even when reading it off the message |
| | | * is what threw: that report is written on the way out of a road which is already |
| | | * failing, and it must not be the throw which unwinds the replay. |
| | | */ |
| | | final CSN delivered = msg.getCSN(); |
| | | try |
| | | { |
| | | // The next operation for which to attempt replay. |
| | |
| | | replayErrorMsg = message.toString(); |
| | | replayFailed = true; |
| | | } |
| | | } finally |
| | | } |
| | | catch (OutOfMemoryError e) |
| | | { |
| | | /* |
| | | * The JVM is out of memory, which is not something to carry on replaying from: the |
| | | * error is left to unwind the replay thread, which ends on it (issue #923). It is |
| | | * not turned into a report of its own here either - the stack trace of an error |
| | | * which can be raised anywhere says little, and the uncaught exception handler of |
| | | * DirectoryThread writes the one line this is worth, with an alert. |
| | | * |
| | | * The other errors of the JVM take the road below. A StackOverflowError is met by |
| | | * the thread which recursed and is gone once the stack has unwound, and the entry |
| | | * being replayed is what raises it rather than the state of this server: ending a |
| | | * thread on it would have one change this replica can not replay cost it a replay |
| | | * thread per delivery, and nothing creates a replay thread to replace one which |
| | | * ends. |
| | | * |
| | | * The change is given back, counted as failing and asked for again on the way out |
| | | * (issue #922). The one thing done here is to make the ack this delivery publishes |
| | | * below say that the change was not applied: a replica which is asking for a change |
| | | * again must not have told an assured write that it is in the data here. It is a |
| | | * constant rather than a message built from the change, because building one asks |
| | | * the JVM for the memory it has just refused. |
| | | */ |
| | | replayErrorMsg = REPLAY_RAN_OUT_OF_MEMORY; |
| | | throw e; |
| | | } |
| | | catch (Error e) |
| | | { |
| | | /* |
| | | * An Error out of the replay - a LinkageError met where a plugin or a backend class |
| | | * is loaded, an AssertionError - unwinds every road this replay has out of here, the |
| | | * ones which give the change back among them. The change is not in the data, so it |
| | | * is reported and given back here, on the road every other failed replay takes: the |
| | | * ack below says it was not applied, the failure counts against the give-up budget |
| | | * of the change, the session is restarted for it to be delivered again (issue #922) |
| | | * and the changes which were waiting for this one are replayed rather than left |
| | | * waiting for a thread to hand them out. |
| | | * |
| | | * It is not given up on where no operation could be built from the message, which is |
| | | * what an Exception at that point means: an Error says that this server could not run |
| | | * the replay, not that the message is one no delivery could ever build an operation |
| | | * from. |
| | | */ |
| | | final LocalizableMessage message = ERR_ERROR_REPLAYING_CHANGE.get( |
| | | msg.getCSN(), getBaseDN(), stackTraceToSingleLineString(e)); |
| | | logger.error(message); |
| | | replayErrorMsg = message.toString(); |
| | | replayFailed = true; |
| | | } |
| | | finally |
| | | { |
| | | if (!dependency) |
| | | { |
| | |
| | | * down, so nothing would reach the server which is waiting for it, and an |
| | | * assured write would wait out its timeout rather than be told what happened. |
| | | */ |
| | | processUpdateDone(msg, replayErrorMsg); |
| | | try |
| | | { |
| | | processUpdateDone(msg, replayErrorMsg); |
| | | } |
| | | catch (OutOfMemoryError e) |
| | | { |
| | | /* |
| | | * The one throw from here which is not caught, on the same terms as the arm |
| | | * above: a JVM which has run out of memory is not something to carry on |
| | | * replaying from, the error is left to end this replay thread, and the uncaught |
| | | * exception handler of DirectoryThread raises the alert #923 is about. Swallowed |
| | | * here instead, it would have this thread go on to the roads below and to the |
| | | * next change of the loop on an exhausted heap, with nothing reported anywhere. |
| | | * |
| | | * What the give-back on the way out of replay() then finds depends on the road |
| | | * the replay took to get here. A replay which failed, or which was unwinding on |
| | | * an OutOfMemoryError of its own, still owns its change: it is given back |
| | | * counted, and the thread ends on this error rather than on the one it stepped |
| | | * over. A replay which committed owns nothing anymore - commit() cleared the |
| | | * owner, and the index the give-back reads, in the same step - so the give-back |
| | | * is a no-op, and rightly so: a change which is in the data is not one to ask |
| | | * for again. What that road steps over is getNextUpdate() below, so the changes |
| | | * parked behind the committed change wait for the next replay of this domain to |
| | | * hand them out. That is the trade #923 asks for: a thread which met this error |
| | | * is not to carry on, not even for them. |
| | | */ |
| | | throw e; |
| | | } |
| | | catch (Throwable ackFailure) |
| | | { |
| | | /* |
| | | * Publishing the ack says nothing about whether the change was applied, so a |
| | | * throw here must not be a road out of the replay. It would step over the |
| | | * give-back of the change and, where the change was committed, over the |
| | | * getNextUpdate() below - the one drain of the changes this thread parked as |
| | | * waiting for it - leaving them waiting for a thread which is not replaying |
| | | * anything anymore. |
| | | * |
| | | * A master which is waiting for the ack waits out its assured timeout either |
| | | * way - and there is one only for an assured write in safe-read mode, which is |
| | | * the one delivery a replica acknowledges. processUpdateDone() runs for every |
| | | * delivery all the same: it accounts for the delivery in the receive window and |
| | | * in the processed-updates counter, and a throw from there is a throw from this |
| | | * bookkeeping, with no ack owed to anybody. Either way it is reported and the |
| | | * replay carries on to the road the change itself decided: applied, failed and |
| | | * asked for again, or given up on. |
| | | * |
| | | * Every step of processUpdateDone() catches what it can meet - the broker keeps |
| | | * a failure to publish to itself and retries it - so what reaches here is what |
| | | * no code of the replication protocol expected: an Error, and the tests drive |
| | | * it as one. |
| | | */ |
| | | try |
| | | { |
| | | logger.error(ERR_ACK_NOT_PUBLISHED, delivered, getBaseDN(), |
| | | stackTraceToSingleLineString(ackFailure)); |
| | | } |
| | | catch (Throwable reportFailure) |
| | | { |
| | | /* |
| | | * Guarded like the report of a give-back which failed: this one runs on the |
| | | * same kind of road - the stack of the throwable is walked to build the line - |
| | | * and a report which can not be built must not become the throw which unwinds |
| | | * the replay past the give-back and past getNextUpdate(). |
| | | * |
| | | * The line is tried once more with the name of the error alone, which walks no |
| | | * stack. Nothing rethrows what was caught here, so an error recorded as |
| | | * suppressed on it would be recorded nowhere, and this is the road on which an |
| | | * operator has the least to go on. A second refusal leaves nothing to say it |
| | | * with, and the replay carries on with what the change decided. |
| | | */ |
| | | try |
| | | { |
| | | logger.error(ERR_ACK_NOT_PUBLISHED, delivered, getBaseDN(), |
| | | ackFailure.getClass().getName()); |
| | | } |
| | | catch (Throwable secondReportFailure) |
| | | { |
| | | // Nothing is left to say it with. |
| | | } |
| | | } |
| | | } |
| | | } |
| | | } |
| | | |
| | |
| | | */ |
| | | private boolean recoverFromReplayFailure(CSN csn, AtomicBoolean replayThreadShutdown) |
| | | { |
| | | return recoverFromReplayFailure(csn, replayThreadShutdown, false); |
| | | } |
| | | |
| | | /** |
| | | * Recovers from a change which could not be replayed, telling apart the replay which was |
| | | * unwound by a JVM out of memory. |
| | | * |
| | | * @param csn |
| | | * the CSN of the change which could not be replayed |
| | | * @param replayThreadShutdown |
| | | * whether the replay thread was asked to stop |
| | | * @param outOfMemory |
| | | * whether what unwound the replay was the JVM running out of memory, in which |
| | | * case the line which says the change is being asked for again is not built - |
| | | * the report of a change which is given up on is, since it is what says this |
| | | * replica has diverged - and the session is restarted without the backoff |
| | | * @return {@code true} when the caller must stop replaying because the session is |
| | | * being restarted or is going away, {@code false} when it may carry on with |
| | | * the changes which follow |
| | | */ |
| | | private boolean recoverFromReplayFailure( |
| | | CSN csn, AtomicBoolean replayThreadShutdown, boolean outOfMemory) |
| | | { |
| | | /* |
| | | * The failure is recorded, and the change given up on, while this thread still owns |
| | | * it: a change which is listed, uncommitted and unowned is what putRemoteUpdate() |
| | |
| | | return true; |
| | | } |
| | | |
| | | logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts()); |
| | | if (!outOfMemory) |
| | | { |
| | | /* |
| | | * Not on the road out of a JVM which has run out of memory: building this line asks |
| | | * it for the memory it has just refused, and the ack of the delivery already says |
| | | * that the change was not applied. The constant that ack carries exists for the same |
| | | * reason. |
| | | */ |
| | | logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts()); |
| | | } |
| | | /* |
| | | * This change is not owned by anyone anymore, so the session has to be restarted for |
| | | * the replication server to deliver it again. Ask for the restart before trying to |
| | |
| | | * A replay thread which is stopping - the number of them is being changed - restarts |
| | | * the session all the same: nothing else would ask for the change it just released, |
| | | * and the ServerState would stay behind it for good. It does not sit through the |
| | | * backoff on its way out, though: the backend is not what is going away. |
| | | * backoff on its way out, though: the backend is not what is going away. Neither does |
| | | * the thread an OutOfMemoryError is ending, for the same reason - and the restart is |
| | | * run rather than left to be asked for again, because that thread will not be there to |
| | | * run it, and a change nobody asks for again holds this domain's ServerState back. |
| | | */ |
| | | runRequestedSessionRestarts(!replayThreadShutdown.get()); |
| | | runRequestedSessionRestarts(!replayThreadShutdown.get() && !outOfMemory); |
| | | return true; |
| | | } |
| | | |
| | |
| | | { |
| | | while (sessionRestartRequested.getAndSet(false)) |
| | | { |
| | | restartSession(wait); |
| | | boolean restarted = false; |
| | | try |
| | | { |
| | | restartSession(wait); |
| | | restarted = true; |
| | | } |
| | | finally |
| | | { |
| | | if (!restarted) |
| | | { |
| | | /* |
| | | * The request is put back where it was taken from. The flag is read and |
| | | * cleared before the restart runs, so a restart which ends abruptly - the |
| | | * session is stopped first, and starting it again creates a listener thread, |
| | | * which the operating system can refuse - would otherwise leave this domain |
| | | * with no session and with nothing left to ask for one. |
| | | * |
| | | * What a request left standing buys is bounded, and the bound is worth |
| | | * stating. Its two readers are the roads out of a failed and of an abandoned |
| | | * replay of this domain, and with no listener thread nothing is delivered |
| | | * anymore: the replays left to run are the changes already taken off the |
| | | * session - the ones waiting in the replay queue, and the ones parked as |
| | | * dependencies. One of those failing finds the request standing and runs the |
| | | * restart, which starts from a clean state, since disableService() drops the |
| | | * listener thread which was never started. Once they are spent, the domain |
| | | * stays down until it is disabled and enabled back, or the server is |
| | | * restarted. That is said where it can be heard: a refused thread is an |
| | | * OutOfMemoryError, and one which leaves recoverFromReplayFailure() or |
| | | * abandonReplay() ends the replay thread it is met on, so the uncaught |
| | | * exception handler of DirectoryThread writes the line and raises the alert, |
| | | * with the start of the listener thread in the trace. |
| | | */ |
| | | sessionRestartRequested.set(true); |
| | | } |
| | | } |
| | | } |
| | | } |
| | | finally |