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]

Reply via email to