From f920f79425a100d571fb4f0048a4e83b53d9d06f Mon Sep 17 00:00:00 2001 From: smjain <49463903+allthingssecurity@users.noreply.github.com> Date: Thu, 24 Sep 2026 20:42:46 +0530 Subject: [PATCH] CAMEL-25002: camel-core - Loop EIP: a negative or huge loop count must not break graceful shutdown LoopProcessor added the evaluated loop count to its pending task counter without checking the sign. A count of zero or less runs no iteration, and nothing compensated a negative add, so after a single message with, for example, loop(header("n")) and n=-1, getPendingExchangesSize() returned -1 for the rest of the route's life. This regressed in CAMEL-16794, which switched to a LongAdder with an unconditional add(count) and dropped the clamp CAMEL-15578 had. DefaultShutdownStrategy adds the pending sizes of the route's services to its inflight count, in int, and only waits while the sum is positive. The same LoopProcessor is a child of several route services and is counted once for each of them. The negative value cancelled real inflight exchanges, so a graceful shutdown stopped the route immediately and the inflight exchange failed with a RejectedExecutionException. A huge loop count (from about 2^29) had the same effect, because the int sum overflowed. A loop that breaks on shutdown also left its remaining iterations pending, as the release added in CAMEL-19738 only ran when the loop stopped due to an exception, so every later shutdown waited for its full timeout. LoopProcessor now only adds a positive count, releases the iterations left exactly once whenever the loop ends (normally, on an exception, or when it breaks on shutdown), and getPendingExchangesSize() no longer wraps or goes negative. DefaultShutdownStrategy sums the pending sizes in a long, ignores negative sizes, and caps the result at Integer.MAX_VALUE. Co-Authored-By: Claude Opus 5.5 --- .../impl/engine/DefaultShutdownStrategy.java | 15 ++- .../apache/camel/processor/LoopProcessor.java | 34 +++-- .../LoopPendingExchangesShutdownTest.java | 125 ++++++++++++++++++ 3 files changed, 154 insertions(+), 20 deletions(-) create mode 100644 core/camel-core/src/test/java/org/apache/camel/processor/LoopPendingExchangesShutdownTest.java 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 d69112208646c..036bdc0e64a45 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 void run() { long loopDelaySeconds = 1; long loopCount = 0; while (!done && !timeoutOccurred.get()) { - int size = 0; + long size = 0; // number of inflights per route final Map 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 @@ protected static int getPendingInflightExchanges(RouteStartupOrder order) { * @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 @@ protected static int getPendingInflightExchanges(RouteStartupOrder order, boolea Set 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 179c0de64ef85..3a01705c934ba 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 boolean deferShutdown(ShutdownRunningTask shutdownRunningTask) { @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 @@ class LoopState implements Runnable { Exchange current; int index; int count; + boolean pendingTasksReleased; public LoopState(Exchange exchange, AsyncCallback callback) throws NoTypeConversionAvailableException { this.exchange = exchange; @@ -126,7 +128,10 @@ public LoopState(Exchange exchange, AsyncCallback callback) throws NoTypeConvers 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 void run() { 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 @@ private void handleException(Exception e) { 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 0000000000000..cc6622f4ae683 --- /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 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 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"); + } + }; + } +}