This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch fix/CAMEL-25082
in repository https://gitbox.apache.org/repos/asf/camel.git

commit 602a4a25411a963190ed859b4c5c6da93b76577d
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 10:58:43 2026 +0200

    CAMEL-25082: camel-disruptor - mark the published copy, not the caller's 
exchange, as ignored on a request/reply timeout
    
    Cause: on a request/reply timeout DisruptorProducer set the
    disruptor.ignoreExchange property on the caller's exchange instead of the
    copy it had published into the ring buffer. The consumer ignored any
    exchange that had the property (containsKey, whatever the value).
    
    Effect: a timed out copy that the consumer had not started yet was still
    processed, and every later copy of the caller's exchange inherited the
    property and was dropped by the consumer: a redelivery after the timeout
    timed out again, and an InOnly fallback to another disruptor endpoint was
    lost silently.
    
    Fix: the producer puts its completed flag (the AtomicBoolean that the reply,
    the timeout and the interrupt already claim) on the published copy before
    publishing it, so the consumer ignores the copy only when the producer no
    longer waits for it. The flag is not copied into the routed exchange, not
    copied back into the caller's exchange with the reply, and not inherited
    by a new copy. The consumer and the reconfiguration buffer check the value
    of the flag. camel-seda is not affected, as SedaProducer removes the timed
    out copy from its queue.
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../component/disruptor/DisruptorConsumer.java     |   8 +-
 .../component/disruptor/DisruptorEndpoint.java     |  20 +++
 .../component/disruptor/DisruptorProducer.java     |  21 ++-
 .../component/disruptor/DisruptorReference.java    |   6 +-
 .../DisruptorTimeoutIgnoreExchangeTest.java        | 141 +++++++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |   6 +
 6 files changed, 183 insertions(+), 19 deletions(-)

diff --git 
a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorConsumer.java
 
b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorConsumer.java
index 4ca64d31b962..adf7823fa77d 100644
--- 
a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorConsumer.java
+++ 
b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorConsumer.java
@@ -136,6 +136,8 @@ public class DisruptorConsumer extends ServiceSupport 
implements Consumer, Suspe
         // send a new copied exchange with new camel context
         // don't copy handovers as they are handled by the Disruptor Event 
Handlers
         final Exchange newExchange = 
ExchangeHelper.copyExchangeWithProperties(exchange, endpoint.getCamelContext());
+        // the flag is only for this consumer, and must not be routed (or 
copied back to the caller)
+        
newExchange.removeProperty(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE);
         // set the from endpoint
         newExchange.getExchangeExtension().setFromEndpoint(endpoint);
         return newExchange;
@@ -145,10 +147,8 @@ public class DisruptorConsumer extends ServiceSupport 
implements Consumer, Suspe
         try {
             Exchange exchange = synchronizedExchange.getExchange();
 
-            final boolean ignore = exchange.hasProperties() && exchange
-                    
.getProperties().containsKey(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE);
-            if (ignore) {
-                // Property was set and it was set to true, so don't process 
Exchange.
+            if (DisruptorEndpoint.isIgnoreExchange(exchange)) {
+                // the producer no longer waits for this exchange (timeout), 
so don't process it
                 LOGGER.trace("Ignoring exchange {}", exchange);
                 return;
             }
diff --git 
a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorEndpoint.java
 
b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorEndpoint.java
index fdc769ec82bf..ab83f1e0dfc1 100644
--- 
a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorEndpoint.java
+++ 
b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorEndpoint.java
@@ -22,6 +22,7 @@ import java.util.HashMap;
 import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.CopyOnWriteArraySet;
+import java.util.concurrent.atomic.AtomicBoolean;
 
 import com.lmax.disruptor.InsufficientCapacityException;
 import org.apache.camel.AsyncEndpoint;
@@ -53,6 +54,11 @@ import org.slf4j.LoggerFactory;
 @UriEndpoint(firstVersion = "2.12.0", scheme = "disruptor,disruptor-vm", title 
= "Disruptor,Disruptor VM",
              remote = false, syntax = "disruptor:name", category = { 
Category.MESSAGING })
 public class DisruptorEndpoint extends DefaultEndpoint implements 
AsyncEndpoint, MultipleConsumersSupport {
+    /**
+     * Property on the exchange published by a producer that waits for the 
reply. Its value is set to true when the
+     * producer no longer waits (timeout or interrupt), so the consumer 
ignores the exchange if it has not started it
+     * yet.
+     */
     public static final String DISRUPTOR_IGNORE_EXCHANGE = 
"disruptor.ignoreExchange";
     private static final Logger LOGGER = 
LoggerFactory.getLogger(DisruptorEndpoint.class);
 
@@ -366,4 +372,18 @@ public class DisruptorEndpoint extends DefaultEndpoint 
implements AsyncEndpoint,
     public int hashCode() {
         return getEndpointUri().hashCode() * 37 + getCamelContext().hashCode();
     }
+
+    /**
+     * Whether the consumer should ignore the given exchange, as the producer 
that published it no longer waits for it.
+     */
+    static boolean isIgnoreExchange(Exchange exchange) {
+        if (!exchange.hasProperties()) {
+            return false;
+        }
+        Object ignore = exchange.getProperty(DISRUPTOR_IGNORE_EXCHANGE);
+        if (ignore instanceof AtomicBoolean flag) {
+            return flag.get();
+        }
+        return Boolean.TRUE.equals(ignore);
+    }
 }
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 51169bbdf1b1..158981fe31e9 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
@@ -91,6 +91,10 @@ public class DisruptorProducer extends DefaultAsyncProducer {
 
                 // we should wait for the reply so install a on completion so 
we know when its complete
                 
copy.getExchangeExtension().addOnCompletion(newOnCompletion(exchange, latch, 
completed));
+                // the consumer ignores the copy if we no longer wait for it 
(timeout or interrupt) before it is processed,
+                // this must be on the published copy, not on the exchange of 
the caller, as later copies of the caller's
+                // exchange (such as a redelivery) would inherit it and be 
ignored as well
+                copy.setProperty(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE, 
completed);
 
                 doPublish(copy);
 
@@ -109,17 +113,8 @@ public class DisruptorProducer extends 
DefaultAsyncProducer {
                     }
                     if (!done) {
                         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);
-
+                            // 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));
                         } else {
                             // the response is being copied into the exchange, 
so wait for the copy to complete
@@ -194,6 +189,8 @@ public class DisruptorProducer extends DefaultAsyncProducer 
{
                     }
                     try {
                         ExchangeHelper.copyResults(exchange, response);
+                        // the flag of the published copy must not be copied 
back to the caller
+                        
exchange.removeProperty(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE);
                     } finally {
                         // always ensure latch is triggered
                         latch.countDown();
@@ -234,6 +231,8 @@ public class DisruptorProducer extends DefaultAsyncProducer 
{
     private Exchange prepareCopy(final Exchange exchange, final boolean copy) 
throws IOException {
         // use a new copy of the exchange to route async
         final Exchange target = ExchangeHelper.createCorrelatedCopy(exchange, 
copy);
+        // a flag from an earlier send must not be inherited, it is set for 
this send only when we wait for the reply
+        target.removeProperty(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE);
         // set a new from endpoint to be the disruptor
         target.getExchangeExtension().setFromEndpoint(endpoint);
         if (copy) {
diff --git 
a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorReference.java
 
b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorReference.java
index ae466834554a..4317c87a047f 100644
--- 
a/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorReference.java
+++ 
b/components/camel-disruptor/src/main/java/org/apache/camel/component/disruptor/DisruptorReference.java
@@ -437,10 +437,8 @@ public class DisruptorReference {
             blockingLatch.await();
             final Exchange exchange = 
event.getSynchronizedExchange().cancelAndGetOriginalExchange();
 
-            final boolean ignoreExchange
-                    = 
exchange.getProperty(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE, false, 
boolean.class);
-            if (ignoreExchange) {
-                // Property was set and it was set to true, so don't process 
Exchange.
+            if (DisruptorEndpoint.isIgnoreExchange(exchange)) {
+                // the producer no longer waits for this exchange (timeout), 
so don't process it
                 LOGGER.trace("Ignoring exchange {}", exchange);
             } else {
                 temporaryExchangeBuffer.offer(exchange);
diff --git 
a/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorTimeoutIgnoreExchangeTest.java
 
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorTimeoutIgnoreExchangeTest.java
new file mode 100644
index 000000000000..5f0e60c90c2a
--- /dev/null
+++ 
b/components/camel-disruptor/src/test/java/org/apache/camel/component/disruptor/DisruptorTimeoutIgnoreExchangeTest.java
@@ -0,0 +1,141 @@
+/*
+ * 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.AtomicInteger;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.ExchangePattern;
+import org.apache.camel.ExchangeTimedOutException;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+/**
+ * When a request/reply send to a disruptor endpoint times out, the consumer 
must ignore the timed out exchange if it
+ * has not started it yet, and later sends of the caller's exchange (a 
redelivery, a fallback) must not be ignored.
+ */
+class DisruptorTimeoutIgnoreExchangeTest extends CamelTestSupport {
+
+    private final CountDownLatch releaseSlow = new CountDownLatch(1);
+    private final CountDownLatch releaseBusy = new CountDownLatch(1);
+    private final AtomicInteger attempts = new AtomicInteger();
+
+    @AfterEach
+    void release() {
+        releaseSlow.countDown();
+        releaseBusy.countDown();
+    }
+
+    @Test
+    void testRedeliveryAfterTimeoutIsProcessed() {
+        // the first attempt times out, the redelivery must reach the consumer 
and get the reply
+        Object reply = template.requestBody("direct:redeliver", "hello");
+        assertEquals("reply to attempt 2", reply);
+    }
+
+    @Test
+    void testInOnlyFallbackAfterTimeoutIsProcessed() throws Exception {
+        getMockEndpoint("mock:fallback").expectedBodiesReceived("hello");
+
+        template.requestBody("direct:fallback", "hello");
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Test
+    void testTimedOutExchangeNotStartedIsIgnored() throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:busy");
+        mock.expectedBodiesReceived("A");
+        // the timed out B must not be processed after A
+        mock.setAssertPeriod(1000);
+
+        // A keeps the (single) consumer busy
+        template.sendBody("disruptor:busy", "A");
+        // B waits in the ring buffer until it times out
+        Exchange out = template.send("disruptor:busy?timeout=250", 
ExchangePattern.InOut,
+                e -> e.getMessage().setBody("B"));
+        assertInstanceOf(ExchangeTimedOutException.class, out.getException());
+        releaseBusy.countDown();
+
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Test
+    void testReplyDoesNotMarkCallerExchange() throws Exception {
+        getMockEndpoint("mock:after").expectedBodiesReceived("echo hello");
+
+        Exchange out = template.send("direct:twice", ExchangePattern.InOut, e 
-> e.getMessage().setBody("hello"));
+
+        assertNull(out.getException());
+        
assertNull(out.getProperty(DisruptorEndpoint.DISRUPTOR_IGNORE_EXCHANGE));
+        MockEndpoint.assertIsSatisfied(context);
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:redeliver")
+                        
.errorHandler(defaultErrorHandler().maximumRedeliveries(2).redeliveryDelay(0))
+                        .to("disruptor:slow?timeout=250");
+
+                from("disruptor:slow?concurrentConsumers=2")
+                        .process(e -> {
+                            int attempt = attempts.incrementAndGet();
+                            if (attempt == 1) {
+                                releaseSlow.await(10, TimeUnit.SECONDS);
+                            }
+                            e.getMessage().setBody("reply to attempt " + 
attempt);
+                        });
+
+                from("direct:fallback")
+                        .doTry()
+                            .to("disruptor:slow2?timeout=250")
+                        .doCatch(ExchangeTimedOutException.class)
+                            .to(ExchangePattern.InOnly, "disruptor:fallback")
+                        .end();
+
+                from("disruptor:slow2")
+                        .process(e -> releaseSlow.await(10, TimeUnit.SECONDS));
+
+                from("disruptor:fallback").to("mock:fallback");
+
+                from("disruptor:busy")
+                        .process(e -> releaseBusy.await(10, TimeUnit.SECONDS))
+                        .to("mock:busy");
+
+                from("direct:twice")
+                        .to("disruptor:echo")
+                        .to(ExchangePattern.InOnly, "disruptor:after");
+
+                from("disruptor:echo").setBody(simple("echo ${body}"));
+
+                from("disruptor:after").to("mock:after");
+            }
+        };
+    }
+}
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 77a23e316e09..c8c2b98a980e 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
@@ -2693,6 +2693,12 @@ with `timeout=0`) now fails the exchange with the 
`InterruptedException`, where
 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.
 
+A request/reply message that timed out (or whose producer was interrupted) 
before a Disruptor consumer started it is
+now ignored by the consumer, as intended. Previously the consumer processed it 
anyway, and instead the exchange of the
+caller was marked as ignored: a redelivery of that exchange after the timeout, 
or a fallback that sends it to another
+Disruptor endpoint, was then dropped by the consumer, so the redelivery timed 
out again and an InOnly fallback message
+was lost.
+
 === 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