| | |
| | | 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; |
| | |
| | | |
| | | import com.forgerock.reactive.ServerConnectionFactoryAdapter; |
| | | |
| | | import io.reactivex.rxjava3.plugins.RxJavaPlugins; |
| | | |
| | | /** |
| | | * Tests the {@code ConnectionFactory} classes. |
| | | */ |
| | |
| | | } |
| | | } |
| | | |
| | | /** |
| | | * 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 { |