Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -86,8 +86,8 @@ public class DatanodeDetails extends NodeImpl implements Comparable<DatanodeDeta
private String revision;
private volatile HddsProtos.NodeOperationalState persistedOpState;
private volatile long persistedOpStateExpiryEpochSec;
private int initialVersion;
private int currentVersion;
private HDDSVersion initialVersion;
private volatile HDDSVersion currentVersion;

private DatanodeDetails(Builder b) {
super(b.hostName, b.networkLocation, NetConstants.NODE_COST_DEFAULT);
Expand Down Expand Up @@ -462,10 +462,7 @@ public static DatanodeDetails.Builder newBuilder(
datanodeDetailsProto.getPersistedOpStateExpiry());
}
if (datanodeDetailsProto.hasCurrentVersion()) {
builder.setCurrentVersion(datanodeDetailsProto.getCurrentVersion());
} else {
// fallback to version 1 if not present
builder.setCurrentVersion(HDDSVersion.SEPARATE_RATIS_PORTS_AVAILABLE.serialize());
builder.setCurrentVersion(HDDSVersion.deserialize(datanodeDetailsProto.getCurrentVersion()));
}
return builder;
}
Expand Down Expand Up @@ -513,14 +510,14 @@ public static DatanodeDetails getFromProtoBuf(
*/
@JsonIgnore
public HddsProtos.DatanodeDetailsProto getProtoBufMessage() {
return toProto(ClientVersion.CURRENT.serialize());
return toProto(ClientVersion.CURRENT);
}

public HddsProtos.DatanodeDetailsProto toProto(int clientVersion) {
public HddsProtos.DatanodeDetailsProto toProto(ClientVersion clientVersion) {
return toProtoBuilder(clientVersion, Collections.emptySet()).build();
}

public HddsProtos.DatanodeDetailsProto toProto(int clientVersion, Set<Port.Name> filterPorts) {
public HddsProtos.DatanodeDetailsProto toProto(ClientVersion clientVersion, Set<Port.Name> filterPorts) {
return toProtoBuilder(clientVersion, filterPorts).build();
}

Expand All @@ -533,7 +530,7 @@ public HddsProtos.DatanodeDetailsProto toProto(int clientVersion, Set<Port.Name>
* @return A {@link HddsProtos.DatanodeDetailsProto.Builder} Object.
*/
public HddsProtos.DatanodeDetailsProto.Builder toProtoBuilder(
int clientVersion, Set<Port.Name> filterPorts) {
ClientVersion clientVersion, Set<Port.Name> filterPorts) {

final HddsProtos.DatanodeIDProto idProto = id.toProto();
final HddsProtos.DatanodeDetailsProto.Builder builder =
Expand Down Expand Up @@ -590,8 +587,7 @@ public HddsProtos.DatanodeDetailsProto.Builder toProtoBuilder(
}
}

builder.setCurrentVersion(currentVersion);

builder.setCurrentVersion(currentVersion.serialize());
return builder;
}

Expand Down Expand Up @@ -622,22 +618,22 @@ public ExtendedDatanodeDetailsProto getExtendedProtoBufMessage() {
* Note: Datanode initial version is not passed to the client due to no use case. See HDDS-9884
* @return the version this datanode was initially created with
*/
public int getInitialVersion() {
public HDDSVersion getInitialVersion() {
return initialVersion;
}

public void setInitialVersion(int initialVersion) {
public void setInitialVersion(HDDSVersion initialVersion) {
this.initialVersion = initialVersion;
}

/**
* @return the version this datanode was last started with
*/
public int getCurrentVersion() {
public HDDSVersion getCurrentVersion() {
return currentVersion;
}

public void setCurrentVersion(int currentVersion) {
public void setCurrentVersion(HDDSVersion currentVersion) {
this.currentVersion = currentVersion;
}

Expand Down Expand Up @@ -721,8 +717,8 @@ public static final class Builder {
private String revision;
private HddsProtos.NodeOperationalState persistedOpState;
private long persistedOpStateExpiryEpochSec = 0;
private int initialVersion;
private int currentVersion = HDDSVersion.SOFTWARE_VERSION.serialize();
private HDDSVersion initialVersion = HDDSVersion.DEFAULT_VERSION;
private HDDSVersion currentVersion = HDDSVersion.DEFAULT_VERSION;

/**
* Default private constructor. To create Builder instance use
Expand Down Expand Up @@ -938,12 +934,12 @@ public Builder setPersistedOpStateExpiry(long expiry) {
return this;
}

public Builder setInitialVersion(int v) {
public Builder setInitialVersion(HDDSVersion v) {
this.initialVersion = v;
return this;
}

public Builder setCurrentVersion(int v) {
public Builder setCurrentVersion(HDDSVersion v) {
this.currentVersion = v;
return this;
}
Expand Down Expand Up @@ -1172,7 +1168,7 @@ public void setRevision(String rev) {

@Override
public HddsProtos.NetworkNode toProtobuf(
int clientVersion) {
ClientVersion clientVersion) {
return HddsProtos.NetworkNode.newBuilder()
.setDatanodeDetails(toProtoBuilder(clientVersion, Collections.emptySet()).build())
.build();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,12 +18,16 @@
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;
import org.apache.hadoop.ozone.ClientVersion;

/**
* Class wraps ozone container info.
Expand Down Expand Up @@ -53,10 +57,15 @@ public static ContainerWithPipeline fromProtobuf(HddsProtos.ContainerWithPipelin
Pipeline.getFromProtobuf(allocatedContainer.getPipeline()));
}

public HddsProtos.ContainerWithPipeline getProtobuf(int 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();
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@
import java.util.Collection;
import java.util.List;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.ozone.ClientVersion;

/**
* The interface defines an inner node in a network topology.
Expand Down Expand Up @@ -92,7 +93,7 @@ Node getLeaf(int leafIndex, List<String> excludedScopes,
Collection<Node> excludedNodes, int ancestorGen);

@Override
HddsProtos.NetworkNode toProtobuf(int clientVersion);
HddsProtos.NetworkNode toProtobuf(ClientVersion clientVersion);

@Override
boolean equals(Object o);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,6 +18,7 @@
package org.apache.hadoop.hdds.scm.net;

import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.ozone.ClientVersion;

/**
* The interface defines a node in a network topology.
Expand Down Expand Up @@ -130,7 +131,7 @@ public interface Node {
boolean isDescendant(String nodePath);

default HddsProtos.NetworkNode toProtobuf(
int clientVersion) {
ClientVersion clientVersion) {
return null;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
import org.apache.commons.lang3.StringUtils;
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.client.ECReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicatedReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicationConfig;
Expand Down Expand Up @@ -67,7 +68,7 @@ public final class Pipeline {
private static final Codec<Pipeline> CODEC = new DelegatedCodec<>(
Proto2Codec.get(HddsProtos.Pipeline.getDefaultInstance()),
Pipeline::getFromProtobufSetCreationTimestamp,
p -> p.getProtobufMessage(ClientVersion.CURRENT.serialize()),
p -> p.getProtobufMessage(ClientVersion.CURRENT),
Pipeline.class,
DelegatedCodec.CopyType.UNSUPPORTED);

Expand Down Expand Up @@ -363,16 +364,41 @@ public ReplicationConfig getReplicationConfig() {
return replicationConfig;
}

public HddsProtos.Pipeline getProtobufMessage(int clientVersion) {
return getProtobufMessage(clientVersion, Collections.emptySet());
public HddsProtos.Pipeline getProtobufMessage(ClientVersion clientVersion) {
return getProtobufMessageInternal(clientVersion, Collections.emptySet(), null);
}

public HddsProtos.Pipeline getProtobufMessage(int clientVersion, Set<DatanodeDetails.Port.Name> filterPorts) {
/**
* 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,
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));
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));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand All @@ -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));
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@

import java.util.Random;
import java.util.concurrent.ThreadLocalRandom;
import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.ozone.test.GenericTestUtils;

Expand Down Expand Up @@ -95,7 +96,8 @@ public static DatanodeDetails createDatanodeDetails(DatanodeID id,
.setIpAddress(ipAddress)
.setNetworkLocation(networkLocation)
.setPersistedOpState(HddsProtos.NodeOperationalState.IN_SERVICE)
.setPersistedOpStateExpiry(0);
.setPersistedOpStateExpiry(0)
.setCurrentVersion(HDDSVersion.SOFTWARE_VERSION);

for (DatanodeDetails.Port.Name name : ALL_PORTS) {
dn.addPort(DatanodeDetails.newPort(name, port));
Expand Down
Loading