/* * The contents of this file are subject to the terms of the Common Development and * Distribution License (the License). You may not use this file except in compliance with the * License. * * You can obtain a copy of the License at legal/CDDLv1.0.txt. See the License for the * specific language governing permission and limitations under the License. * * When distributing Covered Software, include this CDDL Header Notice in each file and include * the License file at legal/CDDLv1.0.txt. If applicable, add the following below the CDDL * Header, with the fields enclosed by brackets [] replaced by your own identifying * information: "Portions Copyright [year] [name of copyright owner]". * * Copyright 2006-2009 Sun Microsystems, Inc. * Portions Copyright 2011-2015 ForgeRock AS. * Portions Copyright 2026 3A Systems, LLC. */ package org.opends.server.replication.server; import java.net.SocketException; import org.forgerock.i18n.LocalizableMessage; import org.opends.server.api.DirectoryThread; import org.forgerock.i18n.slf4j.LocalizedLogger; import org.opends.server.replication.common.ServerStatus; import org.opends.server.replication.protocol.ReplicaOfflineMsg; import org.opends.server.replication.protocol.Session; import org.opends.server.replication.protocol.UpdateMsg; import org.opends.server.replication.service.DSRSShutdownSync; import static org.opends.messages.ReplicationMessages.*; import static org.opends.server.replication.common.ServerStatus.*; import static org.opends.server.util.StaticUtils.*; /** * This class defines a server writer, which is used to send changes to a * directory server. */ public class ServerWriter extends DirectoryThread { private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass(); private final Session session; private final ServerHandler handler; private final ReplicationServerDomain replicationServerDomain; private final DSRSShutdownSync dsrsShutdownSync; /** * Create a ServerWriter. Then ServerWriter then waits on the ServerHandler * for new updates and forward them to the server * * @param session * the Session that will be used to send updates. * @param handler * handler for which the ServerWriter is created. * @param replicationServerDomain * The ReplicationServerDomain of this ServerWriter. * @param dsrsShutdownSync Synchronization object for shutdown of combined DS/RS instances. */ public ServerWriter(Session session, ServerHandler handler, ReplicationServerDomain replicationServerDomain, DSRSShutdownSync dsrsShutdownSync) { // Session may be null for ECLServerWriter. super("Replication server RS(" + handler.getReplicationServerId() + ") writing to " + handler + " at " + (session != null ? session.getReadableRemoteAddress() : "unknown")); this.session = session; this.handler = handler; this.replicationServerDomain = replicationServerDomain; this.dsrsShutdownSync = dsrsShutdownSync; } /** * Run method for the ServerWriter. * Loops waiting for changes from the ReplicationServerDomain and forward them * to the other servers */ @Override public void run() { if (logger.isTraceEnabled()) { logger.trace(getName() + " starting"); } LocalizableMessage errMessage = null; try { /* * Looping here to wait for a pending ReplicaOfflineMsg would achieve nothing: this writer * only stops once its handler has been shut down, which deactivates the consumer, clears * the message queue and closes the session. The shutdown of the domain waits for the * message to be forwarded before it stops the handlers - see * ReplicationServerDomain.shutdown() and OPENDJ-1453. */ while (true) { final UpdateMsg updateMsg = this.handler.take(); if (updateMsg == null) { // this connection is closing errMessage = LocalizableMessage.raw( "Connection closure: null update returned by domain."); break; } if (isUpdateMsgFiltered(updateMsg)) { /* * The message is dropped here and will not be published to this server. When it is the * ReplicaOfflineMsg a shutdown is waiting for - its filter is wider than the one * ReplicationServerDomain.put() applied when it queued the message, so a peer RS can * be given a message this drops - the shutdown must stop waiting for a forward which * will never be reported. */ if (updateMsg instanceof ReplicaOfflineMsg && !handler.isDataServer()) { dsrsShutdownSync.replicaOfflineMsgNotForwarded( replicationServerDomain.getBaseDN(), handler.getServerId()); } } else { // Publish the update to the remote server using a protocol version it supports session.publish(updateMsg); /* * Only the forward to a peer RS ends the wait of the shutdown: what the grace period * buys is the rest of the topology learning that the replica went offline. A directory * server is never handed this message - ReplicationServerDomain.put() does not queue * it for one, and DataServerHandler.updateServerState() drops the one the changelog * cursor of a directory server which is catching up synthesizes from the offline CSN * of the replica (issue #1029) - so the guard says whose forward counts rather than * telling two deliveries apart. */ if (updateMsg instanceof ReplicaOfflineMsg && !handler.isDataServer()) { dsrsShutdownSync.replicaOfflineMsgForwarded( replicationServerDomain.getBaseDN(), updateMsg.getCSN(), handler.getServerId()); } } } } catch (SocketException e) { /* * The remote host has disconnected and this particular Tree is going to * be removed, just ignore the exception and let the thread die as well */ errMessage = handler.getBadlyDisconnectedErrorMessage(); logger.error(errMessage); } catch (Exception e) { /* * An unexpected error happened. * Log an error and close the connection. */ errMessage = ERR_WRITER_UNEXPECTED_EXCEPTION.get(handler + " " + stackTraceToSingleLineString(e)); logger.error(errMessage); } finally { session.close(); replicationServerDomain.stopServer(handler, false); if (logger.isTraceEnabled()) { logger.trace(getName() + " stopped " + errMessage); } } } private boolean isUpdateMsgFiltered(UpdateMsg updateMsg) { if (!updateMsg.isEncodableFor(handler.getProtocolVersion())) { /* * Session.publish() drops what it cannot encode for its peer and returns as if it had sent * it, so this must be caught here: the drop would otherwise be reported as the forward of * a ReplicaOfflineMsg, which the shutdown takes for this peer having been told - and, for * a message no recipient was recorded for, as the end of its wait for the peers which can * still be told - see OPENDJ-1453 and issue #1014. Dropping it here reports what happened * instead. There is nothing to send to this one: the message is not part of the protocol * version it negotiated, and no version of it will ever reach it. *
* The drop is reported rather than traced whoever the consumer is: today the only message * gated on a protocol version is the ReplicaOfflineMsg, which never reaches the writer of * a directory server at all - ReplicationServerDomain.put() does not queue it for one, and * DataServerHandler.updateServerState() drops the copy the changelog cursor of a directory * server which is catching up synthesizes from the offline CSN of the replica (issue * #1029). A message which does get here is one the consumer was to be sent and will not * be, which is what this record says. */ logger.warn(WARN_IGNORING_UPDATE_UNSUPPORTED_BY_PEER, handler.getReplicationServerId(), updateMsg.getCSN(), handler.getBaseDN(), handler.getServerId(), session.getReadableRemoteAddress(), handler.getProtocolVersion()); return true; } if (handler.isDataServer()) { /** * Ignore updates to DS in bad BAD_GENID_STATUS or FULL_UPDATE_STATUS * * The RSD lock should not be taken here as it is acceptable to have a delay * between the time the server has a wrong status and the fact we detect it: * the updates that succeed to pass during this time will have no impact on remote server. * But it is interesting to not saturate uselessly the network * if the updates are not necessary so this check to stop sending updates is interesting anyway. * Not taking the RSD lock allows to have better performances in normal mode (most of the time). */ final ServerStatus dsStatus = handler.getStatus(); if (dsStatus == BAD_GEN_ID_STATUS) { logger.warn(WARN_IGNORING_UPDATE_TO_DS_BADGENID, handler.getReplicationServerId(), updateMsg.getCSN(), handler.getBaseDN(), handler.getServerId(), session.getReadableRemoteAddress(), handler.getGenerationId(), replicationServerDomain.getGenerationId()); return true; } else if (dsStatus == FULL_UPDATE_STATUS) { logger.warn(WARN_IGNORING_UPDATE_TO_DS_FULLUP, handler.getReplicationServerId(), updateMsg.getCSN(), handler.getBaseDN(), handler.getServerId(), session.getReadableRemoteAddress()); return true; } } else { /** * Ignore updates to RS with bad gen id * (no system managed status for a RS) */ final long referenceGenerationId = replicationServerDomain.getGenerationId(); if (referenceGenerationId != handler.getGenerationId() || referenceGenerationId == -1 || handler.getGenerationId() == -1) { logger.error(WARN_IGNORING_UPDATE_TO_RS, handler.getReplicationServerId(), updateMsg.getCSN(), handler.getBaseDN(), handler.getServerId(), session.getReadableRemoteAddress(), handler.getGenerationId(), referenceGenerationId); return true; } } return false; } }