/* * 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 2007-2009 Sun Microsystems, Inc. * Portions Copyright 2013-2016 ForgeRock AS. * Portions Copyright 2026 3A Systems, LLC. */ package org.opends.server.replication.plugin; import java.util.Iterator; import java.util.NoSuchElementException; import java.util.SortedMap; import java.util.SortedSet; import java.util.TreeMap; import java.util.TreeSet; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ConcurrentSkipListSet; import java.util.concurrent.locks.ReentrantLock; import java.util.concurrent.locks.ReentrantReadWriteLock; import net.jcip.annotations.GuardedBy; import org.forgerock.opendj.ldap.DN; import org.opends.server.core.AddOperation; import org.opends.server.core.DeleteOperation; import org.opends.server.core.ModifyDNOperationBasis; import org.opends.server.core.ModifyOperation; import org.opends.server.replication.common.CSN; import org.opends.server.replication.common.ServerState; import org.opends.server.replication.protocol.AddMsg; import org.opends.server.replication.protocol.DeleteMsg; import org.opends.server.replication.protocol.LDAPUpdateMsg; import org.opends.server.replication.protocol.ModifyDNMsg; import org.opends.server.replication.protocol.ModifyMsg; import org.opends.server.replication.protocol.OperationContext; import org.opends.server.types.Operation; /** * This class is used to store the list of remote changes received * from a replication server and that are either currently being replayed * or that are waiting for being replayed. * * It is used to know when the ServerState must be updated and to compute * the dependencies between operations. * * One of this object is instantiated for each ReplicationDomain. */ final class RemotePendingChanges { /** A map used to store the pending changes. */ @GuardedBy("pendingChangesLock") private final SortedMap pendingChanges = new TreeMap<>(); /** * A sorted set containing the list of PendingChanges that have * not been replayed correctly because they are dependent on * another change to be completed. */ @GuardedBy("dependentChangesLock") private final SortedSet dependentChanges = new TreeSet<>(); /** * {@code activeAndDependentChanges} also contains changes discovered to be dependent * on currently in progress changes. */ private final ConcurrentSkipListSet activeAndDependentChanges = new ConcurrentSkipListSet<>(); /** * The change each replay thread is replaying right now, read by the give-back on the way * out of a replay which was unwound. *

* It is an index of what {@link PendingChange#isOwnedBy(Thread)} already says rather than * a second copy of it: it is written in the same locked step wherever a change is taken * over or given back, and every road which acts on what it answers checks the ownership * of the change again. What it buys is the read. That read is made on a road an * {@code OutOfMemoryError} leads to, and looking the change up by walking the pending * changes takes two locks and allocates an iterator - an allocation a JVM which has just * refused one may well refuse again, and the change would then be left listed, * uncommitted and owned by a thread which is not replaying it anymore, which is the wedge * this issue is about (issue #922). *

* A thread is entered here when it takes a change over and removed when it gives it back, * applies it, or parks it as waiting for another change - the parked ones are handed to * whichever thread clears what they wait for, so they are not this one's to give back. *

* The entry of a thread is written by that thread and by nobody else, and that - not the * lock - is what keeps the writes apart: the park in {@link #addDependency(PendingChange)} * clears it under the read lock, where every other writer holds the write lock, and it * races nothing for it. {@link #clear()} is the one writer of every entry, and it holds * both locks. */ private final ConcurrentMap changeBeingReplayed = new ConcurrentHashMap<>(); private final ReentrantReadWriteLock pendingChangesLock = new ReentrantReadWriteLock(true); private final ReentrantReadWriteLock.ReadLock pendingChangesReadLock = pendingChangesLock.readLock(); private final ReentrantReadWriteLock.WriteLock pendingChangesWriteLock = pendingChangesLock.writeLock(); private final ReentrantLock dependentChangesLock = new ReentrantLock(); /** The ServerState that will be updated when LDAPUpdateMsg are fully replayed. */ private final ServerState state; /** * How many of the changes listed here have a replay failure recorded against them. *

* A change is counted from its first failed replay until it leaves this map, and it is * what tells a change which can not be applied from a backend which is serving again: * an outage fails everything in flight, while a change which can never be applied here * fails alone, among changes which replay perfectly well. The session restart backoff * reads it, so that the successful replays which surround a failing change do not keep * resetting the wait this domain has reached on it (issue #889). */ @GuardedBy("pendingChangesLock") private int failingChanges; /** * Creates a new RemotePendingChanges using the provided ServerState. * * @param state The ServerState that will be updated when LDAPUpdateMsg * have been fully replayed. */ public RemotePendingChanges(ServerState state) { this.state = state; } /** * Returns the number of changes waiting to be replayed. * * @return The number of changes waiting to be replayed */ public int getQueueSize() { pendingChangesReadLock.lock(); try { return pendingChanges.size(); } finally { pendingChangesReadLock.unlock(); } } /** * Returns the number of changes actively being replayed. *

* A change whose replay failed counts here until it is applied or given up on: it stays * a dependency of the changes which follow it, since a change which is not in the data * is exactly what they must wait for, and the delivery which comes next takes it over. * * @return the number of changes actively being replayed. */ public int changesInProgressSize() { return activeAndDependentChanges.size(); } /** * Returns the number of changes depending on other changes. * * @return the number of changes depending on other changes. */ public int getDependentChangesSize() { dependentChangesLock.lock(); try { return dependentChanges.size(); } finally { dependentChangesLock.unlock(); } } /** * Add a new LDAPUpdateMsg that was received from the replication server * to the pendingList. *

* A change which is already listed and which a replay thread owns - it is being * replayed, it waits for the change it depends on, or it has been replayed and waits * for the changes before it - is left alone: that copy knows what this one does not, * and replaying it again is exactly what the duplicate check is there to prevent * (OPENDJ-1115). *

* A change which is listed but which no replay thread owns is taken over by this * delivery: it is the one to replay, and the copy which may still wait in the replay * queue is dropped by {@link #markInProgress(LDAPUpdateMsg)}. That is a change whose * replay failed - it is deliberately left out of the ServerState, so the replication * server sends it again over the session which was restarted (issue #889) - and it is * also, harmlessly, a change which was listed a moment ago by a delivery no replay * thread has picked up yet: the two deliveries carry the same change, and the last one * listed is the one replayed. *

* The failures the change went through are kept, whichever delivery replays it: they * are what has this replica eventually give up on a change it can not apply. * * @param update The LDAPUpdateMsg that was received from the replication * server and that will be added to the pending list. * @return {@code false} if the update was already registered in the pending * changes and is owned by a replay thread. */ public boolean putRemoteUpdate(LDAPUpdateMsg update) { pendingChangesWriteLock.lock(); try { final CSN csn = update.getCSN(); final PendingChange listed = pendingChanges.get(csn); if (listed == null) { pendingChanges.put(csn, new PendingChange(csn, null, update)); return true; } if (listed.isCommitted() || listed.isOwned()) { return false; } listed.setMsg(update); return true; } finally { pendingChangesWriteLock.unlock(); } } /** * Mark an update message as committed. *

* A change another replay thread owns is not this thread's to record: it is reported as * a change which is not here, the way one which is not listed anymore is. A thread * decides that a change is in the data, or that it is to be given up on, and records * that decision a turn of this lock later - long enough for the change to have been * handed back, delivered again and taken over in between. Recording it then would * advance the ServerState over a change which is not in the data yet (issue #889) and * have the thread which is applying it right now fail to commit. *

* The give-back on the way out of an unwound replay is what makes this reachable: it * runs wherever the replay was left, so the checks which keep it from taking a change * away from the thread which owns it now belong on every road which reads ownership, * not only on {@link #replayFailed(CSN)} (issue #922). * * @param csn * The CSN of the update message that must be set as committed. * @throws NoSuchElementException * if there is no change with that CSN for this thread to record: it is not * listed as pending anymore, or another replay thread owns it */ public void commit(CSN csn) { pendingChangesWriteLock.lock(); try { PendingChange curChange = pendingChanges.get(csn); if (curChange == null || (curChange.isOwned() && !curChange.isOwnedBy(Thread.currentThread()))) { throw new NoSuchElementException(); } curChange.setCommitted(true); curChange.setOwner(null); changeBeingReplayed.remove(Thread.currentThread(), csn); activeAndDependentChanges.remove(curChange); final Iterator it = pendingChanges.values().iterator(); while (it.hasNext()) { PendingChange pendingChange = it.next(); if (!pendingChange.isCommitted()) { break; } if (pendingChange.getMsg().contributesToDomainState()) { state.update(pendingChange.getCSN()); } if (pendingChange.getReplayFailures() > 0) { // The change is leaving this map, so it is not one of the failing ones anymore. failingChanges--; } it.remove(); } } finally { pendingChangesWriteLock.unlock(); } } /** * Gives up the ownership a replay thread had on the change with the provided CSN, * whose replay failed. *

* The change stays listed here and stays uncommitted: it is the barrier which keeps * the ServerState - and every change which follows it - from moving past a change * which is not in the data (issue #889), and it is what has the replication server * send it again over the restarted session. Only the mark which says a replay thread * owns it is dropped, so that {@link #putRemoteUpdate(LDAPUpdateMsg)} takes the next * delivery instead of discarding it as a duplicate. *

* It also stays listed among the changes the newer ones are checked against: a change * which is not in the data yet is exactly what the changes which follow it must depend * on, whether or not a replay thread owns it right now. *

* The changes another replay thread is applying right now are left alone: they are * about to commit, and forgetting them would have their {@code commit()} fail, the * ServerState stay behind them and the replication server replay them a second time. *

* Only the thread which owns the change gives it back. A release which arrives from * another one is a release of a change which has been taken over since - the failure of * a delivery is reported after the change it carried was handed to the delivery which * follows it - and taking the change away from the thread which is replaying it right * now is the double replay the ownership is there to prevent (issue #922). * * @param csn the CSN of the change whose replay failed */ public void replayFailed(CSN csn) { pendingChangesWriteLock.lock(); try { final PendingChange change = pendingChanges.get(csn); if (change != null && !change.isCommitted() && change.isOwnedBy(Thread.currentThread())) { change.setOwner(null); changeBeingReplayed.remove(Thread.currentThread(), csn); } } finally { pendingChangesWriteLock.unlock(); } } /** How long, and how many times, the replay of one change has been failing. */ static final class ReplayFailure { private final int attempts; private final long failingForMs; private ReplayFailure(int attempts, long failingForMs) { this.attempts = attempts; this.failingForMs = failingForMs; } /** * Returns how many deliveries of the change failed to replay it in a row. A delivery * is attempted several times in place before it is counted here, so this is a count * of deliveries rather than of attempts made on the backend. * * @return the number of failed deliveries, at least 1 */ int getAttempts() { return attempts; } /** * Returns how long the replay of the change has been failing, that is the time * between its first failure and the one which was just recorded. * * @return the duration in milliseconds, 0 for a first failure */ long getFailingForMs() { return failingForMs; } } /** * Records that the replay of the change with the provided CSN failed once more. *

* The failures are kept on the change itself, which stays listed for as long as it has * not been applied, so a change which keeps failing keeps its give-up budget across the * deliveries which take over from one another, and a change which is replayed or given * up on takes its failures away with it (issue #889). * * @param csn * the CSN of the change whose replay failed * @param nowMs * when it failed, on a clock which only moves forward * @return the failures of the change, or {@code null} when there is no change left here * for this thread to give up on: it is not listed as an uncommitted change * anymore, which happens when the domain was disabled while it was being * replayed, or another replay thread owns it, which is that change having been * delivered again and taken over while this thread was on its way to reporting * on it. Spending the give-up budget of a change another thread is applying * would have this replica skip a change which is being written (issue #922) */ public ReplayFailure recordReplayFailure(CSN csn, long nowMs) { pendingChangesWriteLock.lock(); try { final PendingChange change = pendingChanges.get(csn); if (change == null || change.isCommitted() || (change.isOwned() && !change.isOwnedBy(Thread.currentThread()))) { return null; } if (change.getReplayFailures() == 0) { // Its first failure: this change joins the ones which are failing right now. failingChanges++; } change.recordReplayFailure(nowMs); return new ReplayFailure(change.getReplayFailures(), change.getReplayFailingForMs(nowMs)); } finally { pendingChangesWriteLock.unlock(); } } /** * Returns whether the replay of any change listed here is failing right now. *

* A change is failing from its first failed replay until it leaves this map, whether * it leaves it applied or given up on. The session restart backoff reads this: a * change which was replayed only says that this backend is serving again when it is * the last one which was failing, and the successful replays which surround a change * this replica can not apply must not keep resetting the wait it has reached on it. * * @return {@code true} while at least one listed change has a failed replay recorded * against it */ boolean hasFailingChanges() { pendingChangesReadLock.lock(); try { return failingChanges > 0; } finally { pendingChangesReadLock.unlock(); } } /** * Returns how many of the listed changes have a failed replay recorded against them. *

* A change is counted from its first failed replay until it leaves this map, whether it * leaves it applied or given up on. It is what tells that the replay of this domain is * stuck: a domain whose {@code replay-give-up-delay} is unlimited never gives up, so it * never counts a failed change either, and the changes it keeps asking for are only * visible here. * * @return the number of listed changes with a failed replay recorded against them */ public int getFailingChangesSize() { pendingChangesReadLock.lock(); try { return failingChanges; } finally { pendingChangesReadLock.unlock(); } } /** * Forgets every change listed here, without updating the ServerState. *

* Called when the domain is disabled: its ServerState is saved and cleared from * memory, and it is loaded again from the backend when the domain is enabled back, so * the bookkeeping which goes with it must not outlive it. A change which stayed here * would be discarded as a duplicate when the replication server sends it again, and * nothing would ever replay it or record it in the ServerState. */ public void clear() { pendingChangesWriteLock.lock(); dependentChangesLock.lock(); try { pendingChanges.clear(); dependentChanges.clear(); activeAndDependentChanges.clear(); changeBeingReplayed.clear(); failingChanges = 0; } finally { dependentChangesLock.unlock(); pendingChangesWriteLock.unlock(); } } /** * Marks the change of the provided message as being replayed. * * @param msg * the message whose change is being replayed * @return {@code false} if this message is not the delivery which is listed as * pending, which happens when the session was restarted after a failed replay * while this message was still waiting in the replay queue: the replication * server delivered the change again and that delivery took over from this one, * or the domain was disabled and forgot the change, so this copy must not be * replayed. */ public boolean markInProgress(LDAPUpdateMsg msg) { pendingChangesWriteLock.lock(); try { final PendingChange change = pendingChanges.get(msg.getCSN()); if (change == null || change.isCommitted() || change.getLDAPUpdateMsg() != msg) { return false; } /* * Listed as being replayed before it is owned, and not the other way round: the * caller enters the replay - where the give-back on the way out lives - once this * returns, so nothing which allocates must run between the owner being stamped and * that. A change listed here without an owner is the state a failed replay leaves * behind, and the dependency checks which read this set do not read the owner. */ activeAndDependentChanges.add(change); changeBeingReplayed.put(Thread.currentThread(), change.getCSN()); change.setOwner(Thread.currentThread()); return true; } finally { pendingChangesWriteLock.unlock(); } } /** * Returns the CSN of the change the calling thread is replaying, when it still owns one. *

* A thread owns the change it is replaying and the ones it parked as waiting for another * change. The parked ones are left out: they are handed to whichever thread clears the * change they are waiting for, and that thread takes them over, so giving one back here * would have the same change handed to two threads (issue #922). *

* It is a plain read of {@link #changeBeingReplayed}: no lock is taken and nothing is * allocated. This is what the give-back on the way out of an unwound replay asks first, * and it runs on a road an {@code OutOfMemoryError} leads to - a lookup which allocated * could be refused in its turn, and a give-back which does not know which change to give * back leaves it listed, uncommitted and owned by a thread which is not replaying it * anymore, which is the wedge this issue is about. *

* An answer which is out of date is safe: every road which acts on it - {@code commit()}, * {@link #recordReplayFailure(CSN, long)} and {@link #replayFailed(CSN)} - checks the * ownership of the change again under the write lock, and is a no-op for a change this * thread does not own anymore. * * @return the CSN of the change this thread is replaying, or {@code null} when it does * not own one anymore - the road its replay took gave it back, or it was applied */ CSN getChangeOwnedByCurrentThread() { return changeBeingReplayed.get(Thread.currentThread()); } /** * Get the first update in the list that have some dependencies cleared. *

* The change is handed to the calling thread, which owns it from then on: it is * replayed by whichever replay thread cleared the change it was waiting for rather than * by the one which parked it, and a change is given back by the thread which owns it * and by nobody else (issue #922). * * @return The LDAPUpdateMsg to be handled. */ public LDAPUpdateMsg getNextUpdate() { pendingChangesReadLock.lock(); dependentChangesLock.lock(); try { if (!hasChangeToHandOut()) { /* * Nothing is waiting, or what waits is still held back by the changes before it. * This is called at the end of every replay, by every replay thread, so the answer * is looked for under the read lock: taking the write lock here would have a * backlog of waiting changes serialize the replay of the changes which have none. */ return null; } } finally { dependentChangesLock.unlock(); pendingChangesReadLock.unlock(); } /* * There is one to hand out, and handing it out writes its owner, which is written * under the write lock as the rest of the state of a change is. It is looked for again * under that lock: another replay thread may have been handed it in between. */ pendingChangesWriteLock.lock(); dependentChangesLock.lock(); try { if (hasChangeToHandOut()) { final PendingChange firstDependentChange = dependentChanges.first(); /* * Entered as the change this thread is replaying before it is taken out of the * ones which are waiting, for the same reason markInProgress() enters it before it * stamps the owner: an allocation which fails here must leave the change where it * was rather than take it out of the hands which would hand it out again. */ changeBeingReplayed.put(Thread.currentThread(), firstDependentChange.getCSN()); dependentChanges.remove(firstDependentChange); firstDependentChange.setOwner(Thread.currentThread()); return firstDependentChange.getLDAPUpdateMsg(); } return null; } finally { dependentChangesLock.unlock(); pendingChangesWriteLock.unlock(); } } /** * Returns whether the first change waiting for another one can be replayed now, that is * whether every change before it has left the pending changes. */ @GuardedBy("pendingChangesLock, dependentChangesLock") private boolean hasChangeToHandOut() { return !dependentChanges.isEmpty() && !pendingChanges.isEmpty() && pendingChanges.firstKey().isNewerThanOrEqualTo(dependentChanges.first().getCSN()); } /** * Mark the first pendingChange as dependent on the second PendingChange. * @param dependentChange The PendingChange that depend on the second * PendingChange. */ private void addDependency(PendingChange dependentChange) { pendingChangesReadLock.lock(); dependentChangesLock.lock(); try { /* * A change which is not listed as pending anymore - the domain was disabled while * its dependencies were computed - must not be listed as dependent either: * getNextUpdate() reads both and would hand out a change nothing owns. The * replication server sends it again when the domain is enabled back. */ if (pendingChanges.containsKey(dependentChange.getCSN())) { dependentChanges.add(dependentChange); } /* * Whichever of the two it was, this thread is not replaying that change anymore: a * parked one is handed to the thread which clears what it waits for, and one which is * not listed here anymore is gone with the pending changes of a domain which was * disabled. The owner stays as it is - it is what has getNextUpdate() hand the change * over rather than leave it to nobody - and the give-back on the way out of an * unwound replay leaves it alone (issue #922). */ changeBeingReplayed.remove(Thread.currentThread(), dependentChange.getCSN()); } finally { dependentChangesLock.unlock(); pendingChangesReadLock.unlock(); } } private PendingChange getPendingChange(CSN csn) { pendingChangesReadLock.lock(); try { return pendingChanges.get(csn); } finally { pendingChangesReadLock.unlock(); } } /** * Check if the given AddOperation has some dependencies on any * currently running previous operation. * Update the dependency list in the associated PendingChange if * there are some dependencies. * AddOperation depends on * * - DeleteOperation done on the same DN * - ModifyDnOperation with the same target DN as the ADD DN * - ModifyDnOperation with new DN equals to the ADD DN parent * - AddOperation done on the parent DN of the ADD DN * * @param op The AddOperation to be checked. * * @return A boolean indicating if this operation has some dependencies. */ public boolean checkDependencies(AddOperation op) { final CSN csn = OperationContext.getCSN(op); final PendingChange change = getPendingChange(csn); if (change == null) { return false; } boolean hasDependencies = false; final DN targetDN = op.getEntryDN(); for (PendingChange pendingChange : activeAndDependentChanges) { if (pendingChange.getCSN().isNewerThanOrEqualTo(csn)) { // From now on, the dependency should be for newer changes to be dependent on this one, so we can stop for now. break; } final LDAPUpdateMsg pendingMsg = pendingChange.getLDAPUpdateMsg(); if (pendingMsg instanceof DeleteMsg) { if (pendingMsg.getDN().equals(targetDN)) { // it is a deleteOperation on the same DN hasDependencies = true; addDependency(change); } } else if (pendingMsg instanceof AddMsg) { if (pendingMsg.getDN().isSuperiorOrEqualTo(targetDN)) { // it is an addOperation on a parent of the current AddOperation hasDependencies = true; addDependency(change); } } else if (pendingMsg instanceof ModifyDNMsg) { // it is a ModifyDnOperation with the same target DN as the ADD DN // or a ModifyDnOperation with new DN equals to the ADD DN parent? if (pendingMsg.getDN().equals(targetDN)) { hasDependencies = true; addDependency(change); } else { final ModifyDNMsg pendingModDn = (ModifyDNMsg) pendingMsg; if (pendingModDn.newDNIsParent(targetDN)) { hasDependencies = true; addDependency(change); } } } } return hasDependencies; } /** * Check if the given ModifyOperation has some dependencies on any * currently running previous operation. * Update the dependency list in the associated PendingChange if * there are some dependencies. * * ModifyOperation depends on * - AddOperation done on the same DN * - ModifyDNOperation having newDN the same as targetDN * * @param op The ModifyOperation to be checked. * * @return A boolean indicating if this operation has some dependencies. */ public boolean checkDependencies(ModifyOperation op) { final CSN csn = OperationContext.getCSN(op); final PendingChange change = getPendingChange(csn); if (change == null) { return false; } boolean hasDependencies = false; final DN targetDN = change.getLDAPUpdateMsg().getDN(); for (PendingChange pendingChange : activeAndDependentChanges) { if (pendingChange.getCSN().isNewerThanOrEqualTo(csn)) { // From now on, the dependency should be for newer changes to be dependent on this one, so we can stop for now. break; } final LDAPUpdateMsg pendingMsg = pendingChange.getLDAPUpdateMsg(); if (pendingMsg instanceof AddMsg) { if (pendingMsg.getDN().equals(targetDN)) { // it is an addOperation on a same DN hasDependencies = true; addDependency(change); } } else if (pendingMsg instanceof ModifyDNMsg) { if (((ModifyDNMsg) pendingMsg).newDNIsEqual(targetDN)) { hasDependencies = true; addDependency(change); } } else if (pendingMsg instanceof ModifyMsg) { if (pendingMsg.getDN().equals(targetDN)) { // it is another modify on the same DN, they depend hasDependencies = true; addDependency(change); } } } return hasDependencies; } /** * Check if the given ModifyDNMsg has some dependencies on any * currently running previous operation. * Update the dependency list in the associated PendingChange if * there are some dependencies. * * Modify DN Operation depends on * - AddOperation done on the same DN as the target DN of the MODDN operation * - AddOperation done on the new parent of the MODDN operation * - DeleteOperation done on the new DN of the MODDN operation * - ModifyDNOperation done from the new DN of the MODDN operation * * TODO: Consider cases where there is a rename A -> B then rename B -> C. Second change depends on first * * @param msg The ModifyDNMsg to be checked. * * @return A boolean indicating if this operation has some dependencies. */ public boolean checkDependencies(ModifyDNMsg msg) { final CSN csn = msg.getCSN(); final PendingChange change = getPendingChange(csn); if (change == null) { return false; } boolean hasDependencies = false; final DN targetDN = change.getLDAPUpdateMsg().getDN(); for (PendingChange pendingChange : activeAndDependentChanges) { if (pendingChange.getCSN().isNewerThanOrEqualTo(csn)) { // From now on, the dependency should be for newer changes to be dependent on this one, so we can stop for now. break; } final LDAPUpdateMsg pendingMsg = pendingChange.getLDAPUpdateMsg(); if (pendingMsg instanceof DeleteMsg) { // Check if the target of the Delete is the same // as the new DN of this ModifyDN if (msg.newDNIsEqual(pendingMsg.getDN())) { hasDependencies = true; addDependency(change); } } else if (pendingMsg instanceof AddMsg) { // Check if the Add Operation was done on the new parent of // the MODDN operation if (msg.newParentIsEqual(pendingMsg.getDN())) { hasDependencies = true; addDependency(change); } // Check if the AddOperation was done on the same DN as the // target DN of the MODDN operation if (pendingMsg.getDN().equals(targetDN)) { hasDependencies = true; addDependency(change); } } else if (pendingMsg instanceof ModifyDNMsg) { if (msg.newDNIsEqual(pendingMsg.getDN())) { // the ModifyDNOperation was done from the new DN of the MODDN operation hasDependencies = true; addDependency(change); } } } return hasDependencies; } /** * Check if the given DeleteOperation has some dependencies on any * currently running previous operation. * Update the dependency list in the associated PendingChange if * there are some dependencies. * * DeleteOperation depends on * - DeleteOperation done on children DN * - ModifyDnOperation with target DN that are children of the DEL DN * - AddOperation done on the same DN * * * @param op The DeleteOperation to be checked. * * @return A boolean indicating if this operation has some dependencies. */ public boolean checkDependencies(DeleteOperation op) { final CSN csn = OperationContext.getCSN(op); final PendingChange change = getPendingChange(csn); if (change == null) { return false; } boolean hasDependencies = false; final DN targetDN = op.getEntryDN(); for (PendingChange pendingChange : activeAndDependentChanges) { if (pendingChange.getCSN().isNewerThanOrEqualTo(csn)) { // From now on, the dependency should be for newer changes to be dependent on this one, so we can stop for now. break; } final LDAPUpdateMsg pendingMsg = pendingChange.getLDAPUpdateMsg(); if (pendingMsg instanceof DeleteMsg) { /* * Check if the operation to be run is a deleteOperation on a * children of the current DeleteOperation. */ if (pendingMsg.getDN().isSubordinateOrEqualTo(targetDN)) { hasDependencies = true; addDependency(change); } } else if (pendingMsg instanceof AddMsg) { /* * Check if the operation to be run is an addOperation on a * parent of the current DeleteOperation. */ if (pendingMsg.getDN().equals(targetDN)) { hasDependencies = true; addDependency(change); } } else if (pendingMsg instanceof ModifyDNMsg) { final ModifyDNMsg pendingModDn = (ModifyDNMsg) pendingMsg; /* * Check if the operation to be run is an ModifyDNOperation * on a children of the current DeleteOperation */ if (pendingMsg.getDN().isSubordinateOrEqualTo(targetDN) || pendingModDn.newDNIsParent(targetDN)) { hasDependencies = true; addDependency(change); } } } return hasDependencies; } /** * Check the dependencies of a given Operation/UpdateMsg. * * @param op The Operation for which dependencies must be checked. * @param msg The LocalizableMessage for which dependencies must be checked. * @return A boolean indicating if an operation cannot be replayed * because of dependencies. */ public boolean checkDependencies(Operation op, LDAPUpdateMsg msg) { if (op instanceof ModifyOperation) { return checkDependencies((ModifyOperation) op); } else if (op instanceof DeleteOperation) { return checkDependencies((DeleteOperation) op); } else if (op instanceof AddOperation) { return checkDependencies((AddOperation) op); } else if (op instanceof ModifyDNOperationBasis) { return checkDependencies((ModifyDNMsg) msg); } else { return true; // unknown type of operation ?! } } }