mirror of https://github.com/OpenIdentityPlatform/OpenDJ.git

Valery Kharseko
yesterday 2edf7d0363ce6b8d32efbe69c330e5f0e66b31cf
Notify connection event listeners when an internal connection is closed (#787)
4 files modified
1 files added
369 ■■■■■ changed files
opendj-core/src/main/java/org/forgerock/opendj/ldap/Connection.java 7 ●●●●● patch | view | raw | blame | history
opendj-core/src/main/java/org/forgerock/opendj/ldap/InternalConnection.java 47 ●●●● patch | view | raw | blame | history
opendj-core/src/main/java/org/forgerock/opendj/ldap/spi/ConnectionState.java 21 ●●●● patch | view | raw | blame | history
opendj-core/src/test/java/org/forgerock/opendj/ldap/InternalConnectionTestCase.java 228 ●●●●● patch | view | raw | blame | history
opendj-core/src/test/java/org/forgerock/opendj/ldap/spi/ConnectionStateTest.java 66 ●●●●● patch | view | raw | blame | history
opendj-core/src/main/java/org/forgerock/opendj/ldap/Connection.java
@@ -280,13 +280,14 @@
     * Registers the provided connection event listener so that it will be
     * notified when this connection is closed by the application, receives an
     * unsolicited notification, or experiences a fatal error.
     * <p>
     * A listener registered once this connection has already failed and/or been
     * closed is notified of the events it has missed before this method
     * returns.
     *
     * @param listener
     *            The listener which wants to be notified when events occur on
     *            this connection.
     * @throws IllegalStateException
     *             If this connection has already been closed, i.e. if
     *             {@code isClosed() == true}.
     * @throws NullPointerException
     *             If the {@code listener} was {@code null}.
     */
opendj-core/src/main/java/org/forgerock/opendj/ldap/InternalConnection.java
@@ -13,11 +13,11 @@
 *
 * Copyright 2010 Sun Microsystems, Inc.
 * Portions copyright 2011-2016 ForgeRock AS.
 * Portions Copyright 2026 3A Systems, LLC.
 */
package org.forgerock.opendj.ldap;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import org.forgerock.opendj.ldap.requests.AbandonRequest;
@@ -35,6 +35,7 @@
import org.forgerock.opendj.ldap.responses.ExtendedResult;
import org.forgerock.opendj.ldap.responses.Result;
import org.forgerock.opendj.ldap.spi.BindResultLdapPromiseImpl;
import org.forgerock.opendj.ldap.spi.ConnectionState;
import org.forgerock.opendj.ldap.spi.ExtendedResultLdapPromiseImpl;
import org.forgerock.opendj.ldap.spi.ResultLdapPromiseImpl;
import org.forgerock.opendj.ldap.spi.SearchResultLdapPromiseImpl;
@@ -49,7 +50,20 @@
 */
final class InternalConnection extends AbstractAsynchronousConnection {
    private final ServerConnection<Integer> serverConnection;
    private final List<ConnectionEventListener> listeners = new CopyOnWriteArrayList<>();
    /**
     * Tracks whether this connection has been closed and notifies the registered
     * connection event listeners when it is.
     * <p>
     * An internal connection has no transport, so it can never fail nor receive
     * an unsolicited notification: closure by the application is the only event
     * its listeners can observe.
     */
    private final ConnectionState state = new ConnectionState();
    /**
     * Guards the hand-over to the server connection so that it happens exactly
     * once, even if notifying the listeners fails.
     */
    private final AtomicBoolean isClosed = new AtomicBoolean();
    private final AtomicInteger messageID = new AtomicInteger();
    /**
@@ -82,7 +96,7 @@
    @Override
    public void addConnectionEventListener(final ConnectionEventListener listener) {
        Reject.ifNull(listener);
        listeners.add(listener);
        state.addConnectionEventListener(listener);
    }
    @Override
@@ -97,8 +111,21 @@
    @Override
    public void close(final UnbindRequest request, final String reason) {
        final int i = messageID.getAndIncrement();
        serverConnection.handleConnectionClosed(i, request);
        // Closing an already closed connection has no effect.
        if (isClosed.compareAndSet(false, true)) {
            try {
                state.notifyConnectionClosed();
            } finally {
                /*
                 * A listener throwing an exception must not prevent the server
                 * connection from releasing its resources: for an internal
                 * connection this is the only cleanup there is, and a second
                 * close() would be a no-op.
                 */
                final int i = messageID.getAndIncrement();
                serverConnection.handleConnectionClosed(i, request);
            }
        }
    }
    @Override
@@ -133,14 +160,12 @@
    @Override
    public boolean isClosed() {
        // FIXME: this should be true after close has been called.
        return false;
        return state.isClosed();
    }
    @Override
    public boolean isValid() {
        // FIXME: this should be false if this connection is disconnected.
        return true;
        return state.isValid();
    }
    @Override
@@ -166,7 +191,7 @@
    @Override
    public void removeConnectionEventListener(final ConnectionEventListener listener) {
        Reject.ifNull(listener);
        listeners.remove(listener);
        state.removeConnectionEventListener(listener);
    }
    @Override
opendj-core/src/main/java/org/forgerock/opendj/ldap/spi/ConnectionState.java
@@ -12,11 +12,12 @@
 * information: "Portions Copyright [year] [name of copyright owner]".
 *
 * Copyright 2013-2016 ForgeRock AS.
 * Portions Copyright 2026 3A Systems, LLC.
 */
package org.forgerock.opendj.ldap.spi;
import java.util.LinkedList;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import org.forgerock.opendj.ldap.ConnectionEventListener;
import org.forgerock.opendj.ldap.LdapException;
@@ -245,8 +246,15 @@
    /** Whether the connection has failed due to a disconnect notification. */
    private boolean failedDueToDisconnect;
    /** Registered event listeners. */
    private final List<ConnectionEventListener> listeners = new LinkedList<>();
    /**
     * Registered event listeners. Copy on write because a listener is allowed
     * to register or, more commonly, de-register itself while it is being
     * notified: the notification methods are re-entrant on this object's
     * monitor, so any other list implementation would have the notification
     * loop either skip the remaining listeners or fail with a
     * {@link java.util.ConcurrentModificationException}.
     */
    private final List<ConnectionEventListener> listeners = new CopyOnWriteArrayList<>();
    /** Internal state implementation. */
    private volatile State state = State.VALID;
@@ -260,13 +268,14 @@
     * Registers the provided connection event listener so that it will be
     * notified when this connection is closed by the application, receives an
     * unsolicited notification, or experiences a fatal error.
     * <p>
     * A listener registered once the connection has already failed and/or been
     * closed is notified of the events it has missed before this method
     * returns.
     *
     * @param listener
     *            The listener which wants to be notified when events occur on
     *            this connection.
     * @throws IllegalStateException
     *             If this connection has already been closed, i.e. if
     *             {@code isClosed() == true}.
     * @throws NullPointerException
     *             If the {@code listener} was {@code null}.
     */
opendj-core/src/test/java/org/forgerock/opendj/ldap/InternalConnectionTestCase.java
New file
@@ -0,0 +1,228 @@
/*
 * 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 2026 3A Systems, LLC.
 */
package org.forgerock.opendj.ldap;
import static org.fest.assertions.Assertions.assertThat;
import static org.mockito.Matchers.any;
import static org.mockito.Matchers.anyInt;
import static org.mockito.Matchers.same;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
import static org.testng.Assert.fail;
import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
import java.util.concurrent.Future;
import java.util.concurrent.TimeUnit;
import org.forgerock.opendj.ldap.requests.Requests;
import org.forgerock.opendj.ldap.requests.UnbindRequest;
import org.mockito.invocation.InvocationOnMock;
import org.mockito.stubbing.Answer;
import org.testng.annotations.Test;
/** Tests the connection life cycle of the pseudo-connection returned by {@link Connections#newInternalConnection}. */
@SuppressWarnings("javadoc")
public class InternalConnectionTestCase extends SdkTestCase {
    @SuppressWarnings("unchecked")
    private static ServerConnection<Integer> mockServerConnection() {
        return mock(ServerConnection.class);
    }
    /** Asserts that the close was forwarded exactly once, carrying the very request the caller provided. */
    private static void verifyClosedOnce(final ServerConnection<Integer> serverConnection,
            final UnbindRequest request) {
        verify(serverConnection, times(1)).handleConnectionClosed(anyInt(), same(request));
    }
    /** Asserts that the close was forwarded exactly once, for callers which let {@code close()} build the request. */
    private static void verifyClosedOnce(final ServerConnection<Integer> serverConnection) {
        verify(serverConnection, times(1)).handleConnectionClosed(anyInt(), any(UnbindRequest.class));
    }
    @Test
    public void testCloseNotifiesListeners() throws Exception {
        final ServerConnection<Integer> serverConnection = mockServerConnection();
        final Connection connection = Connections.newInternalConnection(serverConnection);
        final MockConnectionEventListener listener = new MockConnectionEventListener();
        connection.addConnectionEventListener(listener);
        final UnbindRequest request = Requests.newUnbindRequest();
        connection.close(request, null);
        assertThat(listener.getInvocationCount()).isEqualTo(1);
        verifyClosedOnce(serverConnection, request);
    }
    @Test
    public void testCloseIsIdempotent() throws Exception {
        final ServerConnection<Integer> serverConnection = mockServerConnection();
        final Connection connection = Connections.newInternalConnection(serverConnection);
        final MockConnectionEventListener listener = new MockConnectionEventListener();
        connection.addConnectionEventListener(listener);
        final UnbindRequest request = Requests.newUnbindRequest();
        connection.close(request, null);
        connection.close(request, null);
        connection.close();
        assertThat(listener.getInvocationCount()).isEqualTo(1);
        verifyClosedOnce(serverConnection, request);
        verifyClosedOnce(serverConnection);
    }
    @Test
    public void testRemovedListenerIsNotNotified() throws Exception {
        final ServerConnection<Integer> serverConnection = mockServerConnection();
        final Connection connection = Connections.newInternalConnection(serverConnection);
        final MockConnectionEventListener listener = new MockConnectionEventListener();
        connection.addConnectionEventListener(listener);
        connection.removeConnectionEventListener(listener);
        connection.close();
        assertThat(listener.getInvocationCount()).isEqualTo(0);
        verifyClosedOnce(serverConnection);
    }
    /** A listener registered after the connection was closed is told about it straight away. */
    @Test
    public void testListenerAddedAfterCloseIsNotifiedImmediately() throws Exception {
        final Connection connection = Connections.newInternalConnection(mockServerConnection());
        connection.close();
        final MockConnectionEventListener listener = new MockConnectionEventListener();
        connection.addConnectionEventListener(listener);
        assertThat(listener.getInvocationCount()).isEqualTo(1);
    }
    @Test
    public void testConnectionStateReflectsClose() throws Exception {
        final Connection connection = Connections.newInternalConnection(mockServerConnection());
        assertThat(connection.isClosed()).isFalse();
        assertThat(connection.isValid()).isTrue();
        connection.close();
        assertThat(connection.isClosed()).isTrue();
        assertThat(connection.isValid()).isFalse();
    }
    /**
     * De-registering a listener from within its own close notification is a common idiom, and it must neither skip
     * the remaining listeners nor break the hand-over to the server connection.
     */
    @Test
    public void testSelfDeregisteringListenerDoesNotStopNotification() throws Exception {
        final ServerConnection<Integer> serverConnection = mockServerConnection();
        final Connection connection = Connections.newInternalConnection(serverConnection);
        final ConnectionEventListener selfDeregistering = mock(ConnectionEventListener.class);
        doAnswer(new Answer<Void>() {
            @Override
            public Void answer(final InvocationOnMock invocation) {
                connection.removeConnectionEventListener(selfDeregistering);
                return null;
            }
        }).when(selfDeregistering).handleConnectionClosed();
        connection.addConnectionEventListener(selfDeregistering);
        // Three more listeners: with a list which does not tolerate mutation while being iterated, two of them
        // would go unnotified and the third would raise a ConcurrentModificationException out of close().
        final List<MockConnectionEventListener> listeners = new ArrayList<>();
        for (int i = 0; i < 3; i++) {
            final MockConnectionEventListener listener = new MockConnectionEventListener();
            listeners.add(listener);
            connection.addConnectionEventListener(listener);
        }
        final UnbindRequest request = Requests.newUnbindRequest();
        connection.close(request, null);
        verify(selfDeregistering, times(1)).handleConnectionClosed();
        for (final MockConnectionEventListener listener : listeners) {
            assertThat(listener.getInvocationCount()).isEqualTo(1);
        }
        verifyClosedOnce(serverConnection, request);
    }
    /**
     * A misbehaving listener must not prevent the server connection from releasing its resources: the forward is the
     * only cleanup an internal connection performs, and a second {@code close()} would be a no-op.
     */
    @Test
    public void testListenerFailureDoesNotPreventServerClose() throws Exception {
        final ServerConnection<Integer> serverConnection = mockServerConnection();
        final Connection connection = Connections.newInternalConnection(serverConnection);
        final ConnectionEventListener listener = mock(ConnectionEventListener.class);
        doThrow(new IllegalStateException("listener failure")).when(listener).handleConnectionClosed();
        connection.addConnectionEventListener(listener);
        final UnbindRequest request = Requests.newUnbindRequest();
        try {
            connection.close(request, null);
            fail("close() should have passed on the listener failure");
        } catch (final IllegalStateException expected) {
            // Expected: no notification site in the SDK guards individual listener callbacks.
        }
        verifyClosedOnce(serverConnection, request);
        assertThat(connection.isClosed()).isTrue();
    }
    /** Concurrent closes notify the listeners and reach the server connection exactly once. */
    @Test
    public void testConcurrentCloseNotifiesAndForwardsOnce() throws Exception {
        final int threadCount = 16;
        final ServerConnection<Integer> serverConnection = mockServerConnection();
        final Connection connection = Connections.newInternalConnection(serverConnection);
        final MockConnectionEventListener listener = new MockConnectionEventListener();
        connection.addConnectionEventListener(listener);
        final CountDownLatch startLatch = new CountDownLatch(1);
        final ExecutorService executor = Executors.newFixedThreadPool(threadCount);
        try {
            final List<Future<Void>> futures = new ArrayList<>(threadCount);
            for (int i = 0; i < threadCount; i++) {
                futures.add(executor.submit(new Callable<Void>() {
                    @Override
                    public Void call() throws Exception {
                        startLatch.await();
                        connection.close();
                        return null;
                    }
                }));
            }
            startLatch.countDown();
            for (final Future<Void> future : futures) {
                future.get(30, TimeUnit.SECONDS);
            }
        } finally {
            executor.shutdownNow();
        }
        assertThat(listener.getInvocationCount()).isEqualTo(1);
        verifyClosedOnce(serverConnection);
    }
}
opendj-core/src/test/java/org/forgerock/opendj/ldap/spi/ConnectionStateTest.java
@@ -12,6 +12,7 @@
 * information: "Portions Copyright [year] [name of copyright owner]".
 *
 * Copyright 2013-2016 ForgeRock AS.
 * Portions Copyright 2026 3A Systems, LLC.
 */
package org.forgerock.opendj.ldap.spi;
@@ -243,6 +244,71 @@
    }
    /**
     * Tests that a listener de-registering itself while being notified of the close, a common idiom, neither skips
     * the remaining listeners nor breaks out of the notification loop.
     */
    @Test
    public void testSelfDeregisteringListenerDuringClose() {
        final ConnectionState state = new ConnectionState();
        final ConnectionEventListener selfDeregistering = mock(ConnectionEventListener.class);
        doAnswer(new Answer<Void>() {
            @Override
            public Void answer(final InvocationOnMock invocation) {
                state.removeConnectionEventListener(selfDeregistering);
                return null;
            }
        }).when(selfDeregistering).handleConnectionClosed();
        final ConnectionEventListener listener2 = mock(ConnectionEventListener.class);
        final ConnectionEventListener listener3 = mock(ConnectionEventListener.class);
        final ConnectionEventListener listener4 = mock(ConnectionEventListener.class);
        state.addConnectionEventListener(selfDeregistering);
        state.addConnectionEventListener(listener2);
        state.addConnectionEventListener(listener3);
        state.addConnectionEventListener(listener4);
        assertThat(state.notifyConnectionClosed()).isTrue();
        assertThat(state.isClosed()).isTrue();
        verify(selfDeregistering).handleConnectionClosed();
        verify(listener2).handleConnectionClosed();
        verify(listener3).handleConnectionClosed();
        verify(listener4).handleConnectionClosed();
    }
    /**
     * Tests that a listener de-registering itself while being notified of an error does not prevent the remaining
     * listeners from being notified.
     */
    @Test
    public void testSelfDeregisteringListenerDuringError() {
        final ConnectionState state = new ConnectionState();
        final ConnectionEventListener selfDeregistering = mock(ConnectionEventListener.class);
        doAnswer(new Answer<Void>() {
            @Override
            public Void answer(final InvocationOnMock invocation) {
                state.removeConnectionEventListener(selfDeregistering);
                return null;
            }
        }).when(selfDeregistering).handleConnectionError(false, ERROR);
        final ConnectionEventListener listener2 = mock(ConnectionEventListener.class);
        final ConnectionEventListener listener3 = mock(ConnectionEventListener.class);
        final ConnectionEventListener listener4 = mock(ConnectionEventListener.class);
        state.addConnectionEventListener(selfDeregistering);
        state.addConnectionEventListener(listener2);
        state.addConnectionEventListener(listener3);
        state.addConnectionEventListener(listener4);
        assertThat(state.notifyConnectionError(false, ERROR)).isTrue();
        verify(selfDeregistering).handleConnectionError(false, ERROR);
        verify(listener2).handleConnectionError(false, ERROR);
        verify(listener3).handleConnectionError(false, ERROR);
        verify(listener4).handleConnectionError(false, ERROR);
    }
    /**
     * Tests that reentrant close from error listener is handled.
     */
    @Test