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 42f9284a04c2 CAMEL-24947: camel-seda - do not copy a late reply into
the exchange after the producer timed out (#26793)
42f9284a04c2 is described below
commit 42f9284a04c2e6ff785aec25ec32ef24f4778bb9
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 17:50:23 2026 +0530
CAMEL-24947: camel-seda - do not copy a late reply into the exchange after
the producer timed out (#26793)
Cause: with waitForTaskToComplete and a timeout, SedaProducer's onDone
checked latch.getCount() == 0 and then copied the reply into the caller's
exchange, while the producer thread could time out, set
ExchangeTimedOutException, count down the latch and return in between. The
check was not atomic with the timeout branch.
Effect: a reply arriving at about the timeout could be returned together
with ExchangeTimedOutException, or be copied into the caller's exchange
after the producer had already returned.
Fix: the reply and the timeout claim the exchange with a compareAndSet on a
shared flag. A reply that loses the claim is ignored; a timeout that loses
the claim waits for the reply copy to complete and returns the reply. The
interrupted reply wait claims the exchange the same way, so a reply after
an interrupt no longer overwrites the InterruptedException.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../apache/camel/component/seda/SedaProducer.java | 52 ++++--
.../component/seda/SedaTimeoutLateReplyTest.java | 195 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 3 +
3 files changed, 237 insertions(+), 13 deletions(-)
diff --git
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java
index 96090c55a377..5037c554b650 100644
---
a/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java
+++
b/components/camel-seda/src/main/java/org/apache/camel/component/seda/SedaProducer.java
@@ -21,6 +21,7 @@ import java.util.concurrent.BlockingQueue;
import java.util.concurrent.CountDownLatch;
import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
import org.apache.camel.AsyncCallback;
import org.apache.camel.Exchange;
@@ -69,13 +70,15 @@ public class SedaProducer extends DefaultAsyncProducer {
// latch that waits until we are complete
final CountDownLatch latch = new CountDownLatch(1);
+ // either the response or the timeout 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(new
SynchronizationAdapter() {
@Override
public void onDone(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());
}
@@ -128,11 +131,15 @@ public class SedaProducer extends DefaultAsyncProducer {
Thread.currentThread().interrupt();
}
if (!done) {
- exchange.setException(new
ExchangeTimedOutException(exchange, timeout));
- // remove timed out Exchange from queue
- endpoint.getQueue().remove(copy);
- // count down to indicate timeout
- latch.countDown();
+ if (completed.compareAndSet(false, true)) {
+ exchange.setException(new
ExchangeTimedOutException(exchange, timeout));
+ // remove timed out Exchange from queue
+ endpoint.getQueue().remove(copy);
+ } 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()) {
@@ -143,13 +150,17 @@ public class SedaProducer extends DefaultAsyncProducer {
latch.await();
} catch (InterruptedException e) {
LOG.debug("Interrupted while waiting for task to complete
at [{}]", endpoint.getEndpointUri());
+ if (completed.compareAndSet(false, true)) {
+ // the task has not completed so fail the exchange (do
not return the request as the reply)
+ exchange.setException(e);
+ // remove the Exchange from queue (if not yet
processed), and a later reply is ignored
+ endpoint.getQueue().remove(copy);
+ } 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();
- // the task has not completed so fail the exchange (do not
return the request as the reply)
- exchange.setException(e);
- // remove the Exchange from queue (if not yet processed)
- endpoint.getQueue().remove(copy);
- // count down to indicate the reply must be ignored
- latch.countDown();
}
}
} else {
@@ -169,6 +180,21 @@ public class SedaProducer extends DefaultAsyncProducer {
return true;
}
+ 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();
+ }
+ }
+
protected Exchange prepareCopy(Exchange exchange, boolean handover) {
// use a new copy of the exchange to route async (and use same message
id)
// if handover we need to do special handover to avoid handing over
diff --git
a/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaTimeoutLateReplyTest.java
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaTimeoutLateReplyTest.java
new file mode 100644
index 000000000000..e2e7c368c1c4
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/component/seda/SedaTimeoutLateReplyTest.java
@@ -0,0 +1,195 @@
+/*
+ * 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.seda;
+
+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.ContextTestSupport;
+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.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 seda 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).
+ */
+public class SedaTimeoutLateReplyTest extends ContextTestSupport {
+
+ 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() {
+ if (armed.getAndSet(false)) {
+ copyStarted.countDown();
+ try {
+ releaseCopy.await(20, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ }
+ }
+ return this;
+ }
+ }
+
+ @Test
+ public void testReplyBeingCopiedWhenTimeoutOccurs() throws Exception {
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ AtomicReference<Thread> caller = new AtomicReference<>();
+ try {
+ Future<Exchange> future = executor.submit(() -> {
+ caller.set(Thread.currentThread());
+ Exchange exchange =
context.getEndpoint("seda:reply").createExchange(ExchangePattern.InOut);
+ exchange.getMessage().setBody("request");
+ // the timeout must not occur before the consumer starts
copying the reply, also on a slow machine
+ return template.send("seda: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());
+ } finally {
+ releaseCopy.countDown();
+ executor.shutdownNow();
+ }
+ }
+
+ @Test
+ public void testReplyBeingCopiedWhenInterrupted() throws Exception {
+ ExecutorService executor = Executors.newSingleThreadExecutor();
+ AtomicReference<Thread> caller = new AtomicReference<>();
+ AtomicBoolean interruptedAfterSend = new AtomicBoolean();
+ try {
+ Future<Exchange> future = executor.submit(() -> {
+ caller.set(Thread.currentThread());
+ Exchange exchange =
context.getEndpoint("seda:reply").createExchange(ExchangePattern.InOut);
+ exchange.getMessage().setBody("request");
+ Exchange answer = template.send("seda: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");
+ } finally {
+ releaseCopy.countDown();
+ executor.shutdownNow();
+ }
+ }
+
+ private static boolean isWaitingForCopy(Thread thread) {
+ if (thread.getState() != Thread.State.WAITING) {
+ return false;
+ }
+ for (StackTraceElement element : thread.getStackTrace()) {
+ if (SedaProducer.class.getName().equals(element.getClassName())
+ && "awaitUninterruptibly".equals(element.getMethodName()))
{
+ return true;
+ }
+ }
+ return false;
+ }
+
+ @Test
+ public void testReplyAfterTimeoutIsIgnored() throws Exception {
+ Exchange exchange =
context.getEndpoint("seda:late").createExchange(ExchangePattern.InOut);
+ exchange.getMessage().setBody("request");
+ Exchange out = template.send("seda: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("seda: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("seda: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 77a23e316e09..6000e7a460a9 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
@@ -2701,6 +2701,9 @@ as successful. When it waits for space in a full queue
(`blockWhenFull=true`, wi
as the message was not added to the queue. When it waits for the reply without
a timeout (`waitForTaskToComplete` with
`timeout=0`), the exchange fails with the `InterruptedException`, instead of
returning the request as the reply.
Camel itself interrupts such threads when a route is forced to stop after the
graceful shutdown timeout.
+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. If the
reply is already being copied when the timeout
+or the interrupt occurs, the producer waits for the copy to complete and
returns the reply.
=== camel-servlet, camel-jetty - the multipart upload whitelist is enforced
against the submitted file name