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

Valery Kharseko
yesterday 6477a7e28301b6d5d9746cf4bfb29ae4b748258d
opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
@@ -37,12 +37,16 @@
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.CopyOnWriteArraySet;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import org.forgerock.i18n.LocalizableMessage;
import org.forgerock.i18n.LocalizableMessageBuilder;
import org.forgerock.i18n.LocalizableMessageDescriptor;
import org.forgerock.i18n.slf4j.LocalizedLogger;
import org.forgerock.opendj.config.server.ConfigChangeResult;
import org.forgerock.opendj.config.server.ConfigException;
@@ -81,6 +85,7 @@
import org.opends.server.types.HostPort;
import org.opends.server.types.SearchFilter;
import org.opends.server.types.VirtualAttributeRule;
import org.opends.server.util.FailureLogThrottle;
/**
 * ReplicationServer Listener. This singleton is the main object of the
@@ -135,6 +140,31 @@
  private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass();
  /**
   * Minimum interval, in minutes, between two warnings about the same failure of the
   * listen thread. Every connection reaching the replication port goes through that
   * thread, so the rate has to be bounded. Each message reports the interval it was
   * logged under, so this one is free to differ from the one a failed handshake is
   * reported with, which happens to be the same five minutes.
   */
  private static final long FAILURE_WARN_INTERVAL_MINUTES = 5;
  /**
   * Time, in milliseconds, the listen thread waits before accepting again after
   * {@link ServerSocket#accept()} failed.
   * <p>
   * Package private for testing: a test which makes {@code accept()} fail asserts that
   * the thread waited instead of spinning, so the wait has to be a value it can name.
   */
  static final long ACCEPT_FAILURE_BACKOFF_MS = 100;
  /** {@link #ACCEPT_FAILURE_BACKOFF_MS}, which is also how close two failures have to be to count as repeating. */
  private static final long ACCEPT_FAILURE_BACKOFF_NANOS =
      TimeUnit.MILLISECONDS.toNanos(ACCEPT_FAILURE_BACKOFF_MS);
  /** Reports a peer this replication server cannot connect to once per outage. */
  private final ConnectFailureReporter connectFailures = new ConnectFailureReporter();
  /** To know whether a domain is enabled for the external changelog. */
  private final ECLEnabledDomainPredicate domainPredicate;
@@ -288,18 +318,75 @@
        socket.getInetAddress().getHostAddress(),
        socket.getLocalPort());
    /*
     * When the last failure of accept() on this socket was done being handled, confined to
     * this thread: a listen port change runs a second listen thread, and each of them times
     * the socket it accepts on.
     *
     * A failure which follows one closely is what the wait below is for. It starts a whole
     * interval in the past, so an isolated failure does not wait, and a successful accept
     * puts it back there: a loop which is serving connections is doing work rather than
     * spinning, and what paced it before is not what the next failure repeats. Without that
     * reset, a stream of connections aborted between the handshake and accept() -- a health
     * check or a port scan -- would charge the whole listen port a wait per probe, and the
     * peers queued behind them would pay it.
     *
     * What the reset gives up is the failure which alternates with a connection: a process
     * out of file descriptors frees one now and then, the accept it lets through resets the
     * clock, and the failure after it is timed as isolated and not waited on. The loop then
     * turns as fast as the connections arrive. That is the trade -- a probe must not cost
     * the peers behind it a wait, and a loop which is accepting is not the silent spin the
     * wait was added for -- and the warning is throttled either way, so the error log holds
     * one line per five minutes of it whichever side of the trade the failure falls on.
     *
     * Read when the previous failure was handled rather than when it happened, so that the
     * wait it was granted is not what makes the next failure look isolated: timing from the
     * failure would leave every second one of a continuous run unwaited, and bound the spin
     * to twice the rate the wait is chosen for.
     */
    long handledAcceptFailureNanos = System.nanoTime() - ACCEPT_FAILURE_BACKOFF_NANOS;
    /*
     * Bound how often a failure of accept() on this socket, and a connection accepted on it
     * which cannot be turned into a session, are warned about. Both are confined to this
     * thread for the same reason the clock above is: a listen port change runs a second
     * listen thread -- switchListenPort() starts it before it stops this one -- and a five
     * minute window opened on the port which was left would otherwise suppress the first
     * failure on the port which replaced it, silence right after the administrator changed
     * the port to get out of trouble.
     */
    final FailureLogThrottle acceptFailures =
        new FailureLogThrottle(FAILURE_WARN_INTERVAL_MINUTES, TimeUnit.MINUTES);
    final FailureLogThrottle sessionSetupFailures =
        new FailureLogThrottle(FAILURE_WARN_INTERVAL_MINUTES, TimeUnit.MINUTES);
    while (!shutdown.get() && !socket.isClosed())
    {
      // Wait on the replicationServer port.
      // Read incoming messages and create LDAP or ReplicationServer listener
      // and Publisher.
      Session session = null;
      try
      {
        Session session;
        Socket newSocket = null;
        final Socket newSocket;
        try
        {
          newSocket = socket.accept();
        }
        catch (Exception e)
        {
          final boolean repeated =
              System.nanoTime() - handledAcceptFailureNanos < ACCEPT_FAILURE_BACKOFF_NANOS;
          handleAcceptFailure(acceptFailures, socket, e, repeated);
          handledAcceptFailureNanos = System.nanoTime();
          continue;
        }
        // A connection served is not a spin: the failures before it stop pacing the loop,
        // and a failure alternating with a connection is therefore never waited on.
        handledAcceptFailureNanos = System.nanoTime() - ACCEPT_FAILURE_BACKOFF_NANOS;
        try
        {
          newSocket.setTcpNoDelay(true);
          newSocket.setKeepAlive(true);
          int timeoutMS = MultimasterReplication.getConnectionTimeoutMS();
@@ -311,12 +398,13 @@
        }
        catch (Exception e)
        {
          // If problems happen during the SSL handshake, it is necessary
          // to close the socket to free the associated resources.
          if (newSocket != null)
          {
            newSocket.close();
          }
          logSessionSetupFailure(sessionSetupFailures, newSocket, e);
          // createServerSession() closes the socket itself when it does not return a
          // session, so this closes the one whose options could not be set, and closes a
          // second time, harmlessly, the one it already released. Through the helper: a
          // close which throws here would escape this catch into the one of the loop,
          // which reports it as a failure to listen, without a throttle.
          close(newSocket);
          continue;
        }
@@ -349,6 +437,16 @@
      }
      catch (Exception e)
      {
        /*
         * A session which reaches here is owned by nothing else: both handlers started
         * above abort a handshake they cannot complete themselves, closing the session
         * and returning, so what lands here is the session of a peer which failed after
         * its handshake -- receive() on a peer which was killed, or which stopped
         * answering past the connection timeout. Leaving it open leaked the socket and
         * its file descriptor for the life of the process, one per such peer, towards the
         * very exhaustion the accept loop above now has to survive.
         */
        close(session);
        // The socket has probably been closed as part of the
        // shutdown or changing the port number process.
        // Just log debug information and loop.
@@ -363,6 +461,128 @@
  }
  /**
   * Reports a failure of {@link ServerSocket#accept()} and, when it is not the first one
   * in a row, waits before accepting again.
   * <p>
   * The listen socket stays open across such a failure, so the loop would come straight
   * back to {@code accept()} and, when the cause is the process running out of file
   * descriptors, spin on it without ever logging anything. The wait is what bounds that
   * spin; the throttle is what bounds the log.
   * <p>
   * Only a failure which repeats is waited on: a connection reset between the handshake
   * and {@code accept()} fails it once, and making the listen thread pause for a
   * connection nobody is waiting for any more would slow down the ones which follow it.
   *
   * @param throttle
   *          The throttle of the calling listen thread, which bounds how often this
   *          failure is warned about.
   * @param socket
   *          The socket connections are accepted on.
   * @param e
   *          The failure.
   * @param repeated
   *          Whether the previous failure of {@code accept()} on that socket was handled
   *          less than {@link #ACCEPT_FAILURE_BACKOFF_MS} ago.
   */
  private void handleAcceptFailure(final FailureLogThrottle throttle, final ServerSocket socket,
      final Exception e, final boolean repeated)
  {
    // Read before the socket is tested rather than after it is reported: a socket closed
    // in between is one this method returns on, while a closed socket names no address at
    // all, and the warning would report the port the administrator needs as "null".
    final Object listenAddress = socket.getLocalSocketAddress();
    if (shutdown.get() || socket.isClosed())
    {
      // The socket was closed to stop this thread or to change the listen port: the loop
      // is about to end, and this failure is how it is told to.
      logger.traceException(e);
      return;
    }
    logThrottledFailure(throttle, WARN_REPLICATION_SERVER_ACCEPT_ERROR, listenAddress, e);
    if (repeated)
    {
      // The wait is not interruptible on purpose: the flag is left alone, as restoring it
      // would make every wait which follows return at once and bring the spin back.
      // Nothing is lost by that -- the only interrupt the listen thread gets, from
      // abortInitialization(), comes after the shutdown flag is set and the socket
      // closed, which is what ends this loop.
      sleep(ACCEPT_FAILURE_BACKOFF_MS);
    }
  }
  /**
   * Reports a connection which was accepted but on which no replication session could be
   * started, and which is therefore about to be closed.
   * <p>
   * A failed SSL handshake is reported by {@link ReplSessionSecurity#createServerSession}
   * itself, so what reaches here is everything else: a trust store which cannot be read,
   * which makes every inbound connection fail this way, and a peer which stops
   * responding during the handshake. Every connection reaching the replication port goes
   * through this path, so the log has to be throttled the way the handshake failure next
   * door is.
   *
   * @param throttle
   *          The throttle of the listen thread which accepted the connection, which is
   *          confined to it: see where it is built.
   * @param socket
   *          The accepted socket, which {@code createServerSession} has usually closed
   *          already, in its own {@code finally}: a closed socket still names the peer it
   *          was connected to, which is what this reads from it.
   * @param e
   *          The failure.
   */
  private void logSessionSetupFailure(final FailureLogThrottle throttle, final Socket socket,
      final Exception e)
  {
    logThrottledFailure(throttle, WARN_REPLICATION_SERVER_SESSION_SETUP_ERROR,
        socket.getRemoteSocketAddress(), e);
  }
  /**
   * Logs a failure of the listen thread as a warning when the provided throttle lets it
   * through, reporting how many failures that line stands for.
   * <p>
   * A failure the throttle suppresses is still recorded, with the information severity:
   * {@code logger.debug} of a localized message publishes to the error log and not to the
   * debug log. What reads that severity is the replication log, whose shipped publisher
   * overrides it on for the {@code SYNC} category; the error log does not publish it by
   * default. So the throttle bounds what the error log holds, and what the replication log
   * holds is one record per failure, which is what the messages say.
   *
   * @param throttle
   *          The throttle bounding how often this failure is warned about.
   * @param message
   *          The message reporting it, which takes the server id, the address, the cause,
   *          the interval and the number of failures suppressed since the previous warning.
   * @param address
   *          The address the failure happened on.
   * @param e
   *          The failure.
   */
  private void logThrottledFailure(final FailureLogThrottle throttle,
      final LocalizableMessageDescriptor.Arg5<Number, Object, Object, Number, Number> message,
      final Object address, final Exception e)
  {
    // The stack trace goes to the trace log, which formats it only when it is enabled.
    logger.traceException(e);
    final long recorded = throttle.record(System.nanoTime());
    // Not a stack trace: the arguments of a message are built whether or not the record
    // they go into is published, and this one is written in a state the server has to
    // survive.
    final LocalizableMessage cause = getExceptionMessage(e);
    if (recorded >= 0)
    {
      logger.warn(message, getServerId(), address, cause, FAILURE_WARN_INTERVAL_MINUTES, recorded);
    }
    else
    {
      logger.debug(message, getServerId(), address, cause, FAILURE_WARN_INTERVAL_MINUTES, -recorded - 1);
    }
  }
  /**
   * This method manages the connection with the other replication servers.
   * It periodically checks that this replication server is indeed connected
   * to all the other replication servers and if not attempts to
@@ -376,8 +596,11 @@
      while (!shutdown.get())
      {
        HostPort localAddress = HostPort.localAddress(getReplicationPort());
        final Set<HostPort> configuredRSAddresses = getConfiguredRSAddresses();
        final Set<DN> domainBaseDNs = new HashSet<>();
        for (ReplicationServerDomain domain : getReplicationServerDomains())
        {
          domainBaseDNs.add(domain.getBaseDN());
          /*
           * If there are N RSs configured then we will usually be connected to
           * N-1 of them, since one of them is usually this RS. However, we
@@ -386,11 +609,17 @@
           */
          final Set<HostPort> connectedRSAddresses =
              getConnectedRSAddresses(domain);
          for (HostPort rsAddress : getConfiguredRSAddresses())
          for (HostPort rsAddress : configuredRSAddresses)
          {
            if (connectedRSAddresses.contains(rsAddress))
            {
              continue; // Skip: already connected.
              // Skip: already connected. The connection may be the one that peer made to
              // this server, which connect() never sees, so this is where a failure
              // reported for a peer which came back on its own is closed: leaving it
              // recorded would silence the next outage of that peer. A handler registered
              // with the domain is a handshake which completed, so this reports a session.
              reportConnectionRestored(rsAddress, domain.getBaseDN(), true);
              continue;
            }
            // FIXME: this will need changing if we ever support listening on
@@ -413,6 +642,12 @@
          }
        }
        // A peer or a domain which is no longer configured is never connected to again,
        // so nothing would ever clear a failure recorded for it. Forget those: a peer
        // taken out of the configuration while it is down, and put back while it still
        // is, would otherwise have its first failure silenced.
        connectFailures.retainAll(configuredRSAddresses, domainBaseDNs);
        // Notify any threads waiting with domain tickets after each iteration.
        synchronized (domainTicketLock)
        {
@@ -448,13 +683,21 @@
  /**
   * Establish a connection to the server with the address and port.
   * <p>
   * Package private for testing: what a peer costs the error log is decided here, and a
   * test of the bookkeeping alone cannot see how it is called.
   *
   * @param remoteServerAddress
   *          The address and port for the server
   * @param baseDN
   *          The baseDN of the connection
   * @return {@code true} if the peer is connected to, {@code false} if it could not be
   *         reached at all or if the handshake offered to it did not complete. The
   *         caller leaves a peer it gets {@code false} for alone for a few passes, which
   *         a peer answering with an abort every second is as much in need of as one
   *         which does not answer.
   */
  private boolean connect(HostPort remoteServerAddress, DN baseDN)
  boolean connect(HostPort remoteServerAddress, DN baseDN)
  {
    boolean sslEncryption = replSessionSecurity.isSslEncryption();
@@ -466,6 +709,7 @@
    Socket socket = new Socket();
    Session session = null;
    final boolean handshakeCompleted;
    try
    {
      socket.setReuseAddress(true);
@@ -490,16 +734,250 @@
      ReplicationServerHandler rsHandler = new ReplicationServerHandler(
          session, config.getQueueSize(), this, config.getWindowSize());
      rsHandler.connect(baseDN, sslEncryption);
      handshakeCompleted = rsHandler.connect(baseDN, sslEncryption);
    }
    catch (Exception e)
    {
      logger.traceException(e);
      if (connectFailures.recordFailure(remoteServerAddress, baseDN))
      {
        // The failure used to be traced and nothing else, so a replication server whose
        // outgoing handshake failed logged nothing at all: the only log naming the problem
        // was the one of the peer, which is the machine whose configuration is right.
        logger.warn(WARN_REPLICATION_SERVER_CONNECT_ERROR, getServerId(), remoteServerAddress,
            baseDN, getExceptionMessage(e));
      }
      close(session);
      close(socket);
      return false;
    }
    return true;
    /*
     * The outage closed here is a failure to connect, and reaching this line is the peer
     * answering on its replication port: everything WARN_REPLICATION_SERVER_CONNECT_ERROR
     * is reported for is above it -- the socket and the session built on it, the handshake
     * throwing nothing of its own. So the outage is closed whatever the handshake did next,
     * and what the handshake did next decides which recovery is reported rather than
     * whether one is.
     *
     * Closing it on the connection alone is what keeps the peers this server never sees
     * connected under the address it dialled reportable, and there is one of those for each
     * narrower reading:
     *
     * the registration with the domain misses a peer which negotiates protocol version 1.
     * ReplicationServerHandler.connect() registers only above V1, the FIXME there being
     * older than this, so such a peer is connected and never registered.
     *
     * the address, and the session left open with it, miss a peer which dials out from an
     * address other than the one it is configured under -- multi homing, NAT. Its inbound
     * handler is registered under the source address of its own connection,
     * ServerHandler.toServerAddressURL() reading the host from the session, so the already
     * connected branch of runConnect() compares the configured address against one it never
     * matches, and the handshake this server offers that same peer aborts on a duplicate
     * server id: abortStart() closes the session, and an open session is never seen here
     * again.
     *
     * An outage left open is not a line too few but a peer gone silent: recordFailure()
     * returns false from then on, so the next real outage of it is not reported at all.
     *
     * What the handshake did next is reported all the same, because it is not the log which
     * can be left to it. ReplicationServerHandler.connect() aborts at seven places, and the
     * three which pass no message log nothing at all, abortStart() logging nothing without
     * one: a StopMsg read where the peer's ReplServerStartMsg was due, a cross connect this
     * server resolves against a peer it is already connected to, and a phase two the peer
     * leaves unanswered. Those would otherwise leave "connected" as the last thing said
     * about a domain which has no session for it. The remaining four have their reason
     * logged here by abortStart(), and they are not all of one kind: this server rejecting
     * the peer, the peer going away while the handshake ran, and whatever else fails on
     * this side of it. So the message reported for an answer without a session tells the
     * two apart by where a reason was logged rather than by naming a side: the operator is
     * sent to the log of the peer only where nothing was logged here, and that is also the
     * abort the peer may have logged nothing about either -- it resolves a cross connect
     * with the same silent abortStart(null), one line up in startFromRemoteRS().
     *
     * The cross connect this server resolves is the one abort of the three where a session
     * for the domain does exist: it is the connection the peer made, which the already
     * connected branch of runConnect() reports on its next pass. Reaching this line with
     * one open needs that registration to land between the snapshot that branch reads and
     * the dial below it, so what it costs is one warning, and the pass after it says what
     * is true.
     */
    reportConnectionRestored(remoteServerAddress, baseDN, handshakeCompleted);
    return handshakeCompleted;
  }
  /**
   * Reports that a peer this replication server had reported it could not connect to can be
   * reached again, and does nothing when what is now true of that peer is what was last
   * reported about it.
   * <p>
   * A peer which answers is no longer the peer the outage was reported for, whether the
   * handshake completed or not: leaving the outage recorded for one which answers and aborts
   * every handshake would silence its next real outage. But the two are not the same
   * recovery, and the difference is one the operator has to be able to read -- a peer which
   * stops the handshake is one this server can reach and has no session with, which is not
   * what "connected" says -- so each moves the peer to a state of its own rather than
   * clearing what is known about it. That is what leaves the session established after an
   * abort still reportable: it is a recovery from the answer without a session, which the
   * clearing form had already consumed.
   *
   * @param remoteServerAddress
   *          The address of the peer which answered.
   * @param baseDN
   *          The base DN of the domain the attempt was for.
   * @param connected
   *          Whether the handshake completed, so that this server is connected to the peer,
   *          rather than aborted after it had answered.
   */
  private void reportConnectionRestored(final HostPort remoteServerAddress, final DN baseDN,
      final boolean connected)
  {
    if (connected)
    {
      if (connectFailures.recordConnected(remoteServerAddress, baseDN))
      {
        logger.info(NOTE_REPLICATION_SERVER_CONNECT_RESTORED, getServerId(), remoteServerAddress, baseDN);
      }
    }
    else if (connectFailures.recordReachableWithoutSession(remoteServerAddress, baseDN))
    {
      logger.warn(WARN_REPLICATION_SERVER_REACHABLE_NO_SESSION, getServerId(), remoteServerAddress, baseDN);
    }
  }
  /**
   * Remembers which replication servers this replication server has already reported it
   * cannot connect to, so that a peer which stays unreachable is reported once instead of
   * on every attempt to reach it.
   * <p>
   * The connect thread retries a failed peer every few seconds, for as long as it is down,
   * and a peer being down is a normal state: one stopped for maintenance would otherwise
   * fill the error log for the duration. The failure is reported when it starts and, through
   * {@link #recordConnected} and {@link #recordReachableWithoutSession}, when it ends, so
   * that neither end of it has to be inferred from a silence.
   * <p>
   * What is held per peer and domain is the last state which was <em>reported</em>, not
   * whether an outage is open. A record which is merely present or absent cannot carry the
   * middle state -- a peer which answers on its replication port and stops the handshake --
   * because reporting it would consume the record, and the session which is established
   * seconds later would then find nothing left to close and go unreported. A peer restarting
   * takes exactly that path: its port answers before its domains are up.
   * <p>
   * Package private for testing.
   */
  static final class ConnectFailureReporter
  {
    /** What the last message about a peer and a domain said about them. */
    private enum Reported
    {
      /** The peer could not be reached at all: {@code WARN_REPLICATION_SERVER_CONNECT_ERROR}. */
      DOWN,
      /**
       * The peer answered and the handshake did not complete:
       * {@code WARN_REPLICATION_SERVER_REACHABLE_NO_SESSION}.
       */
      REACHABLE_NO_SESSION;
    }
    /**
     * What was last reported about each domain of each peer, holding no entry for the peers
     * and domains nothing was reported about, which a connection to them is therefore not to
     * be reported for.
     * <p>
     * Keyed by {@link HostPort} rather than by its string form, which is the raw host
     * while equality is on the normalized one: two spellings of the same peer would
     * otherwise be two keys, and a failure recorded under one would never be cleared by
     * the connection recorded under the other.
     * <p>
     * A HostPort holds the resolution its host had when it was built, and is documented as
     * not meant to be cached. {@link #retainAll} is what makes holding one here bounded: it
     * runs on every pass of the connect thread against freshly built addresses, so a key
     * whose resolution has drifted is dropped within one pass, and the outage under it is
     * reported a second time rather than never reported again.
     * <p>
     * The connect thread is the only one which connects to peers, but a configuration
     * change replaces it, so this is a map two threads may hand over.
     */
    private final ConcurrentMap<HostPort, ConcurrentMap<DN, Reported>> reported =
        new ConcurrentHashMap<>();
    /**
     * Records a failure to connect to the provided peer for the provided domain.
     *
     * @param peer
     *          The address of the replication server which could not be connected to.
     * @param baseDN
     *          The base DN of the domain the connection was for.
     * @return {@code true} if this failure is to be reported, {@code false} if that peer is
     *         already known to be unreachable for that domain.
     */
    boolean recordFailure(final HostPort peer, final DN baseDN)
    {
      return Reported.DOWN != domainsOf(peer).put(baseDN, Reported.DOWN);
    }
    /**
     * Records that the provided peer answered on its replication port for the provided
     * domain and that the handshake did not complete, so that this server has no session
     * with it.
     * <p>
     * Reported only where an outage was: an abort which carries a reason logs that reason on
     * this server, and one which does not is the peer's to log. What this closes is the
     * message which said the peer could not be reached, which is no longer what is wrong
     * with it.
     *
     * @param peer
     *          The address of the replication server which answered.
     * @param baseDN
     *          The base DN of the domain the attempt was for.
     * @return {@code true} if this is to be reported, {@code false} if nothing was reported
     *         about that peer or if it was already reported as answering without a session.
     */
    boolean recordReachableWithoutSession(final HostPort peer, final DN baseDN)
    {
      final ConcurrentMap<DN, Reported> domains = reported.get(peer);
      return domains != null
          && Reported.DOWN == domains.replace(baseDN, Reported.REACHABLE_NO_SESSION);
    }
    /**
     * Records a session established with the provided peer for the provided domain.
     *
     * @param peer
     *          The address of the replication server which was connected to.
     * @param baseDN
     *          The base DN of the domain the connection is for.
     * @return {@code true} if this session ends something which was reported and is
     *         therefore to be reported as well, {@code false} if nothing was reported about
     *         that peer.
     */
    boolean recordConnected(final HostPort peer, final DN baseDN)
    {
      final ConcurrentMap<DN, Reported> domains = reported.get(peer);
      return domains != null && domains.remove(baseDN) != null;
    }
    /**
     * Forgets what was recorded for the peers and the domains which are no longer
     * configured, and which nothing can therefore report a recovery for.
     *
     * @param peers
     *          The addresses of the replication servers which are configured.
     * @param baseDNs
     *          The base DNs of the domains which exist.
     */
    void retainAll(final Set<HostPort> peers, final Set<DN> baseDNs)
    {
      reported.keySet().retainAll(peers);
      for (ConcurrentMap<DN, Reported> domains : reported.values())
      {
        domains.keySet().retainAll(baseDNs);
      }
    }
    private ConcurrentMap<DN, Reported> domainsOf(final HostPort peer)
    {
      return reported.computeIfAbsent(peer, unused -> new ConcurrentHashMap<>());
    }
  }
  /**