This is an automated email from the ASF dual-hosted git repository.
jojochuang pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/master by this push:
new 54da65a9d45 HDDS-15356. Addendum: Make multi-buffer chunk checksum
allocation-free (#11008)
54da65a9d45 is described below
commit 54da65a9d4551b3768419ea02476313824b446ed
Author: Siyao Meng <[email protected]>
AuthorDate: Thu Sep 3 08:16:51 2026 -0700
HDDS-15356. Addendum: Make multi-buffer chunk checksum allocation-free
(#11008)
Generated-by: Codex (GPT-5.6 Sol)
---
.../org/apache/hadoop/ozone/common/Checksum.java | 12 ++++++++
.../apache/hadoop/ozone/common/ChecksumCache.java | 10 ++++++-
.../apache/hadoop/ozone/common/TestChecksum.java | 16 +++++++++++
.../hadoop/ozone/common/TestChecksumCache.java | 33 ++++++++++++++++++++++
4 files changed, 70 insertions(+), 1 deletion(-)
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java
index f4010497774..aefe4399a46 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/Checksum.java
@@ -259,6 +259,8 @@ public ChecksumData computeChecksum(List<ByteString>
byteStrings)
/**
* The default implementation of computeChecksum(ChunkBuffer) that does not
use cache, even if cache is initialized.
* This is a stop-gap solution before the protocol change.
+ * The input must be prepared for reading as described by
+ * {@link #computeChecksum(ChunkBuffer, boolean)}.
* @param data ChunkBuffer
* @return ChecksumData
* @throws OzoneChecksumException
@@ -269,6 +271,10 @@ public ChecksumData computeChecksum(ChunkBuffer data)
}
/**
+ * The input must be prepared for reading. For each buffer returned by
+ * {@link ChunkBuffer#asByteBufferList()}, the bytes in
+ * {@code [position(), limit())} must be the logical checksum data.
+ *
* This method does not advance the positions of {@code data}'s underlying
* buffers. Both the no-cache and cache paths slice via
* {@link ByteBuffer#duplicate()}.
@@ -305,6 +311,7 @@ private List<ByteString> computeChecksumDirect(ChunkBuffer
data,
final int checksumCount = dataLength == 0 ? 0 : 1 + (dataLength - 1) /
bytesPerChecksum;
final List<ByteString> result = new ArrayList<>(checksumCount);
int windowRemaining = bytesPerChecksum;
+ long processed = 0;
algo.reset();
for (ByteBuffer src : data.asByteBufferList()) {
@@ -314,6 +321,7 @@ private List<ByteString> computeChecksumDirect(ChunkBuffer
data,
final int n = Math.min(srcLim - srcPos, windowRemaining);
algo.update(BufferUtils.slice(src, srcPos, n));
srcPos += n;
+ processed += n;
windowRemaining -= n;
if (windowRemaining == 0) {
result.add(algo.finish());
@@ -322,6 +330,10 @@ private List<ByteString> computeChecksumDirect(ChunkBuffer
data,
}
}
}
+ if (processed != dataLength) {
+ throw new IllegalStateException("ChunkBuffer remaining byte count is " +
dataLength
+ + ", but its underlying buffers expose " + processed + " bytes");
+ }
if (windowRemaining < bytesPerChecksum) {
// Unaligned trailing window.
result.add(algo.finish());
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/ChecksumCache.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/ChecksumCache.java
index b59fbdb3f04..b487fef0083 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/ChecksumCache.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/common/ChecksumCache.java
@@ -87,6 +87,14 @@ List<ByteString> computeChecksum(ChunkBuffer data,
+ bytesPerChecksum + " call=" + chksumSize);
}
final int currChunkLength = data.remaining();
+ final List<ByteBuffer> buffers = data.asByteBufferList();
+ long readableByteCount = 0;
+ for (ByteBuffer buffer : buffers) {
+ readableByteCount += buffer.remaining();
+ }
+ Preconditions.checkState(readableByteCount == currChunkLength,
+ "ChunkBuffer remaining byte count is %s, but its underlying buffers
expose %s bytes",
+ currChunkLength, readableByteCount);
if (currChunkLength == prevChunkLength) {
LOG.debug("ChunkBuffer data length same as last time ({}). "
@@ -111,7 +119,7 @@ List<ByteString> computeChecksum(ChunkBuffer data,
algo.reset();
long position = 0;
- for (ByteBuffer src : data.asByteBufferList()) {
+ for (ByteBuffer src : buffers) {
int srcPos = src.position();
final int srcLim = src.limit();
final int srcLen = srcLim - srcPos;
diff --git
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java
index ce597971c04..7ccc5d8547e 100644
---
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java
+++
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksum.java
@@ -20,6 +20,7 @@
import static java.nio.charset.StandardCharsets.UTF_8;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.when;
@@ -117,4 +118,19 @@ public void testChecksumFromNonzeroPosition() throws
Exception {
final Checksum checksum = getChecksum(null, false);
assertEquals(1, checksum.computeChecksum(data).getChecksums().size());
}
+
+ @ParameterizedTest
+ @ValueSource(booleans = {true, false})
+ public void testRejectsIncrementalBufferInWriteMode(boolean
useChecksumCache) {
+ try (ChunkBuffer data = ChunkBuffer.allocate(32, 8)) {
+ data.put(new byte[10]);
+
+ final Checksum checksum = getChecksum(null, useChecksumCache);
+ final IllegalStateException exception = assertThrows(
+ IllegalStateException.class,
+ () -> checksum.computeChecksum(data, useChecksumCache));
+ assertEquals("ChunkBuffer remaining byte count is 22, but its underlying
buffers expose 6 bytes",
+ exception.getMessage());
+ }
+ }
}
diff --git
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksumCache.java
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksumCache.java
index 0d0253c2130..163e12c3b97 100644
---
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksumCache.java
+++
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/common/TestChecksumCache.java
@@ -25,6 +25,7 @@
import org.apache.hadoop.ozone.common.Checksum.StreamingChecksum;
import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
import org.junit.jupiter.params.ParameterizedTest;
import org.junit.jupiter.params.provider.EnumSource;
@@ -111,6 +112,38 @@ void
testPartiallyPositionedMultiBufferMatchesDirectChecksum(
Assertions.assertEquals(expected.getChecksums(), actual.getChecksums());
}
+ @Test
+ void testRejectsInvalidBufferWhenLengthMatchesCache() throws Exception {
+ final Checksum checksum = new Checksum(ChecksumType.SHA256, 10, true);
+ checksum.computeChecksum(ByteBuffer.wrap(new byte[10]), true);
+
+ try (ChunkBuffer data = ChunkBuffer.allocate(20, 10)) {
+ data.put(new byte[10]);
+ final IllegalStateException exception = Assertions.assertThrows(
+ IllegalStateException.class,
+ () -> checksum.computeChecksum(data, true));
+ Assertions.assertEquals(
+ "ChunkBuffer remaining byte count is 10, but its underlying buffers
expose 0 bytes",
+ exception.getMessage());
+ }
+ }
+
+ @Test
+ void testRejectedBufferDoesNotModifyCache() throws Exception {
+ final Checksum cached = new Checksum(ChecksumType.SHA256, 10, true);
+ try (ChunkBuffer data = ChunkBuffer.allocate(80, 32)) {
+ data.put(new byte[33]);
+ Assertions.assertThrows(IllegalStateException.class,
+ () -> cached.computeChecksum(data, true));
+ }
+
+ final ByteBuffer valid = ByteBuffer.wrap(new byte[10]);
+ final ChecksumData expected = new Checksum(ChecksumType.SHA256, 10)
+ .computeChecksum(valid.duplicate());
+ final ChecksumData actual = cached.computeChecksum(valid, true);
+ Assertions.assertEquals(expected, actual);
+ }
+
private static ChunkBuffer split(byte[] data, int length, int bufferSize) {
final List<ByteBuffer> buffers = new ArrayList<>();
for (int offset = 0; offset < length; offset += bufferSize) {
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]