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 551a793e4e1e CAMEL-25502: camel-seda, camel-disruptor - a reply wait
interrupted by a stop is a RouteStoppingException (#27635)
551a793e4e1e is described below
commit 551a793e4e1efbbd5449abb0cf02360ecb9edfb7
Author: Claus Ibsen <[email protected]>
AuthorDate: Sat Oct 10 00:38:52 2026 +0200
CAMEL-25502: camel-seda, camel-disruptor - a reply wait interrupted by a
stop is a RouteStoppingException (#27635)
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.
Claude-Session: https://claude.ai/code/session_01STT6whBgK1AqsSsUKrnE8m
Co-authored-by: Claude Opus 5.5 (1M context) <[email protected]>
---
.../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 dafc371e0f27..f77ec1f362c5 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.