Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -54,7 +54,6 @@
import org.apache.hadoop.hdds.utils.ConnectionFailureUtils;
import org.apache.hadoop.io.retry.RetryPolicy;
import org.apache.hadoop.ozone.common.Checksum;
import org.apache.hadoop.ozone.common.ChecksumData;
import org.apache.hadoop.security.token.Token;
import org.apache.ratis.protocol.exceptions.TimeoutIOException;
import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
Expand Down Expand Up @@ -232,7 +231,7 @@ private int preadOnce(XceiverClientGrpc client, StreamingReader reader, Pipeline
throw new IOException("Uninitialized StreamingReadResponse: " + blockID);
}
client.streamRead(ContainerProtocolCalls.buildReadBlockCommandProto(
blockID, blockOffset, length, responseDataSize, tokenRef.get(), pipeline), response);
blockID, blockOffset, length, responseDataSize, verifyChecksum, tokenRef.get(), pipeline), response);

int copied = 0;
while (copied < length) {
Expand Down Expand Up @@ -475,7 +474,7 @@ synchronized void readBlockImpl(long length) throws IOException {
throw new IOException("Uninitialized StreamingReadResponse: " + blockID);
}
xceiverClient.streamRead(ContainerProtocolCalls.buildReadBlockCommandProto(
blockID, requestedLength, length, responseDataSize, tokenRef.get(), pipelineRef.get()), r);
blockID, requestedLength, length, responseDataSize, verifyChecksum, tokenRef.get(), pipelineRef.get()), r);
}

private void handleExceptions(IOException cause) throws IOException {
Expand Down Expand Up @@ -696,8 +695,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman
try {
ByteBuffer data = readBlock.getData().asReadOnlyByteBuffer();
if (verifyChecksum) {
ChecksumData checksumData = ChecksumData.getFromProtoBuf(readBlock.getChecksumData());
Checksum.verifyChecksum(data, checksumData, 0);
Checksum.validateChecksums(data, readBlock.getOffset(), 0, readBlock.getChunkInfoListList());
Comment thread
peterxcli marked this conversation as resolved.
}
offerToQueue(readBlock);
} catch (Exception e) {
Expand All @@ -709,7 +707,7 @@ public void onNext(ContainerProtos.ContainerCommandResponseProto containerComman
LOG.warn("Failed to process block {} response at offset={}, size={}: {}, {}",
getBlockID().getContainerBlockID(),
offset, data.size(), StringUtils.bytes2Hex(data.asReadOnlyByteBuffer(), 10),
readBlock.getChecksumData(), e);
readBlock.getChunkInfoListList(), e);
if (r != null) {
r.getRequestObserver().onError(e);
}
Expand Down Expand Up @@ -783,10 +781,8 @@ private void setCompleted() {

private void offerToQueue(ReadBlockResponseProto item) {
if (LOG.isTraceEnabled()) {
final ContainerProtos.ChecksumData checksumData = item.getChecksumData();
LOG.trace("{}: enqueue response offset {}, length {}, numChecksums {}, bytesPerChecksum={}",
name, item.getOffset(), item.getData().size(),
checksumData.getChecksumsList().size(), checksumData.getBytesPerChecksum());
LOG.trace("{}: enqueue response offset {}, length {}, chunks {}",
name, item.getOffset(), item.getData().size(), item.getChunkInfoListCount());
}
final boolean offered = responseQueue.offer(item);
Preconditions.assertTrue(offered, () -> "Failed to offer " + item);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@
import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
import org.apache.hadoop.hdds.security.token.OzoneBlockTokenIdentifier;
import org.apache.hadoop.ozone.common.Checksum;
import org.apache.hadoop.ozone.common.OzoneChecksumException;
import org.apache.hadoop.security.token.Token;
import org.apache.ratis.protocol.exceptions.TimeoutIOException;
Expand All @@ -76,6 +77,7 @@
import org.apache.ratis.thirdparty.io.grpc.StatusRuntimeException;
import org.apache.ratis.thirdparty.io.grpc.stub.ClientCallStreamObserver;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;

/**
* Tests for StreamBlockInputStream custom configuration behavior.
Expand Down Expand Up @@ -1138,10 +1140,12 @@ private ReadBlockResponseProto buildReadBlockResponse(byte[] data) {
return ReadBlockResponseProto.newBuilder()
.setOffset(0)
.setData(ByteString.copyFrom(data))
.setChecksumData(ChecksumData.newBuilder()
.setType(ContainerProtos.ChecksumType.NONE)
.setBytesPerChecksum(data.length)
.build())
.addChunkInfoList(ContainerProtos.ChunkInfo.newBuilder()
.setChunkName("chunk").setOffset(0).setLen(data.length)
.setChecksumData(ChecksumData.newBuilder()
.setType(ContainerProtos.ChecksumType.NONE)
.setBytesPerChecksum(data.length)
.build()))
.build();
}

Expand All @@ -1152,10 +1156,12 @@ private ContainerCommandResponseProto buildResponseProto(byte[] data, long offse
.setReadBlock(ReadBlockResponseProto.newBuilder()
.setOffset(offset)
.setData(ByteString.copyFrom(data))
.setChecksumData(ChecksumData.newBuilder()
.setType(ContainerProtos.ChecksumType.NONE)
.setBytesPerChecksum(data.length)
.build())
.addChunkInfoList(ContainerProtos.ChunkInfo.newBuilder()
.setChunkName("chunk").setOffset(offset).setLen(data.length)
.setChecksumData(ChecksumData.newBuilder()
.setType(ContainerProtos.ChecksumType.NONE)
.setBytesPerChecksum(data.length)
.build()))
.build())
.build();
}
Expand All @@ -1167,12 +1173,96 @@ private ContainerCommandResponseProto buildCorruptResponseProto(byte[] data, lon
.setReadBlock(ReadBlockResponseProto.newBuilder()
.setOffset(offset)
.setData(ByteString.copyFrom(data))
.setChecksumData(ChecksumData.newBuilder()
.setType(ContainerProtos.ChecksumType.CRC32)
.setBytesPerChecksum(data.length)
.addChecksums(ByteString.copyFrom(new byte[4]))
.build())
.addChunkInfoList(ContainerProtos.ChunkInfo.newBuilder()
.setChunkName("chunk").setOffset(offset).setLen(data.length)
.setChecksumData(ChecksumData.newBuilder()
.setType(ContainerProtos.ChecksumType.CRC32)
.setBytesPerChecksum(data.length)
.addChecksums(ByteString.copyFrom(new byte[4]))
.build()))
.build())
.build();
}

private static ContainerProtos.ChunkInfo checksummedChunk(byte[] data, int offset, int length) throws Exception {
return ContainerProtos.ChunkInfo.newBuilder().setChunkName("chunk-" + offset).setOffset(offset).setLen(length)
.setChecksumData(new Checksum(ContainerProtos.ChecksumType.CRC32, 4)
.computeChecksum(ByteBuffer.wrap(data, offset, length)).getProtoBufMessage()).build();
}

private static ContainerCommandResponseProto response(byte[] data, long offset,
List<ContainerProtos.ChunkInfo> chunks) {
return ContainerCommandResponseProto.newBuilder().setCmdType(Type.ReadBlock)
.setResult(ContainerProtos.Result.SUCCESS).setReadBlock(ReadBlockResponseProto.newBuilder()
.setOffset(offset).setData(ByteString.copyFrom(data)).addAllChunkInfoList(chunks)).build();
}

@Test
void testChunkChecksumsAndSeekPrefix() throws Exception {
byte[] data = new byte[] {0, 1, 2, 3, 4, 5, 6, 7, 8, 9, 10, 11};
List<ContainerProtos.ChunkInfo> chunks = Arrays.asList(
checksummedChunk(data, 0, 3), checksummedChunk(data, 3, 9));
for (int seek : new int[] {0, 4, 8, 11}) {
int offset = seek == 0 ? 0 : 3 + (seek - 3) / 4 * 4;
ContainerCommandResponseProto response = response(Arrays.copyOfRange(data, offset, data.length), offset,
offset == 0 ? chunks : chunks.subList(1, 2));
assertStreamResponse(response, data, seek, true, false);
}
}

@Test
void testMissingMetadataAndVerificationDisabled() throws Exception {
byte[] data = new byte[] {0, 1, 2, 3};
ContainerCommandResponseProto missing = response(data, 0, Collections.emptyList());
assertStreamResponse(missing, data, 0, true, true);
assertStreamResponse(missing, data, 0, false, false);
assertStreamResponse(buildResponseProto(data, 0), data, 0, true, false);
byte[] corruptData = data.clone();
corruptData[0]++;
ContainerCommandResponseProto corrupt = response(corruptData, 0,
Collections.singletonList(checksummedChunk(data, 0, data.length)));
assertStreamResponse(corrupt, data, 0, true, true);
assertStreamResponse(corrupt, corruptData, 0, false, false);
}

private void assertStreamResponse(ContainerCommandResponseProto response, byte[] expected, int seek,
boolean verifyChecksum, boolean fails) throws Exception {
assertStreamResponse(response, expected, seek, verifyChecksum, fails, false);
assertStreamResponse(response, expected, seek, verifyChecksum, fails, true);
}

private void assertStreamResponse(ContainerCommandResponseProto response, byte[] expected, int seek,
boolean verifyChecksum, boolean fails, boolean positioned) throws Exception {
OzoneClientConfig config = newStreamReadConfig();
config.setChecksumVerify(verifyChecksum);
config.setMaxReadRetryCount(0);
ClientCallStreamObserver<ContainerCommandRequestProto> observer = mock(ClientCallStreamObserver.class);
XceiverClientGrpc client = mockCapturingStreamingReadClient(observer, reader -> reader.onNext(response));
XceiverClientFactory factory = mock(XceiverClientFactory.class);
when(factory.acquireClientForReadData(any(Pipeline.class))).thenReturn(client);
try (StreamBlockInputStream stream = new StreamBlockInputStream(new BlockID(1, 16258), expected.length,
mockStandalonePipeline(), null, factory, NO_REFRESH, config)) {
stream.seek(seek);
byte[] actual = new byte[expected.length - seek];
if (fails) {
IOException error = assertThrows(IOException.class, () -> {
if (positioned) {
stream.readPositioned(seek, ByteBuffer.wrap(actual));
} else {
stream.read(actual);
}
});
assertThat(hasCause(error, OzoneChecksumException.class)).isTrue();
assertArrayEquals(new byte[actual.length], actual);
} else {
assertEquals(actual.length, positioned
? stream.readPositioned(seek, ByteBuffer.wrap(actual)) : stream.read(actual));
assertArrayEquals(Arrays.copyOfRange(expected, seek, expected.length), actual);
}
}
ArgumentCaptor<ContainerCommandRequestProto> request = ArgumentCaptor.forClass(ContainerCommandRequestProto.class);
verify(client).streamRead(request.capture(), any());
assertEquals(verifyChecksum, request.getValue().getReadBlock().getIncludeChecksums());
verify(client, times(1)).completeStreamRead();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -44,7 +44,6 @@
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Type;
import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.WriteChunkResponseProto;
import org.apache.hadoop.ozone.common.ChecksumData;
import org.apache.hadoop.ozone.common.ChunkBufferToByteString;
import org.apache.ratis.thirdparty.com.google.protobuf.ByteString;
import org.apache.ratis.thirdparty.com.google.protobuf.UnsafeByteOperations;
Expand Down Expand Up @@ -344,13 +343,14 @@ public static ContainerCommandResponseProto getReadChunkResponse(
}

public static ContainerCommandResponseProto getReadBlockResponse(
ContainerCommandRequestProto request, ChecksumData checksumData, ByteBuffer data, long offset) {
ContainerCommandRequestProto request, List<ChunkInfo> chunks, ByteBuffer data, long offset) {

ContainerProtos.ReadBlockResponseProto response = ContainerProtos.ReadBlockResponseProto.newBuilder()
.setChecksumData(checksumData.getProtoBufMessage())
ContainerProtos.ReadBlockResponseProto.Builder response = ContainerProtos.ReadBlockResponseProto.newBuilder()
.setData(ByteString.copyFrom(data))
.setOffset(offset)
.build();
.setOffset(offset);
if (request.getReadBlock().getIncludeChecksums()) {
response.addAllChunkInfoList(chunks);
}

return getSuccessResponseBuilder(request)
.setReadBlock(response)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -950,7 +950,7 @@ public static List<Validator> toValidatorList(Validator validator) {
}

public static ContainerCommandRequestProto buildReadBlockCommandProto(
BlockID blockID, long offset, long length, int responseDataSize,
BlockID blockID, long offset, long length, int responseDataSize, boolean includeChecksums,
Token<? extends TokenIdentifier> token, Pipeline pipeline)
throws IOException {
final DatanodeDetails datanode = pipeline.getClosestNode();
Expand All @@ -959,6 +959,7 @@ public static ContainerCommandRequestProto buildReadBlockCommandProto(
.setOffset(offset)
.setLength(length)
.setResponseDataSize(responseDataSize)
.setIncludeChecksums(includeChecksums)
.setBlockID(blockID.getDatanodeBlockIDProtobufBuilder(replicaIndex));
final ContainerCommandRequestProto.Builder builder =
ContainerCommandRequestProto.newBuilder().setCmdType(Type.ReadBlock)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,8 +17,6 @@

package org.apache.hadoop.ozone.common;

import static org.apache.ratis.util.Preconditions.assertSame;

import com.google.common.annotations.VisibleForTesting;
import java.nio.ByteBuffer;
import java.security.MessageDigest;
Expand Down Expand Up @@ -441,37 +439,37 @@ public static void verifyChecksum(List<ByteBuffer> bufferList, int startIndex, C
public static void validateChecksums(ByteBuffer data, long blockOffset, int startIndex,

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Adding those check is nice, but could you clarify the motivation behind refactoring this method? From my perspective, this change seems unnecessary. Please leave this part as it is.

final List<ChunkInfo> chunks) throws OzoneChecksumException {

long firstChunkOffset = blockOffset - chunks.get(startIndex).getOffset();
int bytesPerChecksum = chunks.get(startIndex).getChecksumData().getBytesPerChecksum();
long readLength = data.remaining();
int dataOffset = data.position();

if (readLength <= 0) {
if (!data.hasRemaining()) {
return;
}

assertSame(0, firstChunkOffset % bytesPerChecksum, "blockOffset % bytesPerChecksum");
ContainerProtos.ChunkInfo firstChunk = chunks.get(startIndex);

// verify first chunk, only the first chunk is not start at zero.
int firstChunkIndex = (int) (firstChunkOffset / bytesPerChecksum);
int dataLimit = (int) Math.min(firstChunk.getLen() - firstChunkOffset, readLength);
Checksum.verifySingleChunk(data, dataOffset, dataLimit, firstChunk, firstChunkIndex);
dataOffset += dataLimit;
readLength -= dataLimit;
startIndex++;

while (readLength > 0) {
ContainerProtos.ChunkInfo chunkInfo = chunks.get(startIndex);
dataLimit = (int) Math.min(chunkInfo.getLen(), readLength);

Checksum.verifySingleChunk(data, dataOffset, dataLimit, chunkInfo, 0);

dataOffset += dataLimit;
readLength -= dataLimit;
startIndex++;
int dataOffset = data.position();
int remaining = data.remaining();
long offset = blockOffset;
while (remaining > 0) {
if (startIndex >= chunks.size()) {
throw new OzoneChecksumException("Missing chunk metadata at offset " + offset);
}
ChunkInfo chunk = chunks.get(startIndex++);
long chunkOffset = chunk.getOffset();
long length = chunk.getLen();
if (offset < chunkOffset || offset - chunkOffset >= length
|| (offset != blockOffset && offset != chunkOffset)) {
throw new OzoneChecksumException("Invalid chunk coverage at offset " + offset);
}
long relativeOffset = offset - chunkOffset;
int size = (int) Math.min(remaining, length - relativeOffset);
if (chunk.getChecksumData().getType() != ChecksumType.NONE) {
int bytesPerChecksum = chunk.getChecksumData().getBytesPerChecksum();
if (bytesPerChecksum <= 0 || relativeOffset % bytesPerChecksum != 0
|| (relativeOffset + size != length && (relativeOffset + size) % bytesPerChecksum != 0)) {
throw new OzoneChecksumException("Invalid checksum boundary at offset " + offset);
}
verifySingleChunk(data, dataOffset, size, chunk, (int) (relativeOffset / bytesPerChecksum));
}
dataOffset += size;
offset += size;
remaining -= size;
}

}

private static void verifySingleChunk(ByteBuffer data, int dataOffset, int dataLimit,
Expand Down
Loading
Loading