From 455cd468a59ccfeb4d58fefd364fc33b57ace35e Mon Sep 17 00:00:00 2001 From: AdzerKI Date: Fri, 2 Oct 2026 05:27:11 +0300 Subject: [PATCH] Complete I/O reactor termination when the selector has been closed by another thread closeOpenChannels threw ClosedSelectorException once close(CloseMode) from another thread had closed the selector, so doTerminate reported it through the exception callback and skipped processClosedSessions: event handlers of the closed sessions were never notified of the disconnect. --- .../hc/core5/reactor/SingleCoreIOReactor.java | 9 ++++- .../reactor/TestSingleCoreIOReactor.java | 40 +++++++++++++++++++ 2 files changed, 48 insertions(+), 1 deletion(-) diff --git a/httpcore5/src/main/java/org/apache/hc/core5/reactor/SingleCoreIOReactor.java b/httpcore5/src/main/java/org/apache/hc/core5/reactor/SingleCoreIOReactor.java index e5839b3e8..462f0e0df 100644 --- a/httpcore5/src/main/java/org/apache/hc/core5/reactor/SingleCoreIOReactor.java +++ b/httpcore5/src/main/java/org/apache/hc/core5/reactor/SingleCoreIOReactor.java @@ -37,6 +37,7 @@ import java.net.UnknownHostException; import java.nio.channels.CancelledKeyException; import java.nio.channels.ClosedChannelException; +import java.nio.channels.ClosedSelectorException; import java.nio.channels.SelectionKey; import java.nio.channels.SocketChannel; import java.util.Queue; @@ -434,7 +435,13 @@ private void processConnectionRequest(final SocketChannel socketChannel, final I } private void closeOpenChannels() { - for (final SelectionKey key : selector.keys()) { + final Set keys; + try { + keys = selector.keys(); + } catch (final ClosedSelectorException ignore) { + return; + } + for (final SelectionKey key : keys) { final Object attachment = key.attachment(); if (attachment instanceof InternalChannel) { final InternalChannel channel = (InternalChannel) attachment; diff --git a/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestSingleCoreIOReactor.java b/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestSingleCoreIOReactor.java index 4723f93c6..8b2e2ad97 100644 --- a/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestSingleCoreIOReactor.java +++ b/httpcore5/src/test/java/org/apache/hc/core5/reactor/TestSingleCoreIOReactor.java @@ -26,10 +26,17 @@ */ package org.apache.hc.core5.reactor; +import java.net.InetAddress; import java.net.InetSocketAddress; +import java.nio.channels.ServerSocketChannel; import java.nio.channels.SocketChannel; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; import org.apache.hc.core5.function.Decorator; +import org.apache.hc.core5.io.CloseMode; import org.apache.hc.core5.net.NamedEndpoint; import org.apache.hc.core5.util.Timeout; import org.junit.jupiter.api.Assertions; @@ -91,4 +98,37 @@ void enqueueAfterShutdownThrows() throws Exception { } } + @Test + void terminatesWhenSelectorClosedByAnotherThread() throws Exception { + final CountDownLatch connected = new CountDownLatch(1); + final CountDownLatch reactorClosed = new CountDownLatch(1); + final IOEventHandler handler = Mockito.mock(IOEventHandler.class); + Mockito.doAnswer(invocation -> { + connected.countDown(); + reactorClosed.await(); + return null; + }).when(handler).connected(Mockito.any()); + final List exceptions = new CopyOnWriteArrayList<>(); + try (ServerSocketChannel server = ServerSocketChannel.open()) { + server.bind(new InetSocketAddress(InetAddress.getLoopbackAddress(), 0)); + try (SocketChannel client = SocketChannel.open(server.getLocalAddress()); + SocketChannel channel = server.accept()) { + final SingleCoreIOReactor reactor = new SingleCoreIOReactor( + exceptions::add, (session, attachment) -> handler, IOReactorConfig.DEFAULT, null, null, null, null); + final Thread worker = new Thread(reactor::execute); + worker.start(); + reactor.enqueueChannel(new ChannelEntry(channel, null)); + Assertions.assertTrue(connected.await(5, TimeUnit.SECONDS)); + + reactor.close(CloseMode.IMMEDIATE); + reactorClosed.countDown(); + worker.join(TimeUnit.SECONDS.toMillis(5)); + + Assertions.assertFalse(worker.isAlive()); + Assertions.assertTrue(exceptions.isEmpty(), exceptions::toString); + Mockito.verify(handler).disconnected(Mockito.any()); + } + } + } + }