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 4fff8d5925ef CAMEL-25365: a reload that cuts off in-flight exchanges
says where they were waiting, and why they cannot continue (#27401)
4fff8d5925ef is described below
commit 4fff8d5925ef69a6b6b56f85a10632a105aa9082
Author: Claus Ibsen <[email protected]>
AuthorDate: Tue Oct 6 12:51:54 2026 +0200
CAMEL-25365: a reload that cuts off in-flight exchanges says where they
were waiting, and why they cannot continue (#27401)
* CAMEL-25365: a reload that cuts off in-flight exchanges says where they
were waiting, and why they cannot continue
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01STT6whBgK1AqsSsUKrnE8m
* CAMEL-25365: address review feedback
- direct: the cut-off exchange keeps the InterruptedException as the cause
of the
DirectConsumerNotAvailableException and stays marked as interrupted, so
the error handler
stops routing instead of handling a failure (onException, redelivery,
dead letter channel)
- the forced-shutdown reason names a CamelContext stop, the only stop that
forces
- Waiting at also finds an exchange that came in through direct: (the route
it is in now), and
leaves out entries that are not at a node yet
- tests: await the asserted state, the real shutdown log line (and none
with the verbose
listing), an exchange behind direct:, the three-place cap, the source
suffix and both
rejection reasons
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Signed-off-by: Claus Ibsen <[email protected]>
---------
Signed-off-by: Claus Ibsen <[email protected]>
Co-authored-by: Claude Opus 5.5 (1M context) <[email protected]>
---
.../camel/component/direct/DirectProducer.java | 16 +-
.../camel/impl/engine/DefaultShutdownStrategy.java | 46 +++++
.../errorhandler/RedeliveryErrorHandler.java | 15 +-
.../camel/impl/engine/ShutdownWaitingAtTest.java | 226 +++++++++++++++++++++
.../RedeliveryNotAllowedReasonTest.java | 69 +++++++
5 files changed, 368 insertions(+), 4 deletions(-)
diff --git
a/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java
b/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java
index 3ff25eeacf7c..4f7adf8e5f1c 100644
---
a/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java
+++
b/components/camel-direct/src/main/java/org/apache/camel/component/direct/DirectProducer.java
@@ -99,9 +99,21 @@ public class DirectProducer extends DefaultAsyncProducer {
}
}
} catch (InterruptedException e) {
- LOG.info("Interrupted while processing the exchange");
+ // the only wait here is for a consumer to appear (block=true),
and what interrupts it is a forced shutdown,
+ // such as a dev mode reload that adds the consumer in the same
edit (CAMEL-25365)
+ LOG.info("Interrupted while waiting for a consumer on {}: the
route is being stopped or reloaded",
+ endpoint.getEndpointUri());
Thread.currentThread().interrupt();
- exchange.setException(e);
+ DirectConsumerNotAvailableException cause = new
DirectConsumerNotAvailableException(
+ "No consumers available on endpoint: " + endpoint
+
+ " (interrupted while waiting for one, as the route is being
stopped or reloaded)",
+ exchange);
+ // keep the interruption as the cause, so
onException(InterruptedException.class) still matches
+ cause.initCause(e);
+ exchange.setException(cause);
+ // stay marked as interrupted, as
setException(InterruptedException) did, so the error handler stops
+ // routing instead of handling a failure (onException, redelivery,
dead letter channel, the ERROR log)
+ exchange.getExchangeExtension().setInterrupted(true);
callback.done(true);
return true;
} catch (Exception e) {
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 036bdc0e64a4..0c7b058cbfc1 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
@@ -694,6 +694,10 @@ public class DefaultShutdownStrategy extends
ServiceSupport implements ShutdownS
+ (TimeUnit.SECONDS.convert(timeout,
timeUnit) - (loopCount++ * loopDelaySeconds))
+ " seconds.";
msg += inflightsBuilder.toString();
+ if (!logInflightExchangesOnTimeout) {
+ // the verbose listing is off (as in the dev
profile), so say in this line where they wait
+ msg += waitingAt(context, routes);
+ }
LOG.info(msg);
@@ -801,6 +805,48 @@ public class DefaultShutdownStrategy extends
ServiceSupport implements ShutdownS
return (int) Math.min(Integer.MAX_VALUE, inflight);
}
+ /**
+ * Where the inflight exchanges of the routes are, for the one-line
waiting message: route, node and the line in the
+ * source, such as {@code picked-lines/to1 (aggregator.camel.yaml:23)}. A
dev mode reload waits here for exchanges
+ * blocked on something only the reload itself would add, such as a direct
endpoint whose consumer is part of the
+ * same edit, and without the place the wait is a mystery (CAMEL-25365).
Empty when the inflight repository cannot
+ * be browsed.
+ */
+ static String waitingAt(CamelContext camelContext, List<RouteStartupOrder>
routes) {
+ if (!camelContext.getInflightRepository().isInflightBrowseEnabled()) {
+ return "";
+ }
+ Set<String> routeIds = new HashSet<>();
+ for (RouteStartupOrder route : routes) {
+ routeIds.add(route.getRoute().getId());
+ }
+ Set<String> places = new LinkedHashSet<>();
+ int more = 0;
+ for (InflightRepository.InflightExchange inflight :
camelContext.getInflightRepository().browse()) {
+ // the route the exchange was created by, or the route it is in
now: an exchange that came in through
+ // direct: keeps the route of its caller as its from route
+ if (!routeIds.contains(inflight.getExchange().getFromRouteId())
+ && !routeIds.contains(inflight.getAtRouteId())) {
+ continue;
+ }
+ // between two nodes, or before the route is entered, there is no
place to name yet
+ if (inflight.getAtRouteId() == null || inflight.getNodeId() ==
null) {
+ continue;
+ }
+ String place = inflight.getAtRouteId() + "/" + inflight.getNodeId()
+ + (inflight.getNodeSource() != null ? " (" +
inflight.getNodeSource() + ")" : "");
+ if (places.size() < 3 || places.contains(place)) {
+ places.add(place);
+ } else {
+ more++;
+ }
+ }
+ if (places.isEmpty()) {
+ return "";
+ }
+ return ". Waiting at: " + String.join(", ", places) + (more > 0 ? "
and " + more + " more" : "");
+ }
+
/**
* Logs information about the inflight exchanges
*
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java
index 5aee6db4e8fe..a6728ba7fe36 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/errorhandler/RedeliveryErrorHandler.java
@@ -858,7 +858,7 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
private void runNotAllowed() {
LOG.trace("Run not allowed, will reject executing exchange: {}",
exchange);
if (exchange.getException() == null) {
- exchange.setException(new RejectedExecutionException());
+ exchange.setException(new
RejectedExecutionException(notAllowedReason()));
}
AsyncCallback cb = callback;
taskFactory.release(this);
@@ -1117,7 +1117,7 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
if (!isRunAllowed()) {
LOG.trace("Run not allowed, will reject executing exchange:
{}", exchange);
if (exchange.getException() == null) {
- exchange.setException(new RejectedExecutionException());
+ exchange.setException(new
RejectedExecutionException(notAllowedReason()));
}
AsyncCallback cb = callback;
taskFactory.release(this);
@@ -2208,4 +2208,15 @@ public abstract class RedeliveryErrorHandler extends
ErrorHandlerSupport
return sb.toString();
}
+ /**
+ * Why an exchange cannot go on: its route is being stopped. Without a
message the error reads
+ * {@code RejectedExecutionException - null}, which looks like a fault in
the route, while it is a route stop or a
+ * dev mode reload cutting the exchange off (CAMEL-25365).
+ */
+ String notAllowedReason() {
+ return shutdownStrategy.isForceShutdown()
+ ? "The exchange cannot continue: its route was forced to shut
down, as the graceful shutdown timed out"
+ + " while it was in flight (the CamelContext is being
stopped)"
+ : "The exchange cannot continue: its route is being stopped or
reloaded";
+ }
}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/impl/engine/ShutdownWaitingAtTest.java
b/core/camel-core/src/test/java/org/apache/camel/impl/engine/ShutdownWaitingAtTest.java
new file mode 100644
index 000000000000..3e7e95683c1e
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/impl/engine/ShutdownWaitingAtTest.java
@@ -0,0 +1,226 @@
+/*
+ * 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.impl.engine;
+
+import java.util.List;
+import java.util.Queue;
+import java.util.concurrent.ConcurrentLinkedQueue;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.direct.DirectConsumerNotAvailableException;
+import org.apache.camel.component.log.ConsumingAppender;
+import org.apache.camel.spi.CamelEvent;
+import org.apache.camel.spi.RouteStartupOrder;
+import org.apache.camel.support.EventNotifierSupport;
+import org.apache.logging.log4j.Level;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.core.LoggerContext;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * CAMEL-25365: a route stopped while its exchanges wait for a direct
consumer, as a dev mode reload that adds the
+ * consumer in the same edit does. The shutdown says where they wait, and an
interrupted exchange says why it failed.
+ */
+public class ShutdownWaitingAtTest extends ContextTestSupport {
+
+ private static final String SHUTDOWN_LOGGER =
DefaultShutdownStrategy.class.getName();
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext answer = super.createCamelContext();
+ // the place of a node includes where it is in the source
+ answer.setSourceLocationEnabled(true);
+ return answer;
+ }
+
+ @AfterEach
+ public void removeAppender() {
+ LoggerContext ctx = (LoggerContext) LogManager.getContext(false);
+ ctx.getConfiguration().removeLogger(SHUTDOWN_LOGGER);
+ ctx.updateLoggers();
+ }
+
+ @Test
+ public void theWaitSaysWhereTheExchangesAre() throws Exception {
+ context.getInflightRepository().setInflightBrowseEnabled(true);
+ context.getShutdownStrategy().setTimeout(2);
+ context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS);
+ context.getShutdownStrategy().setLogInflightExchangesOnTimeout(false);
+
+ template.sendBody("seda:start", "A");
+ awaitAt("shipment");
+
+ String at = DefaultShutdownStrategy.waitingAt(context,
context.getCamelContextExtension().getRouteStartupOrder());
+ assertTrue(at.startsWith(". Waiting at: picked-lines/shipment ("), at);
+
+ context.getRouteController().stopRoute("picked-lines");
+ }
+
+ @Test
+ public void theInterruptedExchangeSaysWhy() throws Exception {
+ AtomicReference<Exchange> failure = new AtomicReference<>();
+ context.getManagementStrategy().addEventNotifier(new
EventNotifierSupport() {
+ @Override
+ public void notify(CamelEvent event) {
+ if (event instanceof CamelEvent.ExchangeFailedEvent failed) {
+ failure.compareAndSet(null, failed.getExchange());
+ }
+ }
+ });
+ context.getInflightRepository().setInflightBrowseEnabled(true);
+ context.getShutdownStrategy().setTimeout(1);
+ context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS);
+
+ template.sendBody("seda:start", "A");
+ awaitAt("shipment");
+ // a consumer thread is blocked waiting for the direct consumer; the
forced shutdown interrupts it
+ context.getRouteController().stopRoute("picked-lines");
+
+ await().atMost(5, TimeUnit.SECONDS).until(() -> failure.get() != null);
+ Exchange exchange = failure.get();
+ DirectConsumerNotAvailableException e
+ = assertInstanceOf(DirectConsumerNotAvailableException.class,
exchange.getException());
+ assertTrue(e.getMessage().contains("direct://shipment"),
e.getMessage());
+ assertTrue(e.getMessage().contains("interrupted while waiting for one,
as the route is being stopped or reloaded"),
+ e.getMessage());
+ // still an interruption: onException(InterruptedException.class)
matches through the cause, and the error
+ // handler stops routing instead of handling a failure
+ assertInstanceOf(InterruptedException.class, e.getCause());
+ assertSame(e.getCause(),
exchange.getException(InterruptedException.class));
+ assertTrue(exchange.getExchangeExtension().isInterrupted());
+ }
+
+ @Test
+ public void theShutdownLogSaysWhereTheExchangesAre() throws Exception {
+ Queue<String> messages = new ConcurrentLinkedQueue<>();
+ ConsumingAppender.newAppender(SHUTDOWN_LOGGER,
"ShutdownWaitingAtTest", Level.INFO,
+ event ->
messages.add(event.getMessage().getFormattedMessage()));
+ context.getInflightRepository().setInflightBrowseEnabled(true);
+ context.getShutdownStrategy().setTimeout(2);
+ context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS);
+ context.getShutdownStrategy().setLogInflightExchangesOnTimeout(false);
+
+ template.sendBody("seda:start", "A");
+ awaitAt("shipment");
+ context.getRouteController().stopRoute("picked-lines");
+
+ assertTrue(messages.stream().anyMatch(m -> m.startsWith("Waiting as
there are still 1 inflight")
+ && m.contains(". Waiting at: picked-lines/shipment (")),
messages.toString());
+ }
+
+ @Test
+ public void theVerboseListingLeavesThePlaceOut() throws Exception {
+ Queue<String> messages = new ConcurrentLinkedQueue<>();
+ ConsumingAppender.newAppender(SHUTDOWN_LOGGER,
"ShutdownWaitingAtTest", Level.INFO,
+ event ->
messages.add(event.getMessage().getFormattedMessage()));
+ context.getInflightRepository().setInflightBrowseEnabled(true);
+ context.getShutdownStrategy().setTimeout(2);
+ context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS);
+ // the default: the inflight exchanges are listed in full, so the
one-line place is not added
+
+ template.sendBody("seda:start", "A");
+ awaitAt("shipment");
+ context.getRouteController().stopRoute("picked-lines");
+
+ assertTrue(messages.stream().anyMatch(m -> m.startsWith("Waiting as
there are still 1 inflight")),
+ messages.toString());
+ assertTrue(messages.stream().noneMatch(m -> m.contains("Waiting
at:")), messages.toString());
+ }
+
+ @Test
+ public void anExchangeThatCameInThroughDirectIsFoundInTheStoppedRoute()
throws Exception {
+ context.getInflightRepository().setInflightBrowseEnabled(true);
+ context.getShutdownStrategy().setTimeout(1);
+ context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS);
+
+ // created by the caller route, now waiting in the sub route
+ template.sendBody("seda:caller", "A");
+ awaitAt("deep");
+
+ String at = DefaultShutdownStrategy.waitingAt(context,
startupOrder("sub"));
+ assertTrue(at.startsWith(". Waiting at: sub/deep ("), at);
+ }
+
+ @Test
+ public void atMostThreePlacesAreNamed() throws Exception {
+ context.getInflightRepository().setInflightBrowseEnabled(true);
+ context.getShutdownStrategy().setTimeout(1);
+ context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS);
+
+ for (int i = 1; i <= 4; i++) {
+ template.sendBody("seda:w" + i, "A");
+ awaitAt("wait" + i);
+ }
+
+ String at = DefaultShutdownStrategy.waitingAt(context,
context.getCamelContextExtension().getRouteStartupOrder());
+ assertTrue(at.startsWith(". Waiting at: "), at);
+ assertTrue(at.endsWith(" and 1 more"), at);
+ // three places named, such as w1/wait1 (...), the fourth counted
+ assertEquals(3, at.split("/wait", -1).length - 1, at);
+ }
+
+ private void awaitAt(String nodeId) {
+ // the node id is set once the exchange is at the node, after it was
added to the inflight repository
+ await().atMost(5, TimeUnit.SECONDS).until(() ->
context.getInflightRepository().browse().stream()
+ .anyMatch(inflight -> nodeId.equals(inflight.getNodeId())));
+ }
+
+ private List<RouteStartupOrder> startupOrder(String routeId) {
+ return
context.getCamelContextExtension().getRouteStartupOrder().stream()
+ .filter(order -> routeId.equals(order.getRoute().getId()))
+ .toList();
+ }
+
+ @Test
+ public void nothingWhenTheRepositoryCannotBeBrowsed() {
+ assertEquals("",
+ DefaultShutdownStrategy.waitingAt(context,
context.getCamelContextExtension().getRouteStartupOrder()));
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("seda:start").routeId("picked-lines")
+ .to("direct:shipment").id("shipment");
+
+ from("seda:caller").routeId("caller")
+ .to("direct:sub");
+ from("direct:sub").routeId("sub")
+ .to("direct:missing").id("deep");
+
+ for (int i = 1; i <= 4; i++) {
+ from("seda:w" + i).routeId("w" + i)
+ .to("direct:missing" + i).id("wait" + i);
+ }
+ }
+ };
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/errorhandler/RedeliveryNotAllowedReasonTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/errorhandler/RedeliveryNotAllowedReasonTest.java
new file mode 100644
index 000000000000..305449c47818
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/errorhandler/RedeliveryNotAllowedReasonTest.java
@@ -0,0 +1,69 @@
+/*
+ * 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.errorhandler;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.impl.engine.DefaultShutdownStrategy;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+/**
+ * CAMEL-25365: the RejectedExecutionException of an exchange that cannot
continue says why, instead of
+ * {@code RejectedExecutionException - null}.
+ */
+public class RedeliveryNotAllowedReasonTest extends ContextTestSupport {
+
+ private final ForcedShutdownStrategy shutdownStrategy = new
ForcedShutdownStrategy();
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext answer = super.createCamelContext();
+ answer.setShutdownStrategy(shutdownStrategy);
+ return answer;
+ }
+
+ @Test
+ public void aStoppedOrReloadedRoute() {
+ assertEquals("The exchange cannot continue: its route is being stopped
or reloaded",
+ errorHandler().notAllowedReason());
+ }
+
+ @Test
+ public void aForcedShutdown() {
+ shutdownStrategy.forced = true;
+ assertEquals("The exchange cannot continue: its route was forced to
shut down, as the graceful shutdown timed out"
+ + " while it was in flight (the CamelContext is being
stopped)",
+ errorHandler().notAllowedReason());
+ }
+
+ private DefaultErrorHandler errorHandler() {
+ return new DefaultErrorHandler(context, exchange -> {
+ }, null, null, new RedeliveryPolicy(), null, null, null, null);
+ }
+
+ private static final class ForcedShutdownStrategy extends
DefaultShutdownStrategy {
+
+ private boolean forced;
+
+ @Override
+ public boolean isForceShutdown() {
+ return forced;
+ }
+ }
+}