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
