From ed7832eaf365df6b8ac2eb05ce0c728faa269fb9 Mon Sep 17 00:00:00 2001 From: Ethan Rose Date: Mon, 3 Aug 2026 17:17:22 -0400 Subject: [PATCH] Move version serialization into translator layers --- .../hadoop/hdds/protocol/DatanodeDetails.java | 10 ++--- .../common/helpers/ContainerWithPipeline.java | 3 +- .../apache/hadoop/hdds/scm/net/InnerNode.java | 3 +- .../org/apache/hadoop/hdds/scm/net/Node.java | 3 +- .../hadoop/hdds/scm/pipeline/Pipeline.java | 7 +-- .../hdds/protocol/TestDatanodeDetails.java | 13 +++--- .../hdds/scm/pipeline/TestPipeline.java | 10 ++--- .../DiskBalancerProtocolServer.java | 3 +- .../common/helpers/MoveDataNodePair.java | 4 +- .../hadoop/hdds/scm/net/InnerNodeImpl.java | 3 +- .../StorageContainerLocationProtocol.java | 9 ++-- ...ocationProtocolClientSideTranslatorPB.java | 8 ++-- .../hdds/scm/node/DatanodeUsageInfo.java | 5 ++- .../scm/pipeline/PipelineManagerImpl.java | 2 +- ...ocationProtocolServerSideTranslatorPB.java | 14 +++--- ...ocationProtocolServerSideTranslatorPB.java | 43 ++++++++++--------- .../scm/server/SCMClientProtocolServer.java | 11 ++--- .../hdds/scm/node/TestDatanodeUsageInfo.java | 4 +- .../scm/pipeline/MockPipelineManager.java | 4 +- .../TestPipelineDatanodesIntersection.java | 2 +- .../pipeline/TestPipelinePlacementPolicy.java | 6 +-- .../TestPipelineStateManagerImpl.java | 40 ++++++++--------- .../pipeline/TestRatisPipelineProvider.java | 16 +++---- .../pipeline/TestSimplePipelineProvider.java | 4 +- .../server/TestSCMBlockProtocolServer.java | 4 +- .../scm/cli/ContainerOperationClient.java | 8 ++-- .../hadoop/ozone/fsck/ContainerMapper.java | 2 +- .../debug/om/TestContainerToKeyMapping.java | 5 ++- .../om/helpers/KeyInfoWithVolumeContext.java | 3 +- .../hadoop/ozone/om/helpers/OmKeyInfo.java | 16 +++---- .../ozone/om/helpers/OmKeyLocationInfo.java | 5 ++- .../om/helpers/OmKeyLocationInfoGroup.java | 3 +- .../ozone/om/helpers/OmMultipartPartInfo.java | 2 +- .../ozone/om/helpers/OzoneFileStatus.java | 3 +- .../ozone/om/helpers/RepeatedOmKeyInfo.java | 6 +-- ...ManagerProtocolClientSideTranslatorPB.java | 4 +- .../ozone/om/helpers/TestOmKeyInfo.java | 6 +-- ...stractTestStorageDistributionEndpoint.java | 4 +- .../TestScmApplyTransactionFailure.java | 2 +- .../apache/hadoop/ozone/debug/TestLDBCli.java | 2 +- .../TestContainerCommandReconciliation.java | 8 ++-- .../ozone/om/OmMetadataManagerImpl.java | 2 +- .../om/request/file/OMFileCreateRequest.java | 7 ++- .../file/OMFileCreateRequestWithFSO.java | 3 +- .../request/file/OMRecoverLeaseRequest.java | 7 ++- .../request/key/OMAllocateBlockRequest.java | 4 +- .../om/request/key/OMKeyCreateRequest.java | 6 ++- .../key/OMKeyCreateRequestWithFSO.java | 3 +- .../S3MultipartUploadCommitPartRequest.java | 4 +- .../S3MultipartUploadCompleteRequest.java | 4 +- .../om/service/DirectoryDeletingService.java | 4 +- .../ozone/om/service/KeyDeletingService.java | 2 +- .../om/service/SnapshotDeletingService.java | 4 +- .../OzoneManagerRequestHandler.java | 36 ++++++++-------- .../ozone/om/request/OMRequestTestUtils.java | 4 +- .../file/TestOMRecoverLeaseRequest.java | 2 +- ...tOMDirectoriesPurgeRequestAndResponse.java | 4 +- .../TestOMSnapshotMoveTableKeysResponse.java | 4 +- .../service/TestSnapshotDeletingService.java | 4 +- .../TestOzoneManagerRequestHandler.java | 9 ++-- .../hadoop/ozone/recon/api/NodeEndpoint.java | 2 +- .../ozone/recon/scm/ReconPipelineManager.java | 2 +- .../StorageContainerServiceProviderImpl.java | 2 +- .../hadoop/ozone/recon/api/TestEndpoints.java | 5 ++- 64 files changed, 236 insertions(+), 198 deletions(-) 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 4b0108aa9aeb..e435fea43d3e 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 @@ -513,14 +513,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 filterPorts) { + public HddsProtos.DatanodeDetailsProto toProto(ClientVersion clientVersion, Set filterPorts) { return toProtoBuilder(clientVersion, filterPorts).build(); } @@ -533,7 +533,7 @@ public HddsProtos.DatanodeDetailsProto toProto(int clientVersion, Set * @return A {@link HddsProtos.DatanodeDetailsProto.Builder} Object. */ public HddsProtos.DatanodeDetailsProto.Builder toProtoBuilder( - int clientVersion, Set filterPorts) { + ClientVersion clientVersion, Set filterPorts) { final HddsProtos.DatanodeIDProto idProto = id.toProto(); final HddsProtos.DatanodeDetailsProto.Builder builder = @@ -1172,7 +1172,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(); 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 8db875c181ea..fad799a25d56 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 @@ -24,6 +24,7 @@ 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. @@ -53,7 +54,7 @@ public static ContainerWithPipeline fromProtobuf(HddsProtos.ContainerWithPipelin Pipeline.getFromProtobuf(allocatedContainer.getPipeline())); } - public HddsProtos.ContainerWithPipeline getProtobuf(int clientVersion) { + public HddsProtos.ContainerWithPipeline getProtobuf(ClientVersion clientVersion) { return HddsProtos.ContainerWithPipeline.newBuilder() .setContainerInfo(getContainerInfo().getProtobuf()) .setPipeline(getPipeline().getProtobufMessage(clientVersion, Name.IO_PORTS)) diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/net/InnerNode.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/net/InnerNode.java index ac9ec514a80b..9f1cf8acdf5c 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/net/InnerNode.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/net/InnerNode.java @@ -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. @@ -92,7 +93,7 @@ Node getLeaf(int leafIndex, List excludedScopes, Collection excludedNodes, int ancestorGen); @Override - HddsProtos.NetworkNode toProtobuf(int clientVersion); + HddsProtos.NetworkNode toProtobuf(ClientVersion clientVersion); @Override boolean equals(Object o); diff --git a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/net/Node.java b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/net/Node.java index 849111161bce..177f1a93c50d 100644 --- a/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/net/Node.java +++ b/hadoop-hdds/common/src/main/java/org/apache/hadoop/hdds/scm/net/Node.java @@ -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. @@ -130,7 +131,7 @@ public interface Node { boolean isDescendant(String nodePath); default HddsProtos.NetworkNode toProtobuf( - int clientVersion) { + ClientVersion clientVersion) { return null; } } 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 62c3855f1af5..7f2bc3352118 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 @@ -67,7 +67,7 @@ public final class Pipeline { private static final Codec 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); @@ -363,11 +363,12 @@ public ReplicationConfig getReplicationConfig() { return replicationConfig; } - public HddsProtos.Pipeline getProtobufMessage(int clientVersion) { + public HddsProtos.Pipeline getProtobufMessage(ClientVersion clientVersion) { return getProtobufMessage(clientVersion, Collections.emptySet()); } - public HddsProtos.Pipeline getProtobufMessage(int clientVersion, Set filterPorts) { + public HddsProtos.Pipeline getProtobufMessage(ClientVersion clientVersion, + Set filterPorts) { List members = new ArrayList<>(); List memberReplicaIndexes = new ArrayList<>(); 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 0f592c71d632..6401067c85b2 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 @@ -32,6 +32,7 @@ import org.apache.hadoop.hdds.protocol.DatanodeDetails.Port; import org.apache.hadoop.hdds.protocol.DatanodeDetails.Port.Name; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; +import org.apache.hadoop.ozone.ClientVersion; import org.junit.jupiter.api.Test; /** @@ -44,11 +45,11 @@ void protoIncludesNewPortsOnlyForV1() { DatanodeDetails subject = MockDatanodeDetails.randomDatanodeDetails(); HddsProtos.DatanodeDetailsProto proto = - subject.toProto(DEFAULT_VERSION.serialize()); + subject.toProto(DEFAULT_VERSION); assertPorts(proto, V0_PORTS); HddsProtos.DatanodeDetailsProto protoV1 = - subject.toProto(VERSION_HANDLES_UNKNOWN_DN_PORTS.serialize()); + subject.toProto(VERSION_HANDLES_UNKNOWN_DN_PORTS); assertPorts(protoV1, ALL_PORTS); } @@ -58,11 +59,11 @@ void testRequiredPortsProto() { Set requiredPorts = Stream.of(Port.Name.STANDALONE, Port.Name.RATIS) .collect(Collectors.toSet()); HddsProtos.DatanodeDetailsProto proto = - subject.toProto(subject.getCurrentVersion(), requiredPorts); + subject.toProto(ClientVersion.deserialize(subject.getCurrentVersion()), requiredPorts); assertPorts(proto, ImmutableSet.copyOf(requiredPorts)); HddsProtos.DatanodeDetailsProto ioPortProto = - subject.toProto(subject.getCurrentVersion(), Name.IO_PORTS); + subject.toProto(ClientVersion.deserialize(subject.getCurrentVersion()), Name.IO_PORTS); assertPorts(ioPortProto, ImmutableSet.copyOf(Name.IO_PORTS)); } @@ -74,7 +75,7 @@ public void testNewBuilderCurrentVersion() { Set requiredPorts = Stream.of(Port.Name.STANDALONE, Port.Name.RATIS) .collect(Collectors.toSet()); HddsProtos.DatanodeDetailsProto.Builder protoBuilder = - dn.toProtoBuilder(DEFAULT_VERSION.serialize(), requiredPorts); + dn.toProtoBuilder(DEFAULT_VERSION, requiredPorts); protoBuilder.clearCurrentVersion(); DatanodeDetails dn2 = DatanodeDetails.newBuilder(protoBuilder.build()).build(); assertEquals(HDDSVersion.SEPARATE_RATIS_PORTS_AVAILABLE, @@ -82,7 +83,7 @@ public void testNewBuilderCurrentVersion() { // test that if the current version is set, it is used protoBuilder = - dn.toProtoBuilder(DEFAULT_VERSION.serialize(), requiredPorts); + dn.toProtoBuilder(DEFAULT_VERSION, requiredPorts); DatanodeDetails dn3 = DatanodeDetails.newBuilder(protoBuilder.build()).build(); assertEquals(HDDSVersion.SOFTWARE_VERSION, HDDSVersion.deserialize(dn3.getCurrentVersion())); 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 c9318a9ef580..185f82fbe3cb 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 @@ -47,13 +47,13 @@ public void protoIncludesNewPortsOnlyForV1() throws IOException { Pipeline subject = MockPipeline.createPipeline(3); HddsProtos.Pipeline proto = - subject.getProtobufMessage(DEFAULT_VERSION.serialize()); + subject.getProtobufMessage(DEFAULT_VERSION); for (HddsProtos.DatanodeDetailsProto dn : proto.getMembersList()) { assertPorts(dn, V0_PORTS); } HddsProtos.Pipeline protoV1 = subject.getProtobufMessage( - VERSION_HANDLES_UNKNOWN_DN_PORTS.serialize()); + VERSION_HANDLES_UNKNOWN_DN_PORTS); for (HddsProtos.DatanodeDetailsProto dn : protoV1.getMembersList()) { assertPorts(dn, ALL_PORTS); } @@ -64,14 +64,14 @@ public void getProtobufMessageEC() throws IOException { Pipeline subject = MockPipeline.createPipeline(3); //when EC config is empty/null - HddsProtos.Pipeline protobufMessage = subject.getProtobufMessage(1); + HddsProtos.Pipeline protobufMessage = subject.getProtobufMessage(VERSION_HANDLES_UNKNOWN_DN_PORTS); assertEquals(0, protobufMessage.getEcReplicationConfig().getData()); //when EC config is NOT empty subject = MockPipeline.createEcPipeline(); - protobufMessage = subject.getProtobufMessage(1); + protobufMessage = subject.getProtobufMessage(VERSION_HANDLES_UNKNOWN_DN_PORTS); assertEquals(3, protobufMessage.getEcReplicationConfig().getData()); assertEquals(2, protobufMessage.getEcReplicationConfig().getParity()); @@ -80,7 +80,7 @@ public void getProtobufMessageEC() throws IOException { @Test public void testReplicaIndexesSerialisedCorrectly() { Pipeline pipeline = MockPipeline.createEcPipeline(); - HddsProtos.Pipeline protobufMessage = pipeline.getProtobufMessage(1); + HddsProtos.Pipeline protobufMessage = pipeline.getProtobufMessage(VERSION_HANDLES_UNKNOWN_DN_PORTS); Pipeline reloadedPipeline = Pipeline.getFromProtobuf(protobufMessage); for (DatanodeDetails dn : pipeline.getNodes()) { diff --git a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerProtocolServer.java b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerProtocolServer.java index 770f7b67414b..88028da7a928 100644 --- a/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerProtocolServer.java +++ b/hadoop-hdds/container-service/src/main/java/org/apache/hadoop/ozone/container/diskbalancer/DiskBalancerProtocolServer.java @@ -27,6 +27,7 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DiskBalancerConfigurationProto; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DiskBalancerRunningStatus; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.NodeOperationalState; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.container.common.statemachine.DatanodeStateMachine; import org.apache.hadoop.ozone.container.ozoneimpl.OzoneContainer; import org.slf4j.Logger; @@ -70,7 +71,7 @@ public DatanodeDiskBalancerInfoProto getDiskBalancerInfo(GetDiskBalancerInfoRequ DatanodeDetails datanodeDetails = datanodeStateMachine.getDatanodeDetails(); return DatanodeDiskBalancerInfoProto.newBuilder() - .setNode(datanodeDetails.toProto(request.getClientVersion())) + .setNode(datanodeDetails.toProto(ClientVersion.deserialize(request.getClientVersion()))) .setCurrentVolumeDensitySum(info.getVolumeDataDensity()) .setDiskBalancerConf(DiskBalancerConfigurationProto.newBuilder() .setThreshold(info.getThreshold()) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/container/common/helpers/MoveDataNodePair.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/container/common/helpers/MoveDataNodePair.java index bc43c2a5422c..95f8fee2a19f 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/container/common/helpers/MoveDataNodePair.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/container/common/helpers/MoveDataNodePair.java @@ -35,7 +35,7 @@ public class MoveDataNodePair { private static final Codec CODEC = new DelegatedCodec<>( Proto2Codec.get(MoveDataNodePairProto.getDefaultInstance()), MoveDataNodePair::getFromProtobuf, - pair -> pair.getProtobufMessage(ClientVersion.CURRENT.serialize()), + pair -> pair.getProtobufMessage(ClientVersion.CURRENT), MoveDataNodePair.class, DelegatedCodec.CopyType.SHALLOW); @@ -66,7 +66,7 @@ public DatanodeDetails getSrc() { return src; } - public MoveDataNodePairProto getProtobufMessage(int clientVersion) { + public MoveDataNodePairProto getProtobufMessage(ClientVersion clientVersion) { return MoveDataNodePairProto.newBuilder() .setSrc(src.toProto(clientVersion)) .setTgt(tgt.toProto(clientVersion)) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/net/InnerNodeImpl.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/net/InnerNodeImpl.java index c38a8800b5d6..383dfd20d58a 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/net/InnerNodeImpl.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/net/InnerNodeImpl.java @@ -31,6 +31,7 @@ import java.util.Map; import java.util.Objects; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; +import org.apache.hadoop.ozone.ClientVersion; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -460,7 +461,7 @@ public Node getLeaf(int leafIndex, List excludedScopes, @Override public HddsProtos.NetworkNode toProtobuf( - int clientVersion) { + ClientVersion clientVersion) { HddsProtos.InnerNode.Builder innerNode = HddsProtos.InnerNode.newBuilder() diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java index 505341ffa18d..cfa34891d5e2 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/StorageContainerLocationProtocol.java @@ -46,6 +46,7 @@ import org.apache.hadoop.hdds.scm.container.ReplicationManagerReport; import org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline; import org.apache.hadoop.hdds.scm.pipeline.Pipeline; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.upgrade.UpgradeFinalization.StatusAndMessages; import org.apache.hadoop.security.KerberosInfo; import org.apache.hadoop.security.token.Token; @@ -122,7 +123,7 @@ ContainerWithPipeline getContainerWithPipeline(long containerID) * @throws IOException */ List getContainerReplicas( - long containerId, int clientVersion) throws IOException; + long containerId, ClientVersion clientVersion) throws IOException; /** * Ask SCM the location of a batch of containers. SCM responds with a group of @@ -281,7 +282,7 @@ ContainerListResult listContainer(long startContainerID, */ List queryNode(HddsProtos.NodeOperationalState opState, HddsProtos.NodeState state, HddsProtos.QueryScope queryScope, - String poolName, int clientVersion) throws IOException; + String poolName, ClientVersion clientVersion) throws IOException; HddsProtos.Node queryNode(UUID uuid) throws IOException; @@ -495,7 +496,7 @@ StartContainerBalancerResponseProto startContainerBalancer( * @see org.apache.hadoop.ozone.ClientVersion */ List getDatanodeUsageInfo( - String address, String uuid, int clientVersion) throws IOException; + String address, String uuid, ClientVersion clientVersion) throws IOException; /** * Get usage information of most or least used datanodes. @@ -509,7 +510,7 @@ List getDatanodeUsageInfo( * @see org.apache.hadoop.ozone.ClientVersion */ List getDatanodeUsageInfo( - boolean mostUsed, int count, int clientVersion) throws IOException; + boolean mostUsed, int count, ClientVersion clientVersion) throws IOException; @Deprecated StatusAndMessages finalizeScmUpgrade(String upgradeClientID) diff --git a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java index a2a0b58ac259..c827bca8f523 100644 --- a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java +++ b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/StorageContainerLocationProtocolClientSideTranslatorPB.java @@ -341,7 +341,7 @@ public ContainerWithPipeline getContainerWithPipeline(long containerID) */ @Override public List getContainerReplicas( - long containerID, int clientVersion) throws IOException { + long containerID, ClientVersion clientVersion) throws IOException { Preconditions.checkState(containerID >= 0, "Container ID cannot be negative"); @@ -563,7 +563,7 @@ public Map> getContainersOnDecomNode(DatanodeDetails d public List queryNode( HddsProtos.NodeOperationalState opState, HddsProtos.NodeState nodeState, HddsProtos.QueryScope queryScope, String poolName, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { // TODO : We support only cluster wide query right now. So ignoring checking // queryScope and poolName NodeQueryRequestProto.Builder builder = NodeQueryRequestProto.newBuilder() @@ -1135,7 +1135,7 @@ public ContainerBalancerStatusInfoResponseProto getContainerBalancerStatusInfo() */ @Override public List getDatanodeUsageInfo( - String address, String uuid, int clientVersion) throws IOException { + String address, String uuid, ClientVersion clientVersion) throws IOException { DatanodeUsageInfoRequestProto request = DatanodeUsageInfoRequestProto.newBuilder() @@ -1161,7 +1161,7 @@ public List getDatanodeUsageInfo( */ @Override public List getDatanodeUsageInfo( - boolean mostUsed, int count, int clientVersion) throws IOException { + boolean mostUsed, int count, ClientVersion clientVersion) throws IOException { DatanodeUsageInfoRequestProto request = DatanodeUsageInfoRequestProto.newBuilder() .setMostUsed(mostUsed) diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeUsageInfo.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeUsageInfo.java index d6ac9edfaeb1..977f59278e4b 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeUsageInfo.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/node/DatanodeUsageInfo.java @@ -22,6 +22,7 @@ import org.apache.hadoop.hdds.protocol.DatanodeID; import org.apache.hadoop.hdds.protocol.proto.HddsProtos.DatanodeUsageInfoProto; import org.apache.hadoop.hdds.scm.container.placement.metrics.SCMNodeStat; +import org.apache.hadoop.ozone.ClientVersion; /** * Bundles datanode details with usage statistics. @@ -221,11 +222,11 @@ public int hashCode() { * * @return Protobuf HddsProtos.DatanodeUsageInfo */ - public DatanodeUsageInfoProto toProto(int clientVersion) { + public DatanodeUsageInfoProto toProto(ClientVersion clientVersion) { return toProtoBuilder(clientVersion).build(); } - private DatanodeUsageInfoProto.Builder toProtoBuilder(int clientVersion) { + private DatanodeUsageInfoProto.Builder toProtoBuilder(ClientVersion clientVersion) { DatanodeUsageInfoProto.Builder builder = DatanodeUsageInfoProto.newBuilder(); diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java index e176da2d7114..6ea92bb2a141 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/pipeline/PipelineManagerImpl.java @@ -268,7 +268,7 @@ private void checkIfPipelineCreationIsAllowed( private void addPipelineToManager(Pipeline pipeline) throws IOException { HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); acquireWriteLock(); try { stateManager.addPipeline(pipelineProto); 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 007376670ba4..cb6f5eabf524 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 @@ -53,6 +53,7 @@ import org.apache.hadoop.hdds.server.OzoneProtocolMessageDispatcher; 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.common.BlockGroup; import org.apache.hadoop.ozone.common.DeleteBlockGroupResult; import org.slf4j.Logger; @@ -120,6 +121,7 @@ public SCMBlockLocationResponse send(RpcController controller, private SCMBlockLocationResponse processMessage( SCMBlockLocationRequest request) throws ServiceException { + final ClientVersion clientVersion = ClientVersion.deserialize(request.getVersion()); SCMBlockLocationResponse.Builder response = createSCMBlockResponse( request.getCmdType(), request.getTraceID()); @@ -139,7 +141,7 @@ private SCMBlockLocationResponse processMessage( } } response.setAllocateScmBlockResponse(allocateScmBlock( - request.getAllocateScmBlockRequest(), request.getVersion())); + request.getAllocateScmBlockRequest(), clientVersion)); break; case DeleteScmKeyBlocks: response.setDeleteScmKeyBlocksResponse( @@ -155,12 +157,12 @@ private SCMBlockLocationResponse processMessage( break; case SortDatanodes: response.setSortDatanodesResponse(sortDatanodes( - request.getSortDatanodesRequest(), request.getVersion() + request.getSortDatanodesRequest(), clientVersion )); break; case GetClusterTree: response.setGetClusterTreeResponse( - getClusterTree(request.getVersion())); + getClusterTree(clientVersion)); break; default: // Should never happen @@ -189,7 +191,7 @@ private Status exceptionToResponseStatus(IOException ex) { } public AllocateScmBlockResponseProto allocateScmBlock( - AllocateScmBlockRequestProto request, int clientVersion) + AllocateScmBlockRequestProto request, ClientVersion clientVersion) throws IOException { List allocatedBlocks = impl.allocateBlock(request.getSize(), @@ -261,7 +263,7 @@ public HddsProtos.AddScmResponseProto getAddSCMResponse( } public SortDatanodesResponseProto sortDatanodes( - SortDatanodesRequestProto request, int clientVersion) + SortDatanodesRequestProto request, ClientVersion clientVersion) throws ServiceException { SortDatanodesResponseProto.Builder resp = SortDatanodesResponseProto.newBuilder(); @@ -280,7 +282,7 @@ public SortDatanodesResponseProto sortDatanodes( } } - public GetClusterTreeResponseProto getClusterTree(int clientVersion) + public GetClusterTreeResponseProto getClusterTree(ClientVersion clientVersion) throws IOException { GetClusterTreeResponseProto.Builder resp = GetClusterTreeResponseProto.newBuilder(); 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 a348b49f2763..cc7161581dba 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 @@ -229,7 +229,7 @@ public ScmContainerLocationResponse submitRequest(RpcController controller, // annotated interceptors. boolean checkResponseForECRepConfig = false; if (!ClientVersion.ERASURE_CODING_SUPPORT.isSupportedBy( - request.getVersion())) { + ClientVersion.deserialize(request.getVersion()))) { if (request.getCmdType() == GetContainer || request.getCmdType() == ListContainer || request.getCmdType() == GetContainerWithPipeline @@ -419,6 +419,7 @@ private void disallowECReplicationConfigInGetPipelineResponse( @SuppressWarnings("checkstyle:methodlength") public ScmContainerLocationResponse processRequest( ScmContainerLocationRequest request) throws ServiceException { + final ClientVersion clientVersion = ClientVersion.deserialize(request.getVersion()); try { switch (request.getCmdType()) { case AllocateContainer: @@ -426,7 +427,7 @@ public ScmContainerLocationResponse processRequest( .setCmdType(request.getCmdType()) .setStatus(Status.OK) .setContainerResponse(allocateContainer( - request.getContainerRequest(), request.getVersion())) + request.getContainerRequest(), clientVersion)) .build(); case GetContainer: return ScmContainerLocationResponse.newBuilder() @@ -448,7 +449,7 @@ public ScmContainerLocationResponse processRequest( .setStatus(Status.OK) .setGetContainerWithPipelineResponse(getContainerWithPipeline( request.getGetContainerWithPipelineRequest(), - request.getVersion())) + clientVersion)) .build(); case GetContainerWithPipelineBatch: return ScmContainerLocationResponse.newBuilder() @@ -457,7 +458,7 @@ public ScmContainerLocationResponse processRequest( .setGetContainerWithPipelineBatchResponse( getContainerWithPipelineBatch( request.getGetContainerWithPipelineBatchRequest(), - request.getVersion())) + clientVersion)) .build(); case GetExistContainerWithPipelinesInBatch: return ScmContainerLocationResponse.newBuilder() @@ -466,7 +467,7 @@ public ScmContainerLocationResponse processRequest( .setGetExistContainerWithPipelinesInBatchResponse( getExistContainerWithPipelinesInBatch( request.getGetExistContainerWithPipelinesInBatchRequest(), - request.getVersion())) + clientVersion)) .build(); case ListContainer: return ScmContainerLocationResponse.newBuilder() @@ -480,7 +481,7 @@ public ScmContainerLocationResponse processRequest( .setCmdType(request.getCmdType()) .setStatus(Status.OK) .setNodeQueryResponse(queryNode(request.getNodeQueryRequest(), - request.getVersion())) + clientVersion)) .build(); case SingleNodeQuery: return ScmContainerLocationResponse.newBuilder() @@ -512,14 +513,14 @@ public ScmContainerLocationResponse processRequest( .setCmdType(request.getCmdType()) .setStatus(Status.OK) .setPipelineResponse(allocatePipeline( - request.getPipelineRequest(), request.getVersion())) + request.getPipelineRequest(), clientVersion)) .build(); case ListPipelines: return ScmContainerLocationResponse.newBuilder() .setCmdType(request.getCmdType()) .setStatus(Status.OK) .setListPipelineResponse(listPipelines( - request.getListPipelineRequest(), request.getVersion())) + request.getListPipelineRequest(), clientVersion)) .build(); case ActivatePipeline: return ScmContainerLocationResponse.newBuilder() @@ -624,7 +625,7 @@ public ScmContainerLocationResponse processRequest( .setCmdType(request.getCmdType()) .setStatus(Status.OK) .setGetPipelineResponse(getPipeline( - request.getGetPipelineRequest(), request.getVersion())) + request.getGetPipelineRequest(), clientVersion)) .build(); case GetSafeModeRuleStatuses: return ScmContainerLocationResponse.newBuilder() @@ -680,7 +681,7 @@ public ScmContainerLocationResponse processRequest( .setStatus(Status.OK) .setDatanodeUsageInfoResponse(getDatanodeUsageInfo( request.getDatanodeUsageInfoRequest(), - request.getVersion())) + clientVersion)) .build(); case GetContainerCount: return ScmContainerLocationResponse.newBuilder() @@ -702,7 +703,7 @@ public ScmContainerLocationResponse processRequest( .setStatus(Status.OK) .setGetContainerReplicasResponse(getContainerReplicas( request.getGetContainerReplicasRequest(), - request.getVersion())) + clientVersion)) .build(); case GetFailedDeletedBlocksTransaction: return ScmContainerLocationResponse.newBuilder() @@ -791,7 +792,7 @@ public ScmContainerLocationResponse processRequest( } public GetContainerReplicasResponseProto getContainerReplicas( - GetContainerReplicasRequestProto request, int clientVersion) + GetContainerReplicasRequestProto request, ClientVersion clientVersion) throws IOException { List replicas = impl.getContainerReplicas(request.getContainerID(), clientVersion); @@ -800,7 +801,7 @@ public GetContainerReplicasResponseProto getContainerReplicas( } public ContainerResponseProto allocateContainer(ContainerRequestProto request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { ReplicationConfig replicationConfig = ReplicationConfig.fromProto(request.getReplicationType(), request.getReplicationFactor(), request.getEcReplicationConfig() @@ -834,7 +835,7 @@ public GetContainerTokenResponseProto getContainerToken( public GetContainerWithPipelineResponseProto getContainerWithPipeline( GetContainerWithPipelineRequestProto request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { ContainerWithPipeline container = impl .getContainerWithPipeline(request.getContainerID()); return GetContainerWithPipelineResponseProto.newBuilder() @@ -845,7 +846,7 @@ public GetContainerWithPipelineResponseProto getContainerWithPipeline( public GetContainerWithPipelineBatchResponseProto getContainerWithPipelineBatch( GetContainerWithPipelineBatchRequestProto request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { List containers = impl .getContainerWithPipelineBatch(request.getContainerIDsList()); GetContainerWithPipelineBatchResponseProto.Builder builder = @@ -859,7 +860,7 @@ public GetContainerWithPipelineResponseProto getContainerWithPipeline( public GetExistContainerWithPipelinesInBatchResponseProto getExistContainerWithPipelinesInBatch( GetExistContainerWithPipelinesInBatchRequestProto request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { List containers = impl .getExistContainerWithPipelinesInBatch(request.getContainerIDsList()); GetExistContainerWithPipelinesInBatchResponseProto.Builder builder = @@ -939,7 +940,7 @@ public SCMDeleteContainerResponseProto deleteContainer( public NodeQueryResponseProto queryNode( StorageContainerLocationProtocolProtos.NodeQueryRequestProto request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { HddsProtos.NodeOperationalState opState = null; HddsProtos.NodeState nodeState = null; @@ -989,7 +990,7 @@ public SCMCloseContainerResponseProto closeContainer( public PipelineResponseProto allocatePipeline( StorageContainerLocationProtocolProtos.PipelineRequestProto request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { Pipeline pipeline = impl.createReplicationPipeline( request.getReplicationType(), request.getReplicationFactor(), HddsProtos.NodePool.getDefaultInstance()); @@ -1003,7 +1004,7 @@ public PipelineResponseProto allocatePipeline( } public ListPipelineResponseProto listPipelines( - ListPipelineRequestProto request, int clientVersion) + ListPipelineRequestProto request, ClientVersion clientVersion) throws IOException { ListPipelineResponseProto.Builder builder = ListPipelineResponseProto .newBuilder(); @@ -1016,7 +1017,7 @@ public ListPipelineResponseProto listPipelines( public GetPipelineResponseProto getPipeline( GetPipelineRequestProto request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { GetPipelineResponseProto.Builder builder = GetPipelineResponseProto .newBuilder(); Pipeline pipeline = impl.getPipeline(request.getPipelineID()); @@ -1351,7 +1352,7 @@ public StartMaintenanceNodesResponseProto startMaintenanceNodes( public DatanodeUsageInfoResponseProto getDatanodeUsageInfo( StorageContainerLocationProtocolProtos.DatanodeUsageInfoRequestProto - request, int clientVersion) throws IOException { + request, ClientVersion clientVersion) throws IOException { List infoList; // get info by ip or uuid diff --git a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java index 54434de80ec3..dcf6611a66d7 100644 --- a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java +++ b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMClientProtocolServer.java @@ -113,6 +113,7 @@ import org.apache.hadoop.ipc_.ProtobufRpcEngine; import org.apache.hadoop.ipc_.RPC; import org.apache.hadoop.ipc_.Server; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.audit.AuditAction; import org.apache.hadoop.ozone.audit.AuditEventStatus; @@ -354,7 +355,7 @@ public ContainerWithPipeline getContainerWithPipeline(long containerID) @Override public List getContainerReplicas( - long containerId, int clientVersion) throws IOException { + long containerId, ClientVersion clientVersion) throws IOException { List results = new ArrayList<>(); Map auditMap = new HashMap<>(); auditMap.put("containerId", String.valueOf(containerId)); @@ -677,7 +678,7 @@ public Map> getContainersOnDecomNode(DatanodeDetails d @Override public List queryNode( HddsProtos.NodeOperationalState opState, HddsProtos.NodeState state, - HddsProtos.QueryScope queryScope, String poolName, int clientVersion) + HddsProtos.QueryScope queryScope, String poolName, ClientVersion clientVersion) throws IOException { final Map auditMap = Maps.newHashMap(); auditMap.put("opState", String.valueOf(opState)); @@ -1466,7 +1467,7 @@ public ContainerBalancerStatusInfoResponseProto getContainerBalancerStatusInfo() */ @Override public List getDatanodeUsageInfo( - String address, String uuid, int clientVersion) throws IOException { + String address, String uuid, ClientVersion clientVersion) throws IOException { final Map auditMap = Maps.newHashMap(); auditMap.put("address", address); @@ -1511,7 +1512,7 @@ public List getDatanodeUsageInfo( * @return Usage info such as capacity, SCMUsed, and remaining space. */ private HddsProtos.DatanodeUsageInfoProto getUsageInfoFromDatanodeDetails( - DatanodeDetails node, int clientVersion) { + DatanodeDetails node, ClientVersion clientVersion) { DatanodeUsageInfo usageInfo = scm.getScmNodeManager().getUsageInfo(node); return usageInfo.toProto(clientVersion); } @@ -1530,7 +1531,7 @@ private HddsProtos.DatanodeUsageInfoProto getUsageInfoFromDatanodeDetails( */ @Override public List getDatanodeUsageInfo( - boolean mostUsed, int count, int clientVersion) + boolean mostUsed, int count, ClientVersion clientVersion) throws IOException, IllegalArgumentException { final Map auditMap = Maps.newHashMap(); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDatanodeUsageInfo.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDatanodeUsageInfo.java index 142e2638cf8d..bb02756b4d2e 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDatanodeUsageInfo.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/node/TestDatanodeUsageInfo.java @@ -41,7 +41,7 @@ void testToProtoDoesNotIncludeFilesystemFieldsByDefault() { ); DatanodeUsageInfo info = new DatanodeUsageInfo(dn, stat); - DatanodeUsageInfoProto proto = info.toProto(ClientVersion.CURRENT.serialize()); + DatanodeUsageInfoProto proto = info.toProto(ClientVersion.CURRENT); assertThat(proto.hasFsCapacity()).isFalse(); assertThat(proto.hasFsAvailable()).isFalse(); @@ -59,7 +59,7 @@ void testToProtoIncludesFilesystemFieldsWhenPresent() { DatanodeUsageInfo info = new DatanodeUsageInfo(dn, stat); info.setFilesystemUsage(2000L, 1500L); - DatanodeUsageInfoProto proto = info.toProto(ClientVersion.CURRENT.serialize()); + DatanodeUsageInfoProto proto = info.toProto(ClientVersion.CURRENT); assertThat(proto.hasFsCapacity()).isTrue(); assertThat(proto.hasFsAvailable()).isTrue(); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java index f8abec6c4997..918a310e3129 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/MockPipelineManager.java @@ -87,7 +87,7 @@ public Pipeline createPipeline(ReplicationConfig replicationConfig, } stateManager.addPipeline(pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize())); + ClientVersion.CURRENT)); return pipeline; } @@ -111,7 +111,7 @@ public Pipeline buildECPipeline(ReplicationConfig replicationConfig, public void addEcPipeline(Pipeline pipeline) throws IOException { stateManager.addPipeline(pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize())); + ClientVersion.CURRENT)); } @Override diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineDatanodesIntersection.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineDatanodesIntersection.java index 5a642369d969..87211083109b 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineDatanodesIntersection.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineDatanodesIntersection.java @@ -110,7 +110,7 @@ public void testPipelineDatanodesIntersection(int nodeCount, Pipeline pipeline = provider.create(RatisReplicationConfig.getInstance( ReplicationFactor.THREE)); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); nodeManager.addPipeline(pipeline); List overlapPipelines = RatisPipelineUtils diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java index e3100c217b68..103bd388ff14 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelinePlacementPolicy.java @@ -292,7 +292,7 @@ public void testPickLowestLoadAnchor() throws IOException, TimeoutException { .setNodes(nodes) .build(); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); nodeManager.addPipeline(pipeline); stateManager.addPipeline(pipelineProto); } catch (SCMException e) { @@ -648,7 +648,7 @@ private void insertHeavyNodesIntoNodeManager( .build(); pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); nodeManager.addPipeline(pipeline); stateManager.addPipeline(pipelineProto); pipelineCount++; @@ -791,7 +791,7 @@ private void createPipelineWithReplicationConfig(List dnList, .build(); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); nodeManager.addPipeline(pipeline); stateManager.addPipeline(pipelineProto); } diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineStateManagerImpl.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineStateManagerImpl.java index fc447a74e915..977c4bcb16f2 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineStateManagerImpl.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestPipelineStateManagerImpl.java @@ -110,14 +110,14 @@ private Pipeline createDummyPipeline(HddsProtos.ReplicationType type, public void testAddAndGetPipeline() throws IOException, TimeoutException { Exception e = assertThrows(SCMException.class, () -> stateManager.addPipeline(createDummyPipeline(0) - .getProtobufMessage(ClientVersion.CURRENT.serialize()))); + .getProtobufMessage(ClientVersion.CURRENT))); // replication factor and number of nodes in the pipeline do not match assertThat(e.getMessage()).contains("do not match"); // add a pipeline Pipeline pipeline = createDummyPipeline(1); HddsProtos.Pipeline pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); try { stateManager.addPipeline(pipelineProto); @@ -144,11 +144,11 @@ public void testGetPipelines() throws IOException, TimeoutException { Set pipelines = new HashSet<>(); HddsProtos.Pipeline pipeline = createDummyPipeline(1).getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipeline); pipelines.add(pipeline); pipeline = createDummyPipeline(1).getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipeline); pipelines.add(pipeline); @@ -179,19 +179,19 @@ public void testGetPipelinesByTypeAndFactor() // 5 pipelines in allocated state for each type and factor HddsProtos.Pipeline pipeline = createDummyPipeline(type, factor, factor.getNumber()) - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipeline); pipelines.add(pipeline); // 5 pipelines in open state for each type and factor pipeline = createDummyPipeline(type, factor, factor.getNumber()) - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipeline); pipelines.add(pipeline); // 5 pipelines in closed state for each type and factor pipeline = createDummyPipeline(type, factor, factor.getNumber()) - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipeline); pipelines.add(pipeline); } @@ -232,20 +232,20 @@ public void testGetPipelinesByTypeFactorAndState() // 5 pipelines in allocated state for each type and factor HddsProtos.Pipeline pipeline = createDummyPipeline(type, factor, factor.getNumber()) - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipeline); pipelines.add(pipeline); // 5 pipelines in open state for each type and factor pipeline = createDummyPipeline(type, factor, factor.getNumber()) - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipeline); openPipeline(pipeline); pipelines.add(pipeline); // 5 pipelines in dormant state for each type and factor pipeline = createDummyPipeline(type, factor, factor.getNumber()) - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipeline); openPipeline(pipeline); deactivatePipeline(pipeline); @@ -253,7 +253,7 @@ public void testGetPipelinesByTypeFactorAndState() // 5 pipelines in closed state for each type and factor pipeline = createDummyPipeline(type, factor, factor.getNumber()) - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipeline); finalizePipeline(pipeline); pipelines.add(pipeline); @@ -292,7 +292,7 @@ public void testAddAndGetContainer() throws IOException, TimeoutException { long containerID = 0; Pipeline pipeline = createDummyPipeline(1); HddsProtos.Pipeline pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); pipeline = stateManager.getPipeline(pipeline.getId()); stateManager.addContainerToPipeline(pipeline.getId(), @@ -325,7 +325,7 @@ public void testAddAndGetContainer() throws IOException, TimeoutException { public void testRemovePipeline() throws IOException, TimeoutException { Pipeline pipeline = createDummyPipeline(1); HddsProtos.Pipeline pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); // close the pipeline openPipeline(pipelineProto); @@ -347,7 +347,7 @@ public void testRemoveContainer() throws IOException, TimeoutException { long containerID = 1; Pipeline pipeline = createDummyPipeline(1); HddsProtos.Pipeline pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); // create an open pipeline in stateMap stateManager.addPipeline(pipelineProto); openPipeline(pipelineProto); @@ -387,7 +387,7 @@ public void testRemoveContainer() throws IOException, TimeoutException { public void testFinalizePipeline() throws IOException, TimeoutException { Pipeline pipeline = createDummyPipeline(1); HddsProtos.Pipeline pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); // finalize on ALLOCATED pipeline finalizePipeline(pipelineProto); @@ -398,7 +398,7 @@ public void testFinalizePipeline() throws IOException, TimeoutException { pipeline = createDummyPipeline(1); pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); openPipeline(pipelineProto); // finalize on OPEN pipeline @@ -410,7 +410,7 @@ public void testFinalizePipeline() throws IOException, TimeoutException { pipeline = createDummyPipeline(1); pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); openPipeline(pipelineProto); finalizePipeline(pipelineProto); @@ -426,7 +426,7 @@ public void testFinalizePipeline() throws IOException, TimeoutException { public void testOpenPipeline() throws IOException, TimeoutException { Pipeline pipeline = createDummyPipeline(1); HddsProtos.Pipeline pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); // open on ALLOCATED pipeline openPipeline(pipelineProto); @@ -448,7 +448,7 @@ public void testQueryPipeline() throws IOException, TimeoutException { HddsProtos.ReplicationFactor.THREE, 3); // pipeline in allocated state should not be reported HddsProtos.Pipeline pipelineProto = pipeline - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); assertEquals(0, stateManager .getPipelines(RatisReplicationConfig @@ -470,7 +470,7 @@ public void testQueryPipeline() throws IOException, TimeoutException { .setState(Pipeline.PipelineState.OPEN) .build(); HddsProtos.Pipeline pipelineProto2 = pipeline2 - .getProtobufMessage(ClientVersion.CURRENT.serialize()); + .getProtobufMessage(ClientVersion.CURRENT); // pipeline in open state should be reported stateManager.addPipeline(pipelineProto2); assertEquals(2, stateManager diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java index 04a935f68470..29faf3ea238c 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestRatisPipelineProvider.java @@ -141,14 +141,14 @@ private void createPipelineAndAssertions( assertPipelineProperties(pipeline, factor, REPLICATION_TYPE, Pipeline.PipelineState.ALLOCATED); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); nodeManager.addPipeline(pipeline); Pipeline pipeline1 = provider.create(RatisReplicationConfig .getInstance(factor)); HddsProtos.Pipeline pipelineProto1 = pipeline1.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); assertPipelineProperties(pipeline1, factor, REPLICATION_TYPE, Pipeline.PipelineState.ALLOCATED); // New pipeline should not overlap with the previous created pipeline @@ -190,7 +190,7 @@ public void testCreatePipelineWithFactor() throws Exception { assertPipelineProperties(pipeline, factor, REPLICATION_TYPE, Pipeline.PipelineState.ALLOCATED); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); factor = HddsProtos.ReplicationFactor.ONE; @@ -199,7 +199,7 @@ public void testCreatePipelineWithFactor() throws Exception { assertPipelineProperties(pipeline1, factor, REPLICATION_TYPE, Pipeline.PipelineState.ALLOCATED); HddsProtos.Pipeline pipelineProto1 = pipeline1.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto1); // With enough pipeline quote on datanodes, they should not share // the same set of datanodes. @@ -279,7 +279,7 @@ public void testCreatePipelinesDnExclude() throws Exception { assertPipelineProperties(pipeline, factor, REPLICATION_TYPE, Pipeline.PipelineState.ALLOCATED); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); nodeManager.addPipeline(pipeline); stateManager.addPipeline(pipelineProto); @@ -406,7 +406,7 @@ public void testCreatePipelineWithDefaultLimit() throws Exception { Pipeline p = provider.create( RatisReplicationConfig.getInstance(ReplicationFactor.THREE), new ArrayList<>(), new ArrayList<>()); - stateManager.addPipeline(p.getProtobufMessage(ClientVersion.CURRENT.serialize())); + stateManager.addPipeline(p.getProtobufMessage(ClientVersion.CURRENT)); } // Next pipeline creation should fail with default limit message. @@ -431,7 +431,7 @@ public void testCreatePipelineThrowErrorWithDataNodeLimit(int limit, int pipelin for (int i = 0; i < pipelineCount; i++) { stateManager.addPipeline( provider.create(RatisReplicationConfig.getInstance(ReplicationFactor.THREE), - new ArrayList<>(), new ArrayList<>()).getProtobufMessage(ClientVersion.CURRENT.serialize()) + new ArrayList<>(), new ArrayList<>()).getProtobufMessage(ClientVersion.CURRENT) ); } @@ -458,7 +458,7 @@ private void addPipeline( .setId(PipelineID.randomId()) .build(); HddsProtos.Pipeline pipelineProto = openPipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); nodeManager.addPipeline(openPipeline); diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSimplePipelineProvider.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSimplePipelineProvider.java index 07a064a2b7d0..d2d4e9b6e266 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSimplePipelineProvider.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/pipeline/TestSimplePipelineProvider.java @@ -81,7 +81,7 @@ public void testCreatePipelineWithFactor() throws Exception { Pipeline pipeline = provider.create(StandaloneReplicationConfig.getInstance(factor)); HddsProtos.Pipeline pipelineProto = pipeline.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto); assertEquals(pipeline.getType(), HddsProtos.ReplicationType.STAND_ALONE); assertEquals(pipeline.getReplicationConfig().getRequiredNodes(), factor.getNumber()); @@ -92,7 +92,7 @@ public void testCreatePipelineWithFactor() throws Exception { Pipeline pipeline1 = provider.create(StandaloneReplicationConfig.getInstance(factor)); HddsProtos.Pipeline pipelineProto1 = pipeline1.getProtobufMessage( - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); stateManager.addPipeline(pipelineProto1); assertEquals(pipeline1.getType(), HddsProtos.ReplicationType.STAND_ALONE); assertEquals( diff --git a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestSCMBlockProtocolServer.java b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestSCMBlockProtocolServer.java index 0fe156448eef..fa733d3712a8 100644 --- a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestSCMBlockProtocolServer.java +++ b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestSCMBlockProtocolServer.java @@ -274,7 +274,7 @@ public void testSortDatanodes() throws Exception { .setClient(client) .build(); ScmBlockLocationProtocolProtos.SortDatanodesResponseProto resp = - service.sortDatanodes(request, ClientVersion.CURRENT.serialize()); + service.sortDatanodes(request, ClientVersion.CURRENT); assertEquals(NODE_COUNT, resp.getNodeList().size()); System.out.println("client = " + client); resp.getNodeList().stream().forEach( @@ -290,7 +290,7 @@ public void testSortDatanodes() throws Exception { .addAllNodeNetworkName(nodes) .setClient(client) .build(); - resp = service.sortDatanodes(request, ClientVersion.CURRENT.serialize()); + resp = service.sortDatanodes(request, ClientVersion.CURRENT); System.out.println("client = " + client); assertEquals(0, resp.getNodeList().size()); resp.getNodeList().stream().forEach( diff --git a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java index 9f14744d2a5a..d03a6215a681 100644 --- a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java +++ b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/hdds/scm/cli/ContainerOperationClient.java @@ -252,7 +252,7 @@ public List queryNode( HddsProtos.QueryScope queryScope, String poolName) throws IOException { return storageContainerLocationClient.queryNode(opState, nodeState, - queryScope, poolName, ClientVersion.CURRENT.serialize()); + queryScope, poolName, ClientVersion.CURRENT); } @Override @@ -467,7 +467,7 @@ public ContainerWithPipeline getContainerWithPipeline(long containerId) public List getContainerReplicas(long containerId) throws IOException { List protos = storageContainerLocationClient.getContainerReplicas(containerId, - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); List replicas = new ArrayList<>(); for (HddsProtos.SCMContainerReplicaProto p : protos) { replicas.add(ContainerReplicaInfo.fromProto(p)); @@ -590,14 +590,14 @@ public DeletedBlocksTransactionSummary getDeletedBlockSummary() throws IOExcepti public List getDatanodeUsageInfo( String address, String uuid) throws IOException { return storageContainerLocationClient.getDatanodeUsageInfo(address, - uuid, ClientVersion.CURRENT.serialize()); + uuid, ClientVersion.CURRENT); } @Override public List getDatanodeUsageInfo( boolean mostUsed, int count) throws IOException { return storageContainerLocationClient.getDatanodeUsageInfo(mostUsed, count, - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); } @Override diff --git a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/fsck/ContainerMapper.java b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/fsck/ContainerMapper.java index 294b06f0c02d..42a7726d0399 100644 --- a/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/fsck/ContainerMapper.java +++ b/hadoop-ozone/cli-debug/src/main/java/org/apache/hadoop/ozone/fsck/ContainerMapper.java @@ -88,7 +88,7 @@ public static void main(String[] args) throws IOException { keyValueTableIterator.next(); OmKeyInfo omKeyInfo = keyValue.getValue(); byte[] value = omKeyInfo - .getProtobuf(true, ClientVersion.CURRENT.serialize()) + .getProtobuf(true, ClientVersion.CURRENT) .toByteArray(); OmKeyInfo keyInfo = OmKeyInfo.getFromProtobuf( OzoneManagerProtocolProtos.KeyInfo.parseFrom(value)); diff --git a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/om/TestContainerToKeyMapping.java b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/om/TestContainerToKeyMapping.java index 4cad62cd719c..fd1650683b00 100644 --- a/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/om/TestContainerToKeyMapping.java +++ b/hadoop-ozone/cli-debug/src/test/java/org/apache/hadoop/ozone/debug/om/TestContainerToKeyMapping.java @@ -31,6 +31,7 @@ import org.apache.hadoop.hdds.client.StandaloneReplicationConfig; import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.debug.OzoneDebug; import org.apache.hadoop.ozone.om.OMMetadataManager; import org.apache.hadoop.ozone.om.OmMetadataManagerImpl; @@ -304,7 +305,7 @@ private void createMultipartUpload() throws Exception { // Create part 1 with a block in container 4 OmKeyInfo part1Info = createOBSKeyInfo( mpuKeyName + "/" + uploadId + "/part-1", MPU_PART1_ID, CONTAINER_ID_4); - KeyInfo part1Proto = part1Info.getProtobuf(true, 0); + KeyInfo part1Proto = part1Info.getProtobuf(true, ClientVersion.DEFAULT_VERSION); PartKeyInfo partKeyInfo1 = PartKeyInfo.newBuilder() .setPartName(mpuKeyName + "/" + uploadId + "/part-1") .setPartNumber(1) @@ -314,7 +315,7 @@ private void createMultipartUpload() throws Exception { // Create part 2 with a block in container 4 OmKeyInfo part2Info = createOBSKeyInfo( mpuKeyName + "/" + uploadId + "/part-2", MPU_PART2_ID, CONTAINER_ID_4); - KeyInfo part2Proto = part2Info.getProtobuf(true, 0); + KeyInfo part2Proto = part2Info.getProtobuf(true, ClientVersion.DEFAULT_VERSION); PartKeyInfo partKeyInfo2 = PartKeyInfo.newBuilder() .setPartName(mpuKeyName + "/" + uploadId + "/part-2") .setPartNumber(2) diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/KeyInfoWithVolumeContext.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/KeyInfoWithVolumeContext.java index d6d54d3c174d..8a164efe8ef6 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/KeyInfoWithVolumeContext.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/KeyInfoWithVolumeContext.java @@ -19,6 +19,7 @@ import java.io.IOException; import java.util.Optional; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.GetKeyInfoResponse; /** @@ -55,7 +56,7 @@ public static KeyInfoWithVolumeContext fromProtobuf( .build(); } - public GetKeyInfoResponse toProtobuf(int clientVersion) { + public GetKeyInfoResponse toProtobuf(ClientVersion clientVersion) { GetKeyInfoResponse.Builder builder = GetKeyInfoResponse.newBuilder(); volumeArgs.ifPresent(v -> builder.setVolumeInfo(v.getProtobuf())); userPrincipal.ifPresent(builder::setUserPrincipal); diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyInfo.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyInfo.java index ca911f895c77..7258403e8f18 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyInfo.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyInfo.java @@ -143,7 +143,7 @@ private static Codec newCodec(boolean isOpenKey) { return new DelegatedCodec<>( Proto2Codec.get(KeyInfo.getDefaultInstance()), OmKeyInfo::getFromProtobuf, - k -> k.getProtobuf(true, ClientVersion.CURRENT.serialize(), isOpenKey), + k -> k.getProtobuf(true, ClientVersion.CURRENT, isOpenKey), OmKeyInfo.class); } @@ -724,7 +724,7 @@ protected OmKeyInfo buildObject() { * For network transmit. * @return KeyInfo */ - public KeyInfo getProtobuf(int clientVersion) { + public KeyInfo getProtobuf(ClientVersion clientVersion) { return getProtobuf(false, clientVersion); } @@ -734,7 +734,7 @@ public KeyInfo getProtobuf(int clientVersion) { * @param latestVersion * @return key info. */ - public KeyInfo getNetworkProtobuf(int clientVersion, boolean latestVersion) { + public KeyInfo getNetworkProtobuf(ClientVersion clientVersion, boolean latestVersion) { return getProtobuf(false, null, clientVersion, latestVersion); } @@ -746,7 +746,7 @@ public KeyInfo getNetworkProtobuf(int clientVersion, boolean latestVersion) { * @param latestVersion * @return key info with the user given full key name */ - public KeyInfo getNetworkProtobuf(String fullKeyName, int clientVersion, + public KeyInfo getNetworkProtobuf(String fullKeyName, ClientVersion clientVersion, boolean latestVersion) { return getProtobuf(false, fullKeyName, clientVersion, latestVersion); } @@ -756,7 +756,7 @@ public KeyInfo getNetworkProtobuf(String fullKeyName, int clientVersion, * @param ignorePipeline true for persist to DB, false for network transmit. * @return KeyInfo */ - public KeyInfo getProtobuf(boolean ignorePipeline, int clientVersion) { + public KeyInfo getProtobuf(boolean ignorePipeline, ClientVersion clientVersion) { return getProtobuf(ignorePipeline, null, clientVersion, false, true); } @@ -768,7 +768,7 @@ public KeyInfo getProtobuf(boolean ignorePipeline, int clientVersion) { * @param isOpenKey true for openKeyTable, false for keyTable * @return KeyInfo */ - public KeyInfo getProtobuf(boolean ignorePipeline, int clientVersion, + public KeyInfo getProtobuf(boolean ignorePipeline, ClientVersion clientVersion, boolean isOpenKey) { return getProtobuf(ignorePipeline, null, clientVersion, false, isOpenKey); } @@ -781,7 +781,7 @@ public KeyInfo getProtobuf(boolean ignorePipeline, int clientVersion, * @return key info object */ private KeyInfo getProtobuf(boolean ignorePipeline, String fullKeyName, - int clientVersion, boolean latestVersionBlocks) { + ClientVersion clientVersion, boolean latestVersionBlocks) { return getProtobuf(ignorePipeline, fullKeyName, clientVersion, latestVersionBlocks, true); } @@ -796,7 +796,7 @@ private KeyInfo getProtobuf(boolean ignorePipeline, String fullKeyName, * @return key info object */ private KeyInfo getProtobuf(boolean ignorePipeline, String fullKeyName, - int clientVersion, boolean latestVersionBlocks, + ClientVersion clientVersion, boolean latestVersionBlocks, boolean isOpenKey) { long latestVersion = keyLocationVersions.isEmpty() ? -1 : keyLocationVersions.get(keyLocationVersions.size() - 1).getVersion(); diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyLocationInfo.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyLocationInfo.java index d3fea73b211a..4de931978b92 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyLocationInfo.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyLocationInfo.java @@ -21,6 +21,7 @@ import org.apache.hadoop.hdds.scm.pipeline.Pipeline; import org.apache.hadoop.hdds.scm.storage.BlockLocationInfo; import org.apache.hadoop.hdds.security.token.OzoneBlockTokenIdentifier; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyLocation; import org.apache.hadoop.ozone.protocolPB.OMPBHelper; import org.apache.hadoop.security.token.Token; @@ -88,11 +89,11 @@ public OmKeyLocationInfo build() { } } - public KeyLocation getProtobuf(int clientVersion) { + public KeyLocation getProtobuf(ClientVersion clientVersion) { return getProtobuf(false, clientVersion); } - public KeyLocation getProtobuf(boolean ignorePipeline, int clientVersion) { + public KeyLocation getProtobuf(boolean ignorePipeline, ClientVersion clientVersion) { KeyLocation.Builder builder = KeyLocation.newBuilder() .setBlockID(getBlockID().getProtobuf()) .setLength(getLength()) diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyLocationInfoGroup.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyLocationInfoGroup.java index e2477a4cef10..f54c4108c2c6 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyLocationInfoGroup.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmKeyLocationInfoGroup.java @@ -24,6 +24,7 @@ import java.util.List; import java.util.Map; import java.util.stream.Collectors; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.KeyLocationList; @@ -125,7 +126,7 @@ public List getLocationList(Long versionToFetch) { } public KeyLocationList getProtobuf(boolean ignorePipeline, - int clientVersion) { + ClientVersion clientVersion) { KeyLocationList.Builder builder = KeyLocationList.newBuilder() .setVersion(version).setIsMultipartKey(isMultipartKey); List keyLocationList = diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartPartInfo.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartPartInfo.java index 41bf349f504b..3b750c7e3e96 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartPartInfo.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OmMultipartPartInfo.java @@ -338,7 +338,7 @@ private KeyLocationList getKeyLocationInfosAsProto() { if (keyLocationInfos == null || keyLocationInfos.isEmpty()) { throw new IllegalArgumentException("keyLocationList is required"); } - return keyLocationInfos.get(0).getProtobuf(true, ClientVersion.CURRENT.serialize()); + return keyLocationInfos.get(0).getProtobuf(true, ClientVersion.CURRENT); } private static List getKeyLocationInfosFromProto( diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OzoneFileStatus.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OzoneFileStatus.java index e0d7ebe37b63..8f17d986b26b 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OzoneFileStatus.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/OzoneFileStatus.java @@ -21,6 +21,7 @@ import java.io.IOException; import java.util.Objects; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OzoneFileStatusProto; import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.OzoneFileStatusProto.Builder; @@ -92,7 +93,7 @@ public boolean isFile() { return !isDirectory(); } - public OzoneFileStatusProto getProtobuf(int clientVersion) { + public OzoneFileStatusProto getProtobuf(ClientVersion clientVersion) { Builder builder = OzoneFileStatusProto.newBuilder() .setBlockSize(blockSize) .setIsDirectory(isDirectory); diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/RepeatedOmKeyInfo.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/RepeatedOmKeyInfo.java index 39a3497aa1d7..fe55033b36b2 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/RepeatedOmKeyInfo.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/helpers/RepeatedOmKeyInfo.java @@ -60,7 +60,7 @@ private static Codec newCodec(boolean ignorePipeline, boolean return new DelegatedCodec<>( Proto2Codec.get(RepeatedKeyInfo.getDefaultInstance()), RepeatedOmKeyInfo::getFromProto, - k -> k.getProto(ignorePipeline, ClientVersion.CURRENT.serialize(), isOpenKey), + k -> k.getProto(ignorePipeline, ClientVersion.CURRENT, isOpenKey), RepeatedOmKeyInfo.class); } @@ -151,7 +151,7 @@ public static RepeatedOmKeyInfo getFromProto(RepeatedKeyInfo repeatedKeyInfo) { /** * @param compact true for persistence, false for network transmit */ - public RepeatedKeyInfo getProto(boolean compact, int clientVersion) { + public RepeatedKeyInfo getProto(boolean compact, ClientVersion clientVersion) { return getProto(compact, clientVersion, true); } @@ -160,7 +160,7 @@ public RepeatedKeyInfo getProto(boolean compact, int clientVersion) { * @param clientVersion the client version * @param isOpenKey true for openKeyTable, false for keyTable/deletedTable */ - public RepeatedKeyInfo getProto(boolean compact, int clientVersion, boolean isOpenKey) { + public RepeatedKeyInfo getProto(boolean compact, ClientVersion clientVersion, boolean isOpenKey) { List list = new ArrayList<>(); for (OmKeyInfo k : cloneOmKeyInfoList()) { list.add(k.getProtobuf(compact, clientVersion, isOpenKey)); diff --git a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java index c17af89679d5..92d031342559 100644 --- a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java +++ b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java @@ -855,7 +855,7 @@ private void updateKey(OmKeyArgs args, long clientId, boolean hsync, boolean rec .addAllMetadata(KeyValueUtil.toProtobuf(args.getMetadata())) .addAllKeyLocations(locationInfoList.stream() // TODO use OM version? - .map(info -> info.getProtobuf(ClientVersion.CURRENT.serialize())) + .map(info -> info.getProtobuf(ClientVersion.CURRENT)) .collect(Collectors.toList())); setReplicationConfig(args.getReplicationConfig(), keyArgsBuilder); @@ -1774,7 +1774,7 @@ public OmMultipartCommitUploadPartInfo commitMultipartUploadPart( .addAllMetadata(KeyValueUtil.toProtobuf(omKeyArgs.getMetadata())) .addAllKeyLocations(locationInfoList.stream() // TODO use OM version? - .map(info -> info.getProtobuf(ClientVersion.CURRENT.serialize())) + .map(info -> info.getProtobuf(ClientVersion.CURRENT)) .collect(Collectors.toList())); multipartCommitUploadPartRequest.setClientID(clientId); multipartCommitUploadPartRequest.setKeyArgs(keyArgs.build()); diff --git a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/helpers/TestOmKeyInfo.java b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/helpers/TestOmKeyInfo.java index bc8928f9c4c1..30e69b13d685 100644 --- a/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/helpers/TestOmKeyInfo.java +++ b/hadoop-ozone/common/src/test/java/org/apache/hadoop/ozone/om/helpers/TestOmKeyInfo.java @@ -60,7 +60,7 @@ public void protobufConversion() throws IOException { RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); OmKeyInfo keyAfterSerialization = OmKeyInfo.getFromProtobuf( - key.getProtobuf(ClientVersion.CURRENT.serialize())); + key.getProtobuf(ClientVersion.CURRENT)); assertNotNull(keyAfterSerialization); assertEquals(key, keyAfterSerialization); @@ -78,7 +78,7 @@ public void getProtobufMessageEC() throws IOException { OmKeyInfo key = createOmKeyInfo( RatisReplicationConfig.getInstance(ReplicationFactor.THREE)); OzoneManagerProtocolProtos.KeyInfo omKeyProto = - key.getProtobuf(ClientVersion.CURRENT.serialize()); + key.getProtobuf(ClientVersion.CURRENT); // No EC Config assertFalse(omKeyProto.hasEcReplicationConfig()); @@ -95,7 +95,7 @@ public void getProtobufMessageEC() throws IOException { // EC Config key = createOmKeyInfo(new ECReplicationConfig(3, 2)); assertFalse(key.isHsync()); - omKeyProto = key.getProtobuf(ClientVersion.CURRENT.serialize()); + omKeyProto = key.getProtobuf(ClientVersion.CURRENT); assertEquals(3, omKeyProto.getEcReplicationConfig().getData()); diff --git a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/AbstractTestStorageDistributionEndpoint.java b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/AbstractTestStorageDistributionEndpoint.java index ee63904be59e..0fd5234ba918 100644 --- a/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/AbstractTestStorageDistributionEndpoint.java +++ b/hadoop-ozone/integration-test-recon/src/test/java/org/apache/hadoop/ozone/recon/AbstractTestStorageDistributionEndpoint.java @@ -49,6 +49,7 @@ import org.apache.hadoop.hdds.scm.events.SCMEvents; import org.apache.hadoop.hdds.scm.server.StorageContainerManager; import org.apache.hadoop.hdds.utils.IOUtils; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.MiniOzoneCluster; import org.apache.hadoop.ozone.client.BucketArgs; import org.apache.hadoop.ozone.client.ObjectStore; @@ -229,7 +230,8 @@ protected boolean verifyStorageDistributionAfterKeyCreation() { List reports = storageResponse.getDataNodeUsage(); List scmReports = - scm.getClientProtocolServer().getDatanodeUsageInfo(true, getNumDatanodes(), 1); + scm.getClientProtocolServer().getDatanodeUsageInfo(true, getNumDatanodes(), + ClientVersion.VERSION_HANDLES_UNKNOWN_DN_PORTS); long totalReserved = 0; long totalMinFreeSpace = 0; diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestScmApplyTransactionFailure.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestScmApplyTransactionFailure.java index 687201e0a148..4391b622b57c 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestScmApplyTransactionFailure.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/scm/container/TestScmApplyTransactionFailure.java @@ -94,7 +94,7 @@ public void testAddDuplicatePipelineId() replication, PipelineState.OPEN).get(0); HddsProtos.Pipeline pipelineToCreate = - existing.getProtobufMessage(CURRENT.serialize()); + existing.getProtobufMessage(CURRENT); Throwable ex = assertThrows(SCMException.class, () -> pipelineManager.getStateManager().addPipeline( pipelineToCreate)); diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/debug/TestLDBCli.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/debug/TestLDBCli.java index 94b160e5c8c7..f6cfc6759d7f 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/debug/TestLDBCli.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/debug/TestLDBCli.java @@ -500,7 +500,7 @@ private void prepareKeyTable(int recordsCount) throws IOException { OmKeyInfo value = OMRequestTestUtils.createOmKeyInfo("vol1", "buck1", key, ReplicationConfig.fromProtoTypeAndFactor(STAND_ALONE, HddsProtos.ReplicationFactor.ONE)).build(); - keyTable.put(key.getBytes(UTF_8), value.getProtobuf(ClientVersion.CURRENT.serialize()).toByteArray()); + keyTable.put(key.getBytes(UTF_8), value.getProtobuf(ClientVersion.CURRENT).toByteArray()); // Populate map dbMap.put(key, toMap(value)); } diff --git a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/checksum/TestContainerCommandReconciliation.java b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/checksum/TestContainerCommandReconciliation.java index 3868c691862e..3c8656ab7686 100644 --- a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/checksum/TestContainerCommandReconciliation.java +++ b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/dn/checksum/TestContainerCommandReconciliation.java @@ -511,7 +511,7 @@ public void testDataChecksumReportedAtSCM() throws Exception { // Check non-zero checksum after container close StorageContainerLocationProtocolClientSideTranslatorPB scmClient = cluster.getStorageContainerLocationClient(); List containerReplicas = scmClient.getContainerReplicas(containerID, - ClientVersion.CURRENT.serialize()); + ClientVersion.CURRENT); assertEquals(3, containerReplicas.size()); for (HddsProtos.SCMContainerReplicaProto containerReplica: containerReplicas) { assertNotEquals(0, containerReplica.getDataChecksum()); @@ -545,7 +545,7 @@ public void testDataChecksumReportedAtSCM() throws Exception { scmClient.reconcileContainer(containerID); waitForDataChecksumsAtSCM(containerID, 1); // Check non-zero checksum after container reconciliation - containerReplicas = scmClient.getContainerReplicas(containerID, ClientVersion.CURRENT.serialize()); + containerReplicas = scmClient.getContainerReplicas(containerID, ClientVersion.CURRENT); assertEquals(3, containerReplicas.size()); for (HddsProtos.SCMContainerReplicaProto containerReplica: containerReplicas) { assertNotEquals(0, containerReplica.getDataChecksum()); @@ -559,7 +559,7 @@ public void testDataChecksumReportedAtSCM() throws Exception { } cluster.waitForClusterToBeReady(); waitForDataChecksumsAtSCM(containerID, 1); - containerReplicas = scmClient.getContainerReplicas(containerID, ClientVersion.CURRENT.serialize()); + containerReplicas = scmClient.getContainerReplicas(containerID, ClientVersion.CURRENT); assertEquals(3, containerReplicas.size()); for (HddsProtos.SCMContainerReplicaProto containerReplica: containerReplicas) { assertNotEquals(0, containerReplica.getDataChecksum()); @@ -571,7 +571,7 @@ private void waitForDataChecksumsAtSCM(long containerID, int expectedSize) throw GenericTestUtils.waitFor(() -> { try { Set dataChecksums = cluster.getStorageContainerLocationClient().getContainerReplicas(containerID, - ClientVersion.CURRENT.serialize()).stream() + ClientVersion.CURRENT).stream() .map(HddsProtos.SCMContainerReplicaProto::getDataChecksum) .collect(Collectors.toSet()); LOG.info("Waiting for {} total unique checksums from container {} to be reported to SCM. Currently {} unique" + diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OmMetadataManagerImpl.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OmMetadataManagerImpl.java index 89960027839a..e230b84b8b18 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OmMetadataManagerImpl.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OmMetadataManagerImpl.java @@ -1501,7 +1501,7 @@ public ExpiredOpenKeys getExpiredOpenKeys(Duration expireThreshold, .map(OmKeyLocationInfoGroup::getLocationList) .map(Collection::stream) .orElseGet(Stream::empty) - .map(loc -> loc.getProtobuf(ClientVersion.CURRENT.serialize())) + .map(loc -> loc.getProtobuf(ClientVersion.CURRENT)) .forEach(keyArgs::addKeyLocations); OzoneManagerProtocolClientSideTranslatorPB.setReplicationConfig( diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequest.java index 35e1ac238f7b..0f72da7c58fe 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequest.java @@ -35,6 +35,7 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; import org.apache.hadoop.hdds.utils.UniqueId; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OmUtils; import org.apache.hadoop.ozone.audit.OMAction; import org.apache.hadoop.ozone.om.OMMetadataManager; @@ -135,7 +136,8 @@ public OMRequest preExecute(OzoneManager ozoneManager) throws IOException { .setDataSize(requestedSize); newKeyArgs.addAllKeyLocations(omKeyLocationInfoList.stream() - .map(info -> info.getProtobuf(getOmRequest().getVersion())) + .map(info -> info.getProtobuf( + ClientVersion.deserialize(getOmRequest().getVersion()))) .collect(Collectors.toList())); generateRequiredEncryptionInfo(keyArgs, newKeyArgs, ozoneManager); @@ -279,7 +281,8 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut // Prepare response omResponse.setCreateFileResponse(CreateFileResponse.newBuilder() - .setKeyInfo(omKeyInfo.getNetworkProtobuf(getOmRequest().getVersion(), + .setKeyInfo(omKeyInfo.getNetworkProtobuf( + ClientVersion.deserialize(getOmRequest().getVersion()), keyArgs.getLatestVersionLocation())) .setID(clientID) .setOpenVersion(openVersion).build()) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequestWithFSO.java index 6036fe90dbb7..d41e11bca73a 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMFileCreateRequestWithFSO.java @@ -27,6 +27,7 @@ import java.util.Map; import java.util.stream.Collectors; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.audit.OMAction; import org.apache.hadoop.ozone.om.OMMetadataManager; import org.apache.hadoop.ozone.om.OMMetrics; @@ -207,7 +208,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut // Prepare response. Sets user given full key name in the 'keyName' // attribute in response object. - int clientVersion = getOmRequest().getVersion(); + ClientVersion clientVersion = ClientVersion.deserialize(getOmRequest().getVersion()); omResponse.setCreateFileResponse(CreateFileResponse.newBuilder() .setKeyInfo(omFileInfo.getNetworkProtobuf(keyName, clientVersion, keyArgs.getLatestVersionLocation())) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMRecoverLeaseRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMRecoverLeaseRequest.java index ca1ea07ad6ed..674493b6a426 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMRecoverLeaseRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/file/OMRecoverLeaseRequest.java @@ -38,6 +38,7 @@ import java.util.concurrent.TimeUnit; import org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline; import org.apache.hadoop.hdds.security.token.OzoneBlockTokenSecretManager; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.audit.OMAction; import org.apache.hadoop.ozone.om.OMMetadataManager; @@ -257,8 +258,10 @@ private RecoverLeaseResponse doWork(OzoneManager ozoneManager, } RecoverLeaseResponse.Builder rb = RecoverLeaseResponse.newBuilder(); - rb.setKeyInfo(keyInfo.getNetworkProtobuf(getOmRequest().getVersion(), true)); - rb.setOpenKeyInfo(openKeyInfo.getNetworkProtobuf(getOmRequest().getVersion(), true)); + rb.setKeyInfo(keyInfo.getNetworkProtobuf( + ClientVersion.deserialize(getOmRequest().getVersion()), true)); + rb.setOpenKeyInfo(openKeyInfo.getNetworkProtobuf( + ClientVersion.deserialize(getOmRequest().getVersion()), true)); return rb.build(); } diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMAllocateBlockRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMAllocateBlockRequest.java index 0e11f1d76773..7da2cd53a240 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMAllocateBlockRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMAllocateBlockRequest.java @@ -32,6 +32,7 @@ import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.hdds.utils.db.cache.CacheValue; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.audit.AuditLogger; import org.apache.hadoop.ozone.audit.OMAction; @@ -131,7 +132,8 @@ public OMRequest preExecute(OzoneManager ozoneManager) throws IOException { // Add allocated block info. newAllocatedBlockRequest.setKeyLocation( - omKeyLocationInfoList.get(0).getProtobuf(getOmRequest().getVersion())); + omKeyLocationInfoList.get(0).getProtobuf( + ClientVersion.deserialize(getOmRequest().getVersion()))); return getOmRequest().toBuilder().setUserInfo(userInfo) .setAllocateBlockRequest(newAllocatedBlockRequest).build(); diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequest.java index 929e46222c05..6303f2cbc43e 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequest.java @@ -35,6 +35,7 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList; import org.apache.hadoop.hdds.utils.UniqueId; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OmUtils; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.OzoneManagerVersion; @@ -163,7 +164,7 @@ public OMRequest preExecute(OzoneManager ozoneManager) throws IOException { newKeyArgs.addAllKeyLocations(omKeyLocationInfoList.stream() .map(info -> info.getProtobuf(false, - getOmRequest().getVersion())) + ClientVersion.deserialize(getOmRequest().getVersion()))) .collect(Collectors.toList())); } else { newKeyArgs = keyArgs.toBuilder().setModificationTime(Time.now()); @@ -336,7 +337,8 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut // Prepare response omResponse.setCreateKeyResponse(CreateKeyResponse.newBuilder() - .setKeyInfo(omKeyInfo.getNetworkProtobuf(getOmRequest().getVersion(), + .setKeyInfo(omKeyInfo.getNetworkProtobuf( + ClientVersion.deserialize(getOmRequest().getVersion()), keyArgs.getLatestVersionLocation())) .setID(clientID) .setOpenVersion(openVersion).build()) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequestWithFSO.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequestWithFSO.java index 99fabb46de11..122aaad8e998 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequestWithFSO.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/key/OMKeyCreateRequestWithFSO.java @@ -30,6 +30,7 @@ import java.util.Map; import java.util.stream.Collectors; import org.apache.hadoop.hdds.client.ReplicationConfig; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.audit.OMAction; import org.apache.hadoop.ozone.om.OMMetadataManager; import org.apache.hadoop.ozone.om.OMMetrics; @@ -201,7 +202,7 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut // Prepare response. Sets user given full key name in the 'keyName' // attribute in response object. - int clientVersion = getOmRequest().getVersion(); + ClientVersion clientVersion = ClientVersion.deserialize(getOmRequest().getVersion()); omResponse.setCreateKeyResponse(CreateKeyResponse.newBuilder() .setKeyInfo(omFileInfo.getNetworkProtobuf(keyName, clientVersion, keyArgs.getLatestVersionLocation())) diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCommitPartRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCommitPartRequest.java index 78f1af96cfbd..a2bc6e86a39d 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCommitPartRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCommitPartRequest.java @@ -30,6 +30,7 @@ import org.apache.commons.lang3.StringUtils; import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.hdds.utils.db.cache.CacheValue; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.audit.OMAction; import org.apache.hadoop.ozone.om.OMMetadataManager; @@ -217,7 +218,8 @@ public OMClientResponse validateAndUpdateCache(OzoneManager ozoneManager, Execut OzoneManagerProtocolProtos.PartKeyInfo.newBuilder(); partKeyInfo.setPartName(partName); partKeyInfo.setPartNumber(partNumber); - partKeyInfo.setPartKeyInfo(omKeyInfo.getProtobuf(getOmRequest().getVersion())); + partKeyInfo.setPartKeyInfo(omKeyInfo.getProtobuf( + ClientVersion.deserialize(getOmRequest().getVersion()))); if (multipartKeyInfo.getSchemaVersion() == OmMultipartKeyInfo.LEGACY_SCHEMA_VERSION) { // Add this part information in to multipartKeyInfo. diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java index 841ced7dacce..0da53c46bc1f 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/s3/multipart/S3MultipartUploadCompleteRequest.java @@ -37,6 +37,7 @@ import org.apache.hadoop.hdds.client.ReplicationConfig; import org.apache.hadoop.hdds.utils.db.cache.CacheKey; import org.apache.hadoop.hdds.utils.db.cache.CacheValue; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.audit.OMAction; import org.apache.hadoop.ozone.om.OMMetadataManager; @@ -605,7 +606,8 @@ private OmMultipartKeyInfo.PartKeyInfoMap getPartKeyInfoMap( partKeyInfos.put(entry.getKey(), PartKeyInfo.newBuilder() .setPartName(partInfo.getPartName()) .setPartNumber(partInfo.getPartNumber()) - .setPartKeyInfo(partKeyInfo.getProtobuf(getOmRequest().getVersion())) + .setPartKeyInfo(partKeyInfo.getProtobuf( + ClientVersion.deserialize(getOmRequest().getVersion()))) .build()); } return new OmMultipartKeyInfo.PartKeyInfoMap(partKeyInfos); diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java index 8069f5faf0c0..11b2ba673509 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/DirectoryDeletingService.java @@ -521,14 +521,14 @@ private OzoneManagerProtocolProtos.PurgePathRequest wrapPurgeRequest( for (OmKeyInfo purgeFile : purgeDeletedFiles) { purgePathsRequest.addDeletedSubFiles( - purgeFile.getProtobuf(true, ClientVersion.CURRENT.serialize())); + purgeFile.getProtobuf(true, ClientVersion.CURRENT)); } // Add these directories to deletedDirTable, so that its sub-paths will be // traversed in next iteration to ensure cleanup all sub-children. for (OmKeyInfo dir : markDirsAsDeleted) { purgePathsRequest.addMarkDeletedSubDirs( - dir.getProtobuf(ClientVersion.CURRENT.serialize())); + dir.getProtobuf(ClientVersion.CURRENT)); } return purgePathsRequest.build(); diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/KeyDeletingService.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/KeyDeletingService.java index e37b1406c492..77701ab4260f 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/KeyDeletingService.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/KeyDeletingService.java @@ -333,7 +333,7 @@ private Pair, Boolean> submitPurgeKeysRequest( keyToUpdate.setKey(keyToModify.getKey()); List keyInfos = keyToModify.getValue().getOmKeyInfoList().stream() - .map(k -> k.getProtobuf(ClientVersion.CURRENT.serialize())) + .map(k -> k.getProtobuf(ClientVersion.CURRENT)) .collect(Collectors.toList()); keyToUpdate.addAllKeyInfos(keyInfos); keyToUpdate.setBucketId(keyToModify.getValue().getBucketId()); diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/SnapshotDeletingService.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/SnapshotDeletingService.java index 481a7ce88f92..c1508a89d838 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/SnapshotDeletingService.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/service/SnapshotDeletingService.java @@ -196,14 +196,14 @@ public BackgroundTaskResult call() throws InterruptedException { for (Table.KeyValue> deletedEntry : deletedKeyEntries) { deletedKeys.add(SnapshotMoveKeyInfos.newBuilder().setKey(deletedEntry.getKey()) .addAllKeyInfos(deletedEntry.getValue() - .stream().map(val -> val.getProtobuf(ClientVersion.CURRENT.serialize())) + .stream().map(val -> val.getProtobuf(ClientVersion.CURRENT)) .collect(Collectors.toList())).build()); } // Convert deletedDirEntries to SnapshotMoveKeyInfos. for (Table.KeyValue deletedDirEntry : deletedDirEntries) { deletedDirs.add(SnapshotMoveKeyInfos.newBuilder().setKey(deletedDirEntry.getKey()) - .addKeyInfos(deletedDirEntry.getValue().getProtobuf(ClientVersion.CURRENT.serialize())).build()); + .addKeyInfos(deletedDirEntry.getValue().getProtobuf(ClientVersion.CURRENT)).build()); } // Convert renamedEntries to KeyValue. diff --git a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java index e9670614677d..9b90eb81a5df 100644 --- a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java +++ b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OzoneManagerRequestHandler.java @@ -60,6 +60,7 @@ import org.apache.hadoop.hdds.protocol.proto.HddsProtos.UpgradeFinalizationStatus; import org.apache.hadoop.hdds.scm.protocolPB.OzonePBHelper; import org.apache.hadoop.hdds.utils.FaultInjector; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OzoneAcl; import org.apache.hadoop.ozone.om.OzoneManager; import org.apache.hadoop.ozone.om.exceptions.OMException; @@ -200,6 +201,7 @@ public OMResponse handleReadRequest(OMRequest request) { Type cmdType = request.getCmdType(); OMResponse.Builder responseBuilder = OmResponseUtil.getOMResponseBuilder( request); + final ClientVersion clientVersion = ClientVersion.deserialize(request.getVersion()); try { switch (cmdType) { case CheckVolumeAccess: @@ -229,12 +231,12 @@ public OMResponse handleReadRequest(OMRequest request) { break; case LookupKey: LookupKeyResponse lookupKeyResponse = lookupKey( - request.getLookupKeyRequest(), request.getVersion()); + request.getLookupKeyRequest(), clientVersion); responseBuilder.setLookupKeyResponse(lookupKeyResponse); break; case ListKeys: ListKeysResponse listKeysResponse = listKeys( - request.getListKeysRequest(), request.getVersion()); + request.getListKeysRequest(), clientVersion); responseBuilder.setListKeysResponse(listKeysResponse); break; case ListKeysLight: @@ -254,7 +256,7 @@ public OMResponse handleReadRequest(OMRequest request) { break; case ListOpenFiles: ListOpenFilesResponse listOpenFilesResponse = listOpenFiles( - request.getListOpenFilesRequest(), request.getVersion()); + request.getListOpenFilesRequest(), clientVersion); responseBuilder.setListOpenFilesResponse(listOpenFilesResponse); break; case ServiceList: @@ -274,23 +276,22 @@ public OMResponse handleReadRequest(OMRequest request) { break; case GetFileStatus: GetFileStatusResponse getFileStatusResponse = getOzoneFileStatus( - request.getGetFileStatusRequest(), request.getVersion()); + request.getGetFileStatusRequest(), clientVersion); responseBuilder.setGetFileStatusResponse(getFileStatusResponse); break; case LookupFile: LookupFileResponse lookupFileResponse = - lookupFile(request.getLookupFileRequest(), request.getVersion()); + lookupFile(request.getLookupFileRequest(), clientVersion); responseBuilder.setLookupFileResponse(lookupFileResponse); break; case ListStatus: ListStatusResponse listStatusResponse = - listStatus(request.getListStatusRequest(), request.getVersion()); + listStatus(request.getListStatusRequest(), clientVersion); responseBuilder.setListStatusResponse(listStatusResponse); break; case ListStatusLight: ListStatusLightResponse listStatusLightResponse = - listStatusLight(request.getListStatusRequest(), - request.getVersion()); + listStatusLight(request.getListStatusRequest()); responseBuilder.setListStatusLightResponse(listStatusLightResponse); break; case GetAcl: @@ -343,7 +344,7 @@ public OMResponse handleReadRequest(OMRequest request) { break; case GetKeyInfo: responseBuilder.setGetKeyInfoResponse( - getKeyInfo(request.getGetKeyInfoRequest(), request.getVersion())); + getKeyInfo(request.getGetKeyInfoRequest(), clientVersion)); break; case ListSnapshot: OzoneManagerProtocolProtos.ListSnapshotResponse listSnapshotResponse = @@ -649,7 +650,7 @@ private InfoBucketResponse infoBucket(InfoBucketRequest request) } private LookupKeyResponse lookupKey(LookupKeyRequest request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { LookupKeyResponse.Builder resp = LookupKeyResponse.newBuilder(); KeyArgs keyArgs = request.getKeyArgs(); @@ -669,7 +670,7 @@ private LookupKeyResponse lookupKey(LookupKeyRequest request, } private GetKeyInfoResponse getKeyInfo(GetKeyInfoRequest request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { KeyArgs keyArgs = request.getKeyArgs(); OmKeyArgs omKeyArgs = new OmKeyArgs.Builder() .setVolumeName(keyArgs.getVolumeName()) @@ -766,7 +767,7 @@ private ListBucketsResponse listBuckets(ListBucketsRequest request) return resp.build(); } - private ListKeysResponse listKeys(ListKeysRequest request, int clientVersion) + private ListKeysResponse listKeys(ListKeysRequest request, ClientVersion clientVersion) throws IOException { ListKeysResponse.Builder resp = ListKeysResponse.newBuilder(); @@ -947,7 +948,7 @@ public static OMResponse disallowListTrashWithBucketLayout( @DisallowedUntilLayoutVersion(HBASE_SUPPORT) private ListOpenFilesResponse listOpenFiles(ListOpenFilesRequest req, - int clientVersion) + ClientVersion clientVersion) throws IOException { ListOpenFilesResponse.Builder resp = ListOpenFilesResponse.newBuilder(); @@ -1086,7 +1087,7 @@ private ListMultipartUploadsResponse listMultipartUploads( } private GetFileStatusResponse getOzoneFileStatus( - GetFileStatusRequest request, int clientVersion) throws IOException { + GetFileStatusRequest request, ClientVersion clientVersion) throws IOException { KeyArgs keyArgs = request.getKeyArgs(); OmKeyArgs omKeyArgs = new OmKeyArgs.Builder() .setVolumeName(keyArgs.getVolumeName()) @@ -1191,7 +1192,7 @@ public static OMResponse disallowGetFileStatusWithBucketLayout( } private LookupFileResponse lookupFile(LookupFileRequest request, - int clientVersion) throws IOException { + ClientVersion clientVersion) throws IOException { KeyArgs keyArgs = request.getKeyArgs(); OmKeyArgs omKeyArgs = new OmKeyArgs.Builder() .setVolumeName(keyArgs.getVolumeName()) @@ -1265,7 +1266,7 @@ public static OMResponse disallowLookupFileWithBucketLayout( } private ListStatusResponse listStatus( - ListStatusRequest request, int clientVersion) throws IOException { + ListStatusRequest request, ClientVersion clientVersion) throws IOException { KeyArgs keyArgs = request.getKeyArgs(); OmKeyArgs omKeyArgs = new OmKeyArgs.Builder() .setVolumeName(keyArgs.getVolumeName()) @@ -1289,8 +1290,7 @@ private ListStatusResponse listStatus( return listStatusResponseBuilder.build(); } - private ListStatusLightResponse listStatusLight( - ListStatusRequest request, int clientVersion) throws IOException { + private ListStatusLightResponse listStatusLight(ListStatusRequest request) throws IOException { KeyArgs keyArgs = request.getKeyArgs(); OmKeyArgs omKeyArgs = new OmKeyArgs.Builder() .setVolumeName(keyArgs.getVolumeName()) diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/OMRequestTestUtils.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/OMRequestTestUtils.java index 16f42cb9179f..e8715d63f619 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/OMRequestTestUtils.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/OMRequestTestUtils.java @@ -1381,7 +1381,7 @@ public static OMRequest moveSnapshotTableKeyRequest(UUID snapshotId, .setKey(deletedKey.getKey()) .addAllKeyInfos( deletedKey.getValue().stream() - .map(omKeyInfo -> omKeyInfo.getProtobuf(ClientVersion.CURRENT.serialize())) + .map(omKeyInfo -> omKeyInfo.getProtobuf(ClientVersion.CURRENT)) .collect(Collectors.toList())) .build(); deletedMoveKeys.add(snapshotMoveKeyInfos); @@ -1394,7 +1394,7 @@ public static OMRequest moveSnapshotTableKeyRequest(UUID snapshotId, .setKey(deletedKey.getKey()) .addAllKeyInfos( deletedKey.getValue().stream() - .map(omKeyInfo -> omKeyInfo.getProtobuf(ClientVersion.CURRENT.serialize())) + .map(omKeyInfo -> omKeyInfo.getProtobuf(ClientVersion.CURRENT)) .collect(Collectors.toList())) .build(); deletedDirMoveKeys.add(snapshotMoveKeyInfos); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/file/TestOMRecoverLeaseRequest.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/file/TestOMRecoverLeaseRequest.java index 7ded4c7f3601..7755f9f9d403 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/file/TestOMRecoverLeaseRequest.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/file/TestOMRecoverLeaseRequest.java @@ -377,7 +377,7 @@ private KeyArgs getNewKeyArgs(OmKeyInfo omKeyInfo, long deltaLength) throws IOEx .setDataSize(keyArgs.getDataSize()) .addAllMetadata(KeyValueUtil.toProtobuf(keyArgs.getMetadata())) .addAllKeyLocations(locationInfoList.stream() - .map(info -> info.getProtobuf(ClientVersion.CURRENT.serialize())) + .map(info -> info.getProtobuf(ClientVersion.CURRENT)) .collect(Collectors.toList())); setReplicationConfig(keyArgs.getReplicationConfig(), keyArgsBuilder); return keyArgsBuilder.build(); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMDirectoriesPurgeRequestAndResponse.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMDirectoriesPurgeRequestAndResponse.java index 2b0889ec5798..aea9e41e025a 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMDirectoriesPurgeRequestAndResponse.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/key/TestOMDirectoriesPurgeRequestAndResponse.java @@ -207,14 +207,14 @@ private PurgePathRequest wrapPurgeRequest( for (OmKeyInfo purgeFile : purgeDeletedFiles) { purgePathsRequest.addDeletedSubFiles( - purgeFile.getProtobuf(true, ClientVersion.CURRENT.serialize())); + purgeFile.getProtobuf(true, ClientVersion.CURRENT)); } // Add these directories to deletedDirTable, so that its sub-paths will be // traversed in next iteration to ensure cleanup all sub-children. for (OmKeyInfo dir : markDirsAsDeleted) { purgePathsRequest.addMarkDeletedSubDirs( - dir.getProtobuf(ClientVersion.CURRENT.serialize())); + dir.getProtobuf(ClientVersion.CURRENT)); } return purgePathsRequest.build(); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/response/snapshot/TestOMSnapshotMoveTableKeysResponse.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/response/snapshot/TestOMSnapshotMoveTableKeysResponse.java index 1df73f258612..031fbdd2b107 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/response/snapshot/TestOMSnapshotMoveTableKeysResponse.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/response/snapshot/TestOMSnapshotMoveTableKeysResponse.java @@ -147,13 +147,13 @@ public void testMoveTableKeysToNextSnapshot(boolean nextSnapshotExists) throws E .forEachRemaining(entry -> { deletedTable.add(OzoneManagerProtocolProtos.SnapshotMoveKeyInfos.newBuilder().setKey(entry.getKey()) .addAllKeyInfos(entry.getValue().getOmKeyInfoList().stream().map(omKeyInfo -> omKeyInfo.getProtobuf( - ClientVersion.CURRENT.serialize())).collect(Collectors.toList())).build()); + ClientVersion.CURRENT)).collect(Collectors.toList())).build()); }); snapshot.getMetadataManager().getDeletedDirTable().iterator() .forEachRemaining(entry -> { deletedDirTable.add(OzoneManagerProtocolProtos.SnapshotMoveKeyInfos.newBuilder().setKey(entry.getKey()) - .addKeyInfos(entry.getValue().getProtobuf(ClientVersion.CURRENT.serialize())).build()); + .addKeyInfos(entry.getValue().getProtobuf(ClientVersion.CURRENT)).build()); }); snapshot.getMetadataManager().getSnapshotRenamedTable().iterator().forEachRemaining(entry -> { renamedTable.add(HddsProtos.KeyValue.newBuilder().setKey(entry.getKey()).setValue(entry.getValue()).build()); diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestSnapshotDeletingService.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestSnapshotDeletingService.java index a6af3c0a70a9..7a61087c8acf 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestSnapshotDeletingService.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/service/TestSnapshotDeletingService.java @@ -336,7 +336,7 @@ private List createLargeDeletedKeys(int count) { SnapshotMoveKeyInfos moveKeyInfo = SnapshotMoveKeyInfos.newBuilder() .setKey(largeKeyName) .addAllKeyInfos(keyInfos.stream() - .map(k -> k.getProtobuf(ClientVersion.CURRENT.serialize())) + .map(k -> k.getProtobuf(ClientVersion.CURRENT)) .collect(Collectors.toList())) .build(); deletedKeys.add(moveKeyInfo); @@ -371,7 +371,7 @@ private List createLargeDeletedDirs(int count) { SnapshotMoveKeyInfos moveDirInfo = SnapshotMoveKeyInfos.newBuilder() .setKey(largeDirName) - .addKeyInfos(dirInfo.getProtobuf(ClientVersion.CURRENT.serialize())) + .addKeyInfos(dirInfo.getProtobuf(ClientVersion.CURRENT)) .build(); deletedDirs.add(moveDirInfo); } diff --git a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java index 0601713b7956..9531942e9c16 100644 --- a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java +++ b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOzoneManagerRequestHandler.java @@ -33,6 +33,7 @@ import org.apache.hadoop.hdds.conf.OzoneConfiguration; import org.apache.hadoop.hdds.protocol.proto.HddsProtos; import org.apache.hadoop.metrics2.lib.MutableRate; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.audit.AuditLogger; import org.apache.hadoop.ozone.audit.AuditMessage; import org.apache.hadoop.ozone.om.OMPerformanceMetrics; @@ -77,8 +78,8 @@ private OmKeyInfo getMockedOmKeyInfo() { OzoneManagerProtocolProtos.KeyInfo.newBuilder().setBucketName("bucket").setKeyName("key").setVolumeName( "volume").setDataSize(0).setType(HddsProtos.ReplicationType.RATIS).setCreationTime(0) .setModificationTime(0).build(); - Mockito.when(keyInfo.getProtobuf(Mockito.anyBoolean(), Mockito.anyInt())).thenReturn(info); - Mockito.when(keyInfo.getProtobuf(Mockito.anyInt())).thenReturn(info); + Mockito.when(keyInfo.getProtobuf(Mockito.anyBoolean(), any(ClientVersion.class))).thenReturn(info); + Mockito.when(keyInfo.getProtobuf(any(ClientVersion.class))).thenReturn(info); return keyInfo; } @@ -176,7 +177,7 @@ public void getFileStatusForwardsHeadOpAndStripsLocations() throws IOException { .setVersion(0).build()) .build()) .build(); - Mockito.when(status.getProtobuf(Mockito.anyInt())).thenReturn(proto); + Mockito.when(status.getProtobuf(any(ClientVersion.class))).thenReturn(proto); ArgumentCaptor captor = ArgumentCaptor.forClass(OmKeyArgs.class); Mockito.when(ozoneManager.getFileStatus(captor.capture())).thenReturn(status); @@ -212,7 +213,7 @@ public void getFileStatusHeadOpWithoutKeyInfoIsNoop() throws IOException { OzoneManager ozoneManager = requestHandler.getOzoneManager(); OzoneFileStatus status = Mockito.mock(OzoneFileStatus.class); - Mockito.when(status.getProtobuf(Mockito.anyInt())).thenReturn( + Mockito.when(status.getProtobuf(any(ClientVersion.class))).thenReturn( OzoneManagerProtocolProtos.OzoneFileStatusProto.newBuilder() .setIsDirectory(true).build()); Mockito.when(ozoneManager.getFileStatus(Mockito.any())).thenReturn(status); diff --git a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/NodeEndpoint.java b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/NodeEndpoint.java index 201e0541ad96..35397723d258 100644 --- a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/NodeEndpoint.java +++ b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/api/NodeEndpoint.java @@ -357,7 +357,7 @@ private Response getDecommissionStatusResponse(String uuid, String ipAddress) th Response.ResponseBuilder builder = Response.status(Response.Status.OK); Map responseMap = new HashMap<>(); Stream allNodes = scmClient.queryNode(DECOMMISSIONING, - null, HddsProtos.QueryScope.CLUSTER, "", ClientVersion.CURRENT.serialize()).stream(); + null, HddsProtos.QueryScope.CLUSTER, "", ClientVersion.CURRENT).stream(); List decommissioningNodes = DecommissionUtils.getDecommissioningNodesList(allNodes, uuid, ipAddress); String metricsJson = scmClient.getMetrics("Hadoop:service=StorageContainerManager,name=NodeDecommissionMetrics"); diff --git a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/ReconPipelineManager.java b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/ReconPipelineManager.java index 58df5a67530a..92c84a7bd9be 100644 --- a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/ReconPipelineManager.java +++ b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/scm/ReconPipelineManager.java @@ -166,7 +166,7 @@ public boolean addPipeline(Pipeline pipeline) throws IOException { if (containsPipeline(pipeline.getId())) { return false; } - getStateManager().addPipeline(pipeline.getProtobufMessage(ClientVersion.CURRENT.serialize())); + getStateManager().addPipeline(pipeline.getProtobufMessage(ClientVersion.CURRENT)); return true; } finally { releaseWriteLock(); diff --git a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/impl/StorageContainerServiceProviderImpl.java b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/impl/StorageContainerServiceProviderImpl.java index 41ea241b3d8b..82ab1cb227e3 100644 --- a/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/impl/StorageContainerServiceProviderImpl.java +++ b/hadoop-ozone/recon/src/main/java/org/apache/hadoop/ozone/recon/spi/impl/StorageContainerServiceProviderImpl.java @@ -114,7 +114,7 @@ public List getExistContainerWithPipelinesInBatch( @Override public List getNodes() throws IOException { return scmClient.queryNode(null, null, HddsProtos.QueryScope.CLUSTER, - "", ClientVersion.CURRENT.serialize()); + "", ClientVersion.CURRENT); } @Override diff --git a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/api/TestEndpoints.java b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/api/TestEndpoints.java index d4e73ca32b3c..cb11d56d394f 100644 --- a/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/api/TestEndpoints.java +++ b/hadoop-ozone/recon/src/test/java/org/apache/hadoop/ozone/recon/api/TestEndpoints.java @@ -100,6 +100,7 @@ import org.apache.hadoop.hdds.utils.db.Table; import org.apache.hadoop.hdds.utils.db.TypedTable; import org.apache.hadoop.hdfs.web.URLConnectionFactory; +import org.apache.hadoop.ozone.ClientVersion; import org.apache.hadoop.ozone.OzoneAcl; import org.apache.hadoop.ozone.OzoneConsts; import org.apache.hadoop.ozone.om.OMMetadataManager; @@ -1359,7 +1360,7 @@ public void testExplicitRemovalOfNonExistingNode() { @Test public void testSuccessWhenDecommissionStatus() throws IOException { - when(mockScmClient.queryNode(any(), any(), any(), any(), any(Integer.class))).thenReturn( + when(mockScmClient.queryNode(any(), any(), any(), any(), any(ClientVersion.class))).thenReturn( nodes); // 2 nodes decommissioning when(mockScmClient.getContainersOnDecomNode(any())).thenReturn(containerOnDecom); when(mockScmClient.getMetrics(any())).thenReturn(metrics.get(1)); @@ -1385,7 +1386,7 @@ public void testSuccessWhenDecommissionStatus() throws IOException { @Test public void testSuccessWhenDecommissionStatusWithUUID() throws IOException { - when(mockScmClient.queryNode(any(), any(), any(), any(), any(Integer.class))).thenReturn( + when(mockScmClient.queryNode(any(), any(), any(), any(), any(ClientVersion.class))).thenReturn( getNodeDetailsForUuid("654c4b89-04ef-4015-8a3b-50d0fb0e1684")); // 1 nodes decommissioning when(mockScmClient.getContainersOnDecomNode(any())).thenReturn(containerOnDecom); Response datanodesDecommissionInfo =