From 661dc06886df4738b9add2206a39d74d7dedf476 Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Thu, 10 Sep 2026 11:57:37 +0000
Subject: [PATCH] [#908] Wait for the changes being applied before a domain going down saves its ServerState (#945)
---
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java | 570 ++++++++++++++++++++++++++++++++++++++++----------------
1 files changed, 404 insertions(+), 166 deletions(-)
diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
index add2d58..4ad0c39 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
+++ b/opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
@@ -57,6 +57,7 @@
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.locks.ReentrantReadWriteLock;
import net.jcip.annotations.GuardedBy;
@@ -143,6 +144,7 @@
import org.opends.server.types.ExistingFileBehavior;
import org.opends.server.types.LDIFExportConfig;
import org.opends.server.types.LDIFImportConfig;
+import org.opends.server.types.LockManager;
import org.opends.server.types.Modification;
import org.opends.server.types.Operation;
import org.opends.server.types.OperationType;
@@ -389,6 +391,75 @@
@GuardedBy("serviceStateLock")
private long sessionGeneration;
/**
+ * Held while a replay thread applies a change of this domain, and taken exclusively by
+ * this domain on its way down.
+ * <p>
+ * A change which reached the backend has to be recorded in the ServerState which is saved
+ * when the domain is disabled or shut down, or it ends up in the data and in no
+ * ServerState: the replication server sends it again and a change which is already
+ * applied is replayed a second time (issue #908). The flags which stop the replay -
+ * {@link #disabled}, {@link #shutdown} - are read under this lock as well, so a domain on
+ * its way down sets one of them and then takes this lock: what it waits for is the
+ * changes which were already being applied, and no attempt starts after that.
+ * <p>
+ * What this closes is the window between an operation reaching the backend and its
+ * {@code commit()}, which is the one issue #908 reports. It is not every road by which a
+ * change can be in the data and in no saved ServerState: {@code commit()} only advances
+ * the state over the changes which are committed from the head of the pending list, so a
+ * change applied while an older one is still to be replayed is forgotten by the
+ * {@code clear()} on the way down and sent again - the barrier of issue #889 seen from
+ * the other side, and a thing to fix where that barrier is rather than here. A change
+ * this replica made itself is a third road, and the one the ServerState recovery in
+ * {@code PersistentServerState.loadState()} already repairs, since it only looks for the
+ * CSNs of this server.
+ * <p>
+ * Deliberately not the fair kind. The read lock is taken for every attempt made on the
+ * backend, and a queued writer blocks the readers which come after it even without
+ * fairness; the readers which do barge past it are the ones which see the flag and leave
+ * without applying anything. The flag is read once before the lock for that reason too -
+ * the replay threads are a pool shared by every domain, and a thread which parks on the
+ * lock of a domain going down is a thread no other domain gets its changes replayed by.
+ * <p>
+ * The lock is released at the end of every attempt, so a domain going down waits for one
+ * attempt and the short backoff which follows it inside the loop, rather than for all
+ * the attempts a delivery is given. The backoff between two attempts is a Thread.sleep()
+ * of tens of milliseconds; the one between two session restarts, which is counted in
+ * seconds, is outside every lock and stays there.
+ */
+ private final ReentrantReadWriteLock replayLock = new ReentrantReadWriteLock();
+ private final ReentrantReadWriteLock.ReadLock replayReadLock = replayLock.readLock();
+ private final ReentrantReadWriteLock.WriteLock replayWriteLock = replayLock.writeLock();
+ /**
+ * How long this domain waits for the replay threads which are applying one of its changes
+ * before it saves its ServerState and goes down.
+ * <p>
+ * Derived from the ceiling the server itself puts on an operation which is waiting for an
+ * entry rather than picked: {@link LockManager} gives the subtree lock and the entry lock
+ * {@link LockManager#DEFAULT_LOCK_TIMEOUT} each, so a replayed change whose target is held
+ * by a concurrent local operation - a client deleting the subtree above it, say - is
+ * inside its attempt for twice that before it gives up with BUSY. A bound under that
+ * ceiling would be spent by ordinary lock contention, and the change which is applied
+ * after it would be the one this whole barrier exists to keep out of that window: no
+ * import, no index rebuild and no wedged backend needed.
+ * <p>
+ * A ceiling rather than a guarantee: an operation also waits for the subtree lock of every
+ * entry above its target, one timeout each, so a deep contended chain outlasts this. What
+ * it buys is that the wait is not lost to the contention a serving backend has anyway.
+ * <p>
+ * It is paid in three places - held under serviceStateLock, inside BackendConfigManager's
+ * write lock when a backend is being deregistered, and once per domain by a server going
+ * down - which is why it is bounded at all. It is only ever spent in full by a replay
+ * which is genuinely stuck: the wait ends the moment the attempt does.
+ */
+ private static final long REPLAY_DRAIN_TIMEOUT_IN_MS =
+ 2 * LockManager.DEFAULT_LOCK_TIMEOUT_UNITS.toMillis(LockManager.DEFAULT_LOCK_TIMEOUT) + 1000;
+ /**
+ * How long this domain waits for the replay of its changes on its way down. Only the
+ * tests, which can not hold a replay thread for {@link #REPLAY_DRAIN_TIMEOUT_IN_MS}, set
+ * another value.
+ */
+ private volatile long replayDrainTimeoutInMs = REPLAY_DRAIN_TIMEOUT_IN_MS;
+ /**
* Stands for "the alert about a change this replica gave up on was never sent". The
* time it is compared with only moves forward from an origin which is arbitrary, so
* zero is not far enough in the past to say it.
@@ -556,7 +627,27 @@
Thread.currentThread().interrupt();
}
}
- state.save();
+ /*
+ * A disabled domain saved its ServerState and cleared it from memory, and an import
+ * or a restore is about to replace the data: saving here would write that empty state
+ * over the saved one, which is a REPLACE of ds-sync-state with no value at all - the
+ * replica would come back with no ServerState rather than with the one it saved on
+ * its way down. The disabled flag does not catch the total update this replica is the
+ * target of: preBackendImport() sets ignoreBackendInitializationEvent, so disable() is
+ * not called on that road and only the import says the data is being replaced.
+ *
+ * The direction is asked for rather than ieRunning(), which the save in the loop above
+ * settles for: an export leaves the data and the ServerState of this domain alone, and
+ * the replay of this domain keeps running for the whole of it - a remote-requested
+ * export is dispatched to a thread pool for that very reason. Since the save in the
+ * loop is skipped for either direction, this is the only one which persists the
+ * changes replayed since the export began, and there is no repair for them afterwards:
+ * checkAndUpdateServerState() only repairs the CSNs of this server.
+ */
+ if (!disabled && !importInProgress())
+ {
+ state.save();
+ }
done = true;
}
@@ -2355,6 +2446,15 @@
{
if (shutdown.compareAndSet(false, true))
{
+ /*
+ * Wait for the changes which are being applied before the ServerState is flushed for
+ * the last time: the flush thread stopped below is the last thing which saves it, so a
+ * change which reaches the backend after that save is in the data and in no
+ * ServerState (issue #908). The flag was just set, so this waits for the attempts
+ * which had started already and no new one begins.
+ */
+ awaitReplayDrained();
+
final RSUpdater rsUpdater = this.rsUpdater.get();
if (rsUpdater != null)
{
@@ -2475,7 +2575,205 @@
int retryCount = IN_PLACE_REPLAY_ATTEMPTS;
while (!dependency && !replayDone && retryCount-- > 0)
{
- if (replayThreadShutdown.get() || shutdown.get() || disabled)
+ /*
+ * The flag which says this domain is going down is read before the lock as well
+ * as under it. The replay threads are a pool shared by every domain of this
+ * server, so a thread which took a change of a domain which is going down should
+ * not queue behind the wait for that domain: the changes of every other domain
+ * are behind it in the same pool.
+ *
+ * The read under the lock is the one which decides; the one above it is a
+ * scheduling optimisation for the common case and nothing more. It cannot keep
+ * this thread out of the queue: the flag can be set and the writer can queue
+ * between the two reads, and a reader which arrives behind a queued writer blocks
+ * even on a lock which is not the fair kind.
+ */
+ boolean goingDown = replayThreadShutdown.get() || shutdown.get() || disabled;
+ if (!goingDown)
+ {
+ /*
+ * Every attempt made on the backend is under this lock, and so is the decision
+ * to make one: a domain on its way down takes it exclusively once it has set
+ * the flag read here, so a change which reaches the backend is recorded in the
+ * ServerState which is saved on the way down, or is not applied at all
+ * (issue #908).
+ */
+ replayReadLock.lock();
+ try
+ {
+ goingDown = replayThreadShutdown.get() || shutdown.get() || disabled;
+ if (!goingDown)
+ {
+ if (!firstAttempt)
+ {
+ /*
+ * Every attempt runs an operation of its own. An Operation which already ran
+ * carries the request controls and the access log items of that run, so
+ * re-running the same one stacks one ManageDsaIT control - and one access log
+ * record - per attempt. It also picks up the new state of the UpdateMsg when
+ * conflict resolution rewrote it.
+ * Note: When msg is a DeleteMsg, the DeleteOperation is properly created
+ * with subtreeDelete request control when needed.
+ */
+ nextOp = msg.createOperation(conn);
+ }
+ firstAttempt = false;
+
+ // Try replay the operation
+ op = nextOp;
+ op.setInternalOperation(true);
+ op.setSynchronizationOperation(true);
+
+ // Always add the ManageDSAIT control so that updates to referrals
+ // are processed locally.
+ op.addRequestControl(new LDAPControl(OID_MANAGE_DSAIT_CONTROL));
+
+ // Warning: specific processing ahead. See OPENDJ-2792
+ if (op instanceof ModifyOperation)
+ {
+ ModifyOperation modifyOperation = (ModifyOperation) op;
+ if (modifyOperation.getEntryDN().equals(SET_PERMISSIVE_MODIFY_FOR_DN))
+ {
+ op.addRequestControl(new LDAPControl(OID_PERMISSIVE_MODIFY_CONTROL));
+ }
+ }
+
+ csn = OperationContext.getCSN(op);
+ op.run();
+
+ ResultCode result = op.getResultCode();
+
+ if (result != ResultCode.SUCCESS)
+ {
+ if (result == ResultCode.NO_OPERATION)
+ {
+ // Pre-operation conflict resolution detected that the operation
+ // was a no-op. For example, an add which has already been
+ // replayed, or a modify DN operation on an entry which has been
+ // renamed by a more recent modify DN.
+ // The change is in the data: push it to the serverState.
+ replayDone = true;
+ recordChangeResolved(csn);
+ }
+ else if (result == ResultCode.BUSY)
+ {
+ /*
+ * We probably could not get a lock (OPENDJ-885). Give the server
+ * another chance to process this operation immediately.
+ */
+ Thread.yield();
+ continue;
+ }
+ else if (isServerFailure(result, serverErrorResultCode))
+ {
+ /*
+ * It can happen when a rebuild is performed or the backend is
+ * offline (OPENDJ-49), or when the storage failed to serve the
+ * operation. Give the server another chance to process this
+ * operation after some time.
+ */
+ Thread.sleep(50);
+ continue;
+ }
+ else
+ {
+ ConflictResolution resolution = ConflictResolution.NOTHING_TO_DO;
+ if (op instanceof ModifyOperation)
+ {
+ ModifyOperation castOp = (ModifyOperation) op;
+ dependency = remotePendingChanges.checkDependencies(castOp);
+ ModifyMsg modifyMsg = (ModifyMsg) msg;
+ resolution = dependency ? resolution : solveNamingConflict(castOp, modifyMsg);
+ }
+ else if (op instanceof DeleteOperation)
+ {
+ DeleteOperation castOp = (DeleteOperation) op;
+ dependency = remotePendingChanges.checkDependencies(castOp);
+ resolution = dependency ? resolution : solveNamingConflict(castOp, msg);
+ }
+ else if (op instanceof AddOperation)
+ {
+ AddOperation castOp = (AddOperation) op;
+ AddMsg addMsg = (AddMsg) msg;
+ dependency = remotePendingChanges.checkDependencies(castOp);
+ resolution = dependency ? resolution : solveNamingConflict(castOp, addMsg);
+ }
+ else if (op instanceof ModifyDNOperation)
+ {
+ ModifyDNOperation castOp = (ModifyDNOperation) op;
+ ModifyDNMsg modifyDNMsg = (ModifyDNMsg) msg;
+ dependency = remotePendingChanges.checkDependencies(modifyDNMsg);
+ resolution = dependency ? resolution : solveNamingConflict(castOp, modifyDNMsg);
+ }
+ // else: unknown type of operation ?! there is nothing to replay
+
+ if (!dependency)
+ {
+ switch (resolution)
+ {
+ case NOTHING_TO_DO:
+ // the update became a dummy update and the result
+ // of the conflict resolution phase is to do nothing.
+ // however we still need to push this change to the serverState
+ replayDone = true;
+ recordChangeResolved(csn);
+ break;
+
+ case FAILED:
+ if (serverErrorResultCode.equals(result))
+ {
+ /*
+ * The result code is the one this server puts on an internal error and is
+ * one conflict resolution knows how to solve, so the change was left to it
+ * rather than treated as a failure of the server: it had its chance and
+ * could not solve it, so the storage failing is what is left. Give it the
+ * in-place attempts an UNAVAILABLE gets - a storage busy for a moment must
+ * not cost a session restart - and leave the change out of the ServerState
+ * once they are spent, which the failure of the server below the loop
+ * reports and acts on, reading the result of the attempt which spent the
+ * last of them. A change which is not in the data must not advance the
+ * ServerState (issue #889).
+ */
+ Thread.sleep(50);
+ break;
+ }
+ /*
+ * The operation did not fail on a naming conflict and not on the server
+ * either: the change can not be applied on this replica. Skip it so that the
+ * replica keeps replaying the changes which follow, but report the error in
+ * the ack and tell the administrator that the data now diverge.
+ */
+ final LocalizableMessage errorMsg = ERR_ERROR_REPLAYING_OPERATION.get(
+ op, csn, result, op.getErrorMessage());
+ logger.error(errorMsg);
+ replayErrorMsg = errorMsg.toString();
+ replayDone = true;
+ skipUnreplayableChange(csn, errorMsg);
+ break;
+
+ default:
+ /*
+ * Try replaying the change again: the next attempt creates an operation
+ * reflecting the new state of the UpdateMsg after conflict resolution
+ * modified it, and dependencies might have been replayed by now.
+ */
+ break;
+ }
+ }
+ }
+ }
+ else
+ {
+ replayDone = true;
+ }
+ }
+ }
+ finally
+ {
+ replayReadLock.unlock();
+ }
+ }
+ if (goingDown)
{
/*
* Either this replay thread or this domain is going away, or the domain is
@@ -2500,168 +2798,6 @@
replayDone = true;
break;
}
- if (!firstAttempt)
- {
- /*
- * Every attempt runs an operation of its own. An Operation which already ran
- * carries the request controls and the access log items of that run, so
- * re-running the same one stacks one ManageDsaIT control - and one access log
- * record - per attempt. It also picks up the new state of the UpdateMsg when
- * conflict resolution rewrote it.
- * Note: When msg is a DeleteMsg, the DeleteOperation is properly created
- * with subtreeDelete request control when needed.
- */
- nextOp = msg.createOperation(conn);
- }
- firstAttempt = false;
-
- // Try replay the operation
- op = nextOp;
- op.setInternalOperation(true);
- op.setSynchronizationOperation(true);
-
- // Always add the ManageDSAIT control so that updates to referrals
- // are processed locally.
- op.addRequestControl(new LDAPControl(OID_MANAGE_DSAIT_CONTROL));
-
- // Warning: specific processing ahead. See OPENDJ-2792
- if (op instanceof ModifyOperation)
- {
- ModifyOperation modifyOperation = (ModifyOperation) op;
- if (modifyOperation.getEntryDN().equals(SET_PERMISSIVE_MODIFY_FOR_DN))
- {
- op.addRequestControl(new LDAPControl(OID_PERMISSIVE_MODIFY_CONTROL));
- }
- }
-
- csn = OperationContext.getCSN(op);
- op.run();
-
- ResultCode result = op.getResultCode();
-
- if (result != ResultCode.SUCCESS)
- {
- if (result == ResultCode.NO_OPERATION)
- {
- // Pre-operation conflict resolution detected that the operation
- // was a no-op. For example, an add which has already been
- // replayed, or a modify DN operation on an entry which has been
- // renamed by a more recent modify DN.
- // The change is in the data: push it to the serverState.
- replayDone = true;
- recordChangeResolved(csn);
- }
- else if (result == ResultCode.BUSY)
- {
- /*
- * We probably could not get a lock (OPENDJ-885). Give the server
- * another chance to process this operation immediately.
- */
- Thread.yield();
- continue;
- }
- else if (isServerFailure(result, serverErrorResultCode))
- {
- /*
- * It can happen when a rebuild is performed or the backend is
- * offline (OPENDJ-49), or when the storage failed to serve the
- * operation. Give the server another chance to process this
- * operation after some time.
- */
- Thread.sleep(50);
- continue;
- }
- else
- {
- ConflictResolution resolution = ConflictResolution.NOTHING_TO_DO;
- if (op instanceof ModifyOperation)
- {
- ModifyOperation castOp = (ModifyOperation) op;
- dependency = remotePendingChanges.checkDependencies(castOp);
- ModifyMsg modifyMsg = (ModifyMsg) msg;
- resolution = dependency ? resolution : solveNamingConflict(castOp, modifyMsg);
- }
- else if (op instanceof DeleteOperation)
- {
- DeleteOperation castOp = (DeleteOperation) op;
- dependency = remotePendingChanges.checkDependencies(castOp);
- resolution = dependency ? resolution : solveNamingConflict(castOp, msg);
- }
- else if (op instanceof AddOperation)
- {
- AddOperation castOp = (AddOperation) op;
- AddMsg addMsg = (AddMsg) msg;
- dependency = remotePendingChanges.checkDependencies(castOp);
- resolution = dependency ? resolution : solveNamingConflict(castOp, addMsg);
- }
- else if (op instanceof ModifyDNOperation)
- {
- ModifyDNOperation castOp = (ModifyDNOperation) op;
- ModifyDNMsg modifyDNMsg = (ModifyDNMsg) msg;
- dependency = remotePendingChanges.checkDependencies(modifyDNMsg);
- resolution = dependency ? resolution : solveNamingConflict(castOp, modifyDNMsg);
- }
- // else: unknown type of operation ?! there is nothing to replay
-
- if (!dependency)
- {
- switch (resolution)
- {
- case NOTHING_TO_DO:
- // the update became a dummy update and the result
- // of the conflict resolution phase is to do nothing.
- // however we still need to push this change to the serverState
- replayDone = true;
- recordChangeResolved(csn);
- break;
-
- case FAILED:
- if (serverErrorResultCode.equals(result))
- {
- /*
- * The result code is the one this server puts on an internal error and is
- * one conflict resolution knows how to solve, so the change was left to it
- * rather than treated as a failure of the server: it had its chance and
- * could not solve it, so the storage failing is what is left. Give it the
- * in-place attempts an UNAVAILABLE gets - a storage busy for a moment must
- * not cost a session restart - and leave the change out of the ServerState
- * once they are spent, which the failure of the server below the loop
- * reports and acts on, reading the result of the attempt which spent the
- * last of them. A change which is not in the data must not advance the
- * ServerState (issue #889).
- */
- Thread.sleep(50);
- break;
- }
- /*
- * The operation did not fail on a naming conflict and not on the server
- * either: the change can not be applied on this replica. Skip it so that the
- * replica keeps replaying the changes which follow, but report the error in
- * the ack and tell the administrator that the data now diverge.
- */
- final LocalizableMessage errorMsg = ERR_ERROR_REPLAYING_OPERATION.get(
- op, csn, result, op.getErrorMessage());
- logger.error(errorMsg);
- replayErrorMsg = errorMsg.toString();
- replayDone = true;
- skipUnreplayableChange(csn, errorMsg);
- break;
-
- default:
- /*
- * Try replaying the change again: the next attempt creates an operation
- * reflecting the new state of the UpdateMsg after conflict resolution
- * modified it, and dependencies might have been replayed by now.
- */
- break;
- }
- }
- }
- }
- else
- {
- replayDone = true;
- }
}
if (!replayDone && !dependency)
@@ -3855,11 +3991,29 @@
{
synchronized (serviceStateLock)
{
- state.save();
- state.clearInMemory();
+ /*
+ * The replay is stopped before the ServerState is saved, and not the other way round:
+ * a change a replay thread is applying has to be either recorded in the state which
+ * is about to be saved or not applied at all, or it ends up in the data and in no
+ * ServerState (issue #908). The flag keeps the attempts which have not started from
+ * starting - it is read under the same lock as the attempt it guards - and the wait
+ * below is for the ones which had started already.
+ *
+ * All of it stays under serviceStateLock, so that this and enable() remain the
+ * mutually exclusive pair they have always been: an enable() which ran in the middle
+ * of this would clear the flag and bring a session up, and this would then go on to
+ * cut that session and clear a ServerState which the flush thread - reading a flag
+ * which says the domain is enabled - would write back empty. That is what bounds
+ * REPLAY_DRAIN_TIMEOUT_IN_MS: the wait is held under a lock which a session restart,
+ * a configuration change and the shutdown of this domain take, and it runs inside
+ * BackendConfigManager's write lock when a backend is being deregistered.
+ */
disabled = true;
disableService(); // This will cut the session and wake up the listener
sessionGeneration++;
+ awaitReplayDrained();
+ state.save();
+ state.clearInMemory();
/*
* The ServerState this bookkeeping goes with is now gone from memory and is loaded
* again from the backend when the domain is enabled back, so the changes listed as
@@ -3880,6 +4034,90 @@
}
/**
+ * Waits for the replay threads which are applying a change of this domain to be done
+ * with it.
+ * <p>
+ * Called once {@link #disabled} or {@link #shutdown} has been set, which is what bounds
+ * the wait: a replay thread reads those under {@link #replayReadLock}, the lock this
+ * takes exclusively, so no attempt starts once this returns and what it waits for is the
+ * attempts which were running already. The lock is released before returning for the
+ * same reason - what keeps the replay out is the flag, not the lock.
+ */
+ private void awaitReplayDrained()
+ {
+ boolean drained = false;
+ boolean interrupted = false;
+ try
+ {
+ drained = replayWriteLock.tryLock(replayDrainTimeoutInMs, TimeUnit.MILLISECONDS);
+ }
+ catch (InterruptedException e)
+ {
+ /*
+ * Give up waiting, and put the interrupt back rather than swallow it: whoever
+ * interrupted this thread - the server going down, a thread pool taking its threads
+ * away - is still waiting for it to stop, and this is not the last thing it does.
+ * The cost is on the shutdown road, whose wait for the last ServerState flush is a
+ * Thread.sleep() which ends on its first call once the flag is set: the flush thread
+ * still runs that save, this one just stops waiting for it.
+ */
+ interrupted = true;
+ Thread.currentThread().interrupt();
+ }
+ if (drained)
+ {
+ replayWriteLock.unlock();
+ return;
+ }
+ /*
+ * The change which is being applied may reach the backend without being recorded in the
+ * ServerState which is saved next, so the replication server sends it again and it is
+ * replayed a second time. Better than holding an administrative task - an import, a
+ * restore, a backend being taken offline - for as long as a backend which stopped
+ * answering takes to answer.
+ *
+ * An interrupted wait is reported as what it is: it says nothing about how long the
+ * replay of this domain takes, and the timeout it never spent would have an operator
+ * reading a backend which is slow into it.
+ */
+ if (interrupted)
+ {
+ logger.warn(WARN_REPLAY_DRAIN_INTERRUPTED, getBaseDN());
+ }
+ else
+ {
+ logger.warn(WARN_REPLAY_NOT_DRAINED, getBaseDN(), replayDrainTimeoutInMs);
+ }
+ }
+
+ /**
+ * Returns how long this domain waits for the replay threads which are applying one of
+ * its changes before it saves its ServerState and goes down.
+ *
+ * @return the timeout in milliseconds
+ */
+ @VisibleForTesting
+ public long getReplayDrainTimeout()
+ {
+ return replayDrainTimeoutInMs;
+ }
+
+ /**
+ * Sets how long this domain waits for the replay threads which are applying one of its
+ * changes before it saves its ServerState and goes down.
+ * <p>
+ * Only there for the tests which check what a domain does when that wait runs out: they
+ * can not hold a replay thread for {@link #REPLAY_DRAIN_TIMEOUT_IN_MS}.
+ *
+ * @param timeoutInMs the timeout in milliseconds
+ */
+ @VisibleForTesting
+ public void setReplayDrainTimeout(long timeoutInMs)
+ {
+ replayDrainTimeoutInMs = timeoutInMs;
+ }
+
+ /**
* Do what necessary when the data have changed : load state, load
* generation Id.
* If there is no such information check if there is a
--
Gitblit v1.10.0