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

Valery Kharseko
yesterday cebef54070540c6af46671e7dd467ce2a4c61a1c
opendj-server-legacy/src/main/java/org/opends/server/replication/service/ReplicationDomain.java
@@ -46,6 +46,7 @@
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import net.jcip.annotations.GuardedBy;
import net.jcip.annotations.Immutable;
import org.forgerock.i18n.LocalizableMessage;
@@ -368,6 +369,44 @@
   * session of this ReplicationDomain.
   */
  private final Object sessionLock = new Object();
  /**
   * Serialises the stopping and the starting of the session of this domain, so that a
   * pair of them is atomic: a session stopped so that the configuration it reads can be
   * changed must not be brought back in the middle of that change by something else.
   * <p>
   * Holding it costs something, and knowingly: {@link #enableService()} connects to the
   * replication servers under this lock, so a shutdown, an import or a configuration
   * change which arrives while a replay thread is bringing the session back waits for
   * that connect - up to the configured connection timeout when the replication servers
   * are unreachable, which is the same outage that failed the replay. Every one of those
   * stops the session as its first act, so what they wait for is a session which is
   * about to be stopped again. A wait between a stop and a start belongs outside the
   * lock, so that the waiting is bounded by a connect rather than by a backoff.
   * <p>
   * It comes after the configuration backend's update lock and never before it: a write
   * to the configuration entry of a domain holds that lock while it calls the domain's
   * configuration change listener, which takes this one. So nothing may write a
   * configuration entry while holding this lock - that is why neither the state a domain
   * saves on its way down nor the generationId it stores on its way up falls back to the
   * domain configuration entry when the base entry of the suffix is missing.
   */
  protected final Object serviceStateLock = new Object();
  /**
   * Bumped every time {@link #disableService()} stops the session of this domain or
   * {@link #enableService()} starts it, both of them under {@link #serviceStateLock}. It
   * is the identity of the session: a thread which stops one and lets the lock go - a
   * replay thread waiting out a backoff before it asks for the change it could not apply
   * again - only starts it back if this still is the session it stopped, since a
   * configuration change or the end of an import may have started another one while it
   * was waiting.
   * <p>
   * The starts at domain startup are not counted: {@link #startPublishService()} from
   * the constructor of the domain and {@link #startListenService()} from its start bring
   * the two halves of the first session up before any replay thread exists to hold a
   * claim on it.
   */
  @GuardedBy("serviceStateLock")
  private long sessionGeneration;
  /**
   * The generationId for this replication domain. It is made of a hash of the
@@ -3294,34 +3333,38 @@
   * It can be useful to disable the Replication Service when the
   * repository where the replicated information is stored becomes
   * temporarily unavailable and replicated updates can therefore not
   * be replayed during a while. This method is not MT safe.
   * be replayed during a while.
   */
  public void disableService()
  public final void disableService()
  {
    synchronized (sessionLock)
    synchronized (serviceStateLock)
    {
      /*
       * Stop the broker first in order to prevent the listener from reconnecting - see OPENDJ-457.
       */
      if (broker != null)
      synchronized (sessionLock)
      {
        broker.stop();
      }
        /*
         * Stop the broker first in order to prevent the listener from reconnecting - see OPENDJ-457.
         */
        if (broker != null)
        {
          broker.stop();
        }
      // Stop the listener thread
      if (listenerThread != null)
      {
        listenerThread.initiateShutdown();
        try
        // Stop the listener thread
        if (listenerThread != null)
        {
          listenerThread.join();
          listenerThread.initiateShutdown();
          try
          {
            listenerThread.join();
          }
          catch (InterruptedException e)
          {
            // Give up waiting.
          }
          listenerThread = null;
        }
        catch (InterruptedException e)
        {
          // Give up waiting.
        }
        listenerThread = null;
      }
      sessionGeneration++;
    }
  }
@@ -3339,6 +3382,24 @@
  }
  /**
   * Returns the generation of the session of this domain: bumped by
   * {@link #disableService()} and {@link #enableService()} under
   * {@link #serviceStateLock}, not by the starts at domain startup (see
   * {@code sessionGeneration}).
   * <p>
   * It only says anything while {@link #serviceStateLock} is held, and is meant to be
   * read under the lock which stopped a session and read again under the lock which
   * starts it back: the session stopped is gone when the two differ.
   *
   * @return the generation of the session of this domain
   */
  @GuardedBy("serviceStateLock")
  protected final long getSessionGeneration()
  {
    return sessionGeneration;
  }
  /**
   * Restart the Replication service after a {@link #disableService()}.
   * <p>
   * The Replication Service will restart from the point indicated by the
@@ -3348,34 +3409,56 @@
   * If some data have changed in the repository during the period of time when
   * the Replication Service was disabled, this {@link ServerState} should
   * therefore be updated by the Replication Domain subclass before calling this
   * method. This method is not MT safe.
   * method.
   */
  public void enableService()
  public final void enableService()
  {
    synchronized (sessionLock)
    synchronized (serviceStateLock)
    {
      broker.start();
      startListenService();
      synchronized (sessionLock)
      {
        broker.start();
        startListenService();
      }
      /*
       * Counted once the session really is up: a start which threw leaves the generation
       * where it was, so the thread which stopped this session still owns it and may try
       * to bring it back.
       */
      sessionGeneration++;
    }
  }
  /**
   * Change some ReplicationDomain parameters.
   * <p>
   * The change and the restart it may call for are taken together under
   * {@link #serviceStateLock}, as {@link #readAssuredConfig(ReplicationDomainCfg, boolean)}
   * takes its own: a session brought up between the two would negotiate the broker
   * properties which are half way through being changed.
   *
   * @param config
   *          The new configuration that this domain should now use.
   */
  protected void changeConfig(ReplicationDomainCfg config)
  {
    if (broker != null && broker.changeConfig(config))
    synchronized (serviceStateLock)
    {
      restartService();
      if (broker != null && broker.changeConfig(config))
      {
        restartService();
      }
    }
  }
  /**
   * Applies a configuration change to the attributes which should be included
   * in the ECL.
   * <p>
   * Taken under {@link #serviceStateLock} like every other configuration change: this one
   * comes from the external changelog domain - from the entry of its own, or from the
   * domain configuration change which reads that entry - and it restarts the session as
   * well, so the attributes and the restart go together.
   *
   * @param includeAttributes
   *          attributes to be included with all change records.
@@ -3385,11 +3468,14 @@
  public void changeConfig(Set<String> includeAttributes,
      Set<String> includeAttributesForDeletes)
  {
    final boolean attrsModified = setEclIncludes(
        getServerId(), includeAttributes, includeAttributesForDeletes);
    if (attrsModified && broker != null)
    synchronized (serviceStateLock)
    {
      restartService();
      final boolean attrsModified = setEclIncludes(
          getServerId(), includeAttributes, includeAttributesForDeletes);
      if (attrsModified && broker != null)
      {
        restartService();
      }
    }
  }
@@ -3397,14 +3483,20 @@
   * Stops the session of this domain and starts it again, so that it comes up on the
   * configuration which has just changed.
   * <p>
   * A subclass may leave it alone: a domain which is shutting down, or which was disabled
   * for a total update, owns its session and is not given one back by a configuration
   * change. One which does reports it through {@link #onSessionRestartSuppressed()}.
   * The pair is taken under {@link #serviceStateLock}, so that nothing starts a session
   * back between the stop and the start, and both halves are counted by the session
   * generation. A subclass may leave it alone: a domain which is shutting
   * down, or which was disabled for a total update, owns its session and is not given one
   * back by a configuration change. One which does reports it through
   * {@link #onSessionRestartSuppressed()}.
   */
  protected void restartService()
  {
    disableService();
    enableService();
    synchronized (serviceStateLock)
    {
      disableService();
      enableService();
    }
  }
  /**
@@ -3884,32 +3976,41 @@
  protected void readAssuredConfig(ReplicationDomainCfg config,
      boolean allowReconnection)
  {
    // Disconnect if required: changing configuration values before
    // disconnection would make assured replication used immediately and
    // disconnection could cause some timeouts error.
    final boolean needReconnection = needReconnection(config);
    final boolean needRestart = needReconnection && allowReconnection;
    if (needRestart)
    {
      disableService();
    }
    else if (needReconnection)
    {
      onSessionRestartSuppressed();
    }
    /*
     * Stored whether or not the session was restarted for it, as the fractional
     * configuration is: the assured timeout is the one property a session does not have to
     * be restarted for, so a change carrying it alone - reported as applied and then
     * dropped, before - is applied here. A caller which does not allow the reconnection
     * has no session running assured replication either: the domain is being built, is
     * shutting down, or is disabled for the length of a total update, and the session its
     * enable() starts reads what is stored here.
     * The stop, the change and the start are taken together under serviceStateLock: a
     * session brought up in between, by a replay thread restarting the session after a
     * failed replay, would negotiate an assured configuration which is half way through
     * being changed.
     */
    assuredConfig = config;
    if (needRestart)
    synchronized (serviceStateLock)
    {
      enableService();
      // Disconnect if required: changing configuration values before
      // disconnection would make assured replication used immediately and
      // disconnection could cause some timeouts error.
      final boolean needReconnection = needReconnection(config);
      final boolean needRestart = needReconnection && allowReconnection;
      if (needRestart)
      {
        disableService();
      }
      else if (needReconnection)
      {
        onSessionRestartSuppressed();
      }
      /*
       * Stored whether or not the session was restarted for it, as the fractional
       * configuration is: the assured timeout is the one property a session does not have
       * to be restarted for, so a change carrying it alone - reported as applied and then
       * dropped, before - is applied here. A caller which does not allow the reconnection
       * has no session running assured replication either: the domain is being built, is
       * shutting down, or is disabled for the length of a total update, and the session
       * its enable() starts reads what is stored here.
       */
      assuredConfig = config;
      if (needRestart)
      {
        enableService();
      }
    }
  }