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 c299d06e8ff5 CAMEL-25018: camel-support - Throttling route policies 
must not resume a consumer suspended by the route controller (#26883)
c299d06e8ff5 is described below

commit c299d06e8ff58a68657e6b9c964ef71002982ba5
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 17:55:52 2026 +0530

    CAMEL-25018: camel-support - Throttling route policies must not resume a 
consumer suspended by the route controller (#26883)
    
    ThrottlingInflightRoutePolicy and ThrottlingExceptionRoutePolicy resumed the
    route's consumer with resumeOrStartConsumer, which resumes any suspended
    consumer (and starts any consumer that is not Suspendable), whoever
    suspended or stopped it. The route controller uses the same consumer state:
    for stopRoute and suspendRoute the graceful shutdown suspends the consumer
    first, then waits for the inflight exchanges, and only then stops or
    suspends the route.
    
    - ThrottlingInflightRoutePolicy resumed the consumer on every exchange that
      completed with the inflight count at or below the resume threshold. So the
      first exchange that completed during a graceful stopRoute or suspendRoute
      resumed the consumer, which took new messages while the stop waited for
      the inflight exchanges. Under steady traffic the stop ran into the
      timeout and was forced, failing the exchanges it had just taken in.
    - ThrottlingExceptionRoutePolicy resumed the consumer from the half open
      timer, also when the route had been suspended (e.g. by an operator while
      the circuit was open): the route was Suspended but consumed.
    
    Both policies now only resume a consumer they suspended themselves, and
    not when the route controller has suspended or stopped the route, or is
    suspending or stopping it, or Camel is stopping. The route and context
    status check is a new helper in RoutePolicySupport,
    isResumeOrStartConsumerAllowed.
    
    - ThrottlingInflightRoutePolicy remembers the consumers it suspended, and
      forgets them when the route controller starts, stops, suspends or
      resumes the route. It only suspends (and so remembers) a consumer that
      is started, as ServiceHelper.suspendService also returns true for a
      consumer that is already stopped, such as one the shutdown strategy
      stopped, which would otherwise be claimed and restarted later.
    - ThrottlingExceptionRoutePolicy only resumes the consumer when the circuit
      is open, as it suspended the consumer when it opened the circuit. When
      half open the consumer was already resumed, so closing the circuit no
      longer resumes a consumer that the route controller suspended in the
      meantime. If the resume is not allowed, the circuit state still changes,
      and the consumer is left to the route controller.
    
    A consumer that the policy suspended itself can still be resumed while a
    graceful stop or suspend waits for the inflight exchanges, because the
    route status only changes after that wait.
    
    ThrottlingInflightRoutePolicy still throttles a consumer that is not a
    StatefulService, and throttle() skips the resume check when the policy has
    not suspended any consumer. Adds a 4.23 upgrade guide note that the
    throttling policies only resume a consumer they suspended.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 ...ttlingExceptionRoutePolicySuspendRouteTest.java |  88 +++++++++++
 ...InflightRoutePolicyNonStatefulConsumerTest.java | 167 +++++++++++++++++++++
 ...ThrottlingInflightRoutePolicyStopRouteTest.java | 106 +++++++++++++
 ...lingInflightRoutePolicyStoppedConsumerTest.java |  81 ++++++++++
 .../apache/camel/support/RoutePolicySupport.java   |  22 +++
 .../throttling/ThrottlingExceptionRoutePolicy.java |  12 +-
 .../throttling/ThrottlingInflightRoutePolicy.java  |  56 ++++++-
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  13 ++
 8 files changed, 541 insertions(+), 4 deletions(-)

diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingExceptionRoutePolicySuspendRouteTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingExceptionRoutePolicySuspendRouteTest.java
new file mode 100644
index 000000000000..17e33367a224
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingExceptionRoutePolicySuspendRouteTest.java
@@ -0,0 +1,88 @@
+/*
+ * 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.processor.throttle;
+
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.ServiceStatus;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.throttling.ThrottlingExceptionRoutePolicy;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * When the route is suspended by the route controller while the circuit is 
open, the half open timer of the
+ * {@link ThrottlingExceptionRoutePolicy} must not resume the consumer.
+ */
+class ThrottlingExceptionRoutePolicySuspendRouteTest extends 
ContextTestSupport {
+
+    private ThrottlingExceptionRoutePolicy policy;
+
+    @Test
+    void testSuspendedRouteStaysSuspendedWhenHalfOpen() throws Exception {
+        ServiceSupport consumer = (ServiceSupport) 
context.getRoute("foo").getConsumer();
+
+        // a failure opens the circuit, which suspends the consumer
+        assertThrows(Exception.class, () -> template.sendBody("direct:start", 
"Kaboom"));
+        await().atMost(10, TimeUnit.SECONDS).until(consumer::isSuspended);
+        assertEquals("opened", policy.getStateAsString());
+
+        // the route is suspended (such as by an operator) while the circuit 
is open
+        context.getRouteController().suspendRoute("foo");
+        assertEquals(ServiceStatus.Suspended, 
context.getRouteController().getRouteStatus("foo"));
+        assertEquals("opened", policy.getStateAsString(), "The half open timer 
should not have fired yet");
+
+        // the half open timer fires, but the route must stay suspended
+        await().atMost(10, TimeUnit.SECONDS).until(() -> "half 
opened".equals(policy.getStateAsString()));
+        assertTrue(consumer.isSuspended(), "The consumer of the suspended 
route should not be resumed");
+        assertEquals(ServiceStatus.Suspended, 
context.getRouteController().getRouteStatus("foo"));
+
+        // when the route is resumed then the circuit closes on success
+        MockEndpoint mock = getMockEndpoint("mock:result");
+        mock.expectedBodiesReceived("Hello World");
+        context.getRouteController().resumeRoute("foo");
+        template.sendBody("direct:start", "Hello World");
+        mock.assertIsSatisfied();
+        assertEquals("closed", policy.getStateAsString());
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                // open on the first failure, and only count failures in the 
last 100 millis, so the circuit closes on
+                // the first success after it has been half opened. Half open 
after 3 seconds, so the timer does not
+                // fire before the test has checked that the circuit is open, 
also on a slow machine
+                policy = new ThrottlingExceptionRoutePolicy(1, 100, 3000, 
null);
+
+                
from("direct:start?block=false").routeId("foo").routePolicy(policy)
+                        .filter(body().isEqualTo("Kaboom"))
+                            .throwException(new 
IllegalArgumentException("Forced"))
+                        .end()
+                        .to("mock:result");
+            }
+        };
+    }
+}
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyNonStatefulConsumerTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyNonStatefulConsumerTest.java
new file mode 100644
index 000000000000..f2f35ff32723
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyNonStatefulConsumerTest.java
@@ -0,0 +1,167 @@
+/*
+ * 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.processor.throttle;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.Route;
+import org.apache.camel.StatefulService;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.InflightRepository;
+import org.apache.camel.support.DefaultEndpoint;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.throttling.ThrottlingInflightRoutePolicy;
+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.assertTrue;
+
+/**
+ * The {@link ThrottlingInflightRoutePolicy} must still throttle a consumer 
that is not a {@link StatefulService}: it is
+ * stopped when too many exchanges are inflight, and started again when they 
have completed.
+ */
+class ThrottlingInflightRoutePolicyNonStatefulConsumerTest extends 
ContextTestSupport {
+
+    private final ThrottlingInflightRoutePolicy policy = new 
ThrottlingInflightRoutePolicy();
+
+    @Test
+    void testThrottlesConsumerThatIsNotStatefulService() throws Exception {
+        Route route = context.getRoute("foo");
+        PlainConsumer consumer = (PlainConsumer) route.getConsumer();
+        assertFalse(route.getConsumer() instanceof StatefulService);
+        assertTrue(consumer.running);
+        assertEquals(0, consumer.stops.get());
+        int startsBefore = consumer.starts.get();
+
+        // an exchange completes while more exchanges than the maximum are 
inflight
+        InflightRepository inflight = context.getInflightRepository();
+        Exchange first = new DefaultExchange(context);
+        Exchange second = new DefaultExchange(context);
+        inflight.add(first, "foo");
+        inflight.add(second, "foo");
+        policy.onExchangeDone(route, first);
+
+        assertEquals(1, consumer.stops.get(), "The policy should stop the 
consumer when too many exchanges are inflight");
+        assertFalse(consumer.running);
+
+        // and then the inflight exchanges complete
+        inflight.remove(first, "foo");
+        inflight.remove(second, "foo");
+        policy.onExchangeDone(route, second);
+
+        assertEquals(startsBefore + 1, consumer.starts.get(), "The policy 
should start the consumer it stopped");
+        assertTrue(consumer.running);
+    }
+
+    @Override
+    protected CamelContext createCamelContext() throws Exception {
+        CamelContext context = super.createCamelContext();
+        context.addEndpoint("plain", new DefaultEndpoint() {
+            @Override
+            public Producer createProducer() {
+                throw new UnsupportedOperationException();
+            }
+
+            @Override
+            public Consumer createConsumer(Processor processor) {
+                return new PlainConsumer(this, processor);
+            }
+
+            @Override
+            protected String createEndpointUri() {
+                return "plain";
+            }
+
+            @Override
+            public boolean isSingleton() {
+                return true;
+            }
+        });
+        return context;
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                policy.setMaxInflightExchanges(1);
+
+                from("plain").routeId("foo").routePolicy(policy)
+                        .to("mock:result");
+            }
+        };
+    }
+
+    /**
+     * A consumer that implements only {@link Consumer}, and so is neither a 
{@link StatefulService} nor
+     * {@link org.apache.camel.Suspendable}.
+     */
+    private static final class PlainConsumer implements Consumer {
+
+        private final Endpoint endpoint;
+        private final Processor processor;
+        private final AtomicInteger starts = new AtomicInteger();
+        private final AtomicInteger stops = new AtomicInteger();
+        private volatile boolean running;
+
+        private PlainConsumer(Endpoint endpoint, Processor processor) {
+            this.endpoint = endpoint;
+            this.processor = processor;
+        }
+
+        @Override
+        public void start() {
+            starts.incrementAndGet();
+            running = true;
+        }
+
+        @Override
+        public void stop() {
+            stops.incrementAndGet();
+            running = false;
+        }
+
+        @Override
+        public Endpoint getEndpoint() {
+            return endpoint;
+        }
+
+        @Override
+        public Processor getProcessor() {
+            return processor;
+        }
+
+        @Override
+        public Exchange createExchange(boolean autoRelease) {
+            return endpoint.createExchange();
+        }
+
+        @Override
+        public void releaseExchange(Exchange exchange, boolean autoRelease) {
+            // noop
+        }
+    }
+}
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyStopRouteTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyStopRouteTest.java
new file mode 100644
index 000000000000..a6aa6282bd7e
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyStopRouteTest.java
@@ -0,0 +1,106 @@
+/*
+ * 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.processor.throttle;
+
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.ServiceStatus;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.throttling.ThrottlingInflightRoutePolicy;
+import org.junit.jupiter.api.Test;
+
+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;
+
+/**
+ * A graceful stop of a route suspends the consumer and waits for the inflight 
exchanges. When such an exchange
+ * completes, the {@link ThrottlingInflightRoutePolicy} must not resume the 
consumer, which would take new messages
+ * while the stop waits for the inflight exchanges.
+ */
+class ThrottlingInflightRoutePolicyStopRouteTest extends ContextTestSupport {
+
+    private final CountDownLatch firstStarted = new CountDownLatch(1);
+    private final CountDownLatch releaseFirst = new CountDownLatch(1);
+    private final CountDownLatch releaseOthers = new CountDownLatch(1);
+    private final AtomicInteger counter = new AtomicInteger();
+    private final AtomicBoolean stopping = new AtomicBoolean();
+    private final AtomicInteger startedDuringStop = new AtomicInteger();
+
+    @Test
+    void testGracefulStopRoute() throws Exception {
+        try {
+            assertTrue(firstStarted.await(10, TimeUnit.SECONDS), "The first 
exchange should be started");
+
+            CompletableFuture<Void> stop = CompletableFuture.runAsync(() -> {
+                try {
+                    context.getRouteController().stopRoute("foo", 5, 
TimeUnit.SECONDS);
+                } catch (Exception e) {
+                    throw new RuntimeException(e);
+                }
+            });
+
+            // the graceful stop suspends the consumer and then waits for the 
inflight exchange
+            ServiceSupport consumer = (ServiceSupport) 
context.getRoute("foo").getConsumer();
+            await().atMost(10, TimeUnit.SECONDS).until(consumer::isSuspended);
+            stopping.set(true);
+            releaseFirst.countDown();
+
+            stop.get(20, TimeUnit.SECONDS);
+
+            assertEquals(0, startedDuringStop.get(), "No exchange should be 
started while the route is being stopped");
+            assertFalse(context.getShutdownStrategy().isTimeoutOccurred(), 
"The route should be stopped gracefully");
+            assertEquals(ServiceStatus.Stopped, 
context.getRouteController().getRouteStatus("foo"));
+        } finally {
+            releaseFirst.countDown();
+            releaseOthers.countDown();
+        }
+    }
+
+    private void onExchange(Exchange exchange) throws Exception {
+        if (counter.incrementAndGet() == 1) {
+            firstStarted.countDown();
+            releaseFirst.await(20, TimeUnit.SECONDS);
+        } else if (stopping.get()) {
+            // an exchange that the consumer took while the route is being 
stopped keeps the stop waiting
+            startedDuringStop.incrementAndGet();
+            releaseOthers.await(20, TimeUnit.SECONDS);
+        }
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                ThrottlingInflightRoutePolicy policy = new 
ThrottlingInflightRoutePolicy();
+                policy.setMaxInflightExchanges(10);
+
+                from("timer:foo?period=10").routeId("foo").routePolicy(policy)
+                        .process(e -> onExchange(e));
+            }
+        };
+    }
+}
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyStoppedConsumerTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyStoppedConsumerTest.java
new file mode 100644
index 000000000000..7bc1641bee0f
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingInflightRoutePolicyStoppedConsumerTest.java
@@ -0,0 +1,81 @@
+/*
+ * 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.processor.throttle;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.Route;
+import org.apache.camel.ServiceStatus;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.InflightRepository;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.support.service.ServiceHelper;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.throttling.ThrottlingInflightRoutePolicy;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * The {@link ThrottlingInflightRoutePolicy} must only resume a consumer that 
it suspended itself. A consumer that was
+ * already stopped (such as by the shutdown strategy while it waits for the 
inflight exchanges of the route) must not be
+ * claimed by the policy when too many exchanges are inflight, and then be 
resumed when they complete.
+ */
+class ThrottlingInflightRoutePolicyStoppedConsumerTest extends 
ContextTestSupport {
+
+    private final ThrottlingInflightRoutePolicy policy = new 
ThrottlingInflightRoutePolicy();
+
+    @Test
+    void testDoesNotResumeConsumerStoppedByOthers() throws Exception {
+        Route route = context.getRoute("foo");
+        ServiceSupport consumer = (ServiceSupport) route.getConsumer();
+
+        // the consumer is stopped by someone else, while the route is still 
started
+        ServiceHelper.stopService(consumer);
+        assertTrue(consumer.isStopped());
+        assertEquals(ServiceStatus.Started, 
context.getRouteController().getRouteStatus("foo"));
+
+        // an exchange completes while more exchanges than the maximum are 
inflight
+        InflightRepository inflight = context.getInflightRepository();
+        Exchange first = new DefaultExchange(context);
+        Exchange second = new DefaultExchange(context);
+        inflight.add(first, "foo");
+        inflight.add(second, "foo");
+        policy.onExchangeDone(route, first);
+
+        // and then the inflight exchanges complete
+        inflight.remove(first, "foo");
+        inflight.remove(second, "foo");
+        policy.onExchangeDone(route, second);
+
+        assertTrue(consumer.isStopped(), "The policy should not resume a 
consumer that it did not suspend");
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                policy.setMaxInflightExchanges(1);
+
+                from("seda:foo").routeId("foo").routePolicy(policy)
+                        .to("mock:result");
+            }
+        };
+    }
+}
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/support/RoutePolicySupport.java
 
b/core/camel-support/src/main/java/org/apache/camel/support/RoutePolicySupport.java
index a4f63875534c..149af7f93fda 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/support/RoutePolicySupport.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/support/RoutePolicySupport.java
@@ -18,9 +18,11 @@ package org.apache.camel.support;
 
 import java.util.concurrent.TimeUnit;
 
+import org.apache.camel.CamelContext;
 import org.apache.camel.Consumer;
 import org.apache.camel.Exchange;
 import org.apache.camel.Route;
+import org.apache.camel.ServiceStatus;
 import org.apache.camel.spi.ExceptionHandler;
 import org.apache.camel.spi.RouteController;
 import org.apache.camel.spi.RoutePolicy;
@@ -122,6 +124,26 @@ public abstract class RoutePolicySupport extends 
ServiceSupport implements Route
         return ServiceHelper.resumeService(consumer);
     }
 
+    /**
+     * Whether this policy may resume or start the consumer of the route, 
after it has suspended or stopped the consumer
+     * itself (for example to throttle).
+     * <p/>
+     * This is not allowed when the route controller has suspended or stopped 
the route, or is suspending or stopping
+     * it, or when Camel is shutting down, as the consumer must then stay 
suspended or stopped.
+     *
+     * @param  route the route
+     * @return       <tt>true</tt> if the consumer may be resumed or started, 
<tt>false</tt> otherwise
+     */
+    protected boolean isResumeOrStartConsumerAllowed(Route route) {
+        CamelContext context = route.getCamelContext();
+        if (context.isStopping() || context.isStopped()) {
+            return false;
+        }
+        ServiceStatus status = controller(route).getRouteStatus(route.getId());
+        return status != null && !status.isSuspending() && 
!status.isSuspended() && !status.isStopping()
+                && !status.isStopped();
+    }
+
     public void startRoute(Route route) throws Exception {
         controller(route).startRoute(route.getId());
     }
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/throttling/ThrottlingExceptionRoutePolicy.java
 
b/core/camel-support/src/main/java/org/apache/camel/throttling/ThrottlingExceptionRoutePolicy.java
index 033f5b1d21ad..1f10d577692b 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/throttling/ThrottlingExceptionRoutePolicy.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/throttling/ThrottlingExceptionRoutePolicy.java
@@ -365,7 +365,10 @@ public class ThrottlingExceptionRoutePolicy extends 
RoutePolicySupport implement
     protected void halfOpenCircuit(Route route) {
         try {
             lock.lock();
-            resumeOrStartConsumer(route.getConsumer());
+            // do not resume the consumer if the route controller has 
suspended or stopped the route in the meantime
+            if (isResumeOrStartConsumerAllowed(route)) {
+                resumeOrStartConsumer(route.getConsumer());
+            }
             state.set(STATE_HALF_OPEN);
             logState();
         } catch (Exception e) {
@@ -378,7 +381,12 @@ public class ThrottlingExceptionRoutePolicy extends 
RoutePolicySupport implement
     protected void closeCircuit(Route route) {
         try {
             lock.lock();
-            resumeOrStartConsumer(route.getConsumer());
+            // only resume the consumer when the circuit is open, as then this 
policy has suspended the consumer
+            // (when half open the consumer has already been resumed, so if it 
is suspended now, then by the route
+            // controller), and not if the route controller has suspended or 
stopped the route in the meantime
+            if (state.get() == STATE_OPEN && 
isResumeOrStartConsumerAllowed(route)) {
+                resumeOrStartConsumer(route.getConsumer());
+            }
             failures.set(0);
             success.set(0);
             lastFailure = 0;
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/throttling/ThrottlingInflightRoutePolicy.java
 
b/core/camel-support/src/main/java/org/apache/camel/throttling/ThrottlingInflightRoutePolicy.java
index e2b7b7e8c2a6..cb183f05bf63 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/throttling/ThrottlingInflightRoutePolicy.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/throttling/ThrottlingInflightRoutePolicy.java
@@ -18,6 +18,7 @@ package org.apache.camel.throttling;
 
 import java.util.LinkedHashSet;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.locks.Lock;
 import java.util.concurrent.locks.ReentrantLock;
 
@@ -28,6 +29,7 @@ import org.apache.camel.Exchange;
 import org.apache.camel.LoggingLevel;
 import org.apache.camel.NonManagedService;
 import org.apache.camel.Route;
+import org.apache.camel.StatefulService;
 import org.apache.camel.spi.CamelEvent;
 import org.apache.camel.spi.CamelEvent.ExchangeCompletedEvent;
 import org.apache.camel.spi.CamelLogger;
@@ -66,6 +68,8 @@ public class ThrottlingInflightRoutePolicy extends 
RoutePolicySupport implements
     }
 
     private final Set<Route> routes = new LinkedHashSet<>();
+    // the consumers this policy has suspended (or stopped), as only those can 
be resumed (or started) by this policy
+    private final Set<Consumer> suspendedConsumers = 
ConcurrentHashMap.newKeySet();
     private ContextScopedEventNotifier eventNotifier;
     private CamelContext camelContext;
     private final Lock lock = new ReentrantLock();
@@ -116,6 +120,35 @@ public class ThrottlingInflightRoutePolicy extends 
RoutePolicySupport implements
         routes.add(route);
     }
 
+    @Override
+    public void onStart(Route route) {
+        routeControllerTookOver(route);
+    }
+
+    @Override
+    public void onStop(Route route) {
+        routeControllerTookOver(route);
+    }
+
+    @Override
+    public void onSuspend(Route route) {
+        routeControllerTookOver(route);
+    }
+
+    @Override
+    public void onResume(Route route) {
+        routeControllerTookOver(route);
+    }
+
+    private void routeControllerTookOver(Route route) {
+        // the route controller has started, stopped, suspended or resumed the 
consumer, so it is no longer suspended by
+        // this policy
+        Consumer consumer = route.getConsumer();
+        if (consumer != null) {
+            suspendedConsumers.remove(consumer);
+        }
+    }
+
     @Override
     public void onExchangeDone(Route route, Exchange exchange) {
         // if route scoped then throttle directly
@@ -156,6 +189,12 @@ public class ThrottlingInflightRoutePolicy extends 
RoutePolicySupport implements
             }
         }
 
+        // fast path: nothing to resume unless this policy has suspended a 
consumer (a consumer that this thread has just
+        // suspended above is already in the set, and every later completion 
checks again)
+        if (suspendedConsumers.isEmpty()) {
+            return;
+        }
+
         // reload size in case a race condition with too many at once being 
invoked
         // so we need to ensure that we read the most current size and start 
the consumer if we are already to low
         size = getSize(route, exchange);
@@ -166,7 +205,7 @@ public class ThrottlingInflightRoutePolicy extends 
RoutePolicySupport implements
         if (start) {
             try {
                 lock.lock();
-                startConsumer(size, consumer, resumeInflight);
+                startConsumer(route, size, consumer, resumeInflight);
             } catch (Exception e) {
                 handleException(e);
             } finally {
@@ -274,8 +313,14 @@ public class ThrottlingInflightRoutePolicy extends 
RoutePolicySupport implements
         }
     }
 
-    private void startConsumer(int size, Consumer consumer, int 
resumeInflight) throws Exception {
+    private void startConsumer(Route route, int size, Consumer consumer, int 
resumeInflight) throws Exception {
+        // only resume a consumer this policy has suspended, and not one the 
route controller has suspended or stopped
+        // (such as when it is stopping or suspending the route and waits for 
the inflight exchanges to complete)
+        if (!suspendedConsumers.contains(consumer) || 
!isResumeOrStartConsumerAllowed(route)) {
+            return;
+        }
         boolean started = resumeOrStartConsumer(consumer);
+        suspendedConsumers.remove(consumer);
         if (started) {
             getLogger().log("Throttling consumer: " + size + " <= " + 
resumeInflight
                             + " inflight exchange by resuming consumer: " + 
consumer);
@@ -283,8 +328,15 @@ public class ThrottlingInflightRoutePolicy extends 
RoutePolicySupport implements
     }
 
     private void stopConsumer(int size, Consumer consumer, int maxInflight) 
throws Exception {
+        // only suspend (and so later resume) a consumer that is started: 
suspendOrStopConsumer also returns true for a
+        // consumer that is already stopped, such as one the shutdown strategy 
has stopped, which must stay stopped.
+        // A consumer that is not a StatefulService has no state to check, so 
it is always throttled
+        if (consumer instanceof StatefulService ss && 
!ServiceHelper.isStarted(ss)) {
+            return;
+        }
         boolean stopped = suspendOrStopConsumer(consumer);
         if (stopped) {
+            suspendedConsumers.add(consumer);
             getLogger().log("Throttling consumer: " + size + " > " + 
maxInflight
                             + " inflight exchange by suspending consumer: " + 
consumer);
         }
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 a85c1efe1c2f..8bcae8034fb8 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
@@ -2580,6 +2580,19 @@ for the lifetime of the executor. Callers that inspect 
the returned `Future` wil
 return `true` after the task is done, where it previously stayed live. Callers 
that already cancel the
 `Future` themselves are unaffected.
 
+=== camel-support - the throttling route policies only resume a consumer they 
suspended
+
+`ThrottlingInflightRoutePolicy` now only resumes a consumer that it suspended 
itself. It also no longer resumes the
+consumer while the route is being stopped or suspended (for example while 
`stopRoute` waits for the inflight
+exchanges), or while Camel is stopping. `ThrottlingExceptionRoutePolicy` no 
longer resumes the consumer of a route that
+was suspended through the route controller while the circuit was open.
+
+For a consumer that is not `Suspendable`, `ThrottlingInflightRoutePolicy` 
previously called `start` on the consumer (and
+logged that it resumed the consumer) whenever an exchange completed below the 
resume threshold, even when it had not
+stopped the consumer. It now only starts a consumer that it stopped itself, so 
it no longer restarts a consumer that was
+stopped by the route controller or by other means. Such a consumer is still 
stopped and started to throttle the route,
+as before.
+
 === camel-master
 
 The `backOffMaxAttempts` option now bounds the attempts to start the delegated 
consumer as documented.

Reply via email to