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 446b412daea3 CAMEL-25128: camel-netty - a failed write and the 
correlation timeout should not both complete the exchange (#27037)
446b412daea3 is described below

commit 446b412daea3703292f8c57cbd6af143ec13fd39
Author: allthingssecurity <[email protected]>
AuthorDate: Tue Sep 29 17:26:48 2026 +0530

    CAMEL-25128: camel-netty - a failed write and the correlation timeout 
should not both complete the exchange (#27037)
    
    * CAMEL-25128: camel-netty - a failed write and the correlation timeout 
should not both complete the exchange
    * CAMEL-25128: camel-netty - log the correlation timeout only when it 
completes the exchange
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../camel/component/netty/NettyCamelState.java     |  30 +++-
 .../camel/component/netty/NettyProducer.java       |   6 +-
 .../netty/TimeoutCorrelationManagerSupport.java    |   6 +-
 ...yTimeoutCorrelationManagerWriteFailureTest.java | 185 +++++++++++++++++++++
 4 files changed, 215 insertions(+), 12 deletions(-)

diff --git 
a/components/camel-netty/src/main/java/org/apache/camel/component/netty/NettyCamelState.java
 
b/components/camel-netty/src/main/java/org/apache/camel/component/netty/NettyCamelState.java
index 517fe4c3cc46..2d26a61a7f93 100644
--- 
a/components/camel-netty/src/main/java/org/apache/camel/component/netty/NettyCamelState.java
+++ 
b/components/camel-netty/src/main/java/org/apache/camel/component/netty/NettyCamelState.java
@@ -53,12 +53,22 @@ public final class NettyCamelState {
     }
 
     public void callbackDoneOnce(boolean doneSync) {
-        if (!callbackCalled.getAndSet(true)) {
+        if (markDone()) {
             // this is the first time we call the callback
             callback.done(doneSync);
         }
     }
 
+    /**
+     * Claims the completion of the exchange.
+     *
+     * @return <tt>true</tt> for the first caller only, which must then set 
the outcome on the exchange and call the
+     *         callback, <tt>false</tt> if the exchange is already completed 
and must not be touched anymore
+     */
+    public boolean markDone() {
+        return callbackCalled.compareAndSet(false, true);
+    }
+
     public Exchange getExchange() {
         return exchange;
     }
@@ -68,14 +78,24 @@ public final class NettyCamelState {
     }
 
     public void onExceptionCaughtOnce(boolean doneSync) {
+        onExceptionCaughtOnce(doneSync, null);
+    }
+
+    /**
+     * Completes the exchange with the cause of a failed write, unless an 
exception has already been caught or the
+     * exchange has already been completed, for example by the timeout of a 
correlation manager.
+     */
+    public void onExceptionCaughtOnce(boolean doneSync, Throwable cause) {
         // only trigger callback once if an exception has not already been 
caught
         // (ClientChannelHandler#exceptionCaught vs 
NettyProducer#processWithConnectedChannel)
-        if (exceptionCaught.compareAndSet(false, true)) {
-            // set some general exception as Camel should know the netty write 
operation failed
-            if (exchange.getException() == null) {
+        if (exceptionCaught.compareAndSet(false, true) && markDone()) {
+            if (cause != null) {
+                exchange.setException(cause);
+            } else if (exchange.getException() == null) {
+                // set some general exception as Camel should know the netty 
write operation failed
                 exchange.setException(new IOException("Netty write operation 
failed"));
             }
-            callbackDoneOnce(doneSync);
+            callback.done(doneSync);
         }
     }
 }
diff --git 
a/components/camel-netty/src/main/java/org/apache/camel/component/netty/NettyProducer.java
 
b/components/camel-netty/src/main/java/org/apache/camel/component/netty/NettyProducer.java
index e3fc56afe4e7..af53bd2a532f 100644
--- 
a/components/camel-netty/src/main/java/org/apache/camel/component/netty/NettyProducer.java
+++ 
b/components/camel-netty/src/main/java/org/apache/camel/component/netty/NettyProducer.java
@@ -382,10 +382,8 @@ public class NettyProducer extends DefaultAsyncProducer {
                     } catch (Exception e) {
                         cause = e.getCause();
                     }
-                    if (cause != null) {
-                        exchange.setException(cause);
-                    }
-                    state.onExceptionCaughtOnce(false);
+                    // the exchange may already be completed (by the timeout 
of the correlation manager)
+                    state.onExceptionCaughtOnce(false, cause);
                     return;
                 }
 
diff --git 
a/components/camel-netty/src/main/java/org/apache/camel/component/netty/TimeoutCorrelationManagerSupport.java
 
b/components/camel-netty/src/main/java/org/apache/camel/component/netty/TimeoutCorrelationManagerSupport.java
index e33e47e1b9d7..5f98dd614265 100644
--- 
a/components/camel-netty/src/main/java/org/apache/camel/component/netty/TimeoutCorrelationManagerSupport.java
+++ 
b/components/camel-netty/src/main/java/org/apache/camel/component/netty/TimeoutCorrelationManagerSupport.java
@@ -225,12 +225,12 @@ public abstract class TimeoutCorrelationManagerSupport 
extends ServiceSupport
             return;
         }
 
-        timeoutLogger.log("Timeout of correlation id: " + key);
-
         workerPool.submit(() -> {
             Exchange exchange = value.getExchange();
             AsyncCallback callback = value.getCallback();
-            if (exchange != null && callback != null) {
+            // the exchange may have been completed meanwhile (the write 
failed), then it must not be touched
+            if (exchange != null && callback != null && value.markDone()) {
+                timeoutLogger.log("Timeout of correlation id: " + key);
                 Object timeoutBody = getTimeoutResponse(key, 
exchange.getMessage().getBody());
                 if (timeoutBody != null) {
                     exchange.getMessage().setBody(timeoutBody);
diff --git 
a/components/camel-netty/src/test/java/org/apache/camel/component/netty/NettyTimeoutCorrelationManagerWriteFailureTest.java
 
b/components/camel-netty/src/test/java/org/apache/camel/component/netty/NettyTimeoutCorrelationManagerWriteFailureTest.java
new file mode 100644
index 000000000000..e6a74377e0f6
--- /dev/null
+++ 
b/components/camel-netty/src/test/java/org/apache/camel/component/netty/NettyTimeoutCorrelationManagerWriteFailureTest.java
@@ -0,0 +1,185 @@
+/*
+ * 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.netty;
+
+import java.io.IOException;
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import io.netty.channel.ChannelHandler;
+import io.netty.channel.ChannelHandlerContext;
+import io.netty.channel.ChannelOutboundHandlerAdapter;
+import io.netty.channel.ChannelPromise;
+import io.netty.util.ReferenceCountUtil;
+import org.apache.camel.AsyncProducer;
+import org.apache.camel.BindToRegistry;
+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.support.DefaultExchange;
+import org.apache.camel.support.service.ServiceHelper;
+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.assertTrue;
+
+/**
+ * With a correlation manager extending {@link 
TimeoutCorrelationManagerSupport}, a failed write and the timeout of the
+ * request must complete the exchange only once, whichever comes first.
+ */
+public class NettyTimeoutCorrelationManagerWriteFailureTest extends 
BaseNettyTest {
+
+    // counts down when the worker pool of the correlation manager has 
processed the timeout
+    private final CountDownLatch timeoutProcessed = new CountDownLatch(1);
+    private final ThreadPoolExecutor workerPool = new ThreadPoolExecutor(
+            1, 1, 0, TimeUnit.MILLISECONDS, new LinkedBlockingQueue<>()) {
+        @Override
+        protected void afterExecute(Runnable r, Throwable t) {
+            timeoutProcessed.countDown();
+        }
+    };
+
+    @BindToRegistry("myManager")
+    private final MyCorrelationManager myManager = new MyCorrelationManager();
+
+    @BindToRegistry("failingWrite")
+    private final FailingWriteHandler failingWrite = new FailingWriteHandler();
+
+    private final List<Throwable> completions = new CopyOnWriteArrayList<>();
+    private AsyncProducer producer;
+
+    @AfterEach
+    public void stopProducer() {
+        ServiceHelper.stopService(producer);
+        workerPool.shutdownNow();
+    }
+
+    @Test
+    public void testWriteFailsBeforeTimeout() throws Exception {
+        failingWrite.failImmediately = true;
+
+        Exchange exchange = sendRequest();
+
+        // the correlation timeout fires later for the same request
+        assertTrue(timeoutProcessed.await(5, TimeUnit.SECONDS), "the timeout 
should have been processed");
+        assertEquals(1, completions.size(), "the exchange should be completed 
once: " + completions);
+        assertInstanceOf(IOException.class, completions.get(0));
+        // and the exchange is not changed after it completed
+        assertInstanceOf(IOException.class, exchange.getException());
+    }
+
+    @Test
+    public void testWriteFailsAfterTimeout() throws Exception {
+        failingWrite.failImmediately = false;
+
+        Exchange exchange = sendRequest();
+        assertTrue(timeoutProcessed.await(5, TimeUnit.SECONDS), "the timeout 
should have been processed");
+        assertInstanceOf(ExchangeTimedOutException.class, completions.get(0));
+
+        // now the write that is still pending fails, for example as the 
connection is reset
+        failingWrite.failPendingWrite();
+
+        assertEquals(1, completions.size(), "the exchange should be completed 
once: " + completions);
+        assertInstanceOf(ExchangeTimedOutException.class, 
exchange.getException());
+    }
+
+    private Exchange sendRequest() throws Exception {
+        producer = context.getEndpoint(
+                
"netty:tcp://localhost:{{port}}?sync=true&producerPoolEnabled=false&encoders=#failingWrite"
+                                       + "&correlationManager=#myManager")
+                .createAsyncProducer();
+        ServiceHelper.startService(producer);
+
+        Exchange exchange = new DefaultExchange(context, 
ExchangePattern.InOut);
+        exchange.getMessage().setBody("abc-hello");
+        CountDownLatch done = new CountDownLatch(1);
+        producer.process(exchange, doneSync -> {
+            completions.add(exchange.getException());
+            done.countDown();
+        });
+        assertTrue(done.await(5, TimeUnit.SECONDS), "the exchange should be 
completed");
+        return exchange;
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
from("netty:tcp://localhost:{{port}}?textline=true&sync=true").transform(body().prepend("Bye
 "));
+            }
+        };
+    }
+
+    private final class MyCorrelationManager extends 
TimeoutCorrelationManagerSupport {
+
+        MyCorrelationManager() {
+            setTimeout(200);
+            setTimeoutChecker(50);
+            setWorkerPool(workerPool);
+        }
+
+        @Override
+        public String getRequestCorrelationId(Object request) {
+            return request.toString().substring(0, 3);
+        }
+
+        @Override
+        public String getResponseCorrelationId(Object response) {
+            return response.toString().substring(0, 3);
+        }
+    }
+
+    /**
+     * Fails the write at once, or keeps it pending until {@link 
#failPendingWrite()} is called.
+     */
+    @ChannelHandler.Sharable
+    public static final class FailingWriteHandler extends 
ChannelOutboundHandlerAdapter {
+
+        private final CountDownLatch written = new CountDownLatch(1);
+        private volatile boolean failImmediately;
+        private volatile ChannelPromise pending;
+
+        @Override
+        public void write(ChannelHandlerContext ctx, Object msg, 
ChannelPromise promise) {
+            ReferenceCountUtil.release(msg);
+            if (failImmediately) {
+                promise.setFailure(new IOException("Simulated write failure"));
+            } else {
+                pending = promise;
+                written.countDown();
+            }
+        }
+
+        void failPendingWrite() throws InterruptedException {
+            assertTrue(written.await(5, TimeUnit.SECONDS), "the request should 
have been written");
+            // listeners are notified in the order they were added, so when 
this one runs the write listener of the
+            // producer has run as well
+            CountDownLatch notified = new CountDownLatch(1);
+            pending.addListener(f -> notified.countDown());
+            pending.setFailure(new IOException("Simulated connection reset"));
+            assertTrue(notified.await(5, TimeUnit.SECONDS), "the write 
listeners should have been notified");
+        }
+    }
+}

Reply via email to