Skip to content

Commit 4aee133

Browse files
authored
Merge pull request #2016 from rabbitmq/mergify/bp/v5.x/pr-2015
Honor unlimited negotiated frame_max when capping inbound message size (backport #2015)
2 parents e62665b + de2e6ad commit 4aee133

5 files changed

Lines changed: 181 additions & 1 deletion

File tree

‎src/main/java/com/rabbitmq/client/impl/AMQConnection.java‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -431,12 +431,17 @@ public void start()
431431
int frameMax =
432432
negotiatedMaxValue(this.requestedFrameMax,
433433
connTune.getFrameMax());
434+
435+
if (frameMax < 0) {
436+
throw new IllegalArgumentException("Negotiated frame max cannot be negative: " + frameMax);
437+
}
438+
434439
this._frameMax = frameMax;
435440

436441
// Inbound payload limit: the smaller of frame_max (less framing
437442
// overhead) and the configured message body cap.
438443
_frameHandler.setFrameMax(
439-
Math.min(this.maxInboundMessageBodySize, frameMax));
444+
Utils.inboundFrameMax(this.maxInboundMessageBodySize, frameMax));
440445

441446
int negotiatedHeartbeat =
442447
negotiatedMaxValue(this.requestedHeartbeat,

‎src/main/java/com/rabbitmq/client/impl/Utils.java‎

Lines changed: 11 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -74,6 +74,17 @@ static <T> Consumer<T> noOpConsumer() {
7474
return (Consumer<T>) NO_OP_CONSUMER;
7575
}
7676

77+
/**
78+
* Caps the inbound message body size to the smaller of {@code maxInboundMessageBodySize}
79+
* and the negotiated {@code frameMax}, treating a {@code frameMax} of 0 as "no limit"
80+
* (per the AMQP connection.tune negotiation) rather than a literal, smaller-than-everything
81+
* value.
82+
*/
83+
static int inboundFrameMax(int maxInboundMessageBodySize, int frameMax) {
84+
int effectiveFrameMax = frameMax == 0 ? Integer.MAX_VALUE : frameMax;
85+
return Math.min(maxInboundMessageBodySize, effectiveFrameMax);
86+
}
87+
7788
static int framePayloadLimit(int frameMax) {
7889
if (frameMax <= 0) {
7990
return Integer.MAX_VALUE;
Lines changed: 38 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,38 @@
1+
// Copyright (c) 2007-2026 Broadcom. All Rights Reserved. The term "Broadcom" refers to Broadcom
2+
// Inc. and/or its subsidiaries.
3+
//
4+
// This software, the RabbitMQ Java client library, is triple-licensed under the
5+
// Mozilla Public License 2.0 ("MPL"), the GNU General Public License version 2
6+
// ("GPL") and the Apache License version 2 ("ASL"). For the MPL, please see
7+
// LICENSE-MPL-RabbitMQ. For the GPL, please see LICENSE-GPL2. For the ASL,
8+
// please see LICENSE-APACHE2.
9+
//
10+
// This software is distributed on an "AS IS" basis, WITHOUT WARRANTY OF ANY KIND,
11+
// either express or implied. See the LICENSE file for specific language governing
12+
// rights and limitations of this software.
13+
//
14+
// If you have any questions regarding licensing, please contact us at
15+
// info@rabbitmq.com.
16+
17+
package com.rabbitmq.client.impl;
18+
19+
import static org.assertj.core.api.Assertions.assertThat;
20+
21+
import org.junit.jupiter.params.ParameterizedTest;
22+
import org.junit.jupiter.params.provider.CsvSource;
23+
24+
public class UtilsTest {
25+
26+
@CsvSource({
27+
// maxInboundMessageBodySize, frameMax, expected
28+
"65536,0,65536", // negotiated frame_max of 0 means "no limit", body cap must still apply
29+
"65536,16384,16384", // frame_max smaller than the body cap wins
30+
"16384,65536,16384", // body cap smaller than frame_max wins
31+
"65536,65536,65536",
32+
})
33+
@ParameterizedTest
34+
void inboundFrameMaxHonoursFrameMaxZeroAsUnlimited(
35+
int maxInboundMessageBodySize, int frameMax, int expected) {
36+
assertThat(Utils.inboundFrameMax(maxInboundMessageBodySize, frameMax)).isEqualTo(expected);
37+
}
38+
}

‎src/test/java/com/rabbitmq/client/test/ClientTestSuite.java‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -92,6 +92,7 @@
9292
ByteBufferPublishTest.class,
9393
PublishWithByteBufferTest.class,
9494
InboundFrameMax.class,
95+
UtilsTest.class,
9596
})
9697
public class ClientTestSuite {
9798

‎src/test/java/com/rabbitmq/client/test/InboundFrameMax.java‎

Lines changed: 125 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,7 @@
3131
import java.util.concurrent.CountDownLatch;
3232
import java.util.concurrent.TimeUnit;
3333
import java.util.concurrent.atomic.AtomicReference;
34+
import org.junit.jupiter.api.Test;
3435
import org.junit.jupiter.params.ParameterizedTest;
3536
import org.junit.jupiter.params.provider.ValueSource;
3637

@@ -92,6 +93,130 @@ private void doEnforceInboundFrameMax(int frameMax, int openOkPayloadSize, boole
9293
}
9394
}
9495

96+
// A negotiated frame_max of 0 means "no limit" (client requests 0, the default, and the
97+
// broker also declares 0). The configured message body cap must still be enforced in that
98+
// case, instead of being silently discarded in favor of an effectively unbounded frame size.
99+
@ParameterizedTest
100+
@ValueSource(ints = {8192, 10_000})
101+
void bodySizeCapShouldPassWhenEqualToLimitAndNegotiatedFrameMaxIsUnlimited(
102+
int maxInboundMessageBodySize) throws Exception {
103+
int openOkPayloadSize = maxInboundMessageBodySize - AMQCommand.EMPTY_FRAME_SIZE;
104+
doEnforceInboundFrameMaxWithUnlimitedNegotiatedFrameMax(
105+
maxInboundMessageBodySize, openOkPayloadSize, false);
106+
}
107+
108+
@ParameterizedTest
109+
@ValueSource(ints = {8192, 10_000})
110+
void bodySizeCapShouldNotPassWhenAboveLimitAndNegotiatedFrameMaxIsUnlimited(
111+
int maxInboundMessageBodySize) throws Exception {
112+
int openOkPayloadSize = maxInboundMessageBodySize - AMQCommand.EMPTY_FRAME_SIZE + 1;
113+
doEnforceInboundFrameMaxWithUnlimitedNegotiatedFrameMax(
114+
maxInboundMessageBodySize, openOkPayloadSize, true);
115+
}
116+
117+
private void doEnforceInboundFrameMaxWithUnlimitedNegotiatedFrameMax(
118+
int maxInboundMessageBodySize, int openOkPayloadSize, boolean shouldFail) throws Exception {
119+
120+
CountDownLatch serverDone = new CountDownLatch(1);
121+
AtomicReference<Throwable> serverError = new AtomicReference<>();
122+
123+
try (ServerSocket server = new ServerSocket(0, 1, InetAddress.getByName("127.0.0.1"))) {
124+
int port = server.getLocalPort();
125+
// frame_max of 0 in connection.tune means "no limit"
126+
Thread peer =
127+
new Thread(
128+
() -> runFakeBroker(server, serverDone, serverError, 0, openOkPayloadSize),
129+
"fake-amqp-broker");
130+
peer.setDaemon(true);
131+
peer.start();
132+
133+
ConnectionFactory factory = TestUtils.connectionFactory();
134+
factory.setHost("127.0.0.1");
135+
factory.setPort(port);
136+
// requested frame max of 0 (the client default) is what allows the negotiated
137+
// frame_max to come out to 0 when the broker also declares 0
138+
factory.setRequestedFrameMax(0);
139+
factory.setMaxInboundMessageBodySize(maxInboundMessageBodySize);
140+
factory.setHandshakeTimeout(5000);
141+
factory.setConnectionTimeout(5000);
142+
factory.setRequestedHeartbeat(0);
143+
144+
if (shouldFail) {
145+
assertThatThrownBy(() -> factory.newConnection()).isInstanceOf(IOException.class);
146+
} else {
147+
try (Connection connection = factory.newConnection()) {
148+
assertThat(connection.getFrameMax()).isEqualTo(0);
149+
}
150+
}
151+
}
152+
assertThat(serverDone.await(5, TimeUnit.SECONDS)).isTrue();
153+
if (!shouldFail) {
154+
assertThat(serverError.get()).isNull();
155+
}
156+
}
157+
158+
// connection.tune's frame_max is parsed as a signed 32-bit value; a broker (or a
159+
// misconfigured client requesting a nonzero frame max) can drive the negotiated value
160+
// negative, which must be rejected rather than silently accepted as a bogus limit.
161+
@Test
162+
void negativeNegotiatedFrameMaxShouldBeRejected() throws Exception {
163+
CountDownLatch serverDone = new CountDownLatch(1);
164+
AtomicReference<Throwable> serverError = new AtomicReference<>();
165+
166+
try (ServerSocket server = new ServerSocket(0, 1, InetAddress.getByName("127.0.0.1"))) {
167+
int port = server.getLocalPort();
168+
Thread peer =
169+
new Thread(
170+
() -> runFakeBrokerWithNegativeFrameMax(server, serverDone, serverError),
171+
"fake-amqp-broker");
172+
peer.setDaemon(true);
173+
peer.start();
174+
175+
ConnectionFactory factory = TestUtils.connectionFactory();
176+
factory.setHost("127.0.0.1");
177+
factory.setPort(port);
178+
// a nonzero requested frame max is needed for negotiatedMaxValue to take the
179+
// Math.min(positive, negative) branch, instead of laundering the broker's negative
180+
// value into 0 ("no limit") the way it would if the client requested 0
181+
factory.setRequestedFrameMax(131_072);
182+
factory.setAutomaticRecoveryEnabled(false);
183+
factory.setHandshakeTimeout(5000);
184+
factory.setConnectionTimeout(5000);
185+
factory.setRequestedHeartbeat(0);
186+
187+
// asserting on the message, not just the type, matters here: on the Netty transport,
188+
// an unguarded negative frame max reaches Netty's frame decoder and trips its own
189+
// IllegalArgumentException ("maxFrameLength ... expected: > 0") for an unrelated
190+
// reason, which would otherwise make this test pass without the fix in place
191+
assertThatThrownBy(() -> factory.newConnection())
192+
.isInstanceOf(IllegalArgumentException.class)
193+
.hasMessageContaining("Negotiated frame max cannot be negative");
194+
}
195+
assertThat(serverDone.await(5, TimeUnit.SECONDS)).isTrue();
196+
assertThat(serverError.get()).isNull();
197+
}
198+
199+
private static void runFakeBrokerWithNegativeFrameMax(
200+
ServerSocket server, CountDownLatch done, AtomicReference<Throwable> error) {
201+
try (Socket socket = server.accept()) {
202+
socket.setSoTimeout(5000);
203+
DataInputStream in = new DataInputStream(socket.getInputStream());
204+
DataOutputStream out = new DataOutputStream(socket.getOutputStream());
205+
206+
byte[] header = new byte[8];
207+
in.readFully(header);
208+
writeMethodFrame(out, startPayload());
209+
readFrame(in);
210+
writeMethodFrame(out, tunePayload(-1));
211+
// the client must reject the negotiated frame max before replying, so no further
212+
// frames are expected
213+
} catch (Throwable t) {
214+
error.set(t);
215+
} finally {
216+
done.countDown();
217+
}
218+
}
219+
95220
private static void runFakeBroker(
96221
ServerSocket server,
97222
CountDownLatch done,

0 commit comments

Comments
 (0)