Skip to content
Original file line number Diff line number Diff line change
Expand Up @@ -40,6 +40,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;

@Override
public Integer call() throws Exception {
try (OzoneManagerProtocol client = getClient()) {
Expand All @@ -49,7 +56,12 @@ 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();
}
out().println("Cluster finalization has been started. Monitor progress with `ozone admin upgrade status`");
}
return 0;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,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
Expand Down Expand Up @@ -117,6 +129,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();
}

private ServiceInfoEx serviceInfoWithVersion(OzoneManagerVersion version) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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");

Expand Down Expand Up @@ -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);
Expand Down
14 changes: 14 additions & 0 deletions hadoop-ozone/interface-client/src/main/proto/OMAdminProtocol.proto
Original file line number Diff line number Diff line change
Expand Up @@ -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 = 2;
}

/**
The service for OM admin operations.
*/
Expand All @@ -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);
}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -1657,6 +1658,7 @@ message FinalizeUpgradeResponse {
}

message StartFinalizeUpgradeRequest {
optional bool force = 1 [default = false];
}

message StartFinalizeUpgradeResponse {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -3649,8 +3649,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;
}
Expand All @@ -3662,6 +3661,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 = scmClient.getContainerClient().queryUpgradeStatus();
Expand Down Expand Up @@ -4590,7 +4596,6 @@ public String getOMServiceId() {
return omNodeDetails.getServiceId();
}

@VisibleForTesting
public List<OMNodeDetails> getPeerNodes() {
return new ArrayList<>(peerNodesMap.values());
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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;
}
Expand Down Expand Up @@ -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);
}
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -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();
}
}
Loading