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

Valery Kharseko
16 hours ago eef07575153b2a7f37feaaa8e39977f32631082b
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)
    {