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