This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch camel-4.22.x
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/camel-4.22.x by this push:
new baaed0229952 CAMEL-24626: camel-master - leadership gets its own lock,
and cancelled start tasks leave the task registry
baaed0229952 is described below
commit baaed0229952f6632d93a9fd632e17b9c575b691
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 7 20:34:20 2026 +0200
CAMEL-24626: camel-master - leadership gets its own lock, and cancelled
start tasks leave the task registry
Cherry-pick of #26112 onto camel-4.22.x.
Three follow-ups to CAMEL-24583:
1. BackgroundTask.cancel() now unschedules the task, releases its latch,
marks it Inactive, and removes it from the TaskManagerRegistry, closing a
leak where a cancelled task stayed registered for the lifetime of the
CamelContext. A cancel-during-registration race is also guarded.
2. MasterConsumer now uses a dedicated leadershipLock instead of relying on
the BaseService lock, eliminating a lock inversion with the cluster view's
read/write lock (doStop held the service lock and needed the view's write
lock; event dispatch held the view's read lock and needed the consumer's
service lock).
3. Adds documentation clarifying that a leader exhausting backOffMaxAttempts
consumes nothing until leadership changes, and that 0 retries indefinitely.
Closes #26172
Co-Authored-By: henrik242 <[email protected]>
---
.../camel/catalog/docs/master-component.adoc | 11 +++
.../src/main/docs/master-component.adoc | 11 +++
.../camel/component/master/MasterConsumer.java | 69 +++++++++------
.../master/MasterConsumerLeadershipTest.java | 90 ++++++++++++++++++++
.../consumer/SimpleMessageListenerContainer.java | 7 +-
.../support/task/task/BackgroundTaskTest.java | 98 ++++++++++++++++++++++
.../apache/camel/support/task/BackgroundTask.java | 59 +++++++++++--
7 files changed, 311 insertions(+), 34 deletions(-)
diff --git
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/master-component.adoc
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/master-component.adoc
index cb27680b1b6a..3d19c3d45459 100644
---
a/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/master-component.adoc
+++
b/catalog/camel-catalog/src/generated/resources/org/apache/camel/catalog/docs/master-component.adoc
@@ -35,6 +35,17 @@ include::partial$component-endpoint-headers.adoc[]
== Usage
+=== Starting the delegated consumer
+
+When a node takes the leadership, the delegated consumer is started in the
background, and the start is
+retried when it fails, for instance because the back end it consumes from is
not reachable yet. The retries
+are configured on the component with `backOffDelay` (the delay between
attempts, 5000 millis by default)
+and `backOffMaxAttempts` (the number of attempts, 10 by default).
+
+Once the attempts are used up, the node keeps the leadership but consumes
nothing until the leadership
+changes again, which is logged at ERROR level. Set `backOffMaxAttempts` to 0
to keep retrying for as long as
+the node is the leader instead.
+
=== Using the master endpoint
Prefix any camel endpoint with **master:someName:** where _someName_ is a
logical name and is
diff --git a/components/camel-master/src/main/docs/master-component.adoc
b/components/camel-master/src/main/docs/master-component.adoc
index cb27680b1b6a..3d19c3d45459 100644
--- a/components/camel-master/src/main/docs/master-component.adoc
+++ b/components/camel-master/src/main/docs/master-component.adoc
@@ -35,6 +35,17 @@ include::partial$component-endpoint-headers.adoc[]
== Usage
+=== Starting the delegated consumer
+
+When a node takes the leadership, the delegated consumer is started in the
background, and the start is
+retried when it fails, for instance because the back end it consumes from is
not reachable yet. The retries
+are configured on the component with `backOffDelay` (the delay between
attempts, 5000 millis by default)
+and `backOffMaxAttempts` (the number of attempts, 10 by default).
+
+Once the attempts are used up, the node keeps the leadership but consumes
nothing until the leadership
+changes again, which is logged at ERROR level. Set `backOffMaxAttempts` to 0
to keep retrying for as long as
+the node is the leader instead.
+
=== Using the master endpoint
Prefix any camel endpoint with **master:someName:** where _someName_ is a
logical name and is
diff --git
a/components/camel-master/src/main/java/org/apache/camel/component/master/MasterConsumer.java
b/components/camel-master/src/main/java/org/apache/camel/component/master/MasterConsumer.java
index f4606dd96841..f9d9f79685f2 100644
---
a/components/camel-master/src/main/java/org/apache/camel/component/master/MasterConsumer.java
+++
b/components/camel-master/src/main/java/org/apache/camel/component/master/MasterConsumer.java
@@ -20,6 +20,8 @@ import java.time.Duration;
import java.util.concurrent.Future;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.atomic.AtomicReference;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
import org.apache.camel.Consumer;
import org.apache.camel.Endpoint;
@@ -61,9 +63,14 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
private volatile CamelClusterView view;
private ResumeStrategy resumeStrategy;
private ScheduledExecutorService leaderPool;
- // leadership state and the pending start task are guarded by lock, which
is also held by the
- // service lifecycle methods, so a leadership event cannot interleave with
start/stop of this consumer
+ // leadership state and the pending start task are guarded by
leadershipLock. This is deliberately not
+ // the lock of BaseService: the cluster view dispatches events while
holding its own read lock and then
+ // needs this lock, while doStop holds the service lock and needs the
write lock of the view to remove
+ // the listener. Guarding the leadership with the service lock closes that
into a lock cycle, and it also
+ // makes every leadership event and start attempt wait for whatever
lifecycle operation is in progress
+ private final Lock leadershipLock = new ReentrantLock();
private boolean leadershipTaken;
+ private BackgroundTask leaderTask;
private Future<?> leaderTaskFuture;
public MasterConsumer(MasterEndpoint masterEndpoint, Processor processor,
CamelClusterService clusterService) {
@@ -109,14 +116,19 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
protected void doStop() throws Exception {
super.doStop();
- // a start can still be pending, cancel it first so it cannot start
the delegated consumer
- // after this consumer has been stopped
- leadershipTaken = false;
- cancelLeaderTask(true);
+ leadershipLock.lock();
+ try {
+ // a start can still be pending, cancel it first so it cannot
start the delegated consumer
+ // after this consumer has been stopped
+ leadershipTaken = false;
+ cancelLeaderTask(true);
+ } finally {
+ leadershipLock.unlock();
+ }
- // note: removeEventListener below needs the cluster view lock while
this thread holds the lock of
- // this service, which is the opposite order of an event dispatch.
Nothing that runs under this lock
- // may wait for the view, and the listener bails out before locking
once this consumer is stopping
+ // note: removeEventListener below needs the write lock of the cluster
view, while an event dispatch
+ // takes the read lock of the view and then leadershipLock. This
thread must not hold leadershipLock
+ // here, or the two orders deadlock
if (view != null) {
view.removeEventListener(leadershipListener);
@@ -160,6 +172,7 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
.withBudget(Budgets.iterationTimeBudget()
.withInterval(Duration.ofMillis(masterEndpoint.getComponent().getBackOffDelay()))
.withInitialDelay(Duration.ofSeconds(1))
+ // 0 or less leaves the unlimited default of the
builder in place
.withMaxIterations(masterEndpoint.getComponent().getBackOffMaxAttempts())
// the attempts are bounded by backOffMaxAttempts, not
by the 5s default duration of
// the builder, which would otherwise end the task
before the second attempt
@@ -170,7 +183,7 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
}
private void onLeadershipTaken() {
- lock.lock();
+ leadershipLock.lock();
try {
if (!isRunAllowed()) {
return;
@@ -188,14 +201,15 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
final BackgroundTask task = createTask();
// the consumer is created once and re-used by the start attempts
of this task
final AtomicReference<Consumer> attempt = new AtomicReference<>();
+ leaderTask = task;
leaderTaskFuture = task.schedule(getEndpoint().getCamelContext(),
() -> startDelegatedConsumer(task, attempt));
} finally {
- lock.unlock();
+ leadershipLock.unlock();
}
}
private boolean startDelegatedConsumer(BackgroundTask task,
AtomicReference<Consumer> attempt) {
- lock.lock();
+ leadershipLock.lock();
try {
if (!isRunAllowed()) {
return false;
@@ -212,14 +226,14 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
return true; // no more attempts
}
} finally {
- lock.unlock();
+ leadershipLock.unlock();
}
LOG.info("Leadership taken. Attempt #{} to start consumer: {}",
task.iteration(), delegatedEndpoint);
// the delegate is created and started without holding the lock. It
can block for a long time, and the
- // lock is taken by the service lifecycle and by the cluster view
event dispatch, which must not wait
- // for a broker connect. The leadership is re-checked below before the
consumer is published
+ // lock is taken by the cluster view event dispatch and by doStop,
which must not wait for a broker
+ // connect. The leadership is re-checked below before the consumer is
published
Consumer consumer = attempt.get();
Exception cause = null;
try {
@@ -247,7 +261,7 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
cause = e;
}
- lock.lock();
+ leadershipLock.lock();
try {
if (cause != null) {
// the consumer is kept for the next attempt. It is not
stopped here: a consumer that failed to
@@ -282,12 +296,12 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
cancelLeaderTask(false);
return true; // no more attempts
} finally {
- lock.unlock();
+ leadershipLock.unlock();
}
}
private void onLeadershipLost() {
- lock.lock();
+ leadershipLock.lock();
try {
leadershipTaken = false;
// a start scheduled by the leadership taken event may not have
run yet, cancel it so it
@@ -306,7 +320,7 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
}
LOG.info("Leadership lost. Consumer stopped: {}",
delegatedEndpoint);
} finally {
- lock.unlock();
+ leadershipLock.unlock();
}
}
@@ -315,8 +329,12 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
}
private void cancelLeaderTask(boolean mayInterruptIfRunning) {
- if (leaderTaskFuture != null) {
- leaderTaskFuture.cancel(mayInterruptIfRunning);
+ if (leaderTask != null) {
+ // cancelled through the task and not through its future, so the
task also leaves the
+ // TaskManagerRegistry. Only a run of the task removes it from
there, and once the schedule
+ // is cancelled no run is coming
+ leaderTask.cancel(mayInterruptIfRunning);
+ leaderTask = null;
leaderTaskFuture = null;
}
}
@@ -329,13 +347,12 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
@Override
public void leadershipChanged(CamelClusterView view,
CamelClusterMember leader) {
if (!isRunAllowed()) {
- // checked before taking the lock: this runs on the cluster
view dispatch thread while that
- // view holds its own lock, and a consumer that is stopping
holds this lock while it removes
- // this listener from the view
+ // this runs on the dispatch thread of the cluster view,
holding the lock of that view, so
+ // do no work at all for a consumer that is stopping
return;
}
- lock.lock();
+ leadershipLock.lock();
try {
if (!isRunAllowed()) {
return;
@@ -360,7 +377,7 @@ public class MasterConsumer extends DefaultConsumer
implements ResumeAware<Resum
}
}
} finally {
- lock.unlock();
+ leadershipLock.unlock();
}
}
}
diff --git
a/components/camel-master/src/test/java/org/apache/camel/component/master/MasterConsumerLeadershipTest.java
b/components/camel-master/src/test/java/org/apache/camel/component/master/MasterConsumerLeadershipTest.java
index 59a3a31f95bb..f39510bfbf14 100644
---
a/components/camel-master/src/test/java/org/apache/camel/component/master/MasterConsumerLeadershipTest.java
+++
b/components/camel-master/src/test/java/org/apache/camel/component/master/MasterConsumerLeadershipTest.java
@@ -36,8 +36,10 @@ import org.apache.camel.impl.DefaultCamelContext;
import org.apache.camel.support.DefaultComponent;
import org.apache.camel.support.DefaultConsumer;
import org.apache.camel.support.DefaultEndpoint;
+import org.apache.camel.support.PluginHelper;
import org.apache.camel.support.cluster.AbstractCamelClusterService;
import org.apache.camel.support.cluster.AbstractCamelClusterView;
+import org.apache.camel.support.task.TaskManagerRegistry;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -45,6 +47,8 @@ import org.junit.jupiter.api.Timeout;
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.assertTrue;
/**
* Verifies that the delegated consumer only ever runs while this node holds
the leadership, also when the leadership
@@ -83,6 +87,8 @@ public class MasterConsumerLeadershipTest {
@AfterEach
void tearDown() {
+ // a test that fails while the delegate is parked on a gate would
otherwise wedge the stop below
+ probe.releaseGates();
if (context != null) {
context.stop();
}
@@ -224,6 +230,66 @@ public class MasterConsumerLeadershipTest {
assertEquals(0, probe.started.get());
}
+ @Test
+ @Timeout(60)
+ void testCancellingAPendingStartRemovesTheTaskFromTheRegistry() {
+ TestClusterView view = clusterService.getTestView();
+ TaskManagerRegistry registry =
PluginHelper.getTaskManagerRegistry(context.getCamelContextExtension());
+ MasterComponent master = context.getComponent("master",
MasterComponent.class);
+ // the task must still be retrying when the leadership is lost below,
not exhausted by then
+ master.setBackOffMaxAttempts(1000);
+
+ probe.failStart.set(true);
+ view.setLeader(true);
+
+ // the task adds itself to the registry from its first run
+ await().atMost(20, TimeUnit.SECONDS).until(() ->
probe.startAttempts.get() >= 1);
+ await().atMost(20, TimeUnit.SECONDS).until(() ->
hasLeadershipTask(registry));
+
+ view.setLeader(false);
+
+ // only a run of the task removes it from the registry, and after the
cancel no run is coming
+ await().atMost(20, TimeUnit.SECONDS).untilAsserted(() ->
assertFalse(hasLeadershipTask(registry),
+ "The cancelled start task must not stay in the task
registry"));
+ }
+
+ @Test
+ @Timeout(60)
+ void testEventDispatchIsNotBlockedByALifecycleOperation() throws Exception
{
+ TestClusterView view = clusterService.getTestView();
+
+ view.setLeader(true);
+ await().atMost(10, TimeUnit.SECONDS).untilAsserted(() ->
assertEquals(1, probe.started.get()));
+
+ CountDownLatch suspendEntered = new CountDownLatch(1);
+ CountDownLatch suspendGate = new CountDownLatch(1);
+ probe.suspendEntered.set(suspendEntered);
+ probe.suspendGate.set(suspendGate);
+
+ // suspending holds the service lock of the master consumer for as
long as the delegate takes
+ MasterConsumer consumer = (MasterConsumer)
context.getRoute("master-route").getConsumer();
+ Thread suspender = new Thread(consumer::suspend, "suspend");
+ suspender.start();
+ assertTrue(suspendEntered.await(20, TimeUnit.SECONDS), "The suspend of
the delegate should have started");
+
+ // the cluster view dispatches its events while holding its own lock,
and needs that same lock again
+ // to remove the listener when the consumer stops. An event that waits
here for the service lock of
+ // the consumer is what closes that into a deadlock, so the leadership
must not be guarded by it
+ Thread dispatcher = new Thread(() -> view.setLeader(true),
"leadership-taken");
+ dispatcher.start();
+ try {
+ dispatcher.join(TimeUnit.SECONDS.toMillis(20));
+ assertFalse(dispatcher.isAlive(), "An event dispatch must not wait
for the service lock of the consumer");
+ } finally {
+ suspendGate.countDown();
+ suspender.join(TimeUnit.SECONDS.toMillis(20));
+ }
+ }
+
+ private static boolean hasLeadershipTask(TaskManagerRegistry registry) {
+ return registry.getTasks().stream().anyMatch(task ->
"Leadership".equals(task.getName()));
+ }
+
// ************************************
// Delegated endpoint under observation
// ************************************
@@ -235,11 +301,22 @@ public class MasterConsumerLeadershipTest {
private final AtomicInteger stopped = new AtomicInteger();
private final AtomicBoolean failStart = new AtomicBoolean();
private final AtomicReference<CountDownLatch> startGate = new
AtomicReference<>();
+ private final AtomicReference<CountDownLatch> suspendEntered = new
AtomicReference<>();
+ private final AtomicReference<CountDownLatch> suspendGate = new
AtomicReference<>();
@Override
protected Endpoint createEndpoint(String uri, String remaining,
Map<String, Object> parameters) {
return new ProbeEndpoint(uri, this);
}
+
+ void releaseGates() {
+ List.of(startGate, suspendGate).forEach(gate -> {
+ CountDownLatch latch = gate.get();
+ if (latch != null) {
+ latch.countDown();
+ }
+ });
+ }
}
private static final class ProbeEndpoint extends DefaultEndpoint {
@@ -295,6 +372,19 @@ public class MasterConsumerLeadershipTest {
super.doStop();
component.stopped.incrementAndGet();
}
+
+ @Override
+ protected void doSuspend() throws Exception {
+ super.doSuspend();
+ CountDownLatch entered = component.suspendEntered.get();
+ if (entered != null) {
+ entered.countDown();
+ }
+ CountDownLatch gate = component.suspendGate.get();
+ if (gate != null) {
+ gate.await();
+ }
+ }
}
// ************************************
diff --git
a/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/SimpleMessageListenerContainer.java
b/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/SimpleMessageListenerContainer.java
index 64657e6f2a64..e21f401359b8 100644
---
a/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/SimpleMessageListenerContainer.java
+++
b/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/consumer/SimpleMessageListenerContainer.java
@@ -271,8 +271,11 @@ public class SimpleMessageListenerContainer extends
ServiceSupport
endpoint.getCamelContext().getExecutorServiceManager().shutdown(recoverPool);
recoverPool = null;
}
- if (recoverFuture != null && recoverTask != null &&
recoverTask.isRunning()) {
- recoverFuture.cancel(true);
+ if (recoverTask != null && recoverTask.isRunning()) {
+ // cancelled through the task and not through its future, so the
task also leaves the
+ // TaskManagerRegistry. Only a run of the task removes it from
there, and once the schedule
+ // is cancelled no run is coming
+ recoverTask.cancel(true);
recoverTask = null;
recoverFuture = null;
}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/support/task/task/BackgroundTaskTest.java
b/core/camel-core/src/test/java/org/apache/camel/support/task/task/BackgroundTaskTest.java
index 5d30ef7293a5..039405fb07bf 100644
---
a/core/camel-core/src/test/java/org/apache/camel/support/task/task/BackgroundTaskTest.java
+++
b/core/camel-core/src/test/java/org/apache/camel/support/task/task/BackgroundTaskTest.java
@@ -22,8 +22,10 @@ import java.util.concurrent.Future;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.TimeUnit;
+import org.apache.camel.support.PluginHelper;
import org.apache.camel.support.task.BackgroundTask;
import org.apache.camel.support.task.Task;
+import org.apache.camel.support.task.TaskManagerRegistry;
import org.apache.camel.support.task.Tasks;
import org.apache.camel.support.task.budget.Budgets;
import org.junit.jupiter.api.DisplayName;
@@ -268,4 +270,100 @@ public class BackgroundTaskTest extends TaskTestSupport {
executor.shutdownNow();
}
}
+
+ @DisplayName("Test that a cancelled task is unscheduled and leaves the
task registry")
+ @Test
+ @Timeout(20)
+ void testCancelUnschedulesAndDeregisters() {
+ ScheduledExecutorService executor =
Executors.newSingleThreadScheduledExecutor();
+ try {
+ BackgroundTask task = Tasks.backgroundTask()
+ .withScheduledExecutor(executor)
+ .withBudget(Budgets.iterationTimeBudget()
+ .withInterval(Duration.ofMillis(100))
+ .withInitialDelay(Duration.ZERO)
+ .withUnlimitedDuration()
+ .build())
+ .withName("cancelled")
+ .build();
+
+ TaskManagerRegistry registry =
PluginHelper.getTaskManagerRegistry(camelContext.getCamelContextExtension());
+ Future<?> future = task.schedule(camelContext,
this::booleanSupplier);
+ await().atMost(5, TimeUnit.SECONDS).until(() ->
registry.getTasks().contains(task));
+
+ task.cancel(false);
+
+ assertTrue(future.isCancelled(), "A cancelled task should not stay
scheduled");
+ // a run that had already started may still have re-added itself,
it then removes itself again
+ await().atMost(5, TimeUnit.SECONDS).untilAsserted(() ->
assertFalse(registry.getTasks().contains(task),
+ "A cancelled task should not stay in the task registry"));
+ assertEquals(Task.Status.Inactive, task.getStatus());
+ assertFalse(task.isRunning(), "A cancelled task should not report
itself as running");
+
+ int attempts = taskCount.intValue();
+ await().pollDelay(1, TimeUnit.SECONDS).atMost(5,
TimeUnit.SECONDS).untilAsserted(
+ () -> assertEquals(attempts, taskCount.intValue(), "A
cancelled task should not run again"));
+ } finally {
+ executor.shutdownNow();
+ }
+ }
+
+ @DisplayName("Test that cancelling a task before its first run leaves
nothing behind")
+ @Test
+ @Timeout(20)
+ void testCancelBeforeTheFirstRun() {
+ ScheduledExecutorService executor =
Executors.newSingleThreadScheduledExecutor();
+ try {
+ BackgroundTask task = Tasks.backgroundTask()
+ .withScheduledExecutor(executor)
+ .withBudget(Budgets.iterationTimeBudget()
+ .withInterval(Duration.ofMillis(100))
+ // long enough that the cancel below lands before
the first run
+ .withInitialDelay(Duration.ofSeconds(3))
+ .withUnlimitedDuration()
+ .build())
+ .withName("cancelled-before-first-run")
+ .build();
+
+ TaskManagerRegistry registry =
PluginHelper.getTaskManagerRegistry(camelContext.getCamelContextExtension());
+ Future<?> future = task.schedule(camelContext,
this::booleanSupplier);
+
+ task.cancel(false);
+
+ assertTrue(future.isCancelled(), "A cancelled task should not stay
scheduled");
+ assertFalse(registry.getTasks().contains(task), "A cancelled task
should not stay in the task registry");
+ await().pollDelay(1, TimeUnit.SECONDS).atMost(10,
TimeUnit.SECONDS).untilAsserted(() -> assertEquals(0,
+ taskCount.intValue(), "The supplier of a task cancelled
before its first run should never run"));
+ } finally {
+ executor.shutdownNow();
+ }
+ }
+
+ @DisplayName("Test that cancelling a task that already completed keeps the
outcome of its last run")
+ @Test
+ @Timeout(20)
+ void testCancelKeepsTheOutcomeOfACompletedTask() {
+ ScheduledExecutorService executor =
Executors.newSingleThreadScheduledExecutor();
+ try {
+ BackgroundTask task = Tasks.backgroundTask()
+ .withScheduledExecutor(executor)
+ .withBudget(Budgets.iterationTimeBudget()
+ .withInterval(Duration.ofMillis(100))
+ .withInitialDelay(Duration.ZERO)
+ .withUnlimitedDuration()
+ .build())
+ .withName("completed-then-cancelled")
+ .build();
+
+ task.schedule(camelContext, () -> true);
+ await().atMost(5, TimeUnit.SECONDS).until(() -> task.getStatus()
== Task.Status.Completed);
+
+ // a caller that cancels defensively must not undo the success of
the task
+ task.cancel(false);
+
+ assertEquals(Task.Status.Completed, task.getStatus());
+ } finally {
+ executor.shutdownNow();
+ }
+ }
}
diff --git
a/core/camel-support/src/main/java/org/apache/camel/support/task/BackgroundTask.java
b/core/camel-support/src/main/java/org/apache/camel/support/task/BackgroundTask.java
index 5308be7b9714..41fade98bf95 100644
---
a/core/camel-support/src/main/java/org/apache/camel/support/task/BackgroundTask.java
+++
b/core/camel-support/src/main/java/org/apache/camel/support/task/BackgroundTask.java
@@ -84,6 +84,8 @@ public class BackgroundTask extends AbstractTask implements
BlockingTask {
private final AtomicBoolean completed = new AtomicBoolean();
// only set when scheduled via schedule(), run() cancels the future it
owns itself
private final AtomicReference<Future<?>> scheduledFuture = new
AtomicReference<>();
+ // the context the schedule was made with, so cancel() can deregister
without being handed it again
+ private final AtomicReference<CamelContext> scheduledContext = new
AtomicReference<>();
private volatile boolean registeredByRun;
private volatile boolean attempting;
@@ -97,7 +99,7 @@ public class BackgroundTask extends AbstractTask implements
BlockingTask {
LOG.trace("Current latch value: {}", latch.getCount());
if (latch.getCount() == 0) {
// the task is done and every further run is a no-op, so stop
being rescheduled
- unschedule();
+ unschedule(false);
return;
}
@@ -106,6 +108,12 @@ public class BackgroundTask extends AbstractTask
implements BlockingTask {
registry =
PluginHelper.getTaskManagerRegistry(camelContext.getCamelContextExtension());
if (!registeredByRun) {
registry.addTask(this);
+ if (latch.getCount() == 0) {
+ // cancelled while this run was starting up, so undo the
registration just made
+ registry.removeTask(this);
+ unschedule(false);
+ return;
+ }
}
}
if (!budget.next()) {
@@ -116,7 +124,7 @@ public class BackgroundTask extends AbstractTask implements
BlockingTask {
registry.removeTask(this);
}
latch.countDown();
- unschedule();
+ unschedule(false);
return;
}
@@ -132,7 +140,7 @@ public class BackgroundTask extends AbstractTask implements
BlockingTask {
registry.removeTask(this);
}
latch.countDown();
- unschedule();
+ unschedule(false);
LOG.trace("Task {} succeeded and the current task is
unscheduled: {}", getName(), latch.getCount());
}
} catch (Exception e) {
@@ -163,22 +171,61 @@ public class BackgroundTask extends AbstractTask
implements BlockingTask {
running.set(true);
Future<?> future = service.scheduleWithFixedDelay(() ->
runTaskWrapper(camelContext, supplier),
budget.initialDelay(), budget.interval(),
TimeUnit.MILLISECONDS);
+ scheduledContext.set(camelContext);
scheduledFuture.set(future);
if (latch.getCount() == 0) {
// the task already finished before the future was published, so
it could not unschedule itself
- unschedule();
+ unschedule(false);
}
return future;
}
+ /**
+ * Cancels a task scheduled with {@link #schedule(CamelContext,
BooleanSupplier)} that is no longer needed, and
+ * removes it from the {@link TaskManagerRegistry}. A scheduled task
deregisters itself from one of its runs, which
+ * is not going to happen once the schedule is cancelled, so cancelling
the returned {@link Future} directly leaves
+ * the task behind in the registry.
+ * <p/>
+ * This does not wait for an attempt that is already running: with {@code
mayInterruptIfRunning} false, a supplier
+ * call that is in progress runs to completion after this method returns.
{@link #isRunning()} answers for the
+ * schedule and turns false here even then, so {@link #isAttempting()} is
the one to ask whether an attempt is still
+ * in flight.
+ * <p/>
+ * A task that already completed, failed or exhausted its budget keeps the
outcome of its last run. Only the
+ * schedule of a task that is still {@link Status#Active} is cancelled,
which turns it {@link Status#Inactive}.
+ *
+ * @param mayInterruptIfRunning whether the thread of an attempt that is
currently running should be interrupted
+ */
+ public void cancel(boolean mayInterruptIfRunning) {
+ // any run that has not started yet becomes a no-op
+ latch.countDown();
+ unschedule(mayInterruptIfRunning);
+ if (status == Status.Active) {
+ status = Status.Inactive;
+ completed.set(false);
+ }
+ deregister();
+ running.set(false);
+ }
+
/**
* Cancels the repeating schedule created by {@link
#schedule(CamelContext, BooleanSupplier)}, so a task that has
* nothing left to do does not keep occupying the scheduler for the
lifetime of its executor.
*/
- private void unschedule() {
+ private void unschedule(boolean mayInterruptIfRunning) {
Future<?> future = scheduledFuture.getAndSet(null);
if (future != null) {
- future.cancel(false);
+ future.cancel(mayInterruptIfRunning);
+ }
+ }
+
+ private void deregister() {
+ CamelContext context = scheduledContext.getAndSet(null);
+ if (context != null) {
+ TaskManagerRegistry registry =
PluginHelper.getTaskManagerRegistry(context.getCamelContextExtension());
+ if (registry != null) {
+ registry.removeTask(this);
+ }
}
}