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