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 bac864951fb9 CAMEL-25082: camel-disruptor - mark the published copy,
not the caller's exchange, as ignored on a request/reply timeout (#26980)
bac864951fb9 is described below
commit bac864951fb9da5d6cae569c8e0c0c76bd729ea0
Author: Claus Ibsen <[email protected]>
AuthorDate: Mon Sep 28 22:41:33 2026 +0200
CAMEL-25082: camel-disruptor - mark the published copy, not the caller's
exchange, as ignored on a request/reply timeout (#26980)
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 63241b1bed7c..9ac80ff53e7d 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
@@ -92,6 +92,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);
@@ -110,17 +114,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
@@ -195,6 +190,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();
@@ -235,6 +232,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 0c424094d549..d68b32838222 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
@@ -2832,6 +2832,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