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 0e21c2a190791..5ba2881b39be6 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 @@ void start(RouteHolder route) { } } - 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 0000000000000..bd47206acb22d --- /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 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 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 85f30d1f0fc22..92c5bfd4a9d9b 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 @@ void stop() { void complete(Throwable throwable) { this.cause = throwable; + List> 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 634578742ded8..28190a03dc702 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 void testBackOffTimerRemoveOnComplete() throws Exception { 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 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(); + } + } }