Class LDAPReplicationDomain
- java.lang.Object
-
- org.opends.server.replication.service.ReplicationDomain
-
- org.opends.server.replication.plugin.LDAPReplicationDomain
-
- All Implemented Interfaces:
ConfigurationChangeListener<ReplicationDomainCfg>,AlertGenerator,LocalBackendInitializationListener,ServerShutdownListener
public final class LDAPReplicationDomain extends ReplicationDomain implements ConfigurationChangeListener<ReplicationDomainCfg>, AlertGenerator, LocalBackendInitializationListener, ServerShutdownListener
This class implements the bulk part of the Directory Server side of the replication code. It contains the root method for publishing a change, processing a change received from the replicationServer service, handle conflict resolution, handle protocol messages from the replicationServer.FIXME Move this class to org.opends.server.replication.service or the equivalent package once this code is moved to a maven module.
-
-
Nested Class Summary
-
Nested classes/interfaces inherited from class org.opends.server.replication.service.ReplicationDomain
ReplicationDomain.ImportExportContext
-
-
Field Summary
Fields Modifier and Type Field Description static StringDS_SYNC_CONFLICTThe attribute used to mark conflicting entries.static intIN_PLACE_REPLAY_ATTEMPTSHow many times the replay of a change is attempted straight away, before the session to the replication server is restarted and the change is asked for again: a lock or a storage which is busy for a moment (OPENDJ-885) is waited out here.-
Fields inherited from class org.opends.server.replication.service.ReplicationDomain
broker, config, generationId, serviceStateLock
-
-
Method Summary
All Methods Static Methods Instance Methods Concrete Methods Modifier and Type Method Description voidaddAdditionalMonitoring(MonitorData attributes)Subclasses should use this method to add additional monitoring information in the ReplicationDomain.ConfigChangeResultapplyConfigurationChange(ReplicationDomainCfg configuration)Applies the configuration changes to this change listener.longcountEntries()This method should return the total number of objects in the replicated domain.intdecodeSource(String sourceString)Verifies that the given string represents a valid source from which this server can be initialized.voiddisable()Disable the replication on this domain.voidenable()Enable back the domain after a previous disable.protected voidexportBackend(OutputStream output)This method trigger an export of the replicated data.voidfailNextSessionRestarts(int failures)Has the next session restarts of this domain fail where they would start the session again.Map<String,String>getAlerts()Retrieves information about the set of alerts that this generator may produce.StringgetClassName()Retrieves the fully-qualified name of the Java class for this alert generator implementation.DNgetComponentEntryDN()Retrieves the DN of the configuration entry with which this alert generator is associated.intgetConsecutiveSessionRestarts()Returns how many times in a row the session was restarted while this replica kept failing - seeconsecutiveSessionRestartsfor what does and does not reset it: the count the backoff of the next restart is computed from.longgetReplayDrainTimeout()Returns how long this domain waits for the replay threads which are applying one of its changes before it saves its ServerState and goes down.longgetReplayRetryWarningInterval()Returns how long this domain does not warn again about a change it asks for again.intgetSessionRestartFailuresLeft()Returns how many of the session restart failuresfailNextSessionRestarts(int)asked for have not been spent yet.StringgetShutdownListenerName()Retrieves the human-readable name for this shutdown listener.protected voidimportBackend(InputStream input)This method triggers an import of the replicated data.protected voidinitializeRemote(int target, int requestorID, Task initTask, int initWindow)This is overwritten to allow stopping the (online) export process if the local domain is fractional and the destination is all other servers: This make no sense to have only fractional servers in a replicated topology.booleanisConfigurationChangeAcceptable(ReplicationDomainCfg configuration, List<org.forgerock.i18n.LocalizableMessage> unacceptableReasons)Indicates whether the proposed change to the configuration is acceptable to this change listener.protected voidonSessionRestartSuppressed()Called when what a change carries is negotiated as the session comes up, and the session was not restarted for it.voidperformBackendPostFinalizationProcessing(LocalBackend<?> backend)Performs any processing that may be required whenever a backend is finalized.voidperformBackendPostInitializationProcessing(LocalBackend<?> backend)Performs any processing that may be required after the Initialisation cycle has been completed, that is all listeners have received the initialisation event, and the backend has been put into service,.voidperformBackendPreFinalizationProcessing(LocalBackend<?> backend)Performs any processing that may be required before starting the finalisation cycle, that is invoked before any listener receive the Finalization event.voidperformBackendPreInitializationProcessing(LocalBackend<?> backend)Performs any processing that may be required whenever a backend is initialized for use in the Directory Server.voidprocessServerShutdown(org.forgerock.i18n.LocalizableMessage reason)Indicates that the Directory Server has received a request to stop running and that this shutdown listener should take any action necessary to prepare for it.booleanprocessUpdate(UpdateMsg updateMsg)This method should handle the processing ofUpdateMsgreceive from remote replication entities.voidpublishReplicaOfflineMsg()Publishes a replica offline message if all pending changes for current replica have been sent out.voidpurgeConflictsHistorical(PurgeConflictsHistoricalTask task, long endDate)Check and purge the historical attribute on all eligible entries under this domain.protected byte[]receiveEntryBytes()This is overwritten to allow stopping the (online) import process by the fractional ldif import plugin when it detects that the (imported) remote data set is not consistent with the local fractional configuration.voidrequestSessionRestart()Asks for a session restart the way a replay thread on its way out does, without running it.voidresetReplayRetryWarningThrottle()Lets the next change whose replay fails be warned about straight away.voidresetUnreplayedChangeAlertThrottle()Lets the next change this replica gives up on raise its alert straight away.protected voidrestartService()Stops the session of this domain and starts it again, so that it comes up on the configuration which has just changed.static LDAPReplicationDomainretrievesReplicationDomain(DN baseDN)Retrieves a replication domain based on the baseDN.voidsessionInitiated(ServerStatus initStatus, ServerState rsState)Set the initial status of the domain and perform necessary initializations.voidsetReplayDrainTimeout(long timeoutInMs)Sets how long this domain waits for the replay threads which are applying one of its changes before it saves its ServerState and goes down.voidsetReplayRetryWarningInterval(long intervalInMs)Sets how long this domain does not warn again about a change it asks for again.voidshutdown()Shutdown this ReplicationDomain.voidstart()Starts the Replication Domain.-
Methods inherited from class org.opends.server.replication.service.ReplicationDomain
abortStalledInitializeFromRemote, changeConfig, changeConfig, decodeTarget, disableService, disableServiceUnlessImportInProgress, enableService, getAssuredMode, getAssuredSdAcknowledgedUpdates, getAssuredSdLevel, getAssuredSdSentUpdates, getAssuredSdServerTimeoutUpdates, getAssuredSdTimeoutUpdates, getAssuredSrAcknowledgedUpdates, getAssuredSrNotAcknowledgedUpdates, getAssuredSrReceivedUpdates, getAssuredSrReceivedUpdatesAcked, getAssuredSrReceivedUpdatesNotAcked, getAssuredSrReplayErrorUpdates, getAssuredSrSentUpdates, getAssuredSrServerNotAcknowledgedUpdates, getAssuredSrTimeoutUpdates, getAssuredSrWrongStatusUpdates, getAssuredTimeout, getBaseDN, getEclIncludes, getEclIncludesForDeletes, getGenerationID, getGenerator, getGroupId, getImportExportContext, getLastLocalChange, getLastStatusChangeDate, getRefUrls, getReplicaInfos, getReplicaStates, getReplicationServer, getRsInfos, getRsServerId, getServerId, getServerState, getSessionGeneration, getStatus, hasConnectionError, ieRunning, importInProgress, incProcessedUpdates, initializeFromRemote, initializeRemote, isAssured, isConnected, isListenerShuttingDown, prepareWaitForAckIfAssuredEnabled, processUpdateDone, publish, readAssuredConfig, resetGenerationId, setEclIncludes, setGenerationID, setImportClaimHook, setServiceStopHook, signalNewStatus, startListenService, startPublishService, toString, waitForAckIfAssuredEnabled
-
-
-
-
Field Detail
-
DS_SYNC_CONFLICT
public static final String DS_SYNC_CONFLICT
The attribute used to mark conflicting entries. The value of this attribute should be the dn that this entry was supposed to have when it was marked as conflicting.- See Also:
- Constant Field Values
-
IN_PLACE_REPLAY_ATTEMPTS
public static final int IN_PLACE_REPLAY_ATTEMPTS
How many times the replay of a change is attempted straight away, before the session to the replication server is restarted and the change is asked for again: a lock or a storage which is busy for a moment (OPENDJ-885) is waited out here.- See Also:
- Constant Field Values
-
-
Method Detail
-
getShutdownListenerName
public String getShutdownListenerName()
Description copied from interface:ServerShutdownListenerRetrieves the human-readable name for this shutdown listener.- Specified by:
getShutdownListenerNamein interfaceServerShutdownListener- Returns:
- The human-readable name for this shutdown listener.
-
processServerShutdown
public void processServerShutdown(org.forgerock.i18n.LocalizableMessage reason)
Description copied from interface:ServerShutdownListenerIndicates that the Directory Server has received a request to stop running and that this shutdown listener should take any action necessary to prepare for it.- Specified by:
processServerShutdownin interfaceServerShutdownListener- Parameters:
reason- The human-readable reason for the shutdown.
-
performBackendPreInitializationProcessing
public void performBackendPreInitializationProcessing(LocalBackend<?> backend)
Description copied from interface:LocalBackendInitializationListenerPerforms any processing that may be required whenever a backend is initialized for use in the Directory Server. This method will be invoked after the backend has been initialized but before it has been put into service.- Specified by:
performBackendPreInitializationProcessingin interfaceLocalBackendInitializationListener- Parameters:
backend- The backend that has been initialized and is about to be put into service.
-
performBackendPostFinalizationProcessing
public void performBackendPostFinalizationProcessing(LocalBackend<?> backend)
Description copied from interface:LocalBackendInitializationListenerPerforms any processing that may be required whenever a backend is finalized. This method will be invoked after the backend has been taken out of service but before it has been finalized.- Specified by:
performBackendPostFinalizationProcessingin interfaceLocalBackendInitializationListener- Parameters:
backend- The backend that has been taken out of service and is about to be finalized.
-
performBackendPostInitializationProcessing
public void performBackendPostInitializationProcessing(LocalBackend<?> backend)
Description copied from interface:LocalBackendInitializationListenerPerforms any processing that may be required after the Initialisation cycle has been completed, that is all listeners have received the initialisation event, and the backend has been put into service,.- Specified by:
performBackendPostInitializationProcessingin interfaceLocalBackendInitializationListener- Parameters:
backend- The backend that has been initialized and has been put into service.
-
performBackendPreFinalizationProcessing
public void performBackendPreFinalizationProcessing(LocalBackend<?> backend)
Description copied from interface:LocalBackendInitializationListenerPerforms any processing that may be required before starting the finalisation cycle, that is invoked before any listener receive the Finalization event.- Specified by:
performBackendPreFinalizationProcessingin interfaceLocalBackendInitializationListener- Parameters:
backend- The backend that is about to be finalized.
-
receiveEntryBytes
protected byte[] receiveEntryBytes()
This is overwritten to allow stopping the (online) import process by the fractional ldif import plugin when it detects that the (imported) remote data set is not consistent with the local fractional configuration. Receives bytes related to an entry in the context of an import to initialize the domain (called by ReplLDIFInputStream).- Overrides:
receiveEntryBytesin classReplicationDomain- Returns:
- The bytes. Null when the Done or Err message has been received
-
initializeRemote
protected void initializeRemote(int target, int requestorID, Task initTask, int initWindow) throws DirectoryExceptionThis is overwritten to allow stopping the (online) export process if the local domain is fractional and the destination is all other servers: This make no sense to have only fractional servers in a replicated topology. This prevents from administrator manipulation error that would lead to whole topology data corruption. Process the initialization of some other server or servers in the topology specified by the target argument when this initialization specifying the server that requests the initialization.- Overrides:
initializeRemotein classReplicationDomain- Parameters:
target- The target server that should be initialized.requestorID- The server that initiated the export. It can be the serverID of this server, or the serverID of a remote server.initTask- The task in this server that triggers this initialization and that should be updated with its progress. Null when the export is done following a request coming from a remote server (task is remote).initWindow- The value of the initialization window for flow control between the importer and the exporter.- Throws:
DirectoryException- When an error occurs. No exception raised means success.
-
publishReplicaOfflineMsg
public void publishReplicaOfflineMsg()
Description copied from class:ReplicationDomainPublishes a replica offline message if all pending changes for current replica have been sent out.- Overrides:
publishReplicaOfflineMsgin classReplicationDomain
-
shutdown
public void shutdown()
Shutdown this ReplicationDomain.
-
resetUnreplayedChangeAlertThrottle
public void resetUnreplayedChangeAlertThrottle()
Lets the next change this replica gives up on raise its alert straight away.Only there for the tests which check the alert: they must not be at the mercy of the alert another test raised less than
UNREPLAYED_CHANGE_ALERT_INTERVAL_IN_MSago.
-
resetReplayRetryWarningThrottle
public void resetReplayRetryWarningThrottle()
Lets the next change whose replay fails be warned about straight away.Only there for the tests which check the warning: they must not be at the mercy of the warning another test logged less than
REPLAY_RETRY_WARNING_INTERVAL_IN_MSago. The deliveries folded into no warning are left alone: when they are forgotten is this domain's to decide, and the tests check that it does.
-
getReplayRetryWarningInterval
public long getReplayRetryWarningInterval()
Returns how long this domain does not warn again about a change it asks for again.- Returns:
- the interval in milliseconds
-
setReplayRetryWarningInterval
public void setReplayRetryWarningInterval(long intervalInMs)
Sets how long this domain does not warn again about a change it asks for again.Only there for the tests which check the warning: they can not wait out
REPLAY_RETRY_WARNING_INTERVAL_IN_MSbetween two of them.- Parameters:
intervalInMs- the interval in milliseconds
-
failNextSessionRestarts
public void failNextSessionRestarts(int failures)
Has the next session restarts of this domain fail where they would start the session again.Only there for the tests: nothing asks a restart to fail in production, and what this stands in for - anything thrown out of
ReplicationDomain.enableService(), which stops at no failure the callers ofReplicationBroker.start()report - can not be provoked from the outside.- Parameters:
failures- how many session restarts must fail before one is allowed to run
-
getSessionRestartFailuresLeft
public int getSessionRestartFailuresLeft()
Returns how many of the session restart failuresfailNextSessionRestarts(int)asked for have not been spent yet.Only there for the tests, which read it to tell a restart which failed from one which was never run.
- Returns:
- how many session restarts are still to fail
-
requestSessionRestart
public void requestSessionRestart()
Asks for a session restart the way a replay thread on its way out does, without running it.Only there for the tests, which have no other way to leave a request standing at a time of their choosing without a change released for it: the one a replay thread makes on a domain whose session has no owner is made and run in one go, and one made under an owner stands for a change which is listed, and which the owner forgets or the checkpointer has delivered again.
-
getConsecutiveSessionRestarts
public int getConsecutiveSessionRestarts()
Returns how many times in a row the session was restarted while this replica kept failing - seeconsecutiveSessionRestartsfor what does and does not reset it: the count the backoff of the next restart is computed from.Only there for the tests, which read it to see a restart reach its backoff rather than wait a delay out and hope it began: the count is bumped on the way into the wait, so the restart of
nrestarts in a row is at its backoff, or a few instructions short of it, once this returnsn- and a wake given in those instructions is counted, so the wait sees it all the same.- Returns:
- how many session restarts in a row the domain has run
-
disable
public void disable()
Disable the replication on this domain. The session to the replication server will be stopped. The domain will not be destroyed but call to the pre-operation methods will result in failure. The listener thread will be destroyed. The monitor informations will still be accessible.
-
getReplayDrainTimeout
public long getReplayDrainTimeout()
Returns how long this domain waits for the replay threads which are applying one of its changes before it saves its ServerState and goes down.- Returns:
- the timeout in milliseconds
-
setReplayDrainTimeout
public void setReplayDrainTimeout(long timeoutInMs)
Sets how long this domain waits for the replay threads which are applying one of its changes before it saves its ServerState and goes down.Only there for the tests which check what a domain does when that wait runs out: they can not hold a replay thread for
REPLAY_DRAIN_TIMEOUT_IN_MS.- Parameters:
timeoutInMs- the timeout in milliseconds
-
enable
public void enable()
Enable back the domain after a previous disable. The domain will connect back to a replication Server and will recreate threads to listen for messages from the Synchronization server. The generationId will be retrieved or computed if necessary. The ServerState will also be read again from the local database.
-
exportBackend
protected void exportBackend(OutputStream output) throws DirectoryException
This method trigger an export of the replicated data.- Specified by:
exportBackendin classReplicationDomain- Parameters:
output- The OutputStream where the export should be produced.- Throws:
DirectoryException- When needed.
-
importBackend
protected void importBackend(InputStream input) throws DirectoryException
This method triggers an import of the replicated data.- Specified by:
importBackendin classReplicationDomain- Parameters:
input- The InputStream from which the data are read.- Throws:
DirectoryException- When needed.
-
retrievesReplicationDomain
public static LDAPReplicationDomain retrievesReplicationDomain(DN baseDN) throws DirectoryException
Retrieves a replication domain based on the baseDN.- Parameters:
baseDN- The baseDN of the domain to retrieve- Returns:
- The domain retrieved
- Throws:
DirectoryException- When an error occurred or no domain match the provided baseDN.
-
applyConfigurationChange
public ConfigChangeResult applyConfigurationChange(ReplicationDomainCfg configuration)
Description copied from interface:ConfigurationChangeListenerApplies the configuration changes to this change listener.- Specified by:
applyConfigurationChangein interfaceConfigurationChangeListener<ReplicationDomainCfg>- Parameters:
configuration- The new configuration containing the changes.- Returns:
- Returns information about the result of changing the configuration.
-
restartService
protected void restartService()
Description copied from class:ReplicationDomainStops the session of this domain and starts it again, so that it comes up on the configuration which has just changed.The pair is taken under
ReplicationDomain.serviceStateLock, so that nothing starts a session back between the stop and the start, and both halves are counted by the session generation. A subclass may leave it alone: a domain which is shutting down, or which was disabled for a total update, owns its session and is not given one back by a configuration change, and a total update into this replica reads its entries over the session and starts the next one itself. One which does reports it throughReplicationDomain.onSessionRestartSuppressed().- Overrides:
restartServicein classReplicationDomain
-
onSessionRestartSuppressed
protected void onSessionRestartSuppressed()
Called when what a change carries is negotiated as the session comes up, and the session was not restarted for it.The configuration is stored either way, and the session started next reads it - so this says that the change is not live yet rather than that it was lost. A domain which restarts its session for every change never reaches this; one whose session has an owner - itself while it is shutting down or disabled for a total update, or a total update into this replica reading it - overrides it to tell the administrator what is waiting for that session.
Every restart this domain leaves alone comes through here, whichever step of the change asked for it: the broker properties which are renegotiated, the assured configuration the replication server is told about as the session comes up, the fractional configuration the session filters on, and the attributes the external changelog publishes.
- Overrides:
onSessionRestartSuppressedin classReplicationDomain
-
isConfigurationChangeAcceptable
public boolean isConfigurationChangeAcceptable(ReplicationDomainCfg configuration, List<org.forgerock.i18n.LocalizableMessage> unacceptableReasons)
Description copied from interface:ConfigurationChangeListenerIndicates whether the proposed change to the configuration is acceptable to this change listener.- Specified by:
isConfigurationChangeAcceptablein interfaceConfigurationChangeListener<ReplicationDomainCfg>- Parameters:
configuration- The new configuration containing the changes.unacceptableReasons- A list that can be used to hold messages about why the provided configuration is not acceptable.- Returns:
- Returns
trueif the proposed change is acceptable, orfalseif it is not.
-
getAlerts
public Map<String,String> getAlerts()
Description copied from interface:AlertGeneratorRetrieves information about the set of alerts that this generator may produce. The map returned should be between the notification type for a particular notification and the human-readable description for that notification. This alert generator must not generate any alerts with types that are not contained in this list.- Specified by:
getAlertsin interfaceAlertGenerator- Returns:
- Information about the set of alerts that this generator may produce.
-
getClassName
public String getClassName()
Description copied from interface:AlertGeneratorRetrieves the fully-qualified name of the Java class for this alert generator implementation.- Specified by:
getClassNamein interfaceAlertGenerator- Returns:
- The fully-qualified name of the Java class for this alert generator implementation.
-
getComponentEntryDN
public DN getComponentEntryDN()
Description copied from interface:AlertGeneratorRetrieves the DN of the configuration entry with which this alert generator is associated.- Specified by:
getComponentEntryDNin interfaceAlertGenerator- Returns:
- The DN of the configuration entry with which this alert generator is associated.
-
start
public void start()
Starts the Replication Domain.
-
sessionInitiated
public void sessionInitiated(ServerStatus initStatus, ServerState rsState)
Description copied from class:ReplicationDomainSet the initial status of the domain and perform necessary initializations. This method will be called by the Broker each time the ReplicationBroker establish a new session to a Replication Server. Implementations may override this method when they need to perform additional computing after session establishment. The default implementation should be sufficient for ReplicationDomains that don't need to perform additional computing.- Overrides:
sessionInitiatedin classReplicationDomain- Parameters:
initStatus- The status to enter the state machine with.rsState- The ServerState of the ReplicationServer with which the session was established.
-
countEntries
public long countEntries() throws DirectoryExceptionThis method should return the total number of objects in the replicated domain. This count will be used for reporting.- Specified by:
countEntriesin classReplicationDomain- Returns:
- The number of objects in the replication domain.
- Throws:
DirectoryException- when needed.
-
processUpdate
public boolean processUpdate(UpdateMsg updateMsg)
Description copied from class:ReplicationDomainThis method should handle the processing ofUpdateMsgreceive from remote replication entities.This method will be called by a single thread and should therefore should not be blocking.
- Specified by:
processUpdatein classReplicationDomain- Parameters:
updateMsg- TheUpdateMsgthat was received.- Returns:
- A boolean indicating if the processing is completed at return time.
If
trueis returned, no further processing is necessary. Iffalseis returned, the subclass should call the methodReplicationDomain.processUpdateDone(UpdateMsg, String)and update the ServerState When this processing is complete.
-
addAdditionalMonitoring
public void addAdditionalMonitoring(MonitorData attributes)
Description copied from class:ReplicationDomainSubclasses should use this method to add additional monitoring information in the ReplicationDomain.- Overrides:
addAdditionalMonitoringin classReplicationDomain- Parameters:
attributes- where to additional monitoring attributes
-
decodeSource
public int decodeSource(String sourceString) throws DirectoryException
Verifies that the given string represents a valid source from which this server can be initialized.- Parameters:
sourceString- The string representing the source- Returns:
- The source as a integer value
- Throws:
DirectoryException- if the string is not valid
-
purgeConflictsHistorical
public void purgeConflictsHistorical(PurgeConflictsHistoricalTask task, long endDate) throws DirectoryException
Check and purge the historical attribute on all eligible entries under this domain. The purging logic is the same applied to individual entries during modify operations. This task may be useful in scenarios where a large number of changes are made as a one-off occurrence. Running a purge-historical after the 'ds-cfg-conflicts-historical-purge-delay' period has elapsed would clear out obsolete historical data from all the modified entries reducing the overall database size.- Parameters:
task- the task raising this purge.endDate- the date to stop this task whether the job is done or not.- Throws:
DirectoryException- when an exception happens.
-
-