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");
+ }
+ };
+ }
+}