From 98e0d11fbcb7ed122347270c8ed9d260250fbf23 Mon Sep 17 00:00:00 2001 From: Valera V Harseko Date: Thu, 1 Oct 2026 12:49:16 +0300 Subject: [PATCH] [#1143] Do not report a Notice of Disconnection that cannot reach an already-closed client as an unhandled error ClientConnectionImpl.disconnect(ResultCode, String) subscribed to the notification with no error consumer. When the client had closed its end first, the failed write went to RxJavaPlugins.onError: the server printed an OnErrorNotImplementedException stack trace, and on a worker thread logged ERR_UNCAUGHT_THREAD_EXCEPTION and raised an uncaught-exception alert. Subscribe with an error consumer that only traces the failure. The connection is still closed on both outcomes by doAfterTerminate. Remove the try/catch (OnErrorNotImplementedException) around s.onError() in sendUnsolicitedNotification(): RxJavaPlugins.onError does not throw, so the catch could never see the exception. Fixes #1143 --- .../opendj/grizzly/LDAPServerFilter.java | 24 ++++--- .../grizzly/ConnectionFactoryTestCase.java | 62 +++++++++++++++++++ 2 files changed, 77 insertions(+), 9 deletions(-) diff --git a/opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/LDAPServerFilter.java b/opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/LDAPServerFilter.java index a1a5cc5113..12ff173121 100644 --- a/opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/LDAPServerFilter.java +++ b/opendj-grizzly/src/main/java/org/forgerock/opendj/grizzly/LDAPServerFilter.java @@ -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; /** @@ -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() { + @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) { @@ -645,13 +657,7 @@ public void handleResult(Boolean result) { }).thenOnException(new ExceptionHandler() { @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 diff --git a/opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/ConnectionFactoryTestCase.java b/opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/ConnectionFactoryTestCase.java index 87a652c48a..c5e5499ef0 100644 --- a/opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/ConnectionFactoryTestCase.java +++ b/opendj-grizzly/src/test/java/org/forgerock/opendj/grizzly/ConnectionFactoryTestCase.java @@ -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; @@ -95,6 +97,8 @@ import com.forgerock.reactive.ServerConnectionFactoryAdapter; +import io.reactivex.rxjava3.plugins.RxJavaPlugins; + /** * Tests the {@code ConnectionFactory} classes. */ @@ -665,6 +669,64 @@ public ServerConnection 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 contextHolder = new AtomicReference<>(); + final ServerConnectionFactory mockServer = + mock(ServerConnectionFactory.class); + when(mockServer.handleAccept(any(LDAPClientContext.class))).thenAnswer( + new Answer>() { + @Override + public ServerConnection answer(InvocationOnMock invocation) throws Throwable { + contextHolder.set((LDAPClientContext) invocation.getArguments()[0]); + connectLatch.countDown(); + return mock(ServerConnection.class); + } + }); + + final List unhandledErrors = new CopyOnWriteArrayList<>(); + final io.reactivex.rxjava3.functions.Consumer previousErrorHandler = + RxJavaPlugins.getErrorHandler(); + RxJavaPlugins.setErrorHandler(new io.reactivex.rxjava3.functions.Consumer() { + @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() { + @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 {