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

Reply via email to