diff --git a/fluss-common/src/main/java/org/apache/fluss/utils/ProtoCodecUtils.java b/fluss-common/src/main/java/org/apache/fluss/utils/ProtoCodecUtils.java index 9c83471524..daa37f9b68 100644 --- a/fluss-common/src/main/java/org/apache/fluss/utils/ProtoCodecUtils.java +++ b/fluss-common/src/main/java/org/apache/fluss/utils/ProtoCodecUtils.java @@ -178,6 +178,20 @@ public static int readVarInt(ByteBuf buf) { return result; } + // Reads a length-delimited field's length, rejecting a value that cannot fit in the remaining + // buffer (or is negative) before the caller allocates for it. + public static int readBytesLen(ByteBuf buf) { + int len = readVarInt(buf); + if (len < 0 || len > buf.readableBytes()) { + throw new IllegalArgumentException( + "Protobuf field length " + + len + + " exceeds remaining readable bytes " + + buf.readableBytes()); + } + return len; + } + public static long readVarInt64(ByteBuf buf) { int shift = 0; long result = 0; diff --git a/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufBytesField.java b/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufBytesField.java index 0d28697ec5..0cb0c97cb0 100644 --- a/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufBytesField.java +++ b/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufBytesField.java @@ -35,7 +35,7 @@ public void declaration(PrintWriter w) { @Override public void parse(PrintWriter w) { - w.format("_%sLen = ProtoCodecUtils.readVarInt(_buffer);\n", ccName); + w.format("_%sLen = ProtoCodecUtils.readBytesLen(_buffer);\n", ccName); w.format("%s = new byte[_%sLen];\n", ccName, ccName); w.format("_buffer.readBytes(%s);\n", ccName); } diff --git a/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufRepeatedBytesField.java b/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufRepeatedBytesField.java index 566745bb3e..fbe2bd0f3e 100644 --- a/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufRepeatedBytesField.java +++ b/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufRepeatedBytesField.java @@ -43,7 +43,7 @@ public void parse(PrintWriter w) { w.format( "ProtoCodecUtils.BytesHolder _%sBh = _%sBytesHolder();\n", ccName, ProtoGenUtil.camelCase("new", singularName)); - w.format("_%sBh.len = ProtoCodecUtils.readVarInt(_buffer);\n", ccName); + w.format("_%sBh.len = ProtoCodecUtils.readBytesLen(_buffer);\n", ccName); w.format("_%sBh.b = new byte[_%sBh.len];\n", ccName, ccName); w.format("_buffer.readBytes(_%sBh.b);\n", ccName); } diff --git a/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufRepeatedStringField.java b/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufRepeatedStringField.java index 0e1c981eaf..1b0aeeafdf 100644 --- a/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufRepeatedStringField.java +++ b/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufRepeatedStringField.java @@ -43,7 +43,7 @@ public void parse(PrintWriter w) { w.format( "ProtoCodecUtils.StringHolder _%sSh = _%sStringHolder();\n", ccName, ProtoGenUtil.camelCase("new", singularName)); - w.format("_%sSh.len = ProtoCodecUtils.readVarInt(_buffer);\n", ccName); + w.format("_%sSh.len = ProtoCodecUtils.readBytesLen(_buffer);\n", ccName); w.format( "_%sSh.s = ProtoCodecUtils.readString(_buffer, _buffer.readerIndex(), _%sSh.len);\n", ccName, ccName); diff --git a/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufStringField.java b/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufStringField.java index 5e5fdfc0ff..11862f64cf 100644 --- a/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufStringField.java +++ b/fluss-protogen/fluss-protogen-generator/src/main/java/org/apache/fluss/protogen/generator/generator/ProtobufStringField.java @@ -113,7 +113,7 @@ public void serialize(PrintWriter w) { @Override public void parse(PrintWriter w) { - w.format("_%sLen = ProtoCodecUtils.readVarInt(_buffer);\n", ccName); + w.format("_%sLen = ProtoCodecUtils.readBytesLen(_buffer);\n", ccName); w.format("int _%sBufferIdx = _buffer.readerIndex();\n", ccName); w.format( "%s = ProtoCodecUtils.readString(_buffer, _buffer.readerIndex(), _%sLen);\n", diff --git a/fluss-protogen/fluss-protogen-tests/src/test/java/org/apache/fluss/protogen/tests/FieldLengthBoundTest.java b/fluss-protogen/fluss-protogen-tests/src/test/java/org/apache/fluss/protogen/tests/FieldLengthBoundTest.java new file mode 100644 index 0000000000..5522e31fc3 --- /dev/null +++ b/fluss-protogen/fluss-protogen-tests/src/test/java/org/apache/fluss/protogen/tests/FieldLengthBoundTest.java @@ -0,0 +1,61 @@ +/* + * 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.fluss.protogen.tests; + +import org.apache.fluss.rpc.messages.ApiMessage; +import org.apache.fluss.shaded.netty4.io.netty.buffer.ByteBuf; +import org.apache.fluss.shaded.netty4.io.netty.buffer.Unpooled; + +import org.junit.jupiter.api.Test; + +import static org.assertj.core.api.Assertions.assertThatThrownBy; + +/** Tests that a field declaring more bytes than the buffer holds is rejected, not allocated for. */ +public class FieldLengthBoundTest { + + @Test + public void testBytesField() { + assertOversizedFieldRejected(new B(), 1); // B.payload + } + + @Test + public void testRepeatedBytesField() { + assertOversizedFieldRejected(new B(), 2); // B.extra_items + } + + @Test + public void testStringField() { + assertOversizedFieldRejected(new S(), 1); // S.id + } + + @Test + public void testRepeatedStringField() { + assertOversizedFieldRejected(new S(), 2); // S.names + } + + private void assertOversizedFieldRejected(ApiMessage message, int fieldNumber) { + ByteBuf buf = Unpooled.buffer(); + buf.writeByte((fieldNumber << 3) | 2); // tag: field number + wire type 2 (length-delimited) + buf.writeByte(100); // length prefix claims 100 bytes... + buf.writeByte(0).writeByte(0); // ...but only two follow. + + assertThatThrownBy(() -> message.parseFrom(buf, buf.readableBytes())) + .isInstanceOf(IllegalArgumentException.class) + .hasMessageContaining("remaining readable bytes"); + } +}