rich7420 commented on code in PR #11146: URL: https://github.com/apache/ozone/pull/11146#discussion_r3886088698
########## hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestPerVolumePushReplication.java: ########## @@ -0,0 +1,545 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.ozone.container.replication; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static java.util.Collections.singleton; +import static java.util.Collections.singletonList; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_NODE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_PIPELINE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State.CLOSED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.DECOMMISSIONED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.DEAD; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DATANODE_ADMIN_MONITOR_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DEADNODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HEARTBEAT_PROCESS_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_STALENODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.getDNHostAndPort; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachHealthState; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachOpState; +import static org.apache.hadoop.hdds.scm.pipeline.MockPipeline.createPipeline; +import static org.apache.hadoop.hdds.scm.storage.ContainerProtocolCalls.createContainer; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForContainerClose; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForReplicaCount; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.List; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.BooleanSupplier; +import java.util.stream.Collectors; +import org.apache.hadoop.hdds.HddsConfigKeys; +import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.ECReplicationConfig; +import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.conf.StorageUnit; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.DatanodeID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.scm.ScmConfigKeys; +import org.apache.hadoop.hdds.scm.XceiverClientFactory; +import org.apache.hadoop.hdds.scm.XceiverClientManager; +import org.apache.hadoop.hdds.scm.XceiverClientSpi; +import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; +import org.apache.hadoop.hdds.scm.container.ContainerID; +import org.apache.hadoop.hdds.scm.container.ContainerInfo; +import org.apache.hadoop.hdds.scm.container.ContainerManager; +import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager.ReplicationManagerConfiguration; +import org.apache.hadoop.hdds.scm.node.NodeManager; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.hdds.scm.pipeline.PipelineManager; +import org.apache.hadoop.hdds.scm.server.StorageContainerManager; +import org.apache.hadoop.hdds.utils.IOUtils; +import org.apache.hadoop.ozone.DataTestUtil; +import org.apache.hadoop.ozone.HddsDatanodeService; +import org.apache.hadoop.ozone.MiniOzoneCluster; +import org.apache.hadoop.ozone.OzoneConfigKeys; +import org.apache.hadoop.ozone.UniformDatanodesFactory; +import org.apache.hadoop.ozone.client.OzoneBucket; +import org.apache.hadoop.ozone.client.OzoneClient; +import org.apache.hadoop.ozone.client.OzoneKeyDetails; +import org.apache.hadoop.ozone.container.ContainerTestHelper; +import org.apache.hadoop.ozone.container.common.interfaces.Container; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeConfiguration; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.StateContext; +import org.apache.hadoop.ozone.container.common.volume.HddsVolume; +import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; +import org.apache.hadoop.ozone.container.common.volume.StorageVolume; +import org.apache.hadoop.ozone.dn.DatanodeTestUtils; +import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand; +import org.apache.ozone.test.GenericTestUtils; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.MethodOrderer; +import org.junit.jupiter.api.Order; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; +import org.junit.jupiter.api.TestMethodOrder; +import org.junit.jupiter.api.parallel.Execution; +import org.junit.jupiter.api.parallel.ExecutionMode; +import org.junit.jupiter.api.parallel.ResourceLock; + +/** + * Integration tests for per-volume push replication thread pools (HDDS-15412). + */ +@TestInstance(TestInstance.Lifecycle.PER_CLASS) +@TestMethodOrder(MethodOrderer.OrderAnnotation.class) +@Execution(ExecutionMode.SAME_THREAD) +@ResourceLock("MiniOzoneCluster") +class TestPerVolumePushReplication { + + private static final AtomicLong CONTAINER_ID = new AtomicLong(1_000_000L); + private static final int DATA_VOLUMES = 2; + private static final int DATANODE_COUNT = 7; + private static final String VOLUME = "vol1"; + private static final String BUCKET = "bucket1"; + private static final RatisReplicationConfig RATIS_THREE = + RatisReplicationConfig.getInstance(THREE); + private static final ECReplicationConfig EC_REP = new ECReplicationConfig(3, 2); + + private MiniOzoneCluster cluster; + private XceiverClientFactory clientFactory; + private OzoneBucket bucket; + + @BeforeAll + void setUp() throws Exception { + OzoneConfiguration conf = createSharedConfig(); + cluster = newCluster(conf, DATANODE_COUNT); + cluster.waitForClusterToBeReady(); + clientFactory = new XceiverClientManager(conf); + try (OzoneClient client = cluster.newClient()) { + bucket = DataTestUtil.createVolumeAndBucket(client, VOLUME, BUCKET); + } + } + + @AfterAll + void tearDown() { + IOUtils.closeQuietly(clientFactory, cluster); + } + + @Order(1) + @Test + void testPushAndScmReplicationWithPerVolumeEnabled() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(0); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + long containerId = createClosedContainer(clientFactory, source, 0L); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerId, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerId); + getContainer(cluster, target, containerId); + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + DataTestUtil.createKey(bucket, "pushKey1", RATIS_THREE, "data".getBytes(UTF_8)); + OzoneKeyDetails keyDetails = bucket.getKey("pushKey1"); + long scmContainerId = keyDetails.getOzoneKeyLocations().get(0).getContainerID(); + waitForContainerClose(cluster, scmContainerId); + + ContainerManager containerManager = cluster.getStorageContainerManager().getContainerManager(); + Set<ContainerReplica> replicas = + containerManager.getContainerReplicas(ContainerID.valueOf(scmContainerId)); + DatanodeDetails replicaDn = replicas.iterator().next().getDatanodeDetails(); + cluster.shutdownHddsDatanode(replicaDn); + waitForReplicaCount(scmContainerId, 3, cluster); + } + + @Order(2) + @Test + void testHealthyVolumeReplicationAfterVolumeFailure() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(1); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + MutableVolumeSet volSet = sourceDn.getDatanodeStateMachine().getContainer().getVolumeSet(); + HddsVolume vol0 = (HddsVolume) volSet.getVolumesList().get(0); + HddsVolume vol1 = (HddsVolume) volSet.getVolumesList().get(1); + + long containerOnVol0 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol0, 0L); + long containerOnVol1 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol1, 0L); + + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + triggerAndWaitForVolumeFailure(volSet, vol0); + waitForVolumePoolState(sourceDn, vol0, vol1); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerOnVol1, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerOnVol1); + getContainer(cluster, target, containerOnVol1); + assertEquals(1, volSet.getFailedVolumesList().size()); + + // Task routes via global pool fallback (HDDS-15327); replication fails on bad volume. + ReplicateContainerCommand failedVolCmd = + ReplicateContainerCommand.toTarget(containerOnVol0, target); + ReplicationSupervisor supervisor = + sourceDn.getDatanodeStateMachine().getSupervisor(); + long previousFailures = supervisor.getReplicationFailureCount(); + queuePushAndWaitForFailure(cluster, failedVolCmd, source, target, + containerOnVol0, supervisor, previousFailures); + assertTrue(supervisor.getReplicationFailureCount() >= previousFailures + 1); + + DatanodeTestUtils.restoreBadVolume(vol0); + } + + @Order(3) + @Test + void testDecommissionWithPerVolumePools() throws Exception { Review Comment: Enables per-volume pools but asserts nothing per-volume. This mirrors `TestDecommissionAndMaintenance#testNodeWithOpenPipelineCanBeDecommissionedAndRecommissioned` and would pass with the feature off. What does it catch that the existing decommission test doesn't? ########## hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestPerVolumePushReplication.java: ########## @@ -0,0 +1,545 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.ozone.container.replication; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static java.util.Collections.singleton; +import static java.util.Collections.singletonList; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_NODE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_PIPELINE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State.CLOSED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.DECOMMISSIONED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.DEAD; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DATANODE_ADMIN_MONITOR_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DEADNODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HEARTBEAT_PROCESS_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_STALENODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.getDNHostAndPort; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachHealthState; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachOpState; +import static org.apache.hadoop.hdds.scm.pipeline.MockPipeline.createPipeline; +import static org.apache.hadoop.hdds.scm.storage.ContainerProtocolCalls.createContainer; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForContainerClose; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForReplicaCount; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.List; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.BooleanSupplier; +import java.util.stream.Collectors; +import org.apache.hadoop.hdds.HddsConfigKeys; +import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.ECReplicationConfig; +import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.conf.StorageUnit; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.DatanodeID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.scm.ScmConfigKeys; +import org.apache.hadoop.hdds.scm.XceiverClientFactory; +import org.apache.hadoop.hdds.scm.XceiverClientManager; +import org.apache.hadoop.hdds.scm.XceiverClientSpi; +import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; +import org.apache.hadoop.hdds.scm.container.ContainerID; +import org.apache.hadoop.hdds.scm.container.ContainerInfo; +import org.apache.hadoop.hdds.scm.container.ContainerManager; +import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager.ReplicationManagerConfiguration; +import org.apache.hadoop.hdds.scm.node.NodeManager; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.hdds.scm.pipeline.PipelineManager; +import org.apache.hadoop.hdds.scm.server.StorageContainerManager; +import org.apache.hadoop.hdds.utils.IOUtils; +import org.apache.hadoop.ozone.DataTestUtil; +import org.apache.hadoop.ozone.HddsDatanodeService; +import org.apache.hadoop.ozone.MiniOzoneCluster; +import org.apache.hadoop.ozone.OzoneConfigKeys; +import org.apache.hadoop.ozone.UniformDatanodesFactory; +import org.apache.hadoop.ozone.client.OzoneBucket; +import org.apache.hadoop.ozone.client.OzoneClient; +import org.apache.hadoop.ozone.client.OzoneKeyDetails; +import org.apache.hadoop.ozone.container.ContainerTestHelper; +import org.apache.hadoop.ozone.container.common.interfaces.Container; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeConfiguration; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.StateContext; +import org.apache.hadoop.ozone.container.common.volume.HddsVolume; +import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; +import org.apache.hadoop.ozone.container.common.volume.StorageVolume; +import org.apache.hadoop.ozone.dn.DatanodeTestUtils; +import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand; +import org.apache.ozone.test.GenericTestUtils; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.MethodOrderer; +import org.junit.jupiter.api.Order; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; +import org.junit.jupiter.api.TestMethodOrder; +import org.junit.jupiter.api.parallel.Execution; +import org.junit.jupiter.api.parallel.ExecutionMode; +import org.junit.jupiter.api.parallel.ResourceLock; + +/** + * Integration tests for per-volume push replication thread pools (HDDS-15412). + */ +@TestInstance(TestInstance.Lifecycle.PER_CLASS) +@TestMethodOrder(MethodOrderer.OrderAnnotation.class) +@Execution(ExecutionMode.SAME_THREAD) +@ResourceLock("MiniOzoneCluster") +class TestPerVolumePushReplication { + + private static final AtomicLong CONTAINER_ID = new AtomicLong(1_000_000L); + private static final int DATA_VOLUMES = 2; + private static final int DATANODE_COUNT = 7; + private static final String VOLUME = "vol1"; + private static final String BUCKET = "bucket1"; + private static final RatisReplicationConfig RATIS_THREE = + RatisReplicationConfig.getInstance(THREE); + private static final ECReplicationConfig EC_REP = new ECReplicationConfig(3, 2); + + private MiniOzoneCluster cluster; + private XceiverClientFactory clientFactory; + private OzoneBucket bucket; + + @BeforeAll + void setUp() throws Exception { + OzoneConfiguration conf = createSharedConfig(); + cluster = newCluster(conf, DATANODE_COUNT); + cluster.waitForClusterToBeReady(); + clientFactory = new XceiverClientManager(conf); + try (OzoneClient client = cluster.newClient()) { + bucket = DataTestUtil.createVolumeAndBucket(client, VOLUME, BUCKET); + } + } + + @AfterAll + void tearDown() { + IOUtils.closeQuietly(clientFactory, cluster); + } + + @Order(1) + @Test + void testPushAndScmReplicationWithPerVolumeEnabled() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(0); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + long containerId = createClosedContainer(clientFactory, source, 0L); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerId, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerId); + getContainer(cluster, target, containerId); + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + DataTestUtil.createKey(bucket, "pushKey1", RATIS_THREE, "data".getBytes(UTF_8)); + OzoneKeyDetails keyDetails = bucket.getKey("pushKey1"); + long scmContainerId = keyDetails.getOzoneKeyLocations().get(0).getContainerID(); + waitForContainerClose(cluster, scmContainerId); + + ContainerManager containerManager = cluster.getStorageContainerManager().getContainerManager(); + Set<ContainerReplica> replicas = + containerManager.getContainerReplicas(ContainerID.valueOf(scmContainerId)); + DatanodeDetails replicaDn = replicas.iterator().next().getDatanodeDetails(); + cluster.shutdownHddsDatanode(replicaDn); + waitForReplicaCount(scmContainerId, 3, cluster); + } + + @Order(2) + @Test + void testHealthyVolumeReplicationAfterVolumeFailure() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(1); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + MutableVolumeSet volSet = sourceDn.getDatanodeStateMachine().getContainer().getVolumeSet(); + HddsVolume vol0 = (HddsVolume) volSet.getVolumesList().get(0); + HddsVolume vol1 = (HddsVolume) volSet.getVolumesList().get(1); + + long containerOnVol0 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol0, 0L); + long containerOnVol1 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol1, 0L); + + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + triggerAndWaitForVolumeFailure(volSet, vol0); + waitForVolumePoolState(sourceDn, vol0, vol1); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerOnVol1, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerOnVol1); + getContainer(cluster, target, containerOnVol1); + assertEquals(1, volSet.getFailedVolumesList().size()); + + // Task routes via global pool fallback (HDDS-15327); replication fails on bad volume. + ReplicateContainerCommand failedVolCmd = + ReplicateContainerCommand.toTarget(containerOnVol0, target); + ReplicationSupervisor supervisor = + sourceDn.getDatanodeStateMachine().getSupervisor(); + long previousFailures = supervisor.getReplicationFailureCount(); + queuePushAndWaitForFailure(cluster, failedVolCmd, source, target, + containerOnVol0, supervisor, previousFailures); + assertTrue(supervisor.getReplicationFailureCount() >= previousFailures + 1); + + DatanodeTestUtils.restoreBadVolume(vol0); + } + + @Order(3) + @Test + void testDecommissionWithPerVolumePools() throws Exception { + try (ContainerOperationClient scmClient = new ContainerOperationClient(cluster.getConf())) { + StorageContainerManager scm = cluster.getStorageContainerManager(); + NodeManager nm = scm.getScmNodeManager(); + ContainerManager cm = scm.getContainerManager(); + PipelineManager pm = scm.getPipelineManager(); + + generateData(bucket, 20, "decomKey", RATIS_THREE); + generateData(bucket, 20, "decomEcKey", EC_REP); + + ContainerInfo ratisContainer = waitForKeyContainer(bucket, cm, "decomKey0", 3); + ContainerInfo ecContainer = waitForKeyContainer(bucket, cm, "decomEcKey0", 5); + Pipeline ratisPipeline = pm.getPipeline(ratisContainer.getPipelineID()); + Pipeline ecPipeline = pm.getPipeline(ecContainer.getPipelineID()); + + DatanodeID dnId = ratisPipeline.getNodes().stream() + .filter(node -> ecPipeline.getNodes().contains(node)) + .findFirst() + .orElseThrow(() -> new AssertionError("no intersecting datanode found")) + .getID(); + DatanodeDetails toDecommission = nm.getNode(dnId); + + scmClient.decommissionNodes(singletonList(getDNHostAndPort(toDecommission)), false); + waitForDnToReachOpState(nm, toDecommission, DECOMMISSIONED); + + waitForContainerReplicas(cm, ratisContainer, 4); + waitForContainerReplicas(cm, ecContainer, 6); + + cluster.shutdownHddsDatanode(toDecommission); + waitForDnToReachHealthState(nm, toDecommission, DEAD); + + waitForContainerReplicas(cm, ratisContainer, 3); + waitForContainerReplicas(cm, ecContainer, 5); + + DataTestUtil.createKey(bucket, "sanityKey", RATIS_THREE, + "still healthy".getBytes(StandardCharsets.UTF_8)); + } + } + + private HddsDatanodeService selectHealthyDatanode(int indexAmongHealthy) { + List<HddsDatanodeService> healthy = cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .collect(Collectors.toList()); + if (indexAmongHealthy >= healthy.size()) { + throw new AssertionError("not enough healthy datanodes: requested index " + + indexAmongHealthy + ", found " + healthy.size()); + } + return healthy.get(indexAmongHealthy); + } + + private DatanodeDetails selectOtherHealthyNode(DatanodeDetails source) { + return cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .map(HddsDatanodeService::getDatanodeDetails) + .filter(dn -> !dn.equals(source)) + .findAny() + .orElseThrow(() -> new AssertionError("no target datanode found")); + } + + private boolean isHealthyDatanode(HddsDatanodeService datanode) { Review Comment: This filters on op-state and failed-volume but not liveness, and `shutdownHddsDatanode` (L174, L247) leaves the node `IN_SERVICE` in `getHddsDatanodes()`. With the `@Order`ed shared cluster, could a later test pick a stopped node here and time out? Would a separate cluster per destructive test be safer? ########## hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestPerVolumePushReplication.java: ########## @@ -0,0 +1,545 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.ozone.container.replication; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static java.util.Collections.singleton; +import static java.util.Collections.singletonList; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_NODE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_PIPELINE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State.CLOSED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.DECOMMISSIONED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.DEAD; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DATANODE_ADMIN_MONITOR_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DEADNODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HEARTBEAT_PROCESS_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_STALENODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.getDNHostAndPort; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachHealthState; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachOpState; +import static org.apache.hadoop.hdds.scm.pipeline.MockPipeline.createPipeline; +import static org.apache.hadoop.hdds.scm.storage.ContainerProtocolCalls.createContainer; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForContainerClose; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForReplicaCount; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.List; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.BooleanSupplier; +import java.util.stream.Collectors; +import org.apache.hadoop.hdds.HddsConfigKeys; +import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.ECReplicationConfig; +import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.conf.StorageUnit; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.DatanodeID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.scm.ScmConfigKeys; +import org.apache.hadoop.hdds.scm.XceiverClientFactory; +import org.apache.hadoop.hdds.scm.XceiverClientManager; +import org.apache.hadoop.hdds.scm.XceiverClientSpi; +import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; +import org.apache.hadoop.hdds.scm.container.ContainerID; +import org.apache.hadoop.hdds.scm.container.ContainerInfo; +import org.apache.hadoop.hdds.scm.container.ContainerManager; +import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager.ReplicationManagerConfiguration; +import org.apache.hadoop.hdds.scm.node.NodeManager; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.hdds.scm.pipeline.PipelineManager; +import org.apache.hadoop.hdds.scm.server.StorageContainerManager; +import org.apache.hadoop.hdds.utils.IOUtils; +import org.apache.hadoop.ozone.DataTestUtil; +import org.apache.hadoop.ozone.HddsDatanodeService; +import org.apache.hadoop.ozone.MiniOzoneCluster; +import org.apache.hadoop.ozone.OzoneConfigKeys; +import org.apache.hadoop.ozone.UniformDatanodesFactory; +import org.apache.hadoop.ozone.client.OzoneBucket; +import org.apache.hadoop.ozone.client.OzoneClient; +import org.apache.hadoop.ozone.client.OzoneKeyDetails; +import org.apache.hadoop.ozone.container.ContainerTestHelper; +import org.apache.hadoop.ozone.container.common.interfaces.Container; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeConfiguration; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.StateContext; +import org.apache.hadoop.ozone.container.common.volume.HddsVolume; +import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; +import org.apache.hadoop.ozone.container.common.volume.StorageVolume; +import org.apache.hadoop.ozone.dn.DatanodeTestUtils; +import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand; +import org.apache.ozone.test.GenericTestUtils; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.MethodOrderer; +import org.junit.jupiter.api.Order; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; +import org.junit.jupiter.api.TestMethodOrder; +import org.junit.jupiter.api.parallel.Execution; +import org.junit.jupiter.api.parallel.ExecutionMode; +import org.junit.jupiter.api.parallel.ResourceLock; + +/** + * Integration tests for per-volume push replication thread pools (HDDS-15412). + */ +@TestInstance(TestInstance.Lifecycle.PER_CLASS) +@TestMethodOrder(MethodOrderer.OrderAnnotation.class) +@Execution(ExecutionMode.SAME_THREAD) +@ResourceLock("MiniOzoneCluster") +class TestPerVolumePushReplication { + + private static final AtomicLong CONTAINER_ID = new AtomicLong(1_000_000L); + private static final int DATA_VOLUMES = 2; + private static final int DATANODE_COUNT = 7; + private static final String VOLUME = "vol1"; + private static final String BUCKET = "bucket1"; + private static final RatisReplicationConfig RATIS_THREE = + RatisReplicationConfig.getInstance(THREE); + private static final ECReplicationConfig EC_REP = new ECReplicationConfig(3, 2); + + private MiniOzoneCluster cluster; + private XceiverClientFactory clientFactory; + private OzoneBucket bucket; + + @BeforeAll + void setUp() throws Exception { + OzoneConfiguration conf = createSharedConfig(); + cluster = newCluster(conf, DATANODE_COUNT); + cluster.waitForClusterToBeReady(); + clientFactory = new XceiverClientManager(conf); + try (OzoneClient client = cluster.newClient()) { + bucket = DataTestUtil.createVolumeAndBucket(client, VOLUME, BUCKET); + } + } + + @AfterAll + void tearDown() { + IOUtils.closeQuietly(clientFactory, cluster); + } + + @Order(1) + @Test + void testPushAndScmReplicationWithPerVolumeEnabled() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(0); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + long containerId = createClosedContainer(clientFactory, source, 0L); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerId, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerId); + getContainer(cluster, target, containerId); + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + DataTestUtil.createKey(bucket, "pushKey1", RATIS_THREE, "data".getBytes(UTF_8)); + OzoneKeyDetails keyDetails = bucket.getKey("pushKey1"); + long scmContainerId = keyDetails.getOzoneKeyLocations().get(0).getContainerID(); + waitForContainerClose(cluster, scmContainerId); + + ContainerManager containerManager = cluster.getStorageContainerManager().getContainerManager(); + Set<ContainerReplica> replicas = + containerManager.getContainerReplicas(ContainerID.valueOf(scmContainerId)); + DatanodeDetails replicaDn = replicas.iterator().next().getDatanodeDetails(); + cluster.shutdownHddsDatanode(replicaDn); + waitForReplicaCount(scmContainerId, 3, cluster); + } + + @Order(2) + @Test + void testHealthyVolumeReplicationAfterVolumeFailure() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(1); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + MutableVolumeSet volSet = sourceDn.getDatanodeStateMachine().getContainer().getVolumeSet(); + HddsVolume vol0 = (HddsVolume) volSet.getVolumesList().get(0); + HddsVolume vol1 = (HddsVolume) volSet.getVolumesList().get(1); + + long containerOnVol0 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol0, 0L); + long containerOnVol1 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol1, 0L); + + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + triggerAndWaitForVolumeFailure(volSet, vol0); + waitForVolumePoolState(sourceDn, vol0, vol1); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerOnVol1, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerOnVol1); + getContainer(cluster, target, containerOnVol1); + assertEquals(1, volSet.getFailedVolumesList().size()); + + // Task routes via global pool fallback (HDDS-15327); replication fails on bad volume. + ReplicateContainerCommand failedVolCmd = + ReplicateContainerCommand.toTarget(containerOnVol0, target); + ReplicationSupervisor supervisor = + sourceDn.getDatanodeStateMachine().getSupervisor(); + long previousFailures = supervisor.getReplicationFailureCount(); + queuePushAndWaitForFailure(cluster, failedVolCmd, source, target, + containerOnVol0, supervisor, previousFailures); + assertTrue(supervisor.getReplicationFailureCount() >= previousFailures + 1); + + DatanodeTestUtils.restoreBadVolume(vol0); + } + + @Order(3) + @Test + void testDecommissionWithPerVolumePools() throws Exception { + try (ContainerOperationClient scmClient = new ContainerOperationClient(cluster.getConf())) { + StorageContainerManager scm = cluster.getStorageContainerManager(); + NodeManager nm = scm.getScmNodeManager(); + ContainerManager cm = scm.getContainerManager(); + PipelineManager pm = scm.getPipelineManager(); + + generateData(bucket, 20, "decomKey", RATIS_THREE); + generateData(bucket, 20, "decomEcKey", EC_REP); + + ContainerInfo ratisContainer = waitForKeyContainer(bucket, cm, "decomKey0", 3); + ContainerInfo ecContainer = waitForKeyContainer(bucket, cm, "decomEcKey0", 5); + Pipeline ratisPipeline = pm.getPipeline(ratisContainer.getPipelineID()); + Pipeline ecPipeline = pm.getPipeline(ecContainer.getPipelineID()); + + DatanodeID dnId = ratisPipeline.getNodes().stream() + .filter(node -> ecPipeline.getNodes().contains(node)) + .findFirst() + .orElseThrow(() -> new AssertionError("no intersecting datanode found")) + .getID(); + DatanodeDetails toDecommission = nm.getNode(dnId); + + scmClient.decommissionNodes(singletonList(getDNHostAndPort(toDecommission)), false); + waitForDnToReachOpState(nm, toDecommission, DECOMMISSIONED); + + waitForContainerReplicas(cm, ratisContainer, 4); + waitForContainerReplicas(cm, ecContainer, 6); + + cluster.shutdownHddsDatanode(toDecommission); + waitForDnToReachHealthState(nm, toDecommission, DEAD); + + waitForContainerReplicas(cm, ratisContainer, 3); + waitForContainerReplicas(cm, ecContainer, 5); + + DataTestUtil.createKey(bucket, "sanityKey", RATIS_THREE, + "still healthy".getBytes(StandardCharsets.UTF_8)); + } + } + + private HddsDatanodeService selectHealthyDatanode(int indexAmongHealthy) { + List<HddsDatanodeService> healthy = cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .collect(Collectors.toList()); + if (indexAmongHealthy >= healthy.size()) { + throw new AssertionError("not enough healthy datanodes: requested index " + + indexAmongHealthy + ", found " + healthy.size()); + } + return healthy.get(indexAmongHealthy); + } + + private DatanodeDetails selectOtherHealthyNode(DatanodeDetails source) { + return cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .map(HddsDatanodeService::getDatanodeDetails) + .filter(dn -> !dn.equals(source)) + .findAny() + .orElseThrow(() -> new AssertionError("no target datanode found")); + } + + private boolean isHealthyDatanode(HddsDatanodeService datanode) { + if (datanode.getDatanodeDetails().getPersistedOpState() != IN_SERVICE) { + return false; + } + MutableVolumeSet volumeSet = + datanode.getDatanodeStateMachine().getContainer().getVolumeSet(); + return volumeSet.getFailedVolumesList().isEmpty() + && volumeSet.getVolumesList().size() == DATA_VOLUMES; + } + + private static void assertVolumePools(HddsDatanodeService dn, + int expectedVolumeCount, int expectedPoolSize) { + VolumeReplicationThreadPools pools = + dn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + assertNotNull(pools); + List<? extends StorageVolume> volumes = + dn.getDatanodeStateMachine().getContainer().getVolumeSet().getVolumesList(); + assertEquals(expectedVolumeCount, volumes.size()); + for (StorageVolume volume : volumes) { + String volumeRoot = volume.getStorageDir().getPath(); + assertTrue(pools.hasPool(volumeRoot), "missing pool for " + volumeRoot); + assertEquals(expectedPoolSize, pools.getPoolSize(volumeRoot)); + } + } + + private static void waitForVolumePoolState(HddsDatanodeService sourceDn, + HddsVolume failedVolume, HddsVolume healthyVolume) + throws TimeoutException, InterruptedException { + String failedPath = failedVolume.getStorageDir().getPath(); + String healthyPath = healthyVolume.getStorageDir().getPath(); + GenericTestUtils.waitFor((BooleanSupplier) () -> { + VolumeReplicationThreadPools pools = + sourceDn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + return pools != null + && !pools.hasPool(failedPath) + && pools.hasPool(healthyPath); + }, 100, 60000); + VolumeReplicationThreadPools pools = + sourceDn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + assertNotNull(pools); + assertFalse(pools.hasPool(failedPath)); + assertTrue(pools.hasPool(healthyPath)); + } + + private static void queuePushAndWaitForContainer(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source, + DatanodeDetails target, long containerId) + throws IOException, InterruptedException, TimeoutException { + queueReplicationCommand(cluster, cmd, source); + GenericTestUtils.waitFor((BooleanSupplier) () -> + hasContainer(cluster, target, containerId), 100, 30000); + } + + private static void queuePushAndWaitForFailure(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source, + DatanodeDetails target, long containerId, ReplicationSupervisor supervisor, + long previousFailureCount) + throws IOException, InterruptedException, TimeoutException { + queueReplicationCommand(cluster, cmd, source); + GenericTestUtils.waitFor((BooleanSupplier) () -> + supervisor.getReplicationFailureCount() >= previousFailureCount + 1 + && !hasContainer(cluster, target, containerId), + 100, 30000); + } + + private static void queueReplicationCommand(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source) throws IOException { + DatanodeStateMachine stateMachine = cluster.getHddsDatanode(source).getDatanodeStateMachine(); + StateContext context = stateMachine.getContext(); + context.getTermOfLeaderSCM().ifPresent(cmd::setTerm); Review Comment: Identical to `TestContainerReplication#queueAndWaitForCompletion` (L267-268), and that same-package class also has `CONTAINER_ID` and `createNewClosedContainer`. Could we reuse those instead of re-deriving? ########## hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestPerVolumePushReplication.java: ########## @@ -0,0 +1,545 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.hadoop.ozone.container.replication; + +import static java.nio.charset.StandardCharsets.UTF_8; +import static java.util.Collections.singleton; +import static java.util.Collections.singletonList; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_CONTAINER_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_NODE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_PIPELINE_REPORT_INTERVAL; +import static org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos.ContainerDataProto.State.CLOSED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.DECOMMISSIONED; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState.IN_SERVICE; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeState.DEAD; +import static org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor.THREE; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DATANODE_ADMIN_MONITOR_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_DEADNODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_HEARTBEAT_PROCESS_INTERVAL; +import static org.apache.hadoop.hdds.scm.ScmConfigKeys.OZONE_SCM_STALENODE_INTERVAL; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.getDNHostAndPort; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachHealthState; +import static org.apache.hadoop.hdds.scm.node.NodeTestUtil.waitForDnToReachOpState; +import static org.apache.hadoop.hdds.scm.pipeline.MockPipeline.createPipeline; +import static org.apache.hadoop.hdds.scm.storage.ContainerProtocolCalls.createContainer; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForContainerClose; +import static org.apache.hadoop.ozone.container.OzoneTestHelper.waitForReplicaCount; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertNotNull; +import static org.junit.jupiter.api.Assertions.assertTrue; + +import java.io.IOException; +import java.nio.charset.StandardCharsets; +import java.time.Duration; +import java.util.List; +import java.util.Set; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; +import java.util.concurrent.atomic.AtomicLong; +import java.util.function.BooleanSupplier; +import java.util.stream.Collectors; +import org.apache.hadoop.hdds.HddsConfigKeys; +import org.apache.hadoop.hdds.client.BlockID; +import org.apache.hadoop.hdds.client.ECReplicationConfig; +import org.apache.hadoop.hdds.client.RatisReplicationConfig; +import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.hdds.conf.OzoneConfiguration; +import org.apache.hadoop.hdds.conf.StorageUnit; +import org.apache.hadoop.hdds.protocol.DatanodeDetails; +import org.apache.hadoop.hdds.protocol.DatanodeID; +import org.apache.hadoop.hdds.protocol.datanode.proto.ContainerProtos; +import org.apache.hadoop.hdds.scm.ScmConfigKeys; +import org.apache.hadoop.hdds.scm.XceiverClientFactory; +import org.apache.hadoop.hdds.scm.XceiverClientManager; +import org.apache.hadoop.hdds.scm.XceiverClientSpi; +import org.apache.hadoop.hdds.scm.cli.ContainerOperationClient; +import org.apache.hadoop.hdds.scm.container.ContainerID; +import org.apache.hadoop.hdds.scm.container.ContainerInfo; +import org.apache.hadoop.hdds.scm.container.ContainerManager; +import org.apache.hadoop.hdds.scm.container.ContainerReplica; +import org.apache.hadoop.hdds.scm.container.replication.ReplicationManager.ReplicationManagerConfiguration; +import org.apache.hadoop.hdds.scm.node.NodeManager; +import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.hdds.scm.pipeline.PipelineManager; +import org.apache.hadoop.hdds.scm.server.StorageContainerManager; +import org.apache.hadoop.hdds.utils.IOUtils; +import org.apache.hadoop.ozone.DataTestUtil; +import org.apache.hadoop.ozone.HddsDatanodeService; +import org.apache.hadoop.ozone.MiniOzoneCluster; +import org.apache.hadoop.ozone.OzoneConfigKeys; +import org.apache.hadoop.ozone.UniformDatanodesFactory; +import org.apache.hadoop.ozone.client.OzoneBucket; +import org.apache.hadoop.ozone.client.OzoneClient; +import org.apache.hadoop.ozone.client.OzoneKeyDetails; +import org.apache.hadoop.ozone.container.ContainerTestHelper; +import org.apache.hadoop.ozone.container.common.interfaces.Container; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeConfiguration; +import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; +import org.apache.hadoop.ozone.container.common.statemachine.StateContext; +import org.apache.hadoop.ozone.container.common.volume.HddsVolume; +import org.apache.hadoop.ozone.container.common.volume.MutableVolumeSet; +import org.apache.hadoop.ozone.container.common.volume.StorageVolume; +import org.apache.hadoop.ozone.dn.DatanodeTestUtils; +import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand; +import org.apache.ozone.test.GenericTestUtils; +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.MethodOrderer; +import org.junit.jupiter.api.Order; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.TestInstance; +import org.junit.jupiter.api.TestMethodOrder; +import org.junit.jupiter.api.parallel.Execution; +import org.junit.jupiter.api.parallel.ExecutionMode; +import org.junit.jupiter.api.parallel.ResourceLock; + +/** + * Integration tests for per-volume push replication thread pools (HDDS-15412). + */ +@TestInstance(TestInstance.Lifecycle.PER_CLASS) +@TestMethodOrder(MethodOrderer.OrderAnnotation.class) +@Execution(ExecutionMode.SAME_THREAD) +@ResourceLock("MiniOzoneCluster") +class TestPerVolumePushReplication { + + private static final AtomicLong CONTAINER_ID = new AtomicLong(1_000_000L); + private static final int DATA_VOLUMES = 2; + private static final int DATANODE_COUNT = 7; + private static final String VOLUME = "vol1"; + private static final String BUCKET = "bucket1"; + private static final RatisReplicationConfig RATIS_THREE = + RatisReplicationConfig.getInstance(THREE); + private static final ECReplicationConfig EC_REP = new ECReplicationConfig(3, 2); + + private MiniOzoneCluster cluster; + private XceiverClientFactory clientFactory; + private OzoneBucket bucket; + + @BeforeAll + void setUp() throws Exception { + OzoneConfiguration conf = createSharedConfig(); + cluster = newCluster(conf, DATANODE_COUNT); + cluster.waitForClusterToBeReady(); + clientFactory = new XceiverClientManager(conf); + try (OzoneClient client = cluster.newClient()) { + bucket = DataTestUtil.createVolumeAndBucket(client, VOLUME, BUCKET); + } + } + + @AfterAll + void tearDown() { + IOUtils.closeQuietly(clientFactory, cluster); + } + + @Order(1) + @Test + void testPushAndScmReplicationWithPerVolumeEnabled() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(0); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + long containerId = createClosedContainer(clientFactory, source, 0L); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerId, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerId); + getContainer(cluster, target, containerId); + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + DataTestUtil.createKey(bucket, "pushKey1", RATIS_THREE, "data".getBytes(UTF_8)); + OzoneKeyDetails keyDetails = bucket.getKey("pushKey1"); + long scmContainerId = keyDetails.getOzoneKeyLocations().get(0).getContainerID(); + waitForContainerClose(cluster, scmContainerId); + + ContainerManager containerManager = cluster.getStorageContainerManager().getContainerManager(); + Set<ContainerReplica> replicas = + containerManager.getContainerReplicas(ContainerID.valueOf(scmContainerId)); + DatanodeDetails replicaDn = replicas.iterator().next().getDatanodeDetails(); + cluster.shutdownHddsDatanode(replicaDn); + waitForReplicaCount(scmContainerId, 3, cluster); + } + + @Order(2) + @Test + void testHealthyVolumeReplicationAfterVolumeFailure() throws Exception { + HddsDatanodeService sourceDn = selectHealthyDatanode(1); + DatanodeDetails source = sourceDn.getDatanodeDetails(); + DatanodeDetails target = selectOtherHealthyNode(source); + MutableVolumeSet volSet = sourceDn.getDatanodeStateMachine().getContainer().getVolumeSet(); + HddsVolume vol0 = (HddsVolume) volSet.getVolumesList().get(0); + HddsVolume vol1 = (HddsVolume) volSet.getVolumesList().get(1); + + long containerOnVol0 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol0, 0L); + long containerOnVol1 = findOrCreateContainerOnVolume( + cluster, clientFactory, source, vol1, 0L); + + assertVolumePools(sourceDn, DATA_VOLUMES, 1); + + triggerAndWaitForVolumeFailure(volSet, vol0); + waitForVolumePoolState(sourceDn, vol0, vol1); + + ReplicateContainerCommand cmd = + ReplicateContainerCommand.toTarget(containerOnVol1, target); + queuePushAndWaitForContainer(cluster, cmd, source, target, containerOnVol1); + getContainer(cluster, target, containerOnVol1); + assertEquals(1, volSet.getFailedVolumesList().size()); + + // Task routes via global pool fallback (HDDS-15327); replication fails on bad volume. + ReplicateContainerCommand failedVolCmd = + ReplicateContainerCommand.toTarget(containerOnVol0, target); + ReplicationSupervisor supervisor = + sourceDn.getDatanodeStateMachine().getSupervisor(); + long previousFailures = supervisor.getReplicationFailureCount(); + queuePushAndWaitForFailure(cluster, failedVolCmd, source, target, + containerOnVol0, supervisor, previousFailures); + assertTrue(supervisor.getReplicationFailureCount() >= previousFailures + 1); + + DatanodeTestUtils.restoreBadVolume(vol0); + } + + @Order(3) + @Test + void testDecommissionWithPerVolumePools() throws Exception { + try (ContainerOperationClient scmClient = new ContainerOperationClient(cluster.getConf())) { + StorageContainerManager scm = cluster.getStorageContainerManager(); + NodeManager nm = scm.getScmNodeManager(); + ContainerManager cm = scm.getContainerManager(); + PipelineManager pm = scm.getPipelineManager(); + + generateData(bucket, 20, "decomKey", RATIS_THREE); + generateData(bucket, 20, "decomEcKey", EC_REP); + + ContainerInfo ratisContainer = waitForKeyContainer(bucket, cm, "decomKey0", 3); + ContainerInfo ecContainer = waitForKeyContainer(bucket, cm, "decomEcKey0", 5); + Pipeline ratisPipeline = pm.getPipeline(ratisContainer.getPipelineID()); + Pipeline ecPipeline = pm.getPipeline(ecContainer.getPipelineID()); + + DatanodeID dnId = ratisPipeline.getNodes().stream() + .filter(node -> ecPipeline.getNodes().contains(node)) + .findFirst() + .orElseThrow(() -> new AssertionError("no intersecting datanode found")) + .getID(); + DatanodeDetails toDecommission = nm.getNode(dnId); + + scmClient.decommissionNodes(singletonList(getDNHostAndPort(toDecommission)), false); + waitForDnToReachOpState(nm, toDecommission, DECOMMISSIONED); + + waitForContainerReplicas(cm, ratisContainer, 4); + waitForContainerReplicas(cm, ecContainer, 6); + + cluster.shutdownHddsDatanode(toDecommission); + waitForDnToReachHealthState(nm, toDecommission, DEAD); + + waitForContainerReplicas(cm, ratisContainer, 3); + waitForContainerReplicas(cm, ecContainer, 5); + + DataTestUtil.createKey(bucket, "sanityKey", RATIS_THREE, + "still healthy".getBytes(StandardCharsets.UTF_8)); + } + } + + private HddsDatanodeService selectHealthyDatanode(int indexAmongHealthy) { + List<HddsDatanodeService> healthy = cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .collect(Collectors.toList()); + if (indexAmongHealthy >= healthy.size()) { + throw new AssertionError("not enough healthy datanodes: requested index " + + indexAmongHealthy + ", found " + healthy.size()); + } + return healthy.get(indexAmongHealthy); + } + + private DatanodeDetails selectOtherHealthyNode(DatanodeDetails source) { + return cluster.getHddsDatanodes().stream() + .filter(this::isHealthyDatanode) + .map(HddsDatanodeService::getDatanodeDetails) + .filter(dn -> !dn.equals(source)) + .findAny() + .orElseThrow(() -> new AssertionError("no target datanode found")); + } + + private boolean isHealthyDatanode(HddsDatanodeService datanode) { + if (datanode.getDatanodeDetails().getPersistedOpState() != IN_SERVICE) { + return false; + } + MutableVolumeSet volumeSet = + datanode.getDatanodeStateMachine().getContainer().getVolumeSet(); + return volumeSet.getFailedVolumesList().isEmpty() + && volumeSet.getVolumesList().size() == DATA_VOLUMES; + } + + private static void assertVolumePools(HddsDatanodeService dn, + int expectedVolumeCount, int expectedPoolSize) { + VolumeReplicationThreadPools pools = + dn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + assertNotNull(pools); + List<? extends StorageVolume> volumes = + dn.getDatanodeStateMachine().getContainer().getVolumeSet().getVolumesList(); + assertEquals(expectedVolumeCount, volumes.size()); + for (StorageVolume volume : volumes) { + String volumeRoot = volume.getStorageDir().getPath(); + assertTrue(pools.hasPool(volumeRoot), "missing pool for " + volumeRoot); + assertEquals(expectedPoolSize, pools.getPoolSize(volumeRoot)); + } + } + + private static void waitForVolumePoolState(HddsDatanodeService sourceDn, + HddsVolume failedVolume, HddsVolume healthyVolume) + throws TimeoutException, InterruptedException { + String failedPath = failedVolume.getStorageDir().getPath(); + String healthyPath = healthyVolume.getStorageDir().getPath(); + GenericTestUtils.waitFor((BooleanSupplier) () -> { + VolumeReplicationThreadPools pools = + sourceDn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + return pools != null + && !pools.hasPool(failedPath) + && pools.hasPool(healthyPath); + }, 100, 60000); + VolumeReplicationThreadPools pools = + sourceDn.getDatanodeStateMachine().getSupervisor().getVolumeReplicationThreadPools(); + assertNotNull(pools); + assertFalse(pools.hasPool(failedPath)); + assertTrue(pools.hasPool(healthyPath)); + } + + private static void queuePushAndWaitForContainer(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source, + DatanodeDetails target, long containerId) + throws IOException, InterruptedException, TimeoutException { + queueReplicationCommand(cluster, cmd, source); + GenericTestUtils.waitFor((BooleanSupplier) () -> + hasContainer(cluster, target, containerId), 100, 30000); + } + + private static void queuePushAndWaitForFailure(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source, + DatanodeDetails target, long containerId, ReplicationSupervisor supervisor, + long previousFailureCount) + throws IOException, InterruptedException, TimeoutException { + queueReplicationCommand(cluster, cmd, source); + GenericTestUtils.waitFor((BooleanSupplier) () -> + supervisor.getReplicationFailureCount() >= previousFailureCount + 1 + && !hasContainer(cluster, target, containerId), + 100, 30000); + } + + private static void queueReplicationCommand(MiniOzoneCluster cluster, + ReplicateContainerCommand cmd, DatanodeDetails source) throws IOException { + DatanodeStateMachine stateMachine = cluster.getHddsDatanode(source).getDatanodeStateMachine(); + StateContext context = stateMachine.getContext(); + context.getTermOfLeaderSCM().ifPresent(cmd::setTerm); + context.addCommand(cmd); + } + + private static boolean hasContainer(MiniOzoneCluster cluster, + DatanodeDetails datanode, long containerId) { + try { + return cluster.getHddsDatanode(datanode).getDatanodeStateMachine().getContainer() + .getContainerSet().getContainer(containerId) != null; + } catch (IOException e) { + return false; + } + } + + private static long findOrCreateContainerOnVolume(MiniOzoneCluster cluster, + XceiverClientFactory clientFactory, DatanodeDetails dn, HddsVolume targetVolume, + long dataSize) throws Exception { + for (int attempt = 0; attempt < 30; attempt++) { + long containerId = createClosedContainer(clientFactory, dn, dataSize); + Container<?> container = getContainer(cluster, dn, containerId); + if (targetVolume.equals(container.getContainerData().getVolume())) { + return containerId; + } + } + throw new AssertionError("Could not place container on volume " + targetVolume); + } + + private static long createClosedContainer(XceiverClientFactory clientFactory, + DatanodeDetails dn, long dataSize) throws Exception { + long containerId = CONTAINER_ID.incrementAndGet(); + try (XceiverClientSpi client = clientFactory.acquireClient(createPipeline(singleton(dn)))) { + if (dataSize <= 0) { + createContainer(client, containerId, null, CLOSED, 0); + return containerId; + } + createContainer(client, containerId, null); Review Comment: This branch looks unreachable, since all callers (L157, L188, L190) pass `dataSize=0`. Removing it makes this identical to `TestContainerReplication#createNewClosedContainer`. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
