This is an automated email from the ASF dual-hosted git repository.

errose28 pushed a commit to branch HDDS-14496-zdu
in repository https://gitbox.apache.org/repos/asf/ozone.git


The following commit(s) were added to refs/heads/HDDS-14496-zdu by this push:
     new 72ced386d61 HDDS-16312. SCM needs to choose apparent version for 
reconcile and reconstruct container commands (#11148)
72ced386d61 is described below

commit 72ced386d614aa9f6cddf1974276769ad30aed9b
Author: Ethan Rose <[email protected]>
AuthorDate: Wed Sep 2 10:28:59 2026 -0400

    HDDS-16312. SCM needs to choose apparent version for reconcile and 
reconstruct container commands (#11148)
---
 .../ECReconstructionCommandInfo.java               |  13 +++
 .../commands/ReconcileContainerCommand.java        |  39 ++++++-
 .../commands/ReconstructECContainersCommand.java   |  39 ++++++-
 .../commands/TestReconcileContainerCommand.java    |  88 +++++++++++++++
 .../TestReconstructECContainersCommand.java        |  98 +++++++++++++++++
 .../proto/ScmServerDatanodeHeartbeatProtocol.proto |   2 +
 .../ReconcileContainerEventHandler.java            |  25 ++++-
 .../container/replication/ReplicationManager.java  |  14 +++
 .../hdds/scm/server/StorageContainerManager.java   |   2 +-
 .../TestReconcileContainerEventHandler.java        |  41 ++++++-
 .../replication/TestReplicationManager.java        | 118 ++++++++++++++++++---
 11 files changed, 457 insertions(+), 22 deletions(-)

diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ec/reconstruction/ECReconstructionCommandInfo.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ec/reconstruction/ECReconstructionCommandInfo.java
index 7f86fe8894a..75884c08305 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ec/reconstruction/ECReconstructionCommandInfo.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/ec/reconstruction/ECReconstructionCommandInfo.java
@@ -25,6 +25,7 @@
 import java.util.SortedMap;
 import java.util.TreeMap;
 import java.util.stream.IntStream;
+import org.apache.hadoop.hdds.ComponentVersion;
 import org.apache.hadoop.hdds.StringUtils;
 import org.apache.hadoop.hdds.client.ECReplicationConfig;
 import org.apache.hadoop.hdds.protocol.DatanodeDetails;
@@ -42,6 +43,7 @@ public class ECReconstructionCommandInfo {
   private final ByteString missingContainerIndexes;
   private final long deadlineMsSinceEpoch;
   private final long term;
+  private final ComponentVersion apparentVersion;
 
   public ECReconstructionCommandInfo(ReconstructECContainersCommand cmd) {
     this.containerID = cmd.getContainerID();
@@ -49,6 +51,7 @@ public 
ECReconstructionCommandInfo(ReconstructECContainersCommand cmd) {
     this.missingContainerIndexes = cmd.getMissingContainerIndexes();
     this.deadlineMsSinceEpoch = cmd.getDeadline();
     this.term = cmd.getTerm();
+    this.apparentVersion = cmd.getApparentVersion();
 
     sourceNodeMap = cmd.getSources().stream()
         .collect(toMap(
@@ -97,4 +100,14 @@ public long getTerm() {
     return term;
   }
 
+  /**
+   * This method exists to handle future incompatible changes that may be 
introduced to the EC reconstruction protocol.
+   *
+   * @return the version that this reconstruction request must operate in to 
maintain compatibility with all nodes
+   *    involved.
+   */
+  public ComponentVersion getApparentVersion() {
+    return apparentVersion;
+  }
+
 }
diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReconcileContainerCommand.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReconcileContainerCommand.java
index 6aad68cd877..b8a203ef060 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReconcileContainerCommand.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReconcileContainerCommand.java
@@ -23,10 +23,15 @@
 import java.util.Objects;
 import java.util.Set;
 import java.util.stream.Collectors;
+import org.apache.hadoop.hdds.ComponentVersion;
+import org.apache.hadoop.hdds.HDDSVersion;
 import org.apache.hadoop.hdds.protocol.DatanodeDetails;
 import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
 import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReconcileContainerCommandProto;
 import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
+import org.apache.hadoop.hdds.upgrade.HDDSVersionUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 /**
  * Asks datanodes to reconcile the specified container with other container 
replicas.
@@ -34,6 +39,8 @@
 public class ReconcileContainerCommand extends 
SCMCommand<ReconcileContainerCommandProto> {
 
   private final Set<DatanodeDetails> peerDatanodes;
+  private ComponentVersion apparentVersion;
+  private static final Logger LOG = 
LoggerFactory.getLogger(ReconcileContainerCommand.class);
 
   public ReconcileContainerCommand(long containerID, Set<DatanodeDetails> 
peerDatanodes) {
     // Container ID serves as command ID, since only one reconciliation should 
be in progress at a time.
@@ -41,6 +48,19 @@ public ReconcileContainerCommand(long containerID, 
Set<DatanodeDetails> peerData
     this.peerDatanodes = peerDatanodes;
   }
 
+  public void setApparentVersion(ComponentVersion apparentVersion) {
+    this.apparentVersion = apparentVersion;
+  }
+
+  /**
+   * @return the apparent version that should be used to carry out this
+   *     reconciliation. SCM computes this as the lowest apparent version among
+   *     the nodes involved.
+   */
+  public ComponentVersion getApparentVersion() {
+    return apparentVersion;
+  }
+
   @Override
   public SCMCommandProto.Type getType() {
     return SCMCommandProto.Type.reconcileContainerCommand;
@@ -53,6 +73,8 @@ public ReconcileContainerCommandProto getProto() {
     for (DatanodeDetails dd : peerDatanodes) {
       builder.addPeers(dd.getProtoBufMessage());
     }
+    ComponentVersion version = apparentVersion != null ? apparentVersion : 
HDDSVersion.DEFAULT_VERSION;
+    builder.setApparentVersion(version.serialize());
     return builder.build();
   }
 
@@ -74,14 +96,27 @@ public static ReconcileContainerCommand 
getFromProtobuf(ReconcileContainerComman
         .collect(Collectors.toSet())
         : emptySet();
 
-    return new ReconcileContainerCommand(protoMessage.getContainerID(), 
peerNodes);
+    ReconcileContainerCommand cmd =
+        new ReconcileContainerCommand(protoMessage.getContainerID(), 
peerNodes);
+
+    if (protoMessage.hasApparentVersion()) {
+      cmd.apparentVersion = 
HDDSVersionUtils.deserializeHDDSVersionOrLayoutVersion(
+          protoMessage.getApparentVersion());
+    } else {
+      LOG.warn("Received reconcile command for container {} with no apparent 
version. Falling back to {}",
+          cmd.getContainerID(), HDDSVersion.DEFAULT_VERSION);
+      cmd.apparentVersion = HDDSVersion.DEFAULT_VERSION;
+    }
+
+    return cmd;
   }
 
   @Override
   public String toString() {
     return getType() +
         ": containerId=" + getContainerID() +
-        ", peerNodes=" + peerDatanodes;
+        ", peerNodes=" + peerDatanodes +
+        ", apparentVersion=" + apparentVersion;
   }
 
   @Override
diff --git 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReconstructECContainersCommand.java
 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReconstructECContainersCommand.java
index 59f606608d0..edcd4f132c4 100644
--- 
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReconstructECContainersCommand.java
+++ 
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReconstructECContainersCommand.java
@@ -22,6 +22,8 @@
 import java.util.List;
 import java.util.Objects;
 import java.util.stream.Collectors;
+import org.apache.hadoop.hdds.ComponentVersion;
+import org.apache.hadoop.hdds.HDDSVersion;
 import org.apache.hadoop.hdds.HddsIdFactory;
 import org.apache.hadoop.hdds.client.ECReplicationConfig;
 import org.apache.hadoop.hdds.protocol.DatanodeDetails;
@@ -29,6 +31,9 @@
 import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReconstructECContainersCommandProto;
 import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReconstructECContainersCommandProto.Builder;
 import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto.Type;
+import org.apache.hadoop.hdds.upgrade.HDDSVersionUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 /**
  * SCM command to request reconstruction of EC containers.
@@ -40,6 +45,9 @@ public class ReconstructECContainersCommand
   private final List<DatanodeDetails> targetDatanodes;
   private final ByteString missingContainerIndexes;
   private final ECReplicationConfig ecReplicationConfig;
+  private ComponentVersion apparentVersion;
+
+  private static final Logger LOG = 
LoggerFactory.getLogger(ReconstructECContainersCommand.class);
 
   public ReconstructECContainersCommand(long containerID,
       List<DatanodeDetailsAndReplicaIndex> sources,
@@ -65,6 +73,19 @@ public ReconstructECContainersCommand(long containerID,
     }
   }
 
+  public void setApparentVersion(ComponentVersion apparentVersion) {
+    this.apparentVersion = apparentVersion;
+  }
+
+  /**
+   * @return the apparent version that should be used to carry out this
+   *     reconstruction. SCM computes this as the lowest apparent version among
+   *     the nodes involved.
+   */
+  public ComponentVersion getApparentVersion() {
+    return apparentVersion;
+  }
+
   @Override
   public Type getType() {
     return Type.reconstructECContainersCommand;
@@ -83,6 +104,8 @@ public ReconstructECContainersCommandProto getProto() {
     }
     builder.setMissingContainerIndexes(missingContainerIndexes);
     builder.setEcReplicationConfig(ecReplicationConfig.toProto());
+    ComponentVersion version = apparentVersion != null ? apparentVersion : 
HDDSVersion.DEFAULT_VERSION;
+    builder.setApparentVersion(version.serialize());
     return builder.build();
   }
 
@@ -98,11 +121,22 @@ public static ReconstructECContainersCommand 
getFromProtobuf(
         protoMessage.getTargetsList().stream()
             
.map(DatanodeDetails::getFromProtoBuf).collect(Collectors.toList());
 
-    return new ReconstructECContainersCommand(protoMessage.getContainerID(),
+    ReconstructECContainersCommand cmd = new ReconstructECContainersCommand(
+        protoMessage.getContainerID(),
         srcDatanodeDetails, targetDatanodeDetails,
         protoMessage.getMissingContainerIndexes(),
         new ECReplicationConfig(protoMessage.getEcReplicationConfig()),
         protoMessage.getCmdId());
+    if (protoMessage.hasApparentVersion()) {
+      cmd.apparentVersion = 
HDDSVersionUtils.deserializeHDDSVersionOrLayoutVersion(
+          protoMessage.getApparentVersion());
+    } else {
+      LOG.warn("Received EC reconstruction command for container {} with no 
apparent version. Falling back to {}",
+          cmd.getContainerID(), HDDSVersion.DEFAULT_VERSION);
+      cmd.apparentVersion = HDDSVersion.DEFAULT_VERSION;
+    }
+
+    return cmd;
   }
 
   public long getContainerID() {
@@ -141,7 +175,8 @@ public String toString() {
             .collect(Collectors.joining(", "))).append(']')
         .append(", targets: ").append(getTargetDatanodes())
         .append(", missingIndexes: ").append(
-            Arrays.toString(missingContainerIndexes.toByteArray()));
+            Arrays.toString(missingContainerIndexes.toByteArray()))
+        .append(", apparentVersion: ").append(apparentVersion);
     return sb.toString();
   }
 
diff --git 
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReconcileContainerCommand.java
 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReconcileContainerCommand.java
new file mode 100644
index 00000000000..9b73044448f
--- /dev/null
+++ 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReconcileContainerCommand.java
@@ -0,0 +1,88 @@
+/*
+ * 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.protocol.commands;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import com.google.common.collect.ImmutableSet;
+import org.apache.hadoop.hdds.HDDSVersion;
+import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
+import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReconcileContainerCommandProto;
+import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Test cases to verify {@link ReconcileContainerCommand} serialization.
+ */
+public class TestReconcileContainerCommand {
+
+  @Test
+  public void testApparentVersionRoundTrip() {
+    ReconcileContainerCommand cmd = createCommand();
+    cmd.setApparentVersion(HDDSVersion.ZDU);
+
+    ReconcileContainerCommandProto proto = cmd.getProto();
+    ReconcileContainerCommand deserialized = 
ReconcileContainerCommand.getFromProtobuf(proto);
+
+    assertEquals(HDDSVersion.ZDU, deserialized.getApparentVersion());
+  }
+
+  @Test
+  public void testLayoutFeatureApparentVersionRoundTrip() {
+    // Pre-ZDU datanodes report a layout feature as their apparent version.
+    // The type must survive serialization rather than being erased to an
+    // HDDSVersion.
+    ReconcileContainerCommand cmd = createCommand();
+    cmd.setApparentVersion(HDDSLayoutFeature.HBASE_SUPPORT);
+
+    ReconcileContainerCommandProto proto = cmd.getProto();
+    ReconcileContainerCommand deserialized = 
ReconcileContainerCommand.getFromProtobuf(proto);
+
+    assertEquals(HDDSLayoutFeature.HBASE_SUPPORT, 
deserialized.getApparentVersion());
+  }
+
+  @Test
+  public void testApparentVersionDefaultWhenAbsent() {
+    ReconcileContainerCommand cmd = createCommand();
+
+    // Build a proto without the apparentVersion field to mimic an old server.
+    ReconcileContainerCommandProto proto = 
ReconcileContainerCommandProto.newBuilder(cmd.getProto())
+            .clearApparentVersion()
+            .build();
+
+    ReconcileContainerCommand deserialized = 
ReconcileContainerCommand.getFromProtobuf(proto);
+
+    assertEquals(HDDSVersion.DEFAULT_VERSION, 
deserialized.getApparentVersion());
+  }
+
+  @Test
+  public void testToStringIncludesApparentVersion() {
+    ReconcileContainerCommand cmd = createCommand();
+    cmd.setApparentVersion(HDDSVersion.SOFTWARE_VERSION);
+
+    String str = cmd.toString();
+    assertTrue(str.contains("apparentVersion=" + 
HDDSVersion.SOFTWARE_VERSION));
+  }
+
+  private ReconcileContainerCommand createCommand() {
+    return new ReconcileContainerCommand(1L, ImmutableSet.of(
+        MockDatanodeDetails.randomDatanodeDetails(),
+        MockDatanodeDetails.randomDatanodeDetails()));
+  }
+}
diff --git 
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReconstructECContainersCommand.java
 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReconstructECContainersCommand.java
new file mode 100644
index 00000000000..f63d025cae9
--- /dev/null
+++ 
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReconstructECContainersCommand.java
@@ -0,0 +1,98 @@
+/*
+ * 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.protocol.commands;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import com.google.protobuf.ByteString;
+import java.util.Collections;
+import org.apache.hadoop.hdds.HDDSVersion;
+import org.apache.hadoop.hdds.client.ECReplicationConfig;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
+import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReconstructECContainersCommandProto;
+import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
+import 
org.apache.hadoop.ozone.protocol.commands.ReconstructECContainersCommand.DatanodeDetailsAndReplicaIndex;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Test cases to verify {@link ReconstructECContainersCommand} serialization.
+ */
+public class TestReconstructECContainersCommand {
+
+  @Test
+  public void testApparentVersionRoundTrip() {
+    ReconstructECContainersCommand cmd = createCommand();
+    cmd.setApparentVersion(HDDSVersion.ZDU);
+
+    ReconstructECContainersCommandProto proto = cmd.getProto();
+    ReconstructECContainersCommand deserialized =
+        ReconstructECContainersCommand.getFromProtobuf(proto);
+
+    assertEquals(HDDSVersion.ZDU, deserialized.getApparentVersion());
+  }
+
+  @Test
+  public void testLayoutFeatureApparentVersionRoundTrip() {
+    // Pre-ZDU datanodes report a layout feature as their apparent version.
+    // The type must survive serialization rather than being erased to an
+    // HDDSVersion.
+    ReconstructECContainersCommand cmd = createCommand();
+    cmd.setApparentVersion(HDDSLayoutFeature.HBASE_SUPPORT);
+
+    ReconstructECContainersCommandProto proto = cmd.getProto();
+    ReconstructECContainersCommand deserialized = 
ReconstructECContainersCommand.getFromProtobuf(proto);
+
+    assertEquals(HDDSLayoutFeature.HBASE_SUPPORT, 
deserialized.getApparentVersion());
+  }
+
+  @Test
+  public void testApparentVersionDefaultWhenAbsent() {
+    ReconstructECContainersCommand cmd = createCommand();
+
+    // Build a proto without the apparentVersion field to mimic an old server.
+    ReconstructECContainersCommandProto proto = 
ReconstructECContainersCommandProto.newBuilder(cmd.getProto())
+            .clearApparentVersion()
+            .build();
+
+    ReconstructECContainersCommand deserialized = 
ReconstructECContainersCommand.getFromProtobuf(proto);
+
+    assertEquals(HDDSVersion.DEFAULT_VERSION, 
deserialized.getApparentVersion());
+  }
+
+  @Test
+  public void testToStringIncludesApparentVersion() {
+    ReconstructECContainersCommand cmd = createCommand();
+    cmd.setApparentVersion(HDDSVersion.SOFTWARE_VERSION);
+
+    String str = cmd.toString();
+    assertTrue(str.contains("apparentVersion: " + 
HDDSVersion.SOFTWARE_VERSION));
+  }
+
+  private ReconstructECContainersCommand createCommand() {
+    DatanodeDetails source = MockDatanodeDetails.randomDatanodeDetails();
+    DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
+    byte[] missingIndexes = new byte[]{2};
+    return new ReconstructECContainersCommand(1L,
+        Collections.singletonList(new DatanodeDetailsAndReplicaIndex(source, 
1)),
+        Collections.singletonList(target),
+        ByteString.copyFrom(missingIndexes),
+        new ECReplicationConfig(3, 2));
+  }
+}
diff --git 
a/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto
 
b/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto
index e0a2a390bf3..dc06699f4c0 100644
--- 
a/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto
+++ 
b/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto
@@ -461,6 +461,7 @@ message ReconstructECContainersCommandProto {
   required bytes missingContainerIndexes = 4;
   required ECReplicationConfig ecReplicationConfig = 5;
   required int64 cmdId = 6;
+  optional uint32 apparentVersion = 7;
 }
 
 message DatanodeDetailsAndReplicaIndexProto {
@@ -539,6 +540,7 @@ the container.
 message ReconcileContainerCommandProto {
   required int64 containerID = 1;
   repeated DatanodeDetailsProto peers = 2;
+  optional uint32 apparentVersion = 3;
 }
 
 /**
diff --git 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/reconciliation/ReconcileContainerEventHandler.java
 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/reconciliation/ReconcileContainerEventHandler.java
index 9558e1d51d2..093f404a230 100644
--- 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/reconciliation/ReconcileContainerEventHandler.java
+++ 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/reconciliation/ReconcileContainerEventHandler.java
@@ -19,8 +19,10 @@
 
 import static org.apache.hadoop.hdds.scm.events.SCMEvents.DATANODE_COMMAND;
 
+import java.util.List;
 import java.util.Set;
 import java.util.stream.Collectors;
+import org.apache.hadoop.hdds.ComponentVersion;
 import org.apache.hadoop.hdds.protocol.DatanodeDetails;
 import org.apache.hadoop.hdds.scm.container.ContainerID;
 import org.apache.hadoop.hdds.scm.container.ContainerManager;
@@ -28,6 +30,9 @@
 import org.apache.hadoop.hdds.scm.container.ContainerReplica;
 import 
org.apache.hadoop.hdds.scm.container.reconciliation.ReconciliationEligibilityHandler.EligibilityResult;
 import org.apache.hadoop.hdds.scm.ha.SCMContext;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
+import org.apache.hadoop.hdds.scm.node.NodeManager;
+import org.apache.hadoop.hdds.scm.server.upgrade.ScmVersionManager;
 import org.apache.hadoop.hdds.server.events.EventHandler;
 import org.apache.hadoop.hdds.server.events.EventPublisher;
 import org.apache.hadoop.ozone.protocol.commands.CommandForDatanode;
@@ -46,10 +51,13 @@ public class ReconcileContainerEventHandler implements 
EventHandler<ContainerID>
 
   private final ContainerManager containerManager;
   private final SCMContext scmContext;
+  private final NodeManager nodeManager;
 
-  public ReconcileContainerEventHandler(ContainerManager containerManager, 
SCMContext scmContext) {
+  public ReconcileContainerEventHandler(ContainerManager containerManager, 
SCMContext scmContext,
+      NodeManager nodeManager) {
     this.containerManager = containerManager;
     this.scmContext = scmContext;
+    this.nodeManager = nodeManager;
   }
 
   @Override
@@ -75,12 +83,18 @@ public void onMessage(ContainerID containerID, 
EventPublisher publisher) {
 
       LOG.info("Reconcile container event triggered for container {} with 
peers {}", containerID, allReplicaNodes);
 
+      List<DatanodeInfo> involved = allReplicaNodes.stream()
+          .map(this::getDatanodeInfo)
+          .collect(Collectors.toList());
+      ComponentVersion apparentVersion = 
ScmVersionManager.computeVersionForReplication(involved);
+
       for (DatanodeDetails replica : allReplicaNodes) {
         Set<DatanodeDetails> otherReplicas = allReplicaNodes.stream()
             .filter(other -> !other.equals(replica))
             .collect(Collectors.toSet());
         ReconcileContainerCommand command = new 
ReconcileContainerCommand(containerID.getId(), otherReplicas);
         command.setTerm(scmContext.getTermOfLeader());
+        command.setApparentVersion(apparentVersion);
         publisher.fireEvent(DATANODE_COMMAND, new 
CommandForDatanode<>(replica, command));
       }
     } catch (ContainerNotFoundException ex) {
@@ -89,4 +103,13 @@ public void onMessage(ContainerID containerID, 
EventPublisher publisher) {
       LOG.info("Skip reconciling container {} since current SCM is not 
leader.", containerID);
     }
   }
+
+  private DatanodeInfo getDatanodeInfo(DatanodeDetails dnDetails) {
+    DatanodeInfo datanodeInfo = nodeManager.getNode(dnDetails.getID());
+    if (datanodeInfo == null) {
+      throw new IllegalArgumentException("Datanode " + dnDetails + " not " +
+          "found in NodeManager. Should not happen");
+    }
+    return datanodeInfo;
+  }
 }
diff --git 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java
 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java
index b8bc10758f8..6cbd03a3db3 100644
--- 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java
+++ 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationManager.java
@@ -550,9 +550,23 @@ public void 
sendThrottledReconstructionCommand(ContainerInfo containerInfo,
     }
     DatanodeDetails target = selectAndOptionallyExcludeDatanode(
         rmConf.getReconstructionCommandWeight(), targetWithCmds);
+    command.setApparentVersion(computeVersionForReconstruction(command));
     sendDatanodeCommand(command, containerInfo, target);
   }
 
+  private ComponentVersion 
computeVersionForReconstruction(ReconstructECContainersCommand command) {
+    List<DatanodeInfo> involved = new ArrayList<>();
+    for (ReconstructECContainersCommand.DatanodeDetailsAndReplicaIndex source 
: command.getSources()) {
+      involved.add(getDatanodeInfo(source.getDnDetails()));
+    }
+    for (DatanodeDetails target : command.getTargetDatanodes()) {
+      involved.add(getDatanodeInfo(target));
+    }
+    // Reconstruction is handled as a specific type of replication command 
which uses the same generic logic within the
+    // version manager as computing a version for full container replication.
+    return ScmVersionManager.computeVersionForReplication(involved);
+  }
+
   private ComponentVersion computeVersionForReplication(DatanodeDetails 
source, DatanodeDetails target) {
     return ScmVersionManager.computeVersionForReplication(
         Arrays.asList(getDatanodeInfo(source), getDatanodeInfo(target)));
diff --git 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java
 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java
index 894b9dd2125..1364131778d 100644
--- 
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java
+++ 
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/StorageContainerManager.java
@@ -526,7 +526,7 @@ private void initializeEventHandlers() {
         new ReplicationManagerEventHandler(replicationManager, scmContext);
 
     ReconcileContainerEventHandler reconcileContainerEventHandler =
-        new ReconcileContainerEventHandler(containerManager, scmContext);
+        new ReconcileContainerEventHandler(containerManager, scmContext, 
scmNodeManager);
 
     eventQueue.addHandler(SCMEvents.DATANODE_COMMAND, scmNodeManager);
     eventQueue.addHandler(SCMEvents.RETRIABLE_DATANODE_COMMAND, 
scmNodeManager);
diff --git 
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/reconciliation/TestReconcileContainerEventHandler.java
 
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/reconciliation/TestReconcileContainerEventHandler.java
index cebd376a05e..9d35faf7476 100644
--- 
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/reconciliation/TestReconcileContainerEventHandler.java
+++ 
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/reconciliation/TestReconcileContainerEventHandler.java
@@ -36,9 +36,12 @@
 import java.util.List;
 import java.util.Set;
 import java.util.stream.Collectors;
+import org.apache.hadoop.hdds.ComponentVersion;
+import org.apache.hadoop.hdds.HDDSVersion;
 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.protocol.DatanodeDetails;
 import org.apache.hadoop.hdds.protocol.DatanodeID;
 import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState;
 import 
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ContainerReplicaProto.State;
@@ -54,8 +57,11 @@
 import 
org.apache.hadoop.hdds.scm.container.reconciliation.ReconciliationEligibilityHandler.Result;
 import org.apache.hadoop.hdds.scm.ha.SCMContext;
 import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
+import org.apache.hadoop.hdds.scm.node.NodeManager;
 import org.apache.hadoop.hdds.server.events.EventPublisher;
+import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
 import org.apache.hadoop.ozone.protocol.commands.CommandForDatanode;
+import org.apache.hadoop.ozone.protocol.commands.ReconcileContainerCommand;
 import org.apache.hadoop.ozone.protocol.commands.SCMCommand;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
@@ -72,6 +78,7 @@ public class TestReconcileContainerEventHandler {
   private EventPublisher eventPublisher;
   private ReconcileContainerEventHandler eventHandler;
   private SCMContext scmContext;
+  private NodeManager nodeManager;
 
   private static final ContainerID CONTAINER_ID = ContainerID.valueOf(123L);
   private static final long LEADER_TERM = 3L;
@@ -90,7 +97,14 @@ public void setup() throws Exception {
     when(scmContext.isLeader()).thenReturn(true);
     when(scmContext.getTermOfLeader()).thenReturn(LEADER_TERM);
     eventPublisher = mock(EventPublisher.class);
-    eventHandler = new ReconcileContainerEventHandler(containerManager, 
scmContext);
+    nodeManager = mock(NodeManager.class);
+    // By default every datanode is known and reports the current version, so
+    // the reconcile path can resolve an apparent version. Individual tests
+    // override this for specific nodes.
+    DatanodeInfo defaultNodeInfo = mock(DatanodeInfo.class);
+    
when(defaultNodeInfo.getLastKnownApparentVersion()).thenReturn(HDDSVersion.SOFTWARE_VERSION);
+    
when(nodeManager.getNode(any(DatanodeID.class))).thenReturn(defaultNodeInfo);
+    eventHandler = new ReconcileContainerEventHandler(containerManager, 
scmContext, nodeManager);
   }
 
   /**
@@ -249,6 +263,31 @@ public void testReconcileSentToAllPeers() throws Exception 
{
     assertEquals(allNodeIDs, nodesReceivingCommands);
   }
 
+  @Test
+  public void testReconcileCommandCarriesLowestApparentVersion() throws 
Exception {
+    addContainer(RATIS_THREE_REP, LifeCycleState.CLOSED);
+    Set<ContainerReplica> replicas = addReplicasToContainer(3);
+
+    // Every node defaults to SOFTWARE_VERSION; stub one replica lower so all
+    // dispatched commands must fall back to the lower version.
+    ComponentVersion lower = HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION;
+    DatanodeDetails olderNode = 
replicas.iterator().next().getDatanodeDetails();
+    DatanodeInfo olderInfo = mock(DatanodeInfo.class);
+    when(olderInfo.getLastKnownApparentVersion()).thenReturn(lower);
+    when(nodeManager.getNode(olderNode.getID())).thenReturn(olderInfo);
+
+    eventHandler.onMessage(CONTAINER_ID, eventPublisher);
+
+    verify(eventPublisher, 
times(replicas.size())).fireEvent(eq(DATANODE_COMMAND), 
commandCaptor.capture());
+    int numCmds = 0;
+    for (CommandForDatanode<ReconcileContainerCommandProto> dnCommand : 
commandCaptor.getAllValues()) {
+      ReconcileContainerCommand command = (ReconcileContainerCommand) 
dnCommand.getCommand();
+      assertEquals(lower, command.getApparentVersion());
+      numCmds++;
+    }
+    assertEquals(3, numCmds);
+  }
+
   @ParameterizedTest
   @EnumSource(State.class)
   public void testReconcileFailsWithIneligibleReplicas(State replicaState) 
throws Exception {
diff --git 
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManager.java
 
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManager.java
index 5943efdd4ac..5c558bae902 100644
--- 
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManager.java
+++ 
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManager.java
@@ -112,7 +112,6 @@
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.params.ParameterizedTest;
-import org.junit.jupiter.params.provider.CsvSource;
 import org.junit.jupiter.params.provider.EnumSource;
 import org.junit.jupiter.params.provider.ValueSource;
 import org.mockito.ArgumentCaptor;
@@ -1423,6 +1422,81 @@ public void 
testSendThrottledReconstructionCommandThrowsWhenNoTargets()
         .getEcReconstructionCmdsDeferredTotal());
   }
 
+  /**
+   * The reconstruction command must carry the lowest apparent version among 
all
+   * the involved nodes (sources plus targets). This case places the lowest
+   * version on a source.
+   */
+  @Test
+  public void testReconstructionWhenSourceHasLowestVersion()
+      throws CommandTargetOverloadedException, NodeNotFoundException, 
NotLeaderException {
+    assertReconstructionCarriesLowestApparentVersion(true);
+  }
+
+  /**
+   * The reconstruction command must carry the lowest apparent version among 
all
+   * the involved nodes (sources plus targets). This case places the lowest
+   * version on a target.
+   */
+  @Test
+  public void testReconstructionWhenTargetHasLowestVersion()
+      throws CommandTargetOverloadedException, NodeNotFoundException, 
NotLeaderException {
+    assertReconstructionCarriesLowestApparentVersion(false);
+  }
+
+  private void assertReconstructionCarriesLowestApparentVersion(boolean 
sourceIsLowest)
+      throws CommandTargetOverloadedException, NodeNotFoundException, 
NotLeaderException {
+    Map<DatanodeDetails, Integer> targetNodes = new HashMap<>();
+    DatanodeDetails cmdTarget = MockDatanodeDetails.randomDatanodeDetails();
+    targetNodes.put(cmdTarget, 0);
+    targetNodes.put(MockDatanodeDetails.randomDatanodeDetails(), 5);
+
+    mockReplicationCommandCounts(targetNodes::get, any -> 0);
+
+    ContainerInfo container = ReplicationTestUtil.createContainerInfo(
+        repConfig, 1, HddsProtos.LifeCycleState.CLOSED, 10, 20);
+    ReconstructECContainersCommand command = createReconstructionCommand(
+        container, targetNodes.keySet().toArray(new DatanodeDetails[0]));
+
+    // Every node defaults to SOFTWARE_VERSION; stub one node lower so the
+    // command must fall back to the lower version.
+    ComponentVersion lower = HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION;
+    DatanodeDetails olderNode = sourceIsLowest
+        ? command.getSources().get(0).getDnDetails()
+        : command.getTargetDatanodes().get(0);
+    mockDatanodeWithApparentVersion(olderNode, lower);
+
+    replicationManager.sendThrottledReconstructionCommand(container, command);
+
+    assertEquals(1, commandsSent.size());
+    Pair<DatanodeID, SCMCommand<?>> cmd = commandsSent.iterator().next();
+    ReconstructECContainersCommand sent =
+        (ReconstructECContainersCommand) cmd.getValue();
+    assertEquals(lower, sent.getApparentVersion());
+  }
+
+  @Test
+  public void testReconstructionApparentVersionLookupThrowsWhenNodeNotFound()
+      throws NodeNotFoundException {
+    Map<DatanodeDetails, Integer> targetNodes = new HashMap<>();
+    targetNodes.put(MockDatanodeDetails.randomDatanodeDetails(), 0);
+    targetNodes.put(MockDatanodeDetails.randomDatanodeDetails(), 0);
+
+    mockReplicationCommandCounts(targetNodes::get, any -> 0);
+
+    ContainerInfo container = ReplicationTestUtil.createContainerInfo(
+        repConfig, 1, HddsProtos.LifeCycleState.CLOSED, 10, 20);
+    ReconstructECContainersCommand command = createReconstructionCommand(
+        container, targetNodes.keySet().toArray(new DatanodeDetails[0]));
+
+    // SCM has no information for one of the sources.
+    DatanodeDetails unknownSource = command.getSources().get(0).getDnDetails();
+    when(nodeManager.getNode(unknownSource.getID())).thenReturn(null);
+
+    assertThrows(IllegalArgumentException.class, () ->
+        replicationManager.sendThrottledReconstructionCommand(container, 
command));
+  }
+
   private ReconstructECContainersCommand createReconstructionCommand(
       ContainerInfo containerInfo, DatanodeDetails... targets) {
     List<ReconstructECContainersCommand.DatanodeDetailsAndReplicaIndex> sources
@@ -1773,19 +1847,33 @@ private void mockDatanodeWithApparentVersion(
   }
 
   /**
-   * Regardless of which datanode is newer, and regardless of the command
-   * sending path, the replicate command must carry the lowest apparent version
-   * among the source and target datanodes.
+   * Regardless of the command sending path, the replicate command must carry
+   * the lowest apparent version among the source and target datanodes. This
+   * case places the lowest version on the source.
+   */
+  @ParameterizedTest
+  @ValueSource(booleans = {true, false})
+  public void testReplicationApparentVersionWhenSourceIsLowest(boolean 
throttled)
+      throws CommandTargetOverloadedException, NotLeaderException,
+      NodeNotFoundException {
+    assertReplicationCarriesLowestApparentVersion(throttled, true);
+  }
+
+  /**
+   * Regardless of the command sending path, the replicate command must carry
+   * the lowest apparent version among the source and target datanodes. This
+   * case places the lowest version on the target.
    */
   @ParameterizedTest
-  @CsvSource({
-      "true, true",
-      "true, false",
-      "false, true",
-      "false, false",
-  })
-  public void testApparentVersionIsLowestOfSourceAndTarget(
-      boolean throttled, boolean sourceNewer)
+  @ValueSource(booleans = {true, false})
+  public void testReplicationApparentVersionWhenTargetIsLowest(boolean 
throttled)
+      throws CommandTargetOverloadedException, NotLeaderException,
+      NodeNotFoundException {
+    assertReplicationCarriesLowestApparentVersion(throttled, false);
+  }
+
+  private void assertReplicationCarriesLowestApparentVersion(
+      boolean throttled, boolean sourceIsLowest)
       throws CommandTargetOverloadedException, NotLeaderException,
       NodeNotFoundException {
     ContainerInfo containerInfo =
@@ -1796,8 +1884,8 @@ public void testApparentVersionIsLowestOfSourceAndTarget(
 
     ComponentVersion lower = HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION;
     ComponentVersion higher = HDDSVersion.ZDU;
-    mockDatanodeWithApparentVersion(source, sourceNewer ? higher : lower);
-    mockDatanodeWithApparentVersion(target, sourceNewer ? lower : higher);
+    mockDatanodeWithApparentVersion(source, sourceIsLowest ? lower : higher);
+    mockDatanodeWithApparentVersion(target, sourceIsLowest ? higher : lower);
 
     if (throttled) {
       mockReplicationCommandCounts(dn -> 0, dn -> 0);
@@ -1812,7 +1900,7 @@ public void testApparentVersionIsLowestOfSourceAndTarget(
   }
 
   @Test
-  public void testApparentVersionLookupThrowsWhenNodeNotFound()
+  public void testReplicationApparentVersionLookupThrowsWhenNodeNotFound()
       throws NodeNotFoundException {
     ContainerInfo containerInfo =
         ReplicationTestUtil.createContainerInfo(repConfig, 1,


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to