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 29cb79b27965 CAMEL-25095: camel-jms - do not complete an InOut
exchange twice when the send fails after the timeout or the reply (#26996)
29cb79b27965 is described below
commit 29cb79b27965f1e9270035e2a746895bdb3fee45
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 12:24:11 2026 +0530
CAMEL-25095: camel-jms - do not complete an InOut exchange twice when the
send fails after the timeout or the reply (#26996)
JmsProducer.processInOut registers the reply handler before the JMS send.
When the request timeout or the reply removed the handler while send() was
still running, the exchange was already completed, and a send that then
failed completed it a second time: the on-completions ran as success and
then as failure, the caller got the send exception although the route had
handled the timeout, and the inflight count went negative. CAMEL-24073
only covered a send that fails before the timeout.
ReplyManager.cancelCorrelationId now returns whether it removed the pending
handler. Whoever removes the handler owns the completion of the exchange,
so when the cancel returns false the send failure is logged and the
exchange is left to the timeout or the reply. The return type change is
noted in the 4.23 upgrade guide.
CAMEL-25095: camel-jms - cancel the correlation id the reply handler was
moved to with useMessageIDAsCorrelationID
With useMessageIDAsCorrelationID the MessageSentCallback moves the reply
handler from the provisional correlation id to the JMSMessageID before the
transacted session is committed. When the commit then failed, the send
failure cancelled the provisional id, which was no longer registered, so
the exchange was left to the request timeout (or never completed with
requestTimeout=0) instead of failing with the commit failure. Keep the
registered correlation id in sync with the id the handler was moved to.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../apache/camel/component/jms/JmsProducer.java | 27 ++-
.../camel/component/jms/reply/ReplyManager.java | 10 +-
.../component/jms/reply/ReplyManagerSupport.java | 4 +-
.../jms/JmsInOutSendFailureCallbackTest.java | 204 ++++++++++++++++++++-
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 11 ++
5 files changed, 242 insertions(+), 14 deletions(-)
diff --git
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/JmsProducer.java
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/JmsProducer.java
index f3b2ddf36cc2..5ded35072985 100644
---
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/JmsProducer.java
+++
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/JmsProducer.java
@@ -202,10 +202,21 @@ public class JmsProducer extends DefaultAsyncProducer {
// this is done with the help of the MessageSentCallback
final boolean msgIdAsCorrId =
configuration.isUseMessageIDAsCorrelationID();
final String provisionalCorrelationId = msgIdAsCorrId ?
getUuidGenerator().generateUuid() : null;
+ // the correlation id the reply handler is currently registered under,
which is cancelled if the send fails
+ final String[] registeredCorrelationId = new String[1];
MessageSentCallback messageSentCallback = null;
if (msgIdAsCorrId) {
- messageSentCallback
+ MessageSentCallback updateCorrelationId
= new
UseMessageIdAsCorrelationIdMessageSentCallback(replyManager,
provisionalCorrelationId, timeout);
+ messageSentCallback = (session, message, dest) -> {
+ updateCorrelationId.sent(session, message, dest);
+ // the reply handler has been moved to the JMSMessageID, which
is what must be cancelled if the send
+ // fails afterwards (such as the commit of a transacted
session)
+ String messageId = JmsMessageHelper.getJMSMessageID(message);
+ if (messageId != null) {
+ registeredCorrelationId[0] = messageId;
+ }
+ };
}
final String correlationProperty =
configuration.getCorrelationProperty();
@@ -221,8 +232,6 @@ public class JmsProducer extends DefaultAsyncProducer {
in.setHeader(correlationPropertyToUse,
GENERATED_CORRELATION_ID_PREFIX + getUuidGenerator().generateUuid());
}
- final String[] registeredCorrelationId = new String[1];
-
MessageCreator messageCreator = new MessageCreator() {
public Message createMessage(Session session) throws JMSException {
Message answer =
endpoint.getBinding().makeJmsMessage(exchange, in, session, null);
@@ -262,8 +271,16 @@ public class JmsProducer extends DefaultAsyncProducer {
try {
doSend(exchange, true, destinationName, destination,
messageCreator, messageSentCallback);
} catch (Exception e) {
- // send failed after reply was registered, cancel to prevent
double callback on timeout
- replyManager.cancelCorrelationId(registeredCorrelationId[0]);
+ // send failed after the reply was registered: cancel it to
prevent a second callback from the timeout
+ String registered = registeredCorrelationId[0];
+ if (registered != null &&
!replyManager.cancelCorrelationId(registered)) {
+ // the request timeout (or the reply) removed the correlation
while the send was still running,
+ // and it completes the exchange, so the exchange must not be
completed here a second time
+ LOG.warn("Sending JMS request with correlation id: {} failed
after the request timeout or the reply"
+ + " has already completed the exchange. The send
failure is only logged: {}",
+ registered, e.getMessage(), e);
+ return false;
+ }
throw e;
}
diff --git
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/ReplyManager.java
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/ReplyManager.java
index c95de55f8ede..2dbe522dc88d 100644
---
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/ReplyManager.java
+++
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/ReplyManager.java
@@ -106,10 +106,16 @@ public interface ReplyManager extends
SessionAwareMessageListener {
* <p/>
* This is used when the JMS send fails after the reply has been
registered, to prevent the timeout handler from
* firing a second callback on an already-completed exchange.
+ * <p/>
+ * Whoever removes the correlation owns the completion of the exchange.
When this method returns <tt>false</tt>, the
+ * request timeout or the reply has already removed the correlation (for
example while a slow send was still
+ * running), and it completes the exchange: the caller must then not
complete the exchange as well.
*
- * @param correlationId the correlation id to cancel
+ * @param correlationId the correlation id to cancel
+ * @return <tt>true</tt> if the correlation was pending and
has been cancelled, <tt>false</tt> if it
+ * was not pending (anymore)
*/
- void cancelCorrelationId(String correlationId);
+ boolean cancelCorrelationId(String correlationId);
/**
* Process the reply
diff --git
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/ReplyManagerSupport.java
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/ReplyManagerSupport.java
index 870ad77d7c90..b6708a6650c5 100644
---
a/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/ReplyManagerSupport.java
+++
b/components/camel-jms/src/main/java/org/apache/camel/component/jms/reply/ReplyManagerSupport.java
@@ -137,13 +137,15 @@ public abstract class ReplyManagerSupport extends
ServiceSupport implements Repl
}
@Override
- public void cancelCorrelationId(String correlationId) {
+ public boolean cancelCorrelationId(String correlationId) {
if (correlationId != null && correlation != null) {
ReplyHandler handler = correlation.remove(correlationId);
if (handler != null) {
log.debug("Cancelled reply correlation [{}]", correlationId);
+ return true;
}
}
+ return false;
}
protected abstract ReplyHandler createReplyHandler(
diff --git
a/components/camel-jms/src/test/java/org/apache/camel/component/jms/JmsInOutSendFailureCallbackTest.java
b/components/camel-jms/src/test/java/org/apache/camel/component/jms/JmsInOutSendFailureCallbackTest.java
index e6ac4083d6f3..84b91f75dbbd 100644
---
a/components/camel-jms/src/test/java/org/apache/camel/component/jms/JmsInOutSendFailureCallbackTest.java
+++
b/components/camel-jms/src/test/java/org/apache/camel/component/jms/JmsInOutSendFailureCallbackTest.java
@@ -17,13 +17,20 @@
package org.apache.camel.component.jms;
import java.lang.reflect.InvocationHandler;
+import java.lang.reflect.InvocationTargetException;
+import java.lang.reflect.Method;
import java.lang.reflect.Proxy;
+import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
import jakarta.jms.Connection;
import jakarta.jms.ConnectionFactory;
+import jakarta.jms.Destination;
import jakarta.jms.JMSException;
import jakarta.jms.MessageProducer;
+import jakarta.jms.Queue;
import jakarta.jms.Session;
import org.apache.camel.CamelContext;
@@ -32,8 +39,10 @@ import org.apache.camel.ExchangePattern;
import org.apache.camel.ExchangeTimedOutException;
import org.apache.camel.ProducerTemplate;
import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.support.SynchronizationAdapter;
import org.apache.camel.test.infra.core.CamelContextExtension;
import org.apache.camel.test.infra.core.TransientCamelContextExtension;
+import org.apache.camel.util.StopWatch;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Order;
import org.junit.jupiter.api.Test;
@@ -41,18 +50,36 @@ import org.junit.jupiter.api.extension.RegisterExtension;
import static
org.apache.camel.component.jms.JmsComponent.jmsComponentAutoAcknowledge;
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.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertTrue;
/**
* Tests that when a JMS send fails after the reply correlation has been
registered, the AsyncCallback is invoked
- * exactly once (not a second time by the timeout handler).
+ * exactly once: not a second time by the timeout handler (CAMEL-24073), and
not a second time by the send failure when
+ * the request timeout or the reply has already completed the exchange while
the send was still running. With
+ * useMessageIDAsCorrelationID, a failure after the correlation was moved to
the JMSMessageID (a failing commit of a
+ * transacted session) must cancel the moved correlation and fail the exchange
with that failure.
*
* @see <a
href="https://issues.apache.org/jira/browse/CAMEL-24073">CAMEL-24073</a>
*/
public class JmsInOutSendFailureCallbackTest extends AbstractJMSTest {
+ private static final String FAIL_QUEUE = "JmsInOutSendFailureCallbackTest";
+ private static final String TIMEOUT_QUEUE =
"JmsInOutSendFailureCallbackTest.timeout";
+ private static final String REPLY_QUEUE =
"JmsInOutSendFailureCallbackTest.reply";
+ private static final String COMMIT_QUEUE =
"JmsInOutSendFailureCallbackTest.commit";
+ private static final String COMMIT_FAILURE = "Simulated commit failure:
transaction rolled back";
+
+ // counts the completions of the exchange, the send to TIMEOUT_QUEUE and
REPLY_QUEUE waits for the first one
+ private static final AtomicInteger COMPLETED = new AtomicInteger();
+ private static final AtomicInteger FAILED = new AtomicInteger();
+ private static volatile CountDownLatch exchangeDone = new
CountDownLatch(1);
+
@Order(2)
@RegisterExtension
public static CamelContextExtension camelContextExtension = new
TransientCamelContextExtension();
@@ -78,6 +105,82 @@ public class JmsInOutSendFailureCallbackTest extends
AbstractJMSTest {
"Exception changed after send failure - timeout
handler fired a second callback"));
}
+ @Test
+ public void testCallbackInvokedOnceWhenTimeoutFiresDuringFailingSend()
throws Exception {
+ // the send blocks until the request timeout has completed the
exchange, and then fails
+ Exchange result = template.send("direct:timeoutDuringSend",
ExchangePattern.InOut,
+ p -> p.getIn().setBody("Hello"));
+
+ assertInstanceOf(ExchangeTimedOutException.class,
result.getException(),
+ "Should fail with the timeout that completed the exchange, not
the later send failure");
+ assertCompletedOnce(0, 1);
+ }
+
+ @Test
+ public void
testCallbackInvokedOnceWhenHandledTimeoutFiresDuringFailingSend() throws
Exception {
+ // the route handles the timeout, so the late send failure must not
fail the exchange afterwards
+ Exchange result = template.send("direct:timeoutDuringSendHandled",
ExchangePattern.InOut,
+ p -> p.getIn().setBody("Hello"));
+
+ assertNull(result.getException(), "The timeout was handled by the
route");
+ assertEquals("Timed out", result.getMessage().getBody(String.class));
+ assertCompletedOnce(1, 0);
+ }
+
+ @Test
+ public void testCallbackInvokedOnceWhenReplyArrivesDuringFailingSend()
throws Exception {
+ // the request reaches the broker and is answered, and then the send
fails
+ Exchange result = template.send("direct:replyDuringSend",
ExchangePattern.InOut,
+ p -> p.getIn().setBody("Hello"));
+
+ assertNull(result.getException(), "The reply completed the exchange
before the send failed");
+ assertEquals("Bye World", result.getMessage().getBody(String.class));
+ assertCompletedOnce(1, 0);
+ }
+
+ @Test
+ public void
testCallbackInvokedOnceWhenCommitFailsAfterMessageIdCorrelationUpdate() throws
Exception {
+ // useMessageIDAsCorrelationID moves the reply handler to the
JMSMessageID when the message has been sent, and
+ // only then the transacted session is committed: when the commit
fails the request is rolled back and no reply
+ // comes, so the exchange must fail with the commit failure rather
than wait for the request timeout
+ StopWatch watch = new StopWatch();
+ Exchange result = template.send("direct:commitFailure",
ExchangePattern.InOut,
+ p -> p.getIn().setBody("Hello"));
+ long taken = watch.taken();
+
+ Exception cause = result.getException();
+ assertNotNull(cause, "The exchange should fail with the commit
failure");
+ assertFalse(cause instanceof ExchangeTimedOutException,
+ "Should fail with the commit failure, not
ExchangeTimedOutException: " + cause);
+ assertTrue(hasMessageInChain(cause, COMMIT_FAILURE), "Should fail with
the commit failure: " + cause);
+ assertTrue(taken < 3000, "The exchange should fail promptly, not after
the request timeout, took " + taken + " ms");
+
+ // the reply handler has been cancelled, so the request timeout must
not complete the exchange a second time
+ await().during(4, TimeUnit.SECONDS)
+ .atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(() -> {
+ assertSame(cause, result.getException(), "Exception
changed after the commit failure");
+ assertEquals(1, FAILED.get(), "onFailure calls");
+ });
+ assertCompletedOnce(0, 1);
+ }
+
+ private static boolean hasMessageInChain(Throwable t, String message) {
+ for (Throwable c = t; c != null; c = c.getCause()) {
+ if (c.getMessage() != null && c.getMessage().contains(message)) {
+ return true;
+ }
+ }
+ return false;
+ }
+
+ private void assertCompletedOnce(int expectedCompleted, int
expectedFailed) {
+ assertEquals(expectedCompleted, COMPLETED.get(), "onComplete calls");
+ assertEquals(expectedFailed, FAILED.get(), "onFailure calls");
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(() -> assertEquals(0,
context.getInflightRepository().size(), "inflight exchanges"));
+ }
+
@Override
public String getComponentName() {
return "activemq";
@@ -97,7 +200,36 @@ public class JmsInOutSendFailureCallbackTest extends
AbstractJMSTest {
public void configure() {
from("direct:JmsInOutSendFailureCallbackTest")
.to(ExchangePattern.InOut,
-
"activemq:queue:JmsInOutSendFailureCallbackTest?requestTimeout=2000");
+ "activemq:queue:" + FAIL_QUEUE +
"?requestTimeout=2000");
+
+ from("direct:timeoutDuringSend")
+
.process(JmsInOutSendFailureCallbackTest::countCompletions)
+ .to(ExchangePattern.InOut,
+ "activemq:queue:" + TIMEOUT_QUEUE +
"?requestTimeout=100&requestTimeoutCheckerInterval=50");
+
+ from("direct:timeoutDuringSendHandled")
+
.process(JmsInOutSendFailureCallbackTest::countCompletions)
+ .doTry()
+ .to(ExchangePattern.InOut,
+ "activemq:queue:" + TIMEOUT_QUEUE
+ +
"?requestTimeout=100&requestTimeoutCheckerInterval=50")
+ .doCatch(ExchangeTimedOutException.class)
+ .setBody(constant("Timed out"))
+ .end();
+
+ from("direct:replyDuringSend")
+
.process(JmsInOutSendFailureCallbackTest::countCompletions)
+ .to(ExchangePattern.InOut, "activemq:queue:" +
REPLY_QUEUE + "?requestTimeout=10000");
+
+ from("direct:commitFailure")
+
.process(JmsInOutSendFailureCallbackTest::countCompletions)
+ .to(ExchangePattern.InOut,
+ "activemq:queue:" + COMMIT_QUEUE
+ +
"?useMessageIDAsCorrelationID=true&transactedInOut=true"
+ +
"&requestTimeout=3000&requestTimeoutCheckerInterval=50");
+
+ from("activemq:queue:" + REPLY_QUEUE)
+ .setBody(constant("Bye World"));
}
};
}
@@ -111,6 +243,25 @@ public class JmsInOutSendFailureCallbackTest extends
AbstractJMSTest {
void setUpRequirements() {
context = camelContextExtension.getContext();
template = camelContextExtension.getProducerTemplate();
+ COMPLETED.set(0);
+ FAILED.set(0);
+ exchangeDone = new CountDownLatch(1);
+ }
+
+ private static void countCompletions(Exchange exchange) {
+ exchange.getExchangeExtension().addOnCompletion(new
SynchronizationAdapter() {
+ @Override
+ public void onComplete(Exchange exchange) {
+ COMPLETED.incrementAndGet();
+ exchangeDone.countDown();
+ }
+
+ @Override
+ public void onFailure(Exchange exchange) {
+ FAILED.incrementAndGet();
+ exchangeDone.countDown();
+ }
+ });
}
private static ConnectionFactory
createFailingSendConnectionFactory(ConnectionFactory delegate) {
@@ -128,21 +279,62 @@ public class JmsInOutSendFailureCallbackTest extends
AbstractJMSTest {
}
private static Session wrapSession(Session delegate) {
+ // whether this session has sent to COMMIT_QUEUE, and its commit
should then fail
+ AtomicBoolean failCommit = new AtomicBoolean();
return proxyOf(Session.class, delegate, (proxy, method, args) -> {
- Object result = method.invoke(delegate, args);
- return result instanceof MessageProducer producer ?
wrapProducer(producer) : result;
+ if ("commit".equals(method.getName()) && failCommit.get()) {
+ throw new JMSException(COMMIT_FAILURE);
+ }
+ Object result = invoke(method, delegate, args);
+ if (result instanceof MessageProducer producer) {
+ if (COMMIT_QUEUE.equals(queueName(producer.getDestination())))
{
+ failCommit.set(true);
+ }
+ return wrapProducer(producer);
+ }
+ return result;
});
}
private static MessageProducer wrapProducer(MessageProducer delegate) {
return proxyOf(MessageProducer.class, delegate, (proxy, method, args)
-> {
if ("send".equals(method.getName())) {
- throw new JMSException("Simulated send failure: broker
rejected message");
+ String queue = queueName(delegate.getDestination());
+ if (FAIL_QUEUE.equals(queue)) {
+ throw new JMSException("Simulated send failure: broker
rejected message");
+ } else if (TIMEOUT_QUEUE.equals(queue)) {
+ // a send that blocks longer than the request timeout, and
then fails
+ awaitExchangeDone();
+ throw new JMSException("Simulated send failure after the
request timeout");
+ } else if (REPLY_QUEUE.equals(queue)) {
+ // the message is sent and answered, but the send reports
a failure afterwards
+ invoke(method, delegate, args);
+ awaitExchangeDone();
+ throw new JMSException("Simulated send failure after the
message was sent");
+ }
}
- return method.invoke(delegate, args);
+ return invoke(method, delegate, args);
});
}
+ private static void awaitExchangeDone() throws InterruptedException {
+ if (!exchangeDone.await(20, TimeUnit.SECONDS)) {
+ throw new IllegalStateException("The exchange was not completed
while the send was in progress");
+ }
+ }
+
+ private static String queueName(Destination destination) throws
JMSException {
+ return destination instanceof Queue queue ? queue.getQueueName() :
null;
+ }
+
+ private static Object invoke(Method method, Object delegate, Object[]
args) throws Throwable {
+ try {
+ return method.invoke(delegate, args);
+ } catch (InvocationTargetException e) {
+ throw e.getCause();
+ }
+ }
+
@SuppressWarnings("unchecked")
private static <T> T proxyOf(Class<T> iface, T delegate, InvocationHandler
handler) {
return (T) Proxy.newProxyInstance(iface.getClassLoader(), new
Class<?>[] { iface }, handler);
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 0f7c35cf6615..5e046ad7e58c 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
@@ -4029,6 +4029,17 @@ were sent. As a result:
validator does not recognise a header parameter as an array, so a repeated
header is rejected
like a single-value one.
+=== camel-jms - request/reply completes the exchange once when a send fails
late
+
+An InOut exchange whose JMS send fails after the request timeout, or the
reply, has already completed the exchange
+(for example a send that blocks longer than `requestTimeout` and then fails)
is no longer completed a second time with
+the send exception. The send failure is logged at WARN level instead, and the
exchange keeps the outcome of the timeout
+or the reply.
+
+The `cancelCorrelationId` method of
`org.apache.camel.component.jms.reply.ReplyManager` (added in 4.22) now returns
+a `boolean`: `true` if the pending reply was cancelled, `false` if the request
timeout or the reply had already taken
+it. A custom `ReplyManager` implementation has to change the return type.
+
== ThrottlingExceptionRoutePolicy
`ThrottlingExceptionRoutePolicy.setKeepOpen(true)` now opens the circuit
immediately and synchronously (the consumer