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


##########
hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestPerVolumePushReplication.java:
##########
@@ -0,0 +1,547 @@
+/*
+ * 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.emptySet;
+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.assertj.core.api.Assertions.assertThat;
+
+import java.io.IOException;
+import java.time.Duration;
+import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
+import java.util.concurrent.atomic.AtomicLong;
+import java.util.stream.Collectors;
+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.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.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.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;
+
+/**
+ * Integration tests for per-volume push replication thread pools (HDDS-15412).
+ *
+ * <p>The tests share one cluster and run in {@link Order} sequence, so each 
one restores whatever it broke
+ * (stopped datanode, failed volume) before returning. {@link 
#selectHealthyDatanode} additionally filters out
+ * stopped datanodes so a leaked shutdown cannot silently hand a later test a 
dead node.
+ */
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@TestMethodOrder(MethodOrderer.OrderAnnotation.class)
+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 int PER_VOLUME_STREAMS = 1;
+  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 OzoneClient client;
+  private OzoneBucket bucket;
+
+  @BeforeAll
+  void setUp() throws Exception {
+    OzoneConfiguration conf = createConfig();
+    cluster = newCluster(conf, DATANODE_COUNT);
+    cluster.waitForClusterToBeReady();
+    clientFactory = new XceiverClientManager(conf);
+    // Keep the client open for the lifetime of the class: the bucket handle 
below delegates to it, and all
+    // three tests write keys through that handle.
+    client = cluster.newClient();
+    bucket = DataTestUtil.createVolumeAndBucket(client, VOLUME, BUCKET);
+  }
+
+  @AfterAll
+  void tearDown() {
+    IOUtils.closeQuietly(client, 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);
+
+    assertVolumePools(sourceDn, DATA_VOLUMES, PER_VOLUME_STREAMS);
+
+    // Assert the push was dispatched to the pool of the volume holding the 
container, not the global pool.
+    // Without this the test would also pass with 
hdds.datanode.replication.per.volume.enabled=false.
+    HddsVolume containerVolume = getContainer(cluster, source, 
containerId).getContainerData().getVolume();
+    ThreadPoolExecutor volumePool = volumePoolOf(sourceDn, containerVolume);
+    long completedBefore = volumePool.getCompletedTaskCount();
+
+    ReplicateContainerCommand cmd = 
ReplicateContainerCommand.toTarget(containerId, target);
+    queuePushAndWaitForContainer(cluster, cmd, source, target, containerId);
+    GenericTestUtils.waitFor(() -> volumePool.getCompletedTaskCount() > 
completedBefore, 100, 30000);
+
+    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);
+
+    // SCM learns about replicas asynchronously via ICR, so settle on the full 
replica set before picking one
+    // to stop; otherwise the iterator below can be empty.
+    StorageContainerManager scm = cluster.getStorageContainerManager();
+    ContainerManager containerManager = scm.getContainerManager();
+    ContainerID scmContainer = ContainerID.valueOf(scmContainerId);
+    waitForReplicas(containerManager, scmContainer, 3);
+    Set<ContainerReplica> replicas = 
containerManager.getContainerReplicas(scmContainer);
+    DatanodeDetails replicaDn = 
replicas.iterator().next().getDatanodeDetails();
+
+    cluster.shutdownHddsDatanode(replicaDn);
+    try {
+      // SCM drops a dead node's replicas only once it is declared DEAD, so 
waiting for a count of 3 right after
+      // the shutdown would be satisfied by the stale replica on its very 
first poll. Require 3 replicas none of
+      // which sit on the stopped node, which only re-replication can produce.
+      waitForDnToReachHealthState(scm.getScmNodeManager(), replicaDn, DEAD);
+      GenericTestUtils.waitFor(() -> {
+        Set<ContainerReplica> current = replicasOf(containerManager, 
scmContainer);
+        return current.size() == 3
+            && current.stream().noneMatch(r -> 
r.getDatanodeDetails().equals(replicaDn));
+      }, 500, 60000);
+    } finally {
+      // Restore the cluster so the later tests see all DATANODE_COUNT nodes.
+      cluster.restartHddsDatanode(replicaDn, true);
+    }
+  }
+
+  @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);
+    long containerOnVol1 = findOrCreateContainerOnVolume(cluster, 
clientFactory, source, vol1);
+
+    assertVolumePools(sourceDn, DATA_VOLUMES, PER_VOLUME_STREAMS);
+
+    try {
+      triggerAndWaitForVolumeFailure(volSet, vol0);
+      waitForVolumePoolState(sourceDn, vol0, vol1);
+
+      // The healthy volume keeps its own pool and keeps serving pushes from 
it.
+      ThreadPoolExecutor healthyPool = volumePoolOf(sourceDn, vol1);
+      long completedBefore = healthyPool.getCompletedTaskCount();
+      ReplicateContainerCommand cmd = 
ReplicateContainerCommand.toTarget(containerOnVol1, target);
+      queuePushAndWaitForContainer(cluster, cmd, source, target, 
containerOnVol1);
+      GenericTestUtils.waitFor(() -> healthyPool.getCompletedTaskCount() > 
completedBefore, 100, 30000);
+      assertThat(volSet.getFailedVolumesList()).hasSize(1);
+
+      // Task routes via global pool fallback (HDDS-15327); replication fails 
on bad volume.
+      ReplicateContainerCommand failedVolCmd = 
ReplicateContainerCommand.toTarget(containerOnVol0, target);
+      ReplicationSupervisor supervisor = 
sourceDn.getDatanodeStateMachine().getSupervisor();
+      // Scope the counter to push replication: the cluster-wide counter also 
moves for unrelated
+      // SCM-driven work on this datanode and would satisfy the wait on its 
own.
+      long previousFailures = 
supervisor.getReplicationFailureCount(ReplicationTask.METRIC_NAME);
+      queuePushAndWaitForFailure(cluster, failedVolCmd, source, supervisor, 
previousFailures);
+      
assertThat(supervisor.getReplicationFailureCount(ReplicationTask.METRIC_NAME))
+          .isGreaterThanOrEqualTo(previousFailures + 1);
+      assertThat(hasContainer(cluster, target, containerOnVol0)).isFalse();
+    } finally {
+      // Must run even if an assertion above fails: the volume dir is left 
read-only otherwise, which breaks
+      // the next test and stops the cluster from cleaning up its base dir.
+      DatanodeTestUtils.restoreBadVolume(vol0);
+    }

Review Comment:
   `restoreBadVolume` only restores directory writability; it does not return 
`vol0` from `failedVolumeMap` to the active volume set or recreate its 
replication pool. Consequently the ordered decommission test starts with one 
datanode still missing a volume, and flakes if that datanode is selected from 
the Ratis/EC pipeline intersection because `assertVolumePools` requires two 
pools. Restart this datanode after restoring the directory so the next test 
gets the fully restored cluster promised by the class contract.



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