Skip to content

Commit b491075

Browse files
michaelklishinmergify[bot]
authored andcommitted
Stricter validation of pre-negotiated values in the frame reader
(cherry picked from commit 08790f0)
1 parent a86e8f7 commit b491075

5 files changed

Lines changed: 101 additions & 2 deletions

File tree

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

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -433,6 +433,14 @@ public void start()
433433
connTune.getFrameMax());
434434
this._frameMax = frameMax;
435435

436+
// Bound inbound frames to the negotiated frame_max. EMPTY_FRAME_SIZE is
437+
// the per-frame overhead; +1 because the reader rejects payloads >= the limit.
438+
if (frameMax > 0) {
439+
_frameHandler.setMaxInboundFramePayloadSize(
440+
Math.min(this.maxInboundMessageBodySize,
441+
frameMax - AMQCommand.EMPTY_FRAME_SIZE + 1));
442+
}
443+
436444
int negotiatedHeartbeat =
437445
negotiatedMaxValue(this.requestedHeartbeat,
438446
connTune.getHeartbeat());

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

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -59,6 +59,11 @@ default void finishConnectionNegotiation() {
5959

6060
}
6161

62+
/** Cap inbound frame payloads, applied once frame_max is negotiated. */
63+
default void setMaxInboundFramePayloadSize(int maxPayloadSize) {
64+
65+
}
66+
6267
/**
6368
* Read a {@link Frame} from the underlying data connection.
6469
* @return an incoming Frame, or null if there is none

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

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -351,6 +351,11 @@ public void finishConnectionNegotiation() {
351351
maybeRemoveHandler(HANDLER_PROTOCOL_VERSION_MISMATCH);
352352
}
353353

354+
@Override
355+
public void setMaxInboundFramePayloadSize(int maxPayloadSize) {
356+
this.handler.maxPayloadSize = maxPayloadSize;
357+
}
358+
354359
@Override
355360
public Frame readFrame() {
356361
throw new UnsupportedOperationException();
@@ -506,7 +511,7 @@ InetSocketAddress maybeInetSocketAddress(SocketAddress socketAddress) {
506511

507512
private static class AmqpHandler extends ChannelInboundHandlerAdapter {
508513

509-
private final int maxPayloadSize;
514+
private volatile int maxPayloadSize;
510515
private final Runnable closeSequence;
511516
private final Predicate<ShutdownSignalException> willRecover;
512517
private volatile AMQConnection connection;

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

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ public class SocketFrameHandler implements FrameHandler {
5757
private final DataOutputStream _outputStream;
5858
private final Lock _outputStreamLock = new ReentrantLock();
5959

60-
private final int maxInboundMessageBodySize;
60+
private volatile int maxInboundMessageBodySize;
6161

6262
/** Time to linger before closing the socket forcefully. */
6363
public static final int SOCKET_CLOSING_TIMEOUT = 1;
@@ -193,6 +193,11 @@ public void initialize(AMQConnection connection) {
193193
connection.startMainLoop();
194194
}
195195

196+
@Override
197+
public void setMaxInboundFramePayloadSize(int maxPayloadSize) {
198+
this.maxInboundMessageBodySize = maxPayloadSize;
199+
}
200+
196201
@Override
197202
public Frame readFrame() throws IOException {
198203
_inputStreamLock.lock();
Lines changed: 76 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,76 @@
1+
// Copyright (c) 2026 Broadcom. All Rights Reserved. The term "Broadcom" refers to Broadcom Inc. and/or its subsidiaries.
2+
//
3+
// This software, the RabbitMQ Java client library, is triple-licensed under the
4+
// Mozilla Public License 2.0 ("MPL"), the GNU General Public License version 2
5+
// ("GPL") and the Apache License version 2 ("ASL"). For the MPL, please see
6+
// LICENSE-MPL-RabbitMQ. For the GPL, please see LICENSE-GPL2. For the ASL,
7+
// please see LICENSE-APACHE2.
8+
//
9+
// This software is distributed on an "AS IS" basis, WITHOUT WARRANTY OF ANY KIND,
10+
// either express or implied. See the LICENSE file for specific language governing
11+
// rights and limitations of this software.
12+
//
13+
// If you have any questions regarding licensing, please contact us at
14+
// info@rabbitmq.com.
15+
16+
package com.rabbitmq.client.test;
17+
18+
import static org.assertj.core.api.Assertions.assertThat;
19+
import static org.assertj.core.api.Assertions.assertThatThrownBy;
20+
21+
import com.rabbitmq.client.AMQP;
22+
import com.rabbitmq.client.MalformedFrameException;
23+
import com.rabbitmq.client.impl.AMQCommand;
24+
import com.rabbitmq.client.impl.Frame;
25+
import com.rabbitmq.client.impl.SocketFrameHandler;
26+
import java.io.DataOutputStream;
27+
import java.net.InetAddress;
28+
import java.net.ServerSocket;
29+
import java.net.Socket;
30+
import org.junit.jupiter.api.Test;
31+
32+
/**
33+
* Verifies that once the negotiated frame_max is applied to the inbound reader, a frame whose total
34+
* size exceeds frame_max is rejected, while a frame of exactly frame_max is still accepted.
35+
*/
36+
public class NegotiatedFrameMaxInboundTest {
37+
38+
private static final int FRAME_MAX = 4096;
39+
40+
// Same limit AMQConnection applies once frame_max is negotiated.
41+
private static final int INBOUND_PAYLOAD_LIMIT = FRAME_MAX - AMQCommand.EMPTY_FRAME_SIZE + 1;
42+
43+
@Test
44+
void frameOfExactlyFrameMaxIsAccepted() throws Exception {
45+
int payloadSize = FRAME_MAX - AMQCommand.EMPTY_FRAME_SIZE;
46+
Frame frame = readSingleFrame(payloadSize);
47+
assertThat(frame.getPayload()).hasSize(payloadSize);
48+
assertThat(frame.size()).isEqualTo(FRAME_MAX);
49+
}
50+
51+
@Test
52+
void frameLargerThanFrameMaxIsRejected() {
53+
int payloadSize = FRAME_MAX - AMQCommand.EMPTY_FRAME_SIZE + 1; // total frame size = frame_max + 1
54+
assertThatThrownBy(() -> readSingleFrame(payloadSize))
55+
.isInstanceOf(MalformedFrameException.class);
56+
}
57+
58+
private static Frame readSingleFrame(int payloadSize) throws Exception {
59+
InetAddress loopback = InetAddress.getLoopbackAddress();
60+
try (ServerSocket server = new ServerSocket(0, 1, loopback);
61+
Socket client = new Socket(loopback, server.getLocalPort());
62+
Socket peer = server.accept()) {
63+
DataOutputStream out = new DataOutputStream(peer.getOutputStream());
64+
out.writeByte(AMQP.FRAME_METHOD);
65+
out.writeShort(0);
66+
out.writeInt(payloadSize);
67+
out.write(new byte[payloadSize]);
68+
out.writeByte(AMQP.FRAME_END);
69+
out.flush();
70+
71+
SocketFrameHandler handler = new SocketFrameHandler(client);
72+
handler.setMaxInboundFramePayloadSize(INBOUND_PAYLOAD_LIMIT);
73+
return handler.readFrame();
74+
}
75+
}
76+
}

0 commit comments

Comments
 (0)