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 4443f642315 HDDS-16044. SCM should forward Datanode's current version
to clients for read and write operations (#11023)
4443f642315 is described below
commit 4443f642315e47e0b71a85f731399ed94ce682a8
Author: Ethan Rose <[email protected]>
AuthorDate: Thu Aug 27 16:16:54 2026 -0400
HDDS-16044. SCM should forward Datanode's current version to clients for
read and write operations (#11023)
---
.../org/apache/hadoop/hdds/ComponentVersion.java | 16 +-
.../hadoop/hdds/protocol/DatanodeDetails.java | 2 +-
.../common/helpers/ContainerWithPipeline.java | 12 +-
.../apache/hadoop/hdds/scm/pipeline/Pipeline.java | 31 ++-
.../hadoop/hdds/AbstractComponentVersionTest.java | 19 +-
.../hadoop/hdds/protocol/TestDatanodeDetails.java | 6 +-
.../hadoop/hdds/scm/pipeline/TestPipeline.java | 78 +++++++
.../states/endpoint/HeartbeatEndpointTask.java | 16 +-
.../container/replication/ReplicationManager.java | 50 +++--
.../apache/hadoop/hdds/scm/node/DatanodeInfo.java | 13 ++
.../apache/hadoop/hdds/scm/node/NodeManager.java | 22 --
.../hadoop/hdds/scm/node/NodeStateManager.java | 11 +
.../hadoop/hdds/scm/node/SCMNodeManager.java | 1 +
...lockLocationProtocolServerSideTranslatorPB.java | 56 ++---
...inerLocationProtocolServerSideTranslatorPB.java | 26 ++-
.../hdds/scm/server/upgrade/ScmVersionManager.java | 28 +++
.../replication/TestReplicationManager.java | 9 +-
.../hadoop/hdds/scm/node/TestSCMNodeManager.java | 17 ++
...lockLocationProtocolServerSideTranslatorPB.java | 230 ++++++++++---------
...inerLocationProtocolServerSideTranslatorPB.java | 224 +++++++++++++++++++
.../rpc/TestDatanodeCurrentVersionEndToEnd.java | 244 +++++++++++++++++++++
21 files changed, 882 insertions(+), 229 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 77c048e31b4..e1554f9469d 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
@@ -87,19 +87,11 @@ default Optional<? extends UpgradeAction> action() {
* 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];
- }
+ static <T extends ComponentVersion> T min(T v1, T v2) {
+ if (v1.isSupportedBy(v2)) {
+ return v1;
}
- return lowest;
+ return v2;
}
}
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
index 8925d325942..c354f7f4d93 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/protocol/DatanodeDetails.java
@@ -654,7 +654,7 @@ public void setInitialVersion(HDDSVersion initialVersion) {
}
/**
- * @return the version this datanode was last started with
+ * @return the version this datanode should report to clients
*/
public HDDSVersion getCurrentVersion() {
return currentVersion;
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/common/helpers/ContainerWithPipeline.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/common/helpers/ContainerWithPipeline.java
index fad799a25d5..08890c46d08 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/common/helpers/ContainerWithPipeline.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/container/common/helpers/ContainerWithPipeline.java
@@ -18,9 +18,12 @@
package org.apache.hadoop.hdds.scm.container.common.helpers;
import java.util.Comparator;
+import java.util.Map;
import org.apache.commons.lang3.builder.EqualsBuilder;
import org.apache.commons.lang3.builder.HashCodeBuilder;
+import org.apache.hadoop.hdds.ComponentVersion;
import org.apache.hadoop.hdds.protocol.DatanodeDetails.Port.Name;
+import org.apache.hadoop.hdds.protocol.DatanodeID;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.scm.container.ContainerInfo;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
@@ -54,10 +57,15 @@ public static ContainerWithPipeline
fromProtobuf(HddsProtos.ContainerWithPipelin
Pipeline.getFromProtobuf(allocatedContainer.getPipeline()));
}
- public HddsProtos.ContainerWithPipeline getProtobuf(ClientVersion
clientVersion) {
+ /**
+ * Serializes with a per-datanode currentVersion override (keyed by datanode
id), so read clients see each
+ * datanode's up-to-date version rather than the pipeline's possibly-stale
frozen copy.
+ */
+ public HddsProtos.ContainerWithPipeline getProtobuf(ClientVersion
clientVersion,
+ Map<DatanodeID, ComponentVersion> memberVersions) {
return HddsProtos.ContainerWithPipeline.newBuilder()
.setContainerInfo(getContainerInfo().getProtobuf())
- .setPipeline(getPipeline().getProtobufMessage(clientVersion,
Name.IO_PORTS))
+ .setPipeline(getPipeline().getProtobufMessage(clientVersion,
Name.IO_PORTS, memberVersions))
.build();
}
diff --git
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
index 1b04d1fbfab..8c084d7c8eb 100644
---
a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
+++
b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/pipeline/Pipeline.java
@@ -365,21 +365,40 @@ public ReplicationConfig getReplicationConfig() {
}
public HddsProtos.Pipeline getProtobufMessage(ClientVersion clientVersion) {
- return getProtobufMessage(clientVersion, Collections.emptySet());
+ return getProtobufMessageInternal(clientVersion, Collections.emptySet(),
null);
}
- public HddsProtos.Pipeline getProtobufMessage(ClientVersion clientVersion,
- Set<DatanodeDetails.Port.Name> filterPorts) {
- return getProtobufMessage(clientVersion, filterPorts, null);
+ /**
+ * Write-path override: when {@code datanodeVersion} is non-null it is set
as the currentVersion on <b>every</b>
+ * member proto, so clients target a single pipeline-wide version (typically
the pipeline minimum).
+ */
+ public HddsProtos.Pipeline getProtobufMessage(ClientVersion clientVersion,
Set<DatanodeDetails.Port.Name> filterPorts,
+ ComponentVersion datanodeVersion) {
+ return getProtobufMessageInternal(clientVersion, filterPorts,
+ datanodeVersion == null ? null : nodeId -> datanodeVersion);
}
+ /**
+ * Read-path override: set each member proto's currentVersion from {@code
memberVersions} (keyed by datanode id),
+ * so clients see each datanode's own up-to-date version. Members absent
from the map keep their own version.
+ */
public HddsProtos.Pipeline getProtobufMessage(ClientVersion clientVersion,
Set<DatanodeDetails.Port.Name> filterPorts,
- ComponentVersion versionOverride) {
+ Map<DatanodeID, ComponentVersion> memberVersions) {
+ return getProtobufMessageInternal(clientVersion, filterPorts,
+ memberVersions == null ? null : memberVersions::get);
+ }
+
+ private HddsProtos.Pipeline getProtobufMessageInternal(ClientVersion
clientVersion,
+ Set<DatanodeDetails.Port.Name> filterPorts, Function<DatanodeID,
ComponentVersion> versionOverride) {
List<HddsProtos.DatanodeDetailsProto> members = new ArrayList<>();
List<Integer> memberReplicaIndexes = new ArrayList<>();
for (DatanodeDetails dn : nodeStatus.keySet()) {
- members.add(dn.toProto(clientVersion, filterPorts, versionOverride));
+ HddsProtos.DatanodeDetailsProto.Builder memberBuilder =
dn.toProtoBuilder(clientVersion, filterPorts);
+ if (versionOverride != null) {
+
memberBuilder.setCurrentVersion(versionOverride.apply(dn.getID()).serialize());
+ }
+ members.add(memberBuilder.build());
memberReplicaIndexes.add(replicaIndexes.getOrDefault(dn, 0));
}
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 27671edfd74..821979b5b0d 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,7 +22,6 @@
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;
@@ -133,20 +132,10 @@ public void testDeserializeUnknownVersion() {
}
@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() {
+ public void testMinReturnsLowestKnownVersionInAnyOrder() {
// getValues()[0] is the lowest known version
- assertEquals(getValues()[0], ComponentVersion.min(getValues()));
+ assertEquals(getValues()[0], ComponentVersion.min(getValues()[0],
getValues()[1]));
+ assertEquals(getValues()[0], ComponentVersion.min(getValues()[1],
getValues()[0]));
}
@Test
@@ -157,6 +146,6 @@ public void testMinTreatsUnknownFutureVersionAsHighest() {
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));
+ assertEquals(unknown, ComponentVersion.min(unknown, unknown));
}
}
diff --git
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/protocol/TestDatanodeDetails.java
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/protocol/TestDatanodeDetails.java
index 2a6416b0185..a3a23de34fc 100644
---
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/protocol/TestDatanodeDetails.java
+++
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/protocol/TestDatanodeDetails.java
@@ -75,13 +75,15 @@ public void testNewBuilderCurrentVersion() {
DatanodeDetails dn = MockDatanodeDetails.randomDatanodeDetails();
Set<Port.Name> requiredPorts = Stream.of(Port.Name.STANDALONE,
Port.Name.RATIS)
.collect(Collectors.toSet());
- HddsProtos.DatanodeDetailsProto.Builder protoBuilder =
dn.toProtoBuilder(ClientVersion.CURRENT, requiredPorts);
+ HddsProtos.DatanodeDetailsProto.Builder protoBuilder =
+ dn.toProtoBuilder(ClientVersion.CURRENT, requiredPorts);
protoBuilder.clearCurrentVersion();
DatanodeDetails dn2 =
DatanodeDetails.newBuilder(protoBuilder.build()).build();
assertEquals(HDDSVersion.DEFAULT_VERSION, dn2.getCurrentVersion());
// When the proto field is present, it round-trips correctly.
- protoBuilder = dn.toProtoBuilder(ClientVersion.CURRENT, requiredPorts);
+ protoBuilder =
+ dn.toProtoBuilder(ClientVersion.CURRENT, requiredPorts);
DatanodeDetails dn3 = DatanodeDetails.newBuilder(
protoBuilder.setCurrentVersion(HDDSVersion.SOFTWARE_VERSION.serialize()).build())
.build();
diff --git
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipeline.java
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipeline.java
index 185f82fbe3c..e9431ceef5e 100644
---
a/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipeline.java
+++
b/hadoop-hdds/common/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipeline.java
@@ -32,9 +32,15 @@
import java.io.IOException;
import java.util.Arrays;
+import java.util.HashMap;
+import java.util.Map;
+import org.apache.hadoop.hdds.ComponentVersion;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.DatanodeID;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
+import org.apache.hadoop.ozone.ClientVersion;
import org.junit.jupiter.api.Test;
/**
@@ -59,6 +65,78 @@ public void protoIncludesNewPortsOnlyForV1() throws
IOException {
}
}
+ @Test
+ public void testDefaultPreservesCurrentVersion() {
+ // Read path: each member's own reported currentVersion must survive
serialization.
+ Map<DatanodeID, Integer> expected = new HashMap<>();
+ DatanodeDetails dn0 = randomDatanodeDetails();
+ dn0.setCurrentVersion(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC);
+ DatanodeDetails dn1 = randomDatanodeDetails();
+ dn1.setCurrentVersion(HDDSVersion.STREAM_BLOCK_SUPPORT);
+ DatanodeDetails dn2 = randomDatanodeDetails();
+ dn2.setCurrentVersion(HDDSVersion.SOFTWARE_VERSION);
+ expected.put(dn0.getID(),
HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC.serialize());
+ expected.put(dn1.getID(), HDDSVersion.STREAM_BLOCK_SUPPORT.serialize());
+ expected.put(dn2.getID(), HDDSVersion.SOFTWARE_VERSION.serialize());
+
+ Pipeline subject = MockPipeline.createPipeline(Arrays.asList(dn0, dn1,
dn2));
+
+ assertPerMemberVersions(expected,
subject.getProtobufMessage(ClientVersion.CURRENT));
+ }
+
+ @Test
+ public void testOverrideAllMemberVersions() {
+ // Write path: the explicit datanodeVersion overrides every member's
currentVersion.
+ final HDDSVersion override = HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC;
+ DatanodeDetails dn0 = randomDatanodeDetails();
+ dn0.setCurrentVersion(override);
+ DatanodeDetails dn1 = randomDatanodeDetails();
+ dn1.setCurrentVersion(HDDSVersion.STREAM_BLOCK_SUPPORT);
+ DatanodeDetails dn2 = randomDatanodeDetails();
+ dn2.setCurrentVersion(HDDSVersion.SOFTWARE_VERSION);
+
+ Pipeline subject = MockPipeline.createPipeline(Arrays.asList(dn0, dn1,
dn2));
+
+ HddsProtos.Pipeline proto =
+ subject.getProtobufMessage(ClientVersion.CURRENT, ALL_PORTS, override);
+ for (HddsProtos.DatanodeDetailsProto member : proto.getMembersList()) {
+ assertEquals(override.serialize(), member.getCurrentVersion());
+ }
+ }
+
+ @Test
+ public void testOverrideEachMemberVersion() {
+ // Read path: the translator substitutes each member's live currentVersion
via a per-datanode map,
+ // overriding the stale version frozen into the pipeline member.
+ DatanodeDetails dn0 = randomDatanodeDetails();
+ DatanodeDetails dn1 = randomDatanodeDetails();
+ DatanodeDetails dn2 = randomDatanodeDetails();
+ for (DatanodeDetails dn : Arrays.asList(dn0, dn1, dn2)) {
+ dn.setCurrentVersion(HDDSVersion.DEFAULT_VERSION);
+ }
+
+ Map<DatanodeID, ComponentVersion> overrides = new HashMap<>();
+ overrides.put(dn0.getID(), HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC);
+ overrides.put(dn1.getID(), HDDSVersion.STREAM_BLOCK_SUPPORT);
+ overrides.put(dn2.getID(), HDDSVersion.SOFTWARE_VERSION);
+
+ Pipeline subject = MockPipeline.createPipeline(Arrays.asList(dn0, dn1,
dn2));
+
+ Map<DatanodeID, Integer> expected = new HashMap<>();
+ expected.put(dn0.getID(),
HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC.serialize());
+ expected.put(dn1.getID(), HDDSVersion.STREAM_BLOCK_SUPPORT.serialize());
+ expected.put(dn2.getID(), HDDSVersion.SOFTWARE_VERSION.serialize());
+ assertPerMemberVersions(expected,
subject.getProtobufMessage(DEFAULT_VERSION, ALL_PORTS, overrides));
+ }
+
+ private static void assertPerMemberVersions(Map<DatanodeID, Integer>
expected, HddsProtos.Pipeline proto) {
+ assertEquals(expected.size(), proto.getMembersCount());
+ for (HddsProtos.DatanodeDetailsProto member : proto.getMembersList()) {
+ DatanodeID id = DatanodeDetails.getFromProtoBuf(member).getID();
+ assertEquals(expected.get(id).intValue(), member.getCurrentVersion());
+ }
+ }
+
@Test
public void getProtobufMessageEC() throws IOException {
Pipeline subject = MockPipeline.createPipeline(3);
diff --git
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/HeartbeatEndpointTask.java
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/HeartbeatEndpointTask.java
index 100070b19d5..c752b6bbc5b 100644
---
a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/HeartbeatEndpointTask.java
+++
b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/common/states/endpoint/HeartbeatEndpointTask.java
@@ -35,6 +35,7 @@
import java.util.LinkedList;
import java.util.List;
import java.util.concurrent.Callable;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.conf.ConfigurationSource;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeDetailsProto;
@@ -49,6 +50,7 @@
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.SCMHeartbeatResponseProto;
import org.apache.hadoop.hdds.utils.ConnectionFailureUtils;
import org.apache.hadoop.hdfs.util.EnumCounters;
+import org.apache.hadoop.ozone.HddsDatanodeService;
import
org.apache.hadoop.ozone.container.common.helpers.DeletedContainerBlocksSummary;
import
org.apache.hadoop.ozone.container.common.statemachine.EndpointStateMachine;
import
org.apache.hadoop.ozone.container.common.statemachine.EndpointStateMachine.EndPointStates;
@@ -82,6 +84,8 @@ public class HeartbeatEndpointTask
private int maxContainerActionsPerHB;
private int maxPipelineActionsPerHB;
private final DatanodeVersionManager versionManager;
+ // Test-only override for the version this datanode advertises to clients;
null in production.
+ private final HDDSVersion testCurrentVersion;
private final boolean resolveOnFailureEnabled;
private final int refreshThreshold;
@@ -102,6 +106,15 @@ public HeartbeatEndpointTask(EndpointStateMachine
rpcEndpoint,
HDDS_PIPELINE_ACTION_MAX_LIMIT_DEFAULT);
this.datanodeDetails = context.getParent().getDatanodeDetails();
this.versionManager = context.getParent().getVersionManager();
+
+ if
(conf.isConfigured(HddsDatanodeService.TESTING_DATANODE_VERSION_CURRENT)) {
+ int configuredVersion =
conf.getInt(HddsDatanodeService.TESTING_DATANODE_VERSION_CURRENT,
+ HDDSVersion.SOFTWARE_VERSION.serialize());
+ testCurrentVersion = HDDSVersion.deserialize(configuredVersion);
+ } else {
+ testCurrentVersion = null;
+ }
+
this.resolveOnFailureEnabled =
conf.getBoolean(OZONE_CLIENT_FAILOVER_RESOLVE_NEEDED_KEY,
OZONE_CLIENT_FAILOVER_RESOLVE_NEEDED_DEFAULT);
this.refreshThreshold = Math.max(1,
conf.getInt(HDDS_HEARTBEAT_ADDRESS_REFRESH_MISSED_COUNT_THRESHOLD,
@@ -123,7 +136,8 @@ public EndpointStateMachine.EndPointStates call() throws
Exception {
versionManager.getApparentVersion(),
versionManager.getSoftwareVersion());
- datanodeDetails.setCurrentVersion(versionManager.getVersionForClient());
+ datanodeDetails.setCurrentVersion(
+ testCurrentVersion != null ? testCurrentVersion :
versionManager.getVersionForClient());
DatanodeDetailsProto datanodeDetailsProto =
datanodeDetails.getProtoBufMessage();
requestBuilder = SCMHeartbeatRequestProto.newBuilder()
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 bcc0440b7d2..b8bc10758f8 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
@@ -31,6 +31,7 @@
import java.time.Clock;
import java.time.Duration;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
@@ -41,6 +42,7 @@
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import org.apache.commons.lang3.tuple.Pair;
+import org.apache.hadoop.hdds.ComponentVersion;
import org.apache.hadoop.hdds.HddsConfigKeys;
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.conf.Config;
@@ -78,11 +80,13 @@
import org.apache.hadoop.hdds.scm.events.SCMEvents;
import org.apache.hadoop.hdds.scm.ha.SCMContext;
import org.apache.hadoop.hdds.scm.ha.SCMService;
+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;
import org.apache.hadoop.hdds.scm.pipeline.PipelineNotFoundException;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
+import org.apache.hadoop.hdds.scm.server.upgrade.ScmVersionManager;
import org.apache.hadoop.hdds.server.events.EventPublisher;
import org.apache.hadoop.hdds.utils.HddsServerUtil;
import org.apache.hadoop.ozone.container.replication.ReplicationServer;
@@ -527,17 +531,10 @@ public void sendThrottledReplicationCommand(ContainerInfo
containerInfo,
DatanodeDetails source = selectAndOptionallyExcludeDatanode(
1, sourceWithCmds);
- 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);
- }
+ ReplicateContainerCommand cmd = ReplicateContainerCommand.toTarget(
+ containerID, target, computeVersionForReplication(source, target));
+ cmd.setReplicaIndex(replicaIndex);
+ sendDatanodeCommand(cmd, containerInfo, source);
}
public void sendThrottledReconstructionCommand(ContainerInfo containerInfo,
@@ -556,6 +553,20 @@ public void
sendThrottledReconstructionCommand(ContainerInfo containerInfo,
sendDatanodeCommand(command, containerInfo, target);
}
+ private ComponentVersion computeVersionForReplication(DatanodeDetails
source, DatanodeDetails target) {
+ return ScmVersionManager.computeVersionForReplication(
+ Arrays.asList(getDatanodeInfo(source), getDatanodeInfo(target)));
+ }
+
+ 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;
+ }
+
private DatanodeDetails selectAndOptionallyExcludeDatanode(
int additionalCmdCount, List<Pair<Integer, DatanodeDetails>> datanodes) {
if (datanodes.isEmpty()) {
@@ -634,18 +645,11 @@ public void sendLowPriorityReplicateContainerCommand(
final ContainerInfo container, int replicaIndex, DatanodeDetails source,
DatanodeDetails target, long scmDeadlineEpochMs)
throws NotLeaderException {
- 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);
- }
+ final ReplicateContainerCommand command =
ReplicateContainerCommand.toTarget(
+ container.getContainerID(), target,
computeVersionForReplication(source, target));
+ command.setReplicaIndex(replicaIndex);
+ command.setPriority(ReplicationCommandPriority.LOW);
+ sendDatanodeCommand(command, container, source, scmDeadlineEpochMs);
}
/**
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeInfo.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeInfo.java
index 160f788c90c..75203c1384a 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeInfo.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeInfo.java
@@ -25,6 +25,7 @@
import java.util.concurrent.locks.ReadWriteLock;
import java.util.concurrent.locks.ReentrantReadWriteLock;
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.CommandQueueReportProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerDatanodeProtocolProtos.DatanodeVersionProto;
@@ -127,6 +128,18 @@ public void updateLastKnownVersions(DatanodeVersionProto
version) {
}
}
+ /**
+ * Updates the current version reported by this datanode on its heartbeat in
a thread-safe manner.
+ */
+ public void updateCurrentVersion(HDDSVersion version) {
+ try {
+ lock.writeLock().lock();
+ setCurrentVersion(version);
+ } finally {
+ lock.writeLock().unlock();
+ }
+ }
+
/**
* Returns the last heartbeat time.
*
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 c36419de6ec..a389087e7d2 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
@@ -219,28 +219,6 @@ 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/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
index 266fb0f2d94..cd5d01d05a3 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/NodeStateManager.java
@@ -317,6 +317,17 @@ public void updateLastHeartbeatTime(DatanodeDetails
datanodeDetails)
.updateLastHeartbeatTime();
}
+ /**
+ * Updates the current version reported by the node on its heartbeat.
+ *
+ * @throws NodeNotFoundException if the node is not present
+ */
+ public void updateCurrentVersion(DatanodeDetails datanodeDetails)
+ throws NodeNotFoundException {
+ nodeStateMap.getNodeInfo(datanodeDetails.getID())
+ .updateCurrentVersion(datanodeDetails.getCurrentVersion());
+ }
+
/**
* Updates the last known version of the node.
* @param datanodeDetails DataNode Details
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
index eca5933a2aa..b9b44eb57b9 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/SCMNodeManager.java
@@ -562,6 +562,7 @@ public List<SCMCommand<?>> processHeartbeat(DatanodeDetails
datanodeDetails,
"DatanodeDetails.");
try {
nodeStateManager.updateLastHeartbeatTime(datanodeDetails);
+ nodeStateManager.updateCurrentVersion(datanodeDetails);
metrics.incNumHBProcessed();
updateDatanodeOpState(datanodeDetails);
} catch (NodeNotFoundException e) {
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/ScmBlockLocationProtocolServerSideTranslatorPB.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/ScmBlockLocationProtocolServerSideTranslatorPB.java
index b2fb1f897af..cb6be682783 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/ScmBlockLocationProtocolServerSideTranslatorPB.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/ScmBlockLocationProtocolServerSideTranslatorPB.java
@@ -20,12 +20,12 @@
import com.google.protobuf.RpcController;
import com.google.protobuf.ServiceException;
import java.io.IOException;
+import java.util.ArrayList;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.stream.Collectors;
import org.apache.hadoop.hdds.ComponentVersion;
-import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.annotation.InterfaceAudience;
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
@@ -51,12 +51,13 @@
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.ha.RatisUtil;
import org.apache.hadoop.hdds.scm.net.InnerNode;
-import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
import org.apache.hadoop.hdds.scm.protocolPB.ScmBlockLocationProtocolPB;
import
org.apache.hadoop.hdds.scm.protocolPB.StorageContainerLocationProtocolPB;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
+import org.apache.hadoop.hdds.scm.server.upgrade.ScmVersionManager;
import org.apache.hadoop.hdds.server.OzoneProtocolMessageDispatcher;
import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
import org.apache.hadoop.hdds.utils.ProtocolMessageMetrics;
@@ -219,13 +220,22 @@ public AllocateScmBlockResponseProto allocateScmBlock(
" blocks. Requested " + request.getNumBlocks() + " blocks",
SCMException.ResultCodes.FAILED_TO_ALLOCATE_ENOUGH_BLOCKS);
}
+
+ // Allows skipping version computation for pipelines we have already
encountered in this request.
Map<PipelineID, HddsProtos.Pipeline> pipelineProtoCache = new HashMap<>();
for (AllocatedBlock block : allocatedBlocks) {
Pipeline pipeline = block.getPipeline();
+ if (pipeline.getNodes().isEmpty()) {
+ throw new SCMException("Cannot process allocate block request for
empty pipeline",
+ SCMException.ResultCodes.FAILED_TO_FIND_ACTIVE_PIPELINE);
+ }
HddsProtos.Pipeline pipelineProto =
pipelineProtoCache.get(pipeline.getId());
if (pipelineProto == null) {
- pipelineProto = pipeline.getProtobufMessage(clientVersion,
Name.IO_PORTS,
- computePipelineWriteVersion(pipeline));
+ // Pipeline members are frozen copies rebuilt from the replicated
pipeline proto, so their currentVersion can
+ // be stale; resolve the authoritative value from the live node
registry before computing the minimum.
+ ComponentVersion pipelineVersion =
+
ScmVersionManager.computeVersionForClientWrite(currentNodes(pipeline.getNodes()));
+ pipelineProto = pipeline.getProtobufMessage(clientVersion,
Name.IO_PORTS, pipelineVersion);
pipelineProtoCache.put(pipeline.getId(), pipelineProto);
}
builder.addBlocks(AllocateBlockResponse.newBuilder()
@@ -236,35 +246,17 @@ public AllocateScmBlockResponseProto allocateScmBlock(
return builder.build();
}
- /**
- * Computes the version clients should use for writes to the given pipeline.
- * Before ZDU is finalized, datanodes report apparent versions from the
- * {@link HDDSLayoutFeature} enum, which clients cannot compare against the
- * {@link HDDSVersion} enum they use. In that state we advertise the last
- * {@link HDDSVersion} before {@code ZDU} so clients keep pre-ZDU write
behavior
- * until the cluster finalizes. Once ZDU is finalized every apparent version
is
- * an {@link HDDSVersion} and safe to share: we return the lowest one across
the
- * pipeline, so during a later rolling upgrade clients do not enable a newer
- * write-path feature until every datanode in the pipeline has finalized.
- */
- private ComponentVersion computePipelineWriteVersion(Pipeline pipeline)
throws SCMException {
- List<DatanodeDetails> nodes = pipeline.getNodes();
- if (nodes.isEmpty()) {
- throw new SCMException("Cannot determine the write version for pipeline "
- + pipeline.getId() + " because it has no datanodes",
- SCMException.ResultCodes.NO_SUCH_DATANODE);
- }
- if (!scm.getVersionManager().isAllowed(HDDSVersion.ZDU)) {
- return HDDSVersion.STREAM_BLOCK_SUPPORT; // last HDDSVersion before ZDU
- }
- try {
- return scm.getScmNodeManager()
- .getLowestApparentVersion(nodes.toArray(new DatanodeDetails[0]));
- } catch (NodeNotFoundException e) {
- throw new SCMException("Datanode not found while computing the write
version "
- + "for pipeline " + pipeline.getId() + " during block allocation", e,
- SCMException.ResultCodes.NO_SUCH_DATANODE);
+ private List<DatanodeDetails> currentNodes(List<DatanodeDetails> nodes)
throws SCMException {
+ List<DatanodeDetails> memberVersions = new ArrayList<>();
+ for (DatanodeDetails dn : nodes) {
+ DatanodeInfo live = scm.getScmNodeManager().getNode(dn.getID());
+ if (live == null) {
+ throw new SCMException("Failed to lookup datanode " + dn + " in
NodeManager",
+ SCMException.ResultCodes.NO_SUCH_DATANODE);
+ }
+ memberVersions.add(live);
}
+ return memberVersions;
}
public DeleteScmKeyBlocksResponseProto deleteScmKeyBlocks(
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java
index b9012a744d9..d65b3d03e5d 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocolServerSideTranslatorPB.java
@@ -35,14 +35,17 @@
import com.google.protobuf.ServiceException;
import java.io.IOException;
import java.util.ArrayList;
+import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Optional;
import org.apache.commons.lang3.tuple.Pair;
+import org.apache.hadoop.hdds.ComponentVersion;
import org.apache.hadoop.hdds.annotation.InterfaceAudience;
import org.apache.hadoop.hdds.client.ECReplicationConfig;
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;
import
org.apache.hadoop.hdds.protocol.proto.HddsProtos.TransferLeadershipRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.HddsProtos.TransferLeadershipResponseProto;
@@ -147,6 +150,7 @@
import
org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline;
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.ha.RatisUtil;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
import org.apache.hadoop.hdds.scm.protocolPB.OzonePBHelper;
import
org.apache.hadoop.hdds.scm.protocolPB.StorageContainerLocationProtocolPB;
@@ -817,7 +821,7 @@ public ContainerResponseProto
allocateContainer(ContainerRequestProto request,
);
ContainerWithPipeline cp = impl.allocateContainer(replicationConfig,
request.getOwner());
return ContainerResponseProto.newBuilder()
- .setContainerWithPipeline(cp.getProtobuf(clientVersion))
+ .setContainerWithPipeline(cp.getProtobuf(clientVersion,
currentVersions(cp)))
.setErrorCode(ContainerResponseProto.Error.success)
.build();
@@ -848,7 +852,7 @@ public GetContainerWithPipelineResponseProto
getContainerWithPipeline(
ContainerWithPipeline container = impl
.getContainerWithPipeline(request.getContainerID());
return GetContainerWithPipelineResponseProto.newBuilder()
- .setContainerWithPipeline(container.getProtobuf(clientVersion))
+ .setContainerWithPipeline(container.getProtobuf(clientVersion,
currentVersions(container)))
.build();
}
@@ -861,7 +865,7 @@ public GetContainerWithPipelineResponseProto
getContainerWithPipeline(
GetContainerWithPipelineBatchResponseProto.Builder builder =
GetContainerWithPipelineBatchResponseProto.newBuilder();
for (ContainerWithPipeline container : containers) {
- builder.addContainerWithPipelines(container.getProtobuf(clientVersion));
+ builder.addContainerWithPipelines(container.getProtobuf(clientVersion,
currentVersions(container)));
}
return builder.build();
}
@@ -875,11 +879,25 @@ public GetContainerWithPipelineResponseProto
getContainerWithPipeline(
GetExistContainerWithPipelinesInBatchResponseProto.Builder builder =
GetExistContainerWithPipelinesInBatchResponseProto.newBuilder();
for (ContainerWithPipeline container : containers) {
- builder.addContainerWithPipelines(container.getProtobuf(clientVersion));
+ builder.addContainerWithPipelines(container.getProtobuf(clientVersion,
currentVersions(container)));
}
return builder.build();
}
+ private Map<DatanodeID, ComponentVersion>
currentVersions(ContainerWithPipeline containerWithPipeline)
+ throws SCMException {
+ Map<DatanodeID, ComponentVersion> memberVersions = new HashMap<>();
+ for (DatanodeDetails dn : containerWithPipeline.getPipeline().getNodes()) {
+ DatanodeInfo live = scm.getScmNodeManager().getNode(dn.getID());
+ if (live == null) {
+ throw new SCMException("Failed to lookup datanode " + dn + " in
NodeManager",
+ SCMException.ResultCodes.NO_SUCH_DATANODE);
+ }
+ memberVersions.put(live.getID(), live.getCurrentVersion());
+ }
+ return memberVersions;
+ }
+
public SCMListContainerResponseProto listContainer(
SCMListContainerRequestProto request) throws IOException {
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/upgrade/ScmVersionManager.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/upgrade/ScmVersionManager.java
index 3f79bd0ddcb..93152238c85 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/upgrade/ScmVersionManager.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/upgrade/ScmVersionManager.java
@@ -19,9 +19,13 @@
import com.google.common.annotations.VisibleForTesting;
import java.io.IOException;
+import java.util.List;
import java.util.Map;
+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.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.server.OzoneStorageContainerManager;
import org.apache.hadoop.hdds.scm.server.SCMStorageConfig;
import org.apache.hadoop.hdds.upgrade.HDDSVersionUtils;
@@ -56,6 +60,30 @@ public ScmVersionManager(SCMStorageConfig storage,
upgradeActions = upgradeActionProvider.load();
}
+ public static ComponentVersion
computeVersionForReplication(List<DatanodeInfo> datanodes) {
+ return computeCommonVersion(datanodes.stream()
+ .map(DatanodeInfo::getLastKnownApparentVersion)
+ .collect(Collectors.toList()));
+ }
+
+ public static HDDSVersion computeVersionForClientWrite(List<DatanodeDetails>
datanodes) {
+ return computeCommonVersion(datanodes.stream()
+ .map(DatanodeDetails::getCurrentVersion)
+ .collect(Collectors.toList()));
+ }
+
+ private static <T extends ComponentVersion> T computeCommonVersion(List<T>
dnVersions) {
+ if (dnVersions.isEmpty()) {
+ throw new IllegalArgumentException("No nodes provided");
+ }
+
+ T minVersion = dnVersions.get(0);
+ for (int i = 1; i < dnVersions.size(); i++) {
+ minVersion = ComponentVersion.min(dnVersions.get(i), minVersion);
+ }
+ return minVersion;
+ }
+
@Override
protected void persistApparentVersion(ComponentVersion newVersion) throws
IOException {
storage.setApparentVersion(newVersion.serialize());
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 500ed6119a0..5943efdd4ac 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
@@ -142,7 +142,7 @@ public class TestReplicationManager {
private Set<Pair<DatanodeID, SCMCommand<?>>> commandsSent;
@BeforeEach
- public void setup() throws IOException {
+ public void setup() throws IOException, NodeNotFoundException {
configuration = new OzoneConfiguration();
configuration.set(HDDS_SCM_WAIT_TIME_AFTER_SAFE_MODE_EXIT, "0s");
rmConf =
configuration.getObject(ReplicationManager.ReplicationManagerConfiguration.class);
@@ -1765,12 +1765,11 @@ private ReplicateContainerCommand
captureSentReplicateCommand() {
return (ReplicateContainerCommand) command.getValue();
}
- private DatanodeInfo mockDatanodeWithApparentVersion(
+ private void mockDatanodeWithApparentVersion(
DatanodeDetails dn, ComponentVersion version) {
DatanodeInfo info = mock(DatanodeInfo.class);
when(info.getLastKnownApparentVersion()).thenReturn(version);
when(nodeManager.getNode(dn.getID())).thenReturn(info);
- return info;
}
/**
@@ -1799,8 +1798,6 @@ public void testApparentVersionIsLowestOfSourceAndTarget(
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);
@@ -1826,8 +1823,6 @@ public void
testApparentVersionLookupThrowsWhenNodeNotFound()
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, () ->
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 f247cfb0e40..38986d89423 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
@@ -813,6 +813,23 @@ public void
testDatanodeFinalizationCountsTracksApparentVersionRange()
}
}
+ @Test
+ public void testProcessHeartbeatUpdatesCurrentVersion()
+ throws IOException, AuthenticationException {
+ try (SCMNodeManager nodeManager = createNodeManager(getConf())) {
+ DatanodeDetails datanodeDetails =
+ HddsTestUtils.createRandomDatanodeAndRegister(nodeManager);
+
+ // A later heartbeat reports a new currentVersion; SCM should refresh the
+ // value stored in the datanode's DatanodeInfo.
+
datanodeDetails.setCurrentVersion(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC);
+ nodeManager.processHeartbeat(datanodeDetails);
+
+ assertEquals(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC,
+ nodeManager.getNode(datanodeDetails.getID()).getCurrentVersion());
+ }
+ }
+
@Test
public void testDatanodeFinalizationCountsWithNoHealthyDatanodes()
throws IOException, AuthenticationException {
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/protocol/TestScmBlockLocationProtocolServerSideTranslatorPB.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/protocol/TestScmBlockLocationProtocolServerSideTranslatorPB.java
index e0cd8d389b0..bf6ee7ea2e3 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/protocol/TestScmBlockLocationProtocolServerSideTranslatorPB.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/protocol/TestScmBlockLocationProtocolServerSideTranslatorPB.java
@@ -18,13 +18,13 @@
package org.apache.hadoop.hdds.scm.protocol;
import static
org.apache.hadoop.hdds.protocol.MockDatanodeDetails.randomDatanodeDetails;
+import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.ArgumentMatchers.anyString;
-import static org.mockito.Mockito.lenient;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.times;
import static org.mockito.Mockito.verify;
@@ -33,88 +33,103 @@
import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
+import java.util.HashMap;
import java.util.List;
-import org.apache.hadoop.hdds.ComponentVersion;
+import java.util.Map;
import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.client.ContainerBlockID;
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.DatanodeDetailsProto;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationType;
+import org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.AllocateScmBlockRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.AllocateScmBlockResponseProto;
+import org.apache.hadoop.hdds.scm.HddsTestUtils;
import org.apache.hadoop.hdds.scm.container.common.helpers.AllocatedBlock;
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
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.pipeline.Pipeline;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState;
import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
-import org.apache.hadoop.hdds.scm.server.upgrade.ScmVersionManager;
-import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
import org.apache.hadoop.hdds.utils.ProtocolMessageMetrics;
import org.apache.hadoop.ozone.ClientVersion;
+import org.apache.hadoop.ozone.container.upgrade.UpgradeUtils;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
/**
- * Tests that {@link ScmBlockLocationProtocolServerSideTranslatorPB} forwards a
- * clamped write version (based on the cluster's finalization status) in the
- * {@code currentVersion} of every pipeline member it returns on block
- * allocation, without mutating SCM's in-memory state.
+ * Tests that {@link ScmBlockLocationProtocolServerSideTranslatorPB} forwards
the
+ * pipeline-wide minimum {@code currentVersion} to every pipeline member it
+ * returns on block allocation.
+ * <p>
+ * Pipeline members are frozen copies (rebuilt from the replicated pipeline
+ * proto), so the translator must source each member's current version from the
+ * live node registry, which the heartbeat handler keeps up to date. These
tests
+ * pin the pipeline copies at {@link HDDSVersion#SOFTWARE_VERSION} and
register a
+ * different version per node so a translator that read the stale
+ * pipeline copy would compute the wrong minimum.
*/
class TestScmBlockLocationProtocolServerSideTranslatorPB {
private ScmBlockLocationProtocol impl;
- private NodeManager nodeManager;
- private ScmVersionManager versionManager;
private ScmBlockLocationProtocolServerSideTranslatorPB service;
+ private Map<DatanodeID, DatanodeInfo> registry;
private List<DatanodeDetails> nodes;
+ private NodeManager nodeManager;
@BeforeEach
void setUp() throws Exception {
impl = mock(ScmBlockLocationProtocol.class);
StorageContainerManager scm = mock(StorageContainerManager.class);
+
+ registry = new HashMap<>();
nodeManager = mock(NodeManager.class);
- versionManager = mock(ScmVersionManager.class);
+ when(nodeManager.getNode(any(DatanodeID.class))).thenAnswer(inv ->
registry.get(inv.getArgument(0)));
+ when(scm.getScmNodeManager()).thenReturn(nodeManager);
- nodes = new ArrayList<>();
- for (int i = 0; i < 3; i++) {
- DatanodeDetails dn = randomDatanodeDetails();
- dn.setCurrentVersion(HDDSVersion.SOFTWARE_VERSION);
- nodes.add(dn);
- }
+ // Pipeline copies stay at SOFTWARE_VERSION; the live version lives in the
registry.
+ nodes = registerNodes(3, HDDSVersion.SOFTWARE_VERSION);
Pipeline pipeline = buildPipeline(nodes);
- AllocatedBlock block = blockOn(1L, pipeline);
+ AllocatedBlock block = blockInPipeline(1L, pipeline);
when(impl.allocateBlock(anyLong(), anyInt(), any(ReplicationConfig.class),
anyString(), any(),
anyString())).thenReturn(Collections.singletonList(block));
- when(scm.getScmNodeManager()).thenReturn(nodeManager);
- when(scm.getVersionManager()).thenReturn(versionManager);
- // Exercise the real min-across-nodes helper, driven by the getNode stubs.
-
lenient().when(nodeManager.getLowestApparentVersion(any(DatanodeDetails[].class))).thenCallRealMethod();
-
- // Default: ZDU is finalized and every datanode is at the software version.
- when(versionManager.isAllowed(HDDSVersion.ZDU)).thenReturn(true);
- for (DatanodeDetails dn : nodes) {
- setDatanodeApparentVersion(dn, HDDSVersion.SOFTWARE_VERSION);
- }
service = new ScmBlockLocationProtocolServerSideTranslatorPB(impl, scm,
mock(ProtocolMessageMetrics.class));
}
- private void setDatanodeApparentVersion(DatanodeDetails dn, ComponentVersion
version) {
- DatanodeInfo info = mock(DatanodeInfo.class);
- lenient().when(info.getLastKnownApparentVersion()).thenReturn(version);
- when(nodeManager.getNode(dn.getID())).thenReturn(info);
+ /**
+ * Creates {@code count} pipeline members (left at {@link
HDDSVersion#SOFTWARE_VERSION}) and registers each one in the
+ * mocked node manager with the given live {@code currentVersion}.
+ */
+ private List<DatanodeDetails> registerNodes(int count, HDDSVersion
liveVersion) {
+ List<DatanodeDetails> created = new ArrayList<>();
+ for (int i = 0; i < count; i++) {
+ DatanodeDetails dn = randomDatanodeDetails();
+ // Test setup expects nodes to be created with the latest software
version by default.
+ assertEquals(HDDSVersion.SOFTWARE_VERSION, dn.getCurrentVersion());
+ created.add(dn);
+ setLiveVersion(dn, liveVersion);
+ }
+ return created;
}
- private List<DatanodeDetailsProto> allocateAndGetMembers() throws Exception {
- return allocate(1).getBlocks(0).getPipeline().getMembersList();
+ /**
+ * Sets the version the node registry reports for the given pipeline member,
mirroring a heartbeat update on SCM.
+ */
+ private void setLiveVersion(DatanodeDetails member, HDDSVersion version) {
+ DatanodeInfo info = new DatanodeInfo(member, NodeStatus.inServiceHealthy(),
+ UpgradeUtils.defaultVersionProto(),
HddsTestUtils.ROLL_INTERVAL_MS_DEFAULT);
+ info.setCurrentVersion(version);
+ registry.put(member.getID(), info);
}
private AllocateScmBlockResponseProto allocate(int numBlocks) throws
Exception {
@@ -137,7 +152,7 @@ private Pipeline buildPipeline(List<DatanodeDetails>
pipelineNodes) {
.build();
}
- private AllocatedBlock blockOn(long localId, Pipeline pipeline) {
+ private AllocatedBlock blockInPipeline(long localId, Pipeline pipeline) {
return new AllocatedBlock.Builder()
.setContainerBlockID(new ContainerBlockID(1L, localId))
.setPipeline(pipeline)
@@ -149,7 +164,9 @@ private void setAllocatedBlocks(List<AllocatedBlock>
blocks) throws Exception {
anyString(), any(), anyString())).thenReturn(blocks);
}
- private void assertAllMembersHaveVersion(int expected,
List<DatanodeDetailsProto> members) {
+ private void assertAllMembersHaveVersion(int expected,
+ ScmBlockLocationProtocolProtos.AllocateBlockResponse response) {
+ List<DatanodeDetailsProto> members =
response.getPipeline().getMembersList();
assertEquals(nodes.size(), members.size());
for (DatanodeDetailsProto member : members) {
assertEquals(expected, member.getCurrentVersion());
@@ -157,65 +174,45 @@ private void assertAllMembersHaveVersion(int expected,
List<DatanodeDetailsProto
}
@Test
- void preFinalizedClusterClampsClientVersionDown() throws Exception {
- // Before ZDU is finalized, datanodes report apparent versions from the
- // HDDSLayoutFeature enum. Regardless of what they report, clients must be
- // clamped to the last HDDSVersion before ZDU.
- when(versionManager.isAllowed(HDDSVersion.ZDU)).thenReturn(false);
- for (DatanodeDetails dn : nodes) {
- setDatanodeApparentVersion(dn,
HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION);
- }
-
- assertAllMembersHaveVersion(HDDSVersion.STREAM_BLOCK_SUPPORT.serialize(),
allocateAndGetMembers());
+ public void testMinimumVersionAcrossPipeline() throws Exception {
+ // The pipeline mixes datanodes at different reported versions (as during a
+ // rolling upgrade). The lowest is forwarded to every member so clients do
+ // not enable a feature an non-upgraded datanode cannot handle.
+ setLiveVersion(nodes.get(1), HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC);
+
+
assertAllMembersHaveVersion(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC.serialize(),
+ allocate(1).getBlocks(0));
}
@Test
- void finalizedClusterForwardsRealVersion() throws Exception {
- assertAllMembersHaveVersion(HDDSVersion.SOFTWARE_VERSION.serialize(),
allocateAndGetMembers());
+ public void testUniformPipelineForwardsThatVersion() throws Exception {
+ assertAllMembersHaveVersion(HDDSVersion.SOFTWARE_VERSION.serialize(),
allocate(1).getBlocks(0));
}
@Test
- void finalizedPipelineForwardsLowestApparentVersion() throws Exception {
- // ZDU is finalized, but the pipeline mixes datanodes at different apparent
- // versions (as during a later rolling upgrade). The lowest is forwarded so
- // clients do not enable a feature an un-upgraded datanode cannot handle.
- setDatanodeApparentVersion(nodes.get(1),
HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC);
-
-
assertAllMembersHaveVersion(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC.serialize(),
allocateAndGetMembers());
- }
-
- @Test
- void missingDatanodeInfoFailsAllocation() throws Exception {
- when(nodeManager.getNode(nodes.get(2).getID())).thenReturn(null);
+ public void testEmptyPipelineFailsAllocation() throws Exception {
+ setAllocatedBlocks(Collections.singletonList(blockInPipeline(1L,
buildPipeline(Collections.emptyList()))));
SCMException e = assertThrows(SCMException.class, () -> allocate(1));
- assertEquals(SCMException.ResultCodes.NO_SUCH_DATANODE, e.getResult());
+ assertEquals(SCMException.ResultCodes.FAILED_TO_FIND_ACTIVE_PIPELINE,
e.getResult());
}
@Test
- void emptyPipelineFailsAllocation() throws Exception {
- setAllocatedBlocks(Collections.singletonList(blockOn(1L,
buildPipeline(Collections.emptyList()))));
-
- SCMException e = assertThrows(SCMException.class, () -> allocate(1));
- assertEquals(SCMException.ResultCodes.NO_SUCH_DATANODE, e.getResult());
- }
-
- @Test
- void blocksSharingPipelineAllGetClampedVersion() throws Exception {
- // One straggler datanode clamps the shared pipeline's write version.
- setDatanodeApparentVersion(nodes.get(1),
HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC);
+ public void testBlocksSharingPipelineAllGetMinVersion() throws Exception {
+ // One old datanode should determine the shared pipeline's minimum version.
+ setLiveVersion(nodes.get(1), HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC);
// Three blocks allocated on the same pipeline object; the memoized proto
- // must be returned for each block with the clamped version intact.
+ // must be returned for each block with the minimum version intact.
Pipeline pipeline = buildPipeline(nodes);
- setAllocatedBlocks(Arrays.asList(blockOn(1L, pipeline), blockOn(2L,
pipeline), blockOn(3L, pipeline)));
+ setAllocatedBlocks(Arrays.asList(
+ blockInPipeline(1L, pipeline), blockInPipeline(2L, pipeline),
blockInPipeline(3L, pipeline)));
AllocateScmBlockResponseProto response = allocate(3);
assertEquals(3, response.getBlocksCount());
for (int i = 0; i < response.getBlocksCount(); i++) {
-
assertAllMembersHaveVersion(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC.serialize(),
- response.getBlocks(i).getPipeline().getMembersList());
+
assertAllMembersHaveVersion(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC.serialize(),
response.getBlocks(i));
}
// The write version is memoized per pipeline: each node is looked up once
@@ -226,42 +223,71 @@ void blocksSharingPipelineAllGetClampedVersion() throws
Exception {
}
@Test
- void blocksOnDistinctPipelinesGetOwnVersion() throws Exception {
- // A second, distinct pipeline whose datanodes are at an older
- // version than the software-version pipeline built in setUp().
- List<DatanodeDetails> otherNodes = new ArrayList<>();
- for (int i = 0; i < 3; i++) {
- DatanodeDetails dn = randomDatanodeDetails();
- dn.setCurrentVersion(HDDSVersion.SOFTWARE_VERSION);
- setDatanodeApparentVersion(dn, HDDSVersion.STREAM_BLOCK_SUPPORT);
- otherNodes.add(dn);
- }
+ public void testBlocksOnDistinctPipelinesGetOwnMinVersion() throws Exception
{
+ // A second, distinct pipeline whose datanodes are at an older version than
+ // the software-version pipeline built in setUp().
+ List<DatanodeDetails> otherNodes = registerNodes(3,
HDDSVersion.STREAM_BLOCK_SUPPORT);
- Pipeline finalized = buildPipeline(nodes);
- Pipeline straggler = buildPipeline(otherNodes);
- setAllocatedBlocks(Arrays.asList(blockOn(1L, finalized), blockOn(2L,
straggler)));
+ Pipeline uniform = buildPipeline(nodes);
+ Pipeline oneOld = buildPipeline(otherNodes);
+ setAllocatedBlocks(Arrays.asList(blockInPipeline(1L, uniform),
blockInPipeline(2L, oneOld)));
AllocateScmBlockResponseProto response = allocate(2);
assertEquals(2, response.getBlocksCount());
- assertAllMembersHaveVersion(HDDSVersion.SOFTWARE_VERSION.serialize(),
- response.getBlocks(0).getPipeline().getMembersList());
- assertEquals(otherNodes.size(),
response.getBlocks(1).getPipeline().getMembersCount());
- for (DatanodeDetailsProto member :
response.getBlocks(1).getPipeline().getMembersList()) {
- assertEquals(HDDSVersion.STREAM_BLOCK_SUPPORT.serialize(),
member.getCurrentVersion());
- }
+ assertAllMembersHaveVersion(HDDSVersion.SOFTWARE_VERSION.serialize(),
response.getBlocks(0));
+ assertAllMembersHaveVersion(HDDSVersion.STREAM_BLOCK_SUPPORT.serialize(),
response.getBlocks(1));
}
+ /**
+ * Tests that the original DatanodeDetails object is not modified when the
current version is assigned and returned
+ * to the client.
+ */
@Test
- void doesNotMutateSourcePipelineDatanodes() throws Exception {
- setDatanodeApparentVersion(nodes.get(1), HDDSVersion.STREAM_BLOCK_SUPPORT);
+ public void testDatanodeDetailsVersionOverride() throws Exception {
+ // `nodes` represents the persisted versions in the pipeline manager,
which are configured to be returned on
+ // block allocation.
+ DatanodeDetails oldNode = nodes.get(0);
+ // `registry` contains the latest versions of all datanodes, representing
their heartbeat information.
+ setLiveVersion(oldNode, HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC);
+
+ // The pipeline nodes should all have their default version.
+ for (DatanodeDetails member : nodes) {
+ assertEquals(HDDSVersion.SOFTWARE_VERSION, member.getCurrentVersion());
+ }
+ // The node manager registry should have one node in an older version
after the test setup.
+ for (DatanodeDetails member : registry.values()) {
+ if (member.equals(oldNode)) {
+ assertEquals(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC,
member.getCurrentVersion());
+ } else {
+ assertEquals(HDDSVersion.SOFTWARE_VERSION, member.getCurrentVersion());
+ }
+ }
- allocateAndGetMembers();
+ // Since there is one old node, the whole pipeline should report that as
the version to use.
+
assertAllMembersHaveVersion(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC.serialize(),
+ allocate(1).getBlocks(0));
- // The in-memory DatanodeDetails (shared with SCM internal state) must keep
- // their real software version; only the outgoing proto is overridden.
- for (DatanodeDetails dn : nodes) {
- assertEquals(HDDSVersion.SOFTWARE_VERSION, dn.getCurrentVersion());
+ // Serialization overrides only the outgoing pipeline proto; the source
pipeline copies and the
+ // registered node info keep their own versions.
+ for (DatanodeDetails member : nodes) {
+ assertEquals(HDDSVersion.SOFTWARE_VERSION, member.getCurrentVersion());
+ }
+ for (DatanodeDetails member : registry.values()) {
+ if (member.equals(oldNode)) {
+ assertEquals(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC,
member.getCurrentVersion());
+ } else {
+ assertEquals(HDDSVersion.SOFTWARE_VERSION, member.getCurrentVersion());
+ }
}
+
+ }
+
+ @Test
+ public void testUnknownNodeThrows() {
+ DatanodeDetails unknownNode = nodes.get(1);
+ registry.remove(unknownNode.getID());
+ SCMException ex = assertThrows(SCMException.class, () -> allocate(1));
+ assertThat(ex.getMessage()).contains(unknownNode.toString());
}
}
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/protocol/TestStorageContainerLocationProtocolServerSideTranslatorPB.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/protocol/TestStorageContainerLocationProtocolServerSideTranslatorPB.java
new file mode 100644
index 00000000000..954cf1bf1ef
--- /dev/null
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/protocol/TestStorageContainerLocationProtocolServerSideTranslatorPB.java
@@ -0,0 +1,224 @@
+/*
+ * 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.hdds.scm.protocol;
+
+import static
org.apache.hadoop.hdds.protocol.MockDatanodeDetails.randomDatanodeDetails;
+import static org.assertj.core.api.Assertions.assertThat;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import org.apache.hadoop.hdds.HDDSVersion;
+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.DatanodeDetailsProto;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos.LifeCycleState;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationType;
+import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ContainerRequestProto;
+import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerWithPipelineBatchRequestProto;
+import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerWithPipelineRequestProto;
+import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetExistContainerWithPipelinesInBatchRequestProto;
+import org.apache.hadoop.hdds.scm.HddsTestUtils;
+import
org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline;
+import org.apache.hadoop.hdds.scm.exceptions.SCMException;
+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.pipeline.Pipeline;
+import org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState;
+import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
+import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
+import org.apache.hadoop.hdds.utils.ProtocolMessageMetrics;
+import org.apache.hadoop.ozone.ClientVersion;
+import org.apache.hadoop.ozone.container.upgrade.UpgradeUtils;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Tests that {@link StorageContainerLocationProtocolServerSideTranslatorPB}
forwards each datanode's own current
+ * {@code currentVersion} to read clients on the container-with-pipeline
responses.
+ * <p>
+ * Pipeline members are frozen copies (rebuilt from the replicated pipeline
proto), so their currentVersion can be
+ * stale. The translator must source each member's version from the live node
registry, which the heartbeat handler
+ * keeps up to date. Unlike the write path (which forwards the pipeline-wide
minimum), reads forward each member's
+ * own version, so these tests register a <em>different</em> version per node
and assert every member keeps its own.
+ */
+class TestStorageContainerLocationProtocolServerSideTranslatorPB {
+
+ private StorageContainerLocationProtocol impl;
+ private StorageContainerLocationProtocolServerSideTranslatorPB service;
+ private Map<DatanodeID, DatanodeInfo> registry;
+ private List<DatanodeDetails> nodes;
+
+ @BeforeEach
+ void setUp() throws Exception {
+ impl = mock(StorageContainerLocationProtocol.class);
+ StorageContainerManager scm = mock(StorageContainerManager.class);
+
+ registry = new HashMap<>();
+ NodeManager nodeManager = mock(NodeManager.class);
+ when(nodeManager.getNode(any(DatanodeID.class))).thenAnswer(inv ->
registry.get(inv.getArgument(0)));
+ when(scm.getScmNodeManager()).thenReturn(nodeManager);
+
+ // Pipeline copies stay at SOFTWARE_VERSION; the registry holds a distinct
live version per node.
+ nodes = new ArrayList<>();
+ nodes.add(registerNode(HDDSVersion.SEPARATE_RATIS_PORTS_AVAILABLE));
+ nodes.add(registerNode(HDDSVersion.COMBINED_PUTBLOCK_WRITECHUNK_RPC));
+ nodes.add(registerNode(HDDSVersion.STREAM_BLOCK_SUPPORT));
+
+ service = new StorageContainerLocationProtocolServerSideTranslatorPB(impl,
scm,
+ mock(ProtocolMessageMetrics.class));
+ }
+
+ /**
+ * Creates a pipeline member (left at {@link HDDSVersion#SOFTWARE_VERSION})
and registers it in the mocked node
+ * manager registry with the given live {@code currentVersion}, mirroring
what a heartbeat records on SCM.
+ */
+ private DatanodeDetails registerNode(HDDSVersion liveVersion) {
+ DatanodeDetails dn = randomDatanodeDetails();
+ // Test setup expects nodes to be created with the latest software version
by default.
+ assertEquals(HDDSVersion.SOFTWARE_VERSION, dn.getCurrentVersion());
+
+ // Add the node to the registry with a live version representing what it
is reporting on heartbeat.
+ DatanodeInfo info = new DatanodeInfo(dn, NodeStatus.inServiceHealthy(),
+ UpgradeUtils.defaultVersionProto(),
HddsTestUtils.ROLL_INTERVAL_MS_DEFAULT);
+ info.setCurrentVersion(liveVersion);
+ registry.put(dn.getID(), info);
+ return dn;
+ }
+
+ private ContainerWithPipeline containerWithPipeline() {
+ Pipeline pipeline = Pipeline.newBuilder()
+ .setId(PipelineID.randomId())
+ .setState(PipelineState.OPEN)
+
.setReplicationConfig(RatisReplicationConfig.getInstance(ReplicationFactor.THREE))
+ .setNodes(nodes)
+ .build();
+ return new
ContainerWithPipeline(HddsTestUtils.getContainer(LifeCycleState.CLOSED),
pipeline);
+ }
+
+ /** Asserts each serialized member carries its own live (registry) version.
*/
+ private void
assertCurrentVersionsOverriddenAndMatch(List<DatanodeDetailsProto> members) {
+ assertEquals(nodes.size(), members.size());
+ for (DatanodeDetailsProto member : members) {
+ DatanodeID id = DatanodeDetails.getFromProtoBuf(member).getID();
+ assertEquals(registry.get(id).getCurrentVersion().serialize(),
member.getCurrentVersion());
+ // All datanodes in the `nodes` list are in SOFTWARE_VERSION, but their
live overrides have different versions
+ // meaning that no pipeline should be using SOFTWARE_VERSION as its
overall version.
+ assertNotEquals(HDDSVersion.SOFTWARE_VERSION.serialize(),
member.getCurrentVersion());
+ }
+ }
+
+ @Test
+ public void testGetContainerWithPipelineCurrentVersion() throws Exception {
+
when(impl.getContainerWithPipeline(anyLong())).thenReturn(containerWithPipeline());
+
+ List<DatanodeDetailsProto> members = service.getContainerWithPipeline(
+
GetContainerWithPipelineRequestProto.newBuilder().setContainerID(1L).build(),
ClientVersion.CURRENT)
+ .getContainerWithPipeline().getPipeline().getMembersList();
+
+ assertCurrentVersionsOverriddenAndMatch(members);
+ }
+
+ @Test
+ public void testGetContainerWithPipelineBatchCurrentVersion() throws
Exception {
+
when(impl.getContainerWithPipelineBatch(any())).thenReturn(Collections.singletonList(containerWithPipeline()));
+
+ List<DatanodeDetailsProto> members = service.getContainerWithPipelineBatch(
+
GetContainerWithPipelineBatchRequestProto.newBuilder().addContainerIDs(1L).build(),
ClientVersion.CURRENT)
+ .getContainerWithPipelines(0).getPipeline().getMembersList();
+
+ assertCurrentVersionsOverriddenAndMatch(members);
+ }
+
+ @Test
+ public void testGetExistContainerWithPipelinesInBatchCurrentVersion() throws
Exception {
+ when(impl.getExistContainerWithPipelinesInBatch(any()))
+ .thenReturn(Collections.singletonList(containerWithPipeline()));
+
+ List<DatanodeDetailsProto> members =
service.getExistContainerWithPipelinesInBatch(
+
GetExistContainerWithPipelinesInBatchRequestProto.newBuilder().addContainerIDs(1L).build(),
+ ClientVersion.CURRENT)
+ .getContainerWithPipelines(0).getPipeline().getMembersList();
+
+ assertCurrentVersionsOverriddenAndMatch(members);
+ }
+
+ @Test
+ public void testAllocateContainerCurrentVersion() throws Exception {
+ when(impl.allocateContainer(any(ReplicationConfig.class),
anyString())).thenReturn(containerWithPipeline());
+
+ List<DatanodeDetailsProto> members = service.allocateContainer(
+ ContainerRequestProto.newBuilder()
+ .setReplicationType(ReplicationType.RATIS)
+ .setReplicationFactor(ReplicationFactor.THREE)
+ .setOwner("owner")
+ .build(), ClientVersion.CURRENT)
+ .getContainerWithPipeline().getPipeline().getMembersList();
+
+ assertCurrentVersionsOverriddenAndMatch(members);
+ }
+
+ /**
+ * Tests that the original DatanodeDetails object is not modified when the
current version is assigned and returned
+ * to the client.
+ */
+ @Test
+ public void testDatanodeDetailsVersionOverride() throws Exception {
+ for (DatanodeDetails member : nodes) {
+ assertEquals(HDDSVersion.SOFTWARE_VERSION, member.getCurrentVersion());
+ }
+
when(impl.getContainerWithPipeline(anyLong())).thenReturn(containerWithPipeline());
+
+ GetContainerWithPipelineRequestProto requestProto =
GetContainerWithPipelineRequestProto
+ .newBuilder()
+ .setContainerID(1L)
+ .build();
+
assertCurrentVersionsOverriddenAndMatch(service.getContainerWithPipeline(requestProto,
ClientVersion.CURRENT)
+ .getContainerWithPipeline().getPipeline().getMembersList());
+
+ // Serialization overrides only the outgoing proto; the source pipeline
copies keep their own version.
+ for (DatanodeDetails member : nodes) {
+ assertEquals(HDDSVersion.SOFTWARE_VERSION, member.getCurrentVersion());
+ }
+ }
+
+ @Test
+ public void testUnknownNodeThrows() throws Exception {
+ DatanodeDetails unknownNode = nodes.get(1);
+ registry.remove(unknownNode.getID());
+
when(impl.getContainerWithPipeline(anyLong())).thenReturn(containerWithPipeline());
+
+ SCMException ex = assertThrows(SCMException.class, () ->
service.getContainerWithPipeline(
+
GetContainerWithPipelineRequestProto.newBuilder().setContainerID(1L).build(),
ClientVersion.CURRENT));
+ assertThat(ex.getMessage()).contains(unknownNode.toString());
+ }
+}
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestDatanodeCurrentVersionEndToEnd.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestDatanodeCurrentVersionEndToEnd.java
new file mode 100644
index 00000000000..b7ba832c9c5
--- /dev/null
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/client/rpc/TestDatanodeCurrentVersionEndToEnd.java
@@ -0,0 +1,244 @@
+/*
+ * 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.client.rpc;
+
+import static java.nio.charset.StandardCharsets.UTF_8;
+import static java.util.concurrent.TimeUnit.SECONDS;
+import static org.apache.hadoop.hdds.HddsConfigKeys.HDDS_HEARTBEAT_INTERVAL;
+import static
org.apache.hadoop.ozone.HddsDatanodeService.TESTING_DATANODE_VERSION_CURRENT;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+
+import java.io.IOException;
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.stream.Collectors;
+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.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
+import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
+import org.apache.hadoop.hdds.scm.node.NodeManager;
+import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
+import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
+import org.apache.hadoop.ozone.HddsDatanodeService;
+import org.apache.hadoop.ozone.MiniOzoneCluster;
+import org.apache.hadoop.ozone.client.ObjectStore;
+import org.apache.hadoop.ozone.client.OzoneBucket;
+import org.apache.hadoop.ozone.client.OzoneClient;
+import org.apache.hadoop.ozone.client.io.KeyOutputStream;
+import org.apache.hadoop.ozone.client.io.OzoneOutputStream;
+import org.apache.hadoop.ozone.container.OzoneTestHelper;
+import org.apache.hadoop.ozone.om.helpers.OmKeyArgs;
+import org.apache.hadoop.ozone.om.helpers.OmKeyInfo;
+import org.apache.ozone.test.GenericTestUtils;
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.MethodSource;
+
+/**
+ * End-to-end test that a datanode's reported {@code currentVersion} flows all
+ * the way from the datanode to the client, on both the write and read paths.
+ * <p>
+ * A single {@link MiniOzoneCluster} is shared across all parameter values and
is
+ * never restarted. To exercise a given version, the test reloads the
+ * {@link HddsDatanodeService#TESTING_DATANODE_VERSION_CURRENT} config on every
+ * live datanode. The heartbeat task re-reads that config on every heartbeat,
so
+ * the datanodes start advertising the new version to SCM without a restart;
once
+ * they all report the same version SCM converges on it and stays there (every
+ * heartbeat now carries it), so the client observations below are stable.
+ * <p>
+ * The version is then checked on three client paths:
+ * <ul>
+ * <li><b>write path</b> — the pipeline SCM hands back when a block is
+ * allocated for a new key. SCM forwards the pipeline-wide <em>minimum</em>
+ * currentVersion; with every datanode at the same version the minimum equals
+ * that version.</li>
+ * <li><b>closed-container read</b> — looking up a key whose container is
+ * closed. SCM builds the read pipeline from the container replicas, which
+ * carry each datanode's own currentVersion.</li>
+ * <li><b>open-container read</b> — looking up a key whose container is still
+ * open. SCM returns the open pipeline but refreshes each member's
+ * currentVersion from the live node registry, so it too reflects the current
+ * per-datanode version.</li>
+ * </ul>
+ * <p>
+ * All three paths are exercised for both RATIS (factor three) and EC (rs-3-2)
+ * replication, since the two use different pipeline-provider code paths.
+ */
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+class TestDatanodeCurrentVersionEndToEnd {
+
+ private static final int PROPAGATION_TIMEOUT_MS = 30_000;
+ private static final ReplicationConfig RATIS_THREE =
+ RatisReplicationConfig.getInstance(ReplicationFactor.THREE);
+ private static final ReplicationConfig EC_3_2 = new ECReplicationConfig(3,
2);
+
+ private MiniOzoneCluster cluster;
+ private OzoneClient client;
+ private StorageContainerManager scm;
+ private OzoneBucket bucket;
+ private static final String VOLUME_NAME = "vol1";
+ private static final String BUCKET_NAME = "bucket1";
+ private static final String RATIS_CLOSED_KEY = "closed-key";
+ private static final String RATIS_OPEN_KEY = "open-key";
+ private static final String EC_CLOSED_KEY = "ec-closed-key";
+ private static final String EC_OPEN_KEY = "ec-open-key";
+
+ @BeforeAll
+ void init() throws Exception {
+ OzoneConfiguration conf = new OzoneConfiguration();
+ // Fast heartbeats so a reloaded version reaches SCM quickly and the closed
+ // container converges quickly during setup.
+ conf.setTimeDuration(HDDS_HEARTBEAT_INTERVAL, 1, SECONDS);
+ // Five datanodes so an EC rs-3-2 pipeline (3 data + 2 parity) has enough
+ // members; RATIS/THREE simply uses three of them.
+ cluster = MiniOzoneCluster.newBuilder(conf).setNumDatanodes(5).build();
+ cluster.waitForClusterToBeReady();
+ scm = cluster.getStorageContainerManager();
+ client = cluster.newClient();
+ ObjectStore store = client.getObjectStore();
+ store.createVolume(VOLUME_NAME);
+ store.getVolume(VOLUME_NAME).createBucket(BUCKET_NAME);
+ bucket = store.getVolume(VOLUME_NAME).getBucket(BUCKET_NAME);
+
+ // Pre-create keys whose containers are closed so that read-path lookups
+ // resolve the per-datanode read pipeline (built over the closed replicas),
+ // which reflects each datanode's live currentVersion. Also pre-create keys
+ // whose containers stay OPEN; open-container reads resolve the live open
+ // pipeline, which must reflect current datanode versions too. Cover both
+ // RATIS and EC replication.
+ byte[] data = "current-version".getBytes(UTF_8);
+ createClosedContainerKey(RATIS_CLOSED_KEY, RATIS_THREE, data);
+ createClosedContainerKey(EC_CLOSED_KEY, EC_3_2, data);
+ createOpenContainerKey(RATIS_OPEN_KEY, RATIS_THREE, data);
+ createOpenContainerKey(EC_OPEN_KEY, EC_3_2, data);
+ }
+
+ /** Create a key and wait for all of its containers to close. */
+ private void createClosedContainerKey(String keyName, ReplicationConfig
repl, byte[] data) throws Exception {
+ List<Long> containerIds;
+ try (OzoneOutputStream out = bucket.createKey(keyName, data.length, repl,
new HashMap<>())) {
+ out.write(data);
+ containerIds = ((KeyOutputStream)
out.getOutputStream()).getStreamEntries().stream()
+ .map(entry -> entry.getBlockID().getContainerID())
+ .distinct()
+ .collect(Collectors.toList());
+ }
+ OzoneTestHelper.waitForContainerClose(cluster, containerIds.toArray(new
Long[0]));
+ }
+
+ /** Create a key whose container stays open. */
+ private void createOpenContainerKey(String keyName, ReplicationConfig repl,
byte[] data) throws IOException {
+ try (OzoneOutputStream out = bucket.createKey(keyName, data.length, repl,
new HashMap<>())) {
+ out.write(data);
+ }
+ }
+
+ @AfterAll
+ void shutdown() throws IOException {
+ if (client != null) {
+ client.close();
+ }
+ if (cluster != null) {
+ cluster.shutdown();
+ }
+ }
+
+ static List<HDDSVersion> currentVersions() {
+ // UNKNOWN_VERSION is a placeholder, it is not meant to be returned from
Datanodes to clients.
+ return Arrays.stream(HDDSVersion.values())
+ .filter(v -> !v.equals(HDDSVersion.UNKNOWN_VERSION))
+ .collect(Collectors.toList());
+ }
+
+ @ParameterizedTest
+ @MethodSource("currentVersions")
+ public void testCurrentVersionReachesClient(HDDSVersion version) throws
Exception {
+ reloadDatanodeCurrentVersion(version);
+
+ // Wait for the reloaded version to propagate to SCM via a heartbeat. Once
it
+ // has, it is stable: every datanode is now advertising this version.
+ GenericTestUtils.waitFor(() -> scmReportsForAllDatanodes(version),
+ 100, PROPAGATION_TIMEOUT_MS);
+
+ assertNodesAt(version, writePathPipeline("write-probe", RATIS_THREE),
"RATIS write");
+ assertNodesAt(version, writePathPipeline("ec-write-probe", EC_3_2), "EC
write");
+ assertNodesAt(version, readPathPipeline(RATIS_CLOSED_KEY), "RATIS
closed-container read");
+ assertNodesAt(version, readPathPipeline(EC_CLOSED_KEY), "EC
closed-container read");
+ assertNodesAt(version, readPathPipeline(RATIS_OPEN_KEY), "RATIS
open-container read");
+ assertNodesAt(version, readPathPipeline(EC_OPEN_KEY), "EC open-container
read");
+ }
+
+ /**
+ * Reload the "current version" test config on every live datanode and nudge
a
+ * heartbeat. The heartbeat task re-reads this config each heartbeat, so the
+ * datanodes begin advertising {@code version} without being restarted.
+ */
+ private void reloadDatanodeCurrentVersion(HDDSVersion version) {
+ for (HddsDatanodeService dn : cluster.getHddsDatanodes()) {
+ dn.getConf().setInt(TESTING_DATANODE_VERSION_CURRENT,
version.serialize());
+ // triggerHeartbeat re-reads the config at each invocation.
+ dn.getDatanodeStateMachine().triggerHeartbeat();
+ }
+ }
+
+ /** True once SCM's stored version for every datanode matches {@code
version}. */
+ private boolean scmReportsForAllDatanodes(HDDSVersion version) {
+ NodeManager nodeManager = scm.getScmNodeManager();
+ for (HddsDatanodeService dn : cluster.getHddsDatanodes()) {
+ DatanodeInfo info = nodeManager.getNode(dn.getDatanodeDetails().getID());
+ if (info == null || info.getCurrentVersion() != version) {
+ return false;
+ }
+ }
+ return true;
+ }
+
+ /** Client write: the pipeline SCM hands back for a freshly allocated block.
*/
+ private Pipeline writePathPipeline(String probeKey, ReplicationConfig repl)
throws IOException {
+ try (OzoneOutputStream out = bucket.createKey(probeKey, 1, repl, new
HashMap<>())) {
+ out.write(new byte[] {1});
+ return ((KeyOutputStream) out.getOutputStream())
+ .getStreamEntries().get(0).getPipeline();
+ }
+ }
+
+ /** Client read: the pipeline OM resolves from SCM for the given key. */
+ private Pipeline readPathPipeline(String keyName) throws IOException {
+ OmKeyArgs args = new OmKeyArgs.Builder()
+ .setVolumeName(VOLUME_NAME)
+ .setBucketName(BUCKET_NAME)
+ .setKeyName(keyName)
+ .build();
+ OmKeyInfo keyInfo = cluster.getOzoneManager().lookupKey(args);
+ return
keyInfo.getLatestVersionLocations().getLocationList().get(0).getPipeline();
+ }
+
+ private static void assertNodesAt(HDDSVersion version, Pipeline pipeline,
String path) {
+ for (DatanodeDetails node : pipeline.getNodes()) {
+ assertEquals(version, node.getCurrentVersion(),
+ path + " path: datanode " + node.getID() + " reported the wrong
currentVersion");
+ }
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]