smengcl commented on code in PR #11146:
URL: https://github.com/apache/ozone/pull/11146#discussion_r4100708258


##########
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:
   Yes, it could, and it was reachable today. Addressed in 66bf80e1221.
   
   `selectHealthyDatanode` now rejects stopped nodes via 
`HddsDatanodeService#isStopped`, so a leaked shutdown fails on selection with a 
readable message instead of timing out inside a later wait. On top of that the 
first test restarts the datanode it stops, and `restoreBadVolume` moved into a 
`finally` block, so neither destructive test can leave state behind even when 
an assertion fails partway through. The class javadoc now states that contract.
   
   On a separate cluster per destructive test: I kept the shared cluster. Three 
`MiniOzoneCluster` lifecycles at 7 datanodes each is a large runtime cost for 
this class, and the leaked state was the actual defect rather than the sharing. 
If you would still rather have the isolation, say so and I will switch to 
`MiniOzoneClusterProvider`.
   



-- 
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]

Reply via email to