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 76ded18cb294 CAMEL-25038: camel-disruptor - request/reply: ignore a
late reply, and fail the exchange when interrupted (#26913)
76ded18cb294 is described below
commit 76ded18cb294eb4da752a1f00d6d2a71dbdd7e3d
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 13:20:57 2026 +0530
CAMEL-25038: camel-disruptor - request/reply: ignore a late reply, and fail
the exchange when interrupted (#26913)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../component/disruptor/DisruptorProducer.java | 68 +++++--
.../DisruptorProducerInterruptedTest.java | 108 +++++++++++
.../disruptor/DisruptorTimeoutLateReplyTest.java | 201 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 7 +
4 files changed, 365 insertions(+), 19 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 f55907945b28..51169bbdf1b1 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
@@ -19,6 +19,7 @@ package org.apache.camel.component.disruptor;
import java.io.IOException;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
import com.lmax.disruptor.InsufficientCapacityException;
import org.apache.camel.AsyncCallback;
@@ -85,9 +86,11 @@ public class DisruptorProducer extends DefaultAsyncProducer {
// latch that waits until we are complete
final CountDownLatch latch = new CountDownLatch(1);
+ // either the response, the timeout or an interrupt completes
the exchange, whichever claims it first
+ final AtomicBoolean completed = new AtomicBoolean();
// we should wait for the reply so install a on completion so
we know when its complete
-
copy.getExchangeExtension().addOnCompletion(newOnCompletion(exchange, latch));
+
copy.getExchangeExtension().addOnCompletion(newOnCompletion(exchange, latch,
completed));
doPublish(copy);
@@ -105,21 +108,24 @@ public class DisruptorProducer extends
DefaultAsyncProducer {
Thread.currentThread().interrupt();
}
if (!done) {
- // Remove timed out Exchange from disruptor endpoint.
-
- // We can't actually remove a published exchange from
an active Disruptor.
- // Instead we prevent processing of the exchange by
setting a Property on the exchange and the value
- // would be an AtomicBoolean. This is set by the
Producer and the Consumer would look up that Property and
- // check the AtomicBoolean. If the AtomicBoolean says
that we are good to proceed, it will process the
- // exchange. If false, it will simply disregard the
exchange.
- // But since the Property map is a Concurrent one,
maybe we don't need the AtomicBoolean. Check with Simon.
- // Also check the TimeoutHandler of the new Disruptor
3.0.0, consider making the switch to the latest version.
-
exchange.setProperty(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE, true);
-
- exchange.setException(new
ExchangeTimedOutException(exchange, timeout));
-
- // count down to indicate timeout
- latch.countDown();
+ if (completed.compareAndSet(false, true)) {
+ // Remove timed out Exchange from disruptor
endpoint.
+
+ // We can't actually remove a published exchange
from an active Disruptor.
+ // Instead we prevent processing of the exchange
by setting a Property on the exchange and the value
+ // would be an AtomicBoolean. This is set by the
Producer and the Consumer would look up that Property and
+ // check the AtomicBoolean. If the AtomicBoolean
says that we are good to proceed, it will process the
+ // exchange. If false, it will simply disregard
the exchange.
+ // But since the Property map is a Concurrent one,
maybe we don't need the AtomicBoolean. Check with Simon.
+ // Also check the TimeoutHandler of the new
Disruptor 3.0.0, consider making the switch to the latest version.
+
exchange.setProperty(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE, true);
+
+ exchange.setException(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)
+ awaitUninterruptibly(latch);
+ }
}
} else {
if (LOG.isTraceEnabled()) {
@@ -130,6 +136,15 @@ public class DisruptorProducer extends
DefaultAsyncProducer {
latch.await();
} catch (InterruptedException e) {
LOG.info("Interrupted while waiting for the task to
complete");
+ 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);
+ } 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)
+ awaitUninterruptibly(latch);
+ }
Thread.currentThread().interrupt();
}
}
@@ -149,12 +164,27 @@ public class DisruptorProducer extends
DefaultAsyncProducer {
return true;
}
- private SynchronizationAdapter newOnCompletion(Exchange exchange,
CountDownLatch latch) {
+ private static void awaitUninterruptibly(CountDownLatch latch) {
+ boolean interrupted = false;
+ while (true) {
+ try {
+ latch.await();
+ break;
+ } catch (InterruptedException e) {
+ interrupted = true;
+ }
+ }
+ if (interrupted) {
+ Thread.currentThread().interrupt();
+ }
+ }
+
+ private SynchronizationAdapter newOnCompletion(Exchange exchange,
CountDownLatch latch, AtomicBoolean completed) {
return new SynchronizationAdapter() {
@Override
public void onDone(final Exchange response) {
- // check for timeout, which then already would have invoked
the latch
- if (latch.getCount() == 0) {
+ // check for timeout, which then already has completed the
exchange
+ if (!completed.compareAndSet(false, true)) {
if (LOG.isTraceEnabled()) {
LOG.trace("{}. Timeout occurred so response will be
ignored: {}", this, response.getMessage());
}
diff --git
a/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorProducerInterruptedTest.java
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorProducerInterruptedTest.java
new file mode 100644
index 000000000000..5044ba9601ff
--- /dev/null
+++
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorProducerInterruptedTest.java
@@ -0,0 +1,108 @@
+/*
+ * 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.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+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.builder.RouteBuilder;
+import org.apache.camel.support.SynchronizationAdapter;
+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.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A disruptor producer that is interrupted while it waits for the reply must
fail the exchange, and not return the
+ * request as the reply.
+ */
+class DisruptorProducerInterruptedTest extends CamelTestSupport {
+
+ private final ExecutorService executor =
Executors.newSingleThreadExecutor();
+ private final CountDownLatch releaseConsumer = new CountDownLatch(1);
+ private final CountDownLatch consumerDone = new CountDownLatch(1);
+
+ @AfterEach
+ void releaseAndShutdown() {
+ releaseConsumer.countDown();
+ executor.shutdownNow();
+ }
+
+ @Test
+ void testInterruptedWhileWaitingForReply() throws Exception {
+ AtomicReference<Thread> sender = new AtomicReference<>();
+ Future<Exchange> future = executor.submit(() -> {
+ sender.set(Thread.currentThread());
+ Exchange exchange =
context.getEndpoint("disruptor:slow").createExchange(ExchangePattern.InOut);
+ exchange.getMessage().setBody("request");
+ return template.send("disruptor:slow?timeout=0", exchange);
+ });
+ // wait until the sender waits for the reply in the disruptor producer
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
isWaitingInDisruptorProducer(sender.get()));
+ sender.get().interrupt();
+ Exchange out = future.get(10, TimeUnit.SECONDS);
+
+ // the request must not be returned as the reply
+ assertInstanceOf(InterruptedException.class, out.getException());
+
+ // and the reply from the consumer is ignored when it completes later
+ releaseConsumer.countDown();
+ assertTrue(consumerDone.await(10, TimeUnit.SECONDS));
+ assertInstanceOf(InterruptedException.class, out.getException());
+ assertEquals("request", out.getMessage().getBody());
+ }
+
+ private static boolean isWaitingInDisruptorProducer(Thread thread) {
+ if (thread == null || thread.getState() != Thread.State.WAITING) {
+ return false;
+ }
+ for (StackTraceElement element : thread.getStackTrace()) {
+ if
(DisruptorProducer.class.getName().equals(element.getClassName())) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("disruptor:slow").routeId("slow")
+ .process(e ->
e.getExchangeExtension().addOnCompletion(new SynchronizationAdapter() {
+ @Override
+ public void onDone(Exchange exchange) {
+ consumerDone.countDown();
+ }
+ }))
+ .process(e -> releaseConsumer.await(20,
TimeUnit.SECONDS))
+ .setBody(constant("reply"));
+ }
+ };
+ }
+}
diff --git
a/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorTimeoutLateReplyTest.java
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorTimeoutLateReplyTest.java
new file mode 100644
index 000000000000..66d058a47da9
--- /dev/null
+++
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorTimeoutLateReplyTest.java
@@ -0,0 +1,201 @@
+/*
+ * 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.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicReference;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.ExchangePattern;
+import org.apache.camel.ExchangeTimedOutException;
+import org.apache.camel.SafeCopyProperty;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.support.SynchronizationAdapter;
+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.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A reply that arrives while the disruptor producer times out, or is
interrupted, must either be returned to the
+ * caller, or be ignored. It must never be copied into the caller's exchange
after the producer has returned with the
+ * timeout (or the interruption).
+ */
+class DisruptorTimeoutLateReplyTest extends CamelTestSupport {
+
+ private final ExecutorService executor =
Executors.newSingleThreadExecutor();
+ private final CountDownLatch copyStarted = new CountDownLatch(1);
+ private final CountDownLatch releaseCopy = new CountDownLatch(1);
+ private final CountDownLatch releaseConsumer = new CountDownLatch(1);
+ private final CountDownLatch consumerDone = new CountDownLatch(1);
+
+ /**
+ * Pauses the thread which copies the reply into the caller's exchange, in
the middle of the copy.
+ */
+ private final class PauseCopy implements SafeCopyProperty {
+ private final AtomicBoolean armed = new AtomicBoolean(true);
+
+ @Override
+ public SafeCopyProperty safeCopy() {
+ // the consumer copies the exchange as well, only pause when the
producer's reply copy is in progress
+ if (isCopyingReply() && armed.getAndSet(false)) {
+ copyStarted.countDown();
+ try {
+ releaseCopy.await(20, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ return this;
+ }
+
+ private static boolean isCopyingReply() {
+ for (StackTraceElement element :
Thread.currentThread().getStackTrace()) {
+ if
(element.getClassName().startsWith(DisruptorProducer.class.getName())) {
+ return true;
+ }
+ }
+ return false;
+ }
+ }
+
+ @AfterEach
+ void releaseAndShutdown() {
+ releaseCopy.countDown();
+ releaseConsumer.countDown();
+ executor.shutdownNow();
+ }
+
+ @Test
+ void testReplyBeingCopiedWhenTimeoutOccurs() throws Exception {
+ AtomicReference<Thread> caller = new AtomicReference<>();
+ Future<Exchange> future = executor.submit(() -> {
+ caller.set(Thread.currentThread());
+ Exchange exchange =
context.getEndpoint("disruptor:reply").createExchange(ExchangePattern.InOut);
+ exchange.getMessage().setBody("request");
+ return template.send("disruptor:reply?timeout=1000", exchange);
+ });
+
+ // the consumer is copying its reply into the caller's exchange
+ assertTrue(copyStarted.await(10, TimeUnit.SECONDS));
+ // let the timeout occur while the copy is in progress: either the
producer returns
+ // or it waits (without timeout) for the copy to complete
+ await().atMost(10, TimeUnit.SECONDS)
+ .until(() -> future.isDone() || caller.get().getState() ==
Thread.State.WAITING);
+ boolean returnedBeforeCopyCompleted = future.isDone();
+ releaseCopy.countDown();
+
+ Exchange out = future.get(10, TimeUnit.SECONDS);
+ // the reply won the race, so the caller gets the complete reply and
no timeout
+ assertFalse(returnedBeforeCopyCompleted, "Producer returned while the
reply was copied into the exchange");
+ assertNull(out.getException());
+ assertEquals("reply", out.getMessage().getBody());
+ }
+
+ @Test
+ void testReplyBeingCopiedWhenInterrupted() throws Exception {
+ AtomicReference<Thread> caller = new AtomicReference<>();
+ AtomicBoolean interruptedAfterSend = new AtomicBoolean();
+ Future<Exchange> future = executor.submit(() -> {
+ caller.set(Thread.currentThread());
+ Exchange exchange =
context.getEndpoint("disruptor:reply").createExchange(ExchangePattern.InOut);
+ exchange.getMessage().setBody("request");
+ Exchange answer = template.send("disruptor:reply?timeout=0",
exchange);
+ interruptedAfterSend.set(Thread.currentThread().isInterrupted());
+ return answer;
+ });
+
+ // the consumer is copying its reply into the caller's exchange
+ assertTrue(copyStarted.await(10, TimeUnit.SECONDS));
+ // interrupt the producer while the copy is in progress: either the
producer returns
+ // or it waits (uninterruptibly) for the copy to complete
+ caller.get().interrupt();
+ await().atMost(10, TimeUnit.SECONDS)
+ .until(() -> future.isDone() ||
isWaitingForCopy(caller.get()));
+ boolean returnedBeforeCopyCompleted = future.isDone();
+ releaseCopy.countDown();
+
+ Exchange out = future.get(10, TimeUnit.SECONDS);
+ // the reply won the race, so the caller gets the complete reply, and
the interrupt status is kept
+ assertFalse(returnedBeforeCopyCompleted, "Producer returned while the
reply was copied into the exchange");
+ assertNull(out.getException());
+ assertEquals("reply", out.getMessage().getBody());
+ assertTrue(interruptedAfterSend.get(), "The interrupt status of the
caller should be kept");
+ }
+
+ private static boolean isWaitingForCopy(Thread thread) {
+ if (thread.getState() != Thread.State.WAITING) {
+ return false;
+ }
+ for (StackTraceElement element : thread.getStackTrace()) {
+ if
(DisruptorProducer.class.getName().equals(element.getClassName())
+ && "awaitUninterruptibly".equals(element.getMethodName()))
{
+ return true;
+ }
+ }
+ return false;
+ }
+
+ @Test
+ void testReplyAfterTimeoutIsIgnored() throws Exception {
+ Exchange exchange =
context.getEndpoint("disruptor:late").createExchange(ExchangePattern.InOut);
+ exchange.getMessage().setBody("request");
+ Exchange out = template.send("disruptor:late?timeout=100", exchange);
+ assertInstanceOf(ExchangeTimedOutException.class, out.getException());
+
+ // now the consumer completes after the timeout
+ releaseConsumer.countDown();
+ assertTrue(consumerDone.await(10, TimeUnit.SECONDS));
+
+ assertInstanceOf(ExchangeTimedOutException.class, out.getException());
+ assertEquals("request", out.getMessage().getBody());
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("disruptor:reply").routeId("reply")
+ .setBody(constant("reply"))
+ // copying the reply into the caller's exchange pauses
in the middle of the copy
+ .process(e ->
e.getExchangeExtension().setSafeCopyProperty("pause", new PauseCopy()));
+
+ from("disruptor:late").routeId("late")
+ .process(e -> releaseConsumer.await(20,
TimeUnit.SECONDS))
+ .setBody(constant("late reply"))
+ .process(e ->
e.getExchangeExtension().addOnCompletion(new SynchronizationAdapter() {
+ @Override
+ public void onDone(Exchange exchange) {
+ consumerDone.countDown();
+ }
+ }));
+ }
+ };
+ }
+}
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 c2e8e94988d0..7c178afb2631 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
@@ -2601,6 +2601,13 @@ custom `headerFilterStrategy` is used as-is and is
unaffected.
Routes that relied on one of those headers reaching the wire must set it
through the endpoint
configuration or supply a `headerFilterStrategy` that permits it.
+=== camel-disruptor - an interrupted request/reply producer fails the exchange
+
+A Disruptor producer whose thread is interrupted while it waits for the reply
without a timeout (`waitForTaskToComplete`
+with `timeout=0`) now fails the exchange with the `InterruptedException`,
where previously it reported the exchange as
+successful and returned the request as the reply. A reply that arrives after
the producer timed out, or was interrupted,
+is now ignored; previously it could still be copied into the caller's exchange
after the producer had returned.
+
=== camel-seda - an interrupted producer fails the exchange
A SEDA producer whose thread is interrupted while it waits now fails the
exchange, where previously it reported the send