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 21a37f4e33c9 CAMEL-25122: camel-sjms - do not complete an InOut
exchange twice when the send fails (#27031)
21a37f4e33c9 is described below
commit 21a37f4e33c932964cccb8645ae82ec30e67c797
Author: allthingssecurity <[email protected]>
AuthorDate: Wed Sep 30 22:25:14 2026 +0530
CAMEL-25122: camel-sjms - do not complete an InOut exchange twice when the
send fails (#27031)
SjmsProducer registered the reply handler before sending, and when the send
failed the handler stayed registered, so the request timeout completed the same
exchange a second time, replacing the exception with an
ExchangeTimedOutException. A send that blocked longer than requestTimeout could
race the same way.
The sjms reply manager can now cancel a pending reply, and the producer
ensures the exchange is completed only once, mirroring the camel-jms fixes in
CAMEL-24073 and CAMEL-25095.
Closes #27031
Co-authored-by: Claude Opus 5.5 <[email protected]>
---
.../apache/camel/component/sjms/SjmsProducer.java | 17 +-
.../camel/component/sjms/reply/ReplyManager.java | 19 ++
.../component/sjms/reply/ReplyManagerSupport.java | 9 +
.../producer/InOutSendFailureCallbackTest.java | 250 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 12 +
5 files changed, 306 insertions(+), 1 deletion(-)
diff --git
a/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsProducer.java
b/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsProducer.java
index 8e24ea906c20..21d328757a24 100644
---
a/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsProducer.java
+++
b/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/SjmsProducer.java
@@ -292,6 +292,9 @@ public class SjmsProducer extends DefaultAsyncProducer {
in.setHeader(SjmsConstants.JMS_CORRELATION_ID,
GENERATED_CORRELATION_ID_PREFIX + getUuidGenerator().generateUuid());
}
+ // the correlation id the reply handler is registered under, which is
cancelled if the send fails
+ 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);
@@ -310,7 +313,8 @@ public class SjmsProducer extends DefaultAsyncProducer {
JmsMessageHelper.setJMSReplyTo(answer, replyTo);
String correlationId = determineCorrelationId(answer);
- replyManager.registerReply(replyManager, exchange, callback,
originalCorrelationId, correlationId, timeout);
+ registeredCorrelationId[0] =
replyManager.registerReply(replyManager, exchange, callback,
+ originalCorrelationId, correlationId, timeout);
if (LOG.isDebugEnabled()) {
LOG.debug("Using {}: {}, JMSReplyTo destination: {}, with
request timeout: {} ms.",
@@ -325,6 +329,17 @@ public class SjmsProducer extends DefaultAsyncProducer {
try {
doSend(exchange, true, destinationName, messageCreator);
} catch (Exception e) {
+ // the send failed after the reply was registered: cancel it, as
otherwise the request timeout completes
+ // the exchange a second time
+ 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;
+ }
exchange.setException(e);
callback.done(true);
return true;
diff --git
a/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/reply/ReplyManager.java
b/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/reply/ReplyManager.java
index 4c7b66591032..058fbb902010 100644
---
a/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/reply/ReplyManager.java
+++
b/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/reply/ReplyManager.java
@@ -88,6 +88,25 @@ public interface ReplyManager extends SessionMessageListener
{
*/
void updateCorrelationId(String correlationId, String newCorrelationId,
long requestTimeout);
+ /**
+ * Cancels a pending reply correlation, so the request timeout does not
complete the exchange.
+ * <p/>
+ * This is used when the JMS send fails after the reply has been
registered. 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.
+ * <p/>
+ * The default implementation does not cancel anything and returns
<tt>true</tt>, which keeps the previous behaviour
+ * for a custom implementation. {@link ReplyManagerSupport} overrides it.
+ *
+ * @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)
+ */
+ default boolean cancelCorrelationId(String correlationId) {
+ return true;
+ }
+
/**
* Process the reply
*
diff --git
a/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/reply/ReplyManagerSupport.java
b/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/reply/ReplyManagerSupport.java
index e2800ee1fcde..753bcfdc92be 100644
---
a/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/reply/ReplyManagerSupport.java
+++
b/components/camel-sjms/src/main/java/org/apache/camel/component/sjms/reply/ReplyManagerSupport.java
@@ -127,6 +127,15 @@ public abstract class ReplyManagerSupport extends
ServiceSupport implements Repl
return correlationId;
}
+ @Override
+ public boolean cancelCorrelationId(String correlationId) {
+ if (correlationId != null && correlation != null &&
correlation.remove(correlationId) != null) {
+ log.debug("Cancelled reply correlation [{}]", correlationId);
+ return true;
+ }
+ return false;
+ }
+
@Override
public void onMessage(Message message, Session session) throws
JMSException {
String correlationID = getJMSCorrelationID(message);
diff --git
a/components/camel-sjms/src/test/java/org/apache/camel/component/sjms/producer/InOutSendFailureCallbackTest.java
b/components/camel-sjms/src/test/java/org/apache/camel/component/sjms/producer/InOutSendFailureCallbackTest.java
new file mode 100644
index 000000000000..32448fa5e012
--- /dev/null
+++
b/components/camel-sjms/src/test/java/org/apache/camel/component/sjms/producer/InOutSendFailureCallbackTest.java
@@ -0,0 +1,250 @@
+/*
+ * 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.sjms.producer;
+
+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.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;
+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.sjms.SjmsComponent;
+import org.apache.camel.component.sjms.support.JmsTestSupport;
+import org.apache.camel.support.SynchronizationAdapter;
+import org.junit.jupiter.api.BeforeEach;
+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.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;
+
+/**
+ * When the send of an InOut message fails after the reply has been
registered, the exchange must be completed exactly
+ * once: not a second time by the request timeout, 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.
+ */
+public class InOutSendFailureCallbackTest extends JmsTestSupport {
+
+ private static final String FAIL_QUEUE =
"InOutSendFailureCallbackTest.fail";
+ private static final String TIMEOUT_QUEUE =
"InOutSendFailureCallbackTest.timeout";
+ private static final String REPLY_QUEUE =
"InOutSendFailureCallbackTest.reply";
+
+ private static final AtomicInteger COMPLETED = new AtomicInteger();
+ private static final AtomicInteger FAILED = new AtomicInteger();
+ // the send to TIMEOUT_QUEUE and REPLY_QUEUE waits until the exchange has
been completed
+ private static volatile CountDownLatch exchangeDone = new
CountDownLatch(1);
+
+ public InOutSendFailureCallbackTest() {
+ addSjmsComponent = false;
+ }
+
+ @BeforeEach
+ void resetCounters() {
+ COMPLETED.set(0);
+ FAILED.set(0);
+ exchangeDone = new CountDownLatch(1);
+ }
+
+ @Test
+ public void testSendFailure() {
+ Exchange result = template.send("direct:fail", ExchangePattern.InOut,
e -> e.getIn().setBody("Hello"));
+
+ Exception cause = result.getException();
+ assertNotNull(cause);
+ assertTrue(hasMessageInChain(cause, "Simulated send failure"), "Should
fail with the send failure: " + cause);
+
+ // the request timeout (500 ms) must not complete the exchange a
second time
+ await().during(1500, TimeUnit.MILLISECONDS).atMost(3,
TimeUnit.SECONDS).untilAsserted(() -> {
+ assertSame(cause, result.getException(), "The exception changed
after the send failure");
+ assertEquals(1, FAILED.get(), "onFailure calls");
+ });
+ assertCompletedOnce(0, 1);
+ }
+
+ @Test
+ public void testTimeoutDuringFailingSend() {
+ // the send blocks until the request timeout has completed the
exchange, and then fails
+ Exchange result = template.send("direct:timeout",
ExchangePattern.InOut, e -> e.getIn().setBody("Hello"));
+
+ assertInstanceOf(ExchangeTimedOutException.class,
result.getException(),
+ "Should fail with the timeout that completed the exchange, not
with the later send failure");
+ assertCompletedOnce(0, 1);
+ }
+
+ @Test
+ public void testReplyDuringFailingSend() {
+ // the request reaches the broker and is answered, and then the send
fails
+ Exchange result = template.send("direct:reply", ExchangePattern.InOut,
e -> e.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);
+ }
+
+ 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) {
+ await().atMost(5, TimeUnit.SECONDS)
+ .untilAsserted(() -> assertEquals(0,
context.getInflightRepository().size(), "inflight exchanges"));
+ assertEquals(expectedCompleted, COMPLETED.get(), "onComplete calls");
+ assertEquals(expectedFailed, FAILED.get(), "onFailure calls");
+ }
+
+ @Override
+ protected CamelContext createCamelContext() throws Exception {
+ CamelContext camelContext = super.createCamelContext();
+ SjmsComponent component = new SjmsComponent();
+
component.setConnectionFactory(createFailingSendConnectionFactory(connectionFactory));
+ component.setRequestTimeoutCheckerInterval(50);
+ camelContext.addComponent("sjms", component);
+ return camelContext;
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:fail")
+
.process(InOutSendFailureCallbackTest::countCompletions)
+ .to(ExchangePattern.InOut, "sjms:queue:" + FAIL_QUEUE
+ "?requestTimeout=500");
+
+ from("direct:timeout")
+
.process(InOutSendFailureCallbackTest::countCompletions)
+ .to(ExchangePattern.InOut, "sjms:queue:" +
TIMEOUT_QUEUE + "?requestTimeout=100");
+
+ from("direct:reply")
+
.process(InOutSendFailureCallbackTest::countCompletions)
+ .to(ExchangePattern.InOut, "sjms:queue:" + REPLY_QUEUE
+ "?requestTimeout=10000");
+
+ from("sjms:queue:" + REPLY_QUEUE)
+ .setBody(constant("Bye World"));
+ }
+ };
+ }
+
+ 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) {
+ return proxyOf(ConnectionFactory.class, (proxy, method, args) -> {
+ Object result = invoke(method, delegate, args);
+ return result instanceof Connection connection ?
wrapConnection(connection) : result;
+ });
+ }
+
+ private static Connection wrapConnection(Connection delegate) {
+ return proxyOf(Connection.class, (proxy, method, args) -> {
+ Object result = invoke(method, delegate, args);
+ return result instanceof Session session ? wrapSession(session) :
result;
+ });
+ }
+
+ private static Session wrapSession(Session delegate) {
+ return proxyOf(Session.class, (proxy, method, args) -> {
+ Object result = invoke(method, delegate, args);
+ return result instanceof MessageProducer producer ?
wrapProducer(producer) : result;
+ });
+ }
+
+ private static MessageProducer wrapProducer(MessageProducer delegate) {
+ return proxyOf(MessageProducer.class, (proxy, method, args) -> {
+ if ("send".equals(method.getName())) {
+ String queue = queueName(delegate.getDestination(), args);
+ 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 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 producerDestination, Object[]
args) throws JMSException {
+ Destination destination = producerDestination;
+ if (destination == null && args != null && args.length > 0 && args[0]
instanceof Destination d) {
+ destination = d;
+ }
+ 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, 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 612cc228ccd9..fa1e19bcdef1 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
@@ -1209,6 +1209,18 @@ The container images were removed from Docker Hub on
September 2026.
MinIO is S3-compatible, so existing deployments can migrate to the
`camel-aws2-s3` component by pointing it at the MinIO server, for example:
`aws2-s3://mybucket?overrideEndpoint=true&uriEndpointOverride=http://localhost:9000&forcePathStyle=true&accessKey=...&secretKey=...`
+=== camel-sjms - request/reply completes the exchange once when the send fails
+
+When the send of an InOut message fails, the pending reply is now cancelled,
so the request timeout no longer
+completes the exchange a second time with an `ExchangeTimedOutException` after
it has already failed with the send
+exception. When the 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), the
exchange keeps the outcome of the
+timeout or the reply, and the send failure is logged at WARN level.
+
+`org.apache.camel.component.sjms.reply.ReplyManager` has a new default method
`boolean cancelCorrelationId(String)`,
+which `ReplyManagerSupport` implements. A custom `ReplyManager` implementation
that does not override it keeps the
+previous behaviour.
+
=== camel-tika
The Tika dependency has been upgraded from 3.x to 4.x. Tika 4 removed the
`TikaConfig` class and XML