[#807] Do not drop persistent search notifications through search-phase dedup (#812)
| | |
| | | <org.opends.test.pauseOnFailure>false</org.opends.test.pauseOnFailure> |
| | | <org.opends.test.copyClassesToTestPackage>false</org.opends.test.copyClassesToTestPackage> |
| | | <org.opends.test.timeout>600000</org.opends.test.timeout><!--15 mins--> |
| | | <org.opends.test.trace.pattern>(org\.opends\.server\.replication\.service\..*)|(org\.opends\.server\.replication\.GenerationIdTest)|(org\.opends\.server\.types.\HostPortTest)</org.opends.test.trace.pattern> |
| | | <!-- Matched against the name of the test class, see org.opends.server.TestListener.onStart(). --> |
| | | <org.opends.test.trace.pattern>(org\.opends\.server\.replication\.service\..*)|(org\.opends\.server\.replication\.GenerationIdTest)|(org\.opends\.server\.types\.HostPortTest)|(org\.openidentityplatform\.opendj\.AliasTestCase)</org.opends.test.trace.pattern> |
| | | </systemPropertyVariables> |
| | | <argLine>@{argLine}</argLine> |
| | | <reuseForks>false</reuseForks> |
| | |
| | | * |
| | | * Copyright 2006-2010 Sun Microsystems, Inc. |
| | | * Portions Copyright 2010-2016 ForgeRock AS. |
| | | * Portions Copyright 2026 3A Systems, LLC. |
| | | */ |
| | | package org.forgerock.opendj.reactive; |
| | | |
| | |
| | | // if operation processing encounters a run-time exception after sending the |
| | | // response: the worker thread exception handling code will attempt to send |
| | | // an error result to the client indicating that a problem occurred. |
| | | if (removeOperationInProgress(operation.getMessageID())) { |
| | | // A persistent search is the other way around: its search operation is no longer in |
| | | // progress once the search phase is over, and yet it still owes the client a response if |
| | | // the server terminates it. |
| | | if (removeOperationInProgress(operation.getMessageID()) |
| | | || hasPersistentSearch(operation.getMessageID())) { |
| | | final Response response = operationToResponse(operation); |
| | | final FlowableEmitter<Response> out = getAttachedEmitter(operation); |
| | | if (response != null) { |
| | |
| | | * |
| | | * Copyright 2006-2009 Sun Microsystems, Inc. |
| | | * Portions Copyright 2011-2016 ForgeRock AS. |
| | | * Portions Copyright 2025 3A Systems, LLC. |
| | | * Portions Copyright 2025-2026 3A Systems, LLC. |
| | | */ |
| | | package org.opends.server.api; |
| | | |
| | |
| | | return persistentSearches; |
| | | } |
| | | |
| | | /** |
| | | * Indicates whether a persistent search is registered on this connection for the provided message |
| | | * ID. A persistent search outlives the operation which started it: that operation leaves the set |
| | | * of operations in progress as soon as its search phase is over, but the server can still have a |
| | | * final response to send for it, when the search is terminated on the server side. |
| | | * |
| | | * @param messageID The message ID to look for. |
| | | * @return {@code true} if a persistent search is registered for the provided message ID. |
| | | */ |
| | | protected final boolean hasPersistentSearch(int messageID) |
| | | { |
| | | for (PersistentSearch psearch : persistentSearches) |
| | | { |
| | | if (psearch.getMessageID() == messageID) |
| | | { |
| | | return true; |
| | | } |
| | | } |
| | | return false; |
| | | } |
| | | |
| | | |
| | | |
| | | /** |
| | |
| | | import java.util.Set; |
| | | import java.util.concurrent.ConcurrentLinkedQueue; |
| | | |
| | | import org.forgerock.i18n.LocalizableMessage; |
| | | import org.forgerock.i18n.slf4j.LocalizedLogger; |
| | | import org.forgerock.opendj.config.Configuration; |
| | | import org.forgerock.opendj.config.server.ConfigException; |
| | | import org.forgerock.opendj.ldap.ConditionResult; |
| | |
| | | public abstract class LocalBackend<C extends Configuration> extends Backend<C> |
| | | // should have been BackendCfg instead of Configuration |
| | | { |
| | | private static final LocalizedLogger logger = LocalizedLogger.getLoggerForThisClass(); |
| | | |
| | | /** Indicates whether this is a private backend or one that holds user data. */ |
| | | private boolean isPrivateBackend; |
| | | |
| | |
| | | { |
| | | for (PersistentSearch psearch : persistentSearches) |
| | | { |
| | | psearch.cancel(); |
| | | // Tell the clients that no more changes are coming: this backend will not notify them any |
| | | // more, and a cancelled persistent search which sends nothing leaves them waiting forever. |
| | | final LocalizableMessage reason = WARN_PSEARCH_BACKEND_UNAVAILABLE.get(getBackendID()); |
| | | logger.warn(reason); |
| | | psearch.cancelAndNotifyClient(reason); |
| | | } |
| | | persistentSearches.clear(); |
| | | closeBackend(); |
| | |
| | | import java.util.LinkedHashMap; |
| | | import java.util.List; |
| | | import java.util.Map; |
| | | import java.util.Queue; |
| | | import java.util.Set; |
| | | import java.util.TimeZone; |
| | | import java.util.concurrent.ConcurrentLinkedQueue; |
| | |
| | | { |
| | | final SearchOperation searchOp = pSearch.getSearchOperation(); |
| | | final CookieEntrySender entrySender = searchOp.getAttachment(ENTRY_SENDER_ATTACHMENT); |
| | | entrySender.persistentSearchSendEntry(baseDN, updateMsg); |
| | | if (!entrySender.persistentSearchSendEntry(baseDN, updateMsg)) |
| | | { |
| | | stopPersistentSearch(pSearch); |
| | | } |
| | | } |
| | | } |
| | | catch (DirectoryException e) |
| | |
| | | { |
| | | final SearchOperation searchOp = pSearch.getSearchOperation(); |
| | | final ChangeNumberEntrySender entrySender = searchOp.getAttachment(ENTRY_SENDER_ATTACHMENT); |
| | | entrySender.persistentSearchSendEntry(changeNumber, changeNumberEntry); |
| | | if (!entrySender.persistentSearchSendEntry(changeNumber, changeNumberEntry)) |
| | | { |
| | | stopPersistentSearch(pSearch); |
| | | } |
| | | } |
| | | } |
| | | catch (DirectoryException e) |
| | |
| | | { |
| | | initializePersistentSearch(pSearch); |
| | | |
| | | if (isCookieBased(pSearch.getSearchOperation())) |
| | | final Queue<PersistentSearch> psearches = isCookieBased(pSearch.getSearchOperation()) |
| | | ? cookieBasedPersistentSearches |
| | | : changeNumberBasedPersistentSearches; |
| | | psearches.add(pSearch); |
| | | // Without this, a cancelled persistent search keeps being handed the changes it can no longer |
| | | // report, for as long as this backend lives. |
| | | pSearch.registerCancellationCallback(new PersistentSearch.CancellationCallback() |
| | | { |
| | | cookieBasedPersistentSearches.add(pSearch); |
| | | } |
| | | else |
| | | @Override |
| | | public void persistentSearchCancelled(PersistentSearch psearch) |
| | | { |
| | | changeNumberBasedPersistentSearches.add(pSearch); |
| | | psearches.remove(psearch); |
| | | } |
| | | }); |
| | | super.registerPersistentSearch(pSearch); |
| | | } |
| | | |
| | |
| | | return true; |
| | | } |
| | | |
| | | /** |
| | | * Sends a change reported by the "persistent search" phase, if it matches the base, scope and |
| | | * filter of the current search operation. Contrary to the "initial search" phase, the change goes |
| | | * through the persistent search path: it is not bound by the size and time limits of the search, |
| | | * which are only lifted once the initial phase is over, and a change published in the meantime |
| | | * would be dropped without the client ever hearing about it. |
| | | * |
| | | * @return {@code true} if the persistent search should keep reporting changes, {@code false} |
| | | * otherwise |
| | | */ |
| | | private static boolean sendNotificationIfMatches(SearchOperation searchOp, Entry entry, String cookie) |
| | | throws DirectoryException |
| | | { |
| | | if (matchBaseAndScopeAndFilter(searchOp, entry)) |
| | | { |
| | | return searchOp.returnPersistentSearchEntry(entry, getControls(cookie)); |
| | | } |
| | | // maybe the next entry will match? |
| | | return true; |
| | | } |
| | | |
| | | /** |
| | | * Stops the provided persistent search and tells the client, which would otherwise wait forever |
| | | * for changes on a search which no longer reports any. |
| | | */ |
| | | private static void stopPersistentSearch(PersistentSearch pSearch) |
| | | { |
| | | try |
| | | { |
| | | // Before the cancellation, which deregisters this persistent search from the connection: the |
| | | // search operation left the operations in progress when its initial phase ended, so nothing |
| | | // would be left to hang the response on afterwards. |
| | | pSearch.getSearchOperation().sendSearchResultDone(); |
| | | } |
| | | catch (Exception e) |
| | | { |
| | | logger.traceException(e); |
| | | } |
| | | pSearch.cancel(); |
| | | } |
| | | |
| | | /** Indicates if the provided entry matches the filter, base and scope. */ |
| | | private static boolean matchBaseAndScopeAndFilter(SearchOperation searchOp, Entry entry) throws DirectoryException |
| | | { |
| | |
| | | return sendEntryIfMatches(searchOp, entry, null); |
| | | } |
| | | |
| | | private void persistentSearchSendEntry(long changeNumber, Entry entry) throws DirectoryException |
| | | /** |
| | | * @return {@code true} if the persistent search should keep reporting changes, {@code false} |
| | | * otherwise |
| | | */ |
| | | private boolean persistentSearchSendEntry(long changeNumber, Entry entry) throws DirectoryException |
| | | { |
| | | if (sendEntryData.persistentSearchCanSendEntry(changeNumber)) |
| | | { |
| | | sendEntryIfMatches(searchOp, entry, null); |
| | | return sendNotificationIfMatches(searchOp, entry, null); |
| | | } |
| | | return true; |
| | | } |
| | | } |
| | | |
| | |
| | | return sendEntryIfMatches(searchOp, entry, cookieString); |
| | | } |
| | | |
| | | private void persistentSearchSendEntry(DN baseDN, UpdateMsg updateMsg) |
| | | /** |
| | | * @return {@code true} if the persistent search should keep reporting changes, {@code false} |
| | | * otherwise |
| | | */ |
| | | private boolean persistentSearchSendEntry(DN baseDN, UpdateMsg updateMsg) |
| | | throws DirectoryException |
| | | { |
| | | final CSN csn = updateMsg.getCSN(); |
| | |
| | | final Entry cookieEntry = createEntryFromMsg(baseDN, 0, cookieString, updateMsg); |
| | | // FIXME JNR use this instead of previous line: |
| | | // entry.replaceAttribute(Attributes.create("changelogcookie", cookieString)); |
| | | sendEntryIfMatches(searchOp, cookieEntry, cookieString); |
| | | return sendNotificationIfMatches(searchOp, cookieEntry, cookieString); |
| | | } |
| | | return true; |
| | | } |
| | | |
| | | private String updateCookie(DN baseDN, final CSN csn) |
| | |
| | | * |
| | | * Copyright 2006-2010 Sun Microsystems, Inc. |
| | | * Portions Copyright 2014-2016 ForgeRock AS. |
| | | * Portions Copyright 2026 3A Systems, LLC. |
| | | */ |
| | | package org.opends.server.core; |
| | | |
| | |
| | | import java.util.Set; |
| | | import java.util.concurrent.CopyOnWriteArrayList; |
| | | |
| | | import org.forgerock.i18n.LocalizableMessage; |
| | | import org.forgerock.i18n.slf4j.LocalizedLogger; |
| | | import org.forgerock.opendj.ldap.ResultCode; |
| | | import org.opends.server.controls.EntryChangeNotificationControl; |
| | |
| | | } |
| | | |
| | | /** |
| | | * Cancels this persistent search and tells the client that no more changes will be reported for |
| | | * it. Contrary to {@link #cancel()}, which leaves the search open as far as the client can tell, |
| | | * this is meant for cancellations decided by the server: without a search result done, the client |
| | | * waits forever for changes on a search which no longer exists. |
| | | * |
| | | * @param reason |
| | | * The reason why this persistent search is terminated, reported to the client. |
| | | * @return The result of the cancellation. |
| | | */ |
| | | public synchronized CancelResult cancelAndNotifyClient(LocalizableMessage reason) |
| | | { |
| | | if (isCancelled) |
| | | { |
| | | // Whoever cancelled this search first is responsible for what the client was told: a second |
| | | // search result done for the same message ID would break the protocol. |
| | | return new CancelResult(ResultCode.CANCELLED, null); |
| | | } |
| | | |
| | | try |
| | | { |
| | | searchOperation.setResultCode(ResultCode.UNAVAILABLE); |
| | | searchOperation.appendErrorMessage(reason); |
| | | // The response is sent before the cancellation on purpose: the search operation left the set |
| | | // of operations in progress when its search phase ended, so the connection only knows it as |
| | | // this persistent search, which cancelling deregisters. |
| | | searchOperation.sendSearchResultDone(); |
| | | } |
| | | catch (Exception e) |
| | | { |
| | | // The client may be gone already: the persistent search is cancelled either way. |
| | | logger.traceException(e); |
| | | } |
| | | return cancel(); |
| | | } |
| | | |
| | | /** |
| | | * Gets the message ID associated with this persistent search. |
| | | * |
| | | * @return The message ID associated with this persistent search. |
| | |
| | | { |
| | | try |
| | | { |
| | | if (!searchOperation.returnEntry(entry, entryControls)) |
| | | // Notifications go through their own path: a change must be reported whether or not the |
| | | // entry was already returned by the search phase, and for as long as this search lives. |
| | | if (!searchOperation.returnPersistentSearchEntry(entry, entryControls)) |
| | | { |
| | | cancel(); |
| | | // Send the response first: cancelling deregisters this persistent search, and the search |
| | | // operation is no longer in progress on the connection either, so there would be nothing |
| | | // left to hang the response on. |
| | | searchOperation.sendSearchResultDone(); |
| | | cancel(); |
| | | } |
| | | } |
| | | catch (Exception e) |
| | | { |
| | | logger.traceException(e); |
| | | |
| | | cancel(); |
| | | |
| | | try |
| | | { |
| | | searchOperation.sendSearchResultDone(); |
| | |
| | | { |
| | | logger.traceException(e2); |
| | | } |
| | | |
| | | cancel(); |
| | | } |
| | | } |
| | | |
| | |
| | | boolean evaluateAci); |
| | | |
| | | /** |
| | | * Used as a callback for persistent searches to send an entry which has just changed to the |
| | | * client. Contrary to {@link #returnEntry(Entry, List)}, the entry is not matched against the |
| | | * state kept to dereference aliases during the search phase, and neither the size limit nor the |
| | | * time limit of the search applies to it: a persistent search must report every change it is |
| | | * notified of, for as long as it is alive, whether or not the entry was returned before. |
| | | * |
| | | * @param entry The entry which has changed and should be sent to the client. |
| | | * @param controls The set of controls to include with the entry (may be <CODE>null</CODE> if |
| | | * none are needed). |
| | | * |
| | | * @return <CODE>true</CODE> if the persistent search should keep reporting changes, or |
| | | * <CODE>false</CODE> if it should stop for some reason (e.g. the search has been |
| | | * abandoned). |
| | | */ |
| | | boolean returnPersistentSearchEntry(Entry entry, List<Control> controls); |
| | | |
| | | /** |
| | | * Indicates that the search phase is over and that any further entry comes from a persistent |
| | | * search. State kept to dereference aliases during the search phase is released, and no further |
| | | * entry is matched against it: a persistent search must report every change it is notified of, |
| | | * whether or not the entry was returned by the search phase. |
| | | * search. State kept to dereference aliases during the search phase is released. Entries can |
| | | * still reach {@link #returnEntry(Entry, List)} afterwards, as backends are free to report their |
| | | * own results from another thread, so that method keeps track of this phase being over. |
| | | */ |
| | | void endSearchPhase(); |
| | | |
| | |
| | | import java.util.*; |
| | | import java.util.concurrent.ConcurrentHashMap; |
| | | import java.util.concurrent.atomic.AtomicBoolean; |
| | | import java.util.concurrent.atomic.AtomicInteger; |
| | | |
| | | import org.forgerock.i18n.LocalizedIllegalArgumentException; |
| | | import org.forgerock.i18n.slf4j.LocalizedLogger; |
| | |
| | | /** The proxied authorization target DN for this operation. */ |
| | | private DN proxiedAuthorizationDN; |
| | | |
| | | /** The number of entries that have been sent to the client. */ |
| | | private int entriesSent; |
| | | /** |
| | | * The number of entries that have been sent to the client. Persistent search notifications are |
| | | * sent by the threads which apply the changes, so several of them can be counted concurrently. |
| | | */ |
| | | private final AtomicInteger entriesSent = new AtomicInteger(); |
| | | |
| | | /** |
| | | * The number of search result references that have been sent to the client. |
| | |
| | | @Override |
| | | public final int getEntriesSent() |
| | | { |
| | | return entriesSent; |
| | | return entriesSent.get(); |
| | | } |
| | | |
| | | @Override |
| | |
| | | /** |
| | | * The DNs of the entries already returned by the search phase. An alias may be dereferenced onto |
| | | * an entry which is in the scope of the search as well, and that entry must only be returned once. |
| | | * It is emptied once the search phase is over, because a persistent search must report every |
| | | * change it is notified of, whether or not the entry was returned before. |
| | | * It is emptied once the search phase is over, because it is only meaningful while that phase |
| | | * runs. |
| | | */ |
| | | private final Set<DN> returnedDNs = ConcurrentHashMap.newKeySet(); |
| | | |
| | | /** Whether the search phase is over and only persistent search notifications remain. */ |
| | | /** |
| | | * Whether the search phase is over. Entries still sent through {@link #returnEntry(Entry, List)} |
| | | * after it, as the external changelog backend does for its own notifications, are no longer |
| | | * matched against {@link #returnedDNs}, which has been emptied by then. |
| | | */ |
| | | private volatile boolean searchPhaseOver; |
| | | |
| | | /** Why an entry is being returned to the client. */ |
| | | private enum EntrySource |
| | | { |
| | | /** The entry is a result of the search phase. */ |
| | | SEARCH_PHASE, |
| | | /** |
| | | * The entry reports a change to a persistent search. It must be sent whether or not the same |
| | | * entry was returned before, and neither the size limit nor the time limit of the search |
| | | * applies to it. |
| | | */ |
| | | PSEARCH_NOTIFICATION |
| | | } |
| | | |
| | | @Override |
| | | public final void endSearchPhase() |
| | | { |
| | |
| | | public final boolean returnEntry(Entry entry, List<Control> controls, |
| | | boolean evaluateAci) |
| | | { |
| | | return returnEntry(entry, controls, evaluateAci, null); |
| | | return returnEntry(entry, controls, evaluateAci, EntrySource.SEARCH_PHASE, null); |
| | | } |
| | | |
| | | @Override |
| | | public final boolean returnPersistentSearchEntry(Entry entry, List<Control> controls) |
| | | { |
| | | return returnEntry(entry, controls, true, EntrySource.PSEARCH_NOTIFICATION, null); |
| | | } |
| | | |
| | | /** |
| | |
| | | * @param entry The entry to return. |
| | | * @param controls The controls to attach to the entry. |
| | | * @param evaluateAci Whether the access control handler must be consulted. |
| | | * @param source Whether the entry is a persistent search notification rather than a result |
| | | * of the search phase. |
| | | * @param aliasChain The DNs of the aliases already dereferenced on the way to this entry, or |
| | | * {@code null} if no alias was dereferenced yet. It only spans the current |
| | | * chain, so it cannot grow beyond the length of that chain. |
| | | * @return {@code true} if the search should continue, {@code false} if it should stop. |
| | | */ |
| | | private boolean returnEntry(Entry entry, List<Control> controls, |
| | | boolean evaluateAci, Set<DN> aliasChain) |
| | | boolean evaluateAci, EntrySource source, Set<DN> aliasChain) |
| | | { |
| | | boolean typesOnly = getTypesOnly(); |
| | | |
| | | // Both limits only bound the search phase: they are lifted for the rest of a persistent search |
| | | // once that phase is over, but a notification can reach this point before that happens. |
| | | if (source == EntrySource.SEARCH_PHASE) |
| | | { |
| | | // See if the size limit has been exceeded. If so, then don't send the |
| | | // entry and indicate that the search should end. |
| | | if (getSizeLimit() > 0 && getEntriesSent() >= getSizeLimit()) |
| | |
| | | appendErrorMessage(ERR_SEARCH_TIME_LIMIT_EXCEEDED.get(getTimeLimit())); |
| | | return false; |
| | | } |
| | | } |
| | | |
| | | // Determine whether the provided entry is a subentry and if so whether it |
| | | // should be returned. |
| | |
| | | && !filterIncludesSubentries |
| | | && !isReturnSubentriesOnly()) |
| | | { |
| | | logger.trace("Not sending entry %s: it is a subentry and this search does not ask for " |
| | | + "subentries", entry.getName()); |
| | | return true; |
| | | } |
| | | } |
| | | else if (isReturnSubentriesOnly()) |
| | | { |
| | | // Subentries are visible and normal entries are not. |
| | | logger.trace("Not sending entry %s: this search only asks for subentries", entry.getName()); |
| | | return true; |
| | | } |
| | | |
| | |
| | | SearchResultEntry unfilteredSearchEntry = new SearchResultEntry(entry, controls); |
| | | if (evaluateAci && !getACIHandler().maySend(this, unfilteredSearchEntry)) |
| | | { |
| | | logger.trace("Not sending entry %s: access control forbids it", entry.getName()); |
| | | return true; |
| | | } |
| | | |
| | | //DereferenceAliasesPolicy |
| | | if ( DereferenceAliasesPolicy.ALWAYS.equals(getDerefPolicy()) || DereferenceAliasesPolicy.IN_SEARCHING.equals(getDerefPolicy()) ) { |
| | | if (entry.isAlias() && !baseDN.equals(entry.getName())) { |
| | | return returnAliasedEntry(entry, controls, aliasChain); |
| | | return returnAliasedEntry(entry, controls, source, aliasChain); |
| | | } |
| | | if (!searchPhaseOver && !returnedDNs.add(entry.getName())) { |
| | | if (source == EntrySource.SEARCH_PHASE && !searchPhaseOver && !returnedDNs.add(entry.getName())) { |
| | | // This entry was already returned by the search, through an alias or on its own. |
| | | logger.trace("Not sending entry %s: it was already returned by the search phase", entry.getName()); |
| | | return true; |
| | | } |
| | | } |
| | |
| | | { |
| | | sendSearchEntry(filteredSearchEntry); |
| | | |
| | | entriesSent++; |
| | | entriesSent.incrementAndGet(); |
| | | } |
| | | catch (DirectoryException de) |
| | | { |
| | |
| | | return false; |
| | | } |
| | | } |
| | | else |
| | | { |
| | | logger.trace("Not sending entry %s: a search result entry plugin suppressed it", entry.getName()); |
| | | } |
| | | |
| | | return pluginResult.continueProcessing(); |
| | | } |
| | |
| | | * |
| | | * @param alias The alias entry to dereference. |
| | | * @param controls The controls to attach to the entry. |
| | | * @param source Whether the alias is reported by a persistent search notification rather |
| | | * than by the search phase. |
| | | * @param aliasChain The DNs of the aliases already dereferenced on the way to this alias, or |
| | | * {@code null} if this alias is the first one of the chain. |
| | | * @return {@code true} if the search should continue, {@code false} if it should stop. |
| | | */ |
| | | private boolean returnAliasedEntry(Entry alias, List<Control> controls, Set<DN> aliasChain) |
| | | private boolean returnAliasedEntry(Entry alias, List<Control> controls, |
| | | EntrySource source, Set<DN> aliasChain) |
| | | { |
| | | final DN aliasedDN; |
| | | final Entry aliasedEntry; |
| | |
| | | if (aliasedEntry == null) |
| | | { |
| | | // The alias points to an entry which does not exist: there is nothing to return for it. |
| | | logger.trace("Not dereferencing alias %s: %s does not exist", alias.getName(), aliasedDN); |
| | | return true; |
| | | } |
| | | if (!searchPhaseOver && returnedDNs.contains(aliasedDN)) |
| | | if (source == EntrySource.SEARCH_PHASE && !searchPhaseOver && returnedDNs.contains(aliasedDN)) |
| | | { |
| | | // The aliased entry was already returned by the search. |
| | | logger.trace("Not dereferencing alias %s: %s was already returned by the search phase", |
| | | alias.getName(), aliasedDN); |
| | | return true; |
| | | } |
| | | if (aliasChain == null) |
| | |
| | | if (!aliasChain.add(aliasedDN)) |
| | | { |
| | | // The aliases point at each other: stop before looping forever. |
| | | logger.trace("Not dereferencing alias %s: %s is already part of the alias chain %s", |
| | | alias.getName(), aliasedDN, aliasChain); |
| | | return true; |
| | | } |
| | | return returnEntry(aliasedEntry, controls, true, aliasChain); |
| | | return returnEntry(aliasedEntry, controls, true, source, aliasChain); |
| | | } |
| | | |
| | | private AccessControlHandler<?> getACIHandler() |
| | |
| | | } |
| | | |
| | | @Override |
| | | public boolean returnPersistentSearchEntry(Entry entry, List<Control> controls) |
| | | { |
| | | return getOperation().returnPersistentSearchEntry(entry, controls); |
| | | } |
| | | |
| | | @Override |
| | | public void endSearchPhase() |
| | | { |
| | | getOperation().endSearchPhase(); |
| | |
| | | * |
| | | * Copyright 2006-2010 Sun Microsystems, Inc. |
| | | * Portions Copyright 2010-2016 ForgeRock AS. |
| | | * Portions Copyright 2026 3A Systems, LLC. |
| | | */ |
| | | package org.opends.server.protocols.ldap; |
| | | |
| | |
| | | // if operation processing encounters a run-time exception after sending the |
| | | // response: the worker thread exception handling code will attempt to send |
| | | // an error result to the client indicating that a problem occurred. |
| | | if (removeOperationInProgress(operation.getMessageID())) |
| | | // A persistent search is the other way around: its search operation is no longer in progress |
| | | // once the search phase is over, and yet it still owes the client a response if the server |
| | | // terminates it. |
| | | if (removeOperationInProgress(operation.getMessageID()) |
| | | || hasPersistentSearch(operation.getMessageID())) |
| | | { |
| | | LDAPMessage message = operationToResponseLDAPMessage(operation); |
| | | if (message != null) |
| | |
| | | ERR_SERVICE_DISCOVERY_CONFIG_MANAGER_INIT_MECHANISM_614=Service Discovery Mechanism '%s' initialization failed : %s |
| | | ERR_SERVICE_DISCOVERY_CONFIG_MANAGER_LISTENER_615=Registering Service Discovery Manager's listener failed : %s |
| | | NOTE_IMPORT_MIGRATION_START_616=Migrating %s entries for base DN %s so that they are preserved by the partial import |
| | | WARN_PSEARCH_BACKEND_UNAVAILABLE_617=The persistent search is being terminated because backend %s is \ |
| | | no longer available |
| | |
| | | import org.forgerock.util.Utils; |
| | | import org.opends.server.TestCaseUtils; |
| | | import org.opends.server.core.ModifyOperation; |
| | | import org.opends.server.core.PersistentSearch; |
| | | import org.opends.server.protocols.internal.InternalSearchOperation; |
| | | import org.opends.server.protocols.internal.SearchRequest; |
| | | import org.opends.server.protocols.ldap.LDAPControl; |
| | |
| | | "(objectClass=*)" |
| | | }; |
| | | |
| | | try |
| | | { |
| | | assertEquals(LDAPSearch.run(nullPrintStream(), System.err, args), 11); |
| | | //cancel the persisting persistent search. |
| | | } |
| | | finally |
| | | { |
| | | // Cancel the persistent search itself: search.cancel() only records a cancellation request |
| | | // for the operation, which nothing acts upon now that the thread running it is gone, so the |
| | | // persistent search would stay registered and keep holding the limit set above against |
| | | // whatever runs next in this JVM (a failing test class is rerun in it). |
| | | for (PersistentSearch psearch : search.getClientConnection().getPersistentSearches()) |
| | | { |
| | | if (psearch.getMessageID() == search.getMessageID()) |
| | | { |
| | | psearch.cancel(); |
| | | } |
| | | } |
| | | search.cancel(new CancelRequest(true,LocalizableMessage.EMPTY)); |
| | | |
| | | //Restore the limit configured for the tests. |
| | | ModifyRequest restoreRequest = newModifyRequest("cn=config") |
| | | .addModification(ModificationType.REPLACE, "ds-cfg-max-psearches", "-1"); |
| | | assertEquals(getRootConnection().processModify(restoreRequest).getResultCode(), ResultCode.SUCCESS); |
| | | } |
| | | } |
| | | } |
| | |
| | | import org.forgerock.opendj.ldap.controls.PersistentSearchRequestControl; |
| | | import org.forgerock.opendj.ldap.requests.Requests; |
| | | import org.forgerock.opendj.ldap.requests.SearchRequest; |
| | | import org.forgerock.opendj.ldap.responses.Result; |
| | | import org.forgerock.opendj.ldap.responses.SearchResultEntry; |
| | | import org.forgerock.opendj.ldap.responses.SearchResultReference; |
| | | import org.forgerock.opendj.ldif.ConnectionEntryReader; |
| | |
| | | import org.opends.server.api.LocalBackend; |
| | | import org.opends.server.backends.MemoryBackend; |
| | | import org.opends.server.core.DirectoryServer; |
| | | import org.opends.server.core.PersistentSearch; |
| | | import org.opends.server.protocols.internal.InternalClientConnection; |
| | | import org.opends.server.protocols.internal.InternalSearchOperation; |
| | | import org.opends.server.types.AcceptRejectWarn; |
| | | import org.opends.server.types.Entry; |
| | | import org.testng.annotations.AfterClass; |
| | |
| | | psearch.searchAsync(request, new SearchResultHandler() { |
| | | @Override |
| | | public boolean handleEntry(SearchResultEntry entry) { |
| | | notified.add(entry.getName().toString()); |
| | | // Every notification carries the same DN, so the DN alone cannot tell a lost |
| | | // notification from a duplicated one: record which change is being reported. |
| | | notified.add(entry.getName() + " " + entry.parseAttribute("description").asString()); |
| | | return true; |
| | | } |
| | | |
| | |
| | | |
| | | // searchAsync returns before the server has registered the persistent search, so wait |
| | | // until the backend reports it; otherwise the first modification below can be notified |
| | | // before the search is listening and be missed. |
| | | // before the search is listening and be missed. A failed test is rerun in the same JVM |
| | | // (rerunFailingTestsCount), which can leave the persistent search of the previous run |
| | | // behind, hence the wait for a persistent search on our own base DN. |
| | | final LocalBackend<?> backend = TestCaseUtils.getServerContext() |
| | | .getBackendConfigManager().getLocalBackendById(TestCaseUtils.TEST_BACKEND_ID); |
| | | for (int i = 0; backend.getPersistentSearches().isEmpty() && i < 500; i++) { |
| | | for (int i = 0; !isPersistentSearchRegistered(backend, "ou=psearch,o=test") && i < 500; i++) { |
| | | Thread.sleep(10); |
| | | } |
| | | assertThat(backend.getPersistentSearches()).isNotEmpty(); |
| | | assertThat(isPersistentSearchRegistered(backend, "ou=psearch,o=test")) |
| | | .as("the persistent search was never registered with backend %s", backend.getBackendID()) |
| | | .isTrue(); |
| | | |
| | | // The same entry is modified repeatedly: each change must reach the persistent search. |
| | | for (int i = 1; i <= 3; i++) { |
| | | connection.modify(Requests.newModifyRequest("cn=changing,ou=psearch,o=test") |
| | | .addModification(ModificationType.REPLACE, "description", "change " + i)); |
| | | assertThat(notified.poll(30, TimeUnit.SECONDS)) |
| | | .as("notification for change " + i) |
| | | .isEqualTo("cn=changing,ou=psearch,o=test"); |
| | | .as("notification for change %d, persistent search still registered: %s, " |
| | | + "notifications received afterwards: %s", |
| | | i, isPersistentSearchRegistered(backend, "ou=psearch,o=test"), notified) |
| | | .isEqualTo("cn=changing,ou=psearch,o=test change " + i); |
| | | } |
| | | } |
| | | } |
| | | |
| | | // A persistent search notification is not a search result: it must be reported whether or not |
| | | // the entry was returned before, and it is not bound by the size and time limits of the search. |
| | | // In a real persistent search the search phase is only open for a few instructions after the |
| | | // search is registered with the backend, so the notification path is driven directly here. |
| | | @Test |
| | | public void test_persistent_search_notification_ignores_search_phase_state() throws Exception { |
| | | TestCaseUtils.addEntries( |
| | | "dn: ou=psearch-notify,o=test", |
| | | "objectClass: top", |
| | | "objectClass: organizationalUnit", |
| | | "ou: psearch-notify", |
| | | "" |
| | | ); |
| | | final Entry entry = DirectoryServer.getEntry(DN.valueOf("ou=psearch-notify,o=test")); |
| | | |
| | | final InternalSearchOperation search = new InternalSearchOperation( |
| | | InternalClientConnection.getRootConnection(), |
| | | InternalClientConnection.nextOperationID(), |
| | | InternalClientConnection.nextMessageID(), |
| | | org.opends.server.protocols.internal.Requests |
| | | .newSearchRequest(DN.valueOf("o=test"), SearchScope.WHOLE_SUBTREE) |
| | | .setDereferenceAliasesPolicy(DereferenceAliasesPolicy.ALWAYS)); |
| | | |
| | | // The search phase returns the entry once, and drops it when it reaches it a second time |
| | | // through an alias ... |
| | | assertThat(search.returnEntry(entry, null)).isTrue(); |
| | | assertThat(search.returnEntry(entry, null)).isTrue(); |
| | | assertThat(search.getSearchEntries()).hasSize(1); |
| | | |
| | | // ... but a change reported to a persistent search is never a duplicate, whether the search |
| | | // phase is still open (the entry was just returned by it) or already over. |
| | | assertThat(search.returnPersistentSearchEntry(entry, null)).isTrue(); |
| | | search.endSearchPhase(); |
| | | assertThat(search.returnPersistentSearchEntry(entry, null)).isTrue(); |
| | | assertThat(search.getSearchEntries()).hasSize(3); |
| | | |
| | | // The size limit of the search does not bound a notification: it only bounds the search |
| | | // phase, and is lifted for the rest of a persistent search once that phase is over. |
| | | search.setSizeLimit(1); |
| | | assertThat(search.returnPersistentSearchEntry(entry, null)).isTrue(); |
| | | assertThat(search.getSearchEntries()).hasSize(4); |
| | | // The search phase itself is still bound by it. |
| | | assertThat(search.returnEntry(entry, null)).isFalse(); |
| | | assertThat(search.getResultCode()).isEqualTo(ResultCode.SIZE_LIMIT_EXCEEDED); |
| | | |
| | | // Same for the time limit, checked on its own: with a size limit left in the way the search |
| | | // phase would stop on that one and the time limit would never be reached. |
| | | search.setSizeLimit(0); |
| | | search.setTimeLimit(1); |
| | | search.setTimeLimitExpiration(0); |
| | | assertThat(search.returnPersistentSearchEntry(entry, null)).isTrue(); |
| | | assertThat(search.getSearchEntries()).hasSize(5); |
| | | assertThat(search.returnEntry(entry, null)).isFalse(); |
| | | assertThat(search.getResultCode()).isEqualTo(ResultCode.TIME_LIMIT_EXCEEDED); |
| | | } |
| | | |
| | | // A persistent search which is cancelled because its backend goes away, as happens when the |
| | | // backend is disabled or re-initialized, must be told so: without a search result done the |
| | | // client waits forever for changes on a search which no longer exists. |
| | | @Test |
| | | public void test_persistent_search_is_told_when_its_backend_goes_away() throws Exception { |
| | | final String backendID = "psearchUnavailable"; |
| | | final String baseDN = "o=psearch-unavailable"; |
| | | TestCaseUtils.initializeMemoryBackend(backendID, baseDN, true); |
| | | final MemoryBackend backend = (MemoryBackend) TestCaseUtils.getServerContext() |
| | | .getBackendConfigManager().getLocalBackendById(backendID); |
| | | |
| | | final SearchRequest request = |
| | | Requests.newSearchRequest(baseDN, SearchScope.WHOLE_SUBTREE, "(objectclass=*)") |
| | | .addControl(PersistentSearchRequestControl.newControl( |
| | | true, true, false, PersistentSearchChangeType.MODIFY)); |
| | | |
| | | final LDAPConnectionFactory factory = |
| | | new LDAPConnectionFactory("localhost", TestCaseUtils.getServerLdapPort()); |
| | | try (Connection psearch = factory.getConnection()) { |
| | | psearch.bind("cn=Directory Manager", "password".toCharArray()); |
| | | final LdapPromise<Result> searchDone = psearch.searchAsync(request, new SearchResultHandler() { |
| | | @Override |
| | | public boolean handleEntry(SearchResultEntry entry) { |
| | | return true; |
| | | } |
| | | |
| | | @Override |
| | | public boolean handleReference(SearchResultReference reference) { |
| | | return true; |
| | | } |
| | | }); |
| | | |
| | | for (int i = 0; !isPersistentSearchRegistered(backend, baseDN) && i < 500; i++) { |
| | | Thread.sleep(10); |
| | | } |
| | | assertThat(isPersistentSearchRegistered(backend, baseDN)) |
| | | .as("the persistent search was never registered with backend %s", backendID) |
| | | .isTrue(); |
| | | |
| | | backend.finalizeBackend(); |
| | | |
| | | try { |
| | | final Result result = searchDone.getOrThrow(30, TimeUnit.SECONDS); |
| | | fail("the persistent search should have been terminated, it returned " + result); |
| | | } catch (LdapException e) { |
| | | assertThat(e.getResult().getResultCode()).isEqualTo(ResultCode.UNAVAILABLE); |
| | | assertThat(e.getResult().getDiagnosticMessage()).contains(backendID); |
| | | } |
| | | } finally { |
| | | TestCaseUtils.getServerContext().getBackendConfigManager().deregisterLocalBackend(backend); |
| | | } |
| | | } |
| | | |
| | | /** Whether the provided backend has a persistent search registered for the provided base DN. */ |
| | | private static boolean isPersistentSearchRegistered(LocalBackend<?> backend, String baseDN) { |
| | | for (PersistentSearch psearch : backend.getPersistentSearches()) { |
| | | if (psearch.getSearchOperation().getBaseDN().equals(DN.valueOf(baseDN))) { |
| | | return true; |
| | | } |
| | | } |
| | | return false; |
| | | } |
| | | |
| | | // An alias is dereferenced before its target is reached on its own: the target must still be |
| | | // returned, and exactly once. The original regression was order-sensitive, dropping the target |
| | | // when the alias reached it first, so this pins the alias-before-target order specifically. The |