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 259c390726f7 CAMEL-25129: camel-resilience4j - the bulkhead should 
limit the calls that keep running after a timeout (#27038)
259c390726f7 is described below

commit 259c390726f7e8071d78cfcbd91ecd30c0b89d18
Author: allthingssecurity <[email protected]>
AuthorDate: Tue Sep 29 17:27:22 2026 +0530

    CAMEL-25129: camel-resilience4j - the bulkhead should limit the calls that 
keep running after a timeout (#27038)
    
    * CAMEL-25129: camel-resilience4j - the bulkhead should limit the calls 
that keep running after a timeout
    * CAMEL-25129: camel-resilience4j - note the TimeLimiter metrics in the 
upgrade guide
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../resilience4j/ResilienceProcessor.java          |  23 ++--
 .../ResilienceBulkheadTimeoutTest.java             | 118 +++++++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  17 +++
 3 files changed, 150 insertions(+), 8 deletions(-)

diff --git 
a/components/camel-resilience4j/src/main/java/org/apache/camel/component/resilience4j/ResilienceProcessor.java
 
b/components/camel-resilience4j/src/main/java/org/apache/camel/component/resilience4j/ResilienceProcessor.java
index b5024f2371fe..c35f096ee343 100644
--- 
a/components/camel-resilience4j/src/main/java/org/apache/camel/component/resilience4j/ResilienceProcessor.java
+++ 
b/components/camel-resilience4j/src/main/java/org/apache/camel/component/resilience4j/ResilienceProcessor.java
@@ -588,14 +588,20 @@ public class ResilienceProcessor extends 
BaseProcessorSupport
             Callable<Exchange> callable;
 
             if (timeLimiter != null) {
-                Supplier<CompletableFuture<Exchange>> futureSupplier
-                        = () -> CompletableFuture.supplyAsync(ftask, 
executorService);
+                Supplier<CompletionStage<Exchange>> stage = () -> 
CompletableFuture.supplyAsync(ftask, executorService);
+                if (bulkhead != null) {
+                    // the bulkhead goes inside the time limiter, so the 
permit is held until the task is done,
+                    // and not released when the timeout fires while the task 
is still running
+                    stage = Bulkhead.decorateCompletionStage(bulkhead, stage);
+                }
+                final Supplier<CompletionStage<Exchange>> fstage = stage;
+                Supplier<CompletableFuture<Exchange>> futureSupplier = () -> 
fstage.get().toCompletableFuture();
                 callable = TimeLimiter.decorateFutureSupplier(timeLimiter, 
futureSupplier);
             } else {
                 callable = task;
-            }
-            if (bulkhead != null) {
-                callable = Bulkhead.decorateCallable(bulkhead, callable);
+                if (bulkhead != null) {
+                    callable = Bulkhead.decorateCallable(bulkhead, callable);
+                }
             }
 
             callable = CircuitBreaker.decorateCallable(circuitBreaker, 
callable);
@@ -658,12 +664,13 @@ public class ResilienceProcessor extends 
BaseProcessorSupport
             }
 
             // decorate with resilience4j CompletionStage decorators
-            if (timeLimiter != null) {
-                supplier = TimeLimiter.decorateCompletionStage(timeLimiter, 
scheduledExecutorService, supplier);
-            }
+            // (the bulkhead inside the time limiter, so the permit is held 
until the task is done)
             if (bulkhead != null) {
                 supplier = Bulkhead.decorateCompletionStage(bulkhead, 
supplier);
             }
+            if (timeLimiter != null) {
+                supplier = TimeLimiter.decorateCompletionStage(timeLimiter, 
scheduledExecutorService, supplier);
+            }
             supplier = CircuitBreaker.decorateCompletionStage(circuitBreaker, 
supplier);
 
             // trigger the chain and handle completion
diff --git 
a/components/camel-resilience4j/src/test/java/org/apache/camel/component/resilience4j/ResilienceBulkheadTimeoutTest.java
 
b/components/camel-resilience4j/src/test/java/org/apache/camel/component/resilience4j/ResilienceBulkheadTimeoutTest.java
new file mode 100644
index 000000000000..be5c5401310d
--- /dev/null
+++ 
b/components/camel-resilience4j/src/test/java/org/apache/camel/component/resilience4j/ResilienceBulkheadTimeoutTest.java
@@ -0,0 +1,118 @@
+/*
+ * 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.resilience4j;
+
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.spi.CircuitBreakerConstants;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.AfterEach;
+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.assertTrue;
+
+/**
+ * With timeout and bulkhead enabled, a call that timed out but is still 
running must keep its bulkhead permit, so the
+ * bulkhead limits the calls that run against a slow service, and not only the 
callers that wait for them.
+ */
+public class ResilienceBulkheadTimeoutTest extends CamelTestSupport {
+
+    private final CountDownLatch release = new CountDownLatch(1);
+    private final AtomicInteger running = new AtomicInteger();
+    private final AtomicInteger maxRunning = new AtomicInteger();
+    private final AtomicInteger calls = new AtomicInteger();
+
+    @AfterEach
+    public void releaseCalls() {
+        release.countDown();
+    }
+
+    @Test
+    public void testBulkheadWithTimeout() throws Exception {
+        doTestBulkheadWithTimeout("direct:sync", "cbSync");
+    }
+
+    @Test
+    public void testBulkheadWithTimeoutAsynchronous() throws Exception {
+        doTestBulkheadWithTimeout("direct:async", "cbAsync");
+    }
+
+    private void doTestBulkheadWithTimeout(String uri, String id) throws 
Exception {
+        // the first call times out and gets the fallback, but it is still 
running against the slow service
+        Exchange first = template.request(uri, e -> 
e.getMessage().setBody("Hello 1"));
+        assertEquals("Fallback message", first.getMessage().getBody());
+        assertEquals(Boolean.TRUE, 
first.getProperty(CircuitBreakerConstants.RESPONSE_FROM_FALLBACK));
+        assertEquals(Boolean.FALSE, 
first.getProperty(CircuitBreakerConstants.RESPONSE_REJECTED));
+
+        // so the next calls must be rejected by the bulkhead (it allows 1 
call), and not run as well
+        for (int i = 2; i <= 3; i++) {
+            String body = "Hello " + i;
+            Exchange out = template.request(uri, e -> 
e.getMessage().setBody(body));
+            assertEquals("Fallback message", out.getMessage().getBody());
+            assertEquals(Boolean.TRUE, 
out.getProperty(CircuitBreakerConstants.RESPONSE_REJECTED),
+                    "call " + i + " should be rejected by the bulkhead");
+        }
+        assertEquals(1, calls.get(), "only the first call should have been 
started");
+        assertEquals(1, maxRunning.get(), "the bulkhead allows only 1 call at 
a time");
+        ResilienceProcessor cb = context.getProcessor(id, 
ResilienceProcessor.class);
+        assertEquals(2, cb.getNumberOfBulkheadRejectedCalls());
+        assertEquals(1, cb.getNumberOfTimedOutCalls());
+
+        // when the slow call ends it releases the permit
+        release.countDown();
+        await().atMost(5, TimeUnit.SECONDS)
+                .until(() -> "Bye World".equals(template.requestBody(uri, 
"Hello 4", String.class)));
+        assertEquals(1, maxRunning.get());
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:sync").circuitBreaker().id("cbSync")
+                        
.resilience4jConfiguration().bulkheadEnabled(true).bulkheadMaxConcurrentCalls(1)
+                        .timeoutEnabled(true).timeoutDuration(200).end()
+                        .to("direct:slow")
+                        .onFallback().transform().constant("Fallback 
message").end();
+
+                from("direct:async").circuitBreaker().id("cbAsync")
+                        
.resilience4jConfiguration().asynchronous(true).bulkheadEnabled(true).bulkheadMaxConcurrentCalls(1)
+                        .timeoutEnabled(true).timeoutDuration(200).end()
+                        .to("direct:slow")
+                        .onFallback().transform().constant("Fallback 
message").end();
+
+                from("direct:slow").process(e -> {
+                    calls.incrementAndGet();
+                    maxRunning.accumulateAndGet(running.incrementAndGet(), 
Math::max);
+                    try {
+                        // a slow service
+                        assertTrue(release.await(10, TimeUnit.SECONDS));
+                    } finally {
+                        running.decrementAndGet();
+                    }
+                }).transform().constant("Bye World");
+            }
+        };
+    }
+}
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 89c1cda4e272..bbe263fc2f08 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
@@ -173,6 +173,23 @@ Prior to Camel 4.23 the property was only set when there 
was no fallback and was
 so a fallback that tested it for `null` must now test for `true` or `false` 
instead.
 `CamelCircuitBreakerResponseShortCircuited` is unchanged and remains `true` 
whenever the fallback runs, whatever the cause.
 
+=== camel-resilience4j - bulkhead together with timeout
+
+When the Circuit Breaker EIP with resilience4j has both `bulkheadEnabled` and 
`timeoutEnabled`, the bulkhead permit is
+now held until the protected call has ended, also when the call has timed out 
and the fallback has already been used.
+Prior to Camel 4.23 the permit was released when the timeout fired, while the 
call kept running on the timeout thread
+pool, so `bulkheadMaxConcurrentCalls` did not limit the calls running against 
a slow service. This applies to both the
+synchronous and the `asynchronous` mode, and matches the decorator order 
recommended by resilience4j (the bulkhead
+inside the time limiter) and the behaviour of 
`camel-microprofile-fault-tolerance`.
+
+As a result, while calls that timed out are still running, further calls are 
rejected by the bulkhead (and answered by
+the fallback, with `CamelCircuitBreakerResponseRejected` set to `true`), where 
before they were started. A call that
+never ends keeps its permit.
+
+In the synchronous mode the resilience4j `TimeLimiter` now also sees a call 
that the bulkhead rejects (as an error with
+a `BulkheadFullException`), so the events and metrics of resilience4j itself 
count such a call as a `TimeLimiter`
+error. The Camel circuit breaker counters are not affected.
+
 === Convert Body, Convert Header and Convert Variable EIPs
 
 When a charset is configured, such as `convertBodyTo(String.class, "UTF-8")`, 
then the conversion now uses that charset

Reply via email to