This is an automated email from the ASF dual-hosted git repository. davsclaus pushed a commit to branch fix/CAMEL-25502-seda-disruptor in repository https://gitbox.apache.org/repos/asf/camel.git
commit 0321a390ba57c5c5698391db758b684eef86ec08 Author: Claus Ibsen <[email protected]> AuthorDate: Fri Oct 9 23:51:39 2026 +0200 CAMEL-25502: camel-seda, camel-disruptor - a reply wait interrupted by a stop is a RouteStoppingException A seda or disruptor producer waiting for the reply when its route or the CamelContext was stopped had its wait interrupted, and failed the exchange with ExchangeTimedOutException although the timeout had not passed (or the bare InterruptedException without a timeout). It now sets a RouteStoppingException with the InterruptedException as cause when the route or CamelContext is stopping (ExchangeHelper.isRouteStopping), and the InterruptedException otherwise. A seda producer interrupted while adding to a full queue during a stop sets a RouteStoppingException as well. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]> Claude-Session: https://claude.ai/code/session_01STT6whBgK1AqsSsUKrnE8m --- .../component/disruptor/DisruptorProducer.java | 26 +++++- .../disruptor/DisruptorReplyRouteStopTest.java | 93 ++++++++++++++++++++++ .../apache/camel/component/seda/SedaProducer.java | 39 +++++++-- .../component/seda/SedaReplyRouteStopTest.java | 92 +++++++++++++++++++++ .../org/apache/camel/support/ExchangeHelper.java | 25 ++++++ .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 7 ++ 6 files changed, 272 insertions(+), 10 deletions(-) diff --git a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorProducer.java b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorProducer.java index 9eaa481accb1..78a188ca1b50 100644 --- a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorProducer.java +++ b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorProducer.java @@ -26,6 +26,7 @@ import org.apache.camel.AsyncCallback; import org.apache.camel.Exchange; import org.apache.camel.ExchangePropertyKey; import org.apache.camel.ExchangeTimedOutException; +import org.apache.camel.RouteStoppingException; import org.apache.camel.StreamCache; import org.apache.camel.WaitForTaskToComplete; import org.apache.camel.support.DefaultAsyncProducer; @@ -106,17 +107,22 @@ public class DisruptorProducer extends DefaultAsyncProducer { } // lets see if we can get the task done before the timeout boolean done = false; + InterruptedException interrupted = null; try { done = latch.await(timeout, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { LOG.info("Interrupted while waiting for the task to complete"); + interrupted = e; Thread.currentThread().interrupt(); } if (!done) { if (completed.compareAndSet(false, true)) { // We can't remove a published exchange from an active Disruptor, but the consumer - // ignores the copy, if it has not started it yet, as we have claimed the completed flag - exchange.setException(new ExchangeTimedOutException(exchange, timeout)); + // ignores the copy, if it has not started it yet, as we have claimed the completed flag; + // an interruption is not a timeout: a stop cut it off, or something else interrupted the wait + exchange.setException(interrupted != null + ? interruptedWhileWaiting(exchange, interrupted) + : new ExchangeTimedOutException(exchange, timeout)); } else { // the response is being copied into the exchange, so wait for the copy to complete // (the exchange must not be changed after we have returned) @@ -135,7 +141,7 @@ public class DisruptorProducer extends DefaultAsyncProducer { if (completed.compareAndSet(false, true)) { // the task has not completed so fail the exchange (do not return the request as the reply), // and a later reply is ignored - exchange.setException(e); + exchange.setException(interruptedWhileWaiting(exchange, e)); } else { // the response is being copied into the exchange, so wait for the copy to complete // (the exchange must not be changed after we have returned) @@ -168,6 +174,20 @@ public class DisruptorProducer extends DefaultAsyncProducer { return true; } + /** + * The exception of an exchange whose wait for the reply is interrupted: a cut-off when the route or CamelContext is + * being stopped (CAMEL-25502), else the interruption. + */ + private Exception interruptedWhileWaiting(Exchange exchange, InterruptedException cause) { + if (ExchangeHelper.isRouteStopping(exchange)) { + return new RouteStoppingException( + "Interrupted while waiting for the reply from " + getEndpoint().getEndpointBaseUri() + + ", as the route is being stopped", + cause); + } + return cause; + } + private static void awaitUninterruptibly(CountDownLatch latch) { boolean interrupted = false; while (true) { diff --git a/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorReplyRouteStopTest.java b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorReplyRouteStopTest.java new file mode 100644 index 000000000000..27b0462adbb8 --- /dev/null +++ b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorReplyRouteStopTest.java @@ -0,0 +1,93 @@ +/* + * 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.component.disruptor; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.camel.Exchange; +import org.apache.camel.ExchangePattern; +import org.apache.camel.RouteStoppingException; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.spi.CamelEvent; +import org.apache.camel.support.EventNotifierSupport; +import org.apache.camel.test.junit6.CamelTestSupport; +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.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * CAMEL-25502: a disruptor producer waiting for the reply when its route is stopped is cut off by the stop: a + * RouteStoppingException, not an ExchangeTimedOutException. + */ +public class DisruptorReplyRouteStopTest extends CamelTestSupport { + + private final CountDownLatch release = new CountDownLatch(1); + private final CountDownLatch slowStarted = new CountDownLatch(1); + + @AfterEach + void releaseSlowRoute() { + release.countDown(); + } + + @Test + void theWaitForTheReplyIsCutOffByTheRouteStop() throws Exception { + AtomicReference<Exchange> failure = new AtomicReference<>(); + context.getManagementStrategy().addEventNotifier(new EventNotifierSupport() { + @Override + public void notify(CamelEvent event) { + if (event instanceof CamelEvent.ExchangeFailedEvent failed + && "caller".equals(failed.getExchange().getFromRouteId())) { + failure.compareAndSet(null, failed.getExchange()); + } + } + }); + context.getShutdownStrategy().setTimeout(1); + context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS); + + template.sendBody("seda:start", "A"); + assertTrue(slowStarted.await(5, TimeUnit.SECONDS)); + // the caller waits for the reply of the slow route; the forced stop of the caller interrupts the wait + context.getRouteController().stopRoute("caller"); + + await().atMost(5, TimeUnit.SECONDS).until(() -> failure.get() != null); + RouteStoppingException e = assertInstanceOf(RouteStoppingException.class, failure.get().getException()); + assertTrue(e.getMessage().startsWith("Interrupted while waiting for the reply from disruptor://slow"), + e.getMessage()); + assertInstanceOf(InterruptedException.class, e.getCause()); + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("seda:start").routeId("caller") + .to(ExchangePattern.InOut, "disruptor:slow"); + from("disruptor:slow").routeId("slow") + .process(e -> { + slowStarted.countDown(); + release.await(10, TimeUnit.SECONDS); + }); + } + }; + } +} diff --git a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java index cd4e94808f69..10f6898950d7 100644 --- a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java +++ b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java @@ -27,6 +27,7 @@ import org.apache.camel.AsyncCallback; import org.apache.camel.Exchange; import org.apache.camel.ExchangePropertyKey; import org.apache.camel.ExchangeTimedOutException; +import org.apache.camel.RouteStoppingException; import org.apache.camel.StreamCache; import org.apache.camel.WaitForTaskToComplete; import org.apache.camel.support.DefaultAsyncProducer; @@ -126,14 +127,19 @@ public class SedaProducer extends DefaultAsyncProducer { } // lets see if we can get the task done before the timeout boolean done = false; + InterruptedException interrupted = null; try { done = latch.await(timeout, TimeUnit.MILLISECONDS); } catch (InterruptedException e) { + interrupted = e; Thread.currentThread().interrupt(); } if (!done) { if (completed.compareAndSet(false, true)) { - exchange.setException(new ExchangeTimedOutException(exchange, timeout)); + // an interruption is not a timeout: a stop cut it off, or something else interrupted the wait + exchange.setException(interrupted != null + ? interruptedWhileWaiting(exchange, interrupted) + : new ExchangeTimedOutException(exchange, timeout)); // remove timed out Exchange from queue endpoint.getQueue().remove(copy); } else { @@ -153,7 +159,7 @@ public class SedaProducer extends DefaultAsyncProducer { LOG.debug("Interrupted while waiting for task to complete at [{}]", endpoint.getEndpointUri()); if (completed.compareAndSet(false, true)) { // the task has not completed so fail the exchange (do not return the request as the reply) - exchange.setException(e); + exchange.setException(interruptedWhileWaiting(exchange, e)); // remove the Exchange from queue (if not yet processed), and a later reply is ignored endpoint.getQueue().remove(copy); } else { @@ -292,7 +298,7 @@ public class SedaProducer extends DefaultAsyncProducer { } catch (InterruptedException e) { LOG.debug("Offer interrupted, are we stopping? {}", isStopping() || isStopped()); Thread.currentThread().interrupt(); - throw interruptedWhileAddingToQueue(e); + throw interruptedWhileAddingToQueue(target, e); } } else if (blockWhenFull && offerTimeout == 0) { try { @@ -300,7 +306,7 @@ public class SedaProducer extends DefaultAsyncProducer { } catch (InterruptedException e) { LOG.debug("Put interrupted, are we stopping? {}", isStopping() || isStopped()); Thread.currentThread().interrupt(); - throw interruptedWhileAddingToQueue(e); + throw interruptedWhileAddingToQueue(target, e); } } else if (blockWhenFull && offerTimeout > 0) { try { @@ -313,7 +319,7 @@ public class SedaProducer extends DefaultAsyncProducer { } catch (InterruptedException e) { LOG.debug("Offer interrupted, are we stopping? {}", isStopping() || isStopped()); Thread.currentThread().interrupt(); - throw interruptedWhileAddingToQueue(e); + throw interruptedWhileAddingToQueue(target, e); } } else { queue.add(target); @@ -321,9 +327,28 @@ public class SedaProducer extends DefaultAsyncProducer { return true; } - private static RejectedExecutionException interruptedWhileAddingToQueue(InterruptedException cause) { - // the exchange was not added to the queue, so the exchange must fail + private static RejectedExecutionException interruptedWhileAddingToQueue(Exchange exchange, InterruptedException cause) { + // the exchange was not added to the queue, so the exchange must fail; a stop that interrupts it cuts it off + // (CAMEL-25502) + if (ExchangeHelper.isRouteStopping(exchange)) { + return new RouteStoppingException( + "Interrupted while adding the exchange to the queue, as the route is being stopped", cause); + } return new RejectedExecutionException("Interrupted while adding the exchange to the queue", cause); } + /** + * The exception of an exchange whose wait for the reply is interrupted: a cut-off when the route or CamelContext is + * being stopped (CAMEL-25502), else the interruption. + */ + private Exception interruptedWhileWaiting(Exchange exchange, InterruptedException cause) { + if (ExchangeHelper.isRouteStopping(exchange)) { + return new RouteStoppingException( + "Interrupted while waiting for the reply from " + endpoint.getEndpointBaseUri() + + ", as the route is being stopped", + cause); + } + return cause; + } + } diff --git a/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaReplyRouteStopTest.java b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaReplyRouteStopTest.java new file mode 100644 index 000000000000..02df93984719 --- /dev/null +++ b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaReplyRouteStopTest.java @@ -0,0 +1,92 @@ +/* + * 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.component.seda; + +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicReference; + +import org.apache.camel.ContextTestSupport; +import org.apache.camel.Exchange; +import org.apache.camel.ExchangePattern; +import org.apache.camel.RouteStoppingException; +import org.apache.camel.builder.RouteBuilder; +import org.apache.camel.spi.CamelEvent; +import org.apache.camel.support.EventNotifierSupport; +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.assertInstanceOf; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * CAMEL-25502: a seda producer waiting for the reply when its route is stopped is cut off by the stop: a + * RouteStoppingException, not an ExchangeTimedOutException. + */ +public class SedaReplyRouteStopTest extends ContextTestSupport { + + private final CountDownLatch release = new CountDownLatch(1); + private final CountDownLatch slowStarted = new CountDownLatch(1); + + @AfterEach + public void releaseSlowRoute() { + release.countDown(); + } + + @Test + public void theWaitForTheReplyIsCutOffByTheRouteStop() throws Exception { + AtomicReference<Exchange> failure = new AtomicReference<>(); + context.getManagementStrategy().addEventNotifier(new EventNotifierSupport() { + @Override + public void notify(CamelEvent event) { + if (event instanceof CamelEvent.ExchangeFailedEvent failed + && "caller".equals(failed.getExchange().getFromRouteId())) { + failure.compareAndSet(null, failed.getExchange()); + } + } + }); + context.getShutdownStrategy().setTimeout(1); + context.getShutdownStrategy().setTimeUnit(TimeUnit.SECONDS); + + template.sendBody("seda:start", "A"); + assertTrue(slowStarted.await(5, TimeUnit.SECONDS)); + // the caller waits for the reply of the slow route; the forced stop of the caller interrupts the wait + context.getRouteController().stopRoute("caller"); + + await().atMost(5, TimeUnit.SECONDS).until(() -> failure.get() != null); + RouteStoppingException e = assertInstanceOf(RouteStoppingException.class, failure.get().getException()); + assertTrue(e.getMessage().startsWith("Interrupted while waiting for the reply from seda://slow"), e.getMessage()); + assertInstanceOf(InterruptedException.class, e.getCause()); + } + + @Override + protected RouteBuilder createRouteBuilder() { + return new RouteBuilder() { + @Override + public void configure() { + from("seda:start").routeId("caller") + .to(ExchangePattern.InOut, "seda:slow"); + from("seda:slow").routeId("slow") + .process(e -> { + slowStarted.countDown(); + release.await(10, TimeUnit.SECONDS); + }); + } + }; + } +} diff --git a/core/camel-support/src/main/java/org/apache/camel/support/ExchangeHelper.java b/core/camel-support/src/main/java/org/apache/camel/support/ExchangeHelper.java index 30dca3abe915..bd23381a783e 100644 --- a/core/camel-support/src/main/java/org/apache/camel/support/ExchangeHelper.java +++ b/core/camel-support/src/main/java/org/apache/camel/support/ExchangeHelper.java @@ -51,6 +51,7 @@ import org.apache.camel.NoTypeConversionAvailableException; import org.apache.camel.Route; import org.apache.camel.RuntimeCamelException; import org.apache.camel.SafeCopyProperty; +import org.apache.camel.StatefulService; import org.apache.camel.StreamCache; import org.apache.camel.TypeConversionException; import org.apache.camel.VariableAware; @@ -454,6 +455,30 @@ public final class ExchangeHelper { } } + /** + * Whether the exchange is in flight while it is being stopped: the CamelContext is stopping, or the consumer of the + * route the exchange came from is being suspended or stopped (a route stop, a route reload). A thread waiting for + * such an exchange that is interrupted is then cut off by the stop + * ({@link org.apache.camel.RouteStoppingException}), not failing (CAMEL-25502). + */ + public static boolean isRouteStopping(Exchange exchange) { + CamelContext context = exchange.getContext(); + if (context.isStopping() || context.isStopped() || context.getShutdownStrategy().isForceShutdown()) { + return true; + } + String routeId = exchange.getFromRouteId(); + if (routeId == null) { + return false; + } + Route route = context.getRoute(routeId); + if (route == null) { + // the route is already removed + return true; + } + return route.getConsumer() instanceof StatefulService consumer + && (consumer.isSuspending() || consumer.isSuspended() || consumer.isStopping() || consumer.isStopped()); + } + /** * Returns true if the given exchange pattern (if defined) can support OUT messages * diff --git a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc index f022c9640d33..2d722bde5b2f 100644 --- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc +++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc @@ -1562,6 +1562,13 @@ A `direct` producer that was waiting for a consumer when a forced stop interrupt `RouteStoppingException` (with the `InterruptedException` as its cause) instead of a `DirectConsumerNotAvailableException`. +A `seda` or `disruptor` producer waiting for the reply (InOut, or `waitForTaskToComplete`) that is interrupted while +its route or the CamelContext is being stopped now sets a `RouteStoppingException` (with the `InterruptedException` +as its cause). Before, it set an `ExchangeTimedOutException` (or, without a timeout, the `InterruptedException`), +although the timeout had not passed. An interruption at any other time sets the `InterruptedException` instead of an +`ExchangeTimedOutException`. A `seda` producer interrupted while adding to a full queue (`blockWhenFull`) during a +stop sets a `RouteStoppingException` as well. + Such a cut-off is not recorded in the error registry, as nothing in the route failed; set `camel.errorRegistry.includeRouteStopping=true` to record it as well, for example to see how often stops cut off work that a consumer which does not roll back may lose.
