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()); + } + } + } + }