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

Reply via email to