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 a14d8a1f9f63 CAMEL-25002: camel-core - Loop EIP: a negative or huge 
loop count must not break graceful shutdown (#26860)
a14d8a1f9f63 is described below

commit a14d8a1f9f63620229d8bfe8c4779f66466ed42b
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 13:50:29 2026 +0530

    CAMEL-25002: camel-core - Loop EIP: a negative or huge loop count must not 
break graceful shutdown (#26860)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../camel/impl/engine/DefaultShutdownStrategy.java |  15 ++-
 .../org/apache/camel/processor/LoopProcessor.java  |  34 +++---
 .../LoopPendingExchangesShutdownTest.java          | 125 +++++++++++++++++++++
 3 files changed, 154 insertions(+), 20 deletions(-)

diff --git 
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultShutdownStrategy.java
 
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultShutdownStrategy.java
index d69112208646..036bdc0e64a4 100644
--- 
a/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultShutdownStrategy.java
+++ 
b/core/camel-base-engine/src/main/java/org/apache/camel/impl/engine/DefaultShutdownStrategy.java
@@ -665,13 +665,14 @@ public class DefaultShutdownStrategy extends 
ServiceSupport implements ShutdownS
             long loopDelaySeconds = 1;
             long loopCount = 0;
             while (!done && !timeoutOccurred.get()) {
-                int size = 0;
+                long size = 0;
                 // number of inflights per route
                 final Map<String, Integer> routeInflight = new 
LinkedHashMap<>();
 
                 for (RouteStartupOrder order : routes) {
-                    int inflight = 
context.getInflightRepository().size(order.getRoute().getId());
-                    inflight += getPendingInflightExchanges(order, 
suspendOnly);
+                    long sum = (long) 
context.getInflightRepository().size(order.getRoute().getId())
+                               + getPendingInflightExchanges(order, 
suspendOnly);
+                    int inflight = (int) Math.min(Integer.MAX_VALUE, sum);
                     if (inflight > 0) {
                         String routeId = order.getRoute().getId();
                         routeInflight.put(routeId, inflight);
@@ -781,7 +782,7 @@ public class DefaultShutdownStrategy extends ServiceSupport 
implements ShutdownS
      * @return             number of inflight exchanges
      */
     protected static int getPendingInflightExchanges(RouteStartupOrder order, 
boolean suspendOnly) {
-        int inflight = 0;
+        long inflight = 0;
 
         // the consumer is the 1st service so we always get the consumer
         // the child services are EIPs in the routes which may also have 
pending
@@ -790,12 +791,14 @@ public class DefaultShutdownStrategy extends 
ServiceSupport implements ShutdownS
             Set<Service> children = ServiceHelper.getChildServices(service);
             for (Service child : children) {
                 if (child instanceof ShutdownAware shutdownAware) {
-                    inflight += 
shutdownAware.getPendingExchangesSize(suspendOnly);
+                    // a negative size must not cancel out the pending 
exchanges of other services
+                    inflight += Math.max(0, 
shutdownAware.getPendingExchangesSize(suspendOnly));
                 }
             }
         }
 
-        return inflight;
+        // the same service can be a child of several route services, so the 
sum can exceed an int
+        return (int) Math.min(Integer.MAX_VALUE, inflight);
     }
 
     /**
diff --git 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/LoopProcessor.java
 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/LoopProcessor.java
index 179c0de64ef8..3a01705c934b 100644
--- 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/LoopProcessor.java
+++ 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/LoopProcessor.java
@@ -95,7 +95,8 @@ public class LoopProcessor extends 
BaseDelegateProcessorSupport
 
     @Override
     public int getPendingExchangesSize() {
-        return taskCount.intValue();
+        // the sum of the iterations left in all running loops, which can 
exceed an int
+        return (int) Math.max(0, Math.min(Integer.MAX_VALUE, taskCount.sum()));
     }
 
     @Override
@@ -113,6 +114,7 @@ public class LoopProcessor extends 
BaseDelegateProcessorSupport
         Exchange current;
         int index;
         int count;
+        boolean pendingTasksReleased;
 
         public LoopState(Exchange exchange, AsyncCallback callback) throws 
NoTypeConversionAvailableException {
             this.exchange = exchange;
@@ -126,7 +128,10 @@ public class LoopProcessor extends 
BaseDelegateProcessorSupport
                 String text = expression.evaluate(exchange, String.class);
                 count = ExchangeHelper.convertToMandatoryType(exchange, 
Integer.class, text);
                 // keep track of pending task if loop with fixed value
-                taskCount.add(count);
+                // (a zero or negative count means no iterations, so nothing 
is pending)
+                if (count > 0) {
+                    taskCount.add(count);
+                }
                 exchange.setProperty(ExchangePropertyKey.LOOP_SIZE, count);
             }
         }
@@ -164,13 +169,8 @@ public class LoopProcessor extends 
BaseDelegateProcessorSupport
                     if (LOG.isTraceEnabled()) {
                         LOG.trace("Processing complete for exchangeId: {} >>> 
{}", exchange.getExchangeId(), exchange);
                     }
-                    if (!cont && expression != null) {
-                        // if we should stop due to an exception etc, then 
make sure to dec task count
-                        int gap = count - index;
-                        while (gap-- > 0) {
-                            taskCount.decrement();
-                        }
-                    }
+                    // if we stop early due to an exception, or break on 
shutdown, then make sure to dec task count
+                    releasePendingTasks();
                     callback.done(false);
                 }
             } catch (Exception e) {
@@ -183,14 +183,20 @@ public class LoopProcessor extends 
BaseDelegateProcessorSupport
             if (LOG.isTraceEnabled()) {
                 LOG.trace("Processing failed for exchangeId: {} >>> {}", 
exchange.getExchangeId(), e.getMessage());
             }
-            if (expression != null) {
-                // if we should stop due to an exception etc, then make sure 
to dec task count
+            // if we should stop due to an exception etc, then make sure to 
dec task count
+            releasePendingTasks();
+            exchange.setException(e);
+        }
+
+        private void releasePendingTasks() {
+            // only once, as the exception handling may run after the loop has 
completed (e.g. the callback failed)
+            if (expression != null && !pendingTasksReleased) {
+                pendingTasksReleased = true;
                 int gap = count - index;
-                while (gap-- > 0) {
-                    taskCount.decrement();
+                if (gap > 0) {
+                    taskCount.add(-gap);
                 }
             }
-            exchange.setException(e);
         }
 
         @Override
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/LoopPendingExchangesShutdownTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/LoopPendingExchangesShutdownTest.java
new file mode 100644
index 000000000000..cc6622f4ae68
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/LoopPendingExchangesShutdownTest.java
@@ -0,0 +1,125 @@
+/*
+ * 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;
+
+import java.util.concurrent.Future;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+
+import static java.util.concurrent.TimeUnit.SECONDS;
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+/**
+ * A loop with a zero or negative count runs no iterations and must not leave 
anything behind in the pending task count,
+ * which the shutdown strategy uses to decide whether it should keep waiting 
for inflight exchanges.
+ */
+class LoopPendingExchangesShutdownTest extends ContextTestSupport {
+
+    @Test
+    void testNegativeCountLeavesNoPendingTasks() throws Exception {
+        MockEndpoint loop = getMockEndpoint("mock:loop");
+        loop.expectedMessageCount(3);
+        MockEndpoint done = getMockEndpoint("mock:done");
+        done.expectedBodiesReceived("a", "b", "c");
+
+        template.sendBodyAndHeader("direct:start", "a", "n", -1);
+        template.sendBodyAndHeader("direct:start", "b", "n", 0);
+        template.sendBodyAndHeader("direct:start", "c", "n", 3);
+
+        assertMockEndpointsSatisfied();
+        assertEquals(0, context.getProcessor("myLoop", 
LoopProcessor.class).getPendingExchangesSize());
+    }
+
+    @Test
+    void testGracefulShutdownWaitsForInflightAfterNegativeCount() throws 
Exception {
+        MockEndpoint done = getMockEndpoint("mock:done");
+        done.expectedBodiesReceived("first", "slow");
+
+        // a single message with a negative loop count before the shutdown
+        template.sendBodyAndHeader("direct:start", "first", "n", -1);
+
+        // this message is inflight when the context is stopped, and only 
continues once the shutdown has begun
+        Future<Exchange> slow = template.asyncSend("direct:start", e -> {
+            e.getMessage().setBody("slow");
+            e.getMessage().setHeader("n", 1);
+        });
+        await().atMost(10, SECONDS).until(() -> 
context.getInflightRepository().size("myRoute") == 1);
+
+        // graceful shutdown must wait for the inflight exchange to complete
+        context.stop();
+
+        assertNull(slow.get(10, SECONDS).getException());
+        done.assertIsSatisfied();
+    }
+
+    @Test
+    void testGracefulShutdownWaitsForInflightWithHugeCount() throws Exception {
+        MockEndpoint done = getMockEndpoint("mock:hugeDone");
+        done.expectedBodiesReceived("huge");
+
+        // the loop breaks on shutdown after its first iteration, which only 
completes once the shutdown has begun
+        Future<Exchange> huge = template.asyncSend("direct:huge", e -> {
+            e.getMessage().setBody("huge");
+            e.getMessage().setHeader("n", Integer.MAX_VALUE);
+        });
+        await().atMost(10, SECONDS).until(() -> 
context.getInflightRepository().size("hugeRoute") == 1);
+        LoopProcessor hugeLoop = context.getProcessor("hugeLoop", 
LoopProcessor.class);
+
+        // graceful shutdown must wait for the inflight exchange to complete
+        context.stop();
+
+        assertNull(huge.get(10, SECONDS).getException());
+        done.assertIsSatisfied();
+        // the loop broke out on shutdown, so the iterations left must no 
longer be pending
+        assertEquals(0, hugeLoop.getPendingExchangesSize());
+    }
+
+    private static boolean isStoppingOrStopped(Exchange exchange) {
+        return exchange.getContext().isStopping() || 
exchange.getContext().isStopped();
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:start").routeId("myRoute")
+                        .loop(header("n")).id("myLoop")
+                            .to("mock:loop")
+                        .end()
+                        .process(e -> {
+                            if ("slow".equals(e.getMessage().getBody())) {
+                                await().atMost(10, SECONDS).until(() -> 
isStoppingOrStopped(e));
+                            }
+                        })
+                        .to("mock:done");
+
+                from("direct:huge").routeId("hugeRoute")
+                        .loop(header("n")).breakOnShutdown().id("hugeLoop")
+                            .process(e -> await().atMost(10, SECONDS).until(() 
-> isStoppingOrStopped(e)))
+                        .end()
+                        .to("mock:hugeDone");
+            }
+        };
+    }
+}

Reply via email to