This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new 19025294af48 CAMEL-25062: camel-cluster - ClusteredRoutePolicy must 
not deadlock with a leadership change on stop or removeRoute (#26945)
19025294af48 is described below

commit 19025294af487d4989f6a98528cb75099374cfe1
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 14:11:54 2026 +0530

    CAMEL-25062: camel-cluster - ClusteredRoutePolicy must not deadlock with a 
leadership change on stop or removeRoute (#26945)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../JGroupsRaftClusteredRoutePolicyTest.java       |  13 +-
 .../camel/impl/cluster/ClusteredRoutePolicy.java   | 212 ++++++++++--
 .../cluster/ClusteredRoutePolicyFactoryTest.java   |  61 ++++
 .../ClusteredRoutePolicyLeaderChangeTest.java      |   9 +-
 .../ClusteredRoutePolicyReleaseDeadlockTest.java   | 385 +++++++++++++++++++++
 .../camel/cluster/ClusteredRoutePolicyTest.java    |  11 +
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  16 +
 7 files changed, 664 insertions(+), 43 deletions(-)

diff --git 
a/components/camel-jgroups-raft/src/test/java/org/apache/camel/component/jgroups/raft/cluster/JGroupsRaftClusteredRoutePolicyTest.java
 
b/components/camel-jgroups-raft/src/test/java/org/apache/camel/component/jgroups/raft/cluster/JGroupsRaftClusteredRoutePolicyTest.java
index 48e2d9a88800..f91b85872d1e 100644
--- 
a/components/camel-jgroups-raft/src/test/java/org/apache/camel/component/jgroups/raft/cluster/JGroupsRaftClusteredRoutePolicyTest.java
+++ 
b/components/camel-jgroups-raft/src/test/java/org/apache/camel/component/jgroups/raft/cluster/JGroupsRaftClusteredRoutePolicyTest.java
@@ -17,6 +17,7 @@
 package org.apache.camel.component.jgroups.raft.cluster;
 
 import java.util.ArrayList;
+import java.util.concurrent.TimeUnit;
 
 import org.apache.camel.CamelContext;
 import org.apache.camel.ServiceStatus;
@@ -28,6 +29,7 @@ import org.jgroups.JChannel;
 import org.jgroups.raft.RaftHandle;
 import org.junit.jupiter.api.Test;
 
+import static org.awaitility.Awaitility.await;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 
 public class JGroupsRaftClusteredRoutePolicyTest extends 
JGroupsRaftClusterAbstractTest {
@@ -64,13 +66,13 @@ public class JGroupsRaftClusteredRoutePolicyTest extends 
JGroupsRaftClusterAbstr
         contextB.start();
         contextC.start();
         waitForLeader(50, handleA, handleB, handleC);
-        assertEquals(1, countActiveFromEndpoints(lcc, rn));
+        awaitOneActiveRoute();
 
         contextA.stop();
         // Ensure channel A is fully closed before checking for a new leader
         chA.close();
         waitForLeader(50, handleB, handleC);
-        assertEquals(1, countActiveFromEndpoints(lcc, rn));
+        awaitOneActiveRoute();
 
         contextB.stop();
         // Ensure channel B is fully closed before creating a new channel with 
the same member name
@@ -82,7 +84,7 @@ public class JGroupsRaftClusteredRoutePolicyTest extends 
JGroupsRaftClusterAbstr
         lcc.set(0, contextA);
         contextA.start();
         waitForLeader(50, handleA, handleC);
-        assertEquals(1, countActiveFromEndpoints(lcc, rn));
+        awaitOneActiveRoute();
     }
 
     private CamelContext createContext(String id, RaftHandle rh) throws 
Exception {
@@ -109,6 +111,11 @@ public class JGroupsRaftClusteredRoutePolicyTest extends 
JGroupsRaftClusterAbstr
         return context;
     }
 
+    private void awaitOneActiveRoute() {
+        // the policy starts and stops the routes on its own thread, shortly 
after the leadership change
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, countActiveFromEndpoints(lcc, rn)));
+    }
+
     private int countActiveFromEndpoints(ArrayList<CamelContext> lcc, 
ArrayList<String> rn) {
         int result = 0;
         if (lcc.size() != rn.size()) {
diff --git 
a/core/camel-cluster/src/main/java/org/apache/camel/impl/cluster/ClusteredRoutePolicy.java
 
b/core/camel-cluster/src/main/java/org/apache/camel/impl/cluster/ClusteredRoutePolicy.java
index 43a44af45fb3..c62af73f601c 100644
--- 
a/core/camel-cluster/src/main/java/org/apache/camel/impl/cluster/ClusteredRoutePolicy.java
+++ 
b/core/camel-cluster/src/main/java/org/apache/camel/impl/cluster/ClusteredRoutePolicy.java
@@ -17,11 +17,14 @@
 package org.apache.camel.impl.cluster;
 
 import java.time.Duration;
-import java.util.HashSet;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
 import java.util.concurrent.locks.Lock;
 import java.util.concurrent.locks.ReentrantLock;
 import java.util.stream.Collectors;
@@ -40,12 +43,14 @@ import org.apache.camel.cluster.CamelClusterService;
 import org.apache.camel.cluster.CamelClusterView;
 import org.apache.camel.spi.CamelEvent;
 import org.apache.camel.spi.CamelEvent.CamelContextStartedEvent;
+import org.apache.camel.spi.ThreadPoolProfile;
 import org.apache.camel.support.RoutePolicySupport;
 import org.apache.camel.support.SimpleEventNotifierSupport;
 import org.apache.camel.support.cluster.ClusterServiceHelper;
 import org.apache.camel.support.cluster.ClusterServiceSelectors;
 import org.apache.camel.util.ObjectHelper;
 import org.apache.camel.util.ReferenceCount;
+import org.apache.camel.util.concurrent.ThreadPoolRejectedPolicy;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
@@ -53,6 +58,9 @@ import org.slf4j.LoggerFactory;
 public final class ClusteredRoutePolicy extends RoutePolicySupport implements 
CamelContextAware {
 
     private static final Logger LOG = 
LoggerFactory.getLogger(ClusteredRoutePolicy.class);
+    private static final String THREAD_NAME = "ClusteredRoutePolicy";
+    // how long the policy thread waits for another leadership change before 
it exits
+    private static final long LEADERSHIP_THREAD_KEEP_ALIVE_MILLIS = 1000;
 
     private final AtomicBoolean leader;
     private final Set<Route> autoStartupRoutes;
@@ -66,12 +74,15 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
     private final String namespace;
     private final CamelClusterService.Selector clusterServiceSelector;
     private final Lock lock;
+    private final Lock retainLock;
+    private final AtomicReference<CamelClusterView> clusterView;
+    private final AtomicBoolean leadershipChangePending;
     private CamelClusterService clusterService;
-    private CamelClusterView clusterView;
     private volatile boolean startManagedRoutesEarly;
 
     private Duration initialDelay;
-    private ScheduledExecutorService executorService;
+    private volatile ExecutorService leadershipExecutor;
+    private volatile ScheduledExecutorService initialDelayExecutor;
 
     private CamelContext camelContext;
 
@@ -86,9 +97,13 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
         this.leadershipEventListener = new CamelClusterLeadershipListener();
 
         this.lock = new ReentrantLock();
-        this.stoppedRoutes = new HashSet<>();
-        this.startedRoutes = new HashSet<>();
-        this.autoStartupRoutes = new HashSet<>();
+        this.retainLock = new ReentrantLock();
+        this.leadershipChangePending = new AtomicBoolean();
+        // the routes are started and stopped on the policy thread, while 
routes are added and removed by the caller
+        this.stoppedRoutes = ConcurrentHashMap.newKeySet();
+        this.startedRoutes = ConcurrentHashMap.newKeySet();
+        this.autoStartupRoutes = ConcurrentHashMap.newKeySet();
+        this.clusterView = new AtomicReference<>();
         this.leader = new AtomicBoolean();
         this.contextStarted = new AtomicBoolean();
         this.initialDelay = Duration.ofMillis(0);
@@ -126,8 +141,8 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
             this.camelContext = camelContext;
             this.camelContext.addStartupListener(this.listener);
             
this.camelContext.getManagementStrategy().addEventNotifier(this.listener);
-            this.executorService
-                    = 
camelContext.getExecutorServiceManager().newSingleThreadScheduledExecutor(this, 
"ClusteredRoutePolicy");
+            this.leadershipExecutor = camelContext.getExecutorServiceManager()
+                    .newThreadPool(this, THREAD_NAME, 
newLeadershipThreadPoolProfile());
         } catch (Exception e) {
             throw new RuntimeException(e);
         }
@@ -161,7 +176,12 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
         super.onInit(route);
 
         // Increase number of managed routes by this policy, acquire policy 
view on first run
-        this.refCount.retain();
+        retainLock.lock();
+        try {
+            this.refCount.retain();
+        } finally {
+            retainLock.unlock();
+        }
 
         if (route.isAutoStartup()) {
             autoStartupRoutes.add(route);
@@ -196,13 +216,23 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
     @Override
     public void onRemove(Route route) {
         // Decrease number of managed routes, release view once there are no 
route left
-        refCount.release();
+        retainLock.lock();
+        try {
+            refCount.release();
+        } finally {
+            retainLock.unlock();
+        }
         autoStartupRoutes.remove(route);
     }
 
     @Override
     protected void doShutdown() throws Exception {
-        releaseClusterView();
+        retainLock.lock();
+        try {
+            releaseClusterView();
+        } finally {
+            retainLock.unlock();
+        }
         removeCamelEventListeners();
     }
 
@@ -213,44 +243,58 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
     private void removeCamelEventListeners() {
         if (camelContext != null) {
             camelContext.getManagementStrategy().removeEventNotifier(listener);
-            if (executorService != null) {
-                
camelContext.getExecutorServiceManager().shutdownNow(executorService);
+            ExecutorService executor = leadershipExecutor;
+            if (executor != null) {
+                camelContext.getExecutorServiceManager().shutdownNow(executor);
+            }
+            ScheduledExecutorService scheduler = initialDelayExecutor;
+            if (scheduler != null) {
+                initialDelayExecutor = null;
+                
camelContext.getExecutorServiceManager().shutdownNow(scheduler);
             }
         }
     }
 
+    // The view and the cluster service are called without holding the policy 
lock. The view holds its own lock while
+    // it notifies the listeners, and the policy thread holds the policy lock 
while it starts or stops routes, which
+    // needs the CamelContext route lock. onRemove (and so releaseClusterView) 
can run while the caller holds that route
+    // lock, for example CamelContext.removeRoute, so taking the policy lock 
here could deadlock (CAMEL-24545).
+    //
+    // Retain and release are serialized by the retain lock instead, so a 
route added and a route removed at the same
+    // time on a shared policy cannot release the view that has just been 
retained. Only the threads adding and
+    // removing routes (and the shutdown) take the retain lock, and they never 
wait for the policy lock or the policy
+    // thread while they hold it.
+
     private void retainClusterView() {
-        lock.lock();
         try {
-            clusterView = clusterService.getView(namespace);
-            clusterView.addEventListener(leadershipEventListener);
+            CamelClusterView view = clusterService.getView(namespace);
+            clusterView.set(view);
+            // Take the current leadership right away, as the listener applies 
it asynchronously: onInit uses it to
+            // decide whether the route controller can start the route. No 
route is managed yet, so there is nothing
+            // to start or stop here.
+            leader.set(view.getLocalMember().isLeader());
+            view.addEventListener(leadershipEventListener);
         } catch (Exception e) {
             throw new RuntimeException(e);
-        } finally {
-            lock.unlock();
         }
     }
 
     private void releaseClusterView() {
-        lock.lock();
+        CamelClusterView view = clusterView.getAndSet(null);
         try {
-            // Remove event listener
-            if (clusterView != null) {
-                clusterView.removeEventListener(leadershipEventListener);
+            if (view != null) {
+                // Remove event listener
+                view.removeEventListener(leadershipEventListener);
 
                 // If all the routes have been removed then the view and its
                 // resources can eventually be released.
-                clusterView.getClusterService().releaseView(clusterView);
-                clusterView = null;
+                view.getClusterService().releaseView(view);
             }
         } catch (Exception e) {
             throw new RuntimeException(e);
         } finally {
-            try {
-                setLeader(false);
-            } finally {
-                lock.unlock();
-            }
+            // the routes managed by this policy have been removed or are 
being shut down, so there is nothing to stop
+            leader.set(false);
         }
     }
 
@@ -263,9 +307,25 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
     // Route managements
     // ****************************************************
 
-    private void setLeader(boolean isLeader) {
+    private void setLeader() {
         lock.lock();
         try {
+            // the leadership is read when the change is applied, not when it 
was handed over, so the last change wins
+            CamelClusterView view = clusterView.get();
+            if (view == null) {
+                // the view has been released in the meantime
+                return;
+            }
+
+            if (camelContext.isStopping()) {
+                // The CamelContext stops all its routes, starting with their 
consumers, so neither a leadership taken
+                // nor a leadership lost is applied: a route must not be 
started now, and stopping it here would run a
+                // second shutdown of the route next to the one of the 
CamelContext.
+                LOG.debug("Ignoring leadership change as CamelContext is 
stopping");
+                return;
+            }
+
+            boolean isLeader = view.getLocalMember().isLeader();
             if (isLeader && leader.compareAndSet(false, isLeader)) {
                 LOG.debug("Leadership taken");
                 startManagedRoutes();
@@ -358,11 +418,18 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
                     
startedRoutes.stream().map(Route::getId).collect(Collectors.joining(",")));
         }
 
-        if (startManagedRoutesEarly) {
-            LOG.debug(
-                    "CamelContext is now fully started, can now start managed 
routes eager as we were appointed leader during early startup");
-            startManagedRoutesEarly = false;
-            startManagedRoutes();
+        // under the policy lock, as the policy thread may be handling a 
leadership change and deferring the start of
+        // the routes while the CamelContext is starting
+        lock.lock();
+        try {
+            if (startManagedRoutesEarly) {
+                LOG.debug(
+                        "CamelContext is now fully started, can now start 
managed routes eager as we were appointed leader during early startup");
+                startManagedRoutesEarly = false;
+                startManagedRoutes();
+            }
+        } finally {
+            lock.unlock();
         }
     }
 
@@ -373,10 +440,69 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
     private class CamelClusterLeadershipListener implements 
CamelClusterEventListener.Leadership {
         @Override
         public void leadershipChanged(CamelClusterView view, 
CamelClusterMember leader) {
-            setLeader(clusterView.getLocalMember().isLeader());
+            if (clusterView.get() == null) {
+                // the view has been released
+                return;
+            }
+
+            // The view calls the listeners while it holds its lock, so the 
routes are started and stopped on the
+            // policy thread instead. Starting a route can take a while, and 
it needs the CamelContext route lock,
+            // which a thread removing a route or stopping the CamelContext 
may hold while it releases the view.
+            handOverLeadershipChange();
+        }
+    }
+
+    private void handOverLeadershipChange() {
+        // The policy thread reads the leadership when it applies a change, so 
a change that is still queued covers
+        // this one too, and at most one change is queued. getAndSet is used 
on both sides so the policy thread sees
+        // the leadership that the view has set before it fired this event.
+        if (leadershipChangePending.getAndSet(true)) {
+            return;
+        }
+
+        final ExecutorService executor = leadershipExecutor;
+        if (executor == null) {
+            // Never apply the change on this thread, as it holds the lock of 
the view
+            leadershipChangePending.set(false);
+            LOG.warn("Ignoring leadership change as ClusteredRoutePolicy for 
namespace {} has no CamelContext", namespace);
+            return;
+        }
+
+        try {
+            executor.execute(this::applyLeadershipChange);
+        } catch (RejectedExecutionException e) {
+            leadershipChangePending.set(false);
+            LOG.debug("Ignoring leadership change as ClusteredRoutePolicy for 
namespace {} has been shut down", namespace);
         }
     }
 
+    private void applyLeadershipChange() {
+        // from now on a leadership change queues a new task
+        leadershipChangePending.getAndSet(false);
+        try {
+            setLeader();
+        } catch (Exception e) {
+            LOG.warn("Error applying leadership change of ClusteredRoutePolicy 
for namespace {}. This exception is ignored.",
+                    namespace, e);
+        }
+    }
+
+    private static ThreadPoolProfile newLeadershipThreadPoolProfile() {
+        ThreadPoolProfile profile = new ThreadPoolProfile(THREAD_NAME);
+        // A single thread, so the changes are applied one after the other, 
and no thread while the leadership does
+        // not change, so a policy per route does not keep a thread per route.
+        profile.setPoolSize(0);
+        profile.setMaxPoolSize(1);
+        profile.setKeepAliveTime(LEADERSHIP_THREAD_KEEP_ALIVE_MILLIS);
+        profile.setTimeUnit(TimeUnit.MILLISECONDS);
+        profile.setAllowCoreThreadTimeOut(true);
+        // at most one change is queued (see handOverLeadershipChange)
+        profile.setMaxQueueSize(1);
+        // never run a change on the thread of the view: reject it once the 
policy has been shut down
+        profile.setRejectedPolicy(ThreadPoolRejectedPolicy.Abort);
+        return profile;
+    }
+
     private class CamelContextStartupListener extends 
SimpleEventNotifierSupport
             implements ExtendedStartupListener, NonManagedService {
         @Override
@@ -420,8 +546,18 @@ public final class ClusteredRoutePolicy extends 
RoutePolicySupport implements Ca
                 // Eventually delay the startup of the routes a later time
                 if (initialDelay.toMillis() > 0) {
                     LOG.debug("Policy will be effective in {}", initialDelay);
-                    
executorService.schedule(ClusteredRoutePolicy.this::onCamelContextStarted, 
initialDelay.toMillis(),
-                            TimeUnit.MILLISECONDS);
+                    ScheduledExecutorService scheduler = 
camelContext.getExecutorServiceManager()
+                            
.newSingleThreadScheduledExecutor(ClusteredRoutePolicy.this, THREAD_NAME);
+                    initialDelayExecutor = scheduler;
+                    scheduler.schedule(() -> {
+                        try {
+                            ClusteredRoutePolicy.this.onCamelContextStarted();
+                        } finally {
+                            // the delay applies once, so its thread is not 
needed anymore
+                            initialDelayExecutor = null;
+                            
camelContext.getExecutorServiceManager().shutdown(scheduler);
+                        }
+                    }, initialDelay.toMillis(), TimeUnit.MILLISECONDS);
                 } else {
                     ClusteredRoutePolicy.this.onCamelContextStarted();
                 }
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyFactoryTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyFactoryTest.java
index 9680b9fbb566..4b35a230eeb3 100644
--- 
a/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyFactoryTest.java
+++ 
b/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyFactoryTest.java
@@ -19,6 +19,7 @@ package org.apache.camel.cluster;
 import java.util.Collections;
 import java.util.List;
 import java.util.Optional;
+import java.util.concurrent.TimeUnit;
 
 import org.apache.camel.CamelContext;
 import org.apache.camel.ContextTestSupport;
@@ -30,7 +31,9 @@ import 
org.apache.camel.support.cluster.AbstractCamelClusterService;
 import org.apache.camel.support.cluster.AbstractCamelClusterView;
 import org.junit.jupiter.api.Test;
 
+import static org.awaitility.Awaitility.await;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 public class ClusteredRoutePolicyFactoryTest extends ContextTestSupport {
 
@@ -108,6 +111,9 @@ public class ClusteredRoutePolicyFactoryTest extends 
ContextTestSupport {
     @Test
     public void testClusteredRoutePolicyFactoryAddRouteAlreadyLeader() throws 
Exception {
         cs.getView().setLeader(true);
+        // the policy starts the routes on its own thread
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(
+                () -> assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("foo")));
 
         context.addRoutes(new RouteBuilder() {
             @Override
@@ -133,6 +139,61 @@ public class ClusteredRoutePolicyFactoryTest extends 
ContextTestSupport {
         assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("bar"));
     }
 
+    @Test
+    public void testNoPolicyThreadLeftWhenIdleOrRoutesRemoved() throws 
Exception {
+        cs.getView().setLeader(true);
+
+        for (int cycle = 0; cycle < 3; cycle++) {
+            final String prefix = "route-" + cycle + "-";
+            context.addRoutes(new RouteBuilder() {
+                @Override
+                public void configure() {
+                    for (int i = 0; i < 20; i++) {
+                        from("seda:" + prefix + i).routeId(prefix + i)
+                                .to("mock:result");
+                    }
+                }
+            });
+
+            await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
+                for (int i = 0; i < 20; i++) {
+                    assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus(prefix + i));
+                }
+            });
+
+            // each route has its own policy, and a policy keeps no thread 
while the leadership does not change
+            awaitNoPolicyThread();
+
+            for (int i = 0; i < 20; i++) {
+                context.getRouteController().stopRoute(prefix + i);
+                assertTrue(context.removeRoute(prefix + i));
+            }
+        }
+
+        // the removed routes left no thread behind
+        awaitNoPolicyThread();
+
+        // a leadership change after the policy thread has exited is still 
applied
+        cs.getView().setLeader(false);
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(
+                () -> assertEquals(ServiceStatus.Stopped, 
context.getRouteController().getRouteStatus("foo")));
+        awaitNoPolicyThread();
+    }
+
+    private void awaitNoPolicyThread() {
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(
+                () -> assertEquals(0, countPolicyThreads(), "live 
ClusteredRoutePolicy threads"));
+    }
+
+    private long countPolicyThreads() {
+        String camelId = "(" + context.getName() + ")";
+        return Thread.getAllStackTraces().keySet().stream()
+                .filter(Thread::isAlive)
+                .map(Thread::getName)
+                .filter(name -> name.contains(camelId) && name.endsWith(" - 
ClusteredRoutePolicy"))
+                .count();
+    }
+
     // *********************************
     // Helpers
     // *********************************
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyLeaderChangeTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyLeaderChangeTest.java
index 5a6ea92c8318..1c7f2a2d4906 100644
--- 
a/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyLeaderChangeTest.java
+++ 
b/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyLeaderChangeTest.java
@@ -19,6 +19,7 @@ package org.apache.camel.cluster;
 import java.util.Collections;
 import java.util.List;
 import java.util.Optional;
+import java.util.concurrent.TimeUnit;
 
 import org.apache.camel.CamelContext;
 import org.apache.camel.ContextTestSupport;
@@ -29,6 +30,7 @@ import 
org.apache.camel.support.cluster.AbstractCamelClusterService;
 import org.apache.camel.support.cluster.AbstractCamelClusterView;
 import org.junit.jupiter.api.Test;
 
+import static org.awaitility.Awaitility.await;
 import static org.junit.jupiter.api.Assertions.*;
 
 public class ClusteredRoutePolicyLeaderChangeTest extends ContextTestSupport {
@@ -52,11 +54,14 @@ public class ClusteredRoutePolicyLeaderChangeTest extends 
ContextTestSupport {
     public void testClusteredRoutePolicyOnLeadershipLost() {
         cs.getView().setLeader(true);
 
-        assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("foo"));
+        // the policy starts and stops the routes on its own thread
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(
+                () -> assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("foo")));
 
         cs.getView().setLeader(false);
 
-        assertEquals(ServiceStatus.Stopped, 
context.getRouteController().getRouteStatus("foo"));
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(
+                () -> assertEquals(ServiceStatus.Stopped, 
context.getRouteController().getRouteStatus("foo")));
     }
 
     @Override
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyReleaseDeadlockTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyReleaseDeadlockTest.java
new file mode 100644
index 000000000000..67ed9b4103e3
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyReleaseDeadlockTest.java
@@ -0,0 +1,385 @@
+/*
+ * 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.camel.cluster;
+
+import java.time.Duration;
+import java.util.Collections;
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.Component;
+import org.apache.camel.Consumer;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.Route;
+import org.apache.camel.ServiceStatus;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.impl.cluster.ClusteredRoutePolicyFactory;
+import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.DefaultConsumer;
+import org.apache.camel.support.DefaultEndpoint;
+import org.apache.camel.support.RoutePolicySupport;
+import org.apache.camel.support.cluster.AbstractCamelClusterService;
+import org.apache.camel.support.cluster.AbstractCamelClusterView;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTimeoutPreemptively;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A CamelContext stop or a route removal must not deadlock with a leadership 
change that is starting the routes of a
+ * {@link org.apache.camel.impl.cluster.ClusteredRoutePolicy}.
+ * <p/>
+ * Each route gets its own policy from the factory, so the view notifies one 
listener per route. The route "slow" takes
+ * the leadership first and its start blocks until the test lets it go. A gate 
listener, registered between the
+ * listeners of "slow" and "other", holds the leadership notification until 
the view is being released, so the policies
+ * of "other" and "late" are always notified while a route is being removed. 
The policy of "late" is still registered
+ * then, and has to start its route, which needs the CamelContext route lock 
that removeRoute holds.
+ * <p/>
+ * This test does not use ContextTestSupport, so a regression makes the test 
fail instead of hanging the build in the
+ * tear down.
+ */
+public class ClusteredRoutePolicyReleaseDeadlockTest {
+
+    private static final String NAMESPACE = "my-ns";
+    private static final Duration TIMEOUT = Duration.ofSeconds(20);
+
+    private final AtomicBoolean slowStartArmed = new AtomicBoolean();
+    private final CountDownLatch slowStarting = new CountDownLatch(1);
+    private final CountDownLatch slowProceed = new CountDownLatch(1);
+
+    private final AtomicBoolean gateArmed = new AtomicBoolean();
+    private final CountDownLatch viewReleasing = new CountDownLatch(1);
+
+    private DefaultCamelContext context;
+    private TestClusterService cs;
+
+    @BeforeEach
+    public void setUp() throws Exception {
+        cs = new TestClusterService("my-cluster-service");
+
+        context = new DefaultCamelContext();
+        context.disableJMX();
+        context.addService(cs);
+        context.addComponent("slow", new SlowComponent());
+        
context.addRoutePolicyFactory(ClusteredRoutePolicyFactory.forNamespace(NAMESPACE));
+        context.addRoutes(new RouteBuilder() {
+            @Override
+            public void configure() {
+                // stopping the CamelContext shuts down "late" and "other" 
before "slow"
+                from("slow:start").routeId("slow").startupOrder(1)
+                        .to("mock:slow");
+                // the gate policy comes before the policy created by the 
factory, so its listener is registered
+                // after the one of "slow" and before the one of "other"
+                
from("seda:other").routeId("other").startupOrder(2).routePolicy(new 
GatePolicy())
+                        .to("mock:other");
+                // its listener is registered after the gate
+                from("seda:late").routeId("late").startupOrder(3)
+                        .to("mock:late");
+            }
+        });
+        context.start();
+    }
+
+    @AfterEach
+    public void tearDown() throws InterruptedException {
+        // never leave a test thread blocked on a latch
+        slowProceed.countDown();
+        viewReleasing.countDown();
+        if (context != null && !context.isStopped()) {
+            // stop on another thread, so a regression cannot hang the build 
here
+            Thread stopper = newThread("stop-after-test", context::stop);
+            stopper.start();
+            stopper.join(TIMEOUT.toMillis());
+        }
+    }
+
+    @Test
+    public void testCamelContextStopDuringLeadershipChange() throws Exception {
+        slowStartArmed.set(true);
+        gateArmed.set(true);
+
+        Thread dispatcher = newThread("leadership", () -> 
cs.getView().setLeader(true));
+        dispatcher.start();
+        assertTrue(slowStarting.await(TIMEOUT.toSeconds(), TimeUnit.SECONDS), 
"route slow is not being started");
+
+        Thread operator = newThread("operator", context::stop);
+        operator.start();
+        awaitBlocked(operator);
+
+        // the start of "slow" completes while the CamelContext is being 
stopped
+        slowProceed.countDown();
+
+        assertTimeoutPreemptively(TIMEOUT, () -> {
+            operator.join();
+            dispatcher.join();
+        }, "CamelContext stop deadlocked with the leadership change");
+
+        assertTrue(context.isStopped());
+        // the policies released the view
+        assertFalse(cs.getView().isRunning());
+    }
+
+    @Test
+    public void testRemoveRouteDuringLeadershipChange() throws Exception {
+        slowStartArmed.set(true);
+        gateArmed.set(true);
+
+        Thread dispatcher = newThread("leadership", () -> 
cs.getView().setLeader(true));
+        dispatcher.start();
+        assertTrue(slowStarting.await(TIMEOUT.toSeconds(), TimeUnit.SECONDS), 
"route slow is not being started");
+
+        AtomicReference<Boolean> removed = new AtomicReference<>();
+        Thread operator = newThread("operator", () -> {
+            try {
+                context.getRouteController().stopRoute("other");
+                removed.set(context.removeRoute("other"));
+            } catch (Exception e) {
+                throw new RuntimeException(e);
+            }
+        });
+        operator.start();
+        awaitBlocked(operator);
+
+        // the start of "slow" completes while "other" is being removed
+        slowProceed.countDown();
+
+        assertTimeoutPreemptively(TIMEOUT, () -> {
+            operator.join();
+            dispatcher.join();
+        }, "removeRoute deadlocked with the leadership change");
+
+        assertEquals(Boolean.TRUE, removed.get());
+        assertNull(context.getRoute("other"));
+
+        // the other routes still have the leadership and keep running
+        await().atMost(TIMEOUT).untilAsserted(() -> {
+            assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("slow"));
+            assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("late"));
+        });
+        assertTrue(cs.getView().isRunning());
+    }
+
+    @Test
+    public void testLeadershipChangesStartAndStopRoutes() throws Exception {
+        assertEquals(ServiceStatus.Stopped, 
context.getRouteController().getRouteStatus("slow"));
+        assertEquals(ServiceStatus.Stopped, 
context.getRouteController().getRouteStatus("other"));
+        assertEquals(ServiceStatus.Stopped, 
context.getRouteController().getRouteStatus("late"));
+
+        cs.getView().setLeader(true);
+        awaitRouteStatus(ServiceStatus.Started);
+
+        cs.getView().setLeader(false);
+        awaitRouteStatus(ServiceStatus.Stopped);
+
+        // the policy thread applies the changes in order and ends with the 
last one
+        cs.getView().setLeader(true);
+        cs.getView().setLeader(false);
+        cs.getView().setLeader(true);
+        // the routes may be started for a moment by an earlier change, so 
wait until they stay started
+        
await().during(Duration.ofMillis(500)).atMost(TIMEOUT).untilAsserted(() -> 
assertRouteStatus(ServiceStatus.Started));
+
+        context.getRouteController().stopRoute("other");
+        assertTrue(context.removeRoute("other"));
+        assertTrue(cs.getView().isRunning());
+
+        context.stop();
+        assertFalse(cs.getView().isRunning());
+    }
+
+    private void awaitRouteStatus(ServiceStatus status) {
+        await().atMost(TIMEOUT).untilAsserted(() -> assertRouteStatus(status));
+    }
+
+    private void assertRouteStatus(ServiceStatus status) {
+        assertEquals(status, 
context.getRouteController().getRouteStatus("slow"));
+        assertEquals(status, 
context.getRouteController().getRouteStatus("other"));
+        assertEquals(status, 
context.getRouteController().getRouteStatus("late"));
+    }
+
+    // Waits until the thread waits for a lock or a latch. This makes the race 
likely, not certain, as the thread may
+    // wait for something else first (for example the shutdown strategy), but 
the test does not rely on it: the gate
+    // listener holds the dispatch until the view is being released, so the 
deadlock is reached anyway without the fix.
+    private static void awaitBlocked(Thread thread) {
+        await().atMost(TIMEOUT).until(() -> {
+            Thread.State state = thread.getState();
+            return state == Thread.State.WAITING || state == 
Thread.State.TIMED_WAITING
+                    || state == Thread.State.BLOCKED;
+        });
+    }
+
+    private static Thread newThread(String name, Runnable task) {
+        Thread thread = new Thread(task, 
"ClusteredRoutePolicyReleaseDeadlockTest-" + name);
+        thread.setDaemon(true);
+        return thread;
+    }
+
+    // *********************************
+    // Helpers
+    // *********************************
+
+    private final class SlowComponent extends DefaultComponent {
+        @Override
+        protected Endpoint createEndpoint(String uri, String remaining, 
Map<String, Object> parameters) {
+            return new SlowEndpoint(uri, this);
+        }
+    }
+
+    private final class SlowEndpoint extends DefaultEndpoint {
+        SlowEndpoint(String endpointUri, Component component) {
+            super(endpointUri, component);
+        }
+
+        @Override
+        public Producer createProducer() {
+            throw new UnsupportedOperationException();
+        }
+
+        @Override
+        public Consumer createConsumer(Processor processor) {
+            return new DefaultConsumer(this, processor) {
+                @Override
+                protected void doStart() throws Exception {
+                    // a consumer that takes a while to connect, like a broker 
connection or a subscription
+                    if (slowStartArmed.compareAndSet(true, false)) {
+                        slowStarting.countDown();
+                        slowProceed.await(TIMEOUT.toSeconds(), 
TimeUnit.SECONDS);
+                    }
+                    super.doStart();
+                }
+            };
+        }
+    }
+
+    private final class GatePolicy extends RoutePolicySupport {
+        @Override
+        public void onInit(Route route) {
+            super.onInit(route);
+            // the test service caches its single view, so this is the view of 
the policies, without retaining it
+            
cs.createView(NAMESPACE).addEventListener((CamelClusterEventListener.Leadership)
 (view, leader) -> {
+                if (leader != null && gateArmed.compareAndSet(true, false)) {
+                    try {
+                        viewReleasing.await(TIMEOUT.toSeconds(), 
TimeUnit.SECONDS);
+                    } catch (InterruptedException e) {
+                        Thread.currentThread().interrupt();
+                    }
+                }
+            });
+        }
+    }
+
+    private final class TestClusterView extends AbstractCamelClusterView {
+        private volatile boolean leader;
+        private volatile boolean running;
+
+        TestClusterView(CamelClusterService cluster, String namespace) {
+            super(cluster, namespace);
+        }
+
+        @Override
+        public Optional<CamelClusterMember> getLeader() {
+            return leader ? Optional.of(getLocalMember()) : Optional.empty();
+        }
+
+        @Override
+        public CamelClusterMember getLocalMember() {
+            return new CamelClusterMember() {
+                @Override
+                public boolean isLeader() {
+                    return leader;
+                }
+
+                @Override
+                public boolean isLocal() {
+                    return true;
+                }
+
+                @Override
+                public String getId() {
+                    return getClusterService().getId();
+                }
+            };
+        }
+
+        @Override
+        public List<CamelClusterMember> getMembers() {
+            return Collections.emptyList();
+        }
+
+        @Override
+        public void removeEventListener(CamelClusterEventListener listener) {
+            // a policy is releasing the view: let the leadership notification 
go on
+            viewReleasing.countDown();
+            super.removeEventListener(listener);
+        }
+
+        @Override
+        protected void doStart() {
+            running = true;
+        }
+
+        @Override
+        protected void doStop() {
+            running = false;
+        }
+
+        void setLeader(boolean leader) {
+            this.leader = leader;
+
+            if (isRunAllowed()) {
+                fireLeadershipChangedEvent(getLeader().orElse(null));
+            }
+        }
+
+        boolean isRunning() {
+            return running;
+        }
+    }
+
+    private final class TestClusterService extends 
AbstractCamelClusterService<TestClusterView> {
+        private TestClusterView view;
+
+        TestClusterService(String id) {
+            super(id);
+        }
+
+        @Override
+        protected TestClusterView createView(String namespace) {
+            if (view == null) {
+                view = new TestClusterView(this, namespace);
+            }
+            return view;
+        }
+
+        TestClusterView getView() {
+            return view;
+        }
+    }
+}
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyTest.java
index 6997590032e7..e2426a482bbd 100644
--- 
a/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyTest.java
+++ 
b/core/camel-core/src/test/java/org/apache/camel/cluster/ClusteredRoutePolicyTest.java
@@ -19,6 +19,7 @@ package org.apache.camel.cluster;
 import java.util.Collections;
 import java.util.List;
 import java.util.Optional;
+import java.util.concurrent.TimeUnit;
 
 import org.apache.camel.CamelContext;
 import org.apache.camel.ContextTestSupport;
@@ -30,6 +31,7 @@ import 
org.apache.camel.support.cluster.AbstractCamelClusterService;
 import org.apache.camel.support.cluster.AbstractCamelClusterView;
 import org.junit.jupiter.api.Test;
 
+import static org.awaitility.Awaitility.await;
 import static org.junit.jupiter.api.Assertions.*;
 
 public class ClusteredRoutePolicyTest extends ContextTestSupport {
@@ -69,6 +71,9 @@ public class ClusteredRoutePolicyTest extends 
ContextTestSupport {
     @Test
     public void testClusteredRoutePolicyRemoveAllRoutes() throws Exception {
         cs.getView().setLeader(true);
+        // the policy starts the routes on its own thread
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(
+                () -> assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("foo")));
 
         context.getRouteController().stopRoute("foo");
         context.getRouteController().stopRoute("baz");
@@ -81,6 +86,9 @@ public class ClusteredRoutePolicyTest extends 
ContextTestSupport {
     @Test
     public void testClusteredRoutePolicyDontStartAutoStartFalseRoutes() {
         cs.getView().setLeader(true);
+        // the policy starts the routes on its own thread
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(
+                () -> assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("foo")));
 
         assertEquals(ServiceStatus.Stopped, 
context.getRouteController().getRouteStatus("baz"));
     }
@@ -116,6 +124,9 @@ public class ClusteredRoutePolicyTest extends 
ContextTestSupport {
     @Test
     public void testClusteredRoutePolicyAddRouteAlreadyLeader() throws 
Exception {
         cs.getView().setLeader(true);
+        // the policy starts the routes on its own thread
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(
+                () -> assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("foo")));
 
         context.addRoutes(new RouteBuilder() {
             @Override
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 45581cbb9a5c..b53995cf6da1 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -205,6 +205,22 @@ placeholder based property.
   variable was used as a `null` input. This also applies to a `source` without 
a prefix, as a plain name refers to a
   variable (`source="myVar"`), and to the `source` option of the `xslt`, 
`xslt-saxon` and `xquery` endpoints.
 
+=== ClusteredRoutePolicy - routes are started and stopped on the policy thread
+
+When the leadership changes, `ClusteredRoutePolicy` now starts and stops its 
routes on its own thread, instead of on
+the thread that fires the leadership event while it holds the lock of the 
cluster view. This fixes a deadlock between
+a leadership change that is starting routes and a CamelContext stop or a route 
removal. The thread only exists while
+there is a leadership change to apply, and exits when idle, so a policy per 
route (`ClusteredRoutePolicyFactory`) does
+not keep a thread per route. Leadership changes that reach the policy thread 
while the CamelContext is stopping are
+ignored, as the CamelContext is stopping all the routes.
+
+Most cluster services fire leadership events from their own threads, but some 
fire them on the calling thread, for
+example when a view starts or a listener is added. In all cases the routes are 
now started and stopped on the policy
+thread, shortly after the event. Code that fires the event itself, for example 
a custom `CamelClusterView` in a test,
+can no longer expect the routes to be started or stopped when the event 
returns, and should wait for the route status
+instead. Likewise, `isLeader()` of the policy (JMX attribute `Leader`) 
reflects the view right away, and the routes
+follow asynchronously.
+
 === Apache Avro trusted packages
 
 Camel now uses Apache Avro 1.12.2. Avro validates classes resolved from schemas

Reply via email to