This is an automated email from the ASF dual-hosted git repository.

rmaucher pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/tomcat.git

commit d71071d7f1eea7de42c28fa41fe9747468f8c48d
Author: opencode <[email protected]>
AuthorDate: Fri Oct 9 11:18:44 2026 +0200

    Fix use-after-return of pooled send buffer in async sends
    
    GroupChannel.send() obtains the message buffer from the BufferPool and
    returns it to the pool in a finally block once sendMessage() returns.
    MessageDispatchInterceptor queues async messages for deferred sending
    and only detached the message from that buffer via a deep clone when
    useDeepClone was true, which is the default. With useDeepClone=false,
    the original ChannelData wrapping the pooled buffer was queued as-is,
    so a concurrent send could clear and refill the buffer before the
    dispatch thread serialized the message, silently corrupting cluster
    traffic and the queue size accounting.
    
    When useDeepClone is false, queue a shallow ChannelData clone instead,
    which copies the message bytes into a new, non pooled buffer while
    still avoiding the full re-serialization of the deep clone. Add
    TestMessageDispatchInterceptor with a deterministic reproduction and a
    check that synchronous sends pass the message through unchanged.
---
 .../interceptors/MessageDispatchInterceptor.java   |  13 +-
 .../TestMessageDispatchInterceptor.java            | 141 +++++++++++++++++++++
 webapps/docs/changelog.xml                         |   8 ++
 3 files changed, 160 insertions(+), 2 deletions(-)

diff --git 
a/java/org/apache/catalina/tribes/group/interceptors/MessageDispatchInterceptor.java
 
b/java/org/apache/catalina/tribes/group/interceptors/MessageDispatchInterceptor.java
index f3af235295..ae08a76dd4 100644
--- 
a/java/org/apache/catalina/tribes/group/interceptors/MessageDispatchInterceptor.java
+++ 
b/java/org/apache/catalina/tribes/group/interceptors/MessageDispatchInterceptor.java
@@ -30,6 +30,7 @@ import org.apache.catalina.tribes.Member;
 import org.apache.catalina.tribes.UniqueId;
 import org.apache.catalina.tribes.group.ChannelInterceptorBase;
 import org.apache.catalina.tribes.group.InterceptorPayload;
+import org.apache.catalina.tribes.io.ChannelData;
 import org.apache.catalina.tribes.util.ExecutorFactory;
 import org.apache.catalina.tribes.util.StringManager;
 import org.apache.catalina.tribes.util.TcclThreadFactory;
@@ -58,7 +59,8 @@ public class MessageDispatchInterceptor extends 
ChannelInterceptorBase implement
      */
     protected volatile boolean run = false;
     /**
-     * Whether to use deep clone.
+     * Whether to use deep clone. When false, the message is still detached 
from the caller's buffer by a shallow
+     * clone before it is queued, but the unique id and address are shared 
rather than re-serialized.
      */
     protected boolean useDeepClone = true;
     /**
@@ -113,6 +115,12 @@ public class MessageDispatchInterceptor extends 
ChannelInterceptorBase implement
             // add to queue
             if (useDeepClone) {
                 msg = (ChannelMessage) msg.deepclone();
+            } else if (msg instanceof ChannelData channelData) {
+                // The caller of sendMessage() may reuse the buffer that backs 
the
+                // message as soon as this method returns, e.g. 
GroupChannel.send()
+                // returns it to the BufferPool. Clone, which copies the 
message
+                // data, so the queued message is independent of that buffer.
+                msg = channelData.clone();
             }
             if (!addToQueue(msg, destination, payload)) {
                 throw new 
ChannelException(sm.getString("messageDispatchInterceptor.unableAdd.queue"));
@@ -187,7 +195,8 @@ public class MessageDispatchInterceptor extends 
ChannelInterceptorBase implement
 
 
     /**
-     * Set whether to use deep clone.
+     * Set whether to use deep clone. When false, the queued message is a 
shallow clone of the original with the
+     * message bytes copied to a new buffer rather than a fully re-serialized 
deep clone.
      * @param useDeepClone whether to use deep clone
      */
     public void setUseDeepClone(boolean useDeepClone) {
diff --git 
a/test/org/apache/catalina/tribes/group/interceptors/TestMessageDispatchInterceptor.java
 
b/test/org/apache/catalina/tribes/group/interceptors/TestMessageDispatchInterceptor.java
new file mode 100644
index 0000000000..b3b5efd2a4
--- /dev/null
+++ 
b/test/org/apache/catalina/tribes/group/interceptors/TestMessageDispatchInterceptor.java
@@ -0,0 +1,141 @@
+/*
+ * 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.catalina.tribes.group.interceptors;
+
+import java.nio.charset.StandardCharsets;
+import java.util.Arrays;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import org.apache.catalina.tribes.Channel;
+import org.apache.catalina.tribes.ChannelException;
+import org.apache.catalina.tribes.ChannelMessage;
+import org.apache.catalina.tribes.Member;
+import org.apache.catalina.tribes.group.ChannelInterceptorBase;
+import org.apache.catalina.tribes.group.GroupChannel;
+import org.apache.catalina.tribes.group.InterceptorPayload;
+import org.apache.catalina.tribes.io.ChannelData;
+import org.apache.catalina.tribes.io.XByteBuffer;
+import org.apache.catalina.tribes.membership.MemberImpl;
+
+public class TestMessageDispatchInterceptor {
+
+    /*
+     * Captures the message on the dispatch thread, but only after the caller 
thread has simulated the reuse of the
+     * original (pooled) buffer that GroupChannel.send() returns to the 
BufferPool.
+     */
+    private static class GatedCollector extends ChannelInterceptorBase {
+
+        private final CountDownLatch dispatchStarted = new CountDownLatch(1);
+        private final CountDownLatch bufferReused = new CountDownLatch(1);
+        private final CountDownLatch collected = new CountDownLatch(1);
+
+        private volatile XByteBuffer capturedBuffer;
+        private volatile byte[] capturedData;
+
+        @Override
+        public void sendMessage(Member[] destination, ChannelMessage msg, 
InterceptorPayload payload)
+                throws ChannelException {
+            dispatchStarted.countDown();
+            try {
+                if (!bufferReused.await(10, TimeUnit.SECONDS)) {
+                    throw new ChannelException("Timed out waiting for buffer 
reuse");
+                }
+            } catch (InterruptedException e) {
+                Thread.currentThread().interrupt();
+                throw new ChannelException(e);
+            }
+            capturedBuffer = msg.getMessage();
+            int len = capturedBuffer.getLength();
+            capturedData = Arrays.copyOfRange(capturedBuffer.getBytesDirect(), 
0, len);
+            collected.countDown();
+        }
+    }
+
+    @Test
+    public void testAsyncSendWithoutDeepCloneDetachesPooledBuffer() throws 
Exception {
+        byte[] original = "ORIGINAL".getBytes(StandardCharsets.UTF_8);
+        byte[] reused = "OVERWRITTEN-BY-POOL".getBytes(StandardCharsets.UTF_8);
+
+        GroupChannel channel = new GroupChannel();
+        MessageDispatchInterceptor interceptor = new 
MessageDispatchInterceptor();
+        interceptor.setUseDeepClone(false);
+        GatedCollector collector = new GatedCollector();
+        interceptor.setChannel(channel);
+        interceptor.setNext(collector);
+        interceptor.startQueue();
+        try {
+            ChannelData data = new ChannelData(false);
+            data.setOptions(Channel.SEND_OPTIONS_ASYNCHRONOUS);
+            XByteBuffer buffer = new XByteBuffer(original.length + 128, false);
+            buffer.append(original, 0, original.length);
+            data.setMessage(buffer);
+            Member[] destination = new Member[] { new MemberImpl("localhost", 
4000, -1) };
+
+            interceptor.sendMessage(destination, data, null);
+            Assert.assertTrue("Dispatch thread did not start",
+                    collector.dispatchStarted.await(10, TimeUnit.SECONDS));
+            // Simulate GroupChannel.send() returning the buffer to the pool 
for reuse
+            buffer.clear();
+            buffer.append(reused, 0, reused.length);
+            collector.bufferReused.countDown();
+            Assert.assertTrue("Dispatch thread did not collect the message",
+                    collector.collected.await(10, TimeUnit.SECONDS));
+
+            Assert.assertNotSame("Queued message must not share the pooled 
buffer",
+                    buffer, collector.capturedBuffer);
+            Assert.assertArrayEquals("Queued message content was corrupted by 
buffer reuse",
+                    original, collector.capturedData);
+        } finally {
+            interceptor.stopQueue();
+        }
+    }
+
+    @Test
+    public void testSyncSendPassesMessageThrough() throws Exception {
+        GroupChannel channel = new GroupChannel();
+        MessageDispatchInterceptor interceptor = new 
MessageDispatchInterceptor();
+        RecordingCollector collector = new RecordingCollector();
+        interceptor.setChannel(channel);
+        interceptor.setNext(collector);
+
+        ChannelData data = new ChannelData(false);
+        XByteBuffer buffer = new XByteBuffer(64, false);
+        byte[] payload = "SYNC".getBytes(StandardCharsets.UTF_8);
+        buffer.append(payload, 0, payload.length);
+        data.setMessage(buffer);
+        Member[] destination = new Member[] { new MemberImpl("localhost", 
4000, -1) };
+
+        interceptor.sendMessage(destination, data, null);
+
+        Assert.assertSame("Synchronous sends must pass the original message 
through", data, collector.message);
+    }
+
+    private static class RecordingCollector extends ChannelInterceptorBase {
+
+        private volatile ChannelMessage message;
+
+        @Override
+        public void sendMessage(Member[] destination, ChannelMessage msg, 
InterceptorPayload payload)
+                throws ChannelException {
+            message = msg;
+        }
+    }
+}
diff --git a/webapps/docs/changelog.xml b/webapps/docs/changelog.xml
index e9c553a1c7..52560e692b 100644
--- a/webapps/docs/changelog.xml
+++ b/webapps/docs/changelog.xml
@@ -486,6 +486,14 @@
   </subsection>
   <subsection name="Tribes">
     <changelog>
+      <fix>
+        When the <code>MessageDispatchInterceptor</code> is configured with
+        <code>useDeepClone="false"</code>, asynchronously sent messages are now
+        detached from the caller's buffer with a shallow clone before being
+        queued. Previously, the buffer could be returned to the send buffer 
pool
+        and reused by a concurrent send while the queued message still
+        referenced it, corrupting the transmitted data. (remm)
+      </fix>
       <!-- Entries for backport and removal before 12.0.0-M1 below this line 
-->
     </changelog>
   </subsection>


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to