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 032c0097a916 CAMEL-24583: camel-master - do not start the delegated 
consumer after leadership is lost (CAMEL-24584) (#26052)
032c0097a916 is described below

commit 032c0097a91629bc8caaec63a8ce2028c1740c40
Author: Claus Ibsen <[email protected]>
AuthorDate: Thu Sep 3 07:07:06 2026 +0200

    CAMEL-24583: camel-master - do not start the delegated consumer after 
leadership is lost (CAMEL-24584) (#26052)
    
    MasterConsumer started the delegated consumer from a BackgroundTask 
scheduled after the
    leadership-taken event. A leadership-lost event arriving during that window 
was dropped
    because delegatedConsumer was still null, leaving the consumer running on a 
non-leader node.
    Leadership is now tracked under the consumer lock; the scheduled task 
re-checks it before
    starting, the lost event is dispatched unconditionally to cancel a pending 
start, the delegate
    is created off the lock and published only if leadership still holds, and 
delegatedConsumer is
    published only after a successful start. backOffMaxAttempts now correctly 
bounds the attempts.
    
    Also fixes CAMEL-24584: BackgroundTask.schedule now cancels its repeating 
schedule once the
    task has completed or run out of budget, instead of re-running as a no-op 
for the life of the
    executor.
    
    Closes #26028
    
    (cherry picked from commit f899b39d62ac77d67d63bbc4534d9054bdff58f7)
    
    Co-authored-by: henrik242 <[email protected]>
---
 .../camel/component/master/MasterConsumer.java     | 212 +++++++++---
 .../master/MasterConsumerLeadershipTest.java       | 369 +++++++++++++++++++++
 .../support/task/task/BackgroundTaskTest.java      |  60 ++++
 .../apache/camel/support/task/BackgroundTask.java  |  28 +-
 .../ROOT/pages/camel-4x-upgrade-guide-4_22.adoc    |  15 +
 5 files changed, 632 insertions(+), 52 deletions(-)

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 5dfebcb3e6de..f4606dd96841 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
@@ -17,7 +17,9 @@
 package org.apache.camel.component.master;
 
 import java.time.Duration;
+import java.util.concurrent.Future;
 import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.atomic.AtomicReference;
 
 import org.apache.camel.Consumer;
 import org.apache.camel.Endpoint;
@@ -59,6 +61,10 @@ 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
+    private boolean leadershipTaken;
+    private Future<?> leaderTaskFuture;
 
     public MasterConsumer(MasterEndpoint masterEndpoint, Processor processor, 
CamelClusterService clusterService) {
         super(masterEndpoint, processor);
@@ -103,6 +109,15 @@ 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);
+
+        // 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
+
         if (view != null) {
             view.removeEventListener(leadershipListener);
             clusterService.releaseView(view);
@@ -146,65 +161,126 @@ public class MasterConsumer extends DefaultConsumer 
implements ResumeAware<Resum
                         
.withInterval(Duration.ofMillis(masterEndpoint.getComponent().getBackOffDelay()))
                         .withInitialDelay(Duration.ofSeconds(1))
                         
.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
+                        .withUnlimitedDuration()
                         .build())
                 .withName("Leadership")
                 .build();
     }
 
-    private void onLeadershipTaken() throws Exception {
+    private void onLeadershipTaken() {
         lock.lock();
         try {
             if (!isRunAllowed()) {
                 return;
             }
 
-            if (delegatedConsumer != null) {
+            leadershipTaken = true;
+
+            if (delegatedConsumer != null || isStartPending()) {
                 return;
             }
 
-            final BackgroundTask leaderTask = createTask();
-            leaderTask.schedule(getEndpoint().getCamelContext(), () -> {
-                if (!isRunAllowed()) {
-                    return false;
-                }
-                LOG.info("Leadership taken. Attempt #{} to start consumer: 
{}", leaderTask.iteration(), delegatedEndpoint);
-
-                Exception cause = null;
-                try {
-                    if (delegatedConsumer == null) {
-                        delegatedConsumer = 
delegatedEndpoint.createConsumer(processor);
-                        if (delegatedConsumer instanceof StartupListener) {
-                            
getEndpoint().getCamelContext().addStartupListener((StartupListener) 
delegatedConsumer);
-                        }
-                        if (delegatedConsumer instanceof ResumeAware 
resumeAwareConsumer && resumeStrategy != null) {
-                            LOG.debug("Setting up the resume adapter for the 
resume strategy in consumer");
-                            ResumeAdapter resumeAdapter
-                                    = 
AdapterHelper.eval(clusterService.getCamelContext(), resumeAwareConsumer,
-                                            resumeStrategy);
-                            resumeStrategy.setAdapter(resumeAdapter);
-
-                            LOG.debug("Setting up the resume strategy for 
consumer");
-                            
resumeAwareConsumer.setResumeStrategy(resumeStrategy);
-                        }
-                    }
-                    ServiceHelper.startService(delegatedEndpoint, 
delegatedConsumer);
+            // a task from a previous leadership term may still be scheduled, 
drop it
+            cancelLeaderTask(false);
+
+            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<>();
+            leaderTaskFuture = task.schedule(getEndpoint().getCamelContext(), 
() -> startDelegatedConsumer(task, attempt));
+        } finally {
+            lock.unlock();
+        }
+    }
+
+    private boolean startDelegatedConsumer(BackgroundTask task, 
AtomicReference<Consumer> attempt) {
+        lock.lock();
+        try {
+            if (!isRunAllowed()) {
+                return false;
+            }
+
+            if (!leadershipTaken) {
+                // leadership was lost while this start was pending. Starting 
now would run the consumer on a
+                // node that is not the leader, and no further leadership 
event is coming to stop it again
+                LOG.debug("Leadership lost while the start was pending. Not 
starting consumer: {}", delegatedEndpoint);
+                return true; // no more attempts
+            }
+
+            if (delegatedConsumer != null) {
+                return true; // no more attempts
+            }
+        } finally {
+            lock.unlock();
+        }
+
+        LOG.info("Leadership taken. Attempt #{} to start consumer: {}", 
task.iteration(), delegatedEndpoint);
 
-                } catch (Exception e) {
-                    cause = e;
+        // 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
+        Consumer consumer = attempt.get();
+        Exception cause = null;
+        try {
+            if (consumer == null) {
+                consumer = delegatedEndpoint.createConsumer(processor);
+                // held for the attempts of this task, so the startup listener 
and the resume strategy are
+                // wired once and a retry only starts the consumer again
+                attempt.set(consumer);
+                if (consumer instanceof StartupListener startupListener) {
+                    
getEndpoint().getCamelContext().addStartupListener(startupListener);
+                }
+                if (consumer instanceof ResumeAware resumeAwareConsumer && 
resumeStrategy != null) {
+                    LOG.debug("Setting up the resume adapter for the resume 
strategy in consumer");
+                    ResumeAdapter resumeAdapter
+                            = 
AdapterHelper.eval(clusterService.getCamelContext(), resumeAwareConsumer,
+                                    resumeStrategy);
+                    resumeStrategy.setAdapter(resumeAdapter);
+
+                    LOG.debug("Setting up the resume strategy for consumer");
+                    resumeAwareConsumer.setResumeStrategy(resumeStrategy);
                 }
+            }
+            ServiceHelper.startService(delegatedEndpoint, consumer);
+        } catch (Exception e) {
+            cause = e;
+        }
 
-                if (cause != null) {
-                    String message = "Leadership taken. Attempt #" + 
leaderTask.iteration()
-                                     + " failed to start consumer due to: " + 
cause.getMessage();
-                    getExceptionHandler().handleException(message, cause);
-                    // make the task runner aware of the exception (will retry)
-                    throw new TaskRunFailureException(message, cause);
+        lock.lock();
+        try {
+            if (cause != null) {
+                // the consumer is kept for the next attempt. It is not 
stopped here: a consumer that failed to
+                // start was already stopped by its own start(), and shutting 
it down would also shut down the
+                // processor of the route, which the next attempt and this 
consumer still need
+                String message = "Leadership taken. Attempt #" + 
task.iteration()
+                                 + " failed to start consumer due to: " + 
cause.getMessage();
+                getExceptionHandler().handleException(message, cause);
+                int maxAttempts = 
masterEndpoint.getComponent().getBackOffMaxAttempts();
+                if (maxAttempts > 0 && task.iteration() >= maxAttempts) {
+                    LOG.error("Leadership taken. Giving up after {} attempts 
to start consumer: {}."
+                              + " This node holds the leadership but is not 
consuming, until the leadership changes again.",
+                            task.iteration(), delegatedEndpoint);
                 }
+                // make the task runner aware of the exception (will retry)
+                throw new TaskRunFailureException(message, cause);
+            }
 
-                LOG.info("Leadership taken. Attempt #{} success. Consumer 
started: {}", leaderTask.iteration(),
-                        delegatedEndpoint);
+            if (!leadershipTaken || !isRunAllowed()) {
+                // the leadership went away while the consumer was starting, 
so stop what was just started
+                // instead of publishing it. No leadership event is going to 
do it, delegatedConsumer is unset
+                LOG.info("Leadership lost while the consumer was starting. 
Stopping consumer: {}", delegatedEndpoint);
+                ServiceHelper.stopAndShutdownServices(consumer, 
delegatedEndpoint);
+                attempt.set(null);
                 return true; // no more attempts
-            });
+            }
+
+            delegatedConsumer = consumer;
+            LOG.info("Leadership taken. Attempt #{} success. Consumer started: 
{}", task.iteration(),
+                    delegatedEndpoint);
+            // release the task, a later leadership term schedules a new one
+            cancelLeaderTask(false);
+            return true; // no more attempts
         } finally {
             lock.unlock();
         }
@@ -213,6 +289,15 @@ public class MasterConsumer extends DefaultConsumer 
implements ResumeAware<Resum
     private void onLeadershipLost() {
         lock.lock();
         try {
+            leadershipTaken = false;
+            // a start scheduled by the leadership taken event may not have 
run yet, cancel it so it
+            // cannot start the consumer on a node that is no longer the leader
+            cancelLeaderTask(false);
+
+            if (delegatedConsumer == null) {
+                return;
+            }
+
             LOG.debug("Leadership lost. Stopping consumer: {}", 
delegatedEndpoint);
             try {
                 ServiceHelper.stopAndShutdownServices(delegatedConsumer, 
delegatedEndpoint);
@@ -225,6 +310,17 @@ public class MasterConsumer extends DefaultConsumer 
implements ResumeAware<Resum
         }
     }
 
+    private boolean isStartPending() {
+        return leaderTaskFuture != null && !leaderTaskFuture.isDone();
+    }
+
+    private void cancelLeaderTask(boolean mayInterruptIfRunning) {
+        if (leaderTaskFuture != null) {
+            leaderTaskFuture.cancel(mayInterruptIfRunning);
+            leaderTaskFuture = null;
+        }
+    }
+
     // **************************************
     // Listener
     // **************************************
@@ -233,22 +329,38 @@ 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
                 return;
             }
 
-            if (view.getLocalMember().isLeader()) {
-                try {
-                    onLeadershipTaken();
-                } catch (Exception e) {
-                    getExceptionHandler().handleException("Error starting 
consumer while taking leadership", e);
+            lock.lock();
+            try {
+                if (!isRunAllowed()) {
+                    return;
                 }
-            } else if (delegatedConsumer != null) {
-                try {
-                    onLeadershipLost();
-                } catch (Exception e) {
-                    getExceptionHandler()
-                            .handleException("Error stopping consumer while 
loosing leadership. This exception is ignored.", e);
+
+                // the leadership is read under the same lock that applies it, 
so that two events
+                // dispatched concurrently cannot be applied in the wrong order
+                if (view.getLocalMember().isLeader()) {
+                    try {
+                        onLeadershipTaken();
+                    } catch (Exception e) {
+                        getExceptionHandler().handleException("Error starting 
consumer while taking leadership", e);
+                    }
+                } else {
+                    // dispatched even when there is no consumer yet, as a 
start may be pending
+                    try {
+                        onLeadershipLost();
+                    } catch (Exception e) {
+                        getExceptionHandler()
+                                .handleException("Error stopping consumer 
while loosing leadership. This exception is ignored.",
+                                        e);
+                    }
                 }
+            } finally {
+                lock.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
new file mode 100644
index 000000000000..59a3a31f95bb
--- /dev/null
+++ 
b/components/camel-master/src/test/java/org/apache/camel/component/master/MasterConsumerLeadershipTest.java
@@ -0,0 +1,369 @@
+/*
+ * 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.component.master;
+
+import java.util.List;
+import java.util.Map;
+import java.util.Optional;
+import java.util.UUID;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.Consumer;
+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.cluster.CamelClusterMember;
+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.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 org.junit.jupiter.api.Timeout;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * Verifies that the delegated consumer only ever runs while this node holds 
the leadership, also when the leadership
+ * changes while the start of the delegated consumer is still pending.
+ */
+public class MasterConsumerLeadershipTest {
+
+    private DefaultCamelContext context;
+    private TestClusterService clusterService;
+    private ProbeComponent probe;
+
+    @BeforeEach
+    void setUp() throws Exception {
+        probe = new ProbeComponent();
+        clusterService = new TestClusterService();
+
+        context = new DefaultCamelContext();
+        context.disableJMX();
+        context.addService(clusterService);
+        context.addComponent("probe", probe);
+
+        MasterComponent master = context.getComponent("master", 
MasterComponent.class);
+        // keep the retries short so an exhausted start does not dominate the 
test time
+        master.setBackOffDelay(200);
+        master.setBackOffMaxAttempts(2);
+
+        context.addRoutes(new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("master:ns:probe:test").routeId("master-route").to("mock:result");
+            }
+        });
+
+        context.start();
+    }
+
+    @AfterEach
+    void tearDown() {
+        if (context != null) {
+            context.stop();
+        }
+    }
+
+    @Test
+    @Timeout(60)
+    void testLeadershipLostWhilePendingStartDoesNotStartConsumer() {
+        TestClusterView view = clusterService.getTestView();
+
+        view.setLeader(true);
+        // the start is scheduled with an initial delay, so it is still 
pending here
+        view.setLeader(false);
+
+        // outlast the initial delay of the start task and verify the consumer 
was never even created
+        await().during(3, TimeUnit.SECONDS).atMost(10, TimeUnit.SECONDS)
+                .untilAsserted(() -> assertEquals(0, probe.created.get()));
+
+        // taking the leadership again must still work, which also proves the 
events did reach the consumer
+        view.setLeader(true);
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, probe.started.get()));
+        assertEquals(1, probe.created.get(), "The cancelled start must not 
have created a second consumer");
+    }
+
+    @Test
+    @Timeout(60)
+    void testLeadershipTakenStartsConsumerAndLostStopsIt() {
+        TestClusterView view = clusterService.getTestView();
+
+        view.setLeader(true);
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, probe.started.get()));
+
+        view.setLeader(false);
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, probe.stopped.get()));
+    }
+
+    @Test
+    @Timeout(60)
+    void testConsumerIsRestartedWhenLeadershipFlapsAfterASuccessfulStart() {
+        TestClusterView view = clusterService.getTestView();
+
+        view.setLeader(true);
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, probe.started.get()));
+
+        // the membership flap seen in production: the consumer must stop and 
then come back
+        view.setLeader(false);
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, probe.stopped.get()));
+
+        view.setLeader(true);
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(2, probe.started.get()));
+        assertEquals(1, probe.stopped.get());
+    }
+
+    @Test
+    @Timeout(60)
+    void testLeadershipLostWhileStartIsInProgressStopsTheConsumer() throws 
Exception {
+        TestClusterView view = clusterService.getTestView();
+        CountDownLatch startGate = new CountDownLatch(1);
+        probe.startGate.set(startGate);
+
+        view.setLeader(true);
+        // wait until the start is running and blocked inside the delegated 
consumer
+        await().atMost(10, TimeUnit.SECONDS).until(() -> 
probe.startAttempts.get() == 1);
+
+        // the leadership is lost while the start is in progress, this must 
not be able to interleave
+        Thread loser = new Thread(() -> view.setLeader(false), 
"leadership-lost");
+        loser.start();
+        startGate.countDown();
+        loser.join(TimeUnit.SECONDS.toMillis(20));
+
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> {
+            assertEquals(1, probe.started.get());
+            assertEquals(1, probe.stopped.get(), "The consumer started on a 
node that lost the leadership must be stopped");
+        });
+    }
+
+    @Test
+    @Timeout(60)
+    void testRepeatedLeadershipTakenStartsOnlyOneConsumer() {
+        TestClusterView view = clusterService.getTestView();
+
+        view.setLeader(true);
+        view.setLeader(true);
+
+        await().atMost(10, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, probe.started.get()));
+        await().during(2, TimeUnit.SECONDS).atMost(10, TimeUnit.SECONDS)
+                .untilAsserted(() -> assertEquals(1, probe.created.get()));
+    }
+
+    @Test
+    @Timeout(60)
+    void testStoppingTheConsumerCancelsAPendingStart() {
+        TestClusterView view = clusterService.getTestView();
+
+        view.setLeader(true);
+        // the start is still pending, stopping must cancel it instead of 
letting it start afterwards
+        context.stop();
+
+        await().during(3, TimeUnit.SECONDS).atMost(10, TimeUnit.SECONDS)
+                .untilAsserted(() -> assertEquals(0, probe.created.get()));
+    }
+
+    @Test
+    @Timeout(60)
+    void 
testConsumerStartsAfterLeadershipIsTakenAgainWhenAnEarlierStartFailed() {
+        TestClusterView view = clusterService.getTestView();
+
+        probe.failStart.set(true);
+        view.setLeader(true);
+
+        // every start attempt fails, then the task runs out of budget and 
stops attempting
+        await().atMost(20, TimeUnit.SECONDS).until(() -> 
probe.startAttempts.get() == 2);
+        await().during(1, TimeUnit.SECONDS).atMost(20, TimeUnit.SECONDS)
+                .untilAsserted(() -> assertEquals(2, 
probe.startAttempts.get()));
+        assertEquals(0, probe.started.get());
+
+        // a failed start must not leave state behind that makes a later 
leadership event a no-op,
+        // not even without an intervening leadership lost event
+        probe.failStart.set(false);
+        view.setLeader(true);
+
+        await().atMost(20, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, probe.started.get()));
+    }
+
+    @Test
+    @Timeout(60)
+    void testAllConfiguredStartAttemptsAreMade() {
+        TestClusterView view = clusterService.getTestView();
+        MasterComponent master = context.getComponent("master", 
MasterComponent.class);
+        // the attempts have to outlast the default 5s duration of the 
iteration time budget
+        master.setBackOffDelay(3000);
+        master.setBackOffMaxAttempts(3);
+
+        probe.failStart.set(true);
+        view.setLeader(true);
+
+        // every configured attempt must be made, the task must not end on a 
time budget of its own
+        await().atMost(30, TimeUnit.SECONDS).until(() -> 
probe.startAttempts.get() == 3);
+        assertEquals(0, probe.started.get());
+    }
+
+    // ************************************
+    // Delegated endpoint under observation
+    // ************************************
+
+    private static final class ProbeComponent extends DefaultComponent {
+        private final AtomicInteger created = new AtomicInteger();
+        private final AtomicInteger startAttempts = new AtomicInteger();
+        private final AtomicInteger started = new AtomicInteger();
+        private final AtomicInteger stopped = new AtomicInteger();
+        private final AtomicBoolean failStart = new AtomicBoolean();
+        private final AtomicReference<CountDownLatch> startGate = new 
AtomicReference<>();
+
+        @Override
+        protected Endpoint createEndpoint(String uri, String remaining, 
Map<String, Object> parameters) {
+            return new ProbeEndpoint(uri, this);
+        }
+    }
+
+    private static final class ProbeEndpoint extends DefaultEndpoint {
+        private final ProbeComponent component;
+
+        ProbeEndpoint(String uri, ProbeComponent component) {
+            super(uri, component);
+            this.component = component;
+        }
+
+        @Override
+        public Producer createProducer() {
+            throw new UnsupportedOperationException("Cannot produce from this 
endpoint");
+        }
+
+        @Override
+        public Consumer createConsumer(Processor processor) {
+            component.created.incrementAndGet();
+            return new ProbeConsumer(this, processor, component);
+        }
+
+        @Override
+        public boolean isSingleton() {
+            return true;
+        }
+    }
+
+    private static final class ProbeConsumer extends DefaultConsumer {
+        private final ProbeComponent component;
+
+        ProbeConsumer(Endpoint endpoint, Processor processor, ProbeComponent 
component) {
+            super(endpoint, processor);
+            this.component = component;
+        }
+
+        @Override
+        protected void doStart() throws Exception {
+            super.doStart();
+            // counted before the failure flag is read, so a test can await an 
attempt that has made its decision
+            component.startAttempts.incrementAndGet();
+            CountDownLatch gate = component.startGate.get();
+            if (gate != null) {
+                gate.await();
+            }
+            if (component.failStart.get()) {
+                throw new IllegalStateException("Simulated failure to start 
the delegated consumer");
+            }
+            component.started.incrementAndGet();
+        }
+
+        @Override
+        protected void doStop() throws Exception {
+            super.doStop();
+            component.stopped.incrementAndGet();
+        }
+    }
+
+    // ************************************
+    // Cluster with a leadership we control
+    // ************************************
+
+    private static final class TestClusterService extends 
AbstractCamelClusterService<TestClusterView> {
+        private volatile TestClusterView view;
+
+        TestClusterService() {
+            super("test-cluster-service");
+        }
+
+        TestClusterView getTestView() {
+            return view;
+        }
+
+        @Override
+        protected TestClusterView createView(String namespace) {
+            view = new TestClusterView(this, namespace);
+            return view;
+        }
+    }
+
+    private static final class TestClusterView extends 
AbstractCamelClusterView {
+        private final TestClusterMember localMember = new TestClusterMember();
+
+        TestClusterView(TestClusterService clusterService, String namespace) {
+            super(clusterService, namespace);
+        }
+
+        void setLeader(boolean leader) {
+            localMember.leader = leader;
+            fireLeadershipChangedEvent(leader ? localMember : null);
+        }
+
+        @Override
+        public Optional<CamelClusterMember> getLeader() {
+            return localMember.isLeader() ? Optional.of(localMember) : 
Optional.empty();
+        }
+
+        @Override
+        public CamelClusterMember getLocalMember() {
+            return localMember;
+        }
+
+        @Override
+        public List<CamelClusterMember> getMembers() {
+            return List.of(localMember);
+        }
+    }
+
+    private static final class TestClusterMember implements CamelClusterMember 
{
+        private final String id = UUID.randomUUID().toString();
+        private volatile boolean leader;
+
+        @Override
+        public boolean isLeader() {
+            return leader;
+        }
+
+        @Override
+        public boolean isLocal() {
+            return true;
+        }
+
+        @Override
+        public String getId() {
+            return id;
+        }
+    }
+}
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 e61ff91c32f7..5d30ef7293a5 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
@@ -18,14 +18,19 @@ package org.apache.camel.support.task.task;
 
 import java.time.Duration;
 import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.TimeUnit;
 
 import org.apache.camel.support.task.BackgroundTask;
+import org.apache.camel.support.task.Task;
 import org.apache.camel.support.task.Tasks;
 import org.apache.camel.support.task.budget.Budgets;
 import org.junit.jupiter.api.DisplayName;
 import org.junit.jupiter.api.Test;
 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.assertNotNull;
@@ -208,4 +213,59 @@ public class BackgroundTaskTest extends TaskTestSupport {
         assertTrue(duration.getSeconds() <= 5);
         assertFalse(completed, "The task did not complete because of timeout, 
the return should be false");
     }
+
+    @DisplayName("Test that a scheduled task is unscheduled once it has 
completed")
+    @Test
+    @Timeout(10)
+    void testScheduleStopsWhenCompleted() {
+        ScheduledExecutorService executor = 
Executors.newSingleThreadScheduledExecutor();
+        try {
+            BackgroundTask task = Tasks.backgroundTask()
+                    .withScheduledExecutor(executor)
+                    .withBudget(Budgets.iterationTimeBudget()
+                            .withInterval(Duration.ofMillis(100))
+                            .withInitialDelay(Duration.ZERO)
+                            .withMaxIterations(maxIterations)
+                            .build())
+                    .build();
+
+            Future<?> future = task.schedule(camelContext, () -> {
+                taskCount.increment();
+                return true;
+            });
+
+            await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> 
assertTrue(future.isCancelled(),
+                    "A completed task should not stay scheduled"));
+            assertEquals(1, taskCount.intValue(), "The supplier should have 
run exactly once");
+            assertEquals(Task.Status.Completed, task.getStatus());
+        } finally {
+            executor.shutdownNow();
+        }
+    }
+
+    @DisplayName("Test that a scheduled task is unscheduled once it runs out 
of budget")
+    @Test
+    @Timeout(10)
+    void testScheduleStopsWhenExhausted() {
+        ScheduledExecutorService executor = 
Executors.newSingleThreadScheduledExecutor();
+        try {
+            BackgroundTask task = Tasks.backgroundTask()
+                    .withScheduledExecutor(executor)
+                    .withBudget(Budgets.iterationTimeBudget()
+                            .withInterval(Duration.ofMillis(100))
+                            .withInitialDelay(Duration.ZERO)
+                            .withMaxIterations(maxIterations)
+                            .build())
+                    .build();
+
+            Future<?> future = task.schedule(camelContext, 
this::booleanSupplier);
+
+            await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> 
assertTrue(future.isCancelled(),
+                    "An exhausted task should not stay scheduled"));
+            assertEquals(maxIterations, taskCount.intValue());
+            assertEquals(Task.Status.Exhausted, 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 c8776c352898..5308be7b9714 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
@@ -24,6 +24,7 @@ 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.AtomicReference;
 import java.util.function.BooleanSupplier;
 
 import org.apache.camel.CamelContext;
@@ -81,6 +82,8 @@ public class BackgroundTask extends AbstractTask implements 
BlockingTask {
     private Duration elapsed = Duration.ZERO;
     private final AtomicBoolean running = new AtomicBoolean();
     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<>();
     private volatile boolean registeredByRun;
     private volatile boolean attempting;
 
@@ -93,6 +96,8 @@ public class BackgroundTask extends AbstractTask implements 
BlockingTask {
     private void runTaskWrapper(CamelContext camelContext, BooleanSupplier 
supplier) {
         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();
             return;
         }
 
@@ -111,6 +116,7 @@ public class BackgroundTask extends AbstractTask implements 
BlockingTask {
                 registry.removeTask(this);
             }
             latch.countDown();
+            unschedule();
             return;
         }
 
@@ -126,6 +132,7 @@ public class BackgroundTask extends AbstractTask implements 
BlockingTask {
                     registry.removeTask(this);
                 }
                 latch.countDown();
+                unschedule();
                 LOG.trace("Task {} succeeded and the current task is 
unscheduled: {}", getName(), latch.getCount());
             }
         } catch (Exception e) {
@@ -154,8 +161,25 @@ public class BackgroundTask extends AbstractTask 
implements BlockingTask {
      */
     public Future<?> schedule(CamelContext camelContext, BooleanSupplier 
supplier) {
         running.set(true);
-        return service.scheduleWithFixedDelay(() -> 
runTaskWrapper(camelContext, supplier), budget.initialDelay(),
-                budget.interval(), TimeUnit.MILLISECONDS);
+        Future<?> future = service.scheduleWithFixedDelay(() -> 
runTaskWrapper(camelContext, supplier),
+                budget.initialDelay(), budget.interval(), 
TimeUnit.MILLISECONDS);
+        scheduledFuture.set(future);
+        if (latch.getCount() == 0) {
+            // the task already finished before the future was published, so 
it could not unschedule itself
+            unschedule();
+        }
+        return future;
+    }
+
+    /**
+     * 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() {
+        Future<?> future = scheduledFuture.getAndSet(null);
+        if (future != null) {
+            future.cancel(false);
+        }
     }
 
     @Override
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
index c98771e8780f..435153a8a2df 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_22.adoc
@@ -1799,3 +1799,18 @@ they say:
   `CamelGoogleVertexAIStreamOutputMode` header.
 * `jsonMode=true` sets the response MIME type of the request to 
`application/json`, so the model is
   asked to answer with JSON. The default `false` leaves the request untouched.
+
+=== camel-support
+
+`BackgroundTask.schedule` now cancels the repeating schedule it created once 
the task has completed or has
+run out of budget. Previously the returned `Future` stayed armed and the task 
kept being re-run as a no-op
+for the lifetime of the executor. Callers that inspect the returned `Future` 
will see `isCancelled()`
+return `true` after the task is done, where it previously stayed live. Callers 
that already cancel the
+`Future` themselves are unaffected.
+
+=== camel-master
+
+The `backOffMaxAttempts` option now bounds the attempts to start the delegated 
consumer as documented.
+The retry task previously also carried the default five second duration of its 
budget, which ended the
+task before the second attempt for any `backOffDelay` at or above the default 
of five seconds. A delegate
+that fails to start is therefore retried for longer than before, up to 
`backOffMaxAttempts` times.

Reply via email to