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 686bd4a40588 CAMEL-25109: camel-core - Scheduled poll consumer: fix 
bugs found in a deep review (#27015)
686bd4a40588 is described below

commit 686bd4a405882a2c85595602090359629d5f2d89
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 22:44:28 2026 +0200

    CAMEL-25109: camel-core - Scheduled poll consumer: fix bugs found in a deep 
review (#27015)
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../impl/ScheduledPollConsumerEdgeCasesTest.java   | 226 +++++++++++++++++++++
 .../DefaultScheduledPollConsumerScheduler.java     |   5 +-
 .../camel/support/ScheduledPollConsumer.java       |  67 ++++--
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |   9 +
 4 files changed, 286 insertions(+), 21 deletions(-)

diff --git 
a/core/camel-core/src/test/java/org/apache/camel/impl/ScheduledPollConsumerEdgeCasesTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/impl/ScheduledPollConsumerEdgeCasesTest.java
new file mode 100644
index 000000000000..102438842916
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/impl/ScheduledPollConsumerEdgeCasesTest.java
@@ -0,0 +1,226 @@
+/*
+ * 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;
+
+import java.util.List;
+import java.util.Map;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.CyclicBarrier;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.health.HealthCheck;
+import org.apache.camel.spi.HttpResponseAware;
+import org.apache.camel.support.DefaultScheduledPollConsumerScheduler;
+import org.junit.jupiter.api.Test;
+
+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.assertTrue;
+
+public class ScheduledPollConsumerEdgeCasesTest extends ContextTestSupport {
+
+    private static class MyConsumer extends MockScheduledPollConsumer {
+
+        private CountDownLatch pollStarted;
+        private CountDownLatch pollContinue;
+
+        MyConsumer(Endpoint endpoint, Exception exceptionToThrowOnPoll) {
+            super(endpoint, exceptionToThrowOnPoll);
+        }
+
+        @Override
+        protected int poll() throws Exception {
+            if (pollStarted != null) {
+                pollStarted.countDown();
+                pollContinue.await(10, TimeUnit.SECONDS);
+            }
+            return super.poll();
+        }
+
+        Throwable lastError() {
+            return getLastError();
+        }
+
+        Map<String, Object> lastErrorDetails() {
+            return getLastErrorDetails();
+        }
+    }
+
+    private static class MyHttpException extends Exception implements 
HttpResponseAware {
+        MyHttpException() {
+            super("Service unavailable");
+        }
+
+        @Override
+        public int getHttpResponseCode() {
+            return 503;
+        }
+
+        @Override
+        public void setHttpResponseCode(int code) {
+            // noop
+        }
+
+        @Override
+        public String getHttpResponseStatus() {
+            return "Service Unavailable";
+        }
+
+        @Override
+        public void setHttpResponseStatus(String status) {
+            // noop
+        }
+    }
+
+    @Test
+    public void testChangedDelayUsedAfterRestart() {
+        MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), 
null);
+        consumer.setStartScheduler(false);
+        consumer.setDelay(20);
+        consumer.start();
+        consumer.stop();
+
+        consumer.setDelay(2000);
+        consumer.setInitialDelay(3000);
+        consumer.start();
+        try {
+            DefaultScheduledPollConsumerScheduler scheduler = 
(DefaultScheduledPollConsumerScheduler) consumer.getScheduler();
+            assertEquals(2000, scheduler.getDelay());
+            assertEquals(3000, scheduler.getInitialDelay());
+        } finally {
+            consumer.stop();
+        }
+    }
+
+    @Test
+    public void testConcurrentUnscheduleTask() throws Exception {
+        MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), 
null);
+        List<Throwable> errors = new CopyOnWriteArrayList<>();
+
+        for (int i = 0; i < 30; i++) {
+            DefaultScheduledPollConsumerScheduler scheduler = new 
DefaultScheduledPollConsumerScheduler();
+            scheduler.setCamelContext(context);
+            scheduler.setConcurrentConsumers(8);
+            scheduler.setPoolSize(8);
+            scheduler.setInitialDelay(0);
+            scheduler.setDelay(1);
+            scheduler.onInit(consumer);
+            CountDownLatch latch = new CountDownLatch(8);
+            CyclicBarrier barrier = new CyclicBarrier(8);
+            // each concurrent task unschedules the task at the same time 
(such as when the repeat count is reached)
+            scheduler.scheduleTask(() -> {
+                try {
+                    barrier.await(10, TimeUnit.SECONDS);
+                    scheduler.unscheduleTask();
+                } catch (Throwable e) {
+                    errors.add(e);
+                } finally {
+                    latch.countDown();
+                }
+            });
+            scheduler.start();
+            scheduler.startScheduler();
+            assertTrue(latch.await(10, TimeUnit.SECONDS));
+            scheduler.stop();
+        }
+
+        assertTrue(errors.isEmpty(), errors.toString());
+    }
+
+    @Test
+    public void testPollRunningWhenStoppedDoesNotUpdateState() throws 
Exception {
+        MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), 
null);
+        consumer.setStartScheduler(false);
+        consumer.start();
+
+        consumer.pollStarted = new CountDownLatch(1);
+        consumer.pollContinue = new CountDownLatch(1);
+        Thread thread = new Thread(consumer::run);
+        thread.start();
+        assertTrue(consumer.pollStarted.await(10, TimeUnit.SECONDS));
+
+        // stop while polling, and then let the poll complete
+        consumer.stop();
+        consumer.pollContinue.countDown();
+        thread.join(10000);
+
+        assertFalse(consumer.isFirstPollDone());
+        assertEquals(0, consumer.getSuccessCounter());
+    }
+
+    @Test
+    public void testErrorThresholdWithoutBackoffMultiplier() {
+        MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), new 
IllegalStateException("Forced"));
+        consumer.setBackoffErrorThreshold(2);
+        consumer.start();
+        try {
+            for (int i = 0; i < 5; i++) {
+                consumer.run();
+            }
+            assertEquals(5, consumer.getErrorCounter());
+        } finally {
+            consumer.stop();
+        }
+    }
+
+    @Test
+    public void testErrorCounterKeptWhenBackoffFinished() {
+        MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), new 
IllegalStateException("Forced"));
+        consumer.setBackoffMultiplier(2);
+        consumer.setBackoffErrorThreshold(1);
+        consumer.start();
+        try {
+            // error
+            consumer.run();
+            assertEquals(1, consumer.getErrorCounter());
+            // backoff (skip)
+            consumer.run();
+            assertEquals(1, consumer.getErrorCounter());
+            // backoff finished, poll again which fails again
+            consumer.run();
+            assertEquals(2, consumer.getErrorCounter());
+            // and backoff again (skip)
+            consumer.run();
+            assertEquals(2, consumer.getErrorCounter());
+        } finally {
+            consumer.stop();
+        }
+    }
+
+    @Test
+    public void testLastErrorDetails() {
+        MyConsumer consumer = new MyConsumer(getMockEndpoint("mock:foo"), new 
MyHttpException());
+        consumer.start();
+        consumer.run();
+        assertEquals(503, 
consumer.lastErrorDetails().get(HealthCheck.HTTP_RESPONSE_CODE));
+
+        // another error without http details
+        consumer.setExceptionToThrowOnPoll(new 
IllegalStateException("Forced"));
+        consumer.run();
+        assertTrue(consumer.lastError() instanceof IllegalStateException);
+        assertNull(consumer.lastErrorDetails());
+
+        // the error is cleared when stopped
+        consumer.stop();
+        assertNull(consumer.lastError());
+        assertNull(consumer.lastErrorDetails());
+    }
+}
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultScheduledPollConsumerScheduler.java
 
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultScheduledPollConsumerScheduler.java
index bf6cefeeb986..75f97ad96cdf 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/support/DefaultScheduledPollConsumerScheduler.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/support/DefaultScheduledPollConsumerScheduler.java
@@ -16,9 +16,9 @@
  */
 package org.apache.camel.support;
 
-import java.util.ArrayList;
 import java.util.List;
 import java.util.Locale;
+import java.util.concurrent.CopyOnWriteArrayList;
 import java.util.concurrent.ScheduledExecutorService;
 import java.util.concurrent.ScheduledFuture;
 import java.util.concurrent.TimeUnit;
@@ -45,7 +45,8 @@ public class DefaultScheduledPollConsumerScheduler extends 
ServiceSupport implem
     private Consumer consumer;
     private ScheduledExecutorService scheduledExecutorService;
     private boolean shutdownExecutor;
-    private final List<ScheduledFuture<?>> futures = new ArrayList<>();
+    // thread-safe as the task may be unscheduled by the (concurrent) polling 
threads (such as with repeatCount)
+    private final List<ScheduledFuture<?>> futures = new 
CopyOnWriteArrayList<>();
     private Runnable task;
     private int poolSize = 1;
     private int concurrentConsumers = 1;
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/support/ScheduledPollConsumer.java
 
b/core/camel-support/src/main/java/org/apache/camel/support/ScheduledPollConsumer.java
index 8c765e9313b0..f37151f5972b 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/support/ScheduledPollConsumer.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/support/ScheduledPollConsumer.java
@@ -81,6 +81,12 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
     private volatile Map<String, Object> lastErrorDetails;
     private final AtomicLong counter = new AtomicLong();
     private volatile boolean firstPollDone;
+    // the error counter when the last backoff finished (the error threshold 
of the backoff counts from this)
+    private final AtomicLong backoffErrorBase = new AtomicLong();
+    // incremented when the consumer is stopped, so a poll that is still 
running does not update the state afterwards
+    private final AtomicLong generation = new AtomicLong();
+    // whether this consumer created the scheduler (and therefore configures 
it)
+    private boolean consumerCreatedScheduler;
     private volatile boolean forceReady;
 
     protected ScheduledPollConsumer(Endpoint endpoint, Processor processor) {
@@ -143,6 +149,7 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
     }
 
     private void doRun() {
+        final long currentGeneration = generation.get();
         if (isSuspended()) {
             LOG.trace("Cannot start to poll: {} as its suspended", 
this.getEndpoint());
             return;
@@ -151,8 +158,9 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
         // should we backoff if its enabled, and either the idle or error 
counter is > the threshold
         if (backoffMultiplier > 0
                 // either idle or error threshold could be not in use, so 
check for that and use MAX_VALUE if not in use
-                && idleCounter.longValue() >= (backoffIdleThreshold > 0 ? 
backoffIdleThreshold : Integer.MAX_VALUE)
-                || errorCounter.longValue() >= (backoffErrorThreshold > 0 ? 
backoffErrorThreshold : Integer.MAX_VALUE)) {
+                && (idleCounter.longValue() >= (backoffIdleThreshold > 0 ? 
backoffIdleThreshold : Integer.MAX_VALUE)
+                        || errorCounter.longValue() - 
backoffErrorBase.longValue() >= (backoffErrorThreshold > 0
+                                ? backoffErrorThreshold : Integer.MAX_VALUE))) 
{
             final int currentBackoffCounter = backoffCounter.incrementAndGet();
             if (currentBackoffCounter < backoffMultiplier) {
                 // yes we should backoff
@@ -165,9 +173,10 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
                 }
                 return;
             } else {
-                // we are finished with backoff so reset counters
+                // we are finished with backoff so reset counters, but keep 
the error counter (the consumer is still
+                // failing until the next poll succeeds) and count the error 
threshold from now on
                 idleCounter.set(0);
-                errorCounter.set(0);
+                backoffErrorBase.set(errorCounter.longValue());
                 backoffCounter.set(0);
                 successCounter.set(0);
                 LOG.trace("doRun() backoff finished, resetting counters.");
@@ -220,14 +229,17 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
                                 retryCounter = -1;
                                 LOG.trace("Greedy polling after processing {} 
messages", polledMessages);
 
-                                // clear any error that might be since we have 
successfully polled, otherwise readiness checks might believe the
-                                // consumer to be unhealthy
-                                errorCounter.set(0);
-                                lastError = null;
-                                lastErrorDetails = null;
-
-                                // setting firstPollDone to true if greedy 
polling is enabled
-                                firstPollDone = true;
+                                if (currentGeneration == generation.get()) {
+                                    // clear any error that might be since we 
have successfully polled, otherwise readiness checks might believe the
+                                    // consumer to be unhealthy
+                                    errorCounter.set(0);
+                                    backoffErrorBase.set(0);
+                                    lastError = null;
+                                    lastErrorDetails = null;
+
+                                    // setting firstPollDone to true if greedy 
polling is enabled
+                                    firstPollDone = true;
+                                }
                             }
                         } else {
                             LOG.debug("Cannot begin polling as pollStrategy 
returned false: {}", pollStrategy);
@@ -267,11 +279,19 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
             }
         }
 
+        if (currentGeneration != generation.get()) {
+            // the consumer has been stopped while polling, so the state 
(which has been reset) must not be updated
+            LOG.trace("doRun() done after the consumer was stopped");
+            return;
+        }
+
         if (cause != null) {
             idleCounter.set(0);
             successCounter.set(0);
             errorCounter.incrementAndGet();
             lastError = cause;
+            // the details are for this error only
+            lastErrorDetails = null;
             // enrich last error with http response code if possible
             if (cause instanceof HttpResponseAware httpResponseAware) {
                 int code = httpResponseAware.getHttpResponseCode();
@@ -287,6 +307,7 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
             }
             successCounter.incrementAndGet();
             errorCounter.set(0);
+            backoffErrorBase.set(0);
             lastError = null;
             lastErrorDetails = null;
         }
@@ -641,15 +662,17 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
 
         boolean newScheduler = false;
         if (scheduler == null) {
-            DefaultScheduledPollConsumerScheduler scheduler
-                    = new 
DefaultScheduledPollConsumerScheduler(scheduledExecutorService);
-            scheduler.setDelay(delay);
-            scheduler.setInitialDelay(initialDelay);
-            scheduler.setTimeUnit(timeUnit);
-            scheduler.setUseFixedDelay(useFixedDelay);
-            this.scheduler = scheduler;
+            this.scheduler = new 
DefaultScheduledPollConsumerScheduler(scheduledExecutorService);
+            consumerCreatedScheduler = true;
             newScheduler = true;
         }
+        if (consumerCreatedScheduler && scheduler instanceof 
DefaultScheduledPollConsumerScheduler defaultScheduler) {
+            // (re)apply the options from this consumer, which may have been 
changed while stopped (such as via JMX)
+            defaultScheduler.setDelay(delay);
+            defaultScheduler.setInitialDelay(initialDelay);
+            defaultScheduler.setTimeUnit(timeUnit);
+            defaultScheduler.setUseFixedDelay(useFixedDelay);
+        }
         ObjectHelper.notNull(scheduler, "scheduler", this);
         scheduler.setCamelContext(getEndpoint().getCamelContext());
 
@@ -713,15 +736,21 @@ public abstract class ScheduledPollConsumer extends 
DefaultConsumer
             ServiceHelper.stopAndShutdownServices(scheduler);
         }
 
+        // a poll that is still running must not update the state after it has 
been cleared
+        generation.incrementAndGet();
+
         // clear counters
         backoffCounter.set(0);
         idleCounter.set(0);
         errorCounter.set(0);
+        backoffErrorBase.set(0);
         successCounter.set(0);
         counter.set(0);
         // clear ready state
         firstPollDone = false;
         forceReady = false;
+        lastError = null;
+        lastErrorDetails = null;
 
         super.doStop();
     }
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 3e50f5aa2cba..8fbb0a4c0b5e 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
@@ -3446,6 +3446,15 @@ redelivery now follows 
`negativeAckRedeliveryDelayMicros`, which defaults to 60
 
 A route that wants the previous timing can set 
`negativeAckRedeliveryDelayMicros=10000000`.
 
+=== camel-core - scheduled poll consumer
+
+The `delay`, `initialDelay`, `timeUnit` and `useFixedDelay` of a scheduled 
poll consumer that are changed while the
+consumer is stopped (such as using JMX) are now used when the consumer is 
started again. Previously the values from
+the first start were kept.
+
+When a backoff (`backoffMultiplier`) finishes, the error counter of the 
consumer (used by the health checks and JMX)
+is no longer reset to 0. It is reset when a poll succeeds. The number of 
skipped polls is the same as before.
+
 === camel-seda - multipleConsumers broadcasts to consumers with different uri 
options
 
 With `multipleConsumers=true` every consumer of a SEDA queue now receives a 
copy of each message, also when the consumers

Reply via email to