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 7bfac9f60f62 CAMEL-25014: camel-support - 
ThrottlingExceptionRoutePolicy: keep re-checking an open circuit, and suspend 
the consumer after a restart with keepOpen (#26874)
7bfac9f60f62 is described below

commit 7bfac9f60f62b11b7935c31ed701619a6fecfce7
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 13:47:17 2026 +0530

    CAMEL-25014: camel-support - ThrottlingExceptionRoutePolicy: keep 
re-checking an open circuit, and suspend the consumer after a restart with 
keepOpen (#26874)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../ThrottlingExceptionRoutePolicyReopenTest.java  | 107 +++++++++++++++++++++
 .../throttling/ThrottlingExceptionRoutePolicy.java |  51 ++++++++--
 2 files changed, 150 insertions(+), 8 deletions(-)

diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingExceptionRoutePolicyReopenTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingExceptionRoutePolicyReopenTest.java
new file mode 100644
index 000000000000..1520014dd02d
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/throttle/ThrottlingExceptionRoutePolicyReopenTest.java
@@ -0,0 +1,107 @@
+/*
+ * 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 java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.ContextTestSupport;
+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.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+
+/**
+ * An open circuit is kept open, and re-checked, when the half open handler is 
not ready, and a keepOpen circuit
+ * suspends the consumer again when the route is restarted.
+ */
+class ThrottlingExceptionRoutePolicyReopenTest extends ContextTestSupport {
+
+    private final AtomicInteger handlerCalls = new AtomicInteger();
+    private final AtomicBoolean handlerReady = new AtomicBoolean();
+
+    @Override
+    @BeforeEach
+    public void setUp() throws Exception {
+        super.setUp();
+        context.getShutdownStrategy().setTimeout(1);
+    }
+
+    @Test
+    void testHalfOpenHandlerNotReadyChecksAgain() throws Exception {
+        MockEndpoint result = getMockEndpoint("mock:result");
+        result.whenExchangeReceived(1, e -> e.setException(new 
ThrottlingException("boom")));
+        result.expectedBodiesReceived("fail", "ok");
+
+        template.sendBody("seda:handler", "fail");
+        ServiceSupport consumer = (ServiceSupport) 
context.getRoute("handler").getConsumer();
+        await().atMost(10, TimeUnit.SECONDS).until(consumer::isSuspended);
+
+        // the handler is not ready, so the circuit stays open, and is checked 
again
+        await().atMost(10, TimeUnit.SECONDS).until(() -> handlerCalls.get() >= 
2);
+        await().atMost(10, TimeUnit.SECONDS).until(consumer::isSuspended);
+
+        // once the handler is ready, the next check closes the circuit
+        handlerReady.set(true);
+        await().atMost(10, TimeUnit.SECONDS).until(consumer::isStarted);
+
+        template.sendBody("seda:handler", "ok");
+        MockEndpoint.assertIsSatisfied(context, 10, TimeUnit.SECONDS);
+    }
+
+    @Test
+    void testKeepOpenAfterRouteRestart() throws Exception {
+        ServiceSupport consumer = (ServiceSupport) 
context.getRoute("keepOpen").getConsumer();
+        await().atMost(10, TimeUnit.SECONDS).until(consumer::isSuspended);
+
+        context.getRouteController().stopRoute("keepOpen");
+        context.getRouteController().startRoute("keepOpen");
+
+        // the circuit is still open, so the new consumer must be suspended 
again
+        ServiceSupport restarted = (ServiceSupport) 
context.getRoute("keepOpen").getConsumer();
+        await().atMost(10, TimeUnit.SECONDS).until(restarted::isSuspended);
+
+        getMockEndpoint("mock:keepOpen").expectedMessageCount(0);
+        getMockEndpoint("mock:keepOpen").setAssertPeriod(500);
+        template.sendBody("seda:keepOpen", "should not be consumed");
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                ThrottlingExceptionRoutePolicy handlerPolicy = new 
ThrottlingExceptionRoutePolicy(1, 60000, 200, null);
+                handlerPolicy.setHalfOpenHandler(() -> {
+                    handlerCalls.incrementAndGet();
+                    return handlerReady.get();
+                });
+                
from("seda:handler").routeId("handler").routePolicy(handlerPolicy).to("mock:result");
+
+                ThrottlingExceptionRoutePolicy keepOpenPolicy = new 
ThrottlingExceptionRoutePolicy(1, 60000, 200, null);
+                keepOpenPolicy.setKeepOpen(true);
+                
from("seda:keepOpen").routeId("keepOpen").routePolicy(keepOpenPolicy).to("mock:keepOpen");
+            }
+        };
+    }
+}
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 e453f4be20b3..033f5b1d21ad 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
@@ -177,7 +177,9 @@ public class ThrottlingExceptionRoutePolicy extends 
RoutePolicySupport implement
     public void onStart(Route route) {
         // if keepOpen then start w/ the circuit open
         if (keepOpenBool.get()) {
-            openCircuit(route);
+            // the circuit may still be open from before the route was 
stopped, and then the started consumer must
+            // be suspended again
+            reopenCircuit(route);
         }
     }
 
@@ -276,7 +278,8 @@ public class ThrottlingExceptionRoutePolicy extends 
RoutePolicySupport implement
                             closeCircuit(route);
                         } else {
                             LOG.debug("Opening circuit...");
-                            openCircuit(route);
+                            // keep the circuit open, and check again after 
halfOpenAfter
+                            reopenCircuit(route);
                         }
                     } else {
                         LOG.debug("Half opening circuit...");
@@ -323,9 +326,40 @@ public class ThrottlingExceptionRoutePolicy extends 
RoutePolicySupport implement
         }
     }
 
+    /**
+     * Opens the circuit, also when it is already open: suspends the consumer 
and schedules the next half open check.
+     * Unlike {@link #openCircuit(Route)}, which only acts when the circuit is 
not open yet, this is used to keep an
+     * open circuit open.
+     */
+    protected void reopenCircuit(Route route) {
+        try {
+            lock.lock();
+            suspendOrStopConsumer(route.getConsumer());
+            state.set(STATE_OPEN);
+            openedAt = System.currentTimeMillis();
+            this.addHalfOpenTimer(route);
+            logState();
+        } catch (Exception e) {
+            handleException(e);
+        } finally {
+            lock.unlock();
+        }
+    }
+
     protected void addHalfOpenTimer(Route route) {
-        halfOpenTimer = new Timer();
-        halfOpenTimer.schedule(new HalfOpenTask(route), halfOpenAfter);
+        lock.lock();
+        try {
+            Timer previous = halfOpenTimer;
+            Timer timer = new Timer();
+            timer.schedule(new HalfOpenTask(route, timer), halfOpenAfter);
+            halfOpenTimer = timer;
+            if (previous != null) {
+                // only one half open check at a time (each timer has its own 
thread)
+                previous.cancel();
+            }
+        } finally {
+            lock.unlock();
+        }
     }
 
     protected void halfOpenCircuit(Route route) {
@@ -390,16 +424,17 @@ public class ThrottlingExceptionRoutePolicy extends 
RoutePolicySupport implement
 
     class HalfOpenTask extends TimerTask {
         private final Route route;
+        private final Timer timer;
 
-        HalfOpenTask(Route route) {
+        HalfOpenTask(Route route, Timer timer) {
             this.route = route;
+            this.timer = timer;
         }
 
         @Override
         public void run() {
-            if (halfOpenTimer != null) {
-                halfOpenTimer.cancel();
-            }
+            // cancel the timer of this task only, a newer timer may already 
have been scheduled
+            timer.cancel();
             calculateState(route);
         }
     }

Reply via email to