diff --git a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java index 4d09af278b15..744b1dda6670 100644 --- a/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java +++ b/hadoop-hdds/client/src/main/java/org/apache/hadoop/hdds/scm/storage/BlockOutputStream.java @@ -41,7 +41,9 @@ import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; import java.util.function.Supplier; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.BlockData; @@ -186,7 +188,6 @@ public BlockOutputStream( replicationIndex = pipeline.getReplicaIndex(pipeline.getClosestNode()); KeyValue keyValue = KeyValue.newBuilder().setKey("TYPE").setValue("KEY").build(); - ContainerProtos.DatanodeBlockID.Builder blkIDBuilder = ContainerProtos.DatanodeBlockID.newBuilder() .setContainerID(blockID.getContainerID()) @@ -195,6 +196,8 @@ public BlockOutputStream( if (replicationIndex > 0) { blkIDBuilder.setReplicaIndex(replicationIndex); } + // TODO: Replica to the method parameter + blkIDBuilder.setStorageTypeID(StorageTypeUtils.getID(StorageType.DISK)); this.containerBlockData = BlockData.newBuilder().setBlockID( blkIDBuilder.build()).addMetadata(keyValue); this.pipeline = pipeline; diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/BlockID.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/BlockID.java index 7141a65306dd..588aad88c60e 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/BlockID.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/BlockID.java @@ -19,6 +19,7 @@ import com.fasterxml.jackson.annotation.JsonIgnore; import java.util.Objects; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -34,29 +35,42 @@ public class BlockID { // BlockID object. private final Integer replicaIndex; + // Represents storage type of Block on a particular Datanode. + // Note this variable in the OM side will be null, + // because the OmKeyLocationInfo#getProtobuf Use BlockID#getProtobuf to get Protobuf object, + // BlockID#getProtobuf does not set storageType value in Protobuf. + + // Currently for java class BlockID, OM and Datanode share the same java class object, + // but the protobuf objects HddsProtos.BlockID and ContainerProtos.DatanodeBlockID are + // used for OM and Datanode respectively. + private StorageType storageType; + public BlockID(long containerID, long localID) { - this(containerID, localID, 0, null); + this(containerID, localID, 0, null, null); } - private BlockID(long containerID, long localID, long bcsID, Integer repIndex) { + private BlockID(long containerID, long localID, long bcsID, Integer repIndex, + StorageType storageType) { containerBlockID = new ContainerBlockID(containerID, localID); blockCommitSequenceId = bcsID; this.replicaIndex = repIndex; + this.storageType = storageType; } public BlockID(BlockID blockID) { this(blockID.getContainerID(), blockID.getLocalID(), blockID.getBlockCommitSequenceId(), - blockID.getReplicaIndex()); + blockID.getReplicaIndex(), blockID.getStorageType()); } public BlockID(ContainerBlockID containerBlockID) { - this(containerBlockID, 0, null); + this(containerBlockID, 0, null, null); } - private BlockID(ContainerBlockID containerBlockID, long bcsId, Integer repIndex) { + private BlockID(ContainerBlockID containerBlockID, long bcsId, Integer repIndex, StorageType storageType) { this.containerBlockID = containerBlockID; blockCommitSequenceId = bcsId; this.replicaIndex = repIndex; + this.storageType = storageType; } public long getContainerID() { @@ -84,6 +98,10 @@ public ContainerBlockID getContainerBlockID() { return containerBlockID; } + public StorageType getStorageType() { + return storageType; + } + @Override public String toString() { StringBuilder sb = new StringBuilder(64); @@ -94,7 +112,8 @@ public String toString() { public void appendTo(StringBuilder sb) { containerBlockID.appendTo(sb); sb.append(" bcsId: ").append(blockCommitSequenceId) - .append(" replicaIndex: ").append(replicaIndex); + .append(" replicaIndex: ").append(replicaIndex) + .append(" storageType: ").append(storageType); } @JsonIgnore @@ -103,6 +122,9 @@ public ContainerProtos.DatanodeBlockID getDatanodeBlockIDProtobuf() { if (replicaIndex != null) { blockID.setReplicaIndex(replicaIndex); } + if (storageType != null) { + blockID.setStorageTypeID(StorageTypeUtils.getID(storageType)); + } return blockID.build(); } @@ -116,10 +138,15 @@ public ContainerProtos.DatanodeBlockID.Builder getDatanodeBlockIDProtobufBuilder @JsonIgnore public static BlockID getFromProtobuf(ContainerProtos.DatanodeBlockID blockID) { + StorageType storageType = null; + if (blockID.hasStorageTypeID() && blockID.getStorageTypeID() > 0) { + storageType = StorageTypeUtils.getStorageTypeFromID(blockID.getStorageTypeID()); + } return new BlockID(blockID.getContainerID(), blockID.getLocalID(), blockID.getBlockCommitSequenceId(), - blockID.hasReplicaIndex() ? blockID.getReplicaIndex() : null); + blockID.hasReplicaIndex() ? blockID.getReplicaIndex() : null, + storageType); } @JsonIgnore @@ -133,7 +160,7 @@ public HddsProtos.BlockID getProtobuf() { public static BlockID getFromProtobuf(HddsProtos.BlockID blockID) { return new BlockID( ContainerBlockID.getFromProtobuf(blockID.getContainerBlockID()), - blockID.getBlockCommitSequenceId(), null); + blockID.getBlockCommitSequenceId(), null, null); } @Override @@ -147,12 +174,13 @@ public boolean equals(Object o) { BlockID blockID = (BlockID) o; return this.getContainerBlockID().equals(blockID.getContainerBlockID()) && this.getBlockCommitSequenceId() == blockID.getBlockCommitSequenceId() - && Objects.equals(this.getReplicaIndex(), blockID.getReplicaIndex()); + && Objects.equals(this.getReplicaIndex(), blockID.getReplicaIndex()) + && Objects.equals(this.getStorageType(), blockID.getStorageType()); } @Override public int hashCode() { return Objects.hash(containerBlockID.getContainerID(), containerBlockID.getLocalID(), - blockCommitSequenceId, replicaIndex); + blockCommitSequenceId, replicaIndex, storageType); } } diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTypeUtils.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTypeUtils.java index 6b94d16e57b3..0b56ccd55a20 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTypeUtils.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/client/StorageTypeUtils.java @@ -103,5 +103,13 @@ public static int getIDFromProtobuf(@Nonnull HddsProtos.StorageTypeProto proto) throw new IllegalArgumentException("Illegal Storage Type specified"); } } + + /** + * Returns integer representation of FileSystem StorageType. + * @return storageType int ID value + */ + public static int getID(StorageType storageType) throws IllegalArgumentException { + return getStorageTypeProto(storageType).getNumber(); + } } diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java index 95cc58effaf0..a05aed46c422 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/storage/ContainerProtocolCalls.java @@ -636,6 +636,7 @@ public static PutSmallFileResponseProto writeSmallFile( public static void createRecoveringContainer(XceiverClientSpi client, long containerID, String encodedToken, int replicaIndex) throws IOException { + // TODO StoragePolicy Support EC createContainer(client, containerID, encodedToken, ContainerProtos.ContainerDataProto.State.RECOVERING, replicaIndex); } diff --git a/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/container/ContainerTestHelper.java b/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/container/ContainerTestHelper.java index c2813602edef..96f585a0b529 100644 --- a/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/container/ContainerTestHelper.java +++ b/hadoop-hdds/common/src/test/java/org/apache/hadoop/ozone/container/ContainerTestHelper.java @@ -39,7 +39,9 @@ import java.util.concurrent.ThreadLocalRandom; import org.apache.commons.io.IOUtils; import org.apache.hadoop.conf.StorageUnit; +import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChecksumType; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerCommandRequestProto; @@ -572,15 +574,18 @@ public static ContainerProtos.ContainerCommandRequestProto getFinalizeBlockReque } public static BlockID getTestBlockID(long containerID) { - return getTestBlockID(containerID, null); + return getTestBlockID(containerID, null, null); } - public static BlockID getTestBlockID(long containerID, Integer replicaIndex) { + public static BlockID getTestBlockID(long containerID, Integer replicaIndex, StorageType storageType) { DatanodeBlockID.Builder datanodeBlockID = DatanodeBlockID.newBuilder().setContainerID(containerID) .setLocalID(UniqueId.next()); if (replicaIndex != null) { datanodeBlockID.setReplicaIndex(replicaIndex); } + if (storageType != null) { + datanodeBlockID.setStorageTypeID(StorageTypeUtils.getID(storageType)); + } return BlockID.getFromProtobuf(datanodeBlockID.build()); } 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 f1122a2debf7..6c1e9f1271e4 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 @@ -177,6 +177,7 @@ private boolean canIgnoreException(Result result) { case DELETE_ON_OPEN_CONTAINER: case UNSUPPORTED_REQUEST:// Blame client for sending unsupported request. case MALFORMED_REQUEST:// Blame client for sending malformed request. + case INVALID_ARGUMENT: case CONTAINER_MISSING: case CONTAINER_ALREADY_EXISTS: return true; @@ -194,6 +195,12 @@ public void buildMissingContainerSetAndValidate( @Override public ContainerCommandResponseProto dispatch( ContainerCommandRequestProto msg, DispatcherContext dispatcherContext) { + try { + HddsUtils.getBlockID(msg); + } catch (IllegalArgumentException e) { + return ContainerUtils.logAndReturnError(LOG, + new StorageContainerException(e.getMessage(), e, Result.INVALID_ARGUMENT), msg); + } try { return dispatcher.processRequest(msg, req -> dispatchRequest(msg, dispatcherContext), @@ -568,6 +575,11 @@ private void validateToken( @Override public void validateContainerCommand( ContainerCommandRequestProto msg) throws StorageContainerException { + try { + HddsUtils.getBlockID(msg); + } catch (IllegalArgumentException e) { + throw new StorageContainerException(e.getMessage(), e, Result.INVALID_ARGUMENT); + } try { validateToken(msg); } catch (IOException ioe) { 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 9a3e988e27ad..441f282e3484 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 @@ -1188,7 +1188,8 @@ public CompletableFuture applyTransaction(TransactionContext trx) { if (r.getResult() != ContainerProtos.Result.SUCCESS && r.getResult() != ContainerProtos.Result.CONTAINER_NOT_OPEN && r.getResult() != ContainerProtos.Result.CLOSED_CONTAINER_IO - && r.getResult() != ContainerProtos.Result.CHUNK_FILE_INCONSISTENCY) { + && r.getResult() != ContainerProtos.Result.CHUNK_FILE_INCONSISTENCY + && r.getResult() != ContainerProtos.Result.INVALID_ARGUMENT) { StorageContainerException sce = new StorageContainerException(r.getMessage(), r.getResult()); LOG.error( 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 693b33ab53d9..46f092f583dd 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 @@ -692,6 +692,8 @@ ContainerCommandResponseProto handlePutBlock( ContainerProtos.BlockData data = request.getPutBlock().getBlockData(); BlockData blockData = BlockData.getFromProtoBuf(data); Objects.requireNonNull(blockData, "blockData == null"); + Objects.requireNonNull(blockData.getBlockID()); + BlockUtils.verifyStorageType(kvContainer.getContainerData(), blockData.getBlockID()); boolean endOfBlock = false; if (!request.getPutBlock().hasEof() || request.getPutBlock().getEof()) { @@ -1114,6 +1116,7 @@ ContainerCommandResponseProto handleWriteChunk( WriteChunkRequestProto writeChunk = request.getWriteChunk(); BlockID blockID = BlockID.getFromProtobuf(writeChunk.getBlockID()); + BlockUtils.verifyStorageType(kvContainer.getContainerData(), blockID); ContainerProtos.ChunkInfo chunkInfoProto = writeChunk.getChunkData(); ChunkInfo chunkInfo = ChunkInfo.getFromProtoBuf(chunkInfoProto); @@ -1272,6 +1275,7 @@ ContainerCommandResponseProto handlePutSmallFile( BlockData blockData = BlockData.getFromProtoBuf( putSmallFileReq.getBlock().getBlockData()); Objects.requireNonNull(blockData, "blockData == null"); + BlockUtils.verifyStorageType(kvContainer.getContainerData(), blockData.getBlockID()); ContainerProtos.ChunkInfo chunkInfoProto = putSmallFileReq.getChunkInfo(); ChunkInfo chunkInfo = ChunkInfo.getFromProtoBuf(chunkInfoProto); diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/helpers/BlockUtils.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/helpers/BlockUtils.java index a05f44c1a775..0863a26e2fb0 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/helpers/BlockUtils.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/helpers/BlockUtils.java @@ -20,6 +20,7 @@ import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.CONTAINER_NOT_FOUND; import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.EXPORT_CONTAINER_METADATA_FAILED; import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.IMPORT_CONTAINER_METADATA_FAILED; +import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.INVALID_ARGUMENT; import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.NO_SUCH_BLOCK; import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.UNABLE_TO_READ_METADATA_DB; import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result.UNKNOWN_BCSID; @@ -37,6 +38,7 @@ import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.container.common.helpers.BlockData; +import org.apache.hadoop.ozone.container.common.impl.ContainerData; import org.apache.hadoop.ozone.container.common.interfaces.Container; import org.apache.hadoop.ozone.container.common.interfaces.DBHandle; import org.apache.hadoop.ozone.container.common.utils.ContainerCache; @@ -211,6 +213,23 @@ public static BlockData getBlockData(byte[] bytes) throws IOException { } } + /** + * Verify if request block storage type matches the container storage type. + * + * @param containerData container data + * @param blockID requested block info + * @throws StorageContainerException if storage types do not match + */ + public static void verifyStorageType(ContainerData containerData, BlockID blockID) + throws StorageContainerException { + if (blockID.getStorageType() != null && containerData.getStorageType() != null && + blockID.getStorageType() != containerData.getStorageType()) { + throw new StorageContainerException(String.format( + "Block storage type %s does not match container storage type %s", + blockID.getStorageType(), containerData.getStorageType()), INVALID_ARGUMENT); + } + } + /** * Verify if request block BCSID is supported. * diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/BlockManagerImpl.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/BlockManagerImpl.java index 6e69ed0480f2..4d7bd461ab46 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/BlockManagerImpl.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/keyvalue/impl/BlockManagerImpl.java @@ -187,6 +187,7 @@ public long persistPutBlock(KeyValueContainer container, Objects.requireNonNull(data, "data == null"); Preconditions.checkState(data.getContainerID() >= 0, "Container Id " + "cannot be negative"); + BlockUtils.verifyStorageType(container.getContainerData(), data.getBlockID()); KeyValueContainerData containerData = container.getContainerData(); diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/helpers/TestBlockData.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/helpers/TestBlockData.java index 7b6813a68260..208999d32320 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/helpers/TestBlockData.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/helpers/TestBlockData.java @@ -131,7 +131,7 @@ public void testToString() { final BlockID blockID = new BlockID(5, 123); blockID.setBlockCommitSequenceId(42); final BlockData subject = new BlockData(blockID); - assertEquals("[blockId=conID: 5 locID: 123 bcsId: 42 replicaIndex: null, size=0]", + assertEquals("[blockId=conID: 5 locID: 123 bcsId: 42 replicaIndex: null storageType: null, size=0]", subject.toString()); } } diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestContainerPersistence.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestContainerPersistence.java index c2c2f1b7e89b..bd810fb47247 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestContainerPersistence.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestContainerPersistence.java @@ -57,6 +57,7 @@ import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ChecksumType; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.Result; import org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.hdds.scm.container.common.helpers.StorageContainerException; @@ -137,10 +138,17 @@ private void initSchemaAndVersionInfo(ContainerTestVersionInfo versionInfo) { } @BeforeAll - public static void init() { + public static void init() throws IOException { conf = new OzoneConfiguration(); hddsPath = hddsFile.getPath(); - conf.set(ScmConfigKeys.HDDS_DATANODE_DIR_KEY, hddsPath); + StringBuilder hddsDirs = new StringBuilder(); + for (StorageType storageType : StorageType.values()) { + hddsDirs.append('[').append(storageType.name()).append(']'); + String volume = Files.createTempDirectory( + TestContainerPersistence.class.getSimpleName() + storageType.name()).toFile().getAbsolutePath(); + hddsDirs.append(volume).append(','); + } + conf.set(ScmConfigKeys.HDDS_DATANODE_DIR_KEY, hddsDirs.toString()); conf.set(OzoneConfigKeys.OZONE_METADATA_DIRS, hddsPath); volumeChoosingPolicy = new RoundRobinVolumeChoosingPolicy(); } @@ -194,9 +202,12 @@ private long getTestContainerID() { private KeyValueContainer addContainer(ContainerSet cSet, long cID) throws IOException { - long commitBytesBefore = 0; - long commitBytesAfter = 0; - long commitIncrement = 0; + return addContainer(cSet, cID, StorageType.DISK); + } + + private KeyValueContainer addContainer(ContainerSet cSet, long cID, + StorageType storageType) + throws IOException { KeyValueContainerData data = new KeyValueContainerData(cID, layout, ContainerTestHelper.CONTAINER_MAX_SIZE, UUID.randomUUID().toString(), @@ -204,17 +215,19 @@ private KeyValueContainer addContainer(ContainerSet cSet, long cID) data.addMetadata("VOLUME", "shire"); data.addMetadata("owner)", "bilbo"); KeyValueContainer container = new KeyValueContainer(data, conf); - commitBytesBefore = StorageVolumeUtil.getHddsVolumesList(volumeSet.getVolumesList()).get(0).getCommittedBytes(); + Map committedBytes = new HashMap<>(); + for (HddsVolume volume : StorageVolumeUtil.getHddsVolumesList( + volumeSet.getVolumesList())) { + committedBytes.put(volume, volume.getCommittedBytes()); + } - container.create(volumeSet, volumeChoosingPolicy, SCM_ID, StorageType.DISK); + container.create(volumeSet, volumeChoosingPolicy, SCM_ID, storageType); cSet.addContainer(container); - commitBytesAfter = container.getContainerData() - .getVolume().getCommittedBytes(); - commitIncrement = commitBytesAfter - commitBytesBefore; // did we commit space for the new container? - assertEquals(commitIncrement, - ContainerTestHelper.CONTAINER_MAX_SIZE); + HddsVolume volume = container.getContainerData().getVolume(); + assertEquals(ContainerTestHelper.CONTAINER_MAX_SIZE, + volume.getCommittedBytes() - committedBytes.get(volume)); return container; } @@ -423,7 +436,7 @@ public void testDeleteContainerWithRenaming( // Set up container 1. long testContainerID1 = getTestContainerID(); Container container1 = - addContainer(containerSet, testContainerID1); + addContainer(containerSet, testContainerID1, StorageType.DISK); BlockID container1Block = addBlockToContainer(container1); container1.close(); KeyValueContainerData container1Data = container1.getContainerData(); @@ -432,7 +445,7 @@ public void testDeleteContainerWithRenaming( // Set up container 2. long testContainerID2 = getTestContainerID(); Container container2 = - addContainer(containerSet, testContainerID2); + addContainer(containerSet, testContainerID2, StorageType.DISK); BlockID container2Block = addBlockToContainer(container2); container2.close(); KeyValueContainerData container2Data = container2.getContainerData(); @@ -909,6 +922,66 @@ public void testPutBlock(ContainerTestVersionInfo versionInfo) assertEquals(info.getChecksumData(), readChunk.getChecksumData()); } + /** + * Tests a put block and read block with StorageType. + * @throws IOException + */ + @ContainerTestVersionInfo.ContainerTest + public void testPutBlockWithStorageType(ContainerTestVersionInfo versionInfo) + throws IOException { + initSchemaAndVersionInfo(versionInfo); + for (StorageType storageType : StorageType.values()) { + long testContainerID = getTestContainerID(); + Container container = addContainer(containerSet, testContainerID, storageType); + + BlockID blockID = ContainerTestHelper.getTestBlockID(testContainerID, 0, storageType); + ChunkInfo info = writeChunkHelper(blockID); + BlockData blockData = new BlockData(blockID); + List chunkList = new LinkedList<>(); + chunkList.add(info.getProtoBufMessage()); + blockData.setChunks(chunkList); + blockManager.putBlock(container, blockData); + BlockData readBlockData = blockManager.getBlock(container, blockData.getBlockID()); + ChunkInfo readChunk = ChunkInfo.getFromProtoBuf(readBlockData.getChunks().get(0)); + assertEquals(info.getChecksumData(), readChunk.getChecksumData()); + assertEquals(storageType, readBlockData.getBlockID().getStorageType()); + } + } + + /** + * Tests a put block and read block with StorageType. + * @throws IOException + */ + @ContainerTestVersionInfo.ContainerTest + public void testInvalidPutBlockWithStorageType( + ContainerTestVersionInfo versionInfo) throws IOException { + initSchemaAndVersionInfo(versionInfo); + // If StorageType in the request is not null, + // the PutBlock must put to a Container of the same StorageType. + long testContainerID1 = getTestContainerID(); + Container container1 = addContainer(containerSet, testContainerID1, + StorageType.DISK); + BlockData blockData1 = getBlockData(testContainerID1, StorageType.DISK); + blockManager.putBlock(container1, blockData1); + + long testContainerID2 = getTestContainerID(); + Container container2 = addContainer(containerSet, testContainerID2, StorageType.DISK); + BlockData blockData2 = getBlockData(testContainerID2, StorageType.SSD); + StorageContainerException storageContainerException = + assertThrows(StorageContainerException.class, () -> blockManager.putBlock(container2, blockData2)); + assertEquals(Result.INVALID_ARGUMENT, storageContainerException.getResult()); + } + + private BlockData getBlockData(long containerID, StorageType storageType) throws IOException { + BlockID blockID = ContainerTestHelper.getTestBlockID(containerID, 0, storageType); + ChunkInfo info = writeChunkHelper(blockID); + BlockData blockData = new BlockData(blockID); + List chunkList = new LinkedList<>(); + chunkList.add(info.getProtoBufMessage()); + blockData.setChunks(chunkList); + return blockData; + } + /** * Tests a put block and read block with invalid bcsId. * diff --git a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java index 3d1bb301afda..2a64414a7515 100644 --- a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java +++ b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/impl/TestHddsDispatcher.java @@ -58,6 +58,7 @@ import org.apache.hadoop.conf.StorageUnit; import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.fs.MockSpaceUsageCheckFactory; import org.apache.hadoop.hdds.fs.MockSpaceUsageSource; @@ -549,9 +550,71 @@ public void testCreateContainerRejectsInvalidStorageType() throws IOException { .setStorageTypeID(999))) .build(); + ContainerCommandResponseProto response = + dispatcher.createContainer(requestWithInvalidStorageType); + assertEquals(ContainerProtos.Result.INVALID_ARGUMENT, response.getResult()); + } + + @Test + public void testInvalidStorageTypeDoesNotMarkContainerUnhealthy() throws IOException { + File diskVolume = Files.createTempDirectory(tempDir, "disk").toFile(); + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(ScmConfigKeys.HDDS_DATANODE_DIR_KEY, diskVolume.getAbsolutePath()); + DatanodeDetails dd = randomDatanodeDetails(); + HddsDispatcher dispatcher = createDispatcher(dd, UUID.randomUUID(), conf); + + ContainerCommandRequestProto validWriteChunkRequest = + getWriteChunkRequest(dd.getUuidString(), 1L, 1L, null); + assertEquals(ContainerProtos.Result.SUCCESS, + dispatcher.dispatch(validWriteChunkRequest, null).getResult()); + + ContainerCommandRequestProto writeChunkRequest = + getWriteChunkRequest(dd.getUuidString(), 1L, 2L, null); + WriteChunkRequestProto writeChunk = writeChunkRequest.getWriteChunk(); + ContainerCommandRequestProto requestWithInvalidStorageType = + writeChunkRequest.toBuilder() + .setWriteChunk(writeChunk.toBuilder() + .setBlockID(writeChunk.getBlockID().toBuilder() + .setStorageTypeID(999))) + .build(); + ContainerCommandResponseProto response = dispatcher.dispatch(requestWithInvalidStorageType, null); assertEquals(ContainerProtos.Result.INVALID_ARGUMENT, response.getResult()); + assertEquals(ContainerProtos.ContainerDataProto.State.OPEN, + dispatcher.getContainer(1L).getContainerData().getState()); + } + + @Test + public void testStorageTypeMismatchDoesNotMarkContainerUnhealthy() throws IOException { + File diskVolume = Files.createTempDirectory(tempDir, "disk").toFile(); + OzoneConfiguration conf = new OzoneConfiguration(); + conf.set(ScmConfigKeys.HDDS_DATANODE_DIR_KEY, diskVolume.getAbsolutePath()); + DatanodeDetails dd = randomDatanodeDetails(); + HddsDispatcher dispatcher = createDispatcher(dd, UUID.randomUUID(), conf); + + ContainerCommandRequestProto validWriteChunkRequest = + getWriteChunkRequest(dd.getUuidString(), 1L, 1L, null); + assertEquals(ContainerProtos.Result.SUCCESS, + dispatcher.dispatch(validWriteChunkRequest, null).getResult()); + + ContainerCommandResponseProto putBlockResponse = dispatcher.dispatch( + newPutBlock(1L, 2L, HddsProtos.StorageTypeProto.SSD), null); + assertEquals(ContainerProtos.Result.INVALID_ARGUMENT, putBlockResponse.getResult()); + assertEquals(ContainerProtos.ContainerDataProto.State.OPEN, + dispatcher.getContainer(1L).getContainerData().getState()); + + ContainerCommandResponseProto putSmallFileResponse = dispatcher.dispatch( + newPutSmallFile(1L, 3L, HddsProtos.StorageTypeProto.SSD), null); + assertEquals(ContainerProtos.Result.INVALID_ARGUMENT, putSmallFileResponse.getResult()); + assertEquals(ContainerProtos.ContainerDataProto.State.OPEN, + dispatcher.getContainer(1L).getContainerData().getState()); + + ContainerCommandResponseProto writeChunkResponse = dispatcher.dispatch( + getWriteChunkRequest(dd.getUuidString(), 1L, 4L, HddsProtos.StorageTypeProto.SSD), null); + assertEquals(ContainerProtos.Result.INVALID_ARGUMENT, writeChunkResponse.getResult()); + assertEquals(ContainerProtos.ContainerDataProto.State.OPEN, + dispatcher.getContainer(1L).getContainerData().getState()); } private void assertContainerDoNotExist(HddsDispatcher hddsDispatcher, @@ -870,7 +933,7 @@ private ContainerCommandRequestProto getWriteChunkRequest( new BlockID(containerId, localId).getDatanodeBlockIDProtobufBuilder(); // TODO: Pass the real storage type from the write path once that BlockID support StorageType if (storageType != null) { - blockID.setStorageTypeID(storageType.getNumber()); + blockID.setStorageTypeID(StorageTypeUtils.getIDFromProtobuf(storageType)); } WriteChunkRequestProto.Builder writeChunkRequest = WriteChunkRequestProto