Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -82,10 +82,10 @@

import com.forgerock.reactive.Action;
import com.forgerock.reactive.Completable;
import com.forgerock.reactive.Consumer;
import com.forgerock.reactive.ReactiveHandler;
import com.forgerock.reactive.Stream;

import io.reactivex.rxjava3.exceptions.OnErrorNotImplementedException;
import org.openidentityplatform.rxjava3.internal.util.BackpressureHelper;

/**
Expand Down Expand Up @@ -542,7 +542,19 @@ public void run() throws Exception {
// handleClose() will be invoked once this connection has been closed.
connection.closeSilently();
}
}).subscribe();
}).subscribe(new Action() {
@Override
public void run() throws Exception {
// Nothing to do: the connection is closed on either outcome.
}
}, new Consumer<Throwable>() {
@Override
public void accept(final Throwable error) throws Exception {
// Expected when the client has already closed its end: the notice cannot be delivered,
// and the connection is closed anyway.
logger.traceException(error);
}
});
}

private void notifyConnectionClosedRawUnbind(final LdapRequestEnvelope rawUnbindRequest) {
Expand Down Expand Up @@ -645,13 +657,7 @@ public void handleResult(Boolean result) {
}).thenOnException(new ExceptionHandler<Exception>() {
@Override
public void handleException(Exception exception) {
try {
s.onError(exception);
} catch (Throwable t) {
if (!(t instanceof OnErrorNotImplementedException)) {
throw t;
}
}
s.onError(exception);
}
}).thenOnRuntimeException(new RuntimeExceptionHandler() {
@Override
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,7 +37,9 @@
import java.net.InetSocketAddress;
import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.Callable;
import java.util.concurrent.CopyOnWriteArrayList;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
Expand Down Expand Up @@ -95,6 +97,8 @@

import com.forgerock.reactive.ServerConnectionFactoryAdapter;

import io.reactivex.rxjava3.plugins.RxJavaPlugins;

/**
* Tests the {@code ConnectionFactory} classes.
*/
Expand Down Expand Up @@ -665,6 +669,64 @@ public ServerConnection<Integer> answer(InvocationOnMock invocation)
}
}

/**
* A Notice of Disconnection that cannot be written because the client has already closed its end is an
* expected outcome: it must not be reported as an unhandled error (issue #1143).
*/
@SuppressWarnings("unchecked")
@Test
public void testDisconnectWithNotificationToClosedClientIsNotReportedAsUnhandledError() throws Exception {
final CountDownLatch connectLatch = new CountDownLatch(1);
final AtomicReference<LDAPClientContext> contextHolder = new AtomicReference<>();
final ServerConnectionFactory<LDAPClientContext, Integer> mockServer =
mock(ServerConnectionFactory.class);
when(mockServer.handleAccept(any(LDAPClientContext.class))).thenAnswer(
new Answer<ServerConnection<Integer>>() {
@Override
public ServerConnection<Integer> answer(InvocationOnMock invocation) throws Throwable {
contextHolder.set((LDAPClientContext) invocation.getArguments()[0]);
connectLatch.countDown();
return mock(ServerConnection.class);
}
});

final List<Throwable> unhandledErrors = new CopyOnWriteArrayList<>();
final io.reactivex.rxjava3.functions.Consumer<? super Throwable> previousErrorHandler =
RxJavaPlugins.getErrorHandler();
RxJavaPlugins.setErrorHandler(new io.reactivex.rxjava3.functions.Consumer<Throwable>() {
@Override
public void accept(Throwable error) {
unhandledErrors.add(error);
}
});
LDAPListener listener = new LDAPListener(Collections.singleton(loopbackWithDynamicPort()),
new ServerConnectionFactoryAdapter(Options.defaultOptions().get(LDAP_DECODE_OPTIONS), mockServer));
try {
final InetSocketAddress listenerAddr = listener.getSocketAddresses().iterator().next();
final Connection client = new LDAPConnectionFactory(listenerAddr.getHostName(),
listenerAddr.getPort()).getConnection();
assertThat(connectLatch.await(TEST_TIMEOUT, TimeUnit.SECONDS)).isTrue();
final LDAPClientContext context = contextHolder.get();

// The client leaves first: wait until the server has seen the connection close.
client.close();
waitForCondition(new Callable<Boolean>() {
@Override
public Boolean call() throws Exception {
return context.isClosed();
}
});

// Writing the notice now fails, and does so before disconnect() returns.
context.disconnect(ResultCode.BUSY, "busy");

assertThat(unhandledErrors).isEmpty();
} finally {
RxJavaPlugins.setErrorHandler(previousErrorHandler);
listener.close();
}
}

@Test(description = "Test for OPENDJ-1121: Closing a connection after "
+ "closing the connection factory causes NPE")
public void testFactoryCloseBeforeConnectionClose() throws Exception {
Expand Down
Loading