mirror of https://github.com/OpenIdentityPlatform/OpenDJ.git

Valery Kharseko
16 hours ago eef07575153b2a7f37feaaa8e39977f32631082b
[#952] Keep a failed state write from killing the checkpointer and hanging the shutdown (#977)
2 files modified
1 files added
449 ■■■■■ changed files
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java 88 ●●●● patch | view | raw | blame | history
opendj-server-legacy/src/messages/org/opends/messages/replication.properties 5 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ServerStateFlushTest.java 356 ●●●●● patch | view | raw | blame | history
opendj-server-legacy/src/main/java/org/opends/server/replication/plugin/LDAPReplicationDomain.java
@@ -523,11 +523,23 @@
  /** Published by a configuration change, read by the changelog threads without a lock. */
  private volatile ExternalChangelogDomain eclDomain;
  /** A boolean indicating if the thread used to save the persistentServerState is terminated. */
  private volatile boolean done = true;
  private final ServerStateFlush flushThread;
  /**
   * How long {@link #shutdown()} waits for the state checkpointer to stop before it goes on
   * without it. The checkpointer has one last state to write when it is asked to stop, so it
   * normally stops within a modify: a checkpointer which is still writing after this went to
   * a backend which is not answering, and waiting for it any longer would hang the shutdown
   * of the whole server. This is the budget {@code ServerShutdownMonitor} gives a thread
   * before it starts interrupting them. It is spent by every call of {@link #shutdown()}:
   * once per domain when the server goes down, which shuts the domains down one after
   * another, and again by a second caller of a domain whose checkpointer is stuck. It bounds
   * the shutdown of this domain, not of the server: a write which ignores the interrupt
   * still holds the quiescence of a pluggable backend, so the server shutdown waits for it
   * again when it closes that backend; {@code SchemaBackend} has no such wait.
   */
  private static final long FLUSH_THREAD_SHUTDOWN_TIMEOUT_IN_MS = 30000;
  /** The attribute name used to store the generation id in the backend. */
  private static final String REPLICATION_GENERATION_ID = "ds-sync-generation-id";
  /** The attribute name used to store the fractional include configuration in the backend. */
@@ -614,8 +626,6 @@
    @Override
    public void run()
    {
      done = false;
      while (!isShutdownInitiated())
      {
        try
@@ -623,10 +633,10 @@
          synchronized (this)
          {
            wait(1000);
            if (!disabled && !ieRunning())
            {
              state.save();
            }
          }
          if (!disabled && !ieRunning())
          {
            saveState();
          }
        }
        catch (InterruptedException e)
@@ -654,10 +664,37 @@
       */
      if (!disabled && !importInProgress())
      {
        saveState();
      }
    }
    /**
     * Writes the state of the domain to the backend, keeping a failure to itself.
     * <p>
     * A checkpoint which throws is not a reason to stop checkpointing: the state is still
     * marked as unsaved, so the next checkpoint writes it again. The exit save has no next
     * checkpoint: a domain whose last write failed comes back with the last state it did
     * write - its own CSNs repaired from ds-sync-hist by checkAndUpdateServerState(), those
     * of the other replicas as they were - and replays the changes since. Letting the
     * exception out would end this thread - and with it the checkpointing of this domain for
     * the rest of the life of the server, and the {@link LDAPReplicationDomain#shutdown()}
     * which waits for the thread to stop.
     * <p>
     * The write is run outside the monitor of this thread, which
     * {@link LDAPReplicationDomain#shutdown()} takes to wake it up: holding the monitor
     * across a write which does not come back would block a shutdown before it ever reaches
     * the bounded wait it does for this thread.
     */
    private void saveState()
    {
      try
      {
        state.save();
      }
      done = true;
      catch (RuntimeException e)
      {
        logger.error(ERR_CHECKPOINTING_STATE_FAILED, getBaseDN(), stackTraceToSingleLineString(e));
      }
    }
  }
@@ -2518,13 +2555,10 @@
      }
      // stop the thread in charge of flushing the ServerState.
      if (flushThread != null)
      flushThread.initiateShutdown();
      synchronized (flushThread)
      {
        flushThread.initiateShutdown();
        synchronized (flushThread)
        {
          flushThread.notifyAll();
        }
        flushThread.notifyAll();
      }
      DirectoryServer.deregisterAlertGenerator(this);
@@ -2542,12 +2576,26 @@
      }
    }
    // wait for completion of the ServerStateFlush thread.
    /*
     * Wait for completion of the ServerStateFlush thread, but not for longer than the budget
     * it is given: a thread which is gone - killed by an Error on its way to the backend, say
     * - is never going to report that it is done, and a shutdown which waits for it forever
     * takes the shutdown of the server down with it. join() covers both, and a thread which
     * was never started as well.
     *
     * Every caller waits, the one which lost the race above included: returning at once would
     * let it go on while the checkpointer is still writing. What the loser gets is the budget
     * from its own arrival, which starts before the winner has asked the checkpointer to stop
     * - the winner may still be in awaitReplayDrained() - so its wait may end, and log the
     * whole budget as spent, while the winner is still waiting: a second shutdown() of a
     * domain whose checkpointer is stuck can return before the first one.
     */
    try
    {
      while (!done)
      flushThread.join(FLUSH_THREAD_SHUTDOWN_TIMEOUT_IN_MS);
      if (flushThread.isAlive())
      {
        Thread.sleep(50);
        logger.error(ERR_STATE_CHECKPOINTER_NOT_STOPPED, getBaseDN(), FLUSH_THREAD_SHUTDOWN_TIMEOUT_IN_MS);
      }
    } catch (InterruptedException e)
    {
opendj-server-legacy/src/messages/org/opends/messages/replication.properties
@@ -690,3 +690,8 @@
 in domain "%s": the search of the entry with entryUUID %s did not run (%s). The change is not \
 applied on what a search which read nothing seemed to say about the data, and is not recorded \
 as replayed
ERR_CHECKPOINTING_STATE_FAILED_323=Could not write the replication state of domain "%s" : %s. \
 The state stays unsaved: the next checkpoint writes it again, and a domain which was stopping \
 comes back with the last state it did write and replays the changes since
ERR_STATE_CHECKPOINTER_NOT_STOPPED_324=The state checkpointer of domain "%s" has not stopped within \
 %d ms : the shutdown of the domain goes on without it
opendj-server-legacy/src/test/java/org/opends/server/replication/plugin/ServerStateFlushTest.java
New file
@@ -0,0 +1,356 @@
/*
 * The contents of this file are subject to the terms of the Common Development and
 * Distribution License (the License). You may not use this file except in compliance with the
 * License.
 *
 * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the
 * specific language governing permission and limitations under the License.
 *
 * When distributing Covered Software, include this CDDL Header Notice in each file and include
 * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL
 * Header, with the fields enclosed by brackets [] replaced by your own identifying
 * information: "Portions Copyright [year] [name of copyright owner]".
 *
 * Copyright 2026 3A Systems, LLC.
 */
package org.opends.server.replication.plugin;
import static java.util.concurrent.TimeUnit.*;
import static org.opends.server.util.CollectionUtils.*;
import static org.opends.server.util.StaticUtils.*;
import static org.testng.Assert.*;
import java.util.SortedSet;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import org.forgerock.opendj.ldap.DN;
import org.opends.server.TestCaseUtils;
import org.opends.server.api.DirectoryThread;
import org.opends.server.backends.MemoryBackend;
import org.opends.server.core.ModifyOperation;
import org.opends.server.replication.ReplicationTestCase;
import org.opends.server.replication.common.CSN;
import org.opends.server.replication.server.ReplServerFakeConfiguration;
import org.opends.server.replication.server.ReplicationServer;
import org.opends.server.types.DirectoryException;
import org.opends.server.types.Entry;
import org.opends.server.util.TestTimer;
import org.opends.server.util.TestTimer.CallableVoid;
import org.testng.annotations.AfterMethod;
import org.testng.annotations.BeforeMethod;
import org.testng.annotations.Test;
/**
 * Tests the state checkpointer of a {@link LDAPReplicationDomain} against a backend which
 * fails the write of the replication state.
 */
@SuppressWarnings("javadoc")
public class ServerStateFlushTest extends ReplicationTestCase
{
  private static final String BACKEND_ID = "serverStateFlushTest";
  private static final String BASE_DN_STRING = "o=serverStateFlushTest";
  private static final int DS_ID = 1;
  private static final int RS_ID = 601;
  /** How long a test waits for something the checkpointer does on its own, in seconds. */
  private static final int CHECKPOINT_TIMEOUT_IN_SECS = 20;
  private DN baseDN;
  private StateWriteFailureBackend backend;
  private ReplicationServer replicationServer;
  private LDAPReplicationDomain domain;
  private Thread checkpointer;
  private boolean domainDeleted;
  /** What the backend does with a write of the replication state instead of writing it. */
  private enum StateWrite
  {
    /** Write it, the way the backend normally would. */
    SUCCEEDS,
    /** Throw a {@link RuntimeException} the way an unwrapped driver failure would. */
    THROWS_RUNTIME_EXCEPTION,
    /** Throw an {@link Error}, which no {@code catch} of the checkpointer can be expected to hold. */
    THROWS_ERROR,
    /** Block until the test releases it, the way a write to an unresponsive storage would. */
    BLOCKS;
  }
  /** A memory backend whose write of the replication state fails on demand. */
  private static final class StateWriteFailureBackend extends MemoryBackend
  {
    private final DN stateEntryDN;
    private volatile StateWrite stateWrite = StateWrite.SUCCEEDS;
    private final AtomicInteger failedStateWrites = new AtomicInteger();
    /**
     * Those of {@link #failedStateWrites} the checkpointer made once it was asked to stop:
     * the write its wake-up triggers, and the last one it runs on its way out.
     */
    private final AtomicInteger failedStateWritesAfterShutdown = new AtomicInteger();
    private final CountDownLatch blockedWriteStarted = new CountDownLatch(1);
    private final CountDownLatch blockedWriteReleased = new CountDownLatch(1);
    private StateWriteFailureBackend(DN stateEntryDN)
    {
      this.stateEntryDN = stateEntryDN;
    }
    /**
     * Not synchronized, unlike the method it overrides: a write which blocks must not hold
     * the monitor of the backend, or every other operation on it would block with it.
     */
    @Override
    public void replaceEntry(Entry oldEntry, Entry newEntry, ModifyOperation modifyOperation)
        throws DirectoryException
    {
      if (stateEntryDN.equals(newEntry.getName()))
      {
        switch (stateWrite)
        {
        case THROWS_RUNTIME_EXCEPTION:
          recordFailedStateWrite();
          throw new IllegalStateException("injected failure of a replication state write");
        case THROWS_ERROR:
          recordFailedStateWrite();
          throw new OutOfMemoryError("injected failure of a replication state write");
        case BLOCKS:
          blockedWriteStarted.countDown();
          try
          {
            blockedWriteReleased.await();
          }
          catch (InterruptedException e)
          {
            Thread.currentThread().interrupt();
            return;
          }
          break;
        default:
          break;
        }
      }
      super.replaceEntry(oldEntry, newEntry, modifyOperation);
    }
    private void recordFailedStateWrite()
    {
      failedStateWrites.incrementAndGet();
      final Thread writer = Thread.currentThread();
      if (writer instanceof DirectoryThread && ((DirectoryThread) writer).isShutdownInitiated())
      {
        failedStateWritesAfterShutdown.incrementAndGet();
      }
    }
  }
  @BeforeMethod
  public void setUpDomain() throws Exception
  {
    baseDN = DN.valueOf(BASE_DN_STRING);
    backend = new StateWriteFailureBackend(baseDN);
    backend.setBackendID(BACKEND_ID);
    backend.setBaseDNs(baseDN);
    backend.configureBackend(null, TestCaseUtils.getServerContext());
    backend.openBackend();
    TestCaseUtils.getServerContext().getBackendConfigManager().registerLocalBackend(backend);
    backend.addEntry(createEntry(baseDN), null);
    replicationServer = new ReplicationServer(new ReplServerFakeConfiguration(
        TestCaseUtils.findFreePort(), "serverStateFlushTestDb", 0, RS_ID, 0, 100, null));
    final SortedSet<String> replServers = newTreeSet("localhost:" + replicationServer.getReplicationPort());
    domain = MultimasterReplication.createNewDomain(new DomainFakeCfg(baseDN, DS_ID, replServers));
    domainDeleted = false;
    domain.start();
    checkpointer = checkpointerOf(domain);
    assertNotNull(checkpointer, "the state checkpointer of the domain was not started");
    assertTrue(checkpointer.isAlive(), "the state checkpointer of the domain is not running");
  }
  /**
   * Takes back what {@link #setUpDomain()} put in place, one resource at a time: a setup
   * which failed halfway through must not leave a backend or a domain behind for the next
   * test of the class to trip over.
   * <p>
   * Bounded on its own: the {@code timeOut} of a test covers the test alone, and for a test
   * which leaves the domain to this method, the shutdown under test runs here - a shutdown
   * which hangs must fail the test rather than hold the whole run. The bound has room for
   * the budget the shutdown gives the checkpointer and for the wait for the thread below.
   */
  @AfterMethod(timeOut = 90000)
  public void tearDownDomain() throws Exception
  {
    if (backend != null)
    {
      // Let go of a write the test left blocked, and of the checkpointer waiting on it,
      // before the backend it is writing to is taken away.
      backend.stateWrite = StateWrite.SUCCEEDS;
      backend.blockedWriteReleased.countDown();
    }
    if (domain != null && !domainDeleted)
    {
      // Stops the checkpointer within a tick: waiting for it before this would spend the
      // whole wait on a thread which is running by design.
      deleteDomain();
    }
    domain = null;
    if (checkpointer != null)
    {
      checkpointer.join(SECONDS.toMillis(CHECKPOINT_TIMEOUT_IN_SECS));
      checkpointer = null;
    }
    if (replicationServer != null)
    {
      remove(replicationServer);
      replicationServer = null;
    }
    if (backend != null)
    {
      backend.finalizeBackend();
      TestCaseUtils.getServerContext().getBackendConfigManager().deregisterLocalBackend(backend);
      backend = null;
    }
  }
  /**
   * A write of the state which throws must not stop the checkpointer: the state stays unsaved,
   * and the next checkpoint writes it.
   */
  @Test(timeOut = 120000)
  public void checkpointerKeepsCheckpointingAfterAStateWriteThatThrows() throws Exception
  {
    backend.stateWrite = StateWrite.THROWS_RUNTIME_EXCEPTION;
    final CSN csn = newCSN();
    domain.getServerState().update(csn);
    waitForFailedStateWrites(1);
    backend.stateWrite = StateWrite.SUCCEEDS;
    checkEntryHasAttributeValue(baseDN, "ds-sync-state", csn.toString(), CHECKPOINT_TIMEOUT_IN_SECS,
        "the checkpointer did not write the state after a write which threw");
  }
  /**
   * Shutting the domain down must not wait forever for a checkpointer whose writes all
   * throw, and the checkpointer must hold the exception of the last write it runs on its
   * way out as it does the others: that write used to end the thread before it reported
   * that it was done.
   */
  @Test(timeOut = 120000)
  public void shutdownCompletesWhenEveryStateWriteThrows() throws Exception
  {
    backend.stateWrite = StateWrite.THROWS_RUNTIME_EXCEPTION;
    domain.getServerState().update(newCSN());
    waitForFailedStateWrites(1);
    // A thread which lets an exception out is dead too, so isAlive() alone can not tell the
    // last write being held from the thread dying of it. The handler of the thread outranks
    // the one of its group, which would log the exception and raise an alert instead.
    final AtomicReference<Throwable> uncaught = new AtomicReference<>();
    checkpointer.setUncaughtExceptionHandler((t, e) -> uncaught.set(e));
    deleteDomain();
    assertFalse(checkpointer.isAlive(), "the state checkpointer is still running after the shutdown");
    assertNull(uncaught.get(), "the last state write let an exception out of the checkpointer");
    assertTrue(backend.failedStateWritesAfterShutdown.get() > 0,
        "the checkpointer did not write the state on its way out");
  }
  /**
   * Shutting the domain down must not wait forever for a checkpointer which is gone: an
   * {@link Error} kills the thread whatever it catches.
   */
  @Test(timeOut = 120000)
  public void shutdownCompletesWhenTheCheckpointerDiedOfAnError() throws Exception
  {
    backend.stateWrite = StateWrite.THROWS_ERROR;
    domain.getServerState().update(newCSN());
    waitForCheckpointerToDie();
    deleteDomain();
  }
  /**
   * Shutting the domain down must not wait forever for a checkpointer whose write does not
   * come back.
   */
  @Test(timeOut = 120000)
  public void shutdownCompletesWhileAStateWriteIsStuck() throws Exception
  {
    backend.stateWrite = StateWrite.BLOCKS;
    domain.getServerState().update(newCSN());
    assertTrue(backend.blockedWriteStarted.await(CHECKPOINT_TIMEOUT_IN_SECS, SECONDS),
        "the checkpointer did not start the state write the test blocks on");
    deleteDomain();
  }
  private CSN newCSN()
  {
    return new CSN(System.currentTimeMillis(), 1, DS_ID);
  }
  private void deleteDomain()
  {
    domainDeleted = true;
    MultimasterReplication.deleteDomain(baseDN);
  }
  private void waitForFailedStateWrites(final int count) throws Exception
  {
    newTimer().repeatUntilSuccess(new CallableVoid()
    {
      @Override
      public void call() throws Exception
      {
        assertTrue(backend.failedStateWrites.get() >= count,
            "the state write of the checkpointer did not fail " + count + " time(s)");
      }
    });
  }
  private void waitForCheckpointerToDie() throws Exception
  {
    newTimer().repeatUntilSuccess(new CallableVoid()
    {
      @Override
      public void call() throws Exception
      {
        assertFalse(checkpointer.isAlive(), "the state checkpointer did not die of the injected error");
      }
    });
  }
  private TestTimer newTimer()
  {
    return new TestTimer.Builder()
        .maxSleep(CHECKPOINT_TIMEOUT_IN_SECS, SECONDS)
        .sleepTimes(100, MILLISECONDS)
        .toTimer();
  }
  private static Thread checkpointerOf(LDAPReplicationDomain domain)
  {
    final String name =
        "Replica DS(" + domain.getServerId() + ") state checkpointer for domain \"" + domain.getBaseDN() + "\"";
    final ThreadGroup group = DirectoryThread.DIRECTORY_THREAD_GROUP;
    final Thread[] threads = new Thread[group.activeCount() * 2 + 10];
    final int count = group.enumerate(threads, true);
    for (int i = 0; i < count; i++)
    {
      if (name.equals(threads[i].getName()))
      {
        return threads[i];
      }
    }
    return null;
  }
}