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 53f472a7c04 HDDS-15875. OM should validate software versions of peers
before accepting finalize command (#10777)
53f472a7c04 is described below
commit 53f472a7c04f58a021625cdc6f060bf629e9d4a7
Author: Ethan Rose <[email protected]>
AuthorDate: Thu Jul 30 19:05:26 2026 -0400
HDDS-15875. OM should validate software versions of peers before accepting
finalize command (#10777)
---
.../ozone/admin/upgrade/FinalizeSubCommand.java | 16 +-
.../admin/upgrade/TestFinalizeSubCommand.java | 13 ++
.../hadoop/ozone/om/exceptions/OMException.java | 4 +
.../hadoop/ozone/om/protocol/OMAdminProtocol.java | 8 +
.../ozone/om/protocol/OzoneManagerProtocol.java | 12 ++
.../protocolPB/OMAdminProtocolClientSideImpl.java | 14 ++
...OzoneManagerProtocolClientSideTranslatorPB.java | 10 +
.../hadoop/ozone/om/TestOMUpgradeFinalization.java | 4 +-
.../src/main/proto/OMAdminProtocol.proto | 14 ++
.../src/main/proto/OmClientProtocol.proto | 14 +-
.../org/apache/hadoop/ozone/om/OzoneManager.java | 11 +-
.../upgrade/OMStartFinalizeUpgradeRequest.java | 62 +++++-
.../protocolPB/OMAdminProtocolServerSideImpl.java | 11 +
.../upgrade/TestOMStartFinalizeUpgradeRequest.java | 230 ++++++++++++++++++++-
.../TestOMAdminProtocolServerSideImpl.java | 45 ++++
15 files changed, 453 insertions(+), 15 deletions(-)
diff --git
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/ozone/admin/upgrade/FinalizeSubCommand.java
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/ozone/admin/upgrade/FinalizeSubCommand.java
index 232f3347e9e..1d0bcd981cc 100644
---
a/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/ozone/admin/upgrade/FinalizeSubCommand.java
+++
b/hadoop-ozone/cli-admin/src/main/java/org/apache/hadoop/ozone/admin/upgrade/FinalizeSubCommand.java
@@ -46,6 +46,13 @@ public class FinalizeSubCommand extends AbstractSubcommand
implements Callable<I
@CommandLine.Mixin
private OmAddressOptions.OptionalServiceIdOrHostMixin omAddressOptions;
+ @CommandLine.Option(
+ names = {"--force"},
+ description = "Skip all software version checks before finalizing.",
+ defaultValue = "false",
+ hidden = true)
+ private boolean force;
+
@CommandLine.Option(names = {"--wait"},
defaultValue = "false",
description = "After initiating finalization, poll the cluster status
until the entire cluster (OM, SCM, "
@@ -62,7 +69,14 @@ public Integer call() throws Exception {
"`ozone admin om finalizeupgrade`");
return 1;
}
- client.finalizeUpgrade();
+
+ if (force) {
+ out().println("--force specified: all software version checks will be
skipped before finalizing.");
+ client.forceFinalizeUpgrade();
+ } else {
+ client.finalizeUpgrade();
+ }
+
if (wait) {
out().println("Cluster finalization has been started. Waiting for the
cluster to finalize; "
+ "interrupt with Ctrl-C to stop waiting (finalization continues
on the server).");
diff --git
a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/ozone/admin/upgrade/TestFinalizeSubCommand.java
b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/ozone/admin/upgrade/TestFinalizeSubCommand.java
index 056ce3b3b0c..a836ad11ec7 100644
---
a/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/ozone/admin/upgrade/TestFinalizeSubCommand.java
+++
b/hadoop-ozone/cli-admin/src/test/java/org/apache/hadoop/ozone/admin/upgrade/TestFinalizeSubCommand.java
@@ -101,6 +101,18 @@ public void testCommandRunsAndPrintsOutput() throws
Exception {
String output = outContent.toString(DEFAULT_ENCODING);
assertTrue(output.contains("Cluster finalization has been started"));
verify(omClient).finalizeUpgrade();
+ verify(omClient, never()).forceFinalizeUpgrade();
+ }
+
+ @Test
+ public void testForceFlagIsPassedToClient() throws Exception {
+ new CommandLine(cmd).parseArgs("--force");
+ assertEquals(0, cmd.call());
+
+ String output = outContent.toString(DEFAULT_ENCODING);
+ assertTrue(output.contains("all software version checks will be skipped"));
+ verify(omClient).forceFinalizeUpgrade();
+ verify(omClient, never()).finalizeUpgrade();
}
@Test
@@ -132,6 +144,7 @@ public void testNonZduServerPrintsErrorAndReturnsNonZero()
throws Exception {
String errOutput = errContent.toString(DEFAULT_ENCODING);
assertTrue(errOutput.contains("OM does not support zero downtime
upgrade"));
verify(omClient, never()).finalizeUpgrade();
+ verify(omClient, never()).forceFinalizeUpgrade();
}
@Test
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/exceptions/OMException.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/exceptions/OMException.java
index 4504917e9bc..871a30372e4 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/exceptions/OMException.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/exceptions/OMException.java
@@ -238,9 +238,13 @@ public enum ResultCodes {
DIRECTORY_NOT_EMPTY,
+ @Deprecated
PERSIST_UPGRADE_TO_LAYOUT_VERSION_FAILED,
+ @Deprecated
REMOVE_UPGRADE_TO_LAYOUT_VERSION_FAILED,
+ @Deprecated
UPDATE_LAYOUT_VERSION_FAILED,
+ @Deprecated
LAYOUT_FEATURE_FINALIZATION_FAILED,
// Even though new OM servers do not support prepare for upgrade, clients
may still encounter these result codes
// when talking to an older server, so they are left in place.
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OMAdminProtocol.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OMAdminProtocol.java
index d6f9395dd37..68d2c2fafb7 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OMAdminProtocol.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OMAdminProtocol.java
@@ -19,6 +19,7 @@
import java.io.Closeable;
import java.io.IOException;
+import org.apache.hadoop.ozone.OzoneManagerVersion;
import org.apache.hadoop.ozone.om.OMConfigKeys;
import org.apache.hadoop.ozone.om.helpers.OMNodeDetails;
import org.apache.hadoop.security.KerberosInfo;
@@ -57,4 +58,11 @@ public interface OMAdminProtocol extends Closeable {
* or if the task was triggered successfully (when noWait is true)
*/
boolean triggerSnapshotDefrag(boolean noWait) throws IOException;
+
+ /**
+ * Returns the software version of this OM peer without contacting SCM.
+ * Intended for use by the OM leader to verify that all Ratis group members
+ * are running the same software version before accepting a finalize command.
+ */
+ OzoneManagerVersion getPeerUpgradeStatus() throws IOException;
}
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OzoneManagerProtocol.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OzoneManagerProtocol.java
index ee960ae385a..80438c545e6 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OzoneManagerProtocol.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocol/OzoneManagerProtocol.java
@@ -461,6 +461,9 @@ ListOpenFilesResult listOpenFiles(String path, int maxKeys,
String contToken)
boolean triggerRangerBGSync(boolean noWait) throws IOException;
/**
+ * This command is retained so that new clients can finalize old OM servers.
All new finalize requests should use
+ * `void finalizeUpgrade()`.
+ *
* Initiate metadata upgrade finalization.
* This method when called, initiates finalization of Ozone Manager metadata
* during an upgrade. The status returned contains the status
@@ -507,6 +510,15 @@ ListOpenFilesResult listOpenFiles(String path, int
maxKeys, String contToken)
*/
void finalizeUpgrade() throws IOException;
+ /**
+ * Same as {@link #finalizeUpgrade()}, but the OM skips the peer software
version check before
+ * finalizing. Use this to finalize when a peer is intentionally down or on
a different version.
+ *
+ * @throws IOException If any error occurs. If this happens finalization is
not in progress and the command must be
+ * retried.
+ */
+ void forceFinalizeUpgrade() throws IOException;
+
/**
* Returns the upgrade status of the cluster. This call is received by OM
which will in turn query SCM to get the
* status of it and the datanodes, and return the combined status to the
caller.
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OMAdminProtocolClientSideImpl.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OMAdminProtocolClientSideImpl.java
index d248e03b1bd..0b48e51b298 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OMAdminProtocolClientSideImpl.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OMAdminProtocolClientSideImpl.java
@@ -33,6 +33,7 @@
import org.apache.hadoop.ipc_.RPC;
import org.apache.hadoop.net.NetUtils;
import org.apache.hadoop.ozone.OmUtils;
+import org.apache.hadoop.ozone.OzoneManagerVersion;
import org.apache.hadoop.ozone.om.OMConfigKeys;
import org.apache.hadoop.ozone.om.exceptions.OMLeaderNotReadyException;
import org.apache.hadoop.ozone.om.exceptions.OMNotLeaderException;
@@ -44,6 +45,8 @@
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.CompactResponse;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.DecommissionOMRequest;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.DecommissionOMResponse;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.GetPeerUpgradeStatusRequest;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.GetPeerUpgradeStatusResponse;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.OMConfigurationRequest;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.OMConfigurationResponse;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.OMNodeInfo;
@@ -260,6 +263,17 @@ public boolean triggerSnapshotDefrag(boolean noWait)
throws IOException {
}
}
+ @Override
+ public OzoneManagerVersion getPeerUpgradeStatus() throws IOException {
+ try {
+ GetPeerUpgradeStatusResponse response = rpcProxy.getPeerUpgradeStatus(
+ NULL_RPC_CONTROLLER,
GetPeerUpgradeStatusRequest.newBuilder().build());
+ return OzoneManagerVersion.deserialize(response.getOmSoftwareVersion());
+ } catch (ServiceException e) {
+ throw ProtobufHelper.getRemoteException(e);
+ }
+ }
+
private void throwException(String errorMsg)
throws IOException {
throw new IOException("Request Failed. Error: " + errorMsg);
diff --git
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
index de7e77b9b3f..c17af89679d 100644
---
a/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
+++
b/hadoop-ozone/common/src/main/java/org/apache/hadoop/ozone/om/protocolPB/OzoneManagerProtocolClientSideTranslatorPB.java
@@ -2050,7 +2050,17 @@ public StatusAndMessages finalizeUpgrade(String
upgradeClientID)
@Override
public void finalizeUpgrade() throws IOException {
+ finalizeUpgrade(false);
+ }
+
+ @Override
+ public void forceFinalizeUpgrade() throws IOException {
+ finalizeUpgrade(true);
+ }
+
+ private void finalizeUpgrade(boolean force) throws IOException {
StartFinalizeUpgradeRequest req = StartFinalizeUpgradeRequest.newBuilder()
+ .setForce(force)
.build();
OMRequest omRequest = createOMRequest(Type.StartFinalizeUpgrade)
diff --git
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestOMUpgradeFinalization.java
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestOMUpgradeFinalization.java
index c2b7d0d3b73..0fdae677c61 100644
---
a/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestOMUpgradeFinalization.java
+++
b/hadoop-ozone/integration-test/src/test/java/org/apache/hadoop/ozone/om/TestOMUpgradeFinalization.java
@@ -109,7 +109,7 @@ void testFinalizationFromSnapshot() throws Exception {
// Finalize the running (active) OMs.
LOG.info("Finalizing OMs");
OzoneManagerProtocol omClient =
objectStore.getClientProxy().getOzoneManagerClient();
- omClient.finalizeUpgrade();
+ omClient.forceFinalizeUpgrade();
OMUpgradeTestUtils.waitForFinalization(omClient);
LOG.info("Finalized active OMs");
@@ -175,7 +175,7 @@ private static MiniOzoneHAClusterImpl
newCluster(OzoneConfiguration conf)
}
private static void writeKeysToIncreaseLogIndex(OzoneManagerRatisServer
omRatisServer,
- long targetLogIndex, OzoneBucket bucket) throws IOException {
+ long targetLogIndex,
OzoneBucket bucket) throws IOException {
long logIndex = omRatisServer.getLastAppliedTermIndex().getIndex();
while (logIndex < targetLogIndex) {
createKey(bucket);
diff --git a/hadoop-ozone/interface-client/src/main/proto/OMAdminProtocol.proto
b/hadoop-ozone/interface-client/src/main/proto/OMAdminProtocol.proto
index 4c9a73635bd..60e6e293847 100644
--- a/hadoop-ozone/interface-client/src/main/proto/OMAdminProtocol.proto
+++ b/hadoop-ozone/interface-client/src/main/proto/OMAdminProtocol.proto
@@ -93,6 +93,16 @@ message TriggerSnapshotDefragResponse {
optional bool result = 3;
}
+// Request from an OM leader to a peer OM to read its local upgrade status.
+// Intentionally OM-local: no SCM round-trip is performed by the server.
+message GetPeerUpgradeStatusRequest {
+}
+
+message GetPeerUpgradeStatusResponse {
+ // Serialized OzoneManagerVersion.SOFTWARE_VERSION of the responding OM
binary.
+ optional int32 omSoftwareVersion = 1;
+}
+
/**
The service for OM admin operations.
*/
@@ -113,4 +123,8 @@ service OzoneManagerAdminService {
// RPC request from admin to trigger snapshot defragmentation
rpc triggerSnapshotDefrag(TriggerSnapshotDefragRequest)
returns(TriggerSnapshotDefragResponse);
+
+ // RPC request from the OM leader to a peer OM to read its local upgrade
status.
+ rpc getPeerUpgradeStatus(GetPeerUpgradeStatusRequest)
+ returns(GetPeerUpgradeStatusResponse);
}
diff --git
a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
index fec82cdda36..121a3514bbb 100644
--- a/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
+++ b/hadoop-ozone/interface-client/src/main/proto/OmClientProtocol.proto
@@ -554,12 +554,13 @@ enum Status {
DIRECTORY_NOT_EMPTY = 68;
- PERSIST_UPGRADE_TO_LAYOUT_VERSION_FAILED = 69;
- REMOVE_UPGRADE_TO_LAYOUT_VERSION_FAILED = 70;
- UPDATE_LAYOUT_VERSION_FAILED = 71;
- LAYOUT_FEATURE_FINALIZATION_FAILED = 72;
- PREPARE_FAILED = 73; // Deprecated
- NOT_SUPPORTED_OPERATION_WHEN_PREPARED = 74; // Deprecated
+ PERSIST_UPGRADE_TO_LAYOUT_VERSION_FAILED = 69; // [deprecated = true]
+ REMOVE_UPGRADE_TO_LAYOUT_VERSION_FAILED = 70; // [deprecated = true]
+ UPDATE_LAYOUT_VERSION_FAILED = 71; // [deprecated = true]
+ LAYOUT_FEATURE_FINALIZATION_FAILED = 72; // [deprecated = true]
+ PREPARE_FAILED = 73; // [deprecated = true]
+ NOT_SUPPORTED_OPERATION_WHEN_PREPARED = 74; // [deprecated = true]
+
NOT_SUPPORTED_OPERATION_PRIOR_FINALIZATION = 75;
TENANT_NOT_FOUND = 76;
@@ -1657,6 +1658,7 @@ message FinalizeUpgradeResponse {
}
message StartFinalizeUpgradeRequest {
+ optional bool force = 1 [default = false];
}
message StartFinalizeUpgradeResponse {
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 012e4e857a9..815acb5a638 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
@@ -3650,8 +3650,7 @@ public StatusAndMessages finalizeUpgrade(String
unusedUpgradeClientId)
return FINALIZED_MSG;
}
versionManager.finalizeUpgrade();
- // OM clients currently require STARTING_MSG to be returned when this
method succeeds.
- // TODO This will be removed when OM learns to finalize from SCM.
+ // Old OM clients currently require STARTING_MSG to be returned when this
method succeeds.
return STARTING_MSG;
}
@@ -3663,6 +3662,13 @@ public void finalizeUpgrade() throws IOException {
throw new UnsupportedOperationException();
}
+ @Override
+ public void forceFinalizeUpgrade() throws IOException {
+ // Server-side stub; the real implementation is handled via the Ratis
request path through
+ // OMStartFinalizeUpgradeRequest
+ throw new UnsupportedOperationException();
+ }
+
@Override
public QueryUpgradeStatusResponse queryUpgradeStatus() throws IOException {
HddsProtos.UpgradeStatus scmStatus;
@@ -4606,7 +4612,6 @@ public String getOMServiceId() {
return omNodeDetails.getServiceId();
}
- @VisibleForTesting
public List<OMNodeDetails> getPeerNodes() {
return new ArrayList<>(peerNodesMap.values());
}
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/upgrade/OMStartFinalizeUpgradeRequest.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/upgrade/OMStartFinalizeUpgradeRequest.java
index ef2e016fbbd..29ad641a982 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/upgrade/OMStartFinalizeUpgradeRequest.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/om/request/upgrade/OMStartFinalizeUpgradeRequest.java
@@ -17,19 +17,29 @@
package org.apache.hadoop.ozone.om.request.upgrade;
+import static org.apache.hadoop.hdds.utils.HddsServerUtil.getRemoteUser;
import static
org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos.Type.StartFinalizeUpgrade;
import java.io.IOException;
+import java.util.ArrayList;
import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import org.apache.hadoop.hdds.conf.OzoneConfiguration;
+import org.apache.hadoop.hdds.scm.exceptions.SCMException;
+import org.apache.hadoop.hdds.utils.IOUtils;
import org.apache.hadoop.hdds.utils.db.cache.CacheKey;
import org.apache.hadoop.hdds.utils.db.cache.CacheValue;
import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.OzoneManagerVersion;
import org.apache.hadoop.ozone.audit.AuditLogger;
import org.apache.hadoop.ozone.audit.OMAction;
import org.apache.hadoop.ozone.om.OMMetadataManager;
import org.apache.hadoop.ozone.om.OzoneManager;
import org.apache.hadoop.ozone.om.exceptions.OMException;
import org.apache.hadoop.ozone.om.execution.flowcontrol.ExecutionContext;
+import org.apache.hadoop.ozone.om.helpers.OMNodeDetails;
+import org.apache.hadoop.ozone.om.protocolPB.OMAdminProtocolClientSideImpl;
import org.apache.hadoop.ozone.om.request.OMClientRequest;
import org.apache.hadoop.ozone.om.request.util.OmResponseUtil;
import org.apache.hadoop.ozone.om.response.OMClientResponse;
@@ -61,7 +71,21 @@ public OMRequest preExecute(OzoneManager ozoneManager)
throws IOException {
+ "Superuser privilege is required to start finalize upgrade.",
OMException.ResultCodes.ACCESS_DENIED);
}
}
- ozoneManager.getScmClient().getContainerClient().finalizeUpgrade();
+ boolean force = getOmRequest().getStartFinalizeUpgradeRequest().getForce();
+ if (force) {
+ LOG.warn("Forcing upgrade finalization by skipping OM peer software
version checks");
+ } else {
+ validatePeerOmVersionsBeforeFinalize(ozoneManager.getPeerNodes(),
ozoneManager.getConfiguration());
+ }
+
+ try {
+ ozoneManager.getScmClient().getContainerClient().finalizeUpgrade();
+ } catch (SCMException e) {
+ if (e.getResult() == SCMException.ResultCodes.UNSUPPORTED_OPERATION) {
+ throw new OMException(e.getMessage(), e,
OMException.ResultCodes.NOT_SUPPORTED_OPERATION);
+ }
+ throw e;
+ }
LOG.info("Successfully triggered the finalize upgrade process in SCM");
return omRequest;
}
@@ -94,8 +118,42 @@ public OMClientResponse validateAndUpdateCache(OzoneManager
ozoneManager, Execut
response = new
OMStartFinalizeUpgradeResponse(createErrorOMResponse(responseBuilder, e));
}
- markForAudit(auditLogger, buildAuditMessage(OMAction.UPGRADE_FINALIZE, new
HashMap<>(), exception, userInfo));
+ Map<String, String> auditMap = new HashMap<>();
+ auditMap.put("force",
String.valueOf(getOmRequest().getStartFinalizeUpgradeRequest().getForce()));
+ markForAudit(auditLogger, buildAuditMessage(OMAction.UPGRADE_FINALIZE,
auditMap, exception, userInfo));
return response;
}
+ private static void validatePeerOmVersionsBeforeFinalize(List<OMNodeDetails>
peerNodes,
+ OzoneConfiguration configuration) throws OMException {
+ if (peerNodes.isEmpty()) {
+ return;
+ }
+ OzoneManagerVersion leaderVersion = OzoneManagerVersion.SOFTWARE_VERSION;
+ List<String> failedPeers = new ArrayList<>();
+ for (OMNodeDetails peerDetails : peerNodes) {
+ String peerId = peerDetails.getNodeId();
+ OMAdminProtocolClientSideImpl client = null;
+ try {
+ client =
OMAdminProtocolClientSideImpl.createProxyForSingleOM(configuration,
getRemoteUser(), peerDetails);
+ OzoneManagerVersion peerVersion = client.getPeerUpgradeStatus();
+ if (!peerVersion.equals(leaderVersion)) {
+ LOG.warn("OM 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 OM peer {} to check software version
before finalize.", peerId, e);
+ failedPeers.add(peerId + " (unreachable: " + e.getMessage() + ")");
+ } finally {
+ IOUtils.cleanupWithLogger(LOG, client);
+ }
+ }
+ if (!failedPeers.isEmpty()) {
+ throw new OMException("Finalize rejected: the following OM peers did not
confirm matching software version "
+ + "(expected version=" + leaderVersion + "): " + String.join(", ",
failedPeers),
+ OMException.ResultCodes.NOT_SUPPORTED_OPERATION);
+ }
+ }
+
}
diff --git
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OMAdminProtocolServerSideImpl.java
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OMAdminProtocolServerSideImpl.java
index ba96368cfd5..e245846dc09 100644
---
a/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OMAdminProtocolServerSideImpl.java
+++
b/hadoop-ozone/ozone-manager/src/main/java/org/apache/hadoop/ozone/protocolPB/OMAdminProtocolServerSideImpl.java
@@ -26,6 +26,7 @@
import java.util.ArrayList;
import java.util.List;
import org.apache.hadoop.hdds.utils.db.managed.ManagedCompactRangeOptions;
+import org.apache.hadoop.ozone.OzoneManagerVersion;
import org.apache.hadoop.ozone.om.OzoneManager;
import org.apache.hadoop.ozone.om.exceptions.OMException;
import org.apache.hadoop.ozone.om.helpers.OMNodeDetails;
@@ -37,6 +38,8 @@
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.CompactResponse;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.DecommissionOMRequest;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.DecommissionOMResponse;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.GetPeerUpgradeStatusRequest;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.GetPeerUpgradeStatusResponse;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.OMConfigurationRequest;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.OMConfigurationResponse;
import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.OMNodeInfo;
@@ -154,4 +157,12 @@ public TriggerSnapshotDefragResponse triggerSnapshotDefrag(
.build();
}
}
+
+ @Override
+ public GetPeerUpgradeStatusResponse getPeerUpgradeStatus(RpcController
controller,
+ GetPeerUpgradeStatusRequest request) throws ServiceException {
+ return GetPeerUpgradeStatusResponse.newBuilder()
+ .setOmSoftwareVersion(OzoneManagerVersion.SOFTWARE_VERSION.serialize())
+ .build();
+ }
}
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 a71110a9759..7ffb7d90cf1 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
@@ -21,30 +21,47 @@
import static org.junit.jupiter.api.Assertions.assertNotEquals;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.Mockito.doNothing;
+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 java.io.IOException;
+import java.util.Arrays;
+import java.util.Collections;
+import org.apache.hadoop.hdds.scm.exceptions.SCMException;
import org.apache.hadoop.ozone.OzoneConsts;
+import org.apache.hadoop.ozone.OzoneManagerVersion;
import org.apache.hadoop.ozone.om.exceptions.OMException;
import org.apache.hadoop.ozone.om.execution.flowcontrol.ExecutionContext;
+import org.apache.hadoop.ozone.om.helpers.OMNodeDetails;
+import org.apache.hadoop.ozone.om.protocolPB.OMAdminProtocolClientSideImpl;
import org.apache.hadoop.ozone.om.request.key.OMKeyRequestTests;
import org.apache.hadoop.ozone.om.response.OMClientResponse;
import org.apache.hadoop.ozone.protocol.proto.OzoneManagerProtocolProtos;
import org.apache.hadoop.security.UserGroupInformation;
import org.apache.ratis.protocol.ClientId;
import org.apache.ratis.server.protocol.TermIndex;
+import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
/**
* Tests for OMStartFinalizeUpgradeRequest.
*/
public class TestOMStartFinalizeUpgradeRequest extends OMKeyRequestTests {
-
+
+ @BeforeEach
+ public void stubPeerNodes() {
+ when(ozoneManager.getPeerNodes()).thenReturn(Collections.emptyList());
+ }
+
@Test
public void testPreExecuteCallsScmFinalizeUpgrade() throws IOException {
doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
@@ -62,6 +79,53 @@ public void testPreExecuteCallsScmFinalizeUpgrade() throws
IOException {
verify(scmContainerLocationProtocol).finalizeUpgrade();
}
+ @Test
+ public void testScmFinalizeFailurePropagatesToClient() throws IOException {
+ IOException scmFailure = new IOException("SCM finalize upgrade failed");
+ doThrow(scmFailure).when(scmContainerLocationProtocol).finalizeUpgrade();
+
+ OMStartFinalizeUpgradeRequest request = new
OMStartFinalizeUpgradeRequest(buildRequest());
+
+ // The exception raised by SCM must propagate out of preExecute so the OM
+ // client sees the failure instead of a successful finalize.
+ IOException ex = assertThrows(IOException.class, () ->
request.preExecute(ozoneManager));
+ assertSame(scmFailure, ex);
+
+ verify(scmContainerLocationProtocol).finalizeUpgrade();
+ }
+
+ @Test
+ public void testScmUnsupportedOperationBecomesOmNotSupportedOperation()
throws IOException {
+ SCMException scmFailure =
+ new SCMException("SCM version mismatch",
SCMException.ResultCodes.UNSUPPORTED_OPERATION);
+ doThrow(scmFailure).when(scmContainerLocationProtocol).finalizeUpgrade();
+
+ OMStartFinalizeUpgradeRequest request = new
OMStartFinalizeUpgradeRequest(buildRequest());
+
+ // An SCM UNSUPPORTED_OPERATION is re-mapped to an OM
NOT_SUPPORTED_OPERATION,
+ // preserving the original message and chaining the SCM exception as the
cause.
+ OMException ex = assertThrows(OMException.class, () ->
request.preExecute(ozoneManager));
+ assertEquals(OMException.ResultCodes.NOT_SUPPORTED_OPERATION,
ex.getResult());
+ assertEquals(scmFailure.getMessage(), ex.getMessage());
+ assertSame(scmFailure, ex.getCause());
+
+ verify(scmContainerLocationProtocol).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();
+
+ OMStartFinalizeUpgradeRequest request = new
OMStartFinalizeUpgradeRequest(buildRequest());
+
+ // Only UNSUPPORTED_OPERATION is re-mapped; any other SCM exception
propagates as-is.
+ SCMException ex = assertThrows(SCMException.class, () ->
request.preExecute(ozoneManager));
+ assertSame(scmFailure, ex);
+
+ verify(scmContainerLocationProtocol).finalizeUpgrade();
+ }
+
@Test
public void testValidateAndUpdateCacheAddsFinalizationInProgressKey() throws
IOException {
doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
@@ -79,6 +143,22 @@ public void
testValidateAndUpdateCacheAddsFinalizationInProgressKey() throws IOE
"metric should be 1 after the request");
}
+ @Test
+ public void testAuditMapRecordsForceFlag() throws IOException {
+ doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+ ExecutionContext context = ExecutionContext.of(1, TermIndex.INITIAL_VALUE);
+
+ OMStartFinalizeUpgradeRequest forced = new
OMStartFinalizeUpgradeRequest(buildRequest(true));
+ forced.preExecute(ozoneManager);
+ forced.validateAndUpdateCache(ozoneManager, context);
+ assertEquals("true", forced.getAuditBuilder().getAuditMap().get("force"));
+
+ OMStartFinalizeUpgradeRequest normal = new
OMStartFinalizeUpgradeRequest(buildRequest(false));
+ normal.preExecute(ozoneManager);
+ normal.validateAndUpdateCache(ozoneManager, context);
+ assertEquals("false", normal.getAuditBuilder().getAuditMap().get("force"));
+ }
+
@Test
public void testAccessDeniedWhenUserIsNotAdmin() throws IOException {
when(ozoneManager.isAdminAuthorizationEnabled()).thenReturn(true);
@@ -103,6 +183,148 @@ public void testAccessDeniedWhenUserIsNotAdmin() throws
IOException {
verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
}
+ @Test
+ public void testPeerVersionCheckPassesWhenNoPeers() throws IOException {
+ // @BeforeEach already stubs getPeerNodes() to return an empty list.
+ // preExecute must complete normally and call SCM finalize.
+ doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+
+ OzoneManagerProtocolProtos.OMRequest original = buildRequest();
+ new OMStartFinalizeUpgradeRequest(original).preExecute(ozoneManager);
+
+ verify(scmContainerLocationProtocol).finalizeUpgrade();
+ }
+
+ @Test
+ public void testPeerVersionCheckPassesWhenAllPeersMatch() throws IOException
{
+ doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+
when(ozoneManager.getPeerNodes()).thenReturn(Arrays.asList(buildPeer("om2"),
buildPeer("om3")));
+ OMAdminProtocolClientSideImpl matchingClient =
peerClientWithVersion(OzoneManagerVersion.SOFTWARE_VERSION);
+
+ try (MockedStatic<OMAdminProtocolClientSideImpl> factory =
+ mockStatic(OMAdminProtocolClientSideImpl.class)) {
+ factory.when(() ->
OMAdminProtocolClientSideImpl.createProxyForSingleOM(any(), any(), any()))
+ .thenReturn(matchingClient);
+
+ new
OMStartFinalizeUpgradeRequest(buildRequest()).preExecute(ozoneManager);
+ }
+
+ verify(scmContainerLocationProtocol).finalizeUpgrade();
+ }
+
+ @Test
+ public void testPeerVersionCheckRejectsOneOlderPeer() throws IOException {
+
when(ozoneManager.getPeerNodes()).thenReturn(Arrays.asList(buildPeer("om2"),
buildPeer("om3")));
+ OMAdminProtocolClientSideImpl matchingClient =
peerClientWithVersion(OzoneManagerVersion.SOFTWARE_VERSION);
+ OMAdminProtocolClientSideImpl olderClient =
peerClientWithVersion(OzoneManagerVersion.HBASE_SUPPORT);
+
+ try (MockedStatic<OMAdminProtocolClientSideImpl> factory =
+ mockStatic(OMAdminProtocolClientSideImpl.class)) {
+ factory.when(() ->
OMAdminProtocolClientSideImpl.createProxyForSingleOM(any(), any(), any()))
+ .thenReturn(matchingClient, olderClient);
+
+ OMException ex = assertThrows(OMException.class,
+ () -> new
OMStartFinalizeUpgradeRequest(buildRequest()).preExecute(ozoneManager));
+ assertEquals(OMException.ResultCodes.NOT_SUPPORTED_OPERATION,
ex.getResult());
+ }
+
+ verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ }
+
+ @Test
+ public void testPeerVersionCheckRejectsOneUnknownFuturePeer() throws
IOException {
+
when(ozoneManager.getPeerNodes()).thenReturn(Arrays.asList(buildPeer("om2"),
buildPeer("om3")));
+ OMAdminProtocolClientSideImpl matchingClient =
peerClientWithVersion(OzoneManagerVersion.SOFTWARE_VERSION);
+ OMAdminProtocolClientSideImpl unknownClient =
peerClientWithVersion(OzoneManagerVersion.UNKNOWN_VERSION);
+
+ try (MockedStatic<OMAdminProtocolClientSideImpl> factory =
+ mockStatic(OMAdminProtocolClientSideImpl.class)) {
+ factory.when(() ->
OMAdminProtocolClientSideImpl.createProxyForSingleOM(any(), any(), any()))
+ .thenReturn(matchingClient, unknownClient);
+
+ OMException ex = assertThrows(OMException.class,
+ () -> new
OMStartFinalizeUpgradeRequest(buildRequest()).preExecute(ozoneManager));
+ assertEquals(OMException.ResultCodes.NOT_SUPPORTED_OPERATION,
ex.getResult());
+ }
+
+ verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ }
+
+ @Test
+ public void testPeerVersionCheckRejectsUnreachablePeer() throws IOException {
+
when(ozoneManager.getPeerNodes()).thenReturn(Collections.singletonList(buildPeer("om2")));
+ OMAdminProtocolClientSideImpl unreachableClient =
mock(OMAdminProtocolClientSideImpl.class);
+ when(unreachableClient.getPeerUpgradeStatus()).thenThrow(new
IOException("connection refused"));
+
+ try (MockedStatic<OMAdminProtocolClientSideImpl> factory =
+ mockStatic(OMAdminProtocolClientSideImpl.class)) {
+ factory.when(() ->
OMAdminProtocolClientSideImpl.createProxyForSingleOM(any(), any(), any()))
+ .thenReturn(unreachableClient);
+
+ OMException ex = assertThrows(OMException.class,
+ () -> new
OMStartFinalizeUpgradeRequest(buildRequest()).preExecute(ozoneManager));
+ assertEquals(OMException.ResultCodes.NOT_SUPPORTED_OPERATION,
ex.getResult());
+ }
+
+ verify(scmContainerLocationProtocol, never()).finalizeUpgrade();
+ }
+
+ @Test
+ public void testForceSkipsPeerVersionCheckForUnreachablePeer() throws
IOException {
+ doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+
when(ozoneManager.getPeerNodes()).thenReturn(Collections.singletonList(buildPeer("om2")));
+ OMAdminProtocolClientSideImpl unreachableClient =
mock(OMAdminProtocolClientSideImpl.class);
+ when(unreachableClient.getPeerUpgradeStatus()).thenThrow(new
IOException("connection refused"));
+
+ try (MockedStatic<OMAdminProtocolClientSideImpl> factory =
+ mockStatic(OMAdminProtocolClientSideImpl.class)) {
+ factory.when(() ->
OMAdminProtocolClientSideImpl.createProxyForSingleOM(any(), any(), any()))
+ .thenReturn(unreachableClient);
+
+ // With force=true the peer version check is skipped, so an unreachable
peer does not
+ // prevent finalization and SCM is still asked to begin finalizing.
+ new
OMStartFinalizeUpgradeRequest(buildRequest(true)).preExecute(ozoneManager);
+ }
+
+ verify(unreachableClient, never()).getPeerUpgradeStatus();
+ verify(scmContainerLocationProtocol).finalizeUpgrade();
+ }
+
+ @Test
+ public void testForceSkipsPeerVersionCheckForMismatchedPeer() throws
IOException {
+ doNothing().when(scmContainerLocationProtocol).finalizeUpgrade();
+
when(ozoneManager.getPeerNodes()).thenReturn(Collections.singletonList(buildPeer("om2")));
+ OMAdminProtocolClientSideImpl olderClient =
peerClientWithVersion(OzoneManagerVersion.HBASE_SUPPORT);
+
+ try (MockedStatic<OMAdminProtocolClientSideImpl> factory =
+ mockStatic(OMAdminProtocolClientSideImpl.class)) {
+ factory.when(() ->
OMAdminProtocolClientSideImpl.createProxyForSingleOM(any(), any(), any()))
+ .thenReturn(olderClient);
+
+ // With force=true the peer version check is skipped, so a peer running
a different software
+ // version does not prevent finalization and SCM is still asked to begin
finalizing.
+ new
OMStartFinalizeUpgradeRequest(buildRequest(true)).preExecute(ozoneManager);
+ }
+
+ verify(olderClient, never()).getPeerUpgradeStatus();
+ verify(scmContainerLocationProtocol).finalizeUpgrade();
+ }
+
+ private static OMNodeDetails buildPeer(String nodeId) {
+ return new OMNodeDetails.Builder()
+ .setOMServiceId("testService")
+ .setOMNodeId(nodeId)
+ .setHostAddress("127.0.0.1")
+ .setRpcPort(1)
+ .build();
+ }
+
+ private static OMAdminProtocolClientSideImpl
peerClientWithVersion(OzoneManagerVersion version) throws IOException {
+ OMAdminProtocolClientSideImpl client =
mock(OMAdminProtocolClientSideImpl.class);
+ when(client.getPeerUpgradeStatus()).thenReturn(version);
+ return client;
+ }
+
private OMClientResponse submitRequest() throws IOException {
OzoneManagerProtocolProtos.OMRequest original = buildRequest();
OMStartFinalizeUpgradeRequest request = new
OMStartFinalizeUpgradeRequest(original);
@@ -115,9 +337,15 @@ private OMClientResponse submitRequest() throws
IOException {
}
private OzoneManagerProtocolProtos.OMRequest buildRequest() {
+ return buildRequest(false);
+ }
+
+ private OzoneManagerProtocolProtos.OMRequest buildRequest(boolean force) {
return OzoneManagerProtocolProtos.OMRequest.newBuilder()
.setCmdType(OzoneManagerProtocolProtos.Type.StartFinalizeUpgrade)
.setClientId(ClientId.randomId().toString())
+
.setStartFinalizeUpgradeRequest(OzoneManagerProtocolProtos.StartFinalizeUpgradeRequest.newBuilder()
+ .setForce(force))
.build();
}
}
diff --git
a/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOMAdminProtocolServerSideImpl.java
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOMAdminProtocolServerSideImpl.java
new file mode 100644
index 00000000000..38d953af2e4
--- /dev/null
+++
b/hadoop-ozone/ozone-manager/src/test/java/org/apache/hadoop/ozone/protocolPB/TestOMAdminProtocolServerSideImpl.java
@@ -0,0 +1,45 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.hadoop.ozone.protocolPB;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.Mockito.mock;
+
+import com.google.protobuf.ServiceException;
+import org.apache.hadoop.ozone.OzoneManagerVersion;
+import org.apache.hadoop.ozone.om.OzoneManager;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.GetPeerUpgradeStatusRequest;
+import
org.apache.hadoop.ozone.protocol.proto.OzoneManagerAdminProtocolProtos.GetPeerUpgradeStatusResponse;
+import org.junit.jupiter.api.Test;
+
+/**
+ * Tests for {@link OMAdminProtocolServerSideImpl}.
+ */
+public class TestOMAdminProtocolServerSideImpl {
+
+ @Test
+ public void testGetPeerUpgradeStatusReturnsSoftwareVersion() throws
ServiceException {
+ OzoneManager om = mock(OzoneManager.class);
+
+ OMAdminProtocolServerSideImpl handler = new
OMAdminProtocolServerSideImpl(om);
+ GetPeerUpgradeStatusResponse response =
+ handler.getPeerUpgradeStatus(null,
GetPeerUpgradeStatusRequest.newBuilder().build());
+
+ assertEquals(OzoneManagerVersion.SOFTWARE_VERSION.serialize(),
response.getOmSoftwareVersion());
+ }
+}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]