This is an automated email from the ASF dual-hosted git repository. rmaucher pushed a commit to branch 11.0.x in repository https://gitbox.apache.org/repos/asf/tomcat.git
commit 29a8a486babc9098edb683d1a888d3c5097b4f68 Author: opencode <[email protected]> AuthorDate: Fri Oct 9 11:34:41 2026 +0200 Do not send an RPC reply when the callback returns null RpcCallback.replyRequest is documented to return null if no reply should be sent, but RpcChannel.messageReceived ignored that and sent a reply message carrying the null payload. The requester then received a Response with a null message, which could complete first reply or majority reply collectors early with an empty payload or consume a slot of an all reply collector. Return without sending anything when the callback returns null, and notify an ExtendedRpcCallback of the completion synchronously, passing the original request message and a null response. Add TestRpcChannel covering the no-reply, synchronous notification, asynchronous silent and regular reply paths. --- .../apache/catalina/tribes/group/RpcChannel.java | 8 + .../catalina/tribes/group/TestRpcChannel.java | 163 +++++++++++++++++++++ 2 files changed, 171 insertions(+) diff --git a/java/org/apache/catalina/tribes/group/RpcChannel.java b/java/org/apache/catalina/tribes/group/RpcChannel.java index ca64c5c61f..f5847ed714 100644 --- a/java/org/apache/catalina/tribes/group/RpcChannel.java +++ b/java/org/apache/catalina/tribes/group/RpcChannel.java @@ -170,6 +170,14 @@ public class RpcChannel implements ChannelListener { boolean asyncReply = ((replyMessageOptions & Channel.SEND_OPTIONS_ASYNCHRONOUS) == Channel.SEND_OPTIONS_ASYNCHRONOUS); Serializable reply = callback.replyRequest(rmsg.message, sender); + if (reply == null) { + // The callback is documented to return null when no reply + // should be sent + if (excallback != null && !asyncReply) { + excallback.replySucceeded(msg, null, sender); + } + return; + } ErrorHandler handler = null; final Serializable request = msg; final Serializable response = reply; diff --git a/test/org/apache/catalina/tribes/group/TestRpcChannel.java b/test/org/apache/catalina/tribes/group/TestRpcChannel.java new file mode 100644 index 0000000000..0967b84b3a --- /dev/null +++ b/test/org/apache/catalina/tribes/group/TestRpcChannel.java @@ -0,0 +1,163 @@ +/* + * 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; + +import java.io.Serializable; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.List; + +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.ErrorHandler; +import org.apache.catalina.tribes.Member; +import org.apache.catalina.tribes.UniqueId; +import org.apache.catalina.tribes.membership.MemberImpl; + +public class TestRpcChannel { + + private static class RecordingChannel extends GroupChannel { + + private final List<Serializable> sent = new ArrayList<>(); + + @Override + public UniqueId send(Member[] destination, Serializable msg, int options) throws ChannelException { + sent.add(msg); + return null; + } + + @Override + public UniqueId send(Member[] destination, Serializable msg, int options, ErrorHandler handler) + throws ChannelException { + sent.add(msg); + return null; + } + } + + private static class FixedCallback implements RpcCallback { + + private final Serializable reply; + + FixedCallback(Serializable reply) { + this.reply = reply; + } + + @Override + public Serializable replyRequest(Serializable msg, Member sender) { + return reply; + } + + @Override + public void leftOver(Serializable msg, Member sender) { + // NO-OP + } + } + + private RpcChannel createRpcChannel(RecordingChannel channel, RpcCallback callback) { + return new RpcChannel(new byte[] { 1 }, channel, callback); + } + + @Test + public void testNullReplySendsNothing() throws Exception { + RecordingChannel channel = new RecordingChannel(); + RpcChannel rpcChannel = createRpcChannel(channel, new FixedCallback(null)); + Member sender = new MemberImpl("localhost", 4000, -1); + RpcMessage request = + new RpcMessage(rpcChannel.getRpcId(), "uuid".getBytes(StandardCharsets.UTF_8), "REQUEST"); + + rpcChannel.messageReceived(request, sender); + + Assert.assertTrue("A null reply means no reply should be sent", channel.sent.isEmpty()); + } + + private static class RecordingExtendedCallback extends FixedCallback implements ExtendedRpcCallback { + + private volatile int succeededCount; + private volatile Serializable succeededRequest; + private volatile Serializable succeededResponse; + + RecordingExtendedCallback(Serializable reply) { + super(reply); + } + + @Override + public void replyFailed(Serializable request, Serializable response, Member sender, Exception reason) { + // NO-OP + } + + @Override + public void replySucceeded(Serializable request, Serializable response, Member sender) { + succeededRequest = request; + succeededResponse = response; + succeededCount++; + } + } + + @Test + public void testNullReplySynchronousExtendedCallbackNotified() throws Exception { + RecordingChannel channel = new RecordingChannel(); + RecordingExtendedCallback callback = new RecordingExtendedCallback(null); + RpcChannel rpcChannel = createRpcChannel(channel, callback); + Member sender = new MemberImpl("localhost", 4000, -1); + RpcMessage request = + new RpcMessage(rpcChannel.getRpcId(), "uuid".getBytes(StandardCharsets.UTF_8), "REQUEST"); + + rpcChannel.messageReceived(request, sender); + + Assert.assertTrue(channel.sent.isEmpty()); + Assert.assertEquals(1, callback.succeededCount); + Assert.assertSame(request, callback.succeededRequest); + Assert.assertNull(callback.succeededResponse); + } + + @Test + public void testNullReplyAsynchronousExtendedCallbackNotNotified() throws Exception { + RecordingChannel channel = new RecordingChannel(); + RecordingExtendedCallback callback = new RecordingExtendedCallback(null); + RpcChannel rpcChannel = createRpcChannel(channel, callback); + rpcChannel.setReplyMessageOptions(Channel.SEND_OPTIONS_ASYNCHRONOUS); + Member sender = new MemberImpl("localhost", 4000, -1); + RpcMessage request = + new RpcMessage(rpcChannel.getRpcId(), "uuid".getBytes(StandardCharsets.UTF_8), "REQUEST"); + + rpcChannel.messageReceived(request, sender); + + Assert.assertTrue(channel.sent.isEmpty()); + Assert.assertEquals(0, callback.succeededCount); + } + + @Test + public void testNonNullReplyIsSent() throws Exception { + RecordingChannel channel = new RecordingChannel(); + RpcChannel rpcChannel = createRpcChannel(channel, new FixedCallback("RESPONSE")); + Member sender = new MemberImpl("localhost", 4000, -1); + RpcMessage request = + new RpcMessage(rpcChannel.getRpcId(), "uuid".getBytes(StandardCharsets.UTF_8), "REQUEST"); + + rpcChannel.messageReceived(request, sender); + + Assert.assertEquals(1, channel.sent.size()); + Object sent = channel.sent.get(0); + Assert.assertTrue(sent instanceof RpcMessage); + RpcMessage reply = (RpcMessage) sent; + Assert.assertTrue(reply.reply); + Assert.assertEquals("RESPONSE", reply.message); + } +} --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
