opendj-server-legacy/src/main/java/org/opends/server/backends/pdb/PDBStorage.java
@@ -42,6 +42,7 @@ import java.util.NoSuchElementException; import java.util.Objects; import java.util.Set; import java.util.concurrent.TimeUnit; import org.forgerock.i18n.LocalizableMessage; import org.forgerock.i18n.slf4j.LocalizedLogger; @@ -101,7 +102,37 @@ { private static final int IMPORT_DB_CACHE_SIZE = 32 * MB; private static final double MAX_SLEEP_ON_RETRY_MS = 50.0; /** * Number of attempts a {@link WriteableStorageImpl#write} makes before it propagates the conflict to the caller. * <p> * It is a budget of attempts and not of time, so it is only ever reached by the conflicts that report quickly. * PersistIt reports a write-write conflict only once it has waited on it, up to * {@code SharedResource.DEFAULT_MAX_WAIT_TIME} - a minute, which this backend never lowers - so a conflict slower * to report than {@link #MAX_RETRY_WINDOW_NANOS} spends the whole window inside its first attempt, is granted the * single replay that window's exemption guarantees, and gives up on the window after two attempts rather than * after this many. */ static final int MAX_RETRIES = 10; /** * Wall-clock budget the replays of a {@link WriteableStorageImpl#write} may spend, in nanoseconds. It is checked * between attempts, so an attempt already running is never interrupted, and never before one replay has been * made: the loop returns after at most this window plus two attempts. It bounds the conflicts that are slow to * report, which {@link #MAX_RETRIES} alone does not - an operation whose own work takes seconds would otherwise * multiply that wait by the attempt count. */ static final long MAX_RETRY_WINDOW_NANOS = 10L * 1000L * 1000L * 1000L; //10 s /** * Upper bound of the random delay before the second attempt, in milliseconds; it doubles with every attempt. This * is the bound of the flat sleep this loop took before it was bounded, so the first replay is delayed exactly as * it was and only the later ones back off. */ private static final double BASE_SLEEP_ON_RETRY_MS = 50.0; /** Upper bound the doubled delay is capped at, in milliseconds. */ private static final double MAX_SLEEP_ON_RETRY_MS = 1000.0; private static final String VOLUME_NAME = "dj"; private static final String JOURNAL_NAME = VOLUME_NAME + "_journal"; /** The buffer / page size used by the PersistIt storage. */ @@ -635,8 +666,11 @@ public void write(WriteOperation operation) throws Exception { final Transaction txn = db.getTransaction(); for (;;) final long startedAt = System.nanoTime(); final long giveUpAt = startedAt + retryWindowNanos; for (int attempt = 1;; attempt++) { final RollbackException conflict; txn.begin(); try { @@ -653,8 +687,7 @@ } catch (final RollbackException e) { // retry after random sleep (reduces transactions collision. Drawback: increased latency) Thread.sleep((long) (Math.random() * MAX_SLEEP_ON_RETRY_MS)); conflict = e; } catch (final Exception e) { @@ -665,6 +698,68 @@ { txn.end(); } // decided and slept for outside the try statement: the sleep used to run before the finally ended the // rolled back transaction, holding it open for the whole backoff and lengthening the window every other // writer collides with //System.nanoTime() - giveUpAt is the overflow safe form of the comparison, and attempt > 1 keeps the //window from ending the loop before a single replay: persistit reports a write-write conflict only once //it has waited on it, up to SharedResource.DEFAULT_MAX_WAIT_TIME - a minute, which this backend never //lowers - so one attempt can outlast the window on its own, and it is the attempt after that one which //is likeliest to succeed, the transaction that blocked it having just finished //one clock sample for both, so that the elapsed time reported is the one the give up was decided on final long now = System.nanoTime(); final boolean capSpent = attempt >= maxRetries; if (capSpent || (attempt > 1 && now - giveUpAt >= 0)) { final long elapsedMs = TimeUnit.NANOSECONDS.toMillis(now - startedAt); //which of the two bounds was spent, so that the config change paths - which report this as a bare stack //trace with no message id - say whether raising the attempts or the window is what would have helped final String boundSpent = capSpent ? "attempt cap" : "retry window"; final StorageRuntimeException spent = new StorageRuntimeException( "pdb: backend '" + config.getBackendId() + "' did not apply the transaction after " + attempt + " attempts in " + elapsedMs + " ms, the " + boundSpent + " being spent; the last conflict was " + conflict); // the conflict is suppressed rather than made the cause, because a cause is what every caller strips // this message off with: write(WriteOperation) below unwraps a StorageRuntimeException that carries one // and throws the cause in its place, and EntryContainer.throwAllowedExceptionTypes:1121 rethrows a // StorageRuntimeException unchanged only while getCause() is null, wrapping it a second time otherwise. // Either way the caller would be left holding a bare RollbackException, whose StorageRuntimeException // message is only its class name - which is all ERR_OPEN_ENV_FAIL would then print at startup spent.addSuppressed(conflict); //warned once, at exhaustion only, unlike JDBCStorage which warns on every replay: a conflict is routine //on the ordinary add and modify path of this engine and a line per replay would flood the log. It names //the bound that was spent for the same reason the exception does, and it is the only rendering that can //carry the stack of the conflict: stackTraceToSingleLineString, the form the config change paths report //this exception with, walks the causes and never prints a suppressed exception logger.warn(LocalizableMessage.raw("pdb: giving up on the transaction of backend '%s' after %d attempts" + " in %d ms, the %s being spent: %s", config.getBackendId(), attempt, elapsedMs, boundSpent, stackTraceToSingleLineString(conflict))); throw spent; } if (logger.isTraceEnabled()) { logger.trace("pdb: replaying the transaction after %s, attempt %d of %d", conflict, attempt, maxRetries); } try { // retry after random sleep (reduces transactions collision. Drawback: increased latency), growing with // every attempt so that a contention the first delays did not outlast still has a chance to clear Thread.sleep(retryDelayMillis(attempt)); } catch (final InterruptedException e) { //sleep cleared the interrupt flag: restore it, and report the conflict being replayed rather than the //interrupt, which would hide from the caller what actually went wrong. Wrapped the way the exhausted //loop above wraps it, and for the same reason: a RollbackException carries no message of its own, so //every caller that wraps one reports nothing but its class name Thread.currentThread().interrupt(); final StorageRuntimeException interrupted = new StorageRuntimeException( "pdb: backend '" + config.getBackendId() + "' was interrupted while replaying the transaction after " + attempt + " attempts; the last conflict was " + conflict); interrupted.addSuppressed(conflict); interrupted.addSuppressed(e); throw interrupted; } } } } @@ -915,6 +1010,10 @@ private PDBMonitor monitor; private MemoryQuota memQuota; private StorageStatus storageStatus = StorageStatus.working(); /** Attempt bound of a {@link WriteableStorageImpl#write}, {@link #MAX_RETRIES} outside the tests. */ private final int maxRetries; /** Wall-clock bound of a {@link WriteableStorageImpl#write}, {@link #MAX_RETRY_WINDOW_NANOS} outside the tests. */ private final long retryWindowNanos; /** * Creates a new persistit storage with the provided configuration. @@ -928,7 +1027,34 @@ // FIXME: should be package private once importer is decoupled. public PDBStorage(final PDBBackendCfg cfg, ServerContext serverContext) throws ConfigException { this(cfg, serverContext, MAX_RETRIES, MAX_RETRY_WINDOW_NANOS); } /** * Creates a new persistit storage whose replay bounds are the given ones rather than {@link #MAX_RETRIES} and * {@link #MAX_RETRY_WINDOW_NANOS}. * <p> * Only a test builds one of these, and it does so to stop the two bounds racing each other: with the shipped * values a run of replays spends a random share of the window on backoff alone, so a test of the attempt cap * can be ended by the window on a loaded machine, and a test of the window has to spend seconds of build time * to reach it. * * @param cfg * The configuration. * @param serverContext * This server instance context * @param maxRetries * Number of attempts a write makes before it propagates the conflict to the caller. * @param retryWindowNanos * Wall-clock budget the replays of a write may spend, in nanoseconds. * @throws ConfigException if memory cannot be reserved */ PDBStorage(final PDBBackendCfg cfg, ServerContext serverContext, int maxRetries, long retryWindowNanos) throws ConfigException { this.serverContext = serverContext; this.maxRetries = maxRetries; this.retryWindowNanos = retryWindowNanos; backendDirectory = getBackendDirectory(cfg); config = cfg; cfg.addPDBChangeListener(this); @@ -1095,6 +1221,24 @@ return new ImporterImpl(); } /** * {@inheritDoc} * <p> * A transaction the engine rolled back is replayed, bounded twice: by {@link #MAX_RETRIES} attempts and by the * {@link #MAX_RETRY_WINDOW_NANOS} wall-clock window, whichever is spent first - except that the window alone * never ends the replays before one has been made. It is bounded because the * configuration change paths of the pluggable backend hold an entry container's exclusive lock across this * method, and every reader of that suffix then waits - untimed and uninterruptibly - until it returns, so a * conflict that never clears would park every worker thread of that suffix rather than fail one operation. * <p> * Once the bound is spent the conflict is reported as a {@link StorageRuntimeException} naming the backend, the * attempts spent, the time they took and which of the two bounds ran out. It carries the conflict as a * suppressed exception rather than as its cause: a cause is unwrapped below and thrown in its place, and * {@code EntryContainer.throwAllowedExceptionTypes} likewise passes a {@link StorageRuntimeException} through * untouched only while it has no cause. Given a cause, both hand the caller a bare RollbackException instead, * and the message of a {@link StorageRuntimeException} wrapping one is just its class name - which is all * {@code ERR_OPEN_ENV_FAIL} would report when this happens as a backend starts. */ @Override public void write(final WriteOperation operation) throws Exception { @@ -1108,6 +1252,14 @@ } } /** Returns the randomized delay before the given attempt is replayed, doubling with each attempt up to a cap. */ //package private like the JDBCStorage copy it is taken from, so that a test can pin the growth and the cap static long retryDelayMillis(int attempt) { final double bound = Math.min(MAX_SLEEP_ON_RETRY_MS, BASE_SLEEP_ON_RETRY_MS * (1 << Math.min(attempt - 1, 5))); return (long) (Math.random() * bound); } private Exception unwrap(StorageRuntimeException e) throws Exception { if (e.getCause() != null) opendj-server-legacy/src/main/java/org/opends/server/backends/pluggable/spi/Storage.java
@@ -73,13 +73,24 @@ <T> T read(ReadOperation<T> readOperation) throws Exception; /** * Executes a write operation. In case of a write operation rollback, implementations must ensure * the write operation is retried until it succeeds. * Executes a write operation. In case of a write operation rollback, implementations may replay the write * operation rather than propagate the failure: a {@link WriteOperation} is required to be idempotent for * exactly that reason. A replay must be bounded - by a number of attempts, by a window of time, or by both - * so that a conflict which does not clear reaches the caller instead of being retried forever. The pluggable * backend holds locks across this method, up to the exclusive lock of an entry container, and every thread * waiting on one of those locks waits for as long as this method does. * <p> * A caller that mutates state around this method must handle that bound being spent. Removing an entry from an * in-memory map before the write so that a replay still finds the work to do, or reading configuration back out * of the operation once it returns, both assume the write is applied; when it is not, this method throws with * that state already changed and the transaction not applied, and the caller is the only place that can reconcile * the two. * * @param writeOperation * the write operation to execute * @throws Exception * if a problem occurs with the underlying storage engine * if a problem occurs with the underlying storage engine, including a conflict that outlasted the * replays the implementation makes */ void write(WriteOperation writeOperation) throws Exception; opendj-server-legacy/src/test/java/org/opends/server/backends/pdb/PDBStorageTest.java
@@ -12,6 +12,7 @@ * information: "Portions Copyright [year] [name of copyright owner]". * * Copyright 2016 ForgeRock AS. * Portions Copyright 2026 3A Systems, LLC. */ package org.opends.server.backends.pdb; @@ -21,11 +22,17 @@ import static org.opends.server.util.StaticUtils.*; import static org.forgerock.opendj.ldap.ByteString.*; import java.util.concurrent.atomic.AtomicInteger; import org.forgerock.opendj.config.server.ConfigException; import org.forgerock.opendj.ldap.ByteString; import org.opends.server.DirectoryServerTestCase; import org.opends.server.TestCaseUtils; import org.forgerock.opendj.server.config.server.PDBBackendCfg; import org.opends.server.backends.pluggable.spi.AccessMode; import org.opends.server.backends.pluggable.spi.ReadOperation; import org.opends.server.backends.pluggable.spi.ReadableTransaction; import org.opends.server.backends.pluggable.spi.StorageRuntimeException; import org.opends.server.backends.pluggable.spi.TreeName; import org.opends.server.backends.pluggable.spi.WriteOperation; import org.opends.server.backends.pluggable.spi.WriteableTransaction; @@ -38,10 +45,19 @@ import org.testng.annotations.Test; import com.persistit.Exchange; import com.persistit.exception.RollbackException; public class PDBStorageTest extends DirectoryServerTestCase { /** A window no run of replays can spend, so that a test of the attempt cap is only ever ended by the cap. */ private static final long UNREACHABLE_RETRY_WINDOW_NANOS = 300L * 1000L * 1000L * 1000L; //5 min /** A window a single attempt outlasts, so that a test of the window reaches it without seconds of build time. */ private static final long SHORT_RETRY_WINDOW_NANOS = 200L * 1000L * 1000L; //200 ms /** An attempt long enough to outlast {@link #SHORT_RETRY_WINDOW_NANOS} on its own, in milliseconds. */ private static final long ATTEMPT_LONGER_THAN_SHORT_WINDOW_MS = 300; private final TreeName treeName = new TreeName("dc=test", "test"); private ServerContext serverContext; private PDBStorage storage; @BeforeClass @@ -53,18 +69,72 @@ @BeforeMethod public void setUp() throws ConfigException { ServerContext serverContext = mock(ServerContext.class); serverContext = mock(ServerContext.class); when(serverContext.getMemoryQuota()).thenReturn(new MemoryQuota()); when(serverContext.getDiskSpaceMonitor()).thenReturn(mock(DiskSpaceMonitor.class)); storage = new PDBStorage(createBackendCfg(), serverContext); // the volume is removed on the way in as well as on the way out: a build whose JVM died never ran tearDown(), // and this class shares a fixed db-directory across methods and across builds, so what that run left behind // would still be here to answer this method's reads storage.removeStorageFiles(); storage.open(AccessMode.READ_WRITE); } @AfterMethod public void tearDown() { storage.close(); closeAndRemove(storage); } /** * Closes the storage and removes its volume, keeping whichever of the two failed first. Removing it from a * finally would let a removal failure replace the close() failure (JLS 14.20.2) - and a close() that throws is * exactly the case the removal is here for. */ private static void closeAndRemove(PDBStorage storage) { RuntimeException failure = null; try { storage.close(); } catch (RuntimeException e) { failure = e; } try { storage.removeStorageFiles(); } catch (RuntimeException e) { if (failure == null) { failure = e; } else { failure.addSuppressed(e); } } if (failure != null) { throw failure; } } /** * Replaces the storage under test with one bounded by the given values, so that the bound a test is about is * the one that ends its replays. With the shipped values the two race: the nine backoffs of a full ladder draw * from 50+100+200+400+800+1000x4, so an attempt cap test can be ended by the ten second window instead, and a * window test has to make every attempt outlast seconds of that window to reach it. */ private void reopenWithReplayBounds(int maxRetries, long retryWindowNanos) throws Exception { closeAndRemove(storage); storage = new PDBStorage(createBackendCfg(), serverContext, maxRetries, retryWindowNanos); storage.open(AccessMode.READ_WRITE); } @Test @@ -122,6 +192,260 @@ assertThat(storage.getNewExchange(treeName, true)).isNotSameAs(initial); } @Test public void testWriteGivesUpAfterTheAttemptCap() throws Exception { // the shipped cap, against a window the ladder of backoffs cannot reach: on the shipped window those nine // backoffs draw from up to 5550 ms, so a loaded machine ends this loop on the window and the cap goes untested reopenWithReplayBounds(PDBStorage.MAX_RETRIES, UNREACHABLE_RETRY_WINDOW_NANOS); createTree(); final RollbackException conflict = new RollbackException(); final AtomicInteger attempts = new AtomicInteger(); try { storage.write(new WriteOperation() { @Override public void run(WriteableTransaction txn) throws Exception { attempts.incrementAndGet(); txn.put(treeName, valueOfUtf8("abandoned"), valueOfUtf8("value")); throw conflict; } }); failBecauseExceptionWasNotThrown(StorageRuntimeException.class); } catch (StorageRuntimeException e) { assertThat(e.getSuppressed()).contains(conflict); } assertThat(attempts.get()).isEqualTo(PDBStorage.MAX_RETRIES); assertThat(read("abandoned")).isNull(); } @Test public void testWriteIsReplayedUntilTheConflictClears() throws Exception { createTree(); final AtomicInteger attempts = new AtomicInteger(); storage.write(new WriteOperation() { @Override public void run(WriteableTransaction txn) throws Exception { if (attempts.incrementAndGet() <= 3) { throw new RollbackException(); } txn.put(treeName, valueOfUtf8("applied"), valueOfUtf8("value")); } }); assertThat(attempts.get()).isEqualTo(4); assertThat(read("applied")).isEqualTo(valueOfUtf8("value")); } /** * PersistIt reports a write-write conflict only once it has waited on it - up to * {@code SharedResource.DEFAULT_MAX_WAIT_TIME}, a minute, which this backend never lowers - so a single attempt * can outlast the whole window. Giving up on the window alone would then replay nothing, in the very case where * the replay is likeliest to succeed: the transaction that was blocking this one has just finished. */ @Test public void testWriteIsReplayedOnceWhenTheFirstAttemptOutlastsTheWindow() throws Exception { reopenWithReplayBounds(PDBStorage.MAX_RETRIES, SHORT_RETRY_WINDOW_NANOS); createTree(); final AtomicInteger attempts = new AtomicInteger(); storage.write(new WriteOperation() { @Override public void run(WriteableTransaction txn) throws Exception { if (attempts.incrementAndGet() == 1) { Thread.sleep(ATTEMPT_LONGER_THAN_SHORT_WINDOW_MS); throw new RollbackException(); } txn.put(treeName, valueOfUtf8("outlasted"), valueOfUtf8("written")); } }); assertThat(attempts.get()).isEqualTo(2); assertThat(read("outlasted")).isEqualTo(valueOfUtf8("written")); } @Test public void testExhaustedWriteNamesTheAttemptsItSpent() throws Exception { // the message is the same at any cap, so this one is spent in two backoffs rather than in the shipped ladder final int maxRetries = 3; reopenWithReplayBounds(maxRetries, UNREACHABLE_RETRY_WINDOW_NANOS); createTree(); try { storage.write(new WriteOperation() { @Override public void run(WriteableTransaction txn) throws Exception { throw new RollbackException(); } }); failBecauseExceptionWasNotThrown(StorageRuntimeException.class); } catch (StorageRuntimeException e) { assertThat(e.getMessage()).contains("PDBStorageTest").contains(maxRetries + " attempts"); // and which of the two bounds ran out, since the attempt count alone does not say assertThat(e.getMessage()).contains("attempt cap"); // write() unwraps a StorageRuntimeException that carries a cause, which would replace this message with // the bare RollbackException, and it is the message the config change paths report assertThat(e.getCause()).isNull(); } } @Test public void testWriteGivesUpOnTheWindowWhenAttemptsAreSlow() throws Exception { reopenWithReplayBounds(PDBStorage.MAX_RETRIES, SHORT_RETRY_WINDOW_NANOS); createTree(); final AtomicInteger attempts = new AtomicInteger(); try { storage.write(new WriteOperation() { @Override public void run(WriteableTransaction txn) throws Exception { attempts.incrementAndGet(); // a conflict this slow to report spends the wall clock window long before the attempt cap Thread.sleep(ATTEMPT_LONGER_THAN_SHORT_WINDOW_MS); throw new RollbackException(); } }); failBecauseExceptionWasNotThrown(StorageRuntimeException.class); } catch (StorageRuntimeException e) { // the window is what ended it, and it says so: an assertion on the attempt count alone would also pass for // a give up on attempt 1, which is the regression the attempt > 1 exemption exists to prevent assertThat(e.getMessage()).contains("retry window"); } // one attempt beyond the first: the first spends the window, the exemption grants the replay, and the check // after that replay is the one that gives up assertThat(attempts.get()).isEqualTo(2); } @Test public void testInterruptedWriteReportsTheConflictItWasReplaying() throws Exception { createTree(); final RollbackException conflict = new RollbackException(); final AtomicInteger attempts = new AtomicInteger(); final boolean interruptedAfterwards; try { storage.write(new WriteOperation() { @Override public void run(WriteableTransaction txn) throws Exception { attempts.incrementAndGet(); // interrupted here rather than before the write, where the transaction this attempt begins would // report the interrupt itself and the loop would never reach the backoff being tested Thread.currentThread().interrupt(); throw conflict; } }); failBecauseExceptionWasNotThrown(StorageRuntimeException.class); return; } catch (StorageRuntimeException e) { interruptedAfterwards = Thread.interrupted(); // the conflict, not the interrupt, is what the caller is told about - but through the same shape the // exhausted loop uses, since a bare RollbackException reaches every caller as its own class name assertThat(e.getMessage()).contains("PDBStorageTest").contains("interrupted"); assertThat(e.getSuppressed()).contains(conflict).hasAtLeastOneElementOfType(InterruptedException.class); assertThat(e.getCause()).isNull(); } finally { Thread.interrupted(); } // sleep() cleared the flag, so the caller only learns of the interrupt if the loop restores it assertThat(interruptedAfterwards).isTrue(); // one attempt even though the first backoff is a random 0-49 ms and so is sometimes 0: Thread.sleep() checks // the interrupt flag before it checks for a zero duration, so the replay is never reached assertThat(attempts.get()).isEqualTo(1); } /** * The delay grows with the attempt and stays under the cap, so that a contention the first delays did not * outlast still has a chance to clear without the replays overrunning the window on sleep alone. */ @Test public void testRetryDelayGrowsAndStaysBounded() { long previousBound = 0; for (int attempt = 1; attempt <= PDBStorage.MAX_RETRIES; attempt++) { long bound = 0; for (int i = 0; i < 100; i++) { final long delay = PDBStorage.retryDelayMillis(attempt); assertThat(delay).as("attempt %d", attempt).isGreaterThanOrEqualTo(0).isLessThan(1000); bound = Math.max(bound, delay); } if (attempt == 1) { // the flat sleep this loop took before it was bounded, unchanged: only the later attempts back off assertThat(bound).as("attempt 1 delays past the sleep this loop always took").isLessThan(50); } assertThat(bound).as("attempt %d did not grow past attempt %d", attempt, attempt - 1) .isGreaterThanOrEqualTo(previousBound / 2); previousBound = bound; } // and the growth is real rather than a delay that never leaves the first tier long grown = 0; for (int i = 0; i < 100; i++) { grown = Math.max(grown, PDBStorage.retryDelayMillis(PDBStorage.MAX_RETRIES)); } assertThat(grown).as("the last attempts still sleep within the first attempt's bound").isGreaterThan(500); } private void createTree() throws Exception { storage.write(new WriteOperation() { @Override public void run(WriteableTransaction txn) throws Exception { txn.openTree(treeName, true); } }); } private ByteString read(final String key) throws Exception { return storage.read(new ReadOperation<ByteString>() { @Override public ByteString run(ReadableTransaction txn) throws Exception { return txn.read(treeName, valueOfUtf8(key)); } }); } protected PDBBackendCfg createBackendCfg() { PDBBackendCfg backendCfg = mockCfg(PDBBackendCfg.class);