From 6477a7e28301b6d5d9746cf4bfb29ae4b748258d Mon Sep 17 00:00:00 2001
From: Valery Kharseko <vharseko@3a-systems.ru>
Date: Fri, 11 Sep 2026 12:44:36 +0000
Subject: [PATCH] [#911] Report the replication connections which used to be dropped in silence (#935)

---
 opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java |  504 ++++++++++++++++++++++++++++++++++++++++++++++++++++++-
 1 files changed, 491 insertions(+), 13 deletions(-)

diff --git a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java b/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
index 29a5572..c97a647 100644
--- a/opendj-server-legacy/src/main/java/org/opends/server/replication/server/ReplicationServer.java
+++ b/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<>());
+    }
   }
 
   /**

--
Gitblit v1.10.0