diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java index 249e58df94d4..46d31f8156dd 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/impl/HddsDispatcher.java @@ -80,6 +80,7 @@ import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; import org.apache.ratis.util.UncheckedAutoCloseable; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -807,13 +808,14 @@ private boolean isAllowed(String action) { @Override public StateMachine.DataChannel getStreamDataChannel( - ContainerCommandRequestProto msg) - throws StorageContainerException { + ContainerCommandRequestProto msg, + CheckedConsumer putBlock) + throws StorageContainerException { long containerID = msg.getContainerID(); Container container = getContainer(containerID); if (container != null) { Handler handler = getHandler(getContainerType(container)); - return handler.getStreamDataChannel(container, msg); + return handler.getStreamDataChannel(container, msg, putBlock); } else { throw new StorageContainerException( "ContainerID " + containerID + " does not exist", diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/ContainerDispatcher.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/ContainerDispatcher.java index 2aba8253cefe..af6822443d8b 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/ContainerDispatcher.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/ContainerDispatcher.java @@ -17,6 +17,7 @@ package org.apache.hadoop.ozone.container.common.interfaces; +import java.io.IOException; import java.util.Map; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; @@ -26,6 +27,7 @@ import org.apache.hadoop.ozone.container.common.transport.server.ratis.DispatcherContext; import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +import org.apache.ratis.util.function.CheckedConsumer; /** * Dispatcher acts as the bridge between the transport layer and @@ -87,7 +89,9 @@ void validateContainerCommand( * When uploading using stream, get StreamDataChannel. */ default StateMachine.DataChannel getStreamDataChannel( - ContainerCommandRequestProto msg) throws StorageContainerException { + ContainerCommandRequestProto msg, + CheckedConsumer putBlock) + throws StorageContainerException { throw new UnsupportedOperationException( "getStreamDataChannel not supported."); } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java index fec625f89a08..6fe253ffbc7d 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/interfaces/Handler.java @@ -46,6 +46,7 @@ import org.apache.hadoop.ozone.container.keyvalue.TarContainerPacker; import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +import org.apache.ratis.util.function.CheckedConsumer; /** * Dispatcher sends ContainerCommandRequests to Handler. Each Container Type @@ -93,7 +94,8 @@ public static Handler getHandlerForContainerType( } public abstract StateMachine.DataChannel getStreamDataChannel( - Container container, ContainerCommandRequestProto msg) + Container container, ContainerCommandRequestProto msg, + CheckedConsumer putBlock) throws StorageContainerException; /** diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java index 3de4110a01b3..9a3e988e27ad 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/transport/server/ratis/ContainerStateMachine.java @@ -110,6 +110,7 @@ import org.apache.ratis.util.JavaUtils; import org.apache.ratis.util.LifeCycle; import org.apache.ratis.util.TaskQueue; +import org.apache.ratis.util.function.CheckedConsumer; import org.apache.ratis.util.function.CheckedSupplier; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -726,7 +727,17 @@ private StateMachine.DataChannel getStreamDataChannel( requestProto.getTraceID()); } dispatchCommand(requestProto, context); // stream init - return dispatcher.getStreamDataChannel(requestProto); + final CheckedConsumer putBlock + = requestProto.getCmdType() == Type.StreamInitWithPutBlock ? this::streamPutBlock : null; + return dispatcher.getStreamDataChannel(requestProto, putBlock); + } + + void streamPutBlock(ContainerCommandRequestProto request) throws IOException { + final DispatcherContext context = DispatcherContext.newBuilder(DispatcherContext.Op.STREAM_LINK) + .setStage(DispatcherContext.WriteChunkStage.COMBINED) + .setContainer2BCSIDMap(container2BCSIDMap) + .build(); + dispatchCommand(request, context); } @Override @@ -741,7 +752,7 @@ public CompletableFuture stream(RaftClientRequest request) { .setStage(DispatcherContext.WriteChunkStage.WRITE_DATA) .setContainer2BCSIDMap(container2BCSIDMap) .build(); - DataChannel channel = getStreamDataChannel(requestProto, context); + final DataChannel channel = getStreamDataChannel(requestProto, context); final ExecutorService chunkExecutor = requestProto.hasWriteChunk() ? getChunkExecutor(requestProto.getWriteChunk()) : null; return new LocalStream(channel, chunkExecutor); @@ -773,7 +784,10 @@ public CompletableFuture link(DataStream stream, LogEntryProto entry) { final KeyValueStreamDataChannel kvStreamDataChannel = (KeyValueStreamDataChannel) dataChannel; - kvStreamDataChannel.setLinked(); + if (!kvStreamDataChannel.link()) { + return JavaUtils.completeExceptionally(new IllegalStateException( + "PutBlock was not committed on stream close: " + kvStreamDataChannel)); + } return CompletableFuture.completedFuture(null); } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java index 1b3f399b46fa..941f0479d4b0 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/KeyValueHandler.java @@ -173,6 +173,7 @@ import org.apache.ratis.statemachine.StateMachine; import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.grpc.stub.StreamObserver; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -272,7 +273,8 @@ public KeyValueHandler(ConfigurationSource config, @Override public StateMachine.DataChannel getStreamDataChannel( - Container container, ContainerCommandRequestProto msg) + Container container, ContainerCommandRequestProto msg, + CheckedConsumer putBlock) throws StorageContainerException { KeyValueContainer kvContainer = (KeyValueContainer) container; checkContainerOpen(kvContainer); @@ -282,7 +284,7 @@ public StateMachine.DataChannel getStreamDataChannel( BlockID.getFromProtobuf(msg.getWriteChunk().getBlockID()); return chunkManager.getStreamDataChannel(kvContainer, - blockID, metrics); + blockID, putBlock, metrics); } else { throw new StorageContainerException("Malformed request.", ContainerProtos.Result.IO_EXCEPTION); diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/ChunkManagerDispatcher.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/ChunkManagerDispatcher.java index a83306ff79be..4cb634229d99 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/ChunkManagerDispatcher.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/ChunkManagerDispatcher.java @@ -27,6 +27,7 @@ import java.util.Map; import java.util.Objects; import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.ozone.common.ChunkBuffer; import org.apache.hadoop.ozone.common.ChunkBufferToByteString; @@ -40,6 +41,7 @@ import org.apache.hadoop.ozone.container.keyvalue.interfaces.BlockManager; import org.apache.hadoop.ozone.container.keyvalue.interfaces.ChunkManager; import org.apache.ratis.statemachine.StateMachine; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -79,10 +81,12 @@ public String streamInit(Container container, BlockID blockID) @Override public StateMachine.DataChannel getStreamDataChannel( - Container container, BlockID blockID, ContainerMetrics metrics) + Container container, BlockID blockID, + CheckedConsumer putBlock, + ContainerMetrics metrics) throws StorageContainerException { return selectHandler(container) - .getStreamDataChannel(container, blockID, metrics); + .getStreamDataChannel(container, blockID, putBlock, metrics); } @Override diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/FilePerBlockStrategy.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/FilePerBlockStrategy.java index 9a13507d6b37..01000883f17e 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/FilePerBlockStrategy.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/FilePerBlockStrategy.java @@ -58,6 +58,7 @@ import org.apache.hadoop.ozone.container.keyvalue.interfaces.BlockManager; import org.apache.hadoop.ozone.container.keyvalue.interfaces.ChunkManager; import org.apache.ratis.statemachine.StateMachine; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -111,12 +112,13 @@ public String streamInit(Container container, BlockID blockID) @Override public StateMachine.DataChannel getStreamDataChannel( - Container container, BlockID blockID, ContainerMetrics metrics) - throws StorageContainerException { + Container container, BlockID blockID, + CheckedConsumer putBlock, + ContainerMetrics metrics) throws StorageContainerException { checkLayoutVersion(container); final File chunkFile = getChunkFile(container, blockID); return new KeyValueStreamDataChannel(chunkFile, - container.getContainerData(), metrics); + container.getContainerData(), putBlock, metrics); } @Override diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java index 3218fb4f88d9..5a735812d796 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/KeyValueStreamDataChannel.java @@ -22,14 +22,19 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; +import org.apache.hadoop.hdds.ratis.ContainerCommandRequestMessage; import org.apache.hadoop.hdds.ratis.RatisHelper; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.hdds.scm.storage.BlockDataStreamOutput; import org.apache.hadoop.ozone.container.common.helpers.ContainerMetrics; import org.apache.hadoop.ozone.container.common.impl.ContainerData; +import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.netty.buffer.ByteBuf; import org.apache.ratis.util.ReferenceCountedObject; +import org.apache.ratis.util.function.CheckedConsumer; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -41,12 +46,16 @@ public class KeyValueStreamDataChannel extends StreamDataChannelBase { private final Buffers buffers = new Buffers(BlockDataStreamOutput.PUT_BLOCK_REQUEST_LENGTH_MAX); + private final AtomicReference putBlockRequest + = new AtomicReference<>(); private final AtomicBoolean closed = new AtomicBoolean(); + private final CheckedConsumer putBlock; KeyValueStreamDataChannel(File file, ContainerData containerData, - ContainerMetrics metrics) - throws StorageContainerException { + CheckedConsumer putBlock, + ContainerMetrics metrics) throws StorageContainerException { super(file, containerData, metrics); + this.putBlock = putBlock; } @Override @@ -54,6 +63,11 @@ ContainerProtos.Type getType() { return ContainerProtos.Type.StreamWrite; } + @Override + boolean isPutBlockCommittedOnClose() { + return putBlock != null; + } + @Override public int write(ReferenceCountedObject referenceCounted) throws IOException { @@ -87,6 +101,10 @@ static void writeFully(ByteBuffer b, WriteMethod writeMethod) } } + public ContainerCommandRequestProto getPutBlockRequest() { + return putBlockRequest.get(); + } + void assertOpen() throws IOException { if (closed.get()) { throw new IOException("Already closed: " + this); @@ -97,7 +115,13 @@ void assertOpen() throws IOException { public void close() throws IOException { if (closed.compareAndSet(false, true)) { try { - writeBuffers(); + if (isPutBlockCommittedOnClose()) { + // This path requires the client appends the PutBlock at the + // end of the stream. + closeWithStreamPutBlock(); + } else { + writeBuffers(); + } } finally { super.close(); } @@ -128,6 +152,21 @@ private void writeBuffers() throws IOException { } } + static ContainerCommandRequestProto closeBuffers( + Buffers buffers, WriteMethod writeMethod) throws IOException { + final ReferenceCountedObject ref = buffers.pollAll(); + final ByteBuf buf = ref.retain(); + final ContainerCommandRequestProto putBlockRequestProto; + try { + putBlockRequestProto = readPutBlockRequest(buf); + // write the remaining data + writeFully(buf.nioBuffer(), writeMethod); + } finally { + ref.release(); + } + return putBlockRequestProto; + } + static int readProtoLength(ByteBuf b, int lengthIndex) { final int readerIndex = b.readerIndex(); LOG.debug("{}, lengthIndex = {}, readerIndex = {}", @@ -155,6 +194,65 @@ static void setEndIndex(ByteBuf b) { b.writerIndex(protoIndex); } + static ContainerCommandRequestProto readPutBlockRequest(ByteBuf b) + throws IOException { + // readerIndex protoIndex lengthIndex readerIndex+readableBytes + // V V V V + // format: |--- data ---|--- proto ---|--- proto length (4 bytes) ---| + final int readerIndex = b.readerIndex(); + final int lengthIndex = readerIndex + b.readableBytes() - 4; + if (lengthIndex < readerIndex) { + throw new IOException("Buffer too short for PutBlock length: " + b.readableBytes()); + } + final int protoLength = readProtoLength(b.duplicate(), lengthIndex); + final int protoIndex = lengthIndex - protoLength; + if (protoIndex < readerIndex) { + throw new IOException("Invalid PutBlock proto length: " + protoLength); + } + + final ContainerCommandRequestProto proto; + try { + proto = readPutBlockRequest(b.slice(protoIndex, protoLength).nioBuffer()); + } catch (Throwable t) { + RatisHelper.debug(b, "catch", LOG); + throw new IOException("Failed to readPutBlockRequest from " + b + + ": readerIndex=" + readerIndex + + ", protoIndex=" + protoIndex + + ", protoLength=" + protoLength + + ", lengthIndex=" + lengthIndex, t); + } + + // set index for reading data + b.writerIndex(protoIndex); + + return proto; + } + + private static ContainerCommandRequestProto readPutBlockRequest(ByteBuffer b) + throws IOException { + RatisHelper.debug(b, "readPutBlockRequest", LOG); + final ByteString byteString = ByteString.copyFrom(b); + + final ContainerCommandRequestProto request = + ContainerCommandRequestMessage.toProto(byteString, null); + + if (!request.hasPutBlock()) { + throw new StorageContainerException( + "Malformed PutBlock request. trace ID: " + request.getTraceID(), + ContainerProtos.Result.MALFORMED_REQUEST); + } + return request; + } + + private void closeWithStreamPutBlock() throws IOException { + final ContainerCommandRequestProto proto = + closeBuffers(buffers, super::writeFileChannel); + putBlockRequest.set(proto); + Preconditions.checkState(putBlock != null); + putBlock.accept(proto); + setLinked(); + } + interface WriteMethod { int applyAsInt(ByteBuffer src) throws IOException; } diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java index 43bcea5e9bd9..2cad1b052b21 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/StreamDataChannelBase.java @@ -94,6 +94,8 @@ public final boolean isOpen() { return getChannel().isOpen(); } + abstract boolean isPutBlockCommittedOnClose(); + protected void assertSpaceAvailability(int requested) throws StorageContainerException { ContainerUtils.assertSpaceAvailability(containerData.getContainerID(), containerData.getVolume(), requested); } @@ -102,6 +104,20 @@ public void setLinked() { linked.set(true); } + public boolean link() { + if (isPutBlockCommittedOnClose()) { + // The PutBlock should be commited when the steam is closed + return linked.get(); + } else { + setLinked(); + return true; + } + } + + public boolean isLinked() { + return linked.get(); + } + /** * @return true if {@link org.apache.ratis.statemachine.StateMachine.DataChannel} is already linked. */ diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/interfaces/ChunkManager.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/interfaces/ChunkManager.java index 0fc88a87ae3c..2cfcdf9c8e25 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/interfaces/ChunkManager.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/interfaces/ChunkManager.java @@ -32,6 +32,7 @@ import org.apache.hadoop.ozone.container.common.transport.server.ratis.DispatcherContext; import org.apache.hadoop.ozone.container.keyvalue.KeyValueContainer; import org.apache.ratis.statemachine.StateMachine; +import org.apache.ratis.util.function.CheckedConsumer; /** * Chunk Manager allows read, write, delete and listing of chunks in a container. @@ -114,7 +115,9 @@ default String streamInit(Container container, BlockID blockID) } default StateMachine.DataChannel getStreamDataChannel( - Container container, BlockID blockID, ContainerMetrics metrics) + Container container, BlockID blockID, + CheckedConsumer putBlock, + ContainerMetrics metrics) throws StorageContainerException { return null; } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java index 3c2992f36f4b..638c5707725e 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/keyvalue/impl/TestKeyValueStreamDataChannel.java @@ -23,7 +23,9 @@ import static org.apache.hadoop.ozone.container.keyvalue.impl.KeyValueStreamDataChannel.writeBuffers; import static org.apache.hadoop.ozone.container.keyvalue.impl.KeyValueStreamDataChannel.writeFully; import static org.assertj.core.api.Assertions.assertThat; +import static org.junit.jupiter.api.Assertions.assertArrayEquals; import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertThrows; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; @@ -33,6 +35,7 @@ import java.io.IOException; import java.nio.ByteBuffer; import java.nio.channels.WritableByteChannel; +import java.nio.file.Files; import java.util.ArrayList; import java.util.Collection; import java.util.List; @@ -41,16 +44,15 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.ThreadLocalRandom; +import java.util.concurrent.atomic.AtomicReference; import org.apache.commons.lang3.RandomUtils; import org.apache.hadoop.hdds.fs.SpaceUsageSource; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.BlockData; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.DatanodeBlockID; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.PutBlockRequestProto; -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.ratis.ContainerCommandRequestMessage; -import org.apache.hadoop.hdds.ratis.RatisHelper; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.container.common.helpers.ContainerMetrics; @@ -66,11 +68,13 @@ import org.apache.ratis.protocol.ClientId; import org.apache.ratis.protocol.DataStreamReply; import org.apache.ratis.protocol.RaftClientReply; -import org.apache.ratis.thirdparty.com.google.protobuf.ByteString; import org.apache.ratis.thirdparty.io.netty.buffer.ByteBuf; import org.apache.ratis.thirdparty.io.netty.buffer.Unpooled; import org.apache.ratis.util.ReferenceCountedObject; +import org.apache.ratis.util.function.CheckedConsumer; import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -95,6 +99,56 @@ public class TestKeyValueStreamDataChannel { LOG.info("PUT_BLOCK_PROTO_SIZE = {}", PUT_BLOCK_PROTO_SIZE); } + @ParameterizedTest + @ValueSource(booleans = {false, true}) + public void testClosePutBlockBehavior(boolean commitPutBlockOnClose) throws Exception { + File tempFile = File.createTempFile("test-kv-close-" + commitPutBlockOnClose, ".tmp"); + tempFile.deleteOnExit(); + AtomicReference processed = new AtomicReference<>(); + KeyValueStreamDataChannel channel = newChannel(tempFile, commitPutBlockOnClose, processed); + final byte[] data = RandomUtils.secure().randomBytes(50); + final ByteBuffer putBlockBuf = ContainerCommandRequestMessage.toMessage( + PUT_BLOCK_PROTO, null).getContent().asReadOnlyByteBuffer(); + final ByteBuffer protoLengthBuf = + getProtoLength(putBlockBuf, PUT_BLOCK_REQUEST_LENGTH_MAX); + + write(channel, data); + write(channel, putBlockBuf.duplicate()); + write(channel, protoLengthBuf.duplicate()); + + assertThat(processed.get()).isNull(); + channel.close(); + + if (commitPutBlockOnClose) { + assertEquals(PUT_BLOCK_PROTO, processed.get()); + assertEquals(PUT_BLOCK_PROTO, channel.getPutBlockRequest()); + assertTrue(channel.isLinked()); + } else { + assertThat(processed.get()).isNull(); + assertFalse(channel.isLinked()); + } + assertEquals(data.length, tempFile.length()); + assertArrayEquals(data, Files.readAllBytes(tempFile.toPath())); + } + + @Test + public void testReadPutBlockRequestBufferTooShort() { + final ByteBuf buf = Unpooled.buffer(2); + buf.writeByte(1); + buf.writeByte(2); + assertThrows(IOException.class, () -> KeyValueStreamDataChannel.readPutBlockRequest(buf)); + buf.release(); + } + + @Test + public void testReadPutBlockRequestInvalidProtoLength() { + final ByteBuf buf = Unpooled.buffer(8); + buf.writeInt(1); + buf.writeInt(100); + assertThrows(IOException.class, () -> KeyValueStreamDataChannel.readPutBlockRequest(buf)); + buf.release(); + } + @Test public void testSerialization() throws Exception { final int max = PUT_BLOCK_REQUEST_LENGTH_MAX; @@ -112,53 +166,10 @@ public void testSerialization() throws Exception { buf.writeBytes(putBlockBuf); buf.writeBytes(protoLengthBuf); - final ContainerCommandRequestProto proto = readPutBlockRequest(buf); + final ContainerCommandRequestProto proto = KeyValueStreamDataChannel.readPutBlockRequest(buf); assertEquals(PUT_BLOCK_PROTO, proto); } - static ContainerCommandRequestProto readPutBlockRequest(ByteBuf b) throws IOException { - // readerIndex protoIndex lengthIndex readerIndex+readableBytes - // V V V V - // format: |--- data ---|--- proto ---|--- proto length (4 bytes) ---| - final int readerIndex = b.readerIndex(); - final int lengthIndex = readerIndex + b.readableBytes() - 4; - final int protoLength = KeyValueStreamDataChannel.readProtoLength(b.duplicate(), lengthIndex); - final int protoIndex = lengthIndex - protoLength; - - final ContainerCommandRequestProto proto; - try { - proto = readPutBlockRequest(b.slice(protoIndex, protoLength).nioBuffer()); - } catch (Throwable t) { - RatisHelper.debug(b, "catch", LOG); - throw new IOException("Failed to readPutBlockRequest from " + b - + ": readerIndex=" + readerIndex - + ", protoIndex=" + protoIndex - + ", protoLength=" + protoLength - + ", lengthIndex=" + lengthIndex, t); - } - - // set index for reading data - b.writerIndex(protoIndex); - - return proto; - } - - private static ContainerCommandRequestProto readPutBlockRequest(ByteBuffer b) - throws IOException { - RatisHelper.debug(b, "readPutBlockRequest", LOG); - final ByteString byteString = ByteString.copyFrom(b); - - final ContainerCommandRequestProto request = - ContainerCommandRequestMessage.toProto(byteString, null); - - if (!request.hasPutBlock()) { - throw new StorageContainerException( - "Malformed PutBlock request. trace ID: " + request.getTraceID(), - Result.MALFORMED_REQUEST); - } - return request; - } - @Test public void testVolumeFullCase() throws Exception { File tempFile = File.createTempFile("test-kv-stream", ".tmp"); @@ -170,7 +181,8 @@ public void testVolumeFullCase() throws Exception { when(mockContainerData.getContainerID()).thenReturn(123L); when(mockContainerData.getVolume()).thenReturn(mockVolume); ContainerMetrics mockMetrics = mock(ContainerMetrics.class); - KeyValueStreamDataChannel writeChannel = new KeyValueStreamDataChannel(tempFile, mockContainerData, mockMetrics); + KeyValueStreamDataChannel writeChannel = + new KeyValueStreamDataChannel(tempFile, mockContainerData, null, mockMetrics); assertThrows(StorageContainerException.class, () -> writeChannel.assertSpaceAvailability(1)); final ByteBuffer putBlockBuf = ContainerCommandRequestMessage.toMessage( @@ -309,7 +321,7 @@ static ContainerCommandRequestProto closeBuffers( final ByteBuf buf = ref.retain(); final ContainerCommandRequestProto putBlockRequest; try { - putBlockRequest = readPutBlockRequest(buf); + putBlockRequest = KeyValueStreamDataChannel.readPutBlockRequest(buf); // write the remaining data writeFully(buf.nioBuffer(), writeMethod); } finally { @@ -401,4 +413,33 @@ static CompletableFuture completeExceptionally(Throwable t) { f.completeExceptionally(t); return f; } + + private static KeyValueStreamDataChannel newChannel( + File tempFile, boolean commitPutBlockOnClose, + AtomicReference processed) throws Exception { + HddsVolume mockVolume = mock(HddsVolume.class); + when(mockVolume.getStorageID()).thenReturn("storageId"); + when(mockVolume.getCurrentUsage()).thenReturn(new SpaceUsageSource.Fixed(1000L, 1000L, 0L)); + ContainerData mockContainerData = mock(ContainerData.class); + when(mockContainerData.getContainerID()).thenReturn(123L); + when(mockContainerData.getVolume()).thenReturn(mockVolume); + ContainerMetrics mockMetrics = mock(ContainerMetrics.class); + CheckedConsumer putBlock = + commitPutBlockOnClose ? processed::set : null; + return new KeyValueStreamDataChannel(tempFile, mockContainerData, putBlock, mockMetrics); + } + + private static void write(KeyValueStreamDataChannel channel, byte[] data) + throws IOException { + write(channel, ByteBuffer.wrap(data)); + } + + private static void write(KeyValueStreamDataChannel channel, ByteBuffer buf) + throws IOException { + ReferenceCountedObject ref = + ReferenceCountedObject.wrap(buf, () -> { }, () -> { }); + ref.retain(); + channel.write(ref); + ref.release(); + } } diff --git a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto index 05c94624c990..9438faafc7d9 100644 --- a/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto +++ b/hadoop-hdds/interface-client/src/main/proto/DatanodeClientProtocol.proto @@ -88,6 +88,8 @@ enum Type { GetContainerChecksumInfo = 23; // Allows us to read a block ReadBlock = 24; + // Initializes a stream for writing data with PutBlock commited on close. + StreamInitWithPutBlock = 25; }