This is an automated email from the ASF dual-hosted git repository.
errose28 pushed a commit to branch HDDS-14496-zdu
in repository https://gitbox.apache.org/repos/asf/ozone.git
The following commit(s) were added to refs/heads/HDDS-14496-zdu by this push:
new 81fd3c45d01 HDDS-16356. SCM finalize API should only be reachable by
OM (#11177)
81fd3c45d01 is described below
commit 81fd3c45d01b57d2c2999b7c05f6784b6b339ab1
Author: Ethan Rose <[email protected]>
AuthorDate: Thu Sep 3 10:44:41 2026 -0400
HDDS-16356. SCM finalize API should only be reachable by OM (#11177)
---
.../apache/hadoop/hdds/scm/client/ScmClient.java | 2 -
.../scm/protocol/ScmBlockLocationProtocol.java | 24 ++
.../protocol/StorageContainerLocationProtocol.java | 13 -
...lockLocationProtocolClientSideTranslatorPB.java | 34 +++
...inerLocationProtocolClientSideTranslatorPB.java | 30 ---
.../src/main/proto/ScmAdminProtocol.proto | 26 +-
.../src/main/proto/ScmServerProtocol.proto | 20 ++
...lockLocationProtocolServerSideTranslatorPB.java | 29 +++
...inerLocationProtocolServerSideTranslatorPB.java | 26 --
.../hdds/scm/server/SCMBlockProtocolServer.java | 146 +++++++++++
.../hdds/scm/server/SCMClientProtocolServer.java | 145 -----------
.../scm/server/TestSCMBlockProtocolServer.java | 267 +++++++++++++++++++++
.../scm/server/TestSCMClientProtocolServer.java | 264 --------------------
.../hdds/scm/cli/ContainerOperationClient.java | 5 -
.../dist/src/main/compose/common/security.conf | 4 +-
.../hadoop/hdds/upgrade/HddsUpgradeTestUtils.java | 4 +-
.../TestDNDataDistributionFinalization.java | 15 +-
.../TestScmDataDistributionFinalization.java | 15 +-
.../hadoop/hdds/upgrade/TestScmHAFinalization.java | 12 +-
.../ozone/om/service/TestBlockDeletionService.java | 11 +-
.../org/apache/hadoop/ozone/om/OzoneManager.java | 2 +-
.../upgrade/OMFinalizeUpgradeRequestBase.java | 4 +-
.../ozone/om/upgrade/OMUpgradeFinalizeService.java | 2 +-
.../ozone/om/ScmBlockLocationTestingClient.java | 14 ++
.../upgrade/TestOMStartFinalizeUpgradeRequest.java | 18 +-
.../TestOMStartFinalizeUpgradeRequestBase.java | 38 +--
.../TestOMStartFinalizeUpgradeRequestLegacy.java | 2 +-
.../om/upgrade/TestOMUpgradeFinalizeService.java | 24 +-
28 files changed, 616 insertions(+), 580 deletions(-)
diff --git
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmClient.java
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmClient.java
index 2c9e13fb5d3..e03a67417af 100644
---
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmClient.java
+++
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/client/ScmClient.java
@@ -473,8 +473,6 @@ StatusAndMessages queryUpgradeFinalizationProgress(
String upgradeClientID, boolean force, boolean readonly)
throws IOException;
- HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException;
-
DecommissionScmResponseProto decommissionScm(
String scmId) throws IOException;
diff --git
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/ScmBlockLocationProtocol.java
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/ScmBlockLocationProtocol.java
index a34420b3de0..46b635aabb9 100644
---
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/ScmBlockLocationProtocol.java
+++
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocol/ScmBlockLocationProtocol.java
@@ -23,6 +23,7 @@
import java.util.concurrent.TimeoutException;
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationType;
import org.apache.hadoop.hdds.scm.AddSCMRequest;
@@ -145,4 +146,27 @@ List<DatanodeDetails> sortDatanodes(List<String> nodes,
* @throws IOException
*/
InnerNode getNetworkTopology() throws IOException;
+
+ /**
+ * Triggers HDDS finalization on SCM. Called by OM when it orchestrates
cluster finalization.
+ *
+ * @throws IOException If any error occurs.
+ */
+ void finalizeUpgrade() throws IOException;
+
+ /**
+ * Same as {@link #finalizeUpgrade()}, but SCM skips the peer SCM and
datanode software version
+ * checks before finalizing. Use this to finalize when a peer or datanode is
intentionally down.
+ *
+ * @throws IOException If any error occurs.
+ */
+ void forceFinalizeUpgrade() throws IOException;
+
+ /**
+ * Queries the current HDDS finalization status from SCM.
+ *
+ * @return the upgrade status of SCM and the datanodes.
+ * @throws IOException If any error occurs.
+ */
+ HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException;
}
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 9831ced682b..60eb3f0449a 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
@@ -518,24 +518,11 @@ List<HddsProtos.DatanodeUsageInfoProto>
getDatanodeUsageInfo(
StatusAndMessages finalizeScmUpgrade(String upgradeClientID)
throws IOException;
- void finalizeUpgrade() throws IOException;
-
- /**
- * Same as {@link #finalizeUpgrade()}, but SCM skips the peer SCM and
datanode software version
- * checks before finalizing. Use this to finalize when a peer or datanode is
intentionally down or
- * on a different version.
- *
- * @throws IOException If any error occurs.
- */
- void forceFinalizeUpgrade() throws IOException;
-
@Deprecated
StatusAndMessages queryUpgradeFinalizationProgress(
String upgradeClientID, boolean force, boolean readonly)
throws IOException;
- HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException;
-
/**
* Returns the {@link HDDSVersion#SOFTWARE_VERSION} of the responding SCM
binary. Intended for the
* SCM leader to verify that all peer SCMs run a matching software version
before starting
diff --git
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ScmBlockLocationProtocolClientSideTranslatorPB.java
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ScmBlockLocationProtocolClientSideTranslatorPB.java
index e116dea92d0..963c78e4a95 100644
---
a/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ScmBlockLocationProtocolClientSideTranslatorPB.java
+++
b/hadoop-hdds/framework/src/main/java/org/apache/hadoop/hdds/scm/protocolPB/ScmBlockLocationProtocolClientSideTranslatorPB.java
@@ -45,9 +45,11 @@
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.AllocateScmBlockResponseProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.DeleteScmKeyBlocksRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.DeleteScmKeyBlocksResponseProto;
+import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.FinalizeUpgradeRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.GetClusterTreeRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.GetClusterTreeResponseProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.KeyBlocks;
+import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.QueryUpgradeStatusRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.SCMBlockLocationRequest;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.SCMBlockLocationResponse;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.SortDatanodesRequestProto;
@@ -399,6 +401,38 @@ public InnerNode getNetworkTopology() throws IOException {
InnerNodeImpl.fromProtobuf(resp.getClusterTree()));
}
+ @Override
+ public void finalizeUpgrade() throws IOException {
+ finalizeUpgrade(false);
+ }
+
+ @Override
+ public void forceFinalizeUpgrade() throws IOException {
+ finalizeUpgrade(true);
+ }
+
+ private void finalizeUpgrade(boolean force) throws IOException {
+ FinalizeUpgradeRequestProto request =
FinalizeUpgradeRequestProto.newBuilder()
+ .setForce(force)
+ .build();
+ SCMBlockLocationRequest wrapper =
createSCMBlockRequest(Type.FinalizeUpgrade)
+ .setFinalizeUpgradeRequest(request)
+ .build();
+ handleError(submitRequest(wrapper));
+ }
+
+ @Override
+ public HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException {
+ QueryUpgradeStatusRequestProto request =
+ QueryUpgradeStatusRequestProto.newBuilder().build();
+ SCMBlockLocationRequest wrapper =
createSCMBlockRequest(Type.QueryUpgradeStatus)
+ .setQueryUpgradeStatusRequest(request)
+ .build();
+ final SCMBlockLocationResponse wrappedResponse =
+ handleError(submitRequest(wrapper));
+ return wrappedResponse.getQueryUpgradeStatusResponse().getStatus();
+ }
+
/**
* Sets the parent field for the clusterTree nodes recursively.
*
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 f368c583973..4c95727067f 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
@@ -71,7 +71,6 @@
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.DecommissionScmResponseProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.FinalizeScmUpgradeRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.FinalizeScmUpgradeResponseProto;
-import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.FinalizeUpgradeRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ForceExitSafeModeRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.ForceExitSafeModeResponseProto;
import
org.apache.hadoop.hdds.protocol.proto.StorageContainerLocationProtocolProtos.GetContainerCountRequestProto;
@@ -1197,22 +1196,6 @@ public StatusAndMessages finalizeScmUpgrade(String
upgradeClientID)
status.getMessagesList());
}
- @Override
- public void finalizeUpgrade() throws IOException {
- finalizeUpgrade(false);
- }
-
- @Override
- public void forceFinalizeUpgrade() throws IOException {
- finalizeUpgrade(true);
- }
-
- private void finalizeUpgrade(boolean force) throws IOException {
- FinalizeUpgradeRequestProto req = FinalizeUpgradeRequestProto.newBuilder()
- .setForce(force).build();
- submitRequest(Type.FinalizeUpgrade, builder ->
builder.setFinalizeUpgradeRequest(req));
- }
-
@Override
@Deprecated
public StatusAndMessages queryUpgradeFinalizationProgress(
@@ -1237,19 +1220,6 @@ public StatusAndMessages
queryUpgradeFinalizationProgress(
status.getMessagesList());
}
- @Override
- public HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException {
- StorageContainerLocationProtocolProtos.QueryUpgradeStatusRequestProto req =
- StorageContainerLocationProtocolProtos.QueryUpgradeStatusRequestProto
- .newBuilder()
- .build();
-
- StorageContainerLocationProtocolProtos.QueryUpgradeStatusResponseProto
response =
- submitRequest(Type.QueryUpgradeStatus, builder ->
builder.setQueryUpgradeStatusRequest(req))
- .getQueryUpgradeStatusResponse();
- return response.getStatus();
- }
-
@Override
public HDDSVersion getPeerUpgradeStatus() throws IOException {
StorageContainerLocationProtocolProtos.GetPeerUpgradeStatusRequestProto
req =
diff --git a/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto
b/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto
index b61764d1481..cb3e7edd453 100644
--- a/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto
+++ b/hadoop-hdds/interface-admin/src/main/proto/ScmAdminProtocol.proto
@@ -89,9 +89,7 @@ message ScmContainerLocationRequest {
optional GetDeletedBlocksTxnSummaryRequestProto
getDeletedBlocksTxnSummaryRequest = 50;
optional SCMListContainerIDsRequestProto scmListContainerIDsRequest = 51;
optional SuppressContainerRequestProto suppressContainerRequest = 52;
- optional QueryUpgradeStatusRequestProto queryUpgradeStatusRequest = 53;
- optional FinalizeUpgradeRequestProto finalizeUpgradeRequest = 54;
- optional GetPeerUpgradeStatusRequestProto getPeerUpgradeStatusRequest = 55;
+ optional GetPeerUpgradeStatusRequestProto getPeerUpgradeStatusRequest = 53;
}
message ScmContainerLocationResponse {
@@ -152,9 +150,7 @@ message ScmContainerLocationResponse {
optional GetDeletedBlocksTxnSummaryResponseProto
getDeletedBlocksTxnSummaryResponse = 50;
optional SCMListContainerIDsResponseProto scmListContainerIDsResponse = 51;
optional SuppressContainerResponseProto suppressContainerResponse = 52;
- optional QueryUpgradeStatusResponseProto queryUpgradeStatusResponse = 53;
- optional FinalizeUpgradeResponseProto finalizeUpgradeResponse = 54;
- optional GetPeerUpgradeStatusResponseProto getPeerUpgradeStatusResponse = 55;
+ optional GetPeerUpgradeStatusResponseProto getPeerUpgradeStatusResponse = 53;
enum Status {
OK = 1;
@@ -214,9 +210,7 @@ enum Type {
GetDeletedBlocksTransactionSummary = 46;
ListContainerIDs = 47;
SuppressContainer = 48;
- QueryUpgradeStatus = 49;
- FinalizeUpgrade = 50;
- GetPeerUpgradeStatus = 51;
+ GetPeerUpgradeStatus = 49;
}
/**
@@ -598,20 +592,6 @@ message QueryUpgradeFinalizationProgressResponseProto {
required hadoop.hdds.UpgradeFinalizationStatus status = 1;
}
-message QueryUpgradeStatusRequestProto {
-}
-
-message QueryUpgradeStatusResponseProto {
- required hadoop.hdds.UpgradeStatus status = 1;
-}
-
-message FinalizeUpgradeRequestProto {
- optional bool force = 1 [default = false];
-}
-
-message FinalizeUpgradeResponseProto {
-}
-
// Request from an SCM leader to a peer SCM to read its local upgrade status.
// The resulting status returned currently only contains the SCM's software
version to use in pre-finalization
// safety checks.
diff --git
a/hadoop-hdds/interface-server/src/main/proto/ScmServerProtocol.proto
b/hadoop-hdds/interface-server/src/main/proto/ScmServerProtocol.proto
index 4de70addfcf..557b775e5d8 100644
--- a/hadoop-hdds/interface-server/src/main/proto/ScmServerProtocol.proto
+++ b/hadoop-hdds/interface-server/src/main/proto/ScmServerProtocol.proto
@@ -39,6 +39,8 @@ enum Type {
SortDatanodes = 14;
AddScm = 15;
GetClusterTree = 16;
+ FinalizeUpgrade = 17;
+ QueryUpgradeStatus = 18;
}
message SCMBlockLocationRequest {
@@ -57,6 +59,8 @@ message SCMBlockLocationRequest {
optional SortDatanodesRequestProto sortDatanodesRequest = 14;
optional hadoop.hdds.AddScmRequestProto addScmRequestProto = 15;
optional GetClusterTreeRequestProto getClusterTreeRequest = 16;
+ optional FinalizeUpgradeRequestProto finalizeUpgradeRequest = 17;
+ optional QueryUpgradeStatusRequestProto queryUpgradeStatusRequest = 18;
}
message SCMBlockLocationResponse {
@@ -82,6 +86,8 @@ message SCMBlockLocationResponse {
optional SortDatanodesResponseProto sortDatanodesResponse = 14;
optional hadoop.hdds.AddScmResponseProto addScmResponse = 15;
optional GetClusterTreeResponseProto getClusterTreeResponse = 16;
+ optional FinalizeUpgradeResponseProto finalizeUpgradeResponse = 17;
+ optional QueryUpgradeStatusResponseProto queryUpgradeStatusResponse = 18;
}
/**
@@ -246,6 +252,20 @@ message GetClusterTreeResponseProto {
required InnerNode clusterTree = 1;
}
+message FinalizeUpgradeRequestProto {
+ optional bool force = 1 [default = false];
+}
+
+message FinalizeUpgradeResponseProto {
+}
+
+message QueryUpgradeStatusRequestProto {
+}
+
+message QueryUpgradeStatusResponseProto {
+ required hadoop.hdds.UpgradeStatus status = 1;
+}
+
/**
* Protocol used from OzoneManager to StorageContainerManager.
* See request and response messages for details of the RPC calls.
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 cb6be682783..517eae53f42 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
@@ -38,7 +38,11 @@
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.DeleteKeyBlocksResultProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.DeleteScmKeyBlocksRequestProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.DeleteScmKeyBlocksResponseProto;
+import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.FinalizeUpgradeRequestProto;
+import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.FinalizeUpgradeResponseProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.GetClusterTreeResponseProto;
+import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.QueryUpgradeStatusRequestProto;
+import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.QueryUpgradeStatusResponseProto;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.SCMBlockLocationRequest;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.SCMBlockLocationResponse;
import
org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos.SortDatanodesRequestProto;
@@ -172,6 +176,14 @@ private SCMBlockLocationResponse processMessage(
response.setGetClusterTreeResponse(
getClusterTree(clientVersion));
break;
+ case FinalizeUpgrade:
+ response.setFinalizeUpgradeResponse(
+ finalizeUpgrade(request.getFinalizeUpgradeRequest()));
+ break;
+ case QueryUpgradeStatus:
+ response.setQueryUpgradeStatusResponse(
+ queryUpgradeStatus(request.getQueryUpgradeStatusRequest()));
+ break;
default:
// Should never happen
throw new IOException("Unknown Operation " + request.getCmdType() +
@@ -328,4 +340,21 @@ public GetClusterTreeResponseProto
getClusterTree(ClientVersion clientVersion)
resp.setClusterTree(clusterTree.toProtobuf(clientVersion).getInnerNode());
return resp.build();
}
+
+ public FinalizeUpgradeResponseProto finalizeUpgrade(
+ FinalizeUpgradeRequestProto request) throws IOException {
+ if (request.getForce()) {
+ impl.forceFinalizeUpgrade();
+ } else {
+ impl.finalizeUpgrade();
+ }
+ return FinalizeUpgradeResponseProto.newBuilder().build();
+ }
+
+ public QueryUpgradeStatusResponseProto queryUpgradeStatus(
+ QueryUpgradeStatusRequestProto request) throws IOException {
+ return QueryUpgradeStatusResponseProto.newBuilder()
+ .setStatus(impl.queryUpgradeStatus())
+ .build();
+ }
}
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 d65b3d03e5d..a40b9617102 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
@@ -771,22 +771,6 @@ public ScmContainerLocationResponse processRequest(
.setStatus(Status.OK)
.setSuppressContainerResponse(suppressContainer(request.getSuppressContainerRequest()))
.build();
- case QueryUpgradeStatus:
- return ScmContainerLocationResponse.newBuilder()
- .setCmdType(request.getCmdType())
- .setStatus(Status.OK)
-
.setQueryUpgradeStatusResponse(getQueryUpgradeStatus(request.getQueryUpgradeStatusRequest()))
- .build();
- case FinalizeUpgrade:
- if (request.getFinalizeUpgradeRequest().getForce()) {
- impl.forceFinalizeUpgrade();
- } else {
- impl.finalizeUpgrade();
- }
- return ScmContainerLocationResponse.newBuilder()
- .setCmdType(request.getCmdType())
- .setStatus(Status.OK)
- .build();
case GetPeerUpgradeStatus:
return ScmContainerLocationResponse.newBuilder()
.setCmdType(request.getCmdType())
@@ -1150,16 +1134,6 @@ public FinalizeScmUpgradeResponseProto
getFinalizeScmUpgrade(
.build();
}
- public
StorageContainerLocationProtocolProtos.QueryUpgradeStatusResponseProto
getQueryUpgradeStatus(
- StorageContainerLocationProtocolProtos.QueryUpgradeStatusRequestProto
request) throws IOException {
-
- HddsProtos.UpgradeStatus response = impl.queryUpgradeStatus();
- return
StorageContainerLocationProtocolProtos.QueryUpgradeStatusResponseProto
- .newBuilder()
- .setStatus(response)
- .build();
- }
-
public
StorageContainerLocationProtocolProtos.GetPeerUpgradeStatusResponseProto
getPeerUpgradeStatus(
StorageContainerLocationProtocolProtos.GetPeerUpgradeStatusRequestProto
request) throws IOException {
diff --git
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMBlockProtocolServer.java
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMBlockProtocolServer.java
index 4e3d3b7b3c1..3678925a8e4 100644
---
a/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMBlockProtocolServer.java
+++
b/hadoop-hdds/server-scm/src/main/java/org/apache/hadoop/hdds/scm/server/SCMBlockProtocolServer.java
@@ -35,6 +35,7 @@
import java.io.IOException;
import java.net.InetSocketAddress;
import java.util.ArrayList;
+import java.util.Collections;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
@@ -43,10 +44,12 @@
import java.util.stream.Collectors;
import org.apache.commons.lang3.StringUtils;
import org.apache.hadoop.fs.CommonConfigurationKeysPublic;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
import org.apache.hadoop.hdds.protocol.DatanodeID;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos;
import org.apache.hadoop.hdds.scm.AddSCMRequest;
import org.apache.hadoop.hdds.scm.ScmInfo;
@@ -55,13 +58,17 @@
import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList;
import
org.apache.hadoop.hdds.scm.container.placement.metrics.SCMPerformanceMetrics;
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
+import org.apache.hadoop.hdds.scm.ha.SCMNodeDetails;
import org.apache.hadoop.hdds.scm.net.InnerNode;
import org.apache.hadoop.hdds.scm.net.Node;
import org.apache.hadoop.hdds.scm.net.NodeImpl;
import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol;
import
org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocolServerSideTranslatorPB;
+import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
import org.apache.hadoop.hdds.scm.protocolPB.ScmBlockLocationProtocolPB;
+import
org.apache.hadoop.hdds.scm.protocolPB.StorageContainerLocationProtocolClientSideTranslatorPB.ScmNodeTarget;
+import org.apache.hadoop.hdds.utils.HAUtils;
import org.apache.hadoop.hdds.utils.HddsServerUtil;
import org.apache.hadoop.hdds.utils.ProtocolMessageMetrics;
import org.apache.hadoop.io.IOUtils;
@@ -78,6 +85,7 @@
import org.apache.hadoop.ozone.common.BlockGroup;
import org.apache.hadoop.ozone.common.DeleteBlockGroupResult;
import org.apache.hadoop.ozone.common.DeletedBlock;
+import org.apache.hadoop.security.UserGroupInformation;
import org.apache.hadoop.util.Time;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -442,6 +450,144 @@ public InnerNode getNetworkTopology() {
return (InnerNode) scm.getClusterMap().getNode(ROOT);
}
+ @Override
+ public void finalizeUpgrade() throws IOException {
+ finalizeUpgrade(false);
+ }
+
+ @Override
+ public void forceFinalizeUpgrade() throws IOException {
+ finalizeUpgrade(true);
+ }
+
+ private void finalizeUpgrade(boolean force) throws IOException {
+ final Map<String, String> auditMap = Collections.singletonMap("force",
String.valueOf(force));
+ try {
+ if (force) {
+ LOG.warn("Forcing upgrade finalization by skipping SCM peer and
datanode software version checks");
+ } else {
+ validatePeerScmVersionsBeforeFinalize();
+ validateDatanodeVersionsBeforeFinalize();
+ }
+ scm.getFinalizationManager().finalizeUpgrade();
+
AUDIT.logWriteSuccess(buildAuditMessageForSuccess(SCMAction.FINALIZE_SCM_UPGRADE,
auditMap));
+ } catch (Exception ex) {
+
AUDIT.logWriteFailure(buildAuditMessageForFailure(SCMAction.FINALIZE_SCM_UPGRADE,
auditMap, ex));
+ throw ex;
+ }
+ }
+
+ /**
+ * Verifies that every peer SCM in the Ratis group runs the same software
version as this leader
+ * before finalization begins. Rejects the finalize command if any peer
reports a differing
+ * version or cannot be reached. The resulting exception propagates back to
the OM (which triggered
+ * finalization) and on to the client, leaving nothing finalized.
+ */
+ private void validatePeerScmVersionsBeforeFinalize() throws SCMException {
+ List<SCMNodeDetails> peerNodes =
scm.getSCMHANodeDetails().getPeerNodeDetails();
+ if (peerNodes.isEmpty()) {
+ return;
+ }
+ HDDSVersion leaderVersion = HDDSVersion.SOFTWARE_VERSION;
+ OzoneConfiguration conf = scm.getConfiguration();
+ List<String> failedPeers = new ArrayList<>();
+ for (SCMNodeDetails peer : peerNodes) {
+ String peerId = peer.getNodeId();
+ ScmNodeTarget target = new ScmNodeTarget();
+ target.setNodeId(peerId);
+ StorageContainerLocationProtocol peerClient = null;
+ try {
+ // Contact the peer SCM as this SCM's service (Kerberos keytab)
identity, not the remote client's identity.
+ // This runs inside the finalize RPC handler's doAs context, whose UGI
has no credentials to open a fresh
+ // outbound RPC to the peer SCM.
+ peerClient = HAUtils.getScmContainerClientForNode(conf, target,
UserGroupInformation.getLoginUser());
+ HDDSVersion peerVersion = peerClient.getPeerUpgradeStatus();
+ if (!peerVersion.equals(leaderVersion)) {
+ LOG.warn("SCM peer {} is running software version {} but leader is
running version {}. "
+ + "Rejecting finalize command.", peerId, peerVersion,
leaderVersion);
+ failedPeers.add(peerId + " (version: " + peerVersion + ")");
+ }
+ } catch (IOException e) {
+ LOG.warn("Failed to contact SCM peer {} to check software version
before finalize.", peerId, e);
+ failedPeers.add(peerId + " (unreachable: " + e.getMessage() + ")");
+ } finally {
+ IOUtils.cleanupWithLogger(LOG, peerClient);
+ }
+ }
+ if (!failedPeers.isEmpty()) {
+ throw new SCMException("Finalize rejected: the following SCM peers did
not confirm matching software "
+ + "version (expected version=" + leaderVersion + "): " +
String.join(", ", failedPeers),
+ SCMException.ResultCodes.UNSUPPORTED_OPERATION);
+ }
+ }
+
+ /**
+ * Verifies that every healthy datanode runs the same software version as
this SCM before
+ * finalization begins. Datanodes finalize only after SCM instructs them to,
so this does not
+ * require them to be finalized; it requires their binaries to match SCM's
software version.
+ * Rejects the finalize command if any healthy datanode reports a differing
(or unknown) software
+ * version. The resulting exception propagates back to the OM (which
triggered finalization) and on
+ * to the client, leaving nothing finalized.
+ */
+ private void validateDatanodeVersionsBeforeFinalize() throws SCMException {
+ NodeManager.DatanodeFinalizationCounts counts =
+ scm.getScmNodeManager().getDatanodeFinalizationCounts();
+ if (!counts.allSoftwareVersionsMatchScmVersion()) {
+ LOG.warn("Rejecting finalize command: not all {} healthy datanodes are
running SCM's software "
+ + "version {}.", counts.getTotalHealthyDatanodes(),
HDDSVersion.SOFTWARE_VERSION);
+ throw new SCMException("Finalize rejected: not all healthy datanodes are
running the SCM software version "
+ + HDDSVersion.SOFTWARE_VERSION,
SCMException.ResultCodes.UNSUPPORTED_OPERATION);
+ }
+ }
+
+ @Override
+ public HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException {
+ try {
+ if (scm.getScmContext().isInSafeMode()) {
+ throw new SCMException("Cannot query upgrade status while SCM is in
safe mode. Wait until SCM exits "
+ + "safe mode and try again.",
SCMException.ResultCodes.SAFE_MODE_EXCEPTION);
+ }
+
+ // Set SCM finalization status to return to the client.
+ // Since SCM finalization goes through Ratis, it moves from unfinalized
to finalized immediately with no
+ // in-progress state.
+ boolean scmFinalized = !scm.getVersionManager().needsFinalization();
+ HddsProtos.FinalizationStatus scmFinalizationStatus =
+ scmFinalized ? HddsProtos.FinalizationStatus.FINALIZED :
HddsProtos.FinalizationStatus.UNFINALIZED;
+
+ // Set overall HDDS finalization status (SCM and Datanodes) to return to
the client.
+ NodeManager.DatanodeFinalizationCounts datanodeFinalizationCounts =
+ scm.getScmNodeManager().getDatanodeFinalizationCounts();
+ int finalizedDatanodes =
datanodeFinalizationCounts.getNumFinalizedDatanodes();
+ int healthyDatanodes =
datanodeFinalizationCounts.getTotalHealthyDatanodes();
+ HddsProtos.FinalizationStatus hddsFinalizationStatus;
+ if (!scmFinalized) {
+ // SCM must finish finalizing before Datanodes can start finalizing.
+ hddsFinalizationStatus = HddsProtos.FinalizationStatus.UNFINALIZED;
+ } else if (datanodeFinalizationCounts.allNodesFinalized()) {
+ hddsFinalizationStatus = HddsProtos.FinalizationStatus.FINALIZED;
+ } else {
+ hddsFinalizationStatus = HddsProtos.FinalizationStatus.IN_PROGRESS;
+ }
+
+ HddsProtos.UpgradeStatus result = HddsProtos.UpgradeStatus.newBuilder()
+ .setScmFinalizationStatus(scmFinalizationStatus)
+ .setNumDatanodesFinalized(finalizedDatanodes)
+ .setNumDatanodesTotal(healthyDatanodes)
+ .setHddsFinalizationStatus(hddsFinalizationStatus)
+
.setScmApparentVersion(scm.getVersionManager().getApparentVersion().serialize())
+
.setMinDatanodeApparentVersion(datanodeFinalizationCounts.getMinApparentVersion())
+
.setMaxDatanodeApparentVersion(datanodeFinalizationCounts.getMaxApparentVersion())
+ .build();
+
+
AUDIT.logReadSuccess(buildAuditMessageForSuccess(SCMAction.QUERY_UPGRADE_STATUS,
null));
+ return result;
+ } catch (IOException ex) {
+
AUDIT.logReadFailure(buildAuditMessageForFailure(SCMAction.QUERY_UPGRADE_STATUS,
null, ex));
+ throw ex;
+ }
+ }
+
@Override
public AuditMessage buildAuditMessageForSuccess(
AuditAction op, Map<String, String> auditMap) {
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 d6304413a5a..704213f7925 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
@@ -93,12 +93,10 @@
import org.apache.hadoop.hdds.scm.events.SCMEvents;
import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.exceptions.SCMException.ResultCodes;
-import org.apache.hadoop.hdds.scm.ha.SCMNodeDetails;
import org.apache.hadoop.hdds.scm.ha.SCMRatisServer;
import org.apache.hadoop.hdds.scm.ha.SCMRatisServerImpl;
import org.apache.hadoop.hdds.scm.node.DatanodeInfo;
import org.apache.hadoop.hdds.scm.node.DatanodeUsageInfo;
-import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.node.NodeStatus;
import org.apache.hadoop.hdds.scm.node.states.NodeNotFoundException;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
@@ -107,10 +105,8 @@
import org.apache.hadoop.hdds.scm.pipeline.PipelineNotFoundException;
import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
import
org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocolServerSideTranslatorPB;
-import
org.apache.hadoop.hdds.scm.protocolPB.StorageContainerLocationProtocolClientSideTranslatorPB.ScmNodeTarget;
import
org.apache.hadoop.hdds.scm.protocolPB.StorageContainerLocationProtocolPB;
import org.apache.hadoop.hdds.security.SecurityConfig;
-import org.apache.hadoop.hdds.utils.HAUtils;
import org.apache.hadoop.hdds.utils.HddsServerUtil;
import org.apache.hadoop.hdds.utils.ProtocolMessageMetrics;
import org.apache.hadoop.io.IOUtils;
@@ -1169,102 +1165,11 @@ public StatusAndMessages finalizeScmUpgrade(String
upgradeClientID) {
return new StatusAndMessages(ALREADY_FINALIZED, Collections.emptyList());
}
- @Override
- public void finalizeUpgrade() throws IOException {
- finalizeUpgrade(false);
- }
-
- @Override
- public void forceFinalizeUpgrade() throws IOException {
- finalizeUpgrade(true);
- }
-
- private void finalizeUpgrade(boolean force) throws IOException {
- final Map<String, String> auditMap = Collections.singletonMap("force",
String.valueOf(force));
- try {
- getScm().checkAdminAccess(getRemoteUser(), false);
- if (force) {
- LOG.warn("Forcing upgrade finalization by skipping SCM peer and
datanode software version checks");
- } else {
- validatePeerScmVersionsBeforeFinalize();
- validateDatanodeVersionsBeforeFinalize();
- }
- scm.getFinalizationManager().finalizeUpgrade();
-
AUDIT.logWriteSuccess(buildAuditMessageForSuccess(SCMAction.FINALIZE_SCM_UPGRADE,
auditMap));
- } catch (Exception ex) {
-
AUDIT.logWriteFailure(buildAuditMessageForFailure(SCMAction.FINALIZE_SCM_UPGRADE,
auditMap, ex));
- throw ex;
- }
- }
-
@Override
public HDDSVersion getPeerUpgradeStatus() throws IOException {
return HDDSVersion.SOFTWARE_VERSION;
}
- /**
- * Verifies that every peer SCM in the Ratis group runs the same software
version as this leader
- * before finalization begins. Rejects the finalize command if any peer
reports a differing
- * version or cannot be reached. The resulting exception propagates back to
the OM (which triggered
- * finalization) and on to the client, leaving nothing finalized.
- */
- private void validatePeerScmVersionsBeforeFinalize() throws SCMException {
- List<SCMNodeDetails> peerNodes =
scm.getSCMHANodeDetails().getPeerNodeDetails();
- if (peerNodes.isEmpty()) {
- return;
- }
- HDDSVersion leaderVersion = HDDSVersion.SOFTWARE_VERSION;
- OzoneConfiguration conf = scm.getConfiguration();
- List<String> failedPeers = new ArrayList<>();
- for (SCMNodeDetails peer : peerNodes) {
- String peerId = peer.getNodeId();
- ScmNodeTarget target = new ScmNodeTarget();
- target.setNodeId(peerId);
- StorageContainerLocationProtocol peerClient = null;
- try {
- // Contact the peer SCM as this SCM's service (Kerberos keytab)
identity, not the remote client's identity.
- // This runs inside the finalize RPC handler's doAs context, whose UGI
has no credentials to open a fresh
- // outbound RPC to the peer SCM.
- peerClient = HAUtils.getScmContainerClientForNode(conf, target,
UserGroupInformation.getLoginUser());
- HDDSVersion peerVersion = peerClient.getPeerUpgradeStatus();
- if (!peerVersion.equals(leaderVersion)) {
- LOG.warn("SCM peer {} is running software version {} but leader is
running version {}. "
- + "Rejecting finalize command.", peerId, peerVersion,
leaderVersion);
- failedPeers.add(peerId + " (version: " + peerVersion + ")");
- }
- } catch (IOException e) {
- LOG.warn("Failed to contact SCM peer {} to check software version
before finalize.", peerId, e);
- failedPeers.add(peerId + " (unreachable: " + e.getMessage() + ")");
- } finally {
- IOUtils.cleanupWithLogger(LOG, peerClient);
- }
- }
- if (!failedPeers.isEmpty()) {
- throw new SCMException("Finalize rejected: the following SCM peers did
not confirm matching software "
- + "version (expected version=" + leaderVersion + "): " +
String.join(", ", failedPeers),
- ResultCodes.UNSUPPORTED_OPERATION);
- }
- }
-
- /**
- * Verifies that every healthy datanode runs the same software version as
this SCM before
- * finalization begins. Datanodes finalize only after SCM instructs them to,
so this does not
- * require them to be finalized; it requires their binaries to match SCM's
software version.
- * Rejects the finalize command if any healthy datanode reports a differing
(or unknown) software
- * version. The resulting exception propagates back to the OM (which
triggered finalization) and on
- * to the client, leaving nothing finalized.
- */
- private void validateDatanodeVersionsBeforeFinalize() throws SCMException {
- NodeManager.DatanodeFinalizationCounts counts =
- scm.getScmNodeManager().getDatanodeFinalizationCounts();
- if (!counts.allSoftwareVersionsMatchScmVersion()) {
- LOG.warn("Rejecting finalize command: not all {} healthy datanodes are
running SCM's software "
- + "version {}.", counts.getTotalHealthyDatanodes(),
HDDSVersion.SOFTWARE_VERSION);
- throw new SCMException("Finalize rejected: not all healthy datanodes are
running the SCM software version "
- + HDDSVersion.SOFTWARE_VERSION, ResultCodes.UNSUPPORTED_OPERATION);
- }
- }
-
@Override
@Deprecated
public StatusAndMessages queryUpgradeFinalizationProgress(
@@ -1297,56 +1202,6 @@ public StatusAndMessages
queryUpgradeFinalizationProgress(
}
}
- @Override
- public HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException {
- try {
- getScm().checkAdminAccess(getRemoteUser(), true);
-
- if (scm.getScmContext().isInSafeMode()) {
- throw new SCMException("Cannot query upgrade status while SCM is in
safe mode. Wait until SCM exits "
- + "safe mode and try again.", ResultCodes.SAFE_MODE_EXCEPTION);
- }
-
- // Set SCM finalization status to return to the client.
- // Since SCM finalization goes through Ratis, it moves from unfinalized
to finalized immediately with no
- // in-progress state.
- boolean scmFinalized = !scm.getVersionManager().needsFinalization();
- HddsProtos.FinalizationStatus scmFinalizationStatus =
- scmFinalized ? HddsProtos.FinalizationStatus.FINALIZED :
HddsProtos.FinalizationStatus.UNFINALIZED;
-
- // Set overall HDDS finalization status (SCM and Datanodes) to return to
the client.
- NodeManager.DatanodeFinalizationCounts datanodeFinalizationCounts =
- scm.getScmNodeManager().getDatanodeFinalizationCounts();
- int finalizedDatanodes =
datanodeFinalizationCounts.getNumFinalizedDatanodes();
- int healthyDatanodes =
datanodeFinalizationCounts.getTotalHealthyDatanodes();
- HddsProtos.FinalizationStatus hddsFinalizationStatus;
- if (!scmFinalized) {
- // SCM must finish finalizing before Datanodes can start finalizing.
- hddsFinalizationStatus = HddsProtos.FinalizationStatus.UNFINALIZED;
- } else if (datanodeFinalizationCounts.allNodesFinalized()) {
- hddsFinalizationStatus = HddsProtos.FinalizationStatus.FINALIZED;
- } else {
- hddsFinalizationStatus = HddsProtos.FinalizationStatus.IN_PROGRESS;
- }
-
- HddsProtos.UpgradeStatus result = HddsProtos.UpgradeStatus.newBuilder()
- .setScmFinalizationStatus(scmFinalizationStatus)
- .setNumDatanodesFinalized(finalizedDatanodes)
- .setNumDatanodesTotal(healthyDatanodes)
- .setHddsFinalizationStatus(hddsFinalizationStatus)
-
.setScmApparentVersion(scm.getVersionManager().getApparentVersion().serialize())
-
.setMinDatanodeApparentVersion(datanodeFinalizationCounts.getMinApparentVersion())
-
.setMaxDatanodeApparentVersion(datanodeFinalizationCounts.getMaxApparentVersion())
- .build();
-
-
AUDIT.logReadSuccess(buildAuditMessageForSuccess(SCMAction.QUERY_UPGRADE_STATUS,
null));
- return result;
- } catch (IOException ex) {
-
AUDIT.logReadFailure(buildAuditMessageForFailure(SCMAction.QUERY_UPGRADE_STATUS,
null, ex));
- throw ex;
- }
- }
-
@Override
public StartContainerBalancerResponseProto startContainerBalancer(
Optional<Double> threshold, Optional<Integer> iterations,
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 fa733d3712a..816462a366f 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
@@ -24,13 +24,23 @@
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.junit.jupiter.api.Assertions.fail;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
import com.google.common.collect.ImmutableMap;
import java.io.File;
import java.io.IOException;
+import java.net.InetSocketAddress;
import java.util.ArrayList;
+import java.util.Arrays;
import java.util.Collections;
import java.util.List;
import java.util.Map;
@@ -38,12 +48,15 @@
import java.util.concurrent.ThreadLocalRandom;
import java.util.concurrent.TimeoutException;
import java.util.stream.Collectors;
+import org.apache.hadoop.hdds.ComponentVersion;
+import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.HddsConfigKeys;
import org.apache.hadoop.hdds.client.ContainerBlockID;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
import org.apache.hadoop.hdds.protocol.proto.ScmBlockLocationProtocolProtos;
import org.apache.hadoop.hdds.scm.HddsTestUtils;
@@ -53,23 +66,34 @@
import org.apache.hadoop.hdds.scm.block.SCMBlockDeletingService;
import org.apache.hadoop.hdds.scm.container.common.helpers.AllocatedBlock;
import org.apache.hadoop.hdds.scm.container.common.helpers.ExcludeList;
+import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.ha.SCMContext;
import org.apache.hadoop.hdds.scm.ha.SCMHAManagerStub;
+import org.apache.hadoop.hdds.scm.ha.SCMHANodeDetails;
+import org.apache.hadoop.hdds.scm.ha.SCMNodeDetails;
import org.apache.hadoop.hdds.scm.net.NodeImpl;
import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline;
import org.apache.hadoop.hdds.scm.pipeline.Pipeline.PipelineState;
import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
import
org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocolServerSideTranslatorPB;
+import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
+import org.apache.hadoop.hdds.scm.safemode.SCMSafeModeManager.SafeModeStatus;
+import org.apache.hadoop.hdds.scm.server.upgrade.FinalizationManager;
+import org.apache.hadoop.hdds.scm.server.upgrade.ScmVersionManager;
+import org.apache.hadoop.hdds.utils.HAUtils;
import org.apache.hadoop.hdds.utils.ProtocolMessageMetrics;
import org.apache.hadoop.net.StaticMapping;
import org.apache.hadoop.ozone.ClientVersion;
+import org.apache.hadoop.ozone.audit.SCMAction;
import org.apache.hadoop.ozone.common.BlockGroup;
import org.apache.hadoop.ozone.container.common.SCMTestUtils;
+import org.apache.ozone.test.GenericTestUtils.LogCapturer;
import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
+import org.mockito.MockedStatic;
/**
* Test class for @{@link SCMBlockProtocolServer}.
@@ -338,6 +362,249 @@ void testAllocateBlockWithClientMachine() throws
IOException {
}
}
+ @Test
+ public void testFinalizeProceedsWhenNoPeers() throws IOException {
+ FinalizationManager finalizationManager = mock(FinalizationManager.class);
+ try (SCMBlockProtocolServer testServer =
+ buildTestServer(finalizationManager, Collections.emptyList())) {
+ testServer.finalizeUpgrade();
+ verify(finalizationManager).finalizeUpgrade();
+ }
+ }
+
+ @Test
+ public void testFinalizeProceedsWhenAllPeersMatch() throws IOException {
+ FinalizationManager finalizationManager = mock(FinalizationManager.class);
+ StorageContainerLocationProtocol matching =
peerClient(HDDSVersion.SOFTWARE_VERSION);
+
+ try (SCMBlockProtocolServer testServer =
+ buildTestServer(finalizationManager,
Arrays.asList(peerNode("scm2"), peerNode("scm3")));
+ MockedStatic<HAUtils> haUtils = mockStatic(HAUtils.class)) {
+ haUtils.when(() -> HAUtils.getScmContainerClientForNode(any(), any(),
any())).thenReturn(matching);
+ testServer.finalizeUpgrade();
+ }
+ verify(finalizationManager).finalizeUpgrade();
+ }
+
+ @Test
+ public void testFinalizeRejectsOlderPeerUnlessForced() throws IOException {
+ FinalizationManager finalizationManager = mock(FinalizationManager.class);
+ StorageContainerLocationProtocol matching =
peerClient(HDDSVersion.SOFTWARE_VERSION);
+ StorageContainerLocationProtocol older =
peerClient(HDDSVersion.DEFAULT_VERSION);
+
+ try (SCMBlockProtocolServer testServer =
+ buildTestServer(finalizationManager,
Arrays.asList(peerNode("scm2"), peerNode("scm3")));
+ MockedStatic<HAUtils> haUtils = mockStatic(HAUtils.class)) {
+ haUtils.when(() -> HAUtils.getScmContainerClientForNode(any(), any(),
any())).thenReturn(matching, older);
+ // A peer on an older version is rejected without force.
+ assertThrows(SCMException.class, testServer::finalizeUpgrade);
+ verify(finalizationManager, never()).finalizeUpgrade();
+ // With force the peer version check is skipped and finalization
proceeds.
+ testServer.forceFinalizeUpgrade();
+ }
+ verify(finalizationManager).finalizeUpgrade();
+ }
+
+ @Test
+ public void testFinalizeRejectsUnknownFuturePeerUnlessForced() throws
IOException {
+ FinalizationManager finalizationManager = mock(FinalizationManager.class);
+ StorageContainerLocationProtocol matching =
peerClient(HDDSVersion.SOFTWARE_VERSION);
+ // A version not recognized by this binary deserializes to UNKNOWN_VERSION
in the client translator.
+ StorageContainerLocationProtocol unknown =
peerClient(HDDSVersion.UNKNOWN_VERSION);
+
+ try (SCMBlockProtocolServer testServer =
+ buildTestServer(finalizationManager,
Arrays.asList(peerNode("scm2"), peerNode("scm3")));
+ MockedStatic<HAUtils> haUtils = mockStatic(HAUtils.class)) {
+ haUtils.when(() -> HAUtils.getScmContainerClientForNode(any(), any(),
any()))
+ .thenReturn(matching, unknown);
+ // A peer on an unrecognized future version is rejected without force.
+ assertThrows(SCMException.class, testServer::finalizeUpgrade);
+ verify(finalizationManager, never()).finalizeUpgrade();
+ // With force the peer version check is skipped and finalization
proceeds.
+ testServer.forceFinalizeUpgrade();
+ }
+ verify(finalizationManager).finalizeUpgrade();
+ }
+
+ @Test
+ public void testFinalizeRejectsUnreachablePeerUnlessForced() throws
IOException {
+ FinalizationManager finalizationManager = mock(FinalizationManager.class);
+ StorageContainerLocationProtocol unreachable =
mock(StorageContainerLocationProtocol.class);
+ when(unreachable.getPeerUpgradeStatus()).thenThrow(new
IOException("connection refused"));
+
+ try (SCMBlockProtocolServer testServer =
+ buildTestServer(finalizationManager,
Collections.singletonList(peerNode("scm2")));
+ MockedStatic<HAUtils> haUtils = mockStatic(HAUtils.class)) {
+ haUtils.when(() -> HAUtils.getScmContainerClientForNode(any(), any(),
any())).thenReturn(unreachable);
+ // An unreachable peer is rejected without force.
+ assertThrows(SCMException.class, testServer::finalizeUpgrade);
+ verify(finalizationManager, never()).finalizeUpgrade();
+ // With force the peer version check is skipped and finalization
proceeds.
+ testServer.forceFinalizeUpgrade();
+ }
+ verify(finalizationManager).finalizeUpgrade();
+ }
+
+ @Test
+ public void testFinalizeProceedsWhenAllDatanodesMatchScmVersion() throws
IOException {
+ FinalizationManager finalizationManager = mock(FinalizationManager.class);
+ NodeManager.DatanodeFinalizationCounts datanodeCounts =
NodeManager.DatanodeFinalizationCounts.newBuilder()
+ .setNumFinalizedDatanodes(3)
+ .setTotalHealthyDatanodes(3)
+ .setMinApparentVersion(HDDSVersion.SOFTWARE_VERSION.serialize())
+ .setMaxApparentVersion(HDDSVersion.SOFTWARE_VERSION.serialize())
+ .setAllSoftwareVersionsMatchScm(true)
+ .build();
+ try (SCMBlockProtocolServer testServer =
buildTestServer(finalizationManager,
+ Collections.emptyList(), datanodeCounts)) {
+ testServer.finalizeUpgrade();
+ verify(finalizationManager).finalizeUpgrade();
+ }
+ }
+
+ @Test
+ public void testFinalizeRejectsDatanodeWithMismatchedVersionUnlessForced()
throws IOException {
+ FinalizationManager finalizationManager = mock(FinalizationManager.class);
+ NodeManager.DatanodeFinalizationCounts datanodeCounts =
NodeManager.DatanodeFinalizationCounts.newBuilder()
+ .setNumFinalizedDatanodes(3)
+ .setTotalHealthyDatanodes(3)
+ .setMinApparentVersion(HDDSVersion.DEFAULT_VERSION.serialize())
+ .setMaxApparentVersion(HDDSVersion.SOFTWARE_VERSION.serialize())
+ .setAllSoftwareVersionsMatchScm(false)
+ .build();
+ try (SCMBlockProtocolServer testServer =
buildTestServer(finalizationManager,
+ Collections.emptyList(), datanodeCounts)) {
+ // A datanode on a mismatched version is rejected without force.
+ assertThrows(SCMException.class, testServer::finalizeUpgrade);
+ verify(finalizationManager, never()).finalizeUpgrade();
+ // With force the datanode version check is skipped and finalization
proceeds.
+ testServer.forceFinalizeUpgrade();
+ verify(finalizationManager).finalizeUpgrade();
+ }
+ }
+
+ @Test
+ public void testForceFinalizeAuditRecordsForceFlag() throws IOException {
+ FinalizationManager finalizationManager = mock(FinalizationManager.class);
+ // Drive the failure audit path (logged at ERROR, captured by the default
log4j2 config) so the
+ // recorded audit map can be inspected.
+ doThrow(new
RuntimeException("test")).when(finalizationManager).finalizeUpgrade();
+
+ LogCapturer auditLog = LogCapturer.log4j2("SCMAudit");
+ try (SCMBlockProtocolServer testServer =
buildTestServer(finalizationManager, Collections.emptyList())) {
+ assertThrows(RuntimeException.class, testServer::forceFinalizeUpgrade);
+ } finally {
+ auditLog.stopCapturing();
+ }
+
+ String output = auditLog.getOutput();
+ assertTrue(output.contains(SCMAction.FINALIZE_SCM_UPGRADE.getAction()),
+ "audit log should record the finalize action: " + output);
+ assertTrue(output.contains("\"force\":\"true\""),
+ "audit log should record that force was passed: " + output);
+ }
+
+ @Test
+ public void testQueryUpgradeStatus() throws Exception {
+ // SCM starts already finalized in tests.
+ HddsProtos.UpgradeStatus status = server.queryUpgradeStatus();
+ assertEquals(NODE_COUNT, status.getNumDatanodesFinalized());
+ assertEquals(NODE_COUNT, status.getNumDatanodesTotal());
+ assertEquals(HddsProtos.FinalizationStatus.FINALIZED,
status.getScmFinalizationStatus());
+ }
+
+ @Test
+ public void testQueryUpgradeStatusHddsInProgress() throws Exception {
+ // SCM is finalized but not all datanodes are, so HDDS finalization is
still in progress.
+ ScmVersionManager mockVersionManager = mock(ScmVersionManager.class);
+ when(mockVersionManager.needsFinalization()).thenReturn(false);
+ ComponentVersion apparentVersion = mock(ComponentVersion.class);
+ when(apparentVersion.serialize()).thenReturn(0);
+ when(mockVersionManager.getApparentVersion()).thenReturn(apparentVersion);
+
+ NodeManager.DatanodeFinalizationCounts datanodeCounts =
NodeManager.DatanodeFinalizationCounts.newBuilder()
+ .setNumFinalizedDatanodes(1)
+ .setTotalHealthyDatanodes(3)
+ .build();
+ NodeManager mockNodeManager = mock(NodeManager.class);
+
when(mockNodeManager.getDatanodeFinalizationCounts()).thenReturn(datanodeCounts);
+
+ StorageContainerManager mockScm = mockScmForBlockServer();
+ when(mockScm.getVersionManager()).thenReturn(mockVersionManager);
+ when(mockScm.getScmNodeManager()).thenReturn(mockNodeManager);
+ when(mockScm.getScmContext()).thenReturn(SCMContext.emptyContext());
+
+ try (SCMBlockProtocolServer testServer = new SCMBlockProtocolServer(new
OzoneConfiguration(), mockScm)) {
+ HddsProtos.UpgradeStatus status = testServer.queryUpgradeStatus();
+ assertEquals(HddsProtos.FinalizationStatus.FINALIZED,
status.getScmFinalizationStatus());
+ assertEquals(HddsProtos.FinalizationStatus.IN_PROGRESS,
status.getHddsFinalizationStatus());
+ assertEquals(1, status.getNumDatanodesFinalized());
+ assertEquals(3, status.getNumDatanodesTotal());
+ }
+ }
+
+ @Test
+ public void testQueryUpgradeStatusInSafemode() {
+ // Put SCM into safe mode via the context the server consults.
+ scm.getScmContext().updateSafeModeStatus(SafeModeStatus.INITIAL);
+ assertTrue(scm.getScmContext().isInSafeMode());
+
+ // Querying upgrade status is blocked while SCM is in safe mode.
+ SCMException ex = assertThrows(SCMException.class, () ->
server.queryUpgradeStatus());
+ assertEquals(SCMException.ResultCodes.SAFE_MODE_EXCEPTION, ex.getResult());
+ }
+
+ private SCMBlockProtocolServer buildTestServer(
+ FinalizationManager finalizationManager, List<SCMNodeDetails> peers)
throws IOException {
+ // Default to all datanode versions matching SCM so the SCM peer checks
are exercised in isolation.
+ NodeManager.DatanodeFinalizationCounts datanodeCounts =
NodeManager.DatanodeFinalizationCounts.newBuilder()
+ .setNumFinalizedDatanodes(0)
+ .setTotalHealthyDatanodes(0)
+ .setMinApparentVersion(0)
+ .setMaxApparentVersion(0)
+ .setAllSoftwareVersionsMatchScm(true)
+ .build();
+ return buildTestServer(finalizationManager, peers, datanodeCounts);
+ }
+
+ private SCMBlockProtocolServer buildTestServer(
+ FinalizationManager finalizationManager, List<SCMNodeDetails> peers,
+ NodeManager.DatanodeFinalizationCounts datanodeCounts) throws
IOException {
+ StorageContainerManager mockScm = mockScmForBlockServer();
+ when(mockScm.getFinalizationManager()).thenReturn(finalizationManager);
+ when(mockScm.getConfiguration()).thenReturn(new OzoneConfiguration());
+
+ SCMHANodeDetails haNodeDetails = mock(SCMHANodeDetails.class);
+ when(haNodeDetails.getPeerNodeDetails()).thenReturn(peers);
+ when(mockScm.getSCMHANodeDetails()).thenReturn(haNodeDetails);
+
+ NodeManager mockNodeManager = mock(NodeManager.class);
+
when(mockNodeManager.getDatanodeFinalizationCounts()).thenReturn(datanodeCounts);
+ when(mockScm.getScmNodeManager()).thenReturn(mockNodeManager);
+ return new SCMBlockProtocolServer(new OzoneConfiguration(), mockScm);
+ }
+
+ private StorageContainerManager mockScmForBlockServer() {
+ StorageContainerManager mockScm = mock(StorageContainerManager.class);
+ SCMNodeDetails scmNodeDetails = mock(SCMNodeDetails.class);
+ when(scmNodeDetails.getBlockProtocolServerAddress()).thenReturn(new
InetSocketAddress("localhost", 0));
+ when(scmNodeDetails.getBlockProtocolServerAddressKey()).thenReturn("test");
+ when(mockScm.getScmNodeDetails()).thenReturn(scmNodeDetails);
+ return mockScm;
+ }
+
+ private StorageContainerLocationProtocol peerClient(HDDSVersion version)
throws IOException {
+ StorageContainerLocationProtocol client =
mock(StorageContainerLocationProtocol.class);
+ when(client.getPeerUpgradeStatus()).thenReturn(version);
+ return client;
+ }
+
+ private SCMNodeDetails peerNode(String nodeId) {
+ SCMNodeDetails node = mock(SCMNodeDetails.class);
+ when(node.getNodeId()).thenReturn(nodeId);
+ return node;
+ }
+
private List<String> getNetworkNames() {
return nodeManager.getAllNodes().stream()
.map(NodeImpl::getNetworkName)
diff --git
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestSCMClientProtocolServer.java
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestSCMClientProtocolServer.java
index 564f9433b85..8aa9d5446c7 100644
---
a/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestSCMClientProtocolServer.java
+++
b/hadoop-hdds/server-scm/src/test/java/org/apache/hadoop/hdds/scm/server/TestSCMClientProtocolServer.java
@@ -24,9 +24,7 @@
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
-import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.mockStatic;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
@@ -39,10 +37,8 @@
import java.time.ZoneOffset;
import java.util.ArrayList;
import java.util.Arrays;
-import java.util.Collections;
import java.util.HashSet;
import java.util.List;
-import org.apache.hadoop.hdds.ComponentVersion;
import org.apache.hadoop.hdds.HDDSVersion;
import org.apache.hadoop.hdds.client.RatisReplicationConfig;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
@@ -54,34 +50,24 @@
import org.apache.hadoop.hdds.scm.HddsTestUtils;
import org.apache.hadoop.hdds.scm.container.ContainerInfo;
import org.apache.hadoop.hdds.scm.container.ContainerManagerImpl;
-import org.apache.hadoop.hdds.scm.container.MockNodeManager;
-import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.hdds.scm.ha.SCMContext;
import org.apache.hadoop.hdds.scm.ha.SCMHAManagerStub;
-import org.apache.hadoop.hdds.scm.ha.SCMHANodeDetails;
import org.apache.hadoop.hdds.scm.ha.SCMNodeDetails;
-import org.apache.hadoop.hdds.scm.node.NodeManager;
import org.apache.hadoop.hdds.scm.pipeline.PipelineID;
-import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
import
org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocolServerSideTranslatorPB;
import org.apache.hadoop.hdds.scm.safemode.SCMSafeModeManager;
-import org.apache.hadoop.hdds.scm.safemode.SCMSafeModeManager.SafeModeStatus;
import org.apache.hadoop.hdds.scm.server.upgrade.FinalizationManager;
import org.apache.hadoop.hdds.scm.server.upgrade.ScmVersionManager;
-import org.apache.hadoop.hdds.utils.HAUtils;
import org.apache.hadoop.hdds.utils.ProtocolMessageMetrics;
-import org.apache.hadoop.ozone.audit.SCMAction;
import org.apache.hadoop.ozone.container.common.SCMTestUtils;
import org.apache.hadoop.ozone.upgrade.UpgradeFinalization.StatusAndMessages;
import org.apache.hadoop.security.AccessControlException;
import org.apache.hadoop.security.UserGroupInformation;
-import org.apache.ozone.test.GenericTestUtils.LogCapturer;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.BeforeAll;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.io.TempDir;
-import org.mockito.MockedStatic;
/**
* Unit tests to validate the SCMClientProtocolServer
@@ -294,256 +280,6 @@ public void testGetPeerUpgradeStatusReturnsLocalVersion()
throws IOException {
assertEquals(HDDSVersion.SOFTWARE_VERSION, server.getPeerUpgradeStatus());
}
- @Test
- public void testFinalizeProceedsWhenNoPeers() throws IOException {
- FinalizationManager finalizationManager = mock(FinalizationManager.class);
- try (SCMClientProtocolServer testServer =
- peerCheckServer(finalizationManager, Collections.emptyList())) {
- testServer.finalizeUpgrade();
- verify(finalizationManager).finalizeUpgrade();
- }
- }
-
- @Test
- public void testFinalizeProceedsWhenAllPeersMatch() throws IOException {
- FinalizationManager finalizationManager = mock(FinalizationManager.class);
- StorageContainerLocationProtocol matching =
peerClient(HDDSVersion.SOFTWARE_VERSION);
-
- try (SCMClientProtocolServer testServer =
- peerCheckServer(finalizationManager,
Arrays.asList(peerNode("scm2"), peerNode("scm3")));
- MockedStatic<HAUtils> haUtils = mockStatic(HAUtils.class)) {
- haUtils.when(() -> HAUtils.getScmContainerClientForNode(any(), any(),
any())).thenReturn(matching);
- testServer.finalizeUpgrade();
- }
- verify(finalizationManager).finalizeUpgrade();
- }
-
- @Test
- public void testFinalizeRejectsOlderPeerUnlessForced() throws IOException {
- FinalizationManager finalizationManager = mock(FinalizationManager.class);
- StorageContainerLocationProtocol matching =
peerClient(HDDSVersion.SOFTWARE_VERSION);
- StorageContainerLocationProtocol older =
peerClient(HDDSVersion.DEFAULT_VERSION);
-
- try (SCMClientProtocolServer testServer =
- peerCheckServer(finalizationManager,
Arrays.asList(peerNode("scm2"), peerNode("scm3")));
- MockedStatic<HAUtils> haUtils = mockStatic(HAUtils.class)) {
- haUtils.when(() -> HAUtils.getScmContainerClientForNode(any(), any(),
any())).thenReturn(matching, older);
- // A peer on an older version is rejected without force.
- assertThrows(SCMException.class, testServer::finalizeUpgrade);
- verify(finalizationManager, never()).finalizeUpgrade();
- // With force the peer version check is skipped and finalization
proceeds.
- testServer.forceFinalizeUpgrade();
- }
- verify(finalizationManager).finalizeUpgrade();
- }
-
- @Test
- public void testFinalizeRejectsUnknownFuturePeerUnlessForced() throws
IOException {
- FinalizationManager finalizationManager = mock(FinalizationManager.class);
- StorageContainerLocationProtocol matching =
peerClient(HDDSVersion.SOFTWARE_VERSION);
- // A version not recognized by this binary deserializes to UNKNOWN_VERSION
in the client translator.
- StorageContainerLocationProtocol unknown =
peerClient(HDDSVersion.UNKNOWN_VERSION);
-
- try (SCMClientProtocolServer testServer =
- peerCheckServer(finalizationManager,
Arrays.asList(peerNode("scm2"), peerNode("scm3")));
- MockedStatic<HAUtils> haUtils = mockStatic(HAUtils.class)) {
- haUtils.when(() -> HAUtils.getScmContainerClientForNode(any(), any(),
any()))
- .thenReturn(matching, unknown);
- // A peer on an unrecognized future version is rejected without force.
- assertThrows(SCMException.class, testServer::finalizeUpgrade);
- verify(finalizationManager, never()).finalizeUpgrade();
- // With force the peer version check is skipped and finalization
proceeds.
- testServer.forceFinalizeUpgrade();
- }
- verify(finalizationManager).finalizeUpgrade();
- }
-
- @Test
- public void testFinalizeRejectsUnreachablePeerUnlessForced() throws
IOException {
- FinalizationManager finalizationManager = mock(FinalizationManager.class);
- StorageContainerLocationProtocol unreachable =
mock(StorageContainerLocationProtocol.class);
- when(unreachable.getPeerUpgradeStatus()).thenThrow(new
IOException("connection refused"));
-
- try (SCMClientProtocolServer testServer =
- peerCheckServer(finalizationManager,
Collections.singletonList(peerNode("scm2")));
- MockedStatic<HAUtils> haUtils = mockStatic(HAUtils.class)) {
- haUtils.when(() -> HAUtils.getScmContainerClientForNode(any(), any(),
any())).thenReturn(unreachable);
- // An unreachable peer is rejected without force.
- assertThrows(SCMException.class, testServer::finalizeUpgrade);
- verify(finalizationManager, never()).finalizeUpgrade();
- // With force the peer version check is skipped and finalization
proceeds.
- testServer.forceFinalizeUpgrade();
- }
- verify(finalizationManager).finalizeUpgrade();
- }
-
- @Test
- public void testFinalizeProceedsWhenAllDatanodesMatchScmVersion() throws
IOException {
- FinalizationManager finalizationManager = mock(FinalizationManager.class);
- NodeManager.DatanodeFinalizationCounts datanodeCounts =
NodeManager.DatanodeFinalizationCounts.newBuilder()
- .setNumFinalizedDatanodes(3)
- .setTotalHealthyDatanodes(3)
- .setMinApparentVersion(HDDSVersion.SOFTWARE_VERSION.serialize())
- .setMaxApparentVersion(HDDSVersion.SOFTWARE_VERSION.serialize())
- .setAllSoftwareVersionsMatchScm(true)
- .build();
- try (SCMClientProtocolServer testServer =
peerCheckServer(finalizationManager,
- Collections.emptyList(), datanodeCounts)) {
- testServer.finalizeUpgrade();
- verify(finalizationManager).finalizeUpgrade();
- }
- }
-
- @Test
- public void testFinalizeRejectsDatanodeWithMismatchedVersionUnlessForced()
throws IOException {
- FinalizationManager finalizationManager = mock(FinalizationManager.class);
- NodeManager.DatanodeFinalizationCounts datanodeCounts =
NodeManager.DatanodeFinalizationCounts.newBuilder()
- .setNumFinalizedDatanodes(3)
- .setTotalHealthyDatanodes(3)
- .setMinApparentVersion(HDDSVersion.DEFAULT_VERSION.serialize())
- .setMaxApparentVersion(HDDSVersion.SOFTWARE_VERSION.serialize())
- .setAllSoftwareVersionsMatchScm(false)
- .build();
- try (SCMClientProtocolServer testServer =
peerCheckServer(finalizationManager,
- Collections.emptyList(), datanodeCounts)) {
- // A datanode on a mismatched version is rejected without force.
- assertThrows(SCMException.class, testServer::finalizeUpgrade);
- verify(finalizationManager, never()).finalizeUpgrade();
- // With force the datanode version check is skipped and finalization
proceeds.
- testServer.forceFinalizeUpgrade();
- verify(finalizationManager).finalizeUpgrade();
- }
- }
-
- @Test
- public void testForceFinalizeAuditRecordsForceFlag() throws IOException {
- FinalizationManager finalizationManager = mock(FinalizationManager.class);
- // Drive the failure audit path (logged at ERROR, captured by the default
log4j2 config) so the
- // recorded audit map can be inspected.
- doThrow(new
RuntimeException("test")).when(finalizationManager).finalizeUpgrade();
-
- LogCapturer auditLog = LogCapturer.log4j2("SCMAudit");
- try (SCMClientProtocolServer testServer =
peerCheckServer(finalizationManager, Collections.emptyList())) {
- assertThrows(RuntimeException.class, testServer::forceFinalizeUpgrade);
- } finally {
- auditLog.stopCapturing();
- }
-
- String output = auditLog.getOutput();
- assertTrue(output.contains(SCMAction.FINALIZE_SCM_UPGRADE.getAction()),
- "audit log should record the finalize action: " + output);
- assertTrue(output.contains("\"force\":\"true\""),
- "audit log should record that force was passed: " + output);
- }
-
- private SCMClientProtocolServer peerCheckServer(
- FinalizationManager finalizationManager, List<SCMNodeDetails> peers)
throws IOException {
- // Default to all datanode versions matching SCM so the SCM peer checks
are exercised in isolation.
- NodeManager.DatanodeFinalizationCounts datanodeCounts =
NodeManager.DatanodeFinalizationCounts.newBuilder()
- .setNumFinalizedDatanodes(0)
- .setTotalHealthyDatanodes(0)
- .setMinApparentVersion(0)
- .setMaxApparentVersion(0)
- .setAllSoftwareVersionsMatchScm(true)
- .build();
- return peerCheckServer(finalizationManager, peers, datanodeCounts);
- }
-
- private SCMClientProtocolServer peerCheckServer(
- FinalizationManager finalizationManager, List<SCMNodeDetails> peers,
- NodeManager.DatanodeFinalizationCounts datanodeCounts) throws
IOException {
-
- StorageContainerManager mockScm = mockStorageContainerManager();
- when(mockScm.getFinalizationManager()).thenReturn(finalizationManager);
- when(mockScm.getConfiguration()).thenReturn(new OzoneConfiguration());
-
- SCMHANodeDetails haNodeDetails = mock(SCMHANodeDetails.class);
- when(haNodeDetails.getPeerNodeDetails()).thenReturn(peers);
- when(mockScm.getSCMHANodeDetails()).thenReturn(haNodeDetails);
-
- NodeManager nodeManager = mock(NodeManager.class);
-
when(nodeManager.getDatanodeFinalizationCounts()).thenReturn(datanodeCounts);
- when(mockScm.getScmNodeManager()).thenReturn(nodeManager);
- return new SCMClientProtocolServer(new OzoneConfiguration(), mockScm,
mock(ReconfigurationHandler.class));
- }
-
- private StorageContainerLocationProtocol peerClient(HDDSVersion version)
throws IOException {
- StorageContainerLocationProtocol client =
mock(StorageContainerLocationProtocol.class);
- when(client.getPeerUpgradeStatus()).thenReturn(version);
- return client;
- }
-
- private SCMNodeDetails peerNode(String nodeId) {
- SCMNodeDetails node = mock(SCMNodeDetails.class);
- when(node.getNodeId()).thenReturn(nodeId);
- return node;
- }
-
- @Test
- public void testQueryUpgradeStatus() throws Exception {
- HddsProtos.UpgradeStatus status = server.queryUpgradeStatus();
-
- // SCM starts already finalized in tests
- assertEquals(HddsProtos.FinalizationStatus.FINALIZED,
status.getScmFinalizationStatus());
- // No datanodes registered
- assertEquals(0, status.getNumDatanodesFinalized());
- assertEquals(0, status.getNumDatanodesTotal());
- assertEquals(HddsProtos.FinalizationStatus.FINALIZED,
status.getHddsFinalizationStatus());
- }
-
- @Test
- public void testQueryUpgradeStatusHddsInProgress() throws Exception {
- // SCM is finalized but not all datanodes are, so HDDS finalization is
still in progress.
- ScmVersionManager mockVersionManager = mock(ScmVersionManager.class);
- when(mockVersionManager.needsFinalization()).thenReturn(false);
- ComponentVersion apparentVersion = mock(ComponentVersion.class);
- when(apparentVersion.serialize()).thenReturn(0);
- when(mockVersionManager.getApparentVersion()).thenReturn(apparentVersion);
-
- NodeManager mockNodeManager = new MockNodeManager(false, 0) {
- @Override
- public DatanodeFinalizationCounts getDatanodeFinalizationCounts() {
- return DatanodeFinalizationCounts.newBuilder()
- .setNumFinalizedDatanodes(1)
- .setTotalHealthyDatanodes(3)
- .build();
- }
- };
-
- StorageContainerManager mockScm = mockStorageContainerManager();
- when(mockScm.getVersionManager()).thenReturn(mockVersionManager);
- when(mockScm.getScmNodeManager()).thenReturn(mockNodeManager);
- when(mockScm.getScmContext()).thenReturn(SCMContext.emptyContext());
-
- SCMClientProtocolServer testServer = new SCMClientProtocolServer(
- new OzoneConfiguration(), mockScm, mock(ReconfigurationHandler.class));
- try {
- HddsProtos.UpgradeStatus status = testServer.queryUpgradeStatus();
- assertEquals(HddsProtos.FinalizationStatus.FINALIZED,
status.getScmFinalizationStatus());
- assertEquals(HddsProtos.FinalizationStatus.IN_PROGRESS,
status.getHddsFinalizationStatus());
- assertEquals(1, status.getNumDatanodesFinalized());
- assertEquals(3, status.getNumDatanodesTotal());
- } finally {
- testServer.stop();
- }
- }
-
- @Test
- public void testQueryUpgradeStatusInSafemode() {
- // Put SCM into safe mode via the context the server consults.
- scm.getScmContext().updateSafeModeStatus(SafeModeStatus.INITIAL);
- try {
- assertTrue(scm.getScmContext().isInSafeMode());
-
- // Querying upgrade status is blocked while SCM is in safe mode.
- SCMException ex = assertThrows(SCMException.class, () ->
server.queryUpgradeStatus());
- assertEquals(SCMException.ResultCodes.SAFE_MODE_EXCEPTION,
ex.getResult());
- } finally {
- // Restore for other tests sharing the static SCM instance.
-
scm.getScmContext().updateSafeModeStatus(SafeModeStatus.OUT_OF_SAFE_MODE);
- }
- }
-
private ContainerInfo newContainerWithLastUsedTime(long containerId,
Instant fixedLastUsedInstant) {
return new ContainerInfo.Builder()
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 4edfc16c852..ed5b0b500ea 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
@@ -614,11 +614,6 @@ public StatusAndMessages queryUpgradeFinalizationProgress(
upgradeClientID, force, readonly);
}
- @Override
- public HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException {
- return storageContainerLocationClient.queryUpgradeStatus();
- }
-
@Override
public DecommissionScmResponseProto decommissionScm(
String scmId)
diff --git a/hadoop-ozone/dist/src/main/compose/common/security.conf
b/hadoop-ozone/dist/src/main/compose/common/security.conf
index ef989c2490a..ab3e7101407 100644
--- a/hadoop-ozone/dist/src/main/compose/common/security.conf
+++ b/hadoop-ozone/dist/src/main/compose/common/security.conf
@@ -50,7 +50,7 @@ OZONE-SITE.XML_hdds.grpc.tls.enabled=true
OZONE-SITE.XML_ozone.security.enabled=true
OZONE-SITE.XML_ozone.acl.enabled=true
OZONE-SITE.XML_ozone.acl.authorizer.class=org.apache.hadoop.ozone.security.acl.OzoneNativeAuthorizer
-OZONE-SITE.XML_ozone.administrators="testuser,recon,om"
+OZONE-SITE.XML_ozone.administrators="testuser,recon"
OZONE-SITE.XML_ozone.s3.administrators="testuser,s3g"
OZONE-SITE.XML_ozone.security.http.kerberos.enabled=true
OZONE-SITE.XML_ozone.s3g.secret.http.enabled=true
@@ -95,7 +95,7 @@ CORE-SITE.XML_hadoop.security.authorization=true
HADOOP-POLICY.XML_ozone.om.security.client.protocol.acl=*
HADOOP-POLICY.XML_hdds.security.client.datanode.container.protocol.acl=*
HADOOP-POLICY.XML_hdds.security.client.scm.container.protocol.acl=*
-HADOOP-POLICY.XML_hdds.security.client.scm.block.protocol.acl=*
+HADOOP-POLICY.XML_hdds.security.client.scm.block.protocol.acl=om,scm
HADOOP-POLICY.XML_hdds.security.client.scm.certificate.protocol.acl=*
HTTPFS-SITE.XML_hadoop.http.authentication.type=kerberos
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/HddsUpgradeTestUtils.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/HddsUpgradeTestUtils.java
index 7fcc4cd53a5..e39a68c8010 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/HddsUpgradeTestUtils.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/HddsUpgradeTestUtils.java
@@ -25,7 +25,7 @@
import java.util.List;
import java.util.concurrent.TimeoutException;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
-import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
+import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
import org.apache.hadoop.hdds.utils.db.CodecException;
import org.apache.hadoop.hdds.utils.db.RocksDatabaseException;
@@ -47,7 +47,7 @@ public final class HddsUpgradeTestUtils {
private HddsUpgradeTestUtils() { }
- public static void
waitForFinalizationFromClient(StorageContainerLocationProtocol scmClient)
throws Exception {
+ public static void waitForFinalizationFromClient(ScmBlockLocationProtocol
scmClient) throws Exception {
LambdaTestUtils.await(60_000, 1_000, () -> {
HddsProtos.UpgradeStatus status = scmClient.queryUpgradeStatus();
LOG.info("Waiting for upgrade finalization to complete from client.
Current status is:\n{}", status);
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestDNDataDistributionFinalization.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestDNDataDistributionFinalization.java
index 48d159de691..78a6d1c9448 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestDNDataDistributionFinalization.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestDNDataDistributionFinalization.java
@@ -31,8 +31,9 @@
import java.util.concurrent.TimeUnit;
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.scm.ScmConfig;
-import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
+import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol;
import org.apache.hadoop.hdds.scm.server.SCMStorageConfig;
+import org.apache.hadoop.hdds.utils.HAUtils;
import org.apache.hadoop.ozone.HddsDatanodeService;
import org.apache.hadoop.ozone.MiniOzoneCluster;
import org.apache.hadoop.ozone.MiniOzoneHAClusterImpl;
@@ -57,7 +58,7 @@
* Tests upgrade finalization failure scenarios and corner cases specific to
DN data distribution feature.
*/
public class TestDNDataDistributionFinalization {
- private StorageContainerLocationProtocol scmClient;
+ private ScmBlockLocationProtocol scmBlockClient;
private MiniOzoneHAClusterImpl cluster;
private static final int NUM_DATANODES = 3;
@@ -104,7 +105,7 @@ public void init(OzoneConfiguration conf) throws Exception {
.build());
this.cluster = clusterBuilder.build();
- scmClient = cluster.getStorageContainerLocationClient();
+ scmBlockClient = HAUtils.getScmBlockClient(conf);
cluster.waitForClusterToBeReady();
assertEquals(HDDSLayoutFeature.HBASE_SUPPORT,
cluster.getStorageContainerManager().getVersionManager().getApparentVersion());
@@ -155,8 +156,8 @@ public void testDataDistributionUpgradeScenario() throws
Exception {
validatePreDataDistributionFeatureState();
// Wait for finalization to complete
- scmClient.finalizeUpgrade();
- HddsUpgradeTestUtils.waitForFinalizationFromClient(scmClient);
+ scmBlockClient.finalizeUpgrade();
+ HddsUpgradeTestUtils.waitForFinalizationFromClient(scmBlockClient);
// Verify finalization completed
assertFalse(cluster.getStorageContainerManager().getVersionManager().needsFinalization());
@@ -191,8 +192,8 @@ public void testMissingPendingDeleteMetadataRecalculation()
throws Exception {
out.write(data);
}
bucket.deleteKey(keyName);
- scmClient.finalizeUpgrade();
- HddsUpgradeTestUtils.waitForFinalizationFromClient(scmClient);
+ scmBlockClient.finalizeUpgrade();
+ HddsUpgradeTestUtils.waitForFinalizationFromClient(scmBlockClient);
assertFalse(cluster.getStorageContainerManager().getVersionManager().needsFinalization());
assertTrue(VersionedDatanodeFeatures.isFinalized(HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION));
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestScmDataDistributionFinalization.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestScmDataDistributionFinalization.java
index d92dd4030fb..56b07e029db 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestScmDataDistributionFinalization.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestScmDataDistributionFinalization.java
@@ -56,10 +56,11 @@
import
org.apache.hadoop.hdds.scm.block.SCMDeletedBlockTransactionStatusManager;
import org.apache.hadoop.hdds.scm.ha.SCMHADBTransactionBuffer;
import org.apache.hadoop.hdds.scm.metadata.DBTransactionBuffer;
-import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
+import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol;
import org.apache.hadoop.hdds.scm.server.SCMConfigurator;
import org.apache.hadoop.hdds.scm.server.SCMStorageConfig;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
+import org.apache.hadoop.hdds.utils.HAUtils;
import org.apache.hadoop.hdds.utils.db.CodecException;
import org.apache.hadoop.hdds.utils.db.RocksDatabaseException;
import org.apache.hadoop.hdds.utils.db.Table;
@@ -86,7 +87,7 @@
* Tests upgrade finalization failure scenarios and corner cases specific to
SCM data distribution feature.
*/
public class TestScmDataDistributionFinalization {
- private StorageContainerLocationProtocol scmClient;
+ private ScmBlockLocationProtocol scmBlockClient;
private MiniOzoneHAClusterImpl cluster;
private static final int NUM_DATANODES = 3;
private static final int NUM_SCMS = 3;
@@ -133,7 +134,7 @@ public void init(OzoneConfiguration conf) throws Exception {
.build());
this.cluster = clusterBuilder.build();
- scmClient = cluster.getStorageContainerLocationClient();
+ scmBlockClient = HAUtils.getScmBlockClient(conf);
cluster.waitForClusterToBeReady();
assertEquals(HDDSLayoutFeature.HBASE_SUPPORT,
cluster.getStorageContainerManager().getVersionManager().getApparentVersion());
@@ -165,8 +166,8 @@ public void testFinalizationEmptyClusterDataDistribution()
throws Exception {
init(new OzoneConfiguration());
assertEquals(EMPTY_SUMMARY,
cluster.getStorageContainerLocationClient().getDeletedBlockSummary());
- scmClient.finalizeUpgrade();
- HddsUpgradeTestUtils.waitForFinalizationFromClient(scmClient);
+ scmBlockClient.finalizeUpgrade();
+ HddsUpgradeTestUtils.waitForFinalizationFromClient(scmBlockClient);
// Make sure old leader has caught up and all SCMs have finalized.
waitForScmsToFinalize(cluster.getStorageContainerManagersList());
assertTrue(VersionedDatanodeFeatures.isFinalized(HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION));
@@ -269,8 +270,8 @@ public void
testFinalizationNonEmptyClusterDataDistribution() throws Exception {
flushDBTransactionBuffer(activeSCM);
assertEquals(EMPTY_SUMMARY,
cluster.getStorageContainerLocationClient().getDeletedBlockSummary());
- scmClient.finalizeUpgrade();
- HddsUpgradeTestUtils.waitForFinalizationFromClient(scmClient);
+ scmBlockClient.finalizeUpgrade();
+ HddsUpgradeTestUtils.waitForFinalizationFromClient(scmBlockClient);
// Make sure old leader has caught up and all SCMs have finalized.
waitForScmsToFinalize(cluster.getStorageContainerManagersList());
assertTrue(VersionedDatanodeFeatures.isFinalized(HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION));
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestScmHAFinalization.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestScmHAFinalization.java
index b77741519f3..5d72b3d91e5 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestScmHAFinalization.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/hdds/upgrade/TestScmHAFinalization.java
@@ -28,10 +28,12 @@
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.scm.ScmConfigKeys;
import
org.apache.hadoop.hdds.scm.container.common.helpers.ContainerWithPipeline;
+import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol;
import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
import org.apache.hadoop.hdds.scm.server.SCMConfigurator;
import org.apache.hadoop.hdds.scm.server.SCMStorageConfig;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
+import org.apache.hadoop.hdds.utils.HAUtils;
import org.apache.hadoop.ozone.HddsDatanodeService;
import org.apache.hadoop.ozone.MiniOzoneCluster;
import org.apache.hadoop.ozone.MiniOzoneHAClusterImpl;
@@ -55,6 +57,7 @@ public class TestScmHAFinalization {
LoggerFactory.getLogger(TestScmHAFinalization.class);
private StorageContainerLocationProtocol scmClient;
+ private ScmBlockLocationProtocol scmBlockClient;
private MiniOzoneHAClusterImpl cluster;
private static final int NUM_DATANODES = 3;
private static final int NUM_SCMS = 3;
@@ -80,6 +83,7 @@ public void init(OzoneConfiguration conf, int
numInactiveSCMs) throws Exception
this.cluster = clusterBuilder.build();
scmClient = cluster.getStorageContainerLocationClient();
+ scmBlockClient = HAUtils.getScmBlockClient(cluster.getConf());
cluster.waitForClusterToBeReady();
}
@@ -129,8 +133,8 @@ public void
testFinalizedDatanodesShutDownWithPrefinalizedScm() throws Exception
public void testFinalization() throws Exception {
OzoneConfiguration conf = new OzoneConfiguration();
init(conf, 0);
- scmClient.finalizeUpgrade();
- HddsUpgradeTestUtils.waitForFinalizationFromClient(scmClient);
+ scmBlockClient.finalizeUpgrade();
+ HddsUpgradeTestUtils.waitForFinalizationFromClient(scmBlockClient);
// Ensure all SCMs finalize, indicating the message has been propagated
across them all
waitForScmsToFinalize(cluster.getStorageContainerManagersList());
@@ -168,8 +172,8 @@ public void testSnapshotFinalization() throws Exception {
// Wait for finalization from the client perspective.
// Force finalize skips the peer version checks that require all peers to
be active.
- scmClient.forceFinalizeUpgrade();
- HddsUpgradeTestUtils.waitForFinalizationFromClient(scmClient);
+ scmBlockClient.forceFinalizeUpgrade();
+ HddsUpgradeTestUtils.waitForFinalizationFromClient(scmBlockClient);
// Wait for two running SCMs to finish finalization.
waitForScmsToFinalize(activeScms);
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestBlockDeletionService.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestBlockDeletionService.java
index 8f7a06bb146..c83c443dc92 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestBlockDeletionService.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/service/TestBlockDeletionService.java
@@ -41,11 +41,12 @@
import org.apache.hadoop.hdds.conf.OzoneConfiguration;
import org.apache.hadoop.hdds.scm.block.BlockManager;
import
org.apache.hadoop.hdds.scm.container.placement.metrics.SCMPerformanceMetrics;
-import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
+import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol;
import org.apache.hadoop.hdds.scm.server.SCMStorageConfig;
import org.apache.hadoop.hdds.scm.server.StorageContainerManager;
import org.apache.hadoop.hdds.upgrade.HDDSLayoutFeature;
import org.apache.hadoop.hdds.upgrade.HddsUpgradeTestUtils;
+import org.apache.hadoop.hdds.utils.HAUtils;
import org.apache.hadoop.ozone.MiniOzoneCluster;
import org.apache.hadoop.ozone.UniformDatanodesFactory;
import org.apache.hadoop.ozone.client.OzoneBucket;
@@ -72,7 +73,7 @@ public class TestBlockDeletionService {
private static final String BUCKET_NAME = "bucket1";
private static final int KEY_SIZE = 5 * 1024; // 5 KB
private static MiniOzoneCluster cluster;
- private static StorageContainerLocationProtocol scmClient;
+ private static ScmBlockLocationProtocol scmBlockClient;
private static OzoneBucket bucket;
private static SCMPerformanceMetrics metrics;
@@ -99,7 +100,7 @@ public static void init() throws Exception {
.setApparentVersion(HBASE_SUPPORT).build())
.build();
cluster.waitForClusterToBeReady();
- scmClient = cluster.getStorageContainerLocationClient();
+ scmBlockClient = HAUtils.getScmBlockClient(conf);
assertEquals(HBASE_SUPPORT,
cluster.getStorageContainerManager().getVersionManager().getApparentVersion());
metrics =
cluster.getStorageContainerManager().getBlockProtocolServer().getMetrics();
@@ -139,8 +140,8 @@ public void testDeleteKeyQuotaWithUpgrade() throws
Exception {
GenericTestUtils.waitFor(() -> metrics.getDeleteKeyFailedBlocks() -
initialFailedBlocks == 0, 50, 1000);
// Step 5: wait for finalizing upgrade
- scmClient.finalizeUpgrade();
- HddsUpgradeTestUtils.waitForFinalizationFromClient(scmClient);
+ scmBlockClient.finalizeUpgrade();
+ HddsUpgradeTestUtils.waitForFinalizationFromClient(scmBlockClient);
assertTrue(VersionedDatanodeFeatures.isFinalized(HDDSLayoutFeature.STORAGE_SPACE_DISTRIBUTION));
// POST-UPGRADE
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java
index e9b1a92f832..011c0dcd25a 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/OzoneManager.java
@@ -3780,7 +3780,7 @@ public void forceFinalizeUpgrade() throws IOException {
public QueryUpgradeStatusResponse queryUpgradeStatus() throws IOException {
HddsProtos.UpgradeStatus scmStatus;
try {
- scmStatus = scmClient.getContainerClient().queryUpgradeStatus();
+ scmStatus = scmClient.getBlockClient().queryUpgradeStatus();
} catch (SCMException e) {
// SCM refuses the query while in safe mode
if (e.getResult() == SCMException.ResultCodes.SAFE_MODE_EXCEPTION) {
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/upgrade/OMFinalizeUpgradeRequestBase.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/upgrade/OMFinalizeUpgradeRequestBase.java
index 4467a74905c..c085342cdc3 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/upgrade/OMFinalizeUpgradeRequestBase.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/upgrade/OMFinalizeUpgradeRequestBase.java
@@ -80,9 +80,9 @@ public OMRequest preExecute(OzoneManager ozoneManager) throws
IOException {
try {
if (force) {
-
ozoneManager.getScmClient().getContainerClient().forceFinalizeUpgrade();
+ ozoneManager.getScmClient().getBlockClient().forceFinalizeUpgrade();
} else {
- ozoneManager.getScmClient().getContainerClient().finalizeUpgrade();
+ ozoneManager.getScmClient().getBlockClient().finalizeUpgrade();
}
} catch (SCMException e) {
if (e.getResult() == SCMException.ResultCodes.UNSUPPORTED_OPERATION) {
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/upgrade/OMUpgradeFinalizeService.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/upgrade/OMUpgradeFinalizeService.java
index b168d14c111..042bf92401f 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/upgrade/OMUpgradeFinalizeService.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/upgrade/OMUpgradeFinalizeService.java
@@ -116,7 +116,7 @@ public BackgroundTaskResult call() {
return BackgroundTaskResult.EmptyTaskResult.newResult();
}
- HddsProtos.UpgradeStatus upgradeStatus =
scmClient.getContainerClient().queryUpgradeStatus();
+ HddsProtos.UpgradeStatus upgradeStatus =
scmClient.getBlockClient().queryUpgradeStatus();
if (upgradeStatus.getHddsFinalizationStatus() ==
HddsProtos.FinalizationStatus.FINALIZED) {
LOG.info("The SCM Upgrade has been finalized. OM will now
finalize. Run count {}", run);
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ScmBlockLocationTestingClient.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ScmBlockLocationTestingClient.java
index 823a6405257..fd9a65cdbee 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ScmBlockLocationTestingClient.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/ScmBlockLocationTestingClient.java
@@ -33,6 +33,7 @@
import org.apache.hadoop.hdds.client.ReplicationConfig;
import org.apache.hadoop.hdds.client.StandaloneReplicationConfig;
import org.apache.hadoop.hdds.protocol.DatanodeDetails;
+import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos.ReplicationFactor;
import org.apache.hadoop.hdds.scm.AddSCMRequest;
import org.apache.hadoop.hdds.scm.ScmInfo;
@@ -217,6 +218,19 @@ public int getNumberOfDeletedBlocks() {
return numBlocksDeleted;
}
+ @Override
+ public void finalizeUpgrade() throws IOException {
+ }
+
+ @Override
+ public void forceFinalizeUpgrade() throws IOException {
+ }
+
+ @Override
+ public HddsProtos.UpgradeStatus queryUpgradeStatus() throws IOException {
+ return HddsProtos.UpgradeStatus.getDefaultInstance();
+ }
+
@Override
public void close() throws IOException {
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequest.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequest.java
index 376c80ce7f5..0ff5764f443 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequest.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequest.java
@@ -56,13 +56,13 @@ public void
testForcePreExecuteCallsScmForceFinalizeUpgrade() throws IOException
request.preExecute(ozoneManager);
// A forced request must route to SCM's force path so SCM skips its own
version checks.
- verify(scmContainerLocationProtocol).forceFinalizeUpgrade();
- verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ verify(scmBlockLocationProtocol).forceFinalizeUpgrade();
+ verify(scmBlockLocationProtocol, never()).finalizeUpgrade();
}
@Test
public void testAuditMapRecordsForceFlag() throws IOException {
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
ExecutionContext context = ExecutionContext.of(1, TermIndex.INITIAL_VALUE);
OMStartFinalizeUpgradeRequest forced = new
OMStartFinalizeUpgradeRequest(buildRequest(true));
@@ -78,7 +78,7 @@ public void testAuditMapRecordsForceFlag() throws IOException
{
@Test
public void testForceSkipsPeerVersionCheckForUnreachablePeer() throws
IOException {
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
when(ozoneManager.getPeerNodes()).thenReturn(Collections.singletonList(buildPeer("om2")));
OMAdminProtocolClientSideImpl unreachableClient =
mock(OMAdminProtocolClientSideImpl.class);
when(unreachableClient.getPeerUpgradeStatus()).thenThrow(new
IOException("connection refused"));
@@ -94,13 +94,13 @@ public void
testForceSkipsPeerVersionCheckForUnreachablePeer() throws IOExceptio
}
verify(unreachableClient, never()).getPeerUpgradeStatus();
- verify(scmContainerLocationProtocol).forceFinalizeUpgrade();
- verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ verify(scmBlockLocationProtocol).forceFinalizeUpgrade();
+ verify(scmBlockLocationProtocol, never()).finalizeUpgrade();
}
@Test
public void testForceSkipsPeerVersionCheckForMismatchedPeer() throws
IOException {
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
when(ozoneManager.getPeerNodes()).thenReturn(Collections.singletonList(buildPeer("om2")));
OMAdminProtocolClientSideImpl olderClient =
peerClientWithVersion(OzoneManagerVersion.HBASE_SUPPORT);
@@ -115,8 +115,8 @@ public void
testForceSkipsPeerVersionCheckForMismatchedPeer() throws IOException
}
verify(olderClient, never()).getPeerUpgradeStatus();
- verify(scmContainerLocationProtocol).forceFinalizeUpgrade();
- verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ verify(scmBlockLocationProtocol).forceFinalizeUpgrade();
+ verify(scmBlockLocationProtocol, never()).finalizeUpgrade();
}
private OzoneManagerProtocolProtos.OMRequest buildRequest(boolean force) {
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequestBase.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequestBase.java
index 47cb54c6a77..b5aa2668b27 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequestBase.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequestBase.java
@@ -80,7 +80,7 @@ public void mockPreFinalizedOM() {
@Test
public void testPreExecuteCallsScmFinalizeUpgrade() throws IOException {
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
OMFinalizeUpgradeRequestBase request = newRequest();
OMRequest original = request.getOmRequest();
@@ -92,14 +92,14 @@ public void testPreExecuteCallsScmFinalizeUpgrade() throws
IOException {
assertNotNull(modified.getUserInfo());
// A non-forced initiate request must route to SCM's non-force finalize
path.
- verify(scmContainerLocationProtocol).finalizeUpgrade();
- verify(scmContainerLocationProtocol, never()).forceFinalizeUpgrade();
+ verify(scmBlockLocationProtocol).finalizeUpgrade();
+ verify(scmBlockLocationProtocol, never()).forceFinalizeUpgrade();
}
@Test
public void testScmFinalizeFailurePropagatesToClient() throws IOException {
IOException scmFailure = new IOException("SCM finalize upgrade failed");
- doThrow(scmFailure).when(scmContainerLocationProtocol).finalizeUpgrade();
+ doThrow(scmFailure).when(scmBlockLocationProtocol).finalizeUpgrade();
OMFinalizeUpgradeRequestBase request = newRequest();
@@ -108,14 +108,14 @@ public void testScmFinalizeFailurePropagatesToClient()
throws IOException {
IOException ex = assertThrows(IOException.class, () ->
request.preExecute(ozoneManager));
assertSame(scmFailure, ex);
- verify(scmContainerLocationProtocol).finalizeUpgrade();
+ verify(scmBlockLocationProtocol).finalizeUpgrade();
}
@Test
public void testScmUnsupportedOperationBecomesOmNotSupportedOperation()
throws IOException {
SCMException scmFailure =
new SCMException("SCM version mismatch",
SCMException.ResultCodes.UNSUPPORTED_OPERATION);
- doThrow(scmFailure).when(scmContainerLocationProtocol).finalizeUpgrade();
+ doThrow(scmFailure).when(scmBlockLocationProtocol).finalizeUpgrade();
OMFinalizeUpgradeRequestBase request = newRequest();
@@ -126,13 +126,13 @@ public void
testScmUnsupportedOperationBecomesOmNotSupportedOperation() throws I
assertEquals(scmFailure.getMessage(), ex.getMessage());
assertSame(scmFailure, ex.getCause());
- verify(scmContainerLocationProtocol).finalizeUpgrade();
+ verify(scmBlockLocationProtocol).finalizeUpgrade();
}
@Test
public void testOtherScmExceptionPropagatesUnchanged() throws IOException {
SCMException scmFailure = new SCMException("SCM is in safe mode",
SCMException.ResultCodes.SAFE_MODE_EXCEPTION);
- doThrow(scmFailure).when(scmContainerLocationProtocol).finalizeUpgrade();
+ doThrow(scmFailure).when(scmBlockLocationProtocol).finalizeUpgrade();
OMFinalizeUpgradeRequestBase request = newRequest();
@@ -140,7 +140,7 @@ public void testOtherScmExceptionPropagatesUnchanged()
throws IOException {
SCMException ex = assertThrows(SCMException.class, () ->
request.preExecute(ozoneManager));
assertSame(scmFailure, ex);
- verify(scmContainerLocationProtocol).finalizeUpgrade();
+ verify(scmBlockLocationProtocol).finalizeUpgrade();
}
@Test
@@ -162,23 +162,23 @@ public void testAccessDeniedWhenUserIsNotAdmin() throws
IOException {
"non-admin user should receive ACCESS_DENIED from preExecute");
// SCM must NOT have been called — auth is checked before the SCM call.
- verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ verify(scmBlockLocationProtocol, never()).finalizeUpgrade();
}
@Test
public void testPeerVersionCheckPassesWhenNoPeers() throws IOException {
assertTrue(ozoneManager.getPeerNodes().isEmpty());
// preExecute must complete normally and call SCM finalize.
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
newRequest().preExecute(ozoneManager);
- verify(scmContainerLocationProtocol).finalizeUpgrade();
+ verify(scmBlockLocationProtocol).finalizeUpgrade();
}
@Test
public void testPeerVersionCheckPassesWhenAllPeersMatch() throws IOException
{
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
when(ozoneManager.getPeerNodes()).thenReturn(Arrays.asList(buildPeer("om2"),
buildPeer("om3")));
OMAdminProtocolClientSideImpl matchingClient =
peerClientWithVersion(OzoneManagerVersion.SOFTWARE_VERSION);
@@ -190,7 +190,7 @@ public void testPeerVersionCheckPassesWhenAllPeersMatch()
throws IOException {
newRequest().preExecute(ozoneManager);
}
- verify(scmContainerLocationProtocol).finalizeUpgrade();
+ verify(scmBlockLocationProtocol).finalizeUpgrade();
}
@Test
@@ -208,7 +208,7 @@ public void testPeerVersionCheckRejectsOneOlderPeer()
throws IOException {
assertEquals(OMException.ResultCodes.NOT_SUPPORTED_OPERATION,
ex.getResult());
}
- verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ verify(scmBlockLocationProtocol, never()).finalizeUpgrade();
}
@Test
@@ -226,7 +226,7 @@ public void
testPeerVersionCheckRejectsOneUnknownFuturePeer() throws IOException
assertEquals(OMException.ResultCodes.NOT_SUPPORTED_OPERATION,
ex.getResult());
}
- verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ verify(scmBlockLocationProtocol, never()).finalizeUpgrade();
}
@Test
@@ -244,12 +244,12 @@ public void testPeerVersionCheckRejectsUnreachablePeer()
throws IOException {
assertEquals(OMException.ResultCodes.NOT_SUPPORTED_OPERATION,
ex.getResult());
}
- verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ verify(scmBlockLocationProtocol, never()).finalizeUpgrade();
}
@Test
public void testValidateAndUpdateCacheAddsFinalizationInProgressKey() throws
IOException {
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
assertNull(omMetadataManager.getMetaTable().get(OzoneConsts.FINALIZATION_IN_PROGRESS_KEY),
"key should not exist before the request");
@@ -271,7 +271,7 @@ public void
testValidateAndUpdateCacheAddsFinalizationInProgressKey() throws IOE
@Test
public void testValidateAndUpdateCacheSkipsMarkerWhenAlreadyFinalized()
throws IOException {
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
// Simulate an admin initiating finalize on a cluster that is already
finalized.
OMVersionManager finalizedVersionManager =
OMVersionManagerTestUtils.mockFinalizedOmVersionManager();
when(ozoneManager.getVersionManager()).thenReturn(finalizedVersionManager);
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequestLegacy.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequestLegacy.java
index 28451c6931f..3c4ea8700e2 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequestLegacy.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/request/upgrade/TestOMStartFinalizeUpgradeRequestLegacy.java
@@ -42,7 +42,7 @@ protected OMFinalizeUpgradeRequestBase newRequest() {
@Test
public void testValidateAndUpdateCacheReturnsStartingFinalization() throws
IOException {
- doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ doNothing().when(scmBlockLocationProtocol).finalizeUpgrade();
OMClientResponse response = submitRequest();
OMResponse omResponse = response.getOMResponse();
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/upgrade/TestOMUpgradeFinalizeService.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/upgrade/TestOMUpgradeFinalizeService.java
index e5c7d714529..3ddfa07a389 100644
---
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/upgrade/TestOMUpgradeFinalizeService.java
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/om/upgrade/TestOMUpgradeFinalizeService.java
@@ -32,7 +32,7 @@
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import org.apache.hadoop.hdds.protocol.proto.HddsProtos;
-import org.apache.hadoop.hdds.scm.protocol.StorageContainerLocationProtocol;
+import org.apache.hadoop.hdds.scm.protocol.ScmBlockLocationProtocol;
import org.apache.hadoop.hdds.utils.db.CodecException;
import org.apache.hadoop.hdds.utils.db.RocksDatabaseException;
import org.apache.hadoop.hdds.utils.db.TypedTable;
@@ -63,7 +63,7 @@ public class TestOMUpgradeFinalizeService {
private OMVersionManager versionManager;
private TypedTable<String, String> metaTable;
private ScmClient scmClient;
- private StorageContainerLocationProtocol containerClient;
+ private ScmBlockLocationProtocol blockClient;
private OzoneManagerRatisServer omRatisServer;
private OMUpgradeFinalizeService service;
@@ -88,9 +88,9 @@ void setUp() throws RocksDatabaseException, CodecException {
// OzoneManagerRatisUtils.submitRequest() calls
ozoneManager.getOmRatisServer().submitRequest(...)
when(ozoneManager.getOmRatisServer()).thenReturn(omRatisServer);
- containerClient = mock(StorageContainerLocationProtocol.class);
+ blockClient = mock(ScmBlockLocationProtocol.class);
scmClient = mock(ScmClient.class);
- when(scmClient.getContainerClient()).thenReturn(containerClient);
+ when(scmClient.getBlockClient()).thenReturn(blockClient);
service = new OMUpgradeFinalizeService(ozoneManager, versionManager,
scmClient, INTERVAL_MS);
}
@@ -140,19 +140,19 @@ void
testFinalizationTriggeredWhenScmIsFinalizedAndFinalizationInProgress() thro
.setNumDatanodesFinalized(3)
.setNumDatanodesTotal(3)
.build();
- when(containerClient.queryUpgradeStatus()).thenReturn(scmStatus);
+ when(blockClient.queryUpgradeStatus()).thenReturn(scmStatus);
// Finalization command not given yet
when(metaTable.get(OzoneConsts.FINALIZATION_IN_PROGRESS_KEY)).thenReturn(null);
service.runPeriodicalTaskNow();
- verifyNoInteractions(containerClient);
+ verifyNoInteractions(blockClient);
verifyNoInteractions(omRatisServer);
when(metaTable.get(OzoneConsts.FINALIZATION_IN_PROGRESS_KEY)).thenReturn("ignored");
service.runPeriodicalTaskNow();
- verify(containerClient).queryUpgradeStatus();
+ verify(blockClient).queryUpgradeStatus();
// Implementation submits a FinalizeUpgrade request through Ratis
verify(omRatisServer).submitRequest(any(), any(ClientId.class), anyLong());
}
@@ -173,11 +173,11 @@ void
testFinalizationSkippedWhenScmNotYetFinalized(HddsProtos.FinalizationStatus
.setNumDatanodesFinalized(0)
.setNumDatanodesTotal(3)
.build();
- when(containerClient.queryUpgradeStatus()).thenReturn(scmStatus);
+ when(blockClient.queryUpgradeStatus()).thenReturn(scmStatus);
service.runPeriodicalTaskNow();
- verify(containerClient).queryUpgradeStatus();
+ verify(blockClient).queryUpgradeStatus();
verifyNoInteractions(omRatisServer);
}
@@ -189,12 +189,12 @@ void
testFinalizationSkippedWhenScmNotYetFinalized(HddsProtos.FinalizationStatus
void testExceptionFromScmClientIsHandledGracefully() throws Exception {
when(ozoneManager.isLeaderReady()).thenReturn(true);
when(versionManager.needsFinalization()).thenReturn(true);
- when(containerClient.queryUpgradeStatus()).thenThrow(new IOException("SCM
unavailable"));
+ when(blockClient.queryUpgradeStatus()).thenThrow(new IOException("SCM
unavailable"));
// The catch block in the task swallows the exception.
service.runPeriodicalTaskNow();
- verify(containerClient).queryUpgradeStatus();
+ verify(blockClient).queryUpgradeStatus();
verifyNoInteractions(omRatisServer);
}
@@ -255,7 +255,7 @@ void testExceptionFromRatisSubmitIsHandledGracefully()
throws Exception {
.setNumDatanodesFinalized(3)
.setNumDatanodesTotal(3)
.build();
- when(containerClient.queryUpgradeStatus()).thenReturn(scmStatus);
+ when(blockClient.queryUpgradeStatus()).thenReturn(scmStatus);
when(omRatisServer.submitRequest(any(), any(ClientId.class), anyLong()))
.thenThrow(new ServiceException("Ratis unavailable"));
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]