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]
