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