| | |
| | | * another road asked for with the backoff, or one given back with it, keeps its |
| | | * wait whichever thread runs it. |
| | | */ |
| | | final boolean parkedGivenBack = giveBackParkedChanges( |
| | | final List<CSN> parkedGivenBack = giveBackParkedChanges( |
| | | replayThreadShutdown.get() || t instanceof OutOfMemoryError |
| | | ? SessionRestart.NOW : SessionRestart.AFTER_BACKOFF); |
| | | if (owned != null) |
| | |
| | | */ |
| | | recoverFromReplayFailure(owned, replayThreadShutdown, t instanceof OutOfMemoryError); |
| | | } |
| | | /* |
| | | * The change is handed back and counted as the road it took counts it: the last |
| | | * resort below speaks for a give-back which did not run, and a throw out of what |
| | | * follows - the restart the parked changes are run with, or the line which says it |
| | | * is held - is not one. Left set, that throw would have the last resort report the |
| | | * change as "released without its failure being counted", which is the one line an |
| | | * operator acts on, and it would be false. |
| | | */ |
| | | owned = null; |
| | | } |
| | | if (parkedGivenBack && !replayThreadShutdown.get() && !sessionHasAnOwner()) |
| | | if (!parkedGivenBack.isEmpty() && !replayThreadShutdown.get() && !sessionHasAnOwner()) |
| | | { |
| | | /* |
| | | * The road the change this thread was replaying took may have run the restart the |
| | |
| | | * the state checkpointer runs one restart for every change the threads of the |
| | | * pool hand back on their way out, rather than each of them running one while |
| | | * the configuration change which is stopping them waits. A domain whose session |
| | | * has an owner is left alone the way the give-back left it: nothing was asked |
| | | * for on that road, and a request another thread left standing is not this |
| | | * one's to spend on a restart which is refused where it runs. |
| | | * has an owner is left alone the way the give-back left it - it asked for nothing |
| | | * there, which is why it handed back no change to run a restart for - and a |
| | | * request another thread left standing is not this one's to spend on a restart |
| | | * which is refused where it runs. |
| | | */ |
| | | runRequestedSessionRestarts(); |
| | | if (!runRequestedSessionRestarts()) |
| | | { |
| | | /* |
| | | * The changes above were reported as given back to a replication server which |
| | | * "still owns it and sends it again", and it does not send them yet: a total |
| | | * update is being processed over the session, the restart which brings them back |
| | | * waits for it, and the state checkpointer runs it once it is over - the same |
| | | * hold, and the same line, a change whose replay failed is reported with. One |
| | | * line per change, the way the give-back reports them: these are the changes |
| | | * this one is about, and nothing else says they wait. |
| | | * |
| | | * Built on the road out of a JVM which has run out of memory too, where the line |
| | | * of a failed replay is not: the give-back has already built one line per change |
| | | * on that road, so what this asks the JVM for is not memory it was spared. |
| | | */ |
| | | for (CSN csn : parkedGivenBack) |
| | | { |
| | | logger.info(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE, csn, getBaseDN()); |
| | | } |
| | | } |
| | | } |
| | | } |
| | | catch (Throwable recoveryFailure) |
| | |
| | | * whatever happens here, and the thread this runs on may well be ending on it. A |
| | | * restart which can not run here leaves its request standing, and the state |
| | | * checkpointer of this domain runs it: one which threw asks for itself again on |
| | | * its way out, and one a domain whose session has an owner would refuse is not |
| | | * its way out, one a total update is being processed over the session is not run |
| | | * while it lasts, and one a domain whose session has an owner would refuse is not |
| | | * run at all rather than spent on the refusal. |
| | | * |
| | | * The restart is asked for once the change is released and not before, the way |
| | |
| | | try |
| | | { |
| | | /* |
| | | * Outside the guard above: two roads reach here with a request standing and no |
| | | * change of this thread's to hand back, and both are the parked changes' - the |
| | | * Outside the guard above: three roads reach here with a request standing and no |
| | | * change of this thread's to hand back, and all are the parked changes' - the |
| | | * give-back which released them asks for the restart before it reports them, and |
| | | * a throw out of the report - the JVM which unwound this replay is out of memory |
| | | * - leaves the request standing; and a restart the parked road ran and which |
| | | * threw has asked for one again on its way out. The changes it released are |
| | | * - leaves the request standing; a restart the parked road ran and which threw |
| | | * has asked for one again on its way out; and a throw out of the line which says |
| | | * a total update holds that restart leaves the request standing too, where it is |
| | | * refused until the total update is over. The changes it released are |
| | | * listed, uncommitted and unowned, so the request is what brings them back, and |
| | | * this thread is the one there to run it (issue #954). A give-back which threw |
| | | * before it released anything left the parked changes as they were, owned by this |
| | |
| | | if (replayFailed && recoverFromReplayFailure(msg.getCSN(), replayThreadShutdown)) |
| | | { |
| | | // The ack has been published and the change is given back: the replication server |
| | | // delivers it again, now or - while a total update owns the session - after the |
| | | // import restarts it. There is nothing left to replay here. |
| | | // delivers it again, now or - while a total update is being processed over the |
| | | // session - once it is over, when the restart which was held runs or, on the import |
| | | // direction, when the session is started from the reloaded state. There is nothing |
| | | // left to replay here. |
| | | return; |
| | | } |
| | | |
| | |
| | | * |
| | | * @param csn the CSN of the change which could not be replayed |
| | | * @param failure how long, and over how many deliveries, its replay has been failing |
| | | * @return whether the warning was written, so that a line which qualifies it - the one |
| | | * which says the restart it announced is held - is written with it and folded |
| | | * with it |
| | | */ |
| | | private void logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure failure) |
| | | private boolean logReplayRetryWarning(CSN csn, RemotePendingChanges.ReplayFailure failure) |
| | | { |
| | | final long now = monotonicNowInMs(); |
| | | final long lastLogged = lastReplayRetryWarningTime.get(); |
| | |
| | | { |
| | | logger.warn(WARN_REPLAY_RETRYING_CHANGE, csn, getBaseDN(), failure.getAttempts(), |
| | | failure.getFailingForMs(), foldedReplayRetryWarnings.getAndSet(0)); |
| | | return true; |
| | | } |
| | | else |
| | | { |
| | |
| | | foldedReplayRetryWarnings.incrementAndGet(); |
| | | logger.trace("Could not replay change %s in domain %s: delivery %d, failing for %d ms", |
| | | csn, getBaseDN(), failure.getAttempts(), failure.getFailingForMs()); |
| | | return false; |
| | | } |
| | | } |
| | | |
| | |
| | | return true; |
| | | } |
| | | |
| | | boolean warned = false; |
| | | if (!outOfMemory) |
| | | { |
| | | /* |
| | |
| | | * unlogged: the error ends the replay thread, and the uncaught exception handler of |
| | | * DirectoryThread writes the line and raises the alert for it. |
| | | */ |
| | | logReplayRetryWarning(csn, failure); |
| | | warned = logReplayRetryWarning(csn, failure); |
| | | } |
| | | /* |
| | | * This change is not owned by anyone anymore, so the session has to be restarted for |
| | |
| | | */ |
| | | sessionRestarts.request(replayThreadShutdown.get() || outOfMemory |
| | | ? SessionRestart.NOW : SessionRestart.AFTER_BACKOFF); |
| | | runRequestedSessionRestarts(); |
| | | if (!runRequestedSessionRestarts() && warned) |
| | | { |
| | | /* |
| | | * The warning above said the session is being restarted for the change, and it is not |
| | | * yet: a total update is being processed over that session - almost always an export |
| | | * from this replica, since a total update into it owns the session and is refused |
| | | * above, except by an import which claims its context between that read and this one |
| | | * (issue #1041) - and the restart waits for it, for as long as the total update takes. |
| | | * Said on its own, so that a change which is not delivered again for minutes is not a |
| | | * change nobody asked for. |
| | | * |
| | | * Written where that warning was written and nowhere else: it qualifies that line, so |
| | | * a backend which fails every delivery of a long export would otherwise be one of |
| | | * these per delivery while the warnings they qualify are folded into a count - the |
| | | * very repetition the throttle is there to fold. It is not built at all on the road |
| | | * out of a JVM which has run out of memory, for the reason the warning is not built |
| | | * there: warned is false on it. |
| | | * |
| | | * On the import road of the window above the request is cleared by importBackend() |
| | | * rather than run, so the restart this line announces never runs there. The change is |
| | | * delivered again all the same, which is what an operator reads this for: the import |
| | | * loads the ServerState of the exporter, and the session started at its end asks for |
| | | * everything that state does not cover. |
| | | */ |
| | | logger.info(NOTE_REPLAY_SESSION_RESTART_HELD_BY_TOTAL_UPDATE, csn, getBaseDN()); |
| | | } |
| | | return true; |
| | | } |
| | | |
| | |
| | | */ |
| | | void giveBackChangesParkedByStoppingThread() |
| | | { |
| | | // What it handed back is not read here: this road runs no restart for them, so it has |
| | | // nothing to say about one which waits. |
| | | giveBackParkedChanges(SessionRestart.NOW); |
| | | } |
| | | |
| | | /** |
| | | * Restarts the session as long as changes which could not be replayed are waiting to be |
| | | * delivered again. |
| | | * delivered again - unless a total update is being processed, in which case the requests |
| | | * are left standing for the state checkpointer to run once it is over. |
| | | * <p> |
| | | * A restart stops the session the total update runs over, in either direction. An import |
| | | * into this replica reads its entries from that session and would end on the ones which |
| | | * had arrived - and {@code disabled} does not say an import is running, since |
| | | * {@code preBackendImport()} keeps the backend events this domain is the cause of from |
| | | * disabling it. An export from this replica publishes its entries over it, and |
| | | * {@code exportLDIFEntry()} gives the export up as |
| | | * {@code ERR_INIT_RS_DISCONNECTION_DURING_EXPORT} once the broker has been stopped under |
| | | * it, which leaves the replica it was initializing to be initialized again: minutes on a |
| | | * large backend, spent for a change which would have waited. So the change waits: the |
| | | * request stays standing, the state checkpointer comes for it once a second and runs it |
| | | * as soon as the total update is over ({@link #runPendingSessionRestart()}), and the |
| | | * replication server delivers the change again then. The ServerState waits with it, and |
| | | * the replay of this domain keeps running in the meantime. |
| | | * <p> |
| | | * A total update which begins between this read and any of the stops this call makes is |
| | | * not seen here, and is cut by it: the read and the claim of the import/export context |
| | | * share no lock, which is issue #1041 on the import side. The window is the call rather |
| | | * than a few statements of it: {@code restartSession()} stops the session before it waits |
| | | * its backoff out, so a request taken later in the loop stops the session a backoff wait |
| | | * after the read - up to {@link #MAX_REPLAY_RETRY_DELAY_IN_MS}, and the restart of a |
| | | * request another thread has just been told is held is one of those. Before this the whole |
| | | * of the total update was that window. |
| | | * |
| | | * @return {@code false} when a total update is being processed and the requests were left |
| | | * standing for the state checkpointer, {@code true} otherwise |
| | | */ |
| | | private void runRequestedSessionRestarts() |
| | | private boolean runRequestedSessionRestarts() |
| | | { |
| | | if (ieRunning()) |
| | | { |
| | | return false; |
| | | } |
| | | /* |
| | | * The outer loop is what makes a request which was made while this thread was giving |
| | | * up the recovery its own: the thread which made it found the recovery taken and left |
| | |
| | | replayFailureRecovery.set(false); |
| | | } |
| | | } |
| | | return true; |
| | | } |
| | | |
| | | /** |
| | |
| | | * topology, with the changes it did not replay owned by the replication server and its |
| | | * ServerState stopped behind them. |
| | | * <p> |
| | | * Not run while a total update is being processed, in either direction: a restart stops |
| | | * the session the total update runs over. An import into this replica reads its entries |
| | | * from that session and would end on the ones which had arrived - and {@code disabled} |
| | | * does not say an import is running, since {@code preBackendImport()} keeps the backend |
| | | * events this domain is the cause of from disabling it. An export from this replica |
| | | * publishes its entries over it, and {@code exportLDIFEntry()} gives the export up as |
| | | * {@code ERR_INIT_RS_DISCONNECTION_DURING_EXPORT} once the broker has been stopped |
| | | * under it, which leaves the replica it was initializing to be initialized again. This |
| | | * thread is the one which can afford to wait: the request stays standing, and it comes |
| | | * back here once a second, so the restart is run as soon as the total update is over. |
| | | * The change the restart was asked for waits for as long as the total update takes, and |
| | | * the ServerState with it; the replay of this domain keeps running in the meantime. |
| | | * While a total update is being processed, in either direction, the restart is not run - |
| | | * no restart asked for by a released change is, see {@link #runRequestedSessionRestarts()} |
| | | * - and this thread is the one which can afford to wait for it: the request stays |
| | | * standing, and it comes back here once a second, so the restart is run as soon as the |
| | | * total update is over. |
| | | */ |
| | | private void runPendingSessionRestart() |
| | | { |
| | | if (shutdown.get() || disabled || ieRunning() || !sessionRestarts.isPending()) |
| | | if (shutdown.get() || disabled || !sessionRestarts.isPending()) |
| | | { |
| | | return; |
| | | } |
| | |
| | | * @param restart what the session restart is asked for as: with the backoff a failing |
| | | * backend is owed, or without it on a thread which is stopping or which an |
| | | * OutOfMemoryError is ending |
| | | * @return whether any change was handed back: a change which nobody owns is one only a |
| | | * new delivery brings back, so the caller runs the restart asked for them - on |
| | | * a thread which is not stopping, and on a domain whose session has no owner |
| | | * @return the changes it handed back and asked the restart for, oldest first: a change |
| | | * which nobody owns is one only a new delivery brings back, so the caller runs the |
| | | * restart asked for them - on a thread which is not stopping, and on a domain whose |
| | | * session has no owner - and reports what that restart did for them. Empty when |
| | | * this thread had parked none, and empty on a domain whose session has an owner, |
| | | * where nothing is asked for. The list is the one the release allocated: nothing is |
| | | * allocated for the answer on the road out of a JVM which has run out of memory |
| | | */ |
| | | private boolean giveBackParkedChanges(SessionRestart restart) |
| | | private List<CSN> giveBackParkedChanges(SessionRestart restart) |
| | | { |
| | | final List<CSN> parked = remotePendingChanges.releaseParkedChangesOwnedByCurrentThread(); |
| | | if (parked.isEmpty()) |
| | | { |
| | | return false; |
| | | return parked; |
| | | } |
| | | if (sessionHasAnOwner()) |
| | | { |
| | | // The domain owns its session, or a total update does: both forget the pending |
| | | // changes, and neither leaves a session for this thread to restart. |
| | | return true; |
| | | // changes, and neither leaves a session for this thread to restart. Nothing was asked |
| | | // for here, so the caller has nothing of this road's to run or to report. |
| | | return Collections.emptyList(); |
| | | } |
| | | /* |
| | | * Asked for before the changes are reported: a throw out of the report - the JVM which |
| | |
| | | incProcessedUpdates(); |
| | | logger.info(NOTE_REPLAY_PARKED_CHANGE_GIVEN_BACK, csn, getBaseDN()); |
| | | } |
| | | return true; |
| | | return parked; |
| | | } |
| | | |
| | | /** |