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 dc8b6fa82561 CAMEL-25008, CAMEL-25007: camel-core - Supervising route
controller: fix a deadlock and a lost restart task (#26867)
dc8b6fa82561 is described below
commit dc8b6fa8256110ec36f5927a59295208f16bad37
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 13:48:34 2026 +0530
CAMEL-25008, CAMEL-25007: camel-core - Supervising route controller: fix a
deadlock and a lost restart task (#26867)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../engine/DefaultSupervisingRouteController.java | 4 +-
...ingRouteControllerStartWhileRestartingTest.java | 128 +++++++++++++++++++++
.../camel/util/backoff/BackOffTimerTask.java | 6 +-
.../camel/util/backoff/SimpleBackOffTimerTest.java | 53 +++++++++
4 files changed, 189 insertions(+), 2 deletions(-)
diff --git
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultSupervisingRouteController.java
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultSupervisingRouteController.java
index 0e21c2a19079..5ba2881b39be 100644
---
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultSupervisingRouteController.java
+++
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultSupervisingRouteController.java
@@ -741,7 +741,9 @@ public class DefaultSupervisingRouteController extends
DefaultRouteController im
}
}
- routes.remove(r);
+ // a cancelled task completes again when its
running attempt ends, and by then the route
+ // may have a new restart task, which must not be
removed
+ routes.remove(r, task);
});
return task;
diff --git
a/core/camel-core/src/test/java/org/apache/camel/impl/engine/DefaultSupervisingRouteControllerStartWhileRestartingTest.java
b/core/camel-core/src/test/java/org/apache/camel/impl/engine/DefaultSupervisingRouteControllerStartWhileRestartingTest.java
new file mode 100644
index 000000000000..bd47206acb22
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/impl/engine/DefaultSupervisingRouteControllerStartWhileRestartingTest.java
@@ -0,0 +1,128 @@
+/*
+ * 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.impl.engine;
+
+import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.CamelEvent;
+import org.apache.camel.spi.CamelEvent.RouteRestartingEvent;
+import org.apache.camel.spi.SupervisingRouteController;
+import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.DefaultConsumer;
+import org.apache.camel.support.DefaultEndpoint;
+import org.apache.camel.support.SimpleEventNotifierSupport;
+import org.apache.camel.util.backoff.BackOffTimer;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A route that fails to start manually while its previous restart attempt is
running gets a new restart task, which
+ * must stay registered when the previous attempt ends.
+ */
+class DefaultSupervisingRouteControllerStartWhileRestartingTest extends
ContextTestSupport {
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Test
+ void testStartRouteFailsWhileRestartAttemptIsRunning() throws Exception {
+ CountDownLatch attemptStarted = new CountDownLatch(1);
+ CountDownLatch routeStarted = new CountDownLatch(1);
+ AtomicReference<BackOffTimer.Task> oldTask = new AtomicReference<>();
+
+ context.addComponent("broken", new BrokenComponent());
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("broken:start").routeId("broken").to("mock:result");
+ }
+ });
+
+ SupervisingRouteController src =
context.getRouteController().supervising();
+ src.setBackOffDelay(10);
+ src.setInitialDelay(10);
+ src.setUnhealthyOnRestarting(true);
+
+ context.getManagementStrategy().addEventNotifier(new
SimpleEventNotifierSupport() {
+ @Override
+ public void notify(CamelEvent event) throws Exception {
+ if (event instanceof RouteRestartingEvent rre &&
"broken".equals(rre.getRoute().getRouteId())
+ && oldTask.compareAndSet(null,
src.getRestartingRouteState("broken"))) {
+ // the first restart attempt is about to start the route:
let it be started manually first
+ attemptStarted.countDown();
+ assertTrue(routeStarted.await(10, TimeUnit.SECONDS));
+ }
+ }
+ });
+
+ context.start();
+ assertTrue(attemptStarted.await(10, TimeUnit.SECONDS));
+
+ // the route is started manually, which cancels the old restart task,
and fails, so it gets a new restart task
+ assertThrows(Exception.class, () ->
context.getRouteController().startRoute("broken"));
+ BackOffTimer.Task newTask = src.getRestartingRouteState("broken");
+ assertNotSame(oldTask.get(), newTask);
+
+ // the old restart attempt fails too, and completes the old task again
+ CountDownLatch oldTaskCompleted = new CountDownLatch(1);
+ oldTask.get().whenComplete((task, cause) ->
oldTaskCompleted.countDown());
+ routeStarted.countDown();
+ assertTrue(oldTaskCompleted.await(10, TimeUnit.SECONDS));
+
+ // the new restart task is still the one of the route (so it is
reported, and a stopRoute cancels it)
+ assertSame(newTask, src.getRestartingRouteState("broken"));
+ assertTrue(((DefaultSupervisingRouteController)
src).hasUnhealthyRoutes());
+ }
+
+ private static final class BrokenComponent extends DefaultComponent {
+
+ @Override
+ protected Endpoint createEndpoint(String uri, String remaining,
Map<String, Object> parameters) {
+ return new DefaultEndpoint(uri, this) {
+ @Override
+ public Producer createProducer() {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public Consumer createConsumer(Processor processor) {
+ return new DefaultConsumer(this, processor) {
+ @Override
+ protected void doStart() {
+ throw new IllegalStateException("Cannot connect");
+ }
+ };
+ }
+ };
+ }
+ }
+}
diff --git
a/core/camel-util/src/main/java/org/apache/camel/util/backoff/BackOffTimerTask.java
b/core/camel-util/src/main/java/org/apache/camel/util/backoff/BackOffTimerTask.java
index 85f30d1f0fc2..92c5bfd4a9d9 100644
---
a/core/camel-util/src/main/java/org/apache/camel/util/backoff/BackOffTimerTask.java
+++
b/core/camel-util/src/main/java/org/apache/camel/util/backoff/BackOffTimerTask.java
@@ -217,12 +217,16 @@ public final class BackOffTimerTask implements
BackOffTimer.Task, Runnable {
void complete(Throwable throwable) {
this.cause = throwable;
+ List<BiConsumer<BackOffTimer.Task, Throwable>> copy;
lock.lock();
try {
- consumers.forEach(c -> c.accept(this, throwable));
+ copy = new ArrayList<>(consumers);
} finally {
lock.unlock();
}
+ // call the consumers without holding the lock, as they may take locks
of their own,
+ // which could deadlock with a thread holding such a lock and
cancelling this task
+ copy.forEach(c -> c.accept(this, throwable));
}
// *****************************
diff --git
a/core/camel-util/src/test/java/org/apache/camel/util/backoff/SimpleBackOffTimerTest.java
b/core/camel-util/src/test/java/org/apache/camel/util/backoff/SimpleBackOffTimerTest.java
index 634578742ded..28190a03dc70 100644
---
a/core/camel-util/src/test/java/org/apache/camel/util/backoff/SimpleBackOffTimerTest.java
+++
b/core/camel-util/src/test/java/org/apache/camel/util/backoff/SimpleBackOffTimerTest.java
@@ -17,15 +17,21 @@
package org.apache.camel.util.backoff;
import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
+import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.locks.ReentrantLock;
import org.junit.jupiter.api.Test;
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -209,4 +215,51 @@ public class SimpleBackOffTimerTest {
timer.close();
}
+ @Test
+ void testCancelWhileCompletionCallbackWaitsForLock() throws Exception {
+ // the completion callback takes a lock of the owner of the task (as
the supervising route controller does),
+ // and the owner cancels the task while it holds that lock
+ final ReentrantLock ownerLock = new ReentrantLock();
+ final CountDownLatch ownerHoldsLock = new CountDownLatch(1);
+ final AtomicReference<Thread> timerThread = new AtomicReference<>();
+ final ScheduledExecutorService executor =
Executors.newSingleThreadScheduledExecutor();
+ final ExecutorService owner = Executors.newSingleThreadExecutor();
+ final BackOff backOff = BackOff.builder().delay(10).build();
+ final SimpleBackOffTimer timer = new SimpleBackOffTimer(executor);
+
+ try {
+ BackOffTimer.Task task = timer.schedule(
+ backOff,
+ context -> {
+ timerThread.set(Thread.currentThread());
+ ownerHoldsLock.await();
+ // done, so the task completes
+ return false;
+ });
+ task.whenComplete(
+ (context, throwable) -> {
+ ownerLock.lock();
+ ownerLock.unlock();
+ });
+
+ Future<?> cancel = owner.submit(() -> {
+ ownerLock.lock();
+ try {
+ ownerHoldsLock.countDown();
+ // the completion callback of the task waits for the lock
+ await().atMost(5, TimeUnit.SECONDS)
+ .until(() -> timerThread.get() != null &&
ownerLock.hasQueuedThread(timerThread.get()));
+ task.cancel();
+ } finally {
+ ownerLock.unlock();
+ }
+ });
+
+ assertDoesNotThrow(() -> cancel.get(5, TimeUnit.SECONDS),
"Cancelling the task should not deadlock");
+ } finally {
+ owner.shutdownNow();
+ executor.shutdownNow();
+ timer.close();
+ }
+ }
}