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

Valery Kharseko
2 days ago 9ff409beae28518a8846ac7423eb34f5faec6ee2
opendj-server-legacy/src/main/java/org/opends/server/replication/server/MessageHandler.java
@@ -81,6 +81,19 @@
  /** Specifies whether the consumer is following the producer (is not late). */
  private boolean following;
  /**
   * The state of the domain at the previous wait of {@link #getNextMessage()} on an empty queue,
   * or {@code null} when the queue was not empty then. A change it holds has had a whole wait to
   * reach the queue, see {@link #stopFollowingWhenChangesAreMissing()}. Guarded by
   * {@link #msgQueue}.
   */
  private ServerState domainStateAtPreviousWait;
  /**
   * The state of the domain the changelog was last read again for by
   * {@link #stopFollowingWhenChangesAreMissing()}, or {@code null} when it never was. Guarded by
   * {@link #msgQueue}.
   */
  private ServerState lastMissingChangesDomainState;
  /**
   * Specifies whether the last update message returned by
   * {@link #getNextMessage()} was re-read from the changelog DB (catch-up
   * path) rather than taken from the in-memory {@link #msgQueue}. Only ever
@@ -233,6 +246,25 @@
  }
  /**
   * The state of the domain at the previous wait of {@link #getNextMessage()} on an empty queue,
   * or {@code null} when the queue was not empty then - see {@link #domainStateAtPreviousWait}.
   * <p>
   * Package private for testing: a test which holds a change on its way to the queue hands it
   * over once a wait has seen the domain hold it, so that
   * {@link #stopFollowingWhenChangesAreMissing()} has run with the change on its way - by
   * construction rather than by wall clock.
   *
   * @return the state of the domain at the previous wait on an empty queue, or {@code null}
   */
  ServerState getDomainStateAtPreviousWait()
  {
    synchronized (msgQueue)
    {
      return domainStateAtPreviousWait;
    }
  }
  /**
   * Indicates whether the last update message returned by
   * {@code getNextMessage()} was re-read from the changelog DB (catch-up
   * path) rather than taken from the in-memory queue.
@@ -392,11 +424,17 @@
              {
                return null;
              }
              stopFollowingWhenChangesAreMissing();
            }
          } catch (InterruptedException e)
          {
            return null;
          }
          if (msgQueue.isEmpty())
          {
            // this handler is going back to the changelog, see stopFollowingWhenChangesAreMissing()
            continue;
          }
          UpdateMsg msg = msgQueue.removeFirst();
          if (updateServerState(msg))
          {
@@ -420,6 +458,75 @@
  }
  /**
   * Leaves the in-memory queue path when the domain holds changes this handler was never given.
   * <p>
   * A handler which is following is served by {@link #add(UpdateMsg)} alone, so an empty queue
   * means the consumer has everything the domain received. A change which reached the changelog
   * without reaching that queue breaks this: {@link #fillLateQueue()} is the only reader of the
   * changelog and is not called again while the handler follows, so the change would never be
   * sent and the consumer would be considered up to date forever - see issue #963. Going back to
   * the catch-up path reads the changelog again and delivers it.
   * <p>
   * A change is missing when the domain held it at the previous wait and it is still neither in
   * the queue nor in the state of this handler. The state of the domain is advanced by
   * {@code ReplicationServerDomain.publishUpdateMsg()} slightly before {@code addUpdate()} queues
   * the change, so what the domain received since the previous wait may simply be on its way;
   * comparing with the state of the domain as it was one wait ago leaves such a change alone even
   * when the domain has been ahead of this handler for longer, because of a gap the changelog was
   * already read for.
   * <p>
   * The changelog is read again once per advance of the state of the domain: a gap the re-read
   * does not close is reported once, and read for again when the domain receives something new.
   * <p>
   * Must be called while holding the {@link #msgQueue} monitor.
   */
  private void stopFollowingWhenChangesAreMissing()
  {
    if (!msgQueue.isEmpty() || !following || !isFedByTheDomain())
    {
      domainStateAtPreviousWait = null;
      return;
    }
    final ServerState missingSince = domainStateAtPreviousWait;
    final ServerState domainState = replicationServerDomain.getLatestServerState();
    domainStateAtPreviousWait = domainState;
    if (missingSince == null || serverState.cover(missingSince))
    {
      return;
    }
    if (lastMissingChangesDomainState != null && lastMissingChangesDomainState.cover(missingSince))
    {
      // the changelog has already been read again for these changes and gave nothing: reading it
      // once more would give nothing either until the domain receives something new
      return;
    }
    lastMissingChangesDomainState = domainState;
    following = false;
    logger.warn(WARN_CHANGELOG_READ_AGAIN_FOR_MISSING_CHANGES, replicationServer.getServerId(),
        baseDN, getMonitorInstanceName(), serverState, domainState);
  }
  /**
   * Whether the domain hands every update it receives to this handler.
   * <p>
   * Only for such a handler does the state of the domain being ahead of its own mean that a
   * change was missed, see {@link #stopFollowingWhenChangesAreMissing()}. The state of any other
   * handler is legitimately behind the state of the domain, and the changelog must not be read on
   * its behalf: a peer replication server is handed only the changes of the directory servers
   * connected to this one, so what a third replication server relayed never reaches its queue
   * nor its state, and reading the changelog again would send it what it already holds; a
   * directory server the domain filters out would have the writer drop every change read while
   * its state moved past it.
   *
   * @return {@code true} when {@code ReplicationServerDomain.put()} queues every update it
   *         receives for this handler
   */
  boolean isFedByTheDomain()
  {
    return false;
  }
  /**
   * Fills the late queue with the most recent changes, accepting only the
   * messages from provided replica ids.
   */