diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java index de1b8b61a5e2..a905129c5e6f 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/utils/db/Proto2Codec.java @@ -22,11 +22,9 @@ import com.google.protobuf.Parser; import jakarta.annotation.Nonnull; import java.io.IOException; -import java.io.InputStream; import java.io.OutputStream; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; -import org.apache.hadoop.hdds.utils.IOUtils; import org.apache.ratis.util.function.CheckedFunction; /** @@ -89,13 +87,10 @@ public String toString() { @Override public M fromCodecBuffer(@Nonnull CodecBuffer buffer) throws CodecException { - final InputStream in = buffer.getInputStream(); try { - return parser.parseFrom(in); + return parser.parseFrom(buffer.asReadOnlyByteBuffer()); } catch (InvalidProtocolBufferException e) { throw new CodecException("Failed to parse " + buffer + " for " + getTypeClass(), e); - } finally { - IOUtils.closeQuietly(in); } } diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/CodecTestUtil.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/CodecTestUtil.java index 14ebbd63db27..3488ae30c796 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/CodecTestUtil.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/CodecTestUtil.java @@ -99,6 +99,12 @@ public static void runTest(Codec codec, T original, final T fromWrappedArray = codec.fromCodecBuffer(wrapped); wrapped.release(); assertEquals(original, fromWrappedArray); + + // deserialize from direct CodecBuffer, which is what table get and iterator use + try (CodecBuffer direct = codec.toCodecBuffer(original, CodecBuffer.Allocator.getDirect())) { + assertTrue(direct.isDirect()); + assertEquals(original, codec.fromCodecBuffer(direct)); + } } public static Codec newCodecWithoutCodecBuffer(Codec codec) { diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/Proto2CodecTestBase.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/Proto2CodecTestBase.java index 9de62ce0cc68..c4429a7495d0 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/Proto2CodecTestBase.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/utils/db/Proto2CodecTestBase.java @@ -42,6 +42,19 @@ public void testInvalidProtocolBuffer() { .contains("the input ended unexpectedly"); } + /** + * {@link Proto2Codec#fromCodecBuffer(CodecBuffer)} parses the buffer in place; a buffer which is + * not a valid message must still surface as a {@link CodecException}. + */ + @Test + public void testInvalidProtocolBufferFromCodecBuffer() { + try (CodecBuffer buffer = CodecBuffer.wrap("random".getBytes(UTF_8))) { + final CodecException exception = + assertThrows(CodecException.class, () -> getCodec().fromCodecBuffer(buffer)); + assertInstanceOf(InvalidProtocolBufferException.class, exception.getCause()); + } + } + @Test public void testFromPersistedFormat() { assertThrows(NullPointerException.class, diff --git a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestCodec.java b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestCodec.java index 649c8a46f929..94c560fa6b7e 100644 --- a/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestCodec.java +++ b/hadoop-hdds/framework/src/test/java/org/apache/hadoop/hdds/utils/db/TestCodec.java @@ -36,6 +36,7 @@ import java.util.UUID; import java.util.concurrent.ThreadLocalRandom; import java.util.function.Consumer; +import org.apache.hadoop.hdds.protocol.proto.HddsProtos.KeyValue; import org.apache.hadoop.hdds.utils.db.RDBBatchOperation.Bytes; import org.apache.hadoop.hdds.utils.db.managed.ManagedRocksObjectUtils; import org.junit.jupiter.api.Test; @@ -272,6 +273,19 @@ static void runTestByteStringCodec(ByteString original) throws Exception { runTest(ByteStringCodec.get(), original, original.size()); } + @Test + public void testProto2Codec() throws Exception { + final Codec codec = Proto2Codec.get(KeyValue.getDefaultInstance()); + for (int i = 0; i < NUM_LOOPS; i++) { + final KeyValue original = KeyValue.newBuilder() + .setKey("key" + i) + .setValue("value" + ThreadLocalRandom.current().nextLong()) + .build(); + runTest(codec, original, original.getSerializedSize()); + } + gc(); + } + static Executable tryCatch(Executable executable) { return tryCatch(executable, t -> LOG.info("Good!", t)); } diff --git a/hadoop-ozone/interface-storage/src/test/java/org/apache/hadoop/ozone/om/helpers/TestTransactionInfoCodec.java b/hadoop-ozone/interface-storage/src/test/java/org/apache/hadoop/ozone/om/helpers/TestTransactionInfoCodec.java index 7640fb20fbb8..ad2e77240028 100644 --- a/hadoop-ozone/interface-storage/src/test/java/org/apache/hadoop/ozone/om/helpers/TestTransactionInfoCodec.java +++ b/hadoop-ozone/interface-storage/src/test/java/org/apache/hadoop/ozone/om/helpers/TestTransactionInfoCodec.java @@ -24,6 +24,7 @@ import java.nio.charset.StandardCharsets; import org.apache.hadoop.hdds.utils.TransactionInfo; import org.apache.hadoop.hdds.utils.db.Codec; +import org.apache.hadoop.hdds.utils.db.CodecBuffer; import org.apache.hadoop.hdds.utils.db.Proto2CodecTestBase; import org.junit.jupiter.api.Test; @@ -54,4 +55,14 @@ public void testInvalidProtocolBuffer() { () -> getCodec().fromPersistedFormat("random".getBytes(StandardCharsets.UTF_8))); assertThat(ex).hasMessageContaining("Unexpected split length"); } + + @Override + @Test + public void testInvalidProtocolBufferFromCodecBuffer() { + try (CodecBuffer buffer = CodecBuffer.wrap("random".getBytes(StandardCharsets.UTF_8))) { + IllegalArgumentException ex = assertThrows(IllegalArgumentException.class, + () -> getCodec().fromCodecBuffer(buffer)); + assertThat(ex).hasMessageContaining("Unexpected split length"); + } + } }