This is an automated email from the ASF dual-hosted git repository. dakirily pushed a commit to branch main in repository https://gitbox.apache.org/repos/asf/qpid-broker-j.git
commit ede64d973d8e99d2b876ba9f733b682ffde99f93 Author: Daniil Kirilyuk <[email protected]> AuthorDate: Sat Sep 19 19:45:59 2026 +0200 NO-JIRA: [Broker-J] Network selector resilience improvements --- .../transport/NetworkConnectionScheduler.java | 56 +++++++- .../qpid/server/transport/SelectorThread.java | 5 + .../transport/NetworkConnectionSchedulerTest.java | 148 +++++++++++++++++++++ 3 files changed, 202 insertions(+), 7 deletions(-) diff --git a/broker-core/src/main/java/org/apache/qpid/server/transport/NetworkConnectionScheduler.java b/broker-core/src/main/java/org/apache/qpid/server/transport/NetworkConnectionScheduler.java index 6d0cc37449..1576a53088 100644 --- a/broker-core/src/main/java/org/apache/qpid/server/transport/NetworkConnectionScheduler.java +++ b/broker-core/src/main/java/org/apache/qpid/server/transport/NetworkConnectionScheduler.java @@ -22,8 +22,10 @@ package org.apache.qpid.server.transport; import java.io.IOException; import java.nio.channels.ServerSocketChannel; +import java.util.concurrent.BlockingQueue; import java.util.concurrent.Executors; import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.RejectedExecutionException; import java.util.concurrent.ThreadFactory; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; @@ -102,14 +104,14 @@ public class NetworkConnectionScheduler final int corePoolSize = _poolSize; final int maximumPoolSize = _poolSize; final long keepAliveTime = _threadKeepAliveTimeout; - final java.util.concurrent.BlockingQueue<Runnable> workQueue = new LinkedBlockingQueue<>(); + final BlockingQueue<Runnable> workQueue = new LinkedBlockingQueue<>(); final ThreadFactory factory = _factory; - _executor = new ThreadPoolExecutor(corePoolSize, - maximumPoolSize, - keepAliveTime, - TimeUnit.MINUTES, - workQueue, - QpidByteBuffer.createQpidByteBufferTrackingThreadFactory(factory)); + _executor = new SelectorThreadPoolExecutor(corePoolSize, + maximumPoolSize, + keepAliveTime, + workQueue, + QpidByteBuffer.createQpidByteBufferTrackingThreadFactory( + factory)); _executor.prestartAllCoreThreads(); _executor.allowCoreThreadTimeOut(true); for(int i = 0 ; i < _poolSize; i++) @@ -234,4 +236,44 @@ public class NetworkConnectionScheduler { _selectorThread.addToWork(connection); } + + private final class SelectorThreadPoolExecutor extends ThreadPoolExecutor + { + private SelectorThreadPoolExecutor(final int corePoolSize, + final int maximumPoolSize, + final long keepAliveTime, + final BlockingQueue<Runnable> workQueue, + final ThreadFactory threadFactory) + { + super(corePoolSize, maximumPoolSize, keepAliveTime, TimeUnit.MINUTES, workQueue, threadFactory); + } + + @Override + protected void afterExecute(final Runnable task, final Throwable failure) + { + super.afterExecute(task, failure); + + if (task == _selectorThread && !_selectorThread.isClosed() && !isShutdown()) + { + // executor replaces an abruptly terminated Java worker, but not its long-lived selector task + restoreSelectorTask(task); + } + } + + private void restoreSelectorTask(final Runnable task) + { + try + { + execute(task); + LOGGER.warn("Selector processing task stopped unexpectedly; restored configured capacity"); + } + catch (RejectedExecutionException e) + { + if (!_selectorThread.isClosed() && !isShutdown()) + { + LOGGER.error("Failed to restore selector processing capacity", e); + } + } + } + } } diff --git a/broker-core/src/main/java/org/apache/qpid/server/transport/SelectorThread.java b/broker-core/src/main/java/org/apache/qpid/server/transport/SelectorThread.java index c96ab66246..5f888377e9 100644 --- a/broker-core/src/main/java/org/apache/qpid/server/transport/SelectorThread.java +++ b/broker-core/src/main/java/org/apache/qpid/server/transport/SelectorThread.java @@ -659,6 +659,11 @@ class SelectorThread extends Thread } + boolean isClosed() + { + return _closed.get(); + } + public void addToWork(final NonBlockingConnection connection) { if (_closed.get()) diff --git a/broker-core/src/test/java/org/apache/qpid/server/transport/NetworkConnectionSchedulerTest.java b/broker-core/src/test/java/org/apache/qpid/server/transport/NetworkConnectionSchedulerTest.java new file mode 100644 index 0000000000..d15255211c --- /dev/null +++ b/broker-core/src/test/java/org/apache/qpid/server/transport/NetworkConnectionSchedulerTest.java @@ -0,0 +1,148 @@ +/* + * + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + * + */ +package org.apache.qpid.server.transport; + +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.doAnswer; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import java.util.concurrent.BrokenBarrierException; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.CyclicBarrier; +import java.util.concurrent.ThreadFactory; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; + +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import org.apache.qpid.test.utils.UnitTestBase; + +public class NetworkConnectionSchedulerTest extends UnitTestBase +{ + private static final int SELECTOR_COUNT = 1; + private static final int PROCESSING_THREAD_COUNT = 2; + private static final int POOL_SIZE = SELECTOR_COUNT + PROCESSING_THREAD_COUNT; + private static final long TIMEOUT_SECONDS = 10L; + + private NetworkConnectionScheduler _scheduler; + + @Test + public void testProcessingCapacityRestoredAfterConcurrentTaskFailures() throws Exception + { + final CyclicBarrier failureBarrier = new CyclicBarrier(PROCESSING_THREAD_COUNT); + final CountDownLatch failureTasksStarted = new CountDownLatch(PROCESSING_THREAD_COUNT); + final CountDownLatch workerFailures = new CountDownLatch(PROCESSING_THREAD_COUNT); + final ThreadFactory threadFactory = runnable -> + { + final Thread thread = new Thread(runnable, getTestName()); + thread.setDaemon(true); + thread.setUncaughtExceptionHandler((ignored, failure) -> workerFailures.countDown()); + return thread; + }; + + _scheduler = new NetworkConnectionScheduler(getTestName(), SELECTOR_COUNT, POOL_SIZE, 1L, threadFactory); + _scheduler.start(); + + _scheduler.schedule(createFailingConnection(new NegativeArraySizeException("test failure"), failureBarrier, + failureTasksStarted)); + _scheduler.schedule(createFailingConnection(new StackOverflowError("test failure"), failureBarrier, + failureTasksStarted)); + + assertTrue(failureTasksStarted.await(TIMEOUT_SECONDS, TimeUnit.SECONDS), + "Concurrent failure tasks did not start"); + assertTrue(workerFailures.await(TIMEOUT_SECONDS, TimeUnit.SECONDS), + "Selector processing tasks did not terminate as expected"); + + final CyclicBarrier successfulWorkBarrier = new CyclicBarrier(PROCESSING_THREAD_COUNT); + final CountDownLatch successfulWorkCompleted = new CountDownLatch(PROCESSING_THREAD_COUNT); + for (int i = 0; i < PROCESSING_THREAD_COUNT; i++) + { + _scheduler.schedule(createSuccessfulConnection(successfulWorkBarrier, successfulWorkCompleted)); + } + + assertTrue(successfulWorkCompleted.await(TIMEOUT_SECONDS, TimeUnit.SECONDS), + "Configured selector processing capacity was not restored"); + } + + private NonBlockingConnection createFailingConnection(final Throwable failure, + final CyclicBarrier failureBarrier, + final CountDownLatch failureTasksStarted) + { + final NonBlockingConnection connection = createConnection(); + doAnswer(invocation -> + { + failureTasksStarted.countDown(); + await(failureBarrier); + throw failure; + }).when(connection).doWork(); + return connection; + } + + private NonBlockingConnection createSuccessfulConnection(final CyclicBarrier successfulWorkBarrier, + final CountDownLatch successfulWorkCompleted) + { + final NonBlockingConnection connection = createConnection(); + doAnswer(invocation -> + { + await(successfulWorkBarrier); + successfulWorkCompleted.countDown(); + return true; + }).when(connection).doWork(); + return connection; + } + + private NonBlockingConnection createConnection() + { + final NonBlockingConnection connection = mock(NonBlockingConnection.class); + when(connection.setScheduled()).thenReturn(true); + when(connection.getThreadName()).thenReturn(getTestName()); + when(connection.getScheduler()).thenReturn(_scheduler); + return connection; + } + + private static void await(final CyclicBarrier barrier) + { + try + { + barrier.await(TIMEOUT_SECONDS, TimeUnit.SECONDS); + } + catch (InterruptedException e) + { + Thread.currentThread().interrupt(); + throw new AssertionError("Interrupted while coordinating selector processing tasks", e); + } + catch (BrokenBarrierException | TimeoutException e) + { + throw new AssertionError("Failed to coordinate selector processing tasks", e); + } + } + + @AfterEach + public void tearDown() + { + if (_scheduler != null) + { + _scheduler.close(); + } + } +} --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
