This is an automated email from the ASF dual-hosted git repository. devmadhuu pushed a commit to branch HDDS-11233 in repository https://gitbox.apache.org/repos/asf/ozone.git
commit 8f590d8a2687d83099b1e8fd9613ee24c6bf58a7 Author: Devesh Kumar Singh <[email protected]> AuthorDate: Mon Jul 27 11:10:04 2026 +0530 HDDS-15249. SCM Pipeline Support Create And Get With StorageTier (#10815) --- .../org/apache/hadoop/ozone/OzoneConfigKeys.java | 6 ++ .../common/src/main/resources/ozone-default.xml | 14 +++++ .../hadoop/hdds/scm/block/BlockManagerImpl.java | 14 ++++- .../hdds/scm/container/ContainerManagerImpl.java | 4 +- .../ha/invoker/PipelineStateManagerInvoker.java | 67 ++++++++++++++++------ .../scm/pipeline/BackgroundPipelineCreator.java | 5 +- .../hadoop/hdds/scm/pipeline/PipelineFactory.java | 5 +- .../hadoop/hdds/scm/pipeline/PipelineManager.java | 15 ++++- .../hdds/scm/pipeline/PipelineManagerImpl.java | 26 ++++++--- .../hdds/scm/pipeline/PipelineStateManager.java | 15 ++++- .../scm/pipeline/PipelineStateManagerImpl.java | 29 +++++++++- .../hadoop/hdds/scm/pipeline/PipelineStateMap.java | 63 ++++++++++++++++++-- .../hdds/scm/pipeline/RatisPipelineProvider.java | 14 +++-- .../scm/pipeline/WritableContainerFactory.java | 14 +++-- .../scm/pipeline/WritableContainerProvider.java | 6 +- .../scm/pipeline/WritableECContainerProvider.java | 9 ++- .../pipeline/WritableRatisContainerProvider.java | 48 ++++++++++------ .../hdds/scm/server/SCMClientProtocolServer.java | 5 +- .../hadoop/hdds/scm/block/TestBlockManager.java | 17 +++--- .../scm/container/TestContainerManagerImpl.java | 7 ++- .../scm/container/TestContainerReportHandler.java | 7 ++- .../scm/container/TestContainerStateManager.java | 3 +- .../TestIncrementalContainerReportHandler.java | 3 +- .../hdds/scm/node/TestContainerPlacement.java | 3 +- .../hadoop/hdds/scm/node/TestSCMNodeManager.java | 5 +- .../hdds/scm/pipeline/MockPipelineManager.java | 20 +++++-- .../hdds/scm/pipeline/TestPipelineManagerImpl.java | 50 +++++++++------- .../scm/pipeline/TestRatisPipelineProvider.java | 19 ++++-- .../pipeline/TestWritableECContainerProvider.java | 59 ++++++++++--------- .../TestWritableRatisContainerProvider.java | 35 +++++++---- .../safemode/TestHealthyPipelineSafeModeRule.java | 21 +++---- .../TestOneReplicaPipelineSafeModeRule.java | 3 +- .../hdds/scm/safemode/TestSCMSafeModeManager.java | 9 +-- .../hdds/scm/cli/ContainerOperationClient.java | 1 + .../hadoop/ozone/recon/TestReconAsPassiveScm.java | 4 +- .../org/apache/hadoop/fs/ozone/TestSafeMode.java | 4 +- .../hdds/scm/pipeline/TestMultiRaftSetup.java | 5 +- .../TestRatisPipelineCreateAndDestroy.java | 3 +- .../hdds/scm/storage/TestContainerCommandsEC.java | 5 +- .../hadoop/hdds/upgrade/TestHDDSUpgrade.java | 3 +- 40 files changed, 456 insertions(+), 189 deletions(-) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneConfigKeys.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneConfigKeys.java index 87e6fb86f1e..c2662d91821 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneConfigKeys.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/ozone/OzoneConfigKeys.java @@ -22,6 +22,7 @@ import org.apache.hadoop.hdds.annotation.InterfaceStability; import org.apache.hadoop.hdds.client.ReplicationFactor; import org.apache.hadoop.hdds.client.ReplicationType; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.scm.ScmConfigKeys; import org.apache.hadoop.http.HttpConfig; import org.apache.ratis.util.TimeDuration; @@ -728,6 +729,11 @@ public final class OzoneConfigKeys { "ozone.client.elastic.byte.buffer.pool.max.size"; public static final String OZONE_CLIENT_ELASTIC_BYTE_BUFFER_POOL_MAX_SIZE_DEFAULT = "16GB"; + public static final String OZONE_DEFAULT_STORAGE_TIER_KEY = + "ozone.default.storageTier"; + public static final String OZONE_DEFAULT_STORAGE_TIER_DEFAULT = + StorageTier.DISK.toString(); + /** * There is no need to instantiate this class. */ diff --git a/hadoop-hdds/common/src/main/resources/ozone-default.xml b/hadoop-hdds/common/src/main/resources/ozone-default.xml index 3e51a6d9d04..1ce19f56a4c 100644 --- a/hadoop-hdds/common/src/main/resources/ozone-default.xml +++ b/hadoop-hdds/common/src/main/resources/ozone-default.xml @@ -4434,6 +4434,20 @@ </description> </property> + <property> + <name>ozone.default.storageTier</name> + <value>DISK</value> + <tag>OZONE, MANAGEMENT</tag> + <description> + Default StorageTier used for pipeline creation and block allocation + when a client does not specify a StorageTier (StoragePolicy). + For old-version clients that do not carry a StoragePolicy the + StorageTier is null; this value is used as the fallback so their + write path is unchanged. Supported values are DISK, SSD, ARCHIVE, + RAM_DISK. + </description> + </property> + <property> <name>ozone.client.ec.grpc.retries.enabled</name> <value>true</value> diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/block/BlockManagerImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/block/BlockManagerImpl.java index 108cf6786a9..f5b660f3f3b 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/block/BlockManagerImpl.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/block/BlockManagerImpl.java @@ -29,6 +29,7 @@ import org.apache.hadoop.hdds.client.BlockID; import org.apache.hadoop.hdds.client.ContainerBlockID; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.conf.StorageUnit; import org.apache.hadoop.hdds.scm.ScmConfig; @@ -160,8 +161,19 @@ public AllocatedBlock allocateBlock(final long size, INVALID_BLOCK_SIZE); } + // TODO: Implement pass storageTier(StoragePolicy) from API. + // For the old version client, it will not have a "default StoragePolicy", + // so its StorageTier will be null, we use the "default StorageTier" to + // write data for them. The value of the "default StorageTier" can be set + // via configuration, and the default is StorageTier.DISK. + // + // By default, if the Datanode Volume StorageType is not explicitly + // configured, it will be of type StorageType.DISK and therefore belong + // to a StorageTier.DISK tier, so for old clients the write process is + // unchanged if the Datanode Volume configuration is not changed. + StorageTier storageTier = StorageTier.getDefaultTier(); ContainerInfo containerInfo = writableContainerFactory.getContainer( - size, replicationConfig, owner, excludeList); + size, replicationConfig, owner, excludeList, storageTier); if (containerInfo != null) { return newBlock(containerInfo); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManagerImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManagerImpl.java index 3dc87a40ef1..ad965640376 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManagerImpl.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/ContainerManagerImpl.java @@ -30,6 +30,7 @@ import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ContainerInfoProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleEvent; @@ -192,7 +193,8 @@ public ContainerInfo allocateContainer( if (pipelines.isEmpty()) { try { - pipeline = pipelineManager.createPipeline(replicationConfig); + pipeline = pipelineManager.createPipeline(replicationConfig, + StorageTier.getDefaultTier()); if (replicationConfig.getReplicationType() == HddsProtos.ReplicationType.EC) { pipelineManager.openPipeline(pipeline.getId()); } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/invoker/PipelineStateManagerInvoker.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/invoker/PipelineStateManagerInvoker.java index 4b756078186..3a3ea16e0ee 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/invoker/PipelineStateManagerInvoker.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/ha/invoker/PipelineStateManagerInvoker.java @@ -22,6 +22,7 @@ import java.util.List; import java.util.NavigableSet; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.ha.SCMRatisResponse; @@ -135,10 +136,20 @@ public List<Pipeline> getPipelines(ReplicationConfig arg0, Pipeline.PipelineStat return invoker.getImpl().getPipelines(arg0, arg1); } + @Override + public List<Pipeline> getPipelines(ReplicationConfig arg0, StorageTier arg1) { + return invoker.getImpl().getPipelines(arg0, arg1); + } + + @Override + public List<Pipeline> getPipelines(ReplicationConfig arg0, Pipeline.PipelineState arg1, StorageTier arg2) { + return invoker.getImpl().getPipelines(arg0, arg1, arg2); + } + @Override public List<Pipeline> getPipelines(ReplicationConfig arg0, Pipeline.PipelineState arg1, Collection arg2, - Collection arg3) { - return invoker.getImpl().getPipelines(arg0, arg1, arg2, arg3); + Collection arg3, StorageTier arg4) { + return invoker.getImpl().getPipelines(arg0, arg1, arg2, arg3, arg4); } @Override @@ -238,39 +249,57 @@ public Message invokeLocal(String methodName, Object[] p) throws Exception { returnValue = getImpl().getPipelines(arg11, arg12); break; } - if (p.length == 4 && (p[0] == null || ReplicationConfig.class.isInstance(p[0])) && (p[1] == null || - Pipeline.PipelineState.class.isInstance(p[1])) && (p[2] == null || Collection.class.isInstance(p[2])) && (p[3] - == null || Collection.class.isInstance(p[3]))) { + if (p.length == 2 && (p[0] == null || ReplicationConfig.class.isInstance(p[0])) && (p[1] == null || + StorageTier.class.isInstance(p[1]))) { final ReplicationConfig arg13 = (ReplicationConfig) p[0]; - final Pipeline.PipelineState arg14 = (Pipeline.PipelineState) p[1]; - final Collection arg15 = (Collection) p[2]; - final Collection arg16 = (Collection) p[3]; + final StorageTier arg14 = (StorageTier) p[1]; + returnType = List.class; + returnValue = getImpl().getPipelines(arg13, arg14); + break; + } + if (p.length == 3 && (p[0] == null || ReplicationConfig.class.isInstance(p[0])) && (p[1] == null || + Pipeline.PipelineState.class.isInstance(p[1])) && (p[2] == null || StorageTier.class.isInstance(p[2]))) { + final ReplicationConfig arg15 = (ReplicationConfig) p[0]; + final Pipeline.PipelineState arg16 = (Pipeline.PipelineState) p[1]; + final StorageTier arg17 = (StorageTier) p[2]; + returnType = List.class; + returnValue = getImpl().getPipelines(arg15, arg16, arg17); + break; + } + if (p.length == 5 && (p[0] == null || ReplicationConfig.class.isInstance(p[0])) && (p[1] == null || + Pipeline.PipelineState.class.isInstance(p[1])) && (p[2] == null || Collection.class.isInstance(p[2])) && (p[3] + == null || Collection.class.isInstance(p[3])) && (p[4] == null || StorageTier.class.isInstance(p[4]))) { + final ReplicationConfig arg18 = (ReplicationConfig) p[0]; + final Pipeline.PipelineState arg19 = (Pipeline.PipelineState) p[1]; + final Collection arg20 = (Collection) p[2]; + final Collection arg21 = (Collection) p[3]; + final StorageTier arg22 = (StorageTier) p[4]; returnType = List.class; - returnValue = getImpl().getPipelines(arg13, arg14, arg15, arg16); + returnValue = getImpl().getPipelines(arg18, arg19, arg20, arg21, arg22); break; } throw new IllegalArgumentException("Method not found: " + methodName + " in PipelineStateManager"); case "reinitialize": - final Table arg17 = p.length > 0 ? (Table) p[0] : null; - getImpl().reinitialize(arg17); + final Table arg23 = p.length > 0 ? (Table) p[0] : null; + getImpl().reinitialize(arg23); return Message.EMPTY; case "removeContainerFromPipeline": - final PipelineID arg18 = p.length > 0 ? (PipelineID) p[0] : null; - final ContainerID arg19 = p.length > 1 ? (ContainerID) p[1] : null; - getImpl().removeContainerFromPipeline(arg18, arg19); + final PipelineID arg24 = p.length > 0 ? (PipelineID) p[0] : null; + final ContainerID arg25 = p.length > 1 ? (ContainerID) p[1] : null; + getImpl().removeContainerFromPipeline(arg24, arg25); return Message.EMPTY; case "removePipeline": - final HddsProtos.PipelineID arg20 = p.length > 0 ? (HddsProtos.PipelineID) p[0] : null; - getImpl().removePipeline(arg20); + final HddsProtos.PipelineID arg26 = p.length > 0 ? (HddsProtos.PipelineID) p[0] : null; + getImpl().removePipeline(arg26); return Message.EMPTY; case "updatePipelineState": - final HddsProtos.PipelineID arg21 = p.length > 0 ? (HddsProtos.PipelineID) p[0] : null; - final HddsProtos.PipelineState arg22 = p.length > 1 ? (HddsProtos.PipelineState) p[1] : null; - getImpl().updatePipelineState(arg21, arg22); + final HddsProtos.PipelineID arg27 = p.length > 0 ? (HddsProtos.PipelineID) p[0] : null; + final HddsProtos.PipelineState arg28 = p.length > 1 ? (HddsProtos.PipelineState) p[1] : null; + getImpl().updatePipelineState(arg27, arg28); return Message.EMPTY; default: diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/BackgroundPipelineCreator.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/BackgroundPipelineCreator.java index 97455d56d51..f8e31149d72 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/BackgroundPipelineCreator.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/BackgroundPipelineCreator.java @@ -40,6 +40,7 @@ import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; @@ -223,7 +224,9 @@ private void createPipelines() throws RuntimeException { ReplicationConfig replicationConfig = (ReplicationConfig) it.next(); try { - Pipeline pipeline = pipelineManager.createPipeline(replicationConfig); + // Only create default StorageTier Pipeline + Pipeline pipeline = pipelineManager.createPipeline(replicationConfig, + StorageTier.getDefaultTier()); LOG.info("Created new pipeline {}", pipeline); } catch (IOException ioe) { it.remove(); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineFactory.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineFactory.java index 2268309012b..672798c6783 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineFactory.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineFactory.java @@ -81,11 +81,10 @@ void setProvider( public Pipeline create( ReplicationConfig replicationConfig, List<DatanodeDetails> excludedNodes, - List<DatanodeDetails> favoredNodes) + List<DatanodeDetails> favoredNodes, StorageTier storageTier) throws IOException { - // TODO StoragePolicy replace this StorageTier with the passed StorageTier Pipeline pipeline = providers.get(replicationConfig.getReplicationType()) - .create(replicationConfig, excludedNodes, favoredNodes, StorageTier.getDefaultTier()); + .create(replicationConfig, excludedNodes, favoredNodes, storageTier); checkPipeline(pipeline); return pipeline; } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java index 77c9fec4056..d0b9f4299c9 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManager.java @@ -37,12 +37,14 @@ */ public interface PipelineManager extends Closeable, PipelineManagerMXBean { - Pipeline createPipeline(ReplicationConfig replicationConfig) + Pipeline createPipeline(ReplicationConfig replicationConfig, + StorageTier storageTier) throws IOException; Pipeline createPipeline(ReplicationConfig replicationConfig, List<DatanodeDetails> excludedNodes, - List<DatanodeDetails> favoredNodes) + List<DatanodeDetails> favoredNodes, + StorageTier storageTier) throws IOException; Pipeline buildECPipeline(ReplicationConfig replicationConfig, @@ -75,11 +77,18 @@ List<Pipeline> getPipelines( ReplicationConfig replicationConfig, Pipeline.PipelineState state ); + List<Pipeline> getPipelines( + ReplicationConfig replicationConfig, + Pipeline.PipelineState state, + StorageTier storageTier + ); + List<Pipeline> getPipelines( ReplicationConfig replicationConfig, Pipeline.PipelineState state, Collection<DatanodeDetails> excludeDns, - Collection<PipelineID> excludePipelines + Collection<PipelineID> excludePipelines, + StorageTier storageTier ); /** diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java index 4dcfaad37a5..63c9027b347 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java @@ -211,9 +211,10 @@ public Pipeline buildECPipeline(ReplicationConfig replicationConfig, if (replicationConfig.getReplicationType() != ReplicationType.EC) { throw new IllegalArgumentException("Replication type must be EC"); } + // TODO StoragePolicy Support EC checkIfPipelineCreationIsAllowed(replicationConfig); return pipelineFactory.create(replicationConfig, excludedNodes, - favoredNodes); + favoredNodes, StorageTier.getDefaultTier()); } /** @@ -235,15 +236,17 @@ public void addEcPipeline(Pipeline pipeline) } @Override - public Pipeline createPipeline(ReplicationConfig replicationConfig) + public Pipeline createPipeline(ReplicationConfig replicationConfig, + StorageTier storageTier) throws IOException { return createPipeline(replicationConfig, Collections.emptyList(), - Collections.emptyList()); + Collections.emptyList(), storageTier); } @Override public Pipeline createPipeline(ReplicationConfig replicationConfig, - List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes) + List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes, + StorageTier storageTier) throws IOException { checkIfPipelineCreationIsAllowed(replicationConfig); @@ -252,8 +255,10 @@ public Pipeline createPipeline(ReplicationConfig replicationConfig, try { try { pipeline = pipelineFactory.create(replicationConfig, - excludedNodes, favoredNodes); + excludedNodes, favoredNodes, storageTier); } catch (IOException e) { + LOG.debug("Failed to create pipeline with replicationConfig {} " + + "storageTier {}.", replicationConfig, storageTier, e); metrics.incNumPipelineCreationFailed(); throw e; } @@ -361,13 +366,20 @@ public List<Pipeline> getPipelines(ReplicationConfig config, return stateManager.getPipelines(config, state); } + @Override + public List<Pipeline> getPipelines(ReplicationConfig config, + Pipeline.PipelineState state, StorageTier storageTier) { + return stateManager.getPipelines(config, state, storageTier); + } + @Override public List<Pipeline> getPipelines( ReplicationConfig replicationConfig, Pipeline.PipelineState state, Collection<DatanodeDetails> excludeDns, - Collection<PipelineID> excludePipelines) { + Collection<PipelineID> excludePipelines, StorageTier storageTier) { return stateManager - .getPipelines(replicationConfig, state, excludeDns, excludePipelines); + .getPipelines(replicationConfig, state, excludeDns, excludePipelines, + storageTier); } /** diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManager.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManager.java index 17f7345b0c4..6db53ba369d 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManager.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManager.java @@ -22,6 +22,7 @@ import java.util.List; import java.util.NavigableSet; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.SCMRatisProtocol.RequestType; @@ -87,11 +88,23 @@ List<Pipeline> getPipelines( Pipeline.PipelineState state ); + List<Pipeline> getPipelines( + ReplicationConfig replicationConfig, + StorageTier storageTier + ); + + List<Pipeline> getPipelines( + ReplicationConfig replicationConfig, + Pipeline.PipelineState state, + StorageTier storageTier + ); + List<Pipeline> getPipelines( ReplicationConfig replicationConfig, Pipeline.PipelineState state, Collection<DatanodeDetails> excludeDns, - Collection<PipelineID> excludePipelines + Collection<PipelineID> excludePipelines, + StorageTier storageTier ); int getPipelineCount( diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManagerImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManagerImpl.java index c23c424949c..2607f76d543 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManagerImpl.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateManagerImpl.java @@ -25,6 +25,7 @@ import java.util.concurrent.locks.ReadWriteLock; import java.util.concurrent.locks.ReentrantReadWriteLock; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.container.ContainerID; @@ -167,15 +168,39 @@ public List<Pipeline> getPipelines( } } + @Override + public List<Pipeline> getPipelines( + ReplicationConfig replicationConfig, StorageTier storageTier) { + lock.readLock().lock(); + try { + return pipelineStateMap.getPipelines(replicationConfig, storageTier); + } finally { + lock.readLock().unlock(); + } + } + + @Override + public List<Pipeline> getPipelines( + ReplicationConfig replicationConfig, + Pipeline.PipelineState state, StorageTier storageTier) { + lock.readLock().lock(); + try { + return pipelineStateMap.getPipelines(replicationConfig, state, storageTier); + } finally { + lock.readLock().unlock(); + } + } + @Override public List<Pipeline> getPipelines( ReplicationConfig replicationConfig, Pipeline.PipelineState state, Collection<DatanodeDetails> excludeDns, - Collection<PipelineID> excludePipelines) { + Collection<PipelineID> excludePipelines, StorageTier storageTier) { lock.readLock().lock(); try { return pipelineStateMap - .getPipelines(replicationConfig, state, excludeDns, excludePipelines); + .getPipelines(replicationConfig, state, excludeDns, excludePipelines, + storageTier); } finally { lock.readLock().unlock(); } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java index 2aab03b00b3..82a34c54e16 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineStateMap.java @@ -31,7 +31,9 @@ import java.util.Objects; import java.util.TreeSet; import java.util.function.Predicate; +import java.util.stream.Collectors; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.scm.container.ContainerID; import org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState; @@ -193,6 +195,55 @@ private List<Pipeline> getOpenPipelines(ReplicationConfig replicationConfig) { return pipelines != null && !pipelines.isEmpty() ? new ArrayList<>(pipelines) : Collections.emptyList(); } + /** + * Get list of pipelines corresponding to specified replication type + * and storageTier. + * + * @param replicationConfig - ReplicationConfig + * @param storageTier - Required storageTier + * @return List of pipelines with specified replication type and storageTier + */ + List<Pipeline> getPipelines(ReplicationConfig replicationConfig, + StorageTier storageTier) { + Objects.requireNonNull(replicationConfig, "ReplicationConfig cannot be null"); + Objects.requireNonNull(storageTier, "Pipeline storageTier cannot be null"); + return getPipelines(replicationConfig).stream() + .filter(pipeline -> matchesStorageTier(pipeline, storageTier)) + .collect(Collectors.toList()); + } + + /** + * Get list of pipeline corresponding to specified replication type, + * replication factor, pipeline state and storageTier. + * + * @param replicationConfig - ReplicationConfig + * @param state - Required PipelineState + * @param storageTier - Required storageTier + * @return List of pipelines with specified replication type, + * replication factor, pipeline state and storageTier + */ + List<Pipeline> getPipelines(ReplicationConfig replicationConfig, + PipelineState state, StorageTier storageTier) { + Objects.requireNonNull(replicationConfig, "ReplicationConfig cannot be null"); + Objects.requireNonNull(state, "Pipeline state cannot be null"); + Objects.requireNonNull(storageTier, "Pipeline storageTier cannot be null"); + return getPipelines(replicationConfig, state).stream() + .filter(pipeline -> matchesStorageTier(pipeline, storageTier)) + .collect(Collectors.toList()); + } + + /** + * A pipeline with no supportedStorageTier (e.g. one created by an older + * SCM before this field was introduced, or restored from a legacy DB + * snapshot) is treated as matching any tier. This preserves backward + * compatibility during rolling upgrade: legacy pipelines remain usable + * until they are scrubbed and replaced by tier-tagged ones. + */ + static boolean matchesStorageTier(Pipeline pipeline, StorageTier storageTier) { + final StorageTier pipelineTier = pipeline.getSupportedStorageTier(); + return pipelineTier == null || Objects.equals(pipelineTier, storageTier); + } + /** * Get a count of pipelines with the given replicationConfig and state. * This method is most efficient when getting a count for OPEN pipeline @@ -261,21 +312,25 @@ static Predicate<Pipeline> getPredicate( * @param state - Required PipelineState * @param excludeDns dns to exclude * @param excludePipelines pipelines to exclude + * @param storageTier - Required storageTier * @return List of pipelines with specified replication type, * replication factor and pipeline state */ List<Pipeline> getPipelines(ReplicationConfig replicationConfig, PipelineState state, Collection<DatanodeDetails> excludeDns, - Collection<PipelineID> excludePipelines) { + Collection<PipelineID> excludePipelines, StorageTier storageTier) { Objects.requireNonNull(replicationConfig, "ReplicationConfig cannot be null"); Objects.requireNonNull(state, "Pipeline state cannot be null"); Objects.requireNonNull(excludeDns, "Datanode exclude list cannot be null"); Objects.requireNonNull(excludePipelines, "Pipeline exclude list cannot be null"); + Objects.requireNonNull(storageTier, "Pipeline storageTier cannot be null"); if (state == PipelineState.OPEN) { final List<Pipeline> pipelines = getOpenPipelines(replicationConfig); if (excludeDns.isEmpty() && excludePipelines.isEmpty()) { - return pipelines; + return pipelines.stream() + .filter(pipeline -> matchesStorageTier(pipeline, storageTier)) + .collect(Collectors.toList()); } final Predicate<Pipeline> include = getPredicate(excludeDns, excludePipelines); @@ -287,10 +342,10 @@ List<Pipeline> getPipelines(ReplicationConfig replicationConfig, final List<Pipeline> pipelines = new ArrayList<>(pipelineMap.size() / 2 + 1); // only resize once for (PipelineInfo info : pipelineMap.values()) { final Pipeline pipeline = info.getPipeline(); - if (include.test(pipeline)) { + if (include.test(pipeline) && matchesStorageTier(pipeline, storageTier)) { pipelines.add(pipeline); } - } + } return pipelines; } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java index 8d558b5eee5..8e713cab86b 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/RatisPipelineProvider.java @@ -109,15 +109,18 @@ public RatisPipelineProvider(NodeManager nodeManager, } } - private boolean exceedPipelineNumberLimit(RatisReplicationConfig replicationConfig) { + private boolean exceedPipelineNumberLimit( + RatisReplicationConfig replicationConfig, StorageTier storageTier) { // Apply limits only for replication factor THREE if (replicationConfig.getReplicationFactor() != ReplicationFactor.THREE) { return false; } PipelineStateManager pipelineStateManager = getPipelineStateManager(); - int totalActivePipelines = pipelineStateManager.getPipelines(replicationConfig).size(); - int closedPipelines = pipelineStateManager.getPipelines(replicationConfig, PipelineState.CLOSED).size(); + int totalActivePipelines = pipelineStateManager + .getPipelines(replicationConfig, storageTier).size(); + int closedPipelines = pipelineStateManager + .getPipelines(replicationConfig, PipelineState.CLOSED, storageTier).size(); int openPipelines = totalActivePipelines - closedPipelines; // Check per-datanode pipeline limit if (datanodePipelineLimit > 0) { @@ -130,7 +133,8 @@ private boolean exceedPipelineNumberLimit(RatisReplicationConfig replicationConf // Check global pipeline limit if (pipelineNumberLimit > 0) { int factorOnePipelineCount = pipelineStateManager - .getPipelines(RatisReplicationConfig.getInstance(ReplicationFactor.ONE)).size(); + .getPipelines(RatisReplicationConfig.getInstance(ReplicationFactor.ONE), + storageTier).size(); int allowedOpenPipelines = pipelineNumberLimit - factorOnePipelineCount; return openPipelines >= allowedOpenPipelines; } @@ -155,7 +159,7 @@ public synchronized Pipeline create(RatisReplicationConfig replicationConfig, public synchronized Pipeline create(RatisReplicationConfig replicationConfig, List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes, StorageTier storageTier) throws IOException { - if (exceedPipelineNumberLimit(replicationConfig)) { + if (exceedPipelineNumberLimit(replicationConfig, storageTier)) { String limitInfo = (datanodePipelineLimit > 0) ? String.format("per datanode: %d", datanodePipelineLimit) : String.format(": %d", pipelineNumberLimit); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableContainerFactory.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableContainerFactory.java index b816bc4de7f..85c0cecb40f 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableContainerFactory.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableContainerFactory.java @@ -21,9 +21,11 @@ import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE; import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_CONTAINER_SIZE_DEFAULT; +import jakarta.annotation.Nonnull; import java.io.IOException; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.ConfigurationSource; import org.apache.hadoop.hdds.scm.container.ContainerInfo; import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; @@ -62,17 +64,19 @@ public WritableContainerFactory(StorageContainerManager scm) { } public ContainerInfo getContainer(final long size, - ReplicationConfig repConfig, String owner, ExcludeList excludeList) + ReplicationConfig repConfig, String owner, ExcludeList excludeList, + @Nonnull StorageTier storageTier) throws IOException { switch (repConfig.getReplicationType()) { case STAND_ALONE: return standaloneProvider - .getContainer(size, repConfig, owner, excludeList); + .getContainer(size, repConfig, owner, excludeList, storageTier); case RATIS: - return ratisProvider.getContainer(size, repConfig, owner, excludeList); + return ratisProvider.getContainer(size, repConfig, owner, excludeList, + storageTier); case EC: - return ecProvider.getContainer(size, (ECReplicationConfig)repConfig, - owner, excludeList); + return ecProvider.getContainer(size, (ECReplicationConfig) repConfig, + owner, excludeList, storageTier); default: throw new IOException(repConfig.getReplicationType() + " is an invalid replication type"); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableContainerProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableContainerProvider.java index 8eb3d1233f2..8c674377d41 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableContainerProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableContainerProvider.java @@ -17,8 +17,10 @@ package org.apache.hadoop.hdds.scm.pipeline; +import jakarta.annotation.Nonnull; import java.io.IOException; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.scm.container.ContainerInfo; import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; @@ -45,12 +47,14 @@ public interface WritableContainerProvider<T extends ReplicationConfig> { * @param owner The owner of the container * @param excludeList A set of datanodes, container and pipelines which should * not be considered. + * @param storageTier The storageTier of the container * @return A ContainerInfo which is open and has the capacity to store the * desired block size. * @throws IOException */ ContainerInfo getContainer(long size, T repConfig, - String owner, ExcludeList excludeList) + String owner, ExcludeList excludeList, + @Nonnull StorageTier storageTier) throws IOException; } diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableECContainerProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableECContainerProvider.java index b5faf14ea01..8847a950ac3 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableECContainerProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableECContainerProvider.java @@ -20,6 +20,7 @@ import static org.apache.hadoop.hdds.conf.ConfigTag.SCM; import static org.apache.hadoop.hdds.scm.node.NodeStatus.inServiceHealthy; +import jakarta.annotation.Nonnull; import java.io.IOException; import java.util.ArrayList; import java.util.Collections; @@ -27,6 +28,7 @@ import java.util.NavigableSet; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.Config; import org.apache.hadoop.hdds.conf.ConfigGroup; import org.apache.hadoop.hdds.conf.ConfigTag; @@ -91,8 +93,10 @@ public WritableECContainerProvider(WritableECContainerProviderConfig config, */ @Override public ContainerInfo getContainer(final long size, - ECReplicationConfig repConfig, String owner, ExcludeList excludeList) + ECReplicationConfig repConfig, String owner, ExcludeList excludeList, + @Nonnull StorageTier storageTier) throws IOException { + // TODO StoragePolicy Support EC int maximumPipelines = getMaximumPipelines(repConfig); int openPipelineCount; synchronized (this) { @@ -201,8 +205,9 @@ private ContainerInfo allocateContainer(ReplicationConfig repConfig, excludedNodes = new ArrayList<>(excludeList.getDatanodes()); } + // TODO StoragePolicy Support EC Pipeline newPipeline = pipelineManager.createPipeline(repConfig, - excludedNodes, Collections.emptyList()); + excludedNodes, Collections.emptyList(), StorageTier.getDefaultTier()); // the returned ContainerInfo should not be null (due to not enough space in the Datanodes specifically) because // this is a new pipeline and pipeline creation checks for sufficient space in the Datanodes ContainerInfo container = diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableRatisContainerProvider.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableRatisContainerProvider.java index a61b3289235..edfddfc2b81 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableRatisContainerProvider.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/WritableRatisContainerProvider.java @@ -17,11 +17,13 @@ package org.apache.hadoop.hdds.scm.pipeline; +import jakarta.annotation.Nonnull; import jakarta.annotation.Nullable; import java.io.IOException; import java.util.List; import java.util.stream.Collectors; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.scm.PipelineChoosePolicy; import org.apache.hadoop.hdds.scm.PipelineRequestInformation; import org.apache.hadoop.hdds.scm.container.ContainerInfo; @@ -55,7 +57,8 @@ public WritableRatisContainerProvider( @Override public ContainerInfo getContainer(final long size, - ReplicationConfig repConfig, String owner, ExcludeList excludeList) + ReplicationConfig repConfig, String owner, ExcludeList excludeList, + @Nonnull StorageTier storageTier) throws IOException { /* Here is the high level logic. @@ -80,7 +83,7 @@ public ContainerInfo getContainer(final long size, PipelineRequestInformation.Builder.getBuilder().setSize(size).build(); ContainerInfo containerInfo = - getContainer(repConfig, owner, excludeList, req); + getContainer(repConfig, owner, excludeList, req, storageTier); if (containerInfo != null) { return containerInfo; } @@ -88,19 +91,19 @@ public ContainerInfo getContainer(final long size, try { // TODO: #CLUTIL Remove creation logic when all replication types // and factors are handled by pipeline creator - Pipeline pipeline = pipelineManager.createPipeline(repConfig); + Pipeline pipeline = pipelineManager.createPipeline(repConfig, storageTier); // wait until pipeline is ready pipelineManager.waitPipelineReady(pipeline.getId(), 0); } catch (SCMException se) { - LOG.warn("Pipeline creation failed for repConfig {} " + + LOG.warn("Pipeline creation failed for repConfig: {} storageTier: {} " + "Datanodes may be used up. Try to see if any pipeline is in " + "ALLOCATED state, and then will wait for it to be OPEN", - repConfig, se); + repConfig, storageTier, se); List<Pipeline> allocatedPipelines = findPipelinesByState(repConfig, excludeList, - Pipeline.PipelineState.ALLOCATED); + Pipeline.PipelineState.ALLOCATED, storageTier); if (!allocatedPipelines.isEmpty()) { List<PipelineID> allocatedPipelineIDs = allocatedPipelines.stream() @@ -119,14 +122,14 @@ public ContainerInfo getContainer(final long size, failureReason = se.getMessage(); } } catch (IOException e) { - LOG.warn("Pipeline creation failed for repConfig: {}. " - + "Retrying get pipelines call once.", repConfig, e); + LOG.warn("Pipeline creation failed for repConfig: {} storageTier: {}. " + + "Retrying get pipelines call once.", repConfig, storageTier, e); failureReason = e.getMessage(); } // If Exception occurred or successful creation of pipeline do one // final try to fetch pipelines. - containerInfo = getContainer(repConfig, owner, excludeList, req); + containerInfo = getContainer(repConfig, owner, excludeList, req, storageTier); if (containerInfo != null) { return containerInfo; } @@ -134,24 +137,27 @@ public ContainerInfo getContainer(final long size, // we have tried all strategies we know but somehow we are not able // to get a container for this block. Log that info and throw an exception. LOG.error( - "Unable to allocate a block for the size: {}, repConfig: {}", - size, repConfig); + "Unable to allocate a block for the size: {}, repConfig: {}, storageTier: {}.", + size, repConfig, storageTier); throw new IOException( "Unable to allocate a container to the block of size: " + size - + ", replicationConfig: " + repConfig + ". " + failureReason); + + ", replicationConfig: " + repConfig + + " for storageTier: " + storageTier + ". " + failureReason); } @Nullable private ContainerInfo getContainer(ReplicationConfig repConfig, String owner, - ExcludeList excludeList, PipelineRequestInformation req) { + ExcludeList excludeList, PipelineRequestInformation req, + @Nonnull StorageTier storageTier) { // Acquire pipeline manager lock, to avoid any updates to pipeline // while allocate container happens. This is to avoid scenario like // mentioned in HDDS-5655. pipelineManager.acquireReadLock(); try { List<Pipeline> availablePipelines = findPipelinesByState(repConfig, - excludeList, Pipeline.PipelineState.OPEN); - return selectContainer(availablePipelines, req, owner, excludeList); + excludeList, Pipeline.PipelineState.OPEN, storageTier); + return selectContainer(availablePipelines, req, owner, excludeList, + storageTier); } finally { pipelineManager.releaseReadLock(); } @@ -160,27 +166,31 @@ private ContainerInfo getContainer(ReplicationConfig repConfig, String owner, private List<Pipeline> findPipelinesByState( final ReplicationConfig repConfig, final ExcludeList excludeList, - final Pipeline.PipelineState pipelineState) { + final Pipeline.PipelineState pipelineState, + @Nonnull StorageTier storageTier) { List<Pipeline> pipelines = pipelineManager.getPipelines(repConfig, pipelineState, excludeList.getDatanodes(), - excludeList.getPipelineIds()); + excludeList.getPipelineIds(), storageTier); if (pipelines.isEmpty() && !excludeList.isEmpty()) { // if no pipelines can be found, try finding pipeline without // exclusion - pipelines = pipelineManager.getPipelines(repConfig, pipelineState); + pipelines = pipelineManager.getPipelines(repConfig, pipelineState, + storageTier); } return pipelines; } private @Nullable ContainerInfo selectContainer( List<Pipeline> availablePipelines, PipelineRequestInformation req, - String owner, ExcludeList excludeList) { + String owner, ExcludeList excludeList, + @Nonnull StorageTier storageTier) { while (!availablePipelines.isEmpty()) { Pipeline pipeline = pipelineChoosePolicy.choosePipeline( availablePipelines, req); // look for OPEN containers that match the criteria. + // TODO StoragePolicy replace this StorageType with container actual StorageType final ContainerInfo containerInfo = containerManager.getMatchingContainer( req.getSize(), owner, pipeline, excludeList.getContainerIds()); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java index acb3e64a3ca..a595160e8a6 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java @@ -52,6 +52,7 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.hadoop.fs.CommonConfigurationKeysPublic; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.conf.ReconfigurationHandler; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -860,8 +861,10 @@ public Pipeline createReplicationPipeline(HddsProtos.ReplicationType type, } try { getScm().checkAdminAccess(getRemoteUser(), false); + // TODO: Support Allocate Pipeline Command with StorageTier Pipeline result = scm.getPipelineManager().createPipeline( - ReplicationConfig.fromProtoTypeAndFactor(type, factor)); + ReplicationConfig.fromProtoTypeAndFactor(type, factor), + StorageTier.getDefaultTier()); AUDIT.logWriteSuccess(buildAuditMessageForSuccess( SCMAction.CREATE_PIPELINE, auditMap)); return result; diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/block/TestBlockManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/block/TestBlockManager.java index 45c947cb00a..6e65dfa127d 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/block/TestBlockManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/block/TestBlockManager.java @@ -44,6 +44,7 @@ import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; @@ -194,7 +195,7 @@ public void cleanup() throws Exception { @Test public void testAllocateBlock() throws Exception { - pipelineManager.createPipeline(replicationConfig); + pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); HddsTestUtils.openAllRatisPipelines(pipelineManager); AllocatedBlock block = blockManager.allocateBlock(DEFAULT_BLOCK_SIZE, replicationConfig, OzoneConsts.OZONE, new ExcludeList()); @@ -205,7 +206,7 @@ public void testAllocateBlock() throws Exception { public void testAllocateBlockWithExclusion() throws Exception { try { while (true) { - pipelineManager.createPipeline(replicationConfig); + pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); } } catch (IOException e) { } @@ -272,7 +273,7 @@ void testBlockDistribution() throws Exception { for (int i = 0; i < threadCount; i++) { executors.add(Executors.newSingleThreadExecutor()); } - pipelineManager.createPipeline(replicationConfig); + pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); HddsTestUtils.openAllRatisPipelines(pipelineManager); Map<Long, List<AllocatedBlock>> allocatedBlockMap = new ConcurrentHashMap<>(); @@ -328,7 +329,7 @@ void testBlockDistributionWithMultipleDisks() throws Exception { for (int i = 0; i < threadCount; i++) { executors.add(Executors.newSingleThreadExecutor()); } - pipelineManager.createPipeline(replicationConfig); + pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); HddsTestUtils.openAllRatisPipelines(pipelineManager); Map<Long, List<AllocatedBlock>> allocatedBlockMap = new ConcurrentHashMap<>(); @@ -388,7 +389,7 @@ void testBlockDistributionWithMultipleRaftLogDisks() throws Exception { for (int i = 0; i < threadCount; i++) { executors.add(Executors.newSingleThreadExecutor()); } - pipelineManager.createPipeline(replicationConfig); + pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); HddsTestUtils.openAllRatisPipelines(pipelineManager); Map<Long, List<AllocatedBlock>> allocatedBlockMap = new ConcurrentHashMap<>(); @@ -466,8 +467,8 @@ public void testAllocateBlockSucInSafeMode() throws Exception { public void testMultipleBlockAllocation() throws IOException, TimeoutException, InterruptedException { - pipelineManager.createPipeline(replicationConfig); - pipelineManager.createPipeline(replicationConfig); + pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); + pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); HddsTestUtils.openAllRatisPipelines(pipelineManager); AllocatedBlock allocatedBlock = blockManager @@ -513,7 +514,7 @@ public void testMultipleBlockAllocationWithClosedContainer() for (int i = 0; i < nodeManager.getNodes(NodeStatus.inServiceHealthy()).size() / replicationConfig.getRequiredNodes(); i++) { - pipelineManager.createPipeline(replicationConfig); + pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); } HddsTestUtils.openAllRatisPipelines(pipelineManager); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerManagerImpl.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerManagerImpl.java index a437a17ae7e..d5c5604d56a 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerManagerImpl.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerManagerImpl.java @@ -40,6 +40,7 @@ import java.util.concurrent.TimeoutException; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; @@ -106,7 +107,7 @@ void setUp() throws Exception { .checkSpaceAndRecordAllocation(any(Pipeline.class), any(ContainerID.class)); pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pendingOpsMock = mock(ContainerReplicaPendingOps.class); containerManager = new ContainerManagerImpl(conf, @@ -154,7 +155,7 @@ public void testGetMatchingContainerReturnsNullWhenNotEnoughSpaceInDatanodes() t // create an EC pipeline to test for EC containers ECReplicationConfig ecReplicationConfig = new ECReplicationConfig(3, 2); - pipelineManager.createPipeline(ecReplicationConfig); + pipelineManager.createPipeline(ecReplicationConfig, StorageTier.getDefaultTier()); pipeline = pipelineManager.getPipelines(ecReplicationConfig).iterator().next(); container = containerManager.getMatchingContainer(sizeRequired, "test", pipeline, Collections.emptySet()); assertNull(container); @@ -184,7 +185,7 @@ public void testGetMatchingContainerReturnsContainerWhenEnoughSpaceInDatanodes() // create an EC pipeline to test for EC containers ECReplicationConfig ecReplicationConfig = new ECReplicationConfig(3, 2); - spyPipelineManager.createPipeline(ecReplicationConfig); + spyPipelineManager.createPipeline(ecReplicationConfig, StorageTier.getDefaultTier()); pipeline = spyPipelineManager.getPipelines(ecReplicationConfig).iterator().next(); container = manager.getMatchingContainer(sizeRequired, "test", pipeline, Collections.emptySet()); assertNotNull(container); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReportHandler.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReportHandler.java index b441be1ce12..dd18d552f0f 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReportHandler.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerReportHandler.java @@ -52,6 +52,7 @@ import org.apache.hadoop.fs.StorageType; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -856,7 +857,7 @@ public void openContainerKeyAndBytesUsedUpdatedToMinimumOfAllReplicas() NodeStatus.inServiceHealthy()).iterator(); Pipeline pipeline = pipelineManager.createPipeline( - RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE)); + RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE), StorageTier.getDefaultTier()); final DatanodeDetails datanodeOne = nodeIterator.next(); final DatanodeDetails datanodeTwo = nodeIterator.next(); @@ -1012,7 +1013,7 @@ public void openECContainerKeyAndBytesUsedUpdatedToMinimumOfAllReplicas() final ContainerReportHandler reportHandler = new ContainerReportHandler( nodeManager, containerManager); - Pipeline pipeline = pipelineManager.createPipeline(repConfig); + Pipeline pipeline = pipelineManager.createPipeline(repConfig, StorageTier.getDefaultTier()); Map<Integer, DatanodeDetails> dns = new HashMap<>(); final Iterator<DatanodeDetails> nodeIterator = nodeManager.getNodes( NodeStatus.inServiceHealthy()).iterator(); @@ -1087,7 +1088,7 @@ public void closedECContainerKeyAndBytesUsedUpdatedToMinimumOfAllReplicas() final ContainerReportHandler reportHandler = new ContainerReportHandler( nodeManager, containerManager); - Pipeline pipeline = pipelineManager.createPipeline(repConfig); + Pipeline pipeline = pipelineManager.createPipeline(repConfig, StorageTier.getDefaultTier()); Map<Integer, DatanodeDetails> dns = new HashMap<>(); final Iterator<DatanodeDetails> nodeIterator = nodeManager.getNodes( NodeStatus.inServiceHealthy()).iterator(); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java index 6b1a2cca0b2..525bbb556d4 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestContainerStateManager.java @@ -43,6 +43,7 @@ import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; @@ -104,7 +105,7 @@ public void init() throws IOException, TimeoutException { ReplicationFactor.THREE)) .setNodes(new ArrayList<>()).build(); when(pipelineManager.createPipeline(StandaloneReplicationConfig.getInstance( - ReplicationFactor.THREE))).thenReturn(pipeline); + ReplicationFactor.THREE), StorageTier.getDefaultTier())).thenReturn(pipeline); when(pipelineManager.containsPipeline(any(PipelineID.class))).thenReturn(true); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestIncrementalContainerReportHandler.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestIncrementalContainerReportHandler.java index a42e38a6f4a..23455f1aae3 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestIncrementalContainerReportHandler.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/TestIncrementalContainerReportHandler.java @@ -58,6 +58,7 @@ import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.client.StorageTypeUtils; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -465,7 +466,7 @@ public void testOpenWithUnhealthyReplica() throws IOException { RatisReplicationConfig replicationConfig = RatisReplicationConfig.getInstance(HddsProtos.ReplicationFactor.THREE); - Pipeline pipeline = pipelineManager.createPipeline(replicationConfig); + Pipeline pipeline = pipelineManager.createPipeline(replicationConfig, StorageTier.getDefaultTier()); List<DatanodeDetails> nodes = pipeline.getNodes(); final DatanodeDetails datanodeOne = nodes.get(0); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestContainerPlacement.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestContainerPlacement.java index ef43ca856e0..d3015aca234 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestContainerPlacement.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestContainerPlacement.java @@ -42,6 +42,7 @@ import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.DatanodeID; @@ -112,7 +113,7 @@ public void setUp() throws Exception { pipelineManager = new MockPipelineManager(dbStore, scmhaManager, nodeManager); pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - HddsProtos.ReplicationFactor.THREE)); + HddsProtos.ReplicationFactor.THREE), StorageTier.getDefaultTier()); } @AfterEach diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java index f1f755c699b..6bb1f5baf7b 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java @@ -75,6 +75,7 @@ import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.fs.SpaceUsageSource; import org.apache.hadoop.hdds.protocol.DatanodeDetails; @@ -544,7 +545,7 @@ private void assertPipelineCreationFailsWithNotEnoughNodes( ReplicationConfig.fromProtoTypeAndFactor( HddsProtos.ReplicationType.RATIS, HddsProtos.ReplicationFactor.THREE); - scm.getPipelineManager().createPipeline(ratisThree); + scm.getPipelineManager().createPipeline(ratisThree, StorageTier.getDefaultTier()); }, "3 nodes should not have been found for a pipeline."); assertThat(ex.getMessage()).contains("Required 3. Found " + actualNodeCount); @@ -557,7 +558,7 @@ private void assertPipelineCreationFailsWithExceedingLimit(int limit) { HddsProtos.ReplicationFactor.THREE); SCMException ex = assertThrows( SCMException.class, - () -> scm.getPipelineManager().createPipeline(config), + () -> scm.getPipelineManager().createPipeline(config, StorageTier.getDefaultTier()), "3 nodes should not have been found for a pipeline."); assertThat(ex.getMessage()) .contains("Cannot create pipeline for StorageTier DISK as it would exceed the limit per datanode: " + limit); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java index 12bb97bfb5a..39d78ab1229 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java @@ -65,15 +65,17 @@ public MockPipelineManager(DBStore dbStore, SCMHAManager scmhaManager, NodeManag } @Override - public Pipeline createPipeline(ReplicationConfig replicationConfig) + public Pipeline createPipeline(ReplicationConfig replicationConfig, + StorageTier storageTier) throws IOException { return createPipeline(replicationConfig, Collections.emptyList(), - Collections.emptyList()); + Collections.emptyList(), storageTier); } @Override public Pipeline createPipeline(ReplicationConfig replicationConfig, - List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes) + List<DatanodeDetails> excludedNodes, List<DatanodeDetails> favoredNodes, + StorageTier storageTier) throws IOException { Pipeline pipeline; if (replicationConfig.getReplicationType() @@ -85,7 +87,7 @@ public Pipeline createPipeline(ReplicationConfig replicationConfig, ImmutableList.of(MockDatanodeDetails.randomDatanodeDetails(), MockDatanodeDetails.randomDatanodeDetails(), MockDatanodeDetails.randomDatanodeDetails()), - StorageTier.getDefaultTier()); + storageTier); } stateManager.addPipeline(pipeline.getProtobufMessage( @@ -180,13 +182,19 @@ public List<Pipeline> getPipelines(ReplicationConfig replicationConfig, return stateManager.getPipelines(replicationConfig, state); } + @Override + public List<Pipeline> getPipelines(ReplicationConfig replicationConfig, + final Pipeline.PipelineState state, StorageTier storageTier) { + return stateManager.getPipelines(replicationConfig, state, storageTier); + } + @Override public List<Pipeline> getPipelines(ReplicationConfig replicationConfig, final Pipeline.PipelineState state, final Collection<DatanodeDetails> excludeDns, - final Collection<PipelineID> excludePipelines) { + final Collection<PipelineID> excludePipelines, StorageTier storageTier) { return stateManager.getPipelines(replicationConfig, state, - excludeDns, excludePipelines); + excludeDns, excludePipelines, storageTier); } @Override diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java index 2d6353c92ec..bd542581ae0 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineManagerImpl.java @@ -207,12 +207,12 @@ public void testCreatePipeline() throws Exception { createPipelineManager(true, buffer1); assertTrue(pipelineManager.getPipelines().isEmpty()); Pipeline pipeline1 = pipelineManager.createPipeline( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); assertEquals(1, pipelineManager.getPipelines().size()); assertTrue(pipelineManager.containsPipeline(pipeline1.getId())); Pipeline pipeline2 = pipelineManager.createPipeline( - RatisReplicationConfig.getInstance(ReplicationFactor.ONE)); + RatisReplicationConfig.getInstance(ReplicationFactor.ONE), StorageTier.getDefaultTier()); assertEquals(2, pipelineManager.getPipelines().size()); assertTrue(pipelineManager.containsPipeline(pipeline2.getId())); @@ -236,7 +236,7 @@ public void testCreatePipeline() throws Exception { assertThat(pipelineManager2.getPipelines()).isNotEmpty(); assertEquals(3, pipelineManager.getPipelines().size()); Pipeline pipeline3 = pipelineManager2.createPipeline( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); buffer2.close(); assertEquals(4, pipelineManager2.getPipelines().size()); assertTrue(pipelineManager2.containsPipeline(pipeline3.getId())); @@ -249,7 +249,8 @@ public void testCreatePipelineShouldFailOnFollower() throws Exception { try (PipelineManager pipelineManager = createPipelineManager(false)) { assertTrue(pipelineManager.getPipelines().isEmpty()); assertFailsNotLeader(() -> pipelineManager.createPipeline( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE))); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), + StorageTier.getDefaultTier())); } } @@ -380,7 +381,7 @@ public void testPipelineReport() throws Exception { serviceManager, new EventQueue(), scmContext); Pipeline pipeline = pipelineManager .createPipeline(RatisReplicationConfig - .getInstance(ReplicationFactor.THREE)); + .getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); // pipeline is not healthy until all dns report List<DatanodeDetails> nodes = pipeline.getNodes(); @@ -434,7 +435,7 @@ localNodeManager, pipelineManager, mock(ContainerManager.class), List<DatanodeDetails> nodes = localNodeManager.getNodes(NodeStatus.inServiceHealthy()); assertEquals(nodeCount, nodes.size()); Pipeline pipeline = pipelineManager.createPipeline(RatisReplicationConfig - .getInstance(ReplicationFactor.THREE)); + .getInstance(ReplicationFactor.THREE), StorageTier.DISK); assertEquals(StorageTier.DISK, pipeline.getSupportedStorageTier()); assertFalse(pipelineManager.getPipeline(pipeline.getId()).isHealthy()); @@ -488,7 +489,7 @@ public void testPipelineCreationFailedMetric() throws Exception { for (int i = 0; i < maxPipelineCount; i++) { Pipeline pipeline = pipelineManager .createPipeline(RatisReplicationConfig - .getInstance(ReplicationFactor.THREE)); + .getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); assertNotNull(pipeline); } @@ -504,7 +505,9 @@ public void testPipelineCreationFailedMetric() throws Exception { //This should fail... SCMException e = assertThrows(SCMException.class, - () -> pipelineManager.createPipeline(RatisReplicationConfig.getInstance(ReplicationFactor.THREE))); + () -> pipelineManager.createPipeline( + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), + StorageTier.getDefaultTier())); // pipeline creation failed this time. assertEquals(ResultCodes.FAILED_TO_FIND_SUITABLE_NODE, e.getResult()); @@ -530,7 +533,7 @@ public void testPipelineOpenOnlyWhenLeaderReported() throws Exception { Pipeline pipeline = pipelineManager .createPipeline(RatisReplicationConfig - .getInstance(ReplicationFactor.THREE)); + .getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); // close manager buffer1.close(); pipelineManager.close(); @@ -578,7 +581,7 @@ public void testScrubPipelines() throws Exception { PipelineManagerImpl pipelineManager = createPipelineManager(true); Pipeline allocatedPipeline = pipelineManager .createPipeline(RatisReplicationConfig - .getInstance(ReplicationFactor.THREE)); + .getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); // At this point, pipeline is not at OPEN stage. assertEquals(Pipeline.PipelineState.ALLOCATED, allocatedPipeline.getPipelineState()); @@ -591,7 +594,7 @@ public void testScrubPipelines() throws Exception { Pipeline closedPipeline = pipelineManager .createPipeline(RatisReplicationConfig - .getInstance(ReplicationFactor.THREE)); + .getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(closedPipeline.getId()); pipelineManager.closePipeline(closedPipeline.getId()); @@ -643,7 +646,7 @@ public void testScrubPipelines() throws Exception { public void testScrubOpenWithUnregisteredNodes() throws Exception { PipelineManagerImpl pipelineManager = createPipelineManager(true); Pipeline pipeline = pipelineManager - .createPipeline(new ECReplicationConfig(3, 2)); + .createPipeline(new ECReplicationConfig(3, 2), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline.getId()); // Scrubbing the pipelines should not affect this pipeline @@ -687,7 +690,8 @@ public void testPipelineNotCreatedUntilSafeModePrecheck() throws Exception { PipelineManagerImpl pipelineManager = createPipelineManager(true); assertThrows(IOException.class, - () -> pipelineManager.createPipeline(RatisReplicationConfig.getInstance(ReplicationFactor.THREE)), + () -> pipelineManager.createPipeline(RatisReplicationConfig.getInstance(ReplicationFactor.THREE), + StorageTier.getDefaultTier()), "Pipelines should not have been created"); // No pipeline is created. assertTrue(pipelineManager.getPipelines().isEmpty()); @@ -696,7 +700,7 @@ public void testPipelineNotCreatedUntilSafeModePrecheck() throws Exception { // raised. Pipeline pipeline = pipelineManager .createPipeline(RatisReplicationConfig - .getInstance(ReplicationFactor.ONE)); + .getInstance(ReplicationFactor.ONE), StorageTier.getDefaultTier()); assertTrue(pipelineManager .getPipelines(RatisReplicationConfig .getInstance(ReplicationFactor.ONE)) @@ -745,7 +749,8 @@ public void testAddContainerWithClosedPipelineScmStart() throws Exception { SCMDBDefinition.PIPELINES.getTable(dbStore); Pipeline pipeline = pipelineManager.createPipeline( RatisReplicationConfig - .getInstance(HddsProtos.ReplicationFactor.THREE)); + .getInstance(HddsProtos.ReplicationFactor.THREE), + StorageTier.getDefaultTier()); PipelineID pipelineID = pipeline.getId(); pipelineManager.addContainerToPipeline(pipelineID, ContainerID.valueOf(1)); pipelineManager.getStateManager().updatePipelineState( @@ -768,7 +773,8 @@ public void testAddContainerWithClosedPipeline() throws Exception { SCMDBDefinition.PIPELINES.getTable(dbStore); Pipeline pipeline = pipelineManager.createPipeline( RatisReplicationConfig - .getInstance(HddsProtos.ReplicationFactor.THREE)); + .getInstance(HddsProtos.ReplicationFactor.THREE), + StorageTier.getDefaultTier()); PipelineID pipelineID = pipeline.getId(); pipelineManager.addContainerToPipeline(pipelineID, ContainerID.valueOf(1)); pipelineManager.getStateManager().updatePipelineState( @@ -786,7 +792,8 @@ public void testPipelineCloseFlow() throws IOException { PipelineManagerImpl pipelineManager = createPipelineManager(true); Pipeline pipeline = pipelineManager.createPipeline( RatisReplicationConfig - .getInstance(HddsProtos.ReplicationFactor.THREE)); + .getInstance(HddsProtos.ReplicationFactor.THREE), + StorageTier.getDefaultTier()); PipelineID pipelineID = pipeline.getId(); ContainerManager containerManager = scm.getContainerManager(); ContainerInfo containerInfo = HddsTestUtils. @@ -928,12 +935,12 @@ public void testWaitForAllocatedPipeline() throws IOException { // Throw on pipeline creates, so no new pipelines can be created doThrow(SCMException.class).when(pipelineManagerSpy) - .createPipeline(any(), any(), anyList()); + .createPipeline(any(), any(), anyList(), any(StorageTier.class)); provider = new WritableRatisContainerProvider( pipelineManagerSpy, containerManager, pipelineChoosingPolicy); // Add a single pipeline to manager, (in the allocated state) - allocatedPipeline = pipelineManager.createPipeline(repConfig); + allocatedPipeline = pipelineManager.createPipeline(repConfig, StorageTier.getDefaultTier()); pipelineManager.getStateManager() .updatePipelineState(allocatedPipeline.getId() .getProtobuf(), HddsProtos.PipelineState.PIPELINE_ALLOCATED); @@ -970,7 +977,7 @@ public void testWaitForAllocatedPipeline() throws IOException { ContainerInfo c = provider.getContainer(1, repConfig, - owner, new ExcludeList()); + owner, new ExcludeList(), StorageTier.getDefaultTier()); assertEquals(c, container, "Expected container was returned"); // Confirm that waitOnePipelineReady was called on allocated pipelines @@ -1034,7 +1041,8 @@ private void sendPipelineReport( private static Pipeline assertAllocate(PipelineManagerImpl pipelineManager) { Pipeline pipeline = assertDoesNotThrow( () -> pipelineManager.createPipeline( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE))); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), + StorageTier.getDefaultTier())); assertEquals(1, pipelineManager.getPipelines().size()); assertTrue(pipelineManager.containsPipeline(pipeline.getId())); assertEquals(ALLOCATED, pipeline.getPipelineState()); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java index 95b8b713af3..9ecaa8d5730 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java @@ -248,10 +248,16 @@ public void testPipelineEngagementLimitIsStorageTierAware() throws Exception { } @Test - public void testGlobalPipelineNumberLimitIsNotStorageTierAware() throws Exception { + public void testGlobalPipelineNumberLimitIsStorageTierAware() throws Exception { OzoneConfiguration pipelineLimitConf = new OzoneConfiguration(); pipelineLimitConf.setInt(OZONE_SCM_RATIS_PIPELINE_LIMIT, 1); init(0, pipelineLimitConf, StorageTier.DISK); + // init(0, ...) sets numPipelinePerDatanode=0 which would cause + // filterPipelineEngagement to exclude every datanode. Bump it here + // before creating the spy so filterPipelineEngagement sees a positive + // limit. The provider's datanodePipelineLimit config is still 0, so + // the *global* branch of exceedPipelineNumberLimit is what runs. + nodeManager.setNumPipelinePerDatanode(5); List<DatanodeDetails> diskAndSsdNodes = datanodeList.stream() .map(TestRatisPipelineProvider::datanodeInfoWithDiskAndSsd) .collect(Collectors.toList()); @@ -263,6 +269,7 @@ public void testGlobalPipelineNumberLimitIsNotStorageTierAware() throws Exceptio RatisPipelineProvider localProvider = new RatisPipelineProvider( nodeManagerSpy, stateManager, pipelineLimitConf, new EventQueue(), SCMContext.emptyContext()); + // Fill the DISK tier up to the global limit. for (int i = 0; i < 2; i++) { addPipeline(diskAndSsdNodes.subList(i * 3, i * 3 + 3), Pipeline.PipelineState.OPEN, @@ -270,11 +277,11 @@ nodeManagerSpy, stateManager, pipelineLimitConf, new EventQueue(), StorageTier.DISK); } - SCMException exception = assertThrows(SCMException.class, () -> - localProvider.create( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE), StorageTier.SSD)); - - assertThat(exception.getMessage()).contains("would exceed the limit"); + // SSD is a distinct tier, so the global limit should not block creation. + Pipeline pipeline = localProvider.create( + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), StorageTier.SSD); + assertPipelineProperties(pipeline, ReplicationFactor.THREE, + REPLICATION_TYPE, Pipeline.PipelineState.ALLOCATED, StorageTier.SSD); } private List<DatanodeDetails> createListOfNodes(int count) { diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableECContainerProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableECContainerProvider.java index 6e4447d2e36..29de71bd6b5 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableECContainerProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableECContainerProvider.java @@ -52,6 +52,7 @@ import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; @@ -209,7 +210,7 @@ private Set<ContainerInfo> assertDistinctContainers(int n) Set<ContainerInfo> allocatedContainers = new HashSet<>(); for (int i = 0; i < n; i++) { ContainerInfo container = - provider.getContainer(1, repConfig, OWNER, new ExcludeList()); + provider.getContainer(1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); assertThat(allocatedContainers) .withFailMessage("Provided existing container for request " + i) .doesNotContain(container); @@ -222,7 +223,7 @@ private void assertReusesExisting(Set<ContainerInfo> existing, int n) throws IOException { for (int i = 0; i < 3 * n; i++) { ContainerInfo container = - provider.getContainer(1, repConfig, OWNER, new ExcludeList()); + provider.getContainer(1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); assertThat(existing) .withFailMessage("Provided new container for request " + i) .contains(container); @@ -237,7 +238,7 @@ public void testPiplineLimitIgnoresExcludedPipelines( Set<ContainerInfo> allocatedContainers = new HashSet<>(); for (int i = 0; i < providerConf.getMinimumPipelines(); i++) { ContainerInfo container = provider.getContainer( - 1, repConfig, OWNER, new ExcludeList()); + 1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); allocatedContainers.add(container); } // We have the min limit of pipelines, but then exclude one. It should use @@ -248,7 +249,7 @@ public void testPiplineLimitIgnoresExcludedPipelines( .stream().findFirst().get().getPipelineID(); exclude.addPipeline(excludedID); - ContainerInfo c = provider.getContainer(1, repConfig, OWNER, exclude); + ContainerInfo c = provider.getContainer(1, repConfig, OWNER, exclude, StorageTier.getDefaultTier()); assertNotEquals(excludedID, c.getPipelineID()); assertThat(allocatedContainers).contains(c); } @@ -263,7 +264,7 @@ public void testNewPipelineNotCreatedIfAllPipelinesExcluded( Set<ContainerInfo> allocatedContainers = new HashSet<>(); for (int i = 0; i < providerConf.getMinimumPipelines(); i++) { ContainerInfo container = provider.getContainer( - 1, repConfig, OWNER, new ExcludeList()); + 1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); allocatedContainers.add(container); } // We have the min limit of pipelines, but then exclude them all @@ -272,7 +273,7 @@ public void testNewPipelineNotCreatedIfAllPipelinesExcluded( exclude.addPipeline(c.getPipelineID()); } assertThrows(IOException.class, () -> provider.getContainer( - 1, repConfig, OWNER, exclude)); + 1, repConfig, OWNER, exclude, StorageTier.getDefaultTier())); } @ParameterizedTest @@ -283,7 +284,7 @@ void newPipelineCreatedIfSoftLimitReached(PipelineChoosePolicy policy) providerConf.setMinimumPipelines(1); provider = createSubject(policy); ContainerInfo container = provider.getContainer( - 1, repConfig, OWNER, new ExcludeList()); + 1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); ExcludeList exclude = new ExcludeList(); exclude.addPipeline(container.getPipelineID()); @@ -291,7 +292,7 @@ void newPipelineCreatedIfSoftLimitReached(PipelineChoosePolicy policy) pipelineManager.getPipeline(container.getPipelineID()).getFirstNode()); ContainerInfo newContainer = provider.getContainer( - 1, repConfig, OWNER, exclude); + 1, repConfig, OWNER, exclude, StorageTier.getDefaultTier()); assertNotSame(container, newContainer); } @@ -305,7 +306,7 @@ public void testNewPipelineNotCreatedIfAllContainersExcluded( Set<ContainerInfo> allocatedContainers = new HashSet<>(); for (int i = 0; i < providerConf.getMinimumPipelines(); i++) { ContainerInfo container = provider.getContainer( - 1, repConfig, OWNER, new ExcludeList()); + 1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); allocatedContainers.add(container); } // We have the min limit of pipelines, but then exclude all the associated @@ -315,7 +316,7 @@ public void testNewPipelineNotCreatedIfAllContainersExcluded( exclude.addConatinerId(c.containerID()); } assertThrows(IOException.class, () -> provider.getContainer( - 1, repConfig, OWNER, exclude)); + 1, repConfig, OWNER, exclude, StorageTier.getDefaultTier())); } @ParameterizedTest @@ -327,14 +328,15 @@ public void testUnableToCreateAnyPipelinesThrowsException( @Override public Pipeline createPipeline(ReplicationConfig repConf, List<DatanodeDetails> excludedNodes, - List<DatanodeDetails> favoredNodes) throws IOException { + List<DatanodeDetails> favoredNodes, + StorageTier storageTier) throws IOException { throw new IOException("Cannot create pipelines"); } }; provider = createSubject(policy); IOException ioException = assertThrows(IOException.class, - () -> provider.getContainer(1, repConfig, OWNER, new ExcludeList())); + () -> provider.getContainer(1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier())); assertThat(ioException.getMessage()) .contains("Cannot create pipelines"); } @@ -351,25 +353,26 @@ public void testExistingPipelineReturnedWhenNewCannotBeCreated( @Override public Pipeline createPipeline(ReplicationConfig repConf, List<DatanodeDetails> excludedNodes, - List<DatanodeDetails> favoredNodes) + List<DatanodeDetails> favoredNodes, + StorageTier storageTier) throws IOException { if (throwError) { throw new RocksDatabaseException("Cannot create pipelines"); } throwError = true; - return super.createPipeline(repConfig); + return super.createPipeline(repConfig, storageTier); } }; provider = createSubject(policy); IOException ioException = assertThrows(IOException.class, - () -> provider.getContainer(1, repConfig, OWNER, new ExcludeList())); + () -> provider.getContainer(1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier())); assertThat(ioException.getMessage()) .contains("Cannot create pipelines"); for (int i = 0; i < 5; i++) { ioException = assertThrows(IOException.class, - () -> provider.getContainer(1, repConfig, OWNER, new ExcludeList())); + () -> provider.getContainer(1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier())); assertThat(ioException.getMessage()) .contains("Cannot create pipelines"); } @@ -392,13 +395,14 @@ public void testNewContainerAllocatedAndPipelinesClosedIfNoSpaceInExisting( // We ask for a space of 50 MB, and will actually need 50 MB space. ContainerInfo newContainer = provider.getContainer(50 * 1024 * 1024, repConfig, OWNER, - new ExcludeList()); + new ExcludeList(), StorageTier.getDefaultTier()); assertNotNull(newContainer); assertThat(allocatedContainers).contains(newContainer); // Now get a new container where there is not enough space in the existing // and ensure a new container gets created. newContainer = provider.getContainer( - 128 * 1024 * 1024, repConfig, OWNER, new ExcludeList()); + 128 * 1024 * 1024, repConfig, OWNER, new ExcludeList(), + StorageTier.getDefaultTier()); assertNotNull(newContainer); assertThat(allocatedContainers).doesNotContain(newContainer); // The original pipelines should all be closed, triggered by the lack of @@ -431,7 +435,7 @@ public NavigableSet<ContainerID> getContainersInPipeline(PipelineID pipelineID) // Now attempt to get a container - any attempt to use an existing with // throw PNF and then we must allocate a new one ContainerInfo newContainer = - provider.getContainer(1, repConfig, OWNER, new ExcludeList()); + provider.getContainer(1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); assertNotNull(newContainer); assertThat(allocatedContainers).doesNotContain(newContainer); } @@ -451,7 +455,7 @@ public void testContainerNotFoundWhenAttemptingToUseExisting( }).when(containerManager).getContainer(any(ContainerID.class)); ContainerInfo newContainer = - provider.getContainer(1, repConfig, OWNER, new ExcludeList()); + provider.getContainer(1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); assertNotNull(newContainer); assertThat(allocatedContainers).doesNotContain(newContainer); @@ -473,7 +477,7 @@ public void testPipelineOpenButContainerRemovedFromIt( Set<ContainerInfo> allocatedContainers = new HashSet<>(); for (int i = 0; i < providerConf.getMinimumPipelines(); i++) { ContainerInfo container = provider.getContainer( - 1, repConfig, OWNER, new ExcludeList()); + 1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); assertThat(allocatedContainers).doesNotContain(container); allocatedContainers.add(container); // Remove the container from the pipeline to simulate closing it @@ -481,7 +485,7 @@ public void testPipelineOpenButContainerRemovedFromIt( container.getPipelineID(), container.containerID()); } ContainerInfo newContainer = provider.getContainer( - 1, repConfig, OWNER, new ExcludeList()); + 1, repConfig, OWNER, new ExcludeList(), StorageTier.getDefaultTier()); assertThat(allocatedContainers).doesNotContain(newContainer); for (ContainerInfo c : allocatedContainers) { Pipeline pipeline = pipelineManager.getPipeline(c.getPipelineID()); @@ -518,7 +522,7 @@ public void testExcludedOpenPipelineWithClosedContainerIsClosed( // expecting a new container to be created ContainerInfo containerInfo = provider.getContainer(1, repConfig, OWNER, - excludeList); + excludeList, StorageTier.getDefaultTier()); assertThat(allocated).doesNotContain(containerInfo); for (ContainerInfo c : allocated) { Pipeline pipeline = pipelineManager.getPipeline(c.getPipelineID()); @@ -536,11 +540,12 @@ public void testExcludedNodesPassedToCreatePipelineIfProvided( // EmptyList should be passed if there are no nodes excluded. ContainerInfo container = provider.getContainer( - 1, repConfig, OWNER, excludeList); + 1, repConfig, OWNER, excludeList, StorageTier.getDefaultTier()); assertNotNull(container); verify(pipelineManagerSpy).createPipeline(repConfig, - Collections.emptyList(), Collections.emptyList()); + Collections.emptyList(), Collections.emptyList(), + StorageTier.getDefaultTier()); // If nodes are excluded then the excluded nodes should be passed through to // the create pipeline call. @@ -549,10 +554,10 @@ public void testExcludedNodesPassedToCreatePipelineIfProvided( new ArrayList<>(excludeList.getDatanodes()); container = provider.getContainer( - 1, repConfig, OWNER, excludeList); + 1, repConfig, OWNER, excludeList, StorageTier.getDefaultTier()); assertNotNull(container); verify(pipelineManagerSpy).createPipeline(repConfig, excludedNodes, - Collections.emptyList()); + Collections.emptyList(), StorageTier.getDefaultTier()); } private ContainerInfo createContainer(Pipeline pipeline, diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableRatisContainerProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableRatisContainerProvider.java index a1ba81d0a70..6c43a0d1759 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableRatisContainerProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestWritableRatisContainerProvider.java @@ -24,6 +24,8 @@ import static org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState.OPEN; import static org.junit.jupiter.api.Assertions.assertSame; import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.ArgumentMatchers.eq; import static org.mockito.Mockito.never; import static org.mockito.Mockito.times; import static org.mockito.Mockito.verify; @@ -35,6 +37,7 @@ import java.util.concurrent.atomic.AtomicLong; import org.apache.hadoop.hdds.client.RatisReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.PipelineChoosePolicy; @@ -78,7 +81,8 @@ void returnsExistingContainer() throws Exception { existingPipelines(pipeline); - ContainerInfo container = createSubject().getContainer(CONTAINER_SIZE, REPLICATION_CONFIG, OWNER, NO_EXCLUSION); + ContainerInfo container = createSubject().getContainer( + CONTAINER_SIZE, REPLICATION_CONFIG, OWNER, NO_EXCLUSION, StorageTier.getDefaultTier()); assertSame(existingContainer, container); verifyPipelineNotCreated(); @@ -92,7 +96,8 @@ void skipsPipelineWithoutContainer() throws Exception { Pipeline pipelineWithoutContainer = MockPipeline.createPipeline(3); existingPipelines(pipelineWithoutContainer, pipeline); - ContainerInfo container = createSubject().getContainer(CONTAINER_SIZE, REPLICATION_CONFIG, OWNER, NO_EXCLUSION); + ContainerInfo container = createSubject().getContainer( + CONTAINER_SIZE, REPLICATION_CONFIG, OWNER, NO_EXCLUSION, StorageTier.getDefaultTier()); assertSame(existingContainer, container); verifyPipelineNotCreated(); @@ -102,7 +107,8 @@ void skipsPipelineWithoutContainer() throws Exception { void createsNewContainerIfNoneFound() throws Exception { ContainerInfo newContainer = createNewContainerOnDemand(); - ContainerInfo container = createSubject().getContainer(CONTAINER_SIZE, REPLICATION_CONFIG, OWNER, NO_EXCLUSION); + ContainerInfo container = createSubject().getContainer( + CONTAINER_SIZE, REPLICATION_CONFIG, OWNER, NO_EXCLUSION, StorageTier.getDefaultTier()); assertSame(newContainer, container); verifyPipelineCreated(); @@ -113,7 +119,8 @@ void failsIfContainerCannotBeCreated() throws Exception { throwWhenCreatePipeline(); assertThrows(IOException.class, - () -> createSubject().getContainer(CONTAINER_SIZE, REPLICATION_CONFIG, OWNER, NO_EXCLUSION)); + () -> createSubject().getContainer( + CONTAINER_SIZE, REPLICATION_CONFIG, OWNER, NO_EXCLUSION, StorageTier.getDefaultTier())); verifyPipelineCreated(); } @@ -123,7 +130,8 @@ private void existingPipelines(Pipeline... pipelines) { } private void existingPipelines(List<Pipeline> pipelines) { - when(pipelineManager.getPipelines(REPLICATION_CONFIG, OPEN, emptySet(), emptySet())) + when(pipelineManager.getPipelines(REPLICATION_CONFIG, OPEN, emptySet(), emptySet(), + StorageTier.getDefaultTier())) .thenReturn(pipelines); } @@ -141,10 +149,11 @@ private ContainerInfo pipelineHasContainer(Pipeline pipeline) { private ContainerInfo createNewContainerOnDemand() throws IOException { Pipeline newPipeline = MockPipeline.createPipeline(3); - when(pipelineManager.createPipeline(REPLICATION_CONFIG)) + when(pipelineManager.createPipeline(eq(REPLICATION_CONFIG), any(StorageTier.class))) .thenReturn(newPipeline); - when(pipelineManager.getPipelines(REPLICATION_CONFIG, OPEN, emptySet(), emptySet())) + when(pipelineManager.getPipelines(REPLICATION_CONFIG, OPEN, emptySet(), emptySet(), + StorageTier.getDefaultTier())) .thenReturn(emptyList()) .thenReturn(new ArrayList<>(singletonList(newPipeline))); @@ -152,7 +161,7 @@ private ContainerInfo createNewContainerOnDemand() throws IOException { } private void throwWhenCreatePipeline() throws IOException { - when(pipelineManager.createPipeline(REPLICATION_CONFIG)) + when(pipelineManager.createPipeline(REPLICATION_CONFIG, StorageTier.getDefaultTier())) .thenThrow(new SCMException(SCMException.ResultCodes.FAILED_TO_FIND_SUITABLE_NODE)); } @@ -163,16 +172,18 @@ private WritableRatisContainerProvider createSubject() { private void verifyPipelineCreated() throws IOException { verify(pipelineManager, times(2)) - .getPipelines(REPLICATION_CONFIG, OPEN, emptySet(), emptySet()); + .getPipelines(REPLICATION_CONFIG, OPEN, emptySet(), emptySet(), + StorageTier.getDefaultTier()); verify(pipelineManager) - .createPipeline(REPLICATION_CONFIG); + .createPipeline(REPLICATION_CONFIG, StorageTier.getDefaultTier()); } private void verifyPipelineNotCreated() throws IOException { verify(pipelineManager, times(1)) - .getPipelines(REPLICATION_CONFIG, OPEN, emptySet(), emptySet()); + .getPipelines(REPLICATION_CONFIG, OPEN, emptySet(), emptySet(), + StorageTier.getDefaultTier()); verify(pipelineManager, never()) - .createPipeline(REPLICATION_CONFIG); + .createPipeline(REPLICATION_CONFIG, StorageTier.getDefaultTier()); } } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestHealthyPipelineSafeModeRule.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestHealthyPipelineSafeModeRule.java index b1312512d56..d33438d1aa5 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestHealthyPipelineSafeModeRule.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestHealthyPipelineSafeModeRule.java @@ -30,6 +30,7 @@ import java.util.List; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -154,15 +155,15 @@ public void testHealthyPipelineSafeModeRuleWithPipelines() throws Exception { // Create 3 pipelines Pipeline pipeline1 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline1.getId()); Pipeline pipeline2 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline2.getId()); Pipeline pipeline3 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline3.getId()); // Mark pipeline healthy @@ -247,15 +248,15 @@ public void testHealthyPipelineSafeModeRuleWithMixedPipelines() // Create 3 pipelines Pipeline pipeline1 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.ONE)); + ReplicationFactor.ONE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline1.getId()); Pipeline pipeline2 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline2.getId()); Pipeline pipeline3 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline3.getId()); // Mark pipelines healthy @@ -338,13 +339,13 @@ public void testHealthyPipelineThresholdIncreasesWithMorePipelinesAndReports() // blocked once safe mode prechecks have not passed. Pipeline pipeline1 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); Pipeline pipeline2 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); Pipeline pipeline3 = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); // Start with one healthy open pipeline. Threshold is small at this point. pipelineManager.openPipeline(pipeline1.getId()); @@ -446,7 +447,7 @@ public void testPipelineIgnoredWhenDnIsUnhealthy() throws Exception { // Create a Ratis pipeline with 3 replicas Pipeline pipeline = pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline.getId()); pipeline = pipelineManager.getPipeline(pipeline.getId()); MockRatisPipelineProvider.markPipelineHealthy(pipeline); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestOneReplicaPipelineSafeModeRule.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestOneReplicaPipelineSafeModeRule.java index ab1418a8bf0..9ef13f0dd1a 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestOneReplicaPipelineSafeModeRule.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestOneReplicaPipelineSafeModeRule.java @@ -31,6 +31,7 @@ import java.util.Map; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -238,7 +239,7 @@ private void createPipelines(int count, HddsProtos.ReplicationFactor factor) throws Exception { for (int i = 0; i < count; i++) { Pipeline pipeline = pipelineManager.createPipeline( - RatisReplicationConfig.getInstance(factor)); + RatisReplicationConfig.getInstance(factor), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline.getId()); } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestSCMSafeModeManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestSCMSafeModeManager.java index cf262873eb0..e50ca464f1b 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestSCMSafeModeManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/safemode/TestSCMSafeModeManager.java @@ -41,6 +41,7 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -386,7 +387,7 @@ public void testSafeModeExitRuleWithPipelineAvailabilityCheck( // Create pipeline Pipeline pipeline = pipelineManager.createPipeline( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipelineManager.openPipeline(pipeline.getId()); // Mark pipeline healthy @@ -795,7 +796,7 @@ public void testSafeModePipelineExitRule() throws Exception { Pipeline pipeline = pipelineManager.createPipeline( RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)); + ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipeline = pipelineManager.getPipeline(pipeline.getId()); MockRatisPipelineProvider.markPipelineHealthy(pipeline); @@ -963,7 +964,7 @@ public void testPipelinesNotCreatedUntilPreCheckPasses() throws Exception { assertTrue(scmSafeModeManager.getInSafeMode()); Pipeline pipeline = pipelineManager.createPipeline( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); // Mark pipeline healthy pipeline = pipelineManager.getPipeline(pipeline.getId()); @@ -1082,7 +1083,7 @@ public void testSafeModePeriodicLoggingStopsOnNormalExit() throws Exception { assertThat(logCapturer.getOutput()).contains("SCM SafeMode Status | state=PRE_CHECKS_PASSED"); Pipeline pipeline = pipelineManager.createPipeline( - RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); + RatisReplicationConfig.getInstance(ReplicationFactor.THREE), StorageTier.getDefaultTier()); pipeline = pipelineManager.getPipeline(pipeline.getId()); MockRatisPipelineProvider.markPipelineHealthy(pipeline); firePipelineEvent(pipelineManager, pipeline); diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java index c1973891d0a..2b28903d620 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java @@ -160,6 +160,7 @@ public ContainerWithPipeline createContainer(String owner) XceiverClientSpi client = null; XceiverClientManager clientManager = getXceiverClientManager(); try { + // TODO Support create Container Command with StorageTier ContainerWithPipeline containerWithPipeline = storageContainerLocationClient. allocateContainer(replicationType, replicationFactor, owner); diff --git a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAsPassiveScm.java b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAsPassiveScm.java index 1b47f972382..64738d9c7a8 100644 --- a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAsPassiveScm.java +++ b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/TestReconAsPassiveScm.java @@ -31,6 +31,7 @@ import java.util.Optional; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.scm.XceiverClientGrpc; import org.apache.hadoop.hdds.scm.container.ContainerID; @@ -107,7 +108,8 @@ public void testDatanodeRegistrationAndReports() throws Exception { UnsupportedOperationException exception = assertThrows( UnsupportedOperationException.class, () -> reconPipelineManager - .createPipeline(RatisReplicationConfig.getInstance(ONE))); + .createPipeline(RatisReplicationConfig.getInstance(ONE), + StorageTier.getDefaultTier())); assertTrue(exception.getMessage() .contains("Trying to create pipeline in Recon, which is prohibited!")); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestSafeMode.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestSafeMode.java index 1300c61b1f8..e2cda809fbb 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestSafeMode.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/fs/ozone/TestSafeMode.java @@ -35,6 +35,7 @@ import org.apache.hadoop.fs.SafeMode; import org.apache.hadoop.fs.SafeModeAction; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; import org.apache.hadoop.hdds.utils.IOUtils; @@ -117,7 +118,8 @@ private void testSafeMode(Function<OzoneConfiguration, String> fsRoot) RatisReplicationConfig.getInstance(THREE); assertThrows(IOException.class, () -> cluster.getStorageContainerManager() .getWritableContainerFactory() - .getContainer(MB, replication, OZONE, new ExcludeList())); + .getContainer(MB, replication, OZONE, new ExcludeList(), + StorageTier.getDefaultTier())); } finally { IOUtils.closeQuietly(fs); } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestMultiRaftSetup.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestMultiRaftSetup.java index f780ec184d2..655ad340a29 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestMultiRaftSetup.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestMultiRaftSetup.java @@ -29,6 +29,7 @@ import java.util.stream.Collectors; import org.apache.hadoop.hdds.HddsConfigKeys; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -98,7 +99,7 @@ public void testMultiRaftNotSamePeers() throws Exception { // datanode pipeline limit is set to 2, but only one set of 3 pipelines // will be created. Further pipeline creation should fail assertEquals(1, pipelineManager.getPipelines(RATIS_THREE).size()); - assertThrows(IOException.class, () -> pipelineManager.createPipeline(RATIS_THREE)); + assertThrows(IOException.class, () -> pipelineManager.createPipeline(RATIS_THREE, StorageTier.getDefaultTier())); shutdown(); } @@ -119,7 +120,7 @@ public void testMultiRaft() throws Exception { .filter((dn) -> nodeManager.getPipelinesCount(dn) > 2).collect( Collectors.toList()); assertEquals(1, dns.size()); - assertThrows(IOException.class, () -> pipelineManager.createPipeline(RATIS_THREE)); + assertThrows(IOException.class, () -> pipelineManager.createPipeline(RATIS_THREE, StorageTier.getDefaultTier())); Collection<PipelineID> pipelineIds = nodeManager.getPipelines(dns.get(0)); // Only one dataode should have 3 pipelines in total, 1 RATIS ONE pipeline // and 2 RATIS 3 pipeline diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineCreateAndDestroy.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineCreateAndDestroy.java index 3c3eee25393..bbbd41d4296 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineCreateAndDestroy.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineCreateAndDestroy.java @@ -30,6 +30,7 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.TimeoutException; import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor; @@ -142,7 +143,7 @@ public void testPipelineCreationOnNodeRestart() throws Exception { // try creating another pipeline now SCMException ioe = assertThrows(SCMException.class, () -> pipelineManager.createPipeline(RatisReplicationConfig.getInstance( - ReplicationFactor.THREE)), + ReplicationFactor.THREE), StorageTier.getDefaultTier()), "pipeline creation should fail after shutting down pipeline"); assertEquals(SCMException.ResultCodes.FAILED_TO_FIND_SUITABLE_NODE, ioe.getResult()); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java index a70333d4b21..878b57b6faa 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/storage/TestContainerCommandsEC.java @@ -57,6 +57,7 @@ import org.apache.hadoop.hdds.client.ECReplicationConfig; import org.apache.hadoop.hdds.client.ECReplicationConfig.EcCodec; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.MockDatanodeDetails; @@ -448,7 +449,7 @@ public void testCreateRecoveryContainer() throws Exception { new XceiverClientManager(config)) { ECReplicationConfig replicationConfig = new ECReplicationConfig(3, 2); Pipeline newPipeline = - scm.getPipelineManager().createPipeline(replicationConfig); + scm.getPipelineManager().createPipeline(replicationConfig, StorageTier.getDefaultTier()); scm.getPipelineManager().activatePipeline(newPipeline.getId()); final ContainerInfo container = scm.getContainerManager() .allocateContainer(replicationConfig, "test"); @@ -537,7 +538,7 @@ public void testCreateRecoveryContainerAfterDNRestart() throws Exception { new XceiverClientManager(config)) { ECReplicationConfig replicationConfig = new ECReplicationConfig(3, 2); Pipeline newPipeline = - scm.getPipelineManager().createPipeline(replicationConfig); + scm.getPipelineManager().createPipeline(replicationConfig, StorageTier.getDefaultTier()); scm.getPipelineManager().activatePipeline(newPipeline.getId()); final ContainerInfo container = scm.getContainerManager() .allocateContainer(replicationConfig, "test"); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestHDDSUpgrade.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestHDDSUpgrade.java index c0435088d25..19417f35bbe 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestHDDSUpgrade.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestHDDSUpgrade.java @@ -53,6 +53,7 @@ import org.apache.hadoop.hdds.client.ReplicationConfig; import org.apache.hadoop.hdds.client.ReplicationFactor; import org.apache.hadoop.hdds.client.ReplicationType; +import org.apache.hadoop.hdds.client.StorageTier; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.DatanodeDetails; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; @@ -224,7 +225,7 @@ private void createKey() throws IOException { */ private void testPostUpgradePipelineCreation() throws IOException, TimeoutException { - Pipeline ratisPipeline1 = scmPipelineManager.createPipeline(RATIS_THREE); + Pipeline ratisPipeline1 = scmPipelineManager.createPipeline(RATIS_THREE, StorageTier.getDefaultTier()); scmPipelineManager.openPipeline(ratisPipeline1.getId()); assertEquals(0, scmPipelineManager.getNumberOfContainers(ratisPipeline1.getId())); --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
