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 ad1c396df42 HDDS-15417. Handle mixed datanode versions on the
replication path (#10570)
ad1c396df42 is described below
commit ad1c396df424a785f824b7f3b0b457fed38cbd8c
Author: Zita Dombi <[email protected]>
AuthorDate: Tue Jul 28 15:05:42 2026 +0200
HDDS-15417. Handle mixed datanode versions on the replication path (#10570)
---
.../org/apache/hadoop/hdds/ComponentVersion.java | 22 +++++
.../hadoop/hdds/AbstractComponentVersionTest.java | 29 ++++++
.../container/replication/ReplicationTask.java | 12 +++
.../commands/ReplicateContainerCommand.java | 30 ++++++-
.../ozone/container/common/ContainerTestUtils.java | 13 +++
.../common/statemachine/TestStateContext.java | 10 +--
.../TestReplicateContainerCommandHandler.java | 49 ++++++++--
.../ReplicationSupervisorSchedulingBenchmark.java | 4 +-
.../replication/TestGrpcReplicationService.java | 5 +-
.../replication/TestMeasuredReplicator.java | 20 ++---
.../container/replication/TestPushReplicator.java | 8 +-
.../replication/TestReplicationSupervisor.java | 7 +-
.../commands/TestReplicateContainerCommand.java | 100 +++++++++++++++++++++
.../proto/ScmServerDatanodeHeartbeatProtocol.proto | 1 +
.../container/replication/ReplicationManager.java | 32 +++++--
.../apache/hadoop/hdds/scm/node/NodeManager.java | 22 +++++
.../container/replication/ReplicationTestUtil.java | 5 +-
.../TestContainerReplicaPendingOps.java | 4 +-
.../replication/TestReplicationManager.java | 98 +++++++++++++++++++-
.../TestReplicationManagerScenarios.java | 10 +++
.../hadoop/hdds/scm/node/TestCommandQueue.java | 9 +-
.../hadoop/hdds/scm/node/TestSCMNodeManager.java | 6 +-
.../hdds/scm/TestSCMDatanodeProtocolServer.java | 4 +-
.../replication/TestContainerReplication.java | 11 ++-
24 files changed, 446 insertions(+), 65 deletions(-)
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/ComponentVersion.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/ComponentVersion.java
index 9c1223f35c3..77c048e31b4 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/ComponentVersion.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/ComponentVersion.java
@@ -80,4 +80,26 @@ default boolean isSupportedBy(ComponentVersion other) {
default Optional<? extends UpgradeAction> action() {
return Optional.empty();
}
+
+ /**
+ * Returns the version with the lowest feature set among the given versions,
+ * so callers can pick a version that is mutually supported by all of them.
+ * Comparison is done through {@link #isSupportedBy}, which respects the
+ * negative/unknown-future-version convention, rather than comparing the
+ * opaque {@link #serialize()} values directly.
+ *
+ * @throws IllegalArgumentException if no versions are provided.
+ */
+ static ComponentVersion min(ComponentVersion... versions) {
+ if (versions.length == 0) {
+ throw new IllegalArgumentException("At least one version is required.");
+ }
+ ComponentVersion lowest = versions[0];
+ for (int i = 1; i < versions.length; i++) {
+ if (versions[i].isSupportedBy(lowest)) {
+ lowest = versions[i];
+ }
+ }
+ return lowest;
+ }
}
diff --git
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/AbstractComponentVersionTest.java
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/AbstractComponentVersionTest.java
index 49f54776efb..27671edfd74 100644
---
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/AbstractComponentVersionTest.java
+++
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/AbstractComponentVersionTest.java
@@ -22,6 +22,7 @@
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import org.junit.jupiter.api.Test;
@@ -130,4 +131,32 @@ public void testVersionSerDes() {
public void testDeserializeUnknownVersion() {
assertEquals(getUnknownVersion(), deserialize(Integer.MAX_VALUE));
}
+
+ @Test
+ public void testMinRequiresAtLeastOneVersion() {
+ assertThrows(IllegalArgumentException.class, ComponentVersion::min);
+ }
+
+ @Test
+ public void testMinOfSingleVersionIsItself() {
+ ComponentVersion version = getValues()[0];
+ assertEquals(version, ComponentVersion.min(version));
+ }
+
+ @Test
+ public void testMinReturnsLowestKnownVersion() {
+ // getValues()[0] is the lowest known version
+ assertEquals(getValues()[0], ComponentVersion.min(getValues()));
+ }
+
+ @Test
+ public void testMinTreatsUnknownFutureVersionAsHighest() {
+ ComponentVersion known = getValues()[0];
+ ComponentVersion unknown = getUnknownVersion();
+ // The known version is lower regardless of argument order.
+ assertEquals(known, ComponentVersion.min(known, unknown));
+ assertEquals(known, ComponentVersion.min(unknown, known));
+ // With only the unknown version, it is returned unchanged.
+ assertEquals(unknown, ComponentVersion.min(unknown));
+ }
}
diff --git
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ReplicationTask.java
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ReplicationTask.java
index ce8f535e0c4..e21b47827aa 100644
---
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ReplicationTask.java
+++
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/replication/ReplicationTask.java
@@ -18,6 +18,7 @@
package org.apache.hadoop.ozone.container.replication;
import java.util.Objects;
+import org.apache.hadoop.hdds.ComponentVersion;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
@@ -29,6 +30,7 @@ public class ReplicationTask extends AbstractReplicationTask {
private final ReplicateContainerCommand cmd;
private final ContainerReplicator replicator;
private final String debugString;
+ private final ComponentVersion apparentVersion;
public static final String METRIC_NAME = "ContainerReplications";
public static final String METRIC_DESCRIPTION_SEGMENT = "container
replications";
@@ -43,6 +45,7 @@ public ReplicationTask(ReplicateContainerCommand cmd,
setPriority(cmd.getPriority());
this.cmd = cmd;
this.replicator = replicator;
+ this.apparentVersion = cmd.getApparentVersion();
debugString = cmd.toString();
}
@@ -101,6 +104,15 @@ public void setTransferredBytes(long transferredBytes) {
this.transferredBytes = transferredBytes;
}
+ /**
+ * @return the apparent version to use for this replication task. Replicators
+ * can use this to gate version-dependent protocol features so that the
+ * newer side downgrades to a behavior the older side understands.
+ */
+ public ComponentVersion getApparentVersion() {
+ return apparentVersion;
+ }
+
DatanodeDetails getTarget() {
return cmd.getTargetDatanode();
}
diff --git
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReplicateContainerCommand.java
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReplicateContainerCommand.java
index c3543909653..a91ffc8fd78 100644
---
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReplicateContainerCommand.java
+++
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/protocol/commands/ReplicateContainerCommand.java
@@ -18,12 +18,15 @@
package org.apache.hadoop.ozone.protocol.commands;
import java.util.Objects;
+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.StorageContainerDatanodeProtocolProtos.ReplicateContainerCommandProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicateContainerCommandProto.Builder;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicationCommandPriority;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto.Type;
+import org.apache.hadoop.hdds.upgrade.HDDSVersionUtils;
/**
* SCM command to request push-replication of a container to a target datanode.
@@ -36,10 +39,14 @@ public final class ReplicateContainerCommand
private int replicaIndex = 0;
private ReplicationCommandPriority priority =
ReplicationCommandPriority.NORMAL;
+ private ComponentVersion apparentVersion = HDDSVersion.DEFAULT_VERSION;
public static ReplicateContainerCommand toTarget(long containerID,
- DatanodeDetails target) {
- return new ReplicateContainerCommand(containerID, target);
+ DatanodeDetails target, ComponentVersion apparentVersion) {
+ ReplicateContainerCommand cmd =
+ new ReplicateContainerCommand(containerID, target);
+ cmd.apparentVersion = apparentVersion;
+ return cmd;
}
private ReplicateContainerCommand(long containerID, DatanodeDetails target) {
@@ -63,6 +70,15 @@ public void setPriority(ReplicationCommandPriority priority)
{
this.priority = priority;
}
+ /**
+ * @return the apparent version that should be used to carry out this
+ * replication. SCM computes this as the lowest apparent version among
the
+ * nodes involved.
+ */
+ public ComponentVersion getApparentVersion() {
+ return apparentVersion;
+ }
+
@Override
public Type getType() {
return SCMCommandProto.Type.replicateContainerCommand;
@@ -80,7 +96,8 @@ public ReplicateContainerCommandProto getProto() {
.setContainerID(containerID)
.setReplicaIndex(replicaIndex)
.setTarget(targetDatanode.getProtoBufMessage())
- .setPriority(priority);
+ .setPriority(priority)
+ .setApparentVersion(apparentVersion.serialize());
return builder.build();
}
@@ -100,6 +117,10 @@ public static ReplicateContainerCommand getFromProtobuf(
if (protoMessage.hasPriority()) {
cmd.setPriority(protoMessage.getPriority());
}
+ if (protoMessage.hasApparentVersion()) {
+ cmd.apparentVersion =
HDDSVersionUtils.deserializeHDDSVersionOrLayoutVersion(
+ protoMessage.getApparentVersion());
+ }
return cmd;
}
@@ -129,6 +150,7 @@ public String toString() {
+ ", containerId=" + getContainerID()
+ ", replicaIndex=" + getReplicaIndex()
+ ", targetNode=" + targetDatanode
- + ", priority=" + priority;
+ + ", priority=" + priority
+ + ", apparentVersion=" + apparentVersion;
}
}
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/ContainerTestUtils.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/ContainerTestUtils.java
index a6827c7300c..7172aa3fe50 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/ContainerTestUtils.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/ContainerTestUtils.java
@@ -36,6 +36,7 @@
import java.util.concurrent.atomic.AtomicLong;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.conf.StorageUnit;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.client.BlockID;
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
@@ -95,6 +96,7 @@
import org.apache.hadoop.ozone.container.ozoneimpl.DataScanResult;
import org.apache.hadoop.ozone.container.ozoneimpl.MetadataScanResult;
import org.apache.hadoop.ozone.container.ozoneimpl.OzoneContainer;
+import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
import
org.apache.hadoop.ozone.protocolPB.StorageContainerDatanodeProtocolClientSideTranslatorPB;
import org.apache.hadoop.ozone.protocolPB.StorageContainerDatanodeProtocolPB;
import org.apache.hadoop.security.UserGroupInformation;
@@ -475,4 +477,15 @@ public static void
createBlockMetaData(KeyValueContainerData data, int numOfBloc
}
}
}
+
+ /**
+ * Creates a replicate container command filled in with the latest apparent
+ * version, so callers that do not exercise mixed-version handling don't need
+ * to specify it.
+ */
+ public static ReplicateContainerCommand getReplicateContainerCommand(
+ long containerID, DatanodeDetails target) {
+ return ReplicateContainerCommand.toTarget(
+ containerID, target, HDDSVersion.SOFTWARE_VERSION);
+ }
}
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestStateContext.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestStateContext.java
index 42220c6b99f..ce0b2420291 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestStateContext.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/TestStateContext.java
@@ -61,6 +61,7 @@
import org.apache.hadoop.hdds.scm.net.HostAndPort;
import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
import org.apache.hadoop.hdfs.util.EnumCounters;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.container.common.impl.ContainerSet;
import
org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine.DatanodeStates;
import org.apache.hadoop.ozone.container.common.states.DatanodeState;
@@ -68,7 +69,6 @@
import org.apache.hadoop.ozone.protocol.commands.CloseContainerCommand;
import org.apache.hadoop.ozone.protocol.commands.ClosePipelineCommand;
import org.apache.hadoop.ozone.protocol.commands.ReconcileContainerCommand;
-import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
import org.apache.hadoop.ozone.protocol.commands.SCMCommand;
import org.apache.ozone.test.LambdaTestUtils;
import org.junit.jupiter.api.Test;
@@ -703,10 +703,10 @@ public void testGetReports() {
@Test
public void testCommandQueueSummary() throws IOException {
StateContext ctx = createSubject();
- ctx.addCommand(ReplicateContainerCommand.toTarget(1,
MockDatanodeDetails.randomDatanodeDetails()));
+ ctx.addCommand(ContainerTestUtils.getReplicateContainerCommand(1,
MockDatanodeDetails.randomDatanodeDetails()));
ctx.addCommand(new ClosePipelineCommand(PipelineID.randomId()));
- ctx.addCommand(ReplicateContainerCommand.toTarget(2,
MockDatanodeDetails.randomDatanodeDetails()));
- ctx.addCommand(ReplicateContainerCommand.toTarget(3,
MockDatanodeDetails.randomDatanodeDetails()));
+ ctx.addCommand(ContainerTestUtils.getReplicateContainerCommand(2,
MockDatanodeDetails.randomDatanodeDetails()));
+ ctx.addCommand(ContainerTestUtils.getReplicateContainerCommand(3,
MockDatanodeDetails.randomDatanodeDetails()));
ctx.addCommand(new ClosePipelineCommand(PipelineID.randomId()));
ctx.addCommand(new CloseContainerCommand(1, PipelineID.randomId()));
ctx.addCommand(new ReconcileContainerCommand(4, Collections.emptySet()));
@@ -773,7 +773,7 @@ private static StateContext createSubject() throws
IOException {
}
private static SCMCommand<?> someCommand() {
- return ReplicateContainerCommand.toTarget(1,
MockDatanodeDetails.randomDatanodeDetails());
+ return ContainerTestUtils.getReplicateContainerCommand(1,
MockDatanodeDetails.randomDatanodeDetails());
}
}
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/commandhandler/TestReplicateContainerCommandHandler.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/commandhandler/TestReplicateContainerCommandHandler.java
index 74edeae7aff..1830b0af391 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/commandhandler/TestReplicateContainerCommandHandler.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/common/statemachine/commandhandler/TestReplicateContainerCommandHandler.java
@@ -25,10 +25,12 @@
import java.util.HashMap;
import java.util.Map;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
import org.apache.hadoop.metrics2.impl.MetricsCollectorImpl;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.container.common.helpers.CommandHandlerMetrics;
import
org.apache.hadoop.ozone.container.common.statemachine.SCMConnectionManager;
import org.apache.hadoop.ozone.container.common.statemachine.StateContext;
@@ -39,6 +41,7 @@
import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
/**
* Test cases to verify {@link ReplicateContainerCommandHandler}.
@@ -59,6 +62,17 @@ public void setUp() {
stateContext = mock(StateContext.class);
}
+ /**
+ * Stubs {@link ReplicationSupervisor#addTask} and returns a captor for the
+ * {@link ReplicationTask} the handler submits.
+ */
+ private ArgumentCaptor<ReplicationTask> captureSubmittedTask() {
+ ArgumentCaptor<ReplicationTask> captor =
+ ArgumentCaptor.forClass(ReplicationTask.class);
+ doNothing().when(supervisor).addTask(captor.capture());
+ return captor;
+ }
+
@Test
public void testMetrics() {
ReplicateContainerCommandHandler commandHandler =
@@ -71,20 +85,22 @@ public void testMetrics() {
DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
ReplicateContainerCommand command =
- ReplicateContainerCommand.toTarget(1, target);
+ ContainerTestUtils.getReplicateContainerCommand(1, target);
commandHandler.handle(command, ozoneContainer, stateContext,
connectionManager);
String metricsName = ReplicationTask.METRIC_NAME;
assertEquals(commandHandler.getMetricsName(), metricsName);
when(supervisor.getReplicationRequestCount(metricsName)).thenReturn(1L);
assertEquals(commandHandler.getInvocationCount(), 1);
- commandHandler.handle(ReplicateContainerCommand.toTarget(2, target),
+ commandHandler.handle(ContainerTestUtils.getReplicateContainerCommand(2,
target),
ozoneContainer, stateContext, connectionManager);
- commandHandler.handle(ReplicateContainerCommand.toTarget(3, target),
+ commandHandler.handle(ContainerTestUtils.getReplicateContainerCommand(3,
target),
ozoneContainer, stateContext, connectionManager);
- commandHandler.handle(ReplicateContainerCommand.toTarget(4, target),
+ commandHandler.handle(
+ ContainerTestUtils.getReplicateContainerCommand(4, target),
ozoneContainer, stateContext, connectionManager);
- commandHandler.handle(ReplicateContainerCommand.toTarget(5, target),
+ commandHandler.handle(
+ ContainerTestUtils.getReplicateContainerCommand(5, target),
ozoneContainer, stateContext, connectionManager);
when(supervisor.getReplicationRequestCount(metricsName)).thenReturn(5L);
@@ -103,4 +119,27 @@ public void testMetrics() {
metrics.unRegister();
}
}
+
+ /**
+ * The datanode follows the apparent version decided by SCM and carried in
the
+ * command, without recomputing it. The submitted task should expose exactly
+ * the command's apparent version.
+ */
+ @Test
+ public void testApparentVersionTakenFromCommand() {
+ ArgumentCaptor<ReplicationTask> captor = captureSubmittedTask();
+
+ ReplicateContainerCommandHandler handler =
+ new ReplicateContainerCommandHandler(supervisor, pushReplicator);
+
+ DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
+ ReplicateContainerCommand cmd = ReplicateContainerCommand.toTarget(
+ 1, target, HDDSVersion.SEPARATE_RATIS_PORTS_AVAILABLE);
+
+ handler.handle(cmd, ozoneContainer, stateContext, connectionManager);
+
+ ReplicationTask task = captor.getValue();
+ assertEquals(HDDSVersion.SEPARATE_RATIS_PORTS_AVAILABLE,
+ task.getApparentVersion());
+ }
}
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/ReplicationSupervisorSchedulingBenchmark.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/ReplicationSupervisorSchedulingBenchmark.java
index e3d0c4126da..980b2c27f37 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/ReplicationSupervisorSchedulingBenchmark.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/ReplicationSupervisorSchedulingBenchmark.java
@@ -17,7 +17,7 @@
package org.apache.hadoop.ozone.container.replication;
-import static
org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand.toTarget;
+import static
org.apache.hadoop.ozone.container.common.ContainerTestUtils.getReplicateContainerCommand;
import static org.assertj.core.api.Assertions.assertThat;
import java.util.HashMap;
@@ -77,7 +77,7 @@ public void test() throws InterruptedException {
//schedule 100 container replication
for (int i = 0; i < 100; i++) {
- rs.addTask(new ReplicationTask(toTarget(i, target), replicator));
+ rs.addTask(new ReplicationTask(getReplicateContainerCommand(i, target),
replicator));
}
rs.shutdownAfterFinish();
final long executionTime = Time.monotonicNow() - start;
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcReplicationService.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcReplicationService.java
index 079c41fcbec..e688d65653b 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcReplicationService.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestGrpcReplicationService.java
@@ -19,7 +19,6 @@
import static org.apache.hadoop.ozone.OzoneConsts.GB;
import static
org.apache.hadoop.ozone.container.common.impl.ContainerImplTestUtils.newContainerSet;
-import static
org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand.toTarget;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyLong;
@@ -170,7 +169,7 @@ public void testUpload() {
PushReplicator pushReplicator = new PushReplicator(conf, source, uploader);
ReplicationTask task =
- new ReplicationTask(toTarget(CONTAINER_ID, datanode), pushReplicator);
+ new
ReplicationTask(ContainerTestUtils.getReplicateContainerCommand(CONTAINER_ID,
datanode), pushReplicator);
pushReplicator.replicate(task);
@@ -200,7 +199,7 @@ public void copyData(long containerId, OutputStream
destination,
PushReplicator pushReplicator = new PushReplicator(conf, source, uploader);
ReplicationTask task = new ReplicationTask(
- toTarget(CONTAINER_ID, datanode), pushReplicator);
+ ContainerTestUtils.getReplicateContainerCommand(CONTAINER_ID,
datanode), pushReplicator);
// WHEN
pushReplicator.replicate(task);
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestMeasuredReplicator.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestMeasuredReplicator.java
index f0423bcc0fc..a381b06b1fe 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestMeasuredReplicator.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestMeasuredReplicator.java
@@ -17,7 +17,7 @@
package org.apache.hadoop.ozone.container.replication;
-import static
org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand.toTarget;
+import static
org.apache.hadoop.ozone.container.common.ContainerTestUtils.getReplicateContainerCommand;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
@@ -69,9 +69,9 @@ public void closeReplicator() throws Exception {
@Test
public void measureFailureSuccessAndBytes() {
//WHEN
- measuredReplicator.replicate(new ReplicationTask(toTarget(1, TARGET),
replicator));
- measuredReplicator.replicate(new ReplicationTask(toTarget(2, TARGET),
replicator));
- measuredReplicator.replicate(new ReplicationTask(toTarget(3, TARGET),
replicator));
+ measuredReplicator.replicate(new
ReplicationTask(getReplicateContainerCommand(1, TARGET), replicator));
+ measuredReplicator.replicate(new
ReplicationTask(getReplicateContainerCommand(2, TARGET), replicator));
+ measuredReplicator.replicate(new
ReplicationTask(getReplicateContainerCommand(3, TARGET), replicator));
//THEN
//even containers should be failed
@@ -89,9 +89,9 @@ public void measureFailureSuccessAndBytes() {
public void testReplicationTime() throws Exception {
//WHEN
//will wait at least the 300ms
- measuredReplicator.replicate(new ReplicationTask(toTarget(101, TARGET),
replicator));
- measuredReplicator.replicate(new ReplicationTask(toTarget(201, TARGET),
replicator));
- measuredReplicator.replicate(new ReplicationTask(toTarget(300, TARGET),
replicator));
+ measuredReplicator.replicate(new
ReplicationTask(getReplicateContainerCommand(101, TARGET), replicator));
+ measuredReplicator.replicate(new
ReplicationTask(getReplicateContainerCommand(201, TARGET), replicator));
+ measuredReplicator.replicate(new
ReplicationTask(getReplicateContainerCommand(300, TARGET), replicator));
//THEN
//even containers should be failed
@@ -109,7 +109,7 @@ public void testReplicationTime() throws Exception {
public void testFailureTimeSuccessExcluded() {
//WHEN
//will wait at least the 15ms
- measuredReplicator.replicate(new ReplicationTask(toTarget(15, TARGET),
replicator));
+ measuredReplicator.replicate(new
ReplicationTask(getReplicateContainerCommand(15, TARGET), replicator));
//THEN
@@ -121,7 +121,7 @@ public void testFailureTimeSuccessExcluded() {
public void testSuccessTimeFailureExcluded() {
//WHEN
//will wait at least the 10ms
- measuredReplicator.replicate(new ReplicationTask(toTarget(10, TARGET),
replicator));
+ measuredReplicator.replicate(new
ReplicationTask(getReplicateContainerCommand(10, TARGET), replicator));
//THEN
@@ -132,7 +132,7 @@ public void testSuccessTimeFailureExcluded() {
@Test
public void testReplicationQueueTimeMetrics() {
final Instant queued = Instant.now().minus(1, ChronoUnit.SECONDS);
- ReplicationTask task = new ReplicationTask(toTarget(100, TARGET),
replicator) {
+ ReplicationTask task = new
ReplicationTask(getReplicateContainerCommand(100, TARGET), replicator) {
@Override
public Instant getQueued() {
return queued;
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestPushReplicator.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestPushReplicator.java
index a4463410cea..97112c15d41 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestPushReplicator.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestPushReplicator.java
@@ -18,8 +18,8 @@
package org.apache.hadoop.ozone.container.replication;
import static org.apache.commons.io.output.NullOutputStream.NULL_OUTPUT_STREAM;
+import static
org.apache.hadoop.ozone.container.common.ContainerTestUtils.getReplicateContainerCommand;
import static
org.apache.hadoop.ozone.container.replication.CopyContainerCompression.NO_COMPRESSION;
-import static
org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand.toTarget;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.mockito.Mockito.any;
import static org.mockito.Mockito.doAnswer;
@@ -68,7 +68,7 @@ void uploadCompletesNormally(CopyContainerCompression
compression)
SpyOutputStream output = new SpyOutputStream(NULL_OUTPUT_STREAM);
ContainerReplicator subject = createSubject(containerID, target,
output, completion, compression);
- ReplicationTask task = new ReplicationTask(toTarget(containerID, target),
+ ReplicationTask task = new
ReplicationTask(getReplicateContainerCommand(containerID, target),
subject);
// WHEN
@@ -89,7 +89,7 @@ void uploadFailsWithException() throws IOException {
fut -> fut.completeExceptionally(new Exception("testing"));
ContainerReplicator subject = createSubject(containerID, target,
output, completion, NO_COMPRESSION);
- ReplicationTask task = new ReplicationTask(toTarget(containerID, target),
+ ReplicationTask task = new
ReplicationTask(getReplicateContainerCommand(containerID, target),
subject);
// WHEN
@@ -111,7 +111,7 @@ void packFailsWithException() throws IOException {
};
ContainerReplicator subject = createSubject(containerID, target,
output, completion, NO_COMPRESSION);
- ReplicationTask task = new ReplicationTask(toTarget(containerID, target),
+ ReplicationTask task = new
ReplicationTask(getReplicateContainerCommand(containerID, target),
subject);
// WHEN
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestReplicationSupervisor.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestReplicationSupervisor.java
index 026ac6ac76a..5f8e9a754cd 100644
---
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestReplicationSupervisor.java
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/container/replication/TestReplicationSupervisor.java
@@ -28,7 +28,6 @@
import static
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicationCommandPriority.NORMAL;
import static
org.apache.hadoop.ozone.container.common.impl.ContainerImplTestUtils.newContainerSet;
import static
org.apache.hadoop.ozone.container.replication.AbstractReplicationTask.Status.DONE;
-import static
org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand.toTarget;
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertTrue;
@@ -74,6 +73,7 @@
import org.apache.hadoop.metrics2.impl.MetricsCollectorImpl;
import org.apache.hadoop.ozone.container.checksum.DNContainerOperationClient;
import org.apache.hadoop.ozone.container.checksum.ReconcileContainerTask;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.container.common.impl.ContainerLayoutVersion;
import org.apache.hadoop.ozone.container.common.impl.ContainerSet;
import
org.apache.hadoop.ozone.container.common.statemachine.DatanodeConfiguration;
@@ -800,8 +800,7 @@ private ECReconstructionCoordinatorTask
createECTaskWithCoordinator(long contain
}
private ReplicateContainerCommand createCommand(long containerId) {
- ReplicateContainerCommand cmd =
- ReplicateContainerCommand.toTarget(containerId, datanode);
+ ReplicateContainerCommand cmd =
ContainerTestUtils.getReplicateContainerCommand(containerId, datanode);
cmd.setTerm(CURRENT_TERM);
return cmd;
}
@@ -1049,7 +1048,7 @@ private void scheduleTasks(
List<DatanodeDetails> datanodes, ReplicationSupervisor rs) {
for (int i = 0; i < 10; i++) {
DatanodeDetails target = datanodes.get(i % datanodes.size());
- rs.addTask(new ReplicationTask(toTarget(i, target), noopReplicator));
+ rs.addTask(new
ReplicationTask(ContainerTestUtils.getReplicateContainerCommand(i, target),
noopReplicator));
}
}
}
diff --git
a/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReplicateContainerCommand.java
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReplicateContainerCommand.java
new file mode 100644
index 00000000000..ff47ccdb04c
--- /dev/null
+++
b/hadoop-hdds/container-service/src/test/java/org/apache/hadoop/ozone/protocol/commands/TestReplicateContainerCommand.java
@@ -0,0 +1,100 @@
+/*
+ * 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.apache.hadoop.hdds.protocol.MockDatanodeDetails.randomDatanodeDetails;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+import org.apache.hadoop.hdds.HDDSVersion;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
+import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.ReplicateContainerCommandProto;
+import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Test cases to verify {@link ReplicateContainerCommand} serialization.
+ */
+public class TestReplicateContainerCommand {
+
+ @Test
+ public void testApparentVersionRoundTrip() {
+ DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
+
+ ReplicateContainerCommand cmd =
+ ReplicateContainerCommand.toTarget(1L, target, HDDSVersion.ZDU);
+
+ ReplicateContainerCommandProto proto = cmd.getProto();
+ ReplicateContainerCommand deserialized =
+ ReplicateContainerCommand.getFromProtobuf(proto);
+
+ assertEquals(HDDSVersion.ZDU, deserialized.getApparentVersion());
+ assertEquals(target.getID(),
+ deserialized.getTargetDatanode().getID());
+ }
+
+ @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.
+ DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
+
+ ReplicateContainerCommand cmd =
+ ReplicateContainerCommand.toTarget(1L, target,
+ HDDSLayoutFeature.HBASE_SUPPORT);
+
+ ReplicateContainerCommandProto proto = cmd.getProto();
+ ReplicateContainerCommand deserialized =
+ ReplicateContainerCommand.getFromProtobuf(proto);
+
+ assertEquals(HDDSLayoutFeature.HBASE_SUPPORT,
+ deserialized.getApparentVersion());
+ }
+
+ @Test
+ public void testApparentVersionDefaultWhenAbsent() {
+ DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
+
+ ReplicateContainerCommand cmd =
+ ReplicateContainerCommand.toTarget(1L, target, HDDSVersion.ZDU);
+
+ ReplicateContainerCommandProto proto =
+ ReplicateContainerCommandProto.newBuilder()
+ .setContainerID(1L)
+ .setCmdId(cmd.getId())
+ .setTarget(target.getProtoBufMessage())
+ .build();
+
+ ReplicateContainerCommand deserialized =
+ ReplicateContainerCommand.getFromProtobuf(proto);
+
+ assertEquals(HDDSVersion.DEFAULT_VERSION,
+ deserialized.getApparentVersion());
+ }
+
+ @Test
+ public void testToStringIncludesApparentVersion() {
+ ReplicateContainerCommand cmd =
+ ReplicateContainerCommand.toTarget(1L, randomDatanodeDetails(),
HDDSVersion.SOFTWARE_VERSION);
+
+ String str = cmd.toString();
+ assertTrue(str.contains("apparentVersion=" +
HDDSVersion.SOFTWARE_VERSION));
+ }
+}
diff --git
a/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto
b/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto
index 4d9bb74a04f..e0a2a390bf3 100644
---
a/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto
+++
b/hadoop-hdds/interface-server/src/main/proto/ScmServerDatanodeHeartbeatProtocol.proto
@@ -447,6 +447,7 @@ message ReplicateContainerCommandProto {
optional int32 replicaIndex = 4;
optional DatanodeDetailsProto target = 5;
optional ReplicationCommandPriority priority = 6 [default = NORMAL];
+ optional uint32 apparentVersion = 7;
}
/**
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 f890fb6a082..bcc0440b7d2 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
@@ -527,10 +527,17 @@ public void sendThrottledReplicationCommand(ContainerInfo
containerInfo,
DatanodeDetails source = selectAndOptionallyExcludeDatanode(
1, sourceWithCmds);
- ReplicateContainerCommand cmd =
- ReplicateContainerCommand.toTarget(containerID, target);
- cmd.setReplicaIndex(replicaIndex);
- sendDatanodeCommand(cmd, containerInfo, source);
+ try {
+ ReplicateContainerCommand cmd = ReplicateContainerCommand.toTarget(
+ containerID, target,
+ nodeManager.getLowestApparentVersion(source, target));
+ cmd.setReplicaIndex(replicaIndex);
+ sendDatanodeCommand(cmd, containerInfo, source);
+ } catch (NodeNotFoundException e) {
+ throw new IllegalArgumentException("Datanode not found in NodeManager
while sending replication "
+ + "command for container " + containerID + " from source " + source
+ " to target " + target
+ + ". Should not happen", e);
+ }
}
public void sendThrottledReconstructionCommand(ContainerInfo containerInfo,
@@ -627,11 +634,18 @@ public void sendLowPriorityReplicateContainerCommand(
final ContainerInfo container, int replicaIndex, DatanodeDetails source,
DatanodeDetails target, long scmDeadlineEpochMs)
throws NotLeaderException {
- final ReplicateContainerCommand command = ReplicateContainerCommand
- .toTarget(container.getContainerID(), target);
- command.setReplicaIndex(replicaIndex);
- command.setPriority(ReplicationCommandPriority.LOW);
- sendDatanodeCommand(command, container, source, scmDeadlineEpochMs);
+ try {
+ final ReplicateContainerCommand command =
ReplicateContainerCommand.toTarget(
+ container.getContainerID(), target,
+ nodeManager.getLowestApparentVersion(source, target));
+ command.setReplicaIndex(replicaIndex);
+ command.setPriority(ReplicationCommandPriority.LOW);
+ sendDatanodeCommand(command, container, source, scmDeadlineEpochMs);
+ } catch (NodeNotFoundException e) {
+ throw new IllegalArgumentException("Datanode not found in NodeManager
while sending replication "
+ + "command for container " + container.getContainerID() + " from
source " + source
+ + " to target " + target + ". Should not happen", e);
+ }
}
/**
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
index 3f232cb8a51..0678d992c20 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeManager.java
@@ -206,6 +206,28 @@ default DatanodeFinalizationCounts
getDatanodeFinalizationCounts() {
.build();
}
+ /**
+ * Returns the lowest apparent version among the given datanodes,
+ * so every node involved in an operation uses the same, mutually-supported
+ * version.
+ *
+ * @throws NodeNotFoundException if SCM has no record of one of the nodes;
+ * callers must not proceed with an operation involving a node SCM does
+ * not know about.
+ */
+ default ComponentVersion getLowestApparentVersion(DatanodeDetails... nodes)
+ throws NodeNotFoundException {
+ ComponentVersion[] versions = new ComponentVersion[nodes.length];
+ for (int i = 0; i < nodes.length; i++) {
+ DatanodeInfo info = getNode(nodes[i].getID());
+ if (info == null) {
+ throw new NodeNotFoundException(nodes[i].getID());
+ }
+ versions[i] = info.getLastKnownApparentVersion();
+ }
+ return ComponentVersion.min(versions);
+ }
+
/**
* Returns the aggregated node stats.
* @return the aggregated node stats.
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationTestUtil.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationTestUtil.java
index 63f4c6c9f33..f715c766b08 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationTestUtil.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/ReplicationTestUtil.java
@@ -55,6 +55,7 @@
import org.apache.hadoop.hdds.scm.net.Node;
import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.utils.HddsServerUtil;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.protocol.commands.DeleteContainerCommand;
import
org.apache.hadoop.ozone.protocol.commands.ReconstructECContainersCommand;
import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
@@ -457,8 +458,8 @@ public static void
mockRMSendThrottleReplicateCommand(ReplicationManager mock,
}
List<DatanodeDetails> sources = invocationOnMock.getArgument(1);
ContainerInfo containerInfo = invocationOnMock.getArgument(0);
- ReplicateContainerCommand command = ReplicateContainerCommand
- .toTarget(containerInfo.getContainerID(),
+ ReplicateContainerCommand command = ContainerTestUtils
+ .getReplicateContainerCommand(containerInfo.getContainerID(),
invocationOnMock.getArgument(2));
command.setReplicaIndex(invocationOnMock.getArgument(3));
commandsSent.add(Pair.of(sources.get(0), command));
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestContainerReplicaPendingOps.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestContainerReplicaPendingOps.java
index 217e752278f..9cf626c02b8 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestContainerReplicaPendingOps.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestContainerReplicaPendingOps.java
@@ -42,8 +42,8 @@
import org.apache.hadoop.hdds.protocol.DatanodeID;
import org.apache.hadoop.hdds.protocol.MockDatanodeDetails;
import org.apache.hadoop.hdds.scm.container.ContainerID;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.protocol.commands.DeleteContainerCommand;
-import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
import org.apache.hadoop.ozone.protocol.commands.SCMCommand;
import org.apache.ozone.test.MockClock;
import org.junit.jupiter.api.AfterEach;
@@ -84,7 +84,7 @@ public void setup() {
dn2 = MockDatanodeDetails.randomDatanodeDetails();
dn3 = MockDatanodeDetails.randomDatanodeDetails();
- addCmd = ReplicateContainerCommand.toTarget(1, dn3);
+ addCmd = ContainerTestUtils.getReplicateContainerCommand(1, dn3);
deleteCmd = new DeleteContainerCommand(1, false);
}
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 bcb3dea9768..500ed6119a0 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
@@ -65,6 +65,8 @@
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.function.Function;
import org.apache.commons.lang3.tuple.Pair;
+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;
@@ -89,6 +91,7 @@
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.ha.SCMContext;
import org.apache.hadoop.hdds.scm.ha.SCMServiceManager;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.node.NodeStatus;
import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
@@ -96,6 +99,8 @@
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
import org.apache.hadoop.hdds.security.token.ContainerTokenGenerator;
import org.apache.hadoop.hdds.server.events.EventPublisher;
+import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.protocol.commands.DeleteContainerCommand;
import
org.apache.hadoop.ozone.protocol.commands.ReconstructECContainersCommand;
import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
@@ -107,6 +112,7 @@
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;
@@ -160,6 +166,15 @@ public void setup() throws IOException {
return null;
}).when(nodeManager).addDatanodeCommand(any(), any());
+ // By default every datanode is known and reports the current version, so
+ // the replication 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);
+
clock = new MockClock(Instant.now(), ZoneId.systemDefault());
containerReplicaPendingOps =
new ContainerReplicaPendingOps(clock, null);
@@ -1207,7 +1222,7 @@ public void testReplicateContainerCommandToTarget()
// command will be pushed from source to target
DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
DatanodeDetails source = MockDatanodeDetails.randomDatanodeDetails();
- ReplicateContainerCommand command = ReplicateContainerCommand.toTarget(
+ ReplicateContainerCommand command =
ContainerTestUtils.getReplicateContainerCommand(
containerInfo.getContainerID(), target);
command.setReplicaIndex(1);
replicationManager.sendDatanodeCommand(command, containerInfo, source);
@@ -1739,4 +1754,85 @@ private void mockReplicationCommandCounts(
});
}
+ /**
+ * Captures the ReplicateContainerCommand SCM sends to a datanode.
+ */
+ private ReplicateContainerCommand captureSentReplicateCommand() {
+ ArgumentCaptor<SCMCommand<?>> command =
+ ArgumentCaptor.forClass(SCMCommand.class);
+ verify(nodeManager).addDatanodeCommand(any(DatanodeID.class),
+ command.capture());
+ return (ReplicateContainerCommand) command.getValue();
+ }
+
+ private DatanodeInfo mockDatanodeWithApparentVersion(
+ DatanodeDetails dn, ComponentVersion version) {
+ DatanodeInfo info = mock(DatanodeInfo.class);
+ when(info.getLastKnownApparentVersion()).thenReturn(version);
+ when(nodeManager.getNode(dn.getID())).thenReturn(info);
+ return info;
+ }
+
+ /**
+ * 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.
+ */
+ @ParameterizedTest
+ @CsvSource({
+ "true, true",
+ "true, false",
+ "false, true",
+ "false, false",
+ })
+ public void testApparentVersionIsLowestOfSourceAndTarget(
+ boolean throttled, boolean sourceNewer)
+ throws CommandTargetOverloadedException, NotLeaderException,
+ NodeNotFoundException {
+ ContainerInfo containerInfo =
+ ReplicationTestUtil.createContainerInfo(repConfig, 1,
+ HddsProtos.LifeCycleState.CLOSED, 10, 20);
+ DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
+ DatanodeDetails source = MockDatanodeDetails.randomDatanodeDetails();
+
+ ComponentVersion lower = HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION;
+ ComponentVersion higher = HDDSVersion.ZDU;
+ mockDatanodeWithApparentVersion(source, sourceNewer ? higher : lower);
+ mockDatanodeWithApparentVersion(target, sourceNewer ? lower : higher);
+ when(nodeManager.getLowestApparentVersion(source, target))
+ .thenCallRealMethod();
+
+ if (throttled) {
+ mockReplicationCommandCounts(dn -> 0, dn -> 0);
+ replicationManager.sendThrottledReplicationCommand(containerInfo,
+ Collections.singletonList(source), target, 1);
+ } else {
+
replicationManager.sendLowPriorityReplicateContainerCommand(containerInfo,
+ 1, source, target, clock.millis() + rmConf.getEventTimeout());
+ }
+
+ assertEquals(lower, captureSentReplicateCommand().getApparentVersion());
+ }
+
+ @Test
+ public void testApparentVersionLookupThrowsWhenNodeNotFound()
+ throws NodeNotFoundException {
+ ContainerInfo containerInfo =
+ ReplicationTestUtil.createContainerInfo(repConfig, 1,
+ HddsProtos.LifeCycleState.CLOSED, 10, 20);
+ DatanodeDetails target = MockDatanodeDetails.randomDatanodeDetails();
+ DatanodeDetails source = MockDatanodeDetails.randomDatanodeDetails();
+
+ mockDatanodeWithApparentVersion(source, HDDSVersion.SOFTWARE_VERSION);
+ // SCM has no information for the target.
+ when(nodeManager.getNode(target.getID())).thenReturn(null);
+ when(nodeManager.getLowestApparentVersion(source, target))
+ .thenCallRealMethod();
+
+ // We must not proceed with a replication command for a node we don't know.
+ assertThrows(IllegalArgumentException.class, () ->
+ replicationManager.sendLowPriorityReplicateContainerCommand(
+ containerInfo, 0, source, target,
+ clock.millis() + rmConf.getEventTimeout()));
+ }
}
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerScenarios.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerScenarios.java
index 7662ed5ba78..15cce574ee9 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerScenarios.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/container/replication/TestReplicationManagerScenarios.java
@@ -45,6 +45,7 @@
import java.util.Set;
import java.util.stream.Stream;
import org.apache.commons.lang3.tuple.Pair;
+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;
@@ -63,6 +64,7 @@
import org.apache.hadoop.hdds.scm.container.ContainerReplica;
import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport;
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.node.NodeStatus;
import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
@@ -201,6 +203,14 @@ public void setup() throws IOException,
NodeNotFoundException {
return NODE_STATUS_MAP.getOrDefault(dn,
NodeStatus.inServiceHealthy());
});
+ // Replication commands look up each involved datanode's apparent version,
+ // so return an info reporting the current version for every node.
+ DatanodeInfo defaultNodeInfo = mock(DatanodeInfo.class);
+ when(defaultNodeInfo.getLastKnownApparentVersion())
+ .thenReturn(HDDSVersion.SOFTWARE_VERSION);
+ when(nodeManager.getNode(any(DatanodeID.class)))
+ .thenReturn(defaultNodeInfo);
+
final HashMap<SCMCommandProto.Type, Integer> countMap = new HashMap<>();
for (SCMCommandProto.Type type : SCMCommandProto.Type.values()) {
countMap.put(type, 0);
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestCommandQueue.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestCommandQueue.java
index d0cfa0be699..a70cbca4477 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestCommandQueue.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestCommandQueue.java
@@ -28,6 +28,7 @@
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMCommandProto;
import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.protocol.commands.CloseContainerCommand;
import org.apache.hadoop.ozone.protocol.commands.CreatePipelineCommand;
import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
@@ -49,11 +50,11 @@ public void testSummaryUpdated() {
new CreatePipelineCommand(PipelineID.randomId(),
HddsProtos.ReplicationType.RATIS,
HddsProtos.ReplicationFactor.THREE, Collections.emptyList());
- SCMCommand<?> replicationCommand = ReplicateContainerCommand.toTarget(
- containerID, MockDatanodeDetails.randomDatanodeDetails());
+ SCMCommand<?> replicationCommand = ContainerTestUtils
+ .getReplicateContainerCommand(containerID,
MockDatanodeDetails.randomDatanodeDetails());
- ReplicateContainerCommand lowReplicationCommand = ReplicateContainerCommand
- .toTarget(containerID, MockDatanodeDetails.randomDatanodeDetails());
+ ReplicateContainerCommand lowReplicationCommand = ContainerTestUtils
+ .getReplicateContainerCommand(containerID,
MockDatanodeDetails.randomDatanodeDetails());
lowReplicationCommand.setPriority(StorageContainerDatanodeProtocolProtos
.ReplicationCommandPriority.LOW);
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java
index 2a253746727..c689d30c3b6 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestSCMNodeManager.java
@@ -102,13 +102,13 @@
import org.apache.hadoop.hdds.server.events.EventPublisher;
import org.apache.hadoop.hdds.server.events.EventQueue;
import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.container.upgrade.UpgradeUtils;
import org.apache.hadoop.ozone.protocol.commands.CloseContainerCommand;
import org.apache.hadoop.ozone.protocol.commands.CommandForDatanode;
import org.apache.hadoop.ozone.protocol.commands.CreatePipelineCommand;
import org.apache.hadoop.ozone.protocol.commands.DeleteBlocksCommand;
import org.apache.hadoop.ozone.protocol.commands.RegisteredCommand;
-import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
import org.apache.hadoop.ozone.protocol.commands.SCMCommand;
import
org.apache.hadoop.ozone.protocol.commands.SetNodeOperationalStateCommand;
import
org.apache.hadoop.security.authentication.client.AuthenticationException;
@@ -941,8 +941,8 @@ scmStorageConfig, eventPublisher, new
NetworkTopologyImpl(conf),
verify(eventPublisher,
times(1)).fireEvent(NEW_NODE, node1);
for (int i = 0; i < 3; i++) {
- nodeManager.addDatanodeCommand(node1.getID(), ReplicateContainerCommand
- .toTarget(1, MockDatanodeDetails.randomDatanodeDetails()));
+ nodeManager.addDatanodeCommand(node1.getID(), ContainerTestUtils
+ .getReplicateContainerCommand(1,
MockDatanodeDetails.randomDatanodeDetails()));
}
for (int i = 0; i < 5; i++) {
nodeManager.addDatanodeCommand(node1.getID(),
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMDatanodeProtocolServer.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMDatanodeProtocolServer.java
index 1c9a138e33f..cc20985db50 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMDatanodeProtocolServer.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/TestSCMDatanodeProtocolServer.java
@@ -26,6 +26,7 @@
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos;
import org.apache.hadoop.hdds.scm.server.OzoneStorageContainerManager;
import org.apache.hadoop.hdds.scm.server.SCMDatanodeProtocolServer;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.protocol.commands.ReplicateContainerCommand;
import org.junit.jupiter.api.Test;
@@ -40,7 +41,8 @@ public void ensureTermAndDeadlineOnCommands()
OzoneStorageContainerManager scm =
mock(OzoneStorageContainerManager.class);
- ReplicateContainerCommand command = ReplicateContainerCommand.toTarget(1,
randomDatanodeDetails());
+ ReplicateContainerCommand command = ContainerTestUtils
+ .getReplicateContainerCommand(1, randomDatanodeDetails());
command.setTerm(5L);
command.setDeadline(1234L);
StorageContainerDatanodeProtocolProtos.SCMCommandProto proto =
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerReplication.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerReplication.java
index 2fbd737625a..ced0f999942 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerReplication.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/container/replication/TestContainerReplication.java
@@ -51,6 +51,7 @@
import org.apache.hadoop.ozone.MiniOzoneCluster;
import org.apache.hadoop.ozone.OzoneConfigKeys;
import org.apache.hadoop.ozone.container.ContainerTestHelper;
+import org.apache.hadoop.ozone.container.common.ContainerTestUtils;
import org.apache.hadoop.ozone.container.common.interfaces.Container;
import org.apache.hadoop.ozone.container.common.interfaces.DBHandle;
import
org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine;
@@ -114,7 +115,7 @@ void testPush(CopyContainerCompression compression) throws
Exception {
long containerID = createNewClosedContainer(source);
DatanodeDetails target = selectOtherNode(source);
ReplicateContainerCommand cmd =
- ReplicateContainerCommand.toTarget(containerID, target);
+ ContainerTestUtils.getReplicateContainerCommand(containerID, target);
queueAndWaitForCompletion(cmd, source,
ReplicationSupervisor::getReplicationSuccessCount);
@@ -136,7 +137,7 @@ void pushSucceedsWhenSourceBytesUsedIsNegative(long
containerSize) throws Except
poisonBytesUsed(source, containerID, containerSize);
ReplicateContainerCommand cmd =
- ReplicateContainerCommand.toTarget(containerID, target);
+ ContainerTestUtils.getReplicateContainerCommand(containerID, target);
queueAndWaitForCompletion(cmd, source,
ReplicationSupervisor::getReplicationSuccessCount);
@@ -177,8 +178,7 @@ void pushUnknownContainer() throws Exception {
.getDatanodeDetails();
DatanodeDetails target = selectOtherNode(source);
ReplicateContainerCommand cmd =
- ReplicateContainerCommand.toTarget(CONTAINER_ID.incrementAndGet(),
- target);
+
ContainerTestUtils.getReplicateContainerCommand(CONTAINER_ID.incrementAndGet(),
target);
queueAndWaitForCompletion(cmd, source,
ReplicationSupervisor::getReplicationFailureCount);
@@ -223,8 +223,7 @@ void testPushWithOverAllocatedContainer(String testName,
Long containerSize)
assertEquals(originalSize, containerSize);
// Create replication command to push container to target
- ReplicateContainerCommand cmd =
- ReplicateContainerCommand.toTarget(containerID, target);
+ ReplicateContainerCommand cmd =
ContainerTestUtils.getReplicateContainerCommand(containerID, target);
// Execute push replication
queueAndWaitForCompletion(cmd, source,
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]