Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<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);
Expand Down Expand Up @@ -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
Expand All @@ -790,12 +791,14 @@ protected static int getPendingInflightExchanges(RouteStartupOrder order, boolea
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);
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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;
Expand All @@ -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);
}
}
Expand Down Expand Up @@ -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) {
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -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");
}
};
}
}
Loading