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

Reply via email to